use crate::metadata_store::{RabsMetadataStore, SqlEngine, SqlMetadataStore, SqlValue, StoreError};
use crate::startup_reconciliation::{
Drift, FilesystemReality, ServingDecision, StartupReport, reconcile_startup,
};
use rabs_protocol::generation::WorkerIncarnationId;
use rabs_protocol::wire_time::PeerId;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OpenSession {
pub worker: String,
pub incarnation_hex: String,
pub incarnation: WorkerIncarnationId,
pub started_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StaleOperation {
pub id_hex: String,
pub kind: String,
pub state: String,
pub updated_seq: u64,
pub lag: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProposedResolution {
pub target: String,
pub action: String,
pub remediation: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WorkerReconcileReport {
pub worker: String,
pub sessions_ended: u32,
pub repaired: Vec<Drift>,
pub reported_orphans: Vec<String>,
pub serving: ServingDecision,
pub stale_operations: Vec<StaleOperation>,
pub proposals: Vec<ProposedResolution>,
}
fn text_at(row: &[SqlValue], i: usize) -> String {
match row.get(i) {
Some(SqlValue::Text(t)) => t.clone(),
Some(other) => format!("{other:?}"),
None => String::new(),
}
}
fn int_at(row: &[SqlValue], i: usize) -> u64 {
match row.get(i) {
Some(SqlValue::Int(v)) => (*v).max(0) as u64,
_ => 0,
}
}
fn open_sessions_for_worker<E: SqlEngine>(
store: &mut SqlMetadataStore<E>,
worker: &str,
) -> Result<Vec<OpenSession>, StoreError> {
let rows = store.engine_mut().query(
"SELECT worker, incarnation, started_seq FROM worker_sessions \
WHERE worker = ?1 AND ended_seq IS NULL",
&[SqlValue::Text(worker.to_owned())],
)?;
rows.into_iter()
.map(|row| {
let Some(SqlValue::Blob(bytes)) = row.get(1) else {
return Err(StoreError::Corruption(
"worker session incarnation shape".into(),
));
};
let incarnation_bytes: [u8; 16] = bytes.as_slice().try_into().map_err(|_| {
StoreError::Corruption("worker session incarnation not 16 bytes".into())
})?;
Ok(OpenSession {
worker: text_at(&row, 0),
incarnation_hex: bytes.iter().map(|b| format!("{b:02x}")).collect(),
incarnation: WorkerIncarnationId(u128::from_be_bytes(incarnation_bytes)),
started_seq: int_at(&row, 2),
})
})
.collect()
}
fn stale_operations<E: SqlEngine>(
store: &mut SqlMetadataStore<E>,
min_seq_lag: u64,
) -> Result<Vec<StaleOperation>, StoreError> {
let rows = store.engine_mut().query(
"SELECT id_hex, kind, state, updated_seq, \
(SELECT COALESCE(MAX(updated_seq), 0) FROM operations \
WHERE state NOT IN ('committed', 'failed', 'abandoned')) - updated_seq AS lag \
FROM operations \
WHERE state NOT IN ('committed', 'failed', 'abandoned') \
ORDER BY lag DESC",
&[],
)?;
Ok(rows
.into_iter()
.filter_map(|row| {
let lag = int_at(&row, 4);
if lag <= min_seq_lag {
return None;
}
Some(StaleOperation {
id_hex: text_at(&row, 0),
kind: text_at(&row, 1),
state: text_at(&row, 2),
updated_seq: int_at(&row, 3),
lag,
})
})
.collect())
}
fn expired_looking_pins<E: SqlEngine>(
store: &mut SqlMetadataStore<E>,
now_seq: u64,
) -> Result<Vec<(String, String, u64)>, StoreError> {
let rows = store.engine_mut().query(
"SELECT id_hex, owner, expires_at_seq FROM pins \
WHERE released = 0 AND expires_at_seq IS NOT NULL AND expires_at_seq < ?1",
&[SqlValue::Int(now_seq as i64)],
)?;
Ok(rows
.into_iter()
.map(|row| (text_at(&row, 0), text_at(&row, 1), int_at(&row, 2)))
.collect())
}
fn operation_proposal(op: &StaleOperation) -> ProposedResolution {
ProposedResolution {
target: format!("operation:{}", op.id_hex),
action: format!("update_operation_state({}, \"abandoned\")", op.id_hex),
remediation: format!(
"non-terminal '{}' untouched for {} seqs; abandoning writes one \
row and releases nothing",
op.state, op.lag
),
}
}
fn pin_proposal(pin_hex: &str, owner: &str, expires_at: u64, now_seq: u64) -> ProposedResolution {
ProposedResolution {
target: format!("pin:{pin_hex}"),
action: format!("release_pin({pin_hex}, owner={owner:?})"),
remediation: format!(
"expires_at_seq {expires_at} < current {now_seq}; H041: stays \
protecting until reconciliation is CONFIRMED — confirm first"
),
}
}
pub fn reconcile_worker<E: SqlEngine>(
store: &mut SqlMetadataStore<E>,
filesystem: &dyn FilesystemReality,
worker: &str,
now_seq: u64,
min_operation_seq_lag: u64,
) -> Result<WorkerReconcileReport, StoreError> {
let sessions = open_sessions_for_worker(store, worker)?;
let active_authority = store.active_authority()?.map(|row| row.digest);
if !sessions.is_empty() && active_authority.is_none() {
return Err(StoreError::NotActiveAuthority);
}
let worker_peer_id = PeerId(worker.to_owned());
let mut sessions_ended = 0u32;
for session in sessions {
let authority = active_authority
.as_ref()
.ok_or(StoreError::NotActiveAuthority)?;
let ended = if store.release_worker_session(
authority,
&worker_peer_id,
session.incarnation,
session.started_seq,
now_seq,
)? {
true
} else {
store.end_worker_session(&session.worker, session.started_seq, now_seq)?
};
if ended {
sessions_ended += 1;
}
}
let startup: StartupReport = reconcile_startup(store, filesystem)?;
let mut reported_orphans = Vec::new();
for drift in &startup.reported {
if let Drift::OrphanPathReported { store_path } = drift {
reported_orphans.push(store_path.clone());
}
}
let stale = stale_operations(store, min_operation_seq_lag)?;
let mut proposals: Vec<ProposedResolution> = stale.iter().map(operation_proposal).collect();
for (pin_hex, owner, expires_at) in expired_looking_pins(store, now_seq)? {
proposals.push(pin_proposal(&pin_hex, &owner, expires_at, now_seq));
}
Ok(WorkerReconcileReport {
worker: worker.to_owned(),
sessions_ended,
repaired: startup.repaired,
reported_orphans,
serving: startup.serving,
stale_operations: stale,
proposals,
})
}
pub fn find_stale_state<E: SqlEngine>(
store: &mut SqlMetadataStore<E>,
now_seq: u64,
min_operation_seq_lag: u64,
) -> Result<Vec<ProposedResolution>, StoreError> {
let mut proposals: Vec<ProposedResolution> = stale_operations(store, min_operation_seq_lag)?
.iter()
.map(operation_proposal)
.collect();
for (pin_hex, owner, expires_at) in expired_looking_pins(store, now_seq)? {
proposals.push(pin_proposal(&pin_hex, &owner, expires_at, now_seq));
}
Ok(proposals)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{AuthorityRow, RusqliteEngine, SqlMetadataStore};
use crate::publication::authority_digest;
use crate::startup_reconciliation::SetFilesystem;
use rabs_protocol::authority::{ClusterId, CoordinatorAuthority, CoordinatorIncarnationId};
use rabs_protocol::generation::WorkerBootGeneration;
use rabs_protocol::worker_fence::{WorkerAdmission, WorkerSessionOffer};
fn seeded_store() -> SqlMetadataStore<RusqliteEngine> {
let mut store =
SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).expect("store");
let coordinator = CoordinatorAuthority {
cluster_id: ClusterId("cluster-a".to_owned()),
credential_generation: 1,
term: 3,
incarnation_id: CoordinatorIncarnationId(77),
};
store
.acquire_authority(&AuthorityRow {
digest: authority_digest(&coordinator),
cluster_id: "cluster-a".to_owned(),
incarnation: 77,
term: 3,
acquired_seq: 1,
})
.expect("acquire authority");
store
.record_worker_session("worker-a", 0x1111, 10)
.expect("session 1");
assert_eq!(
store
.admit_worker_session(
&authority_digest(&coordinator),
&WorkerSessionOffer {
worker_peer_id: PeerId("worker-a".to_owned()),
boot_generation: WorkerBootGeneration(1),
incarnation: WorkerIncarnationId(0x2222),
reenrollment_proof: None,
},
20,
)
.expect("session 2"),
WorkerAdmission::AdmitNewGeneration
);
store
.create_operation(
&authority_digest(&coordinator),
0xAA00,
"build",
"running",
100,
)
.expect("stale operation");
store
.create_operation(
&authority_digest(&coordinator),
0xAB00,
"build",
"running",
5_000,
)
.expect("watermark operation");
store
}
#[test]
fn reconcile_worker_ends_sessions_and_reports_stale_operations() {
let mut store = seeded_store();
let fs = SetFilesystem::default();
let report = reconcile_worker(&mut store, &fs, "worker-a", 6_000, 100).expect("reconcile");
assert_eq!(report.sessions_ended, 2, "both open sessions ended");
assert!(
open_sessions_for_worker(&mut store, "worker-a")
.unwrap()
.is_empty(),
"no open session rows survive"
);
assert_eq!(
store
.worker_incarnation_fence(&PeerId("worker-a".to_owned()))
.unwrap()
.unwrap()
.active_incarnation,
None,
"reconciliation clears only the exact active incarnation"
);
assert_eq!(report.stale_operations.len(), 1);
let stale = &report.stale_operations[0];
assert_eq!(stale.id_hex, u128_hex(0xAA00));
assert_eq!(stale.lag, 4_900);
assert!(
report
.proposals
.iter()
.any(|p| p.target == format!("operation:{}", u128_hex(0xAA00))),
"every finding gets a named safe-resolution proposal"
);
}
#[test]
fn reconcile_without_authority_preserves_journal_and_fence() {
let mut store = seeded_store();
let active = store
.active_authority()
.expect("authority query")
.expect("seeded authority")
.digest;
store.release_authority(&active).expect("release authority");
assert_eq!(
reconcile_worker(
&mut store,
&SetFilesystem::default(),
"worker-a",
6_000,
100,
),
Err(StoreError::NotActiveAuthority)
);
assert_eq!(
open_sessions_for_worker(&mut store, "worker-a")
.expect("open sessions")
.len(),
2,
"failed recovery must not end only the journals"
);
assert_eq!(
store
.worker_incarnation_fence(&PeerId("worker-a".to_owned()))
.expect("worker fence")
.expect("seeded fence")
.active_incarnation,
Some(WorkerIncarnationId(0x2222)),
"failed recovery must leave the matching fence untouched"
);
}
#[test]
fn applying_the_proposal_resolves_the_finding() {
let mut store = seeded_store();
let _fs = SetFilesystem::default();
let before = find_stale_state(&mut store, 6_000, 100).expect("doctor");
assert_eq!(before.len(), 1, "seeded staleness is found");
store
.update_operation_state(0xAA00, "abandoned", 6_000)
.expect("apply proposal");
let after = find_stale_state(&mut store, 6_100, 100).expect("doctor again");
assert!(
after.is_empty(),
"seeded stale state is FOUND + RESOLVED: {after:?}"
);
}
#[test]
fn reconcile_pass_is_idempotent_for_sessions_and_doctor_finds_nothing_else() {
let mut store = seeded_store();
let fs = SetFilesystem::default();
let first = reconcile_worker(&mut store, &fs, "worker-a", 6_000, 100).unwrap();
assert_eq!(first.sessions_ended, 2);
let second = reconcile_worker(&mut store, &fs, "worker-a", 6_500, 100).unwrap();
assert_eq!(second.sessions_ended, 0, "no sessions left to end");
assert!(matches!(
second.serving,
ServingDecision::Allowed | ServingDecision::Refused(_)
));
}
fn u128_hex(v: u128) -> String {
format!("{v:032x}")
}
}