package store import ( "bytes" "context" "crypto/sha256" "database/sql" "errors" "fmt" "math" "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 // 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 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 // 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, immutableHash []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, immutable_request_sha256, 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, &immutableHash, &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 || len(immutableHash) != 32 { return nil, ErrInvalidSegmentRecord } var issue domain.UUID copy(issue[:], encodedIssue) var immutable [32]byte copy(immutable[:], immutableHash) spec, err := decompressCommandSpec(stored, rawBytes) 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 } 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(), ImmutableSHA256: immutable, ExecutionSpec: spec, ScriptPresent: scriptPresent, ScriptContent: script} if expiry.Valid { value := time.Unix(0, expiry.Int64).UTC() candidate.QueueExpiryTime = &value } 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. 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 }