1use std::collections::HashMap;
93use std::path::{Path, PathBuf};
94use std::sync::{Arc, OnceLock};
95use std::time::Duration;
96
97use async_trait::async_trait;
98use tokio::sync::Notify;
99use tokio::time::Instant;
100
101use car_feedback_core::spool::{
102 IdentityLane, Spool, SpoolEntryId, SpoolEntrySummary, SpoolState, TerminalReason,
103};
104use car_parslee::feedback_transport::{
105 DrainAction, DrainEligibility, FeedbackTransport, FeedbackTransportError, SubmitOutcome,
106 TransportActionableReason,
107};
108
109use crate::feedback::FEEDBACK_OUTBOX_DIR;
110use crate::session::ServerState;
111
112#[derive(Debug, Clone)]
116pub struct DrainConfig {
117 pub min_upload_interval: Duration,
119 pub initial_backoff: Duration,
121 pub max_backoff: Duration,
123 pub held_recheck_interval: Duration,
127 pub park_check_interval: Duration,
134 pub max_retriable_attempts: u32,
141}
142
143impl Default for DrainConfig {
144 fn default() -> Self {
145 DrainConfig {
146 min_upload_interval: Duration::from_secs(15),
147 initial_backoff: Duration::from_secs(30),
148 max_backoff: Duration::from_secs(15 * 60),
149 held_recheck_interval: Duration::from_secs(15 * 60),
150 park_check_interval: Duration::from_secs(60),
151 max_retriable_attempts: 24,
152 }
153 }
154}
155
156#[derive(Debug, Default)]
163pub struct RetryLedger {
164 failures: HashMap<SpoolEntryId, u32>,
165}
166
167impl RetryLedger {
168 fn record_failure(&mut self, id: &SpoolEntryId, cap: u32) {
171 let count = self.failures.entry(id.clone()).or_insert(0);
172 *count = count.saturating_add(1);
173 if *count == cap {
174 tracing::warn!(
175 target: "car::feedback",
176 entry = %id, attempts = cap,
177 "feedback entry hit this daemon's retriable-attempt cap; it stays queued \
178 (export it with `car feedback --export`) and retries again after the \
179 daemon restarts"
180 );
181 }
182 }
183
184 fn is_exhausted(&self, id: &SpoolEntryId, cap: u32) -> bool {
186 self.failures.get(id).is_some_and(|count| *count >= cap)
187 }
188
189 fn retain_queued(&mut self, queued: &[SpoolEntrySummary]) {
192 self.failures
193 .retain(|id, _| queued.iter().any(|entry| &entry.id == id));
194 }
195}
196
197#[async_trait]
200pub trait DrainTransport: Send + Sync {
201 async fn drain_eligible(&self, lane: &IdentityLane) -> DrainEligibility;
202 async fn submit(
203 &self,
204 bundle: &car_feedback_core::bundle::RedactedBundle,
205 lane: &IdentityLane,
206 client_submission_id: &str,
207 ) -> Result<SubmitOutcome, FeedbackTransportError>;
208}
209
210#[async_trait]
211impl DrainTransport for FeedbackTransport {
212 async fn drain_eligible(&self, lane: &IdentityLane) -> DrainEligibility {
213 FeedbackTransport::drain_eligible(self, lane).await
214 }
215 async fn submit(
216 &self,
217 bundle: &car_feedback_core::bundle::RedactedBundle,
218 lane: &IdentityLane,
219 client_submission_id: &str,
220 ) -> Result<SubmitOutcome, FeedbackTransportError> {
221 FeedbackTransport::submit_report(self, bundle, lane, client_submission_id).await
222 }
223}
224
225struct LazyLiveTransport {
235 inner: tokio::sync::OnceCell<FeedbackTransport>,
236 #[cfg(test)]
237 constructor_error: Option<String>,
238}
239
240impl LazyLiveTransport {
241 fn new() -> Self {
242 Self {
243 inner: tokio::sync::OnceCell::new(),
244 #[cfg(test)]
245 constructor_error: None,
246 }
247 }
248
249 #[cfg(test)]
250 fn failing(error: &str) -> Self {
251 Self {
252 inner: tokio::sync::OnceCell::new(),
253 constructor_error: Some(error.to_string()),
254 }
255 }
256
257 async fn get(&self) -> Result<&FeedbackTransport, String> {
258 #[cfg(test)]
259 if let Some(error) = &self.constructor_error {
260 return Err(error.clone());
261 }
262 self.inner
263 .get_or_try_init(|| async { FeedbackTransport::live() })
264 .await
265 }
266}
267
268#[async_trait]
269impl DrainTransport for LazyLiveTransport {
270 async fn drain_eligible(&self, lane: &IdentityLane) -> DrainEligibility {
271 match self.get().await {
272 Ok(transport) => DrainTransport::drain_eligible(transport, lane).await,
273 Err(error) => DrainEligibility::Hold {
274 reason: format!("feedback transport unavailable: {error}"),
275 },
276 }
277 }
278 async fn submit(
279 &self,
280 bundle: &car_feedback_core::bundle::RedactedBundle,
281 lane: &IdentityLane,
282 client_submission_id: &str,
283 ) -> Result<SubmitOutcome, FeedbackTransportError> {
284 match self.get().await {
285 Ok(transport) => {
286 DrainTransport::submit(transport, bundle, lane, client_submission_id).await
287 }
288 Err(error) => Err(FeedbackTransportError::FetchFailed(format!(
292 "feedback transport unavailable: {error}"
293 ))),
294 }
295 }
296}
297
298#[derive(Debug, Clone, PartialEq, Eq)]
300pub enum RoundOutcome {
301 Idle,
303 AllHeld,
309 Progressed,
312 Backoff { min_delay: Duration },
314 RetryAfter(Duration),
316}
317
318#[derive(Clone)]
320pub struct FeedbackDrainHandle {
321 notify: Arc<Notify>,
322}
323
324impl FeedbackDrainHandle {
325 pub fn wake(&self) {
326 self.notify.notify_one();
327 }
328}
329
330static DRAIN_HANDLE: OnceLock<FeedbackDrainHandle> = OnceLock::new();
333
334pub fn wake_feedback_drain() {
338 if let Some(handle) = DRAIN_HANDLE.get() {
339 handle.wake();
340 }
341}
342
343pub fn spawn_feedback_drain(state: &ServerState) {
349 let car_home = match crate::feedback::car_home_dir(state) {
350 Ok(home) => home,
351 Err(e) => {
352 tracing::warn!(target: "car::feedback", error = %e, "feedback drain not started");
353 return;
354 }
355 };
356 let transport: Arc<dyn DrainTransport> = Arc::new(LazyLiveTransport::new());
361 let handle = spawn_feedback_drain_with(car_home, transport, DrainConfig::default());
362 let _ = DRAIN_HANDLE.set(handle);
364}
365
366pub fn spawn_feedback_drain_with(
369 car_home: PathBuf,
370 transport: Arc<dyn DrainTransport>,
371 config: DrainConfig,
372) -> FeedbackDrainHandle {
373 let notify = Arc::new(Notify::new());
374 let handle = FeedbackDrainHandle {
375 notify: notify.clone(),
376 };
377 let spool_root = car_home.join(FEEDBACK_OUTBOX_DIR);
378 tokio::spawn(async move {
379 let mut last_upload: Option<Instant> = None;
380 let mut ledger = RetryLedger::default();
383 loop {
384 let mut backoff = config.initial_backoff;
386 loop {
387 let outcome = run_drain_round(
388 &spool_root,
389 transport.as_ref(),
390 &config,
391 &mut last_upload,
392 &mut ledger,
393 )
394 .await;
395 match outcome {
396 RoundOutcome::Idle => break,
397 RoundOutcome::Progressed => {
398 backoff = config.initial_backoff;
399 }
400 RoundOutcome::AllHeld => {
401 wait_or_wake(¬ify, config.held_recheck_interval).await;
404 }
405 RoundOutcome::Backoff { min_delay } => {
406 let delay = backoff.max(min_delay).min(config.max_backoff);
407 wait_or_wake(¬ify, delay).await;
408 backoff = (backoff * 2).min(config.max_backoff);
409 }
410 RoundOutcome::RetryAfter(delay) => {
411 let delay = delay.min(config.max_backoff);
420 wait_or_wake(¬ify, delay).await;
421 }
422 }
423 }
424 loop {
430 tokio::select! {
431 _ = notify.notified() => break,
432 _ = tokio::time::sleep(config.park_check_interval) => {
433 if spool_has_pending_entries(&spool_root) {
434 break;
435 }
436 }
437 }
438 }
439 }
440 });
441 handle
442}
443
444fn spool_has_pending_entries(spool_root: &Path) -> bool {
457 let Ok(dirents) = std::fs::read_dir(spool_root) else {
458 return false;
460 };
461 let has_published_dir = dirents.flatten().any(|dirent| {
462 !dirent.file_name().to_string_lossy().starts_with(".tmp-")
463 && dirent.file_type().map(|t| t.is_dir()).unwrap_or(false)
464 });
465 if !has_published_dir {
466 return false;
467 }
468 let Ok(spool) = Spool::open(spool_root) else {
469 return false;
470 };
471 spool
472 .list()
473 .map(|rows| {
474 rows.iter()
475 .any(|row| matches!(row.state, SpoolState::Queued | SpoolState::Sending))
476 })
477 .unwrap_or(false)
478}
479
480async fn wait_or_wake(notify: &Notify, delay: Duration) {
481 tokio::select! {
482 _ = tokio::time::sleep(delay) => {}
483 _ = notify.notified() => {}
484 }
485}
486
487pub async fn run_drain_round(
493 spool_root: &Path,
494 transport: &dyn DrainTransport,
495 config: &DrainConfig,
496 last_upload: &mut Option<Instant>,
497 ledger: &mut RetryLedger,
498) -> RoundOutcome {
499 let spool = match Spool::open(spool_root) {
502 Ok(s) => s,
503 Err(e) => {
504 tracing::warn!(target: "car::feedback", error = %e, "feedback spool unavailable");
505 return RoundOutcome::Backoff {
506 min_delay: config.initial_backoff,
507 };
508 }
509 };
510 let entries = match spool.list() {
511 Ok(rows) => rows,
512 Err(e) => {
513 tracing::warn!(target: "car::feedback", error = %e, "feedback spool list failed");
514 return RoundOutcome::Backoff {
515 min_delay: config.initial_backoff,
516 };
517 }
518 };
519
520 let mut queued: Vec<_> = Vec::new();
524 for entry in entries {
525 match &entry.state {
526 SpoolState::Queued => queued.push(entry),
527 SpoolState::Sending => {
528 if let Err(e) = spool.mark_queued(&entry.id) {
529 tracing::warn!(
530 target: "car::feedback",
531 entry = %entry.id, error = %e,
532 "stale Sending entry could not be recovered"
533 );
534 } else {
535 queued.push(entry);
536 }
537 }
538 SpoolState::Acknowledged { .. }
539 | SpoolState::TerminalActionable { .. }
540 | SpoolState::TerminalRejected { .. } => {}
541 }
542 }
543
544 if queued.is_empty() {
545 return RoundOutcome::Idle;
547 }
548
549 ledger.retain_queued(&queued);
556 let (queued, capped): (Vec<_>, Vec<_>) = queued
557 .into_iter()
558 .partition(|entry| !ledger.is_exhausted(&entry.id, config.max_retriable_attempts));
559 for entry in &capped {
560 tracing::debug!(
561 target: "car::feedback",
562 entry = %entry.id,
563 "feedback entry past this daemon's retriable-attempt cap; skipped this round"
564 );
565 }
566 if queued.is_empty() {
567 return RoundOutcome::AllHeld;
570 }
571
572 let mut eligibility: HashMap<&'static str, DrainEligibility> = HashMap::new();
579 let mut progressed = false;
580 let mut retriable: Option<Duration> = None;
581
582 for entry in queued {
583 if matches!(entry.lane, IdentityLane::Anonymous) {
587 tracing::debug!(
588 target: "car::feedback",
589 entry = %entry.id,
590 "anonymous feedback entry held locally; v1 has no anonymous drain"
591 );
592 continue;
593 }
594 let cache_key = "authenticated";
595 let verdict = match eligibility.get(&cache_key) {
596 Some(v) => v.clone(),
597 None => {
598 let v = transport.drain_eligible(&entry.lane).await;
599 eligibility.insert(cache_key, v.clone());
600 v
601 }
602 };
603 if let DrainEligibility::Hold { reason } = verdict {
604 tracing::debug!(
605 target: "car::feedback",
606 entry = %entry.id, reason = %reason,
607 "feedback entry held queued"
608 );
609 continue;
610 }
611
612 if let Some(last) = *last_upload {
614 let since = last.elapsed();
615 if since < config.min_upload_interval {
616 tokio::time::sleep(config.min_upload_interval - since).await;
617 }
618 }
619
620 if let Err(e) = spool.mark_sending(&entry.id) {
621 tracing::warn!(
622 target: "car::feedback",
623 entry = %entry.id, error = %e,
624 "mark_sending failed; skipping entry this round"
625 );
626 continue;
627 }
628
629 let bundle = match spool.load_bundle(&entry.id) {
630 Ok(b) => b,
631 Err(e) => {
632 let _ = spool.mark_terminal(
635 &entry.id,
636 TerminalReason::Rejected {
637 message: format!("stored bundle unreadable: {e}"),
638 },
639 );
640 progressed = true;
641 continue;
642 }
643 };
644
645 let attempt = transport
646 .submit(&bundle, &entry.lane, &entry.client_submission_id)
647 .await;
648 *last_upload = Some(Instant::now());
649
650 match attempt {
651 Ok(SubmitOutcome { action, omitted }) => {
652 report_omitted(&entry.id, &omitted);
653 match action {
654 DrainAction::Acknowledge { server_id } => {
655 if apply(&spool, &entry.id, |s| {
656 s.mark_acknowledged(&entry.id, &server_id)
657 }) {
658 progressed = true;
659 }
660 }
661 DrainAction::Requeue { backoff } => {
662 apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
663 ledger.record_failure(&entry.id, config.max_retriable_attempts);
664 let delay = Duration::from_secs(backoff);
665 retriable = Some(retriable.map_or(delay, |d| d.max(delay)));
666 }
667 DrainAction::RequeueAfter { secs } => {
668 apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
669 return RoundOutcome::RetryAfter(Duration::from_secs(secs));
673 }
674 DrainAction::TerminalActionable { reason } => {
675 let reason = match reason {
676 TransportActionableReason::AuthRequired => TerminalReason::AuthRequired,
677 TransportActionableReason::ReconsentRequired => {
678 TerminalReason::ReconsentRequired
679 }
680 TransportActionableReason::Forbidden => TerminalReason::Forbidden,
681 };
682 if apply(&spool, &entry.id, |s| s.mark_terminal(&entry.id, reason)) {
683 progressed = true;
684 }
685 }
686 DrainAction::TerminalRejected { message } => {
687 let message = match omitted_suffix(&omitted) {
691 Some(suffix) => format!("{message}{suffix}"),
692 None => message,
693 };
694 if apply(&spool, &entry.id, |s| {
695 s.mark_terminal(&entry.id, TerminalReason::Rejected { message })
696 }) {
697 progressed = true;
698 }
699 }
700 }
701 }
702 Err(FeedbackTransportError::AnonymousNotYetSupported)
705 | Err(FeedbackTransportError::NoSession) => {
706 apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
707 }
708 Err(FeedbackTransportError::Unauthorized) => {
713 apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
714 }
715 Err(FeedbackTransportError::FetchFailed(e)) => {
716 tracing::warn!(
717 target: "car::feedback",
718 entry = %entry.id, error = %e,
719 "feedback submit transport failure; will retry"
720 );
721 apply(&spool, &entry.id, |s| s.mark_queued(&entry.id));
722 ledger.record_failure(&entry.id, config.max_retriable_attempts);
723 let delay = config.initial_backoff;
724 retriable = Some(retriable.map_or(delay, |d| d.max(delay)));
725 }
726 }
727 }
728
729 if let Some(min_delay) = retriable {
730 RoundOutcome::Backoff { min_delay }
731 } else if progressed {
732 RoundOutcome::Progressed
733 } else {
734 RoundOutcome::AllHeld
735 }
736}
737
738fn omitted_suffix(omitted: &[String]) -> Option<String> {
741 if omitted.is_empty() {
742 None
743 } else {
744 Some(format!(" (omitted: {})", omitted.join("; ")))
745 }
746}
747
748fn report_omitted(id: &SpoolEntryId, omitted: &[String]) {
755 if let Some(suffix) = omitted_suffix(omitted) {
756 tracing::warn!(
757 target: "car::feedback",
758 entry = %id,
759 "feedback upload sent with omissions{suffix}"
760 );
761 }
762}
763
764fn apply(
767 spool: &Spool,
768 id: &SpoolEntryId,
769 transition: impl FnOnce(&Spool) -> std::io::Result<()>,
770) -> bool {
771 match transition(spool) {
772 Ok(()) => true,
773 Err(e) => {
774 tracing::warn!(
775 target: "car::feedback",
776 entry = %id, error = %e,
777 "spool transition failed"
778 );
779 false
780 }
781 }
782}
783
784#[cfg(test)]
785mod tests {
786 use super::*;
787 use car_feedback_core::bundle::{collect, CollectInputs, RedactedBundle};
788 use std::sync::atomic::{AtomicUsize, Ordering};
789 use std::sync::Mutex;
790 use tempfile::TempDir;
791
792 struct MockTransport {
795 eligible_calls: AtomicUsize,
796 submit_calls: AtomicUsize,
797 eligibility: Mutex<DrainEligibility>,
798 script: Mutex<Vec<Result<DrainAction, FeedbackTransportError>>>,
800 omitted: Mutex<Vec<String>>,
803 submit_at: Mutex<Vec<Instant>>,
805 submitted_ids: Mutex<Vec<String>>,
806 }
807
808 impl MockTransport {
809 fn new(action: DrainAction) -> Self {
810 MockTransport {
811 eligible_calls: AtomicUsize::new(0),
812 submit_calls: AtomicUsize::new(0),
813 eligibility: Mutex::new(DrainEligibility::Eligible),
814 script: Mutex::new(vec![Ok(action)]),
815 omitted: Mutex::new(Vec::new()),
816 submit_at: Mutex::new(Vec::new()),
817 submitted_ids: Mutex::new(Vec::new()),
818 }
819 }
820
821 fn holding(reason: &str) -> Self {
822 let t = Self::new(DrainAction::Acknowledge {
823 server_id: "unused".into(),
824 });
825 *t.eligibility.lock().unwrap() = DrainEligibility::Hold {
826 reason: reason.to_string(),
827 };
828 t
829 }
830
831 fn total_calls(&self) -> usize {
832 self.eligible_calls.load(Ordering::SeqCst) + self.submit_calls.load(Ordering::SeqCst)
833 }
834 }
835
836 #[async_trait]
837 impl DrainTransport for MockTransport {
838 async fn drain_eligible(&self, _lane: &IdentityLane) -> DrainEligibility {
839 self.eligible_calls.fetch_add(1, Ordering::SeqCst);
840 self.eligibility.lock().unwrap().clone()
841 }
842 async fn submit(
843 &self,
844 _bundle: &RedactedBundle,
845 _lane: &IdentityLane,
846 client_submission_id: &str,
847 ) -> Result<SubmitOutcome, FeedbackTransportError> {
848 self.submit_calls.fetch_add(1, Ordering::SeqCst);
849 self.submit_at.lock().unwrap().push(Instant::now());
850 self.submitted_ids
851 .lock()
852 .unwrap()
853 .push(client_submission_id.to_string());
854 let mut script = self.script.lock().unwrap();
855 let next = if script.len() > 1 {
856 script.remove(0)
857 } else {
858 script[0].clone()
859 };
860 next.map(|action| SubmitOutcome {
861 action,
862 omitted: self.omitted.lock().unwrap().clone(),
863 })
864 }
865 }
866
867 fn bundle() -> RedactedBundle {
868 let tmp = TempDir::new().unwrap();
869 collect(CollectInputs {
870 description: "the command deck window went blank".to_string(),
871 state_root: Some(tmp.path().to_path_buf()),
872 ..CollectInputs::default()
873 })
874 .unwrap()
875 }
876
877 fn auth_lane() -> IdentityLane {
878 IdentityLane::Authenticated {
879 org_id: "org_abc".to_string(),
880 }
881 }
882
883 fn fast_config() -> DrainConfig {
884 DrainConfig {
885 min_upload_interval: Duration::from_secs(15),
886 initial_backoff: Duration::from_secs(1),
887 max_backoff: Duration::from_secs(8),
888 held_recheck_interval: Duration::from_secs(60),
889 park_check_interval: Duration::from_secs(60),
890 max_retriable_attempts: DrainConfig::default().max_retriable_attempts,
891 }
892 }
893
894 fn ledger() -> RetryLedger {
897 RetryLedger::default()
898 }
899
900 fn enqueue(root: &Path, lane: IdentityLane) -> SpoolEntryId {
901 let spool = Spool::open(root).unwrap();
902 spool.enqueue(&bundle(), lane, "title").unwrap()
903 }
904
905 fn persisted_state(root: &Path, id: &SpoolEntryId) -> SpoolState {
908 Spool::open(root)
909 .unwrap()
910 .list()
911 .unwrap()
912 .into_iter()
913 .find(|e| &e.id == id)
914 .expect("entry on disk")
915 .state
916 }
917
918 #[tokio::test]
921 async fn lazy_live_transport_construction_failure_holds_without_submitting() {
922 let transport = LazyLiveTransport::failing("fixture construction failure");
923 let verdict = transport.drain_eligible(&auth_lane()).await;
924 assert_eq!(
925 verdict,
926 DrainEligibility::Hold {
927 reason: "feedback transport unavailable: fixture construction failure".to_string()
928 }
929 );
930 }
931
932 #[tokio::test]
933 async fn empty_outbox_produces_no_transport_calls() {
934 let tmp = TempDir::new().unwrap();
939 let spool_root = tmp.path().join("feedback-outbox");
940 let transport = MockTransport::new(DrainAction::Acknowledge {
941 server_id: "never".into(),
942 });
943 let mut last = None;
944 let outcome = run_drain_round(
945 &spool_root,
946 &transport,
947 &fast_config(),
948 &mut last,
949 &mut ledger(),
950 )
951 .await;
952 assert_eq!(outcome, RoundOutcome::Idle);
953 assert_eq!(
954 transport.total_calls(),
955 0,
956 "empty outbox must touch nothing"
957 );
958 }
959
960 #[tokio::test]
963 async fn acknowledged_entry_persists_the_server_id_on_disk() {
964 let tmp = TempDir::new().unwrap();
965 let root = tmp.path().join("feedback-outbox");
966 let id = enqueue(&root, auth_lane());
967 let transport = MockTransport::new(DrainAction::Acknowledge {
968 server_id: "row-7".into(),
969 });
970 let mut last = None;
971 let outcome =
972 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
973 assert_eq!(outcome, RoundOutcome::Progressed);
974 assert_eq!(
975 persisted_state(&root, &id),
976 SpoolState::Acknowledged {
977 server_id: "row-7".to_string()
978 }
979 );
980 let sent = transport.submitted_ids.lock().unwrap().clone();
982 assert_eq!(sent.len(), 1);
983 assert!(!sent[0].is_empty());
984 }
985
986 #[tokio::test]
989 async fn hold_verdict_keeps_entries_queued_then_capability_flip_drains_them() {
990 let tmp = TempDir::new().unwrap();
995 let root = tmp.path().join("feedback-outbox");
996 let id = enqueue(&root, auth_lane());
997 let transport =
998 MockTransport::holding("capability does not advertise authenticated intake");
999 let mut last = None;
1000
1001 let outcome =
1002 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1003 assert_eq!(outcome, RoundOutcome::AllHeld);
1004 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 0);
1005 assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
1006
1007 *transport.eligibility.lock().unwrap() = DrainEligibility::Eligible;
1008 *transport.script.lock().unwrap() = vec![Ok(DrainAction::Acknowledge {
1009 server_id: "row-1".into(),
1010 })];
1011 let outcome =
1012 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1013 assert_eq!(outcome, RoundOutcome::Progressed);
1014 assert_eq!(
1015 persisted_state(&root, &id),
1016 SpoolState::Acknowledged {
1017 server_id: "row-1".to_string()
1018 }
1019 );
1020 }
1021
1022 #[tokio::test]
1023 async fn capability_probe_is_cached_per_round_across_orgs() {
1024 let tmp = TempDir::new().unwrap();
1029 let root = tmp.path().join("feedback-outbox");
1030 enqueue(&root, auth_lane());
1031 enqueue(
1032 &root,
1033 IdentityLane::Authenticated {
1034 org_id: "org_other".to_string(),
1035 },
1036 );
1037 let transport = MockTransport::holding("capability unavailable");
1038 let mut last = None;
1039 let outcome =
1040 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1041 assert_eq!(outcome, RoundOutcome::AllHeld);
1042 assert_eq!(
1043 transport.eligible_calls.load(Ordering::SeqCst),
1044 1,
1045 "one capability-backed eligibility probe per lane kind per round"
1046 );
1047 }
1048
1049 #[tokio::test]
1050 async fn anonymous_entries_hold_queued_without_a_submit() {
1051 let tmp = TempDir::new().unwrap();
1055 let root = tmp.path().join("feedback-outbox");
1056 let id = enqueue(&root, IdentityLane::Anonymous);
1057 let transport = MockTransport::holding("server does not accept anonymous feedback");
1058 let mut last = None;
1059 let outcome =
1060 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1061 assert_eq!(outcome, RoundOutcome::AllHeld);
1062 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 0);
1063 assert_eq!(
1064 transport.eligible_calls.load(Ordering::SeqCst),
1065 0,
1066 "anonymous v1 entries must not trigger a capability probe"
1067 );
1068 assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
1069 }
1070
1071 #[tokio::test]
1074 async fn corrupt_bundle_settles_terminal_instead_of_retrying_forever() {
1075 let tmp = TempDir::new().unwrap();
1076 let root = tmp.path().join("feedback-outbox");
1077 let id = enqueue(&root, auth_lane());
1078 std::fs::write(root.join(id.as_str()).join("bundle.json"), b"{corrupt").unwrap();
1079 let transport = MockTransport::new(DrainAction::Acknowledge {
1080 server_id: "must-not-submit".into(),
1081 });
1082 let mut last = None;
1083 let outcome =
1084 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1085 assert_eq!(outcome, RoundOutcome::Progressed);
1086 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 0);
1087 match persisted_state(&root, &id) {
1088 SpoolState::TerminalRejected { message } => {
1089 assert!(message.contains("stored bundle unreadable"), "{message}");
1090 }
1091 state => panic!("corrupt bundle must settle terminal, got {state:?}"),
1092 }
1093 }
1094
1095 #[tokio::test]
1096 async fn submit_unauthorized_returns_entry_to_queued_hold() {
1097 let tmp = TempDir::new().unwrap();
1098 let root = tmp.path().join("feedback-outbox");
1099 let id = enqueue(&root, auth_lane());
1100 let transport = MockTransport::new(DrainAction::Acknowledge {
1101 server_id: "unused".into(),
1102 });
1103 *transport.script.lock().unwrap() = vec![Err(FeedbackTransportError::Unauthorized)];
1104 let mut last = None;
1105 let outcome =
1106 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1107 assert_eq!(outcome, RoundOutcome::AllHeld);
1108 assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
1109 }
1110
1111 #[tokio::test]
1112 async fn terminal_actionable_and_rejected_persist_and_never_retry() {
1113 let tmp = TempDir::new().unwrap();
1114 let root = tmp.path().join("feedback-outbox");
1115 let auth_id = enqueue(&root, auth_lane());
1116 let transport = MockTransport::new(DrainAction::TerminalActionable {
1117 reason: TransportActionableReason::ReconsentRequired,
1118 });
1119 let mut last = None;
1120 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1121 assert!(matches!(
1122 persisted_state(&root, &auth_id),
1123 SpoolState::TerminalActionable {
1124 reason: car_feedback_core::spool::ActionableReason::ReconsentRequired
1125 }
1126 ));
1127
1128 let before = transport.submit_calls.load(Ordering::SeqCst);
1130 let outcome =
1131 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1132 assert_eq!(outcome, RoundOutcome::Idle);
1133 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), before);
1134
1135 last = None;
1138 let rejected_id = enqueue(&root, auth_lane());
1139 *transport.script.lock().unwrap() = vec![Ok(DrainAction::TerminalRejected {
1140 message: "HTTP 400: description invalid".into(),
1141 })];
1142 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1143 assert_eq!(
1144 persisted_state(&root, &rejected_id),
1145 SpoolState::TerminalRejected {
1146 message: "HTTP 400: description invalid".to_string()
1147 }
1148 );
1149 }
1150
1151 #[tokio::test]
1152 async fn retriable_failure_requeues_durably_and_reports_backoff() {
1153 let tmp = TempDir::new().unwrap();
1154 let root = tmp.path().join("feedback-outbox");
1155 let id = enqueue(&root, auth_lane());
1156 let transport = MockTransport::new(DrainAction::Requeue { backoff: 30 });
1157 let mut last = None;
1158 let outcome =
1159 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1160 assert_eq!(
1161 outcome,
1162 RoundOutcome::Backoff {
1163 min_delay: Duration::from_secs(30)
1164 }
1165 );
1166 assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
1167 }
1168
1169 #[tokio::test(start_paused = true)]
1170 async fn retry_after_stops_the_round_and_is_honored_before_the_next_upload() {
1171 let tmp = TempDir::new().unwrap();
1175 let root = tmp.path().join("feedback-outbox");
1176 let first = enqueue(&root, auth_lane());
1177 let second = enqueue(&root, auth_lane());
1178 let transport = MockTransport::new(DrainAction::Acknowledge {
1179 server_id: "row".into(),
1180 });
1181 *transport.script.lock().unwrap() = vec![
1182 Ok(DrainAction::RequeueAfter { secs: 40 }),
1183 Ok(DrainAction::Acknowledge {
1184 server_id: "row-a".into(),
1185 }),
1186 Ok(DrainAction::Acknowledge {
1187 server_id: "row-b".into(),
1188 }),
1189 ];
1190 let mut last = None;
1191 let outcome =
1192 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1193 assert_eq!(outcome, RoundOutcome::RetryAfter(Duration::from_secs(40)));
1194 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 1);
1197 assert_eq!(persisted_state(&root, &first), SpoolState::Queued);
1198 assert_eq!(persisted_state(&root, &second), SpoolState::Queued);
1199
1200 tokio::time::sleep(Duration::from_secs(40)).await;
1202 let outcome =
1203 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1204 assert_eq!(outcome, RoundOutcome::Progressed);
1205 let stamps = transport.submit_at.lock().unwrap().clone();
1206 assert!(stamps.len() >= 2);
1207 assert!(
1208 stamps[1].duration_since(stamps[0]) >= Duration::from_secs(40),
1209 "second upload ran {:?} after the 429 — Retry-After not honored",
1210 stamps[1].duration_since(stamps[0])
1211 );
1212 }
1213
1214 #[tokio::test(start_paused = true)]
1215 async fn uploads_pace_at_most_one_per_min_interval() {
1216 let tmp = TempDir::new().unwrap();
1218 let root = tmp.path().join("feedback-outbox");
1219 enqueue(&root, auth_lane());
1220 enqueue(&root, auth_lane());
1221 let transport = MockTransport::new(DrainAction::Acknowledge {
1222 server_id: "row".into(),
1223 });
1224 let mut last = None;
1225 let outcome =
1226 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1227 assert_eq!(outcome, RoundOutcome::Progressed);
1228 let stamps = transport.submit_at.lock().unwrap().clone();
1229 assert_eq!(stamps.len(), 2);
1230 assert!(
1231 stamps[1].duration_since(stamps[0]) >= Duration::from_secs(15),
1232 "uploads {:?} apart — pacing not applied",
1233 stamps[1].duration_since(stamps[0])
1234 );
1235 }
1236
1237 #[tokio::test]
1240 async fn stale_sending_entry_from_a_crashed_drain_recovers_and_resends_same_id() {
1241 let tmp = TempDir::new().unwrap();
1245 let root = tmp.path().join("feedback-outbox");
1246 let id = enqueue(&root, auth_lane());
1247 let spool = Spool::open(&root).unwrap();
1248 spool.mark_sending(&id).unwrap();
1249 let original_csid = spool
1250 .list()
1251 .unwrap()
1252 .into_iter()
1253 .find(|e| e.id == id)
1254 .unwrap()
1255 .client_submission_id;
1256 drop(spool); let transport = MockTransport::new(DrainAction::Acknowledge {
1259 server_id: "row-1".into(),
1260 });
1261 let mut last = None;
1262 let outcome =
1263 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1264 assert_eq!(outcome, RoundOutcome::Progressed);
1265 assert_eq!(
1266 persisted_state(&root, &id),
1267 SpoolState::Acknowledged {
1268 server_id: "row-1".to_string()
1269 }
1270 );
1271 assert_eq!(
1272 transport.submitted_ids.lock().unwrap().as_slice(),
1273 &[original_csid],
1274 "the recovered entry must re-send its ORIGINAL idempotency key"
1275 );
1276 }
1277
1278 #[tokio::test(start_paused = true)]
1281 async fn spawned_drain_parks_idle_and_drains_on_wake() {
1282 let tmp = TempDir::new().unwrap();
1283 let car_home = tmp.path().to_path_buf();
1284 let root = car_home.join(FEEDBACK_OUTBOX_DIR);
1285 let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
1286 server_id: "row-1".into(),
1287 }));
1288 let handle = spawn_feedback_drain_with(car_home, transport.clone(), fast_config());
1289
1290 tokio::time::sleep(Duration::from_secs(3600)).await;
1295 assert_eq!(
1296 transport.total_calls(),
1297 0,
1298 "an empty-outbox park (incl. its local disk ticks) must never touch the transport"
1299 );
1300
1301 let id = enqueue(&root, auth_lane());
1303 handle.wake();
1304 for _ in 0..200 {
1306 tokio::task::yield_now().await;
1307 if transport.submit_calls.load(Ordering::SeqCst) > 0 {
1308 break;
1309 }
1310 tokio::time::sleep(Duration::from_millis(50)).await;
1311 }
1312 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 1);
1313 assert_eq!(
1314 persisted_state(&root, &id),
1315 SpoolState::Acknowledged {
1316 server_id: "row-1".to_string()
1317 }
1318 );
1319 }
1320
1321 #[tokio::test(start_paused = true)]
1324 async fn cli_direct_enqueue_while_parked_drains_within_one_tick() {
1325 let tmp = TempDir::new().unwrap();
1330 let car_home = tmp.path().to_path_buf();
1331 let root = car_home.join(FEEDBACK_OUTBOX_DIR);
1332 let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
1333 server_id: "row-cli".into(),
1334 }));
1335 let _handle = spawn_feedback_drain_with(car_home, transport.clone(), fast_config());
1336
1337 tokio::time::sleep(Duration::from_secs(1)).await;
1339 assert_eq!(transport.total_calls(), 0);
1340
1341 let id = enqueue(&root, auth_lane());
1343
1344 for _ in 0..200 {
1346 tokio::task::yield_now().await;
1347 if transport.submit_calls.load(Ordering::SeqCst) > 0 {
1348 break;
1349 }
1350 tokio::time::sleep(Duration::from_secs(1)).await;
1351 }
1352 assert_eq!(
1353 transport.submit_calls.load(Ordering::SeqCst),
1354 1,
1355 "a parked drain must catch a direct spool enqueue via the local tick"
1356 );
1357 assert_eq!(
1358 persisted_state(&root, &id),
1359 SpoolState::Acknowledged {
1360 server_id: "row-cli".to_string()
1361 }
1362 );
1363 }
1364
1365 #[test]
1371 fn park_probe_sees_only_pending_entries() {
1372 let tmp = TempDir::new().unwrap();
1373 let root = tmp.path().join(FEEDBACK_OUTBOX_DIR);
1374 assert!(!spool_has_pending_entries(&root));
1376 std::fs::create_dir_all(root.join(".tmp-half-written")).unwrap();
1377 std::fs::write(root.join("stray-file"), b"x").unwrap();
1378 assert!(
1379 !spool_has_pending_entries(&root),
1380 "staging dirs and files don't count"
1381 );
1382 std::fs::create_dir_all(root.join("00000000000000000000-not-an-entry")).unwrap();
1384 assert!(
1385 !spool_has_pending_entries(&root),
1386 "a stray directory is not a pending entry"
1387 );
1388 let settled = enqueue(&root, auth_lane());
1390 {
1391 let spool = Spool::open(&root).unwrap();
1392 spool.mark_sending(&settled).unwrap();
1393 spool.mark_acknowledged(&settled, "row-1").unwrap();
1394 }
1395 assert_eq!(
1396 persisted_state(&root, &settled),
1397 SpoolState::Acknowledged {
1398 server_id: "row-1".to_string()
1399 }
1400 );
1401 assert!(
1402 !spool_has_pending_entries(&root),
1403 "a settled entry must not break the park"
1404 );
1405 let queued = enqueue(&root, auth_lane());
1407 assert!(spool_has_pending_entries(&root));
1408 Spool::open(&root).unwrap().mark_sending(&queued).unwrap();
1410 assert!(spool_has_pending_entries(&root));
1411 }
1412
1413 #[tokio::test(start_paused = true)]
1421 async fn absurd_retry_after_is_clamped_to_the_backoff_ceiling() {
1422 let tmp = TempDir::new().unwrap();
1423 let car_home = tmp.path().to_path_buf();
1424 let root = car_home.join(FEEDBACK_OUTBOX_DIR);
1425 let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
1426 server_id: "row-1".into(),
1427 }));
1428 *transport.script.lock().unwrap() = vec![
1429 Ok(DrainAction::RequeueAfter { secs: 999_999_999 }),
1430 Ok(DrainAction::Acknowledge {
1431 server_id: "row-1".into(),
1432 }),
1433 ];
1434 let handle = spawn_feedback_drain_with(car_home, transport.clone(), fast_config());
1435 tokio::time::sleep(Duration::from_secs(1)).await;
1436 let id = enqueue(&root, auth_lane());
1437 handle.wake();
1438 for _ in 0..200 {
1440 tokio::task::yield_now().await;
1441 if transport.submit_calls.load(Ordering::SeqCst) >= 1 {
1442 break;
1443 }
1444 tokio::time::sleep(Duration::from_millis(50)).await;
1445 }
1446 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 1);
1447 assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
1448
1449 tokio::time::sleep(Duration::from_secs(60)).await;
1452 assert_eq!(
1453 transport.submit_calls.load(Ordering::SeqCst),
1454 2,
1455 "the drain must retry within the backoff ceiling, not the header's decades"
1456 );
1457 assert_eq!(
1458 persisted_state(&root, &id),
1459 SpoolState::Acknowledged {
1460 server_id: "row-1".to_string()
1461 }
1462 );
1463 }
1464
1465 #[tokio::test(start_paused = true)]
1471 async fn submit_wake_ends_a_retry_after_wait_early() {
1472 let tmp = TempDir::new().unwrap();
1473 let car_home = tmp.path().to_path_buf();
1474 let root = car_home.join(FEEDBACK_OUTBOX_DIR);
1475 let transport = Arc::new(MockTransport::new(DrainAction::Acknowledge {
1476 server_id: "row".into(),
1477 }));
1478 *transport.script.lock().unwrap() = vec![
1479 Ok(DrainAction::RequeueAfter { secs: 300 }),
1480 Ok(DrainAction::Acknowledge {
1481 server_id: "row-a".into(),
1482 }),
1483 Ok(DrainAction::Acknowledge {
1484 server_id: "row-b".into(),
1485 }),
1486 ];
1487 let mut config = fast_config();
1490 config.max_backoff = Duration::from_secs(600);
1491 let handle = spawn_feedback_drain_with(car_home, transport.clone(), config);
1492 tokio::time::sleep(Duration::from_secs(1)).await;
1493 let first = enqueue(&root, auth_lane());
1494 handle.wake();
1495 for _ in 0..200 {
1496 tokio::task::yield_now().await;
1497 if transport.submit_calls.load(Ordering::SeqCst) >= 1 {
1498 break;
1499 }
1500 tokio::time::sleep(Duration::from_millis(50)).await;
1501 }
1502 assert_eq!(
1503 transport.submit_calls.load(Ordering::SeqCst),
1504 1,
1505 "the 429 landed"
1506 );
1507 assert_eq!(persisted_state(&root, &first), SpoolState::Queued);
1508
1509 let second = enqueue(&root, auth_lane());
1511 handle.wake();
1512 tokio::time::sleep(Duration::from_secs(30)).await;
1513 assert!(
1514 transport.submit_calls.load(Ordering::SeqCst) >= 2,
1515 "a submit wake must end the Retry-After wait (calls: {})",
1516 transport.submit_calls.load(Ordering::SeqCst)
1517 );
1518 assert_eq!(
1519 persisted_state(&root, &first),
1520 SpoolState::Acknowledged {
1521 server_id: "row-a".to_string()
1522 },
1523 "the oldest queued entry is retried first"
1524 );
1525 tokio::time::sleep(Duration::from_secs(30)).await;
1527 assert_eq!(
1528 persisted_state(&root, &second),
1529 SpoolState::Acknowledged {
1530 server_id: "row-b".to_string()
1531 }
1532 );
1533 }
1534
1535 #[tokio::test]
1538 async fn rejected_upload_with_omissions_persists_the_omitted_note() {
1539 let tmp = TempDir::new().unwrap();
1540 let root = tmp.path().join("feedback-outbox");
1541 let id = enqueue(&root, auth_lane());
1542 let transport = MockTransport::new(DrainAction::TerminalRejected {
1543 message: "HTTP 400: description invalid".into(),
1544 });
1545 *transport.omitted.lock().unwrap() = vec!["screenshot dropped (413 fallback)".to_string()];
1546 let mut last = None;
1547 run_drain_round(&root, &transport, &fast_config(), &mut last, &mut ledger()).await;
1548 match persisted_state(&root, &id) {
1549 SpoolState::TerminalRejected { message } => {
1550 assert!(
1551 message.contains("screenshot dropped (413 fallback)"),
1552 "the omitted note must persist in the durable rejection message: {message}"
1553 );
1554 assert!(message.contains("HTTP 400"));
1555 }
1556 other => panic!("expected TerminalRejected, got {other:?}"),
1557 }
1558 }
1559
1560 #[tokio::test(start_paused = true)]
1571 async fn retriable_failures_stop_after_the_per_process_cap_and_the_entry_stays_queued() {
1572 let tmp = TempDir::new().unwrap();
1573 let root = tmp.path().join("feedback-outbox");
1574 let id = enqueue(&root, auth_lane());
1575 let transport = MockTransport::new(DrainAction::Requeue { backoff: 1 });
1576 *transport.script.lock().unwrap() = vec![
1577 Ok(DrainAction::Requeue { backoff: 1 }),
1578 Err(FeedbackTransportError::FetchFailed(
1579 "connection reset".into(),
1580 )),
1581 Ok(DrainAction::Requeue { backoff: 1 }),
1582 ];
1583 let mut config = fast_config();
1584 config.max_retriable_attempts = 3;
1585 let mut last = None;
1586 let mut ledger = RetryLedger::default();
1587
1588 for round in 1..=3 {
1589 let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
1590 assert!(
1591 matches!(outcome, RoundOutcome::Backoff { .. }),
1592 "round {round}: {outcome:?}"
1593 );
1594 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), round);
1595 assert_eq!(persisted_state(&root, &id), SpoolState::Queued);
1596 }
1597
1598 let probes_before = transport.eligible_calls.load(Ordering::SeqCst);
1600 for _ in 0..3 {
1601 let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
1602 assert_eq!(outcome, RoundOutcome::AllHeld);
1603 }
1604 assert_eq!(
1605 transport.submit_calls.load(Ordering::SeqCst),
1606 3,
1607 "no submit past the cap"
1608 );
1609 assert_eq!(
1610 transport.eligible_calls.load(Ordering::SeqCst),
1611 probes_before,
1612 "a capped entry costs no network — not even the eligibility probe"
1613 );
1614 assert_eq!(
1615 persisted_state(&root, &id),
1616 SpoolState::Queued,
1617 "capped ⇒ still Queued on disk: never a terminal state, never pruned"
1618 );
1619
1620 let mut restarted = RetryLedger::default();
1622 let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut restarted).await;
1623 assert!(matches!(outcome, RoundOutcome::Backoff { .. }));
1624 assert_eq!(transport.submit_calls.load(Ordering::SeqCst), 4);
1625 }
1626
1627 #[tokio::test(start_paused = true)]
1632 async fn retry_after_does_not_consume_the_attempt_budget() {
1633 let tmp = TempDir::new().unwrap();
1634 let root = tmp.path().join("feedback-outbox");
1635 let id = enqueue(&root, auth_lane());
1636 let transport = MockTransport::new(DrainAction::Acknowledge {
1637 server_id: "row".into(),
1638 });
1639 *transport.script.lock().unwrap() = vec![
1640 Ok(DrainAction::RequeueAfter { secs: 1 }),
1641 Ok(DrainAction::RequeueAfter { secs: 1 }),
1642 Ok(DrainAction::Acknowledge {
1643 server_id: "row-1".into(),
1644 }),
1645 ];
1646 let mut config = fast_config();
1647 config.max_retriable_attempts = 1;
1648 let mut last = None;
1649 let mut ledger = RetryLedger::default();
1650
1651 for _ in 0..2 {
1652 let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
1653 assert_eq!(outcome, RoundOutcome::RetryAfter(Duration::from_secs(1)));
1654 tokio::time::sleep(Duration::from_secs(1)).await;
1655 }
1656 let outcome = run_drain_round(&root, &transport, &config, &mut last, &mut ledger).await;
1657 assert_eq!(outcome, RoundOutcome::Progressed);
1658 assert_eq!(
1659 persisted_state(&root, &id),
1660 SpoolState::Acknowledged {
1661 server_id: "row-1".to_string()
1662 }
1663 );
1664 }
1665}