feather-reader 0.4.6

A minimalist, atproto-native RSS/Atom reader in Rust — your feed subscriptions live in your own PDS.
Documentation
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
1653
1654
1655
1656
1657
1658
1659
1660
1661
1662
1663
1664
1665
1666
1667
1668
1669
1670
1671
1672
1673
1674
1675
1676
1677
1678
1679
1680
1681
1682
1683
1684
1685
1686
1687
1688
1689
1690
1691
1692
1693
1694
1695
1696
1697
1698
1699
1700
1701
1702
1703
1704
1705
1706
1707
1708
1709
1710
1711
1712
1713
1714
1715
1716
1717
1718
1719
1720
1721
1722
1723
1724
1725
1726
1727
1728
1729
1730
1731
1732
1733
1734
1735
1736
1737
1738
1739
1740
1741
1742
1743
1744
1745
1746
1747
1748
1749
1750
1751
1752
1753
1754
1755
1756
1757
1758
1759
1760
1761
1762
1763
1764
1765
1766
1767
1768
1769
1770
1771
1772
1773
1774
1775
1776
1777
1778
1779
1780
1781
1782
1783
1784
1785
1786
1787
1788
1789
1790
1791
1792
1793
1794
1795
1796
1797
1798
1799
1800
1801
1802
1803
1804
1805
1806
1807
1808
1809
1810
1811
1812
1813
1814
1815
1816
1817
1818
1819
1820
1821
1822
1823
1824
1825
1826
1827
1828
1829
1830
1831
1832
1833
1834
1835
1836
1837
1838
1839
1840
1841
1842
1843
1844
1845
1846
1847
1848
1849
1850
1851
1852
1853
1854
1855
1856
1857
1858
1859
1860
1861
1862
1863
1864
1865
1866
1867
1868
1869
1870
1871
1872
1873
1874
1875
1876
1877
1878
1879
1880
1881
1882
1883
1884
1885
1886
1887
1888
1889
1890
1891
1892
1893
1894
1895
1896
1897
1898
1899
1900
1901
1902
1903
1904
1905
1906
1907
1908
1909
1910
1911
1912
1913
1914
1915
1916
1917
1918
1919
1920
1921
1922
1923
1924
1925
1926
1927
1928
1929
1930
1931
1932
1933
1934
1935
1936
1937
1938
1939
1940
1941
1942
1943
1944
1945
1946
1947
1948
1949
1950
1951
1952
1953
1954
1955
1956
1957
1958
1959
1960
1961
1962
1963
1964
1965
1966
1967
1968
1969
1970
1971
1972
1973
1974
1975
1976
1977
1978
1979
1980
1981
1982
1983
1984
1985
1986
1987
1988
1989
1990
1991
1992
1993
1994
1995
1996
1997
1998
1999
2000
2001
2002
2003
2004
2005
2006
2007
2008
2009
2010
2011
2012
2013
2014
2015
2016
2017
2018
2019
2020
2021
2022
2023
2024
2025
2026
2027
2028
2029
2030
2031
2032
2033
2034
2035
2036
2037
2038
2039
2040
2041
2042
2043
2044
2045
2046
2047
2048
2049
2050
2051
2052
2053
2054
2055
2056
2057
2058
2059
2060
2061
2062
2063
2064
2065
2066
2067
2068
2069
2070
2071
2072
2073
2074
2075
2076
2077
2078
2079
2080
2081
2082
2083
2084
2085
2086
2087
2088
2089
2090
2091
2092
2093
2094
2095
2096
2097
2098
2099
2100
2101
2102
2103
2104
2105
2106
2107
2108
2109
2110
2111
2112
2113
2114
2115
2116
2117
2118
2119
2120
2121
2122
2123
2124
2125
2126
2127
2128
2129
2130
2131
2132
2133
2134
2135
2136
2137
2138
2139
2140
2141
2142
2143
2144
2145
2146
2147
2148
2149
2150
2151
2152
2153
2154
2155
2156
2157
2158
2159
2160
2161
2162
2163
2164
2165
2166
2167
2168
2169
2170
2171
2172
2173
2174
2175
2176
2177
2178
2179
2180
2181
2182
2183
2184
2185
2186
2187
2188
2189
2190
2191
2192
2193
2194
2195
2196
2197
2198
2199
2200
2201
2202
2203
2204
2205
2206
2207
2208
2209
2210
2211
2212
2213
2214
2215
2216
2217
2218
2219
2220
2221
2222
2223
2224
2225
2226
2227
2228
2229
2230
2231
2232
2233
2234
2235
2236
2237
2238
2239
2240
2241
2242
2243
2244
2245
2246
2247
2248
2249
2250
2251
2252
2253
2254
2255
2256
2257
2258
2259
2260
2261
2262
2263
2264
2265
2266
2267
2268
2269
2270
2271
2272
2273
2274
2275
2276
2277
2278
2279
2280
2281
2282
2283
2284
2285
2286
2287
2288
2289
2290
2291
2292
2293
2294
2295
2296
2297
2298
2299
2300
2301
2302
2303
2304
2305
2306
2307
2308
2309
2310
2311
2312
2313
2314
2315
2316
2317
2318
2319
2320
2321
2322
2323
2324
2325
2326
2327
2328
2329
2330
2331
2332
2333
2334
2335
2336
2337
2338
2339
2340
2341
2342
2343
2344
2345
2346
2347
2348
2349
2350
2351
2352
2353
2354
2355
2356
2357
2358
2359
2360
2361
2362
2363
2364
2365
2366
2367
2368
2369
2370
2371
2372
2373
2374
2375
2376
//! Background schedulers — the **poll schedulers**, the **read-state flusher**,
//! and the sweeps and probe that ride alongside them.
//!
//! These are the long-lived `tokio` tasks that turn FeatherReader from a
//! request/response web app into a live reader: eight in all — the six
//! offset-bearing loops in `Loop::ALL` (RSS poller, publication poller, three
//! sweepers, adoption probe) plus the read-state flusher and the metrics
//! flusher. All are spawned from `main` after the [`AppState`] is built, behind
//! a config flag so tests and local dev can disable them, and all are
//! **graceful-shutdown-aware**: they select on a shutdown signal and drain
//! before returning.
//!
//! ## Poll scheduler ([`run_poller`])
//!
//! A single interval loop that, on each tick, asks the store for the feeds
//! whose `next_poll` is **due** ([`store::due_feeds_of_kind`], RSS only; standard.site
//! publications have their own loop, [`run_publication_poller`]) and polls each with
//! [`feed::poll_feed`] — which already does the conditional GET (`ETag` /
//! `Last-Modified`) and returns a [`feed::PollOutcome`]. The scheduler owns
//! **cadence**: `poll_feed` deliberately leaves `next_poll = None`, so after
//! each poll the scheduler computes the next-poll time from the feed's
//! `fetchHint` cadence (or the configured default) — and on failure honours the
//! **backoff** the outcome carries. Polls are **staggered / rate-limited** with
//! a bounded [`Semaphore`] and a small per-launch delay, so a batch of due
//! feeds does not stampede.
//!
//! ## Read-state flusher ([`run_flusher`])
//!
//! A **debounced** loop (default ~60 s, `Config`-tunable via the env) that
//! scans the store for **dirty** per-feed read cursors ([`store::dirty_cursors`])
//! across every DID, coalesces each DID's dirty cursors into **one**
//! `com.atproto.repo.applyWrites` batch via
//! `SidecarClient::flush_read_states` (sent as several calls past 200 cursors
//! or 128 KiB — `atproto::apply_writes_chunked`), and — only on success — clears the
//! `dirty` flag ([`store::clear_cursor_dirty`]). Dozens of articles read in one
//! sitting collapse into one write per feed (one record per feed, keyed by a
//! feed-derived rkey), and several feeds' cursors ride one round-trip. It also
//! flushes **once more on graceful shutdown** so a Ctrl-C never strands unsynced
//! read-state.
//!
//! ## Relay adoption probe ([`run_adoption_probe`])
//!
//! One unauthenticated GET per configured relay per day, counting the repos on
//! the public atproto network that hold `community.lexicon.rss.subscription`
//! (see `design/NETWORK-SPEC.md` §4 and [`feather_reader::network`]). It holds no
//! personal data, writes one `network_stat` row per relay, and cannot affect the
//! reader: every failure is a `warn!` that leaves the previous observation in
//! place. `FEATHERREADER_ADOPTION_INTERVAL_SECS=0` disables it on its own.

use std::sync::Arc;
use std::time::Duration;

use chrono::{SecondsFormat, Utc};
use sqlx::Row;
use tokio::sync::{watch, Semaphore};
use tokio::time::{interval, MissedTickBehavior};
use tracing::{debug, error, info, warn};

use feather_reader::feed::{self, PollOutcome};
use std::collections::HashSet;

use feather_reader::lexicon::nsid;
use feather_reader::network::RelayClient;
use feather_reader::readstate::{flush_did, fnv1a_64};
use feather_reader::store::{self, Feed, Pool};
use feather_reader::AppState;

// ---------------------------------------------------------------------------
// Tunables (env-overridable so tests/dev can move fast; sane defaults)
// ---------------------------------------------------------------------------

/// How often the poll scheduler wakes to look for due feeds. This is the *loop*
/// cadence, not the per-feed poll interval — a feed is only fetched when its own
/// `next_poll` is due. Overridable via `FEATHERREADER_POLL_TICK_SECS`.
const DEFAULT_POLL_TICK: Duration = Duration::from_secs(60);

/// Max feeds pulled off the due queue per tick — bounds the burst of work a
/// single wake can schedule. Overridable via `FEATHERREADER_POLL_BATCH`.
const DEFAULT_POLL_BATCH: i64 = 50;

/// Max feeds fetched concurrently — the rate limit. Overridable via
/// `FEATHERREADER_POLL_CONCURRENCY`.
const DEFAULT_POLL_CONCURRENCY: usize = 4;

/// Small delay between *launching* each feed fetch, so a batch of due feeds is
/// staggered rather than fired in one instant (polite to the network + to any
/// single upstream). Overridable via `FEATHERREADER_POLL_STAGGER_MS`.
const DEFAULT_POLL_STAGGER: Duration = Duration::from_millis(250);

/// The read-state flush debounce window — a given DID's dirty cursors are
/// flushed at most once per this interval. ~60 s per the design. Overridable via
/// `FEATHERREADER_FLUSH_DEBOUNCE_SECS`.
const DEFAULT_FLUSH_DEBOUNCE: Duration = Duration::from_secs(60);

/// How often the invite-code TTL sweep runs, expiring `active` codes past their
/// `expires_at`. Hourly is plenty — expiry is coarse-grained and `redeem_code`
/// already rejects a past-expiry code at redeem time regardless of this sweep, so
/// this is just housekeeping. Overridable via `FEATHERREADER_CODE_SWEEP_SECS`.
const DEFAULT_CODE_SWEEP: Duration = Duration::from_secs(3600);

/// How often the retention sweep runs, deleting shared-cache entries older than
/// `config.retention_days`. Daily is plenty — the window is coarse (days) and the
/// per-feed `max_entries_per_feed` trim already bounds any single feed on every
/// poll. Overridable via `FEATHERREADER_RETENTION_SWEEP_SECS`.
const DEFAULT_RETENTION_SWEEP: Duration = Duration::from_secs(24 * 60 * 60);

/// Delay before the FIRST adoption probe after boot. Unlike the local sweeps,
/// this tick is an outbound request to somebody else's relay, and
/// `deploy/container-entrypoint.sh` tears the machine down (and Fly recreates it)
/// the moment any child exits — so an immediate first tick would probe the relay
/// once per crash-loop restart rather than once per day.
const ADOPTION_STARTUP_DELAY: Duration = Duration::from_secs(5 * 60);

