From 467ca65f606ec9dad6dfec40662f9dd37064fb24 Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 07:34:01 +0000 Subject: [PATCH] feat: dispatch queued commands after reconciliation --- internal/server/session/agent_server.go | 77 +++++++++- internal/server/session/session.go | 3 + internal/server/store/command_test.go | 53 +++++++ internal/server/store/dispatch.go | 137 ++++++++++++++++++ internal/server/store/migrations.go | 2 +- test/coverage.toml | 12 ++ .../clientagent_integration_test.go | 55 +++++++ 7 files changed, 337 insertions(+), 2 deletions(-) create mode 100644 internal/server/store/dispatch.go diff --git a/internal/server/session/agent_server.go b/internal/server/session/agent_server.go index 12f0e8e..cd915fc 100644 --- a/internal/server/session/agent_server.go +++ b/internal/server/session/agent_server.go @@ -28,6 +28,7 @@ var ( ErrUnexpectedOrigin = errors.New("agent connections must not send an Origin header") ErrUnexpectedMessage = errors.New("agent message is not a binary protobuf envelope") ErrStaleSession = errors.New("agent message is for an unknown or fenced session") + ErrDispatchDataFull = errors.New("session dispatch data lane is full") ) // AgentServer is the transport edge for agent WebSocket sessions. Command @@ -137,6 +138,11 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w queue := NewWriterQueue(16, 64, 8) defer queue.Close() + capacity := &CapacityShadow{} + if !capacity.UpdateAdvertised(0, 0, hello.GetMaxRunningCommands(), hello.GetMaxQueuedCommands()) { + server.close(connection, websocket.StatusPolicyViolation, "invalid initial client capacity") + return + } encodedSessionID := encodeSessionID(sessionID) welcome, err := proto.Marshal(&rvboxv1.AgentEnvelope{ SessionId: encodedSessionID, SessionGeneration: registration.Generation, @@ -194,14 +200,80 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w SessionId: encodedSessionID, SessionGeneration: registration.Generation, Payload: &rvboxv1.AgentEnvelope_ReconcileResult{ReconcileResult: result}, }) - if err != nil || queue.EnqueueControl(Frame{Kind: FrameControl, Payload: encoded}) != nil { + resultWritten := make(chan struct{}) + if err != nil || queue.EnqueueControl(Frame{Kind: FrameControl, Payload: encoded, Written: resultWritten}) != nil { server.close(connection, websocket.StatusInternalError, "could not queue reconciliation result") return } + select { + case <-resultWritten: + case <-sessionContext.Done(): + return + } + lane := capacity.Reserve() + dispatched, dispatchErr := false, error(nil) + if lane != DispatchNone { + dispatched, dispatchErr = server.enqueueNextDispatch(sessionContext, queue, hello.GetClientId(), hello.GetPlatform(), encodedSessionID, registration.Generation) + if !dispatched { + capacity.Release(lane) + } + } + if dispatchErr != nil && !errors.Is(dispatchErr, ErrDispatchDataFull) { + server.close(connection, websocket.StatusInternalError, "could not queue command dispatch") + return + } + continue + } + if advertised := envelope.GetClientCapacity(); advertised != nil { + if !capacity.UpdateAdvertised(advertised.GetRunningCommands(), advertised.GetQueuedCommands(), advertised.GetMaxRunningCommands(), advertised.GetMaxQueuedCommands()) { + server.close(connection, websocket.StatusPolicyViolation, "invalid client capacity") + return + } } } } +// 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) (bool, error) { + candidate, err := server.Store.ClaimNextDispatch(ctx, clientID, generation, server.now()) + if err != nil || candidate == nil { + return false, err + } + requeue := func(cause error) (bool, error) { + _, rollbackErr := server.Store.RequeueDispatch(ctx, candidate.IssueUUID, clientID, generation) + if rollbackErr != nil { + return false, rollbackErr + } + return false, cause + } + spec := &rvboxv1.ExecutionSpec{} + if err := proto.Unmarshal(candidate.ExecutionSpec, spec); err != nil { + return requeue(fmt.Errorf("decode persisted execution spec: %w", err)) + } + if err := agentproto.ValidateExecutionSpec(spec, server.limits(), platform); err != nil { + return requeue(fmt.Errorf("validate persisted execution spec: %w", err)) + } + dispatch := &rvboxv1.CommandDispatch{ + IssueUuid: candidate.IssueUUID.String(), CommandRevision: candidate.Revision, TargetSessionGeneration: generation, + IssueTime: timestamppb.New(candidate.IssueTime), Spec: spec, + } + if candidate.QueueExpiryTime != nil { + dispatch.QueueExpiryTime = timestamppb.New(*candidate.QueueExpiryTime) + } + encoded, err := proto.Marshal(&rvboxv1.AgentEnvelope{ + SessionId: sessionID, SessionGeneration: generation, Payload: &rvboxv1.AgentEnvelope_CommandDispatch{CommandDispatch: dispatch}, + }) + if err != nil { + return requeue(err) + } + if !queue.EnqueueData(Frame{Kind: FrameData, Payload: encoded}) { + return requeue(ErrDispatchDataFull) + } + return true, nil +} + func (server *AgentServer) writeLoop(ctx context.Context, connection *websocket.Conn, queue *WriterQueue, heartbeat *synchronizedHeartbeat, started time.Time) { for { frameContext, cancel := context.WithTimeout(ctx, server.heartbeatPollInterval()) @@ -214,6 +286,9 @@ func (server *AgentServer) writeLoop(ctx context.Context, connection *websocket. if err != nil { return } + if frame.Written != nil { + close(frame.Written) + } continue } if !errors.Is(err, context.DeadlineExceeded) { diff --git a/internal/server/session/session.go b/internal/server/session/session.go index fcd4dc4..0390218 100644 --- a/internal/server/session/session.go +++ b/internal/server/session/session.go @@ -23,6 +23,9 @@ const ( type Frame struct { Kind FrameKind Payload []byte + // Written is closed by the sole socket writer after a successful write. + // It is used only for protocol barriers such as reconciliation-before-work. + Written chan<- struct{} } // WriterQueue is owned by one socket writer. Data saturation leaves the work diff --git a/internal/server/store/command_test.go b/internal/server/store/command_test.go index cddf628..5473ad2 100644 --- a/internal/server/store/command_test.go +++ b/internal/server/store/command_test.go @@ -71,3 +71,56 @@ func TestQueueCommandDurableIdempotency_HP_DISPATCH_01(t *testing.T) { t.Fatalf("missing client error = %v", err) } } + +func TestClaimDispatchExpiresAndFencesRequeue_HP_DISPATCH_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-client", Platform: 2, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\`, SupportedShells: []byte{1}, + ClientInstanceID: [16]byte{3}, SessionID: [16]byte{4}, ConnectedAt: time.Date(2026, time.September, 6, 12, 0, 0, 0, time.UTC), + }); err != nil { + t.Fatal(err) + } + now := time.Date(2026, time.September, 6, 12, 0, 1, 0, time.UTC) + expiredIssue, err := domain.ParseUUIDv7("019c46f1-1d02-7000-8000-000000000093") + if err != nil { + t.Fatal(err) + } + queuedIssue, err := domain.ParseUUIDv7("019c46f1-1d02-7000-8000-000000000094") + if err != nil { + t.Fatal(err) + } + for _, input := range []QueueCommandInput{ + {IssueUUID: expiredIssue, ClientID: "win-client", IssueTime: now.Add(-time.Minute), ReceiptTime: now.Add(-time.Minute), QueueExpiryTime: timePtr(now.Add(-time.Second)), ImmutableSHA256: sha256.Sum256([]byte("expired")), ExecutionSpec: []byte("expired spec")}, + {IssueUUID: queuedIssue, ClientID: "win-client", IssueTime: now, ReceiptTime: now, ImmutableSHA256: sha256.Sum256([]byte("queued")), ExecutionSpec: []byte("queued spec")}, + } { + if _, err := opened.QueueCommand(ctx, input); err != nil { + t.Fatal(err) + } + } + candidate, err := opened.ClaimNextDispatch(ctx, "win-client", 7, now) + if err != nil || candidate == nil || candidate.IssueUUID != queuedIssue || candidate.Revision != 1 || string(candidate.ExecutionSpec) != "queued spec" { + t.Fatalf("dispatch candidate = %#v, %v", candidate, err) + } + if requeued, err := opened.RequeueDispatch(ctx, queuedIssue, "win-client", 6); err != nil || requeued { + t.Fatalf("stale requeue = %t, %v", requeued, err) + } + if requeued, err := opened.RequeueDispatch(ctx, queuedIssue, "win-client", 7); err != nil || !requeued { + t.Fatalf("owned requeue = %t, %v", requeued, err) + } + candidate, err = opened.ClaimNextDispatch(ctx, "win-client", 8, now) + if err != nil || candidate == nil || candidate.IssueUUID != queuedIssue { + t.Fatalf("reclaimed candidate = %#v, %v", candidate, err) + } + var lifecycle, revision uint64 + if err := opened.DB().QueryRow(`SELECT lifecycle, revision FROM commands WHERE issue_uuid = ?`, expiredIssue[:]).Scan(&lifecycle, &revision); err != nil || lifecycle != 10 || revision != 2 { + t.Fatalf("expired command lifecycle/revision = %d/%d, %v", lifecycle, revision, err) + } +} + +func timePtr(value time.Time) *time.Time { return &value } diff --git a/internal/server/store/dispatch.go b/internal/server/store/dispatch.go new file mode 100644 index 0000000..0675b37 --- /dev/null +++ b/internal/server/store/dispatch.go @@ -0,0 +1,137 @@ +package store + +import ( + "bytes" + "context" + "database/sql" + "errors" + "fmt" + "math" + "time" + + "github.com/klauspost/compress/zstd" + "github.com/rvbox/rvbox/internal/domain" +) + +const maxStoredExecutionSpecBytes = 1 << 20 + +// DispatchCandidate is the immutable work reconstructed after an atomic +// queued-to-dispatched transition. The session writer must use the supplied +// generation and call RequeueDispatch if its network write never starts. +type DispatchCandidate struct { + IssueUUID domain.UUID + Revision uint64 + IssueTime time.Time + QueueExpiryTime *time.Time + ExecutionSpec []byte +} + +// ClaimNextDispatch expires old queued rows and atomically assigns the oldest +// remaining row to one live client-session generation. It deliberately records +// the dispatch before the wire write so uncertain writes are reconciled rather +// than accidentally delivered twice. +func (store *Store) ClaimNextDispatch(ctx context.Context, clientID string, generation uint64, now time.Time) (*DispatchCandidate, error) { + if clientID == "" || generation == 0 || now.IsZero() { + return nil, errors.New("invalid dispatch claim") + } + store.writeMu.Lock() + defer store.writeMu.Unlock() + database, err := store.openDatabase() + if err != nil { + return nil, err + } + tx, err := database.BeginTx(ctx, nil) + if err != nil { + return nil, err + } + defer tx.Rollback() + // Expiry is terminal. The revision advance distinguishes it from a delayed + // queued snapshot or a late dispatch attempt. + if _, err := tx.ExecContext(ctx, `UPDATE commands SET lifecycle = 10, terminal_time = ?, revision = revision + 1 +WHERE client_id = ? AND lifecycle = 1 AND queue_expiry_time IS NOT NULL AND queue_expiry_time <= ?`, now.UTC().UnixNano(), clientID, now.UTC().UnixNano()); err != nil { + return nil, err + } + var encodedIssue, stored []byte + var revision uint64 + var issueTime int64 + var expiry sql.NullInt64 + var rawBytes uint64 + err = tx.QueryRowContext(ctx, `SELECT issue_uuid, revision, issue_time, queue_expiry_time, execution_spec, execution_spec_raw_bytes +FROM commands WHERE client_id = ? AND lifecycle = 1 ORDER BY issue_time, issue_uuid LIMIT 1`, clientID).Scan(&encodedIssue, &revision, &issueTime, &expiry, &stored, &rawBytes) + if errors.Is(err, sql.ErrNoRows) { + if err := tx.Commit(); err != nil { + return nil, err + } + return nil, nil + } + if err != nil { + return nil, err + } + if len(encodedIssue) != 16 { + return nil, ErrInvalidSegmentRecord + } + var issue domain.UUID + copy(issue[:], encodedIssue) + spec, err := decompressCommandSpec(stored, rawBytes) + 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 + } + changed, err := result.RowsAffected() + if err != nil { + return nil, err + } + if changed != 1 { + return nil, errors.New("queued command changed during dispatch claim") + } + if err := tx.Commit(); err != nil { + return nil, err + } + candidate := &DispatchCandidate{IssueUUID: issue, Revision: revision, IssueTime: time.Unix(0, issueTime).UTC(), ExecutionSpec: spec} + if expiry.Valid { + value := time.Unix(0, expiry.Int64).UTC() + candidate.QueueExpiryTime = &value + } + return candidate, 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. +func (store *Store) RequeueDispatch(ctx context.Context, issueUUID domain.UUID, clientID string, generation uint64) (bool, error) { + if isZeroUUID([16]byte(issueUUID)) || clientID == "" || generation == 0 { + return false, errors.New("invalid dispatch requeue") + } + store.writeMu.Lock() + defer store.writeMu.Unlock() + database, err := store.openDatabase() + if err != nil { + return false, err + } + result, err := database.ExecContext(ctx, `UPDATE commands SET lifecycle = 1, target_session_generation = NULL +WHERE issue_uuid = ? AND client_id = ? AND lifecycle = 2 AND target_session_generation = ?`, issueUUID[:], clientID, generation) + if err != nil { + return false, err + } + changed, err := result.RowsAffected() + return changed == 1, err +} + +func decompressCommandSpec(stored []byte, rawBytes uint64) ([]byte, error) { + if rawBytes == 0 || rawBytes > maxStoredExecutionSpecBytes || rawBytes > uint64(math.MaxInt) { + return nil, errors.New("invalid stored execution spec size") + } + decoder, err := zstd.NewReader(nil, zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(maxStoredExecutionSpecBytes+1)) + if err != nil { + return nil, err + } + defer decoder.Close() + decoded, err := decoder.DecodeAll(stored, nil) + if err != nil || uint64(len(decoded)) != rawBytes || uint64(len(decoded)) > maxStoredExecutionSpecBytes { + return nil, fmt.Errorf("invalid stored execution spec") + } + return bytes.Clone(decoded), nil +} diff --git a/internal/server/store/migrations.go b/internal/server/store/migrations.go index 2305c0d..4e658b9 100644 --- a/internal/server/store/migrations.go +++ b/internal/server/store/migrations.go @@ -75,7 +75,7 @@ CREATE TABLE commands ( client_id TEXT NOT NULL REFERENCES clients(client_id) ON DELETE RESTRICT, issue_time INTEGER NOT NULL, server_receipt_time INTEGER NOT NULL, queue_expiry_time INTEGER, terminal_time INTEGER, lifecycle INTEGER NOT NULL CHECK(lifecycle BETWEEN 1 AND 11), - revision INTEGER NOT NULL CHECK(revision > 0), exit_code INTEGER, + revision INTEGER NOT NULL CHECK(revision > 0), target_session_generation INTEGER CHECK(target_session_generation IS NULL OR target_session_generation > 0), exit_code INTEGER, retention_status INTEGER NOT NULL DEFAULT 1 CHECK(retention_status BETWEEN 1 AND 3), last_event_seq INTEGER NOT NULL DEFAULT 0 CHECK(last_event_seq >= 0), retained_compressed_bytes INTEGER NOT NULL DEFAULT 0 CHECK(retained_compressed_bytes >= 0), diff --git a/test/coverage.toml b/test/coverage.toml index 6992e3d..82eaae8 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -200,6 +200,18 @@ layer = "unit" status = "implemented" tests = ["internal/server/store/command_test.go:TestQueueCommandDurableIdempotency_HP_DISPATCH_01"] +[[requirements]] +id = "HP-DISPATCH-02" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/command_test.go:TestClaimDispatchExpiresAndFencesRequeue_HP_DISPATCH_02"] + +[[requirements]] +id = "HP-DISPATCH-03" +layer = "integration" +status = "implemented" +tests = ["test/integration/clientagent/clientagent_integration_test.go:TestWebSocketDispatchAfterReconciliation_HP_DISPATCH_03"] + [[requirements]] id = "HP-SES-05" layer = "integration" diff --git a/test/integration/clientagent/clientagent_integration_test.go b/test/integration/clientagent/clientagent_integration_test.go index 1a78d6c..25dd13b 100644 --- a/test/integration/clientagent/clientagent_integration_test.go +++ b/test/integration/clientagent/clientagent_integration_test.go @@ -2,6 +2,7 @@ package clientagent_test import ( "context" + "crypto/sha256" "net/http/httptest" "path/filepath" "strings" @@ -11,8 +12,10 @@ import ( rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" "github.com/rvbox/rvbox/internal/agentproto" "github.com/rvbox/rvbox/internal/client/agent" + "github.com/rvbox/rvbox/internal/domain" "github.com/rvbox/rvbox/internal/server/session" "github.com/rvbox/rvbox/internal/server/store" + "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/timestamppb" ) @@ -55,3 +58,55 @@ func TestWebSocketHelloWelcome_HP_SES_06(t *testing.T) { t.Fatalf("discard local terminal issues = %v", got) } } + +func TestWebSocketDispatchAfterReconciliation_HP_DISPATCH_03(t *testing.T) { + ctx := context.Background() + persistence, err := store.Open(ctx, store.Options{DataDir: filepath.Join(t.TempDir(), "state"), BusyTimeout: time.Second}) + if err != nil { + t.Fatal(err) + } + defer persistence.Close() + server := httptest.NewServer(&session.AgentServer{Store: persistence, Registry: session.NewRegistry(), Path: "/v1/agent", Limits: agentproto.DefaultLimits()}) + defer server.Close() + address := "ws" + strings.TrimPrefix(server.URL, "http") + "/v1/agent" + transport, err := agent.DialWebSocket(ctx, address, nil) + if err != nil { + t.Fatal(err) + } + defer transport.Close() + hello := &rvboxv1.ClientHello{ClientId: "win-dispatch-client", SupportedProtocol: &rvboxv1.ProtocolRange{Major: 1, MinMinor: 0, MaxMinor: 0}, DaemonVersion: "test", Platform: rvboxv1.Platform_PLATFORM_WINDOWS, Architecture: "amd64", DaemonCwd: `C:\`, SupportedShells: []rvboxv1.ShellType{rvboxv1.ShellType_SHELL_POWERSHELL}, ClientInstanceId: "019c46f1-1d02-7000-8000-000000000074", MaxRunningCommands: 1, MaxQueuedCommands: 1, SentAt: timestamppb.Now()} + accepted, err := agent.Handshake(ctx, transport, hello, agentproto.DefaultLimits()) + if err != nil { + t.Fatal(err) + } + issue, err := domain.ParseUUIDv7("019c46f1-1d02-7000-8000-000000000075") + if err != nil { + t.Fatal(err) + } + spec := &rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_POWERSHELL, Source: &rvboxv1.ExecutionSpec_CommandText{CommandText: "Write-Output dispatch"}} + serialized, err := proto.Marshal(spec) + if err != nil { + t.Fatal(err) + } + requestHash := sha256.Sum256([]byte("dispatch request")) + if _, err := persistence.QueueCommand(ctx, store.QueueCommandInput{IssueUUID: issue, ClientID: hello.GetClientId(), IssueTime: time.Now().UTC(), ReceiptTime: time.Now().UTC(), ImmutableSHA256: requestHash, ExecutionSpec: serialized}); err != nil { + t.Fatal(err) + } + if _, err := agent.Reconcile(ctx, transport, accepted, &rvboxv1.ReconcileSnapshot{}, agentproto.DefaultLimits()); err != nil { + t.Fatal(err) + } + readContext, cancel := context.WithTimeout(ctx, time.Second) + defer cancel() + encoded, err := transport.Read(readContext) + if err != nil { + t.Fatal(err) + } + envelope, err := agentproto.DecodeEnvelope(encoded, agentproto.DefaultLimits(), rvboxv1.Platform_PLATFORM_WINDOWS) + if err != nil || envelope.GetSessionId() != accepted.ID || envelope.GetSessionGeneration() != accepted.Generation { + t.Fatalf("dispatch envelope = %#v, %v", envelope, err) + } + dispatch := envelope.GetCommandDispatch() + if dispatch == nil || dispatch.GetIssueUuid() != issue.String() || dispatch.GetCommandRevision() != 1 || dispatch.GetTargetSessionGeneration() != accepted.Generation || !proto.Equal(dispatch.GetSpec(), spec) { + t.Fatalf("dispatch = %#v", dispatch) + } +}