Skip to main content

nmbrs_runtime/
adapter.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Adapter traits: the tiered interface that database/protocol drivers implement (SRD 38).
5//!
6//! The tiered interface separates init-time template analysis from cycle-time execution:
7//! - `DriverAdapter::map_op()` — called once per template at activity startup
8//! - `OpDispenser::execute()` — called per-cycle to bind values and execute
9
10use std::any::Any;
11use std::fmt;
12use std::sync::OnceLock;
13
14// Re-export so adapter crates can write `use nmbrs_runtime::adapter::ExecCtx`
15// alongside `use nmbrs_runtime::adapter::{OpDispenser, ResolvedFields, ...}`.
16pub use crate::fixture::ExecCtx;
17
18// Re-export the engine-neutral kernel surface so adapter `map_op`
19// impls can name the `parent: Arc<dyn Kernel>` parameter type, and
20// resolve a name the way a scope does (`KernelLookup`), without each
21// adapter crate taking a direct polydat dependency. SRD-68 §"Adapter
22// API surface" pins this as the canonical import path for adapters.
23// The kernels an adapter receives may be on any engine.
24pub use polydat::Kernel;
25pub use polydat::kernel::interp::{KernelLookup, Lookup};
26
27// Re-export the binder API so `binders_for` impls in adapter
28// crates can name `Binder` / `BinderSlot` / `PortType` through
29// the existing `nmbrs_runtime::adapter` import path, no direct
30// polydat dependency required. Same reason as the Kernel
31// re-export above.
32pub use polydat::ast::PortType;
33
34/// Verify typed binders against a kernel of any engine: each slot's
35/// wire resolves as an output of the kernel's program, else as an input
36/// (the coordinate / extern wires an op-template kernel declares). The
37/// engine-neutral form of `polydat::binder::verify_against_kernel`,
38/// which takes an interpreter kernel. Violation messages are joined
39/// with `"; "`.
40pub fn verify_binders(binders: &[Binder], kernel: &dyn Kernel) -> Result<(), String> {
41    let violations = polydat::binder::verify_binders(binders, |name: &str| {
42        kernel
43            .output_type(name)
44            .or_else(|| kernel.input_port_type(name))
45    });
46    if violations.is_empty() {
47        Ok(())
48    } else {
49        Err(violations
50            .into_iter()
51            .map(|v| v.message)
52            .collect::<Vec<_>>()
53            .join("; "))
54    }
55}
56pub use polydat::binder::{Binder, BinderSlot};
57
58/// Boxed, `Send` future returned by [`DriverAdapter::map_op`] — yields a
59/// boxed [`OpDispenser`] (or an error message). The lifetime ties the
60/// future to the borrowed `&self` / template references.
61pub type MapOpFuture<'a> = std::pin::Pin<
62    Box<dyn std::future::Future<Output = Result<Box<dyn OpDispenser>, String>> + Send + 'a>,
63>;
64
65/// Boxed, `Send` future produced by an adapter/driver `create` factory —
66/// yields a connected [`DriverAdapter`] (or an error message).
67pub type CreateAdapterFuture = std::pin::Pin<
68    Box<dyn std::future::Future<Output = Result<std::sync::Arc<dyn DriverAdapter>, String>> + Send>,
69>;
70
71/// Trait for adapter-specific result bodies.
72///
73/// The adapter defines its own concrete result type and implements
74/// this trait. Internal adapter code can downcast via `as_any()` to
75/// access native types (e.g., CQL rows, HTTP response structs).
76/// External consumers call `to_json()` for a universal representation.
77///
78/// The `element_count()` and `byte_count()` methods support result
79/// traversal — the framework uses these to verify that result data
80/// was fully received and to record consumption metrics.
81pub trait ResultBody: Send + Sync + fmt::Debug {
82    /// Serialize this result to JSON for logging, capture, verification.
83    fn to_json(&self) -> serde_json::Value;
84
85    /// Downcast support for adapter-internal use.
86    fn as_any(&self) -> &dyn Any;
87
88    /// Count of logical elements in this result (rows, records, items).
89    /// Used by the result traverser for metrics. Default 1 (one result).
90    fn element_count(&self) -> u64 {
91        1
92    }
93
94    /// Size in bytes of the result payload, if known.
95    /// Used by the result traverser for throughput metrics.
96    fn byte_count(&self) -> Option<u64> {
97        None
98    }
99
100    /// SRD-66 §"Surface 1" body projection — return the body's
101    /// canonical text representation for use by result-binding
102    /// expressions (the magic `body: Str` extern). Default
103    /// implementation falls back to `serde_json::to_string`
104    /// over `to_json()`; adapters with native text bodies (e.g.
105    /// `TextBody`, the CQL describe row) override to return
106    /// the underlying text directly.
107    fn to_text(&self) -> String {
108        serde_json::to_string(&self.to_json()).unwrap_or_default()
109    }
110}
111
112/// Simple text result body.
113#[derive(Debug, Clone)]
114pub struct TextBody(pub String);
115
116impl ResultBody for TextBody {
117    fn to_json(&self) -> serde_json::Value {
118        serde_json::Value::String(self.0.clone())
119    }
120    fn as_any(&self) -> &dyn Any {
121        self
122    }
123    fn to_text(&self) -> String {
124        self.0.clone()
125    }
126}
127
128/// JSON result body. Adapters that natively produce JSON
129/// (HTTP/REST endpoints, the Jolokia JMX bridge, JSON-RPC) wrap
130/// the parsed value here so `verify:` field assertions and
131/// result-binding extractors can address nested keys (`status`,
132/// `value`, etc.) without a re-parse step.
133///
134/// `element_count` follows the JSON-array convention: array
135/// bodies report `array.len()`, scalar/object bodies report `1`.
136/// That makes `await_empty` polling work naturally against
137/// endpoints whose "still running" state is a non-empty array
138/// and "done" state is `[]`.
139#[derive(Debug, Clone)]
140pub struct JsonBody(pub serde_json::Value);
141
142impl ResultBody for JsonBody {
143    fn to_json(&self) -> serde_json::Value {
144        self.0.clone()
145    }
146    fn as_any(&self) -> &dyn Any {
147        self
148    }
149    fn element_count(&self) -> u64 {
150        match &self.0 {
151            serde_json::Value::Array(arr) => arr.len() as u64,
152            _ => 1,
153        }
154    }
155    fn to_text(&self) -> String {
156        serde_json::to_string(&self.0).unwrap_or_default()
157    }
158}
159
160/// The result of a successful operation.
161///
162/// If you have an `OpResult`, the operation succeeded. Failure is
163/// represented by `ExecutionError`, not by a flag on the result.
164/// Protocol-specific status codes (HTTP, CQL) live inside the
165/// adapter's `ResultBody` implementation, not on the generic result.
166#[derive(Default)]
167pub struct OpResult {
168    /// Adapter-specific response body. The adapter owns the native
169    /// type; consumers call `.to_json()` for a universal view.
170    /// Adapter-internal code can downcast via `.as_any()`.
171    /// `None` for operations with no meaningful result (e.g., DDL).
172    pub body: Option<Box<dyn ResultBody>>,
173    /// If true, this op was conditionally skipped (via `if:` field).
174    /// The activity loop counts this as a skip, not a success or error.
175    pub skipped: bool,
176}
177
178impl OpResult {
179    /// Create a skipped result (no execution).
180    pub fn skipped() -> Self {
181        Self {
182            body: None,
183            skipped: true,
184        }
185    }
186}
187
188impl fmt::Debug for OpResult {
189    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
190        f.debug_struct("OpResult")
191            .field("body", &self.body.as_ref().map(|b| b.to_json()))
192            .finish()
193    }
194}
195
196/// Execution error with scope delamination.
197///
198/// Distinguishes between per-op errors (template-specific, retryable)
199/// and adapter-level errors (connection-wide, affect all ops).
200#[derive(Debug)]
201pub enum ExecutionError {
202    /// Per-op failure: this specific operation failed. Template-specific,
203    /// may be retried with the same resolved fields.
204    Op(AdapterError),
205    /// Adapter-level failure: the driver connection or session is
206    /// degraded. Affects all ops. The activity may need to pause or stop.
207    Adapter(AdapterError),
208}
209
210impl ExecutionError {
211    /// Access the inner AdapterError regardless of scope.
212    pub fn error(&self) -> &AdapterError {
213        match self {
214            ExecutionError::Op(e) | ExecutionError::Adapter(e) => e,
215        }
216    }
217
218    /// Whether this is an adapter-level (connection-wide) error.
219    pub fn is_adapter_level(&self) -> bool {
220        matches!(self, ExecutionError::Adapter(_))
221    }
222}
223
224impl fmt::Display for ExecutionError {
225    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
226        match self {
227            ExecutionError::Op(e) => write!(f, "[op] {e}"),
228            ExecutionError::Adapter(e) => write!(f, "[adapter] {e}"),
229        }
230    }
231}
232
233impl std::error::Error for ExecutionError {}
234
235/// Error from an adapter operation.
236#[derive(Debug)]
237pub struct AdapterError {
238    /// Error classification name (for error handler routing).
239    pub error_name: String,
240    /// Human-readable error message.
241    pub message: String,
242    /// Hint to the executor: is this error worth retrying?
243    pub retryable: bool,
244}
245
246impl fmt::Display for AdapterError {
247    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
248        write!(f, "[{}] {}", self.error_name, self.message)
249    }
250}
251
252impl std::error::Error for AdapterError {}
253
254/// A protocol-specific driver adapter. Constructed once per activity,
255/// shared across fibers via Arc.
256///
257/// The adapter owns the driver connection (session, client, pool) and
258/// provides OpDispensers that pre-process each op template at init time.
259pub trait DriverAdapter: Send + Sync + 'static {
260    /// Human-readable adapter name (e.g., "cql", "http", "stdout").
261    fn name(&self) -> &str;
262
263    /// Map an op template into a dispenser. Called once per unique op
264    /// template at activity startup — before any cycles execute.
265    ///
266    /// Init-time work: parse the op, prepare statements, pre-compute
267    /// bind-point resolution, validate field names, attach metrics.
268    ///
269    /// `parent` is the phase scope kernel — the Polydat context the op
270    /// template's matter (phase `bindings:`, `result:` block, etc.)
271    /// should attach to. Adapters that need their own Polydat context
272    /// for op-template-scope name resolution clone the Arc and
273    /// either retain it directly (no op-level matter) or use it as
274    /// the parent their canonical kernel is bound under
275    /// (`polydat::kernel::bind_under` / `ScopeModule::instantiate_under`) (op-level matter present); adapters
276    /// with no Polydat needs ignore the parameter. The Arc lets the
277    /// dispenser own a long-lived reference to the canonical
278    /// kernel without re-cloning state. See SRD-68 §"Adapter API
279    /// surface".
280    /// Construct an [`OpDispenser`] for one op template — the
281    /// per-op dispenser-initialization stack frame.
282    ///
283    /// This is where the currying stack completes: prepare any
284    /// protocol-level handles (CQL: prepare statement, read
285    /// parameter metadata), build the per-cycle data pullers,
286    /// fold everything into the dispenser the runtime will then
287    /// call repeatedly with `execute(cycle, ctx)`. Nothing about
288    /// init-time state needs to outlive this call — once `map_op`
289    /// returns, the dispenser holds whatever it needs and the
290    /// rest is dropped.
291    ///
292    /// ## Per-op binder verification (the typed-lvalue contract)
293    ///
294    /// Per-op compulsion: for any op-template field this
295    /// dispenser will bind through a typed-parameter API
296    /// (anything beyond pure text concatenation into a `Str`
297    /// lvalue), the implementor MUST construct the appropriate
298    /// [`Binder`] shape from its protocol-side metadata
299    /// (positional / named / single — see [`polydat::binder`])
300    /// and verify it against `parent` via
301    /// [`polydat::binder::verify_against_kernel`] before
302    /// returning the dispenser. A verification failure surfaces
303    /// as a `map_op` `Err`; construction stops before any cycle
304    /// runs and the operator sees the rvalue→lvalue mismatch
305    /// (one violation per slot, no silent passthrough).
306    ///
307    /// This is stack-local work: the binder is constructed,
308    /// verified, used to wire up the typed binding path in the
309    /// dispenser, and dropped. It is NOT cached on the adapter,
310    /// returned as a sidecar, or otherwise persisted past the
311    /// `map_op` return. Per-op-template, not adapter-wide.
312    ///
313    /// Adapters whose every op-template field is a text template
314    /// (stdout, http body templates, testkit captures) can skip
315    /// the verify step — every wire reference's lvalue is `Str`
316    /// and the rvalue→lvalue rule permits any rvalue into a `Str`
317    /// lvalue, so verification is always a no-op for them. The
318    /// runtime guard in `wires::substitute_via_wires` remains the
319    /// catch-all safety net for that path.
320    fn map_op<'a>(
321        &'a self,
322        template: &'a nmbrs_workload::model::ParsedOp,
323        parent: std::sync::Arc<dyn Kernel>,
324    ) -> MapOpFuture<'a>;
325
326    /// Default metric names to display on the status line for this adapter.
327    ///
328    /// Each entry is a metric name (matching `adapter_metrics()` labels)
329    /// and a display label. Workloads can override this via a `status:`
330    /// field on phases or ops. Default: empty (no adapter-specific status).
331    fn default_status_metrics(&self) -> Vec<StatusMetric> {
332        Vec::new()
333    }
334
335    /// Preferred display mode for this adapter.
336    ///
337    /// When multiple adapters are involved in a workload, the runtime
338    /// uses the most restrictive (lowest) mode. Adapters that use
339    /// raw terminal output (plotter) return `Off` to prevent the TUI
340    /// from entering alternate screen.
341    ///
342    /// - `Auto`: adapter is compatible with TUI (default)
343    /// - `Off`: adapter requires raw stderr/stdout, TUI must not activate
344    fn display_preference(&self) -> DisplayPreference {
345        DisplayPreference::Auto
346    }
347
348    /// Declare the set of op-field names this adapter knows how to
349    /// interpret. SRD 30 §"Core-first field processing" requires
350    /// that after the core runtime strips its own fields, every
351    /// remaining key in `ParsedOp.op` must be in this list — an
352    /// unknown field is a hard error, not a silent pass-through.
353    ///
354    /// Return `None` (the default) to opt out of strict validation
355    /// during the transition; existing adapters that haven't been
356    /// audited remain permissive. Adapters that return `Some(...)`
357    /// get the "unknown field" guard automatically — core rejects
358    /// templates with fields the adapter doesn't claim.
359    ///
360    /// Returning an empty slice `Some(&[])` is a valid declaration
361    /// for an adapter that consumes no op fields (e.g. `stdout`
362    /// rendering the raw bindings only).
363    fn known_op_fields(&self) -> Option<&'static [&'static str]> {
364        None
365    }
366
367    /// Adapter-specific params keys allowed under an op's
368    /// top-level `params` (not `op`) section. Returned keys
369    /// extend the core's [`crate::validation::CORE_OP_PARAMS`]
370    /// allow-list at op-validation time. Default: empty —
371    /// adapters that don't need extras simply rely on the
372    /// core vocab. Override to declare adapter-only params
373    /// (e.g. CQL's `cl:` consistency-level overrides if those
374    /// were lifted from `op` to `params` someday).
375    fn known_op_params(&self) -> &'static [&'static str] {
376        &[]
377    }
378
379    /// Declare adapter-specific dynamic controls (SRD 23) on a
380    /// subcomponent attached to the activity's component. Called
381    /// once per activity, *after* the activity declares its own
382    /// `concurrency` / `rate` controls. The default no-op fits
383    /// adapters that have no adapter-level dynamic knobs.
384    ///
385    /// Convention: each adapter that overrides this attaches a
386    /// single subcomponent named after itself (e.g. `cql`,
387    /// `http`) under `parent`, and declares all of its dynamic
388    /// controls there. That keeps controls reachable from every
389    /// descendant scope via the standard SRD-16 walk-up while
390    /// giving each adapter a stable component path label.
391    ///
392    /// The trait method takes an `&dyn DriverAdapter` (i.e.
393    /// `&self`), so adapters that hold per-instance state
394    /// (handles, atomics) can wire up control appliers that
395    /// write into that state. Multiple ops in one activity see
396    /// the same adapter instance and the same controls — there
397    /// is exactly one subcomponent per adapter per activity.
398    fn declare_controls(
399        &self,
400        _parent: &std::sync::Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
401    ) {
402    }
403
404    /// Async teardown hook fired by the resource pool when a
405    /// shared adapter's last reference detaches (or at session
406    /// shutdown). Default: no-op — drop is enough for adapters
407    /// whose underlying driver closes synchronously in its
408    /// destructor.
409    ///
410    /// Override to await a driver-specific close handshake.
411    /// The CQL adapter's override calls
412    /// `Session::close().await` (which wraps
413    /// `cass_session_close`) inside a 5-second timeout so a
414    /// hung node doesn't pin the runtime.
415    ///
416    /// The future returned here borrows `&self` for the
417    /// duration of the await; callers keep the adapter Arc
418    /// alive across the await so the borrow stays valid.
419    fn shutdown<'a>(
420        &'a self,
421    ) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
422        Box::pin(async {})
423    }
424
425    /// SRD-104 — the adapter's **accessor payload**: a type-erased handle
426    /// (`Arc<dyn Any + Send + Sync>`) a kernel node can obtain by fingerprint
427    /// through its kernel tree's resource scope (`ResourceScope::lookup`). The resource pool surfaces this
428    /// through [`crate::resource_pool::SharedResource::accessor_payload`]
429    /// (the pool-shared wrapper delegates here) and stores it on the entry at
430    /// init. Default `None` — an adapter opts in only when it wants kernels
431    /// to reach a live handle (the first consumer is the CQL session handle,
432    /// SRD-103). Built over the adapter's own connected session, so the
433    /// payload and the op-execution path share one resource.
434    fn accessor_payload(&self) -> Option<std::sync::Arc<dyn Any + Send + Sync>> {
435        None
436    }
437}
438
439/// Adapter display preference for TUI activation.
440#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
441pub enum DisplayPreference {
442    /// TUI must not activate — adapter uses raw terminal output.
443    Off = 0,
444    /// Adapter is compatible with TUI (default for most adapters).
445    Auto = 1,
446}
447
448/// A metric to display on the activity status line.
449pub struct StatusMetric {
450    /// Metric name matching the `name` label in `adapter_metrics()` samples.
451    pub metric_name: String,
452    /// Short display label for the status line (e.g., "rows/s").
453    pub display: String,
454    /// How to render the value: "rate" (count/elapsed), "count", "latency".
455    pub render: StatusRender,
456}
457
458/// How to render a status metric value.
459pub enum StatusRender {
460    /// Show as rate: count / elapsed seconds (e.g., "1.2K/s").
461    Rate,
462    /// Show as raw count.
463    Count,
464    /// Show as latency with auto-scaled units.
465    Latency,
466}
467
468/// A per-template op factory. Created at init time by the adapter's
469/// `map_op()`, called per-cycle to bind values and execute operations.
470///
471/// The dispenser captures template-specific state (prepared statement,
472/// field names, bind-point indices, metrics) so the per-cycle path is
473/// minimal: bind resolved values and execute.
474///
475/// Dispensers are shared across fibers and must be thread-safe.
476pub trait OpDispenser: Send + Sync {
477    /// Execute an operation for the given cycle.
478    ///
479    /// The `ctx` bundle carries:
480    /// - `ctx.fields` — op-field substitution view for the inner adapter
481    ///   (positional / by-name access matching the prepared statement).
482    /// - `ctx.pulls`  — wrapper-facing handle-indexed view of Polydat values
483    ///   (used by validation / conditional / throttle wrappers; adapters
484    ///   ignore this).
485    ///
486    /// See SRD 32 §"`ExecCtx` — cycle-time bundle" for the design.
487    fn execute<'a>(
488        &'a self,
489        cycle: u64,
490        ctx: &'a crate::fixture::ExecCtx<'a>,
491    ) -> std::pin::Pin<
492        Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
493    >;
494
495    /// One-line description of the op this dispenser
496    /// represents — typically the statement text /
497    /// request shape with bind-point placeholders left
498    /// unresolved (`{table}`, `{key}`, `?`, …).
499    ///
500    /// Used by the runtime's error-capture path to attach
501    /// the actual op identity to a phase-stop diagnostic.
502    /// Without this, an error like `[validation_failed]
503    /// op 'indexes_present_cass' …` only names the op
504    /// template; the operator can't see what statement
505    /// the dispenser is actually firing without grepping
506    /// the workload yaml. With this hooked, the captured
507    /// reason carries the rendered statement so the
508    /// operator can immediately tell whether (a) it's the
509    /// wrong dialect's branch firing, (b) the bindpoints
510    /// resolved to something unexpected, (c) the
511    /// statement is malformed in the workload, etc.
512    ///
513    /// Returns `None` for dispensers whose op shape isn't
514    /// usefully one-lineable (composed wrappers that
515    /// delegate, dispensers whose body is multi-paragraph
516    /// HTTP, etc.). The runtime falls through to the
517    /// inner-dispenser chain when this returns None.
518    /// Default: delegate to inner dispenser if any, else
519    /// None.
520    fn describe(&self) -> Option<String> {
521        self.inner_dispenser().and_then(|inner| inner.describe())
522    }
523
524    /// Render the actual op the dispenser would fire for the given
525    /// cycle — the dryrun-equivalent view of the statement after
526    /// every bind point is interpolated.
527    ///
528    /// Pairs with [`Self::describe`]: `describe()` returns the
529    /// op-template (placeholders intact) so the operator can match
530    /// the failure to the workload yaml; `describe_resolved(wires)`
531    /// returns what was *actually* sent for this cycle, so the
532    /// operator can immediately tell whether bindpoints resolved to
533    /// expected values, whether the wrong dialect's branch fired,
534    /// whether quoting / escaping is broken, etc.
535    ///
536    /// SRD-68 Push 5: takes the dispenser's bound `WireSource`
537    /// (same surface adapters use at cycle time) so the rendered
538    /// view comes from the canonical resolution path — no
539    /// synthesis-layer ResolvedFields detour.
540    ///
541    /// Returns `None` when the dispenser can't usefully render the
542    /// resolved form (e.g. opaque request bodies, dispensers without
543    /// per-cycle interpolation). The default delegates to the inner
544    /// dispenser; leaf dispensers should override when they have a
545    /// useful per-cycle rendering.
546    fn describe_resolved(&self, wires: &dyn crate::wires::WireSource) -> Option<String> {
547        self.inner_dispenser()
548            .and_then(|inner| inner.describe_resolved(wires))
549    }
550
551    /// SRD-68 invariant I-3 — the dispenser's canonical Polydat Kernel,
552    /// established at construction by `map_op` from its parent — the
553    /// parent itself, or a kernel bound under it. Any engine: the
554    /// executor finds the program to bind per fiber by `program_id`.
555    /// Returns `None` for dispensers that don't own a kernel
556    /// (adapters with no Polydat needs, or wrappers that delegate to
557    /// an inner dispenser).
558    ///
559    /// The executor walks `Some` returns at fiber spawn to
560    /// materialise per-fiber subscope kernels — one per dispenser,
561    /// indexed parallel to the dispenser registry. At cycle time
562    /// the firing fiber's slot for this dispenser is handed in
563    /// via `ExecCtx::wires` so cycle-time reads stay on the SRD-68
564    /// I-1 single resolution surface.
565    ///
566    /// Default: delegate to inner dispenser if any, else `None`.
567    fn canonical_kernel(&self) -> Option<&std::sync::Arc<dyn Kernel>> {
568        self.inner_dispenser()
569            .and_then(|inner| inner.canonical_kernel())
570    }
571
572    /// Snapshot adapter-specific metrics for inclusion in the capture
573    /// snapshot.
574    ///
575    /// Called by the metrics scheduler alongside the standard activity
576    /// metrics. Adapters return additional `(family_name, labels,
577    /// MetricValue)` triples that represent adapter-internal state
578    /// (e.g., rows/s for batched CQL). These appear in the summary
579    /// report.
580    ///
581    /// The OpenMetrics-shaped runtime model lives under `nmbrs_metrics::snapshot`
582    /// — adapters typically build `MetricValue::Counter` /
583    /// `MetricValue::Histogram` / `MetricValue::Gauge` directly.
584    /// Default: no additional metrics.
585    fn adapter_metrics(
586        &self,
587    ) -> Vec<(
588        String,
589        nmbrs_metrics::labels::Labels,
590        nmbrs_metrics::snapshot::MetricValue,
591    )> {
592        if let Some(inner) = self.inner_dispenser() {
593            inner.adapter_metrics()
594        } else {
595            Vec::new()
596        }
597    }
598
599    /// Adapter-specific status line entries (cumulative, non-destructive read).
600    ///
601    /// Unlike `adapter_metrics()` which snapshots delta timers, this method
602    /// returns cumulative counters safe to read from the progress thread
603    /// without interfering with the metrics pipeline.
604    /// Returns `(display_name, cumulative_count)` pairs.
605    ///
606    /// A name with a **leading underscore** (`_batch_writes`) is an INTERNAL
607    /// counter: published only to back a derived display metric (the
608    /// `rows/batch` average divides `rows_inserted` by `_batch_writes`), it is
609    /// looked up by name for that computation but never rendered as its own
610    /// `<name>/s` throughput chip. See
611    /// [`crate::readout_context::is_internal_counter`] — the single predicate
612    /// every chip-rendering surface filters on.
613    /// Default: delegates to inner dispenser (for wrapper chains).
614    fn status_counters(&self) -> Vec<(&str, u64)> {
615        if let Some(inner) = self.inner_dispenser() {
616            inner.status_counters()
617        } else {
618            Vec::new()
619        }
620    }
621
622    /// The op's uniform per-invocation cursor consumption — how many
623    /// consecutive wire ordinals one `execute` call reads and covers.
624    ///
625    /// The executor drives the phase cursor with `Σ rows_per_op` over
626    /// the stanza's ops (instead of the raw stanza length) and hands
627    /// each op a contiguous sub-run of exactly this size, so a batch
628    /// op that reads `N` rows per call also advances the cursor by
629    /// `N` — consecutive stanzas then cover disjoint ordinal runs
630    /// (SRD-22 "phase extent = cursor exhaustion", cover-once). A
631    /// batch dispenser overrides this to return its fixed stride `N`;
632    /// every ordinary op keeps the default `1` (identical to the
633    /// pre-batching model).
634    ///
635    /// The default **delegates to the inner dispenser** so a wrapped
636    /// batch op (retry / result / metrics layers) still reports its
637    /// leaf stride — the executor sees the outermost wrapper, and the
638    /// stride must reach it through the chain (mirrors `describe` /
639    /// `adapter_metrics`). Leaves with no inner fall through to `1`.
640    ///
641    /// This is the NOMINAL stride (what sets `reserve(N)`); at the
642    /// cursor tail the reservation may be shorter, and the executor
643    /// passes the actual (possibly-short) run length via
644    /// [`crate::fixture::ExecCtx::run_len`] so the op inserts exactly
645    /// the ordinals reserved — never over-reading, never dropping the
646    /// remainder.
647    fn rows_per_op(&self) -> usize {
648        self.inner_dispenser()
649            .map(|inner| inner.rows_per_op())
650            .unwrap_or(1)
651    }
652
653    /// Returns the wrapped dispenser when this is a
654    /// wrapper, `None` when this is a leaf (the adapter's
655    /// base dispenser — e.g. CQL raw / prepared / batch,
656    /// HTTP, stdout, plotter, …).
657    ///
658    /// Default returns `None`, which is correct for every
659    /// leaf. Adapter-base dispensers should rely on the
660    /// default — they have no inner to expose, and asking
661    /// them to write `fn inner_dispenser(&self) -> None`
662    /// is pure boilerplate.
663    ///
664    /// **Wrappers MUST override** to return
665    /// `Some(self.inner.as_ref())`. Several pieces of
666    /// cross-cutting machinery walk this chain:
667    /// - `adapter_metrics` and `status_counters` delegate
668    ///   through wrapper layers via the default
669    ///   implementations on this trait, which call
670    ///   `inner_dispenser()` to find the wrapped layer.
671    /// - `describe()` walks inward to surface the runtime
672    ///   op shape (CQL statement text) for error-context
673    ///   dumps; missing this on a wrapper silently breaks
674    ///   the walk and the error loses its op-shape line.
675    ///
676    /// The [`WrappingDispenser`] marker trait below is
677    /// the type-system signal that flags "this is a
678    /// wrapper"; wrappers should implement both. Future
679    /// composition machinery (SRD-32a) will use the
680    /// `WrappingDispenser` bound to require the override
681    /// at the registration boundary.
682    fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
683        None
684    }
685}
686
687/// Marker trait for dispensers that wrap another
688/// `OpDispenser`. Implementing it is the type-level
689/// commitment that this layer has overridden
690/// [`OpDispenser::inner_dispenser`] to expose its inner
691/// dispenser. Without that override, cross-cutting
692/// machinery (`describe()`, `adapter_metrics`,
693/// `status_counters`) silently stops walking at the
694/// wrapper, which has been a real source of bugs.
695///
696/// Pure marker — no methods. The chain-walking code
697/// already takes `&dyn OpDispenser` and calls
698/// `inner_dispenser()`; this trait exists to make the
699/// wrapper-vs-leaf split visible at the type level and
700/// to give SRD-32a's wrapper registry a bound to require
701/// at composition time.
702///
703/// Leaves do not implement this trait — their
704/// `inner_dispenser` falls through to the default
705/// `None`-returning impl on `OpDispenser`, no boilerplate.
706pub trait WrappingDispenser: OpDispenser {}
707
708/// Resolved field values for a single cycle. Produced by the GK
709/// synthesis pipeline, consumed by the OpDispenser.
710///
711/// Fields are indexed by name. Typed values are always available;
712/// string rendering is deferred until first access to avoid wasted
713/// work for adapters that bind typed values natively (e.g., CQL).
714///
715/// SRD-68 Push 5: this struct survives as an in-adapter rendering
716/// container (stdout/testkit/plotter build one locally from
717/// [`crate::wires::resolve_op_fields_via_wires`] when they need
718/// name-keyed value access). Adapters that bind typed values
719/// natively (CQL prepared params, vector arguments) read directly
720/// through `wires.get(name)` and don't construct this type.
721pub struct ResolvedFields {
722    /// Field names in op template declaration order.
723    pub names: Vec<String>,
724    /// Typed values, parallel to `names`.
725    pub values: Vec<polydat::ast::Value>,
726    /// Lazily rendered string representations, parallel to `names`.
727    strings: OnceLock<Vec<String>>,
728}
729
730impl fmt::Debug for ResolvedFields {
731    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
732        f.debug_struct("ResolvedFields")
733            .field("names", &self.names)
734            .field("values", &self.values)
735            .finish()
736    }
737}
738
739impl Clone for ResolvedFields {
740    fn clone(&self) -> Self {
741        Self {
742            names: self.names.clone(),
743            values: self.values.clone(),
744            strings: self.strings.clone(),
745        }
746    }
747}
748
749impl ResolvedFields {
750    /// Create with names and typed values. Strings are lazily rendered.
751    pub fn new(names: Vec<String>, values: Vec<polydat::ast::Value>) -> Self {
752        Self {
753            names,
754            values,
755            strings: OnceLock::new(),
756        }
757    }
758
759    /// Access the lazily-rendered string representations.
760    /// Computed once on first call, then cached.
761    pub fn strings(&self) -> &[String] {
762        self.strings
763            .get_or_init(|| self.values.iter().map(|v| v.to_display_string()).collect())
764    }
765
766    /// Get a field value by name as a string.
767    pub fn get_str(&self, name: &str) -> Option<&str> {
768        self.names
769            .iter()
770            .position(|n| n == name)
771            .map(|i| self.strings()[i].as_str())
772    }
773
774    /// Get a field value by name as a typed Value.
775    pub fn get_value(&self, name: &str) -> Option<&polydat::ast::Value> {
776        self.names
777            .iter()
778            .position(|n| n == name)
779            .map(|i| &self.values[i])
780    }
781
782    /// Get a string by index (triggers lazy rendering if needed).
783    pub fn str_at(&self, index: usize) -> &str {
784        &self.strings()[index]
785    }
786
787    /// Return a copy with the named field removed.
788    /// Used by `ConditionalDispenser` to strip internal fields
789    /// before the adapter sees them.
790    pub fn without(&self, name: &str) -> Self {
791        let mut names = Vec::new();
792        let mut values = Vec::new();
793        for (i, n) in self.names.iter().enumerate() {
794            if n != name {
795                names.push(n.clone());
796                values.push(self.values[i].clone());
797            }
798        }
799        Self::new(names, values)
800    }
801
802    /// Serialize all fields to JSON for diagnostic/logging use.
803    pub fn to_json(&self) -> serde_json::Value {
804        let map: serde_json::Map<String, serde_json::Value> = self
805            .names
806            .iter()
807            .zip(self.values.iter())
808            .map(|(name, value)| {
809                let json_val = match value {
810                    polydat::ast::Value::U64(v) => serde_json::Value::Number((*v).into()),
811                    polydat::ast::Value::F64(v) => serde_json::Number::from_f64(*v)
812                        .map(serde_json::Value::Number)
813                        .unwrap_or(serde_json::Value::Null),
814                    polydat::ast::Value::Bool(v) => serde_json::Value::Bool(*v),
815                    _ => serde_json::Value::String(value.to_display_string()),
816                };
817                (name.clone(), json_val)
818            })
819            .collect();
820        serde_json::Value::Object(map)
821    }
822}
823
824/// Capture point declaration in an op template.
825///
826/// Parsed from `[name]`, `[source as alias]`, or `[(Type)name]` syntax.
827#[derive(Debug, Clone)]
828pub struct CaptureDecl {
829    /// The field name in the operation result.
830    pub source_name: String,
831    /// The name under which the value is stored in the capture context.
832    pub as_name: String,
833    /// Optional type qualifier for validation.
834    pub type_qualifier: Option<String>,
835}
836
837#[cfg(test)]
838mod tests {
839    use super::*;
840
841    #[test]
842    fn resolved_fields_lazy_strings() {
843        let fields = ResolvedFields::new(
844            vec!["a".into(), "b".into()],
845            vec![polydat::ast::Value::U64(42), polydat::ast::Value::F64(3.5)],
846        );
847        // Strings not computed yet
848        assert!(fields.strings.get().is_none());
849        // First access triggers rendering
850        assert_eq!(fields.get_str("a"), Some("42"));
851        assert!(fields.strings.get().is_some());
852        assert_eq!(fields.get_str("b"), Some("3.5"));
853    }
854
855    #[test]
856    fn resolved_fields_get_value() {
857        let fields = ResolvedFields::new(vec!["x".into()], vec![polydat::ast::Value::F64(3.5)]);
858        match fields.get_value("x") {
859            Some(polydat::ast::Value::F64(v)) => assert!((v - 3.5).abs() < 1e-10),
860            other => panic!("expected F64(3.5), got {other:?}"),
861        }
862        // get_value doesn't trigger string rendering
863        assert!(fields.strings.get().is_none());
864    }
865
866    #[test]
867    fn execution_error_display() {
868        let op_err = ExecutionError::Op(AdapterError {
869            error_name: "Timeout".into(),
870            message: "timed out".into(),
871            retryable: true,
872        });
873        assert!(format!("{op_err}").contains("[op]"));
874        assert!(!op_err.is_adapter_level());
875
876        let adapter_err = ExecutionError::Adapter(AdapterError {
877            error_name: "ConnectionRefused".into(),
878            message: "refused".into(),
879            retryable: false,
880        });
881        assert!(format!("{adapter_err}").contains("[adapter]"));
882        assert!(adapter_err.is_adapter_level());
883    }
884}
885
886// =========================================================================
887// Adapter Registration (inventory-based, link-time collection)
888// =========================================================================
889
890/// An adapter module's registration, submitted at link time via `inventory`.
891///
892/// Each adapter crate (nmbrs-adapter-stdout, nmbrs-adapter-cql, etc.)
893/// submits one of these. The shared runner collects all submissions to
894/// build the adapter dispatch table without any explicit adapter list.
895pub struct AdapterRegistration {
896    /// Driver names this adapter responds to (e.g., `&["stdout"]` or `&["cql", "cassandra"]`).
897    pub names: fn() -> &'static [&'static str],
898    /// Extra param names this adapter accepts (for CLI validation).
899    pub known_params: fn() -> &'static [&'static str],
900    /// Display preference for TUI activation. Checked at startup before
901    /// any adapter is constructed — no connection overhead. Takes the
902    /// run params so an adapter whose terminal use depends on config can
903    /// decide (e.g. `stdout` wants the TUI off when it writes op output
904    /// to the console, but is TUI-compatible when `filename=` redirects
905    /// to a file).
906    pub display_preference: fn(&std::collections::HashMap<String, String>) -> DisplayPreference,
907    /// SRD-23 — the dynamic controls this adapter can declare, as static
908    /// [`ControlDesc`](crate::control_catalog::ControlDesc) capability
909    /// descriptors. This is the *discovery* surface (`nmbrs describe controls`
910    /// / `describe adapter=<name>`), read without constructing the adapter;
911    /// the adapter's [`declare_controls`](DriverAdapter::declare_controls)
912    /// *derives* the live control from the same descriptors so the two cannot
913    /// drift. Default `|| &[]` for adapters with no dynamic knobs.
914    pub supported_controls: fn() -> &'static [crate::control_catalog::ControlDesc],
915    /// Async factory: given params, create the adapter.
916    /// Returns a boxed future so async connect is supported (e.g., CQL).
917    pub create: fn(std::collections::HashMap<String, String>) -> CreateAdapterFuture,
918}
919
920inventory::collect!(AdapterRegistration);
921
922/// Look up an adapter by driver name from all link-time registrations.
923pub fn find_adapter_registration(driver: &str) -> Option<&'static AdapterRegistration> {
924    inventory::iter::<AdapterRegistration>
925        .into_iter()
926        .find(|&reg| (reg.names)().contains(&driver))
927        .map(|v| v as _)
928}
929
930/// List all registered driver names.
931pub fn registered_driver_names() -> Vec<&'static str> {
932    let mut names = Vec::new();
933    for reg in inventory::iter::<AdapterRegistration> {
934        names.extend_from_slice((reg.names)());
935    }
936    names
937}
938
939/// Look up the display preference for a driver name without constructing the adapter.
940///
941/// Returns `Auto` if the driver is not registered (unknown adapters are
942/// assumed TUI-compatible; construction will fail later with a clear error).
943pub fn adapter_display_preference(
944    driver: &str,
945    params: &std::collections::HashMap<String, String>,
946) -> DisplayPreference {
947    find_adapter_registration(driver)
948        .map(|reg| (reg.display_preference)(params))
949        .unwrap_or(DisplayPreference::Auto)
950}
951
952/// Collect all extra known params from registered adapters,
953/// unioned with driver-implementation–specific params
954/// contributed via [`DriverImpl`] entries (e.g. CQL's `hosts`,
955/// `port`, `keyspace`, ...) so they don't trip the
956/// "unrecognized parameter" guard at the CLI layer.
957pub fn registered_adapter_params() -> Vec<&'static str> {
958    let mut params = Vec::new();
959    for reg in inventory::iter::<AdapterRegistration> {
960        params.extend_from_slice((reg.known_params)());
961    }
962    for entry in inventory::iter::<DriverImpl> {
963        params.extend_from_slice((entry.known_params)());
964    }
965    params
966}
967
968// =========================================================================
969// Driver implementations
970// =========================================================================
971//
972// Some adapters (e.g. `cql`) are backed by more than one
973// internal driver implementation — for CQL these are `scylla`
974// (pure Rust) and `cassandra-cpp` (Apache Cassandra C++ driver via FFI).
975// The adapter is a single user-facing concept; the driver is an
976// internal implementation detail surfaced through a per-adapter
977// selector parameter (e.g. `cqldriver=scylla`).
978
979/// One driver implementation registered for an adapter.
980///
981/// Each driver contributes a [`DriverImpl`] to the inventory.
982/// At session start the runner picks one — by user override
983/// (e.g. `cqldriver=scylla`) or by default ranking — and calls
984/// its `create` factory. The resulting [`DriverAdapter`] wears
985/// the **adapter's** name (e.g. `"cql"`), not the driver's; the
986/// driver name is internal and never appears in op-level
987/// `adapter:` fields.
988///
989/// ```ignore
990/// // In a driver module under nmbrs-adapter-cql:
991/// inventory::submit! {
992///     nmbrs_runtime::adapter::DriverImpl {
993///         adapter: "cql",
994///         driver: "scylla",
995///         default_rank: 200,
996///         create: |params| Box::pin(async move { ... }),
997///         known_params: || &["hosts", "port", ...],
998///     }
999/// }
1000/// ```
1001pub struct DriverImpl {
1002    /// The adapter this driver implements (e.g. `"cql"`).
1003    /// Matches the [`AdapterRegistration::names`] of the
1004    /// adapter the user selects with `adapter=…`.
1005    pub adapter: &'static str,
1006    /// The driver identifier. Surfaces in
1007    /// `<adapter>driver=<driver>` for user-facing selection
1008    /// (e.g. `cqldriver=scylla`). Internal — never appears as
1009    /// a top-level adapter name.
1010    pub driver: &'static str,
1011    /// Lower wins when no user override is given. Drivers pick
1012    /// their own rank — convention is to space ranks by 100s
1013    /// so future drivers can slot between them.
1014    pub default_rank: u32,
1015    /// Async factory for this driver. Same shape as
1016    /// [`AdapterRegistration::create`] — given workload params,
1017    /// connect and return a boxed `DriverAdapter` whose
1018    /// [`DriverAdapter::name`] is the *adapter* name.
1019    pub create: fn(std::collections::HashMap<String, String>) -> CreateAdapterFuture,
1020    /// Driver-specific param names — unioned into the
1021    /// adapter's known-params surface for CLI validation.
1022    pub known_params: fn() -> &'static [&'static str],
1023}
1024
1025inventory::collect!(DriverImpl);
1026
1027/// SRD-35 Push B: declares that a `(adapter, driver)` pair
1028/// supports pool-shared instances. Drivers that submit one
1029/// of these opt into the resource pool's `Shared` policy
1030/// path: instead of a fresh `Arc<dyn DriverAdapter>` per
1031/// phase, the pool caches the adapter under
1032/// [`Self::resource_key`] and reuses it across every phase
1033/// whose params produce the same key.
1034///
1035/// Drivers without a `SharedDriverRegistration` continue
1036/// to use Push A's `LegacyAdapterResource` shim, which is
1037/// `PerPhase`-isolated. Migration is opt-in per driver.
1038pub struct SharedDriverRegistration {
1039    /// Adapter name this registration applies to (e.g.
1040    /// `"cql"`). Pairs with [`Self::driver`] for lookup.
1041    pub adapter: &'static str,
1042    /// Driver identifier (e.g. `"cassandra-cpp"`,
1043    /// `"scylla"`). Distinguishes registrations when the
1044    /// same adapter has multiple driver implementations.
1045    pub driver: &'static str,
1046    /// Strongest sharing the driver supports. Default
1047    /// `Shared` for typical pool-friendly drivers; stricter
1048    /// only if the driver type can't be safely shared.
1049    pub share_capability: crate::resource_pool::ShareCapability,
1050    /// Pure function that derives the resource key from
1051    /// the adapter's params. Two phases with the same key
1052    /// share an instance under `Shared` policy. SRD-35
1053    /// §"Instance-shaping vs shell-shaping params" — only
1054    /// instance-shaping params (`hosts`, `keyspace`, …)
1055    /// belong here; per-statement and per-phase knobs MUST
1056    /// NOT.
1057    pub resource_key: fn(
1058        &std::collections::HashMap<String, String>,
1059    ) -> Result<crate::resource_pool::ResourceKey, String>,
1060}
1061
1062inventory::collect!(SharedDriverRegistration);
1063
1064/// Look up the shared registration for a `(adapter, driver)`
1065/// pair. Returns `None` when the driver hasn't migrated to
1066/// the pool-shared shape — the executor falls back to the
1067/// `LegacyAdapterResource` shim under `PerPhase` policy in
1068/// that case.
1069pub fn find_shared_driver(
1070    adapter: &str,
1071    driver: &str,
1072) -> Option<&'static SharedDriverRegistration> {
1073    inventory::iter::<SharedDriverRegistration>
1074        .into_iter()
1075        .find(|e| e.adapter == adapter && e.driver == driver)
1076}
1077
1078/// Default order of drivers registered for `adapter`, sorted by
1079/// ascending [`DriverImpl::default_rank`]. Used as the fallback
1080/// when the user doesn't set the adapter's driver-selector
1081/// parameter.
1082pub fn default_drivers(adapter: &str) -> Vec<&'static str> {
1083    let mut entries: Vec<&'static DriverImpl> = inventory::iter::<DriverImpl>
1084        .into_iter()
1085        .filter(|e| e.adapter == adapter)
1086        .collect();
1087    entries.sort_by_key(|e| e.default_rank);
1088    entries.into_iter().map(|e| e.driver).collect()
1089}
1090
1091/// Find a driver implementation by `(adapter, driver)`. Used
1092/// by the runner to pick the right factory after resolving the
1093/// driver-selector parameter.
1094pub fn find_driver(adapter: &str, driver: &str) -> Option<&'static DriverImpl> {
1095    inventory::iter::<DriverImpl>
1096        .into_iter()
1097        .find(|e| e.adapter == adapter && e.driver == driver)
1098}
1099
1100/// Union of every driver-specific known-param for `adapter`.
1101/// Surfaced in CLI validation so unknown-param warnings don't
1102/// fire for driver-private knobs (e.g. CQL's `hosts`, `port`,
1103/// `keyspace`).
1104pub fn adapter_driver_params(adapter: &str) -> Vec<&'static str> {
1105    let mut params: Vec<&'static str> = Vec::new();
1106    for e in inventory::iter::<DriverImpl>
1107        .into_iter()
1108        .filter(|e| e.adapter == adapter)
1109    {
1110        params.extend_from_slice((e.known_params)());
1111    }
1112    params
1113}
1114
1115/// Pick a driver for `adapter` (user-supplied selector or
1116/// default ranking) and instantiate it.
1117///
1118/// `selector_param` is the workload param a user sets to override
1119/// the default — e.g. `"cqldriver"` for the `cql` adapter.
1120/// Single-name semantics by convention; comma-separated lists
1121/// are accepted (the runner walks them in order and takes the
1122/// first that's compiled in).
1123/// Sentinel driver name for adapters that don't have
1124/// multiple `DriverImpl` registrations (HTTP, stdout,
1125/// testkit, openapi — anything that registers only via
1126/// [`AdapterRegistration`] with a direct `create`
1127/// factory). The shared-pool path keys
1128/// `SharedDriverRegistration` by `(adapter, driver)`;
1129/// single-engine adapters submit theirs with this sentinel
1130/// so the executor's lookup finds them after
1131/// [`resolve_driver_name`] returns the same sentinel.
1132pub const DEFAULT_DRIVER_NAME: &str = "default";
1133
1134/// Resolve the driver name for an `(adapter, params)` pair
1135/// without instantiating anything. Used by the resource
1136/// pool's Push B path: the executor needs the resolved
1137/// driver name to look up a [`SharedDriverRegistration`]
1138/// before deciding which attach path to use. Mirrors the
1139/// resolution logic in [`instantiate_with_driver`] —
1140/// user-supplied selector first (comma-separated, in
1141/// order), then the rank-sorted default list.
1142///
1143/// Returns [`DEFAULT_DRIVER_NAME`] when no `DriverImpl` is
1144/// registered for the adapter (single-engine adapters that
1145/// only register via `AdapterRegistration`). Single-engine
1146/// adapters that opt into pool sharing submit their
1147/// `SharedDriverRegistration` with the same sentinel.
1148pub fn resolve_driver_name(
1149    adapter: &str,
1150    selector_param: &str,
1151    params: &std::collections::HashMap<String, String>,
1152) -> Option<&'static str> {
1153    let user_order: Option<Vec<&str>> = params.get(selector_param).map(|s| {
1154        s.split(',')
1155            .map(str::trim)
1156            .filter(|s| !s.is_empty())
1157            .collect()
1158    });
1159    let default_order = default_drivers(adapter);
1160    let order: Vec<&str> = match &user_order {
1161        Some(v) => v.clone(),
1162        None => default_order.to_vec(),
1163    };
1164    for driver in &order {
1165        if let Some(entry) = find_driver(adapter, driver) {
1166            return Some(entry.driver);
1167        }
1168    }
1169    // Single-engine adapter (no DriverImpl registered, just
1170    // a direct AdapterRegistration). Fall back to the
1171    // sentinel so SharedDriverRegistration lookup can find
1172    // a `(adapter, "default")` entry.
1173    Some(DEFAULT_DRIVER_NAME)
1174}
1175
1176pub async fn instantiate_with_driver(
1177    adapter: &str,
1178    selector_param: &str,
1179    params: std::collections::HashMap<String, String>,
1180) -> Result<std::sync::Arc<dyn DriverAdapter>, String> {
1181    let user_order: Option<Vec<&str>> = params.get(selector_param).map(|s| {
1182        s.split(',')
1183            .map(str::trim)
1184            .filter(|s| !s.is_empty())
1185            .collect()
1186    });
1187    let default_order = default_drivers(adapter);
1188    let order: Vec<&str> = match &user_order {
1189        Some(v) => v.clone(),
1190        None => default_order.clone(),
1191    };
1192    if order.is_empty() {
1193        return Err(format!(
1194            "adapter '{adapter}': no driver implementations registered. \
1195             Build the binary with at least one driver feature enabled."
1196        ));
1197    }
1198    for driver in &order {
1199        if let Some(entry) = find_driver(adapter, driver) {
1200            return (entry.create)(params).await;
1201        }
1202    }
1203    Err(format!(
1204        "adapter '{adapter}': no driver in {selector_param}='{}' is compiled in; \
1205         available drivers: [{}]",
1206        order.join(","),
1207        default_order.join(", "),
1208    ))
1209}