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}