lgwks_bot 2.2.0

Capability-gated automation bots on a change-detecting ECS schedule: Observe, Evaluate, Execute, and Query, with an async runtime facade.
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
//! The durable, file-backed per-step record store behind a resumable run.
//!
//! [`HostBuilder::run_store`] installs one
//! of these. A step marked [`Scope::remember`]
//! consults it before running its future and appends to it after; that is the
//! whole of the durability claim, and it covers exactly the steps an author asked
//! to be durable.
//!
//! # Why this is not `journal::FileJournal`
//!
//! The *frame grammar* is the estate's and is shared verbatim: [`journal::frame`]
//! holds the length prefix, the 32-byte head, the torn-tail scan and the refusal
//! of a frame no writer produces, and both this store and `journal::file` call
//! those same functions. What cannot be reused is the *record*: `journal::file`
//! frames [`EffectEvent`] and its head chains effect
//! positions, so a step record smuggled through it would claim to be an effect
//! event and chain against a sequence that has nothing to do with steps. So this
//! module states its own record and reuses the frame grammar; the encoding is
//! `lgwks_std::wire` in both. It shares `journal::file`'s storage-owner thread
//! too, so a durable step's `sync_all` costs the device, not the executor.
//!
//! # What a record is
//!
//! One record is the run, the step's key digest, the step's path, the owning
//! tenant, and the archived bytes of the value the step returned. The path is
//! stored beside the digest so an undecodable value is attributable; the digest
//! is what the lookup is on, so two steps cannot read each other's record even
//! if an author gives them the same name.
//!
//! # Bounds
//!
//! Three ceilings, all declared as constants and all reported as typed refusals
//! rather than enforced by truncation: per-record bytes, records per run, and
//! total file bytes. A refusal never deletes, compacts or rewrites a committed
//! record — INV-BOT-14's rule, applied to this store.

use std::collections::HashMap;
use std::fmt;
use std::fs::{File, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};

use lgwks_std::hash::{Digest, Hasher};

use crate::journal::frame::SaturatingFrom;
use lgwks_std::wire::{WireError, from_bytes, to_bytes};

use crate::effect::RunId;
use crate::journal::frame::{self, Cursor, HEAD_BYTES};
use crate::journal::owner::{self, Stage, StorageGate, StorageOwner, SubmitError};
use crate::script::run_store::{Appended, RunRecords, StagedRecord, StoredValue};
use crate::script::{FlowError, StepKey};

use super::definition::{DefinitionIdentity, Drift, STORE_FORMAT};

/// The most bytes one archived step record may occupy.
///
/// A step that returns more than this is refused rather than written: the value
/// belongs in the effect journal or in the caller's own storage, not in a record
/// whose whole purpose is to be replayed into memory on resume.
pub const MAX_RECORD_BYTES: usize = 256 * 1024;

/// The most records one run's store retains.
pub const MAX_RECORDS_PER_RUN: u64 = 65_536;

/// The most bytes one run's store file may hold.
pub const MAX_STORE_BYTES: u64 = 64 * 1024 * 1024;

/// The constant prefix of every run store's file.
///
/// A file that does not begin with it is not this format and is refused rather
/// than scanned: reading a foreign file as records is how a caller is handed a
/// value nobody wrote. The last byte of the header is the format's version, so a
/// future change is a refusal rather than a misreading. `\x01` was the record
/// without a [`DefinitionIdentity`]; a file in that format still *is* a run
/// store, so it is refused as [`StoreError::FormatVersion`] naming both versions
/// rather than as `NotAStore`, which would claim the bytes were never this
/// store's. See [`check_format_version`] for the rule and the `definition` module
/// for why reading one as the unversioned identity would be the worse answer.
const STORE_MAGIC: &[u8; 16] = b"lgwks-runstore\x00\x00";

/// The header as this crate writes it: the constant prefix and the version.
///
/// Built once so the writer and the reader cannot disagree about which byte is
/// which, which is the only thing a hand-written `OpenOptions` chain would get
/// wrong.
const STORE_HEADER: [u8; 16] = {
    let mut header = *STORE_MAGIC;
    header[15] = STORE_FORMAT;
    header
};

/// The index of the version byte in [`STORE_HEADER`].
///
/// Spelled as a name because the writer, the reader and the refusal all have to
/// agree on *which* byte carries the version, and a literal `15` in three places
/// is three chances for one of them to mean a different byte. It is the last
/// byte, which is what leaves the first [`STORE_MAGIC`] free of version bytes.
const VERSION_BYTE: usize = STORE_MAGIC.len() - 1;

/// Refuse a header whose format this build does not read.
///
/// Two refusals, and the order is the argument: a header whose *first* bytes are
/// not this store's magic was never a run store ([`StoreError::NotAStore`]),
/// while a header whose first bytes match and whose version byte does not is a
/// run store this build cannot read ([`StoreError::FormatVersion`], naming both
/// versions). Reporting the second as the first would tell an operator their
/// data was never theirs, when in fact it was written by an earlier version of
/// this very crate — which is the one conclusion that must never be drawn from a
/// refusal, because it is what makes someone delete the file.
///
/// [`STORE_MAGIC`] is constant across versions and the version byte is not, so
/// "the first bytes match" is exactly "this is some version of this format",
/// and the one comparison below is what keeps those two facts from collapsing
/// into one check.
///
/// # Errors
///
/// [`StoreError::NotAStore`] when the magic does not match, and
/// [`StoreError::FormatVersion`] naming the found and expected versions when it
/// does.
fn check_format_version(header: [u8; STORE_HEADER.len()]) -> Result<(), StoreError> {
    if header[..VERSION_BYTE] != STORE_MAGIC[..VERSION_BYTE] {
        let refusal = Err(StoreError::NotAStore);
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_format_version: returning an error to the caller");
        return refusal;
    }
    let found = header[VERSION_BYTE];
    if found == STORE_FORMAT {
        return Ok(());
    }
    Err(StoreError::FormatVersion {
        found,
        expected: STORE_FORMAT,
    })
}

/// The chain digest a store starts from: 32 zero bytes, hashed through the same
/// framing as any record so the first head is a real chain step and not a
/// constant every file shares.
fn genesis_head() -> Digest {
    let mut hasher = Hasher::new();
    hasher.write_framed(b"lgwks-runstore/genesis");
    hasher.finalize()
}

/// Which ceiling an append or a replay hit.
///
/// Named rather than collapsed into one "capacity exceeded" because the repair
/// differs per arm: one is a value too big, one is a run that grew, one is a
/// store that filled. A caller that cannot tell them apart raises every cap.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum StoreLimitKind {
    /// One archived record exceeded [`MAX_RECORD_BYTES`].
    RecordBytes,
    /// The run already holds [`MAX_RECORDS_PER_RUN`] records.
    Records,
    /// The store file would exceed [`MAX_STORE_BYTES`].
    StoreBytes,
}

