From e79f882993f91056110b976738f2e7a0c216924d Mon Sep 17 00:00:00 2001 From: cabbage Date: Fri, 11 Sep 2026 07:36:11 +0000 Subject: [PATCH] fix: fence reconciliation admission race --- docs/architecture.md | 7 ++- docs/implementation-plan.v1.md | 6 ++ docs/testing.md | 5 +- internal/server/session/agent_server.go | 8 ++- internal/server/session/agent_server_test.go | 58 +++++++++++++++++++- internal/server/store/reconcile.go | 46 ++++++++++++++-- internal/server/store/reconcile_test.go | 25 +++++++++ scripts/windows/native-test | 27 +++++++++ scripts/windows/test-host | 23 +++++++- test/windowsnative/native_fixture_test.go | 4 +- 10 files changed, 193 insertions(+), 16 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 744318c..3013c7b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -97,8 +97,11 @@ Scheduler is not used. 4. The server queues work while a client is offline and dispatches it only when the active session advertises capacity. A replacement session immediately performs bidirectional reconciliation between the server's non-terminal set - and the client's complete retained-command set. The server returns explicit - local terminate/discard decisions before new dispatch begins. + and the client's complete retained-command set. That comparison uses one + server receipt-time cutoff captured before session registration, so commands + admitted during the handshake remain fresh queued work rather than false + missing-client contradictions. The server returns explicit local + terminate/discard decisions before new dispatch begins. ## Command model diff --git a/docs/implementation-plan.v1.md b/docs/implementation-plan.v1.md index cad15ad..5be5bc0 100644 --- a/docs/implementation-plan.v1.md +++ b/docs/implementation-plan.v1.md @@ -1751,6 +1751,12 @@ a lost final `EventAck` followed by server retention. Write this idempotent result before dispatch. A client with unresolved essential-store corruption cannot complete reconciliation or accept work. +Capture one server receipt-time cutoff immediately before registering the live +session, and use that identical cutoff when building `ReconcileRequest` and +applying its `ReconcileSnapshot`. A command admitted after that cutoff is new +work: it must remain queued for post-reconciliation dispatch, never be treated +as absent client evidence merely because its admission raced the handshake. + Implement reconciliation as this explicit matrix: | Server state | Client evidence | Durable result | diff --git a/docs/testing.md b/docs/testing.md index 85a03b9..45bbd89 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -231,7 +231,10 @@ the real `rvbox.exe --install-service` path and proves completion through SCM. `run` is for reconfiguration/restart scenarios after that first installation. Neither action invokes the GUI-subsystem executable directly with the normal Guest Control account. `collect` obtains only bounded/redacted artifacts, and -`reset` restores the exact clean baseline and leaves the VM powered off. +`reset` restores the exact clean baseline and leaves the VM powered off. It +first permits a bounded ACPI shutdown; if that hangs, it force-powers off only +the exact leased disposable VM before snapshot restoration. That intentional +state loss is confined to the test isolation boundary. This service-driven protocol is required because VirtualBox Guest Control 7.2.16 does not reliably complete a direct GUI-subsystem `rvbox.exe` run; diff --git a/internal/server/session/agent_server.go b/internal/server/session/agent_server.go index ce06686..3c8277f 100644 --- a/internal/server/session/agent_server.go +++ b/internal/server/session/agent_server.go @@ -102,6 +102,10 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w return } hello := helloEnvelope.GetClientHello() + // This cutoff precedes registration, which is when control callers can + // first observe the live client. Reconciliation must compare only the + // state included in its request; later control commands are fresh work. + reconcileBoundary := server.now() instanceID, err := domain.ParseUUIDv7(hello.GetClientInstanceId()) if err != nil { server.close(connection, websocket.StatusPolicyViolation, "invalid client instance ID") @@ -192,7 +196,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w server.close(connection, websocket.StatusInternalError, "could not queue welcome") return } - targets, err := server.Store.ReconcileTargets(parent, hello.GetClientId()) + targets, err := server.Store.ReconcileTargetsAt(parent, hello.GetClientId(), reconcileBoundary) if err != nil { server.close(connection, websocket.StatusInternalError, "could not build reconciliation request") return @@ -237,7 +241,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w return } if snapshot := envelope.GetReconcileSnapshot(); snapshot != nil { - result, reconcileErr := server.Store.ReconcileClientSnapshotForSession(sessionContext, hello.GetClientId(), registration.Generation, snapshot) + result, reconcileErr := server.Store.ReconcileClientSnapshotForSessionAt(sessionContext, hello.GetClientId(), registration.Generation, snapshot, reconcileBoundary) if reconcileErr != nil { server.close(connection, websocket.StatusPolicyViolation, "reconciliation failed") return diff --git a/internal/server/session/agent_server_test.go b/internal/server/session/agent_server_test.go index 988c1a9..42f176d 100644 --- a/internal/server/session/agent_server_test.go +++ b/internal/server/session/agent_server_test.go @@ -98,6 +98,57 @@ func TestAgentServerRegistrationAndReplacement_HP_SES_05(t *testing.T) { } } +func TestAgentServerDispatchesCommandAdmittedDuringReconcile_HP_SES_11(t *testing.T) { + server, agent, cleanup := newTestAgentServerWithStore(t) + defer cleanup() + instance, err := domain.NewUUIDv7() + if err != nil { + t.Fatal(err) + } + connection, welcome := dialAndHello(t, server, "client-race", instance.String()) + defer connection.CloseNow() + issue, err := domain.NewUUIDv7() + if err != nil { + t.Fatal(err) + } + spec, err := proto.Marshal(&rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_CMD, Source: &rvboxv1.ExecutionSpec_CommandText{CommandText: "echo dispatch"}}) + if err != nil { + t.Fatal(err) + } + now := time.Now().UTC() + if _, err := agent.Store.QueueCommand(context.Background(), store.QueueCommandInput{IssueUUID: issue, ClientID: "client-race", IssueTime: now, ReceiptTime: now, ImmutableSHA256: sha256.Sum256([]byte("during-reconcile")), ExecutionSpec: spec}); err != nil { + t.Fatal(err) + } + snapshot, err := proto.Marshal(&rvboxv1.AgentEnvelope{ + SessionId: welcome.GetSessionId(), SessionGeneration: welcome.GetSessionGeneration(), + Payload: &rvboxv1.AgentEnvelope_ReconcileSnapshot{ReconcileSnapshot: &rvboxv1.ReconcileSnapshot{}}, + }) + if err != nil { + t.Fatal(err) + } + if err := connection.Write(context.Background(), websocket.MessageBinary, snapshot); err != nil { + t.Fatal(err) + } + readContext, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + for index := 0; index < 2; index++ { + _, payload, err := connection.Read(readContext) + if err != nil { + t.Fatal(err) + } + var envelope rvboxv1.AgentEnvelope + if err := proto.Unmarshal(payload, &envelope); err != nil { + t.Fatal(err) + } + if index == 0 && envelope.GetReconcileResult() == nil { + t.Fatalf("first post-snapshot envelope = %T, want reconcile result", envelope.Payload) + } + if index == 1 && (envelope.GetCommandDispatch() == nil || envelope.GetCommandDispatch().GetIssueUuid() != issue.String()) { + t.Fatalf("second post-snapshot envelope = %+v, want dispatch %s", envelope.Payload, issue) + } + } +} + func TestAgentServerRejectsHostileWireInputs_BH_SES_04(t *testing.T) { server, cleanup := newTestAgentServer(t) defer cleanup() @@ -174,6 +225,11 @@ func TestWireEventAppendCarriesClientBinding_HP_EVENT_01(t *testing.T) { } func newTestAgentServer(t *testing.T) (*httptest.Server, func()) { + server, _, cleanup := newTestAgentServerWithStore(t) + return server, cleanup +} + +func newTestAgentServerWithStore(t *testing.T) (*httptest.Server, *AgentServer, func()) { t.Helper() dataDirectory := filepath.Join(t.TempDir(), "store") if err := os.Mkdir(dataDirectory, 0o700); err != nil { @@ -190,7 +246,7 @@ func newTestAgentServer(t *testing.T) (*httptest.Server, func()) { HeartbeatIdle: time.Second, LivenessTimeout: 2 * time.Second, } httpServer := httptest.NewServer(agent) - return httpServer, func() { + return httpServer, agent, func() { httpServer.Close() if err := persistence.Close(); err != nil { t.Errorf("close persistence: %v", err) diff --git a/internal/server/store/reconcile.go b/internal/server/store/reconcile.go index 82d708d..cddc0d9 100644 --- a/internal/server/store/reconcile.go +++ b/internal/server/store/reconcile.go @@ -16,11 +16,26 @@ import ( // client. It is a read-only snapshot used to tell a reconnecting agent which // UUIDs and durable cursors must be compared before fresh dispatch is enabled. func (store *Store) ReconcileTargets(ctx context.Context, clientID string) ([]*rvboxv1.ReconcileTarget, error) { + return store.ReconcileTargetsAt(ctx, clientID, time.Time{}) +} + +// ReconcileTargetsAt returns the durable non-terminal view at the supplied +// receipt-time boundary. A connecting session uses one boundary for both the +// request and its reply: commands admitted after it are fresh work, not absent +// client evidence to be reconciled. +func (store *Store) ReconcileTargetsAt(ctx context.Context, clientID string, receiptBoundary time.Time) ([]*rvboxv1.ReconcileTarget, error) { if clientID == "" { return nil, errors.New("client ID is required") } - rows, err := store.db.QueryContext(ctx, `SELECT issue_uuid, last_event_seq, revision, immutable_request_sha256 -FROM commands WHERE client_id = ? AND lifecycle BETWEEN 1 AND 4 ORDER BY issue_time, issue_uuid`, clientID) + query := `SELECT issue_uuid, last_event_seq, revision, immutable_request_sha256 +FROM commands WHERE client_id = ? AND lifecycle BETWEEN 1 AND 4` + args := []any{clientID} + if !receiptBoundary.IsZero() { + query += ` AND server_receipt_time <= ?` + args = append(args, receiptBoundary.UTC().UnixNano()) + } + query += ` ORDER BY issue_time, issue_uuid` + rows, err := store.db.QueryContext(ctx, query, args...) if err != nil { return nil, err } @@ -62,10 +77,19 @@ func (store *Store) ReconcileClientSnapshot(ctx context.Context, clientID string // incidented. Retained non-terminal rows are retargeted to this generation so // late events from the previous connection cannot advance the command. func (store *Store) ReconcileClientSnapshotForSession(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot) (*rvboxv1.ReconcileResult, error) { + return store.ReconcileClientSnapshotForSessionAt(ctx, clientID, generation, snapshot, time.Time{}) +} + +// ReconcileClientSnapshotForSessionAt reconciles exactly the server state that +// was included in the matching ReconcileRequest. Work admitted after the +// boundary is intentionally left for the dispatch loop once reconciliation +// completes; treating it as absent client evidence would lose a valid command +// during the Hello/reconcile race. +func (store *Store) ReconcileClientSnapshotForSessionAt(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot, receiptBoundary time.Time) (*rvboxv1.ReconcileResult, error) { if generation == 0 { return nil, errors.New("session generation is required") } - return store.reconcileClientSnapshot(ctx, clientID, generation, snapshot) + return store.reconcileClientSnapshotAt(ctx, clientID, generation, snapshot, receiptBoundary) } type reconcileServerRow struct { @@ -78,6 +102,10 @@ type reconcileServerRow struct { } func (store *Store) reconcileClientSnapshot(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot) (*rvboxv1.ReconcileResult, error) { + return store.reconcileClientSnapshotAt(ctx, clientID, generation, snapshot, time.Time{}) +} + +func (store *Store) reconcileClientSnapshotAt(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot, receiptBoundary time.Time) (*rvboxv1.ReconcileResult, error) { if clientID == "" { return nil, errors.New("client ID is required") } @@ -112,7 +140,7 @@ func (store *Store) reconcileClientSnapshot(ctx context.Context, clientID string return nil, err } defer tx.Rollback() - serverRows, tombstones, err := loadReconcileRows(ctx, tx, clientID) + serverRows, tombstones, err := loadReconcileRows(ctx, tx, clientID, receiptBoundary) if err != nil { store.writeMu.Unlock() return nil, err @@ -215,8 +243,14 @@ type reconcileIncident struct { dataLoss bool } -func loadReconcileRows(ctx context.Context, tx *sql.Tx, clientID string) (map[string]reconcileServerRow, map[string][]byte, error) { - rows, err := tx.QueryContext(ctx, `SELECT issue_uuid, lifecycle, revision, last_event_seq, immutable_request_sha256, target_session_generation FROM commands WHERE client_id = ?`, clientID) +func loadReconcileRows(ctx context.Context, tx *sql.Tx, clientID string, receiptBoundary time.Time) (map[string]reconcileServerRow, map[string][]byte, error) { + query := `SELECT issue_uuid, lifecycle, revision, last_event_seq, immutable_request_sha256, target_session_generation FROM commands WHERE client_id = ?` + args := []any{clientID} + if !receiptBoundary.IsZero() { + query += ` AND server_receipt_time <= ?` + args = append(args, receiptBoundary.UTC().UnixNano()) + } + rows, err := tx.QueryContext(ctx, query, args...) if err != nil { return nil, nil, err } diff --git a/internal/server/store/reconcile_test.go b/internal/server/store/reconcile_test.go index 6712ddb..58947cf 100644 --- a/internal/server/store/reconcile_test.go +++ b/internal/server/store/reconcile_test.go @@ -71,6 +71,31 @@ func TestReconcileSnapshotMutatesMissingAndRetargets_HP_SES_13(t *testing.T) { } } +func TestReconcileSessionBoundaryLeavesNewlyQueuedCommandForDispatch_HP_RECONCILE_08(t *testing.T) { + t.Parallel() + ctx := context.Background() + opened, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "state"), BusyTimeout: time.Second}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = opened.Close() }) + boundary := time.Date(2026, time.September, 11, 12, 0, 0, 0, time.UTC) + if _, err := opened.RegisterClientSession(ctx, ClientRegistration{ClientID: "boundary-client", Platform: 3, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\\`, SupportedShells: []byte{1}, ClientInstanceID: [16]byte{31}, SessionID: [16]byte{32}, ConnectedAt: boundary}); err != nil { + t.Fatal(err) + } + issue := mustReconcileIssue(t, "019c46f1-1d02-7000-8000-0000000000d1") + if _, err := opened.QueueCommand(ctx, QueueCommandInput{IssueUUID: issue, ClientID: "boundary-client", IssueTime: boundary.Add(time.Nanosecond), ReceiptTime: boundary.Add(time.Nanosecond), ImmutableSHA256: sha256.Sum256([]byte("new-after-reconcile-boundary")), ExecutionSpec: []byte("echo queued")}); err != nil { + t.Fatal(err) + } + if _, err := opened.ReconcileClientSnapshotForSessionAt(ctx, "boundary-client", 1, &rvboxv1.ReconcileSnapshot{}, boundary); err != nil { + t.Fatal(err) + } + view, err := opened.GetCommandView(ctx, "boundary-client", issue) + if err != nil || view.Lifecycle != uint32(rvboxv1.CommandLifecycle_COMMAND_QUEUED) { + t.Fatalf("post-boundary command = %#v, %v", view, err) + } +} + func mustReconcileIssue(t *testing.T, value string) domain.UUID { t.Helper() issue, err := domain.ParseUUIDv7(value) diff --git a/scripts/windows/native-test b/scripts/windows/native-test index 0c41c3c..fc2a017 100755 --- a/scripts/windows/native-test +++ b/scripts/windows/native-test @@ -179,6 +179,32 @@ assert_stdin_close() { fail "stdin command did not reach terminal success" } +assert_signal_term() { + issued=$(rvc run --background --shell cmd "$client_id" 'ping -t 127.0.0.1 >NUL' 2>&1) || fail "signal command admission failed: $issued" + issue=$(printf '%s\n' "$issued" | awk 'NR == 1 { print $1 }') + case $issue in ????????-????-7???-????-????????????) ;; *) fail "signal command returned invalid issue UUID: $issued" ;; esac + attempt=0 + while [ "$attempt" -lt 30 ]; do + result=$(rvc stat "$client_id" "$issue" 2>/dev/null || true) + if printf '%s\n' "$result" | grep -q 'lifecycle=COMMAND_RUNNING'; then break; fi + attempt=$((attempt + 1)); sleep 1 + done + [ "$attempt" -lt 30 ] || fail "signal command did not reach running state" + rvc kill TERM "$client_id" "$issue" >/dev/null || fail "TERM request failed" + attempt=0 + while [ "$attempt" -lt 45 ]; do + result=$(rvc stat "$client_id" "$issue" 2>/dev/null || true) + if printf '%s\n' "$result" | grep -q 'lifecycle=COMMAND_TERMINATED'; then + printf '%s\n' "$result" >"$fixture_dir/signal-term.stat" + printf 'passed signal-term issue=%s\n' "$issue" + return 0 + fi + case $result in *'lifecycle=COMMAND_FAILED'*|*'lifecycle=COMMAND_REJECTED'*|*'lifecycle=COMMAND_SUCCEEDED'*) fail "TERM command reached wrong terminal state: $result" ;; esac + attempt=$((attempt + 1)); sleep 1 + done + fail "TERM command did not reach terminal state" +} + collect() { if [ "$vm_prepared" = yes ]; then "$repo_root/scripts/windows/test-host" collect --run-id "$run_id" || true @@ -224,6 +250,7 @@ case $action in "$repo_root/scripts/windows/test-host" run --run-id "$run_id" --endpoint "$endpoint_host:$port" wait_client assert_stdin_close + assert_signal_term assert_context active-user no active-user assert_context active-user-elevated yes active-user-elevated "$repo_root/scripts/windows/test-host" run --run-id "$run_id" --fail-contexts ACTIVE_USER_ELEVATED diff --git a/scripts/windows/test-host b/scripts/windows/test-host index 517bf78..303d9e0 100755 --- a/scripts/windows/test-host +++ b/scripts/windows/test-host @@ -540,7 +540,22 @@ case "$action" in attempt=$((attempt + 1)) sleep 1 done - [ "$(state)" = poweroff ] || fail "guest did not power off before reset" + # This is the exact named disposable fixture, with the matching + # run lease. Reset is its isolation boundary, not a graceful-stop + # diagnostic command: after a bounded ACPI attempt, discard only + # this guest's state so snapshot restore cannot strand the lane. + if [ "$(state)" != poweroff ]; then + printf 'reset: ACPI shutdown timed out; forcing disposable VM poweroff\n' >&2 + VBoxManage controlvm "$vm" poweroff >/dev/null + attempt=0 + while [ "$attempt" -lt 30 ]; do + [ "$(state)" = poweroff ] && break + attempt=$((attempt + 1)) + sleep 1 + done + [ "$(state)" = poweroff ] || fail "guest did not power off after forced reset" + printf 'reset-force-poweroff vm=%s\n' "$vm" + fi fi VBoxManage snapshot "$vm" restore "$snapshot" >/dev/null assert_identity @@ -580,7 +595,9 @@ accelerated_stage() { printf 'stage: accelerated publish directory unavailable\n' >&2 return 1 } - stage_http_dir=$(mktemp -d "${TMPDIR:-/tmp}/rvbox-http-stage.XXXXXX") + # Keep all disposable test bytes inside the owning run directory; never + # consume global /tmp, which may belong to a different test or user. + stage_http_dir=$(mktemp -d "$run_root/.accelerated-stage.XXXXXX") stage_http_xz=$stage_http_dir/rvbox.exe.xz xz -T0 -3 -c "$bundle/rvbox.exe" >"$stage_http_xz" stage_http_hash=$(sha256sum "$stage_http_xz" | awk '{print $1}') @@ -615,7 +632,7 @@ copy_stage_file() { target_name=$2 copy_attempt=1 while [ "$copy_attempt" -le 4 ]; do - if scp -q "$source_file" "$RVBOX_TEST_VBOX_HOST:$host_stage/$target_name"; then + if scp -q -o BatchMode=yes -o ConnectTimeout=10 -o ServerAliveInterval=10 -o ServerAliveCountMax=2 "$source_file" "$RVBOX_TEST_VBOX_HOST:$host_stage/$target_name"; then return 0 fi copy_attempt=$((copy_attempt + 1)) diff --git a/test/windowsnative/native_fixture_test.go b/test/windowsnative/native_fixture_test.go index 14e7329..1b9b75f 100644 --- a/test/windowsnative/native_fixture_test.go +++ b/test/windowsnative/native_fixture_test.go @@ -29,6 +29,8 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) { "ACTIVE_USER_ELEVATED,ACTIVE_SYSTEM", "assert_context local-service no local-service", "assert_stdin_close", + "assert_signal_term", + "rvc kill TERM \"$client_id\" \"$issue\"", "rvc append \"$client_id\" \"$issue\" RVBOX_NATIVE_INPUT", "rvc close-stdin \"$client_id\" \"$issue\"", "clean --purge --yes", @@ -40,7 +42,7 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) { } } testHost := read("scripts/windows/test-host") - for _, required := range []string{"xz -T0 -3", "accelerated_stage", "copy_stage_file", "--proxy", "--anyauth", "--continue-at", "rvbox.exe.xz", "retry_limit=4", "ConnectTimeout=10", "bundle transfer did not reach the expected SHA-256 manifest", "verified transfer_sha256"} { + for _, required := range []string{"xz -T0 -3", "accelerated_stage", "copy_stage_file", "--proxy", "--anyauth", "--continue-at", "rvbox.exe.xz", "retry_limit=4", "ConnectTimeout=10", "bundle transfer did not reach the expected SHA-256 manifest", "verified transfer_sha256", "reset-force-poweroff", "$run_root/.accelerated-stage.XXXXXX"} { if !strings.Contains(testHost, required) { t.Fatalf("native test-host is missing compressed transfer contract %q", required) }