skardi 0.6.0

High performance query engine for both offline compute and online serving
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
//! The backend abstraction (design §Backend abstraction) and its
//! milestone-1 implementation: [`AgeClient`] — Apache AGE, openCypher
//! inside Postgres.
//!
//! One deliberate deviation from the design text, recorded here: the
//! design named `tokio-postgres`; this implementation rides **sqlx**,
//! because the workspace already ships sqlx-postgres for the relational
//! Postgres provider — one Postgres stack beats two (same TLS, same
//! pooling, same error surface).
//!
//! Read-only is backend-enforced: every call runs inside a Postgres
//! `READ ONLY` transaction — the keyword guard upstream is UX, this is
//! the boundary. Every call is bounded: `statement_timeout` server-side
//! plus a client-side wrap, and a row cap enforced while streaming, so a
//! whole-graph `RETURN` fails loudly at `max_rows + 1` instead of
//! buffering the world.

use std::io::ErrorKind as IoErrorKind;
use std::time::Duration;

use std::sync::atomic::{AtomicU64, Ordering};

use async_trait::async_trait;
use futures::StreamExt;
use futures::stream::BoxStream;
use serde_json::Value;
use sqlx::postgres::{PgConnectOptions, PgPool, PgPoolOptions};
use sqlx::{Connection, Either, Executor, Row};

use super::error::{GraphError, json_kind};
use super::value::parse_agtype;

/// Per-query operational bounds (design §Security and operational
/// bounds), carried per source.
#[derive(Debug, Clone, Copy)]
pub struct QueryBounds {
    pub timeout: Duration,
    pub max_rows: usize,
}

/// One result row: one JSON value per declared column.
pub type GraphRow = Vec<Value>;
/// The stream `execute` returns. Milestone 1 implementations may buffer
/// internally up to `max_rows` — the trait shape is what must not change
/// (design §Backend abstraction).
pub type GraphRowStream = BoxStream<'static, Result<GraphRow, GraphError>>;

/// What a backend must provide. Hides AGE (Postgres wire) vs Neo4j Bolt
/// vs Kuzu details behind one seam.
///
/// # Example
/// ```
/// use async_trait::async_trait;
/// use futures::StreamExt;
/// use serde_json::Value;
/// use skardi::sources::providers::graph::client::{
///     GraphClient, GraphRowStream, QueryBounds,
/// };
/// use skardi::sources::providers::graph::error::GraphError;
///
/// /// A canned backend, the shape tests use.
/// #[derive(Debug)]
/// struct Fixed(Vec<Vec<Value>>);
///
/// #[async_trait]
/// impl GraphClient for Fixed {
///     async fn execute(
///         &self,
///         _cypher: &str,
///         _params: &Value,
///         _arity: usize,
///         _bounds: QueryBounds,
///         limit: Option<usize>,
///     ) -> Result<GraphRowStream, GraphError> {
///         let mut rows = self.0.clone();
///         if let Some(l) = limit {
///             rows.truncate(l);
///         }
///         Ok(futures::stream::iter(rows.into_iter().map(Ok)).boxed())
///     }
///
///     async fn labels(
///         &self,
///         _bounds: QueryBounds,
///         _limit: Option<usize>,
///     ) -> Result<Vec<(String, String)>, GraphError> {
///         Ok(vec![("Person".into(), "vertex".into())])
///     }
/// }
/// ```
#[async_trait]
pub trait GraphClient: Send + Sync + std::fmt::Debug {
    /// Run read-only Cypher inside a backend-enforced read transaction,
    /// bounded by `bounds`. `arity` is the declared column count — AGE's
    /// `cypher()` call must declare its result arity, which is why the
    /// declared-columns mode is required on this backend. `limit` is
    /// pushed BOTH ways: into the statement as a real SQL LIMIT
    /// (min(limit, max_rows + 1), so the backend and the wire are
    /// bounded even when the caller passes none — an overflow costs
    /// max_rows + 1 rows, not a full scan plus a drain) AND enforced at
    /// consumption as defense in depth. Hitting `limit` is a clean early
    /// stop; the row cap stays the loud overflow signal.
    async fn execute(
        &self,
        cypher: &str,
        params: &Value,
        arity: usize,
        bounds: QueryBounds,
        limit: Option<usize>,
    ) -> Result<GraphRowStream, GraphError>;

    /// Label roster for `graph_schema`: one `(label, kind)` row per
    /// label, kinds `vertex` / `edge`. Names only — never property
    /// values (design §Agent and LLM interaction). Bounded like every
    /// other query — "every query is bounded" has no catalog exemption —
    /// and `limit` is the SQL LIMIT, pushed into the catalog fetch as a
    /// clean early stop.
    async fn labels(
        &self,
        bounds: QueryBounds,
        limit: Option<usize>,
    ) -> Result<Vec<(String, String)>, GraphError>;
}

/// Cancellation safety for the hand-rolled transaction: [`AgeClient::execute`]
/// manages BEGIN/ROLLBACK manually (rollback must precede DEALLOCATE on
/// the same session — sqlx's `Transaction` guard cannot order that), and
/// the cost of going manual is that NOTHING rolls back if the scan
/// future is dropped mid-transaction — DataFusion drops scan streams on
/// client disconnect, plan short-circuits, and sibling-partition errors,
/// and sqlx's pool return path never rolls back (it pings). A connection
/// returned `idle in transaction` mostly self-heals on reuse, but it
/// holds a snapshot while idle and a cancellation burst can pin every
/// pooled session. This guard closes it: if dropped while ARMED, the
/// connection is moved into a spawned task that rolls back AND
/// deallocates the call's prepared statement, and only then returns it
/// to the pool; the clean path defuses and runs its ordered cleanup
/// inline. The DEALLOCATE matters as much as the rollback: the guard is
/// what RESCUES the connection back into the pool, so without it every
/// cancelled parameterized query would strand one session-level
/// `skq_p_*` statement permanently — monotonic growth on long-lived
/// sessions, the exact failure the live leak sweep exists to prevent.
struct OpenTxnGuard {
    conn: Option<sqlx::pool::PoolConnection<sqlx::Postgres>>,
    /// The per-call PREPARE name, armed as soon as it is known.
    prepared: Option<String>,
}

