Luigit
repositories / smith

smith

There are many coding harnesses - but this one is fast

owned by admin

smith-core/src/lib.rs

Raw
//! Provider-independent Smith agent core (`SMH-SPEC-SPEC0001`, Sessions).
//!
//! Session persistence, CBOR framing, configurable runtime, replay traces,
//! and tool lifecycle events.

#![forbid(unsafe_code)]

/// Provider-independent agent loop.
pub mod agent;
/// Bounded, cancellable shell execution.
pub mod bash;
/// Checksummed session framing.
pub mod frame;
/// Branching session state.
pub mod session;
/// Durable session storage.
pub mod store;
/// Tool registry, validation, and recorded effects.
pub mod tools;
/// Provider-free replay evidence.
pub mod trace;

pub use agent::{
    Agent, AgentEvent, CompactionSummary, Delivery, QueuedInput, SecretProxy, TurnOutcome,
};
pub use bash::{BashResult, execute};
pub use frame::{Frame, Recovery, encode_entry, read_frames, write_frame};
pub use session::{
    Branch, COMPACTION_KIND, Entry, EntryContent, EntryContentRef, EntryFrame, EntryKind, Session,
    SessionHeader, Span,
};
pub use store::{
    OpenSession, SessionWriter, create_session, fork_session, load_session, open_session,
    repair_session,
};
pub use tools::{ToolOutcome, ToolRegistry, ToolSession};
pub use trace::{TraceKind, TraceRecord, replay};

/// Readiness signals raised by real child processes, for tests that must
/// act while a command runs.
#[cfg(test)]
#[expect(
    clippy::unwrap_used,
    reason = "tests may panic on invariant violations"
)]
mod process_signal {
    use std::net::{TcpListener, TcpStream};
    use std::sync::mpsc::{self, Receiver};
    use std::time::Duration;

    /// Upper bound on any single wait for a child; generous for loaded CI.
    const DEADLINE: Duration = Duration::from_secs(10);

    /// A loopback listener a shell command connects to once it runs.
    pub struct ProcessSignal {
        port: u16,
        connected: Receiver<TcpStream>,
    }

    impl ProcessSignal {
        pub fn new() -> Self {
            let listener = TcpListener::bind("127.0.0.1:0").unwrap();
            let port = listener.local_addr().unwrap().port();
            let (tx, connected) = mpsc::sync_channel(1);
            std::thread::spawn(move || {
                if let Ok((socket, _)) = listener.accept() {
                    let _ = tx.send(socket);
                }
            });
            Self { port, connected }
        }

        /// Bash that raises the signal and keeps the connection open for
        /// as long as the raising process lives.
        pub fn raise(&self) -> String {
            format!("exec 3<>/dev/tcp/127.0.0.1/{}", self.port)
        }

        /// The raising process's connection, once it has connected.
        pub fn wait(&self) -> TcpStream {
            let socket = self.connected.recv_timeout(DEADLINE);
            assert!(socket.is_ok(), "no process connected within {DEADLINE:?}");
            socket.unwrap()
        }
    }

    /// Whether every process holding `socket` has exited, as seen by end of
    /// stream before `DEADLINE`.
    pub fn closed(socket: &mut TcpStream) -> bool {
        use std::io::Read;
        socket.set_read_timeout(Some(DEADLINE)).unwrap();
        matches!(socket.read(&mut [0u8; 1]), Ok(0))
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use smith::message::{Message, Role};

    #[test]
    fn session_creates_and_appends() {
        let mut session = Session::new();
        let msg = Message::with_text(Role::User, "hello");
        let entry = *session.append(&EntryContent::Message(msg)).unwrap();
        assert_eq!(session.active_branch().entries(), vec![entry]);
        assert_eq!(session.selected_id(), Some(entry.id));
    }
}