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::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/// Information about the responses from the engine.
55#[derive(Clone, Debug)]
56pub struct ResponseInformation {
57    /// The responses from the engine.
58    responses: Arc<RwLock<IndexMap<uuid::Uuid, WebSocketResponse>>>,
59}
60
61impl ResponseInformation {
62    /// Basic constructor.
63    pub fn new(responses: Arc<RwLock<IndexMap<uuid::Uuid, WebSocketResponse>>>) -> Self {
64        Self { responses }
65    }
66
67    /// Add a new response from the engine.
68    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    // Replaces `engine_req_tx: mpsc::Sender<ToEngineReq>`
76    // from the original native connection type.
77    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    /// The default planes for the scene.
84    #[builder(default)]
85    default_planes: Arc<RwLock<Option<DefaultPlanes>>>,
86    /// If the server sends session data, it'll be copied to here.
87    session_data: Arc<RwLock<Option<ModelingSessionData>>>,
88
89    /// Request ID returned by the HTTP request that upgraded to this WebSocket.
90    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    /// Mock connection that doesn't actually connect to anything.
188    /// Used for testing.
189    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    /// Take the ids of async commands that have accumulated so far and clear them.
212    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    /// Take the responses that have accumulated so far and clear them.
217    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        // Clear any batched commands leftover from previous scenes.
230        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        // Flush the batch queue, so clear is run right away.
251        // Otherwise the hooks below won't work.
252        self.flush_batch(batch_context, false, source_range).await?;
253
254        // Do the after clear scene hook.
255        self.clear_scene_post_hook(batch_context, id_generator, source_range, geometry_only)
256            .await?;
257
258        Ok(())
259    }
260
261    /// Ensure a specific async command has been completed.
262    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            // Look it up if we don't have it.
271            self.ids_of_async_commands()
272                .read()
273                .await
274                .get(&id)
275                .cloned()
276                .unwrap_or_default()
277        };
278
279        // The previous 60s ceiling here was too aggressive for long-running
280        // engine commands - notably large STEP / B-rep imports, which the
281        // engine itself routinely takes several minutes to process. When the
282        // ceiling fired first the user got a generic "async command timed
283        // out" message and the eventual engine response (success OR error)
284        // was discarded, masking the real outcome. 600s (10 min) gives the
285        // engine room to finish or to surface its own error.
286        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                // Yield to the event loop so that we don’t block the UI thread.
292                // No seriously WE DO NOT WANT TO PAUSE THE WHOLE APP ON THE JS SIDE.
293                #[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            // If the response is an error, return it.
309            // Parsing will do that and we can ignore the result, we don't care.
310            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    /// Ensure ALL async commands have been completed.
323    pub async fn ensure_async_commands_completed(&self, batch_context: &EngineBatchContext) -> Result<(), KclError> {
324        // Check if all async commands have been completed.
325        let ids = self.take_ids_of_async_commands().await;
326
327        // Try to get them from the responses.
328        for (id, source_range) in ids {
329            self.ensure_async_command_completed(id, Some(source_range)).await?;
330        }
331
332        // Make sure we check for all async tasks as well.
333        // The reason why we ignore the error here is that, if a model fillets an edge
334        // we previously called something on, it might no longer exist. In which case,
335        // the artifact graph won't care either if its gone since you can't select it
336        // anymore anyways.
337        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        // Flush the batch to make sure nothing remains.
345        self.flush_batch(batch_context, true, SourceRange::default()).await?;
346
347        Ok(())
348    }
349
350    /// Set the visibility of edges.
351    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    /// Re-run the command to apply the settings.
370    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        // Set the edge visibility.
382        self.set_edge_visibility(batch_context, settings.highlight_edges, source_range, id_generator)
383            .await?;
384
385        // Send the command to show the grid.
386
387        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        // Set up user's color choices.
397        self.set_user_colors(batch_context, settings, source_range, id_generator)
398            .await?;
399
400        // We do not have commands for changing ssao on the fly.
401
402        // Flush the batch queue, so the settings are applied right away.
403        self.flush_batch(batch_context, false, source_range).await?;
404
405        Ok(())
406    }
407
408    // Add a modeling command to the batch but don't fire it right away.
409    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        // Add cmd to the batch.
422        batch_context.push(req, source_range).await;
423        self.stats().commands_batched.fetch_add(1, Relaxed);
424
425        Ok(())
426    }
427
428    // Add a vector of modeling commands to the batch but don't fire it right away.
429    // This allows you to force them all to be added together in the same order.
430    // When we are running things in parallel this prevents race conditions that might come
431    // if specific commands are run before others.
432    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        // Add cmds to the batch.
439        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    /// Add a command to the batch that needs to be executed at the very end.
450    /// This for stuff like fillets or chamfers where if we execute too soon the
451    /// engine will eat the ID and we can't reference it for other commands.
452    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        // Add cmd to the batch end.
465        batch_context.insert_end(id, req, source_range).await;
466        self.stats().commands_batched.fetch_add(1, Relaxed);
467        Ok(())
468    }
469
470    /// Send the modeling cmd and wait for the response.
471    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        // Add the command to the batch.
481        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        // Flush the batch queue.
491        self.run_batch(requests, source_range).await
492    }
493
494    /// Send the modeling cmd async and don't wait for the response.
495    /// Add it to our list of async commands.
496    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        // Add the command ID to the list of async commands.
503        self.ids_of_async_commands().write().await.insert(id, source_range);
504
505        // Fire off the command now, but don't wait for the response, we don't care about it.
506        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    /// Run the batch for the specific commands.
522    async fn run_batch(
523        &self,
524        orig_requests: Vec<(WebSocketRequest, SourceRange)>,
525        source_range: SourceRange,
526    ) -> Result<OkWebSocketResponseData, crate::errors::KclError> {
527        // Return early if we have no commands to send.
528        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            // We can unwrap here because we know the batch has only one element.
553            orig_requests.first().unwrap().0.clone()
554        } else {
555            batched_requests
556        };
557
558        // Create the map of original command IDs to source range.
559        // This is for the wasm side, kurt needs it for selections.
560        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        // We pop off the responses to cleanup our mappings.
582        match final_req {
583            WebSocketRequest::ModelingCmdBatchReq(ModelingBatch {
584                ref requests,
585                batch_id,
586                responses: _,
587            }) => {
588                // Get the last command ID.
589                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 we have a batch response, we want to return the specific id we care about.
596                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                    // We should never get here.
600                    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                // You are probably wondering why we can't just return the source range we were
608                // passed with the function. Well this is actually really important.
609                // If this is the last command in the batch and there is only one and we've reached
610                // the end of the file, this will trigger a flush batch function, but it will just
611                // send default or the end of the file as it's source range not the origin of the
612                // request so we need the original request source range in case the engine returns
613                // an error.
614                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    /// Force flush the batch queue.
633    pub async fn flush_batch(
634        &self,
635        batch_context: &EngineBatchContext,
636        // Whether or not to flush the end commands as well.
637        // We only do this at the very end of the file.
638        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        // Create new default planes.
662        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            // Set the color.
683            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                // We should never get here.
718                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        // Flush the batch queue, so these planes are created right away.
731        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        // The last response we are looking for.
773        id: uuid::Uuid,
774        // The mapping of source ranges to command IDs.
775        id_to_source_range: HashMap<uuid::Uuid, SourceRange>,
776        // Allows us to print which command failed
777        id_to_command: HashMap<uuid::Uuid, ModelingCmdEndpoint>,
778        // The response from the engine.
779        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        // Iterate over the responses and check for errors.
784        // Any error takes precedent over any Ok.
785        #[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                    // This is the response we care about.
794                    // Keep looking for errors after locating it.
795                    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                    // Get the source range for the command.
806                    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                // Return an error that we did not get an error or the response we wanted.
834                // This should never happen but who knows.
835                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        // Hide/show the grid.
880        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        // Hide/show the grid scale text.
902        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        // Remake the default planes, since they would have been removed after the scene was cleared.
952        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    /// Request ID returned by the HTTP request that upgraded to this WebSocket.
983    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/// State of the connection to the engine.
993#[derive(Debug, PartialEq)]
994pub enum SocketHealth {
995    Active,
996    Inactive,
997}