/// Why a record could not be read, written, or resumed.
///
/// A read failure is never absence (INV-BOT-7): every arm here is an error the
/// caller must see, never an empty index that would cause a finished step to
/// re-run as though it had never run.
#[derive(Debug)]
#[non_exhaustive]
pub enum StoreError {
    /// A declared ceiling refused the write, or the replay. The store is
    /// unchanged: a refusal is not a compaction.
    Limit {
        /// Which ceiling.
        kind: StoreLimitKind,
        /// What would have been needed.
        requested: u64,
        /// The ceiling.
        limit: u64,
    },
    /// The store file could not be opened, read or written.
    Storage {
        /// What the device said.
        cause: std::io::Error,
    },
    /// The file exists but is not a run store, or its header is short.
    NotAStore,
    /// The file *is* a run store, in a format this build does not read.
    ///
    /// Distinct from [`NotAStore`][Self::NotAStore] because the two say opposite
    /// things about what is on the disk: one says these bytes were never this
    /// store's, the other says they were, and were written by a version of this
    /// crate whose records mean something this build cannot reconstruct. A
    /// pre-version store holds records with no
    /// [`DefinitionIdentity`] in them, and the only
    /// ways to read those are to invent an identity they never carried — which
    /// makes every pre-version resume look exactly compatible — or to discard
    /// evidence a running system is relying on. So the store is refused with both
    /// versions named, and the caller is told which this build wants rather than
    /// being handed a corrupt-store error that says nothing about either.
    ///
    /// This is a breaking change for every existing run store, and it is
    /// deliberately not softened into one: the format has never shipped a version
    /// that could lose a record, so there is nothing to convert, and a deployment
    /// that needs its records keeps its own copy and re-runs.
    FormatVersion {
        /// The version byte the file declares.
        found: u8,
        /// The version this build reads.
        expected: u8,
    },
    /// A committed frame does not follow from the records before it. Those
    /// bytes are acknowledged evidence and are refused, never trimmed.
    Corrupt {
        /// The frame's index, from zero.
        at: u64,
    },
    /// The value could not be archived.
    Encoding {
        /// What the codec said.
        cause: WireError,
    },
    /// The run id belongs to a tenant other than the one resuming it.
    ///
    /// A refusal rather than an absence: the run *is* in this store, and reading
    /// its records for the wrong tenant is a cross-tenant read dressed up as a
    /// miss.
    ForeignTenant {
        /// The tenant that minted the run.
        owner: String,
        /// The tenant that asked to resume it.
        asked: String,
    },
    /// The run id names no run this store holds, so it cannot be attributed to
    /// any tenant here.
    ///
    /// Distinct from [`ForeignTenant`][Self::ForeignTenant] because the two say
    /// opposite things: one says the run exists and belongs elsewhere, the other
    /// says this store has never seen it. A resume in either case is refused —
    /// recording a run this host cannot attribute would put one tenant's work
    /// under another's name.
    UnknownRun {
        /// The run id that was asked for.
        run: String,
        /// The store that could not attribute it.
        tenant: String,
    },
    /// The run's records were written under a different definition identity, so
    /// replaying them would answer a question this build never asked.
    ///
    /// Refused before any step runs and with nothing written, because the two
    /// available answers are both worse: replaying is returning a value the
    /// current definition never produced, and re-running is repeating effects a
    /// previous attempt already performed. Which axis disagrees is named, so the
    /// caller knows whether it needs a new run, a decision, a migration or an
    /// edit.
    Incompatible {
        /// The run whose records disagree.
        run: String,
        /// Which axis, and what each side holds.
        drift: Drift,
    },
}

impl StoreError {
    /// A storage refusal carrying the device's own error.
    ///
    /// A constructor rather than seven literals, because every `map_err` in this
    /// module wraps the same fact — the device said no — and a literal per call
    /// site would be seven chances to name a different variant for the same
    /// event.
    ///
    /// `pub(crate)` because the repair ledger beside this store reports device
    /// refusals in the same vocabulary: one error type for both file-backed
    /// stores means a caller handling a device refusal handles it once.
    pub(crate) fn storage(cause: std::io::Error) -> Self {
        Self::Storage { cause }
    }
}

impl fmt::Display for StoreError {
    /// The refusal, naming the ceiling, the device, or the owning tenant.
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        match *self {
            Self::Limit {
                kind,
                requested,
                limit,
            } => {
                let what = match kind {
                    StoreLimitKind::RecordBytes => "one record's bytes",
                    StoreLimitKind::Records => "records in one run",
                    StoreLimitKind::StoreBytes => "bytes in one run's store",
                };
                write!(
                    formatter,
                    "{what}: {requested} exceeds the ceiling of {limit}"
                )
            }
            Self::Storage { ref cause } => {
                write!(formatter, "the run store's device refused: {cause}")
            }
            Self::NotAStore => {
                formatter.write_str("this file is not a run store; refusing to read it as records")
            }
            Self::FormatVersion { found, expected } => write!(
                formatter,
                "this is a run store in format version {found}, and this build reads version \
                 {expected}; refusing to read records whose definition identity this version \
                 cannot reconstruct"
            ),
            Self::Corrupt { at } => write!(
                formatter,
                "run store frame {at} does not follow from the records before it; \
                 the bytes are refused, not trimmed"
            ),
            Self::Encoding { ref cause } => {
                write!(formatter, "the step's value could not be archived: {cause}")
            }
            Self::ForeignTenant {
                ref owner,
                ref asked,
            } => write!(
                formatter,
                "run belongs to tenant {owner:?}, not {asked:?}; refusing to read its records"
            ),
            Self::UnknownRun {
                ref run,
                ref tenant,
            } => write!(
                formatter,
                "no records for run {run} in tenant {tenant:?}'s store; refusing to resume a \
                 run this store cannot attribute to it"
            ),
            Self::Incompatible { ref run, ref drift } => write!(
                formatter,
                "run {run} was recorded under a different definition: {}; refusing to replay \
                 it under this one",
                drift
            ),
        }
    }
}

impl std::error::Error for StoreError {
    /// The device's or the codec's own error.
    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
        match *self {
            Self::Storage { ref cause } => Some(cause),
            Self::Encoding { ref cause } => Some(cause),
            Self::Limit { .. }
            | Self::NotAStore
            | Self::FormatVersion { .. }
            | Self::Corrupt { .. }
            | Self::ForeignTenant { .. }
            | Self::UnknownRun { .. }
            | Self::Incompatible { .. } => None,
        }
    }
}

impl From<StoreError> for FlowError {
    /// A store refusal is permanent: repeating the step would ask the same store
    /// the same question and get the same refusal.
    ///
    /// Two arms and not one, because the refusal's *kind* is what the caller acts
    /// on. A drift is already the flow's own typed [`FlowError::Incompatible`]
    /// carrying its axis, so it stays that rather than being buried one level
    /// down; every other store refusal — an unreadable device, a ceiling, a
    /// foreign tenant — reaches the caller as `FlowError::Store` wrapping the
    /// store's own error, so a read failure is never reported as a definition
    /// drift and never as a rendered string (INV-BOT-7).
    fn from(error: StoreError) -> Self {
        match error {
            StoreError::Incompatible { drift, .. } => Self::incompatible("", drift),
            other => Self::Store {
                at: Arc::from(""),
                source: Box::new(other),
            },
        }
    }
}

// ── The store ───────────────────────────────────────────────────────────────

/// One durable, file-backed set of per-step records.
///
/// Cloning shares the file and the index through an `Arc`, so the host's handles
/// see one store rather than one per clone — the rule the admission budget
/// follows, and for the same reason: a second store over one file would be a
/// second writer.
#[derive(Clone)]
pub struct RunStore {
    /// Where the store lives, and the index every clone shares.
    inner: Arc<StoreInner>,
}

/// The one state every clone of a store handle shares, plus the fault the
/// read-failure injector arms.
struct StoreInner {
    /// The thread that holds the file, and the one door an append reaches the
    /// disk through. Its writes are `write_all` plus `sync_all`, so they belong
    /// off the executor — INV-BOT-50.
    owner: StorageOwner<Arc<Mutex<Index>>, Appended>,
    /// Where the file is, for diagnostics.
    path: PathBuf,
    /// The index every clone reads, and the one the owner folds into: per run, per
    /// step key, the record.
    index: Arc<Mutex<Index>>,
    /// Set once by [`RunStore::fail_next_index_read`], and taken by the first
    /// *record* read that observes it.
    ///
    /// `AtomicBool` rather than a `Mutex<Option<..>>` so arming and firing are
    /// two words wide and a read takes no second lock: the read path must not gain
    /// a lock acquisition purely so a fault can be scheduled, and the injection is
    /// a bounded one-shot — `swap`, not `load` — so two racing reads cannot both
    /// take it and the injector cannot leave the store armed forever for a caller
    /// who never expected a fault.
    ///
    /// Taken by [`Self::step_readable`] alone, and never by [`Self::index`]: the
    /// host's own admission reads (`tenant_of`, `definition_of`) run before the
    /// body and would otherwise spend the fault on themselves, which would test
    /// the host's pre-flight rather than the durable step's replay check. A fault
    /// that fires on the wrong read is a fault that proved the wrong thing.
    unreadable: AtomicBool,
}

