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
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
//! Bulk multi-step processing on top of [`taquba_workflow`].
//!
//! `taquba-bulk` runs one [`Pipeline`] over many input items in parallel,
//! inside a single process, with per-item memoization, retry, streamed
//! output, and a rolled-up cost report. It is the per-batch orchestrator for
//! workloads that fan out 10-1000x per run: bulk LLM jobs (classify, look up,
//! draft, check, refine over thousands of tickets), document/OCR pipelines,
//! data enrichment, parameter sweeps. The pipeline contract is workload
//! agnostic.
//!
//! # Execution model: one item, one run, one step
//!
//! Each input item becomes one [`taquba_workflow`] run whose single step
//! invokes [`Pipeline::run`]. The pipeline's own logical steps live inside
//! that method as [`BulkCtx::memoized`] or
//! [`BulkCtx::memoized_by_content`] calls. Taquba delivers at-least-once, so
//! a step may run again if its lease expires before it acks; memoization makes
//! that replay cheap, because each completed logical step returns its cached
//! result instead of repeating a paid call. A pipeline error retries with
//! backoff and then dead-letters the item (terminating it failed); the rest of
//! the batch is unaffected.
//!
//! [`BulkCtx::memoized`] is [`taquba_workflow`]'s per-step memo store
//! applied at a finer granularity: the item's single step holds one memo
//! entry per logical phase, so the phases of [`Pipeline::run`] resume
//! individually even though the workflow sees one step.
//!
//! # Content-addressed memoization
//!
//! Use [`BulkCtx::memoized_by_content`] when the natural memo key is a
//! serialized input value rather than a caller-supplied string:
//!
//! ```ignore
//! #[derive(serde::Serialize)]
//! struct LookupKey<'a> {
//! operation: &'static str,
//! query: &'a str,
//! }
//!
//! let key = LookupKey {
//! operation: "lookup",
//! query: &ctx.input.body,
//! };
//! let response = ctx
//! .memoized_by_content(&key, async {
//! Ok::<_, StepError>(lookup(&ctx.input.body).await?)
//! })
//! .await?;
//! ```
//!
//! The helper serializes the key as MessagePack, hashes it with SHA-256,
//! and uses the digest inside the item's existing workflow memo namespace.
//! The entry remains scoped to one item run; this is not a cross-item cache.
//! Include an operation name in the serialized key when multiple logical
//! operations may receive the same input shape.
//!
//! # Single process, remote work per step
//!
//! The orchestrator is single-process by design: SlateDB allows one writer
//! per store, so all producers and workers for a batch share one
//! `Arc<Queue>` (see the Taquba docs). That is not a throughput ceiling for
//! bulk work. Each step's expensive operation is a call to a remote service
//! (an LLM API, an OCR service), so the process is I/O-bound and one host
//! sustains hundreds of concurrent items. The remote call runs elsewhere and
//! its response is memoized on return.
//!
//! # Quick start
//!
//! ```no_run
//! use std::sync::Arc;
//! use serde::{Deserialize, Serialize};
//! use taquba::{Queue, object_store::memory::InMemory};
//! use taquba_bulk::{Bulk, BulkCtx, CostReport, Pipeline, StepError};
//!
//! #[derive(Serialize, Deserialize)]
//! struct Ticket { id: String, body: String }
//!
//! #[derive(Serialize, Deserialize)]
//! struct Processed { id: String, classification: String }
//!
//! struct TicketPipeline;
//!
//! impl Pipeline for TicketPipeline {
//! type Input = Ticket;
//! type Output = Processed;
//! type Error = StepError;
//!
//! async fn run(&self, ctx: &BulkCtx<Ticket>) -> Result<Processed, StepError> {
//! let classification = ctx
//! .memoized_with_cached_cost("classify", async {
//! let cost = CostReport::new();
//! cost.record("llm_calls", 1.0);
//! Ok::<_, StepError>(("billing".to_string(), cost))
//! })
//! .await?;
//! Ok(Processed { id: ctx.input.id.clone(), classification })
//! }
//! }
//!
//! # async fn run() -> Result<(), Box<dyn std::error::Error>> {
//! let store = Arc::new(InMemory::new());
//! let queue = Arc::new(Queue::open(store.clone(), "db").await?);
//!
//! let bulk = Bulk::builder(queue, store, TicketPipeline)
//! .key_fn(|t| t.id.clone())
//! .max_concurrent(200)
//! .build();
//!
//! let inputs = vec![
//! Ticket { id: "t1".into(), body: "help".into() },
//! Ticket { id: "t2".into(), body: "refund".into() },
//! ];
//! let report = bulk.run(inputs).await?;
//! println!("{}/{} succeeded", report.succeeded, report.total);
//! # Ok(()) }
//! ```
//!
//! # Cost tracking
//!
//! Pipelines report arbitrary named metrics via [`BulkCtx::record_cost`]
//! (token counts, paid-API units, compute-seconds, dollars). Per-item totals
//! roll up into [`ProgressSnapshot::cost`] and [`BulkReport::cost`], so the
//! batch cost is visible live and in the final report. See [`CostReport`].
//! When counters are produced inside a memoized closure, return `(value,
//! cost)` from [`BulkCtx::memoized_with_cached_cost`] or
//! [`BulkCtx::memoized_by_content_with_cached_cost`] so the same counters are
//! recorded on a cache hit.
//!
//! # Failure policy
//!
//! Per-item failures are recorded, not fatal: each failed item is written to
//! the output sink with its error and its run id is collected on
//! [`BulkReport::failed_run_ids`]. Set [`BulkBuilder::fail_threshold`] to
//! turn the whole run into an [`Error::FailureThresholdExceeded`] when the
//! share of failures crosses a percentage, so a silent mass failure
//! surfaces.
//!
//! # Replay
//!
//! Because memo entries are retained, re-submitting a failed item's input
//! with the same run id resumes from its last cached step rather than
//! recomputing. [`BulkReport::failed_run_ids`] is the set to replay.
//!
//! # Input and output
//!
//! Line-delimited JSON: [`read_jsonl`] decodes inputs and
//! [`JsonlSink`] writes one result record per line. Both sides are traits
//! ([`OutputSink`]), so other formats can be added without touching the
//! runner. The default sink is [`NullSink`], for pipelines whose results are
//! side effects.
pub use ;
pub use CostReport;
pub use ;
pub use ;
pub use ;
pub use ;
/// Re-exported from [`taquba_workflow`]: the error type a [`Pipeline`]
/// returns, with [`StepError::transient`] / [`StepError::permanent`]
/// controlling retry versus immediate dead-letter.
pub use ;