/// Delay before each local loop's FIRST tick after boot.
///
/// All four used to fire immediately. Combined with the container supervisor —
/// which tears the machine down the moment any child exits, and Fly restarts it
/// — "once per boot" becomes "once per crash-loop restart", and the loops all
/// pile onto the same instant while the machine is still opening its database
/// and warming its caches. The adoption probe already reasoned about exactly
/// this ([`ADOPTION_STARTUP_DELAY`]); its three siblings did not.
///
/// The values are deliberately DISTINCT rather than jittered. There is exactly
/// one machine, so there is no fleet to de-synchronise; what matters is that the
/// loops do not land together, and fixed offsets give that property while
/// staying reproducible in a test. They are also short enough to be irrelevant
/// to an hourly poller and a daily sweep.
///
/// `FEATHERREADER_STARTUP_DELAY_SECS` scales all of them (0 restores the old
/// immediate-first-tick behaviour), for dev loops and integration tests that
/// cannot wait.
const POLLER_STARTUP_DELAY: Duration = Duration::from_secs(30);
const PENDING_SWEEP_STARTUP_DELAY: Duration = Duration::from_secs(45);
const CODE_SWEEP_STARTUP_DELAY: Duration = Duration::from_secs(60);
const RETENTION_STARTUP_DELAY: Duration = Duration::from_secs(90);
/// The standard.site publication poller's first tick (0.4.0). Between the
/// code sweeper's and the retention sweeper's, so no two loops share a second.
const PUBLICATION_POLLER_STARTUP_DELAY: Duration = Duration::from_secs(75);

/// Which background loop an offset belongs to.
///
/// **An enum, not a string key.** The first cut of this used `&'static str`
/// names looked up with `.unwrap_or(POLLER_STARTUP_DELAY)`, and a review showed
/// that was strictly WORSE than the per-loop constants it replaced: mistyping
/// `offset_for("pending-sweeper")` compiled, passed all 706 tests, and silently
/// moved that loop onto the poller's tick — the everything-at-once collision the
/// offsets exist to prevent. A wrong constant name used to be a compile error;
/// a wrong string was a silent production change.
///
/// With an enum and an exhaustive `match` the table is total by construction,
/// there is no fallback to be wrong, and a typo is a compile error again.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
enum Loop {
    Poller,
    PublicationPoller,
    PendingSweep,
    CodeSweep,
    Retention,
    Adoption,
}

impl Loop {
    /// The registry: every offset-bearing loop, in the order `spawn` starts them.
    ///
    /// A variant missing from here is never STARTED — a visible absence — rather
    /// than started with the wrong offset, which was silent.
    const ALL: [Loop; 6] = [
        Loop::Poller,
        Loop::PublicationPoller,
        Loop::PendingSweep,
        Loop::CodeSweep,
        Loop::Retention,
        Loop::Adoption,
    ];

    /// Start this loop, with its own offset.
    ///
    /// **The variant names the loop AND the offset, in one place.** The previous
    /// shape passed `offset_for(Loop::X)` as a positional argument at five
    /// near-identical `tokio::spawn` lines, ~600 lines from the loop it
    /// configured — so writing `offset_for(Loop::Poller)` at the pending-sweeper
    /// site compiled, passed 707 tests and clippy, and silently collided the two
    /// loops at boot. That is the third form of this same bug; the first two were
    /// a wrong constant and a mistyped string key.
    ///
    /// **What is still NOT prevented:** pairing a variant with the wrong `run_*`
    /// in an arm below — `Loop::PendingSweep => run_poller(..)` compiles and the
    /// suite passes. That is a different failure (one loop never starts, another
    /// runs twice) and it is narrower: one exhaustive match in one place, rather
    /// than five spawn sites scattered through the file. Catching it would need
    /// each loop to report that it started, i.e. production instrumentation for a
    /// test — not obviously worth it, but it is a hole, not an absence of one.
    fn spawn_at(
        self,
        state: AppState,
        shutdown: watch::Receiver<()>,
        startup: Duration,
    ) -> tokio::task::JoinHandle<()> {
        match self {
            Loop::Poller => tokio::spawn(run_poller(state, shutdown, startup)),
            Loop::PublicationPoller => {
                tokio::spawn(run_publication_poller(state, shutdown, startup))
            }
            Loop::PendingSweep => tokio::spawn(run_pending_sweeper(state, shutdown, startup)),
            Loop::CodeSweep => tokio::spawn(run_code_sweeper(state, shutdown, startup)),
            Loop::Retention => tokio::spawn(run_retention_sweeper(state, shutdown, startup)),
            Loop::Adoption => tokio::spawn(run_adoption_probe(state, shutdown, startup)),
        }
    }

    /// This loop's startup offset. Exhaustive: adding a variant without an
    /// offset does not compile.
    const fn startup_offset(self) -> Duration {
        match self {
            Loop::Poller => POLLER_STARTUP_DELAY,
            Loop::PublicationPoller => PUBLICATION_POLLER_STARTUP_DELAY,
            Loop::PendingSweep => PENDING_SWEEP_STARTUP_DELAY,
            Loop::CodeSweep => CODE_SWEEP_STARTUP_DELAY,
            Loop::Retention => RETENTION_STARTUP_DELAY,
            Loop::Adoption => ADOPTION_STARTUP_DELAY,
        }
    }
}

/// The loop that polls `kind`, or `None` for a kind nothing polls.
///
/// Exhaustive on purpose: in a test build, a new `FeedKind` does not compile
/// until it is given a loop here or explicitly none, and
/// `every_pollable_kind_has_a_poller` checks the POLLABLE kinds against that.
/// **It is a record, not a mechanism**: each loop selects its own kind
/// directly, and nothing in production reads this map.
#[cfg(test)]
const fn poller_for(kind: feed::FeedKind) -> Option<Loop> {
    match kind {
        feed::FeedKind::Rss => Some(Loop::Poller),
        feed::FeedKind::Publication => Some(Loop::PublicationPoller),
        feed::FeedKind::Unsupported => None,
    }
}

/// The startup-override variable. The one place its name is written in
/// production code.
///
/// **A constant, because the name was untested glue.** [`startup_plan`] is the
/// production path and no test can reach it without `set_var`, so mistyping the
/// key by one character left the whole suite AND clippy green while silently
/// disabling the override for every loop. The test below spells the name
/// independently, so the two have to agree.
const STARTUP_DELAY_ENV: &str = "FEATHERREADER_STARTUP_DELAY_SECS";

/// The `(loop, resolved startup offset)` pairs `spawn` starts, in registry order.
///
/// **The offset decision is a VALUE a test can read, not a call buried in a
/// match arm.** It used to live inside `spawn_with`, as `offset_for(self)`.
/// Writing `offset_for(Loop::Poller)` there compiled, passed 697 lib + 15 bin
/// and clippy, and booted every loop at 30 s — the all-at-once collision this
/// file's whole history exists to prevent. Nothing could see it, because no test
/// can observe an argument passed inside a spawn arm.
///
/// Pulling it out here does not make the wrong thing unwritable — nothing can —
/// but it moves the decision somewhere a test can compare against an
/// independently spelled table. `spawn_at` now only chooses the `run_*`, which
/// is the one hole this file still documents rather than closes.
fn startup_plan() -> Vec<(Loop, Duration)> {
    startup_plan_from(std::env::var(STARTUP_DELAY_ENV).ok())
}

/// [`startup_plan`] with the environment value passed in.
///
/// Split for the same reason `offset_from` and `startup_delay_from` are: so the
/// composition is testable without `set_var`, which is a documented data race
/// against the ~39 `env::var` reads in this binary.
fn startup_plan_from(raw: Option<String>) -> Vec<(Loop, Duration)> {
    let mut plan = Vec::new();
    for_each_loop(|l| plan.push((l, offset_from(l, raw.clone()))));
    plan
}

/// The historical `offset_for`, with the environment value passed in.
///
/// Split for the same reason `startup_delay_from` is: so the COMPOSITION —
/// this loop's offset, then the ceiling — is testable without reading the
/// environment. Deleting the ceiling here disables
/// `FEATHERREADER_STARTUP_DELAY_SECS` for every loop at once, and asserting it
/// through `offset_for` could not catch that without `set_var`, which is a
/// documented data race against the ~39 `env::var` reads in this binary.
fn offset_from(which: Loop, raw: Option<String>) -> Duration {
    startup_delay_from(which.startup_offset(), raw)
}

/// Apply the `FEATHERREADER_STARTUP_DELAY_SECS` override to a startup delay,
/// with the environment value passed in.
///
/// The variable is a CEILING, not a replacement: it can only shorten the wait,
/// so setting it cannot accidentally push a production loop out further than its
/// constant intends.
///
/// Takes the raw value rather than reading it, so the ceiling behaviour is
/// testable without `std::env::set_var` — a documented data race against the ~39
/// `std::env::var` reads elsewhere in this binary, and this was the only
/// `set_var` in `src/`, in a 650-test multithreaded runner.
///
/// Split out so the ceiling behaviour is testable without `std::env::set_var`,
/// which is a documented data race against the ~39 `std::env::var` reads
/// elsewhere in this binary — and this was the only `set_var` in `src/`, in a
/// 650-test multithreaded runner. Latent today because nothing else reads that
/// key; a flaky crash the moment something does.
fn startup_delay_from(default: Duration, raw: Option<String>) -> Duration {
    match raw.as_deref().map(str::trim) {
        Some(v) => match v.parse::<u64>() {
            Ok(secs) => {
                let requested = Duration::from_secs(secs);
                if requested > default {
                    // The variable is a CEILING, which is a good property and a
                    // surprising one: setting it to 300 changes nothing. Saying
                    // so beats leaving an operator to wonder why.
                    info!(
                        requested_secs = secs,
                        effective_secs = default.as_secs(),
                        "FEATHERREADER_STARTUP_DELAY_SECS is a ceiling and can only \
                         SHORTEN a startup delay; using the built-in value"
                    );
                    return default;
                }
                requested
            }
            Err(_) => {
                warn!(
                    value = v,
                    "FEATHERREADER_STARTUP_DELAY_SECS is not a number; ignoring it"
                );
                default
            }
        },
        None => default,
    }
}

/// An `Interval` whose first tick is `delay` from now, then every `period`,
/// skipping missed ticks rather than bursting to catch up.
fn delayed_interval(delay: Duration, period: Duration) -> tokio::time::Interval {
    let mut ticker = tokio::time::interval_at(tokio::time::Instant::now() + delay, period);
    ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
    ticker
}

/// Read a `Duration` (in seconds) from the environment, or fall back.
fn env_duration_secs(key: &str, default: Duration) -> Duration {
    match std::env::var(key)
        .ok()
        .and_then(|v| v.trim().parse::<u64>().ok())
    {
        Some(secs) if secs > 0 => Duration::from_secs(secs),
        _ => default,
    }
}

/// Read a `u64`/`usize`/`i64` scalar from the environment, or fall back.
fn env_scalar<T: std::str::FromStr>(key: &str, default: T) -> T {
    std::env::var(key)
        .ok()
        .and_then(|v| v.trim().parse::<T>().ok())
        .unwrap_or(default)
}

// ---------------------------------------------------------------------------
// Config flag — is the background machinery enabled?
// ---------------------------------------------------------------------------

/// Whether the background schedulers should run.
///
/// Defaults to **on** for a real deployment, but is disabled when
/// `FEATHERREADER_DISABLE_SCHEDULER` is truthy (`1`/`true`/`yes`/`on`) — the
/// seam tests and pure-web local runs use so they don't spin poll/flush loops.
/// Kept here (not in `Config`) so this task owns its own flag and touches no
/// other module.
pub fn schedulers_enabled() -> bool {
    match std::env::var("FEATHERREADER_DISABLE_SCHEDULER") {
        Ok(v) => !matches!(
            v.trim().to_ascii_lowercase().as_str(),
            "1" | "true" | "yes" | "on"
        ),
        Err(_) => true,
    }
}

// ---------------------------------------------------------------------------
// Spawn helper — wire both tasks to a shared shutdown signal
// ---------------------------------------------------------------------------

