use crate::config::{Action, Event, Workflow};
use crate::prompt::{ChildDispatchRequest, Deps, Error, child_dispatch, compactor, procedure};
use std::path::Path;
pub(in crate::prompt::dispatch) fn run_flush(
workspace: &Path,
agent_id: &str,
worktree: &Path,
workflow: &Workflow,
deps: &Deps<'_>,
) -> Result<(), Error> {
let Some(cfg) = workflow.compaction.as_ref() else {
return Ok(());
};
let state = compactor::state(worktree, agent_id, deps.clock.now_unix(), false, deps.git)?;
if !compactor::due(Some(cfg), &state)? {
return Ok(());
}
let Some(point) = compaction_point(worktree, agent_id, cfg, &state, deps.git)? else {
return Ok(());
};
for action in flush_actions(workflow) {
execute_flush(
&action,
workspace,
agent_id,
worktree,
point.as_deref(),
deps,
)?;
}
Ok(())
}
fn compaction_point(
worktree: &Path,
agent_id: &str,
cfg: &crate::config::CompactionConfig,
state: &compactor::checkpoint::CheckpointState,
git: &dyn crate::template::GitRunner,
) -> Result<Option<Option<String>>, Error> {
if let Some(budget) = cfg.intermediate.keep_recent_tokens {
return Ok(compactor::checkpoint::tail::point(worktree, agent_id, budget, git)?.map(Some));
}
let keep = cfg.intermediate.keep_recent.unwrap_or(0);
if keep == 0 {
return Ok(Some(None));
}
if state.commits_since_checkpoint <= keep {
return Ok(None);
}
let rev = format!("HEAD~{keep}");
let sha = git
.run_capture(worktree, &["rev-parse", &rev])
.map_err(|source| Error::Git {
op: "compaction point rev-parse",
source,
})?;
Ok(Some(Some(sha.trim().to_string())))
}
fn flush_actions(workflow: &Workflow) -> Vec<Action> {
let bound = workflow.actions_for(Event::WorkerFlush);
if bound.is_empty() {
vec![Action::Dispatch {
role: compactor::COMPACTOR_ROLE.to_string(),
with: None,
mode: None,
}]
} else {
bound
}
}
fn execute_flush(
action: &Action,
workspace: &Path,
agent_id: &str,
worktree: &Path,
point: Option<&str>,
deps: &Deps<'_>,
) -> Result<(), Error> {
let unsupported = || Error::ActionUnsupported {
action: format!("{action:?}"),
event: Event::WorkerFlush.as_str(),
};
let Action::Dispatch { role, .. } = action else {
return Err(unsupported());
};
let Some(build) = procedure::checkpoint_goal(role) else {
return Err(unsupported());
};
let goal = build(worktree, agent_id)?;
dispatch_at_point(role, &goal, workspace, agent_id, worktree, point, deps)
}
fn dispatch_at_point(
role: &str,
goal: &str,
workspace: &Path,
agent_id: &str,
worktree: &Path,
point: Option<&str>,
deps: &Deps<'_>,
) -> Result<(), Error> {
let req = ChildDispatchRequest {
repo: workspace,
parent_branch: agent_id,
parent_worktree: worktree,
role,
goal,
name: None,
fork_point: point,
cwd: None,
pins: crate::prompt::PinnedDocs::none(),
};
child_dispatch::run_procedure(
&req,
deps.git,
deps.clock,
deps.id_gen,
deps.launcher,
deps.rng,
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::CompactionConfig;
use crate::prompt::compactor::checkpoint::CheckpointState;
use crate::template::GitRunner;
use std::path::PathBuf;
fn cfg(yaml: &str) -> CompactionConfig {
serde_yaml_ng::from_str(yaml).unwrap()
}
fn state(commits: u32) -> CheckpointState {
CheckpointState {
commits_since_checkpoint: commits,
seconds_since_checkpoint: 0,
flush_requested: false,
is_checkpoint_child: false,
compaction_in_flight: false,
last_usage: None,
}
}
struct RevGit(Option<&'static str>);
impl GitRunner for RevGit {
fn run(&self, _d: &Path, _a: &[&str]) -> std::io::Result<()> {
unreachable!("compaction_point only captures")
}
fn run_capture(&self, _d: &Path, args: &[&str]) -> std::io::Result<String> {
assert_eq!(args[0], "rev-parse");
self.0
.map(|s| format!("{s}\n"))
.ok_or_else(|| std::io::Error::other("boom"))
}
}
#[test]
fn no_retained_tail_is_the_tip() {
let c = cfg("intermediate:\n trigger: every_n_commits\n n: 3\n");
let p = compaction_point(&PathBuf::from("/x"), "p1", &c, &state(5), &RevGit(None)).unwrap();
assert_eq!(p, Some(None));
}
#[test]
fn a_span_inside_the_retained_tail_skips_the_flush() {
let c = cfg("intermediate:\n trigger: every_t_seconds\n n: 1\n keep_recent: 4\n");
let p = compaction_point(&PathBuf::from("/x"), "p1", &c, &state(4), &RevGit(None)).unwrap();
assert_eq!(p, None);
}
#[test]
fn a_retained_tail_puts_the_point_behind_the_tip() {
let c = cfg("intermediate:\n trigger: every_n_commits\n n: 3\n keep_recent: 2\n");
let p = compaction_point(
&PathBuf::from("/x"),
"p1",
&c,
&state(5),
&RevGit(Some("abc")),
)
.unwrap();
assert_eq!(p, Some(Some("abc".into())));
}
#[test]
fn a_token_tail_takes_the_point_from_the_usage_walk() {
let c = cfg("intermediate:\n trigger: on_flush\n keep_recent_tokens: 20000\n");
let dir = tempfile::TempDir::new().unwrap();
let p = compaction_point(dir.path(), "p1", &c, &state(5), &RevGit(None)).unwrap();
assert_eq!(p, None);
}
#[test]
fn a_rev_parse_failure_surfaces_as_git_error() {
let c = cfg("intermediate:\n trigger: every_n_commits\n n: 3\n keep_recent: 2\n");
let err =
compaction_point(&PathBuf::from("/x"), "p1", &c, &state(5), &RevGit(None)).unwrap_err();
assert!(
matches!(
err,
Error::Git {
op: "compaction point rev-parse",
..
}
),
"{err:?}"
);
}
}