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