From e1e4c3e7004d1b127ab78f34a85aa08d589349cb Mon Sep 17 00:00:00 2001 From: rcourtman <8825017+rcourtman@users.noreply.github.com> Date: Tue, 1 Sep 2026 22:51:56 +0100 Subject: [PATCH 1/3] Wait for the server to observe runner disconnect before replay reconnect The cancellation replay test stopped the first action runner and started a second one as soon as the client goroutine exited. The server observes the socket close on its own reader, so under GOMAXPROCS=1 with the race detector the second startRunner saw the stale session as still connected and the replay was dispatched to the dead socket, timing out after 30s on the sharded release preflight while passing on multi-core hosts. Wait until the server reports the agent disconnected before reconnecting, as the agentexec server tests already do. --- .../operation_receipt_websocket_integration_test.go | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/internal/hostagent/operation_receipt_websocket_integration_test.go b/internal/hostagent/operation_receipt_websocket_integration_test.go index 322f4c2e3..221da7730 100644 --- a/internal/hostagent/operation_receipt_websocket_integration_test.go +++ b/internal/hostagent/operation_receipt_websocket_integration_test.go @@ -308,6 +308,17 @@ func TestRealServerActionRunnerCancellationPersistsAndReplaysProxmoxReceiptAfter case <-time.After(3 * time.Second): t.Fatal("action runner did not stop") } + // The client has exited, but the server observes the socket close on + // its own reader goroutine. Reconnecting before that lands would let + // startRunner see the stale session as "connected" and dispatch the + // replay to a dead socket. + deadline := time.Now().Add(3 * time.Second) + for server.IsAgentConnected(admission.AgentID) && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + if server.IsAgentConnected(admission.AgentID) { + t.Fatal("server did not observe the action runner disconnect") + } } request := agentexec.ProxmoxGuestLifecyclePayload{ From b763b8068090196120f73fb3ca9160f7c3910d6b Mon Sep 17 00:00:00 2001 From: rcourtman <8825017+rcourtman@users.noreply.github.com> Date: Tue, 1 Sep 2026 23:05:22 +0100 Subject: [PATCH 2/3] Build the FIFO lifecycle installer test only on unix agent_state_dir_lifecycle_test.go calls syscall.Mkfifo, which does not exist on Windows, so scripts/installtests has failed to compile in the Windows leg of Unified Agent Native Verification since 53267e149d and the install.ps1 contract tests there have not run. Every test in the file drives install.sh through bash and systemd, so tag the file unix-only, matching the other lifecycle lab files. GOOS=windows go vet now passes. --- scripts/installtests/agent_state_dir_lifecycle_test.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/scripts/installtests/agent_state_dir_lifecycle_test.go b/scripts/installtests/agent_state_dir_lifecycle_test.go index 0c325d950..696f83988 100644 --- a/scripts/installtests/agent_state_dir_lifecycle_test.go +++ b/scripts/installtests/agent_state_dir_lifecycle_test.go @@ -1,3 +1,5 @@ +//go:build unix + package installtests import ( 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 3/3] 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) {