Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -71,3 +71,4 @@ CODEBUDDY.md

# go core binary
swanlab/bin/
core/bin/
223 changes: 174 additions & 49 deletions core/cmd/swanlab-core/main.go
Original file line number Diff line number Diff line change
@@ -1,16 +1,19 @@
// Command swanlab-core 是 SwanLab Go core 的进程入口。
//
// 当前为发布脚手架阶段:仅建立进程生命周期骨架——版本上报(--version)、
// 端点监听、父进程退出监控与信号处理;gRPC 服务端(CoreService /
// CoreSyncService / ProbeService)在后续迭代中接入,届时替换 serve 循环。
// 启动与端点约定:
//
// 端点约定:
// --port-filename <路径> listen 成功后原子写入端点回报文件(SDK 发现端点依据)
// --parent-pid <pid> 监控指定的父进程 PID,父进程退出则 core 联动退出
// --listen unix:///path/to/uds Linux/macOS 本地通信(手动调试入口)
// --listen tcp://127.0.0.1:port Windows 回环地址
//
// --listen unix:///path/to/uds Linux/macOS 进程内通信(默认路径由 Python SDK 分配)
// --listen tcp://127.0.0.1:port Windows 回环地址(uds 不可用,named pipe 支持后续提供)
// 未显式指定 --listen 时按平台自选端点:
// POSIX 优先使用 port-filename 同目录下的 core.sock(UDS),listen 失败时记录 warning 并回退至 127.0.0.1 随机回环端口;
// Windows 默认直接使用随机回环端口。
//
// 生命周期:父进程退出(process 包监控)或收到 SIGINT/SIGTERM 时优雅退出,
// 防止 Python SDK 崩溃后 core 沦为孤儿进程。
// 生命周期管理:
// 汇集 Teardown RPC、系统信号、父进程退出监控或 Serve 异常,统一触发 Controller 的
// 关闭序列(优先 GracefulStop,超时强制 Stop);退出时清理自己创建的 socket 与 port-file 文件。
package main

