1mod base;
2mod lease;
3mod recovery;
4mod turn;
5
6use std::ops::{Deref, DerefMut};
7use std::path::Path;
8use std::sync::Arc;
9
10use crate::config::LlmSettings;
11use crate::context::SkillEntry;
12use crate::model::{ChatMessage, OBSERVATION_ROLE};
13
14pub use base::{
15 sessions_dir, validate_session_id, CompactionRecord, InterruptionRecord, SessionHistoryRecord,
16 SessionMetadata, SessionToolResult,
17};
18pub use turn::{
19 TurnEvent, TurnKind, TurnLifecycleRecord, TurnOutcome, TurnPhase, TurnState, TurnStatus,
20};
21
22#[derive(Debug)]
23pub struct SessionError(String);
24
25impl SessionError {
26 fn new(message: impl Into<String>) -> Self {
27 Self(message.into())
28 }
29}
30
31impl std::fmt::Display for SessionError {
32 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
33 formatter.write_str(&self.0)
34 }
35}
36
37impl std::error::Error for SessionError {}
38
39impl From<base::SessionError> for SessionError {
40 fn from(error: base::SessionError) -> Self {
41 Self(error.to_string())
42 }
43}
44
45#[derive(Debug, Clone)]
48pub struct Session {
49 inner: base::Session,
50 _lease: Arc<lease::SessionLease>,
51}
52
53impl Deref for Session {
54 type Target = base::Session;
55
56 fn deref(&self) -> &Self::Target {
57 &self.inner
58 }
59}
60
61impl DerefMut for Session {
62 fn deref_mut(&mut self) -> &mut Self::Target {
63 &mut self.inner
64 }
65}
66
67impl Session {
68 fn wrap_created(home: &Path, inner: base::Session) -> Result<Self, SessionError> {
69 let lease =
70 Arc::new(lease::SessionLease::acquire(home, &inner.id).map_err(SessionError::new)?);
71 Ok(Self {
72 inner,
73 _lease: lease,
74 })
75 }
76
77 pub fn create(
78 home: &Path,
79 cwd: &Path,
80 boot_system_prompt: String,
81 llm: LlmSettings,
82 ) -> Result<Self, SessionError> {
83 let inner = base::Session::create(home, cwd, boot_system_prompt, llm)
84 .map_err(SessionError::from)?;
85 Self::wrap_created(home, inner)
86 }
87
88 pub fn create_with_secret(
89 home: &Path,
90 cwd: &Path,
91 boot_system_prompt: String,
92 llm: LlmSettings,
93 secret: Option<&str>,
94 ) -> Result<Self, SessionError> {
95 let inner = base::Session::create_with_secret(home, cwd, boot_system_prompt, llm, secret)
96 .map_err(SessionError::from)?;
97 Self::wrap_created(home, inner)
98 }
99
100 pub fn create_with_skills_and_secret(
101 home: &Path,
102 cwd: &Path,
103 boot_system_prompt: String,
104 llm: LlmSettings,
105 skills: Vec<SkillEntry>,
106 secret: Option<&str>,
107 ) -> Result<Self, SessionError> {
108 let inner = base::Session::create_with_skills_and_secret(
109 home,
110 cwd,
111 boot_system_prompt,
112 llm,
113 skills,
114 secret,
115 )
116 .map_err(SessionError::from)?;
117 Self::wrap_created(home, inner)
118 }
119
120 pub fn resume(home: &Path, id: &str) -> Result<Self, SessionError> {
121 Self::resume_with_secret(home, id, None)
122 }
123
124 pub fn resume_with_secret(
125 home: &Path,
126 id: &str,
127 external_secret: Option<&str>,
128 ) -> Result<Self, SessionError> {
129 base::validate_session_id(id).map_err(SessionError::from)?;
130 let lease = Arc::new(lease::SessionLease::acquire(home, id).map_err(SessionError::new)?);
131 recovery::recover_journals(home, id).map_err(SessionError::new)?;
132 let inner = base::Session::resume_with_secret(home, id, external_secret)
133 .map_err(SessionError::from)?;
134 Ok(Self {
135 inner,
136 _lease: lease,
137 })
138 }
139
140 pub fn list(home: &Path) -> Result<Vec<SessionMetadata>, SessionError> {
141 base::Session::list(home).map_err(Into::into)
142 }
143
144 pub fn list_with_secret(
145 home: &Path,
146 external_secret: Option<&str>,
147 ) -> Result<Vec<SessionMetadata>, SessionError> {
148 base::Session::list_with_secret(home, external_secret).map_err(Into::into)
149 }
150
151 pub fn provider_messages(&self) -> Vec<ChatMessage> {
152 crate::tool_pruning::prune_old_tool_outputs(&self.inner.provider_messages())
153 }
154
155 pub fn append_message(&mut self, message: ChatMessage) -> Result<(), SessionError> {
159 if message.role == "user" {
160 return self.append_user_or_retry(message);
161 }
162
163 let starts_background = message.role == OBSERVATION_ROLE
164 && self
165 .latest_turn()
166 .map_err(SessionError::new)?
167 .is_none_or(|turn| !turn.is_pending());
168 if starts_background {
169 self.start_turn(TurnKind::Background)
170 .map_err(SessionError::new)?;
171 }
172
173 self.inner
174 .append_message(message.clone())
175 .map_err(SessionError::from)?;
176
177 let Some(turn) = self
178 .latest_turn()
179 .map_err(SessionError::new)?
180 .filter(|turn| turn.is_pending())
181 else {
182 return Ok(());
183 };
184
185 if message.role == "assistant" {
186 if message.tool_calls.is_empty() {
187 self.finish_turn(&turn.turn_id, TurnOutcome::Completed, None)
188 .map_err(SessionError::new)?;
189 } else {
190 self.set_turn_phase(&turn.turn_id, TurnPhase::ExecutingTools)
191 .map_err(SessionError::new)?;
192 }
193 } else if message.role == "tool" {
194 self.set_turn_phase(&turn.turn_id, TurnPhase::ProviderStream)
195 .map_err(SessionError::new)?;
196 }
197 Ok(())
198 }
199
200 pub fn append_interruption(
201 &mut self,
202 interruption: InterruptionRecord,
203 ) -> Result<(), SessionError> {
204 self.inner
205 .append_interruption(interruption)
206 .map_err(SessionError::from)?;
207 if let Some(turn) = self.latest_turn().map_err(SessionError::new)? {
208 if turn.is_pending() {
209 self.finish_turn(&turn.turn_id, TurnOutcome::Interrupted, None)
210 .map_err(SessionError::new)?;
211 }
212 }
213 Ok(())
214 }
215
216 pub fn append_compaction(
217 &mut self,
218 summary: String,
219 first_kept_message: usize,
220 tokens_before: usize,
221 ) -> Result<(), SessionError> {
222 self.inner
223 .append_compaction(summary, first_kept_message, tokens_before)
224 .map_err(SessionError::from)?;
225 if let Some(turn) = self
226 .latest_turn()
227 .map_err(SessionError::new)?
228 .filter(|turn| turn.is_pending())
229 {
230 self.set_turn_phase(&turn.turn_id, TurnPhase::ProviderStream)
231 .map_err(SessionError::new)?;
232 }
233 Ok(())
234 }
235
236 fn append_user_or_retry(&mut self, message: ChatMessage) -> Result<(), SessionError> {
237 let pending = self
238 .latest_turn()
239 .map_err(SessionError::new)?
240 .filter(|turn| turn.is_pending());
241 let is_retry = message.content.as_deref().map(str::trim) == Some("!retry");
242
243 if let Some(turn) = pending {
244 if !is_retry {
245 return Err(SessionError::new(
246 "session has an unresolved turn; send !retry before starting a new message",
247 ));
248 }
249 if turn.status == TurnStatus::Active {
250 self.fail_turn_retryably(
251 &turn.turn_id,
252 "previous turn stopped before completion".to_owned(),
253 )
254 .map_err(SessionError::new)?;
255 }
256 self.resume_retryable_turn(&turn.turn_id)
257 .map_err(SessionError::new)?;
258 return Ok(());
259 }
260 if is_retry {
261 return Err(SessionError::new("session has no pending turn to retry"));
262 }
263
264 let turn_id = self.start_turn(TurnKind::User).map_err(SessionError::new)?;
265 if let Err(error) = self.inner.append_message(message) {
266 let _ = self.finish_turn(
267 &turn_id,
268 TurnOutcome::TerminalFailure,
269 Some(error.to_string()),
270 );
271 return Err(error.into());
272 }
273 Ok(())
274 }
275}
276
277#[cfg(test)]
278mod wrapper_tests {
279 use std::fs;
280 use std::sync::atomic::{AtomicU64, Ordering};
281
282 use super::*;
283
284 static COUNTER: AtomicU64 = AtomicU64::new(0);
285
286 fn session() -> (std::path::PathBuf, Session) {
287 let home = std::env::temp_dir().join(format!(
288 "lucy-session-wrapper-{}-{}",
289 std::process::id(),
290 COUNTER.fetch_add(1, Ordering::Relaxed)
291 ));
292 fs::create_dir(&home).expect("home");
293 let session = Session::create_with_secret(
294 &home,
295 &std::env::current_dir().expect("cwd"),
296 "prompt".to_owned(),
297 LlmSettings {
298 base_url: "http://localhost".to_owned(),
299 model: "model".to_owned(),
300 api_key_env: "LUCY_WRAPPER_TEST_KEY".to_owned(),
301 effort: None,
302 },
303 None,
304 )
305 .expect("session");
306 (home, session)
307 }
308
309 #[test]
310 fn retry_control_does_not_append_a_second_user_message() {
311 let (home, mut session) = session();
312 session
313 .append_message(ChatMessage::user("original".to_owned()))
314 .expect("original");
315 assert!(session
316 .append_message(ChatMessage::user("replacement".to_owned()))
317 .is_err());
318 session
319 .append_message(ChatMessage::user("!retry".to_owned()))
320 .expect("retry");
321 assert_eq!(session.messages.len(), 1);
322 assert_eq!(session.messages[0].content.as_deref(), Some("original"));
323 assert_eq!(
324 session.latest_turn().expect("state").expect("turn").status,
325 TurnStatus::Active
326 );
327 fs::remove_dir_all(home).expect("cleanup");
328 }
329
330 #[test]
331 fn final_assistant_message_completes_the_tracked_turn() {
332 let (home, mut session) = session();
333 session
334 .append_message(ChatMessage::user("work".to_owned()))
335 .expect("user");
336 session
337 .append_message(ChatMessage::assistant("done".to_owned(), Vec::new()))
338 .expect("assistant");
339 assert_eq!(
340 session.latest_turn().expect("state").expect("turn").status,
341 TurnStatus::Completed
342 );
343 fs::remove_dir_all(home).expect("cleanup");
344 }
345 #[test]
346 fn second_mutable_session_handle_is_rejected_until_the_lease_is_released() {
347 let (home, session) = session();
348 let id = session.id.clone();
349 let error = Session::resume(&home, &id).expect_err("second writer");
350 assert_eq!(error.to_string(), "session is already open for writing");
351 drop(session);
352 Session::resume(&home, &id).expect("writer after release");
353 fs::remove_dir_all(home).expect("cleanup");
354 }
355
356 #[test]
357 fn resume_recovers_only_an_unterminated_trailing_fragment() {
358 use std::io::Write;
359 let (home, session) = session();
360 let id = session.id.clone();
361 let path = session.path.clone();
362 drop(session);
363 let mut file = std::fs::OpenOptions::new()
364 .append(true)
365 .open(&path)
366 .expect("append partial record");
367 file.write_all(b"{\"record\":\"message\",\"timestamp\":9,\"message\":")
368 .expect("partial record");
369 file.sync_data().expect("partial checkpoint");
370
371 let resumed = Session::resume(&home, &id).expect("recover trailing fragment");
372 assert!(std::fs::read(&path).expect("transcript").ends_with(b"\n"));
373 assert!(home
374 .join(".lucy/recovery")
375 .read_dir()
376 .expect("evidence")
377 .next()
378 .is_some());
379 drop(resumed);
380 fs::remove_dir_all(home).expect("cleanup");
381 }
382
383 #[test]
384 fn complete_middle_corruption_remains_fatal() {
385 use std::io::Write;
386 let (home, session) = session();
387 let id = session.id.clone();
388 let path = session.path.clone();
389 drop(session);
390 let mut file = std::fs::OpenOptions::new()
391 .append(true)
392 .open(&path)
393 .expect("append corrupt record");
394 file.write_all(b"not-json\n").expect("corrupt record");
395 file.sync_data().expect("corrupt checkpoint");
396 assert!(Session::resume(&home, &id).is_err());
397 fs::remove_dir_all(home).expect("cleanup");
398 }
399 #[test]
400 fn pruned_provider_context_is_stable_across_resume_without_mutating_raw_history() {
401 let (home, mut session) = session();
402 session
403 .append_message(ChatMessage::user("run a large command".to_owned()))
404 .expect("user");
405 session
406 .append_message(ChatMessage::assistant(
407 "running".to_owned(),
408 vec![crate::model::ChatToolCall {
409 id: "large".to_owned(),
410 name: "cmd".to_owned(),
411 arguments: "{}".to_owned(),
412 }],
413 ))
414 .expect("assistant");
415 session
416 .append_message(ChatMessage::tool(
417 "large".to_owned(),
418 "cmd".to_owned(),
419 "x".repeat(100_000),
420 ))
421 .expect("tool");
422 let id = session.id.clone();
423 let first = session.provider_messages();
424 assert_eq!(
425 session
426 .messages
427 .last()
428 .and_then(|message| message.content.as_deref())
429 .map(str::len),
430 Some(100_000)
431 );
432 drop(session);
433
434 let resumed = Session::resume(&home, &id).expect("resume");
435 assert_eq!(resumed.provider_messages(), first);
436 assert_eq!(
437 resumed
438 .messages
439 .last()
440 .and_then(|message| message.content.as_deref())
441 .map(str::len),
442 Some(100_000)
443 );
444 drop(resumed);
445 fs::remove_dir_all(home).expect("cleanup");
446 }
447}