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