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