1use 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#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
30pub enum ObserveMode {
31 #[default]
33 FinalOnly,
34 Events,
36 Ledger,
38}
39
40#[derive(Clone, Debug, PartialEq, Eq)]
62pub struct RealizeRequest {
63 pub term: Term,
65 pub result_shape: Option<Ref>,
67 pub required_capabilities: Vec<CapabilityName>,
69 pub deadline: Option<Duration>,
71 pub consistency: Consistency,
73 pub mode: EvalMode,
75 pub answer_limit: Option<usize>,
77 pub buffer_limit: Option<usize>,
79 pub observe: ObserveMode,
81}
82
83impl RealizeRequest {
84 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 pub fn observing(mut self, observe: ObserveMode) -> Self {
101 self.observe = observe;
102 self
103 }
104
105 pub fn requiring(mut self, capability: CapabilityName) -> Self {
107 self.required_capabilities.push(capability);
108 self
109 }
110
111 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 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#[derive(Debug)]
164pub struct BufferedEventSource {
165 events: Mutex<VecDeque<Event>>,
166}
167
168impl BufferedEventSource {
169 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
273pub 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
311pub 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
324pub 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}