vgi-bridge 0.16.0

The per-community VGI bridge: holds the community's forge credentials, takes git-ns/bridge jobs from its VTC over TSP or DIDComm, runs the forge adapters, reports results, events and drift, and posts the commit-trust check where no forge feature can be trusted to.
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
//! Durable state, in one redb file.
//!
//! What the bridge must remember across a restart, and why:
//!
//! - **The job ledger** ([`JobRecord`]): `jobId` idempotency (a repeated job
//!   is answered, never run twice), the result until the VTC acknowledges it,
//!   and jobs still to run.
//! - **Namespaces** ([`NamespaceRecord`]): the binding the VTC's namespace id
//!   maps to, the capabilities found while binding (and any later
//!   `CapabilityChanged`), and on GitHub the managed repository set and the
//!   required-workflow pin — which the adapter refuses to act without after a
//!   restart, so [`crate::Bridge::restore`] hands them back first.
//! - **Repositories** ([`RepoRecord`]): forge id, owners and the roles the
//!   bridge projected — the projection an `inspect` compares against, since a
//!   job carries none.
//! - **Pending flows** ([`PendingFlow`]): binds and account links waiting for
//!   the person, keyed by their single-use `state`.
//! - **The outbox** ([`OutboxEntry`]): results and events not yet
//!   acknowledged by the VTC.
//! - **Sealed secrets**: see [`crate::seal`]. Only ciphertext is written.
//! - **The provenance ledger** ([`BranchLedger`]): who pushed to each
//!   `dependabot/*` branch since it was created. The Dependabot re-sign acts
//!   only on an unbroken record, so losing it means Dependabot pull requests
//!   open before the loss are not re-signed (`@dependabot recreate` starts
//!   them over).
//!
//! Every value is JSON in a string-keyed table: small, inspectable in a
//! support session, and free of a schema migration story for records this
//! size. Writes are one transaction each, committed durably (fsync) before
//! the caller goes on — in particular a job is on disk before the bridge
//! answers `accepted: true`.

use std::collections::BTreeSet;
use std::path::Path;
use std::sync::Arc;

use anyhow::{Context, Result};
use redb::{Database, ReadableDatabase, ReadableTable, TableDefinition};
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use vgi_forge::{Capabilities, ForgeAccount, NamespaceBinding, Resource, RoleAssignment};
use zeroize::Zeroizing;

use crate::seal::MasterKey;

/// One logical collection.
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
#[non_exhaustive]
pub enum Table {
    /// [`JobRecord`] by `jobId`.
    Jobs,
    /// [`NamespaceRecord`] by the VTC's namespace id.
    Namespaces,
    /// [`RepoRecord`] by `<host>#<forge id>`.
    Repos,
    /// [`PendingFlow`] by its `state`.
    Pending,
    /// [`OutboxEntry`] by `result:<jobId>` / `event:<id>`.
    Outbox,
    /// Webhook delivery ids already handled, with when.
    Deliveries,
    /// Sealed secrets by name.
    Secrets,
    /// Small bookkeeping values (last token rotation, …) by name.
    Meta,
    /// [`BranchLedger`] by `<host>#<repository id>#<branch>`: the Dependabot
    /// re-sign's provenance ledger.
    Branches,
    /// VTA mode, this host only: the version each app-state record was at
    /// when this host last wrote or read it (`u64` by remote key). What
    /// tells a record this host mirrored (and another host has since
    /// deleted) from one it never did.
    Mirror,
}

impl Table {
    const ALL: [Table; 10] = [
        Table::Jobs,
        Table::Namespaces,
        Table::Repos,
        Table::Pending,
        Table::Outbox,
        Table::Deliveries,
        Table::Secrets,
        Table::Meta,
        Table::Branches,
        Table::Mirror,
    ];

