Skip to main content

mobius_gateway/
storage_usage.rs

1//! Bounded storage measurements shared by paired clients and telemetry collectors.
2use serde::{Deserialize, Serialize};
3use std::path::{Path, PathBuf};
4
5/// File counts and logical bytes; not an estimate of reclaimable disk space.
6#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
7pub struct StorageSize {
8    /// Sum of regular file lengths and symlink metadata lengths; targets are not followed.
9    pub bytes: u64,
10    /// Number of regular files and symlink entries.
11    pub files: u64,
12    /// False when a limit or inaccessible entry prevented measurement, or a symlink
13    /// appeared within or replaced the charged blob directory.
14    pub complete: bool,
15}
16/// A stable storage ownership category.
17#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
18pub struct StorageCategory {
19    /// Stable category ID, never a path.
20    pub id: String,
21    /// Measured size.
22    pub size: StorageSize,
23}
24/// Logical files eligible for selective cleanup while preserving chat history.
25#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
26pub struct SessionStorageUsage {
27    /// Chat ID for navigation and later cleanup.
28    pub session_id: String,
29    /// Optional owning project ID.
30    pub project_id: Option<String>,
31    /// User uploads, including shared content once per reference.
32    pub uploads: StorageSize,
33    /// Agent artifacts eligible for cleanup, including shared content once per reference.
34    pub artifacts: StorageSize,
35    /// Hidden agent observations and fork grants.
36    pub observations: StorageSize,
37}
38/// A project measurement; project totals can overlap gateway storage or other projects.
39#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
40pub struct ProjectStorageUsage {
41    /// Stable project ID.
42    pub project_id: String,
43    /// Files under the registered project directory.
44    pub size: StorageSize,
45}
46/// Comprehensive usage report. Logical references and project totals are not additive.
47#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
48pub struct StorageUsage {
49    /// Schema version independent of the gateway wire version.
50    pub version: u32,
51    /// Measurement time in Unix seconds.
52    pub measured_at: i64,
53    /// Non-overlapping categories within the gateway state directory.
54    pub categories: Vec<StorageCategory>,
55    /// Per-chat logical ownership.
56    pub sessions: Vec<SessionStorageUsage>,
57    /// Total chats measured before the telemetry transport bounds its detail rows.
58    pub session_count: usize,
59    /// Distinct registered projects.
60    pub projects: Vec<ProjectStorageUsage>,
61    /// Total registered project directories measured.
62    pub project_count: usize,
63    /// Detail rows were omitted by the measurement budget or telemetry transport limit.
64    pub details_truncated: bool,
65    /// Aggregate bytes and file counts of gateway categories only.
66    /// Completeness also covers project discovery and measurement for upload admission.
67    pub gateway_total: StorageSize,
68    /// Maximum content blob and project file bytes; absent for self-hosted gateways.
69    pub limit_bytes: Option<u64>,
70    /// Content blob and project file bytes charged against the allowance.
71    /// Nested project paths are charged once; separate file copies count separately.
72    pub used_bytes: u64,
73}
74impl StorageUsage {
75    pub(crate) fn bound_telemetry_details(&mut self) {
76        self.sessions.sort_by_key(|session| {
77            std::cmp::Reverse(
78                session
79                    .uploads
80                    .bytes
81                    .saturating_add(session.artifacts.bytes)
82                    .saturating_add(session.observations.bytes),
83            )
84        });
85        self.projects
86            .sort_by_key(|project| std::cmp::Reverse(project.size.bytes));
87        self.details_truncated |= self.sessions.len() > 64 || self.projects.len() > 64;
88        self.sessions.truncate(64);
89        self.projects.truncate(64);
90    }
91}
92
93pub(crate) struct MeasurementBudget {
94    started: std::time::Instant,
95    entries: usize,
96}
97impl MeasurementBudget {
98    pub(crate) fn new() -> Self {
99        Self {
100            started: std::time::Instant::now(),
101            entries: 0,
102        }
103    }
104    pub(crate) fn deadline(&self) -> tokio::time::Instant {
105        tokio::time::Instant::from_std(self.started) + std::time::Duration::from_secs(2)
106    }
107    async fn wait<F: std::future::Future>(&self, future: F) -> Option<F::Output> {
108        if self.exhausted() {
109            return None;
110        }
111        tokio::time::timeout_at(self.deadline(), future).await.ok()
112    }
113    fn exhausted(&self) -> bool {
114        self.entries >= 100_000 || self.started.elapsed().as_secs() >= 2
115    }
116}
117
118fn measure(path: &Path, budget: &mut MeasurementBudget) -> StorageSize {
119    measure_content(path, budget, None).0
120}
121
122fn measure_content(
123    path: &Path,
124    budget: &mut MeasurementBudget,
125    blobs: Option<&Path>,
126) -> (StorageSize, u64) {
127    measure_paths(vec![(path.to_path_buf(), 0)], budget, blobs)
128}
129
130fn measure_paths(
131    mut pending: Vec<(PathBuf, u32)>,
132    budget: &mut MeasurementBudget,
133    blobs: Option<&Path>,
134) -> (StorageSize, u64) {
135    let mut used_bytes = 0_u64;
136    let mut size = StorageSize {
137        complete: true,
138        ..Default::default()
139    };
140    while let Some((path, depth)) = pending.pop() {
141        budget.entries += 1;
142        if budget.exhausted() {
143            size.complete = false;
144            break;
145        }
146        if depth > 64 {
147            size.complete = false;
148            continue;
149        }
150        let metadata = match std::fs::symlink_metadata(&path) {
151            Ok(metadata) => metadata,
152            Err(e) if e.kind() == std::io::ErrorKind::NotFound && depth == 0 => continue,
153            Err(_) => {
154                size.complete = false;
155                continue;
156            }
157        };
158        if metadata.is_symlink() {
159            size.bytes = size.bytes.saturating_add(metadata.len());
160            size.files = size.files.saturating_add(1);
161            // Browser profiles legitimately contain runtime links. Blob storage must remain
162            // regular files: a linked blob or ancestor makes the charged total untrustworthy.
163            if blobs.is_some_and(|blobs| path.starts_with(blobs) || blobs.starts_with(&path)) {
164                size.complete = false;
165            }
166            continue;
167        }
168        if metadata.is_file() {
169            if blobs.is_some_and(|blobs| path.parent() == Some(blobs)) {
170                used_bytes = used_bytes.saturating_add(metadata.len());
171            }
172            size.bytes = size.bytes.saturating_add(metadata.len());
173            size.files = size.files.saturating_add(1);
174        }
175        if metadata.is_dir() {
176            match std::fs::read_dir(path) {
177                Ok(children) => {
178                    for child in children {
179                        if budget.exhausted() || pending.len() >= 100_000 {
180                            size.complete = false;
181                            break;
182                        }
183                        match child {
184                            Ok(child) if pending.len() < 100_000 => {
185                                pending.push((child.path(), depth + 1))
186                            }
187                            _ => size.complete = false,
188                        }
189                    }
190                }
191                Err(_) => size.complete = false,
192            }
193        }
194    }
195    (size, used_bytes)
196}
197
198fn gateway_categories(root: &Path, budget: &mut MeasurementBudget) -> (Vec<StorageCategory>, u64) {
199    let mut used_bytes = 0_u64;
200    let mut categories = std::collections::BTreeMap::<String, StorageSize>::new();
201    let entries = match std::fs::read_dir(root) {
202        Ok(entries) => entries,
203        Err(_) => {
204            return (
205                vec![StorageCategory {
206                    id: "other".into(),
207                    size: StorageSize::default(),
208                }],
209                0,
210            );
211        }
212    };
213    for entry in entries {
214        if budget.exhausted() {
215            categories.entry("other".into()).or_default().complete = false;
216            break;
217        }
218        let Ok(entry) = entry else {
219            categories.entry("other".into()).or_default().complete = false;
220            continue;
221        };
222        let name = entry.file_name();
223        let name = name.to_string_lossy();
224        let id = match name.as_ref() {
225            "session-files" => "session_files",
226            name if name.starts_with("checkpoints.sqlite3") => "session_history",
227            name if name.starts_with("bots.sqlite3") => "bots_and_routines",
228            _ => "other",
229        };
230        let blobs = root.join("session-files/blobs");
231        let (measured, charged) = measure_content(&entry.path(), budget, Some(&blobs));
232        used_bytes = used_bytes.saturating_add(charged);
233        let size = categories.entry(id.into()).or_insert_with(|| StorageSize {
234            complete: true,
235            ..Default::default()
236        });
237        size.bytes = size.bytes.saturating_add(measured.bytes);
238        size.files = size.files.saturating_add(measured.files);
239        size.complete &= measured.complete;
240    }
241    (
242        categories
243            .into_iter()
244            .map(|(id, size)| StorageCategory { id, size })
245            .collect(),
246        used_bytes,
247    )
248}
249
250fn project_sizes(
251    mut projects: Vec<(String, PathBuf)>,
252    budget: &mut MeasurementBudget,
253) -> (Vec<ProjectStorageUsage>, u64) {
254    // Paths sort parents before their children, irrespective of project ID order.
255    projects.sort_by(|(_, left), (_, right)| left.cmp(right));
256    let mut charged_root = None::<PathBuf>;
257    let mut used_bytes = 0_u64;
258    let projects = projects
259        .into_iter()
260        .map(|(project_id, path)| {
261            let size = measure(&path, budget);
262            if charged_root
263                .as_ref()
264                .is_none_or(|root| !path.starts_with(root))
265            {
266                used_bytes = used_bytes.saturating_add(size.bytes);
267                charged_root = Some(path);
268            }
269            ProjectStorageUsage { project_id, size }
270        })
271        .collect();
272    (projects, used_bytes)
273}
274
275pub(crate) async fn measure_usage(
276    root: PathBuf,
277    checkpoints: std::sync::Arc<dyn mobius::backend::checkpoint::CheckpointStore>,
278    files: mobius::backend::session_files::SessionFileStore,
279    bots: std::sync::Arc<crate::bots::BotStore>,
280    limit_bytes: Option<u64>,
281    tls: Option<crate::config::TlsConfig>,
282    mut budget: MeasurementBudget,
283) -> crate::Result<StorageUsage> {
284    use mobius::backend::checkpoint::SessionPageRequest;
285    use mobius::backend::session_files::SessionFileOrigin::{Artifact, Observation, Upload};
286    let mut sessions = Vec::new();
287    let mut projects = std::collections::BTreeMap::new();
288    let mut cursor = None;
289    let mut complete = true;
290    'pages: loop {
291        let Some(page) = budget
292            .wait(checkpoints.list_sessions_page(SessionPageRequest {
293                owner_id: None,
294                cursor,
295                limit: 128,
296            }))
297            .await
298        else {
299            complete = false;
300            break;
301        };
302        let page = page?;
303        for summary in page.sessions {
304            if budget.exhausted() {
305                complete = false;
306                break 'pages;
307            }
308            budget.entries += 1;
309            let project_id = summary.session_context.workspace_id;
310            if let Some(id) = &project_id
311                && !projects.contains_key(id)
312            {
313                let Some(metadata) = budget
314                    .wait(checkpoints.session_metadata(&summary.session_id))
315                    .await
316                else {
317                    complete = false;
318                    break 'pages;
319                };
320                if let Some(metadata) = metadata? {
321                    let bots = std::sync::Arc::clone(&bots);
322                    let root = root.clone();
323                    let tls = tls.clone();
324                    let spec = tokio::task::spawn_blocking(move || {
325                        crate::config::ChatSpec::from_metadata_if_present(
326                            &metadata,
327                            &bots,
328                            &root,
329                            tls.as_ref(),
330                        )
331                    });
332                    let Some(spec) = budget.wait(spec).await else {
333                        complete = false;
334                        break 'pages;
335                    };
336                    match spec.map_err(|error| {
337                        crate::Error::Config(format!("storage metadata task failed: {error}"))
338                    })? {
339                        Ok(Some(spec)) => {
340                            if let Some(path) = spec.workspace {
341                                projects.insert(id.clone(), path);
342                            }
343                        }
344                        Ok(None) => {}
345                        Err(_) => complete = false,
346                    }
347                }
348            }
349            let mut uploads = StorageSize {
350                complete: true,
351                ..Default::default()
352            };
353            let mut artifacts = uploads;
354            let mut observations = uploads;
355            let records = budget
356                .wait(
357                    files.list_cleanup_files(&summary.session_id, &[Upload, Artifact, Observation]),
358                )
359                .await;
360            match records {
361                Some(Ok(records)) => {
362                    budget.entries += records.len();
363                    for (origin, file) in records {
364                        let size = match origin {
365                            Upload => &mut uploads,
366                            Artifact => &mut artifacts,
367                            Observation => &mut observations,
368                        };
369                        size.bytes = size.bytes.saturating_add(file.size);
370                        size.files += 1;
371                    }
372                }
373                _ => {
374                    complete = false;
375                    uploads.complete = false;
376                    artifacts.complete = false;
377                    observations.complete = false;
378                }
379            }
380            sessions.push(SessionStorageUsage {
381                session_id: summary.session_id,
382                project_id,
383                uploads,
384                artifacts,
385                observations,
386            });
387        }
388        cursor = page.next_cursor;
389        if cursor.is_none() {
390            break;
391        }
392    }
393    let deadline = budget.deadline();
394    let measurement = tokio::task::spawn_blocking(move || {
395        let (categories, used_bytes) = gateway_categories(&root, &mut budget);
396        let (projects, project_bytes) = project_sizes(projects.into_iter().collect(), &mut budget);
397        (
398            categories,
399            projects,
400            used_bytes.saturating_add(project_bytes),
401        )
402    });
403    let (categories, projects, used_bytes) = tokio::time::timeout_at(deadline, measurement)
404        .await
405        .map_err(|_| crate::Error::Config("storage measurement exceeded its time budget".into()))?
406        .map_err(|error| crate::Error::Config(format!("storage measurement failed: {error}")))?;
407    let mut gateway_total = StorageSize {
408        complete: complete && projects.iter().all(|project| project.size.complete),
409        ..Default::default()
410    };
411    for category in &categories {
412        gateway_total.bytes = gateway_total.bytes.saturating_add(category.size.bytes);
413        gateway_total.files = gateway_total.files.saturating_add(category.size.files);
414        gateway_total.complete &= category.size.complete;
415    }
416    Ok(StorageUsage {
417        version: 1,
418        measured_at: chrono::Utc::now().timestamp(),
419        session_count: sessions.len(),
420        project_count: projects.len(),
421        details_truncated: !complete,
422        categories,
423        sessions,
424        projects,
425        gateway_total,
426        limit_bytes,
427        used_bytes,
428    })
429}
430
431#[cfg(test)]
432mod tests {
433    #[test]
434    fn measures_regular_files_without_following_symlinks() {
435        let root = tempfile::tempdir().unwrap();
436        std::fs::write(root.path().join("file"), b"data").unwrap();
437        #[cfg(unix)]
438        std::os::unix::fs::symlink(root.path(), root.path().join("loop")).unwrap();
439        let size = super::measure(root.path(), &mut super::MeasurementBudget::new());
440        #[cfg(unix)]
441        assert_eq!(
442            (size.bytes, size.files),
443            (
444                4 + std::fs::symlink_metadata(root.path().join("loop"))
445                    .unwrap()
446                    .len(),
447                2
448            )
449        );
450        #[cfg(not(unix))]
451        assert_eq!((size.bytes, size.files), (4, 1));
452        assert!(size.complete);
453    }
454    #[test]
455    fn depth_limit_skips_only_the_deep_subtree() {
456        let root = tempfile::tempdir().unwrap();
457        let file = root.path().join("file");
458        std::fs::write(&file, b"data").unwrap();
459        let (size, _) = super::measure_paths(
460            vec![(file, 1), (root.path().join("deep"), 65)],
461            &mut super::MeasurementBudget::new(),
462            None,
463        );
464        assert_eq!((size.bytes, size.files, size.complete), (4, 1, false));
465    }
466
467    #[test]
468    fn charged_blobs_are_counted_during_the_category_walk() {
469        let root = tempfile::tempdir().unwrap();
470        let blobs = root.path().join("session-files/blobs");
471        std::fs::create_dir_all(&blobs).unwrap();
472        std::fs::write(blobs.join("blob"), b"payload").unwrap();
473        std::fs::write(root.path().join("session-files/metadata"), b"meta").unwrap();
474        let (categories, used) =
475            super::gateway_categories(root.path(), &mut super::MeasurementBudget::new());
476        assert_eq!(used, 7);
477        assert_eq!(categories[0].size.bytes, 11);
478    }
479
480    #[test]
481    fn quota_includes_project_files_and_copies_but_counts_nested_paths_once() {
482        let root = tempfile::tempdir().unwrap();
483        let state = root.path().join("state");
484        let blobs = state.join("session-files/blobs");
485        let workspace = root.path().join("work");
486        let nested = workspace.join("nested");
487        let attachments = workspace.join(".mobius/attachments");
488        let sibling = root.path().join("work-other");
489        for path in [&blobs, &nested, &attachments, &sibling] {
490            std::fs::create_dir_all(path).unwrap();
491        }
492        std::fs::write(blobs.join("blob"), b"payload").unwrap();
493        std::fs::write(state.join("session-files/metadata"), b"meta").unwrap();
494        std::fs::write(workspace.join("file"), b"payload").unwrap();
495        std::fs::write(attachments.join("copy"), b"payload").unwrap();
496        std::fs::write(nested.join("file"), b"data").unwrap();
497        std::fs::write(sibling.join("file"), b"other").unwrap();
498        let mut budget = super::MeasurementBudget::new();
499        let (_, blob_bytes) = super::gateway_categories(&state, &mut budget);
500        let (projects, project_bytes) = super::project_sizes(
501            vec![
502                ("nested".into(), nested),
503                ("parent".into(), workspace.clone()),
504                ("same-path".into(), workspace),
505                ("sibling".into(), sibling),
506            ],
507            &mut budget,
508        );
509        assert_eq!(blob_bytes + project_bytes, 30);
510        assert!(projects.iter().all(|project| project.size.complete));
511        let sizes: std::collections::BTreeMap<_, _> = projects
512            .into_iter()
513            .map(|project| (project.project_id, (project.size.bytes, project.size.files)))
514            .collect();
515        assert_eq!(sizes["parent"], (18, 3));
516        assert_eq!(sizes["same-path"], (18, 3));
517        assert_eq!(sizes["nested"], (4, 1));
518        assert_eq!(sizes["sibling"], (5, 1));
519    }
520
521    #[cfg(unix)]
522    #[test]
523    fn browser_profile_symlinks_preserve_complete_usage_without_charging_targets() {
524        let root = tempfile::tempdir().unwrap();
525        let outside = tempfile::tempdir().unwrap();
526        std::fs::write(outside.path().join("unrelated"), b"outside data").unwrap();
527        let profile = root.path().join("desktop/profile");
528        std::fs::create_dir_all(&profile).unwrap();
529        std::fs::write(profile.join("Preferences"), b"preferences").unwrap();
530        let link = profile.join("SingletonSocket");
531        std::os::unix::fs::symlink(outside.path().join("missing-runtime-socket"), &link).unwrap();
532        let (categories, used) =
533            super::gateway_categories(root.path(), &mut super::MeasurementBudget::new());
534        assert_eq!(used, 0);
535        assert!(categories.iter().all(|category| category.size.complete));
536        assert_eq!(categories[0].size.files, 2);
537        assert_eq!(
538            categories[0].size.bytes,
539            11 + std::fs::symlink_metadata(link).unwrap().len()
540        );
541    }
542
543    #[cfg(unix)]
544    #[test]
545    fn charged_blob_symlinks_and_linked_blob_roots_make_usage_incomplete() {
546        let root = tempfile::tempdir().unwrap();
547        let outside = tempfile::tempdir().unwrap();
548        std::fs::write(outside.path().join("blob"), b"outside payload").unwrap();
549        let blobs = root.path().join("session-files/blobs");
550        std::fs::create_dir_all(&blobs).unwrap();
551        std::fs::write(blobs.join("regular"), b"data").unwrap();
552        std::os::unix::fs::symlink(outside.path().join("blob"), blobs.join("linked")).unwrap();
553        let (categories, used) =
554            super::gateway_categories(root.path(), &mut super::MeasurementBudget::new());
555        assert_eq!(used, 4);
556        assert!(!categories[0].size.complete);
557        std::fs::remove_dir_all(&blobs).unwrap();
558        std::os::unix::fs::symlink(outside.path(), &blobs).unwrap();
559        let (categories, used) =
560            super::gateway_categories(root.path(), &mut super::MeasurementBudget::new());
561        assert_eq!(used, 0);
562        assert!(!categories[0].size.complete);
563    }
564
565    #[tokio::test]
566    async fn expired_measurement_budget_skips_async_reads() {
567        let budget = super::MeasurementBudget {
568            started: std::time::Instant::now() - std::time::Duration::from_secs(3),
569            entries: 0,
570        };
571        assert!(
572            budget
573                .wait(async { panic!("expired operation was polled") })
574                .await
575                .is_none()
576        );
577    }
578}