From e49d1818a5fe64d4d34d210e187943900e2fcdff Mon Sep 17 00:00:00 2001 From: dawn Date: Thu, 27 Aug 2026 21:34:26 +0900 Subject: [PATCH] spindle/{config,mill}: dont quarantine executor if lease cancel was ACKed but not done yet after timeout Signed-off-by: dawn --- spindle/config/config.go | 26 ++++++----- spindle/mill/lease.go | 11 +++++ spindle/mill/mill.go | 90 +++++++++++++++++++++++++------------ spindle/mill/mill_test.go | 93 ++++++++++++++++++++++++++++++++++++--- spindle/server.go | 8 ++-- 5 files changed, 179 insertions(+), 49 deletions(-) diff --git a/spindle/config/config.go b/spindle/config/config.go index d79a86a01..f183c7412 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -162,18 +162,20 @@ const ( // fields are selectively active depending on the role type Mill struct { - URL string `env:"URL"` // mill websocket endpoint dialled by the executor - SharedSecret string `env:"SHARED_SECRET"` // the executor's token for dialing the mill - MaxPending int `env:"MAX_PENDING, default=100"` // mill pending job queue limit - ReconnectGrace time.Duration `env:"RECONNECT_GRACE, default=45s"` // reconnect window before leases are failed - DrainTimeout time.Duration `env:"DRAIN_TIMEOUT, default=20m"` - Seats int `env:"SEATS, default=4"` // executor seats advertised to the mill - Labels []string `env:"LABELS"` // executor capability labels - ArtifactStore string `env:"ARTIFACT_STORE"` // store shared by mill and its executors - JumpListenAddr string `env:"JUMP_LISTEN_ADDR"` - JumpHostKeyPath string `env:"JUMP_HOST_KEY_PATH"` - DebugExecutorPort uint32 `env:"DEBUG_EXECUTOR_PORT, default=2223"` - MaxJumpConnections int `env:"MAX_JUMP_CONNECTIONS, default=128"` + URL string `env:"URL"` // mill websocket endpoint dialled by the executor + SharedSecret string `env:"SHARED_SECRET"` // the executor's token for dialing the mill + MaxPending int `env:"MAX_PENDING, default=100"` // mill pending job queue limit + ReconnectGrace time.Duration `env:"RECONNECT_GRACE, default=45s"` // reconnect window before leases are failed + CancelAckTimeout time.Duration `env:"CANCEL_ACK_TIMEOUT, default=10s"` + CancelTeardownTimeout time.Duration `env:"CANCEL_TEARDOWN_TIMEOUT, default=10m"` + DrainTimeout time.Duration `env:"DRAIN_TIMEOUT, default=20m"` + Seats int `env:"SEATS, default=4"` // executor seats advertised to the mill + Labels []string `env:"LABELS"` // executor capability labels + ArtifactStore string `env:"ARTIFACT_STORE"` // store shared by mill and its executors + JumpListenAddr string `env:"JUMP_LISTEN_ADDR"` + JumpHostKeyPath string `env:"JUMP_HOST_KEY_PATH"` + DebugExecutorPort uint32 `env:"DEBUG_EXECUTOR_PORT, default=2223"` + MaxJumpConnections int `env:"MAX_JUMP_CONNECTIONS, default=128"` } type Tracing struct { Endpoint string `env:"ENDPOINT"` diff --git a/spindle/mill/lease.go b/spindle/mill/lease.go index b63fd9565..4f211a8bd 100644 --- a/spindle/mill/lease.go +++ b/spindle/mill/lease.go @@ -143,6 +143,17 @@ func (l *RemoteLease) requestCancel(reason string) cancelAction { return cancelRemote } +// reconnect can replay the cancel, only the first ack starts the teardown timer +func (l *RemoteLease) markCancelAcked() bool { + l.mu.Lock() + defer l.mu.Unlock() + if l.cancelAcked { + return false + } + l.cancelAcked = true + return true +} + func (l *RemoteLease) cancelRequested() (bool, string) { l.mu.Lock() defer l.mu.Unlock() diff --git a/spindle/mill/mill.go b/spindle/mill/mill.go index 2d720a131..930c794cf 100644 --- a/spindle/mill/mill.go +++ b/spindle/mill/mill.go @@ -30,12 +30,14 @@ import ( ) const ( - defaultReconnectGrace = 45 * time.Second - defaultJobTimeout = 24 * time.Hour - defaultBidTimeout = 5 * time.Second - defaultTopK = 3 - defaultMaxPending = 100 - defaultQuarantineStrikes = 3 + defaultReconnectGrace = 45 * time.Second + defaultJobTimeout = 24 * time.Hour + defaultBidTimeout = 5 * time.Second + defaultTopK = 3 + defaultMaxPending = 100 + defaultQuarantineStrikes = 3 + defaultCancelAckTimeout = 10 * time.Second + defaultCancelTeardownTimeout = 10 * time.Minute ) // repeated protocol errors quarantine the node @@ -47,13 +49,15 @@ func protoErrf(format string, args ...any) error { type Config struct { // mill appends live-tailed executor lines here so logview can follow running remote jobs - LogDir string - MaxPending int - ReconnectGrace time.Duration - JobTimeout time.Duration - BidTimeout time.Duration - TopK int - CancelTimeout time.Duration + LogDir string + MaxPending int + ReconnectGrace time.Duration + JobTimeout time.Duration + BidTimeout time.Duration + TopK int + CancelAckTimeout time.Duration + // teardown can take minutes after the cancel is acknowledged + CancelTeardownTimeout time.Duration // how many protocol-violating session deaths in a row quarantine the node QuarantineStrikes int } @@ -96,8 +100,11 @@ func New(l *slog.Logger, cfg Config) *Mill { if cfg.MaxPending <= 0 { cfg.MaxPending = defaultMaxPending } - if cfg.CancelTimeout <= 0 { - cfg.CancelTimeout = 10 * time.Second + if cfg.CancelAckTimeout <= 0 { + cfg.CancelAckTimeout = defaultCancelAckTimeout + } + if cfg.CancelTeardownTimeout <= 0 { + cfg.CancelTeardownTimeout = defaultCancelTeardownTimeout } if cfg.QuarantineStrikes <= 0 { cfg.QuarantineStrikes = defaultQuarantineStrikes @@ -1126,7 +1133,7 @@ func (m *Mill) sendCancel(sess *millSession, lease *RemoteLease, reason string) }}); err != nil { return } - time.AfterFunc(m.cfg.CancelTimeout, func() { m.checkCancelDeadline(lease) }) + time.AfterFunc(m.cfg.CancelAckTimeout, func() { m.checkCancelAck(lease) }) } func (m *Mill) sessionForNode(nodeID string) *millSession { @@ -1438,35 +1445,62 @@ func (m *Mill) onCancelAck(sess *millSession, ca *millv1.CancelAck) { if lease == nil || lease.nodeID != sess.nodeID { return } - lease.mu.Lock() - lease.cancelAcked = true - lease.mu.Unlock() + if !lease.markCancelAcked() { + return + } + time.AfterFunc(m.cfg.CancelTeardownTimeout, func() { m.checkCancelTeardown(lease) }) } -func (m *Mill) checkCancelDeadline(lease *RemoteLease) { +func (m *Mill) checkCancelAck(lease *RemoteLease) { lease.mu.Lock() isDone := (lease.state == leaseDone) isAcked := lease.cancelAcked lease.mu.Unlock() - if isDone && isAcked { + if isDone || isAcked { return } - m.l.Warn("node failed to comply with cancel request within deadline, quarantining", "node", lease.nodeID, "lease", lease.id, "done", isDone, "acked", isAcked) + m.mu.Lock() + sess := m.sessions[lease.nodeID] + disconnected := sess == nil || sess.disconnected + m.mu.Unlock() + if disconnected { + // reconnect replays the cancel, the grace timer handles nodes that never return + return + } - reason := fmt.Sprintf("cancel noncompliance for lease %s (done: %t, acked: %t)", lease.id, isDone, isAcked) + m.l.Warn("node did not acknowledge cancel request within deadline, quarantining", "node", lease.nodeID, "lease", lease.id, "deadline", m.cfg.CancelAckTimeout) + + reason := fmt.Sprintf("cancel unacknowledged for lease %s after %s", lease.id, m.cfg.CancelAckTimeout) if m.db != nil { if err := m.db.QuarantineExecutor(lease.nodeID, reason); err != nil { m.l.Error("failed to quarantine executor", "node", lease.nodeID, "err", err) } } - m.mu.Lock() - sess := m.sessions[lease.nodeID] - m.mu.Unlock() - if sess != nil { - sess.close() + sess.close() +} + +// don't fail the node's other jobs because one teardown got stuck +func (m *Mill) checkCancelTeardown(lease *RemoteLease) { + if lease.getState() == leaseDone { + return + } + + _, cancelReason := lease.cancelRequested() + reason := fmt.Sprintf("%s (node teardown exceeded %s)", cancelReason, m.cfg.CancelTeardownTimeout) + m.l.Warn("acked cancel exceeded teardown deadline; finishing lease without the node's terminal", "node", lease.nodeID, "lease", lease.id, "deadline", m.cfg.CancelTeardownTimeout) + + status := string(models.StatusKindCancelled) + if lease.orphaned { + if err := m.finishOrphan(lease, status, &reason, nil); err != nil { + m.l.Error("finish orphan after cancel teardown deadline failed", "lease", lease.id, "err", err) + } + return + } + if err := m.finishLiveLease(lease, status, reason); err != nil { + m.l.Error("finish lease after cancel teardown deadline failed", "lease", lease.id, "err", err) } } diff --git a/spindle/mill/mill_test.go b/spindle/mill/mill_test.go index 928a50c08..4e4034daa 100644 --- a/spindle/mill/mill_test.go +++ b/spindle/mill/mill_test.go @@ -783,13 +783,47 @@ func TestClaimedSweep(t *testing.T) { } } -func TestCancelDeadline(t *testing.T) { - m, _ := restoreTestMill(t, Config{CancelTimeout: 50 * time.Millisecond}) +func TestUnacknowledgedCancelQuarantinesNode(t *testing.T) { + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: 50 * time.Millisecond}) + registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil) sessClosed := make(chan struct{}) - sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error { + sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) + sess.closeTransport = func() error { + close(sessClosed) return nil - }), discardLogger()) + } + m.attachSession(sess) + + lease := newLease("lease-1", "node-1", "inc-1", "dummy") + lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.setState(leaseRunning) + m.mu.Lock() + m.leases[lease.id] = lease + m.mu.Unlock() + + m.destroy(lease.wid) + + select { + case <-sessClosed: + case <-time.After(2 * time.Second): + t.Fatal("session was not closed after the cancel went unacknowledged") + } + if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || ok { + t.Fatalf("executor still admitted after ignoring a cancel: ok = %v, err = %v", ok, err) + } +} + +func TestAckedCancelFinishesOnlyItsLeaseAfterTeardownDeadline(t *testing.T) { + m, bdb := restoreTestMill(t, Config{ + ReconnectGrace: time.Minute, + CancelAckTimeout: 30 * time.Millisecond, + CancelTeardownTimeout: 150 * time.Millisecond, + }) + registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil) + + sessClosed := make(chan struct{}) + sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) sess.closeTransport = func() error { close(sessClosed) return nil @@ -799,16 +833,63 @@ func TestCancelDeadline(t *testing.T) { lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) + if err := m.persistLease(lease, leaseRowRunning); err != nil { + t.Fatalf("persistLease: %v", err) + } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock() m.destroy(lease.wid) + m.onCancelAck(sess, &millv1.CancelAck{LeaseId: lease.id}) + + time.Sleep(80 * time.Millisecond) + select { + case <-sessClosed: + t.Fatal("session was closed while an acked cancel was still tearing down") + default: + } + + select { + case res := <-lease.terminal: + if res.GetStatus() != millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED { + t.Fatalf("terminal status = %v, want CANCELLED", res.GetStatus()) + } + case <-time.After(2 * time.Second): + t.Fatal("teardown deadline passed without the mill finishing the lease") + } select { case <-sessClosed: - case <-time.After(1 * time.Second): - t.Fatal("session was not closed after cancel deadline expiration") + t.Fatal("teardown deadline closed the session, failing the node's other jobs") + default: + } + if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || !ok { + t.Fatalf("node quarantined for a slow teardown: ok = %v, err = %v", ok, err) + } +} + +func TestCancelAckDeadlineSparesDisconnectedNode(t *testing.T) { + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) + registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil) + + sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) + m.attachSession(sess) + + lease := newLease("lease-1", "node-1", "inc-1", "dummy") + lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.setState(leaseRunning) + m.mu.Lock() + m.leases[lease.id] = lease + m.mu.Unlock() + + m.destroy(lease.wid) + m.detachSession(sess) + + m.checkCancelAck(lease) + + if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || !ok { + t.Fatalf("node quarantined for a cancel it could not answer: ok = %v, err = %v", ok, err) } } diff --git a/spindle/server.go b/spindle/server.go index 874542d0e..7bc64c804 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -590,9 +590,11 @@ func Run(ctx context.Context) error { // on a mill host, engines place jobs on executors instead of running // them. all names share one Mill m = mill.New(log.SubLogger(logger, "mill"), mill.Config{ - LogDir: cfg.Server.LogDir, - MaxPending: cfg.Mill.MaxPending, - ReconnectGrace: cfg.Mill.ReconnectGrace, + LogDir: cfg.Server.LogDir, + MaxPending: cfg.Mill.MaxPending, + ReconnectGrace: cfg.Mill.ReconnectGrace, + CancelAckTimeout: cfg.Mill.CancelAckTimeout, + CancelTeardownTimeout: cfg.Mill.CancelTeardownTimeout, }) engines = map[string]models.Engine{ "nixery": mill.NewEngine("nixery", m), -- 2.51.2