Ajout du reporting de position audio

Ajoute un reporting de la position audio toutes les secondes depuis le navigateur vers le serveur pour une meilleure synchronisation.

- Implémente un intervalle pour envoyer la position courante et la durée de l'audio
- Ajoute les endpoints POST /api/webrenderer/{id}/position
- Met à jour la logique de gestion de la position dans le registre
- Supprime le tracking de position dans le pipeline audio (déplacé côté client)
This commit is contained in:
2026-02-28 12:45:55 +01:00
parent f9056f92e0
commit ca8702fd7f
5 changed files with 63 additions and 27 deletions

View File

@@ -71,6 +71,30 @@ export function useWebRenderer() {
let onConnectedCallback: (() => void) | null = null;
let sseUnsubscribe: (() => void) | null = null;
let pendingCanPlay: (() => void) | null = null;
let positionInterval: ReturnType<typeof setInterval> | null = null;
// ── Reporting de position ─────────────────────────────────────────────────
function startPositionReporting(): void {
stopPositionReporting();
positionInterval = setInterval(() => {
if (!audioEl || !instanceId) return;
const pos = audioEl.currentTime;
const dur = isFinite(audioEl.duration) ? audioEl.duration : null;
void fetch(`/api/webrenderer/${instanceId}/position`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ position_sec: pos, duration_sec: dur }),
});
}, 1000);
}
function stopPositionReporting(): void {
if (positionInterval !== null) {
clearInterval(positionInterval);
positionInterval = null;
}
}
// ── Connexion / déconnexion du flux audio ─────────────────────────────────
@@ -124,6 +148,7 @@ export function useWebRenderer() {
el.removeEventListener("canplay", onCanPlay);
el.play().then(() => {
console.log("[WebRenderer] play() resolved OK");
startPositionReporting();
}).catch((e: unknown) => {
console.warn("[WebRenderer] play() rejected:", e);
});
@@ -134,6 +159,7 @@ export function useWebRenderer() {
function stopStream(): void {
console.log("[WebRenderer] stopStream called, readyState=", audioEl?.readyState);
stopPositionReporting();
if (!audioEl) return;
if (pendingCanPlay) {
audioEl.removeEventListener("canplay", pendingCanPlay);
@@ -230,6 +256,7 @@ export function useWebRenderer() {
onUnmounted(() => {
sseUnsubscribe?.();
sseUnsubscribe = null;
stopPositionReporting();
stopStream();
void unregister();
if (audioEl) {

View File

@@ -15,7 +15,7 @@ use pmocontrol::ControlPoint;
#[cfg(feature = "pmoserver")]
use crate::error::WebRendererError;
#[cfg(feature = "pmoserver")]
use crate::register::{register_handler, unregister_handler};
use crate::register::{position_update_handler, register_handler, unregister_handler};
#[cfg(feature = "pmoserver")]
use crate::registry::RendererRegistry;
#[cfg(feature = "pmoserver")]
@@ -48,9 +48,10 @@ impl WebRendererExt for pmoserver::Server {
)
.await;
// GET /api/webrenderer/{id}/stream + DELETE /api/webrenderer/{id}
// GET /api/webrenderer/{id}/stream + DELETE /api/webrenderer/{id} + POST /api/webrenderer/{id}/position
let dynamic_router = Router::new()
.route("/{id}/stream", get(stream_handler))
.route("/{id}/position", post(position_update_handler))
.route("/{id}", delete(unregister_handler))
.with_state(registry.clone());
self.add_router("/api/webrenderer", dynamic_router).await;
@@ -58,6 +59,7 @@ impl WebRendererExt for pmoserver::Server {
tracing::info!("WebRenderer server-side streaming endpoints registered");
tracing::info!(" POST /api/webrenderer/register");
tracing::info!(" GET /api/webrenderer/{{id}}/stream");
tracing::info!(" POST /api/webrenderer/{{id}}/position");
tracing::info!(" DELETE /api/webrenderer/{{id}}");
Ok(())
}

View File

@@ -7,7 +7,7 @@
use std::sync::Arc;
use pmoaudio::{AudioSegment, PositionTrackerNode, ResamplingNode, ToI24Node};
use pmoaudio::{AudioSegment, ResamplingNode, ToI24Node};
use pmometadata::{MemoryTrackMetadata, TrackMetadata};
use pmoaudio_ext::sinks::{
DirectOggFlacHandle, DirectOggFlacSink,
@@ -84,13 +84,9 @@ impl InstancePipeline {
let (sink, flac_handle) = DirectOggFlacSink::new(EncoderOptions::default());
// Nœud de suivi de position : lit le timestamp des chunks sortant vers le sink
let (mut position_tracker, position_handle) = PositionTrackerNode::new();
position_tracker.register(sink.boxed());
// Nœud de conversion de profondeur : tout type entier → I24
let mut to_i24 = ToI24Node::new();
to_i24.register(position_tracker.boxed());
to_i24.register(sink.boxed());
// Nœud de rééchantillonnage : n'importe quel sample rate → 96 kHz
let mut resampler = ResamplingNode::new(DIRECT_OGG_FLAC_SAMPLE_RATE);
@@ -113,25 +109,6 @@ impl InstancePipeline {
debug!("Sink task terminated");
});
// Task de mise à jour de la position (toutes les secondes)
{
let state_pos = state.clone();
let pos_stop = stop_token.clone();
tokio::spawn(async move {
loop {
tokio::select! {
_ = pos_stop.cancelled() => break,
_ = tokio::time::sleep(std::time::Duration::from_secs(1)) => {
let pos = position_handle.current_position_sec();
if pos > 0.0 {
state_pos.write().position = Some(seconds_to_upnp_time(pos));
}
}
}
}
});
}
// Task pipeline : reçoit les commandes UPnP et pilote les sources
let stop_token_clone = stop_token.clone();
let control_tx_clone = pipeline_handle.control_tx.clone();

View File

@@ -61,6 +61,22 @@ pub async fn register_handler(
}
}
#[derive(Debug, Deserialize)]
pub struct PositionUpdateRequest {
pub position_sec: f64,
pub duration_sec: Option<f64>,
}
/// POST /api/webrenderer/{id}/position
pub async fn position_update_handler(
State(registry): State<Arc<RendererRegistry>>,
Path(instance_id): Path<String>,
Json(req): Json<PositionUpdateRequest>,
) -> impl IntoResponse {
registry.update_position(&instance_id, req.position_sec, req.duration_sec);
StatusCode::NO_CONTENT
}
/// DELETE /api/webrenderer/{id}
pub async fn unregister_handler(
State(registry): State<Arc<RendererRegistry>>,

View File

@@ -151,6 +151,20 @@ impl RendererRegistry {
.map(|i| i.device_instance.clone())
}
/// Met à jour la position et la durée depuis le navigateur (audio.currentTime).
pub fn update_position(&self, instance_id: &str, position_sec: f64, duration_sec: Option<f64>) {
let instances = self.instances.read();
if let Some(instance) = instances.get(instance_id) {
let mut s = instance.state.write();
s.position = Some(crate::pipeline::seconds_to_upnp_time(position_sec));
if let Some(dur) = duration_sec {
if dur > 0.0 {
s.duration = Some(crate::pipeline::seconds_to_upnp_time(dur));
}
}
}
}
pub fn schedule_unregister(self: &Arc<Self>, instance_id: &str) {
use tokio_util::sync::CancellationToken;