1use std::collections::HashMap;
4
5use agentos_sidecar_client::wire;
6use tokio::sync::{broadcast, watch};
7
8use crate::agent_os::AgentOs;
9use crate::agent_os::ProcessEntry;
10use crate::error::{ClientError, ClientResult};
11use crate::process::{ProcessOutput, ProcessStream};
12
13#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
14pub struct ContextDescriptor {
15 pub context_id: String,
16 pub state: String,
17 pub language: Option<String>,
18 pub created_at_ms: u64,
19 pub last_started_at_ms: Option<u64>,
20 pub last_completed_at_ms: Option<u64>,
21}
22
23#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
24pub struct ProcessDescriptor {
25 pub pid: u32,
26 pub state: String,
27 pub language: Option<String>,
28 pub started_at_ms: u64,
29}
30
31pub type CodeExecutionResult = wire::ExecutionCompletedResponse;
32pub type ExecutionOutputEvent = wire::ExecutionOutputEvent;
33pub type ExecutionCompletedEvent = wire::ExecutionCompletedEvent;
34pub type TypeScriptDiagnostic = wire::TypeScriptDiagnostic;
35
36#[derive(Debug, Clone, Default)]
37pub struct LanguageExecutionOptions {
38 pub context_id: Option<String>,
39 pub cwd: Option<String>,
40 pub env: HashMap<String, String>,
41 pub args: Vec<String>,
42 pub stdin: Option<Vec<u8>>,
43 pub timeout_ms: Option<u64>,
44 pub pty: Option<ExecutionPtyOptions>,
45 pub output: ExecutionOutputOptions,
46}
47
48#[derive(Debug, Clone, Default)]
49pub struct LanguageSpawnOptions {
50 pub cwd: Option<String>,
51 pub env: HashMap<String, String>,
52 pub args: Vec<String>,
53 pub stdin: Option<Vec<u8>>,
54 pub timeout_ms: Option<u64>,
55 pub pty: Option<ExecutionPtyOptions>,
56 pub retain_events: bool,
57}
58
59#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
60pub enum OutputCapture {
61 #[default]
62 None,
63 Stderr,
64 All,
65}
66
67#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
68pub struct ExecutionOutputOptions {
69 pub capture: OutputCapture,
70 pub retain_events: bool,
71}
72
73#[derive(Debug, Clone, Copy, Default)]
74pub struct ExecutionPtyOptions {
75 pub cols: Option<u16>,
76 pub rows: Option<u16>,
77}
78
79#[derive(Debug, Clone, Default)]
80pub struct InlineExecutionOptions {
81 pub process: LanguageExecutionOptions,
82 pub inputs: Option<serde_json::Map<String, serde_json::Value>>,
83}
84
85#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
86pub enum JavaScriptModuleFormat {
87 #[default]
88 Module,
89 CommonJs,
90}
91
92#[derive(Debug, Clone, Default)]
93pub struct JavaScriptExecutionOptions {
94 pub inline: InlineExecutionOptions,
95 pub format: JavaScriptModuleFormat,
96 pub file_path: Option<String>,
97}
98
99#[derive(Debug, Clone, Default)]
100pub struct TypeScriptExecutionOptions {
101 pub inline: InlineExecutionOptions,
102 pub file_path: Option<String>,
103 pub tsconfig_path: Option<String>,
104 pub compiler_options: Option<serde_json::Map<String, serde_json::Value>>,
105}
106
107#[derive(Debug, Clone, Default)]
108pub struct TypeScriptCheckOptions {
109 pub context_id: Option<String>,
110 pub cwd: Option<String>,
111 pub file_path: Option<String>,
112 pub tsconfig_path: Option<String>,
113 pub compiler_options: Option<serde_json::Map<String, serde_json::Value>>,
114 pub timeout_ms: Option<u64>,
115 pub output: ExecutionOutputOptions,
116}
117
118#[derive(Debug, Clone, Default)]
119pub struct NpmProjectInstallOptions {
120 pub context_id: Option<String>,
121 pub cwd: Option<String>,
122 pub env: HashMap<String, String>,
123 pub timeout_ms: Option<u64>,
124 pub frozen: Option<bool>,
125 pub output: ExecutionOutputOptions,
126}
127
128#[derive(Debug, Clone, Default)]
129pub struct NpmPackageInstallOptions {
130 pub context_id: Option<String>,
131 pub cwd: Option<String>,
132 pub env: HashMap<String, String>,
133 pub timeout_ms: Option<u64>,
134 pub dev: Option<bool>,
135 pub global: Option<bool>,
136 pub output: ExecutionOutputOptions,
137}
138
139#[derive(Debug, Clone, Default)]
140pub struct PythonInstallOptions {
141 pub context_id: Option<String>,
142 pub cwd: Option<String>,
143 pub env: HashMap<String, String>,
144 pub timeout_ms: Option<u64>,
145 pub upgrade: Option<bool>,
146 pub requirements_file: Option<String>,
147 pub index_url: Option<String>,
148 pub extra_index_urls: Vec<String>,
149 pub output: ExecutionOutputOptions,
150}
151
152#[derive(Debug, Clone)]
153pub struct TypeScriptCheckResult {
154 pub result: CodeExecutionResult,
155 pub has_errors: Option<bool>,
156 pub diagnostics: Vec<TypeScriptDiagnostic>,
157}
158
159#[derive(Debug, Clone)]
160enum ExecutionSubmission {
161 Completed(CodeExecutionResult),
162 Background(wire::ExecutionDescriptor),
163}
164
165#[derive(Debug, Clone)]
166pub struct CodeEvaluationResult {
167 pub result: CodeExecutionResult,
168 pub value: Option<serde_json::Value>,
169}
170
171fn identity(options: &LanguageExecutionOptions) -> wire::ExecutionIdentityOptions {
172 wire::ExecutionIdentityOptions {
173 context_id: options.context_id.clone(),
174 }
175}
176
177fn process(options: &LanguageExecutionOptions) -> wire::ProcessExecutionOptions {
178 wire::ProcessExecutionOptions {
179 identity: identity(options),
180 output: output(options.output),
181 operation_id: None,
182 background: Some(false),
183 cwd: options.cwd.clone(),
184 env: (!options.env.is_empty()).then(|| options.env.clone()),
185 args: options.args.clone(),
186 stdin: options.stdin.clone(),
187 timeout_ms: options.timeout_ms,
188 pty: options.pty.map(|pty| wire::ExecutionPtyOptions {
189 cols: pty.cols,
190 rows: pty.rows,
191 }),
192 }
193}
194
195fn background_process(
196 options: &LanguageSpawnOptions,
197 operation_id: String,
198) -> wire::ProcessExecutionOptions {
199 wire::ProcessExecutionOptions {
200 identity: wire::ExecutionIdentityOptions { context_id: None },
201 output: wire::ExecutionOutputOptions {
202 capture: Some(wire::ExecutionOutputCapture::None),
203 retain_events: Some(options.retain_events),
204 },
205 operation_id: Some(operation_id),
206 background: Some(true),
207 cwd: options.cwd.clone(),
208 env: (!options.env.is_empty()).then(|| options.env.clone()),
209 args: options.args.clone(),
210 stdin: options.stdin.clone(),
211 timeout_ms: options.timeout_ms,
212 pty: options.pty.map(|pty| wire::ExecutionPtyOptions {
213 cols: pty.cols,
214 rows: pty.rows,
215 }),
216 }
217}
218
219fn output(options: ExecutionOutputOptions) -> wire::ExecutionOutputOptions {
220 wire::ExecutionOutputOptions {
221 capture: Some(match options.capture {
222 OutputCapture::None => wire::ExecutionOutputCapture::None,
223 OutputCapture::Stderr => wire::ExecutionOutputCapture::Stderr,
224 OutputCapture::All => wire::ExecutionOutputCapture::All,
225 }),
226 retain_events: Some(options.retain_events),
227 }
228}
229
230fn json_inputs(options: &InlineExecutionOptions) -> ClientResult<Option<String>> {
231 options
232 .inputs
233 .as_ref()
234 .map(serde_json::to_string)
235 .transpose()
236 .map_err(|error| ClientError::Sidecar(format!("failed to serialize inputs: {error}")))
237}
238
239fn context_descriptor(descriptor: wire::ExecutionDescriptor) -> ContextDescriptor {
240 ContextDescriptor {
241 context_id: descriptor.execution_id,
242 state: match descriptor.state {
243 wire::ExecutionState::Creating | wire::ExecutionState::Idle => "idle",
244 wire::ExecutionState::Running => "running",
245 wire::ExecutionState::Resetting => "resetting",
246 wire::ExecutionState::Deleting => "deleting",
247 wire::ExecutionState::Failed => "failed",
248 }
249 .to_owned(),
250 language: descriptor.retained_language.map(|language| match language {
251 wire::RetainedExecutionLanguage::JavaScript => String::from("javascript"),
252 wire::RetainedExecutionLanguage::Python => String::from("python"),
253 }),
254 created_at_ms: descriptor.created_at_ms,
255 last_started_at_ms: descriptor.last_started_at_ms,
256 last_completed_at_ms: descriptor.last_completed_at_ms,
257 }
258}
259
260fn execution_request_options(
261 payload: &wire::RequestPayload,
262) -> Option<(
263 &wire::ExecutionIdentityOptions,
264 &wire::ExecutionOutputOptions,
265)> {
266 Some(match payload {
267 wire::RequestPayload::ShellExecutionRequest(request) => {
268 (&request.process.identity, &request.process.output)
269 }
270 wire::RequestPayload::ArgvExecutionRequest(request) => {
271 (&request.process.identity, &request.process.output)
272 }
273 wire::RequestPayload::JavaScriptExecutionRequest(request) => {
274 (&request.process.identity, &request.process.output)
275 }
276 wire::RequestPayload::JavaScriptEvaluationRequest(request) => {
277 (&request.process.identity, &request.process.output)
278 }
279 wire::RequestPayload::JavaScriptFileExecutionRequest(request) => {
280 (&request.process.identity, &request.process.output)
281 }
282 wire::RequestPayload::TypeScriptExecutionRequest(request) => {
283 (&request.process.identity, &request.process.output)
284 }
285 wire::RequestPayload::TypeScriptEvaluationRequest(request) => {
286 (&request.process.identity, &request.process.output)
287 }
288 wire::RequestPayload::TypeScriptFileExecutionRequest(request) => {
289 (&request.process.identity, &request.process.output)
290 }
291 wire::RequestPayload::TypeScriptCheckRequest(request) => {
292 (&request.identity, &request.output)
293 }
294 wire::RequestPayload::TypeScriptProjectCheckRequest(request) => {
295 (&request.identity, &request.output)
296 }
297 wire::RequestPayload::NpmProjectInstallRequest(request) => {
298 (&request.identity, &request.output)
299 }
300 wire::RequestPayload::NpmPackageInstallRequest(request) => {
301 (&request.identity, &request.output)
302 }
303 wire::RequestPayload::NpmScriptExecutionRequest(request) => {
304 (&request.process.identity, &request.process.output)
305 }
306 wire::RequestPayload::NpmPackageExecutionRequest(request) => {
307 (&request.process.identity, &request.process.output)
308 }
309 wire::RequestPayload::PythonExecutionRequest(request) => {
310 (&request.process.identity, &request.process.output)
311 }
312 wire::RequestPayload::PythonEvaluationRequest(request) => {
313 (&request.process.identity, &request.process.output)
314 }
315 wire::RequestPayload::PythonFileExecutionRequest(request) => {
316 (&request.process.identity, &request.process.output)
317 }
318 wire::RequestPayload::PythonModuleExecutionRequest(request) => {
319 (&request.process.identity, &request.process.output)
320 }
321 wire::RequestPayload::PythonInstallRequest(request) => (&request.identity, &request.output),
322 _ => return None,
323 })
324}
325
326impl AgentOs {
327 fn execution_ownership(&self) -> wire::OwnershipScope {
328 let inner = self.inner();
329 wire::OwnershipScope::VmOwnership(wire::VmOwnership {
330 connection_id: inner.connection_id.clone(),
331 session_id: inner.session_id.clone(),
332 vm_id: inner.vm_id.clone(),
333 })
334 }
335
336 async fn submit_execution(
337 &self,
338 payload: wire::RequestPayload,
339 background: bool,
340 ) -> ClientResult<ExecutionSubmission> {
341 if let Some((identity, output)) = execution_request_options(&payload) {
342 if background && identity.context_id.is_some() {
343 return Err(ClientError::Sidecar(String::from(
344 "spawned language processes cannot use context_id",
345 )));
346 }
347 if output.retain_events == Some(true) && identity.context_id.is_none() && !background {
348 return Err(ClientError::Sidecar(String::from(
349 "retain_events requires a context or spawned language process",
350 )));
351 }
352 }
353 let mut events = self.transport().subscribe_wire_events();
354 let accepted = match self
355 .transport()
356 .request_wire(self.execution_ownership(), payload)
357 .await?
358 {
359 wire::ResponsePayload::ExecutionAcceptedResponse(response) => response,
360 wire::ResponsePayload::RejectedResponse(rejected) => {
361 return Err(ClientError::from_rejection(rejected))
362 }
363 response => {
364 return Err(ClientError::Sidecar(format!(
365 "unexpected execution response: {response:?}"
366 )))
367 }
368 };
369 if background {
370 return Ok(ExecutionSubmission::Background(
371 accepted.execution.ok_or_else(|| {
372 ClientError::Sidecar(String::from(
373 "spawned process admission returned no descriptor",
374 ))
375 })?,
376 ));
377 }
378 wait_for_completion_event(&mut events, &accepted.operation_id).await?;
379 Ok(ExecutionSubmission::Completed(
380 self.wait_execution(&accepted.operation_id).await?,
381 ))
382 }
383
384 pub async fn exec(
385 &self,
386 command: impl Into<String>,
387 options: LanguageExecutionOptions,
388 ) -> ClientResult<CodeExecutionResult> {
389 completed_submission(
390 self.submit_execution(
391 wire::RequestPayload::ShellExecutionRequest(wire::ShellExecutionRequest {
392 process: process(&options),
393 command: command.into(),
394 }),
395 false,
396 )
397 .await?,
398 )
399 }
400
401 pub async fn exec_argv(
402 &self,
403 command: impl Into<String>,
404 args: Vec<String>,
405 mut options: LanguageExecutionOptions,
406 ) -> ClientResult<CodeExecutionResult> {
407 options.args = args;
408 completed_submission(
409 self.submit_execution(
410 wire::RequestPayload::ArgvExecutionRequest(wire::ArgvExecutionRequest {
411 process: process(&options),
412 command: command.into(),
413 }),
414 false,
415 )
416 .await?,
417 )
418 }
419
420 pub async fn execute_javascript(
421 &self,
422 source: impl Into<String>,
423 options: JavaScriptExecutionOptions,
424 ) -> ClientResult<CodeExecutionResult> {
425 let inputs = json_inputs(&options.inline)?;
426 completed_submission(
427 self.submit_execution(
428 wire::RequestPayload::JavaScriptExecutionRequest(
429 wire::JavaScriptExecutionRequest {
430 process: process(&options.inline.process),
431 source: source.into(),
432 format: Some(match options.format {
433 JavaScriptModuleFormat::Module => wire::JavaScriptModuleFormat::Module,
434 JavaScriptModuleFormat::CommonJs => {
435 wire::JavaScriptModuleFormat::CommonJs
436 }
437 }),
438 file_path: options.file_path,
439 inputs,
440 },
441 ),
442 false,
443 )
444 .await?,
445 )
446 }
447
448 pub async fn evaluate_javascript(
449 &self,
450 expression: impl Into<String>,
451 options: JavaScriptExecutionOptions,
452 ) -> ClientResult<CodeEvaluationResult> {
453 let inputs = json_inputs(&options.inline)?;
454 let submission = self
455 .submit_execution(
456 wire::RequestPayload::JavaScriptEvaluationRequest(
457 wire::JavaScriptEvaluationRequest {
458 process: process(&options.inline.process),
459 expression: expression.into(),
460 format: Some(match options.format {
461 JavaScriptModuleFormat::Module => wire::JavaScriptModuleFormat::Module,
462 JavaScriptModuleFormat::CommonJs => {
463 wire::JavaScriptModuleFormat::CommonJs
464 }
465 }),
466 file_path: options.file_path,
467 inputs,
468 },
469 ),
470 false,
471 )
472 .await?;
473 evaluation_result(submission)
474 }
475
476 pub async fn execute_javascript_file(
477 &self,
478 path: impl Into<String>,
479 options: LanguageExecutionOptions,
480 ) -> ClientResult<CodeExecutionResult> {
481 completed_submission(
482 self.submit_execution(
483 wire::RequestPayload::JavaScriptFileExecutionRequest(
484 wire::JavaScriptFileExecutionRequest {
485 process: process(&options),
486 path: path.into(),
487 },
488 ),
489 false,
490 )
491 .await?,
492 )
493 }
494
495 pub async fn spawn_javascript(
496 &self,
497 source: impl Into<String>,
498 options: LanguageSpawnOptions,
499 ) -> ClientResult<ProcessDescriptor> {
500 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
501 background_submission(
502 self,
503 self.submit_execution(
504 wire::RequestPayload::JavaScriptExecutionRequest(
505 wire::JavaScriptExecutionRequest {
506 process: background_process(&options, operation_id),
507 source: source.into(),
508 format: Some(wire::JavaScriptModuleFormat::Module),
509 file_path: None,
510 inputs: None,
511 },
512 ),
513 true,
514 )
515 .await?,
516 "javascript",
517 )
518 }
519
520 pub async fn spawn_javascript_file(
521 &self,
522 path: impl Into<String>,
523 options: LanguageSpawnOptions,
524 ) -> ClientResult<ProcessDescriptor> {
525 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
526 background_submission(
527 self,
528 self.submit_execution(
529 wire::RequestPayload::JavaScriptFileExecutionRequest(
530 wire::JavaScriptFileExecutionRequest {
531 process: background_process(&options, operation_id),
532 path: path.into(),
533 },
534 ),
535 true,
536 )
537 .await?,
538 "javascript",
539 )
540 }
541
542 pub async fn execute_typescript(
543 &self,
544 source: impl Into<String>,
545 options: TypeScriptExecutionOptions,
546 ) -> ClientResult<CodeExecutionResult> {
547 let inputs = json_inputs(&options.inline)?;
548 completed_submission(
549 self.submit_execution(
550 wire::RequestPayload::TypeScriptExecutionRequest(
551 wire::TypeScriptExecutionRequest {
552 process: process(&options.inline.process),
553 source: source.into(),
554 file_path: options.file_path,
555 tsconfig_path: options.tsconfig_path,
556 compiler_options: options
557 .compiler_options
558 .as_ref()
559 .map(serde_json::to_string)
560 .transpose()
561 .map_err(|error| ClientError::Sidecar(error.to_string()))?,
562 inputs,
563 },
564 ),
565 false,
566 )
567 .await?,
568 )
569 }
570
571 pub async fn evaluate_typescript(
572 &self,
573 expression: impl Into<String>,
574 options: TypeScriptExecutionOptions,
575 ) -> ClientResult<CodeEvaluationResult> {
576 let inputs = json_inputs(&options.inline)?;
577 let submission = self
578 .submit_execution(
579 wire::RequestPayload::TypeScriptEvaluationRequest(
580 wire::TypeScriptEvaluationRequest {
581 process: process(&options.inline.process),
582 expression: expression.into(),
583 file_path: options.file_path,
584 tsconfig_path: options.tsconfig_path,
585 compiler_options: options
586 .compiler_options
587 .as_ref()
588 .map(serde_json::to_string)
589 .transpose()
590 .map_err(|error| ClientError::Sidecar(error.to_string()))?,
591 inputs,
592 },
593 ),
594 false,
595 )
596 .await?;
597 evaluation_result(submission)
598 }
599
600 pub async fn execute_typescript_file(
601 &self,
602 path: impl Into<String>,
603 options: TypeScriptExecutionOptions,
604 ) -> ClientResult<CodeExecutionResult> {
605 completed_submission(
606 self.submit_execution(
607 wire::RequestPayload::TypeScriptFileExecutionRequest(
608 wire::TypeScriptFileExecutionRequest {
609 process: process(&options.inline.process),
610 path: path.into(),
611 tsconfig_path: options.tsconfig_path,
612 compiler_options: options
613 .compiler_options
614 .as_ref()
615 .map(serde_json::to_string)
616 .transpose()
617 .map_err(|error| ClientError::Sidecar(error.to_string()))?,
618 },
619 ),
620 false,
621 )
622 .await?,
623 )
624 }
625
626 pub async fn spawn_typescript(
627 &self,
628 source: impl Into<String>,
629 options: LanguageSpawnOptions,
630 ) -> ClientResult<ProcessDescriptor> {
631 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
632 background_submission(
633 self,
634 self.submit_execution(
635 wire::RequestPayload::TypeScriptExecutionRequest(
636 wire::TypeScriptExecutionRequest {
637 process: background_process(&options, operation_id),
638 source: source.into(),
639 file_path: None,
640 tsconfig_path: None,
641 compiler_options: None,
642 inputs: None,
643 },
644 ),
645 true,
646 )
647 .await?,
648 "javascript",
649 )
650 }
651
652 pub async fn spawn_typescript_file(
653 &self,
654 path: impl Into<String>,
655 options: LanguageSpawnOptions,
656 ) -> ClientResult<ProcessDescriptor> {
657 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
658 background_submission(
659 self,
660 self.submit_execution(
661 wire::RequestPayload::TypeScriptFileExecutionRequest(
662 wire::TypeScriptFileExecutionRequest {
663 process: background_process(&options, operation_id),
664 path: path.into(),
665 tsconfig_path: None,
666 compiler_options: None,
667 },
668 ),
669 true,
670 )
671 .await?,
672 "javascript",
673 )
674 }
675
676 pub async fn check_typescript(
677 &self,
678 source: impl Into<String>,
679 options: TypeScriptCheckOptions,
680 ) -> ClientResult<TypeScriptCheckResult> {
681 let compiler_options = options
682 .compiler_options
683 .as_ref()
684 .map(serde_json::to_string)
685 .transpose()
686 .map_err(|error| ClientError::Sidecar(error.to_string()))?;
687 let submission = self
688 .submit_execution(
689 wire::RequestPayload::TypeScriptCheckRequest(wire::TypeScriptCheckRequest {
690 identity: wire::ExecutionIdentityOptions {
691 context_id: options.context_id,
692 },
693 output: output(options.output),
694 source: source.into(),
695 cwd: options.cwd,
696 file_path: options.file_path,
697 tsconfig_path: options.tsconfig_path,
698 compiler_options,
699 timeout_ms: options.timeout_ms,
700 }),
701 false,
702 )
703 .await?;
704 let result = completed_submission(submission)?;
705 typescript_check_result(result)
706 }
707
708 pub async fn check_typescript_project(
709 &self,
710 options: TypeScriptCheckOptions,
711 ) -> ClientResult<TypeScriptCheckResult> {
712 let submission = self
713 .submit_execution(
714 wire::RequestPayload::TypeScriptProjectCheckRequest(
715 wire::TypeScriptProjectCheckRequest {
716 identity: wire::ExecutionIdentityOptions {
717 context_id: options.context_id,
718 },
719 output: output(options.output),
720 cwd: options.cwd,
721 tsconfig_path: options.tsconfig_path,
722 timeout_ms: options.timeout_ms,
723 },
724 ),
725 false,
726 )
727 .await?;
728 let result = completed_submission(submission)?;
729 typescript_check_result(result)
730 }
731
732 pub async fn install_npm_project(
733 &self,
734 options: NpmProjectInstallOptions,
735 ) -> ClientResult<CodeExecutionResult> {
736 completed_submission(
737 self.submit_execution(
738 wire::RequestPayload::NpmProjectInstallRequest(wire::NpmProjectInstallRequest {
739 identity: wire::ExecutionIdentityOptions {
740 context_id: options.context_id,
741 },
742 output: output(options.output),
743 cwd: options.cwd,
744 env: (!options.env.is_empty()).then_some(options.env),
745 timeout_ms: options.timeout_ms,
746 frozen: options.frozen,
747 }),
748 false,
749 )
750 .await?,
751 )
752 }
753
754 pub async fn install_npm_packages(
755 &self,
756 packages: Vec<String>,
757 options: NpmPackageInstallOptions,
758 ) -> ClientResult<CodeExecutionResult> {
759 completed_submission(
760 self.submit_execution(
761 wire::RequestPayload::NpmPackageInstallRequest(wire::NpmPackageInstallRequest {
762 identity: wire::ExecutionIdentityOptions {
763 context_id: options.context_id,
764 },
765 output: output(options.output),
766 cwd: options.cwd,
767 env: (!options.env.is_empty()).then_some(options.env),
768 timeout_ms: options.timeout_ms,
769 packages,
770 dev: options.dev,
771 global: options.global,
772 }),
773 false,
774 )
775 .await?,
776 )
777 }
778
779 pub async fn execute_npm_script(
780 &self,
781 script: impl Into<String>,
782 options: LanguageExecutionOptions,
783 ) -> ClientResult<CodeExecutionResult> {
784 completed_submission(
785 self.submit_execution(
786 wire::RequestPayload::NpmScriptExecutionRequest(wire::NpmScriptExecutionRequest {
787 process: process(&options),
788 script: script.into(),
789 }),
790 false,
791 )
792 .await?,
793 )
794 }
795
796 pub async fn execute_npm_package(
797 &self,
798 package_spec: impl Into<String>,
799 binary: Option<String>,
800 options: LanguageExecutionOptions,
801 ) -> ClientResult<CodeExecutionResult> {
802 completed_submission(
803 self.submit_execution(
804 wire::RequestPayload::NpmPackageExecutionRequest(
805 wire::NpmPackageExecutionRequest {
806 process: process(&options),
807 package_spec: package_spec.into(),
808 binary,
809 },
810 ),
811 false,
812 )
813 .await?,
814 )
815 }
816
817 pub async fn execute_python(
818 &self,
819 source: impl Into<String>,
820 options: InlineExecutionOptions,
821 ) -> ClientResult<CodeExecutionResult> {
822 let inputs = json_inputs(&options)?;
823 completed_submission(
824 self.submit_execution(
825 wire::RequestPayload::PythonExecutionRequest(wire::PythonExecutionRequest {
826 process: process(&options.process),
827 source: source.into(),
828 inputs,
829 }),
830 false,
831 )
832 .await?,
833 )
834 }
835
836 pub async fn evaluate_python(
837 &self,
838 expression: impl Into<String>,
839 options: InlineExecutionOptions,
840 ) -> ClientResult<CodeEvaluationResult> {
841 let inputs = json_inputs(&options)?;
842 let submission = self
843 .submit_execution(
844 wire::RequestPayload::PythonEvaluationRequest(wire::PythonEvaluationRequest {
845 process: process(&options.process),
846 expression: expression.into(),
847 inputs,
848 }),
849 false,
850 )
851 .await?;
852 evaluation_result(submission)
853 }
854
855 pub async fn execute_python_file(
856 &self,
857 path: impl Into<String>,
858 options: LanguageExecutionOptions,
859 ) -> ClientResult<CodeExecutionResult> {
860 completed_submission(
861 self.submit_execution(
862 wire::RequestPayload::PythonFileExecutionRequest(
863 wire::PythonFileExecutionRequest {
864 process: process(&options),
865 path: path.into(),
866 },
867 ),
868 false,
869 )
870 .await?,
871 )
872 }
873
874 pub async fn execute_python_module(
875 &self,
876 module: impl Into<String>,
877 options: LanguageExecutionOptions,
878 ) -> ClientResult<CodeExecutionResult> {
879 completed_submission(
880 self.submit_execution(
881 wire::RequestPayload::PythonModuleExecutionRequest(
882 wire::PythonModuleExecutionRequest {
883 process: process(&options),
884 module: module.into(),
885 },
886 ),
887 false,
888 )
889 .await?,
890 )
891 }
892
893 pub async fn spawn_python(
894 &self,
895 source: impl Into<String>,
896 options: LanguageSpawnOptions,
897 ) -> ClientResult<ProcessDescriptor> {
898 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
899 background_submission(
900 self,
901 self.submit_execution(
902 wire::RequestPayload::PythonExecutionRequest(wire::PythonExecutionRequest {
903 process: background_process(&options, operation_id),
904 source: source.into(),
905 inputs: None,
906 }),
907 true,
908 )
909 .await?,
910 "python",
911 )
912 }
913
914 pub async fn spawn_python_file(
915 &self,
916 path: impl Into<String>,
917 options: LanguageSpawnOptions,
918 ) -> ClientResult<ProcessDescriptor> {
919 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
920 background_submission(
921 self,
922 self.submit_execution(
923 wire::RequestPayload::PythonFileExecutionRequest(
924 wire::PythonFileExecutionRequest {
925 process: background_process(&options, operation_id),
926 path: path.into(),
927 },
928 ),
929 true,
930 )
931 .await?,
932 "python",
933 )
934 }
935
936 pub async fn spawn_python_module(
937 &self,
938 module: impl Into<String>,
939 options: LanguageSpawnOptions,
940 ) -> ClientResult<ProcessDescriptor> {
941 let operation_id = format!("process-{}", uuid::Uuid::new_v4());
942 background_submission(
943 self,
944 self.submit_execution(
945 wire::RequestPayload::PythonModuleExecutionRequest(
946 wire::PythonModuleExecutionRequest {
947 process: background_process(&options, operation_id),
948 module: module.into(),
949 },
950 ),
951 true,
952 )
953 .await?,
954 "python",
955 )
956 }
957
958 pub async fn install_python_packages(
959 &self,
960 packages: Vec<String>,
961 options: PythonInstallOptions,
962 ) -> ClientResult<CodeExecutionResult> {
963 if !packages.is_empty() && options.requirements_file.is_some() {
964 return Err(ClientError::Sidecar(String::from(
965 "install_python_packages cannot combine packages with requirements_file",
966 )));
967 }
968 completed_submission(
969 self.submit_execution(
970 wire::RequestPayload::PythonInstallRequest(wire::PythonInstallRequest {
971 identity: wire::ExecutionIdentityOptions {
972 context_id: options.context_id,
973 },
974 output: output(options.output),
975 cwd: options.cwd,
976 env: (!options.env.is_empty()).then_some(options.env),
977 timeout_ms: options.timeout_ms,
978 packages,
979 upgrade: options.upgrade,
980 requirements_file: options.requirements_file,
981 index_url: options.index_url,
982 extra_index_urls: options.extra_index_urls,
983 }),
984 false,
985 )
986 .await?,
987 )
988 }
989
990 pub async fn create_context(&self, context_id: &str) -> ClientResult<()> {
991 match self
992 .transport()
993 .request_wire(
994 self.execution_ownership(),
995 wire::RequestPayload::CreateContextRequest(wire::CreateContextRequest {
996 context_id: context_id.to_owned(),
997 }),
998 )
999 .await?
1000 {
1001 wire::ResponsePayload::ExecutionDescriptorResponse(_) => Ok(()),
1002 wire::ResponsePayload::RejectedResponse(rejected) => {
1003 Err(ClientError::from_rejection(rejected))
1004 }
1005 response => Err(ClientError::Sidecar(format!(
1006 "unexpected create_context response: {response:?}"
1007 ))),
1008 }
1009 }
1010
1011 pub async fn get_context(&self, context_id: &str) -> ClientResult<ContextDescriptor> {
1012 match self
1013 .transport()
1014 .request_wire(
1015 self.execution_ownership(),
1016 wire::RequestPayload::GetExecutionRequest(wire::GetExecutionRequest {
1017 execution_id: context_id.to_owned(),
1018 }),
1019 )
1020 .await?
1021 {
1022 wire::ResponsePayload::ExecutionDescriptorResponse(response) => {
1023 Ok(context_descriptor(response.execution))
1024 }
1025 wire::ResponsePayload::RejectedResponse(rejected) => {
1026 Err(ClientError::from_rejection(rejected))
1027 }
1028 response => Err(ClientError::Sidecar(format!(
1029 "unexpected get_context response: {response:?}"
1030 ))),
1031 }
1032 }
1033
1034 pub async fn list_contexts(&self) -> ClientResult<Vec<ContextDescriptor>> {
1035 match self
1036 .transport()
1037 .request_wire(
1038 self.execution_ownership(),
1039 wire::RequestPayload::ListExecutionsRequest,
1040 )
1041 .await?
1042 {
1043 wire::ResponsePayload::ExecutionListResponse(response) => Ok(response
1044 .executions
1045 .into_iter()
1046 .map(context_descriptor)
1047 .collect()),
1048 wire::ResponsePayload::RejectedResponse(rejected) => {
1049 Err(ClientError::from_rejection(rejected))
1050 }
1051 response => Err(ClientError::Sidecar(format!(
1052 "unexpected list_contexts response: {response:?}"
1053 ))),
1054 }
1055 }
1056
1057 async fn wait_execution(&self, execution_id: &str) -> ClientResult<CodeExecutionResult> {
1058 let mut events = self.transport().subscribe_wire_events();
1059 let first = self
1060 .transport()
1061 .request_wire(
1062 self.execution_ownership(),
1063 wire::RequestPayload::WaitExecutionRequest(wire::WaitExecutionRequest {
1064 execution_id: execution_id.to_owned(),
1065 }),
1066 )
1067 .await?;
1068 let response = match first {
1069 wire::ResponsePayload::RejectedResponse(rejected)
1070 if rejected.code == "execution_busy" =>
1071 {
1072 wait_for_completion_event(&mut events, execution_id).await?;
1073 self.transport()
1074 .request_wire(
1075 self.execution_ownership(),
1076 wire::RequestPayload::WaitExecutionRequest(wire::WaitExecutionRequest {
1077 execution_id: execution_id.to_owned(),
1078 }),
1079 )
1080 .await?
1081 }
1082 response => response,
1083 };
1084 match response {
1085 wire::ResponsePayload::ExecutionCompletedResponse(response) => Ok(response),
1086 wire::ResponsePayload::RejectedResponse(rejected) => {
1087 Err(ClientError::from_rejection(rejected))
1088 }
1089 response => Err(ClientError::Sidecar(format!(
1090 "unexpected wait_execution response: {response:?}"
1091 ))),
1092 }
1093 }
1094
1095 pub async fn reset_context(&self, context_id: &str) -> ClientResult<()> {
1096 self.execution_descriptor_request(wire::RequestPayload::ResetExecutionRequest(
1097 wire::ResetExecutionRequest {
1098 execution_id: context_id.to_owned(),
1099 },
1100 ))
1101 .await?;
1102 Ok(())
1103 }
1104
1105 async fn execution_descriptor_request(
1106 &self,
1107 payload: wire::RequestPayload,
1108 ) -> ClientResult<wire::ExecutionDescriptor> {
1109 match self
1110 .transport()
1111 .request_wire(self.execution_ownership(), payload)
1112 .await?
1113 {
1114 wire::ResponsePayload::ExecutionDescriptorResponse(response) => Ok(response.execution),
1115 wire::ResponsePayload::RejectedResponse(rejected) => {
1116 Err(ClientError::from_rejection(rejected))
1117 }
1118 response => Err(ClientError::Sidecar(format!(
1119 "unexpected execution lifecycle response: {response:?}"
1120 ))),
1121 }
1122 }
1123
1124 pub async fn delete_context(&self, context_id: &str) -> ClientResult<()> {
1125 match self
1126 .transport()
1127 .request_wire(
1128 self.execution_ownership(),
1129 wire::RequestPayload::DeleteExecutionRequest(wire::DeleteExecutionRequest {
1130 execution_id: context_id.to_owned(),
1131 }),
1132 )
1133 .await?
1134 {
1135 wire::ResponsePayload::ExecutionDeletedResponse(_) => Ok(()),
1136 wire::ResponsePayload::RejectedResponse(rejected) => {
1137 Err(ClientError::from_rejection(rejected))
1138 }
1139 response => Err(ClientError::Sidecar(format!(
1140 "unexpected delete_context response: {response:?}"
1141 ))),
1142 }
1143 }
1144}
1145
1146async fn wait_for_completion_event(
1147 events: &mut broadcast::Receiver<(wire::OwnershipScope, wire::EventPayload)>,
1148 execution_id: &str,
1149) -> ClientResult<()> {
1150 loop {
1151 match events.recv().await {
1152 Ok((_, wire::EventPayload::ExecutionCompletedEvent(event)))
1153 if event.execution_id == execution_id =>
1154 {
1155 return Ok(())
1156 }
1157 Ok(_) | Err(broadcast::error::RecvError::Lagged(_)) => {}
1158 Err(broadcast::error::RecvError::Closed) => {
1159 return Err(ClientError::Sidecar(String::from(
1160 "execution event stream closed before completion",
1161 )))
1162 }
1163 }
1164 }
1165}
1166
1167fn evaluation_result(submission: ExecutionSubmission) -> ClientResult<CodeEvaluationResult> {
1168 let ExecutionSubmission::Completed(result) = submission else {
1169 return Err(ClientError::Sidecar(String::from(
1170 "evaluation unexpectedly returned a background process",
1171 )));
1172 };
1173 let value = result
1174 .evaluation_value
1175 .as_deref()
1176 .map(serde_json::from_str::<serde_json::Value>)
1177 .transpose()
1178 .map_err(|error| {
1179 ClientError::Sidecar(format!("failed to decode evaluation result: {error}"))
1180 })?;
1181 Ok(CodeEvaluationResult { result, value })
1182}
1183
1184fn typescript_check_result(result: CodeExecutionResult) -> ClientResult<TypeScriptCheckResult> {
1185 if result.outcome != wire::ExecutionOutcome::Succeeded {
1186 return Ok(TypeScriptCheckResult {
1187 result,
1188 has_errors: None,
1189 diagnostics: Vec::new(),
1190 });
1191 }
1192 let data: serde_json::Value = result
1193 .type_script_check_result
1194 .as_deref()
1195 .map(serde_json::from_str)
1196 .transpose()
1197 .map_err(|error| {
1198 ClientError::Sidecar(format!(
1199 "failed to decode TypeScript checker result: {error}"
1200 ))
1201 })?
1202 .ok_or_else(|| {
1203 ClientError::Sidecar(String::from(
1204 "TypeScript checker returned no diagnostic result",
1205 ))
1206 })?;
1207 let data = data.as_object().ok_or_else(|| {
1208 ClientError::Sidecar(String::from(
1209 "TypeScript checker returned an invalid diagnostic result",
1210 ))
1211 })?;
1212 let has_errors = data
1213 .get("hasErrors")
1214 .and_then(serde_json::Value::as_bool)
1215 .ok_or_else(|| {
1216 ClientError::Sidecar(String::from(
1217 "TypeScript checker returned an invalid hasErrors value",
1218 ))
1219 })?;
1220 let diagnostics = data
1221 .get("diagnostics")
1222 .and_then(serde_json::Value::as_array)
1223 .ok_or_else(|| {
1224 ClientError::Sidecar(String::from(
1225 "TypeScript checker returned invalid diagnostics",
1226 ))
1227 })?
1228 .iter()
1229 .map(|diagnostic| {
1230 let code = diagnostic
1231 .get("code")
1232 .and_then(serde_json::Value::as_u64)
1233 .and_then(|code| u32::try_from(code).ok())
1234 .ok_or_else(|| {
1235 ClientError::Sidecar(String::from(
1236 "TypeScript checker returned an invalid diagnostic code",
1237 ))
1238 })?;
1239 let category = diagnostic
1240 .get("category")
1241 .and_then(serde_json::Value::as_str)
1242 .ok_or_else(|| {
1243 ClientError::Sidecar(String::from(
1244 "TypeScript checker returned an invalid diagnostic category",
1245 ))
1246 })?
1247 .to_owned();
1248 let message = diagnostic
1249 .get("message")
1250 .and_then(serde_json::Value::as_str)
1251 .ok_or_else(|| {
1252 ClientError::Sidecar(String::from(
1253 "TypeScript checker returned an invalid diagnostic message",
1254 ))
1255 })?
1256 .to_owned();
1257 Ok(wire::TypeScriptDiagnostic {
1258 code,
1259 category,
1260 message,
1261 file_path: diagnostic
1262 .get("filePath")
1263 .and_then(serde_json::Value::as_str)
1264 .map(str::to_owned),
1265 line: diagnostic
1266 .get("line")
1267 .and_then(serde_json::Value::as_u64)
1268 .and_then(|line| u32::try_from(line).ok()),
1269 column: diagnostic
1270 .get("column")
1271 .and_then(serde_json::Value::as_u64)
1272 .and_then(|column| u32::try_from(column).ok()),
1273 })
1274 })
1275 .collect::<ClientResult<Vec<_>>>()?;
1276 Ok(TypeScriptCheckResult {
1277 result,
1278 has_errors: Some(has_errors),
1279 diagnostics,
1280 })
1281}
1282
1283fn completed_submission(submission: ExecutionSubmission) -> ClientResult<CodeExecutionResult> {
1284 match submission {
1285 ExecutionSubmission::Completed(result) => Ok(result),
1286 ExecutionSubmission::Background(_) => Err(ClientError::Sidecar(String::from(
1287 "attached operation unexpectedly returned a background process",
1288 ))),
1289 }
1290}
1291
1292fn background_submission(
1293 client: &AgentOs,
1294 submission: ExecutionSubmission,
1295 language: &str,
1296) -> ClientResult<ProcessDescriptor> {
1297 match submission {
1298 ExecutionSubmission::Background(descriptor) => {
1299 let pid = descriptor.pid.ok_or_else(|| {
1300 ClientError::Sidecar(String::from(
1301 "spawned process admission returned no numeric pid",
1302 ))
1303 })?;
1304 let process_id = descriptor.process_id.clone().ok_or_else(|| {
1305 ClientError::Sidecar(String::from(
1306 "spawned process admission returned no process routing id",
1307 ))
1308 })?;
1309 let (stdout_tx, _) = broadcast::channel::<Vec<u8>>(1024);
1310 let (stderr_tx, _) = broadcast::channel::<Vec<u8>>(1024);
1311 let (output_tx, _) = broadcast::channel::<ProcessOutput>(1024);
1312 let (exit_tx, _) = watch::channel::<Option<i32>>(None);
1313 let (kernel_pid_tx, _) = watch::channel(Some(pid));
1314 let entry = ProcessEntry {
1315 command: format!("{language} source"),
1316 args: Vec::new(),
1317 stdout_tx: stdout_tx.clone(),
1318 stderr_tx: stderr_tx.clone(),
1319 output_tx: output_tx.clone(),
1320 exit_tx: exit_tx.clone(),
1321 process_id: process_id.clone(),
1322 kernel_pid: kernel_pid_tx,
1323 output_tasks: Vec::new(),
1324 started_at: descriptor.created_at_ms as i64,
1325 };
1326 let _ = client.inner().processes.insert(pid, entry);
1327 let mut events = client.transport().subscribe_wire_events();
1328 let operation_id = descriptor.execution_id.clone();
1329 tokio::spawn(async move {
1330 loop {
1331 let Ok((_, event)) = events.recv().await else {
1332 break;
1333 };
1334 match event {
1335 wire::EventPayload::ExecutionOutputEvent(output)
1336 if output.execution_id == operation_id =>
1337 {
1338 let (stream, tx) = match output.channel {
1339 wire::ExecutionStreamChannel::Stdout => {
1340 (ProcessStream::Stdout, &stdout_tx)
1341 }
1342 wire::ExecutionStreamChannel::Stderr
1343 | wire::ExecutionStreamChannel::Pty => {
1344 (ProcessStream::Stderr, &stderr_tx)
1345 }
1346 };
1347 let _ = tx.send(output.chunk.clone());
1348 let _ = output_tx.send(ProcessOutput {
1349 pid,
1350 stream,
1351 data: output.chunk,
1352 });
1353 }
1354 wire::EventPayload::ExecutionCompletedEvent(completed)
1355 if completed.execution_id == operation_id =>
1356 {
1357 let _ = exit_tx.send(Some(completed.exit_code.unwrap_or(1)));
1358 break;
1359 }
1360 _ => {}
1361 }
1362 }
1363 });
1364 Ok(ProcessDescriptor {
1365 pid,
1366 state: String::from("running"),
1367 language: Some(language.to_owned()),
1368 started_at_ms: descriptor.created_at_ms,
1369 })
1370 }
1371 ExecutionSubmission::Completed(_) => Err(ClientError::Sidecar(String::from(
1372 "spawned process unexpectedly returned an attached result",
1373 ))),
1374 }
1375}