import (
Expand All @@ -21,17 +24,25 @@ import (
"net"
"os"
"os/signal"
"path/filepath"
"runtime"
"strconv"
"strings"
"syscall"
"time"

"google.golang.org/grpc"

"github.com/swanhubx/swanlab/core/internal/manager"
"github.com/swanhubx/swanlab/core/internal/pkg/console"
"github.com/swanhubx/swanlab/core/internal/pkg/portinfo"
"github.com/swanhubx/swanlab/core/internal/pkg/process"
"github.com/swanhubx/swanlab/core/internal/server"
"github.com/swanhubx/swanlab/core/internal/service"
)

// version 与 commit 由构建管线通过 -ldflags -X 注入(见 core/hatch.py),
// 缺省值仅供本地 go run / go build 使用。
// 缺省值供本地 go run / go build 使用。
var (
version = "dev"
commit = "unknown"
Expand All @@ -49,17 +60,30 @@ const (
exitRunError = 1
)

// 自选端点与收尾参数。
const (
coreSocketName = "core.sock"
loopbackAddr = "127.0.0.1:0"
shutdownGrace = 10 * time.Second
)

func main() {
os.Exit(run(os.Args[1:]))
}

func run(args []string) int {
fs := flag.NewFlagSet("swanlab-core", flag.ContinueOnError)
printVersion := fs.Bool("version", false, "打印版本信息后退出")
printVersion := fs.Bool("version", false, "print version and exit")
listenAddr := fs.String("listen", os.Getenv(envListenAddr),
"监听端点,格式 unix://<uds路径> 或 tcp://<地址:端口>;Windows 仅支持 tcp:// 回环地址")
"listen endpoint, unix://<uds path> or tcp://<addr:port>; auto-selected per platform when unset")
portFilename := fs.String("port-filename", "",
"endpoint report file; atomically written once listen succeeds, for callers to poll")
parentPID := fs.Int("parent-pid", envInt(envParentPID),
"预期父进程 PID,父进程退出时 core 随之退出;未指定时取启动瞬间的实际父进程")
"expected parent PID; core exits when the parent exits, defaults to the actual parent at startup")
detach := fs.Bool("detach", false,
"detach from parent process (reserved; accepted but ignored in this build)")
idleTimeout := fs.Duration("idle-timeout", 0,
"detached idle timeout (reserved; accepted but ignored in this build)")
if err := fs.Parse(args); err != nil {
if errors.Is(err, flag.ErrHelp) {
return 0
Expand All @@ -71,87 +95,188 @@ func run(args []string) int {
fmt.Printf("swanlab-core %s (commit %s)\n", version, commit)
return 0
}
if *listenAddr == "" {
console.Error("未指定监听端点:通过 --listen 或环境变量 " + envListenAddr + " 传入")
if *detach || *idleTimeout != 0 {
console.Warning("--detach/--idle-timeout accepted but ignored: detached service is not implemented; core stays bound to the parent process")
}
if *listenAddr == "" && *portFilename == "" {
console.Error("no listen endpoint: pass --listen (manual debug) or --port-filename (SDK startup convention)")
return exitUsageError
}

ln, err := listen(*listenAddr)
// 自建资源记录,退出时清理自己创建的部分。
var socketPath string
selfPID := os.Getpid() // 退出时凭它确认 port-file 属本实例
wrotePortFile := false
defer func() {
cleanupSocket(socketPath)
if wrotePortFile {
cleanupPortFile(*portFilename, selfPID)
}
}()

ln, err := openEndpoint(*listenAddr, *portFilename)
if err != nil {
console.Error("监听失败:", err)
console.Error("listen failed:", err)
return exitRunError
}
defer ln.Close()
defer func() { _ = ln.Close() }()
if addr, ok := ln.Addr().(*net.UnixAddr); ok && !strings.HasPrefix(addr.Name, "@") {
socketPath = addr.Name
}

// 父进程监控:显式传入的 PID 优先(Python SDK 启动约定),未传时回退为
// 监控启动瞬间的实际父进程(本地终端运行场景)。监控建立失败按约定终止启动。
// 父进程监控:显式传入的 PID 生效(启动约定);未传时监控启动瞬间的
// 实际父进程(本地终端运行场景)。监控建立失败终止启动。
pid := *parentPID
if pid <= 0 {
pid = os.Getppid()
}
parentExited, err := process.NotifyOnParentExit(pid)
if err != nil {
console.Error("父进程监控建立失败,终止启动:", err)
console.Error("failed to watch parent process, aborting startup:", err)
return exitRunError
}

// 装配服务核心分层:
// - server.Controller: 负责 gRPC 进程宿主生命周期与优雅退出控制;
// - manager.Manager: 负责 Run 领域会话管理与路由注册;
// - service.CoreService: 负责 gRPC 请求接入与协议映射。
grpcServer := grpc.NewServer()
ctrl := server.NewController(grpcServer, shutdownGrace)
service.NewCoreService(ctrl, manager.New()).Register(grpcServer)

// listen 与 server 初始化成功后写 port-file。
if *portFilename != "" {
info := portinfo.Info{Protocol: portinfo.ProtocolVersion, PID: selfPID}
switch addr := ln.Addr().(type) {
case *net.UnixAddr:
info.UnixPath = addr.Name
case *net.TCPAddr:
info.SockPort = addr.Port
default:
console.Error("unrecognized listener address type:", ln.Addr())
return exitRunError
}
if err2 := portinfo.WriteFile(*portFilename, &info); err2 != nil {
console.Error("failed to write port-file:", err2)
return exitRunError
}
wrotePortFile = true
}

ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()

console.Infof("swanlab-core %s listening on %s (parent pid %d)", version, *listenAddr, pid)
console.Infof("swanlab-core %s listening on %s (parent pid %d)", version, ln.Addr(), pid)

serveErr := make(chan error, 1)
go func() {
serveErr <- serve(ln)
serveErr <- grpcServer.Serve(ln)
}()

exitCode := 0
var cause string
var serveFailure error
served := false
select {
case <-ctx.Done():
console.Info("收到退出信号,正在关闭")
console.Info("shutdown signal received, stopping")
cause = "signal"
case <-parentExited:
console.Warning("父进程已退出,core 随之退出")
case err := <-serveErr:
if err != nil {
console.Error("监听异常退出:", err)
return exitRunError
}
console.Warning("parent process exited, stopping core")
cause = "parent-exit"
case serveFailure = <-serveErr:
served = true
cause = "serve-error"
}
ctrl.Shutdown(cause)
// Shutdown 的关闭序列会让 Serve 返回(GracefulStop 超时转 Stop);
// 上面的 select 未消费 serveErr 时,在此收取其结果。
if !served {
serveFailure = <-serveErr
}
<-ctrl.Done()
if serveFailure != nil {
console.Error("gRPC Serve exited with error:", serveFailure)
exitCode = exitRunError
}
return 0
return exitCode
}

// listen 按协议前缀创建监听器。uds 仅在非 Windows 平台可用;Windows 使用
// TCP 回环地址兜底(named pipe 接入后在此分支扩展)。
func listen(addr string) (net.Listener, error) {
scheme, rest, ok := strings.Cut(addr, "://")
// openEndpoint 创建监听器。显式 --listen 生效(手动调试,失败不回退,
// tcp:// 限定回环);未传时按平台自选:POSIX 尝试 port-filename 同目录
// 下的 UDS,失败记录 warning 后回退随机回环 TCP;Windows 使用随机回环端口。
func openEndpoint(listenAddr, portFilename string) (net.Listener, error) {
if listenAddr == "" {
if runtime.GOOS == "windows" {
return net.Listen("tcp", loopbackAddr)
}
sockPath := filepath.Join(filepath.Dir(portFilename), coreSocketName)
ln, err := net.Listen("unix", sockPath)
if err != nil {
console.Warningf("unix listen on %s failed: %v; falling back to loopback tcp %s", sockPath, err, loopbackAddr)
return net.Listen("tcp", loopbackAddr)
}
return ln, nil
}
scheme, rest, ok := strings.Cut(listenAddr, "://")
if !ok {
return nil, fmt.Errorf("监听端点缺少协议前缀(unix:// 或 tcp://): %s", addr)
return nil, fmt.Errorf("listen endpoint missing scheme prefix (unix:// or tcp://): %s", listenAddr)
}
switch scheme {
case "unix":
if runtime.GOOS == "windows" {
return nil, errors.New("windows 平台不支持 unix:// 端点,请使用 tcp://127.0.0.1:<端口>")
return nil, errors.New("unix:// endpoints are not supported on Windows; use tcp://127.0.0.1:<port>")
}
return net.Listen("unix", rest)
case "tcp":
if err := requireLoopbackHost(rest); err != nil {
return nil, err
}
return net.Listen("tcp", rest)
default:
return nil, fmt.Errorf("不支持的监听协议 %q(仅 unix:// 或 tcp://)", scheme)
return nil, fmt.Errorf("unsupported listen scheme %q (only unix:// or tcp://)", scheme)
}
}

// serve 接受连接后立即关闭。脚手架阶段仅验证端点连通性,
// gRPC 服务端就绪后由此接入 serve 逻辑。
func serve(ln net.Listener) error {
for {
conn, err := ln.Accept()
if err != nil {
if errors.Is(err, net.ErrClosed) {
return nil
}
return err
}
_ = conn.Close()
// cleanupSocket 删除自己创建的 UDS socket 文件;Linux 抽象 socket(@ 前缀)
// 不占文件系统,无需清理。
func cleanupSocket(path string) {
if path == "" || strings.HasPrefix(path, "@") {
return
}
_ = os.Remove(path)
}

// cleanupPortFile 在 pid 匹配时删除 port-file。同一路径被后启动实例
// 覆盖后,本实例退出不得误删他者文件;文件缺失、损坏或 pid 不匹配时不删。
func cleanupPortFile(path string, pid int) {
if path == "" || pid <= 0 {
return
}
info, err := portinfo.ParseFile(path)
if err != nil {
return
}
if info.PID != pid {
return
}
_ = os.Remove(path)
}

// requireLoopbackHost 限制显式 tcp:// 监听地址为回环,避免服务暴露到
// 网络;放行 127.0.0.1/::1/localhost。
func requireLoopbackHost(addr string) error {
host, _, err := net.SplitHostPort(addr)
if err != nil {
return fmt.Errorf("parse tcp listen address %q: %w", addr, err)
}
if host == "localhost" {
return nil
}
if ip := net.ParseIP(host); ip != nil && ip.IsLoopback() {
return nil
}
return fmt.Errorf("tcp listen address %q is not loopback; use 127.0.0.1 or ::1", addr)
}

// envInt 解析整型环境变量,缺失或非法时返回 0。
Expand Down
Loading
Loading