impl OpenTxnGuard {
    fn conn(&mut self) -> &mut sqlx::PgConnection {
        self.conn.as_mut().expect("armed until defused")
    }

    fn defuse(mut self) -> sqlx::pool::PoolConnection<sqlx::Postgres> {
        self.conn.take().expect("defused exactly once")
    }
}

impl Drop for OpenTxnGuard {
    fn drop(&mut self) {
        if let Some(mut conn) = self.conn.take() {
            let prepared = self.prepared.take();
            // Drop is sync; the cleanup needs the runtime. Execution
            // always runs inside tokio here — the fallback (no runtime:
            // plain drop, sqlx pings a possibly-in-transaction session
            // back into the pool) is only reachable from exotic test
            // harnesses.
            if let Ok(handle) = tokio::runtime::Handle::try_current() {
                handle.spawn(async move {
                    // NOTE: this ROLLBACK queues BEHIND the statement
                    // still running on this session, so the connection
                    // does not return to the pool until the server-side
                    // statement_timeout expires (up to
                    // query_timeout_seconds — measured ~23s in review).
                    // A burst of cancellations can therefore pin every
                    // pooled session for that long, with new queries
                    // queueing on acquire_timeout — bounded, but not
                    // prompt. Cancelling the running query first
                    // (pg_cancel_backend from a second connection) is
                    // the fix if that ever bites.
                    //
                    // ROLLBACK first (clears an aborted transaction, where
                    // DEALLOCATE is refused), then drop the statement —
                    // the same order the clean path uses.
                    let _ = conn.execute("ROLLBACK").await;
                    if let Some(name) = prepared {
                        let _ = conn.execute(format!("DEALLOCATE {name}").as_str()).await;
                    }
                });
            }
        }
    }
}

/// Ceiling on the REGISTRATION preflight (dial + per-connection setup +
/// the graph-existence probe). `query_timeout_seconds` means "how long
/// may a traversal run" — config permits a day — and must not double as
/// "how long may server startup block on a dead host": the preflight is
/// three tiny statements, so 30s is generous. The trade this ceiling
/// makes explicit: a backend that WOULD answer slower than the bound
/// classifies Unavailable and boots degraded (the first scan carries the
/// real diagnosis) — acceptable for a probe this small, unlike the
/// serial multi-minute boot stall the uncapped bound allowed. The pool's
/// own acquire_timeout stays at the full query bound.
const PREFLIGHT_TIMEOUT_CEILING: Duration = Duration::from_secs(30);

/// Apache AGE over the workspace's sqlx-postgres stack.
#[derive(Debug)]
pub struct AgeClient {
    pool: PgPool,
    graph_name: String,
    source_name: String,
    /// Uniquifies per-call PREPARE names (see [`Self::build_sql`]) so
    /// pooled-connection reuse can never collide.
    prepare_seq: AtomicU64,
}

impl AgeClient {
    /// Connect a pool whose every connection has AGE loaded and
    /// `ag_catalog` on the search path. Credentials arrive as env-var
    /// NAMES (values never in YAML); when set they override whatever the
    /// URL carries.
    pub async fn connect(
        source_name: &str,
        connection_string: &str,
        graph_name: &str,
        username_env: Option<&str>,
        password_env: Option<&str>,
        max_connections: u32,
        acquire_timeout: Duration,
    ) -> Result<Self, GraphError> {
        let options = build_options(source_name, connection_string, username_env, password_env)?;
        // PREFLIGHT on a single direct connection, before any pool
        // exists: sqlx treats an `after_connect` error as a failed
        // attempt and RETRIES until acquire_timeout, discarding the
        // underlying cause — a bad search_path or a missing AGE would
        // surface as a 30-second "pool timed out" pointing at pool
        // sizing. Running the same per-connection setup here first means
        // auth failures, setup failures, and the graph probe all fail
        // FAST with their real, named error (the eager-fail contract in
        // mod.rs's docstring). The whole preflight is bounded by
        // min(acquire_timeout, PREFLIGHT_TIMEOUT_CEILING): a blackholed
        // address or a stuck TLS/auth handshake must not hold startup
        // hostage — it degrades like any other unreachable backend — and
        // an operator raising query_timeout_seconds for heavy traversals
        // must not thereby raise the boot stall (see the ceiling's doc).
        let preflight = async {
            let mut conn = sqlx::postgres::PgConnection::connect_with(&options)
                .await
                .map_err(|e| connect_error(source_name, &e))?;
            per_connection_setup(&mut conn)
                .await
                .map_err(|e| backend_error(source_name, &e))?;
            // The graph name gets the same eager treatment as the URL
            // and the credential: a typo would otherwise fail LATE and
            // split — per-query raw backend errors on cypher_query, and
            // zero rows WITHOUT error on graph_schema (the catalog join
            // just misses), which an agent reads as "this graph is
            // empty".
            let exists: Option<i32> =
                sqlx::query_scalar("SELECT 1 FROM ag_catalog.ag_graph WHERE name = $1")
                    .bind(graph_name)
                    .fetch_optional(&mut conn)
                    .await
                    .map_err(|e| backend_error(source_name, &e))?;
            if exists.is_none() {
                return Err(GraphError::InvalidConfig {
                    name: source_name.to_string(),
                    reason: format!(
                        "graph '{graph_name}' does not exist in this database \
                         (ag_catalog.ag_graph has no such entry; create it with \
                         SELECT create_graph(...) or fix graph_name)"
                    ),
                });
            }
            Ok::<_, GraphError>(())
        };
        let preflight_timeout = acquire_timeout.min(PREFLIGHT_TIMEOUT_CEILING);
        tokio::time::timeout(preflight_timeout, preflight)
            .await
            .map_err(|_| GraphError::Unavailable {
                source_name: source_name.to_string(),
                reason: format!(
                    "the registration preflight did not complete within {}s",
                    preflight_timeout.as_secs()
                ),
            })??;
        Ok(Self::from_pool(
            build_pool(options, max_connections, acquire_timeout),
            source_name,
            graph_name,
        ))
    }

