1use std::collections::HashMap;
2use std::sync::Arc;
3use std::sync::atomic::Ordering::Relaxed;
4
5use anyhow::Result;
6pub use engine_transport::EngineTransport;
7use indexmap::IndexMap;
8use kcmc::ModelingCmd;
9use kcmc::each_cmd as mcmd;
10use kcmc::shared::Color;
11use kcmc::websocket::BatchResponse;
12use kcmc::websocket::ModelingCmdReq;
13use kcmc::websocket::ModelingSessionData;
14use kcmc::websocket::OkWebSocketResponseData;
15use kcmc::websocket::WebSocketRequest;
16use kcmc::websocket::WebSocketResponse;
17use kittycad_modeling_cmds::ModelingCmdEndpoint;
18use kittycad_modeling_cmds::length_unit::LengthUnit;
19use kittycad_modeling_cmds::ok_response::OkModelingCmdResponse;
20use kittycad_modeling_cmds::websocket::ModelingBatch;
21use kittycad_modeling_cmds::{self as kcmc};
22use tokio::sync::RwLock;
23use uuid::Uuid;
24use web_time::Instant;
25
26use crate::ExecutorSettings;
27use crate::SourceRange;
28use crate::engine::AsyncTasks;
29use crate::engine::DEFAULT_PLANE_INFO;
30use crate::engine::EngineBatchContext;
31use crate::engine::EngineStats;
32use crate::engine::GRID_OBJECT_ID;
33use crate::engine::GRID_SCALE_TEXT_OBJECT_ID;
34use crate::engine::GridScaleBehavior;
35use crate::engine::PlaneName;
36use crate::errors::KclError;
37use crate::errors::KclErrorDetails;
38use crate::execution::DefaultPlanes;
39use crate::execution::IdGenerator;
40use crate::execution::PlaneInfo;
41use crate::settings::types::default_backface_color;
42use crate::settings::types::default_backface_color_struct;
43
44pub enum TransportCloseError {}
45
46mod engine_transport;
47mod mock_transport;
48#[cfg(target_arch = "wasm32")]
49pub mod wasm_transport;
50#[cfg(not(target_arch = "wasm32"))]
51pub mod ws_transport;
52
53#[derive(Clone, Debug)]
55pub struct ResponseInformation {
56 responses: Arc<RwLock<IndexMap<uuid::Uuid, WebSocketResponse>>>,
58}
59
60impl ResponseInformation {
61 pub fn new(responses: Arc<RwLock<IndexMap<uuid::Uuid, WebSocketResponse>>>) -> Self {
63 Self { responses }
64 }
65
66 pub async fn add(&self, id: Uuid, response: WebSocketResponse) {
68 self.responses.write().await.insert(id, response);
69 }
70}
71
72#[derive(bon::Builder)]
73pub struct EngineManager {
74 pub transport: Arc<Box<dyn EngineTransport>>,
77 responses: ResponseInformation,
78 pending_errors: Arc<RwLock<Vec<String>>>,
79 socket_health: Arc<RwLock<SocketHealth>>,
80 ids_of_async_commands: Arc<RwLock<IndexMap<Uuid, SourceRange>>>,
81
82 #[builder(default)]
84 default_planes: Arc<RwLock<Option<DefaultPlanes>>>,
85 session_data: Arc<RwLock<Option<ModelingSessionData>>>,
87
88 websocket_upgrade_request_id: Option<String>,
90
91 #[builder(default)]
92 stats: EngineStats,
93
94 #[builder(default)]
95 async_tasks: AsyncTasks,
96}
97
98impl std::fmt::Debug for EngineManager {
99 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
100 f.debug_struct("EngineManager")
101 .field("responses", &self.responses)
102 .field("pending_errors", &self.pending_errors)
103 .field("socket_health", &self.socket_health)
104 .field("ids_of_async_commands", &self.ids_of_async_commands)
105 .field("default_planes", &self.default_planes)
106 .field("session_data", &self.session_data)
107 .field("websocket_upgrade_request_id", &self.websocket_upgrade_request_id)
108 .field("stats", &self.stats)
109 .field("async_tasks", &self.async_tasks)
110 .finish()
111 }
112}
113
114impl EngineManager {
115 #[cfg(target_arch = "wasm32")]
116 pub fn new_wasm_transport(
117 manager: wasm_transport::EngineCommandManager,
118 response_context: Arc<wasm_transport::ResponseContext>,
119 ) -> Self {
120 let session_data: Arc<RwLock<Option<ModelingSessionData>>> = Arc::new(RwLock::new(None));
121 let ids_of_async_commands: Arc<RwLock<IndexMap<Uuid, SourceRange>>> = Arc::new(RwLock::new(IndexMap::new()));
122 let socket_health = Arc::new(RwLock::new(SocketHealth::Active));
123 let pending_errors = Arc::new(RwLock::new(Vec::new()));
124 let responses = response_context.response_information();
125
126 Self {
127 transport: Arc::new(Box::new(wasm_transport::WasmTransport::new(manager))),
128 responses,
129 pending_errors,
130 socket_health,
131 ids_of_async_commands,
132 default_planes: Default::default(),
133 session_data,
134 websocket_upgrade_request_id: None,
135 stats: Default::default(),
136 async_tasks: Default::default(),
137 }
138 }
139
140 #[cfg(not(target_arch = "wasm32"))]
141 pub async fn new_websocket_transport(ws: reqwest::Upgraded, heartbeats: Option<u64>) -> Self {
142 Self::new_websocket_transport_with_request_id(ws, heartbeats, None).await
143 }
144
145 #[cfg(not(target_arch = "wasm32"))]
146 pub(crate) async fn new_websocket_transport_with_request_id(
147 ws: reqwest::Upgraded,
148 heartbeats: Option<u64>,
149 request_id: Option<String>,
150 ) -> Self {
151 use crate::engine::engine_manager::ws_transport::WebSocketTransport;
152
153 let session_data: Arc<RwLock<Option<ModelingSessionData>>> = Arc::new(RwLock::new(None));
154 let ids_of_async_commands: Arc<RwLock<IndexMap<Uuid, SourceRange>>> = Arc::new(RwLock::new(IndexMap::new()));
155 let socket_health = Arc::new(RwLock::new(SocketHealth::Active));
156 let pending_errors = Arc::new(RwLock::new(Vec::new()));
157 let responses = ResponseInformation {
158 responses: Arc::new(RwLock::new(IndexMap::new())),
159 };
160
161 let transport = WebSocketTransport::spawn(
162 ws,
163 heartbeats,
164 responses.clone(),
165 Arc::clone(&session_data),
166 Arc::clone(&pending_errors),
167 Arc::clone(&socket_health),
168 request_id.clone(),
169 )
170 .await;
171
172 Self {
173 transport: Arc::new(Box::new(transport)),
174 responses,
175 pending_errors,
176 socket_health,
177 ids_of_async_commands,
178 default_planes: Default::default(),
179 session_data,
180 websocket_upgrade_request_id: request_id,
181 stats: Default::default(),
182 async_tasks: Default::default(),
183 }
184 }
185
186 pub fn new_mock() -> Self {
189 let session_data: Arc<RwLock<Option<ModelingSessionData>>> = Arc::new(RwLock::new(None));
190 let ids_of_async_commands: Arc<RwLock<IndexMap<Uuid, SourceRange>>> = Arc::new(RwLock::new(IndexMap::new()));
191 let socket_health = Arc::new(RwLock::new(SocketHealth::Active));
192 let pending_errors = Arc::new(RwLock::new(Vec::new()));
193 let responses = ResponseInformation {
194 responses: Arc::new(RwLock::new(IndexMap::new())),
195 };
196 Self {
197 transport: Arc::new(Box::new(mock_transport::MockTransport::new(responses.clone()))),
198 responses,
199 pending_errors,
200 socket_health,
201 ids_of_async_commands,
202 default_planes: Default::default(),
203 session_data,
204 websocket_upgrade_request_id: None,
205 stats: Default::default(),
206 async_tasks: Default::default(),
207 }
208 }
209
210 async fn take_ids_of_async_commands(&self) -> IndexMap<Uuid, SourceRange> {
212 std::mem::take(&mut *self.ids_of_async_commands().write().await)
213 }
214
215 pub async fn take_responses(&self) -> IndexMap<Uuid, WebSocketResponse> {
217 std::mem::take(&mut *self.responses().write().await)
218 }
219
220 pub async fn clear_scene(
221 &self,
222 batch_context: &EngineBatchContext,
223 id_generator: &mut IdGenerator,
224 source_range: SourceRange,
225 geometry_only: bool,
226 ) -> Result<(), crate::errors::KclError> {
227 self.clear_queues(batch_context).await;
229
230 self.batch_modeling_cmd(
231 batch_context,
232 id_generator.next_uuid(),
233 source_range,
234 &ModelingCmd::SceneClearAll(mcmd::SceneClearAll::default()),
235 )
236 .await?;
237
238 self.flush_batch(batch_context, false, source_range).await?;
241
242 self.clear_scene_post_hook(batch_context, id_generator, source_range, geometry_only)
244 .await?;
245
246 Ok(())
247 }
248
249 pub async fn ensure_async_command_completed(
251 &self,
252 id: uuid::Uuid,
253 source_range: Option<SourceRange>,
254 ) -> Result<OkWebSocketResponseData, KclError> {
255 let source_range = if let Some(source_range) = source_range {
256 source_range
257 } else {
258 self.ids_of_async_commands()
260 .read()
261 .await
262 .get(&id)
263 .cloned()
264 .unwrap_or_default()
265 };
266
267 const ASYNC_CMD_TIMEOUT_SECS: u64 = 600;
275 let current_time = Instant::now();
276 while current_time.elapsed().as_secs() < ASYNC_CMD_TIMEOUT_SECS {
277 let responses = self.responses().read().await.clone();
278 let Some(resp) = responses.get(&id) else {
279 #[cfg(target_arch = "wasm32")]
282 {
283 let duration = web_time::Duration::from_millis(1);
284 wasm_timer::Delay::new(duration).await.map_err(|err| {
285 KclError::new_internal(KclErrorDetails::new(
286 format!("Failed to sleep: {:?}", err),
287 vec![source_range],
288 ))
289 })?;
290 }
291 #[cfg(not(target_arch = "wasm32"))]
292 tokio::task::yield_now().await;
293 continue;
294 };
295
296 let response = self.parse_websocket_response(resp.clone(), source_range)?;
299 return Ok(response);
300 }
301
302 Err(KclError::new_engine(KclErrorDetails::new(
303 format!(
304 "async command timed out after {ASYNC_CMD_TIMEOUT_SECS}s (client-side ceiling, not an engine error)"
305 ),
306 vec![source_range],
307 )))
308 }
309
310 pub async fn ensure_async_commands_completed(&self, batch_context: &EngineBatchContext) -> Result<(), KclError> {
312 let ids = self.take_ids_of_async_commands().await;
314
315 for (id, source_range) in ids {
317 self.ensure_async_command_completed(id, Some(source_range)).await?;
318 }
319
320 if let Err(err) = self.async_tasks().join_all().await {
326 crate::log::logln!(
327 "Error waiting for async tasks (this is typically fine and just means that an edge became something else): {:?}",
328 err
329 );
330 }
331
332 self.flush_batch(batch_context, true, SourceRange::default()).await?;
334
335 Ok(())
336 }
337
338 async fn set_edge_visibility(
340 &self,
341 batch_context: &EngineBatchContext,
342 visible: bool,
343 source_range: SourceRange,
344 id_generator: &mut IdGenerator,
345 ) -> Result<(), crate::errors::KclError> {
346 self.batch_modeling_cmd(
347 batch_context,
348 id_generator.next_uuid(),
349 source_range,
350 &ModelingCmd::from(mcmd::EdgeLinesVisible::builder().hidden(!visible).build()),
351 )
352 .await?;
353
354 Ok(())
355 }
356
357 pub async fn reapply_settings(
359 &self,
360 batch_context: &EngineBatchContext,
361 settings: &crate::ExecutorSettings,
362 source_range: SourceRange,
363 id_generator: &mut IdGenerator,
364 grid_scale_unit: GridScaleBehavior,
365 ) -> Result<(), crate::errors::KclError> {
366 if settings.geometry_only {
367 return Ok(());
368 }
369 self.set_edge_visibility(batch_context, settings.highlight_edges, source_range, id_generator)
371 .await?;
372
373 self.modify_grid(
376 batch_context,
377 !settings.show_grid,
378 grid_scale_unit,
379 source_range,
380 id_generator,
381 )
382 .await?;
383
384 self.set_user_colors(batch_context, settings, source_range, id_generator)
386 .await?;
387
388 self.flush_batch(batch_context, false, source_range).await?;
392
393 Ok(())
394 }
395
396 pub async fn batch_modeling_cmd(
398 &self,
399 batch_context: &EngineBatchContext,
400 id: uuid::Uuid,
401 source_range: SourceRange,
402 cmd: &ModelingCmd,
403 ) -> Result<(), crate::errors::KclError> {
404 let req = WebSocketRequest::ModelingCmdReq(ModelingCmdReq {
405 cmd: cmd.clone(),
406 cmd_id: id.into(),
407 });
408
409 batch_context.push(req, source_range).await;
411 self.stats().commands_batched.fetch_add(1, Relaxed);
412
413 Ok(())
414 }
415
416 pub async fn batch_modeling_cmds(
421 &self,
422 batch_context: &EngineBatchContext,
423 source_range: SourceRange,
424 cmds: &[ModelingCmdReq],
425 ) -> Result<(), crate::errors::KclError> {
426 let mut extended_cmds = Vec::with_capacity(cmds.len());
428 for cmd in cmds {
429 extended_cmds.push((WebSocketRequest::ModelingCmdReq(cmd.clone()), source_range));
430 }
431 self.stats().commands_batched.fetch_add(extended_cmds.len(), Relaxed);
432 batch_context.extend(extended_cmds).await;
433
434 Ok(())
435 }
436
437 pub async fn batch_end_cmd(
441 &self,
442 batch_context: &EngineBatchContext,
443 id: uuid::Uuid,
444 source_range: SourceRange,
445 cmd: &ModelingCmd,
446 ) -> Result<(), crate::errors::KclError> {
447 let req = WebSocketRequest::ModelingCmdReq(ModelingCmdReq {
448 cmd: cmd.clone(),
449 cmd_id: id.into(),
450 });
451
452 batch_context.insert_end(id, req, source_range).await;
454 self.stats().commands_batched.fetch_add(1, Relaxed);
455 Ok(())
456 }
457
458 pub async fn send_modeling_cmd(
460 &self,
461 batch_context: &EngineBatchContext,
462 id: uuid::Uuid,
463 source_range: SourceRange,
464 cmd: &ModelingCmd,
465 ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
466 let mut requests = batch_context.take_batch().await;
467
468 requests.push((
470 WebSocketRequest::ModelingCmdReq(ModelingCmdReq {
471 cmd: cmd.clone(),
472 cmd_id: id.into(),
473 }),
474 source_range,
475 ));
476 self.stats().commands_batched.fetch_add(1, Relaxed);
477
478 self.run_batch(requests, source_range).await
480 }
481
482 pub async fn async_modeling_cmd(
485 &self,
486 id: uuid::Uuid,
487 source_range: SourceRange,
488 cmd: &ModelingCmd,
489 ) -> Result<(), crate::errors::KclError> {
490 self.ids_of_async_commands().write().await.insert(id, source_range);
492
493 self.transport
495 .inner_fire_modeling_cmd(
496 id,
497 source_range,
498 WebSocketRequest::ModelingCmdReq(ModelingCmdReq {
499 cmd: cmd.clone(),
500 cmd_id: id.into(),
501 }),
502 HashMap::from([(id, source_range)]),
503 )
504 .await?;
505
506 Ok(())
507 }
508
509 async fn run_batch(
511 &self,
512 orig_requests: Vec<(WebSocketRequest, SourceRange)>,
513 source_range: SourceRange,
514 ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
515 if orig_requests.is_empty() {
517 return Ok(OkWebSocketResponseData::Modeling {
518 modeling_response: OkModelingCmdResponse::Empty {},
519 });
520 }
521
522 let requests: Vec<ModelingCmdReq> = orig_requests
523 .iter()
524 .filter_map(|(val, _)| match val {
525 WebSocketRequest::ModelingCmdReq(ModelingCmdReq { cmd, cmd_id }) => Some(ModelingCmdReq {
526 cmd: cmd.clone(),
527 cmd_id: *cmd_id,
528 }),
529 _ => None,
530 })
531 .collect();
532
533 let batched_requests = WebSocketRequest::ModelingCmdBatchReq(ModelingBatch {
534 requests,
535 batch_id: uuid::Uuid::new_v4().into(),
536 responses: true,
537 });
538
539 let final_req = if orig_requests.len() == 1 {
540 orig_requests.first().unwrap().0.clone()
542 } else {
543 batched_requests
544 };
545
546 let mut id_to_source_range = HashMap::new();
549
550 let mut id_to_command = HashMap::new();
551 for (req, range) in orig_requests.iter() {
552 match req {
553 WebSocketRequest::ModelingCmdReq(ModelingCmdReq { cmd, cmd_id }) => {
554 let id = Uuid::from(*cmd_id);
555 id_to_source_range.insert(id, *range);
556 id_to_command.insert(id, ModelingCmdEndpoint::from(cmd));
557 }
558 _ => {
559 return Err(KclError::new_engine(KclErrorDetails::new(
560 format!("The request is not a modeling command: {req:?}"),
561 vec![*range],
562 )));
563 }
564 }
565 }
566
567 self.stats().batches_sent.fetch_add(1, Relaxed);
568
569 match final_req {
571 WebSocketRequest::ModelingCmdBatchReq(ModelingBatch {
572 ref requests,
573 batch_id,
574 responses: _,
575 }) => {
576 let last_id = requests.last().unwrap().cmd_id;
578 let ws_resp = self
579 .inner_send_modeling_cmd(batch_id.into(), source_range, final_req, id_to_source_range.clone())
580 .await?;
581 let response = self.parse_websocket_response(ws_resp, source_range)?;
582
583 if let OkWebSocketResponseData::ModelingBatch { responses } = response {
585 self.parse_batch_responses(last_id.into(), id_to_source_range, id_to_command, responses)
586 } else {
587 Err(KclError::new_engine(KclErrorDetails::new(
589 format!("Failed to get batch response: {response:?}"),
590 vec![source_range],
591 )))
592 }
593 }
594 WebSocketRequest::ModelingCmdReq(ModelingCmdReq { cmd: _, cmd_id }) => {
595 let source_range = id_to_source_range.get(cmd_id.as_ref()).cloned().ok_or_else(|| {
603 KclError::new_engine(KclErrorDetails::new(
604 format!("Failed to get source range for command ID: {cmd_id:?}"),
605 vec![],
606 ))
607 })?;
608 let ws_resp = self
609 .inner_send_modeling_cmd(cmd_id.into(), source_range, final_req, id_to_source_range)
610 .await?;
611 self.parse_websocket_response(ws_resp, source_range)
612 }
613 _ => Err(KclError::new_engine(KclErrorDetails::new(
614 format!("The final request is not a modeling command: {final_req:?}"),
615 vec![source_range],
616 ))),
617 }
618 }
619
620 pub async fn flush_batch(
622 &self,
623 batch_context: &EngineBatchContext,
624 batch_end: bool,
627 source_range: SourceRange,
628 ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
629 let all_requests = if batch_end {
630 let mut requests = batch_context.take_batch().await;
631 requests.extend(batch_context.take_batch_end().await.values().cloned());
632 requests
633 } else {
634 batch_context.take_batch().await
635 };
636
637 self.run_batch(all_requests, source_range).await
638 }
639
640 async fn make_default_plane(
641 &self,
642 batch_context: &EngineBatchContext,
643 plane_id: uuid::Uuid,
644 info: &PlaneInfo,
645 color: Option<Color>,
646 source_range: SourceRange,
647 id_generator: &mut IdGenerator,
648 ) -> Result<uuid::Uuid, KclError> {
649 let default_size = 100.0;
651
652 self.batch_modeling_cmd(
653 batch_context,
654 plane_id,
655 source_range,
656 &ModelingCmd::from(
657 mcmd::MakePlane::builder()
658 .clobber(false)
659 .origin(info.origin.into())
660 .size(LengthUnit(default_size))
661 .x_axis(info.x_axis.into())
662 .y_axis(info.y_axis.into())
663 .hide(true)
664 .build(),
665 ),
666 )
667 .await?;
668
669 if let Some(color) = color {
670 self.batch_modeling_cmd(
672 batch_context,
673 id_generator.next_uuid(),
674 source_range,
675 &ModelingCmd::from(mcmd::PlaneSetColor::builder().color(color).plane_id(plane_id).build()),
676 )
677 .await?;
678 }
679
680 Ok(plane_id)
681 }
682
683 async fn new_default_planes(
684 &self,
685 batch_context: &EngineBatchContext,
686 id_generator: &mut IdGenerator,
687 source_range: SourceRange,
688 geometry_only: bool,
689 ) -> Result<DefaultPlanes, KclError> {
690 let plane_opacity = 0.1;
691 let plane_color =
692 |red, green, blue| (!geometry_only).then(|| Color::from_rgba(red, green, blue, plane_opacity));
693 let plane_settings: Vec<(PlaneName, Uuid, Option<Color>)> = vec![
694 (PlaneName::Xy, id_generator.next_uuid(), plane_color(0.7, 0.28, 0.28)),
695 (PlaneName::Yz, id_generator.next_uuid(), plane_color(0.28, 0.7, 0.28)),
696 (PlaneName::Xz, id_generator.next_uuid(), plane_color(0.28, 0.28, 0.7)),
697 (PlaneName::NegXy, id_generator.next_uuid(), None),
698 (PlaneName::NegYz, id_generator.next_uuid(), None),
699 (PlaneName::NegXz, id_generator.next_uuid(), None),
700 ];
701
702 let mut planes = HashMap::new();
703 for (name, plane_id, color) in plane_settings {
704 let info = DEFAULT_PLANE_INFO.get(&name).ok_or_else(|| {
705 KclError::new_engine(KclErrorDetails::new(
707 format!("Failed to get default plane info for: {name:?}"),
708 vec![source_range],
709 ))
710 })?;
711 planes.insert(
712 name,
713 self.make_default_plane(batch_context, plane_id, info, color, source_range, id_generator)
714 .await?,
715 );
716 }
717
718 self.flush_batch(batch_context, false, source_range).await?;
720
721 Ok(DefaultPlanes {
722 xy: planes[&PlaneName::Xy],
723 neg_xy: planes[&PlaneName::NegXy],
724 xz: planes[&PlaneName::Xz],
725 neg_xz: planes[&PlaneName::NegXz],
726 yz: planes[&PlaneName::Yz],
727 neg_yz: planes[&PlaneName::NegYz],
728 })
729 }
730
731 fn parse_websocket_response(
732 &self,
733 response: WebSocketResponse,
734 source_range: SourceRange,
735 ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
736 match response {
737 WebSocketResponse::Success(success) => Ok(success.resp),
738 WebSocketResponse::Failure(fail) => {
739 let _request_id = fail.request_id;
740 if fail.errors.is_empty() {
741 return Err(KclError::new_engine(KclErrorDetails::new(
742 "Failure response with no error details".to_owned(),
743 vec![source_range],
744 )));
745 }
746 Err(KclError::new_engine(KclErrorDetails::new(
747 fail.errors
748 .iter()
749 .map(|e| e.message.clone())
750 .collect::<Vec<_>>()
751 .join("\n"),
752 vec![source_range],
753 )))
754 }
755 }
756 }
757
758 fn parse_batch_responses(
759 &self,
760 id: uuid::Uuid,
762 id_to_source_range: HashMap<uuid::Uuid, SourceRange>,
764 id_to_command: HashMap<uuid::Uuid, ModelingCmdEndpoint>,
766 responses: HashMap<kcmc::id::ModelingCmdId, BatchResponse>,
768 ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
769 let mut any_err: Option<crate::errors::KclError> = None;
770 let mut target_ok: Option<OkWebSocketResponseData> = None;
771 #[expect(
774 clippy::iter_over_hash_type,
775 reason = "modeling command uses a HashMap and keys are random, so we don't really have a choice"
776 )]
777 for (cmd_id, resp) in responses.iter() {
778 let cmd_id = Uuid::from(*cmd_id);
779 match resp {
780 BatchResponse::Success { response } if cmd_id == id => {
781 target_ok = Some(OkWebSocketResponseData::Modeling {
784 modeling_response: response.clone(),
785 });
786 }
787 BatchResponse::Success { .. } => continue,
788 BatchResponse::Failure { errors } => {
789 let command = id_to_command
790 .get(&cmd_id)
791 .map(ModelingCmdEndpoint::to_string)
792 .unwrap_or("[missing entry]".to_string());
793 let source_range = id_to_source_range.get(&cmd_id).cloned().ok_or_else(|| {
795 KclError::new_engine(KclErrorDetails::new(
796 format!("Failed to get source range for command {command} with ID: {cmd_id:?}"),
797 vec![],
798 ))
799 })?;
800 if errors.is_empty() {
801 any_err = Some(KclError::new_engine(KclErrorDetails::new(
802 format!("Failure response for batch with no error details at command {command}"),
803 vec![source_range],
804 )));
805 break;
806 }
807 let errors = errors.iter().map(|e| e.message.clone()).collect::<Vec<_>>().join("\n");
808 any_err = Some(KclError::new_engine(KclErrorDetails::new(
809 format!("command {command} resulted in errors: \n {errors}"),
810 vec![source_range],
811 )));
812 break;
813 }
814 }
815 }
816
817 match (any_err, target_ok) {
818 (Some(err), _) => Err(err),
819 (None, Some(ok)) => Ok(ok),
820 (None, None) => {
821 Err(KclError::new_engine(KclErrorDetails::new(
824 format!("Failed to find response for command ID: {id:?}"),
825 vec![],
826 )))
827 }
828 }
829 }
830
831 async fn set_user_colors(
832 &self,
833 batch_context: &EngineBatchContext,
834 settings: &ExecutorSettings,
835 source_range: SourceRange,
836 id_generator: &mut IdGenerator,
837 ) -> Result<(), KclError> {
838 let bf = settings
839 .default_backface_color
840 .clone()
841 .unwrap_or(default_backface_color());
842 let backface = csscolorparser::parse(&bf)
843 .map(|color| kcmc::shared::Color::from_rgba(color.r, color.g, color.b, color.a))
844 .unwrap_or(default_backface_color_struct());
845 self.batch_modeling_cmd(
846 batch_context,
847 id_generator.next_uuid(),
848 source_range,
849 &ModelingCmd::from(
850 mcmd::SetDefaultSystemProperties::builder()
851 .backface_color(backface)
852 .build(),
853 ),
854 )
855 .await?;
856 Ok(())
857 }
858
859 async fn modify_grid(
860 &self,
861 batch_context: &EngineBatchContext,
862 hidden: bool,
863 grid_scale_behavior: GridScaleBehavior,
864 source_range: SourceRange,
865 id_generator: &mut IdGenerator,
866 ) -> Result<(), KclError> {
867 self.batch_modeling_cmd(
869 batch_context,
870 id_generator.next_uuid(),
871 source_range,
872 &ModelingCmd::from(
873 mcmd::ObjectVisible::builder()
874 .hidden(hidden)
875 .object_id(*GRID_OBJECT_ID)
876 .build(),
877 ),
878 )
879 .await?;
880
881 self.batch_modeling_cmd(
882 batch_context,
883 id_generator.next_uuid(),
884 source_range,
885 &grid_scale_behavior.into_modeling_cmd(),
886 )
887 .await?;
888
889 self.batch_modeling_cmd(
891 batch_context,
892 id_generator.next_uuid(),
893 source_range,
894 &ModelingCmd::from(
895 mcmd::ObjectVisible::builder()
896 .hidden(hidden)
897 .object_id(*GRID_SCALE_TEXT_OBJECT_ID)
898 .build(),
899 ),
900 )
901 .await?;
902
903 Ok(())
904 }
905
906 pub async fn clear_queues(&self, batch_context: &EngineBatchContext) {
907 batch_context.clear().await;
908 self.ids_of_async_commands().write().await.clear();
909 self.async_tasks().clear().await;
910 }
911
912 fn responses(&self) -> Arc<RwLock<IndexMap<Uuid, WebSocketResponse>>> {
913 self.responses.responses.clone()
914 }
915
916 fn ids_of_async_commands(&self) -> Arc<RwLock<IndexMap<Uuid, SourceRange>>> {
917 self.ids_of_async_commands.clone()
918 }
919
920 fn async_tasks(&self) -> AsyncTasks {
921 self.async_tasks.clone()
922 }
923
924 pub fn stats(&self) -> &EngineStats {
925 &self.stats
926 }
927
928 pub fn get_default_planes(&self) -> Arc<RwLock<Option<DefaultPlanes>>> {
929 self.default_planes.clone()
930 }
931
932 async fn clear_scene_post_hook(
933 &self,
934 batch_context: &EngineBatchContext,
935 id_generator: &mut IdGenerator,
936 source_range: SourceRange,
937 geometry_only: bool,
938 ) -> Result<(), KclError> {
939 let new_planes = self
941 .new_default_planes(batch_context, id_generator, source_range, geometry_only)
942 .await?;
943 *self.default_planes.write().await = Some(new_planes);
944
945 self.transport.start_new_session(source_range).await?;
946
947 Ok(())
948 }
949
950 async fn inner_send_modeling_cmd(
951 &self,
952 id: uuid::Uuid,
953 source_range: SourceRange,
954 cmd: WebSocketRequest,
955 id_to_source_range: HashMap<Uuid, SourceRange>,
956 ) -> Result<WebSocketResponse, KclError> {
957 let response = self
958 .transport
959 .inner_send_modeling_cmd(id, source_range, cmd, id_to_source_range)
960 .await?;
961
962 self.responses.add(id, response.clone()).await;
963 Ok(response)
964 }
965
966 pub async fn get_session_data(&self) -> Option<ModelingSessionData> {
967 self.session_data.read().await.clone()
968 }
969
970 pub fn websocket_upgrade_request_id(&self) -> Option<&str> {
972 self.websocket_upgrade_request_id.as_deref()
973 }
974
975 pub async fn close(&self) {
976 let _ = self.transport.close().await;
977 }
978}
979
980#[derive(Debug, PartialEq)]
982pub enum SocketHealth {
983 Active,
984 Inactive,
985}