271 lines
9.4 KiB
Go
271 lines
9.4 KiB
Go
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
|
|
}
|