1use super::{
3 LocalStore, Service, acquisition, closure, copy, discovery, intent, inventory, source,
4};
5use crate::indexed::budget::{Budgeted, SliceBudget};
6use crate::pipeline::ShardMap;
7use crate::store::{BlobKey, BlobStore, ContentIndex, StoreError};
8use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
9use crate::{
10 Addressing, Batch, Clock, Cursor, Key, NamespaceStore, Partition, Precondition, RepoId, Value,
11};
12use mkit_core::hash::{Hash, to_hex};
13use serde::{Deserialize, Serialize};
14use std::sync::Arc;
15
16fn bad() -> StoreError {
17 StoreError::Corrupt("invalid preservation work".into())
18}
19fn value<T: Serialize>(v: &T) -> Result<Value, StoreError> {
20 intent::encode(v).map_err(|_| bad())
21}
22fn decode<T: serde::de::DeserializeOwned>(v: &Value) -> Result<T, StoreError> {
23 intent::decode(v).map_err(|_| bad())
24}
25pub(super) fn key(tag: &[u8], action: &Hash, tail: &[u8]) -> Key {
26 Key::new(
27 [
28 b"b\0\xffpreservation\0".as_slice(),
29 tag,
30 b"\0",
31 action,
32 tail,
33 ]
34 .concat(),
35 )
36}
37pub(super) fn range(tag: &[u8], action: &Hash) -> (Key, Key) {
38 let start = key(tag, action, &[]);
39 let end = prefix_end(&start);
40 (start, end)
41}
42fn known_holder_key(id: &Hash, object: &Hash, repo: &RepoId) -> Key {
43 key(
44 b"known-holder",
45 id,
46 &[
47 object.as_slice(),
48 repo.namespace.as_str().as_bytes(),
49 b"\0",
50 repo.name.as_str().as_bytes(),
51 ]
52 .concat(),
53 )
54}
55fn holder_context() -> Result<Value, StoreError> {
56 value(
57 &serde_json::json!({"contextComplete":false,"signerMetadata":"unavailable_in_existing_source"}),
58 )
59}
60fn prefix_end(start: &Key) -> Key {
61 let mut bytes = start.as_bytes().to_vec();
62 while bytes.last() == Some(&255) {
63 bytes.pop();
64 }
65 if let Some(last) = bytes.last_mut() {
66 *last += 1;
67 }
68 Key::new(bytes)
69}
70#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
71pub(super) enum Phase {
72 Seed,
73 Acquire,
74 Closure,
75 Discover,
76 Retain,
77 Purging,
78 Purged,
79}
80#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
81pub(super) enum Verification {
82 CanonicalPending,
83 ManifestClosurePending,
84 SourceCorrupt,
85 Verified,
86}
87#[derive(Debug, Clone, Serialize, Deserialize)]
88#[serde(deny_unknown_fields)]
89pub(super) struct State {
90 pub version: u8,
91 pub phase: Phase,
92 pub retain_until: u64,
93 pub hold: bool,
94 pub purged: bool,
95 pub purge_after: Option<Vec<u8>>,
96 pub resume_phase: Phase,
97 pub next_purge_at: u64,
98 pub verification: Verification,
99 pub discovery_complete: bool,
100 pub seed: usize,
101 pub discovery: Option<discovery::DiscoveryState>,
102 pub current: Option<Hash>,
103 pub verified_objects: u64,
104}
105impl State {
106 pub(super) fn acquisition_complete(&self) -> bool {
107 self.verification == Verification::Verified
108 }
109}
110#[derive(Debug, Clone, Default, Serialize, Deserialize)]
111#[serde(deny_unknown_fields)]
112pub(super) struct ObjectInfo {
113 pub kind: u8,
114 pub size: u64,
115 pub copied: u64,
116 pub chunks: u32,
117 pub verified: bool,
118 pub source_failed: bool,
119 pub holders: Option<Vec<u8>>,
120 pub holders_done: bool,
121 pub namespace_after: Option<Vec<u8>>,
122}
123pub struct Work<N, B, P> {
125 pub metadata: N,
126 pub purge: Option<crate::purge::PurgeConfig>,
128 pub serving: B,
129 pub preserved: P,
130 pub root: Partition,
131 pub shards: Arc<dyn ShardMap>,
132 pub addressing: Addressing,
133 pub retention_ms: u64,
134 pub discovery_margin_ms: u64,
135 pub profile: acquisition::Profile,
136 pub clock: Arc<dyn Clock>,
137}
138impl<N, B, P> std::fmt::Debug for Work<N, B, P> {
139 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
140 f.debug_struct("PreservationWork")
141 .field("root", &self.root)
142 .finish_non_exhaustive()
143 }
144}
145impl<N: NamespaceStore + Clone, B: BlobStore, P: BlobStore> Work<N, B, P> {
146 pub async fn plan_legal_hold<S: NamespaceStore>(
152 &self,
153 store: &S,
154 id: Hash,
155 enabled: bool,
156 now: u64,
157 ) -> Result<Batch, StoreError> {
158 let service = Service::new(
159 LocalStore::new(store, &self.root, store),
160 self.root.clone(),
161 self.shards.clone(),
162 )
163 .with_purge(self.purge.clone());
164 let (record, _) = service
165 .record(store, &id)
166 .await
167 .map_err(|_| bad())?
168 .ok_or_else(bad)?;
169 let (mut state, old) = self.state(store, &id, record.created).await?;
170 if enabled && (state.purged || state.phase == Phase::Purging) {
171 return Err(StoreError::unavailable(
172 "preservation purge already owns request",
173 ));
174 }
175 state.hold = enabled;
176 let state_key = key(b"state", &id, &[]);
177 Ok(Batch::new()
178 .require(old.map_or_else(
179 || Precondition::Absent(state_key.clone()),
180 |old| Precondition::Equals(state_key.clone(), old),
181 ))
182 .require(Precondition::NotAfter(
183 now.saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS),
184 ))
185 .put(state_key, value(&state)?)
186 .put(
187 crate::store::keys::timer(now.saturating_add(1), kinds::TAKEDOWN_WORK.get(), &id),
188 Value::default(),
189 ))
190 }
191 pub(super) async fn state<S: NamespaceStore>(
192 &self,
193 store: &S,
194 id: &Hash,
195 created: u64,
196 ) -> Result<(State, Option<Value>), StoreError> {
197 let raw = store.get(&self.root, &key(b"state", id, &[])).await?;
198 let state = if let Some(raw) = &raw {
199 decode(raw)?
200 } else {
201 State {
202 version: 1,
203 phase: Phase::Seed,
204 retain_until: created
205 .checked_add(self.retention_ms)
206 .filter(|n| *n <= i64::MAX.unsigned_abs())
207 .ok_or_else(bad)?,
208 hold: false,
209 purged: false,
210 purge_after: None,
211 resume_phase: Phase::Seed,
212 next_purge_at: 0,
213 verification: Verification::CanonicalPending,
214 discovery_complete: false,
215 seed: 0,
216 discovery: None,
217 current: None,
218 verified_objects: 0,
219 }
220 };
221 if state.version != 1
222 || self.retention_ms == 0
223 || state.retain_until < created
224 || state.hold && matches!(state.phase, Phase::Purging | Phase::Purged)
225 {
226 return Err(bad());
227 }
228 Ok((state, raw))
229 }
230 pub(super) async fn info<S: NamespaceStore>(
231 &self,
232 store: &S,
233 id: &Hash,
234 object: &Hash,
235 ) -> Result<ObjectInfo, StoreError> {
236 store
237 .get(&self.root, &key(b"object", id, object))
238 .await?
239 .map(|v| decode(&v))
240 .transpose()
241 .map(Option::unwrap_or_default)
242 }
243 async fn plan_cache_purge<S: NamespaceStore>(
244 &self,
245 store: &S,
246 id: &Hash,
247 object: &Hash,
248 repo: &RepoId,
249 now: u64,
250 local_budget: &crate::purge::SliceBudget,
251 ) -> Result<Batch, StoreError> {
252 if self.purge.is_none() {
253 return Ok(Batch::new());
254 }
255 let marker = known_holder_key(id, object, repo);
256 if store.get(&self.root, &marker).await?.is_some() {
257 return Ok(Batch::new());
258 }
259 let operation = format!(
260 "discovery:{}",
261 to_hex(&mkit_core::hash::hash(
262 &[id.as_slice(), object.as_slice()].concat()
263 ))
264 );
265 let batch = crate::purge::automatic::plan_repository(
266 self.purge.as_ref(),
267 store,
268 &self.root,
269 repo,
270 crate::purge::Trigger::Takedown,
271 &operation,
272 now,
273 )
274 .await?;
275 crate::purge::automatic::invalidate_repository(
278 self.purge.as_ref(),
279 &self.root,
280 repo,
281 crate::purge::Trigger::Takedown,
282 &operation,
283 local_budget,
284 )
285 .await;
286 Ok(batch
287 .require(Precondition::Absent(marker.clone()))
288 .put(marker, holder_context()?))
289 }
290 fn enqueue(mut batch: Batch, id: &Hash, object: &Hash, kind: u8) -> Batch {
291 let row = key(b"todo", id, object);
292 batch.writes.retain(|write| match write {
293 crate::Write::Put(key, _) | crate::Write::Delete(key) => key != &row,
294 });
295 batch.put(row, Value::new(vec![kind]))
296 }
297 #[allow(clippy::too_many_lines)] pub(super) async fn step<S: NamespaceStore>(
299 &self,
300 store: &S,
301 id: Hash,
302 now: u64,
303 budget: &SliceBudget,
304 ) -> Result<Fired, StoreError> {
305 let local_budget = crate::purge::SliceBudget::with_parent(64, budget.clone());
306 let service = Service::new(
307 LocalStore::new(store, &self.root, store),
308 self.root.clone(),
309 self.shards.clone(),
310 )
311 .with_purge(self.purge.clone());
312 service
313 .resume_with_local_budget(id, now, budget, &local_budget)
314 .await
315 .map_err(|_| StoreError::unavailable("denial activation pending"))?;
316 let (record, _) = service
317 .record(store, &id)
318 .await
319 .map_err(|_| bad())?
320 .ok_or_else(bad)?;
321 let (mut state, old) = self.state(store, &id, record.created).await?;
322 if state.seed > record.actions.len() {
323 return Err(bad());
324 }
325 let serving = Budgeted::new(&self.serving, budget);
326 let preserved = Budgeted::new(&self.preserved, budget);
327 let mut batch = Batch::new();
328 let mut event = "PreservationCheckpoint";
329 let mut audit_targets = vec![to_hex(&id)];
330 if !state.hold
331 && state.phase != Phase::Purging
332 && (now >= state.retain_until && !state.purged
333 || state.purged && now >= state.next_purge_at)
334 {
335 state.purge_after = None;
336 state.resume_phase = state.phase;
337 state.phase = Phase::Purging;
338 event = "PreservationPurgeStarted";
339 } else if state.phase == Phase::Seed {
340 if let Some(pack) = record.pack {
341 let (start, end) = range(b"pack-todo", &id);
342 if state.seed == 0 {
343 batch = batch
344 .put(
345 key(b"pack-todo", &id, &pack),
346 value(&inventory::InventoryCursor::default())?,
347 )
348 .put(key(b"pack-seen", &id, &pack), Value::default());
349 state.seed = 1;
350 } else if let Some((row, raw)) = store
351 .scan(&self.root, &start, &end, None, 1)
352 .await?
353 .entries
354 .first()
355 {
356 let pack: Hash = row
357 .as_bytes()
358 .strip_prefix(start.as_bytes())
359 .ok_or_else(bad)?
360 .try_into()
361 .map_err(|_| bad())?;
362 let (next, entries, done) = inventory::next(store, &pack, decode(raw)?).await?;
363 for (object, entry) in entries {
364 if entry.kind == 0 && inventory::has_seal(store, &object).await? {
365 let seen = key(b"pack-seen", &id, &object);
366 if store.get(&self.root, &seen).await?.is_none() {
367 batch = batch.put(seen, Value::default()).put(
368 key(b"pack-todo", &id, &object),
369 value(&inventory::InventoryCursor::default())?,
370 );
371 }
372 } else {
373 batch = Self::enqueue(batch, &id, &object, entry.kind);
374 }
375 }
376 if done {
377 batch = batch
378 .delete(row.clone())
379 .put(key(b"discover", &id, &pack), Value::new(vec![1]));
380 } else {
381 batch = batch.put(row.clone(), value(&next)?);
382 }
383 } else {
384 state.phase = Phase::Acquire;
385 }
386 } else {
387 for reference in record.actions.iter().skip(state.seed).take(32) {
388 batch = Self::enqueue(batch, &id, &reference.object, 0)
389 .put(key(b"discover", &id, &reference.object), Value::default());
390 state.seed += 1;
391 }
392 if state.seed == record.actions.len() {
393 state.phase = Phase::Acquire;
394 }
395 }
396 } else if state.phase == Phase::Acquire && state.purged {
397 state.phase = Phase::Discover;
398 } else if state.phase == Phase::Acquire {
399 let (start, end) = range(b"todo", &id);
400 let page = store.scan(&self.root, &start, &end, None, 1).await?;
401 if let Some((todo, expected)) = page.entries.first() {
402 let object: Hash = todo
403 .as_bytes()
404 .strip_prefix(start.as_bytes())
405 .ok_or_else(bad)?
406 .try_into()
407 .map_err(|_| bad())?;
408 let mut info = self.info(store, &id, &object).await?;
409 let repo = intent::repository(&record.repository).map_err(|_| bad())?;
410 let checkpoint_key = key(b"source", &id, &object);
411 let old_checkpoint = store.get(&self.root, &checkpoint_key).await?;
412 let checkpoint = old_checkpoint
413 .as_ref()
414 .map(decode::<source::Checkpoint>)
415 .transpose()?
416 .unwrap_or_else(|| source::Checkpoint::new(object));
417 let prefix = key(b"source-frame", &id, &object);
418 if info.source_failed {
419 batch = batch.delete(todo.clone());
421 } else if checkpoint.next.is_some() {
422 let next = source::step(
423 store,
424 self.shards.as_ref(),
425 &repo,
426 &prefix,
427 &self.profile,
428 checkpoint,
429 )
430 .await?;
431 if next.checkpoint.next.is_none() {
432 event = "PreservationSourceSelected";
433 }
434 batch.writes.extend(next.batch.writes);
435 batch.preconditions.extend(next.batch.preconditions);
436 batch = batch
437 .require(old_checkpoint.map_or_else(
438 || Precondition::Absent(checkpoint_key.clone()),
439 |raw| Precondition::Equals(checkpoint_key.clone(), raw),
440 ))
441 .put(checkpoint_key, value(&next.checkpoint)?);
442 } else {
443 let source = acquisition::resolve_selected(
444 &serving,
445 store,
446 self.shards.as_ref(),
447 &repo,
448 object,
449 &self.profile,
450 &self.root,
451 &prefix,
452 )
453 .await;
454 let source = match source {
455 Ok(source) => Some(source),
456 Err(error) if error.code() == crate::Code::DataLoss => {
457 let row = Key::new([prefix.as_bytes(), &0u32.to_be_bytes()].concat());
458 let raw = store.get(&self.root, &row).await?.ok_or_else(bad)?;
459 let (_, selected) = source::decode_frame(&raw)?;
460 audit_targets.extend([to_hex(&object), to_hex(&selected.pack)]);
461 info.source_failed = true;
462 info.verified = false;
463 state.verification = Verification::SourceCorrupt;
464 batch = batch
466 .require(Precondition::Equals(
467 checkpoint_key,
468 old_checkpoint.ok_or_else(bad)?,
469 ))
470 .require(Precondition::Equals(row, raw))
471 .delete(todo.clone())
472 .put(key(b"object", &id, &object), value(&info)?)
473 .put(key(b"discover", &id, &object), Value::default());
474 event = "PreservationSourceCorrupt";
475 None
476 }
477 Err(_) => {
478 return Err(StoreError::unavailable("preservation source unavailable"));
479 }
480 };
481 if let Some(source) = source {
482 if expected.as_bytes().len() != 1
483 || expected.as_bytes()[0] != 0 && expected.as_bytes()[0] != source.kind
484 || record.pack.is_none() && !matches!(source.kind, 1 | 5)
485 {
486 return Err(bad());
487 }
488 let size = u64::try_from(source.canonical.len()).map_err(|_| bad())?;
489 if info.copied > size
490 || info.kind != 0 && (info.kind != source.kind || info.size != size)
491 {
492 return Err(bad());
493 }
494 if source.kind == 5 && state.verification != Verification::SourceCorrupt {
496 state.verification = Verification::ManifestClosurePending;
497 }
498 info.kind = source.kind;
499 info.size = size;
500 for _ in 0..8 {
501 let offset = usize::try_from(info.copied).map_err(|_| bad())?;
502 if offset == source.canonical.len() {
503 break;
504 }
505 let end = source
506 .canonical
507 .len()
508 .min(offset.saturating_add(copy::PIECE_BYTES));
509 budget.charge()?;
511 budget.charge()?;
512 let piece = copy::plan(
513 &id,
514 &object,
515 info.copied,
516 &source.canonical[offset..end],
517 )?;
518 let piece_key = key(
519 b"piece",
520 &id,
521 &[object.as_slice(), &info.copied.to_be_bytes()].concat(),
522 );
523 let mut intent = crate::admin::plan_system(
524 store,
525 &self.root,
526 "system:timer",
527 "system:timer/PreservationPieceIntent",
528 &[to_hex(&id)],
529 now,
530 )
531 .await
532 .map_err(|_| bad())?;
533 intent = intent
534 .require(Precondition::Equals(
535 key(b"state", &id, &[]),
536 old.clone().ok_or_else(bad)?,
537 ))
538 .require(Precondition::NotAfter(if state.hold {
539 u64::try_from(self.clock.now_ms())
540 .map_err(|_| bad())?
541 .saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS)
542 } else {
543 state.retain_until
544 }))
545 .put(piece_key, value(&piece)?);
546 if store.apply(&self.root, intent).await?
547 != crate::BatchOutcome::Committed
548 {
549 return Err(StoreError::unavailable("preservation intent raced"));
550 }
551 copy::write(
552 &preserved,
553 &id,
554 &object,
555 info.copied,
556 &source.canonical[offset..end],
557 )
558 .await?;
559 info.copied = u64::try_from(end).map_err(|_| bad())?;
560 }
561 if info.copied == size {
562 let count = if source.kind == 5 {
563 u32::from_le_bytes(
564 source
565 .canonical
566 .get(18..22)
567 .ok_or_else(bad)?
568 .try_into()
569 .map_err(|_| bad())?,
570 )
571 } else {
572 0
573 };
574 if info.chunks > count {
575 return Err(bad());
576 }
577 for index in info.chunks..count.min(info.chunks.saturating_add(64)) {
578 let offset = 22 + usize::try_from(index).map_err(|_| bad())? * 32;
579 let chunk: Hash = source
580 .canonical
581 .get(offset..offset + 32)
582 .ok_or_else(bad)?
583 .try_into()
584 .map_err(|_| bad())?;
585 let child = self.info(store, &id, &chunk).await?;
586 if child.kind != 0 && child.kind != 1 {
587 return Err(bad());
588 }
589 let done = child.verified || child.source_failed;
590 if !done {
591 batch = Self::enqueue(batch, &id, &chunk, 1);
592 }
593 info.chunks += 1;
594 }
595 if info.chunks == count {
596 info.verified = true;
597 state.verified_objects =
598 state.verified_objects.checked_add(1).ok_or_else(bad)?;
599 batch = batch
600 .delete(todo.clone())
601 .put(key(b"discover", &id, &object), Value::default());
602 if source.kind == 5 {
603 batch = batch.put(
604 key(b"closure", &id, &object),
605 value(&closure::Checkpoint::default())?,
606 );
607 }
608 }
609 }
610 batch = batch.put(key(b"object", &id, &object), value(&info)?);
611 }
612 }
613 } else {
614 if state.verification == Verification::CanonicalPending {
615 state.verification = Verification::Verified;
616 }
617 state.phase = if state.verification == Verification::ManifestClosurePending {
618 Phase::Closure
619 } else {
620 Phase::Discover
621 };
622 event = match state.verification {
623 Verification::Verified => "PreservationVerified",
624 Verification::SourceCorrupt => "PreservationAcquisitionIncomplete",
625 _ => "PreservationClosurePending",
626 };
627 }
628 } else if state.phase == Phase::Closure && state.purged {
629 state.phase = Phase::Discover;
630 } else if state.phase == Phase::Closure {
631 let (start, end) = range(b"closure", &id);
632 let page = store.scan(&self.root, &start, &end, None, 1).await?;
633 if let Some((row, raw)) = page.entries.first() {
634 let manifest: Hash = row
635 .as_bytes()
636 .strip_prefix(start.as_bytes())
637 .ok_or_else(bad)?
638 .try_into()
639 .map_err(|_| bad())?;
640 let info = self.info(store, &id, &manifest).await?;
641 let next = closure::step(
642 store,
643 &preserved,
644 &self.root,
645 &id,
646 &manifest,
647 &info,
648 decode(raw)?,
649 )
650 .await?;
651 batch = if next.complete {
652 batch.delete(row.clone())
653 } else {
654 batch.put(row.clone(), value(&next.checkpoint)?)
655 };
656 } else {
657 state.verification = Verification::Verified;
658 state.phase = Phase::Discover;
659 event = "PreservationVerified";
660 }
661 } else if state.phase == Phase::Discover {
662 let (start, end) = range(b"discover", &id);
663 let page = store.scan(&self.root, &start, &end, None, 1).await?;
664 if let Some((todo, kind)) = page.entries.first() {
665 let object: Hash = todo
666 .as_bytes()
667 .strip_prefix(start.as_bytes())
668 .ok_or_else(bad)?
669 .try_into()
670 .map_err(|_| bad())?;
671 let is_pack = kind.as_bytes() == [1];
672 let mut info = self.info(store, &id, &object).await?;
673 if info.holders_done {
674 if state.current.is_some_and(|old| old != object) {
675 return Err(bad());
676 }
677 state.current = Some(object);
678 let mut checkpoint = state.discovery.take().map_or_else(
679 || {
680 discovery::DiscoveryState::new(
681 &self.addressing,
682 &intent::repository(&record.repository)
683 .map_err(|_| bad())?
684 .namespace,
685 record.created,
686 self.discovery_margin_ms,
687 )
688 },
689 Ok,
690 )?;
691 let mut candidates_done = false;
692 if checkpoint.traversed() && !checkpoint.exhaustive() {
693 let start = key(b"known-ns", &id, &object);
694 let cursor = info.namespace_after.clone().map(Cursor::new);
695 let page = store
696 .scan(&self.root, &start, &prefix_end(&start), cursor.as_ref(), 1)
697 .await?;
698 if let Some((candidate, _)) = page.entries.first() {
699 let name = std::str::from_utf8(
700 candidate
701 .as_bytes()
702 .strip_prefix(start.as_bytes())
703 .ok_or_else(bad)?,
704 )
705 .map_err(|_| bad())?;
706 let namespace = if name == "root" {
707 crate::NamespaceKey::deployment_default()
708 } else {
709 crate::NamespaceKey::from_namespace(
710 &mkit_core::repo_identity::Namespace::parse(name)
711 .map_err(|_| bad())?,
712 )
713 };
714 checkpoint.next_candidate(&namespace)?;
715 info.namespace_after = Some(candidate.as_bytes().to_vec());
716 batch = batch.put(key(b"object", &id, &object), value(&info)?);
717 } else {
718 candidates_done = true;
719 state.discovery = None;
720 state.current = None;
721 batch = batch.delete(todo.clone());
722 }
723 }
724 if !candidates_done {
725 let next = discovery::step(
726 store,
727 self.shards.as_ref(),
728 &self.root,
729 &id,
730 &object,
731 is_pack,
732 now,
733 checkpoint,
734 )
735 .await?;
736 if let Some(repo) = &next.repository {
737 let purge = self
738 .plan_cache_purge(store, &id, &object, repo, now, &local_budget)
739 .await?;
740 batch.preconditions.extend(purge.preconditions);
741 batch.writes.extend(purge.writes);
742 }
743 batch.preconditions.extend(next.batch.preconditions);
744 batch.writes.extend(next.batch.writes);
745 if next.complete {
746 batch = batch.delete(todo.clone());
747 state.current = None;
748 } else {
749 state.discovery = Some(next.state);
750 }
751 }
752 } else {
753 let cursor = info.holders.clone().map(Cursor::new);
754 let holders = ContentIndex::new(crate::store::BorrowedStore(store))
755 .holders(&object, cursor.as_ref(), 1)
756 .await?;
757 for holder in holders.holders {
758 batch = batch.put(
759 key(
760 b"known-ns",
761 &id,
762 &[object.as_slice(), holder.ns.as_str().as_bytes()].concat(),
763 ),
764 Value::default(),
765 );
766 let repo = RepoId {
767 namespace: holder.ns,
768 name: holder.repo,
769 };
770 let purge = self
771 .plan_cache_purge(store, &id, &object, &repo, now, &local_budget)
772 .await?;
773 batch.preconditions.extend(purge.preconditions);
774 batch.writes.extend(purge.writes);
775 if self.purge.is_none() {
776 batch =
777 batch.put(known_holder_key(&id, &object, &repo), holder_context()?);
778 }
779 }
780 info.holders = holders.next.map(|c| c.as_bytes().to_vec());
781 info.holders_done = info.holders.is_none();
782 batch = batch.put(key(b"object", &id, &object), value(&info)?);
783 }
784 } else {
785 state.phase = if state.purged {
786 Phase::Purged
787 } else {
788 Phase::Retain
789 };
790 state.discovery_complete = state.acquisition_complete()
791 && !matches!(&self.addressing, Addressing::Multi(multi) if matches!(multi.namespace_policy, crate::policy::NamespacePolicy::Any { .. }));
792 event = if state.discovery_complete {
793 "PreservationDiscoveryComplete"
794 } else {
795 "PreservationDiscoveryIncomplete"
796 };
797 }
798 } else if state.phase == Phase::Purging {
799 let (start, end) = range(b"piece", &id);
800 let page = store
801 .scan(
802 &self.root,
803 &start,
804 &end,
805 state.purge_after.clone().map(Cursor::new).as_ref(),
806 1,
807 )
808 .await?;
809 if let Some((row, raw)) = page.entries.first() {
810 let piece: copy::Piece = decode(raw)?;
811 if row
812 != &key(
813 b"piece",
814 &id,
815 &[piece.object.as_slice(), &piece.offset.to_be_bytes()].concat(),
816 )
817 {
818 return Err(bad());
819 }
820 if copy::read(&preserved, &id, &piece).await?.is_some() {
822 budget.charge()?;
823 budget.charge()?;
824 preserved.delete(&BlobKey::pack(piece.storage)).await?;
825 }
826 state.purge_after = Some(row.as_bytes().to_vec());
828 event = "PreservationPiecePurged";
829 } else {
830 state.phase = match state.resume_phase {
831 Phase::Acquire | Phase::Closure => Phase::Discover,
832 Phase::Retain => Phase::Purged,
833 other => other,
834 };
835 state.purged = true;
836 state.next_purge_at = now.saturating_add(3_600_000);
837 event = "PreservationPurgeComplete";
838 }
839 }
840 let audit = crate::admin::plan_system(
841 store,
842 &self.root,
843 "system:timer",
844 &format!("system:timer/{event}"),
845 &audit_targets,
846 now,
847 )
848 .await
849 .map_err(|_| bad())?;
850 batch.preconditions.extend(audit.preconditions);
851 batch.writes.extend(audit.writes);
852 let state_key = key(b"state", &id, &[]);
853 batch = batch
854 .require(old.map_or_else(
855 || Precondition::Absent(state_key.clone()),
856 |old| Precondition::Equals(state_key.clone(), old),
857 ))
858 .require(Precondition::NotAfter(
859 u64::try_from(self.clock.now_ms())
860 .map_err(|_| bad())?
861 .saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS),
862 ))
863 .put(state_key, value(&state)?);
864 let delay = if state.phase == Phase::Purged || state.phase == Phase::Retain && state.hold {
865 3_600_000
866 } else if state.phase == Phase::Retain {
867 state
868 .retain_until
869 .saturating_sub(now)
870 .clamp(1_000, 3_600_000)
871 } else {
872 1_000
873 };
874 Ok(Fired::Reschedule {
875 due_at_ms: now.saturating_add(delay),
876 value: Value::default(),
877 batch,
878 })
879 }
880}
881impl<S: NamespaceStore, N: NamespaceStore + Clone, B: BlobStore, P: BlobStore> TimerHandler<S>
882 for Work<N, B, P>
883{
884 fn kind(&self) -> TimerKind {
885 kinds::TAKEDOWN_WORK
886 }
887 fn max_per_tick(&self) -> Option<u32> {
888 Some(1)
889 }
890 fn fire<'a>(
891 &'a self,
892 ctx: &'a TimerCtx<'a, S>,
893 timer: &'a DueTimer,
894 ) -> crate::BoxFuture<'a, Result<Fired, StoreError>> {
895 Box::pin(async move {
896 if ctx.partition != &self.root || timer.kind != kinds::TAKEDOWN_WORK {
897 return Err(bad());
898 }
899 let id = timer.reference.as_ref().try_into().map_err(|_| bad())?;
900 let local = LocalStore::new(ctx.store, ctx.partition, &self.metadata);
901 let budget = SliceBudget::new(self.profile.slice_calls());
902 let store = Budgeted::new(&local, &budget);
903 self.step(&store, id, ctx.now_ms, &budget).await
904 })
905 }
906}
907
908#[cfg(all(test, feature = "memory"))]
909#[path = "work_tests.rs"]
910mod tests;