use std::collections::HashMap;
use rabs_protocol::generation::{
ActionGeneration, ActionGenerationId, AttemptAuthority, LeaseRenewal, WorkerBootGeneration,
WorkerIncarnationId,
};
use rabs_protocol::release_authorization::ReleaseVerdict;
use rabs_protocol::result_identity::{DigestAlgorithm, TypedDigest};
use rabs_protocol::serving::ServingValidity;
use rabs_protocol::wire_time::PeerId;
use rabs_protocol::worker_fence::{
WorkerAdmission, WorkerIncarnationFenceRecord, WorkerLeaseBindingRejection, WorkerSessionOffer,
};
pub const SCHEMA_VERSION: u32 = 25;
pub struct Migration {
pub version: u32,
pub statements: &'static [&'static str],
}
pub const MIGRATIONS: &[Migration] = &[
Migration {
version: 1,
statements: &[
"CREATE TABLE schema_epochs (version INTEGER PRIMARY KEY, applied_seq INTEGER NOT NULL)",
"CREATE TABLE coordinator_authorities (key TEXT PRIMARY KEY, algo TEXT NOT NULL, \
domain TEXT NOT NULL, bytes BLOB NOT NULL, cluster_id TEXT NOT NULL, \
incarnation BLOB NOT NULL, term INTEGER NOT NULL, acquired_seq INTEGER NOT NULL, \
released INTEGER NOT NULL)",
"CREATE TABLE action_entries (key TEXT PRIMARY KEY, algo TEXT NOT NULL, \
domain TEXT NOT NULL, bytes BLOB NOT NULL, key_epoch INTEGER NOT NULL, \
projection_epoch INTEGER NOT NULL)",
"CREATE TABLE action_generations (id_hex TEXT PRIMARY KEY, id BLOB NOT NULL, \
action_key TEXT NOT NULL, authority_key TEXT NOT NULL, tombstoned INTEGER NOT NULL)",
"CREATE TABLE generation_high_water (kind TEXT PRIMARY KEY, value BLOB NOT NULL)",
"CREATE TABLE action_attempts (id_hex TEXT PRIMARY KEY, id BLOB NOT NULL, \
generation_hex TEXT NOT NULL, worker TEXT NOT NULL, seq INTEGER NOT NULL)",
"CREATE TABLE execution_leases (id_hex TEXT PRIMARY KEY, id BLOB NOT NULL, \
attempt_hex TEXT NOT NULL, renewal_seq INTEGER NOT NULL, \
expires_at_seq INTEGER NOT NULL, released INTEGER NOT NULL)",
"CREATE TABLE action_publications (action_key TEXT PRIMARY KEY, \
descriptor_algo TEXT NOT NULL, descriptor_domain TEXT NOT NULL, \
descriptor_bytes BLOB NOT NULL, manifest_algo TEXT NOT NULL, \
manifest_domain TEXT NOT NULL, manifest_bytes BLOB NOT NULL, \
winner_generation_hex TEXT NOT NULL, winner_attempt_hex TEXT NOT NULL, \
result_kind TEXT NOT NULL, pin_hex TEXT NOT NULL)",
"CREATE TABLE action_serving_states (action_key TEXT PRIMARY KEY, \
disposition TEXT NOT NULL, version INTEGER NOT NULL)",
"CREATE TABLE objects (key TEXT PRIMARY KEY, algo TEXT NOT NULL, \
domain TEXT NOT NULL, bytes BLOB NOT NULL, logical_size INTEGER NOT NULL)",
"CREATE TABLE object_locations (object_key TEXT NOT NULL, store_path TEXT NOT NULL, \
verified_seq INTEGER, PRIMARY KEY (object_key, store_path))",
"CREATE TABLE pins (id_hex TEXT PRIMARY KEY, id BLOB NOT NULL, root_key TEXT NOT NULL, \
owner TEXT NOT NULL, class TEXT NOT NULL, expires_at_seq INTEGER, \
released INTEGER NOT NULL)",
"CREATE TABLE observed_input_recipes (action_key TEXT PRIMARY KEY, \
recipe_algo TEXT NOT NULL, recipe_domain TEXT NOT NULL, recipe_bytes BLOB NOT NULL)",
"CREATE TABLE key_breakdowns (action_key TEXT NOT NULL, component TEXT NOT NULL, \
algo TEXT NOT NULL, domain TEXT NOT NULL, bytes BLOB NOT NULL, \
PRIMARY KEY (action_key, component))",
"CREATE TABLE trust_states (action_key TEXT PRIMARY KEY, state TEXT NOT NULL, \
reason TEXT NOT NULL)",
"CREATE TABLE quarantines (scope TEXT NOT NULL, subject TEXT NOT NULL, \
reason TEXT NOT NULL, PRIMARY KEY (scope, subject))",
"CREATE TABLE verification_samples (action_key TEXT NOT NULL, attempt_hex TEXT NOT NULL, \
passed INTEGER NOT NULL, seq INTEGER NOT NULL, \
PRIMARY KEY (action_key, attempt_hex, seq))",
"CREATE TABLE gc_runs (id INTEGER PRIMARY KEY, seq INTEGER NOT NULL, \
pinned_roots INTEGER NOT NULL, located_objects INTEGER NOT NULL)",
],
},
Migration {
version: 2,
statements: &[
"CREATE TABLE action_evidence_index (action_key TEXT NOT NULL, \
evidence_algo TEXT NOT NULL, evidence_domain TEXT NOT NULL, \
evidence_bytes BLOB NOT NULL, generation_hex TEXT NOT NULL, \
attempt_hex TEXT NOT NULL, PRIMARY KEY (action_key, evidence_domain, evidence_bytes))",
],
},
Migration {
version: 3,
statements: &[
"CREATE TABLE object_edges (parent_key TEXT NOT NULL, child_key TEXT NOT NULL, \
kind TEXT NOT NULL, PRIMARY KEY (parent_key, child_key, kind))",
"ALTER TABLE object_locations ADD COLUMN encoding TEXT NOT NULL DEFAULT 'raw'",
"ALTER TABLE object_locations ADD COLUMN quarantined INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE pins ADD COLUMN evidence TEXT",
"ALTER TABLE pins ADD COLUMN renewal_seq INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE pins ADD COLUMN durable INTEGER NOT NULL DEFAULT 1",
"ALTER TABLE pins ADD COLUMN reason TEXT NOT NULL DEFAULT ''",
"ALTER TABLE gc_runs ADD COLUMN reachable_objects INTEGER NOT NULL DEFAULT 0",
],
},
Migration {
version: 4,
statements: &[
"CREATE TABLE gc_receipts (id INTEGER PRIMARY KEY, seq INTEGER NOT NULL, \
mode TEXT NOT NULL, planned INTEGER NOT NULL, reclaimed INTEGER NOT NULL, \
skipped INTEGER NOT NULL, truncated INTEGER NOT NULL)",
],
},
Migration {
version: 5,
statements: &["CREATE TABLE gc_tombstones (object_key TEXT NOT NULL, \
store_path TEXT NOT NULL, marked_seq INTEGER NOT NULL, \
grace_until_seq INTEGER NOT NULL, PRIMARY KEY (object_key, store_path))"],
},
Migration {
version: 6,
statements: &[
"CREATE TABLE eviction_tombstones (action_key TEXT PRIMARY KEY, \
semantic_algo TEXT NOT NULL, semantic_domain TEXT NOT NULL, \
semantic_bytes BLOB NOT NULL, observable_algo TEXT NOT NULL, \
observable_domain TEXT NOT NULL, observable_bytes BLOB NOT NULL, \
evicted_seq INTEGER NOT NULL)",
],
},
Migration {
version: 7,
statements: &[
"CREATE TABLE operator_resets (generation INTEGER PRIMARY KEY, \
applied_seq INTEGER NOT NULL)",
],
},
Migration {
version: 8,
statements: &[
"CREATE TABLE peer_authority_high_water (peer_id TEXT PRIMARY KEY, \
term INTEGER NOT NULL, observed_seq INTEGER NOT NULL)",
"CREATE TABLE worker_incarnation_fences (worker TEXT PRIMARY KEY, \
incarnation BLOB NOT NULL)",
"CREATE TABLE edge_incarnation_fences (edge_id TEXT PRIMARY KEY, \
incarnation BLOB NOT NULL)",
"CREATE TABLE edge_handoffs (edge_id TEXT PRIMARY KEY, \
active_incarnation BLOB NOT NULL, predecessor_incarnation BLOB NOT NULL, \
begun_seq INTEGER NOT NULL, resolved INTEGER NOT NULL)",
"CREATE TABLE action_trust_evaluations (action_key TEXT NOT NULL, \
version INTEGER NOT NULL, state TEXT NOT NULL, reason TEXT NOT NULL, \
evaluated_seq INTEGER NOT NULL, PRIMARY KEY (action_key, version))",
"CREATE TABLE operations (id_hex TEXT PRIMARY KEY, id BLOB NOT NULL, \
kind TEXT NOT NULL, state TEXT NOT NULL, updated_seq INTEGER NOT NULL)",
"CREATE TABLE edge_subscribers (edge_id TEXT NOT NULL, subscriber TEXT NOT NULL, \
registered_seq INTEGER NOT NULL, PRIMARY KEY (edge_id, subscriber))",
"CREATE TABLE manifests (key TEXT PRIMARY KEY, algo TEXT NOT NULL, \
domain TEXT NOT NULL, bytes BLOB NOT NULL, kind TEXT NOT NULL, \
entry_count INTEGER NOT NULL)",
"CREATE TABLE worker_sessions (worker TEXT NOT NULL, incarnation BLOB NOT NULL, \
started_seq INTEGER NOT NULL, ended_seq INTEGER, \
PRIMARY KEY (worker, started_seq))",
"CREATE TABLE worker_capabilities (worker TEXT NOT NULL, capability TEXT NOT NULL, \
PRIMARY KEY (worker, capability))",
"CREATE TABLE worker_health_samples (worker TEXT NOT NULL, seq INTEGER NOT NULL, \
healthy INTEGER NOT NULL, detail TEXT NOT NULL, PRIMARY KEY (worker, seq))",
"CREATE TABLE decision_receipts (kind TEXT NOT NULL, subject TEXT NOT NULL, \
seq INTEGER NOT NULL, decision TEXT NOT NULL, reason TEXT NOT NULL, \
PRIMARY KEY (kind, subject, seq))",
"CREATE TABLE provenance_edges (from_key TEXT NOT NULL, to_key TEXT NOT NULL, \
kind TEXT NOT NULL, PRIMARY KEY (from_key, to_key, kind))",
"CREATE TABLE determinism_audits (action_key TEXT NOT NULL, \
attempt_hex TEXT NOT NULL, seq INTEGER NOT NULL, verdict TEXT NOT NULL, \
PRIMARY KEY (action_key, attempt_hex, seq))",
"CREATE TABLE materialization_records (id_hex TEXT PRIMARY KEY, id BLOB NOT NULL, \
root_key TEXT NOT NULL, dest_path TEXT NOT NULL, state TEXT NOT NULL, \
updated_seq INTEGER NOT NULL)",
],
},
Migration {
version: 9,
statements: &[
"ALTER TABLE action_serving_states ADD COLUMN state_revision INTEGER NOT NULL \
DEFAULT 0",
"ALTER TABLE action_serving_states ADD COLUMN authority_key TEXT NOT NULL \
DEFAULT ''",
"ALTER TABLE action_serving_states ADD COLUMN evaluated_at_micros INTEGER NOT NULL \
DEFAULT 0",
"ALTER TABLE action_serving_states ADD COLUMN max_age_micros INTEGER",
"ALTER TABLE action_serving_states ADD COLUMN clock_uncertainty_micros INTEGER \
NOT NULL DEFAULT 0",
"ALTER TABLE action_serving_states ADD COLUMN clock_epoch INTEGER NOT NULL \
DEFAULT 0",
"CREATE TABLE serving_blocking_quarantines (action_key TEXT NOT NULL, \
scope TEXT NOT NULL, subject TEXT NOT NULL, PRIMARY KEY (action_key, scope, subject))",
],
},
Migration {
version: 10,
statements: &[
"CREATE TABLE divergence_incidents (action_key TEXT NOT NULL, \
seq INTEGER NOT NULL, class TEXT NOT NULL, \
committed_manifest_key TEXT NOT NULL, candidate_manifest_key TEXT NOT NULL, \
candidate_evidence_key TEXT NOT NULL, candidate_pin_hex TEXT NOT NULL, \
generation_hex TEXT NOT NULL, attempt_hex TEXT NOT NULL, detail TEXT NOT NULL, \
PRIMARY KEY (action_key, seq))",
],
},
Migration {
version: 11,
statements: &[
"ALTER TABLE action_evidence_index ADD COLUMN manifest_key TEXT NOT NULL DEFAULT ''",
],
},
Migration {
version: 12,
statements: &["ALTER TABLE object_locations ADD COLUMN durable INTEGER NOT NULL DEFAULT 0"],
},
Migration {
version: 13,
statements: &[
"CREATE TABLE provisional_ancestry (consumer_action_key TEXT NOT NULL, \
producer_action_key TEXT NOT NULL, role TEXT NOT NULL, \
virtual_path BLOB NOT NULL, object_key TEXT NOT NULL, \
adopted INTEGER NOT NULL, \
PRIMARY KEY (consumer_action_key, producer_action_key, role, virtual_path))",
"CREATE TABLE adoption_edges (producer_action_key TEXT NOT NULL, \
role TEXT NOT NULL, virtual_path BLOB NOT NULL, \
from_object_key TEXT NOT NULL, to_object_key TEXT NOT NULL, \
PRIMARY KEY (producer_action_key, role, virtual_path, from_object_key))",
],
},
Migration {
version: 14,
statements: &[
"CREATE TABLE provisional_pins (pin_key TEXT PRIMARY KEY, \
authority_key TEXT NOT NULL, action_key TEXT NOT NULL, \
generation_hex TEXT NOT NULL, attempt_hex TEXT NOT NULL, \
lease_hex TEXT NOT NULL, role INTEGER NOT NULL, \
virtual_path BLOB NOT NULL, obj_algo TEXT NOT NULL, \
obj_domain TEXT NOT NULL, obj_bytes BLOB NOT NULL, \
object_key TEXT NOT NULL, protective_pin_hex TEXT NOT NULL, \
renewal_seq INTEGER NOT NULL DEFAULT 0, adopted_object_key TEXT, \
invalidated_reason TEXT, released INTEGER NOT NULL DEFAULT 0)",
"CREATE TABLE provisional_pin_grants (pin_key TEXT NOT NULL, \
grantee_kind TEXT NOT NULL, grantee_id TEXT NOT NULL, \
granted_seq INTEGER NOT NULL, \
PRIMARY KEY (pin_key, grantee_kind, grantee_id))",
],
},
Migration {
version: 15,
statements: &[
"CREATE TABLE provisional_obligations (consumer_worker TEXT NOT NULL, \
consumer_attempt_hex TEXT NOT NULL, pin_key TEXT NOT NULL, \
producer_action_key TEXT NOT NULL, producer_generation_hex TEXT NOT NULL, \
producer_attempt_hex TEXT NOT NULL, role INTEGER NOT NULL, \
virtual_path BLOB NOT NULL, object_key TEXT NOT NULL, \
status TEXT NOT NULL DEFAULT 'open', resolution_object_key TEXT, \
created_seq INTEGER NOT NULL, \
PRIMARY KEY (consumer_worker, consumer_attempt_hex, pin_key))",
],
},
Migration {
version: 16,
statements: &[
"ALTER TABLE provisional_pins ADD COLUMN toolchain_contract_key \
TEXT NOT NULL DEFAULT ''",
"ALTER TABLE provisional_pins ADD COLUMN event_contract_key \
TEXT NOT NULL DEFAULT ''",
"CREATE TABLE provisional_pin_lineage (descendant_pin_key TEXT NOT NULL, \
ancestor_pin_key TEXT NOT NULL, \
PRIMARY KEY (descendant_pin_key, ancestor_pin_key))",
"CREATE INDEX idx_provisional_pin_lineage_ancestor \
ON provisional_pin_lineage (ancestor_pin_key)",
],
},
Migration {
version: 17,
statements: &["ALTER TABLE provisional_pin_lineage ADD COLUMN \
min_hops INTEGER NOT NULL DEFAULT 1"],
},
Migration {
version: 18,
statements: &[
"CREATE TABLE provisional_install_journal (pin_key TEXT NOT NULL, \
consumer_worker TEXT NOT NULL, consumer_attempt_hex TEXT NOT NULL, \
installed_path BLOB NOT NULL, obj_algo TEXT NOT NULL, \
obj_domain TEXT NOT NULL, obj_bytes BLOB NOT NULL, object_key TEXT NOT NULL, \
installed_seq INTEGER NOT NULL, state TEXT NOT NULL DEFAULT 'installed', \
PRIMARY KEY (pin_key, consumer_attempt_hex, installed_path))",
"CREATE INDEX idx_provisional_install_state \
ON provisional_install_journal (state)",
],
},
Migration {
version: 19,
statements: &[
"CREATE TABLE native_child_bindings (parent_action_key TEXT NOT NULL, \
child_action_key TEXT NOT NULL, bound_seq INTEGER NOT NULL, \
state TEXT NOT NULL DEFAULT 'bound', \
PRIMARY KEY (parent_action_key, child_action_key))",
],
},
Migration {
version: 20,
statements: &[
"ALTER TABLE worker_incarnation_fences ADD COLUMN \
highest_boot_generation BLOB NOT NULL DEFAULT X'0000000000000000'",
"ALTER TABLE worker_incarnation_fences ADD COLUMN \
active INTEGER NOT NULL DEFAULT 1",
"ALTER TABLE worker_incarnation_fences ADD COLUMN \
operator_reenrollment_generation BLOB NOT NULL DEFAULT X'0000000000000000'",
],
},
Migration {
version: 21,
statements: &[
"ALTER TABLE worker_incarnation_fences ADD COLUMN \
clone_ambiguous INTEGER NOT NULL DEFAULT 1",
"ALTER TABLE action_generations ADD COLUMN per_key_ordinal BLOB",
"ALTER TABLE action_attempts ADD COLUMN worker_boot_generation BLOB",
"ALTER TABLE action_attempts ADD COLUMN worker_incarnation BLOB",
"ALTER TABLE action_attempts ADD COLUMN execution_lease_hex TEXT",
],
},
Migration {
version: 22,
statements: &["ALTER TABLE peer_authority_high_water ADD COLUMN incarnation BLOB"],
},
Migration {
version: 23,
statements: &[
"ALTER TABLE peer_authority_high_water ADD COLUMN credential_generation INTEGER",
],
},
Migration {
version: 24,
statements: &["CREATE INDEX idx_pins_released ON pins (released)"],
},
Migration {
version: 25,
statements: &["CREATE TABLE release_verdicts ( \
build TEXT PRIMARY KEY, \
corpus TEXT NOT NULL, \
replayed INTEGER NOT NULL, \
explained INTEGER NOT NULL, \
evaluated_at_micros INTEGER NOT NULL, \
max_age_micros INTEGER, \
clock_uncertainty_micros INTEGER NOT NULL, \
clock_epoch INTEGER NOT NULL)"],
},
];
impl std::fmt::Display for StoreError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "store error: {self:?}")
}
}
impl std::error::Error for StoreError {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StoreError {
Backend(String),
AuthorityHeld {
holder: String,
},
NotActiveAuthority,
GenerationIdNotAboveHighWater,
GenerationIdentityExhausted,
GenerationOrdinalExhausted,
DuplicateAttempt,
DuplicateLease,
UnknownGeneration,
GenerationTombstoned,
UnknownLease,
LeaseAttemptMismatch,
AttemptAuthorityMismatch,
LegacyUnboundAuthority,
UnknownWorkerFence,
WorkerLeaseRejected(WorkerLeaseBindingRejection),
LeaseRenewalMismatch,
LeaseReleased,
LeaseExpired,
NonMonotonicRenewal,
UnknownPin,
PinOwnerMismatch,
PinReleased,
NonMonotonicPinRenewal,
StaleOperatorReset,
StalePeerAuthority,
StaleEdgeIncarnation,
EdgeHandoffActive,
AdoptionEdgeConflict,
EdgeHandoffPredecessorMismatch,
UnknownEdgeHandoff,
NonMonotonicTrustEvaluation,
DuplicateOperation,
UnknownOperation,
ManifestDivergence,
AppendConflict(String),
UnknownMaterialization,
StaleServingRevision,
QuarantineRequiresRepair,
UnknownQuarantineReference,
DomainNotInterned(String),
Corruption(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CommitOutcome {
Committed,
IdempotentDuplicate,
ConflictQuarantined,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AuthorityRow {
pub digest: TypedDigest,
pub cluster_id: String,
pub incarnation: u128,
pub term: u64,
pub acquired_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ActionEntryRow {
pub action_key: TypedDigest,
pub key_epoch: u32,
pub projection_epoch: u32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PublicationRow {
pub action_key: TypedDigest,
pub descriptor_digest: TypedDigest,
pub manifest_digest: TypedDigest,
pub evidence_digest: TypedDigest,
pub winner_generation: u128,
pub winner_attempt: u128,
pub result_kind: ResultKindTag,
pub pin_id: u128,
pub pin_owner: String,
pub provisional_ancestors: Vec<ProvisionalAncestorRow>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalAncestorRow {
pub producer_action_key: String,
pub role: String,
pub virtual_path: Vec<u8>,
pub object_key: String,
pub adopted: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalPinRecord {
pub pin_key: String,
pub authority_key: String,
pub action_key: String,
pub generation_hex: String,
pub attempt_hex: String,
pub lease_hex: String,
pub role_tag: i64,
pub virtual_path: Vec<u8>,
pub object: TypedDigest,
pub object_key: String,
pub protective_pin_hex: String,
pub renewal_seq: u64,
pub adopted_object_key: Option<String>,
pub invalidated_reason: Option<String>,
pub released: bool,
pub toolchain_contract_key: String,
pub event_contract_key: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalPinInsert {
pub pin_key: String,
pub authority_key: String,
pub action_key: String,
pub generation: u128,
pub attempt: u128,
pub lease: u128,
pub role_tag: i64,
pub virtual_path: Vec<u8>,
pub object: TypedDigest,
pub protective_pin_id: u128,
pub reason: String,
pub toolchain_contract_key: String,
pub event_contract_key: String,
pub ancestor_pin_keys: Vec<(String, u64)>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalObligationRow {
pub consumer_worker: String,
pub consumer_attempt_hex: String,
pub pin_key: String,
pub producer_action_key: String,
pub producer_generation_hex: String,
pub producer_attempt_hex: String,
pub role_tag: i64,
pub virtual_path: Vec<u8>,
pub object_key: String,
pub status: String,
pub resolution_object_key: Option<String>,
pub created_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalObligationInsert {
pub consumer_worker: String,
pub consumer_attempt: u128,
pub pin_key: String,
pub producer_action_key: String,
pub producer_generation: u128,
pub producer_attempt: u128,
pub role_tag: i64,
pub virtual_path: Vec<u8>,
pub object_key: String,
pub created_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalInstallRecord {
pub pin_key: String,
pub consumer_worker: String,
pub consumer_attempt_hex: String,
pub installed_path: Vec<u8>,
pub object: TypedDigest,
pub object_key: String,
pub installed_seq: u64,
pub state: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalInstallInsert {
pub pin_key: String,
pub consumer_worker: String,
pub consumer_attempt: u128,
pub installed_path: Vec<u8>,
pub object: TypedDigest,
pub installed_seq: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ResultKindTag {
Success,
DeterministicFailure,
}
impl ResultKindTag {
const fn as_str(self) -> &'static str {
match self {
Self::Success => "success",
Self::DeterministicFailure => "deterministic-failure",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QuarantineScope {
Location,
LogicalObject,
ActionEntry,
}
impl QuarantineScope {
const fn as_str(self) -> &'static str {
match self {
Self::Location => "location",
Self::LogicalObject => "logical-object",
Self::ActionEntry => "action-entry",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct GenerationState {
pub tombstoned: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LeaseState {
pub released: bool,
pub renewal_seq: u64,
pub expires_at_own_monotonic_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PinRow {
pub root_key: String,
pub owner: String,
pub class: String,
pub expires_at_seq: Option<u64>,
pub released: bool,
pub renewal_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GcTombstoneRow {
pub object_key: String,
pub store_path: String,
pub marked_seq: u64,
pub grace_until_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GcReceiptRow {
pub seq: u64,
pub mode: String,
pub planned: u64,
pub reclaimed: u64,
pub skipped: u64,
pub truncated: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GcSnapshot {
pub pinned_roots: Vec<String>,
pub located_objects: Vec<String>,
pub reachable_from_pins: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReconciliationRow {
pub object_key: String,
pub store_path: String,
pub verified_seq: Option<u64>,
pub encoding: String,
pub quarantined: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TrustEvaluationRow {
pub version: u32,
pub state: String,
pub reason: String,
pub evaluated_seq: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct EdgeHandoffRow {
pub active_incarnation: u128,
pub predecessor_incarnation: u128,
pub begun_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ServingRecordRow {
pub disposition: String,
pub state_revision: u64,
pub authority_key: String,
pub validity: ServingValidity,
pub blocking: Vec<(String, String)>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DivergenceIncidentRow {
pub action_key: String,
pub seq: u64,
pub class: String,
pub committed_manifest_key: String,
pub candidate_manifest_key: String,
pub candidate_evidence_key: String,
pub candidate_pin_hex: String,
pub generation_hex: String,
pub attempt_hex: String,
pub detail: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct VerificationSampleRow {
pub attempt_hex: String,
pub passed: bool,
pub seq: u64,
}
pub trait RabsMetadataStore {
fn schema_version(&mut self) -> Result<u32, StoreError>;
fn query(&mut self, sql: &str, params: &[SqlValue]) -> Result<Vec<Vec<SqlValue>>, StoreError>;
fn intern_domain(&mut self, domain: &'static str);
fn acquire_authority(&mut self, row: &AuthorityRow) -> Result<(), StoreError>;
fn release_authority(&mut self, digest: &TypedDigest) -> Result<(), StoreError>;
fn active_authority(&mut self) -> Result<Option<AuthorityRow>, StoreError>;
fn upsert_action_entry(&mut self, row: &ActionEntryRow) -> Result<(), StoreError>;
fn lookup_action(&mut self, key: &TypedDigest) -> Result<Option<ActionEntryRow>, StoreError>;
fn create_generation(
&mut self,
authority: &TypedDigest,
id: u128,
action_key: &TypedDigest,
) -> Result<(), StoreError>;
fn create_bound_generation(
&mut self,
authority: &TypedDigest,
generation: &ActionGeneration,
action_key: &TypedDigest,
) -> Result<(), StoreError>;
fn allocate_bound_generation(
&mut self,
authority: &TypedDigest,
action_key: &TypedDigest,
) -> Result<ActionGeneration, StoreError>;
fn tombstone_generation(&mut self, id: u128) -> Result<(), StoreError>;
fn close_generations_for_other_authorities(
&mut self,
active: &TypedDigest,
) -> Result<u64, StoreError>;
fn record_attempt(
&mut self,
id: u128,
generation: u128,
worker: &str,
seq: u64,
) -> Result<(), StoreError>;
fn admit_attempt_lease(
&mut self,
authority: &AttemptAuthority,
recorded_seq: u64,
expires_at_own_monotonic_ms: u64,
) -> Result<(), StoreError>;
fn renew_attempt_lease(
&mut self,
authority: &AttemptAuthority,
renewal: LeaseRenewal,
expires_at_own_monotonic_ms: u64,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError>;
fn release_lease(&mut self, id: u128) -> Result<(), StoreError>;
fn commit_publication(
&mut self,
authority: &TypedDigest,
attempt_authority: Option<(&AttemptAuthority, &dyn Fn() -> u64)>,
row: &PublicationRow,
) -> Result<CommitOutcome, StoreError>;
fn append_evidence(
&mut self,
action: &TypedDigest,
manifest_key: &str,
evidence: &TypedDigest,
generation: u128,
attempt: u128,
) -> Result<(), StoreError>;
fn append_evidence_for_attempt(
&mut self,
authority: &AttemptAuthority,
manifest_key: &str,
evidence: &TypedDigest,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError>;
fn has_publication(&mut self, action: &TypedDigest) -> Result<bool, StoreError>;
fn published_manifest_key(
&mut self,
action: &TypedDigest,
) -> Result<Option<String>, StoreError>;
fn published_manifest_key_str(
&mut self,
action_key: &str,
) -> Result<Option<String>, StoreError>;
fn generation_state(&mut self, id: u128) -> Result<Option<GenerationState>, StoreError>;
fn attempt_exists(&mut self, id: u128, generation: u128) -> Result<bool, StoreError>;
fn lease_state(&mut self, id: u128) -> Result<Option<LeaseState>, StoreError>;
fn validate_attempt_lease(
&mut self,
authority: &AttemptAuthority,
own_monotonic_now_ms: u64,
) -> Result<LeaseState, StoreError>;
fn object_located(&mut self, object: &TypedDigest) -> Result<bool, StoreError>;
fn object_locations(
&mut self,
object: &TypedDigest,
) -> Result<Vec<(String, String, bool)>, StoreError>;
fn record_object(&mut self, id: &TypedDigest, logical_size: u64) -> Result<(), StoreError>;
fn add_location(
&mut self,
object: &TypedDigest,
store_path: &str,
verified_seq: Option<u64>,
encoding: &str,
durable: bool,
) -> Result<(), StoreError>;
fn object_durably_located(&mut self, object: &TypedDigest) -> Result<bool, StoreError>;
fn set_location_quarantined(
&mut self,
object: &TypedDigest,
store_path: &str,
quarantined: bool,
) -> Result<(), StoreError>;
fn add_object_edge(
&mut self,
parent: &TypedDigest,
child: &TypedDigest,
kind: &str,
) -> Result<(), StoreError>;
#[allow(clippy::too_many_arguments)]
fn create_pin(
&mut self,
id: u128,
root: &TypedDigest,
owner: &str,
class: &str,
expires_at_seq: Option<u64>,
evidence: Option<&str>,
durable: bool,
reason: &str,
) -> Result<(), StoreError>;
fn renew_pin(&mut self, id: u128, renewal_seq: u64) -> Result<(), StoreError>;
fn pin_row(&mut self, id: u128) -> Result<Option<PinRow>, StoreError>;
fn release_pin(&mut self, id: u128, owner: &str) -> Result<(), StoreError>;
fn put_recipe(&mut self, action: &TypedDigest, recipe: &TypedDigest) -> Result<(), StoreError>;
fn put_key_breakdown(
&mut self,
action: &TypedDigest,
component: &str,
digest: &TypedDigest,
) -> Result<(), StoreError>;
fn set_trust(
&mut self,
action: &TypedDigest,
state: &str,
reason: &str,
) -> Result<(), StoreError>;
fn add_quarantine(
&mut self,
scope: QuarantineScope,
subject: &str,
reason: &str,
) -> Result<(), StoreError>;
fn record_verification_sample(
&mut self,
action: &TypedDigest,
attempt: u128,
passed: bool,
seq: u64,
) -> Result<(), StoreError>;
fn list_evidence_keys(&mut self, action: &TypedDigest) -> Result<Vec<String>, StoreError>;
fn list_evidence_keys_for_manifest(
&mut self,
manifest_key: &str,
) -> Result<Vec<String>, StoreError>;
fn list_verification_samples(
&mut self,
action: &TypedDigest,
) -> Result<Vec<VerificationSampleRow>, StoreError>;
fn attempt_worker_by_hex(&mut self, attempt_hex: &str) -> Result<Option<String>, StoreError>;
fn gc_snapshot(&mut self, seq: u64) -> Result<GcSnapshot, StoreError>;
fn reconciliation_scan(&mut self) -> Result<Vec<ReconciliationRow>, StoreError>;
fn remove_location_by_key(
&mut self,
object_key: &str,
store_path: &str,
) -> Result<bool, StoreError>;
fn record_gc_receipt(&mut self, receipt: &GcReceiptRow) -> Result<(), StoreError>;
fn add_gc_tombstone(
&mut self,
object_key: &str,
store_path: &str,
marked_seq: u64,
grace_until_seq: u64,
) -> Result<(), StoreError>;
fn due_gc_tombstones(&mut self, now_seq: u64) -> Result<Vec<GcTombstoneRow>, StoreError>;
fn remove_gc_tombstone(
&mut self,
object_key: &str,
store_path: &str,
) -> Result<bool, StoreError>;
fn list_publications(&mut self) -> Result<Vec<(String, String)>, StoreError>;
fn publications_missing_their_generation(&mut self) -> Result<Vec<String>, StoreError>;
fn pin_released_by_hex(&mut self, pin_hex: &str) -> Result<Option<bool>, StoreError>;
fn has_serving_state_key(&mut self, action_key: &str) -> Result<bool, StoreError>;
fn has_evidence_key(&mut self, action_key: &str) -> Result<bool, StoreError>;
fn authority_count(&mut self) -> Result<u64, StoreError>;
fn generation_count(&mut self) -> Result<u64, StoreError>;
fn has_generation_high_water(&mut self) -> Result<bool, StoreError>;
fn record_eviction_tombstone(
&mut self,
action: &TypedDigest,
semantic: &TypedDigest,
observable: &TypedDigest,
evicted_seq: u64,
) -> Result<(), StoreError>;
fn eviction_tombstone(
&mut self,
action: &TypedDigest,
) -> Result<Option<(TypedDigest, TypedDigest)>, StoreError>;
fn consume_eviction_tombstone(&mut self, action: &TypedDigest) -> Result<bool, StoreError>;
fn record_operator_reset(&mut self, generation: u64, seq: u64) -> Result<(), StoreError>;
fn apply_operator_reset_to_peer(
&mut self,
authority: &TypedDigest,
peer_id: &str,
reset_generation: u64,
credential_generation: u64,
term: u64,
incarnation: u128,
) -> Result<(), StoreError>;
fn highest_operator_reset(&mut self) -> Result<Option<u64>, StoreError>;
fn serving_disposition_key(&mut self, action_key: &str) -> Result<Option<String>, StoreError>;
fn set_serving_disposition_key(
&mut self,
action_key: &str,
disposition: &str,
) -> Result<(), StoreError>;
fn set_serving_disposition_for_attempt(
&mut self,
authority: &AttemptAuthority,
disposition: &str,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError>;
fn record_peer_authority_high_water(
&mut self,
authority: &TypedDigest,
peer_id: &str,
credential_generation: u64,
term: u64,
incarnation: u128,
observed_seq: u64,
) -> Result<(), StoreError>;
fn peer_authority_high_water(
&mut self,
peer_id: &str,
) -> Result<Option<(u64, u64, u64)>, StoreError>;
fn admit_worker_session(
&mut self,
authority: &TypedDigest,
offer: &WorkerSessionOffer,
started_seq: u64,
) -> Result<WorkerAdmission, StoreError>;
fn release_worker_session(
&mut self,
authority: &TypedDigest,
worker: &PeerId,
incarnation: WorkerIncarnationId,
started_seq: u64,
ended_seq: u64,
) -> Result<bool, StoreError>;
fn worker_incarnation_fence(
&mut self,
worker: &PeerId,
) -> Result<Option<WorkerIncarnationFenceRecord>, StoreError>;
fn advance_edge_fence(
&mut self,
authority: &TypedDigest,
edge_id: &str,
incarnation: u128,
) -> Result<(), StoreError>;
fn edge_fence(&mut self, edge_id: &str) -> Result<Option<u128>, StoreError>;
fn begin_edge_handoff(
&mut self,
authority: &TypedDigest,
edge_id: &str,
active_incarnation: u128,
predecessor_incarnation: u128,
begun_seq: u64,
) -> Result<(), StoreError>;
fn resolve_edge_handoff(
&mut self,
authority: &TypedDigest,
edge_id: &str,
active_incarnation: u128,
) -> Result<(), StoreError>;
fn active_edge_handoff(&mut self, edge_id: &str) -> Result<Option<EdgeHandoffRow>, StoreError>;
fn append_trust_evaluation(
&mut self,
authority: &TypedDigest,
action: &TypedDigest,
row: &TrustEvaluationRow,
) -> Result<(), StoreError>;
fn latest_trust_evaluation(
&mut self,
action: &TypedDigest,
) -> Result<Option<TrustEvaluationRow>, StoreError>;
fn create_operation(
&mut self,
authority: &TypedDigest,
id: u128,
kind: &str,
state: &str,
seq: u64,
) -> Result<(), StoreError>;
fn update_operation_state(&mut self, id: u128, state: &str, seq: u64)
-> Result<(), StoreError>;
fn operation_state(&mut self, id: u128) -> Result<Option<String>, StoreError>;
fn register_edge_subscriber(
&mut self,
edge_id: &str,
subscriber: &str,
registered_seq: u64,
) -> Result<(), StoreError>;
fn remove_edge_subscriber(
&mut self,
edge_id: &str,
subscriber: &str,
) -> Result<bool, StoreError>;
fn list_edge_subscribers(&mut self, edge_id: &str) -> Result<Vec<String>, StoreError>;
fn record_manifest(
&mut self,
manifest: &TypedDigest,
kind: &str,
entry_count: u64,
) -> Result<(), StoreError>;
fn manifest_meta(
&mut self,
manifest: &TypedDigest,
) -> Result<Option<(String, u64)>, StoreError>;
fn record_worker_session(
&mut self,
worker: &str,
incarnation: u128,
started_seq: u64,
) -> Result<(), StoreError>;
fn end_worker_session(
&mut self,
worker: &str,
started_seq: u64,
ended_seq: u64,
) -> Result<bool, StoreError>;
fn record_worker_capability(
&mut self,
worker: &str,
capability: &str,
) -> Result<(), StoreError>;
fn list_worker_capabilities(&mut self, worker: &str) -> Result<Vec<String>, StoreError>;
fn record_worker_health_sample(
&mut self,
worker: &str,
seq: u64,
healthy: bool,
detail: &str,
) -> Result<(), StoreError>;
fn record_decision_receipt(
&mut self,
kind: &str,
subject: &str,
seq: u64,
decision: &str,
reason: &str,
) -> Result<(), StoreError>;
fn add_provenance_edge(
&mut self,
from: &TypedDigest,
to: &TypedDigest,
kind: &str,
) -> Result<(), StoreError>;
fn record_determinism_audit(
&mut self,
action: &TypedDigest,
attempt: u128,
seq: u64,
verdict: &str,
) -> Result<(), StoreError>;
fn create_materialization(
&mut self,
id: u128,
root: &TypedDigest,
dest_path: &str,
state: &str,
seq: u64,
) -> Result<(), StoreError>;
fn update_materialization_state(
&mut self,
id: u128,
state: &str,
seq: u64,
) -> Result<(), StoreError>;
fn materialization_state(&mut self, id: u128) -> Result<Option<String>, StoreError>;
fn put_serving_record(
&mut self,
authority: &TypedDigest,
action_key: &str,
disposition: &str,
state_revision: u64,
validity: &ServingValidity,
blocking: &[(QuarantineScope, String)],
) -> Result<(), StoreError>;
fn serving_record(&mut self, action_key: &str) -> Result<Option<ServingRecordRow>, StoreError>;
fn record_release_verdict(&mut self, verdict: &ReleaseVerdict) -> Result<(), StoreError>;
fn release_verdict(&mut self, build: &str) -> Result<Option<ReleaseVerdict>, StoreError>;
fn record_divergence_incident(
&mut self,
authority: &TypedDigest,
row: &DivergenceIncidentRow,
) -> Result<(), StoreError>;
fn list_divergence_incidents(
&mut self,
action_key: &str,
) -> Result<Vec<DivergenceIncidentRow>, StoreError>;
fn record_adoption_edge(
&mut self,
authority: &TypedDigest,
producer_action_key: &str,
role: &str,
virtual_path: &[u8],
from_object_key: &str,
to_object_key: &str,
) -> Result<(), StoreError>;
fn has_adoption_edge(
&mut self,
producer_action_key: &str,
role: &str,
virtual_path: &[u8],
from_object_key: &str,
to_object_key: &str,
) -> Result<bool, StoreError>;
fn record_provisional_consumption(
&mut self,
consumption: &ProvisionalObligationInsert,
) -> Result<(), StoreError>;
fn list_open_provisional_obligations(
&mut self,
consumer_worker: &str,
consumer_attempt_hex: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError>;
fn resolve_provisional_obligations(
&mut self,
pin_key: &str,
resolution_object_key: &str,
) -> Result<usize, StoreError>;
fn cancel_provisional_obligations(&mut self, pin_key: &str) -> Result<usize, StoreError>;
fn count_open_provisional_obligations(&mut self, pin_key: &str) -> Result<usize, StoreError>;
fn insert_provisional_install(
&mut self,
install: &ProvisionalInstallInsert,
) -> Result<(), StoreError>;
fn list_provisional_installs_for_pins(
&mut self,
pin_keys: &[String],
) -> Result<Vec<ProvisionalInstallRecord>, StoreError>;
fn list_provisional_installs_by_state(
&mut self,
state: &str,
) -> Result<Vec<ProvisionalInstallRecord>, StoreError>;
fn set_provisional_install_state(
&mut self,
pin_key: &str,
consumer_attempt_hex: &str,
installed_path: &[u8],
state: &str,
) -> Result<(), StoreError>;
fn bind_native_children(
&mut self,
parent_action_key: &str,
child_action_keys: &[String],
bound_seq: u64,
) -> Result<(), StoreError>;
fn list_native_child_bindings(
&mut self,
parent_action_key: &str,
) -> Result<Vec<(String, String)>, StoreError>;
fn set_native_child_binding_state(
&mut self,
parent_action_key: &str,
child_action_key: &str,
state: &str,
) -> Result<(), StoreError>;
fn list_open_provisional_pins_for_action_generation(
&mut self,
action_key: &str,
generation_hex: &str,
) -> Result<Vec<ProvisionalPinRecord>, StoreError>;
fn list_open_provisional_pins_for_action(
&mut self,
action_key: &str,
) -> Result<Vec<ProvisionalPinRecord>, StoreError>;
fn list_open_provisional_pins_for_authority(
&mut self,
authority_key: &str,
) -> Result<Vec<ProvisionalPinRecord>, StoreError>;
fn list_provisional_obligations_for_pin(
&mut self,
pin_key: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError>;
fn list_provisional_pin_ancestors(
&mut self,
descendant_pin_key: &str,
) -> Result<Vec<(String, u64)>, StoreError>;
fn provisional_pin_closure_depth(
&mut self,
descendant_pin_key: &str,
) -> Result<u64, StoreError>;
fn list_provisional_pin_descendants(
&mut self,
ancestor_pin_key: &str,
) -> Result<Vec<String>, StoreError>;
fn list_open_provisional_obligations_by_attempt(
&mut self,
consumer_attempt_hex: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError>;
fn list_provisional_obligations_by_attempt_all(
&mut self,
consumer_attempt_hex: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError>;
fn list_provisional_ancestors(
&mut self,
consumer_action_key: &str,
) -> Result<Vec<ProvisionalAncestorRow>, StoreError>;
fn provisional_pin_row(
&mut self,
pin_key: &str,
) -> Result<Option<ProvisionalPinRecord>, StoreError>;
fn insert_provisional_pin(&mut self, pin: &ProvisionalPinInsert) -> Result<(), StoreError>;
fn record_provisional_grant(
&mut self,
pin_key: &str,
grantee_kind: &str,
grantee_id: &str,
granted_seq: u64,
) -> Result<(), StoreError>;
fn list_provisional_grants(
&mut self,
pin_key: &str,
) -> Result<Vec<(String, String, u64)>, StoreError>;
fn renew_provisional_pin(&mut self, pin_key: &str, renewal_seq: u64) -> Result<(), StoreError>;
fn close_provisional_pin(
&mut self,
pin_key: &str,
invalidation_reason: Option<&str>,
) -> Result<(), StoreError>;
fn adopt_provisional_pin(
&mut self,
pin_key: &str,
committed_object_key: &str,
) -> Result<(), StoreError>;
fn record_served_consumer(
&mut self,
action_key: &str,
consumer: &str,
) -> Result<(), StoreError>;
fn list_served_consumers(&mut self, action_key: &str) -> Result<Vec<String>, StoreError>;
fn differential_snapshot(&mut self) -> Result<Vec<String>, StoreError>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SqlValue {
Null,
Int(i64),
Text(String),
Blob(Vec<u8>),
}
pub trait SqlEngine {
fn execute(&mut self, sql: &str, params: &[SqlValue]) -> Result<usize, StoreError>;
fn query(&mut self, sql: &str, params: &[SqlValue]) -> Result<Vec<Vec<SqlValue>>, StoreError>;
}
pub struct SqlMetadataStore<E: SqlEngine> {
engine: E,
domains: HashMap<String, &'static str>,
}
fn hex(bytes: &[u8]) -> String {
let mut out = String::with_capacity(bytes.len() * 2);
for b in bytes {
use std::fmt::Write;
let _ = write!(out, "{b:02x}");
}
out
}
#[must_use]
pub fn digest_key(d: &TypedDigest) -> String {
format!("{}:{}", d.domain, hex(&d.bytes))
}
const fn algo_tag(a: DigestAlgorithm) -> &'static str {
match a {
DigestAlgorithm::Sha256V1 => "sha256-v1",
}
}
fn u128_blob(v: u128) -> Vec<u8> {
v.to_be_bytes().to_vec()
}
fn u64_blob(v: u64) -> Vec<u8> {
v.to_be_bytes().to_vec()
}
fn u128_hex(v: u128) -> String {
hex(&v.to_be_bytes())
}
impl<E: SqlEngine> SqlMetadataStore<E> {
pub fn open(engine: E) -> Result<Self, StoreError> {
let mut store = Self {
engine,
domains: HashMap::new(),
};
store.apply_migrations()?;
Ok(store)
}
pub fn operation_update_high_water(&mut self) -> Result<u64, StoreError> {
let rows = self
.engine
.query("SELECT COALESCE(MAX(updated_seq), 1) FROM operations", &[])?;
let [row] = rows.as_slice() else {
return Err(StoreError::Corruption(
"operation update high-water row count".into(),
));
};
let [value] = row.as_slice() else {
return Err(StoreError::Corruption(
"operation update high-water shape".into(),
));
};
Ok(expect_u64(value, "operation update high-water")?.max(1))
}
pub(crate) fn engine_mut(&mut self) -> &mut E {
&mut self.engine
}
fn apply_migrations(&mut self) -> Result<(), StoreError> {
let applied: u32 = {
let has_epochs = !self
.engine
.query(
"SELECT name FROM sqlite_master WHERE type = 'table' \
AND name = 'schema_epochs'",
&[],
)?
.is_empty();
if has_epochs {
let rows = self
.engine
.query("SELECT MAX(version) FROM schema_epochs", &[])?;
match rows.first().and_then(|r| r.first()) {
Some(SqlValue::Int(v)) => u32::try_from(*v)
.map_err(|_| StoreError::Corruption("negative schema version".into()))?,
_ => 0,
}
} else {
0
}
};
let Some(first_pending) = MIGRATIONS
.iter()
.position(|migration| migration.version > applied)
else {
return Ok(());
};
self.engine.execute("BEGIN", &[])?;
let result = (|| -> Result<(), StoreError> {
for migration in &MIGRATIONS[first_pending..] {
for statement in migration.statements {
self.engine.execute(statement, &[])?;
}
self.engine.execute(
"INSERT INTO schema_epochs (version, applied_seq) VALUES (?1, ?2)",
&[
SqlValue::Int(i64::from(migration.version)),
SqlValue::Int(0),
],
)?;
}
self.engine.execute("COMMIT", &[])?;
Ok(())
})();
if result.is_err() {
let _ = self.engine.execute("ROLLBACK", &[]);
}
result
}
fn intern(&mut self, domain: &'static str) {
self.domains.entry(domain.to_owned()).or_insert(domain);
}
fn restore_domain(&self, domain: &str) -> Result<&'static str, StoreError> {
self.domains
.get(domain)
.copied()
.ok_or_else(|| StoreError::DomainNotInterned(domain.to_owned()))
}
fn restore_digest(
&self,
algo: &SqlValue,
domain: &SqlValue,
bytes: &SqlValue,
) -> Result<TypedDigest, StoreError> {
let (SqlValue::Text(algo), SqlValue::Text(domain), SqlValue::Blob(bytes)) =
(algo, domain, bytes)
else {
return Err(StoreError::Corruption("digest column shape".into()));
};
if algo != "sha256-v1" {
return Err(StoreError::Corruption(format!("unknown algorithm {algo}")));
}
let bytes: [u8; 32] = bytes
.as_slice()
.try_into()
.map_err(|_| StoreError::Corruption("digest not 32 bytes".into()))?;
Ok(TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain: self.restore_domain(domain)?,
bytes,
})
}
fn digest_params(d: &TypedDigest) -> [SqlValue; 3] {
[
SqlValue::Text(algo_tag(d.algorithm).to_owned()),
SqlValue::Text(d.domain.to_owned()),
SqlValue::Blob(d.bytes.to_vec()),
]
}
fn in_txn<T>(
&mut self,
body: impl FnOnce(&mut E) -> Result<T, StoreError>,
) -> Result<T, StoreError> {
self.engine.execute("BEGIN", &[])?;
match body(&mut self.engine) {
Ok(v) => {
self.engine.execute("COMMIT", &[])?;
Ok(v)
}
Err(e) => {
let _ = self.engine.execute("ROLLBACK", &[]);
Err(e)
}
}
}
fn active_authority_key(engine: &mut E) -> Result<Option<String>, StoreError> {
let rows = engine.query(
"SELECT key FROM coordinator_authorities WHERE released = 0",
&[],
)?;
match rows.as_slice() {
[] => Ok(None),
[row] => match row.first() {
Some(SqlValue::Text(k)) => Ok(Some(k.clone())),
_ => Err(StoreError::Corruption("authority key shape".into())),
},
_ => Err(StoreError::Corruption(
"more than one active authority".into(),
)),
}
}
fn require_active(engine: &mut E, authority: &TypedDigest) -> Result<(), StoreError> {
match Self::active_authority_key(engine)? {
Some(active) if active == digest_key(authority) => Ok(()),
_ => Err(StoreError::NotActiveAuthority),
}
}
fn attempt_authority_digest(authority: &AttemptAuthority) -> TypedDigest {
rabs_key::authority_binding::coordinator_authority_digest(&authority.coordinator)
}
fn require_attempt_context(
engine: &mut E,
authority: &AttemptAuthority,
) -> Result<TypedDigest, StoreError> {
let coordinator = Self::attempt_authority_digest(authority);
Self::require_active(engine, &coordinator)?;
if authority.action_generation.created_under_authority_digest != coordinator {
return Err(StoreError::AttemptAuthorityMismatch);
}
let rows = engine.query(
"SELECT action_key, authority_key, tombstoned, per_key_ordinal \
FROM action_generations WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(
authority.action_generation.generation_id.0,
))],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownGeneration);
};
if rows.len() != 1 {
return Err(StoreError::Corruption(
"duplicate action generation rows".into(),
));
}
let [action_key, authority_key, tombstoned, ordinal] = row.as_slice() else {
return Err(StoreError::Corruption(
"action generation binding shape".into(),
));
};
if expect_u64(tombstoned, "generation tombstoned")? != 0 {
return Err(StoreError::GenerationTombstoned);
}
let stored_ordinal = match ordinal {
SqlValue::Null => return Err(StoreError::LegacyUnboundAuthority),
value => expect_u64_blob(value, "generation ordinal")?,
};
if expect_text(action_key, "generation action key")? != digest_key(&authority.action_key)
|| expect_text(authority_key, "generation authority")? != digest_key(&coordinator)
|| stored_ordinal != authority.action_generation.per_key_ordinal
{
return Err(StoreError::AttemptAuthorityMismatch);
}
Ok(coordinator)
}
fn require_worker_lease_binding(
engine: &mut E,
worker: &PeerId,
boot_generation: WorkerBootGeneration,
incarnation: WorkerIncarnationId,
) -> Result<(), StoreError> {
let rows = engine.query(
"SELECT highest_boot_generation, incarnation, active, \
operator_reenrollment_generation, clone_ambiguous \
FROM worker_incarnation_fences WHERE worker = ?1",
&[SqlValue::Text(worker.0.clone())],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownWorkerFence);
};
if rows.len() != 1 {
return Err(StoreError::Corruption("duplicate worker fence rows".into()));
}
let fence = decode_worker_incarnation_fence(&worker.0, row)?;
fence
.validate_lease_binding(worker, boot_generation, incarnation)
.map_err(StoreError::WorkerLeaseRejected)
}
fn revoke_worker_leases(engine: &mut E, worker: &str) -> Result<(), StoreError> {
engine.execute(
"UPDATE execution_leases SET released = 1 \
WHERE released = 0 AND attempt_hex IN \
(SELECT id_hex FROM action_attempts WHERE worker = ?1)",
&[SqlValue::Text(worker.to_owned())],
)?;
Ok(())
}
fn require_live_attempt(
engine: &mut E,
authority: &AttemptAuthority,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError> {
let state = Self::bound_lease_state(engine, authority)?;
if state.released {
return Err(StoreError::LeaseReleased);
}
if own_monotonic_now_ms() >= state.expires_at_own_monotonic_ms {
return Err(StoreError::LeaseExpired);
}
if state.renewal_seq != authority.lease_renewal_seq.0 {
return Err(StoreError::LeaseRenewalMismatch);
}
Ok(())
}
fn append_evidence_row(
engine: &mut E,
action: &TypedDigest,
manifest_key: &str,
evidence: &TypedDigest,
generation: u128,
attempt: u128,
) -> Result<(), StoreError> {
let [algo, domain, bytes] = Self::digest_params(evidence);
engine.execute(
"INSERT OR IGNORE INTO action_evidence_index \
(action_key, evidence_algo, evidence_domain, evidence_bytes, \
generation_hex, attempt_hex, manifest_key) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
&[
SqlValue::Text(digest_key(action)),
algo,
domain,
bytes,
SqlValue::Text(u128_hex(generation)),
SqlValue::Text(u128_hex(attempt)),
SqlValue::Text(manifest_key.to_owned()),
],
)?;
Ok(())
}
fn require_quarantine_preserved(current: &str, proposed: &str) -> Result<(), StoreError> {
if (current == "quarantined" && proposed != "quarantined")
|| (current == "presentation-quarantined"
&& !matches!(proposed, "presentation-quarantined" | "quarantined"))
{
return Err(StoreError::QuarantineRequiresRepair);
}
Ok(())
}
fn set_serving_disposition_row(
engine: &mut E,
action_key: &str,
disposition: &str,
) -> Result<(), StoreError> {
let rows = engine.query(
"SELECT disposition, state_revision FROM action_serving_states WHERE action_key = ?1",
&[SqlValue::Text(action_key.to_owned())],
)?;
if let Some(row) = rows.first() {
let [stored_disposition, stored_revision] = row.as_slice() else {
return Err(StoreError::Corruption(
"serving disposition row shape".into(),
));
};
let stored_disposition = expect_text(stored_disposition, "disposition")?;
Self::require_quarantine_preserved(&stored_disposition, disposition)?;
let revision = expect_u64(stored_revision, "state_revision")?;
let revision = if stored_disposition == disposition {
revision
} else {
revision
.checked_add(1)
.ok_or_else(|| StoreError::Corruption("serving revision exhausted".into()))?
};
let revision = to_seq(revision, "state_revision")?;
engine.execute(
"UPDATE action_serving_states SET disposition = ?2, state_revision = ?3 \
WHERE action_key = ?1",
&[
SqlValue::Text(action_key.to_owned()),
SqlValue::Text(disposition.to_owned()),
SqlValue::Int(revision),
],
)?;
} else {
engine.execute(
"INSERT INTO action_serving_states (action_key, disposition, version, state_revision) \
VALUES (?1, ?2, 1, 0)",
&[
SqlValue::Text(action_key.to_owned()),
SqlValue::Text(disposition.to_owned()),
],
)?;
}
Ok(())
}
fn bound_lease_state(
engine: &mut E,
authority: &AttemptAuthority,
) -> Result<LeaseState, StoreError> {
Self::require_attempt_context(engine, authority)?;
let rows = engine.query(
"SELECT l.attempt_hex, l.released, l.renewal_seq, l.expires_at_seq, \
a.generation_hex, a.worker, a.worker_boot_generation, \
a.worker_incarnation, a.execution_lease_hex \
FROM execution_leases l \
LEFT JOIN action_attempts a ON a.id_hex = l.attempt_hex \
WHERE l.id_hex = ?1",
&[SqlValue::Text(u128_hex(authority.execution_lease_id.0))],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownLease);
};
if rows.len() != 1 {
return Err(StoreError::Corruption(
"duplicate execution lease rows".into(),
));
}
let [
attempt_hex,
released,
renewal_seq,
expires_at_seq,
generation_hex,
worker,
boot,
incarnation,
attempt_lease_hex,
] = row.as_slice()
else {
return Err(StoreError::Corruption("bound lease row shape".into()));
};
if expect_text(attempt_hex, "lease attempt")? != u128_hex(authority.attempt_id.0) {
return Err(StoreError::LeaseAttemptMismatch);
}
let (stored_boot, stored_incarnation, stored_attempt_lease) =
match (boot, incarnation, attempt_lease_hex) {
(SqlValue::Null, _, _) | (_, SqlValue::Null, _) | (_, _, SqlValue::Null) => {
return Err(StoreError::LegacyUnboundAuthority);
}
(boot, incarnation, attempt_lease) => (
expect_u64_blob(boot, "attempt worker boot generation")?,
expect_u128(incarnation, "attempt worker incarnation")?,
expect_text(attempt_lease, "attempt execution lease")?,
),
};
if expect_text(generation_hex, "attempt generation")?
!= u128_hex(authority.action_generation.generation_id.0)
|| expect_text(worker, "attempt worker")? != authority.worker_peer_id.0
|| stored_boot != authority.worker_boot_generation.0
|| stored_incarnation != authority.worker_incarnation_id.0
|| stored_attempt_lease != u128_hex(authority.execution_lease_id.0)
{
return Err(StoreError::AttemptAuthorityMismatch);
}
Self::require_worker_lease_binding(
engine,
&authority.worker_peer_id,
authority.worker_boot_generation,
authority.worker_incarnation_id,
)?;
Ok(LeaseState {
released: expect_u64(released, "released")? != 0,
renewal_seq: expect_u64(renewal_seq, "renewal_seq")?,
expires_at_own_monotonic_ms: expect_u64(expires_at_seq, "expires_at_seq")?,
})
}
fn generation_high_water(engine: &mut E) -> Result<u128, StoreError> {
fn decode(value: &SqlValue) -> Result<u128, StoreError> {
match value {
SqlValue::Blob(b) => {
let bytes: [u8; 16] = b
.as_slice()
.try_into()
.map_err(|_| StoreError::Corruption("high-water not 16 bytes".into()))?;
Ok(u128::from_be_bytes(bytes))
}
SqlValue::Null => Ok(0),
_ => Err(StoreError::Corruption("high-water shape".into())),
}
}
let stored = engine.query(
"SELECT value FROM generation_high_water WHERE kind = 'action-generation'",
&[],
)?;
let stored = match stored.first().and_then(|r| r.first()) {
None => 0,
Some(value) => decode(value)?,
};
let observed = engine.query("SELECT MAX(id) FROM action_generations", &[])?;
let observed = match observed.first().and_then(|r| r.first()) {
None => 0,
Some(value) => decode(value)?,
};
Ok(stored.max(observed))
}
fn map_provisional_pin_row(
&self,
row: &[SqlValue],
) -> Result<ProvisionalPinRecord, StoreError> {
let [
pin,
authority,
action,
generation,
attempt,
lease,
role,
path,
algo,
domain,
bytes,
object,
protective,
renewal,
adopted,
invalidation,
released,
toolchain_contract,
event_contract,
] = row
else {
return Err(StoreError::Corruption("provisional pin shape".into()));
};
Ok(ProvisionalPinRecord {
pin_key: expect_text(pin, "provisional pin key")?,
authority_key: expect_text(authority, "provisional authority")?,
action_key: expect_text(action, "provisional action")?,
generation_hex: expect_text(generation, "provisional generation")?,
attempt_hex: expect_text(attempt, "provisional attempt")?,
lease_hex: expect_text(lease, "provisional lease")?,
role_tag: match role {
SqlValue::Int(v) => *v,
_ => return Err(StoreError::Corruption("provisional role shape".into())),
},
virtual_path: match path {
SqlValue::Blob(b) => b.clone(),
_ => return Err(StoreError::Corruption("provisional path shape".into())),
},
object: self.restore_digest(algo, domain, bytes)?,
object_key: expect_text(object, "provisional object")?,
protective_pin_hex: expect_text(protective, "provisional protective pin")?,
renewal_seq: expect_u64(renewal, "provisional renewal")?,
adopted_object_key: expect_opt_text(adopted, "provisional adoption")?,
invalidated_reason: expect_opt_text(invalidation, "provisional invalidation")?,
released: expect_u64(released, "provisional released")? != 0,
toolchain_contract_key: expect_text(
toolchain_contract,
"provisional toolchain contract",
)?,
event_contract_key: expect_text(event_contract, "provisional event contract")?,
})
}
fn map_obligation_row(row: &[SqlValue]) -> Result<ProvisionalObligationRow, StoreError> {
let [
worker,
attempt_hex,
pin,
action,
generation,
producer,
role,
path,
object,
status,
resolution,
created,
] = row
else {
return Err(StoreError::Corruption("obligation row shape".into()));
};
Ok(ProvisionalObligationRow {
consumer_worker: expect_text(worker, "obligation worker")?,
consumer_attempt_hex: expect_text(attempt_hex, "obligation attempt")?,
pin_key: expect_text(pin, "obligation pin")?,
producer_action_key: expect_text(action, "obligation action")?,
producer_generation_hex: expect_text(generation, "obligation generation")?,
producer_attempt_hex: expect_text(producer, "obligation producer")?,
role_tag: match role {
SqlValue::Int(v) => *v,
_ => return Err(StoreError::Corruption("obligation role shape".into())),
},
virtual_path: match path {
SqlValue::Blob(b) => b.clone(),
_ => return Err(StoreError::Corruption("obligation path shape".into())),
},
object_key: expect_text(object, "obligation object")?,
status: expect_text(status, "obligation status")?,
resolution_object_key: expect_opt_text(resolution, "obligation resolution")?,
created_seq: expect_u64(created, "obligation seq")?,
})
}
fn map_install_row(&self, row: &[SqlValue]) -> Result<ProvisionalInstallRecord, StoreError> {
let [
pin,
worker,
attempt_hex,
path,
algo,
domain,
bytes,
object_key,
seq,
state,
] = row
else {
return Err(StoreError::Corruption("install journal shape".into()));
};
Ok(ProvisionalInstallRecord {
pin_key: expect_text(pin, "install pin")?,
consumer_worker: expect_text(worker, "install worker")?,
consumer_attempt_hex: expect_text(attempt_hex, "install attempt")?,
installed_path: match path {
SqlValue::Blob(b) => b.clone(),
_ => return Err(StoreError::Corruption("install path shape".into())),
},
object: self.restore_digest(algo, domain, bytes)?,
object_key: expect_text(object_key, "install object")?,
installed_seq: expect_u64(seq, "install seq")?,
state: expect_text(state, "install state")?,
})
}
}
fn expect_u64(v: &SqlValue, what: &str) -> Result<u64, StoreError> {
match v {
SqlValue::Int(i) => {
u64::try_from(*i).map_err(|_| StoreError::Corruption(format!("negative {what}")))
}
_ => Err(StoreError::Corruption(format!("{what} shape"))),
}
}
fn expect_u128(v: &SqlValue, what: &str) -> Result<u128, StoreError> {
match v {
SqlValue::Blob(b) => {
let bytes: [u8; 16] = b
.as_slice()
.try_into()
.map_err(|_| StoreError::Corruption(format!("{what} not 16 bytes")))?;
Ok(u128::from_be_bytes(bytes))
}
_ => Err(StoreError::Corruption(format!("{what} shape"))),
}
}
fn expect_u64_blob(v: &SqlValue, what: &str) -> Result<u64, StoreError> {
match v {
SqlValue::Blob(b) => {
let bytes: [u8; 8] = b
.as_slice()
.try_into()
.map_err(|_| StoreError::Corruption(format!("{what} not 8 bytes")))?;
Ok(u64::from_be_bytes(bytes))
}
_ => Err(StoreError::Corruption(format!("{what} shape"))),
}
}
fn decode_worker_incarnation_fence(
worker: &str,
row: &[SqlValue],
) -> Result<WorkerIncarnationFenceRecord, StoreError> {
let [
highest_boot_generation,
incarnation,
active,
operator_reenrollment_generation,
clone_ambiguous,
] = row
else {
return Err(StoreError::Corruption("worker fence shape".into()));
};
let incarnation = WorkerIncarnationId(expect_u128(incarnation, "worker incarnation")?);
let active_incarnation = match expect_u64(active, "worker fence active")? {
0 => None,
1 => Some(incarnation),
other => {
return Err(StoreError::Corruption(format!(
"worker fence active flag {other}"
)));
}
};
Ok(WorkerIncarnationFenceRecord {
worker_peer_id: PeerId(worker.to_owned()),
highest_boot_generation: WorkerBootGeneration(expect_u64_blob(
highest_boot_generation,
"worker boot generation",
)?),
active_incarnation,
clone_ambiguous: match expect_u64(clone_ambiguous, "clone ambiguous")? {
0 => false,
1 => true,
other => {
return Err(StoreError::Corruption(format!(
"clone ambiguous flag {other}"
)));
}
},
operator_reenrollment_generation: expect_u64_blob(
operator_reenrollment_generation,
"worker reenrollment generation",
)?,
})
}
fn expect_text(v: &SqlValue, what: &str) -> Result<String, StoreError> {
match v {
SqlValue::Text(t) => Ok(t.clone()),
_ => Err(StoreError::Corruption(format!("{what} shape"))),
}
}
fn expect_opt_text(v: &SqlValue, what: &str) -> Result<Option<String>, StoreError> {
match v {
SqlValue::Null => Ok(None),
SqlValue::Text(t) => Ok(Some(t.clone())),
_ => Err(StoreError::Corruption(format!("{what} shape"))),
}
}
fn to_seq(v: u64, what: &str) -> Result<i64, StoreError> {
i64::try_from(v).map_err(|_| StoreError::Corruption(format!("{what} out of range")))
}
fn evidence_key_rows(rows: &[Vec<SqlValue>]) -> Result<Vec<String>, StoreError> {
rows.iter()
.map(|row| {
let [domain, bytes] = row.as_slice() else {
return Err(StoreError::Corruption("evidence row shape".into()));
};
let SqlValue::Blob(bytes) = bytes else {
return Err(StoreError::Corruption("evidence bytes shape".into()));
};
Ok(format!(
"{}:{}",
expect_text(domain, "evidence domain")?,
hex(bytes)
))
})
.collect()
}
impl<E: SqlEngine> RabsMetadataStore for SqlMetadataStore<E> {
fn query(&mut self, sql: &str, params: &[SqlValue]) -> Result<Vec<Vec<SqlValue>>, StoreError> {
self.engine.query(sql, params)
}
fn schema_version(&mut self) -> Result<u32, StoreError> {
let rows = self
.engine
.query("SELECT MAX(version) FROM schema_epochs", &[])?;
match rows.first().and_then(|r| r.first()) {
Some(SqlValue::Int(v)) => u32::try_from(*v)
.map_err(|_| StoreError::Corruption("negative schema version".into())),
_ => Err(StoreError::Corruption("missing schema_epochs".into())),
}
}
fn intern_domain(&mut self, domain: &'static str) {
self.intern(domain);
}
fn acquire_authority(&mut self, row: &AuthorityRow) -> Result<(), StoreError> {
self.intern(row.digest.domain);
let key = digest_key(&row.digest);
let [algo, domain, bytes] = SqlMetadataStore::<E>::digest_params(&row.digest);
let cluster = row.cluster_id.clone();
let incarnation = u128_blob(row.incarnation);
let term = i64::try_from(row.term)
.map_err(|_| StoreError::Corruption("term out of range".into()))?;
let acquired = i64::try_from(row.acquired_seq)
.map_err(|_| StoreError::Corruption("acquired_seq out of range".into()))?;
self.in_txn(move |engine| {
match SqlMetadataStore::<E>::active_authority_key(engine)? {
Some(active) if active == key => return Ok(()), Some(active) => return Err(StoreError::AuthorityHeld { holder: active }),
None => {}
}
engine.execute(
"INSERT OR REPLACE INTO coordinator_authorities \
(key, algo, domain, bytes, cluster_id, incarnation, term, acquired_seq, released) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 0)",
&[
SqlValue::Text(key),
algo,
domain,
bytes,
SqlValue::Text(cluster),
SqlValue::Blob(incarnation),
SqlValue::Int(term),
SqlValue::Int(acquired),
],
)?;
Ok(())
})
}
fn release_authority(&mut self, digest: &TypedDigest) -> Result<(), StoreError> {
let key = digest_key(digest);
self.in_txn(move |engine| {
engine.execute(
"UPDATE coordinator_authorities SET released = 1 WHERE key = ?1",
&[SqlValue::Text(key)],
)?;
Ok(())
})
}
fn active_authority(&mut self) -> Result<Option<AuthorityRow>, StoreError> {
let rows = self.engine.query(
"SELECT algo, domain, bytes, cluster_id, incarnation, term, acquired_seq \
FROM coordinator_authorities WHERE released = 0",
&[],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
if rows.len() > 1 {
return Err(StoreError::Corruption(
"more than one active authority".into(),
));
}
let [algo, domain, bytes, cluster, incarnation, term, acquired] = row.as_slice() else {
return Err(StoreError::Corruption("authority row shape".into()));
};
let digest = self.restore_digest(algo, domain, bytes)?;
let SqlValue::Text(cluster) = cluster else {
return Err(StoreError::Corruption("cluster_id shape".into()));
};
let SqlValue::Blob(incarnation) = incarnation else {
return Err(StoreError::Corruption("incarnation shape".into()));
};
let incarnation_bytes: [u8; 16] = incarnation
.as_slice()
.try_into()
.map_err(|_| StoreError::Corruption("incarnation not 16 bytes".into()))?;
Ok(Some(AuthorityRow {
digest,
cluster_id: cluster.clone(),
incarnation: u128::from_be_bytes(incarnation_bytes),
term: expect_u64(term, "term")?,
acquired_seq: expect_u64(acquired, "acquired_seq")?,
}))
}
fn upsert_action_entry(&mut self, row: &ActionEntryRow) -> Result<(), StoreError> {
self.intern(row.action_key.domain);
let key = digest_key(&row.action_key);
let [algo, domain, bytes] = SqlMetadataStore::<E>::digest_params(&row.action_key);
let key_epoch = i64::from(row.key_epoch);
let projection_epoch = i64::from(row.projection_epoch);
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO action_entries \
(key, algo, domain, bytes, key_epoch, projection_epoch) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
&[
SqlValue::Text(key),
algo,
domain,
bytes,
SqlValue::Int(key_epoch),
SqlValue::Int(projection_epoch),
],
)?;
Ok(())
})
}
fn lookup_action(&mut self, key: &TypedDigest) -> Result<Option<ActionEntryRow>, StoreError> {
let rows = self.engine.query(
"SELECT algo, domain, bytes, key_epoch, projection_epoch \
FROM action_entries WHERE key = ?1",
&[SqlValue::Text(digest_key(key))],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [algo, domain, bytes, key_epoch, projection_epoch] = row.as_slice() else {
return Err(StoreError::Corruption("action entry shape".into()));
};
let action_key = self.restore_digest(algo, domain, bytes)?;
let key_epoch = u32::try_from(expect_u64(key_epoch, "key_epoch")?)
.map_err(|_| StoreError::Corruption("key_epoch out of range".into()))?;
let projection_epoch = u32::try_from(expect_u64(projection_epoch, "projection_epoch")?)
.map_err(|_| StoreError::Corruption("projection_epoch out of range".into()))?;
Ok(Some(ActionEntryRow {
action_key,
key_epoch,
projection_epoch,
}))
}
fn create_generation(
&mut self,
authority: &TypedDigest,
id: u128,
action_key: &TypedDigest,
) -> Result<(), StoreError> {
let authority = authority.clone();
let action = digest_key(action_key);
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let high_water = SqlMetadataStore::<E>::generation_high_water(engine)?;
if id <= high_water {
return Err(StoreError::GenerationIdNotAboveHighWater);
}
engine.execute(
"INSERT INTO action_generations (id_hex, id, action_key, authority_key, tombstoned) \
VALUES (?1, ?2, ?3, ?4, 0)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(action),
SqlValue::Text(digest_key(&authority)),
],
)?;
engine.execute(
"INSERT OR REPLACE INTO generation_high_water (kind, value) \
VALUES ('action-generation', ?1)",
&[SqlValue::Blob(u128_blob(id))],
)?;
Ok(())
})
}
fn create_bound_generation(
&mut self,
authority: &TypedDigest,
generation: &ActionGeneration,
action_key: &TypedDigest,
) -> Result<(), StoreError> {
if generation.created_under_authority_digest != *authority {
return Err(StoreError::AttemptAuthorityMismatch);
}
let authority = authority.clone();
let action = digest_key(action_key);
let generation = generation.clone();
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let id = generation.generation_id.0;
let high_water = SqlMetadataStore::<E>::generation_high_water(engine)?;
if id <= high_water {
return Err(StoreError::GenerationIdNotAboveHighWater);
}
engine.execute(
"INSERT INTO action_generations \
(id_hex, id, action_key, authority_key, tombstoned, per_key_ordinal) \
VALUES (?1, ?2, ?3, ?4, 0, ?5)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(action),
SqlValue::Text(digest_key(&authority)),
SqlValue::Blob(u64_blob(generation.per_key_ordinal)),
],
)?;
engine.execute(
"INSERT OR REPLACE INTO generation_high_water (kind, value) \
VALUES ('action-generation', ?1)",
&[SqlValue::Blob(u128_blob(id))],
)?;
Ok(())
})
}
fn allocate_bound_generation(
&mut self,
authority: &TypedDigest,
action_key: &TypedDigest,
) -> Result<ActionGeneration, StoreError> {
let authority = authority.clone();
let action = digest_key(action_key);
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let id = SqlMetadataStore::<E>::generation_high_water(engine)?
.checked_add(1)
.ok_or(StoreError::GenerationIdentityExhausted)?;
let malformed = engine.query(
"SELECT per_key_ordinal FROM action_generations \
WHERE action_key = ?1 AND per_key_ordinal IS NOT NULL \
AND (typeof(per_key_ordinal) != 'blob' OR length(per_key_ordinal) != 8) \
LIMIT 1",
&[SqlValue::Text(action.clone())],
)?;
if !malformed.is_empty() {
return Err(StoreError::Corruption(
"generation ordinal representation".into(),
));
}
let rows = engine.query(
"SELECT per_key_ordinal FROM action_generations \
WHERE action_key = ?1 AND per_key_ordinal IS NOT NULL \
ORDER BY per_key_ordinal DESC LIMIT 1",
&[SqlValue::Text(action.clone())],
)?;
let previous = match rows.as_slice() {
[] => 0,
[row] => match row.as_slice() {
[value] => expect_u64_blob(value, "generation ordinal")?,
_ => return Err(StoreError::Corruption("generation ordinal shape".into())),
},
_ => {
return Err(StoreError::Corruption(
"generation ordinal row count".into(),
));
}
};
let ordinal = previous
.checked_add(1)
.ok_or(StoreError::GenerationOrdinalExhausted)?;
let generation = ActionGeneration {
generation_id: ActionGenerationId(id),
per_key_ordinal: ordinal,
created_under_authority_digest: authority.clone(),
};
engine.execute(
"INSERT INTO action_generations \
(id_hex, id, action_key, authority_key, tombstoned, per_key_ordinal) \
VALUES (?1, ?2, ?3, ?4, 0, ?5)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(action),
SqlValue::Text(digest_key(&authority)),
SqlValue::Blob(u64_blob(ordinal)),
],
)?;
engine.execute(
"INSERT OR REPLACE INTO generation_high_water (kind, value) \
VALUES ('action-generation', ?1)",
&[SqlValue::Blob(u128_blob(id))],
)?;
Ok(generation)
})
}
fn close_generations_for_other_authorities(
&mut self,
active: &TypedDigest,
) -> Result<u64, StoreError> {
let active = digest_key(active);
self.in_txn(move |engine| {
let closed = engine.execute(
"UPDATE action_generations SET tombstoned = 1 \
WHERE tombstoned = 0 AND authority_key != ?1",
&[SqlValue::Text(active)],
)?;
u64::try_from(closed)
.map_err(|_| StoreError::Corruption("generation close count out of range".into()))
})
}
fn tombstone_generation(&mut self, id: u128) -> Result<(), StoreError> {
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE action_generations SET tombstoned = 1 WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
if changed == 0 {
return Err(StoreError::UnknownGeneration);
}
Ok(())
})
}
fn record_attempt(
&mut self,
id: u128,
generation: u128,
worker: &str,
seq: u64,
) -> Result<(), StoreError> {
let worker = worker.to_owned();
let seq =
i64::try_from(seq).map_err(|_| StoreError::Corruption("seq out of range".into()))?;
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT id_hex FROM action_attempts WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
if !existing.is_empty() {
return Err(StoreError::DuplicateAttempt);
}
let generation_exists = engine.query(
"SELECT id_hex FROM action_generations WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(generation))],
)?;
if generation_exists.is_empty() {
return Err(StoreError::UnknownGeneration);
}
engine.execute(
"INSERT INTO action_attempts (id_hex, id, generation_hex, worker, seq) \
VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(u128_hex(generation)),
SqlValue::Text(worker),
SqlValue::Int(seq),
],
)?;
Ok(())
})
}
fn admit_attempt_lease(
&mut self,
authority: &AttemptAuthority,
recorded_seq: u64,
expires_at_own_monotonic_ms: u64,
) -> Result<(), StoreError> {
let authority = authority.clone();
let recorded = i64::try_from(recorded_seq)
.map_err(|_| StoreError::Corruption("recorded_seq out of range".into()))?;
let renewal = i64::try_from(authority.lease_renewal_seq.0)
.map_err(|_| StoreError::Corruption("renewal_seq out of range".into()))?;
let expires = i64::try_from(expires_at_own_monotonic_ms)
.map_err(|_| StoreError::Corruption("expires_at_seq out of range".into()))?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_attempt_context(engine, &authority)?;
SqlMetadataStore::<E>::require_worker_lease_binding(
engine,
&authority.worker_peer_id,
authority.worker_boot_generation,
authority.worker_incarnation_id,
)?;
let attempt_id = authority.attempt_id.0;
if !engine
.query(
"SELECT id_hex FROM action_attempts WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(attempt_id))],
)?
.is_empty()
{
return Err(StoreError::DuplicateAttempt);
}
let lease_id = authority.execution_lease_id.0;
if !engine
.query(
"SELECT id_hex FROM execution_leases WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(lease_id))],
)?
.is_empty()
{
return Err(StoreError::DuplicateLease);
}
engine.execute(
"INSERT INTO action_attempts \
(id_hex, id, generation_hex, worker, seq, \
worker_boot_generation, worker_incarnation, execution_lease_hex) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
&[
SqlValue::Text(u128_hex(attempt_id)),
SqlValue::Blob(u128_blob(attempt_id)),
SqlValue::Text(u128_hex(authority.action_generation.generation_id.0)),
SqlValue::Text(authority.worker_peer_id.0.clone()),
SqlValue::Int(recorded),
SqlValue::Blob(u64_blob(authority.worker_boot_generation.0)),
SqlValue::Blob(u128_blob(authority.worker_incarnation_id.0)),
SqlValue::Text(u128_hex(lease_id)),
],
)?;
engine.execute(
"INSERT INTO execution_leases \
(id_hex, id, attempt_hex, renewal_seq, expires_at_seq, released) \
VALUES (?1, ?2, ?3, ?4, ?5, 0)",
&[
SqlValue::Text(u128_hex(lease_id)),
SqlValue::Blob(u128_blob(lease_id)),
SqlValue::Text(u128_hex(attempt_id)),
SqlValue::Int(renewal),
SqlValue::Int(expires),
],
)?;
Ok(())
})
}
fn renew_attempt_lease(
&mut self,
authority: &AttemptAuthority,
renewal: LeaseRenewal,
expires_at_own_monotonic_ms: u64,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError> {
let authority = authority.clone();
let expires = i64::try_from(expires_at_own_monotonic_ms)
.map_err(|_| StoreError::Corruption("expires_at_seq out of range".into()))?;
self.in_txn(move |engine| {
let state = SqlMetadataStore::<E>::bound_lease_state(engine, &authority)?;
if state.released {
return Err(StoreError::LeaseReleased);
}
let now = own_monotonic_now_ms();
if now >= state.expires_at_own_monotonic_ms || expires_at_own_monotonic_ms <= now {
return Err(StoreError::LeaseExpired);
}
if renewal.lease != authority.execution_lease_id {
return Err(StoreError::LeaseAttemptMismatch);
}
if authority.lease_renewal_seq.0 != state.renewal_seq {
return Err(StoreError::LeaseRenewalMismatch);
}
if renewal.seq.0 <= state.renewal_seq {
return Err(StoreError::NonMonotonicRenewal);
}
let renewal_seq = i64::try_from(renewal.seq.0)
.map_err(|_| StoreError::Corruption("renewal_seq out of range".into()))?;
engine.execute(
"UPDATE execution_leases SET renewal_seq = ?1, expires_at_seq = ?2 \
WHERE id_hex = ?3",
&[
SqlValue::Int(renewal_seq),
SqlValue::Int(expires),
SqlValue::Text(u128_hex(authority.execution_lease_id.0)),
],
)?;
let now = own_monotonic_now_ms();
if now >= state.expires_at_own_monotonic_ms || now >= expires_at_own_monotonic_ms {
return Err(StoreError::LeaseExpired);
}
Ok(())
})
}
fn release_lease(&mut self, id: u128) -> Result<(), StoreError> {
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE execution_leases SET released = 1 WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
if changed == 0 {
return Err(StoreError::UnknownLease);
}
Ok(())
})
}
fn commit_publication(
&mut self,
authority: &TypedDigest,
attempt_authority: Option<(&AttemptAuthority, &dyn Fn() -> u64)>,
row: &PublicationRow,
) -> Result<CommitOutcome, StoreError> {
self.intern(row.action_key.domain);
self.intern(row.descriptor_digest.domain);
self.intern(row.manifest_digest.domain);
self.intern(row.evidence_digest.domain);
let authority = authority.clone();
let attempt_authority = attempt_authority.map(|(attempt, now)| (attempt.clone(), now));
let row = row.clone();
self.in_txn(move |engine| {
let lease_clock_and_deadline = match &attempt_authority {
Some((attempt, own_monotonic_now_ms)) => {
if SqlMetadataStore::<E>::attempt_authority_digest(attempt) != authority
|| row.action_key != attempt.action_key
|| row.winner_generation != attempt.action_generation.generation_id.0
|| row.winner_attempt != attempt.attempt_id.0
{
return Err(StoreError::AttemptAuthorityMismatch);
}
let state = SqlMetadataStore::<E>::bound_lease_state(engine, attempt)?;
if state.released {
return Err(StoreError::LeaseReleased);
}
if own_monotonic_now_ms() >= state.expires_at_own_monotonic_ms {
return Err(StoreError::LeaseExpired);
}
if state.renewal_seq != attempt.lease_renewal_seq.0 {
return Err(StoreError::LeaseRenewalMismatch);
}
Some((own_monotonic_now_ms, state.expires_at_own_monotonic_ms))
}
None => {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
None
}
};
let require_live_lease = || {
if let Some((clock, deadline)) = lease_clock_and_deadline
&& clock() >= deadline
{
return Err(StoreError::LeaseExpired);
}
Ok(())
};
let action = digest_key(&row.action_key);
let existing = engine.query(
"SELECT descriptor_domain, descriptor_bytes FROM action_publications \
WHERE action_key = ?1",
&[SqlValue::Text(action.clone())],
)?;
require_live_lease()?;
if let Some(existing_row) = existing.first() {
let [domain, bytes] = existing_row.as_slice() else {
return Err(StoreError::Corruption("publication row shape".into()));
};
let (SqlValue::Text(domain), SqlValue::Blob(bytes)) = (domain, bytes) else {
return Err(StoreError::Corruption("publication digest shape".into()));
};
let same = domain == row.descriptor_digest.domain
&& bytes.as_slice() == row.descriptor_digest.bytes.as_slice();
if same {
return Ok(CommitOutcome::IdempotentDuplicate);
}
engine.execute(
"INSERT OR REPLACE INTO quarantines (scope, subject, reason) \
VALUES ('action-entry', ?1, 'publication descriptor conflict')",
&[SqlValue::Text(action)],
)?;
require_live_lease()?;
return Ok(CommitOutcome::ConflictQuarantined);
}
let [d_algo, d_domain, d_bytes] =
SqlMetadataStore::<E>::digest_params(&row.descriptor_digest);
let [m_algo, m_domain, m_bytes] =
SqlMetadataStore::<E>::digest_params(&row.manifest_digest);
engine.execute(
"INSERT INTO action_publications \
(action_key, descriptor_algo, descriptor_domain, descriptor_bytes, \
manifest_algo, manifest_domain, manifest_bytes, \
winner_generation_hex, winner_attempt_hex, result_kind, pin_hex) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
&[
SqlValue::Text(action.clone()),
d_algo,
d_domain,
d_bytes,
m_algo,
m_domain,
m_bytes,
SqlValue::Text(u128_hex(row.winner_generation)),
SqlValue::Text(u128_hex(row.winner_attempt)),
SqlValue::Text(row.result_kind.as_str().to_owned()),
SqlValue::Text(u128_hex(row.pin_id)),
],
)?;
let serving_exists = engine.query(
"SELECT 1 FROM action_serving_states WHERE action_key = ?1",
&[SqlValue::Text(action.clone())],
)?;
if serving_exists.is_empty() {
engine.execute(
"INSERT INTO action_serving_states (action_key, disposition, version) \
VALUES (?1, 'servable', 1)",
&[SqlValue::Text(action.clone())],
)?;
} else {
Self::set_serving_disposition_row(engine, &action, "servable")?;
}
let [e_algo, e_domain, e_bytes] =
SqlMetadataStore::<E>::digest_params(&row.evidence_digest);
engine.execute(
"INSERT OR IGNORE INTO action_evidence_index \
(action_key, evidence_algo, evidence_domain, evidence_bytes, \
generation_hex, attempt_hex, manifest_key) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
&[
SqlValue::Text(action),
e_algo,
e_domain,
e_bytes,
SqlValue::Text(u128_hex(row.winner_generation)),
SqlValue::Text(u128_hex(row.winner_attempt)),
SqlValue::Text(digest_key(&row.manifest_digest)),
],
)?;
engine.execute(
"INSERT INTO pins (id_hex, id, root_key, owner, class, expires_at_seq, released, \
evidence, renewal_seq, durable, reason) \
VALUES (?1, ?2, ?3, ?4, 'action-publication', NULL, 0, ?5, 0, 1, \
'publication reachability root')",
&[
SqlValue::Text(u128_hex(row.pin_id)),
SqlValue::Blob(u128_blob(row.pin_id)),
SqlValue::Text(digest_key(&row.manifest_digest)),
SqlValue::Text(row.pin_owner.clone()),
SqlValue::Text(digest_key(&row.action_key)),
],
)?;
for ancestor in &row.provisional_ancestors {
engine.execute(
"INSERT OR IGNORE INTO provisional_ancestry \
(consumer_action_key, producer_action_key, role, virtual_path, \
object_key, adopted) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
&[
SqlValue::Text(digest_key(&row.action_key)),
SqlValue::Text(ancestor.producer_action_key.clone()),
SqlValue::Text(ancestor.role.clone()),
SqlValue::Blob(ancestor.virtual_path.clone()),
SqlValue::Text(ancestor.object_key.clone()),
SqlValue::Int(i64::from(ancestor.adopted)),
],
)?;
}
require_live_lease()?;
Ok(CommitOutcome::Committed)
})
}
fn append_evidence(
&mut self,
action: &TypedDigest,
manifest_key: &str,
evidence: &TypedDigest,
generation: u128,
attempt: u128,
) -> Result<(), StoreError> {
self.intern(evidence.domain);
self.in_txn(|engine| {
Self::append_evidence_row(engine, action, manifest_key, evidence, generation, attempt)
})
}
fn append_evidence_for_attempt(
&mut self,
authority: &AttemptAuthority,
manifest_key: &str,
evidence: &TypedDigest,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError> {
self.intern(evidence.domain);
self.in_txn(|engine| {
Self::require_live_attempt(engine, authority, own_monotonic_now_ms)?;
Self::append_evidence_row(
engine,
&authority.action_key,
manifest_key,
evidence,
authority.action_generation.generation_id.0,
authority.attempt_id.0,
)?;
Self::require_live_attempt(engine, authority, own_monotonic_now_ms)
})
}
fn has_publication(&mut self, action: &TypedDigest) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT action_key FROM action_publications WHERE action_key = ?1",
&[SqlValue::Text(digest_key(action))],
)?;
Ok(!rows.is_empty())
}
fn published_manifest_key(
&mut self,
action: &TypedDigest,
) -> Result<Option<String>, StoreError> {
self.published_manifest_key_str(&digest_key(action))
}
fn published_manifest_key_str(
&mut self,
action_key: &str,
) -> Result<Option<String>, StoreError> {
let rows = self.engine.query(
"SELECT manifest_domain, manifest_bytes FROM action_publications WHERE action_key = ?1",
&[SqlValue::Text(action_key.to_owned())],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [SqlValue::Text(domain), SqlValue::Blob(bytes)] = row.as_slice() else {
return Err(StoreError::Corruption("publication manifest shape".into()));
};
Ok(Some(format!("{}:{}", domain, hex(bytes))))
}
fn generation_state(&mut self, id: u128) -> Result<Option<GenerationState>, StoreError> {
let rows = self.engine.query(
"SELECT tombstoned FROM action_generations WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(v) => Ok(Some(GenerationState {
tombstoned: expect_u64(v, "tombstoned")? != 0,
})),
}
}
fn attempt_exists(&mut self, id: u128, generation: u128) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT id_hex FROM action_attempts WHERE id_hex = ?1 AND generation_hex = ?2",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Text(u128_hex(generation)),
],
)?;
Ok(!rows.is_empty())
}
fn lease_state(&mut self, id: u128) -> Result<Option<LeaseState>, StoreError> {
let rows = self.engine.query(
"SELECT released, renewal_seq, expires_at_seq FROM execution_leases WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [released, renewal_seq, expires_at_seq] = row.as_slice() else {
return Err(StoreError::Corruption("lease state shape".into()));
};
Ok(Some(LeaseState {
released: expect_u64(released, "released")? != 0,
renewal_seq: expect_u64(renewal_seq, "renewal_seq")?,
expires_at_own_monotonic_ms: expect_u64(expires_at_seq, "expires_at_seq")?,
}))
}
fn validate_attempt_lease(
&mut self,
authority: &AttemptAuthority,
own_monotonic_now_ms: u64,
) -> Result<LeaseState, StoreError> {
let state = SqlMetadataStore::<E>::bound_lease_state(&mut self.engine, authority)?;
if state.released {
return Err(StoreError::LeaseReleased);
}
if own_monotonic_now_ms >= state.expires_at_own_monotonic_ms {
return Err(StoreError::LeaseExpired);
}
if state.renewal_seq != authority.lease_renewal_seq.0 {
return Err(StoreError::LeaseRenewalMismatch);
}
Ok(state)
}
fn object_located(&mut self, object: &TypedDigest) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT object_key FROM object_locations \
WHERE object_key = ?1 AND quarantined = 0 LIMIT 1",
&[SqlValue::Text(digest_key(object))],
)?;
Ok(!rows.is_empty())
}
fn object_locations(
&mut self,
object: &TypedDigest,
) -> Result<Vec<(String, String, bool)>, StoreError> {
let rows = self.engine.query(
"SELECT store_path, encoding, durable FROM object_locations \
WHERE object_key = ?1 AND quarantined = 0 ORDER BY store_path",
&[SqlValue::Text(digest_key(object))],
)?;
rows.iter()
.map(|row| {
let [path, encoding, durable] = row.as_slice() else {
return Err(StoreError::Corruption("object location shape".into()));
};
Ok((
expect_text(path, "store_path")?,
expect_text(encoding, "encoding")?,
expect_u64(durable, "durable")? != 0,
))
})
.collect()
}
fn record_object(&mut self, id: &TypedDigest, logical_size: u64) -> Result<(), StoreError> {
self.intern(id.domain);
let key = digest_key(id);
let [algo, domain, bytes] = SqlMetadataStore::<E>::digest_params(id);
let size = i64::try_from(logical_size)
.map_err(|_| StoreError::Corruption("logical_size out of range".into()))?;
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO objects (key, algo, domain, bytes, logical_size) \
VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(key),
algo,
domain,
bytes,
SqlValue::Int(size),
],
)?;
Ok(())
})
}
fn add_location(
&mut self,
object: &TypedDigest,
store_path: &str,
verified_seq: Option<u64>,
encoding: &str,
durable: bool,
) -> Result<(), StoreError> {
let key = digest_key(object);
let path = store_path.to_owned();
let encoding = encoding.to_owned();
let verified = match verified_seq {
None => SqlValue::Null,
Some(v) => SqlValue::Int(
i64::try_from(v)
.map_err(|_| StoreError::Corruption("verified_seq out of range".into()))?,
),
};
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT quarantined FROM object_locations WHERE object_key = ?1 AND store_path = ?2",
&[SqlValue::Text(key.clone()), SqlValue::Text(path.clone())],
)?;
let quarantined = if let Some(row) = existing.first() {
let [value] = row.as_slice() else {
return Err(StoreError::Corruption("location quarantine shape".into()));
};
expect_u64(value, "quarantined")? != 0
} else {
false
};
engine.execute(
"INSERT OR REPLACE INTO object_locations \
(object_key, store_path, verified_seq, encoding, quarantined, durable) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
&[
SqlValue::Text(key),
SqlValue::Text(path),
verified,
SqlValue::Text(encoding),
SqlValue::Int(i64::from(quarantined)),
SqlValue::Int(i64::from(durable)),
],
)?;
Ok(())
})
}
fn object_durably_located(&mut self, object: &TypedDigest) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT object_key FROM object_locations \
WHERE object_key = ?1 AND quarantined = 0 AND durable = 1 LIMIT 1",
&[SqlValue::Text(digest_key(object))],
)?;
Ok(!rows.is_empty())
}
fn set_location_quarantined(
&mut self,
object: &TypedDigest,
store_path: &str,
quarantined: bool,
) -> Result<(), StoreError> {
let key = digest_key(object);
let path = store_path.to_owned();
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE object_locations SET quarantined = ?1 \
WHERE object_key = ?2 AND store_path = ?3",
&[
SqlValue::Int(i64::from(quarantined)),
SqlValue::Text(key),
SqlValue::Text(path),
],
)?;
if changed == 0 {
return Err(StoreError::Corruption("unknown location".into()));
}
Ok(())
})
}
fn add_object_edge(
&mut self,
parent: &TypedDigest,
child: &TypedDigest,
kind: &str,
) -> Result<(), StoreError> {
let parent = digest_key(parent);
let child = digest_key(child);
let kind = kind.to_owned();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO object_edges (parent_key, child_key, kind) \
VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(parent),
SqlValue::Text(child),
SqlValue::Text(kind),
],
)?;
Ok(())
})
}
fn create_pin(
&mut self,
id: u128,
root: &TypedDigest,
owner: &str,
class: &str,
expires_at_seq: Option<u64>,
evidence: Option<&str>,
durable: bool,
reason: &str,
) -> Result<(), StoreError> {
let root = digest_key(root);
let owner = owner.to_owned();
let class = class.to_owned();
let evidence = match evidence {
None => SqlValue::Null,
Some(e) => SqlValue::Text(e.to_owned()),
};
let reason = reason.to_owned();
let expires = match expires_at_seq {
None => SqlValue::Null,
Some(v) => SqlValue::Int(
i64::try_from(v)
.map_err(|_| StoreError::Corruption("expires_at_seq out of range".into()))?,
),
};
self.in_txn(move |engine| {
engine.execute(
"INSERT INTO pins (id_hex, id, root_key, owner, class, expires_at_seq, released, \
evidence, renewal_seq, durable, reason) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, 0, ?7, 0, ?8, ?9)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(root),
SqlValue::Text(owner),
SqlValue::Text(class),
expires,
evidence,
SqlValue::Int(i64::from(durable)),
SqlValue::Text(reason),
],
)?;
Ok(())
})
}
fn renew_pin(&mut self, id: u128, renewal_seq: u64) -> Result<(), StoreError> {
self.in_txn(move |engine| {
let rows = engine.query(
"SELECT renewal_seq, released FROM pins WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownPin);
};
let [stored_seq, released] = row.as_slice() else {
return Err(StoreError::Corruption("pin row shape".into()));
};
if expect_u64(released, "released")? != 0 {
return Err(StoreError::PinReleased);
}
if renewal_seq <= expect_u64(stored_seq, "renewal_seq")? {
return Err(StoreError::NonMonotonicPinRenewal);
}
let renewal = i64::try_from(renewal_seq)
.map_err(|_| StoreError::Corruption("renewal_seq out of range".into()))?;
engine.execute(
"UPDATE pins SET renewal_seq = ?1 WHERE id_hex = ?2",
&[SqlValue::Int(renewal), SqlValue::Text(u128_hex(id))],
)?;
Ok(())
})
}
fn release_pin(&mut self, id: u128, owner: &str) -> Result<(), StoreError> {
let owner = owner.to_owned();
self.in_txn(move |engine| {
let rows = engine.query(
"SELECT owner FROM pins WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownPin);
};
match row.first() {
Some(SqlValue::Text(stored)) if *stored == owner => {}
Some(SqlValue::Text(_)) => return Err(StoreError::PinOwnerMismatch),
_ => return Err(StoreError::Corruption("pin owner shape".into())),
}
engine.execute(
"UPDATE pins SET released = 1 WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
Ok(())
})
}
fn put_recipe(&mut self, action: &TypedDigest, recipe: &TypedDigest) -> Result<(), StoreError> {
self.intern(recipe.domain);
let action = digest_key(action);
let [algo, domain, bytes] = SqlMetadataStore::<E>::digest_params(recipe);
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO observed_input_recipes \
(action_key, recipe_algo, recipe_domain, recipe_bytes) VALUES (?1, ?2, ?3, ?4)",
&[SqlValue::Text(action), algo, domain, bytes],
)?;
Ok(())
})
}
fn put_key_breakdown(
&mut self,
action: &TypedDigest,
component: &str,
digest: &TypedDigest,
) -> Result<(), StoreError> {
self.intern(digest.domain);
let action = digest_key(action);
let component = component.to_owned();
let [algo, domain, bytes] = SqlMetadataStore::<E>::digest_params(digest);
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO key_breakdowns \
(action_key, component, algo, domain, bytes) VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(action),
SqlValue::Text(component),
algo,
domain,
bytes,
],
)?;
Ok(())
})
}
fn set_trust(
&mut self,
action: &TypedDigest,
state: &str,
reason: &str,
) -> Result<(), StoreError> {
let action = digest_key(action);
let state = state.to_owned();
let reason = reason.to_owned();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO trust_states (action_key, state, reason) \
VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(action),
SqlValue::Text(state),
SqlValue::Text(reason),
],
)?;
Ok(())
})
}
fn add_quarantine(
&mut self,
scope: QuarantineScope,
subject: &str,
reason: &str,
) -> Result<(), StoreError> {
let subject = subject.to_owned();
let reason = reason.to_owned();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO quarantines (scope, subject, reason) \
VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(scope.as_str().to_owned()),
SqlValue::Text(subject),
SqlValue::Text(reason),
],
)?;
Ok(())
})
}
fn record_verification_sample(
&mut self,
action: &TypedDigest,
attempt: u128,
passed: bool,
seq: u64,
) -> Result<(), StoreError> {
let action = digest_key(action);
let seq =
i64::try_from(seq).map_err(|_| StoreError::Corruption("seq out of range".into()))?;
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO verification_samples (action_key, attempt_hex, passed, seq) \
VALUES (?1, ?2, ?3, ?4)",
&[
SqlValue::Text(action),
SqlValue::Text(u128_hex(attempt)),
SqlValue::Int(i64::from(passed)),
SqlValue::Int(seq),
],
)?;
Ok(())
})
}
fn list_evidence_keys(&mut self, action: &TypedDigest) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT evidence_domain, evidence_bytes FROM action_evidence_index \
WHERE action_key = ?1 ORDER BY evidence_domain, evidence_bytes",
&[SqlValue::Text(digest_key(action))],
)?;
evidence_key_rows(&rows)
}
fn list_evidence_keys_for_manifest(
&mut self,
manifest_key: &str,
) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT evidence_domain, evidence_bytes FROM action_evidence_index \
WHERE manifest_key = ?1 ORDER BY evidence_domain, evidence_bytes",
&[SqlValue::Text(manifest_key.to_owned())],
)?;
evidence_key_rows(&rows)
}
fn list_verification_samples(
&mut self,
action: &TypedDigest,
) -> Result<Vec<VerificationSampleRow>, StoreError> {
let rows = self.engine.query(
"SELECT attempt_hex, passed, seq FROM verification_samples \
WHERE action_key = ?1 ORDER BY attempt_hex, seq",
&[SqlValue::Text(digest_key(action))],
)?;
rows.iter()
.map(|row| {
let [attempt_hex, passed, seq] = row.as_slice() else {
return Err(StoreError::Corruption("verification sample shape".into()));
};
Ok(VerificationSampleRow {
attempt_hex: expect_text(attempt_hex, "attempt_hex")?,
passed: expect_u64(passed, "passed")? != 0,
seq: expect_u64(seq, "seq")?,
})
})
.collect()
}
fn attempt_worker_by_hex(&mut self, attempt_hex: &str) -> Result<Option<String>, StoreError> {
let rows = self.engine.query(
"SELECT worker FROM action_attempts WHERE id_hex = ?1",
&[SqlValue::Text(attempt_hex.to_owned())],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(v) => Ok(Some(expect_text(v, "worker")?)),
}
}
fn gc_snapshot(&mut self, seq: u64) -> Result<GcSnapshot, StoreError> {
let seq =
i64::try_from(seq).map_err(|_| StoreError::Corruption("seq out of range".into()))?;
self.in_txn(move |engine| {
let pinned = engine.query(
"SELECT root_key FROM pins WHERE released = 0 ORDER BY root_key",
&[],
)?;
let located = engine.query(
"SELECT DISTINCT object_key FROM object_locations ORDER BY object_key",
&[],
)?;
let text_col = |rows: Vec<Vec<SqlValue>>| -> Result<Vec<String>, StoreError> {
rows.into_iter()
.map(|row| match row.into_iter().next() {
Some(SqlValue::Text(t)) => Ok(t),
_ => Err(StoreError::Corruption("gc column shape".into())),
})
.collect()
};
let pinned_roots = text_col(pinned)?;
let located_objects = text_col(located)?;
let edges = engine.query(
"SELECT parent_key, child_key FROM object_edges ORDER BY parent_key, child_key",
&[],
)?;
let mut children: std::collections::BTreeMap<String, Vec<String>> =
std::collections::BTreeMap::new();
for row in edges {
let [SqlValue::Text(parent), SqlValue::Text(child)] = row.as_slice() else {
return Err(StoreError::Corruption("edge row shape".into()));
};
children
.entry(parent.clone())
.or_default()
.push(child.clone());
}
let mut reachable: std::collections::BTreeSet<String> =
std::collections::BTreeSet::new();
let mut frontier: Vec<String> = pinned_roots.clone();
while let Some(key) = frontier.pop() {
if !reachable.insert(key.clone()) {
continue;
}
if let Some(next) = children.get(&key) {
frontier.extend(next.iter().cloned());
}
}
let reachable_from_pins: Vec<String> = reachable.into_iter().collect();
let pinned_count = i64::try_from(pinned_roots.len())
.map_err(|_| StoreError::Corruption("pin count".into()))?;
let located_count = i64::try_from(located_objects.len())
.map_err(|_| StoreError::Corruption("location count".into()))?;
let reachable_count = i64::try_from(reachable_from_pins.len())
.map_err(|_| StoreError::Corruption("reachable count".into()))?;
engine.execute(
"INSERT INTO gc_runs (seq, pinned_roots, located_objects, reachable_objects) \
VALUES (?1, ?2, ?3, ?4)",
&[
SqlValue::Int(seq),
SqlValue::Int(pinned_count),
SqlValue::Int(located_count),
SqlValue::Int(reachable_count),
],
)?;
Ok(GcSnapshot {
pinned_roots,
located_objects,
reachable_from_pins,
})
})
}
fn reconciliation_scan(&mut self) -> Result<Vec<ReconciliationRow>, StoreError> {
let rows = self.engine.query(
"SELECT object_key, store_path, verified_seq, encoding, quarantined \
FROM object_locations ORDER BY object_key, store_path",
&[],
)?;
rows.into_iter()
.map(|row| {
let [key, path, verified, encoding, quarantined] = row.as_slice() else {
return Err(StoreError::Corruption("location row shape".into()));
};
let (SqlValue::Text(key), SqlValue::Text(path), SqlValue::Text(encoding)) =
(key, path, encoding)
else {
return Err(StoreError::Corruption("location column shape".into()));
};
let verified_seq = match verified {
SqlValue::Null => None,
other => Some(expect_u64(other, "verified_seq")?),
};
Ok(ReconciliationRow {
object_key: key.clone(),
store_path: path.clone(),
verified_seq,
encoding: encoding.clone(),
quarantined: expect_u64(quarantined, "quarantined")? != 0,
})
})
.collect()
}
fn remove_location_by_key(
&mut self,
object_key: &str,
store_path: &str,
) -> Result<bool, StoreError> {
let object_key = object_key.to_owned();
let store_path = store_path.to_owned();
self.in_txn(move |engine| {
let removed = engine.execute(
"DELETE FROM object_locations WHERE object_key = ?1 AND store_path = ?2",
&[SqlValue::Text(object_key), SqlValue::Text(store_path)],
)?;
Ok(removed > 0)
})
}
fn record_gc_receipt(&mut self, receipt: &GcReceiptRow) -> Result<(), StoreError> {
let to_int = |v: u64, what: &str| -> Result<i64, StoreError> {
i64::try_from(v).map_err(|_| StoreError::Corruption(format!("{what} out of range")))
};
let seq = to_int(receipt.seq, "seq")?;
let planned = to_int(receipt.planned, "planned")?;
let reclaimed = to_int(receipt.reclaimed, "reclaimed")?;
let skipped = to_int(receipt.skipped, "skipped")?;
let mode = receipt.mode.clone();
let truncated = i64::from(receipt.truncated);
self.in_txn(move |engine| {
engine.execute(
"INSERT INTO gc_receipts (seq, mode, planned, reclaimed, skipped, truncated) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
&[
SqlValue::Int(seq),
SqlValue::Text(mode),
SqlValue::Int(planned),
SqlValue::Int(reclaimed),
SqlValue::Int(skipped),
SqlValue::Int(truncated),
],
)?;
Ok(())
})
}
fn add_gc_tombstone(
&mut self,
object_key: &str,
store_path: &str,
marked_seq: u64,
grace_until_seq: u64,
) -> Result<(), StoreError> {
let object_key = object_key.to_owned();
let store_path = store_path.to_owned();
let marked = i64::try_from(marked_seq)
.map_err(|_| StoreError::Corruption("marked_seq out of range".into()))?;
let grace = i64::try_from(grace_until_seq)
.map_err(|_| StoreError::Corruption("grace_until_seq out of range".into()))?;
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT object_key FROM gc_tombstones WHERE object_key = ?1 AND store_path = ?2",
&[
SqlValue::Text(object_key.clone()),
SqlValue::Text(store_path.clone()),
],
)?;
if !existing.is_empty() {
return Ok(());
}
engine.execute(
"INSERT INTO gc_tombstones (object_key, store_path, marked_seq, grace_until_seq) \
VALUES (?1, ?2, ?3, ?4)",
&[
SqlValue::Text(object_key),
SqlValue::Text(store_path),
SqlValue::Int(marked),
SqlValue::Int(grace),
],
)?;
Ok(())
})
}
fn due_gc_tombstones(&mut self, now_seq: u64) -> Result<Vec<GcTombstoneRow>, StoreError> {
let now = i64::try_from(now_seq).unwrap_or(i64::MAX);
let rows = self.engine.query(
"SELECT object_key, store_path, marked_seq, grace_until_seq FROM gc_tombstones \
WHERE grace_until_seq <= ?1 ORDER BY object_key, store_path",
&[SqlValue::Int(now)],
)?;
rows.into_iter()
.map(|row| {
let [key, path, marked, grace] = row.as_slice() else {
return Err(StoreError::Corruption("tombstone row shape".into()));
};
let (SqlValue::Text(key), SqlValue::Text(path)) = (key, path) else {
return Err(StoreError::Corruption("tombstone column shape".into()));
};
Ok(GcTombstoneRow {
object_key: key.clone(),
store_path: path.clone(),
marked_seq: expect_u64(marked, "marked_seq")?,
grace_until_seq: expect_u64(grace, "grace_until_seq")?,
})
})
.collect()
}
fn remove_gc_tombstone(
&mut self,
object_key: &str,
store_path: &str,
) -> Result<bool, StoreError> {
let object_key = object_key.to_owned();
let store_path = store_path.to_owned();
self.in_txn(move |engine| {
let removed = engine.execute(
"DELETE FROM gc_tombstones WHERE object_key = ?1 AND store_path = ?2",
&[SqlValue::Text(object_key), SqlValue::Text(store_path)],
)?;
Ok(removed > 0)
})
}
fn list_publications(&mut self) -> Result<Vec<(String, String)>, StoreError> {
let rows = self.engine.query(
"SELECT action_key, pin_hex FROM action_publications ORDER BY action_key",
&[],
)?;
rows.into_iter()
.map(|row| match row.as_slice() {
[SqlValue::Text(action), SqlValue::Text(pin)] => Ok((action.clone(), pin.clone())),
_ => Err(StoreError::Corruption("publication list shape".into())),
})
.collect()
}
fn publications_missing_their_generation(&mut self) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT p.action_key FROM action_publications p \
LEFT JOIN action_generations g ON g.id_hex = p.winner_generation_hex \
WHERE g.id_hex IS NULL ORDER BY p.action_key",
&[],
)?;
rows.into_iter()
.map(|row| match row.as_slice() {
[SqlValue::Text(action)] => Ok(action.clone()),
_ => Err(StoreError::Corruption(
"orphaned-publication list shape".into(),
)),
})
.collect()
}
fn pin_row(&mut self, id: u128) -> Result<Option<PinRow>, StoreError> {
let rows = self.engine.query(
"SELECT root_key, owner, class, expires_at_seq, released, renewal_seq \
FROM pins WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [root_key, owner, class, expires, released, renewal] = row.as_slice() else {
return Err(StoreError::Corruption("pin row shape".into()));
};
let (SqlValue::Text(root_key), SqlValue::Text(owner), SqlValue::Text(class)) =
(root_key, owner, class)
else {
return Err(StoreError::Corruption("pin column shape".into()));
};
let expires_at_seq = match expires {
SqlValue::Null => None,
other => Some(expect_u64(other, "expires_at_seq")?),
};
Ok(Some(PinRow {
root_key: root_key.clone(),
owner: owner.clone(),
class: class.clone(),
expires_at_seq,
released: expect_u64(released, "released")? != 0,
renewal_seq: expect_u64(renewal, "renewal_seq")?,
}))
}
fn pin_released_by_hex(&mut self, pin_hex: &str) -> Result<Option<bool>, StoreError> {
let rows = self.engine.query(
"SELECT released FROM pins WHERE id_hex = ?1",
&[SqlValue::Text(pin_hex.to_owned())],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(v) => Ok(Some(expect_u64(v, "released")? != 0)),
}
}
fn has_serving_state_key(&mut self, action_key: &str) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT action_key FROM action_serving_states WHERE action_key = ?1",
&[SqlValue::Text(action_key.to_owned())],
)?;
Ok(!rows.is_empty())
}
fn has_evidence_key(&mut self, action_key: &str) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT action_key FROM action_evidence_index WHERE action_key = ?1 LIMIT 1",
&[SqlValue::Text(action_key.to_owned())],
)?;
Ok(!rows.is_empty())
}
fn authority_count(&mut self) -> Result<u64, StoreError> {
let rows = self
.engine
.query("SELECT COUNT(*) FROM coordinator_authorities", &[])?;
match rows.first().and_then(|r| r.first()) {
Some(v) => expect_u64(v, "authority count"),
None => Err(StoreError::Corruption("count shape".into())),
}
}
fn generation_count(&mut self) -> Result<u64, StoreError> {
let rows = self
.engine
.query("SELECT COUNT(*) FROM action_generations", &[])?;
match rows.first().and_then(|r| r.first()) {
Some(v) => expect_u64(v, "generation count"),
None => Err(StoreError::Corruption("count shape".into())),
}
}
fn has_generation_high_water(&mut self) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT kind FROM generation_high_water WHERE kind = 'action-generation'",
&[],
)?;
Ok(!rows.is_empty())
}
fn record_eviction_tombstone(
&mut self,
action: &TypedDigest,
semantic: &TypedDigest,
observable: &TypedDigest,
evicted_seq: u64,
) -> Result<(), StoreError> {
self.intern(semantic.domain);
self.intern(observable.domain);
let action = digest_key(action);
let [s_algo, s_domain, s_bytes] = SqlMetadataStore::<E>::digest_params(semantic);
let [o_algo, o_domain, o_bytes] = SqlMetadataStore::<E>::digest_params(observable);
let seq = i64::try_from(evicted_seq)
.map_err(|_| StoreError::Corruption("evicted_seq out of range".into()))?;
self.in_txn(move |engine| {
engine.execute(
"INSERT OR REPLACE INTO eviction_tombstones \
(action_key, semantic_algo, semantic_domain, semantic_bytes, \
observable_algo, observable_domain, observable_bytes, evicted_seq) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
&[
SqlValue::Text(action),
s_algo,
s_domain,
s_bytes,
o_algo,
o_domain,
o_bytes,
SqlValue::Int(seq),
],
)?;
Ok(())
})
}
fn eviction_tombstone(
&mut self,
action: &TypedDigest,
) -> Result<Option<(TypedDigest, TypedDigest)>, StoreError> {
let rows = self.engine.query(
"SELECT semantic_algo, semantic_domain, semantic_bytes, \
observable_algo, observable_domain, observable_bytes \
FROM eviction_tombstones WHERE action_key = ?1",
&[SqlValue::Text(digest_key(action))],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [s_algo, s_domain, s_bytes, o_algo, o_domain, o_bytes] = row.as_slice() else {
return Err(StoreError::Corruption("eviction tombstone shape".into()));
};
let semantic = self.restore_digest(s_algo, s_domain, s_bytes)?;
let observable = self.restore_digest(o_algo, o_domain, o_bytes)?;
Ok(Some((semantic, observable)))
}
fn consume_eviction_tombstone(&mut self, action: &TypedDigest) -> Result<bool, StoreError> {
let action = digest_key(action);
self.in_txn(move |engine| {
let removed = engine.execute(
"DELETE FROM eviction_tombstones WHERE action_key = ?1",
&[SqlValue::Text(action)],
)?;
Ok(removed > 0)
})
}
fn record_operator_reset(&mut self, generation: u64, seq: u64) -> Result<(), StoreError> {
let generation_int = i64::try_from(generation)
.map_err(|_| StoreError::Corruption("reset generation out of range".into()))?;
let seq =
i64::try_from(seq).map_err(|_| StoreError::Corruption("seq out of range".into()))?;
self.in_txn(move |engine| {
let rows = engine.query("SELECT MAX(generation) FROM operator_resets", &[])?;
if let Some(SqlValue::Int(highest)) = rows.first().and_then(|r| r.first())
&& generation_int <= *highest
{
return Err(StoreError::StaleOperatorReset);
}
engine.execute(
"INSERT INTO operator_resets (generation, applied_seq) VALUES (?1, ?2)",
&[SqlValue::Int(generation_int), SqlValue::Int(seq)],
)?;
Ok(())
})
}
fn highest_operator_reset(&mut self) -> Result<Option<u64>, StoreError> {
let rows = self
.engine
.query("SELECT MAX(generation) FROM operator_resets", &[])?;
match rows.first().and_then(|r| r.first()) {
Some(SqlValue::Int(v)) => {
Ok(Some(u64::try_from(*v).map_err(|_| {
StoreError::Corruption("negative reset generation".into())
})?))
}
_ => Ok(None),
}
}
fn serving_disposition_key(&mut self, action_key: &str) -> Result<Option<String>, StoreError> {
let rows = self.engine.query(
"SELECT disposition FROM action_serving_states WHERE action_key = ?1",
&[SqlValue::Text(action_key.to_owned())],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(SqlValue::Text(d)) => Ok(Some(d.clone())),
Some(_) => Err(StoreError::Corruption("disposition shape".into())),
}
}
fn set_serving_disposition_key(
&mut self,
action_key: &str,
disposition: &str,
) -> Result<(), StoreError> {
self.in_txn(|engine| Self::set_serving_disposition_row(engine, action_key, disposition))
}
fn set_serving_disposition_for_attempt(
&mut self,
authority: &AttemptAuthority,
disposition: &str,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<(), StoreError> {
self.in_txn(|engine| {
Self::require_live_attempt(engine, authority, own_monotonic_now_ms)?;
Self::set_serving_disposition_row(
engine,
&digest_key(&authority.action_key),
disposition,
)?;
Self::require_live_attempt(engine, authority, own_monotonic_now_ms)
})
}
fn apply_operator_reset_to_peer(
&mut self,
authority: &TypedDigest,
peer_id: &str,
reset_generation: u64,
credential_generation: u64,
term: u64,
incarnation: u128,
) -> Result<(), StoreError> {
let authority = authority.clone();
let peer = peer_id.to_owned();
let reset_int = i64::try_from(reset_generation)
.map_err(|_| StoreError::Corruption("reset generation out of range".into()))?;
let term_int = to_seq(term, "term")?;
let generation_int = to_seq(credential_generation, "credential_generation")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let ledger = engine.query("SELECT MAX(generation) FROM operator_resets", &[])?;
if let Some(SqlValue::Int(highest)) = ledger.first().and_then(|r| r.first())
&& reset_int <= *highest
{
return Err(StoreError::StaleOperatorReset);
}
engine.execute(
"INSERT INTO operator_resets (generation, applied_seq) VALUES (?1, ?2)",
&[SqlValue::Int(reset_int), SqlValue::Int(reset_int)],
)?;
engine.execute(
"INSERT OR REPLACE INTO peer_authority_high_water \
(peer_id, term, observed_seq, incarnation, credential_generation) \
VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(peer),
SqlValue::Int(term_int),
SqlValue::Int(reset_int),
SqlValue::Blob(incarnation.to_be_bytes().to_vec()),
SqlValue::Int(generation_int),
],
)?;
Ok(())
})
}
fn record_peer_authority_high_water(
&mut self,
authority: &TypedDigest,
peer_id: &str,
credential_generation: u64,
term: u64,
incarnation: u128,
observed_seq: u64,
) -> Result<(), StoreError> {
let authority = authority.clone();
let peer = peer_id.to_owned();
let term_int = to_seq(term, "term")?;
let generation_int = to_seq(credential_generation, "credential_generation")?;
let observed = to_seq(observed_seq, "observed_seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let rows = engine.query(
"SELECT term, incarnation, credential_generation FROM peer_authority_high_water WHERE peer_id = ?1",
&[SqlValue::Text(peer.clone())],
)?;
if let Some(row) = rows.first() {
let stored = expect_u64(
row.first()
.ok_or_else(|| StoreError::Corruption("peer high-water shape".into()))?,
"peer term",
)?;
let stored_generation = match row.get(2) {
Some(SqlValue::Null) => return Err(StoreError::StalePeerAuthority),
Some(value) => expect_u64(value, "peer credential generation")?,
None => return Err(StoreError::Corruption("peer high-water shape".into())),
};
if (credential_generation, term) < (stored_generation, stored) {
return Err(StoreError::StalePeerAuthority);
}
if (credential_generation, term) == (stored_generation, stored) {
if !matches!(row.get(1), Some(SqlValue::Blob(bytes)) if bytes.as_slice() == incarnation.to_be_bytes()) {
return Err(StoreError::StalePeerAuthority);
}
return Ok(());
}
}
engine.execute(
"INSERT OR REPLACE INTO peer_authority_high_water \
(peer_id, term, observed_seq, incarnation, credential_generation) VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(peer),
SqlValue::Int(term_int),
SqlValue::Int(observed),
SqlValue::Blob(incarnation.to_be_bytes().to_vec()),
SqlValue::Int(generation_int),
],
)?;
Ok(())
})
}
fn peer_authority_high_water(
&mut self,
peer_id: &str,
) -> Result<Option<(u64, u64, u64)>, StoreError> {
let rows = self.engine.query(
"SELECT credential_generation, term, observed_seq FROM peer_authority_high_water WHERE peer_id = ?1",
&[SqlValue::Text(peer_id.to_owned())],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [generation, term, observed] = row.as_slice() else {
return Err(StoreError::Corruption("peer high-water shape".into()));
};
if matches!(generation, SqlValue::Null) {
return Err(StoreError::StalePeerAuthority);
}
Ok(Some((
expect_u64(generation, "peer credential generation")?,
expect_u64(term, "peer term")?,
expect_u64(observed, "observed_seq")?,
)))
}
fn admit_worker_session(
&mut self,
authority: &TypedDigest,
offer: &WorkerSessionOffer,
started_seq: u64,
) -> Result<WorkerAdmission, StoreError> {
let authority = authority.clone();
let offer = offer.clone();
let worker = offer.worker_peer_id.0.clone();
let started = to_seq(started_seq, "started_seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let rows = engine.query(
"SELECT highest_boot_generation, incarnation, active, \
operator_reenrollment_generation, clone_ambiguous \
FROM worker_incarnation_fences \
WHERE worker = ?1",
&[SqlValue::Text(worker.clone())],
)?;
let (admission, fence, revoke_prior_leases) = match rows.as_slice() {
[] => (
WorkerAdmission::AdmitNewGeneration,
WorkerIncarnationFenceRecord {
worker_peer_id: offer.worker_peer_id.clone(),
highest_boot_generation: offer.boot_generation,
active_incarnation: Some(offer.incarnation),
clone_ambiguous: false,
operator_reenrollment_generation: offer
.reenrollment_proof
.unwrap_or_default(),
},
false,
),
[row] => {
let prior_incarnation =
WorkerIncarnationId(expect_u128(&row[1], "worker incarnation")?);
let mut fence = decode_worker_incarnation_fence(&worker, row)?;
let admission = fence.evaluate(&offer);
let revoke_prior_leases = match admission {
WorkerAdmission::AdmitNewGeneration => {
fence.highest_boot_generation = offer.boot_generation;
fence.active_incarnation = Some(offer.incarnation);
fence.clone_ambiguous = false;
true
}
WorkerAdmission::AdmitReconnect => false,
WorkerAdmission::AdmitResume => {
fence.active_incarnation = Some(offer.incarnation);
prior_incarnation != offer.incarnation
}
WorkerAdmission::AdmitViaReenrollment => {
let proof = offer.reenrollment_proof.ok_or_else(|| {
StoreError::Corruption(
"reenrollment admission without proof".into(),
)
})?;
fence.highest_boot_generation = offer.boot_generation;
fence.active_incarnation = Some(offer.incarnation);
fence.clone_ambiguous = false;
fence.operator_reenrollment_generation = proof;
true
}
WorkerAdmission::RejectCloneAmbiguity => {
engine.execute(
"UPDATE worker_incarnation_fences \
SET clone_ambiguous = 1 WHERE worker = ?1",
&[SqlValue::Text(worker.clone())],
)?;
SqlMetadataStore::<E>::revoke_worker_leases(engine, &worker)?;
return Ok(admission);
}
WorkerAdmission::RejectStaleBootGeneration
| WorkerAdmission::RejectIdentityMismatch => return Ok(admission),
};
(admission, fence, revoke_prior_leases)
}
_ => {
return Err(StoreError::Corruption("duplicate worker fence rows".into()));
}
};
if revoke_prior_leases {
SqlMetadataStore::<E>::revoke_worker_leases(engine, &worker)?;
}
let active_incarnation = fence.active_incarnation.ok_or_else(|| {
StoreError::Corruption("admitted worker fence is inactive".into())
})?;
let existing = engine.query(
"SELECT incarnation, ended_seq FROM worker_sessions \
WHERE worker = ?1 AND started_seq = ?2",
&[SqlValue::Text(worker.clone()), SqlValue::Int(started)],
)?;
let insert_session = match existing.as_slice() {
[] => true,
[row] => {
let [stored_incarnation, ended_seq] = row.as_slice() else {
return Err(StoreError::Corruption("worker session shape".into()));
};
if expect_u128(stored_incarnation, "session incarnation")?
== active_incarnation.0
&& matches!(ended_seq, SqlValue::Null)
{
false
} else {
return Err(StoreError::AppendConflict("worker_sessions".into()));
}
}
_ => {
return Err(StoreError::Corruption(
"duplicate worker session rows".into(),
));
}
};
engine.execute(
"INSERT OR REPLACE INTO worker_incarnation_fences \
(worker, incarnation, highest_boot_generation, active, \
operator_reenrollment_generation, clone_ambiguous) \
VALUES (?1, ?2, ?3, 1, ?4, ?5)",
&[
SqlValue::Text(worker.clone()),
SqlValue::Blob(u128_blob(active_incarnation.0)),
SqlValue::Blob(u64_blob(fence.highest_boot_generation.0)),
SqlValue::Blob(u64_blob(fence.operator_reenrollment_generation)),
SqlValue::Int(i64::from(fence.clone_ambiguous)),
],
)?;
if insert_session {
engine.execute(
"INSERT INTO worker_sessions \
(worker, incarnation, started_seq, ended_seq) VALUES (?1, ?2, ?3, NULL)",
&[
SqlValue::Text(worker),
SqlValue::Blob(u128_blob(active_incarnation.0)),
SqlValue::Int(started),
],
)?;
}
Ok(admission)
})
}
fn release_worker_session(
&mut self,
authority: &TypedDigest,
worker: &PeerId,
incarnation: WorkerIncarnationId,
started_seq: u64,
ended_seq: u64,
) -> Result<bool, StoreError> {
let authority = authority.clone();
let worker = worker.0.clone();
let started = to_seq(started_seq, "started_seq")?;
let ended = to_seq(ended_seq, "ended_seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let rows = engine.query(
"SELECT highest_boot_generation, incarnation, active, \
operator_reenrollment_generation, clone_ambiguous \
FROM worker_incarnation_fences \
WHERE worker = ?1",
&[SqlValue::Text(worker.clone())],
)?;
let fence = match rows.as_slice() {
[] => return Ok(false),
[row] => decode_worker_incarnation_fence(&worker, row)?,
_ => {
return Err(StoreError::Corruption("duplicate worker fence rows".into()));
}
};
if fence.active_incarnation != Some(incarnation) {
return Ok(false);
}
let ended_session = engine.execute(
"UPDATE worker_sessions SET ended_seq = ?1 \
WHERE worker = ?2 AND incarnation = ?3 AND started_seq = ?4 \
AND ended_seq IS NULL",
&[
SqlValue::Int(ended),
SqlValue::Text(worker.clone()),
SqlValue::Blob(u128_blob(incarnation.0)),
SqlValue::Int(started),
],
)?;
if ended_session == 0 {
return Ok(false);
}
let remaining = engine.query(
"SELECT COUNT(*) FROM worker_sessions \
WHERE worker = ?1 AND incarnation = ?2 AND ended_seq IS NULL",
&[
SqlValue::Text(worker.clone()),
SqlValue::Blob(u128_blob(incarnation.0)),
],
)?;
let remaining = match remaining.as_slice() {
[row] => match row.as_slice() {
[value] => expect_u64(value, "open worker session count")?,
_ => {
return Err(StoreError::Corruption(
"open worker session count shape".into(),
));
}
},
_ => {
return Err(StoreError::Corruption(
"open worker session count shape".into(),
));
}
};
if remaining > 0 {
return Ok(true);
}
let cleared = engine.execute(
"UPDATE worker_incarnation_fences SET active = 0 \
WHERE worker = ?1 AND incarnation = ?2 AND active = 1",
&[
SqlValue::Text(worker),
SqlValue::Blob(u128_blob(incarnation.0)),
],
)?;
if cleared != 1 {
return Err(StoreError::Corruption(
"worker session ended without clearing its fence".into(),
));
}
Ok(true)
})
}
fn worker_incarnation_fence(
&mut self,
worker: &PeerId,
) -> Result<Option<WorkerIncarnationFenceRecord>, StoreError> {
let rows = self.engine.query(
"SELECT highest_boot_generation, incarnation, active, \
operator_reenrollment_generation, clone_ambiguous \
FROM worker_incarnation_fences \
WHERE worker = ?1",
&[SqlValue::Text(worker.0.clone())],
)?;
match rows.as_slice() {
[] => Ok(None),
[row] => decode_worker_incarnation_fence(&worker.0, row).map(Some),
_ => Err(StoreError::Corruption("duplicate worker fence rows".into())),
}
}
fn advance_edge_fence(
&mut self,
authority: &TypedDigest,
edge_id: &str,
incarnation: u128,
) -> Result<(), StoreError> {
let authority = authority.clone();
let edge = edge_id.to_owned();
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let rows = engine.query(
"SELECT incarnation FROM edge_incarnation_fences WHERE edge_id = ?1",
&[SqlValue::Text(edge.clone())],
)?;
if let Some(row) = rows.first() {
let stored = expect_u128(
row.first()
.ok_or_else(|| StoreError::Corruption("edge fence shape".into()))?,
"edge fence",
)?;
if incarnation < stored {
return Err(StoreError::StaleEdgeIncarnation);
}
if incarnation == stored {
return Ok(()); }
}
engine.execute(
"INSERT OR REPLACE INTO edge_incarnation_fences (edge_id, incarnation) \
VALUES (?1, ?2)",
&[SqlValue::Text(edge), SqlValue::Blob(u128_blob(incarnation))],
)?;
Ok(())
})
}
fn edge_fence(&mut self, edge_id: &str) -> Result<Option<u128>, StoreError> {
let rows = self.engine.query(
"SELECT incarnation FROM edge_incarnation_fences WHERE edge_id = ?1",
&[SqlValue::Text(edge_id.to_owned())],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(v) => Ok(Some(expect_u128(v, "edge fence")?)),
}
}
fn begin_edge_handoff(
&mut self,
authority: &TypedDigest,
edge_id: &str,
active_incarnation: u128,
predecessor_incarnation: u128,
begun_seq: u64,
) -> Result<(), StoreError> {
let authority = authority.clone();
let edge = edge_id.to_owned();
let begun = to_seq(begun_seq, "begun_seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let rows = engine.query(
"SELECT active_incarnation, predecessor_incarnation, resolved \
FROM edge_handoffs WHERE edge_id = ?1",
&[SqlValue::Text(edge.clone())],
)?;
if let Some(row) = rows.first() {
let [active, predecessor, resolved] = row.as_slice() else {
return Err(StoreError::Corruption("edge handoff shape".into()));
};
if expect_u64(resolved, "resolved")? == 0 {
if expect_u128(active, "active incarnation")? == active_incarnation
&& expect_u128(predecessor, "predecessor incarnation")?
== predecessor_incarnation
{
return Ok(()); }
return Err(StoreError::EdgeHandoffActive);
}
}
let fence = engine.query(
"SELECT incarnation FROM edge_incarnation_fences WHERE edge_id = ?1",
&[SqlValue::Text(edge.clone())],
)?;
if let Some(v) = fence.first().and_then(|r| r.first())
&& expect_u128(v, "edge fence")? != predecessor_incarnation
{
return Err(StoreError::EdgeHandoffPredecessorMismatch);
}
if active_incarnation <= predecessor_incarnation {
return Err(StoreError::StaleEdgeIncarnation);
}
engine.execute(
"INSERT OR REPLACE INTO edge_handoffs \
(edge_id, active_incarnation, predecessor_incarnation, begun_seq, resolved) \
VALUES (?1, ?2, ?3, ?4, 0)",
&[
SqlValue::Text(edge),
SqlValue::Blob(u128_blob(active_incarnation)),
SqlValue::Blob(u128_blob(predecessor_incarnation)),
SqlValue::Int(begun),
],
)?;
Ok(())
})
}
fn resolve_edge_handoff(
&mut self,
authority: &TypedDigest,
edge_id: &str,
active_incarnation: u128,
) -> Result<(), StoreError> {
let authority = authority.clone();
let edge = edge_id.to_owned();
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let rows = engine.query(
"SELECT active_incarnation, resolved FROM edge_handoffs WHERE edge_id = ?1",
&[SqlValue::Text(edge.clone())],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownEdgeHandoff);
};
let [active, resolved] = row.as_slice() else {
return Err(StoreError::Corruption("edge handoff shape".into()));
};
if expect_u64(resolved, "resolved")? != 0
|| expect_u128(active, "active incarnation")? != active_incarnation
{
return Err(StoreError::UnknownEdgeHandoff);
}
engine.execute(
"UPDATE edge_handoffs SET resolved = 1 WHERE edge_id = ?1",
&[SqlValue::Text(edge.clone())],
)?;
engine.execute(
"INSERT OR REPLACE INTO edge_incarnation_fences (edge_id, incarnation) \
VALUES (?1, ?2)",
&[
SqlValue::Text(edge),
SqlValue::Blob(u128_blob(active_incarnation)),
],
)?;
Ok(())
})
}
fn active_edge_handoff(&mut self, edge_id: &str) -> Result<Option<EdgeHandoffRow>, StoreError> {
let rows = self.engine.query(
"SELECT active_incarnation, predecessor_incarnation, begun_seq \
FROM edge_handoffs WHERE edge_id = ?1 AND resolved = 0",
&[SqlValue::Text(edge_id.to_owned())],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [active, predecessor, begun] = row.as_slice() else {
return Err(StoreError::Corruption("edge handoff shape".into()));
};
Ok(Some(EdgeHandoffRow {
active_incarnation: expect_u128(active, "active incarnation")?,
predecessor_incarnation: expect_u128(predecessor, "predecessor incarnation")?,
begun_seq: expect_u64(begun, "begun_seq")?,
}))
}
fn append_trust_evaluation(
&mut self,
authority: &TypedDigest,
action: &TypedDigest,
row: &TrustEvaluationRow,
) -> Result<(), StoreError> {
let authority = authority.clone();
let action = digest_key(action);
let row = row.clone();
let evaluated = to_seq(row.evaluated_seq, "evaluated_seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let stored = engine.query(
"SELECT MAX(version) FROM action_trust_evaluations WHERE action_key = ?1",
&[SqlValue::Text(action.clone())],
)?;
if let Some(SqlValue::Int(max)) = stored.first().and_then(|r| r.first())
&& u64::from(row.version) <= u64::try_from(*max).unwrap_or(0)
{
return Err(StoreError::NonMonotonicTrustEvaluation);
}
engine.execute(
"INSERT INTO action_trust_evaluations \
(action_key, version, state, reason, evaluated_seq) \
VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(action),
SqlValue::Int(i64::from(row.version)),
SqlValue::Text(row.state),
SqlValue::Text(row.reason),
SqlValue::Int(evaluated),
],
)?;
Ok(())
})
}
fn latest_trust_evaluation(
&mut self,
action: &TypedDigest,
) -> Result<Option<TrustEvaluationRow>, StoreError> {
let rows = self.engine.query(
"SELECT version, state, reason, evaluated_seq FROM action_trust_evaluations \
WHERE action_key = ?1 ORDER BY version DESC LIMIT 1",
&[SqlValue::Text(digest_key(action))],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [version, state, reason, evaluated] = row.as_slice() else {
return Err(StoreError::Corruption("trust evaluation shape".into()));
};
Ok(Some(TrustEvaluationRow {
version: u32::try_from(expect_u64(version, "version")?)
.map_err(|_| StoreError::Corruption("version out of range".into()))?,
state: expect_text(state, "state")?,
reason: expect_text(reason, "reason")?,
evaluated_seq: expect_u64(evaluated, "evaluated_seq")?,
}))
}
fn create_operation(
&mut self,
authority: &TypedDigest,
id: u128,
kind: &str,
state: &str,
seq: u64,
) -> Result<(), StoreError> {
let authority = authority.clone();
let kind = kind.to_owned();
let state = state.to_owned();
let seq = to_seq(seq, "seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let existing = engine.query(
"SELECT id_hex FROM operations WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
if !existing.is_empty() {
return Err(StoreError::DuplicateOperation);
}
engine.execute(
"INSERT INTO operations (id_hex, id, kind, state, updated_seq) \
VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(kind),
SqlValue::Text(state),
SqlValue::Int(seq),
],
)?;
Ok(())
})
}
fn update_operation_state(
&mut self,
id: u128,
state: &str,
seq: u64,
) -> Result<(), StoreError> {
let state = state.to_owned();
let seq = to_seq(seq, "seq")?;
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE operations SET state = ?1, updated_seq = ?2 WHERE id_hex = ?3",
&[
SqlValue::Text(state),
SqlValue::Int(seq),
SqlValue::Text(u128_hex(id)),
],
)?;
if changed == 0 {
return Err(StoreError::UnknownOperation);
}
Ok(())
})
}
fn operation_state(&mut self, id: u128) -> Result<Option<String>, StoreError> {
let rows = self.engine.query(
"SELECT state FROM operations WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(v) => Ok(Some(expect_text(v, "operation state")?)),
}
}
fn register_edge_subscriber(
&mut self,
edge_id: &str,
subscriber: &str,
registered_seq: u64,
) -> Result<(), StoreError> {
let edge = edge_id.to_owned();
let subscriber = subscriber.to_owned();
let seq = to_seq(registered_seq, "registered_seq")?;
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO edge_subscribers (edge_id, subscriber, registered_seq) \
VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(edge),
SqlValue::Text(subscriber),
SqlValue::Int(seq),
],
)?;
Ok(())
})
}
fn remove_edge_subscriber(
&mut self,
edge_id: &str,
subscriber: &str,
) -> Result<bool, StoreError> {
let edge = edge_id.to_owned();
let subscriber = subscriber.to_owned();
self.in_txn(move |engine| {
let changed = engine.execute(
"DELETE FROM edge_subscribers WHERE edge_id = ?1 AND subscriber = ?2",
&[SqlValue::Text(edge), SqlValue::Text(subscriber)],
)?;
Ok(changed > 0)
})
}
fn list_edge_subscribers(&mut self, edge_id: &str) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT subscriber FROM edge_subscribers WHERE edge_id = ?1 ORDER BY subscriber",
&[SqlValue::Text(edge_id.to_owned())],
)?;
rows.iter()
.map(|row| {
expect_text(
row.first()
.ok_or_else(|| StoreError::Corruption("subscriber shape".into()))?,
"subscriber",
)
})
.collect()
}
fn record_manifest(
&mut self,
manifest: &TypedDigest,
kind: &str,
entry_count: u64,
) -> Result<(), StoreError> {
self.intern(manifest.domain);
let key = digest_key(manifest);
let [algo, domain, bytes] = SqlMetadataStore::<E>::digest_params(manifest);
let kind = kind.to_owned();
let count = to_seq(entry_count, "entry_count")?;
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT kind, entry_count FROM manifests WHERE key = ?1",
&[SqlValue::Text(key.clone())],
)?;
if let Some(row) = existing.first() {
let [stored_kind, stored_count] = row.as_slice() else {
return Err(StoreError::Corruption("manifest row shape".into()));
};
if expect_text(stored_kind, "manifest kind")? == kind
&& expect_u64(stored_count, "entry_count")? == entry_count
{
return Ok(()); }
return Err(StoreError::ManifestDivergence);
}
engine.execute(
"INSERT INTO manifests (key, algo, domain, bytes, kind, entry_count) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
&[
SqlValue::Text(key),
algo,
domain,
bytes,
SqlValue::Text(kind),
SqlValue::Int(count),
],
)?;
Ok(())
})
}
fn manifest_meta(
&mut self,
manifest: &TypedDigest,
) -> Result<Option<(String, u64)>, StoreError> {
let rows = self.engine.query(
"SELECT kind, entry_count FROM manifests WHERE key = ?1",
&[SqlValue::Text(digest_key(manifest))],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [kind, count] = row.as_slice() else {
return Err(StoreError::Corruption("manifest row shape".into()));
};
Ok(Some((
expect_text(kind, "manifest kind")?,
expect_u64(count, "entry_count")?,
)))
}
fn record_worker_session(
&mut self,
worker: &str,
incarnation: u128,
started_seq: u64,
) -> Result<(), StoreError> {
let worker = worker.to_owned();
let started = to_seq(started_seq, "started_seq")?;
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT incarnation FROM worker_sessions WHERE worker = ?1 AND started_seq = ?2",
&[SqlValue::Text(worker.clone()), SqlValue::Int(started)],
)?;
if let Some(row) = existing.first() {
let stored = expect_u128(
row.first()
.ok_or_else(|| StoreError::Corruption("worker session shape".into()))?,
"session incarnation",
)?;
if stored == incarnation {
return Ok(()); }
return Err(StoreError::AppendConflict("worker_sessions".into()));
}
engine.execute(
"INSERT INTO worker_sessions (worker, incarnation, started_seq, ended_seq) \
VALUES (?1, ?2, ?3, NULL)",
&[
SqlValue::Text(worker),
SqlValue::Blob(u128_blob(incarnation)),
SqlValue::Int(started),
],
)?;
Ok(())
})
}
fn end_worker_session(
&mut self,
worker: &str,
started_seq: u64,
ended_seq: u64,
) -> Result<bool, StoreError> {
let worker = worker.to_owned();
let started = to_seq(started_seq, "started_seq")?;
let ended = to_seq(ended_seq, "ended_seq")?;
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE worker_sessions SET ended_seq = ?1 \
WHERE worker = ?2 AND started_seq = ?3 AND ended_seq IS NULL",
&[
SqlValue::Int(ended),
SqlValue::Text(worker),
SqlValue::Int(started),
],
)?;
Ok(changed > 0)
})
}
fn record_worker_capability(
&mut self,
worker: &str,
capability: &str,
) -> Result<(), StoreError> {
let worker = worker.to_owned();
let capability = capability.to_owned();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO worker_capabilities (worker, capability) \
VALUES (?1, ?2)",
&[SqlValue::Text(worker), SqlValue::Text(capability)],
)?;
Ok(())
})
}
fn list_worker_capabilities(&mut self, worker: &str) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT capability FROM worker_capabilities WHERE worker = ?1 ORDER BY capability",
&[SqlValue::Text(worker.to_owned())],
)?;
rows.iter()
.map(|row| {
expect_text(
row.first()
.ok_or_else(|| StoreError::Corruption("capability shape".into()))?,
"capability",
)
})
.collect()
}
fn record_worker_health_sample(
&mut self,
worker: &str,
seq: u64,
healthy: bool,
detail: &str,
) -> Result<(), StoreError> {
let worker = worker.to_owned();
let seq = to_seq(seq, "seq")?;
let detail = detail.to_owned();
let healthy_int = i64::from(healthy);
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT healthy, detail FROM worker_health_samples \
WHERE worker = ?1 AND seq = ?2",
&[SqlValue::Text(worker.clone()), SqlValue::Int(seq)],
)?;
if let Some(row) = existing.first() {
let [stored_healthy, stored_detail] = row.as_slice() else {
return Err(StoreError::Corruption("health sample shape".into()));
};
if expect_u64(stored_healthy, "healthy")? == u64::from(healthy)
&& expect_text(stored_detail, "detail")? == detail
{
return Ok(()); }
return Err(StoreError::AppendConflict("worker_health_samples".into()));
}
engine.execute(
"INSERT INTO worker_health_samples (worker, seq, healthy, detail) \
VALUES (?1, ?2, ?3, ?4)",
&[
SqlValue::Text(worker),
SqlValue::Int(seq),
SqlValue::Int(healthy_int),
SqlValue::Text(detail),
],
)?;
Ok(())
})
}
fn record_decision_receipt(
&mut self,
kind: &str,
subject: &str,
seq: u64,
decision: &str,
reason: &str,
) -> Result<(), StoreError> {
let kind = kind.to_owned();
let subject = subject.to_owned();
let seq = to_seq(seq, "seq")?;
let decision = decision.to_owned();
let reason = reason.to_owned();
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT decision, reason FROM decision_receipts \
WHERE kind = ?1 AND subject = ?2 AND seq = ?3",
&[
SqlValue::Text(kind.clone()),
SqlValue::Text(subject.clone()),
SqlValue::Int(seq),
],
)?;
if let Some(row) = existing.first() {
let [stored_decision, stored_reason] = row.as_slice() else {
return Err(StoreError::Corruption("decision receipt shape".into()));
};
if expect_text(stored_decision, "decision")? == decision
&& expect_text(stored_reason, "reason")? == reason
{
return Ok(()); }
return Err(StoreError::AppendConflict("decision_receipts".into()));
}
engine.execute(
"INSERT INTO decision_receipts (kind, subject, seq, decision, reason) \
VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(kind),
SqlValue::Text(subject),
SqlValue::Int(seq),
SqlValue::Text(decision),
SqlValue::Text(reason),
],
)?;
Ok(())
})
}
fn add_provenance_edge(
&mut self,
from: &TypedDigest,
to: &TypedDigest,
kind: &str,
) -> Result<(), StoreError> {
let from = digest_key(from);
let to = digest_key(to);
let kind = kind.to_owned();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO provenance_edges (from_key, to_key, kind) \
VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(from),
SqlValue::Text(to),
SqlValue::Text(kind),
],
)?;
Ok(())
})
}
fn record_determinism_audit(
&mut self,
action: &TypedDigest,
attempt: u128,
seq: u64,
verdict: &str,
) -> Result<(), StoreError> {
let action = digest_key(action);
let seq = to_seq(seq, "seq")?;
let verdict = verdict.to_owned();
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT verdict FROM determinism_audits \
WHERE action_key = ?1 AND attempt_hex = ?2 AND seq = ?3",
&[
SqlValue::Text(action.clone()),
SqlValue::Text(u128_hex(attempt)),
SqlValue::Int(seq),
],
)?;
if let Some(row) = existing.first() {
let stored = expect_text(
row.first()
.ok_or_else(|| StoreError::Corruption("determinism audit shape".into()))?,
"verdict",
)?;
if stored == verdict {
return Ok(()); }
return Err(StoreError::AppendConflict("determinism_audits".into()));
}
engine.execute(
"INSERT INTO determinism_audits (action_key, attempt_hex, seq, verdict) \
VALUES (?1, ?2, ?3, ?4)",
&[
SqlValue::Text(action),
SqlValue::Text(u128_hex(attempt)),
SqlValue::Int(seq),
SqlValue::Text(verdict),
],
)?;
Ok(())
})
}
fn create_materialization(
&mut self,
id: u128,
root: &TypedDigest,
dest_path: &str,
state: &str,
seq: u64,
) -> Result<(), StoreError> {
let root = digest_key(root);
let dest = dest_path.to_owned();
let state = state.to_owned();
let seq = to_seq(seq, "seq")?;
self.in_txn(move |engine| {
let existing = engine.query(
"SELECT root_key, dest_path FROM materialization_records WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
if let Some(row) = existing.first() {
let [stored_root, stored_dest] = row.as_slice() else {
return Err(StoreError::Corruption("materialization shape".into()));
};
if expect_text(stored_root, "root_key")? == root
&& expect_text(stored_dest, "dest_path")? == dest
{
return Ok(()); }
return Err(StoreError::AppendConflict("materialization_records".into()));
}
engine.execute(
"INSERT INTO materialization_records \
(id_hex, id, root_key, dest_path, state, updated_seq) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
&[
SqlValue::Text(u128_hex(id)),
SqlValue::Blob(u128_blob(id)),
SqlValue::Text(root),
SqlValue::Text(dest),
SqlValue::Text(state),
SqlValue::Int(seq),
],
)?;
Ok(())
})
}
fn update_materialization_state(
&mut self,
id: u128,
state: &str,
seq: u64,
) -> Result<(), StoreError> {
let state = state.to_owned();
let seq = to_seq(seq, "seq")?;
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE materialization_records SET state = ?1, updated_seq = ?2 \
WHERE id_hex = ?3",
&[
SqlValue::Text(state),
SqlValue::Int(seq),
SqlValue::Text(u128_hex(id)),
],
)?;
if changed == 0 {
return Err(StoreError::UnknownMaterialization);
}
Ok(())
})
}
fn materialization_state(&mut self, id: u128) -> Result<Option<String>, StoreError> {
let rows = self.engine.query(
"SELECT state FROM materialization_records WHERE id_hex = ?1",
&[SqlValue::Text(u128_hex(id))],
)?;
match rows.first().and_then(|r| r.first()) {
None => Ok(None),
Some(v) => Ok(Some(expect_text(v, "materialization state")?)),
}
}
fn put_serving_record(
&mut self,
authority: &TypedDigest,
action_key: &str,
disposition: &str,
state_revision: u64,
validity: &ServingValidity,
blocking: &[(QuarantineScope, String)],
) -> Result<(), StoreError> {
let authority = authority.clone();
let action_key = action_key.to_owned();
let disposition = disposition.to_owned();
let revision = to_seq(state_revision, "state_revision")?;
let evaluated_at = validity.evaluated_at_unix_micros;
let max_age = validity
.maximum_age_micros
.map(|v| to_seq(v, "max_age_micros"))
.transpose()?;
let uncertainty = to_seq(
validity.clock_uncertainty_micros,
"clock_uncertainty_micros",
)?;
let epoch = to_seq(validity.coordinator_clock_epoch, "clock_epoch")?;
let blocking: Vec<(String, String)> = blocking
.iter()
.map(|(scope, subject)| (scope.as_str().to_owned(), subject.clone()))
.collect();
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let stored = engine.query(
"SELECT state_revision, disposition FROM action_serving_states WHERE action_key = ?1",
&[SqlValue::Text(action_key.clone())],
)?;
let stored_revision = match stored.first().and_then(|r| r.first()) {
None => None,
Some(v) => Some(expect_u64(v, "state_revision")?),
};
if state_revision == 0 || stored_revision.is_some_and(|s| state_revision <= s) {
return Err(StoreError::StaleServingRevision);
}
if let Some(row) = stored.first() {
let [_, current] = row.as_slice() else {
return Err(StoreError::Corruption("serving revision row shape".into()));
};
Self::require_quarantine_preserved(&expect_text(current, "disposition")?, &disposition)?;
}
let existing_blocking = engine.query(
"SELECT scope, subject FROM serving_blocking_quarantines WHERE action_key = ?1",
&[SqlValue::Text(action_key.clone())],
)?;
for row in &existing_blocking {
let [scope, subject] = row.as_slice() else {
return Err(StoreError::Corruption("blocking reference shape".into()));
};
let reference = (
expect_text(scope, "scope")?,
expect_text(subject, "subject")?,
);
if !blocking.contains(&reference) {
return Err(StoreError::QuarantineRequiresRepair);
}
}
for (scope, subject) in &blocking {
let exists = engine.query(
"SELECT scope FROM quarantines WHERE scope = ?1 AND subject = ?2",
&[
SqlValue::Text(scope.clone()),
SqlValue::Text(subject.clone()),
],
)?;
if exists.is_empty() {
return Err(StoreError::UnknownQuarantineReference);
}
}
engine.execute(
"INSERT OR REPLACE INTO action_serving_states \
(action_key, disposition, version, state_revision, authority_key, \
evaluated_at_micros, max_age_micros, clock_uncertainty_micros, clock_epoch) \
VALUES (?1, ?2, 1, ?3, ?4, ?5, ?6, ?7, ?8)",
&[
SqlValue::Text(action_key.clone()),
SqlValue::Text(disposition),
SqlValue::Int(revision),
SqlValue::Text(digest_key(&authority)),
SqlValue::Int(evaluated_at),
max_age.map_or(SqlValue::Null, SqlValue::Int),
SqlValue::Int(uncertainty),
SqlValue::Int(epoch),
],
)?;
engine.execute(
"DELETE FROM serving_blocking_quarantines WHERE action_key = ?1",
&[SqlValue::Text(action_key.clone())],
)?;
for (scope, subject) in &blocking {
engine.execute(
"INSERT OR IGNORE INTO serving_blocking_quarantines \
(action_key, scope, subject) VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(action_key.clone()),
SqlValue::Text(scope.clone()),
SqlValue::Text(subject.clone()),
],
)?;
}
Ok(())
})
}
fn serving_record(&mut self, action_key: &str) -> Result<Option<ServingRecordRow>, StoreError> {
let rows = self.engine.query(
"SELECT disposition, state_revision, authority_key, evaluated_at_micros, \
max_age_micros, clock_uncertainty_micros, clock_epoch \
FROM action_serving_states WHERE action_key = ?1",
&[SqlValue::Text(action_key.to_owned())],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [
disposition,
revision,
authority_key,
evaluated_at,
max_age,
uncertainty,
epoch,
] = row.as_slice()
else {
return Err(StoreError::Corruption("serving record shape".into()));
};
let evaluated_at = match evaluated_at {
SqlValue::Int(v) => *v,
_ => return Err(StoreError::Corruption("evaluated_at_micros shape".into())),
};
let max_age = match max_age {
SqlValue::Null => None,
SqlValue::Int(_) => Some(expect_u64(max_age, "max_age_micros")?),
_ => return Err(StoreError::Corruption("max_age_micros shape".into())),
};
let blocking_rows = self.engine.query(
"SELECT scope, subject FROM serving_blocking_quarantines \
WHERE action_key = ?1 ORDER BY scope, subject",
&[SqlValue::Text(action_key.to_owned())],
)?;
let blocking = blocking_rows
.iter()
.map(|row| {
let [scope, subject] = row.as_slice() else {
return Err(StoreError::Corruption("blocking reference shape".into()));
};
Ok((
expect_text(scope, "scope")?,
expect_text(subject, "subject")?,
))
})
.collect::<Result<Vec<_>, StoreError>>()?;
Ok(Some(ServingRecordRow {
disposition: expect_text(disposition, "disposition")?,
state_revision: expect_u64(revision, "state_revision")?,
authority_key: expect_text(authority_key, "authority_key")?,
validity: ServingValidity {
evaluated_at_unix_micros: evaluated_at,
maximum_age_micros: max_age,
clock_uncertainty_micros: expect_u64(uncertainty, "clock_uncertainty_micros")?,
coordinator_clock_epoch: expect_u64(epoch, "clock_epoch")?,
},
blocking,
}))
}
fn record_release_verdict(&mut self, verdict: &ReleaseVerdict) -> Result<(), StoreError> {
let replayed = to_seq(verdict.replayed, "replayed")?;
let explained = to_seq(verdict.explained, "explained")?;
let uncertainty = to_seq(
verdict.validity.clock_uncertainty_micros,
"clock_uncertainty_micros",
)?;
let epoch = to_seq(verdict.validity.coordinator_clock_epoch, "clock_epoch")?;
let max_age = match verdict.validity.maximum_age_micros {
None => SqlValue::Null,
Some(age) => SqlValue::Int(to_seq(age, "max_age_micros")?),
};
let build = verdict.build.clone();
let corpus = verdict.corpus.clone();
let evaluated_at = verdict.validity.evaluated_at_unix_micros;
self.in_txn(move |engine| {
engine.execute(
"DELETE FROM release_verdicts WHERE build = ?1",
&[SqlValue::Text(build.clone())],
)?;
engine.execute(
"INSERT INTO release_verdicts (build, corpus, replayed, explained, \
evaluated_at_micros, max_age_micros, clock_uncertainty_micros, clock_epoch) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
&[
SqlValue::Text(build),
SqlValue::Text(corpus),
SqlValue::Int(replayed),
SqlValue::Int(explained),
SqlValue::Int(evaluated_at),
max_age,
SqlValue::Int(uncertainty),
SqlValue::Int(epoch),
],
)?;
Ok(())
})
}
fn release_verdict(&mut self, build: &str) -> Result<Option<ReleaseVerdict>, StoreError> {
let rows = self.engine.query(
"SELECT corpus, replayed, explained, evaluated_at_micros, max_age_micros, \
clock_uncertainty_micros, clock_epoch FROM release_verdicts WHERE build = ?1",
&[SqlValue::Text(build.to_owned())],
)?;
let Some(row) = rows.first() else {
return Ok(None);
};
let [
corpus,
replayed,
explained,
evaluated_at,
max_age,
uncertainty,
epoch,
] = row.as_slice()
else {
return Err(StoreError::Corruption("release verdict shape".into()));
};
let evaluated_at = match evaluated_at {
SqlValue::Int(v) => *v,
_ => return Err(StoreError::Corruption("evaluated_at_micros shape".into())),
};
let maximum_age_micros = match max_age {
SqlValue::Null => None,
SqlValue::Int(_) => Some(expect_u64(max_age, "max_age_micros")?),
_ => return Err(StoreError::Corruption("max_age_micros shape".into())),
};
Ok(Some(ReleaseVerdict {
build: build.to_owned(),
corpus: expect_text(corpus, "corpus")?,
replayed: expect_u64(replayed, "replayed")?,
explained: expect_u64(explained, "explained")?,
validity: ServingValidity {
evaluated_at_unix_micros: evaluated_at,
maximum_age_micros,
clock_uncertainty_micros: expect_u64(uncertainty, "clock_uncertainty_micros")?,
coordinator_clock_epoch: expect_u64(epoch, "clock_epoch")?,
},
}))
}
fn record_divergence_incident(
&mut self,
authority: &TypedDigest,
row: &DivergenceIncidentRow,
) -> Result<(), StoreError> {
let authority = authority.clone();
let row = row.clone();
let seq = to_seq(row.seq, "seq")?;
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let existing = engine.query(
"SELECT class, committed_manifest_key, candidate_manifest_key, \
candidate_evidence_key, candidate_pin_hex, generation_hex, attempt_hex, \
detail FROM divergence_incidents WHERE action_key = ?1 AND seq = ?2",
&[SqlValue::Text(row.action_key.clone()), SqlValue::Int(seq)],
)?;
if let Some(stored) = existing.first() {
let same = stored.len() == 8
&& expect_text(&stored[0], "class")? == row.class
&& expect_text(&stored[1], "committed_manifest_key")?
== row.committed_manifest_key
&& expect_text(&stored[2], "candidate_manifest_key")?
== row.candidate_manifest_key
&& expect_text(&stored[3], "candidate_evidence_key")?
== row.candidate_evidence_key
&& expect_text(&stored[4], "candidate_pin_hex")? == row.candidate_pin_hex
&& expect_text(&stored[5], "generation_hex")? == row.generation_hex
&& expect_text(&stored[6], "attempt_hex")? == row.attempt_hex
&& expect_text(&stored[7], "detail")? == row.detail;
if same {
return Ok(()); }
return Err(StoreError::AppendConflict("divergence_incidents".into()));
}
engine.execute(
"INSERT INTO divergence_incidents (action_key, seq, class, \
committed_manifest_key, candidate_manifest_key, candidate_evidence_key, \
candidate_pin_hex, generation_hex, attempt_hex, detail) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
&[
SqlValue::Text(row.action_key),
SqlValue::Int(seq),
SqlValue::Text(row.class),
SqlValue::Text(row.committed_manifest_key),
SqlValue::Text(row.candidate_manifest_key),
SqlValue::Text(row.candidate_evidence_key),
SqlValue::Text(row.candidate_pin_hex),
SqlValue::Text(row.generation_hex),
SqlValue::Text(row.attempt_hex),
SqlValue::Text(row.detail),
],
)?;
Ok(())
})
}
fn list_divergence_incidents(
&mut self,
action_key: &str,
) -> Result<Vec<DivergenceIncidentRow>, StoreError> {
let rows = self.engine.query(
"SELECT seq, class, committed_manifest_key, candidate_manifest_key, \
candidate_evidence_key, candidate_pin_hex, generation_hex, attempt_hex, detail \
FROM divergence_incidents WHERE action_key = ?1 ORDER BY seq",
&[SqlValue::Text(action_key.to_owned())],
)?;
rows.iter()
.map(|row| {
let [
seq,
class,
committed,
candidate,
evidence,
pin,
generation,
attempt,
detail,
] = row.as_slice()
else {
return Err(StoreError::Corruption("divergence incident shape".into()));
};
Ok(DivergenceIncidentRow {
action_key: action_key.to_owned(),
seq: expect_u64(seq, "seq")?,
class: expect_text(class, "class")?,
committed_manifest_key: expect_text(committed, "committed_manifest_key")?,
candidate_manifest_key: expect_text(candidate, "candidate_manifest_key")?,
candidate_evidence_key: expect_text(evidence, "candidate_evidence_key")?,
candidate_pin_hex: expect_text(pin, "candidate_pin_hex")?,
generation_hex: expect_text(generation, "generation_hex")?,
attempt_hex: expect_text(attempt, "attempt_hex")?,
detail: expect_text(detail, "detail")?,
})
})
.collect()
}
fn record_adoption_edge(
&mut self,
authority: &TypedDigest,
producer_action_key: &str,
role: &str,
virtual_path: &[u8],
from_object_key: &str,
to_object_key: &str,
) -> Result<(), StoreError> {
let authority = authority.clone();
let producer = producer_action_key.to_owned();
let role = role.to_owned();
let path = virtual_path.to_vec();
let from = from_object_key.to_owned();
let to = to_object_key.to_owned();
self.in_txn(move |engine| {
SqlMetadataStore::<E>::require_active(engine, &authority)?;
let existing = engine.query(
"SELECT to_object_key FROM adoption_edges \
WHERE producer_action_key = ?1 AND role = ?2 AND virtual_path = ?3 \
AND from_object_key = ?4",
&[
SqlValue::Text(producer.clone()),
SqlValue::Text(role.clone()),
SqlValue::Blob(path.clone()),
SqlValue::Text(from.clone()),
],
)?;
if let Some(row) = existing.first() {
let [existing_to] = row.as_slice() else {
return Err(StoreError::Corruption("adoption edge shape".into()));
};
return if expect_text(existing_to, "adoption target")? == to {
Ok(())
} else {
Err(StoreError::AdoptionEdgeConflict)
};
}
engine.execute(
"INSERT INTO adoption_edges (producer_action_key, role, virtual_path, \
from_object_key, to_object_key) VALUES (?1, ?2, ?3, ?4, ?5)",
&[
SqlValue::Text(producer),
SqlValue::Text(role),
SqlValue::Blob(path),
SqlValue::Text(from),
SqlValue::Text(to),
],
)?;
Ok(())
})
}
fn has_adoption_edge(
&mut self,
producer_action_key: &str,
role: &str,
virtual_path: &[u8],
from_object_key: &str,
to_object_key: &str,
) -> Result<bool, StoreError> {
let rows = self.engine.query(
"SELECT producer_action_key FROM adoption_edges \
WHERE producer_action_key = ?1 AND role = ?2 AND virtual_path = ?3 \
AND from_object_key = ?4 AND to_object_key = ?5 LIMIT 1",
&[
SqlValue::Text(producer_action_key.to_owned()),
SqlValue::Text(role.to_owned()),
SqlValue::Blob(virtual_path.to_vec()),
SqlValue::Text(from_object_key.to_owned()),
SqlValue::Text(to_object_key.to_owned()),
],
)?;
Ok(!rows.is_empty())
}
fn list_provisional_ancestors(
&mut self,
consumer_action_key: &str,
) -> Result<Vec<ProvisionalAncestorRow>, StoreError> {
let rows = self.engine.query(
"SELECT producer_action_key, role, virtual_path, object_key, adopted \
FROM provisional_ancestry WHERE consumer_action_key = ?1 \
ORDER BY producer_action_key, role, virtual_path",
&[SqlValue::Text(consumer_action_key.to_owned())],
)?;
rows.iter()
.map(|row| {
let [producer, role, path, object, adopted] = row.as_slice() else {
return Err(StoreError::Corruption("provisional ancestry shape".into()));
};
let SqlValue::Blob(path) = path else {
return Err(StoreError::Corruption("ancestry path shape".into()));
};
Ok(ProvisionalAncestorRow {
producer_action_key: expect_text(producer, "ancestry producer")?,
role: expect_text(role, "ancestry role")?,
virtual_path: path.clone(),
object_key: expect_text(object, "ancestry object")?,
adopted: expect_u64(adopted, "ancestry adopted")? != 0,
})
})
.collect()
}
fn provisional_pin_row(
&mut self,
pin_key: &str,
) -> Result<Option<ProvisionalPinRecord>, StoreError> {
let rows = self.engine.query(
"SELECT pin_key, authority_key, action_key, generation_hex, attempt_hex, \
lease_hex, role, virtual_path, obj_algo, obj_domain, obj_bytes, object_key, \
protective_pin_hex, renewal_seq, adopted_object_key, invalidated_reason, \
released, toolchain_contract_key, event_contract_key \
FROM provisional_pins WHERE pin_key = ?1",
&[SqlValue::Text(pin_key.to_owned())],
)?;
rows.first()
.map(|row| self.map_provisional_pin_row(row))
.transpose()
}
fn insert_provisional_pin(&mut self, pin: &ProvisionalPinInsert) -> Result<(), StoreError> {
let row = pin.clone();
self.intern(row.object.domain);
self.in_txn(move |engine| {
engine.execute(
"INSERT INTO provisional_pins (pin_key, authority_key, action_key, \
generation_hex, attempt_hex, lease_hex, role, virtual_path, obj_algo, \
obj_domain, obj_bytes, object_key, protective_pin_hex, renewal_seq, \
adopted_object_key, invalidated_reason, released, \
toolchain_contract_key, event_contract_key) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, 0, \
NULL, NULL, 0, ?14, ?15)",
&[
SqlValue::Text(row.pin_key.clone()),
SqlValue::Text(row.authority_key.clone()),
SqlValue::Text(row.action_key.clone()),
SqlValue::Text(u128_hex(row.generation)),
SqlValue::Text(u128_hex(row.attempt)),
SqlValue::Text(u128_hex(row.lease)),
SqlValue::Int(row.role_tag),
SqlValue::Blob(row.virtual_path.clone()),
SqlValue::Text(algo_tag(row.object.algorithm).to_owned()),
SqlValue::Text(row.object.domain.to_owned()),
SqlValue::Blob(row.object.bytes.to_vec()),
SqlValue::Text(digest_key(&row.object)),
SqlValue::Text(u128_hex(row.protective_pin_id)),
SqlValue::Text(row.toolchain_contract_key.clone()),
SqlValue::Text(row.event_contract_key.clone()),
],
)?;
engine.execute(
"INSERT INTO pins (id_hex, id, root_key, owner, class, expires_at_seq, \
released, evidence, renewal_seq, durable, reason) \
VALUES (?1, ?2, ?3, 'coordinator', 'provisional-metadata', NULL, 0, \
NULL, 0, 1, ?4)",
&[
SqlValue::Text(u128_hex(row.protective_pin_id)),
SqlValue::Blob(row.protective_pin_id.to_be_bytes().to_vec()),
SqlValue::Text(digest_key(&row.object)),
SqlValue::Text(row.reason),
],
)?;
for (ancestor_key, min_hops) in &row.ancestor_pin_keys {
engine.execute(
"INSERT OR IGNORE INTO provisional_pin_lineage \
(descendant_pin_key, ancestor_pin_key, min_hops) VALUES (?1, ?2, ?3)",
&[
SqlValue::Text(row.pin_key.clone()),
SqlValue::Text(ancestor_key.clone()),
SqlValue::Int(*min_hops as i64),
],
)?;
}
Ok(())
})
}
fn list_provisional_pin_ancestors(
&mut self,
descendant_pin_key: &str,
) -> Result<Vec<(String, u64)>, StoreError> {
let rows = self.engine.query(
"SELECT ancestor_pin_key, min_hops FROM provisional_pin_lineage \
WHERE descendant_pin_key = ?1 ORDER BY ancestor_pin_key",
&[SqlValue::Text(descendant_pin_key.to_owned())],
)?;
rows.into_iter()
.map(|row| match row.as_slice() {
[SqlValue::Text(k), SqlValue::Int(h)] => {
Ok((k.clone(), u64::try_from(*h).unwrap_or(0)))
}
_ => Err(StoreError::Corruption("pin lineage ancestor shape".into())),
})
.collect()
}
fn provisional_pin_closure_depth(
&mut self,
descendant_pin_key: &str,
) -> Result<u64, StoreError> {
let rows = self.engine.query(
"SELECT COALESCE(MAX(min_hops), 0) FROM provisional_pin_lineage \
WHERE descendant_pin_key = ?1",
&[SqlValue::Text(descendant_pin_key.to_owned())],
)?;
match rows.first().map(Vec::as_slice) {
Some([SqlValue::Int(h)]) => Ok(u64::try_from(*h).unwrap_or(0)),
_ => Err(StoreError::Corruption("pin lineage depth shape".into())),
}
}
fn list_provisional_pin_descendants(
&mut self,
ancestor_pin_key: &str,
) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT descendant_pin_key FROM provisional_pin_lineage \
WHERE ancestor_pin_key = ?1 ORDER BY descendant_pin_key",
&[SqlValue::Text(ancestor_pin_key.to_owned())],
)?;
rows.into_iter()
.map(|row| match row.first() {
Some(SqlValue::Text(k)) => Ok(k.clone()),
_ => Err(StoreError::Corruption(
"pin lineage descendant shape".into(),
)),
})
.collect()
}
fn record_provisional_grant(
&mut self,
pin_key: &str,
grantee_kind: &str,
grantee_id: &str,
granted_seq: u64,
) -> Result<(), StoreError> {
let (pin_key, kind, id) = (
pin_key.to_owned(),
grantee_kind.to_owned(),
grantee_id.to_owned(),
);
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO provisional_pin_grants \
(pin_key, grantee_kind, grantee_id, granted_seq) VALUES (?1, ?2, ?3, ?4)",
&[
SqlValue::Text(pin_key),
SqlValue::Text(kind),
SqlValue::Text(id),
SqlValue::Int(granted_seq as i64),
],
)?;
Ok(())
})
}
fn list_provisional_grants(
&mut self,
pin_key: &str,
) -> Result<Vec<(String, String, u64)>, StoreError> {
let rows = self.engine.query(
"SELECT grantee_kind, grantee_id, granted_seq FROM provisional_pin_grants \
WHERE pin_key = ?1 ORDER BY grantee_kind, grantee_id",
&[SqlValue::Text(pin_key.to_owned())],
)?;
rows.iter()
.map(|row| {
let [kind, id, seq] = row.as_slice() else {
return Err(StoreError::Corruption("provisional grant shape".into()));
};
Ok((
expect_text(kind, "grant kind")?,
expect_text(id, "grant id")?,
expect_u64(seq, "grant seq")?,
))
})
.collect()
}
fn renew_provisional_pin(&mut self, pin_key: &str, renewal_seq: u64) -> Result<(), StoreError> {
let pin_key = pin_key.to_owned();
self.in_txn(move |engine| {
let rows = engine.query(
"SELECT renewal_seq, released FROM provisional_pins WHERE pin_key = ?1",
&[SqlValue::Text(pin_key.clone())],
)?;
let Some(row) = rows.first() else {
return Err(StoreError::UnknownPin);
};
let [stored, released] = row.as_slice() else {
return Err(StoreError::Corruption("provisional renewal shape".into()));
};
if expect_u64(released, "released")? != 0 {
return Err(StoreError::PinReleased);
}
if renewal_seq <= expect_u64(stored, "renewal")? {
return Err(StoreError::NonMonotonicPinRenewal);
}
engine.execute(
"UPDATE provisional_pins SET renewal_seq = ?1 WHERE pin_key = ?2",
&[SqlValue::Int(renewal_seq as i64), SqlValue::Text(pin_key)],
)?;
Ok(())
})
}
fn close_provisional_pin(
&mut self,
pin_key: &str,
invalidation_reason: Option<&str>,
) -> Result<(), StoreError> {
let pin_key = pin_key.to_owned();
let reason = invalidation_reason.map(str::to_owned);
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE provisional_pins SET released = 1, invalidated_reason = \
COALESCE(?1, invalidated_reason) WHERE pin_key = ?2 AND released = 0",
&[
reason.map(SqlValue::Text).unwrap_or(SqlValue::Null),
SqlValue::Text(pin_key.clone()),
],
)?;
if changed == 0
&& engine
.query(
"SELECT 1 FROM provisional_pins WHERE pin_key = ?1",
&[SqlValue::Text(pin_key)],
)?
.is_empty()
{
return Err(StoreError::UnknownPin);
}
Ok(())
})
}
fn adopt_provisional_pin(
&mut self,
pin_key: &str,
committed_object_key: &str,
) -> Result<(), StoreError> {
let pin_key = pin_key.to_owned();
let committed = committed_object_key.to_owned();
self.in_txn(move |engine| {
let changed = engine.execute(
"UPDATE provisional_pins SET adopted_object_key = ?1 \
WHERE pin_key = ?2 AND adopted_object_key IS NULL",
&[SqlValue::Text(committed), SqlValue::Text(pin_key.clone())],
)?;
if changed == 0
&& engine
.query(
"SELECT 1 FROM provisional_pins WHERE pin_key = ?1",
&[SqlValue::Text(pin_key)],
)?
.is_empty()
{
return Err(StoreError::UnknownPin);
}
Ok(())
})
}
fn record_provisional_consumption(
&mut self,
consumption: &ProvisionalObligationInsert,
) -> Result<(), StoreError> {
let row = consumption.clone();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO provisional_obligations (consumer_worker, \
consumer_attempt_hex, pin_key, producer_action_key, \
producer_generation_hex, producer_attempt_hex, role, virtual_path, \
object_key, status, resolution_object_key, created_seq) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, 'open', NULL, ?10)",
&[
SqlValue::Text(row.consumer_worker),
SqlValue::Text(u128_hex(row.consumer_attempt)),
SqlValue::Text(row.pin_key),
SqlValue::Text(row.producer_action_key),
SqlValue::Text(u128_hex(row.producer_generation)),
SqlValue::Text(u128_hex(row.producer_attempt)),
SqlValue::Int(row.role_tag),
SqlValue::Blob(row.virtual_path),
SqlValue::Text(row.object_key),
SqlValue::Int(to_seq(row.created_seq, "consumption seq")?),
],
)?;
Ok(())
})
}
fn list_open_provisional_obligations(
&mut self,
consumer_worker: &str,
consumer_attempt_hex: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError> {
let rows = self.engine.query(
"SELECT consumer_worker, consumer_attempt_hex, pin_key, producer_action_key, \
producer_generation_hex, producer_attempt_hex, role, virtual_path, object_key, \
status, resolution_object_key, created_seq FROM provisional_obligations \
WHERE consumer_worker = ?1 AND consumer_attempt_hex = ?2 AND status != 'resolved' \
ORDER BY status, pin_key",
&[
SqlValue::Text(consumer_worker.to_owned()),
SqlValue::Text(consumer_attempt_hex.to_owned()),
],
)?;
rows.iter()
.map(|row| {
let [
worker,
attempt_hex,
pin,
action,
generation,
producer,
role,
path,
object,
status,
resolution,
created,
] = row.as_slice()
else {
return Err(StoreError::Corruption("obligation row shape".into()));
};
Ok(ProvisionalObligationRow {
consumer_worker: expect_text(worker, "obligation worker")?,
consumer_attempt_hex: expect_text(attempt_hex, "obligation attempt")?,
pin_key: expect_text(pin, "obligation pin")?,
producer_action_key: expect_text(action, "obligation action")?,
producer_generation_hex: expect_text(generation, "obligation generation")?,
producer_attempt_hex: expect_text(producer, "obligation producer")?,
role_tag: match role {
SqlValue::Int(v) => *v,
_ => return Err(StoreError::Corruption("obligation role shape".into())),
},
virtual_path: match path {
SqlValue::Blob(b) => b.clone(),
_ => return Err(StoreError::Corruption("obligation path shape".into())),
},
object_key: expect_text(object, "obligation object")?,
status: expect_text(status, "obligation status")?,
resolution_object_key: expect_opt_text(resolution, "obligation resolution")?,
created_seq: expect_u64(created, "obligation seq")?,
})
})
.collect()
}
fn resolve_provisional_obligations(
&mut self,
pin_key: &str,
resolution_object_key: &str,
) -> Result<usize, StoreError> {
let (pin_key, resolved) = (pin_key.to_owned(), resolution_object_key.to_owned());
self.in_txn(move |engine| {
engine.execute(
"UPDATE provisional_obligations SET status = 'resolved', \
resolution_object_key = ?1 WHERE pin_key = ?2 AND status = 'open'",
&[SqlValue::Text(resolved), SqlValue::Text(pin_key)],
)
})
}
fn cancel_provisional_obligations(&mut self, pin_key: &str) -> Result<usize, StoreError> {
let pin_key = pin_key.to_owned();
self.in_txn(move |engine| {
engine.execute(
"UPDATE provisional_obligations SET status = 'cancelled' \
WHERE pin_key = ?1 AND status = 'open'",
&[SqlValue::Text(pin_key)],
)
})
}
fn record_served_consumer(
&mut self,
action_key: &str,
consumer: &str,
) -> Result<(), StoreError> {
let action_key = action_key.to_owned();
let consumer = consumer.to_owned();
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO provenance_edges (from_key, to_key, kind) \
VALUES (?1, ?2, 'served-to')",
&[SqlValue::Text(action_key), SqlValue::Text(consumer)],
)?;
Ok(())
})
}
fn list_served_consumers(&mut self, action_key: &str) -> Result<Vec<String>, StoreError> {
let rows = self.engine.query(
"SELECT to_key FROM provenance_edges \
WHERE from_key = ?1 AND kind = 'served-to' ORDER BY to_key",
&[SqlValue::Text(action_key.to_owned())],
)?;
rows.iter()
.map(|row| {
expect_text(
row.first()
.ok_or_else(|| StoreError::Corruption("served consumer shape".into()))?,
"consumer",
)
})
.collect()
}
fn differential_snapshot(&mut self) -> Result<Vec<String>, StoreError> {
const DUMPS: &[(&str, &str)] = &[
(
"schema_epochs",
"SELECT version, applied_seq FROM schema_epochs ORDER BY version",
),
(
"coordinator_authorities",
"SELECT key, algo, domain, bytes, cluster_id, incarnation, term, acquired_seq, \
released FROM coordinator_authorities ORDER BY key",
),
(
"action_entries",
"SELECT key, algo, domain, bytes, key_epoch, projection_epoch \
FROM action_entries ORDER BY key",
),
(
"action_generations",
"SELECT id_hex, id, action_key, authority_key, tombstoned, per_key_ordinal \
FROM action_generations ORDER BY id_hex",
),
(
"generation_high_water",
"SELECT kind, value FROM generation_high_water ORDER BY kind",
),
(
"action_attempts",
"SELECT id_hex, id, generation_hex, worker, seq, worker_boot_generation, \
worker_incarnation, execution_lease_hex FROM action_attempts \
ORDER BY id_hex",
),
(
"execution_leases",
"SELECT id_hex, id, attempt_hex, renewal_seq, expires_at_seq, released \
FROM execution_leases ORDER BY id_hex",
),
(
"action_publications",
"SELECT action_key, descriptor_algo, descriptor_domain, descriptor_bytes, \
manifest_algo, manifest_domain, manifest_bytes, winner_generation_hex, \
winner_attempt_hex, result_kind, pin_hex FROM action_publications \
ORDER BY action_key",
),
(
"action_serving_states",
"SELECT action_key, disposition, version, state_revision, authority_key, \
evaluated_at_micros, max_age_micros, clock_uncertainty_micros, clock_epoch \
FROM action_serving_states ORDER BY action_key",
),
(
"serving_blocking_quarantines",
"SELECT action_key, scope, subject FROM serving_blocking_quarantines \
ORDER BY action_key, scope, subject",
),
(
"objects",
"SELECT key, algo, domain, bytes, logical_size FROM objects ORDER BY key",
),
(
"object_locations",
"SELECT object_key, store_path, verified_seq, encoding, quarantined, durable \
FROM object_locations ORDER BY object_key, store_path",
),
(
"object_edges",
"SELECT parent_key, child_key, kind FROM object_edges \
ORDER BY parent_key, child_key, kind",
),
(
"pins",
"SELECT id_hex, id, root_key, owner, class, expires_at_seq, released, \
evidence, renewal_seq, durable, reason FROM pins ORDER BY id_hex",
),
(
"observed_input_recipes",
"SELECT action_key, recipe_algo, recipe_domain, recipe_bytes \
FROM observed_input_recipes ORDER BY action_key",
),
(
"key_breakdowns",
"SELECT action_key, component, algo, domain, bytes FROM key_breakdowns \
ORDER BY action_key, component",
),
(
"trust_states",
"SELECT action_key, state, reason FROM trust_states ORDER BY action_key",
),
(
"quarantines",
"SELECT scope, subject, reason FROM quarantines ORDER BY scope, subject",
),
(
"verification_samples",
"SELECT action_key, attempt_hex, passed, seq FROM verification_samples \
ORDER BY action_key, attempt_hex, seq",
),
(
"gc_runs",
"SELECT id, seq, pinned_roots, located_objects, reachable_objects \
FROM gc_runs ORDER BY id",
),
(
"operator_resets",
"SELECT generation, applied_seq FROM operator_resets ORDER BY generation",
),
(
"eviction_tombstones",
"SELECT action_key, semantic_algo, semantic_domain, semantic_bytes, \
observable_algo, observable_domain, observable_bytes, evicted_seq \
FROM eviction_tombstones ORDER BY action_key",
),
(
"gc_tombstones",
"SELECT object_key, store_path, marked_seq, grace_until_seq \
FROM gc_tombstones ORDER BY object_key, store_path",
),
(
"gc_receipts",
"SELECT id, seq, mode, planned, reclaimed, skipped, truncated \
FROM gc_receipts ORDER BY id",
),
(
"action_evidence_index",
"SELECT action_key, evidence_algo, evidence_domain, evidence_bytes, \
generation_hex, attempt_hex, manifest_key FROM action_evidence_index \
ORDER BY action_key, evidence_domain, evidence_bytes",
),
(
"peer_authority_high_water",
"SELECT peer_id, term, observed_seq, incarnation, credential_generation FROM peer_authority_high_water \
ORDER BY peer_id",
),
(
"worker_incarnation_fences",
"SELECT worker, incarnation, highest_boot_generation, active, \
operator_reenrollment_generation, clone_ambiguous \
FROM worker_incarnation_fences \
ORDER BY worker",
),
(
"edge_incarnation_fences",
"SELECT edge_id, incarnation FROM edge_incarnation_fences ORDER BY edge_id",
),
(
"edge_handoffs",
"SELECT edge_id, active_incarnation, predecessor_incarnation, begun_seq, \
resolved FROM edge_handoffs ORDER BY edge_id",
),
(
"action_trust_evaluations",
"SELECT action_key, version, state, reason, evaluated_seq \
FROM action_trust_evaluations ORDER BY action_key, version",
),
(
"operations",
"SELECT id_hex, id, kind, state, updated_seq FROM operations ORDER BY id_hex",
),
(
"edge_subscribers",
"SELECT edge_id, subscriber, registered_seq FROM edge_subscribers \
ORDER BY edge_id, subscriber",
),
(
"manifests",
"SELECT key, algo, domain, bytes, kind, entry_count FROM manifests \
ORDER BY key",
),
(
"worker_sessions",
"SELECT worker, incarnation, started_seq, ended_seq FROM worker_sessions \
ORDER BY worker, started_seq",
),
(
"worker_capabilities",
"SELECT worker, capability FROM worker_capabilities \
ORDER BY worker, capability",
),
(
"native_child_bindings",
"SELECT parent_action_key, child_action_key, bound_seq, state \
FROM native_child_bindings \
ORDER BY parent_action_key, child_action_key",
),
(
"worker_health_samples",
"SELECT worker, seq, healthy, detail FROM worker_health_samples \
ORDER BY worker, seq",
),
(
"decision_receipts",
"SELECT kind, subject, seq, decision, reason FROM decision_receipts \
ORDER BY kind, subject, seq",
),
(
"provenance_edges",
"SELECT from_key, to_key, kind FROM provenance_edges \
ORDER BY from_key, to_key, kind",
),
(
"determinism_audits",
"SELECT action_key, attempt_hex, seq, verdict FROM determinism_audits \
ORDER BY action_key, attempt_hex, seq",
),
(
"materialization_records",
"SELECT id_hex, id, root_key, dest_path, state, updated_seq \
FROM materialization_records ORDER BY id_hex",
),
(
"divergence_incidents",
"SELECT action_key, seq, class, committed_manifest_key, \
candidate_manifest_key, candidate_evidence_key, candidate_pin_hex, \
generation_hex, attempt_hex, detail FROM divergence_incidents \
ORDER BY action_key, seq",
),
(
"provisional_ancestry",
"SELECT consumer_action_key, producer_action_key, role, virtual_path, \
object_key, adopted FROM provisional_ancestry \
ORDER BY consumer_action_key, producer_action_key, role, virtual_path",
),
(
"adoption_edges",
"SELECT producer_action_key, role, virtual_path, from_object_key, \
to_object_key FROM adoption_edges \
ORDER BY producer_action_key, role, virtual_path, from_object_key",
),
(
"provisional_pins",
"SELECT pin_key, authority_key, action_key, generation_hex, attempt_hex, \
lease_hex, role, virtual_path, obj_algo, obj_domain, obj_bytes, object_key, \
protective_pin_hex, renewal_seq, adopted_object_key, invalidated_reason, \
released, toolchain_contract_key, event_contract_key \
FROM provisional_pins ORDER BY pin_key",
),
(
"provisional_pin_grants",
"SELECT pin_key, grantee_kind, grantee_id, granted_seq \
FROM provisional_pin_grants ORDER BY pin_key, grantee_kind, grantee_id",
),
(
"provisional_pin_lineage",
"SELECT descendant_pin_key, ancestor_pin_key, min_hops \
FROM provisional_pin_lineage \
ORDER BY descendant_pin_key, ancestor_pin_key",
),
(
"provisional_install_journal",
"SELECT pin_key, consumer_worker, consumer_attempt_hex, installed_path, \
obj_algo, obj_domain, obj_bytes, object_key, installed_seq, state \
FROM provisional_install_journal \
ORDER BY pin_key, consumer_attempt_hex, installed_path",
),
(
"provisional_obligations",
"SELECT consumer_worker, consumer_attempt_hex, pin_key, producer_action_key, \
producer_generation_hex, producer_attempt_hex, role, virtual_path, object_key, \
status, resolution_object_key, created_seq FROM provisional_obligations \
ORDER BY consumer_worker, consumer_attempt_hex, pin_key",
),
];
let mut lines = Vec::new();
for (table, sql) in DUMPS {
for row in self.engine.query(sql, &[])? {
let mut line = String::new();
line.push_str(table);
for value in &row {
line.push('|');
match value {
SqlValue::Null => line.push_str("NULL"),
SqlValue::Int(i) => {
use std::fmt::Write;
let _ = write!(line, "{i}");
}
SqlValue::Text(t) => line.push_str(t),
SqlValue::Blob(b) => line.push_str(&hex(b)),
}
}
lines.push(line);
}
}
Ok(lines)
}
fn count_open_provisional_obligations(&mut self, pin_key: &str) -> Result<usize, StoreError> {
let rows = self.engine.query(
"SELECT COUNT(*) FROM provisional_obligations \
WHERE pin_key = ?1 AND status = 'open'",
&[SqlValue::Text(pin_key.to_owned())],
)?;
let Some(row) = rows.first() else {
return Ok(0);
};
expect_u64(
row.first()
.ok_or_else(|| StoreError::Corruption("obligation count shape".into()))?,
"obligation count",
)
.map(|v| v as usize)
}
fn list_open_provisional_pins_for_action_generation(
&mut self,
action_key: &str,
generation_hex: &str,
) -> Result<Vec<ProvisionalPinRecord>, StoreError> {
let rows = self.engine.query(
"SELECT pin_key, authority_key, action_key, generation_hex, attempt_hex, \
lease_hex, role, virtual_path, obj_algo, obj_domain, obj_bytes, object_key, \
protective_pin_hex, renewal_seq, adopted_object_key, invalidated_reason, \
released, toolchain_contract_key, event_contract_key \
FROM provisional_pins \
WHERE action_key = ?1 AND generation_hex = ?2 AND released = 0 ORDER BY pin_key",
&[
SqlValue::Text(action_key.to_owned()),
SqlValue::Text(generation_hex.to_owned()),
],
)?;
rows.iter()
.map(|row| self.map_provisional_pin_row(row))
.collect()
}
fn list_open_provisional_pins_for_action(
&mut self,
action_key: &str,
) -> Result<Vec<ProvisionalPinRecord>, StoreError> {
let rows = self.engine.query(
"SELECT pin_key, authority_key, action_key, generation_hex, attempt_hex, \
lease_hex, role, virtual_path, obj_algo, obj_domain, obj_bytes, object_key, \
protective_pin_hex, renewal_seq, adopted_object_key, invalidated_reason, \
released, toolchain_contract_key, event_contract_key \
FROM provisional_pins WHERE action_key = ?1 AND released = 0 \
ORDER BY pin_key",
&[SqlValue::Text(action_key.to_owned())],
)?;
rows.iter()
.map(|row| self.map_provisional_pin_row(row))
.collect()
}
fn list_open_provisional_pins_for_authority(
&mut self,
authority_key: &str,
) -> Result<Vec<ProvisionalPinRecord>, StoreError> {
let rows = self.engine.query(
"SELECT pin_key, authority_key, action_key, generation_hex, attempt_hex, \
lease_hex, role, virtual_path, obj_algo, obj_domain, obj_bytes, object_key, \
protective_pin_hex, renewal_seq, adopted_object_key, invalidated_reason, \
released, toolchain_contract_key, event_contract_key \
FROM provisional_pins WHERE authority_key = ?1 AND released = 0 \
ORDER BY pin_key",
&[SqlValue::Text(authority_key.to_owned())],
)?;
rows.iter()
.map(|row| self.map_provisional_pin_row(row))
.collect()
}
fn list_open_provisional_obligations_by_attempt(
&mut self,
consumer_attempt_hex: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError> {
let rows = self.engine.query(
"SELECT consumer_worker, consumer_attempt_hex, pin_key, producer_action_key, \
producer_generation_hex, producer_attempt_hex, role, virtual_path, object_key, \
status, resolution_object_key, created_seq FROM provisional_obligations \
WHERE consumer_attempt_hex = ?1 AND status != 'resolved' \
ORDER BY pin_key",
&[SqlValue::Text(consumer_attempt_hex.to_owned())],
)?;
rows.iter()
.map(|row| Self::map_obligation_row(row))
.collect()
}
fn list_provisional_obligations_by_attempt_all(
&mut self,
consumer_attempt_hex: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError> {
let rows = self.engine.query(
"SELECT consumer_worker, consumer_attempt_hex, pin_key, producer_action_key, \
producer_generation_hex, producer_attempt_hex, role, virtual_path, object_key, \
status, resolution_object_key, created_seq FROM provisional_obligations \
WHERE consumer_attempt_hex = ?1 \
ORDER BY status, pin_key",
&[SqlValue::Text(consumer_attempt_hex.to_owned())],
)?;
rows.iter()
.map(|row| Self::map_obligation_row(row))
.collect()
}
fn insert_provisional_install(
&mut self,
install: &ProvisionalInstallInsert,
) -> Result<(), StoreError> {
let row = install.clone();
self.intern(row.object.domain);
self.in_txn(move |engine| {
engine.execute(
"INSERT OR IGNORE INTO provisional_install_journal \
(pin_key, consumer_worker, consumer_attempt_hex, installed_path, \
obj_algo, obj_domain, obj_bytes, object_key, installed_seq, state) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, 'installed')",
&[
SqlValue::Text(row.pin_key),
SqlValue::Text(row.consumer_worker),
SqlValue::Text(u128_hex(row.consumer_attempt)),
SqlValue::Blob(row.installed_path),
SqlValue::Text(algo_tag(row.object.algorithm).to_owned()),
SqlValue::Text(row.object.domain.to_owned()),
SqlValue::Blob(row.object.bytes.to_vec()),
SqlValue::Text(digest_key(&row.object)),
SqlValue::Int(row.installed_seq as i64),
],
)?;
Ok(())
})
}
fn list_provisional_installs_for_pins(
&mut self,
pin_keys: &[String],
) -> Result<Vec<ProvisionalInstallRecord>, StoreError> {
if pin_keys.is_empty() {
return Ok(Vec::new());
}
let placeholders = (0..pin_keys.len())
.map(|i| format!("?{}", i + 1))
.collect::<Vec<_>>()
.join(", ");
let params: Vec<SqlValue> = pin_keys
.iter()
.map(|key| SqlValue::Text(key.clone()))
.collect();
let rows = self.engine.query(
&format!(
"SELECT pin_key, consumer_worker, consumer_attempt_hex, installed_path, \
obj_algo, obj_domain, obj_bytes, object_key, installed_seq, state \
FROM provisional_install_journal WHERE pin_key IN ({placeholders}) \
ORDER BY pin_key, consumer_attempt_hex, installed_path"
),
¶ms,
)?;
rows.iter()
.map(|row| Self::map_install_row(self, row))
.collect()
}
fn list_provisional_installs_by_state(
&mut self,
state: &str,
) -> Result<Vec<ProvisionalInstallRecord>, StoreError> {
let rows = self.engine.query(
"SELECT pin_key, consumer_worker, consumer_attempt_hex, installed_path, \
obj_algo, obj_domain, obj_bytes, object_key, installed_seq, state \
FROM provisional_install_journal WHERE state = ?1 \
ORDER BY pin_key, consumer_attempt_hex, installed_path",
&[SqlValue::Text(state.to_owned())],
)?;
rows.iter()
.map(|row| Self::map_install_row(self, row))
.collect()
}
fn set_provisional_install_state(
&mut self,
pin_key: &str,
consumer_attempt_hex: &str,
installed_path: &[u8],
state: &str,
) -> Result<(), StoreError> {
let (pin_key, attempt_hex, path, state) = (
pin_key.to_owned(),
consumer_attempt_hex.to_owned(),
installed_path.to_vec(),
state.to_owned(),
);
self.in_txn(move |engine| {
let affected = engine.execute(
"UPDATE provisional_install_journal SET state = ?4 \
WHERE pin_key = ?1 AND consumer_attempt_hex = ?2 AND installed_path = ?3",
&[
SqlValue::Text(pin_key),
SqlValue::Text(attempt_hex),
SqlValue::Blob(path),
SqlValue::Text(state),
],
)?;
if affected == 0 {
return Err(StoreError::UnknownPin);
}
Ok(())
})
}
fn bind_native_children(
&mut self,
parent_action_key: &str,
child_action_keys: &[String],
bound_seq: u64,
) -> Result<(), StoreError> {
let parent = parent_action_key.to_owned();
let children: Vec<String> = child_action_keys.to_vec();
self.in_txn(move |engine| {
for child in children {
engine.execute(
"INSERT OR IGNORE INTO native_child_bindings \
(parent_action_key, child_action_key, bound_seq, state) \
VALUES (?1, ?2, ?3, 'bound')",
&[
SqlValue::Text(parent.clone()),
SqlValue::Text(child),
SqlValue::Int(bound_seq as i64),
],
)?;
}
Ok(())
})
}
fn list_native_child_bindings(
&mut self,
parent_action_key: &str,
) -> Result<Vec<(String, String)>, StoreError> {
let rows = self.engine.query(
"SELECT child_action_key, state FROM native_child_bindings \
WHERE parent_action_key = ?1 ORDER BY child_action_key",
&[SqlValue::Text(parent_action_key.to_owned())],
)?;
rows.into_iter()
.map(|row| match row.as_slice() {
[SqlValue::Text(child), SqlValue::Text(state)] => {
Ok((child.clone(), state.clone()))
}
_ => Err(StoreError::Corruption("native child binding shape".into())),
})
.collect()
}
fn set_native_child_binding_state(
&mut self,
parent_action_key: &str,
child_action_key: &str,
state: &str,
) -> Result<(), StoreError> {
let (parent, child, state) = (
parent_action_key.to_owned(),
child_action_key.to_owned(),
state.to_owned(),
);
self.in_txn(move |engine| {
let affected = engine.execute(
"UPDATE native_child_bindings SET state = ?3 \
WHERE parent_action_key = ?1 AND child_action_key = ?2",
&[
SqlValue::Text(parent),
SqlValue::Text(child),
SqlValue::Text(state),
],
)?;
if affected == 0 {
return Err(StoreError::UnknownPin);
}
Ok(())
})
}
fn list_provisional_obligations_for_pin(
&mut self,
pin_key: &str,
) -> Result<Vec<ProvisionalObligationRow>, StoreError> {
let rows = self.engine.query(
"SELECT consumer_worker, consumer_attempt_hex, pin_key, producer_action_key, \
producer_generation_hex, producer_attempt_hex, role, virtual_path, object_key, \
status, resolution_object_key, created_seq FROM provisional_obligations \
WHERE pin_key = ?1 ORDER BY consumer_worker, consumer_attempt_hex",
&[SqlValue::Text(pin_key.to_owned())],
)?;
rows.iter()
.map(|row| Self::map_obligation_row(row))
.collect()
}
}
pub struct RusqliteEngine {
conn: rusqlite::Connection,
}
impl RusqliteEngine {
pub fn open(path: &std::path::Path) -> Result<Self, StoreError> {
rusqlite::Connection::open(path)
.map(|conn| Self { conn })
.map_err(|e| StoreError::Backend(e.to_string()))
}
pub fn open_in_memory() -> Result<Self, StoreError> {
rusqlite::Connection::open_in_memory()
.map(|conn| Self { conn })
.map_err(|e| StoreError::Backend(e.to_string()))
}
}
fn to_rusqlite(v: &SqlValue) -> rusqlite::types::Value {
match v {
SqlValue::Null => rusqlite::types::Value::Null,
SqlValue::Int(i) => rusqlite::types::Value::Integer(*i),
SqlValue::Text(t) => rusqlite::types::Value::Text(t.clone()),
SqlValue::Blob(b) => rusqlite::types::Value::Blob(b.clone()),
}
}
impl SqlEngine for RusqliteEngine {
fn execute(&mut self, sql: &str, params: &[SqlValue]) -> Result<usize, StoreError> {
let bound = rusqlite::params_from_iter(params.iter().map(to_rusqlite));
self.conn
.execute(sql, bound)
.map_err(|e| StoreError::Backend(e.to_string()))
}
fn query(&mut self, sql: &str, params: &[SqlValue]) -> Result<Vec<Vec<SqlValue>>, StoreError> {
let mut statement = self
.conn
.prepare(sql)
.map_err(|e| StoreError::Backend(e.to_string()))?;
let column_count = statement.column_count();
let bound = rusqlite::params_from_iter(params.iter().map(to_rusqlite));
let mut rows = statement
.query(bound)
.map_err(|e| StoreError::Backend(e.to_string()))?;
let mut out = Vec::new();
while let Some(row) = rows
.next()
.map_err(|e| StoreError::Backend(e.to_string()))?
{
let mut cols = Vec::with_capacity(column_count);
for i in 0..column_count {
let value: rusqlite::types::Value =
row.get(i).map_err(|e| StoreError::Backend(e.to_string()))?;
cols.push(match value {
rusqlite::types::Value::Null => SqlValue::Null,
rusqlite::types::Value::Integer(v) => SqlValue::Int(v),
rusqlite::types::Value::Real(f) => {
return Err(StoreError::Corruption(format!(
"unexpected float column {f}"
)));
}
rusqlite::types::Value::Text(t) => SqlValue::Text(t),
rusqlite::types::Value::Blob(b) => SqlValue::Blob(b),
});
}
out.push(cols);
}
Ok(out)
}
}
pub struct FsqliteEngine {
conn: fsqlite::AsyncConnection,
}
impl FsqliteEngine {
pub fn open(path: &std::path::Path) -> Result<Self, StoreError> {
fsqlite::AsyncConnection::open_sync(path.display().to_string())
.map(|conn| Self { conn })
.map_err(|e| StoreError::Backend(e.to_string()))
}
}
impl Drop for FsqliteEngine {
fn drop(&mut self) {
if let Err(error) = self.conn.close_without_checkpoint_sync() {
use std::io::Write as _;
let _ = writeln!(std::io::stderr(), "FrankenSQLite close failed: {error}");
}
}
}
fn to_fsqlite(v: &SqlValue) -> fsqlite::SqliteValue {
match v {
SqlValue::Null => fsqlite::SqliteValue::Null,
SqlValue::Int(i) => fsqlite::SqliteValue::Integer(*i),
SqlValue::Text(t) => fsqlite::SqliteValue::Text(t.as_str().into()),
SqlValue::Blob(b) => fsqlite::SqliteValue::Blob(b.clone().into()),
}
}
fn from_fsqlite(v: &fsqlite::SqliteValue) -> Result<SqlValue, StoreError> {
match v {
fsqlite::SqliteValue::Null => Ok(SqlValue::Null),
fsqlite::SqliteValue::Integer(i) => Ok(SqlValue::Int(*i)),
fsqlite::SqliteValue::Float(f) => Err(StoreError::Corruption(format!(
"unexpected float column {f}"
))),
fsqlite::SqliteValue::Text(t) => Ok(SqlValue::Text(t.as_str().to_owned())),
fsqlite::SqliteValue::Blob(b) => Ok(SqlValue::Blob(b.as_ref().to_vec())),
}
}
impl SqlEngine for FsqliteEngine {
fn execute(&mut self, sql: &str, params: &[SqlValue]) -> Result<usize, StoreError> {
if params.is_empty() {
return self
.conn
.execute_sync(sql)
.map_err(|e| StoreError::Backend(e.to_string()));
}
let bound: Vec<fsqlite::SqliteValue> = params.iter().map(to_fsqlite).collect();
self.conn
.execute_with_params_sync(sql, &bound)
.map_err(|e| StoreError::Backend(e.to_string()))
}
fn query(&mut self, sql: &str, params: &[SqlValue]) -> Result<Vec<Vec<SqlValue>>, StoreError> {
let bound: Vec<fsqlite::SqliteValue> = params.iter().map(to_fsqlite).collect();
let rows = self
.conn
.query_with_params_sync(sql, &bound)
.map_err(|e| StoreError::Backend(e.to_string()))?;
rows.iter()
.map(|row| row.values().iter().map(from_fsqlite).collect())
.collect()
}
}
#[cfg(test)]
mod tests {
use super::*;
use rabs_protocol::authority::{ClusterId, CoordinatorAuthority, CoordinatorIncarnationId};
use rabs_protocol::generation::{
ActionGenerationId, AttemptId, ExecutionLeaseId, LeaseRenewalSeq,
};
use std::sync::atomic::{AtomicU64, Ordering};
static DB_COUNTER: AtomicU64 = AtomicU64::new(0);
fn fresh_path(tag: &str) -> std::path::PathBuf {
let n = DB_COUNTER.fetch_add(1, Ordering::SeqCst);
std::env::temp_dir().join(format!("rabs-h009-{}-{}-{}.db", std::process::id(), tag, n))
}
#[test]
fn fsqlite_drop_finishes_rollback_before_immediate_reopen() {
let path = fresh_path("fsqlite-drop-rollback");
assert!(!path.exists());
{
let mut engine = FsqliteEngine::open(&path).unwrap();
engine
.execute("CREATE TABLE lifecycle (id INTEGER PRIMARY KEY)", &[])
.unwrap();
engine
.execute("INSERT INTO lifecycle VALUES (?1)", &[SqlValue::Int(1)])
.unwrap();
engine.execute("BEGIN", &[]).unwrap();
engine
.execute("INSERT INTO lifecycle VALUES (?1)", &[SqlValue::Int(2)])
.unwrap();
}
let mut reopened = FsqliteEngine::open(&path).unwrap();
assert_eq!(
reopened
.query("SELECT id FROM lifecycle ORDER BY id", &[])
.unwrap(),
vec![vec![SqlValue::Int(1)]]
);
reopened
.execute("INSERT INTO lifecycle VALUES (?1)", &[SqlValue::Int(3)])
.unwrap();
assert_eq!(
reopened
.query("SELECT id FROM lifecycle ORDER BY id", &[])
.unwrap(),
vec![vec![SqlValue::Int(1)], vec![SqlValue::Int(3)]]
);
}
#[test]
fn fsqlite_drop_preserves_committed_wal_before_reopen() {
let path = fresh_path("fsqlite-drop-wal");
assert!(!path.exists());
let mut wal_path = path.as_os_str().to_os_string();
wal_path.push("-wal");
let wal_path = std::path::PathBuf::from(wal_path);
assert!(!wal_path.exists());
let wal_before;
{
let mut engine = FsqliteEngine::open(&path).unwrap();
engine.query("PRAGMA journal_mode=WAL", &[]).unwrap();
engine.execute("PRAGMA wal_autocheckpoint=0", &[]).unwrap();
engine
.execute("CREATE TABLE lifecycle (id INTEGER PRIMARY KEY)", &[])
.unwrap();
engine
.execute("INSERT INTO lifecycle VALUES (?1)", &[SqlValue::Int(7)])
.unwrap();
wal_before = std::fs::read(&wal_path).unwrap();
assert!(wal_before.len() > 32);
}
assert_eq!(std::fs::read(&wal_path).unwrap(), wal_before);
let mut reopened = FsqliteEngine::open(&path).unwrap();
assert_eq!(
reopened.query("SELECT id FROM lifecycle", &[]).unwrap(),
vec![vec![SqlValue::Int(7)]]
);
}
struct MigrationProbe<E> {
inner: E,
executed: Vec<String>,
fail_at: Option<usize>,
}
impl<E: SqlEngine> SqlEngine for MigrationProbe<E> {
fn execute(&mut self, sql: &str, params: &[SqlValue]) -> Result<usize, StoreError> {
let index = self.executed.len();
self.executed.push(sql.to_owned());
if self.fail_at == Some(index) {
return Err(StoreError::Backend("injected migration failure".into()));
}
self.inner.execute(sql, params)
}
fn query(
&mut self,
sql: &str,
params: &[SqlValue],
) -> Result<Vec<Vec<SqlValue>>, StoreError> {
self.inner.query(sql, params)
}
}
fn migration_probe<E>(inner: E, fail_at: Option<usize>) -> MigrationProbe<E> {
MigrationProbe {
inner,
executed: Vec::new(),
fail_at,
}
}
fn assert_one_migration_commit<E: SqlEngine>(store: &mut SqlMetadataStore<MigrationProbe<E>>) {
assert_eq!(store.schema_version().unwrap(), SCHEMA_VERSION);
let statements = &store.engine.executed;
assert_eq!(statements.iter().filter(|sql| *sql == "BEGIN").count(), 1);
assert_eq!(statements.iter().filter(|sql| *sql == "COMMIT").count(), 1);
assert!(!statements.iter().any(|sql| sql == "ROLLBACK"));
assert_eq!(
store
.engine
.query("SELECT version FROM schema_epochs ORDER BY version", &[])
.unwrap(),
(1..=SCHEMA_VERSION)
.map(|version| vec![SqlValue::Int(i64::from(version))])
.collect::<Vec<_>>()
);
}
fn migration_batch_behavior<E: SqlEngine>(open: fn(&std::path::Path) -> Result<E, StoreError>) {
let fresh = fresh_path("migration-single-commit");
let mut store =
SqlMetadataStore::open(migration_probe(open(&fresh).unwrap(), None)).unwrap();
assert_one_migration_commit(&mut store);
drop(store);
let store = SqlMetadataStore::open(migration_probe(open(&fresh).unwrap(), None)).unwrap();
assert!(
store.engine.executed.is_empty(),
"current-schema reopen must not write"
);
drop(store);
for version in [0, 19, 20] {
let pending_writes: usize = MIGRATIONS
.iter()
.filter(|migration| migration.version > version)
.map(|migration| migration.statements.len() + 1)
.sum();
for fail_at in [pending_writes - 1, pending_writes, pending_writes + 1] {
let path = fresh_path("migration-rollback");
let mut engine = open(&path).unwrap();
match version {
19 => seed_v19_worker_fence(&mut engine),
20 => seed_v20_attempt_lease(&mut engine),
_ => {}
}
let queries = [
"SELECT * FROM worker_incarnation_fences ORDER BY worker",
"SELECT * FROM action_generations ORDER BY id_hex",
"SELECT * FROM action_attempts ORDER BY id_hex",
];
let prior_rows = if version == 0 {
Vec::new()
} else {
queries
.iter()
.map(|sql| engine.query(sql, &[]).unwrap())
.collect::<Vec<_>>()
};
let mut failed = SqlMetadataStore {
engine: migration_probe(engine, Some(fail_at)),
domains: HashMap::new(),
};
assert_eq!(
failed.apply_migrations(),
Err(StoreError::Backend("injected migration failure".into()))
);
assert_eq!(
failed.engine.executed.last().map(String::as_str),
Some("ROLLBACK")
);
failed.engine.execute("BEGIN", &[]).unwrap();
failed.engine.execute("ROLLBACK", &[]).unwrap();
drop(failed);
let mut reopened = open(&path).unwrap();
if version == 0 {
assert!(
reopened
.query(
"SELECT name FROM sqlite_master WHERE name = 'schema_epochs'",
&[]
)
.unwrap()
.is_empty()
);
} else {
assert_eq!(
reopened
.query("SELECT MAX(version) FROM schema_epochs", &[])
.unwrap(),
vec![vec![SqlValue::Int(i64::from(version))]]
);
for (sql, expected) in queries.iter().zip(&prior_rows) {
assert_eq!(
&reopened.query(sql, &[]).unwrap(),
expected,
"version {version}, failure {fail_at}, query {sql}"
);
}
}
let mut retried = SqlMetadataStore::open(migration_probe(reopened, None)).unwrap();
assert_one_migration_commit(&mut retried);
if version == 19 {
assert_v19_worker_fence_migrated(&mut retried);
} else if version == 20 {
assert_v20_attempt_lease_fails_closed(&mut retried);
}
}
}
}
#[test]
fn pending_migrations_commit_once_and_rollback_reference() {
migration_batch_behavior(RusqliteEngine::open);
}
#[test]
fn pending_migrations_commit_once_and_rollback_frankensqlite() {
migration_batch_behavior(FsqliteEngine::open);
}
fn digest(domain: &'static str, tag: u8) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes: [tag; 32],
}
}
fn authority(tag: u8) -> AuthorityRow {
AuthorityRow {
digest: digest("rabs.authority.sha256.v1", tag),
cluster_id: "cluster-a".to_owned(),
incarnation: u128::from(tag),
term: u64::from(tag),
acquired_seq: 1,
}
}
fn coordinator_authority() -> CoordinatorAuthority {
CoordinatorAuthority {
cluster_id: ClusterId("cluster-a".to_owned()),
credential_generation: 1,
term: 1,
incarnation_id: CoordinatorIncarnationId(1),
}
}
fn bound_authority_row() -> AuthorityRow {
let coordinator = coordinator_authority();
AuthorityRow {
digest: rabs_key::authority_binding::coordinator_authority_digest(&coordinator),
cluster_id: "cluster-a".to_owned(),
incarnation: 1,
term: 1,
acquired_seq: 1,
}
}
fn bound_attempt_authority() -> AttemptAuthority {
let coordinator = coordinator_authority();
let created_under = rabs_key::authority_binding::coordinator_authority_digest(&coordinator);
AttemptAuthority {
coordinator,
action_key: digest("rabs.action-key.sha256.v1", 7),
action_generation: ActionGeneration {
generation_id: ActionGenerationId(11),
per_key_ordinal: 1,
created_under_authority_digest: created_under,
},
attempt_id: AttemptId(22),
execution_lease_id: ExecutionLeaseId(30),
lease_renewal_seq: LeaseRenewalSeq(1),
worker_peer_id: PeerId("lease-worker".to_owned()),
worker_boot_generation: WorkerBootGeneration(3),
worker_incarnation_id: WorkerIncarnationId(99),
}
}
fn worker_offer(worker: &str, generation: u64, incarnation: u128) -> WorkerSessionOffer {
WorkerSessionOffer {
worker_peer_id: PeerId(worker.to_owned()),
boot_generation: WorkerBootGeneration(generation),
incarnation: WorkerIncarnationId(incarnation),
reenrollment_proof: None,
}
}
fn publication(action_tag: u8, descriptor_tag: u8, pin: u128) -> PublicationRow {
PublicationRow {
action_key: digest("rabs.action-key.sha256.v1", action_tag),
descriptor_digest: digest("rabs.descriptor.sha256.v1", descriptor_tag),
manifest_digest: digest("rabs.result-manifest.sha256.v1", descriptor_tag),
evidence_digest: digest("rabs.evidence-bundle.sha256.v1", descriptor_tag),
winner_generation: 10,
winner_attempt: 20,
result_kind: ResultKindTag::Success,
pin_id: pin,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: Vec::new(),
}
}
fn behavioral_pass(store: &mut dyn RabsMetadataStore) -> Vec<String> {
assert_eq!(store.schema_version().unwrap(), SCHEMA_VERSION);
let bound_authority = bound_authority_row();
store.acquire_authority(&bound_authority).unwrap();
store.acquire_authority(&bound_authority).unwrap(); assert!(matches!(
store.acquire_authority(&authority(2)),
Err(StoreError::AuthorityHeld { .. })
));
assert_eq!(store.active_authority().unwrap().unwrap(), bound_authority);
let action = ActionEntryRow {
action_key: digest("rabs.action-key.sha256.v1", 7),
key_epoch: 3,
projection_epoch: 4,
};
store.upsert_action_entry(&action).unwrap();
assert_eq!(
store.lookup_action(&action.action_key).unwrap(),
Some(action.clone())
);
assert_eq!(
store
.lookup_action(&digest("rabs.action-key.sha256.v1", 99))
.unwrap(),
None
);
let active = bound_authority_row().digest;
let wrong = digest("rabs.authority.sha256.v1", 2);
assert_eq!(
store.create_generation(&wrong, 10, &action.action_key),
Err(StoreError::NotActiveAuthority)
);
store
.create_generation(&active, 10, &action.action_key)
.unwrap();
store.tombstone_generation(10).unwrap();
assert_eq!(
store.create_generation(&active, 10, &action.action_key),
Err(StoreError::GenerationIdNotAboveHighWater)
);
assert_eq!(
store.create_generation(&active, 9, &action.action_key),
Err(StoreError::GenerationIdNotAboveHighWater)
);
let attempt_authority = bound_attempt_authority();
store
.create_bound_generation(
&active,
&attempt_authority.action_generation,
&action.action_key,
)
.unwrap();
store.record_attempt(20, 11, "worker-a", 5).unwrap();
assert_eq!(
store.record_attempt(20, 11, "worker-a", 6),
Err(StoreError::DuplicateAttempt)
);
assert_eq!(
store.record_attempt(21, 999, "worker-a", 6),
Err(StoreError::UnknownGeneration)
);
store
.admit_worker_session(
&active,
&WorkerSessionOffer {
worker_peer_id: attempt_authority.worker_peer_id.clone(),
boot_generation: attempt_authority.worker_boot_generation,
incarnation: attempt_authority.worker_incarnation_id,
reenrollment_proof: None,
},
4,
)
.unwrap();
store
.admit_attempt_lease(&attempt_authority, 5, 100)
.unwrap();
let renewal = LeaseRenewal {
lease: attempt_authority.execution_lease_id,
seq: LeaseRenewalSeq(2),
};
store
.renew_attempt_lease(&attempt_authority, renewal, 200, &|| 10)
.unwrap();
assert_eq!(
store.renew_attempt_lease(&attempt_authority, renewal, 300, &|| 10),
Err(StoreError::LeaseRenewalMismatch)
);
let mut current_authority = attempt_authority.clone();
current_authority.lease_renewal_seq = LeaseRenewalSeq(2);
assert_eq!(
store.renew_attempt_lease(
¤t_authority,
LeaseRenewal {
lease: current_authority.execution_lease_id,
seq: LeaseRenewalSeq(2),
},
300,
&|| 10,
),
Err(StoreError::NonMonotonicRenewal)
);
let mut unknown = current_authority.clone();
unknown.execution_lease_id = rabs_protocol::generation::ExecutionLeaseId(31);
assert_eq!(
store.renew_attempt_lease(
&unknown,
LeaseRenewal {
lease: unknown.execution_lease_id,
seq: LeaseRenewalSeq(3),
},
300,
&|| 10,
),
Err(StoreError::UnknownLease)
);
store.release_lease(30).unwrap();
assert_eq!(
store.renew_attempt_lease(
¤t_authority,
LeaseRenewal {
lease: current_authority.execution_lease_id,
seq: LeaseRenewalSeq(3),
},
300,
&|| 10,
),
Err(StoreError::LeaseReleased)
);
let publication_row = publication(7, 1, 40);
assert_eq!(
store.commit_publication(&wrong, None, &publication_row),
Err(StoreError::NotActiveAuthority)
);
assert_eq!(
store
.commit_publication(&active, None, &publication_row)
.unwrap(),
CommitOutcome::Committed
);
assert_eq!(
store
.commit_publication(&active, None, &publication_row)
.unwrap(),
CommitOutcome::IdempotentDuplicate
);
assert_eq!(
store
.commit_publication(&active, None, &publication(7, 2, 41))
.unwrap(),
CommitOutcome::ConflictQuarantined
);
assert!(store.has_publication(&publication_row.action_key).unwrap());
assert!(
!store
.has_publication(&digest("rabs.action-key.sha256.v1", 99))
.unwrap()
);
let extra_evidence = digest("rabs.evidence-bundle.sha256.v1", 90);
let manifest_key = digest_key(&publication_row.manifest_digest);
store
.append_evidence(
&publication_row.action_key,
&manifest_key,
&extra_evidence,
11,
20,
)
.unwrap();
store
.append_evidence(
&publication_row.action_key,
"rabs.canonical-result-manifest.sha256.v1:ff",
&extra_evidence,
99,
98,
)
.unwrap();
let evidence_line = store
.differential_snapshot()
.unwrap()
.into_iter()
.find(|l| l.starts_with("action_evidence_index|") && l.contains("|5a5a"))
.expect("extra evidence row present");
assert!(
evidence_line.contains(&u128_hex(11)) && evidence_line.contains(&u128_hex(20)),
"first-writer attribution preserved: {evidence_line}"
);
assert!(
!evidence_line.ends_with(":ff"),
"re-append must not rebind the manifest: {evidence_line}"
);
let for_manifest = store
.list_evidence_keys_for_manifest(&manifest_key)
.unwrap();
assert!(for_manifest.contains(&digest_key(&extra_evidence)));
assert!(
store
.list_evidence_keys_for_manifest("rabs.canonical-result-manifest.sha256.v1:ff")
.unwrap()
.is_empty()
);
let object = digest("rabs.object.sha256.v1", 50);
store.record_object(&object, 4096).unwrap();
store
.add_location(&object, "/cas/aa/bb", Some(7), "raw", true)
.unwrap();
store
.add_location(&object, "/cas/cc/dd", None, "zstd", false)
.unwrap();
store
.create_pin(
60,
&object,
"operator",
"administrative",
Some(500),
Some("manual hold"),
false,
"operator investigation",
)
.unwrap();
assert_eq!(
store.release_pin(60, "someone-else"),
Err(StoreError::PinOwnerMismatch)
);
assert_eq!(
store.release_pin(61, "operator"),
Err(StoreError::UnknownPin)
);
store.renew_pin(60, 2).unwrap();
assert_eq!(
store.renew_pin(60, 2),
Err(StoreError::NonMonotonicPinRenewal)
);
assert_eq!(store.renew_pin(62, 1), Err(StoreError::UnknownPin));
let child = digest("rabs.object.sha256.v1", 51);
let grandchild = digest("rabs.object.sha256.v1", 52);
let orphan = digest("rabs.object.sha256.v1", 53);
store.record_object(&child, 16).unwrap();
store.record_object(&grandchild, 16).unwrap();
store.record_object(&orphan, 16).unwrap();
store
.add_object_edge(&object, &child, "manifest-entry")
.unwrap();
store.add_object_edge(&child, &grandchild, "chunk").unwrap();
store
.add_object_edge(&grandchild, &object, "back-ref")
.unwrap();
store
.set_location_quarantined(&object, "/cas/aa/bb", true)
.unwrap();
store
.add_location(&object, "/cas/aa/bb", Some(8), "raw", true)
.unwrap();
assert!(
store.object_located(&object).unwrap(),
"second clean copy remains"
);
assert!(
!store.object_durably_located(&object).unwrap(),
"refreshing the durable copy must not release its quarantine"
);
store
.set_location_quarantined(&object, "/cas/cc/dd", true)
.unwrap();
store
.add_location(&object, "/cas/cc/dd", None, "zstd", false)
.unwrap();
assert!(
!store.object_located(&object).unwrap(),
"all copies quarantined: object unavailable"
);
store
.set_location_quarantined(&object, "/cas/aa/bb", false)
.unwrap();
assert!(store.object_located(&object).unwrap());
store
.put_recipe(&action.action_key, &digest("rabs.recipe.sha256.v1", 70))
.unwrap();
store
.put_key_breakdown(
&action.action_key,
"toolchain",
&digest("rabs.component.sha256.v1", 71),
)
.unwrap();
store
.set_trust(&action.action_key, "trusted", "verified twice")
.unwrap();
store
.add_quarantine(QuarantineScope::Location, "/cas/cc/dd", "checksum mismatch")
.unwrap();
store
.record_verification_sample(&action.action_key, 20, true, 8)
.unwrap();
let snapshot = store.gc_snapshot(9).unwrap();
assert_eq!(snapshot.pinned_roots.len(), 2);
assert_eq!(snapshot.located_objects, vec![digest_key(&object)]);
assert_eq!(snapshot.reachable_from_pins.len(), 4);
assert!(snapshot.reachable_from_pins.contains(&digest_key(&object)));
assert!(snapshot.reachable_from_pins.contains(&digest_key(&child)));
assert!(
snapshot
.reachable_from_pins
.contains(&digest_key(&grandchild))
);
assert!(!snapshot.reachable_from_pins.contains(&digest_key(&orphan)));
let scan = store.reconciliation_scan().unwrap();
assert_eq!(scan.len(), 2);
assert_eq!(scan[0].verified_seq, Some(8));
assert_eq!(scan[0].encoding, "raw");
assert!(!scan[0].quarantined);
assert_eq!(scan[1].verified_seq, None);
assert_eq!(scan[1].encoding, "zstd");
assert!(scan[1].quarantined);
assert_eq!(
store.record_peer_authority_high_water(&wrong, "peer-1", 2, 5, 11, 100),
Err(StoreError::NotActiveAuthority)
);
store
.record_peer_authority_high_water(&active, "peer-1", 2, 5, 11, 100)
.unwrap();
store
.record_peer_authority_high_water(&active, "peer-1", 2, 5, 11, 101)
.unwrap();
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 2, 5, 12, 102),
Err(StoreError::StalePeerAuthority)
);
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((2, 5, 100)),
"a conflicting incarnation must not change the durable mark"
);
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 2, 4, 11, 102),
Err(StoreError::StalePeerAuthority)
);
store
.record_peer_authority_high_water(&active, "peer-1", 2, 6, 12, 103)
.unwrap();
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 2, 6, 11, 104),
Err(StoreError::StalePeerAuthority)
);
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((2, 6, 103))
);
assert_eq!(store.peer_authority_high_water("peer-2").unwrap(), None);
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 1, 999, 12, 105),
Err(StoreError::StalePeerAuthority),
"a larger term cannot revive an older credential generation"
);
store
.record_peer_authority_high_water(&active, "peer-1", 3, 1, 13, 106)
.unwrap();
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((3, 1, 106)),
"a new credential generation opens a new term namespace"
);
let first = worker_offer("worker-a", 3, 11);
assert_eq!(
store.admit_worker_session(&wrong, &first, 1_100),
Err(StoreError::NotActiveAuthority)
);
assert_eq!(
store.admit_worker_session(&active, &first, 1_100),
Ok(WorkerAdmission::AdmitNewGeneration)
);
assert_eq!(
store.admit_worker_session(&active, &first, 1_100),
Ok(WorkerAdmission::AdmitReconnect),
"same session admission is idempotent"
);
assert_eq!(
store.admit_worker_session(&active, &first, 1_101),
Ok(WorkerAdmission::AdmitReconnect),
"one incarnation may hold multiple reconnect sessions"
);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 2, u128::MAX), 1_102,),
Ok(WorkerAdmission::RejectStaleBootGeneration)
);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 3, 22), 1_103),
Ok(WorkerAdmission::RejectCloneAmbiguity)
);
assert!(
!store
.release_worker_session(
&active,
&PeerId("worker-a".into()),
WorkerIncarnationId(22),
1_100,
1_104,
)
.unwrap()
);
assert!(
store
.release_worker_session(
&active,
&PeerId("worker-a".into()),
WorkerIncarnationId(11),
1_100,
1_104,
)
.unwrap()
);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 3, 22), 1_105),
Ok(WorkerAdmission::RejectCloneAmbiguity),
"releasing one reconnect must not clear another open session's fence"
);
assert!(
store
.release_worker_session(
&active,
&PeerId("worker-a".into()),
WorkerIncarnationId(11),
1_101,
1_106,
)
.unwrap()
);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 3, 22), 1_107),
Ok(WorkerAdmission::RejectCloneAmbiguity),
"session release cannot adjudicate a detected clone"
);
let mut rolled_back = worker_offer("worker-a", 1, 99);
rolled_back.reenrollment_proof = Some(1);
assert_eq!(
store.admit_worker_session(&active, &rolled_back, 1_108),
Ok(WorkerAdmission::RejectStaleBootGeneration),
"operator proof cannot lower the durable global high-water"
);
let mut reenrolled = worker_offer("worker-a", 3, 99);
reenrolled.reenrollment_proof = Some(1);
assert_eq!(
store.admit_worker_session(&active, &reenrolled, 1_109),
Ok(WorkerAdmission::AdmitViaReenrollment)
);
let mut replay = worker_offer("worker-a", 3, 100);
replay.reenrollment_proof = Some(1);
assert_eq!(
store.admit_worker_session(&active, &replay, 1_110),
Ok(WorkerAdmission::RejectCloneAmbiguity),
"a consumed operator proof cannot pick a second clone"
);
assert_eq!(
store.worker_incarnation_fence(&PeerId("worker-a".into())),
Ok(Some(WorkerIncarnationFenceRecord {
worker_peer_id: PeerId("worker-a".into()),
highest_boot_generation: WorkerBootGeneration(3),
active_incarnation: Some(WorkerIncarnationId(99)),
clone_ambiguous: true,
operator_reenrollment_generation: 1,
}))
);
assert_eq!(
store.worker_incarnation_fence(&PeerId("worker-b".into())),
Ok(None)
);
let max_generation = worker_offer("worker-max", u64::MAX, u128::MAX);
assert_eq!(
store.admit_worker_session(&active, &max_generation, 1_111),
Ok(WorkerAdmission::AdmitNewGeneration)
);
assert_eq!(
store
.worker_incarnation_fence(&PeerId("worker-max".into()))
.unwrap()
.unwrap()
.highest_boot_generation,
WorkerBootGeneration(u64::MAX),
);
store.advance_edge_fence(&active, "edge-1", 10).unwrap();
assert_eq!(
store.advance_edge_fence(&active, "edge-1", 9),
Err(StoreError::StaleEdgeIncarnation)
);
assert_eq!(
store.begin_edge_handoff(&active, "edge-1", 11, 9, 200),
Err(StoreError::EdgeHandoffPredecessorMismatch)
);
assert_eq!(
store.begin_edge_handoff(&active, "edge-1", 10, 10, 200),
Err(StoreError::StaleEdgeIncarnation)
);
store
.begin_edge_handoff(&active, "edge-1", 11, 10, 200)
.unwrap();
store
.begin_edge_handoff(&active, "edge-1", 11, 10, 200)
.unwrap(); assert_eq!(
store.begin_edge_handoff(&active, "edge-1", 12, 10, 201),
Err(StoreError::EdgeHandoffActive)
);
assert_eq!(
store.active_edge_handoff("edge-1").unwrap(),
Some(EdgeHandoffRow {
active_incarnation: 11,
predecessor_incarnation: 10,
begun_seq: 200,
})
);
assert_eq!(
store.edge_fence("edge-1").unwrap(),
Some(10),
"fence advances only at resolve"
);
assert_eq!(
store.resolve_edge_handoff(&active, "edge-1", 12),
Err(StoreError::UnknownEdgeHandoff)
);
store.resolve_edge_handoff(&active, "edge-1", 11).unwrap();
assert_eq!(store.edge_fence("edge-1").unwrap(), Some(11));
assert_eq!(store.active_edge_handoff("edge-1").unwrap(), None);
assert_eq!(
store.begin_edge_handoff(&active, "edge-1", 12, 10, 202),
Err(StoreError::EdgeHandoffPredecessorMismatch)
);
store
.begin_edge_handoff(&active, "edge-1", 12, 11, 202)
.unwrap();
let evaluation = TrustEvaluationRow {
version: 1,
state: "trusted".to_owned(),
reason: "verified once".to_owned(),
evaluated_seq: 300,
};
assert_eq!(
store.append_trust_evaluation(&wrong, &action.action_key, &evaluation),
Err(StoreError::NotActiveAuthority)
);
store
.append_trust_evaluation(&active, &action.action_key, &evaluation)
.unwrap();
assert_eq!(
store.append_trust_evaluation(&active, &action.action_key, &evaluation),
Err(StoreError::NonMonotonicTrustEvaluation)
);
let demotion = TrustEvaluationRow {
version: 2,
state: "suspect".to_owned(),
reason: "divergent recompute".to_owned(),
evaluated_seq: 301,
};
store
.append_trust_evaluation(&active, &action.action_key, &demotion)
.unwrap();
assert_eq!(
store.latest_trust_evaluation(&action.action_key).unwrap(),
Some(demotion)
);
assert_eq!(
store.create_operation(&wrong, 400, "transfer", "running", 310),
Err(StoreError::NotActiveAuthority)
);
store
.create_operation(&active, 400, "transfer", "running", 310)
.unwrap();
assert_eq!(
store.create_operation(&active, 400, "transfer", "running", 311),
Err(StoreError::DuplicateOperation)
);
store.update_operation_state(400, "done", 312).unwrap();
assert_eq!(
store.update_operation_state(401, "done", 312),
Err(StoreError::UnknownOperation)
);
assert_eq!(store.operation_state(400).unwrap(), Some("done".to_owned()));
store
.register_edge_subscriber("edge-1", "sub-a", 320)
.unwrap();
store
.register_edge_subscriber("edge-1", "sub-a", 999)
.unwrap(); store
.register_edge_subscriber("edge-1", "sub-b", 321)
.unwrap();
assert_eq!(
store.list_edge_subscribers("edge-1").unwrap(),
vec!["sub-a".to_owned(), "sub-b".to_owned()]
);
assert!(store.remove_edge_subscriber("edge-1", "sub-b").unwrap());
assert!(!store.remove_edge_subscriber("edge-1", "sub-b").unwrap());
let manifest = digest("rabs.result-manifest.sha256.v1", 80);
store.record_manifest(&manifest, "tree", 12).unwrap();
store.record_manifest(&manifest, "tree", 12).unwrap();
assert_eq!(
store.record_manifest(&manifest, "tree", 13),
Err(StoreError::ManifestDivergence)
);
assert_eq!(
store.manifest_meta(&manifest).unwrap(),
Some(("tree".to_owned(), 12))
);
store.record_worker_session("worker-a", 3, 330).unwrap();
store.record_worker_session("worker-a", 3, 330).unwrap();
assert_eq!(
store.record_worker_session("worker-a", 4, 330),
Err(StoreError::AppendConflict("worker_sessions".into()))
);
assert!(store.end_worker_session("worker-a", 330, 340).unwrap());
assert!(!store.end_worker_session("worker-a", 330, 341).unwrap());
store.record_worker_capability("worker-a", "nix").unwrap();
store.record_worker_capability("worker-a", "nix").unwrap();
store.record_worker_capability("worker-a", "cargo").unwrap();
assert_eq!(
store.list_worker_capabilities("worker-a").unwrap(),
vec!["cargo".to_owned(), "nix".to_owned()]
);
store
.record_worker_health_sample("worker-a", 350, true, "ok")
.unwrap();
store
.record_worker_health_sample("worker-a", 350, true, "ok")
.unwrap();
assert_eq!(
store.record_worker_health_sample("worker-a", 350, false, "ok"),
Err(StoreError::AppendConflict("worker_health_samples".into()))
);
store
.record_decision_receipt("gc", "run-1", 360, "reclaim", "under pressure")
.unwrap();
store
.record_decision_receipt("gc", "run-1", 360, "reclaim", "under pressure")
.unwrap();
assert_eq!(
store.record_decision_receipt("gc", "run-1", 360, "skip", "under pressure"),
Err(StoreError::AppendConflict("decision_receipts".into()))
);
store
.add_provenance_edge(&object, &child, "derived-from")
.unwrap();
store
.add_provenance_edge(&object, &child, "derived-from")
.unwrap();
store
.record_determinism_audit(&action.action_key, 20, 370, "deterministic")
.unwrap();
store
.record_determinism_audit(&action.action_key, 20, 370, "deterministic")
.unwrap();
assert_eq!(
store.record_determinism_audit(&action.action_key, 20, 370, "divergent"),
Err(StoreError::AppendConflict("determinism_audits".into()))
);
store
.create_materialization(500, &object, "/work/out", "staging", 380)
.unwrap();
store
.create_materialization(500, &object, "/work/out", "staging", 380)
.unwrap();
assert_eq!(
store.create_materialization(500, &object, "/work/other", "staging", 381),
Err(StoreError::AppendConflict("materialization_records".into()))
);
store
.update_materialization_state(500, "complete", 382)
.unwrap();
assert_eq!(
store.update_materialization_state(501, "complete", 382),
Err(StoreError::UnknownMaterialization)
);
assert_eq!(
store.materialization_state(500).unwrap(),
Some("complete".to_owned())
);
store.release_authority(&active).unwrap();
assert_eq!(store.active_authority().unwrap(), None);
store.acquire_authority(&authority(2)).unwrap();
store.differential_snapshot().unwrap()
}
#[test]
fn h009_reference_backend_full_behavioral_pass() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
let snapshot = behavioral_pass(&mut store);
assert!(!snapshot.is_empty());
}
#[test]
fn h009_frankensqlite_backend_full_behavioral_pass() {
let engine = FsqliteEngine::open(&fresh_path("fsqlite-pass")).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
let snapshot = behavioral_pass(&mut store);
assert!(!snapshot.is_empty());
}
#[test]
fn h009_differential_harness_reference_vs_frankensqlite() {
let reference_engine = RusqliteEngine::open(&fresh_path("ref-diff")).unwrap();
let mut reference = SqlMetadataStore::open(reference_engine).unwrap();
let candidate_engine = FsqliteEngine::open(&fresh_path("fsq-diff")).unwrap();
let mut candidate = SqlMetadataStore::open(candidate_engine).unwrap();
let reference_snapshot = behavioral_pass(&mut reference);
let candidate_snapshot = behavioral_pass(&mut candidate);
assert_eq!(reference_snapshot, candidate_snapshot);
}
#[test]
fn h009_generation_high_water_survives_reopen() {
let path = fresh_path("reopen");
{
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
store.acquire_authority(&authority(1)).unwrap();
let action = ActionEntryRow {
action_key: digest("rabs.action-key.sha256.v1", 7),
key_epoch: 0,
projection_epoch: 0,
};
store.upsert_action_entry(&action).unwrap();
store
.create_generation(
&digest("rabs.authority.sha256.v1", 1),
42,
&action.action_key,
)
.unwrap();
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_eq!(store.schema_version().unwrap(), SCHEMA_VERSION);
store.acquire_authority(&authority(1)).unwrap();
assert_eq!(
store.create_generation(
&digest("rabs.authority.sha256.v1", 1),
42,
&digest("rabs.action-key.sha256.v1", 7),
),
Err(StoreError::GenerationIdNotAboveHighWater)
);
}
fn allocation_behavior<E: SqlEngine>(store: &mut SqlMetadataStore<E>) {
let active = authority(1);
let action = digest("rabs.action-key.sha256.v1", 7);
let other = digest("rabs.action-key.sha256.v1", 8);
store.acquire_authority(&active).unwrap();
store
.create_generation(&active.digest, 899, &action)
.unwrap();
let old = ActionGeneration {
generation_id: ActionGenerationId(900),
per_key_ordinal: 255,
created_under_authority_digest: active.digest.clone(),
};
store
.create_bound_generation(&active.digest, &old, &action)
.unwrap();
store.tombstone_generation(900).unwrap();
let next = store
.allocate_bound_generation(&active.digest, &action)
.unwrap();
assert_eq!(next.generation_id, ActionGenerationId(901));
assert_eq!(
next.per_key_ordinal, 256,
"BE blob ordering crosses byte boundary"
);
assert_eq!(next.created_under_authority_digest, active.digest);
let independent = store
.allocate_bound_generation(&active.digest, &other)
.unwrap();
assert_eq!(independent.generation_id, ActionGenerationId(902));
assert_eq!(independent.per_key_ordinal, 1);
assert_eq!(
store.allocate_bound_generation(&authority(2).digest, &action),
Err(StoreError::NotActiveAuthority)
);
assert_eq!(store.generation_count().unwrap(), 4);
for corrupt in [
SqlValue::Blob(vec![0; 7]),
SqlValue::Text("00000000".into()),
] {
store
.engine
.execute(
"UPDATE action_generations SET per_key_ordinal = ?1 WHERE id_hex = ?2",
&[corrupt, SqlValue::Text(u128_hex(900))],
)
.unwrap();
assert!(matches!(
store.allocate_bound_generation(&active.digest, &action),
Err(StoreError::Corruption(_))
));
assert_eq!(store.generation_count().unwrap(), 4);
assert_eq!(
SqlMetadataStore::<E>::generation_high_water(&mut store.engine).unwrap(),
902
);
}
store
.engine
.execute(
"UPDATE action_generations SET per_key_ordinal = ?1 WHERE id_hex = ?2",
&[SqlValue::Blob(u64_blob(255)), SqlValue::Text(u128_hex(900))],
)
.unwrap();
let next = store
.allocate_bound_generation(&active.digest, &action)
.unwrap();
assert_eq!(next.generation_id, ActionGenerationId(903));
assert_eq!(next.per_key_ordinal, 257, "256 must sort above 255");
let exhausted = ActionGeneration {
generation_id: ActionGenerationId(904),
per_key_ordinal: u64::MAX,
created_under_authority_digest: active.digest.clone(),
};
store
.create_bound_generation(&active.digest, &exhausted, &action)
.unwrap();
assert_eq!(
store.allocate_bound_generation(&active.digest, &action),
Err(StoreError::GenerationOrdinalExhausted)
);
assert_eq!(store.generation_count().unwrap(), 6);
assert_eq!(
SqlMetadataStore::<E>::generation_high_water(&mut store.engine).unwrap(),
904
);
store
.create_generation(&active.digest, u128::MAX, &other)
.unwrap();
assert_eq!(
store.allocate_bound_generation(&active.digest, &other),
Err(StoreError::GenerationIdentityExhausted)
);
assert_eq!(store.generation_count().unwrap(), 7);
}
#[test]
fn allocated_generations_enforce_fences_and_bounds_on_both_sql_engines() {
let reference = RusqliteEngine::open_in_memory().unwrap();
allocation_behavior(&mut SqlMetadataStore::open(reference).unwrap());
let candidate = FsqliteEngine::open(&fresh_path("allocation-fsq")).unwrap();
allocation_behavior(&mut SqlMetadataStore::open(candidate).unwrap());
}
#[test]
fn allocated_generation_identity_and_ordinal_survive_reopen_and_new_authority() {
let path = fresh_path("allocation-reopen");
let action = digest("rabs.action-key.sha256.v1", 7);
let old_id = 1_u128 << 100;
{
let mut store = SqlMetadataStore::open(RusqliteEngine::open(&path).unwrap()).unwrap();
let active = authority(1);
store.acquire_authority(&active).unwrap();
let old = ActionGeneration {
generation_id: ActionGenerationId(old_id),
per_key_ordinal: 17,
created_under_authority_digest: active.digest.clone(),
};
store
.create_bound_generation(&active.digest, &old, &action)
.unwrap();
let next = store
.allocate_bound_generation(&active.digest, &action)
.unwrap();
assert_eq!(next.generation_id, ActionGenerationId(old_id + 1));
assert_eq!(next.per_key_ordinal, 18);
store.tombstone_generation(next.generation_id.0).unwrap();
store.release_authority(&active.digest).unwrap();
}
let mut store = SqlMetadataStore::open(RusqliteEngine::open(&path).unwrap()).unwrap();
let active = authority(2);
store.acquire_authority(&active).unwrap();
store
.close_generations_for_other_authorities(&active.digest)
.unwrap();
let next = store
.allocate_bound_generation(&active.digest, &action)
.unwrap();
assert_eq!(next.generation_id, ActionGenerationId(old_id + 2));
assert_eq!(next.per_key_ordinal, 19);
assert_eq!(next.created_under_authority_digest, active.digest);
assert_eq!(store.generation_count().unwrap(), 3);
}
#[test]
fn h009_domain_restore_is_fail_closed() {
let path = fresh_path("domains");
{
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
store.acquire_authority(&authority(1)).unwrap();
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_eq!(
store.active_authority(),
Err(StoreError::DomainNotInterned(
"rabs.authority.sha256.v1".to_owned()
))
);
}
fn h038_fence_seed(store: &mut dyn RabsMetadataStore) {
store.acquire_authority(&authority(1)).unwrap();
let active = digest("rabs.authority.sha256.v1", 1);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 7, 70), 700),
Ok(WorkerAdmission::AdmitNewGeneration)
);
store.advance_edge_fence(&active, "edge-1", 9).unwrap();
store
.record_peer_authority_high_water(&active, "peer-1", 2, 5, 11, 100)
.unwrap();
}
fn h038_fence_check_after_reopen(store: &mut dyn RabsMetadataStore) {
store.acquire_authority(&authority(1)).unwrap();
let active = digest("rabs.authority.sha256.v1", 1);
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 2, 5, 12, 101),
Err(StoreError::StalePeerAuthority),
"reopening must not let a clone reuse an accepted term"
);
store
.record_peer_authority_high_water(&active, "peer-1", 2, 5, 11, 102)
.unwrap();
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((2, 5, 100))
);
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 1, 999, 11, 103),
Err(StoreError::StalePeerAuthority),
"credential rollback remains fenced after reopening"
);
store
.record_peer_authority_high_water(&active, "peer-1", 3, 1, 12, 104)
.unwrap();
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 2, 999, 11, 105),
Err(StoreError::StalePeerAuthority)
);
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((3, 1, 104))
);
assert_eq!(
store
.worker_incarnation_fence(&PeerId("worker-a".into()))
.unwrap()
.unwrap()
.highest_boot_generation,
WorkerBootGeneration(7)
);
assert_eq!(store.edge_fence("edge-1").unwrap(), Some(9));
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 6, 60), 701),
Ok(WorkerAdmission::RejectStaleBootGeneration)
);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 7, 71), 702),
Ok(WorkerAdmission::RejectCloneAmbiguity)
);
assert_eq!(
store.advance_edge_fence(&active, "edge-1", 8),
Err(StoreError::StaleEdgeIncarnation)
);
assert_eq!(
store.admit_worker_session(&active, &worker_offer("worker-a", 8, 80), 703),
Ok(WorkerAdmission::RejectCloneAmbiguity),
"a self-reported boot increment cannot select the legitimate clone"
);
let mut resolved = worker_offer("worker-a", 8, 80);
resolved.reenrollment_proof = Some(1);
assert_eq!(
store.admit_worker_session(&active, &resolved, 704),
Ok(WorkerAdmission::AdmitViaReenrollment)
);
}
fn seed_v19_worker_fence<E: SqlEngine>(engine: &mut E) {
for migration in MIGRATIONS
.iter()
.filter(|migration| migration.version <= 19)
{
engine.execute("BEGIN", &[]).unwrap();
for statement in migration.statements {
engine.execute(statement, &[]).unwrap();
}
engine
.execute(
"INSERT INTO schema_epochs (version, applied_seq) VALUES (?1, 0)",
&[SqlValue::Int(i64::from(migration.version))],
)
.unwrap();
engine.execute("COMMIT", &[]).unwrap();
}
engine
.execute(
"INSERT INTO worker_incarnation_fences (worker, incarnation) VALUES (?1, ?2)",
&[
SqlValue::Text("legacy-worker".to_owned()),
SqlValue::Blob(u128_blob(0xCAFE)),
],
)
.unwrap();
}
fn assert_v19_worker_fence_migrated(store: &mut dyn RabsMetadataStore) {
assert_eq!(store.schema_version(), Ok(SCHEMA_VERSION));
assert_eq!(
store.worker_incarnation_fence(&PeerId("legacy-worker".to_owned())),
Ok(Some(WorkerIncarnationFenceRecord {
worker_peer_id: PeerId("legacy-worker".to_owned()),
highest_boot_generation: WorkerBootGeneration(0),
active_incarnation: Some(WorkerIncarnationId(0xCAFE)),
clone_ambiguous: true,
operator_reenrollment_generation: 0,
})),
"a pre-S022 row migrates as an ambiguous active generation-zero fence"
);
}
fn seed_v20_attempt_lease<E: SqlEngine>(engine: &mut E) {
for migration in MIGRATIONS
.iter()
.filter(|migration| migration.version <= 20)
{
engine.execute("BEGIN", &[]).unwrap();
for statement in migration.statements {
engine.execute(statement, &[]).unwrap();
}
engine
.execute(
"INSERT INTO schema_epochs (version, applied_seq) VALUES (?1, 0)",
&[SqlValue::Int(i64::from(migration.version))],
)
.unwrap();
engine.execute("COMMIT", &[]).unwrap();
}
let authority = bound_authority_row().digest;
let action = digest("rabs.action-key.sha256.v1", 7);
engine
.execute(
"INSERT INTO action_generations \
(id_hex, id, action_key, authority_key, tombstoned) \
VALUES (?1, ?2, ?3, ?4, 0)",
&[
SqlValue::Text(u128_hex(11)),
SqlValue::Blob(u128_blob(11)),
SqlValue::Text(digest_key(&action)),
SqlValue::Text(digest_key(&authority)),
],
)
.unwrap();
engine
.execute(
"INSERT INTO generation_high_water (kind, value) \
VALUES ('action-generation', ?1)",
&[SqlValue::Blob(u128_blob(11))],
)
.unwrap();
engine
.execute(
"INSERT INTO action_attempts (id_hex, id, generation_hex, worker, seq) \
VALUES (?1, ?2, ?3, 'lease-worker', 5)",
&[
SqlValue::Text(u128_hex(22)),
SqlValue::Blob(u128_blob(22)),
SqlValue::Text(u128_hex(11)),
],
)
.unwrap();
engine
.execute(
"INSERT INTO execution_leases \
(id_hex, id, attempt_hex, renewal_seq, expires_at_seq, released) \
VALUES (?1, ?2, ?3, 1, 100, 0)",
&[
SqlValue::Text(u128_hex(30)),
SqlValue::Blob(u128_blob(30)),
SqlValue::Text(u128_hex(22)),
],
)
.unwrap();
engine
.execute(
"INSERT INTO worker_incarnation_fences \
(worker, incarnation, highest_boot_generation, active, \
operator_reenrollment_generation) \
VALUES ('lease-worker', ?1, ?2, 1, ?3)",
&[
SqlValue::Blob(u128_blob(99)),
SqlValue::Blob(u64_blob(3)),
SqlValue::Blob(u64_blob(0)),
],
)
.unwrap();
}
fn assert_v20_attempt_lease_fails_closed(store: &mut dyn RabsMetadataStore) {
assert_eq!(store.schema_version(), Ok(SCHEMA_VERSION));
store.acquire_authority(&bound_authority_row()).unwrap();
assert_eq!(
store.worker_incarnation_fence(&PeerId("lease-worker".to_owned())),
Ok(Some(WorkerIncarnationFenceRecord {
worker_peer_id: PeerId("lease-worker".to_owned()),
highest_boot_generation: WorkerBootGeneration(3),
active_incarnation: Some(WorkerIncarnationId(99)),
clone_ambiguous: true,
operator_reenrollment_generation: 0,
}))
);
let legacy = bound_attempt_authority();
assert_eq!(
store.validate_attempt_lease(&legacy, 10),
Err(StoreError::LegacyUnboundAuthority)
);
assert_eq!(
store.renew_attempt_lease(
&legacy,
LeaseRenewal {
lease: legacy.execution_lease_id,
seq: LeaseRenewalSeq(2),
},
200,
&|| 10,
),
Err(StoreError::LegacyUnboundAuthority)
);
let mut legacy_row = publication(7, 1, 40);
legacy_row.winner_generation = 11;
legacy_row.winner_attempt = 22;
assert_eq!(
store.commit_publication(
&bound_authority_row().digest,
Some((&legacy, &|| 10)),
&legacy_row,
),
Err(StoreError::LegacyUnboundAuthority)
);
assert!(!store.has_publication(&legacy.action_key).unwrap());
let mut selected = worker_offer("lease-worker", 3, 100);
selected.reenrollment_proof = Some(1);
assert_eq!(
store.admit_worker_session(&bound_authority_row().digest, &selected, 200),
Ok(WorkerAdmission::AdmitViaReenrollment)
);
let mut replacement = legacy;
replacement.action_generation.generation_id = ActionGenerationId(12);
replacement.action_generation.per_key_ordinal = 2;
replacement.attempt_id = AttemptId(23);
replacement.execution_lease_id = ExecutionLeaseId(31);
replacement.worker_incarnation_id = WorkerIncarnationId(100);
store
.upsert_action_entry(&ActionEntryRow {
action_key: replacement.action_key.clone(),
key_epoch: 1,
projection_epoch: 1,
})
.unwrap();
store
.create_bound_generation(
&bound_authority_row().digest,
&replacement.action_generation,
&replacement.action_key,
)
.unwrap();
store.admit_attempt_lease(&replacement, 201, 300).unwrap();
assert_eq!(
store.validate_attempt_lease(&replacement, 10),
Ok(LeaseState {
released: false,
renewal_seq: 1,
expires_at_own_monotonic_ms: 300,
})
);
}
#[test]
fn s022_migrates_populated_v19_worker_fence_reference() {
let path = fresh_path("s022-v19-ref");
{
let mut engine = RusqliteEngine::open(&path).unwrap();
seed_v19_worker_fence(&mut engine);
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_v19_worker_fence_migrated(&mut store);
}
#[test]
fn s022_migrates_populated_v19_worker_fence_frankensqlite() {
let path = fresh_path("s022-v19-fsq");
{
let mut engine = FsqliteEngine::open(&path).unwrap();
seed_v19_worker_fence(&mut engine);
}
let engine = FsqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_v19_worker_fence_migrated(&mut store);
}
#[test]
fn t038_migrates_populated_v20_attempt_lease_fail_closed_reference() {
let path = fresh_path("t038-v20-ref");
{
let mut engine = RusqliteEngine::open(&path).unwrap();
seed_v20_attempt_lease(&mut engine);
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_v20_attempt_lease_fails_closed(&mut store);
}
#[test]
fn t038_migrates_populated_v20_attempt_lease_fail_closed_frankensqlite() {
let path = fresh_path("t038-v20-fsq");
{
let mut engine = FsqliteEngine::open(&path).unwrap();
seed_v20_attempt_lease(&mut engine);
}
let engine = FsqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_v20_attempt_lease_fails_closed(&mut store);
}
#[test]
fn s022_and_h038_fences_survive_reopen_reference() {
let path = fresh_path("h038-fence-ref");
{
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
h038_fence_seed(&mut store);
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
h038_fence_check_after_reopen(&mut store);
}
#[test]
fn s022_and_h038_fences_survive_reopen_frankensqlite() {
let path = fresh_path("h038-fence-fsq");
{
let engine = FsqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
h038_fence_seed(&mut store);
}
let engine = FsqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
h038_fence_check_after_reopen(&mut store);
}
fn s020_operator_reset_consumption<E: SqlEngine>(
open: fn(&std::path::Path) -> Result<E, StoreError>,
tag: &str,
) {
let path = fresh_path(tag);
let mut store = SqlMetadataStore::open(open(&path).unwrap()).unwrap();
store.acquire_authority(&authority(1)).unwrap();
let active = digest("rabs.authority.sha256.v1", 1);
store
.record_peer_authority_high_water(&active, "peer-1", 2, 5, 11, 100)
.unwrap();
store
.apply_operator_reset_to_peer(&active, "peer-1", 1, 2, 5, 99)
.unwrap();
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((2, 5, 1))
);
store
.record_peer_authority_high_water(&active, "peer-1", 2, 5, 99, 201)
.unwrap();
assert_eq!(
store.record_peer_authority_high_water(&active, "peer-1", 2, 5, 11, 202),
Err(StoreError::StalePeerAuthority),
"the fenced pre-reset incarnation must stay refused"
);
assert_eq!(
store.apply_operator_reset_to_peer(&active, "peer-1", 1, 2, 6, 99),
Err(StoreError::StaleOperatorReset)
);
store
.apply_operator_reset_to_peer(&active, "peer-1", 2, 3, 7, 100)
.unwrap();
assert_eq!(store.highest_operator_reset().unwrap(), Some(2));
assert_eq!(
store.peer_authority_high_water("peer-1").unwrap(),
Some((3, 7, 2))
);
}
#[test]
fn s020_operator_reset_reopens_fenced_peer_reference() {
s020_operator_reset_consumption(RusqliteEngine::open, "s020-reset-ref");
}
#[test]
fn s020_operator_reset_reopens_fenced_peer_frankensqlite() {
s020_operator_reset_consumption(FsqliteEngine::open, "s020-reset-fsq");
}
fn table_names<E: SqlEngine>(store: &mut SqlMetadataStore<E>) -> Vec<String> {
store
.engine_mut()
.query(
"SELECT name FROM sqlite_master WHERE type = 'table' ORDER BY name",
&[],
)
.unwrap()
.into_iter()
.map(|row| match row.into_iter().next() {
Some(SqlValue::Text(t)) => t,
other => panic!("table name shape: {other:?}"),
})
.collect()
}
#[test]
fn h038_full_authoritative_table_set_no_failure_table() {
let expected: Vec<String> = [
"action_attempts",
"action_entries",
"action_evidence_index",
"action_generations",
"action_publications",
"action_serving_states",
"action_trust_evaluations",
"adoption_edges",
"coordinator_authorities",
"decision_receipts",
"determinism_audits",
"divergence_incidents",
"edge_handoffs",
"edge_incarnation_fences",
"edge_subscribers",
"eviction_tombstones",
"execution_leases",
"gc_receipts",
"gc_runs",
"gc_tombstones",
"generation_high_water",
"key_breakdowns",
"manifests",
"materialization_records",
"native_child_bindings",
"object_edges",
"object_locations",
"objects",
"observed_input_recipes",
"operations",
"operator_resets",
"peer_authority_high_water",
"pins",
"provenance_edges",
"provisional_ancestry",
"provisional_install_journal",
"provisional_obligations",
"provisional_pin_grants",
"provisional_pin_lineage",
"provisional_pins",
"quarantines",
"release_verdicts",
"schema_epochs",
"serving_blocking_quarantines",
"trust_states",
"verification_samples",
"worker_capabilities",
"worker_health_samples",
"worker_incarnation_fences",
"worker_sessions",
]
.iter()
.map(|s| (*s).to_owned())
.collect();
assert!(
expected.iter().all(|name| !name.contains("failure")),
"deterministic failures are ResultKind publications; no failure table may exist"
);
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_eq!(table_names(&mut store), expected);
let engine = FsqliteEngine::open(&fresh_path("h038-tables")).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_eq!(table_names(&mut store), expected);
}
#[test]
fn t011_a_release_verdict_round_trips_and_a_rerun_supersedes_it() {
use rabs_protocol::release_authorization::{
ReleaseAuthorization, ReleaseAuthorizationMode, authorization,
};
let mut verdict = ReleaseVerdict {
build: "rabs-build-1".to_owned(),
corpus: "corpus-abc".to_owned(),
replayed: 1_200,
explained: 3,
validity: ServingValidity {
evaluated_at_unix_micros: 5_000,
maximum_age_micros: Some(100_000),
clock_uncertainty_micros: 25,
coordinator_clock_epoch: 4,
},
};
let stores: [Box<dyn RabsMetadataStore>; 2] = [
Box::new(SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap()),
Box::new(
SqlMetadataStore::open(
FsqliteEngine::open(&fresh_path("t011-release-verdicts")).unwrap(),
)
.unwrap(),
),
];
for mut store in stores {
assert_eq!(store.release_verdict("rabs-build-1").unwrap(), None);
store.record_release_verdict(&verdict).unwrap();
assert_eq!(
store.release_verdict("rabs-build-1").unwrap().as_ref(),
Some(&verdict),
"every field must survive the round trip, including the validity window"
);
assert_eq!(store.release_verdict("rabs-build-2").unwrap(), None);
let mut rerun = verdict.clone();
rerun.corpus = "corpus-def".to_owned();
rerun.replayed = 1_500;
rerun.validity.evaluated_at_unix_micros = 9_000;
store.record_release_verdict(&rerun).unwrap();
assert_eq!(
store.release_verdict("rabs-build-1").unwrap(),
Some(rerun.clone()),
"the rerun must replace the earlier verdict, not sit beside it"
);
let stored = store.release_verdict("rabs-build-1").unwrap();
assert_eq!(
authorization(stored.as_ref(), "rabs-build-1", 9_500, 4),
ReleaseAuthorization::Authorized { replayed: 1_500 }
);
let expired = authorization(stored.as_ref(), "rabs-build-1", 500_000, 4);
assert_eq!(expired, ReleaseAuthorization::Expired);
assert!(!expired.permits_serving(ReleaseAuthorizationMode::Required));
}
verdict.validity.maximum_age_micros = None;
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
store.record_release_verdict(&verdict).unwrap();
let stored = store.release_verdict("rabs-build-1").unwrap();
assert_eq!(
stored.as_ref().unwrap().validity.maximum_age_micros,
None,
"NULL max-age must read back as no bound"
);
assert_eq!(
authorization(stored.as_ref(), "rabs-build-1", i64::MAX, 4),
ReleaseAuthorization::Authorized { replayed: 1_200 }
);
}
#[test]
fn h038_deterministic_failure_publishes_as_result_kind() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
store.acquire_authority(&authority(1)).unwrap();
let active = digest("rabs.authority.sha256.v1", 1);
let mut row = publication(9, 3, 90);
row.result_kind = ResultKindTag::DeterministicFailure;
assert_eq!(
store.commit_publication(&active, None, &row).unwrap(),
CommitOutcome::Committed
);
let dump = store.differential_snapshot().unwrap();
assert!(
dump.iter().any(|line| {
line.starts_with("action_publications|") && line.contains("|deterministic-failure|")
}),
"failure publication must appear in action_publications with its result kind"
);
}
}