Skip to main content

faucet_core/
idempotency.rs

1//! Exactly-once / idempotent delivery primitives.
2//!
3//! The pipeline issues a monotonic **commit token** for every page that carries
4//! a bookmark. The token is persisted in the [`StateStore`](crate::state::StateStore)
5//! value next to the bookmark and committed inside the sink's own transaction,
6//! so a crash between "sink durably wrote" and "state persisted" is resolved on
7//! resume by skipping pages the sink already committed. See
8//! `docs/superpowers/specs/2026-06-09-exactly-once-delivery-design.md`.
9
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12
13/// Delivery guarantee for a pipeline run.
14#[derive(
15    Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, schemars::JsonSchema,
16)]
17#[serde(rename_all = "snake_case")]
18pub enum DeliveryMode {
19    /// Today's behaviour: a page may be re-delivered after a crash between the
20    /// sink write and the bookmark persist. Downstream must tolerate duplicates.
21    #[default]
22    AtLeastOnce,
23    /// The sink durably records a per-page commit token atomically with the
24    /// data; on resume the pipeline skips already-committed pages. Requires a
25    /// state store, an idempotent sink, and a deterministic-replay source.
26    ExactlyOnce,
27}
28
29/// How faithfully a [`Source`](crate::Source) **replays** its record stream
30/// when resumed from a bookmark.
31///
32/// This is the source-side capability the effectively-once *atomic-watermark*
33/// mechanism depends on: after a crash the pipeline re-anchors the source at a
34/// persisted position, and correctness requires that nothing before that
35/// position is re-emitted and nothing after it is skipped.
36#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
37#[serde(rename_all = "snake_case")]
38pub enum ReplayGuarantee {
39    /// Resuming from a bookmark may replay a *different* record stream
40    /// (query-based sources whose upstream can mutate, sources without
41    /// per-page bookmarks). The default.
42    #[default]
43    NonDeterministic,
44    /// The source emits a complete resume position (bookmark) on **every**
45    /// page, and resuming from any such bookmark continues the record stream
46    /// at exactly that position — no record before the bookmark is re-emitted
47    /// and none after it is skipped (immutable-log sources: CDC WAL/binlog/
48    /// change streams, Kafka partitions).
49    Deterministic,
50}
51
52/// The strongest delivery guarantee a [`Sink`](crate::Sink) can uphold.
53#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
54#[serde(rename_all = "snake_case")]
55pub enum SinkGuarantee {
56    /// Plain writes: a replayed page is written again. The default.
57    #[default]
58    AtLeastOnce,
59    /// The sink can dedup by key (`write_mode: upsert` with a configured
60    /// `key`): re-applying a record with the same key converges instead of
61    /// duplicating.
62    KeyedUpsert,
63    /// The sink can commit a page's rows **and** a commit token in one atomic
64    /// transaction ([`Sink::write_batch_idempotent`](crate::Sink)).
65    AtomicWatermark,
66}
67
68/// The mechanism through which a pipeline achieves effectively-once delivery.
69#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
70#[serde(rename_all = "snake_case")]
71pub enum EffectivelyOnceMechanism {
72    /// Deterministic-replay source + sink that commits data and a per-page
73    /// commit token atomically; on resume already-committed pages are skipped
74    /// (or the stream is re-anchored at the sink's recorded position).
75    AtomicWatermark,
76    /// The sink dedups by key (`write_mode: upsert`); replayed records
77    /// converge on the same keyed row. Works with any source.
78    KeyedUpsert,
79}
80
81/// The end-to-end guarantee a *pipeline* provides for a given
82/// source × sink × config combination.
83///
84/// Deliberately no `ExactlyOnce` variant — distributed-consensus exactly-once
85/// is not achievable here; effectively-once (idempotent at-least-once: each
86/// record is *observably applied* once) is the ceiling, and
87/// `delivery: exactly_once` in config is precisely documented as requesting it.
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "snake_case", tag = "guarantee", content = "via")]
90pub enum DeliveryGuarantee {
91    /// A crash between the sink write and the bookmark persist may re-deliver
92    /// a page. Downstream must tolerate duplicates.
93    AtLeastOnce,
94    /// Idempotent at-least-once: each record is observably applied once.
95    EffectivelyOnce(EffectivelyOnceMechanism),
96}
97
98impl std::fmt::Display for DeliveryGuarantee {
99    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
100        match self {
101            Self::AtLeastOnce => write!(f, "at-least-once"),
102            Self::EffectivelyOnce(EffectivelyOnceMechanism::AtomicWatermark) => {
103                write!(f, "effectively-once (atomic watermark)")
104            }
105            Self::EffectivelyOnce(EffectivelyOnceMechanism::KeyedUpsert) => {
106                write!(f, "effectively-once (keyed upsert)")
107            }
108        }
109    }
110}
111
112/// Inputs to [`derive_delivery_guarantee`] — the facts about a concrete
113/// source × sink × config combination the derivation keys off.
114#[derive(Debug, Clone, Copy, Default)]
115pub struct GuaranteeInputs {
116    /// The source's replay capability.
117    pub replay: ReplayGuarantee,
118    /// Whether the sink commits data + token atomically
119    /// (`Sink::supports_idempotent_writes`).
120    pub sink_atomic: bool,
121    /// Whether the sink is *configured* to dedup by key — `write_mode: upsert`
122    /// (or `delete`) with a non-empty `key` (`Sink::dedups_by_key`).
123    pub keyed_upsert_configured: bool,
124    /// Whether a durable (non-memory) state store is configured. The
125    /// atomic-watermark mechanism persists its cross-restart sequence here.
126    pub durable_state: bool,
127    /// Whether a DLQ is configured (incompatible with the atomic-watermark
128    /// mechanism in this version).
129    pub dlq: bool,
130}
131
132/// Derive the end-to-end [`DeliveryGuarantee`] a pipeline actually provides.
133///
134/// Preference order: the atomic-watermark mechanism (strongest bookkeeping,
135/// no keyed-schema requirement) when the topology supports it, then keyed
136/// upsert, then at-least-once. A sink that is both atomic and keyed reports
137/// atomic-watermark when the source replays deterministically, and falls back
138/// to keyed upsert otherwise.
139pub fn derive_delivery_guarantee(i: &GuaranteeInputs) -> DeliveryGuarantee {
140    if i.sink_atomic && i.replay == ReplayGuarantee::Deterministic && i.durable_state && !i.dlq {
141        return DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::AtomicWatermark);
142    }
143    if i.keyed_upsert_configured {
144        return DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::KeyedUpsert);
145    }
146    DeliveryGuarantee::AtLeastOnce
147}
148
149/// Reserved key marking the exactly-once state wrapper object.
150const EO_MARKER: &str = "__faucet_eo";
151const EO_BOOKMARK: &str = "bookmark";
152const EO_SEQ: &str = "seq";
153
154/// Width of the zero-padded decimal token. `u64::MAX` is 20 digits, so 20 makes
155/// lexicographic order match numeric order for the full `u64` range.
156const TOKEN_WIDTH: usize = 20;
157
158/// Separator between the numeric sequence and the embedded resume bookmark in
159/// a commit token. The prefix before it is always the fixed-width sequence.
160const TOKEN_BOOKMARK_SEP: char = '#';
161
162/// Render a page sequence as a fixed-width, lexicographically-ordered token.
163pub fn format_token(seq: u64) -> String {
164    format!("{seq:0TOKEN_WIDTH$}")
165}
166
167/// Render a commit token that carries the page's **resume bookmark** alongside
168/// the sequence: `"{seq:020}#{bookmark-json}"`.
169///
170/// Sinks store the token opaquely, so the committed watermark doubles as a
171/// durable record of *where the stream stood* when the page committed. On
172/// resume the pipeline recovers that position from the sink
173/// ([`parse_token_parts`]) and re-anchors the source there — closing the
174/// crash window between "sink durably committed" and "state store persisted"
175/// without requiring the source to replay identical page boundaries.
176pub fn format_token_with_bookmark(seq: u64, bookmark: Option<&Value>) -> String {
177    match bookmark {
178        Some(bm) => format!("{seq:0TOKEN_WIDTH$}{TOKEN_BOOKMARK_SEP}{bm}"),
179        None => format_token(seq),
180    }
181}
182
183/// Parse the numeric sequence from a token produced by [`format_token`] or
184/// [`format_token_with_bookmark`]. Returns `None` on garbage.
185pub fn parse_token(s: &str) -> Option<u64> {
186    let seq = match s.split_once(TOKEN_BOOKMARK_SEP) {
187        Some((prefix, _)) => prefix,
188        None => s,
189    };
190    seq.trim().parse::<u64>().ok()
191}
192
193/// Parse a stored commit token into `(seq, embedded_bookmark)`.
194///
195/// Tokens written before bookmarks were embedded (bare `format_token` output)
196/// parse with `bookmark = None`. A bookmark suffix that is not valid JSON also
197/// yields `None` for the bookmark — the sequence alone still drives the
198/// skip-on-resume path.
199pub fn parse_token_parts(s: &str) -> Option<(u64, Option<Value>)> {
200    match s.split_once(TOKEN_BOOKMARK_SEP) {
201        Some((prefix, suffix)) => {
202            let seq = prefix.trim().parse::<u64>().ok()?;
203            Some((seq, serde_json::from_str(suffix).ok()))
204        }
205        None => Some((s.trim().parse::<u64>().ok()?, None)),
206    }
207}
208
209/// Wrap a bookmark + sequence into the exactly-once state value.
210pub fn wrap_state(bookmark: Option<&Value>, seq: u64) -> Value {
211    serde_json::json!({
212        EO_MARKER: 1,
213        EO_BOOKMARK: bookmark.cloned().unwrap_or(Value::Null),
214        EO_SEQ: seq,
215    })
216}
217
218/// Whether `value` is the exactly-once wrapper (`{"__faucet_eo": 1, …}`).
219pub fn is_eo_envelope(value: &Value) -> bool {
220    value.get(EO_MARKER).and_then(Value::as_u64) == Some(1)
221}
222
223/// Unwrap a stored state value into `(bookmark, seq)`. A versioned envelope
224/// (#736) is looked through first.
225///
226/// A value that is the exactly-once wrapper object unwraps to its inner
227/// bookmark + seq. Anything else is treated as a legacy/at-least-once **bare
228/// bookmark** with `seq = 0` — so switching an existing pipeline to
229/// `exactly_once` resumes cleanly (the sink's own watermark is authoritative).
230pub fn unwrap_state(value: &Value) -> (Option<Value>, u64) {
231    let peeled = crate::state_version::peel_versioned(value);
232    let value = &peeled;
233    if let Value::Object(map) = value
234        && map.get(EO_MARKER).and_then(Value::as_u64) == Some(1)
235    {
236        let bookmark = match map.get(EO_BOOKMARK) {
237            None | Some(Value::Null) => None,
238            Some(v) => Some(v.clone()),
239        };
240        let seq = map.get(EO_SEQ).and_then(Value::as_u64).unwrap_or(0);
241        return (bookmark, seq);
242    }
243    // Legacy bare bookmark.
244    let bookmark = if value.is_null() {
245        None
246    } else {
247        Some(value.clone())
248    };
249    (bookmark, 0)
250}
251
252/// Canonical watermark table the SQL sinks UPSERT the commit token into.
253pub const COMMIT_TOKEN_TABLE: &str = "_faucet_commit_token";
254/// Watermark column holding the pipeline state-key (`{name}::{row_id}`).
255pub const COMMIT_TOKEN_SCOPE_COL: &str = "scope";
256/// Watermark column holding the latest committed token.
257pub const COMMIT_TOKEN_TOKEN_COL: &str = "token";
258
259/// Iceberg snapshot summary property names.
260pub const ICEBERG_SCOPE_PROP: &str = "faucet.commit-scope";
261pub const ICEBERG_TOKEN_PROP: &str = "faucet.commit-token";
262
263/// Suffix appended to a destination's name to form the staging relation an
264/// in-flight `write_mode: overwrite` run writes into.
265///
266/// Reserved: the overwrite lifecycle creates, truncates, and **drops** whatever
267/// carries this suffix, so a real destination named `<x>__faucet_ovw` would be
268/// destroyed by `abort_overwrite`. Every sink derives its staging name from this
269/// one constant so the reserved set stays enumerable by a collision check or an
270/// orphan sweep (#654 M12).
271pub const OVERWRITE_STAGING_SUFFIX: &str = "__faucet_ovw";
272/// Suffix of the transient relation the *outgoing* destination is renamed to
273/// mid-swap, on backends whose only atomic publish is a rename (MySQL, whose
274/// DDL auto-commits so a transaction cannot span the swap).
275///
276/// Extends [`OVERWRITE_STAGING_SUFFIX`], so a sweep matching the base suffix as
277/// a prefix catches this one too.
278pub const OVERWRITE_STAGING_OLD_SUFFIX: &str = "__faucet_ovw_old";
279
280/// Fit a watermark scope into a length-capped, indexable key column.
281///
282/// The scope is the pipeline state key — `{name}::{row}` for a root, plus
283/// `::{parent_record_key}` for a child — and a child's key comes from *record
284/// data*, so it has no length bound. The SQL sinks store it as a PRIMARY KEY, and
285/// a key column cannot be unbounded (MySQL's index limit, SQL Server's 900-byte
286/// key budget), so an over-long scope either errors or — under a non-strict MySQL
287/// `sql_mode` — is **truncated**, silently collapsing two distinct rows onto one
288/// watermark so one row's committed token suppresses the other's pages (#456 L1).
289///
290/// Scopes at or under `max` are returned verbatim, so every watermark written
291/// before this existed still resolves. A longer one is replaced by a
292/// deterministic, collision-resistant digest form (`__h:<64-hex>`), which is
293/// stable across restarts — the only property the watermark needs.
294/// Length of the `__h:` + 16 hex-digit suffix appended to a shortened scope.
295const SCOPE_DIGEST_LEN: usize = 4 + 16;
296
297pub fn scope_key(scope: &str, max: usize) -> String {
298    if scope.len() <= max {
299        return scope.to_owned();
300    }
301    // FNV-1a rather than a crypto hash: `sha2` is an optional dependency of this
302    // crate (masking / transform-hash / encryption) and this module is always
303    // compiled. The requirement is determinism, not preimage resistance — the same
304    // choice the backfill progress marker makes.
305    let mut h: u64 = 0xcbf2_9ce4_8422_2325;
306    for b in scope.as_bytes() {
307        h ^= u64::from(*b);
308        h = h.wrapping_mul(0x1000_0000_01b3);
309    }
310    // Keep as much of the readable head as fits, so an operator inspecting the
311    // watermark table can still tell which pipeline a row belongs to. Truncate on
312    // a char boundary — a scope may hold non-ASCII.
313    let room = max.saturating_sub(SCOPE_DIGEST_LEN);
314    let mut head = 0usize;
315    for (i, _) in scope.char_indices() {
316        if i > room {
317            break;
318        }
319        head = i;
320    }
321    let shortened = format!("{}__h:{h:016x}", &scope[..head]);
322    tracing::debug!(
323        scope_len = scope.len(),
324        max,
325        "exactly-once scope exceeds the sink's key-column width; shortening with a digest"
326    );
327    shortened
328}
329
330#[cfg(test)]
331mod tests {
332    use super::*;
333    use serde_json::json;
334
335    #[test]
336    fn token_round_trips_and_orders_lexicographically() {
337        assert_eq!(format_token(42).len(), TOKEN_WIDTH);
338        assert_eq!(parse_token(&format_token(42)), Some(42));
339        assert_eq!(parse_token(&format_token(0)), Some(0));
340        assert_eq!(parse_token(&format_token(u64::MAX)), Some(u64::MAX));
341        assert!(format_token(9) < format_token(10));
342        assert!(format_token(2) < format_token(1000));
343    }
344
345    #[test]
346    fn parse_token_rejects_garbage() {
347        assert_eq!(parse_token("abc"), None);
348        assert_eq!(parse_token(""), None);
349    }
350
351    #[test]
352    fn wrap_then_unwrap_preserves_bookmark_and_seq() {
353        let bm = json!({"lsn": "0/16B2D58"});
354        let wrapped = wrap_state(Some(&bm), 7);
355        let (got_bm, got_seq) = unwrap_state(&wrapped);
356        assert_eq!(got_bm, Some(bm));
357        assert_eq!(got_seq, 7);
358    }
359
360    #[test]
361    fn wrap_none_bookmark_unwraps_to_none() {
362        let wrapped = wrap_state(None, 3);
363        let (got_bm, got_seq) = unwrap_state(&wrapped);
364        assert_eq!(got_bm, None);
365        assert_eq!(got_seq, 3);
366    }
367
368    #[test]
369    fn legacy_bare_bookmark_unwraps_with_seq_zero() {
370        let (bm, seq) = unwrap_state(&json!("2024-12-01"));
371        assert_eq!(bm, Some(json!("2024-12-01")));
372        assert_eq!(seq, 0);
373        let (bm2, seq2) = unwrap_state(&json!({"updated_at": "2024-12-01"}));
374        assert_eq!(bm2, Some(json!({"updated_at": "2024-12-01"})));
375        assert_eq!(seq2, 0);
376    }
377
378    #[test]
379    fn object_with_non_sentinel_marker_is_treated_as_bare_bookmark() {
380        // A legacy/user object that merely contains the key must NOT be misread
381        // as an EO wrapper — only the typed sentinel `1` counts.
382        let v = json!({"__faucet_eo": null, "offset": 500});
383        let (bm, seq) = unwrap_state(&v);
384        assert_eq!(bm, Some(v));
385        assert_eq!(seq, 0);
386    }
387
388    #[test]
389    fn null_value_unwraps_to_none_seq_zero() {
390        let (bm, seq) = unwrap_state(&json!(null));
391        assert_eq!(bm, None);
392        assert_eq!(seq, 0);
393    }
394
395    #[test]
396    fn token_with_bookmark_round_trips() {
397        let bm = json!({"partition_offsets": [{"topic": "t", "partition": 0, "offset": 42}]});
398        let token = format_token_with_bookmark(7, Some(&bm));
399        assert!(token.starts_with(&format_token(7)));
400        assert_eq!(parse_token(&token), Some(7));
401        let (seq, parsed_bm) = parse_token_parts(&token).unwrap();
402        assert_eq!(seq, 7);
403        assert_eq!(parsed_bm, Some(bm));
404    }
405
406    #[test]
407    fn token_with_no_bookmark_is_bare_and_back_compatible() {
408        assert_eq!(format_token_with_bookmark(3, None), format_token(3));
409        let (seq, bm) = parse_token_parts(&format_token(3)).unwrap();
410        assert_eq!((seq, bm), (3, None));
411    }
412
413    #[test]
414    fn token_with_bookmark_orders_lexicographically_on_prefix() {
415        // The fixed-width numeric prefix keeps lexicographic order meaningful
416        // even with an embedded bookmark (kafka side-topic folding compares
417        // parsed sequences, but SQL MAX() naturally works too).
418        let a = format_token_with_bookmark(9, Some(&json!({"o": 1})));
419        let b = format_token_with_bookmark(10, Some(&json!({"o": 2})));
420        assert!(a < b);
421    }
422
423    #[test]
424    fn parse_token_parts_tolerates_garbage() {
425        assert_eq!(parse_token_parts("abc"), None);
426        assert_eq!(parse_token_parts(""), None);
427        // Bad JSON suffix: sequence survives, bookmark is dropped.
428        let (seq, bm) = parse_token_parts("00000000000000000005#{not json").unwrap();
429        assert_eq!((seq, bm), (5, None));
430        // parse_token ignores the suffix entirely.
431        assert_eq!(parse_token("00000000000000000005#{not json"), Some(5));
432    }
433
434    #[test]
435    fn derive_guarantee_prefers_atomic_then_keyed_then_at_least_once() {
436        use ReplayGuarantee::*;
437        let base = GuaranteeInputs {
438            replay: Deterministic,
439            sink_atomic: true,
440            keyed_upsert_configured: false,
441            durable_state: true,
442            dlq: false,
443        };
444        assert_eq!(
445            derive_delivery_guarantee(&base),
446            DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::AtomicWatermark)
447        );
448        // Atomic path degrades without deterministic replay…
449        let non_det = GuaranteeInputs {
450            replay: NonDeterministic,
451            ..base
452        };
453        assert_eq!(
454            derive_delivery_guarantee(&non_det),
455            DeliveryGuarantee::AtLeastOnce
456        );
457        // …but keyed upsert rescues it, source-independent.
458        let keyed = GuaranteeInputs {
459            keyed_upsert_configured: true,
460            ..non_det
461        };
462        assert_eq!(
463            derive_delivery_guarantee(&keyed),
464            DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::KeyedUpsert)
465        );
466        // A DLQ or missing durable state disables atomic; keyed still applies.
467        let dlq = GuaranteeInputs {
468            dlq: true,
469            keyed_upsert_configured: true,
470            ..base
471        };
472        assert_eq!(
473            derive_delivery_guarantee(&dlq),
474            DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::KeyedUpsert)
475        );
476        let mem_state = GuaranteeInputs {
477            durable_state: false,
478            ..base
479        };
480        assert_eq!(
481            derive_delivery_guarantee(&mem_state),
482            DeliveryGuarantee::AtLeastOnce
483        );
484    }
485
486    #[test]
487    fn guarantee_display_is_human_readable() {
488        assert_eq!(DeliveryGuarantee::AtLeastOnce.to_string(), "at-least-once");
489        assert_eq!(
490            DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::AtomicWatermark)
491                .to_string(),
492            "effectively-once (atomic watermark)"
493        );
494        assert_eq!(
495            DeliveryGuarantee::EffectivelyOnce(EffectivelyOnceMechanism::KeyedUpsert).to_string(),
496            "effectively-once (keyed upsert)"
497        );
498    }
499
500    #[test]
501    fn capability_enums_default_to_weakest() {
502        assert_eq!(
503            ReplayGuarantee::default(),
504            ReplayGuarantee::NonDeterministic
505        );
506        assert_eq!(SinkGuarantee::default(), SinkGuarantee::AtLeastOnce);
507    }
508
509    #[test]
510    fn delivery_mode_serde_is_snake_case_and_defaults_at_least_once() {
511        assert_eq!(DeliveryMode::default(), DeliveryMode::AtLeastOnce);
512        assert_eq!(
513            serde_json::to_string(&DeliveryMode::ExactlyOnce).unwrap(),
514            "\"exactly_once\""
515        );
516        let m: DeliveryMode = serde_json::from_str("\"at_least_once\"").unwrap();
517        assert_eq!(m, DeliveryMode::AtLeastOnce);
518    }
519
520    #[test]
521    fn overwrite_staging_suffixes_are_pinned() {
522        // These are on-the-wire relation names: a live destination staged under
523        // the old spelling would be stranded (and the new spelling's relation
524        // dropped) if either value ever moved.
525        assert_eq!(OVERWRITE_STAGING_SUFFIX, "__faucet_ovw");
526        assert_eq!(OVERWRITE_STAGING_OLD_SUFFIX, "__faucet_ovw_old");
527        assert_ne!(OVERWRITE_STAGING_SUFFIX, OVERWRITE_STAGING_OLD_SUFFIX);
528        // An orphan sweep matching the base suffix as a prefix must also catch
529        // the mid-swap relation.
530        assert!(OVERWRITE_STAGING_OLD_SUFFIX.starts_with(OVERWRITE_STAGING_SUFFIX));
531    }
532}
533
534#[cfg(test)]
535mod scope_key_tests {
536    use super::*;
537
538    /// #456 L1: the SQL sinks store the scope as a length-capped PRIMARY KEY, so
539    /// a long child scope (its key comes from record data and has no bound) either
540    /// errored or — under a non-strict MySQL sql_mode — truncated, collapsing two
541    /// rows onto one watermark.
542    #[test]
543    fn scope_key_passes_short_scopes_through_and_shortens_long_ones() {
544        // Backwards compatible: anything that fit before is returned verbatim, so
545        // watermarks written before this existed still resolve.
546        assert_eq!(scope_key("pipe::row", 255), "pipe::row");
547        let exactly = "x".repeat(255);
548        assert_eq!(scope_key(&exactly, 255), exactly);
549
550        // Over the cap: shortened, and within the cap.
551        let long = format!("pipe::row::{}", "k".repeat(400));
552        let key = scope_key(&long, 255);
553        assert!(key.len() <= 255, "len {}", key.len());
554        assert_ne!(key, long);
555        // Keeps a readable head so the row is still attributable.
556        assert!(key.starts_with("pipe::row::"), "{key}");
557        assert!(key.contains("__h:"), "{key}");
558    }
559
560    #[test]
561    fn scope_key_is_deterministic_and_distinguishes_scopes() {
562        let a = format!("pipe::row::{}", "a".repeat(400));
563        let b = format!("pipe::row::{}", "b".repeat(400));
564        // Stable across calls — the watermark must resolve after a restart.
565        assert_eq!(scope_key(&a, 255), scope_key(&a, 255));
566        // Two distinct scopes must not collide onto one watermark. Under plain
567        // truncation both of these would become the same 255-char prefix.
568        assert_ne!(scope_key(&a, 255), scope_key(&b, 255));
569        assert_eq!(&a[..255], &format!("pipe::row::{}", "a".repeat(400))[..255]);
570    }
571
572    #[test]
573    fn scope_key_truncates_on_a_char_boundary() {
574        // A multi-byte head must not be split mid-character (that would panic).
575        let long = format!("pipé::{}", "é".repeat(400));
576        let key = scope_key(&long, 255);
577        assert!(key.len() <= 255);
578        assert!(key.contains("__h:"), "{key}");
579    }
580}