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>,
}
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,
_ => 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(),
);
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 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;
}
}
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,
});
}
apply_gate_filter(daemon, session, &mut selected);
selected
}
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"));
}
}
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)
}
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>,
cfg_table: &std::collections::BTreeMap<String, String>,
) -> Option<String> {
let client = client.unwrap_or("claude-code");
scoped
.resolve(client)
.map(str::to_string)
.or_else(|| cfg_table.get(client).cloned())
}
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);
}
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);
}
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: resolve_scoped(&sel.model, session.client.as_deref(), &cfg.fork_models),
mode: resolve_scoped(&sel.mode, session.client.as_deref(), &cfg.fork_modes),
});
} 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);
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: def.as_ref().and_then(|d| {
resolve_scoped(&d.model, session.client.as_deref(), &cfg.fork_models)
}),
mode: def.as_ref().and_then(|d| {
resolve_scoped(&d.mode, session.client.as_deref(), &cfg.fork_modes)
}),
}
})
.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))
}