use std::sync::Arc;
use super::types::{AsyncRaftProposer, ReplicatedEntry};
use crate::control::state::SharedState;
const BACKOFF_MS: [u64; 5] = [10, 25, 50, 100, 200];
pub async fn propose_replicated_entry(
state: &SharedState,
proposer: &Arc<AsyncRaftProposer>,
entry: ReplicatedEntry,
) -> crate::Result<(Vec<u8>, crate::types::Lsn)> {
let idempotency_key = entry.idempotency_key;
let data = entry.to_bytes();
let vshard_id = entry.vshard_id;
let mut payload = None;
let mut last_err: Option<crate::Error> = None;
for (attempt, backoff_ms) in BACKOFF_MS.iter().enumerate() {
match proposer(vshard_id, idempotency_key, data.clone()).await {
Ok(p) => {
payload = Some(p);
break;
}
Err(crate::Error::RetryableLeaderChange {
group_id,
log_index,
}) => {
state
.raft_propose_leader_change_retries
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::warn!(
attempt,
group_id,
log_index,
"raft entry overwritten by leader change — re-proposing"
);
last_err = Some(crate::Error::RetryableLeaderChange {
group_id,
log_index,
});
tokio::time::sleep(std::time::Duration::from_millis(*backoff_ms)).await;
continue;
}
Err(other) => {
return Err(crate::Error::Dispatch {
detail: format!("raft propose failed: {other}"),
});
}
}
}
payload.ok_or_else(|| {
last_err.unwrap_or_else(|| crate::Error::Dispatch {
detail: "raft propose retries exhausted".into(),
})
})
}