feat(artifacts): add fenced reconciliation recovery

This commit is contained in:
2026-08-28 00:44:27 +03:00
parent 8784964fb2
commit d2849ea3fe
21 changed files with 3322 additions and 97 deletions
+7
View File
@@ -611,6 +611,13 @@ impl From<RegistryError> for ApiError {
"artifact source failed integrity verification",
json!({ "error_code": "artifact_source_integrity" }),
),
RegistryError::ArtifactClaimInProgress => Self::conflict_with_context(
"artifact reconciliation is in progress",
json!({
"error_code": "artifact_claim_in_progress",
"recovery": "retry"
}),
),
RegistryError::InvalidArtifactSource { field } => Self::validation_with_context(
"artifact source metadata is invalid",
json!({
+1
View File
@@ -5,6 +5,7 @@ pub mod error;
pub mod import_guidance;
pub mod pool_metrics;
pub mod rate_limit;
pub mod reconciliation;
pub mod request_context;
pub mod routes;
pub mod service;
+4 -1
View File
@@ -4,6 +4,7 @@ use admin_api::{
app::build_app,
auth::{AuthSettings, BootstrapAdminConfig},
pool_metrics::spawn_postgres_pool_metrics,
reconciliation::{open_reconciliation_store, spawn_artifact_reconciliation},
service::AdminServiceBuilder,
state::AppState,
};
@@ -202,6 +203,7 @@ async fn run(
let secret_crypto =
verified_startup_secret_crypto(&registry, config.runtime.master_key.expose_secret())
.await?;
let artifact_store = open_reconciliation_store(config.storage_root.clone()).await?;
let outbound_http_policy = crank_runtime::OutboundHttpPolicy::try_new_with_limits(
config.runtime.outbound.allowed_hosts.clone(),
config.runtime.outbound.denied_hosts.clone(),
@@ -216,7 +218,7 @@ async fn run(
let identity_provider =
PasswordIdentityProvider::new(registry.clone(), auth_settings.password_pepper.clone());
let service = AdminServiceBuilder::new(
registry,
registry.clone(),
config.storage_root.clone(),
auth_settings,
secret_crypto,
@@ -229,6 +231,7 @@ async fn run(
if config.demo_seed {
service.seed_demo_assets().await?;
}
spawn_artifact_reconciliation(registry, artifact_store).await;
spawn_invocation_log_cleanup(service.clone(), config.invocation_log_retention_days);
let state = AppState {
service,
+783
View File
@@ -0,0 +1,783 @@
use std::{
path::PathBuf,
time::{Duration, SystemTime},
};
use async_trait::async_trait;
use crank_artifacts::{
ArtifactError, ArtifactStore, ReconciliationCursor, ReconciliationMutation,
ReconciliationNamespace, ReconciliationPresence, ReconciliationRegistration,
ReconciliationScanStop,
};
use crank_registry::{
ArtifactClaimFinalization, ArtifactClaimFinalizeOutcome, ArtifactClaimOutcome,
ArtifactClaimRecheckOutcome, ArtifactClaimToken, ArtifactExpiredClaimProbe,
ClaimArtifactReconciliationRequest, ClaimExpiredArtifactReconciliationRequest,
PostgresRegistry,
};
use time::OffsetDateTime;
use tokio::time::MissedTickBehavior;
use tracing::{info, warn};
pub const RECONCILIATION_INTERVAL: Duration = Duration::from_secs(15 * 60);
pub const RECONCILIATION_GRACE: Duration = Duration::from_secs(24 * 60 * 60);
pub const RECONCILIATION_LEASE: Duration = Duration::from_secs(5 * 60);
pub const RECONCILIATION_TRAVERSAL_LIMIT: usize = 4_096;
pub const RECONCILIATION_CANDIDATE_LIMIT: usize = 32;
pub const RECONCILIATION_MUTATION_LIMIT: usize = 32;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum Phase {
Recovery,
Sweep,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct StepLimits {
traversal_syscalls: usize,
candidates: usize,
mutations: usize,
}
#[derive(Debug)]
struct Step<C> {
cursor: Option<C>,
scan_complete: bool,
traversal_syscalls: usize,
candidates: usize,
mutation_attempts: usize,
classifications: ReconciliationClassificationCounters,
physical_mutation_possible: bool,
retryable: bool,
}
/// Bounded, identity-free classifications observed during one or more scans.
///
/// The fixed fields are the complete telemetry vocabulary: no digest, path,
/// token, source or workspace value can be attached to an item classification.
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct ReconciliationClassificationCounters {
pub scanner_malformed: usize,
pub scanner_unsafe: usize,
pub recovery_present: usize,
pub recovery_safety: usize,
pub recovery_retryable: usize,
pub registration_integrity: usize,
pub registration_safety: usize,
pub registration_retryable: usize,
pub claim_metadata_conflict: usize,
pub mutation_integrity: usize,
pub mutation_safety: usize,
pub mutation_retryable: usize,
}
impl ReconciliationClassificationCounters {
fn checked_add(self, other: Self) -> Option<Self> {
Some(Self {
scanner_malformed: self
.scanner_malformed
.checked_add(other.scanner_malformed)?,
scanner_unsafe: self.scanner_unsafe.checked_add(other.scanner_unsafe)?,
recovery_present: self.recovery_present.checked_add(other.recovery_present)?,
recovery_safety: self.recovery_safety.checked_add(other.recovery_safety)?,
recovery_retryable: self
.recovery_retryable
.checked_add(other.recovery_retryable)?,
registration_integrity: self
.registration_integrity
.checked_add(other.registration_integrity)?,
registration_safety: self
.registration_safety
.checked_add(other.registration_safety)?,
registration_retryable: self
.registration_retryable
.checked_add(other.registration_retryable)?,
claim_metadata_conflict: self
.claim_metadata_conflict
.checked_add(other.claim_metadata_conflict)?,
mutation_integrity: self
.mutation_integrity
.checked_add(other.mutation_integrity)?,
mutation_safety: self.mutation_safety.checked_add(other.mutation_safety)?,
mutation_retryable: self
.mutation_retryable
.checked_add(other.mutation_retryable)?,
})
}
fn record_registration_error(&mut self, error: ArtifactError) -> bool {
match error {
ArtifactError::Integrity
| ArtifactError::EmptySource
| ArtifactError::SourceTooLarge => {
self.registration_integrity += 1;
false
}
ArtifactError::UnsafeRoot
| ArtifactError::InvalidReference
| ArtifactError::NotFound => {
self.registration_safety += 1;
false
}
ArtifactError::Storage => {
self.registration_retryable += 1;
true
}
}
}
fn record_mutation_error(&mut self, error: ArtifactError) -> bool {
match error {
ArtifactError::Integrity
| ArtifactError::EmptySource
| ArtifactError::SourceTooLarge => self.mutation_integrity += 1,
ArtifactError::UnsafeRoot
| ArtifactError::InvalidReference
| ArtifactError::NotFound => self.mutation_safety += 1,
ArtifactError::Storage => self.mutation_retryable += 1,
}
true
}
fn record_recovery_presence(&mut self, presence: ReconciliationPresence) -> bool {
match presence {
ReconciliationPresence::Final | ReconciliationPresence::Quarantine => {
self.recovery_present += 1;
false
}
ReconciliationPresence::Unsafe => {
self.recovery_safety += 1;
false
}
ReconciliationPresence::Retryable => {
self.recovery_retryable += 1;
true
}
ReconciliationPresence::Absent => false,
}
}
fn record_mutation_presence(&mut self, presence: ReconciliationPresence) -> bool {
match presence {
ReconciliationPresence::Absent => false,
ReconciliationPresence::Unsafe => {
self.mutation_safety += 1;
true
}
ReconciliationPresence::Final
| ReconciliationPresence::Quarantine
| ReconciliationPresence::Retryable => {
self.mutation_retryable += 1;
true
}
}
}
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct ReconciliationTickReport {
pub traversal_syscalls: usize,
pub candidates: usize,
pub mutation_attempts: usize,
pub classifications: ReconciliationClassificationCounters,
pub recovery_completed: bool,
pub cycle_completed: bool,
pub retryable: bool,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
pub enum ReconciliationCoordinatorError {
#[error("artifact reconciliation backend is unavailable")]
Backend,
#[error("artifact reconciliation backend exceeded its bounded contract")]
Bounds,
}
#[async_trait]
trait ReconciliationBackend: Send + Sync {
type Cursor: Send;
async fn step(
&self,
phase: Phase,
cursor: Option<Self::Cursor>,
limits: StepLimits,
) -> Result<Step<Self::Cursor>, ReconciliationCoordinatorError>;
}
struct ReconciliationCoordinator<B: ReconciliationBackend> {
backend: B,
phase: Phase,
cursor: Option<B::Cursor>,
tick_limits: StepLimits,
}
#[derive(Clone)]
struct PostgresReconciliationBackend {
registry: PostgresRegistry,
store: ArtifactStore,
}
enum PostgresReconciliationCursor {
ExpiredClaim(ArtifactExpiredClaimProbe),
Filesystem(ReconciliationCursor),
}
impl std::fmt::Debug for PostgresReconciliationCursor {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("PostgresReconciliationCursor(..)")
}
}
impl PostgresReconciliationBackend {
fn new(registry: PostgresRegistry, store: ArtifactStore) -> Self {
Self { registry, store }
}
async fn database_now(&self) -> Result<OffsetDateTime, ReconciliationCoordinatorError> {
self.registry
.artifact_reconciliation_now()
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)
}
async fn presence(
&self,
artifact_ref: crank_artifacts::ArtifactRef,
) -> Result<ReconciliationPresence, ReconciliationCoordinatorError> {
let store = self.store.clone();
let result =
tokio::task::spawn_blocking(move || store.reconciliation_presence(&artifact_ref))
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
Ok(match result {
Ok(presence) => presence,
Err(ArtifactError::Storage) => ReconciliationPresence::Retryable,
Err(_) => ReconciliationPresence::Unsafe,
})
}
}
fn reconciliation_lease_expires_at(claimed_at: OffsetDateTime) -> OffsetDateTime {
claimed_at
+ time::Duration::seconds(
i64::try_from(RECONCILIATION_LEASE.as_secs())
.expect("the fixed reconciliation lease fits i64"),
)
}
impl<B: ReconciliationBackend> ReconciliationCoordinator<B> {
fn new(backend: B) -> Self {
Self {
backend,
phase: Phase::Recovery,
cursor: None,
tick_limits: StepLimits {
traversal_syscalls: RECONCILIATION_TRAVERSAL_LIMIT,
candidates: RECONCILIATION_CANDIDATE_LIMIT,
mutations: RECONCILIATION_MUTATION_LIMIT,
},
}
}
#[cfg(test)]
fn with_tick_limits(backend: B, tick_limits: StepLimits) -> Self {
Self {
backend,
phase: Phase::Recovery,
cursor: None,
tick_limits,
}
}
async fn tick(&mut self) -> Result<ReconciliationTickReport, ReconciliationCoordinatorError> {
let mut report = ReconciliationTickReport::default();
loop {
let limits = StepLimits {
traversal_syscalls: self
.tick_limits
.traversal_syscalls
.saturating_sub(report.traversal_syscalls),
candidates: self
.tick_limits
.candidates
.saturating_sub(report.candidates),
mutations: self
.tick_limits
.mutations
.saturating_sub(report.mutation_attempts),
};
if limits.traversal_syscalls == 0 || limits.candidates == 0 || limits.mutations == 0 {
break;
}
let step = match self
.backend
.step(self.phase, self.cursor.take(), limits)
.await
{
Ok(step) => step,
Err(error) => {
self.cursor = None;
self.phase = Phase::Recovery;
return Err(error);
}
};
if step.traversal_syscalls > limits.traversal_syscalls
|| step.candidates > limits.candidates
|| step.mutation_attempts > limits.mutations
|| step.mutation_attempts > step.candidates
{
self.cursor = None;
self.phase = Phase::Recovery;
return Err(ReconciliationCoordinatorError::Bounds);
}
report.traversal_syscalls += step.traversal_syscalls;
report.candidates += step.candidates;
report.mutation_attempts += step.mutation_attempts;
report.classifications = report
.classifications
.checked_add(step.classifications)
.ok_or_else(|| {
self.cursor = None;
self.phase = Phase::Recovery;
ReconciliationCoordinatorError::Bounds
})?;
report.retryable |= step.retryable;
let made_progress =
step.traversal_syscalls != 0 || step.candidates != 0 || step.mutation_attempts != 0;
if step.physical_mutation_possible {
// A cursor is valid only while the namespace is unchanged. A
// possible rename/unlink includes ambiguous fsync outcomes.
drop(step.cursor);
self.cursor = None;
self.phase = Phase::Recovery;
} else {
self.cursor = step.cursor;
if step.scan_complete {
self.cursor = None;
match self.phase {
Phase::Recovery => {
report.recovery_completed = true;
self.phase = Phase::Sweep;
}
Phase::Sweep => {
report.cycle_completed = true;
self.phase = Phase::Recovery;
break;
}
}
}
}
if step.retryable || (!made_progress && !step.scan_complete) {
break;
}
}
Ok(report)
}
}
#[async_trait]
impl ReconciliationBackend for PostgresReconciliationBackend {
type Cursor = PostgresReconciliationCursor;
async fn step(
&self,
phase: Phase,
cursor: Option<Self::Cursor>,
limits: StepLimits,
) -> Result<Step<Self::Cursor>, ReconciliationCoordinatorError> {
let mut recovered_candidates = 0;
let mut classifications = ReconciliationClassificationCounters::default();
let mut recovery_retryable = false;
let mut filesystem_cursor = None;
if phase == Phase::Recovery {
let after = match cursor {
Some(PostgresReconciliationCursor::ExpiredClaim(probe)) => Some(probe),
Some(PostgresReconciliationCursor::Filesystem(cursor)) => {
filesystem_cursor = Some(cursor);
None
}
None => None,
};
if filesystem_cursor.is_none() {
let database_now = self.database_now().await?;
let probe_limit = u32::try_from(limits.candidates).unwrap_or(u32::MAX);
let probes = self
.registry
.list_expired_artifact_reconciliation_probes_after(
database_now,
after.as_ref(),
probe_limit,
)
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
let page_full = probes.len() == probe_limit as usize;
let mut last_processed = after;
for probe in probes {
recovered_candidates += 1;
last_processed = Some(probe.clone());
let presence = self.presence(probe.artifact_ref().clone()).await?;
if presence != ReconciliationPresence::Absent {
recovery_retryable |= classifications.record_recovery_presence(presence);
continue;
}
let claimed_at = self.database_now().await?;
let claim = self
.registry
.claim_expired_artifact_reconciliation(
&probe,
ClaimExpiredArtifactReconciliationRequest {
token: ArtifactClaimToken::generate(),
claimed_at,
lease_expires_at: reconciliation_lease_expires_at(claimed_at),
},
)
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
let Some(claim) = claim else {
continue;
};
// Recheck after the DB CAS. A concurrent filesystem put
// does not consult PostgreSQL, so the first probe alone
// cannot prove absence at finalization time.
let presence = self.presence(claim.artifact_ref().clone()).await?;
if presence != ReconciliationPresence::Absent {
classifications.record_recovery_presence(presence);
recovery_retryable = true;
continue;
}
let recovered_at = self.database_now().await?;
let result = self
.registry
.recover_artifact_reconciliation_claim(
&claim,
ArtifactClaimFinalization::AlreadyAbsent,
recovered_at,
)
.await;
let recovered = matches!(result, Ok(ArtifactClaimFinalizeOutcome::Unavailable));
if recovered {
let presence = self.presence(claim.artifact_ref().clone()).await?;
if presence != ReconciliationPresence::Absent {
classifications.record_recovery_presence(presence);
recovery_retryable = true;
}
}
return Ok(Step {
cursor: last_processed.map(PostgresReconciliationCursor::ExpiredClaim),
scan_complete: false,
traversal_syscalls: 0,
candidates: recovered_candidates,
mutation_attempts: 1,
classifications,
physical_mutation_possible: false,
retryable: recovery_retryable || !recovered,
});
}
if page_full {
return Ok(Step {
cursor: last_processed.map(PostgresReconciliationCursor::ExpiredClaim),
scan_complete: false,
traversal_syscalls: 0,
candidates: recovered_candidates,
mutation_attempts: 0,
classifications,
physical_mutation_possible: false,
retryable: recovery_retryable,
});
}
}
} else if let Some(PostgresReconciliationCursor::Filesystem(cursor)) = cursor {
filesystem_cursor = Some(cursor);
}
let scan_candidate_limit = limits.candidates.saturating_sub(recovered_candidates);
if scan_candidate_limit == 0 {
return Ok(Step {
cursor: filesystem_cursor.map(PostgresReconciliationCursor::Filesystem),
scan_complete: false,
traversal_syscalls: 0,
candidates: recovered_candidates,
mutation_attempts: 0,
classifications,
physical_mutation_possible: false,
retryable: recovery_retryable,
});
}
let namespace = match phase {
Phase::Recovery => ReconciliationNamespace::Quarantine,
Phase::Sweep => ReconciliationNamespace::Final,
};
let store = self.store.clone();
let (mut scan, candidates) = tokio::task::spawn_blocking(move || {
store.scan_reconciliation_namespace(
namespace,
filesystem_cursor,
limits.traversal_syscalls,
scan_candidate_limit,
)
})
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
let mut step = Step {
cursor: scan
.continuation
.take()
.map(PostgresReconciliationCursor::Filesystem),
scan_complete: scan.stop == ReconciliationScanStop::Complete,
traversal_syscalls: scan.traversal_syscalls,
candidates: recovered_candidates,
mutation_attempts: 0,
classifications: classifications
.checked_add(ReconciliationClassificationCounters {
scanner_malformed: scan.malformed,
scanner_unsafe: scan.unsafe_entries,
..ReconciliationClassificationCounters::default()
})
.ok_or(ReconciliationCoordinatorError::Bounds)?,
physical_mutation_possible: false,
retryable: recovery_retryable || scan.stop == ReconciliationScanStop::Retryable,
};
for candidate in candidates {
if step.candidates == limits.candidates || step.mutation_attempts == limits.mutations {
break;
}
step.candidates += 1;
let registration_grace = match phase {
Phase::Recovery => Duration::ZERO,
Phase::Sweep => RECONCILIATION_GRACE,
};
let store = self.store.clone();
let registration = tokio::task::spawn_blocking(move || {
let registration = store.register_reconciliation(
&candidate,
registration_grace,
SystemTime::now(),
);
(candidate, registration)
})
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
let (candidate, registration) = registration;
let artifact = match registration {
Ok(ReconciliationRegistration::Registered(artifact)) => artifact,
Ok(ReconciliationRegistration::NotEligible) => continue,
Err(error) => {
step.retryable |= step.classifications.record_registration_error(error);
continue;
}
};
let claimed_at = self.database_now().await?;
let lease_expires_at = reconciliation_lease_expires_at(claimed_at);
let detached_grace = time::Duration::seconds(
i64::try_from(RECONCILIATION_GRACE.as_secs())
.expect("the fixed reconciliation grace fits i64"),
);
let claim = self
.registry
.claim_artifact_reconciliation(ClaimArtifactReconciliationRequest {
artifact: &artifact,
token: ArtifactClaimToken::generate(),
claimed_at,
lease_expires_at,
detached_grace,
})
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
let claim = match claim {
ArtifactClaimOutcome::Claimed(claim) => claim,
ArtifactClaimOutcome::HeldByOther
| ArtifactClaimOutcome::ActiveReference
| ArtifactClaimOutcome::DetachedReferenceInGrace => continue,
ArtifactClaimOutcome::MetadataConflict => {
step.classifications.claim_metadata_conflict += 1;
continue;
}
};
let rechecked_at = self.database_now().await?;
let recheck = self
.registry
.recheck_artifact_reconciliation_claim(&claim, rechecked_at, detached_grace)
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?;
if recheck != ArtifactClaimRecheckOutcome::Mutate {
continue;
}
step.mutation_attempts += 1;
step.physical_mutation_possible = true;
let store = self.store.clone();
let mutation = match tokio::task::spawn_blocking(move || match phase {
Phase::Recovery => store.delete_quarantined_reconciliation(candidate),
Phase::Sweep => store.quarantine_reconciliation(candidate),
})
.await
{
Ok(mutation) => mutation,
Err(_) => Err(ArtifactError::Storage),
};
let mut observation = match mutation {
Ok(ReconciliationMutation::Quarantined)
| Ok(ReconciliationMutation::AlreadyQuarantined) => {
ArtifactClaimFinalization::Quarantined
}
Ok(ReconciliationMutation::Deleted) => ArtifactClaimFinalization::Deleted,
Ok(ReconciliationMutation::AlreadyAbsent) => {
ArtifactClaimFinalization::AlreadyAbsent
}
Ok(ReconciliationMutation::Retryable) => {
step.classifications.mutation_retryable += 1;
step.retryable = true;
ArtifactClaimFinalization::Retryable
}
Err(error) => {
step.retryable |= step.classifications.record_mutation_error(error);
ArtifactClaimFinalization::Retryable
}
};
if matches!(
observation,
ArtifactClaimFinalization::Deleted | ArtifactClaimFinalization::AlreadyAbsent
) {
// `unavailable` means both canonical locations are confirmed
// absent. A concurrent republish can recreate final bytes
// without consulting PostgreSQL, so unlink alone is not proof.
let presence = self.presence(claim.artifact_ref().clone()).await?;
observation = match presence {
ReconciliationPresence::Absent => ArtifactClaimFinalization::AlreadyAbsent,
presence => {
step.retryable |= step.classifications.record_mutation_presence(presence);
ArtifactClaimFinalization::Retryable
}
};
}
let observed_at = self.database_now().await?;
let finalization = match phase {
Phase::Recovery => {
self.registry
.recover_artifact_reconciliation_claim(&claim, observation, observed_at)
.await
}
Phase::Sweep => {
self.registry
.finalize_artifact_reconciliation_claim(&claim, observation, observed_at)
.await
}
};
let finalized_as_expected =
match observation {
ArtifactClaimFinalization::Deleted
| ArtifactClaimFinalization::AlreadyAbsent => matches!(
finalization.as_ref(),
Ok(ArtifactClaimFinalizeOutcome::Unavailable)
),
ArtifactClaimFinalization::Quarantined
| ArtifactClaimFinalization::Retryable => matches!(
finalization.as_ref(),
Ok(ArtifactClaimFinalizeOutcome::Retained)
),
};
if !finalized_as_expected {
step.retryable = true;
}
if matches!(
finalization.as_ref(),
Ok(ArtifactClaimFinalizeOutcome::Unavailable)
) {
let presence = self.presence(claim.artifact_ref().clone()).await?;
step.retryable |= step.classifications.record_mutation_presence(presence);
}
// The scan cursor and every remaining observation came from the
// pre-mutation namespace. Return immediately so the coordinator
// drops them and restarts a recovery cycle from `None`.
break;
}
Ok(step)
}
}
/// Opens the protected root without running filesystem syscalls on Tokio's
/// async executor.
pub async fn open_reconciliation_store(
root: PathBuf,
) -> Result<ArtifactStore, ReconciliationCoordinatorError> {
tokio::task::spawn_blocking(move || ArtifactStore::open(root))
.await
.map_err(|_| ReconciliationCoordinatorError::Backend)?
.map_err(|_| ReconciliationCoordinatorError::Backend)
}
/// Runs the immediate bounded startup cycle, then schedules 15-minute ticks.
/// Recovery remains ahead of the final sweep even when it spans several ticks.
pub async fn spawn_artifact_reconciliation(registry: PostgresRegistry, store: ArtifactStore) {
let backend = PostgresReconciliationBackend::new(registry, store);
let mut coordinator = ReconciliationCoordinator::new(backend);
observe_tick(coordinator.tick().await);
tokio::spawn(async move {
let mut interval = tokio::time::interval(RECONCILIATION_INTERVAL);
interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
// The immediate tick was run above as part of startup composition.
interval.tick().await;
loop {
interval.tick().await;
observe_tick(coordinator.tick().await);
}
});
}
fn observe_tick(result: Result<ReconciliationTickReport, ReconciliationCoordinatorError>) {
match result {
Ok(report) => info!(
name: "admin.artifact_reconciliation.tick",
traversal_syscalls = report.traversal_syscalls,
candidates = report.candidates,
mutation_attempts = report.mutation_attempts,
scanner_malformed = report.classifications.scanner_malformed,
scanner_unsafe = report.classifications.scanner_unsafe,
recovery_present = report.classifications.recovery_present,
recovery_safety = report.classifications.recovery_safety,
recovery_retryable = report.classifications.recovery_retryable,
registration_integrity = report.classifications.registration_integrity,
registration_safety = report.classifications.registration_safety,
registration_retryable = report.classifications.registration_retryable,
claim_metadata_conflict = report.classifications.claim_metadata_conflict,
mutation_integrity = report.classifications.mutation_integrity,
mutation_safety = report.classifications.mutation_safety,
mutation_retryable = report.classifications.mutation_retryable,
recovery_completed = report.recovery_completed,
cycle_completed = report.cycle_completed,
retryable = report.retryable,
"artifact reconciliation tick completed"
),
Err(error) => warn!(
name: "admin.artifact_reconciliation.failed",
error_category = match error {
ReconciliationCoordinatorError::Backend => "backend",
ReconciliationCoordinatorError::Bounds => "bounds",
},
"artifact reconciliation tick will be retried"
),
}
}
#[cfg(test)]
#[path = "reconciliation/tests.rs"]
mod tests;
+690
View File
@@ -0,0 +1,690 @@
use std::{
collections::VecDeque,
fs,
os::unix::fs::{PermissionsExt, symlink},
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
},
};
use async_trait::async_trait;
use crank_artifacts::{ReconciliationRegistration, RegisteredArtifact};
use crank_registry::MigrationAuthority;
use super::*;
static NEXT_TEST_ROOT: AtomicUsize = AtomicUsize::new(0);
struct TestRoot(PathBuf);
impl TestRoot {
fn new() -> Self {
let path = std::env::temp_dir().join(format!(
"crank-admin-reconciliation-{}-{}",
std::process::id(),
NEXT_TEST_ROOT.fetch_add(1, Ordering::Relaxed)
));
fs::create_dir(&path).unwrap();
fs::set_permissions(&path, fs::Permissions::from_mode(0o700)).unwrap();
Self(path)
}
}
impl Drop for TestRoot {
fn drop(&mut self) {
let _ = fs::set_permissions(&self.0, fs::Permissions::from_mode(0o700));
let _ = fs::remove_dir_all(&self.0);
}
}
fn artifact_path(root: &Path, artifact: &RegisteredArtifact) -> PathBuf {
let digest = artifact.artifact_ref().digest_hex();
root.join("sha256").join(&digest[..2]).join(digest)
}
fn age_artifact(root: &Path, artifact: &RegisteredArtifact) {
fs::File::open(artifact_path(root, artifact))
.unwrap()
.set_modified(SystemTime::now() - RECONCILIATION_GRACE - Duration::from_secs(60))
.unwrap();
}
fn remove_registered_artifact(store: &ArtifactStore, artifact: &RegisteredArtifact) {
let (_, candidates) = store
.scan_reconciliation_namespace(ReconciliationNamespace::Final, None, 4_096, 32)
.unwrap();
let candidate = candidates
.into_iter()
.find(|candidate| {
matches!(
store.register_reconciliation(candidate, Duration::ZERO, SystemTime::now()),
Ok(ReconciliationRegistration::Registered(found)) if found == *artifact
)
})
.expect("the final artifact is discoverable");
assert!(matches!(
store.quarantine_reconciliation(candidate),
Ok(ReconciliationMutation::Quarantined | ReconciliationMutation::AlreadyQuarantined)
));
let (_, candidates) = store
.scan_reconciliation_namespace(ReconciliationNamespace::Quarantine, None, 4_096, 32)
.unwrap();
let candidate = candidates
.into_iter()
.find(|candidate| {
matches!(
store.register_reconciliation(candidate, Duration::ZERO, SystemTime::now()),
Ok(ReconciliationRegistration::Registered(found)) if found == *artifact
)
})
.expect("the quarantined artifact is discoverable");
assert!(matches!(
store.delete_quarantined_reconciliation(candidate),
Ok(ReconciliationMutation::Deleted | ReconciliationMutation::AlreadyAbsent)
));
}
#[derive(Clone, Copy)]
struct StepTemplate {
scan_complete: bool,
traversal_syscalls: usize,
candidates: usize,
mutation_attempts: usize,
classifications: ReconciliationClassificationCounters,
physical_mutation_possible: bool,
retryable: bool,
return_cursor: bool,
error: bool,
}
impl StepTemplate {
fn complete() -> Self {
Self {
scan_complete: true,
traversal_syscalls: 1,
candidates: 0,
mutation_attempts: 0,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: false,
retryable: false,
return_cursor: false,
error: false,
}
}
}
struct CursorEvidence {
alive: Arc<AtomicUsize>,
}
impl CursorEvidence {
fn new(alive: Arc<AtomicUsize>, maximum: &AtomicUsize) -> Self {
let current = alive.fetch_add(1, Ordering::SeqCst) + 1;
maximum.fetch_max(current, Ordering::SeqCst);
Self { alive }
}
}
impl Drop for CursorEvidence {
fn drop(&mut self) {
self.alive.fetch_sub(1, Ordering::SeqCst);
}
}
struct ScriptedBackend {
steps: Mutex<VecDeque<StepTemplate>>,
calls: Mutex<Vec<(Phase, bool, StepLimits)>>,
alive: Arc<AtomicUsize>,
maximum: AtomicUsize,
}
impl ScriptedBackend {
fn new(steps: impl IntoIterator<Item = StepTemplate>) -> Self {
Self {
steps: Mutex::new(steps.into_iter().collect()),
calls: Mutex::new(Vec::new()),
alive: Arc::new(AtomicUsize::new(0)),
maximum: AtomicUsize::new(0),
}
}
}
#[async_trait]
impl ReconciliationBackend for Arc<ScriptedBackend> {
type Cursor = CursorEvidence;
async fn step(
&self,
phase: Phase,
cursor: Option<Self::Cursor>,
limits: StepLimits,
) -> Result<Step<Self::Cursor>, ReconciliationCoordinatorError> {
self.calls
.lock()
.unwrap()
.push((phase, cursor.is_some(), limits));
drop(cursor);
let template = self.steps.lock().unwrap().pop_front().unwrap();
if template.error {
return Err(ReconciliationCoordinatorError::Backend);
}
let cursor = template
.return_cursor
.then(|| CursorEvidence::new(self.alive.clone(), &self.maximum));
Ok(Step {
cursor,
scan_complete: template.scan_complete,
traversal_syscalls: template.traversal_syscalls,
candidates: template.candidates,
mutation_attempts: template.mutation_attempts,
classifications: template.classifications,
physical_mutation_possible: template.physical_mutation_possible,
retryable: template.retryable,
})
}
}
#[tokio::test]
async fn recovery_always_completes_before_final_sweep() {
let backend = Arc::new(ScriptedBackend::new([
StepTemplate::complete(),
StepTemplate::complete(),
]));
let mut coordinator = ReconciliationCoordinator::new(backend.clone());
let report = coordinator.tick().await.unwrap();
assert!(report.recovery_completed);
assert!(report.cycle_completed);
let phases = backend
.calls
.lock()
.unwrap()
.iter()
.map(|call| call.0)
.collect::<Vec<_>>();
assert_eq!(phases, vec![Phase::Recovery, Phase::Sweep]);
}
#[tokio::test]
async fn possible_mutation_discards_cursor_and_restarts_recovery() {
let paged = StepTemplate {
scan_complete: false,
traversal_syscalls: 2,
candidates: 1,
mutation_attempts: 0,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: false,
retryable: true,
return_cursor: true,
error: false,
};
let mutation = StepTemplate {
scan_complete: false,
traversal_syscalls: 1,
candidates: 1,
mutation_attempts: 1,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: true,
retryable: false,
return_cursor: true,
error: false,
};
let backend = Arc::new(ScriptedBackend::new([
StepTemplate::complete(),
paged,
mutation,
StepTemplate::complete(),
StepTemplate::complete(),
]));
let mut coordinator = ReconciliationCoordinator::new(backend.clone());
let first = coordinator.tick().await.unwrap();
assert!(first.retryable);
assert_eq!(backend.alive.load(Ordering::SeqCst), 1);
let second = coordinator.tick().await.unwrap();
assert!(second.cycle_completed);
let calls = backend.calls.lock().unwrap();
assert_eq!(
calls
.iter()
.map(|(phase, had_cursor, _)| (*phase, *had_cursor))
.collect::<Vec<_>>(),
vec![
(Phase::Recovery, false),
(Phase::Sweep, false),
(Phase::Sweep, true),
(Phase::Recovery, false),
(Phase::Sweep, false),
]
);
assert_eq!(backend.alive.load(Ordering::SeqCst), 0);
assert_eq!(backend.maximum.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn backend_error_discards_cursor_and_restarts_recovery() {
let paged = StepTemplate {
scan_complete: false,
traversal_syscalls: 1,
candidates: 1,
mutation_attempts: 0,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: false,
retryable: true,
return_cursor: true,
error: false,
};
let failure = StepTemplate {
error: true,
..StepTemplate::complete()
};
let backend = Arc::new(ScriptedBackend::new([
StepTemplate::complete(),
paged,
failure,
StepTemplate::complete(),
StepTemplate::complete(),
]));
let mut coordinator = ReconciliationCoordinator::new(backend.clone());
assert!(coordinator.tick().await.unwrap().retryable);
assert_eq!(backend.alive.load(Ordering::SeqCst), 1);
assert_eq!(
coordinator.tick().await,
Err(ReconciliationCoordinatorError::Backend)
);
assert_eq!(backend.alive.load(Ordering::SeqCst), 0);
assert!(coordinator.tick().await.unwrap().cycle_completed);
let calls = backend.calls.lock().unwrap();
assert_eq!(
calls
.iter()
.map(|(phase, had_cursor, _)| (*phase, *had_cursor))
.collect::<Vec<_>>(),
vec![
(Phase::Recovery, false),
(Phase::Sweep, false),
(Phase::Sweep, true),
(Phase::Recovery, false),
(Phase::Sweep, false),
]
);
}
#[tokio::test]
async fn backend_cannot_exceed_tick_limits() {
let backend = Arc::new(ScriptedBackend::new([StepTemplate {
scan_complete: false,
traversal_syscalls: RECONCILIATION_TRAVERSAL_LIMIT + 1,
candidates: 0,
mutation_attempts: 0,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: false,
retryable: false,
return_cursor: false,
error: false,
}]));
let mut coordinator = ReconciliationCoordinator::new(backend);
assert_eq!(
coordinator.tick().await,
Err(ReconciliationCoordinatorError::Bounds)
);
}
#[tokio::test]
async fn tick_stops_at_the_exact_candidate_and_mutation_limits() {
let backend = Arc::new(ScriptedBackend::new(
(0..RECONCILIATION_MUTATION_LIMIT).map(|_| StepTemplate {
scan_complete: false,
traversal_syscalls: 1,
candidates: 1,
mutation_attempts: 1,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: true,
retryable: false,
return_cursor: true,
error: false,
}),
));
let mut coordinator = ReconciliationCoordinator::new(backend.clone());
let report = coordinator.tick().await.unwrap();
assert_eq!(report.candidates, RECONCILIATION_CANDIDATE_LIMIT);
assert_eq!(report.mutation_attempts, RECONCILIATION_MUTATION_LIMIT);
assert_eq!(report.traversal_syscalls, RECONCILIATION_MUTATION_LIMIT);
assert_eq!(
backend.calls.lock().unwrap().len(),
RECONCILIATION_MUTATION_LIMIT
);
assert_eq!(backend.alive.load(Ordering::SeqCst), 0);
assert_eq!(backend.maximum.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn coordinator_retains_at_most_one_cursor_between_ticks() {
let backend = Arc::new(ScriptedBackend::new((0..128).map(|_| StepTemplate {
scan_complete: false,
traversal_syscalls: 1,
candidates: 0,
mutation_attempts: 0,
classifications: ReconciliationClassificationCounters::default(),
physical_mutation_possible: false,
retryable: true,
return_cursor: true,
error: false,
})));
let mut coordinator = ReconciliationCoordinator::new(backend.clone());
for _ in 0..128 {
coordinator.tick().await.unwrap();
assert_eq!(backend.alive.load(Ordering::SeqCst), 1);
}
assert_eq!(backend.maximum.load(Ordering::SeqCst), 1);
drop(coordinator);
assert_eq!(backend.alive.load(Ordering::SeqCst), 0);
}
#[test]
fn registration_failures_have_a_closed_redacted_classification() {
let mut counters = ReconciliationClassificationCounters::default();
assert!(!counters.record_registration_error(ArtifactError::Integrity));
assert!(!counters.record_registration_error(ArtifactError::UnsafeRoot));
assert!(counters.record_registration_error(ArtifactError::Storage));
assert!(!counters.record_registration_error(ArtifactError::NotFound));
assert!(!counters.record_registration_error(ArtifactError::InvalidReference));
assert!(!counters.record_registration_error(ArtifactError::EmptySource));
assert!(!counters.record_registration_error(ArtifactError::SourceTooLarge));
assert_eq!(
counters,
ReconciliationClassificationCounters {
scanner_malformed: 0,
scanner_unsafe: 0,
recovery_present: 0,
recovery_safety: 0,
recovery_retryable: 0,
registration_integrity: 3,
registration_safety: 3,
registration_retryable: 1,
claim_metadata_conflict: 0,
mutation_integrity: 0,
mutation_safety: 0,
mutation_retryable: 0,
}
);
let debug = format!("{counters:?}");
assert!(!debug.contains("sha256:"));
assert!(!debug.contains('/'));
}
#[test]
fn mutation_failures_are_classified_and_always_retryable() {
let mut counters = ReconciliationClassificationCounters::default();
assert!(counters.record_mutation_error(ArtifactError::Integrity));
assert!(counters.record_mutation_error(ArtifactError::UnsafeRoot));
assert!(counters.record_mutation_error(ArtifactError::Storage));
assert!(counters.record_mutation_error(ArtifactError::NotFound));
assert!(counters.record_mutation_error(ArtifactError::InvalidReference));
assert!(counters.record_mutation_error(ArtifactError::EmptySource));
assert!(counters.record_mutation_error(ArtifactError::SourceTooLarge));
assert_eq!(counters.mutation_integrity, 3);
assert_eq!(counters.mutation_safety, 3);
assert_eq!(counters.mutation_retryable, 1);
}
#[tokio::test]
async fn production_backend_recovers_absence_before_sweep_and_restarts_after_mutation() {
let database_url =
crank_test_support::postgres_schema_url("admin_reconciliation_backend").await;
let pool = sqlx::PgPool::connect(&database_url).await.unwrap();
MigrationAuthority::apply(&pool).await.unwrap();
let registry = PostgresRegistry::connect(&database_url).await.unwrap();
let root = TestRoot::new();
let store = ArtifactStore::open(&root.0).unwrap();
let now = OffsetDateTime::now_utc();
let expired = store.put_registered(b"expired physical absence\n").unwrap();
let expired_claim = registry
.claim_artifact_reconciliation(ClaimArtifactReconciliationRequest {
artifact: &expired,
token: ArtifactClaimToken::generate(),
claimed_at: now - time::Duration::minutes(10),
lease_expires_at: now - time::Duration::minutes(5),
detached_grace: time::Duration::hours(24),
})
.await
.unwrap();
assert!(matches!(expired_claim, ArtifactClaimOutcome::Claimed(_)));
remove_registered_artifact(&store, &expired);
sqlx::query(
"update artifact_blobs
set claim_expires_at = clock_timestamp() - interval '1 second'
where digest = $1",
)
.bind(expired.artifact_ref().digest_hex())
.execute(registry.pool())
.await
.unwrap();
let corrupted = store
.put_registered(b"integrity-blocked-candidate\n")
.unwrap();
let corrupted_path = artifact_path(&root.0, &corrupted);
let corrupted_size = fs::metadata(&corrupted_path).unwrap().len() as usize;
fs::set_permissions(&corrupted_path, fs::Permissions::from_mode(0o600)).unwrap();
fs::write(&corrupted_path, vec![b'x'; corrupted_size]).unwrap();
fs::set_permissions(&corrupted_path, fs::Permissions::from_mode(0o400)).unwrap();
age_artifact(&root.0, &corrupted);
let later = store.put_registered(b"later-valid-candidate\n").unwrap();
age_artifact(&root.0, &later);
let metadata_conflict = store
.put_registered(b"metadata-conflict-candidate\n")
.unwrap();
age_artifact(&root.0, &metadata_conflict);
sqlx::query(
"insert into artifact_blobs
(digest, artifact_ref, size_bytes, storage_lifecycle, created_at, updated_at)
values ($1, $2, $3, 'available', $4, $4)",
)
.bind(metadata_conflict.artifact_ref().digest_hex())
.bind(metadata_conflict.artifact_ref().as_str())
.bind(i64::try_from(metadata_conflict.size_bytes()).unwrap() + 1)
.bind(now)
.execute(registry.pool())
.await
.unwrap();
let shard = artifact_path(&root.0, &later)
.parent()
.unwrap()
.to_path_buf();
fs::write(shard.join("not-a-digest"), b"malformed").unwrap();
let unsafe_name = format!(
"{}{}",
&later.artifact_ref().digest_hex()[..2],
"f".repeat(62)
);
assert_ne!(unsafe_name, later.artifact_ref().digest_hex());
symlink("not-a-digest", shard.join(unsafe_name)).unwrap();
let backend = PostgresReconciliationBackend::new(registry.clone(), store.clone());
let mut coordinator = ReconciliationCoordinator::with_tick_limits(
backend,
StepLimits {
traversal_syscalls: RECONCILIATION_TRAVERSAL_LIMIT,
candidates: RECONCILIATION_CANDIDATE_LIMIT,
mutations: 1,
},
);
// The single mutation slot is consumed by expired-claim recovery. The old
// final candidate proves Sweep did not run ahead of Recovery.
let first = coordinator.tick().await.unwrap();
assert_eq!(first.mutation_attempts, 1);
assert!(!first.recovery_completed);
assert_eq!(
store.reconciliation_presence(later.artifact_ref()).unwrap(),
ReconciliationPresence::Final
);
let lifecycle: String =
sqlx::query_scalar("select storage_lifecycle from artifact_blobs where digest = $1")
.bind(expired.artifact_ref().digest_hex())
.fetch_one(registry.pool())
.await
.unwrap();
assert_eq!(lifecycle, "unavailable");
// The real final sweep classifies blocked entries and still reaches the
// later valid item. Its physical mutation invalidates the final cursor.
let second = coordinator.tick().await.unwrap();
assert!(second.recovery_completed);
assert_eq!(second.mutation_attempts, 1);
assert_eq!(
store.reconciliation_presence(later.artifact_ref()).unwrap(),
ReconciliationPresence::Quarantine
);
// Model the next scheduled/restarted lease window. Because the coordinator
// reset to Recovery and discarded its final cursor, quarantine is handled
// before another final sweep.
sqlx::query("update artifact_blobs set claim_expires_at = $1 where digest = $2")
.bind(OffsetDateTime::now_utc() - time::Duration::seconds(1))
.bind(later.artifact_ref().digest_hex())
.execute(registry.pool())
.await
.unwrap();
let third = coordinator.tick().await.unwrap();
assert_eq!(third.mutation_attempts, 1);
assert_eq!(
store.reconciliation_presence(later.artifact_ref()).unwrap(),
ReconciliationPresence::Absent
);
let fourth = coordinator.tick().await.unwrap();
assert!(fourth.cycle_completed);
let classifications = first
.classifications
.checked_add(second.classifications)
.and_then(|value| value.checked_add(third.classifications))
.and_then(|value| value.checked_add(fourth.classifications))
.unwrap();
assert!(classifications.scanner_malformed > 0);
assert!(classifications.scanner_unsafe > 0);
assert!(classifications.registration_integrity > 0);
assert!(classifications.claim_metadata_conflict > 0);
assert_eq!(
store
.reconciliation_presence(corrupted.artifact_ref())
.unwrap(),
ReconciliationPresence::Unsafe
);
assert_eq!(classifications.registration_safety, 0);
assert_eq!(classifications.registration_retryable, 0);
}
#[tokio::test]
async fn expired_claim_cursor_prevents_unsafe_prefix_starvation_across_ticks() {
let database_url =
crank_test_support::postgres_schema_url("admin_reconciliation_starvation").await;
let pool = sqlx::PgPool::connect(&database_url).await.unwrap();
MigrationAuthority::apply(&pool).await.unwrap();
let registry = PostgresRegistry::connect(&database_url).await.unwrap();
let root = TestRoot::new();
let store = ArtifactStore::open(&root.0).unwrap();
let now = OffsetDateTime::now_utc();
for index in 0..=RECONCILIATION_CANDIDATE_LIMIT {
let artifact = store
.put_registered(format!("unsafe expired claim {index}\n").as_bytes())
.unwrap();
let outcome = registry
.claim_artifact_reconciliation(ClaimArtifactReconciliationRequest {
artifact: &artifact,
token: ArtifactClaimToken::generate(),
claimed_at: now - time::Duration::minutes(20),
lease_expires_at: now - time::Duration::minutes(15),
detached_grace: time::Duration::hours(24),
})
.await
.unwrap();
assert!(matches!(outcome, ArtifactClaimOutcome::Claimed(_)));
sqlx::query(
"update artifact_blobs
set claim_expires_at = clock_timestamp() - interval '10 minutes'
where digest = $1",
)
.bind(artifact.artifact_ref().digest_hex())
.execute(registry.pool())
.await
.unwrap();
fs::set_permissions(
artifact_path(&root.0, &artifact),
fs::Permissions::from_mode(0o600),
)
.unwrap();
}
let absent = store
.put_registered(b"absent after unsafe prefix\n")
.unwrap();
let outcome = registry
.claim_artifact_reconciliation(ClaimArtifactReconciliationRequest {
artifact: &absent,
token: ArtifactClaimToken::generate(),
claimed_at: now - time::Duration::minutes(15),
lease_expires_at: now - time::Duration::minutes(10),
detached_grace: time::Duration::hours(24),
})
.await
.unwrap();
assert!(matches!(outcome, ArtifactClaimOutcome::Claimed(_)));
sqlx::query(
"update artifact_blobs
set claim_expires_at = clock_timestamp() - interval '5 minutes'
where digest = $1",
)
.bind(absent.artifact_ref().digest_hex())
.execute(registry.pool())
.await
.unwrap();
remove_registered_artifact(&store, &absent);
let backend = PostgresReconciliationBackend::new(registry.clone(), store);
let mut coordinator = ReconciliationCoordinator::new(backend);
let first = coordinator.tick().await.unwrap();
assert_eq!(first.candidates, RECONCILIATION_CANDIDATE_LIMIT);
assert_eq!(first.mutation_attempts, 0);
assert_eq!(
first.classifications.recovery_safety,
RECONCILIATION_CANDIDATE_LIMIT
);
let first_lifecycle: String =
sqlx::query_scalar("select storage_lifecycle from artifact_blobs where digest = $1")
.bind(absent.artifact_ref().digest_hex())
.fetch_one(registry.pool())
.await
.unwrap();
assert_eq!(first_lifecycle, "available");
let second = coordinator.tick().await.unwrap();
assert!(second.classifications.recovery_safety > 0);
assert!(second.mutation_attempts > 0);
let second_lifecycle: String =
sqlx::query_scalar("select storage_lifecycle from artifact_blobs where digest = $1")
.bind(absent.artifact_ref().digest_hex())
.fetch_one(registry.pool())
.await
.unwrap();
assert_eq!(second_lifecycle, "unavailable");
}