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 pub pending: HashMap<Uuid, oneshot::Sender<RendererMessage>>,
30 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 pub fn new() -> Self {
53 Self::default()
54 }
55
56 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 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 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 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 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 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 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}