From d06845bc3c5f7063d51fbfe55d9a2fa242ae6971 Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Mon, 21 Sep 2026 06:44:34 +0000 Subject: [PATCH] =?UTF-8?q?fix:=20=E9=9A=94=E7=A6=BB=20Pod=20=E4=B8=8E=20V?= =?UTF-8?q?M=20=E5=B9=B6=E5=8F=91=E5=AE=B9=E9=87=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/gitea-dynamic-runner/controller.go | 4 +--- docs/runner-protocol-roadmap.md | 6 ++--- internal/taskscheduler/gate.go | 31 -------------------------- internal/taskscheduler/gate_test.go | 29 ------------------------ internal/taskscheduler/poller.go | 13 ----------- 5 files changed, 4 insertions(+), 79 deletions(-) delete mode 100644 internal/taskscheduler/gate.go delete mode 100644 internal/taskscheduler/gate_test.go diff --git a/cmd/gitea-dynamic-runner/controller.go b/cmd/gitea-dynamic-runner/controller.go index 252382c..ff74ff4 100644 --- a/cmd/gitea-dynamic-runner/controller.go +++ b/cmd/gitea-dynamic-runner/controller.go @@ -91,7 +91,6 @@ func runController(ctx context.Context) error { return err } registry := runnerfacade.NewRegistry() - gate := taskscheduler.NewSingleFlightGate() var podExecutorBackend *podbackend.Backend var vmExecutorBackend *opensandboxbackend.Backend giteaClient := giteaactions.NewClient(giteaactions.DefaultHTTPClient(), config.GiteaURL, config.GiteaUUID, config.GiteaToken) @@ -114,7 +113,6 @@ func runController(ctx context.Context) error { return err } } - gate.Release() return nil }, } @@ -131,7 +129,7 @@ func runController(ctx context.Context) error { labels = append(labels, string(taskassignment.BackendVM)) } poller := taskscheduler.Poller{ - Client: giteaClient, Gate: gate, + Client: giteaClient, Scheduler: &taskscheduler.Scheduler{TrustDomain: config.TrustDomain, Dispatcher: assignmentqueue.Publisher{ JetStream: producerJS, SubjectBase: config.SubjectBase, }}, diff --git a/docs/runner-protocol-roadmap.md b/docs/runner-protocol-roadmap.md index 01dc6a2..27da1bf 100644 --- a/docs/runner-protocol-roadmap.md +++ b/docs/runner-protocol-roadmap.md @@ -80,9 +80,9 @@ facade,并严格校验 facade 的 SPIFFE ID。这样无需修改 runner 或把 - consumer 在 executor 使用上述 facade 成功 claim task 后确认 assignment;无需把完整 task 写入 Pod annotation、OpenSandbox metadata 或环境变量。 - Pod 与 VM 共享 task/executor 协议,只有环境创建和销毁实现不同。 -- 首次生产 canary 使用 controller 进程内 single-flight gate:只有官方 runner 的终态 - `UpdateTask` 已被 Gitea 接受后才允许 Fetch 下一条任务。它把未知故障收敛为停止领取, - 而不是在 backlog 下连续创建 executor;后续容量调度必须以 backend 实际运行资源为准。 +- scheduler 在 assignment 持久化到 JetStream 后即可继续领取;Pod 与 VM 分别由 durable + consumer 的 capacity 限制并发,不共享全局执行槽位。未知后端故障由对应 consumer 的 + NAK/redelivery 收敛,不能阻塞另一种 backend。 - 两种 backend 都注入同一份 runner bootstrap 环境;Pod 仍由 homelab Kubernetes 原生 创建,只有 VM 经 OpenSandbox 创建,bootstrap 机制不改变 backend 边界。 diff --git a/internal/taskscheduler/gate.go b/internal/taskscheduler/gate.go deleted file mode 100644 index b88c2d7..0000000 --- a/internal/taskscheduler/gate.go +++ /dev/null @@ -1,31 +0,0 @@ -package taskscheduler - -import "context" - -// SingleFlightGate keeps at most one fetched task in flight. Release is -// idempotent so repeated terminal updates cannot increase capacity. -type SingleFlightGate struct { - token chan struct{} -} - -func NewSingleFlightGate() *SingleFlightGate { - gate := &SingleFlightGate{token: make(chan struct{}, 1)} - gate.token <- struct{}{} - return gate -} - -func (g *SingleFlightGate) Acquire(ctx context.Context) error { - select { - case <-g.token: - return nil - case <-ctx.Done(): - return ctx.Err() - } -} - -func (g *SingleFlightGate) Release() { - select { - case g.token <- struct{}{}: - default: - } -} diff --git a/internal/taskscheduler/gate_test.go b/internal/taskscheduler/gate_test.go deleted file mode 100644 index 25815d0..0000000 --- a/internal/taskscheduler/gate_test.go +++ /dev/null @@ -1,29 +0,0 @@ -package taskscheduler - -import ( - "context" - "testing" - "time" -) - -func TestSingleFlightGateBlocksUntilTerminalRelease(t *testing.T) { - gate := NewSingleFlightGate() - if err := gate.Acquire(context.Background()); err != nil { - t.Fatal(err) - } - blocked, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) - defer cancel() - if err := gate.Acquire(blocked); err == nil { - t.Fatal("second task acquired capacity before release") - } - gate.Release() - gate.Release() - if err := gate.Acquire(context.Background()); err != nil { - t.Fatal(err) - } - blockedAgain, cancelAgain := context.WithTimeout(context.Background(), 20*time.Millisecond) - defer cancelAgain() - if err := gate.Acquire(blockedAgain); err == nil { - t.Fatal("duplicate terminal release increased capacity") - } -} diff --git a/internal/taskscheduler/poller.go b/internal/taskscheduler/poller.go index 6447e73..b03faaf 100644 --- a/internal/taskscheduler/poller.go +++ b/internal/taskscheduler/poller.go @@ -26,7 +26,6 @@ type PollerConfig struct { type Poller struct { Client PollClient Scheduler *Scheduler - Gate *SingleFlightGate Config PollerConfig OnError func(error) } @@ -51,14 +50,7 @@ func (p Poller) Run(ctx context.Context) error { errorBackoff = 5 * time.Second } var tasksVersion int64 - haveLease := false for { - if p.Gate != nil && !haveLease { - if err := p.Gate.Acquire(ctx); err != nil { - return nil - } - haveLease = true - } response, err := p.Client.FetchTask(ctx, tasksVersion) if err != nil { if ctx.Err() != nil { @@ -80,10 +72,6 @@ func (p Poller) Run(ctx context.Context) error { tasksVersion = response.GetTasksVersion() task := response.GetTask() if task == nil { - if p.Gate != nil { - p.Gate.Release() - haveLease = false - } if !wait(ctx, emptyBackoff) { return nil } @@ -91,7 +79,6 @@ func (p Poller) Run(ctx context.Context) error { } for { if err := p.Scheduler.Run(ctx, task); err == nil { - haveLease = false break } else { if ctx.Err() != nil {