use crate::config::{Action, Event, Workflow};
use crate::prompt::{ChildDispatchRequest, Deps, Error, child_dispatch, compactor};
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, 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,
cfg: &crate::config::CompactionConfig,
state: &compactor::checkpoint::CheckpointState,
git: &dyn crate::template::GitRunner,
) -> Result<Option<Option<String>>, Error> {
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> {
match action {
Action::Dispatch { role, .. } if role == compactor::COMPACTOR_ROLE => {
dispatch_compactor(workspace, agent_id, worktree, point, deps)
}
other => Err(Error::ActionUnsupported {
action: format!("{other:?}"),
event: Event::WorkerFlush.as_str(),
}),
}
}
fn dispatch_compactor(
workspace: &Path,
agent_id: &str,
worktree: &Path,
point: Option<&str>,
deps: &Deps<'_>,
) -> Result<(), Error> {
let goal = compactor::compactor_goal(worktree, agent_id)?;
let req = ChildDispatchRequest {
repo: workspace,
parent_branch: agent_id,
parent_worktree: worktree,
role: compactor::COMPACTOR_ROLE,
goal: &goal,
name: None,
fork_point: point,
pins: crate::prompt::PinnedDocs::none(),
};
child_dispatch::run_procedure(&req, deps.git, deps.clock, deps.id_gen, deps.launcher)
}
#[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_compactor: false,
}
}
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"), &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"), &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"), &c, &state(5), &RevGit(Some("abc"))).unwrap();
assert_eq!(p, Some(Some("abc".into())));
}
#[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"), &c, &state(5), &RevGit(None)).unwrap_err();
assert!(
matches!(
err,
Error::Git {
op: "compaction point rev-parse",
..
}
),
"{err:?}"
);
}
}