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 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 pub pending: HashMap<Uuid, oneshot::Sender<RendererMessage>>,
35 pub messages_sent: AtomicU64,
37 pub messages_received: AtomicU64,
38 pub(crate) timed_out: bool,
41 pub(crate) connection_id: Uuid,
45}
46
47struct 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 pub fn new() -> Self {
69 Self::default()
70 }
71
72 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 pub async fn register(&self, session: RendererSession) -> Result<RendererId> {
85 let id = session.id.clone();
86
87 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 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 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 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 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 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 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 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 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 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 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 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}