/// The index, and the committed length it accounts for.
#[derive(Debug)]
struct Index {
    /// Per run: the owning tenant and its steps.
    runs: HashMap<RunId, RunIndex>,
    /// The file length every indexed record accounts for.
    ///
    /// Held beside the index rather than recomputed from it, because the append
    /// fence compares *bytes* against the device and a sum over the index would
    /// have to reproduce the frame overhead exactly — a second definition of the
    /// file's length that could drift from the first.
    committed: u64,
    /// The chain head the next frame follows.
    ///
    /// Carried beside the index because it cannot be recovered from it: a head is
    /// over the *archived* record, so recomputing one would re-archive a value
    /// and risk hashing something other than the bytes on the disk.
    tail: Digest,
}

/// One run's records.
#[derive(Debug)]
struct RunIndex {
    /// The tenant that minted the run, and the only one allowed to read it.
    tenant: String,
    /// The definition every record under this run was written under.
    ///
    /// One per run rather than one per record because a run is one attempt to
    /// execute one definition: a record written under a different definition
    /// than the run's first is not a step that changed, it is a different run
    /// that arrived under the same id, and the first record is what the rest are
    /// compared against.
    definition: DefinitionIdentity,
    /// Per step key hex: the step's path and its archived value.
    steps: HashMap<String, Held>,
}

/// One step's record as the index holds it.
#[derive(Debug)]
struct Held {
    /// The step path, kept so a decode failure is attributable.
    path: String,
    /// The archived value, kept so a resume does not re-read the file per step.
    bytes: Vec<u8>,
}

impl fmt::Debug for RunStore {
    /// The path and the indexed run count. The values are not rendered: a report
    /// or a log line must not carry a step's output into a reader that asked
    /// where the store is.
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        formatter
            .debug_struct("RunStore")
            .field("path", &self.inner.path)
            .field(
                "runs",
                &self.inner.index.lock().map_or(0, |index| index.runs.len()),
            )
            .field(
                "committed",
                &self.inner.index.lock().map_or(0, |index| index.committed),
            )
            .finish()
    }
}

impl RunStore {
    /// Open (or create) the store at `path`.
    ///
    /// The file is created with a header if it does not exist, and replayed if
    /// it does. A torn final record is dropped, because an interrupted append is
    /// a prefix cut short and was never anyone's answer; every earlier record is
    /// retained.
    ///
    /// # Errors
    ///
    /// [`StoreError::Storage`] when the device refuses, or when a torn tail
    /// cannot be trimmed; [`StoreError::NotAStore`] when `path` exists and is
    /// not a run store; [`StoreError::Corrupt`] when a committed frame does not
    /// follow from the ones before it.
    pub fn open(path: impl Into<PathBuf>) -> Result<Self, StoreError> {
        Self::open_impl(path.into(), false)
    }

    /// Open a store whose device does not answer a flush until its gate is
    /// released.
    ///
    /// A fault injector, and public for the reason
    /// [`FileJournal::open_with_stalled_storage`](crate::journal::FileJournal::open_with_stalled_storage)
    /// is: "what does this run do while its record store has stopped answering" is
    /// a question an operator has to be able to ask on a real process, and a probe
    /// that only exists inside the crate's own test binary cannot answer it. The
    /// bytes written are real and the frames are the ones the store always writes;
    /// only the device's answer is held back, which is exactly what a stalled
    /// device does. Nothing is acknowledged until the release.
    ///
    /// # Errors
    ///
    /// As [`RunStore::open`].
    pub fn open_with_stalled_device(path: impl Into<PathBuf>) -> Result<Self, StoreError> {
        Self::open_impl(path.into(), true)
    }

    /// One constructor for both: a stalled device is a property of the storage
    /// owner, not a second way to open a store.
    ///
    /// # Errors
    ///
    /// As [`RunStore::open`].
    fn open_impl(path: PathBuf, stalled: bool) -> Result<Self, StoreError> {
        let existed = path.exists();
        let mut file = OpenOptions::new()
            .read(true)
            .write(true)
            .create(true)
            .truncate(false)
            .open(&path)
            .map_err(StoreError::storage)?;
        let header = u64::saturating_from(STORE_HEADER.len());
        if !existed || file.metadata().map_err(StoreError::storage)?.len() == 0 {
            file.write_all(&STORE_HEADER)
                .and_then(|()| file.sync_all())
                .map_err(StoreError::storage)?;
        }
        let index = replay(&mut file)?;
        // Trim the incomplete tail rather than appending after it: a write past
        // an interrupted frame leaves a file whose prefix is a prefix of a
        // record and whose middle is another record, which no reader could fold.
        if file.metadata().map_err(StoreError::storage)?.len() != index.committed {
            file.set_len(index.committed).map_err(StoreError::storage)?;
            file.sync_all().map_err(StoreError::storage)?;
        }
        debug_assert!(
            index.committed >= header,
            "the header is part of the committed length"
        );
        // The index is shared, not copied: the owner folds each committed record
        // into it inside the same ordered step as the write, and a lookup reads it
        // from an executor thread. That is not a lock across an `await` — a lookup
        // takes it for a hash-map probe and gives it straight back, and the write
        // it can wait for is the *file's*, which is what the owner thread exists to
        // absorb. Parking on a mutex a few hundred nanoseconds wide is not
        // parking on an `fsync`; it is the second thing the fix is about.
        let index = Arc::new(Mutex::new(index));
        let owner =
            StorageOwner::spawn(file, Arc::clone(&index), stalled).map_err(StoreError::storage)?;
        Ok(Self {
            inner: Arc::new(StoreInner {
                owner,
                path,
                index,
                unreadable: AtomicBool::new(false),
            }),
        })
    }

    /// Make the next *durable step's* record read fail, as an unreadable device would.
    ///
    /// A fault injector, and public for the reason
    /// [`RunStore::open_with_stalled_device`](Self::open_with_stalled_device) is:
    /// "what does a durable step do when its store cannot be read" is a question
    /// INV-BOT-7 answers in a fault nobody can schedule on a real process, and a
    /// probe that exists only inside the crate's own test binary cannot answer it.
    /// The records themselves are untouched — only the *read* of the in-memory
    /// index fails — so the store is exactly one it holds records for, and the
    /// question under test is what the step does with an answer it did not get.
    ///
    /// Aimed at the step rather than at every reader, because the host reads the
    /// same index before the body runs: a fault spent on the host's admission
    /// pre-flight would report `Disposition::Refused` at admission and would never
    /// reach the check whose behaviour this exists to observe.
    ///
    /// One-shot, and bounded by construction: the first step read that observes it
    /// takes it with a `swap`, so a caller that arms it and never reads is not left
    /// with a store that fails forever, and two concurrent reads cannot both take
    /// the same fault.
    pub fn fail_next_index_read(&self) {
        self.inner.unreadable.store(true, Ordering::SeqCst);
    }

    /// A handle that can release this store's parked device, independently of
    /// the store.
    ///
    /// Returned rather than only as [`Self::release_device`] because the caller
    /// that most needs to un-stick the device is the one awaiting a record on it,
    /// and that caller holds the store's borrow for the whole wait.
    #[must_use]
    pub fn storage_gate(&self) -> StorageGate {
        self.inner.owner.gate()
    }

    /// Let the outstanding flush proceed, and every flush after it.
    ///
    /// After this the handle is an ordinary store: the stall is over and is not
    /// re-armed. A caller awaiting a record on this store cannot reach this
    /// method — the borrow the append holds is the same one — and needs the gate
    /// [`Self::storage_gate`] hands back instead.
    pub fn release_device(&self) {
        self.storage_gate().release();
    }

    /// Arm a refusal of the next batch's covering `sync_all`.
    ///
    /// The seam the storage owner's `fail_next_flush` reaches through, so the
    /// all-or-nothing answer of a failed group commit — every member refused, no
    /// fold run, one poison latched, and a reopen that reads the file rather than
    /// the handle — is proved on the shipped store rather than a copy of it. A
    /// test-only door for the same reason `StorageOwner::fail_next_commit` is: no
    /// filesystem refuses a flush on demand.
    #[cfg(test)]
    pub(crate) fn fail_next_flush(&self) {
        self.inner.owner.fail_next_flush();
    }