/// Spawn the poll scheduler and the read-state flusher as detached `tokio`
/// tasks, each wired to the same graceful-shutdown signal.
///
/// Returns immediately with the [`JoinHandle`](tokio::task::JoinHandle)s so the
/// caller *may* await them at shutdown; `main` typically fires-and-forgets since
/// the shutdown channel is what actually stops them. A no-op (returns an empty
/// vec) when [`schedulers_enabled`] is false.
///
/// `shutdown` is a `watch` receiver that fires when the process is asked to stop
/// (the same signal `axum::serve` uses for graceful shutdown). Each task takes
/// its own clone of the receiver.
pub fn spawn(state: AppState, shutdown: watch::Receiver<()>) -> Vec<tokio::task::JoinHandle<()>> {
    if !schedulers_enabled() {
        info!("background schedulers disabled (FEATHERREADER_DISABLE_SCHEDULER)");
        // Recorded so a handler can tell "never started" from "started and
        // stopped ticking" — identical from the outside, opposite responses.
        state.runtime_health.set_schedulers_enabled(false);
        return Vec::new();
    }
    state.runtime_health.set_schedulers_enabled(true);

    info!(
        "spawning background schedulers (RSS poller + publication poller + sweepers + adoption probe + read-state flusher)"
    );

    // **Driven off the registry, not five hand-written lines.** Every offset-
    // bearing loop is started here by iterating `Loop::ALL`, so a loop cannot be
    // given another's offset — the variant chooses both.
    let mut handles: Vec<tokio::task::JoinHandle<()>> = Vec::new();
    for (l, startup) in startup_plan() {
        handles.push(l.spawn_at(state.clone(), shutdown.clone(), startup));
    }

    // The two loops with NO startup offset. `run_metrics_flusher` deliberately
    // fires immediately at boot; `run_flusher` swallows its first tick. Neither
    // is in `Loop`, so neither is covered by the distinctness invariant — stated
    // here because the test's message would otherwise read as covering all loops.
    handles.push({
        let state = state.clone();
        let shutdown = shutdown.clone();
        tokio::spawn(async move { run_metrics_flusher(state, shutdown).await })
    });
    // Last: consumes the un-cloned `state` / `shutdown` by move.
    handles.push(tokio::spawn(
        async move { run_flusher(state, shutdown).await },
    ));

    handles
}

/// Call `f` once for every offset-bearing loop, in registry order.
///
/// **The iteration itself, extracted so a test can watch it.** `spawn` iterated
/// `Loop::ALL` inline and nothing in the tree reached `spawn` — so a `.filter()`
/// dropping one loop compiled, passed the whole suite, passed clippy, and the
/// pending-login sweeper simply never started while nonce rows accumulated
/// unbounded.
///
/// That is the FOURTH form of one defect in this file. Each fix closed the seam
/// a level down — a wrong constant, then a wrong string key, then a wrong
/// positional argument — while the untested glue moved a level up. This is the
/// level `spawn` actually decides at.
fn for_each_loop(mut f: impl FnMut(Loop)) {
    for l in Loop::ALL {
        f(l);
    }
}

/// Resolve when the `watch` channel fires (the shutdown broadcast) or its sender
/// is dropped — the shared "time to stop" signal both loops select on.
async fn shutdown_fired(rx: &mut watch::Receiver<()>) {
    let _ = rx.changed().await;
}

// ---------------------------------------------------------------------------
// Poll scheduler
// ---------------------------------------------------------------------------

/// The poll-scheduler loop. Wakes on an interval, selects due feeds, and polls
/// each (conditional-GET + backoff via [`feed::poll_feed`]), staggered and
/// concurrency-bounded. Returns when `shutdown` resolves.
pub async fn run_poller(state: AppState, mut shutdown: watch::Receiver<()>, startup: Duration) {
    let tick = env_duration_secs("FEATHERREADER_POLL_TICK_SECS", DEFAULT_POLL_TICK);
    let batch = env_scalar::<i64>("FEATHERREADER_POLL_BATCH", DEFAULT_POLL_BATCH).max(1);
    let concurrency =
        env_scalar::<usize>("FEATHERREADER_POLL_CONCURRENCY", DEFAULT_POLL_CONCURRENCY).max(1);
    // The stagger default is sub-second, so read the ms knob directly.
    let stagger = std::env::var("FEATHERREADER_POLL_STAGGER_MS")
        .ok()
        .and_then(|v| v.trim().parse::<u64>().ok())
        .map(Duration::from_millis)
        .unwrap_or(DEFAULT_POLL_STAGGER);

    info!(
        ?tick,
        batch,
        concurrency,
        ?stagger,
        default_interval = ?state.config.poll_interval,
        "poll scheduler started"
    );

    let client = match feed::build_client() {
        Ok(c) => c,
        Err(err) => {
            error!(%err, "poll scheduler: failed to build HTTP client; poller will not run");
            return;
        }
    };
    let limiter = Arc::new(Semaphore::new(concurrency));

    // Not an immediate first tick: see `POLLER_STARTUP_DELAY`. Missed ticks are
    // skipped rather than burst through, so a slow poll round does not queue up.
    let mut ticker = delayed_interval(startup, tick);

    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("poll scheduler: shutdown signal received, stopping");
                break;
            }
            _ = ticker.tick() => {
                if let Err(err) =
                    poll_due_once(&state, &client, &limiter, batch, stagger, &shutdown).await
                {
                    // A store-level error is worth logging, but must not kill the
                    // loop — the next tick retries.
                    error!(%err, "poll scheduler: tick failed");
                }
                // Heartbeat, stamped on COMPLETION — including after a failed
                // tick, which is the honest reading: the loop is alive and
                // erroring, which is a different condition from the loop being
                // wedged, and `/health` reports them differently. A tick that
                // hangs forever never reaches here, which is the point.
                state.runtime_health.poll_tick_completed(Utc::now().timestamp());
            }
        }
    }
}

/// Whether new fetching is paused at the DB-size watermark. Above it, nothing
/// new is fetched — RSS or publication — so a small box cannot be filled to a
/// crash by the pollers. Shared by both loops: the publication loop, the bigger
/// writer, once ran without it (found in review).
async fn over_watermark(state: &AppState) -> bool {
    let watermark = state.config.db_size_watermark_bytes;
    if watermark > 0 {
        match store::db_size_bytes(&state.db).await {
            Ok(size) if size >= watermark => {
                warn!(
                    db_size_bytes = size,
                    watermark_bytes = watermark,
                    "DB size at/above watermark: pausing new polling until it drops (retention/prune)"
                );
                // No VACUUM here. `db_size_bytes` is already freelist-aware
                // (page_count - freelist_count), so the daily retention sweep's
                // DELETE lowers the measured size and lifts the watermark WITHOUT
                // a full-file rewrite — and the sweeper reclaims after it prunes
                // (see `run_retention_sweeper`). Running VACUUM on every poll tick
                // while over the watermark was both redundant (freelist accounting
                // already reflects freed pages) and dangerous: a full VACUUM needs
                // free disk ~= the live DB size to write the new file, which is
                // exactly what's scarce under the disk pressure that tripped the
                // watermark.
                //
                // Recorded, not just logged. This is one of the two states that
                // stop feeds updating, and until now the per-tick `warn!` was its
                // ONLY trace — so `/stats` showed `overdue` climbing and
                // `polled_last_hour` falling with nothing to say which of the two
                // causes was responsible. See `runtime_health`.
                state.runtime_health.set_watermark(true);
                return true;
            }
            Ok(_) => state.runtime_health.set_watermark(false),
            Err(err) => warn!(%err, "could not read DB size for watermark check; polling anyway"),
        }
    }
    false
}

/// One poll round: select due feeds and poll each, concurrency-bounded and
/// staggered. Feed-level failures are handled per-feed (rescheduled with the
/// outcome's backoff); only a store-level failure to *select* propagates.
async fn poll_due_once(
    state: &AppState,
    client: &reqwest::Client,
    limiter: &Arc<Semaphore>,
    batch: i64,
    stagger: Duration,
    shutdown: &watch::Receiver<()>,
) -> anyhow::Result<()> {
    // DB-size watermark: above it, stop pulling NEW content so a small box can't
    // be filled to a crash by the poller. Reads/serving continue; only fetching
    // is paused. `<= 0` disables the watermark.
    if over_watermark(state).await {
        return Ok(());
    }

    let now = now_rfc3339();
    // RSS only: publications have their own loop (`run_publication_poller`),
    // so a slow publication read can never hold this tick.
    let due = store::due_feeds_of_kind(&state.db, &now, feed::FeedKind::Rss, batch).await?;
    if due.is_empty() {
        debug!("poll scheduler: no feeds due");
        return Ok(());
    }
    info!(count = due.len(), "poll scheduler: polling due feeds");

    let mut handles = Vec::with_capacity(due.len());
    let mut abandoned = 0usize;
    for feed in due {
        // **Stop LAUNCHING once shutdown is asked for.**
        //
        // The loop above only checked shutdown between ticks, so once inside a
        // tick this ran to completion: a full batch is 50 feeds at a 250 ms
        // stagger — 12.5 s just to launch — against Fly's default 5 s
        // `kill_timeout`. Every feed in the batch has already been leased an hour
        // forward by `poll_and_reschedule`, and SIGKILL rolls nothing back, so a
        // routine deploy landing mid-tick silently pushed up to 50 feeds out by
        // an hour. Nobody would attribute that: it presents as "some feeds are
        // behind after a deploy", and `/stats` cannot show it as overdue because
        // `next_poll` was moved FORWARD.
        //
        // Feeds not launched keep whatever `next_poll` they had, so they stay due
        // and the next boot picks them up immediately.
        if shutdown.has_changed().unwrap_or(true) {
            abandoned += 1;
            continue;
        }
        let pool = state.db.clone();
        let client = client.clone();
        let default_interval = state.config.poll_interval;
        let config = Arc::clone(&state.config);
        // Acquire a permit *before* launching so at most `concurrency`
        // fetches are ever in flight; the permit is released when the task
        // ends.
        let permit = match Arc::clone(limiter).acquire_owned().await {
            Ok(p) => p,
            Err(_) => break, // semaphore closed — shutting down
        };
        handles.push(tokio::spawn(async move {
            let _permit = permit; // held for the duration of this poll
            poll_and_reschedule(&pool, &client, &feed, default_interval, &config).await;
        }));
        // Stagger launches so a batch doesn't fire in one instant.
        if !stagger.is_zero() {
            tokio::time::sleep(stagger).await;
        }
    }

    if abandoned > 0 {
        info!(
            abandoned,
            "poll scheduler: shutdown requested mid-round; these feeds were not \
             launched and stay due"
        );
    }

    // Drain the batch so the next tick starts from a clean slate.
    for h in handles {
        if let Err(err) = h.await {
            warn!(%err, "poll scheduler: a feed poll task panicked");
        }
    }
    Ok(())
}

/// The standard.site publication poller (0.4.0): its own loop, so the RSS
/// poller's tick never waits on a publication.
///
/// **One publication at a time, selected only when it is about to be read.**
/// An earlier shape ran publications as tasks off the RSS tick: they held the
/// tick for the sum of their reads, and once detached, a row whose lease ran
/// out while its task waited was queued again, without bound, and a queued row
/// was deferred an interval by every deploy (all found in review). Selecting
/// the next due row only when ready to read it has none of those: nothing waits
/// in a queue, a row not yet read keeps its `next_poll`, and the read in flight
/// at shutdown is the only one, bounded by `publication_read_deadline`.
///
/// Two things this does NOT do, on purpose: it has no `/health` heartbeat (the
/// RSS poller's is the one reported, and overdue publications still show in
/// `/stats`), and its one read runs ALONGSIDE the RSS poller's
/// `DEFAULT_POLL_CONCURRENCY`, not within it — one more fetch in flight than
/// before, bounded by the walk budget. (A subscribe from the form also runs one
/// publication read inline, outside both loops, bounded per request by the
/// per-IP rate limit on `/subscriptions`.)
pub async fn run_publication_poller(
    state: AppState,
    mut shutdown: watch::Receiver<()>,
    startup: Duration,
) {
    let client = match feed::build_client() {
        Ok(c) => c,
        Err(err) => {
            error!(%err, "publication poller: failed to build HTTP client; it will not run");
            return;
        }
    };
    // Publications became pollable in 0.4.0, and every row of a newly admitted
    // kind is due at once. Spread the never-polled ones across one interval so
    // they arrive as a trickle, not as one pass of N back-to-back reads.
    match store::stagger_unscheduled(
        &state.db,
        feed::FeedKind::Publication,
        state.config.poll_interval,
    )
    .await
    {
        Ok(0) => {}
        Ok(n) => info!(
            scheduled = n,
            "staggered the first poll of never-polled publications"
        ),
        Err(err) => {
            warn!(%err, "could not stagger never-polled publications; they will all be due at once")
        }
    }

    let tick = env_duration_secs("FEATHERREADER_POLL_TICK_SECS", DEFAULT_POLL_TICK);
    let mut ticker = delayed_interval(startup, tick);
    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("publication poller: shutdown signal received, stopping");
                break;
            }
            _ = ticker.tick() => {
                if let Err(err) = poll_publications_once(&state, &client, &shutdown).await {
                    error!(%err, "publication poller: pass failed");
                }
            }
        }
    }
}

