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