From 3f84d3b2f110435f0966ee551d64d853e55afb05 Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 11:11:36 +0000 Subject: [PATCH] feat: persist and transfer script command payloads --- cmd/rvc/main.go | 89 ++++++- cmd/rvc/main_test.go | 15 ++ internal/client/agent/dispatch.go | 12 +- internal/client/agent/handshake_test.go | 37 +++ internal/client/agent/script.go | 47 ++++ internal/client/spool/script.go | 11 +- internal/client/spool/script_test.go | 4 + internal/client/supervisor/windows/launch.go | 235 ++++++++++++++++++ .../client/supervisor/windows/launch_test.go | 151 +++++++++++ internal/server/control/service.go | 19 +- internal/server/control/service_test.go | 14 +- internal/server/session/agent_server.go | 140 ++++++++++- internal/server/session/agent_server_test.go | 51 ++++ internal/server/store/command.go | 62 ++++- internal/server/store/command_test.go | 50 ++++ internal/server/store/dispatch.go | 132 +++++++++- internal/server/store/payload_recovery.go | 114 +++++++++ internal/server/store/recovery.go | 3 + 18 files changed, 1155 insertions(+), 31 deletions(-) create mode 100644 internal/client/agent/script.go create mode 100644 internal/client/supervisor/windows/launch.go create mode 100644 internal/client/supervisor/windows/launch_test.go create mode 100644 internal/server/store/payload_recovery.go diff --git a/cmd/rvc/main.go b/cmd/rvc/main.go index d78ce24..c79a485 100644 --- a/cmd/rvc/main.go +++ b/cmd/rvc/main.go @@ -40,9 +40,16 @@ func run(args []string, output, diagnostics io.Writer) error { if err != nil { return err } + requestID, args, err := globalRequestID(args) + if err != nil { + return err + } if len(args) == 0 { return errors.New("a command is required (stat or run)") } + if requestID != "" && !readOnlyCommand(args) { + args = append([]string{args[0], "--request-id", requestID}, args[1:]...) + } ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() connection, err := dial(ctx, socket) @@ -95,6 +102,44 @@ func globalSocket(args []string) (string, []string, error) { return socket, remaining, nil } +// globalRequestID removes the optional request ID wherever the CLI's global +// option parser finds it. It is injected before the subcommand's positional +// arguments so Go's flag package cannot mistake it for command text. Read-only +// commands intentionally discard the value in run(). +func globalRequestID(args []string) (string, []string, error) { + var requestID string + remaining := make([]string, 0, len(args)) + for index := 0; index < len(args); index++ { + if args[index] == "--request-id" { + if index+1 >= len(args) { + return "", nil, errors.New("--request-id requires a UUIDv7") + } + if requestID != "" { + return "", nil, errors.New("--request-id may be supplied only once") + } + requestID = args[index+1] + index++ + continue + } + if strings.HasPrefix(args[index], "--request-id=") { + if requestID != "" { + return "", nil, errors.New("--request-id may be supplied only once") + } + requestID = strings.TrimPrefix(args[index], "--request-id=") + continue + } + remaining = append(remaining, args[index]) + } + return requestID, remaining, nil +} + +func readOnlyCommand(args []string) bool { + if len(args) == 0 || args[0] == "stat" { + return true + } + return args[0] == "storage" && len(args) > 1 && args[1] == "incidents" +} + func dial(ctx context.Context, socket string) (*grpc.ClientConn, error) { dialer := func(ctx context.Context, _ string) (net.Conn, error) { return (&net.Dialer{}).DialContext(ctx, "unix", socket) @@ -165,6 +210,7 @@ func runCommand(ctx context.Context, client rvboxv1.ControlClient, args []string cwd := flags.String("cwd", "", "command working directory") shell := flags.String("shell", "", "shell (sh, bash, cmd, powershell)") requestID := flags.String("request-id", "", "canonical UUIDv7 used for idempotent admission") + elevated := flags.Bool("elevated", false, "request the Windows elevated execution policy") queueTTL := flags.Duration("queue-ttl", -1, "queue TTL; zero means no expiry") scriptPath := flags.String("script", "", "script file path") envValues := repeatedFlag{} @@ -188,7 +234,7 @@ func runCommand(ctx context.Context, client rvboxv1.ControlClient, args []string if _, err := domain.ParseUUIDv7(*requestID); err != nil { return fmt.Errorf("--request-id: %w", err) } - spec := &rvboxv1.ExecutionSpec{Cwd: *cwd} + spec := &rvboxv1.ExecutionSpec{Cwd: *cwd, Elevated: *elevated} spec.ShellType, _ = parseShell(*shell) for _, value := range envValues { parts := strings.SplitN(value, "=", 2) @@ -383,12 +429,20 @@ func killCommand(ctx context.Context, client rvboxv1.ControlClient, args []strin } func storage(ctx context.Context, client rvboxv1.ControlClient, args []string, output io.Writer) error { + flags := flag.NewFlagSet("storage", flag.ContinueOnError) + flags.SetOutput(io.Discard) + requestID := flags.String("request-id", "", "canonical UUIDv7") + includeResolved := flags.Bool("all", false, "include resolved incidents") + if err := flags.Parse(args); err != nil { + return err + } + args = flags.Args() if len(args) == 0 { return errors.New("storage requires incidents, repair, or acknowledge") } switch args[0] { case "incidents": - response, err := client.ListStorageIncidents(ctx, &rvboxv1.ListStorageIncidentsRequest{IncludeResolved: contains(args[1:], "--all")}) + response, err := client.ListStorageIncidents(ctx, &rvboxv1.ListStorageIncidentsRequest{IncludeResolved: *includeResolved}) if err != nil { return err } @@ -400,12 +454,15 @@ func storage(ctx context.Context, client rvboxv1.ControlClient, args []string, o if len(args) < 2 { return errors.New("storage mutation requires INCIDENT_ID") } - requestID, err := generatedRequestID() - if err != nil { - return err + if *requestID == "" { + var err error + *requestID, err = generatedRequestID() + if err != nil { + return err + } } if args[0] == "repair" { - response, callErr := client.RepairStorageIncident(ctx, &rvboxv1.RepairStorageIncidentRequest{IncidentId: args[1], RequestId: requestID}) + response, callErr := client.RepairStorageIncident(ctx, &rvboxv1.RepairStorageIncidentRequest{IncidentId: args[1], RequestId: *requestID}) if callErr != nil { return callErr } @@ -416,7 +473,7 @@ func storage(ctx context.Context, client rvboxv1.ControlClient, args []string, o if note == "" { return errors.New("storage acknowledge requires a note") } - response, callErr := client.AcknowledgeStorageIncident(ctx, &rvboxv1.AcknowledgeStorageIncidentRequest{IncidentId: args[1], RequestId: requestID, Note: note}) + response, callErr := client.AcknowledgeStorageIncident(ctx, &rvboxv1.AcknowledgeStorageIncidentRequest{IncidentId: args[1], RequestId: *requestID, Note: note}) if callErr != nil { return callErr } @@ -428,17 +485,27 @@ func storage(ctx context.Context, client rvboxv1.ControlClient, args []string, o } func clientCommand(ctx context.Context, client rvboxv1.ControlClient, args []string, output io.Writer) error { + flags := flag.NewFlagSet("client", flag.ContinueOnError) + flags.SetOutput(io.Discard) + requestID := flags.String("request-id", "", "canonical UUIDv7") + if err := flags.Parse(args); err != nil { + return err + } + args = flags.Args() if len(args) != 3 || args[0] != "takeover" { return errors.New("client takeover requires CLIENT INSTANCE_ID") } if _, err := domain.ParseUUIDv7(args[2]); err != nil { return err } - requestID, err := generatedRequestID() - if err != nil { - return err + if *requestID == "" { + var err error + *requestID, err = generatedRequestID() + if err != nil { + return err + } } - response, err := client.AuthorizeClientTakeover(ctx, &rvboxv1.AuthorizeClientTakeoverRequest{ClientId: args[1], ClientInstanceId: args[2], RequestId: requestID}) + response, err := client.AuthorizeClientTakeover(ctx, &rvboxv1.AuthorizeClientTakeoverRequest{ClientId: args[1], ClientInstanceId: args[2], RequestId: *requestID}) if err != nil { return err } diff --git a/cmd/rvc/main_test.go b/cmd/rvc/main_test.go index cda7a03..a141c40 100644 --- a/cmd/rvc/main_test.go +++ b/cmd/rvc/main_test.go @@ -32,3 +32,18 @@ func TestGlobalSocketAndCLIValueParsing_HP_CTL_11(t *testing.T) { t.Fatalf("zero queue duration = %v", got) } } + +func TestGlobalRequestIDIsInjectedOnlyForMutations_HP_CTL_12(t *testing.T) { + t.Parallel() + id := "019c46f1-1d02-7000-8000-0000000000f1" + requestID, remaining, err := globalRequestID([]string{"--request-id", id, "run", "client", "echo"}) + if err != nil || requestID != id || len(remaining) != 3 || remaining[0] != "run" { + t.Fatalf("global request ID = %q %#v %v", requestID, remaining, err) + } + if !readOnlyCommand([]string{"stat", "client"}) || !readOnlyCommand([]string{"storage", "incidents"}) || readOnlyCommand([]string{"run", "client", "echo"}) { + t.Fatal("read-only command classification is incorrect") + } + if _, _, err := globalRequestID([]string{"--request-id", id, "--request-id", id}); err == nil { + t.Fatal("duplicate global request IDs accepted") + } +} diff --git a/internal/client/agent/dispatch.go b/internal/client/agent/dispatch.go index dbfaa52..ee556fb 100644 --- a/internal/client/agent/dispatch.go +++ b/internal/client/agent/dispatch.go @@ -35,5 +35,15 @@ func PersistDispatch(ctx context.Context, store *spool.Store, session Session, d if err != nil { return spool.Acceptance{}, err } - return store.AcceptCommand(ctx, spool.Command{IssueUUID: issue, ImmutableSHA256: immutable, Revision: dispatch.GetCommandRevision(), Phase: uint32(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED), ExecutionSpec: spec}, now) + acceptance, err := store.AcceptCommand(ctx, spool.Command{IssueUUID: issue, ImmutableSHA256: immutable, Revision: dispatch.GetCommandRevision(), Phase: uint32(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED), ExecutionSpec: spec}, now) + if err != nil { + return spool.Acceptance{}, err + } + if descriptor := dispatch.GetSpec().GetScript(); descriptor != nil { + _, err = store.BeginScript(ctx, issue, spool.ScriptDescriptor{SizeBytes: descriptor.GetSizeBytes(), SHA256: bytesToDigest(descriptor.GetSha256())}) + if err != nil { + return spool.Acceptance{}, err + } + } + return acceptance, nil } diff --git a/internal/client/agent/handshake_test.go b/internal/client/agent/handshake_test.go index fbbbce2..7efbd56 100644 --- a/internal/client/agent/handshake_test.go +++ b/internal/client/agent/handshake_test.go @@ -3,6 +3,7 @@ package agent import ( "bytes" "context" + "crypto/sha256" "errors" "path/filepath" "testing" @@ -103,6 +104,42 @@ func TestPersistDispatchUsesImmutableHash_HP_DISPATCH_04(t *testing.T) { } } +func TestScriptFramesApplyDurablyAndReplaySafely_HP_SCRIPT_04(t *testing.T) { + store, err := spool.Open(context.Background(), spool.Options{DataDir: filepath.Join(t.TempDir(), "script-spool"), BusyTimeout: time.Second}) + if err != nil { + t.Fatal(err) + } + defer store.Close() + body := []byte("Write-Output 'hello'\r\n") + digest := sha256.Sum256(body) + dispatch := &rvboxv1.CommandDispatch{IssueUuid: "019c46f1-1d02-7000-8000-000000000067", CommandRevision: 1, TargetSessionGeneration: 9, IssueTime: timestamppb.Now(), ImmutableRequestSha256: []byte("12345678901234567890123456789012"), Spec: &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_POWERSHELL, Source: &rvboxv1.ExecutionSpec_Script{Script: &rvboxv1.ScriptDescriptor{Filename: "hello.ps1", SizeBytes: uint64(len(body)), Sha256: digest[:]}}}} + accepted, err := PersistDispatch(context.Background(), store, Session{Generation: 9}, dispatch, time.Now().UTC(), agentproto.DefaultLimits()) + if err != nil || accepted.Duplicate { + t.Fatalf("script PersistDispatch = %#v, %v", accepted, err) + } + chunkDigest := sha256.Sum256(body[:8]) + status, err := ApplyScriptChunk(context.Background(), store, Session{Generation: 9}, &rvboxv1.ScriptChunk{IssueUuid: dispatch.GetIssueUuid(), Offset: 0, Data: body[:8], Sha256: chunkDigest[:]}, agentproto.DefaultLimits()) + if err != nil || status.ReceivedBytes != 8 { + t.Fatalf("first script chunk = %#v, %v", status, err) + } + duplicate, err := ApplyScriptChunk(context.Background(), store, Session{Generation: 9}, &rvboxv1.ScriptChunk{IssueUuid: dispatch.GetIssueUuid(), Offset: 0, Data: body[:8], Sha256: chunkDigest[:]}, agentproto.DefaultLimits()) + if err != nil || !duplicate.Duplicate || duplicate.ReceivedBytes != 8 { + t.Fatalf("replayed script chunk = %#v, %v", duplicate, err) + } + restDigest := sha256.Sum256(body[8:]) + if _, err := ApplyScriptChunk(context.Background(), store, Session{Generation: 9}, &rvboxv1.ScriptChunk{IssueUuid: dispatch.GetIssueUuid(), Offset: 8, Data: body[8:], Sha256: restDigest[:]}, agentproto.DefaultLimits()); err != nil { + t.Fatal(err) + } + committed, err := ApplyScriptCommit(context.Background(), store, Session{Generation: 9}, &rvboxv1.ScriptCommit{IssueUuid: dispatch.GetIssueUuid(), SizeBytes: uint64(len(body)), Sha256: digest[:]}, agentproto.DefaultLimits()) + if err != nil || !committed.Committed || committed.ReceivedBytes != uint64(len(body)) { + t.Fatalf("script commit = %#v, %v", committed, err) + } + replayed, err := ApplyScriptChunk(context.Background(), store, Session{Generation: 9}, &rvboxv1.ScriptChunk{IssueUuid: dispatch.GetIssueUuid(), Offset: 8, Data: body[8:], Sha256: restDigest[:]}, agentproto.DefaultLimits()) + if err != nil || !replayed.Duplicate || !replayed.Committed { + t.Fatalf("post-commit chunk replay = %#v, %v", replayed, err) + } +} + func TestSendCommandEventCanonicalDigest_HP_EVENT_02(t *testing.T) { transport := &fakeTransport{} event := &rvboxv1.CommandEvent{IssueUuid: "019c46f1-1d02-7000-8000-000000000065", EventSeq: 1, ObservedAt: timestamppb.Now(), Payload: &rvboxv1.CommandEvent_Lifecycle{Lifecycle: &rvboxv1.LifecycleChange{Lifecycle: rvboxv1.CommandLifecycle_COMMAND_RUNNING, CommandRevision: 1}}} diff --git a/internal/client/agent/script.go b/internal/client/agent/script.go new file mode 100644 index 0000000..a5ac875 --- /dev/null +++ b/internal/client/agent/script.go @@ -0,0 +1,47 @@ +package agent + +import ( + "context" + "crypto/sha256" + "errors" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" + "github.com/rvbox/rvbox/internal/agentproto" + "github.com/rvbox/rvbox/internal/client/spool" + "github.com/rvbox/rvbox/internal/domain" +) + +var ErrInvalidScriptTransfer = errors.New("invalid script transfer") + +// ApplyScriptChunk validates one server script frame and appends it to the +// durable client spool. The spool accepts only the next contiguous range or +// an exact replay, so reconnects cannot create holes or duplicate bytes. +func ApplyScriptChunk(ctx context.Context, store *spool.Store, session Session, chunk *rvboxv1.ScriptChunk, limits agentproto.Limits) (spool.ScriptStatus, error) { + if store == nil || session.Generation == 0 || chunk == nil || chunk.GetOffset() > limits.MaxScriptBytes || uint64(len(chunk.GetData())) > limits.MaxRawChunkBytes || len(chunk.GetSha256()) != sha256.Size { + return spool.ScriptStatus{}, ErrInvalidScriptTransfer + } + issue, err := domain.ParseUUIDv7(chunk.GetIssueUuid()) + if err != nil || sha256.Sum256(chunk.GetData()) != bytesToDigest(chunk.GetSha256()) { + return spool.ScriptStatus{}, ErrInvalidScriptTransfer + } + return store.AppendScriptChunk(ctx, issue, chunk.GetOffset(), chunk.GetData(), bytesToDigest(chunk.GetSha256())) +} + +// ApplyScriptCommit marks a fully received script launch-eligible only after +// the spool verifies the declared size and whole-body SHA-256. +func ApplyScriptCommit(ctx context.Context, store *spool.Store, session Session, commit *rvboxv1.ScriptCommit, limits agentproto.Limits) (spool.ScriptStatus, error) { + if store == nil || session.Generation == 0 || commit == nil || commit.GetSizeBytes() > limits.MaxScriptBytes || len(commit.GetSha256()) != sha256.Size { + return spool.ScriptStatus{}, ErrInvalidScriptTransfer + } + issue, err := domain.ParseUUIDv7(commit.GetIssueUuid()) + if err != nil { + return spool.ScriptStatus{}, ErrInvalidScriptTransfer + } + return store.CommitScript(ctx, issue, spool.ScriptDescriptor{SizeBytes: commit.GetSizeBytes(), SHA256: bytesToDigest(commit.GetSha256())}) +} + +func bytesToDigest(value []byte) [sha256.Size]byte { + var digest [sha256.Size]byte + copy(digest[:], value) + return digest +} diff --git a/internal/client/spool/script.go b/internal/client/spool/script.go index cd0e2f4..8769318 100644 --- a/internal/client/spool/script.go +++ b/internal/client/spool/script.go @@ -102,7 +102,7 @@ func (store *Store) AppendScriptChunk(ctx context.Context, issueUUID domain.UUID if err != nil { return ScriptStatus{}, err } - if row.Committed { + if row.Committed && offset >= row.ReceivedBytes { return ScriptStatus{}, ErrScriptTerminal } if offset > row.ReceivedBytes { @@ -123,8 +123,17 @@ func (store *Store) AppendScriptChunk(ctx context.Context, issueUUID domain.UUID if end > row.ReceivedBytes || !bytes.Equal(raw[offset:end], data) { return ScriptStatus{}, ErrScriptOverlap } + // A reconnect can replay a chunk after the client has already committed + // the complete script. Exact byte-for-byte replays are safe and must be + // idempotent; a new byte after commit remains forbidden. + if row.Committed { + return ScriptStatus{ReceivedBytes: row.ReceivedBytes, Committed: true, Duplicate: true}, tx.Commit() + } return ScriptStatus{ReceivedBytes: row.ReceivedBytes, Duplicate: true}, tx.Commit() } + if row.Committed { + return ScriptStatus{}, ErrScriptTerminal + } updatedRaw := append(raw, data...) stored, err := compressScript(updatedRaw) if err != nil { diff --git a/internal/client/spool/script_test.go b/internal/client/spool/script_test.go index 31a9bd5..20280f8 100644 --- a/internal/client/spool/script_test.go +++ b/internal/client/spool/script_test.go @@ -63,6 +63,10 @@ func TestScriptUploadDurableExactReplayAndCommit_HP_SCRIPT_01(t *testing.T) { if err != nil || !duplicate.Duplicate || !duplicate.Committed { t.Fatalf("ScriptCommit replay = %#v, %v", duplicate, err) } + duplicate, err = store.AppendScriptChunk(ctx, issue, 0, raw[:len(first)], sha256.Sum256(raw[:len(first)])) + if err != nil || !duplicate.Duplicate || !duplicate.Committed || duplicate.ReceivedBytes != uint64(len(raw)) { + t.Fatalf("post-commit ScriptChunk replay = %#v, %v", duplicate, err) + } if _, err := store.AppendScriptChunk(ctx, issue, uint64(len(raw)), []byte("!"), sha256.Sum256([]byte("!"))); !errors.Is(err, ErrScriptTerminal) { t.Fatalf("post-commit chunk error = %v, want ErrScriptTerminal", err) } diff --git a/internal/client/supervisor/windows/launch.go b/internal/client/supervisor/windows/launch.go new file mode 100644 index 0000000..fa70357 --- /dev/null +++ b/internal/client/supervisor/windows/launch.go @@ -0,0 +1,235 @@ +package windows + +import ( + "bytes" + "crypto/sha256" + "errors" + "fmt" + "strings" + "unicode/utf8" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +// The Windows supervisor never resolves a shell through PATH. ShellPaths is +// the already-defaulted configuration snapshot supplied by the daemon; the +// native launch adapter re-stats the selected path immediately before use. +type ShellPaths struct { + CMD string + PowerShell string +} + +var ( + ErrUnsupportedShell = errors.New("shell is unsupported on Windows") + ErrInvalidShellPath = errors.New("Windows shell path must be an absolute executable path") + ErrInvalidWrapperSource = errors.New("invalid generated wrapper source") + ErrWrapperTooLarge = errors.New("generated wrapper exceeds the configured limit") + ErrWrapperDigestMismatch = errors.New("script body does not match its descriptor digest") + ErrWrapperLengthMismatch = errors.New("script body does not match its descriptor length") + ErrInvalidWrapperPath = errors.New("generated wrapper path is invalid") + ErrInvalidWorkingDirectory = errors.New("Windows working directory is invalid") +) + +// WrapperEncoding documents how the private generated file is encoded. CMD +// receives the source bytes unchanged; PowerShell gets a UTF-8 BOM so the +// inbox Windows PowerShell implementation does not interpret non-ASCII text +// using the legacy system code page. +type WrapperEncoding uint8 + +const ( + WrapperEncodingUTF8 WrapperEncoding = iota + 1 + WrapperEncodingUTF8BOM +) + +// ShellPlan is the immutable invocation contract for one supported Windows +// shell. Arguments intentionally exclude the application name: the native +// adapter passes ApplicationName separately to CreateProcessAsUser while +// CommandLine contains the complete argv-compatible command line. +type ShellPlan struct { + Type rvboxv1.ShellType + ApplicationName string + Arguments []string + WrapperExtension string + Encoding WrapperEncoding +} + +// Wrapper is the private generated script body. Filename is never derived +// from caller input; the native materializer appends Extension to an +// unpredictable per-command basename. +type Wrapper struct { + Extension string + Encoding WrapperEncoding + Bytes []byte + SHA256 [sha256.Size]byte +} + +// LaunchPlan contains only deterministic launch inputs. Handles, ACLs, Job +// Objects, token selection and release barriers belong to the platform-native +// supervisor and are deliberately absent here. +type LaunchPlan struct { + ApplicationName string + Arguments []string + CommandLine string + WrapperPath string + WorkingDirectory string + Environment []uint16 +} + +// ResolveShell chooses exactly the configured executable for shellType. It +// does not consult PATH, COMSPEC, file associations, environment overrides, +// or a wrapper's extension. +func ResolveShell(shellType rvboxv1.ShellType, paths ShellPaths) (ShellPlan, error) { + plan, err := shellTemplate(shellType) + if err != nil { + return ShellPlan{}, err + } + if shellType == rvboxv1.ShellType_SHELL_CMD { + plan.ApplicationName = paths.CMD + } else { + plan.ApplicationName = paths.PowerShell + } + if err := ValidateWindowsExecutablePath(plan.ApplicationName); err != nil { + return ShellPlan{}, fmt.Errorf("%w: %v", ErrInvalidShellPath, err) + } + return plan, nil +} + +// BuildWrapper validates and copies command/script bytes. For script sources, +// descriptor length and SHA-256 are checked before any materialization. The +// returned bytes are safe to hand to a private generated-file writer. +func BuildWrapper(spec *rvboxv1.ExecutionSpec, scriptBody []byte, maxBytes uint64) (Wrapper, error) { + if spec == nil { + return Wrapper{}, ErrInvalidWrapperSource + } + plan, err := shellTemplate(spec.GetShellType()) + if err != nil { + return Wrapper{}, err + } + var source []byte + switch value := spec.Source.(type) { + case *rvboxv1.ExecutionSpec_CommandText: + if !utf8.ValidString(value.CommandText) { + return Wrapper{}, fmt.Errorf("%w: command text is not UTF-8", ErrInvalidWrapperSource) + } + source = []byte(value.CommandText) + case *rvboxv1.ExecutionSpec_Script: + if value.Script == nil || len(value.Script.Sha256) != sha256.Size { + return Wrapper{}, fmt.Errorf("%w: script descriptor is incomplete", ErrInvalidWrapperSource) + } + if uint64(len(scriptBody)) != value.Script.SizeBytes { + return Wrapper{}, ErrWrapperLengthMismatch + } + digest := sha256.Sum256(scriptBody) + if !bytes.Equal(digest[:], value.Script.Sha256) { + return Wrapper{}, ErrWrapperDigestMismatch + } + source = append([]byte(nil), scriptBody...) + default: + return Wrapper{}, fmt.Errorf("%w: exactly one command or script source is required", ErrInvalidWrapperSource) + } + if uint64(len(source)) > maxBytes && maxBytes != 0 { + return Wrapper{}, ErrWrapperTooLarge + } + if bytes.IndexByte(source, 0) >= 0 { + return Wrapper{}, fmt.Errorf("%w: source contains NUL", ErrInvalidWrapperSource) + } + if plan.Encoding == WrapperEncodingUTF8BOM { + source = append([]byte{0xef, 0xbb, 0xbf}, source...) + } + if uint64(len(source)) > maxBytes && maxBytes != 0 { + return Wrapper{}, ErrWrapperTooLarge + } + return Wrapper{Extension: plan.WrapperExtension, Encoding: plan.Encoding, Bytes: source, SHA256: sha256.Sum256(source)}, nil +} + +func shellTemplate(shellType rvboxv1.ShellType) (ShellPlan, error) { + switch shellType { + case rvboxv1.ShellType_SHELL_CMD: + return ShellPlan{Type: shellType, Arguments: []string{"/D", "/S", "/C"}, WrapperExtension: ".cmd", Encoding: WrapperEncodingUTF8}, nil + case rvboxv1.ShellType_SHELL_POWERSHELL: + return ShellPlan{Type: shellType, Arguments: []string{"-NoLogo", "-NoProfile", "-NonInteractive", "-File"}, WrapperExtension: ".ps1", Encoding: WrapperEncodingUTF8BOM}, nil + default: + return ShellPlan{}, fmt.Errorf("%w: %s", ErrUnsupportedShell, shellType) + } +} + +// BuildLaunchPlan appends the generated wrapper path to the fixed shell +// argument vector and quotes the complete vector with the reviewed Windows +// routine. Environment is expected to be BuildEnvironmentBlock output. +func (plan ShellPlan) BuildLaunchPlan(wrapperPath, workingDirectory string, environment []uint16) (LaunchPlan, error) { + if err := ValidateWindowsExecutablePath(plan.ApplicationName); err != nil { + return LaunchPlan{}, fmt.Errorf("%w: %v", ErrInvalidShellPath, err) + } + if !validAbsoluteWindowsPath(wrapperPath) || strings.ToLower(strings.TrimSpace(wrapperPath)) != strings.ToLower(wrapperPath) || !strings.EqualFold(extensionOfWindowsPath(wrapperPath), plan.WrapperExtension) { + return LaunchPlan{}, ErrInvalidWrapperPath + } + if !validAbsoluteWindowsPath(workingDirectory) { + return LaunchPlan{}, ErrInvalidWorkingDirectory + } + arguments := append([]string(nil), plan.Arguments...) + arguments = append(arguments, wrapperPath) + full := append([]string{plan.ApplicationName}, arguments...) + commandLine, err := BuildCommandLine(full) + if err != nil { + return LaunchPlan{}, err + } + return LaunchPlan{ApplicationName: plan.ApplicationName, Arguments: arguments, CommandLine: commandLine, WrapperPath: wrapperPath, WorkingDirectory: workingDirectory, Environment: append([]uint16(nil), environment...)}, nil +} + +// ValidateWindowsExecutablePath applies the lexical checks that are possible +// before a native re-stat. The Windows adapter must still verify that the +// path names the same regular executable immediately before launch. +func ValidateWindowsExecutablePath(path string) error { + if !validAbsoluteWindowsPath(path) || strings.TrimSpace(path) != path { + return ErrInvalidShellPath + } + extension := strings.ToLower(extensionOfWindowsPath(path)) + switch extension { + case ".exe", ".com", ".cmd", ".bat": + return nil + default: + return fmt.Errorf("unsupported executable extension %q", extension) + } +} + +func validAbsoluteWindowsPath(value string) bool { + if value == "" || strings.IndexByte(value, 0) >= 0 || !utf8.ValidString(value) || strings.ContainsAny(value, "\r\n\t") { + return false + } + value = strings.ReplaceAll(value, "/", `\`) + if strings.HasPrefix(value, `\\`) { + parts := strings.Split(strings.TrimPrefix(value, `\\`), `\`) + if len(parts) < 2 || parts[0] == "" || parts[1] == "" { + return false + } + } else if len(value) < 3 || !((value[0] >= 'A' && value[0] <= 'Z') || (value[0] >= 'a' && value[0] <= 'z')) || value[1] != ':' || value[2] != '\\' { + return false + } + for _, part := range strings.Split(value, `\`) { + if part == "" || part == "." || part == ".." { + continue + } + if strings.HasSuffix(part, " ") || strings.HasSuffix(part, ".") { + return false + } + for _, character := range part { + if character < 0x20 || character == 0x7f { + return false + } + } + } + return true +} + +func extensionOfWindowsPath(path string) string { + path = strings.ReplaceAll(path, "/", `\`) + index := strings.LastIndexByte(path, '\\') + if index >= 0 { + path = path[index+1:] + } + dot := strings.LastIndexByte(path, '.') + if dot < 0 { + return "" + } + return path[dot:] +} diff --git a/internal/client/supervisor/windows/launch_test.go b/internal/client/supervisor/windows/launch_test.go new file mode 100644 index 0000000..2d30bc1 --- /dev/null +++ b/internal/client/supervisor/windows/launch_test.go @@ -0,0 +1,151 @@ +package windows + +import ( + "bytes" + "crypto/sha256" + "errors" + "strings" + "testing" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestResolveShellUsesOnlyConfiguredAbsoluteExecutable_BH_WINLAUNCH_01(t *testing.T) { + t.Parallel() + paths := ShellPaths{CMD: `C:\Trusted Tools\cmd.exe`, PowerShell: `C:\Trusted Tools\powershell.exe`} + cmd, err := ResolveShell(rvboxv1.ShellType_SHELL_CMD, paths) + if err != nil { + t.Fatal(err) + } + if cmd.ApplicationName != paths.CMD || !equalStrings(cmd.Arguments, []string{"/D", "/S", "/C"}) || cmd.WrapperExtension != ".cmd" { + t.Fatalf("cmd plan = %+v", cmd) + } + powershell, err := ResolveShell(rvboxv1.ShellType_SHELL_POWERSHELL, paths) + if err != nil { + t.Fatal(err) + } + if powershell.ApplicationName != paths.PowerShell || !equalStrings(powershell.Arguments, []string{"-NoLogo", "-NoProfile", "-NonInteractive", "-File"}) || powershell.WrapperExtension != ".ps1" || powershell.Encoding != WrapperEncodingUTF8BOM { + t.Fatalf("PowerShell plan = %+v", powershell) + } + // A hostile request PATH or COMSPEC is not an input to resolution. The + // configured absolute path remains the only application name. + if strings.Contains(cmd.ApplicationName, "untrusted") || strings.Contains(powershell.ApplicationName, "untrusted") { + t.Fatal("request environment influenced shell resolution") + } +} + +func TestResolveShellRejectsUnsupportedOrUnsafePaths_BH_WINLAUNCH_01(t *testing.T) { + t.Parallel() + cases := []struct { + name string + shell rvboxv1.ShellType + paths ShellPaths + }{ + {"unspecified", rvboxv1.ShellType_SHELL_TYPE_UNSPECIFIED, ShellPaths{CMD: `C:\Windows\System32\cmd.exe`}}, + {"unix shell", rvboxv1.ShellType_SHELL_BASH, ShellPaths{CMD: `C:\Windows\System32\cmd.exe`}}, + {"relative", rvboxv1.ShellType_SHELL_CMD, ShellPaths{CMD: `cmd.exe`}}, + {"drive relative", rvboxv1.ShellType_SHELL_CMD, ShellPaths{CMD: `C:cmd.exe`}}, + {"wrong extension", rvboxv1.ShellType_SHELL_CMD, ShellPaths{CMD: `C:\Windows\System32\cmd.dll`}}, + {"trailing space", rvboxv1.ShellType_SHELL_CMD, ShellPaths{CMD: "C:\\Windows\\System32\\cmd.exe "}}, + {"NUL", rvboxv1.ShellType_SHELL_CMD, ShellPaths{CMD: "C:\\Windows\\System32\\cmd.exe\x00"}}, + } + for _, test := range cases { + t.Run(test.name, func(t *testing.T) { + t.Parallel() + if _, err := ResolveShell(test.shell, test.paths); err == nil { + t.Fatal("unsafe shell configuration accepted") + } + }) + } +} + +func TestBuildWrapperValidatesAndEncodesSources_HP_WINLAUNCH_02(t *testing.T) { + t.Parallel() + command := &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_CMD, Source: &rvboxv1.ExecutionSpec_CommandText{CommandText: "echo hello\r\n"}} + wrapper, err := BuildWrapper(command, nil, 1024) + if err != nil { + t.Fatal(err) + } + if wrapper.Extension != ".cmd" || wrapper.Encoding != WrapperEncodingUTF8 || !bytes.Equal(wrapper.Bytes, []byte("echo hello\r\n")) { + t.Fatalf("cmd wrapper = %+v", wrapper) + } + if wrapper.SHA256 != sha256.Sum256(wrapper.Bytes) { + t.Fatal("wrapper digest does not cover materialized bytes") + } + + body := []byte("Write-Output 'héllo'\n") + digest := sha256.Sum256(body) + powershell := &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_POWERSHELL, Source: &rvboxv1.ExecutionSpec_Script{Script: &rvboxv1.ScriptDescriptor{Filename: "display.ps1", SizeBytes: uint64(len(body)), Sha256: digest[:]}}} + encoded, err := BuildWrapper(powershell, body, 1024) + if err != nil { + t.Fatal(err) + } + if encoded.Extension != ".ps1" || encoded.Encoding != WrapperEncodingUTF8BOM || !bytes.HasPrefix(encoded.Bytes, []byte{0xef, 0xbb, 0xbf}) || !bytes.HasSuffix(encoded.Bytes, body) { + t.Fatalf("PowerShell wrapper bytes = %x", encoded.Bytes) + } +} + +func TestBuildWrapperRejectsMismatchNULAndLimit_BH_WINLAUNCH_02(t *testing.T) { + t.Parallel() + body := []byte("echo expected") + digest := sha256.Sum256(body) + script := func(size uint64, hash []byte, value []byte) *rvboxv1.ExecutionSpec { + return &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_CMD, Source: &rvboxv1.ExecutionSpec_Script{Script: &rvboxv1.ScriptDescriptor{Filename: "x.cmd", SizeBytes: size, Sha256: hash}}} + } + cases := []struct { + name string + spec *rvboxv1.ExecutionSpec + body []byte + want error + }{ + {"length", script(uint64(len(body)+1), digest[:], body), body, ErrWrapperLengthMismatch}, + {"digest", script(uint64(len(body)), make([]byte, sha256.Size), body), body, ErrWrapperDigestMismatch}, + {"NUL", &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_CMD, Source: &rvboxv1.ExecutionSpec_CommandText{CommandText: "echo\x00bad"}}, nil, ErrInvalidWrapperSource}, + {"limit", &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_CMD, Source: &rvboxv1.ExecutionSpec_CommandText{CommandText: "123456"}}, nil, ErrWrapperTooLarge}, + } + for _, test := range cases { + t.Run(test.name, func(t *testing.T) { + t.Parallel() + _, err := BuildWrapper(test.spec, test.body, map[bool]uint64{true: 5, false: 1024}[test.name == "limit"]) + if !errors.Is(err, test.want) { + t.Fatalf("error = %v, want %v", err, test.want) + } + }) + } +} + +func TestBuildLaunchPlanQuotesWrapperAndKeepsWorkingDirectory(t *testing.T) { + t.Parallel() + plan, err := ResolveShell(rvboxv1.ShellType_SHELL_CMD, ShellPaths{CMD: `C:\Windows\System32\cmd.exe`}) + if err != nil { + t.Fatal(err) + } + launch, err := plan.BuildLaunchPlan(`C:\ProgramData\RVBox\work\issue one\wrapper.cmd`, `C:\ProgramData\RVBox\work\issue one`, []uint16{0}) + if err != nil { + t.Fatal(err) + } + if launch.ApplicationName != plan.ApplicationName || launch.WorkingDirectory == "" || !strings.Contains(launch.CommandLine, `"C:\ProgramData\RVBox\work\issue one\wrapper.cmd"`) { + t.Fatalf("launch plan = %+v", launch) + } + if len(launch.Environment) != 1 || launch.Environment[0] != 0 { + t.Fatalf("environment was not copied: %v", launch.Environment) + } + if _, err := plan.BuildLaunchPlan(`C:\ProgramData\RVBox\work\wrapper.ps1`, `C:\ProgramData\RVBox\work`, nil); !errors.Is(err, ErrInvalidWrapperPath) { + t.Fatalf("extension mismatch error = %v", err) + } + if _, err := plan.BuildLaunchPlan(`relative\wrapper.cmd`, `C:\ProgramData\RVBox\work`, nil); !errors.Is(err, ErrInvalidWrapperPath) { + t.Fatalf("relative wrapper error = %v", err) + } +} + +func equalStrings(left, right []string) bool { + if len(left) != len(right) { + return false + } + for index := range left { + if left[index] != right[index] { + return false + } + } + return true +} diff --git a/internal/server/control/service.go b/internal/server/control/service.go index 510d65e..cf422e7 100644 --- a/internal/server/control/service.go +++ b/internal/server/control/service.go @@ -4,6 +4,7 @@ package control import ( + "bytes" "context" "crypto/rand" "crypto/sha256" @@ -241,15 +242,20 @@ func (service *Service) RunCommand(ctx context.Context, request *rvboxv1.RunComm return nil, controlError(codes.InvalidArgument, rvboxv1.ControlError_INVALID_ARGUMENT, "target client platform is unspecified") } } - if _, script := spec.Source.(*rvboxv1.ExecutionSpec_Script); script { - return nil, status.Error(codes.Unimplemented, "script command admission is not enabled until payload dispatch is implemented") - } - if len(request.GetScriptContent()) != 0 { + _, scriptSource := spec.Source.(*rvboxv1.ExecutionSpec_Script) + if !scriptSource && len(request.GetScriptContent()) != 0 { return nil, controlError(codes.InvalidArgument, rvboxv1.ControlError_INVALID_ARGUMENT, "script_content requires a script source") } if err := agentproto.ValidateExecutionSpec(spec, service.limits, platform); err != nil { return nil, controlError(codes.InvalidArgument, rvboxv1.ControlError_INVALID_ARGUMENT, err.Error()) } + if scriptSource { + descriptor := spec.GetScript() + digest := sha256.Sum256(request.GetScriptContent()) + if descriptor == nil || uint64(len(request.GetScriptContent())) != descriptor.GetSizeBytes() || !bytes.Equal(digest[:], descriptor.GetSha256()) { + return nil, controlError(codes.InvalidArgument, rvboxv1.ControlError_INVALID_ARGUMENT, "script_content does not match the script descriptor") + } + } if !advertisedShell(client.SupportedShells, spec.GetShellType()) { return nil, controlError(codes.InvalidArgument, rvboxv1.ControlError_INVALID_ARGUMENT, "target client does not advertise the requested shell") } @@ -262,7 +268,9 @@ func (service *Service) RunCommand(ctx context.Context, request *rvboxv1.RunComm if err != nil { return nil, err } - canonical := &rvboxv1.RunCommandRequest{TargetClientId: request.GetTargetClientId(), Spec: spec, QueueTtl: request.GetQueueTtl()} + canonical := proto.Clone(request).(*rvboxv1.RunCommandRequest) + canonical.Spec = spec + canonical.RequestId = "" encoded, err := proto.MarshalOptions{Deterministic: true}.Marshal(canonical) if err != nil { return nil, controlError(codes.Internal, rvboxv1.ControlError_INTERNAL, "canonicalize command request") @@ -271,6 +279,7 @@ func (service *Service) RunCommand(ctx context.Context, request *rvboxv1.RunComm queued, err := service.store.QueueCommand(ctx, store.QueueCommandInput{ IssueUUID: issue, ClientID: request.GetTargetClientId(), IssueTime: now, ReceiptTime: now, QueueExpiryTime: expiry, ImmutableSHA256: hash, ExecutionSpec: mustMarshal(spec), + ScriptPresent: scriptSource, ScriptContent: append([]byte(nil), request.GetScriptContent()...), }) if err != nil { return nil, mapStoreError(err) diff --git a/internal/server/control/service_test.go b/internal/server/control/service_test.go index ef2ba9b..1396a4c 100644 --- a/internal/server/control/service_test.go +++ b/internal/server/control/service_test.go @@ -90,9 +90,17 @@ func TestRunCommandRequestIDIdempotencyAndValidation_BH_CONTROL_02(t *testing.T) if _, err := service.RunCommand(ctx, badID); status.Code(err) != codes.InvalidArgument { t.Fatalf("bad request ID code = %v", status.Code(err)) } - script := &rvboxv1.RunCommandRequest{TargetClientId: "win-a", RequestId: fixedIssue(0xa3).String(), Spec: &rvboxv1.ExecutionSpec{Source: &rvboxv1.ExecutionSpec_Script{Script: &rvboxv1.ScriptDescriptor{Filename: "x.ps1", SizeBytes: 1, Sha256: sha256.New().Sum(nil)}}}, ScriptContent: []byte("x")} - if _, err := service.RunCommand(ctx, script); status.Code(err) != codes.Unimplemented { - t.Fatalf("script admission code = %v", status.Code(err)) + scriptBody := []byte("x") + scriptDigest := sha256.Sum256(scriptBody) + script := &rvboxv1.RunCommandRequest{TargetClientId: "win-a", RequestId: fixedIssue(0xa3).String(), Spec: &rvboxv1.ExecutionSpec{Source: &rvboxv1.ExecutionSpec_Script{Script: &rvboxv1.ScriptDescriptor{Filename: "x.ps1", SizeBytes: 1, Sha256: scriptDigest[:]}}}, ScriptContent: scriptBody} + if queued, err := service.RunCommand(ctx, script); err != nil || queued.GetIssueUuid() != script.GetRequestId() { + t.Fatalf("script admission = %#v, %v", queued, err) + } + badScript := proto.Clone(script).(*rvboxv1.RunCommandRequest) + badScript.RequestId = fixedIssue(0xa7).String() + badScript.ScriptContent = []byte("y") + if _, err := service.RunCommand(ctx, badScript); status.Code(err) != codes.InvalidArgument { + t.Fatalf("mismatched script code = %v", status.Code(err)) } badTTL := proto.Clone(request).(*rvboxv1.RunCommandRequest) badTTL.RequestId = fixedIssue(0xa4).String() diff --git a/internal/server/session/agent_server.go b/internal/server/session/agent_server.go index 49709a6..9e49787 100644 --- a/internal/server/session/agent_server.go +++ b/internal/server/session/agent_server.go @@ -3,10 +3,12 @@ package session import ( "context" "crypto/rand" + "crypto/sha256" "encoding/base64" "errors" "fmt" "net/http" + "sort" "sync" "time" @@ -143,13 +145,26 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w reservations := make(map[string]DispatchLane) stdinSent := make(map[string]struct{}) signalSent := make(map[string]struct{}) + scriptTransfers := make(map[string]*scriptTransfer) + // A dispatch may have been durably claimed immediately before a server or + // client reconnect. Reconstruct every still-active script from the command + // store before enabling the dispatcher so an uncertain upload is replayed + // from its immutable body instead of being silently lost. + pendingScripts, err := server.Store.PendingScriptDispatches(parent, hello.GetClientId()) + if err != nil { + server.close(connection, websocket.StatusInternalError, "could not restore script transfers") + return + } + for _, pending := range pendingScripts { + scriptTransfers[pending.IssueUUID.String()] = &scriptTransfer{Body: append([]byte(nil), pending.Body...), Digest: pending.Digest} + } reconciled := make(chan struct{}) var reconcileOnce sync.Once if !capacity.UpdateAdvertised(0, 0, hello.GetMaxRunningCommands(), hello.GetMaxQueuedCommands()) { server.close(connection, websocket.StatusPolicyViolation, "invalid initial client capacity") return } - go server.dispatchLoop(sessionContext, queue, handle.DispatchWake(), reconciled, &capacityMu, capacity, reservations, stdinSent, signalSent, hello.GetClientId(), hello.GetPlatform(), encodeSessionID(sessionID), registration.Generation, cancel) + go server.dispatchLoop(sessionContext, queue, handle.DispatchWake(), reconciled, &capacityMu, capacity, reservations, stdinSent, signalSent, scriptTransfers, hello.GetClientId(), hello.GetPlatform(), encodeSessionID(sessionID), registration.Generation, cancel) encodedSessionID := encodeSessionID(sessionID) welcome, err := proto.Marshal(&rvboxv1.AgentEnvelope{ SessionId: encodedSessionID, SessionGeneration: registration.Generation, @@ -252,6 +267,26 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w continue } if event := envelope.GetCommandEvent(); event != nil { + var scriptIssue domain.UUID + if scriptStatus := event.GetScriptStatus(); scriptStatus != nil { + var parseErr error + scriptIssue, parseErr = domain.ParseUUIDv7(event.GetIssueUuid()) + if parseErr != nil { + server.close(connection, websocket.StatusPolicyViolation, "invalid script status") + return + } + capacityMu.Lock() + transfer := scriptTransfers[scriptIssue.String()] + validStatus := transfer != nil && scriptStatus.GetReceivedBytes() <= uint64(len(transfer.Body)) && scriptStatus.GetReceivedBytes() <= transfer.NextOffset + if validStatus && scriptStatus.GetComplete() { + validStatus = scriptStatus.GetReceivedBytes() == uint64(len(transfer.Body)) && transfer.CommitQueued + } + capacityMu.Unlock() + if !validStatus { + server.close(connection, websocket.StatusPolicyViolation, "invalid script progress") + return + } + } appendEvent, eventErr := eventAppendFromWire(event, hello.GetClientId(), registration.Generation, server.now()) if eventErr != nil { server.close(connection, websocket.StatusPolicyViolation, "invalid command event") @@ -281,6 +316,18 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w capacityMu.Unlock() handle.SignalDispatch() } + if scriptStatus := event.GetScriptStatus(); scriptStatus != nil { + capacityMu.Lock() + transfer := scriptTransfers[scriptIssue.String()] + if scriptStatus.GetReceivedBytes() > transfer.AcknowledgedOffset { + transfer.AcknowledgedOffset = scriptStatus.GetReceivedBytes() + } + if scriptStatus.GetComplete() { + transfer.Completed = true + } + capacityMu.Unlock() + handle.SignalDispatch() + } } } } @@ -337,7 +384,7 @@ func eventType(event *rvboxv1.CommandEvent) uint16 { // enqueueNextDispatch records the queued-to-dispatched transition before // exposing work to the network. A full data lane is a pre-write failure, so // only the owning generation can put the command back into the queue. -func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *WriterQueue, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, beforeEnqueue func(domain.UUID)) (domain.UUID, bool, error) { +func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *WriterQueue, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, beforeEnqueue func(*store.DispatchCandidate)) (domain.UUID, bool, error) { candidate, err := server.Store.ClaimNextDispatch(ctx, clientID, generation, server.now()) if err != nil || candidate == nil { return domain.UUID{}, false, err @@ -373,7 +420,7 @@ func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *Write return candidate.IssueUUID, sent, requeueErr } if beforeEnqueue != nil { - beforeEnqueue(candidate.IssueUUID) + beforeEnqueue(candidate) } if !queue.EnqueueData(Frame{Kind: FrameData, Payload: encoded}) { sent, requeueErr := requeue(ErrDispatchDataFull) @@ -382,6 +429,75 @@ func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *Write return candidate.IssueUUID, true, nil } +const scriptSendWindowBytes uint64 = 1 << 20 + +type scriptTransfer struct { + Body []byte + Digest [sha256.Size]byte + NextOffset uint64 + AcknowledgedOffset uint64 + CommitQueued bool + Completed bool +} + +// enqueueNextScript sends one bounded script frame. The dispatcher never +// queues more than a 1 MiB unacknowledged window, and the data lane remains +// bounded; a reconnect simply starts the exact payload from offset zero. +func (server *AgentServer) enqueueNextScript(queue *WriterQueue, sessionID string, generation uint64, transfers map[string]*scriptTransfer, sentMu *sync.Mutex, chunkLimit uint64) (bool, error) { + if chunkLimit == 0 { + return false, errors.New("script chunk limit is zero") + } + sentMu.Lock() + keys := make([]string, 0, len(transfers)) + for key := range transfers { + keys = append(keys, key) + } + sort.Strings(keys) + for _, key := range keys { + transfer := transfers[key] + if transfer == nil || transfer.Completed || transfer.NextOffset-transfer.AcknowledgedOffset >= scriptSendWindowBytes { + continue + } + var envelope *rvboxv1.AgentEnvelope + if transfer.NextOffset < uint64(len(transfer.Body)) { + end := transfer.NextOffset + chunkLimit + if end > uint64(len(transfer.Body)) { + end = uint64(len(transfer.Body)) + } + chunk := transfer.Body[transfer.NextOffset:end] + digest := sha256.Sum256(chunk) + envelope = &rvboxv1.AgentEnvelope{SessionId: sessionID, SessionGeneration: generation, Payload: &rvboxv1.AgentEnvelope_ScriptChunk{ScriptChunk: &rvboxv1.ScriptChunk{IssueUuid: key, Offset: transfer.NextOffset, Data: append([]byte(nil), chunk...), Sha256: digest[:]}}} + } else if !transfer.CommitQueued { + transfer.CommitQueued = true + envelope = &rvboxv1.AgentEnvelope{SessionId: sessionID, SessionGeneration: generation, Payload: &rvboxv1.AgentEnvelope_ScriptCommit{ScriptCommit: &rvboxv1.ScriptCommit{IssueUuid: key, SizeBytes: uint64(len(transfer.Body)), Sha256: transfer.Digest[:]}}} + } else { + continue + } + encoded, err := proto.Marshal(envelope) + if err != nil { + if transfer.CommitQueued && transfer.NextOffset == uint64(len(transfer.Body)) { + transfer.CommitQueued = false + } + sentMu.Unlock() + return false, err + } + if !queue.EnqueueData(Frame{Kind: FrameData, Payload: encoded}) { + if transfer.CommitQueued && transfer.NextOffset == uint64(len(transfer.Body)) { + transfer.CommitQueued = false + } + sentMu.Unlock() + return false, ErrDispatchDataFull + } + if transfer.NextOffset < uint64(len(transfer.Body)) { + transfer.NextOffset += uint64(len(envelope.GetScriptChunk().GetData())) + } + sentMu.Unlock() + return true, nil + } + sentMu.Unlock() + return false, nil +} + // enqueueNextStdin exposes one durable input intent on the essential control // lane. Session-local sent tracking suppresses duplicate frames while a live // connection remains usable; reconnecting naturally replays unacknowledged @@ -472,7 +588,7 @@ func (server *AgentServer) enqueueNextSignal(ctx context.Context, queue *WriterQ // dispatchLoop is the per-session serialized dispatcher. It waits for a // complete reconciliation result before consuming queued work, then coalesces // wakeups from local control RPCs, capacity advertisements, and acceptances. -func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, wake <-chan struct{}, reconciled <-chan struct{}, capacityMu *sync.Mutex, capacity *CapacityShadow, reservations map[string]DispatchLane, stdinSent, signalSent map[string]struct{}, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, cancel context.CancelFunc) { +func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, wake <-chan struct{}, reconciled <-chan struct{}, capacityMu *sync.Mutex, capacity *CapacityShadow, reservations map[string]DispatchLane, stdinSent, signalSent map[string]struct{}, scriptTransfers map[string]*scriptTransfer, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, cancel context.CancelFunc) { select { case <-reconciled: case <-ctx.Done(): @@ -496,6 +612,14 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, if signalQueued { continue } + scriptQueued, scriptErr := server.enqueueNextScript(queue, sessionID, generation, scriptTransfers, capacityMu, server.limits().MaxRawChunkBytes) + if scriptErr != nil && !errors.Is(scriptErr, ErrDispatchDataFull) { + server.closeForDispatchFailure(cancel) + return + } + if scriptQueued { + continue + } capacityMu.Lock() lane := capacity.Reserve() capacityMu.Unlock() @@ -503,9 +627,12 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, break } inserted := false - issue, sent, err := server.enqueueNextDispatch(ctx, queue, clientID, platform, sessionID, generation, func(issue domain.UUID) { + issue, sent, err := server.enqueueNextDispatch(ctx, queue, clientID, platform, sessionID, generation, func(candidate *store.DispatchCandidate) { capacityMu.Lock() - reservations[issue.String()] = lane + reservations[candidate.IssueUUID.String()] = lane + if candidate.ScriptPresent { + scriptTransfers[candidate.IssueUUID.String()] = &scriptTransfer{Body: append([]byte(nil), candidate.ScriptContent...), Digest: sha256.Sum256(candidate.ScriptContent)} + } inserted = true capacityMu.Unlock() }) @@ -513,6 +640,7 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, capacityMu.Lock() if inserted { delete(reservations, issue.String()) + delete(scriptTransfers, issue.String()) } capacity.Release(lane) capacityMu.Unlock() diff --git a/internal/server/session/agent_server_test.go b/internal/server/session/agent_server_test.go index 6564489..0cce729 100644 --- a/internal/server/session/agent_server_test.go +++ b/internal/server/session/agent_server_test.go @@ -1,13 +1,16 @@ package session import ( + "bytes" "context" + "crypto/sha256" "errors" "net/http" "net/http/httptest" "os" "path/filepath" "strings" + "sync" "testing" "time" @@ -20,6 +23,54 @@ import ( "google.golang.org/protobuf/types/known/timestamppb" ) +func TestScriptTransferUsesChecksummedChunksAndCommit_HP_SCRIPT_03(t *testing.T) { + t.Parallel() + body := bytes.Repeat([]byte("abcd"), 8) + digest := sha256.Sum256(body) + transfer := map[string]*scriptTransfer{"019c46f1-1d02-7000-8000-0000000000b1": {Body: body, Digest: digest}} + queue := NewWriterQueue(4, 8, 2) + defer queue.Close() + for offset := uint64(0); offset < uint64(len(body)); { + queued, err := (&AgentServer{}).enqueueNextScript(queue, "session", 7, transfer, &sync.Mutex{}, 5) + if err != nil || !queued { + t.Fatalf("chunk at offset %d: queued=%t err=%v", offset, queued, err) + } + frame, err := queue.Next(context.Background()) + if err != nil { + t.Fatal(err) + } + var envelope rvboxv1.AgentEnvelope + if err := proto.Unmarshal(frame.Payload, &envelope); err != nil { + t.Fatal(err) + } + chunk := envelope.GetScriptChunk() + if chunk == nil || chunk.GetOffset() != offset || !bytes.Equal(chunk.GetData(), body[offset:offset+uint64(len(chunk.GetData()))]) { + t.Fatalf("chunk = %+v at offset %d", chunk, offset) + } + chunkDigest := sha256.Sum256(chunk.GetData()) + if !bytes.Equal(chunk.GetSha256(), chunkDigest[:]) { + t.Fatal("chunk digest mismatch") + } + offset += uint64(len(chunk.GetData())) + } + queued, err := (&AgentServer{}).enqueueNextScript(queue, "session", 7, transfer, &sync.Mutex{}, 5) + if err != nil || !queued { + t.Fatalf("commit queued=%t err=%v", queued, err) + } + frame, err := queue.Next(context.Background()) + if err != nil { + t.Fatal(err) + } + var envelope rvboxv1.AgentEnvelope + if err := proto.Unmarshal(frame.Payload, &envelope); err != nil { + t.Fatal(err) + } + commit := envelope.GetScriptCommit() + if commit == nil || commit.GetSizeBytes() != uint64(len(body)) || !bytes.Equal(commit.GetSha256(), digest[:]) { + t.Fatalf("commit = %+v", commit) + } +} + func TestAgentServerRegistrationAndReplacement_HP_SES_05(t *testing.T) { server, cleanup := newTestAgentServer(t) defer cleanup() diff --git a/internal/server/store/command.go b/internal/server/store/command.go index 634c456..708df49 100644 --- a/internal/server/store/command.go +++ b/internal/server/store/command.go @@ -1,7 +1,9 @@ package store import ( + "bytes" "context" + "crypto/sha256" "crypto/subtle" "database/sql" "errors" @@ -9,7 +11,9 @@ import ( "time" "github.com/klauspost/compress/zstd" + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" "github.com/rvbox/rvbox/internal/domain" + "google.golang.org/protobuf/proto" ) var ( @@ -28,6 +32,12 @@ type QueueCommandInput struct { QueueExpiryTime *time.Time ImmutableSHA256 [32]byte ExecutionSpec []byte + // ScriptPresent distinguishes a valid zero-byte script from a command-text + // request with no script payload. When true, ScriptContent is checked + // against the descriptor embedded in ExecutionSpec and stored as a + // command-owned compressed payload. + ScriptPresent bool + ScriptContent []byte } type QueueCommandResult struct { @@ -47,7 +57,18 @@ func (store *Store) QueueCommand(ctx context.Context, input QueueCommandInput) ( if err != nil { return QueueCommandResult{}, err } - charge, err := EstimateCharge(ChargeInput{EncodedBytes: uint64(len(stored)), SQLiteRows: 1, IndexEntries: 2}) + scriptStored, scriptDescriptor, err := validateAndCompressScript(input) + if err != nil { + return QueueCommandResult{}, err + } + rows, indexes := uint64(1), uint64(2) + encodedBytes := uint64(len(stored)) + if input.ScriptPresent { + rows++ + indexes++ + encodedBytes += uint64(len(scriptStored)) + } + charge, err := EstimateCharge(ChargeInput{EncodedBytes: encodedBytes, SQLiteRows: rows, IndexEntries: indexes}) if err != nil { return QueueCommandResult{}, err } @@ -94,7 +115,7 @@ FROM clients JOIN storage_counters ON storage_counters.singleton = 1 WHERE clien decision, err := CheckReservation(store.quotaLimits, ReservationState{ ClientTotalCharged: clientCharged, ServerTotalCharged: serverCharged, CloseoutRemaining: store.quotaLimits.CloseoutReserveBytes, FilesystemFreeBytes: freeBytes, - }, ReservationRequest{ChargedBytes: charge, PhysicalBytes: uint64(len(stored))}) + }, ReservationRequest{ChargedBytes: charge, PhysicalBytes: encodedBytes}) if err != nil { return QueueCommandResult{}, err } @@ -109,10 +130,17 @@ closeout_remaining_bytes, immutable_request_sha256, execution_spec, execution_spec_raw_bytes, execution_spec_stored_bytes, execution_spec_compression ) VALUES (?, ?, ?, ?, ?, 1, 1, ?, ?, ?, ?, ?, ?, ?, 2)`, input.IssueUUID[:], input.ClientID, input.IssueTime.UTC().UnixNano(), input.ReceiptTime.UTC().UnixNano(), expiry, - len(stored), charge, decision.CloseoutRemaining, input.ImmutableSHA256[:], stored, len(input.ExecutionSpec), len(stored)) + encodedBytes, charge, decision.CloseoutRemaining, input.ImmutableSHA256[:], stored, len(input.ExecutionSpec), len(stored)) if err != nil { return QueueCommandResult{}, err } + if input.ScriptPresent { + if _, err := tx.ExecContext(ctx, `INSERT INTO command_payloads ( +issue_uuid, kind, raw_bytes, stored_bytes, compression, sha256, inline_data, segment_path +) VALUES (?, 'script', ?, ?, 2, ?, ?, NULL)`, input.IssueUUID[:], scriptDescriptor.GetSizeBytes(), len(scriptStored), scriptDescriptor.GetSha256(), scriptStored); err != nil { + return QueueCommandResult{}, err + } + } if _, err := tx.ExecContext(ctx, `UPDATE clients SET charged_bytes = ? WHERE client_id = ?`, decision.ClientTotalCharged, input.ClientID); err != nil { return QueueCommandResult{}, err } @@ -129,9 +157,37 @@ func validCommandInput(input QueueCommandInput) bool { if isZeroUUID([16]byte(input.IssueUUID)) || input.ClientID == "" || len(input.ClientID) > 128 || input.IssueTime.IsZero() || input.ReceiptTime.IsZero() || len(input.ExecutionSpec) == 0 || allZero(input.ImmutableSHA256[:]) { return false } + if !input.ScriptPresent && len(input.ScriptContent) != 0 { + return false + } return input.QueueExpiryTime == nil || input.QueueExpiryTime.After(input.IssueTime) } +const maxStoredScriptBytes = 10 << 20 + +func validateAndCompressScript(input QueueCommandInput) ([]byte, *rvboxv1.ScriptDescriptor, error) { + if !input.ScriptPresent { + return nil, nil, nil + } + var spec rvboxv1.ExecutionSpec + if err := proto.Unmarshal(input.ExecutionSpec, &spec); err != nil { + return nil, nil, fmt.Errorf("decode execution spec for script payload: %w", err) + } + descriptor := spec.GetScript() + if descriptor == nil || descriptor.GetSizeBytes() > maxStoredScriptBytes || len(descriptor.GetSha256()) != sha256.Size || uint64(len(input.ScriptContent)) != descriptor.GetSizeBytes() { + return nil, nil, errors.New("script payload does not match execution descriptor") + } + digest := sha256.Sum256(input.ScriptContent) + if !bytes.Equal(digest[:], descriptor.GetSha256()) { + return nil, nil, errors.New("script payload digest does not match execution descriptor") + } + stored, err := compressCommandSpec(input.ScriptContent) + if err != nil { + return nil, nil, err + } + return stored, descriptor, nil +} + func compressCommandSpec(spec []byte) ([]byte, error) { encoder, err := zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1)) if err != nil { diff --git a/internal/server/store/command_test.go b/internal/server/store/command_test.go index f368db8..32b0b9c 100644 --- a/internal/server/store/command_test.go +++ b/internal/server/store/command_test.go @@ -5,12 +5,15 @@ import ( "context" "crypto/sha256" "errors" + "fmt" "path/filepath" "testing" "time" "github.com/klauspost/compress/zstd" + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" "github.com/rvbox/rvbox/internal/domain" + "google.golang.org/protobuf/proto" ) func TestQueueCommandDurableIdempotency_HP_DISPATCH_01(t *testing.T) { @@ -123,6 +126,45 @@ func TestClaimDispatchExpiresAndFencesRequeue_HP_DISPATCH_02(t *testing.T) { } } +func TestQueueAndClaimScriptPayloadIsDurable_HP_SCRIPT_02(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() }) + if _, err := opened.RegisterClientSession(ctx, ClientRegistration{ClientID: "win-script", Platform: 2, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\`, SupportedShells: []byte{1}, ClientInstanceID: [16]byte{11}, SessionID: [16]byte{12}, ConnectedAt: time.Now().UTC()}); err != nil { + t.Fatal(err) + } + body := []byte("Write-Output 'script payload'\r\n") + digest := sha256.Sum256(body) + spec, err := proto.Marshal(&rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_POWERSHELL, Source: &rvboxv1.ExecutionSpec_Script{Script: &rvboxv1.ScriptDescriptor{Filename: "script.ps1", SizeBytes: uint64(len(body)), Sha256: digest[:]}}}) + if err != nil { + t.Fatal(err) + } + issue := fixedStoreIssue(0xb1) + now := time.Now().UTC() + if _, err := opened.QueueCommand(ctx, QueueCommandInput{IssueUUID: issue, ClientID: "win-script", IssueTime: now, ReceiptTime: now, ImmutableSHA256: sha256.Sum256([]byte("script-request")), ExecutionSpec: spec, ScriptPresent: true, ScriptContent: body}); err != nil { + t.Fatal(err) + } + var payload []byte + if err := opened.DB().QueryRowContext(ctx, `SELECT inline_data FROM command_payloads WHERE issue_uuid = ? AND kind = 'script'`, issue[:]).Scan(&payload); err != nil { + t.Fatal(err) + } + if bytes.Equal(payload, body) { + t.Fatal("script payload was stored uncompressed") + } + candidate, err := opened.ClaimNextDispatch(ctx, "win-script", 7, now.Add(time.Second)) + if err != nil || candidate == nil || !bytes.Equal(candidate.ScriptContent, body) { + t.Fatalf("script dispatch candidate = %#v, %v", candidate, err) + } + pending, err := opened.PendingScriptDispatches(ctx, "win-script") + if err != nil || len(pending) != 1 || pending[0].IssueUUID != issue || !bytes.Equal(pending[0].Body, body) || pending[0].Digest != digest { + t.Fatalf("reconstructed script transfer = %#v, %v", pending, err) + } +} + func TestRecordCommandAcceptanceFencesGeneration_HP_DISPATCH_05(t *testing.T) { t.Parallel() ctx := context.Background() @@ -156,3 +198,11 @@ func TestRecordCommandAcceptanceFencesGeneration_HP_DISPATCH_05(t *testing.T) { } func timePtr(value time.Time) *time.Time { return &value } + +func fixedStoreIssue(last byte) domain.UUID { + value, err := domain.ParseUUIDv7(fmt.Sprintf("019c46f1-1d02-7000-8000-0000000000%02x", last)) + if err != nil { + panic(err) + } + return value +} diff --git a/internal/server/store/dispatch.go b/internal/server/store/dispatch.go index 73c9da1..962dadb 100644 --- a/internal/server/store/dispatch.go +++ b/internal/server/store/dispatch.go @@ -3,6 +3,7 @@ package store import ( "bytes" "context" + "crypto/sha256" "database/sql" "errors" "fmt" @@ -10,7 +11,9 @@ import ( "time" "github.com/klauspost/compress/zstd" + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" "github.com/rvbox/rvbox/internal/domain" + "google.golang.org/protobuf/proto" ) const maxStoredExecutionSpecBytes = 1 << 20 @@ -25,6 +28,19 @@ type DispatchCandidate struct { QueueExpiryTime *time.Time ImmutableSHA256 [32]byte ExecutionSpec []byte + ScriptPresent bool + ScriptContent []byte +} + +// PendingScript is the server-owned immutable script body for a command that +// has already crossed the queued boundary. It is used to reconstruct the +// session-local transfer window after a daemon or server reconnect; the +// command row remains the source of truth for whether the work is still +// non-terminal. +type PendingScript struct { + IssueUUID domain.UUID + Body []byte + Digest [sha256.Size]byte } // ClaimNextDispatch expires old queued rows and atomically assigns the oldest @@ -79,6 +95,10 @@ FROM commands WHERE client_id = ? AND lifecycle = 1 ORDER BY issue_time, issue_u if err != nil { return nil, err } + script, scriptPresent, err := loadDispatchScript(ctx, tx, issue, spec) + if err != nil { + return nil, err + } result, err := tx.ExecContext(ctx, `UPDATE commands SET lifecycle = 2, target_session_generation = ? WHERE issue_uuid = ? AND client_id = ? AND lifecycle = 1`, generation, issue[:], clientID) if err != nil { return nil, err @@ -93,7 +113,7 @@ FROM commands WHERE client_id = ? AND lifecycle = 1 ORDER BY issue_time, issue_u if err := tx.Commit(); err != nil { return nil, err } - candidate := &DispatchCandidate{IssueUUID: issue, Revision: revision, IssueTime: time.Unix(0, issueTime).UTC(), ImmutableSHA256: immutable, ExecutionSpec: spec} + candidate := &DispatchCandidate{IssueUUID: issue, Revision: revision, IssueTime: time.Unix(0, issueTime).UTC(), ImmutableSHA256: immutable, ExecutionSpec: spec, ScriptPresent: scriptPresent, ScriptContent: script} if expiry.Valid { value := time.Unix(0, expiry.Int64).UTC() candidate.QueueExpiryTime = &value @@ -101,6 +121,116 @@ FROM commands WHERE client_id = ? AND lifecycle = 1 ORDER BY issue_time, issue_u return candidate, nil } +func (store *Store) PendingScriptDispatches(ctx context.Context, clientID string) ([]PendingScript, error) { + if clientID == "" { + return nil, errors.New("client ID is required") + } + database, err := store.openDatabase() + if err != nil { + return nil, err + } + rows, err := database.QueryContext(ctx, `SELECT issue_uuid, execution_spec, execution_spec_raw_bytes +FROM commands WHERE client_id = ? AND lifecycle BETWEEN 2 AND 4 ORDER BY issue_time, issue_uuid`, clientID) + if err != nil { + return nil, err + } + defer rows.Close() + type pendingRow struct { + issue domain.UUID + storedSpec []byte + rawSpecSize uint64 + } + var candidates []pendingRow + for rows.Next() { + var encodedIssue, stored []byte + var rawBytes uint64 + if err := rows.Scan(&encodedIssue, &stored, &rawBytes); err != nil { + return nil, err + } + if len(encodedIssue) != 16 { + return nil, ErrInvalidSegmentRecord + } + var issue domain.UUID + copy(issue[:], encodedIssue) + candidates = append(candidates, pendingRow{issue: issue, storedSpec: bytes.Clone(stored), rawSpecSize: rawBytes}) + } + if err := rows.Close(); err != nil { + return nil, err + } + var pending []PendingScript + for _, candidate := range candidates { + spec, err := decompressCommandSpec(candidate.storedSpec, candidate.rawSpecSize) + if err != nil { + return nil, err + } + body, present, err := loadDispatchScript(ctx, database, candidate.issue, spec) + if err != nil { + return nil, err + } + if !present { + continue + } + pending = append(pending, PendingScript{IssueUUID: candidate.issue, Body: body, Digest: sha256.Sum256(body)}) + } + if err := rows.Err(); err != nil { + return nil, err + } + return pending, nil +} + +type queryRower interface { + QueryRowContext(context.Context, string, ...any) *sql.Row +} + +func loadDispatchScript(ctx context.Context, tx queryRower, issue domain.UUID, encodedSpec []byte) ([]byte, bool, error) { + var spec rvboxv1.ExecutionSpec + if err := proto.Unmarshal(encodedSpec, &spec); err != nil { + // Existing store callers may use opaque test bytes. They cannot carry a + // script descriptor, so there is no payload to load. + return nil, false, nil + } + descriptor := spec.GetScript() + if descriptor == nil { + return nil, false, nil + } + var stored []byte + var rawBytes, storedBytes uint64 + var compression uint32 + var digest []byte + if err := tx.QueryRowContext(ctx, `SELECT inline_data, raw_bytes, stored_bytes, compression, sha256 +FROM command_payloads WHERE issue_uuid = ? AND kind = 'script'`, issue[:]).Scan(&stored, &rawBytes, &storedBytes, &compression, &digest); err != nil { + return nil, false, fmt.Errorf("load persisted script payload: %w", err) + } + if compression != 2 || storedBytes != uint64(len(stored)) || rawBytes != descriptor.GetSizeBytes() || len(digest) != 32 || !bytes.Equal(digest, descriptor.GetSha256()) { + return nil, false, ErrInvalidSegmentRecord + } + body, err := decompressStoredPayload(stored, rawBytes, maxStoredScriptBytes) + if err != nil { + return nil, false, err + } + computed := sha256.Sum256(body) + if !bytes.Equal(computed[:], descriptor.GetSha256()) { + return nil, false, ErrInvalidSegmentRecord + } + return body, true, nil +} + +func decompressStoredPayload(stored []byte, rawBytes, maximum uint64) ([]byte, error) { + if rawBytes > maximum || rawBytes > uint64(math.MaxInt) { + return nil, ErrInvalidSegmentRecord + } + decoder, err := zstd.NewReader(nil, zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(maximum+1)) + if err != nil { + return nil, ErrInvalidSegmentRecord + } + defer decoder.Close() + decoded, err := decoder.DecodeAll(stored, nil) + if err != nil || uint64(len(decoded)) != rawBytes || uint64(len(decoded)) > maximum { + return nil, ErrInvalidSegmentRecord + } + return bytes.Clone(decoded), nil +} + // RequeueDispatch reverses only a known pre-acceptance write failure. A later // session generation, acceptance, or reconciliation decision cannot be rolled // back by a stale writer. diff --git a/internal/server/store/payload_recovery.go b/internal/server/store/payload_recovery.go new file mode 100644 index 0000000..a9e477d --- /dev/null +++ b/internal/server/store/payload_recovery.go @@ -0,0 +1,114 @@ +package store + +import ( + "bytes" + "context" + "crypto/sha256" + "database/sql" + + "github.com/klauspost/compress/zstd" + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" + "google.golang.org/protobuf/proto" +) + +// checkCommandPayloads verifies the inline compressed payload table and its +// relationship to script descriptors. A valid SQLite page is not enough: a +// short/changed zstd frame must make the owning scope dirty before dispatch. +func checkCommandPayloads(ctx context.Context, database *sql.DB) error { + rows, err := database.QueryContext(ctx, `SELECT issue_uuid, kind, raw_bytes, stored_bytes, compression, sha256, inline_data, segment_path FROM command_payloads`) + if err != nil { + return err + } + defer rows.Close() + for rows.Next() { + var issue, digest, stored []byte + var kind string + var rawBytes, storedBytes uint64 + var compression uint32 + var segmentPath sql.NullString + if err := rows.Scan(&issue, &kind, &rawBytes, &storedBytes, &compression, &digest, &stored, &segmentPath); err != nil { + return err + } + if len(issue) != 16 || kind != "script" || compression != 2 || rawBytes > maxStoredScriptBytes || storedBytes != uint64(len(stored)) || len(digest) != sha256.Size || segmentPath.Valid || len(stored) == 0 { + return ErrCommittedRangeMissing + } + decoded, err := decodeCommandPayload(stored, rawBytes, maxStoredScriptBytes) + if err != nil { + return err + } + computed := sha256.Sum256(decoded) + if !bytes.Equal(computed[:], digest) { + return ErrCommittedRangeMissing + } + } + if err := rows.Err(); err != nil { + return err + } + + commands, err := database.QueryContext(ctx, `SELECT issue_uuid, execution_spec FROM commands`) + if err != nil { + return err + } + defer commands.Close() + for commands.Next() { + var issue, encoded []byte + if err := commands.Scan(&issue, &encoded); err != nil { + return err + } + if len(issue) != 16 { + return ErrCommittedRangeMissing + } + var spec rvboxv1.ExecutionSpec + if err := proto.Unmarshal(encoded, &spec); err != nil { + // Some low-level store tests intentionally use opaque bytes. Such + // records cannot claim a script source, and dispatch validation will + // reject them before any user code runs. + continue + } + descriptor := spec.GetScript() + var payloadCount uint64 + if err := database.QueryRowContext(ctx, `SELECT count(*) FROM command_payloads WHERE issue_uuid = ? AND kind = 'script'`, issue).Scan(&payloadCount); err != nil { + return err + } + if descriptor == nil { + if payloadCount != 0 { + return ErrCommittedRangeMissing + } + continue + } + if payloadCount != 1 || descriptor.GetSizeBytes() > maxStoredScriptBytes || len(descriptor.GetSha256()) != sha256.Size { + return ErrCommittedRangeMissing + } + var rawBytes, storedBytes uint64 + var compression uint32 + var digest, stored []byte + if err := database.QueryRowContext(ctx, `SELECT raw_bytes, stored_bytes, compression, sha256, inline_data FROM command_payloads WHERE issue_uuid = ? AND kind = 'script'`, issue).Scan(&rawBytes, &storedBytes, &compression, &digest, &stored); err != nil { + return err + } + if rawBytes != descriptor.GetSizeBytes() || compression != 2 || len(digest) != sha256.Size || !bytes.Equal(digest, descriptor.GetSha256()) { + return ErrCommittedRangeMissing + } + decoded, err := decodeCommandPayload(stored, rawBytes, maxStoredScriptBytes) + computed := sha256.Sum256(decoded) + if err != nil || !bytes.Equal(computed[:], descriptor.GetSha256()) { + return ErrCommittedRangeMissing + } + } + return commands.Err() +} + +func decodeCommandPayload(stored []byte, expected, maximum uint64) ([]byte, error) { + if expected > maximum { + return nil, ErrCommittedRangeMissing + } + decoder, err := zstd.NewReader(nil, zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(maximum+1)) + if err != nil { + return nil, ErrCommittedRangeMissing + } + defer decoder.Close() + decoded, err := decoder.DecodeAll(stored, nil) + if err != nil || uint64(len(decoded)) != expected || uint64(len(decoded)) > maximum { + return nil, ErrCommittedRangeMissing + } + return decoded, nil +} diff --git a/internal/server/store/recovery.go b/internal/server/store/recovery.go index cfd38d0..debcd05 100644 --- a/internal/server/store/recovery.go +++ b/internal/server/store/recovery.go @@ -122,6 +122,9 @@ FROM output_segments ORDER BY issue_uuid, ordinal`) return report, err } } + if err := checkCommandPayloads(ctx, database); err != nil { + return report, err + } if err := checkQuotaCounters(ctx, database); err != nil { return report, err }