use crate::error::store_error;
use crate::records::{self, LeaderVal, PlanFinalityRepr, SplitProgressRecord, SplitSpecRecord};
use crate::store::{CasOutcome, CoordinationStore, Keyspace, Revision};
use crate::task::Task;
use spate_core::coordination::{
CoordinationError, CoordinationErrorKind, PlanContext, PlanFinality, SplitPlan, SplitPlanner,
};
use spate_core::metrics::ReplanOutcome;
use tokio::time::Instant;
pub(crate) type PlannerOutput = (Box<dyn SplitPlanner>, Result<SplitPlan, CoordinationError>);
pub(crate) struct PlanRun {
pub(crate) handle: tokio::task::JoinHandle<PlannerOutput>,
plan: records::PlanRecord,
plan_rev: Revision,
generation: u64,
started: Instant,
}
impl<S: CoordinationStore> Task<S> {
pub(crate) async fn try_elect(&mut self) -> Result<(), CoordinationError> {
let generation = self.plan.as_ref().map_or(0, |(p, _)| p.generation) + 1;
let val = records::encode_val(&LeaderVal {
schema: records::SCHEMA,
owner: self.instance.clone(),
nonce: self.nonce.clone(),
generation,
});
match self
.store
.create(Keyspace::Ephemeral, records::LEADER_KEY, val)
.await
{
Ok(CasOutcome::Won(rev)) => {
tracing::info!(generation, "elected planner leader");
self.leadership = Some(rev);
self.metrics(|m| m.set_leader(true));
self.mark_assignment_dirty();
self.bump_generation(generation).await?;
Ok(())
}
Ok(CasOutcome::Lost) => Ok(()), Err(e) => {
tracing::warn!(error = %e, "election write failed; retrying on observation");
Ok(())
}
}
}
async fn bump_generation(&mut self, generation: u64) -> Result<(), CoordinationError> {
for _ in 0..3 {
let Some((plan, rev)) = &self.plan else {
return Ok(()); };
if plan.generation >= generation {
self.demote().await;
return Ok(());
}
let mut bumped = plan.clone();
bumped.generation = generation;
bumped.updated_at_ms = records::now_ms();
match self
.store
.update(Keyspace::Durable, records::PLAN_KEY, bumped.encode(), *rev)
.await
{
Ok(CasOutcome::Won(new_rev)) => {
self.plan_rev_seen = self.plan_rev_seen.max(new_rev.0);
self.plan = Some((bumped, new_rev));
self.plan_now = true;
return Ok(());
}
Ok(CasOutcome::Lost) => {
let entry = self
.store
.get(Keyspace::Durable, records::PLAN_KEY)
.await
.map_err(|e| store_error("re-reading the plan record", &e))?;
let Some(entry) = entry else {
return Err(crate::error::fatal(
"plan record vanished mid-election; the store prefix was \
tampered with",
));
};
let plan = records::PlanRecord::parse(&entry.value, &self.fingerprint)?;
self.plan_rev_seen = self.plan_rev_seen.max(entry.revision.0);
self.plan = Some((plan, entry.revision));
}
Err(e) if matches!(e, crate::store::StoreError::Retryable(_)) => {
tracing::warn!(error = %e, "generation bump failed; retrying");
}
Err(e) => return Err(store_error("bumping the plan generation", &e)),
}
}
self.demote().await;
Ok(())
}
pub(crate) async fn demote(&mut self) {
if let Some(rev) = self.leadership.take() {
self.metrics(|m| m.set_leader(false));
let _ = self
.store
.delete(Keyspace::Ephemeral, records::LEADER_KEY, Some(rev))
.await;
}
}
pub(crate) fn maybe_start_plan(&mut self) -> Result<Option<PlanRun>, CoordinationError> {
if !self.plan_now || self.leadership.is_none() {
return Ok(None);
}
self.plan_now = false;
let Some((plan, plan_rev)) = self.plan.clone() else {
return Ok(None);
};
let generation = plan.generation;
let cursor: Option<Vec<u8>> = match &plan.planner_state {
Some(encoded) => Some(records::b64_decode(encoded).map_err(|e| {
crate::error::fatal(format!("plan record: corrupt planner cursor ({e})"))
})?),
None => None,
};
let Some(mut planner) = self.planner.take() else {
return Ok(None); };
let started = Instant::now();
let handle = tokio::task::spawn_blocking(move || {
let ctx = PlanContext::new(cursor.as_deref(), generation);
let result = planner.plan(ctx);
(planner, result)
});
Ok(Some(PlanRun {
handle,
plan,
plan_rev,
generation,
started,
}))
}
pub(crate) async fn finish_plan(
&mut self,
joined: Result<PlannerOutput, tokio::task::JoinError>,
run: PlanRun,
) -> Result<(), CoordinationError> {
let (planner, result) = match joined {
Ok(parts) => parts,
Err(join_error) => {
return Err(crate::error::fatal(format!(
"the planner panicked: {join_error}"
)));
}
};
self.planner = Some(planner);
let split_plan = match result {
Ok(plan) => plan,
Err(e) if e.kind == CoordinationErrorKind::Retryable => {
tracing::warn!(error = %e, "planner failed; next replan tick retries");
self.metrics(|m| m.replan(ReplanOutcome::Error, run.started.elapsed()));
return Ok(());
}
Err(e) => return Err(e),
};
let mut created = 0u64;
for planned in &split_plan.splits {
let spec_record = SplitSpecRecord::planned(&planned.spec, self.fp, run.generation);
match self
.store
.create(
Keyspace::Durable,
&records::spec_key(&planned.spec.id),
spec_record.encode(),
)
.await
{
Ok(_) => {} Err(e) => {
tracing::warn!(split = %planned.spec.id, error = %e,
"spec seeding failed; next replan tick retries");
self.metrics(|m| m.replan(ReplanOutcome::Error, run.started.elapsed()));
return Ok(());
}
}
let progress =
SplitProgressRecord::planned(&planned.spec.id, self.fp, planned.seed.as_ref());
match self
.store
.create(
Keyspace::Durable,
&records::split_key(&planned.spec.id),
progress.encode(),
)
.await
{
Ok(CasOutcome::Won(rev)) => {
created += 1;
self.attach_spec(planned.spec.id.as_str(), spec_record);
self.upsert_progress(planned.spec.id.as_str(), progress, rev)?;
}
Ok(CasOutcome::Lost) => {} Err(e) => {
tracing::warn!(split = %planned.spec.id, error = %e,
"split seeding failed; next replan tick retries");
self.metrics(|m| m.replan(ReplanOutcome::Error, run.started.elapsed()));
return Ok(());
}
}
}
self.metrics(|m| m.planned(created));
let listed = match self
.store
.list(Keyspace::Durable, records::SPLIT_PREFIX)
.await
{
Ok(entries) => entries.len() as u64,
Err(e) => {
tracing::warn!(error = %e, "planned recount failed; next replan tick retries");
self.metrics(|m| m.replan(ReplanOutcome::Error, run.started.elapsed()));
return Ok(());
}
};
let finality = PlanFinalityRepr::from(split_plan.finality);
let finality_changed = run.plan.finality != finality;
let count_changed = run.plan.planned != listed;
let mut published = run.plan.clone();
published.planned = listed;
published.finality = finality;
if let Some(state) = &split_plan.planner_state {
published.planner_state = Some(records::b64_encode(state));
}
published.updated_at_ms = records::now_ms();
match self
.store
.update(
Keyspace::Durable,
records::PLAN_KEY,
published.encode(),
run.plan_rev,
)
.await
{
Ok(CasOutcome::Won(rev)) => {
self.plan_rev_seen = self.plan_rev_seen.max(rev.0);
self.plan = Some((published, rev));
let outcome = if created > 0 || finality_changed || count_changed {
ReplanOutcome::Ok
} else {
ReplanOutcome::Noop
};
self.metrics(|m| m.replan(outcome, run.started.elapsed()));
if split_plan.finality == PlanFinality::Final {
tracing::info!(
planned = self.plan.as_ref().map_or(0, |(p, _)| p.planned),
"plan is final"
);
}
Ok(())
}
Ok(CasOutcome::Lost) => {
tracing::warn!("plan publish fenced; a successor leads");
self.metrics(|m| m.replan(ReplanOutcome::Error, run.started.elapsed()));
self.demote().await;
Ok(())
}
Err(e) => {
tracing::warn!(error = %e, "plan publish failed; next replan tick retries");
self.metrics(|m| m.replan(ReplanOutcome::Error, run.started.elapsed()));
Ok(())
}
}
}
}