orion-server 1.8.0

Turn business logic into live REST/Kafka services, declared as JSON
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
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
1501
1502
1503
1504
1505
1506
1507
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
1564
1565
1566
1567
1568
1569
1570
1571
1572
1573
1574
1575
1576
1577
1578
1579
1580
1581
1582
1583
1584
1585
1586
1587
1588
1589
1590
1591
1592
1593
1594
1595
1596
1597
1598
1599
1600
1601
1602
1603
1604
1605
1606
1607
1608
1609
1610
1611
1612
1613
1614
1615
1616
1617
1618
1619
1620
1621
1622
1623
1624
1625
1626
1627
1628
1629
1630
1631
1632
1633
1634
1635
1636
1637
1638
1639
1640
1641
1642
1643
1644
1645
1646
1647
1648
1649
1650
1651
1652
//! Managed OAuth2 for `http` connectors (#268): token acquisition, in-memory
//! caching, single-flight refresh, and rotation persistence.
//!
//! The lifecycle machinery here is **grant-agnostic** — cache, refresh margin,
//! single-flight, negative cache, persistence, cluster adoption are written
//! once; each grant contributes only its token-request shape. A future grant
//! (`jwt-bearer`, token exchange, device code) is a new
//! [`OAuth2Grant`] value plus a request builder, not new machinery.
//!
//! ## Lifecycle
//!
//! Tokens are acquired lazily on first use, cached per connector, and
//! refreshed `refresh_margin_secs` before expiry. A comfortably-fresh token
//! under an unchanged config is served lock-free (`ArcSwap`); everything else
//! falls to the per-connector `Mutex`, which makes refreshes
//! **single-flight**: concurrent requests wait on the in-flight result instead
//! of racing it — the race that, under refresh-token rotation, makes the
//! losers invalidate the winner's new token.
//!
//! ## Rotation persistence
//!
//! A `refresh_token` response carrying a new refresh token is persisted to the
//! `connector_oauth_state` table (encrypted at rest like `config_json`)
//! **together with the access token and its expiry** — that is what lets other
//! cluster nodes *adopt* a fresh token instead of re-racing the rotation. The
//! state row carries a fingerprint of the connector's oauth2 block: editing
//! the connector invalidates stale state, which is also the burned-token
//! recovery story (re-seed the config → new fingerprint → seed wins). If the
//! persist write fails, the live token stays in memory and the next refresh
//! retries persistence — loudly logged, since only a restart during a
//! persistent DB failure can lose the rotated value.
//!
//! ## Failure taxonomy (F6/F42)
//!
//! A token endpoint that cannot be reached is a **retryable** dependency
//! failure. A token endpoint that *rejects* the request (`invalid_grant`,
//! `invalid_client`) is a **non-retryable** configuration failure with a 30 s
//! negative cache, so a burned refresh token is never retry-looped against the
//! IdP — and it deliberately does not trip the API's circuit breaker, because
//! a credential failure is not evidence about the API's health.

use std::collections::HashMap;
use std::sync::{Arc, OnceLock};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};

use sha2::{Digest, Sha256};

use super::config::{HttpConnectorConfig, OAuth2ClientAuth, OAuth2Config, OAuth2Grant};
use crate::storage::repositories::connectors::ConnectorRepository;

/// How long a token-endpoint rejection (`invalid_grant`, …) is answered from
/// memory before the IdP is asked again.
const NEGATIVE_CACHE: Duration = Duration::from_secs(30);
/// Token-endpoint request timeout.
const TOKEN_TIMEOUT: Duration = Duration::from_secs(10);
/// A token response larger than this is not a token, it is a problem.
const MAX_TOKEN_RESPONSE_BYTES: usize = 65_536;
/// `expires_in` absent from a token response → this conservative lifetime.
const DEFAULT_TOKEN_TTL_SECS: u64 = 300;
/// Floor on the effective lifetime, so a tiny `expires_in` cannot turn every
/// request into a refresh stampede against the IdP.
const MIN_TOKEN_TTL_SECS: u64 = 30;
/// Cluster refresh lease TTL. Comfortably above [`TOKEN_TIMEOUT`], so the
/// holder finishes (or dies) before the lease can change hands.
const LEASE_TTL_SECS: u64 = 30;
/// How a lease loser waits for the winner: poll the state row this many times…
const ADOPT_POLL_ATTEMPTS: u32 = 10;
/// …this far apart.
const ADOPT_POLL_INTERVAL: Duration = Duration::from_millis(500);

/// What the token manager needs at runtime, injected once from bootstrap.
pub struct OAuthRuntimeDeps {
    /// The shared engine client (SSRF-pinned, no auto-redirects).
    pub http_client: reqwest::Client,
    /// Where rotation state persists.
    pub repo: Arc<dyn ConnectorRepository>,
    /// Cross-node single-flight for refresh-token rotation. `None` on a
    /// single node, where the per-connector mutex already serialises.
    pub lease: Option<Arc<crate::cluster::JobLeaseGate>>,
}

/// A token-acquisition failure, classified for the estate's retry taxonomy.
#[derive(Debug, Clone)]
pub enum OAuthError {
    /// The auth block cannot be acted on (unknown grant, missing seed, the
    /// runtime not initialised). Non-retryable.
    Config(String),
    /// The IdP rejected the request (`invalid_grant`, `invalid_client`, …).
    /// Non-retryable; negative-cached.
    Rejected(String),
    /// The IdP could not be reached, answered a server error, or answered
    /// garbage. Retryable.
    Transport(String),
    /// Another node holds the refresh lease and its result did not appear in
    /// time. Retryable.
    NotReady(String),
}

impl OAuthError {
    /// Whether retrying can help — the dependency-health half of the estate's
    /// F42 split. Rejections and config problems do not fix themselves.
    pub fn retryable(&self) -> bool {
        matches!(self, OAuthError::Transport(_) | OAuthError::NotReady(_))
    }
}

impl std::fmt::Display for OAuthError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            OAuthError::Config(m)
            | OAuthError::Rejected(m)
            | OAuthError::Transport(m)
            | OAuthError::NotReady(m) => f.write_str(m),
        }
    }
}

/// The persisted half of the lifecycle — one `connector_oauth_state` row's
/// `state_json`, decrypted. Epoch milliseconds rather than an `Instant` so it
/// means the same thing on every node.
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
struct PersistedState {
    access_token: String,
    expires_at_epoch_ms: u64,
    refresh_token: String,
}

/// One connector's cached token.
#[derive(Clone)]
struct CachedToken {
    access_token: String,
    expires_at: Instant,
    expires_at_epoch_ms: u64,
}

/// One connector's in-memory lifecycle state, all mutated under the entry's
/// mutex — which is exactly what makes the refresh single-flight.
struct EntryState {
    /// Fingerprint of the oauth2 block this state was minted under. A
    /// mismatch (the connector was edited) resets everything below.
    fingerprint: String,
    /// The invalidation generation `token` was acquired under — the same
    /// stamp the published [`FastToken`] carries, so both paths refuse a
    /// token a 401 has already rejected.
    generation: u64,
    token: Option<CachedToken>,
    /// The freshest known refresh token (seed, persisted, or rotated).
    refresh_token: Option<String>,
    /// A recent IdP rejection: answered from memory for [`NEGATIVE_CACHE`].
    rejected: Option<(Instant, String)>,
}

impl EntryState {
    fn fresh(fingerprint: String) -> Self {
        Self {
            fingerprint,
            generation: 0,
            token: None,
            refresh_token: None,
            rejected: None,
        }
    }
}

