1use std::collections::BTreeSet;
2
3use sim_codec::{Input, decode_with_codec};
4use sim_kernel::{
5 CapabilityName, CapabilitySet, Consistency, Cx, Diagnostic, Error, EvalMode, EvalReply,
6 EvalRequest, Expr, ReadPolicy, Result, Symbol,
7};
8use sim_lib_agent_runner_core::{ModelResponse, OutputContract, terminal_model_content};
9use sim_lib_stream_core::{DevCassette, DevEvent};
10use sim_value::{access::field, build::entry};
11
12use crate::{
13 AuthorTask, CompiledIntent, IntentStatus, RankedContractCard, RouteAttempt, RouteAttemptStatus,
14 RoutePolicy, RouteTarget,
15 author::{author_model_request, project_contracts_with_cards},
16};
17
18pub struct AuthorOutcome {
20 pub checked_form: Option<Expr>,
22 pub realized: Option<Expr>,
24 pub attempts: Vec<RouteAttempt>,
26 pub diagnostics: Vec<Diagnostic>,
28 pub cassette: DevCassette,
30}
31
32pub fn authorized_capabilities(projection_cards: &[RankedContractCard]) -> Vec<Symbol> {
34 projection_cards
35 .iter()
36 .flat_map(|ranked| ranked.card.capability_symbols.iter().cloned())
37 .collect::<BTreeSet<_>>()
38 .into_iter()
39 .collect()
40}
41
42pub fn run_author_task(
44 cx: &mut Cx,
45 task: &AuthorTask,
46 policy: &RoutePolicy<'_>,
47) -> Result<AuthorOutcome> {
48 let (projection, projection_cards) =
49 project_contracts_with_cards(&task.contract_cards, &task.projection_caps);
50 let mut diagnostics = projection.diagnostics.clone();
51 let ceiling = authorized_capabilities(&projection_cards);
52 let model_request = match author_model_request(cx, task, &projection) {
53 Ok(request) => request,
54 Err(err) => {
55 diagnostics.push(Diagnostic::info(err.to_string()));
56 return author_outcome(task, None, None, Vec::new(), diagnostics);
57 }
58 };
59
60 let mut attempts = Vec::new();
61 if policy.ladder.is_empty() {
62 diagnostics.push(Diagnostic::info(format!(
63 "author task {} has no route targets",
64 task.name
65 )));
66 return author_outcome(task, None, None, attempts, diagnostics);
67 }
68
69 for target in &policy.ladder {
70 if let Some(reason) = target_skip_reason(target, &ceiling) {
71 attempts.push(RouteAttempt {
72 target: target.id.clone(),
73 status: RouteAttemptStatus::Skipped,
74 reason: Some(reason),
75 });
76 continue;
77 }
78
79 for _ in 0..policy.escalate_after {
80 match run_target_once(cx, task, policy, target, &model_request, &ceiling) {
81 Ok(accepted) => {
82 diagnostics.extend(accepted.diagnostics);
83 attempts.push(RouteAttempt {
84 target: target.id.clone(),
85 status: RouteAttemptStatus::Accepted,
86 reason: None,
87 });
88 return author_outcome(
89 task,
90 Some(accepted.checked_form),
91 Some(accepted.realized),
92 attempts,
93 diagnostics,
94 );
95 }
96 Err(reason) => {
97 attempts.push(RouteAttempt {
98 target: target.id.clone(),
99 status: RouteAttemptStatus::Failed,
100 reason: Some(reason),
101 });
102 }
103 }
104 }
105 }
106
107 diagnostics.push(Diagnostic::info(format!(
108 "author task {} exhausted route policy without a checked form",
109 task.name
110 )));
111 author_outcome(task, None, None, attempts, diagnostics)
112}
113
114struct AcceptedAuthorRun {
115 checked_form: Expr,
116 realized: Expr,
117 diagnostics: Vec<Diagnostic>,
118}
119
120fn run_target_once(
121 cx: &mut Cx,
122 task: &AuthorTask,
123 policy: &RoutePolicy<'_>,
124 target: &RouteTarget<'_>,
125 model_request: &sim_lib_agent_runner_core::ModelRequest,
126 ceiling: &[Symbol],
127) -> std::result::Result<AcceptedAuthorRun, String> {
128 let model_reply = target
129 .fabric
130 .realize(cx, model_eval_request(model_request.clone()))
131 .map_err(|err| format!("model request failed: {err}"))?;
132 let response = decode_model_response(cx, model_reply)
133 .map_err(|err| format!("model response failed: {err}"))?;
134 let checked_form = decode_terminal_form(cx, task, &response)
135 .map_err(|err| format!("terminal output failed: {err}"))?;
136 let required = required_capabilities_for_form(&checked_form);
137 if let Some(reason) = outside_ceiling_reason(&required, ceiling) {
138 return Err(reason);
139 }
140
141 let narrowed = diminished_capabilities(cx.capabilities(), ceiling);
142 let realize_request = realize_eval_request(checked_form.clone(), &required);
143 let EvalReply {
144 value, diagnostics, ..
145 } = cx
146 .with_capabilities(narrowed, |scoped| {
147 target.fabric.realize(scoped, realize_request)
148 })
149 .map_err(|err| format!("realize failed: {err}"))?;
150 let realized = value
151 .object()
152 .as_expr(cx)
153 .map_err(|err| format!("realize value failed: {err}"))?;
154 let matched = task
155 .return_shape
156 .check_expr(cx, &realized)
157 .map_err(|err| format!("realized shape check failed: {err}"))?;
158 if !matched.accepted {
159 return Err(format!(
160 "realized form failed return Shape: {}",
161 shape_diagnostics(&matched.diagnostics)
162 ));
163 }
164
165 let verify_report = policy
166 .verify_catalog
167 .verify_answer(cx, &author_intent(task), &realized)
168 .map_err(|err| format!("verifier check failed: {err}"))?;
169 if !verify_report.accepted() {
170 let reasons = verify_report
171 .failed
172 .iter()
173 .map(|failure| format!("{}: {}", failure.id, failure.reason))
174 .collect::<Vec<_>>()
175 .join("; ");
176 return Err(format!("verifier check failed: {reasons}"));
177 }
178
179 Ok(AcceptedAuthorRun {
180 checked_form,
181 realized,
182 diagnostics,
183 })
184}
185
186fn decode_model_response(cx: &mut Cx, reply: EvalReply) -> Result<ModelResponse> {
187 ModelResponse::try_from(reply.value.object().as_expr(cx)?)
188}
189
190fn decode_terminal_form(cx: &mut Cx, task: &AuthorTask, response: &ModelResponse) -> Result<Expr> {
191 let output = OutputContract::for_shape(
192 task.target_codec.clone(),
193 task.return_shape_expr.clone(),
194 task.return_shape.as_ref(),
195 task.strict_grammar,
196 );
197 let input = match terminal_model_content(response) {
198 Ok(Expr::String(text)) => Input::Text(text.clone()),
199 Ok(Expr::Bytes(bytes)) => Input::Bytes(bytes.clone()),
200 Ok(Expr::Map(_)) => match field(terminal_model_content(response)?, "text") {
201 Some(Expr::String(text)) => Input::Text(text.clone()),
202 _ => {
203 return Err(Error::Eval(
204 "terminal content map must carry text".to_owned(),
205 ));
206 }
207 },
208 Ok(other) => {
209 return Err(Error::Eval(format!(
210 "terminal content must be text or bytes, found {other:?}"
211 )));
212 }
213 Err(err) => return Err(err),
214 };
215 let decoded = decode_with_codec(cx, &output.codec, input, ReadPolicy::default())?;
216 let matched = task.return_shape.check_expr(cx, &decoded)?;
217 if !matched.accepted {
218 return Err(Error::Eval(format!(
219 "decoded form failed return Shape: {}",
220 shape_diagnostics(&matched.diagnostics)
221 )));
222 }
223 Ok(decoded)
224}
225
226fn model_eval_request(model_request: sim_lib_agent_runner_core::ModelRequest) -> EvalRequest {
227 EvalRequest {
228 expr: Expr::from(model_request),
229 result_shape: None,
230 required_capabilities: Vec::new(),
231 deadline: None,
232 consistency: Consistency::default(),
233 mode: EvalMode::default(),
234 answer_limit: None,
235 stream_buffer: None,
236 stream: false,
237 trace: false,
238 }
239}
240
241fn realize_eval_request(expr: Expr, required: &[Symbol]) -> EvalRequest {
242 EvalRequest {
243 expr,
244 result_shape: None,
245 required_capabilities: required.iter().map(capability_name).collect(),
246 deadline: None,
247 consistency: Consistency::default(),
248 mode: EvalMode::default(),
249 answer_limit: None,
250 stream_buffer: None,
251 stream: false,
252 trace: false,
253 }
254}
255
256fn author_outcome(
257 task: &AuthorTask,
258 checked_form: Option<Expr>,
259 realized: Option<Expr>,
260 attempts: Vec<RouteAttempt>,
261 diagnostics: Vec<Diagnostic>,
262) -> Result<AuthorOutcome> {
263 let cassette = author_cassette(task, &attempts, &diagnostics)?;
264 Ok(AuthorOutcome {
265 checked_form,
266 realized,
267 attempts,
268 diagnostics,
269 cassette,
270 })
271}
272
273fn author_cassette(
274 task: &AuthorTask,
275 attempts: &[RouteAttempt],
276 diagnostics: &[Diagnostic],
277) -> Result<DevCassette> {
278 let mut events = Vec::new();
279 if attempts.is_empty() {
280 events.push(DevEvent::refusal(
281 task.name.clone(),
282 Expr::Map(vec![
283 entry("target", Expr::String("none".to_owned())),
284 entry(
285 "reason",
286 Expr::String(
287 first_diagnostic(diagnostics)
288 .unwrap_or("no route attempt")
289 .to_owned(),
290 ),
291 ),
292 ]),
293 )?);
294 } else {
295 for attempt in attempts {
296 let kind = if matches!(attempt.status, RouteAttemptStatus::Accepted) {
297 CassetteKind::Validate
298 } else {
299 CassetteKind::Refusal
300 };
301 events.push(cassette_event(task, attempt, kind)?);
302 }
303 }
304 DevCassette::from_events(
305 Symbol::qualified("forge-author", task.name.as_qualified_str()),
306 events,
307 )
308}
309
310fn cassette_event(
311 task: &AuthorTask,
312 attempt: &RouteAttempt,
313 kind: CassetteKind,
314) -> Result<DevEvent> {
315 let payload = Expr::Map(vec![
316 entry("target", Expr::String(attempt.target.clone())),
317 entry("status", Expr::Symbol(route_status_symbol(&attempt.status))),
318 entry(
319 "reason",
320 attempt
321 .reason
322 .as_ref()
323 .map(|reason| Expr::String(reason.clone()))
324 .unwrap_or(Expr::Nil),
325 ),
326 ]);
327 match kind {
328 CassetteKind::Validate => DevEvent::validate(task.name.clone(), payload),
329 CassetteKind::Refusal => DevEvent::refusal(task.name.clone(), payload),
330 }
331}
332
333enum CassetteKind {
334 Validate,
335 Refusal,
336}
337
338fn target_skip_reason(target: &RouteTarget<'_>, ceiling: &[Symbol]) -> Option<String> {
339 target
340 .required_capabilities
341 .iter()
342 .find(|required| !capability_allowed(required, ceiling))
343 .map(|required| format!("target requires capability {required} outside projection ceiling"))
344}
345
346fn outside_ceiling_reason(required: &[Symbol], ceiling: &[Symbol]) -> Option<String> {
347 required
348 .iter()
349 .find(|required| !capability_allowed(required, ceiling))
350 .map(|required| format!("form requires capability {required} outside projection ceiling"))
351}
352
353fn capability_allowed(required: &Symbol, ceiling: &[Symbol]) -> bool {
354 let required = capability_name(required);
355 ceiling
356 .iter()
357 .map(capability_name)
358 .any(|allowed| allowed == required)
359}
360
361fn diminished_capabilities(current: &CapabilitySet, ceiling: &[Symbol]) -> CapabilitySet {
362 let allowed = ceiling
363 .iter()
364 .map(capability_name)
365 .fold(CapabilitySet::new(), CapabilitySet::grant);
366 current.intersect(&allowed)
367}
368
369fn capability_name(symbol: &Symbol) -> CapabilityName {
370 match symbol.namespace.as_deref() {
371 Some("capability") => CapabilityName::new(symbol.name.as_ref()),
372 _ => CapabilityName::new(symbol.to_string()),
373 }
374}
375
376fn required_capabilities_for_form(expr: &Expr) -> Vec<Symbol> {
377 let mut required = BTreeSet::new();
378 collect_required_capabilities(expr, &mut required);
379 required.into_iter().collect()
380}
381
382fn collect_required_capabilities(expr: &Expr, required: &mut BTreeSet<Symbol>) {
383 match expr {
384 Expr::Symbol(symbol) => {
385 if symbol.namespace.as_deref() == Some("capability") {
386 required.insert(symbol.clone());
387 }
388 }
389 Expr::List(items) | Expr::Vector(items) | Expr::Set(items) | Expr::Block(items) => {
390 for item in items {
391 collect_required_capabilities(item, required);
392 }
393 }
394 Expr::Map(entries) => {
395 for (key, value) in entries {
396 collect_required_capabilities(key, required);
397 collect_required_capabilities(value, required);
398 }
399 }
400 Expr::Call { operator, args } => {
401 collect_required_capabilities(operator, required);
402 for arg in args {
403 collect_required_capabilities(arg, required);
404 }
405 }
406 Expr::Infix {
407 operator: _,
408 left,
409 right,
410 } => {
411 collect_required_capabilities(left, required);
412 collect_required_capabilities(right, required);
413 }
414 Expr::Prefix { operator: _, arg } | Expr::Postfix { operator: _, arg } => {
415 collect_required_capabilities(arg, required);
416 }
417 Expr::Quote { expr, .. } | Expr::Annotated { expr, .. } => {
418 collect_required_capabilities(expr, required);
419 }
420 Expr::Extension { payload, .. } => collect_required_capabilities(payload, required),
421 Expr::Nil
422 | Expr::Bool(_)
423 | Expr::Number(_)
424 | Expr::Local(_)
425 | Expr::String(_)
426 | Expr::Bytes(_) => {}
427 }
428}
429
430fn author_intent(task: &AuthorTask) -> CompiledIntent {
431 CompiledIntent {
432 name: task.name.clone(),
433 verifiers: task.verifiers.clone(),
434 status: IntentStatus::Verified,
435 ..CompiledIntent::default()
436 }
437}
438
439fn shape_diagnostics(diagnostics: &[Diagnostic]) -> String {
440 if diagnostics.is_empty() {
441 "no diagnostics".to_owned()
442 } else {
443 diagnostics
444 .iter()
445 .map(|diagnostic| diagnostic.message.clone())
446 .collect::<Vec<_>>()
447 .join("; ")
448 }
449}
450
451fn first_diagnostic(diagnostics: &[Diagnostic]) -> Option<&str> {
452 diagnostics
453 .first()
454 .map(|diagnostic| diagnostic.message.as_str())
455}
456
457fn route_status_symbol(status: &RouteAttemptStatus) -> Symbol {
458 let name = match status {
459 RouteAttemptStatus::Skipped => "skipped",
460 RouteAttemptStatus::Failed => "failed",
461 RouteAttemptStatus::Accepted => "accepted",
462 };
463 Symbol::qualified("forge-route", name)
464}