Merge pull request 'feat: améliorer la gestion des métadonnées de flux en continu avec dc:date' (#84) from push-zmzoxmlkqmmz into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 10m14s

Reviewed-on: #84
This commit was merged in pull request #84.
This commit is contained in:
2026-03-26 00:47:12 +01:00
8 changed files with 82 additions and 54 deletions

View File

@@ -25,6 +25,14 @@ host:
friendly_name_prefix: "PMOMusic-Dev1"
```
## Méthode de travail
Avant de modifier quoi que ce soit sur une fonctionnalité non triviale :
1. Lire et comprendre le flux complet des données concernées, de bout en bout
2. Identifier précisément où ça casse et pourquoi
3. Faire une seule modification ciblée
Ne pas avancer par tâtonnements ("vibe programming") — cela produit des allers-retours, des bugs introduits puis annulés, et du code inutilement compliqué.
## Architecture
- Projet Rust multi-crates avec workspaces
- Crates principales : pmoupnp, pmomediaserver, pmomediarenderer, pmoconfig

2
Cargo.lock generated
View File

@@ -4,7 +4,7 @@ version = 4
[[package]]
name = "PMOMusic"
version = "0.3.29"
version = "0.3.30"
dependencies = [
"axum 0.8.7",
"console-subscriber",

View File

@@ -1,6 +1,6 @@
[package]
name = "PMOMusic"
version = "0.3.29"
version = "0.3.30"
edition = "2024"
[dependencies]

View File

@@ -37,7 +37,7 @@ tokio-util = { version = "0.7", optional = true }
async-trait = { version = "0.1", optional = true }
tokio-stream = { version = "0.1", features = ["sync"], optional = true }
async-stream = { version = "0.3", optional = true }
chrono = { version = "0.4", features = ["serde"], optional = true }
chrono = { version = "0.4", features = ["serde"] }
[dev-dependencies]
percent-encoding = "2.3"
@@ -45,4 +45,4 @@ percent-encoding = "2.3"
[features]
default = []
# Active l'API REST pmoserver
pmoserver = ["dep:pmoserver", "dep:utoipa", "dep:axum", "dep:tokio", "dep:tokio-util", "dep:async-trait", "dep:tokio-stream", "dep:async-stream", "dep:chrono"]
pmoserver = ["dep:pmoserver", "dep:utoipa", "dep:axum", "dep:tokio", "dep:tokio-util", "dep:async-trait", "dep:tokio-stream", "dep:async-stream"]

View File

@@ -1360,12 +1360,21 @@ impl MusicRenderer {
pub fn set_last_metadata(&self, metadata: Option<TrackMetadata>) {
let mut state = self.state.lock().unwrap();
let metadata_changed = state.last_metadata != metadata;
state.last_metadata = metadata;
if metadata_changed {
state.track_start_time = Some(SystemTime::now());
// Pour les flux continus: utiliser dc:date comme track_start_time réel de diffusion.
// Toute source stream peut encoder son start_time en RFC 3339 dans dc:date.
// Fallback sur now() si absent ou non parseable.
let start_time = metadata
.as_ref()
.filter(|m| m.is_continuous_stream)
.and_then(|m| m.date.as_deref())
.and_then(parse_rfc3339_to_system_time)
.unwrap_or_else(SystemTime::now);
state.track_start_time = Some(start_time);
// Reset duration cache when track changes (for streams)
state.current_track_duration = None;
}
state.last_metadata = metadata;
}
/// Gets the timestamp when the current track started playing.
@@ -1785,6 +1794,20 @@ fn parse_didl_duration(didl_xml: &str) -> Option<String> {
duration
}
/// Parse un datetime RFC 3339 en SystemTime.
///
/// Utilisé pour calibrer track_start_time sur le début réel de diffusion
/// d'un segment, encodé dans dc:date du DIDL par la source stream.
fn parse_rfc3339_to_system_time(s: &str) -> Option<SystemTime> {
let dt = chrono::DateTime::parse_from_rfc3339(s).ok()?;
let secs = dt.timestamp();
if secs < 0 {
return None;
}
Some(std::time::UNIX_EPOCH + std::time::Duration::from_secs(secs as u64))
}
/// Transport control façade that dispatches to whichever backend can fulfill
/// the request, returning a standardized error if the backend lacks support.
impl TransportControl for MusicRendererBackend {

View File

@@ -54,10 +54,11 @@ pub struct CachedMetadata {
pub protocol_info: String,
pub sample_frequency: Option<String>,
pub nr_audio_channels: Option<String>,
pub duration: Option<String>, // Calculé depuis end_time
pub duration: Option<String>, // Calculé depuis end_time - start_time
// TTL = end_time de l'API Radio France
pub end_time: Option<u64>, // Unix timestamp
// Bornes temporelles de l'émission/morceau courant
pub start_time: Option<u64>, // Unix timestamp
pub end_time: Option<u64>, // Unix timestamp (sert aussi de TTL)
/// Timestamp de la dernière récupération (pour TTL minimum)
pub fetched_at: u64,
@@ -89,7 +90,8 @@ impl CachedMetadata {
let (stream_url, protocol_info, sample_frequency, nr_audio_channels, duration) =
Self::build_stream_resource(live, &station.slug, server_base_url);
// 4. TTL = end_time
// 4. Bornes temporelles + TTL
let start_time = live.now.start_time;
let end_time = live.now.end_time;
let fetched_at = SystemTime::now()
.duration_since(UNIX_EPOCH)
@@ -111,6 +113,7 @@ impl CachedMetadata {
sample_frequency,
nr_audio_channels,
duration,
start_time,
end_time,
fetched_at,
})
@@ -294,24 +297,16 @@ impl CachedMetadata {
Option<String>,
Option<String>,
) {
// Calculer la durée restante (maintenant -> end_time)
let duration = if let Some(end) = metadata.now.end_time {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs();
if end > now {
let duration_secs = end - now;
// Durée totale = end_time - start_time (valeur fixe, indépendante de now)
let duration = match (metadata.now.start_time, metadata.now.end_time) {
(Some(start), Some(end)) if end > start => {
let duration_secs = end - start;
let hours = duration_secs / 3600;
let minutes = (duration_secs % 3600) / 60;
let seconds = duration_secs % 60;
Some(format!("{}:{:02}:{:02}", hours, minutes, seconds))
} else {
None
}
} else {
None
_ => None,
};
// URL du proxy PMOMusic
@@ -368,45 +363,29 @@ impl CachedMetadata {
/// La playlist et l'item ont EXACTEMENT les mêmes métadonnées
#[cfg(feature = "playlist")]
pub fn to_didl(&self, playlist_id: &str, parent_id: &str) -> Container {
// Calculer la duration dynamiquement (temps restant jusqu'à end_time)
let duration = if let Some(end) = self.end_time {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs();
if end > now {
let duration_secs = end - now;
// Durée totale = end_time - start_time (valeur fixe)
let duration = match (self.start_time, self.end_time) {
(Some(start), Some(end)) if end > start => {
let duration_secs = end - start;
let hours = duration_secs / 3600;
let minutes = (duration_secs % 3600) / 60;
let seconds = duration_secs % 60;
let dur = format!("{}:{:02}:{:02}", hours, minutes, seconds);
#[cfg(feature = "logging")]
tracing::debug!(
"Duration calculated for {}: {} (end_time: {}, now: {}, remaining: {}s)",
self.slug,
dur,
end,
now,
duration_secs
"Duration for {}: {} (start: {}, end: {}, total: {}s)",
self.slug, dur, start, end, duration_secs
);
Some(dur)
} else {
}
_ => {
#[cfg(feature = "logging")]
tracing::warn!(
"Duration expired for {}: end_time {} < now {}",
self.slug,
end,
now
"No duration for {}: start_time={:?}, end_time={:?}",
self.slug, self.start_time, self.end_time
);
None
}
} else {
#[cfg(feature = "logging")]
tracing::warn!("No end_time for {}, duration will be None", self.slug);
None
};
let item = Item {
@@ -421,7 +400,13 @@ impl CachedMetadata {
genre: self.genre.clone(),
album_art: self.album_art.clone(),
album_art_pk: self.album_art_pk.clone(),
date: None,
// dc:date: début réel de diffusion du segment en RFC 3339 UTC.
// Les timestamps Radio France sont des Unix timestamps UTC.
// Pmocontrol utilise cette valeur pour calibrer track_start_time.
date: self.start_time.and_then(|start| {
chrono::DateTime::from_timestamp(start as i64, 0)
.map(|dt| dt.to_rfc3339())
}),
original_track_number: None,
resources: vec![Resource {
protocol_info: self.protocol_info.clone(),
@@ -797,6 +782,7 @@ mod tests {
sample_frequency: None,
nr_audio_channels: None,
duration: None,
start_time: None,
end_time: Some(0), // Dans le passé
fetched_at: now,
};

View File

@@ -49,10 +49,21 @@ pub struct PullResponse {
impl PullResponse {
/// Get the current step at depth 1 (the "now playing" item)
pub fn current_step(&self) -> Option<&PullStep> {
// levels[0].items[0] is the current item at depth 1 (most relevant)
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
// Le niveau liste tous les steps de la session (passés et futurs) dans l'ordre.
// On cherche le step dont start <= now < end.
let level = self.levels.first()?;
let step_id = level.items.first()?;
self.steps.get(step_id)
level.items.iter().find_map(|step_id| {
let step = self.steps.get(step_id)?;
match (step.start, step.end) {
(Some(start), Some(end)) if start <= now && now < end => Some(step),
_ => None,
}
})
}
}

View File

@@ -1 +1 @@
0.3.29
0.3.30