feat: dispatch queued commands after reconciliation
This commit is contained in:
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 }
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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),
|
||||
|
||||
Reference in New Issue
Block a user