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, 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    pub sender: mpsc::Sender<ServerMessage>,
26    pub instances: HashMap<InstanceId, GraphicInstance>,
27    /// Requests awaiting a correlated reply from this renderer, keyed by the
28    /// `requestId` sent out on the ServerMessage.
29    pub pending: HashMap<Uuid, oneshot::Sender<RendererMessage>>,
30    /// Metrics tracking (for observability)
31    pub messages_sent: AtomicU64,
32    pub messages_received: AtomicU64,
33}
34
35pub struct RendererRegistry {
36    sessions: DashMap<RendererId, RendererSession>,
37    max_pending: usize,
38}
39
40impl Default for RendererRegistry {
41    fn default() -> Self {
42        Self {
43            sessions: DashMap::new(),
44            max_pending: 100,
45        }
46    }
47}
48
49impl RendererRegistry {
50    /// Creates a new RendererRegistry with default settings (max_pending=100).
51    /// Use `with_max_pending()` to customize the limit.
52    pub fn new() -> Self {
53        Self::default()
54    }
55
56    /// Creates a new RendererRegistry with a custom max pending requests limit.
57    pub fn with_max_pending(max_pending: usize) -> Self {
58        Self {
59            sessions: DashMap::new(),
60            max_pending,
61        }
62    }
63
64    pub async fn register(&self, session: RendererSession) -> RendererId {
65        let id = session.id;
66        self.sessions.insert(id, session);
67        id
68    }
69
70    pub async fn unregister(&self, id: RendererId) {
71        self.sessions.remove(&id);
72    }
73
74    pub async fn get_info(&self, id: RendererId) -> Result<RendererInfo> {
75        self.sessions
76            .get(&id)
77            .map(|entry| session_to_info(&entry))
78            .ok_or_else(|| AppError::NotFound(format!("renderer '{id}'")))
79    }
80
81    pub async fn list_info(&self) -> Vec<RendererInfo> {
82        self.sessions
83            .iter()
84            .map(|entry| session_to_info(entry.value()))
85            .collect()
86    }
87
88    /// Sends a command to a renderer and waits for its correlated result.
89    /// The OGraf spec requires load()/playAction()/etc HTTP responses to
90    /// reflect what actually happened inside the GraphicInstance, so this
91    /// is the only way commands reach a renderer — no fire-and-forget.
92    pub async fn send_and_await(
93        &self,
94        renderer_id: RendererId,
95        build: impl FnOnce(Uuid) -> ServerMessage,
96        timeout: Duration,
97    ) -> Result<RendererMessage> {
98        let request_id = Uuid::new_v4();
99        let msg = build(request_id);
100        let (tx, rx) = oneshot::channel();
101
102        // DashMap's get_mut holds a shard-level lock, not a global lock,
103        // so one busy renderer doesn't stall commands to other renderers
104        // (assuming they hash to different shards).
105        let sender = {
106            let mut session = self
107                .sessions
108                .get_mut(&renderer_id)
109                .ok_or_else(|| AppError::RendererNotConnected(renderer_id.to_string()))?;
110
111            // Check if renderer has too many pending requests (DoS protection)
112            if session.pending.len() >= self.max_pending {
113                return Err(AppError::RendererOverloaded(format!(
114                    "Renderer has {} pending requests (max: {})",
115                    session.pending.len(),
116                    self.max_pending
117                )));
118            }
119
120            session.pending.insert(request_id, tx);
121            session.messages_sent.fetch_add(1, Ordering::Relaxed);
122            session.sender.clone()
123        };
124
125        sender
126            .send(msg)
127            .await
128            .map_err(|_| AppError::RendererNotConnected(renderer_id.to_string()))?;
129
130        match tokio::time::timeout(timeout, rx).await {
131            Ok(Ok(reply)) => Ok(reply),
132            Ok(Err(_)) => Err(AppError::RendererNotConnected(renderer_id.to_string())),
133            Err(_) => {
134                if let Some(mut session) = self.sessions.get_mut(&renderer_id) {
135                    session.pending.remove(&request_id);
136                }
137                Err(AppError::Timeout(renderer_id.to_string()))
138            }
139        }
140    }
141
142    /// Applies a renderer-confirmed result to session state (so GET renderer
143    /// endpoints reflect renderer-confirmed truth, not what was merely sent)
144    /// and wakes up the HTTP handler awaiting it via `send_and_await`, if any.
145    pub async fn resolve(
146        &self,
147        renderer_id: RendererId,
148        request_id: Uuid,
149        message: RendererMessage,
150    ) {
151        if let Some(mut session) = self.sessions.get_mut(&renderer_id) {
152            session.messages_received.fetch_add(1, Ordering::Relaxed);
153            apply_result(&mut session, &message);
154
155            if let Some(tx) = session.pending.remove(&request_id) {
156                let _ = tx.send(message);
157            }
158        }
159    }
160}
161
162fn apply_result(session: &mut RendererSession, message: &RendererMessage) {
163    // Only update state on success — Composed Method pattern makes this
164    // rule explicit rather than repeating `if status_code < 400` in each arm.
165    if !message.is_success() {
166        return;
167    }
168
169    match message {
170        RendererMessage::LoadResult {
171            instance_id,
172            graphic_id,
173            data,
174            ..
175        } => {
176            session.instances.insert(
177                *instance_id,
178                GraphicInstance {
179                    instance_id: *instance_id,
180                    graphic_id: graphic_id.clone(),
181                    data: data.clone(),
182                    loaded_at: Utc::now(),
183                    state: InstanceState::Loaded,
184                    current_step: None,
185                },
186            );
187        }
188        RendererMessage::PlayActionResult {
189            instance_id,
190            current_step,
191            ..
192        } => {
193            if let Some(instance) = session.instances.get_mut(instance_id) {
194                instance.state = InstanceState::Playing;
195                instance.current_step = Some(*current_step);
196            }
197        }
198        RendererMessage::StopActionResult { instance_id, .. } => {
199            if let Some(instance) = session.instances.get_mut(instance_id) {
200                instance.state = InstanceState::Stopped;
201            }
202        }
203        RendererMessage::UpdateActionResult {
204            instance_id,
205            data: Some(data),
206            ..
207        } => {
208            if let Some(instance) = session.instances.get_mut(instance_id) {
209                instance.data = Some(data.clone());
210            }
211        }
212        RendererMessage::ClearResult { instance_id, .. } => {
213            session.instances.remove(instance_id);
214        }
215        _ => {}
216    }
217}
218
219fn session_to_info(s: &RendererSession) -> RendererInfo {
220    // `instances` is a HashMap (keyed by id for O(1) lookup on actions),
221    // which has no defined iteration order — sorting by `loaded_at` here
222    // reconstructs actual load order, which is also visual stacking order
223    // in typical renderers (newer instances paint on top). Newest first,
224    // so callers can treat this list as "top of stack first".
225    let mut instances: Vec<_> = s.instances.values().cloned().collect();
226    instances.sort_by(|a, b| b.loaded_at.cmp(&a.loaded_at));
227
228    let uptime_seconds = s
229        .connected_at
230        .signed_duration_since(Utc::now())
231        .num_seconds()
232        .unsigned_abs();
233
234    let metrics = RendererMetrics {
235        pending_requests: s.pending.len(),
236        messages_sent: s.messages_sent.load(Ordering::Relaxed),
237        messages_received: s.messages_received.load(Ordering::Relaxed),
238        uptime_seconds,
239    };
240
241    RendererInfo {
242        id: s.id,
243        name: s.name.clone(),
244        connected_at: s.connected_at,
245        render_target: s.render_target.clone(),
246        render_target_schema: s.render_target_schema.clone(),
247        instances,
248        metrics: Some(metrics),
249    }
250}