diff --git a/docs/runner-protocol-roadmap.md b/docs/runner-protocol-roadmap.md index 23285d9..5b66a88 100644 --- a/docs/runner-protocol-roadmap.md +++ b/docs/runner-protocol-roadmap.md @@ -18,6 +18,10 @@ dynamic-runner scheduler └──────────────────────► Gitea ``` +controller 使用单一 Go 二进制;默认在同一进程启用 `scheduler`、`pod-worker` 和 +`vm-worker`,也可通过 `--components` 只启用其中一部分。组件是独立应用服务边界, +共享进程不意味着共享后端状态或把 assignment 降级为内存 channel。 + 这与“收到 webhook 后临时注册另一个 act_runner”不同。`FetchTask` 已经完成任务分配, 不能再期待 Gitea 把同一个 task 分配给随后启动的 runner。协议调度器必须让 executor 执行已经领取的 task,并继续完成日志、状态、心跳、取消和最终结果上报。 diff --git a/internal/controller/components.go b/internal/controller/components.go new file mode 100644 index 0000000..1fcdf48 --- /dev/null +++ b/internal/controller/components.go @@ -0,0 +1,75 @@ +// Package controller composes independently runnable scheduler and backend workers. +package controller + +import ( + "context" + "errors" + "fmt" + "slices" + "strings" + + "golang.org/x/sync/errgroup" +) + +type ComponentName string + +const ( + Scheduler ComponentName = "scheduler" + PodWorker ComponentName = "pod-worker" + VMWorker ComponentName = "vm-worker" +) + +var defaultComponents = []ComponentName{Scheduler, PodWorker, VMWorker} + +// Selection parses --components. An empty value enables all components. +type Selection []ComponentName + +func ParseSelection(value string) (Selection, error) { + if strings.TrimSpace(value) == "" || strings.TrimSpace(value) == "all" { + return append(Selection(nil), defaultComponents...), nil + } + var selected Selection + for _, raw := range strings.Split(value, ",") { + name := ComponentName(strings.TrimSpace(raw)) + if !slices.Contains(defaultComponents, name) { + return nil, fmt.Errorf("unknown controller component %q", name) + } + if !slices.Contains(selected, name) { + selected = append(selected, name) + } + } + if len(selected) == 0 { + return nil, errors.New("at least one controller component is required") + } + return selected, nil +} + +type Component interface { + Run(context.Context) error +} + +type Registry map[ComponentName]Component + +// Run starts exactly the selected components in one process. The first real +// failure cancels its peers; ordinary context cancellation is graceful. +func Run(ctx context.Context, selection Selection, registry Registry) error { + group, groupContext := errgroup.WithContext(ctx) + for _, name := range selection { + component, ok := registry[name] + if !ok || component == nil { + return fmt.Errorf("component %q is not configured", name) + } + name, component := name, component + group.Go(func() error { + err := component.Run(groupContext) + if errors.Is(err, context.Canceled) && groupContext.Err() != nil { + return nil + } + if err != nil { + return fmt.Errorf("component %s: %w", name, err) + } + return nil + }) + } + return group.Wait() +} diff --git a/internal/controller/components_test.go b/internal/controller/components_test.go new file mode 100644 index 0000000..a5cce8c --- /dev/null +++ b/internal/controller/components_test.go @@ -0,0 +1,63 @@ +package controller + +import ( + "context" + "errors" + "sync" + "testing" +) + +func TestParseSelectionDefaultsToAll(t *testing.T) { + for _, input := range []string{"", "all"} { + selection, err := ParseSelection(input) + if err != nil { + t.Fatal(err) + } + if len(selection) != 3 || selection[0] != Scheduler || selection[1] != PodWorker || selection[2] != VMWorker { + t.Fatalf("selection = %v", selection) + } + } +} + +func TestParseSelectionAllowsOneOrMoreComponents(t *testing.T) { + selection, err := ParseSelection("vm-worker,scheduler,vm-worker") + if err != nil { + t.Fatal(err) + } + if len(selection) != 2 || selection[0] != VMWorker || selection[1] != Scheduler { + t.Fatalf("selection = %v", selection) + } + if _, err := ParseSelection("webhook"); err == nil { + t.Fatal("expected obsolete component to be rejected") + } +} + +type componentFunc func(context.Context) error + +func (f componentFunc) Run(ctx context.Context) error { return f(ctx) } + +func TestRunStartsSelectedComponentsAndCancelsPeers(t *testing.T) { + started := make(chan ComponentName, 2) + peerStopped := make(chan struct{}) + var once sync.Once + registry := Registry{ + Scheduler: componentFunc(func(context.Context) error { + started <- Scheduler + return errors.New("poll failed") + }), + PodWorker: componentFunc(func(ctx context.Context) error { + started <- PodWorker + <-ctx.Done() + once.Do(func() { close(peerStopped) }) + return ctx.Err() + }), + } + err := Run(context.Background(), Selection{Scheduler, PodWorker}, registry) + if err == nil || !errors.Is(err, context.Canceled) && err.Error() != "component scheduler: poll failed" { + t.Fatalf("Run() error = %v", err) + } + <-peerStopped + if len(started) != 2 { + t.Fatalf("started components = %d", len(started)) + } +}