/// The lock-free fast path: the last published token and the config it was
/// minted under. Written only by the slow path (holding `TokenEntry::state`),
/// read without any lock on every request — the common case is a cached token
/// comfortably inside its lifetime under an unchanged config, and it must not
/// queue behind an in-flight refresh or pay a serialize+hash fingerprint.
struct FastToken {
    config: Option<OAuth2Config>,
    token: Option<CachedToken>,
    /// The invalidation generation this token was acquired under. Stale
    /// against [`TokenEntry::invalidated`] means a 401 arrived since, and the
    /// fast path must not serve it. See [`OAuthTokenManager::invalidate`].
    generation: u64,
}

struct TokenEntry {
    fast: arc_swap::ArcSwap<FastToken>,
    state: tokio::sync::Mutex<EntryState>,
    /// Bumped by [`OAuthTokenManager::invalidate`] on every API 401, without
    /// taking `state`.
    ///
    /// The invalidation used to be attempted under `state.try_lock()` and
    /// simply skipped when another task held it — which is exactly when it
    /// matters, because the task holding it is usually the one about to
    /// publish a token. A skipped invalidation left the rejected token
    /// fast-path eligible until its refresh margin, so the connector answered
    /// 401 after 401 while a perfectly good refresh sat one call away.
    ///
    /// A counter rather than a flag, and read *before* a fetch rather than
    /// after: that is what makes an invalidation arriving mid-flight
    /// survive. See `access_token`.
    invalidated: std::sync::atomic::AtomicU64,
}

/// The process-wide token manager, held on the
/// [`ConnectorRegistry`](super::ConnectorRegistry) beside the circuit
/// breakers — per-connector runtime state keyed off the same identity.
/// The context every grant acquires against: what the connector is, how it is
/// configured, and the SSRF opt-out that config carries.
///
/// One shape for all three grants, so the dispatch reads as three ways of
/// doing the same thing rather than three unrelated argument lists — and so
/// `refresh_flow`, which needs two fields the others do not, does not become a
/// seven-parameter signature where transposing `connector` and `fingerprint`
/// compiles.
#[derive(Clone, Copy)]
struct GrantCtx<'a> {
    deps: &'a OAuthRuntimeDeps,
    connector: &'a str,
    cfg: &'a OAuth2Config,
    /// Ties cached and persisted state to the config that produced it; only
    /// the refresh-token grant persists a row, so only it reads this.
    fingerprint: &'a str,
    /// How early a token counts as stale. Read by the refresh grant, which can
    /// adopt another node's persisted access token if it is still fresh.
    margin: Duration,
    /// The connector's own SSRF opt-out.
    allow_private_urls: bool,
}

pub struct OAuthTokenManager {
    entries: tokio::sync::RwLock<HashMap<String, Arc<TokenEntry>>>,
    deps: OnceLock<OAuthRuntimeDeps>,
}

impl Default for OAuthTokenManager {
    fn default() -> Self {
        Self::new()
    }
}

impl OAuthTokenManager {
    pub fn new() -> Self {
        Self {
            entries: tokio::sync::RwLock::new(HashMap::new()),
            deps: OnceLock::new(),
        }
    }

    /// Inject the runtime dependencies, once, from bootstrap. A second call
    /// (a test harness building over the same registry) is a no-op.
    pub fn init(&self, deps: OAuthRuntimeDeps) {
        let _ = self.deps.set(deps);
    }

    /// Drop the cached access token for `connector`, so the next call
    /// refetches. The refresh token is kept — it is still valid; only the
    /// access token was rejected. Called when the *API* answers 401 on an
    /// oauth2 connector: self-healing after IdP-side revocation.
    ///
    /// Cannot fail and cannot wait. It takes no lock at all: bumping the
    /// generation is what invalidates, and both the fast path and the slow
    /// path compare against it. Clearing the published token as well is a
    /// courtesy — it saves the next caller a trip through the mutex — but it
    /// is not what makes this correct, because an acquisition already in
    /// flight will publish over it. The generation is what that publish
    /// cannot outrun.
    pub async fn invalidate(&self, connector: &str) {
        let Some(entry) = self.entries.read().await.get(connector).cloned() else {
            return;
        };
        // Release: a reader that sees this bump must also see everything
        // written before the 401 that prompted it.
        entry
            .invalidated
            .fetch_add(1, std::sync::atomic::Ordering::Release);
        let fast = entry.fast.load();
        entry.fast.store(Arc::new(FastToken {
            config: fast.config.clone(),
            token: None,
            generation: fast.generation,
        }));
    }

    /// The current access token for `connector`, acquiring or refreshing as
    /// needed. `allow_private_urls` is the connector's own SSRF opt-out and
    /// covers the token endpoint too.
    pub async fn access_token(
        &self,
        connector: &str,
        cfg: &OAuth2Config,
        allow_private_urls: bool,
    ) -> Result<String, OAuthError> {
        let deps = self.deps.get().ok_or_else(|| {
            OAuthError::Config(
                "OAuth2 runtime is not initialised (managed OAuth2 is unavailable in \
                 this context)"
                    .to_string(),
            )
        })?;
        let grant = OAuth2Grant::parse(&cfg.grant).ok_or_else(|| {
            OAuthError::Config(format!(
                "connector '{connector}': unknown OAuth2 grant '{}' (expected {})",
                cfg.grant,
                OAuth2Grant::VALUES
            ))
        })?;
        let margin = Duration::from_secs(cfg.refresh_margin_secs.min(3600));
        let entry = self.entry(connector).await;

        // The invalidation generation, read once and used for both paths.
        // Read *before* anything else so that an invalidation racing this call
        // is either already visible here (and honoured) or lands after the
        // token this call publishes is stamped — in which case the stamp is
        // stale and the *next* caller refetches. Either way it is not lost,
        // which is the whole difference from the `try_lock` this replaces.
        let generation = entry.invalidated.load(std::sync::atomic::Ordering::Acquire);

        // Fast path: exactly the condition the slow path answers from cache
        // (unchanged config, token comfortably fresh, no 401 since it was
        // acquired), without the mutex — so requests never serialize behind an
        // in-flight refresh while the current token is still valid past the
        // margin.
        {
            let fast = entry.fast.load();
            if fast.generation == generation
                && fast.config.as_ref() == Some(cfg)
                && let Some(token) = fresh_token(&fast.token, margin)
            {
                return Ok(token);
            }
        }

        let fingerprint = fingerprint(cfg);
        let mut st = entry.state.lock().await;

        // The connector was edited since this state was minted: everything —
        // token, refresh token, negative cache — belongs to the old config.
        if st.fingerprint != fingerprint {
            *st = EntryState::fresh(fingerprint.clone());
        }

        // A 401 since this token was cached means it is not a token any more,
        // however fresh its expiry looks. Dropped here rather than by
        // `invalidate` itself, because `invalidate` deliberately takes no
        // lock — the generation is the message, this is where it is read.
        if st.generation < generation {
            st.token = None;
        }

        // Re-check under the lock: a queued caller usually finds the winner's
        // freshly published token here.
        if let Some(token) = fresh_token(&st.token, margin) {
            return Ok(token);
        }
        if let Some((at, message)) = &st.rejected {
            if at.elapsed() < NEGATIVE_CACHE {
                return Err(OAuthError::Rejected(message.clone()));
            }
            st.rejected = None;
        }

        let ctx = GrantCtx {
            deps,
            connector,
            cfg,
            fingerprint: &fingerprint,
            margin,
            allow_private_urls,
        };
        let result = match grant {
            OAuth2Grant::ClientCredentials => self.acquire_client_credentials(&ctx, &mut st).await,
            OAuth2Grant::AccountCredentials => {
                self.acquire_account_credentials(&ctx, &mut st).await
            }
            OAuth2Grant::RefreshToken => self.refresh_flow(&ctx, &mut st).await,
        };
        if let Err(e) = &result
            && !e.retryable()
        {
            st.rejected = Some((Instant::now(), e.to_string()));
        }
        // Publish for the fast path, whatever happened, while the lock is
        // still held (its writers are ordered by the mutex).
        //
        // Stamped with the generation read at the *top* of this call, not a
        // fresh load: an invalidation that arrived while the token request was
        // in flight rejected the token this is about to publish, and stamping
        // it with the newer generation would hide that. Stale-stamping is the
        // safe direction — the fast path skips it and the next caller
        // refetches.
        st.generation = generation;
        entry.fast.store(Arc::new(FastToken {
            config: Some(cfg.clone()),
            token: st.token.clone(),
            generation,
        }));
        result
    }

