Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion .github/workflows/recording-reliability.yml
Original file line number Diff line number Diff line change
Expand Up @@ -47,4 +47,6 @@ jobs:
__tests__/unit/media-server-progress.test.ts \
__tests__/unit/media-processing-budget.test.ts \
__tests__/unit/playback-source.test.ts \
__tests__/unit/upload-progress-playback.test.ts
__tests__/unit/upload-progress-playback.test.ts \
__tests__/unit/recording-prepare.test.ts \
__tests__/unit/s3-bucket-connections.test.ts
121 changes: 121 additions & 0 deletions apps/desktop-gpui/src/upload.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ use std::time::{Duration, Instant};

use cap_enc_ffmpeg::segmented_stream::{SegmentCompletedEvent, SegmentMediaType};
use cap_project::{RecordingMeta, S3UploadMeta, SharingMeta, UploadMeta, VideoUploadInfo};
use cap_recording::upload_preparation::{Preparation, Segment};
use futures_util::{StreamExt as _, stream::FuturesUnordered};
use reqwest::StatusCode;
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -517,6 +518,13 @@ async fn run_segment_upload(
}

trait SegmentTransport: Send + Sync {
fn prepare(
&self,
_video_id: &str,
_segments: Vec<Segment>,
) -> impl Future<Output = Result<Option<Vec<Segment>>, String>> + Send {
std::future::ready(Ok(None))
}
fn prefetch(
&self,
video_id: &str,
Expand Down Expand Up @@ -545,6 +553,41 @@ trait SegmentTransport: Send + Sync {

struct LiveSegmentTransport;
impl SegmentTransport for LiveSegmentTransport {
async fn prepare(
&self,
video_id: &str,
segments: Vec<Segment>,
) -> Result<Option<Vec<Segment>>, String> {
tokio::time::timeout(Duration::from_secs(20), async {
#[derive(Deserialize)]
struct Response {
version: u32,
prepared: Vec<Segment>,
}
let response = auth::authed_request(
reqwest::Method::POST,
"/api/recording/prepare",
Some(json!({ "videoId": video_id, "segments": segments })),
)
.await
.map_err(|error| error.to_string())?;
let status = response.status();
if status == StatusCode::NOT_FOUND || status == StatusCode::METHOD_NOT_ALLOWED {
return Ok(None);
}
if !status.is_success() {
return Err(format!("Recording preparation returned {status}"));
}
let response = response
.json::<Response>()
.await
.map_err(|error| error.to_string())?;
Ok((response.version == 1 && response.prepared.len() <= 32)
.then_some(response.prepared))
})
.await
.map_err(|error| error.to_string())?
}
async fn prefetch(
&self,
video_id: &str,
Expand Down Expand Up @@ -589,6 +632,13 @@ async fn upload_segments(
let mut uploads = FuturesUnordered::new();
let mut events_closed = false;
let mut last_manifest_upload: Option<Instant> = None;
let mut preparation = Preparation::default();
let mut preparation_enabled = true;
let mut preparation_batch = Vec::new();
let mut preparation_interval = tokio::time::interval(Duration::from_secs(30));
preparation_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut preparation_request = None;
let mut next_preparation_request = tokio::time::Instant::now();
let mut next_prefetch = SEGMENT_URL_PREFETCH + 1;
let prefetched = checked_segment_step(&cancel, || async {
Ok(transport.prefetch(video_id, 1, SEGMENT_URL_PREFETCH).await)
Expand All @@ -612,6 +662,25 @@ async fn upload_segments(
None => (None, None),
};
tokio::select! {
_ = preparation_interval.tick(), if preparation_enabled && preparation_request.is_none() => {
if tokio::time::Instant::now() < next_preparation_request { continue; }
preparation_batch = preparation.next_batch(manifest.video_segments.iter().map(|segment| segment.index), manifest.audio_segments.iter().map(|segment| segment.index));
if !preparation_batch.is_empty() {
preparation_request = Some(Box::pin(transport.prepare(video_id, preparation_batch.clone())));
}
}
response = async { preparation_request.as_mut().unwrap().await }, if preparation_request.is_some() => {
preparation_request = None;
match response {
Ok(Some(prepared)) => preparation.acknowledge(&preparation_batch, &prepared),
Ok(None) => preparation_enabled = false,
Err(error) => {
preparation.request_failed();
tracing::debug!(%error, "Optional recording preparation unavailable");
}
}
next_preparation_request = tokio::time::Instant::now() + preparation.retry_delay();
}
permission = async { permission.unwrap().await }, if !authorized => {
permission.map_err(|_| "Instant completion was not authorized".to_string())?;
authorized = true;
Expand Down Expand Up @@ -667,6 +736,7 @@ async fn upload_segments(
}
}

drop(preparation_request);
if cancel.load(Ordering::Acquire) {
return Err("Instant recording upload cancelled".to_string());
}
Expand Down Expand Up @@ -2354,6 +2424,9 @@ mod tests {
}
#[derive(Default)]
struct FakeSegmentTransport {
delay_preparation: AtomicBool,
preparation_started: tokio::sync::Notify,
preparation_dropped: AtomicBool,
manifests: Mutex<Vec<bool>>,
completed: std::sync::atomic::AtomicUsize,
uploaded: std::sync::atomic::AtomicUsize,
Expand All @@ -2367,6 +2440,20 @@ mod tests {
prefetch_response: tokio::sync::Notify,
}
impl SegmentTransport for FakeSegmentTransport {
async fn prepare(&self, _: &str, _: Vec<Segment>) -> Result<Option<Vec<Segment>>, String> {
if !self.delay_preparation.load(Ordering::Acquire) {
return Ok(None);
}
struct Finish<'a>(&'a AtomicBool);
impl Drop for Finish<'_> {
fn drop(&mut self) {
self.0.store(true, Ordering::Release);
}
}
let _finish = Finish(&self.preparation_dropped);
self.preparation_started.notify_one();
std::future::pending().await
}
async fn prefetch(
&self,
_: &str,
Expand Down Expand Up @@ -2412,6 +2499,40 @@ mod tests {
Ok(())
}
}
#[tokio::test]
async fn stopped_upload_drops_optional_preparation_without_waiting_for_it() {
let transport = FakeSegmentTransport::default();
transport.delay_preparation.store(true, Ordering::Release);
let (sender, events) = flume::unbounded();
sender
.send(segment_event(0, 0.0, true, SegmentMediaType::Video))
.unwrap();
sender
.send(segment_event(1, 1.0, false, SegmentMediaType::Video))
.unwrap();
let upload = upload_segments(
&transport,
"preparation",
events,
Arc::new(AtomicBool::new(false)),
None,
);
tokio::pin!(upload);
tokio::select! {
_ = transport.preparation_started.notified() => {}
result = &mut upload => panic!("Upload ended before preparation: {result:?}"),
_ = tokio::time::sleep(Duration::from_secs(35)) => panic!("Preparation did not start"),
}
drop(sender);
tokio::time::timeout(Duration::from_secs(1), upload)
.await
.unwrap()
.unwrap();
assert!(transport.preparation_dropped.load(Ordering::Acquire));
assert_eq!(transport.completed.load(Ordering::Acquire), 1);
assert_eq!(transport.manifests.lock().unwrap().last(), Some(&true));
}

fn closed_segment_events() -> flume::Receiver<SegmentCompletedEvent> {
let (sender, receiver) = flume::unbounded();
sender
Expand Down
36 changes: 36 additions & 0 deletions apps/desktop/src-tauri/src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -329,6 +329,42 @@ pub struct Organization {
pub brand_colors: OrganizationBrandColors,
}

pub(crate) async fn prepare_recording_segments(
app: &AppHandle,
video_id: &str,
segments: &[crate::upload::preparation::Segment],
) -> Result<Option<Vec<crate::upload::preparation::Segment>>, AuthedApiError> {
#[derive(Deserialize)]
struct Response {
version: u32,
prepared: Vec<crate::upload::preparation::Segment>,
}

let response = app
.authed_api_request("/api/recording/prepare", |client, url| {
client
.post(url)
.timeout(std::time::Duration::from_secs(20))
.json(&serde_json::json!({ "videoId": video_id, "segments": segments }))
})
.await?;
if matches!(response.status().as_u16(), 404 | 405) {
return Ok(None);
}
if !response.status().is_success() {
return Err(format!(
"Optional recording preparation unavailable ({})",
response.status()
)
.into());
}
let response: Response = crate::upload::lifecycle::cancellable(response.json()).await??;
if response.version != 1 || response.prepared.len() > 32 {
return Ok(None);
}
Ok(Some(response.prepared))
}

pub async fn verify_recording_complete(
app: &AppHandle,
video_id: &str,
Expand Down
8 changes: 8 additions & 0 deletions apps/desktop/src-tauri/src/upload.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ use tokio_util::io::ReaderStream;
use tracing::{Span, debug, error, info, info_span, instrument, trace, warn};

pub(crate) mod lifecycle;
pub(crate) mod preparation;
pub(crate) mod resume;
use tracing_futures::Instrument;

Expand Down Expand Up @@ -1337,6 +1338,12 @@ impl SegmentUploader {
})?;

let state = Arc::new(Mutex::new(SegmentUploadState::new()));
let preparation = preparation::start(
app.clone(),
video_id.clone(),
state.clone(),
session.clone(),
);
let semaphore = Arc::new(tokio::sync::Semaphore::new(6));
let read_semaphore = Arc::new(tokio::sync::Semaphore::new(12));
let consecutive_failures = Arc::new(std::sync::atomic::AtomicU32::new(0));
Expand Down Expand Up @@ -1666,6 +1673,7 @@ impl SegmentUploader {
}

drain_segment_upload_tasks(&state, &mut in_flight).await;
preparation.stop().await;

if bridge_handle.join().is_err() {
state
Expand Down
96 changes: 96 additions & 0 deletions apps/desktop/src-tauri/src/upload/preparation.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
use super::{SegmentUploadState, lifecycle};
use crate::{api, web_api::inherit_upload_context};
use cap_recording::upload_preparation::Preparation;
pub(crate) use cap_recording::upload_preparation::Segment;
use std::{
sync::{Arc, Mutex},
time::Duration,
};
use tauri::AppHandle;

pub(super) struct Task(tokio::task::JoinHandle<()>);

impl Task {
pub(super) async fn stop(mut self) {
self.0.abort();
if let Err(error) = (&mut self.0).await
&& !error.is_cancelled()
{
tracing::warn!(%error, "Optional recording preparation stopped");
}
}
}

impl Drop for Task {
fn drop(&mut self) {
self.0.abort();
}
}

pub(super) fn start(
app: AppHandle,
video_id: String,
state: Arc<Mutex<SegmentUploadState>>,
session: Arc<lifecycle::Session>,
) -> Task {
let context = session.context();
Task(tokio::spawn(inherit_upload_context(context, async move {
let mut preparation = Preparation::default();
let mut next_request = tokio::time::Instant::now();
let mut interval = tokio::time::interval(Duration::from_secs(30));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = session.cancelled() => return,
_ = interval.tick() => {}
}
if tokio::time::Instant::now() < next_request {
continue;
}
let batch = {
let state = state.lock().unwrap_or_else(|error| error.into_inner());
preparation.next_batch(
state.uploaded_video_segments.keys().copied(),
state.uploaded_audio_segments.keys().copied(),
)
};
if batch.is_empty() {
continue;
}
let result = tokio::select! {
_ = session.cancelled() => return,
result = api::prepare_recording_segments(&app, &video_id, &batch) => result,
};
match result {
Ok(Some(prepared)) => {
preparation.acknowledge(&batch, &prepared);
}
Ok(None) => return,
Err(error) => {
preparation.request_failed();
tracing::debug!(%error, "Optional recording preparation unavailable");
}
}
next_request = tokio::time::Instant::now() + preparation.retry_delay();
}
})))
}

#[cfg(test)]
mod tests {
use super::*;

#[tokio::test]
async fn stopping_preparation_cancels_an_in_flight_request() {
let (started, ready) = tokio::sync::oneshot::channel();
let handle = tokio::spawn(async move {
started.send(()).unwrap();
std::future::pending::<()>().await;
});
let abort = handle.abort_handle();
let task = Task(handle);
ready.await.unwrap();
task.stop().await;
assert!(abort.is_finished());
}
}
Loading
Loading