Skip to main content

aven_core/attachments/
lifecycle.rs

1#![allow(dead_code)]
2
3use std::fs;
4use std::path::{Path, PathBuf};
5use std::time::Duration;
6
7use anyhow::{Context, Result, bail};
8use chrono::{DateTime, SecondsFormat, Utc};
9use sqlx::{Row, SqliteConnection};
10
11use crate::attachments::storage::object_path;
12use crate::attachments::validation::validate_sha256;
13use crate::db::begin_immediate;
14use crate::ids::new_id;
15
16pub const DEFAULT_LOCAL_GRACE: Duration = Duration::from_secs(7 * 24 * 60 * 60);
17pub const DEFAULT_ORIGINAL_QUOTA_BYTES: i64 = 10 * 1024 * 1024 * 1024;
18pub const DEFAULT_PREVIEW_QUOTA_BYTES: u64 = 512 * 1024 * 1024;
19pub const DEFAULT_MAINTENANCE_LIMIT: usize = 128;
20const LEASE_TTL: Duration = Duration::from_secs(10 * 60);
21
22pub trait Clock: Send + Sync {
23    fn now(&self) -> DateTime<Utc>;
24}
25
26pub struct SystemClock;
27
28impl Clock for SystemClock {
29    fn now(&self) -> DateTime<Utc> {
30        Utc::now()
31    }
32}
33
34#[derive(Debug, Clone, Copy)]
35pub struct LifecyclePolicy {
36    pub grace: Duration,
37    pub quota_bytes: i64,
38    pub preview_quota_bytes: u64,
39    pub maintenance_limit: usize,
40}
41
42impl Default for LifecyclePolicy {
43    fn default() -> Self {
44        Self {
45            grace: DEFAULT_LOCAL_GRACE,
46            quota_bytes: DEFAULT_ORIGINAL_QUOTA_BYTES,
47            preview_quota_bytes: DEFAULT_PREVIEW_QUOTA_BYTES,
48            maintenance_limit: DEFAULT_MAINTENANCE_LIMIT,
49        }
50    }
51}
52
53#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
54pub struct ByteCount {
55    pub count: u64,
56    pub bytes: u64,
57}
58
59#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
60pub struct LifecycleReport {
61    pub referenced: ByteCount,
62    pub protected: ByteCount,
63    pub grace_period: ByteCount,
64    pub eligible: ByteCount,
65    pub staging: ByteCount,
66    pub trash: ByteCount,
67    pub reservations: ByteCount,
68    pub quota: ByteCount,
69    pub inconsistencies: ByteCount,
70}
71
72#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
73pub struct PruneSummary {
74    pub eligible: ByteCount,
75    pub pruned: ByteCount,
76}
77
78fn timestamp(now: DateTime<Utc>) -> String {
79    now.to_rfc3339_opts(SecondsFormat::Secs, true)
80}
81
82fn cutoff(now: DateTime<Utc>, grace: Duration) -> Result<String> {
83    let grace = chrono::Duration::from_std(grace)?;
84    Ok(timestamp(now - grace))
85}
86
87fn trash_dir(blob_dir: &Path) -> PathBuf {
88    blob_dir.join("trash")
89}
90
91fn staging_dir(blob_dir: &Path) -> PathBuf {
92    blob_dir.join("objects").join("sha256")
93}
94
95pub async fn reconcile_liveness(conn: &mut SqliteConnection, clock: &dyn Clock) -> Result<()> {
96    let mut tx = begin_immediate(conn).await?;
97    reconcile_liveness_in_transaction(&mut tx, clock).await?;
98    tx.commit().await?;
99    Ok(())
100}
101
102pub(crate) async fn reconcile_liveness_in_transaction(
103    conn: &mut SqliteConnection,
104    clock: &dyn Clock,
105) -> Result<()> {
106    let now = timestamp(clock.now());
107    sqlx::query("DELETE FROM blob_leases WHERE expires_at <= ?")
108        .bind(&now)
109        .execute(&mut *conn)
110        .await?;
111    sqlx::query("DELETE FROM blob_upload_reservations WHERE expires_at <= ?")
112        .bind(&now)
113        .execute(&mut *conn)
114        .await?;
115    sqlx::query(
116        "INSERT OR IGNORE INTO blob_lifecycle(sha256, unreferenced_at)
117         SELECT sha256, NULL FROM blob_inventory",
118    )
119    .execute(&mut *conn)
120    .await?;
121    sqlx::query(
122        "UPDATE blob_lifecycle SET unreferenced_at = NULL
123         WHERE EXISTS (
124           SELECT 1 FROM task_attachments ta
125           JOIN tasks t ON t.workspace_id = ta.workspace_id AND t.id = ta.task_id
126           WHERE ta.sha256 = blob_lifecycle.sha256 AND ta.deleted = 0 AND t.deleted = 0
127         ) OR EXISTS (
128           SELECT 1 FROM server_blob_references sbr
129           LEFT JOIN server_task_tombstones st
130             ON st.workspace_id = sbr.workspace_id AND st.task_id = sbr.task_id
131           WHERE sbr.sha256 = blob_lifecycle.sha256 AND sbr.deleted = 0
132             AND COALESCE(st.deleted, 0) = 0
133         )",
134    )
135    .execute(&mut *conn)
136    .await?;
137    sqlx::query(
138        "UPDATE blob_lifecycle SET unreferenced_at = ?
139         WHERE unreferenced_at IS NULL
140           AND NOT EXISTS (
141             SELECT 1 FROM task_attachments ta
142             JOIN tasks t ON t.workspace_id = ta.workspace_id AND t.id = ta.task_id
143             WHERE ta.sha256 = blob_lifecycle.sha256 AND ta.deleted = 0 AND t.deleted = 0
144           )
145           AND NOT EXISTS (
146             SELECT 1 FROM server_blob_references sbr
147             LEFT JOIN server_task_tombstones st
148               ON st.workspace_id = sbr.workspace_id AND st.task_id = sbr.task_id
149             WHERE sbr.sha256 = blob_lifecycle.sha256 AND sbr.deleted = 0
150               AND COALESCE(st.deleted, 0) = 0
151           )",
152    )
153    .bind(&now)
154    .execute(&mut *conn)
155    .await?;
156    Ok(())
157}
158
159async fn is_protected(conn: &mut SqliteConnection, sha256: &str, now: &str) -> Result<bool> {
160    Ok(sqlx::query_scalar::<_, bool>(
161        "SELECT
162           EXISTS(
163             SELECT 1 FROM task_attachments ta
164             JOIN tasks t ON t.workspace_id = ta.workspace_id AND t.id = ta.task_id
165             WHERE ta.sha256 = ? AND ta.deleted = 0 AND t.deleted = 0
166           ) OR EXISTS(
167             SELECT 1 FROM server_blob_references sbr
168             LEFT JOIN server_task_tombstones st
169               ON st.workspace_id = sbr.workspace_id AND st.task_id = sbr.task_id
170             WHERE sbr.sha256 = ? AND sbr.deleted = 0 AND COALESCE(st.deleted, 0) = 0
171           ) OR EXISTS(
172             SELECT 1 FROM changes
173             WHERE server_seq IS NULL AND op_type = 'attachment_add'
174               AND json_extract(payload, '$.sha256') = ?
175           ) OR EXISTS(
176             SELECT 1 FROM blob_leases WHERE sha256 = ? AND expires_at > ?
177           ) OR EXISTS(
178             SELECT 1 FROM blob_upload_reservations WHERE sha256 = ? AND expires_at > ?
179           )",
180    )
181    .bind(sha256)
182    .bind(sha256)
183    .bind(sha256)
184    .bind(sha256)
185    .bind(now)
186    .bind(sha256)
187    .bind(now)
188    .fetch_one(&mut *conn)
189    .await?)
190}
191
192pub async fn acquire_lease(
193    conn: &mut SqliteConnection,
194    sha256: &str,
195    kind: &str,
196    clock: &dyn Clock,
197) -> Result<String> {
198    validate_sha256(sha256)?;
199    if !matches!(kind, "staging" | "read" | "backup" | "transfer") {
200        bail!("error attachment-lease-kind-invalid");
201    }
202    let lease_id = new_id();
203    let now = clock.now();
204    let expires = now + chrono::Duration::from_std(LEASE_TTL)?;
205    sqlx::query(
206        "INSERT INTO blob_leases(lease_id, sha256, kind, created_at, expires_at)
207         VALUES (?, ?, ?, ?, ?)",
208    )
209    .bind(&lease_id)
210    .bind(sha256)
211    .bind(kind)
212    .bind(timestamp(now))
213    .bind(timestamp(expires))
214    .execute(&mut *conn)
215    .await?;
216    Ok(lease_id)
217}
218
219pub async fn release_lease(conn: &mut SqliteConnection, lease_id: &str) -> Result<()> {
220    sqlx::query("DELETE FROM blob_leases WHERE lease_id = ?")
221        .bind(lease_id)
222        .execute(&mut *conn)
223        .await?;
224    Ok(())
225}
226
227pub async fn reserve_upload(
228    conn: &mut SqliteConnection,
229    workspace_id: &str,
230    sha256: &str,
231    byte_size: i64,
232    quota_bytes: i64,
233    clock: &dyn Clock,
234) -> Result<Option<String>> {
235    validate_sha256(sha256)?;
236    let mut tx = begin_immediate(conn).await?;
237    let existing: bool = sqlx::query_scalar(
238        "SELECT EXISTS(
239           SELECT 1 FROM server_blob_references sbr
240           LEFT JOIN server_task_tombstones st
241             ON st.workspace_id = sbr.workspace_id AND st.task_id = sbr.task_id
242           WHERE sbr.workspace_id = ? AND sbr.sha256 = ? AND sbr.deleted = 0
243             AND COALESCE(st.deleted, 0) = 0
244         )",
245    )
246    .bind(workspace_id)
247    .bind(sha256)
248    .fetch_one(&mut *tx)
249    .await?;
250    if existing {
251        tx.commit().await?;
252        return Ok(None);
253    }
254    let used: i64 = sqlx::query_scalar(
255        "SELECT COALESCE(SUM(byte_size), 0) FROM (
256           SELECT sbr.sha256, MAX(sbr.byte_size) AS byte_size
257           FROM server_blob_references sbr
258           LEFT JOIN server_task_tombstones st
259             ON st.workspace_id = sbr.workspace_id AND st.task_id = sbr.task_id
260           WHERE sbr.workspace_id = ? AND sbr.deleted = 0 AND COALESCE(st.deleted, 0) = 0
261           GROUP BY sbr.sha256
262         )",
263    )
264    .bind(workspace_id)
265    .fetch_one(&mut *tx)
266    .await?;
267    let reserved: i64 = sqlx::query_scalar(
268        "SELECT COALESCE(SUM(byte_size), 0) FROM blob_upload_reservations
269         WHERE workspace_id = ? AND sha256 != ? AND expires_at > ?",
270    )
271    .bind(workspace_id)
272    .bind(sha256)
273    .bind(timestamp(clock.now()))
274    .fetch_one(&mut *tx)
275    .await?;
276    if used.saturating_add(reserved).saturating_add(byte_size) > quota_bytes {
277        tx.rollback().await?;
278        bail!("error attachment-quota-exceeded");
279    }
280    let reservation_id = new_id();
281    let now = clock.now();
282    let expires = now + chrono::Duration::from_std(LEASE_TTL)?;
283    sqlx::query(
284        "INSERT INTO blob_upload_reservations(
285           reservation_id, workspace_id, sha256, byte_size, created_at, expires_at
286         ) VALUES (?, ?, ?, ?, ?, ?)
287         ON CONFLICT(workspace_id, sha256) DO UPDATE SET
288           reservation_id = excluded.reservation_id,
289           byte_size = excluded.byte_size,
290           created_at = excluded.created_at,
291           expires_at = excluded.expires_at",
292    )
293    .bind(&reservation_id)
294    .bind(workspace_id)
295    .bind(sha256)
296    .bind(byte_size)
297    .bind(timestamp(now))
298    .bind(timestamp(expires))
299    .execute(&mut *tx)
300    .await?;
301    tx.commit().await?;
302    Ok(Some(reservation_id))
303}
304
305pub async fn release_reservation(conn: &mut SqliteConnection, reservation_id: &str) -> Result<()> {
306    sqlx::query("DELETE FROM blob_upload_reservations WHERE reservation_id = ?")
307        .bind(reservation_id)
308        .execute(&mut *conn)
309        .await?;
310    Ok(())
311}
312
313pub async fn local_unique_bytes(conn: &mut SqliteConnection) -> Result<i64> {
314    Ok(sqlx::query_scalar(
315        "SELECT COALESCE(SUM(byte_size), 0) FROM blob_inventory WHERE available = 1",
316    )
317    .fetch_one(&mut *conn)
318    .await?)
319}
320
321pub async fn ensure_local_capacity(
322    conn: &mut SqliteConnection,
323    blob_dir: &Path,
324    sha256: &str,
325    byte_size: i64,
326    policy: LifecyclePolicy,
327    clock: &dyn Clock,
328) -> Result<Option<String>> {
329    let used = local_unique_bytes(conn).await?;
330    if used.saturating_add(byte_size) > policy.quota_bytes {
331        prune(conn, blob_dir, policy, true, clock).await?;
332    }
333    let now = clock.now();
334    let mut tx = begin_immediate(conn).await?;
335    let existing: bool = sqlx::query_scalar(
336        "SELECT EXISTS(SELECT 1 FROM blob_inventory WHERE sha256 = ? AND available = 1)",
337    )
338    .bind(sha256)
339    .fetch_one(&mut *tx)
340    .await?;
341    if existing {
342        tx.commit().await?;
343        return Ok(None);
344    }
345    let used: i64 = sqlx::query_scalar(
346        "SELECT COALESCE(SUM(byte_size), 0) FROM blob_inventory WHERE available = 1",
347    )
348    .fetch_one(&mut *tx)
349    .await?;
350    let reserved: i64 = sqlx::query_scalar(
351        "SELECT COALESCE(SUM(byte_size), 0) FROM blob_upload_reservations
352         WHERE workspace_id = '__local__' AND sha256 != ? AND expires_at > ?",
353    )
354    .bind(sha256)
355    .bind(timestamp(now))
356    .fetch_one(&mut *tx)
357    .await?;
358    if used.saturating_add(reserved).saturating_add(byte_size) > policy.quota_bytes {
359        tx.rollback().await?;
360        bail!("error attachment-quota-exceeded");
361    }
362    let reservation_id = new_id();
363    let expires = now + chrono::Duration::from_std(LEASE_TTL)?;
364    sqlx::query(
365        "INSERT INTO blob_upload_reservations(
366           reservation_id, workspace_id, sha256, byte_size, created_at, expires_at
367         ) VALUES (?, '__local__', ?, ?, ?, ?)
368         ON CONFLICT(workspace_id, sha256) DO UPDATE SET
369           reservation_id = excluded.reservation_id,
370           byte_size = excluded.byte_size,
371           created_at = excluded.created_at,
372           expires_at = excluded.expires_at",
373    )
374    .bind(&reservation_id)
375    .bind(sha256)
376    .bind(byte_size)
377    .bind(timestamp(now))
378    .bind(timestamp(expires))
379    .execute(&mut *tx)
380    .await?;
381    tx.commit().await?;
382    Ok(Some(reservation_id))
383}
384
385pub async fn reconcile_trash(conn: &mut SqliteConnection, blob_dir: &Path) -> Result<()> {
386    let trash = trash_dir(blob_dir);
387    if !trash.exists() {
388        return Ok(());
389    }
390    for entry in fs::read_dir(&trash)? {
391        let entry = entry?;
392        if !entry.file_type()?.is_file() {
393            continue;
394        }
395        let Some(sha256) = entry.file_name().to_str().map(str::to_owned) else {
396            continue;
397        };
398        if validate_sha256(&sha256).is_err() {
399            continue;
400        }
401        let available: bool = sqlx::query_scalar(
402            "SELECT COALESCE((SELECT available FROM blob_inventory WHERE sha256 = ?), 0)",
403        )
404        .bind(&sha256)
405        .fetch_one(&mut *conn)
406        .await?;
407        if available {
408            let target = object_path(blob_dir, &sha256)?;
409            if !target.exists() {
410                if let Some(parent) = target.parent() {
411                    fs::create_dir_all(parent)?;
412                }
413                fs::rename(entry.path(), target)?;
414            } else {
415                fs::remove_file(entry.path())?;
416            }
417        } else {
418            fs::remove_file(entry.path())?;
419        }
420    }
421    Ok(())
422}
423
424pub async fn reconcile_staging(blob_dir: &Path) -> Result<ByteCount> {
425    let mut removed = ByteCount::default();
426    let dir = staging_dir(blob_dir);
427    if !dir.exists() {
428        return Ok(removed);
429    }
430    for entry in fs::read_dir(dir)? {
431        let entry = entry?;
432        let name = entry.file_name();
433        if !name.to_string_lossy().starts_with(".aven-stage-") {
434            continue;
435        }
436        let metadata = entry.metadata()?;
437        let stale = metadata
438            .modified()?
439            .elapsed()
440            .is_ok_and(|age| age >= LEASE_TTL);
441        if !stale {
442            continue;
443        }
444        removed.count += 1;
445        removed.bytes += metadata.len();
446        fs::remove_file(entry.path())?;
447    }
448    Ok(removed)
449}
450
451pub async fn reconcile_orphan_objects(
452    conn: &mut SqliteConnection,
453    blob_dir: &Path,
454    grace: Duration,
455    clock: &dyn Clock,
456) -> Result<ByteCount> {
457    let mut removed = ByteCount::default();
458    let dir = staging_dir(blob_dir);
459    if !dir.exists() {
460        return Ok(removed);
461    }
462    let cutoff = clock.now() - chrono::Duration::from_std(grace)?;
463    let now = timestamp(clock.now());
464    for entry in fs::read_dir(dir)? {
465        let entry = entry?;
466        if !entry.file_type()?.is_file() {
467            continue;
468        }
469        let Some(sha256) = entry.file_name().to_str().map(str::to_owned) else {
470            continue;
471        };
472        if validate_sha256(&sha256).is_err() {
473            continue;
474        }
475        let tracked: bool =
476            sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM blob_inventory WHERE sha256 = ?)")
477                .bind(&sha256)
478                .fetch_one(&mut *conn)
479                .await?;
480        let metadata = entry.metadata()?;
481        if tracked
482            || DateTime::<Utc>::from(metadata.modified()?) > cutoff
483            || is_protected(conn, &sha256, &now).await?
484        {
485            continue;
486        }
487        let trash = trash_dir(blob_dir);
488        fs::create_dir_all(&trash)?;
489        let trashed = trash.join(&sha256);
490        fs::rename(entry.path(), &trashed)?;
491        fs::remove_file(trashed)?;
492        removed.count += 1;
493        removed.bytes += metadata.len();
494    }
495    Ok(removed)
496}
497
498pub async fn prune(
499    conn: &mut SqliteConnection,
500    blob_dir: &Path,
501    policy: LifecyclePolicy,
502    apply: bool,
503    clock: &dyn Clock,
504) -> Result<PruneSummary> {
505    if apply {
506        reconcile_trash(conn, blob_dir).await?;
507        reconcile_staging(blob_dir).await?;
508    }
509    reconcile_liveness(conn, clock).await?;
510    if apply {
511        reconcile_orphan_objects(conn, blob_dir, policy.grace, clock).await?;
512    }
513    let now = timestamp(clock.now());
514    let cutoff = cutoff(clock.now(), policy.grace)?;
515    let rows = sqlx::query(
516        "SELECT bi.sha256, bi.byte_size
517         FROM blob_inventory bi
518         JOIN blob_lifecycle bl ON bl.sha256 = bi.sha256
519         WHERE bi.available = 1 AND bl.unreferenced_at IS NOT NULL
520           AND bl.unreferenced_at <= ?
521         ORDER BY bl.unreferenced_at, bi.sha256 LIMIT ?",
522    )
523    .bind(&cutoff)
524    .bind(i64::try_from(policy.maintenance_limit)?)
525    .fetch_all(&mut *conn)
526    .await?;
527    let mut summary = PruneSummary::default();
528    for row in rows {
529        let sha256: String = row.get("sha256");
530        let byte_size: i64 = row.get("byte_size");
531        if is_protected(conn, &sha256, &now).await? {
532            continue;
533        }
534        summary.eligible.count += 1;
535        summary.eligible.bytes += u64::try_from(byte_size)?;
536        if !apply {
537            continue;
538        }
539        let mut tx = begin_immediate(conn).await?;
540        let still_eligible: bool = sqlx::query_scalar(
541            "SELECT EXISTS(
542               SELECT 1 FROM blob_inventory bi
543               JOIN blob_lifecycle bl ON bl.sha256 = bi.sha256
544               WHERE bi.sha256 = ? AND bi.available = 1
545                 AND bl.unreferenced_at IS NOT NULL AND bl.unreferenced_at <= ?
546             )",
547        )
548        .bind(&sha256)
549        .bind(&cutoff)
550        .fetch_one(&mut *tx)
551        .await?;
552        if !still_eligible || is_protected(&mut tx, &sha256, &now).await? {
553            tx.rollback().await?;
554            continue;
555        }
556        let source = object_path(blob_dir, &sha256)?;
557        let trash = trash_dir(blob_dir);
558        fs::create_dir_all(&trash)?;
559        let trashed = trash.join(&sha256);
560        if source.exists() {
561            fs::rename(&source, &trashed).with_context(|| "could not move attachment to trash")?;
562        }
563        if let Err(error) = sqlx::query(
564            "UPDATE blob_inventory SET available = 0, last_verified_at = ? WHERE sha256 = ?",
565        )
566        .bind(&now)
567        .bind(&sha256)
568        .execute(&mut *tx)
569        .await
570        {
571            if trashed.exists() {
572                let _ = fs::rename(&trashed, &source);
573            }
574            return Err(error.into());
575        }
576        tx.commit().await?;
577        if trashed.exists() {
578            fs::remove_file(&trashed)?;
579        }
580        summary.pruned.count += 1;
581        summary.pruned.bytes += u64::try_from(byte_size)?;
582    }
583    prune_preview_cache(blob_dir, policy.preview_quota_bytes)?;
584    Ok(summary)
585}
586
587pub fn prune_preview_cache(blob_dir: &Path, quota: u64) -> Result<ByteCount> {
588    let root = blob_dir.join("cache").join("previews");
589    if !root.exists() {
590        return Ok(ByteCount::default());
591    }
592    let mut files = Vec::new();
593    let mut dirs = vec![root];
594    while let Some(dir) = dirs.pop() {
595        for entry in fs::read_dir(dir)? {
596            let entry = entry?;
597            if entry.file_type()?.is_dir() {
598                dirs.push(entry.path());
599            } else if entry.file_type()?.is_file() {
600                let metadata = entry.metadata()?;
601                files.push((metadata.modified()?, metadata.len(), entry.path()));
602            }
603        }
604    }
605    let mut total: u64 = files.iter().map(|(_, size, _)| *size).sum();
606    files.sort_by_key(|(modified, _, path)| (*modified, path.clone()));
607    let mut removed = ByteCount::default();
608    for (_, size, path) in files {
609        if total <= quota {
610            break;
611        }
612        fs::remove_file(path)?;
613        total -= size;
614        removed.count += 1;
615        removed.bytes += size;
616    }
617    Ok(removed)
618}
619
620pub async fn lifecycle_report(
621    conn: &mut SqliteConnection,
622    blob_dir: &Path,
623    policy: LifecyclePolicy,
624    clock: &dyn Clock,
625) -> Result<LifecycleReport> {
626    reconcile_liveness(conn, clock).await?;
627    let now = timestamp(clock.now());
628    let cutoff = cutoff(clock.now(), policy.grace)?;
629    let rows = sqlx::query(
630        "SELECT bi.sha256, bi.byte_size, bi.available, bl.unreferenced_at,
631           EXISTS(
632             SELECT 1 FROM task_attachments ta JOIN tasks t
633             ON t.workspace_id = ta.workspace_id AND t.id = ta.task_id
634             WHERE ta.sha256 = bi.sha256 AND ta.deleted = 0 AND t.deleted = 0
635           ) OR EXISTS(
636             SELECT 1 FROM server_blob_references sbr
637             LEFT JOIN server_task_tombstones st
638               ON st.workspace_id = sbr.workspace_id AND st.task_id = sbr.task_id
639             WHERE sbr.sha256 = bi.sha256 AND sbr.deleted = 0 AND COALESCE(st.deleted, 0) = 0
640           ) AS referenced
641         FROM blob_inventory bi LEFT JOIN blob_lifecycle bl ON bl.sha256 = bi.sha256",
642    )
643    .fetch_all(&mut *conn)
644    .await?;
645    let mut report = LifecycleReport::default();
646    for row in rows {
647        let sha256: String = row.get("sha256");
648        let bytes = u64::try_from(row.get::<i64, _>("byte_size"))?;
649        let available = row.get::<i64, _>("available") != 0;
650        let referenced = row.get::<i64, _>("referenced") != 0;
651        let unreferenced_at: Option<String> = row.get("unreferenced_at");
652        if available {
653            report.quota.count += 1;
654            report.quota.bytes += bytes;
655        }
656        if referenced {
657            report.referenced.count += 1;
658            report.referenced.bytes += bytes;
659        } else if is_protected(conn, &sha256, &now).await? {
660            report.protected.count += 1;
661            report.protected.bytes += bytes;
662        } else if unreferenced_at
663            .as_deref()
664            .is_some_and(|at| at <= cutoff.as_str())
665        {
666            report.eligible.count += 1;
667            report.eligible.bytes += bytes;
668        } else {
669            report.grace_period.count += 1;
670            report.grace_period.bytes += bytes;
671        }
672        if referenced && unreferenced_at.is_some() || !referenced && unreferenced_at.is_none() {
673            report.inconsistencies.count += 1;
674            report.inconsistencies.bytes += bytes;
675        }
676    }
677    let (reservation_count, reservation_bytes): (i64, i64) = sqlx::query_as(
678        "SELECT COUNT(*), COALESCE(SUM(byte_size), 0)
679         FROM blob_upload_reservations WHERE expires_at > ?",
680    )
681    .bind(&now)
682    .fetch_one(&mut *conn)
683    .await?;
684    report.reservations = ByteCount {
685        count: u64::try_from(reservation_count)?,
686        bytes: u64::try_from(reservation_bytes)?,
687    };
688    for entry in fs::read_dir(staging_dir(blob_dir))
689        .into_iter()
690        .flatten()
691        .flatten()
692    {
693        let name = entry.file_name().to_string_lossy().to_string();
694        let metadata = entry.metadata()?;
695        if name.starts_with(".aven-stage-") {
696            report.staging.count += 1;
697            report.staging.bytes += metadata.len();
698        } else if validate_sha256(&name).is_ok() {
699            let tracked: bool =
700                sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM blob_inventory WHERE sha256 = ?)")
701                    .bind(&name)
702                    .fetch_one(&mut *conn)
703                    .await?;
704            if !tracked {
705                report.staging.count += 1;
706                report.staging.bytes += metadata.len();
707                report.inconsistencies.count += 1;
708                report.inconsistencies.bytes += metadata.len();
709            }
710        }
711    }
712    for entry in fs::read_dir(trash_dir(blob_dir))
713        .into_iter()
714        .flatten()
715        .flatten()
716    {
717        if entry.file_type()?.is_file() {
718            report.trash.count += 1;
719            report.trash.bytes += entry.metadata()?.len();
720        }
721    }
722    Ok(report)
723}
724
725#[cfg(test)]
726mod tests {
727    use std::sync::{Arc, Mutex};
728
729    use chrono::DateTime;
730    use sqlx::Connection as _;
731    use sqlx::sqlite::SqliteConnectOptions;
732
733    use crate::attachments::storage::{object_path, upsert_inventory_available};
734    use crate::db::open_db;
735
736    use super::*;
737
738    #[derive(Clone)]
739    struct TestClock(Arc<Mutex<DateTime<Utc>>>);
740
741    impl TestClock {
742        fn at(value: &str) -> Self {
743            Self(Arc::new(Mutex::new(
744                DateTime::parse_from_rfc3339(value).unwrap().to_utc(),
745            )))
746        }
747
748        fn advance(&self, duration: chrono::Duration) {
749            let mut now = self.0.lock().unwrap();
750            *now += duration;
751        }
752    }
753
754    impl Clock for TestClock {
755        fn now(&self) -> DateTime<Utc> {
756            *self.0.lock().unwrap()
757        }
758    }
759
760    async fn insert_task(conn: &mut SqliteConnection, task_id: &str) {
761        sqlx::query(
762            "INSERT INTO tasks(
763               workspace_id, id, title, description, project_id, status, priority,
764               created_at, updated_at, queue_activity_at
765             ) VALUES ('0000000000000000', ?, 'task', '', 'project', 'inbox', 'none',
766                       '2026-01-01T00:00:00Z', '2026-01-01T00:00:00Z', '2026-01-01T00:00:00Z')",
767        )
768        .bind(task_id)
769        .execute(conn)
770        .await
771        .unwrap();
772    }
773
774    async fn insert_attachment(
775        conn: &mut SqliteConnection,
776        attachment_id: &str,
777        task_id: &str,
778        sha256: &str,
779        deleted: bool,
780    ) {
781        sqlx::query(
782            "INSERT INTO task_attachments(
783               workspace_id, attachment_id, task_id, sha256, byte_size, media_type, width, height,
784               created_at, deleted, deleted_at
785             ) VALUES ('0000000000000000', ?, ?, ?, 4, 'image/png', 1, 1,
786                       '2026-01-01T00:00:00Z', ?, ?)",
787        )
788        .bind(attachment_id)
789        .bind(task_id)
790        .bind(sha256)
791        .bind(i64::from(deleted))
792        .bind(deleted.then_some("2026-01-01T00:00:00Z"))
793        .execute(conn)
794        .await
795        .unwrap();
796    }
797
798    #[tokio::test]
799    async fn final_live_reference_starts_grace_once_and_restore_clears_it() {
800        let temp = tempfile::tempdir().unwrap();
801        let pool = open_db(&temp.path().join("test.sqlite")).await.unwrap();
802        let mut conn = pool.acquire().await.unwrap();
803        let hash = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
804        upsert_inventory_available(&mut conn, hash, 4, "image/png")
805            .await
806            .unwrap();
807        insert_task(&mut conn, "0000000000000001").await;
808        insert_task(&mut conn, "0000000000000002").await;
809        insert_attachment(
810            &mut conn,
811            "0000000000000011",
812            "0000000000000001",
813            hash,
814            false,
815        )
816        .await;
817        insert_attachment(
818            &mut conn,
819            "0000000000000012",
820            "0000000000000002",
821            hash,
822            false,
823        )
824        .await;
825        let clock = TestClock::at("2026-07-01T00:00:00Z");
826
827        reconcile_liveness(&mut conn, &clock).await.unwrap();
828        sqlx::query("UPDATE task_attachments SET deleted = 1, deleted_at = 'x' WHERE attachment_id = '0000000000000011'")
829            .execute(&mut *conn).await.unwrap();
830        reconcile_liveness(&mut conn, &clock).await.unwrap();
831        let value: Option<String> =
832            sqlx::query_scalar("SELECT unreferenced_at FROM blob_lifecycle WHERE sha256 = ?")
833                .bind(hash)
834                .fetch_one(&mut *conn)
835                .await
836                .unwrap();
837        assert_eq!(value, None, "one live reference keeps the hash live");
838
839        sqlx::query("UPDATE task_attachments SET deleted = 1, deleted_at = 'x' WHERE attachment_id = '0000000000000012'")
840            .execute(&mut *conn).await.unwrap();
841        reconcile_liveness(&mut conn, &clock).await.unwrap();
842        let first: String =
843            sqlx::query_scalar("SELECT unreferenced_at FROM blob_lifecycle WHERE sha256 = ?")
844                .bind(hash)
845                .fetch_one(&mut *conn)
846                .await
847                .unwrap();
848        clock.advance(chrono::Duration::days(1));
849        reconcile_liveness(&mut conn, &clock).await.unwrap();
850        let second: String =
851            sqlx::query_scalar("SELECT unreferenced_at FROM blob_lifecycle WHERE sha256 = ?")
852                .bind(hash)
853                .fetch_one(&mut *conn)
854                .await
855                .unwrap();
856        assert_eq!(first, second, "grace starts exactly once");
857
858        sqlx::query("UPDATE task_attachments SET deleted = 0, deleted_at = NULL WHERE attachment_id = '0000000000000012'")
859            .execute(&mut *conn).await.unwrap();
860        reconcile_liveness(&mut conn, &clock).await.unwrap();
861        let restored: Option<String> =
862            sqlx::query_scalar("SELECT unreferenced_at FROM blob_lifecycle WHERE sha256 = ?")
863                .bind(hash)
864                .fetch_one(&mut *conn)
865                .await
866                .unwrap();
867        assert_eq!(restored, None);
868    }
869
870    #[tokio::test]
871    async fn lease_protects_expired_unreferenced_blob_until_release() {
872        let temp = tempfile::tempdir().unwrap();
873        let db_path = temp.path().join("test.sqlite");
874        let pool = open_db(&db_path).await.unwrap();
875        let mut conn = pool.acquire().await.unwrap();
876        let blob_dir = temp.path().join("blobs");
877        let hash = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
878        upsert_inventory_available(&mut conn, hash, 4, "image/png")
879            .await
880            .unwrap();
881        let path = object_path(&blob_dir, hash).unwrap();
882        fs::create_dir_all(path.parent().unwrap()).unwrap();
883        fs::write(&path, b"blob").unwrap();
884        let clock = TestClock::at("2026-07-10T00:00:00Z");
885        reconcile_liveness(&mut conn, &clock).await.unwrap();
886        clock.advance(chrono::Duration::days(8));
887        let lease = acquire_lease(&mut conn, hash, "backup", &clock)
888            .await
889            .unwrap();
890        let policy = LifecyclePolicy::default();
891        let blocked = prune(&mut conn, &blob_dir, policy, true, &clock)
892            .await
893            .unwrap();
894        assert_eq!(blocked.pruned.count, 0);
895        assert!(path.exists());
896
897        release_lease(&mut conn, &lease).await.unwrap();
898        let pruned = prune(&mut conn, &blob_dir, policy, true, &clock)
899            .await
900            .unwrap();
901        assert_eq!(pruned.pruned.count, 1);
902        assert!(!path.exists());
903    }
904
905    #[tokio::test]
906    async fn quota_is_unique_by_hash_and_reservations_are_workspace_scoped() {
907        let temp = tempfile::tempdir().unwrap();
908        let pool = open_db(&temp.path().join("test.sqlite")).await.unwrap();
909        let mut conn = pool.acquire().await.unwrap();
910        let hash = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc";
911        let clock = TestClock::at("2026-07-01T00:00:00Z");
912        let first = reserve_upload(&mut conn, "workspace-a", hash, 8, 8, &clock)
913            .await
914            .unwrap();
915        assert!(first.is_some());
916        let replacement = reserve_upload(&mut conn, "workspace-a", hash, 8, 8, &clock)
917            .await
918            .unwrap();
919        assert!(replacement.is_some());
920        let other = reserve_upload(&mut conn, "workspace-b", hash, 8, 8, &clock)
921            .await
922            .unwrap();
923        assert!(other.is_some());
924        let count: i64 = sqlx::query_scalar("SELECT count(*) FROM blob_upload_reservations")
925            .fetch_one(&mut *conn)
926            .await
927            .unwrap();
928        assert_eq!(count, 2);
929    }
930
931    #[tokio::test]
932    async fn concurrent_attach_wins_prune_recheck() {
933        let temp = tempfile::tempdir().unwrap();
934        let db_path = temp.path().join("test.sqlite");
935        let pool = open_db(&db_path).await.unwrap();
936        let mut conn = pool.acquire().await.unwrap();
937        let blob_dir = temp.path().join("blobs");
938        let hash = "abababababababababababababababababababababababababababababababab";
939        upsert_inventory_available(&mut conn, hash, 4, "image/png")
940            .await
941            .unwrap();
942        insert_task(&mut conn, "0000000000000003").await;
943        let path = object_path(&blob_dir, hash).unwrap();
944        fs::create_dir_all(path.parent().unwrap()).unwrap();
945        fs::write(&path, b"blob").unwrap();
946        let clock = TestClock::at("2026-07-01T00:00:00Z");
947        reconcile_liveness(&mut conn, &clock).await.unwrap();
948        clock.advance(chrono::Duration::days(8));
949        let options = SqliteConnectOptions::new()
950            .filename(&db_path)
951            .busy_timeout(Duration::from_secs(5));
952        let prune_conn = SqliteConnection::connect_with(&options).await.unwrap();
953
954        let mut tx = begin_immediate(&mut conn).await.unwrap();
955        insert_attachment(&mut tx, "0000000000000013", "0000000000000003", hash, false).await;
956        let prune_dir = blob_dir.clone();
957        let prune_clock = clock.clone();
958        let pruning = tokio::spawn(async move {
959            let mut prune_conn = prune_conn;
960            prune(
961                &mut prune_conn,
962                &prune_dir,
963                LifecyclePolicy::default(),
964                true,
965                &prune_clock,
966            )
967            .await
968            .unwrap()
969        });
970        tokio::time::sleep(Duration::from_millis(50)).await;
971        tx.commit().await.unwrap();
972
973        let summary = pruning.await.unwrap();
974        assert_eq!(summary.pruned.count, 0);
975        assert!(path.exists());
976    }
977
978    #[tokio::test]
979    async fn accepted_server_reference_keeps_blob_live() {
980        let temp = tempfile::tempdir().unwrap();
981        let pool = open_db(&temp.path().join("test.sqlite")).await.unwrap();
982        let mut conn = pool.acquire().await.unwrap();
983        let hash = "acacacacacacacacacacacacacacacacacacacacacacacacacacacacacacacac";
984        upsert_inventory_available(&mut conn, hash, 4, "image/png")
985            .await
986            .unwrap();
987        sqlx::query(
988            "INSERT INTO server_blob_references(
989               workspace_id, attachment_id, task_id, sha256, byte_size
990             ) VALUES ('workspace', '0000000000000042', '0000000000000004', ?, 4)",
991        )
992        .bind(hash)
993        .execute(&mut *conn)
994        .await
995        .unwrap();
996        let clock = TestClock::at("2026-07-01T00:00:00Z");
997
998        reconcile_liveness(&mut conn, &clock).await.unwrap();
999        let unreferenced_at: Option<String> =
1000            sqlx::query_scalar("SELECT unreferenced_at FROM blob_lifecycle WHERE sha256 = ?")
1001                .bind(hash)
1002                .fetch_one(&mut *conn)
1003                .await
1004                .unwrap();
1005        assert_eq!(unreferenced_at, None);
1006    }
1007
1008    #[tokio::test]
1009    async fn local_quota_boundary_is_hash_idempotent() {
1010        let temp = tempfile::tempdir().unwrap();
1011        let db_path = temp.path().join("test.sqlite");
1012        let pool = open_db(&db_path).await.unwrap();
1013        let mut conn = pool.acquire().await.unwrap();
1014        let blob_dir = temp.path().join("blobs");
1015        let existing = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
1016        let new_hash = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
1017        upsert_inventory_available(&mut conn, existing, 8, "image/png")
1018            .await
1019            .unwrap();
1020        let clock = TestClock::at("2026-07-01T00:00:00Z");
1021        let policy = LifecyclePolicy {
1022            quota_bytes: 8,
1023            ..LifecyclePolicy::default()
1024        };
1025
1026        ensure_local_capacity(&mut conn, &blob_dir, existing, 8, policy, &clock)
1027            .await
1028            .unwrap();
1029        let error = ensure_local_capacity(&mut conn, &blob_dir, new_hash, 1, policy, &clock)
1030            .await
1031            .unwrap_err();
1032        assert_eq!(error.to_string(), "error attachment-quota-exceeded");
1033    }
1034
1035    #[tokio::test]
1036    async fn interrupted_atomic_create_object_is_reconciled_after_grace() {
1037        let temp = tempfile::tempdir().unwrap();
1038        let db_path = temp.path().join("test.sqlite");
1039        let pool = open_db(&db_path).await.unwrap();
1040        let mut conn = pool.acquire().await.unwrap();
1041        let blob_dir = temp.path().join("blobs");
1042        let hash = "cdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcdcd";
1043        let path = object_path(&blob_dir, hash).unwrap();
1044        fs::create_dir_all(path.parent().unwrap()).unwrap();
1045        fs::write(&path, b"orphan").unwrap();
1046        let clock = TestClock::at("2030-07-01T00:00:00Z");
1047        let policy = LifecyclePolicy {
1048            grace: Duration::ZERO,
1049            ..LifecyclePolicy::default()
1050        };
1051
1052        prune(&mut conn, &blob_dir, policy, true, &clock)
1053            .await
1054            .unwrap();
1055        assert!(!path.exists());
1056    }
1057
1058    #[tokio::test]
1059    async fn interrupted_trash_move_restores_available_object() {
1060        let temp = tempfile::tempdir().unwrap();
1061        let db_path = temp.path().join("test.sqlite");
1062        let pool = open_db(&db_path).await.unwrap();
1063        let mut conn = pool.acquire().await.unwrap();
1064        let blob_dir = temp.path().join("blobs");
1065        let hash = "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd";
1066        upsert_inventory_available(&mut conn, hash, 4, "image/png")
1067            .await
1068            .unwrap();
1069        let source = object_path(&blob_dir, hash).unwrap();
1070        let trash = trash_dir(&blob_dir).join(hash);
1071        fs::create_dir_all(source.parent().unwrap()).unwrap();
1072        fs::create_dir_all(trash.parent().unwrap()).unwrap();
1073        fs::write(&source, b"blob").unwrap();
1074        fs::rename(&source, &trash).unwrap();
1075
1076        reconcile_trash(&mut conn, &blob_dir).await.unwrap();
1077        assert!(source.exists());
1078        assert!(!trash.exists());
1079    }
1080}