    /// The table's name (redb's, and the VTA mirror's key segment).
    pub fn name(self) -> &'static str {
        match self {
            Table::Jobs => "jobs",
            Table::Namespaces => "namespaces",
            Table::Repos => "repos",
            Table::Pending => "pending",
            Table::Outbox => "outbox",
            Table::Deliveries => "deliveries",
            Table::Secrets => "secrets",
            Table::Meta => "meta",
            Table::Branches => "branches",
            Table::Mirror => "mirror",
        }
    }

    /// The table named `name`.
    pub fn from_name(name: &str) -> Option<Table> {
        Table::ALL.into_iter().find(|t| t.name() == name)
    }

    fn def(self) -> TableDefinition<'static, &'static str, &'static [u8]> {
        TableDefinition::new(self.name())
    }
}

/// Where a job is.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub enum JobState {
    /// Recorded, not started (or interrupted by a restart: run again — every
    /// job is convergent).
    Queued,
    /// Running now.
    Running,
    /// A `begin*` job waiting for the person.
    Waiting,
    /// Done; `result` holds the result payload.
    Finished,
}

/// One job, as the ledger holds it.
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct JobRecord {
    /// The VTC's id for it.
    pub job_id: String,
    /// SHA-256 of the payload's canonical JSON: a repeat with other content
    /// is `jobIdReused`.
    pub digest: String,
    /// The namespace it acts in.
    pub namespace: String,
    /// Its `kind`.
    pub kind: String,
    /// The payload, while the job may still run. Cleared once finished — the
    /// VTC is the source of truth and resends the full desired state (spec:
    /// the bridge SHOULD NOT keep `desiredRoles` beyond the result).
    pub payload: Option<Value>,
    /// Where it is.
    pub state: JobState,
    /// For `begin*` jobs: the `next` the job response carried.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub next: Option<Value>,
    /// The result payload, once finished.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub result: Option<Value>,
    /// Unix seconds.
    pub received_at: i64,
    /// The job document's `issuedAt`, in Unix seconds (absent in records
    /// written before it was kept).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub issued_at: Option<i64>,
    /// Unix seconds, once finished.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub finished_at: Option<i64>,
}

impl JobRecord {
    /// A job just received, queued to run.
    pub fn queued(
        job_id: impl Into<String>,
        digest: impl Into<String>,
        namespace: impl Into<String>,
        kind: impl Into<String>,
        payload: Value,
        received_at: i64,
    ) -> Self {
        JobRecord {
            job_id: job_id.into(),
            digest: digest.into(),
            namespace: namespace.into(),
            kind: kind.into(),
            payload: Some(payload),
            state: JobState::Queued,
            next: None,
            result: None,
            received_at,
            issued_at: None,
            finished_at: None,
        }
    }
}

/// Whether a namespace's binding has completed.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub enum NamespaceState {
    /// A `beginBind` is waiting for the admin.
    Pending,
    /// Bound: the adapter acts on it.
    Bound,
}

/// The required-workflow pin, as GitHub's adapter hands it out.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct PinRecord {
    /// `<org>/.vgi`'s forge id.
    pub repository_id: u64,
    /// The pinned commit.
    pub sha: String,
    /// The check name.
    pub check: String,
}

/// One namespace the VTC bound through this bridge.
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct NamespaceRecord {
    /// The VTC's namespace id.
    pub id: String,
    /// `host/owner`.
    pub resource: Resource,
    /// Pending or bound.
    pub state: NamespaceState,
    /// The completed binding (owner id, kind, installation).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub binding: Option<NamespaceBinding>,
    /// Capabilities found while binding, updated by `CapabilityChanged`.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub capabilities: Option<Capabilities>,
    /// GitHub: whether org rulesets (a required workflow) are available.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub required_workflow: Option<bool>,
    /// GitHub: the required-workflow pin.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub pin: Option<PinRecord>,
    /// GitHub: whether the installation carries the bridge-posted check
    /// (its permissions and event subscriptions). `None`: not probed yet.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub bridge_checks: Option<bool>,
    /// Forge ids of the repositories the bridge manages here.
    #[serde(default)]
    pub managed: BTreeSet<u64>,
}

