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