/// Read every due publication, one at a time, until none is due or shutdown is
/// asked for. Each is selected only when it is about to be read, and leased by
/// `poll_and_reschedule` before its fetch, as every feed is.
async fn poll_publications_once(
    state: &AppState,
    client: &reqwest::Client,
    shutdown: &watch::Receiver<()>,
) -> anyhow::Result<()> {
    poll_publications_with(state, shutdown, |feeds, interval| async move {
        poll_group_and_reschedule(&state.db, client, &feeds, interval, &state.config).await;
    })
    .await
}

/// The most publications from one repo read together. A group shares one read
/// and one deadline, so it is bounded; the rest of a big repo's due
/// publications are read next, in another group.
const MAX_PUBLICATION_GROUP: usize = 16;

/// [`poll_publications_once`] with the read injected, so the pass's own rules
/// — which rows next, grouped how, how often, with which interval, when to
/// stop — are testable without a network, the way `poll_and_reschedule_with`
/// is.
async fn poll_publications_with<F, Fut>(
    state: &AppState,
    shutdown: &watch::Receiver<()>,
    mut poll: F,
) -> anyhow::Result<()>
where
    F: FnMut(Vec<Feed>, Duration) -> Fut,
    Fut: std::future::Future<Output = ()>,
{
    // **A row is read at most once per pass, and a row that is still due does
    // not stop the pass.** The next rows are re-selected after every read, and
    // a row whose next poll could not be written (a full or read-only volume)
    // stays due at the head of the order. Re-reading it was a tight loop
    // (thousands of reads a second); returning at it starved every publication
    // behind it, every pass (both found in review). So the pass skips what it
    // has read and moves on.
    let mut read = std::collections::HashSet::new();
    let repo_of =
        |f: &Feed| feather_reader::standard_site::AtUri::parse(&f.url).map(|u| u.authority);
    loop {
        if shutdown.has_changed().unwrap_or(true) || over_watermark(state).await {
            return Ok(());
        }
        let now = now_rfc3339();
        let due: Vec<Feed> = store::due_feeds_of_kind(
            &state.db,
            &now,
            feed::FeedKind::Publication,
            // Past every row already read, so rows that stay due after
            // their read cannot fill the window (found in review).
            i64::try_from(read.len() + 1_000).unwrap_or(i64::MAX),
        )
        .await?
        .into_iter()
        .filter(|f| !read.contains(&f.id))
        .collect();
        let Some(head) = due.first() else {
            return Ok(());
        };
        // **Every due publication in the head's repo, read together**: one
        // walk of that repo's documents instead of one per publication.
        let repo = repo_of(head);
        let group: Vec<Feed> = due
            .iter()
            .filter(|f| repo.is_some() && repo_of(f) == repo)
            .take(MAX_PUBLICATION_GROUP)
            .cloned()
            .collect();
        let group = if group.is_empty() {
            vec![head.clone()]
        } else {
            group
        };
        read.extend(group.iter().map(|f| f.id));
        poll(group, state.config.poll_interval).await;
    }
}

/// Lease every feed in a one-repo group, read the group once, and settle each
/// feed with its own outcome — `poll_and_reschedule_with`, for a group.
async fn poll_group_and_reschedule(
    pool: &Pool,
    client: &reqwest::Client,
    feeds: &[Feed],
    default_interval: Duration,
    config: &feather_reader::config::Config,
) {
    for f in feeds {
        if let Err(err) = store::set_next_poll(pool, &f.url, cadence_for(f, default_interval)).await
        {
            warn!(feed = %f.url, %err, "failed to lease next_poll before fetching");
        }
    }
    let outcomes = feed::poll_publication_group(pool, client, config, feeds).await;
    for (f, outcome) in feeds.iter().zip(outcomes) {
        match outcome {
            Ok(o) => feed::settle_poll(pool, &f.url, &o, cadence_for(f, default_interval)).await,
            Err(err) => {
                error!(feed = %f.url, %err, "polled: store error");
                if let Err(err) =
                    store::set_next_poll(pool, &f.url, cadence_for(f, default_interval)).await
                {
                    error!(feed = %f.url, %err, "failed to persist next_poll");
                }
            }
        }
    }
}

/// Poll one feed and persist its **next** poll time.
///
/// [`feed::poll_feed`] never returns `Err` for a merely-broken feed (only for a
/// broken local store), and it deliberately leaves `next_poll` unset — cadence
/// is the scheduler's job. So on every outcome we compute and store the next
/// poll time: the feed's cadence on success/not-modified, the outcome's backoff
/// on failure.
async fn poll_and_reschedule(
    pool: &Pool,
    client: &reqwest::Client,
    feed: &Feed,
    default_interval: Duration,
    config: &feather_reader::config::Config,
) {
    // By kind: an RSS feed is fetched over HTTP, a standard.site publication is
    // read from its author's PDS. One entry point, so the choice lives with
    // `FeedKind` and not here.
    poll_and_reschedule_with(pool, feed, default_interval, |pool, feed| {
        feed::poll_feed_by_kind(pool, client, config, feed)
    })
    .await;
}

/// [`poll_and_reschedule`] with the fetch injected, so the ordering guarantee
/// below can be tested without a network — the same shape the OAuth
/// orchestrators use.
async fn poll_and_reschedule_with<'a, F, Fut>(
    pool: &'a Pool,
    feed: &'a Feed,
    default_interval: Duration,
    poll: F,
) where
    F: FnOnce(&'a Pool, &'a Feed) -> Fut,
    Fut: std::future::Future<Output = anyhow::Result<PollOutcome>>,
{
    // **Lease the feed forward BEFORE fetching it.**
    //
    // Nothing used to be written to the feed row until after `poll_feed`
    // returned. `due_feeds` orders by `next_poll ASC` and the poller's first
    // tick fires immediately, so a feed whose fetch or parse takes the PROCESS
    // down was re-selected first on every restart — forever, with no escape
    // short of editing the database by hand.
    //
    // The trigger is plausible on a 512 MB box: `MAX_BODY_BYTES` is 8 MiB and
    // concurrency is 4, so 32 MiB of raw bodies can be in flight, and `feed_rs`
    // builds an in-memory model several times the wire size alongside the
    // sanitized `Vec<NewEntry>`. `deploy/container-entrypoint.sh` tears the
    // machine down the moment any child exits, and Fly restarts it — which is
    // what turns "one bad poll" into a loop.
    //
    // Writing the optimistic next time FIRST converts that permanent loop into
    // a single restart: the killer feed goes to the BACK of the due queue
    // instead of the front, every other feed gets polled, and the instance
    // heals itself. The value is the cadence the feed would have got had the
    // poll succeeded, so the common case — the poll returns and overwrites this
    // — is unchanged.
    //
    // The cost is one extra tiny UPDATE per feed per poll on a single-writer
    // database. Against an unrecoverable instance, that is not a close call.
    if let Err(err) =
        store::set_next_poll(pool, &feed.url, cadence_for(feed, default_interval)).await
    {
        // Non-fatal: the poll is still worth attempting. It just means a crash
        // during THIS fetch is not protected.
        warn!(feed = %feed.url, %err, "failed to lease next_poll before fetching; \
                                       a crash during this poll would re-select this feed first");
    }

    // One shared sequence for every outcome the poll produced: settle the
    // error columns and reschedule. `feed::settle_poll` is the only copy — a
    // second, inlined copy here is how `web::add_subscription` came to have a
    // third that did half the job.
    let outcome = poll(pool, feed).await;
    match &outcome {
        Ok(o) => feed::settle_poll(pool, &feed.url, o, cadence_for(feed, default_interval)).await,
        Err(err) => {
            // Store-level error for this feed — log and reschedule on the normal
            // cadence so we retry rather than getting stuck re-polling instantly.
            error!(feed = %feed.url, %err, "polled: store error");
            if let Err(err) =
                store::set_next_poll(pool, &feed.url, cadence_for(feed, default_interval)).await
            {
                error!(feed = %feed.url, %err, "failed to persist next_poll");
            }
        }
    }
}

/// The per-feed poll cadence. Honours the feed's `fetchHint` when the feed row
/// carries one; otherwise the configured default interval.
///
/// The `fetchHint` cadence hint (`realtime`/`hourly`/`daily`/`weekly`) lives on
/// the PDS-side `subscription` record. It is not yet projected onto the local
/// [`Feed`] row, so this maps the known values when present and otherwise falls
/// back to the config default — the mapping is factored out so wiring the
/// projected hint later is a one-line change.
fn cadence_for(feed: &Feed, default_interval: Duration) -> Duration {
    // `fetchHint` is not yet projected onto the local `feeds` row, so there is no
    // hint to read yet — this resolves to the configured default. The mapping is
    // routed through `cadence_from_hint` so wiring the projected hint later is a
    // one-line change here (pass `feed`'s hint instead of `None`).
    let hint: Option<&str> = feed_fetch_hint(feed);
    match hint {
        Some(h) => cadence_from_hint(h, default_interval),
        None => default_interval,
    }
}

/// The feed's `fetchHint`, if the local row carries one. The local `feeds` row
/// does not yet project the PDS-side hint, so this currently always returns
/// `None` — the single place to change when the hint column lands.
fn feed_fetch_hint(_feed: &Feed) -> Option<&str> {
    None
}

/// Map a `fetchHint` known-value to a poll cadence. Referenced by
/// [`cadence_for`] once the hint is projected onto the feed row; retained now so
/// the mapping is defined in one place and unit-tested.
fn cadence_from_hint(hint: &str, default_interval: Duration) -> Duration {
    match hint.trim().to_ascii_lowercase().as_str() {
        "realtime" => Duration::from_secs(5 * 60),
        "hourly" => Duration::from_secs(60 * 60),
        "daily" => Duration::from_secs(24 * 60 * 60),
        "weekly" => Duration::from_secs(7 * 24 * 60 * 60),
        _ => default_interval,
    }
}

// ---------------------------------------------------------------------------
// Invite-code TTL sweeper
// ---------------------------------------------------------------------------

/// The invite-code TTL sweep loop. On a periodic tick (hourly by default) it
/// flips every `active` invite code past its `expires_at` to `expired`
/// ([`store::expire_old_codes`]), keeping the closed-beta table tidy. Returns
/// when `shutdown` resolves. Failures are logged and never kill the loop — a
/// missed sweep is harmless because `redeem_code` re-checks expiry itself.
pub async fn run_code_sweeper(
    state: AppState,
    mut shutdown: watch::Receiver<()>,
    startup: Duration,
) {
    let period = env_duration_secs("FEATHERREADER_CODE_SWEEP_SECS", DEFAULT_CODE_SWEEP);
    info!(?period, "invite-code TTL sweeper started");

    // Delayed first tick (see `CODE_SWEEP_STARTUP_DELAY`) rather than the
    // immediate one this used to have — a long-stale set of codes is still swept
    // a minute into the boot, and `redeem_code` re-checks expiry itself, so the
    // sweep was never on the correctness path to begin with.
    let mut ticker = delayed_interval(startup, period);
    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("invite-code TTL sweeper: shutdown signal received, stopping");
                break;
            }
            _ = ticker.tick() => {
                match store::expire_old_codes(&state.db).await {
                    Ok(0) => debug!("invite-code TTL sweeper: nothing to expire"),
                    Ok(n) => info!(expired = n, "invite-code TTL sweeper: expired codes"),
                    Err(err) => error!(%err, "invite-code TTL sweeper: sweep failed"),
                }
            }
        }
    }
}

// ---------------------------------------------------------------------------
// Retention sweeper
// ---------------------------------------------------------------------------