impl NamespaceRecord {
    /// A namespace whose bind has started.
    pub fn pending(id: impl Into<String>, resource: Resource) -> Self {
        NamespaceRecord {
            id: id.into(),
            resource,
            state: NamespaceState::Pending,
            binding: None,
            capabilities: None,
            required_workflow: None,
            pin: None,
            bridge_checks: None,
            managed: BTreeSet::new(),
        }
    }
}

/// One repository the bridge manages.
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct RepoRecord {
    /// The VTC namespace id.
    pub namespace: String,
    /// Where it is now.
    pub resource: Resource,
    /// The forge's id — what everything is keyed on.
    pub forge_id: u64,
    /// Owners with linked accounts (for the owner-review guard and the
    /// projection).
    #[serde(default)]
    pub owners: Vec<ForgeAccount>,
    /// The forge roles the bridge last projected: the roles it manages.
    #[serde(default)]
    pub roles: Vec<RoleAssignment>,
    /// Whether roles have been projected at all (an empty set is a real
    /// projection; `false` means "unknown, do not report role drift").
    #[serde(default)]
    pub roles_known: bool,
    /// The required check, once bootstrapped.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub required_check: Option<String>,
    /// Archived through the bridge.
    #[serde(default)]
    pub archived: bool,
    /// Digest of the drift last reported, so a sweep that finds the same
    /// drift again sends nothing.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub last_drift: Option<String>,
    /// The check-source guard the last inspection found in force, in the
    /// VTC's words ([`crate::status::Guard`]). Reported, never acted on.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub guard: Option<String>,
    /// The last check the bridge posted on the repository itself.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub last_check: Option<LastCheck>,
    /// The role map, as the forge applies it, under which `roles` were last
    /// projected. `None` on a record written before the bridge kept it,
    /// which is taken to be the default map (`crate::rolemap`): the only
    /// map a released bridge applied until then.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub role_map: Option<vgi_forge::RoleMap>,
}

/// A check the bridge posted (GitHub fallback mode), as it reports it.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct LastCheck {
    /// The commit the check run is on.
    pub sha: String,
    /// `success` or `failure`.
    pub conclusion: String,
    /// Unix seconds, when it was completed.
    pub at: i64,
}

impl LastCheck {
    /// A check completed on `sha` at `at`.
    pub fn new(sha: impl Into<String>, conclusion: impl Into<String>, at: i64) -> Self {
        LastCheck {
            sha: sha.into(),
            conclusion: conclusion.into(),
            at,
        }
    }
}

impl RepoRecord {
    /// A record for a repository first seen now.
    pub fn new(namespace: impl Into<String>, resource: Resource, forge_id: u64) -> Self {
        RepoRecord {
            namespace: namespace.into(),
            resource,
            forge_id,
            owners: Vec::new(),
            roles: Vec::new(),
            roles_known: false,
            required_check: None,
            archived: false,
            last_drift: None,
            guard: None,
            last_check: None,
            role_map: None,
        }
    }

    /// The store key: `<host>#<forge id>`.
    pub fn key(&self) -> String {
        repo_key(self.resource.host(), self.forge_id)
    }
}

/// The key a repository is stored under.
pub fn repo_key(host: &str, forge_id: u64) -> String {
    format!("{host}#{forge_id}")
}

/// One verified `push` to a `dependabot/*` branch, as GitHub reported it.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct PushRecord {
    /// The branch's value before the push (all zeros when it was created).
    pub before: String,
    /// Its value after.
    pub after: String,
    /// Who GitHub says pushed: login…
    pub sender_login: String,
    /// …and numeric id.
    pub sender_id: u64,
    /// The push created the branch.
    pub created: bool,
    /// GitHub's delivery id.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub delivery_id: Option<String>,
    /// Unix seconds, when the bridge recorded it.
    pub at: i64,
}

/// A push the bridge made itself (a re-sign), recorded before it was sent.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct OwnPush {
    /// The head it replaced (the force-push's lease).
    pub before: String,
    /// The re-signed head.
    pub after: String,
    /// Unix seconds.
    pub at: i64,
}

