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}