1- use std:: collections:: HashSet ;
1+ use std:: collections:: { HashMap , HashSet } ;
22use std:: path:: PathBuf ;
33use std:: sync:: Arc ;
44use tokio:: sync:: watch;
@@ -13,8 +13,9 @@ const SWEEP_INTERVAL_SECS: u64 = 3600;
1313/// the O(repos) amplification the admission-control work exists to prevent.
1414const REPOS_PER_PASS : usize = 100 ;
1515
16- /// Maximum objects to pin per repo in a single pass — prevents one large
17- /// repo from monopolizing the blocking pool or the hourly budget.
16+ /// Maximum objects to pin per backend per repo in a single pass — prevents one
17+ /// large repo from monopolizing the blocking pool or the hourly budget. Applied
18+ /// after filtering out already-pinned objects so the cap reflects actual work.
1819const MAX_OBJECTS_PER_REPO : usize = 50_000 ;
1920
2021/// Spawn the periodic reconciliation sweep background task.
@@ -30,7 +31,6 @@ pub fn spawn(
3031 let node_seed = * node_keypair. to_seed ( ) ;
3132 let mut cursor = 0usize ;
3233
33- // Run the first pass immediately on startup, then periodically.
3434 loop {
3535 let start = std:: time:: Instant :: now ( ) ;
3636 match run_pass (
@@ -40,6 +40,7 @@ pub fn spawn(
4040 & node_seed,
4141 & node_did,
4242 & mut cursor,
43+ & mut shutdown_rx,
4344 )
4445 . await
4546 {
@@ -57,6 +58,12 @@ pub fn spawn(
5758 }
5859 }
5960
61+ // Check shutdown before sleeping.
62+ if * shutdown_rx. borrow ( ) {
63+ tracing:: info!( "reconciliation sweep: shutdown signal received, exiting" ) ;
64+ return ;
65+ }
66+
6067 tokio:: select! {
6168 _ = tokio:: time:: sleep( std:: time:: Duration :: from_secs( SWEEP_INTERVAL_SECS ) ) => { }
6269 _ = shutdown_rx. changed( ) => {
@@ -78,9 +85,8 @@ async fn run_pass(
7885 node_seed : & [ u8 ; 32 ] ,
7986 node_did : & gitlawb_core:: did:: Did ,
8087 cursor : & mut usize ,
88+ shutdown_rx : & mut watch:: Receiver < bool > ,
8189) -> anyhow:: Result < ( usize , usize , usize ) > {
82- // Use the canonical/deduplicated listing so mirror rows never bypass
83- // visibility rules. The dedup CTE also excludes quarantined repos.
8490 let all = db. list_all_repos_deduped ( ) . await ?;
8591
8692 if all. is_empty ( ) {
@@ -98,6 +104,12 @@ async fn run_pass(
98104 let mut total_gaps_filled = 0usize ;
99105
100106 for repo in batch {
107+ // Cooperative shutdown: exit between repos if signal received.
108+ if * shutdown_rx. borrow ( ) {
109+ tracing:: info!( "reconciliation sweep: shutdown signal received mid-pass, exiting" ) ;
110+ break ;
111+ }
112+
101113 let repo_slug = format ! (
102114 "{}/{}" ,
103115 crate :: db:: normalize_owner_key( & repo. owner_did) ,
@@ -128,10 +140,6 @@ async fn run_pass(
128140 let is_public = repo. is_public ;
129141 let object_list = tokio:: task:: spawn_blocking ( move || -> anyhow:: Result < Vec < String > > {
130142 let all_objs = crate :: git:: push_delta:: list_all_objects ( & disk_clone) ?;
131- // Always compute the reachable, visibility-allowed blob set so
132- // ordinary public repos recover their blobs too (not just commits
133- // and trees). replicable_blob_set handles the no-rule case
134- // correctly — all reachable blobs are allowed.
135143 let allowed = crate :: git:: visibility_pack:: replicable_blob_set (
136144 & disk_clone,
137145 & rules_clone,
@@ -145,7 +153,7 @@ async fn run_pass(
145153 } )
146154 . await ;
147155
148- let mut object_list = match object_list {
156+ let object_list = match object_list {
149157 Ok ( Ok ( list) ) => list,
150158 Ok ( Err ( e) ) => {
151159 tracing:: warn!( repo = %repo_slug, err = %e, "full-scan failed, skipping" ) ;
@@ -161,45 +169,94 @@ async fn run_pass(
161169 continue ;
162170 }
163171
164- // Enforce per-repo object cap so a huge repo cannot monopolize the
165- // blocking pool or run past the hourly interval.
166- let truncated = object_list. len ( ) > MAX_OBJECTS_PER_REPO ;
167- if truncated {
168- object_list. truncate ( MAX_OBJECTS_PER_REPO ) ;
172+ // Pre-cap the object list before batch-filtering to keep queries bounded.
173+ let candidates: Vec < String > = if object_list. len ( ) > MAX_OBJECTS_PER_REPO {
169174 tracing:: warn!(
170175 repo = %repo_slug,
171176 cap = MAX_OBJECTS_PER_REPO ,
172- "reconciliation per-repo object cap reached, truncating"
177+ total = object_list. len( ) ,
178+ "reconciliation per-repo candidate list truncated to cap"
173179 ) ;
174- }
175-
176- let has_path_scoped = crate :: git:: visibility_pack:: has_path_scoped_rule ( & rules) ;
180+ object_list. into_iter ( ) . take ( MAX_OBJECTS_PER_REPO ) . collect ( )
181+ } else {
182+ object_list
183+ } ;
177184
178185 // ── Phase 1: Public-object pinning (IPFS + Pinata) ────────────────
179186 // Each backend independently tracks its own completion state, so we
180- // pass the full replicable set to both. A row in pinned_cids from
181- // IPFS does not imply a Pinata upload succeeded, and vice versa.
187+ // compute the actually-missing set per backend and cap independently.
188+
189+ // Recheck quarantine before attempting any external pinning.
190+ match db. is_repo_quarantined ( & repo. id ) . await {
191+ Ok ( true ) => {
192+ tracing:: warn!( repo = %repo_slug, "repo quarantined, skipping public-object pinning" ) ;
193+ // Phase 2 (encrypted) is also skipped — a quarantined repo's
194+ // withheld blobs should not be published either.
195+ continue ;
196+ }
197+ Ok ( false ) => { }
198+ Err ( e) => {
199+ tracing:: warn!( repo = %repo_slug, err = %e, "quarantine check failed, skipping" ) ;
200+ continue ;
201+ }
202+ }
203+
204+ // Compute IPFS-missing set, capped per-repo.
205+ let already_ipfs = db. filter_ipfs_pinned_oids ( & candidates) . await ?;
206+ let ipfs_missing_set: HashSet < & str > = candidates
207+ . iter ( )
208+ . map ( |s| s. as_str ( ) )
209+ . collect :: < HashSet < _ > > ( )
210+ . difference ( & already_ipfs. iter ( ) . map ( |s| s. as_str ( ) ) . collect ( ) )
211+ . copied ( )
212+ . collect ( ) ;
213+ let mut ipfs_candidates: Vec < String > =
214+ ipfs_missing_set. into_iter ( ) . map ( String :: from) . collect ( ) ;
215+ if ipfs_candidates. len ( ) > MAX_OBJECTS_PER_REPO {
216+ ipfs_candidates. truncate ( MAX_OBJECTS_PER_REPO ) ;
217+ tracing:: warn!(
218+ repo = %repo_slug,
219+ cap = MAX_OBJECTS_PER_REPO ,
220+ "IPFS per-repo missing cap reached, truncating"
221+ ) ;
222+ }
223+
224+ // Compute Pinata-missing set, capped per-repo.
225+ let already_pinata = db. filter_pinata_pinned_oids ( & candidates) . await ?;
226+ let pinata_missing_set: HashSet < & str > = candidates
227+ . iter ( )
228+ . map ( |s| s. as_str ( ) )
229+ . collect :: < HashSet < _ > > ( )
230+ . difference ( & already_pinata. iter ( ) . map ( |s| s. as_str ( ) ) . collect ( ) )
231+ . copied ( )
232+ . collect ( ) ;
233+ let mut pinata_candidates: Vec < String > =
234+ pinata_missing_set. into_iter ( ) . map ( String :: from) . collect ( ) ;
235+ if pinata_candidates. len ( ) > MAX_OBJECTS_PER_REPO {
236+ pinata_candidates. truncate ( MAX_OBJECTS_PER_REPO ) ;
237+ tracing:: warn!(
238+ repo = %repo_slug,
239+ cap = MAX_OBJECTS_PER_REPO ,
240+ "Pinata per-repo missing cap reached, truncating"
241+ ) ;
242+ }
243+
182244 let pinned_ipfs =
183- crate :: ipfs_pin:: pin_new_objects ( & config. ipfs_api , & disk, object_list. clone ( ) , db)
184- . await ;
245+ crate :: ipfs_pin:: pin_new_objects ( & config. ipfs_api , & disk, ipfs_candidates, db) . await ;
185246
186247 let pinned_pinata = crate :: pinata:: pin_new_objects (
187248 http_client,
188249 & config. pinata_upload_url ,
189250 & config. pinata_jwt ,
190251 & disk,
191- object_list ,
252+ pinata_candidates ,
192253 db,
193254 )
194255 . await ;
195256
196257 let repo_filled = pinned_ipfs. len ( ) + pinned_pinata. len ( ) ;
197258 if repo_filled > 0 {
198259 total_gaps_filled += repo_filled;
199- // Approximate gaps count for observability — objects that were
200- // missing from at least one backend.
201- let missing_ipfs = pinned_ipfs. len ( ) ;
202- let missing_pinata = pinned_pinata. len ( ) ;
203260 let deduped = pinned_ipfs
204261 . iter ( )
205262 . chain ( & pinned_pinata)
@@ -211,8 +268,8 @@ async fn run_pass(
211268
212269 tracing:: info!(
213270 repo = %repo_slug,
214- ipfs = missing_ipfs ,
215- pinata = missing_pinata ,
271+ ipfs = pinned_ipfs . len ( ) ,
272+ pinata = pinned_pinata . len ( ) ,
216273 total = repo_filled,
217274 "reconciliation sweep filled public-object gaps"
218275 ) ;
@@ -221,6 +278,21 @@ async fn run_pass(
221278 // ── Phase 2: Encrypted recovery-copy resealing (withheld blobs) ──
222279 // Only relevant when path-scoped visibility rules exist — without them
223280 // no blobs are withheld and withheld_blob_recipients returns empty.
281+
282+ // Recheck quarantine before encrypted pinning.
283+ let quarantined = match db. is_repo_quarantined ( & repo. id ) . await {
284+ Ok ( q) => q,
285+ Err ( e) => {
286+ tracing:: warn!( repo = %repo_slug, err = %e, "quarantine recheck failed, skipping encrypted pin" ) ;
287+ continue ;
288+ }
289+ } ;
290+ if quarantined {
291+ tracing:: warn!( repo = %repo_slug, "repo quarantined, skipping encrypted pinning" ) ;
292+ continue ;
293+ }
294+
295+ let has_path_scoped = crate :: git:: visibility_pack:: has_path_scoped_rule ( & rules) ;
224296 if has_path_scoped && !config. ipfs_api . is_empty ( ) {
225297 let p = disk. clone ( ) ;
226298 let owner = repo. owner_did . clone ( ) ;
@@ -242,18 +314,35 @@ async fn run_pass(
242314 & rec,
243315 )
244316 . await ;
245- if !sealed. is_empty ( ) && !config. irys_url . is_empty ( ) {
317+
318+ // Anchor ALL existing encrypted blobs for this repo, not
319+ // just the ones encrypted this pass. This ensures that if
320+ // a prior manifest anchor failed the retry will include
321+ // previously-encrypted blobs too.
322+ let all_existing = db. list_all_encrypted_blobs ( & repo. id ) . await ?;
323+ if !all_existing. is_empty ( ) && !config. irys_url . is_empty ( ) {
246324 let owner_short = crate :: db:: normalize_owner_key ( & repo. owner_did ) ;
247325 let slug = format ! ( "{}/{}" , owner_short, repo. name) ;
248326 let ts = chrono:: Utc :: now ( ) . to_rfc3339 ( ) ;
249327 let node_did_str = node_did. to_string ( ) ;
250- let blobs: Vec < ( String , String ) > = sealed;
328+
329+ // Merge existing blobs with freshly-sealed ones,
330+ // preferring later entries (newly-sealed) on conflict.
331+ let mut blob_map: HashMap < String , String > = HashMap :: new ( ) ;
332+ for ( oid, cid) in & all_existing {
333+ blob_map. insert ( oid. clone ( ) , cid. clone ( ) ) ;
334+ }
335+ for ( oid, cid) in & sealed {
336+ blob_map. insert ( oid. clone ( ) , cid. clone ( ) ) ;
337+ }
338+ let merged: Vec < ( String , String ) > = blob_map. into_iter ( ) . collect ( ) ;
339+
251340 let manifest = crate :: arweave:: EncryptedManifest {
252341 repo : & slug,
253342 owner_did : & repo. owner_did ,
254343 node_did : & node_did_str,
255344 timestamp : & ts,
256- blobs : & blobs ,
345+ blobs : & merged ,
257346 } ;
258347 if let Err ( e) = crate :: arweave:: anchor_encrypted_manifest (
259348 http_client,
@@ -265,7 +354,7 @@ async fn run_pass(
265354 tracing:: warn!(
266355 repo = %slug,
267356 err = %e,
268- "encrypted manifest anchor failed (best-effort )"
357+ "encrypted manifest anchor failed (will retry next pass )"
269358 ) ;
270359 }
271360 }
0 commit comments