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