Skip to main content

kcl_lib/engine/
engine_manager.rs

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/// Information about the responses from the engine.
54#[derive(Clone, Debug)]
55pub struct ResponseInformation {
56    /// The responses from the engine.
57    responses: Arc<RwLock<IndexMap<uuid::Uuid, WebSocketResponse>>>,
58}
59
60impl ResponseInformation {
61    /// Basic constructor.
62    pub fn new(responses: Arc<RwLock<IndexMap<uuid::Uuid, WebSocketResponse>>>) -> Self {
63        Self { responses }
64    }
65
66    /// Add a new response from the engine.
67    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    // Replaces `engine_req_tx: mpsc::Sender<ToEngineReq>`
75    // from the original native connection type.
76    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    /// The default planes for the scene.
83    #[builder(default)]
84    default_planes: Arc<RwLock<Option<DefaultPlanes>>>,
85    /// If the server sends session data, it'll be copied to here.
86    session_data: Arc<RwLock<Option<ModelingSessionData>>>,
87
88    /// Request ID returned by the HTTP request that upgraded to this WebSocket.
89    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    /// Mock connection that doesn't actually connect to anything.
187    /// Used for testing.
188    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    /// Take the ids of async commands that have accumulated so far and clear them.
211    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    /// Take the responses that have accumulated so far and clear them.
216    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        // Clear any batched commands leftover from previous scenes.
228        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        // Flush the batch queue, so clear is run right away.
239        // Otherwise the hooks below won't work.
240        self.flush_batch(batch_context, false, source_range).await?;
241
242        // Do the after clear scene hook.
243        self.clear_scene_post_hook(batch_context, id_generator, source_range, geometry_only)
244            .await?;
245
246        Ok(())
247    }
248
249    /// Ensure a specific async command has been completed.
250    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            // Look it up if we don't have it.
259            self.ids_of_async_commands()
260                .read()
261                .await
262                .get(&id)
263                .cloned()
264                .unwrap_or_default()
265        };
266
267        // The previous 60s ceiling here was too aggressive for long-running
268        // engine commands - notably large STEP / B-rep imports, which the
269        // engine itself routinely takes several minutes to process. When the
270        // ceiling fired first the user got a generic "async command timed
271        // out" message and the eventual engine response (success OR error)
272        // was discarded, masking the real outcome. 600s (10 min) gives the
273        // engine room to finish or to surface its own error.
274        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                // Yield to the event loop so that we don’t block the UI thread.
280                // No seriously WE DO NOT WANT TO PAUSE THE WHOLE APP ON THE JS SIDE.
281                #[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            // If the response is an error, return it.
297            // Parsing will do that and we can ignore the result, we don't care.
298            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    /// Ensure ALL async commands have been completed.
311    pub async fn ensure_async_commands_completed(&self, batch_context: &EngineBatchContext) -> Result<(), KclError> {
312        // Check if all async commands have been completed.
313        let ids = self.take_ids_of_async_commands().await;
314
315        // Try to get them from the responses.
316        for (id, source_range) in ids {
317            self.ensure_async_command_completed(id, Some(source_range)).await?;
318        }
319
320        // Make sure we check for all async tasks as well.
321        // The reason why we ignore the error here is that, if a model fillets an edge
322        // we previously called something on, it might no longer exist. In which case,
323        // the artifact graph won't care either if its gone since you can't select it
324        // anymore anyways.
325        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        // Flush the batch to make sure nothing remains.
333        self.flush_batch(batch_context, true, SourceRange::default()).await?;
334
335        Ok(())
336    }
337
338    /// Set the visibility of edges.
339    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    /// Re-run the command to apply the settings.
358    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        // Set the edge visibility.
370        self.set_edge_visibility(batch_context, settings.highlight_edges, source_range, id_generator)
371            .await?;
372
373        // Send the command to show the grid.
374
375        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        // Set up user's color choices.
385        self.set_user_colors(batch_context, settings, source_range, id_generator)
386            .await?;
387
388        // We do not have commands for changing ssao on the fly.
389
390        // Flush the batch queue, so the settings are applied right away.
391        self.flush_batch(batch_context, false, source_range).await?;
392
393        Ok(())
394    }
395
396    // Add a modeling command to the batch but don't fire it right away.
397    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        // Add cmd to the batch.
410        batch_context.push(req, source_range).await;
411        self.stats().commands_batched.fetch_add(1, Relaxed);
412
413        Ok(())
414    }
415
416    // Add a vector of modeling commands to the batch but don't fire it right away.
417    // This allows you to force them all to be added together in the same order.
418    // When we are running things in parallel this prevents race conditions that might come
419    // if specific commands are run before others.
420    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        // Add cmds to the batch.
427        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    /// Add a command to the batch that needs to be executed at the very end.
438    /// This for stuff like fillets or chamfers where if we execute too soon the
439    /// engine will eat the ID and we can't reference it for other commands.
440    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        // Add cmd to the batch end.
453        batch_context.insert_end(id, req, source_range).await;
454        self.stats().commands_batched.fetch_add(1, Relaxed);
455        Ok(())
456    }
457
458    /// Send the modeling cmd and wait for the response.
459    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        // Add the command to the batch.
469        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        // Flush the batch queue.
479        self.run_batch(requests, source_range).await
480    }
481
482    /// Send the modeling cmd async and don't wait for the response.
483    /// Add it to our list of async commands.
484    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        // Add the command ID to the list of async commands.
491        self.ids_of_async_commands().write().await.insert(id, source_range);
492
493        // Fire off the command now, but don't wait for the response, we don't care about it.
494        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    /// Run the batch for the specific commands.
510    async fn run_batch(
511        &self,
512        orig_requests: Vec<(WebSocketRequest, SourceRange)>,
513        source_range: SourceRange,
514    ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
515        // Return early if we have no commands to send.
516        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            // We can unwrap here because we know the batch has only one element.
541            orig_requests.first().unwrap().0.clone()
542        } else {
543            batched_requests
544        };
545
546        // Create the map of original command IDs to source range.
547        // This is for the wasm side, kurt needs it for selections.
548        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        // We pop off the responses to cleanup our mappings.
570        match final_req {
571            WebSocketRequest::ModelingCmdBatchReq(ModelingBatch {
572                ref requests,
573                batch_id,
574                responses: _,
575            }) => {
576                // Get the last command ID.
577                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 we have a batch response, we want to return the specific id we care about.
584                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                    // We should never get here.
588                    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                // You are probably wondering why we can't just return the source range we were
596                // passed with the function. Well this is actually really important.
597                // If this is the last command in the batch and there is only one and we've reached
598                // the end of the file, this will trigger a flush batch function, but it will just
599                // send default or the end of the file as it's source range not the origin of the
600                // request so we need the original request source range in case the engine returns
601                // an error.
602                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    /// Force flush the batch queue.
621    pub async fn flush_batch(
622        &self,
623        batch_context: &EngineBatchContext,
624        // Whether or not to flush the end commands as well.
625        // We only do this at the very end of the file.
626        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        // Create new default planes.
650        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            // Set the color.
671            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                // We should never get here.
706                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        // Flush the batch queue, so these planes are created right away.
719        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        // The last response we are looking for.
761        id: uuid::Uuid,
762        // The mapping of source ranges to command IDs.
763        id_to_source_range: HashMap<uuid::Uuid, SourceRange>,
764        // Allows us to print which command failed
765        id_to_command: HashMap<uuid::Uuid, ModelingCmdEndpoint>,
766        // The response from the engine.
767        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        // Iterate over the responses and check for errors.
772        // Any error takes precedent over any Ok.
773        #[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                    // This is the response we care about.
782                    // Keep looking for errors after locating it.
783                    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                    // Get the source range for the command.
794                    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                // Return an error that we did not get an error or the response we wanted.
822                // This should never happen but who knows.
823                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        // Hide/show the grid.
868        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        // Hide/show the grid scale text.
890        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        // Remake the default planes, since they would have been removed after the scene was cleared.
940        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    /// Request ID returned by the HTTP request that upgraded to this WebSocket.
971    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/// State of the connection to the engine.
981#[derive(Debug, PartialEq)]
982pub enum SocketHealth {
983    Active,
984    Inactive,
985}