From 0a0b6f9707843906de9220a73cf8cdb7ad46126c Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 07:22:17 +0000 Subject: [PATCH] feat: apply terminal reconciliation cleanup --- internal/client/agent/handshake_test.go | 25 ++++++++ internal/client/agent/reconcile.go | 54 +++++++++++++++++ internal/client/spool/events.go | 77 +++++++++++++++++++++++++ internal/client/spool/spool_test.go | 34 +++++++++++ test/coverage.toml | 12 ++++ 5 files changed, 202 insertions(+) create mode 100644 internal/client/agent/reconcile.go diff --git a/internal/client/agent/handshake_test.go b/internal/client/agent/handshake_test.go index 9178402..6052c2c 100644 --- a/internal/client/agent/handshake_test.go +++ b/internal/client/agent/handshake_test.go @@ -8,6 +8,7 @@ import ( rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" "github.com/rvbox/rvbox/internal/agentproto" + "github.com/rvbox/rvbox/internal/domain" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/timestamppb" ) @@ -64,11 +65,35 @@ func TestReconcileSnapshot_HP_SES_08(t *testing.T) { } } +func TestApplyReconcileResult_HP_SES_12(t *testing.T) { + t.Parallel() + active := "019c46f1-1d02-7000-8000-000000000062" + terminal := "019c46f1-1d02-7000-8000-000000000063" + store := &reconcileStore{} + terminated, err := ApplyReconcileResult(context.Background(), store, &rvboxv1.ReconcileResult{ + TerminateLocalIssueUuids: []string{active}, + DiscardLocalTerminalIssueUuids: []string{terminal}, + }, time.Date(2026, time.September, 6, 0, 0, 0, 0, time.UTC)) + if err != nil || len(terminated) != 1 || terminated[0].String() != active || len(store.discarded) != 1 || store.discarded[0].String() != terminal { + t.Fatalf("ApplyReconcileResult = terminated=%v discarded=%v err=%v", terminated, store.discarded, err) + } + if _, err := ApplyReconcileResult(context.Background(), store, &rvboxv1.ReconcileResult{TerminateLocalIssueUuids: []string{active}, DiscardLocalTerminalIssueUuids: []string{active}}, time.Now()); !errors.Is(err, ErrInvalidReconcileResult) { + t.Fatalf("overlapping result error = %v", err) + } +} + type fakeTransport struct { written []byte read []byte } +type reconcileStore struct{ discarded []domain.UUID } + +func (store *reconcileStore) DiscardTerminal(_ context.Context, issue domain.UUID, _ time.Time) error { + store.discarded = append(store.discarded, issue) + return nil +} + func (transport *fakeTransport) Write(_ context.Context, value []byte) error { transport.written = append([]byte(nil), value...) return nil diff --git a/internal/client/agent/reconcile.go b/internal/client/agent/reconcile.go new file mode 100644 index 0000000..ffab928 --- /dev/null +++ b/internal/client/agent/reconcile.go @@ -0,0 +1,54 @@ +package agent + +import ( + "context" + "errors" + "time" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" + "github.com/rvbox/rvbox/internal/domain" +) + +var ErrInvalidReconcileResult = errors.New("invalid reconciliation result") + +// ReconcileStore is the narrowly scoped durable operation required after the +// server has confirmed a terminal command. Keeping the interface small lets the +// network loop remain independent of SQLite and of the Windows supervisor. +type ReconcileStore interface { + DiscardTerminal(context.Context, domain.UUID, time.Time) error +} + +// ApplyReconcileResult commits server-authorized terminal discards before the +// session starts new work. It returns the active command IDs that the caller +// must first terminate through the supervisor; marking those terminal is not +// safe until the owned process and output readers have stopped. +func ApplyReconcileResult(ctx context.Context, store ReconcileStore, result *rvboxv1.ReconcileResult, now time.Time) ([]domain.UUID, error) { + if store == nil || result == nil || now.IsZero() { + return nil, ErrInvalidReconcileResult + } + seen := make(map[domain.UUID]struct{}, len(result.GetTerminateLocalIssueUuids())+len(result.GetDiscardLocalTerminalIssueUuids())) + terminate := make([]domain.UUID, 0, len(result.GetTerminateLocalIssueUuids())) + for _, set := range [][]string{result.GetTerminateLocalIssueUuids(), result.GetDiscardLocalTerminalIssueUuids()} { + for _, text := range set { + issue, err := domain.ParseUUIDv7(text) + if err != nil { + return nil, ErrInvalidReconcileResult + } + if _, duplicate := seen[issue]; duplicate { + return nil, ErrInvalidReconcileResult + } + seen[issue] = struct{}{} + } + } + for _, text := range result.GetDiscardLocalTerminalIssueUuids() { + issue, _ := domain.ParseUUIDv7(text) + if err := store.DiscardTerminal(ctx, issue, now); err != nil { + return nil, err + } + } + for _, text := range result.GetTerminateLocalIssueUuids() { + issue, _ := domain.ParseUUIDv7(text) + terminate = append(terminate, issue) + } + return terminate, nil +} diff --git a/internal/client/spool/events.go b/internal/client/spool/events.go index 19fadaf..aeb4c6a 100644 --- a/internal/client/spool/events.go +++ b/internal/client/spool/events.go @@ -345,6 +345,83 @@ func (store *Store) CleanupTerminal(ctx context.Context, issueUUID domain.UUID, return tx.Commit() } +// DiscardTerminal removes a server-confirmed terminal command even when its +// local event rows have not been acknowledged. The compact tombstone is kept so +// a late duplicate dispatch cannot execute it again. It is intentionally only +// for a ReconcileResult that authoritatively tells the client to discard local +// terminal state; ordinary terminal cleanup must use CleanupTerminal. +func (store *Store) DiscardTerminal(ctx context.Context, issueUUID domain.UUID, acknowledgedAt time.Time) error { + if !validUUID(issueUUID) || acknowledgedAt.IsZero() { + return errors.New("invalid terminal discard") + } + tx, err := store.db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + + var hash []byte + var revision uint64 + var phase uint32 + var terminal int + var totalCharged uint64 + 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 errors.Is(err, sql.ErrNoRows) { + var tombstoneHash []byte + err = tx.QueryRowContext(ctx, `SELECT immutable_sha256 FROM command_tombstones WHERE issue_uuid = ?`, issueUUID[:]).Scan(&tombstoneHash) + if errors.Is(err, sql.ErrNoRows) { + return ErrUnknownCommand + } + if err != nil { + return err + } + if len(tombstoneHash) != 32 { + return fmt.Errorf("invalid terminal tombstone") + } + return tx.Commit() + } + if err != nil { + return err + } + if terminal == 0 { + return errors.New("command is not terminal") + } + clientTotal, err := clientTotalCharge(ctx, tx) + if err != nil { + return err + } + if totalCharged > clientTotal { + return fmt.Errorf("client spool aggregate charge counter mismatch") + } + var existingHash []byte + err = tx.QueryRowContext(ctx, `SELECT immutable_sha256 FROM command_tombstones WHERE issue_uuid = ?`, issueUUID[:]).Scan(&existingHash) + if err == nil && (len(existingHash) != len(hash) || string(existingHash) != string(hash)) { + return ErrCommandConflict + } + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + if errors.Is(err, sql.ErrNoRows) { + 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 { + return err + } + if err := updateClientTotalCharge(ctx, tx, clientTotal-totalCharged); err != nil { + return err + } + trim, err := tombstonesToTrim(ctx, tx, store.tombstoneLimit) + if err != nil { + return err + } + if _, err := tx.ExecContext(ctx, `DELETE FROM command_tombstones WHERE issue_uuid IN (SELECT issue_uuid FROM command_tombstones ORDER BY acknowledged_at, issue_uuid LIMIT ?)`, trim); err != nil { + return err + } + return tx.Commit() +} + type queryer interface { QueryContext(context.Context, string, ...any) (*sql.Rows, error) } diff --git a/internal/client/spool/spool_test.go b/internal/client/spool/spool_test.go index b025d42..2de09f4 100644 --- a/internal/client/spool/spool_test.go +++ b/internal/client/spool/spool_test.go @@ -154,6 +154,40 @@ func TestTerminalCleanupTombstonesAndConflicts_BH_CLIENT_02(t *testing.T) { } } +func TestDiscardServerConfirmedTerminalDropsPendingEvents_HP_CLIENT_10(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-000000000013") + command := testCommand(issue, []byte("server-confirmed terminal")) + now := time.Date(2026, time.September, 6, 12, 0, 0, 0, time.UTC) + if _, err := store.AcceptCommand(ctx, command, now); err != nil { + t.Fatal(err) + } + if _, err := store.AppendEvent(ctx, issue, EventInput{Kind: 1, Compression: 1, Payload: []byte("unacknowledged local terminal output"), CreatedAt: now}); err != nil { + t.Fatal(err) + } + if err := store.MarkTerminal(ctx, issue, 5); err != nil { + t.Fatal(err) + } + if err := store.DiscardTerminal(ctx, issue, now); err != nil { + t.Fatal(err) + } + if pending, err := store.PendingEvents(ctx, issue); err != nil || len(pending) != 0 { + t.Fatalf("discarded command pending events = %#v, %v", pending, err) + } + snapshot, err := store.ReconcileSnapshot(ctx) + if err != nil || len(snapshot.GetRetainedCommands()) != 1 || !snapshot.GetRetainedCommands()[0].GetTombstoned() { + t.Fatalf("discarded command snapshot = %#v, %v", snapshot, err) + } + if err := store.DiscardTerminal(ctx, issue, now.Add(time.Second)); err != nil { + t.Fatalf("repeated discard = %v", err) + } + if _, err := store.AcceptCommand(ctx, command, now); !errors.Is(err, ErrAlreadyExecuted) { + t.Fatalf("duplicate after discard error = %v", err) + } +} + func openTestStore(t *testing.T, ctx context.Context, directory string, tombstoneLimit uint64) *Store { t.Helper() store, err := Open(ctx, Options{DataDir: directory, BusyTimeout: time.Second, TombstoneLimit: tombstoneLimit}) diff --git a/test/coverage.toml b/test/coverage.toml index 7a2da43..b07a6fd 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -128,6 +128,12 @@ layer = "unit" status = "implemented" tests = ["internal/client/spool/spool_test.go:TestTerminalCleanupTombstonesAndConflicts_BH_CLIENT_02"] +[[requirements]] +id = "HP-CLIENT-10" +layer = "unit" +status = "implemented" +tests = ["internal/client/spool/spool_test.go:TestDiscardServerConfirmedTerminalDropsPendingEvents_HP_CLIENT_10"] + [[requirements]] id = "HP-WINCTX-01" layer = "unit" @@ -182,6 +188,12 @@ layer = "unit" status = "implemented" tests = ["internal/client/spool/snapshot_test.go:TestReconcileSnapshotRetainedAndTombstoned_HP_SES_09"] +[[requirements]] +id = "HP-SES-12" +layer = "unit" +status = "implemented" +tests = ["internal/client/agent/handshake_test.go:TestApplyReconcileResult_HP_SES_12"] + [[requirements]] id = "HP-SES-05" layer = "integration"