Skip to main content

nmbrs_runtime/
op_modifier.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Per-op field modifiers via initializer-time currying.
5//!
6//! Adapter-agnostic implementation of the monadic-compose
7//! "enhancer chain" pattern from upstream nosqlbench
8//! (`Cqld4BaseOpDispenser.getEnhancedStmtFunc` →
9//! `ParsedOp.enhanceFuncOptionally`). The dispenser's initializer
10//! resolves universal-field names through the Polydat scope ONCE,
11//! captures the resolved values into modifier structs, and stores
12//! the resulting `ModifierChain<T>` on the dispenser. Critical path
13//! is one call: `chain.apply(&mut target)`.
14//!
15//! See [`SRD 73`](../../../docs/SRD/73_op_field_modifiers.md)
16//! for the design rationale and CQL-specific surface.
17//!
18//! # Two-phase contract
19//!
20//! - **Initializer (one-time):** the adapter walks its declared
21//!   universal-field list, calls `parent.lookup(name)` on the GK
22//!   kernel handed to `DriverAdapter::map_op`, and for each
23//!   bound name pushes an `OpFieldModifier<T>` that has CAPTURED
24//!   the resolved value as an owned field. Names not bound in
25//!   scope contribute nothing — the driver's native default stays
26//!   in force for that knob.
27//!
28//! - **Critical path (per cycle):** the dispenser calls
29//!   `chain.apply(&mut stmt)` on the constructed engine
30//!   statement immediately before binding values / sending. Each
31//!   active modifier applies its captured setter; no Polydat access,
32//!   no name resolution, no map lookup.
33//!
34//! # Trace sink (optional, lazy)
35//!
36//! A `ModifierTraceSink` may be installed at the session level.
37//! When present, the chain calls `sink.modifier_applied(...)`
38//! after each `apply`, handing the sink a `&dyn Fn() ->
39//! serde_json::Value` closure. The closure is invoked only if
40//! the sink decides to record the event (sink-internal gate),
41//! so JSON serialization is paid only when a consumer will read
42//! the value. With no sink installed, `apply` runs through a
43//! tight loop with zero closure construction.
44//!
45//! Three concepts are kept ORTHOGONAL — see SRD 73
46//! §"Tracing terminology":
47//!
48//! 1. The CQL query-tracing **subsystem** (rows in
49//!    `system_traces.*` on the cluster). Engaged per-op via the
50//!    `cql_trace` universal field. A DATA SOURCE.
51//! 2. The Rust `tracing` crate's **log severity** filter. Not
52//!    used by nmbrs (we use [`crate::observer`] /
53//!    [`crate::trace_router`] instead). Orthogonal to (1).
54//! 3. nmbrs event-log **emissions** — checkpoint JSONL etc.
55//!    Pluggable via `ModifierTraceSink`. Orthogonal to (1) and (2).
56
57use std::sync::Arc;
58use std::sync::OnceLock;
59
60use nmbrs_metrics::labels::Labels;
61
62/// A conditional per-op modifier.
63///
64/// Modifiers are constructed only when the user actually bound
65/// the corresponding field in the Polydat scope — the chain pre-filters
66/// at build time, so anything reachable through
67/// [`ModifierChain::apply`] is by construction active. There is
68/// no `is_active()` method.
69///
70/// The target type `T` is the engine's per-statement type (e.g.
71/// `scylla::statement::Statement` or
72/// `cassandra_cpp::Statement`). Each engine module provides its
73/// own `impl OpFieldModifier<EngineStatement> for FooMod` types
74/// that translate a captured Rust value into the engine's setter
75/// call.
76pub trait OpFieldModifier<T>: Send + Sync + 'static {
77    /// User-facing field name. Matches the op-template key and
78    /// the adapter's universal-field selector list. Returned as
79    /// `&'static str` so trace sinks can keep cheap references.
80    fn field_name(&self) -> &'static str;
81
82    /// Mutate the target with the captured value. Must not look
83    /// anything up — all state needed for the mutation is already
84    /// captured in the modifier struct's fields.
85    fn apply(&self, target: &mut T);
86
87    /// Structured diagnostic representation of the captured
88    /// value, used by `ModifierTraceSink` consumers. Called
89    /// LAZILY — only when a sink is installed AND the sink
90    /// decides to record this event.
91    fn diagnostic_value(&self) -> serde_json::Value;
92}
93
94/// A composed chain of `OpFieldModifier<T>` — the moral
95/// equivalent of upstream NB Java's specialized `LongFunction<S>`.
96///
97/// Built once in the dispenser initializer; called per cycle on
98/// the critical path. Carries an optional `ModifierTraceSink` for
99/// cross-cutting observation; the sink hot-path is gated lazily
100/// (see [`Self::apply`]).
101pub struct ModifierChain<T> {
102    op_label: String,
103    active: Vec<Box<dyn OpFieldModifier<T>>>,
104    event_sink: Option<Arc<dyn ModifierTraceSink>>,
105}
106
107impl<T: 'static> ModifierChain<T> {
108    /// Construct a chain. `active` MUST already exclude inactive
109    /// modifiers; the caller (typically a per-adapter builder
110    /// like `build_cql_modifier_chain`) is responsible for
111    /// dropping `None` results from `parent.lookup(...)` so this
112    /// vec contains only modifiers the user actually bound.
113    pub fn new(
114        op_label: impl Into<String>,
115        active: Vec<Box<dyn OpFieldModifier<T>>>,
116        event_sink: Option<Arc<dyn ModifierTraceSink>>,
117    ) -> Self {
118        Self {
119            op_label: op_label.into(),
120            active,
121            event_sink,
122        }
123    }
124
125    /// True when the user did not bind any universal field for
126    /// this op. Lets the caller skip the `apply` call entirely
127    /// on a fully-default path — the most common case in practice.
128    pub fn is_empty(&self) -> bool {
129        self.active.is_empty()
130    }
131
132    /// Number of active modifiers. For diagnostics / tests.
133    pub fn len(&self) -> usize {
134        self.active.len()
135    }
136
137    /// The op label this chain was built for. Used by sinks for
138    /// event correlation.
139    pub fn op_label(&self) -> &str {
140        &self.op_label
141    }
142
143    /// Critical-path entry. Apply every active modifier to the
144    /// target.
145    ///
146    /// Two arms:
147    ///
148    /// - **No sink installed:** tight loop over `active`. No
149    ///   closure construction, no JSON serialization, no virtual
150    ///   dispatch beyond each modifier's own `apply`.
151    /// - **Sink installed:** after each `apply`, the sink is
152    ///   handed a `&dyn Fn() -> serde_json::Value` that the sink
153    ///   may or may not invoke. JSON serialization is paid only
154    ///   when the sink actually wants the value (e.g.
155    ///   trace-router has at least one subscriber to this
156    ///   adapter's traces).
157    pub fn apply(&self, target: &mut T) {
158        match &self.event_sink {
159            None => {
160                for m in &self.active {
161                    m.apply(target);
162                }
163            }
164            Some(sink) => {
165                for m in &self.active {
166                    m.apply(target);
167                    sink.modifier_applied(&self.op_label, m.field_name(), &|| m.diagnostic_value());
168                }
169            }
170        }
171    }
172}
173
174impl<T: 'static> std::fmt::Debug for ModifierChain<T> {
175    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
176        f.debug_struct("ModifierChain")
177            .field("op_label", &self.op_label)
178            .field("active_count", &self.active.len())
179            .field("event_sink", &self.event_sink.is_some())
180            .finish()
181    }
182}
183
184/// Cross-cutting observer for modifier application.
185///
186/// One sink per session, installed via
187/// [`install_session_sink`]. The chain hands the sink the op
188/// label, the field name, and a closure that produces the
189/// diagnostic JSON on demand. Sinks SHOULD check their own
190/// filter state before invoking the closure — `diagnostic_value`
191/// can do real work (string formatting, struct serialization)
192/// and we want that paid only when a consumer will read it.
193///
194/// Built-in implementations:
195///
196/// - [`TraceRouterSink`] — routes to [`crate::trace_router`]
197///   with a `{component: "op_modifier", op: <label>, field:
198///   <name>}` label set. Gates on [`crate::trace_router::enabled`]
199///   so the closure is invoked only when at least one trace
200///   target is configured.
201///
202/// Future:
203///
204/// - `JsonEventSink` — writes structured records to the SRD-44a
205///   checkpoint JSONL. To be added when the JSON event-log
206///   subscription path lands.
207pub trait ModifierTraceSink: Send + Sync {
208    /// Invoked from [`ModifierChain::apply`] after a modifier
209    /// runs. `value_fn` produces the JSON diagnostic on demand;
210    /// the sink decides whether to invoke it.
211    fn modifier_applied(
212        &self,
213        op: &str,
214        field: &'static str,
215        value_fn: &dyn Fn() -> serde_json::Value,
216    );
217}
218
219/// Built-in sink that routes modifier events through
220/// [`crate::trace_router`].
221///
222/// Each event becomes a trace-router log line with labels
223/// `{component: "op_modifier", op: <op_label>, field:
224/// <field_name>}`, formatted as `"<field>=<json_value>"`.
225/// Operators select which adapters / ops / fields to record via
226/// the `--trace=<spec>` CLI surface.
227///
228/// The sink checks [`crate::trace_router::enabled`] before
229/// invoking `value_fn`, so the JSON serialization cost is paid
230/// only when the trace router has at least one subscribed
231/// target.
232pub struct TraceRouterSink;
233
234impl ModifierTraceSink for TraceRouterSink {
235    fn modifier_applied(
236        &self,
237        op: &str,
238        field: &'static str,
239        value_fn: &dyn Fn() -> serde_json::Value,
240    ) {
241        // Cheap gate: an atomic load. When no trace-router
242        // target is configured we never compute the JSON.
243        if !crate::trace_router::enabled() {
244            return;
245        }
246        let value = value_fn();
247        let labels = Labels::of("component", "op_modifier")
248            .with("op", op.to_string())
249            .with("field", field);
250        let message = format!("{field}={value}");
251        crate::trace_router::log(&labels, &message);
252    }
253}
254
255/// Session-global modifier trace sink. Set once by the runner
256/// after parsing session config; adapters call
257/// [`session_sink`] to fetch the optional handle and pass it to
258/// `ModifierChain::new`.
259static SESSION_MODIFIER_SINK: OnceLock<Arc<dyn ModifierTraceSink>> = OnceLock::new();
260
261/// Install the session-global trace sink. Idempotent — first
262/// caller wins (the runner installs at session-init).
263pub fn install_session_sink(sink: Arc<dyn ModifierTraceSink>) {
264    let _ = SESSION_MODIFIER_SINK.set(sink);
265}
266
267/// Fetch the session-global trace sink, if one was installed.
268/// Adapters call this from their `map_op` to pass into
269/// `ModifierChain::new`.
270pub fn session_sink() -> Option<Arc<dyn ModifierTraceSink>> {
271    SESSION_MODIFIER_SINK.get().cloned()
272}
273
274// =========================================================================
275// Tests
276// =========================================================================
277
278#[cfg(test)]
279mod tests {
280    use super::*;
281    use std::sync::Mutex;
282    use std::sync::atomic::{AtomicUsize, Ordering};
283
284    /// Synthetic target type — stands in for a CQL Statement.
285    #[derive(Default, Debug, PartialEq)]
286    struct FakeStmt {
287        timeout_ms: Option<u64>,
288        consistency: Option<String>,
289        page_size: Option<i32>,
290    }
291
292    struct TimeoutMod {
293        ms: u64,
294    }
295    impl OpFieldModifier<FakeStmt> for TimeoutMod {
296        fn field_name(&self) -> &'static str {
297            "request_timeout_ms"
298        }
299        fn apply(&self, t: &mut FakeStmt) {
300            t.timeout_ms = Some(self.ms);
301        }
302        fn diagnostic_value(&self) -> serde_json::Value {
303            serde_json::Value::from(self.ms)
304        }
305    }
306
307    struct ConsistencyMod {
308        cl: String,
309    }
310    impl OpFieldModifier<FakeStmt> for ConsistencyMod {
311        fn field_name(&self) -> &'static str {
312            "consistency"
313        }
314        fn apply(&self, t: &mut FakeStmt) {
315            t.consistency = Some(self.cl.clone());
316        }
317        fn diagnostic_value(&self) -> serde_json::Value {
318            serde_json::Value::String(self.cl.clone())
319        }
320    }
321
322    #[test]
323    fn empty_chain_is_noop() {
324        let chain: ModifierChain<FakeStmt> = ModifierChain::new("op1", vec![], None);
325        assert!(chain.is_empty());
326        let mut stmt = FakeStmt::default();
327        chain.apply(&mut stmt);
328        assert_eq!(stmt, FakeStmt::default());
329    }
330
331    #[test]
332    fn single_modifier_applies_captured_value() {
333        let chain: ModifierChain<FakeStmt> = ModifierChain::new(
334            "op_drop_index",
335            vec![Box::new(TimeoutMod { ms: 300_000 })],
336            None,
337        );
338        assert_eq!(chain.len(), 1);
339        let mut stmt = FakeStmt::default();
340        chain.apply(&mut stmt);
341        assert_eq!(stmt.timeout_ms, Some(300_000));
342        assert_eq!(stmt.consistency, None);
343        assert_eq!(stmt.page_size, None);
344    }
345
346    #[test]
347    fn multiple_modifiers_apply_in_order() {
348        let chain: ModifierChain<FakeStmt> = ModifierChain::new(
349            "op_select",
350            vec![
351                Box::new(TimeoutMod { ms: 5_000 }),
352                Box::new(ConsistencyMod {
353                    cl: "LOCAL_QUORUM".to_string(),
354                }),
355            ],
356            None,
357        );
358        let mut stmt = FakeStmt::default();
359        chain.apply(&mut stmt);
360        assert_eq!(stmt.timeout_ms, Some(5_000));
361        assert_eq!(stmt.consistency, Some("LOCAL_QUORUM".to_string()));
362    }
363
364    /// Sink that records each event into a Vec for inspection.
365    /// Models the "sink enabled" case where the closure IS
366    /// invoked.
367    struct RecordingSink {
368        records: Mutex<Vec<(String, &'static str, serde_json::Value)>>,
369        invocation_count: AtomicUsize,
370    }
371    impl ModifierTraceSink for RecordingSink {
372        fn modifier_applied(
373            &self,
374            op: &str,
375            field: &'static str,
376            value_fn: &dyn Fn() -> serde_json::Value,
377        ) {
378            self.invocation_count.fetch_add(1, Ordering::Relaxed);
379            let value = value_fn(); // sink chooses to invoke
380            self.records
381                .lock()
382                .unwrap()
383                .push((op.to_string(), field, value));
384        }
385    }
386
387    /// Sink that NEVER invokes the closure — models the
388    /// "filter-disabled" hot path. We use a counter on the
389    /// modifier's diagnostic_value to prove the closure was not
390    /// called.
391    struct GatedSink {
392        skipped: AtomicUsize,
393    }
394    impl ModifierTraceSink for GatedSink {
395        fn modifier_applied(
396            &self,
397            _op: &str,
398            _field: &'static str,
399            _value_fn: &dyn Fn() -> serde_json::Value,
400        ) {
401            // Sink's filter says "don't record" — closure is
402            // never invoked, no JSON work is done.
403            self.skipped.fetch_add(1, Ordering::Relaxed);
404        }
405    }
406
407    /// Modifier whose `diagnostic_value` increments a counter,
408    /// so a test can prove whether the closure was invoked.
409    struct CountingMod {
410        diag_calls: Arc<AtomicUsize>,
411    }
412    impl OpFieldModifier<FakeStmt> for CountingMod {
413        fn field_name(&self) -> &'static str {
414            "request_timeout_ms"
415        }
416        fn apply(&self, t: &mut FakeStmt) {
417            t.timeout_ms = Some(42);
418        }
419        fn diagnostic_value(&self) -> serde_json::Value {
420            self.diag_calls.fetch_add(1, Ordering::Relaxed);
421            serde_json::Value::from(42u64)
422        }
423    }
424
425    #[test]
426    fn recording_sink_sees_all_fired_modifiers() {
427        let sink = Arc::new(RecordingSink {
428            records: Mutex::new(Vec::new()),
429            invocation_count: AtomicUsize::new(0),
430        });
431        let chain: ModifierChain<FakeStmt> = ModifierChain::new(
432            "op_drop_index",
433            vec![
434                Box::new(TimeoutMod { ms: 300_000 }),
435                Box::new(ConsistencyMod {
436                    cl: "ONE".to_string(),
437                }),
438            ],
439            Some(sink.clone()),
440        );
441        let mut stmt = FakeStmt::default();
442        chain.apply(&mut stmt);
443
444        assert_eq!(sink.invocation_count.load(Ordering::Relaxed), 2);
445        let records = sink.records.lock().unwrap();
446        assert_eq!(records.len(), 2);
447        assert_eq!(records[0].0, "op_drop_index");
448        assert_eq!(records[0].1, "request_timeout_ms");
449        assert_eq!(records[0].2, serde_json::json!(300_000));
450        assert_eq!(records[1].1, "consistency");
451        assert_eq!(records[1].2, serde_json::json!("ONE"));
452    }
453
454    #[test]
455    fn gated_sink_does_not_invoke_diagnostic_closure() {
456        let diag_calls = Arc::new(AtomicUsize::new(0));
457        let sink = Arc::new(GatedSink {
458            skipped: AtomicUsize::new(0),
459        });
460
461        let chain: ModifierChain<FakeStmt> = ModifierChain::new(
462            "op_select",
463            vec![Box::new(CountingMod {
464                diag_calls: diag_calls.clone(),
465            })],
466            Some(sink.clone()),
467        );
468
469        let mut stmt = FakeStmt::default();
470        for _ in 0..1000 {
471            chain.apply(&mut stmt);
472        }
473
474        // Sink was called 1000 times — modifier_applied fired
475        // for every apply.
476        assert_eq!(sink.skipped.load(Ordering::Relaxed), 1000);
477        // But the diagnostic closure was NEVER invoked, because
478        // the sink chose not to call it. This is the laziness
479        // contract: JSON work is paid only when a consumer
480        // actually wants the value.
481        assert_eq!(diag_calls.load(Ordering::Relaxed), 0);
482    }
483
484    #[test]
485    fn no_sink_hot_path_does_not_invoke_diagnostic_closure() {
486        let diag_calls = Arc::new(AtomicUsize::new(0));
487        let chain: ModifierChain<FakeStmt> = ModifierChain::new(
488            "op_select",
489            vec![Box::new(CountingMod {
490                diag_calls: diag_calls.clone(),
491            })],
492            None, // no sink — None arm of the match
493        );
494
495        let mut stmt = FakeStmt::default();
496        for _ in 0..1000 {
497            chain.apply(&mut stmt);
498        }
499
500        // 1000 applies but zero diagnostic invocations. The
501        // None branch of the match in `apply` skips the sink
502        // entirely.
503        assert_eq!(diag_calls.load(Ordering::Relaxed), 0);
504        // And the modifier itself DID run all 1000 times.
505        assert_eq!(stmt.timeout_ms, Some(42));
506    }
507
508    #[test]
509    fn diagnostic_value_returns_json() {
510        // Sanity-check that diagnostic_value produces typed
511        // JSON values the way per-engine modifiers will.
512        let m = TimeoutMod { ms: 300_000 };
513        assert_eq!(m.diagnostic_value(), serde_json::json!(300_000));
514
515        let m = ConsistencyMod {
516            cl: "LOCAL_QUORUM".to_string(),
517        };
518        assert_eq!(m.diagnostic_value(), serde_json::json!("LOCAL_QUORUM"));
519    }
520
521    #[test]
522    fn debug_impl_reports_active_count_without_revealing_state() {
523        let chain: ModifierChain<FakeStmt> =
524            ModifierChain::new("op1", vec![Box::new(TimeoutMod { ms: 100 })], None);
525        let s = format!("{:?}", chain);
526        assert!(s.contains("op_label"));
527        assert!(s.contains("active_count"));
528        assert!(s.contains("event_sink"));
529    }
530}