/// The retention sweep loop. On a periodic tick (daily by default) it deletes
/// shared-cache entries older than `config.retention_days`
/// ([`store::prune_old_entries`]) — the mechanism that makes the README/wiki
/// "90-day rolling window" claim TRUE — and, after a sweep that actually deleted
/// rows, calls [`store::reclaim`] so the freed pages return to the OS (otherwise
/// the file never shrinks and the DB-size watermark can stay latched). Orphaned
/// entry ids are scrubbed from the affected `read_cursor` id-sets inside the
/// prune itself.
///
/// The loop runs if EITHER knob is on. `retention_days == 0` disables only the
/// rolling window; `retention_hard_days` still evicts everything past the
/// ceiling, and that is deliberate — the ceiling is what bounds the shared cache
/// for entries a reader pinned by starring or marking unread, and the per-feed
/// trim now spares those. Only when both are zero does the loop log once and
/// return, spawning no ticker; that configuration has no bound at all and says
/// so. Failures are logged and never kill the loop — a missed sweep just means
/// the window is enforced on the next tick.
pub async fn run_retention_sweeper(
    state: AppState,
    mut shutdown: watch::Receiver<()>,
    startup: Duration,
) {
    let days = state.config.retention_days as i64;
    let hard_days = state.config.retention_hard_days as i64;
    // The third window: the ceiling for kinds the rolling window does not apply
    // to (a publication). It is on by default and independent of the other two,
    // so the "disabled entirely" branch below has to consider it or a
    // `RETENTION_DAYS=0 RETENTION_HARD_DAYS=0` instance would silently stop
    // reaping publications as well.
    let publication_days = state.config.publication_retention_days as i64;
    if days <= 0 && hard_days <= 0 && publication_days <= 0 {
        info!(
            "retention sweeper: retention_days=0, retention_hard_days=0 and \
             publication_retention_days=0, retention disabled entirely (no rolling \
             window, NO ceiling — the shared cache is unbounded in this \
             configuration)"
        );
        return;
    }
    let period = env_duration_secs(
        "FEATHERREADER_RETENTION_SWEEP_SECS",
        DEFAULT_RETENTION_SWEEP,
    );
    info!(
        retention_days = days,
        retention_hard_days = hard_days,
        publication_retention_days = publication_days,
        ?period,
        "retention sweeper started"
    );

    // Delayed first tick (see `RETENTION_STARTUP_DELAY`). This is the heaviest
    // of the local loops — it takes the single write lock for the whole delete —
    // so firing it into a boot that is still opening the database and warming
    // caches was the worst timing available.
    let mut ticker = delayed_interval(startup, period);
    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("retention sweeper: shutdown signal received, stopping");
                break;
            }
            _ = ticker.tick() => {
                match store::prune_old_entries(&state.db, days, hard_days, publication_days).await {
                    Ok(0) => debug!("retention sweeper: nothing past the retention window"),
                    Ok(n) => {
                        info!(
                            pruned = n,
                            retention_days = days,
                            retention_hard_days = hard_days,
                            publication_retention_days = publication_days,
                            "retention sweeper: pruned old entries"
                        );
                        // Return the freed pages to the OS so the file actually
                        // shrinks and the DB-size watermark can fall back.
                        if let Err(err) = store::reclaim(&state.db).await {
                            warn!(%err, "retention sweeper: reclaim after prune failed");
                        }
                    }
                    Err(err) => error!(%err, "retention sweeper: prune failed"),
                }
            }
        }
    }
}

// ---------------------------------------------------------------------------
// Relay adoption probe
// ---------------------------------------------------------------------------

/// The relay adoption probe loop (`design/NETWORK-SPEC.md` §4). On a periodic
/// tick (daily by default) it asks every configured relay how many repos hold
/// `community.lexicon.rss.subscription`, logs the number — **the log is the
/// metric**; there is no `/metrics` endpoint — and upserts one `network_stat`
/// row per relay.
///
/// Two kill switches: `FEATHERREADER_ADOPTION_INTERVAL_SECS=0` (or an empty
/// `FEATHERREADER_RELAY_HOSTS`) stops just this loop, and the pre-existing
/// `FEATHERREADER_DISABLE_SCHEDULER` stops it with every other background task
/// (the five other offset loops in [`Loop::ALL`], the read-state flusher and the
/// metrics flusher).
///
/// It deviates from the other offset loops in exactly one way: `interval_at` with
/// an [`ADOPTION_STARTUP_DELAY`] (5 min) instead of a seconds-scale offset, because this
/// tick is a request to a third party and the container supervisor turns "once
/// per boot" into "once per crash-loop restart". Nothing it does can fail the
/// process: every error path is a `warn!` that leaves the previous observation
/// in place.
pub async fn run_adoption_probe(
    state: AppState,
    mut shutdown: watch::Receiver<()>,
    startup: Duration,
) {
    let period = state.config.adoption_interval;
    if period.is_zero() {
        info!("adoption probe: disabled (FEATHERREADER_ADOPTION_INTERVAL_SECS=0)");
        return;
    }
    // Rejected entries are surfaced HERE rather than at parse time: config is
    // read before `init_tracing` (steps 1 and 2 in `main`), so a warning emitted during
    // parsing would go nowhere. Warn whether or not any usable host survived —
    // a typo the operator never hears about is the failure mode this replaced a
    // boot abort with, and it must not be silent as well as non-fatal.
    for bad in &state.config.relay_host_errors {
        warn!(
            entry = %bad,
            "adoption probe: ignoring unusable FEATHERREADER_RELAY_HOSTS entry"
        );
    }
    if state.config.relay_hosts.is_empty() {
        info!("adoption probe: no relay hosts configured, probe disabled");
        return;
    }
    // A typo'd relay host must disable an optional metric, never block boot —
    // so the client is built here, in the task, not in `main`.
    let client = match RelayClient::new(state.http.clone(), &state.config.relay_hosts) {
        Ok(client) => client,
        Err(err) => {
            warn!(%err, "adoption probe: unusable FEATHERREADER_RELAY_HOSTS, probe disabled");
            return;
        }
    };

    let period = jittered(period, &state.config.public_url);
    info!(
        ?period,
        relays = client.hosts().len(),
        "adoption probe started"
    );

    let mut ticker = tokio::time::interval_at(tokio::time::Instant::now() + startup, period);
    ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("adoption probe: shutdown signal received, stopping");
                break;
            }
            _ = ticker.tick() => probe_adoption_once(&state, &client).await,
        }
    }
}

/// One probe run. Returns `()` — no error can reach the loop body, because none
/// of them is actionable: a relay outage, a 429, a malformed body, and a SQLite
/// write failure all leave the previous row in place and change nothing else.
/// (Deliberately unlike `poll_due_once`, which returns `Result` because a
/// store-level failure there IS a real signal.)
async fn probe_adoption_once(state: &AppState, client: &RelayClient) {
    let report = client.count_repos_with_collection(nsid::SUBSCRIPTION).await;

    for failure in &report.failures {
        warn!(
            host = %failure.host,
            reason = %failure.reason,
            "adoption probe: relay query failed; keeping previous observation"
        );
    }

    for obs in &report.observations {
        // §4.4: this log line IS the operator-facing metric.
        info!(
            key = store::ADOPTION_STAT_KEY,
            source = %obs.source,
            repos = obs.repos,
            truncated = obs.truncated,
            "adoption probe: observed"
        );
        let stat = store::NetworkStat {
            key: store::ADOPTION_STAT_KEY.to_string(),
            source: obs.source.clone(),
            // Bounded by MAX_PAGES × the page limit, far inside i64.
            value: obs.repos as i64,
            truncated: obs.truncated,
            observed_at: obs.observed_at.clone(),
        };
        if let Err(err) = store::record_network_stat(&state.db, &stat).await {
            warn!(%err, source = %obs.source, "adoption probe: failed to persist observation");
        }
    }

    // §4.1: two relays disagreeing is itself worth logging — non-archival
    // relays index different host sets, so this is information, not an error.
    if report.disagrees() {
        let counts: Vec<(String, u64)> = report
            .observations
            .iter()
            .map(|o| (o.source.clone(), o.repos))
            .collect();
        info!(
            ?counts,
            "adoption probe: relays disagree; the max is surfaced"
        );
    }
}

/// Spread the probe cadence ±10% so many self-hosted instances do not
/// synchronise on the relay.
///
/// Seeded from the instance's public URL through the FNV hash already in this
/// file, so it is **stable across restarts** (a restart must never re-roll into
/// a tighter cadence) and unit-testable — no `rand` dependency, no RNG in the
/// loop. The result is clamped to at least one second.
fn jittered(period: Duration, seed: &str) -> Duration {
    // 0..=200 → −100..=+100 tenths of a percent… i.e. ±10%.
    let basis = (fnv1a_64(seed.as_bytes()) % 201) as i64 - 100;
    let secs = period.as_secs_f64() * (1.0 + basis as f64 / 1000.0);
    Duration::from_secs_f64(secs.max(1.0))
}

// ---------------------------------------------------------------------------
// Pending-login sweeper
// ---------------------------------------------------------------------------

/// How often abandoned logins and stale nonces are swept.
const PENDING_SWEEP_SECS: u64 = 900;

/// How long an untouched DPoP nonce is kept. A server nonce lasts minutes; a day
/// is generous and keeps the table to the origins actually in use.
const NONCE_MAX_AGE_SECS: i64 = 24 * 60 * 60;

/// Delete expired pending logins.
///
/// An abandoned login — the user is redirected to their PDS and closes the tab —
/// leaves an `oauth_state` row behind. `take_pending` only ever consumes rows
/// that come BACK, so nothing else removes these, and each one holds a sealed
/// DPoP private key and a PKCE verifier. Without this the table grows without
/// bound and accumulates secret material that can no longer be used for
/// anything.
///
/// Runs on both backends: the rows are written by the Rust login path, and a
/// deployment that flips back to the sidecar still has whatever it left behind.
pub async fn run_pending_sweeper(
    state: AppState,
    mut shutdown: watch::Receiver<()>,
    startup: Duration,
) {
    let period = Duration::from_secs(PENDING_SWEEP_SECS);
    info!(?period, "pending-login sweeper started");

    // Delayed first tick, like its siblings (see `PENDING_SWEEP_STARTUP_DELAY`).
    let mut ticker = delayed_interval(startup, period);
    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("pending-login sweeper: shutdown signal received, stopping");
                break;
            }
            _ = ticker.tick() => {
                let now = Utc::now().timestamp();
                match feather_reader::oauth::store::sweep_expired_pending(&state.db, now).await {
                    Ok(0) => debug!("pending-login sweeper: nothing to expire"),
                    Ok(n) => info!(swept = n, "pending-login sweeper: removed abandoned logins"),
                    Err(err) => error!(%err, "pending-login sweeper: sweep failed"),
                }
                // Same volume, same pre-auth write primitive, and a stale nonce
                // is worthless — the server issues a new one with the next
                // challenge.
                match feather_reader::oauth::store::sweep_stale_nonces(
                    &state.db,
                    now - NONCE_MAX_AGE_SECS,
                )
                .await
                {
                    Ok(0) => debug!("nonce sweeper: nothing stale"),
                    Ok(n) => info!(swept = n, "nonce sweeper: removed stale DPoP nonces"),
                    Err(err) => error!(%err, "nonce sweeper: sweep failed"),
                }
            }
        }
    }
}

// ---------------------------------------------------------------------------
// Repo-timing flusher
// ---------------------------------------------------------------------------

/// How often buffered repo timings are written to SQLite.
///
/// Frequent enough that a crash loses little, rare enough that the write is
/// nowhere near the request path. Recording itself only touches memory.
const METRICS_FLUSH_SECS: u64 = 30;

/// Periodically persist buffered repo timings, and once more on shutdown.
///
/// The shutdown flush is the one that matters for a CUTOVER: throwing the
/// switch means a restart, and unflushed samples from the outgoing backend
/// would be lost at precisely the moment they became the thing worth comparing
/// against.
pub async fn run_metrics_flusher(state: AppState, mut shutdown: watch::Receiver<()>) {
    let period = Duration::from_secs(METRICS_FLUSH_SECS);
    info!(?period, "repo-timing flusher started");
    let mut ticker = tokio::time::interval(period);
    ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);

    loop {
        tokio::select! {
            _ = ticker.tick() => flush_metrics_once(&state).await,
            _ = shutdown.changed() => {
                flush_metrics_once(&state).await;
                info!("repo-timing flusher stopped (final flush done)");
                return;
            }
        }
    }
}

/// One flush. A metrics write must never be able to take anything else down, so
/// a failure is logged and the loop continues.
async fn flush_metrics_once(state: &AppState) {
    if let Err(err) =
        feather_reader::metrics::flush(&state.metrics, &state.db, Utc::now().timestamp()).await
    {
        warn!(%err, "could not persist repo timings");
    }
}

// ---------------------------------------------------------------------------
// Read-state flusher
// ---------------------------------------------------------------------------

