Skip to main content

kcode_kennedy_sessions/
lib.rs

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