    /// Connect WITHOUT the preflight probe: URL parsing and env-var
    /// resolution still fail here (those are config errors, never
    /// transient), but no network I/O happens — the lazy pool dials on
    /// first acquire. Used by degraded registration
    /// (`register_graph_tables`): an unreachable shared backend must not
    /// take server startup — and every unrelated source — down with it;
    /// the first scan surfaces the real backend error instead.
    pub fn connect_degraded(
        source_name: &str,
        connection_string: &str,
        graph_name: &str,
        username_env: Option<&str>,
        password_env: Option<&str>,
        max_connections: u32,
        acquire_timeout: Duration,
    ) -> Result<Self, GraphError> {
        let options = build_options(source_name, connection_string, username_env, password_env)?;
        Ok(Self::from_pool(
            build_pool(options, max_connections, acquire_timeout),
            source_name,
            graph_name,
        ))
    }

    fn from_pool(pool: PgPool, source_name: &str, graph_name: &str) -> Self {
        tracing::debug!(
            source = source_name,
            graph = graph_name,
            "graph source connected"
        );
        Self {
            pool,
            graph_name: graph_name.to_string(),
            source_name: source_name.to_string(),
            prepare_seq: AtomicU64::new(0),
        }
    }

    /// The pool, for the live leak-regression test ONLY (the
    /// `pg_prepared_statements` sweep is session-local, so it must run
    /// on the very sessions `execute` used). Not API.
    #[doc(hidden)]
    pub fn pool_for_tests(&self) -> &PgPool {
        &self.pool
    }
}

/// URL parse + credential env resolution — PURE (no network I/O), so its
/// errors are config errors and hard-fail in every connect path,
/// degraded included.
fn build_options(
    source_name: &str,
    connection_string: &str,
    username_env: Option<&str>,
    password_env: Option<&str>,
) -> Result<PgConnectOptions, GraphError> {
    let mut options: PgConnectOptions =
        connection_string
            .parse()
            .map_err(|e: sqlx::Error| GraphError::InvalidConfig {
                name: source_name.to_string(),
                reason: format!("connection_string does not parse: {e}"),
            })?;
    if let Some(env) = username_env {
        let user = read_env(source_name, "username_env", env)?;
        options = options.username(&user);
    }
    if let Some(env) = password_env {
        let pass = read_env(source_name, "password_env", env)?;
        options = options.password(&pass);
    }
    Ok(options)
}

/// The lazy pool shared by both connect paths: connections dial on
/// demand, each running [`per_connection_setup`] via `after_connect`.
fn build_pool(
    options: PgConnectOptions,
    max_connections: u32,
    acquire_timeout: Duration,
) -> PgPool {
    PgPoolOptions::new()
        .max_connections(max_connections)
        // Queueing on a saturated pool is bounded by the SAME knob
        // as the query itself — sqlx's 30s default is unrelated to
        // query_timeout_seconds and would surface as a generic
        // driver failure after possibly LONGER than the configured
        // bound ("every query is bounded" covers the queue too).
        .acquire_timeout(acquire_timeout)
        .after_connect(|conn, _meta| Box::pin(per_connection_setup(conn)))
        .connect_lazy_with(options)
}

/// Per-connection session state, shared by the registration preflight
/// and the pool's `after_connect` hook so the two can never drift.
///
/// `LOAD 'age'` is BEST-EFFORT, deliberately: Postgres restricts LOAD
/// of a library outside `$libdir/plugins` to superusers (AGE installs
/// to `$libdir/age.so`), so requiring it would force every graph source
/// onto a superuser credential — the exact opposite of the design's
/// least-privilege recommendation, and it would put the module's ONLY
/// enforcing layer (the backend READ ONLY transaction) behind maximum
/// privilege. The supported deployment (the official apache/age image)
/// ships `shared_preload_libraries = age`, where the LOAD is a no-op;
/// where AGE is genuinely absent, the registration preflight's
/// `ag_catalog.ag_graph` probe is what fails, with a named error. The
/// search_path, by contrast, is required state — its failure is real.
async fn per_connection_setup(conn: &mut sqlx::PgConnection) -> Result<(), sqlx::Error> {
    let _ = conn.execute("LOAD 'age'").await;
    conn.execute("SET search_path = ag_catalog, \"$user\", public")
        .await?;
    Ok(())
}

impl AgeClient {
    /// The `cypher()` invocation, spoken over the SIMPLE query protocol —
    /// three AGE realities verified live make this the one sound
    /// spelling:
    ///
    /// - `cypher()`'s params argument must be a prepared-statement
    ///   parameter (AGE rejects a constant with "third argument … must be
    ///   a parameter"), so a parameterized call rides the documented
    ///   `PREPARE name(agtype) AS …; EXECUTE name('…');` pattern — with a
    ///   per-call unique name so pooled connections never collide.
    /// - `::text` on agtype is AGE's scalar-only cast (a vertex fails
    ///   with "unsupported argument agtype"), and for scalars it strips
    ///   JSON quoting. The simple protocol instead delivers every column
    ///   through `agtype_out` — the uniform annotated-JSON text the
    ///   decoder expects.
    /// - Everything is constant SQL text: the graph name is
    ///   identifier-validated at config load AND quote-escaped (belt and
    ///   braces); the caller's Cypher rides a dollar-quoted string whose
    ///   tag provably does not occur in it; params are serde_json-encoded
    ///   and single-quote-escaped — values cannot break out of the
    ///   serialization.
    ///
    /// Returns the statement batch and the prepared name to DEALLOCATE
    /// (when params were bound).
    fn build_sql(
        &self,
        cypher: &str,
        params: &Value,
        arity: usize,
        fetch: usize,
    ) -> (String, Option<String>) {
        build_cypher_sql(
            &self.graph_name,
            self.prepare_seq.fetch_add(1, Ordering::Relaxed),
            cypher,
            params,
            arity,
            fetch,
        )
    }
}

