use crate::types::*;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use yoagent_state::{
ActorRef, GitEventStore, Goal, NodeId, YoAgentModelCalled, YoAgentModelFinished,
YoAgentRunFinished, YoAgentRunStarted, YoAgentState, YoAgentStateAdapter, YoAgentStateSink,
YoAgentToolCalled, YoAgentToolFinished,
};
pub use yoagent_state::{GoalId, RunId, StateError};
#[derive(Debug, Clone)]
pub enum GoalRef {
Existing(GoalId),
New { title: String },
}
type Summarizer = std::sync::Arc<dyn Fn(&str) -> String + Send + Sync>;
pub struct GaspRecorder {
state: YoAgentState<GitEventStore>,
store: GitEventStore,
actor: ActorRef,
goal: GoalId,
summarize: Summarizer,
}
impl GaspRecorder {
pub async fn init(
root: impl AsRef<std::path::Path>,
agent_id: &str,
worker_id: &str,
goal: GoalRef,
) -> Result<Self, StateError> {
let store = yoagent_state::init_agent_repo(root, agent_id, worker_id)?;
Self::with_store(store, agent_id, goal).await
}
pub async fn open(
root: impl Into<std::path::PathBuf>,
agent_id: &str,
worker_id: &str,
goal: GoalRef,
) -> Result<Self, StateError> {
let store = GitEventStore::open(root, worker_id)?;
Self::with_store(store, agent_id, goal).await
}
async fn with_store(
store: GitEventStore,
agent_id: &str,
goal: GoalRef,
) -> Result<Self, StateError> {
let actor = ActorRef::agent(agent_id);
let state = YoAgentState::load(store.clone()).await?;
if let Some(stale) = state.resume_open_run().await? {
tracing::warn!(run = %stale, "closing stale open run as interrupted");
state
.record_run_finished(actor.clone(), stale, "interrupted")
.await
.map_err(|e| {
StateError::Validation(format!(
"found an open run and could not close it ({e}); if another \
worker is live on this repo, do not share it — GASP repos \
are single-writer"
))
})?;
}
let goal = match goal {
GoalRef::Existing(id) => {
if state.get_node(NodeId::new(id.as_str())).await.is_none() {
return Err(StateError::Validation(format!(
"goal {id} does not exist in this repo's graph"
)));
}
id
}
GoalRef::New { title } => {
let id = GoalId::generate();
state
.record_goal(Goal::new(id.clone(), title.clone(), title, actor.clone()))
.await?;
id
}
};
commit_scaffolding(&store)?;
Ok(Self {
state,
store,
actor,
goal,
summarize: std::sync::Arc::new(|text: &str| summarize(text)),
})
}
pub fn with_summarizer(
mut self,
summarize: impl Fn(&str) -> String + Send + Sync + 'static,
) -> Self {
self.summarize = std::sync::Arc::new(summarize);
self
}
pub fn goal(&self) -> &GoalId {
&self.goal
}
pub fn recording_sender(
&self,
task: impl Into<String>,
forward: Option<mpsc::UnboundedSender<AgentEvent>>,
) -> (
mpsc::UnboundedSender<AgentEvent>,
JoinHandle<Result<Option<RunId>, StateError>>,
) {
let (tx, rx) = mpsc::unbounded_channel();
let sink = YoAgentStateAdapter::new(self.state.clone(), self.actor.clone());
let store = self.store.clone();
let goal = self.goal.clone();
let task = task.into();
let summarize = self.summarize.clone();
let handle = tokio::spawn(consume(sink, store, goal, task, summarize, rx, forward));
(tx, handle)
}
}
struct RunTracking {
run_id: RunId,
task: String,
started: bool,
finished: bool,
turn: usize,
outcome: String,
}
async fn consume(
sink: YoAgentStateAdapter<GitEventStore>,
store: GitEventStore,
goal: GoalId,
task: String,
summarize: Summarizer,
mut rx: mpsc::UnboundedReceiver<AgentEvent>,
forward: Option<mpsc::UnboundedSender<AgentEvent>>,
) -> Result<Option<RunId>, StateError> {
let mut tracking = RunTracking {
run_id: RunId::generate(),
task,
started: false,
finished: false,
turn: 0,
outcome: "interrupted".to_string(),
};
let mut recording_error: Option<StateError> = None;
while let Some(event) = rx.recv().await {
if let Some(fwd) = &forward {
let _ = fwd.send(event.clone());
}
if recording_error.is_some() {
continue; }
if let Err(e) = record_event(&sink, &summarize, &mut tracking, &event).await {
tracing::error!(
run = %tracking.run_id,
error = %e,
"GASP recording failed; recording stops but event forwarding continues"
);
recording_error = Some(e);
}
}
if let Some(e) = recording_error {
let _ = store.release_lease();
return Err(e);
}
if !tracking.started {
return Ok(None);
}
if !tracking.finished {
sink.on_run_finished(YoAgentRunFinished {
run_id: tracking.run_id.clone(),
outcome: tracking.outcome.clone(),
metadata: serde_json::json!({}),
})
.await?;
}
store.commit_run(&tracking.run_id, &goal, &tracking.outcome, &[])?;
let _ = store.release_lease();
Ok(Some(tracking.run_id))
}
async fn record_event(
sink: &YoAgentStateAdapter<GitEventStore>,
summarize: &Summarizer,
tracking: &mut RunTracking,
event: &AgentEvent,
) -> Result<(), StateError> {
match event {
AgentEvent::AgentStart => {
sink.on_run_started(YoAgentRunStarted {
run_id: tracking.run_id.clone(),
task: tracking.task.clone(),
metadata: serde_json::json!({}),
})
.await?;
tracking.started = true;
}
AgentEvent::MessageEnd {
message:
AgentMessage::Llm(Message::Assistant {
content,
model,
stop_reason,
..
}),
} if tracking.started => {
tracking.turn += 1;
sink.on_model_called(YoAgentModelCalled {
run_id: tracking.run_id.clone(),
model: model.clone(),
prompt_summary: if tracking.turn == 1 {
summarize(&tracking.task)
} else {
format!("turn {}", tracking.turn)
},
})
.await?;
let text = content
.iter()
.find_map(|c| match c {
Content::Text { text } if !text.is_empty() => Some(text.as_str()),
_ => None,
})
.unwrap_or("(no text)");
sink.on_model_finished(YoAgentModelFinished {
run_id: tracking.run_id.clone(),
model: model.clone(),
output_summary: summarize(text),
})
.await?;
tracking.outcome = outcome_for(stop_reason).to_string();
}
AgentEvent::ToolExecutionStart {
tool_name, args, ..
} if tracking.started => {
sink.on_tool_called(YoAgentToolCalled {
run_id: tracking.run_id.clone(),
tool: tool_name.clone(),
input_summary: summarize(&args.to_string()),
})
.await?;
}
AgentEvent::ToolExecutionEnd {
tool_name,
result,
is_error,
..
} if tracking.started => {
let text = result
.content
.iter()
.find_map(|c| match c {
Content::Text { text } => Some(text.as_str()),
_ => None,
})
.unwrap_or("(no output)");
sink.on_tool_finished(YoAgentToolFinished {
run_id: tracking.run_id.clone(),
tool: tool_name.clone(),
output_summary: summarize(text),
success: !is_error,
})
.await?;
}
AgentEvent::InputRejected { .. } if tracking.started => {
tracking.outcome = "rejected".to_string();
}
AgentEvent::AgentEnd { .. } if tracking.started && !tracking.finished => {
sink.on_run_finished(YoAgentRunFinished {
run_id: tracking.run_id.clone(),
outcome: tracking.outcome.clone(),
metadata: serde_json::json!({}),
})
.await?;
tracking.finished = true;
}
_ => {}
}
Ok(())
}
fn commit_scaffolding(store: &GitEventStore) -> Result<(), StateError> {
let events = store.events_path();
let root = events
.parent()
.and_then(|p| p.parent())
.ok_or_else(|| StateError::Store("events path has no repo root".into()))?
.to_path_buf();
let run = |args: &[&str]| -> Result<std::process::Output, StateError> {
std::process::Command::new("git")
.args(args)
.current_dir(&root)
.output()
.map_err(|e| StateError::Store(format!("git {}: {e}", args.join(" "))))
};
run(&[
"add",
"--",
"AGENT.md",
"identity",
".gitignore",
"state/events.jsonl",
])?;
let staged = run(&["diff", "--cached", "--quiet"])?;
if !staged.status.success() {
let out = run(&["commit", "-q", "-m", "gasp: agent scaffolding"])?;
if !out.status.success() {
return Err(StateError::Store(format!(
"scaffolding commit failed: {}",
String::from_utf8_lossy(&out.stderr)
)));
}
}
Ok(())
}
fn summarize(text: &str) -> String {
let one_line = text.split_whitespace().collect::<Vec<_>>().join(" ");
if one_line.chars().count() <= 200 {
one_line
} else {
let truncated: String = one_line.chars().take(200).collect();
format!("{truncated}…")
}
}
fn outcome_for(stop_reason: &StopReason) -> &'static str {
match stop_reason {
StopReason::Stop | StopReason::ToolUse => "completed",
StopReason::Length => "truncated",
StopReason::Error => "error",
StopReason::Aborted => "aborted",
StopReason::Refusal => "refused",
}
}
#[cfg(test)]
mod tests {
use super::summarize;
#[test]
fn summarize_collapses_and_truncates_on_char_boundaries() {
assert_eq!(summarize("a\nb\t c"), "a b c");
let long: String = "ö".repeat(300);
let s = summarize(&long);
assert_eq!(s.chars().count(), 201);
assert!(s.ends_with('…'));
let exact: String = "x".repeat(200);
assert_eq!(summarize(&exact), exact);
}
}