diff --git a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md index 349bb33c8..1e657a771 100644 --- a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md +++ b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md @@ -7792,6 +7792,12 @@ kept running on the Proxmox host). Three coupled guarantees: 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. + Inbound results are correlated to the exact session that carried the + request, so once that session's reader exits (socket drop or read error) + the server signals the session done and every dispatch still waiting on + it fails immediately with a `disconnected before ... receipt` error + instead of waiting out its operation timeout. The caller recovers the + outcome by replaying the same request ID on the runner's next session. 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 @@ -7815,7 +7821,8 @@ Proofs: `internal/agentexec/server_websocket_test.go` `TestExecuteCommand_AbandonedCommandSendsCancel`), `internal/agentexec/server_websocket_test.go` (`TestTypedOperations_AbandonedDispatchSendsExactlyOneCancel`, -`TestTypedOperation_TimeoutSendsCancelAndExpiredContextNeverDispatches`), +`TestTypedOperation_TimeoutSendsCancelAndExpiredContextNeverDispatches`, +`TestTypedOperation_SocketDropAfterSendUnblocksDispatch`), `internal/hostagent/operation_receipt_websocket_integration_test.go` (`TestRealServerActionRunnerCancellationPersistsAndReplaysProxmoxReceiptAfterReconnect`), `internal/hostagent/command_client_test.go` diff --git a/internal/agentexec/server.go b/internal/agentexec/server.go index 41473283f..7f5bc48e7 100644 --- a/internal/agentexec/server.go +++ b/internal/agentexec/server.go @@ -1580,6 +1580,10 @@ func (s *Server) readLoop(ac *agentConn) { for _, ch := range closeChs { close(ch) } + // Results are correlated to this exact session, so nothing can answer + // a request it carried once its reader is gone. Release in-flight + // dispatches now instead of leaving them to wait out their timeout. + ac.signalDone() if err := ac.conn.Close(); err != nil { log.Debug().Err(err).Str("agent_id", agentID).Msg("Failed to close connection during read-loop cleanup") } diff --git a/internal/agentexec/server_websocket_test.go b/internal/agentexec/server_websocket_test.go index 51559ae97..e622e3a9b 100644 --- a/internal/agentexec/server_websocket_test.go +++ b/internal/agentexec/server_websocket_test.go @@ -2152,6 +2152,38 @@ func TestTypedOperation_TimeoutSendsCancelAndExpiredContextNeverDispatches(t *te }) } +// A typed request is correlated to the session that carried it, so once that +// socket drops nothing can deliver its receipt. The dispatch must report the +// disconnect instead of waiting out the operation timeout. +func TestTypedOperation_SocketDropAfterSendUnblocksDispatch(t *testing.T) { + admission := AgentAdmission{TokenID: "typed-token", AgentID: "typed-a1", Hostname: "typed-host1", RuntimeRole: RuntimeRoleActionRunner, ActionCapability: ActionCapabilityTypedV1} + s := NewServerWithAdmissionValidator(func(token, _, _ string) (AgentAdmission, bool) { + return admission, token == admission.TokenID + }, func(AgentAdmission) bool { return true }) + ts := newWSServer(t, s) + defer ts.Close() + agent := registerTypedCancelTestAgent(t, s, ts.URL) + + request := ProxmoxGuestLifecyclePayload{RequestID: "pve.dropped.1", ActionID: "pve", Operation: "shutdown", GuestKind: "vm", VMID: 101, ExpectedStatus: "running", Timeout: 30} + if err := BindProxmoxGuestLifecyclePayload(&request); err != nil { + t.Fatal(err) + } + done := make(chan error, 1) + go func() { _, err := s.ExecuteProxmoxGuestLifecycle(context.Background(), "typed-a1", request); done <- err }() + if sent, ok := agent.nextMessage(3 * time.Second); !ok || sent.Type != MsgTypeProxmoxGuestLifecycle { + t.Fatalf("dispatched request = %+v", sent) + } + _ = agent.conn.Close() + select { + case err := <-done: + if err == nil || !strings.Contains(err.Error(), "disconnected before") { + t.Fatalf("dropped-session dispatch error = %v", err) + } + case <-time.After(5 * time.Second): + t.Fatal("dispatch kept waiting on a session whose socket already closed") + } +} + // Probe the scheduling window after publication but before the registration // acknowledgement. No sleeps or repeated runs are needed to expose this order. func TestRegistrationReservesFirstFrameBeforePublishingSession(t *testing.T) { diff --git a/internal/hostagent/operation_receipt_websocket_integration_test.go b/internal/hostagent/operation_receipt_websocket_integration_test.go index 221da7730..1d6a421e1 100644 --- a/internal/hostagent/operation_receipt_websocket_integration_test.go +++ b/internal/hostagent/operation_receipt_websocket_integration_test.go @@ -1,6 +1,7 @@ package hostagent import ( + "bytes" "context" "encoding/json" "errors" @@ -253,8 +254,23 @@ func TestRealServerActionRunnerCancellationPersistsAndReplaysProxmoxReceiptAfter })) defer httpServer.Close() + // A runner that drops its session must come back well inside the replay + // window rather than after the production reconnect backoff. + origReconnectDelay := reconnectDelay + reconnectDelay = 50 * time.Millisecond + t.Cleanup(func() { reconnectDelay = origReconnectDelay }) + stateDir := t.TempDir() - logger := zerolog.Nop() + // Keep runner logs so a failure shows why a session dropped. They are only + // printed from the cleanup below, never from runner goroutines that can + // outlive the test. + runnerLogs := &lockedLogBuffer{} + t.Cleanup(func() { + if t.Failed() { + t.Logf("action runner log:\n%s", runnerLogs.String()) + } + }) + logger := zerolog.New(runnerLogs).With().Timestamp().Logger() mutationStarted := make(chan struct{}) var startOnce sync.Once var mutationMu sync.Mutex @@ -373,7 +389,23 @@ func TestRealServerActionRunnerCancellationPersistsAndReplaysProxmoxReceiptAfter second, cancelSecond, secondDone := startRunner(t) defer stopRunner(t, second, cancelSecond, secondDone) - replayed, err := server.ExecuteProxmoxGuestLifecycle(context.Background(), admission.AgentID, request) + // The server publishes the session before the runner finishes activating, + // so the replay can land on a session the runner then drops. That request + // can never be answered; send the same request again on the session the + // runner reconnects with. Every attempt must replay the durable receipt. + var replayed *agentexec.ProxmoxGuestLifecycleResultPayload + var err error + for attempt := 1; ; attempt++ { + replayed, err = server.ExecuteProxmoxGuestLifecycle(context.Background(), admission.AgentID, request) + if err == nil || attempt == 3 || !strings.Contains(err.Error(), "disconnected before") { + break + } + t.Logf("replay attempt %d lost its session (%v), retrying after reconnect", attempt, err) + deadline := time.Now().Add(3 * time.Second) + for !server.IsAgentConnected(admission.AgentID) && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + } if err != nil || replayed == nil || !replayed.MutationStarted || replayed.MutationCompleted || replayed.Error != canceled.Error { t.Fatalf("replayed result = %+v, err=%v", replayed, err) } @@ -404,3 +436,21 @@ func TestRealServerActionRunnerCancellationPersistsAndReplaysProxmoxReceiptAfter t.Fatalf("cross-agent receipt query error = %v", err) } } + +// lockedLogBuffer collects zerolog output from concurrent runner goroutines. +type lockedLogBuffer struct { + mu sync.Mutex + buf bytes.Buffer +} + +func (b *lockedLogBuffer) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.Write(p) +} + +func (b *lockedLogBuffer) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.String() +}