/// The provenance ledger of one `dependabot/*` branch (§9): every push the
/// bridge has seen to it since it was created, and the bridge's own pushes.
/// A branch's deletion clears it.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct BranchLedger {
    /// `host/owner/repo`, as last seen.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub repo: Option<Resource>,
    /// The branch.
    #[serde(default)]
    pub branch: String,
    /// Pushes, in the order they arrived (GitHub does not promise delivery
    /// order; the chain is rebuilt from `before`/`after`).
    #[serde(default)]
    pub pushes: Vec<PushRecord>,
    /// The bridge's own re-sign pushes.
    #[serde(default)]
    pub own: Vec<OwnPush>,
    /// More pushes arrived than the ledger keeps: the branch is never clean
    /// again (until deleted).
    #[serde(default)]
    pub overflow: bool,
    /// The open pull request from this branch, once one was seen — so a push
    /// that arrives after the pull request's delivery can resume the re-sign.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub pull_request: Option<u64>,
    /// Unix seconds, when anything last changed here (the oldest untouched
    /// ledger of a repository is evicted first).
    #[serde(default)]
    pub touched: i64,
    /// The newest `repository.pushed_at` of a delivery recorded here: an
    /// older delivery may add a record but never resets or clears the
    /// ledger.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub last_pushed_at: Option<i64>,
    /// Why the last re-sign of a head stopped at its commits (the head, and
    /// the reason), for the check's summary.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub last_skip: Option<(String, String)>,
}

/// A flow waiting for a person, by its `state`.
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", tag = "type")]
#[non_exhaustive]
pub enum PendingFlow {
    /// A namespace bind.
    Bind {
        /// The `beginBind` job.
        job_id: String,
        /// The VTC namespace id.
        namespace: String,
        /// `host/owner` being bound.
        resource: Resource,
        /// Unix seconds.
        expires_at: i64,
    },
    /// An account link.
    Link {
        /// The `beginAccountLink` job.
        job_id: String,
        /// The VTC namespace id.
        namespace: String,
        /// The forge host.
        host: String,
        /// The member's DID, for the adapter's member-bound `state`.
        member: String,
        /// Unix seconds.
        expires_at: i64,
        /// A device flow is being polled (the device code is a sealed
        /// secret, `pending/<state>/device`).
        #[serde(default)]
        device: Option<DevicePoll>,
    },
    /// A GitHub App registration through the manifest flow.
    Manifest {
        /// The forge host.
        host: String,
        /// The org the App is registered under, or `None` for the account
        /// of whoever registers it.
        owner: Option<String>,
        /// Unix seconds.
        expires_at: i64,
    },
}

impl PendingFlow {
    /// When it lapses.
    pub fn expires_at(&self) -> i64 {
        match self {
            PendingFlow::Bind { expires_at, .. }
            | PendingFlow::Link { expires_at, .. }
            | PendingFlow::Manifest { expires_at, .. } => *expires_at,
        }
    }

    /// The job it belongs to, if any.
    pub fn job_id(&self) -> Option<&str> {
        match self {
            PendingFlow::Bind { job_id, .. } | PendingFlow::Link { job_id, .. } => Some(job_id),
            PendingFlow::Manifest { .. } => None,
        }
    }
}

/// How a device-flow link is polled.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct DevicePoll {
    /// Seconds between polls.
    pub interval: u64,
    /// Seconds the code lives.
    pub expires_in: u64,
}

/// A result or event the VTC has not acknowledged yet.
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct OutboxEntry {
    /// `result` or `event`.
    pub kind: OutboxKind,
    /// The payload; a fresh, signed document is built around it on every
    /// send (a stale `issuedAt` would be refused).
    pub payload: Value,
    /// Ids of the documents sent so far (the most recent last), to match the
    /// VTC's response by `threadId`.
    #[serde(default)]
    pub doc_ids: Vec<String>,
    /// Unix seconds of the last send.
    pub last_sent: i64,
    /// Sends so far.
    pub attempts: u32,
}