    async fn entry(&self, connector: &str) -> Arc<TokenEntry> {
        if let Some(entry) = self.entries.read().await.get(connector) {
            return Arc::clone(entry);
        }
        let mut entries = self.entries.write().await;
        Arc::clone(entries.entry(connector.to_string()).or_insert_with(|| {
            Arc::new(TokenEntry {
                fast: arc_swap::ArcSwap::new(Arc::new(FastToken {
                    config: None,
                    token: None,
                    generation: 0,
                })),
                state: tokio::sync::Mutex::new(EntryState::fresh(String::new())),
                invalidated: std::sync::atomic::AtomicU64::new(0),
            })
        }))
    }

    async fn acquire_client_credentials(
        &self,
        ctx: &GrantCtx<'_>,
        st: &mut EntryState,
    ) -> Result<String, OAuthError> {
        let GrantCtx {
            deps,
            connector,
            cfg,
            allow_private_urls,
            ..
        } = *ctx;
        let mut params: Vec<(&str, String)> =
            vec![("grant_type", "client_credentials".to_string())];
        if !cfg.scopes.is_empty() {
            params.push(("scope", cfg.scopes.join(" ")));
        }
        if let Some(a) = &cfg.audience {
            params.push(("audience", a.clone()));
        }
        if let Some(r) = &cfg.resource {
            params.push(("resource", r.clone()));
        }
        for (k, v) in &cfg.extra_params {
            params.push((k.as_str(), v.clone()));
        }
        let response = request_token(deps, connector, cfg, allow_private_urls, params).await?;
        Ok(cache_token(st, &response).access_token)
    }

    /// Zoom Server-to-Server OAuth: `grant_type=account_credentials` plus the
    /// account id, with Basic client auth (Zoom's expectation, and Orion's
    /// default).
    ///
    /// Everything below `acquire_*` is grant-agnostic, so this inherits with
    /// no new code: lazy acquisition, the per-connector `ArcSwap` fast path,
    /// early refresh, single-flight, the retryable/non-retryable failure split
    /// and its negative cache, 401-from-API self-healing, the token-request
    /// metric, and the admin connector probe.
    ///
    /// It deliberately does **not** inherit rotation persistence or the
    /// cluster job lease — those live only in `refresh_flow`. Zoom re-acquires
    /// from static credentials, so there is no `connector_oauth_state` row and
    /// no lease to take: correct by construction, not by omission.
    async fn acquire_account_credentials(
        &self,
        ctx: &GrantCtx<'_>,
        st: &mut EntryState,
    ) -> Result<String, OAuthError> {
        let GrantCtx {
            deps,
            connector,
            cfg,
            allow_private_urls,
            ..
        } = *ctx;
        let mut params: Vec<(&str, String)> =
            vec![("grant_type", "account_credentials".to_string())];
        // Validation requires it for this grant, so an absent value here means
        // a hand-edited row: fail as a non-retryable config error rather than
        // sending Zoom a request it will reject on every attempt.
        let Some(account_id) = &cfg.account_id else {
            return Err(OAuthError::Config(
                "oauth2 account_credentials requires 'account_id'".to_string(),
            ));
        };
        params.push(("account_id", account_id.clone()));
        if !cfg.scopes.is_empty() {
            params.push(("scope", cfg.scopes.join(" ")));
        }
        for (k, v) in &cfg.extra_params {
            params.push((k.as_str(), v.clone()));
        }
        let response = request_token(deps, connector, cfg, allow_private_urls, params).await?;
        Ok(cache_token(st, &response).access_token)
    }

    async fn refresh_flow(
        &self,
        ctx: &GrantCtx<'_>,
        st: &mut EntryState,
    ) -> Result<String, OAuthError> {
        let GrantCtx {
            deps,
            connector,
            cfg,
            fingerprint,
            margin,
            allow_private_urls,
        } = *ctx;
        // First use in this process: the freshest refresh token may be a
        // rotation another node (or a previous run) persisted — and its
        // access token may still be perfectly good.
        if st.refresh_token.is_none() {
            if let Some(persisted) = load_state(deps, connector, fingerprint).await {
                st.refresh_token = Some(persisted.refresh_token.clone());
                if let Some(token) = adopt(st, &persisted, margin) {
                    return Ok(token);
                }
            } else {
                let seed = cfg.refresh_token.as_deref().unwrap_or("").trim();
                if seed.is_empty() {
                    return Err(OAuthError::Config(format!(
                        "connector '{connector}': the refresh_token grant needs a \
                         'refresh_token' seed in the auth block"
                    )));
                }
                st.refresh_token = Some(seed.to_string());
            }
        }

        // Cross-node single-flight: rotation must happen on one node. A lease
        // loser adopts the winner's persisted token instead of re-racing.
        if let Some(lease) = &deps.lease {
            let job = format!("oauth-refresh:{connector}");
            if !lease.try_acquire(&job, LEASE_TTL_SECS).await {
                for _ in 0..ADOPT_POLL_ATTEMPTS {
                    tokio::time::sleep(ADOPT_POLL_INTERVAL).await;
                    if let Some(persisted) = load_state(deps, connector, fingerprint).await {
                        st.refresh_token = Some(persisted.refresh_token.clone());
                        if let Some(token) = adopt(st, &persisted, margin) {
                            return Ok(token);
                        }
                    }
                }
                return Err(OAuthError::NotReady(format!(
                    "connector '{connector}': another node is refreshing the OAuth2 \
                     token and its result has not appeared yet"
                )));
            }
        }

        let current_rt = st
            .refresh_token
            .clone()
            .expect("ensured above: persisted, adopted, or seeded");
        let mut params: Vec<(&str, String)> = vec![
            ("grant_type", "refresh_token".to_string()),
            ("refresh_token", current_rt.clone()),
        ];
        for (k, v) in &cfg.extra_params {
            params.push((k.as_str(), v.clone()));
        }
        let response = request_token(deps, connector, cfg, allow_private_urls, params).await?;
        let cached = cache_token(st, &response);

        // Rotation: the response's refresh token (when present) replaces the
        // one just spent. Persist token + rotation before anyone else asks.
        let next_rt = response.refresh_token.clone().unwrap_or(current_rt);
        st.refresh_token = Some(next_rt.clone());
        let persisted = PersistedState {
            access_token: cached.access_token.clone(),
            expires_at_epoch_ms: cached.expires_at_epoch_ms,
            refresh_token: next_rt,
        };
        match serde_json::to_string(&persisted) {
            Ok(json) => {
                if let Err(e) = deps
                    .repo
                    .put_oauth_state(connector, fingerprint, &json)
                    .await
                {
                    // The live values are still in memory, so the next
                    // refresh retries persistence — but a restart during a
                    // persistent failure here loses the rotation.
                    crate::metrics::record_error("oauth_state_persist");
                    tracing::error!(
                        connector,
                        error = %e,
                        "failed to persist rotated OAuth2 refresh token; \
                         it survives only in memory until the next refresh"
                    );
                }
            }
            Err(e) => {
                crate::metrics::record_error("oauth_state_persist");
                tracing::error!(connector, error = %e, "OAuth2 state did not serialize");
            }
        }
        Ok(cached.access_token)
    }
}

