Skip to main content

sim_kernel/
realize.rs

1//! The `realize` surface: the location-transparent distributed eval contract.
2//!
3//! The kernel defines the realize request, observe modes, and event-draining
4//! contract that server and agent code target instead of transport-specific
5//! APIs; libraries supply the concrete transports behind it.
6
7use std::{collections::VecDeque, sync::Mutex, time::Duration};
8
9use crate::{
10    capability::CapabilityName,
11    env::Cx,
12    error::{Error, Result},
13    eval::{Consistency, EvalFabric, EvalMode, EvalReply, EvalRequest},
14    event::{Event, EventKind, EventSource},
15    event_ledger::EventLedger,
16    expr::Expr,
17    handle_store::HandleStore,
18    id::{CORE_SEQUENCE_CLASS_ID, Symbol},
19    object::{ClassRef, Object, ShapeRef},
20    ref_id::{HandleId, Ref},
21    ref_resolver::{RefResolver, ResolvedRef, TemporaryRefResolver, value_from_ref},
22    seq::{Sequence, SequenceItem, sequence_item_from_event},
23    shape_check::check_shape_value,
24    term::Term,
25    value::Value,
26};
27
28/// How much of an evaluation a caller wants to observe.
29#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
30pub enum ObserveMode {
31    /// Observe only the final value (the default).
32    #[default]
33    FinalOnly,
34    /// Stream intermediate events as they occur.
35    Events,
36    /// Collect the full event ledger alongside the final value.
37    Ledger,
38}
39
40/// A reference-based realize request: the portable form of an [`EvalRequest`].
41///
42/// Where [`EvalRequest`] carries an in-process [`Expr`] and live values, a
43/// [`RealizeRequest`] carries a [`Term`] and [`Ref`]s, so it can cross a
44/// transport boundary. The two convert via
45/// [`from_eval_request`](RealizeRequest::from_eval_request) and
46/// [`to_eval_request`](RealizeRequest::to_eval_request). See the README
47/// section "Distributed evaluation".
48///
49/// # Examples
50///
51/// ```
52/// use sim_kernel::realize::{ObserveMode, RealizeRequest};
53/// use sim_kernel::term::Term;
54/// use sim_kernel::Symbol;
55///
56/// let request = RealizeRequest::new(Term::Local(Symbol::new("x")))
57///     .observing(ObserveMode::Events);
58/// assert_eq!(request.observe, ObserveMode::Events);
59/// assert!(request.required_capabilities.is_empty());
60/// ```
61#[derive(Clone, Debug, PartialEq, Eq)]
62pub struct RealizeRequest {
63    /// The term to evaluate.
64    pub term: Term,
65    /// Optional reference to a shape the result must satisfy.
66    pub result_shape: Option<Ref>,
67    /// Capabilities the evaluation requires.
68    pub required_capabilities: Vec<CapabilityName>,
69    /// Optional wall-clock deadline for the evaluation.
70    pub deadline: Option<Duration>,
71    /// Where the request may be answered from.
72    pub consistency: Consistency,
73    /// Which evaluation discipline to run under.
74    pub mode: EvalMode,
75    /// Optional cap on the number of answers (logic mode).
76    pub answer_limit: Option<usize>,
77    /// Optional buffer size for streamed events.
78    pub buffer_limit: Option<usize>,
79    /// How much of the evaluation to observe.
80    pub observe: ObserveMode,
81}
82
83impl RealizeRequest {
84    /// Creates a request to realize `term` with all defaults.
85    pub fn new(term: Term) -> Self {
86        Self {
87            term,
88            result_shape: None,
89            required_capabilities: Vec::new(),
90            deadline: None,
91            consistency: Consistency::default(),
92            mode: EvalMode::default(),
93            answer_limit: None,
94            buffer_limit: None,
95            observe: ObserveMode::default(),
96        }
97    }
98
99    /// Sets the [`ObserveMode`], returning the request (builder style).
100    pub fn observing(mut self, observe: ObserveMode) -> Self {
101        self.observe = observe;
102        self
103    }
104
105    /// Adds a required capability, returning the request (builder style).
106    pub fn requiring(mut self, capability: CapabilityName) -> Self {
107        self.required_capabilities.push(capability);
108        self
109    }
110
111    /// Builds a [`RealizeRequest`] from an in-process [`EvalRequest`].
112    ///
113    /// Live values become [`Ref`]s and the [`Expr`] becomes a [`Term`] so the
114    /// request can cross a transport boundary.
115    pub fn from_eval_request(cx: &mut Cx, request: &EvalRequest) -> Result<Self> {
116        Ok(Self {
117            term: term_from_eval_expr(cx, &request.expr)?,
118            result_shape: request
119                .result_shape
120                .as_ref()
121                .map(|shape| handle_ref_for_value(cx, shape))
122                .transpose()?,
123            required_capabilities: request.required_capabilities.clone(),
124            deadline: request.deadline,
125            consistency: request.consistency,
126            mode: request.mode,
127            answer_limit: request.answer_limit,
128            buffer_limit: request.stream_buffer,
129            observe: observe_from_eval_request(request),
130        })
131    }
132
133    /// Resolves this request back into an in-process [`EvalRequest`].
134    ///
135    /// The inverse of [`from_eval_request`](RealizeRequest::from_eval_request):
136    /// [`Term`] and [`Ref`]s are resolved against `cx` into an [`Expr`] and
137    /// live values.
138    pub fn to_eval_request(&self, cx: &mut Cx) -> Result<EvalRequest> {
139        Ok(EvalRequest {
140            expr: eval_expr_from_term(cx, &self.term)?,
141            result_shape: self
142                .result_shape
143                .as_ref()
144                .map(|reference| shape_from_ref(cx, reference))
145                .transpose()?,
146            required_capabilities: self.required_capabilities.clone(),
147            deadline: self.deadline,
148            consistency: self.consistency,
149            mode: self.mode,
150            answer_limit: self.answer_limit,
151            stream_buffer: self.buffer_limit,
152            stream: self.observe == ObserveMode::Events,
153            trace: self.observe != ObserveMode::FinalOnly,
154        })
155    }
156}
157
158// sim-non-citizen(reason = "buffered event-source handle; descriptor data is the event sequence payload", kind = "handle", descriptor = "core/EventSource")
159/// An in-memory [`EventSource`] backed by a buffered queue of events.
160///
161/// Returned by [`realize_events`] to replay a completed evaluation's events
162/// (diagnostics, trace, final value, done) without a live transport.
163#[derive(Debug)]
164pub struct BufferedEventSource {
165    events: Mutex<VecDeque<Event>>,
166}
167
168impl BufferedEventSource {
169    /// Wraps an ordered list of events as a drainable source.
170    pub fn new(events: Vec<Event>) -> Self {
171        Self {
172            events: Mutex::new(events.into()),
173        }
174    }
175}
176
177impl EventSource for BufferedEventSource {
178    fn next(&self, _cx: &mut Cx) -> Result<Option<Event>> {
179        Ok(self
180            .events
181            .lock()
182            .map_err(|_| Error::PoisonedLock("event source"))?
183            .pop_front())
184    }
185
186    fn close(&self, _cx: &mut Cx) -> Result<()> {
187        self.events
188            .lock()
189            .map_err(|_| Error::PoisonedLock("event source"))?
190            .clear();
191        Ok(())
192    }
193}
194
195impl Object for BufferedEventSource {
196    fn display(&self, _cx: &mut Cx) -> Result<String> {
197        Ok("#<event-source>".to_owned())
198    }
199
200    fn as_any(&self) -> &dyn std::any::Any {
201        self
202    }
203}
204
205impl crate::ObjectCompat for BufferedEventSource {
206    fn class(&self, cx: &mut Cx) -> Result<ClassRef> {
207        cx.factory().class_stub(
208            CORE_SEQUENCE_CLASS_ID,
209            Symbol::qualified("core", "EventSource"),
210        )
211    }
212    fn as_sequence(&self) -> Option<&dyn Sequence> {
213        Some(self)
214    }
215}
216
217impl Sequence for BufferedEventSource {
218    fn next_item(&self, cx: &mut Cx) -> Result<Option<SequenceItem>> {
219        while let Some(event) = EventSource::next(self, cx)? {
220            let done = matches!(event.kind, EventKind::Done);
221            if let Some(item) = sequence_item_from_event(cx, event)? {
222                return Ok(Some(item));
223            }
224            if done {
225                return Ok(None);
226            }
227        }
228        Ok(None)
229    }
230
231    fn close(&self, cx: &mut Cx) -> Result<()> {
232        EventSource::close(self, cx)
233    }
234
235    fn peek_item(&self, cx: &mut Cx) -> Result<Option<SequenceItem>> {
236        let events = self
237            .events
238            .lock()
239            .map_err(|_| Error::PoisonedLock("event source"))?
240            .iter()
241            .cloned()
242            .collect::<Vec<_>>();
243        for event in events {
244            let done = matches!(event.kind, EventKind::Done);
245            if let Some(item) = sequence_item_from_event(cx, event)? {
246                return Ok(Some(item));
247            }
248            if done {
249                return Ok(None);
250            }
251        }
252        Ok(None)
253    }
254
255    fn is_done(&self, _cx: &mut Cx) -> Result<bool> {
256        let events = self
257            .events
258            .lock()
259            .map_err(|_| Error::PoisonedLock("event source"))?;
260        for event in events.iter() {
261            match event.kind {
262                EventKind::Chunk { .. } | EventKind::Final(_) | EventKind::Failed(_) => {
263                    return Ok(false);
264                }
265                EventKind::Done => return Ok(true),
266                _ => {}
267            }
268        }
269        Ok(true)
270    }
271}
272
273/// Realizes `request` against `target` and returns its events as a source.
274///
275/// Runs the evaluation through the [`EvalFabric`], records the resulting
276/// diagnostics, optional trace, final value, and a terminating `done` into an
277/// event ledger, and hands them back as a [`BufferedEventSource`].
278pub fn realize_events(
279    cx: &mut Cx,
280    target: &dyn EvalFabric,
281    request: EvalRequest,
282) -> Result<BufferedEventSource> {
283    let eventful_request = RealizeRequest::from_eval_request(cx, &request)?;
284    let request_ref = ref_for_realize_request(cx, &eventful_request)?;
285    let run = Ref::Handle(HandleId::fresh());
286    let mut ledger = EventLedger::new();
287    ledger.started(run.clone(), request_ref)?;
288
289    let result_shape = request.result_shape.clone();
290    let reply = target.realize(cx, request)?;
291    if let Some(shape) = result_shape {
292        check_shape_value(cx, &shape, None, reply.value.clone())?;
293    }
294    for diagnostic in &reply.diagnostics {
295        ledger.push(run.clone(), EventKind::Diagnostic(diagnostic.clone()))?;
296    }
297    if let Some(trace) = &reply.trace {
298        ledger.push(
299            run.clone(),
300            EventKind::Trace(handle_ref_for_value(cx, trace)?),
301        )?;
302    }
303    ledger.final_value(run.clone(), handle_ref_for_value(cx, &reply.value)?)?;
304    ledger.done(run.clone())?;
305
306    Ok(BufferedEventSource::new(
307        ledger.events_for_run(&run).to_vec(),
308    ))
309}
310
311/// Realizes `request` against `target` and collects only the final reply.
312///
313/// A convenience over [`realize_events`] plus [`drain_events_to_reply`] for
314/// callers that want the [`EvalReply`] rather than the event stream.
315pub fn realize_final(
316    cx: &mut Cx,
317    target: &dyn EvalFabric,
318    request: EvalRequest,
319) -> Result<EvalReply> {
320    let events = realize_events(cx, target, request)?;
321    drain_events_to_reply(cx, &events)
322}
323
324/// Drains an [`EventSource`] into a single [`EvalReply`].
325///
326/// Accumulates diagnostics and an optional trace, captures the final value,
327/// and stops at `done`. A `Failed` event is turned into an error, and a
328/// stream that ends without a final value is an error.
329pub fn drain_events_to_reply(cx: &mut Cx, source: &dyn EventSource) -> Result<EvalReply> {
330    let mut diagnostics = Vec::new();
331    let mut trace = None;
332    let mut value = None;
333
334    while let Some(event) = source.next(cx)? {
335        match event.kind {
336            EventKind::Diagnostic(diagnostic) => diagnostics.push(diagnostic),
337            EventKind::Trace(reference) => trace = Some(value_from_ref(cx, &reference)?),
338            EventKind::Final(reference) => value = Some(value_from_ref(cx, &reference)?),
339            EventKind::Failed(reference) => return Err(error_from_failed_ref(cx, &reference)),
340            EventKind::Done => break,
341            EventKind::Started { .. }
342            | EventKind::Claim { .. }
343            | EventKind::Chunk { .. }
344            | EventKind::EffectRequested { .. }
345            | EventKind::EffectResolved { .. }
346            | EventKind::Capture { .. }
347            | EventKind::Card { .. } => {}
348        }
349    }
350
351    Ok(EvalReply {
352        value: value
353            .ok_or_else(|| Error::Eval("eventful realize ended without Final".to_owned()))?,
354        diagnostics,
355        trace,
356    })
357}
358
359fn observe_from_eval_request(request: &EvalRequest) -> ObserveMode {
360    if request.stream {
361        ObserveMode::Events
362    } else if request.trace {
363        ObserveMode::Ledger
364    } else {
365        ObserveMode::FinalOnly
366    }
367}
368
369fn term_from_eval_expr(cx: &mut Cx, expr: &Expr) -> Result<Term> {
370    match Term::try_from(expr.clone()) {
371        Ok(term) => Ok(term),
372        Err(_) => {
373            let value = cx.factory().expr(expr.clone())?;
374            Ok(Term::Ref(handle_ref_for_value(cx, &value)?))
375        }
376    }
377}
378
379fn eval_expr_from_term(cx: &mut Cx, term: &Term) -> Result<Expr> {
380    let Term::Ref(reference) = term else {
381        return Ok(Expr::from(term.clone()));
382    };
383    match TemporaryRefResolver::new().resolve_ref(cx, reference)? {
384        ResolvedRef::Symbol(symbol) => Ok(Expr::Symbol(symbol)),
385        ResolvedRef::Datum(datum) => Ok(Expr::from(datum)),
386        ResolvedRef::Value(value) => value.object().as_expr(cx),
387        ResolvedRef::Coordinate(_) | ResolvedRef::Missing(_) => Ok(Expr::from(term.clone())),
388    }
389}
390
391fn shape_from_ref(cx: &mut Cx, reference: &Ref) -> Result<ShapeRef> {
392    match TemporaryRefResolver::new().resolve_ref(cx, reference)? {
393        ResolvedRef::Symbol(symbol) => cx.resolve_shape(&symbol),
394        ResolvedRef::Value(value) => {
395            if value.object().as_shape().is_some() {
396                Ok(value)
397            } else if let Some(class) = value.object().as_class() {
398                class.instance_shape(cx)
399            } else {
400                Err(Error::TypeMismatch {
401                    expected: "shape ref",
402                    found: "non-shape",
403                })
404            }
405        }
406        ResolvedRef::Datum(_) => Err(Error::TypeMismatch {
407            expected: "shape ref",
408            found: "non-shape",
409        }),
410        ResolvedRef::Coordinate(_) | ResolvedRef::Missing(_) => Err(Error::Eval(format!(
411            "unresolved result shape ref {reference:?}"
412        ))),
413    }
414}
415
416fn ref_for_realize_request(cx: &mut Cx, request: &RealizeRequest) -> Result<Ref> {
417    let value = realize_request_value(cx, request)?;
418    handle_ref_for_value(cx, &value)
419}
420
421fn realize_request_value(cx: &mut Cx, request: &RealizeRequest) -> Result<Value> {
422    let capabilities = request
423        .required_capabilities
424        .iter()
425        .map(|capability| cx.factory().string(capability.as_str().to_owned()))
426        .collect::<Result<Vec<_>>>()?;
427    let term = cx.factory().expr(Expr::from(request.term.clone()))?;
428    let result_shape = optional_ref_value(cx, request.result_shape.as_ref())?;
429    let requires = cx.factory().list(capabilities)?;
430    let deadline = match request.deadline {
431        Some(deadline) => cx.factory().string(format!("{}ms", deadline.as_millis()))?,
432        None => cx.factory().nil()?,
433    };
434    let consistency = cx.factory().symbol(request.consistency.as_symbol())?;
435    let mode = cx.factory().symbol(request.mode.as_symbol())?;
436    let answer_limit = optional_usize_value(cx, request.answer_limit)?;
437    let buffer_limit = optional_usize_value(cx, request.buffer_limit)?;
438    let observe = cx.factory().symbol(observe_symbol(request.observe))?;
439
440    cx.factory().table(vec![
441        (Symbol::new("term"), term),
442        (Symbol::new("result-shape"), result_shape),
443        (Symbol::new("requires"), requires),
444        (Symbol::new("deadline"), deadline),
445        (Symbol::new("consistency"), consistency),
446        (Symbol::new("mode"), mode),
447        (Symbol::new("answer-limit"), answer_limit),
448        (Symbol::new("buffer-limit"), buffer_limit),
449        (Symbol::new("observe"), observe),
450    ])
451}
452
453fn optional_ref_value(cx: &mut Cx, reference: Option<&Ref>) -> Result<Value> {
454    match reference {
455        Some(reference) => cx.factory().expr(Expr::from(Term::Ref(reference.clone()))),
456        None => cx.factory().nil(),
457    }
458}
459
460fn optional_usize_value(cx: &mut Cx, value: Option<usize>) -> Result<Value> {
461    match value {
462        Some(value) => cx.factory().string(value.to_string()),
463        None => cx.factory().nil(),
464    }
465}
466
467fn observe_symbol(observe: ObserveMode) -> Symbol {
468    Symbol::new(match observe {
469        ObserveMode::FinalOnly => "final-only",
470        ObserveMode::Events => "events",
471        ObserveMode::Ledger => "ledger",
472    })
473}
474
475fn handle_ref_for_value(cx: &mut Cx, value: &Value) -> Result<Ref> {
476    Ok(Ref::Handle(cx.handles_mut().intern(value.clone())))
477}
478
479fn error_from_failed_ref(cx: &mut Cx, reference: &Ref) -> Error {
480    match value_from_ref(cx, reference) {
481        Ok(value) => match value.object().display(cx) {
482            Ok(message) => Error::Eval(message),
483            Err(err) => err,
484        },
485        Err(err) => err,
486    }
487}