#[async_trait]
impl GraphClient for AgeClient {
    async fn execute(
        &self,
        cypher: &str,
        params: &Value,
        arity: usize,
        bounds: QueryBounds,
        limit: Option<usize>,
    ) -> Result<GraphRowStream, GraphError> {
        // The fetch bound rides the OUTER statement as a real SQL LIMIT
        // (never inside the Cypher text): the backend and the wire are
        // bounded even when the caller passes no limit, a RowCapExceeded
        // costs max_rows + 1 rows instead of a full scan plus a full
        // drain before ROLLBACK, and the consumption-side checks below
        // stay as defense in depth. Same formula as labels(): a SQL
        // LIMIT at or under the cap is a clean stop; otherwise fetch one
        // past the cap to PROVE an overflow.
        let fetch = match limit {
            Some(l) if l <= bounds.max_rows => l,
            _ => bounds.max_rows.saturating_add(1),
        };
        let (sql, prepared) = self.build_sql(cypher, params, arity, fetch);
        let source = self.source_name.clone();
        let started = std::time::Instant::now();
        // Acquire OUTSIDE the client-side timeout window: acquire and
        // execution are sequential costs, and acquire_timeout already
        // bounds the queue with its own typed PoolTimedOut mapping.
        // Sharing one window would let 25s of pool contention leave a
        // 30s statement 10s of budget — reported as the query's timeout
        // with the server-side bound (the authoritative one) mostly
        // unspent.
        let acquired = self
            .pool
            .acquire()
            .await
            .map_err(|e| map_query_error(&source, bounds, &e))?;
        let run = async {
            // A MANUAL transaction on an explicitly acquired connection,
            // not sqlx's Transaction guard: cleanup must be able to
            // ROLLBACK FIRST and then DEALLOCATE on the SAME session — a
            // backend error (invalid Cypher, runtime error,
            // statement_timeout) puts the transaction in aborted state,
            // where DEALLOCATE itself is refused, while rollback does NOT
            // clear session-level prepared statements. Rollback-then-
            // deallocate is the only order that cleans up after backend
            // errors. The [`OpenTxnGuard`] covers what MANUAL cannot: a
            // scan future dropped mid-transaction still rolls back AND
            // deallocates before the connection re-enters the pool.
            let mut guard = OpenTxnGuard {
                conn: Some(acquired),
                prepared: prepared.clone(),
            };
            // The security boundary: the SERVER refuses writes in a read
            // transaction, whatever slipped past the keyword guard.
            guard
                .conn()
                .execute("BEGIN TRANSACTION READ ONLY")
                .await
                .map_err(|e| backend_error(&source, &e))?;
            let collected: Result<Vec<GraphRow>, GraphError> = async {
                let conn = guard.conn();
                // Server-side: runaway traversals die in the backend, not
                // in a client that gave up waiting.
                conn.execute(
                    format!(
                        "SET LOCAL statement_timeout = '{}ms'",
                        bounds.timeout.as_millis()
                    )
                    .as_str(),
                )
                .await
                .map_err(|e| backend_error(&source, &e))?;

                // Stream (simple protocol — see build_sql) and cap:
                // max_rows + 1 proves the overflow without buffering past
                // it.
                let mut rows: Vec<GraphRow> = Vec::new();
                let mut stream = conn.fetch_many(sqlx::raw_sql(&sql));
                while let Some(step) = stream.next().await {
                    let step = step.map_err(|e| map_query_error(&source, bounds, &e))?;
                    // PREPARE contributes a rowless result; only rows count.
                    let Either::Right(row) = step else { continue };
                    // A SQL LIMIT is a clean early stop — enough rows is
                    // success, unlike the cap below.
                    if limit.is_some_and(|l| rows.len() >= l) {
                        break;
                    }
                    if rows.len() >= bounds.max_rows {
                        return Err(GraphError::RowCapExceeded {
                            max_rows: bounds.max_rows,
                        });
                    }
                    let row_idx = rows.len();
                    let mut values = Vec::with_capacity(arity);
                    for col in 0..arity {
                        // Unchecked: simple-protocol cells are agtype_out
                        // text, whose OID sqlx cannot map to String.
                        let text: Option<String> = row
                            .try_get_unchecked(col)
                            .map_err(|e| backend_error(&source, &e))?;
                        let value = match text {
                            None => Value::Null,
                            Some(t) => {
                                parse_agtype(&t).map_err(|reason| GraphError::MalformedCell {
                                    row: row_idx,
                                    column: col,
                                    reason,
                                })?
                            }
                        };
                        values.push(value);
                    }
                    rows.push(values);
                }
                Ok(rows)
            }
            .await;
            // The clean path reclaims the connection — from here the
            // guard's drop no longer fires and cleanup runs inline, in
            // its required order.
            let mut conn = guard.defuse();
            // ROLLBACK unconditionally and FIRST: it clears an aborted
            // transaction (making the DEALLOCATE below executable) and a
            // read-only transaction had nothing to undo anyway. Its
            // RESULT is checked LAST — an early `?` here would skip the
            // DEALLOCATE below and leak the statement on the very success
            // path this cleanup exists for.
            let rollback = conn.execute("ROLLBACK").await;
            let mut dealloc = Ok(Default::default());
            if let Some(name) = &prepared {
                // On EVERY exit: prepared statements are SESSION-level,
                // the pid+seq names never reuse, and skipping this on an
                // error path would monotonically accumulate statements on
                // the pooled connection. Best-effort on the error path
                // (the original error stays the one reported — and when
                // the PREPARE itself failed there is nothing to drop); a
                // failed DEALLOCATE on the success path is itself an
                // error. Targeted, never ALL — sqlx's own statement cache
                // lives on the same connection.
                dealloc = conn.execute(format!("DEALLOCATE {name}").as_str()).await;
            }
            let rows = collected?;
            // Success path only, and only AFTER both cleanups ran:
            // rollback failure first (it happened first), then dealloc.
            rollback.map_err(|e| backend_error(&source, &e))?;
            dealloc.map_err(|e| backend_error(&source, &e))?;
            Ok(rows)
        };
        // The client-side timeout (and any drop of the scan future)
        // ABANDONS the run mid-flight; [`OpenTxnGuard`] then rolls the
        // open transaction back AND deallocates this call's prepared
        // statement before the connection re-enters the pool — the
        // guard is what RESCUES the session, so it must also be what
        // cleans it. The wrap covers only the transaction body (acquire
        // happened above, bounded by its own acquire_timeout), so the
        // server-side statement_timeout keeps its full budget and stays
        // the authoritative bound; this one covers a backend that stops
        // answering entirely — which is why it maps to BackendSilent,
        // NOT Timeout: a 57014 is the server answering, silence is an
        // availability failure, and the registration/recovery
        // classification keys on that distinction. saturating_add: the
        // config caps the timeout, but arithmetic here must not be the
        // thing that panics if that invariant ever moves.
        let rows = tokio::time::timeout(bounds.timeout.saturating_add(Duration::from_secs(5)), run)
            .await
            .map_err(|_| GraphError::BackendSilent {
                seconds: bounds
                    .timeout
                    .saturating_add(Duration::from_secs(5))
                    .as_secs(),
            })??;
        tracing::debug!(
            source = %self.source_name,
            elapsed_ms = started.elapsed().as_millis() as u64,
            rows = rows.len(),
            "cypher_query executed"
        );
        Ok(futures::stream::iter(rows.into_iter().map(Ok)).boxed())
    }

