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}