From b1044cd8a48671f5f16e01c1ea6639cfe768bbbc Mon Sep 17 00:00:00 2001 From: rcourtman <8825017+rcourtman@users.noreply.github.com> Date: Tue, 1 Sep 2026 23:22:55 +0100 Subject: [PATCH] Let replayed request ids wait for the in-flight handler instead of dropping Since 60d0651a88 every typed request registers a per-connection cancellable slot that its handler goroutine releases in a deferred cleanup after sending its result. The server replays a request id when it wants the durable receipt again, and that replay can reach the reader before the previous handler's deferred release runs. launchCancellableRequest treated that as a duplicate and dropped it, so the server waited out the operation's full timeout for a result the agent already held. The Linux x64 native-verification leg failed this way on 12 of the last 25 main runs, always on a "replay 1" dispatch of host update, storage cleanup, or Docker lifecycle. Give each slot a done channel that closes on release. A replay whose id is still registered on the same connection now waits for that release and then runs, answering from the durable receipt. Invalid ids and over-capacity requests are still dropped. A unit test pins the wait-then-run behaviour and the agent-lifecycle contract records the replay rule. --- .../v6/internal/subsystems/agent-lifecycle.md | 8 +++- internal/hostagent/command_client_test.go | 47 +++++++++++++++++++ internal/hostagent/commands.go | 45 ++++++++++++++++-- 3 files changed, 96 insertions(+), 4 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md index d9f8ad823..dfa0dd614 100644 --- a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md +++ b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md @@ -7415,6 +7415,11 @@ kept running on the Proxmox host). Three coupled guarantees: connection-generation-scoped state table. Cancellation or connection teardown before handler registration leaves a tombstone that registration consumes atomically, so provider handoff cannot start after abandonment. + A replay of a request ID that arrives while the previous handler still + owns its slot waits for that slot to be released and then runs, so it + answers from the durable receipt; it is never dropped as a duplicate, + because the server replays exact request IDs to recover receipts and a + dropped replay would leave the server waiting out the full timeout. 3. The unified agent's command client tracks in-flight `execute_command`/`read_file` executions and durable host update, storage-cleanup, Proxmox guest lifecycle, and container lifecycle/update @@ -7445,7 +7450,8 @@ Proofs: `internal/agentexec/server_websocket_test.go` (`TestCommandClient_handleCancelCommand_CancelsRegisteredRequest`, `TestCommandClient_handleCancelCommand_UnknownRequestIsNoOp`, `TestCommandClient_CancellationBeforeRegistrationIsConsumedAndConnectionScoped`, -`TestCommandClient_StaleCleanupCannotEraseReusedRequestCancellation`), +`TestCommandClient_StaleCleanupCannotEraseReusedRequestCancellation`, +`TestCommandClient_ReplayedRequestWaitsForInFlightHandlerInsteadOfDropping`), `internal/hostagent/proxmox_guest_lifecycle_test.go` (`TestProxmoxGuestLifecycleCancellationBeforeHandlerRegistrationSkipsProviderAndPersistsReceipt`), and `internal/hostagent/commands_execute_unix_test.go` (timeout and cancel diff --git a/internal/hostagent/command_client_test.go b/internal/hostagent/command_client_test.go index 30021d8d0..501448745 100644 --- a/internal/hostagent/command_client_test.go +++ b/internal/hostagent/command_client_test.go @@ -6,6 +6,7 @@ import ( "io" "reflect" "testing" + "time" "github.com/gorilla/websocket" "github.com/rcourtman/pulse-go-rewrite/internal/agentexec" @@ -467,3 +468,49 @@ func TestCommandClientActionRunnerMessageCatalogRejectsGenericAuthority(t *testi } } } + +func TestCommandClient_ReplayedRequestWaitsForInFlightHandlerInsteadOfDropping(t *testing.T) { + c := &CommandClient{logger: zerolog.Nop()} + conn := &websocket.Conn{} + const requestID = "typed-replay" + + release := make(chan struct{}) + firstRunning := make(chan struct{}) + secondRan := make(chan struct{}) + c.launchCancellableRequest(conn, requestID, "typed", func() { + close(firstRunning) + <-release + }) + select { + case <-firstRunning: + case <-time.After(2 * time.Second): + t.Fatal("first handler did not start") + } + + // The replay arrives while the first handler still owns the slot. It must + // not be dropped; it runs once the first handler releases the slot, so it + // can answer from the durable receipt. + c.launchCancellableRequest(conn, requestID, "typed", func() { close(secondRan) }) + select { + case <-secondRan: + t.Fatal("replay ran while the first handler still owned the slot") + case <-time.After(50 * time.Millisecond): + } + if c.inflightCancellableRequest(conn, requestID) == nil { + t.Fatal("first handler lost its slot before finishing") + } + + close(release) + select { + case <-secondRan: + case <-time.After(2 * time.Second): + t.Fatal("replay was dropped instead of running after the first handler finished") + } + deadline := time.Now().Add(2 * time.Second) + for c.inflightCancellableRequest(conn, requestID) != nil && time.Now().Before(deadline) { + time.Sleep(5 * time.Millisecond) + } + if c.inflightCancellableRequest(conn, requestID) != nil { + t.Fatal("replay handler did not release the slot") + } +} diff --git a/internal/hostagent/commands.go b/internal/hostagent/commands.go index ff5769fd0..5c54070f1 100644 --- a/internal/hostagent/commands.go +++ b/internal/hostagent/commands.go @@ -157,6 +157,14 @@ type cancellableRequestKey struct { type cancellableRequestState struct { cancel context.CancelFunc canceled bool + // done closes when the request releases its slot, so a replay of the same + // request id that arrives while the previous handler is still finishing + // can wait for it instead of being dropped. + done chan struct{} +} + +func newCancellableRequestState() *cancellableRequestState { + return &cancellableRequestState{done: make(chan struct{})} } // NewCommandClient creates a new command execution client @@ -618,7 +626,19 @@ func computeReconnectDelay(failures int) time.Duration { func (c *CommandClient) launchCancellableRequest(conn *websocket.Conn, requestID, operation string, handle func()) { state := c.noteCancellableRequest(conn, requestID) if state == nil { - c.logger.Warn().Str("request_id", requestID).Str("operation", operation).Msg("Dropping duplicate, invalid, or over-capacity cancellable request") + // The server replays a request id it already dispatched when it wants + // the durable receipt again, and that replay can arrive on the reader + // before the previous handler goroutine has released its slot. Wait for + // that handler instead of dropping the replay, which would leave the + // server waiting out its full timeout for a result that already exists. + if inflight := c.inflightCancellableRequest(conn, requestID); inflight != nil { + go func() { + <-inflight.done + c.launchCancellableRequest(conn, requestID, operation, handle) + }() + return + } + c.logger.Warn().Str("request_id", requestID).Str("operation", operation).Msg("Dropping invalid or over-capacity cancellable request") return } go func() { @@ -627,6 +647,20 @@ func (c *CommandClient) launchCancellableRequest(conn *websocket.Conn, requestID }() } +// inflightCancellableRequest returns the state currently registered for a +// request id on this connection, or nil when the slot is free or the id is +// invalid. +func (c *CommandClient) inflightCancellableRequest(conn *websocket.Conn, requestID string) *cancellableRequestState { + requestID = strings.TrimSpace(requestID) + if requestID == "" || len(requestID) > 128 { + return nil + } + key := cancellableRequestKey{connection: conn, requestID: requestID} + c.activeCommandsMu.Lock() + defer c.activeCommandsMu.Unlock() + return c.cancellableRequests[key] +} + func (c *CommandClient) handleMessages(ctx context.Context, conn *websocket.Conn) error { for { select { @@ -1277,7 +1311,7 @@ func (c *CommandClient) noteCancellableRequest(conn *websocket.Conn, requestID s if _, exists := c.cancellableRequests[key]; exists || len(c.cancellableRequests) >= maxCancellableRequestsPerConnection { return nil } - state := &cancellableRequestState{} + state := newCancellableRequestState() c.cancellableRequests[key] = state return state } @@ -1298,7 +1332,7 @@ func (c *CommandClient) registerActiveCommand(conn *websocket.Conn, requestID st cancel() return nil, false } - state = &cancellableRequestState{} + state = newCancellableRequestState() c.cancellableRequests[key] = state } if state.cancel != nil { @@ -1319,10 +1353,15 @@ func (c *CommandClient) registerActiveCommand(conn *websocket.Conn, requestID st func (c *CommandClient) finishCancellableRequest(conn *websocket.Conn, requestID string, state *cancellableRequestState) { key := cancellableRequestKey{connection: conn, requestID: strings.TrimSpace(requestID)} c.activeCommandsMu.Lock() + released := false if current := c.cancellableRequests[key]; state != nil && current == state { delete(c.cancellableRequests, key) + released = true } c.activeCommandsMu.Unlock() + if released && state.done != nil { + close(state.done) + } } func (c *CommandClient) clearCancellableRequests(conn *websocket.Conn) {