    /// Open a store under `dir`, named by the tenant that owns it.
    ///
    /// The form [`HostBuilder::run_store`](crate::task::HostBuilder::run_store)
    /// uses. The tenant is part of the file name, so two tenants pointed at one
    /// directory never share a file — and a store installed under two names is
    /// two stores, which is why the cross-tenant refusal below is enforced on the
    /// record rather than on the path.
    ///
    /// # Errors
    ///
    /// As [`RunStore::open`].
    pub fn open_in(dir: &Path, tenant: &str) -> Result<Self, StoreError> {
        std::fs::create_dir_all(dir).map_err(StoreError::storage)?;
        Self::open(dir.join(format!("{tenant}.runstore")))
    }

    /// Where this store lives.
    #[must_use]
    pub fn path(&self) -> &Path {
        &self.inner.path
    }

    /// How many records `run` holds.
    #[must_use]
    pub fn record_count(&self, run: RunId) -> usize {
        self.index()
            .runs
            .get(&run)
            .map_or(0, |held| held.steps.len())
    }

    /// The archived bytes committed for `key` under `run`, or `None` when there
    /// are none.
    ///
    /// The reader's door onto a durable step's value, and the other half of
    /// [`RunStore::record_count`]: a resumed run needs to read back *what* a
    /// previous instance recorded, not merely *that* it recorded something, and
    /// without this the only way to see a record's contents would be to re-run the
    /// step that wrote it.
    ///
    /// `Err` is a read failure — a run another tenant owns — never a miss, so a
    /// caller never reads a refusal as absence (INV-BOT-7).
    ///
    /// # Errors
    ///
    /// [`FlowError`] when `run` belongs to a tenant other than `tenant`.
    pub fn lookup(
        &self,
        tenant: &str,
        run: RunId,
        key: StepKey,
    ) -> Result<Option<Vec<u8>>, FlowError> {
        Ok(RunRecords::lookup(self, tenant, run, key)?.map(|stored| stored.into_bytes()))
    }

    /// The bytes this store has committed, header included.
    #[must_use]
    pub fn committed_bytes(&self) -> u64 {
        self.index().committed
    }

    /// How many `fsync`s this store has paid, and how many records it has staged.
    ///
    /// The evidence that a group commit formed rather than a claim that one could:
    /// the first number counts the `sync_all` calls the storage owner actually
    /// made, the second counts the records it actually staged. A store that flushed
    /// once per record has a ratio of one; a store that batched has a ratio below
    /// one, and the smaller it is the more records each flush carried.
    ///
    /// Both are counted on the owner thread rather than by the callers, so a caller
    /// that gave up on its answer is still counted as a staged record — leaving it
    /// out would report a better batching factor than the store achieved.
    #[must_use]
    pub fn flush_counts(&self) -> (u64, u64) {
        self.inner.owner.flush_counts()
    }

    /// Whether this store has a record under `run`.
    ///
    /// The check `Host::resume` makes before it admits a run: a run id this
    /// store never wrote is one it cannot attribute to a tenant, and resuming it
    /// would be recording another host's run under this host's tenant.
    #[must_use]
    pub fn knows_run(&self, run: RunId) -> bool {
        self.index().runs.contains_key(&run)
    }

    /// Who minted `run`, when the store knows.
    ///
    /// The answer a caller checks before deciding a run id belongs to someone
    /// else, so it names the tenant from the record rather than from the path.
    #[must_use]
    pub fn tenant_of(&self, run: RunId) -> Option<String> {
        let index = self.index();
        index.runs.get(&run).map(|held| held.tenant.clone())
    }

    /// The definition `run`'s records were written under.
    ///
    /// The answer a resume compares its own identity against, through
    /// [`Self::check_definition`]. `None` for a run this store never wrote, which
    /// is not a refusal to read — a process killed before its first record left
    /// nothing to disagree with.
    #[must_use]
    pub fn definition_of(&self, run: RunId) -> Option<DefinitionIdentity> {
        self.index()
            .runs
            .get(&run)
            .map(|held| held.definition.clone())
    }

    /// Refuse `identity` for `run` if the records disagree with it.
    ///
    /// The pre-flight a resume makes before admission: a run whose recorded
    /// definition, input, step count or value schema is not the one being resumed
    /// is refused with [`StoreError::Incompatible`] and nothing is written, so no
    /// step body is polled and no record is added. Called both at the host, before
    /// the body is ever constructed, and by every durable step before it replays a
    /// value — the second because a host with no store is not the only door, and a
    /// store handed a run from elsewhere must not answer from another definition's
    /// records either.
    ///
    /// # Errors
    ///
    /// [`StoreError::Incompatible`] naming the axis that disagrees. A run this
    /// store has no record for is not refused.
    pub fn check_definition(
        &self,
        run: RunId,
        identity: &DefinitionIdentity,
    ) -> Result<(), StoreError> {
        let Some(recorded) = self.definition_of(run) else {
            return Ok(());
        };
        match identity.drift_from(&recorded) {
            None => Ok(()),
            Some(drift) => Err(StoreError::Incompatible {
                run: run.id().to_hex(),
                drift,
            }),
        }
    }

    /// Turn a borrowed record into one that can cross onto the storage owner's
    /// thread.
    ///
    /// The stored record owns its tenant and path, so the hop is the only part that
    /// has to be `'static`, and it takes the record by value rather than borrowing
    /// the caller's strings for the length of a flush.
    fn owned(&self, record: StagedRecord<'_>) -> Staged {
        Staged {
            record: Stored {
                run: record.run(),
                key: record.key().as_bytes().to_vec(),
                tenant: record.tenant().to_owned(),
                path: record.path().to_owned(),
                definition: StoredDefinition::of(record.definition()),
                value: record.bytes().to_vec(),
            },
            key: record.key(),
        }
    }

    /// The index's lock, recovering a poison.
    ///
    /// Recovering is right here for the reason it is right in `journal::owner`: no
    /// arm in the append path panics while the lock is held — every one returns its
    /// `Result` first — so a poison is a bug in an unrelated thread, and
    /// propagating it would brick a store that is still perfectly readable.
    fn index(&self) -> MutexGuard<'_, Index> {
        owner::lock(&self.inner.index)
    }

    /// Refuse a durable step's read when the injector is armed.
    ///
    /// The injector is aimed at the [`RunRecords`] doors rather than at the index
    /// itself, and that placement is the whole design. Both the host's admission
    /// pre-flight and the durable step's replay check read the same in-memory
    /// index through [`Self::check_definition`], so a fault aimed at "the index"
    /// would be spent by whichever of the two ran first — which for a resumed run
    /// is always the host, and a fault the host absorbs tests the host's refusal
    /// rather than the step's, proving nothing about `remember_at`. The
    /// [`RunRecords`] trait is reached only from `Scope::remember`, so arming it
    /// here makes the fault land on the check under test whatever else the run
    /// path has already read.
    ///
    /// One-shot, and bounded by construction: the first step read that observes it
    /// takes it with a `swap`, so a caller that arms it and never reads is not left
    /// with a store that fails forever, and two concurrent reads cannot both take
    /// the same fault.
    fn step_readable(&self) -> Result<(), StoreError> {
        if self.inner.unreadable.swap(false, Ordering::SeqCst) {
            let refusal = Err(StoreError::storage(std::io::Error::other(
                "the run store's device refused a read of its index",
            )));
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "step_readable: returning an error to the caller");
            return refusal;
        }
        Ok(())
    }
}

