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.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}