1#![forbid(unsafe_code)]
4
5mod services;
6#[cfg(test)]
7mod tests;
8
9pub use kcode_kennedy_session_objects::ResolvedObject;
10pub use kcode_kennedy_turn_admission::PendingTurnAdmission;
11pub use kcode_telegram_session_coordinator::validate_file_name as validate_delivery_file_name;
12pub use services::{Api as Service, LocalServices as Capabilities};
13
14use std::{
15 collections::{BTreeMap, BTreeSet, HashMap},
16 future::Future,
17 sync::{Arc, Weak},
18 time::{Duration, Instant},
19};
20
21use anyhow::Context as _;
22use chrono::{DateTime, Utc};
23use kcode_agent_runtime::SessionHost as _;
24use kcode_commit_session::{CommitReceipt, CommitRequest};
25use kcode_dev_tools::{
26 ATTACH_OBJECT_WEB_LIB_TOOL, CALL_RUST_BIN_TOOL, RUST_BIN_TOOLS, RUST_LIB_TOOLS, WEB_LIB_TOOLS,
27 proposed_write_snapshot,
28};
29use kcode_dev_tools_chatend::{
30 FreeformWrite, SourceSnapshot, apply_snapshot, decode_freeform_write, prepare_freeform_write,
31};
32use kcode_history_ingress_context::{
33 Outcome as HistoryIngressContextOutcome, RecoveryOutcome as ContextRecoveryOutcome,
34};
35use kcode_kennedy_kweb_loader::{load_durable_batch, node_from_value};
36use kcode_kennedy_kweb_plan::{
37 Mutation as KwebMutation, Plan as KwebPlan, referenced_pending_nodes,
38};
39use kcode_kennedy_session_ingress::{is_terminal_external_response, restore_pending_turn};
40use kcode_kennedy_session_presentation::{RenderRequest, render};
41use kcode_kennedy_session_tool_contracts::{
42 DecodedTool, ManagedObjectArguments, ValidationRequest, decode, decode_managed_objects,
43 decode_note_to_self, validate,
44};
45use kcode_kennedy_session_tool_presentation::invocation_box_content;
46use kcode_kennedy_subagent_context::Context as SubagentContext;
47use kcode_kennedy_turn_admission::{
48 AdmissionKind, stage_turn_admission, validate_authoritative_user_event,
49};
50use kcode_kweb_context::{Context as KwebContext, Node as KwebNode};
51use kcode_kweb_db::NodeId;
52use kcode_server_object_envelopes::encode_file;
53use kcode_session_history::{
54 ErrorKind as HistoryErrorKind, LaunchSession as HistoryLaunchSession, NewSession,
55 Session as HistorySession,
56 chatend::{
57 BoxContent, BoxId, BoxOwner, CacheExpectation, ContextProjection, Event, EventId,
58 EventKind, PreparedProviderResume, ProviderContext, ProviderToolDefinition, SessionKind,
59 SessionMetadata,
60 },
61};
62use kcode_session_runtime_budget::{RoundBudget, RuntimeBudget, TimeBudget, TimeBudgetKind};
63use kcode_speaker_system::KTOOLS as SPEECH_CLASSIFICATION_TOOLS;
64use serde::{Deserialize, Serialize};
65use serde_json::{Value, json};
66use sha2::{Digest, Sha256};
67use uuid::Uuid;
68
69const BROWSER_CONVERSATION_REQUEST_TIMEOUT: Duration = Duration::from_secs(225 * 60);
70const HISTORY_INGRESS_REQUEST_TIMEOUT: Duration = Duration::from_secs(225 * 60);
71const HISTORY_INGRESS_ATTEMPT_DURATION: Duration = Duration::from_secs(45 * 60);
72const WAKEUP_REQUEST_TIMEOUT: Duration = Duration::from_secs(225 * 60);
73const SELF_TIME_HARD_STOP_ALLOWANCE: Duration = Duration::from_secs(15 * 60);
74const MAX_MEDIA_ENRICHMENT_BYTES: u64 = 20 * 1024 * 1024;
75const MAX_LAUNCH_INTENTS_PER_USER_TURN: usize = 10;
76const LAUNCH_SESSION_TOOL: &str = "LaunchSession";
77const KWEB_TOOL_INSTANCE: &str = "kweb";
78const TASK_BOARD_TOOLS: [&str; 8] = [
79 "CreateTaskCategory",
80 "GetTaskCategory",
81 "RemoveTaskCategory",
82 "CreateTask",
83 "GetTask",
84 "UpdateTask",
85 "RemoveTask",
86 "GetTopTaskOrphan",
87];
88const CONTEXT_OVERFLOW_WARNING_BOX_NAME: &str = "Context overflow warning";
89const 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";
90const INGRESS_FORCE_COMMIT_NOTE: &str = "ingress_force_commit";
91const PENDING_TURN_ADMISSION_KIND: &str = "pendingTurnAdmissionKind";
92const PENDING_TURN_ADMISSION_USER: &str = "user";
93const PENDING_TURN_ADMISSION_SOURCE: &str = "source";
94const CHECKPOINT_STATE_VERSION: u64 = 5;
95
96#[derive(Debug)]
97struct IngressTimeExpired;
98
99impl std::fmt::Display for IngressTimeExpired {
100 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
101 formatter.write_str("history ingress time expired before EndSession")
102 }
103}
104
105impl std::error::Error for IngressTimeExpired {}
106
107pub fn is_ingress_time_expired(error: &anyhow::Error) -> bool {
108 error.is::<IngressTimeExpired>()
109}
110
111fn ingress_time_remaining_at(deadline: &mut Option<Instant>, now: Instant) -> anyhow::Result<u64> {
112 let Some(current) = *deadline else {
113 *deadline = Some(
114 now.checked_add(HISTORY_INGRESS_ATTEMPT_DURATION)
115 .context("history ingress deadline overflow")?,
116 );
117 return Ok(HISTORY_INGRESS_ATTEMPT_DURATION.as_secs());
118 };
119 if now >= current {
120 return Err(anyhow::Error::new(IngressTimeExpired));
121 }
122 Ok(current.duration_since(now).as_secs())
123}
124
125#[derive(Clone, Debug)]
126pub struct RuntimeModel {
127 pub model: String,
128 pub reasoning_effort: String,
129 pub context_window_tokens: u64,
130}
131
132impl RuntimeModel {
133 pub fn from_intelligence(runtime: kcode_intelligence_router::RuntimeModel) -> Self {
134 Self {
135 model: runtime.model,
136 reasoning_effort: runtime.reasoning_effort,
137 context_window_tokens: runtime.context_window_tokens,
138 }
139 }
140
141 fn attribution(&self) -> String {
142 format!("{}-{}", self.model, self.reasoning_effort)
143 }
144}
145
146#[derive(Clone, Debug, PartialEq, Eq)]
147pub enum AgentMode {
148 Conversation,
149 FreeTime,
150 Wakeup,
151 Ingress { record_id: Option<String> },
152}
153
154#[derive(Clone, Copy, Debug, Eq, PartialEq)]
155pub enum TurnDeadlineKind {
156 Telegram,
157 SelfTimeHardStop,
158}
159
160#[derive(Clone, Copy, Debug, Eq, PartialEq)]
161pub struct TurnDeadline {
162 pub kind: TurnDeadlineKind,
163 pub at: DateTime<Utc>,
164}
165
166#[derive(Clone, Debug)]
167pub struct SessionOptions {
168 pub session_type: String,
169 pub root_node_ids: Vec<String>,
170 pub reference_root_node_ids: Vec<String>,
171 pub channel: Value,
172 pub free_time: Value,
173 pub orchestration: Value,
174 pub provenance_id: Option<String>,
175 pub mode: AgentMode,
176 pub source_session_type: Option<String>,
177 pub group_context: Value,
178 pub rust_lib_session_id: Option<String>,
179}
180
181impl SessionOptions {
182 pub fn conversation(session_type: impl Into<String>, roots: Vec<String>) -> Self {
183 Self {
184 session_type: session_type.into(),
185 root_node_ids: roots,
186 reference_root_node_ids: Vec::new(),
187 channel: Value::Null,
188 free_time: Value::Null,
189 orchestration: json!({"owner":"backend","status":"idle"}),
190 provenance_id: None,
191 mode: AgentMode::Conversation,
192 source_session_type: None,
193 group_context: Value::Null,
194 rust_lib_session_id: None,
195 }
196 }
197}
198
199fn restore_session_type(options: &mut SessionOptions, state: &Value) {
200 if !matches!(&options.mode, AgentMode::Ingress { .. }) {
201 options.session_type = state
202 .get("sessionType")
203 .and_then(Value::as_str)
204 .unwrap_or(&options.session_type)
205 .to_owned();
206 }
207}
208
209fn restore_commit_receipt(restored: Option<&Value>) -> anyhow::Result<Option<CommitReceipt>> {
210 restored
211 .and_then(|state| state.get("commitReceipt"))
212 .filter(|receipt| !receipt.is_null())
213 .cloned()
214 .map(serde_json::from_value)
215 .transpose()
216 .context("decoding the stored session commit receipt")
217}
218
219fn journal_kweb_plan(journal: &HistorySession) -> Option<&Value> {
220 journal
221 .state()
222 .current_ingress_attempt_events()
223 .iter()
224 .rev()
225 .find_map(|event| {
226 let EventKind::KwebPlanChanged { operation } = &event.kind else {
227 return None;
228 };
229 operation.get("plan")
230 })
231}
232
233#[derive(Clone, Debug, Deserialize)]
234#[serde(rename_all = "camelCase", deny_unknown_fields)]
235struct LaunchSessionArguments {
236 directive: String,
237 context_node_ids: Vec<String>,
238}
239
240#[derive(Clone, Debug, Deserialize, Serialize)]
241#[serde(rename_all = "camelCase")]
242struct LaunchIntent {
243 invocation_id: String,
244 user_turn_id: EventId,
245 started_at: String,
246 parent_session_id: String,
247 effective_context_tokens: u64,
248 root_node_ids: Vec<String>,
249 reference_root_node_ids: Vec<String>,
250 context_node_ids: Vec<String>,
251}
252
253#[derive(Serialize)]
254#[serde(rename_all = "camelCase")]
255struct LaunchSuccess<'a> {
256 session_id: &'a str,
257 command_id: &'a str,
258}
259
260fn decode_launch_session_arguments(value: &Value) -> anyhow::Result<LaunchSessionArguments> {
261 let arguments: LaunchSessionArguments =
262 serde_json::from_value(value.clone()).context("LaunchSession arguments are invalid")?;
263 anyhow::ensure!(
264 !arguments.directive.trim().is_empty(),
265 "LaunchSession directive must not be blank"
266 );
267 validate_canonical_distinct_ids(&arguments.context_node_ids, "context node")?;
268 Ok(arguments)
269}
270
271fn validate_canonical_distinct_ids(ids: &[String], label: &str) -> anyhow::Result<()> {
272 let mut seen = BTreeSet::new();
273 for id in ids {
274 canonical_id(id).with_context(|| format!("LaunchSession {label} ID is invalid"))?;
275 anyhow::ensure!(
276 seen.insert(id.as_str()),
277 "LaunchSession {label} ID {id} is duplicated"
278 );
279 }
280 Ok(())
281}
282
283fn validate_loaded_launch_context(
284 context: &KwebContext,
285 context_node_ids: &[String],
286) -> anyhow::Result<()> {
287 for id in context_node_ids {
288 anyhow::ensure!(
289 context.contains_full_node(id),
290 "LaunchSession context node {id} is not fully loaded in the parent context"
291 );
292 }
293 Ok(())
294}
295
296fn authoritative_user_box_event(journal: &HistorySession, event: &Event) -> Option<EventId> {
297 let EventKind::BoxCreated {
298 box_id,
299 owner: BoxOwner::User,
300 ..
301 } = &event.kind
302 else {
303 return None;
304 };
305 if *box_id != BoxId(event.id.0) {
306 return None;
307 }
308 journal
309 .state()
310 .box_state(*box_id)
311 .filter(|state| matches!(state.owner, BoxOwner::User))
312 .map(|_| event.id)
313}
314
315fn unique_user_box_event(
316 journal: &HistorySession,
317 events: &[Event],
318) -> anyhow::Result<Option<EventId>> {
319 let ids = events
320 .iter()
321 .filter_map(|event| authoritative_user_box_event(journal, event))
322 .collect::<Vec<_>>();
323 anyhow::ensure!(
324 ids.len() <= 1,
325 "multiple authoritative user inputs appeared in one launch-authority interval"
326 );
327 Ok(ids.into_iter().next())
328}
329
330#[derive(Clone)]
331struct RecoveredUserTurn {
332 id: EventId,
333 external_event_id: Option<String>,
334}
335
336enum RecoveredAdmission {
337 Unmarked(RecoveredUserTurn),
338 MarkedUser(RecoveredUserTurn),
339 MarkedSource,
340}
341
342fn classify_recovered_admission(
343 journal: &HistorySession,
344 event: &Event,
345) -> anyhow::Result<Option<RecoveredAdmission>> {
346 let Some(id) = authoritative_user_box_event(journal, event) else {
347 return Ok(None);
348 };
349 let Some(state) = journal.state().box_state(BoxId(id.0)) else {
350 return Ok(None);
351 };
352 let metadata = &state.canonical.content.metadata;
353 let record = RecoveredUserTurn {
354 id,
355 external_event_id: metadata
356 .get("externalEventId")
357 .and_then(Value::as_str)
358 .map(str::to_owned),
359 };
360 match metadata.get(PENDING_TURN_ADMISSION_KIND) {
361 None => Ok(Some(RecoveredAdmission::Unmarked(record))),
362 Some(Value::String(kind)) if kind == PENDING_TURN_ADMISSION_USER => {
363 validate_authoritative_user_event(journal, event.id)?;
364 Ok(Some(RecoveredAdmission::MarkedUser(record)))
365 }
366 Some(Value::String(kind)) if kind == PENDING_TURN_ADMISSION_SOURCE => {
367 Ok(Some(RecoveredAdmission::MarkedSource))
368 }
369 Some(_) => anyhow::bail!("pending turn admission marker is invalid"),
370 }
371}
372
373fn recovered_user_turns(
374 journal: &HistorySession,
375 events: &[Event],
376) -> anyhow::Result<Vec<RecoveredAdmission>> {
377 let mut found_unmarked = false;
378 let mut recovered = Vec::new();
379 for event in events {
380 match classify_recovered_admission(journal, event)? {
381 Some(RecoveredAdmission::Unmarked(record)) => {
382 anyhow::ensure!(
383 !found_unmarked,
384 "multiple unmarked authoritative user inputs appeared in one launch-authority interval"
385 );
386 found_unmarked = true;
387 recovered.push(RecoveredAdmission::Unmarked(record));
388 }
389 Some(admission) => recovered.push(admission),
390 None => {}
391 }
392 }
393 Ok(recovered)
394}
395
396fn validate_user_turn_id(journal: &HistorySession, id: EventId) -> anyhow::Result<()> {
397 let event = journal
398 .state()
399 .event(id)
400 .context("restored launch user-turn event does not exist")?;
401 anyhow::ensure!(
402 matches!(
403 classify_recovered_admission(journal, event)?,
404 Some(RecoveredAdmission::Unmarked(_) | RecoveredAdmission::MarkedUser(_))
405 ),
406 "restored launch user-turn event is not an authoritative user BoxCreated event"
407 );
408 Ok(())
409}
410
411fn consume_launch_bootstrap_marker(
412 orchestration: &mut Value,
413 launch_bootstrap_pending: &mut bool,
414) -> bool {
415 if !*launch_bootstrap_pending {
416 return false;
417 }
418 *launch_bootstrap_pending = false;
419 if !orchestration.is_object() {
420 *orchestration = json!({});
421 }
422 orchestration["launchBootstrapPending"] = json!(false);
423 true
424}
425
426fn grant_marked_user_launch_authority(
427 id: EventId,
428 orchestration: &mut Value,
429 launch_bootstrap_pending: &mut bool,
430 launch_user_turn_id: &mut Option<EventId>,
431) {
432 consume_launch_bootstrap_marker(orchestration, launch_bootstrap_pending);
433 *launch_user_turn_id = Some(id);
434}
435
436fn user_turn_launch_authority(
437 user_turn: Option<EventId>,
438 orchestration: &mut Value,
439 launch_bootstrap_pending: &mut bool,
440) -> Option<EventId> {
441 if consume_launch_bootstrap_marker(orchestration, launch_bootstrap_pending) {
442 None
443 } else {
444 user_turn
445 }
446}
447
448fn reconcile_recovered_launch_authority(
449 pending_turn: bool,
450 recovered_user_turn: Option<EventId>,
451 launch_provenance: &Value,
452 orchestration: &mut Value,
453 launch_bootstrap_pending: &mut bool,
454 launch_user_turn_id: &mut Option<EventId>,
455) {
456 if let Some(recovered) = recovered_user_turn {
457 if *launch_bootstrap_pending || !launch_provenance.is_null() {
458 consume_launch_bootstrap_marker(orchestration, launch_bootstrap_pending);
459 *launch_user_turn_id = None;
460 } else {
461 *launch_user_turn_id = Some(recovered);
462 }
463 }
464 if !pending_turn || *launch_bootstrap_pending {
465 *launch_user_turn_id = None;
466 }
467}
468
469struct RecoveredAdmissionReplayState<'a> {
470 pending_turn: &'a bool,
471 launch_provenance: &'a Value,
472 orchestration: &'a mut Value,
473 launch_bootstrap_pending: &'a mut bool,
474 launch_user_turn_id: &'a mut Option<EventId>,
475 rounds_used: &'a mut u64,
476 pending_external_event_id: &'a mut Option<String>,
477}
478
479fn replay_recovered_admissions(
480 recovered: Vec<RecoveredAdmission>,
481 state: RecoveredAdmissionReplayState<'_>,
482) {
483 for admission in recovered {
484 match admission {
485 RecoveredAdmission::Unmarked(record) => {
486 reconcile_recovered_launch_authority(
487 *state.pending_turn,
488 Some(record.id),
489 state.launch_provenance,
490 state.orchestration,
491 state.launch_bootstrap_pending,
492 state.launch_user_turn_id,
493 );
494 *state.rounds_used = 0;
495 *state.pending_external_event_id = record.external_event_id;
496 }
497 RecoveredAdmission::MarkedUser(record) => {
498 grant_marked_user_launch_authority(
499 record.id,
500 state.orchestration,
501 state.launch_bootstrap_pending,
502 state.launch_user_turn_id,
503 );
504 *state.rounds_used = 0;
505 *state.pending_external_event_id = record.external_event_id;
506 }
507 RecoveredAdmission::MarkedSource => {}
508 }
509 }
510}
511
512fn completed_invocation_ids(journal: &HistorySession) -> BTreeSet<String> {
513 journal
514 .state()
515 .events
516 .iter()
517 .filter_map(|event| {
518 let EventKind::ToolCompleted {
519 invocation_id: Some(id),
520 ..
521 } = &event.kind
522 else {
523 return None;
524 };
525 Some(id.clone())
526 })
527 .collect()
528}
529
530fn pruned_launch_intents(
531 intents: &[LaunchIntent],
532 current_turn: Option<EventId>,
533 completed: &BTreeSet<String>,
534) -> Vec<LaunchIntent> {
535 intents
536 .iter()
537 .filter(|intent| {
538 Some(intent.user_turn_id) == current_turn || !completed.contains(&intent.invocation_id)
539 })
540 .cloned()
541 .collect()
542}
543
544fn invocation_arguments<'a>(
545 journal: &'a HistorySession,
546 invocation_id: &str,
547) -> anyhow::Result<&'a Value> {
548 journal
549 .state()
550 .events
551 .iter()
552 .find_map(|event| {
553 let EventKind::ToolInvoked {
554 tool_name,
555 arguments,
556 invocation_id: Some(id),
557 ..
558 } = &event.kind
559 else {
560 return None;
561 };
562 (tool_name == LAUNCH_SESSION_TOOL && id == invocation_id).then_some(arguments)
563 })
564 .with_context(|| {
565 format!(
566 "launch intent {} has no matching ToolInvoked event",
567 invocation_id
568 )
569 })
570}
571
572fn launch_success_json(session_id: &str, command_id: &str) -> anyhow::Result<String> {
573 serde_json::to_string(&LaunchSuccess {
574 session_id,
575 command_id,
576 })
577 .context("serializing LaunchSession result")
578}
579
580fn tool_invocation_content(name: &str, arguments: &Value) -> anyhow::Result<BoxContent> {
581 if name == LAUNCH_SESSION_TOOL {
582 return Ok(BoxContent::text("LaunchSession"));
583 }
584 invocation_box_content(name, arguments)
585}
586
587pub struct Session {
588 api: Service,
589 subagent_codex_prompt: String,
590 runtime: RuntimeModel,
591 journal: HistorySession,
592 plan: KwebPlan,
593 pub session_type: String,
594 pub channel: Value,
595 pub free_time: Value,
596 pub orchestration: Value,
597 pub provenance_id: Option<String>,
598 pub rust_lib_session_id: String,
599 pub root_node_ids: Vec<String>,
600 pub reference_root_node_ids: Vec<String>,
601 pub started_at: String,
602 pub transcript: Vec<Value>,
603 pub pending_turn: bool,
604 pub pending_external_event_id: Option<String>,
605 pub completed: bool,
606 pub rounds_used: u64,
607 commit_receipt: Option<CommitReceipt>,
608 commit_author: String,
609 mode: AgentMode,
610 source_session_type: Option<String>,
611 group_context: Value,
612 context: KwebContext,
613 free_time_end_reason: Option<String>,
614 fatal_persistence_error: Option<String>,
615 active_provider_deadline: Option<DateTime<Utc>>,
616 active_turn_deadline: Option<TurnDeadline>,
617 provider_affinity: Option<ProviderAffinityState>,
618 next_thread_reset_reason: Option<String>,
619 ingress_deadline: Option<Instant>,
620 previous_ingress_attempt_timed_out: bool,
621 launch_provenance: Value,
622 launch_context_node_ids: Vec<String>,
623 launch_user_turn_id: Option<EventId>,
624 launch_intents: Vec<LaunchIntent>,
625 launch_bootstrap_pending: bool,
626 turn_lease_slot: Option<Weak<TurnLeaseToken>>,
627}
628
629struct TurnLeaseToken;
630
631struct TurnLease {
632 token: Arc<TurnLeaseToken>,
633}
634
635impl TurnLease {
636 fn acquire(slot: &mut Option<Weak<TurnLeaseToken>>) -> anyhow::Result<Self> {
637 anyhow::ensure!(
638 slot.as_ref().and_then(Weak::upgrade).is_none(),
639 "a stepped session turn is already active"
640 );
641 let token = Arc::new(TurnLeaseToken);
642 *slot = Some(Arc::downgrade(&token));
643 Ok(Self { token })
644 }
645
646 fn validate(&self, slot: &Option<Weak<TurnLeaseToken>>) -> anyhow::Result<()> {
647 anyhow::ensure!(
648 slot.as_ref()
649 .and_then(Weak::upgrade)
650 .is_some_and(|active| Arc::ptr_eq(&self.token, &active)),
651 "the stepped session turn is stale"
652 );
653 Ok(())
654 }
655}
656
657struct PrimaryTurnState {
658 accounting: Option<kcode_intelligence_chatend::TopLevelCall>,
659 pending_freeform_write: Option<PendingFreeformWrite>,
660 deadline_after_response: bool,
661 operation_id: Uuid,
662 prepared_cache: Option<PreparedCacheObservation>,
663 provider_synchronized_after: Option<EventId>,
664 restart_fresh_reason: Option<String>,
665 exact_tool_result: bool,
666 used_tool: bool,
667 finish_requested: bool,
668 emitted_response: bool,
669 pending_capture: Option<Value>,
670}
671
672#[must_use]
673pub struct SessionTurn {
674 lease: TurnLease,
675 user_id: String,
676 completed_rounds: u64,
677 round_limit: u64,
678 state: PrimaryTurnState,
679 at_yielded_boundary: bool,
680 admission_poisoned: bool,
681}
682
683enum PendingInferenceAction {
684 Start {
685 runtime: kcode_agent_runtime::AgentRuntime,
686 request: kcode_agent_runtime::SessionInferenceRequest,
687 },
688 Next {
689 inference: kcode_agent_runtime::SessionInference,
690 },
691 Respond {
692 inference: kcode_agent_runtime::SessionInference,
693 call_id: String,
694 result: kcode_codex_runtime_v2::ToolResult,
695 stop: bool,
696 },
697}
698
699#[must_use]
700pub struct PendingSessionInference {
701 turn: SessionTurn,
702 action: Box<PendingInferenceAction>,
703}
704
705enum SessionInferenceWakeKind {
706 Event {
707 inference: kcode_agent_runtime::SessionInference,
708 event: anyhow::Result<Option<kcode_agent_runtime::SessionInferenceEvent>>,
709 },
710 StartFailed(anyhow::Error),
711 RespondFailed {
712 inference: kcode_agent_runtime::SessionInference,
713 error: anyhow::Error,
714 },
715 RespondedStop {
716 inference: kcode_agent_runtime::SessionInference,
717 },
718}
719
720#[must_use]
721pub struct SessionInferenceWake {
722 turn: SessionTurn,
723 kind: SessionInferenceWakeKind,
724}
725
726#[must_use]
727pub enum TurnBoundary {
728 Await(PendingSessionInference),
729 Yield(SessionTurn),
730 Complete(Option<String>),
731}
732
733impl PendingSessionInference {
734 pub async fn wait(self) -> SessionInferenceWake {
735 let Self { turn, action } = self;
736 let kind = match *action {
737 PendingInferenceAction::Start { runtime, request } => {
738 match runtime.start_session_inference(request).await {
739 Ok(mut inference) => {
740 let event = inference.next_event().await;
741 SessionInferenceWakeKind::Event { inference, event }
742 }
743 Err(error) => SessionInferenceWakeKind::StartFailed(error),
744 }
745 }
746 PendingInferenceAction::Next { mut inference } => {
747 let event = inference.next_event().await;
748 SessionInferenceWakeKind::Event { inference, event }
749 }
750 PendingInferenceAction::Respond {
751 mut inference,
752 call_id,
753 result,
754 stop,
755 } => match inference.respond(&call_id, result).await {
756 Ok(()) if stop => SessionInferenceWakeKind::RespondedStop { inference },
757 Ok(()) => {
758 let event = inference.next_event().await;
759 SessionInferenceWakeKind::Event { inference, event }
760 }
761 Err(error) => SessionInferenceWakeKind::RespondFailed { inference, error },
762 },
763 };
764 SessionInferenceWake { turn, kind }
765 }
766}
767
768#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
769#[serde(rename_all = "camelCase")]
770struct ProviderAffinityState {
771 continuation: kcode_intelligence_router::AgentContinuation,
772 synchronized_event_id: EventId,
773 material_fingerprint: String,
774}
775
776#[derive(Debug, Eq, PartialEq)]
777enum NativeProviderResumePreparation {
778 Continue { marker_lines: Vec<String> },
779 RestartFresh { reason: String },
780}
781
782fn apply_prepared_provider_resume(
783 provider_affinity: &mut Option<ProviderAffinityState>,
784 next_thread_reset_reason: &mut Option<String>,
785 prepared: PreparedProviderResume,
786) -> NativeProviderResumePreparation {
787 match prepared.thread_reset_reason {
788 Some(reason) => {
789 *provider_affinity = None;
790 *next_thread_reset_reason = Some(reason.clone());
791 NativeProviderResumePreparation::RestartFresh { reason }
792 }
793 None => NativeProviderResumePreparation::Continue {
794 marker_lines: prepared.marker_lines,
795 },
796 }
797}
798
799fn restore_provider_affinity(
800 restored: Option<&Value>,
801 fresh_ingress_attempt: bool,
802) -> anyhow::Result<Option<ProviderAffinityState>> {
803 let state_version = restored
804 .and_then(|state| state.get("stateVersion"))
805 .and_then(Value::as_u64);
806 if fresh_ingress_attempt || state_version != Some(CHECKPOINT_STATE_VERSION) {
807 return Ok(None);
808 }
809 restored
810 .and_then(|state| state.get("providerAffinity"))
811 .filter(|value| !value.is_null())
812 .cloned()
813 .map(serde_json::from_value)
814 .transpose()
815 .context("restored provider affinity is invalid")
816}
817
818#[derive(Clone, Copy, Debug, Eq, PartialEq)]
819enum InputStage {
820 Accepted,
821}
822
823#[derive(Clone, Copy, Debug, Eq, PartialEq)]
824enum ContextRecovery {
825 NotNeeded,
826 Recovered,
827 Irreducible,
828}
829
830fn kweb_slot_box_ids(journal: &HistorySession) -> Vec<BoxId> {
831 journal
832 .state()
833 .tools
834 .get(KWEB_TOOL_INSTANCE)
835 .map(|tool| tool.slots.iter().map(|slot| slot.box_id).collect())
836 .unwrap_or_default()
837}
838
839fn load_box_changes(before: &[BoxId], after: &[BoxId], stale: &[BoxId]) -> Vec<BoxId> {
840 let before = before.iter().copied().collect::<BTreeSet<_>>();
841 let stale = stale.iter().copied().collect::<BTreeSet<_>>();
842 let mut emitted = BTreeSet::new();
843 after
844 .iter()
845 .copied()
846 .filter(|id| (!before.contains(id) || stale.contains(id)) && emitted.insert(*id))
847 .collect()
848}
849
850fn render_load_nodes_result(
851 journal: &HistorySession,
852 changed_box_ids: &[BoxId],
853 footer_lines: &[String],
854) -> anyhow::Result<String> {
855 let changed_box_ids = changed_box_ids
856 .iter()
857 .map(ToString::to_string)
858 .collect::<Vec<_>>();
859 let projected_boxes = journal
860 .state()
861 .projection_with_footer_lines(footer_lines)
862 .items
863 .into_iter()
864 .filter(|item| !item.marker)
865 .map(|item| (item.box_id.to_string(), item.text))
866 .collect::<Vec<_>>();
867 render(RenderRequest::LoadNodes {
868 changed_box_ids: &changed_box_ids,
869 projected_boxes: &projected_boxes,
870 })
871}
872
873fn provider_tool_result_with_context_footer(footer: &str, result: &str) -> String {
874 render(RenderRequest::ProviderFooter { result, footer })
875 .expect("provider-footer rendering is infallible")
876}
877
878fn completes_before_provider_resume(outcome: &kcode_agent_runtime::SessionToolOutcome) -> bool {
879 outcome.stop || (outcome.ok && outcome.finish_after_round)
880}
881
882fn append_slow_tool_duration(text: &mut String, elapsed: Duration) {
883 *text = render(RenderRequest::SlowTool { text, elapsed })
884 .expect("slow-tool rendering is infallible");
885}
886
887fn log_primary_thread_observation(
888 operation_id: Uuid,
889 round: u64,
890 requested_model: &str,
891 prepared: &PreparedCacheObservation,
892 provider_thread_id: Option<&str>,
893 input_tokens: u64,
894 cached_input_tokens: u64,
895) {
896 tracing::info!(
897 affinity_scope = "primary",
898 %operation_id,
899 round,
900 provider = prepared.provider,
901 requested_model,
902 model = prepared.model,
903 thread_action = prepared.thread_action,
904 provider_thread_id = provider_thread_id.unwrap_or(""),
905 thread_reset_reason = prepared.thread_reset_reason.as_deref().unwrap_or(""),
906 projection_hash = prepared.projection_hash,
907 provider_input_hash = prepared.provider_input_hash,
908 provider_input_bytes = prepared.provider_input_bytes,
909 input_tokens,
910 cached_input_tokens,
911 "Provider thread-affinity observation"
912 );
913}
914
915fn render_web_search_result(
916 result: &kcode_intelligence_router::SearchResponse,
917) -> anyhow::Result<String> {
918 let sources = result
919 .sources
920 .iter()
921 .map(|source| (source.title.clone(), source.url.clone()))
922 .collect::<Vec<_>>();
923 render(RenderRequest::WebSearch {
924 answer: &result.answer,
925 sources: &sources,
926 })
927}
928
929fn render_web_fetch_result(
930 result: &kcode_intelligence_router::FetchResponse,
931) -> anyhow::Result<String> {
932 render(RenderRequest::WebFetch {
933 url: &result.url,
934 title: result.title.as_deref(),
935 content_type: &result.content_type,
936 truncated: result.truncated,
937 content: &result.content,
938 })
939}
940
941fn render_media_annotation_result(
942 object_id: &str,
943 file_name: &str,
944 content_type: &str,
945 result: &kcode_intelligence_router::AnnotationResponse,
946) -> anyhow::Result<String> {
947 render(RenderRequest::MediaAnnotation {
948 object_id,
949 file_name,
950 content_type,
951 model: &result.model,
952 complete: result.complete,
953 incomplete_reason: result.incomplete_reason.as_deref(),
954 text: &result.text,
955 })
956}
957
958fn render_audio_transcription_result(
959 object_id: &str,
960 file_name: &str,
961 content_type: &str,
962 result: &kcode_intelligence_router::TranscriptionResponse,
963) -> anyhow::Result<String> {
964 render(RenderRequest::AudioTranscription {
965 object_id,
966 file_name,
967 content_type,
968 model: &result.model,
969 text: &result.text,
970 })
971}
972
973fn render_document_extraction_result(
974 object_id: &str,
975 file_name: &str,
976 result: &kcode_intelligence_router::DocumentExtraction,
977) -> anyhow::Result<String> {
978 render(RenderRequest::DocumentExtraction {
979 object_id,
980 file_name,
981 format: &result.format,
982 characters: result.characters,
983 truncated: result.truncated,
984 text: &result.text,
985 })
986}
987
988struct ToolCall {
989 name: String,
990 arguments: Value,
991}
992
993#[derive(Deserialize)]
994#[serde(rename_all = "camelCase", deny_unknown_fields)]
995struct TaskId {
996 task_id: String,
997}
998
999#[derive(Deserialize)]
1000#[serde(rename_all = "camelCase", deny_unknown_fields)]
1001struct CategoryId {
1002 category_id: String,
1003}
1004
1005#[derive(Deserialize)]
1006#[serde(rename_all = "camelCase", deny_unknown_fields)]
1007struct CategoryCall {
1008 category_id: String,
1009 #[serde(default)]
1010 offset: u64,
1011 #[serde(default = "task_page_limit")]
1012 limit: u32,
1013}
1014
1015#[derive(Deserialize)]
1016#[serde(deny_unknown_fields)]
1017struct EmptyCall {}
1018
1019fn task_page_limit() -> u32 {
1020 50
1021}
1022
1023struct RecordedToolInvocation {
1024 invocation_id: String,
1025 tool_instance: String,
1026 tool_name: String,
1027}
1028
1029fn record_tool_completion_event(
1030 journal: &mut HistorySession,
1031 invocation: Option<&RecordedToolInvocation>,
1032 outcome: Value,
1033) -> anyhow::Result<EventId> {
1034 let (tool_instance, tool_name, invocation_id) = invocation
1035 .map(|invocation| {
1036 (
1037 invocation.tool_instance.clone(),
1038 invocation.tool_name.clone(),
1039 Some(invocation.invocation_id.clone()),
1040 )
1041 })
1042 .unwrap_or_else(|| ("call_ktool".into(), "call_ktool".into(), None));
1043 journal.record(
1044 now(),
1045 EventKind::ToolCompleted {
1046 tool_instance,
1047 tool_name,
1048 outcome,
1049 invocation_id,
1050 },
1051 )
1052}
1053
1054fn ensure_tool_result_box(
1055 journal: &mut HistorySession,
1056 invocation: Option<&RecordedToolInvocation>,
1057 text: &str,
1058 ok: bool,
1059) -> anyhow::Result<String> {
1060 let Some(invocation) = invocation else {
1061 journal.create_box(
1062 now(),
1063 "Kennedy tool result",
1064 BoxOwner::Controller,
1065 BoxContent::text(text),
1066 )?;
1067 return Ok(text.to_owned());
1068 };
1069 let matches = journal
1070 .state()
1071 .boxes
1072 .values()
1073 .filter(|state| {
1074 matches!(state.owner, BoxOwner::Controller)
1075 && state
1076 .canonical
1077 .content
1078 .metadata
1079 .get("toolInvocationId")
1080 .and_then(Value::as_str)
1081 == Some(invocation.invocation_id.as_str())
1082 })
1083 .map(|state| {
1084 (
1085 state.id,
1086 state.canonical.content.text.clone(),
1087 state
1088 .canonical
1089 .content
1090 .metadata
1091 .get("toolResultOk")
1092 .and_then(Value::as_bool),
1093 )
1094 })
1095 .collect::<Vec<_>>();
1096 anyhow::ensure!(
1097 matches.len() <= 1,
1098 "tool invocation {} has duplicate durable result boxes",
1099 invocation.invocation_id
1100 );
1101 if let Some((_box_id, stored_text, stored_ok)) = matches.into_iter().next() {
1102 anyhow::ensure!(
1103 stored_ok == Some(ok),
1104 "tool invocation {} has a result box with a conflicting outcome",
1105 invocation.invocation_id
1106 );
1107 return Ok(stored_text);
1108 }
1109
1110 let mut content = BoxContent::text(text);
1111 content.metadata = json!({
1112 "toolInvocationId":invocation.invocation_id,
1113 "toolInstance":invocation.tool_instance,
1114 "toolName":invocation.tool_name,
1115 "toolResultOk":ok,
1116 });
1117 journal.create_box(now(), "Kennedy tool result", BoxOwner::Controller, content)?;
1118 Ok(text.to_owned())
1119}
1120
1121fn complete_launch_reconciliation(
1122 journal: &mut HistorySession,
1123 invocation: &RecordedToolInvocation,
1124 result: Result<kcode_session_history::SessionLaunch, kcode_session_history::Error>,
1125) -> anyhow::Result<()> {
1126 if completed_invocation_ids(journal).contains(&invocation.invocation_id) {
1127 return Ok(());
1128 }
1129 let (ok, text) = match result {
1130 Ok(launch) => (
1131 true,
1132 launch_success_json(&launch.session_id, &launch.command_id)?,
1133 ),
1134 Err(error)
1135 if matches!(
1136 error.kind,
1137 HistoryErrorKind::InvalidInput | HistoryErrorKind::Conflict
1138 ) =>
1139 {
1140 (false, format!("LaunchSession failed: {}", error.message))
1141 }
1142 Err(error) => {
1143 anyhow::bail!(
1144 "LaunchSession reconciliation remains unresolved ({}): {}",
1145 error.kind.code(),
1146 error.message
1147 );
1148 }
1149 };
1150 let text = ensure_tool_result_box(journal, Some(invocation), &text, ok)?;
1151 record_tool_completion_event(journal, Some(invocation), json!({"ok":ok,"result":text}))?;
1152 Ok(())
1153}
1154
1155struct PendingFreeformWrite {
1156 request: FreeformWrite,
1157 call_box_id: BoxId,
1158}
1159
1160struct ToolOutcome {
1161 text: String,
1162 store_result: bool,
1163 ok: bool,
1164 end_session: bool,
1165 freeform_write: Option<FreeformWrite>,
1166 managed_source_snapshot: Option<SourceSnapshot>,
1167 exact_result: bool,
1168}
1169
1170fn result_displays_snapshot(result: &str, snapshot: &SourceSnapshot) -> bool {
1171 result == snapshot.text
1172}
1173
1174fn subagent_managed_write_fits(
1175 context: &SubagentContext,
1176 call: &ToolCall,
1177 budget: &kcode_agent_runtime::ContextBudget,
1178) -> bool {
1179 let Some(snapshot) = proposed_write_snapshot(&call.name, &call.arguments) else {
1180 return true;
1181 };
1182 let state = context.source_state(&snapshot);
1183 budget.fits_state(state.key, state.text)
1184}
1185
1186struct KennedySubagentHost<'a> {
1187 session: &'a mut Session,
1188 context: SubagentContext,
1189 captures: HashMap<String, FreeformWrite>,
1190}
1191
1192struct KennedySessionHost<'a, C> {
1193 session: &'a mut Session,
1194 checkpoint: &'a mut C,
1195 state: &'a mut PrimaryTurnState,
1196}
1197
1198struct PreparedCacheObservation {
1199 cacheable_prefix_bytes: u64,
1200 expectation: CacheExpectation,
1201 material_fingerprint: String,
1202 projection_hash: String,
1203 logical_input: String,
1204 provider_input_hash: String,
1205 provider_input_bytes: u64,
1206 thread_action: String,
1207 thread_reset_reason: Option<String>,
1208 estimated_input_tokens: u64,
1209 raw_estimated_input_tokens: u64,
1210 provider: String,
1211 model: String,
1212}
1213
1214fn is_kweb_mutation(name: &str) -> bool {
1215 matches!(
1216 name,
1217 "ConnectNodes" | "ConsolidateFanout" | "SetFixedConnection" | "CreateNode" | "UpdateNode"
1218 )
1219}
1220
1221fn subagent_unavailable_reason(name: &str) -> Option<&'static str> {
1222 match name {
1223 LAUNCH_SESSION_TOOL => Some(
1224 "LaunchSession is unavailable inside a subagent. Only an authorized genuine parent user turn may launch a session.",
1225 ),
1226 "RunSubagent" => {
1227 Some("RunSubagent is unavailable inside a subagent. Only Kennedy may launch subagents.")
1228 }
1229 "EndSession" => Some(
1230 "EndSession is unavailable inside a subagent. A child cannot control the parent session lifecycle.",
1231 ),
1232 "DehydrateBoxes" | "SummarizeBox" | "HydrateBox" | "BoxesIntoObjects" => {
1233 Some("Parent box controls are unavailable inside a box-free subagent context.")
1234 }
1235 _ => None,
1236 }
1237}
1238
1239fn ensure_plan_node_known(plan: &KwebPlan, context: &KwebContext, id: &str) -> anyhow::Result<()> {
1240 if id.starts_with("pending:") {
1241 anyhow::ensure!(
1242 plan.contains_pending(id),
1243 "pending node {id} is not part of this session"
1244 );
1245 } else {
1246 canonical_id(id)?;
1247 anyhow::ensure!(
1248 context.contains_full_node(id),
1249 "node {id} is not loaded; call LoadNodes first"
1250 );
1251 }
1252 Ok(())
1253}
1254
1255fn kweb_mutation(
1256 decoded: DecodedTool,
1257 context: &KwebContext,
1258 plan: &KwebPlan,
1259 journal: &mut HistorySession,
1260) -> anyhow::Result<KwebMutation> {
1261 Ok(match decoded {
1262 DecodedTool::ConnectNodes(nodes) => KwebMutation::ConnectNodes(nodes),
1263 DecodedTool::ConsolidateFanout {
1264 parent,
1265 fanout,
1266 aggregator,
1267 } => KwebMutation::ConsolidateFanout {
1268 parent,
1269 fanout,
1270 aggregator,
1271 },
1272 DecodedTool::SetFixedConnection {
1273 parent,
1274 child,
1275 slot,
1276 } => KwebMutation::SetFixedConnection {
1277 parent,
1278 child,
1279 slot,
1280 },
1281 DecodedTool::CreateNode {
1282 parents,
1283 owner,
1284 short_name,
1285 short_description,
1286 long_description,
1287 } => {
1288 for id in parents.iter().chain(std::iter::once(&owner)) {
1289 if id != "self" && id != "unowned" {
1290 ensure_plan_node_known(plan, context, id)?;
1291 }
1292 }
1293 KwebMutation::CreateNode {
1294 pending_id: journal.allocate_pending_node(now())?.to_string(),
1295 parents,
1296 owner,
1297 short_name,
1298 short_description,
1299 long_description,
1300 }
1301 }
1302 DecodedTool::UpdateNode {
1303 id,
1304 owner,
1305 short_name,
1306 short_description,
1307 long_description,
1308 } => KwebMutation::UpdateNode {
1309 id,
1310 owner,
1311 short_name,
1312 short_description,
1313 long_description,
1314 },
1315 _ => anyhow::bail!("decoded contract did not match a Kweb mutation"),
1316 })
1317}
1318
1319fn connect_nodes_result_with_counts(
1320 result: String,
1321 plan: &KwebPlan,
1322 ids: &[String],
1323) -> anyhow::Result<String> {
1324 let (updates, creates) = plan.context_projection();
1325 let mut seen = BTreeSet::new();
1326 let mut counts = Vec::new();
1327 for id in ids {
1328 if !seen.insert(id.as_str()) {
1329 continue;
1330 }
1331 let count = updates
1332 .get(id)
1333 .map(|node| node.recent_connections.len())
1334 .or_else(|| {
1335 creates
1336 .iter()
1337 .find(|create| create.pending_id == id.as_str())
1338 .map(|create| create.data.recent_connections.len())
1339 })
1340 .with_context(|| format!("ConnectNodes did not stage touched node {id}"))?;
1341 counts.push(format!("{id}: {count}"));
1342 }
1343 Ok(format!(
1344 "{result}\nPost-call recent connection counts: {}.",
1345 counts.join(", ")
1346 ))
1347}
1348
1349fn execute_kweb_mutation(
1350 name: &str,
1351 decoded: DecodedTool,
1352 context: &KwebContext,
1353 plan: &mut KwebPlan,
1354 journal: &mut HistorySession,
1355) -> anyhow::Result<(String, Vec<String>)> {
1356 let mutation = kweb_mutation(decoded, context, plan, journal)
1357 .with_context(|| format!("decoded contract for {name} did not match its Kweb mutation"))?;
1358 let connect_nodes = match &mutation {
1359 KwebMutation::ConnectNodes(ids) => Some(ids.clone()),
1360 _ => None,
1361 };
1362 let referenced = referenced_pending_nodes(&mutation);
1363 let mut result = plan.apply(context, mutation)?;
1364 if let Some(ids) = connect_nodes {
1365 result = connect_nodes_result_with_counts(result, plan, &ids)?;
1366 }
1367 Ok((result, referenced))
1368}
1369
1370fn inference_error_receipt(
1371 error: &anyhow::Error,
1372) -> Option<kcode_intelligence_router::UsageReceipt> {
1373 error.chain().find_map(|cause| {
1374 cause
1375 .downcast_ref::<kcode_intelligence_router::Error>()
1376 .and_then(|error| error.receipt().cloned())
1377 })
1378}
1379
1380async fn record_unavailable_inference<C, F>(
1381 host: &mut KennedySessionHost<'_, C>,
1382 inference: &mut kcode_agent_runtime::SessionInference,
1383 round: u64,
1384) -> anyhow::Result<()>
1385where
1386 C: FnMut(Value) -> F + Send,
1387 F: Future<Output = anyhow::Result<()>> + Send,
1388{
1389 let receipt = inference.finish_unavailable()?;
1390 host.record(kcode_agent_runtime::SessionEvent::ProviderReceipt {
1391 round,
1392 usage: None,
1393 receipt,
1394 continuation: None,
1395 })
1396 .await
1397}
1398
1399impl Session {
1400 pub fn mark_previous_ingress_attempt_timed_out(&mut self) {
1401 if matches!(self.mode, AgentMode::Ingress { .. })
1402 && !self.previous_ingress_attempt_timed_out
1403 {
1404 self.invalidate_active_stepped_turn();
1405 self.previous_ingress_attempt_timed_out = true;
1406 }
1407 }
1408
1409 fn ingress_time_remaining(&mut self) -> anyhow::Result<Option<u64>> {
1410 if !matches!(self.mode, AgentMode::Ingress { .. }) {
1411 return Ok(None);
1412 }
1413 ingress_time_remaining_at(&mut self.ingress_deadline, Instant::now()).map(Some)
1414 }
1415
1416 fn runtime_budget(&self) -> RuntimeBudget {
1417 let Some(provider_deadline) = self.active_provider_deadline else {
1418 return RuntimeBudget::default();
1419 };
1420 let mut time_limits = vec![TimeBudget {
1421 kind: TimeBudgetKind::ProviderCall,
1422 remaining: remaining_until(provider_deadline),
1423 }];
1424 if matches!(self.mode, AgentMode::FreeTime)
1425 && let Some(work_deadline) = deadline(&self.free_time)
1426 {
1427 time_limits.push(TimeBudget {
1428 kind: TimeBudgetKind::SelfTimeWork,
1429 remaining: remaining_until(work_deadline),
1430 });
1431 }
1432 if let Some(outer) = self.active_turn_deadline {
1433 time_limits.push(TimeBudget {
1434 kind: match outer.kind {
1435 TurnDeadlineKind::Telegram => TimeBudgetKind::TelegramTurn,
1436 TurnDeadlineKind::SelfTimeHardStop => TimeBudgetKind::SelfTimeHardStop,
1437 },
1438 remaining: remaining_until(outer.at),
1439 });
1440 }
1441 RuntimeBudget {
1442 rounds: Some(RoundBudget {
1443 used: self.rounds_used,
1444 limit: kcode_agent_runtime::DEFAULT_ROUND_LIMIT,
1445 }),
1446 time_limits,
1447 }
1448 }
1449
1450 fn projection(&self) -> ContextProjection {
1451 self.journal
1452 .state()
1453 .projection_with_footer_lines(&self.runtime_budget().footer_lines())
1454 }
1455
1456 fn provider_material_fingerprint(&self, tool_description: &str) -> String {
1457 let material = json!({
1458 "model":self.runtime.model,
1459 "reasoningEffort":self.runtime.reasoning_effort,
1460 "tool":"call_ktool",
1461 "toolDescription":tool_description,
1462 });
1463 hex::encode(Sha256::digest(
1464 serde_json::to_vec(&material).expect("provider material always serializes"),
1465 ))
1466 }
1467
1468 fn begin_provider_call_budget(&mut self, timeout: Option<Duration>) {
1469 self.active_provider_deadline = timeout.and_then(|timeout| {
1470 chrono::Duration::from_std(timeout)
1471 .ok()
1472 .map(|timeout| Utc::now() + timeout)
1473 });
1474 }
1475
1476 fn clear_turn_deadlines(&mut self) {
1477 self.active_provider_deadline = None;
1478 self.active_turn_deadline = None;
1479 }
1480
1481 fn invalidate_active_stepped_turn(&mut self) {
1482 self.turn_lease_slot = None;
1483 self.clear_turn_deadlines();
1484 }
1485
1486 fn synchronize_provider_known_events(&mut self) {
1487 if let Some(affinity) = self.provider_affinity.as_mut()
1488 && let Some(event) = self.journal.state().events.last()
1489 {
1490 affinity.synchronized_event_id = event.id;
1491 }
1492 }
1493
1494 pub async fn new(
1495 api: Service,
1496 system_prompt: String,
1497 subagent_codex_prompt: String,
1498 runtime: RuntimeModel,
1499 started_at: String,
1500 mut options: SessionOptions,
1501 restored: Option<&Value>,
1502 ) -> anyhow::Result<Self> {
1503 if let Some(state) = restored {
1504 restore_session_type(&mut options, state);
1505 options.channel = state.get("channel").cloned().unwrap_or(options.channel);
1506 options.free_time = state.get("freeTime").cloned().unwrap_or(options.free_time);
1507 options.orchestration = state
1508 .get("orchestration")
1509 .cloned()
1510 .unwrap_or(options.orchestration);
1511 }
1512 if options.group_context.is_null() {
1513 options.group_context = options
1514 .channel
1515 .get("groupContext")
1516 .cloned()
1517 .unwrap_or(Value::Null);
1518 }
1519 options
1520 .reference_root_node_ids
1521 .retain(|id| !options.root_node_ids.contains(id));
1522 options.reference_root_node_ids.sort();
1523 options.reference_root_node_ids.dedup();
1524
1525 DateTime::parse_from_rfc3339(&started_at).context("session start timestamp is invalid")?;
1526 if let Some(restored_started_at) = restored
1527 .and_then(|state| state.get("startedAt"))
1528 .and_then(Value::as_str)
1529 {
1530 anyhow::ensure!(
1531 restored_started_at == started_at,
1532 "restored session start timestamp changed"
1533 );
1534 }
1535 let rust_lib_session_id = restored
1536 .and_then(|state| state.get("rustLibSessionId"))
1537 .and_then(Value::as_str)
1538 .map(str::to_owned)
1539 .or(options.rust_lib_session_id.clone())
1540 .unwrap_or_else(|| format!("kennedy:{}", Uuid::new_v4()));
1541 let history_session_id = restored
1542 .and_then(|state| state.get("sessionId"))
1543 .and_then(Value::as_str)
1544 .map(str::to_owned);
1545 let source_session_type = options.source_session_type.clone().or_else(|| {
1546 restored
1547 .and_then(|state| state.get("sourceSessionType"))
1548 .and_then(Value::as_str)
1549 .map(str::to_owned)
1550 });
1551 let session_id = history_session_id
1552 .clone()
1553 .unwrap_or_else(|| Uuid::new_v4().to_string());
1554 let metadata = SessionMetadata {
1555 session_id: session_id.clone(),
1556 kind: session_kind(&options.session_type, &options.mode),
1557 created_at: started_at.clone(),
1558 effective_context_tokens: runtime.context_window_tokens,
1559 channel: options.channel.clone(),
1560 };
1561 let mut journal = if history_session_id.is_some() {
1562 api.history_session(metadata, &runtime.model)
1563 .with_context(|| {
1564 format!(
1565 "opening authoritative session {session_id} (legacy snapshots are intentionally unsupported)"
1566 )
1567 })?
1568 } else {
1569 api.create_history_session(NewSession {
1570 kind: metadata.kind,
1571 created_at: metadata.created_at,
1572 effective_context_tokens: metadata.effective_context_tokens,
1573 channel: metadata.channel,
1574 })?
1575 };
1576 let checkpoint_event_count = restored
1577 .and_then(|state| state.get("eventCount"))
1578 .and_then(Value::as_u64)
1579 .map(|count| usize::try_from(count).context("checkpoint event count is too large"))
1580 .transpose()?
1581 .unwrap_or(journal.state().events.len());
1582 anyhow::ensure!(
1583 checkpoint_event_count <= journal.state().events.len(),
1584 "checkpoint event count is ahead of the durable journal"
1585 );
1586 let fresh_ingress_attempt = matches!(options.mode, AgentMode::Ingress { .. })
1587 && !journal.is_sealed()
1588 && journal.state().history_ingress_started;
1589 if fresh_ingress_attempt {
1590 journal.reset_history_ingress_attempt(now())?;
1591 }
1592 let launch_context_node_ids = restored
1593 .and_then(|state| state.get("launchContextNodeIds"))
1594 .cloned()
1595 .map(serde_json::from_value::<Vec<String>>)
1596 .transpose()
1597 .context("restored launch context node IDs are invalid")?
1598 .unwrap_or_default();
1599 validate_canonical_distinct_ids(&launch_context_node_ids, "context node")?;
1600 let mut context = KwebContext::with_fixed_connections(
1601 options.root_node_ids.clone(),
1602 api.loads_fixed_connections(),
1603 )
1604 .map_err(anyhow::Error::new)?;
1605 restore_kweb_context(&journal, &mut context)?;
1606 let plan = if fresh_ingress_attempt {
1607 KwebPlan::default()
1608 } else {
1609 KwebPlan::restore(
1610 restored.and_then(|state| state.get("kwebPlan")),
1611 journal_kweb_plan(&journal),
1612 )?
1613 };
1614 let transcript = kcode_kennedy_session_ingress::transcript_from_journal(&journal);
1615 let (pending_turn, mut pending_external_event_id) =
1616 restore_pending_turn(restored, &transcript);
1617 let mut rounds_used = (!fresh_ingress_attempt)
1618 .then(|| {
1619 restored
1620 .and_then(|state| state.get("roundsUsed"))
1621 .and_then(Value::as_u64)
1622 })
1623 .flatten()
1624 .unwrap_or_default();
1625
1626 let needs_initialization = !journal
1627 .state()
1628 .boxes
1629 .values()
1630 .any(|state| matches!(state.owner, BoxOwner::System));
1631 let commit_receipt = restore_commit_receipt(restored)?;
1632 let commit_author = restored
1633 .and_then(|state| state.get("commitAuthor"))
1634 .and_then(Value::as_str)
1635 .map(str::to_owned)
1636 .unwrap_or_else(|| runtime.attribution());
1637 let provider_affinity = restore_provider_affinity(restored, fresh_ingress_attempt)?;
1638 let next_thread_reset_reason = (!fresh_ingress_attempt)
1639 .then(|| {
1640 restored
1641 .and_then(|state| state.get("nextThreadResetReason"))
1642 .and_then(Value::as_str)
1643 .map(str::to_owned)
1644 })
1645 .flatten();
1646 if let Some(receipt) = &commit_receipt {
1647 journal.mark_completed(receipt.session_object_id.to_string());
1648 }
1649 let completed =
1650 journal.state().completed_session_object.is_some() || commit_receipt.is_some();
1651 let launch_provenance = restored
1652 .and_then(|state| state.get("launchProvenance"))
1653 .cloned()
1654 .unwrap_or(Value::Null);
1655 let mut launch_bootstrap_pending = restored
1656 .and_then(|state| state.get("orchestration"))
1657 .and_then(|value| value.get("launchBootstrapPending"))
1658 .and_then(Value::as_bool)
1659 .unwrap_or(!launch_provenance.is_null());
1660 let mut launch_user_turn_id = restored
1661 .and_then(|state| state.get("launchUserTurnId"))
1662 .filter(|value| !value.is_null())
1663 .cloned()
1664 .map(serde_json::from_value::<EventId>)
1665 .transpose()
1666 .context("restored launch user-turn ID is invalid")?;
1667 if let Some(id) = launch_user_turn_id {
1668 validate_user_turn_id(&journal, id)?;
1669 }
1670 let recovered_admissions =
1671 if pending_turn && checkpoint_event_count < journal.state().events.len() {
1672 recovered_user_turns(&journal, &journal.state().events[checkpoint_event_count..])?
1673 } else {
1674 Vec::new()
1675 };
1676 replay_recovered_admissions(
1677 recovered_admissions,
1678 RecoveredAdmissionReplayState {
1679 pending_turn: &pending_turn,
1680 launch_provenance: &launch_provenance,
1681 orchestration: &mut options.orchestration,
1682 launch_bootstrap_pending: &mut launch_bootstrap_pending,
1683 launch_user_turn_id: &mut launch_user_turn_id,
1684 rounds_used: &mut rounds_used,
1685 pending_external_event_id: &mut pending_external_event_id,
1686 },
1687 );
1688 let launch_intents = restored
1689 .and_then(|state| state.get("launchIntents"))
1690 .cloned()
1691 .map(serde_json::from_value::<Vec<LaunchIntent>>)
1692 .transpose()
1693 .context("restored launch intents are invalid")?
1694 .unwrap_or_default();
1695 let mut session = Self {
1696 api,
1697 subagent_codex_prompt,
1698 runtime,
1699 journal,
1700 plan,
1701 session_type: options.session_type,
1702 channel: options.channel,
1703 free_time: options.free_time,
1704 orchestration: options.orchestration,
1705 provenance_id: options.provenance_id,
1706 rust_lib_session_id,
1707 root_node_ids: options.root_node_ids,
1708 reference_root_node_ids: options.reference_root_node_ids,
1709 started_at,
1710 transcript,
1711 pending_turn,
1712 pending_external_event_id,
1713 completed,
1714 rounds_used,
1715 commit_receipt,
1716 commit_author,
1717 mode: options.mode,
1718 source_session_type,
1719 group_context: options.group_context,
1720 context,
1721 free_time_end_reason: None,
1722 fatal_persistence_error: None,
1723 active_provider_deadline: None,
1724 active_turn_deadline: None,
1725 provider_affinity,
1726 next_thread_reset_reason,
1727 ingress_deadline: None,
1728 previous_ingress_attempt_timed_out: false,
1729 launch_provenance,
1730 launch_context_node_ids,
1731 launch_user_turn_id,
1732 launch_intents,
1733 launch_bootstrap_pending,
1734 turn_lease_slot: None,
1735 };
1736 session.validate_launch_intents()?;
1737
1738 if session.journal.is_sealed() {
1739 session.provider_affinity = None;
1740 session.next_thread_reset_reason = None;
1741 anyhow::ensure!(
1742 !matches!(session.mode, AgentMode::Conversation),
1743 "a read-only conversation has an unexpectedly sealed session log"
1744 );
1745 if session.commit_receipt.is_none() {
1746 session.finalize_kweb_session()?;
1747 }
1748 session.completed = true;
1749 return Ok(session);
1750 }
1751
1752 session.reconcile_launch_intents().await?;
1753 session.prune_launch_intents();
1754 session.repair_unfinished_tools()?;
1755
1756 if needs_initialization {
1757 session.journal.create_box(
1758 now(),
1759 "Kennedy system prompt",
1760 BoxOwner::System,
1761 BoxContent::text(&system_prompt),
1762 )?;
1763 if session.session_type == "telegram-group" && !session.group_context.is_null() {
1764 session.journal.create_box(
1765 now(),
1766 "Telegram group context",
1767 BoxOwner::Controller,
1768 BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
1769 &session.group_context,
1770 )),
1771 )?;
1772 }
1773 let identifiers = session.initial_context_identifiers();
1774 let invocation =
1775 session.record_tool_invocation("LoadNodes", json!({"identifiers":&identifiers}))?;
1776 let result =
1777 load_durable_batch(session.api.kmap(), &mut session.context, &identifiers)?;
1778 for id in &session.launch_context_node_ids {
1779 anyhow::ensure!(
1780 session.context.contains_full_node(id),
1781 "launched child context node {id} is unavailable"
1782 );
1783 }
1784 session.sync_cache_safe_kweb_boxes()?;
1785 session.record_tool_completion(
1786 Some(&invocation),
1787 json!({"ok":true,"automatic":true,"identifiers":identifiers,"result":result}),
1788 )?;
1789 } else {
1790 if !session.launch_context_node_ids.is_empty() {
1791 let identifiers = session.initial_context_identifiers();
1792 load_durable_batch(session.api.kmap(), &mut session.context, &identifiers)?;
1793 for id in &session.launch_context_node_ids {
1794 anyhow::ensure!(
1795 session.context.contains_full_node(id),
1796 "launched child context node {id} is unavailable"
1797 );
1798 }
1799 }
1800 session.sync_cache_safe_kweb_boxes()?;
1801 }
1802 if fresh_ingress_attempt {
1803 session.revalidate_loaded_nodes().await?;
1804 session.pending_turn = true;
1805 }
1806 if matches!(session.mode, AgentMode::Ingress { .. })
1807 && !session.completed
1808 && !session.journal.state().history_ingress_started
1809 {
1810 session.prepare_history_ingress(&system_prompt).await?;
1811 }
1812 Ok(session)
1813 }
1814
1815 fn initial_context_identifiers(&self) -> Vec<String> {
1816 let mut identifiers = self.root_node_ids.clone();
1817 if !self.launch_context_node_ids.is_empty() {
1818 for id in self
1819 .reference_root_node_ids
1820 .iter()
1821 .chain(self.launch_context_node_ids.iter())
1822 {
1823 if !identifiers.contains(id) {
1824 identifiers.push(id.clone());
1825 }
1826 }
1827 }
1828 identifiers
1829 }
1830
1831 fn launch_session_authorized(&self) -> bool {
1832 self.pending_turn
1833 && self.launch_user_turn_id.is_some()
1834 && !self.launch_bootstrap_pending
1835 && matches!(self.mode, AgentMode::Conversation)
1836 && matches!(
1837 self.session_type.as_str(),
1838 "conversation" | "telegram" | "telegram-group"
1839 )
1840 }
1841
1842 fn validate_launch_intents(&self) -> anyhow::Result<()> {
1843 let mut seen = BTreeSet::new();
1844 for intent in &self.launch_intents {
1845 Uuid::parse_str(&intent.invocation_id).with_context(|| {
1846 format!("launch intent {} has an invalid UUID", intent.invocation_id)
1847 })?;
1848 anyhow::ensure!(
1849 seen.insert(intent.invocation_id.as_str()),
1850 "duplicate launch intent {}",
1851 intent.invocation_id
1852 );
1853 validate_user_turn_id(&self.journal, intent.user_turn_id)?;
1854 DateTime::parse_from_rfc3339(&intent.started_at)
1855 .context("launch intent timestamp is invalid")?;
1856 anyhow::ensure!(
1857 intent.parent_session_id == self.journal.state().metadata.session_id,
1858 "launch intent parent session changed"
1859 );
1860 anyhow::ensure!(
1861 intent.effective_context_tokens > 0,
1862 "launch intent effective context size is invalid"
1863 );
1864 validate_canonical_distinct_ids(&intent.root_node_ids, "root node")?;
1865 validate_canonical_distinct_ids(
1866 &intent.reference_root_node_ids,
1867 "reference root node",
1868 )?;
1869 validate_canonical_distinct_ids(&intent.context_node_ids, "context node")?;
1870 let _ = invocation_arguments(&self.journal, &intent.invocation_id)?;
1871 }
1872 Ok(())
1873 }
1874
1875 fn prune_launch_intents(&mut self) {
1876 let completed = completed_invocation_ids(&self.journal);
1877 self.launch_intents =
1878 pruned_launch_intents(&self.launch_intents, self.launch_user_turn_id, &completed);
1879 }
1880
1881 fn unfinished_launch_intents(&self) -> Vec<&LaunchIntent> {
1882 let completed = completed_invocation_ids(&self.journal);
1883 self.launch_intents
1884 .iter()
1885 .filter(|intent| !completed.contains(&intent.invocation_id))
1886 .collect()
1887 }
1888
1889 fn ensure_no_unfinished_launch_intents(&self) -> anyhow::Result<()> {
1890 let unfinished = self.unfinished_launch_intents();
1891 anyhow::ensure!(
1892 unfinished.is_empty(),
1893 "session has an unfinished intent-backed LaunchSession invocation"
1894 );
1895 Ok(())
1896 }
1897
1898 fn repair_unfinished_tools(&mut self) -> anyhow::Result<()> {
1899 self.ensure_no_unfinished_launch_intents()?;
1900 self.journal.repair_unfinished_tools(now())?;
1901 Ok(())
1902 }
1903
1904 fn prepare_launch_intent(
1905 &mut self,
1906 invocation: &RecordedToolInvocation,
1907 arguments: &LaunchSessionArguments,
1908 ) -> anyhow::Result<LaunchIntent> {
1909 anyhow::ensure!(
1910 self.launch_session_authorized(),
1911 "LaunchSession is unavailable without a genuine current user turn in an eligible conversation"
1912 );
1913 validate_loaded_launch_context(&self.context, &arguments.context_node_ids)?;
1914 if let Some(existing) = self
1915 .launch_intents
1916 .iter()
1917 .find(|intent| intent.invocation_id == invocation.invocation_id)
1918 {
1919 return Ok(existing.clone());
1920 }
1921 let user_turn_id = self
1922 .launch_user_turn_id
1923 .context("LaunchSession user-turn authority is missing")?;
1924 let current_count = self
1925 .launch_intents
1926 .iter()
1927 .filter(|intent| intent.user_turn_id == user_turn_id)
1928 .count();
1929 anyhow::ensure!(
1930 current_count < MAX_LAUNCH_INTENTS_PER_USER_TURN,
1931 "LaunchSession permits at most ten new sessions per genuine user turn"
1932 );
1933 let intent = LaunchIntent {
1934 invocation_id: invocation.invocation_id.clone(),
1935 user_turn_id,
1936 started_at: now(),
1937 parent_session_id: self.journal.state().metadata.session_id.clone(),
1938 effective_context_tokens: self.runtime.context_window_tokens,
1939 root_node_ids: self.root_node_ids.clone(),
1940 reference_root_node_ids: self.reference_root_node_ids.clone(),
1941 context_node_ids: arguments.context_node_ids.clone(),
1942 };
1943 self.launch_intents.push(intent.clone());
1944 Ok(intent)
1945 }
1946
1947 fn launch_request(intent: &LaunchIntent, directive: &str) -> HistoryLaunchSession {
1948 let provenance = json!({
1949 "kind":"synthetic-launch-bootstrap",
1950 "denyLaunchSession":true,
1951 "parentSessionId":intent.parent_session_id,
1952 "parentInvocationId":intent.invocation_id,
1953 "parentUserTurnId":intent.user_turn_id,
1954 });
1955 HistoryLaunchSession {
1956 session_id: intent.invocation_id.clone(),
1957 started_at: intent.started_at.clone(),
1958 effective_context_tokens: intent.effective_context_tokens,
1959 channel: json!({"kind":"browser"}),
1960 state: json!({
1961 "sessionId":intent.invocation_id,
1962 "chatendMetadata":{
1963 "sessionId":intent.invocation_id,
1964 "kind":SessionKind::Conversation,
1965 "createdAt":intent.started_at,
1966 "effectiveContextTokens":intent.effective_context_tokens,
1967 "channel":{"kind":"browser"},
1968 },
1969 "sessionType":"conversation",
1970 "channel":{"kind":"browser"},
1971 "freeTime":Value::Null,
1972 "orchestration":{
1973 "owner":"backend",
1974 "status":"idle",
1975 "launchBootstrapPending":true,
1976 },
1977 "launchProvenance":provenance,
1978 "rootNodeIds":intent.root_node_ids,
1979 "referenceRootNodeIds":intent.reference_root_node_ids,
1980 "launchContextNodeIds":intent.context_node_ids,
1981 "startedAt":intent.started_at,
1982 "pendingTurn":false,
1983 "completed":false,
1984 }),
1985 initial_message: json!({
1986 "text":directive,
1987 "metadata":{
1988 "launchProvenance":provenance,
1989 },
1990 }),
1991 }
1992 }
1993
1994 async fn lower_launch(
1995 &self,
1996 intent: &LaunchIntent,
1997 arguments: &LaunchSessionArguments,
1998 ) -> Result<kcode_session_history::SessionLaunch, kcode_session_history::Error> {
1999 self.api
2000 .launch_session(Self::launch_request(intent, &arguments.directive))
2001 .await
2002 }
2003
2004 async fn reconcile_launch_intents(&mut self) -> anyhow::Result<()> {
2005 let completed = completed_invocation_ids(&self.journal);
2006 let unfinished = self
2007 .launch_intents
2008 .iter()
2009 .filter(|intent| !completed.contains(&intent.invocation_id))
2010 .cloned()
2011 .collect::<Vec<_>>();
2012 for intent in unfinished {
2013 let raw = invocation_arguments(&self.journal, &intent.invocation_id)?.clone();
2014 let arguments = decode_launch_session_arguments(&raw)?;
2015 let invocation = RecordedToolInvocation {
2016 invocation_id: intent.invocation_id.clone(),
2017 tool_instance: tool_instance_for_invocation(
2018 LAUNCH_SESSION_TOOL,
2019 &intent.invocation_id,
2020 ),
2021 tool_name: LAUNCH_SESSION_TOOL.into(),
2022 };
2023 let result = self.lower_launch(&intent, &arguments).await;
2024 complete_launch_reconciliation(&mut self.journal, &invocation, result)?;
2025 }
2026 Ok(())
2027 }
2028
2029 async fn prepare_history_ingress(&mut self, prompt: &str) -> anyhow::Result<()> {
2030 let cost_at_ingress = self.projection().status;
2031 if !self.journal.state().source_terminated {
2032 self.journal.record(
2033 now(),
2034 EventKind::SourceTerminated {
2035 reason: "history_ingress".into(),
2036 },
2037 )?;
2038 }
2039 let system_box = self
2040 .journal
2041 .state()
2042 .boxes
2043 .values()
2044 .find(|state| matches!(state.owner, BoxOwner::System))
2045 .map(|state| state.id)
2046 .context("session has no system-prompt box")?;
2047 let recorded_at = now();
2048 self.journal
2049 .update_box(recorded_at.clone(), system_box, BoxContent::text(prompt))?;
2050 self.journal.rehydrate_box(recorded_at, system_box)?;
2051 let ingress_kind = session_kind(&self.session_type, &self.mode);
2052 if self.journal.state().metadata.effective_context_tokens
2053 != self.runtime.context_window_tokens
2054 || self.journal.state().metadata.kind != ingress_kind
2055 {
2056 self.journal
2057 .configure_context(ingress_kind, self.runtime.context_window_tokens);
2058 }
2059 self.journal.create_box(
2060 now(),
2061 "Session cost at ingress",
2062 BoxOwner::Controller,
2063 BoxContent::text(cost_summary(
2064 "session cost before history ingress",
2065 cost_at_ingress.estimated_cost_usd_nanos,
2066 cost_at_ingress.unpriced_provider_calls,
2067 )),
2068 )?;
2069 self.revalidate_loaded_nodes().await?;
2070 match kcode_history_ingress_context::prepare(&mut self.journal, now())? {
2071 HistoryIngressContextOutcome::Ready => {}
2072 HistoryIngressContextOutcome::OverCapacity {
2073 estimated_tokens,
2074 target_tokens,
2075 } => {
2076 self.journal.record(
2077 now(),
2078 EventKind::Note {
2079 label: INGRESS_FORCE_COMMIT_NOTE.into(),
2080 value: json!({
2081 "reason":"fully_dehydrated_context_above_initial_target",
2082 "estimatedTokens":estimated_tokens,
2083 "initialTargetTokens":target_tokens,
2084 }),
2085 },
2086 )?;
2087 self.clear_launch_turn_authority();
2088 self.pending_turn = false;
2089 self.finalize_kweb_session()?;
2090 self.completed = true;
2091 return Ok(());
2092 }
2093 }
2094 self.journal
2095 .record(now(), EventKind::HistoryIngressStarted)?;
2096 self.pending_turn = true;
2097 Ok(())
2098 }
2099
2100 async fn revalidate_loaded_nodes(&mut self) -> anyhow::Result<()> {
2101 let direct = self.context.loaded_node_ids().to_vec();
2102 load_durable_batch(self.api.kmap(), &mut self.context, &direct)?;
2103 self.sync_cache_safe_kweb_boxes()?;
2104 Ok(())
2105 }
2106
2107 fn stage_user_input(&mut self, text: &str, metadata: &Value) -> Option<InputStage> {
2108 let recorded_at = now();
2109 let result = (|| -> anyhow::Result<Option<InputStage>> {
2110 let Some(staged) = kcode_kennedy_session_ingress::stage_user_input(
2111 &mut self.journal,
2112 text,
2113 metadata,
2114 &recorded_at,
2115 )?
2116 else {
2117 return Ok(None);
2118 };
2119 self.transcript.push(staged.transcript);
2120 self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
2121 Ok(Some(InputStage::Accepted))
2122 })();
2123 match result {
2124 Ok(stage) => stage,
2125 Err(error) => {
2126 self.fatal_persistence_error = Some(error.to_string());
2127 tracing::error!(error=%error, "Could not durably stage session input");
2128 Some(InputStage::Accepted)
2129 }
2130 }
2131 }
2132
2133 pub fn append_final_user_message(&mut self, text: &str, metadata: &Value) -> bool {
2134 let accepted = self.stage_user_input(text, metadata).is_some();
2135 if accepted {
2136 self.invalidate_active_stepped_turn();
2137 }
2138 accepted
2139 }
2140
2141 pub fn stage_source_message(
2142 &mut self,
2143 kennedy: bool,
2144 text: &str,
2145 metadata: Value,
2146 ) -> anyhow::Result<()> {
2147 let staged = kcode_kennedy_session_ingress::stage_source_input(
2148 &mut self.journal,
2149 kennedy,
2150 text,
2151 metadata,
2152 &now(),
2153 )?;
2154 self.invalidate_active_stepped_turn();
2155 self.transcript.push(staged.transcript);
2156 self.recover_context_overflow(staged.external_event_id.as_deref(), &[])?;
2157 Ok(())
2158 }
2159
2160 pub fn answer_for_external_event(&self, id: &str) -> Option<&Value> {
2161 self.transcript.iter().rev().find(|entry| {
2162 is_terminal_external_response(entry)
2163 && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
2164 })
2165 }
2166
2167 pub fn responses_for_external_event(&self, id: &str) -> Vec<&Value> {
2168 self.transcript
2169 .iter()
2170 .filter(|entry| {
2171 matches!(
2172 entry.get("role").and_then(Value::as_str),
2173 Some("kennedy" | "system")
2174 ) && entry.get("externalEventId").and_then(Value::as_str) == Some(id)
2175 })
2176 .collect()
2177 }
2178
2179 pub fn resolve_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
2180 let api = self.api.clone();
2181 kcode_kennedy_session_objects::resolve_object(
2182 &mut self.journal,
2183 object_id,
2184 move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
2185 )
2186 }
2187
2188 fn resolve_media_object(&mut self, object_id: &str) -> anyhow::Result<ResolvedObject> {
2189 let api = self.api.clone();
2190 kcode_kennedy_session_objects::resolve_media_object(
2191 &mut self.journal,
2192 object_id,
2193 MAX_MEDIA_ENRICHMENT_BYTES,
2194 move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
2195 )
2196 }
2197
2198 fn resolve_image_object(
2199 &mut self,
2200 object_id: &str,
2201 ) -> anyhow::Result<(Vec<u8>, String, String)> {
2202 let resolved = self.resolve_media_object(object_id)?;
2203 anyhow::ensure!(
2204 resolved.media_type.starts_with("image/"),
2205 "GenerateImage reference {object_id} is not an image"
2206 );
2207 Ok((resolved.bytes, resolved.file_name, resolved.media_type))
2208 }
2209
2210 fn recover_context_overflow(
2211 &mut self,
2212 external_event_id: Option<&str>,
2213 pinned_box_ids: &[BoxId],
2214 ) -> anyhow::Result<ContextRecovery> {
2215 let projection = self.projection();
2216 let target_tokens = self.journal.state().active_context_limit();
2217 if projection.estimated_tokens <= target_tokens {
2218 return Ok(ContextRecovery::NotNeeded);
2219 }
2220 let projection_hash = hex::encode(Sha256::digest(projection.render().as_bytes()));
2221 let already_irreducible = self
2222 .journal
2223 .state()
2224 .events
2225 .iter()
2226 .rev()
2227 .find_map(|event| match &event.kind {
2228 EventKind::Note { label, value } if label == "context_overflow_recovery" => {
2229 Some(value)
2230 }
2231 _ => None,
2232 })
2233 .is_some_and(|value| {
2234 value.get("irreducible").and_then(Value::as_bool) == Some(true)
2235 && value.get("limitTokens").and_then(Value::as_u64) == Some(target_tokens)
2236 && value.get("projectionHash").and_then(Value::as_str)
2237 == Some(projection_hash.as_str())
2238 });
2239 if already_irreducible {
2240 return Ok(ContextRecovery::Irreducible);
2241 }
2242
2243 let before_tokens = projection.estimated_tokens;
2244 let mut metadata = json!({
2245 "transcriptRole":"system",
2246 "contextOverflowWarning":true,
2247 "projectedTokens":before_tokens,
2248 "limitTokens":target_tokens,
2249 });
2250 if let Some(id) = external_event_id {
2251 metadata["externalEventId"] = json!(id);
2252 }
2253 let warning_box_id = self.journal.create_box(
2254 now(),
2255 CONTEXT_OVERFLOW_WARNING_BOX_NAME,
2256 BoxOwner::Controller,
2257 BoxContent {
2258 text: CONTEXT_OVERFLOW_WARNING.into(),
2259 objects: Vec::new(),
2260 metadata,
2261 },
2262 )?;
2263 let mut transcript = json!({
2264 "role":"system",
2265 "content":CONTEXT_OVERFLOW_WARNING,
2266 "contextOverflowWarning":true,
2267 });
2268 if let Some(id) = external_event_id {
2269 transcript["externalEventId"] = json!(id);
2270 }
2271 self.transcript.push(transcript);
2272
2273 let mut pins = pinned_box_ids.to_vec();
2274 if !pins.contains(&warning_box_id) {
2275 pins.push(warning_box_id);
2276 }
2277 let outcome = kcode_history_ingress_context::recover(&mut self.journal, now(), &pins)?;
2278 let (dehydrated_box_ids, estimated_tokens, target_tokens, irreducible) = match outcome {
2279 ContextRecoveryOutcome::Recovered {
2280 dehydrated_box_ids,
2281 estimated_tokens,
2282 target_tokens,
2283 } => (dehydrated_box_ids, estimated_tokens, target_tokens, false),
2284 ContextRecoveryOutcome::OverCapacity {
2285 dehydrated_box_ids,
2286 estimated_tokens,
2287 target_tokens,
2288 } => (dehydrated_box_ids, estimated_tokens, target_tokens, true),
2289 };
2290 let final_projection_hash =
2291 hex::encode(Sha256::digest(self.projection().render().as_bytes()));
2292 self.journal.record(
2293 now(),
2294 EventKind::Note {
2295 label: "context_overflow_recovery".into(),
2296 value: json!({
2297 "beforeTokens":before_tokens,
2298 "estimatedTokens":estimated_tokens,
2299 "limitTokens":target_tokens,
2300 "dehydratedBoxIds":dehydrated_box_ids,
2301 "irreducible":irreducible,
2302 "projectionHash":final_projection_hash,
2303 }),
2304 },
2305 )?;
2306 if irreducible {
2307 if matches!(self.mode, AgentMode::Ingress { .. }) {
2308 self.request_ingress_force_commit(
2309 "irreducible_context_overflow",
2310 estimated_tokens,
2311 )?;
2312 } else if !self.journal.state().source_terminated {
2313 self.journal.record(
2314 now(),
2315 EventKind::SourceTerminated {
2316 reason: "irreducible_context_overflow".into(),
2317 },
2318 )?;
2319 }
2320 Ok(ContextRecovery::Irreducible)
2321 } else {
2322 Ok(ContextRecovery::Recovered)
2323 }
2324 }
2325
2326 fn request_ingress_force_commit(
2327 &mut self,
2328 reason: &str,
2329 projected_tokens: u64,
2330 ) -> anyhow::Result<()> {
2331 if self.ingress_force_commit_requested() {
2332 return Ok(());
2333 }
2334 self.journal.record(
2335 now(),
2336 EventKind::Note {
2337 label: INGRESS_FORCE_COMMIT_NOTE.into(),
2338 value: json!({
2339 "reason":reason,
2340 "projectedTokens":projected_tokens,
2341 "limitTokens":self.journal.state().ingress_context_limit(),
2342 }),
2343 },
2344 )?;
2345 Ok(())
2346 }
2347
2348 fn ingress_force_commit_requested(&self) -> bool {
2349 self.journal
2350 .state()
2351 .current_ingress_attempt_events()
2352 .iter()
2353 .rev()
2354 .any(|event| {
2355 matches!(
2356 &event.kind,
2357 EventKind::Note { label, .. } if label == INGRESS_FORCE_COMMIT_NOTE
2358 )
2359 })
2360 }
2361
2362 pub fn requires_history_ingress(&self) -> bool {
2363 matches!(self.mode, AgentMode::Conversation) && self.journal.state().source_terminated
2364 }
2365
2366 pub fn stage_free_time_opening(&mut self) -> bool {
2367 if self.pending_turn {
2368 return false;
2369 }
2370 self.invalidate_active_stepped_turn();
2371 self.launch_user_turn_id = None;
2372 self.prune_launch_intents();
2373 let mut blocks = vec![
2374 render(RenderRequest::FreeTimeOpening {
2375 free_time: &self.free_time,
2376 })
2377 .expect("free-time opening rendering is infallible"),
2378 ];
2379 if let Some(message) = self
2380 .free_time
2381 .get("handoffMessage")
2382 .and_then(Value::as_str)
2383 .filter(|message| !message.trim().is_empty())
2384 {
2385 blocks.push(format!(
2386 "Message from the previous self-time session:\n\n{message}"
2387 ));
2388 }
2389 let Some(stage) = self.stage_user_input(&blocks.join("\n\n"), &json!({"kind":"self-time"}))
2390 else {
2391 return false;
2392 };
2393 self.pending_turn = matches!(stage, InputStage::Accepted);
2394 true
2395 }
2396
2397 pub fn stage_wakeup_opening(&mut self) -> anyhow::Result<bool> {
2398 if self.pending_turn {
2399 return Ok(false);
2400 }
2401 self.invalidate_active_stepped_turn();
2402 self.launch_user_turn_id = None;
2403 self.prune_launch_intents();
2404 let marker = self
2405 .channel
2406 .get("wakeupMarker")
2407 .and_then(Value::as_str)
2408 .context("wakeup session is missing its acquired time marker")?;
2409 let marker = DateTime::parse_from_rfc3339(marker)
2410 .context("wakeup session has an invalid acquired time marker")?
2411 .with_timezone(&Utc);
2412 let text = render(RenderRequest::WakeupOpening { marker })?;
2413 let Some(stage) = self.stage_user_input(
2414 &text,
2415 &json!({"kind":"wakeup","wakeupMarker":marker.to_rfc3339()}),
2416 ) else {
2417 return Ok(false);
2418 };
2419 self.pending_turn = matches!(stage, InputStage::Accepted);
2420 Ok(true)
2421 }
2422
2423 pub fn begin_user_turn(&mut self, text: &str, metadata: &Value) -> bool {
2424 if self.pending_turn {
2425 return false;
2426 }
2427 self.invalidate_active_stepped_turn();
2428 self.launch_user_turn_id = None;
2429 self.prune_launch_intents();
2430 let first_event = self.journal.state().events.len();
2431 let Some(stage) = self.stage_user_input(text, metadata) else {
2432 return false;
2433 };
2434 debug_assert_eq!(stage, InputStage::Accepted);
2435 let user_turn =
2436 unique_user_box_event(&self.journal, &self.journal.state().events[first_event..])
2437 .ok()
2438 .flatten();
2439 self.launch_user_turn_id = user_turn_launch_authority(
2440 user_turn,
2441 &mut self.orchestration,
2442 &mut self.launch_bootstrap_pending,
2443 );
2444 self.rounds_used = 0;
2445 self.pending_turn = true;
2446 self.pending_external_event_id = metadata
2447 .get("externalEventId")
2448 .and_then(Value::as_str)
2449 .map(str::to_owned);
2450 true
2451 }
2452
2453 fn clear_launch_turn_authority(&mut self) {
2454 self.launch_user_turn_id = None;
2455 self.prune_launch_intents();
2456 }
2457
2458 pub fn reset_exhausted_turn_rounds_for_retry(&mut self) {
2459 if matches!(self.mode, AgentMode::Conversation)
2460 && self.rounds_used >= kcode_agent_runtime::DEFAULT_ROUND_LIMIT
2461 {
2462 self.invalidate_active_stepped_turn();
2463 self.rounds_used = 0;
2464 }
2465 }
2466
2467 pub fn interrupt_current_turn(&mut self) -> anyhow::Result<()> {
2468 self.ensure_no_unfinished_launch_intents()?;
2469 self.invalidate_active_stepped_turn();
2470 self.provider_affinity = None;
2471 self.next_thread_reset_reason = Some("prior_provider_turn_interrupted".into());
2472 self.repair_unfinished_tools()?;
2473 let notice = "The user stopped this agent turn.";
2474 let mut metadata = json!({"transcriptRole":"system","userStopped":true});
2475 let mut transcript_entry = json!({
2476 "role":"system",
2477 "content":notice,
2478 "userStopped":true,
2479 });
2480 if let Some(external_event_id) = &self.pending_external_event_id {
2481 metadata["externalEventId"] = json!(external_event_id);
2482 transcript_entry["externalEventId"] = json!(external_event_id);
2483 }
2484 self.journal.create_box(
2485 now(),
2486 "Turn stopped",
2487 BoxOwner::Controller,
2488 BoxContent {
2489 text: notice.into(),
2490 objects: Vec::new(),
2491 metadata,
2492 },
2493 )?;
2494 self.transcript.push(transcript_entry);
2495 self.pending_turn = false;
2496 self.pending_external_event_id = None;
2497 self.clear_launch_turn_authority();
2498 self.orchestration =
2499 json!({"owner":"backend","status":"idle","lastOutcome":"user-stopped"});
2500 Ok(())
2501 }
2502
2503 pub fn begin_pending_turn(
2504 &mut self,
2505 operation_id: Uuid,
2506 turn_deadline: Option<TurnDeadline>,
2507 ) -> anyhow::Result<Option<SessionTurn>> {
2508 if let Some(error) = self.fatal_persistence_error.take() {
2509 anyhow::bail!("session journal write failed: {error}");
2510 }
2511 if !self.pending_turn {
2512 return Ok(None);
2513 }
2514 let user_id = self
2515 .root_node_ids
2516 .first()
2517 .context("session has no user root for intelligence accounting")?
2518 .clone();
2519 let lease = TurnLease::acquire(&mut self.turn_lease_slot)?;
2520 self.active_provider_deadline = None;
2521 self.active_turn_deadline = turn_deadline;
2522 Ok(Some(SessionTurn {
2523 lease,
2524 user_id,
2525 completed_rounds: self.rounds_used,
2526 round_limit: kcode_agent_runtime::DEFAULT_ROUND_LIMIT,
2527 state: PrimaryTurnState {
2528 accounting: None,
2529 pending_freeform_write: None,
2530 deadline_after_response: false,
2531 operation_id,
2532 prepared_cache: None,
2533 provider_synchronized_after: None,
2534 restart_fresh_reason: None,
2535 exact_tool_result: false,
2536 used_tool: false,
2537 finish_requested: false,
2538 emitted_response: false,
2539 pending_capture: None,
2540 },
2541 at_yielded_boundary: false,
2542 admission_poisoned: false,
2543 }))
2544 }
2545
2546 pub async fn admit_pending_turn<C, F>(
2547 &mut self,
2548 turn: &mut SessionTurn,
2549 admission: PendingTurnAdmission,
2550 recorded_at: &str,
2551 checkpoint: &mut C,
2552 ) -> anyhow::Result<bool>
2553 where
2554 C: FnMut(Value) -> F + Send,
2555 F: Future<Output = anyhow::Result<()>> + Send,
2556 {
2557 turn.lease.validate(&self.turn_lease_slot)?;
2558 anyhow::ensure!(self.pending_turn, "the session has no current pending turn");
2559 anyhow::ensure!(
2560 turn.at_yielded_boundary,
2561 "the stepped session turn is not at a yielded boundary"
2562 );
2563 anyhow::ensure!(
2564 !turn.admission_poisoned,
2565 "the stepped session turn admission handle is poisoned"
2566 );
2567
2568 let expected_kind = match &admission {
2569 PendingTurnAdmission::User { .. } => AdmissionKind::User,
2570 PendingTurnAdmission::Source { .. } => AdmissionKind::Source,
2571 };
2572 let result = self
2573 .admit_pending_turn_staged(turn, admission, expected_kind, recorded_at, checkpoint)
2574 .await;
2575 if result.is_err() {
2576 turn.admission_poisoned = true;
2577 }
2578 result
2579 }
2580
2581 async fn admit_pending_turn_staged<C, F>(
2582 &mut self,
2583 turn: &mut SessionTurn,
2584 admission: PendingTurnAdmission,
2585 expected_kind: AdmissionKind,
2586 recorded_at: &str,
2587 checkpoint: &mut C,
2588 ) -> anyhow::Result<bool>
2589 where
2590 C: FnMut(Value) -> F + Send,
2591 F: Future<Output = anyhow::Result<()>> + Send,
2592 {
2593 let Some(staged) = stage_turn_admission(&mut self.journal, admission, recorded_at)? else {
2594 return Ok(false);
2595 };
2596
2597 let kind = staged.kind;
2598 let external_event_id = staged.external_event_id;
2599 let user_turn_id = staged.user_turn_id;
2600 self.transcript.push(staged.transcript);
2601 self.recover_context_overflow(external_event_id.as_deref(), &[])?;
2602
2603 anyhow::ensure!(
2604 kind == expected_kind,
2605 "staged turn admission kind differs from the requested kind"
2606 );
2607 match kind {
2608 AdmissionKind::User => {
2609 let user_turn_id =
2610 user_turn_id.context("staged User admission has no user-turn ID")?;
2611 grant_marked_user_launch_authority(
2612 user_turn_id,
2613 &mut self.orchestration,
2614 &mut self.launch_bootstrap_pending,
2615 &mut self.launch_user_turn_id,
2616 );
2617 self.prune_launch_intents();
2618 self.rounds_used = 0;
2619 turn.completed_rounds = 0;
2620 self.pending_external_event_id = external_event_id;
2621 }
2622 AdmissionKind::Source => {
2623 anyhow::ensure!(
2624 user_turn_id.is_none(),
2625 "staged Source admission has a user-turn ID"
2626 );
2627 }
2628 }
2629
2630 checkpoint(self.snapshot()?).await?;
2631 Ok(true)
2632 }
2633
2634 async fn finish_stepped_turn<C, F>(
2635 &mut self,
2636 turn: SessionTurn,
2637 result: Option<String>,
2638 checkpoint: &mut C,
2639 ) -> anyhow::Result<TurnBoundary>
2640 where
2641 C: FnMut(Value) -> F + Send,
2642 F: Future<Output = anyhow::Result<()>> + Send,
2643 {
2644 turn.lease.validate(&self.turn_lease_slot)?;
2645 self.turn_lease_slot = None;
2646 self.clear_turn_deadlines();
2647 let output = match self.mode {
2648 AgentMode::Conversation => {
2649 if self.journal.state().source_terminated {
2650 self.provider_affinity = None;
2651 self.next_thread_reset_reason = None;
2652 self.pending_turn = false;
2653 self.pending_external_event_id = None;
2654 self.clear_launch_turn_authority();
2655 checkpoint(self.snapshot()?).await?;
2656 None
2657 } else {
2658 let Some(answer) = result else {
2659 if self
2660 .pending_external_event_id
2661 .as_deref()
2662 .and_then(|id| self.answer_for_external_event(id))
2663 .is_some()
2664 {
2665 self.pending_turn = false;
2666 self.pending_external_event_id = None;
2667 self.clear_launch_turn_authority();
2668 checkpoint(self.snapshot()?).await?;
2669 return Ok(TurnBoundary::Complete(None));
2670 }
2671 anyhow::bail!(
2672 "Kennedy ended a conversational turn without an assistant response"
2673 );
2674 };
2675 self.pending_turn = false;
2676 self.pending_external_event_id = None;
2677 self.clear_launch_turn_authority();
2678 checkpoint(self.snapshot()?).await?;
2679 Some(answer)
2680 }
2681 }
2682 AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. } => {
2683 self.pending_turn = false;
2684 self.pending_external_event_id = None;
2685 self.clear_launch_turn_authority();
2686 self.finalize_kweb_session()?;
2687 self.completed = true;
2688 checkpoint(self.snapshot()?).await?;
2689 None
2690 }
2691 };
2692 Ok(TurnBoundary::Complete(output))
2693 }
2694
2695 pub async fn advance_pending_turn<C, F>(
2696 &mut self,
2697 mut turn: SessionTurn,
2698 checkpoint: &mut C,
2699 ) -> anyhow::Result<TurnBoundary>
2700 where
2701 C: FnMut(Value) -> F + Send,
2702 F: Future<Output = anyhow::Result<()>> + Send,
2703 {
2704 turn.lease.validate(&self.turn_lease_slot)?;
2705 anyhow::ensure!(
2706 !turn.admission_poisoned,
2707 "the stepped session turn admission handle is poisoned"
2708 );
2709 if turn.at_yielded_boundary {
2710 turn.at_yielded_boundary = false;
2711 }
2712 let result = self.advance_pending_turn_validated(turn, checkpoint).await;
2713 if result.is_err() {
2714 self.clear_turn_deadlines();
2715 }
2716 result
2717 }
2718
2719 async fn advance_pending_turn_validated<C, F>(
2720 &mut self,
2721 mut turn: SessionTurn,
2722 checkpoint: &mut C,
2723 ) -> anyhow::Result<TurnBoundary>
2724 where
2725 C: FnMut(Value) -> F + Send,
2726 F: Future<Output = anyhow::Result<()>> + Send,
2727 {
2728 if turn.completed_rounds >= turn.round_limit {
2729 let runtime = self.api.agent_runtime();
2730 let operation_id = turn.state.operation_id;
2731 let answer = {
2732 let mut host = KennedySessionHost {
2733 session: self,
2734 checkpoint,
2735 state: &mut turn.state,
2736 };
2737 runtime
2738 .run_session(
2739 kcode_agent_runtime::SessionRunRequest {
2740 user_id: turn.user_id.clone(),
2741 operation_id,
2742 rounds_used: turn.round_limit,
2743 round_limit: turn.round_limit,
2744 },
2745 &mut host,
2746 )
2747 .await?
2748 };
2749 return self.finish_stepped_turn(turn, answer, checkpoint).await;
2750 }
2751
2752 let round = turn.completed_rounds + 1;
2753 let operation_id = turn.state.operation_id;
2754 let runtime = self.api.agent_runtime();
2755 let prepared = {
2756 let mut host = KennedySessionHost {
2757 session: self,
2758 checkpoint,
2759 state: &mut turn.state,
2760 };
2761 match host.prepare_round(round).await? {
2762 kcode_agent_runtime::RoundPreparation::Run(prepared) => {
2763 let manifest_hash = hex::encode(Sha256::digest(prepared.input.as_bytes()));
2764 host.record(kcode_agent_runtime::SessionEvent::InferenceSubmitted {
2765 round,
2766 manifest_hash,
2767 model: prepared.model.clone(),
2768 })
2769 .await?;
2770 prepared
2771 }
2772 kcode_agent_runtime::RoundPreparation::Complete(answer) => {
2773 return self.finish_stepped_turn(turn, answer, checkpoint).await;
2774 }
2775 }
2776 };
2777 turn.completed_rounds = round;
2778 Ok(TurnBoundary::Await(PendingSessionInference {
2779 turn,
2780 action: Box::new(PendingInferenceAction::Start {
2781 runtime,
2782 request: kcode_agent_runtime::SessionInferenceRequest {
2783 user_id: self
2784 .root_node_ids
2785 .first()
2786 .context("session has no user root for intelligence accounting")?
2787 .clone(),
2788 operation_id,
2789 round,
2790 prepared,
2791 },
2792 }),
2793 }))
2794 }
2795
2796 pub async fn apply_inference_wake<C, F>(
2797 &mut self,
2798 wake: SessionInferenceWake,
2799 checkpoint: &mut C,
2800 ) -> anyhow::Result<TurnBoundary>
2801 where
2802 C: FnMut(Value) -> F + Send,
2803 F: Future<Output = anyhow::Result<()>> + Send,
2804 {
2805 wake.turn.lease.validate(&self.turn_lease_slot)?;
2806 let result = self.apply_inference_wake_validated(wake, checkpoint).await;
2807 if result.is_err() {
2808 self.clear_turn_deadlines();
2809 }
2810 result
2811 }
2812
2813 async fn apply_inference_wake_validated<C, F>(
2814 &mut self,
2815 wake: SessionInferenceWake,
2816 checkpoint: &mut C,
2817 ) -> anyhow::Result<TurnBoundary>
2818 where
2819 C: FnMut(Value) -> F + Send,
2820 F: Future<Output = anyhow::Result<()>> + Send,
2821 {
2822 let SessionInferenceWake { mut turn, kind } = wake;
2823 let round = turn.completed_rounds;
2824 let operation_id = turn.state.operation_id;
2825 match kind {
2826 SessionInferenceWakeKind::StartFailed(error) => {
2827 let mut host = KennedySessionHost {
2828 session: self,
2829 checkpoint,
2830 state: &mut turn.state,
2831 };
2832 if let Some(receipt) = inference_error_receipt(&error) {
2833 host.record(kcode_agent_runtime::SessionEvent::ProviderReceipt {
2834 round,
2835 usage: None,
2836 receipt: Box::new(receipt),
2837 continuation: None,
2838 })
2839 .await?;
2840 }
2841 Err(error)
2842 }
2843 SessionInferenceWakeKind::RespondFailed {
2844 mut inference,
2845 error,
2846 } => {
2847 let mut host = KennedySessionHost {
2848 session: self,
2849 checkpoint,
2850 state: &mut turn.state,
2851 };
2852 record_unavailable_inference(&mut host, &mut inference, round).await?;
2853 Err(error)
2854 }
2855 SessionInferenceWakeKind::RespondedStop { inference } => {
2856 drop(inference);
2857 self.finish_stepped_turn(turn, None, checkpoint).await
2858 }
2859 SessionInferenceWakeKind::Event {
2860 mut inference,
2861 event,
2862 } => {
2863 let event = match event {
2864 Ok(Some(event)) => event,
2865 Ok(None) => {
2866 let mut host = KennedySessionHost {
2867 session: self,
2868 checkpoint,
2869 state: &mut turn.state,
2870 };
2871 record_unavailable_inference(&mut host, &mut inference, round).await?;
2872 anyhow::bail!("provider ended without a terminal turn event");
2873 }
2874 Err(error) => {
2875 let mut host = KennedySessionHost {
2876 session: self,
2877 checkpoint,
2878 state: &mut turn.state,
2879 };
2880 if let Some(receipt) = inference_error_receipt(&error) {
2881 host.record(kcode_agent_runtime::SessionEvent::ProviderReceipt {
2882 round,
2883 usage: None,
2884 receipt: Box::new(receipt),
2885 continuation: None,
2886 })
2887 .await?;
2888 }
2889 return Err(error);
2890 }
2891 };
2892
2893 match event {
2894 kcode_agent_runtime::SessionInferenceEvent::ProviderInput { context } => {
2895 let mut host = KennedySessionHost {
2896 session: self,
2897 checkpoint,
2898 state: &mut turn.state,
2899 };
2900 host.record(kcode_agent_runtime::SessionEvent::ProviderInput {
2901 round,
2902 context,
2903 })
2904 .await?;
2905 Ok(TurnBoundary::Await(PendingSessionInference {
2906 turn,
2907 action: Box::new(PendingInferenceAction::Next { inference }),
2908 }))
2909 }
2910 kcode_agent_runtime::SessionInferenceEvent::UsageUpdated { usage } => {
2911 let mut host = KennedySessionHost {
2912 session: self,
2913 checkpoint,
2914 state: &mut turn.state,
2915 };
2916 host.record(kcode_agent_runtime::SessionEvent::UsageUpdated {
2917 round,
2918 usage,
2919 })
2920 .await?;
2921 Ok(TurnBoundary::Await(PendingSessionInference {
2922 turn,
2923 action: Box::new(PendingInferenceAction::Next { inference }),
2924 }))
2925 }
2926 kcode_agent_runtime::SessionInferenceEvent::ToolCall { call_id, call } => {
2927 turn.state.used_tool = true;
2928 let call =
2929 call.map_err(|error| anyhow::anyhow!("Invalid Ktool call: {error}"));
2930 let mut host = KennedySessionHost {
2931 session: self,
2932 checkpoint,
2933 state: &mut turn.state,
2934 };
2935 let outcome = match host.execute_tool(call, operation_id).await {
2936 Ok(outcome) => outcome,
2937 Err(error) => {
2938 record_unavailable_inference(&mut host, &mut inference, round)
2939 .await?;
2940 return Err(error);
2941 }
2942 };
2943 let resume = match host.prepare_provider_resume(outcome).await {
2944 Ok(resume) => resume,
2945 Err(error) => {
2946 record_unavailable_inference(&mut host, &mut inference, round)
2947 .await?;
2948 return Err(error);
2949 }
2950 };
2951 match resume {
2952 kcode_agent_runtime::ProviderResume::Continue(mut outcome) => {
2953 host.state.finish_requested |=
2954 outcome.ok && outcome.finish_after_round;
2955 host.state.emitted_response |=
2956 outcome.ok && outcome.emitted_response;
2957 host.state.pending_capture = outcome.capture.take();
2958 let stop = outcome.stop;
2959 let result = if outcome.ok {
2960 kcode_codex_runtime_v2::ToolResult::success(outcome.text)
2961 } else {
2962 kcode_codex_runtime_v2::ToolResult::failure(outcome.text)
2963 };
2964 Ok(TurnBoundary::Await(PendingSessionInference {
2965 turn,
2966 action: Box::new(PendingInferenceAction::Respond {
2967 inference,
2968 call_id,
2969 result,
2970 stop,
2971 }),
2972 }))
2973 }
2974 kcode_agent_runtime::ProviderResume::Complete(answer) => {
2975 record_unavailable_inference(&mut host, &mut inference, round)
2976 .await?;
2977 self.finish_stepped_turn(turn, answer, checkpoint).await
2978 }
2979 kcode_agent_runtime::ProviderResume::RestartFresh => {
2980 record_unavailable_inference(&mut host, &mut inference, round)
2981 .await?;
2982 turn.at_yielded_boundary = true;
2983 Ok(TurnBoundary::Yield(turn))
2984 }
2985 }
2986 }
2987 kcode_agent_runtime::SessionInferenceEvent::Completed {
2988 answer,
2989 usage,
2990 receipt,
2991 continuation,
2992 } => {
2993 drop(inference);
2994 let mut host = KennedySessionHost {
2995 session: self,
2996 checkpoint,
2997 state: &mut turn.state,
2998 };
2999 host.record(kcode_agent_runtime::SessionEvent::ProviderReceipt {
3000 round,
3001 usage,
3002 receipt,
3003 continuation,
3004 })
3005 .await?;
3006 let capture = host.state.pending_capture.take();
3007 let control = if let Some(capture) = capture {
3008 host.complete_capture(capture, answer).await?
3009 } else {
3010 let completion = kcode_agent_runtime::RoundCompletion {
3011 answer,
3012 used_tool: host.state.used_tool,
3013 finish_requested: host.state.finish_requested,
3014 emitted_response: host.state.emitted_response,
3015 };
3016 host.complete_round(completion).await?
3017 };
3018 host.state.used_tool = false;
3019 host.state.finish_requested = false;
3020 host.state.emitted_response = false;
3021 match control {
3022 kcode_agent_runtime::SessionControl::Continue => {
3023 turn.at_yielded_boundary = true;
3024 Ok(TurnBoundary::Yield(turn))
3025 }
3026 kcode_agent_runtime::SessionControl::Complete(answer) => {
3027 self.finish_stepped_turn(turn, answer, checkpoint).await
3028 }
3029 }
3030 }
3031 }
3032 }
3033 }
3034 }
3035
3036 pub async fn run_pending_turn<C, F>(
3037 &mut self,
3038 operation_id: Uuid,
3039 turn_deadline: Option<TurnDeadline>,
3040 mut checkpoint: C,
3041 ) -> anyhow::Result<Option<String>>
3042 where
3043 C: FnMut(Value) -> F + Send,
3044 F: Future<Output = anyhow::Result<()>> + Send,
3045 {
3046 let Some(turn) = self.begin_pending_turn(operation_id, turn_deadline)? else {
3047 return Ok(None);
3048 };
3049 let mut boundary = self.advance_pending_turn(turn, &mut checkpoint).await?;
3050 loop {
3051 boundary = match boundary {
3052 TurnBoundary::Await(pending) => {
3053 let wake = pending.wait().await;
3054 self.apply_inference_wake(wake, &mut checkpoint).await?
3055 }
3056 TurnBoundary::Yield(turn) => {
3057 self.advance_pending_turn(turn, &mut checkpoint).await?
3058 }
3059 TurnBoundary::Complete(answer) => return Ok(answer),
3060 };
3061 }
3062 }
3063
3064 fn project_descendant<T>(
3065 &mut self,
3066 outcome: Result<kcode_intelligence_router::Accounted<T>, services::ApiError>,
3067 ) -> anyhow::Result<T> {
3068 match outcome {
3069 Ok(accounted) => {
3070 kcode_intelligence_chatend::record_descendant_receipt(
3071 &mut self.journal,
3072 &accounted.receipt,
3073 )?;
3074 Ok(accounted.value)
3075 }
3076 Err(error) => {
3077 if let Some(receipt) = &error.receipt {
3078 kcode_intelligence_chatend::record_descendant_receipt(
3079 &mut self.journal,
3080 receipt,
3081 )?;
3082 }
3083 Err(error.into())
3084 }
3085 }
3086 }
3087
3088 async fn run_subagent(
3089 &mut self,
3090 model: String,
3091 reasoning_effort: Option<String>,
3092 context_node_ids: Vec<String>,
3093 task: String,
3094 parent_operation_id: Uuid,
3095 ) -> anyhow::Result<String> {
3096 let reasoning_effort =
3097 reasoning_effort.unwrap_or_else(|| self.runtime.reasoning_effort.clone());
3098 let mut selected_node_descriptions = Vec::with_capacity(context_node_ids.len());
3099 for node_id in &context_node_ids {
3100 selected_node_descriptions.push(self.api.kmap_node(node_id)?.data.long_description);
3101 }
3102 let user_id = self
3103 .root_node_ids
3104 .first()
3105 .context("session has no user root for subagent intelligence accounting")?
3106 .clone();
3107 let timeout = self.agent_request_timeout();
3108 let runtime = self.api.agent_runtime();
3109 let provider = runtime.resolve_model(&model).await?.provider;
3110 let first_event = self.journal.state().events.len();
3111 let cost_before = self.projection().status;
3112 let subagent_context = SubagentContext::new(
3113 self.root_node_ids.clone(),
3114 self.api.loads_fixed_connections(),
3115 provider,
3116 self.subagent_codex_prompt.clone(),
3117 selected_node_descriptions,
3118 )?;
3119 let initial_sections = subagent_context.initial_sections().to_vec();
3120 let result = {
3121 let mut host = KennedySubagentHost {
3122 session: self,
3123 context: subagent_context,
3124 captures: HashMap::new(),
3125 };
3126 runtime
3127 .run(
3128 kcode_agent_runtime::RunRequest {
3129 user_id,
3130 parent_operation_id,
3131 model,
3132 reasoning_effort,
3133 context: initial_sections,
3134 task,
3135 timeout,
3136 start_metadata: json!({"contextNodeIds":context_node_ids}),
3137 },
3138 &mut host,
3139 )
3140 .await
3141 };
3142 match result {
3143 Ok(result) => {
3144 let cost_after = self.projection().status;
3145 Ok(format!(
3146 "{}\n\n[{}]",
3147 result.answer,
3148 cost_summary(
3149 "subagent cost",
3150 cost_after
3151 .estimated_cost_usd_nanos
3152 .saturating_sub(cost_before.estimated_cost_usd_nanos),
3153 cost_after
3154 .unpriced_provider_calls
3155 .saturating_sub(cost_before.unpriced_provider_calls),
3156 )
3157 ))
3158 }
3159 Err(error) => {
3160 let may_have_effects =
3161 self.journal.state().events[first_event..]
3162 .iter()
3163 .any(|event| {
3164 matches!(
3165 &event.kind,
3166 EventKind::Note { label, .. } if label == "subagent_tool_call"
3167 )
3168 });
3169 if may_have_effects {
3170 Err(error.context(
3171 "the subagent failed after making Ktool calls; some tool effects may already have occurred",
3172 ))
3173 } else {
3174 Err(error)
3175 }
3176 }
3177 }
3178 }
3179
3180 async fn complete_subagent_freeform_write(
3181 &mut self,
3182 context: &mut SubagentContext,
3183 request: FreeformWrite,
3184 contents: String,
3185 budget: &kcode_agent_runtime::ContextBudget,
3186 ) -> anyhow::Result<kcode_agent_runtime::ToolOutcome> {
3187 let kind = request.kind();
3188 let freeform_tool = request.write_tool();
3189 anyhow::ensure!(
3190 context.source_is_open(kind, request.name()),
3191 "{} {:?} is not open in this subagent context. Call {} first.",
3192 kind.label(),
3193 request.name(),
3194 kind.open_tool()
3195 );
3196 let backend_arguments = request.capture_subagent(&mut self.journal, &now(), contents)?;
3197 let preview = self
3198 .api
3199 .managed_source_execute(
3200 &self.rust_lib_session_id,
3201 request.preview_tool(),
3202 backend_arguments.clone(),
3203 Vec::new(),
3204 )
3205 .await?;
3206 let preview = preview
3207 .snapshot
3208 .context("subagent freeform write preview omitted its source snapshot")?;
3209 let preview_state = context.source_state(&preview);
3210 anyhow::ensure!(
3211 budget.fits_state(preview_state.key, preview_state.text),
3212 "{freeform_tool} was not run because its resulting source state would exceed the subagent context limit"
3213 );
3214 let execution = self
3215 .api
3216 .managed_source_execute(
3217 &self.rust_lib_session_id,
3218 freeform_tool,
3219 backend_arguments,
3220 Vec::new(),
3221 )
3222 .await?;
3223 let snapshot = execution
3224 .snapshot
3225 .context("subagent freeform write omitted its resulting source snapshot")?;
3226 let state = context.apply_source_snapshot(snapshot);
3227 Ok(kcode_agent_runtime::ToolOutcome {
3228 text: execution.text,
3229 ok: true,
3230 state_updates: state.update.into_iter().collect(),
3231 displayed_state_keys: Vec::new(),
3232 capture: None,
3233 })
3234 }
3235
3236 async fn complete_freeform_write(
3237 &mut self,
3238 pending: PendingFreeformWrite,
3239 contents: String,
3240 ) -> anyhow::Result<ToolOutcome> {
3241 let request = pending.request;
3242 let freeform_tool = request.write_tool();
3243 let backend_arguments =
3244 request.capture(&mut self.journal, &now(), pending.call_box_id, contents)?;
3245 let preview_result = self
3246 .api
3247 .managed_source_execute(
3248 &self.rust_lib_session_id,
3249 request.preview_tool(),
3250 backend_arguments.clone(),
3251 Vec::new(),
3252 )
3253 .await;
3254 let preview = match preview_result {
3255 Ok(preview) => preview,
3256 Err(error) => {
3257 return Ok(ToolOutcome {
3258 text: format!("{freeform_tool} failed: {error}"),
3259 store_result: true,
3260 ok: false,
3261 end_session: false,
3262 freeform_write: None,
3263 managed_source_snapshot: None,
3264 exact_result: false,
3265 });
3266 }
3267 };
3268 let _preview = preview
3269 .snapshot
3270 .context("freeform write preview omitted the resulting source snapshot")?;
3271 request.source_box_id(&self.journal)?;
3272
3273 let execution_result = self
3274 .api
3275 .managed_source_execute(
3276 &self.rust_lib_session_id,
3277 freeform_tool,
3278 backend_arguments,
3279 Vec::new(),
3280 )
3281 .await;
3282 let execution = match execution_result {
3283 Ok(execution) => execution,
3284 Err(error) => {
3285 return Ok(ToolOutcome {
3286 text: format!("{freeform_tool} failed: {error}"),
3287 store_result: true,
3288 ok: false,
3289 end_session: false,
3290 freeform_write: None,
3291 managed_source_snapshot: None,
3292 exact_result: false,
3293 });
3294 }
3295 };
3296 let snapshot = execution
3297 .snapshot
3298 .context("freeform write omitted its resulting source snapshot")?;
3299 apply_snapshot(&mut self.journal, &now(), snapshot)?;
3300 Ok(ToolOutcome {
3301 text: execution.text,
3302 store_result: false,
3303 ok: true,
3304 end_session: false,
3305 freeform_write: None,
3306 managed_source_snapshot: None,
3307 exact_result: false,
3308 })
3309 }
3310
3311 async fn send_telegram_dm(&mut self, arguments: &Value) -> anyhow::Result<String> {
3312 let request = kcode_telegram_session_coordinator::parse_private_request(arguments)?;
3313 let attachments = self.telegram_delivery_attachments(request.attachments)?;
3314 let caller_holds_user_lock = self.session_type == "telegram"
3315 && self.channel.get("telegramUserId").and_then(Value::as_i64)
3316 == Some(request.telegram_user_id);
3317 self.api
3318 .telegram()
3319 .send_private(kcode_telegram_session_coordinator::PrivateDelivery {
3320 telegram_user_id: request.telegram_user_id,
3321 message: request.message,
3322 attachments,
3323 caller_holds_user_lock,
3324 })
3325 .await
3326 }
3327
3328 async fn send_telegram_group_message(&mut self, arguments: &Value) -> anyhow::Result<String> {
3329 let request = kcode_telegram_session_coordinator::parse_group_request(arguments)?;
3330 let attachments = self.telegram_delivery_attachments(request.attachments)?;
3331 self.api
3332 .telegram()
3333 .send_group(kcode_telegram_session_coordinator::GroupDelivery {
3334 root_node_id: request.root_node_id,
3335 message: request.message,
3336 attachments,
3337 })
3338 .await
3339 }
3340
3341 fn telegram_delivery_attachments(
3342 &mut self,
3343 requests: Vec<kcode_telegram_session_coordinator::AttachmentRequest>,
3344 ) -> anyhow::Result<Vec<kcode_telegram_session_coordinator::Attachment>> {
3345 let api = self.api.clone();
3346 kcode_kennedy_session_objects::delivery_attachments(
3347 &mut self.journal,
3348 requests,
3349 move |canonical_id| api.kmap_file(canonical_id).map_err(Into::into),
3350 )
3351 }
3352
3353 async fn execute_tool(
3354 &mut self,
3355 call: &ToolCall,
3356 operation_id: Uuid,
3357 ) -> anyhow::Result<ToolOutcome> {
3358 self.assert_tool_allowed(&call.name)?;
3359 anyhow::ensure!(
3360 call.name != LAUNCH_SESSION_TOOL,
3361 "LaunchSession requires the checkpointed launch dispatch lane"
3362 );
3363 let decoded = decode(&call.name, &call.arguments)?;
3364 let mut end_session = false;
3365 let mut store_result = true;
3366 let mut freeform_write = None;
3367 let mut managed_source_snapshot = None;
3368 let text = match (call.name.as_str(), decoded) {
3369 ("NoteToSelf", None) => {
3370 decode_note_to_self(&call.arguments)?;
3371 store_result = false;
3372 "Note saved.".into()
3373 }
3374 ("SendTelegramDM", _) => self.send_telegram_dm(&call.arguments).await?,
3375 ("SendTelegramGroupMessage", _) => {
3376 self.send_telegram_group_message(&call.arguments).await?
3377 }
3378 (
3379 "RunSubagent",
3380 Some(DecodedTool::RunSubagent {
3381 model,
3382 reasoning_effort,
3383 context_node_ids,
3384 task,
3385 }),
3386 ) => {
3387 let first_event = self.journal.state().events.len();
3388 match self
3389 .run_subagent(
3390 model,
3391 reasoning_effort,
3392 context_node_ids,
3393 task,
3394 operation_id,
3395 )
3396 .await
3397 {
3398 Ok(response) => response,
3399 Err(error) => {
3400 let may_have_effects = self.journal.state().events[first_event..]
3401 .iter()
3402 .any(|event| {
3403 matches!(
3404 &event.kind,
3405 EventKind::Note { label, .. }
3406 if label == "subagent_tool_call"
3407 )
3408 });
3409 if may_have_effects {
3410 return Err(error.context(
3411 "the subagent failed after making Ktool calls; some tool effects may already have occurred",
3412 ));
3413 }
3414 return Err(error);
3415 }
3416 }
3417 }
3418 ("EndSession", Some(DecodedTool::EndSession { message })) => {
3419 anyhow::ensure!(
3420 !matches!(self.mode, AgentMode::Conversation),
3421 "EndSession is only available during an autonomous or history-ingress session"
3422 );
3423 end_session = true;
3424 if matches!(self.mode, AgentMode::FreeTime)
3425 && let Some(message) = message.filter(|message| !message.trim().is_empty())
3426 {
3427 self.free_time["nextSessionMessage"] = json!(message);
3428 }
3429 "Session ending.".into()
3430 }
3431 ("DehydrateBoxes", Some(DecodedTool::BoxIds(ids))) => {
3432 self.journal.dehydrate_boxes(now(), &ids)?;
3433 format!(
3434 "Dehydrated boxes {}.",
3435 ids.iter()
3436 .map(ToString::to_string)
3437 .collect::<Vec<_>>()
3438 .join(", ")
3439 )
3440 }
3441 ("SummarizeBox", Some(DecodedTool::SummarizeBox { box_id, summary })) => {
3442 self.journal.summarize_box(now(), box_id, summary)?;
3443 format!("Summarized box {box_id}.")
3444 }
3445 ("HydrateBox", Some(DecodedTool::BoxId(id))) => {
3446 self.journal.rehydrate_box(now(), id)?;
3447 let external_event_id = self.pending_external_event_id.clone();
3448 match self.recover_context_overflow(external_event_id.as_deref(), &[id])? {
3449 ContextRecovery::NotNeeded => format!("Hydrated box {id}."),
3450 ContextRecovery::Recovered => {
3451 format!("Hydrated box {id}.\n\n{CONTEXT_OVERFLOW_WARNING}")
3452 }
3453 ContextRecovery::Irreducible => anyhow::bail!(CONTEXT_OVERFLOW_WARNING),
3454 }
3455 }
3456 ("BoxesIntoObjects", Some(DecodedTool::BoxIds(ids))) => {
3457 kcode_kennedy_box_text_objects::stage_box_text_objects(
3458 &mut self.journal,
3459 &ids,
3460 &now(),
3461 )?
3462 }
3463 ("LoadNodes", Some(DecodedTool::LoadNodes(identifiers))) => {
3464 load_durable_batch(self.api.kmap(), &mut self.context, &identifiers)?;
3465 let changed = self.sync_cache_safe_kweb_boxes()?;
3466 store_result = false;
3467 render_load_nodes_result(
3468 &self.journal,
3469 &changed,
3470 &self.runtime_budget().footer_lines(),
3471 )?
3472 }
3473 (
3474 "EmitObject",
3475 Some(DecodedTool::EmitObject {
3476 object_id,
3477 file_name,
3478 }),
3479 ) => {
3480 anyhow::ensure!(
3481 matches!(self.mode, AgentMode::Conversation),
3482 "EmitObject is only available in a conversation"
3483 );
3484 let object = self.resolve_object(&object_id)?;
3485 let file_name = file_name.unwrap_or_else(|| object.file_name.clone());
3486 if let Some(maximum) = self.channel.get("maxObjectBytes").and_then(Value::as_u64) {
3487 anyhow::ensure!(
3488 !object.bytes.is_empty(),
3489 "object {object_id} is empty and cannot be sent through this channel"
3490 );
3491 anyhow::ensure!(
3492 object.bytes.len() as u64 <= maximum,
3493 "object {object_id} is {} bytes, over this channel's {maximum}-byte limit",
3494 object.bytes.len()
3495 );
3496 }
3497 let descriptor = json!({
3498 "objectId":object_id,
3499 "fileName":file_name,
3500 "mediaType":object.media_type,
3501 "byteLength":object.bytes.len(),
3502 });
3503 let mut metadata = json!({
3504 "outputKind":"object",
3505 "attachments":[descriptor.clone()],
3506 });
3507 if let Some(external_event_id) = &self.pending_external_event_id {
3508 metadata["externalEventId"] = json!(external_event_id);
3509 }
3510 let content = BoxContent {
3511 text: String::new(),
3512 objects: vec![object_id.clone()],
3513 metadata,
3514 };
3515 self.journal
3516 .create_box(now(), "Kennedy message", BoxOwner::Kennedy, content)?;
3517 let mut transcript = json!({
3518 "role":"kennedy",
3519 "content":"",
3520 "objects":[object_id],
3521 "attachments":[descriptor],
3522 });
3523 if let Some(external_event_id) = &self.pending_external_event_id {
3524 transcript["externalEventId"] = json!(external_event_id);
3525 }
3526 self.transcript.push(transcript);
3527 store_result = false;
3528 "Object emitted to the user.".into()
3529 }
3530 ("WebSearch", Some(DecodedTool::WebSearch { question, model })) => {
3531 let user_id = self
3532 .root_node_ids
3533 .first()
3534 .context("session has no user root for intelligence accounting")?
3535 .clone();
3536 let outcome = self
3537 .api
3538 .search(
3539 &user_id,
3540 kcode_intelligence_router::SearchRequest {
3541 question,
3542 model,
3543 operation_id: Uuid::new_v4(),
3544 parent_operation_id: Some(operation_id),
3545 },
3546 )
3547 .await;
3548 let result = self.project_descendant(outcome)?;
3549 render_web_search_result(&result)?
3550 }
3551 ("WebFetch", Some(DecodedTool::WebFetch(url))) => {
3552 let user_id = self
3553 .root_node_ids
3554 .first()
3555 .context("session has no user root for intelligence accounting")?;
3556 let result = self
3557 .api
3558 .fetch(
3559 user_id,
3560 kcode_intelligence_router::FetchRequest {
3561 url,
3562 operation_id: Uuid::new_v4(),
3563 parent_operation_id: Some(operation_id),
3564 },
3565 )
3566 .await?;
3567 render_web_fetch_result(&result)?
3568 }
3569 ("StageTelegramGroupMedia", Some(DecodedTool::StageTelegramGroupMedia(message_id))) => {
3570 let media_ref = kcode_telegram_session_coordinator::group_media_reference(
3571 &self.group_context,
3572 message_id,
3573 )?;
3574 let chat_id = media_ref.chat_id;
3575 let api = self.api.clone();
3576 let staged = kcode_kennedy_session_objects::stage_telegram_group_media(
3577 &mut self.journal,
3578 kcode_kennedy_session_objects::TelegramStageRequest {
3579 chat_id,
3580 message_id,
3581 maximum_bytes: MAX_MEDIA_ENRICHMENT_BYTES,
3582 transport_metadata: media_ref.transport_metadata(),
3583 recorded_at: now(),
3584 },
3585 || api.telegram().group_message_media(chat_id, message_id),
3586 |media_type| {
3587 kcode_telegram_session_coordinator::group_media_file_name(
3588 &media_ref, media_type,
3589 )
3590 },
3591 )?;
3592 render(RenderRequest::StagedTelegramMedia {
3593 pending_id: &staged.descriptor.pending_id,
3594 kind: &staged.kind,
3595 file_name: &staged.descriptor.file_name,
3596 media_type: &staged.descriptor.media_type,
3597 size_bytes: staged.descriptor.size_bytes,
3598 message_id,
3599 reused: staged.reused,
3600 })?
3601 }
3602 (
3603 "TranscribeAudio",
3604 Some(DecodedTool::MediaEnrichment {
3605 object_id,
3606 model,
3607 prompt,
3608 }),
3609 ) => {
3610 let object = self.resolve_media_object(&object_id)?;
3611 validate(ValidationRequest::TranscribableAudio(&object.media_type))?;
3612 validate(ValidationRequest::TranscriptionModel(&model))?;
3613 let user_id = self
3614 .root_node_ids
3615 .first()
3616 .context("session has no user root for intelligence accounting")?
3617 .clone();
3618 let outcome = self
3619 .api
3620 .transcribe_audio(
3621 &user_id,
3622 &model,
3623 &prompt,
3624 object.bytes,
3625 object.file_name.clone(),
3626 &object.media_type,
3627 None,
3628 operation_id,
3629 )
3630 .await;
3631 let result = self.project_descendant(outcome)?;
3632 render_audio_transcription_result(
3633 &object.object_id,
3634 &object.file_name,
3635 &object.media_type,
3636 &result,
3637 )?
3638 }
3639 (
3640 "AnnotateMedia",
3641 Some(DecodedTool::MediaEnrichment {
3642 object_id,
3643 model,
3644 prompt,
3645 }),
3646 ) => {
3647 let media = self.resolve_media_object(&object_id)?;
3648 validate(ValidationRequest::Annotation {
3649 model: &model,
3650 media_type: &media.media_type,
3651 })?;
3652 let user_id = self
3653 .root_node_ids
3654 .first()
3655 .context("session has no user root for intelligence accounting")?
3656 .clone();
3657 let outcome = self
3658 .api
3659 .annotate_media(
3660 &user_id,
3661 &model,
3662 &prompt,
3663 media.bytes,
3664 media.file_name.clone(),
3665 &media.media_type,
3666 operation_id,
3667 )
3668 .await;
3669 let result = self.project_descendant(outcome)?;
3670 render_media_annotation_result(
3671 &media.object_id,
3672 &media.file_name,
3673 &media.media_type,
3674 &result,
3675 )?
3676 }
3677 (
3678 "GenerateImage",
3679 Some(DecodedTool::GenerateImage {
3680 model,
3681 prompt,
3682 reference_object_ids,
3683 }),
3684 ) => {
3685 let mut references = Vec::with_capacity(reference_object_ids.len());
3686 for object_id in &reference_object_ids {
3687 references.push(self.resolve_image_object(object_id)?);
3688 }
3689 let user_id = self
3690 .root_node_ids
3691 .first()
3692 .context("session has no user root for intelligence accounting")?
3693 .clone();
3694 let outcome = self
3695 .api
3696 .generate_image(&user_id, &model, &prompt, references, operation_id)
3697 .await;
3698 let result = self.project_descendant(outcome)?;
3699 let size = result.bytes.len();
3700 let file_name =
3701 format!("generated-image.{}", image_extension(&result.content_type));
3702 let object_id = self.api.save_generated_image(
3703 result.bytes,
3704 &file_name,
3705 &result.content_type,
3706 &result.model,
3707 )?;
3708 format!(
3709 "Generated image.\nObject: {object_id}\nFile: {file_name}\nContent type: {}\nSize: {size} bytes\nModel: {}\nUse EmitObject with {object_id} to deliver it.",
3710 result.content_type, result.model
3711 )
3712 }
3713 ("ExtractDocumentText", Some(DecodedTool::ObjectId(object_id))) => {
3714 let object = self.resolve_media_object(&object_id)?;
3715 validate(ValidationRequest::ExtractableDocument {
3716 media_type: &object.media_type,
3717 file_name: &object.file_name,
3718 })?;
3719 let result = self
3720 .api
3721 .extract_document(object.bytes, object.file_name.clone(), &object.media_type)
3722 .await?;
3723 render_document_extraction_result(&object.object_id, &object.file_name, &result)?
3724 }
3725 (name, None) if SPEECH_CLASSIFICATION_TOOLS.contains(&name) => {
3726 self.api
3727 .execute_speech_classification_tool(name, call.arguments.clone())
3728 .await?
3729 }
3730 (name, None) if TASK_BOARD_TOOLS.contains(&name) => {
3731 self.execute_task_board_tool(name, &call.arguments).await?
3732 }
3733 (name, Some(decoded)) if is_kweb_mutation(name) => {
3734 let (text, _) = execute_kweb_mutation(
3735 name,
3736 decoded,
3737 &self.context,
3738 &mut self.plan,
3739 &mut self.journal,
3740 )?;
3741 self.sync_cache_safe_kweb_boxes()?;
3742 text
3743 }
3744 (name, None)
3745 if RUST_LIB_TOOLS.contains(&name)
3746 || WEB_LIB_TOOLS.contains(&name)
3747 || RUST_BIN_TOOLS.contains(&name) =>
3748 {
3749 if let Some(request) = prepare_freeform_write(&self.journal, name, &call.arguments)?
3750 {
3751 store_result = false;
3752 let acknowledgement = request.acknowledgement();
3753 freeform_write = Some(request);
3754 acknowledgement
3755 } else {
3756 let object_ids = if name == CALL_RUST_BIN_TOOL {
3757 decode_managed_objects(ManagedObjectArguments::RustBinary(&call.arguments))?
3758 } else if name == ATTACH_OBJECT_WEB_LIB_TOOL {
3759 decode_managed_objects(ManagedObjectArguments::WebLibraryAttachment(
3760 &call.arguments,
3761 ))?
3762 } else {
3763 Vec::new()
3764 };
3765 let mut objects = Vec::with_capacity(object_ids.len());
3766 for object_id in object_ids {
3767 objects.push(self.resolve_object(&object_id)?.bytes);
3768 }
3769 let execution = self
3770 .api
3771 .managed_source_execute(
3772 &self.rust_lib_session_id,
3773 name,
3774 call.arguments.clone(),
3775 objects,
3776 )
3777 .await?;
3778 if let Some(snapshot) = execution.snapshot {
3779 managed_source_snapshot = Some(snapshot);
3780 store_result = false;
3781 }
3782 execution.text
3783 }
3784 }
3785 (name, Some(_)) => {
3786 anyhow::bail!("decoded contract for {name} did not match its dispatch lane")
3787 }
3788 (name, None) => anyhow::bail!("Tool {name} is not available"),
3789 };
3790 Ok(ToolOutcome {
3791 text,
3792 store_result,
3793 ok: true,
3794 end_session,
3795 freeform_write,
3796 managed_source_snapshot,
3797 exact_result: false,
3798 })
3799 }
3800
3801 async fn execute_task_board_tool(
3802 &self,
3803 name: &str,
3804 arguments: &Value,
3805 ) -> anyhow::Result<String> {
3806 let board = self
3807 .api
3808 .task_board()
3809 .context("task board is not configured")?
3810 .clone();
3811 let name = name.to_owned();
3812 let arguments = arguments.clone();
3813 let user_id = self
3814 .root_node_ids
3815 .first()
3816 .context("session has no user root for task-category lookup")?
3817 .clone();
3818 tokio::task::spawn_blocking(move || -> anyhow::Result<String> {
3819 let output = match name.as_str() {
3820 "CreateTaskCategory" => serde_json::to_string_pretty(
3821 &board.create_category(serde_json::from_value(arguments)?)?,
3822 )?,
3823 "GetTaskCategory" => {
3824 let call: CategoryCall = serde_json::from_value(arguments)?;
3825 serde_json::to_string_pretty(&board.category(
3826 &call.category_id,
3827 kcode_task_board::BrowsePage {
3828 user_id,
3829 offset: call.offset,
3830 limit: call.limit,
3831 },
3832 )?)?
3833 }
3834 "RemoveTaskCategory" => {
3835 let call: CategoryId = serde_json::from_value(arguments)?;
3836 board.remove_category(&call.category_id)?;
3837 format!("Removed category {}.", call.category_id)
3838 }
3839 "CreateTask" => serde_json::to_string_pretty(
3840 &board.create_task(serde_json::from_value(arguments)?)?,
3841 )?,
3842 "GetTask" => {
3843 let call: TaskId = serde_json::from_value(arguments)?;
3844 serde_json::to_string_pretty(&board.task(&call.task_id)?)?
3845 }
3846 "UpdateTask" => serde_json::to_string_pretty(
3847 &board.update_task(serde_json::from_value(arguments)?)?,
3848 )?,
3849 "RemoveTask" => {
3850 let call: TaskId = serde_json::from_value(arguments)?;
3851 board.remove_task(&call.task_id)?;
3852 format!("Removed task {}.", call.task_id)
3853 }
3854 "GetTopTaskOrphan" => {
3855 let _: EmptyCall = serde_json::from_value(arguments)?;
3856 serde_json::to_string_pretty(&board.top_orphan()?)?
3857 }
3858 _ => anyhow::bail!("Tool {name} is not a task-board operation"),
3859 };
3860 Ok(output)
3861 })
3862 .await
3863 .context("task-board worker stopped")?
3864 }
3865
3866 fn assert_tool_allowed(&self, name: &str) -> anyhow::Result<()> {
3867 let write = matches!(
3868 name,
3869 "ConnectNodes"
3870 | "ConsolidateFanout"
3871 | "SetFixedConnection"
3872 | "CreateNode"
3873 | "UpdateNode"
3874 );
3875 anyhow::ensure!(
3876 !write || !matches!(self.mode, AgentMode::Conversation),
3877 "{name} requires the global Kweb write lane and is unavailable in a read-only conversation"
3878 );
3879 if name == "EndSession" {
3880 anyhow::ensure!(
3881 !matches!(self.mode, AgentMode::Conversation),
3882 "EndSession is unavailable in a conversation"
3883 );
3884 }
3885 if name == LAUNCH_SESSION_TOOL {
3886 anyhow::ensure!(
3887 self.launch_session_authorized(),
3888 "LaunchSession is unavailable without a genuine current user turn in an eligible conversation"
3889 );
3890 }
3891 Ok(())
3892 }
3893
3894 fn sync_cache_safe_kweb_boxes(&mut self) -> anyhow::Result<Vec<BoxId>> {
3895 let before = kweb_slot_box_ids(&self.journal);
3896 let (updates, creates) = self.plan.context_projection();
3897 let stale = self
3898 .context
3899 .sync_load_chatend(&mut self.journal, now(), &updates, &creates)
3900 .map_err(anyhow::Error::new)?;
3901 let after = kweb_slot_box_ids(&self.journal);
3902 Ok(load_box_changes(&before, &after, &stale))
3903 }
3904
3905 fn record_tool_invocation(
3906 &mut self,
3907 name: &str,
3908 arguments: Value,
3909 ) -> anyhow::Result<RecordedToolInvocation> {
3910 let invocation_id = Uuid::new_v4().to_string();
3911 let invocation = RecordedToolInvocation {
3912 tool_instance: tool_instance_for_invocation(name, &invocation_id),
3913 invocation_id,
3914 tool_name: name.into(),
3915 };
3916 self.journal.record(
3917 now(),
3918 EventKind::ToolInvoked {
3919 tool_instance: invocation.tool_instance.clone(),
3920 tool_name: invocation.tool_name.clone(),
3921 arguments,
3922 invocation_id: Some(invocation.invocation_id.clone()),
3923 },
3924 )?;
3925 Ok(invocation)
3926 }
3927
3928 fn record_tool_completion(
3929 &mut self,
3930 invocation: Option<&RecordedToolInvocation>,
3931 outcome: Value,
3932 ) -> anyhow::Result<EventId> {
3933 record_tool_completion_event(&mut self.journal, invocation, outcome)
3934 }
3935
3936 fn finalize_kweb_session(&mut self) -> anyhow::Result<()> {
3937 self.provider_affinity = None;
3938 self.next_thread_reset_reason = None;
3939 if self.commit_receipt.is_some() {
3940 return Ok(());
3941 }
3942 self.repair_unfinished_tools()?;
3943 self.journal.seal()?;
3944 let archive = self.journal.archive_bytes()?;
3945 let object_locations = self
3946 .journal
3947 .objects()
3948 .iter()
3949 .map(|(id, location)| (id.clone(), location.clone()))
3950 .collect::<Vec<_>>();
3951 let mut objects = BTreeMap::new();
3952 for (id, location) in object_locations {
3953 let pending_id = id.to_string();
3954 let transport_kind =
3955 kcode_kennedy_session_objects::staged_descriptor(&self.journal, &id)?
3956 .transport_kind;
3957 let bytes = encode_file(
3958 &pending_id,
3959 location.metadata.file_name.as_deref(),
3960 &location.metadata.media_type,
3961 transport_kind.as_deref(),
3962 self.journal.read_object(&id)?,
3963 )
3964 .with_context(|| format!("encoding staged object {pending_id}"))?;
3965 anyhow::ensure!(
3966 objects.insert(pending_id.clone(), bytes).is_none(),
3967 "duplicate staged object {pending_id}"
3968 );
3969 }
3970 let material = self.plan.commit_material()?;
3971 let result = self.api.commit_kweb_session(CommitRequest {
3972 idempotency_key: self.journal.state().metadata.session_id.clone(),
3973 author: self.commit_author.clone(),
3974 source_created_at: DateTime::parse_from_rfc3339(&self.started_at)
3975 .context("session start timestamp is invalid")?
3976 .with_timezone(&Utc),
3977 archive,
3978 objects,
3979 creates: material.creates,
3980 updates: material.updates,
3981 })?;
3982 self.journal
3983 .mark_completed(result.session_object_id.to_string());
3984 self.commit_receipt = Some(result);
3985 Ok(())
3986 }
3987
3988 fn prepare_free_time_round(&mut self) -> anyhow::Result<bool> {
3989 if !matches!(self.mode, AgentMode::FreeTime) {
3990 return Ok(false);
3991 }
3992 let Some(deadline) = deadline(&self.free_time) else {
3993 return Ok(false);
3994 };
3995 if Utc::now() >= deadline {
3996 self.free_time_end_reason = Some("deadline".into());
3997 self.journal.create_box(
3998 now(),
3999 "Self-time timer",
4000 BoxOwner::Controller,
4001 BoxContent::text(
4002 "The self-time deadline has arrived. Finish without starting more tool work.",
4003 ),
4004 )?;
4005 return Ok(true);
4006 }
4007 Ok(false)
4008 }
4009
4010 fn agent_request_timeout(&self) -> Option<Duration> {
4011 if matches!(self.mode, AgentMode::Conversation) && self.session_type == "conversation" {
4012 return Some(BROWSER_CONVERSATION_REQUEST_TIMEOUT);
4013 }
4014 if matches!(self.mode, AgentMode::Ingress { .. }) {
4015 return Some(HISTORY_INGRESS_REQUEST_TIMEOUT);
4016 }
4017 if matches!(self.mode, AgentMode::Wakeup) {
4018 return Some(WAKEUP_REQUEST_TIMEOUT);
4019 }
4020 if matches!(self.mode, AgentMode::FreeTime) {
4021 let deadline = deadline(&self.free_time)?;
4022 return Some(Duration::from_secs(
4023 (deadline - Utc::now()).num_seconds().max(1) as u64
4024 + SELF_TIME_HARD_STOP_ALLOWANCE.as_secs(),
4025 ));
4026 }
4027 None
4028 }
4029
4030 pub fn refresh_telegram_group_context(
4031 &mut self,
4032 group_context: &Value,
4033 current_message_id: Option<&str>,
4034 ) -> anyhow::Result<()> {
4035 if self.session_type != "telegram-group" {
4036 return Ok(());
4037 }
4038 self.invalidate_active_stepped_turn();
4039 self.channel["groupContext"] = group_context.clone();
4040 self.group_context = group_context.clone();
4041 self.journal.create_box(
4042 now(),
4043 "Telegram group update",
4044 BoxOwner::Controller,
4045 BoxContent::text(kcode_telegram_session_coordinator::format_group_context(
4046 group_context,
4047 )),
4048 )?;
4049 self.recover_context_overflow(current_message_id, &[])?;
4050 Ok(())
4051 }
4052
4053 pub fn finalize_free_time(&mut self, reason: &str) -> anyhow::Result<()> {
4054 anyhow::ensure!(
4055 matches!(reason, "tool" | "deadline" | "hard-stop" | "user-stop"),
4056 "invalid self-time completion reason"
4057 );
4058 self.invalidate_active_stepped_turn();
4059 self.free_time["sliceEndedReason"] = json!(reason);
4060 self.free_time["sliceEndedAt"] = json!(now());
4061 self.pending_turn = false;
4062 self.pending_external_event_id = None;
4063 self.clear_launch_turn_authority();
4064 Ok(())
4065 }
4066
4067 pub fn commit_current_write_session(&mut self) -> anyhow::Result<()> {
4068 anyhow::ensure!(
4069 matches!(
4070 self.mode,
4071 AgentMode::FreeTime | AgentMode::Wakeup | AgentMode::Ingress { .. }
4072 ),
4073 "a read-only conversation cannot be committed as a Kweb write session"
4074 );
4075 self.invalidate_active_stepped_turn();
4076 self.finalize_kweb_session()?;
4077 self.completed = true;
4078 Ok(())
4079 }
4080
4081 pub fn snapshot(&self) -> anyhow::Result<Value> {
4082 let projection = self.projection();
4083 let submitted = self
4084 .journal
4085 .state()
4086 .current_ingress_attempt_events()
4087 .iter()
4088 .rev()
4089 .find_map(|event| {
4090 let EventKind::ProviderInputSubmitted { round, context, .. } = &event.kind else {
4091 return None;
4092 };
4093 Some((event.recorded_at.as_str(), *round, context))
4094 });
4095 let (chatend_text, chatend_text_source, structured_material) = match submitted {
4096 Some((submitted_at, round, submitted)) => (
4097 submitted.input.clone(),
4098 "submitted",
4099 json!({
4100 "provider":submitted.provider,
4101 "model":submitted.model,
4102 "reasoningEffort":submitted.reasoning_effort,
4103 "baseInstructions":submitted.base_instructions,
4104 "developerInstructions":submitted.developer_instructions,
4105 "tools":submitted.tools,
4106 "round":round,
4107 "submittedAt":submitted_at,
4108 }),
4109 ),
4110 None => (projection.render(), "reconstructed", Value::Null),
4111 };
4112 let session_status = projection.status.clone();
4113 let completed_invocations = completed_invocation_ids(&self.journal);
4114 let launch_intents = pruned_launch_intents(
4115 &self.launch_intents,
4116 self.launch_user_turn_id,
4117 &completed_invocations,
4118 );
4119 Ok(json!({
4120 "format":"kennedy-chatend",
4121 "version":1,
4122 "stateVersion":CHECKPOINT_STATE_VERSION,
4123 "sessionId":self.journal.state().metadata.session_id,
4124 "chatendMetadata":self.journal.state().metadata,
4125 "sessionType":self.session_type,
4126 "sourceSessionType":self.source_session_type,
4127 "channel":self.channel,
4128 "freeTime":self.free_time,
4129 "orchestration":self.orchestration,
4130 "provenanceId":self.provenance_id,
4131 "launchProvenance":self.launch_provenance,
4132 "launchContextNodeIds":self.launch_context_node_ids,
4133 "launchUserTurnId":self.launch_user_turn_id,
4134 "launchIntents":launch_intents,
4135 "rustLibSessionId":self.rust_lib_session_id,
4136 "rootNodeIds":self.root_node_ids,
4137 "referenceRootNodeIds":self.reference_root_node_ids,
4138 "startedAt":self.started_at,
4139 "transcript":self.transcript,
4140 "pendingTurn":self.pending_turn,
4141 "pendingExternalEventId":self.pending_external_event_id,
4142 "roundsUsed":self.rounds_used,
4143 "providerAffinity":self.provider_affinity,
4144 "nextThreadResetReason":self.next_thread_reset_reason,
4145 "completed":self.completed,
4146 "sessionObjectId":self.journal.state().completed_session_object,
4147 "commitReceipt":self.commit_receipt,
4148 "commitAuthor":self.commit_author,
4149 "providerModel":self.runtime.model,
4150 "kwebPlan":self.plan.checkpoint_value()?,
4151 "boxCount":self.journal.state().boxes.len(),
4152 "eventCount":self.journal.state().events.len(),
4153 "boxes":self.journal.state().boxes,
4154 "events":self.journal.state().events,
4155 "context":projection,
4156 "sessionStatus":session_status,
4157 "chatendText":chatend_text,
4158 "chatendTextSource":chatend_text_source,
4159 "structuredMaterial":structured_material,
4160 }))
4161 }
4162
4163 pub async fn release_managed_sources(&self) {
4164 self.api
4165 .release_managed_sources(&self.rust_lib_session_id)
4166 .await;
4167 }
4168}
4169
4170impl<C, F> kcode_agent_runtime::SessionHost for KennedySessionHost<'_, C>
4171where
4172 C: FnMut(Value) -> F + Send,
4173 F: Future<Output = anyhow::Result<()>> + Send,
4174{
4175 fn prepare_round<'a>(
4176 &'a mut self,
4177 round: u64,
4178 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::RoundPreparation> {
4179 Box::pin(async move {
4180 self.session.rounds_used = round;
4181 self.state.deadline_after_response = self.session.prepare_free_time_round()?;
4182 let external_event_id = self.session.pending_external_event_id.clone();
4183 if self
4184 .session
4185 .recover_context_overflow(external_event_id.as_deref(), &[])?
4186 == ContextRecovery::Irreducible
4187 || (matches!(self.session.mode, AgentMode::Ingress { .. })
4188 && self.session.ingress_force_commit_requested())
4189 {
4190 return Ok(kcode_agent_runtime::RoundPreparation::Complete(None));
4191 }
4192 let ingress_time_remaining = self.session.ingress_time_remaining()?;
4193 let timeout = self.session.agent_request_timeout();
4194 self.session.begin_provider_call_budget(timeout);
4195 let tool_description = call_ktool_description(self.session.launch_session_authorized());
4196 let material_fingerprint = self
4197 .session
4198 .provider_material_fingerprint(&tool_description);
4199 let mut thread_reset_reason = self.session.next_thread_reset_reason.take();
4200 let mut continuation = None;
4201 let mut resume_after = None;
4202 if let Some(affinity) = &self.session.provider_affinity {
4203 if affinity.material_fingerprint == material_fingerprint {
4204 continuation = Some(affinity.continuation.clone());
4205 resume_after = Some(affinity.synchronized_event_id);
4206 } else {
4207 self.session.provider_affinity = None;
4208 thread_reset_reason = Some("provider_material_changed".into());
4209 }
4210 }
4211 if continuation.is_some() {
4212 self.session.provider_affinity = None;
4213 self.session.next_thread_reset_reason =
4214 Some("prior_provider_turn_ambiguous".into());
4215 }
4216 let footer_lines = self.session.runtime_budget().footer_lines();
4217 let prepared = if let Some(remaining_seconds) = ingress_time_remaining {
4218 self.session
4219 .journal
4220 .prepare_provider_projection_with_ingress_time(
4221 now(),
4222 &footer_lines,
4223 &material_fingerprint,
4224 resume_after,
4225 remaining_seconds,
4226 self.session.previous_ingress_attempt_timed_out,
4227 )?
4228 } else {
4229 self.session.journal.prepare_provider_projection(
4230 now(),
4231 &footer_lines,
4232 &material_fingerprint,
4233 resume_after,
4234 )?
4235 };
4236 if let Some(reason) = prepared.thread_reset_reason.clone() {
4237 self.session.provider_affinity = None;
4238 continuation = None;
4239 thread_reset_reason = Some(reason);
4240 }
4241 if continuation.is_none() && thread_reset_reason.is_some() {
4242 self.session.next_thread_reset_reason = thread_reset_reason.clone();
4243 }
4244 let input = prepared.projection.render();
4245 let projection_hash = hex::encode(Sha256::digest(input.as_bytes()));
4246 let provider_input_hash =
4247 hex::encode(Sha256::digest(prepared.provider_input.as_bytes()));
4248 let provider_input_bytes = prepared.provider_input.len() as u64;
4249 let thread_action = if continuation.is_some() {
4250 "resume"
4251 } else {
4252 "start"
4253 }
4254 .to_owned();
4255 self.state.prepared_cache = Some(PreparedCacheObservation {
4256 cacheable_prefix_bytes: prepared.cacheable_prefix_bytes,
4257 expectation: prepared.expectation,
4258 material_fingerprint,
4259 projection_hash,
4260 logical_input: input.clone(),
4261 provider_input_hash,
4262 provider_input_bytes,
4263 thread_action,
4264 thread_reset_reason,
4265 estimated_input_tokens: prepared.projection.estimated_tokens,
4266 raw_estimated_input_tokens: prepared.projection.raw_estimated_tokens,
4267 provider: String::new(),
4268 model: self.session.runtime.model.clone(),
4269 });
4270 Ok(kcode_agent_runtime::RoundPreparation::Run(
4271 kcode_agent_runtime::PreparedRound {
4272 input,
4273 provider_input: prepared.provider_input,
4274 continuation,
4275 model: self.session.runtime.model.clone(),
4276 reasoning_effort: self.session.runtime.reasoning_effort.clone(),
4277 tool_description,
4278 timeout,
4279 },
4280 ))
4281 })
4282 }
4283
4284 fn record<'a>(
4285 &'a mut self,
4286 event: kcode_agent_runtime::SessionEvent,
4287 ) -> kcode_agent_runtime::HostFuture<'a, ()> {
4288 Box::pin(async move {
4289 match event {
4290 kcode_agent_runtime::SessionEvent::InferenceSubmitted {
4291 manifest_hash,
4292 model,
4293 ..
4294 } => {
4295 let prepared =
4296 self.state.prepared_cache.as_ref().context(
4297 "inference was submitted before provider context preparation",
4298 )?;
4299 anyhow::ensure!(
4300 prepared.projection_hash == manifest_hash,
4301 "provider input hash changed after context preparation"
4302 );
4303 self.state.accounting = Some(kcode_intelligence_chatend::TopLevelCall::new(
4304 manifest_hash.clone(),
4305 model,
4306 ));
4307 self.session.journal.record(
4308 now(),
4309 EventKind::InferenceSubmitted {
4310 manifest_hash,
4311 estimated_input_tokens: prepared.estimated_input_tokens,
4312 raw_estimated_input_tokens: Some(prepared.raw_estimated_input_tokens),
4313 },
4314 )?;
4315 }
4316 kcode_agent_runtime::SessionEvent::ProviderInput { round, context } => {
4317 let prepared = self
4318 .state
4319 .prepared_cache
4320 .as_ref()
4321 .context("provider context arrived before context preparation")?;
4322 anyhow::ensure!(
4323 hex::encode(Sha256::digest(context.input.as_bytes()))
4324 == prepared.provider_input_hash,
4325 "provider submitted transport input different from the prepared continuation delta"
4326 );
4327 let provider = context.provider.clone();
4328 let model = context.model.clone();
4329 let synchronized_after = self.session.journal.record(
4330 now(),
4331 EventKind::ProviderInputSubmitted {
4332 round,
4333 context: ProviderContext {
4334 input: prepared.logical_input.clone(),
4335 provider: context.provider,
4336 model: context.model,
4337 reasoning_effort: context.reasoning_effort,
4338 base_instructions: context.base_instructions,
4339 developer_instructions: context.developer_instructions,
4340 tools: context
4341 .tools
4342 .into_iter()
4343 .map(|tool| ProviderToolDefinition {
4344 name: tool.name,
4345 description: tool.description,
4346 input_schema: tool.input_schema,
4347 })
4348 .collect(),
4349 },
4350 transport_input_hash: Some(prepared.provider_input_hash.clone()),
4351 transport_input_bytes: Some(prepared.provider_input_bytes),
4352 thread_action: Some(prepared.thread_action.clone()),
4353 thread_reset_reason: prepared.thread_reset_reason.clone(),
4354 cacheable_prefix_bytes: prepared.cacheable_prefix_bytes,
4355 material_fingerprint: prepared.material_fingerprint.clone(),
4356 cache_expectation: prepared.expectation.label().into(),
4357 planned_invalidation_reason: prepared
4358 .expectation
4359 .planned_reason()
4360 .map(str::to_owned),
4361 },
4362 )?;
4363 self.state.provider_synchronized_after = Some(synchronized_after);
4364 if let Some(prepared) = self.state.prepared_cache.as_mut() {
4365 prepared.provider = provider;
4366 prepared.model = model;
4367 }
4368 }
4369 kcode_agent_runtime::SessionEvent::UsageUpdated { usage, .. } => {
4370 self.state
4371 .accounting
4372 .as_mut()
4373 .context("provider usage arrived before inference submission")?
4374 .usage_updated(&mut self.session.journal, &now(), &usage)?;
4375 }
4376 kcode_agent_runtime::SessionEvent::ProviderReceipt {
4377 usage,
4378 receipt,
4379 continuation,
4380 ..
4381 } => {
4382 self.state
4383 .accounting
4384 .take()
4385 .context("provider receipt arrived before inference submission")?
4386 .completed(&mut self.session.journal, &now(), usage.as_ref())?;
4387 let prepared = self
4388 .state
4389 .prepared_cache
4390 .take()
4391 .context("provider receipt arrived before context preparation")?;
4392 if let Some(reason) = self.state.restart_fresh_reason.take() {
4393 anyhow::ensure!(
4394 continuation.is_none(),
4395 "restart-fresh receipt unexpectedly retained a native continuation"
4396 );
4397 self.session.provider_affinity = None;
4398 self.session.next_thread_reset_reason = Some(reason);
4399 } else if let Some(continuation) = continuation {
4400 anyhow::ensure!(
4401 receipt.provider_thread_id.as_deref()
4402 == Some(continuation.thread_id.as_str()),
4403 "provider receipt thread differs from continuation state"
4404 );
4405 let synchronized_event_id = self
4406 .session
4407 .journal
4408 .state()
4409 .events
4410 .last()
4411 .context("provider completion did not create a journal event")?
4412 .id;
4413 self.session.provider_affinity = Some(ProviderAffinityState {
4414 continuation,
4415 synchronized_event_id,
4416 material_fingerprint: prepared.material_fingerprint.clone(),
4417 });
4418 self.session.next_thread_reset_reason = None;
4419 } else {
4420 self.session.provider_affinity = None;
4421 self.session.next_thread_reset_reason = Some(
4422 if prepared.thread_action == "resume" {
4423 "provider_thread_resume_unavailable"
4424 } else {
4425 "provider_continuation_unavailable"
4426 }
4427 .into(),
4428 );
4429 }
4430 log_primary_thread_observation(
4431 self.state.operation_id,
4432 self.session.rounds_used,
4433 &self.session.runtime.model,
4434 &prepared,
4435 receipt.provider_thread_id.as_deref(),
4436 usage.as_ref().map_or(0, |usage| usage.input_tokens),
4437 usage.as_ref().map_or(0, |usage| usage.cached_input_tokens),
4438 );
4439 }
4440 }
4441 let snapshot = self.session.snapshot()?;
4442 (self.checkpoint)(snapshot).await
4443 })
4444 }
4445
4446 fn execute_tool<'a>(
4447 &'a mut self,
4448 call: anyhow::Result<kcode_agent_runtime::ToolCall>,
4449 operation_id: Uuid,
4450 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionToolOutcome> {
4451 Box::pin(async move {
4452 if let Some(pending) = &self.state.pending_freeform_write {
4453 let text = format!(
4454 "{} is awaiting the complete file contents; no other Ktool can run before that output.",
4455 pending.request.write_tool()
4456 );
4457 self.session
4458 .record_tool_completion(None, json!({"ok":false,"result":text}))?;
4459 return Ok(kcode_agent_runtime::SessionToolOutcome {
4460 text,
4461 ok: false,
4462 capture: Some(json!(true)),
4463 stop: false,
4464 finish_after_round: false,
4465 emitted_response: false,
4466 });
4467 }
4468 let tool_started_at = std::time::Instant::now();
4469 let mut created_call_box_id = None;
4470 let mut recorded_invocation = None;
4471 let transcript_start = self.session.transcript.len();
4472 let mut emitted_response = false;
4473 let mut outcome = match call {
4474 Ok(call) => {
4475 let call = ToolCall {
4476 name: call.name,
4477 arguments: call.arguments,
4478 };
4479 let call_name = format!("Kennedy tool call: {}", call.name);
4480 let call_content = tool_invocation_content(&call.name, &call.arguments)?;
4481 recorded_invocation = Some(
4482 self.session
4483 .record_tool_invocation(&call.name, call.arguments.clone())?,
4484 );
4485 created_call_box_id = Some(self.session.journal.create_box(
4486 now(),
4487 call_name,
4488 BoxOwner::Kennedy,
4489 call_content,
4490 )?);
4491 let external_event_id = self.session.pending_external_event_id.clone();
4492 if self
4493 .session
4494 .recover_context_overflow(external_event_id.as_deref(), &[])?
4495 == ContextRecovery::Irreducible
4496 {
4497 ToolOutcome {
4498 text: CONTEXT_OVERFLOW_WARNING.into(),
4499 store_result: false,
4500 ok: false,
4501 end_session: false,
4502 freeform_write: None,
4503 managed_source_snapshot: None,
4504 exact_result: false,
4505 }
4506 } else if call.name == LAUNCH_SESSION_TOOL {
4507 let invocation = recorded_invocation
4508 .as_ref()
4509 .context("LaunchSession invocation was not recorded")?;
4510 let result =
4511 (|| -> anyhow::Result<(LaunchSessionArguments, LaunchIntent)> {
4512 self.session.assert_tool_allowed(LAUNCH_SESSION_TOOL)?;
4513 let arguments = decode_launch_session_arguments(&call.arguments)?;
4514 let intent =
4515 self.session.prepare_launch_intent(invocation, &arguments)?;
4516 Ok((arguments, intent))
4517 })();
4518 match result {
4519 Ok((arguments, intent)) => {
4520 (self.checkpoint)(self.session.snapshot()?).await?;
4521 match self.session.lower_launch(&intent, &arguments).await {
4522 Ok(launch) => ToolOutcome {
4523 text: launch_success_json(
4524 &launch.session_id,
4525 &launch.command_id,
4526 )?,
4527 store_result: true,
4528 ok: true,
4529 end_session: false,
4530 freeform_write: None,
4531 managed_source_snapshot: None,
4532 exact_result: true,
4533 },
4534 Err(error)
4535 if matches!(
4536 error.kind,
4537 HistoryErrorKind::InvalidInput
4538 | HistoryErrorKind::Conflict
4539 ) =>
4540 {
4541 ToolOutcome {
4542 text: format!(
4543 "LaunchSession failed: {}",
4544 error.message
4545 ),
4546 store_result: true,
4547 ok: false,
4548 end_session: false,
4549 freeform_write: None,
4550 managed_source_snapshot: None,
4551 exact_result: false,
4552 }
4553 }
4554 Err(error) => {
4555 return Err(anyhow::anyhow!(
4556 "LaunchSession remains unresolved ({}): {}",
4557 error.kind.code(),
4558 error.message
4559 ));
4560 }
4561 }
4562 }
4563 Err(error) => ToolOutcome {
4564 text: format!("LaunchSession failed: {error}"),
4565 store_result: true,
4566 ok: false,
4567 end_session: false,
4568 freeform_write: None,
4569 managed_source_snapshot: None,
4570 exact_result: false,
4571 },
4572 }
4573 } else {
4574 match self.session.execute_tool(&call, operation_id).await {
4575 Ok(outcome) => {
4576 emitted_response = call.name == "EmitObject" && outcome.ok;
4577 outcome
4578 }
4579 Err(error) => ToolOutcome {
4580 text: format!("{} failed: {error}", call.name),
4581 store_result: call.name != "LoadNodes",
4582 ok: false,
4583 end_session: false,
4584 freeform_write: None,
4585 managed_source_snapshot: None,
4586 exact_result: false,
4587 },
4588 }
4589 }
4590 }
4591 Err(error) => ToolOutcome {
4592 text: error.to_string(),
4593 store_result: true,
4594 ok: false,
4595 end_session: false,
4596 freeform_write: None,
4597 managed_source_snapshot: None,
4598 exact_result: false,
4599 },
4600 };
4601 if let Some(snapshot) = outcome.managed_source_snapshot.take() {
4602 apply_snapshot(&mut self.session.journal, &now(), snapshot)?;
4603 outcome.store_result = false;
4604 }
4605 if !outcome.exact_result {
4606 append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
4607 }
4608 self.state.exact_tool_result = outcome.exact_result;
4609 let capture = if let Some(request) = outcome.freeform_write.take() {
4610 self.state.pending_freeform_write = Some(PendingFreeformWrite {
4611 request,
4612 call_box_id: created_call_box_id
4613 .context("freeform write call box was not created")?,
4614 });
4615 Some(json!(true))
4616 } else {
4617 None
4618 };
4619 if outcome.store_result {
4620 outcome.text = ensure_tool_result_box(
4621 &mut self.session.journal,
4622 recorded_invocation.as_ref(),
4623 &outcome.text,
4624 outcome.ok,
4625 )?;
4626 }
4627 let external_event_id = self.session.pending_external_event_id.clone();
4628 let recovery = self
4629 .session
4630 .recover_context_overflow(external_event_id.as_deref(), &[])?;
4631 let context_warning_added =
4632 self.session.transcript[transcript_start..]
4633 .iter()
4634 .any(|entry| {
4635 entry.get("contextOverflowWarning").and_then(Value::as_bool) == Some(true)
4636 });
4637 let mut provider_text = outcome.text.clone();
4638 if !outcome.exact_result
4639 && context_warning_added
4640 && !provider_text.contains(CONTEXT_OVERFLOW_WARNING)
4641 {
4642 if !provider_text.is_empty() {
4643 provider_text.push_str("\n\n");
4644 }
4645 provider_text.push_str(CONTEXT_OVERFLOW_WARNING);
4646 }
4647 self.session.record_tool_completion(
4648 recorded_invocation.as_ref(),
4649 json!({"ok":outcome.ok,"result":outcome.text}),
4650 )?;
4651 let stop = recovery == ContextRecovery::Irreducible
4652 || (matches!(self.session.mode, AgentMode::Ingress { .. })
4653 && self.session.ingress_force_commit_requested())
4654 || (!matches!(self.session.mode, AgentMode::Ingress { .. })
4655 && self.session.journal.state().source_terminated);
4656 Ok(kcode_agent_runtime::SessionToolOutcome {
4657 text: provider_text,
4658 ok: outcome.ok,
4659 capture,
4660 stop,
4661 finish_after_round: outcome.end_session,
4662 emitted_response,
4663 })
4664 })
4665 }
4666
4667 fn prepare_provider_resume<'a>(
4668 &'a mut self,
4669 mut outcome: kcode_agent_runtime::SessionToolOutcome,
4670 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ProviderResume> {
4671 Box::pin(async move {
4672 if completes_before_provider_resume(&outcome) {
4673 (self.checkpoint)(self.session.snapshot()?).await?;
4674 return Ok(kcode_agent_runtime::ProviderResume::Complete(None));
4675 }
4676
4677 let ingress_time = self
4678 .session
4679 .ingress_time_remaining()?
4680 .map(|remaining| (remaining, self.session.previous_ingress_attempt_timed_out));
4681 let synchronized_after = self
4682 .state
4683 .provider_synchronized_after
4684 .context("provider resume was prepared before its input was recorded")?;
4685 let prepared = self.session.journal.prepare_provider_resume(
4686 now(),
4687 synchronized_after,
4688 ingress_time,
4689 )?;
4690 match apply_prepared_provider_resume(
4691 &mut self.session.provider_affinity,
4692 &mut self.session.next_thread_reset_reason,
4693 prepared,
4694 ) {
4695 NativeProviderResumePreparation::Continue { marker_lines } => {
4696 self.state.provider_synchronized_after = Some(
4697 self.session
4698 .journal
4699 .state()
4700 .events
4701 .last()
4702 .context("provider resume preparation left no journal event")?
4703 .id,
4704 );
4705 if self.state.exact_tool_result {
4706 self.state.exact_tool_result = false;
4707 } else {
4708 let mut footer_lines = marker_lines;
4709 footer_lines.extend(self.session.runtime_budget().footer_lines());
4710 outcome.text = provider_tool_result_with_context_footer(
4711 &footer_lines.join("\n"),
4712 &outcome.text,
4713 );
4714 }
4715 (self.checkpoint)(self.session.snapshot()?).await?;
4716 Ok(kcode_agent_runtime::ProviderResume::Continue(outcome))
4717 }
4718 NativeProviderResumePreparation::RestartFresh { reason } => {
4719 self.state.exact_tool_result = false;
4720 self.state.restart_fresh_reason = Some(reason);
4721 (self.checkpoint)(self.session.snapshot()?).await?;
4722 Ok(kcode_agent_runtime::ProviderResume::RestartFresh)
4723 }
4724 }
4725 })
4726 }
4727
4728 fn complete_capture<'a>(
4729 &'a mut self,
4730 _capture: Value,
4731 contents: String,
4732 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
4733 Box::pin(async move {
4734 let pending = self
4735 .state
4736 .pending_freeform_write
4737 .take()
4738 .context("provider completed without a pending freeform write")?;
4739 let result_metadata = pending.request.clone();
4740 let outcome = self
4741 .session
4742 .complete_freeform_write(pending, contents)
4743 .await?;
4744 if outcome.store_result {
4745 self.session.journal.create_box(
4746 now(),
4747 "Kennedy tool result",
4748 BoxOwner::Controller,
4749 BoxContent::text(&outcome.text),
4750 )?;
4751 }
4752 self.session.journal.record(
4753 now(),
4754 EventKind::Note {
4755 label: "write_file_freeform_result".into(),
4756 value: result_metadata.result_record(outcome.ok, &outcome.text),
4757 },
4758 )?;
4759 let external_event_id = self.session.pending_external_event_id.clone();
4760 let recovery = self
4761 .session
4762 .recover_context_overflow(external_event_id.as_deref(), &[])?;
4763 let snapshot = self.session.snapshot()?;
4764 (self.checkpoint)(snapshot).await?;
4765 if recovery == ContextRecovery::Irreducible
4766 || (matches!(self.session.mode, AgentMode::Ingress { .. })
4767 && self.session.ingress_force_commit_requested())
4768 || (!matches!(self.session.mode, AgentMode::Ingress { .. })
4769 && self.session.journal.state().source_terminated)
4770 || self.state.deadline_after_response
4771 {
4772 return Ok(kcode_agent_runtime::SessionControl::Complete(None));
4773 }
4774 self.session.journal.create_box(
4775 now(),
4776 controller_box_name(&self.session.mode),
4777 BoxOwner::Controller,
4778 BoxContent::text(controller_message(
4779 &self.session.mode,
4780 &self.session.free_time,
4781 )),
4782 )?;
4783 Ok(kcode_agent_runtime::SessionControl::Continue)
4784 })
4785 }
4786
4787 fn complete_round<'a>(
4788 &'a mut self,
4789 completion: kcode_agent_runtime::RoundCompletion,
4790 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::SessionControl> {
4791 Box::pin(async move {
4792 let answer = completion.answer.trim().to_owned();
4793 let mut completion_recovery = ContextRecovery::NotNeeded;
4794 if !answer.is_empty() {
4795 let mut content = BoxContent::text(answer.clone());
4796 if let Some(id) = &self.session.pending_external_event_id {
4797 content.metadata["externalEventId"] = json!(id);
4798 }
4799 self.session.journal.create_box(
4800 now(),
4801 "Kennedy message",
4802 BoxOwner::Kennedy,
4803 content,
4804 )?;
4805 let mut transcript = json!({"role":"kennedy","content":answer});
4806 if let Some(id) = &self.session.pending_external_event_id {
4807 transcript["externalEventId"] = json!(id);
4808 }
4809 self.session.transcript.push(transcript);
4810 self.session.synchronize_provider_known_events();
4811 let external_event_id = self.session.pending_external_event_id.clone();
4812 completion_recovery = self
4813 .session
4814 .recover_context_overflow(external_event_id.as_deref(), &[])?;
4815 }
4816 let snapshot = self.session.snapshot()?;
4817 (self.checkpoint)(snapshot).await?;
4818 if completion_recovery == ContextRecovery::Irreducible
4819 || (matches!(self.session.mode, AgentMode::Ingress { .. })
4820 && self.session.ingress_force_commit_requested())
4821 || (!matches!(self.session.mode, AgentMode::Ingress { .. })
4822 && self.session.journal.state().source_terminated)
4823 {
4824 return Ok(kcode_agent_runtime::SessionControl::Complete(None));
4825 }
4826 if completion.finish_requested || self.state.deadline_after_response {
4827 return Ok(kcode_agent_runtime::SessionControl::Complete(
4828 (!answer.is_empty()).then_some(answer),
4829 ));
4830 }
4831 if matches!(self.session.mode, AgentMode::Conversation) && !answer.is_empty() {
4832 return Ok(kcode_agent_runtime::SessionControl::Complete(Some(answer)));
4833 }
4834 if matches!(self.session.mode, AgentMode::Conversation) && completion.emitted_response {
4835 return Ok(kcode_agent_runtime::SessionControl::Complete(None));
4836 }
4837 let solo_ingress_response =
4838 matches!(self.session.mode, AgentMode::Ingress { .. }) && !answer.is_empty();
4839 anyhow::ensure!(
4840 completion.used_tool || solo_ingress_response,
4841 "provider completed without a response or tool call"
4842 );
4843 self.session.journal.create_box(
4844 now(),
4845 controller_box_name(&self.session.mode),
4846 BoxOwner::Controller,
4847 BoxContent::text(controller_message(
4848 &self.session.mode,
4849 &self.session.free_time,
4850 )),
4851 )?;
4852 Ok(kcode_agent_runtime::SessionControl::Continue)
4853 })
4854 }
4855}
4856
4857impl kcode_agent_runtime::Host for KennedySubagentHost<'_> {
4858 fn render_tool_call(&mut self, call: &kcode_agent_runtime::ToolCall) -> anyhow::Result<String> {
4859 Ok(tool_invocation_content(&call.name, &call.arguments)?.text)
4860 }
4861
4862 fn execute_tool<'a>(
4863 &'a mut self,
4864 call: kcode_agent_runtime::ToolCall,
4865 operation_id: Uuid,
4866 budget: kcode_agent_runtime::ContextBudget,
4867 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
4868 Box::pin(async move {
4869 let call = ToolCall {
4870 name: call.name,
4871 arguments: call.arguments,
4872 };
4873 if let Some(reason) = subagent_unavailable_reason(&call.name) {
4874 return Ok(kcode_agent_runtime::ToolOutcome::failure(reason));
4875 }
4876 if budget.estimated_tokens() > budget.max_input_tokens() {
4877 return Ok(kcode_agent_runtime::ToolOutcome::failure(
4878 "The Ktool call was not run because its retained invocation would exceed the subagent context limit.",
4879 ));
4880 }
4881 if !subagent_managed_write_fits(&self.context, &call, &budget) {
4882 return Ok(kcode_agent_runtime::ToolOutcome::failure(
4883 "The managed-source write was not run because its resulting current state would exceed the subagent context limit.",
4884 ));
4885 }
4886
4887 let tool_started_at = std::time::Instant::now();
4888
4889 if call.name == "LoadNodes" {
4890 let Some(DecodedTool::LoadNodes(identifiers)) =
4891 decode(&call.name, &call.arguments)?
4892 else {
4893 return Ok(kcode_agent_runtime::ToolOutcome::failure(
4894 "LoadNodes did not match its tool contract.",
4895 ));
4896 };
4897 load_durable_batch(
4898 self.session.api.kmap(),
4899 self.context.kweb_mut(),
4900 &identifiers,
4901 )?;
4902 let (updates, creates) = self.session.plan.context_projection();
4903 let changes = self.context.reconcile_kweb(&updates, &creates)?;
4904 let displayed_state_keys = changes.displayed_state_keys();
4905 let mut text = if changes.is_empty() {
4906 "LoadNodes completed. The subagent Kweb projection was already current.".into()
4907 } else {
4908 changes.display_text()
4909 };
4910 append_slow_tool_duration(&mut text, tool_started_at.elapsed());
4911 return Ok(kcode_agent_runtime::ToolOutcome {
4912 text,
4913 ok: true,
4914 state_updates: changes.updates,
4915 displayed_state_keys,
4916 capture: None,
4917 });
4918 }
4919
4920 if is_kweb_mutation(&call.name) {
4921 self.session.assert_tool_allowed(&call.name)?;
4922 let decoded = decode(&call.name, &call.arguments)?
4923 .with_context(|| format!("{} did not match its tool contract", call.name))?;
4924 let prior_create_count = self.session.plan.create_count();
4925 let (mut text, referenced_pending) = execute_kweb_mutation(
4926 &call.name,
4927 decoded,
4928 self.context.kweb(),
4929 &mut self.session.plan,
4930 &mut self.session.journal,
4931 )?;
4932 self.context.include_staged_nodes(
4933 referenced_pending
4934 .into_iter()
4935 .chain(self.session.plan.pending_ids_from(prior_create_count)),
4936 );
4937 let (updates, creates) = self.session.plan.context_projection();
4938 let changes = self.context.reconcile_kweb(&updates, &creates)?;
4939 append_slow_tool_duration(&mut text, tool_started_at.elapsed());
4940 return Ok(kcode_agent_runtime::ToolOutcome {
4941 text,
4942 ok: true,
4943 state_updates: changes.updates,
4944 displayed_state_keys: Vec::new(),
4945 capture: None,
4946 });
4947 }
4948
4949 if let Some(request) = decode_freeform_write(&call.name, &call.arguments)? {
4950 if !self.context.source_is_open(request.kind(), request.name()) {
4951 return Ok(kcode_agent_runtime::ToolOutcome::failure(format!(
4952 "{} {:?} is not open in this subagent context. Call {} first.",
4953 request.kind().label(),
4954 request.name(),
4955 request.kind().open_tool()
4956 )));
4957 }
4958 let acknowledgement = request.acknowledgement();
4959 let id = Uuid::new_v4().to_string();
4960 self.captures.insert(id.clone(), request);
4961 return Ok(kcode_agent_runtime::ToolOutcome {
4962 text: acknowledgement,
4963 ok: true,
4964 state_updates: Vec::new(),
4965 displayed_state_keys: Vec::new(),
4966 capture: Some(Value::String(id)),
4967 });
4968 }
4969
4970 let mut outcome = match self.session.execute_tool(&call, operation_id).await {
4971 Ok(outcome) => outcome,
4972 Err(error) => {
4973 let mut text = format!("{} failed: {error}", call.name);
4974 append_slow_tool_duration(&mut text, tool_started_at.elapsed());
4975 return Ok(kcode_agent_runtime::ToolOutcome::failure(text));
4976 }
4977 };
4978 let displays_managed_snapshot = outcome
4979 .managed_source_snapshot
4980 .as_ref()
4981 .is_some_and(|snapshot| result_displays_snapshot(&outcome.text, snapshot));
4982 append_slow_tool_duration(&mut outcome.text, tool_started_at.elapsed());
4983 let (state_updates, displayed_state_keys) =
4984 if let Some(snapshot) = outcome.managed_source_snapshot.take() {
4985 let state = self.context.apply_source_snapshot(snapshot);
4986 let displayed = displays_managed_snapshot.then_some(state.key);
4987 (
4988 state.update.into_iter().collect(),
4989 displayed.into_iter().collect(),
4990 )
4991 } else {
4992 (Vec::new(), Vec::new())
4993 };
4994 let capture = outcome.freeform_write.take().map(|request| {
4995 let id = Uuid::new_v4().to_string();
4996 self.captures.insert(id.clone(), request);
4997 Value::String(id)
4998 });
4999 Ok(kcode_agent_runtime::ToolOutcome {
5000 text: outcome.text,
5001 ok: outcome.ok,
5002 state_updates,
5003 displayed_state_keys,
5004 capture,
5005 })
5006 })
5007 }
5008
5009 fn complete_capture<'a>(
5010 &'a mut self,
5011 capture: Value,
5012 contents: String,
5013 budget: kcode_agent_runtime::ContextBudget,
5014 ) -> kcode_agent_runtime::HostFuture<'a, kcode_agent_runtime::ToolOutcome> {
5015 Box::pin(async move {
5016 let id = capture
5017 .as_str()
5018 .context("subagent freeform capture token is invalid")?;
5019 let request = self
5020 .captures
5021 .remove(id)
5022 .context("subagent freeform capture token is unknown")?;
5023 self.session
5024 .complete_subagent_freeform_write(&mut self.context, request, contents, &budget)
5025 .await
5026 })
5027 }
5028
5029 fn record(&mut self, event: kcode_agent_runtime::AuditEvent) -> anyhow::Result<()> {
5030 kcode_intelligence_chatend::record_subagent_event(&mut self.session.journal, &now(), &event)
5031 }
5032}
5033
5034fn cost_summary(label: &str, estimated_cost_usd_nanos: u64, unpriced_calls: u64) -> String {
5035 render(RenderRequest::CostSummary {
5036 label,
5037 estimated_cost_usd_nanos,
5038 unpriced_calls,
5039 })
5040 .expect("cost-summary rendering is infallible")
5041}
5042
5043fn restore_kweb_context(journal: &HistorySession, context: &mut KwebContext) -> anyhow::Result<()> {
5044 let Some(tool) = journal.state().tools.get(KWEB_TOOL_INSTANCE) else {
5045 return Ok(());
5046 };
5047 let mut nodes = BTreeMap::new();
5048 for slot in &tool.slots {
5049 let state = journal
5050 .state()
5051 .box_state(slot.box_id)
5052 .context("Kweb slot references a missing box")?;
5053 if let Some(node) = state.canonical.content.metadata.get("storedNode") {
5054 let node = match serde_json::from_value::<KwebNode>(node.clone()) {
5055 Ok(node) => node,
5056 Err(_) => node_from_value(node).context("decoding a stored Kweb context node")?,
5057 };
5058 nodes.insert(node.id.clone(), node);
5059 }
5060 }
5061 let mut direct = journal
5062 .state()
5063 .current_ingress_attempt_events()
5064 .iter()
5065 .flat_map(|event| {
5066 let EventKind::ToolInvoked {
5067 tool_name,
5068 arguments,
5069 ..
5070 } = &event.kind
5071 else {
5072 return Vec::new();
5073 };
5074 match tool_name.as_str() {
5075 "LoadNodes" => arguments
5076 .get("identifiers")
5077 .and_then(Value::as_array)
5078 .into_iter()
5079 .flatten()
5080 .filter_map(Value::as_str)
5081 .map(str::to_owned)
5082 .collect(),
5083 "LoadNode" => arguments
5084 .get("identifier")
5085 .and_then(Value::as_str)
5086 .map(str::to_owned)
5087 .into_iter()
5088 .collect(),
5089 _ => Vec::new(),
5090 }
5091 })
5092 .collect::<Vec<_>>();
5093 if direct.is_empty() {
5094 direct = context.root_node_ids().to_vec();
5095 }
5096 context
5097 .restore(nodes.into_values(), direct)
5098 .map_err(anyhow::Error::new)
5099}
5100
5101fn session_kind(session_type: &str, mode: &AgentMode) -> SessionKind {
5102 if matches!(mode, AgentMode::Ingress { .. }) {
5103 return SessionKind::HistoryIngress;
5104 }
5105 match session_type {
5106 "conversation" => SessionKind::Conversation,
5107 "telegram" => SessionKind::Telegram,
5108 "telegram-group" => SessionKind::TelegramGroup,
5109 "free-time" => SessionKind::SelfTime,
5110 "wakeup" => SessionKind::Other("wakeup".into()),
5111 "audio" => SessionKind::AudioIngress,
5112 other => SessionKind::Other(other.into()),
5113 }
5114}
5115
5116fn tool_instance_for_invocation(name: &str, invocation_id: &str) -> String {
5117 if name == "LoadNodes" {
5118 return KWEB_TOOL_INSTANCE.into();
5119 }
5120 format!("{name}:{invocation_id}")
5121}
5122
5123fn canonical_id(value: &str) -> anyhow::Result<String> {
5124 value
5125 .parse::<NodeId>()
5126 .with_context(|| format!("{value:?} is not a canonical node ID"))?;
5127 Ok(value.into())
5128}
5129
5130fn image_extension(media_type: &str) -> &'static str {
5131 match media_type
5132 .split(';')
5133 .next()
5134 .unwrap_or(media_type)
5135 .trim()
5136 .to_ascii_lowercase()
5137 .as_str()
5138 {
5139 "image/jpeg" => "jpg",
5140 "image/webp" => "webp",
5141 _ => "png",
5142 }
5143}
5144
5145fn call_ktool_description(include_launch_session: bool) -> String {
5146 let mut description = render(RenderRequest::CallKtoolDescription)
5147 .expect("Ktool-description rendering is infallible");
5148 if include_launch_session {
5149 description.push_str(
5150 "\n\nLaunchSession is available for this genuine user turn. Call it with exactly directive (a nonblank string) and contextNodeIds (an ordered array of distinct fully loaded canonical node IDs). It launches an ordinary browser conversation and returns exactly sessionId and commandId.",
5151 );
5152 }
5153 description
5154}
5155
5156fn now() -> String {
5157 Utc::now().to_rfc3339()
5158}
5159
5160fn deadline(value: &Value) -> Option<DateTime<Utc>> {
5161 value
5162 .get("deadlineAt")
5163 .and_then(Value::as_str)
5164 .and_then(|value| DateTime::parse_from_rfc3339(value).ok())
5165 .map(|value| value.with_timezone(&Utc))
5166}
5167
5168fn remaining_until(deadline: DateTime<Utc>) -> Duration {
5169 (deadline - Utc::now()).to_std().unwrap_or(Duration::ZERO)
5170}
5171
5172fn controller_box_name(mode: &AgentMode) -> &'static str {
5173 match mode {
5174 AgentMode::Conversation => "Turn continuation",
5175 AgentMode::FreeTime => "Self-time continuation",
5176 AgentMode::Wakeup => "Wakeup continuation",
5177 AgentMode::Ingress { .. } => "History-ingress continuation",
5178 }
5179}
5180
5181fn controller_message(mode: &AgentMode, free_time: &Value) -> String {
5182 let mode = match mode {
5183 AgentMode::Conversation => "conversation",
5184 AgentMode::FreeTime => "free-time",
5185 AgentMode::Wakeup => "wakeup",
5186 AgentMode::Ingress { .. } => "ingress",
5187 };
5188 render(RenderRequest::ControllerMessage { mode, free_time })
5189 .expect("known controller modes render successfully")
5190}