impl OutboxEntry {
    /// A result not sent yet.
    pub fn result(payload: Value) -> Self {
        OutboxEntry {
            kind: OutboxKind::Result,
            payload,
            doc_ids: Vec::new(),
            last_sent: 0,
            attempts: 0,
        }
    }
}

/// What an [`OutboxEntry`] carries.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub enum OutboxKind {
    /// A `git-ns/bridge/result`.
    Result,
    /// A `git-ns/bridge/event`.
    Event,
}

/// The store.
///
/// In VTA mode ([`Store::with_mirror`]) it is a cache of the state held in
/// the VTA's app-state ([`crate::appstate`]): writes to the mirrored tables
/// are marked for the mirror task, and secrets are held in memory and in the
/// VTA only — never sealed into the file.
#[derive(Clone)]
pub struct Store {
    db: Arc<Database>,
    key: Arc<MasterKey>,
    mirror: Option<Arc<crate::appstate::Mirror>>,
}

impl std::fmt::Debug for Store {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.debug_struct("Store").finish_non_exhaustive()
    }
}

impl Store {
    /// Open (or create) the store at `path`, sealing secrets under `key`.
    /// redb takes an exclusive lock: one bridge process per store.
    pub fn open(path: &Path, key: MasterKey) -> Result<Self> {
        let db = Database::create(path).with_context(|| {
            format!(
                "opening the state store {} (is another bridge running on it?)",
                path.display()
            )
        })?;
        Self::with_db(db, key)
    }

    /// An in-memory store (tests).
    pub fn in_memory(key: MasterKey) -> Result<Self> {
        let db = Database::builder().create_with_backend(redb::backends::InMemoryBackend::new())?;
        Self::with_db(db, key)
    }

    fn with_db(db: Database, key: MasterKey) -> Result<Self> {
        let w = db.begin_write()?;
        for t in Table::ALL {
            w.open_table(t.def())?;
        }
        w.commit()?;
        Ok(Store {
            db: Arc::new(db),
            key: Arc::new(key),
            mirror: None,
        })
    }

    /// This store as the local cache of VTA mode: mirrored writes are marked
    /// on `mirror`, and secrets live in it (memory) rather than in the file.
    pub fn with_mirror(mut self, mirror: Arc<crate::appstate::Mirror>) -> Self {
        self.mirror = Some(mirror);
        self
    }

    /// The VTA mirror, in VTA mode.
    pub fn mirror(&self) -> Option<&Arc<crate::appstate::Mirror>> {
        self.mirror.as_ref()
    }

    /// Wait until every change is in the VTA (at most `timeout`). Always
    /// `true` outside VTA mode, where the file is the durable copy.
    pub async fn flush(&self, timeout: std::time::Duration) -> bool {
        match &self.mirror {
            Some(m) => m.flush(timeout).await,
            None => true,
        }
    }

    fn touched(&self, table: Table, key: &str) {
        if let Some(m) = &self.mirror
            && crate::appstate::record_is_mirrored(table, key)
        {
            m.mark(crate::appstate::Dirty::Record(table, key.to_string()));
        }
    }

    /// Write a record pulled from the VTA, without marking it for the
    /// mirror.
    pub(crate) fn put_cached<T: Serialize>(
        &self,
        table: Table,
        key: &str,
        value: &T,
    ) -> Result<()> {
        let bytes = serde_json::to_vec(value)?;
        let w = self.db.begin_write()?;
        w.open_table(table.def())?.insert(key, bytes.as_slice())?;
        w.commit()?;
        Ok(())
    }

    /// Delete a record without marking it for the mirror (what the VTA
    /// already dropped).
    pub(crate) fn delete_cached(&self, table: Table, key: &str) -> Result<()> {
        let w = self.db.begin_write()?;
        w.open_table(table.def())?.remove(key)?;
        w.commit()?;
        Ok(())
    }

