1use serde::{Deserialize, Serialize};
3use std::path::{Path, PathBuf};
4
5#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
7pub struct StorageSize {
8 pub bytes: u64,
10 pub files: u64,
12 pub complete: bool,
15}
16#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
18pub struct StorageCategory {
19 pub id: String,
21 pub size: StorageSize,
23}
24#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
26pub struct SessionStorageUsage {
27 pub session_id: String,
29 pub project_id: Option<String>,
31 pub uploads: StorageSize,
33 pub artifacts: StorageSize,
35 pub observations: StorageSize,
37}
38#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
40pub struct ProjectStorageUsage {
41 pub project_id: String,
43 pub size: StorageSize,
45}
46#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
48pub struct StorageUsage {
49 pub version: u32,
51 pub measured_at: i64,
53 pub categories: Vec<StorageCategory>,
55 pub sessions: Vec<SessionStorageUsage>,
57 pub session_count: usize,
59 pub projects: Vec<ProjectStorageUsage>,
61 pub project_count: usize,
63 pub details_truncated: bool,
65 pub gateway_total: StorageSize,
68 pub limit_bytes: Option<u64>,
70 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 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 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}