Skip to main content

lash_lashlang_runtime/
process.rs

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(), &current_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            // `emit_finished` fires after an actual execution, whose outcome is
974            // Success/Failure/Cancelled — abandonment is written out-of-band by the
975            // sweep, never returned by a run. Map it defensively to Failed.
976            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
1159/// Builds the trace runtime's read-only foreground skeleton from the workflow graph.
1160pub 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}