    /// One record as JSON.
    pub(crate) fn get_raw(&self, table: Table, key: &str) -> Result<Option<Value>> {
        self.get(table, key)
    }

    /// Every key in a table.
    pub(crate) fn keys(&self, table: Table) -> Result<Vec<String>> {
        let r = self.db.begin_read()?;
        let t = r.open_table(table.def())?;
        let mut out = Vec::new();
        for row in t.iter()? {
            out.push(row?.0.value().to_string());
        }
        Ok(out)
    }

    /// Whether the file holds any record of `table`.
    pub fn is_empty(&self, table: Table) -> Result<bool> {
        let r = self.db.begin_read()?;
        let t = r.open_table(table.def())?;
        Ok(t.iter()?.next().is_none())
    }

    #[cfg(test)]
    pub(crate) fn raw_secret_rows(&self) -> Result<Vec<String>> {
        self.keys(Table::Secrets)
    }

    /// Read one record.
    pub fn get<T: DeserializeOwned>(&self, table: Table, key: &str) -> Result<Option<T>> {
        let r = self.db.begin_read()?;
        let t = r.open_table(table.def())?;
        match t.get(key)? {
            Some(v) => Ok(Some(
                serde_json::from_slice(v.value())
                    .with_context(|| format!("decoding {table:?}/{key}"))?,
            )),
            None => Ok(None),
        }
    }

    /// Write one record.
    pub fn put<T: Serialize>(&self, table: Table, key: &str, value: &T) -> Result<()> {
        let bytes = serde_json::to_vec(value)?;
        let w = self.db.begin_write()?;
        w.open_table(table.def())?.insert(key, bytes.as_slice())?;
        w.commit()?;
        self.touched(table, key);
        Ok(())
    }

    /// Write one record only if the key is free. `false` if it was taken.
    pub fn put_new<T: Serialize>(&self, table: Table, key: &str, value: &T) -> Result<bool> {
        let bytes = serde_json::to_vec(value)?;
        let w = self.db.begin_write()?;
        {
            let mut t = w.open_table(table.def())?;
            if t.get(key)?.is_some() {
                return Ok(false);
            }
            t.insert(key, bytes.as_slice())?;
        }
        w.commit()?;
        self.touched(table, key);
        Ok(true)
    }

    /// Read-modify-write one record in a single transaction. `f` returns the
    /// new value (`None` deletes) and what to hand back.
    pub fn update<T, R>(
        &self,
        table: Table,
        key: &str,
        f: impl FnOnce(Option<T>) -> Result<(Option<T>, R)>,
    ) -> Result<R>
    where
        T: Serialize + DeserializeOwned,
    {
        let w = self.db.begin_write()?;
        let out = {
            let mut t = w.open_table(table.def())?;
            let current = match t.get(key)? {
                Some(v) => Some(serde_json::from_slice(v.value())?),
                None => None,
            };
            let (next, out) = f(current)?;
            match next {
                Some(v) => {
                    let bytes = serde_json::to_vec(&v)?;
                    t.insert(key, bytes.as_slice())?;
                }
                None => {
                    t.remove(key)?;
                }
            }
            out
        };
        w.commit()?;
        self.touched(table, key);
        Ok(out)
    }

    /// Close job `job_id` with `result` and queue it for sending, in **one**
    /// transaction: a crash can leave the job either running (it runs again)
    /// or finished with its result queued — never finished with the result
    /// lost. `false` (and nothing written) if the job is unknown or already
    /// finished.
    pub fn finish_job(&self, job_id: &str, result: &Value, finished_at: i64) -> Result<bool> {
        let w = self.db.begin_write()?;
        {
            let mut jobs = w.open_table(Table::Jobs.def())?;
            let Some(current) = jobs.get(job_id)? else {
                return Ok(false);
            };
            let mut rec: JobRecord = serde_json::from_slice(current.value())?;
            drop(current);
            if rec.state == JobState::Finished {
                return Ok(false);
            }
            rec.state = JobState::Finished;
            rec.result = Some(result.clone());
            rec.payload = None;
            rec.finished_at = Some(finished_at);
            let bytes = serde_json::to_vec(&rec)?;
            jobs.insert(job_id, bytes.as_slice())?;
            let entry = OutboxEntry::result(result.clone());
            let bytes = serde_json::to_vec(&entry)?;
            w.open_table(Table::Outbox.def())?
                .insert(format!("result:{job_id}").as_str(), bytes.as_slice())?;
        }
        w.commit()?;
        Ok(true)
    }

