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