Skip to main content

mkit_server/takedown/
work.rs

1//! One durable acquisition, discovery or independently owned retention step.
2use 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}
123/// All runtime dependencies; preservation is a separately provisioned blob store.
124pub struct Work<N, B, P> {
125    pub metadata: N,
126    /// Automatic cache purge settings shared with signed and late intake.
127    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    /// Prepare guarded legal-hold state for the signed admin framework.
147    /// Commit with the operator's audit and replay result; no public route exists yet.
148    ///
149    /// # Errors
150    /// Unknown requests, ended purge ownership, storage errors or invalid retention.
151    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        // Invalidating early is safe on a failed checkpoint; the returned
276        // durable responsibility is committed with the holder/discovery cursor.
277        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)] // Each branch performs one bounded checkpoint transition.
298    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                    // A later manifest may reference the same failed member.
420                    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                            // Keep selected source rows as provenance; remove only its pending work.
465                            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                        // Reassembly validation runs over preserved bytes after every child is acquired.
495                        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                            // Reserve PUT and its immutable-existence HEAD before the sink runs.
510                            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                // Purging is durable before DELETE; legal hold cannot enter this phase.
821                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                // Keep the intent: a delayed PUT remains discoverable on every later pass.
827                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;