    /// Delete one record. `true` if it existed.
    pub fn delete(&self, table: Table, key: &str) -> Result<bool> {
        let w = self.db.begin_write()?;
        let existed = w.open_table(table.def())?.remove(key)?.is_some();
        w.commit()?;
        if existed {
            self.touched(table, key);
        }
        Ok(existed)
    }

    /// Every record in a table, by key.
    pub fn list<T: DeserializeOwned>(&self, table: Table) -> Result<Vec<(String, T)>> {
        let r = self.db.begin_read()?;
        let t = r.open_table(table.def())?;
        let mut out = Vec::new();
        for row in t.iter()? {
            let (k, v) = row?;
            let key = k.value().to_string();
            let value = serde_json::from_slice(v.value())
                .with_context(|| format!("decoding {table:?}/{key}"))?;
            out.push((key, value));
        }
        Ok(out)
    }

    /// Seal `value` under `name` and store it.
    pub fn put_secret(&self, name: &str, value: &[u8]) -> Result<()> {
        if let Some(m) = &self.mirror {
            m.secrets
                .lock()
                .expect("lock")
                .insert(name.to_string(), Zeroizing::new(value.to_vec()));
            if crate::appstate::secret_is_mirrored(name) {
                m.mark(crate::appstate::Dirty::Secret(name.to_string()));
            }
            return Ok(());
        }
        let sealed = self.key.seal(name, value)?;
        let w = self.db.begin_write()?;
        w.open_table(Table::Secrets.def())?
            .insert(name, sealed.as_slice())?;
        w.commit()?;
        Ok(())
    }

    /// Open the secret `name`, if stored.
    pub fn get_secret(&self, name: &str) -> Result<Option<Zeroizing<Vec<u8>>>> {
        if let Some(m) = &self.mirror {
            return Ok(m.secrets.lock().expect("lock").get(name).cloned());
        }
        let r = self.db.begin_read()?;
        let t = r.open_table(Table::Secrets.def())?;
        match t.get(name)? {
            Some(v) => Ok(Some(self.key.open(name, v.value())?)),
            None => Ok(None),
        }
    }

    /// The secret `name` as UTF-8, if stored.
    pub fn get_secret_string(&self, name: &str) -> Result<Option<Zeroizing<String>>> {
        match self.get_secret(name)? {
            Some(bytes) => Ok(Some(Zeroizing::new(
                String::from_utf8(bytes.to_vec())
                    .map_err(|_| anyhow::anyhow!("secret `{name}` is not UTF-8"))?,
            ))),
            None => Ok(None),
        }
    }

    /// Remove the secret `name`.
    pub fn delete_secret(&self, name: &str) -> Result<()> {
        if let Some(m) = &self.mirror {
            let existed = m.secrets.lock().expect("lock").remove(name).is_some();
            if existed && crate::appstate::secret_is_mirrored(name) {
                m.mark(crate::appstate::Dirty::Secret(name.to_string()));
            }
            return Ok(());
        }
        self.delete(Table::Secrets, name).map(|_| ())
    }

