1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
// Project: scalo
// File: src/tiered_sink/mod.rs
// Purpose: Tiered sink with disk spillover for resilient message delivery
// Language: Rust
//
// License: Apache-2.0
// Copyright: (c) 2026 HYPERI PTY LIMITED
//! Tiered sink with automatic disk spillover for resilient message delivery.
//!
//! This module provides a wrapper around any async sink (Kafka, S3, HTTP, etc.)
//! that automatically spills messages to disk when the primary sink is unavailable
//! or backpressuring, then drains them back when the sink recovers.
//!
//! ## Design
//!
//! ```text
//! ┌─────────────────────────────────────┐
//! │ TieredSink │
//! │ │
//! Message ────────►│ try_send() to primary sink │
//! │ │ │
//! │ ▼ │
//! │ ┌─────────┐ │
//! │ │ Success │──► Done (hot path) │
//! │ └────┬────┘ │
//! │ │ Err(Full/Unavailable) │
//! │ ▼ │
//! │ ┌─────────┐ │
//! │ │ Spool │──► Disk (cold path) │
//! │ └────┬────┘ │
//! │ │ │
//! │ Background drain task │
//! │ (when primary recovers) │
//! └─────────────────────────────────────┘
//! ```
//!
//! ## Features
//!
//! - **Hot path first**: Always tries primary sink with timeout
//! - **Automatic spillover**: Writes to disk only when primary fails
//! - **Circuit breaker**: Avoids hammering a dead sink
//! - **Background drain**: Recovers spooled messages when sink is healthy
//! - **Configurable ordering**: Interleaved (default) or strict FIFO
//! - **Multiple compression codecs**: Zstd (default, level 1), LZ4, Snappy, None
//!
//! ## Example
//!
//! TieredSink wraps any [`TransportSender`](crate::transport::TransportSender) --
//! the same senders the transport factory produces. Records go straight to the
//! sender on the happy path (no encode); only when the downstream fails are they
//! serialised and spilled to disk, then drained back on recovery.
//!
//! ```rust,ignore
//! use scalo::tiered_sink::{TieredSink, TieredSinkConfig};
//! use scalo::transport::AnySender;
//!
//! let sender = AnySender::from_config("transport.output").await?;
//! let config = TieredSinkConfig::new("/var/spool/myapp.queue");
//! let tiered = TieredSink::new(sender, config).await?;
//!
//! // Automatically spills to disk if the downstream is down, drains on recovery.
//! tiered.send(&record).await?;
//! ```
pub use ;
pub use CompressionCodec;
pub use ;
pub use TieredSinkError;
pub use TieredSink;
/// Result type for tiered sink operations.
pub type Result<T> = Result;