use crate::engine::model::{UnloadPolicy, Workflow};
use crate::state::{Kind, now_ms};
use serde_json::json;
pub(crate) const PIN_PREFIX: &str = "_pins/";
#[derive(Debug, Clone)]
pub(crate) struct Retiring {
pub hash: String,
pub deadline_ms: Option<u64>,
}
impl super::reactor::Runtime {
pub(crate) fn ensure_pin(&mut self, wf: &Workflow) {
if !self.pin_written.insert(wf.hash.clone()) {
return;
}
if let Err(e) = self.durable.put(
Kind::Memory,
&format!("{PIN_PREFIX}{}", wf.hash),
json!({"name": wf.name, "definition": wf.definition}),
None,
) {
self.log.warn(
"workflow.pin_fail",
json!({"workflow": wf.name, "hash": &wf.hash[..12.min(wf.hash.len())], "err": e.to_string()}),
);
self.pin_written.remove(&wf.hash);
}
}
pub(crate) fn restore_pins(&mut self) {
let missing: Vec<(String, String)> = self
.runs
.values()
.filter(|r| !r.status.is_terminal())
.filter(|r| {
let current = self
.workflows
.get(&r.workflow)
.is_some_and(|w| w.hash == r.workflow_hash);
!current && !self.pinned.contains_key(&r.workflow_hash)
})
.map(|r| (r.workflow.clone(), r.workflow_hash.clone()))
.collect();
for (name, hash) in missing {
let doc = self
.durable
.get(Kind::Memory, &format!("{PIN_PREFIX}{hash}"))
.ok()
.flatten()
.and_then(|env| env.state.get("definition").cloned());
match doc.map(|d| crate::engine::model::parse_workflow(&d)) {
Some(Ok(wf)) if wf.hash == hash => {
self.log.info(
"workflow.pin_restored",
json!({"workflow": name, "hash": &hash[..12.min(hash.len())]}),
);
self.pin_written.insert(hash.clone());
self.pinned.insert(hash, std::sync::Arc::new(wf));
}
other => {
self.log.warn(
"workflow.pin_missing",
json!({"workflow": name, "hash": &hash[..12.min(hash.len())],
"err": match other { Some(Err(e)) => e.join("; "), _ => "no durable pin".into() }}),
);
}
}
}
}
pub(crate) fn retire_workflow(&mut self, wf: &Workflow, reason: &str) {
for s in wf.start_steps() {
if s.kind != "subscribe" {
continue;
}
let (Some(server), Some(uri)) = (s.field_str("server"), s.field_str("uri")) else {
continue;
};
let still_wanted = self
.workflows
.values()
.filter(|w| w.name != wf.name && w.armed)
.flat_map(|w| w.start_steps())
.any(|o| {
o.kind == "subscribe"
&& o.field_str("server") == Some(server)
&& o.field_str("uri") == Some(uri)
});
if !still_wanted && let Some(c) = self.mcp.get(server) {
match c.unsubscribe(uri) {
Ok(()) => self.log.info(
"workflow.unsubscribed",
json!({"workflow": wf.name, "server": server, "uri": uri}),
),
Err(e) => self.log.warn(
"workflow.unsubscribe_fail",
json!({"workflow": wf.name, "server": server, "uri": uri, "err": e.to_string()}),
),
}
}
}
let live: Vec<String> = self
.runs
.values()
.filter(|r| {
!r.status.is_terminal() && r.workflow == wf.name && r.workflow_hash == wf.hash
})
.map(|r| r.id.clone())
.collect();
if live.is_empty() {
self.log.info(
"workflow.unloaded",
json!({"workflow": wf.name, "hash": &wf.hash[..12.min(wf.hash.len())], "reason": reason, "live_runs": 0}),
);
return;
}
self.pinned
.insert(wf.hash.clone(), std::sync::Arc::new(wf.clone()));
let policy = wf.unload.policy;
self.log.info(
"workflow.retiring",
json!({"workflow": wf.name, "hash": &wf.hash[..12.min(wf.hash.len())],
"reason": reason, "policy": policy.as_str(), "live_runs": live.len(),
"timeout_ms": wf.unload.timeout_ms}),
);
match policy {
UnloadPolicy::Cancel => {
for id in &live {
self.cancel_run(id, "workflow retired (unload: cancel)");
}
self.retiring.insert(
wf.hash.clone(),
Retiring {
hash: wf.hash.clone(),
deadline_ms: None,
},
);
}
UnloadPolicy::Drain | UnloadPolicy::Detach => {
let deadline_ms = match policy {
UnloadPolicy::Drain => wf.unload.timeout_ms.map(|t| now_ms() + t),
_ => None,
};
self.retiring.insert(
wf.hash.clone(),
Retiring {
hash: wf.hash.clone(),
deadline_ms,
},
);
}
}
}
pub(crate) fn retire_tick(&mut self) {
if self.retiring.is_empty() {
return;
}
let now = now_ms();
let overdue: Vec<String> = self
.retiring
.values()
.filter(|r| r.deadline_ms.is_some_and(|d| now >= d))
.map(|r| r.hash.clone())
.collect();
for hash in overdue {
let victims: Vec<String> = self
.runs
.values()
.filter(|r| !r.status.is_terminal() && r.workflow_hash == hash)
.map(|r| r.id.clone())
.collect();
for id in &victims {
self.cancel_run(id, "workflow retired (unload drain timeout)");
}
if let Some(r) = self.retiring.get_mut(&hash) {
r.deadline_ms = None; }
}
}
pub(crate) fn retire_sweep(&mut self) {
if self.pinned.is_empty() && self.retiring.is_empty() {
return;
}
let referenced: std::collections::BTreeSet<String> = self
.runs
.values()
.filter(|r| !r.status.is_terminal())
.map(|r| r.workflow_hash.clone())
.collect();
let dropped: Vec<(String, String)> = self
.pinned
.iter()
.filter(|(hash, _)| !referenced.contains(*hash))
.map(|(hash, wf)| (hash.clone(), wf.name.clone()))
.collect();
for (hash, name) in &dropped {
self.log.info(
"workflow.unloaded",
json!({"workflow": name, "hash": &hash[..12.min(hash.len())], "live_runs": 0}),
);
let _ = self
.durable
.delete(Kind::Memory, &format!("{PIN_PREFIX}{hash}"));
self.pinned.remove(hash);
self.pin_written.remove(hash);
}
self.retiring.retain(|hash, _| referenced.contains(hash));
}
}