1use crate::{
7 capability::CapabilityName,
8 datum::Datum,
9 datum_store::DatumStore,
10 env::Cx,
11 error::{Error, Result},
12 expr::NumberLiteral,
13 id::Symbol,
14 ref_id::{ContentId, Coordinate, HandleId, Ref},
15 term::OpKey,
16};
17
18pub const EFFECT_REPLAY_VERSION: &str = "sim6-effect-replay-v1";
20
21#[derive(Clone, Debug, PartialEq, Eq)]
23pub struct Effect {
24 pub id: Ref,
26 pub kind: Symbol,
28 pub subject: Ref,
30 pub input: Ref,
32 pub result_shape: Ref,
34 pub resume_op: OpKey,
36 pub abort_op: OpKey,
38 pub requires: Vec<CapabilityName>,
40 pub replay_key: Option<ContentId>,
42}
43
44impl Effect {
45 pub fn new(
47 kind: Symbol,
48 subject: Ref,
49 input: Ref,
50 result_shape: Ref,
51 resume_op: OpKey,
52 abort_op: OpKey,
53 ) -> Self {
54 Self {
55 id: Ref::Handle(HandleId::fresh()),
56 kind,
57 subject,
58 input,
59 result_shape,
60 resume_op,
61 abort_op,
62 requires: Vec::new(),
63 replay_key: None,
64 }
65 }
66
67 pub fn with_id(mut self, id: Ref) -> Self {
69 self.id = id;
70 self
71 }
72
73 pub fn requiring(mut self, capability: CapabilityName) -> Self {
75 self.requires.push(capability);
76 self
77 }
78
79 pub fn with_requirements(mut self, requires: Vec<CapabilityName>) -> Self {
81 self.requires = requires;
82 self
83 }
84
85 pub fn with_replay_key(mut self, implementation: Option<Ref>) -> Result<Self> {
87 self.replay_key = Some(effect_replay_key(&self, implementation)?);
88 Ok(self)
89 }
90
91 pub fn ensure_replay_key(&mut self, implementation: Option<Ref>) -> Result<ContentId> {
93 if let Some(key) = &self.replay_key {
94 return Ok(key.clone());
95 }
96 let key = effect_replay_key(self, implementation)?;
97 self.replay_key = Some(key.clone());
98 Ok(key)
99 }
100}
101
102#[derive(Clone, Debug, PartialEq, Eq)]
104pub struct EffectRecord {
105 pub effect: Ref,
107 pub requested_event: Ref,
109 pub resolved_event: Option<Ref>,
111 pub result: Option<Ref>,
113 pub aborted: bool,
115}
116
117pub fn resolve_effect<F>(cx: &mut Cx, mut effect: Effect, perform: F) -> Result<Ref>
148where
149 F: FnOnce(&mut Cx, &Effect) -> Result<Ref>,
150{
151 let preimage = effect_replay_preimage(&effect, None);
152 let replay_key = match effect.replay_key.clone() {
153 Some(key) => key,
154 None => cx.datum_store_mut().intern(preimage)?,
155 };
156 effect.replay_key = Some(replay_key.clone());
157
158 let cassette_result = cx.with_effect_ledger(|cx, ledger| {
159 ledger.record_requested(cx.datum_store_mut(), effect.clone())?;
160 Ok(ledger.cassette_result(&replay_key).cloned())
161 })?;
162
163 if let Err(err) = cx.require_all(&effect.requires) {
164 record_effect_failure(cx, effect.id.clone(), &err)?;
165 return Err(err);
166 }
167
168 if let Some(result) = cassette_result {
169 cx.with_effect_ledger(|cx, ledger| {
170 ledger.record_resolved(cx.datum_store_mut(), effect.id.clone(), result.clone())?;
171 Ok(())
172 })?;
173 return Ok(result);
174 }
175
176 match perform(cx, &effect) {
177 Ok(result) => {
178 cx.with_effect_ledger(|cx, ledger| {
179 ledger.record_resolved(cx.datum_store_mut(), effect.id.clone(), result.clone())?;
180 Ok(())
181 })?;
182 Ok(result)
183 }
184 Err(err) => {
185 record_effect_failure(cx, effect.id, &err)?;
186 Err(err)
187 }
188 }
189}
190
191pub fn effect_replay_key(effect: &Effect, implementation: Option<Ref>) -> Result<ContentId> {
193 effect_replay_preimage(effect, implementation).content_id()
194}
195
196pub fn effect_replay_preimage(effect: &Effect, implementation: Option<Ref>) -> Datum {
198 let mut requires = effect.requires.clone();
199 requires.sort();
200 requires.dedup();
201 let mut fields = vec![
202 (
203 Symbol::new("version"),
204 Datum::String(EFFECT_REPLAY_VERSION.to_owned()),
205 ),
206 (Symbol::new("kind"), Datum::Symbol(effect.kind.clone())),
207 (Symbol::new("subject"), ref_datum(effect.subject.clone())),
208 (Symbol::new("input"), ref_datum(effect.input.clone())),
209 (
210 Symbol::new("result-shape"),
211 ref_datum(effect.result_shape.clone()),
212 ),
213 (
214 Symbol::new("resume-op"),
215 op_key_datum(effect.resume_op.clone()),
216 ),
217 (
218 Symbol::new("abort-op"),
219 op_key_datum(effect.abort_op.clone()),
220 ),
221 (
222 Symbol::new("requires"),
223 Datum::List(
224 requires
225 .into_iter()
226 .map(|capability| Datum::String(capability.as_str().to_owned()))
227 .collect(),
228 ),
229 ),
230 ];
231 if let Some(implementation) = implementation {
232 fields.push((Symbol::new("implementation"), ref_datum(implementation)));
233 }
234 Datum::Node {
235 tag: core_symbol("EffectReplayKey"),
236 fields,
237 }
238}
239
240pub fn effect_control_prompt_kind() -> Symbol {
242 effect_symbol("control-prompt")
243}
244
245pub fn effect_control_capture_kind() -> Symbol {
247 effect_symbol("control-capture")
248}
249
250pub fn effect_control_abort_kind() -> Symbol {
252 effect_symbol("control-abort")
253}
254
255pub fn effect_control_resume_kind() -> Symbol {
257 effect_symbol("control-resume")
258}
259
260pub fn effect_resume_op_key() -> OpKey {
262 OpKey::new(effect_symbol("control"), Symbol::new("resume"), 1)
263}
264
265pub fn effect_abort_op_key() -> OpKey {
267 OpKey::new(effect_symbol("control"), Symbol::new("abort"), 1)
268}
269
270#[cfg(test)]
271fn effect_test_kind(name: &str) -> Symbol {
272 effect_symbol(name)
273}
274
275fn record_effect_failure(cx: &mut Cx, effect: Ref, err: &Error) -> Result<()> {
276 let error_ref = error_ref(cx, err)?;
277 cx.with_effect_ledger(|cx, ledger| {
278 ledger.record_failed(cx.datum_store_mut(), effect, error_ref)?;
279 Ok(())
280 })
281}
282
283fn error_ref(cx: &mut Cx, err: &Error) -> Result<Ref> {
284 let id = cx
285 .datum_store_mut()
286 .intern(Datum::String(err.to_string()))?;
287 Ok(Ref::Content(id))
288}
289
290fn ref_datum(reference: Ref) -> Datum {
291 match reference {
292 Ref::Symbol(symbol) => Datum::Node {
293 tag: core_symbol("ref"),
294 fields: vec![
295 (Symbol::new("kind"), Datum::Symbol(core_symbol("symbol"))),
296 (Symbol::new("symbol"), Datum::Symbol(symbol)),
297 ],
298 },
299 Ref::Content(content) => Datum::Node {
300 tag: core_symbol("ref"),
301 fields: vec![
302 (Symbol::new("kind"), Datum::Symbol(core_symbol("content"))),
303 (Symbol::new("content"), content_id_datum(content)),
304 ],
305 },
306 Ref::Handle(handle) => Datum::Node {
307 tag: core_symbol("ref"),
308 fields: vec![
309 (Symbol::new("kind"), Datum::Symbol(core_symbol("handle"))),
310 (Symbol::new("id"), handle_id_datum(handle)),
311 ],
312 },
313 Ref::Coord(coordinate) => coordinate_datum(coordinate),
314 }
315}
316
317fn coordinate_datum(coordinate: Coordinate) -> Datum {
318 Datum::Node {
319 tag: core_symbol("ref"),
320 fields: vec![
321 (Symbol::new("kind"), Datum::Symbol(core_symbol("coord"))),
322 (Symbol::new("space"), Datum::Symbol(coordinate.space)),
323 (Symbol::new("ordinal"), content_id_datum(coordinate.ordinal)),
324 ],
325 }
326}
327
328fn content_id_datum(content: ContentId) -> Datum {
329 Datum::Node {
330 tag: core_symbol("content-id"),
331 fields: vec![
332 (Symbol::new("algorithm"), Datum::Symbol(content.algorithm)),
333 (Symbol::new("bytes"), Datum::Bytes(content.bytes.to_vec())),
334 ],
335 }
336}
337
338fn handle_id_datum(handle: HandleId) -> Datum {
339 Datum::Bytes(handle.0.to_be_bytes().to_vec())
340}
341
342fn op_key_datum(op: OpKey) -> Datum {
343 Datum::Node {
344 tag: core_symbol("op-key"),
345 fields: vec![
346 (Symbol::new("namespace"), Datum::Symbol(op.namespace)),
347 (Symbol::new("name"), Datum::Symbol(op.name)),
348 (
349 Symbol::new("version"),
350 Datum::Number(NumberLiteral {
351 domain: core_symbol("u16"),
352 canonical: op.version.to_string(),
353 }),
354 ),
355 ],
356 }
357}
358
359fn effect_symbol(name: &str) -> Symbol {
360 Symbol::qualified("effect", name)
361}
362
363fn core_symbol(name: &str) -> Symbol {
364 Symbol::qualified("core", name)
365}
366
367#[cfg(test)]
368mod tests {
369 use std::sync::{
370 Arc,
371 atomic::{AtomicUsize, Ordering},
372 };
373
374 use super::*;
375 use crate::EventKind;
376
377 use crate::testing::bare_cx as cx;
378
379 fn effect(input: Ref) -> Effect {
380 Effect::new(
381 effect_test_kind("tool-call"),
382 Ref::Symbol(Symbol::qualified("test", "tool")),
383 input,
384 Ref::Symbol(core_symbol("Any")),
385 effect_resume_op_key(),
386 effect_abort_op_key(),
387 )
388 }
389
390 #[test]
391 fn same_replay_preimage_gives_same_key() {
392 let left = effect(Ref::Symbol(Symbol::qualified("test", "input")));
393 let right = effect(Ref::Symbol(Symbol::qualified("test", "input")));
394
395 assert_eq!(
396 effect_replay_key(&left, None).unwrap(),
397 effect_replay_key(&right, None).unwrap()
398 );
399 }
400
401 #[test]
402 fn changed_input_gives_different_key() {
403 let left = effect(Ref::Symbol(Symbol::qualified("test", "left")));
404 let right = effect(Ref::Symbol(Symbol::qualified("test", "right")));
405
406 assert_ne!(
407 effect_replay_key(&left, None).unwrap(),
408 effect_replay_key(&right, None).unwrap()
409 );
410 }
411
412 #[test]
413 fn resolving_effect_emits_requested_and_resolved_events() {
414 let mut cx = cx();
415 let result = Ref::Symbol(Symbol::qualified("test", "result"));
416
417 let actual = resolve_effect(&mut cx, effect(Ref::Symbol(Symbol::new("input"))), {
418 let result = result.clone();
419 move |_cx, _effect| Ok(result)
420 })
421 .unwrap();
422
423 assert_eq!(actual, result);
424 let records = cx.effect_ledger().records();
425 assert_eq!(records.len(), 1);
426 assert_eq!(records[0].result, Some(result.clone()));
427 let events = cx.effect_ledger().events_for_run();
428 assert!(matches!(events[0].kind, EventKind::EffectRequested { .. }));
429 assert!(matches!(events[1].kind, EventKind::EffectResolved { .. }));
430 }
431
432 #[test]
433 fn missing_capability_denies_effect_before_performer_runs() {
434 let mut cx = cx();
435 let calls = Arc::new(AtomicUsize::new(0));
436 let err = resolve_effect(
437 &mut cx,
438 effect(Ref::Symbol(Symbol::new("input")))
439 .requiring(CapabilityName::new("test.required")),
440 {
441 let calls = calls.clone();
442 move |_cx, _effect| {
443 calls.fetch_add(1, Ordering::SeqCst);
444 Ok(Ref::Symbol(Symbol::new("unreachable")))
445 }
446 },
447 )
448 .unwrap_err();
449
450 assert!(
451 matches!(err, Error::CapabilityDenied { capability } if capability.as_str() == "test.required")
452 );
453 assert_eq!(calls.load(Ordering::SeqCst), 0);
454 assert!(cx.effect_ledger().records()[0].aborted);
455 }
456
457 #[test]
458 fn cassette_result_is_used_when_replay_key_matches() {
459 let mut cx = cx();
460 let mut effect = effect(Ref::Symbol(Symbol::new("input")));
461 let key = effect.ensure_replay_key(None).unwrap();
462 let cassette = Ref::Symbol(Symbol::qualified("test", "cassette-result"));
463 cx.effect_ledger_mut()
464 .insert_cassette_result(key, cassette.clone());
465 let calls = Arc::new(AtomicUsize::new(0));
466
467 let actual = resolve_effect(&mut cx, effect, {
468 let calls = calls.clone();
469 move |_cx, _effect| {
470 calls.fetch_add(1, Ordering::SeqCst);
471 Ok(Ref::Symbol(Symbol::new("performed")))
472 }
473 })
474 .unwrap();
475
476 assert_eq!(actual, cassette);
477 assert_eq!(calls.load(Ordering::SeqCst), 0);
478 }
479}