/// The store as the script module reaches it: a typed lookup and a typed append
/// over archived bytes, with no `lgwks_bot::task` in its own surface.
///
/// The implementation is here; the trait lives in
/// [`script`] because the task module is a *consumer* of a
/// durable step, and a step must be markable without the host that runs it.
impl RunRecords for RunStore {
    fn lookup(
        &self,
        tenant: &str,
        run: RunId,
        key: StepKey,
    ) -> Result<Option<StoredValue>, FlowError> {
        self.step_readable()?;
        let index = self.index();
        let Some(held) = index.runs.get(&run) else {
            return Ok(None);
        };
        if held.tenant != tenant {
            let refusal = Err(StoreError::ForeignTenant {
                owner: held.tenant.clone(),
                asked: tenant.to_owned(),
            }
            .into());
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "lookup: returning an error to the caller");
            return refusal;
        }
        Ok(held
            .steps
            .get(&key.to_hex())
            .map(|step| StoredValue::new(step.path.clone(), step.bytes.clone())))
    }

    fn append(
        &self,
        tenant: &str,
        run: RunId,
        key: StepKey,
        path: &str,
        definition: &DefinitionIdentity,
        bytes: Vec<u8>,
    ) -> Result<Appended, FlowError> {
        let staged = self.owned(RunRecords::stage(
            self, tenant, run, key, path, definition, bytes,
        ));
        self.commit(staged.record, staged.key)
            .map_err(FlowError::from)
    }

    /// Whether the records for `run` were written under `definition`.
    ///
    /// The check every durable step makes before it returns a recorded value. A
    /// run this store holds nothing for is compatible by definition: that is a
    /// first attempt, and refusing it would make the first record unreachable.
    ///
    /// # Errors
    ///
    /// Whatever [`Self::check_definition`] reports, unchanged: a drift as
    /// [`FlowError::Incompatible`], and a device refusal as itself. Neither is ever
    /// collapsed into the other's answer.
    fn compatibility(
        &self,
        run: RunId,
        definition: &DefinitionIdentity,
    ) -> Result<bool, FlowError> {
        self.step_readable()?;
        self.check_definition(run, definition)
            .map(|()| true)
            .map_err(FlowError::from)
    }

    /// Which axis of `run`'s recorded identity disagrees with `definition`.
    ///
    /// The same comparison [`Self::compatibility`] made, returning the axis rather
    /// than the boolean. Two doors rather than one returning a richer type because
    /// the boolean is asked on every durable step — including the ones that pass,
    /// where the axis is never wanted — and building a [`Drift`] that names two
    /// digests or two schema ids on a step that will replay cleanly is work spent
    /// only to be thrown away.
    ///
    /// # Errors
    ///
    /// [`FlowError::Failed`] when this store holds no record for `run`, so there is
    /// no recorded identity to compare against, or when the two doors disagree —
    /// `compatibility` refused this pair while this comparison finds nothing to
    /// refuse. Naming an axis in either case would mean inventing one.
    fn drift(&self, run: RunId, definition: &DefinitionIdentity) -> Result<Drift, FlowError> {
        let Some(recorded) = self.definition_of(run) else {
            let refusal = Err(FlowError::failed(format!(
                "this store holds no definition identity for run {}, so it cannot name a drift \
             axis; refusing to invent one",
                run.id().to_hex()
            )));
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "drift: returning an error to the caller");
            return refusal;
        };
        match definition.drift_from(&recorded) {
            Some(drift) => Ok(drift),
            None => Err(FlowError::failed(format!(
                "run {} was refused as incompatible and yet agrees with the declared \
                 definition; refusing to name an axis for a disagreement that is not there",
                run.id().to_hex()
            ))),
        }
    }

    /// The awaited door, which is the one a durable step takes.
    ///
    /// The `sync_all` behind it is the device's cost, and this is what keeps that
    /// cost off the executor: the step's future waits for the owner's answer
    /// instead of sitting through the flush (INV-BOT-50).
    fn append_async<'a>(
        &'a self,
        record: StagedRecord<'a>,
    ) -> crate::BoxFuture<'a, Result<Appended, FlowError>> {
        let staged = self.owned(record);
        Box::pin(async move {
            self.commit_async(staged.record, staged.key)
                .await
                .map_err(FlowError::from)
        })
    }
}

/// Map one owner refusal onto this store's vocabulary.
///
/// One conversion rather than a match per call site, because every `?` on an owner
/// refusal wraps the same fact — the device, or the handle, said no — and a literal
/// per call site would be four chances to name a different variant for one event.
/// The refusal's own message distinguishes a full queue from a poison from an
/// undelivered answer, so nothing is lost by carrying it as the device's error.
impl From<SubmitError> for StoreError {
    fn from(cause: SubmitError) -> Self {
        Self::Storage {
            cause: cause.into_io(),
        }
    }
}

/// One staged append: the record to write and the step key it is for.
///
/// Its own type so the two append doors build the record once between them, and
/// the key travels with it rather than being threaded beside it at each call site.
struct Staged {
    /// The record as it will be archived.
    record: Stored,
    /// The step key the already-recorded and conflict checks are on.
    key: StepKey,
}

/// The in-file chain head of one run, held with the index.
impl RunStore {
    /// The chain step of `stored`: the head of the whole file, which every run
    /// shares, so a frame's head cannot be valid under one run's history and
    /// invalid under the file's.
    ///
    /// One chain for the file rather than one per run, deliberately. A per-run
    /// chain would need the head to be re-found by replaying only that run's
    /// frames, and a caller that read one run's records could then not verify
    /// them without reading every other tenant's records too. One file, one
    /// chain, and the tenant check above is what keeps the reads separated.
    ///
    /// The whole of this runs on the storage owner's thread — the in-memory checks,
    /// the fence, the write, the `sync_all` and the fold into the index — so two
    /// appends to this store are two ordered steps rather than two racing ones,
    /// and none of them parks the executor the `sync_all` does. INV-BOT-50.
    fn commit(&self, stored: Stored, key: StepKey) -> Result<Appended, StoreError> {
        self.inner
            .owner
            .submit(move |file, index| append_on_owner(file, index, &stored, key))
            .map_err(StoreError::from)
    }

    /// The same append, awaited rather than sat through.
    ///
    /// The one a durable step takes, because a step's body is a future and the
    /// `sync_all` its record costs is the device's, not the executor's.
    fn commit_async<'a>(
        &'a self,
        stored: Stored,
        key: StepKey,
    ) -> crate::BoxFuture<'a, Result<Appended, StoreError>> {
        Box::pin(async move {
            self.inner
                .owner
                .submit_async(move |file, index| append_on_owner(file, index, &stored, key))
                .await
                .map_err(StoreError::from)
        })
    }
}

/// The ordered append, over the file and the index it folds into.
///
/// A free function taking the shared state rather than a method on `&self`, because
/// it runs on the storage owner's thread with a closure that must own everything it
/// touches: the `Arc` the store's state lives in and the record itself.
///
/// The lock is taken for the stage half only. The settle half is a fold into the
/// same index, and it runs on this same thread under the owner's ordering — after
/// the batch's one `sync_all`, never before it — so no append can observe a
/// half-committed index and no `.await` is ever reached holding it.
///
/// # Errors
///
/// Every [`StoreError`] an append can produce, carried as the device's own error so
/// the owner's one reply channel serves both stores.
fn append_on_owner(
    file: &mut File,
    shared: &Arc<Mutex<Index>>,
    stored: &Stored,
    key: StepKey,
) -> std::io::Result<Stage<Appended, Arc<Mutex<Index>>>> {
    let mut index = owner::lock(shared);
    decide_and_write(file, &mut index, stored, key)
}

