Skip to main content

kcode_kennedy_sessions/
lib.rs

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