/// The read-state flusher loop. On a debounced interval it flushes every DID's
/// dirty read cursors to the PDS in batches; on shutdown it flushes once more so
/// no read-state is stranded. Returns when `shutdown` resolves.
pub async fn run_flusher(state: AppState, mut shutdown: watch::Receiver<()>) {
    let debounce = env_duration_secs("FEATHERREADER_FLUSH_DEBOUNCE_SECS", DEFAULT_FLUSH_DEBOUNCE);
    info!(?debounce, "read-state flusher started");

    // Which DIDs we have already reported as parked, so the log line is once
    // per DID per process rather than once per minute forever.
    let mut parked: HashSet<String> = HashSet::new();

    let mut ticker = interval(debounce);
    ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
    // The first `tick()` completes immediately; swallow it so the debounce window
    // is respected before the first flush.
    ticker.tick().await;

    loop {
        tokio::select! {
            _ = shutdown_fired(&mut shutdown) => {
                info!("read-state flusher: shutdown signal received, final flush");
                // Final drain so Ctrl-C never strands unsynced read-state.
                if let Err(err) = flush_all_dirty(&state, &mut parked).await {
                    error!(%err, "read-state flusher: final flush failed");
                }
                break;
            }
            _ = ticker.tick() => {
                if let Err(err) = flush_all_dirty(&state, &mut parked).await {
                    error!(%err, "read-state flusher: flush round failed");
                }
            }
        }
    }
}

/// Flush every DID that has dirty cursors. Coalesces each DID's dirty cursors
/// into one `applyWrites` batch (chunked by the client), then clears the `dirty` flag on the ones
/// that flushed successfully.
async fn flush_all_dirty(state: &AppState, parked: &mut HashSet<String>) -> anyhow::Result<()> {
    let dids = dids_with_dirty_cursors(&state.db).await?;
    if dids.is_empty() {
        debug!("read-state flusher: nothing dirty");
        return Ok(());
    }
    debug!(
        dids = dids.len(),
        "read-state flusher: flushing dirty cursors"
    );

    for did in dids {
        // **A DID with no session is PARKED, not failed (#117).**
        //
        // Attempting the flush anyway is what produced the production incident:
        // `Repo::session` fails before any network call, the enclosing
        // `flush_read_states` records an error, and the whole thing repeats
        // every 60s forever. 20 failures in the first 20 minutes, and nothing
        // about it could ever have succeeded — the user is signed out.
        //
        // The cursors stay DIRTY on purpose. The reads are not discarded; they
        // wait, and flush on the user's next sign-in. Clearing the flag here
        // would turn a stalled sync into silent data loss, which is strictly
        // worse than the bug being fixed.
        match state.repo().has_session(&did).await {
            Ok(false) => {
                // Once per DID per process: enough to diagnose, not enough to
                // drown the log or mask a real failure during the soak.
                if parked.insert(did.clone()) {
                    info!(
                        %did,
                        "read-state flusher: no OAuth session; parking this DID's \
                         read-state until it signs in again"
                    );
                }
                continue;
            }
            Ok(true) => {
                // It had a session and may have just got one back — stop
                // suppressing its log line, so a LATER park is reported.
                parked.remove(&did);
            }
            Err(err) => {
                // The precondition check itself failed (a DB problem, not an
                // absent session). Fall through and let the flush attempt
                // produce the real error rather than silently skipping.
                warn!(%did, %err, "read-state flusher: session check failed; attempting anyway");
            }
        }

        if let Err(err) = flush_did(state, &did).await {
            // One DID's PDS hiccup must not block the others — its cursors stay
            // dirty and retry next round. This arm is now genuinely transient
            // failures only; the permanent case is parked above.
            warn!(%did, %err, "read-state flusher: DID flush failed; will retry");
        }
    }
    Ok(())
}

/// Every DID that currently has at least one dirty read cursor.
///
/// The store exposes `dirty_cursors(did)` (per-DID, the flusher's hot query) but
/// not the DID enumeration the *global* flusher needs, so this runs the small
/// `SELECT DISTINCT did` directly against the pool. Kept in this module so the
/// scheduler owns its own query and touches no other file.
async fn dids_with_dirty_cursors(pool: &Pool) -> anyhow::Result<Vec<String>> {
    let rows = sqlx::query("SELECT DISTINCT did FROM read_cursor WHERE dirty = 1")
        .fetch_all(pool)
        .await?;
    Ok(rows
        .into_iter()
        .map(|r| r.get::<String, _>("did"))
        .collect())
}