/// The decision, the write, and the fold this record owes its batch's flush.
///
/// Every check that can refuse in memory runs before a byte moves, and the length
/// fence runs inside the owner's ordered step where no other append can overtake it
/// — which is why the file's `metadata` is read here rather than through the
/// read-only view: the authoritative answer is the one taken at the moment of the
/// write.
///
/// The three ceilings and the tenant and duplicate checks are decided *here*, at
/// stage time, exactly as they were before group commit: a refusal is the same
/// refusal and refuses the same single record. The fold into the index is what the
/// settle half does, after the batch's one `sync_all` returns, so an index never
/// names a record the device did not take.
///
/// # Errors
///
/// Every [`StoreError`] an append can produce. Nothing is written by any of them,
/// and a refusal leaves the index and the file byte-identical.
fn decide_and_write(
    file: &mut File,
    index: &mut Index,
    stored: &Stored,
    key: StepKey,
) -> std::io::Result<Stage<Appended, Arc<Mutex<Index>>>> {
    let owned = index.runs.get(&stored.run);
    if let Some(owned) = owned {
        if owned.tenant != stored.tenant {
            let refusal = Err(refusal(StoreError::ForeignTenant {
                owner: owned.tenant.clone(),
                asked: stored.tenant.clone(),
            }));
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
            return refusal;
        }
        if let Some(step) = owned.steps.get(&key.to_hex()) {
            // Already recorded. An identical record is the same fact seen
            // twice — a step whose value was returned and whose caller
            // recorded it again — and a different one is a conflict. Neither
            // writes a byte.
            return Ok(Stage::Settled(Ok(if step.bytes == stored.value {
                Appended::AlreadyRecorded
            } else {
                Appended::Conflicting
            })));
        }
        if u64::saturating_from(owned.steps.len()) >= MAX_RECORDS_PER_RUN {
            let refusal = Err(refusal(StoreError::Limit {
                kind: StoreLimitKind::Records,
                requested: u64::saturating_from(owned.steps.len()).saturating_add(1),
                limit: MAX_RECORDS_PER_RUN,
            }));
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
            return refusal;
        }
    }

    let previous = tail_head(index);
    let (framed, head) = frame(stored, &previous).map_err(refusal)?;

    let staged = u64::saturating_from(framed.len());
    let next = index
        .committed
        .checked_add(staged)
        .ok_or_else(|| refusal(store_full(u64::MAX)))?;
    if next > MAX_STORE_BYTES {
        let refusal = Err(refusal(store_full(next)));
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
        return refusal;
    }

    // The fence, checked before a byte is written: the length this handle
    // indexed must be the length on the disk, so a file another writer moved
    // under it is refused here rather than forked. It runs on the owner thread,
    // so the check and the write it guards cannot be split by another append.
    let on_disk = file
        .metadata()
        .map_err(StoreError::storage)
        .map_err(refusal)?
        .len();
    if on_disk != index.committed {
        let refusal = Err(refusal(StoreError::Corrupt {
            at: u64::saturating_from(index.runs.get(&stored.run).map_or(0, |run| run.steps.len())),
        }));
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
        return refusal;
    }

    // The chain head and the committed length move *before* the write, and are the
    // batch's own bookkeeping rather than a claim the device took anything: the next
    // member of this batch must chain over this record, and the next batch's length
    // fence must expect these bytes. A head that lagged a write would fork the chain.
    index.committed = next;
    index.tail = head;

    file.write_all(&framed)
        .map_err(StoreError::storage)
        .map_err(refusal)?;

    // The readable half of the fold waits for the flush: a record a lookup can find
    // is a record whose bytes the device has acknowledged, so publishing it earlier
    // would let a concurrent reader see a value a crash could take back.
    let held = Held {
        path: stored.path.clone(),
        bytes: stored.value.clone(),
    };
    let key = hex_of(&stored.key);
    let run = stored.run;
    let tenant = stored.tenant.clone();
    let definition = stored.definition.to_identity();
    Ok(Stage::Unsynced {
        answer: Appended::Recorded,
        bytes: framed.len(),
        settle: Box::new(move |index: &mut Arc<Mutex<Index>>| {
            let mut index = owner::lock(index);
            let entry = index.runs.entry(run).or_insert_with(|| RunIndex {
                tenant: tenant.clone(),
                definition: definition.clone(),
                steps: HashMap::new(),
            });
            entry.steps.insert(key, held);
        }),
    })
}

/// A store refusal as the device error the owner's one reply channel carries.
///
/// One conversion rather than a `map_err` per arm, because every `?` on a store
/// refusal wraps the same fact — this store said no — and a literal per site would
/// be seven chances to name a different message for one event. The refusal's own
/// `Display` carries the ceiling or the tenant that names it.
fn refusal(cause: StoreError) -> std::io::Error {
    std::io::Error::other(cause.to_string())
}

/// The chain head the next frame in this file follows.
///
/// The index carries the tail head beside the committed length, because the head
/// cannot be recomputed from the index: a frame's head is over the *archived*
/// record, and re-archiving a value to hash it would be a second encoding of the
/// same bytes with its own chance of disagreeing with what is on the disk.
fn tail_head(index: &Index) -> Digest {
    index.tail
}

/// The record as it is archived and as it is checked.
#[derive(
    Debug, Clone, lgwks_std::wire::Archive, lgwks_std::wire::Serialize, lgwks_std::wire::Deserialize,
)]
#[rkyv(
    attr(non_exhaustive),
    crate = lgwks_std::wire::rkyv,
    compare(PartialEq),
    derive(Debug)
)]
struct Stored {
    /// The run this record belongs to.
    #[rkyv(attr(doc = "The run this record belongs to."))]
    run: RunId,
    /// The step key's digest.
    #[rkyv(attr(doc = "The step key's digest."))]
    key: Vec<u8>,
    /// The tenant that minted the run.
    #[rkyv(attr(doc = "The tenant that minted the run."))]
    tenant: String,
    /// The step's path, for attribution.
    #[rkyv(attr(doc = "The step's path, for attribution."))]
    path: String,
    /// The definition the record was written under, as its fields.
    ///
    /// Archived as fields rather than as one digest because a drift refusal has
    /// to *name* the two sides: a head over the digest could detect the
    /// disagreement but could not report which axis it was, and the fields are
    /// what a caller reads to decide whether it needs a new run or a migration.
    #[rkyv(attr(doc = "The definition the record was written under, as its fields."))]
    definition: StoredDefinition,
    /// The archived value the step returned.
    #[rkyv(attr(doc = "The archived value the step returned."))]
    value: Vec<u8>,
}

/// A [`DefinitionIdentity`] in the shape a record archives.
///
/// Named rather than archived directly so the record's on-disk shape is stated
/// once, in the store that writes it, and a future revision of the definition
/// vocabulary is a new struct here rather than a change to a type every consumer
/// of [`Drift`] is already matching.
#[derive(
    Debug, Clone, lgwks_std::wire::Archive, lgwks_std::wire::Serialize, lgwks_std::wire::Deserialize,
)]
#[rkyv(
    attr(non_exhaustive),
    crate = lgwks_std::wire::rkyv,
    compare(PartialEq),
    derive(Debug)
)]
struct StoredDefinition {
    /// The task name.
    #[rkyv(attr(doc = "The task name."))]
    name: String,
    /// The declared definition revision.
    #[rkyv(attr(doc = "The declared definition revision."))]
    revision: u64,
    /// The input digest.
    #[rkyv(attr(doc = "The input digest."))]
    input: Vec<u8>,
    /// How many durable steps the definition declares.
    #[rkyv(attr(doc = "How many durable steps the definition declares."))]
    steps: u64,
    /// The declared durable-value schema id.
    #[rkyv(attr(doc = "The declared durable-value schema id."))]
    codec: String,
}

/// The 32 bytes a record archived, as the digest the identity carries back.
///
/// A fixed width rather than a slice, because a record whose input field is not
/// 32 bytes is a corrupt frame and the checked decoder is what says so; the
/// conversion cannot silently truncate, because `as` is forbidden in this
/// workspace and this is the honest alternative to it.
fn digest_from_record(bytes: &[u8]) -> Digest {
    let mut digest = [0u8; 32];
    let copied = bytes.len().min(digest.len());
    digest[..copied].copy_from_slice(&bytes[..copied]);
    Digest::from_bytes(digest)
}

impl StoredDefinition {
    /// The record's archived form of a live identity.
    fn of(identity: &DefinitionIdentity) -> Self {
        Self {
            name: identity.name().to_owned(),
            revision: identity.revision(),
            input: identity.input().as_bytes().to_vec(),
            steps: u64::saturating_from(identity.steps()),
            codec: identity.codec().to_owned(),
        }
    }

    /// The live identity this record was written under.
    fn to_identity(&self) -> DefinitionIdentity {
        DefinitionIdentity::new(
            &self.name,
            self.revision,
            digest_from_record(&self.input),
            usize::saturating_from(self.steps),
        )
        .with_codec(&self.codec)
    }
}