/// The cached token when it is still comfortably inside its lifetime.
fn fresh_token(token: &Option<CachedToken>, margin: Duration) -> Option<String> {
    let token = token.as_ref()?;
    (Instant::now() + margin < token.expires_at).then(|| token.access_token.clone())
}

/// Store a token response in the entry and hand back the cached form — the
/// caller reads the expiry from it rather than re-deriving it from `st`.
fn cache_token(st: &mut EntryState, response: &TokenResponse) -> CachedToken {
    let ttl = response
        .expires_in
        .unwrap_or(DEFAULT_TOKEN_TTL_SECS)
        .max(MIN_TOKEN_TTL_SECS);
    let token = CachedToken {
        access_token: response.access_token.clone(),
        expires_at: Instant::now() + Duration::from_secs(ttl),
        expires_at_epoch_ms: epoch_ms_now().saturating_add(ttl * 1000),
    };
    st.token = Some(token.clone());
    token
}

/// Adopt a persisted access token when it is still fresh under `margin`.
fn adopt(st: &mut EntryState, persisted: &PersistedState, margin: Duration) -> Option<String> {
    let now_ms = epoch_ms_now();
    let margin_ms = margin.as_millis() as u64;
    let remaining_ms = persisted.expires_at_epoch_ms.checked_sub(now_ms)?;
    if remaining_ms <= margin_ms {
        return None;
    }
    st.token = Some(CachedToken {
        access_token: persisted.access_token.clone(),
        expires_at: Instant::now() + Duration::from_millis(remaining_ms),
        expires_at_epoch_ms: persisted.expires_at_epoch_ms,
    });
    Some(persisted.access_token.clone())
}

/// The persisted state for `connector`, if any exists **for this exact
/// config** — a fingerprint mismatch means the connector was edited and the
/// stored state belongs to the old credentials.
async fn load_state(
    deps: &OAuthRuntimeDeps,
    connector: &str,
    fingerprint: &str,
) -> Option<PersistedState> {
    match deps.repo.get_oauth_state(connector).await {
        Ok(Some(row)) if row.fingerprint == fingerprint => {
            serde_json::from_str(&row.state_json).ok()
        }
        Ok(_) => None,
        Err(e) => {
            tracing::warn!(connector, error = %e, "could not read OAuth2 state");
            None
        }
    }
}

/// A fingerprint of the oauth2 block: what ties cached and persisted state to
/// the exact credentials/endpoint they were minted under.
fn fingerprint(cfg: &OAuth2Config) -> String {
    let serialized = serde_json::to_string(cfg).unwrap_or_default();
    let mut hasher = Sha256::new();
    hasher.update(serialized.as_bytes());
    hex::encode(hasher.finalize())
}

fn epoch_ms_now() -> u64 {
    SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .map(|d| d.as_millis() as u64)
        .unwrap_or_default()
}

/// The RFC 6749 §5.1 fields this build reads.
#[derive(Debug, Clone)]
pub(crate) struct TokenResponse {
    pub access_token: String,
    pub token_type: Option<String>,
    pub expires_in: Option<u64>,
    pub refresh_token: Option<String>,
    /// OIDC (#307). The connector grants never see one; the inbound
    /// authorization-code grant is the whole reason it is read.
    pub id_token: Option<String>,
    /// The scopes actually granted, which an IdP may narrow from what was
    /// asked for.
    pub scope: Option<String>,
}

/// The four things a token request needs, independent of which grant asked for
/// it and of whether the caller is a connector at all.
///
/// [`request_token_inner`] used to take `&OAuth2Config` whole, which made the
/// exchange unreachable from the inbound sign-in flow (#307) without adding
/// authorization-code fields to that struct. That would not have been free:
/// [`fingerprint`] hashes the serialized `OAuth2Config`, and a persisted
/// `connector_oauth_state` row is adopted only on a fingerprint match, so a new
/// field would orphan every stored rotation state on upgrade. Narrowing the
/// parameter instead leaves the fingerprint untouched.
#[derive(Clone, Copy)]
pub(crate) struct TokenEndpoint<'a> {
    pub token_url: &'a str,
    pub client_id: &'a str,
    pub client_secret: &'a str,
    /// `basic` or `body`; anything else is a config error, named.
    pub client_auth: &'a str,
}

impl<'a> TokenEndpoint<'a> {
    fn from_config(cfg: &'a OAuth2Config) -> Self {
        Self {
            token_url: &cfg.token_url,
            client_id: &cfg.client_id,
            client_secret: &cfg.client_secret,
            client_auth: &cfg.client_auth,
        }
    }
}

/// Exchange an authorization code for tokens (RFC 6749 §4.1.3) — the inbound
/// half of OAuth2, where Orion is the relying party rather than the client.
///
/// Shares every safety property of the connector grants because it is the same
/// function underneath: the SSRF check on the token URL, the response size
/// caps, the RFC 6749 §5.2 error classification and the Bearer-only check.
/// `connector` is the metric label — the channel name, here — so one dashboard
/// covers outbound and inbound token traffic.
pub(crate) async fn exchange_code(
    http_client: &reqwest::Client,
    label: &str,
    endpoint: TokenEndpoint<'_>,
    allow_private_urls: bool,
    params: Vec<(&str, String)>,
) -> Result<TokenResponse, OAuthError> {
    let result = request_token_inner(http_client, endpoint, allow_private_urls, params).await;
    let outcome = match &result {
        Ok(_) => "ok",
        Err(OAuthError::Rejected(_)) => "rejected",
        Err(_) => "transport_error",
    };
    crate::metrics::record_oauth_token_request(label, outcome);
    result
}

/// One form-encoded token request (RFC 6749 §4.4.2 / §6), classified per the
/// failure taxonomy and counted in `orion_oauth_token_requests_total`.
async fn request_token(
    deps: &OAuthRuntimeDeps,
    connector: &str,
    cfg: &OAuth2Config,
    allow_private_urls: bool,
    params: Vec<(&str, String)>,
) -> Result<TokenResponse, OAuthError> {
    let result = request_token_inner(
        &deps.http_client,
        TokenEndpoint::from_config(cfg),
        allow_private_urls,
        params,
    )
    .await;
    let outcome = match &result {
        Ok(_) => "ok",
        Err(OAuthError::Rejected(_)) => "rejected",
        Err(_) => "transport_error",
    };
    crate::metrics::record_oauth_token_request(connector, outcome);
    result.map_err(|e| prefix_connector(connector, e))
}

/// Put the connector name on the message once, at the boundary.
fn prefix_connector(connector: &str, e: OAuthError) -> OAuthError {
    let wrap = |m: String| format!("connector '{connector}': {m}");
    match e {
        OAuthError::Config(m) => OAuthError::Config(wrap(m)),
        OAuthError::Rejected(m) => OAuthError::Rejected(wrap(m)),
        OAuthError::Transport(m) => OAuthError::Transport(wrap(m)),
        OAuthError::NotReady(m) => OAuthError::NotReady(wrap(m)),
    }
}

