kcode_k1_chat_thread_durable_state/
lib.rs1#![forbid(unsafe_code)]
2
3use kcode_k1_access_kmap::K1AccessKmap;
4use kcode_k1_chat_persistence::Session;
5pub use kcode_k1_chat_state::BoxId;
6use kcode_k1_chat_state::USER_MESSAGE_TYPE;
7use kcode_k1_chat_thread_actions::ChatThreadActions;
8pub use kcode_k1_chat_thread_actions::{AccessContext, AccessPolicy, ProfileId};
9pub use kcode_k1_chat_thread_durable_turn::{
10 BoxValue, ChatBox, PreparedCall, PreparedSteer, Status, ToolCallId,
11};
12use kcode_k1_chat_thread_durable_turn::{DurableTurn, RestartError, ShimOutput};
13use std::sync::Arc;
14
15#[derive(Clone, Debug, Eq, PartialEq)]
16pub enum TransitionError {
17 Unauthorized,
18 NotStalled,
19 NotRestartable,
20 Internal(String),
21}
22
23pub struct DurableThread {
24 turn: DurableTurn,
25 actions: ChatThreadActions,
26 authorized: bool,
27}
28
29impl DurableThread {
30 pub fn recover(session: Session, kmap: Arc<K1AccessKmap>) -> Result<Self, String> {
31 Ok(Self {
32 turn: DurableTurn::recover(session)?,
33 actions: ChatThreadActions::new(kmap),
34 authorized: false,
35 })
36 }
37
38 pub fn boxes(&self) -> &[ChatBox] {
39 self.turn.boxes()
40 }
41
42 pub fn status(&self) -> Status {
43 self.turn.status()
44 }
45
46 pub fn accept_box(
47 &mut self,
48 box_type: String,
49 contents: String,
50 hidden_type: String,
51 hidden_contents: String,
52 ) -> Result<(), String> {
53 self.turn
54 .accept(box_type, contents, hidden_type, hidden_contents)
55 }
56
57 pub fn accept_user(
58 &mut self,
59 context: AccessContext,
60 profile_id: ProfileId,
61 policy: AccessPolicy,
62 contents: String,
63 ) -> Result<(), TransitionError> {
64 let installed = self.bind_authorization(context, profile_id, policy)?;
65 match self.turn.accept(
66 USER_MESSAGE_TYPE.into(),
67 contents,
68 String::new(),
69 String::new(),
70 ) {
71 Ok(()) => Ok(()),
72 Err(error) => {
73 if installed {
74 self.clear_authorization();
75 }
76 Err(TransitionError::Internal(error))
77 }
78 }
79 }
80
81 pub fn accept_return(
82 &mut self,
83 id: ToolCallId,
84 result: Result<String, String>,
85 ) -> Result<(), String> {
86 self.turn.accept_tool_return(id, result)
87 }
88
89 pub fn prepare_stage(
90 &mut self,
91 job: u64,
92 text: String,
93 boxes: Vec<BoxValue>,
94 ) -> Result<Vec<PreparedCall>, String> {
95 self.turn.prepare_stage(job, text, boxes)
96 }
97
98 pub fn launch_action(&mut self, name: &str, arguments: &str) -> Result<String, String> {
99 self.actions.launch(name, arguments)
100 }
101
102 pub fn accept_tool_message(&mut self, id: ToolCallId, contents: String) -> Result<(), String> {
103 self.turn.accept_tool_message(id, contents)
104 }
105
106 pub fn accept_tool_return(
107 &mut self,
108 id: ToolCallId,
109 result: Result<String, String>,
110 ) -> Result<(), String> {
111 self.turn.accept_tool_return(id, result)
112 }
113
114 pub fn accept_tool_return_v2(
115 &mut self,
116 id: ToolCallId,
117 result: Result<String, String>,
118 metadata_type: String,
119 metadata_contents: String,
120 ) -> Result<(), String> {
121 self.turn
122 .accept_tool_return_v2(id, result, metadata_type, metadata_contents)
123 }
124
125 pub fn prepare_steer(&mut self, job: u64) -> Result<Option<PreparedSteer>, String> {
126 self.turn.prepare_steer(job)
127 }
128
129 pub fn prepared_input(&self, prepared: &PreparedSteer) -> Result<String, String> {
130 self.turn.validate_steer(prepared)?;
131 render_input(prepared.values())
132 }
133
134 pub fn commit_steer(&mut self, prepared: PreparedSteer) -> Result<(), String> {
135 self.turn.commit_steer(prepared)
136 }
137
138 pub fn begin_input(&mut self) -> Result<Option<(u64, String)>, String> {
139 let Some(start) = self.turn.begin()? else {
140 return Ok(None);
141 };
142 Ok(Some((start.job, render_input(&start.values)?)))
143 }
144
145 pub fn complete(&mut self, job: u64, output: ShimOutput<BoxValue>) -> Result<bool, String> {
146 self.turn.complete(job, output)
147 }
148
149 pub fn fail(&mut self, job: u64, error: String, restartable: bool) {
150 self.turn.fail(job, error, restartable);
151 self.clear_authorization();
152 }
153
154 pub fn restart(
155 &mut self,
156 context: AccessContext,
157 profile_id: ProfileId,
158 policy: AccessPolicy,
159 ) -> Result<(), TransitionError> {
160 let installed = self.bind_authorization(context, profile_id, policy)?;
161 if let Err(error) = self.turn.restart().map_err(|error| match error {
162 RestartError::NotStalled => TransitionError::NotStalled,
163 RestartError::ProviderActionAccepted => TransitionError::NotRestartable,
164 }) {
165 if installed {
166 self.clear_authorization();
167 }
168 return Err(error);
169 }
170 Ok(())
171 }
172
173 pub fn clear_authorization(&mut self) {
174 self.actions.clear_authorization();
175 self.authorized = false;
176 }
177
178 fn bind_authorization(
179 &mut self,
180 context: AccessContext,
181 profile_id: ProfileId,
182 policy: AccessPolicy,
183 ) -> Result<bool, TransitionError> {
184 let installed = !self.authorized;
185 if self
186 .actions
187 .bind_authorization(context, profile_id, policy)
188 .is_err()
189 {
190 if installed {
191 self.actions.clear_authorization();
192 }
193 return Err(TransitionError::Unauthorized);
194 }
195 self.authorized = true;
196 Ok(installed)
197 }
198}
199
200fn render_input(values: &[BoxValue]) -> Result<String, String> {
201 let mut output = String::new();
202 for value in values {
203 let BoxValue::History(section) = value else {
204 return Err("Codex provider input contains a non-history value".into());
205 };
206 if section.is_empty() {
207 continue;
208 }
209 if !output.is_empty() && !output.ends_with('\n') {
210 output.push('\n');
211 }
212 output.push_str(section);
213 }
214 Ok(output)
215}