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