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