1use std::sync::Arc;
2use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
3
4use crate::runtime::sandboxed_skill::scan_skill_dir;
5use crate::runtime::skill_watcher::SkillWatcher;
6use async_stream::try_stream;
7use deepstrike_core::governance::quota::ResourceQuota;
8use deepstrike_core::mm::memory::{
9 MemoryAuthor, MemoryKind, MemoryProvenance, MemoryQuery, MemoryRecall, MemoryRecord,
10 MemoryScope, MemoryTrustLevel, validate_memory_write,
11};
12use deepstrike_core::runtime::kernel::wire::{CancellationReason, MemoryPolicy};
13use deepstrike_core::runtime::kernel::{KernelObservation, KernelPressureAction};
14use deepstrike_core::runtime::session::SessionEvent;
15use deepstrike_core::scheduler::policy::SchedulerPolicyConfig;
16use deepstrike_core::types::message::{CoreMessage, ToolCall};
17use deepstrike_core::types::milestone::MilestoneCheckResult;
18use deepstrike_core::types::signal::{
19 RuntimeSignal as KernelSignal, SignalSource as KernelSignalSource,
20 SignalType as KernelSignalType, Urgency,
21};
22use deepstrike_core::types::task::RuntimeTask;
23use futures::StreamExt;
24
25use crate::governance::Governance;
26use crate::knowledge::KnowledgeSource;
27use crate::memory::MemoryStore;
28use crate::providers::{LLMProvider, StreamEvent};
29use crate::run_event::RunEvent;
30use crate::runtime::archive::ArchiveStore;
31use crate::runtime::canonical_kernel::CanonicalKernel;
32use crate::runtime::canonical_runner_runtime::{
33 CanonicalRunnerOptions, CanonicalRunnerRuntime, PersistPayloadFn, PersistedPayload,
34 canonical_kernel_action, canonical_kernel_apply,
35};
36use crate::runtime::execution_plane::{
37 ExecutionPlane, LocalExecutionPlane, PermissionRequest, PermissionRequestHandler,
38 PermissionResponse, RunContext, ToolSuspendHandler,
39};
40use crate::runtime::host_projection::{HostAction, HostEffect};
41use crate::runtime::os_profile::{
42 GovernancePolicy, OsProfile, SignalPolicy, assert_native_profile,
43};
44use crate::runtime::payload_store::{FilePayloadStore, PayloadStore};
45use crate::runtime::provider_replay::{peek_provider_replay, seed_provider_replay_from_events};
46use crate::runtime::replay::{
47 is_mid_run, replay_messages_with_cap, replay_messages_with_cap_and_loader,
48};
49use crate::runtime::session_log::{InMemorySessionLog, SessionEntry, SessionLog};
50use crate::runtime::{InMemoryKernelJournal, KernelJournal};
51use crate::{Error, Result};
52use crate::{SignalDeliveryReceipt, SignalSource};
53use deepstrike_core::context::task_state::TaskUpdate;
54use deepstrike_core::runtime::repair::repair_llm_completed;
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
60pub enum MilestonePolicy {
61 #[default]
63 RequireVerifier,
64 Terminate,
66 AutoPass,
69}
70
71#[derive(Debug, Clone)]
72pub struct MilestoneEvaluationContext {
73 pub phase_id: String,
74 pub criteria: Vec<String>,
75 pub required_evidence: Vec<String>,
76}
77
78pub type MilestoneEvaluationHandler = std::sync::Arc<
79 dyn Fn(
80 MilestoneEvaluationContext,
81 ) -> futures::future::BoxFuture<'static, Result<MilestoneCheckResult>>
82 + Send
83 + Sync,
84>;
85
86#[derive(Debug, Clone)]
91pub struct TurnMetrics {
92 pub turn: u32,
93 pub tools_exposed: usize,
94 pub tools_called: usize,
95 pub active_skill: Option<String>,
96 pub input_tokens: u32,
97 pub cache_read_tokens: u32,
98 pub cache_creation_tokens: u32,
99 pub cache_read_tokens_by_slot: Option<crate::providers::CacheReadBySlot>,
101}
102
103pub type OnTurnMetricsHandler = std::sync::Arc<dyn Fn(TurnMetrics) + Send + Sync>;
105
106#[derive(Debug, Clone, Default)]
108pub struct KernelReliability {
109 pub provider_recovery_attempts: Option<u8>,
110 pub output_recovery_attempts: Option<u8>,
111 pub max_input_bytes: Option<u32>,
112}
113
114pub struct RuntimeOptions {
116 pub provider: Box<dyn LLMProvider>,
117 pub artifact_set_digest: Option<String>,
119 pub execution_plane: Option<Box<dyn ExecutionPlane>>,
120 pub session_log: Option<Arc<dyn SessionLog>>,
122 pub compression_store: Option<Arc<dyn ArchiveStore>>,
123 pub payload_store: Option<Arc<dyn PayloadStore>>,
125 pub kernel_reliability: Option<KernelReliability>,
127 pub session_id: Option<String>,
129 pub max_tokens: u32,
130 pub max_turns: Option<u32>,
131 pub timeout_ms: Option<u64>,
132 pub extensions: Option<serde_json::Value>,
133 pub agent_id: Option<String>,
134 pub memory_scope: Option<MemoryScope>,
136 pub pre_query_memory: Option<std::sync::Arc<dyn Fn(&str) -> Vec<MemoryQuery> + Send + Sync>>,
142 pub system_prompt: Option<String>,
143 pub initial_memory: Vec<String>,
144 pub skill_dir: Option<std::path::PathBuf>,
145 pub memory_store: Option<Box<dyn MemoryStore>>,
146 pub knowledge_source: Option<Box<dyn KnowledgeSource>>,
147 pub signal_source: Option<Box<dyn SignalSource>>,
148 pub governance: Option<Arc<tokio::sync::Mutex<Governance>>>,
149 pub os_profile: Option<OsProfile>,
150 pub governance_policy: Option<GovernancePolicy>,
151 pub signal_policy: Option<SignalPolicy>,
152 pub scheduler_policy: Option<SchedulerPolicyConfig>,
153 pub resource_quota: Option<ResourceQuota>,
154 pub memory_policy: Option<MemoryPolicy>,
156 pub tokenizer: Option<String>,
157 pub enable_plan_tool: Option<bool>,
158 pub on_tool_suspend: Option<ToolSuspendHandler>,
159 pub on_permission_request: Option<PermissionRequestHandler>,
160 pub milestone_policy: MilestonePolicy,
162 pub milestone_contract: Option<deepstrike_core::types::milestone::MilestoneContract>,
163 pub run_spec: Option<deepstrike_core::types::agent::AgentRunSpec>,
164 pub allowed_tool_ids: Option<Vec<String>>,
173 pub baseline_tool_ids: Option<Vec<String>>,
179 pub on_turn_metrics: Option<OnTurnMetricsHandler>,
182 pub stable_core_tool_ids: Vec<String>,
185 pub on_milestone_evaluate: Option<MilestoneEvaluationHandler>,
186}
187
188fn build_run_spec(
193 explicit: Option<deepstrike_core::types::agent::AgentRunSpec>,
194 allowed_tool_ids: Option<&[String]>,
195 baseline_tool_ids: Option<&[String]>,
196 verification_contract_id: Option<&str>,
197 agent_id: Option<&str>,
198 session_id: &str,
199 goal: &str,
200) -> Option<deepstrike_core::types::agent::AgentRunSpec> {
201 use deepstrike_core::types::agent::{AgentIdentity, AgentRole, AgentRunSpec};
202 let profile = allowed_tool_ids.filter(|ids| !ids.is_empty());
203 let mut spec = match (explicit, profile) {
204 (Some(mut spec), Some(ids)) => {
205 spec.capability_filter.allowed_ids = ids.iter().map(|s| s.as_str().into()).collect();
206 Some(spec)
207 }
208 (Some(spec), None) => Some(spec),
209 (None, Some(ids)) => {
210 let mut spec = AgentRunSpec::new(
211 AgentIdentity::new(agent_id.unwrap_or("root"), session_id),
212 AgentRole::Custom,
213 goal.to_string(),
214 );
215 spec.capability_filter.allowed_ids = ids.iter().map(|s| s.as_str().into()).collect();
216 Some(spec)
217 }
218 (None, None) => Some(AgentRunSpec::new(
219 AgentIdentity::new(agent_id.unwrap_or("root"), session_id),
220 AgentRole::Custom,
221 goal.to_string(),
222 )),
223 };
224 if let Some(spec) = spec.as_mut() {
225 if let Some(baseline) = baseline_tool_ids {
226 spec.exposure_baseline = Some(baseline.iter().map(|s| s.as_str().into()).collect());
227 } else if spec.exposure_baseline.is_none() {
228 spec.exposure_baseline = Some(Vec::new());
229 }
230 }
231 if let (Some(spec), Some(contract_id)) = (spec.as_mut(), verification_contract_id)
232 && spec.verification_contract_id.is_none()
233 {
234 spec.verification_contract_id = Some(contract_id.into());
235 }
236 spec
237}
238
239fn utf8_prefix(value: &str, max_bytes: usize) -> &str {
240 let mut end = value.len().min(max_bytes);
241 while end > 0 && !value.is_char_boundary(end) {
242 end -= 1;
243 }
244 &value[..end]
245}
246
247pub struct RuntimeRunner {
249 opts: RuntimeOptions,
250 plane: Box<dyn ExecutionPlane>,
251 kernel_journal: Arc<dyn KernelJournal>,
252 interrupted: AtomicBool,
253 cancellation_reason: AtomicU8,
254 active_kernel:
255 std::sync::Mutex<Option<std::sync::Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>>>,
256 memory_write_timestamps: tokio::sync::Mutex<std::collections::VecDeque<u64>>,
257 local_page_out_cache: std::sync::Mutex<Vec<CoreMessage>>,
258}
259
260impl RuntimeRunner {
261 pub fn new(opts: RuntimeOptions) -> Self {
262 Self::new_with_kernel_journal(opts, Arc::new(InMemoryKernelJournal::new()))
263 }
264
265 pub fn new_with_kernel_journal(
267 mut opts: RuntimeOptions,
268 kernel_journal: Arc<dyn KernelJournal>,
269 ) -> Self {
270 if opts.session_log.is_none() {
271 opts.session_log = Some(Arc::new(InMemorySessionLog::new()));
272 }
273 if opts.payload_store.is_none() {
274 opts.payload_store = Some(Arc::new(FilePayloadStore::new(".payloads")));
275 }
276 let plane = opts
277 .execution_plane
278 .take()
279 .unwrap_or_else(|| Box::new(LocalExecutionPlane::new()));
280 Self {
281 opts,
282 plane,
283 kernel_journal,
284 interrupted: AtomicBool::new(false),
285 cancellation_reason: AtomicU8::new(0),
286 active_kernel: std::sync::Mutex::new(None),
287 memory_write_timestamps: tokio::sync::Mutex::new(std::collections::VecDeque::new()),
288 local_page_out_cache: std::sync::Mutex::new(Vec::new()),
289 }
290 }
291
292 pub fn interrupt(&self) {
293 self.interrupt_with_reason(CancellationReason::User);
294 }
295
296 pub fn interrupt_with_reason(&self, reason: CancellationReason) {
297 self.cancellation_reason
298 .store(cancellation_reason_code(reason), Ordering::Relaxed);
299 self.interrupted.store(true, Ordering::Relaxed);
300 }
301
302 pub fn execution_plane(&self) -> &dyn ExecutionPlane {
303 self.plane.as_ref()
304 }
305
306 pub async fn write_memory(
307 &self,
308 memory: MemoryRecord,
309 session_id: Option<&str>,
310 agent_id: Option<&str>,
311 ) -> Result<()> {
312 let Some(store) = &self.opts.memory_store else {
313 return Ok(());
314 };
315 let Some(agent_id) = agent_id.or(self.opts.agent_id.as_deref()) else {
316 return Ok(());
317 };
318
319 let turn = self.active_kernel_turn().await;
320 let validation = match self.opts.memory_policy.as_ref() {
321 Some(policy) if policy.validation_enabled == Some(false) => Ok(()),
322 Some(policy) => {
323 let mut validation = deepstrike_core::mm::memory::MemoryValidation::default();
324 if let Some(max_content_bytes) = policy.max_content_bytes {
325 validation.max_size_bytes = max_content_bytes;
326 }
327 if let Some(max_name_length) = policy.max_name_length {
328 validation.max_name_length = max_name_length as usize;
329 }
330 validation.validate(&memory)
331 }
332 None => validate_memory_write(&memory),
333 };
334 if let Err(error) = validation {
335 self.append_memory_syscall_observations(
336 session_id,
337 vec![KernelObservation::MemoryValidationFailed {
338 turn,
339 record_id: memory.record_id.clone(),
340 error: format!("{error:?}"),
341 }],
342 )
343 .await;
344 return Ok(());
345 }
346
347 let now_ms = std::time::SystemTime::now()
348 .duration_since(std::time::UNIX_EPOCH)
349 .unwrap_or_default()
350 .as_millis() as u64;
351 let write_limit = self
352 .opts
353 .resource_quota
354 .as_ref()
355 .and_then(|quota| quota.memory_writes_per_window);
356 let mut quota_guard = if write_limit.is_some() {
357 Some(self.memory_write_timestamps.lock().await)
358 } else {
359 None
360 };
361 if let (Some((max_writes, window_ms)), Some(timestamps)) =
362 (write_limit, quota_guard.as_mut())
363 {
364 let cutoff = now_ms.saturating_sub(window_ms);
365 while timestamps
366 .front()
367 .is_some_and(|timestamp| *timestamp < cutoff)
368 {
369 timestamps.pop_front();
370 }
371 if window_ms == 0 || timestamps.len() >= max_writes as usize {
372 drop(quota_guard);
373 self.append_memory_syscall_observations(
374 session_id,
375 vec![KernelObservation::MemoryValidationFailed {
376 turn,
377 record_id: memory.record_id.clone(),
378 error: format!(
379 "memory write quota exceeded: max {max_writes} writes per {window_ms}ms"
380 ),
381 }],
382 )
383 .await;
384 return Ok(());
385 }
386 }
387
388 store.put(agent_id, memory.clone()).await?;
389 if let Some(timestamps) = quota_guard.as_mut() {
390 timestamps.push_back(now_ms);
391 }
392 drop(quota_guard);
393 self.append_memory_syscall_observations(
394 session_id,
395 vec![KernelObservation::MemoryWritten {
396 turn,
397 record_id: memory.record_id,
398 scope: memory.scope,
399 memory_kind: memory.kind,
400 name: memory.name,
401 size_bytes: memory.content.len() as u32,
402 }],
403 )
404 .await;
405 Ok(())
406 }
407
408 pub async fn query_memory(
409 &self,
410 query: MemoryQuery,
411 session_id: Option<&str>,
412 agent_id: Option<&str>,
413 ) -> Result<Vec<MemoryRecall>> {
414 let Some(store) = &self.opts.memory_store else {
415 return Ok(Vec::new());
416 };
417 let Some(agent_id) = agent_id.or(self.opts.agent_id.as_deref()) else {
418 return Ok(Vec::new());
419 };
420
421 let turn = self.active_kernel_turn().await;
422 let mut canonical_query = query;
423 if let Some(top_k) = self
424 .opts
425 .memory_policy
426 .as_ref()
427 .and_then(|policy| policy.retrieval_top_k)
428 {
429 canonical_query.top_k = canonical_query.top_k.min(top_k as usize);
430 }
431 let hits = store.search(agent_id, &canonical_query).await?;
432 self.append_memory_syscall_observations(
433 session_id,
434 vec![KernelObservation::MemoryQueried {
435 turn,
436 scope: canonical_query.scope.clone(),
437 query: canonical_query.query.clone(),
438 requested_k: canonical_query.top_k,
439 requires_async_response: true,
440 }],
441 )
442 .await;
443 self.log_memory_retrieval_result(session_id, hits.clone())
444 .await;
445 Ok(hits)
446 }
447
448 async fn extract_session_memories(
449 &self,
450 session: &deepstrike_core::memory::durable::SessionData,
451 scope: &MemoryScope,
452 ) -> Result<Vec<MemoryRecord>> {
453 let transcript = session
454 .messages
455 .iter()
456 .map(|message| {
457 format!(
458 "[{:?}] {}",
459 message.role,
460 message.content.as_text().unwrap_or_default()
461 )
462 })
463 .collect::<Vec<_>>()
464 .join("\n")
465 .chars()
466 .take(8_000)
467 .collect::<String>();
468 let prompt = format!(
469 "{transcript}\n\nReturn {{\"memories\":[{{\"name\":\"stable-kebab-key\",\"kind\":\"user|feedback|project|reference\",\"content\":\"fact\",\"description\":\"why durable\",\"confidence\":0.0,\"links\":[],\"pinned\":false,\"ttl_days\":null,\"evidence_refs\":[]}}]}} with at most 10 items. Return {{\"memories\":[]}} when nothing is durable."
470 );
471 let context = rendered_context_from_messages(vec![
472 CoreMessage::system(
473 "Extract durable, reusable facts from this completed session. Return only JSON; do not include transient progress or guesses.",
474 ),
475 CoreMessage::user(prompt),
476 ]);
477 let state = self.opts.provider.create_run_state();
478 let mut stream = self
479 .opts
480 .provider
481 .stream(&context, &[], None, state.as_ref())
482 .await?;
483 let mut output = String::new();
484 while let Some(event) = stream.next().await {
485 if let StreamEvent::TextDelta { delta } = event? {
486 output.push_str(&delta);
487 }
488 }
489 Ok(crate::memory::parse_extracted_memories(
490 &output, session, scope,
491 ))
492 }
493
494 async fn log_memory_retrieval_result(&self, session_id: Option<&str>, hits: Vec<MemoryRecall>) {
495 let Some(session_id) = session_id.or(self.opts.session_id.as_deref()) else {
496 return;
497 };
498 self.log(session_id, SessionEvent::MemoryRetrievalResult { hits })
501 .await;
502 }
503
504 #[cfg(test)]
508 pub(crate) fn active_pending_effect_count(&self) -> Option<usize> {
509 self.active_kernel
510 .lock()
511 .unwrap()
512 .as_ref()
513 .and_then(|kernel| {
514 kernel
515 .try_lock()
516 .ok()
517 .map(|runtime| runtime.pending_effect_count())
518 })
519 }
520
521 async fn active_kernel_turn(&self) -> u32 {
522 let active = self.active_kernel.lock().unwrap().clone();
523 match active {
524 Some(kernel) => kernel.lock().await.turn(),
525 None => 0,
526 }
527 }
528
529 fn create_canonical_runtime(
530 &self,
531 operation_id: String,
532 session_id: &str,
533 ) -> Result<CanonicalRunnerRuntime> {
534 let provider_policy = self.opts.provider.runtime_policy();
535 let effective_max_turns = self
536 .opts
537 .max_turns
538 .or(provider_policy.max_turns)
539 .unwrap_or(25);
540 let effective_timeout = self.opts.timeout_ms.or(provider_policy.timeout_ms);
541 let payload_store = self
542 .opts
543 .payload_store
544 .clone()
545 .expect("runtime constructor installs a payload store");
546 let payload_session = session_id.to_string();
547 let persist_payload: PersistPayloadFn =
548 Arc::new(move |_call_id, content, preview_bytes| {
549 let payload_store = payload_store.clone();
550 let payload_session = payload_session.clone();
551 Box::pin(async move {
552 let digest = deepstrike_core::runtime::kernel::wire::canonical_digest(
553 content.as_bytes(),
554 )
555 .as_str()
556 .to_string();
557 let payload_ref = format!(
558 "payload:{}",
559 digest
560 .trim_start_matches("sha256:")
561 .chars()
562 .take(32)
563 .collect::<String>()
564 );
565 payload_store.persist(&payload_session, &payload_ref, &content)?;
566 Ok(PersistedPayload {
567 payload_ref,
568 digest,
569 original_size: content.len().to_string(),
570 preview: utf8_prefix(&content, preview_bytes).to_string(),
571 })
572 })
573 });
574 let mut runtime = CanonicalRunnerRuntime::new(
575 CanonicalKernel::default(),
576 self.kernel_journal.clone(),
577 operation_id,
578 CanonicalRunnerOptions {
579 max_context_tokens: self.opts.max_tokens,
580 max_turns: Some(effective_max_turns),
581 max_total_tokens: None,
582 max_wall_ms: effective_timeout,
583 artifact_set_digest: self.opts.artifact_set_digest.clone(),
584 memory_binding_id: self
585 .opts
586 .agent_id
587 .clone()
588 .unwrap_or_else(|| format!("memory:{session_id}")),
589 persist_payload: Some(persist_payload),
590 },
591 )?;
592 if let Some(contract) = self.opts.milestone_contract.as_ref() {
593 runtime.remember_milestone_contract(contract);
594 }
595 Ok(runtime)
596 }
597
598 async fn append_memory_syscall_observations(
599 &self,
600 session_id: Option<&str>,
601 observations: Vec<KernelObservation>,
602 ) {
603 let Some(session_id) = session_id.or(self.opts.session_id.as_deref()) else {
604 return;
605 };
606 for obs in observations {
607 match obs {
608 KernelObservation::MemoryWritten {
609 turn,
610 record_id,
611 scope,
612 memory_kind,
613 name,
614 size_bytes,
615 } => {
616 self.log(
617 session_id,
618 SessionEvent::MemoryWritten {
619 turn,
620 record_id,
621 scope,
622 memory_kind,
623 name,
624 size_bytes,
625 },
626 )
627 .await;
628 }
629 KernelObservation::MemoryQueried {
630 turn,
631 scope,
632 query,
633 requested_k,
634 requires_async_response,
635 } => {
636 self.log(
637 session_id,
638 SessionEvent::MemoryQueried {
639 turn,
640 scope,
641 query,
642 requested_k,
643 requires_async_response,
644 },
645 )
646 .await;
647 }
648 KernelObservation::MemoryValidationFailed {
649 turn,
650 record_id,
651 error,
652 } => {
653 self.log(
654 session_id,
655 SessionEvent::MemoryValidationFailed {
656 turn,
657 record_id,
658 error,
659 },
660 )
661 .await;
662 }
663 _ => {}
664 }
665 }
666 }
667
668 pub async fn execute(&self, goal: &str) -> Result<String> {
669 collect_text(self.run_streaming(goal, &[], None, None).await?).await
670 }
671
672 pub async fn execute_with_criteria(&self, goal: &str, criteria: &[String]) -> Result<String> {
673 collect_text(self.run_streaming(goal, criteria, None, None).await?).await
674 }
675
676 pub async fn run_streaming<'a>(
677 &'a self,
678 goal: &'a str,
679 criteria: &'a [String],
680 extensions: Option<&'a serde_json::Value>,
681 session_id: Option<&'a str>,
682 ) -> Result<std::pin::Pin<Box<dyn futures::Stream<Item = Result<RunEvent>> + 'a>>> {
683 self.run_streaming_with_attachments(goal, criteria, extensions, session_id, &[])
684 .await
685 }
686
687 pub async fn run_streaming_with_attachments<'a>(
690 &'a self,
691 goal: &'a str,
692 criteria: &'a [String],
693 extensions: Option<&'a serde_json::Value>,
694 session_id: Option<&'a str>,
695 attachments: &'a [deepstrike_core::types::message::ContentPart],
696 ) -> Result<std::pin::Pin<Box<dyn futures::Stream<Item = Result<RunEvent>> + 'a>>> {
697 let session_id = session_id
698 .map(str::to_string)
699 .or_else(|| self.opts.session_id.clone())
700 .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
701
702 let prior = self.read_entries(&session_id).await?;
703 let mut mid_run = is_mid_run(&prior);
704 if !mid_run {
705 if let Some(operation_id) = prior.iter().rev().find_map(|entry| match &entry.event {
706 SessionEvent::RunStarted { run_id, .. } => Some(run_id.clone()),
707 _ => None,
708 }) {
709 if self.kernel_journal.head(&operation_id).await?.is_some() {
710 let mut authoritative =
711 self.create_canonical_runtime(operation_id, &session_id)?;
712 authoritative.restore().await?;
713 mid_run = !authoritative.is_terminal();
714 }
715 }
716 }
717
718 let operation_id = if mid_run {
719 prior
720 .iter()
721 .rev()
722 .find_map(|entry| match &entry.event {
723 SessionEvent::RunStarted { run_id, .. } => Some(run_id.clone()),
724 _ => None,
725 })
726 .ok_or_else(|| {
727 Error::Other(format!(
728 "mid-run session has no run_started identity: {session_id}"
729 ))
730 })?
731 } else {
732 let run_id = uuid::Uuid::new_v4().to_string();
733 self.log(
734 &session_id,
735 SessionEvent::RunStarted {
736 run_id: run_id.clone(),
737 goal: goal.to_string(),
738 criteria: criteria.to_vec(),
739 agent_id: self.opts.agent_id.clone(),
740 system_prompt: self.opts.system_prompt.clone(),
741 attachments: attachments.to_vec(),
742 },
743 )
744 .await;
745 run_id
746 };
747
748 let goal_owned = goal.to_string();
749 let criteria_owned = criteria.to_vec();
750 let extensions_owned = extensions.cloned();
751 let attachments_owned = attachments.to_vec();
752 let prior_events = if prior.is_empty() { None } else { Some(prior) };
753
754 Ok(Box::pin(self.execute_inner(
755 session_id,
756 operation_id,
757 goal_owned,
758 criteria_owned,
759 extensions_owned,
760 prior_events,
761 mid_run,
762 attachments_owned,
763 )))
764 }
765
766 pub async fn wake_streaming(
767 &self,
768 session_id: &str,
769 extensions: Option<&serde_json::Value>,
770 ) -> Result<std::pin::Pin<Box<dyn futures::Stream<Item = Result<RunEvent>> + '_>>> {
771 let prior = self.read_entries(session_id).await?;
772 let (start_index, start) = prior
773 .iter()
774 .enumerate()
775 .rev()
776 .find(|(_, entry)| matches!(entry.event, SessionEvent::RunStarted { .. }))
777 .ok_or_else(|| Error::Other(format!("no run_started for session: {session_id}")))?;
778 let (operation_id, goal, criteria, attachments) = match &start.event {
779 SessionEvent::RunStarted {
780 run_id,
781 goal,
782 criteria,
783 attachments,
784 ..
785 } => (
786 run_id.clone(),
787 goal.clone(),
788 criteria.clone(),
789 attachments.clone(),
790 ),
791 _ => unreachable!(),
792 };
793
794 if prior[start_index + 1..]
795 .iter()
796 .any(|entry| matches!(entry.event, SessionEvent::RunTerminal { .. }))
797 {
798 if self.kernel_journal.head(&operation_id).await?.is_none() {
799 return Err(Error::Other(
800 "run_terminal projection has no canonical journal".into(),
801 ));
802 }
803 let mut authoritative =
804 self.create_canonical_runtime(operation_id.clone(), session_id)?;
805 authoritative.restore().await?;
806 if authoritative.is_terminal() {
807 return Ok(Box::pin(futures::stream::empty()));
808 }
809 }
810
811 Ok(Box::pin(self.execute_inner(
812 session_id.to_string(),
813 operation_id,
814 goal,
815 criteria,
816 extensions.cloned(),
817 Some(prior),
818 true,
819 attachments,
820 )))
821 }
822
823 pub async fn wake(&self, session_id: &str) -> Result<String> {
824 collect_text(self.wake_streaming(session_id, None).await?).await
825 }
826
827 fn execute_inner(
828 &self,
829 session_id: String,
830 operation_id: String,
831 goal: String,
832 criteria: Vec<String>,
833 extensions: Option<serde_json::Value>,
834 prior_events: Option<Vec<SessionEntry>>,
835 resume_mid_run: bool,
836 attachments: Vec<deepstrike_core::types::message::ContentPart>,
837 ) -> impl futures::Stream<Item = Result<RunEvent>> + '_ {
838 try_stream! {
839 self.interrupted.store(false, Ordering::Relaxed);
840 self.cancellation_reason.store(0, Ordering::Relaxed);
841
842 if let Some(ks) = &self.opts.knowledge_source {
843 ks.init().await?;
844 }
845
846 let mut runtime = self.create_canonical_runtime(operation_id, &session_id)?;
847 if resume_mid_run {
848 runtime.restore().await?;
849 }
850 let kernel = std::sync::Arc::new(tokio::sync::Mutex::new(runtime));
851 {
852 let mut active = self.active_kernel.lock().unwrap();
853 *active = Some(kernel.clone());
854 }
855
856 struct ActiveKernelGuard<'a> {
857 runner: &'a RuntimeRunner,
858 }
859 impl<'a> Drop for ActiveKernelGuard<'a> {
860 fn drop(&mut self) {
861 if let Ok(mut active) = self.runner.active_kernel.lock() {
862 *active = None;
863 }
864 }
865 }
866 let _guard = ActiveKernelGuard { runner: self };
867
868 let mut pending_observations = Vec::new();
869 let mut pending_page_out_starts = std::collections::VecDeque::new();
870 let mut active_page_out_start = None;
871 let skill_watcher = self.opts.skill_dir.as_deref().and_then(SkillWatcher::start);
872
873 if !resume_mid_run {
874 if self.opts.kernel_reliability.is_some()
875 || self.opts.scheduler_policy.is_some()
876 {
877 let mut config = serde_json::Map::new();
878 if let Some(reliability) = self.opts.kernel_reliability.as_ref() {
879 config.insert(
880 "reliability".into(),
881 serde_json::json!({
882 "provider_recovery_attempts": reliability.provider_recovery_attempts,
883 "output_recovery_attempts": reliability.output_recovery_attempts,
884 "max_input_bytes": reliability.max_input_bytes,
885 }),
886 );
887 }
888 if let Some(policy) = self.opts.scheduler_policy {
889 let policy = serde_json::to_value(policy).map_err(|error| {
890 Error::Other(format!("scheduler policy is not serializable: {error}"))
891 })?;
892 config.insert("scheduler_policy".into(), policy);
893 }
894 kernel_apply(
895 &kernel,
896 &mut pending_observations,
897 serde_json::json!({
898 "kind": "configure_run",
899 "config": config,
900 }),
901 )
902 .await?;
903 }
904
905 if let Some(tokenizer_name) = &self.opts.tokenizer {
906 kernel_apply(
907 &kernel,
908 &mut pending_observations,
909 serde_json::json!({ "kind": "set_tokenizer", "name": tokenizer_name }),
910 ).await?;
911 }
912 if let Some(enabled) = self.opts.enable_plan_tool {
913 kernel_apply(
914 &kernel,
915 &mut pending_observations,
916 serde_json::json!({ "kind": "set_plan_tool_enabled", "enabled": enabled }),
917 ).await?;
918 }
919
920 kernel_apply(
921 &kernel,
922 &mut pending_observations,
923 serde_json::json!({ "kind": "set_tools", "tools": self.plane.schemas() }),
924 ).await?;
925
926 if self.opts.memory_store.is_some() && self.opts.agent_id.is_some() {
927 kernel_apply(
928 &kernel,
929 &mut pending_observations,
930 serde_json::json!({ "kind": "set_memory_enabled", "enabled": true }),
931 ).await?;
932 }
933 if self.opts.knowledge_source.is_some() {
934 kernel_apply(
935 &kernel,
936 &mut pending_observations,
937 serde_json::json!({ "kind": "set_knowledge_enabled", "enabled": true }),
938 ).await?;
939 }
940
941 if let Some(sp) = &self.opts.system_prompt {
942 let tokens = ((sp.len() / 4) as u32).max(1);
943 kernel_apply(
944 &kernel,
945 &mut pending_observations,
946 serde_json::json!({
947 "kind": "add_system_message",
948 "content": sp,
949 "tokens": tokens,
950 }),
951 ).await?;
952 }
953 for mem in &self.opts.initial_memory {
954 let tokens = ((mem.len() / 4) as u32).max(1);
955 kernel_apply(
956 &kernel,
957 &mut pending_observations,
958 serde_json::json!({
959 "kind": "add_knowledge_message",
960 "content": mem,
961 "tokens": tokens,
962 "pinned": false,
963 }),
964 ).await?;
965 }
966
967 if let Some(skill_dir) = &self.opts.skill_dir {
968 kernel_apply(
969 &kernel,
970 &mut pending_observations,
971 serde_json::json!({
972 "kind": "set_available_skills",
973 "skills": scan_skill_dir(skill_dir),
974 }),
975 ).await?;
976 }
977
978 if !self.opts.stable_core_tool_ids.is_empty() {
980 kernel_apply(
981 &kernel,
982 &mut pending_observations,
983 serde_json::json!({
984 "kind": "set_stable_core_tools",
985 "tool_ids": self.opts.stable_core_tool_ids,
986 }),
987 ).await?;
988 }
989
990 if let Some(milestones) = self.opts.milestone_contract.clone() {
991 kernel_apply(
992 &kernel,
993 &mut pending_observations,
994 serde_json::json!({
995 "kind": "load_milestone_contract",
996 "contract": milestones,
997 }),
998 ).await?;
999 }
1000
1001 let max_bytes = {
1002 let k = kernel.lock().await;
1003 k.recovery_content_bytes()
1004 };
1005
1006 if let Some(ref events) = prior_events {
1007 seed_provider_replay_from_events(self.opts.provider.as_ref(), events);
1008
1009 let messages = if let Some(ref store) = self.opts.compression_store {
1010 let store_clone = store.clone();
1011 replay_messages_with_cap_and_loader(events, max_bytes, move |archive_ref| {
1012 store_clone.read(archive_ref).map_err(|_| {
1013 deepstrike_core::context::fault::ContextFault::MissingArchive {
1014 session_id: String::new(),
1015 seq: 0,
1016 }
1017 })
1018 })
1019 } else {
1020 replay_messages_with_cap(events, max_bytes)
1021 };
1022
1023 kernel_apply(
1024 &kernel,
1025 &mut pending_observations,
1026 serde_json::json!({ "kind": "preload_history", "messages": messages }),
1027 ).await?;
1028 }
1029 } else if let Some(ref events) = prior_events {
1030 seed_provider_replay_from_events(self.opts.provider.as_ref(), events);
1031 }
1032
1033 let ext = merge_extensions(self.opts.extensions.as_ref(), extensions.as_ref());
1034 let provider_state = self.opts.provider.create_run_state();
1035 let mut next_archive_start = next_archived_seq_start(prior_events.as_deref());
1036 let mut active_skill: Option<String> = None;
1038 let session_start_ms = std::time::SystemTime::now()
1039 .duration_since(std::time::UNIX_EPOCH)
1040 .unwrap_or_default()
1041 .as_millis() as u64;
1042
1043 if !resume_mid_run {
1044 let os_profile = assert_native_profile(self.opts.os_profile.clone())?;
1045 let governance_policy = self
1046 .opts
1047 .governance_policy
1048 .clone()
1049 .unwrap_or(os_profile.governance_policy);
1050 kernel_apply(
1051 &kernel,
1052 &mut pending_observations,
1053 governance_policy.into_host_fact(),
1054 ).await?;
1055
1056 let signal_policy = self
1057 .opts
1058 .signal_policy
1059 .unwrap_or(os_profile.signal_policy);
1060 kernel_apply(
1061 &kernel,
1062 &mut pending_observations,
1063 serde_json::json!({
1064 "kind": "set_signal_policy",
1065 "policy": signal_policy.into_kernel(),
1066 }),
1067 ).await?;
1068
1069 if let Some(quota) = self.opts.resource_quota.clone() {
1070 kernel_apply(
1071 &kernel,
1072 &mut pending_observations,
1073 serde_json::json!({ "kind": "set_resource_quota", "quota": quota }),
1074 ).await?;
1075 }
1076
1077 if let Some(policy) = self.opts.memory_policy.clone() {
1078 kernel_apply(
1079 &kernel,
1080 &mut pending_observations,
1081 memory_policy_host_fact(policy),
1082 ).await?;
1083 }
1084
1085 if !resume_mid_run && !attachments.is_empty() {
1087 kernel_apply(
1088 &kernel,
1089 &mut pending_observations,
1090 serde_json::json!({
1091 "kind": "add_history_message",
1092 "message": CoreMessage::user_multimodal(attachments.clone()),
1093 }),
1094 ).await?;
1095 }
1096
1097 if !resume_mid_run {
1101 if let (Some(pre), Some(store), Some(agent_id)) = (
1102 self.opts.pre_query_memory.clone(),
1103 self.opts.memory_store.as_ref(),
1104 self.opts.agent_id.as_deref(),
1105 ) {
1106 let queries = pre(goal.as_str());
1107 let mut recalled = Vec::new();
1108 for q in &queries {
1109 if q.query.trim().is_empty() {
1110 continue;
1111 }
1112 if let Ok(hits) = store.search(agent_id, q).await {
1113 for hit in hits {
1114 recalled.push(format!(
1115 "[memory record_id={} trust={} score={:.3}] {}",
1116 hit.record.record_id,
1117 match hit.record.provenance.trust {
1118 MemoryTrustLevel::Untrusted => "untrusted",
1119 MemoryTrustLevel::UserAsserted => "user_asserted",
1120 MemoryTrustLevel::HostVerified => "host_verified",
1121 },
1122 hit.score,
1123 hit.record.content
1124 ));
1125 }
1126 }
1127 }
1128 if !recalled.is_empty() {
1129 kernel_apply(
1130 &kernel,
1131 &mut pending_observations,
1132 serde_json::json!({
1133 "kind": "add_history_message",
1134 "message": CoreMessage::user(recalled.join("\n")),
1135 }),
1136 ).await?;
1137 }
1138 }
1139 }
1140 }
1141
1142 let mut action = if resume_mid_run {
1143 let mut runtime = kernel.lock().await;
1144 let action = runtime.resume_action()?.ok_or_else(|| {
1145 Error::Other(
1146 "restored canonical operation has no pending effect or terminal".into(),
1147 )
1148 })?;
1149 pending_observations.extend(runtime.drain_host_observations());
1150 action
1151 } else {
1152 let run_spec = build_run_spec(
1156 self.opts.run_spec.clone(),
1157 self.opts.allowed_tool_ids.as_deref(),
1158 self.opts.baseline_tool_ids.as_deref(),
1159 self.opts
1160 .milestone_contract
1161 .as_ref()
1162 .map(|_| "rust-default"),
1163 self.opts.agent_id.as_deref(),
1164 &session_id,
1165 &goal,
1166 );
1167 kernel_start_agent(
1168 &kernel,
1169 &mut pending_observations,
1170 RuntimeTask::new(&goal).with_criteria(criteria),
1171 run_spec,
1172 ).await?
1173 };
1174
1175 let mut last_skill_version: u64 = skill_watcher.as_ref().map(|w| w.version()).unwrap_or(0);
1176
1177 while !kernel.lock().await.is_terminal() {
1178 if let (Some(watcher), Some(skill_dir)) =
1180 (&skill_watcher, &self.opts.skill_dir)
1181 {
1182 let cur = watcher.version();
1183 if cur != last_skill_version {
1184 last_skill_version = cur;
1185 kernel_apply(
1186 &kernel,
1187 &mut pending_observations,
1188 serde_json::json!({
1189 "kind": "set_available_skills",
1190 "skills": scan_skill_dir(skill_dir),
1191 }),
1192 ).await?;
1193 }
1194 }
1195
1196 next_archive_start = self
1197 .append_observations(
1198 &session_id,
1199 &kernel,
1200 &mut pending_observations,
1201 &mut pending_page_out_starts,
1202 next_archive_start,
1203 )
1204 .await;
1205
1206 if self.interrupted.load(Ordering::Relaxed) {
1207 let operation_id = kernel.lock().await.operation_id().to_string();
1208 kernel_apply(
1209 &kernel,
1210 &mut pending_observations,
1211 serde_json::json!({
1212 "kind": "cancel_operation",
1213 "operation_id": operation_id,
1214 "reason": cancellation_reason_from_code(self.cancellation_reason.load(Ordering::Relaxed)),
1215 "pending_call_ids": pending_call_ids(&action),
1216 }),
1217 ).await?;
1218 break;
1219 }
1220
1221 if let Some(ss) = &self.opts.signal_source {
1222 if let Some(claim) = ss.claim_signal().await? {
1223 let urgency = match claim.signal.urgency.as_str() {
1224 "low" => Urgency::Low,
1225 "high" => Urgency::High,
1226 "critical" => Urgency::Critical,
1227 _ => Urgency::Normal,
1228 };
1229 let source = match claim.signal.source.as_str() {
1230 "cron" => KernelSignalSource::Cron,
1231 "gateway" => KernelSignalSource::Gateway,
1232 "heartbeat" => KernelSignalSource::Heartbeat,
1233 _ => KernelSignalSource::Custom,
1234 };
1235 let signal_type = match claim.signal.signal_type.as_str() {
1236 "job" => KernelSignalType::Job,
1237 "alert" => KernelSignalType::Alert,
1238 _ => KernelSignalType::Event,
1239 };
1240 let summary = claim
1241 .signal
1242 .payload
1243 .get("goal")
1244 .and_then(serde_json::Value::as_str)
1245 .unwrap_or("signal");
1246 let mut kernel_sig = KernelSignal::new(
1247 source,
1248 signal_type,
1249 urgency,
1250 summary,
1251 )
1252 .with_payload(claim.signal.payload.clone())
1253 .with_timestamp(
1254 std::time::SystemTime::now()
1255 .duration_since(std::time::UNIX_EPOCH)
1256 .unwrap_or_default()
1257 .as_millis() as u64,
1258 );
1259 kernel_sig.id = claim.signal_id.as_str().into();
1262 if let Some(dedupe_key) = &claim.signal.dedupe_key {
1263 kernel_sig = kernel_sig.with_dedupe(dedupe_key.clone());
1264 }
1265 if let Some(recipient) = &claim.signal.recipient {
1266 kernel_sig = kernel_sig.with_recipient(recipient.clone());
1267 }
1268 if let Some(deadline_ms) = claim.signal.deadline_ms {
1269 kernel_sig = kernel_sig.with_deadline(deadline_ms);
1270 }
1271 if let Some(coalesce_key) = &claim.signal.coalesce_key {
1272 kernel_sig = kernel_sig.with_coalesce(coalesce_key.clone());
1273 }
1274 kernel_sig.coalesced_count = claim.signal.coalesced_count.max(1);
1275 let observation_start = pending_observations.len();
1280 let signal_action = kernel_transition(
1281 &kernel,
1282 &mut pending_observations,
1283 serde_json::json!({
1284 "kind": "deliver_signal",
1285 "delivery_id": claim.delivery_id,
1286 "attempt": claim.delivery_attempt,
1287 "signal": kernel_sig,
1288 }),
1289 )
1290 .await;
1291 let receipt = SignalDeliveryReceipt {
1292 delivery_id: claim.delivery_id.clone(),
1293 lease_token: claim.lease_token.clone(),
1294 };
1295 let signal_action = match signal_action {
1296 Ok(action) => action,
1297 Err(error) => {
1298 let _ = ss.nack_signal(&receipt).await?;
1299 Err(error)?
1300 }
1301 };
1302 let disposition_matches = pending_observations[observation_start..]
1303 .iter()
1304 .filter(|observation| match observation {
1305 KernelObservation::SignalDeliveryDisposed {
1306 delivery_id,
1307 attempt,
1308 ..
1309 } => {
1310 delivery_id == &claim.delivery_id
1311 && attempt == &claim.delivery_attempt
1312 }
1313 _ => false,
1314 })
1315 .count()
1316 == 1;
1317 if !disposition_matches {
1318 let _ = ss.nack_signal(&receipt).await?;
1319 Err(crate::Error::Other(
1320 "kernel did not return the matching signal delivery disposition".into(),
1321 ))?;
1322 }
1323 if !ss.ack_signal(&receipt).await? {
1324 let _ = ss.nack_signal(&receipt).await?;
1325 Err(crate::Error::Other(
1326 "signal lease was lost before acknowledgement".into(),
1327 ))?;
1328 }
1329 if let Some(sig_action) = signal_action {
1330 action = sig_action;
1331 }
1332 }
1334 }
1335 if kernel.lock().await.is_terminal() {
1336 break;
1337 }
1338
1339 match &action.effect {
1340 HostEffect::CallProvider { context, tools, context_effect } => {
1341 let provider_effect_id = action.effect_id.clone();
1342 let mut final_text = String::new();
1343 let mut final_tool_calls: Vec<ToolCall> = Vec::new();
1344 let mut turn_tokens: u32 = 0;
1345 let mut turn_input_tokens: u32 = 0;
1346 let mut turn_cache_read_tokens: u32 = 0;
1347 let mut turn_cache_creation_tokens: u32 = 0;
1348 let mut turn_cache_read_by_slot: Option<crate::providers::CacheReadBySlot> = None;
1349 let mut turn_stop_reason: Option<String> = None;
1350 let (filtered_tools, filtered_context_storage);
1354 let (provider_tools, provider_context): (&[_], &_) = if let Some(policy) = self.opts.governance_policy.as_ref() {
1355 if policy.surface_denied_in_system {
1356 let (allowed, denied) = crate::runtime::governance_filter_schema(tools, policy);
1357 if !denied.is_empty() {
1358 filtered_tools = allowed;
1359 let mut cloned = context.clone();
1360 let note = format!("[governance] the following tools are denied for this run and will fail if called: {}.", denied.join(", "));
1361 cloned.system_knowledge = if cloned.system_knowledge.is_empty() {
1362 note
1363 } else {
1364 format!("{}\n\n{}", cloned.system_knowledge, note)
1365 };
1366 filtered_context_storage = cloned;
1367 (&filtered_tools[..], &filtered_context_storage)
1368 } else {
1369 (&tools[..], context)
1370 }
1371 } else { (&tools[..], context) }
1372 } else { (&tools[..], context) };
1373 let tools_exposed = provider_tools.len();
1376
1377 let mut route = self.opts.provider.context_route();
1380 let prepared_request = self.opts.provider.prepare_context_request(
1381 provider_context, provider_tools, ext.as_ref(), provider_state.as_ref(),
1382 )?;
1383 route.as_object_mut().ok_or_else(|| Error::Other("provider context route must be an object".into()))?
1384 .insert("request_fingerprint_scope".into(), prepared_request.get("scope")
1385 .cloned().unwrap_or_else(|| serde_json::json!("adapter_input")));
1386 let material = serde_json::json!({
1387 "route": route, "request": prepared_request,
1388 "context": provider_context, "tools": provider_tools,
1389 "options": ext, "provider_state": provider_state,
1390 });
1391 let material_bytes = deepstrike_core::runtime::kernel::wire::record::canonical_bytes(&material)
1392 .map_err(|error| Error::Other(error.to_string()))?;
1393 let fingerprint = deepstrike_core::evolution::ContentDigest::from_bytes(material_bytes.as_slice());
1394 let counted_material = deepstrike_core::runtime::kernel::wire::record::canonical_bytes(
1396 prepared_request.get("body").unwrap_or(&prepared_request),
1397 ).map_err(|error| Error::Other(error.to_string()))?;
1398 let estimated = counted_material.as_slice().len().div_ceil(4) as u64;
1399 let preparation = deepstrike_core::context::execution::prepare_context_dispatch(
1400 &deepstrike_core::context::execution::ContextDispatchRequest {
1401 effect: context_effect.clone(), request_fingerprint: fingerprint.clone(),
1402 provider_route: route,
1403 prompt_measurement: deepstrike_core::context::execution::ContextPromptMeasurement {
1404 request_fingerprint: fingerprint, input_tokens: estimated,
1405 source: deepstrike_core::context::measurement::MeasurementSource::Heuristic,
1406 confidence: deepstrike_core::context::measurement::MeasurementConfidence::LowConfidence,
1407 },
1408 },
1409 ).map_err(|error| Error::Other(error.to_string()))?;
1410 if let Some(log) = &self.opts.session_log {
1411 log.append(&session_id, SessionEvent::ContextPrepared {
1412 turn: kernel.lock().await.turn(),
1413 effect_id: provider_effect_id.clone(), preparation: Box::new(preparation),
1414 }).await.map_err(Error::Io)?;
1415 }
1416
1417 let mut provider_stream = match self
1418 .opts
1419 .provider
1420 .stream_prepared(&prepared_request, provider_context, provider_tools, ext.as_ref(), provider_state.as_ref())
1421 .await
1422 {
1423 Ok(s) => s,
1424 Err(e) => {
1425 let msg = provider_error_message(&e);
1432 action = kernel_action(
1433 &kernel,
1434 &mut pending_observations,
1435 provider_error_event(&provider_effect_id, &e),
1436 ).await?;
1437 if matches!(&action.effect, HostEffect::Done { .. }) {
1444 yield RunEvent::Error(msg);
1445 }
1446 continue;
1447 }
1448 };
1449
1450 let mut stream_error: Option<crate::Error> = None;
1456 while let Some(evt) = provider_stream.next().await {
1457 if self.interrupted.load(Ordering::Relaxed) {
1458 break;
1459 }
1460 let evt = match evt {
1461 Ok(evt) => evt,
1462 Err(e) => {
1463 stream_error = Some(e);
1464 break;
1465 }
1466 };
1467 match evt {
1468 StreamEvent::TextDelta { delta } => {
1469 final_text.push_str(&delta);
1470 yield RunEvent::TextDelta(delta);
1471 }
1472 StreamEvent::ThinkingDelta { delta } => {
1473 yield RunEvent::ThinkingDelta(delta);
1474 }
1475 StreamEvent::ToolCall { id, name, arguments } => {
1476 yield RunEvent::ToolCall { id: id.clone(), name: name.clone() };
1477 final_tool_calls.push(ToolCall {
1478 id: compact_str::CompactString::new(&id),
1479 name: compact_str::CompactString::new(&name),
1480 arguments,
1481 });
1482 }
1483 StreamEvent::Usage {
1484 total_tokens,
1485 input_tokens,
1486 cache_read_input_tokens,
1487 cache_creation_input_tokens,
1488 cache_read_input_tokens_by_slot,
1489 stop_reason,
1490 ..
1491 } => {
1492 turn_tokens = total_tokens;
1493 turn_input_tokens = input_tokens;
1495 turn_cache_read_tokens = cache_read_input_tokens;
1496 turn_cache_creation_tokens = cache_creation_input_tokens;
1497 turn_cache_read_by_slot = cache_read_input_tokens_by_slot;
1498 if stop_reason.is_some() { turn_stop_reason = stop_reason; }
1500 }
1501 StreamEvent::Done => {}
1502 }
1503 }
1504
1505 if self.interrupted.load(Ordering::Relaxed) {
1506 let operation_id = kernel.lock().await.operation_id().to_string();
1507 action = kernel_action(
1508 &kernel,
1509 &mut pending_observations,
1510 serde_json::json!({
1511 "kind": "cancel_operation",
1512 "operation_id": operation_id,
1513 "reason": cancellation_reason_from_code(self.cancellation_reason.load(Ordering::Relaxed)),
1514 "pending_call_ids": [provider_effect_id],
1515 }),
1516 ).await?;
1517 break;
1518 }
1519
1520 if let Some(error) = stream_error {
1521 let msg = provider_error_message(&error);
1522 action = kernel_action(
1528 &kernel,
1529 &mut pending_observations,
1530 provider_error_event(&provider_effect_id, &error),
1531 ).await?;
1532 if matches!(&action.effect, HostEffect::Done { .. }) {
1533 yield RunEvent::Error(msg);
1534 }
1535 continue;
1536 }
1537
1538 let mut assistant = CoreMessage {
1539 role: deepstrike_core::types::message::Role::Assistant,
1540 content: deepstrike_core::types::message::Content::Text(final_text.clone()),
1541 tool_calls: final_tool_calls.clone(),
1542 };
1543
1544 self.opts.provider.commit_stream_replay(&final_text, &final_tool_calls);
1545 let mut provider_replay = peek_provider_replay(
1546 self.opts.provider.as_ref(),
1547 &final_text,
1548 &final_tool_calls,
1549 );
1550 repair_llm_completed(&mut assistant, &mut provider_replay);
1551
1552 action = kernel_action(
1553 &kernel,
1554 &mut pending_observations,
1555 serde_json::json!({
1556 "kind": "provider_result",
1557 "effect_id": provider_effect_id,
1558 "message": assistant,
1559 "stop_reason": turn_stop_reason,
1561 }),
1562 ).await?;
1563 self.log(
1564 &session_id,
1565 SessionEvent::LlmCompleted {
1566 turn: kernel.lock().await.turn(),
1567 message: assistant,
1568 provider_replay,
1569 },
1570 )
1571 .await;
1572
1573 if let Some(ref sink) = self.opts.on_turn_metrics {
1577 sink(TurnMetrics {
1578 turn: kernel.lock().await.turn(),
1579 tools_exposed,
1580 tools_called: final_tool_calls.len(),
1581 active_skill: active_skill.clone(),
1582 input_tokens: turn_input_tokens,
1583 cache_read_tokens: turn_cache_read_tokens,
1584 cache_creation_tokens: turn_cache_creation_tokens,
1585 cache_read_tokens_by_slot: turn_cache_read_by_slot.clone(),
1586 });
1587 }
1588 if let Some(skill_call) =
1589 final_tool_calls.iter().find(|c| c.name.as_str() == "skill")
1590 {
1591 if let Some(name) = skill_call.arguments.get("name").and_then(|v| v.as_str()) {
1592 active_skill = Some(name.to_string());
1593 }
1594 }
1595 }
1596 HostEffect::RequestApproval { requests } => {
1597 let approval_effect_id = action.effect_id.clone();
1598 let mut approved_calls = Vec::new();
1599 let mut denied_calls = Vec::new();
1600 for request in requests {
1601 let arguments = request.arguments.to_string();
1602 self.log(
1603 &session_id,
1604 SessionEvent::PermissionRequested {
1605 turn: kernel.lock().await.turn(),
1606 tool: request.tool.clone(),
1607 arguments: arguments.clone(),
1608 reason: Some(request.reason.clone()),
1609 },
1610 )
1611 .await;
1612 yield RunEvent::PermissionRequest {
1613 call_id: request.call_id.clone(),
1614 tool_name: request.tool.clone(),
1615 arguments: arguments.clone(),
1616 reason: request.reason.clone(),
1617 };
1618
1619 let response = match &self.opts.on_permission_request {
1620 Some(handler) => match handler(PermissionRequest {
1621 call_id: request.call_id.clone(),
1622 tool_name: request.tool.clone(),
1623 arguments,
1624 reason: request.reason.clone(),
1625 })
1626 .await
1627 {
1628 Ok(response) => response,
1629 Err(err) => PermissionResponse {
1630 approved: false,
1631 responder: "permission_handler".to_string(),
1632 reason: Some(format!("permission handler failed: {err}")),
1633 },
1634 },
1635 None => PermissionResponse {
1636 approved: false,
1637 responder: "policy_gate".to_string(),
1638 reason: Some("no permission handler configured".to_string()),
1639 },
1640 };
1641 if response.approved {
1642 approved_calls.push(request.call_id.clone());
1643 } else {
1644 denied_calls.push(request.call_id.clone());
1645 }
1646 let responder = if response.responder.is_empty() {
1647 "host".to_string()
1648 } else {
1649 response.responder
1650 };
1651 self.log(
1652 &session_id,
1653 SessionEvent::PermissionResolved {
1654 turn: kernel.lock().await.turn(),
1655 approved: response.approved,
1656 responder: responder.clone(),
1657 },
1658 )
1659 .await;
1660 yield RunEvent::PermissionResolved {
1661 call_id: request.call_id.clone(),
1662 tool_name: request.tool.clone(),
1663 approved: response.approved,
1664 responder,
1665 reason: response.reason,
1666 };
1667 }
1668 action = kernel_action(
1669 &kernel,
1670 &mut pending_observations,
1671 serde_json::json!({
1672 "kind": "approval_result",
1673 "effect_id": approval_effect_id,
1674 "approved_calls": approved_calls,
1675 "denied_calls": denied_calls,
1676 }),
1677 ).await?;
1678 }
1679 HostEffect::SpawnWorkflow { nodes, .. } => {
1680 let workflow_effect_id = action.effect_id.clone();
1684 let failures: Vec<deepstrike_core::runtime::kernel::WorkflowSpawnFailure> = nodes
1685 .into_iter()
1686 .map(|node| deepstrike_core::runtime::kernel::WorkflowSpawnFailure {
1687 agent_id: node.agent_id.clone(),
1688 error: "Rust RuntimeRunner has no workflow orchestrator".to_string(),
1689 })
1690 .collect();
1691 action = kernel_action(
1692 &kernel,
1693 &mut pending_observations,
1694 serde_json::json!({
1695 "kind": "workflow_spawn_result",
1696 "effect_id": workflow_effect_id,
1697 "started_agent_ids": [],
1698 "failures": failures,
1699 }),
1700 ).await?;
1701 }
1702 HostEffect::PreemptSubAgents { .. } => {
1703 let preempt_effect_id = action.effect_id.clone();
1706 action = kernel_action(
1707 &kernel,
1708 &mut pending_observations,
1709 serde_json::json!({
1710 "kind": "preempt_result",
1711 "effect_id": preempt_effect_id,
1712 }),
1713 ).await?;
1714 }
1715 HostEffect::PersistMemory { memory } => {
1716 let effect_id = action.effect_id.clone();
1717 let error = match (
1718 self.opts.memory_store.as_ref(),
1719 self.opts.agent_id.as_deref(),
1720 ) {
1721 (Some(store), Some(agent_id)) => {
1722 let mut memory = memory.clone();
1723 if let Some(scope) = self.opts.memory_scope.as_ref() {
1724 memory.scope = scope.clone();
1725 }
1726 memory.provenance.session_id = Some(session_id.clone());
1727 store
1728 .put(agent_id, memory)
1729 .await
1730 .err()
1731 .map(|error| error.to_string())
1732 }
1733 _ => Some(
1734 "memory persistence is unavailable without memory_store and agent_id"
1735 .to_string(),
1736 ),
1737 };
1738 action = kernel_action(
1739 &kernel,
1740 &mut pending_observations,
1741 serde_json::json!({
1742 "kind": "memory_persist_result",
1743 "effect_id": effect_id,
1744 "error": error,
1745 }),
1746 ).await?;
1747 }
1748 HostEffect::QueryMemory { query, requested_k } => {
1749 let effect_id = action.effect_id.clone();
1750 let (hits, error) = match (
1751 self.opts.memory_store.as_ref(),
1752 self.opts.agent_id.as_deref(),
1753 ) {
1754 (Some(store), Some(agent_id)) => {
1755 let mut query = query.clone();
1756 query.top_k = *requested_k;
1757 if let Some(scope) = self.opts.memory_scope.as_ref() {
1758 query.scope = scope.clone();
1759 }
1760 match store.search(agent_id, &query).await {
1761 Ok(hits) => (hits, None),
1762 Err(error) => (Vec::new(), Some(error.to_string())),
1763 }
1764 }
1765 _ => (
1766 Vec::new(),
1767 Some(
1768 "memory query is unavailable without memory_store and agent_id"
1769 .to_string(),
1770 ),
1771 ),
1772 };
1773 if error.is_none() {
1774 self.log_memory_retrieval_result(Some(&session_id), hits.clone())
1775 .await;
1776 }
1777 action = kernel_action(
1778 &kernel,
1779 &mut pending_observations,
1780 serde_json::json!({
1781 "kind": "memory_query_result",
1782 "effect_id": effect_id,
1783 "hits": hits,
1784 "error": error,
1785 }),
1786 ).await?;
1787 }
1788 HostEffect::ArchivePageOut { archived, tier, action: pressure_action, .. } => {
1789 let effect_id = action.effect_id.clone();
1790 let archived = archived.clone();
1791 let tier = tier.clone();
1792 let action_name = action_str_of(*pressure_action);
1793 let archive_start = *active_page_out_start.get_or_insert_with(|| {
1794 pending_page_out_starts.pop_front().unwrap_or(next_archive_start)
1795 });
1796 let archive_result = if let Some(store) = &self.opts.compression_store {
1797 store.write(&session_id, archive_start, &archived)
1798 .map(|path| (!path.is_empty()).then_some(path))
1799 } else {
1800 Ok(None)
1801 };
1802 let (archive_ref, error) = match archive_result {
1803 Ok(archive_ref) => {
1804 self.local_page_out_cache.lock().unwrap().extend(archived.clone());
1805 if tier == "semantic" {
1806 self.archive_semantic_page_out(archived, Some(action_name)).await;
1807 }
1808 (archive_ref, None)
1809 }
1810 Err(error) => (None, Some(error.to_string())),
1811 };
1812 if error.is_none() {
1813 active_page_out_start = None;
1814 }
1815 action = kernel_action(
1816 &kernel,
1817 &mut pending_observations,
1818 serde_json::json!({
1819 "kind": "page_out_archive_result",
1820 "effect_id": effect_id,
1821 "archive_ref": archive_ref,
1822 "error": error,
1823 }),
1824 ).await?;
1825 }
1826 HostEffect::LoadPayload { handle_id, payload_ref } => {
1827 let effect_id = action.effect_id.clone();
1828 let content = self
1829 .opts
1830 .payload_store
1831 .as_ref()
1832 .expect("runtime constructor installs a payload store")
1833 .load(&session_id, payload_ref)?;
1834 let event = match content {
1835 Some(content) => serde_json::json!({
1836 "kind": "payload_loaded",
1837 "effect_id": effect_id,
1838 "handle_id": handle_id,
1839 "digest": deepstrike_core::runtime::kernel::wire::canonical_digest(
1840 content.as_bytes(),
1841 )
1842 .as_str(),
1843 "original_size": content.len(),
1844 "content": content,
1845 }),
1846 None => serde_json::json!({
1847 "kind": "payload_load_failed",
1848 "effect_id": effect_id,
1849 "error": format!("payload is unavailable: {payload_ref}"),
1850 }),
1851 };
1852 action = kernel_action(
1853 &kernel,
1854 &mut pending_observations,
1855 event,
1856 ).await?;
1857 }
1858 HostEffect::ExecuteTool { calls } => {
1859 let tool_effect_id = action.effect_id.clone();
1860 let tool_calls = calls.clone();
1861 self.log(
1862 &session_id,
1863 SessionEvent::ToolRequested {
1864 turn: kernel.lock().await.turn(),
1865 calls: tool_calls.clone(),
1866 },
1867 )
1868 .await;
1869
1870 if let Some(gov) = &self.opts.governance {
1871 let mut g = gov.lock().await;
1872 if let Some(aid) = &self.opts.agent_id {
1873 g.set_identity(aid, &session_id);
1874 }
1875 }
1876
1877 let run_ctx = RunContext {
1878 agent_id: self.opts.agent_id.as_deref(),
1879 memory_scope: self.opts.memory_scope.as_ref(),
1880 skill_dir: self.opts.skill_dir.as_deref(),
1881 memory_store: self.opts.memory_store.as_deref(),
1882 knowledge_source: self.opts.knowledge_source.as_deref(),
1883 governance: self.opts.governance.clone(),
1884 on_tool_suspend: self.opts.on_tool_suspend.clone(),
1885 on_permission_request: self.opts.on_permission_request.clone(),
1886 };
1887
1888 let mut tool_results = Vec::new();
1889 let mut normal_calls = Vec::new();
1890 let mut plan_calls = Vec::new();
1891
1892 for call in &tool_calls {
1893 if call.name == "update_plan" {
1894 plan_calls.push(call);
1895 } else {
1896 normal_calls.push(call.clone());
1897 }
1898 }
1899
1900 for call in plan_calls {
1901 let update = parse_update_plan_args(&call.arguments);
1902 kernel_apply(
1903 &kernel,
1904 &mut pending_observations,
1905 serde_json::json!({ "kind": "update_task", "update": update }),
1906 ).await?;
1907 tool_results.push(deepstrike_core::types::message::ToolResult {
1908 call_id: call.id.clone(),
1909 output: deepstrike_core::types::message::Content::Text("success".to_string()),
1910 durable_content: None,
1911 is_error: false,
1912 is_fatal: false,
1913 error_kind: None,
1914 });
1915 yield RunEvent::ToolResult {
1916 call_id: call.id.to_string(),
1917 content: "success".to_string(),
1918 is_error: false,
1919 is_fatal: false,
1920 error_kind: None,
1921 };
1922 }
1923
1924 if !normal_calls.is_empty() {
1925 let plane_stream = self.plane.execute_all(&normal_calls, run_ctx);
1926 let mut stream = plane_stream;
1927 while let Some(evt) = stream.next().await {
1928 match evt? {
1929 RunEvent::ToolResult {
1930 call_id,
1931 content,
1932 is_error,
1933 is_fatal,
1934 error_kind,
1935 } => {
1936 tool_results.push(deepstrike_core::types::message::ToolResult {
1937 call_id: compact_str::CompactString::new(&call_id),
1938 output: deepstrike_core::types::message::Content::Text(content),
1939 durable_content: None,
1940 is_error,
1941 is_fatal,
1942 error_kind,
1943 });
1944 }
1945 RunEvent::ToolArgumentRepaired { call_id, name, original_arguments, repaired_arguments } => {
1946 self.log(
1947 &session_id,
1948 SessionEvent::ToolArgumentRepaired {
1949 turn: kernel.lock().await.turn(),
1950 tool: name.clone(),
1951 original_arguments: original_arguments.clone(),
1952 repaired_arguments: repaired_arguments.clone(),
1953 },
1954 )
1955 .await;
1956 yield RunEvent::ToolArgumentRepaired {
1957 call_id,
1958 name,
1959 original_arguments,
1960 repaired_arguments,
1961 };
1962 }
1963 RunEvent::ToolDenied { call_id, tool_name, reason } => {
1964 self.log(
1965 &session_id,
1966 SessionEvent::ToolDenied {
1967 turn: kernel.lock().await.turn(),
1968 call_id: call_id.clone(),
1969 tool_name: tool_name.clone(),
1970 reason: reason.clone(),
1971 },
1972 )
1973 .await;
1974 yield RunEvent::ToolDenied { call_id, tool_name, reason };
1975 }
1976 RunEvent::PermissionRequest { call_id, tool_name, arguments, reason } => {
1977 let turn = kernel.lock().await.turn();
1978 self.log(
1979 &session_id,
1980 SessionEvent::PermissionRequested {
1981 turn,
1982 tool: tool_name.clone(),
1983 arguments: arguments.clone(),
1984 reason: Some(reason.clone()),
1985 },
1986 )
1987 .await;
1988 yield RunEvent::PermissionRequest { call_id, tool_name, arguments, reason };
1989 }
1990 RunEvent::PermissionResolved { call_id, tool_name, approved, responder, reason } => {
1991 let turn = kernel.lock().await.turn();
1992 self.log(
1993 &session_id,
1994 SessionEvent::PermissionResolved {
1995 turn,
1996 approved,
1997 responder: responder.clone(),
1998 },
1999 )
2000 .await;
2001 yield RunEvent::PermissionResolved { call_id, tool_name, approved, responder, reason };
2002 }
2003 other => yield other,
2004 }
2005 }
2006 let names: Vec<String> = normal_calls.iter().map(|c| c.name.to_string()).collect();
2007 kernel_apply(
2008 &kernel,
2009 &mut pending_observations,
2010 serde_json::json!({
2011 "kind": "update_task",
2012 "update": TaskUpdate {
2013 progress: Some(format!("Executed tools: {}", names.join(", "))),
2014 ..Default::default()
2015 },
2016 }),
2017 ).await?;
2018 }
2019
2020 self.log(
2021 &session_id,
2022 SessionEvent::ToolCompleted {
2023 turn: kernel.lock().await.turn(),
2024 results: tool_results.clone(),
2025 },
2026 )
2027 .await;
2028
2029 action = kernel_action(
2030 &kernel,
2031 &mut pending_observations,
2032 serde_json::json!({
2033 "kind": "tool_results",
2034 "effect_id": tool_effect_id,
2035 "results": tool_results,
2036 }),
2037 ).await?;
2038 }
2039 HostEffect::EvaluateMilestone {
2040 phase_id,
2041 criteria,
2042 required_evidence,
2043 ..
2044 } => {
2045 let milestone_effect_id = action.effect_id.clone();
2046 let policy = self.opts.milestone_policy;
2047 if policy == MilestonePolicy::AutoPass {
2048 let result = MilestoneCheckResult::pass(phase_id.clone());
2049 action = kernel_action(
2050 &kernel,
2051 &mut pending_observations,
2052 serde_json::json!({
2053 "kind": "milestone_result",
2054 "effect_id": milestone_effect_id,
2055 "result": result,
2056 }),
2057 ).await?;
2058 next_archive_start = self
2059 .append_observations(
2060 &session_id,
2061 &kernel,
2062 &mut pending_observations,
2063 &mut pending_page_out_starts,
2064 next_archive_start,
2065 )
2066 .await;
2067 } else if let Some(handler) = &self.opts.on_milestone_evaluate {
2068 let context = MilestoneEvaluationContext {
2069 phase_id: phase_id.clone(),
2070 criteria: criteria.clone(),
2071 required_evidence: required_evidence.clone(),
2072 };
2073 let check_future = handler(context);
2074 let result = check_future.await?;
2075 action = kernel_action(
2076 &kernel,
2077 &mut pending_observations,
2078 serde_json::json!({
2079 "kind": "milestone_result",
2080 "effect_id": milestone_effect_id,
2081 "result": result,
2082 }),
2083 ).await?;
2084 next_archive_start = self
2085 .append_observations(
2086 &session_id,
2087 &kernel,
2088 &mut pending_observations,
2089 &mut pending_page_out_starts,
2090 next_archive_start,
2091 )
2092 .await;
2093 } else {
2094 let result = MilestoneCheckResult::fail(
2107 phase_id.clone(),
2108 "milestone unverified: no verifier configured and no host evaluation hook (fail-closed)",
2109 );
2110 let _unverified = kernel_action(
2111 &kernel,
2112 &mut pending_observations,
2113 serde_json::json!({
2114 "kind": "milestone_result",
2115 "effect_id": milestone_effect_id,
2116 "result": result,
2117 }),
2118 ).await?;
2119 next_archive_start = self
2120 .append_observations(
2121 &session_id,
2122 &kernel,
2123 &mut pending_observations,
2124 &mut pending_page_out_starts,
2125 next_archive_start,
2126 )
2127 .await;
2128 self.log(
2129 &session_id,
2130 SessionEvent::RunTerminal {
2131 reason: "milestone_pending".to_string(),
2132 turns_used: kernel.lock().await.turn().max(1),
2133 total_tokens: 0,
2134 },
2135 )
2136 .await;
2137 yield RunEvent::Done {
2138 iterations: kernel.lock().await.turn().max(1),
2139 total_tokens: 0,
2140 status: "milestone_pending".to_string(),
2141 };
2142 return;
2143 }
2144 }
2145 HostEffect::Done { result } => {
2146 let status = format!("{:?}", result.termination).to_lowercase();
2147 let turns_used = result.turns_used.max(1);
2148 let total_tokens = result.total_tokens_used;
2149
2150 next_archive_start = self
2151 .append_observations(
2152 &session_id,
2153 &kernel,
2154 &mut pending_observations,
2155 &mut pending_page_out_starts,
2156 next_archive_start,
2157 )
2158 .await;
2159
2160 self.log(
2161 &session_id,
2162 SessionEvent::RunTerminal {
2163 reason: status.clone(),
2164 turns_used,
2165 total_tokens,
2166 },
2167 )
2168 .await;
2169
2170 if let (Some(store), Some(agent_id)) =
2171 (&self.opts.memory_store, &self.opts.agent_id)
2172 {
2173 let new_msgs = kernel.lock().await.drain_new_messages();
2174 if !new_msgs.is_empty() {
2175 let now_ms = std::time::SystemTime::now()
2176 .duration_since(std::time::UNIX_EPOCH)
2177 .unwrap_or_default()
2178 .as_millis() as u64;
2179 let session = deepstrike_core::memory::durable::SessionData {
2180 session_id: session_id.clone(),
2181 agent_id: agent_id.clone(),
2182 messages: new_msgs,
2183 metadata: serde_json::Value::Null,
2184 created_at_ms: session_start_ms,
2185 updated_at_ms: now_ms,
2186 };
2187 let _ = store.save_session(session.clone()).await;
2188 if let Some(scope) = self.opts.memory_scope.as_ref() {
2189 if let Ok(memories) = self.extract_session_memories(&session, scope).await {
2190 for memory in memories {
2191 let _ = self.write_memory(memory, Some(&session_id), Some(agent_id)).await;
2192 }
2193 }
2194 }
2195 }
2196 }
2197
2198 yield RunEvent::Done {
2199 iterations: turns_used,
2200 total_tokens,
2201 status,
2202 };
2203 return;
2204 }
2205 }
2206 }
2207
2208 next_archive_start = self
2209 .append_observations(
2210 &session_id,
2211 &kernel,
2212 &mut pending_observations,
2213 &mut pending_page_out_starts,
2214 next_archive_start,
2215 )
2216 .await;
2217
2218 let (status, turns_used, total_tokens) = match &action.effect {
2222 HostEffect::Done { result } => (
2223 format!("{:?}", result.termination).to_lowercase(),
2224 result.turns_used.max(1),
2225 result.total_tokens_used,
2226 ),
2227 _ => ("error".to_string(), kernel.lock().await.turn().max(1), 0),
2228 };
2229
2230 self.log(
2231 &session_id,
2232 SessionEvent::RunTerminal {
2233 reason: status.clone(),
2234 turns_used,
2235 total_tokens,
2236 },
2237 )
2238 .await;
2239
2240 if let HostEffect::Done { .. } = &action.effect {
2241 if let (Some(store), Some(agent_id)) =
2242 (&self.opts.memory_store, &self.opts.agent_id)
2243 {
2244 let new_msgs = kernel.lock().await.drain_new_messages();
2245 if !new_msgs.is_empty() {
2246 let now_ms = std::time::SystemTime::now()
2247 .duration_since(std::time::UNIX_EPOCH)
2248 .unwrap_or_default()
2249 .as_millis() as u64;
2250 let session = deepstrike_core::memory::durable::SessionData {
2251 session_id: session_id.clone(),
2252 agent_id: agent_id.clone(),
2253 messages: new_msgs,
2254 metadata: serde_json::Value::Null,
2255 created_at_ms: session_start_ms,
2256 updated_at_ms: now_ms,
2257 };
2258 let _ = store.save_session(session.clone()).await;
2259 if let Some(scope) = self.opts.memory_scope.as_ref() {
2260 if let Ok(memories) = self.extract_session_memories(&session, scope).await {
2261 for memory in memories {
2262 let _ = self.write_memory(memory, Some(&session_id), Some(agent_id)).await;
2263 }
2264 }
2265 }
2266 }
2267 }
2268 }
2269
2270 yield RunEvent::Done {
2271 iterations: turns_used,
2272 total_tokens,
2273 status,
2274 };
2275 }
2276 }
2277
2278 pub(crate) async fn append_observations(
2279 &self,
2280 session_id: &str,
2281 kernel_mutex: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2282 observations: &mut Vec<KernelObservation>,
2283 pending_page_out_starts: &mut std::collections::VecDeque<u64>,
2284 mut next_archive_start: u64,
2285 ) -> u64 {
2286 let drained = std::mem::take(observations);
2287 let (turn, preserved_refs, summary_tokens_by_index) = {
2288 let kernel = kernel_mutex.lock().await;
2289 let summary_tokens_by_index = drained
2290 .iter()
2291 .map(|obs| match obs {
2292 KernelObservation::Compressed { summary, .. } => {
2293 summary.as_ref().map(|s| kernel.count_tokens(s))
2294 }
2295 _ => None,
2296 })
2297 .collect::<Vec<_>>();
2298 (
2299 kernel.turn(),
2300 kernel.preserved_refs(),
2301 summary_tokens_by_index,
2302 )
2303 };
2304
2305 for (index, obs) in drained.into_iter().enumerate() {
2306 match obs {
2307 KernelObservation::Compressed {
2308 turn: _,
2309 action,
2310 rho_after: _,
2311 summary,
2312 archived_count,
2313 invalidates_prefix_at: _,
2314 } => {
2315 let Some(log) = &self.opts.session_log else {
2316 continue;
2317 };
2318 let latest = log.latest_seq(session_id).await.unwrap_or(-1) as u64;
2319 if latest < next_archive_start {
2320 continue;
2321 }
2322 let end = latest;
2323 if archived_count > 0 {
2324 pending_page_out_starts.push_back(next_archive_start);
2325 }
2326
2327 let summary_tokens = summary_tokens_by_index.get(index).copied().flatten();
2328 let action_str = action_str_of(action);
2329
2330 if let Ok(compressed_seq) = log
2331 .append(
2332 session_id,
2333 SessionEvent::Compressed {
2334 turn,
2335 archived_seq_range: (next_archive_start, end),
2336 action: Some(action_str),
2337 summary: summary.clone(),
2338 summary_tokens,
2339 preserved_refs: preserved_refs.clone(),
2340 },
2341 )
2342 .await
2343 {
2344 next_archive_start = compressed_seq + 1;
2345 }
2346 }
2347 KernelObservation::PageOutArchived {
2348 turn,
2349 action,
2350 summary,
2351 tier,
2352 message_count,
2353 archive_ref,
2354 } => {
2355 self.log(
2356 session_id,
2357 SessionEvent::PageOut {
2358 turn,
2359 action: Some(action_str_of(action)),
2360 summary,
2361 tier_hint: Some(tier),
2362 message_count,
2363 archive_ref,
2364 },
2365 )
2366 .await;
2367 }
2368 KernelObservation::PageOutArchiveFailed { .. } => {}
2369 KernelObservation::PayloadResidencyChanged { .. }
2371 | KernelObservation::PayloadLoadFailed { .. } => {}
2372 KernelObservation::Rollbacked {
2373 turn,
2374 checkpoint_history_len,
2375 reason,
2376 } => {
2377 self.log(
2378 session_id,
2379 SessionEvent::Rollbacked {
2380 turn,
2381 checkpoint_history_len,
2382 reason,
2383 },
2384 )
2385 .await;
2386 }
2387 KernelObservation::CapabilityChanged {
2388 turn,
2389 added,
2390 removed,
2391 change_kind,
2392 capability_id,
2393 version,
2394 mounted_by,
2395 mount_reason,
2396 } => {
2397 self.log(
2398 session_id,
2399 SessionEvent::CapabilityChanged {
2400 turn,
2401 added,
2402 removed,
2403 change_kind,
2404 capability_id,
2405 version,
2406 mounted_by,
2407 mount_reason,
2408 },
2409 )
2410 .await;
2411 }
2412 KernelObservation::MilestoneAdvanced {
2413 turn,
2414 phase_id,
2415 capabilities_unlocked,
2416 } => {
2417 self.log(
2418 session_id,
2419 SessionEvent::MilestoneAdvanced {
2420 turn,
2421 phase_id,
2422 capabilities_unlocked,
2423 },
2424 )
2425 .await;
2426 }
2427 KernelObservation::MilestoneBlocked {
2428 turn,
2429 phase_id,
2430 reason,
2431 } => {
2432 self.log(
2433 session_id,
2434 SessionEvent::MilestoneBlocked {
2435 turn,
2436 phase_id,
2437 reason,
2438 },
2439 )
2440 .await;
2441 }
2442 KernelObservation::Renewed { .. } => {}
2443 KernelObservation::ContextBudgetExceeded { .. } => {}
2444 KernelObservation::KnowledgeSwept { .. } => {}
2445 KernelObservation::KnowledgeBudgetExceeded { .. } => {}
2446 KernelObservation::RepeatFuseTripped { .. } => {}
2447 KernelObservation::CriteriaGateFired { .. } => {}
2448 KernelObservation::CheckpointTaken { turn, history_len } => {
2449 self.log(
2450 session_id,
2451 SessionEvent::CheckpointTaken { turn, history_len },
2452 )
2453 .await;
2454 }
2455 KernelObservation::EntropySample {
2456 turn,
2457 score,
2458 rho,
2459 repeat_pressure,
2460 failure_rate,
2461 rollbacks_in_window,
2462 window_turns,
2463 } => {
2464 self.log(
2465 session_id,
2466 SessionEvent::EntropySample {
2467 turn,
2468 score,
2469 rho,
2470 repeat_pressure,
2471 failure_rate,
2472 rollbacks_in_window,
2473 window_turns,
2474 },
2475 )
2476 .await;
2477 }
2478 KernelObservation::EntropyAlert {
2479 turn,
2480 score,
2481 threshold,
2482 } => {
2483 self.log(
2484 session_id,
2485 SessionEvent::EntropyAlert {
2486 turn,
2487 score,
2488 threshold,
2489 },
2490 )
2491 .await;
2492 }
2493 KernelObservation::AgentProcessChanged { .. } => {}
2494 KernelObservation::ChildSupervised { .. }
2497 | KernelObservation::LocalRunnableTrace { .. } => {}
2498 KernelObservation::WorkflowBatchSpawned { .. } => {}
2501 KernelObservation::WorkflowSpawnFailed { .. } => {}
2502 KernelObservation::WorkflowCompleted { .. } => {}
2503 KernelObservation::NodesRejected { .. } => {}
2504 KernelObservation::AgentPreempted { .. } => {}
2505 KernelObservation::AgentPreemptFailed { .. } => {}
2506 KernelObservation::MemoryWriteFailed { .. } => {}
2507 KernelObservation::MemoryQueryFailed { .. } => {}
2508 KernelObservation::MemoryRecalled { .. }
2511 | KernelObservation::PromotionSuggested { .. } => {}
2512 KernelObservation::ToolGated { .. } => {}
2515 KernelObservation::SignalDeliveryDisposed { .. } => {}
2519 KernelObservation::SignalDisplaced { .. }
2520 | KernelObservation::SignalExpired { .. }
2521 | KernelObservation::SignalsPending { .. } => {}
2522 KernelObservation::BudgetExceeded {
2523 turn,
2524 operation_id,
2525 reservation_id,
2526 budget,
2527 } => {
2528 self.log(
2529 session_id,
2530 SessionEvent::BudgetExceeded {
2531 turn,
2532 operation_id,
2533 reservation_id,
2534 budget,
2535 },
2536 )
2537 .await;
2538 }
2539 KernelObservation::BudgetUsageReported {
2540 operation_id,
2541 reservation_id,
2542 tokens,
2543 subagents,
2544 rounds,
2545 } => {
2546 self.log(
2547 session_id,
2548 SessionEvent::BudgetUsageReported {
2549 turn,
2550 operation_id,
2551 reservation_id,
2552 tokens,
2553 subagents,
2554 rounds,
2555 },
2556 )
2557 .await;
2558 }
2559 KernelObservation::OperationCancelled {
2560 turn,
2561 operation_id,
2562 reason,
2563 pending_call_ids,
2564 } => {
2565 self.log(
2566 session_id,
2567 SessionEvent::OperationCancelled {
2568 turn,
2569 operation_id,
2570 reason,
2571 pending_call_ids,
2572 },
2573 )
2574 .await;
2575 }
2576 KernelObservation::LivePolicyChanged { .. } => {}
2579 KernelObservation::Suspended { .. }
2580 | KernelObservation::ApprovalResolutionFailed { .. } => {}
2581 KernelObservation::Resumed { .. } => {}
2582 KernelObservation::WorkflowNodesSubmitted { .. } => {}
2585 KernelObservation::RoundPaced { .. } => {}
2588 KernelObservation::MemoryWritten {
2589 turn,
2590 record_id,
2591 scope,
2592 memory_kind,
2593 name,
2594 size_bytes,
2595 } => {
2596 self.log(
2597 session_id,
2598 SessionEvent::MemoryWritten {
2599 turn,
2600 record_id,
2601 scope,
2602 memory_kind,
2603 name,
2604 size_bytes,
2605 },
2606 )
2607 .await;
2608 }
2609 KernelObservation::MemoryQueried {
2610 turn,
2611 scope,
2612 query,
2613 requested_k,
2614 requires_async_response,
2615 } => {
2616 self.log(
2617 session_id,
2618 SessionEvent::MemoryQueried {
2619 turn,
2620 scope,
2621 query,
2622 requested_k,
2623 requires_async_response,
2624 },
2625 )
2626 .await;
2627 }
2628 KernelObservation::MemoryValidationFailed {
2630 turn,
2631 record_id,
2632 error,
2633 } => {
2634 self.log(
2635 session_id,
2636 SessionEvent::MemoryValidationFailed {
2637 turn,
2638 record_id,
2639 error,
2640 },
2641 )
2642 .await;
2643 }
2644 KernelObservation::ControlRequestRejected { .. } => {}
2647 KernelObservation::StepPublishedEffects { .. } => {}
2648 }
2649 }
2650 next_archive_start
2651 }
2652
2653 async fn read_entries(&self, session_id: &str) -> Result<Vec<SessionEntry>> {
2654 if let Some(log) = &self.opts.session_log {
2655 log.read(session_id, 0, None).await.map_err(Error::Io)
2656 } else {
2657 Ok(Vec::new())
2658 }
2659 }
2660
2661 async fn log(&self, session_id: &str, event: SessionEvent) {
2662 if let Some(log) = &self.opts.session_log {
2663 let _ = log.append(session_id, event).await;
2664 }
2665 }
2666
2667 async fn archive_semantic_page_out(&self, archived: Vec<CoreMessage>, action: Option<String>) {
2668 let (Some(_store), Some(agent_id), Some(scope)) = (
2669 &self.opts.memory_store,
2670 &self.opts.agent_id,
2671 &self.opts.memory_scope,
2672 ) else {
2673 return;
2674 };
2675
2676 let summary = match self.summarize_for_long_term_memory(&archived).await {
2677 Ok(s) => s,
2678 Err(_) => return, };
2680
2681 let now = std::time::SystemTime::now()
2685 .duration_since(std::time::UNIX_EPOCH)
2686 .unwrap_or_default()
2687 .as_millis() as u64;
2688 let name = format!("page-out-{now}");
2689 let request = MemoryRecord {
2690 record_id: format!("{}:{}:project:{}", scope.tenant_id, scope.namespace, name),
2691 scope: scope.clone(),
2692 name,
2693 kind: MemoryKind::Project,
2694 content: summary,
2695 description: format!(
2696 "auto summary of {} archive",
2697 action.as_deref().unwrap_or("compaction")
2698 ),
2699 provenance: MemoryProvenance {
2700 session_id: self.opts.session_id.clone(),
2701 author: MemoryAuthor::Extraction,
2702 trust: MemoryTrustLevel::Untrusted,
2703 evidence_refs: Vec::new(),
2704 },
2705 created_at: now,
2706 updated_at: now,
2707 last_recalled_at: None,
2708 recall_count: 0,
2709 confidence: 0.6,
2710 links: Vec::new(),
2711 pinned: false,
2712 ttl_days: None,
2713 };
2714 let _ = self.write_memory(request, None, Some(agent_id)).await;
2715 }
2716
2717 async fn summarize_for_long_term_memory(
2718 &self,
2719 archived: &[CoreMessage],
2720 ) -> crate::Result<String> {
2721 let transcript = archived
2722 .iter()
2723 .map(|m| {
2724 let role_str = match m.role {
2725 deepstrike_core::types::message::Role::System => "system",
2726 deepstrike_core::types::message::Role::User => "user",
2727 deepstrike_core::types::message::Role::Assistant => "assistant",
2728 deepstrike_core::types::message::Role::Tool => "tool",
2729 };
2730 let content_str = message_content_as_text(&m.content);
2731 format!("{}: {}", role_str, content_str)
2732 })
2733 .collect::<Vec<_>>()
2734 .join("\n");
2735
2736 let system_prompt_opt = self.opts.system_prompt.as_deref();
2737 let system_text = match system_prompt_opt {
2738 Some(sp) => format!(
2739 "{}\n\nSummarize the following conversation for long-term memory. Preserve key facts, decisions, and open questions.",
2740 sp
2741 ),
2742 None => "Summarize the following conversation for long-term memory. Preserve key facts, decisions, and open questions.".to_string(),
2743 };
2744
2745 let context = deepstrike_core::context::renderer::InternalRenderedContext {
2746 system_text,
2747 system_stable: String::new(),
2748 system_knowledge: String::new(),
2749 turns: vec![deepstrike_core::types::message::CoreMessage {
2750 role: deepstrike_core::types::message::Role::User,
2751 content: deepstrike_core::types::message::Content::Text(transcript.clone()),
2752 tool_calls: vec![],
2753 }],
2754 state_turn: None,
2755 frozen_prefix_len: None,
2756 budget_overflow: None,
2757 };
2758
2759 let synth_state = self.opts.provider.create_run_state();
2760 let mut stream = self
2761 .opts
2762 .provider
2763 .stream(&context, &[], None, synth_state.as_ref())
2764 .await?;
2765
2766 let mut synthesis_text = String::new();
2767 while let Some(evt) = stream.next().await {
2768 if let Ok(StreamEvent::TextDelta { delta }) = evt {
2769 synthesis_text.push_str(&delta);
2770 }
2771 }
2772
2773 let text = synthesis_text.trim();
2774 if text.is_empty() {
2775 Ok(transcript.chars().take(2000).collect())
2776 } else {
2777 Ok(text.to_string())
2778 }
2779 }
2780}
2781
2782fn message_content_as_text(content: &deepstrike_core::types::message::Content) -> String {
2783 match content {
2784 deepstrike_core::types::message::Content::Text(s) => s.clone(),
2785 deepstrike_core::types::message::Content::Parts(parts) => parts
2786 .iter()
2787 .filter_map(|p| match p {
2788 deepstrike_core::types::message::ContentPart::Text { text } => Some(text.as_str()),
2789 deepstrike_core::types::message::ContentPart::ToolResult { output, .. } => {
2790 Some(output.as_str())
2791 }
2792 _ => None,
2793 })
2794 .collect::<Vec<_>>()
2795 .join("\n"),
2796 }
2797}
2798
2799fn action_str_of(action: KernelPressureAction) -> String {
2800 match action {
2801 KernelPressureAction::None => "none".to_string(),
2802 KernelPressureAction::SnipCompact => "snip_compact".to_string(),
2803 KernelPressureAction::MicroCompact => "micro_compact".to_string(),
2804 KernelPressureAction::ContextCollapse => "context_collapse".to_string(),
2805 KernelPressureAction::AutoCompact => "auto_compact".to_string(),
2806 }
2807}
2808
2809pub(crate) async fn kernel_apply(
2810 kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2811 pending_observations: &mut Vec<KernelObservation>,
2812 event: serde_json::Value,
2813) -> Result<()> {
2814 let mut runtime = kernel.lock().await;
2815 canonical_kernel_apply(&mut runtime, pending_observations, event).await
2816}
2817
2818async fn kernel_transition(
2819 kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2820 pending_observations: &mut Vec<KernelObservation>,
2821 event: serde_json::Value,
2822) -> Result<Option<HostAction>> {
2823 let mut runtime = kernel.lock().await;
2824 let action = runtime.apply_host_event(event).await?;
2825 pending_observations.extend(runtime.drain_host_observations());
2826 Ok(action)
2827}
2828
2829async fn kernel_action(
2830 kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2831 pending_observations: &mut Vec<KernelObservation>,
2832 event: serde_json::Value,
2833) -> Result<HostAction> {
2834 let mut runtime = kernel.lock().await;
2835 canonical_kernel_action(&mut runtime, pending_observations, event).await
2836}
2837
2838async fn kernel_start_agent(
2839 kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2840 pending_observations: &mut Vec<KernelObservation>,
2841 task: RuntimeTask,
2842 run_spec: Option<deepstrike_core::types::agent::AgentRunSpec>,
2843) -> Result<HostAction> {
2844 let task = serde_json::to_value(task)
2845 .map_err(|error| Error::Other(format!("canonical task is not serializable: {error}")))?;
2846 let run_spec = run_spec
2847 .map(serde_json::to_value)
2848 .transpose()
2849 .map_err(|error| {
2850 Error::Other(format!("canonical run spec is not serializable: {error}"))
2851 })?;
2852 let mut runtime = kernel.lock().await;
2853 let action = runtime.start_agent_value(task, run_spec).await?;
2854 pending_observations.extend(runtime.drain_host_observations());
2855 action.ok_or_else(|| Error::Other("canonical agent root must return one host action".into()))
2856}
2857
2858pub async fn collect_text(
2859 mut stream: std::pin::Pin<Box<dyn futures::Stream<Item = Result<RunEvent>> + '_>>,
2860) -> Result<String> {
2861 let mut text = String::new();
2862 while let Some(evt) = stream.next().await {
2863 if let RunEvent::TextDelta(d) = evt? {
2864 text.push_str(&d);
2865 }
2866 }
2867 Ok(text)
2868}
2869
2870fn merge_extensions(
2871 base: Option<&serde_json::Value>,
2872 over: Option<&serde_json::Value>,
2873) -> Option<serde_json::Value> {
2874 match (base, over) {
2875 (Some(b), Some(o)) => {
2876 let mut merged = b.clone();
2877 if let (Some(m), Some(obj)) = (merged.as_object_mut(), o.as_object()) {
2878 for (k, v) in obj {
2879 m.insert(k.clone(), v.clone());
2880 }
2881 }
2882 Some(merged)
2883 }
2884 (Some(b), None) => Some(b.clone()),
2885 (None, Some(o)) => Some(o.clone()),
2886 (None, None) => None,
2887 }
2888}
2889
2890fn cancellation_reason_code(reason: CancellationReason) -> u8 {
2891 match reason {
2892 CancellationReason::User => 0,
2893 CancellationReason::Deadline => 1,
2894 CancellationReason::LeaseLost => 2,
2895 CancellationReason::HostShutdown => 3,
2896 }
2897}
2898
2899fn provider_error_message(error: &crate::Error) -> String {
2900 match error {
2901 crate::Error::ProviderFailure(error) => error.message.clone(),
2902 _ => error.to_string(),
2903 }
2904}
2905
2906fn provider_error_event(effect_id: &str, error: &crate::Error) -> serde_json::Value {
2907 let mut event = serde_json::json!({
2908 "kind": "provider_error",
2909 "effect_id": effect_id,
2910 "message": provider_error_message(error),
2911 });
2912 if let crate::Error::ProviderFailure(error) = error {
2913 let object = event
2914 .as_object_mut()
2915 .expect("provider error event is an object");
2916 object.insert("error_kind".into(), error.kind.as_str().into());
2917 object.insert("retryable".into(), error.retryable.into());
2918 if let Some(status) = error.http_status {
2919 object.insert("http_status".into(), status.into());
2920 }
2921 if let Some(code) = error.provider_code.as_ref() {
2922 object.insert("provider_code".into(), code.clone().into());
2923 }
2924 }
2925 event
2926}
2927
2928fn cancellation_reason_from_code(code: u8) -> CancellationReason {
2929 match code {
2930 1 => CancellationReason::Deadline,
2931 2 => CancellationReason::LeaseLost,
2932 3 => CancellationReason::HostShutdown,
2933 _ => CancellationReason::User,
2934 }
2935}
2936
2937fn pending_call_ids(action: &HostAction) -> Vec<String> {
2938 match &action.effect {
2939 HostEffect::CallProvider { .. } => vec![action.effect_id.clone()],
2940 HostEffect::ExecuteTool { calls } => calls.iter().map(|call| call.id.to_string()).collect(),
2941 HostEffect::RequestApproval { requests } => requests
2942 .iter()
2943 .map(|request| request.call_id.clone())
2944 .collect(),
2945 HostEffect::SpawnWorkflow { nodes, .. } => {
2946 nodes.iter().map(|node| node.agent_id.clone()).collect()
2947 }
2948 HostEffect::PreemptSubAgents { agent_ids, .. } => agent_ids.clone(),
2949 HostEffect::Done { .. } => Vec::new(),
2950 _ => vec![action.effect_id.clone()],
2951 }
2952}
2953
2954fn memory_policy_host_fact(policy: MemoryPolicy) -> serde_json::Value {
2956 let mut value = serde_json::to_value(policy).expect("canonical memory policy serializes");
2957 value
2958 .as_object_mut()
2959 .expect("canonical memory policy is a JSON object")
2960 .insert("kind".into(), serde_json::json!("set_memory_policy"));
2961 value
2962}
2963
2964fn next_archived_seq_start(events: Option<&[SessionEntry]>) -> u64 {
2965 let mut next = 0u64;
2966 for entry in events.unwrap_or_default() {
2967 if let SessionEvent::Compressed {
2968 archived_seq_range, ..
2969 } = &entry.event
2970 {
2971 next = next.max(archived_seq_range.1 + 1);
2972 }
2973 }
2974 next
2975}
2976
2977fn rendered_context_from_messages(
2978 messages: Vec<CoreMessage>,
2979) -> deepstrike_core::context::renderer::InternalRenderedContext {
2980 let mut system_parts = Vec::new();
2981 let mut turns = Vec::new();
2982 for message in messages {
2983 if message.role == deepstrike_core::types::message::Role::System {
2984 if let Some(text) = message.content.as_text() {
2985 system_parts.push(text.to_owned());
2986 }
2987 } else {
2988 turns.push(message);
2989 }
2990 }
2991 let system_text = system_parts.join("\n\n");
2992 deepstrike_core::context::renderer::InternalRenderedContext {
2993 system_text: system_text.clone(),
2994 system_stable: system_text,
2995 system_knowledge: String::new(),
2996 turns,
2997 state_turn: None,
2998 frozen_prefix_len: None,
2999 budget_overflow: None,
3000 }
3001}
3002
3003fn parse_update_plan_args(val: &serde_json::Value) -> TaskUpdate {
3004 let plan = val.get("plan").and_then(|v| {
3005 v.as_array().map(|arr| {
3006 arr.iter()
3007 .filter_map(|x| x.as_str().map(|s| s.to_string()))
3008 .collect()
3009 })
3010 });
3011 let current_step = val
3012 .get("current_step")
3013 .or_else(|| val.get("currentStep"))
3014 .and_then(|v| v.as_u64().map(|x| x as usize));
3015 let progress = val
3016 .get("progress")
3017 .and_then(|v| v.as_str().map(|s| s.to_string()));
3018 let scratchpad = val
3019 .get("scratchpad")
3020 .and_then(|v| v.as_str().map(|s| s.to_string()));
3021 let blocked_on = val
3022 .get("blocked_on")
3023 .or_else(|| val.get("blockedOn"))
3024 .and_then(|v| {
3025 v.as_array().map(|arr| {
3026 arr.iter()
3027 .filter_map(|x| x.as_str().map(|s| s.to_string()))
3028 .collect()
3029 })
3030 });
3031 let preserved_refs = val
3032 .get("preserved_refs")
3033 .or_else(|| val.get("preservedRefs"))
3034 .and_then(|v| {
3035 v.as_array().map(|arr| {
3036 arr.iter()
3037 .filter_map(|x| x.as_str().map(|s| s.to_string()))
3038 .collect()
3039 })
3040 });
3041 TaskUpdate {
3042 plan,
3043 current_step,
3044 progress,
3045 scratchpad,
3046 blocked_on,
3047 preserved_refs,
3048 directives: None,
3051 }
3052}