    async fn labels(
        &self,
        bounds: QueryBounds,
        limit: Option<usize>,
    ) -> Result<Vec<(String, String)>, GraphError> {
        // Same bounds discipline as execute — a catalog read against a
        // wedged backend must not hang forever, and the design's "every
        // query is bounded" makes no catalog exemption: READ ONLY
        // transaction, server-side statement_timeout, row cap, and the
        // client-side wrap. The cap and the SQL LIMIT both ride the
        // query's own LIMIT clause, so at most max_rows + 1 label rows
        // ever cross the wire — never the whole catalog first.
        let fetch = match limit {
            // A SQL LIMIT at or under the cap is a clean early stop.
            Some(l) if l <= bounds.max_rows => l,
            // Otherwise fetch one past the cap to PROVE an overflow.
            _ => bounds.max_rows.saturating_add(1),
        };
        let source = self.source_name.clone();
        // Acquire OUTSIDE the client-side timeout window — the same
        // discipline (and rationale) as execute(): acquire and execution
        // are sequential costs with their own typed bound
        // (ConnectionAcquireTimeout); sharing one window would let pool
        // contention eat the catalog read's budget and misreport the
        // wait as the query's timeout. map_query_error, not
        // backend_error: a saturated or unreachable pool must surface
        // from graph_schema as the same typed ConnectionAcquireTimeout
        // cypher_query reports — and the availability classification
        // keys on the typed variant.
        let mut acquired = self
            .pool
            .acquire()
            .await
            .map_err(|e| map_query_error(&source, bounds, &e))?;
        let run = async {
            let mut tx = acquired
                .begin()
                .await
                .map_err(|e| backend_error(&source, &e))?;
            // Same protocol discipline as execute(): the SETs ride the
            // SIMPLE protocol as constant text. `sqlx::query` would take
            // the extended protocol and plant a server-side prepared
            // statement in the per-connection cache — keyed on the
            // interpolated text, so each distinct timeout value would
            // cache another copy on every pooled session.
            (&mut *tx)
                .execute("SET TRANSACTION READ ONLY")
                .await
                .map_err(|e| backend_error(&source, &e))?;
            (&mut *tx)
                .execute(
                    format!(
                        "SET LOCAL statement_timeout = '{}ms'",
                        bounds.timeout.as_millis()
                    )
                    .as_str(),
                )
                .await
                .map_err(|e| backend_error(&source, &e))?;
            // ag_catalog is the exact source for labels; AGE's own
            // catch-all labels (`_ag_label_vertex` / `_ag_label_edge`)
            // are implementation noise for an agent and are filtered.
            let rows = sqlx::query(
                "SELECT l.name, l.kind::text \
                 FROM ag_catalog.ag_label l \
                 JOIN ag_catalog.ag_graph g ON g.graphid = l.graph \
                 WHERE g.name = $1 AND l.name NOT LIKE '\\_ag\\_label\\_%' \
                 ORDER BY l.name LIMIT $2",
            )
            .bind(&self.graph_name)
            .bind(i64::try_from(fetch).unwrap_or(i64::MAX))
            .fetch_all(&mut *tx)
            .await
            .map_err(|e| map_query_error(&source, bounds, &e))?;
            if limit.is_none_or(|l| l > bounds.max_rows) && rows.len() > bounds.max_rows {
                return Err(GraphError::RowCapExceeded {
                    max_rows: bounds.max_rows,
                });
            }
            tx.rollback()
                .await
                .map_err(|e| backend_error(&source, &e))?;
            tracing::debug!(source = %source, labels = rows.len(), "graph_schema catalog read");
            rows.into_iter()
                .map(|row| {
                    let name: String = row.try_get(0).map_err(|e| backend_error(&source, &e))?;
                    let kind: String = row.try_get(1).map_err(|e| backend_error(&source, &e))?;
                    let kind = match kind.as_str() {
                        "v" => "vertex".to_string(),
                        "e" => "edge".to_string(),
                        other => other.to_string(),
                    };
                    Ok((name, kind))
                })
                .collect()
        };
        // BackendSilent, not Timeout — same reasoning as execute()'s
        // wrap: the server answering slowly is 57014; silence is an
        // availability failure.
        tokio::time::timeout(bounds.timeout.saturating_add(Duration::from_secs(5)), run)
            .await
            .map_err(|_| GraphError::BackendSilent {
                seconds: bounds
                    .timeout
                    .saturating_add(Duration::from_secs(5))
                    .as_secs(),
            })?
    }
}

