Skip to main content

ograf_core/store/
renderers.rs

1use chrono::Utc;
2use dashmap::DashMap;
3use std::{
4    collections::HashMap,
5    sync::atomic::{AtomicU64, Ordering},
6    time::Duration,
7};
8use tokio::sync::{mpsc, oneshot};
9use uuid::Uuid;
10
11use crate::{
12    error::{AppError, Result},
13    models::{
14        GraphicInstance, InstanceId, InstanceState, RenderTarget, RendererId, RendererInfo,
15        RendererMessage, RendererMetrics, RendererStatus, ServerMessage,
16    },
17};
18
19pub struct RendererSession {
20    pub id: RendererId,
21    pub name: String,
22    pub connected_at: chrono::DateTime<Utc>,
23    pub render_target: RenderTarget,
24    pub render_target_schema: Option<serde_json::Value>,
25    /// `hello.capabilities.description` / `.customActions` /
26    /// `.renderCharacteristics` — passed through to the spec's RendererInfo.
27    pub description: Option<String>,
28    pub custom_actions: Option<serde_json::Value>,
29    pub render_characteristics: Option<serde_json::Value>,
30    pub sender: mpsc::Sender<ServerMessage>,
31    pub instances: HashMap<InstanceId, GraphicInstance>,
32    /// Requests awaiting a correlated reply from this renderer, keyed by the
33    /// `requestId` sent out on the ServerMessage.
34    pub pending: HashMap<Uuid, oneshot::Sender<RendererMessage>>,
35    /// Metrics tracking (for observability)
36    pub messages_sent: AtomicU64,
37    pub messages_received: AtomicU64,
38    /// Set when a request to this renderer times out, cleared by its next
39    /// reply — reported as `status: WARNING`.
40    pub(crate) timed_out: bool,
41    /// Internal connection identifier for cleanup safety. Not exposed in API.
42    /// Ensures a refused or late-cleaning-up connection doesn't unregister
43    /// a different session that took over the name.
44    pub(crate) connection_id: Uuid,
45}
46
47/// What's left of a session once it ends — enough to keep listing the
48/// renderer (with `status: ERROR`) until it reconnects.
49struct DepartedRenderer {
50    info: RendererInfo,
51}
52
53pub struct RendererRegistry {
54    sessions: DashMap<RendererId, RendererSession>,
55    departed: DashMap<RendererId, DepartedRenderer>,
56    max_pending: usize,
57}
58
59impl Default for RendererRegistry {
60    fn default() -> Self {
61        Self::with_max_pending(100)
62    }
63}
64
65impl RendererRegistry {
66    /// Creates a new RendererRegistry with default settings (max_pending=100).
67    /// Use `with_max_pending()` to customize the limit.
68    pub fn new() -> Self {
69        Self::default()
70    }
71
72    /// Creates a new RendererRegistry with a custom max pending requests limit.
73    pub fn with_max_pending(max_pending: usize) -> Self {
74        Self {
75            sessions: DashMap::new(),
76            departed: DashMap::new(),
77            max_pending,
78        }
79    }
80
81    /// Registers a renderer session. Returns Ok(id) on success, or Err if
82    /// the name is already in use (first-wins policy). Atomic check-and-insert
83    /// ensures no race conditions with concurrent registration attempts.
84    pub async fn register(&self, session: RendererSession) -> Result<RendererId> {
85        let id = session.id.clone();
86
87        // Atomic check-and-insert: insert only if key doesn't exist (first-wins)
88        use dashmap::mapref::entry::Entry;
89        match self.sessions.entry(id.clone()) {
90            Entry::Vacant(entry) => {
91                entry.insert(session);
92                self.departed.remove(&id);
93                Ok(id)
94            }
95            Entry::Occupied(_) => Err(AppError::Conflict(format!(
96                "renderer name '{}' already connected",
97                id
98            ))),
99        }
100    }
101
102    /// Unregisters a renderer session, but only if the connection_id matches.
103    /// This prevents a refused or late-cleaning-up connection from unregistering
104    /// a different session that took over the name.
105    /// The renderer stays listed afterwards, as disconnected.
106    pub async fn unregister(&self, id: &str, connection_id: Uuid) {
107        let removed = self
108            .sessions
109            .remove_if(id, |_, session| session.connection_id == connection_id);
110        if let Some((id, session)) = removed {
111            let info = departed_info(&session);
112            self.departed.insert(id, DepartedRenderer { info });
113        }
114    }
115
116    /// Connected, or disconnected since this process started.
117    pub async fn get_info(&self, id: &str) -> Result<RendererInfo> {
118        if let Some(entry) = self.sessions.get(id) {
119            return Ok(self.session_to_info(&entry));
120        }
121        self.departed
122            .get(id)
123            .map(|entry| entry.info.clone())
124            .ok_or_else(|| AppError::NotFound(format!("renderer '{id}'")))
125    }
126
127    /// Connected renderers plus ones that disconnected since this process
128    /// started — check [`RendererInfo::is_connected`] to tell them apart.
129    pub async fn list_info(&self) -> Vec<RendererInfo> {
130        let connected = self
131            .sessions
132            .iter()
133            .map(|entry| self.session_to_info(entry.value()));
134        let departed = self
135            .departed
136            .iter()
137            .filter(|entry| !self.sessions.contains_key(entry.key()))
138            .map(|entry| entry.info.clone());
139        connected.chain(departed).collect()
140    }
141
142    /// Whether any connected renderer has an instance of `graphic_id` —
143    /// a deleted graphic's files stay until this is false.
144    pub fn graphic_in_use(&self, graphic_id: &str) -> bool {
145        self.sessions.iter().any(|session| {
146            session
147                .instances
148                .values()
149                .any(|instance| instance.graphic_id == graphic_id)
150        })
151    }
152
153    /// Sends a command to a renderer and waits for its correlated result.
154    /// The OGraf spec requires load()/playAction()/etc HTTP responses to
155    /// reflect what actually happened inside the GraphicInstance, so this
156    /// is the only way commands reach a renderer — no fire-and-forget.
157    pub async fn send_and_await(
158        &self,
159        renderer_id: &str,
160        build: impl FnOnce(Uuid) -> ServerMessage,
161        timeout: Duration,
162    ) -> Result<RendererMessage> {
163        let request_id = Uuid::new_v4();
164        let msg = build(request_id);
165        let (tx, rx) = oneshot::channel();
166
167        // DashMap's get_mut holds a shard-level lock, not a global lock,
168        // so one busy renderer doesn't stall commands to other renderers
169        // (assuming they hash to different shards).
170        let sender = {
171            let mut session = self
172                .sessions
173                .get_mut(renderer_id)
174                .ok_or_else(|| AppError::RendererNotConnected(renderer_id.to_string()))?;
175
176            // Check if renderer has too many pending requests (DoS protection)
177            if session.pending.len() >= self.max_pending {
178                return Err(AppError::RendererOverloaded(format!(
179                    "Renderer has {} pending requests (max: {})",
180                    session.pending.len(),
181                    self.max_pending
182                )));
183            }
184
185            session.pending.insert(request_id, tx);
186            session.messages_sent.fetch_add(1, Ordering::Relaxed);
187            session.sender.clone()
188        };
189
190        sender
191            .send(msg)
192            .await
193            .map_err(|_| AppError::RendererNotConnected(renderer_id.to_string()))?;
194
195        match tokio::time::timeout(timeout, rx).await {
196            Ok(Ok(reply)) => Ok(reply),
197            Ok(Err(_)) => Err(AppError::RendererNotConnected(renderer_id.to_string())),
198            Err(_) => {
199                if let Some(mut session) = self.sessions.get_mut(renderer_id) {
200                    session.pending.remove(&request_id);
201                    session.timed_out = true;
202                }
203                Err(AppError::Timeout(renderer_id.to_string()))
204            }
205        }
206    }
207
208    /// Applies a renderer-confirmed result to session state (so GET renderer
209    /// endpoints reflect renderer-confirmed truth, not what was merely sent)
210    /// and wakes up the HTTP handler awaiting it via `send_and_await`, if any.
211    pub async fn resolve(
212        &self,
213        renderer_id: &str,
214        request_id: Uuid,
215        message: RendererMessage,
216    ) {
217        if let Some(mut session) = self.sessions.get_mut(renderer_id) {
218            session.messages_received.fetch_add(1, Ordering::Relaxed);
219            session.timed_out = false;
220            apply_result(&mut session, &message);
221
222            if let Some(tx) = session.pending.remove(&request_id) {
223                let _ = tx.send(message);
224            }
225        }
226    }
227}
228
229fn apply_result(session: &mut RendererSession, message: &RendererMessage) {
230    // Only update state on success — Composed Method pattern makes this
231    // rule explicit rather than repeating `if status_code < 400` in each arm.
232    if !message.is_success() {
233        return;
234    }
235
236    match message {
237        RendererMessage::LoadResult {
238            instance_id,
239            graphic_id,
240            data,
241            ..
242        } => {
243            session.instances.insert(
244                *instance_id,
245                GraphicInstance {
246                    instance_id: *instance_id,
247                    graphic_id: graphic_id.clone(),
248                    data: data.clone(),
249                    loaded_at: Utc::now(),
250                    state: InstanceState::Loaded,
251                    current_step: None,
252                },
253            );
254        }
255        RendererMessage::PlayActionResult {
256            instance_id,
257            current_step,
258            ..
259        } => {
260            if let Some(instance) = session.instances.get_mut(instance_id) {
261                instance.state = InstanceState::Playing;
262                instance.current_step = Some(*current_step);
263            }
264        }
265        RendererMessage::StopActionResult { instance_id, .. } => {
266            if let Some(instance) = session.instances.get_mut(instance_id) {
267                instance.state = InstanceState::Stopped;
268            }
269        }
270        RendererMessage::UpdateActionResult {
271            instance_id,
272            data: Some(data),
273            ..
274        } => {
275            if let Some(instance) = session.instances.get_mut(instance_id) {
276                instance.data = Some(data.clone());
277            }
278        }
279        RendererMessage::ClearResult { instance_id, .. } => {
280            session.instances.remove(instance_id);
281        }
282        _ => {}
283    }
284}
285
286impl RendererRegistry {
287    fn session_to_info(&self, s: &RendererSession) -> RendererInfo {
288        session_to_info(s, self.status_of(s))
289    }
290
291    /// `WARNING` once requests pile up (half of `max_pending`) or the last
292    /// one timed out — the early sign of a renderer that stopped answering
293    /// while its WebSocket stays open.
294    fn status_of(&self, s: &RendererSession) -> RendererStatus {
295        let pending = s.pending.len();
296        if pending > 0 && pending >= (self.max_pending / 2).max(1) {
297            RendererStatus::warning(format!("{pending} pending requests"))
298        } else if s.timed_out {
299            RendererStatus::warning("last request timed out")
300        } else {
301            RendererStatus::ok()
302        }
303    }
304}
305
306fn departed_info(s: &RendererSession) -> RendererInfo {
307    let now = Utc::now();
308    let mut info = session_to_info(s, RendererStatus::error(format!("disconnected since {}", now.to_rfc3339())));
309    info.connected_at = None;
310    info.disconnected_at = Some(now);
311    info.instances = Vec::new();
312    info.metrics = None;
313    info
314}
315
316fn session_to_info(s: &RendererSession, status: RendererStatus) -> RendererInfo {
317    // `instances` is a HashMap (keyed by id for O(1) lookup on actions),
318    // which has no defined iteration order — sorting by `loaded_at` here
319    // reconstructs actual load order, which is also visual stacking order
320    // in typical renderers (newer instances paint on top). Newest first,
321    // so callers can treat this list as "top of stack first".
322    let mut instances: Vec<_> = s.instances.values().cloned().collect();
323    instances.sort_by(|a, b| b.loaded_at.cmp(&a.loaded_at));
324
325    let uptime_seconds = s
326        .connected_at
327        .signed_duration_since(Utc::now())
328        .num_seconds()
329        .unsigned_abs();
330
331    let metrics = RendererMetrics {
332        pending_requests: s.pending.len(),
333        messages_sent: s.messages_sent.load(Ordering::Relaxed),
334        messages_received: s.messages_received.load(Ordering::Relaxed),
335        uptime_seconds,
336    };
337
338    RendererInfo {
339        id: s.id.clone(),
340        name: s.name.clone(),
341        description: s.description.clone(),
342        status,
343        connected_at: Some(s.connected_at),
344        disconnected_at: None,
345        render_target: Some(s.render_target.clone()),
346        render_target_schema: s.render_target_schema.clone(),
347        custom_actions: s.custom_actions.clone(),
348        render_characteristics: s.render_characteristics.clone(),
349        instances,
350        metrics: Some(metrics),
351    }
352}