use std::future::Future;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use tokio::sync::{Mutex as AsyncMutex, Notify};
use super::backoff::Backoff;
use super::compile::{CandidateCompiler, CompileError};
use super::lkg::{LastKnownGood, LastKnownGoodError};
use super::settings::ConvergenceSettings;
use super::status::{Clock, Rejection, RevisionReport, RevisionStatus, SnapshotSource};
use crate::backends::BackendFailure;
use crate::backends::control_plane::{ControlPlaneError, ControlPlaneStore};
use crate::desired_state::{LoadedRevision, RevisionId};
use crate::policy::{ActivationRefusal, PolicyView};
use crate::state::{AppState, ConfigSnapshot};
use crate::telemetry;
pub trait SnapshotSink: Send + Sync {
fn admit(&self, _snapshot: &ConfigSnapshot) -> Result<(), ActivationRefusal> {
Ok(())
}
fn publish(&self, snapshot: ConfigSnapshot);
fn generation(&self) -> u64;
}
impl SnapshotSink for AppState {
fn admit(&self, snapshot: &ConfigSnapshot) -> Result<(), ActivationRefusal> {
self.policy().plan(&PolicyView::of(&snapshot.config))?;
Ok(())
}
fn publish(&self, snapshot: ConfigSnapshot) {
self.policy().install(PolicyView::of(&snapshot.config));
AppState::publish(self, snapshot);
}
fn generation(&self) -> u64 {
self.config().generation
}
}
#[derive(Debug, Default)]
pub struct ChangeSignal {
notify: Notify,
refresh: AtomicBool,
}
impl ChangeSignal {
pub fn new() -> Self {
Self::default()
}
pub fn notify(&self) {
self.notify.notify_one();
}
pub fn force_refresh(&self) {
self.refresh.store(true, Ordering::Release);
self.notify.notify_one();
}
async fn notified(&self) -> bool {
self.notify.notified().await;
self.refresh.swap(false, Ordering::Acquire)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Outcome {
Published {
revision: RevisionId,
generation: u64,
took: Duration,
},
AlreadyConverged { revision: Option<RevisionId> },
Empty,
Rejected {
revision: Option<RevisionId>,
reason: &'static str,
},
}
impl Outcome {
pub const fn as_str(&self) -> &'static str {
match self {
Self::Published { .. } => "published",
Self::AlreadyConverged { .. } => "converged",
Self::Empty => "empty",
Self::Rejected { .. } => "rejected",
}
}
}
pub const INCOMPATIBLE_REASON: &str = "incompatible";
pub const REVISION_REASONS: &[&str] = &[
"unavailable",
"conflict",
"not_found",
"invalid",
"denied",
"corrupt",
INCOMPATIBLE_REASON,
"secret",
"projection",
"validation",
"pricing",
"clock",
"snapshot",
"unsupported",
"migration",
"refused",
"withdrawn",
"ungoverned",
"invalid_policy",
];
pub const fn category_reason(category: crate::backends::FailureCategory) -> &'static str {
use crate::backends::FailureCategory;
match category {
FailureCategory::Unavailable => "unavailable",
FailureCategory::Conflict => "conflict",
FailureCategory::NotFound => "not_found",
FailureCategory::Invalid => "invalid",
FailureCategory::Denied => "denied",
FailureCategory::Corrupt => "corrupt",
}
}
#[derive(Debug, thiserror::Error)]
enum AttemptError {
#[error(transparent)]
Store(#[from] ControlPlaneError),
#[error(transparent)]
Compile(#[from] CompileError),
}
impl AttemptError {
fn reason(&self) -> &'static str {
match self {
Self::Store(ControlPlaneError::Incompatible { .. }) => INCOMPATIBLE_REASON,
Self::Store(error) => category_reason(error.category()),
Self::Compile(error) => error.reason(),
}
}
fn revision(&self) -> Option<RevisionId> {
match self {
Self::Store(ControlPlaneError::Corrupt { revision, .. })
| Self::Store(ControlPlaneError::Incompatible { revision, .. })
| Self::Store(ControlPlaneError::RevisionNotFound(revision))
| Self::Store(ControlPlaneError::TooLarge { revision, .. }) => Some(*revision),
Self::Store(_) => None,
Self::Compile(error) => Some(error.revision()),
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum BootstrapError {
#[error(
"the control plane is reachable but has published no revision; a stateful replica \
has nothing to serve until desired state exists"
)]
Empty,
#[error(
"the control plane is unreachable and no last-known-good snapshot is available: {source}"
)]
Unavailable {
#[source]
source: ControlPlaneError,
},
#[error("the control plane refused to yield desired state: {source}")]
Store {
#[source]
source: ControlPlaneError,
},
#[error("the desired revision cannot be served: {source}")]
Rejected {
#[source]
source: Box<CompileError>,
},
#[error("the last-known-good snapshot could not be restored: {source}")]
Cache {
#[source]
source: Box<LastKnownGoodError>,
},
}
pub struct Reconciler {
store: Arc<dyn ControlPlaneStore>,
compiler: Arc<dyn CandidateCompiler>,
sink: Arc<dyn SnapshotSink>,
status: Arc<RevisionStatus>,
settings: ConvergenceSettings,
cache: Option<LastKnownGood>,
clock: Arc<dyn Clock>,
active: Mutex<Option<RevisionId>>,
attempt_lock: AsyncMutex<()>,
refresh_pending: AtomicBool,
backoff: Mutex<Backoff>,
export_failing: AtomicBool,
}
impl Reconciler {
pub fn new(
store: Arc<dyn ControlPlaneStore>,
compiler: Arc<dyn CandidateCompiler>,
sink: Arc<dyn SnapshotSink>,
settings: ConvergenceSettings,
cache: Option<LastKnownGood>,
clock: Arc<dyn Clock>,
) -> Self {
let status = Arc::new(RevisionStatus::new(Box::new(ArcClock(Arc::clone(&clock)))));
Self {
store,
compiler,
sink,
status,
backoff: Mutex::new(Backoff::new(settings.backoff)),
settings,
cache,
clock,
active: Mutex::new(None),
attempt_lock: AsyncMutex::new(()),
refresh_pending: AtomicBool::new(false),
export_failing: AtomicBool::new(false),
}
}
pub fn report(&self) -> RevisionReport {
self.status.report()
}
pub fn status(&self) -> &Arc<RevisionStatus> {
&self.status
}
pub async fn bootstrap(&self) -> Result<RevisionId, BootstrapError> {
let span = telemetry::revision_convergence_span(telemetry::CONVERGENCE_BOOT);
let result = self.bootstrap_inner().await;
let outcome = match &result {
Ok(revision) => Outcome::Published {
revision: *revision,
generation: self.status.report().generation,
took: self.status.report().last_convergence.unwrap_or_default(),
},
Err(BootstrapError::Empty) => Outcome::Empty,
Err(_) => Outcome::Rejected {
revision: None,
reason: self
.status
.report()
.last_rejection
.map_or("boot", |rejection| rejection.reason),
},
};
telemetry::finish_revision_convergence(
&span,
telemetry::CONVERGENCE_BOOT,
&outcome,
&self.status.report(),
);
result
}
async fn bootstrap_inner(&self) -> Result<RevisionId, BootstrapError> {
match self.attempt().await {
Ok(Some(published)) => Ok(published),
Ok(None) => Err(BootstrapError::Empty),
Err(error) => {
self.record_failure(&error);
match error {
AttemptError::Store(source)
if source.retryable()
|| matches!(source, ControlPlaneError::Incompatible { .. }) =>
{
self.restore_from_cache(source).await
}
AttemptError::Store(source) => Err(BootstrapError::Store { source }),
AttemptError::Compile(source) => Err(BootstrapError::Rejected {
source: Box::new(source),
}),
}
}
}
}
pub async fn converge_once(&self, trigger: &'static str) -> Outcome {
self.converge_once_with_refresh(trigger, false).await
}
pub async fn force_refresh_once(&self, trigger: &'static str) -> Outcome {
self.converge_once_with_refresh(trigger, true).await
}
async fn converge_once_with_refresh(
&self,
trigger: &'static str,
force_refresh: bool,
) -> Outcome {
if force_refresh {
self.refresh_pending.store(true, Ordering::Release);
}
let span = telemetry::revision_convergence_span(trigger);
let outcome = match self.attempt().await {
Ok(Some(revision)) => {
let report = self.status.report();
Outcome::Published {
revision,
generation: report.generation,
took: report.last_convergence.unwrap_or_default(),
}
}
Ok(None) => match *self.active.lock().expect("not poisoned") {
Some(revision) => Outcome::AlreadyConverged {
revision: Some(revision),
},
None => Outcome::Empty,
},
Err(error) => {
let reason = self.record_failure(&error);
Outcome::Rejected {
revision: error.revision(),
reason,
}
}
};
telemetry::finish_revision_convergence(&span, trigger, &outcome, &self.status.report());
outcome
}
pub async fn run(&self, signal: Arc<ChangeSignal>, shutdown: impl Future<Output = ()> + Send) {
let shutdown = std::pin::pin!(shutdown);
let mut shutdown = shutdown;
loop {
let delay = {
let backoff = self.backoff.lock().expect("not poisoned");
if backoff.failures() == 0 {
self.settings.poll_interval
} else {
backoff.delay()
}
};
tokio::select! {
biased;
() = &mut shutdown => {
tracing::debug!("revision convergence stopped");
return;
}
force_refresh = signal.notified() => {
self.converge_once_with_refresh(
telemetry::CONVERGENCE_NOTIFIED,
force_refresh,
).await;
}
() = tokio::time::sleep(delay) => {
self.converge_once(telemetry::CONVERGENCE_POLLED).await;
}
}
}
}
async fn attempt(&self) -> Result<Option<RevisionId>, AttemptError> {
let _attempt = self.attempt_lock.lock().await;
let force_refresh = self.refresh_pending.swap(false, Ordering::AcqRel);
let result = self.attempt_inner(force_refresh).await;
if force_refresh && result.is_err() {
self.refresh_pending.store(true, Ordering::Release);
}
result
}
async fn attempt_inner(&self, force_refresh: bool) -> Result<Option<RevisionId>, AttemptError> {
let started = self.clock.now();
let desired = self.store.desired_revision().await?;
self.status.observe_desired(desired);
let active = *self.active.lock().expect("not poisoned");
if desired.is_none() || (!force_refresh && desired == active) {
self.backoff.lock().expect("not poisoned").succeed();
return Ok(None);
}
let Some(revision) = self.store.load_desired_revision().await? else {
self.backoff.lock().expect("not poisoned").succeed();
return Ok(None);
};
self.status.observe_desired(Some(revision.id()));
self.publish(revision, SnapshotSource::ControlPlane, started)
.await
.map(Some)
.map_err(AttemptError::from)
}
async fn publish(
&self,
revision: LoadedRevision,
source: SnapshotSource,
started: std::time::Instant,
) -> Result<RevisionId, CompileError> {
let id = revision.id();
let generation = self.sink.generation().saturating_add(1);
let snapshot = self.compiler.compile(&revision, generation).await?;
self.sink
.admit(&snapshot)
.map_err(|source| CompileError::Activation {
revision: id,
source,
})?;
self.status.observe_loaded(id);
let materialized = snapshot.secrets().len();
let pricing = snapshot.pricing().map(|pricing| {
(
pricing.book(),
pricing.checksum(),
pricing.catalog(),
pricing.is_approved(),
pricing.targets().len(),
pricing.effective().ends(),
)
});
self.sink.publish(snapshot);
*self.active.lock().expect("not poisoned") = Some(id);
self.backoff.lock().expect("not poisoned").succeed();
let took = self.clock.now().saturating_duration_since(started);
self.status.record_published(id, generation, source, took);
tracing::info!(
revision = %id,
generation,
source = source.as_str(),
took_ms = took.as_millis(),
materialized,
"published desired revision"
);
if let Some((book, checksum, catalog, approved, targets, until)) = pricing {
tracing::info!(
revision = %id,
generation,
price_book = %book,
price_book_checksum = %checksum,
catalog = %catalog,
approved,
priced_targets = targets,
priced_until = until.map(|until| until.millis()),
"published approved pricing"
);
}
self.export(&revision);
Ok(id)
}
fn export(&self, revision: &LoadedRevision) {
let Some(cache) = &self.cache else {
return;
};
match cache.export(revision) {
Ok(()) => {
if self.export_failing.swap(false, Ordering::Relaxed) {
tracing::info!(
path = %cache.path().display(),
"last-known-good snapshot is writable again"
);
}
telemetry::record_last_known_good("exported");
}
Err(error) => {
telemetry::record_last_known_good("export_failed");
if !self.export_failing.swap(true, Ordering::Relaxed) {
tracing::warn!(
path = %cache.path().display(),
error = %error,
"the last-known-good snapshot could not be written; the replica keeps \
serving, but a cold boot during a control-plane outage will have no \
cached state"
);
}
}
}
}
async fn restore_from_cache(
&self,
source: ControlPlaneError,
) -> Result<RevisionId, BootstrapError> {
let Some(cache) = &self.cache else {
return Err(Self::uncached(source));
};
let restored = match cache.load() {
Ok(restored) => restored,
Err(LastKnownGoodError::Integrity(integrity)) if integrity.is_incompatible() => {
telemetry::record_last_known_good(INCOMPATIBLE_REASON);
tracing::warn!(
error = %integrity,
"the last-known-good snapshot was written by a build this one cannot read; \
it cannot stand in for desired state"
);
return Err(Self::uncached(source));
}
Err(source) => {
return Err(BootstrapError::Cache {
source: Box::new(source),
});
}
};
let Some(revision) = restored else {
return Err(Self::uncached(source));
};
let started = self.clock.now();
let id = self
.publish(revision, SnapshotSource::LastKnownGood, started)
.await
.map_err(|source| BootstrapError::Rejected {
source: Box::new(source),
})?;
telemetry::record_last_known_good("restored");
tracing::warn!(
revision = %id,
error = %source,
"desired state is not usable on this replica; booted from the signed \
last-known-good snapshot, which may be older than desired state"
);
Ok(id)
}
fn uncached(source: ControlPlaneError) -> BootstrapError {
if source.retryable() {
BootstrapError::Unavailable { source }
} else {
BootstrapError::Store { source }
}
}
fn record_failure(&self, error: &AttemptError) -> &'static str {
let reason = error.reason();
let (failures, delay) = {
let mut backoff = self.backoff.lock().expect("not poisoned");
let delay = backoff.fail();
(backoff.failures(), delay)
};
self.status.record_rejection(
Rejection {
revision: error.revision(),
reason,
detail: error.to_string(),
},
failures,
);
telemetry::record_revision_rejection(reason);
tracing::warn!(
reason,
failures,
retry_in_ms = delay.as_millis(),
error = %error,
"desired revision was not applied; the active revision keeps serving"
);
reason
}
}
struct ArcClock(Arc<dyn Clock>);
impl Clock for ArcClock {
fn now(&self) -> std::time::Instant {
self.0.now()
}
}