/// Render the `cypher()` invocation batch (see [`AgeClient::build_sql`]
/// for the protocol constraints that shape it). A free function so the
/// exact SQL text — quote doubling, dollar tags, arity columns, the
/// PREPARE/EXECUTE split — is pinned by unit tests without a pool.
fn build_cypher_sql(
    graph_name: &str,
    seq: u64,
    cypher: &str,
    params: &Value,
    arity: usize,
    fetch: usize,
) -> (String, Option<String>) {
    let tag = dollar_tag(cypher);
    let graph = graph_name.replace('\'', "''");
    let cols: Vec<String> = (0..arity)
        .map(|i| format!("c{i} ag_catalog.agtype"))
        .collect();
    let outs: Vec<String> = (0..arity).map(|i| format!("c{i}")).collect();
    let has_params = params.as_object().is_some_and(|m| !m.is_empty());
    if has_params {
        let name = format!("skq_p_{}_{seq}", std::process::id());
        let literal = params.to_string().replace('\'', "''");
        let batch = format!(
            "PREPARE {name}(ag_catalog.agtype) AS \
             SELECT {} FROM ag_catalog.cypher('{graph}', {tag}{cypher}{tag}, $1) \
             AS t({}) LIMIT {fetch}; \
             EXECUTE {name}('{literal}');",
            outs.join(", "),
            cols.join(", ")
        );
        (batch, Some(name))
    } else {
        (
            format!(
                "SELECT {} FROM ag_catalog.cypher('{graph}', {tag}{cypher}{tag}) \
                 AS t({}) LIMIT {fetch}",
                outs.join(", "),
                cols.join(", ")
            ),
            None,
        )
    }
}

/// A dollar-quote tag that provably does not occur in `text`: pick the
/// first `$skqN$` suffix the text does not contain — no escape rules to
/// get wrong.
fn dollar_tag(text: &str) -> String {
    // The collision test runs against `text + "$"`, not `text`: the
    // closing delimiter starts with `$`, so a query ENDING in the tag's
    // interior (e.g. `RETURN $skq` with a parameter named `skq`) forms
    // the full tag at the text/closing-tag boundary and would close the
    // literal early. Appending the `$` makes the check cover exactly the
    // stream Postgres scans.
    //
    // ONE pass: collect the tag-shaped substrings actually present, then
    // take the first free suffix — a probe stuffed with `$skq$ $skq1$ …`
    // costs one scan, not one scan per collision.
    let probe = format!("{text}$");
    let mut base_used = false;
    let mut used: std::collections::HashSet<u32> = std::collections::HashSet::new();
    for (i, _) in probe.match_indices("$skq") {
        let rest = &probe[i + "$skq".len()..];
        if let Some(end) = rest.find('$') {
            let digits = &rest[..end];
            if digits.is_empty() {
                base_used = true;
            } else if let Ok(n) = digits.parse::<u32>() {
                used.insert(n);
            }
        }
    }
    if !base_used {
        return "$skq$".to_string();
    }
    let mut n = 1u32;
    while used.contains(&n) {
        n += 1;
    }
    format!("$skq{n}$")
}

fn read_env(source: &str, field: &str, env: &str) -> Result<String, GraphError> {
    std::env::var(env).map_err(|_| GraphError::InvalidConfig {
        name: source.to_string(),
        reason: format!("{field}: ${env} is not set in this environment"),
    })
}

fn backend_error(source: &str, e: &sqlx::Error) -> GraphError {
    // db.message(), never e.to_string(), for database errors: Postgres
    // errors carry `position` and statement context — the caller's
    // Cypher, i.e. values — and sqlx's Display dropping them today is
    // an implementation detail of sqlx, not a guarantee of ours. Taking
    // message() makes the "query text never flows into errors" rule
    // THIS module's own (the 300-byte cap remains as backstop only).
    let (code, message) = match e {
        sqlx::Error::Database(db) => (db.code().map(|c| c.to_string()), db.message().to_string()),
        _ => (None, e.to_string()),
    };
    GraphError::backend(source, code.as_deref().unwrap_or("io"), &message)
}

