1use std::collections::BTreeMap;
2use std::future::Future;
3use std::pin::Pin;
4use std::sync::Arc;
5use std::sync::atomic::{AtomicU64, Ordering};
6
7use lash_core::ToolChildExecutionTraceHook;
8use lash_trace::{
9 TraceBranchSelection, TraceContext, TraceEvent, TraceLabelMetadata,
10 TraceLashlangChildExecution, TraceLashlangExecutionEvent, TraceLashlangExecutionIdentity,
11 TraceLashlangMap, TraceLashlangMapEdge, TraceLashlangMapNode, TraceLashlangStatus, TraceRecord,
12 TraceRuntimeScope, TraceRuntimeSubject, TraceSink,
13};
14use lashlang::{ExecutionHost, ExecutionHostError};
15use sha2::{Digest, Sha256};
16use tokio_util::sync::CancellationToken;
17
18use crate::{
19 LASHLANG_ENGINE_KIND, LashlangProcessEngine, LashlangProcessInput,
20 bridge::{
21 lashlang_value_to_json, process_event_payload, protocol_tool_reply_to_lashlang_value,
22 sleep_duration_ms,
23 },
24 lashlang_host_environment_satisfies_requirements, prepare_lashlang_process_start,
25 resolve_lashlang_module_operation,
26};
27
28static SEGMENT_BOUNDARY_DECLINED_TOTAL: AtomicU64 = AtomicU64::new(0);
29
30fn record_segment_boundary_decline(error: &dyn std::fmt::Display, message: &'static str) {
31 let declined_total = SEGMENT_BOUNDARY_DECLINED_TOTAL
32 .fetch_add(1, Ordering::Relaxed)
33 .saturating_add(1);
34 tracing::warn!(error = %error, declined_total, "{message}");
35}
36
37fn trace_lifecycle_for_segment(
38 is_initial_segment: bool,
39 output: &lash_core::ProcessRunOutcome,
40) -> (bool, bool) {
41 (
42 is_initial_segment,
43 matches!(output, lash_core::ProcessRunOutcome::Terminal(_)),
44 )
45}
46
47#[derive(serde::Serialize, serde::Deserialize)]
48struct LashlangSegmentState {
49 vm: lashlang::VmContinuation,
50 sleep_sequence: u64,
51 event_sequence: u64,
52 signal_send_sequence: u64,
53 signal_wait_ordinals: BTreeMap<String, u64>,
54}
55
56fn lashlang_program_hash(input: &LashlangProcessInput) -> String {
57 let identity = serde_json::to_vec(&(
58 "lashlang-bytecode",
59 lashlang::BYTECODE_FORMAT_VERSION,
60 &input.module_ref,
61 &input.process_ref,
62 &input.host_requirements_ref,
63 &input.process_name,
64 ))
65 .expect("lashlang program identity should serialize");
66 format!("sha256:{:x}", Sha256::digest(identity))
67}
68
69fn validate_lashlang_program_hash(
70 persisted: Option<&str>,
71 current: &str,
72) -> Result<(), Box<lash_core::ProcessAwaitOutput>> {
73 if let Some(persisted) = persisted
74 && persisted != current
75 {
76 return Err(Box::new(process_lashlang_failure(
77 "restate_segment_program_hash_mismatch",
78 format!(
79 "lashlang segment program identity mismatch: persisted {persisted}, current {current}"
80 ),
81 None,
82 )));
83 }
84 Ok(())
85}
86
87pub async fn run_lashlang_process(
88 engine: LashlangProcessEngine,
89 mut context: lash_core::ProcessEngineRunContext<'_>,
90 payload: serde_json::Value,
91) -> lash_core::ProcessRunOutcome {
92 let handover = context.take_handover();
93 let is_initial_segment = handover.is_none();
94 let persisted_program_hash = handover
95 .as_ref()
96 .and_then(|handover| handover.program_hash.clone());
97 let segment_controller = context.scoped_effect_controller();
98 let phase_probe = context.turn_phase_probe();
99 let input = match LashlangProcessInput::from_payload(payload) {
100 Ok(input) => input,
101 Err(err) => {
102 return process_lashlang_failure(
103 "process_payload_invalid",
104 format!("invalid lashlang process payload: {err}"),
105 None,
106 )
107 .into();
108 }
109 };
110 let artifact = {
111 let _phase = context.named_phase("rlm_process.load_artifact");
112 match engine
113 .artifact_store
114 .get_module_artifact(&input.module_ref)
115 .await
116 {
117 Ok(Some(artifact)) => artifact,
118 Ok(None) => {
119 return process_lashlang_failure(
120 "process_module_artifact_missing",
121 format!("missing lashlang module artifact `{}`", input.module_ref),
122 None,
123 )
124 .into();
125 }
126 Err(err) => {
127 return process_lashlang_failure(
128 "process_module_artifact_load_failed",
129 format!(
130 "failed to load lashlang module artifact `{}`: {err}",
131 input.module_ref
132 ),
133 None,
134 )
135 .into();
136 }
137 }
138 };
139 if artifact.host_requirements_ref != input.host_requirements_ref {
140 return process_lashlang_failure(
141 "process_host_requirements_mismatch",
142 format!(
143 "lashlang process `{}` requested surface {}, artifact has {}",
144 input.process_name, input.host_requirements_ref, artifact.host_requirements_ref
145 ),
146 None,
147 )
148 .into();
149 }
150 if artifact.process_ref(&input.process_name) != Some(&input.process_ref) {
151 return process_lashlang_failure(
152 "process_ref_mismatch",
153 format!(
154 "lashlang module `{}` does not export process `{}` as requested ref {:?}",
155 input.module_ref, input.process_name, input.process_ref
156 ),
157 None,
158 )
159 .into();
160 }
161 let (tool_catalog, host_environment) = {
162 let _phase = context.named_phase("rlm_process.resolve_environment");
163 let tool_catalog = match context.resolved_tool_catalog() {
164 Ok(tool_catalog) => tool_catalog,
165 Err(err) => {
166 return process_lashlang_failure(
167 "process_tool_catalog_failed",
168 err.to_string(),
169 None,
170 )
171 .into();
172 }
173 };
174 let surface = engine
175 .surface
176 .clone()
177 .for_process_registry(context.process_registry_available());
178 let host_environment = match surface.host_environment(&tool_catalog) {
179 Ok(host_environment) => host_environment,
180 Err(err) => {
181 return process_lashlang_failure("process_host_environment_invalid", err, None)
182 .into();
183 }
184 };
185 if let Err(err) = lashlang_host_environment_satisfies_requirements(
186 &artifact.host_requirements,
187 &host_environment,
188 ) {
189 return process_lashlang_failure(
190 "process_host_environment_incompatible",
191 format!(
192 "lashlang process `{}` is incompatible with this host surface: {err}",
193 input.process_name
194 ),
195 None,
196 )
197 .into();
198 }
199 (tool_catalog, host_environment)
200 };
201 let compiled = {
202 let _phase = context.named_phase("rlm_process.compile");
203 let compiled = match engine.process_cache.lock() {
204 Ok(mut cache) => {
205 cache.get_or_compile(&artifact, &input.process_ref, &input.host_requirements_ref)
206 }
207 Err(_) => Err(lashlang::RuntimeError::ValueError {
208 message: "lashlang compiled process cache lock poisoned".to_string(),
209 }),
210 };
211 match compiled {
212 Ok(compiled) => compiled,
213 Err(err) => {
214 return process_lashlang_failure(
215 "process_compile_failed",
216 format!("failed to compile process `{}`: {err}", input.process_name),
217 None,
218 )
219 .into();
220 }
221 }
222 };
223 let segment_state: Option<LashlangSegmentState> = match handover {
224 Some(handover) => match serde_json::from_slice(&handover.engine_state) {
225 Ok(state) => Some(state),
226 Err(err) => {
227 return process_lashlang_failure(
228 "process_segment_handover_invalid",
229 format!("invalid lashlang segment handover: {err}"),
230 None,
231 )
232 .into();
233 }
234 },
235 None => None,
236 };
237 let current_program_hash = lashlang_program_hash(&input);
238 if let Err(output) =
239 validate_lashlang_program_hash(persisted_program_hash.as_deref(), ¤t_program_hash)
240 {
241 return (*output).into();
242 }
243 let process_id = context.registration().id.clone();
244 let session_id = context.session_id().to_string();
245 let lashlang_execution_trace = LashlangProcessExecutionTrace::new(
246 engine.execution_sink.clone(),
247 engine.trace_context.clone(),
248 session_id,
249 process_id.clone(),
250 artifact.module_ref.clone(),
251 input.process_ref.clone(),
252 input.process_name.clone(),
253 );
254 if is_initial_segment {
255 lashlang_execution_trace.emit_started(&artifact);
256 }
257 let registry = context.registry();
258 let cancellation = context.cancellation_token();
259 let (ctx, guard, mut state) = {
260 let _phase = context.named_phase("rlm_process.build_context");
261 let runtime_context = match context.into_runtime_context(tool_catalog) {
262 Ok(runtime_context) => runtime_context,
263 Err(err) => {
264 return process_lashlang_failure(
265 "process_run_context_failed",
266 err.to_string(),
267 None,
268 )
269 .into();
270 }
271 };
272 let (ctx, guard) = runtime_context.into_parts();
273 let mut globals = lashlang::Record::with_capacity(input.args.len());
274 for (name, value) in input.args {
275 globals.insert(name, lashlang::from_json(value));
276 }
277 let state = lashlang::State::from_snapshot(lashlang::Snapshot { globals });
278 (ctx, guard, state)
279 };
280 let sleep_sequence = segment_state
281 .as_ref()
282 .map_or(0, |state| state.sleep_sequence);
283 let event_sequence = segment_state
284 .as_ref()
285 .map_or(0, |state| state.event_sequence);
286 let signal_send_sequence = segment_state
287 .as_ref()
288 .map_or(0, |state| state.signal_send_sequence);
289 let signal_wait_ordinals = segment_state
290 .as_ref()
291 .map_or_else(BTreeMap::new, |state| state.signal_wait_ordinals.clone());
292 let host = LashlangProcessHost {
293 ctx,
294 host_environment,
295 artifact_store: engine.artifact_store(),
296 registry,
297 process_id: process_id.clone(),
298 lashlang_execution_trace: lashlang_execution_trace.clone(),
299 sleep_sequence: AtomicU64::new(sleep_sequence),
300 event_sequence: AtomicU64::new(event_sequence),
301 signal_send_sequence: AtomicU64::new(signal_send_sequence),
302 signal_wait_ordinals: tokio::sync::Mutex::new(signal_wait_ordinals),
303 };
304 let env = lashlang::ExecutionEnvironment::new(&host).process();
305 let output = {
306 let _phase = host.ctx.named_phase("rlm_process.execute");
307 execute_lashlang(
308 compiled,
309 &mut state,
310 &env,
311 cancellation.clone(),
312 segment_controller.controller(),
313 &host,
314 (segment_state, current_program_hash),
315 )
316 .await
317 };
318 drop(env);
319 drop(host);
320 {
321 let _phase =
322 lash_core::runtime::RuntimeNamedPhase::begin(phase_probe, "rlm_process.shutdown");
323 guard.shutdown().await;
324 }
325 if trace_lifecycle_for_segment(is_initial_segment, &output).1
326 && let lash_core::ProcessRunOutcome::Terminal(output) = &output
327 {
328 lashlang_execution_trace.emit_finished(output);
329 }
330 output
331}
332
333async fn execute_lashlang(
334 compiled: Arc<lashlang::CompiledProgram>,
335 state: &mut lashlang::State,
336 env: &lashlang::ExecutionEnvironment<'_, LashlangProcessHost<'_>>,
337 cancellation: CancellationToken,
338 controller: &dyn lash_core::RuntimeEffectController,
339 host: &LashlangProcessHost<'_>,
340 segment: (Option<LashlangSegmentState>, String),
341) -> lash_core::ProcessRunOutcome {
342 let (segment_state, program_hash) = segment;
343 let mut vm = if let Some(segment_state) = segment_state {
344 match lashlang::Vm::resume_from(segment_state.vm, compiled.as_ref(), env) {
345 Ok(vm) => vm,
346 Err(err) => {
347 return process_lashlang_failure(
348 "process_segment_resume_failed",
349 format!("failed to resume lashlang segment: {err}"),
350 None,
351 )
352 .into();
353 }
354 }
355 } else {
356 lashlang::Vm::from_state(compiled.as_ref(), state, env)
357 };
358 let mut progress = lash_core::SegmentProgress::default();
359 loop {
360 let execution = if env.trace_runtime_errors() {
361 tokio::select! {
362 _ = cancellation.cancelled() => {
363 return process_lashlang_cancelled("lashlang process was cancelled").into();
364 }
365 result = vm.run_process_traced_until_effect() => {
366 result.map_err(|failure| {
367 let error = failure.error.clone();
368 env.observe_runtime_failure(failure);
369 error
370 })
371 }
372 }
373 } else {
374 tokio::select! {
375 _ = cancellation.cancelled() => {
376 return process_lashlang_cancelled("lashlang process was cancelled").into();
377 }
378 result = vm.run_process_until_effect() => result,
379 }
380 };
381 if cancellation.is_cancelled() {
382 return process_lashlang_cancelled("lashlang process was cancelled").into();
383 }
384 match execution {
385 Ok(lashlang::VmRunOutcome::Complete(output)) => {
386 vm.flush_profile(compiled.as_ref(), env);
387 return process_lashlang_execution_result(Ok(output)).into();
388 }
389 Err(err) => {
390 vm.flush_profile(compiled.as_ref(), env);
391 return process_lashlang_execution_result(Err(err)).into();
392 }
393 Ok(lashlang::VmRunOutcome::EffectCompleted) => {
394 progress.effects_executed += 1;
395 let Some(reason) = controller.wants_segment_boundary(&progress) else {
396 continue;
397 };
398 match vm.suspend() {
399 Ok(continuation) => {
400 let segment_state = LashlangSegmentState {
401 vm: continuation,
402 sleep_sequence: host.sleep_sequence.load(Ordering::Relaxed),
403 event_sequence: host.event_sequence.load(Ordering::Relaxed),
404 signal_send_sequence: host.signal_send_sequence.load(Ordering::Relaxed),
405 signal_wait_ordinals: host.signal_wait_ordinals.lock().await.clone(),
406 };
407 match serde_json::to_vec(&segment_state) {
408 Ok(engine_state) => {
409 return lash_core::ProcessRunOutcome::SegmentBoundary(
410 lash_core::SegmentHandover {
411 reason,
412 program_hash: Some(program_hash.clone()),
413 engine_state,
414 },
415 );
416 }
417 Err(err) => {
418 record_segment_boundary_decline(
419 &err,
420 "lashlang segment continuation was not serializable; continuing",
421 );
422 }
423 }
424 }
425 Err(err) => {
426 record_segment_boundary_decline(
427 &err,
428 "lashlang segment boundary declined at non-capturable point",
429 );
430 }
431 }
432 }
433 }
434 }
435}
436
437struct LashlangProcessHost<'run> {
438 ctx: lash_core::RuntimeExecutionContext<'run>,
439 host_environment: lashlang::LashlangHostEnvironment,
440 artifact_store: Arc<dyn lashlang::LashlangArtifactStore>,
441 registry: Arc<dyn lash_core::ProcessRegistry>,
442 process_id: String,
443 lashlang_execution_trace: LashlangProcessExecutionTrace,
444 sleep_sequence: AtomicU64,
445 event_sequence: AtomicU64,
446 signal_send_sequence: AtomicU64,
447 signal_wait_ordinals: tokio::sync::Mutex<BTreeMap<String, u64>>,
448}
449
450type ProcessHostAbilityFuture<'a> =
451 Pin<Box<dyn Future<Output = Result<lashlang::AbilityResult, ExecutionHostError>> + Send + 'a>>;
452
453impl LashlangProcessHost<'_> {
454 fn resource_payload(
455 &self,
456 args: &[lashlang::Value],
457 ) -> Result<serde_json::Value, ExecutionHostError> {
458 let mut payload = if let [lashlang::Value::Record(record)] = args {
459 lashlang_value_to_json(&lashlang::Value::Record(Arc::clone(record)))?
460 } else {
461 serde_json::json!({
462 "args": args
463 .iter()
464 .map(lashlang_value_to_json)
465 .collect::<Result<Vec<_>, _>>()?,
466 })
467 };
468 payload
469 .as_object_mut()
470 .ok_or_else(|| ExecutionHostError::new("module operation payload must be an object"))?;
471 Ok(payload)
472 }
473
474 fn resource_tool_call_id(
475 &self,
476 host_operation: &str,
477 call_site: &lashlang::LashlangExecutionCallSite,
478 batch_index: Option<usize>,
479 ) -> String {
480 let mut call_id = format!(
481 "lashlang:{}:resource:{}:{}:{}",
482 self.process_id, host_operation, call_site.site.node_id, call_site.occurrence
483 );
484 if let Some(batch_index) = batch_index {
485 call_id.push_str(&format!(":child:{batch_index}"));
486 }
487 call_id
488 }
489
490 fn prepare_resource_invocation(
491 &self,
492 operation: String,
493 receiver: lashlang::Value,
494 args: Vec<lashlang::Value>,
495 call_site: Option<lashlang::LashlangExecutionCallSite>,
496 batch_index: Option<usize>,
497 ) -> Result<(String, lash_core::ToolInvocation), ExecutionHostError> {
498 let receiver = match &receiver {
499 lashlang::Value::Resource(receiver) => receiver,
500 _ => {
501 return Err(ExecutionHostError::new(format!(
502 "module operation `{operation}` requires a module authority receiver"
503 )));
504 }
505 };
506 let host_operation =
507 resolve_lashlang_module_operation(&self.host_environment, receiver, &operation)?;
508 let tool_id = lash_core::ToolId::from(host_operation.as_str());
509 let manifest = self.ctx.callable_tool_manifest_by_id(&tool_id).ok_or_else(|| {
510 ExecutionHostError::new(format!(
511 "module operation `{}` resolved to unavailable host operation `{host_operation}`",
512 operation
513 ))
514 })?;
515 let payload = self.resource_payload(&args)?;
516 let call_site = call_site.ok_or_else(|| {
517 ExecutionHostError::new(format!(
518 "module operation `{operation}` resolved to host operation `{host_operation}` but has no deterministic lashlang execution call site"
519 ))
520 })?;
521 let call_id = self.resource_tool_call_id(&host_operation, &call_site, batch_index);
522 let mut invocation = lash_core::ToolInvocation::new(call_id, manifest.id.clone(), payload);
523 if let Some(hook) = self
524 .lashlang_execution_trace
525 .tool_child_execution_trace_hook(call_site)
526 {
527 invocation = invocation.with_child_execution_trace_hook(hook);
528 }
529 Ok((host_operation, invocation))
530 }
531
532 async fn resource_operation(
533 &self,
534 operation: String,
535 receiver: lashlang::Value,
536 args: Vec<lashlang::Value>,
537 call_site: Option<lashlang::LashlangExecutionCallSite>,
538 ) -> Result<lashlang::Value, ExecutionHostError> {
539 let (_, invocation) =
540 self.prepare_resource_invocation(operation, receiver, args, call_site, None)?;
541 let lash_core::ToolInvocation {
542 id,
543 tool_id,
544 args,
545 execution_grant: _,
546 child_execution_trace_hook,
547 } = invocation;
548 let reply = if let Some(call_site) = child_execution_trace_hook {
549 self.ctx
550 .call_tool_by_id_with_child_execution_trace_hook(id, tool_id, args, 0, call_site)
551 .await
552 } else {
553 self.ctx.call_tool_by_id(id, tool_id, args, 0).await
554 };
555 protocol_tool_reply_to_lashlang_value(reply)
556 }
557
558 async fn resource_operation_batch(
559 &self,
560 batch: lashlang::ResourceOperationBatch,
561 ) -> lashlang::ResourceOperationBatchResult {
562 let mut results = vec![None; batch.operations.len()];
563 let mut positions = Vec::new();
564 let mut invocations = Vec::new();
565 for (index, operation) in batch.operations.into_iter().enumerate() {
566 match self.prepare_resource_invocation(
567 operation.operation,
568 operation.receiver,
569 operation.args,
570 operation.call_site,
571 Some(index),
572 ) {
573 Ok((_, invocation)) => {
574 positions.push(index);
575 invocations.push(invocation);
576 }
577 Err(error) => {
578 results[index] = Some(lashlang::ResourceOperationResult::Error(error));
579 }
580 }
581 }
582
583 for (index, reply) in positions
584 .into_iter()
585 .zip(self.ctx.call_tool_batch(invocations).await)
586 {
587 results[index] = Some(lashlang::ResourceOperationResult::from_result(
588 protocol_tool_reply_to_lashlang_value(reply),
589 ));
590 }
591
592 lashlang::ResourceOperationBatchResult {
593 results: results
594 .into_iter()
595 .map(|result| result.expect("every batch result slot should be filled"))
596 .collect(),
597 }
598 }
599
600 async fn await_handle(
601 &self,
602 handle: lashlang::Value,
603 ) -> Result<lashlang::Value, ExecutionHostError> {
604 let reply = {
605 let _phase = self.ctx.named_phase("rlm_process.await_handle");
606 self.ctx
607 .await_tool_handle(
608 uuid::Uuid::new_v4().to_string(),
609 lashlang_value_to_json(&handle)?,
610 )
611 .await
612 };
613 protocol_tool_reply_to_lashlang_value(reply)
614 }
615
616 async fn cancel_handle(
617 &self,
618 handle: lashlang::Value,
619 ) -> Result<lashlang::Value, ExecutionHostError> {
620 let reply = self
621 .ctx
622 .cancel_tool_handle(
623 uuid::Uuid::new_v4().to_string(),
624 lashlang_value_to_json(&handle)?,
625 )
626 .await;
627 protocol_tool_reply_to_lashlang_value(reply)
628 }
629
630 async fn start_process(
631 &self,
632 start: lashlang::ProcessStart,
633 ) -> Result<lashlang::Value, ExecutionHostError> {
634 let prepared = {
635 let _phase = self.ctx.named_phase("rlm_process.prepare_start");
636 let parent_start_seed = format!("parent-process:{}", self.process_id);
637 prepare_lashlang_process_start(
638 Arc::clone(&self.artifact_store),
639 &parent_start_seed,
640 start,
641 )
642 .await
643 .map_err(ExecutionHostError::new)?
644 };
645 let reply = {
646 let _phase = self.ctx.named_phase("rlm_process.start");
647 self.ctx
648 .start_child_process(prepared.registration, LASHLANG_ENGINE_KIND, prepared.label)
649 .await
650 };
651 protocol_tool_reply_to_lashlang_value(reply)
652 }
653
654 async fn process_event(&self, event: lashlang::ProcessEvent) -> Result<(), ExecutionHostError> {
655 let event_type = match event.kind {
656 lashlang::ProcessEventKind::Yield => "process.yield",
657 lashlang::ProcessEventKind::Wake => "process.wake",
658 };
659 let ordinal = self.event_sequence.fetch_add(1, Ordering::Relaxed);
660 self.ctx
661 .append_process_event(
662 Arc::clone(&self.registry),
663 &self.process_id,
664 lash_core::ProcessEventAppendRequest::new(
665 event_type,
666 process_event_payload(&event.value)?,
667 )
668 .with_replay_key(format!("process:{}:event:{ordinal}", self.process_id)),
669 )
670 .await
671 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
672 Ok(())
673 }
674
675 async fn sleep(&self, sleep: lashlang::Sleep) -> Result<lashlang::Value, ExecutionHostError> {
676 let duration_ms = sleep_duration_ms(sleep.kind, &sleep.value)?;
677 let sequence = self.sleep_sequence.fetch_add(1, Ordering::Relaxed);
678 let scope = format!("process:{}", self.process_id);
679 self.ctx
680 .sleep_process(&scope, sequence, duration_ms)
681 .await
682 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
683 Ok(lashlang::Value::Null)
684 }
685
686 async fn wait_signal(&self, name: String) -> Result<lashlang::Value, ExecutionHostError> {
687 let event_type = lash_core::process_signal_event_type(&name)
688 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
689 let event_ordinal = {
690 let mut ordinals = self.signal_wait_ordinals.lock().await;
691 let ordinal = ordinals.entry(name.clone()).or_insert(0);
692 *ordinal += 1;
693 *ordinal
694 };
695 let key = lash_core::process_signal_wait_key(&self.process_id, &name, event_ordinal);
696 let waiting_replay_key = format!(
697 "process:{}:waiting:signal.{}:{event_ordinal}",
698 self.process_id, name
699 );
700 let since_ms = self
701 .wait_since_ms(&key, &waiting_replay_key)
702 .await
703 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
704 let wait = lash_core::WaitState {
705 since_ms,
706 kind: lash_core::WaitKind::Signal {
707 name: name.clone(),
708 event_type: event_type.clone(),
709 key: key.clone(),
710 ordinal: event_ordinal,
711 },
712 };
713 self.registry
714 .set_process_wait(&self.process_id, wait.clone())
715 .await
716 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
717 self.registry
718 .append_event(
719 &self.process_id,
720 lash_core::ProcessEventAppendRequest::new(
721 "process.waiting",
722 serde_json::json!({ "wait": wait }),
723 )
724 .with_replay_key(waiting_replay_key),
725 )
726 .await
727 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
728 let payload = self
729 .ctx
730 .await_process_signal_event(&self.process_id, &name, event_ordinal)
731 .await
732 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
733 self.registry
734 .clear_process_wait(&self.process_id)
735 .await
736 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
737 self.registry
738 .append_event(
739 &self.process_id,
740 lash_core::ProcessEventAppendRequest::new(
741 "process.resumed",
742 serde_json::json!({
743 "signal": name,
744 "key": key,
745 "ordinal": event_ordinal,
746 }),
747 )
748 .with_replay_key(format!(
749 "process:{}:resumed:signal.{}:{event_ordinal}",
750 self.process_id, name
751 )),
752 )
753 .await
754 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
755 Ok(lashlang::from_json(payload))
756 }
757
758 async fn wait_since_ms(
759 &self,
760 key: &str,
761 waiting_replay_key: &str,
762 ) -> Result<u64, lash_core::PluginError> {
763 if let Some(since_ms) =
764 self.registry
765 .get_process(&self.process_id)
766 .await
767 .and_then(|record| {
768 let wait = record.wait?;
769 match &wait.kind {
770 lash_core::WaitKind::Signal { key: wait_key, .. } if wait_key == key => {
771 Some(wait.since_ms)
772 }
773 _ => None,
774 }
775 })
776 {
777 return Ok(since_ms);
778 }
779
780 for event in self
781 .registry
782 .events_after(&self.process_id, 0)
783 .await?
784 .into_iter()
785 .rev()
786 {
787 if event.event_type != "process.waiting"
788 || event.invocation.replay_key() != Some(waiting_replay_key)
789 {
790 continue;
791 }
792 let Some(wait_value) = event.payload.get("wait") else {
793 continue;
794 };
795 if let Ok(wait) = serde_json::from_value::<lash_core::WaitState>(wait_value.clone()) {
796 return Ok(wait.since_ms);
797 }
798 }
799 Ok(lash_core::current_epoch_ms())
800 }
801
802 async fn signal_run(
803 &self,
804 signal: lashlang::ProcessSignal,
805 ) -> Result<lashlang::Value, ExecutionHostError> {
806 let target = process_id_from_lashlang_handle(&signal.run)?;
807 let payload = lashlang_value_to_json(&signal.payload)?;
808 let sequence = self.signal_send_sequence.fetch_add(1, Ordering::Relaxed);
809 let signal_id = format!(
810 "lashlang:{}:signal.{}:{sequence}",
811 self.process_id, signal.name
812 );
813 self.ctx
814 .signal_process_by_id(
815 Arc::clone(&self.registry),
816 &target,
817 &signal.name,
818 signal_id,
819 payload,
820 )
821 .await
822 .map_err(|err| ExecutionHostError::new(err.to_string()))?;
823 Ok(lashlang::Value::Null)
824 }
825
826 fn perform_selected_ability<'a>(
827 &'a self,
828 op: lashlang::AbilityOp,
829 ) -> ProcessHostAbilityFuture<'a> {
830 match op {
831 lashlang::AbilityOp::ResourceOperation(operation) => Box::pin(async move {
832 self.resource_operation(
833 operation.operation,
834 operation.receiver,
835 operation.args,
836 operation.call_site,
837 )
838 .await
839 .map(lashlang::AbilityResult::Value)
840 }),
841 lashlang::AbilityOp::ResourceOperationBatch(batch) => Box::pin(async move {
842 Ok(lashlang::AbilityResult::ResourceOperationBatch(
843 self.resource_operation_batch(batch).await,
844 ))
845 }),
846 lashlang::AbilityOp::Await(handle) => Box::pin(async move {
847 self.await_handle(handle)
848 .await
849 .map(lashlang::AbilityResult::Value)
850 }),
851 lashlang::AbilityOp::Cancel(handle) => Box::pin(async move {
852 self.cancel_handle(handle)
853 .await
854 .map(lashlang::AbilityResult::Value)
855 }),
856 lashlang::AbilityOp::StartProcess(start) => Box::pin(async move {
857 self.start_process(*start)
858 .await
859 .map(lashlang::AbilityResult::Value)
860 }),
861 lashlang::AbilityOp::ProcessEvent(event) => Box::pin(async move {
862 self.process_event(event).await?;
863 Ok(lashlang::AbilityResult::Unit)
864 }),
865 lashlang::AbilityOp::Sleep(sleep) => {
866 Box::pin(async move { self.sleep(sleep).await.map(lashlang::AbilityResult::Value) })
867 }
868 lashlang::AbilityOp::WaitSignal { name } => Box::pin(async move {
869 self.wait_signal(name)
870 .await
871 .map(lashlang::AbilityResult::Value)
872 }),
873 lashlang::AbilityOp::SignalRun(signal) => Box::pin(async move {
874 self.signal_run(signal)
875 .await
876 .map(lashlang::AbilityResult::Value)
877 }),
878 lashlang::AbilityOp::Print(_) => Box::pin(async {
879 Err(ExecutionHostError::new(
880 "`print` is not available inside lashlang process bodies",
881 ))
882 }),
883 lashlang::AbilityOp::Finish(value) | lashlang::AbilityOp::Fail(value) => {
884 Box::pin(async move { Ok(lashlang::AbilityResult::Value(value)) })
885 }
886 }
887 }
888}
889
890impl lashlang::ExecutionHost for LashlangProcessHost<'_> {
891 fn perform(
892 &self,
893 op: lashlang::AbilityOp,
894 ) -> impl Future<Output = Result<lashlang::AbilityResult, ExecutionHostError>> + Send {
895 self.perform_selected_ability(op)
896 }
897
898 fn observe_lashlang_execution(&self, observation: lashlang::LashlangExecutionObservation) {
899 self.lashlang_execution_trace.emit_observation(observation);
900 }
901}
902
903#[derive(Clone)]
904struct LashlangProcessExecutionTrace {
905 sink: Option<Arc<dyn TraceSink>>,
906 base_context: TraceContext,
907 session_id: String,
908 process_id: String,
909 module_ref: lashlang::ModuleRef,
910 process_ref: lashlang::ProcessRef,
911 process_name: String,
912}
913
914impl LashlangProcessExecutionTrace {
915 fn new(
916 sink: Option<Arc<dyn TraceSink>>,
917 base_context: TraceContext,
918 session_id: String,
919 process_id: String,
920 module_ref: lashlang::ModuleRef,
921 process_ref: lashlang::ProcessRef,
922 process_name: String,
923 ) -> Self {
924 Self {
925 sink,
926 base_context,
927 session_id,
928 process_id,
929 module_ref,
930 process_ref,
931 process_name,
932 }
933 }
934
935 fn identity(&self) -> TraceLashlangExecutionIdentity {
936 TraceLashlangExecutionIdentity {
937 scope: TraceRuntimeScope::new(self.session_id.clone()),
938 subject: TraceRuntimeSubject::Process {
939 process_id: self.process_id.clone(),
940 },
941 module_ref: self.module_ref.to_string(),
942 entry_kind: "process".to_string(),
943 entry_ref: Some(lashlang::process_ref_key(&self.process_ref)),
944 entry_name: self.process_name.clone(),
945 }
946 }
947
948 fn event_key(&self, suffix: impl std::fmt::Display) -> String {
949 format!("lashlang_execution:{}:{suffix}", self.process_id)
950 }
951
952 fn emit_started(&self, artifact: &lashlang::ModuleArtifact) {
953 self.emit(TraceLashlangExecutionEvent::ExecutionStarted {
954 event_key: self.event_key("started"),
955 identity: self.identity(),
956 execution_map: trace_lashlang_process_map(
957 artifact,
958 &self.process_ref,
959 &self.process_name,
960 ),
961 });
962 }
963
964 fn emit_finished(&self, output: &lash_core::ProcessAwaitOutput) {
965 let (status, error) = match output {
966 lash_core::ProcessAwaitOutput::Success { .. } => (TraceLashlangStatus::Completed, None),
967 lash_core::ProcessAwaitOutput::Failure { message, .. } => {
968 (TraceLashlangStatus::Failed, Some(message.clone()))
969 }
970 lash_core::ProcessAwaitOutput::Cancelled { message, .. } => {
971 (TraceLashlangStatus::Cancelled, Some(message.clone()))
972 }
973 lash_core::ProcessAwaitOutput::Abandoned { .. } => (
977 TraceLashlangStatus::Failed,
978 Some("process abandoned".to_string()),
979 ),
980 };
981 self.emit(TraceLashlangExecutionEvent::ExecutionFinished {
982 event_key: self.event_key("finished"),
983 identity: self.identity(),
984 status,
985 error,
986 });
987 }
988
989 fn emit_observation(&self, observation: lashlang::LashlangExecutionObservation) {
990 if self.sink.is_none() {
991 return;
992 }
993 let identity = self.identity();
994 let event = match observation {
995 lashlang::LashlangExecutionObservation::NodeStarted { site, occurrence } => {
996 TraceLashlangExecutionEvent::NodeStarted {
997 event_key: self
998 .event_key(format!("node:{}:{occurrence}:started", site.node_id)),
999 identity,
1000 node_id: site.node_id,
1001 node_kind: site.node_kind,
1002 label: site.label,
1003 occurrence,
1004 }
1005 }
1006 lashlang::LashlangExecutionObservation::NodeCompleted { site, occurrence } => {
1007 TraceLashlangExecutionEvent::NodeCompleted {
1008 event_key: self
1009 .event_key(format!("node:{}:{occurrence}:completed", site.node_id)),
1010 identity,
1011 node_id: site.node_id,
1012 node_kind: site.node_kind,
1013 label: site.label,
1014 occurrence,
1015 }
1016 }
1017 lashlang::LashlangExecutionObservation::NodeFailed {
1018 site,
1019 occurrence,
1020 error,
1021 } => TraceLashlangExecutionEvent::NodeFailed {
1022 event_key: self.event_key(format!("node:{}:{occurrence}:failed", site.node_id)),
1023 identity,
1024 node_id: site.node_id,
1025 node_kind: site.node_kind,
1026 label: site.label,
1027 occurrence,
1028 error,
1029 },
1030 lashlang::LashlangExecutionObservation::BranchSelected {
1031 site,
1032 occurrence,
1033 edge_id,
1034 selected,
1035 } => TraceLashlangExecutionEvent::BranchSelected {
1036 event_key: self
1037 .event_key(format!("branch:{}:{occurrence}:{edge_id}", site.node_id)),
1038 identity,
1039 node_id: site.node_id,
1040 occurrence,
1041 edge_id,
1042 selected: match selected {
1043 lashlang::ProcessBranchSelection::Then => TraceBranchSelection::Then,
1044 lashlang::ProcessBranchSelection::Else => TraceBranchSelection::Else,
1045 },
1046 },
1047 lashlang::LashlangExecutionObservation::ChildStarted {
1048 site,
1049 occurrence,
1050 child,
1051 } => TraceLashlangExecutionEvent::ChildStarted {
1052 event_key: self.event_key(format!(
1053 "child:{}:{occurrence}:{}",
1054 site.node_id, child.process_id
1055 )),
1056 identity,
1057 parent_node_id: site.node_id,
1058 occurrence,
1059 child: TraceLashlangChildExecution {
1060 scope: TraceRuntimeScope::new(self.session_id.clone()),
1061 subject: TraceRuntimeSubject::Process {
1062 process_id: child.process_id,
1063 },
1064 module_ref: Some(child.module_ref.to_string()),
1065 entry_ref: Some(lashlang::process_ref_key(&child.process_ref)),
1066 entry_name: Some(child.process_name),
1067 },
1068 },
1069 };
1070 self.emit(event);
1071 }
1072
1073 fn tool_child_execution_trace_hook(
1074 &self,
1075 call_site: lashlang::LashlangExecutionCallSite,
1076 ) -> Option<ToolChildExecutionTraceHook> {
1077 self.sink.as_ref()?;
1078 let trace = self.clone();
1079 let parent_node_id = call_site.site.node_id;
1080 let occurrence = call_site.occurrence;
1081 Some(ToolChildExecutionTraceHook::new(move |started| {
1082 let child = TraceLashlangChildExecution {
1083 scope: TraceRuntimeScope::new(trace.session_id.clone()),
1084 subject: TraceRuntimeSubject::Process {
1085 process_id: started.process_id,
1086 },
1087 module_ref: None,
1088 entry_ref: None,
1089 entry_name: started.child_entry_name,
1090 };
1091 let child_graph_key = child.graph_key();
1092 trace.emit(TraceLashlangExecutionEvent::ChildStarted {
1093 event_key: trace.event_key(format!(
1094 "child:{parent_node_id}:{occurrence}:{child_graph_key}"
1095 )),
1096 identity: trace.identity(),
1097 parent_node_id: parent_node_id.clone(),
1098 occurrence,
1099 child,
1100 });
1101 }))
1102 }
1103
1104 fn emit(&self, event: TraceLashlangExecutionEvent) {
1105 let Some(sink) = &self.sink else {
1106 return;
1107 };
1108 let mut context = self.base_context.clone();
1109 context.session_id = Some(self.session_id.clone());
1110 let _ = sink.append(&TraceRecord::new(
1111 context,
1112 TraceEvent::LashlangExecution { event },
1113 ));
1114 }
1115}
1116
1117fn trace_lashlang_process_map(
1118 artifact: &lashlang::ModuleArtifact,
1119 process_ref: &lashlang::ProcessRef,
1120 process_name: &str,
1121) -> TraceLashlangMap {
1122 let source = lashlang::canonical_program_source_with_requirements(
1123 &artifact.canonical_ir,
1124 &artifact.host_requirements,
1125 );
1126 let graph = source
1127 .ok()
1128 .and_then(|source| lashlang::workflow_graph_from_source(&source).ok());
1129 let Some(process) = graph.as_ref().and_then(|graph| graph.process(process_name)) else {
1130 return TraceLashlangMap {
1131 module_ref: artifact.module_ref.to_string(),
1132 entry_kind: "process".to_string(),
1133 entry_ref: Some(lashlang::process_ref_key(process_ref)),
1134 entry_name: process_name.to_string(),
1135 nodes: Vec::new(),
1136 edges: Vec::new(),
1137 };
1138 };
1139 let mut nodes = Vec::new();
1140 let mut edges = Vec::new();
1141 let mut primary_runtime_ids = BTreeMap::new();
1142 append_trace_workflow_subgraph(
1143 artifact,
1144 &process.body,
1145 &mut nodes,
1146 &mut edges,
1147 &mut primary_runtime_ids,
1148 );
1149 TraceLashlangMap {
1150 module_ref: artifact.module_ref.to_string(),
1151 entry_kind: "process".to_string(),
1152 entry_ref: Some(lashlang::process_ref_key(process_ref)),
1153 entry_name: process_name.to_string(),
1154 nodes,
1155 edges,
1156 }
1157}
1158
1159pub fn trace_lashlang_main_map(artifact: &lashlang::ModuleArtifact) -> TraceLashlangMap {
1161 let graph = lashlang::canonical_program_source_with_requirements(
1162 &artifact.canonical_ir,
1163 &artifact.host_requirements,
1164 )
1165 .ok()
1166 .and_then(|source| lashlang::workflow_graph_from_source(&source).ok());
1167 let Some(graph) = graph else {
1168 return TraceLashlangMap {
1169 module_ref: artifact.module_ref.to_string(),
1170 entry_kind: "main".to_string(),
1171 entry_ref: None,
1172 entry_name: "main".to_string(),
1173 nodes: Vec::new(),
1174 edges: Vec::new(),
1175 };
1176 };
1177 let mut nodes = Vec::new();
1178 let mut edges = Vec::new();
1179 let mut primary_runtime_ids = BTreeMap::new();
1180 append_trace_workflow_subgraph(
1181 artifact,
1182 &graph.main,
1183 &mut nodes,
1184 &mut edges,
1185 &mut primary_runtime_ids,
1186 );
1187 TraceLashlangMap {
1188 module_ref: artifact.module_ref.to_string(),
1189 entry_kind: "main".to_string(),
1190 entry_ref: None,
1191 entry_name: "main".to_string(),
1192 nodes,
1193 edges,
1194 }
1195}
1196
1197fn append_trace_workflow_subgraph(
1198 artifact: &lashlang::ModuleArtifact,
1199 graph: &lashlang::WorkflowSubgraph,
1200 nodes: &mut Vec<TraceLashlangMapNode>,
1201 edges: &mut Vec<TraceLashlangMapEdge>,
1202 primary_runtime_ids: &mut BTreeMap<String, String>,
1203) {
1204 for node in &graph.nodes {
1205 let label_metadata =
1206 (node.name_source == lashlang::WorkflowNodeNameSource::Label).then(|| {
1207 TraceLabelMetadata {
1208 title: node.name.clone(),
1209 description: node.description.clone(),
1210 }
1211 });
1212 for site in &node.execution_sites {
1213 let Some(runtime_site) =
1214 lashlang::runtime_execution_site_for_workflow_site(artifact, site)
1215 else {
1216 continue;
1217 };
1218 primary_runtime_ids
1219 .entry(node.id.to_string())
1220 .or_insert_with(|| runtime_site.node_id.clone());
1221 if nodes.iter().any(|node| node.id == runtime_site.node_id) {
1222 continue;
1223 }
1224 nodes.push(TraceLashlangMapNode {
1225 id: runtime_site.node_id,
1226 kind: runtime_site.node_kind,
1227 label: runtime_site.label,
1228 label_metadata: label_metadata.clone(),
1229 });
1230 }
1231 match &node.kind {
1232 lashlang::WorkflowNodeKind::Container(lashlang::WorkflowContainer::If {
1233 then_graph,
1234 else_graph,
1235 ..
1236 }) => {
1237 if let Some(graph) = then_graph {
1238 append_trace_workflow_subgraph(
1239 artifact,
1240 graph,
1241 nodes,
1242 edges,
1243 primary_runtime_ids,
1244 );
1245 }
1246 if let Some(graph) = else_graph {
1247 append_trace_workflow_subgraph(
1248 artifact,
1249 graph,
1250 nodes,
1251 edges,
1252 primary_runtime_ids,
1253 );
1254 }
1255 }
1256 lashlang::WorkflowNodeKind::Container(lashlang::WorkflowContainer::For {
1257 body: Some(graph),
1258 ..
1259 }) => {
1260 append_trace_workflow_subgraph(artifact, graph, nodes, edges, primary_runtime_ids);
1261 }
1262 lashlang::WorkflowNodeKind::Container(lashlang::WorkflowContainer::While {
1263 body: Some(graph),
1264 ..
1265 }) => {
1266 append_trace_workflow_subgraph(artifact, graph, nodes, edges, primary_runtime_ids);
1267 }
1268 lashlang::WorkflowNodeKind::Container(
1269 lashlang::WorkflowContainer::ListComprehension {
1270 element: Some(graph),
1271 ..
1272 },
1273 ) => {
1274 append_trace_workflow_subgraph(artifact, graph, nodes, edges, primary_runtime_ids);
1275 }
1276 _ => {}
1277 }
1278 }
1279 for edge in &graph.edges {
1280 let (Some(from), Some(to)) = (
1281 primary_runtime_ids.get(edge.from.as_str()),
1282 primary_runtime_ids.get(edge.to.as_str()),
1283 ) else {
1284 continue;
1285 };
1286 let label = match &edge.kind {
1287 lashlang::WorkflowEdgeKind::Sequence => "sequence".to_string(),
1288 lashlang::WorkflowEdgeKind::DataDependency { variable, version } => {
1289 format!("{variable}@{version}")
1290 }
1291 };
1292 edges.push(TraceLashlangMapEdge {
1293 id: edge.id.clone(),
1294 from: from.clone(),
1295 to: to.clone(),
1296 label,
1297 });
1298 }
1299}
1300
1301fn process_lashlang_execution_result(
1302 result: Result<lashlang::ExecutionOutcome, lashlang::RuntimeError>,
1303) -> lash_core::ProcessAwaitOutput {
1304 match result {
1305 Ok(lashlang::ExecutionOutcome::Finished(value)) => lash_core::ProcessAwaitOutput::Success {
1306 value: lashlang_value_to_json(&value)
1307 .unwrap_or_else(|err| serde_json::json!({ "error": err.to_string() })),
1308 control: None,
1309 },
1310 Ok(lashlang::ExecutionOutcome::Failed(value)) => process_lashlang_failure(
1311 "process_failed",
1312 value.to_string(),
1313 Some(
1314 lashlang_value_to_json(&value)
1315 .unwrap_or_else(|err| serde_json::json!({ "error": err.to_string() })),
1316 ),
1317 ),
1318 Ok(lashlang::ExecutionOutcome::Continued) => lash_core::ProcessAwaitOutput::Success {
1319 value: serde_json::Value::Null,
1320 control: None,
1321 },
1322 Err(err) => process_lashlang_failure("process_runtime_error", err.to_string(), None),
1323 }
1324}
1325
1326fn process_lashlang_failure(
1327 code: &str,
1328 message: impl Into<String>,
1329 raw: Option<serde_json::Value>,
1330) -> lash_core::ProcessAwaitOutput {
1331 lash_core::ProcessAwaitOutput::Failure {
1332 class: lash_core::ToolFailureClass::Execution,
1333 code: code.to_string(),
1334 message: message.into(),
1335 raw,
1336 control: None,
1337 }
1338}
1339
1340fn process_lashlang_cancelled(message: impl Into<String>) -> lash_core::ProcessAwaitOutput {
1341 lash_core::ProcessAwaitOutput::Cancelled {
1342 message: message.into(),
1343 raw: None,
1344 control: None,
1345 }
1346}
1347
1348fn process_id_from_lashlang_handle(handle: &lashlang::Value) -> Result<String, ExecutionHostError> {
1349 let value = lashlang_value_to_json(handle)?;
1350 let Some(object) = value.as_object() else {
1351 return Err(ExecutionHostError::new(
1352 "signal_run expects a process handle",
1353 ));
1354 };
1355 if object.get("__handle__").and_then(serde_json::Value::as_str) != Some("process") {
1356 return Err(ExecutionHostError::new(
1357 "signal_run expects a process handle",
1358 ));
1359 }
1360 object
1361 .get("id")
1362 .and_then(serde_json::Value::as_str)
1363 .map(ToOwned::to_owned)
1364 .ok_or_else(|| ExecutionHostError::new("signal_run process handle is missing `id`"))
1365}
1366
1367pub fn lashlang_process_event_types() -> Vec<lash_core::ProcessEventType> {
1368 vec![
1369 lash_core::ProcessEventType {
1370 name: "process.yield".to_string(),
1371 payload_schema: lash_core::LashSchema::any(),
1372 semantics: lash_core::ProcessEventSemanticsSpec::default(),
1373 },
1374 lash_core::ProcessEventType {
1375 name: "process.wake".to_string(),
1376 payload_schema: lash_core::LashSchema::any(),
1377 semantics: lash_core::ProcessEventSemanticsSpec {
1378 wake: Some(lash_core::ProcessWakeSpec {
1379 when: None,
1380 input: lash_core::ProcessValueSelector::Pointer("/text".to_string()),
1381 dedupe_key: lash_core::ProcessWakeDedupeKey::EventIdentity,
1382 }),
1383 ..lash_core::ProcessEventSemanticsSpec::default()
1384 },
1385 },
1386 ]
1387}
1388
1389pub fn lashlang_process_signal_event_types(
1390 process: &lashlang::ProcessDecl,
1391) -> Vec<lash_core::ProcessEventType> {
1392 process
1393 .signals
1394 .iter()
1395 .map(|signal| lash_core::ProcessEventType {
1396 name: lash_core::process_signal_event_type(signal.name.as_str())
1397 .expect("lashlang process signal declarations use parser-validated names"),
1398 payload_schema: lash_core::LashSchema::new(lashlang_type_expr_schema(&signal.ty)),
1399 semantics: lash_core::ProcessEventSemanticsSpec::default(),
1400 })
1401 .collect()
1402}
1403
1404pub fn lashlang_type_expr_schema(ty: &lashlang::TypeExpr) -> serde_json::Value {
1405 match ty {
1406 lashlang::TypeExpr::Any
1407 | lashlang::TypeExpr::Dict
1408 | lashlang::TypeExpr::Ref(_)
1409 | lashlang::TypeExpr::Process { .. }
1410 | lashlang::TypeExpr::TriggerHandle(_) => serde_json::json!({}),
1411 lashlang::TypeExpr::Str => serde_json::json!({ "type": "string" }),
1412 lashlang::TypeExpr::Int => serde_json::json!({ "type": "integer" }),
1413 lashlang::TypeExpr::Float => serde_json::json!({ "type": "number" }),
1414 lashlang::TypeExpr::Bool => serde_json::json!({ "type": "boolean" }),
1415 lashlang::TypeExpr::Null => serde_json::json!({ "type": "null" }),
1416 lashlang::TypeExpr::Enum(values) => serde_json::json!({
1417 "enum": values.iter().map(|value| value.as_str()).collect::<Vec<_>>()
1418 }),
1419 lashlang::TypeExpr::List(item) => serde_json::json!({
1420 "type": "array",
1421 "items": lashlang_type_expr_schema(item),
1422 }),
1423 lashlang::TypeExpr::Object(fields) => {
1424 let mut properties = serde_json::Map::new();
1425 let mut required = Vec::new();
1426 for field in fields {
1427 properties.insert(field.name.to_string(), lashlang_type_expr_schema(&field.ty));
1428 if !field.optional {
1429 required.push(serde_json::Value::String(field.name.to_string()));
1430 }
1431 }
1432 let mut schema = serde_json::Map::new();
1433 schema.insert(
1434 "type".to_string(),
1435 serde_json::Value::String("object".to_string()),
1436 );
1437 schema.insert(
1438 "properties".to_string(),
1439 serde_json::Value::Object(properties),
1440 );
1441 if !required.is_empty() {
1442 schema.insert("required".to_string(), serde_json::Value::Array(required));
1443 }
1444 schema.insert(
1445 "additionalProperties".to_string(),
1446 serde_json::Value::Bool(true),
1447 );
1448 serde_json::Value::Object(schema)
1449 }
1450 lashlang::TypeExpr::Union(variants) => serde_json::json!({
1451 "anyOf": variants.iter().map(lashlang_type_expr_schema).collect::<Vec<_>>()
1452 }),
1453 }
1454}
1455
1456#[cfg(test)]
1457mod segment_trace_tests {
1458 use super::{
1459 SEGMENT_BOUNDARY_DECLINED_TOTAL, record_segment_boundary_decline,
1460 trace_lifecycle_for_segment, validate_lashlang_program_hash,
1461 };
1462 use std::sync::atomic::Ordering;
1463
1464 #[test]
1465 fn multi_segment_trace_emits_one_started_and_one_finished() {
1466 let boundary = || {
1467 lash_core::ProcessRunOutcome::SegmentBoundary(lash_core::SegmentHandover {
1468 reason: lash_core::BoundaryReason::JournalBudget,
1469 program_hash: Some("program-v1".to_string()),
1470 engine_state: Vec::new(),
1471 })
1472 };
1473 let terminal = lash_core::ProcessRunOutcome::from(lash_core::ProcessAwaitOutput::Success {
1474 value: serde_json::Value::Null,
1475 control: None,
1476 });
1477 let lifecycle = [
1478 trace_lifecycle_for_segment(true, &boundary()),
1479 trace_lifecycle_for_segment(false, &boundary()),
1480 trace_lifecycle_for_segment(false, &terminal),
1481 ];
1482 assert_eq!(lifecycle.iter().filter(|(started, _)| *started).count(), 1);
1483 assert_eq!(
1484 lifecycle.iter().filter(|(_, finished)| *finished).count(),
1485 1
1486 );
1487 }
1488
1489 #[test]
1490 fn resume_rejects_changed_bytecode_program_hash_with_typed_failure() {
1491 let output = validate_lashlang_program_hash(Some("sha256:old"), "sha256:current")
1492 .expect_err("changed bytecode identity must fail closed");
1493 assert!(matches!(
1494 *output,
1495 lash_core::ProcessAwaitOutput::Failure { code, .. }
1496 if code == "restate_segment_program_hash_mismatch"
1497 ));
1498 }
1499
1500 #[test]
1501 fn declined_boundary_is_warned_and_counted() {
1502 let before = SEGMENT_BOUNDARY_DECLINED_TOTAL.load(Ordering::Relaxed);
1503 record_segment_boundary_decline(&"projected state", "test boundary decline");
1504 assert_eq!(
1505 SEGMENT_BOUNDARY_DECLINED_TOTAL.load(Ordering::Relaxed),
1506 before + 1
1507 );
1508 }
1509}