use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use async_trait::async_trait;
use tokio::sync::Notify;
use tokio::time::Instant;
use car_feedback_core::spool::{
IdentityLane, Spool, SpoolEntryId, SpoolEntrySummary, SpoolState, TerminalReason,
};
use car_parslee::feedback_transport::{
DrainAction, DrainEligibility, FeedbackTransport, FeedbackTransportError, SubmitOutcome,
TransportActionableReason,
};
use crate::feedback::FEEDBACK_OUTBOX_DIR;
use crate::session::ServerState;
#[derive(Debug, Clone)]
pub struct DrainConfig {
pub min_upload_interval: Duration,
pub initial_backoff: Duration,
pub max_backoff: Duration,
pub held_recheck_interval: Duration,
pub park_check_interval: Duration,
pub max_retriable_attempts: u32,
}
impl Default for DrainConfig {
fn default() -> Self {
DrainConfig {
min_upload_interval: Duration::from_secs(15),
initial_backoff: Duration::from_secs(30),
max_backoff: Duration::from_secs(15 * 60),
held_recheck_interval: Duration::from_secs(15 * 60),
park_check_interval: Duration::from_secs(60),
max_retriable_attempts: 24,
}
}
}
#[derive(Debug, Default)]
pub struct RetryLedger {
failures: HashMap<SpoolEntryId, u32>,
}
impl RetryLedger {
fn record_failure(&mut self, id: &SpoolEntryId, cap: u32) {
let count = self.failures.entry(id.clone()).or_insert(0);
*count = count.saturating_add(1);
if *count == cap {
tracing::warn!(
target: "car::feedback",
entry = %id, attempts = cap,
"feedback entry hit this daemon's retriable-attempt cap; it stays queued \
(export it with `car feedback --export`) and retries again after the \
daemon restarts"
);
}
}
fn is_exhausted(&self, id: &SpoolEntryId, cap: u32) -> bool {
self.failures.get(id).is_some_and(|count| *count >= cap)
}
fn retain_queued(&mut self, queued: &[SpoolEntrySummary]) {
self.failures
.retain(|id, _| queued.iter().any(|entry| &entry.id == id));
}
}
#[async_trait]
pub trait DrainTransport: Send + Sync {
async fn drain_eligible(&self, lane: &IdentityLane) -> DrainEligibility;
async fn submit(
&self,
bundle: &car_feedback_core::bundle::RedactedBundle,
lane: &IdentityLane,
client_submission_id: &str,
) -> Result<SubmitOutcome, FeedbackTransportError>;
}
#[async_trait]
impl DrainTransport for FeedbackTransport {
async fn drain_eligible(&self, lane: &IdentityLane) -> DrainEligibility {
FeedbackTransport::drain_eligible(self, lane).await
}
async fn submit(
&self,
bundle: &car_feedback_core::bundle::RedactedBundle,
lane: &IdentityLane,
client_submission_id: &str,
) -> Result<SubmitOutcome, FeedbackTransportError> {
FeedbackTransport::submit_report(self, bundle, lane, client_submission_id).await
}
}
struct LazyLiveTransport {
inner: tokio::sync::OnceCell<FeedbackTransport>,
#[cfg(test)]
constructor_error: Option<String>,
}
impl LazyLiveTransport {
fn new() -> Self {
Self {
inner: tokio::sync::OnceCell::new(),
#[cfg(test)]
constructor_error: None,
}
}
#[cfg(test)]
fn failing(error: &str) -> Self {
Self {
inner: tokio::sync::OnceCell::new(),
constructor_error: Some(error.to_string()),
}
}
async fn get(&self) -> Result<&FeedbackTransport, String> {
#[cfg(test)]
if let Some(error) = &self.constructor_error {
return Err(error.clone());
}
self.inner
.get_or_try_init(|| async { FeedbackTransport::live() })
.await
}
}
#[async_trait]
impl DrainTransport for LazyLiveTransport {
async fn drain_eligible(&self, lane: &IdentityLane) -> DrainEligibility {
match self.get().await {
Ok(transport) => DrainTransport::drain_eligible(transport, lane).await,
Err(error) => DrainEligibility::Hold {
reason: format!("feedback transport unavailable: {error}"),
},
}
}
async fn submit(
&self,
bundle: &car_feedback_core::bundle::RedactedBundle,
lane: &IdentityLane,
client_submission_id: &str,
) -> Result<SubmitOutcome, FeedbackTransportError> {
match self.get().await {
Ok(transport) => {
DrainTransport::submit(transport, bundle, lane, client_submission_id).await
}
Err(error) => Err(FeedbackTransportError::FetchFailed(format!(
"feedback transport unavailable: {error}"
))),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RoundOutcome {
Idle,
AllHeld,
Progressed,
Backoff { min_delay: Duration },
RetryAfter(Duration),
}
#[derive(Clone)]
pub struct FeedbackDrainHandle {
notify: Arc<Notify>,
}
impl FeedbackDrainHandle {
pub fn wake(&self) {
self.notify.notify_one();
}
}
static DRAIN_HANDLE: OnceLock<FeedbackDrainHandle> = OnceLock::new();
pub fn wake_feedback_drain() {
if let Some(handle) = DRAIN_HANDLE.get() {
handle.wake();
}
}
pub fn spawn_feedback_drain(state: &ServerState) {
let car_home = match crate::feedback::car_home_dir(state) {
Ok(home) => home,
Err(e) => {
tracing::warn!(target: "car::feedback", error = %e, "feedback drain not started");
return;
}
};
let transport: Arc<dyn DrainTransport> = Arc::new(LazyLiveTransport::new());
let handle = spawn_feedback_drain_with(car_home, transport, DrainConfig::default());
let _ = DRAIN_HANDLE.set(handle);
}
pub fn spawn_feedback_drain_with(
car_home: PathBuf,
transport: Arc<dyn DrainTransport>,
config: DrainConfig,
) -> FeedbackDrainHandle {
let notify = Arc::new(Notify::new());
let handle = FeedbackDrainHandle {
notify: notify.clone(),
};
let spool_root = car_home.join(FEEDBACK_OUTBOX_DIR);
tokio::spawn(async move {
let mut last_upload: Option<Instant> = None;
let mut ledger = RetryLedger::default();
loop {
let mut backoff = config.initial_backoff;
loop {
let outcome = run_drain_round(
&spool_root,
transport.as_ref(),
&config,
&mut last_upload,
&mut ledger,
)
.await;
match outcome {
RoundOutcome::Idle => break,
RoundOutcome::Progressed => {
backoff = config.initial_backoff;
}
RoundOutcome::AllHeld => {
wait_or_wake(¬ify, config.held_recheck_interval).await;
}
RoundOutcome::Backoff { min_delay } => {
let delay = backoff.max(min_delay).min(config.max_backoff);
wait_or_wake(¬ify, delay).await;
backoff = (backoff * 2).min(config.max_backoff);
}
RoundOutcome::RetryAfter(delay) => {
let delay = delay.min(config.max_backoff);
wait_or_wake(¬ify, delay).await;
}
}
}
loop {
tokio::select! {
_ = notify.notified() => break,
_ = tokio::time::sleep(config.park_check_interval) => {
if spool_has_pending_entries(&spool_root) {
break;
}
}
}
}
}
});
handle
}
fn spool_has_pending_entries(spool_root: &Path) -> bool {
let Ok(dirents) = std::fs::read_dir(spool_root) else {
return false;
};
let has_published_dir = dirents.flatten().any(|dirent| {
!dirent.file_name().to_string_lossy().starts_with(".tmp-")
&& dirent.file_type().map(|t| t.is_dir()).unwrap_or(false)
});
if !has_published_dir {
return false;
}
let Ok(spool) = Spool::open(spool_root) else {
return false;
};
spool
.list()
.map(|rows| {
rows.iter()
.any(|row| matches!(row.state, SpoolState::Queued | SpoolState::Sending))
})
.unwrap_or(false)
}
async fn wait_or_wake(notify: &Notify, delay: Duration) {
tokio::select! {
_ = tokio::time::sleep(delay) => {}
_ = notify.notified() => {}
}
}
pub async fn run_drain_round(
spool_root: &Path,
transport: &dyn DrainTransport,
config: &DrainConfig,
last_upload: &mut Option<Instant>,
ledger: &mut RetryLedger,
) -> RoundOutcome {
let spool = match Spool::open(spool_root) {
Ok(s) => s,
Err(e) => {
tracing::warn!(target: "car::feedback", error = %e, "feedback spool unavailable");
return RoundOutcome::Backoff {
min_delay: config.initial_backoff,
};
}
};
let entries = match spool.list() {
Ok(rows) => rows,
Err(e) => {
tracing::warn!(target: "car::feedback", error = %e, "feedback spool list failed");
return RoundOutcome::Backoff {
min_delay: config.initial_backoff,
};
}
};
let mut queued: Vec<_> = Vec::new();
for entry in entries {
match &entry.state {
SpoolState::Queued => queued.push(entry),
SpoolState::Sending => {
if let Err(e) = spool.mark_queued(&entry.id) {
tracing::warn!(
target: "car::feedback",
entry = %entry.id, error = %e,
"stale Sending entry could not be recovered"
);
} else {
queued.push(entry);
}
}
SpoolState::Acknowledged { .. }
| SpoolState::TerminalActionable { .. }
| SpoolState::TerminalRejected { .. } => {}
}
}
if queued.is_empty() {
return RoundOutcome::Idle;
}
ledger.retain_queued(&queued);
let (queued, capped): (Vec<_>, Vec<_>) = queued
.into_iter()
.partition(|entry| !ledger.is_exhausted(&entry.id, config.max_retriable_attempts));
for entry in &capped {
tracing::debug!(
target: "car::feedback",
entry = %entry.id,
"feedback entry past this daemon's retriable-attempt cap; skipped this round"
);
}
if queued.is_empty() {
return RoundOutcome::AllHeld;
}
let mut eligibility: HashMap<&'static str, DrainEligibility> = HashMap::new();
let mut progressed = false;
let mut retriable: Option<Duration> = None;
for entry in queued {
if matches!(entry.lane, IdentityLane::Anonymous) {
tracing::debug!(
target: "car::feedback",
entry = %entry.id,
"anonymous feedback entry held locally; v1 has no anonymous drain"
);
continue;
}
let cache_key = "authenticated";
let verdict = match eligibility.get(&cache_key) {
Some(v) => v.clone(),
None => {
let v = transport.drain_eligible(&entry.lane).await;
eligibility.insert(cache_key, v.clone());
v
}
};
if let DrainEligibility::Hold { reason } = verdict {
tracing::debug!(
target: "car::feedback",
entry = %entry.id, reason = %reason,
"feedback entry held queued"
);
continue;
}
if let Some(last) = *last_upload {
let since = last.elapsed();
if since < config.min_upload_interval {
tokio::time::sleep(config.min_upload_interval - since).await;
}
}
if let Err(e) = spool.mark_sending(&entry.id) {
tracing::warn!(
target: "car::feedback",
entry = %entry.id, error = %e,
"mark_sending failed; skipping entry this round"
);
continue;
}
let bundle = match spool.load_bundle(&entry.id) {
Ok(b) => b,
Err(e) => {
let _ = spool.mark_terminal(
&entry.id,
TerminalReason::Rejected {
message: format!("stored bundle unreadable: {e}"),
},
);
progressed = true;
continue;
}
};
let attempt = transport
.submit(&bundle, &entry.lane, &entry.client_submission_id)
.await;
*last_upload = Some(Instant::now());
match attempt {
Ok(SubmitOutcome { action, omitted }) => {
report_omitted(&entry.id, &omitted);
match action {
DrainAction::Acknowledge { server_id } => {
if apply(&spool, &entry.id, |s| {
s.mark_acknowledged(&entry.id, &server_id)
}) {
progressed = true;
}
}
DrainAction::Requeue { backoff } => {
apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
ledger.record_failure(&entry.id, config.max_retriable_attempts);
let delay = Duration::from_secs(backoff);
retriable = Some(retriable.map_or(delay, |d| d.max(delay)));
}
DrainAction::RequeueAfter { secs } => {
apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
return RoundOutcome::RetryAfter(Duration::from_secs(secs));
}
DrainAction::TerminalActionable { reason } => {
let reason = match reason {
TransportActionableReason::AuthRequired => TerminalReason::AuthRequired,
TransportActionableReason::ReconsentRequired => {
TerminalReason::ReconsentRequired
}
TransportActionableReason::Forbidden => TerminalReason::Forbidden,
};
if apply(&spool, &entry.id, |s| s.mark_terminal(&entry.id, reason)) {
progressed = true;
}
}
DrainAction::TerminalRejected { message } => {
let message = match omitted_suffix(&omitted) {
Some(suffix) => format!("{message}{suffix}"),
None => message,
};
if apply(&spool, &entry.id, |s| {
s.mark_terminal(&entry.id, TerminalReason::Rejected { message })
}) {
progressed = true;
}
}
}
}
Err(FeedbackTransportError::AnonymousNotYetSupported)
| Err(FeedbackTransportError::NoSession) => {
apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
}
Err(FeedbackTransportError::Unauthorized) => {
apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
}
Err(FeedbackTransportError::FetchFailed(e)) => {
tracing::warn!(
target: "car::feedback",
entry = %entry.id, error = %e,
"feedback submit transport failure; will retry"
);
apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
ledger.record_failure(&entry.id, config.max_retriable_attempts);
let delay = config.initial_backoff;
retriable = Some(retriable.map_or(delay, |d| d.max(delay)));
}
}
}
if let Some(min_delay) = retriable {
RoundOutcome::Backoff { min_delay }
} else if progressed {
RoundOutcome::Progressed
} else {
RoundOutcome::AllHeld
}
}
fn omitted_suffix(omitted: &[String]) -> Option<String> {
if omitted.is_empty() {
None
} else {
Some(format!(" (omitted: {})", omitted.join("; ")))
}
}
fn report_omitted(id: &SpoolEntryId, omitted: &[String]) {
if let Some(suffix) = omitted_suffix(omitted) {
tracing::warn!(
target: "car::feedback",
entry = %id,
"feedback upload sent with omissions{suffix}"
);
}
}
fn apply(
spool: &Spool,
id: &SpoolEntryId,
transition: impl FnOnce(&Spool) -> std::io::Result<()>,
) -> bool {
match transition(spool) {
Ok(()) => true,
Err(e) => {
tracing::warn!(
target: "car::feedback",
entry = %id, error = %e,
"spool transition failed"
);
false
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use car_feedback_core::bundle::{collect, CollectInputs, RedactedBundle};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Mutex;
use tempfile::TempDir;
struct MockTransport {
eligible_calls: AtomicUsize,
submit_calls: AtomicUsize,
eligibility: Mutex<DrainEligibility>,
script: Mutex<Vec<Result<DrainAction, FeedbackTransportError>>>,
omitted: Mutex<Vec<String>>,
submit_at: Mutex<Vec<Instant>>,
submitted_ids: Mutex<Vec<String>>,
}
impl MockTransport {
fn new(action: DrainAction) -> Self {
MockTransport {
eligible_calls: AtomicUsize::new(0),
submit_calls: AtomicUsize::new(0),
eligibility: Mutex::new(DrainEligibility::Eligible),
script: Mutex::new(vec![Ok(action)]),
omitted: Mutex::new(Vec::new()),
submit_at: Mutex::new(Vec::new()),
submitted_ids: Mutex::new(Vec::new()),
}
}
fn holding(reason: &str) -> Self {
let t = Self::new(DrainAction::Acknowledge {
server_id: "unused".into(),
});
*t.eligibility.lock().unwrap() = DrainEligibility::Hold {
reason: reason.to_string(),
};
t
}
fn total_calls(&self) -> usize {
self.eligible_calls.load(Ordering::SeqCst) + self.submit_calls.load(Ordering::SeqCst)
}
}
#[async_trait]
impl DrainTransport for MockTransport {
async fn drain_eligible(&self, _lane: &IdentityLane) -> DrainEligibility {
self.eligible_calls.fetch_add(1, Ordering::SeqCst);
self.eligibility.lock().unwrap().clone()
}
async fn submit(
&self,
_bundle: &RedactedBundle,
_lane: &IdentityLane,
client_submission_id: &str,
) -> Result<SubmitOutcome, FeedbackTransportError> {
self.submit_calls.fetch_add(1, Ordering::SeqCst);
self.submit_at.lock().unwrap().push(Instant::now());
self.submitted_ids
.lock()
.unwrap()
.push(client_submission_id.to_string());
let mut script = self.script.lock().unwrap();
let next = if script.len() > 1 {
script.remove(0)
} else {
script[0].clone()
};
next.map(|action| SubmitOutcome {
action,
omitted: self.omitted.lock().unwrap().clone(),
})
}
}
fn bundle() -> RedactedBundle {
let tmp = TempDir::new().unwrap();
collect(CollectInputs {
description: "the command deck window went blank".to_string(),
state_root: Some(tmp.path().to_path_buf()),
..CollectInputs::default()
})
.unwrap()
}
fn auth_lane() -> IdentityLane {
IdentityLane::Authenticated {
org_id: "org_abc".to_string(),
}
}
fn fast_config() -> DrainConfig {
DrainConfig {
min_upload_interval: Duration::from_secs(15),
initial_backoff: Duration::from_secs(1),
max_backoff: Duration::from_secs(8),
held_recheck_interval: Duration::from_secs(60),
park_check_interval: Duration::from_secs(60),
max_retriable_attempts: DrainConfig::default().max_retriable_attempts,
}
}
fn ledger() -> RetryLedger {
RetryLedger::default()
}
fn enqueue(root: &Path, lane: IdentityLane) -> SpoolEntryId {
let spool = Spool::open(root).unwrap();
spool.enqueue(&bundle(), lane, "title").unwrap()
}
fn persisted_state(root: &Path, id: &SpoolEntryId) -> SpoolState {
Spool::open(root)
.unwrap()
.list()
.unwrap()
.into_iter()
.find(|e| &e.id == id)
.expect("entry on disk")
.state
}
#[tokio::test]
async fn lazy_live_transport_construction_failure_holds_without_submitting() {
let transport = LazyLiveTransport::failing("fixture construction failure");
let verdict = transport.drain_eligible(&auth_lane()).await;
assert_eq!(
verdict,
DrainEligibility::Hold {
reason: "feedback transport unavailable: fixture construction failure".to_string()
}
);
}
#[tokio::test]
async fn empty_outbox_produces_no_transport_calls() {
let tmp = TempDir::new().unwrap();
let spool_root = tmp.path().join("feedback-outbox");
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "never".into(),
});
let mut last = None;
let outcome = run_drain_round(
&spool_root,
&transport,
&fast_config(),
&mut last,
&mut ledger(),
)
.await;
assert_eq!(outcome, RoundOutcome::Idle);
assert_eq!(
transport.total_calls(),
0,
"empty outbox must touch nothing"
);
}
#[tokio::test]
async fn acknowledged_entry_persists_the_server_id_on_disk() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "row-7".into(),
});
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Progressed);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-7".to_string()
}
);
let sent = transport.submitted_ids.lock().unwrap().clone();
assert_eq!(sent.len(), 1);
assert!(!sent[0].is_empty());
}
#[tokio::test]
async fn hold_verdict_keeps_entries_queued_then_capability_flip_drains_them() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport =
MockTransport::holding("capability does not advertise authenticated intake");
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::AllHeld);
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 0);
assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
*transport.eligibility.lock().unwrap() = DrainEligibility::Eligible;
*transport.script.lock().unwrap() = vec![Ok(DrainAction::Acknowledge {
server_id: "row-1".into(),
})];
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Progressed);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-1".to_string()
}
);
}
#[tokio::test]
async fn capability_probe_is_cached_per_round_across_orgs() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
enqueue(&root, auth_lane());
enqueue(
&root,
IdentityLane::Authenticated {
org_id: "org_other".to_string(),
},
);
let transport = MockTransport::holding("capability unavailable");
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::AllHeld);
assert_eq!(
transport.eligible_calls.load(Ordering::SeqCst),
1,
"one capability-backed eligibility probe per lane kind per round"
);
}
#[tokio::test]
async fn anonymous_entries_hold_queued_without_a_submit() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, IdentityLane::Anonymous);
let transport = MockTransport::holding("server does not accept anonymous feedback");
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::AllHeld);
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 0);
assert_eq!(
transport.eligible_calls.load(Ordering::SeqCst),
0,
"anonymous v1 entries must not trigger a capability probe"
);
assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
}
#[tokio::test]
async fn corrupt_bundle_settles_terminal_instead_of_retrying_forever() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
std::fs::write(root.join(id.as_str()).join("bundle.json"), b"{corrupt").unwrap();
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "must-not-submit".into(),
});
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Progressed);
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 0);
match persisted_state(&root, &id) {
SpoolState::TerminalRejected { message } => {
assert!(message.contains("stored bundle unreadable"), "{message}");
}
state => panic!("corrupt bundle must settle terminal, got {state:?}"),
}
}
#[tokio::test]
async fn submit_unauthorized_returns_entry_to_queued_hold() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "unused".into(),
});
*transport.script.lock().unwrap() = vec![Err(FeedbackTransportError::Unauthorized)];
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::AllHeld);
assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
}
#[tokio::test]
async fn terminal_actionable_and_rejected_persist_and_never_retry() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let auth_id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::TerminalActionable {
reason: TransportActionableReason::ReconsentRequired,
});
let mut last = None;
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert!(matches!(
persisted_state(&root, &auth_id),
SpoolState::TerminalActionable {
reason: car_feedback_core::spool::ActionableReason::ReconsentRequired
}
));
let before = transport.submit_calls.load(Ordering::SeqCst);
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Idle);
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), before);
last = None;
let rejected_id = enqueue(&root, auth_lane());
*transport.script.lock().unwrap() = vec![Ok(DrainAction::TerminalRejected {
message: "HTTP 400: description invalid".into(),
})];
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(
persisted_state(&root, &rejected_id),
SpoolState::TerminalRejected {
message: "HTTP 400: description invalid".to_string()
}
);
}
#[tokio::test]
async fn retriable_failure_requeues_durably_and_reports_backoff() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Requeue { backoff: 30 });
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(
outcome,
RoundOutcome::Backoff {
min_delay: Duration::from_secs(30)
}
);
assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
}
#[tokio::test(start_paused = true)]
async fn retry_after_stops_the_round_and_is_honored_before_the_next_upload() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let first = enqueue(&root, auth_lane());
let second = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "row".into(),
});
*transport.script.lock().unwrap() = vec![
Ok(DrainAction::RequeueAfter { secs: 40 }),
Ok(DrainAction::Acknowledge {
server_id: "row-a".into(),
}),
Ok(DrainAction::Acknowledge {
server_id: "row-b".into(),
}),
];
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::RetryAfter(Duration::from_secs(40)));
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 1);
assert_eq!(persisted_state(&root, &first), SpoolState::Queued);
assert_eq!(persisted_state(&root, &second), SpoolState::Queued);
tokio::time::sleep(Duration::from_secs(40)).await;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Progressed);
let stamps = transport.submit_at.lock().unwrap().clone();
assert!(stamps.len() >= 2);
assert!(
stamps[1].duration_since(stamps[0]) >= Duration::from_secs(40),
"second upload ran {:?} after the 429 — Retry-After not honored",
stamps[1].duration_since(stamps[0])
);
}
#[tokio::test(start_paused = true)]
async fn uploads_pace_at_most_one_per_min_interval() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
enqueue(&root, auth_lane());
enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "row".into(),
});
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Progressed);
let stamps = transport.submit_at.lock().unwrap().clone();
assert_eq!(stamps.len(), 2);
assert!(
stamps[1].duration_since(stamps[0]) >= Duration::from_secs(15),
"uploads {:?} apart — pacing not applied",
stamps[1].duration_since(stamps[0])
);
}
#[tokio::test]
async fn stale_sending_entry_from_a_crashed_drain_recovers_and_resends_same_id() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let spool = Spool::open(&root).unwrap();
spool.mark_sending(&id).unwrap();
let original_csid = spool
.list()
.unwrap()
.into_iter()
.find(|e| e.id == id)
.unwrap()
.client_submission_id;
drop(spool);
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "row-1".into(),
});
let mut last = None;
let outcome =
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
assert_eq!(outcome, RoundOutcome::Progressed);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-1".to_string()
}
);
assert_eq!(
transport.submitted_ids.lock().unwrap().as_slice(),
&[original_csid],
"the recovered entry must re-send its ORIGINAL idempotency key"
);
}
#[tokio::test(start_paused = true)]
async fn spawned_drain_parks_idle_and_drains_on_wake() {
let tmp = TempDir::new().unwrap();
let car_home = tmp.path().to_path_buf();
let root = car_home.join(FEEDBACK_OUTBOX_DIR);
let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
server_id: "row-1".into(),
}));
let handle = spawn_feedback_drain_with(car_home, transport.clone(), fast_config());
tokio::time::sleep(Duration::from_secs(3600)).await;
assert_eq!(
transport.total_calls(),
0,
"an empty-outbox park (incl. its local disk ticks) must never touch the transport"
);
let id = enqueue(&root, auth_lane());
handle.wake();
for _ in 0..200 {
tokio::task::yield_now().await;
if transport.submit_calls.load(Ordering::SeqCst) > 0 {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 1);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-1".to_string()
}
);
}
#[tokio::test(start_paused = true)]
async fn cli_direct_enqueue_while_parked_drains_within_one_tick() {
let tmp = TempDir::new().unwrap();
let car_home = tmp.path().to_path_buf();
let root = car_home.join(FEEDBACK_OUTBOX_DIR);
let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
server_id: "row-cli".into(),
}));
let _handle = spawn_feedback_drain_with(car_home, transport.clone(), fast_config());
tokio::time::sleep(Duration::from_secs(1)).await;
assert_eq!(transport.total_calls(), 0);
let id = enqueue(&root, auth_lane());
for _ in 0..200 {
tokio::task::yield_now().await;
if transport.submit_calls.load(Ordering::SeqCst) > 0 {
break;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
assert_eq!(
transport.submit_calls.load(Ordering::SeqCst),
1,
"a parked drain must catch a direct spool enqueue via the local tick"
);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-cli".to_string()
}
);
}
#[test]
fn park_probe_sees_only_pending_entries() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join(FEEDBACK_OUTBOX_DIR);
assert!(!spool_has_pending_entries(&root));
std::fs::create_dir_all(root.join(".tmp-half-written")).unwrap();
std::fs::write(root.join("stray-file"), b"x").unwrap();
assert!(
!spool_has_pending_entries(&root),
"staging dirs and files don't count"
);
std::fs::create_dir_all(root.join("00000000000000000000-not-an-entry")).unwrap();
assert!(
!spool_has_pending_entries(&root),
"a stray directory is not a pending entry"
);
let settled = enqueue(&root, auth_lane());
{
let spool = Spool::open(&root).unwrap();
spool.mark_sending(&settled).unwrap();
spool.mark_acknowledged(&settled, "row-1").unwrap();
}
assert_eq!(
persisted_state(&root, &settled),
SpoolState::Acknowledged {
server_id: "row-1".to_string()
}
);
assert!(
!spool_has_pending_entries(&root),
"a settled entry must not break the park"
);
let queued = enqueue(&root, auth_lane());
assert!(spool_has_pending_entries(&root));
Spool::open(&root).unwrap().mark_sending(&queued).unwrap();
assert!(spool_has_pending_entries(&root));
}
#[tokio::test(start_paused = true)]
async fn absurd_retry_after_is_clamped_to_the_backoff_ceiling() {
let tmp = TempDir::new().unwrap();
let car_home = tmp.path().to_path_buf();
let root = car_home.join(FEEDBACK_OUTBOX_DIR);
let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
server_id: "row-1".into(),
}));
*transport.script.lock().unwrap() = vec![
Ok(DrainAction::RequeueAfter { secs: 999_999_999 }),
Ok(DrainAction::Acknowledge {
server_id: "row-1".into(),
}),
];
let handle = spawn_feedback_drain_with(car_home, transport.clone(), fast_config());
tokio::time::sleep(Duration::from_secs(1)).await;
let id = enqueue(&root, auth_lane());
handle.wake();
for _ in 0..200 {
tokio::task::yield_now().await;
if transport.submit_calls.load(Ordering::SeqCst) >= 1 {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 1);
assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
tokio::time::sleep(Duration::from_secs(60)).await;
assert_eq!(
transport.submit_calls.load(Ordering::SeqCst),
2,
"the drain must retry within the backoff ceiling, not the header's decades"
);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-1".to_string()
}
);
}
#[tokio::test(start_paused = true)]
async fn submit_wake_ends_a_retry_after_wait_early() {
let tmp = TempDir::new().unwrap();
let car_home = tmp.path().to_path_buf();
let root = car_home.join(FEEDBACK_OUTBOX_DIR);
let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
server_id: "row".into(),
}));
*transport.script.lock().unwrap() = vec![
Ok(DrainAction::RequeueAfter { secs: 300 }),
Ok(DrainAction::Acknowledge {
server_id: "row-a".into(),
}),
Ok(DrainAction::Acknowledge {
server_id: "row-b".into(),
}),
];
let mut config = fast_config();
config.max_backoff = Duration::from_secs(600);
let handle = spawn_feedback_drain_with(car_home, transport.clone(), config);
tokio::time::sleep(Duration::from_secs(1)).await;
let first = enqueue(&root, auth_lane());
handle.wake();
for _ in 0..200 {
tokio::task::yield_now().await;
if transport.submit_calls.load(Ordering::SeqCst) >= 1 {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert_eq!(
transport.submit_calls.load(Ordering::SeqCst),
1,
"the 429 landed"
);
assert_eq!(persisted_state(&root, &first), SpoolState::Queued);
let second = enqueue(&root, auth_lane());
handle.wake();
tokio::time::sleep(Duration::from_secs(30)).await;
assert!(
transport.submit_calls.load(Ordering::SeqCst) >= 2,
"a submit wake must end the Retry-After wait (calls: {})",
transport.submit_calls.load(Ordering::SeqCst)
);
assert_eq!(
persisted_state(&root, &first),
SpoolState::Acknowledged {
server_id: "row-a".to_string()
},
"the oldest queued entry is retried first"
);
tokio::time::sleep(Duration::from_secs(30)).await;
assert_eq!(
persisted_state(&root, &second),
SpoolState::Acknowledged {
server_id: "row-b".to_string()
}
);
}
#[tokio::test]
async fn rejected_upload_with_omissions_persists_the_omitted_note() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::TerminalRejected {
message: "HTTP 400: description invalid".into(),
});
*transport.omitted.lock().unwrap() = vec!["screenshot dropped (413 fallback)".to_string()];
let mut last = None;
run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
match persisted_state(&root, &id) {
SpoolState::TerminalRejected { message } => {
assert!(
message.contains("screenshot dropped (413 fallback)"),
"the omitted note must persist in the durable rejection message: {message}"
);
assert!(message.contains("HTTP 400"));
}
other => panic!("expected TerminalRejected, got {other:?}"),
}
}
#[tokio::test(start_paused = true)]
async fn retriable_failures_stop_after_the_per_process_cap_and_the_entry_stays_queued() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Requeue { backoff: 1 });
*transport.script.lock().unwrap() = vec![
Ok(DrainAction::Requeue { backoff: 1 }),
Err(FeedbackTransportError::FetchFailed(
"connection reset".into(),
)),
Ok(DrainAction::Requeue { backoff: 1 }),
];
let mut config = fast_config();
config.max_retriable_attempts = 3;
let mut last = None;
let mut ledger = RetryLedger::default();
for round in 1..=3 {
let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
assert!(
matches!(outcome, RoundOutcome::Backoff { .. }),
"round {round}: {outcome:?}"
);
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), round);
assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
}
let probes_before = transport.eligible_calls.load(Ordering::SeqCst);
for _ in 0..3 {
let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
assert_eq!(outcome, RoundOutcome::AllHeld);
}
assert_eq!(
transport.submit_calls.load(Ordering::SeqCst),
3,
"no submit past the cap"
);
assert_eq!(
transport.eligible_calls.load(Ordering::SeqCst),
probes_before,
"a capped entry costs no network — not even the eligibility probe"
);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Queued,
"capped ⇒ still Queued on disk: never a terminal state, never pruned"
);
let mut restarted = RetryLedger::default();
let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut restarted).await;
assert!(matches!(outcome, RoundOutcome::Backoff { .. }));
assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 4);
}
#[tokio::test(start_paused = true)]
async fn retry_after_does_not_consume_the_attempt_budget() {
let tmp = TempDir::new().unwrap();
let root = tmp.path().join("feedback-outbox");
let id = enqueue(&root, auth_lane());
let transport = MockTransport::new(DrainAction::Acknowledge {
server_id: "row".into(),
});
*transport.script.lock().unwrap() = vec![
Ok(DrainAction::RequeueAfter { secs: 1 }),
Ok(DrainAction::RequeueAfter { secs: 1 }),
Ok(DrainAction::Acknowledge {
server_id: "row-1".into(),
}),
];
let mut config = fast_config();
config.max_retriable_attempts = 1;
let mut last = None;
let mut ledger = RetryLedger::default();
for _ in 0..2 {
let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
assert_eq!(outcome, RoundOutcome::RetryAfter(Duration::from_secs(1)));
tokio::time::sleep(Duration::from_secs(1)).await;
}
let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
assert_eq!(outcome, RoundOutcome::Progressed);
assert_eq!(
persisted_state(&root, &id),
SpoolState::Acknowledged {
server_id: "row-1".to_string()
}
);
}
}