/// Current time as an RFC3339 string (UTC, second precision) — the shape the
/// store's timestamp columns use.
fn now_rfc3339() -> String {
    Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true)
}

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

    /// A feed row that is due right now (`next_poll` NULL), plus its store.
    async fn due_feed(url: &str) -> (Pool, Feed) {
        let pool = store::init_url("sqlite::memory:").await.unwrap();
        store::upsert_feed(
            &pool,
            &store::NewFeed {
                url: url.to_string(),
                ..Default::default()
            },
        )
        .await
        .unwrap();
        let feed = store::get_feed_by_url(&pool, url).await.unwrap().unwrap();
        assert!(
            feed.next_poll.is_none(),
            "fixture must start due (NULL next_poll sorts FIRST in due_feeds)"
        );
        (pool, feed)
    }

    /// **`next_poll` must already be in the future when the fetch is invoked.**
    ///
    /// This is the ordering that converts a process-killing feed from a permanent
    /// crash loop into a single restart. A test cannot kill the process, so it
    /// asserts the observable form of the same fact: at the moment the fetch
    /// begins, the durable state a restart would read has already moved on.
    #[tokio::test]
    async fn next_poll_moves_before_the_fetch_is_invoked() {
        let url = "https://killer.example/feed.xml";
        let (pool, feed) = due_feed(url).await;

        let seen_at_fetch = std::sync::Arc::new(std::sync::Mutex::new(None::<Option<String>>));
        let probe = std::sync::Arc::clone(&seen_at_fetch);

        poll_and_reschedule_with(&pool, &feed, Duration::from_secs(3600), |pool, feed| {
            let probe = std::sync::Arc::clone(&probe);
            let url = feed.url.clone();
            async move {
                // What a restart happening RIGHT NOW would find.
                let row = store::get_feed_by_url(pool, &url).await.unwrap().unwrap();
                *probe.lock().unwrap() = Some(row.next_poll);
                // Then take the process down, as far as this test can simulate it:
                // never produce an outcome the caller could reschedule from.
                Err(anyhow::anyhow!("the fetch killed the process"))
            }
        })
        .await;

        let at_fetch = seen_at_fetch.lock().unwrap().clone().expect("fetch ran");
        let at_fetch = at_fetch.expect(
            "next_poll was still NULL when the fetch began: a crash here re-selects \
             this feed FIRST on every restart, forever",
        );
        assert!(
            at_fetch > now_rfc3339(),
            "next_poll was leased to {at_fetch}, which is not in the future"
        );
    }

    /// The lease is optimistic, not final: a poll that returns overwrites it.
    /// A failure must land on its backoff, not sit at the full cadence.
    #[tokio::test]
    async fn a_returning_poll_overwrites_the_lease() {
        let url = "https://slow.example/feed.xml";
        let (pool, feed) = due_feed(url).await;

        // A cadence far in the future, so "the lease survived" is unmistakable.
        poll_and_reschedule_with(&pool, &feed, Duration::from_secs(86_400), |_, _| async {
            Ok(PollOutcome::Failed {
                backoff: Duration::from_secs(300),
                kind: feather_reader::feed::FailureKind::Fetch,
                detail: "test".to_string(),
            })
        })
        .await;

        let after = store::get_feed_by_url(&pool, url)
            .await
            .unwrap()
            .unwrap()
            .next_poll
            .expect("next_poll must be set after a poll");
        // **A window around the first-failure backoff, not "less than an
        // hour".** The old assertion admitted any value in the whole hour —
        // and the `backoff: 300` this test passes is discarded by
        // `settle_poll`, which computes from the error count, so the test's
        // own input never reached the calculation it named. The value here
        // is `feed::backoff_for(1)` = 5 min; the bin cannot see that private
        // function, so the window is written out.
        let after_dt = chrono::DateTime::parse_from_rfc3339(&after)
            .expect("next_poll is RFC3339")
            .with_timezone(&chrono::Utc);
        let delay = (after_dt - chrono::Utc::now()).num_seconds();
        assert!(
            (240..=360).contains(&delay),
            "a first failure must land on its 5-minute backoff, not the lease: next_poll={after} ({delay}s out)"
        );
    }

    /// The startup-delay override can only SHORTEN the wait. A deployment that
    /// sets an enormous value must not push a production loop further out than
    /// the constant intends.
    ///
    /// Tested through `startup_delay_from` rather than the environment: a
    /// `set_var` here would be the only one in `src/`, racing ~39 `var` reads on
    /// other test threads.
    #[test]
    fn the_startup_delay_override_is_a_ceiling() {
        let d = POLLER_STARTUP_DELAY;
        let at = |v: &str| startup_delay_from(d, Some(v.to_string()));

        assert_eq!(at("0"), Duration::ZERO);
        assert_eq!(at("5"), Duration::from_secs(5));
        assert_eq!(at(" 5 "), Duration::from_secs(5), "surrounding space");
        assert_eq!(at("99999"), d, "the override lengthened the wait");
        assert_eq!(at("not-a-number"), d);
        assert_eq!(at(""), d);
        assert_eq!(startup_delay_from(d, None), d);
    }

    /// **The loops must not land on the same instant at boot.**
    ///
    /// The original version compared the CONSTANTS with no call site involved,
    /// so pointing all five loops at `POLLER_STARTUP_DELAY` passed. The second
    /// version asserted on a string-keyed table and made things worse — see the
    /// `Loop` doc. This reads the exhaustive mapping the code actually uses.
    ///
    /// Asserts against `startup_offset()` directly, NOT `offset_for()`: the
    /// latter applies the `FEATHERREADER_STARTUP_DELAY_SECS` ceiling, which is a
    /// documented dev setting, and asserting through it made the suite fail
    /// under `FEATHERREADER_STARTUP_DELAY_SECS=0`. Keeping env out of this
    /// module's tests is why `startup_delay_from` was split out in the first
    /// place.
    #[test]
    fn the_startup_delays_are_distinct() {
        let offsets: Vec<Duration> = Loop::ALL.iter().map(|l| l.startup_offset()).collect();
        let unique: std::collections::HashSet<Duration> = offsets.iter().copied().collect();
        assert_eq!(
            unique.len(),
            Loop::ALL.len(),
            "two loops share a startup delay: {offsets:?}",
        );
        assert!(
            offsets.iter().all(|d| *d > Duration::ZERO),
            "a loop still fires immediately at boot: {offsets:?}",
        );
    }

    /// **`spawn` really visits every loop — the iteration, not a const array.**
    ///
    /// The previous test asserted `Loop::ALL.len() == 5`, which says nothing
    /// about whether `spawn` reads it. A `.filter()` dropping a loop was green.
    /// This drives the same function `spawn` drives.
    #[test]
    fn every_loop_is_visited_exactly_once() {
        let mut seen = Vec::new();
        for_each_loop(|l| seen.push(l));
        assert_eq!(
            seen,
            Loop::ALL.to_vec(),
            "the iteration `spawn` uses does not visit every loop exactly once, \
             in registry order",
        );
    }

    /// The startup-override key is spelled the same in the code and here.
    ///
    /// Independently written on purpose: mistyping it in `offset_for` silently
    /// disabled the override for every loop with a green suite and green clippy.
    #[test]
    fn the_startup_override_env_key_is_the_documented_one() {
        assert_eq!(STARTUP_DELAY_ENV, "FEATHERREADER_STARTUP_DELAY_SECS");
    }

    /// The registry holds each loop exactly once.
    ///
    /// A pure statement about `Loop::ALL`. It says nothing about `spawn` — see
    /// `spawn_starts_every_registered_loop` for that, and read the note there
    /// before trusting a test in this file that has "spawn" in its name.
    #[test]
    fn the_registry_lists_each_loop_exactly_once() {
        let unique: std::collections::HashSet<Loop> = Loop::ALL.iter().copied().collect();
        assert_eq!(
            unique.len(),
            Loop::ALL.len(),
            "a loop is listed twice in the registry and would be started twice",
        );
    }

    /// **`spawn` starts one task per registered loop — by calling `spawn`.**
    ///
    /// This test previously carried this name while asserting only
    /// `Loop::ALL.len() == 5` plus uniqueness. It never called `spawn`, so the
    /// callback `spawn` passes to `for_each_loop` was an untested decision
    /// point, and dropping a loop there was silent:
    ///
    /// ```ignore
    /// for_each_loop(|l| {
    ///     if l != Loop::PendingSweep {           // 697 lib + 13 bin green,
    ///         handles.push(l.spawn_with(..));    // clippy -D warnings clean
    ///     }
    /// });
    /// ```
    ///
    /// The pending-login sweeper never starts and nonce rows grow without
    /// bound. That is the FOURTH form of this file's recurring defect — after a
    /// wrong constant, a mistyped string key, and a wrong positional argument —
    /// and the round that introduced `for_each_loop` to close the third form
    /// also introduced the mis-named test that hid this one.
    ///
    /// The `+ 2` is the two loops with no startup offset: the metrics flusher
    /// (which fires immediately at boot, deliberately) and the read-state
    /// flusher. Neither is in `Loop`.
    #[tokio::test]
    async fn spawn_starts_every_registered_loop() {
        // `spawn` returns early with an empty Vec when the kill switch is set,
        // which would make the assertion below vacuously wrong rather than
        // failing for a real reason. Say so instead of measuring nothing.
        assert!(
            schedulers_enabled(),
            "FEATHERREADER_DISABLE_SCHEDULER is set in this test process, so \
             `spawn` returns no handles and this test cannot measure anything",
        );
        let state = rust_state().await;
        let (tx, rx) = watch::channel(());
        let handles = spawn(state, rx);
        assert_eq!(
            handles.len(),
            Loop::ALL.len() + 2,
            "`spawn` started {} tasks for {} registered loops + 2 unoffset ones \
             — it is not starting one task per registry entry",
            handles.len(),
            Loop::ALL.len(),
        );
        // Shut them down rather than leaking tasks into the rest of the suite.
        drop(tx);
        for h in handles {
            let _ = h.await;
        }
    }

    /// The ceiling still applies on the way to a loop — the one thing
    /// `offset_for` adds over the raw table.
    ///
    /// Pinned because a mutation deleting `startup_delay(..)` from `offset_for`
    /// — disabling `FEATHERREADER_STARTUP_DELAY_SECS` for every loop at once —
    /// passed the whole suite. Uses `startup_delay_from` so no environment
    /// variable is read.
    #[test]
    fn the_startup_ceiling_applies_to_every_loop() {
        for l in Loop::ALL {
            assert_eq!(
                offset_from(l, Some("0".into())),
                Duration::ZERO,
                "{l:?} ignored the startup-delay ceiling",
            );
        }
    }

    /// **Each variant gets ITS OWN offset — spelled out, not derived.**
    ///
    /// This replaces `assert_eq!(offset_from(l, None), l.startup_offset())`,
    /// which was a tautology: `offset_from` is
    /// `startup_delay_from(which.startup_offset(), raw)` and
    /// `startup_delay_from(d, None)` is `d`, so both sides reduced to the same
    /// expression and the assertion could not fail for ANY mapping. Swapping
    /// two variants' arms in `startup_offset` passed the whole suite.
    ///
    /// The table below is written independently of the `match`, so the two have
    /// to agree — the same reason `the_startup_override_env_key_is_the_documented_one`
    /// spells the env key out by hand. Spelling the seconds here rather than
    /// naming the constants is the point: naming them would reintroduce the
    /// tautology one level up.
    /// The startup offset each loop is documented to run on, spelled out
    /// independently of `Loop::startup_offset`'s match arms so the two have to
    /// agree. Naming the constants here instead would reintroduce the tautology.
    const DOCUMENTED_OFFSETS: [(Loop, u64); 6] = [
        (Loop::Poller, 30),
        (Loop::PublicationPoller, 75),
        (Loop::PendingSweep, 45),
        (Loop::CodeSweep, 60),
        (Loop::Retention, 90),
        (Loop::Adoption, 300),
    ];

    #[test]
    fn every_loop_is_on_its_documented_offset() {
        // Guard by SET, not by length. `documented.len() == Loop::ALL.len()`
        // counts rows, so duplicating one row silently drops a variant from
        // coverage — verified: duplicating the Poller row and drifting
        // `PENDING_SWEEP_STARTUP_DELAY` to 47 s passed the whole suite, since
        // 47 is still distinct and non-zero.
        let listed: std::collections::HashSet<Loop> =
            DOCUMENTED_OFFSETS.iter().map(|(l, _)| *l).collect();
        let registered: std::collections::HashSet<Loop> = Loop::ALL.iter().copied().collect();
        assert_eq!(
            listed, registered,
            "this table and the registry do not cover the same loops",
        );
        for (l, secs) in DOCUMENTED_OFFSETS {
            assert_eq!(
                l.startup_offset(),
                Duration::from_secs(secs),
                "{l:?} is not on its documented {secs}s offset",
            );
        }
    }

    /// **The offsets are not just correct in the table — each loop is handed
    /// ITS OWN on the way to being spawned.**
    ///
    /// This restores coverage a previous round deleted as a "tautology". The
    /// deleted assertion was `offset_from(l, None) == l.startup_offset()`, and
    /// calling it tautological was WRONG: it is tautological only with respect
    /// to changing `startup_offset`'s arms, while independently pinning that
    /// `offset_from` routes through `which` at all. With it gone,
    ///
    /// ```ignore
    /// fn offset_from(which: Loop, raw: Option<String>) -> Duration {
    ///     let _ = which;
    ///     startup_delay_from(Loop::Poller.startup_offset(), raw)
    /// }
    /// ```
    ///
    /// passed 697 lib + 15 bin and clippy, booting every loop at 30 s.
    ///
    /// Asserting against the independently spelled seconds rather than against
    /// `l.startup_offset()` is what keeps this non-tautological — the form that
    /// invited the deletion in the first place.
    #[test]
    fn the_startup_plan_hands_each_loop_its_own_offset() {
        let expected: Vec<(Loop, Duration)> = DOCUMENTED_OFFSETS
            .iter()
            .map(|(l, secs)| (*l, Duration::from_secs(*secs)))
            .collect();
        assert_eq!(
            startup_plan_from(None),
            expected,
            "the plan `spawn` starts from does not pair every loop with its own \
             documented offset, in registry order",
        );
    }

    #[test]
    fn cadence_from_hint_maps_known_values() {
        let d = Duration::from_secs(3600);
        assert_eq!(cadence_from_hint("hourly", d), Duration::from_secs(3600));
        assert_eq!(cadence_from_hint("daily", d), Duration::from_secs(86_400));
        assert_eq!(cadence_from_hint("weekly", d), Duration::from_secs(604_800));
        assert_eq!(cadence_from_hint("realtime", d), Duration::from_secs(300));
        assert_eq!(cadence_from_hint("bogus", d), d);
    }

    #[test]
    fn jitter_stays_within_ten_percent_and_is_seed_stable() {
        let period = Duration::from_secs(86_400);
        let seeds = [
            "https://feather-reader.com",
            "http://localhost:8080",
            "https://reader.example.org",
        ];
        for seed in seeds {
            let j = jittered(period, seed);
            assert!(
                j >= Duration::from_secs(77_760) && j <= Duration::from_secs(95_040),
                "{seed}: {j:?} escaped ±10% of a day"
            );
            // Stable across "restarts": the same seed always yields the same
            // cadence, so a crash loop cannot walk the interval tighter.
            assert_eq!(j, jittered(period, seed));
        }
        // Different instances land on different cadences.
        assert_ne!(jittered(period, seeds[0]), jittered(period, seeds[1]));
    }

    #[test]
    fn jitter_never_returns_a_sub_second_period() {
        assert!(jittered(Duration::from_secs(1), "x") >= Duration::from_secs(1));
        assert!(jittered(Duration::from_millis(1), "x") >= Duration::from_secs(1));
    }

    // ── 0.4.0: publications have their own loop ───────────────────────────

    async fn state_with_due(urls: &[&str]) -> AppState {
        let state = rust_state().await;
        for url in urls {
            store::upsert_feed(
                &state.db,
                &store::NewFeed {
                    url: url.to_string(),
                    ..Default::default()
                },
            )
            .await
            .unwrap();
        }
        state
    }

    const PUB_A: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
    const PUB_B: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lac";
    /// In ANOTHER repo from PUB_A and PUB_B.
    const PUB_C: &str = "at://did:plc:anotherrepoaaaaaaaaaaaaa/site.standard.publication/3lad";

    /// The RSS tick never selects a publication: a slow publication read once
    /// held it, and every RSS feed behind it, for as long as it took.
    #[tokio::test]
    async fn the_rss_tick_never_selects_a_publication() {
        let state = state_with_due(&[PUB_A]).await;
        let client = feed::build_client().unwrap();
        let (_tx, shutdown) = watch::channel(());
        poll_due_once(
            &state,
            &client,
            &Arc::new(Semaphore::new(4)),
            50,
            Duration::ZERO,
            &shutdown,
        )
        .await
        .unwrap();
        let row = store::get_feed_by_url(&state.db, PUB_A)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(row.next_poll, None, "the RSS tick touched a publication");
    }

    /// A pass reads each due publication once and stops: nothing is queued,
    /// so nothing can be queued twice. (The PLC directory is unreachable, so
    /// each read fails fast and is rescheduled with backoff.)
    #[tokio::test]
    async fn a_publication_pass_reads_each_due_row_once_and_stops() {
        let mut state = state_with_due(&[PUB_A, PUB_B]).await;
        let mut config = (*state.config).clone();
        config.oauth.plc_directory = "http://plc.nowhere.invalid".into();
        state.config = Arc::new(config);
        let client = feed::build_client().unwrap();
        let (_tx, shutdown) = watch::channel(());
        tokio::time::timeout(
            Duration::from_secs(30),
            poll_publications_once(&state, &client, &shutdown),
        )
        .await
        .expect("the pass did not stop")
        .unwrap();
        for url in [PUB_A, PUB_B] {
            let row = store::get_feed_by_url(&state.db, url)
                .await
                .unwrap()
                .unwrap();
            assert!(row.next_poll.is_some(), "{url} was not read");
            assert_eq!(row.consecutive_errors, 1, "{url} was read other than once");
        }
    }

    /// Once shutdown is asked for, a pass reads nothing more — and a row it did
    /// not read keeps its `next_poll`, so the next boot reads it at once instead
    /// of an interval later.
    #[tokio::test]
    async fn a_publication_pass_reads_nothing_after_shutdown() {
        let state = state_with_due(&[PUB_A]).await;
        let client = feed::build_client().unwrap();
        let (tx, shutdown) = watch::channel(());
        tx.send(()).unwrap();
        poll_publications_once(&state, &client, &shutdown)
            .await
            .unwrap();
        let row = store::get_feed_by_url(&state.db, PUB_A)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(
            row.next_poll, None,
            "a row the pass never read was deferred"
        );
    }

    async fn with_config(
        mut state: AppState,
        f: impl FnOnce(&mut feather_reader::config::Config),
    ) -> AppState {
        let mut config = (*state.config).clone();
        f(&mut config);
        state.config = Arc::new(config);
        state
    }

    /// Third review of #225: over the DB-size watermark the RSS poller pauses,
    /// and the publication loop — the biggest writer — kept reading.
    #[tokio::test]
    async fn a_publication_pass_respects_the_watermark() {
        let state = with_config(state_with_due(&[PUB_A]).await, |c| {
            c.db_size_watermark_bytes = 1;
            c.oauth.plc_directory = "http://plc.nowhere.invalid".into();
        })
        .await;
        let client = feed::build_client().unwrap();
        let (_tx, shutdown) = watch::channel(());
        poll_publications_once(&state, &client, &shutdown)
            .await
            .unwrap();
        let row = store::get_feed_by_url(&state.db, PUB_A)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(
            row.next_poll, None,
            "a publication was read over the watermark"
        );
    }

    /// Fourth review of #225: a row still due after its read made the pass
    /// return, and it sorts first, so every pass read only it — every other
    /// publication starved behind it for as long as its writes failed.
    #[tokio::test]
    async fn a_stuck_row_does_not_starve_the_rest() {
        let state = state_with_due(&[PUB_A, PUB_C]).await;
        let (_tx, shutdown) = watch::channel(());
        let mut seen = Vec::new();
        // A read that leaves the rows due, as one whose writes fail does.
        poll_publications_with(&state, &shutdown, |feeds, _| {
            seen.extend(feeds.iter().map(|f| f.url.clone()));
            async {}
        })
        .await
        .unwrap();
        seen.sort();
        let mut want = vec![PUB_A.to_string(), PUB_C.to_string()];
        want.sort();
        assert_eq!(seen, want, "every due row, once each");
    }

    /// A healthy read is rescheduled one configured interval out. Every other
    /// loop test fails its read, so only the backoff path ran — a pass handing
    /// `Duration::ZERO` would re-read every healthy publication every tick.
    #[tokio::test]
    async fn a_pass_reschedules_with_the_configured_interval() {
        let state = with_config(state_with_due(&[PUB_A]).await, |c| {
            c.poll_interval = Duration::from_secs(1234);
        })
        .await;
        let (_tx, shutdown) = watch::channel(());
        let mut intervals = Vec::new();
        poll_publications_with(&state, &shutdown, |feeds, interval| {
            intervals.push(interval);
            let db = state.db.clone();
            async move {
                for f in feeds {
                    store::set_next_poll(&db, &f.url, interval).await.unwrap();
                }
            }
        })
        .await
        .unwrap();
        assert_eq!(intervals, vec![Duration::from_secs(1234)]);
    }

    /// The watermark is checked before EVERY read, not once per pass: a pass
    /// can be long, and each read can store a lot.
    #[tokio::test]
    async fn the_watermark_is_checked_before_every_read() {
        // Two REPOS, so two reads: rows of one repo are read together.
        let state = state_with_due(&[PUB_A, PUB_C]).await;
        let size = store::db_size_bytes(&state.db).await.unwrap();
        let state = with_config(state, |c| c.db_size_watermark_bytes = size + 64 * 1024).await;
        let (_tx, shutdown) = watch::channel(());
        let mut reads = 0;
        poll_publications_with(&state, &shutdown, |feeds, interval| {
            reads += 1;
            let db = state.db.clone();
            async move {
                // The first read grows the database past the watermark.
                sqlx::query("CREATE TABLE IF NOT EXISTS ballast (b BLOB)")
                    .execute(&db)
                    .await
                    .unwrap();
                sqlx::query("INSERT INTO ballast VALUES (zeroblob(1048576))")
                    .execute(&db)
                    .await
                    .unwrap();
                for f in feeds {
                    store::set_next_poll(&db, &f.url, interval).await.unwrap();
                }
            }
        })
        .await
        .unwrap();
        assert_eq!(reads, 1, "a second repo was read over the watermark");
    }

    /// 0.4.0 step 2b: due publications in one repo are read together — one
    /// walk of its documents — and a publication in another repo separately.
    #[tokio::test]
    async fn publications_in_one_repo_are_read_together() {
        let state = state_with_due(&[PUB_A, PUB_B, PUB_C]).await;
        let (_tx, shutdown) = watch::channel(());
        let mut calls: Vec<Vec<String>> = Vec::new();
        poll_publications_with(&state, &shutdown, |feeds, _| {
            let mut urls: Vec<String> = feeds.iter().map(|f| f.url.clone()).collect();
            urls.sort();
            calls.push(urls);
            async {}
        })
        .await
        .unwrap();
        calls.sort();
        let mut want = vec![
            vec![PUB_A.to_string(), PUB_B.to_string()],
            vec![PUB_C.to_string()],
        ];
        want.sort();
        assert_eq!(calls, want, "one read per repo");
    }

    /// Review of #228: the pass selected 1,000 due rows a round, so with more
    /// than that still due after their reads (failing writes), every row past
    /// the first 1,000 was skipped for the whole pass.
    #[tokio::test]
    async fn more_than_a_thousand_stuck_rows_are_all_read() {
        // 1,005 publications, each in its own repo (a valid did:plc each).
        let alphabet: Vec<char> = "abcdefghijklmnopqrstuvwxyz234567".chars().collect();
        let did = |mut n: usize| {
            let mut s = String::new();
            for _ in 0..24 {
                s.push(alphabet[n % 32]);
                n /= 32;
            }
            format!("did:plc:{s}")
        };
        let urls: Vec<String> = (0..1005)
            .map(|i| format!("at://{}/site.standard.publication/3lab", did(i)))
            .collect();
        let refs: Vec<&str> = urls.iter().map(String::as_str).collect();
        let state = state_with_due(&refs).await;
        let (_tx, shutdown) = watch::channel(());
        let mut read = 0usize;
        // Reads that leave every row due, as failing writes do.
        poll_publications_with(&state, &shutdown, |feeds, _| {
            read += feeds.len();
            async {}
        })
        .await
        .unwrap();
        assert_eq!(read, 1005, "rows past the selection window were skipped");
    }

    /// The RSS tick pauses at the watermark too. It always did; it had no test
    /// until the check moved into `over_watermark`, shared with the
    /// publication loop.
    #[tokio::test]
    async fn the_rss_tick_respects_the_watermark() {
        let rss = "https://rss.example/feed.xml";
        let state = with_config(state_with_due(&[rss]).await, |c| {
            c.db_size_watermark_bytes = 1
        })
        .await;
        let client = feed::build_client().unwrap();
        let (_tx, shutdown) = watch::channel(());
        poll_due_once(
            &state,
            &client,
            &Arc::new(Semaphore::new(4)),
            50,
            Duration::ZERO,
            &shutdown,
        )
        .await
        .unwrap();
        let row = store::get_feed_by_url(&state.db, rss)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(
            row.next_poll, None,
            "an RSS feed was polled over the watermark"
        );
    }

    /// Third review of #225: a row whose next poll cannot be written stays
    /// due, and the pass re-selected it at once — 11,571 reads in 10 s. A pass
    /// reads a row at most once, whatever the store does.
    #[tokio::test]
    async fn a_pass_never_reads_the_same_row_twice() {
        let state = with_config(state_with_due(&[PUB_A]).await, |c| {
            c.oauth.plc_directory = "http://plc.nowhere.invalid".into();
        })
        .await;
        sqlx::query(
            "CREATE TRIGGER no_lease BEFORE UPDATE OF next_poll ON feeds \
             WHEN NEW.kind = 'publication' BEGIN SELECT RAISE(ABORT, 'disk full'); END",
        )
        .execute(&state.db)
        .await
        .unwrap();
        let client = feed::build_client().unwrap();
        let (_tx, shutdown) = watch::channel(());
        tokio::time::timeout(
            Duration::from_secs(10),
            poll_publications_once(&state, &client, &shutdown),
        )
        .await
        .expect("the pass kept re-reading a row it could not reschedule")
        .unwrap();
        let row = store::get_feed_by_url(&state.db, PUB_A)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(
            row.consecutive_errors, 1,
            "read {} times in one pass",
            row.consecutive_errors
        );
    }

    /// Every kind the store calls pollable has a loop that polls it. A kind
    /// added to `FeedKind::POLLABLE` with no poller would be counted as due and
    /// overdue by `/stats`, and read by nothing.
    #[test]
    fn every_pollable_kind_has_a_poller() {
        for kind in feed::FeedKind::POLLABLE {
            let poller = poller_for(*kind);
            assert!(
                poller.is_some_and(|l| Loop::ALL.contains(&l)),
                "{kind:?} has no running poller"
            );
        }
        assert_eq!(poller_for(feed::FeedKind::Unsupported), None);
    }

    // ── #117: orphaned dirty read-state ──────────────────────────────────────

    /// An `AppState` on the rust backend with an empty in-memory store.
    ///
    /// `key_path` is per-test: the default is the RELATIVE
    /// `oauth-signing-key.json`, so a rust-backend test would write real
    /// encrypted key material into the working directory and later runs would
    /// fail to decrypt it under a fresh key.
    async fn rust_state() -> AppState {
        let db = store::init_url("sqlite::memory:").await.unwrap();
        AppState::new(
            feather_reader::config::Config {
                repo_backend: feather_reader::metrics::Backend::Rust,
                oauth: feather_reader::config::OauthConfig {
                    key_path: std::env::temp_dir().join(format!(
                        "fr-sched-oauth-key-{}-{:p}.json",
                        std::process::id(),
                        &db as *const _
                    )),
                    encryption_key: Some("a".repeat(43)),
                    ..feather_reader::config::OauthConfig::default()
                },
                ..feather_reader::config::Config::default()
            },
            db,
        )
        .unwrap()
    }

    /// A dirty cursor for `did`, as a mark-read would leave it.
    async fn dirty_cursor_for(state: &AppState, did: &str) {
        store::upsert_cursor(
            &state.db,
            &ReadCursor {
                did: did.to_string(),
                feed_url: "https://example.com/feed.xml".into(),
                read_through: None,
                read_ids: "[\"1\"]".into(),
                unread_ids: "[]".into(),
                dirty: true,
                pds_created: false,
                updated_at: now_rfc3339(),
            },
        )
        .await
        .unwrap();
    }

    /// **#117 — a DID with no OAuth session is PARKED, not retried.**
    ///
    /// The inverse of the characterization test this replaces. Before the fix,
    /// five rounds produced five recorded failures and a `warn!` each; nothing
    /// about them could ever have succeeded, because the user is signed out.
    ///
    /// Asserts the two properties the fix has to hold together:
    ///   1. no error is recorded, however many rounds run — the soak is not
    ///      polluted by a condition that is not a failure; and
    ///   2. the cursor stays DIRTY — the reads are parked, not discarded.
    ///
    /// (2) is the one worth guarding. The cheapest way to silence the loop is to
    /// clear the flag, and that would turn a stalled sync into silent data loss.
    #[tokio::test]
    async fn a_did_with_no_session_is_parked_not_retried() {
        let state = rust_state().await;
        let did = "did:plc:orphanedreadstate00000000";
        dirty_cursor_for(&state, did).await;

        let mut parked = HashSet::new();
        const ROUNDS: usize = 5;
        for round in 1..=ROUNDS {
            flush_all_dirty(&state, &mut parked)
                .await
                .expect("a parked DID must not abort the sweep");
            assert_eq!(
                store::dirty_cursors(&state.db, did).await.unwrap().len(),
                1,
                "round {round}: the parked cursor was cleared — the reads are now lost",
            );
        }

        let err = state
            .metrics
            .snapshot()
            .into_iter()
            .find(|r| r.op == "flush_read_states")
            .map(|r| r.stats.err_count)
            .unwrap_or(0);
        assert_eq!(err, 0, "a parked DID was counted as {err} flush failures");
        assert_eq!(parked.len(), 1, "the DID should be recorded as parked once");
    }

    /// **The reads survive the gap: parking holds them until the user returns.**
    ///
    /// This is the test that makes parking defensible rather than merely quiet.
    /// A cursor parked while signed out must still be there — and still flush —
    /// once a session exists again.
    ///
    /// The flush itself fails here (the fixture PDS is unreachable), which is
    /// the point: what is asserted is that the DID is no longer SKIPPED, so the
    /// attempt is made at all. A fix that parked permanently would pass the test
    /// above and fail this one.
    #[tokio::test]
    async fn a_parked_cursor_is_retried_once_the_user_signs_in_again() {
        let state = rust_state().await;
        let did = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
        dirty_cursor_for(&state, did).await;
        let mut parked = HashSet::new();

        flush_all_dirty(&state, &mut parked).await.unwrap();
        assert!(
            parked.contains(did),
            "precondition: parked while signed out"
        );
        assert_eq!(
            state
                .metrics
                .snapshot()
                .into_iter()
                .find(|r| r.op == "flush_read_states")
                .map(|r| r.stats.ok_count + r.stats.err_count)
                .unwrap_or(0),
            0,
            "precondition: no flush was attempted while parked",
        );

        // The user signs back in.
        let runtime = state.oauth.as_deref().expect("oauth runtime");
        feather_reader::oauth::store::put_session(
            &state.db,
            &runtime.codec,
            &feather_reader::oauth::store::OAuthSession {
                sub: did.into(),
                issuer: "https://auth.invalid".into(),
                aud: "https://pds.invalid".into(),
                dpop_key_jwk: feather_reader::oauth::keys::SigningKey::generate("session-dpop")
                    .to_jwk_json()
                    .unwrap(),
                access_token: "at".into(),
                refresh_token: "rt".into(),
                token_type: "DPoP".into(),
                granted_scope: "atproto".into(),
                expires_at: Some(Utc::now().timestamp() + 3600),
            },
        )
        .await
        .unwrap();

        flush_all_dirty(&state, &mut parked).await.unwrap();

        assert!(
            !parked.contains(did),
            "the DID is still marked parked after regaining a session",
        );
        let attempts = state
            .metrics
            .snapshot()
            .into_iter()
            .find(|r| r.op == "flush_read_states")
            .map(|r| r.stats.ok_count + r.stats.err_count)
            .unwrap_or(0);
        assert_eq!(
            attempts, 1,
            "the parked read-state was never re-attempted after sign-in",
        );
        assert_eq!(
            store::dirty_cursors(&state.db, did).await.unwrap().len(),
            1,
            "the unflushed cursor must remain dirty after a failed attempt",
        );
    }
}