From 62f8c605c3bd184c8b6cb67b0eb2e035503e5902 Mon Sep 17 00:00:00 2001
From: Richie McIlroy <33632126+richiemcilroy@users.noreply.github.com>
Date: Wed, 26 Aug 2026 19:31:05 +0100
Subject: [PATCH 1/5] fix: prevent muxer reader shutdown deadlocks
---
crates/cap-muxer/src/main.rs | 104 ++++++++++++++++++++++++++--
crates/recording/tests/oop_muxer.rs | 73 +++++++++++++++++++
2 files changed, 172 insertions(+), 5 deletions(-)
diff --git a/crates/cap-muxer/src/main.rs b/crates/cap-muxer/src/main.rs
index 91add0159e6..af74843ed9c 100644
--- a/crates/cap-muxer/src/main.rs
+++ b/crates/cap-muxer/src/main.rs
@@ -103,6 +103,7 @@ fn resolve_queue_capacity() -> usize {
struct FrameQueueInner {
frames: VecDeque,
bytes: usize,
+ cancelled: bool,
reader_done: bool,
reader_err: Option,
capacity_bytes: usize,
@@ -120,6 +121,7 @@ impl FrameQueue {
state: Mutex::new(FrameQueueInner {
frames: VecDeque::with_capacity(256),
bytes: 0,
+ cancelled: false,
reader_done: false,
reader_err: None,
capacity_bytes,
@@ -129,19 +131,36 @@ impl FrameQueue {
}
}
- fn push(&self, frame: Frame) {
+ fn push(&self, frame: Frame) -> bool {
let frame_bytes = frame_size_hint(&frame);
let mut guard = self.state.lock().unwrap_or_else(PoisonError::into_inner);
- while guard.bytes + frame_bytes > guard.capacity_bytes && !guard.frames.is_empty() {
+ while !guard.cancelled
+ && guard.bytes + frame_bytes > guard.capacity_bytes
+ && !guard.frames.is_empty()
+ {
guard = self
.space_cv
.wait(guard)
.unwrap_or_else(PoisonError::into_inner);
}
+ if guard.cancelled {
+ return false;
+ }
guard.bytes = guard.bytes.saturating_add(frame_bytes);
guard.frames.push_back(frame);
drop(guard);
self.data_cv.notify_one();
+ true
+ }
+
+ fn cancel(&self) {
+ let mut guard = self.state.lock().unwrap_or_else(PoisonError::into_inner);
+ guard.cancelled = true;
+ guard.frames.clear();
+ guard.bytes = 0;
+ drop(guard);
+ self.data_cv.notify_all();
+ self.space_cv.notify_all();
}
fn mark_reader_done(&self, err: Option) {
@@ -165,6 +184,9 @@ impl FrameQueue {
self.space_cv.notify_one();
return PopResult::Frame(frame);
}
+ if guard.cancelled {
+ return PopResult::Drained;
+ }
if guard.reader_done {
let err = guard.reader_err.take();
return match err {
@@ -203,6 +225,11 @@ fn frame_size_hint(frame: &Frame) -> usize {
}
}
+fn cancel_reader(queue: &FrameQueue, reader_handle: std::thread::JoinHandle<()>) {
+ queue.cancel();
+ drop(reader_handle);
+}
+
fn run() -> Result<(), MuxerError> {
ffmpeg::init().map_err(|e| MuxerError::Init(anyhow::Error::from(e)))?;
@@ -222,7 +249,11 @@ fn run() -> Result<(), MuxerError> {
let mut reader = BufReader::with_capacity(1024 * 1024, stdin.lock());
loop {
match read_frame(&mut reader) {
- Ok(frame) => queue.push(frame),
+ Ok(frame) => {
+ if !queue.push(frame) {
+ return;
+ }
+ }
Err(ProtocolError::Io(ref ioe))
if ioe.kind() == io::ErrorKind::UnexpectedEof =>
{
@@ -265,10 +296,9 @@ fn run() -> Result<(), MuxerError> {
}
}
+ cancel_reader(&queue, reader_handle);
let finish_result = state.finish();
- let _ = reader_handle.join();
-
result?;
finish_result?;
@@ -795,6 +825,8 @@ fn write_ready_packet(
#[cfg(test)]
mod tests {
use super::*;
+ use std::sync::mpsc::{self, RecvTimeoutError};
+ use std::time::Duration;
fn packet(stream_index: u8, pts: i64, dts: i64, duration: u64) -> Packet {
Packet {
@@ -862,4 +894,66 @@ mod tests {
assert_eq!(nominal_audio_duration_input_tb(&init), Some(1_024));
}
+
+ #[test]
+ fn cancellation_releases_producer_blocked_by_full_queue() {
+ let queue = Arc::new(FrameQueue::new(frame_size_hint(&Frame::Finish)));
+ assert!(queue.push(Frame::Finish));
+
+ let (started_tx, started_rx) = mpsc::channel();
+ let (result_tx, result_rx) = mpsc::channel();
+ let producer = {
+ let queue = Arc::clone(&queue);
+ std::thread::spawn(move || {
+ started_tx.send(()).unwrap();
+ result_tx.send(queue.push(Frame::Finish)).unwrap();
+ })
+ };
+
+ started_rx.recv_timeout(Duration::from_secs(1)).unwrap();
+ assert!(matches!(
+ result_rx.recv_timeout(Duration::from_millis(50)),
+ Err(RecvTimeoutError::Timeout)
+ ));
+
+ queue.cancel();
+
+ assert!(!result_rx.recv_timeout(Duration::from_secs(1)).unwrap());
+ producer.join().unwrap();
+ assert!(matches!(queue.pop(), PopResult::Drained));
+ }
+
+ #[test]
+ fn cancellation_rejects_late_pushes() {
+ let queue = FrameQueue::new(frame_size_hint(&Frame::Finish));
+
+ queue.cancel();
+
+ assert!(!queue.push(Frame::Finish));
+ assert!(matches!(queue.pop(), PopResult::Drained));
+ }
+
+ #[test]
+ fn cancellation_does_not_wait_for_reader_input_to_close() {
+ let queue = Arc::new(FrameQueue::new(frame_size_hint(&Frame::Finish)));
+ let (input_closed_tx, input_closed_rx) = mpsc::channel();
+ let reader_handle = std::thread::spawn(move || {
+ let _ = input_closed_rx.recv();
+ });
+ let (shutdown_done_tx, shutdown_done_rx) = mpsc::channel();
+ let shutdown_thread = {
+ let queue = Arc::clone(&queue);
+ std::thread::spawn(move || {
+ cancel_reader(&queue, reader_handle);
+ shutdown_done_tx.send(()).unwrap();
+ })
+ };
+
+ shutdown_done_rx
+ .recv_timeout(Duration::from_secs(1))
+ .unwrap();
+
+ input_closed_tx.send(()).unwrap();
+ shutdown_thread.join().unwrap();
+ }
}
diff --git a/crates/recording/tests/oop_muxer.rs b/crates/recording/tests/oop_muxer.rs
index 0c3302df9d4..ee57f5c3c7a 100644
--- a/crates/recording/tests/oop_muxer.rs
+++ b/crates/recording/tests/oop_muxer.rs
@@ -64,6 +64,79 @@ fn subprocess_spawns_and_finishes_cleanly_without_packets() {
assert_eq!(report.packets_written, 0);
}
+#[test]
+fn subprocess_exits_after_finish_or_abort_without_waiting_for_stdin_eof() {
+ use cap_muxer_protocol::{Frame, write_frame};
+ use std::io::Write;
+ use std::process::{Command, Stdio};
+ use std::time::{Duration, Instant};
+
+ let binary = setup_muxer_binary();
+ for (frame, expected_exit) in [(Frame::Finish, 0), (Frame::Abort("test".into()), 40)] {
+ let mut child = Command::new(&binary)
+ .stdin(Stdio::piped())
+ .stdout(Stdio::null())
+ .stderr(Stdio::null())
+ .spawn()
+ .unwrap();
+ let mut stdin = child.stdin.take().unwrap();
+ write_frame(&mut stdin, &frame).unwrap();
+ stdin.flush().unwrap();
+ let started = Instant::now();
+ let status = loop {
+ if let Some(status) = child.try_wait().unwrap() {
+ break status;
+ }
+ if started.elapsed() >= Duration::from_secs(3) {
+ child.kill().unwrap();
+ child.wait().unwrap();
+ panic!("muxer waited for stdin EOF after receiving {frame:?}");
+ }
+ std::thread::sleep(Duration::from_millis(20));
+ };
+ drop(stdin);
+ assert_eq!(status.code(), Some(expected_exit));
+ }
+}
+
+#[cfg(unix)]
+#[test]
+fn finish_preserves_packets_when_a_slow_muxer_takes_more_than_five_seconds() {
+ use cap_muxer_protocol::{Frame, read_frame};
+ use std::io::Cursor;
+ use std::os::unix::fs::PermissionsExt;
+ use std::time::{Duration, Instant};
+
+ let directory = TempDir::new().unwrap();
+ let binary = directory.path().join("slow-muxer");
+ std::fs::write(&binary, "#!/bin/sh\nsleep 6\ncat > \"$0.received\"\n").unwrap();
+ std::fs::set_permissions(&binary, std::fs::Permissions::from_mode(0o700)).unwrap();
+ let config = minimal_video_config(&directory.path().join("video"), Vec::new());
+ let mut subprocess = MuxerSubprocess::spawn(binary, config, None).unwrap();
+ let payload = vec![42; 128 * 1024];
+ subprocess
+ .write_video_packet(0, 0, 1, true, &payload)
+ .unwrap();
+ let started = Instant::now();
+ let report = subprocess.finish().unwrap();
+ assert_eq!(report.exit_code, Some(0));
+ assert!(started.elapsed() >= Duration::from_secs(5));
+
+ let received = std::fs::read(directory.path().join("slow-muxer.received")).unwrap();
+ let mut reader = Cursor::new(received.as_slice());
+ assert!(matches!(
+ read_frame(&mut reader).unwrap(),
+ Frame::InitVideo(_)
+ ));
+ assert!(matches!(read_frame(&mut reader).unwrap(), Frame::Start(_)));
+ let Frame::Packet(packet) = read_frame(&mut reader).unwrap() else {
+ panic!("expected the buffered video packet");
+ };
+ assert_eq!(packet.data, payload);
+ assert!(matches!(read_frame(&mut reader).unwrap(), Frame::Finish));
+ assert_eq!(reader.position(), received.len() as u64);
+}
+
#[test]
fn subprocess_survives_kill_and_parent_reports_crashed() {
let bin = setup_muxer_binary();
From 2ee22a63f04c787ecaeade49e25e9c11c0fe0cb9 Mon Sep 17 00:00:00 2001
From: Richie McIlroy <33632126+richiemcilroy@users.noreply.github.com>
Date: Wed, 26 Aug 2026 19:31:05 +0100
Subject: [PATCH 2/5] fix: recover fragments missing from recording manifests
---
crates/recording/src/recovery.rs | 86 +++++++++++++++++++++++++++++++-
1 file changed, 84 insertions(+), 2 deletions(-)
diff --git a/crates/recording/src/recovery.rs b/crates/recording/src/recovery.rs
index 13dfc6734d2..1f91bc29cec 100644
--- a/crates/recording/src/recovery.rs
+++ b/crates/recording/src/recovery.rs
@@ -371,7 +371,7 @@ impl RecoveryManager {
f.get("file_size").and_then(|s| s.as_u64())
};
- let result: Vec = entries
+ let mut result: Vec = entries
.iter()
.filter(|f| {
f.get("is_complete")
@@ -430,6 +430,29 @@ impl RecoveryManager {
})
.collect();
+ if manifest_type == "m4s_segments" && init_segment.is_some() {
+ let listed: std::collections::HashSet<_> = entries
+ .iter()
+ .filter_map(|entry| entry.get("path").and_then(|path| path.as_str()))
+ .map(|path| dir.join(path))
+ .collect();
+ result.extend(Self::probe_m4s_fragments_with_init(dir).into_iter().filter(
+ |path| {
+ !listed.contains(path)
+ && Self::m4s_fragment_index(path).is_some()
+ && path
+ .symlink_metadata()
+ .is_ok_and(|metadata| metadata.is_file())
+ && tail_is_complete(path).unwrap_or(false)
+ },
+ ));
+ result.sort_by(|a, b| {
+ Self::m4s_fragment_index(a)
+ .cmp(&Self::m4s_fragment_index(b))
+ .then_with(|| a.cmp(b))
+ });
+ }
+
if !result.is_empty() {
return FragmentsInfo {
fragments: result,
@@ -693,6 +716,15 @@ impl RecoveryManager {
fragments
}
+ fn m4s_fragment_index(path: &Path) -> Option {
+ path.file_name()?
+ .to_str()?
+ .strip_prefix("segment_")?
+ .strip_suffix(".m4s")?
+ .parse()
+ .ok()
+ }
+
fn probe_single_file(path: &Path) -> Option {
if !path.exists() {
return None;
@@ -1666,10 +1698,60 @@ fn replace_file(src: &Path, dst: &Path) -> Result<(), RecoveryError> {
#[cfg(test)]
mod tests {
- use super::{replace_file, start_time_or_display_fallback, valid_recovered_audio};
+ use super::{
+ RecoveryManager, replace_file, start_time_or_display_fallback, valid_recovered_audio,
+ };
use std::fs;
use tempfile::tempdir;
+ #[test]
+ fn recovery_includes_complete_fragments_missing_from_the_last_manifest() {
+ let dir = tempdir().unwrap();
+ let mut fragment = Vec::new();
+ for name in [b"moof", b"mdat"] {
+ fragment.extend_from_slice(&72u32.to_be_bytes());
+ fragment.extend_from_slice(name);
+ fragment.extend_from_slice(&[0; 64]);
+ }
+ fs::write(dir.path().join("init.mp4"), [0; 128]).unwrap();
+ for name in [
+ "segment_001.m4s",
+ "segment_002.m4s",
+ "segment_999.m4s",
+ "segment_1000.m4s",
+ ] {
+ fs::write(dir.path().join(name), &fragment).unwrap();
+ }
+ fs::write(
+ dir.path().join("segment_1001.m4s"),
+ &fragment[..fragment.len() - 1],
+ )
+ .unwrap();
+ let manifest = serde_json::to_vec(&serde_json::json!({
+ "version": 5,
+ "type": "m4s_segments",
+ "init_segment": "init.mp4",
+ "segments": [
+ {"path": "segment_001.m4s", "is_complete": true, "file_size": fragment.len()},
+ {"path": "segment_002.m4s", "is_complete": true, "file_size": fragment.len() + 1}
+ ]
+ }))
+ .unwrap();
+ fs::write(dir.path().join("manifest.json"), &manifest).unwrap();
+
+ let recovered = RecoveryManager::find_complete_fragments_with_init(dir.path());
+ assert_eq!(
+ recovered.fragments,
+ ["segment_001.m4s", "segment_999.m4s", "segment_1000.m4s"]
+ .map(|name| dir.path().join(name))
+ );
+ assert_eq!(
+ fs::read(dir.path().join("manifest.json")).unwrap(),
+ manifest
+ );
+ assert!(dir.path().join("segment_1001.m4s").is_file());
+ }
+
#[test]
fn replace_file_overwrites_existing_destination() {
let dir = tempdir().unwrap();
From 2cd5af796e8ef36d0ce7a7934d5d34508f2ea62d Mon Sep 17 00:00:00 2001
From: Richie McIlroy <33632126+richiemcilroy@users.noreply.github.com>
Date: Wed, 26 Aug 2026 19:31:05 +0100
Subject: [PATCH 3/5] fix: warn about low storage and preserve recordings
---
apps/desktop-gpui/src/controls_window.rs | 14 +
apps/desktop-gpui/src/deeplink.rs | 30 +-
apps/desktop-gpui/src/library.rs | 140 ++++-
apps/desktop-gpui/src/main_window.rs | 153 +++++-
apps/desktop-gpui/src/recording.rs | 659 +++++++++++++++++------
apps/desktop-gpui/src/session.rs | 113 +++-
apps/desktop-gpui/src/upload.rs | 77 ++-
crates/utils/src/disk_space.rs | 129 +++++
8 files changed, 1087 insertions(+), 228 deletions(-)
diff --git a/apps/desktop-gpui/src/controls_window.rs b/apps/desktop-gpui/src/controls_window.rs
index 50f6af5ef2c..cf21f136c06 100644
--- a/apps/desktop-gpui/src/controls_window.rs
+++ b/apps/desktop-gpui/src/controls_window.rs
@@ -114,6 +114,8 @@ impl ControlsWindow {
let stopping = session.phase == Phase::Stopping;
let label: SharedString = if starting {
"Starting".into()
+ } else if stopping {
+ "Saving…".into()
} else {
Self::format_elapsed(session.elapsed()).into()
};
@@ -373,6 +375,18 @@ impl Render for ControlsWindow {
.font_family("Geist")
// `body { font-weight: 500 }` (`ui-solid/src/main.css:189-192`).
.font_weight(FontWeight::MEDIUM)
+ .when(self.session.read(cx).storage_warning, |this| {
+ this.child(
+ div()
+ .mb(px(8.))
+ .rounded(px(8.))
+ .bg(self.theme.red_2)
+ .p(px(8.))
+ .text_size(px(11.))
+ .text_color(self.theme.red_11)
+ .child("Low storage. Cap will stop soon to save your recording."),
+ )
+ })
.child(self.render_bar(cx))
}
}
diff --git a/apps/desktop-gpui/src/deeplink.rs b/apps/desktop-gpui/src/deeplink.rs
index a25a98299b6..ce119e51bfc 100644
--- a/apps/desktop-gpui/src/deeplink.rs
+++ b/apps/desktop-gpui/src/deeplink.rs
@@ -369,19 +369,23 @@ impl DeepLinkAction {
let feeds = feeds.read(cx);
(feeds.camera_actor(), feeds.mic_actor())
};
- app_windows::begin_recording(
- recording::StartConfig {
- mode,
- target: capture_target,
- microphone: mic_label,
- camera,
- system_audio: capture_system_audio,
- excluded_windows: Vec::new(),
- camera_feed,
- mic_feed,
- },
- cx,
- );
+ let main = cx.global::().main;
+ main.update(cx, |view, _, cx| {
+ view.start_recording_config(
+ recording::StartConfig {
+ mode,
+ target: capture_target,
+ microphone: mic_label,
+ camera,
+ system_audio: capture_system_audio,
+ excluded_windows: Vec::new(),
+ camera_feed,
+ mic_feed,
+ },
+ cx,
+ )
+ })
+ .ok();
});
Ok(())
}
diff --git a/apps/desktop-gpui/src/library.rs b/apps/desktop-gpui/src/library.rs
index e98b5f218ff..dd53165448f 100644
--- a/apps/desktop-gpui/src/library.rs
+++ b/apps/desktop-gpui/src/library.rs
@@ -300,10 +300,7 @@ impl RecordingItem {
/// there are no uploads in this app, so `MultipartUpload` /
/// `SinglePartUpload` can never be the reason to keep polling.
pub fn is_active(&self) -> bool {
- matches!(
- self.status,
- RecordingStatus::InProgress | RecordingStatus::NeedsRemux
- )
+ self.status == RecordingStatus::InProgress
}
/// `studioCompleteCheck()`: the only rows whose whole body is clickable.
@@ -339,7 +336,10 @@ fn recording_item(path: PathBuf, meta: RecordingMeta, sort_time_millis: f64) ->
},
StudioRecordingMeta::SingleSegment { .. } => RecordingStatus::Complete,
},
- RecordingMetaInner::Instant(InstantRecordingMeta::InProgress { .. }) => {
+ RecordingMetaInner::Instant(InstantRecordingMeta::InProgress { recording: false }) => {
+ RecordingStatus::NeedsRemux
+ }
+ RecordingMetaInner::Instant(InstantRecordingMeta::InProgress { recording: true }) => {
RecordingStatus::InProgress
}
RecordingMetaInner::Instant(InstantRecordingMeta::Failed { error }) => {
@@ -411,15 +411,35 @@ pub fn find_incomplete_recordings_in(
list_recordings_in(dirs)
.into_iter()
.filter(|item| {
- item.mode == RecordingMode::Studio
- && matches!(
- item.status,
- RecordingStatus::InProgress | RecordingStatus::NeedsRemux
- )
- && active_recording != Some(item.path.as_path())
- && is_recording_after_recovery_cutoff(&item.pretty_name)
+ matches!(
+ item.status,
+ RecordingStatus::InProgress | RecordingStatus::NeedsRemux
+ ) && active_recording != Some(item.path.as_path())
+ && is_recording_after_recovery_cutoff(&item.pretty_name, item.sort_time_millis)
})
.filter_map(|item| {
+ if item.mode == RecordingMode::Instant {
+ let display = item.path.join("content/display");
+ if !display.join("init.mp4").is_file() {
+ return None;
+ }
+ let segment_count = std::fs::read_dir(display)
+ .ok()?
+ .filter_map(Result::ok)
+ .filter(|entry| {
+ entry
+ .path()
+ .extension()
+ .is_some_and(|extension| extension == "m4s")
+ })
+ .count();
+ return (segment_count > 0).then_some(IncompleteRecordingItem {
+ project_path: item.path,
+ pretty_name: item.pretty_name,
+ segment_count,
+ estimated_duration_secs: 0.0,
+ });
+ }
let incomplete = RecoveryManager::inspect_recording(&item.path)?;
(!incomplete.recoverable_segments.is_empty()).then_some(IncompleteRecordingItem {
project_path: item.path,
@@ -435,11 +455,15 @@ pub fn find_incomplete_recordings() -> Vec {
find_incomplete_recordings_in(&known_recordings_dirs(), None)
}
-fn is_recording_after_recovery_cutoff(pretty_name: &str) -> bool {
+fn is_recording_after_recovery_cutoff(pretty_name: &str, sort_time_millis: f64) -> bool {
let Some(date) = pretty_name
.strip_prefix("Cap ")
.and_then(|name| name.split(" at ").next())
.and_then(|value| chrono::NaiveDate::parse_from_str(value, "%Y-%m-%d").ok())
+ .or_else(|| {
+ chrono::DateTime::from_timestamp_millis(sort_time_millis as i64)
+ .map(|timestamp| timestamp.date_naive())
+ })
else {
return false;
};
@@ -475,6 +499,10 @@ fn recover_incomplete_recording_in(
let meta = RecordingMeta::load_for_project(&canonical_path)
.map_err(|error| format!("Failed to load recording metadata: {error}"))?;
+ if matches!(meta.inner, RecordingMetaInner::Instant(_)) {
+ return crate::recording::recover_instant_recording(&canonical_path)
+ .map_err(|error| format!("Could not save the instant recording: {error:#}"));
+ }
let Some(studio) = meta.studio_meta() else {
return Err("Only incomplete studio recordings can be recovered".to_string());
};
@@ -487,6 +515,8 @@ fn recover_incomplete_recording_in(
let incomplete = RecoveryManager::inspect_recording(&canonical_path)
.ok_or_else(|| "No recoverable segments found".to_string())?;
+ crate::recording::ensure_finalization_storage(&canonical_path)
+ .map_err(|error| format!("{error:#}"))?;
let recovered = RecoveryManager::recover(&incomplete)
.map_err(|error| format!("Failed to recover recording: {error}"))?;
let display = match &recovered.meta {
@@ -1470,7 +1500,10 @@ mod tests {
let remux = by_name("Studio remux");
assert_eq!(remux.status, RecordingStatus::NeedsRemux);
- assert!(remux.is_active(), "NeedsRemux keeps the 2s poll running");
+ assert!(
+ !remux.is_active(),
+ "deferred finalization does not keep polling"
+ );
let complete = by_name("Instant complete");
assert_eq!(complete.mode, RecordingMode::Instant);
@@ -1601,6 +1634,38 @@ mod tests {
std::fs::remove_dir_all(root).unwrap();
}
+ #[test]
+ fn interrupted_instant_recordings_are_recoverable_but_active_recordings_are_excluded() {
+ for recording in [false, true] {
+ let root = temp_dir("instant-storage-recovery");
+ let metadata = serde_json::json!({
+ "pretty_name": "Custom instant recording",
+ "sharing": null,
+ "recording": recording,
+ })
+ .to_string();
+ let bundle = write_bundle(&root, "deferred-instant", &metadata);
+ let display = bundle.join("content/display");
+ std::fs::create_dir_all(&display).unwrap();
+ std::fs::write(display.join("init.mp4"), [0; 8]).unwrap();
+ std::fs::write(display.join("segment_001.m4s"), [0; 16]).unwrap();
+ let items = list_recordings_in(std::slice::from_ref(&root));
+ if !recording {
+ assert_eq!(items[0].status, RecordingStatus::NeedsRemux);
+ assert!(!items[0].is_active());
+ }
+ let recoverable = find_incomplete_recordings_in(std::slice::from_ref(&root), None);
+ assert_eq!(recoverable.len(), 1);
+ assert_eq!(recoverable[0].project_path, bundle);
+ assert!(
+ find_incomplete_recordings_in(std::slice::from_ref(&root), Some(&bundle))
+ .is_empty()
+ );
+ assert!(display.join("segment_001.m4s").is_file());
+ std::fs::remove_dir_all(root).unwrap();
+ }
+ }
+
#[test]
fn recovery_rejects_paths_outside_the_recording_library() {
let root = temp_dir("recovery-path");
@@ -1646,24 +1711,49 @@ mod tests {
}
}
+ let mut interrupted = RecordingMeta::load_for_project(©).unwrap();
+ if matches!(interrupted.inner, RecordingMetaInner::Instant(_)) {
+ interrupted.inner =
+ RecordingMetaInner::Instant(InstantRecordingMeta::InProgress { recording: true });
+ interrupted.save_for_project().unwrap();
+ }
+
let dirs = std::slice::from_ref(&recordings);
let found = find_incomplete_recordings_in(dirs, None);
assert_eq!(found.len(), 1);
- assert_eq!(found[0].segment_count, 1);
+ assert!(found[0].segment_count > 0);
let recovered = recover_incomplete_recording_in(dirs, ©).unwrap();
let meta = RecordingMeta::load_for_project(&recovered).unwrap();
- assert!(matches!(
- meta.studio_meta().unwrap().status(),
- StudioRecordingStatus::Complete
- ));
- assert!(
- recovered
- .join("content/segments/segment-0/display.mp4")
- .is_file()
- );
- assert!(recovered.join("project-config.json").is_file());
+ match &meta.inner {
+ RecordingMetaInner::Studio(studio) => {
+ assert!(matches!(studio.status(), StudioRecordingStatus::Complete));
+ assert!(
+ recovered
+ .join("content/segments/segment-0/display.mp4")
+ .is_file()
+ );
+ assert!(recovered.join("project-config.json").is_file());
+ }
+ RecordingMetaInner::Instant(InstantRecordingMeta::Complete { fps, sample_rate }) => {
+ assert_eq!(*fps, 30);
+ assert_eq!(*sample_rate, Some(48_000));
+ let output = recovered.join("content/output.mp4");
+ assert!(output.is_file());
+ assert!(ffmpeg::format::input(&output).unwrap().duration() >= 2_900_000);
+ assert!(recovered.join("content/display/segment_001.m4s").is_file());
+ assert!(recovered.join("content/audio/segment_001.m4s").is_file());
+ let status = std::process::Command::new("ffmpeg")
+ .args(["-v", "error", "-xerror", "-i"])
+ .arg(output)
+ .args(["-f", "null", "-"])
+ .status()
+ .unwrap();
+ assert!(status.success());
+ }
+ _ => panic!("Recovery left incomplete metadata"),
+ }
assert!(bundle_thumbnail_path(&recovered).is_file());
assert!(find_incomplete_recordings_in(dirs, None).is_empty());
diff --git a/apps/desktop-gpui/src/main_window.rs b/apps/desktop-gpui/src/main_window.rs
index a1447b90c94..1b1344d7ae2 100644
--- a/apps/desktop-gpui/src/main_window.rs
+++ b/apps/desktop-gpui/src/main_window.rs
@@ -330,6 +330,7 @@ pub struct MainWindow {
/// The app-wide recording session; the lifecycle itself lives there so the
/// controls bar window can drive the same recording.
session: Entity,
+ checking_storage: bool,
/// The Recents scan, or `None` while the first one is in flight -- which
/// is the query's `isLoading`, and draws the same three skeleton cards.
recents: Option>,
@@ -405,7 +406,38 @@ impl MainWindow {
) -> Self {
crate::theme::bind_window(window, cx);
let theme = Theme::for_window(window, cx, true);
- cx.observe(&session, |_, _, cx| cx.notify()).detach();
+ let mut previous_phase = Phase::Idle;
+ cx.observe_in(&session, window, move |this, session, window, cx| {
+ let phase = session.read(cx).phase;
+ if phase == Phase::Idle && previous_phase != Phase::Idle {
+ this.scan_incomplete_recordings(window, cx, std::time::Duration::ZERO);
+ if let Some(notice) = session.read(cx).storage_notice.clone() {
+ let main = cx.global::().main;
+ cx.spawn(async move |_, cx| {
+ let receiver = cx.update(|cx| {
+ app_windows::show_main_window(cx);
+ cx.activate(true);
+ main.update(cx, |_, window, cx| {
+ window.prompt(
+ gpui::PromptLevel::Warning,
+ "Low storage",
+ Some(¬ice),
+ &[gpui::PromptButton::cancel("OK")],
+ cx,
+ )
+ })
+ });
+ if let Ok(receiver) = receiver {
+ let _ = receiver.await;
+ }
+ })
+ .detach();
+ }
+ }
+ previous_phase = phase;
+ cx.notify();
+ })
+ .detach();
// Track the app-scoped feeds: the camera bubble's close button
// deselects the camera there, and this window's selection has to
@@ -457,6 +489,7 @@ impl MainWindow {
_search_events: search_events,
enumerating: true,
session,
+ checking_storage: false,
recents: None,
recents_task: None,
library: None,
@@ -1010,6 +1043,17 @@ impl MainWindow {
let Some(recording) = self.incomplete_recording.clone() else {
return;
};
+ if !recover
+ && !crate::platform::confirm_dialog(
+ "Cap",
+ "Are you sure you want to delete this recording?",
+ "Yes",
+ "No",
+ false,
+ )
+ {
+ return;
+ }
self.recovery_pending = true;
self.recovery_error = None;
@@ -1039,7 +1083,11 @@ impl MainWindow {
this.refresh_open_library(window, cx);
this.scan_incomplete_recordings(window, cx, std::time::Duration::ZERO);
if recover {
- cx.defer(move |cx| app_windows::open_editor(project_path, cx));
+ if project_path.join("content/output.mp4").is_file() {
+ cx.reveal_path(&project_path.join("content/output.mp4"));
+ } else {
+ cx.defer(move |cx| app_windows::open_editor(project_path, cx));
+ }
}
}
Err(error) => {
@@ -1620,7 +1668,7 @@ impl MainWindow {
_window: &mut Window,
cx: &mut Context,
) {
- if self.session.read(cx).phase != Phase::Idle {
+ if self.session.read(cx).phase != Phase::Idle || self.checking_storage {
return;
}
// An armed editor recording target forces Studio before this window's
@@ -1670,7 +1718,104 @@ impl MainWindow {
mic_feed,
};
- cx.defer(move |cx: &mut gpui::App| app_windows::begin_recording(config, cx));
+ self.start_recording_config(config, cx);
+ }
+
+ pub(crate) fn start_recording_config(
+ &mut self,
+ config: recording::StartConfig,
+ cx: &mut Context,
+ ) {
+ self.check_storage_before_start(config, false, cx);
+ }
+
+ fn check_storage_before_start(
+ &mut self,
+ config: recording::StartConfig,
+ acknowledged: bool,
+ cx: &mut Context,
+ ) {
+ if self.checking_storage || self.session.read(cx).phase != Phase::Idle {
+ return;
+ }
+ self.checking_storage = true;
+ let main = cx.global::().main;
+ cx.spawn(async move |this, cx| {
+ let result = cx.background_executor().spawn(async {
+ recording::available_recording_storage()
+ }).await;
+ if !this.update(cx, |this, cx| {
+ if this.session.read(cx).phase != Phase::Idle {
+ this.checking_storage = false;
+ return false;
+ }
+ true
+ }).unwrap_or(false) {
+ return;
+ }
+ let storage = match result {
+ Ok(storage) => storage,
+ Err(error) => {
+ this.update(cx, |this, cx| {
+ this.checking_storage = false;
+ this.session.update(cx, |session, cx| {
+ session.error = Some(format!("Could not check recording storage: {error}"));
+ cx.notify();
+ });
+ cx.defer(app_windows::show_main_window);
+ if this.session.read(cx).editor_recording_target().is_some() {
+ cx.defer(app_windows::abort_editor_recording_flow);
+ }
+ }).ok();
+ return;
+ }
+ };
+ let can_start = storage.status() != cap_utils::disk_space::DiskSpaceStatus::Exhausted;
+ if storage.status() == cap_utils::disk_space::DiskSpaceStatus::Ok || acknowledged && can_start {
+ this.update(cx, |this, cx| {
+ this.checking_storage = false;
+ cx.defer(move |cx| app_windows::begin_recording(config, cx));
+ }).ok();
+ return;
+ }
+ let available = storage.available_bytes as f64 / 1_073_741_824.0;
+ let detail = if can_start {
+ format!("Only {available:.2} GB is available on your recording drive. Cap will stop automatically if storage gets too low, preserving your recording.")
+ } else {
+ format!("Only {available:.2} GB is available on your recording drive. Free up space so at least 512 MB is available before recording.")
+ };
+ let buttons = if can_start {
+ vec![gpui::PromptButton::ok("Record anyway"), gpui::PromptButton::cancel("Go back")]
+ } else {
+ vec![gpui::PromptButton::cancel("OK")]
+ };
+ let receiver = cx.update(|cx| {
+ if RecordingSession::global(cx).read(cx).phase != Phase::Idle {
+ return Err(anyhow::anyhow!("A recording has already started."));
+ }
+ app_windows::show_main_window(cx);
+ cx.activate(true);
+ main.update(cx, |_, window, cx| {
+ window.prompt(gpui::PromptLevel::Warning, "Low storage", Some(&detail), &buttons, cx)
+ })
+ });
+ let confirmed = match receiver {
+ Ok(receiver) => receiver.await == Ok(0) && can_start,
+ Err(_) => false,
+ };
+ this.update(cx, |this, cx| {
+ this.checking_storage = false;
+ if this.session.read(cx).phase != Phase::Idle {
+ return;
+ }
+ if confirmed {
+ this.check_storage_before_start(config, true, cx);
+ } else if this.session.read(cx).editor_recording_target().is_some() {
+ cx.defer(app_windows::abort_editor_recording_flow);
+ }
+ cx.notify();
+ }).ok();
+ }).detach();
}
/// `await commands.focusWindow(target.id)` at the end of
diff --git a/apps/desktop-gpui/src/recording.rs b/apps/desktop-gpui/src/recording.rs
index 1bb497e37d4..a8f647b86cd 100644
--- a/apps/desktop-gpui/src/recording.rs
+++ b/apps/desktop-gpui/src/recording.rs
@@ -7,8 +7,9 @@
//! runs on the tokio runtime (`gpui_tokio`), never on gpui's main thread --
//! kameo actors and the capture pipeline both assume tokio.
+use std::io::Write as _;
use std::path::PathBuf;
-use std::sync::Arc;
+use std::sync::{Arc, Mutex};
use anyhow::{Context as _, anyhow};
use cap_recording::{
@@ -144,14 +145,22 @@ impl ActiveRecording {
Ok(())
}
- pub async fn stop(self) -> anyhow::Result {
+ pub async fn stop(self, preserve_local: bool) -> anyhow::Result {
let mut instant_upload = self.instant_upload;
match self.handle {
Handle::Studio(handle) => {
let completed = handle.stop().await?;
let project_path = completed.project_path.clone();
+ let needs_remux = matches!(
+ completed.meta.status(),
+ cap_project::StudioRecordingStatus::NeedsRemux
+ );
tokio::task::spawn_blocking(move || {
+ if needs_remux {
+ ensure_finalization_storage(&project_path)?;
+ }
cap_recording::recovery::RecoveryManager::remux_if_needed(&project_path)
+ .map_err(anyhow::Error::from)
})
.await
.context("studio finalize task")?
@@ -175,102 +184,251 @@ impl ActiveRecording {
Ok(completed.project_path)
}
Handle::Instant(handle) => {
- let completed = handle.stop().await?;
- let project_path = completed.project_path.clone();
- let upload = instant_upload
- .as_mut()
- .ok_or_else(|| anyhow!("instant recording has no upload session"))?;
- let segmented = upload.is_segmented();
- let segment_upload_result = upload.finish_segments().await;
-
- let display_dir = project_path.join("content/display");
- let audio_dir = project_path.join("content/audio");
- let output_path = project_path.join("content/output.mp4");
- if display_dir.is_dir() {
- let muxed = output_path.clone();
+ let result = async {
+ let stopped = handle.stop().await;
+ let metadata_result = match instant_upload.as_ref() {
+ Some(upload) => {
+ mark_instant_recording_stopped(&self.project_dir, upload.metadata_lock())
+ }
+ None => Err(anyhow!("instant recording has no upload session")),
+ };
+ if let Err(error) = &metadata_result {
+ tracing::warn!(%error, "Could not mark the instant recording as stopped");
+ }
+ let completed = stopped?;
+ metadata_result?;
+ let project_path = completed.project_path.clone();
+ let upload = instant_upload
+ .as_mut()
+ .ok_or_else(|| anyhow!("instant recording has no upload session"))?;
+ let segmented = upload.is_segmented();
+
+ let display_dir = project_path.join("content/display");
+ let audio_dir = project_path.join("content/audio");
+ let output_path = project_path.join("content/output.mp4");
+ if display_dir.is_dir() {
+ let muxed = output_path.clone();
+ let project_path = project_path.clone();
+ tokio::task::spawn_blocking(move || {
+ ensure_finalization_storage(&project_path)?;
+ cap_recording::recovery::RecoveryManager::finalize_instant_output(
+ &display_dir,
+ &audio_dir,
+ &muxed,
+ )
+ .map_err(anyhow::Error::from)
+ })
+ .await
+ .context("instant finalize task")?
+ .context("instant finalize")?;
+ } else if !output_path.is_file() {
+ return Err(anyhow!("instant recording has no finalized output"));
+ }
+
+ persist_instant_meta(&completed, upload.video(), upload.metadata_lock())?;
+
+ // The Tauri app builds the instant thumbnail by concatenating
+ // `content/display`'s init segment with the first media
+ // segment (`create_screenshot_source_from_segments`); by this
+ // point `finalize_instant_output` has already muxed the whole
+ // thing into `content/output.mp4`, which is the same first
+ // frame without the temporary file. The blur bridge is *not*
+ // applied here: `project_config_from_recording` is the studio
+ // arm of `handle_recording_finish` only.
+ let project_path = completed.project_path.clone();
tokio::task::spawn_blocking(move || {
- cap_recording::recovery::RecoveryManager::finalize_instant_output(
- &display_dir,
- &audio_dir,
- &muxed,
- )
+ write_bundle_thumbnail(&project_path, &output_path);
})
.await
- .context("instant finalize task")?
- .context("instant finalize")?;
- } else if !output_path.is_file() {
- return Err(anyhow!("instant recording has no finalized output"));
- }
+ .context("instant thumbnail task")?;
+
+ if let Err(error) = upload.finish_segments().await {
+ persist_instant_upload_failure(
+ &completed.project_path,
+ &error,
+ upload.metadata_lock(),
+ )?;
+ return Err(anyhow!(error));
+ }
- persist_instant_meta(&completed, upload.video())?;
+ let upload_result = if segmented {
+ upload.finish_screenshot(&completed.project_path).await
+ } else {
+ crate::upload::upload_exported_video(
+ completed.project_path.clone(),
+ None,
+ |_| {},
+ Arc::new(std::sync::atomic::AtomicBool::new(false)),
+ )
+ .await
+ .and_then(|result| match result {
+ crate::upload::UploadResult::Success(_) => Ok(()),
+ crate::upload::UploadResult::NotAuthenticated => Err(
+ "Your session has expired. Please sign in again to upload this recording."
+ .to_string(),
+ ),
+ crate::upload::UploadResult::UpgradeRequired => {
+ Err("Instant recording requires an upgraded plan.".to_string())
+ }
+ })
+ };
+ if let Err(error) = upload_result {
+ persist_instant_upload_failure(
+ &completed.project_path,
+ &error,
+ upload.metadata_lock(),
+ )?;
+ return Err(anyhow!(error));
+ }
- // The Tauri app builds the instant thumbnail by concatenating
- // `content/display`'s init segment with the first media
- // segment (`create_screenshot_source_from_segments`); by this
- // point `finalize_instant_output` has already muxed the whole
- // thing into `content/output.mp4`, which is the same first
- // frame without the temporary file. The blur bridge is *not*
- // applied here: `project_config_from_recording` is the studio
- // arm of `handle_recording_finish` only.
- let project_path = completed.project_path.clone();
- tokio::task::spawn_blocking(move || {
- write_bundle_thumbnail(&project_path, &output_path);
- })
- .await
- .context("instant thumbnail task")?;
+ persist_instant_upload_complete(&completed.project_path, upload.metadata_lock())?;
- if let Err(error) = segment_upload_result {
- persist_instant_upload_failure(&completed.project_path, &error)?;
- return Err(anyhow!(error));
- }
+ if !preserve_local
+ && crate::store::GeneralSettings::load().delete_instant_recordings_after_upload
+ {
+ let directory = completed.project_path.clone();
+ tokio::task::spawn_blocking(move || std::fs::remove_dir_all(directory))
+ .await
+ .context("instant upload cleanup task")?
+ .context("deleting uploaded instant recording")?;
+ }
- let upload_result = if segmented {
- upload.finish_screenshot(&completed.project_path).await
- } else {
- crate::upload::upload_exported_video(
- completed.project_path.clone(),
- None,
- |_| {},
- Arc::new(std::sync::atomic::AtomicBool::new(false)),
- )
- .await
- .and_then(|result| match result {
- crate::upload::UploadResult::Success(_) => Ok(()),
- crate::upload::UploadResult::NotAuthenticated => Err(
- "Your session has expired. Please sign in again to upload this recording."
- .to_string(),
- ),
- crate::upload::UploadResult::UpgradeRequired => {
- Err("Instant recording requires an upgraded plan.".to_string())
- }
- })
- };
- if let Err(error) = upload_result {
- persist_instant_upload_failure(&completed.project_path, &error)?;
- return Err(anyhow!(error));
+ Ok(completed.project_path)
}
-
- let mut meta =
- cap_project::RecordingMeta::load_for_project(&completed.project_path)
- .map_err(|error| anyhow!("loading instant recording metadata: {error}"))?;
- meta.upload = Some(cap_project::UploadMeta::Complete);
- meta.save_for_project()
- .map_err(|error| anyhow!("saving completed instant upload: {error}"))?;
-
- if crate::store::GeneralSettings::load().delete_instant_recordings_after_upload {
- let directory = completed.project_path.clone();
- tokio::task::spawn_blocking(move || std::fs::remove_dir_all(directory))
- .await
- .context("instant upload cleanup task")?
- .context("deleting uploaded instant recording")?;
+ .await;
+ if result.is_err()
+ && let Some(upload) = instant_upload.as_mut()
+ {
+ upload.abort_segments().await;
}
-
- Ok(completed.project_path)
+ result
}
}
}
}
+fn with_instant_metadata_lock(
+ metadata_lock: &Mutex<()>,
+ update: impl FnOnce() -> anyhow::Result,
+) -> anyhow::Result {
+ let _guard = metadata_lock
+ .lock()
+ .unwrap_or_else(|error| error.into_inner());
+ update()
+}
+
+fn save_instant_metadata(meta: &cap_project::RecordingMeta) -> anyhow::Result<()> {
+ let contents = serde_json::to_vec_pretty(meta)?;
+ write_instant_metadata(&meta.project_path, |file| file.write_all(&contents))
+}
+
+fn write_instant_metadata(
+ project_path: &std::path::Path,
+ write: impl FnOnce(&mut std::fs::File) -> std::io::Result<()>,
+) -> anyhow::Result<()> {
+ let temporary = project_path.join(format!(
+ ".recording-meta-{}.tmp",
+ crate::store::new_uuid_v4()
+ ));
+ let mut file = std::fs::OpenOptions::new()
+ .write(true)
+ .create_new(true)
+ .open(&temporary)?;
+ let result = (|| {
+ write(&mut file)?;
+ file.sync_all()?;
+ drop(file);
+ std::fs::rename(&temporary, project_path.join("recording-meta.json"))
+ })();
+ if result.is_err() {
+ let _ = std::fs::remove_file(&temporary);
+ }
+ result.context("saving instant recording metadata")
+}
+
+fn mark_instant_recording_stopped(
+ project_path: &std::path::Path,
+ metadata_lock: &Mutex<()>,
+) -> anyhow::Result<()> {
+ with_instant_metadata_lock(metadata_lock, || {
+ let mut meta = cap_project::RecordingMeta::load_for_project(project_path)
+ .map_err(|error| anyhow!("loading recording metadata: {error}"))?;
+ meta.inner = cap_project::RecordingMetaInner::Instant(
+ cap_project::InstantRecordingMeta::InProgress { recording: false },
+ );
+ save_instant_metadata(&meta)?;
+ Ok(())
+ })
+}
+
+pub fn available_recording_storage() -> std::io::Result {
+ Ok(cap_utils::disk_space::RecordingStorage {
+ available_bytes: cap_utils::disk_space::free_bytes_for_path(&recordings_dir())?,
+ recording_bytes: 0,
+ })
+}
+
+pub(crate) fn ensure_finalization_storage(project_path: &std::path::Path) -> anyhow::Result<()> {
+ let storage = cap_utils::disk_space::recording_storage(project_path)
+ .context("checking storage before saving the recording")?;
+ if !storage.can_finalize() {
+ return Err(anyhow!(
+ "Low storage. Your recording files are preserved at {}. Free up space, then recover the recording in Cap.",
+ project_path.display()
+ ));
+ }
+ Ok(())
+}
+
+pub fn recover_instant_recording(project_path: &std::path::Path) -> anyhow::Result {
+ use cap_project::{InstantRecordingMeta, RecordingMeta, RecordingMetaInner};
+
+ let mut meta = RecordingMeta::load_for_project(project_path)
+ .map_err(|error| anyhow!("loading recording metadata: {error}"))?;
+ if !matches!(
+ meta.inner,
+ RecordingMetaInner::Instant(InstantRecordingMeta::InProgress { .. })
+ ) {
+ return Err(anyhow!(
+ "This instant recording is not waiting to be saved."
+ ));
+ }
+ ensure_finalization_storage(project_path)?;
+ let output = project_path.join("content/output.mp4");
+ cap_recording::recovery::RecoveryManager::finalize_instant_output(
+ &project_path.join("content/display"),
+ &project_path.join("content/audio"),
+ &output,
+ )?;
+ let input = ffmpeg::format::input(&output)?;
+ let video = input
+ .streams()
+ .best(ffmpeg::media::Type::Video)
+ .ok_or_else(|| anyhow!("The recovered recording has no video track."))?;
+ let fps = f64::from(video.avg_frame_rate());
+ if !fps.is_finite() || fps <= 0.0 {
+ return Err(anyhow!(
+ "The recovered recording has an invalid frame rate."
+ ));
+ }
+ let sample_rate = input
+ .streams()
+ .best(ffmpeg::media::Type::Audio)
+ .map(|audio| {
+ ffmpeg::codec::context::Context::from_parameters(audio.parameters())
+ .and_then(|context| context.decoder().audio())
+ .map(|decoder| decoder.rate())
+ })
+ .transpose()?;
+ meta.inner = RecordingMetaInner::Instant(InstantRecordingMeta::Complete {
+ fps: fps.round() as u32,
+ sample_rate,
+ });
+ save_instant_metadata(&meta)?;
+ write_bundle_thumbnail(project_path, &output);
+ Ok(project_path.to_path_buf())
+}
+
/// The first segment's display track, which is what
/// `handle_recording_finish` hands `create_screenshot`
/// (`apps/desktop/src-tauri/src/recording.rs:3415-3429`).
@@ -433,103 +591,124 @@ pub fn apply_camera_blur_to_project_config(
fn persist_instant_meta(
completed: &instant_recording::CompletedRecording,
upload: &cap_project::VideoUploadInfo,
+ metadata_lock: &Mutex<()>,
) -> anyhow::Result<()> {
- use cap_project::{
- InstantRecordingMeta, Platform, ProjectConfiguration, RecordingMeta, RecordingMetaInner,
- };
+ with_instant_metadata_lock(metadata_lock, || {
+ use cap_project::{
+ InstantRecordingMeta, Platform, ProjectConfiguration, RecordingMeta, RecordingMetaInner,
+ };
- let pretty_name = completed
- .project_path
- .file_stem()
- .and_then(|name| name.to_str())
- .filter(|name| !name.is_empty())
- .unwrap_or("Cap Recording")
- .to_string();
- let meta = match &completed.meta {
- InstantRecordingMeta::InProgress { .. } => InstantRecordingMeta::Failed {
- error: "instant recording stopped before completion".to_string(),
- },
- other => other.clone(),
- };
+ let pretty_name = completed
+ .project_path
+ .file_stem()
+ .and_then(|name| name.to_str())
+ .filter(|name| !name.is_empty())
+ .unwrap_or("Cap Recording")
+ .to_string();
+ let meta = match &completed.meta {
+ InstantRecordingMeta::InProgress { .. } => InstantRecordingMeta::Failed {
+ error: "instant recording stopped before completion".to_string(),
+ },
+ other => other.clone(),
+ };
- let previous_upload = RecordingMeta::load_for_project(&completed.project_path)
- .ok()
- .and_then(|meta| meta.upload);
-
- RecordingMeta {
- platform: Some(Platform::default()),
- project_path: completed.project_path.clone(),
- pretty_name,
- sharing: Some(cap_project::SharingMeta {
- id: upload.id.clone(),
- link: upload.link.clone(),
- content_hash: None,
- }),
- inner: RecordingMetaInner::Instant(meta),
- upload: previous_upload,
- }
- .save_for_project()
- .map_err(|e| anyhow!("saving instant recording meta: {e}"))?;
+ let previous_upload = RecordingMeta::load_for_project(&completed.project_path)
+ .ok()
+ .and_then(|meta| meta.upload);
+
+ let meta = RecordingMeta {
+ platform: Some(Platform::default()),
+ project_path: completed.project_path.clone(),
+ pretty_name,
+ sharing: Some(cap_project::SharingMeta {
+ id: upload.id.clone(),
+ link: upload.link.clone(),
+ content_hash: None,
+ }),
+ inner: RecordingMetaInner::Instant(meta),
+ upload: previous_upload,
+ };
+ save_instant_metadata(&meta)?;
- ProjectConfiguration::default()
- .write(&completed.project_path)
- .map_err(|e| anyhow!("saving instant project config: {e}"))?;
- Ok(())
+ ProjectConfiguration::default()
+ .write(&completed.project_path)
+ .map_err(|e| anyhow!("saving instant project config: {e}"))?;
+ Ok(())
+ })
}
-fn persist_instant_upload_failure(
+pub(crate) fn persist_instant_upload_failure(
project_path: &std::path::Path,
error: &str,
+ metadata_lock: &Mutex<()>,
) -> anyhow::Result<()> {
- let mut meta = cap_project::RecordingMeta::load_for_project(project_path)
- .map_err(|load_error| anyhow!("loading failed instant recording metadata: {load_error}"))?;
- meta.upload = Some(cap_project::UploadMeta::Failed {
- error: error.to_string(),
- });
- meta.save_for_project()
- .map_err(|save_error| anyhow!("saving failed instant upload: {save_error}"))?;
- Ok(())
+ with_instant_metadata_lock(metadata_lock, || {
+ let mut meta =
+ cap_project::RecordingMeta::load_for_project(project_path).map_err(|load_error| {
+ anyhow!("loading failed instant recording metadata: {load_error}")
+ })?;
+ meta.upload = Some(cap_project::UploadMeta::Failed {
+ error: error.to_string(),
+ });
+ save_instant_metadata(&meta)?;
+ Ok(())
+ })
+}
+
+fn persist_instant_upload_complete(
+ project_path: &std::path::Path,
+ metadata_lock: &Mutex<()>,
+) -> anyhow::Result<()> {
+ with_instant_metadata_lock(metadata_lock, || {
+ let mut meta = cap_project::RecordingMeta::load_for_project(project_path)
+ .map_err(|error| anyhow!("loading instant recording metadata: {error}"))?;
+ meta.upload = Some(cap_project::UploadMeta::Complete);
+ save_instant_metadata(&meta)?;
+ Ok(())
+ })
}
fn persist_in_progress_instant_meta(
project_path: &std::path::Path,
video: &cap_project::VideoUploadInfo,
segmented: bool,
+ metadata_lock: &Mutex<()>,
) -> anyhow::Result<()> {
- let upload = if segmented {
- cap_project::UploadMeta::SegmentUpload {
- video_id: video.id.clone(),
- pre_created_video: video.clone(),
- recording_dir: project_path.to_path_buf(),
- }
- } else {
- cap_project::UploadMeta::MultipartUpload {
- video_id: video.id.clone(),
- file_path: project_path.join("content/output.mp4"),
- pre_created_video: video.clone(),
- recording_dir: project_path.to_path_buf(),
- }
- };
- let pretty_name = project_path
- .file_stem()
- .and_then(|name| name.to_str())
- .filter(|name| !name.is_empty())
- .unwrap_or("Cap Recording")
- .to_string();
-
- cap_project::RecordingMeta {
- platform: Some(cap_project::Platform::default()),
- project_path: project_path.to_path_buf(),
- pretty_name,
- sharing: None,
- inner: cap_project::RecordingMetaInner::Instant(
- cap_project::InstantRecordingMeta::InProgress { recording: true },
- ),
- upload: Some(upload),
- }
- .save_for_project()
- .map_err(|error| anyhow!("saving in-progress instant recording metadata: {error}"))?;
- Ok(())
+ with_instant_metadata_lock(metadata_lock, || {
+ let upload = if segmented {
+ cap_project::UploadMeta::SegmentUpload {
+ video_id: video.id.clone(),
+ pre_created_video: video.clone(),
+ recording_dir: project_path.to_path_buf(),
+ }
+ } else {
+ cap_project::UploadMeta::MultipartUpload {
+ video_id: video.id.clone(),
+ file_path: project_path.join("content/output.mp4"),
+ pre_created_video: video.clone(),
+ recording_dir: project_path.to_path_buf(),
+ }
+ };
+ let pretty_name = project_path
+ .file_stem()
+ .and_then(|name| name.to_str())
+ .filter(|name| !name.is_empty())
+ .unwrap_or("Cap Recording")
+ .to_string();
+
+ let meta = cap_project::RecordingMeta {
+ platform: Some(cap_project::Platform::default()),
+ project_path: project_path.to_path_buf(),
+ pretty_name,
+ sharing: None,
+ inner: cap_project::RecordingMetaInner::Instant(
+ cap_project::InstantRecordingMeta::InProgress { recording: true },
+ ),
+ upload: Some(upload),
+ };
+ save_instant_metadata(&meta)?;
+ Ok(())
+ })
}
pub async fn start(config: StartConfig) -> anyhow::Result {
@@ -770,10 +949,21 @@ async fn start_attempt_with_upload(
let video = pre_created_video
.ok_or_else(|| anyhow!("instant recording has no reserved upload"))?;
let segment_rx = handle.take_segment_rx();
- persist_in_progress_instant_meta(&project_dir, &video, segment_rx.is_some())?;
+ let metadata_lock = Arc::new(Mutex::new(()));
+ persist_in_progress_instant_meta(
+ &project_dir,
+ &video,
+ segment_rx.is_some(),
+ &metadata_lock,
+ )?;
instant_upload = Some(
- crate::upload::start_instant_upload(video, project_dir.clone(), segment_rx)
- .map_err(anyhow::Error::msg)?,
+ crate::upload::start_instant_upload(
+ video,
+ project_dir.clone(),
+ segment_rx,
+ metadata_lock,
+ )
+ .map_err(anyhow::Error::msg)?,
);
Handle::Instant(handle)
}
@@ -970,21 +1160,21 @@ fn create_project_dir(
.with_context(|| format!("creating recordings dir {}", base.display()))?;
match cap_utils::disk_space::free_bytes_for_path(&base) {
- Ok(bytes) if bytes <= cap_utils::disk_space::LOW_DISK_STOP_BYTES => {
+ Ok(bytes) if bytes <= cap_utils::disk_space::RECORDING_DISK_RESERVE_BYTES => {
return Err(anyhow!(
- "Not enough disk space to start recording ({:.2} GB free). Free up at least {} MB and try again.",
+ "Low storage: only {:.2} GB is available. Free up at least {} MB before recording.",
bytes as f64 / 1_073_741_824.0,
- cap_utils::disk_space::LOW_DISK_STOP_BYTES / (1024 * 1024)
+ cap_utils::disk_space::RECORDING_DISK_RESERVE_BYTES / (1024 * 1024)
));
}
- Ok(bytes) if bytes <= cap_utils::disk_space::LOW_DISK_WARN_BYTES => {
+ Ok(bytes) if bytes <= cap_utils::disk_space::RECORDING_DISK_WARN_BYTES => {
tracing::warn!(
bytes_remaining = bytes,
"Starting recording with low disk space"
);
}
Ok(_) => {}
- Err(error) => tracing::warn!("Failed to check disk space before recording: {error}"),
+ Err(error) => return Err(anyhow!("Could not check recording storage: {error}")),
}
let target_name = target.title().unwrap_or_else(|| "Unknown".into());
@@ -1071,6 +1261,7 @@ mod tests {
use crate::store::BlurMode;
use chrono::TimeZone as _;
use serde_json::Value;
+ use std::sync::{Arc, Barrier};
#[test]
fn recording_project_names_honor_mode_target_and_custom_datetime_formats() {
@@ -1126,6 +1317,128 @@ mod tests {
dir
}
+ #[test]
+ fn stopped_instant_metadata_preserves_recovery_and_upload_state() {
+ let dir = temp_project("instant-stopped");
+ std::fs::write(
+ dir.join("recording-meta.json"),
+ r#"{"pretty_name":"Storage test","sharing":null,"recording":true,"upload":{"state":"Failed","error":"offline"}}"#,
+ )
+ .unwrap();
+ mark_instant_recording_stopped(&dir, &std::sync::Mutex::new(())).unwrap();
+ let meta = cap_project::RecordingMeta::load_for_project(&dir).unwrap();
+ assert!(matches!(
+ meta.inner,
+ cap_project::RecordingMetaInner::Instant(
+ cap_project::InstantRecordingMeta::InProgress { recording: false }
+ )
+ ));
+ assert!(matches!(
+ meta.upload,
+ Some(cap_project::UploadMeta::Failed { error }) if error == "offline"
+ ));
+ std::fs::remove_dir_all(dir).unwrap();
+ }
+
+ #[test]
+ fn failed_instant_metadata_write_preserves_the_previous_recording_state() {
+ let dir = temp_project("instant-metadata-write-failure");
+ let original = br#"{"pretty_name":"Storage test","sharing":null,"recording":true}"#;
+ let path = dir.join("recording-meta.json");
+ std::fs::write(&path, original).unwrap();
+
+ let result = write_instant_metadata(&dir, |file| {
+ std::io::Write::write_all(file, b"{partial")?;
+ Err(std::io::Error::other("simulated disk full"))
+ });
+
+ assert!(result.is_err());
+ assert_eq!(std::fs::read(&path).unwrap(), original);
+ assert_eq!(std::fs::read_dir(&dir).unwrap().count(), 1);
+ std::fs::remove_dir_all(dir).unwrap();
+ }
+
+ #[test]
+ fn concurrent_instant_metadata_updates_stay_valid_and_recoverable() {
+ for finalized in [false, true] {
+ let dir = temp_project("instant-metadata-race");
+ let initial = serde_json::json!({
+ "pretty_name": "Race test",
+ "sharing": null,
+ "recording": true,
+ "upload": { "state": "Failed", "error": "old" },
+ });
+ std::fs::write(
+ dir.join("recording-meta.json"),
+ serde_json::to_vec(&initial).unwrap(),
+ )
+ .unwrap();
+
+ let metadata_lock = Arc::new(std::sync::Mutex::new(()));
+ let barrier = Arc::new(Barrier::new(2));
+ let stopped_dir = dir.clone();
+ let stopped_lock = metadata_lock.clone();
+ let stopped_barrier = barrier.clone();
+ let failed_dir = dir.clone();
+ let failed_lock = metadata_lock.clone();
+ let failed_barrier = barrier.clone();
+ std::thread::scope(|scope| {
+ let stopped = scope.spawn(move || {
+ stopped_barrier.wait();
+ if finalized {
+ let completed = instant_recording::CompletedRecording {
+ project_path: stopped_dir,
+ display_source: ScreenCaptureTarget::CameraOnly,
+ meta: cap_project::InstantRecordingMeta::Complete {
+ fps: 30,
+ sample_rate: Some(48_000),
+ },
+ health: cap_recording::RecordingHealth::Healthy,
+ };
+ let video = cap_project::VideoUploadInfo {
+ id: "test".into(),
+ link: "https://example.invalid/s/test".into(),
+ config: cap_project::S3UploadMeta { id: "test".into() },
+ };
+ persist_instant_meta(&completed, &video, &stopped_lock)
+ } else {
+ mark_instant_recording_stopped(&stopped_dir, &stopped_lock)
+ }
+ });
+ let failed = scope.spawn(move || {
+ failed_barrier.wait();
+ persist_instant_upload_failure(&failed_dir, "offline", &failed_lock)
+ });
+ stopped.join().unwrap().unwrap();
+ failed.join().unwrap().unwrap();
+ });
+
+ let meta = cap_project::RecordingMeta::load_for_project(&dir).unwrap();
+ assert!(matches!(
+ (finalized, meta.inner),
+ (
+ false,
+ cap_project::RecordingMetaInner::Instant(
+ cap_project::InstantRecordingMeta::InProgress { recording: false },
+ )
+ ) | (
+ true,
+ cap_project::RecordingMetaInner::Instant(
+ cap_project::InstantRecordingMeta::Complete {
+ fps: 30,
+ sample_rate: Some(48_000)
+ },
+ )
+ )
+ ));
+ assert!(matches!(
+ meta.upload,
+ Some(cap_project::UploadMeta::Failed { error }) if error == "offline"
+ ));
+ std::fs::remove_dir_all(dir).unwrap();
+ }
+ }
+
/// The bridge writes exactly `camera.backgroundBlur.mode` and nothing else
/// -- the whole reason it is a raw-JSON merge and not a typed round trip.
/// `ProjectConfiguration` here models a dozen sections this app has no
diff --git a/apps/desktop-gpui/src/session.rs b/apps/desktop-gpui/src/session.rs
index 8063df86938..95a8c039f85 100644
--- a/apps/desktop-gpui/src/session.rs
+++ b/apps/desktop-gpui/src/session.rs
@@ -8,7 +8,8 @@
use std::time::{Duration, Instant};
-use gpui::{App, AppContext as _, Context, Entity, Global};
+use cap_utils::disk_space::{DiskSpaceStatus, RecordingStorage};
+use gpui::{App, AppContext as _, Context, Entity, Global, Task};
use crate::recording::{self, ActiveRecording, StartConfig};
@@ -27,6 +28,10 @@ pub struct RecordingSession {
active: Option,
/// Why the last start attempt failed, for the main window to surface.
pub error: Option,
+ pub storage_warning: bool,
+ pub storage_notice: Option,
+ stopped_for_low_storage: bool,
+ stopped_elapsed: Option,
/// The config of the live (or last) recording, kept for restart.
last_config: Option,
/// True while the controls bar window is open, so the main window knows to
@@ -53,6 +58,7 @@ pub struct RecordingSession {
started_at: Option,
paused_accum: Duration,
paused_since: Option,
+ storage_monitor: Option>,
}
struct SessionGlobal(Entity);
@@ -64,6 +70,10 @@ impl RecordingSession {
phase: Phase::Idle,
active: None,
error: None,
+ storage_warning: false,
+ storage_notice: None,
+ stopped_for_low_storage: false,
+ stopped_elapsed: None,
last_config: None,
controls_open: false,
mic_muted: false,
@@ -72,6 +82,7 @@ impl RecordingSession {
started_at: None,
paused_accum: Duration::ZERO,
paused_since: None,
+ storage_monitor: None,
});
cx.set_global(SessionGlobal(session.clone()));
session
@@ -92,6 +103,9 @@ impl RecordingSession {
/// Elapsed recording time, excluding paused stretches -- what the bar's
/// timer shows.
pub fn elapsed(&self) -> Duration {
+ if let Some(elapsed) = self.stopped_elapsed {
+ return elapsed;
+ }
let Some(started_at) = self.started_at else {
return Duration::ZERO;
};
@@ -165,8 +179,13 @@ impl RecordingSession {
if self.phase != Phase::Idle {
return;
}
+ self.storage_monitor = None;
self.phase = Phase::Starting;
self.error = None;
+ self.storage_warning = false;
+ self.storage_notice = None;
+ self.stopped_for_low_storage = false;
+ self.stopped_elapsed = None;
self.last_config = Some(config.clone());
cx.notify();
@@ -177,6 +196,7 @@ impl RecordingSession {
match result {
Ok(Ok(active)) => {
tracing::info!(dir = %active.project_dir.display(), "recording started");
+ let project_dir = active.project_dir.clone();
this.active = Some(active);
this.phase = Phase::Recording { paused: false };
// A fresh mic lock always starts unmuted.
@@ -184,6 +204,7 @@ impl RecordingSession {
this.started_at = Some(Instant::now());
this.paused_accum = Duration::ZERO;
this.paused_since = None;
+ this.monitor_storage(project_dir, cx);
}
Ok(Err(error)) => {
tracing::error!("recording failed to start: {error:#}");
@@ -203,24 +224,100 @@ impl RecordingSession {
.detach();
}
+ fn monitor_storage(&mut self, project_dir: std::path::PathBuf, cx: &mut Context) {
+ self.storage_monitor = None;
+ let task = cx.spawn(async move |this, cx| {
+ let mut check_failed = false;
+ loop {
+ cx.background_executor().timer(Duration::from_secs(2)).await;
+ let is_current = |session: &Self| {
+ matches!(session.phase, Phase::Recording { .. })
+ && session
+ .active
+ .as_ref()
+ .is_some_and(|active| active.project_dir == project_dir)
+ };
+ if !this.update(cx, |this, _| is_current(this)).unwrap_or(false) {
+ return;
+ }
+ let path = project_dir.clone();
+ let result = cx
+ .background_executor()
+ .spawn(async move { cap_utils::disk_space::free_bytes_for_path(&path) })
+ .await;
+ let available_bytes = match result {
+ Ok(bytes) => {
+ check_failed = false;
+ bytes
+ }
+ Err(error) => {
+ if !check_failed {
+ tracing::warn!(%error, "Could not check recording storage");
+ check_failed = true;
+ }
+ continue;
+ }
+ };
+ if !this
+ .update(cx, |this, cx| {
+ if !is_current(this) {
+ return false;
+ }
+ match (RecordingStorage {
+ available_bytes,
+ recording_bytes: 0,
+ })
+ .status()
+ {
+ DiskSpaceStatus::Exhausted => {
+ this.storage_warning = true;
+ this.stopped_for_low_storage = true;
+ this.storage_notice =
+ Some("Low storage. Stopping and saving your recording…".into());
+ this.stop(cx);
+ false
+ }
+ status => {
+ let warning = status == DiskSpaceStatus::Low;
+ if warning != this.storage_warning {
+ this.storage_warning = warning;
+ cx.notify();
+ }
+ true
+ }
+ }
+ })
+ .unwrap_or(false)
+ {
+ return;
+ }
+ }
+ });
+ self.storage_monitor = Some(task);
+ }
+
pub fn stop(&mut self, cx: &mut Context) {
if !matches!(self.phase, Phase::Recording { .. }) {
return;
}
+ self.storage_monitor = None;
let Some(active) = self.active.take() else {
return;
};
let instant_share_url = active.instant_share_url().map(ToString::to_string);
+ let low_storage = self.stopped_for_low_storage;
if let Some(link) = &instant_share_url
+ && !low_storage
&& !crate::store::GeneralSettings::load().disable_auto_open_links
{
let separator = if link.contains('?') { '&' } else { '?' };
cx.open_url(&format!("{link}{separator}recordingStopped=1"));
}
+ self.stopped_elapsed = Some(self.elapsed());
self.phase = Phase::Stopping;
cx.notify();
- let task = gpui_tokio::Tokio::spawn(cx, active.stop());
+ let task = gpui_tokio::Tokio::spawn(cx, active.stop(low_storage));
cx.spawn(async move |this, cx| {
let result = task.await;
this.update(cx, |this, cx| {
@@ -233,7 +330,10 @@ impl RecordingSession {
// never does that, so it goes with the placeholder it
// was standing in for.
tracing::info!(dir = %project_dir.display(), "recording finished");
- if this.mode() == Some(crate::recording::RecordingMode::Studio) {
+ if low_storage {
+ this.storage_notice = Some("Recording stopped because storage is low. Your recording was saved.".into());
+ }
+ if this.mode() == Some(crate::recording::RecordingMode::Studio) && !low_storage {
this.finished_studio = Some(project_dir);
} else if let Some(link) = instant_share_url {
cx.write_to_clipboard(gpui::ClipboardItem::new_string(link));
@@ -241,10 +341,12 @@ impl RecordingSession {
}
Ok(Err(error)) => {
tracing::error!("recording failed to stop cleanly: {error:#}");
+ this.storage_notice = low_storage.then(|| "Recording stopped because storage is low. Your recording files were kept. Free up space, then recover the recording below.".into());
this.error = Some(format!("{error:#}"));
}
Err(join_error) => {
tracing::error!("recording stop task died: {join_error}");
+ this.storage_notice = low_storage.then(|| "Recording stopped because storage is low. Your recording files were kept. Free up space, then recover the recording below.".into());
this.error = Some("Stop task failed.".into());
}
}
@@ -294,6 +396,7 @@ impl RecordingSession {
if !matches!(self.phase, Phase::Recording { .. }) {
return;
}
+ self.storage_monitor = None;
let Some(active) = self.active.take() else {
return;
};
@@ -316,6 +419,7 @@ impl RecordingSession {
if !matches!(self.phase, Phase::Recording { .. }) {
return;
}
+ self.storage_monitor = None;
let Some(active) = self.active.take() else {
return;
};
@@ -344,8 +448,11 @@ impl RecordingSession {
}
fn finish(&mut self, cx: &mut Context) {
+ self.storage_monitor = None;
self.phase = Phase::Idle;
self.mic_muted = false;
+ self.storage_warning = false;
+ self.stopped_elapsed = None;
self.started_at = None;
self.paused_accum = Duration::ZERO;
self.paused_since = None;
diff --git a/apps/desktop-gpui/src/upload.rs b/apps/desktop-gpui/src/upload.rs
index b69c7d56b17..811b371c4b4 100644
--- a/apps/desktop-gpui/src/upload.rs
+++ b/apps/desktop-gpui/src/upload.rs
@@ -39,6 +39,7 @@ pub struct InstantUpload {
pub video: VideoUploadInfo,
segment_upload: Option>>,
cancel: Arc,
+ metadata_lock: Arc>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
@@ -143,6 +144,7 @@ pub fn start_instant_upload(
video: VideoUploadInfo,
project_path: PathBuf,
segment_rx: Option>,
+ metadata_lock: Arc>,
) -> Result {
let cancel = Arc::new(AtomicBool::new(false));
let segment_upload = if let Some(segment_rx) = segment_rx {
@@ -160,8 +162,16 @@ pub fn start_instant_upload(
let upload_video = video.clone();
let upload_cancel = cancel.clone();
+ let upload_metadata_lock = metadata_lock.clone();
Some(tokio::spawn(async move {
- run_segment_upload(upload_video, project_path, events_rx, upload_cancel).await
+ run_segment_upload(
+ upload_video,
+ project_path,
+ events_rx,
+ upload_cancel,
+ upload_metadata_lock,
+ )
+ .await
}))
} else {
None
@@ -171,6 +181,7 @@ pub fn start_instant_upload(
video,
segment_upload,
cancel,
+ metadata_lock,
})
}
@@ -183,6 +194,10 @@ impl InstantUpload {
self.segment_upload.is_some()
}
+ pub(crate) fn metadata_lock(&self) -> &Mutex<()> {
+ &self.metadata_lock
+ }
+
pub async fn finish_segments(&mut self) -> Result<(), String> {
let Some(upload) = self.segment_upload.take() else {
return Ok(());
@@ -202,12 +217,16 @@ impl InstantUpload {
}
pub async fn cancel(mut self) -> Result<(), String> {
+ self.abort_segments().await;
+ delete_instant_video(&self.video.id).await
+ }
+
+ pub(crate) async fn abort_segments(&mut self) {
self.cancel.store(true, Ordering::Release);
if let Some(upload) = self.segment_upload.take() {
upload.abort();
+ let _ = upload.await;
}
-
- delete_instant_video(&self.video.id).await
}
}
@@ -234,18 +253,15 @@ async fn run_segment_upload(
project_path: PathBuf,
events: flume::Receiver,
cancel: Arc,
+ metadata_lock: Arc>,
) -> Result<(), String> {
let result = upload_segments(&video.id, events, cancel.clone()).await;
if let Err(error) = &result
&& !cancel.load(Ordering::Acquire)
- && let Ok(mut meta) = RecordingMeta::load_for_project(&project_path)
+ && let Err(save_error) =
+ crate::recording::persist_instant_upload_failure(&project_path, error, &metadata_lock)
{
- meta.upload = Some(UploadMeta::Failed {
- error: error.clone(),
- });
- if let Err(save_error) = meta.save_for_project() {
- tracing::error!("Failed to persist instant upload failure: {save_error}");
- }
+ tracing::error!("Failed to persist instant upload failure: {save_error}");
}
result
}
@@ -1404,6 +1420,47 @@ fn urlencoding(value: &str) -> String {
mod tests {
use super::*;
+ #[test]
+ fn abort_segments_waits_until_the_upload_task_has_stopped() {
+ struct Dropped(Arc);
+ impl Drop for Dropped {
+ fn drop(&mut self) {
+ self.0.store(true, Ordering::Release);
+ }
+ }
+
+ tokio::runtime::Builder::new_current_thread()
+ .build()
+ .unwrap()
+ .block_on(async {
+ let dropped = Arc::new(AtomicBool::new(false));
+ let upload_dropped = dropped.clone();
+ let (started_tx, started_rx) = tokio::sync::oneshot::channel();
+ let task = tokio::spawn(async move {
+ let _dropped = Dropped(upload_dropped);
+ let _ = started_tx.send(());
+ std::future::pending::<()>().await;
+ Ok(())
+ });
+ started_rx.await.unwrap();
+ let mut upload = InstantUpload {
+ video: VideoUploadInfo {
+ id: "test".into(),
+ link: "https://example.invalid/s/test".into(),
+ config: S3UploadMeta { id: "test".into() },
+ },
+ segment_upload: Some(task),
+ cancel: Arc::new(AtomicBool::new(false)),
+ metadata_lock: Arc::new(Mutex::new(())),
+ };
+
+ upload.abort_segments().await;
+
+ assert!(dropped.load(Ordering::Acquire));
+ assert!(upload.cancel.load(Ordering::Acquire));
+ });
+ }
+
fn segment_event(
index: u32,
duration: f64,
diff --git a/crates/utils/src/disk_space.rs b/crates/utils/src/disk_space.rs
index 066f4734610..62769c92887 100644
--- a/crates/utils/src/disk_space.rs
+++ b/crates/utils/src/disk_space.rs
@@ -3,6 +3,60 @@ use std::path::Path;
pub const LOW_DISK_WARN_BYTES: u64 = 200 * 1024 * 1024;
pub const LOW_DISK_STOP_BYTES: u64 = 50 * 1024 * 1024;
+pub const RECORDING_DISK_WARN_BYTES: u64 = 2 * 1024 * 1024 * 1024;
+pub const RECORDING_DISK_RESERVE_BYTES: u64 = 512 * 1024 * 1024;
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub struct RecordingStorage {
+ pub available_bytes: u64,
+ pub recording_bytes: u64,
+}
+
+impl RecordingStorage {
+ pub fn finalization_bytes(self) -> u64 {
+ self.recording_bytes.saturating_mul(2)
+ }
+
+ pub fn status(self) -> DiskSpaceStatus {
+ let remaining = self
+ .available_bytes
+ .saturating_sub(self.finalization_bytes());
+ if remaining <= RECORDING_DISK_RESERVE_BYTES {
+ DiskSpaceStatus::Exhausted
+ } else if remaining <= RECORDING_DISK_WARN_BYTES {
+ DiskSpaceStatus::Low
+ } else {
+ DiskSpaceStatus::Ok
+ }
+ }
+
+ pub fn can_finalize(self) -> bool {
+ self.available_bytes
+ > self
+ .finalization_bytes()
+ .saturating_add(RECORDING_DISK_RESERVE_BYTES / 2)
+ }
+}
+
+pub fn recording_storage(path: &Path) -> io::Result {
+ let mut pending = vec![path.to_path_buf()];
+ let mut recording_bytes = 0u64;
+ while let Some(directory) = pending.pop() {
+ for entry in std::fs::read_dir(directory)? {
+ let entry = entry?;
+ let file_type = entry.file_type()?;
+ if file_type.is_dir() {
+ pending.push(entry.path());
+ } else if file_type.is_file() {
+ recording_bytes = recording_bytes.saturating_add(entry.metadata()?.len());
+ }
+ }
+ }
+ Ok(RecordingStorage {
+ available_bytes: free_bytes_for_path(path)?,
+ recording_bytes,
+ })
+}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DiskSpaceStatus {
@@ -92,6 +146,81 @@ fn resolve_existing_ancestor(path: &Path) -> std::path::PathBuf {
mod tests {
use super::*;
+ #[test]
+ fn recording_storage_warns_and_stops_before_the_drive_is_full() {
+ let status = |available_bytes| {
+ RecordingStorage {
+ available_bytes,
+ recording_bytes: 0,
+ }
+ .status()
+ };
+ assert_eq!(status(RECORDING_DISK_WARN_BYTES + 1), DiskSpaceStatus::Ok);
+ assert_eq!(status(RECORDING_DISK_WARN_BYTES), DiskSpaceStatus::Low);
+ assert_eq!(
+ status(RECORDING_DISK_RESERVE_BYTES + 1),
+ DiskSpaceStatus::Low
+ );
+ assert_eq!(
+ status(RECORDING_DISK_RESERVE_BYTES),
+ DiskSpaceStatus::Exhausted
+ );
+ assert_eq!(status(139_771_904), DiskSpaceStatus::Exhausted);
+ assert_eq!(status(0), DiskSpaceStatus::Exhausted);
+ }
+
+ #[test]
+ fn recording_storage_saturates_large_finalization_requirements() {
+ let storage = RecordingStorage {
+ available_bytes: u64::MAX,
+ recording_bytes: u64::MAX,
+ };
+ assert_eq!(storage.status(), DiskSpaceStatus::Exhausted);
+ assert!(!storage.can_finalize());
+ }
+
+ #[test]
+ fn recording_storage_reserves_space_for_finalization() {
+ let recording_bytes = 1024 * 1024 * 1024;
+ let mut storage = RecordingStorage {
+ available_bytes: recording_bytes * 2 + RECORDING_DISK_WARN_BYTES + 1,
+ recording_bytes,
+ };
+ assert_eq!(storage.status(), DiskSpaceStatus::Ok);
+ storage.available_bytes -= 1;
+ assert_eq!(storage.status(), DiskSpaceStatus::Low);
+ storage.available_bytes = recording_bytes * 2 + RECORDING_DISK_RESERVE_BYTES;
+ assert_eq!(storage.status(), DiskSpaceStatus::Exhausted);
+ assert!(storage.can_finalize());
+ storage.available_bytes = recording_bytes * 2;
+ assert!(!storage.can_finalize());
+ }
+
+ #[test]
+ fn recording_storage_accounts_for_every_track() {
+ let directory = tempfile::tempdir().unwrap();
+ let segment = directory.path().join("segment-0");
+ std::fs::create_dir(&segment).unwrap();
+ std::fs::write(segment.join("display.m4s"), [0; 128]).unwrap();
+ std::fs::write(segment.join("camera.m4s"), [0; 64]).unwrap();
+ std::fs::write(directory.path().join("recording-meta.json"), [0; 16]).unwrap();
+ assert_eq!(
+ recording_storage(directory.path()).unwrap().recording_bytes,
+ 208
+ );
+ }
+
+ #[cfg(unix)]
+ #[test]
+ fn recording_storage_does_not_follow_symbolic_links() {
+ let directory = tempfile::tempdir().unwrap();
+ std::os::unix::fs::symlink(directory.path(), directory.path().join("loop")).unwrap();
+ assert_eq!(
+ recording_storage(directory.path()).unwrap().recording_bytes,
+ 0
+ );
+ }
+
#[test]
fn status_from_bytes() {
assert_eq!(DiskSpaceStatus::from_bytes(u64::MAX), DiskSpaceStatus::Ok);
From 9c95e5118b3d2b86d9279667b5e0684112b9e51b Mon Sep 17 00:00:00 2001
From: Richie McIlroy <33632126+richiemcilroy@users.noreply.github.com>
Date: Wed, 26 Aug 2026 20:07:28 +0100
Subject: [PATCH 4/5] fix: reserve finalization space while recording
---
apps/desktop-gpui/src/session.rs | 24 ++---
crates/utils/src/disk_space.rs | 166 ++++++++++++++++++++++++++++++-
2 files changed, 177 insertions(+), 13 deletions(-)
diff --git a/apps/desktop-gpui/src/session.rs b/apps/desktop-gpui/src/session.rs
index 95a8c039f85..e543690e04b 100644
--- a/apps/desktop-gpui/src/session.rs
+++ b/apps/desktop-gpui/src/session.rs
@@ -8,7 +8,7 @@
use std::time::{Duration, Instant};
-use cap_utils::disk_space::{DiskSpaceStatus, RecordingStorage};
+use cap_utils::disk_space::{DiskSpaceStatus, RecordingStorageMonitor};
use gpui::{App, AppContext as _, Context, Entity, Global, Task};
use crate::recording::{self, ActiveRecording, StartConfig};
@@ -228,6 +228,7 @@ impl RecordingSession {
self.storage_monitor = None;
let task = cx.spawn(async move |this, cx| {
let mut check_failed = false;
+ let mut storage_monitor = RecordingStorageMonitor::default();
loop {
cx.background_executor().timer(Duration::from_secs(2)).await;
let is_current = |session: &Self| {
@@ -241,14 +242,18 @@ impl RecordingSession {
return;
}
let path = project_dir.clone();
- let result = cx
+ let (next_monitor, result) = cx
.background_executor()
- .spawn(async move { cap_utils::disk_space::free_bytes_for_path(&path) })
+ .spawn(async move {
+ let result = storage_monitor.sample(&path);
+ (storage_monitor, result)
+ })
.await;
- let available_bytes = match result {
- Ok(bytes) => {
+ storage_monitor = next_monitor;
+ let storage = match result {
+ Ok(storage) => {
check_failed = false;
- bytes
+ storage
}
Err(error) => {
if !check_failed {
@@ -263,12 +268,7 @@ impl RecordingSession {
if !is_current(this) {
return false;
}
- match (RecordingStorage {
- available_bytes,
- recording_bytes: 0,
- })
- .status()
- {
+ match storage.status() {
DiskSpaceStatus::Exhausted => {
this.storage_warning = true;
this.stopped_for_low_storage = true;
diff --git a/crates/utils/src/disk_space.rs b/crates/utils/src/disk_space.rs
index 62769c92887..bcf29e55457 100644
--- a/crates/utils/src/disk_space.rs
+++ b/crates/utils/src/disk_space.rs
@@ -1,5 +1,6 @@
+use std::collections::HashMap;
use std::io;
-use std::path::Path;
+use std::path::{Path, PathBuf};
pub const LOW_DISK_WARN_BYTES: u64 = 200 * 1024 * 1024;
pub const LOW_DISK_STOP_BYTES: u64 = 50 * 1024 * 1024;
@@ -58,6 +59,96 @@ pub fn recording_storage(path: &Path) -> io::Result {
})
}
+#[derive(Default)]
+pub struct RecordingStorageMonitor {
+ fragments: HashMap,
+}
+
+#[derive(Clone, Copy, Default)]
+struct FinalizedFragments {
+ last_index: Option,
+ bytes: u64,
+}
+
+impl RecordingStorageMonitor {
+ pub fn sample(&mut self, path: &Path) -> io::Result {
+ const MAX_CACHED_DIRECTORIES: usize = 256;
+ let mut pending = vec![path.to_path_buf()];
+ let mut recording_bytes = 0u64;
+ while let Some(directory) = pending.pop() {
+ let cached = self.fragments.get(&directory).copied().unwrap_or_default();
+ let mut new_fragment_bytes = 0u64;
+ let mut newest_fragment = None;
+ let mut other_bytes = 0u64;
+ for entry in std::fs::read_dir(&directory)? {
+ let entry = entry?;
+ let file_type = match entry.file_type() {
+ Ok(file_type) => file_type,
+ Err(error) if error.kind() == io::ErrorKind::NotFound => continue,
+ Err(error) => return Err(error),
+ };
+ if file_type.is_dir() {
+ pending.push(entry.path());
+ } else if file_type.is_file() {
+ let name = entry.file_name();
+ let fragment_index = name.to_str().and_then(|name| {
+ name.strip_prefix("segment_")?
+ .strip_suffix(".m4s")?
+ .parse::()
+ .ok()
+ });
+ if let Some(index) = fragment_index {
+ if cached.last_index.is_some_and(|last| index <= last) {
+ continue;
+ }
+ let bytes = match entry.metadata() {
+ Ok(metadata) => metadata.len(),
+ Err(error) if error.kind() == io::ErrorKind::NotFound => continue,
+ Err(error) => return Err(error),
+ };
+ new_fragment_bytes = new_fragment_bytes.saturating_add(bytes);
+ if newest_fragment.is_none_or(|(newest, _)| index > newest) {
+ newest_fragment = Some((index, bytes));
+ }
+ } else {
+ match entry.metadata() {
+ Ok(metadata) => {
+ other_bytes = other_bytes.saturating_add(metadata.len());
+ }
+ Err(error) if error.kind() == io::ErrorKind::NotFound => {}
+ Err(error) => return Err(error),
+ }
+ }
+ }
+ }
+ recording_bytes = recording_bytes
+ .saturating_add(cached.bytes)
+ .saturating_add(new_fragment_bytes)
+ .saturating_add(other_bytes);
+ if let Some((index, bytes)) = newest_fragment
+ && (self.fragments.contains_key(&directory)
+ || self.fragments.len() < MAX_CACHED_DIRECTORIES)
+ {
+ // DASH fragments are immutable once a later fragment exists. Keep
+ // measuring the newest fragment and temporary files until then.
+ let _ = self.fragments.insert(
+ directory,
+ FinalizedFragments {
+ last_index: index.checked_sub(1),
+ bytes: cached
+ .bytes
+ .saturating_add(new_fragment_bytes.saturating_sub(bytes)),
+ },
+ );
+ }
+ }
+ Ok(RecordingStorage {
+ available_bytes: free_bytes_for_path(path)?,
+ recording_bytes,
+ })
+ }
+}
+
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DiskSpaceStatus {
Ok,
@@ -210,6 +301,79 @@ mod tests {
);
}
+ #[test]
+ fn monitor_counts_growing_fragments_temporary_files_and_other_tracks() {
+ let directory = tempfile::tempdir().unwrap();
+ let display = directory.path().join("display");
+ std::fs::create_dir(&display).unwrap();
+ std::fs::write(display.join("segment_999.m4s"), [0; 128]).unwrap();
+ std::fs::write(display.join("segment_1000.m4s"), [0; 64]).unwrap();
+ std::fs::write(display.join("segment_1001.m4s.tmp"), [0; 32]).unwrap();
+ std::fs::write(directory.path().join("camera.mp4"), [0; 16]).unwrap();
+ let mut monitor = RecordingStorageMonitor::default();
+ assert_eq!(
+ monitor.sample(directory.path()).unwrap().recording_bytes,
+ 240
+ );
+ assert_eq!(monitor.fragments.len(), 1);
+
+ std::fs::write(display.join("segment_1000.m4s"), [0; 96]).unwrap();
+ std::fs::rename(
+ display.join("segment_1001.m4s.tmp"),
+ display.join("segment_1001.m4s"),
+ )
+ .unwrap();
+ std::fs::write(directory.path().join("camera.mp4"), [0; 32]).unwrap();
+ assert_eq!(
+ monitor.sample(directory.path()).unwrap().recording_bytes,
+ 288
+ );
+ assert_eq!(
+ monitor.sample(directory.path()).unwrap().recording_bytes,
+ recording_storage(directory.path()).unwrap().recording_bytes
+ );
+
+ std::fs::remove_file(display.join("segment_999.m4s")).unwrap();
+ assert_eq!(
+ monitor.sample(directory.path()).unwrap().recording_bytes,
+ 288
+ );
+ }
+
+ #[test]
+ fn monitor_reserves_finalization_space_for_large_recordings() {
+ let directory = tempfile::tempdir().unwrap();
+ let path = directory.path().join("segment_000.m4s");
+ let file = std::fs::File::create(path).unwrap();
+ file.set_len(1024 * 1024 * 1024).unwrap();
+ let mut storage = RecordingStorageMonitor::default()
+ .sample(directory.path())
+ .unwrap();
+ storage.available_bytes = 2 * storage.recording_bytes + RECORDING_DISK_RESERVE_BYTES;
+ assert_eq!(storage.recording_bytes, 1024 * 1024 * 1024);
+ assert_eq!(storage.status(), DiskSpaceStatus::Exhausted);
+ assert!(storage.can_finalize());
+ }
+
+ #[test]
+ fn monitor_bounds_cached_directories_without_missing_uncached_tracks() {
+ let directory = tempfile::tempdir().unwrap();
+ for index in 0..260 {
+ let track = directory.path().join(index.to_string());
+ std::fs::create_dir(&track).unwrap();
+ std::fs::write(track.join("segment_000.m4s"), [0; 4]).unwrap();
+ std::fs::write(track.join("segment_001.m4s"), [0; 8]).unwrap();
+ }
+ let mut monitor = RecordingStorageMonitor::default();
+ for _ in 0..2 {
+ assert_eq!(
+ monitor.sample(directory.path()).unwrap().recording_bytes,
+ 3120
+ );
+ assert_eq!(monitor.fragments.len(), 256);
+ }
+ }
+
#[cfg(unix)]
#[test]
fn recording_storage_does_not_follow_symbolic_links() {
From cb7554f6ab1a4fae331bfe98eecf9ced28d2de41 Mon Sep 17 00:00:00 2001
From: Richie McIlroy <33632126+richiemcilroy@users.noreply.github.com>
Date: Wed, 26 Aug 2026 20:12:03 +0100
Subject: [PATCH 5/5] fix: reject incomplete fragments during recovery
---
apps/desktop-gpui/src/library.rs | 59 ++++++++++--
crates/recording/src/recovery.rs | 159 ++++++++++++++++++++++++++-----
2 files changed, 186 insertions(+), 32 deletions(-)
diff --git a/apps/desktop-gpui/src/library.rs b/apps/desktop-gpui/src/library.rs
index dd53165448f..11dc7f883ba 100644
--- a/apps/desktop-gpui/src/library.rs
+++ b/apps/desktop-gpui/src/library.rs
@@ -426,12 +426,7 @@ pub fn find_incomplete_recordings_in(
let segment_count = std::fs::read_dir(display)
.ok()?
.filter_map(Result::ok)
- .filter(|entry| {
- entry
- .path()
- .extension()
- .is_some_and(|extension| extension == "m4s")
- })
+ .filter(|entry| RecoveryManager::is_m4s_complete(&entry.path()))
.count();
return (segment_count > 0).then_some(IncompleteRecordingItem {
project_path: item.path,
@@ -1271,6 +1266,16 @@ mod tests {
dir
}
+ fn complete_m4s_fragment() -> Vec {
+ let mut fragment = Vec::new();
+ for name in [b"moof", b"mdat"] {
+ fragment.extend_from_slice(&72u32.to_be_bytes());
+ fragment.extend_from_slice(name);
+ fragment.extend_from_slice(&[0; 64]);
+ }
+ fragment
+ }
+
#[test]
fn recents_are_newest_first_and_capped_at_nine() {
let root = temp_dir("cap");
@@ -1544,9 +1549,10 @@ mod tests {
);
if fragments {
let display = bundle.join("content/segments/segment-0/display");
+ let fragment = complete_m4s_fragment();
std::fs::create_dir_all(&display).unwrap();
std::fs::write(display.join("init.mp4"), vec![0u8; 128]).unwrap();
- std::fs::write(display.join("segment_001.m4s"), vec![1u8; 256]).unwrap();
+ std::fs::write(display.join("segment_001.m4s"), &fragment).unwrap();
std::fs::write(
display.join("manifest.json"),
serde_json::to_vec(&serde_json::json!({
@@ -1556,7 +1562,7 @@ mod tests {
"segments": [{
"path": "segment_001.m4s",
"is_complete": true,
- "file_size": 256
+ "file_size": fragment.len()
}]
}))
.unwrap(),
@@ -1648,7 +1654,13 @@ mod tests {
let display = bundle.join("content/display");
std::fs::create_dir_all(&display).unwrap();
std::fs::write(display.join("init.mp4"), [0; 8]).unwrap();
- std::fs::write(display.join("segment_001.m4s"), [0; 16]).unwrap();
+ let fragment = complete_m4s_fragment();
+ std::fs::write(display.join("segment_001.m4s"), &fragment).unwrap();
+ std::fs::write(
+ display.join("segment_002.m4s"),
+ &fragment[..fragment.len() - 1],
+ )
+ .unwrap();
let items = list_recordings_in(std::slice::from_ref(&root));
if !recording {
assert_eq!(items[0].status, RecordingStatus::NeedsRemux);
@@ -1657,15 +1669,44 @@ mod tests {
let recoverable = find_incomplete_recordings_in(std::slice::from_ref(&root), None);
assert_eq!(recoverable.len(), 1);
assert_eq!(recoverable[0].project_path, bundle);
+ assert_eq!(recoverable[0].segment_count, 1);
assert!(
find_incomplete_recordings_in(std::slice::from_ref(&root), Some(&bundle))
.is_empty()
);
assert!(display.join("segment_001.m4s").is_file());
+ assert!(display.join("segment_002.m4s").is_file());
std::fs::remove_dir_all(root).unwrap();
}
}
+ #[test]
+ fn interrupted_instant_recordings_with_only_partial_fragments_are_not_recoverable() {
+ let root = temp_dir("instant-partial-recovery");
+ let metadata = serde_json::json!({
+ "pretty_name": "Custom instant recording",
+ "sharing": null,
+ "recording": false,
+ })
+ .to_string();
+ let bundle = write_bundle(&root, "deferred-instant", &metadata);
+ let display = bundle.join("content/display");
+ std::fs::create_dir_all(&display).unwrap();
+ std::fs::write(display.join("init.mp4"), [0; 8]).unwrap();
+ let fragment = complete_m4s_fragment();
+ std::fs::write(
+ display.join("segment_001.m4s"),
+ &fragment[..fragment.len() - 1],
+ )
+ .unwrap();
+
+ let recoverable = find_incomplete_recordings_in(std::slice::from_ref(&root), None);
+
+ assert!(recoverable.is_empty());
+ assert!(display.join("segment_001.m4s").is_file());
+ std::fs::remove_dir_all(root).unwrap();
+ }
+
#[test]
fn recovery_rejects_paths_outside_the_recording_library() {
let root = temp_dir("recovery-path");
diff --git a/crates/recording/src/recovery.rs b/crates/recording/src/recovery.rs
index 1f91bc29cec..60e92c15d6d 100644
--- a/crates/recording/src/recovery.rs
+++ b/crates/recording/src/recovery.rs
@@ -398,6 +398,15 @@ impl RecoveryManager {
return None;
}
+ if path
+ .extension()
+ .is_some_and(|extension| extension.eq_ignore_ascii_case("m4s"))
+ && !Self::is_m4s_complete(&path)
+ {
+ warn!("Fragment {} has an incomplete M4S tail", path.display());
+ return None;
+ }
+
if Self::is_video_file(&path) {
if init_segment.is_some() {
Some(path)
@@ -443,7 +452,6 @@ impl RecoveryManager {
&& path
.symlink_metadata()
.is_ok_and(|metadata| metadata.is_file())
- && tail_is_complete(path).unwrap_or(false)
},
));
result.sort_by(|a, b| {
@@ -462,7 +470,10 @@ impl RecoveryManager {
}
}
- if let Some(init_segment) = manifest_init_segment {
+ if let Some(init_segment) = manifest_init_segment.or_else(|| {
+ let path = dir.join("init.mp4");
+ path.is_file().then_some(path)
+ }) {
let fragments = Self::probe_m4s_fragments_with_init(dir);
if !fragments.is_empty() {
return FragmentsInfo {
@@ -482,8 +493,6 @@ impl RecoveryManager {
dir: &Path,
health_tx: Option<&HealthSender>,
) -> Vec<(u32, PathBuf, Vec)> {
- const MIN_VALID_FRAGMENT_SIZE: u64 = 100;
-
let Ok(entries) = std::fs::read_dir(dir) else {
return Vec::new();
};
@@ -532,13 +541,8 @@ impl RecoveryManager {
.strip_prefix("segment_")
.and_then(|s| s.strip_suffix(".m4s"))
.and_then(|s| s.parse().ok())?;
- let metadata = std::fs::metadata(&p).ok()?;
- if metadata.len() < MIN_VALID_FRAGMENT_SIZE {
- debug!(
- "Skipping tiny respawn fragment {} ({} bytes)",
- p.display(),
- metadata.len()
- );
+ if !Self::is_m4s_complete(&p) {
+ debug!("Skipping incomplete respawn fragment {}", p.display());
return None;
}
Some((idx, p))
@@ -673,6 +677,7 @@ impl RecoveryManager {
.and_then(|e| e.to_str())
.map(|e| e.to_lowercase());
match ext.as_deref() {
+ Some("m4s") if !Self::is_m4s_complete(p) => false,
Some("mp4") | Some("m4s") => match probe_video_can_decode(p) {
Ok(true) => true,
Ok(false) => {
@@ -695,8 +700,6 @@ impl RecoveryManager {
}
fn probe_m4s_fragments_with_init(dir: &Path) -> Vec {
- const MIN_VALID_FRAGMENT_SIZE: u64 = 100;
-
let Ok(entries) = std::fs::read_dir(dir) else {
return Vec::new();
};
@@ -704,18 +707,28 @@ impl RecoveryManager {
let mut fragments: Vec<_> = entries
.filter_map(|e| e.ok())
.map(|e| e.path())
- .filter(|p| p.extension().is_some_and(|e| e.eq_ignore_ascii_case("m4s")))
- .filter(|p| {
- std::fs::metadata(p)
- .map(|metadata| metadata.len() >= MIN_VALID_FRAGMENT_SIZE)
- .unwrap_or(false)
- })
+ .filter(|p| Self::is_m4s_complete(p))
.collect();
- fragments.sort();
+ fragments.sort_by(|a, b| {
+ Self::m4s_fragment_index(a)
+ .cmp(&Self::m4s_fragment_index(b))
+ .then_with(|| a.cmp(b))
+ });
fragments
}
+ pub fn is_m4s_complete(path: &Path) -> bool {
+ const MIN_VALID_FRAGMENT_SIZE: u64 = 100;
+
+ path.extension()
+ .is_some_and(|extension| extension.eq_ignore_ascii_case("m4s"))
+ && std::fs::symlink_metadata(path).is_ok_and(|metadata| {
+ metadata.is_file() && metadata.len() >= MIN_VALID_FRAGMENT_SIZE
+ })
+ && tail_is_complete(path).unwrap_or(false)
+ }
+
fn m4s_fragment_index(path: &Path) -> Option {
path.file_name()?
.to_str()?
@@ -1704,15 +1717,20 @@ mod tests {
use std::fs;
use tempfile::tempdir;
- #[test]
- fn recovery_includes_complete_fragments_missing_from_the_last_manifest() {
- let dir = tempdir().unwrap();
+ fn complete_m4s_fragment() -> Vec {
let mut fragment = Vec::new();
for name in [b"moof", b"mdat"] {
fragment.extend_from_slice(&72u32.to_be_bytes());
fragment.extend_from_slice(name);
fragment.extend_from_slice(&[0; 64]);
}
+ fragment
+ }
+
+ #[test]
+ fn recovery_includes_complete_fragments_missing_from_the_last_manifest() {
+ let dir = tempdir().unwrap();
+ let fragment = complete_m4s_fragment();
fs::write(dir.path().join("init.mp4"), [0; 128]).unwrap();
for name in [
"segment_001.m4s",
@@ -1752,6 +1770,101 @@ mod tests {
assert!(dir.path().join("segment_1001.m4s").is_file());
}
+ #[test]
+ fn recovery_reconstructs_only_complete_fragments_when_manifest_has_no_valid_entries() {
+ let dir = tempdir().unwrap();
+ let fragment = complete_m4s_fragment();
+ let partial = &fragment[..fragment.len() - 1];
+ fs::write(dir.path().join("init.mp4"), [0; 128]).unwrap();
+ fs::write(dir.path().join("segment_001.m4s"), &fragment).unwrap();
+ fs::write(dir.path().join("segment_999.m4s"), &fragment).unwrap();
+ fs::write(dir.path().join("segment_1000.m4s"), &fragment).unwrap();
+ fs::write(dir.path().join("segment_002.m4s"), partial).unwrap();
+ let manifest = serde_json::to_vec(&serde_json::json!({
+ "version": 5,
+ "type": "m4s_segments",
+ "init_segment": "init.mp4",
+ "segments": [
+ {"path": "segment_001.m4s", "is_complete": false, "file_size": fragment.len()},
+ {"path": "segment_999.m4s", "is_complete": false, "file_size": fragment.len()},
+ {"path": "segment_1000.m4s", "is_complete": false, "file_size": fragment.len()},
+ {"path": "segment_002.m4s", "is_complete": false, "file_size": partial.len()}
+ ]
+ }))
+ .unwrap();
+ fs::write(dir.path().join("manifest.json"), &manifest).unwrap();
+
+ let recovered = RecoveryManager::find_complete_fragments_with_init(dir.path());
+
+ assert_eq!(
+ recovered.fragments,
+ ["segment_001.m4s", "segment_999.m4s", "segment_1000.m4s"]
+ .map(|name| dir.path().join(name))
+ );
+ assert_eq!(recovered.init_segment, Some(dir.path().join("init.mp4")));
+ assert_eq!(
+ fs::read(dir.path().join("manifest.json")).unwrap(),
+ manifest
+ );
+ assert!(dir.path().join("segment_002.m4s").is_file());
+ }
+
+ #[test]
+ fn recovery_rejects_partial_only_fragments_when_manifest_has_no_valid_entries() {
+ let dir = tempdir().unwrap();
+ let fragment = complete_m4s_fragment();
+ let partial = &fragment[..fragment.len() - 1];
+ fs::write(dir.path().join("init.mp4"), [0; 128]).unwrap();
+ fs::write(dir.path().join("segment_001.m4s"), partial).unwrap();
+ let manifest = serde_json::to_vec(&serde_json::json!({
+ "version": 5,
+ "type": "m4s_segments",
+ "init_segment": "init.mp4",
+ "segments": [
+ {"path": "segment_001.m4s", "is_complete": false, "file_size": partial.len()}
+ ]
+ }))
+ .unwrap();
+ fs::write(dir.path().join("manifest.json"), &manifest).unwrap();
+
+ let recovered = RecoveryManager::find_complete_fragments_with_init(dir.path());
+
+ assert!(recovered.fragments.is_empty());
+ assert_eq!(
+ fs::read(dir.path().join("manifest.json")).unwrap(),
+ manifest
+ );
+ assert!(dir.path().join("segment_001.m4s").is_file());
+ }
+
+ #[test]
+ fn recovery_finds_complete_fragments_without_a_readable_manifest() {
+ for manifest in [None, Some("{interrupted")] {
+ let dir = tempdir().unwrap();
+ let fragment = complete_m4s_fragment();
+ fs::write(dir.path().join("init.mp4"), [0; 128]).unwrap();
+ fs::write(dir.path().join("segment_001.m4s"), &fragment).unwrap();
+ fs::write(
+ dir.path().join("segment_002.m4s"),
+ &fragment[..fragment.len() - 1],
+ )
+ .unwrap();
+ if let Some(manifest) = manifest {
+ fs::write(dir.path().join("manifest.json"), manifest).unwrap();
+ }
+
+ let recovered = RecoveryManager::find_complete_fragments_with_init(dir.path());
+
+ assert_eq!(recovered.fragments, [dir.path().join("segment_001.m4s")]);
+ assert_eq!(recovered.init_segment, Some(dir.path().join("init.mp4")));
+ assert_eq!(
+ fs::read_to_string(dir.path().join("manifest.json")).ok(),
+ manifest.map(str::to_string)
+ );
+ assert!(dir.path().join("segment_002.m4s").is_file());
+ }
+ }
+
#[test]
fn replace_file_overwrites_existing_destination() {
let dir = tempdir().unwrap();