#[cfg(feature = "sync-sender-qwp-ws")]
use std::collections::VecDeque;
use std::fs;
use std::io;
use std::path::{Path, PathBuf};
#[cfg(feature = "sync-sender-qwp-ws")]
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
#[cfg(feature = "sync-sender-qwp-ws")]
use std::sync::{Arc, Mutex};
#[cfg(feature = "sync-sender-qwp-ws")]
use std::thread;
#[cfg(feature = "sync-sender-qwp-ws")]
use std::time::{Duration, Instant};
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ErrorCode;
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::buffer::SymbolGlobalDict;
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::conf::QwpWsConfig;
use crate::ingress::conf::QwpWsManagedSlotExclusion;
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::tls::TlsSettings;
#[cfg(feature = "sync-sender-qwp-ws")]
use super::qwp_ws::{QwpWsConnectKind, TrafficGate, try_dup_recovered};
#[cfg(feature = "sync-sender-qwp-ws")]
use super::qwp_ws_driver::{
BlockingQwpWsTransport, CloseStepOutcome, DEFAULT_EVENT_CAPACITY, DriverError, PublicationLog,
QwpWsPublicationStore, QwpWsSendCore, ReconnectPolicy,
};
#[cfg(feature = "sync-sender-qwp-ws")]
use super::qwp_ws_sfa_queue::{SfaQueueError, SfaQueueOptions};
#[cfg(feature = "sync-sender-qwp-ws")]
use super::qwp_ws_sfa_slot::SfaSlotQueue;
pub(crate) const FAILED_SENTINEL_NAME: &str = ".failed";
#[cfg(feature = "sync-sender-qwp-ws")]
const LAST_ERROR_NAME: &str = ".last_error";
#[cfg(feature = "sync-sender-qwp-ws")]
const ORPHAN_IDLE_PARK: Duration = Duration::from_millis(50);
#[cfg(feature = "sync-sender-qwp-ws")]
const ORPHAN_POOL_GRACEFUL_DRAIN: Duration = Duration::from_millis(2500);
#[cfg(feature = "sync-sender-qwp-ws")]
const ORPHAN_POOL_STOP_GRACE: Duration = Duration::from_millis(500);
#[cfg(feature = "sync-sender-qwp-ws")]
const ORPHAN_POOL_CLOSE_POLL: Duration = Duration::from_millis(10);
pub(crate) fn scan_orphan_slots(
sf_dir: &Path,
own_sender_id: &str,
managed_exclusions: &[QwpWsManagedSlotExclusion],
) -> Vec<PathBuf> {
let Ok(entries) = fs::read_dir(sf_dir) else {
return Vec::new();
};
let mut orphans = Vec::new();
for entry in entries.flatten() {
let slot_path = entry.path();
if !slot_path.is_dir() {
continue;
}
if entry.file_name().to_str() == Some(own_sender_id) {
continue;
}
if entry.file_name().to_str().is_some_and(|name| {
managed_exclusions
.iter()
.any(|exclusion| exclusion.matches(name))
}) {
continue;
}
if !is_candidate_orphan(&slot_path) {
continue;
}
orphans.push(slot_path);
}
orphans
}
pub(crate) fn is_candidate_orphan(slot_dir: &Path) -> bool {
!has_failed_sentinel(slot_dir) && has_any_sfa_file(slot_dir)
}
pub(crate) fn has_failed_sentinel(slot_dir: &Path) -> bool {
slot_dir.join(FAILED_SENTINEL_NAME).exists()
}
pub(crate) fn mark_failed(slot_dir: &Path, reason: &str) -> io::Result<()> {
fs::write(slot_dir.join(FAILED_SENTINEL_NAME), reason)
}
pub(crate) fn has_any_sfa_file(slot_dir: &Path) -> bool {
let Ok(entries) = fs::read_dir(slot_dir) else {
return false;
};
entries
.flatten()
.any(|entry| entry.path().file_name().is_some_and(is_sfa_file_name))
}
fn is_sfa_file_name(name: &std::ffi::OsStr) -> bool {
name.to_str().is_some_and(|name| name.ends_with(".sfa"))
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[derive(Clone)]
pub(crate) struct OrphanDrainerConfig {
host: String,
port: String,
use_tls: bool,
tls_settings: Option<TlsSettings>,
qwp_ws: QwpWsConfig,
auth_header: Option<String>,
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl OrphanDrainerConfig {
pub(crate) fn new(
host: &str,
port: &str,
use_tls: bool,
tls_settings: Option<TlsSettings>,
qwp_ws: &QwpWsConfig,
auth_header: Option<String>,
) -> Self {
Self {
host: host.to_owned(),
port: port.to_owned(),
use_tls,
tls_settings,
qwp_ws: qwp_ws.clone(),
auth_header,
}
}
fn queue_options(&self, slot_dir: PathBuf) -> Result<SfaQueueOptions, String> {
let max_bytes = usize::try_from(self.qwp_ws.sf_max_total_bytes())
.map_err(|_| "sf_max_total_bytes value is too large for this platform".to_owned())?;
Ok(SfaQueueOptions {
slot_dir,
segment_size_bytes: *self.qwp_ws.sf_max_segment_bytes,
max_bytes,
periodic_sync_interval: self.qwp_ws.periodic_sync_interval(),
})
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
pub(crate) struct OrphanDrainerPool {
stop: Arc<AtomicBool>,
threads: Vec<(thread::JoinHandle<()>, Arc<TrafficGate>)>,
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl OrphanDrainerPool {
pub(crate) fn start(
candidates: Vec<PathBuf>,
max_background_drainers: usize,
config: OrphanDrainerConfig,
) -> Option<Self> {
if candidates.is_empty() || max_background_drainers == 0 {
return None;
}
let pending = Arc::new(Mutex::new(VecDeque::from(candidates)));
let stop = Arc::new(AtomicBool::new(false));
let worker_count = max_background_drainers.min(pending.lock().ok()?.len());
let mut threads = Vec::with_capacity(worker_count);
for _ in 0..worker_count {
let pending = Arc::clone(&pending);
let stop = Arc::clone(&stop);
let config = config.clone();
let traffic_gate = Arc::new(TrafficGate::default());
let worker_gate = Arc::clone(&traffic_gate);
let worker = thread::spawn(move || {
while !stop.load(Ordering::Acquire) {
let Some(slot_dir) = pop_pending_orphan(&pending) else {
break;
};
drain_orphan_to_completion(slot_dir, &config, &stop, &worker_gate);
}
});
threads.push((worker, traffic_gate));
}
Some(Self { stop, threads })
}
pub(crate) fn close(&mut self) {
self.close_with_timeouts(ORPHAN_POOL_GRACEFUL_DRAIN, ORPHAN_POOL_STOP_GRACE);
}
fn close_with_timeouts(&mut self, graceful_drain: Duration, stop_grace: Duration) {
self.join_finished_threads();
self.wait_for_finished_threads(graceful_drain);
self.stop.store(true, Ordering::Release);
self.join_finished_threads();
for (_, traffic_gate) in &self.threads {
if let Err(err) = traffic_gate.shutdown() {
log::warn!("could not shut down QWP/WebSocket orphan drainer traffic: {err}");
}
}
self.wait_for_finished_threads(stop_grace);
if !self.threads.is_empty() {
log::error!(
"{} QWP/WebSocket orphan drainer worker(s) did not stop within {:?}; \
detaching them while their active slot locks remain retained",
self.threads.len(),
stop_grace
);
}
self.threads.clear();
}
fn wait_for_finished_threads(&mut self, timeout: Duration) {
let Some(deadline) = Instant::now().checked_add(timeout) else {
return;
};
while !self.threads.is_empty() {
self.join_finished_threads();
if self.threads.is_empty() {
return;
}
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return;
}
thread::sleep(remaining.min(ORPHAN_POOL_CLOSE_POLL));
}
}
fn join_finished_threads(&mut self) {
let mut running = Vec::with_capacity(self.threads.len());
for (thread, traffic_gate) in self.threads.drain(..) {
if thread.is_finished() {
let _ = thread.join();
} else {
running.push((thread, traffic_gate));
}
}
self.threads = running;
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl Drop for OrphanDrainerPool {
fn drop(&mut self) {
self.close();
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
pub(crate) struct ManualOrphanDrainers {
pending: VecDeque<PendingOrphanSlot>,
active: Vec<ActiveOrphanDrainer>,
next_active: usize,
max_active: usize,
config: OrphanDrainerConfig,
}
#[cfg(feature = "sync-sender-qwp-ws")]
struct PendingOrphanSlot {
slot_dir: PathBuf,
not_before: Option<Instant>,
backoff: Duration,
}
#[cfg(feature = "sync-sender-qwp-ws")]
struct ActiveOrphanDrainer {
drainer: OrphanDrainer,
gated_until: Option<Instant>,
backoff: Duration,
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl ManualOrphanDrainers {
pub(crate) fn new(
candidates: Vec<PathBuf>,
max_active: usize,
config: OrphanDrainerConfig,
) -> Option<Self> {
if candidates.is_empty() || max_active == 0 {
return None;
}
let initial_backoff = *config.qwp_ws.reconnect_initial_backoff;
Some(Self {
pending: candidates
.into_iter()
.map(|slot_dir| PendingOrphanSlot {
slot_dir,
not_before: None,
backoff: initial_backoff,
})
.collect(),
active: Vec::new(),
next_active: 0,
max_active,
config,
})
}
pub(crate) fn drive_once(&mut self) -> bool {
if self.active.len() < self.max_active && self.activate_one() {
return true;
}
self.drive_active_once()
}
fn initial_backoff(&self) -> Duration {
*self.config.qwp_ws.reconnect_initial_backoff
}
fn max_backoff(&self) -> Duration {
*self.config.qwp_ws.reconnect_max_backoff
}
fn activate_one(&mut self) -> bool {
let now = Instant::now();
for _ in 0..self.pending.len() {
let Some(slot) = self.pending.pop_front() else {
return false;
};
if slot.not_before.is_some_and(|not_before| now < not_before) {
self.pending.push_back(slot);
continue;
}
return self.open_pending(slot);
}
false
}
fn open_pending(&mut self, slot: PendingOrphanSlot) -> bool {
if has_failed_sentinel(&slot.slot_dir) {
return true;
}
match OrphanDrainer::open(slot.slot_dir.clone(), &self.config) {
OrphanOpenOutcome::Drainer(drainer) => {
self.active.push(ActiveOrphanDrainer {
drainer: *drainer,
gated_until: None,
backoff: self.initial_backoff(),
});
true
}
OrphanOpenOutcome::AlreadyDrained
| OrphanOpenOutcome::FailedSentinel
| OrphanOpenOutcome::Locked
| OrphanOpenOutcome::Stopped => true,
OrphanOpenOutcome::RetryLater(reason) => {
let _ = record_last_error(&slot.slot_dir, &reason);
self.requeue(slot.slot_dir, slot.backoff);
true
}
OrphanOpenOutcome::Unrecoverable(reason) => {
let _ = mark_failed(&slot.slot_dir, &reason);
true
}
}
}
fn requeue(&mut self, slot_dir: PathBuf, backoff: Duration) {
self.pending.push_back(PendingOrphanSlot {
slot_dir,
not_before: Instant::now().checked_add(backoff),
backoff: next_orphan_retry_backoff(backoff, self.max_backoff()),
});
}
fn drive_active_once(&mut self) -> bool {
if self.active.is_empty() {
return false;
}
let initial_backoff = self.initial_backoff();
let max_backoff = self.max_backoff();
let active_count = self.active.len();
for _ in 0..active_count {
let index = self.next_active % self.active.len();
self.next_active = (index + 1) % self.active.len();
let now = Instant::now();
if let Some(gated_until) = self.active[index].gated_until {
if now < gated_until {
continue;
}
self.active[index].gated_until = None;
}
match self.active[index].drainer.drive_once(None) {
OrphanDriveOutcome::Idle => {
self.active[index].backoff = initial_backoff;
}
OrphanDriveOutcome::Progress => {
self.active[index].backoff = initial_backoff;
return true;
}
OrphanDriveOutcome::Waiting {
sleep_for,
deadline,
} => {
let until = now.checked_add(sleep_for);
self.active[index].gated_until = match (until, deadline) {
(Some(until), Some(deadline)) => Some(until.min(deadline)),
(Some(until), None) => Some(until),
(None, deadline) => deadline,
};
return true;
}
OrphanDriveOutcome::Drained => {
let entry = self.active.remove(index);
entry.drainer.clear_last_error();
self.normalize_next_active();
return true;
}
OrphanDriveOutcome::RetryLater(reason) => {
let entry = &mut self.active[index];
entry.drainer.record_last_error(&reason);
entry.gated_until = now.checked_add(entry.backoff);
entry.backoff = next_orphan_retry_backoff(entry.backoff, max_backoff);
return true;
}
OrphanDriveOutcome::Terminal(reason) => {
let entry = self.active.remove(index);
entry.drainer.record_last_error(&reason);
self.normalize_next_active();
return true;
}
OrphanDriveOutcome::Unrecoverable(reason) => {
let entry = self.active.remove(index);
let _ = mark_failed(&entry.drainer.slot_dir, &reason);
self.normalize_next_active();
return true;
}
OrphanDriveOutcome::Stopped => {
self.active.remove(index);
self.normalize_next_active();
return true;
}
}
}
false
}
fn normalize_next_active(&mut self) {
if self.active.is_empty() {
self.next_active = 0;
} else {
self.next_active %= self.active.len();
}
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
struct OrphanDrainer {
slot_dir: PathBuf,
store: QwpWsPublicationStore<SfaSlotQueue>,
send_core: QwpWsSendCore<BlockingQwpWsTransport>,
}
#[cfg(feature = "sync-sender-qwp-ws")]
enum OrphanOpenOutcome {
Drainer(Box<OrphanDrainer>),
AlreadyDrained,
FailedSentinel,
Locked,
RetryLater(String),
Stopped,
Unrecoverable(String),
}
#[cfg(feature = "sync-sender-qwp-ws")]
enum OrphanDriveOutcome {
Idle,
Progress,
Drained,
Stopped,
Waiting {
sleep_for: Duration,
deadline: Option<Instant>,
},
RetryLater(String),
Terminal(String),
Unrecoverable(String),
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl OrphanDrainer {
fn open(slot_dir: PathBuf, config: &OrphanDrainerConfig) -> OrphanOpenOutcome {
Self::open_inner(slot_dir, config, None, None)
}
fn open_with_stop(
slot_dir: PathBuf,
config: &OrphanDrainerConfig,
stop: &AtomicBool,
traffic_gate: &Arc<TrafficGate>,
) -> OrphanOpenOutcome {
Self::open_inner(slot_dir, config, Some(stop), Some(traffic_gate))
}
fn open_inner(
slot_dir: PathBuf,
config: &OrphanDrainerConfig,
stop: Option<&AtomicBool>,
traffic_gate: Option<&Arc<TrafficGate>>,
) -> OrphanOpenOutcome {
if orphan_stop_requested(stop) {
return OrphanOpenOutcome::Stopped;
}
if has_failed_sentinel(&slot_dir) {
return OrphanOpenOutcome::FailedSentinel;
}
let options = match config.queue_options(slot_dir.clone()) {
Ok(options) => options,
Err(err) => return OrphanOpenOutcome::RetryLater(err),
};
let mut queue = match SfaSlotQueue::open_replay_only_existing(options) {
Ok(queue) => queue,
Err(SfaQueueError::SlotInUse { .. }) => return OrphanOpenOutcome::Locked,
Err(SfaQueueError::SanitizedResidue { .. }) => {
let retry_options = match config.queue_options(slot_dir.clone()) {
Ok(options) => options,
Err(err) => return OrphanOpenOutcome::RetryLater(err),
};
match SfaSlotQueue::open_replay_only_existing(retry_options) {
Ok(queue) => queue,
Err(SfaQueueError::SlotInUse { .. }) => {
return OrphanOpenOutcome::Locked;
}
Err(
err @ (SfaQueueError::SanitizedResidue { .. }
| SfaQueueError::Recovery { .. }
| SfaQueueError::CorruptSegments { .. }),
) => {
return OrphanOpenOutcome::Unrecoverable(format!("{err:?}"));
}
Err(err) => return OrphanOpenOutcome::RetryLater(format!("{err:?}")),
}
}
Err(err @ (SfaQueueError::Recovery { .. } | SfaQueueError::CorruptSegments { .. })) => {
return OrphanOpenOutcome::Unrecoverable(format!("{err:?}"));
}
Err(err) => return OrphanOpenOutcome::RetryLater(format!("{err:?}")),
};
if orphan_stop_requested(stop) {
return OrphanOpenOutcome::Stopped;
}
if has_failed_sentinel(&slot_dir) {
return OrphanOpenOutcome::FailedSentinel;
}
if orphan_queue_drained(&queue) {
if let Err(err) = queue.close() {
let reason = format!("{err:?}");
return retry_open_later(reason, stop);
}
let _ = clear_last_error(&slot_dir);
return OrphanOpenOutcome::AlreadyDrained;
}
let delta_dict_enabled = queue.is_delta_dict_enabled();
let recovered_dict_count = queue.recovered_symbol_dict_count();
let recovered_dict = if delta_dict_enabled {
let entries = match try_dup_recovered(queue.recovered_symbol_dict_entries()) {
Ok(entries) => entries,
Err(err) => {
let reason = format!("recovered symbol dictionary allocation failed: {err}");
log::warn!(
"QWP/WebSocket orphan slot {}: {reason}; retrying adoption later",
slot_dir.display()
);
return retry_open_later(reason, stop);
}
};
match SymbolGlobalDict::new().seed(&entries, recovered_dict_count) {
Ok(()) => Some((entries, recovered_dict_count)),
Err(err) => {
log::warn!(
"QWP/WebSocket orphan slot {}: persisted symbol dictionary was \
rejected ({err}); draining with the dictionary the stored \
frames carry instead -- a frame that needs an id they cannot \
supply is rejected as resend-required, but the rest still \
drain.",
slot_dir.display()
);
Some((Vec::new(), 0))
}
}
} else {
None
};
if orphan_stop_requested(stop) {
return OrphanOpenOutcome::Stopped;
}
drop(queue.take_persisted_symbol_dict());
let transport = match BlockingQwpWsTransport::connect(
config.host.clone(),
config.port.clone(),
config.use_tls,
config.tls_settings.clone(),
QwpWsConnectKind::BackgroundDrainer,
config.qwp_ws.clone(),
config.auth_header.clone(),
Arc::new(AtomicUsize::new(0)),
traffic_gate.cloned(),
) {
Ok(transport) => transport,
Err(err) => return retry_open_later(err.to_string(), stop),
};
if orphan_stop_requested(stop) {
return OrphanOpenOutcome::Stopped;
}
let store = QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY);
let mut send_core = QwpWsSendCore::new_with_durable_ack_and_rejection_limit(
transport,
ReconnectPolicy::bounded(
*config.qwp_ws.reconnect_max_duration,
*config.qwp_ws.reconnect_initial_backoff,
*config.qwp_ws.reconnect_max_backoff,
),
*config.qwp_ws.request_durable_ack,
*config.qwp_ws.max_frame_rejections,
*config.qwp_ws.poison_min_escalation_window,
);
if let Some((recovered_dict_entries, recovered_dict_count)) = recovered_dict {
send_core.enable_delta_dict_owned(recovered_dict_entries, recovered_dict_count);
}
OrphanOpenOutcome::Drainer(Box::new(Self {
slot_dir,
store,
send_core,
}))
}
fn drive_once(&mut self, stop: Option<&AtomicBool>) -> OrphanDriveOutcome {
if orphan_stop_requested(stop) {
return OrphanDriveOutcome::Stopped;
}
match self.send_core.close_drain_ready_step(&mut self.store) {
Ok(CloseStepOutcome::Drained) => OrphanDriveOutcome::Drained,
Ok(CloseStepOutcome::Terminal) => self.terminal_outcome(),
Ok(CloseStepOutcome::Waiting {
sleep_for,
deadline,
}) => OrphanDriveOutcome::Waiting {
sleep_for,
deadline,
},
Ok(CloseStepOutcome::Progress) => OrphanDriveOutcome::Progress,
Ok(CloseStepOutcome::Idle) => OrphanDriveOutcome::Idle,
Err(err) => OrphanDriveOutcome::RetryLater(driver_error_message(err)),
}
}
fn record_last_error(&self, reason: &str) {
let _ = record_last_error(&self.slot_dir, reason);
}
fn clear_last_error(&self) {
let _ = clear_last_error(&self.slot_dir);
}
fn terminal_outcome(&self) -> OrphanDriveOutcome {
let reason = self.terminal_message();
if terminal_error_is_proven_local_unrecoverable(self.store.terminal_error()) {
OrphanDriveOutcome::Unrecoverable(reason)
} else {
OrphanDriveOutcome::Terminal(reason)
}
}
fn terminal_message(&self) -> String {
if let Some(err) = self.store.terminal_error() {
return err.to_string();
}
if let Some(err) = self.store.terminal_sender_error() {
return format!("{err:?}");
}
"orphan drainer reached terminal state".to_owned()
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn terminal_error_is_proven_local_unrecoverable(error: Option<&crate::Error>) -> bool {
error.is_some_and(|err| {
err.code() == ErrorCode::StoreResendRequired
})
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn orphan_reconnect_deadline_expired(deadline: Option<Instant>) -> bool {
deadline.is_some_and(|deadline| Instant::now() >= deadline)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn sleep_before_orphan_reconnect(
deadline: Option<Instant>,
backoff: Duration,
stop: Option<&AtomicBool>,
) -> bool {
if stop.is_some_and(|stop| stop.load(Ordering::Acquire)) {
return false;
}
if backoff.is_zero() {
return stop.is_none_or(|stop| !stop.load(Ordering::Acquire))
&& !orphan_reconnect_deadline_expired(deadline);
}
let sleep_for = match deadline {
Some(deadline) => {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return false;
}
backoff.min(remaining)
}
None => backoff,
};
let sleep_start = Instant::now();
while stop.is_none_or(|stop| !stop.load(Ordering::Acquire)) {
let remaining = sleep_for.saturating_sub(sleep_start.elapsed());
if remaining.is_zero() {
break;
}
thread::sleep(remaining.min(Duration::from_millis(50)));
}
stop.is_none_or(|stop| !stop.load(Ordering::Acquire))
&& !orphan_reconnect_deadline_expired(deadline)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn next_orphan_retry_backoff(current: Duration, max: Duration) -> Duration {
current.checked_mul(2).unwrap_or(Duration::MAX).min(max)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn drain_orphan_to_completion(
slot_dir: PathBuf,
config: &OrphanDrainerConfig,
stop: &AtomicBool,
traffic_gate: &Arc<TrafficGate>,
) {
let initial_backoff = *config.qwp_ws.reconnect_initial_backoff;
let max_backoff = *config.qwp_ws.reconnect_max_backoff;
let mut retry_backoff = initial_backoff;
let mut drainer = loop {
match OrphanDrainer::open_with_stop(slot_dir.clone(), config, stop, traffic_gate) {
OrphanOpenOutcome::Drainer(drainer) => break drainer,
OrphanOpenOutcome::AlreadyDrained
| OrphanOpenOutcome::FailedSentinel
| OrphanOpenOutcome::Locked
| OrphanOpenOutcome::Stopped => return,
OrphanOpenOutcome::RetryLater(reason) => {
record_last_error_unless_stopped(&slot_dir, &reason, stop);
if !sleep_before_orphan_reconnect(None, retry_backoff, Some(stop)) {
return;
}
retry_backoff = next_orphan_retry_backoff(retry_backoff, max_backoff);
}
OrphanOpenOutcome::Unrecoverable(reason) => {
mark_orphan_failed_unless_stopped(&slot_dir, &reason, stop);
return;
}
}
};
retry_backoff = initial_backoff;
while !stop.load(Ordering::Acquire) {
match drainer.drive_once(Some(stop)) {
OrphanDriveOutcome::Drained => {
drainer.clear_last_error();
return;
}
OrphanDriveOutcome::Waiting {
sleep_for,
deadline,
} => {
if !sleep_before_orphan_reconnect(deadline, sleep_for, Some(stop))
&& stop.load(Ordering::Acquire)
{
return;
}
retry_backoff = initial_backoff;
}
OrphanDriveOutcome::RetryLater(reason) => {
if !stop.load(Ordering::Acquire) {
drainer.record_last_error(&reason);
}
if !sleep_before_orphan_reconnect(None, retry_backoff, Some(stop)) {
return;
}
retry_backoff = next_orphan_retry_backoff(retry_backoff, max_backoff);
}
OrphanDriveOutcome::Terminal(reason) => {
record_last_error_unless_stopped(&slot_dir, &reason, stop);
return;
}
OrphanDriveOutcome::Unrecoverable(reason) => {
mark_orphan_failed_unless_stopped(&slot_dir, &reason, stop);
return;
}
OrphanDriveOutcome::Stopped => return,
OrphanDriveOutcome::Progress => {
retry_backoff = initial_backoff;
}
OrphanDriveOutcome::Idle => {
retry_backoff = initial_backoff;
thread::sleep(ORPHAN_IDLE_PARK);
}
}
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn pop_pending_orphan(pending: &Mutex<VecDeque<PathBuf>>) -> Option<PathBuf> {
pending.lock().ok()?.pop_front()
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn mark_orphan_failed_unless_stopped(slot_dir: &Path, reason: &str, stop: &AtomicBool) {
if !stop.load(Ordering::Acquire) {
let _ = mark_failed(slot_dir, reason);
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn record_last_error(slot_dir: &Path, reason: &str) -> io::Result<()> {
fs::write(slot_dir.join(LAST_ERROR_NAME), reason)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn clear_last_error(slot_dir: &Path) -> io::Result<()> {
match fs::remove_file(slot_dir.join(LAST_ERROR_NAME)) {
Ok(()) => Ok(()),
Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(()),
Err(err) => Err(err),
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn record_last_error_unless_stopped(slot_dir: &Path, reason: &str, stop: &AtomicBool) {
if !stop.load(Ordering::Acquire) {
let _ = record_last_error(slot_dir, reason);
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn retry_open_later(reason: String, stop: Option<&AtomicBool>) -> OrphanOpenOutcome {
if orphan_stop_requested(stop) {
OrphanOpenOutcome::Stopped
} else {
OrphanOpenOutcome::RetryLater(reason)
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn orphan_stop_requested(stop: Option<&AtomicBool>) -> bool {
stop.is_some_and(|stop| stop.load(Ordering::Acquire))
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn orphan_queue_drained(queue: &SfaSlotQueue) -> bool {
match PublicationLog::published_fsn(queue) {
None => true,
Some(published) => {
PublicationLog::completed_fsn(queue).is_some_and(|completed| completed >= published)
}
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn driver_error_message(err: DriverError) -> String {
format!("{err:?}")
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::conf::{ConfigSetting, QwpWsConfig};
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
use crate::ingress::sender::qwp_ws_sfa_manifest::{SfManifest, SfaAckWatermark};
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
use crate::ingress::sender::qwp_ws_sfa_queue::SfaStorageStep;
#[cfg(all(feature = "sync-sender-qwp-ws", unix))]
use crate::ingress::sender::qwp_ws_sfa_segment::initial_segment_path;
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
use crate::ingress::sender::qwp_ws_sfa_segment::{SfaSegment, scan_file, spare_segment_path};
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
use crate::ingress::sender::qwp_ws_sfa_slot::SfaSlotOptions;
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::{QwpWsErrorCategory, QwpWsErrorPolicy, QwpWsSenderError};
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
use std::net::TcpListener;
#[test]
fn scan_returns_no_orphans_for_missing_root() {
let temp = TempDir::new().unwrap();
assert!(scan_orphan_slots(&temp.path().join("missing"), "default", &[]).is_empty());
}
#[test]
fn scan_filters_own_slot_failed_slots_and_empty_dirs() {
let temp = TempDir::new().unwrap();
let own = temp.path().join("default");
let failed = temp.path().join("failed");
let empty = temp.path().join("empty");
let orphan = temp.path().join("orphan");
fs::create_dir(&own).unwrap();
fs::create_dir(&failed).unwrap();
fs::create_dir(&empty).unwrap();
fs::create_dir(&orphan).unwrap();
fs::write(own.join("sf-initial.sfa"), b"own").unwrap();
fs::write(failed.join("sf-initial.sfa"), b"failed").unwrap();
mark_failed(&failed, "previous drainer failed").unwrap();
fs::write(orphan.join("sf-initial.sfa"), b"orphan").unwrap();
fs::write(temp.path().join("top-level.sfa"), b"not a slot").unwrap();
assert_eq!(scan_orphan_slots(temp.path(), "default", &[]), vec![orphan]);
}
#[test]
fn scan_filters_canonical_managed_slot_ranges_only() {
let temp = TempDir::new().unwrap();
for name in [
"default-ingest-0",
"default-ingest-1",
"default-ingest-2",
"default-ingest-02",
"default-ingest-x",
"unmanaged-1",
"default",
] {
let slot = temp.path().join(name);
fs::create_dir(&slot).unwrap();
fs::write(slot.join("sf-initial.sfa"), b"queued").unwrap();
}
let exclusions = vec![QwpWsManagedSlotExclusion::new(
"default-ingest-".to_owned(),
2,
)];
let mut actual = scan_orphan_slots(temp.path(), "unmanaged-1", &exclusions);
actual.sort();
assert_eq!(
actual,
vec![
temp.path().join("default"),
temp.path().join("default-ingest-02"),
temp.path().join("default-ingest-2"),
temp.path().join("default-ingest-x"),
]
);
}
#[test]
fn failed_sentinel_is_written_with_reason() {
let temp = TempDir::new().unwrap();
mark_failed(temp.path(), "connect failed").unwrap();
assert!(has_failed_sentinel(temp.path()));
assert_eq!(
fs::read_to_string(temp.path().join(FAILED_SENTINEL_NAME)).unwrap(),
"connect failed"
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn test_config() -> OrphanDrainerConfig {
let qwp_ws = QwpWsConfig {
sf_max_segment_bytes: ConfigSetting::new_default(256),
sf_max_total_bytes: ConfigSetting::new_default(Some(1024)),
..QwpWsConfig::default()
};
OrphanDrainerConfig::new("127.0.0.1", "1", false, None, &qwp_ws, None)
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn background_pool_close_allows_graceful_finish_before_stop() {
let stop = Arc::new(AtomicBool::new(false));
let observed_stop = Arc::new(AtomicBool::new(true));
let (ran_tx, ran_rx) = std::sync::mpsc::channel();
let worker_stop = Arc::clone(&stop);
let worker_observed_stop = Arc::clone(&observed_stop);
let worker = thread::spawn(move || {
thread::sleep(Duration::from_millis(20));
worker_observed_stop.store(worker_stop.load(Ordering::Acquire), Ordering::Release);
ran_tx.send(()).unwrap();
});
let mut pool = OrphanDrainerPool {
stop,
threads: vec![(worker, Arc::new(TrafficGate::default()))],
};
pool.close_with_timeouts(Duration::from_secs(1), Duration::from_millis(10));
ran_rx.recv_timeout(Duration::from_secs(1)).unwrap();
assert!(!observed_stop.load(Ordering::Acquire));
assert!(pool.stop.load(Ordering::Acquire));
assert!(pool.threads.is_empty());
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn background_pool_close_detaches_worker_after_stop_grace() {
let stop = Arc::new(AtomicBool::new(false));
let (entered_tx, entered_rx) = std::sync::mpsc::channel();
let (stop_seen_tx, stop_seen_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let worker_stop = Arc::clone(&stop);
let worker = thread::spawn(move || {
entered_tx.send(()).unwrap();
while !worker_stop.load(Ordering::Acquire) {
thread::sleep(Duration::from_millis(1));
}
stop_seen_tx.send(()).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(1));
});
entered_rx.recv_timeout(Duration::from_millis(500)).unwrap();
let mut pool = OrphanDrainerPool {
stop,
threads: vec![(worker, Arc::new(TrafficGate::default()))],
};
let started = Instant::now();
pool.close_with_timeouts(Duration::from_millis(10), Duration::from_millis(10));
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_millis(500),
"orphan pool close took {elapsed:?}"
);
assert!(pool.stop.load(Ordering::Acquire));
assert!(pool.threads.is_empty());
stop_seen_rx
.recv_timeout(Duration::from_millis(500))
.unwrap();
release_tx.send(()).unwrap();
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn detached_background_worker_retains_slot_until_exit() {
let temp = TempDir::new().unwrap();
let options = SfaSlotOptions {
sf_dir: temp.path().to_path_buf(),
sender_id: "orphan".to_owned(),
segment_size_bytes: 256,
max_bytes: 1024,
periodic_sync_interval: None,
};
let stop = Arc::new(AtomicBool::new(false));
let (entered_tx, entered_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let (exited_tx, exited_rx) = std::sync::mpsc::channel();
let worker_stop = Arc::clone(&stop);
let worker_options = options.clone();
let worker = thread::spawn(move || {
let queue = SfaSlotQueue::open(worker_options).unwrap();
entered_tx.send(()).unwrap();
while !worker_stop.load(Ordering::Acquire) {
thread::yield_now();
}
release_rx.recv().unwrap();
drop(queue);
exited_tx.send(()).unwrap();
});
entered_rx.recv_timeout(Duration::from_secs(10)).unwrap();
let mut pool = OrphanDrainerPool {
stop,
threads: vec![(worker, Arc::new(TrafficGate::default()))],
};
pool.close_with_timeouts(Duration::from_millis(10), Duration::from_millis(10));
assert!(matches!(
SfaSlotQueue::open(options.clone()),
Err(SfaQueueError::SlotInUse { .. })
));
release_tx.send(()).unwrap();
exited_rx.recv_timeout(Duration::from_secs(10)).unwrap();
let mut reopened = SfaSlotQueue::open(options).unwrap();
reopened.close().unwrap();
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn orphan_reconnect_sleep_observes_stop_request() {
let stop = Arc::new(AtomicBool::new(false));
let worker_stop = Arc::clone(&stop);
let stopper = thread::spawn(move || {
thread::sleep(Duration::from_millis(20));
worker_stop.store(true, Ordering::Release);
});
let started = Instant::now();
assert!(!sleep_before_orphan_reconnect(
None,
Duration::from_secs(5),
Some(&stop)
));
assert!(
started.elapsed() < Duration::from_secs(1),
"orphan reconnect sleep ignored stop request"
);
stopper.join().unwrap();
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn orphan_retry_backoff_doubles_and_caps() {
let max = Duration::from_millis(500);
let mut backoff = Duration::from_millis(100);
backoff = next_orphan_retry_backoff(backoff, max);
assert_eq!(backoff, Duration::from_millis(200));
backoff = next_orphan_retry_backoff(backoff, max);
assert_eq!(backoff, Duration::from_millis(400));
backoff = next_orphan_retry_backoff(backoff, max);
assert_eq!(backoff, max);
assert_eq!(next_orphan_retry_backoff(backoff, max), max);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn orphan_terminal_classification_marks_only_local_resend_required_unrecoverable() {
let local = crate::Error::new(ErrorCode::StoreResendRequired, "torn dictionary");
assert!(terminal_error_is_proven_local_unrecoverable(Some(&local)));
let parse_rejection = QwpWsSenderError {
category: QwpWsErrorCategory::ParseError,
applied_policy: QwpWsErrorPolicy::Terminal,
status: Some(0x05),
message: Some("bad frame".to_owned()),
message_sequence: Some(0),
from_fsn: 0,
to_fsn: 0,
};
let parse = crate::Error::new(ErrorCode::ServerRejection, "bad frame")
.with_qwp_ws_rejection(parse_rejection);
assert!(!terminal_error_is_proven_local_unrecoverable(Some(&parse)));
assert!(!terminal_error_is_proven_local_unrecoverable(None));
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
fn slot_options(sf_dir: &Path, sender_id: &str) -> SfaSlotOptions {
SfaSlotOptions {
sf_dir: sf_dir.to_path_buf(),
sender_id: sender_id.to_owned(),
segment_size_bytes: 256,
max_bytes: 1024,
periodic_sync_interval: None,
}
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
fn create_queued_orphan(sf_dir: &Path, sender_id: &str) -> PathBuf {
let slot_dir = sf_dir.join(sender_id);
let mut queue = SfaSlotQueue::open(slot_options(sf_dir, sender_id)).unwrap();
PublicationLog::try_publish(&mut queue, b"orphaned frame").unwrap();
queue.close().unwrap();
slot_dir
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn manual_drainer_consumes_already_drained_slot_without_network() {
let temp = TempDir::new().unwrap();
let slot_dir = temp.path().join("orphan");
fs::create_dir(&slot_dir).unwrap();
let mut drainers =
ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, test_config()).unwrap();
assert!(drainers.drive_once());
assert!(!has_failed_sentinel(&slot_dir));
assert!(!drainers.drive_once());
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn orphan_adoption_inherits_periodic_sync_interval_from_config() {
let temp = TempDir::new().unwrap();
let sf_dir = temp.path().join("sf-root");
let slot_dir = sf_dir.join("orphan");
{
let mut queue = SfaSlotQueue::open(slot_options(&sf_dir, "orphan")).unwrap();
PublicationLog::try_publish(&mut queue, b"orphaned").unwrap();
queue.close().unwrap();
}
let builder = crate::ingress::SenderBuilder::from_conf(format!(
"ws::addr=127.0.0.1:1;sf_dir={};sender_id=primary;drain_orphans=on;\
sf_max_segment_bytes=256;sf_max_total_bytes=1024;\
sf_durability=periodic;sf_sync_interval_millis=123;",
sf_dir.display()
))
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
let config = OrphanDrainerConfig::new("127.0.0.1", "1", false, None, qwp_ws, None);
let queue_options = config.queue_options(slot_dir).unwrap();
assert_eq!(
queue_options.periodic_sync_interval,
Some(Duration::from_millis(123))
);
let mut adopted = SfaSlotQueue::open_replay_only_existing(queue_options).unwrap();
let sync_step = PublicationLog::take_storage_maintenance_step(&mut adopted, false)
.unwrap()
.expect("adopted periodic queue must schedule a sync step");
assert!(matches!(sync_step, SfaStorageStep::SyncPublished(_)));
let sync_result = sync_step.perform().unwrap();
PublicationLog::finish_storage_maintenance(&mut adopted, sync_result, false).unwrap();
PublicationLog::complete_storage_maintenance(&mut adopted).unwrap();
adopted.close().unwrap();
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn manual_drainer_honors_failed_sentinel_created_after_scan() {
let temp = TempDir::new().unwrap();
let slot_dir = temp.path().join("orphan");
fs::create_dir(&slot_dir).unwrap();
let mut drainers =
ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, test_config()).unwrap();
mark_failed(&slot_dir, "created after scan").unwrap();
assert!(drainers.drive_once());
assert!(!drainers.drive_once());
assert_eq!(
fs::read_to_string(slot_dir.join(FAILED_SENTINEL_NAME)).unwrap(),
"created after scan"
);
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn manual_drainer_skips_locked_slot_without_failed_sentinel() {
let temp = TempDir::new().unwrap();
let sf_dir = temp.path().join("sf-root");
let slot_dir = sf_dir.join("locked");
let _held = SfaSlotQueue::open(slot_options(&sf_dir, "locked")).unwrap();
assert!(
matches!(
OrphanDrainer::open(slot_dir.clone(), &test_config()),
OrphanOpenOutcome::Locked
),
"held slot must surface as Locked before already-drained handling"
);
let mut drainers =
ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, test_config()).unwrap();
assert!(drainers.drive_once());
assert!(!has_failed_sentinel(&slot_dir));
assert!(!drainers.drive_once());
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn manual_drainer_does_not_poison_orphan_after_connect_failure() {
let temp = TempDir::new().unwrap();
let sf_dir = temp.path().join("sf-root");
let slot_dir = sf_dir.join("orphan");
{
let mut queue = SfaSlotQueue::open(slot_options(&sf_dir, "orphan")).unwrap();
PublicationLog::try_publish(&mut queue, b"orphaned frame").unwrap();
queue.close().unwrap();
}
assert_eq!(
scan_orphan_slots(&sf_dir, "primary", &[]),
vec![slot_dir.clone()]
);
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener);
let mut config = test_config();
config.port = port.to_string();
let mut drainers = ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, config).unwrap();
assert!(drainers.drive_once());
assert!(
!has_failed_sentinel(&slot_dir),
"transient connect failure must leave orphan slot recoverable"
);
assert!(slot_dir.join(LAST_ERROR_NAME).exists());
assert_eq!(scan_orphan_slots(&sf_dir, "primary", &[]), vec![slot_dir]);
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn manual_drainer_gates_connect_retries_so_caller_can_park() {
let temp = TempDir::new().unwrap();
let sf_dir = temp.path().join("sf-root");
let slot_dir = create_queued_orphan(&sf_dir, "orphan");
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener);
let mut config = test_config();
config.port = port.to_string();
config.qwp_ws.reconnect_initial_backoff =
ConfigSetting::new_default(Duration::from_secs(3600));
let mut drainers = ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, config).unwrap();
assert!(
drainers.drive_once(),
"the first connect attempt is a unit of work"
);
for _ in 0..8 {
assert!(
!drainers.drive_once(),
"a gated slot must let the caller park instead of connect-storming"
);
}
assert!(!has_failed_sentinel(&slot_dir));
assert!(slot_dir.join(LAST_ERROR_NAME).exists());
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn manual_drainer_retries_connect_after_backoff_elapses() {
let temp = TempDir::new().unwrap();
let sf_dir = temp.path().join("sf-root");
let slot_dir = create_queued_orphan(&sf_dir, "orphan");
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener);
let mut config = test_config();
config.port = port.to_string();
config.qwp_ws.reconnect_initial_backoff =
ConfigSetting::new_default(Duration::from_millis(10));
let mut drainers = ManualOrphanDrainers::new(vec![slot_dir], 1, config).unwrap();
assert!(drainers.drive_once());
thread::sleep(Duration::from_millis(50));
assert!(
drainers.drive_once(),
"an expired backoff gate must allow the next connect attempt"
);
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn manual_drainer_connect_is_invisible_to_foreground_events() {
let temp = TempDir::new().unwrap();
let sf_dir = temp.path().join("sf-root");
let slot_dir = create_queued_orphan(&sf_dir, "orphan");
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener);
let source = Arc::new(crate::ingress::conn_events::ConnectionEventSource::disabled());
let mut config = test_config();
config.port = port.to_string();
config.qwp_ws.conn_events = Some(Arc::clone(&source));
let mut drainers = ManualOrphanDrainers::new(vec![slot_dir], 1, config).unwrap();
assert!(drainers.drive_once());
assert_eq!(
source.next_attempt(),
1,
"manual orphan connects must not consume foreground attempt numbers"
);
}
#[cfg(all(feature = "sync-sender-qwp-ws", any(unix, windows)))]
#[test]
fn manual_drainer_retries_once_after_sanitizing_sealed_residue() {
let temp = TempDir::new().unwrap();
let slot_dir = temp.path().join("orphan");
fs::create_dir(&slot_dir).unwrap();
let sealed_path = spare_segment_path(&slot_dir, 0);
let active_path = spare_segment_path(&slot_dir, 1);
let mut sealed = SfaSegment::create_new_manifested(&sealed_path, 0, 256, 0).unwrap();
sealed.try_append(b"orphaned").unwrap();
drop(sealed);
let append_offset = scan_file(&sealed_path).unwrap().append_offset as usize;
let mut bytes = fs::read(&sealed_path).unwrap();
bytes[append_offset] = 0xa5;
fs::write(&sealed_path, bytes).unwrap();
drop(SfaSegment::create_new_manifested(&active_path, 1, 256, 0).unwrap());
drop(SfManifest::create(&slot_dir, 0, 1).unwrap());
let mut watermark = SfaAckWatermark::open(&slot_dir).unwrap();
watermark.write(-1).unwrap();
watermark.sync_data().unwrap();
drop(watermark);
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener);
let mut config = test_config();
config.port = port.to_string();
let mut drainers = ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, config).unwrap();
assert!(drainers.drive_once());
assert_eq!(scan_file(&sealed_path).unwrap().torn_tail_bytes, 0);
assert!(
!has_failed_sentinel(&slot_dir),
"the first sanitized-residue incident must be retried, not poisoned"
);
assert!(slot_dir.join(LAST_ERROR_NAME).exists());
}
#[cfg(all(feature = "sync-sender-qwp-ws", unix))]
#[test]
fn manual_drainer_marks_recovery_failure_failed_without_network() {
let temp = TempDir::new().unwrap();
let slot_dir = temp.path().join("orphan");
fs::create_dir(&slot_dir).unwrap();
let segment_path = initial_segment_path(&slot_dir);
let mut bytes = vec![0u8; 64];
bytes[..4].copy_from_slice(&0xdead_beefu32.to_le_bytes());
fs::write(&segment_path, bytes).unwrap();
let mut drainers =
ManualOrphanDrainers::new(vec![slot_dir.clone()], 1, test_config()).unwrap();
assert!(drainers.drive_once());
assert!(has_failed_sentinel(&slot_dir));
assert!(
!segment_path.exists(),
"manifest-less corruption must leave the .sfa scan"
);
assert!(
segment_path.with_extension("sfa.corrupt").exists(),
"quarantine must preserve the corrupt bytes for diagnosis"
);
assert!(!drainers.drive_once());
}
}