Skip to main content

deepstrike_sdk/runtime/
runner.rs

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/// Controls what the runner does when the state machine returns
57/// `EvaluateMilestone` — i.e., the LLM finished a turn but a milestone phase
58/// has not yet been evaluated.
59#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
60pub enum MilestonePolicy {
61    /// Wait for a verifier callback or suspend if none is configured (default).
62    #[default]
63    RequireVerifier,
64    /// Terminate the run immediately with `status = "milestone_pending"`.
65    Terminate,
66    /// Unconditionally pass every milestone phase.  Useful in unit tests and
67    /// capability-unlock–only scenarios where the criteria check is a no-op.
68    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/// P0-C tool-gating telemetry: per-LLM-turn metrics, delivered to [`RuntimeOptions::on_turn_metrics`].
87/// Pure observation — no behavior change. `tools_exposed` vs `tools_called` quantifies over-exposure;
88/// consecutive equal `active_skill` values measure skill dwell `D`; the cache split gives the
89/// prompt-cache hit baseline. Mirrors the node SDK `TurnMetrics`.
90#[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    /// I1: pro-rata per-slot attribution of `cache_read_tokens` (Anthropic only). Mirrors Node.
100    pub cache_read_tokens_by_slot: Option<crate::providers::CacheReadBySlot>,
101}
102
103/// Sink for per-turn [`TurnMetrics`]. Synchronous, infallible — it must never affect the run.
104pub type OnTurnMetricsHandler = std::sync::Arc<dyn Fn(TurnMetrics) + Send + Sync>;
105
106/// Canonical recovery and input-bound overrides. Omitted fields retain core defaults.
107#[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
114/// Configuration for a `RuntimeRunner` (aligned with Node/Python `RuntimeOptions`).
115pub 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    /// Storage for canonical opaque external payload locators.
121    pub payload_store: Option<Arc<dyn PayloadStore>>,
122    /// Bounded recovery and replay policy. Omitted fields retain kernel defaults.
123    pub kernel_reliability: Option<KernelReliability>,
124    /// When set, `execute` reuses this session id.
125    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    /// Required by host-generated memory queries and semantic page-out writes.
132    pub memory_scope: Option<MemoryScope>,
133    /// I4: optional run-start memory pre-fetch hook. The runner calls this once per run, before
134    /// the first LLM turn, with the goal string; each returned query becomes a `memory_store.search`
135    /// and the resulting hits page into the knowledge partition before turn 1. Mirrors the Node
136    /// SDK `preQueryMemory`. Sync-only in Rust today — async hosts can pre-compute. Errs-open
137    /// when `memory_store` or `agent_id` is missing.
138    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    /// Opt-in long-term memory policy (`set_memory_policy`), enforced at the kernel memory traps.
152    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    /// How to handle `EvaluateMilestone` actions. Default: `RequireVerifier`.
158    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    /// The run's **exposure ceiling** — the outer bound on what this run may EVER advertise. Not a
162    /// static profile: an INTERSECTION applied every turn (`exposed ⊆ ceiling`), so `baseline_tool_ids`,
163    /// `stable_core_tool_ids`, and skill `allowed_tools` all narrow *within* it and none can widen
164    /// past it. The kernel meta-tools (skill/memory/knowledge/update_plan/read_result) are exempt on
165    /// the id axis; the KIND axis still applies. Lowers to the same `capability_filter` sub-agents
166    /// use; byte-stable across the run, so it never busts the prompt-cache prefix. Augments
167    /// `run_spec`'s filter when both are set; synthesizes a minimal top-level spec otherwise.
168    /// `None`/empty means no ceiling. Exposure still starts from the minimal baseline.
169    pub allowed_tool_ids: Option<Vec<String>>,
170    /// The PRE-ACTIVATION exposure surface, selected from under the `allowed_tool_ids` ceiling
171    /// (`AgentRunSpec::exposure_baseline`). Makes narrow→wide progressive disclosure expressible:
172    /// `exposed = meta ∪ ((baseline ∪ stable_core ∪ ⋃ active skills' allowed_tools) ∩ ceiling)`.
173    /// `None` and `Some(vec![])` both select the minimal surface (meta-tools + stable-core only).
174    /// Entries outside the ceiling silently intersect away.
175    pub baseline_tool_ids: Option<Vec<String>>,
176    /// P0-C: optional per-turn metrics sink for tool-gating telemetry (see [`TurnMetrics`]). Pure
177    /// observation; invoked once per LLM turn. Panics are not caught — keep the sink trivial.
178    pub on_turn_metrics: Option<OnTurnMetricsHandler>,
179    /// P1-B/D stable-core: tool ids always exposed under skill gating. Empty ⇒ skills narrow to
180    /// exactly their declared tools + meta-tools. Opt-in: no skill declaring tools ⇒ never engages.
181    pub stable_core_tool_ids: Vec<String>,
182    pub on_milestone_evaluate: Option<MilestoneEvaluationHandler>,
183}
184
185/// P0-A: compute the effective top-level run spec from an optional explicit `run_spec`, an optional
186/// `allowed_tool_ids` ceiling, and an optional `baseline_tool_ids` pre-activation surface. Each
187/// augments an explicit spec, or synthesizes a minimal `custom`-role spec when none is given. Every
188/// run carries the canonical exposure baseline; omitted and empty both mean the minimal surface.
189fn 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
244/// Orchestrates the agentic turn loop via the runtime kernel + session event log.
245pub 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    /// Construct a runner with an explicit durable canonical journal implementation.
263    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        // The session-log record is the durable audit artifact; the kernel needs no
493        // acknowledgment (the former kernel event was a no-op and was removed).
494        self.log(session_id, SessionEvent::MemoryRetrievalResult { hits })
495            .await;
496    }
497
498    /// Test-only probe: the live kernel's `pending_effects` size. Valid while the run stream is
499    /// suspended at a `yield` (the active-kernel guard is still in scope); `None` once the run
500    /// generator has finished. Used by the effect-leak regressions (R-B27).
501    #[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    /// Like [`Self::run_streaming`], but seeds multimodal `attachments` into kernel history
681    /// before the first render (parity with Node/Python `run({ attachments })`).
682    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                // P1-B/D: configure stable-core tool ids (always exposed under skill gating).
972                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            // P0-C: the skill loaded and in effect going into the current turn → per-turn metric.
1030            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                // Multimodal upload: seed attachments before the canonical root start (Node/Python parity).
1079                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                // I4: pre-fetch memory into the knowledge partition before the first LLM turn.
1091                // Mirrors Node/WASM/Python preQueryMemory. Errs-open: missing memory_store/agent_id
1092                // or a faulty closure silently skip the pre-fetch.
1093                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                // P0-A: fold an explicit `run_spec`, the `allowed_tool_ids` ceiling, and/or the
1146                // `baseline_tool_ids` pre-activation surface into the kernel run spec (reuses the
1147                // existing run_spec wire — no new ABI).
1148                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                // Hot-reload: refresh skill catalog if the watcher detected changes.
1172                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                        // §7.7 · the claim's signal id is the business identity, kept verbatim; it
1253                        // no longer has to parse as a UUID.
1254                        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                        // Kernel-routed (parity with node/py): the kernel's attention policy decides
1269                        // the disposition (dedup / queue / interrupt / preempt) and emits
1270                        // `signal_delivery_disposed`; an actionable disposition yields the next action to
1271                        // adopt (e.g. a forced Reason turn on Critical), queued/observed yields none.
1272                        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                        // Critical attention/preemption is distinct from operation cancellation.
1326                    }
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                        // I5: governance schema-level pre-filter. When a GovernancePolicy is loaded
1344                        // and `surface_denied_in_system` is true (default), drop denied tools from
1345                        // the schema before the provider sees them.
1346                        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                        // P0-C: snapshot the exposed-tool count now — `tools` borrows `action`, which is
1367                        // reassigned before the metrics emit below.
1368                        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                                // Reactive recovery is now a kernel decision. Forward the raw
1379                                // provider error and dispatch whatever the kernel returns:
1380                                // CallProvider to retry with a freshly compacted context, or Done to
1381                                // terminate with an honest ContextOverflow. The classify + compact +
1382                                // retry + give-up policy lives in the kernel (one place), not
1383                                // duplicated across the four SDK runners.
1384                                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                                // Withholding (query.ts parity): surface the raw provider error only
1391                                // when the kernel could NOT recover (it returned a terminal). On a
1392                                // recovered retry (CallProvider) the error stays hidden. `continue`
1393                                // re-enters the loop: a recovered turn persists its compaction
1394                                // archive at the loop's normal append point, and a terminal Done
1395                                // exits through `is_terminal()` into the run_terminal emit.
1396                                if matches!(&action.effect, HostEffect::Done { .. }) {
1397                                    yield RunEvent::Error(msg);
1398                                }
1399                                continue;
1400                            }
1401                        };
1402
1403                        // R-B27/R-B29 sibling: an exception raised AFTER the first chunk must not
1404                        // escape through `?`. Doing so leaves the `call_provider` effect pending
1405                        // forever (nothing ever resolves it) and skips the kernel's reactive
1406                        // recovery ladder. Capture it here and feed `provider_error` below, exactly
1407                        // like the stream-open error path and like node/python.
1408                        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                                    // P0-C: capture input + prompt-cache split for the hit-rate baseline.
1447                                    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                                    // Phase 4: keep the last non-empty stop_reason for output-cap recovery.
1452                                    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                            // Same contract as the stream-open failure above: hand the raw provider
1476                            // error to the kernel, which resolves the pending provider effect and
1477                            // decides recover-and-retry (`CallProvider`) vs honest terminal (`Done`).
1478                            // Surface the error to the caller only when the kernel gave up, so a
1479                            // recovered turn does not emit a phantom failure.
1480                            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                                // Phase 4: stop_reason drives the kernel's max-output-tokens recovery.
1514                                "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                        // P0-C: per-turn tool-gating telemetry. `active_skill` reflects the skill in
1528                        // effect GOING INTO this turn; a `skill` call here only takes effect next turn
1529                        // — emit first, then advance.
1530                        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                        // This runner has no workflow child orchestrator. Report each
1635                        // requested spawn as a completed failure instead of treating the
1636                        // action as an observation or leaving the effect unresolved.
1637                        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                        // RuntimeRunner does not launch external child runners, so
1658                        // there is no host process to cancel before acknowledging.
1659                        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                            // R-B27: no verifier and no evaluation hook. The run still suspends
2051                            // with `milestone_pending`, but the `evaluate_milestone` effect MUST be
2052                            // resolved first — returning without a result leaves a dangling entry
2053                            // in the kernel's `pending_effects`, which becomes an unresolvable
2054                            // pending item once logical-checkpoint recovery lands.
2055                            //
2056                            // `MilestoneCheckResult` has no error channel on the wire today
2057                            // (Phase 1 adds one), so the most conservative shape the current
2058                            // contract can express is `passed = false` with an explanatory
2059                            // `reason`: fail-closed, the phase does not advance and no capability
2060                            // is unlocked. The returned action is intentionally dropped — this
2061                            // branch terminates the run regardless.
2062                            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            // I0a: when the loop exits without a clean kernel-done, preserve preempt intent
2175            // (interrupted flag set) in the run_terminal reason — otherwise an interrupt-curtailed
2176            // run reports "error" indistinguishable from a real crash. Mirrors Node/WASM/Python.
2177            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                // Payload residency is already durable in the canonical transaction record.
2326                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                // Local process supervision and scheduling traces are durable canonical audit
2451                // facts. The Rust SDK does not maintain a second session-log projection for them.
2452                KernelObservation::ChildSupervised { .. }
2453                | KernelObservation::LocalRunnableTrace { .. } => {}
2454                // W0-ABI workflow lifecycle. The rust SDK has no workflow drive yet
2455                // (node/python only), so these are observed-but-ignored here.
2456                KernelObservation::WorkflowBatchSpawned { .. } => {}
2457                KernelObservation::WorkflowSpawnFailed { .. } => {}
2458                KernelObservation::WorkflowCompleted { .. } => {}
2459                KernelObservation::NodesRejected { .. } => {}
2460                KernelObservation::AgentPreempted { .. } => {}
2461                KernelObservation::AgentPreemptFailed { .. } => {}
2462                KernelObservation::MemoryWriteFailed { .. } => {}
2463                KernelObservation::MemoryQueryFailed { .. } => {}
2464                // M3/M4 lifecycle observations. Durable-store mirroring is a Node/Python SDK
2465                // concern; this Rust session-log loop does not persist them (parity follow-up).
2466                KernelObservation::MemoryRecalled { .. }
2467                | KernelObservation::PromotionSuggested { .. } => {}
2468                // Governance flagged a tool call for user approval. The kernel does
2469                // not block it; the SDK-side human-approval workflow is a follow-up.
2470                KernelObservation::ToolGated { .. } => {}
2471                // In-kernel signal routing decision. The rust SDK does not yet drive
2472                // signals through the kernel attention policy; observation is logged
2473                // by the generic observation path elsewhere if needed.
2474                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                // §13.2 · live policy patches are not exposed through this SDK's public runner, so
2533                // no host path currently produces this canonical observation.
2534                KernelObservation::LivePolicyChanged { .. } => {}
2535                KernelObservation::Suspended { .. }
2536                | KernelObservation::ApprovalResolutionFailed { .. } => {}
2537                KernelObservation::Resumed { .. } => {}
2538                // R3-1: submission bookkeeping — the rust SDK has no workflow driver, so the
2539                // base-index observation has no session record to enrich here.
2540                KernelObservation::WorkflowNodesSubmitted { .. } => {}
2541                // ③ loop-agent pacing: the rust SDK has no loop driver yet; the decision also
2542                // rides LoopResult.pace_decision for embedders that want it.
2543                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                // Phase 7 / M3: no dedicated session kinds yet in rust SDK.
2585                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                // Rejections are already durable in the kernel transaction record. Call-specific
2601                // APIs inspect the observation directly; the generic runner has no host effect.
2602                KernelObservation::ControlRequestRejected { .. } => {}
2603                KernelObservation::StepPublishedEffects { .. } => {}
2604            }
2605        }
2606        next_archive_start
2607    }
2608
2609    async fn read_entries(&self, session_id: &str) -> Result<Vec<SessionEntry>> {
2610        if let Some(log) = &self.opts.session_log {
2611            log.read(session_id, 0, None).await.map_err(Error::Io)
2612        } else {
2613            Ok(Vec::new())
2614        }
2615    }
2616
2617    async fn log(&self, session_id: &str, event: SessionEvent) {
2618        if let Some(log) = &self.opts.session_log {
2619            let _ = log.append(session_id, event).await;
2620        }
2621    }
2622
2623    async fn archive_semantic_page_out(&self, archived: Vec<Message>, action: Option<String>) {
2624        let (Some(_store), Some(agent_id), Some(scope)) = (
2625            &self.opts.memory_store,
2626            &self.opts.agent_id,
2627            &self.opts.memory_scope,
2628        ) else {
2629            return;
2630        };
2631
2632        let summary = match self.summarize_for_long_term_memory(&archived).await {
2633            Ok(s) => s,
2634            Err(_) => return, // non-fatal
2635        };
2636
2637        // P2 write-funnel: route through the ONE gated write_memory syscall so validation,
2638        // the rolling write quota, dedup, and the memory_written audit all apply. Score is
2639        // advisory (0.6) — an automatic summary must never outrank curated content.
2640        let now = std::time::SystemTime::now()
2641            .duration_since(std::time::UNIX_EPOCH)
2642            .unwrap_or_default()
2643            .as_millis() as u64;
2644        let name = format!("page-out-{now}");
2645        let request = MemoryRecord {
2646            record_id: format!("{}:{}:project:{}", scope.tenant_id, scope.namespace, name),
2647            scope: scope.clone(),
2648            name,
2649            kind: MemoryKind::Project,
2650            content: summary,
2651            description: format!(
2652                "auto summary of {} archive",
2653                action.as_deref().unwrap_or("compaction")
2654            ),
2655            provenance: MemoryProvenance {
2656                session_id: self.opts.session_id.clone(),
2657                author: MemoryAuthor::Extraction,
2658                trust: MemoryTrustLevel::Untrusted,
2659                evidence_refs: Vec::new(),
2660            },
2661            created_at: now,
2662            updated_at: now,
2663            last_recalled_at: None,
2664            recall_count: 0,
2665            confidence: 0.6,
2666            links: Vec::new(),
2667            pinned: false,
2668            ttl_days: None,
2669        };
2670        let _ = self.write_memory(request, None, Some(agent_id)).await;
2671    }
2672
2673    async fn summarize_for_long_term_memory(&self, archived: &[Message]) -> crate::Result<String> {
2674        let transcript = archived
2675            .iter()
2676            .map(|m| {
2677                let role_str = match m.role {
2678                    deepstrike_core::types::message::Role::System => "system",
2679                    deepstrike_core::types::message::Role::User => "user",
2680                    deepstrike_core::types::message::Role::Assistant => "assistant",
2681                    deepstrike_core::types::message::Role::Tool => "tool",
2682                };
2683                let content_str = message_content_as_text(&m.content);
2684                format!("{}: {}", role_str, content_str)
2685            })
2686            .collect::<Vec<_>>()
2687            .join("\n");
2688
2689        let system_prompt_opt = self.opts.system_prompt.as_deref();
2690        let system_text = match system_prompt_opt {
2691            Some(sp) => format!(
2692                "{}\n\nSummarize the following conversation for long-term memory. Preserve key facts, decisions, and open questions.",
2693                sp
2694            ),
2695            None => "Summarize the following conversation for long-term memory. Preserve key facts, decisions, and open questions.".to_string(),
2696        };
2697
2698        let context = deepstrike_core::context::renderer::RenderedContext {
2699            system_text,
2700            system_stable: String::new(),
2701            system_knowledge: String::new(),
2702            turns: vec![deepstrike_core::types::message::Message {
2703                role: deepstrike_core::types::message::Role::User,
2704                content: deepstrike_core::types::message::Content::Text(transcript.clone()),
2705                tool_calls: vec![],
2706                token_count: None,
2707            }],
2708            state_turn: None,
2709            frozen_prefix_len: None,
2710            budget_overflow: None,
2711        };
2712
2713        let synth_state = self.opts.provider.create_run_state();
2714        let mut stream = self
2715            .opts
2716            .provider
2717            .stream(&context, &[], None, synth_state.as_ref())
2718            .await?;
2719
2720        let mut synthesis_text = String::new();
2721        while let Some(evt) = stream.next().await {
2722            if let Ok(StreamEvent::TextDelta { delta }) = evt {
2723                synthesis_text.push_str(&delta);
2724            }
2725        }
2726
2727        let text = synthesis_text.trim();
2728        if text.is_empty() {
2729            Ok(transcript.chars().take(2000).collect())
2730        } else {
2731            Ok(text.to_string())
2732        }
2733    }
2734}
2735
2736fn message_content_as_text(content: &deepstrike_core::types::message::Content) -> String {
2737    match content {
2738        deepstrike_core::types::message::Content::Text(s) => s.clone(),
2739        deepstrike_core::types::message::Content::Parts(parts) => parts
2740            .iter()
2741            .filter_map(|p| match p {
2742                deepstrike_core::types::message::ContentPart::Text { text } => Some(text.as_str()),
2743                deepstrike_core::types::message::ContentPart::ToolResult { output, .. } => {
2744                    Some(output.as_str())
2745                }
2746                _ => None,
2747            })
2748            .collect::<Vec<_>>()
2749            .join("\n"),
2750    }
2751}
2752
2753fn action_str_of(action: KernelPressureAction) -> String {
2754    match action {
2755        KernelPressureAction::None => "none".to_string(),
2756        KernelPressureAction::SnipCompact => "snip_compact".to_string(),
2757        KernelPressureAction::MicroCompact => "micro_compact".to_string(),
2758        KernelPressureAction::ContextCollapse => "context_collapse".to_string(),
2759        KernelPressureAction::AutoCompact => "auto_compact".to_string(),
2760    }
2761}
2762
2763pub(crate) async fn kernel_apply(
2764    kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2765    pending_observations: &mut Vec<KernelObservation>,
2766    event: serde_json::Value,
2767) -> Result<()> {
2768    let mut runtime = kernel.lock().await;
2769    canonical_kernel_apply(&mut runtime, pending_observations, event).await
2770}
2771
2772async fn kernel_transition(
2773    kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2774    pending_observations: &mut Vec<KernelObservation>,
2775    event: serde_json::Value,
2776) -> Result<Option<HostAction>> {
2777    let mut runtime = kernel.lock().await;
2778    let action = runtime.apply_host_event(event).await?;
2779    pending_observations.extend(runtime.drain_host_observations());
2780    Ok(action)
2781}
2782
2783async fn kernel_action(
2784    kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2785    pending_observations: &mut Vec<KernelObservation>,
2786    event: serde_json::Value,
2787) -> Result<HostAction> {
2788    let mut runtime = kernel.lock().await;
2789    canonical_kernel_action(&mut runtime, pending_observations, event).await
2790}
2791
2792async fn kernel_start_agent(
2793    kernel: &Arc<tokio::sync::Mutex<CanonicalRunnerRuntime>>,
2794    pending_observations: &mut Vec<KernelObservation>,
2795    task: RuntimeTask,
2796    run_spec: Option<deepstrike_core::types::agent::AgentRunSpec>,
2797) -> Result<HostAction> {
2798    let task = serde_json::to_value(task)
2799        .map_err(|error| Error::Other(format!("canonical task is not serializable: {error}")))?;
2800    let run_spec = run_spec
2801        .map(serde_json::to_value)
2802        .transpose()
2803        .map_err(|error| {
2804            Error::Other(format!("canonical run spec is not serializable: {error}"))
2805        })?;
2806    let mut runtime = kernel.lock().await;
2807    let action = runtime.start_agent_value(task, run_spec).await?;
2808    pending_observations.extend(runtime.drain_host_observations());
2809    action.ok_or_else(|| Error::Other("canonical agent root must return one host action".into()))
2810}
2811
2812pub async fn collect_text(
2813    mut stream: std::pin::Pin<Box<dyn futures::Stream<Item = Result<RunEvent>> + '_>>,
2814) -> Result<String> {
2815    let mut text = String::new();
2816    while let Some(evt) = stream.next().await {
2817        if let RunEvent::TextDelta(d) = evt? {
2818            text.push_str(&d);
2819        }
2820    }
2821    Ok(text)
2822}
2823
2824fn merge_extensions(
2825    base: Option<&serde_json::Value>,
2826    over: Option<&serde_json::Value>,
2827) -> Option<serde_json::Value> {
2828    match (base, over) {
2829        (Some(b), Some(o)) => {
2830            let mut merged = b.clone();
2831            if let (Some(m), Some(obj)) = (merged.as_object_mut(), o.as_object()) {
2832                for (k, v) in obj {
2833                    m.insert(k.clone(), v.clone());
2834                }
2835            }
2836            Some(merged)
2837        }
2838        (Some(b), None) => Some(b.clone()),
2839        (None, Some(o)) => Some(o.clone()),
2840        (None, None) => None,
2841    }
2842}
2843
2844fn cancellation_reason_code(reason: CancellationReason) -> u8 {
2845    match reason {
2846        CancellationReason::User => 0,
2847        CancellationReason::Deadline => 1,
2848        CancellationReason::LeaseLost => 2,
2849        CancellationReason::HostShutdown => 3,
2850    }
2851}
2852
2853fn provider_error_message(error: &crate::Error) -> String {
2854    match error {
2855        crate::Error::ProviderFailure(error) => error.message.clone(),
2856        _ => error.to_string(),
2857    }
2858}
2859
2860fn provider_error_event(effect_id: &str, error: &crate::Error) -> serde_json::Value {
2861    let mut event = serde_json::json!({
2862        "kind": "provider_error",
2863        "effect_id": effect_id,
2864        "message": provider_error_message(error),
2865    });
2866    if let crate::Error::ProviderFailure(error) = error {
2867        let object = event
2868            .as_object_mut()
2869            .expect("provider error event is an object");
2870        object.insert("error_kind".into(), error.kind.as_str().into());
2871        object.insert("retryable".into(), error.retryable.into());
2872        if let Some(status) = error.http_status {
2873            object.insert("http_status".into(), status.into());
2874        }
2875        if let Some(code) = error.provider_code.as_ref() {
2876            object.insert("provider_code".into(), code.clone().into());
2877        }
2878    }
2879    event
2880}
2881
2882fn cancellation_reason_from_code(code: u8) -> CancellationReason {
2883    match code {
2884        1 => CancellationReason::Deadline,
2885        2 => CancellationReason::LeaseLost,
2886        3 => CancellationReason::HostShutdown,
2887        _ => CancellationReason::User,
2888    }
2889}
2890
2891fn pending_call_ids(action: &HostAction) -> Vec<String> {
2892    match &action.effect {
2893        HostEffect::CallProvider { .. } => vec![action.effect_id.clone()],
2894        HostEffect::ExecuteTool { calls } => calls.iter().map(|call| call.id.to_string()).collect(),
2895        HostEffect::RequestApproval { requests } => requests
2896            .iter()
2897            .map(|request| request.call_id.clone())
2898            .collect(),
2899        HostEffect::SpawnWorkflow { nodes, .. } => {
2900            nodes.iter().map(|node| node.agent_id.clone()).collect()
2901        }
2902        HostEffect::PreemptSubAgents { agent_ids, .. } => agent_ids.clone(),
2903        HostEffect::Done { .. } => Vec::new(),
2904        _ => vec![action.effect_id.clone()],
2905    }
2906}
2907
2908/// Map the ergonomic [`MemoryPolicy`] onto the SDK-owned bootstrap fact.
2909fn memory_policy_host_fact(policy: MemoryPolicy) -> serde_json::Value {
2910    let mut value = serde_json::to_value(policy).expect("canonical memory policy serializes");
2911    value
2912        .as_object_mut()
2913        .expect("canonical memory policy is a JSON object")
2914        .insert("kind".into(), serde_json::json!("set_memory_policy"));
2915    value
2916}
2917
2918fn next_archived_seq_start(events: Option<&[SessionEntry]>) -> u64 {
2919    let mut next = 0u64;
2920    for entry in events.unwrap_or_default() {
2921        if let SessionEvent::Compressed {
2922            archived_seq_range, ..
2923        } = &entry.event
2924        {
2925            next = next.max(archived_seq_range.1 + 1);
2926        }
2927    }
2928    next
2929}
2930
2931fn rendered_context_from_messages(
2932    messages: Vec<Message>,
2933) -> deepstrike_core::context::renderer::RenderedContext {
2934    let mut system_parts = Vec::new();
2935    let mut turns = Vec::new();
2936    for message in messages {
2937        if message.role == deepstrike_core::types::message::Role::System {
2938            if let Some(text) = message.content.as_text() {
2939                system_parts.push(text.to_owned());
2940            }
2941        } else {
2942            turns.push(message);
2943        }
2944    }
2945    let system_text = system_parts.join("\n\n");
2946    deepstrike_core::context::renderer::RenderedContext {
2947        system_text: system_text.clone(),
2948        system_stable: system_text,
2949        system_knowledge: String::new(),
2950        turns,
2951        state_turn: None,
2952        frozen_prefix_len: None,
2953        budget_overflow: None,
2954    }
2955}
2956
2957fn parse_update_plan_args(val: &serde_json::Value) -> TaskUpdate {
2958    let plan = val.get("plan").and_then(|v| {
2959        v.as_array().map(|arr| {
2960            arr.iter()
2961                .filter_map(|x| x.as_str().map(|s| s.to_string()))
2962                .collect()
2963        })
2964    });
2965    let current_step = val
2966        .get("current_step")
2967        .or_else(|| val.get("currentStep"))
2968        .and_then(|v| v.as_u64().map(|x| x as usize));
2969    let progress = val
2970        .get("progress")
2971        .and_then(|v| v.as_str().map(|s| s.to_string()));
2972    let scratchpad = val
2973        .get("scratchpad")
2974        .and_then(|v| v.as_str().map(|s| s.to_string()));
2975    let blocked_on = val
2976        .get("blocked_on")
2977        .or_else(|| val.get("blockedOn"))
2978        .and_then(|v| {
2979            v.as_array().map(|arr| {
2980                arr.iter()
2981                    .filter_map(|x| x.as_str().map(|s| s.to_string()))
2982                    .collect()
2983            })
2984        });
2985    let preserved_refs = val
2986        .get("preserved_refs")
2987        .or_else(|| val.get("preservedRefs"))
2988        .and_then(|v| {
2989            v.as_array().map(|arr| {
2990                arr.iter()
2991                    .filter_map(|x| x.as_str().map(|s| s.to_string()))
2992                    .collect()
2993            })
2994        });
2995    TaskUpdate {
2996        plan,
2997        current_step,
2998        progress,
2999        scratchpad,
3000        blocked_on,
3001        preserved_refs,
3002        // Directives are promoted in-kernel from acted-on signals; the SDK update path leaves them
3003        // untouched here (use `..` semantics) unless a future control plane curates them explicitly.
3004        directives: None,
3005    }
3006}