Skip to main content

kcode_kennedy_sessions/
lib.rs

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