From 3ecab55faaf12e63e8af079070707e264a76cc3c Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 07:15:02 +0000 Subject: [PATCH] feat: snapshot retained client reconciliation state --- internal/client/spool/events.go | 6 ++- internal/client/spool/migrations.go | 2 + internal/client/spool/snapshot.go | 73 ++++++++++++++++++++++++++ internal/client/spool/snapshot_test.go | 52 ++++++++++++++++++ test/coverage.toml | 6 +++ 5 files changed, 137 insertions(+), 2 deletions(-) create mode 100644 internal/client/spool/snapshot.go create mode 100644 internal/client/spool/snapshot_test.go diff --git a/internal/client/spool/events.go b/internal/client/spool/events.go index e09b038..19fadaf 100644 --- a/internal/client/spool/events.go +++ b/internal/client/spool/events.go @@ -299,8 +299,10 @@ func (store *Store) CleanupTerminal(ctx context.Context, issueUUID domain.UUID, defer tx.Rollback() var hash []byte var terminal int + var revision uint64 + var phase uint32 var totalCharged uint64 - err = tx.QueryRowContext(ctx, `SELECT immutable_sha256, terminal, total_charged_bytes FROM commands WHERE issue_uuid = ?`, issueUUID[:]).Scan(&hash, &terminal, &totalCharged) + err = tx.QueryRowContext(ctx, `SELECT immutable_sha256, command_revision, phase, terminal, total_charged_bytes FROM commands WHERE issue_uuid = ?`, issueUUID[:]).Scan(&hash, &revision, &phase, &terminal, &totalCharged) if err == sql.ErrNoRows { return ErrUnknownCommand } @@ -324,7 +326,7 @@ func (store *Store) CleanupTerminal(ctx context.Context, issueUUID domain.UUID, if totalCharged > clientTotal { return fmt.Errorf("client spool aggregate charge counter mismatch") } - if _, err := tx.ExecContext(ctx, `INSERT INTO command_tombstones(issue_uuid, immutable_sha256, acknowledged_at) VALUES (?, ?, ?)`, issueUUID[:], hash, acknowledgedAt.UnixNano()); err != nil { + if _, err := tx.ExecContext(ctx, `INSERT INTO command_tombstones(issue_uuid, immutable_sha256, command_revision, terminal_lifecycle, acknowledged_at) VALUES (?, ?, ?, ?, ?)`, issueUUID[:], hash, revision, phase, acknowledgedAt.UnixNano()); err != nil { return err } if _, err := tx.ExecContext(ctx, `DELETE FROM commands WHERE issue_uuid = ?`, issueUUID[:]); err != nil { diff --git a/internal/client/spool/migrations.go b/internal/client/spool/migrations.go index 6d78f5f..3900bee 100644 --- a/internal/client/spool/migrations.go +++ b/internal/client/spool/migrations.go @@ -107,6 +107,8 @@ INSERT INTO spool_counters(singleton, client_total_charged_bytes, charge_version CREATE TABLE command_tombstones ( issue_uuid BLOB PRIMARY KEY CHECK(length(issue_uuid) = 16), immutable_sha256 BLOB NOT NULL CHECK(length(immutable_sha256) = 32), + command_revision INTEGER NOT NULL CHECK(command_revision > 0), + terminal_lifecycle INTEGER NOT NULL CHECK(terminal_lifecycle BETWEEN 5 AND 11), acknowledged_at INTEGER NOT NULL ) STRICT; CREATE INDEX tombstones_fifo ON command_tombstones(acknowledged_at, issue_uuid); diff --git a/internal/client/spool/snapshot.go b/internal/client/spool/snapshot.go new file mode 100644 index 0000000..890112c --- /dev/null +++ b/internal/client/spool/snapshot.go @@ -0,0 +1,73 @@ +package spool + +import ( + "context" + "fmt" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +// ReconcileSnapshot is a consistent durable view. The caller sends it only +// after welcome and before admitting new dispatches on that session. +func (store *Store) ReconcileSnapshot(ctx context.Context) (*rvboxv1.ReconcileSnapshot, error) { + result := &rvboxv1.ReconcileSnapshot{} + rows, err := store.db.QueryContext(ctx, `SELECT issue_uuid, command_revision, phase, next_event_seq, immutable_sha256 FROM commands ORDER BY issue_uuid`) + if err != nil { + return nil, err + } + for rows.Next() { + command, err := scanReconcileCommand(rows, false) + if err != nil { + _ = rows.Close() + return nil, err + } + result.RetainedCommands = append(result.RetainedCommands, command) + } + if err := rows.Close(); err != nil { + return nil, err + } + tombstones, err := store.db.QueryContext(ctx, `SELECT issue_uuid, command_revision, terminal_lifecycle, immutable_sha256 FROM command_tombstones ORDER BY acknowledged_at, issue_uuid`) + if err != nil { + return nil, err + } + defer tombstones.Close() + for tombstones.Next() { + var owner, digest []byte + var revision uint64 + var lifecycle uint32 + if err := tombstones.Scan(&owner, &revision, &lifecycle, &digest); err != nil { + return nil, err + } + if _, valid := copiedUUID(owner); !valid || len(digest) != 32 || revision == 0 || !isTerminalPhase(lifecycle) { + return nil, ErrScriptState + } + result.RetainedCommands = append(result.RetainedCommands, &rvboxv1.ReconcileCommandState{IssueUuid: uuidText(owner), CommandRevision: revision, Lifecycle: rvboxv1.CommandLifecycle(lifecycle), Tombstoned: true, ImmutableRequestSha256: append([]byte(nil), digest...)}) + } + if err := tombstones.Err(); err != nil { + return nil, err + } + return result, nil +} + +type reconcileRow interface { + Scan(...any) error +} + +func scanReconcileCommand(row reconcileRow, tombstoned bool) (*rvboxv1.ReconcileCommandState, error) { + var owner, digest []byte + var revision uint64 + var lifecycle uint32 + var nextSequence uint64 + if err := row.Scan(&owner, &revision, &lifecycle, &nextSequence, &digest); err != nil { + return nil, err + } + if _, valid := copiedUUID(owner); !valid || len(digest) != 32 || revision == 0 || lifecycle == 0 || lifecycle > 11 || nextSequence == 0 { + return nil, fmt.Errorf("%w: invalid reconciliation row", ErrScriptState) + } + return &rvboxv1.ReconcileCommandState{IssueUuid: uuidText(owner), CommandRevision: revision, Lifecycle: rvboxv1.CommandLifecycle(lifecycle), LastClientEventSeq: nextSequence - 1, Tombstoned: tombstoned, ImmutableRequestSha256: append([]byte(nil), digest...)}, nil +} + +func uuidText(value []byte) string { + parsed, _ := copiedUUID(value) + return parsed.String() +} diff --git a/internal/client/spool/snapshot_test.go b/internal/client/spool/snapshot_test.go new file mode 100644 index 0000000..efba920 --- /dev/null +++ b/internal/client/spool/snapshot_test.go @@ -0,0 +1,52 @@ +package spool + +import ( + "context" + "path/filepath" + "testing" + "time" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestReconcileSnapshotRetainedAndTombstoned_HP_SES_09(t *testing.T) { + t.Parallel() + ctx := context.Background() + store := openTestStore(t, ctx, filepath.Join(t.TempDir(), "spool"), DefaultTombstoneLimit) + issue := testUUID(t, "019c46f1-1d02-7000-8000-000000000081") + command := testCommand(issue, []byte("reconcile")) + if _, err := store.AcceptCommand(ctx, command, time.Now().UTC()); err != nil { + t.Fatal(err) + } + if _, err := store.AppendEvent(ctx, issue, EventInput{Kind: 3, Compression: 1, Payload: []byte("running"), CreatedAt: time.Now().UTC()}); err != nil { + t.Fatal(err) + } + if _, err := store.AssignSendWindow(ctx, issue, 1, 1<<20); err != nil { + t.Fatal(err) + } + snapshot, err := store.ReconcileSnapshot(ctx) + if err != nil || len(snapshot.RetainedCommands) != 1 { + t.Fatalf("retained snapshot = %#v, %v", snapshot, err) + } + retained := snapshot.RetainedCommands[0] + if retained.GetIssueUuid() != issue.String() || retained.GetLastClientEventSeq() != 1 || retained.GetTombstoned() || retained.GetLifecycle() != rvboxv1.CommandLifecycle_COMMAND_QUEUED { + t.Fatalf("retained command = %#v", retained) + } + if err := store.Ack(ctx, issue, 1); err != nil { + t.Fatal(err) + } + if err := store.MarkTerminal(ctx, issue, uint32(rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED)); err != nil { + t.Fatal(err) + } + if err := store.CleanupTerminal(ctx, issue, time.Now().UTC()); err != nil { + t.Fatal(err) + } + snapshot, err = store.ReconcileSnapshot(ctx) + if err != nil || len(snapshot.RetainedCommands) != 1 { + t.Fatalf("tombstone snapshot = %#v, %v", snapshot, err) + } + tombstone := snapshot.RetainedCommands[0] + if !tombstone.GetTombstoned() || tombstone.GetLifecycle() != rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED || tombstone.GetCommandRevision() != 1 || tombstone.GetLastClientEventSeq() != 0 { + t.Fatalf("tombstone command = %#v", tombstone) + } +} diff --git a/test/coverage.toml b/test/coverage.toml index bc48ac5..7a2da43 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -176,6 +176,12 @@ layer = "unit" status = "implemented" tests = ["internal/client/agent/handshake_test.go:TestReconcileSnapshot_HP_SES_08"] +[[requirements]] +id = "HP-SES-09" +layer = "unit" +status = "implemented" +tests = ["internal/client/spool/snapshot_test.go:TestReconcileSnapshotRetainedAndTombstoned_HP_SES_09"] + [[requirements]] id = "HP-SES-05" layer = "integration"