/// Map a DIAL failure (the registration preflight's `connect_with`).
/// Only a transport-level `io` error that says "no server answered"
/// qualifies for degraded registration. A `Database` error here is the
/// SERVER answering (bad credentials, `28P01`), and a TLS failure is a
/// client configuration problem; both stay [`backend_error`] so they
/// hard-fail registration instead of degrading.
///
/// The `Io` variant is NOT uniformly "nobody answered": with this
/// workspace's `tls-rustls` feature the TLS handshake runs through
/// `RustlsSocket::complete_io()`, and rustls wraps TLS protocol errors —
/// certificate-verification failures included — as
/// `io::Error(InvalidData)`, which sqlx's `handshake()` propagates into
/// `Error::Io` (`sqlx::Error::Tls` covers only pre-handshake
/// configuration failures: bad CA file, invalid hostname, client-cert
/// setup). `InvalidData` therefore means we WERE talking to something
/// whose TLS setup is broken — a misconfiguration that must refuse boot,
/// exactly like a bad credential, never sit degraded.
fn connect_error(source: &str, e: &sqlx::Error) -> GraphError {
    match e {
        sqlx::Error::Io(io) if io.kind() != IoErrorKind::InvalidData => GraphError::Unavailable {
            source_name: source.to_string(),
            reason: io.to_string(),
        },
        other => backend_error(source, other),
    }
}

/// Postgres cancels a statement-timeout overrun with SQLSTATE 57014 —
/// surface it as the typed timeout, not a generic backend error.
fn map_query_error(source: &str, bounds: QueryBounds, e: &sqlx::Error) -> GraphError {
    if let sqlx::Error::Database(db) = e
        && db.code().as_deref() == Some("57014")
    {
        return GraphError::Timeout {
            seconds: bounds.timeout.as_secs(),
        };
    }
    // A pool-acquire timeout is NOT the statement timeout: sqlx retries
    // a refused dial until the acquire deadline, so an unreachable
    // backend surfaces as PoolTimedOut and the query never ran. The
    // honest variant keeps "narrow the traversal" off an error it
    // cannot help.
    if matches!(e, sqlx::Error::PoolTimedOut) {
        return GraphError::ConnectionAcquireTimeout {
            seconds: bounds.timeout.as_secs(),
        };
    }
    backend_error(source, e)
}

