369 lines
15 KiB
Go
369 lines
15 KiB
Go
package spool
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/sha256"
|
|
"errors"
|
|
"net/url"
|
|
"os"
|
|
"path/filepath"
|
|
"testing"
|
|
"time"
|
|
|
|
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
|
|
|
|
"github.com/rvbox/rvbox/internal/domain"
|
|
"google.golang.org/protobuf/proto"
|
|
)
|
|
|
|
func TestSQLiteFileURIWindowsDrivePath(t *testing.T) {
|
|
query := url.Values{"_pragma": {"journal_mode(WAL)"}}
|
|
if got, want := sqliteFileURI(`C:\ProgramData\RVBox\test-state\spool.db`, query), "file:///C:/ProgramData/RVBox/test-state/spool.db?_pragma=journal_mode%28WAL%29"; got != want {
|
|
t.Fatalf("sqliteFileURI() = %q, want %q", got, want)
|
|
}
|
|
}
|
|
|
|
func TestExecutionSpecIsDurableAndQuotaCounted_HP_DISPATCH_07(t *testing.T) {
|
|
t.Parallel()
|
|
ctx := context.Background()
|
|
directory := filepath.Join(t.TempDir(), "spool")
|
|
store := openTestStore(t, ctx, directory, DefaultTombstoneLimit)
|
|
issue := testUUID(t, "019c46f1-1d02-7000-8000-000000000002")
|
|
spec := []byte("deterministic execution specification")
|
|
command := testCommand(issue, []byte("spec-hash"))
|
|
command.ExecutionSpec = spec
|
|
acceptedAt := time.Date(2026, time.September, 6, 12, 0, 0, 0, time.UTC)
|
|
accepted, err := store.AcceptCommand(ctx, command, acceptedAt)
|
|
if err != nil || accepted.Duplicate {
|
|
t.Fatalf("spec acceptance = %#v, %v", accepted, err)
|
|
}
|
|
if !bytes.Equal(accepted.Command.ExecutionSpec, spec) {
|
|
t.Fatalf("accepted spec = %q, want %q", accepted.Command.ExecutionSpec, spec)
|
|
}
|
|
var base, total, specCharge uint64
|
|
if err := store.db.QueryRowContext(ctx, `SELECT base_charged_bytes, total_charged_bytes FROM commands WHERE issue_uuid = ?`, issue[:]).Scan(&base, &total); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := store.db.QueryRowContext(ctx, `SELECT charged_bytes FROM command_specs WHERE issue_uuid = ?`, issue[:]).Scan(&specCharge); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if specCharge == 0 || total != base+specCharge {
|
|
t.Fatalf("spec charge accounting = base=%d total=%d spec=%d", base, total, specCharge)
|
|
}
|
|
duplicate, err := store.AcceptCommand(ctx, command, acceptedAt.Add(time.Second))
|
|
if err != nil || !duplicate.Duplicate || !bytes.Equal(duplicate.Command.ExecutionSpec, spec) {
|
|
t.Fatalf("spec replay = %#v, %v", duplicate, err)
|
|
}
|
|
if err := store.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
store = openTestStore(t, ctx, directory, DefaultTombstoneLimit)
|
|
reopened, err := store.AcceptCommand(ctx, command, acceptedAt.Add(2*time.Second))
|
|
if err != nil || !reopened.Duplicate || !bytes.Equal(reopened.Command.ExecutionSpec, spec) {
|
|
t.Fatalf("spec replay after reopen = %#v, %v", reopened, err)
|
|
}
|
|
report, err := store.Check(ctx)
|
|
if err != nil || report.SpecsChecked != 1 {
|
|
t.Fatalf("spec recovery report = %#v, %v", report, err)
|
|
}
|
|
}
|
|
|
|
func TestExecutionSpecCorruptionMarksSpoolDirty_BH_DISPATCH_07(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-000000000003")
|
|
command := testCommand(issue, []byte("spec-corrupt"))
|
|
command.ExecutionSpec = []byte("payload")
|
|
if _, err := store.AcceptCommand(ctx, command, time.Now().UTC()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := store.db.ExecContext(ctx, `UPDATE command_specs SET payload = ? WHERE issue_uuid = ?`, []byte("tampered"), issue[:]); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := store.Check(ctx); !errors.Is(err, ErrStoredPayloadChecksum) {
|
|
t.Fatalf("corrupt execution spec Check error = %v, want ErrStoredPayloadChecksum", err)
|
|
}
|
|
}
|
|
|
|
func TestStoreDurableIdentityAndSingleInstance_HP_CLIENT_01(t *testing.T) {
|
|
t.Parallel()
|
|
ctx := context.Background()
|
|
directory := filepath.Join(t.TempDir(), "spool")
|
|
store := openTestStore(t, ctx, directory, DefaultTombstoneLimit)
|
|
identity := store.ClientInstanceID()
|
|
if identity == (domain.UUID{}) {
|
|
t.Fatal("client instance ID is zero")
|
|
}
|
|
if _, err := Open(ctx, Options{DataDir: directory, BusyTimeout: time.Second}); !errors.Is(err, ErrAlreadyOpen) {
|
|
t.Fatalf("second Open error = %v, want ErrAlreadyOpen", err)
|
|
}
|
|
if err := store.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
reopened := openTestStore(t, ctx, directory, DefaultTombstoneLimit)
|
|
if got := reopened.ClientInstanceID(); got != identity {
|
|
t.Fatalf("reopened client instance ID = %s, want %s", got, identity)
|
|
}
|
|
}
|
|
|
|
func TestStoreRejectsCorruptIdentity_BH_CLIENT_01(t *testing.T) {
|
|
t.Parallel()
|
|
directory := t.TempDir()
|
|
if err := os.Chmod(directory, 0o700); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := os.WriteFile(filepath.Join(directory, "client-instance-id"), []byte("not-an-id\n"), 0o600); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
_, err := Open(context.Background(), Options{DataDir: directory, BusyTimeout: time.Second})
|
|
if !errors.Is(err, ErrIdentityCorrupt) {
|
|
t.Fatalf("Open corrupt identity error = %v, want ErrIdentityCorrupt", err)
|
|
}
|
|
data, err := os.ReadFile(filepath.Join(directory, "client-instance-id"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if string(data) != "not-an-id\n" {
|
|
t.Fatalf("corrupt identity was rewritten as %q", data)
|
|
}
|
|
}
|
|
|
|
func TestSpoolAcceptanceSequencingAndAck_HP_CLIENT_07(t *testing.T) {
|
|
t.Parallel()
|
|
ctx := context.Background()
|
|
directory := filepath.Join(t.TempDir(), "spool")
|
|
store := openTestStore(t, ctx, directory, DefaultTombstoneLimit)
|
|
issue := testUUID(t, "019c46f1-1d02-7000-8000-000000000001")
|
|
command := testCommand(issue, []byte("immutable command"))
|
|
acceptedAt := time.Date(2026, time.September, 6, 12, 0, 0, 0, time.UTC)
|
|
accepted, err := store.AcceptCommand(ctx, command, acceptedAt)
|
|
if err != nil || accepted.Duplicate {
|
|
t.Fatalf("first acceptance = %#v, %v", accepted, err)
|
|
}
|
|
duplicate, err := store.AcceptCommand(ctx, command, acceptedAt.Add(time.Second))
|
|
if err != nil || !duplicate.Duplicate || duplicate.Command.Revision != command.Revision {
|
|
t.Fatalf("duplicate acceptance = %#v, %v", duplicate, err)
|
|
}
|
|
first, err := store.AppendEvent(ctx, issue, EventInput{Kind: 1, Compression: 1, RawBytes: 5, Payload: []byte("first"), CreatedAt: acceptedAt})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
second, err := store.AppendEvent(ctx, issue, EventInput{Kind: 2, Compression: 1, RawBytes: 6, Payload: []byte("second"), CreatedAt: acceptedAt.Add(time.Nanosecond)})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if first.LocalOrdinal != 1 || first.EventSeq != 0 || second.LocalOrdinal != 2 || second.EventSeq != 0 {
|
|
t.Fatalf("unassigned events = %#v, %#v", first, second)
|
|
}
|
|
assigned, err := store.AssignSendWindow(ctx, issue, 1, 1<<20)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(assigned) != 1 || assigned[0].LocalOrdinal != 1 || assigned[0].EventSeq != 1 {
|
|
t.Fatalf("first send window = %#v", assigned)
|
|
}
|
|
assigned, err = store.AssignSendWindow(ctx, issue, 1, 1<<20)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(assigned) != 2 || assigned[1].LocalOrdinal != 2 || assigned[1].EventSeq != 2 {
|
|
t.Fatalf("replayed send window = %#v", assigned)
|
|
}
|
|
if err := store.Ack(ctx, issue, 1); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
pending, err := store.PendingEvents(ctx, issue)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(pending) != 1 || pending[0].EventSeq != 2 || string(pending[0].Payload) != "second" {
|
|
t.Fatalf("pending after cumulative ack = %#v", pending)
|
|
}
|
|
if err := store.Ack(ctx, issue, 3); !errors.Is(err, ErrInvalidEventAck) {
|
|
t.Fatalf("ack beyond send window error = %v, want ErrInvalidEventAck", err)
|
|
}
|
|
if err := store.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
store = openTestStore(t, ctx, directory, DefaultTombstoneLimit)
|
|
pending, err = store.PendingEvents(ctx, issue)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(pending) != 1 || pending[0].EventSeq != 2 {
|
|
t.Fatalf("pending after reopen = %#v", pending)
|
|
}
|
|
}
|
|
|
|
func TestSendWindowPinsUnacknowledgedBytes_BH_CLIENT_13(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-000000000014")
|
|
command := testCommand(issue, []byte("send-window"))
|
|
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)
|
|
}
|
|
for index, payload := range [][]byte{[]byte("first"), []byte("second")} {
|
|
if _, err := store.AppendEvent(ctx, issue, EventInput{Kind: uint32(index + 1), Compression: 1, Payload: payload, CreatedAt: now.Add(time.Duration(index) * time.Second)}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
assigned, err := store.AssignSendWindow(ctx, issue, 1, uint64(len("first")))
|
|
if err != nil || len(assigned) != 1 || assigned[0].EventSeq != 1 {
|
|
t.Fatalf("first bounded window = %#v, %v", assigned, err)
|
|
}
|
|
assigned, err = store.AssignSendWindow(ctx, issue, 1, uint64(len("first")))
|
|
if err != nil || len(assigned) != 1 || assigned[0].EventSeq != 1 {
|
|
t.Fatalf("full pinned window assigned more data = %#v, %v", assigned, err)
|
|
}
|
|
if err := store.Ack(ctx, issue, 1); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
assigned, err = store.AssignSendWindow(ctx, issue, 1, uint64(len("second")))
|
|
if err != nil || len(assigned) != 1 || assigned[0].EventSeq != 2 {
|
|
t.Fatalf("window after cumulative ack = %#v, %v", assigned, err)
|
|
}
|
|
}
|
|
|
|
func TestAppendLifecycleAtomicallyUpdatesPhaseAndEvent_HP_CLIENT_11(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-000000000021")
|
|
command := testCommand(issue, []byte("lifecycle"))
|
|
command.Phase = uint32(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED)
|
|
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)
|
|
}
|
|
effective := rvboxv1.WindowsExecutionContext_WINDOWS_EXECUTION_CONTEXT_ACTIVE_USER
|
|
identity := &rvboxv1.WindowsExecutionIdentity{EffectiveContext: &effective, SessionId: ptrUint32(1), SessionUserSid: "S-1-5-21-user", EffectiveUserSid: "S-1-5-21-user", AttemptedContexts: []rvboxv1.WindowsExecutionContext{effective}, SelectionDetail: "selected"}
|
|
if _, err := store.AppendLifecycleWithIdentity(ctx, issue, uint32(rvboxv1.CommandLifecycle_COMMAND_RUNNING), 1, "launch authorized", now, identity); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := store.AppendLifecycle(ctx, issue, uint32(rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED), 1, "exit 0", now.Add(time.Second)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var phase uint32
|
|
var terminal int
|
|
if err := store.db.QueryRow(`SELECT phase, terminal FROM commands WHERE issue_uuid = ?`, issue[:]).Scan(&phase, &terminal); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if phase != uint32(rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED) || terminal != 1 {
|
|
t.Fatalf("phase/terminal = %d/%d", phase, terminal)
|
|
}
|
|
if _, err := store.AppendLifecycle(ctx, issue, uint32(rvboxv1.CommandLifecycle_COMMAND_RUNNING), 1, "illegal", now.Add(2*time.Second)); err == nil {
|
|
t.Fatal("terminal lifecycle regressed")
|
|
}
|
|
if _, err := store.AssignSendWindow(ctx, issue, 8, 1<<20); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
events, err := store.PendingEvents(ctx, issue)
|
|
if err != nil || len(events) != 2 {
|
|
t.Fatalf("lifecycle events = %#v, %v", events, err)
|
|
}
|
|
var decoded rvboxv1.CommandEvent
|
|
if err := proto.Unmarshal(events[0].Payload, &decoded); err != nil || decoded.GetLifecycle().GetWindowsExecutionIdentity().GetEffectiveContext() != effective {
|
|
t.Fatalf("running identity event = %v, %v", decoded.GetLifecycle(), err)
|
|
}
|
|
}
|
|
|
|
func ptrUint32(value uint32) *uint32 { return &value }
|
|
|
|
func TestTerminalCleanupTombstonesAndConflicts_BH_CLIENT_02(t *testing.T) {
|
|
t.Parallel()
|
|
ctx := context.Background()
|
|
store := openTestStore(t, ctx, filepath.Join(t.TempDir(), "spool"), 1)
|
|
first := testCommand(testUUID(t, "019c46f1-1d02-7000-8000-000000000011"), []byte("first"))
|
|
second := testCommand(testUUID(t, "019c46f1-1d02-7000-8000-000000000012"), []byte("second"))
|
|
now := time.Date(2026, time.September, 6, 12, 0, 0, 0, time.UTC)
|
|
for _, command := range []Command{first, second} {
|
|
if _, err := store.AcceptCommand(ctx, command, now); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := store.MarkTerminal(ctx, command.IssueUUID, 5); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := store.CleanupTerminal(ctx, command.IssueUUID, now); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
now = now.Add(time.Nanosecond)
|
|
}
|
|
if _, err := store.AcceptCommand(ctx, second, now); !errors.Is(err, ErrAlreadyExecuted) {
|
|
t.Fatalf("exact tombstone replay error = %v, want ErrAlreadyExecuted", err)
|
|
}
|
|
conflict := second
|
|
conflict.ImmutableSHA256 = sha256.Sum256([]byte("changed"))
|
|
if _, err := store.AcceptCommand(ctx, conflict, now); !errors.Is(err, ErrCommandConflict) {
|
|
t.Fatalf("conflicting tombstone replay error = %v, want ErrCommandConflict", err)
|
|
}
|
|
if _, err := store.AcceptCommand(ctx, first, now); err != nil {
|
|
t.Fatalf("oldest tombstone was not rotated: %v", err)
|
|
}
|
|
}
|
|
|
|
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})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() {
|
|
if err := store.Close(); err != nil {
|
|
t.Error(err)
|
|
}
|
|
})
|
|
return store
|
|
}
|
|
|
|
func testUUID(t *testing.T, value string) domain.UUID {
|
|
t.Helper()
|
|
parsed, err := domain.ParseUUIDv7(value)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return parsed
|
|
}
|
|
|
|
func testCommand(issue domain.UUID, immutable []byte) Command {
|
|
return Command{IssueUUID: issue, ImmutableSHA256: sha256.Sum256(immutable), Revision: 1, Phase: 1}
|
|
}
|