use crate::daemon::{now, Daemon};
use autofork_core::config::Config;
use autofork_core::frontmatter::{ForkParse, ForkRunOn};
use autofork_core::moments::{match_moments, ForkMoment};
use autofork_core::protocol::WakeFork;
use autofork_core::schedule::{effective_priorities, resolve_deps, Selected};
use autofork_core::store::SessionRow;
use autofork_core::tags::tags_allowed;
use autofork_core::wake::{
build_release_payload, build_wake_forks, build_wake_payload, DueFork, HeldFork,
};
use std::path::{Path, PathBuf};
use std::sync::Arc;
#[derive(Clone)]
pub struct SelectedFork {
pub name: String,
pub path: PathBuf,
pub trigger: String,
pub overlap: bool,
pub after: Vec<String>,
pub priority: i64,
pub tags: Vec<String>,
pub chain: bool,
pub gate: bool,
pub model: autofork_core::frontmatter::ClientScoped,
pub mode: autofork_core::frontmatter::ClientScoped,
pub latch_key: Option<String>,
pub detail: Option<String>,
pub consumes: Option<(String, String)>,
}
fn latch_key_for(trigger: &ForkRunOn, pause_epoch: i64) -> Option<String> {
match trigger {
ForkRunOn::Idle { .. } => Some(format!("idle-pause:{pause_epoch}")),
ForkRunOn::ContextTokens(_) | ForkRunOn::ContextUsedPct(_) | ForkRunOn::ContextLeft(_) => {
Some(trigger.label())
}
ForkRunOn::Every { .. } => None,
ForkRunOn::Changed { .. } | ForkRunOn::Event { .. } => None,
_ => None,
}
}
impl Selected for SelectedFork {
fn name(&self) -> &str {
&self.name
}
fn after(&self) -> Vec<&str> {
self.after.iter().map(|a| a.as_str()).collect()
}
fn priority(&self) -> i64 {
self.priority
}
}
pub fn refresh_roster(daemon: &Arc<Daemon>, session_id: &str, cwd: &Path) {
let (entries, _) = autofork_core::discovery::discover_forks(
cwd,
Some(&daemon.user_forks_root()),
daemon.claude_dir().as_deref(),
daemon.agents_dir().as_deref(),
);
let store = daemon.store.lock().unwrap();
let t = now();
for entry in entries {
if let Ok(true) = store.queue_fork(session_id, &entry.name, &entry.path, t) {
tracing::info!(fork = %entry.name, session = session_id, "fork rostered");
}
}
}
pub fn select_forks(
daemon: &Arc<Daemon>,
session: &SessionRow,
cfg: &Config,
moments: &[ForkMoment],
) -> Vec<SelectedFork> {
refresh_roster(daemon, &session.session_id, &session.cwd);
let roster = {
let store = daemon.store.lock().unwrap();
store.roster(&session.session_id).unwrap_or_default()
};
let effective_enable = session
.enable_tags
.as_deref()
.or(cfg.enable_tags.as_deref());
let effective_disable = session
.disable_tags
.as_deref()
.or(cfg.disable_tags.as_deref());
let pending: Vec<(String, String, String)> = {
let store = daemon.store.lock().unwrap();
store
.pending_triggers(&session.session_id)
.unwrap_or_default()
};
let mut selected: Vec<SelectedFork> = Vec::new();
let t = now();
for entry in roster {
let Ok(content) = std::fs::read_to_string(&entry.fork_path) else {
continue;
};
let ForkParse::Fork(parsed) = parse_fork(&entry.fork_name, &content) else {
continue;
};
if !tags_allowed(&parsed.def.tags, effective_enable, effective_disable) {
continue;
}
let Some(trigger) = match_moments(
&parsed.def,
moments,
cfg.default_idle_deadline_secs,
entry.ran_at,
session.created_at,
) else {
continue;
};
if let (Some(throttle), Some(ran_at)) = (parsed.def.throttle_secs, entry.ran_at) {
if (t - ran_at).max(0) < throttle as i64 {
tracing::debug!(fork = %entry.fork_name, "throttled, skipping");
continue;
}
}
if !parsed.def.tags.is_empty() && !cfg.tag_throttles.is_empty() {
let store = daemon.store.lock().unwrap();
let mut hit = None;
for tag in &parsed.def.tags {
let Some(&window) = cfg.tag_throttles.get(tag) else {
continue;
};
if let Ok(Some(last)) =
store.last_run_for_tags(&session.project_root, std::slice::from_ref(tag))
{
if (t - last).max(0) < window as i64 {
hit = Some(tag.clone());
break;
}
}
}
drop(store);
if let Some(tag) = hit {
tracing::debug!(fork = %entry.fork_name, %tag, "tag-throttled, skipping");
continue;
}
}
let label = trigger.label();
let latch_key = latch_key_for(&trigger, session.pause_epoch);
if let Some(key) = &latch_key {
let latched = {
let store = daemon.store.lock().unwrap();
store
.is_latched(&session.session_id, &entry.fork_name, key)
.unwrap_or(false)
};
if latched {
continue;
}
}
if !parsed.def.overlap {
let live = {
let store = daemon.store.lock().unwrap();
store
.live_spawn_newer_than(
&session.session_id,
&entry.fork_name,
t - overlap_spawn_max_age_secs(),
)
.unwrap_or(false)
};
if live {
tracing::info!(session = %session.session_id, fork = %entry.fork_name,
"a run is still in flight and overlap is false, skipping");
continue;
}
}
if cfg.runaway_limit > 0 && !matches!(trigger, ForkRunOn::Every { .. }) {
let recent = {
let store = daemon.store.lock().unwrap();
store
.count_runs_since(
&session.session_id,
&entry.fork_name,
t - runaway_window_secs(),
)
.unwrap_or(0)
};
if recent >= cfg.runaway_limit as i64 {
tracing::warn!(session = %session.session_id, fork = %entry.fork_name,
runs = recent, limit = cfg.runaway_limit,
"runaway breaker: fork hit its hourly run cap, skipping \
(raise `runaway_limit` in config if this rate is intended)");
continue;
}
}
let (detail, consumes) = match &trigger {
ForkRunOn::Changed { pattern } => (
pending_detail(&pending, "changed", pattern),
Some(("changed".to_string(), pattern.clone())),
),
ForkRunOn::Event { name } => (
pending_detail(&pending, "event", name),
Some(("event".to_string(), name.clone())),
),
_ => (None, None),
};
selected.push(SelectedFork {
name: entry.fork_name.clone(),
path: entry.fork_path.clone(),
trigger: label,
overlap: parsed.def.overlap,
after: parsed.def.after.clone(),
priority: parsed.def.priority,
tags: parsed.def.tags.clone(),
chain: parsed.def.chain,
gate: parsed.def.gate,
model: parsed.def.model.clone(),
mode: parsed.def.mode.clone(),
latch_key,
detail,
consumes,
});
}
apply_gate_filter(daemon, session, &mut selected);
selected
}
fn pending_detail(pending: &[(String, String, String)], kind: &str, key: &str) -> Option<String> {
pending
.iter()
.find(|(k, s, _)| k == kind && s == key)
.map(|(_, _, detail)| detail.clone())
.filter(|d| !d.is_empty())
}
pub fn reserve_fast_path(session: &SessionRow, selected: &mut Vec<SelectedFork>) {
if session.client.as_deref() == Some("codex") {
selected.retain(|s| !(s.chain && s.trigger == "idle:0"));
}
}
pub fn build_final_runs(daemon: &Arc<Daemon>, session: &SessionRow) -> Vec<WakeFork> {
let cfg = daemon.cfg_for(Some(&session.project_root));
let defs = collect_roster_defs(daemon, session);
let deadlines =
autofork_core::moments::idle_deadlines(defs.iter(), cfg.default_idle_deadline_secs);
let moments: Vec<ForkMoment> = deadlines
.into_iter()
.map(|deadline_secs| ForkMoment::Idle { deadline_secs })
.collect();
let mut selected = select_forks(daemon, session, &cfg, &moments);
selected.retain(|s| is_idle_trigger(&s.trigger));
if selected.is_empty() {
return Vec::new();
}
tracing::info!(
session = %session.session_id,
forks = ?selected.iter().map(|s| s.name.as_str()).collect::<Vec<_>>(),
"flush-on-close: issuing final runs"
);
let deps = resolve_deps(&selected);
let eff = effective_priorities(&selected, &deps);
let n = selected.len();
let mut emitted: Vec<usize> = Vec::with_capacity(n);
let mut done = vec![false; n];
while emitted.len() < n {
let mut pick: Option<usize> = None;
for i in 0..n {
if done[i] || deps[i].iter().any(|&j| !done[j]) {
continue;
}
if pick.map(|p| eff[i] < eff[p]).unwrap_or(true) {
pick = Some(i);
}
}
let Some(i) = pick else { break };
done[i] = true;
emitted.push(i);
}
let t = now();
{
let store = daemon.store.lock().unwrap();
for sel in &selected {
let _ = store.touch_fork_ran(&session.session_id, &sel.name, t);
let tags_joined = (!sel.tags.is_empty()).then(|| sel.tags.join(","));
let _ = store.record_issued_run(
&session.session_id,
&sel.name,
&sel.trigger,
tags_joined.as_deref(),
t,
);
if let Some(key) = &sel.latch_key {
let _ = store.try_latch_fire(&session.session_id, &sel.name, key, t);
}
}
}
let conv = conversation_id(session);
let root_str = session.project_root.to_string_lossy();
emitted
.into_iter()
.map(|i| {
let sel = &selected[i];
let report_preds: Vec<String> =
deps[i].iter().map(|&j| selected[j].name.clone()).collect();
let (model, model_fallbacks) = split_candidates(resolve_scoped(
&sel.model,
session.client.as_deref(),
session.model.as_deref(),
&cfg.fork_models,
));
let mode = resolve_scoped(
&sel.mode,
session.client.as_deref(),
session.model.as_deref(),
&cfg.fork_modes,
)
.into_iter()
.next();
let due = DueFork {
name: sel.name.clone(),
path: sel.path.to_string_lossy().into_owned(),
trigger: format!("{} (at close)", sel.trigger),
overlap: sel.overlap,
after: report_preds,
skill: autofork_core::discovery::skill_sibling(&sel.path)
.map(|p| p.to_string_lossy().into_owned()),
chain: false, model,
model_fallbacks,
mode,
detail: None,
};
build_wake_forks(&session.session_id, &conv, &root_str, &[due]).remove(0)
})
.collect()
}
fn collect_roster_defs(
daemon: &Arc<Daemon>,
session: &SessionRow,
) -> Vec<autofork_core::frontmatter::ForkDef> {
refresh_roster(daemon, &session.session_id, &session.cwd);
let roster = {
let store = daemon.store.lock().unwrap();
store.roster(&session.session_id).unwrap_or_default()
};
roster
.into_iter()
.filter_map(|entry| {
let content = std::fs::read_to_string(&entry.fork_path).ok()?;
match parse_fork(&entry.fork_name, &content) {
ForkParse::Fork(p) => Some(p.def),
_ => None,
}
})
.collect()
}
fn is_idle_trigger(label: &str) -> bool {
label == "idle" || label.starts_with("idle:")
}
pub(crate) fn gate_grace_secs() -> i64 {
std::env::var("AUTOFORK_GATE_GRACE_SECS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(180)
}
pub(crate) fn overlap_spawn_max_age_secs() -> i64 {
std::env::var("AUTOFORK_OVERLAP_SPAWN_MAX_AGE_SECS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(30 * 60)
}
pub(crate) fn runaway_window_secs() -> i64 {
std::env::var("AUTOFORK_RUNAWAY_WINDOW_SECS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(3600)
}
fn apply_gate_filter(daemon: &Arc<Daemon>, session: &SessionRow, selected: &mut Vec<SelectedFork>) {
let selected_gate = selected.iter().any(|s| s.gate);
let mut holding = selected_gate;
if !holding {
let gate = {
let store = daemon.store.lock().unwrap();
store
.get_session(&session.session_id)
.ok()
.flatten()
.and_then(|s| s.active_gate)
};
if let Some(g) = gate {
let store = daemon.store.lock().unwrap();
let live = store
.live_spawn_exists(&session.session_id, &g)
.unwrap_or(false);
let recent = store
.last_issued_at(&session.session_id, &g)
.ok()
.flatten()
.is_some_and(|at| now() - at < gate_grace_secs());
if live || recent {
holding = true;
} else {
tracing::warn!(session = %session.session_id, gate = %g,
"active gate has no live spawn and no recent wake — lifting it");
let _ = store.clear_active_gate(&session.session_id);
}
}
}
if holding {
let before = selected.len();
selected.retain(|s| s.gate || !is_idle_trigger(&s.trigger));
if selected.len() < before {
tracing::info!(session = %session.session_id, held = before - selected.len(),
"gate active: holding other idle forks");
}
}
}
fn parse_fork(name: &str, content: &str) -> ForkParse {
autofork_core::frontmatter::parse_fork_file(name, content)
}
fn resolve_scoped(
scoped: &autofork_core::frontmatter::ClientScoped,
client: Option<&str>,
parent_model: Option<&str>,
cfg_table: &std::collections::BTreeMap<String, autofork_core::config::ModelPolicy>,
) -> Vec<String> {
let client = client.unwrap_or("claude-code");
let own = scoped.resolve(client);
if !own.is_empty() {
return own.to_vec();
}
cfg_table
.get(client)
.map(|policy| policy.resolve(parent_model).to_vec())
.unwrap_or_default()
}
fn split_candidates(mut c: Vec<String>) -> (Option<String>, Vec<String>) {
if c.is_empty() {
(None, Vec::new())
} else {
let first = c.remove(0);
(Some(first), c)
}
}
fn conversation_id(session: &SessionRow) -> String {
session
.transcript_path
.as_deref()
.and_then(|p| p.file_stem())
.map(|s| s.to_string_lossy().into_owned())
.unwrap_or_else(|| session.session_id.clone())
}
pub fn build_wake(
daemon: &Arc<Daemon>,
session: &SessionRow,
selected: Vec<SelectedFork>,
) -> Option<(String, Vec<WakeFork>)> {
if selected.is_empty() {
return None;
}
tracing::info!(
session = %session.session_id,
forks = ?selected.iter().map(|s| s.name.as_str()).collect::<Vec<_>>(),
"issuing wake"
);
let cfg = daemon.cfg_for(Some(&session.project_root));
let deps = resolve_deps(&selected);
let eff = effective_priorities(&selected, &deps);
let min_eff = eff.iter().min().copied().unwrap_or(0);
let t = now();
let mut roots: Vec<DueFork> = Vec::new();
let mut held: Vec<HeldFork> = Vec::new();
{
let store = daemon.store.lock().unwrap();
for (i, sel) in selected.iter().enumerate() {
let _ = store.touch_fork_ran(&session.session_id, &sel.name, t);
let tags_joined = (!sel.tags.is_empty()).then(|| sel.tags.join(","));
let _ = store.record_issued_run(
&session.session_id,
&sel.name,
&sel.trigger,
tags_joined.as_deref(),
t,
);
if let Some(key) = &sel.latch_key {
let _ = store.try_latch_fire(&session.session_id, &sel.name, key, t);
}
if let Some((kind, key)) = &sel.consumes {
let _ = store.clear_pending_trigger(&session.session_id, kind, key);
}
let report_preds: Vec<String> =
deps[i].iter().map(|&j| selected[j].name.clone()).collect();
let mut gate = report_preds.clone();
if eff[i] > min_eff {
for (j, other) in selected.iter().enumerate() {
if j != i && eff[j] < eff[i] && !gate.contains(&other.name) {
gate.push(other.name.clone());
}
}
}
if sel.gate {
let _ = store.set_active_gate(&session.session_id, &sel.name);
}
let (model, fallbacks) = split_candidates(resolve_scoped(
&sel.model,
session.client.as_deref(),
session.model.as_deref(),
&cfg.fork_models,
));
let mode = resolve_scoped(
&sel.mode,
session.client.as_deref(),
session.model.as_deref(),
&cfg.fork_modes,
)
.into_iter()
.next();
if gate.is_empty() {
roots.push(DueFork {
name: sel.name.clone(),
path: sel.path.to_string_lossy().into_owned(),
trigger: sel.trigger.clone(),
overlap: sel.overlap,
after: Vec::new(),
skill: autofork_core::discovery::skill_sibling(&sel.path)
.map(|p| p.to_string_lossy().into_owned()),
chain: sel.chain,
model: model.clone(),
model_fallbacks: fallbacks.clone(),
mode: mode.clone(),
detail: sel.detail.clone(),
});
} else {
let _ = store.insert_pending_dep(
&session.session_id,
&sel.name,
&sel.path,
&sel.trigger,
sel.overlap,
&gate,
&report_preds,
t,
);
held.push(HeldFork {
name: sel.name.clone(),
after: gate,
});
}
}
}
daemon.note_wake_issued(&session.session_id);
let conv = conversation_id(session);
let root_str = session.project_root.to_string_lossy();
let payload = build_wake_payload(&session.session_id, &conv, &root_str, &roots, &held);
let forks = build_wake_forks(&session.session_id, &conv, &root_str, &roots);
Some((payload, forks))
}
pub fn release_due(daemon: &Arc<Daemon>, session: &SessionRow) -> Option<(String, Vec<WakeFork>)> {
let released: Vec<autofork_core::store::PendingDep> = {
let store = daemon.store.lock().unwrap();
let active_gate = store
.get_session(&session.session_id)
.ok()
.flatten()
.and_then(|s| s.active_gate);
let pending = store.list_pending_deps(&session.session_id).ok()?;
pending
.into_iter()
.filter(|dep| active_gate.as_deref().is_none_or(|g| g == dep.fork_name))
.filter(|dep| {
dep.preds.iter().all(|pred| {
store
.fork_completed_since(&session.session_id, pred, dep.created_at)
.unwrap_or(false)
})
})
.collect()
};
if released.is_empty() {
return None;
}
tracing::info!(
session = %session.session_id,
forks = ?released.iter().map(|d| d.fork_name.as_str()).collect::<Vec<_>>(),
"releasing held dependents"
);
let def_of =
|dep: &autofork_core::store::PendingDep| -> Option<autofork_core::frontmatter::ForkDef> {
std::fs::read_to_string(&dep.fork_path).ok().and_then(|c| {
match parse_fork(&dep.fork_name, &c) {
ForkParse::Fork(p) => Some(p.def),
_ => None,
}
})
};
let cfg = daemon.cfg_for(Some(&session.project_root));
let due: Vec<DueFork> = released
.iter()
.map(|dep| {
let def = def_of(dep);
let (model, model_fallbacks) = split_candidates(
def.as_ref()
.map(|d| {
resolve_scoped(
&d.model,
session.client.as_deref(),
session.model.as_deref(),
&cfg.fork_models,
)
})
.unwrap_or_default(),
);
let mode = def
.as_ref()
.map(|d| {
resolve_scoped(
&d.mode,
session.client.as_deref(),
session.model.as_deref(),
&cfg.fork_modes,
)
})
.unwrap_or_default()
.into_iter()
.next();
DueFork {
name: dep.fork_name.clone(),
path: dep.fork_path.to_string_lossy().into_owned(),
trigger: dep.trigger_label.clone(),
overlap: dep.overlap,
after: dep.report_preds.clone(),
skill: autofork_core::discovery::skill_sibling(&dep.fork_path)
.map(|p| p.to_string_lossy().into_owned()),
chain: def.as_ref().map(|d| d.chain).unwrap_or(false),
model,
model_fallbacks,
mode,
detail: None,
}
})
.collect();
{
let store = daemon.store.lock().unwrap();
for dep in &released {
let _ = store.delete_pending_dep(&session.session_id, &dep.fork_name);
if def_of(dep).map(|d| d.gate).unwrap_or(false) {
let _ = store.set_active_gate(&session.session_id, &dep.fork_name);
}
}
}
daemon.note_wake_issued(&session.session_id);
let conv = conversation_id(session);
let root_str = session.project_root.to_string_lossy();
let payload = build_release_payload(&session.session_id, &conv, &root_str, &due);
let forks = build_wake_forks(&session.session_id, &conv, &root_str, &due);
Some((payload, forks))
}