/// Numeric `params` mapping (design §SQL surface): a JSON number with no
/// fraction or exponent binds as Integer; with either, as Float — write
/// `1.0` to force Float. serde_json preserves the distinction (`is_i64`
/// vs `is_f64`), so validation is a kind check, not a coercion.
pub fn validate_params(params: &Value) -> Result<(), GraphError> {
    match params {
        Value::Object(_) => Ok(()),
        other => Err(GraphError::InvalidParams {
            found: json_kind(other).to_string(),
        }),
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn dial_io_errors_split_on_talking_versus_nobody_answered() {
        // rustls wraps TLS protocol errors — certificate verification
        // included — as io::Error(InvalidData), and sqlx's handshake
        // propagates that into Error::Io. That is a CONFIGURATION
        // failure of a server we reached: it must hard-fail (a Backend
        // classification), never sit degraded behind boot.
        let tls = sqlx::Error::Io(std::io::Error::new(
            IoErrorKind::InvalidData,
            "invalid peer certificate: UnknownIssuer",
        ));
        assert!(
            !matches!(connect_error("kg", &tls), GraphError::Unavailable { .. }),
            "a TLS handshake failure is not an outage"
        );
        // A refused/unroutable dial IS "nobody answered" — the degraded
        // qualifier.
        for kind in [
            IoErrorKind::ConnectionRefused,
            IoErrorKind::TimedOut,
            IoErrorKind::ConnectionReset,
        ] {
            let e = sqlx::Error::Io(std::io::Error::new(kind, "dial failed"));
            assert!(
                matches!(connect_error("kg", &e), GraphError::Unavailable { .. }),
                "{kind:?} must degrade"
            );
        }
        // A non-Io variant (a Database answer, pool bookkeeping) never
        // takes the Unavailable path either.
        assert!(!matches!(
            connect_error("kg", &sqlx::Error::PoolTimedOut),
            GraphError::Unavailable { .. }
        ));
    }

    #[test]
    fn dollar_tags_never_collide_with_the_text() {
        assert_eq!(dollar_tag("MATCH (n) RETURN n"), "$skq$");
        let hostile = "RETURN '$skq$ $skq1$'";
        let tag = dollar_tag(hostile);
        assert!(!hostile.contains(&tag), "{tag}");

        // The BOUNDARY case: a query ending in the tag's interior forms
        // the full tag against the closing delimiter's leading `$` —
        // `RETURN $skq` + `$…` would close `$skq$…$skq$` early. The probe
        // must scan text+"$", exactly the stream Postgres sees.
        let boundary = "RETURN $skq";
        let tag = dollar_tag(boundary);
        assert_ne!(tag, "$skq$", "boundary composition must bump the tag");
        assert!(!format!("{boundary}$").contains(&tag), "{tag}");
        let (sql, _) = build_cypher_sql("kg", 0, boundary, &serde_json::json!({}), 1, 11);
        // The emitted literal parses as ONE dollar-quoted string: the
        // closing tag is found exactly once past the opening.
        let open = sql.find(&tag).expect("opening tag");
        let close = sql[open + tag.len()..].find(&tag).expect("closing tag");
        assert_eq!(
            &sql[open + tag.len()..open + tag.len() + close],
            boundary,
            "the whole query is the literal body: {sql}"
        );
    }

    #[test]
    fn parameterless_sql_is_a_single_select_with_no_prepare() {
        let (sql, prepared) = build_cypher_sql(
            "kg",
            0,
            "MATCH (n) RETURN n.a, n.b",
            &serde_json::json!({}),
            2,
            101,
        );
        assert!(prepared.is_none(), "no params → nothing to DEALLOCATE");
        assert!(
            sql.starts_with("SELECT c0, c1 FROM ag_catalog.cypher('kg', $skq$"),
            "{sql}"
        );
        // The fetch bound is a REAL SQL LIMIT on the outer statement —
        // the backend and the wire are bounded even with no caller LIMIT.
        assert!(
            sql.ends_with("AS t(c0 ag_catalog.agtype, c1 ag_catalog.agtype) LIMIT 101"),
            "{sql}"
        );
        assert!(!sql.contains("PREPARE"), "{sql}");
    }

    #[test]
    fn parameterized_sql_prepares_executes_and_names_the_statement() {
        let (sql, prepared) = build_cypher_sql(
            "kg",
            7,
            "MATCH (n) WHERE n.x = $x RETURN n",
            &serde_json::json!({"x": 1}),
            1,
            5,
        );
        assert!(
            sql.contains("ag_catalog.agtype) LIMIT 5;"),
            "the PREPARE body carries the fetch LIMIT: {sql}"
        );
        let name = prepared.expect("params → PREPARE name to DEALLOCATE");
        assert!(name.starts_with("skq_p_"), "{name}");
        assert!(name.ends_with("_7"), "the seq uniquifies: {name}");
        assert!(
            sql.contains(&format!("PREPARE {name}(ag_catalog.agtype)")),
            "{sql}"
        );
        assert!(
            sql.contains(&format!("EXECUTE {name}('{{\"x\":1}}');")),
            "{sql}"
        );
    }

    #[test]
    fn hostile_values_stay_inside_their_literals() {
        // Graph names double their quotes; param values ride serde_json
        // encoding plus WHOLE-literal quote doubling — the SQL text can
        // never fall out of its string literal.
        let (sql, _) = build_cypher_sql(
            "kg",
            0,
            "RETURN $s",
            &serde_json::json!({"s": "O'Brien '; DROP TABLE x; --"}),
            1,
            10,
        );
        assert!(sql.contains("O''Brien ''; DROP TABLE x; --"), "{sql}");

        let (sql, _) = build_cypher_sql("g'name", 0, "RETURN 1", &serde_json::json!({}), 1, 10);
        assert!(sql.contains("cypher('g''name'"), "{sql}");
    }

    #[test]
    fn params_must_be_an_object() {
        assert!(validate_params(&serde_json::json!({"a": 1})).is_ok());
        let err = validate_params(&serde_json::json!([1])).unwrap_err();
        assert!(err.to_string().contains("an array"), "{err}");
    }

    #[tokio::test]
    async fn a_blackholed_backend_fails_the_preflight_as_unavailable_within_the_bound() {
        // 10.255.255.1 is unroutable in practice: the dial neither
        // connects nor is refused quickly, so without the preflight
        // timeout this connect would hang for the OS TCP timeout
        // (minutes) — holding startup hostage and never reaching the
        // degraded branch. Bounded at 1s it must surface as Unavailable,
        // the degraded qualifier. (Networks that RST unroutable
        // addresses take the io-error path to the same variant, so the
        // assertion holds either way.)
        let started = std::time::Instant::now();
        let err = AgeClient::connect(
            "kg",
            "postgres://10.255.255.1:5432/none",
            "knowledge",
            None,
            None,
            1,
            Duration::from_secs(1),
        )
        .await
        .expect_err("a blackhole cannot be reached");
        assert!(
            matches!(err, GraphError::Unavailable { .. }),
            "classified as an availability failure: {err}"
        );
        assert!(
            started.elapsed() < Duration::from_secs(10),
            "bounded by the configured timeout, not the OS's"
        );
    }

    #[test]
    fn only_transport_failures_classify_as_unavailable() {
        // A refused dial (no server answered) is the degraded qualifier.
        let io: sqlx::Error =
            std::io::Error::new(std::io::ErrorKind::ConnectionRefused, "refused").into();
        let err = connect_error("kg", &io);
        assert!(matches!(err, GraphError::Unavailable { .. }), "{err}");
        assert!(err.to_string().contains("unreachable"), "{err}");
        // Anything else — a server-answered error, pool semantics, TLS —
        // is a configuration problem that must hard-fail registration.
        let err = connect_error("kg", &sqlx::Error::PoolTimedOut);
        assert!(matches!(err, GraphError::Backend { .. }), "{err}");
    }

    #[test]
    fn a_pool_acquire_timeout_is_not_misreported_as_a_statement_timeout() {
        // sqlx retries a refused dial until the acquire deadline, so an
        // UNREACHABLE backend surfaces as PoolTimedOut — the query never
        // started. Mapping that to the statement Timeout would blame
        // the query ("narrow the traversal") for a connection problem.
        let bounds = QueryBounds {
            timeout: std::time::Duration::from_secs(3),
            max_rows: 10,
        };
        let err = map_query_error("kg", bounds, &sqlx::Error::PoolTimedOut);
        let msg = err.to_string();
        assert!(
            matches!(err, GraphError::ConnectionAcquireTimeout { seconds: 3 }),
            "{msg}"
        );
        assert!(msg.contains("never started"), "{msg}");
        assert!(!msg.contains("narrow the traversal"), "{msg}");
    }

    #[tokio::test]
    async fn connect_degraded_hard_fails_config_errors_but_never_dials() {
        // Unparseable URL: a config error, hard-failed even on the
        // degraded path.
        let err = AgeClient::connect_degraded(
            "kg",
            "not a url",
            "g",
            None,
            None,
            4,
            Duration::from_secs(1),
        )
        .unwrap_err();
        assert!(err.to_string().contains("does not parse"), "{err}");
        // Unset credential env var: same — config, not transient.
        let err = AgeClient::connect_degraded(
            "kg",
            "postgres://127.0.0.1:1/none",
            "g",
            Some("SKARDI_TEST_GRAPH_DEFINITELY_UNSET"),
            None,
            4,
            Duration::from_secs(1),
        )
        .unwrap_err();
        assert!(err.to_string().contains("is not set"), "{err}");
        // A closed port is NOT probed: the lazy pool builds fine and the
        // backend is first touched on acquire (degraded registration's
        // whole point).
        AgeClient::connect_degraded(
            "kg",
            "postgres://127.0.0.1:1/none",
            "g",
            None,
            None,
            4,
            Duration::from_secs(1),
        )
        .expect("no dial at connect_degraded");
    }
}