impl Stored {
    /// This record's chain head over the head before it.
    ///
    /// Framed exactly as the step key is framed, so a record's head cannot be
    /// confused with the digest of the bytes it holds: the run, the key, the
    /// tenant, the path and every definition field are length-framed, and the
    /// value is framed last so no two different records hash the same byte
    /// string. The definition's own bytes go in through the identity's framing
    /// so a head over the record commits to the identity the same way the
    /// identity's digest commits to it.
    fn head_from(&self, previous: &Digest, _archived: &[u8]) -> Digest {
        let mut hasher = Hasher::new();
        hasher.write_framed(previous.as_bytes());
        hasher.write_framed(&self.run_key_bytes());
        hasher.write_framed(&self.key);
        hasher.write_framed(self.tenant.as_bytes());
        hasher.write_framed(self.path.as_bytes());
        hasher.write_framed(self.definition.name.as_bytes());
        hasher.write_framed(&self.definition.revision.to_le_bytes());
        hasher.write_framed(&self.definition.input);
        hasher.write_framed(&self.definition.steps.to_le_bytes());
        hasher.write_framed(self.definition.codec.as_bytes());
        hasher.write_framed(&self.value);
        hasher.finalize()
    }

    /// The run id's 16 bytes, as a field the framing can hold.
    ///
    /// The hex spelling rather than the big-endian integer, because
    /// `RunId::to_hex` is the one conversion this crate defines for the role
    /// and a second one here would be a second definition to keep in step.
    fn run_key_bytes(&self) -> Vec<u8> {
        self.run.id().to_hex().into_bytes()
    }
}

/// The frame for one record: length, payload, head.
///
/// The head is stored rather than recomputed on faith. That is what separates a
/// torn tail (dropped, because a write leaves a prefix and a prefix has no head)
/// from a frame that does not follow (refused, because those bytes were
/// acknowledged).
fn frame(stored: &Stored, previous: &Digest) -> Result<(Vec<u8>, Digest), StoreError> {
    // Archive, bound, chain and lay out — the same step the effect journal takes.
    // Only the record's own archiving and its step-chaining head are this store's.
    frame::frame_record(
        stored,
        previous,
        MAX_RECORD_BYTES,
        |stored| {
            to_bytes::<WireError>(stored)
                .map(|bytes| bytes.as_ref().to_vec())
                .map_err(|cause| StoreError::Encoding { cause })
        },
        Stored::head_from,
        record_too_large,
    )
}

/// The refusal for a store that has no room left, named by what it would have
/// needed.
///
/// One constructor because both call sites discovered the same fact by different
/// routes — an addition that overflowed and an addition that came out past the
/// ceiling — and a literal per site would be two chances to name a different bound
/// for the same ceiling.
fn store_full(requested: u64) -> StoreError {
    StoreError::Limit {
        kind: StoreLimitKind::StoreBytes,
        requested,
        limit: MAX_STORE_BYTES,
    }
}

/// The refusal for a record whose archived bytes are past this store's ceiling.
///
/// One constructor, because the message and the two numbers are the same fact
/// whichever call site discovers it, and a literal per site would be two chances to
/// name a different bound for the same ceiling.
fn record_too_large(len: usize) -> StoreError {
    StoreError::Limit {
        kind: StoreLimitKind::RecordBytes,
        requested: u64::saturating_from(len),
        limit: u64::saturating_from(MAX_RECORD_BYTES),
    }
}

/// Read `buf` in full, reporting whether all of it arrived.
fn read_full(reader: &mut impl Read, buf: &mut [u8]) -> Result<bool, StoreError> {
    let mut filled = 0usize;
    while filled < buf.len() {
        match reader.read(&mut buf[filled..]) {
            Ok(0) => return Ok(false),
            Ok(read) => filled = filled.saturating_add(read),
            Err(ref error) if error.kind() == std::io::ErrorKind::Interrupted => {}
            Err(cause) => {
                let refusal = Err(StoreError::storage(cause));
                lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "read_full: returning an error to the caller");
                return refusal;
            }
        }
    }
    Ok(true)
}

/// One complete, decoded frame read from the store.
struct Framed {
    /// The record the frame's payload decoded to.
    stored: Stored,
    /// The payload bytes, which the chain head is computed over.
    payload: Vec<u8>,
    /// The chain head the frame recorded for itself.
    head: [u8; HEAD_BYTES],
    /// The payload length the frame declared.
    declared: usize,
}

/// Read the frame at ordinal `at`, or `None` where the file stops holding a
/// whole frame.
///
/// A partial length prefix, or a payload or head cut short, is an append that
/// never finished: it was never anyone's answer, so the scan stops and the
/// caller trims to `committed`. A complete prefix naming a frame this store
/// never writes (a write leaves a prefix of a length it did finish computing,
/// and that length was always legal) or a payload that does not decode cannot
/// be an interrupted append and is refused. So is a cut whose bytes hold a frame
/// this store acknowledged under another length: `cursor` is where the frame
/// starts and the head it chains from, which is what the stored head is checked
/// against (#262).
fn next_frame(file: &mut File, cursor: &Cursor<'_>) -> Result<Option<Framed>, StoreError> {
    let at = cursor.at;
    let corrupt = || StoreError::Corrupt { at };
    let Some(raw) = frame::read_raw(
        file,
        cursor,
        MAX_RECORD_BYTES,
        StoreError::storage,
        corrupt,
        |previous, payload| {
            let aligned = frame::decodable(payload);
            from_bytes::<Stored, WireError>(aligned.as_slice())
                .ok()
                .map(|stored| stored.head_from(previous, payload))
        },
    )?
    else {
        return Ok(None);
    };
    let stored = from_bytes::<Stored, WireError>(&raw.payload).map_err(|error| {
        lgwks_std::trace::debug!(?error, at, "next_frame: the payload did not decode");
        StoreError::Corrupt { at }
    })?;
    Ok(Some(Framed {
        stored,
        declared: raw.payload.len(),
        payload: raw.payload,
        head: raw.head,
    }))
}

/// Read every committed frame, returning the index it implies.
///
/// The scan is bounded by the file's own length: every iteration consumes at
/// least one byte, and a file that grows underneath the scan cannot extend it,
/// because the reader stops at the length it saw when it began. A frame whose
/// stored head does not follow is [`StoreError::Corrupt`] and is not truncated.
fn replay(file: &mut File) -> Result<Index, StoreError> {
    let total = file.metadata().map_err(StoreError::storage)?.len();
    file.seek(SeekFrom::Start(0)).map_err(StoreError::storage)?;

    let mut header = [0u8; STORE_HEADER.len()];
    if !read_full(file, &mut header)? {
        let refusal = Err(StoreError::NotAStore);
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
        return refusal;
    }
    // The two refusals are separated, and the version one is checked first,
    // because they are different facts and only one of them is a version the
    // writer could have meant. Every version of this format shares the whole
    // magic *except* its last byte, which is that format's version — so a file
    // whose first fifteen bytes match is this store's, at some version, and a
    // file whose first fifteen do not match was never a run store at all. That
    // is the whole reason `STORE_MAGIC` is 15 bytes of constant plus a version
    // byte rather than one 16-byte constant.
    check_format_version(header)?;
    let mut index = Index {
        runs: HashMap::new(),
        committed: u64::saturating_from(STORE_HEADER.len()),
        tail: genesis_head(),
    };
    let mut previous = genesis_head();
    let mut at = 0u64;
    while let Some(Framed {
        stored,
        payload,
        head,
        declared,
    }) = next_frame(file, &Cursor::new(at, index.committed, &previous))?
    {
        if stored.head_from(&previous, &payload) != Digest::from_bytes(head) {
            let refusal = Err(StoreError::Corrupt { at });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
            return refusal;
        }
        let frame_len = frame::framed_len(declared);
        // A frame that runs past the end of the file is an interrupted append.
        // Readable as a whole above, so this only catches a header that lied
        // about a frame already acknowledged, which is rot rather than a crash.
        if index
            .committed
            .checked_add(frame_len)
            .is_none_or(|end| end > total)
        {
            break;
        }
        index.committed = index.committed.saturating_add(frame_len);
        index.tail = stored.head_from(&previous, &payload);
        previous = stored.head_from(&previous, &payload);
        at = at.saturating_add(1);
        if at > MAX_RECORDS_PER_RUN {
            let refusal = Err(StoreError::Limit {
                kind: StoreLimitKind::Records,
                requested: at,
                limit: MAX_RECORDS_PER_RUN,
            });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
            return refusal;
        }
        let run = index.runs.entry(stored.run).or_insert_with(|| RunIndex {
            tenant: stored.tenant.clone(),
            definition: stored.definition.to_identity(),
            steps: HashMap::new(),
        });
        run.steps.insert(
            hex_of(&stored.key),
            Held {
                path: stored.path.clone(),
                bytes: stored.value.clone(),
            },
        );
    }
    Ok(index)
}

