From 7e902d1103ee5c38708a123c4faf82de2d745144 Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 07:18:52 +0000 Subject: [PATCH] feat: reconcile retained client commands --- internal/server/session/agent_server.go | 9 ++- internal/server/store/reconcile.go | 63 +++++++++++++++++++ .../clientagent_integration_test.go | 16 ++++- 3 files changed, 81 insertions(+), 7 deletions(-) create mode 100644 internal/server/store/reconcile.go diff --git a/internal/server/session/agent_server.go b/internal/server/session/agent_server.go index 37b6616..12f0e8e 100644 --- a/internal/server/session/agent_server.go +++ b/internal/server/session/agent_server.go @@ -185,15 +185,14 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w return } if snapshot := envelope.GetReconcileSnapshot(); snapshot != nil { - // A retained command requires the command-store reconciliation phase, - // which is deliberately not substituted with a no-op result. - if len(snapshot.GetRetainedCommands()) != 0 { - server.close(connection, websocket.StatusPolicyViolation, "reconciliation for retained commands is unavailable") + result, reconcileErr := server.Store.ReconcileClientSnapshot(sessionContext, hello.GetClientId(), snapshot) + if reconcileErr != nil { + server.close(connection, websocket.StatusPolicyViolation, "reconciliation failed") return } encoded, err := proto.Marshal(&rvboxv1.AgentEnvelope{ SessionId: encodedSessionID, SessionGeneration: registration.Generation, - Payload: &rvboxv1.AgentEnvelope_ReconcileResult{ReconcileResult: &rvboxv1.ReconcileResult{}}, + Payload: &rvboxv1.AgentEnvelope_ReconcileResult{ReconcileResult: result}, }) if err != nil || queue.EnqueueControl(Frame{Kind: FrameControl, Payload: encoded}) != nil { server.close(connection, websocket.StatusInternalError, "could not queue reconciliation result") diff --git a/internal/server/store/reconcile.go b/internal/server/store/reconcile.go new file mode 100644 index 0000000..7afdf89 --- /dev/null +++ b/internal/server/store/reconcile.go @@ -0,0 +1,63 @@ +package store + +import ( + "context" + "database/sql" + "errors" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" + "github.com/rvbox/rvbox/internal/domain" +) + +// ReconcileClientSnapshot compares client evidence with this client's durable +// server rows. It intentionally performs no lifecycle mutation yet: callers +// receive only the safe local terminate/discard instructions. +func (store *Store) ReconcileClientSnapshot(ctx context.Context, clientID string, snapshot *rvboxv1.ReconcileSnapshot) (*rvboxv1.ReconcileResult, error) { + if snapshot == nil { + return nil, errors.New("missing client reconciliation snapshot") + } + result := &rvboxv1.ReconcileResult{} + for _, client := range snapshot.GetRetainedCommands() { + issue, err := domain.ParseUUIDv7(client.GetIssueUuid()) + if err != nil { + return nil, err + } + input := domain.ReconcileInput{ClientEvidence: domain.ClientEvidenceRetained, ClientLifecycle: client.GetLifecycle(), ClientRevision: domain.CommandRevision(client.GetCommandRevision()), ClientLastEventSeq: client.GetLastClientEventSeq(), ImmutableHashMatches: true} + if client.GetTombstoned() { + input.ClientEvidence = domain.ClientEvidenceTombstone + } + var lifecycle uint32 + var revision, lastSequence uint64 + var hash []byte + err = store.db.QueryRowContext(ctx, `SELECT lifecycle, revision, last_event_seq, immutable_request_sha256 FROM commands WHERE issue_uuid = ? AND client_id = ?`, issue[:], clientID).Scan(&lifecycle, &revision, &lastSequence, &hash) + if err == nil { + input.ServerPresent = true + input.ServerLifecycle = rvboxv1.CommandLifecycle(lifecycle) + input.ServerRevision = domain.CommandRevision(revision) + input.ServerLastEventSeq = lastSequence + input.ImmutableHashMatches = len(hash) == len(client.GetImmutableRequestSha256()) && string(hash) == string(client.GetImmutableRequestSha256()) + } else if !errors.Is(err, sql.ErrNoRows) { + return nil, err + } else { + var tombstoneHash []byte + err = store.db.QueryRowContext(ctx, `SELECT immutable_sha256 FROM command_tombstones WHERE issue_uuid = ? AND client_id = ?`, issue[:], clientID).Scan(&tombstoneHash) + if err == nil { + input.ServerHasTombstone = true + input.ImmutableHashMatches = len(tombstoneHash) == len(client.GetImmutableRequestSha256()) && string(tombstoneHash) == string(client.GetImmutableRequestSha256()) + } else if !errors.Is(err, sql.ErrNoRows) { + return nil, err + } + } + decision, err := domain.DecideReconciliation(input) + if err != nil && !errors.Is(err, domain.ErrReconcileContradiction) { + return nil, err + } + switch decision.Action { + case domain.ReconcileTerminateLocal: + result.TerminateLocalIssueUuids = append(result.TerminateLocalIssueUuids, issue.String()) + case domain.ReconcileDiscardLocalTerminal: + result.DiscardLocalTerminalIssueUuids = append(result.DiscardLocalTerminalIssueUuids, issue.String()) + } + } + return result, nil +} diff --git a/test/integration/clientagent/clientagent_integration_test.go b/test/integration/clientagent/clientagent_integration_test.go index 8d91854..1a78d6c 100644 --- a/test/integration/clientagent/clientagent_integration_test.go +++ b/test/integration/clientagent/clientagent_integration_test.go @@ -39,7 +39,19 @@ func TestWebSocketHelloWelcome_HP_SES_06(t *testing.T) { if result.ID == "" || result.Generation != 1 || result.Protocol.GetMajor() != 1 { t.Fatalf("handshake result = %#v", result) } - if reconciled, err := agent.Reconcile(ctx, transport, result, &rvboxv1.ReconcileSnapshot{}, agentproto.DefaultLimits()); err != nil || reconciled == nil { - t.Fatalf("empty reconciliation = %#v, %v", reconciled, err) + activeIssue := "019c46f1-1d02-7000-8000-000000000072" + terminalIssue := "019c46f1-1d02-7000-8000-000000000073" + reconciled, err := agent.Reconcile(ctx, transport, result, &rvboxv1.ReconcileSnapshot{RetainedCommands: []*rvboxv1.ReconcileCommandState{ + {IssueUuid: activeIssue, Lifecycle: rvboxv1.CommandLifecycle_COMMAND_RUNNING, CommandRevision: 1, ImmutableRequestSha256: make([]byte, 32)}, + {IssueUuid: terminalIssue, Lifecycle: rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, CommandRevision: 1, Tombstoned: true, ImmutableRequestSha256: make([]byte, 32)}, + }}, agentproto.DefaultLimits()) + if err != nil { + t.Fatal(err) + } + if got := reconciled.GetTerminateLocalIssueUuids(); len(got) != 1 || got[0] != activeIssue { + t.Fatalf("terminate local issues = %v", got) + } + if got := reconciled.GetDiscardLocalTerminalIssueUuids(); len(got) != 1 || got[0] != terminalIssue { + t.Fatalf("discard local terminal issues = %v", got) } }