async fn request_token_inner(
    http_client: &reqwest::Client,
    endpoint: TokenEndpoint<'_>,
    allow_private_urls: bool,
    mut params: Vec<(&str, String)>,
) -> Result<TokenResponse, OAuthError> {
    // The token endpoint gets the same SSRF treatment as the connector's own
    // endpoint — it is admin-authored config, but the check is cheap and the
    // opt-out is the connector's existing one.
    if !allow_private_urls
        && let Err(msg) = crate::validation::validate_url_not_private(endpoint.token_url).await
    {
        return Err(OAuthError::Config(format!("SSRF protection: {msg}")));
    }

    let client_auth = OAuth2ClientAuth::parse(endpoint.client_auth).ok_or_else(|| {
        OAuthError::Config(format!(
            "unknown client_auth '{}' (expected {})",
            endpoint.client_auth,
            OAuth2ClientAuth::VALUES
        ))
    })?;
    let mut req = http_client.post(endpoint.token_url).timeout(TOKEN_TIMEOUT);
    match client_auth {
        OAuth2ClientAuth::Basic => {
            req = req.basic_auth(endpoint.client_id, Some(endpoint.client_secret));
        }
        OAuth2ClientAuth::Body => {
            params.push(("client_id", endpoint.client_id.to_string()));
            params.push(("client_secret", endpoint.client_secret.to_string()));
        }
    }

    // Form-encoded per RFC 6749 §4 — the same encoder `http_call`'s
    // `body_format: "form"` uses (#261), built manually because the shared
    // client is compiled without reqwest's form feature. Scoped: the
    // serializer is not `Send` and must not live across the await.
    let form_body = {
        let mut ser = url::form_urlencoded::Serializer::new(String::new());
        for (k, v) in &params {
            ser.append_pair(k, v);
        }
        ser.finish()
    };
    let response = req
        .header("content-type", "application/x-www-form-urlencoded")
        // RFC 6749 §5.1 says the response is JSON, and every IdP that follows
        // it ignores this header. GitHub does not: without an explicit
        // `Accept` its token endpoint answers `application/x-www-form-urlencoded`,
        // which parses to `null` below and surfaces as the wrong error
        // entirely ("carried no access_token", classified retryable).
        .header("accept", "application/json")
        .body(form_body)
        .send()
        .await
        .map_err(|e| OAuthError::Transport(format!("token request failed: {e}")))?;
    let status = response.status();
    // Bounded *while streaming* (`http_body`), and read on every status: the
    // body is the authority on whether the request was rejected (see below),
    // so this cannot skip it for a non-2xx. The declared-length check alone
    // was no cap at all against a token endpoint that omits `Content-Length`.
    let body = crate::http_body::read_bounded(response, MAX_TOKEN_RESPONSE_BYTES)
        .await
        .map_err(|e| OAuthError::Transport(format!("token response {e}")))?;
    let json: serde_json::Value = serde_json::from_slice(&body).unwrap_or_default();

    // RFC 6749 §5.2 pairs the error body with a 400, and most IdPs comply.
    // GitHub answers `200` with `{"error": "bad_verification_code"}` for a
    // spent or forged authorization code. Classified by status alone that
    // falls through to "carried no access_token" — an `OAuthError::Transport`,
    // which `retryable()` reports as **true**, so a permanently dead code
    // would be retried. The body is the authority on whether the request was
    // rejected; the status only says how politely.
    let rejected_with_200 = status.is_success() && json.get("error").is_some();

    if !status.is_success() || rejected_with_200 {
        // RFC 6749 §5.2: the error code is bounded vocabulary and safe to
        // surface; the free-text description is logged, not returned.
        let code = json
            .get("error")
            .and_then(|e| e.as_str())
            .unwrap_or("no error code");
        if let Some(desc) = json.get("error_description").and_then(|d| d.as_str()) {
            tracing::warn!(
                status = status.as_u16(),
                code,
                desc,
                "OAuth2 token request rejected"
            );
        }
        if status.is_client_error() || rejected_with_200 {
            return Err(OAuthError::Rejected(format!(
                "token endpoint rejected the request ({status}, {code}) — check the \
                 credentials, grant, and (for refresh_token) whether the seed is \
                 still valid; re-seed the connector to recover a burned token"
            )));
        }
        return Err(OAuthError::Transport(format!(
            "token endpoint answered {status} ({code})"
        )));
    }

    let access_token = json
        .get("access_token")
        .and_then(|t| t.as_str())
        .filter(|t| !t.is_empty())
        .ok_or_else(|| {
            OAuthError::Transport("token response carried no access_token".to_string())
        })?;
    let token_type = json.get("token_type").and_then(|t| t.as_str());
    if let Some(token_type) = token_type
        && !token_type.eq_ignore_ascii_case("bearer")
    {
        return Err(OAuthError::Config(format!(
            "token endpoint issued a '{token_type}' token; only Bearer is supported"
        )));
    }
    let string_field = |name: &str| {
        json.get(name)
            .and_then(|t| t.as_str())
            .filter(|t| !t.is_empty())
            .map(str::to_string)
    };
    Ok(TokenResponse {
        access_token: access_token.to_string(),
        token_type: token_type.map(str::to_string),
        expires_in: json.get("expires_in").and_then(|e| e.as_u64()),
        refresh_token: string_field("refresh_token"),
        id_token: string_field("id_token"),
        scope: string_field("scope"),
    })
}