/// Lowercase hex, for a digest used as a map key.
///
/// `lgwks_std::hash::Digest::to_hex` takes a `Digest`, and this is the same 32
/// bytes before the head has been computed; formatting them here rather than
/// reaching into the digest type keeps the key exactly the bytes on the disk.
fn hex_of(bytes: &[u8]) -> String {
    const DIGITS: &[u8; 16] = b"0123456789abcdef";
    let mut hex = Vec::with_capacity(bytes.len().saturating_mul(2));
    for byte in bytes.iter().copied() {
        hex.push(DIGITS[usize::from(byte >> 4) & 0x0f]);
        hex.push(DIGITS[usize::from(byte) & 0x0f]);
    }
    // Every byte pushed above is an ASCII digit, so this is well-formed UTF-8
    // and the conversion cannot fail; `String::from_utf8` reports the invariant
    // rather than trusting it.
    hex.into_iter().map(char::from).collect()
}

/// A frame whose length prefix lies is refused, and an append that was cut is
/// repaired (#262). The file is built frame by frame with the store's own framing,
/// then one prefix is changed: the bytes are the acknowledged ones, so nothing but
/// the stored head can say whether the length is the writer's.
#[cfg(test)]
mod tests {
    use super::*;
    use crate::journal::frame::probe::{Scratch, declared_at, frame_starts, with_prefix};

    type TestResult = Result<(), Box<dyn std::error::Error>>;

    /// The store file of `count` records, written to `path`, and its bytes.
    fn written(path: &Path, count: u8) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
        let run = RunId::from_hex(&format!("a{}", "0".repeat(31)))?;
        let identity =
            DefinitionIdentity::new("tail", 1, Digest::from_bytes([1; 32]), usize::from(count));
        let mut bytes = STORE_HEADER.to_vec();
        let mut previous = genesis_head();
        for step in 0..count {
            let stored = Stored {
                run,
                key: vec![step; 32],
                tenant: "acme".to_owned(),
                path: format!("s{step}"),
                definition: StoredDefinition::of(&identity),
                value: vec![step; 40 + usize::from(step)],
            };
            let (framed, head) = frame(&stored, &previous)?;
            bytes.extend_from_slice(&framed);
            previous = head;
        }
        std::fs::write(path, &bytes)?;
        Ok(bytes)
    }

    /// Open `path` and require a refusal as `Corrupt` at frame `at`, with the file
    /// byte-identical to `before`.
    fn require_refused(path: &Path, before: &[u8], at: u64, why: &str) -> TestResult {
        let outcome: TestResult = match RunStore::open(path) {
            Err(StoreError::Corrupt { at: named }) => {
                assert_eq!(named, at, "{why}");
                Ok(())
            }
            Err(other) => Err(format!("{why}: expected Corrupt, got {other}").into()),
            Ok(opened) => Err(format!(
                "{why}: reopened at {} committed bytes, so an acknowledged frame was lost",
                opened.committed_bytes()
            )
            .into()),
        };
        outcome?;
        assert_eq!(std::fs::read(path)?, before, "{why}: refused bytes move");
        Ok(())
    }

    /// The acknowledged final record whose prefix grows is the defect: `extra` up to
    /// 32 leaves the payload whole and the head short, more ends inside the payload.
    #[test]
    fn a_lengthened_acknowledged_final_record_is_refused_not_trimmed() -> TestResult {
        let scratch = Scratch::new("store-lengthened")?;
        let bytes = written(scratch.path(), 3)?;
        let last = frame_starts(&bytes, STORE_HEADER.len())?[2];
        let declared = declared_at(&bytes, last);
        for extra in (1u32..=40).chain([100, 255, 1024, 4096]) {
            let lied = with_prefix(&bytes, last, declared + extra);
            std::fs::write(scratch.path(), &lied)?;
            require_refused(scratch.path(), &lied, 2, &format!("final record L+{extra}"))?;
        }
        Ok(())
    }

    /// An inflated length on a record with records behind it: the true record is still
    /// first behind the prefix, so it authenticates and the frames after it are kept.
    #[test]
    fn an_inflated_record_with_records_behind_it_is_refused_untouched() -> TestResult {
        let scratch = Scratch::new("store-inflated")?;
        let bytes = written(scratch.path(), 3)?;
        let middle = frame_starts(&bytes, STORE_HEADER.len())?[1];
        let remaining = u32::try_from(bytes.len() - middle)?;
        for declared in [remaining, remaining + 1, remaining + 31, remaining + 500] {
            let lied = with_prefix(&bytes, middle, declared);
            std::fs::write(scratch.path(), &lied)?;
            require_refused(scratch.path(), &lied, 1, &format!("declared {declared}"))?;
        }
        Ok(())
    }

    /// A lengthened prefix over a record whose own bytes are also damaged still
    /// refuses, because the record behind it authenticates.
    #[test]
    fn a_damaged_cut_record_with_an_acknowledged_one_behind_it_is_refused() -> TestResult {
        let scratch = Scratch::new("store-damaged")?;
        let bytes = written(scratch.path(), 3)?;
        let middle = frame_starts(&bytes, STORE_HEADER.len())?[1];
        let remaining = u32::try_from(bytes.len() - middle)?;
        let mut lied = with_prefix(&bytes, middle, remaining + 7);
        if let Some(byte) = lied.get_mut(middle + 12) {
            *byte ^= 0x55;
        }
        std::fs::write(scratch.path(), &lied)?;
        require_refused(scratch.path(), &lied, 1, "damaged middle, lengthened")
    }

    /// The control: an append cut at any place in the final record is still repaired,
    /// back to exactly the acknowledged prefix. Sampled at each boundary and between
    /// them, since each repair is an `fsync`.
    #[test]
    fn an_append_cut_inside_the_final_record_is_repaired() -> TestResult {
        let scratch = Scratch::new("store-cut")?;
        let bytes = written(scratch.path(), 3)?;
        let last = frame_starts(&bytes, STORE_HEADER.len())?[2];
        let kept = u64::try_from(last)?;
        let whole = bytes.len() - last;
        for cut in [
            1,
            3,
            4,
            5,
            whole >> 1,
            whole - 33,
            whole - 32,
            whole - 31,
            whole - 1,
        ] {
            std::fs::write(
                scratch.path(),
                bytes.get(..last + cut).ok_or("cut past the end")?,
            )?;
            let store = RunStore::open(scratch.path())?;
            assert_eq!(
                store.committed_bytes(),
                kept,
                "cut {cut}: the two records survive"
            );
            drop(store);
            assert_eq!(std::fs::metadata(scratch.path())?.len(), kept, "cut {cut}");
        }
        Ok(())
    }

    /// The worst tail a repair can be asked to examine, a full-ceiling prefix over
    /// noise, is decided in bounded time: the first search tries a payload length per
    /// byte and each try decodes, so the cost is the thing to bound.
    #[test]
    fn a_ceiling_sized_noise_tail_is_trimmed_in_bounded_time() -> TestResult {
        let scratch = Scratch::new("store-noise")?;
        let mut bytes = written(scratch.path(), 1)?;
        let kept = u64::try_from(bytes.len())?;
        bytes.extend_from_slice(&u32::try_from(MAX_RECORD_BYTES)?.to_be_bytes());
        let mut state = 0x9e37_79b9_7f4a_7c15u64;
        for _ in 0..MAX_RECORD_BYTES + HEAD_BYTES - 1 {
            state = state
                .wrapping_mul(6_364_136_223_846_793_005)
                .wrapping_add(1);
            bytes.push(state.to_be_bytes()[0]);
        }
        std::fs::write(scratch.path(), &bytes)?;
        let started = std::time::Instant::now();
        let store = RunStore::open(scratch.path())?;
        let elapsed = started.elapsed();
        assert_eq!(
            store.committed_bytes(),
            kept,
            "noise holds no acknowledged record"
        );
        assert!(elapsed.as_secs() < 5, "the search took {elapsed:?}");
        Ok(())
    }
}