Skip to main content

kcode_kennedy_sessions/
lib.rs

1//! Kennedy's complete logical session lifecycle and agent orchestration.
2
3#![forbid(unsafe_code)]
4
5mod services;
6
7pub use kcode_kennedy_session_objects::ResolvedObject;
8pub use kcode_telegram_session_coordinator::validate_file_name as validate_delivery_file_name;
9pub use services::{Api as Service, LocalServices as Capabilities};
10
11use std::{
12    collections::{BTreeMap, HashMap},
13    future::Future,
14    time::Duration,
15};
16
17use anyhow::Context as _;
18use chrono::{DateTime, Utc};
19use kcode_commit_session::{CommitReceipt, CommitRequest, PlannedNode};
20use kcode_dev_tools::{
21    ATTACH_OBJECT_WEB_LIB_TOOL, CALL_RUST_BIN_TOOL, RUST_BIN_TOOLS, RUST_LIB_TOOLS, WEB_LIB_TOOLS,
22    WRITE_RUST_BIN_TOOL, WRITE_RUST_LIB_TOOL, WRITE_WEB_LIB_TOOL, proposed_write_snapshot,
23};
24use kcode_dev_tools_chatend::{
25    FreeformWrite, SourceSnapshot, apply_snapshot, decode_freeform_write, prepare_freeform_write,
26};
27use kcode_history_ingress_context::{
28    Outcome as HistoryIngressContextOutcome, RecoveryOutcome as ContextRecoveryOutcome,
29};
30use kcode_kennedy_kweb_loader::{load_durable_batch, node_from_value};
31use kcode_kennedy_session_ingress::{is_terminal_external_response, restore_pending_turn};
32use kcode_kennedy_session_presentation::{RenderRequest, render};
33use kcode_kennedy_session_tool_contracts::{
34    DecodedTool, ManagedObjectArguments, ValidationRequest, decode, decode_managed_objects,
35    validate,
36};
37use kcode_kennedy_subagent_context::Context as SubagentContext;
38use kcode_kweb_context::{
39    Context as KwebContext, Node as KwebNode, NodeDraft, StagedCreate as KwebStagedCreate,
40};
41use kcode_kweb_db::NodeId;
42use kcode_server_object_envelopes::encode_file;
43use kcode_session_history::{
44    NewSession, Session as HistorySession,
45    chatend::{
46        BoxContent, BoxId, BoxOwner, ContextProjection, EventId, EventKind, SessionKind,
47        SessionMetadata,
48    },
49};
50use kcode_session_runtime_budget::{RoundBudget, RuntimeBudget, TimeBudget, TimeBudgetKind};
51use kcode_speaker_system::KTOOLS as SPEECH_CLASSIFICATION_TOOLS;
52use serde::{Deserialize, Serialize};
53use serde_json::{Value, json};
54use sha2::{Digest, Sha256};
55use uuid::Uuid;
56
57const BROWSER_CONVERSATION_REQUEST_TIMEOUT: Duration = Duration::from_secs(225 * 60);
58const HISTORY_INGRESS_REQUEST_TIMEOUT: Duration = Duration::from_secs(225 * 60);
59const WAKEUP_REQUEST_TIMEOUT: Duration = Duration::from_secs(225 * 60);
60const SELF_TIME_HARD_STOP_ALLOWANCE: Duration = Duration::from_secs(15 * 60);
61const MAX_MEDIA_ENRICHMENT_BYTES: u64 = 20 * 1024 * 1024;
62const KWEB_TOOL_INSTANCE: &str = "kweb";
63const TASK_BOARD_TOOLS: [&str; 8] = [
64    "CreateTaskCategory",
65    "GetTaskCategory",
66    "RemoveTaskCategory",
67    "CreateTask",
68    "GetTask",
69    "UpdateTask",
70    "RemoveTask",
71    "GetTopTaskOrphan",
72];
73const CONTEXT_OVERFLOW_WARNING_BOX_NAME: &str = "Context overflow warning";
74const CONTEXT_OVERFLOW_WARNING: &str = "Context size was exceeded, some context has been dehydrated. The session is now at risk of destabilizing, please perform any cleanup tasks and end the session";
75const INGRESS_FORCE_COMMIT_NOTE: &str = "ingress_force_commit";
76
77/// Application-selected primary model facts used mechanically by a session.
78#[derive(Clone, Debug)]
79pub struct RuntimeModel {
80    pub model: String,
81    pub reasoning_effort: String,
82    pub context_window_tokens: u64,
83}
84
85impl RuntimeModel {
86    pub fn from_intelligence(runtime: kcode_intelligence_router::RuntimeModel) -> Self {
87        Self {
88            model: runtime.model,
89            reasoning_effort: runtime.reasoning_effort,
90            context_window_tokens: runtime.context_window_tokens,
91        }
92    }
93
94    fn attribution(&self) -> String {
95        format!("{}-{}", self.model, self.reasoning_effort)
96    }
97}
98
99#[derive(Clone, Debug, PartialEq, Eq)]
100pub enum AgentMode {
101    Conversation,
102    FreeTime,
103    Wakeup,
104    Ingress { record_id: Option<String> },
105}
106
107/// Kind of finite enclosing deadline supplied by the workflow owner.
108#[derive(Clone, Copy, Debug, Eq, PartialEq)]
109pub enum TurnDeadlineKind {
110    /// Hard deadline for an inbound Telegram event.
111    Telegram,
112    /// Bounded cleanup deadline after user-selected self time.
113    SelfTimeHardStop,
114}
115
116/// Absolute enclosing deadline that remains transient during a turn.
117#[derive(Clone, Copy, Debug, Eq, PartialEq)]
118pub struct TurnDeadline {
119    /// Semantic category rendered in the provider footer.
120    pub kind: TurnDeadlineKind,
121    /// Absolute UTC deadline owned and enforced by orchestration.
122    pub at: DateTime<Utc>,
123}
124
125#[derive(Clone, Debug)]
126pub struct SessionOptions {
127    pub session_type: String,
128    pub root_node_ids: Vec<String>,
129    pub reference_root_node_ids: Vec<String>,
130    pub channel: Value,
131    pub free_time: Value,
132    pub orchestration: Value,
133    pub provenance_id: Option<String>,
134    pub mode: AgentMode,
135    pub source_session_type: Option<String>,
136    pub group_context: Value,
137    pub rust_lib_session_id: Option<String>,
138}
139
140impl SessionOptions {
141    pub fn conversation(session_type: impl Into<String>, roots: Vec<String>) -> Self {
142        Self {
143            session_type: session_type.into(),
144            root_node_ids: roots,
145            reference_root_node_ids: Vec::new(),
146            channel: Value::Null,
147            free_time: Value::Null,
148            orchestration: json!({"owner":"backend","status":"idle"}),
149            provenance_id: None,
150            mode: AgentMode::Conversation,
151            source_session_type: None,
152            group_context: Value::Null,
153            rust_lib_session_id: None,
154        }
155    }
156}
157
158fn restore_session_type(options: &mut SessionOptions, state: &Value) {
159    if !matches!(&options.mode, AgentMode::Ingress { .. }) {
160        options.session_type = state
161            .get("sessionType")
162            .and_then(Value::as_str)
163            .unwrap_or(&options.session_type)
164            .to_owned();
165    }
166}
167
168fn restore_commit_receipt(restored: Option<&Value>) -> anyhow::Result<Option<CommitReceipt>> {
169    restored
170        .and_then(|state| state.get("commitReceipt"))
171        .filter(|receipt| !receipt.is_null())
172        .cloned()
173        .map(serde_json::from_value)
174        .transpose()
175        .context("decoding the stored session commit receipt")
176}
177
178#[derive(Clone, Debug, Default, Deserialize, Serialize)]
179#[serde(rename_all = "camelCase")]
180struct KwebPlan {
181    creates: Vec<StagedNodeCreate>,
182    updates: BTreeMap<String, PlannedNode>,
183}
184
185#[derive(Clone, Debug, Deserialize, Serialize)]
186#[serde(rename_all = "camelCase")]
187struct StagedNodeCreate {
188    pending_id: String,
189    data: PlannedNode,
190}
191
192struct CreateNodeArguments {
193    parents: Vec<String>,
194    owner: String,
195    short_name: String,
196    short_description: String,
197    long_description: String,
198}
199
200impl KwebPlan {
201    fn restore(restored: Option<&Value>, journal: &HistorySession) -> anyhow::Result<Self> {
202        if let Some(plan) = restored.and_then(|state| state.get("kwebPlan")) {
203            return serde_json::from_value(plan.clone()).context("decoding the staged Kweb plan");
204        }
205        let latest = journal.state().events.iter().rev().find_map(|event| {
206            let EventKind::KwebPlanChanged { operation } = &event.kind else {
207                return None;
208            };
209            operation.get("plan")
210        });
211        latest
212            .cloned()
213            .map(serde_json::from_value)
214            .transpose()
215            .context("decoding the staged Kweb plan")
216            .map(Option::unwrap_or_default)
217    }
218
219    fn created(&self, id: &str) -> Option<&PlannedNode> {
220        self.creates
221            .iter()
222            .find(|create| create.pending_id == id)
223            .map(|create| &create.data)
224    }
225
226    fn created_mut(&mut self, id: &str) -> Option<&mut PlannedNode> {
227        self.creates
228            .iter_mut()
229            .find(|create| create.pending_id == id)
230            .map(|create| &mut create.data)
231    }
232}
233
234pub struct Session {
235    api: Service,
236    runtime: RuntimeModel,
237    journal: HistorySession,
238    plan: KwebPlan,
239    pub session_type: String,
240    pub channel: Value,
241    pub free_time: Value,
242    pub orchestration: Value,
243    pub provenance_id: Option<String>,
244    pub rust_lib_session_id: String,
245    pub root_node_ids: Vec<String>,
246    pub reference_root_node_ids: Vec<String>,
247    pub started_at: String,
248    pub transcript: Vec<Value>,
249    pub pending_turn: bool,
250    pub pending_external_event_id: Option<String>,
251    pub completed: bool,
252    pub rounds_used: u64,
253    commit_receipt: Option<CommitReceipt>,
254    commit_author: String,
255    mode: AgentMode,
256    source_session_type: Option<String>,
257    group_context: Value,
258    context: KwebContext,
259    free_time_end_reason: Option<String>,
260    fatal_persistence_error: Option<String>,
261    active_provider_deadline: Option<DateTime<Utc>>,
262    active_turn_deadline: Option<TurnDeadline>,
263}
264
265#[derive(Clone, Copy, Debug, Eq, PartialEq)]
266enum InputStage {
267    Accepted,
268}
269
270#[derive(Clone, Copy, Debug, Eq, PartialEq)]
271enum ContextRecovery {
272    NotNeeded,
273    Recovered,
274    Irreducible,
275}
276
277fn render_load_nodes_result(
278    journal: &HistorySession,
279    changed_box_ids: &[BoxId],
280    footer_lines: &[String],
281) -> anyhow::Result<String> {
282    let changed_box_ids = changed_box_ids
283        .iter()
284        .map(ToString::to_string)
285        .collect::<Vec<_>>();
286    let projected_boxes = journal
287        .state()
288        .projection_with_footer_lines(footer_lines)
289        .items
290        .into_iter()
291        .filter(|item| !item.marker)
292        .map(|item| (item.box_id.to_string(), item.text))
293        .collect::<Vec<_>>();
294    render(RenderRequest::LoadNodes {
295        changed_box_ids: &changed_box_ids,
296        projected_boxes: &projected_boxes,
297    })
298}
299
300fn provider_tool_result_with_context_footer(footer: &str, result: &str) -> String {
301    render(RenderRequest::ProviderFooter { result, footer })
302        .expect("provider-footer rendering is infallible")
303}
304
305fn append_slow_tool_duration(text: &mut String, elapsed: Duration) {
306    *text = render(RenderRequest::SlowTool { text, elapsed })
307        .expect("slow-tool rendering is infallible");
308}
309
310fn render_web_search_result(
311    result: &kcode_intelligence_router::SearchResponse,
312) -> anyhow::Result<String> {
313    let sources = result
314        .sources
315        .iter()
316        .map(|source| (source.title.clone(), source.url.clone()))
317        .collect::<Vec<_>>();
318    render(RenderRequest::WebSearch {
319        answer: &result.answer,
320        sources: &sources,
321    })
322}
323
324fn render_web_fetch_result(
325    result: &kcode_intelligence_router::FetchResponse,
326) -> anyhow::Result<String> {
327    render(RenderRequest::WebFetch {
328        url: &result.url,
329        title: result.title.as_deref(),
330        content_type: &result.content_type,
331        truncated: result.truncated,
332        content: &result.content,
333    })
334}
335
336fn render_media_annotation_result(
337    object_id: &str,
338    file_name: &str,
339    content_type: &str,
340    result: &kcode_intelligence_router::AnnotationResponse,
341) -> anyhow::Result<String> {
342    render(RenderRequest::MediaAnnotation {
343        object_id,
344        file_name,
345        content_type,
346        model: &result.model,
347        complete: result.complete,
348        incomplete_reason: result.incomplete_reason.as_deref(),
349        text: &result.text,
350    })
351}
352
353fn render_audio_transcription_result(
354    object_id: &str,
355    file_name: &str,
356    content_type: &str,
357    result: &kcode_intelligence_router::TranscriptionResponse,
358) -> anyhow::Result<String> {
359    render(RenderRequest::AudioTranscription {
360        object_id,
361        file_name,
362        content_type,
363        model: &result.model,
364        text: &result.text,
365    })
366}
367
368fn render_document_extraction_result(
369    object_id: &str,
370    file_name: &str,
371    result: &kcode_intelligence_router::DocumentExtraction,
372) -> anyhow::Result<String> {
373    render(RenderRequest::DocumentExtraction {
374        object_id,
375        file_name,
376        format: &result.format,
377        characters: result.characters,
378        truncated: result.truncated,
379        text: &result.text,
380    })
381}
382
383struct ToolCall {
384    name: String,
385    arguments: Value,
386}
387
388#[derive(Deserialize)]
389#[serde(rename_all = "camelCase", deny_unknown_fields)]
390struct TaskId {
391    task_id: String,
392}
393
394#[derive(Deserialize)]
395#[serde(rename_all = "camelCase", deny_unknown_fields)]
396struct CategoryId {
397    category_id: String,
398}
399
400#[derive(Deserialize)]
401#[serde(rename_all = "camelCase", deny_unknown_fields)]
402struct CategoryCall {
403    category_id: String,
404    #[serde(default)]
405    offset: u64,
406    #[serde(default = "task_page_limit")]
407    limit: u32,
408}
409
410#[derive(Deserialize)]
411#[serde(deny_unknown_fields)]
412struct EmptyCall {}
413
414fn task_page_limit() -> u32 {
415    50
416}
417
418struct RecordedToolInvocation {
419    invocation_id: String,
420    tool_instance: String,
421    tool_name: String,
422}
423
424struct PendingFreeformWrite {
425    request: FreeformWrite,
426    call_box_id: BoxId,
427}
428
429fn tool_call_box_content(call: &ToolCall) -> anyhow::Result<BoxContent> {
430    if matches!(
431        call.name.as_str(),
432        WRITE_RUST_LIB_TOOL | WRITE_WEB_LIB_TOOL | WRITE_RUST_BIN_TOOL
433    ) {
434        let name = call
435            .arguments
436            .get("name")
437            .and_then(Value::as_str)
438            .map(|name| name.chars().take(255).collect::<String>());
439        let file_count = call
440            .arguments
441            .get("files")
442            .and_then(Value::as_array)
443            .map(Vec::len);
444        return Ok(BoxContent {
445            text: serde_json::to_string_pretty(&json!({
446                "name":call.name,
447                "arguments":{
448                    "name":name,
449                    "fileCount":file_count,
450                    "completeFileContents":"omitted from active context; retained in the durable tool invocation"
451                }
452            }))?,
453            objects: Vec::new(),
454            metadata: json!({
455                "compactedToolInvocation":true,
456                "toolName":call.name,
457            }),
458        });
459    }
460    Ok(BoxContent::text(serde_json::to_string_pretty(
461        &json!({"name":call.name,"arguments":call.arguments}),
462    )?))
463}
464
465struct ToolOutcome {
466    text: String,
467    store_result: bool,
468    ok: bool,
469    end_session: bool,
470    freeform_write: Option<FreeformWrite>,
471    managed_source_snapshot: Option<SourceSnapshot>,
472}
473
474fn result_displays_snapshot(result: &str, snapshot: &SourceSnapshot) -> bool {
475    result == snapshot.text
476}
477
478fn subagent_managed_write_fits(
479    context: &SubagentContext,
480    call: &ToolCall,
481    budget: &kcode_agent_runtime::ContextBudget,
482) -> bool {
483    let Some(snapshot) = proposed_write_snapshot(&call.name, &call.arguments) else {
484        return true;
485    };
486    let state = context.source_state(&snapshot);
487    budget.fits_state(state.key, state.text)
488}
489
490struct KennedySubagentHost<'a> {
491    session: &'a mut Session,
492    context: SubagentContext,
493    captures: HashMap<String, FreeformWrite>,
494}
495
496struct KennedySessionHost<'a, C> {
497    session: &'a mut Session,
498    checkpoint: &'a mut C,
499    accounting: Option<kcode_intelligence_chatend::TopLevelCall>,
500    pending_freeform_write: Option<PendingFreeformWrite>,
501    deadline_after_response: bool,
502}
503
504fn is_kweb_mutation(name: &str) -> bool {
505    matches!(
506        name,
507        "ConnectNodes" | "ConsolidateFanout" | "SetFixedConnection" | "CreateNode" | "UpdateNode"
508    )
509}
510
511fn referenced_pending_nodes(decoded: &DecodedTool) -> Vec<String> {
512    let mut ids = Vec::new();
513    match decoded {
514        DecodedTool::ConnectNodes(nodes) => ids.extend(nodes.iter()),
515        DecodedTool::ConsolidateFanout {
516            parent,
517            fanout,
518            aggregator,
519        } => {
520            ids.push(parent);
521            ids.extend(fanout.iter());
522            ids.push(aggregator);
523        }
524        DecodedTool::SetFixedConnection {
525            parent,
526            child,
527            slot: _,
528        } => {
529            ids.push(parent);
530            ids.extend(child.iter());
531        }
532        DecodedTool::CreateNode { parents, owner, .. } => {
533            ids.extend(parents.iter());
534            ids.push(owner);
535        }
536        DecodedTool::UpdateNode { id, owner, .. } => {
537            ids.push(id);
538            ids.push(owner);
539        }
540        _ => {}
541    }
542    ids.into_iter()
543        .filter(|id| id.starts_with("pending:"))
544        .cloned()
545        .collect()
546}
547
548fn subagent_unavailable_reason(name: &str) -> Option<&'static str> {
549    match name {
550        "RunSubagent" => {
551            Some("RunSubagent is unavailable inside a subagent. Only Kennedy may launch subagents.")
552        }
553        "EndSession" => Some(
554            "EndSession is unavailable inside a subagent. A child cannot control the parent session lifecycle.",
555        ),
556        "DehydrateBoxes" | "SummarizeBox" | "HydrateBox" | "BoxesIntoObjects" => {
557            Some("Parent box controls are unavailable inside a box-free subagent context.")
558        }
559        _ => None,
560    }
561}
562
563fn execute_kweb_mutation(
564    name: &str,
565    decoded: DecodedTool,
566    context: &KwebContext,
567    plan: &mut KwebPlan,
568    journal: &mut HistorySession,
569) -> anyhow::Result<String> {
570    match (name, decoded) {
571        ("ConnectNodes", DecodedTool::ConnectNodes(ids)) => {
572            Session::connect_nodes(plan, context, ids)
573        }
574        (
575            "ConsolidateFanout",
576            DecodedTool::ConsolidateFanout {
577                parent,
578                fanout,
579                aggregator,
580            },
581        ) => Session::consolidate_fanout(plan, context, parent, fanout, aggregator),
582        (
583            "SetFixedConnection",
584            DecodedTool::SetFixedConnection {
585                parent,
586                child,
587                slot,
588            },
589        ) => Session::set_fixed_connection(plan, context, parent, child, slot),
590        (
591            "CreateNode",
592            DecodedTool::CreateNode {
593                parents,
594                owner,
595                short_name,
596                short_description,
597                long_description,
598            },
599        ) => Session::create_node(
600            plan,
601            context,
602            journal,
603            CreateNodeArguments {
604                parents,
605                owner,
606                short_name,
607                short_description,
608                long_description,
609            },
610        ),
611        (
612            "UpdateNode",
613            DecodedTool::UpdateNode {
614                id,
615                owner,
616                short_name,
617                short_description,
618                long_description,
619            },
620        ) => Session::update_node(
621            plan,
622            context,
623            id,
624            owner,
625            short_name,
626            short_description,
627            long_description,
628        ),
629        _ => anyhow::bail!("decoded contract for {name} did not match its Kweb mutation"),
630    }
631}
632
633fn kweb_plan_projection(plan: &KwebPlan) -> (BTreeMap<String, NodeDraft>, Vec<KwebStagedCreate>) {
634    let updates = plan
635        .updates
636        .iter()
637        .map(|(id, node)| (id.clone(), kweb_node_draft(node)))
638        .collect();
639    let creates = plan
640        .creates
641        .iter()
642        .map(|create| KwebStagedCreate {
643            pending_id: create.pending_id.clone(),
644            data: kweb_node_draft(&create.data),
645        })
646        .collect();
647    (updates, creates)
648}
649
650impl Session {
651    fn runtime_budget(&self) -> RuntimeBudget {
652        let Some(provider_deadline) = self.active_provider_deadline else {
653            return RuntimeBudget::default();
654        };
655        let mut time_limits = vec![TimeBudget {
656            kind: TimeBudgetKind::ProviderCall,
657            remaining: remaining_until(provider_deadline),
658        }];
659        if matches!(self.mode, AgentMode::FreeTime)
660            && let Some(work_deadline) = deadline(&self.free_time)
661        {
662            time_limits.push(TimeBudget {
663                kind: TimeBudgetKind::SelfTimeWork,
664                remaining: remaining_until(work_deadline),
665            });
666        }
667        if let Some(outer) = self.active_turn_deadline {
668            time_limits.push(TimeBudget {
669                kind: match outer.kind {
670                    TurnDeadlineKind::Telegram => TimeBudgetKind::TelegramTurn,
671                    TurnDeadlineKind::SelfTimeHardStop => TimeBudgetKind::SelfTimeHardStop,
672                },
673                remaining: remaining_until(outer.at),
674            });
675        }
676        RuntimeBudget {
677            rounds: Some(RoundBudget {
678                used: self.rounds_used,
679                limit: kcode_agent_runtime::DEFAULT_ROUND_LIMIT,
680            }),
681            time_limits,
682        }
683    }
684
685    fn projection(&self) -> ContextProjection {
686        self.journal
687            .state()
688            .projection_with_footer_lines(&self.runtime_budget().footer_lines())
689    }
690
691    fn begin_provider_call_budget(&mut self, timeout: Option<Duration>) {
692        self.active_provider_deadline = timeout.and_then(|timeout| {
693            chrono::Duration::from_std(timeout)
694                .ok()
695                .map(|timeout| Utc::now() + timeout)
696        });
697    }
698
699    pub async fn new(
700        api: Service,
701        system_prompt: String,
702        runtime: RuntimeModel,
703        mut options: SessionOptions,
704        restored: Option<&Value>,
705    ) -> anyhow::Result<Self> {
706        if let Some(state) = restored {
707            restore_session_type(&mut options, state);
708            options.channel = state.get("channel").cloned().unwrap_or(options.channel);
709            options.free_time = state.get("freeTime").cloned().unwrap_or(options.free_time);
710            options.orchestration = state
711                .get("orchestration")
712                .cloned()
713                .unwrap_or(options.orchestration);
714        }
715        if options.group_context.is_null() {
716            options.group_context = options
717                .channel
718                .get("groupContext")
719                .cloned()
720                .unwrap_or(Value::Null);
721        }
722        options
723            .reference_root_node_ids
724            .retain(|id| !options.root_node_ids.contains(id));
725        options.reference_root_node_ids.sort();
726        options.reference_root_node_ids.dedup();
727
728        let started_at = restored
729            .and_then(|state| state.get("startedAt"))
730            .and_then(Value::as_str)
731            .map(str::to_owned)
732            .unwrap_or_else(|| Utc::now().to_rfc3339());
733        let rust_lib_session_id = restored
734            .and_then(|state| state.get("rustLibSessionId"))
735            .and_then(Value::as_str)
736            .map(str::to_owned)
737            .or(options.rust_lib_session_id.clone())
738            .unwrap_or_else(|| format!("kennedy:{}", Uuid::new_v4()));
739        let history_session_id = restored
740            .and_then(|state| state.get("sessionId"))
741            .and_then(Value::as_str)
742            .map(str::to_owned);
743        let source_session_type = options.source_session_type.clone().or_else(|| {
744            restored
745                .and_then(|state| state.get("sourceSessionType"))
746                .and_then(Value::as_str)
747                .map(str::to_owned)
748        });
749        let session_id = history_session_id
750            .clone()
751            .unwrap_or_else(|| Uuid::new_v4().to_string());
752        let metadata = SessionMetadata {
753            session_id: session_id.clone(),
754            kind: session_kind(&options.session_type, &options.mode),
755            created_at: started_at.clone(),
756            effective_context_tokens: runtime.context_window_tokens,
757            channel: options.channel.clone(),
758        };
759        let mut journal = if history_session_id.is_some() {
760            api.history_session(metadata, &runtime.model)
761                .with_context(|| {
762                    format!(
763                        "opening authoritative session {session_id} (legacy snapshots are intentionally unsupported)"
764                    )
765                })?
766        } else {
767            api.create_history_session(NewSession {
768                kind: metadata.kind,
769                created_at: metadata.created_at,
770                effective_context_tokens: metadata.effective_context_tokens,
771                channel: metadata.channel,
772            })?
773        };
774        let mut context = KwebContext::with_fixed_connections(
775            options.root_node_ids.clone(),
776            api.loads_fixed_connections(),
777        )
778        .map_err(anyhow::Error::new)?;
779        restore_kweb_context(&journal, &mut context)?;
780        let plan = KwebPlan::restore(restored, &journal)?;
781        let transcript = kcode_kennedy_session_ingress::transcript_from_journal(&journal);
782        let (pending_turn, pending_external_event_id) = restore_pending_turn(restored, &transcript);
783
784        let needs_initialization = !journal
785            .state()
786            .boxes
787            .values()
788            .any(|state| matches!(state.owner, BoxOwner::System));
789        let commit_receipt = restore_commit_receipt(restored)?;
790        let commit_author = restored
791            .and_then(|state| state.get("commitAuthor"))
792            .and_then(Value::as_str)
793            .map(str::to_owned)
794            .unwrap_or_else(|| runtime.attribution());
795        if let Some(receipt) = &commit_receipt {
796            journal.mark_completed(receipt.session_object_id.to_string());
797        }
798        let completed =
799            journal.state().completed_session_object.is_some() || commit_receipt.is_some();
800        let mut session = Self {
801            api,
802            runtime,
803            journal,
804            plan,
805            session_type: options.session_type,
806            channel: options.channel,
807            free_time: options.free_time,
808            orchestration: options.orchestration,
809            provenance_id: options.provenance_id,
810            rust_lib_session_id,
811            root_node_ids: options.root_node_ids,
812            reference_root_node_ids: options.reference_root_node_ids,
813            started_at,
814            transcript,
815            pending_turn,
816            pending_external_event_id,
817            completed,
818            rounds_used: restored
819                .and_then(|state| state.get("roundsUsed"))
820                .and_then(Value::as_u64)
821                .unwrap_or_default(),
822            commit_receipt,
823            commit_author,
824            mode: options.mode,
825            source_session_type,
826            group_context: options.group_context,
827            context,
828            free_time_end_reason: None,
829            fatal_persistence_error: None,
830            active_provider_deadline: None,
831            active_turn_deadline: None,
832        };
833
834        if matches!(session.mode, AgentMode::Ingress { .. }) && !session.journal.is_sealed() {
835            session.journal.repair_unfinished_tools(now())?;
836        }
837        if session.journal.is_sealed() {
838            anyhow::ensure!(
839                !matches!(session.mode, AgentMode::Conversation),
840                "a read-only conversation has an unexpectedly sealed session log"
841            );
842            if session.commit_receipt.is_none() {
843                session.finalize_kweb_session()?;
844            }
845            session.completed = true;
846            return Ok(session);
847        }
848
849        if needs_initialization {
850            session.journal.create_box(
851                now(),
852                "Kennedy system prompt",
853                BoxOwner::System,
854                BoxContent::text(&system_prompt),
855            )?;
856            if session.session_type == "telegram-group" && !session.group_context.is_null() {
857                session.journal.create_box(
858                    now(),
859                    "Telegram group context",
860                    BoxOwner::Controller,
861                    BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
862                        &session.group_context,
863                    )),
864                )?;
865            }
866            let roots = session.root_node_ids.clone();
867            let invocation =
868                session.record_tool_invocation("LoadNodes", json!({"identifiers":&roots}))?;
869            let result = load_durable_batch(session.api.kmap(), &mut session.context, &roots)?;
870            session.sync_kweb_boxes()?;
871            session.record_tool_completion(
872                Some(&invocation),
873                json!({"ok":true,"automatic":true,"identifiers":roots,"result":result}),
874            )?;
875        } else {
876            session.sync_kweb_boxes()?;
877        }
878        if matches!(session.mode, AgentMode::Ingress { .. })
879            && !session.completed
880            && !session.journal.state().history_ingress_started
881        {
882            session.prepare_history_ingress(&system_prompt).await?;
883        }
884        Ok(session)
885    }
886
887    async fn prepare_history_ingress(&mut self, prompt: &str) -> anyhow::Result<()> {
888        let cost_at_ingress = self.projection().status;
889        if !self.journal.state().source_terminated {
890            self.journal.record(
891                now(),
892                EventKind::SourceTerminated {
893                    reason: "history_ingress".into(),
894                },
895            )?;
896        }
897        let system_box = self
898            .journal
899            .state()
900            .boxes
901            .values()
902            .find(|state| matches!(state.owner, BoxOwner::System))
903            .map(|state| state.id)
904            .context("session has no system-prompt box")?;
905        self.journal
906            .update_box(now(), system_box, BoxContent::text(prompt))?;
907        let ingress_kind = session_kind(&self.session_type, &self.mode);
908        if self.journal.state().metadata.effective_context_tokens
909            != self.runtime.context_window_tokens
910            || self.journal.state().metadata.kind != ingress_kind
911        {
912            self.journal
913                .configure_context(ingress_kind, self.runtime.context_window_tokens);
914        }
915        self.journal.create_box(
916            now(),
917            "Session cost at ingress",
918            BoxOwner::Controller,
919            BoxContent::text(cost_summary(
920                "session cost before history ingress",
921                cost_at_ingress.estimated_cost_usd_nanos,
922                cost_at_ingress.unpriced_provider_calls,
923            )),
924        )?;
925        self.revalidate_loaded_nodes().await?;
926        match kcode_history_ingress_context::prepare(&mut self.journal, now())? {
927            HistoryIngressContextOutcome::Ready => {}
928            HistoryIngressContextOutcome::OverCapacity {
929                estimated_tokens,
930                target_tokens,
931            } => {
932                self.journal.record(
933                    now(),
934                    EventKind::Note {
935                        label: INGRESS_FORCE_COMMIT_NOTE.into(),
936                        value: json!({
937                            "reason":"fully_dehydrated_context_above_initial_target",
938                            "estimatedTokens":estimated_tokens,
939                            "initialTargetTokens":target_tokens,
940                        }),
941                    },
942                )?;
943                self.pending_turn = false;
944                self.finalize_kweb_session()?;
945                self.completed = true;
946                return Ok(());
947            }
948        }
949        self.journal
950            .record(now(), EventKind::HistoryIngressStarted)?;
951        self.pending_turn = true;
952        Ok(())
953    }
954
955    async fn revalidate_loaded_nodes(&mut self) -> anyhow::Result<()> {
956        let direct = self.context.loaded_node_ids().to_vec();
957        load_durable_batch(self.api.kmap(), &mut self.context, &direct)?;
958        self.sync_kweb_boxes()?;
959        Ok(())
960    }
961
962    fn stage_user_input(&mut self, text: &str, metadata: &Value) -> Option<InputStage> {
963        let recorded_at = now();
964        let result = (|| -> anyhow::Result<Option<InputStage>> {
965            let Some(staged) = kcode_kennedy_session_ingress::stage_user_input(
966                &mut self.journal,
967                text,
968                metadata,
969                &recorded_at,
970            )?
971            else {
972                return Ok(None);
973            };
974            self.transcript.push(staged.transcript);
975            self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
976            Ok(Some(InputStage::Accepted))
977        })();
978        match result {
979            Ok(stage) => stage,
980            Err(error) => {
981                self.fatal_persistence_error = Some(error.to_string());
982                tracing::error!(error=%error, "Could not durably stage session input");
983                Some(InputStage::Accepted)
984            }
985        }
986    }
987
988    pub fn append_final_user_message(&mut self, text: &str, metadata: &Value) -> bool {
989        self.stage_user_input(text, metadata).is_some()
990    }
991
992    pub fn stage_source_message(
993        &mut self,
994        kennedy: bool,
995        text: &str,
996        metadata: Value,
997    ) -> anyhow::Result<()> {
998        let staged = kcode_kennedy_session_ingress::stage_source_input(
999            &mut self.journal,
1000            kennedy,
1001            text,
1002            metadata,
1003            &now(),
1004        )?;
1005        self.transcript.push(staged.transcript);
1006        self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
1007        Ok(())
1008    }
1009
1010    pub fn answer_for_external_event(&self, id: &str) -> Option<&Value> {
1011        self.transcript.iter().rev().find(|entry| {
1012            is_terminal_external_response(entry)
1013                && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
1014        })
1015    }
1016
1017    pub fn responses_for_external_event(&self, id: &str) -> Vec<&Value> {
1018        self.transcript
1019            .iter()
1020            .filter(|entry| {
1021                matches!(
1022                    entry.get("role").and_then(Value::as_str),
1023                    Some("kennedy" | "system")
1024                ) && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
1025            })
1026            .collect()
1027    }
1028
1029    pub fn resolve_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
1030        let api = self.api.clone();
1031        kcode_kennedy_session_objects::resolve_object(
1032            &mut self.journal,
1033            object_id,
1034            move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
1035        )
1036    }
1037
1038    fn resolve_media_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
1039        let api = self.api.clone();
1040        kcode_kennedy_session_objects::resolve_media_object(
1041            &mut self.journal,
1042            object_id,
1043            MAX_MEDIA_ENRICHMENT_BYTES,
1044            move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
1045        )
1046    }
1047
1048    fn resolve_image_object(
1049        &mut self,
1050        object_id: &str,
1051    ) -> anyhow::Result<(Vec<u8>, String, String)> {
1052        let resolved = self.resolve_media_object(object_id)?;
1053        anyhow::ensure!(
1054            resolved.media_type.starts_with("image/"),
1055            "GenerateImage reference {object_id} is not an image"
1056        );
1057        Ok((resolved.bytes, resolved.file_name, resolved.media_type))
1058    }
1059
1060    fn recover_context_overflow(
1061        &mut self,
1062        external_event_id: Option<&str>,
1063        pinned_box_ids: &[BoxId],
1064    ) -> anyhow::Result<ContextRecovery> {
1065        let projection = self.projection();
1066        let target_tokens = self.journal.state().active_context_limit();
1067        if projection.estimated_tokens <= target_tokens {
1068            return Ok(ContextRecovery::NotNeeded);
1069        }
1070        let projection_hash = hex::encode(Sha256::digest(projection.render().as_bytes()));
1071        let already_irreducible = self
1072            .journal
1073            .state()
1074            .events
1075            .iter()
1076            .rev()
1077            .find_map(|event| match &event.kind {
1078                EventKind::Note { label, value } if label == "context_overflow_recovery" => {
1079                    Some(value)
1080                }
1081                _ => None,
1082            })
1083            .is_some_and(|value| {
1084                value.get("irreducible").and_then(Value::as_bool) == Some(true)
1085                    && value.get("limitTokens").and_then(Value::as_u64) == Some(target_tokens)
1086                    && value.get("projectionHash").and_then(Value::as_str)
1087                        == Some(projection_hash.as_str())
1088            });
1089        if already_irreducible {
1090            return Ok(ContextRecovery::Irreducible);
1091        }
1092
1093        let before_tokens = projection.estimated_tokens;
1094        let mut metadata = json!({
1095            "transcriptRole":"system",
1096            "contextOverflowWarning":true,
1097            "projectedTokens":before_tokens,
1098            "limitTokens":target_tokens,
1099        });
1100        if let Some(id) = external_event_id {
1101            metadata["externalEventId"] = json!(id);
1102        }
1103        let warning_box_id = self.journal.create_box(
1104            now(),
1105            CONTEXT_OVERFLOW_WARNING_BOX_NAME,
1106            BoxOwner::Controller,
1107            BoxContent {
1108                text: CONTEXT_OVERFLOW_WARNING.into(),
1109                objects: Vec::new(),
1110                metadata,
1111            },
1112        )?;
1113        let mut transcript = json!({
1114            "role":"system",
1115            "content":CONTEXT_OVERFLOW_WARNING,
1116            "contextOverflowWarning":true,
1117        });
1118        if let Some(id) = external_event_id {
1119            transcript["externalEventId"] = json!(id);
1120        }
1121        self.transcript.push(transcript);
1122
1123        let mut pins = pinned_box_ids.to_vec();
1124        if !pins.contains(&warning_box_id) {
1125            pins.push(warning_box_id);
1126        }
1127        let outcome = kcode_history_ingress_context::recover(&mut self.journal, now(), &pins)?;
1128        let (dehydrated_box_ids, estimated_tokens, target_tokens, irreducible) = match outcome {
1129            ContextRecoveryOutcome::Recovered {
1130                dehydrated_box_ids,
1131                estimated_tokens,
1132                target_tokens,
1133            } => (dehydrated_box_ids, estimated_tokens, target_tokens, false),
1134            ContextRecoveryOutcome::OverCapacity {
1135                dehydrated_box_ids,
1136                estimated_tokens,
1137                target_tokens,
1138            } => (dehydrated_box_ids, estimated_tokens, target_tokens, true),
1139        };
1140        let final_projection_hash =
1141            hex::encode(Sha256::digest(self.projection().render().as_bytes()));
1142        self.journal.record(
1143            now(),
1144            EventKind::Note {
1145                label: "context_overflow_recovery".into(),
1146                value: json!({
1147                    "beforeTokens":before_tokens,
1148                    "estimatedTokens":estimated_tokens,
1149                    "limitTokens":target_tokens,
1150                    "dehydratedBoxIds":dehydrated_box_ids,
1151                    "irreducible":irreducible,
1152                    "projectionHash":final_projection_hash,
1153                }),
1154            },
1155        )?;
1156        if irreducible {
1157            if matches!(self.mode, AgentMode::Ingress { .. }) {
1158                self.request_ingress_force_commit(
1159                    "irreducible_context_overflow",
1160                    estimated_tokens,
1161                )?;
1162            } else if !self.journal.state().source_terminated {
1163                self.journal.record(
1164                    now(),
1165                    EventKind::SourceTerminated {
1166                        reason: "irreducible_context_overflow".into(),
1167                    },
1168                )?;
1169            }
1170            Ok(ContextRecovery::Irreducible)
1171        } else {
1172            Ok(ContextRecovery::Recovered)
1173        }
1174    }
1175
1176    fn request_ingress_force_commit(
1177        &mut self,
1178        reason: &str,
1179        projected_tokens: u64,
1180    ) -> anyhow::Result<()> {
1181        if self.ingress_force_commit_requested() {
1182            return Ok(());
1183        }
1184        self.journal.record(
1185            now(),
1186            EventKind::Note {
1187                label: INGRESS_FORCE_COMMIT_NOTE.into(),
1188                value: json!({
1189                    "reason":reason,
1190                    "projectedTokens":projected_tokens,
1191                    "limitTokens":self.journal.state().ingress_context_limit(),
1192                }),
1193            },
1194        )?;
1195        Ok(())
1196    }
1197
1198    fn ingress_force_commit_requested(&self) -> bool {
1199        self.journal.state().events.iter().rev().any(|event| {
1200            matches!(
1201                &event.kind,
1202                EventKind::Note { label, .. } if label == INGRESS_FORCE_COMMIT_NOTE
1203            )
1204        })
1205    }
1206
1207    pub fn requires_history_ingress(&self) -> bool {
1208        matches!(self.mode, AgentMode::Conversation) && self.journal.state().source_terminated
1209    }
1210
1211    pub fn stage_free_time_opening(&mut self) -> bool {
1212        if self.pending_turn {
1213            return false;
1214        }
1215        let mut blocks = vec![
1216            render(RenderRequest::FreeTimeOpening {
1217                free_time: &self.free_time,
1218            })
1219            .expect("free-time opening rendering is infallible"),
1220        ];
1221        if let Some(message) = self
1222            .free_time
1223            .get("handoffMessage")
1224            .and_then(Value::as_str)
1225            .filter(|message| !message.trim().is_empty())
1226        {
1227            blocks.push(format!(
1228                "Message from the previous self-time session:\n\n{message}"
1229            ));
1230        }
1231        let Some(stage) = self.stage_user_input(&blocks.join("\n\n"), &json!({"kind":"self-time"}))
1232        else {
1233            return false;
1234        };
1235        self.pending_turn = matches!(stage, InputStage::Accepted);
1236        true
1237    }
1238
1239    pub fn stage_wakeup_opening(&mut self) -> anyhow::Result<bool> {
1240        if self.pending_turn {
1241            return Ok(false);
1242        }
1243        let marker = self
1244            .channel
1245            .get("wakeupMarker")
1246            .and_then(Value::as_str)
1247            .context("wakeup session is missing its acquired time marker")?;
1248        let marker = DateTime::parse_from_rfc3339(marker)
1249            .context("wakeup session has an invalid acquired time marker")?
1250            .with_timezone(&Utc);
1251        let text = render(RenderRequest::WakeupOpening { marker })?;
1252        let Some(stage) = self.stage_user_input(
1253            &text,
1254            &json!({"kind":"wakeup","wakeupMarker":marker.to_rfc3339()}),
1255        ) else {
1256            return Ok(false);
1257        };
1258        self.pending_turn = matches!(stage, InputStage::Accepted);
1259        Ok(true)
1260    }
1261
1262    pub fn begin_user_turn(&mut self, text: &str, metadata: &Value) -> bool {
1263        if self.pending_turn {
1264            return false;
1265        }
1266        let Some(stage) = self.stage_user_input(text, metadata) else {
1267            return false;
1268        };
1269        debug_assert_eq!(stage, InputStage::Accepted);
1270        self.rounds_used = 0;
1271        self.pending_turn = true;
1272        self.pending_external_event_id = metadata
1273            .get("externalEventId")
1274            .and_then(Value::as_str)
1275            .map(str::to_owned);
1276        true
1277    }
1278
1279    pub fn reset_exhausted_turn_rounds_for_retry(&mut self) {
1280        if matches!(self.mode, AgentMode::Conversation)
1281            && self.rounds_used >= kcode_agent_runtime::DEFAULT_ROUND_LIMIT
1282        {
1283            self.rounds_used = 0;
1284        }
1285    }
1286
1287    pub fn interrupt_current_turn(&mut self) -> anyhow::Result<()> {
1288        self.journal.repair_unfinished_tools(now())?;
1289        let notice = "The user stopped this agent turn.";
1290        let mut metadata = json!({"transcriptRole":"system","userStopped":true});
1291        let mut transcript_entry = json!({
1292            "role":"system",
1293            "content":notice,
1294            "userStopped":true,
1295        });
1296        if let Some(external_event_id) = &self.pending_external_event_id {
1297            metadata["externalEventId"] = json!(external_event_id);
1298            transcript_entry["externalEventId"] = json!(external_event_id);
1299        }
1300        self.journal.create_box(
1301            now(),
1302            "Turn stopped",
1303            BoxOwner::Controller,
1304            BoxContent {
1305                text: notice.into(),
1306                objects: Vec::new(),
1307                metadata,
1308            },
1309        )?;
1310        self.transcript.push(transcript_entry);
1311        self.pending_turn = false;
1312        self.pending_external_event_id = None;
1313        self.orchestration =
1314            json!({"owner":"backend","status":"idle","lastOutcome":"user-stopped"});
1315        Ok(())
1316    }
1317
1318    pub async fn run_pending_turn<C, F>(
1319        &mut self,
1320        operation_id: Uuid,
1321        turn_deadline: Option<TurnDeadline>,
1322        mut checkpoint: C,
1323    ) -> anyhow::Result<Option<String>>
1324    where
1325        C: FnMut(Value) -> F + Send,
1326        F: Future<Output = anyhow::Result<()>> + Send,
1327    {
1328        if let Some(error) = self.fatal_persistence_error.take() {
1329            anyhow::bail!("session journal write failed: {error}");
1330        }
1331        if !self.pending_turn {
1332            return Ok(None);
1333        }
1334        self.active_turn_deadline = turn_deadline;
1335        let runtime = self.api.agent_runtime();
1336        let user_id = self
1337            .root_node_ids
1338            .first()
1339            .context("session has no user root for intelligence accounting")?
1340            .clone();
1341        let request = kcode_agent_runtime::SessionRunRequest {
1342            user_id,
1343            operation_id,
1344            rounds_used: self.rounds_used,
1345            round_limit: kcode_agent_runtime::DEFAULT_ROUND_LIMIT,
1346        };
1347        let mut host = KennedySessionHost {
1348            session: self,
1349            checkpoint: &mut checkpoint,
1350            accounting: None,
1351            pending_freeform_write: None,
1352            deadline_after_response: false,
1353        };
1354        let result = runtime.run_session(request, &mut host).await;
1355        let round_limit = result
1356            .as_ref()
1357            .is_err_and(kcode_agent_runtime::is_session_round_limit);
1358        let result: anyhow::Result<Option<String>> = async {
1359            if round_limit && matches!(host.session.mode, AgentMode::Ingress { .. }) {
1360                host.session.request_ingress_force_commit(
1361                    "agent_loop_round_limit",
1362                    host.session.projection().estimated_tokens,
1363                )?;
1364                (host.checkpoint)(host.session.snapshot()?).await?;
1365                Ok(None)
1366            } else {
1367                result
1368            }
1369        }
1370        .await;
1371        drop(host);
1372        self.active_provider_deadline = None;
1373        self.active_turn_deadline = None;
1374        let result = result?;
1375        match self.mode {
1376            AgentMode::Conversation => {
1377                if self.journal.state().source_terminated {
1378                    self.pending_turn = false;
1379                    self.pending_external_event_id = None;
1380                    checkpoint(self.snapshot()?).await?;
1381                    return Ok(None);
1382                }
1383                let Some(answer) = result else {
1384                    if self
1385                        .pending_external_event_id
1386                        .as_deref()
1387                        .and_then(|id| self.answer_for_external_event(id))
1388                        .is_some()
1389                    {
1390                        self.pending_turn = false;
1391                        self.pending_external_event_id = None;
1392                        checkpoint(self.snapshot()?).await?;
1393                        return Ok(None);
1394                    }
1395                    anyhow::bail!(
1396                        "Kennedy ended a conversational turn without an assistant response"
1397                    );
1398                };
1399                self.pending_turn = false;
1400                self.pending_external_event_id = None;
1401                checkpoint(self.snapshot()?).await?;
1402                Ok(Some(answer))
1403            }
1404            AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. } => {
1405                self.pending_turn = false;
1406                self.pending_external_event_id = None;
1407                self.finalize_kweb_session()?;
1408                self.completed = true;
1409                checkpoint(self.snapshot()?).await?;
1410                Ok(None)
1411            }
1412        }
1413    }
1414
1415    fn project_descendant<T>(
1416        &mut self,
1417        outcome: Result<kcode_intelligence_router::Accounted<T>, services::ApiError>,
1418    ) -> anyhow::Result<T> {
1419        match outcome {
1420            Ok(accounted) => {
1421                kcode_intelligence_chatend::record_descendant_receipt(
1422                    &mut self.journal,
1423                    &accounted.receipt,
1424                )?;
1425                Ok(accounted.value)
1426            }
1427            Err(error) => {
1428                if let Some(receipt) = &error.receipt {
1429                    kcode_intelligence_chatend::record_descendant_receipt(
1430                        &mut self.journal,
1431                        receipt,
1432                    )?;
1433                }
1434                Err(error.into())
1435            }
1436        }
1437    }
1438
1439    async fn run_subagent(
1440        &mut self,
1441        model: String,
1442        reasoning_effort: Option<String>,
1443        context_node_ids: Vec<String>,
1444        task: String,
1445        parent_operation_id: Uuid,
1446    ) -> anyhow::Result<String> {
1447        let reasoning_effort =
1448            reasoning_effort.unwrap_or_else(|| self.runtime.reasoning_effort.clone());
1449        let mut context = Vec::with_capacity(context_node_ids.len());
1450        for node_id in &context_node_ids {
1451            context.push(self.api.kmap_node(node_id)?.data.long_description);
1452        }
1453        let user_id = self
1454            .root_node_ids
1455            .first()
1456            .context("session has no user root for subagent intelligence accounting")?
1457            .clone();
1458        let timeout = self.agent_request_timeout();
1459        let runtime = self.api.agent_runtime();
1460        let first_event = self.journal.state().events.len();
1461        let cost_before = self.projection().status;
1462        let subagent_context = SubagentContext::new(
1463            self.root_node_ids.clone(),
1464            self.api.loads_fixed_connections(),
1465        )?;
1466        let result = {
1467            let mut host = KennedySubagentHost {
1468                session: self,
1469                context: subagent_context,
1470                captures: HashMap::new(),
1471            };
1472            runtime
1473                .run(
1474                    kcode_agent_runtime::RunRequest {
1475                        user_id,
1476                        parent_operation_id,
1477                        model,
1478                        reasoning_effort,
1479                        context,
1480                        task,
1481                        timeout,
1482                        start_metadata: json!({"contextNodeIds":context_node_ids}),
1483                    },
1484                    &mut host,
1485                )
1486                .await
1487        };
1488        match result {
1489            Ok(result) => {
1490                let cost_after = self.projection().status;
1491                Ok(format!(
1492                    "{}\n\n[{}]",
1493                    result.answer,
1494                    cost_summary(
1495                        "subagent cost",
1496                        cost_after
1497                            .estimated_cost_usd_nanos
1498                            .saturating_sub(cost_before.estimated_cost_usd_nanos),
1499                        cost_after
1500                            .unpriced_provider_calls
1501                            .saturating_sub(cost_before.unpriced_provider_calls),
1502                    )
1503                ))
1504            }
1505            Err(error) => {
1506                let may_have_effects =
1507                    self.journal.state().events[first_event..]
1508                        .iter()
1509                        .any(|event| {
1510                            matches!(
1511                                &event.kind,
1512                                EventKind::Note { label, .. } if label == "subagent_tool_call"
1513                            )
1514                        });
1515                if may_have_effects {
1516                    Err(error.context(
1517                        "the subagent failed after making Ktool calls; some tool effects may already have occurred",
1518                    ))
1519                } else {
1520                    Err(error)
1521                }
1522            }
1523        }
1524    }
1525
1526    async fn complete_subagent_freeform_write(
1527        &mut self,
1528        context: &mut SubagentContext,
1529        request: FreeformWrite,
1530        contents: String,
1531        budget: &kcode_agent_runtime::ContextBudget,
1532    ) -> anyhow::Result<kcode_agent_runtime::ToolOutcome> {
1533        let kind = request.kind();
1534        let freeform_tool = request.write_tool();
1535        anyhow::ensure!(
1536            context.source_is_open(kind, request.name()),
1537            "{} {:?} is not open in this subagent context. Call {} first.",
1538            kind.label(),
1539            request.name(),
1540            kind.open_tool()
1541        );
1542        let backend_arguments = request.capture_subagent(&mut self.journal, &now(), contents)?;
1543        let preview = self
1544            .api
1545            .managed_source_execute(
1546                &self.rust_lib_session_id,
1547                request.preview_tool(),
1548                backend_arguments.clone(),
1549                Vec::new(),
1550            )
1551            .await?;
1552        let preview = preview
1553            .snapshot
1554            .context("subagent freeform write preview omitted its source snapshot")?;
1555        let preview_state = context.source_state(&preview);
1556        anyhow::ensure!(
1557            budget.fits_state(preview_state.key, preview_state.text),
1558            "{freeform_tool} was not run because its resulting source state would exceed the subagent context limit"
1559        );
1560        let execution = self
1561            .api
1562            .managed_source_execute(
1563                &self.rust_lib_session_id,
1564                freeform_tool,
1565                backend_arguments,
1566                Vec::new(),
1567            )
1568            .await?;
1569        let snapshot = execution
1570            .snapshot
1571            .context("subagent freeform write omitted its resulting source snapshot")?;
1572        let state = context.apply_source_snapshot(snapshot);
1573        Ok(kcode_agent_runtime::ToolOutcome {
1574            text: execution.text,
1575            ok: true,
1576            state_updates: state.update.into_iter().collect(),
1577            displayed_state_keys: Vec::new(),
1578            capture: None,
1579        })
1580    }
1581
1582    async fn complete_freeform_write(
1583        &mut self,
1584        pending: PendingFreeformWrite,
1585        contents: String,
1586    ) -> anyhow::Result<ToolOutcome> {
1587        let request = pending.request;
1588        let freeform_tool = request.write_tool();
1589        let backend_arguments =
1590            request.capture(&mut self.journal, &now(), pending.call_box_id, contents)?;
1591        let preview_result = self
1592            .api
1593            .managed_source_execute(
1594                &self.rust_lib_session_id,
1595                request.preview_tool(),
1596                backend_arguments.clone(),
1597                Vec::new(),
1598            )
1599            .await;
1600        let preview = match preview_result {
1601            Ok(preview) => preview,
1602            Err(error) => {
1603                return Ok(ToolOutcome {
1604                    text: format!("{freeform_tool} failed: {error}"),
1605                    store_result: true,
1606                    ok: false,
1607                    end_session: false,
1608                    freeform_write: None,
1609                    managed_source_snapshot: None,
1610                });
1611            }
1612        };
1613        let _preview = preview
1614            .snapshot
1615            .context("freeform write preview omitted the resulting source snapshot")?;
1616        request.source_box_id(&self.journal)?;
1617
1618        let execution_result = self
1619            .api
1620            .managed_source_execute(
1621                &self.rust_lib_session_id,
1622                freeform_tool,
1623                backend_arguments,
1624                Vec::new(),
1625            )
1626            .await;
1627        let execution = match execution_result {
1628            Ok(execution) => execution,
1629            Err(error) => {
1630                return Ok(ToolOutcome {
1631                    text: format!("{freeform_tool} failed: {error}"),
1632                    store_result: true,
1633                    ok: false,
1634                    end_session: false,
1635                    freeform_write: None,
1636                    managed_source_snapshot: None,
1637                });
1638            }
1639        };
1640        let snapshot = execution
1641            .snapshot
1642            .context("freeform write omitted the resulting source snapshot")?;
1643        apply_snapshot(&mut self.journal, &now(), snapshot)?;
1644        Ok(ToolOutcome {
1645            text: execution.text,
1646            store_result: false,
1647            ok: true,
1648            end_session: false,
1649            freeform_write: None,
1650            managed_source_snapshot: None,
1651        })
1652    }
1653
1654    async fn send_telegram_dm(&mut self, arguments: &Value) -> anyhow::Result<String> {
1655        let request = kcode_telegram_session_coordinator::parse_private_request(arguments)?;
1656        let attachments = self.telegram_delivery_attachments(request.attachments)?;
1657        let caller_holds_user_lock = self.session_type == "telegram"
1658            && self.channel.get("telegramUserId").and_then(Value::as_i64)
1659                == Some(request.telegram_user_id);
1660        self.api
1661            .telegram()
1662            .send_private(kcode_telegram_session_coordinator::PrivateDelivery {
1663                telegram_user_id: request.telegram_user_id,
1664                message: request.message,
1665                attachments,
1666                caller_holds_user_lock,
1667            })
1668            .await
1669    }
1670
1671    async fn send_telegram_group_message(&mut self, arguments: &Value) -> anyhow::Result<String> {
1672        let request = kcode_telegram_session_coordinator::parse_group_request(arguments)?;
1673        let attachments = self.telegram_delivery_attachments(request.attachments)?;
1674        self.api
1675            .telegram()
1676            .send_group(kcode_telegram_session_coordinator::GroupDelivery {
1677                root_node_id: request.root_node_id,
1678                message: request.message,
1679                attachments,
1680            })
1681            .await
1682    }
1683
1684    fn telegram_delivery_attachments(
1685        &mut self,
1686        requests: Vec<kcode_telegram_session_coordinator::AttachmentRequest>,
1687    ) -> anyhow::Result<Vec<kcode_telegram_session_coordinator::Attachment>> {
1688        let api = self.api.clone();
1689        kcode_kennedy_session_objects::delivery_attachments(
1690            &mut self.journal,
1691            requests,
1692            move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
1693        )
1694    }
1695
1696    async fn execute_tool(
1697        &mut self,
1698        call: &ToolCall,
1699        operation_id: Uuid,
1700    ) -> anyhow::Result<ToolOutcome> {
1701        self.assert_tool_allowed(&call.name)?;
1702        let decoded = decode(&call.name, &call.arguments)?;
1703        let mut end_session = false;
1704        let mut store_result = true;
1705        let mut freeform_write = None;
1706        let mut managed_source_snapshot = None;
1707        let text = match (call.name.as_str(), decoded) {
1708            ("SendTelegramDM", _) => self.send_telegram_dm(&call.arguments).await?,
1709            ("SendTelegramGroupMessage", _) => {
1710                self.send_telegram_group_message(&call.arguments).await?
1711            }
1712            (
1713                "RunSubagent",
1714                Some(DecodedTool::RunSubagent {
1715                    model,
1716                    reasoning_effort,
1717                    context_node_ids,
1718                    task,
1719                }),
1720            ) => {
1721                let first_event = self.journal.state().events.len();
1722                match self
1723                    .run_subagent(
1724                        model,
1725                        reasoning_effort,
1726                        context_node_ids,
1727                        task,
1728                        operation_id,
1729                    )
1730                    .await
1731                {
1732                    Ok(response) => response,
1733                    Err(error) => {
1734                        let may_have_effects = self.journal.state().events[first_event..]
1735                            .iter()
1736                            .any(|event| {
1737                                matches!(
1738                                    &event.kind,
1739                                    EventKind::Note { label, .. }
1740                                        if label == "subagent_tool_call"
1741                                )
1742                            });
1743                        if may_have_effects {
1744                            return Err(error.context(
1745                                "the subagent failed after making Ktool calls; some tool effects may already have occurred",
1746                            ));
1747                        }
1748                        return Err(error);
1749                    }
1750                }
1751            }
1752            ("EndSession", Some(DecodedTool::EndSession { message })) => {
1753                anyhow::ensure!(
1754                    !matches!(self.mode, AgentMode::Conversation),
1755                    "EndSession is only available during an autonomous or history-ingress session"
1756                );
1757                end_session = true;
1758                if matches!(self.mode, AgentMode::FreeTime)
1759                    && let Some(message) = message.filter(|message| !message.trim().is_empty())
1760                {
1761                    self.free_time["nextSessionMessage"] = json!(message);
1762                }
1763                "Session ending.".into()
1764            }
1765            ("DehydrateBoxes", Some(DecodedTool::BoxIds(ids))) => {
1766                self.journal.dehydrate_boxes(now(), &ids)?;
1767                format!(
1768                    "Dehydrated boxes {}.",
1769                    ids.iter()
1770                        .map(ToString::to_string)
1771                        .collect::<Vec<_>>()
1772                        .join(", ")
1773                )
1774            }
1775            ("SummarizeBox", Some(DecodedTool::SummarizeBox { box_id, summary })) => {
1776                self.journal.summarize_box(now(), box_id, summary)?;
1777                format!("Summarized box {box_id}.")
1778            }
1779            ("HydrateBox", Some(DecodedTool::BoxId(id))) => {
1780                self.journal.rehydrate_box(now(), id)?;
1781                let external_event_id = self.pending_external_event_id.clone();
1782                match self.recover_context_overflow(external_event_id.as_deref(), &[id])? {
1783                    ContextRecovery::NotNeeded => format!("Hydrated box {id}."),
1784                    ContextRecovery::Recovered => {
1785                        format!("Hydrated box {id}.\n\n{CONTEXT_OVERFLOW_WARNING}")
1786                    }
1787                    ContextRecovery::Irreducible => anyhow::bail!(CONTEXT_OVERFLOW_WARNING),
1788                }
1789            }
1790            ("BoxesIntoObjects", Some(DecodedTool::BoxIds(ids))) => {
1791                kcode_kennedy_box_text_objects::stage_box_text_objects(
1792                    &mut self.journal,
1793                    &ids,
1794                    &now(),
1795                )?
1796            }
1797            ("LoadNodes", Some(DecodedTool::LoadNodes(identifiers))) => {
1798                load_durable_batch(self.api.kmap(), &mut self.context, &identifiers)?;
1799                let changed = self.sync_kweb_boxes()?;
1800                store_result = false;
1801                render_load_nodes_result(
1802                    &self.journal,
1803                    &changed,
1804                    &self.runtime_budget().footer_lines(),
1805                )?
1806            }
1807            (
1808                "EmitObject",
1809                Some(DecodedTool::EmitObject {
1810                    object_id,
1811                    file_name,
1812                }),
1813            ) => {
1814                anyhow::ensure!(
1815                    matches!(self.mode, AgentMode::Conversation),
1816                    "EmitObject is only available in a conversation"
1817                );
1818                let object = self.resolve_object(&object_id)?;
1819                let file_name = file_name.unwrap_or_else(|| object.file_name.clone());
1820                if let Some(maximum) = self.channel.get("maxObjectBytes").and_then(Value::as_u64) {
1821                    anyhow::ensure!(
1822                        !object.bytes.is_empty(),
1823                        "object {object_id} is empty and cannot be sent through this channel"
1824                    );
1825                    anyhow::ensure!(
1826                        object.bytes.len() as u64 <= maximum,
1827                        "object {object_id} is {} bytes, over this channel's {maximum}-byte limit",
1828                        object.bytes.len()
1829                    );
1830                }
1831                let descriptor = json!({
1832                    "objectId":object_id,
1833                    "fileName":file_name,
1834                    "mediaType":object.media_type,
1835                    "byteLength":object.bytes.len(),
1836                });
1837                let mut metadata = json!({
1838                    "outputKind":"object",
1839                    "attachments":[descriptor.clone()],
1840                });
1841                if let Some(external_event_id) = &self.pending_external_event_id {
1842                    metadata["externalEventId"] = json!(external_event_id);
1843                }
1844                let content = BoxContent {
1845                    text: String::new(),
1846                    objects: vec![object_id.clone()],
1847                    metadata,
1848                };
1849                self.journal
1850                    .create_box(now(), "Kennedy message", BoxOwner::Kennedy, content)?;
1851                let mut transcript = json!({
1852                    "role":"kennedy",
1853                    "content":"",
1854                    "objects":[object_id],
1855                    "attachments":[descriptor],
1856                });
1857                if let Some(external_event_id) = &self.pending_external_event_id {
1858                    transcript["externalEventId"] = json!(external_event_id);
1859                }
1860                self.transcript.push(transcript);
1861                store_result = false;
1862                "Object emitted to the user.".into()
1863            }
1864            ("WebSearch", Some(DecodedTool::WebSearch { question, model })) => {
1865                let user_id = self
1866                    .root_node_ids
1867                    .first()
1868                    .context("session has no user root for intelligence accounting")?
1869                    .clone();
1870                let outcome = self
1871                    .api
1872                    .search(
1873                        &user_id,
1874                        kcode_intelligence_router::SearchRequest {
1875                            question,
1876                            model,
1877                            operation_id: Uuid::new_v4(),
1878                            parent_operation_id: Some(operation_id),
1879                        },
1880                    )
1881                    .await;
1882                let result = self.project_descendant(outcome)?;
1883                render_web_search_result(&result)?
1884            }
1885            ("WebFetch", Some(DecodedTool::WebFetch(url))) => {
1886                let user_id = self
1887                    .root_node_ids
1888                    .first()
1889                    .context("session has no user root for intelligence accounting")?;
1890                let result = self
1891                    .api
1892                    .fetch(
1893                        user_id,
1894                        kcode_intelligence_router::FetchRequest {
1895                            url,
1896                            operation_id: Uuid::new_v4(),
1897                            parent_operation_id: Some(operation_id),
1898                        },
1899                    )
1900                    .await?;
1901                render_web_fetch_result(&result)?
1902            }
1903            ("StageTelegramGroupMedia", Some(DecodedTool::StageTelegramGroupMedia(message_id))) => {
1904                let media_ref = kcode_telegram_session_coordinator::group_media_reference(
1905                    &self.group_context,
1906                    message_id,
1907                )?;
1908                let chat_id = media_ref.chat_id;
1909                let api = self.api.clone();
1910                let staged = kcode_kennedy_session_objects::stage_telegram_group_media(
1911                    &mut self.journal,
1912                    kcode_kennedy_session_objects::TelegramStageRequest {
1913                        chat_id,
1914                        message_id,
1915                        maximum_bytes: MAX_MEDIA_ENRICHMENT_BYTES,
1916                        transport_metadata: media_ref.transport_metadata(),
1917                        recorded_at: now(),
1918                    },
1919                    || api.telegram().group_message_media(chat_id, message_id),
1920                    |media_type| {
1921                        kcode_telegram_session_coordinator::group_media_file_name(
1922                            &media_ref, media_type,
1923                        )
1924                    },
1925                )?;
1926                render(RenderRequest::StagedTelegramMedia {
1927                    pending_id: &staged.descriptor.pending_id,
1928                    kind: &staged.kind,
1929                    file_name: &staged.descriptor.file_name,
1930                    media_type: &staged.descriptor.media_type,
1931                    size_bytes: staged.descriptor.size_bytes,
1932                    message_id,
1933                    reused: staged.reused,
1934                })?
1935            }
1936            (
1937                "TranscribeAudio",
1938                Some(DecodedTool::MediaEnrichment {
1939                    object_id,
1940                    model,
1941                    prompt,
1942                }),
1943            ) => {
1944                let object = self.resolve_media_object(&object_id)?;
1945                validate(ValidationRequest::TranscribableAudio(&object.media_type))?;
1946                validate(ValidationRequest::TranscriptionModel(&model))?;
1947                let user_id = self
1948                    .root_node_ids
1949                    .first()
1950                    .context("session has no user root for intelligence accounting")?
1951                    .clone();
1952                let outcome = self
1953                    .api
1954                    .transcribe_audio(
1955                        &user_id,
1956                        &model,
1957                        &prompt,
1958                        object.bytes,
1959                        object.file_name.clone(),
1960                        &object.media_type,
1961                        None,
1962                        operation_id,
1963                    )
1964                    .await;
1965                let result = self.project_descendant(outcome)?;
1966                render_audio_transcription_result(
1967                    &object.object_id,
1968                    &object.file_name,
1969                    &object.media_type,
1970                    &result,
1971                )?
1972            }
1973            (
1974                "AnnotateMedia",
1975                Some(DecodedTool::MediaEnrichment {
1976                    object_id,
1977                    model,
1978                    prompt,
1979                }),
1980            ) => {
1981                let media = self.resolve_media_object(&object_id)?;
1982                validate(ValidationRequest::Annotation {
1983                    model: &model,
1984                    media_type: &media.media_type,
1985                })?;
1986                let user_id = self
1987                    .root_node_ids
1988                    .first()
1989                    .context("session has no user root for intelligence accounting")?
1990                    .clone();
1991                let outcome = self
1992                    .api
1993                    .annotate_media(
1994                        &user_id,
1995                        &model,
1996                        &prompt,
1997                        media.bytes,
1998                        media.file_name.clone(),
1999                        &media.media_type,
2000                        operation_id,
2001                    )
2002                    .await;
2003                let result = self.project_descendant(outcome)?;
2004                render_media_annotation_result(
2005                    &media.object_id,
2006                    &media.file_name,
2007                    &media.media_type,
2008                    &result,
2009                )?
2010            }
2011            (
2012                "GenerateImage",
2013                Some(DecodedTool::GenerateImage {
2014                    model,
2015                    prompt,
2016                    reference_object_ids,
2017                }),
2018            ) => {
2019                let mut references = Vec::with_capacity(reference_object_ids.len());
2020                for object_id in &reference_object_ids {
2021                    references.push(self.resolve_image_object(object_id)?);
2022                }
2023                let user_id = self
2024                    .root_node_ids
2025                    .first()
2026                    .context("session has no user root for intelligence accounting")?
2027                    .clone();
2028                let outcome = self
2029                    .api
2030                    .generate_image(&user_id, &model, &prompt, references, operation_id)
2031                    .await;
2032                let result = self.project_descendant(outcome)?;
2033                let size = result.bytes.len();
2034                let file_name =
2035                    format!("generated-image.{}", image_extension(&result.content_type));
2036                let object_id = self.api.save_generated_image(
2037                    result.bytes,
2038                    &file_name,
2039                    &result.content_type,
2040                    &result.model,
2041                )?;
2042                format!(
2043                    "Generated image.\nObject: {object_id}\nFile: {file_name}\nContent type: {}\nSize: {size} bytes\nModel: {}\nUse EmitObject with {object_id} to deliver it.",
2044                    result.content_type, result.model
2045                )
2046            }
2047            ("ExtractDocumentText", Some(DecodedTool::ObjectId(object_id))) => {
2048                let object = self.resolve_media_object(&object_id)?;
2049                validate(ValidationRequest::ExtractableDocument {
2050                    media_type: &object.media_type,
2051                    file_name: &object.file_name,
2052                })?;
2053                let result = self
2054                    .api
2055                    .extract_document(object.bytes, object.file_name.clone(), &object.media_type)
2056                    .await?;
2057                render_document_extraction_result(&object.object_id, &object.file_name, &result)?
2058            }
2059            (name, None) if SPEECH_CLASSIFICATION_TOOLS.contains(&name) => {
2060                self.api
2061                    .execute_speech_classification_tool(name, call.arguments.clone())
2062                    .await?
2063            }
2064            (name, None) if TASK_BOARD_TOOLS.contains(&name) => {
2065                self.execute_task_board_tool(name, &call.arguments).await?
2066            }
2067            (name, Some(decoded)) if is_kweb_mutation(name) => {
2068                let text = execute_kweb_mutation(
2069                    name,
2070                    decoded,
2071                    &self.context,
2072                    &mut self.plan,
2073                    &mut self.journal,
2074                )?;
2075                self.sync_kweb_boxes()?;
2076                text
2077            }
2078            (name, None)
2079                if RUST_LIB_TOOLS.contains(&name)
2080                    || WEB_LIB_TOOLS.contains(&name)
2081                    || RUST_BIN_TOOLS.contains(&name) =>
2082            {
2083                if let Some(request) = prepare_freeform_write(&self.journal, name, &call.arguments)?
2084                {
2085                    store_result = false;
2086                    let acknowledgement = request.acknowledgement();
2087                    freeform_write = Some(request);
2088                    acknowledgement
2089                } else {
2090                    let object_ids = if name == CALL_RUST_BIN_TOOL {
2091                        decode_managed_objects(ManagedObjectArguments::RustBinary(&call.arguments))?
2092                    } else if name == ATTACH_OBJECT_WEB_LIB_TOOL {
2093                        decode_managed_objects(ManagedObjectArguments::WebLibraryAttachment(
2094                            &call.arguments,
2095                        ))?
2096                    } else {
2097                        Vec::new()
2098                    };
2099                    let mut objects = Vec::with_capacity(object_ids.len());
2100                    for object_id in object_ids {
2101                        objects.push(self.resolve_object(&object_id)?.bytes);
2102                    }
2103                    let execution = self
2104                        .api
2105                        .managed_source_execute(
2106                            &self.rust_lib_session_id,
2107                            name,
2108                            call.arguments.clone(),
2109                            objects,
2110                        )
2111                        .await?;
2112                    if let Some(snapshot) = execution.snapshot {
2113                        managed_source_snapshot = Some(snapshot);
2114                        store_result = false;
2115                    }
2116                    execution.text
2117                }
2118            }
2119            (name, Some(_)) => {
2120                anyhow::bail!("decoded contract for {name} did not match its dispatch lane")
2121            }
2122            (name, None) => anyhow::bail!("Tool {name} is not available"),
2123        };
2124        Ok(ToolOutcome {
2125            text,
2126            store_result,
2127            ok: true,
2128            end_session,
2129            freeform_write,
2130            managed_source_snapshot,
2131        })
2132    }
2133
2134    async fn execute_task_board_tool(
2135        &self,
2136        name: &str,
2137        arguments: &Value,
2138    ) -> anyhow::Result<String> {
2139        let board = self
2140            .api
2141            .task_board()
2142            .context("task board is not configured")?
2143            .clone();
2144        let name = name.to_owned();
2145        let arguments = arguments.clone();
2146        let user_id = self
2147            .root_node_ids
2148            .first()
2149            .context("session has no user root for task-category lookup")?
2150            .clone();
2151        tokio::task::spawn_blocking(move || -> anyhow::Result<String> {
2152            let output = match name.as_str() {
2153                "CreateTaskCategory" => serde_json::to_string_pretty(
2154                    &board.create_category(serde_json::from_value(arguments)?)?,
2155                )?,
2156                "GetTaskCategory" => {
2157                    let call: CategoryCall = serde_json::from_value(arguments)?;
2158                    serde_json::to_string_pretty(&board.category(
2159                        &call.category_id,
2160                        kcode_task_board::BrowsePage {
2161                            user_id,
2162                            offset: call.offset,
2163                            limit: call.limit,
2164                        },
2165                    )?)?
2166                }
2167                "RemoveTaskCategory" => {
2168                    let call: CategoryId = serde_json::from_value(arguments)?;
2169                    board.remove_category(&call.category_id)?;
2170                    format!("Removed category {}.", call.category_id)
2171                }
2172                "CreateTask" => serde_json::to_string_pretty(
2173                    &board.create_task(serde_json::from_value(arguments)?)?,
2174                )?,
2175                "GetTask" => {
2176                    let call: TaskId = serde_json::from_value(arguments)?;
2177                    serde_json::to_string_pretty(&board.task(&call.task_id)?)?
2178                }
2179                "UpdateTask" => serde_json::to_string_pretty(
2180                    &board.update_task(serde_json::from_value(arguments)?)?,
2181                )?,
2182                "RemoveTask" => {
2183                    let call: TaskId = serde_json::from_value(arguments)?;
2184                    board.remove_task(&call.task_id)?;
2185                    format!("Removed task {}.", call.task_id)
2186                }
2187                "GetTopTaskOrphan" => {
2188                    let _: EmptyCall = serde_json::from_value(arguments)?;
2189                    serde_json::to_string_pretty(&board.top_orphan()?)?
2190                }
2191                _ => anyhow::bail!("Tool {name} is not a task-board operation"),
2192            };
2193            Ok(output)
2194        })
2195        .await
2196        .context("task-board worker stopped")?
2197    }
2198
2199    fn assert_tool_allowed(&self, name: &str) -> anyhow::Result<()> {
2200        let write = matches!(
2201            name,
2202            "ConnectNodes"
2203                | "ConsolidateFanout"
2204                | "SetFixedConnection"
2205                | "CreateNode"
2206                | "UpdateNode"
2207        );
2208        anyhow::ensure!(
2209            !write || !matches!(self.mode, AgentMode::Conversation),
2210            "{name} requires the global Kweb write lane and is unavailable in a read-only conversation"
2211        );
2212        if name == "EndSession" {
2213            anyhow::ensure!(
2214                !matches!(self.mode, AgentMode::Conversation),
2215                "EndSession is unavailable in a conversation"
2216            );
2217        }
2218        Ok(())
2219    }
2220
2221    fn sync_kweb_boxes(&mut self) -> anyhow::Result<Vec<BoxId>> {
2222        let (updates, creates) = kweb_plan_projection(&self.plan);
2223        self.context
2224            .sync_chatend(&mut self.journal, now(), &updates, &creates)
2225            .map_err(anyhow::Error::new)
2226    }
2227
2228    fn record_tool_invocation(
2229        &mut self,
2230        name: &str,
2231        arguments: Value,
2232    ) -> anyhow::Result<RecordedToolInvocation> {
2233        let invocation = RecordedToolInvocation {
2234            invocation_id: Uuid::new_v4().to_string(),
2235            tool_instance: tool_instance(name),
2236            tool_name: name.into(),
2237        };
2238        self.journal.record(
2239            now(),
2240            EventKind::ToolInvoked {
2241                tool_instance: invocation.tool_instance.clone(),
2242                tool_name: invocation.tool_name.clone(),
2243                arguments,
2244                invocation_id: Some(invocation.invocation_id.clone()),
2245            },
2246        )?;
2247        Ok(invocation)
2248    }
2249
2250    fn record_tool_completion(
2251        &mut self,
2252        invocation: Option<&RecordedToolInvocation>,
2253        outcome: Value,
2254    ) -> anyhow::Result<EventId> {
2255        let (tool_instance, tool_name, invocation_id) = invocation
2256            .map(|invocation| {
2257                (
2258                    invocation.tool_instance.clone(),
2259                    invocation.tool_name.clone(),
2260                    Some(invocation.invocation_id.clone()),
2261                )
2262            })
2263            .unwrap_or_else(|| ("call_ktool".into(), "call_ktool".into(), None));
2264        self.journal.record(
2265            now(),
2266            EventKind::ToolCompleted {
2267                tool_instance,
2268                tool_name,
2269                outcome,
2270                invocation_id,
2271            },
2272        )
2273    }
2274
2275    fn node_data(plan: &KwebPlan, context: &KwebContext, id: &str) -> anyhow::Result<PlannedNode> {
2276        if let Some(data) = plan.created(id) {
2277            return Ok(data.clone());
2278        }
2279        if let Some(data) = plan.updates.get(id) {
2280            return Ok(data.clone());
2281        }
2282        let node = context
2283            .node(id)
2284            .with_context(|| format!("Kweb context does not contain node {id}"))?;
2285        Ok(planned_node(node))
2286    }
2287
2288    fn put_node_data(plan: &mut KwebPlan, id: &str, data: PlannedNode) -> anyhow::Result<()> {
2289        if let Some(created) = plan.created_mut(id) {
2290            *created = data;
2291        } else {
2292            canonical_id(id)?;
2293            plan.updates.insert(id.to_owned(), data);
2294        }
2295        Ok(())
2296    }
2297
2298    fn connect_nodes(
2299        plan: &mut KwebPlan,
2300        context: &KwebContext,
2301        ids: Vec<String>,
2302    ) -> anyhow::Result<String> {
2303        for id in &ids {
2304            Self::ensure_known_node(plan, context, id)?;
2305        }
2306        for id in &ids {
2307            let mut data = Self::node_data(plan, context, id)?;
2308            let mut recent = ids
2309                .iter()
2310                .filter(|other| *other != id)
2311                .cloned()
2312                .collect::<Vec<_>>();
2313            for other in data.recent_connections {
2314                if &other != id && !recent.contains(&other) {
2315                    recent.push(other);
2316                }
2317            }
2318            data.recent_connections = recent;
2319            Self::put_node_data(plan, id, data)?;
2320        }
2321        Ok(format!(
2322            "Staged connections among nodes {}.",
2323            ids.join(", ")
2324        ))
2325    }
2326
2327    fn consolidate_fanout(
2328        plan: &mut KwebPlan,
2329        context: &KwebContext,
2330        parent: String,
2331        fanout: Vec<String>,
2332        aggregator: String,
2333    ) -> anyhow::Result<String> {
2334        for id in std::iter::once(&parent)
2335            .chain(std::iter::once(&aggregator))
2336            .chain(fanout.iter())
2337        {
2338            Self::ensure_known_node(plan, context, id)?;
2339        }
2340        let mut parent_data = Self::node_data(plan, context, &parent)?;
2341        parent_data
2342            .recent_connections
2343            .retain(|id| !fanout.contains(id));
2344        if !parent_data.recent_connections.contains(&aggregator) {
2345            parent_data.recent_connections.push(aggregator.clone());
2346        }
2347        let mut aggregator_data = Self::node_data(plan, context, &aggregator)?;
2348        for id in fanout {
2349            if !aggregator_data.recent_connections.contains(&id) {
2350                aggregator_data.recent_connections.push(id);
2351            }
2352        }
2353        Self::put_node_data(plan, &parent, parent_data)?;
2354        Self::put_node_data(plan, &aggregator, aggregator_data)?;
2355        Ok(format!(
2356            "Staged fanout consolidation from node {parent} into node {aggregator}."
2357        ))
2358    }
2359
2360    fn set_fixed_connection(
2361        plan: &mut KwebPlan,
2362        context: &KwebContext,
2363        parent: String,
2364        child: Option<String>,
2365        slot: usize,
2366    ) -> anyhow::Result<String> {
2367        Self::ensure_known_node(plan, context, &parent)?;
2368        if let Some(child) = &child {
2369            Self::ensure_known_node(plan, context, child)?;
2370            anyhow::ensure!(child != &parent, "a node cannot connect to itself");
2371        }
2372        let mut data = Self::node_data(plan, context, &parent)?;
2373        if let Some(child) = child.clone() {
2374            anyhow::ensure!(
2375                slot <= data.fixed_connections.len() + 1,
2376                "fixed connection positions must remain contiguous"
2377            );
2378            data.fixed_connections.retain(|id| id != &child);
2379            if slot - 1 < data.fixed_connections.len() {
2380                data.fixed_connections[slot - 1] = child;
2381            } else {
2382                data.fixed_connections.push(child);
2383            }
2384        } else if slot > 0 && slot - 1 < data.fixed_connections.len() {
2385            data.fixed_connections.remove(slot - 1);
2386        }
2387        Self::put_node_data(plan, &parent, data)?;
2388        Ok(match child {
2389            Some(child) => {
2390                format!("Staged node {child} in fixed slot {slot} of node {parent}.")
2391            }
2392            None => format!("Cleared fixed slot {slot} of node {parent} in the staged plan."),
2393        })
2394    }
2395
2396    fn create_node(
2397        plan: &mut KwebPlan,
2398        context: &KwebContext,
2399        journal: &mut HistorySession,
2400        arguments: CreateNodeArguments,
2401    ) -> anyhow::Result<String> {
2402        let CreateNodeArguments {
2403            parents,
2404            owner,
2405            short_name,
2406            short_description,
2407            long_description,
2408        } = arguments;
2409        for id in parents.iter().chain(std::iter::once(&owner)) {
2410            if id != "self" && id != "unowned" {
2411                Self::ensure_known_node(plan, context, id)?;
2412            }
2413        }
2414        let pending = journal.allocate_pending_node(now())?.to_string();
2415        plan.creates.push(StagedNodeCreate {
2416            pending_id: pending.clone(),
2417            data: PlannedNode {
2418                short_name,
2419                short_description,
2420                long_description,
2421                owner,
2422                fixed_connections: Vec::new(),
2423                recent_connections: parents.clone(),
2424                objects: Vec::new(),
2425                attach_session_archive: true,
2426            },
2427        });
2428        for parent in parents {
2429            let mut data = Self::node_data(plan, context, &parent)?;
2430            data.recent_connections.retain(|id| id != &pending);
2431            data.recent_connections.insert(0, pending.clone());
2432            Self::put_node_data(plan, &parent, data)?;
2433        }
2434        Ok(format!("Created staged node {pending}."))
2435    }
2436
2437    fn update_node(
2438        plan: &mut KwebPlan,
2439        context: &KwebContext,
2440        id: String,
2441        owner: String,
2442        short_name: String,
2443        short_description: String,
2444        long_description: String,
2445    ) -> anyhow::Result<String> {
2446        Self::ensure_known_node(plan, context, &id)?;
2447        if owner != "self" && owner != "unowned" {
2448            Self::ensure_known_node(plan, context, &owner)?;
2449        }
2450        let mut data = Self::node_data(plan, context, &id)?;
2451        data.owner = owner;
2452        data.short_name = short_name;
2453        data.short_description = short_description;
2454        data.long_description = long_description;
2455        data.attach_session_archive = true;
2456        Self::put_node_data(plan, &id, data)?;
2457        Ok(format!("Staged the update to node {id}."))
2458    }
2459
2460    fn ensure_known_node(plan: &KwebPlan, context: &KwebContext, id: &str) -> anyhow::Result<()> {
2461        if id.starts_with("pending:") {
2462            anyhow::ensure!(
2463                plan.created(id).is_some(),
2464                "pending node {id} is not part of this session"
2465            );
2466        } else {
2467            canonical_id(id)?;
2468            anyhow::ensure!(
2469                context.contains_full_node(id) || plan.updates.contains_key(id),
2470                "node {id} is not loaded; call LoadNodes first"
2471            );
2472        }
2473        Ok(())
2474    }
2475
2476    fn finalize_kweb_session(&mut self) -> anyhow::Result<()> {
2477        if self.commit_receipt.is_some() {
2478            return Ok(());
2479        }
2480        self.journal.repair_unfinished_tools(now())?;
2481        self.journal.seal()?;
2482        let archive = self.journal.archive_bytes()?;
2483        let object_locations = self
2484            .journal
2485            .objects()
2486            .iter()
2487            .map(|(id, location)| (id.clone(), location.clone()))
2488            .collect::<Vec<_>>();
2489        let mut objects = BTreeMap::new();
2490        for (id, location) in object_locations {
2491            let pending_id = id.to_string();
2492            let transport_kind =
2493                kcode_kennedy_session_objects::staged_descriptor(&self.journal, &id)?
2494                    .transport_kind;
2495            let bytes = encode_file(
2496                &pending_id,
2497                location.metadata.file_name.as_deref(),
2498                &location.metadata.media_type,
2499                transport_kind.as_deref(),
2500                self.journal.read_object(&id)?,
2501            )
2502            .with_context(|| format!("encoding staged object {pending_id}"))?;
2503            anyhow::ensure!(
2504                objects.insert(pending_id.clone(), bytes).is_none(),
2505                "duplicate staged object {pending_id}"
2506            );
2507        }
2508        let mut creates = BTreeMap::new();
2509        for create in &self.plan.creates {
2510            anyhow::ensure!(
2511                creates
2512                    .insert(create.pending_id.clone(), create.data.clone())
2513                    .is_none(),
2514                "duplicate staged node {}",
2515                create.pending_id
2516            );
2517        }
2518        let updates = self
2519            .plan
2520            .updates
2521            .iter()
2522            .map(|(node_id, data)| {
2523                node_id
2524                    .parse::<NodeId>()
2525                    .with_context(|| format!("{node_id:?} is not a canonical node ID"))
2526                    .map(|node_id| (node_id, data.clone()))
2527            })
2528            .collect::<anyhow::Result<BTreeMap<_, _>>>()?;
2529        let result = self.api.commit_kweb_session(CommitRequest {
2530            idempotency_key: self.journal.state().metadata.session_id.clone(),
2531            author: self.commit_author.clone(),
2532            source_created_at: DateTime::parse_from_rfc3339(&self.started_at)
2533                .context("session start timestamp is invalid")?
2534                .with_timezone(&Utc),
2535            archive,
2536            objects,
2537            creates,
2538            updates,
2539        })?;
2540        self.journal
2541            .mark_completed(result.session_object_id.to_string());
2542        self.commit_receipt = Some(result);
2543        Ok(())
2544    }
2545
2546    fn prepare_free_time_round(&mut self) -> anyhow::Result<bool> {
2547        if !matches!(self.mode, AgentMode::FreeTime) {
2548            return Ok(false);
2549        }
2550        let Some(deadline) = deadline(&self.free_time) else {
2551            return Ok(false);
2552        };
2553        if Utc::now() >= deadline {
2554            self.free_time_end_reason = Some("deadline".into());
2555            self.journal.create_box(
2556                now(),
2557                "Self-time timer",
2558                BoxOwner::Controller,
2559                BoxContent::text(
2560                    "The self-time deadline has arrived. Finish without starting more tool work.",
2561                ),
2562            )?;
2563            return Ok(true);
2564        }
2565        Ok(false)
2566    }
2567
2568    fn refresh_runtime_prompt(&mut self) -> anyhow::Result<()> {
2569        let Some(system_box) = self
2570            .journal
2571            .state()
2572            .boxes
2573            .values()
2574            .find(|state| matches!(state.owner, BoxOwner::System))
2575            .map(|state| state.id)
2576        else {
2577            return Ok(());
2578        };
2579        let current = self
2580            .journal
2581            .state()
2582            .boxes
2583            .get(&system_box)
2584            .context("system-prompt box disappeared")?
2585            .canonical
2586            .content
2587            .text
2588            .clone();
2589        let marker = "\n\nCurrent runtime\n\n";
2590        let Some((prefix, _)) = current.rsplit_once(marker) else {
2591            return Ok(());
2592        };
2593        let refreshed = format!(
2594            "{prefix}{marker}{}",
2595            render(RenderRequest::RuntimeDescription {
2596                model: &self.runtime.model,
2597                reasoning_effort: &self.runtime.reasoning_effort,
2598                current_time: Utc::now(),
2599            })?
2600        );
2601        if refreshed != current {
2602            self.journal
2603                .update_box(now(), system_box, BoxContent::text(refreshed))?;
2604        }
2605        Ok(())
2606    }
2607
2608    fn agent_request_timeout(&self) -> Option<Duration> {
2609        if matches!(self.mode, AgentMode::Conversation) && self.session_type == "conversation" {
2610            return Some(BROWSER_CONVERSATION_REQUEST_TIMEOUT);
2611        }
2612        if matches!(self.mode, AgentMode::Ingress { .. }) {
2613            return Some(HISTORY_INGRESS_REQUEST_TIMEOUT);
2614        }
2615        if matches!(self.mode, AgentMode::Wakeup) {
2616            return Some(WAKEUP_REQUEST_TIMEOUT);
2617        }
2618        if matches!(self.mode, AgentMode::FreeTime) {
2619            let deadline = deadline(&self.free_time)?;
2620            return Some(Duration::from_secs(
2621                (deadline - Utc::now()).num_seconds().max(1) as u64
2622                    + SELF_TIME_HARD_STOP_ALLOWANCE.as_secs(),
2623            ));
2624        }
2625        None
2626    }
2627
2628    pub fn refresh_telegram_group_context(
2629        &mut self,
2630        group_context: &Value,
2631        current_message_id: Option<&str>,
2632    ) -> anyhow::Result<()> {
2633        if self.session_type != "telegram-group" {
2634            return Ok(());
2635        }
2636        self.channel["groupContext"] = group_context.clone();
2637        self.group_context = group_context.clone();
2638        self.journal.create_box(
2639            now(),
2640            "Telegram group update",
2641            BoxOwner::Controller,
2642            BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
2643                group_context,
2644            )),
2645        )?;
2646        self.recover_context_overflow(current_message_id, &[])?;
2647        Ok(())
2648    }
2649
2650    pub fn finalize_free_time(&mut self, reason: &str) -> anyhow::Result<()> {
2651        anyhow::ensure!(
2652            matches!(reason, "tool" | "deadline" | "hard-stop" | "user-stop"),
2653            "invalid self-time completion reason"
2654        );
2655        self.free_time["sliceEndedReason"] = json!(reason);
2656        self.free_time["sliceEndedAt"] = json!(now());
2657        self.pending_turn = false;
2658        self.pending_external_event_id = None;
2659        Ok(())
2660    }
2661
2662    pub fn commit_current_write_session(&mut self) -> anyhow::Result<()> {
2663        anyhow::ensure!(
2664            matches!(
2665                self.mode,
2666                AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. }
2667            ),
2668            "a read-only conversation cannot be committed as a Kweb write session"
2669        );
2670        self.finalize_kweb_session()?;
2671        self.completed = true;
2672        Ok(())
2673    }
2674
2675    pub fn snapshot(&self) -> anyhow::Result<Value> {
2676        let projection = self.projection();
2677        let chatend_text = projection.render();
2678        let session_status = projection.status.clone();
2679        Ok(json!({
2680            "format":"kennedy-chatend",
2681            "version":1,
2682            "stateVersion":3,
2683            "sessionId":self.journal.state().metadata.session_id,
2684            "chatendMetadata":self.journal.state().metadata,
2685            "sessionType":self.session_type,
2686            "sourceSessionType":self.source_session_type,
2687            "channel":self.channel,
2688            "freeTime":self.free_time,
2689            "orchestration":self.orchestration,
2690            "provenanceId":self.provenance_id,
2691            "rustLibSessionId":self.rust_lib_session_id,
2692            "rootNodeIds":self.root_node_ids,
2693            "referenceRootNodeIds":self.reference_root_node_ids,
2694            "startedAt":self.started_at,
2695            "transcript":self.transcript,
2696            "pendingTurn":self.pending_turn,
2697            "pendingExternalEventId":self.pending_external_event_id,
2698            "roundsUsed":self.rounds_used,
2699            "completed":self.completed,
2700            "sessionObjectId":self.journal.state().completed_session_object,
2701            "commitReceipt":self.commit_receipt,
2702            "commitAuthor":self.commit_author,
2703            "providerModel":self.runtime.model,
2704            "kwebPlan":self.plan,
2705            "boxCount":self.journal.state().boxes.len(),
2706            "eventCount":self.journal.state().events.len(),
2707            "boxes":self.journal.state().boxes,
2708            "events":self.journal.state().events,
2709            "context":projection,
2710            "sessionStatus":session_status,
2711            "chatendText":chatend_text,
2712        }))
2713    }
2714
2715    pub async fn release_managed_sources(&self) {
2716        self.api
2717            .release_managed_sources(&self.rust_lib_session_id)
2718            .await;
2719    }
2720}
2721
2722impl<C, F> kcode_agent_runtime::SessionHost for KennedySessionHost<'_, C>
2723where
2724    C: FnMut(Value) -> F + Send,
2725    F: Future<Output = anyhow::Result<()>> + Send,
2726{
2727    fn prepare_round<'a>(
2728        &'a mut self,
2729        round: u64,
2730    ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::RoundPreparation> {
2731        Box::pin(async move {
2732            self.session.rounds_used = round;
2733            self.session.refresh_runtime_prompt()?;
2734            self.deadline_after_response = self.session.prepare_free_time_round()?;
2735            let external_event_id = self.session.pending_external_event_id.clone();
2736            if self
2737                .session
2738                .recover_context_overflow(external_event_id.as_deref(), &[])?
2739                == ContextRecovery::Irreducible
2740                || (matches!(self.session.mode, AgentMode::Ingress { .. })
2741                    && self.session.ingress_force_commit_requested())
2742            {
2743                return Ok(kcode_agent_runtime::RoundPreparation::Complete(None));
2744            }
2745            let timeout = self.session.agent_request_timeout();
2746            self.session.begin_provider_call_budget(timeout);
2747            let input = self.session.projection().render();
2748            Ok(kcode_agent_runtime::RoundPreparation::Run(
2749                kcode_agent_runtime::PreparedRound {
2750                    input,
2751                    model: self.session.runtime.model.clone(),
2752                    reasoning_effort: self.session.runtime.reasoning_effort.clone(),
2753                    tool_description: call_ktool_description(),
2754                    timeout,
2755                },
2756            ))
2757        })
2758    }
2759
2760    fn record<'a>(
2761        &'a mut self,
2762        event: kcode_agent_runtime::SessionEvent,
2763    ) -> kcode_agent_runtime::HostFuture<'a, ()> {
2764        Box::pin(async move {
2765            match event {
2766                kcode_agent_runtime::SessionEvent::InferenceSubmitted {
2767                    manifest_hash,
2768                    model,
2769                    ..
2770                } => {
2771                    let projection = self.session.projection();
2772                    self.accounting = Some(kcode_intelligence_chatend::TopLevelCall::new(
2773                        manifest_hash.clone(),
2774                        model,
2775                    ));
2776                    self.session.journal.record(
2777                        now(),
2778                        EventKind::InferenceSubmitted {
2779                            manifest_hash,
2780                            estimated_input_tokens: projection.estimated_tokens,
2781                            raw_estimated_input_tokens: Some(projection.raw_estimated_tokens),
2782                        },
2783                    )?;
2784                }
2785                kcode_agent_runtime::SessionEvent::ProviderInput { input, .. } => {
2786                    let input = kcode_session_runtime_budget::without_footer_lines(&input);
2787                    self.session.journal.record(
2788                        now(),
2789                        EventKind::Note {
2790                            label: "provider_input".into(),
2791                            value: Value::String(input),
2792                        },
2793                    )?;
2794                }
2795                kcode_agent_runtime::SessionEvent::UsageUpdated { usage, .. } => {
2796                    self.accounting
2797                        .as_mut()
2798                        .context("provider usage arrived before inference submission")?
2799                        .usage_updated(&mut self.session.journal, &now(), &usage)?;
2800                }
2801                kcode_agent_runtime::SessionEvent::ProviderReceipt { usage, .. } => {
2802                    self.accounting
2803                        .take()
2804                        .context("provider receipt arrived before inference submission")?
2805                        .completed(&mut self.session.journal, &now(), usage.as_ref())?;
2806                }
2807            }
2808            let snapshot = self.session.snapshot()?;
2809            (self.checkpoint)(snapshot).await
2810        })
2811    }
2812
2813    fn execute_tool<'a>(
2814        &'a mut self,
2815        call: anyhow::Result<kcode_agent_runtime::ToolCall>,
2816        operation_id: Uuid,
2817    ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionToolOutcome> {
2818        Box::pin(async move {
2819            if let Some(pending) = &self.pending_freeform_write {
2820                let text = format!(
2821                    "{} is awaiting the complete file contents; no other Ktool can run before that output.",
2822                    pending.request.write_tool()
2823                );
2824                self.session
2825                    .record_tool_completion(None, json!({"ok":false,"result":text}))?;
2826                let footer = self.session.projection().footer;
2827                let text = provider_tool_result_with_context_footer(&footer, &text);
2828                (self.checkpoint)(self.session.snapshot()?).await?;
2829                return Ok(kcode_agent_runtime::SessionToolOutcome {
2830                    text,
2831                    ok: false,
2832                    capture: Some(json!(true)),
2833                    stop: false,
2834                    finish_after_round: false,
2835                    emitted_response: false,
2836                });
2837            }
2838            let tool_started_at = std::time::Instant::now();
2839            let mut created_call_box_id = None;
2840            let mut recorded_invocation = None;
2841            let transcript_start = self.session.transcript.len();
2842            let mut emitted_response = false;
2843            let mut outcome = match call {
2844                Ok(call) => {
2845                    let call = ToolCall {
2846                        name: call.name,
2847                        arguments: call.arguments,
2848                    };
2849                    let call_name = format!("Kennedy tool call: {}", call.name);
2850                    let call_content = tool_call_box_content(&call)?;
2851                    recorded_invocation = Some(
2852                        self.session
2853                            .record_tool_invocation(&call.name, call.arguments.clone())?,
2854                    );
2855                    created_call_box_id = Some(self.session.journal.create_box(
2856                        now(),
2857                        call_name,
2858                        BoxOwner::Kennedy,
2859                        call_content,
2860                    )?);
2861                    let external_event_id = self.session.pending_external_event_id.clone();
2862                    if self
2863                        .session
2864                        .recover_context_overflow(external_event_id.as_deref(), &[])?
2865                        == ContextRecovery::Irreducible
2866                    {
2867                        ToolOutcome {
2868                            text: CONTEXT_OVERFLOW_WARNING.into(),
2869                            store_result: false,
2870                            ok: false,
2871                            end_session: false,
2872                            freeform_write: None,
2873                            managed_source_snapshot: None,
2874                        }
2875                    } else {
2876                        match self.session.execute_tool(&call, operation_id).await {
2877                            Ok(outcome) => {
2878                                emitted_response = call.name == "EmitObject" && outcome.ok;
2879                                outcome
2880                            }
2881                            Err(error) => ToolOutcome {
2882                                text: format!("{} failed: {error}", call.name),
2883                                store_result: call.name != "LoadNodes",
2884                                ok: false,
2885                                end_session: false,
2886                                freeform_write: None,
2887                                managed_source_snapshot: None,
2888                            },
2889                        }
2890                    }
2891                }
2892                Err(error) => ToolOutcome {
2893                    text: error.to_string(),
2894                    store_result: true,
2895                    ok: false,
2896                    end_session: false,
2897                    freeform_write: None,
2898                    managed_source_snapshot: None,
2899                },
2900            };
2901            if let Some(snapshot) = outcome.managed_source_snapshot.take() {
2902                apply_snapshot(&mut self.session.journal, &now(), snapshot)?;
2903                outcome.store_result = false;
2904            }
2905            append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
2906            let capture = if let Some(request) = outcome.freeform_write.take() {
2907                self.pending_freeform_write = Some(PendingFreeformWrite {
2908                    request,
2909                    call_box_id: created_call_box_id
2910                        .context("freeform write call box was not created")?,
2911                });
2912                Some(json!(true))
2913            } else {
2914                None
2915            };
2916            if outcome.store_result {
2917                self.session.journal.create_box(
2918                    now(),
2919                    "Kennedy tool result",
2920                    BoxOwner::Controller,
2921                    BoxContent::text(&outcome.text),
2922                )?;
2923            }
2924            let external_event_id = self.session.pending_external_event_id.clone();
2925            let recovery = self
2926                .session
2927                .recover_context_overflow(external_event_id.as_deref(), &[])?;
2928            let context_warning_added =
2929                self.session.transcript[transcript_start..]
2930                    .iter()
2931                    .any(|entry| {
2932                        entry.get("contextOverflowWarning").and_then(Value::as_bool) == Some(true)
2933                    });
2934            let mut provider_text = outcome.text.clone();
2935            if context_warning_added && !provider_text.contains(CONTEXT_OVERFLOW_WARNING) {
2936                if !provider_text.is_empty() {
2937                    provider_text.push_str("\n\n");
2938                }
2939                provider_text.push_str(CONTEXT_OVERFLOW_WARNING);
2940            }
2941            let footer = self.session.projection().footer;
2942            let provider_text = provider_tool_result_with_context_footer(&footer, &provider_text);
2943            self.session.record_tool_completion(
2944                recorded_invocation.as_ref(),
2945                json!({"ok":outcome.ok,"result":outcome.text}),
2946            )?;
2947            let stop = recovery == ContextRecovery::Irreducible
2948                || (matches!(self.session.mode, AgentMode::Ingress { .. })
2949                    && self.session.ingress_force_commit_requested())
2950                || (!matches!(self.session.mode, AgentMode::Ingress { .. })
2951                    && self.session.journal.state().source_terminated);
2952            let snapshot = self.session.snapshot()?;
2953            (self.checkpoint)(snapshot).await?;
2954            Ok(kcode_agent_runtime::SessionToolOutcome {
2955                text: provider_text,
2956                ok: outcome.ok,
2957                capture,
2958                stop,
2959                finish_after_round: outcome.end_session,
2960                emitted_response,
2961            })
2962        })
2963    }
2964
2965    fn complete_capture<'a>(
2966        &'a mut self,
2967        _capture: Value,
2968        contents: String,
2969    ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
2970        Box::pin(async move {
2971            let pending = self
2972                .pending_freeform_write
2973                .take()
2974                .context("provider completed without a pending freeform write")?;
2975            let result_metadata = pending.request.clone();
2976            let outcome = self
2977                .session
2978                .complete_freeform_write(pending, contents)
2979                .await?;
2980            if outcome.store_result {
2981                self.session.journal.create_box(
2982                    now(),
2983                    "Kennedy tool result",
2984                    BoxOwner::Controller,
2985                    BoxContent::text(&outcome.text),
2986                )?;
2987            }
2988            self.session.journal.record(
2989                now(),
2990                EventKind::Note {
2991                    label: "write_file_freeform_result".into(),
2992                    value: result_metadata.result_record(outcome.ok, &outcome.text),
2993                },
2994            )?;
2995            let external_event_id = self.session.pending_external_event_id.clone();
2996            let recovery = self
2997                .session
2998                .recover_context_overflow(external_event_id.as_deref(), &[])?;
2999            let snapshot = self.session.snapshot()?;
3000            (self.checkpoint)(snapshot).await?;
3001            if recovery == ContextRecovery::Irreducible
3002                || (matches!(self.session.mode, AgentMode::Ingress { .. })
3003                    && self.session.ingress_force_commit_requested())
3004                || (!matches!(self.session.mode, AgentMode::Ingress { .. })
3005                    && self.session.journal.state().source_terminated)
3006                || self.deadline_after_response
3007            {
3008                return Ok(kcode_agent_runtime::SessionControl::Complete(None));
3009            }
3010            self.session.journal.create_box(
3011                now(),
3012                controller_box_name(&self.session.mode),
3013                BoxOwner::Controller,
3014                BoxContent::text(controller_message(
3015                    &self.session.mode,
3016                    &self.session.free_time,
3017                )),
3018            )?;
3019            Ok(kcode_agent_runtime::SessionControl::Continue)
3020        })
3021    }
3022
3023    fn complete_round<'a>(
3024        &'a mut self,
3025        completion: kcode_agent_runtime::RoundCompletion,
3026    ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
3027        Box::pin(async move {
3028            let answer = completion.answer.trim().to_owned();
3029            let mut completion_recovery = ContextRecovery::NotNeeded;
3030            if !answer.is_empty() {
3031                let mut content = BoxContent::text(answer.clone());
3032                if let Some(id) = &self.session.pending_external_event_id {
3033                    content.metadata["externalEventId"] = json!(id);
3034                }
3035                self.session.journal.create_box(
3036                    now(),
3037                    "Kennedy message",
3038                    BoxOwner::Kennedy,
3039                    content,
3040                )?;
3041                let mut transcript = json!({"role":"kennedy","content":answer});
3042                if let Some(id) = &self.session.pending_external_event_id {
3043                    transcript["externalEventId"] = json!(id);
3044                }
3045                self.session.transcript.push(transcript);
3046                let external_event_id = self.session.pending_external_event_id.clone();
3047                completion_recovery = self
3048                    .session
3049                    .recover_context_overflow(external_event_id.as_deref(), &[])?;
3050            }
3051            let snapshot = self.session.snapshot()?;
3052            (self.checkpoint)(snapshot).await?;
3053            if completion_recovery == ContextRecovery::Irreducible
3054                || (matches!(self.session.mode, AgentMode::Ingress { .. })
3055                    && self.session.ingress_force_commit_requested())
3056                || (!matches!(self.session.mode, AgentMode::Ingress { .. })
3057                    && self.session.journal.state().source_terminated)
3058            {
3059                return Ok(kcode_agent_runtime::SessionControl::Complete(None));
3060            }
3061            if completion.finish_requested || self.deadline_after_response {
3062                return Ok(kcode_agent_runtime::SessionControl::Complete(
3063                    (!answer.is_empty()).then_some(answer),
3064                ));
3065            }
3066            if matches!(self.session.mode, AgentMode::Conversation) && !answer.is_empty() {
3067                return Ok(kcode_agent_runtime::SessionControl::Complete(Some(answer)));
3068            }
3069            if matches!(self.session.mode, AgentMode::Conversation) && completion.emitted_response {
3070                return Ok(kcode_agent_runtime::SessionControl::Complete(None));
3071            }
3072            let solo_ingress_response =
3073                matches!(self.session.mode, AgentMode::Ingress { .. }) && !answer.is_empty();
3074            anyhow::ensure!(
3075                completion.used_tool || solo_ingress_response,
3076                "provider completed without a response or tool call"
3077            );
3078            self.session.journal.create_box(
3079                now(),
3080                controller_box_name(&self.session.mode),
3081                BoxOwner::Controller,
3082                BoxContent::text(controller_message(
3083                    &self.session.mode,
3084                    &self.session.free_time,
3085                )),
3086            )?;
3087            Ok(kcode_agent_runtime::SessionControl::Continue)
3088        })
3089    }
3090}
3091
3092impl kcode_agent_runtime::Host for KennedySubagentHost<'_> {
3093    fn render_tool_call(&mut self, call: &kcode_agent_runtime::ToolCall) -> anyhow::Result<String> {
3094        Ok(tool_call_box_content(&ToolCall {
3095            name: call.name.clone(),
3096            arguments: call.arguments.clone(),
3097        })?
3098        .text)
3099    }
3100
3101    fn execute_tool<'a>(
3102        &'a mut self,
3103        call: kcode_agent_runtime::ToolCall,
3104        operation_id: Uuid,
3105        budget: kcode_agent_runtime::ContextBudget,
3106    ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
3107        Box::pin(async move {
3108            let call = ToolCall {
3109                name: call.name,
3110                arguments: call.arguments,
3111            };
3112            if let Some(reason) = subagent_unavailable_reason(&call.name) {
3113                return Ok(kcode_agent_runtime::ToolOutcome::failure(reason));
3114            }
3115            if budget.estimated_tokens() > budget.max_input_tokens() {
3116                return Ok(kcode_agent_runtime::ToolOutcome::failure(
3117                    "The Ktool call was not run because its retained invocation would exceed the subagent context limit.",
3118                ));
3119            }
3120            if !subagent_managed_write_fits(&self.context, &call, &budget) {
3121                return Ok(kcode_agent_runtime::ToolOutcome::failure(
3122                    "The managed-source write was not run because its resulting current state would exceed the subagent context limit.",
3123                ));
3124            }
3125
3126            let tool_started_at = std::time::Instant::now();
3127
3128            if call.name == "LoadNodes" {
3129                let Some(DecodedTool::LoadNodes(identifiers)) =
3130                    decode(&call.name, &call.arguments)?
3131                else {
3132                    return Ok(kcode_agent_runtime::ToolOutcome::failure(
3133                        "LoadNodes did not match its tool contract.",
3134                    ));
3135                };
3136                load_durable_batch(
3137                    self.session.api.kmap(),
3138                    self.context.kweb_mut(),
3139                    &identifiers,
3140                )?;
3141                let (updates, creates) = kweb_plan_projection(&self.session.plan);
3142                let changes = self.context.reconcile_kweb(&updates, &creates)?;
3143                let displayed_state_keys = changes.displayed_state_keys();
3144                let mut text = if changes.is_empty() {
3145                    "LoadNodes completed. The subagent Kweb projection was already current.".into()
3146                } else {
3147                    changes.display_text()
3148                };
3149                append_slow_tool_duration(&mut text, tool_started_at.elapsed());
3150                return Ok(kcode_agent_runtime::ToolOutcome {
3151                    text,
3152                    ok: true,
3153                    state_updates: changes.updates,
3154                    displayed_state_keys,
3155                    capture: None,
3156                });
3157            }
3158
3159            if is_kweb_mutation(&call.name) {
3160                self.session.assert_tool_allowed(&call.name)?;
3161                let decoded = decode(&call.name, &call.arguments)?
3162                    .with_context(|| format!("{} did not match its tool contract", call.name))?;
3163                let referenced_pending = referenced_pending_nodes(&decoded);
3164                let prior_create_count = self.session.plan.creates.len();
3165                let mut text = execute_kweb_mutation(
3166                    &call.name,
3167                    decoded,
3168                    self.context.kweb(),
3169                    &mut self.session.plan,
3170                    &mut self.session.journal,
3171                )?;
3172                self.context.include_staged_nodes(
3173                    referenced_pending.into_iter().chain(
3174                        self.session.plan.creates[prior_create_count..]
3175                            .iter()
3176                            .map(|create| create.pending_id.clone()),
3177                    ),
3178                );
3179                let (updates, creates) = kweb_plan_projection(&self.session.plan);
3180                let changes = self.context.reconcile_kweb(&updates, &creates)?;
3181                append_slow_tool_duration(&mut text, tool_started_at.elapsed());
3182                return Ok(kcode_agent_runtime::ToolOutcome {
3183                    text,
3184                    ok: true,
3185                    state_updates: changes.updates,
3186                    displayed_state_keys: Vec::new(),
3187                    capture: None,
3188                });
3189            }
3190
3191            if let Some(request) = decode_freeform_write(&call.name, &call.arguments)? {
3192                if !self.context.source_is_open(request.kind(), request.name()) {
3193                    return Ok(kcode_agent_runtime::ToolOutcome::failure(format!(
3194                        "{} {:?} is not open in this subagent context. Call {} first.",
3195                        request.kind().label(),
3196                        request.name(),
3197                        request.kind().open_tool()
3198                    )));
3199                }
3200                let acknowledgement = request.acknowledgement();
3201                let id = Uuid::new_v4().to_string();
3202                self.captures.insert(id.clone(), request);
3203                return Ok(kcode_agent_runtime::ToolOutcome {
3204                    text: acknowledgement,
3205                    ok: true,
3206                    state_updates: Vec::new(),
3207                    displayed_state_keys: Vec::new(),
3208                    capture: Some(Value::String(id)),
3209                });
3210            }
3211
3212            let mut outcome = match self.session.execute_tool(&call, operation_id).await {
3213                Ok(outcome) => outcome,
3214                Err(error) => {
3215                    let mut text = format!("{} failed: {error}", call.name);
3216                    append_slow_tool_duration(&mut text, tool_started_at.elapsed());
3217                    return Ok(kcode_agent_runtime::ToolOutcome::failure(text));
3218                }
3219            };
3220            let displays_managed_snapshot = outcome
3221                .managed_source_snapshot
3222                .as_ref()
3223                .is_some_and(|snapshot| result_displays_snapshot(&outcome.text, snapshot));
3224            append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
3225            let (state_updates, displayed_state_keys) =
3226                if let Some(snapshot) = outcome.managed_source_snapshot.take() {
3227                    let state = self.context.apply_source_snapshot(snapshot);
3228                    let displayed = displays_managed_snapshot.then_some(state.key);
3229                    (
3230                        state.update.into_iter().collect(),
3231                        displayed.into_iter().collect(),
3232                    )
3233                } else {
3234                    (Vec::new(), Vec::new())
3235                };
3236            let capture = outcome.freeform_write.take().map(|request| {
3237                let id = Uuid::new_v4().to_string();
3238                self.captures.insert(id.clone(), request);
3239                Value::String(id)
3240            });
3241            Ok(kcode_agent_runtime::ToolOutcome {
3242                text: outcome.text,
3243                ok: outcome.ok,
3244                state_updates,
3245                displayed_state_keys,
3246                capture,
3247            })
3248        })
3249    }
3250
3251    fn complete_capture<'a>(
3252        &'a mut self,
3253        capture: Value,
3254        contents: String,
3255        budget: kcode_agent_runtime::ContextBudget,
3256    ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
3257        Box::pin(async move {
3258            let id = capture
3259                .as_str()
3260                .context("subagent freeform capture token is invalid")?;
3261            let request = self
3262                .captures
3263                .remove(id)
3264                .context("subagent freeform capture token is unknown")?;
3265            self.session
3266                .complete_subagent_freeform_write(&mut self.context, request, contents, &budget)
3267                .await
3268        })
3269    }
3270
3271    fn record(&mut self, event: kcode_agent_runtime::AuditEvent) -> anyhow::Result<()> {
3272        kcode_intelligence_chatend::record_subagent_event(&mut self.session.journal, &now(), &event)
3273    }
3274}
3275
3276fn cost_summary(label: &str, estimated_cost_usd_nanos: u64, unpriced_calls: u64) -> String {
3277    render(RenderRequest::CostSummary {
3278        label,
3279        estimated_cost_usd_nanos,
3280        unpriced_calls,
3281    })
3282    .expect("cost-summary rendering is infallible")
3283}
3284
3285fn restore_kweb_context(journal: &HistorySession, context: &mut KwebContext) -> anyhow::Result<()> {
3286    let Some(tool) = journal.state().tools.get(KWEB_TOOL_INSTANCE) else {
3287        return Ok(());
3288    };
3289    let mut nodes = BTreeMap::new();
3290    for slot in &tool.slots {
3291        let state = journal
3292            .state()
3293            .box_state(slot.box_id)
3294            .context("Kweb slot references a missing box")?;
3295        if let Some(node) = state.canonical.content.metadata.get("storedNode") {
3296            let node = match serde_json::from_value::<KwebNode>(node.clone()) {
3297                Ok(node) => node,
3298                Err(_) => node_from_value(node).context("decoding a stored Kweb context node")?,
3299            };
3300            nodes.insert(node.id.clone(), node);
3301        }
3302    }
3303    let mut direct = journal
3304        .state()
3305        .events
3306        .iter()
3307        .flat_map(|event| {
3308            let EventKind::ToolInvoked {
3309                tool_name,
3310                arguments,
3311                ..
3312            } = &event.kind
3313            else {
3314                return Vec::new();
3315            };
3316            match tool_name.as_str() {
3317                "LoadNodes" => arguments
3318                    .get("identifiers")
3319                    .and_then(Value::as_array)
3320                    .into_iter()
3321                    .flatten()
3322                    .filter_map(Value::as_str)
3323                    .map(str::to_owned)
3324                    .collect(),
3325                "LoadNode" => arguments
3326                    .get("identifier")
3327                    .and_then(Value::as_str)
3328                    .map(str::to_owned)
3329                    .into_iter()
3330                    .collect(),
3331                _ => Vec::new(),
3332            }
3333        })
3334        .collect::<Vec<_>>();
3335    if direct.is_empty() {
3336        direct = context.root_node_ids().to_vec();
3337    }
3338    context
3339        .restore(nodes.into_values(), direct)
3340        .map_err(anyhow::Error::new)
3341}
3342
3343fn kweb_node_draft(node: &PlannedNode) -> NodeDraft {
3344    NodeDraft {
3345        short_name: node.short_name.clone(),
3346        short_description: node.short_description.clone(),
3347        long_description: node.long_description.clone(),
3348        owner: node.owner.clone(),
3349        fixed_connections: node.fixed_connections.clone(),
3350        recent_connections: node.recent_connections.clone(),
3351        objects: node.objects.clone(),
3352    }
3353}
3354
3355fn planned_node(node: &KwebNode) -> PlannedNode {
3356    PlannedNode {
3357        short_name: node.short_name.clone(),
3358        short_description: node.short_description.clone(),
3359        long_description: node.long_description.clone(),
3360        owner: node.owner.clone(),
3361        fixed_connections: node
3362            .fixed_connections
3363            .iter()
3364            .map(|connection| connection.id.clone())
3365            .collect(),
3366        recent_connections: node
3367            .recent_connections
3368            .iter()
3369            .map(|connection| connection.id.clone())
3370            .collect(),
3371        objects: node.objects.clone(),
3372        attach_session_archive: true,
3373    }
3374}
3375
3376fn session_kind(session_type: &str, mode: &AgentMode) -> SessionKind {
3377    if matches!(mode, AgentMode::Ingress { .. }) {
3378        return SessionKind::HistoryIngress;
3379    }
3380    match session_type {
3381        "conversation" => SessionKind::Conversation,
3382        "telegram" => SessionKind::Telegram,
3383        "telegram-group" => SessionKind::TelegramGroup,
3384        "free-time" => SessionKind::SelfTime,
3385        "wakeup" => SessionKind::Other("wakeup".into()),
3386        "audio" => SessionKind::AudioIngress,
3387        other => SessionKind::Other(other.into()),
3388    }
3389}
3390
3391fn tool_instance(name: &str) -> String {
3392    if name == "LoadNodes" {
3393        return KWEB_TOOL_INSTANCE.into();
3394    }
3395    format!("{name}:{}", Uuid::new_v4())
3396}
3397
3398fn canonical_id(value: &str) -> anyhow::Result<String> {
3399    value
3400        .parse::<NodeId>()
3401        .with_context(|| format!("{value:?} is not a canonical node ID"))?;
3402    Ok(value.into())
3403}
3404
3405fn image_extension(media_type: &str) -> &'static str {
3406    match media_type
3407        .split(';')
3408        .next()
3409        .unwrap_or(media_type)
3410        .trim()
3411        .to_ascii_lowercase()
3412        .as_str()
3413    {
3414        "image/jpeg" => "jpg",
3415        "image/webp" => "webp",
3416        _ => "png",
3417    }
3418}
3419
3420fn call_ktool_description() -> String {
3421    render(RenderRequest::CallKtoolDescription).expect("Ktool-description rendering is infallible")
3422}
3423
3424fn now() -> String {
3425    Utc::now().to_rfc3339()
3426}
3427
3428fn deadline(value: &Value) -> Option<DateTime<Utc>> {
3429    value
3430        .get("deadlineAt")
3431        .and_then(Value::as_str)
3432        .and_then(|value| DateTime::parse_from_rfc3339(value).ok())
3433        .map(|value| value.with_timezone(&Utc))
3434}
3435
3436fn remaining_until(deadline: DateTime<Utc>) -> Duration {
3437    (deadline - Utc::now()).to_std().unwrap_or(Duration::ZERO)
3438}
3439
3440fn controller_box_name(mode: &AgentMode) -> &'static str {
3441    match mode {
3442        AgentMode::Conversation => "Turn continuation",
3443        AgentMode::FreeTime => "Self-time continuation",
3444        AgentMode::Wakeup => "Wakeup continuation",
3445        AgentMode::Ingress { .. } => "History-ingress continuation",
3446    }
3447}
3448
3449fn controller_message(mode: &AgentMode, free_time: &Value) -> String {
3450    let mode = match mode {
3451        AgentMode::Conversation => "conversation",
3452        AgentMode::FreeTime => "free-time",
3453        AgentMode::Wakeup => "wakeup",
3454        AgentMode::Ingress { .. } => "ingress",
3455    };
3456    render(RenderRequest::ControllerMessage { mode, free_time })
3457        .expect("known controller modes render successfully")
3458}
3459
3460#[cfg(test)]
3461mod tests {
3462    use std::time::{SystemTime, UNIX_EPOCH};
3463
3464    use kcode_session_history::{Config as HistoryConfig, NewSession, SessionHistory};
3465
3466    use super::*;
3467
3468    fn test_journal(label: &str) -> (std::path::PathBuf, HistorySession) {
3469        let root = std::env::temp_dir().join(format!(
3470            "kcode-kennedy-sessions-{label}-{}-{}",
3471            std::process::id(),
3472            SystemTime::now()
3473                .duration_since(UNIX_EPOCH)
3474                .unwrap()
3475                .as_nanos()
3476        ));
3477        let history = SessionHistory::open(HistoryConfig {
3478            directory: root.join("sessions"),
3479            completed_list: root.join("completed.jsonl"),
3480            provider_cost_compatibility: None,
3481        })
3482        .unwrap();
3483        let journal = history
3484            .create_session(NewSession {
3485                kind: SessionKind::SelfTime,
3486                created_at: "2026-08-05T00:00:00Z".into(),
3487                effective_context_tokens: 10_000,
3488                channel: Value::Null,
3489            })
3490            .unwrap();
3491        (root, journal)
3492    }
3493
3494    #[test]
3495    fn subagents_reject_parent_controls_but_allow_delivery_effects() {
3496        for unavailable in [
3497            "RunSubagent",
3498            "EndSession",
3499            "DehydrateBoxes",
3500            "SummarizeBox",
3501            "HydrateBox",
3502            "BoxesIntoObjects",
3503        ] {
3504            assert!(subagent_unavailable_reason(unavailable).is_some());
3505        }
3506        for delegated in [
3507            "EmitObject",
3508            "SendTelegramDM",
3509            "SendTelegramGroupMessage",
3510            "LoadNodes",
3511            "ExtractDocumentText",
3512        ] {
3513            assert_eq!(subagent_unavailable_reason(delegated), None);
3514        }
3515    }
3516
3517    #[test]
3518    fn only_complete_snapshot_results_claim_to_display_managed_state() {
3519        let snapshot = SourceSnapshot {
3520            kind: kcode_dev_tools::ManagedSourceKind::RustLibrary,
3521            name: "example".into(),
3522            text: "complete source".into(),
3523        };
3524
3525        assert!(result_displays_snapshot("complete source", &snapshot));
3526        assert!(!result_displays_snapshot(
3527            "Wrote file src/lib.rs in Rust library example.",
3528            &snapshot
3529        ));
3530    }
3531
3532    #[test]
3533    fn child_kweb_mutation_changes_plan_without_creating_parent_boxes() {
3534        let (root, mut journal) = test_journal("subagent-kweb");
3535        let node_id = "AAAAAAAE".to_owned();
3536        let mut context = KwebContext::new(vec![node_id.clone()]).unwrap();
3537        context
3538            .apply_load(
3539                KwebNode {
3540                    id: node_id.clone(),
3541                    short_name: "Old".into(),
3542                    short_description: "Old summary".into(),
3543                    long_description: "Old details".into(),
3544                    owner: "self".into(),
3545                    fixed_connections: Vec::new(),
3546                    recent_connections: Vec::new(),
3547                    objects: Vec::new(),
3548                    last_modified_by: "test".into(),
3549                    last_modified_at: None,
3550                },
3551                Vec::new(),
3552            )
3553            .unwrap();
3554        let mut plan = KwebPlan::default();
3555
3556        let result = execute_kweb_mutation(
3557            "UpdateNode",
3558            DecodedTool::UpdateNode {
3559                id: node_id.clone(),
3560                owner: "self".into(),
3561                short_name: "New".into(),
3562                short_description: "New summary".into(),
3563                long_description: "New details".into(),
3564            },
3565            &context,
3566            &mut plan,
3567            &mut journal,
3568        )
3569        .unwrap();
3570
3571        assert_eq!(result, format!("Staged the update to node {node_id}."));
3572        assert_eq!(plan.updates[&node_id].long_description, "New details");
3573        assert!(journal.state().boxes.is_empty());
3574        std::fs::remove_dir_all(root).unwrap();
3575    }
3576}