Skip to main content

lucy/
session.rs

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/// A compatibility wrapper around Lucy's append-only session. Every
46/// mutable handle owns the same process-shared, OS-backed writer lease.
47#[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    /// Append a semantic message while maintaining one explicit logical turn.
156    /// `!retry` is a control input: it resumes the existing turn without
157    /// becoming another provider-visible user message.
158    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}