diff --git a/devolutions-gateway/src/api/jrec.rs b/devolutions-gateway/src/api/jrec.rs index 5212eb095..3690204cc 100644 --- a/devolutions-gateway/src/api/jrec.rs +++ b/devolutions-gateway/src/api/jrec.rs @@ -28,6 +28,7 @@ use zip::write::SimpleFileOptions; use crate::DgwState; use crate::api::heartbeat::recording_storage_health; +use crate::artifacts::JrecArtifacts; use crate::extract::{JrecToken, RecordingDeleteScope, RecordingsReadScope}; use crate::http::{HttpError, HttpErrorBuilder}; use crate::recording::{PushOutcome, RecordingMessageSender}; @@ -36,7 +37,7 @@ use crate::token::{JrecTokenClaims, RecordingFileType, RecordingOperation}; /// Read chunk size when streaming a finished session ZIP from the temp file. const ZIP_CHUNK_SIZE: usize = 64 * 1024; -/// Maximum files in a session ZIP (`recording.json` + clips). +/// Maximum files in a session ZIP (`recording.json` + clips and artifacts). /// /// Reconnect windows only mint a small number of clips per session in practice; /// this bound blocks pathological manifests without rejecting normal multi-clip packages. @@ -633,6 +634,8 @@ where #[serde(rename_all = "camelCase")] struct RecordingZipManifest { files: Vec, + #[serde(default)] + artifacts: JrecArtifacts, } #[derive(Debug, Deserialize)] @@ -672,7 +675,7 @@ impl RecordingZipPlan { } } -/// Snapshots `recording.json` and the clip files it references at call time. +/// Snapshots `recording.json` and the clip and artifact files it references at call time. async fn snapshot_recording_zip_plan(recording_dir: &Utf8Path) -> Result { let manifest_path = recording_dir.join("recording.json"); let manifest_bytes = tokio::fs::read(&manifest_path).await.map_err(|error| { @@ -695,22 +698,27 @@ async fn snapshot_recording_zip_plan(recording_dir: &Utf8Path) -> Result, +} + +impl JrecArtifacts { + pub(crate) fn is_empty(&self) -> bool { + let Self { ai_analysis } = self; + ai_analysis.is_empty() + } + + pub(crate) fn into_file_names(self) -> impl IntoIterator { + let Self { ai_analysis } = self; + ai_analysis.into_iter().map(|artifact| artifact.file_name) + } + + pub(crate) fn of_kind_mut(&mut self, kind: ArtifactKind) -> &mut Vec { + match kind { + ArtifactKind::AiAnalysis => &mut self.ai_analysis, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct JrecArtifact { + pub(crate) file_name: String, +} + +/// Kind of a non-recording artifact, used as its key in the manifest `artifacts` object. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ArtifactKind { + AiAnalysis, +} + +impl ArtifactKind { + pub const fn as_str(self) -> &'static str { + match self { + ArtifactKind::AiAnalysis => "ai-analysis", + } + } + + pub(crate) const fn file_type(self) -> RecordingFileType { + match self { + ArtifactKind::AiAnalysis => RecordingFileType::SessionRecordingLog, + } + } + + /// Name of the artifact at `index` in its kind's list, such as `ai-analysis-0.slog`. + pub(crate) fn file_name(self, index: usize) -> String { + format!("{}-{index}.{}", self.as_str(), self.file_type().extension()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn manifest_key_is_the_kind_name() { + let kind = ArtifactKind::AiAnalysis; + let mut artifacts = JrecArtifacts::default(); + artifacts.of_kind_mut(kind).push(JrecArtifact { + file_name: "file".to_owned(), + }); + + let json = serde_json::to_value(&artifacts).expect("serialize artifacts"); + let keys: Vec<_> = json.as_object().expect("object").keys().collect(); + assert_eq!(keys, [kind.as_str()]); + } +} diff --git a/devolutions-gateway/src/lib.rs b/devolutions-gateway/src/lib.rs index 1fb08ff85..cd53900b2 100644 --- a/devolutions-gateway/src/lib.rs +++ b/devolutions-gateway/src/lib.rs @@ -13,6 +13,7 @@ extern crate tracing; pub mod openapi; pub mod api; +pub mod artifacts; pub mod cli; pub mod config; pub mod credential; diff --git a/devolutions-gateway/src/recording.intent.md b/devolutions-gateway/src/recording.intent.md new file mode 100644 index 000000000..14bb0371e --- /dev/null +++ b/devolutions-gateway/src/recording.intent.md @@ -0,0 +1,10 @@ +# Recording +Gateway recording is a continuous stream of bytes pushed from the client who pocesses the valid recording token. +The url is `/jet/jrec/push/{sessionId}?fileType={fileType}`. +We expect the `fileType` to be one of the following file types: `webm`, `cast`, `trp` and `slog`, which must be specified. +When connection is established with request of recordings for the session, if recording is not enabled within a short period of time, the connection will be closed with indication of violation of the recording policy. + + +# Artifacts +Artifacts are files that are not recordings, currently only have `ai-analysis` with combination to `slog` file type. +Currently, we do not have an api endpoint to push artifacts, but we will have one in the future. \ No newline at end of file diff --git a/devolutions-gateway/src/recording.rs b/devolutions-gateway/src/recording.rs index f07913868..59905df5b 100644 --- a/devolutions-gateway/src/recording.rs +++ b/devolutions-gateway/src/recording.rs @@ -20,6 +20,7 @@ use typed_builder::TypedBuilder; use uuid::Uuid; use video_streamer::SignalWriter; +use crate::artifacts::{ArtifactKind, JrecArtifact, JrecArtifacts}; use crate::job_queue::JobQueueHandle; use crate::session::SessionMessageSender; use crate::token::{JrecTokenClaims, RecordingFileType}; @@ -42,6 +43,8 @@ struct JrecManifest { start_time: i64, duration: i64, files: Vec, + #[serde(default, skip_serializing_if = "JrecArtifacts::is_empty")] + artifacts: JrecArtifacts, } impl JrecManifest { @@ -237,7 +240,6 @@ pub enum OnGoingRecordingState { #[derive(Debug, Clone)] struct OnGoingRecording { state: OnGoingRecordingState, - manifest: JrecManifest, manifest_path: Utf8PathBuf, session_must_be_recorded: bool, disconnected_ttl: Duration, @@ -250,6 +252,11 @@ enum RecordingManagerMessage { disconnected_ttl: Duration, channel: oneshot::Sender, }, + AddArtifact { + id: Uuid, + kind: ArtifactKind, + channel: oneshot::Sender>, + }, Disconnect { id: Uuid, }, @@ -288,6 +295,11 @@ impl fmt::Debug for RecordingManagerMessage { .field("file_type", file_type) .field("disconnected_ttl", disconnected_ttl) .finish_non_exhaustive(), + RecordingManagerMessage::AddArtifact { id, kind, channel: _ } => f + .debug_struct("AddArtifact") + .field("id", id) + .field("kind", kind) + .finish_non_exhaustive(), RecordingManagerMessage::Disconnect { id } => f.debug_struct("Disconnect").field("id", id).finish(), RecordingManagerMessage::GetState { id, channel: _ } => { f.debug_struct("GetState").field("id", id).finish_non_exhaustive() @@ -340,6 +352,19 @@ impl RecordingMessageSender { .context("couldn't receive recording file path for this recording") } + /// Adds an empty artifact of `kind` to the session and returns its path, for the caller to write into. + /// + /// Fails when the session has no recording. + pub async fn add_artifact(&self, id: Uuid, kind: ArtifactKind) -> anyhow::Result { + let (tx, rx) = oneshot::channel(); + self.channel + .send(RecordingManagerMessage::AddArtifact { id, kind, channel: tx }) + .await + .ok() + .context("couldn't send AddArtifact message")?; + rx.await.context("couldn't receive AddArtifact result")? + } + async fn disconnect(&self, id: Uuid) -> anyhow::Result<()> { self.channel .send(RecordingManagerMessage::Disconnect { id }) @@ -519,7 +544,7 @@ impl RecordingManagerTask { let recording_path = self.recordings_path.join(id.to_string()); let manifest_path = recording_path.join("recording.json"); - let (manifest, recording_file) = if recording_path.exists() { + let recording_file = if recording_path.exists() { debug!(path = %recording_path, "Recording directory already exists"); let mut existing_manifest = @@ -541,7 +566,7 @@ impl RecordingManagerTask { .save_to_file(&manifest_path) .context("override existing manifest")?; - (existing_manifest, recording_file) + recording_file } else { debug!(path = %recording_path, "Create recording directory"); @@ -564,13 +589,14 @@ impl RecordingManagerTask { start_time, duration: 0, files: vec![first_file], + artifacts: JrecArtifacts::default(), }; initial_manifest .save_to_file(&manifest_path) .context("write initial manifest to disk")?; - (initial_manifest, recording_file) + recording_file }; let active_recording_count = self.rx.active_recordings.insert(id); @@ -592,7 +618,6 @@ impl RecordingManagerTask { id, OnGoingRecording { state: OnGoingRecordingState::Connected, - manifest, manifest_path, session_must_be_recorded, disconnected_ttl, @@ -625,14 +650,14 @@ impl RecordingManagerTask { ongoing.state = OnGoingRecordingState::LastSeen { timestamp: end_time }; - let current_file = ongoing - .manifest - .files - .last_mut() - .context("no recording file (this is a bug)")?; + // Re-read from disk: an artifact may have been added since this recording connected. + let mut manifest = JrecManifest::read_from_file(&ongoing.manifest_path) + .with_context(|| format!("read manifest at {}", ongoing.manifest_path))?; + + let current_file = manifest.files.last_mut().context("no recording file (this is a bug)")?; current_file.duration = end_time - current_file.start_time; - ongoing.manifest.duration = end_time - ongoing.manifest.start_time; + manifest.duration = end_time - manifest.start_time; let recording_file_path = ongoing .manifest_path @@ -642,8 +667,7 @@ impl RecordingManagerTask { debug!(path = %ongoing.manifest_path, "Write updated manifest to disk"); - ongoing - .manifest + manifest .save_to_file(&ongoing.manifest_path) .with_context(|| format!("write manifest at {}", ongoing.manifest_path))?; @@ -675,6 +699,34 @@ impl RecordingManagerTask { Ok(()) } + async fn handle_add_artifact(&self, id: Uuid, kind: ArtifactKind) -> anyhow::Result { + let recording_path = self.recordings_path.join(id.to_string()); + let manifest_path = recording_path.join("recording.json"); + + // A recording deleted while its artifact was being generated must stay deleted. + if !manifest_path.exists() { + anyhow::bail!("session {id} has no recording"); + } + + let mut manifest = JrecManifest::read_from_file(&manifest_path).context("read manifest from disk")?; + + let artifacts = manifest.artifacts.of_kind_mut(kind); + let file_name = kind.file_name(artifacts.len()); + let artifact_path = recording_path.join(&file_name); + + fs::File::create(&artifact_path) + .await + .with_context(|| format!("create {artifact_path}"))?; + + artifacts.push(JrecArtifact { file_name }); + + manifest + .save_to_file(&manifest_path) + .context("write manifest to disk")?; + + Ok(artifact_path) + } + fn handle_remove(&mut self, id: Uuid) { if let Some(ongoing) = self.ongoing_recordings.get(&id) { let now = time::OffsetDateTime::now_utc().unix_timestamp(); @@ -802,6 +854,9 @@ async fn recording_manager_task( Err(e) => error!(error = format!("{e:#}"), "handle_connect"), } }, + RecordingManagerMessage::AddArtifact { id, kind, channel } => { + let _ = channel.send(manager.handle_add_artifact(id, kind).await); + }, RecordingManagerMessage::Disconnect { id } => { if let Err(e) = manager.handle_disconnect(id).await { error!(error = format!("{e:#}"), "handle_disconnect"); @@ -852,14 +907,18 @@ async fn recording_manager_task( Some(recording) => { let recordings_folder = recording.manifest_path.parent().expect("a parent"); - let files = recording - .manifest - .files - .iter() - .map(|file| recordings_folder.join(&file.file_name)) - .collect(); + match JrecManifest::read_from_file(&recording.manifest_path) { + Ok(manifest) => { + let files = manifest + .files + .iter() + .map(|file| recordings_folder.join(&file.file_name)) + .collect(); - let _ = channel.send(files); + let _ = channel.send(files); + } + Err(e) => error!(error = format!("{e:#}"), %id, "Failed to read recording manifest"), + } } None => { warn!(%id, "No recording found for provided ID"); @@ -964,3 +1023,245 @@ async fn remux(input_path: Utf8PathBuf) { Ok(()) } } + +#[cfg(test)] +mod tests { + use devolutions_gateway_task::ShutdownHandle; + use serde_json::json; + + use super::*; + const MASTER_MANIFEST: &str = r#"{ + "sessionId": "22fcd533-5e72-4db7-aa0f-29952dbbca9f", + "startTime": 1, + "duration": 5, + "files": [ + { + "fileName": "recording-0.webm", + "startTime": 1, + "duration": 5 + } + ] +}"#; + + const WEBM: RecordingFileType = RecordingFileType::WebM; + const SLOG: RecordingFileType = RecordingFileType::SessionRecordingLog; + const AI_ANALYSIS: ArtifactKind = ArtifactKind::AiAnalysis; + + struct Harness { + _dir: tempfile::TempDir, + recordings_path: Utf8PathBuf, + sender: RecordingMessageSender, + _shutdown_handle: ShutdownHandle, + } + + impl Harness { + fn start() -> Self { + let dir = tempfile::tempdir().expect("temp dir"); + let recordings_path = Utf8PathBuf::from_path_buf(dir.path().join("recordings")).expect("utf8 path"); + let (sender, receiver) = recording_message_channel(); + // No session runs through this Gateway, as when it is used solely as a recording server. + let (session_manager_handle, _) = crate::session::session_manager_channel(); + let (job_queue_handle, _) = JobQueueHandle::new(); + let task = RecordingManagerTask::new( + receiver, + recordings_path.clone(), + session_manager_handle, + job_queue_handle, + ); + let (shutdown_handle, shutdown_signal) = ShutdownHandle::new(); + tokio::spawn(recording_manager_task(task, shutdown_signal)); + + Self { + _dir: dir, + recordings_path, + sender, + _shutdown_handle: shutdown_handle, + } + } + + fn manifest_path(&self, id: Uuid) -> Utf8PathBuf { + self.recordings_path.join(id.to_string()).join("recording.json") + } + + fn write_manifest(&self, id: Uuid, json: &str) { + std::fs::create_dir_all(self.recordings_path.join(id.to_string())).expect("create recording dir"); + std::fs::write(self.manifest_path(id), json).expect("write manifest"); + } + + fn read_manifest(&self, id: Uuid) -> serde_json::Value { + serde_json::from_slice(&std::fs::read(self.manifest_path(id)).expect("read manifest")) + .expect("parse manifest") + } + + async fn try_connect(&self, id: Uuid, file_type: RecordingFileType) -> anyhow::Result { + self.sender.connect(id, file_type, Duration::from_secs(60)).await + } + + async fn connect(&self, id: Uuid, file_type: RecordingFileType) -> String { + let path = self.try_connect(id, file_type).await.expect("connect"); + path.file_name().expect("file name").to_owned() + } + + async fn disconnect(&self, id: Uuid) { + self.sender.disconnect(id).await.expect("disconnect"); + // Messages are processed in order, so this waits for the disconnect to be handled. + self.sender.get_count().await.expect("sync with manager"); + } + + async fn push(&self, id: Uuid, file_type: RecordingFileType) -> String { + let file_name = self.connect(id, file_type).await; + self.disconnect(id).await; + file_name + } + + async fn add_artifact(&self, id: Uuid, content: &str) -> String { + let path = self.sender.add_artifact(id, AI_ANALYSIS).await.expect("add artifact"); + std::fs::write(&path, content).expect("write artifact"); + path.file_name().expect("file name").to_owned() + } + + fn read_file(&self, id: Uuid, file_name: &str) -> String { + std::fs::read_to_string(self.recordings_path.join(id.to_string()).join(file_name)).expect("read file") + } + } + + #[tokio::test] + async fn slog_recording_stays_in_files() { + let harness = Harness::start(); + let id = Uuid::new_v4(); + + let first = harness.push(id, SLOG).await; + harness.push(id, WEBM).await; + let third = harness.push(id, SLOG).await; + + assert_eq!(first, "recording-0.slog"); + assert_eq!(third, "recording-2.slog"); + let manifest = harness.read_manifest(id); + assert_eq!(manifest["files"][2]["fileName"], "recording-2.slog"); + assert!(manifest.get("artifacts").is_none()); + } + + #[tokio::test] + async fn ai_analysis_goes_to_artifacts() { + let harness = Harness::start(); + let id = Uuid::new_v4(); + + harness.push(id, WEBM).await; + let first_artifact = harness.add_artifact(id, "first").await; + let second_artifact = harness.add_artifact(id, "second").await; + + assert_eq!(first_artifact, "ai-analysis-0.slog"); + assert_eq!(second_artifact, "ai-analysis-1.slog"); + assert_eq!(harness.read_file(id, &first_artifact), "first"); + assert_eq!(harness.read_file(id, &second_artifact), "second"); + let manifest = harness.read_manifest(id); + let files = manifest["files"].as_array().expect("files"); + assert_eq!(files.len(), 1); + assert_eq!(files[0]["fileName"], "recording-0.webm"); + assert_eq!( + manifest["artifacts"], + json!({ "ai-analysis": [{ "fileName": "ai-analysis-0.slog" }, { "fileName": "ai-analysis-1.slog" }] }) + ); + } + + #[tokio::test] + async fn artifact_needs_a_recording() { + let harness = Harness::start(); + + assert!(harness.sender.add_artifact(Uuid::new_v4(), AI_ANALYSIS).await.is_err()); + assert!(!harness.recordings_path.exists()); + } + + #[tokio::test] + async fn artifact_keeps_recording_timing() { + let harness = Harness::start(); + let id = Uuid::new_v4(); + harness.write_manifest(id, MASTER_MANIFEST); + + harness.add_artifact(id, "analysis").await; + + let manifest = harness.read_manifest(id); + assert_eq!(manifest["startTime"], 1); + assert_eq!(manifest["duration"], 5); + assert_eq!(manifest["files"][0]["duration"], 5); + + harness.connect(id, WEBM).await; + assert_eq!(harness.read_manifest(id)["startTime"], 1); + } + + #[tokio::test] + async fn artifact_added_during_a_recording_survives_its_disconnect() { + let harness = Harness::start(); + let id = Uuid::new_v4(); + let expected_artifacts = + json!({ "ai-analysis": [{ "fileName": "ai-analysis-0.slog" }, { "fileName": "ai-analysis-1.slog" }] }); + + harness.connect(id, WEBM).await; + harness.add_artifact(id, "first").await; + harness.add_artifact(id, "second").await; + assert!(harness.try_connect(id, WEBM).await.is_err()); + + harness.disconnect(id).await; + assert_eq!(harness.read_manifest(id)["artifacts"], expected_artifacts); + + assert_eq!(harness.connect(id, WEBM).await, "recording-1.webm"); + harness.disconnect(id).await; + + let manifest = harness.read_manifest(id); + let file_names: Vec<_> = manifest["files"] + .as_array() + .expect("files") + .iter() + .map(|file| file["fileName"].clone()) + .collect(); + assert_eq!(file_names, [json!("recording-0.webm"), json!("recording-1.webm")]); + assert_eq!(manifest["artifacts"], expected_artifacts); + } + + #[tokio::test] + async fn shadow_lists_the_ongoing_recording_files() { + let harness = Harness::start(); + let id = Uuid::new_v4(); + harness.connect(id, WEBM).await; + harness.add_artifact(id, "analysis").await; + + let files = harness.sender.list_files(id).await.expect("list files"); + + let file_names: Vec<_> = files.iter().filter_map(|path| path.file_name()).collect(); + assert_eq!(file_names, ["recording-0.webm"]); + } + + #[tokio::test] + async fn artifact_stays_out_of_the_recording_lifecycle() { + let harness = Harness::start(); + let finished = Uuid::new_v4(); + let recorded = Uuid::new_v4(); + harness.write_manifest(finished, MASTER_MANIFEST); + + harness.add_artifact(finished, "analysis").await; + assert!(harness.sender.get_state(finished).await.expect("get state").is_none()); + assert!(!harness.sender.active_recordings.contains(finished)); + assert_eq!(harness.sender.get_count().await.expect("count"), 0); + + harness.push(recorded, WEBM).await; + harness.add_artifact(recorded, "analysis").await; + assert!(matches!( + harness.sender.get_state(recorded).await.expect("get state"), + Some(OnGoingRecordingState::LastSeen { .. }) + )); + assert!(harness.sender.active_recordings.contains(recorded)); + assert_eq!(harness.sender.get_count().await.expect("count"), 1); + } + + #[test] + fn manifest_without_artifacts_round_trips_byte_for_byte() { + let dir = tempfile::tempdir().expect("temp dir"); + let path = dir.path().join("recording.json"); + std::fs::write(&path, MASTER_MANIFEST).expect("write manifest"); + + let manifest = JrecManifest::read_from_file(&path).expect("read manifest"); + manifest.save_to_file(&path).expect("save manifest"); + + assert_eq!(std::fs::read_to_string(&path).expect("read back"), MASTER_MANIFEST); + } +}