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