/// The auth to actually apply for one request to `connector`: static variants
/// pass through untouched; `oauth2` resolves to a Bearer token through the
/// manager. The typed [`OAuthError`] is surfaced so the caller can map it
/// into its own error vocabulary (the engine maps retryable → `Io`).
pub async fn effective_auth<'a>(
    manager: &OAuthTokenManager,
    connector: &str,
    http: &'a HttpConnectorConfig,
) -> Result<Option<std::borrow::Cow<'a, super::config::AuthConfig>>, OAuthError> {
    use super::config::AuthConfig;
    match &http.auth {
        None => Ok(None),
        Some(AuthConfig::OAuth2(cfg)) => {
            let token = manager
                .access_token(connector, cfg, http.allow_private_urls)
                .await?;
            Ok(Some(std::borrow::Cow::Owned(AuthConfig::Bearer { token })))
        }
        Some(other) => Ok(Some(std::borrow::Cow::Borrowed(other))),
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::connector::config::AuthConfig;
    use crate::connector::test_support::StubConnectorRepo;
    use axum::extract::State;
    use axum::http::{HeaderMap, StatusCode};
    use serde_json::{Value, json};
    use std::sync::Mutex;
    use std::sync::atomic::{AtomicU64, Ordering};

    /// What the fake IdP records and how it answers — the observable half of
    /// every test below.
    struct IdpState {
        hits: AtomicU64,
        /// The `refresh_token` presented on the most recent refresh.
        last_refresh_token: Mutex<Option<String>>,
        /// Whether the last request carried HTTP Basic client auth.
        last_had_basic: Mutex<bool>,
        last_form: Mutex<Vec<(String, String)>>,
        /// The canned answer: `(status, body)`.
        respond: Mutex<(u16, Value)>,
    }

    impl IdpState {
        fn ok(body: Value) -> Arc<Self> {
            Arc::new(Self {
                hits: AtomicU64::new(0),
                last_refresh_token: Mutex::new(None),
                last_had_basic: Mutex::new(false),
                last_form: Mutex::new(Vec::new()),
                respond: Mutex::new((200, body)),
            })
        }
        fn hits(&self) -> u64 {
            self.hits.load(Ordering::SeqCst)
        }
        fn set_response(&self, status: u16, body: Value) {
            *self.respond.lock().expect("test") = (status, body);
        }
    }

    /// An in-process token endpoint speaking just enough RFC 6749.
    async fn fake_idp(state: Arc<IdpState>) -> String {
        async fn token(
            State(st): State<Arc<IdpState>>,
            headers: HeaderMap,
            body: String,
        ) -> (StatusCode, axum::Json<Value>) {
            let form: Vec<(String, String)> = url::form_urlencoded::parse(body.as_bytes())
                .map(|(k, v)| (k.into_owned(), v.into_owned()))
                .collect();
            st.hits.fetch_add(1, Ordering::SeqCst);
            *st.last_had_basic.lock().expect("test") = headers
                .get("authorization")
                .and_then(|v| v.to_str().ok())
                .is_some_and(|v| v.starts_with("Basic "));
            *st.last_refresh_token.lock().expect("test") = form
                .iter()
                .find(|(k, _)| k == "refresh_token")
                .map(|(_, v)| v.clone());
            *st.last_form.lock().expect("test") = form;
            let (status, body) = st.respond.lock().expect("test").clone();
            (
                StatusCode::from_u16(status).expect("test status"),
                axum::Json(body),
            )
        }
        let app = axum::Router::new()
            .route("/token", axum::routing::post(token))
            .with_state(state);
        let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
            .await
            .expect("bind");
        let addr = listener.local_addr().expect("addr");
        tokio::spawn(async move {
            axum::serve(listener, app).await.expect("serve");
        });
        format!("http://{addr}/token")
    }

    fn manager_with(repo: Arc<StubConnectorRepo>) -> OAuthTokenManager {
        let manager = OAuthTokenManager::new();
        manager.init(OAuthRuntimeDeps {
            http_client: reqwest::Client::new(),
            repo,
            lease: None,
        });
        manager
    }

    fn cc_config(token_url: &str) -> OAuth2Config {
        OAuth2Config {
            grant: "client_credentials".to_string(),
            token_url: token_url.to_string(),
            client_id: "cid".to_string(),
            client_secret: "csecret".to_string(),
            client_auth: "basic".to_string(),
            refresh_token: None,
            scopes: vec!["api.read".to_string(), "api.write".to_string()],
            audience: Some("https://api.example.com".to_string()),
            resource: None,
            account_id: None,
            extra_params: HashMap::new(),
            refresh_margin_secs: 1,
        }
    }

    fn rt_config(token_url: &str, seed: &str) -> OAuth2Config {
        OAuth2Config {
            grant: "refresh_token".to_string(),
            refresh_token: Some(seed.to_string()),
            scopes: Vec::new(),
            audience: None,
            ..cc_config(token_url)
        }
    }

    /// Zoom Server-to-Server: the account id replaces the audience, and there
    /// is no seed to carry.
    fn ac_config(token_url: &str) -> OAuth2Config {
        OAuth2Config {
            grant: "account_credentials".to_string(),
            account_id: Some("zoom-acct-1".to_string()),
            audience: None,
            scopes: Vec::new(),
            ..cc_config(token_url)
        }
    }

    // ---- client_credentials: acquisition, caching, request shape ----

    #[tokio::test]
    async fn a_cached_token_serves_every_call_until_the_margin() {
        let idp = IdpState::ok(json!({
            "access_token": "tok-1", "token_type": "Bearer", "expires_in": 3600
        }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = cc_config(&url);

        for _ in 0..5 {
            let token = manager
                .access_token("crm", &cfg, true)
                .await
                .expect("token");
            assert_eq!(token, "tok-1");
        }
        assert_eq!(idp.hits(), 1, "one acquisition serves the cache window");

        // The request carried the RFC 6749 shape: Basic client auth, the
        // grant, the space-joined scopes, and the audience.
        assert!(*idp.last_had_basic.lock().expect("test"));
        let form = idp.last_form.lock().expect("test").clone();
        let get = |k: &str| {
            form.iter()
                .find(|(key, _)| key == k)
                .map(|(_, v)| v.clone())
        };
        assert_eq!(get("grant_type").as_deref(), Some("client_credentials"));
        assert_eq!(get("scope").as_deref(), Some("api.read api.write"));
        assert_eq!(get("audience").as_deref(), Some("https://api.example.com"));
        assert_eq!(
            get("client_id"),
            None,
            "basic auth puts nothing in the body"
        );
    }

    // ---- account_credentials (Zoom Server-to-Server) ----

    /// #273: the grant's whole request shape —
    /// `grant_type=account_credentials`, the account id, and Basic client
    /// auth, which is what Zoom expects and Orion's default.
    #[tokio::test]
    async fn account_credentials_sends_the_zoom_request_shape() {
        let idp = IdpState::ok(json!({
            "access_token": "zoom-tok", "token_type": "Bearer", "expires_in": 3600
        }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = ac_config(&url);

        let token = manager
            .access_token("zoom", &cfg, true)
            .await
            .expect("token");
        assert_eq!(token, "zoom-tok");

        assert!(
            *idp.last_had_basic.lock().expect("test"),
            "Zoom authenticates the client with HTTP Basic"
        );
        let form = idp.last_form.lock().expect("test").clone();
        let get = |k: &str| {
            form.iter()
                .find(|(key, _)| key == k)
                .map(|(_, v)| v.clone())
        };
        assert_eq!(get("grant_type").as_deref(), Some("account_credentials"));
        assert_eq!(get("account_id").as_deref(), Some("zoom-acct-1"));
        assert_eq!(
            get("client_id"),
            None,
            "basic auth puts nothing in the body"
        );
    }

    /// The grant inherits the cache for free — it goes through the same
    /// `cache_token` path as `client_credentials`, with no rotation state.
    #[tokio::test]
    async fn an_account_credentials_token_is_cached_like_any_other() {
        let idp = IdpState::ok(json!({
            "access_token": "zoom-tok", "expires_in": 3600
        }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let repo = Arc::new(StubConnectorRepo::with(vec![]));
        let manager = manager_with(Arc::clone(&repo) as Arc<StubConnectorRepo>);
        let cfg = ac_config(&url);

        for _ in 0..5 {
            assert_eq!(
                manager
                    .access_token("zoom", &cfg, true)
                    .await
                    .expect("token"),
                "zoom-tok"
            );
        }
        assert_eq!(idp.hits(), 1, "one acquisition serves the cache window");
        assert!(
            repo.get_oauth_state("zoom").await.expect("repo").is_none(),
            "re-acquired from static credentials, so there is no rotation state \
             to persist and no cluster lease to take"
        );
    }

    /// ⚠️ The upgrade hazard, pinned. `fingerprint` hashes the whole
    /// serialized `OAuth2Config`, and a persisted state row is adopted only on
    /// a fingerprint match. If `account_id` serialized as `"account_id": null`
    /// when unset, **every existing refresh_token connector** would change
    /// fingerprint on upgrade, discard its persisted state and fall back to
    /// the config seed — which, after any prior rotation, is a spent refresh
    /// token. `skip_serializing_if` is what prevents that, and this test is
    /// what stops someone "tidying" it away.
    #[test]
    fn an_unset_account_id_does_not_move_an_existing_fingerprint() {
        let cfg = rt_config("https://idp.example/token", "rt-seed");
        let serialized = serde_json::to_string(&cfg).expect("json");
        assert!(
            !serialized.contains("account_id"),
            "an unset account_id must not serialize at all, or it changes the \
             fingerprint of every stored connector: {serialized}"
        );
    }

    #[tokio::test]
    async fn client_auth_body_moves_the_credentials_into_the_form() {
        let idp = IdpState::ok(json!({ "access_token": "t", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = OAuth2Config {
            client_auth: "body".to_string(),
            ..cc_config(&url)
        };
        manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");
        assert!(!*idp.last_had_basic.lock().expect("test"));
        let form = idp.last_form.lock().expect("test").clone();
        assert!(form.iter().any(|(k, v)| k == "client_id" && v == "cid"));
        assert!(
            form.iter()
                .any(|(k, v)| k == "client_secret" && v == "csecret")
        );
    }

    /// The single-flight guarantee: concurrent first requests produce one
    /// token acquisition, not a race.
    #[tokio::test]
    async fn concurrent_requests_share_one_acquisition() {
        let idp = IdpState::ok(json!({ "access_token": "t", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = Arc::new(manager_with(Arc::new(StubConnectorRepo::with(vec![]))));
        let cfg = Arc::new(cc_config(&url));

        let mut handles = Vec::new();
        for _ in 0..8 {
            let manager = Arc::clone(&manager);
            let cfg = Arc::clone(&cfg);
            handles.push(tokio::spawn(async move {
                manager.access_token("crm", &cfg, true).await
            }));
        }
        for h in handles {
            h.await.expect("join").expect("token");
        }
        assert_eq!(idp.hits(), 1, "eight callers, one token request");
    }

    // ---- 401 invalidation cannot be dropped ----

    /// The race the `try_lock` version lost.
    ///
    /// `invalidate` used to attempt `state.try_lock()` and skip when another
    /// task held it — which is precisely when it matters, because the task
    /// holding it is the one about to publish a token. Here the lock is held
    /// for the whole invalidation, so the old code did nothing at all and the
    /// rejected token stayed fast-path eligible until its refresh margin: a
    /// connector answering 401 after 401 with a good refresh one call away.
    #[tokio::test]
    async fn invalidation_is_not_lost_when_the_state_lock_is_held() {
        let idp = IdpState::ok(json!({ "access_token": "at-1", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = cc_config(&url);

        let first = manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");
        assert_eq!(first, "at-1");
        assert_eq!(idp.hits(), 1);

        // The API answers 401 while something else owns the entry's state —
        // an acquisition for another config, a refresh, anything.
        let entry = manager.entry("crm").await;
        let held = entry.state.lock().await;
        manager.invalidate("crm").await;
        drop(held);

        idp.set_response(200, json!({ "access_token": "at-2", "expires_in": 3600 }));
        let second = manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");

        assert_eq!(
            second, "at-2",
            "the rejected token must not be served again"
        );
        assert_eq!(idp.hits(), 2, "the invalidation must force a token request");
    }

    /// An invalidation that lands *while* a token request is in flight rejects
    /// the token that request is about to publish.
    ///
    /// This is why the generation is read before the fetch and stamped onto
    /// the result, rather than re-read after it: a fresh read would say "no
    /// invalidations since I published", which is false — one arrived in the
    /// window, and the token being published is exactly the one it rejected.
    /// Stale-stamping is the safe direction; the next caller refetches.
    #[tokio::test]
    async fn an_invalidation_during_a_fetch_rejects_the_token_it_publishes() {
        let idp = IdpState::ok(json!({ "access_token": "at-1", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = Arc::new(manager_with(Arc::new(StubConnectorRepo::with(vec![]))));
        let cfg = Arc::new(cc_config(&url));

        // Stand where the in-flight acquisition stands: the generation has
        // been read, the token not yet published.
        let entry = manager.entry("crm").await;
        let observed = entry.invalidated.load(std::sync::atomic::Ordering::Acquire);

        // The 401 arrives now, before the publish.
        manager.invalidate("crm").await;

        // …and the acquisition publishes, stamped with what it observed.
        {
            let mut st = entry.state.lock().await;
            st.generation = observed;
            st.token = Some(CachedToken {
                access_token: "at-stale".to_string(),
                expires_at: Instant::now() + Duration::from_secs(3600),
                expires_at_epoch_ms: 0,
            });
            entry.fast.store(Arc::new(FastToken {
                config: Some((*cfg).clone()),
                token: st.token.clone(),
                generation: observed,
            }));
        }

        idp.set_response(200, json!({ "access_token": "at-2", "expires_in": 3600 }));
        let next = manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");

        assert_eq!(
            next, "at-2",
            "a token published after the 401 that rejected it must not be served"
        );
        assert_eq!(idp.hits(), 1, "exactly one refetch");
    }

    /// The generation does not make every call refetch: with no invalidation,
    /// the fast path still answers from cache.
    #[tokio::test]
    async fn the_generation_does_not_disturb_the_cached_fast_path() {
        let idp = IdpState::ok(json!({ "access_token": "at-1", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = cc_config(&url);

        for _ in 0..5 {
            assert_eq!(
                manager
                    .access_token("crm", &cfg, true)
                    .await
                    .expect("token"),
                "at-1"
            );
        }
        assert_eq!(idp.hits(), 1, "five calls, one token request");
    }

    // ---- refresh_token: seed, rotation, persistence, adoption ----

    #[tokio::test]
    async fn rotation_persists_and_the_next_refresh_uses_the_new_token() {
        let idp = IdpState::ok(json!({
            "access_token": "at-1", "expires_in": 3600, "refresh_token": "rt-2"
        }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let repo = Arc::new(StubConnectorRepo::with(vec![]));
        let manager = manager_with(Arc::clone(&repo) as Arc<StubConnectorRepo>);
        let cfg = rt_config(&url, "rt-1");

        manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");
        assert_eq!(
            idp.last_refresh_token.lock().expect("test").as_deref(),
            Some("rt-1"),
            "the first refresh presents the seed"
        );

        // The rotation was persisted with the fingerprint.
        let row = repo
            .get_oauth_state("crm")
            .await
            .expect("read")
            .expect("state row");
        assert_eq!(row.fingerprint, fingerprint(&cfg));
        let persisted: PersistedState = serde_json::from_str(&row.state_json).expect("state json");
        assert_eq!(persisted.refresh_token, "rt-2");
        assert_eq!(persisted.access_token, "at-1");

        // Invalidate the access token: the next refresh presents the
        // *rotated* token, not the burned seed.
        manager.invalidate("crm").await;
        idp.set_response(
            200,
            json!({ "access_token": "at-2", "expires_in": 3600, "refresh_token": "rt-3" }),
        );
        manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");
        assert_eq!(
            idp.last_refresh_token.lock().expect("test").as_deref(),
            Some("rt-2")
        );
    }

    /// A fresh process (or another node) adopts persisted state instead of
    /// spending a refresh — zero IdP traffic.
    #[tokio::test]
    async fn persisted_state_is_adopted_without_touching_the_idp() {
        let idp = IdpState::ok(json!({ "access_token": "never", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let repo = Arc::new(StubConnectorRepo::with(vec![]));
        let cfg = rt_config(&url, "rt-seed");
        let state = PersistedState {
            access_token: "adopted".to_string(),
            expires_at_epoch_ms: epoch_ms_now() + 3_600_000,
            refresh_token: "rt-live".to_string(),
        };
        repo.put_oauth_state(
            "crm",
            &fingerprint(&cfg),
            &serde_json::to_string(&state).expect("json"),
        )
        .await
        .expect("seed state");

        let manager = manager_with(Arc::clone(&repo) as Arc<StubConnectorRepo>);
        let token = manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");
        assert_eq!(token, "adopted");
        assert_eq!(idp.hits(), 0, "adoption costs no token request");
    }

    /// Editing the connector discards state minted under the old config —
    /// which is exactly the burned-token recovery story: re-seed, refresh.
    #[tokio::test]
    async fn a_fingerprint_mismatch_falls_back_to_the_seed() {
        let idp = IdpState::ok(json!({
            "access_token": "at", "expires_in": 3600, "refresh_token": "rt-next"
        }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let repo = Arc::new(StubConnectorRepo::with(vec![]));
        let stale = PersistedState {
            access_token: "stale".to_string(),
            expires_at_epoch_ms: epoch_ms_now() + 3_600_000,
            refresh_token: "rt-stale".to_string(),
        };
        repo.put_oauth_state(
            "crm",
            "an-old-fingerprint",
            &serde_json::to_string(&stale).expect("json"),
        )
        .await
        .expect("seed state");

        let manager = manager_with(Arc::clone(&repo) as Arc<StubConnectorRepo>);
        let cfg = rt_config(&url, "rt-reseeded");
        manager
            .access_token("crm", &cfg, true)
            .await
            .expect("token");
        assert_eq!(
            idp.last_refresh_token.lock().expect("test").as_deref(),
            Some("rt-reseeded"),
            "stale state must lose to the freshly-seeded config"
        );
    }

    // ---- failure taxonomy ----

    /// `invalid_grant` is non-retryable and negative-cached: a burned token
    /// is never retry-looped against the IdP.
    #[tokio::test]
    async fn a_rejection_is_negative_cached_and_non_retryable() {
        let idp = IdpState::ok(json!({})); // replaced below
        idp.set_response(
            400,
            json!({ "error": "invalid_grant", "error_description": "revoked" }),
        );
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = rt_config(&url, "rt-burned");

        let err = manager
            .access_token("crm", &cfg, true)
            .await
            .expect_err("burned token");
        assert!(!err.retryable(), "{err}");
        assert!(err.to_string().contains("invalid_grant"), "{err}");
        assert!(err.to_string().contains("re-seed"), "{err}");
        assert!(
            !err.to_string().contains("revoked"),
            "the free-text description is logged, never surfaced: {err}"
        );

        for _ in 0..3 {
            let err = manager
                .access_token("crm", &cfg, true)
                .await
                .expect_err("still cached");
            assert!(!err.retryable(), "{err}");
        }
        assert_eq!(idp.hits(), 1, "the rejection is answered from memory");
    }

    #[tokio::test]
    async fn a_server_error_is_retryable_and_not_cached() {
        let idp = IdpState::ok(json!({}));
        idp.set_response(503, json!({ "error": "temporarily_unavailable" }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let cfg = cc_config(&url);

        for _ in 0..2 {
            let err = manager
                .access_token("crm", &cfg, true)
                .await
                .expect_err("5xx");
            assert!(err.retryable(), "{err}");
        }
        assert_eq!(idp.hits(), 2, "transport failures are not negative-cached");
    }

    #[tokio::test]
    async fn config_mistakes_are_named_without_a_request() {
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let mut cfg = cc_config("http://127.0.0.1:9/token");
        cfg.grant = "password".to_string();
        let err = manager
            .access_token("crm", &cfg, true)
            .await
            .expect_err("ROPC is not a thing here");
        assert!(err.to_string().contains("password"), "{err}");
        assert!(err.to_string().contains(OAuth2Grant::VALUES), "{err}");

        let cfg = OAuth2Config {
            refresh_token: None,
            ..rt_config("http://127.0.0.1:9/token", "x")
        };
        let err = manager
            .access_token("crm", &cfg, true)
            .await
            .expect_err("no seed");
        assert!(err.to_string().contains("refresh_token"), "{err}");
    }

    /// An uninitialised manager (a bare registry in a unit test) refuses
    /// clearly instead of panicking.
    #[tokio::test]
    async fn an_uninitialised_manager_refuses_cleanly() {
        let manager = OAuthTokenManager::new();
        let err = manager
            .access_token("crm", &cc_config("http://x/token"), true)
            .await
            .expect_err("no deps");
        assert!(err.to_string().contains("not initialised"), "{err}");
        assert!(!err.retryable());
    }

    // ---- effective_auth: the seam the request paths use ----

    #[tokio::test]
    async fn effective_auth_passes_static_variants_through() {
        let manager = OAuthTokenManager::new();
        let http = crate::connector::HttpConnectorConfig {
            auth: Some(AuthConfig::Bearer {
                token: "static".to_string(),
            }),
            ..http_config_base()
        };
        let auth = effective_auth(&manager, "crm", &http)
            .await
            .expect("static auth needs no runtime");
        assert!(matches!(
            auth.as_deref(),
            Some(AuthConfig::Bearer { token }) if token == "static"
        ));

        let no_auth = crate::connector::HttpConnectorConfig {
            auth: None,
            ..http_config_base()
        };
        assert!(
            effective_auth(&manager, "crm", &no_auth)
                .await
                .expect("no auth")
                .is_none()
        );
    }

    #[tokio::test]
    async fn effective_auth_resolves_oauth2_to_a_bearer() {
        let idp = IdpState::ok(json!({ "access_token": "resolved", "expires_in": 3600 }));
        let url = fake_idp(Arc::clone(&idp)).await;
        let manager = manager_with(Arc::new(StubConnectorRepo::with(vec![])));
        let http = crate::connector::HttpConnectorConfig {
            auth: Some(AuthConfig::OAuth2(Box::new(cc_config(&url)))),
            allow_private_urls: true,
            ..http_config_base()
        };
        let auth = effective_auth(&manager, "crm", &http)
            .await
            .expect("resolves");
        assert!(matches!(
            auth.as_deref(),
            Some(AuthConfig::Bearer { token }) if token == "resolved"
        ));
    }

    fn http_config_base() -> crate::connector::HttpConnectorConfig {
        serde_json::from_value(json!({
            "url": "https://api.example.com",
            "method": "GET"
        }))
        .expect("base http config")
    }

    /// A token endpoint that streams without declaring a length cannot make
    /// this buffer past `MAX_TOKEN_RESPONSE_BYTES`.
    ///
    /// The declared-length check was already here and is no cap on its own: a
    /// chunked response skips it entirely, and the check that followed ran on
    /// a body `bytes()` had already read to its end. This path also runs on
    /// every token refresh for every OAuth2 connector, so the flood is
    /// repeatable on a schedule Orion sets itself.
    #[tokio::test]
    async fn a_flooding_token_endpoint_is_cut_off_at_the_cap() {
        const CHUNK: usize = 64 * 1024;
        const CHUNKS: usize = 128; // 8 MiB against a 64 KiB cap
        let (url, server) = crate::http_body::flood_server(CHUNK, CHUNKS).await;

        let err = request_token_inner(
            &reqwest::Client::new(),
            TokenEndpoint {
                token_url: &url,
                client_id: "cid",
                client_secret: "csecret",
                client_auth: "basic",
            },
            // The flood server is on loopback; the SSRF check would refuse it
            // before a body was read.
            true,
            vec![("grant_type", "client_credentials".to_string())],
        )
        .await
        .expect_err("must refuse an oversized token response");

        assert!(
            matches!(err, OAuthError::Transport(ref m) if m.contains("token response")),
            "unexpected error: {err:?}"
        );
        crate::http_body::assert_stopped_early(
            server.await.expect("test server"),
            CHUNK * CHUNKS,
            "the token request",
        );
    }
}