    /// Names of the stored secrets (never their values).
    pub fn secret_names(&self) -> Result<Vec<String>> {
        if let Some(m) = &self.mirror {
            return Ok(m.secrets.lock().expect("lock").keys().cloned().collect());
        }
        let r = self.db.begin_read()?;
        let t = r.open_table(Table::Secrets.def())?;
        let mut out = Vec::new();
        for row in t.iter()? {
            out.push(row?.0.value().to_string());
        }
        Ok(out)
    }
}

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

    fn store() -> Store {
        Store::in_memory(MasterKey::generate().unwrap()).unwrap()
    }

    #[test]
    fn records_round_trip_and_update_is_atomic() {
        let s = store();
        let ns = NamespaceRecord::pending("ns_1", Resource::parse("github.com/acme").unwrap());
        s.put(Table::Namespaces, "ns_1", &ns).unwrap();
        let back: NamespaceRecord = s.get(Table::Namespaces, "ns_1").unwrap().unwrap();
        assert_eq!(back.state, NamespaceState::Pending);
        let n = s
            .update::<NamespaceRecord, _>(Table::Namespaces, "ns_1", |r| {
                let mut r = r.unwrap();
                r.managed.insert(9);
                let n = r.managed.len();
                Ok((Some(r), n))
            })
            .unwrap();
        assert_eq!(n, 1);
        assert!(!s.put_new(Table::Namespaces, "ns_1", &back).unwrap());
        assert_eq!(
            s.list::<NamespaceRecord>(Table::Namespaces).unwrap().len(),
            1
        );
        assert!(s.delete(Table::Namespaces, "ns_1").unwrap());
        assert!(
            s.get::<NamespaceRecord>(Table::Namespaces, "ns_1")
                .unwrap()
                .is_none()
        );
    }

    #[test]
    fn finishing_a_job_records_the_result_and_queues_it_together() {
        let s = store();
        let job = JobRecord::queued("j1", "d", "ns", "inspect", serde_json::json!({}), 0);
        s.put(Table::Jobs, "j1", &job).unwrap();
        let result = serde_json::json!({ "jobId": "j1", "outcome": "succeeded" });
        assert!(s.finish_job("j1", &result, 5).unwrap());
        let rec: JobRecord = s.get(Table::Jobs, "j1").unwrap().unwrap();
        assert_eq!(rec.state, JobState::Finished);
        assert_eq!(rec.result.as_ref(), Some(&result));
        assert!(rec.payload.is_none());
        let queued: OutboxEntry = s.get(Table::Outbox, "result:j1").unwrap().unwrap();
        assert_eq!(queued.payload, result);
        // Exactly once: a second close changes nothing, even after the
        // entry was acknowledged.
        s.delete(Table::Outbox, "result:j1").unwrap();
        assert!(!s.finish_job("j1", &serde_json::json!({}), 6).unwrap());
        assert!(
            s.get::<OutboxEntry>(Table::Outbox, "result:j1")
                .unwrap()
                .is_none()
        );
        assert!(!s.finish_job("unknown", &result, 6).unwrap());
    }

    #[test]
    fn secrets_are_stored_sealed() {
        let s = store();
        s.put_secret("github/github.com/app", b"pem-bytes").unwrap();
        assert_eq!(
            &*s.get_secret("github/github.com/app").unwrap().unwrap(),
            b"pem-bytes"
        );
        let raw =
            s.db.begin_read()
                .unwrap()
                .open_table(Table::Secrets.def())
                .unwrap()
                .get("github/github.com/app")
                .unwrap()
                .unwrap()
                .value()
                .to_vec();
        assert!(
            !raw.windows(9).any(|w| w == b"pem-bytes"),
            "ciphertext only"
        );
        assert_eq!(s.secret_names().unwrap(), ["github/github.com/app"]);
    }

    #[test]
    fn a_file_store_survives_reopening() {
        let dir = tempfile::tempdir().unwrap();
        let path = dir.path().join("state.redb");
        let key = MasterKey::generate().unwrap();
        let text = key.to_text();
        {
            let s = Store::open(&path, key).unwrap();
            s.put_secret("x", b"y").unwrap();
            s.put(Table::Deliveries, "d-1", &1_i64).unwrap();
        }
        let s = Store::open(&path, MasterKey::from_text(&text).unwrap()).unwrap();
        assert_eq!(&*s.get_secret("x").unwrap().unwrap(), b"y");
        assert_eq!(s.get::<i64>(Table::Deliveries, "d-1").unwrap(), Some(1));
    }
}