mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-02 05:52:26 +02:00
feat(dvr): recording sidecars, and a shared retention pool that exists
Task 16, steps 4 and 5. A finished recording now has a `<file>.tuliprox-recording.json` beside it holding what the file is: materialization id, media identity, kind, relative path, size, completion time. Nothing else. It lives in the recording directory, which is not a place for owners, quota, headers, URLs or resume validators, and a test walks the encoded fields to keep it that way. The plan names the sidecar `.tuliprox-recording.json`. Taken literally that is one file per directory, and the organised layouts put many recordings in one -- it would describe only whichever finished last. It is suffixed onto the recording's own name instead. Written staged-and-renamed before the repository is told the recording completed, so an operator finding an orphan has the record even when the commit is what failed. Rewriting is normal: finalization is idempotent and the sidecar write has to be too. Failing to write one is logged, not fatal -- refusing to complete a real recording because a descriptive file could not be written trades the thing for the note about it. Startup reports orphans by count. Deliberately not by path, and deliberately without creating anything: a sidecar describes a file, it is not evidence anyone was entitled to play it. The return type carries no principal, so there is nothing to build an entry from. Step 1 turned up a defect. `RetentionOwner::Shared` was a variant nothing constructed: `from_recording_owner` read `meta.owner`, which is always a user, so every shared recording was retained out of its creator's personal budget -- and shared recordings by different creators never shared a pool at all. The doc comment claimed the opposite. Retention now keys on visibility, the same way quota already did.
This commit is contained in:
@@ -21,6 +21,7 @@ pub mod recording_rule_scheduler;
|
||||
pub mod recording_rule_service;
|
||||
pub mod recording_security;
|
||||
pub mod recording_service;
|
||||
pub mod recording_sidecar;
|
||||
pub mod recording_source_resolution;
|
||||
pub mod recording_supervisor;
|
||||
pub mod recording_transfer;
|
||||
|
||||
@@ -5,7 +5,10 @@
|
||||
//! first with a stable task-id tie-break.
|
||||
|
||||
use super::recording_quota::QuotaRecordingTaskView;
|
||||
use shared::model::{recording::RecordingOwner, UserId};
|
||||
use shared::model::{
|
||||
recording::{RecordingMetadata, RecordingVisibility},
|
||||
UserId,
|
||||
};
|
||||
use std::collections::HashMap;
|
||||
|
||||
/// Retention configuration derived from `RecordingRetentionConfig`.
|
||||
@@ -113,9 +116,16 @@ pub fn normalize_channel_name(name: &str) -> String {
|
||||
}
|
||||
|
||||
impl RetentionOwner {
|
||||
pub fn from_recording_owner(owner: &RecordingOwner) -> Self {
|
||||
let RecordingOwner::User(uid) = owner;
|
||||
Self::Private(uid.clone())
|
||||
/// Which retention budget a recording is charged to.
|
||||
///
|
||||
/// Visibility decides it, exactly as it does for quota. Reading the owner
|
||||
/// instead put every shared recording in its creator's personal budget, so
|
||||
/// `Shared` was a variant nothing could produce.
|
||||
pub fn from_metadata(meta: &RecordingMetadata) -> Self {
|
||||
match meta.visibility {
|
||||
RecordingVisibility::Shared => Self::Shared,
|
||||
RecordingVisibility::Private => Self::Private(meta.owner_id().clone()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -153,7 +163,7 @@ fn group_for<V: QuotaRecordingTaskView>(task: &V) -> Option<(RetentionGroupKey,
|
||||
}
|
||||
let completed_at = meta.completed_at?;
|
||||
let channel = ChannelKey::from_metadata(meta.channel_id.as_deref(), meta.channel_name.as_deref());
|
||||
let owner = RetentionOwner::from_recording_owner(&meta.owner);
|
||||
let owner = RetentionOwner::from_metadata(meta);
|
||||
Some((RetentionGroupKey { owner, channel }, completed_at))
|
||||
}
|
||||
|
||||
@@ -539,6 +549,28 @@ mod tests {
|
||||
assert_eq!(candidates.len(), 0, "each library holds one recording and is allowed one");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn private_retention_never_reaches_the_shared_library() {
|
||||
// Step 1: a user's own budget must not evict the shared copy, and the
|
||||
// shared budget must not evict anyone's private one. They are separate
|
||||
// pools that happen to name the same channel.
|
||||
let config = RetentionConfig { keep_last_per_channel: Some(1), delete_after_days: None };
|
||||
let tasks = vec![
|
||||
completed("alice-old", RecordingOwner::User(UserId::from("web:alice")), Some("c1"), Some("Alpha"), 1_000),
|
||||
completed("alice-new", RecordingOwner::User(UserId::from("web:alice")), Some("c1"), Some("Alpha"), 2_000),
|
||||
shared_completed("shared-old", Some("c1"), Some("Alpha"), 1_000),
|
||||
shared_completed("shared-new", Some("c1"), Some("Alpha"), 2_000),
|
||||
];
|
||||
|
||||
let candidates = compute_candidates(&tasks, &config, 3_000);
|
||||
let evicted: Vec<&str> = candidates.iter().map(|candidate| candidate.uuid.as_str()).collect();
|
||||
|
||||
// One over budget in each pool, and each pool gives up its own oldest.
|
||||
assert!(evicted.contains(&"alice-old"), "alice keeps her newest and gives up her oldest");
|
||||
assert!(evicted.contains(&"shared-old"), "the shared library does the same, separately");
|
||||
assert_eq!(evicted.len(), 2, "neither pool spends the other's budget");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn shared_owner_groups_only_by_channel() {
|
||||
// Two shared recordings on the same channel — count
|
||||
|
||||
@@ -0,0 +1,229 @@
|
||||
//! The physical record written beside a finished recording.
|
||||
//!
|
||||
//! A sidecar describes the *file*, never the people who asked for it. It exists
|
||||
//! so an operator, or a rebuild, can tell what an orphaned recording on disk
|
||||
//! actually is. It is not an access grant: rights live only in the repository,
|
||||
//! and a sidecar found on disk can rediscover a file but never re-create the
|
||||
//! library entry that would let somebody play it.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use shared::model::RecordingKind;
|
||||
use std::{
|
||||
io,
|
||||
path::{Path, PathBuf},
|
||||
};
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
/// Suffix appended to the recording's own filename.
|
||||
///
|
||||
/// The organised layouts put many recordings in one directory, so a single
|
||||
/// fixed name per directory would describe only whichever finished last.
|
||||
pub const SIDECAR_SUFFIX: &str = ".tuliprox-recording.json";
|
||||
|
||||
/// What a finished recording is, as told by the file next to it.
|
||||
///
|
||||
/// Every field here is a physical fact. Owner, visibility, quota, headers, URL
|
||||
/// and resume validators are deliberately absent: this file sits in the
|
||||
/// recording directory, which is not a place to put anything user-specific or
|
||||
/// secret.
|
||||
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct RecordingSidecar {
|
||||
pub materialization_id: String,
|
||||
pub media_identity: String,
|
||||
pub kind: RecordingKind,
|
||||
pub relative_path: String,
|
||||
pub size_bytes: u64,
|
||||
pub completed_at: i64,
|
||||
}
|
||||
|
||||
/// The sidecar that belongs to a recording file.
|
||||
pub fn sidecar_path(recording: &Path) -> PathBuf {
|
||||
let mut name = recording.as_os_str().to_os_string();
|
||||
name.push(SIDECAR_SUFFIX);
|
||||
PathBuf::from(name)
|
||||
}
|
||||
|
||||
/// `true` when this path is a sidecar rather than a recording.
|
||||
pub fn is_sidecar(path: &Path) -> bool {
|
||||
path.file_name().and_then(|name| name.to_str()).is_some_and(|name| name.ends_with(SIDECAR_SUFFIX))
|
||||
}
|
||||
|
||||
/// Write the sidecar beside its recording, replacing any earlier one.
|
||||
///
|
||||
/// Staged and renamed so a crash mid-write cannot leave a half-parsed file
|
||||
/// where a valid one used to be. Rewriting an existing sidecar is normal: a
|
||||
/// finalization that runs twice must not fail the second time.
|
||||
pub async fn write_sidecar(recording: &Path, sidecar: &RecordingSidecar) -> io::Result<PathBuf> {
|
||||
let target = sidecar_path(recording);
|
||||
let staging = sidecar_path(recording).with_extension("writing");
|
||||
let encoded = serde_json::to_vec_pretty(sidecar).map_err(io::Error::other)?;
|
||||
|
||||
let mut file = tokio::fs::File::create(&staging).await?;
|
||||
file.write_all(&encoded).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
tokio::fs::rename(&staging, &target).await?;
|
||||
Ok(target)
|
||||
}
|
||||
|
||||
/// Read one sidecar. A file that does not parse is not a sidecar.
|
||||
pub async fn read_sidecar(path: &Path) -> io::Result<RecordingSidecar> {
|
||||
let bytes = tokio::fs::read(path).await?;
|
||||
serde_json::from_slice(&bytes).map_err(io::Error::other)
|
||||
}
|
||||
|
||||
/// Every sidecar under `root` whose materialization is not already known.
|
||||
///
|
||||
/// Returns physical descriptions only. Nothing here creates a library entry or
|
||||
/// grants anyone access -- that is the point of the return type: an orphan is
|
||||
/// something for an operator to look at, not something to hand to a user.
|
||||
/// Unreadable or unparseable files are skipped rather than failing the scan;
|
||||
/// one bad file must not hide every good one.
|
||||
pub async fn scan_orphans(root: &Path, known_materializations: &[String]) -> io::Result<Vec<RecordingSidecar>> {
|
||||
let mut found = Vec::new();
|
||||
let mut directories = vec![root.to_path_buf()];
|
||||
while let Some(directory) = directories.pop() {
|
||||
let mut entries = match tokio::fs::read_dir(&directory).await {
|
||||
Ok(entries) => entries,
|
||||
Err(error) if error.kind() == io::ErrorKind::NotFound => continue,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
while let Some(entry) = entries.next_entry().await? {
|
||||
let path = entry.path();
|
||||
// `file_type` does not follow symlinks, so a link cannot walk the
|
||||
// scan out of the recording root.
|
||||
let file_type = entry.file_type().await?;
|
||||
if file_type.is_dir() {
|
||||
directories.push(path);
|
||||
} else if file_type.is_file() && is_sidecar(&path) {
|
||||
if let Ok(sidecar) = read_sidecar(&path).await {
|
||||
if !known_materializations.iter().any(|known| known == &sidecar.materialization_id) {
|
||||
found.push(sidecar);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
found.sort_by(|left, right| left.materialization_id.cmp(&right.materialization_id));
|
||||
Ok(found)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{is_sidecar, read_sidecar, scan_orphans, sidecar_path, write_sidecar, RecordingSidecar};
|
||||
use shared::model::RecordingKind;
|
||||
use std::path::Path;
|
||||
use tempfile::TempDir;
|
||||
|
||||
fn sidecar(id: &str) -> RecordingSidecar {
|
||||
RecordingSidecar {
|
||||
materialization_id: id.to_string(),
|
||||
media_identity: format!("programme-{id}"),
|
||||
kind: RecordingKind::Vod,
|
||||
relative_path: format!("{id}.mp4"),
|
||||
size_bytes: 4_096,
|
||||
completed_at: 1_700_000_000,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_sidecar_is_named_after_its_recording() {
|
||||
// A fixed name per directory would describe only the last recording to
|
||||
// finish, and the organised layouts share directories.
|
||||
let first = sidecar_path(Path::new("/rec/Channel/one.ts"));
|
||||
let second = sidecar_path(Path::new("/rec/Channel/two.ts"));
|
||||
assert_ne!(first, second, "two recordings in one directory need two sidecars");
|
||||
assert!(is_sidecar(&first));
|
||||
assert!(!is_sidecar(Path::new("/rec/Channel/one.ts")));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_sidecar_round_trips() {
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
let recording = dir.path().join("film.mp4");
|
||||
let written = write_sidecar(&recording, &sidecar("mat-a")).await.expect("write");
|
||||
assert_eq!(read_sidecar(&written).await.expect("read"), sidecar("mat-a"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn writing_twice_replaces_rather_than_fails() {
|
||||
// Finalization is idempotent, so the sidecar write has to be too.
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
let recording = dir.path().join("film.mp4");
|
||||
write_sidecar(&recording, &sidecar("mat-a")).await.expect("first");
|
||||
let mut grown = sidecar("mat-a");
|
||||
grown.size_bytes = 8_192;
|
||||
let written = write_sidecar(&recording, &grown).await.expect("second");
|
||||
assert_eq!(read_sidecar(&written).await.expect("read").size_bytes, 8_192);
|
||||
assert!(!sidecar_path(&recording).with_extension("writing").exists(), "no staging file is left behind");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_sidecar_carries_no_user_or_transport_detail() {
|
||||
// It lives in the recording directory, which is not a place for owners,
|
||||
// credentials or provider URLs.
|
||||
let encoded = serde_json::to_string(&sidecar("mat-a")).expect("encode");
|
||||
for forbidden in ["owner", "user", "visibility", "quota", "header", "url", "etag", "token", "password"] {
|
||||
assert!(!encoded.to_lowercase().contains(forbidden), "{forbidden} must not appear in {encoded}");
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_unknown_sidecar_is_reported_as_an_orphan() {
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
let nested = dir.path().join("Channel/Season 01");
|
||||
std::fs::create_dir_all(&nested).expect("dirs");
|
||||
write_sidecar(&dir.path().join("known.mp4"), &sidecar("mat-known")).await.expect("write");
|
||||
write_sidecar(&nested.join("lost.mp4"), &sidecar("mat-lost")).await.expect("write");
|
||||
|
||||
let orphans = scan_orphans(dir.path(), &["mat-known".to_string()]).await.expect("scan");
|
||||
|
||||
assert_eq!(orphans.len(), 1, "only the file the repository does not know about");
|
||||
assert_eq!(orphans[0].materialization_id, "mat-lost");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_file_that_is_not_a_sidecar_is_ignored() {
|
||||
// Including one that is unreadable: a single bad file must not hide
|
||||
// every good one from the operator.
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
std::fs::write(dir.path().join("film.mp4"), b"not json").expect("recording");
|
||||
std::fs::write(dir.path().join("broken.mp4.tuliprox-recording.json"), b"{ not json").expect("broken");
|
||||
write_sidecar(&dir.path().join("good.mp4"), &sidecar("mat-good")).await.expect("write");
|
||||
|
||||
let orphans = scan_orphans(dir.path(), &[]).await.expect("scan");
|
||||
|
||||
assert_eq!(orphans.len(), 1);
|
||||
assert_eq!(orphans[0].materialization_id, "mat-good");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rediscovering_an_orphan_describes_the_file_and_nothing_more() {
|
||||
// The whole risk of scanning the recording directory is that a file
|
||||
// found there becomes a way in. A sidecar carries no principal, so
|
||||
// there is nothing here to build a library entry from: an orphan can be
|
||||
// identified and no more. Access comes from the repository or not at
|
||||
// all.
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
write_sidecar(&dir.path().join("lost.mp4"), &sidecar("mat-lost")).await.expect("write");
|
||||
|
||||
let orphans = scan_orphans(dir.path(), &[]).await.expect("scan");
|
||||
|
||||
let found = &orphans[0];
|
||||
let encoded = serde_json::to_value(found).expect("encode");
|
||||
let fields: Vec<&str> = encoded.as_object().expect("object").keys().map(String::as_str).collect();
|
||||
assert_eq!(
|
||||
fields,
|
||||
["materialization_id", "media_identity", "kind", "relative_path", "size_bytes", "completed_at"],
|
||||
"a sidecar describes a file; anything more would be a claim about who may have it"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn scanning_a_missing_root_is_not_an_error() {
|
||||
let dir = TempDir::new().expect("tempdir");
|
||||
let orphans = scan_orphans(&dir.path().join("nothing-here"), &[]).await.expect("scan");
|
||||
assert!(orphans.is_empty(), "an absent recording directory just has no orphans");
|
||||
}
|
||||
}
|
||||
@@ -47,6 +47,7 @@ pub async fn run_startup_reconciliation(ctx: &RecordingCtx) {
|
||||
}
|
||||
let stuck = recover_stuck_deletions(ctx).await;
|
||||
let drift = reconcile_rule_drift(ctx).await;
|
||||
report_orphan_recordings(ctx).await;
|
||||
SupervisorHealth::stamp(&supervisor_health().reconciliation_last_run, now_ts());
|
||||
if stuck > 0 || drift > 0 {
|
||||
info!("DVR startup reconciliation: repaired {stuck} interrupted deletion(s), {drift} rule drift item(s)");
|
||||
@@ -67,6 +68,40 @@ fn media_is_still_referenced(all_tasks: &[RecordingTask], subject: &RecordingTas
|
||||
})
|
||||
}
|
||||
|
||||
/// Tell the operator about recordings on disk the repository does not know.
|
||||
///
|
||||
/// Reporting only. A sidecar describes a file; it is not evidence that anyone
|
||||
/// was ever entitled to play it, so rediscovering one must never put it back in
|
||||
/// somebody's library. Access comes from the repository or not at all.
|
||||
async fn report_orphan_recordings(ctx: &RecordingCtx) {
|
||||
let Some(root) = recording_config(&ctx.app_config)
|
||||
.map(|cfg| PathBuf::from(cfg.directory))
|
||||
.filter(|dir| !dir.as_os_str().is_empty())
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let (_revision, tasks) = ctx.recordings.committed_snapshot().await;
|
||||
let known: Vec<String> = tasks
|
||||
.iter()
|
||||
.map(|task| {
|
||||
let persisted = crate::recording::recording_queue::RecordingQueue::to_persisted(task);
|
||||
tuliprox_repository::recording_repository::materialization_id_for(&persisted)
|
||||
})
|
||||
.collect();
|
||||
match crate::recording::recording_sidecar::scan_orphans(&root, &known).await {
|
||||
Ok(orphans) if orphans.is_empty() => {}
|
||||
Ok(orphans) => {
|
||||
// Count only: the paths belong to whoever recorded them.
|
||||
warn!(
|
||||
target: "recording::audit",
|
||||
"recording_orphan_files: {} recording(s) on disk are not in the library and are not playable through it",
|
||||
orphans.len()
|
||||
);
|
||||
}
|
||||
Err(error) => warn!("Could not scan the recording directory for orphans: {error}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// Finish or undo every deletion the previous process left half-done.
|
||||
async fn recover_stuck_deletions(ctx: &RecordingCtx) -> usize {
|
||||
let (_revision, all_tasks) = ctx.recordings.committed_snapshot().await;
|
||||
|
||||
@@ -19,6 +19,7 @@ use crate::{
|
||||
mutate_optional, PersistedRecordingTask, QueueMutationError, RecordingControl, RecordingQueue,
|
||||
RecordingTask, RecordingTaskState, RecordingWaitOutcome,
|
||||
},
|
||||
recording_sidecar,
|
||||
recording_url::build_stable_recording_url,
|
||||
recording_worker::{recording_partial_path, run_recording, RecordingExecutionResult},
|
||||
},
|
||||
@@ -987,6 +988,32 @@ async fn requeue_active_download_for_capacity_wait(
|
||||
Ok(result.unwrap_or(false))
|
||||
}
|
||||
|
||||
/// Write the sidecar for the recording that just finished.
|
||||
///
|
||||
/// Best effort by design: the recording is on disk and committed either way,
|
||||
/// and refusing to complete a transfer because a descriptive file could not be
|
||||
/// written would trade a real recording for a diagnostic.
|
||||
async fn write_completion_sidecar(download_queue: &RecordingQueue, uuid: &str, measured_bytes: u64) {
|
||||
let Some(task) = download_queue.active.read().await.as_ref().filter(|task| task.uuid == uuid).cloned() else {
|
||||
return;
|
||||
};
|
||||
let Some(relative_path) = task.recording.relative_path.clone() else {
|
||||
return;
|
||||
};
|
||||
let persisted = RecordingQueue::to_persisted(&task);
|
||||
let sidecar = recording_sidecar::RecordingSidecar {
|
||||
materialization_id: tuliprox_repository::recording_repository::materialization_id_for(&persisted),
|
||||
media_identity: persisted.media_identity,
|
||||
kind: task.kind,
|
||||
relative_path,
|
||||
size_bytes: measured_bytes,
|
||||
completed_at: chrono::Utc::now().timestamp(),
|
||||
};
|
||||
if let Err(error) = recording_sidecar::write_sidecar(&task.file_path, &sidecar).await {
|
||||
warn!("Could not write the sidecar for {}: {error}", task.file_path.display());
|
||||
}
|
||||
}
|
||||
|
||||
async fn promote_next_download(
|
||||
download_queue: &RecordingQueue,
|
||||
) -> Result<Option<(String, String)>, QueueMutationError> {
|
||||
@@ -1418,6 +1445,11 @@ pub async fn ensure_recording_worker_running(
|
||||
None => 0,
|
||||
}
|
||||
};
|
||||
// Beside the file, before the repository is told
|
||||
// it is complete: an operator finding an orphan
|
||||
// needs the record even if the commit is what
|
||||
// failed.
|
||||
write_completion_sidecar(&dq, &worker_uuid, measured_bytes).await;
|
||||
let committed = finish_active_and_promote(&dq, &worker_uuid, |fd| {
|
||||
fd.finished = true;
|
||||
fd.paused = false;
|
||||
|
||||
@@ -331,7 +331,7 @@ impl StoredRecords {
|
||||
/// one file. Distinct from the entry id, which is what a user addresses.
|
||||
/// Falls back to the task's own uuid when no identity was supplied, which
|
||||
/// keeps such a task on a file of its own rather than colliding with others.
|
||||
fn materialization_id_for(task: &PersistedRecordingTask) -> String {
|
||||
pub fn materialization_id_for(task: &PersistedRecordingTask) -> String {
|
||||
if task.media_identity.is_empty() {
|
||||
return format!("mat-uuid:{}", task.uuid);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user