use anyhow::{Context, Result};
use chrono::{DateTime, Duration, Months, Utc};
use serde_json::{json, Value};
use khive_mcp::serve::enforce_strict_actor_mode;
use khive_mcp::server::KhiveMcpServer;
use khive_mcp::tools::request::RequestParams;
use khive_runtime::{KhiveRuntime, Namespace, RuntimeConfig};
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter, SortDir};
use khive_storage::types::{PageRequest, SqlValue};
use crate::dbpath::resolve_db_override;
#[derive(Debug, Default)]
pub struct DrainSummary {
pub scanned: u64,
pub fired: u64,
pub advanced: u64,
pub failed: u64,
pub skipped_not_due: u64,
}
pub async fn run_pending_events(
db: Option<&str>,
namespace: &str,
verbose: bool,
) -> Result<DrainSummary> {
let mut cfg = RuntimeConfig::default();
if let Some(db_path) = resolve_db_override(db) {
cfg.db_path = db_path;
}
cfg.default_namespace = Namespace::parse(namespace).map_err(|e| anyhow::anyhow!("{e}"))?;
let rt = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
enforce_strict_actor_mode(rt.config().actor_id.as_deref(), &rt.config().packs)?;
let server = KhiveMcpServer::new(rt.clone()).map_err(|e| anyhow::anyhow!("{e}"))?;
let now = Utc::now();
let mut summary = DrainSummary::default();
let namespaces = discover_pending_namespaces(&rt, now).await?;
if verbose {
eprintln!(
"[pending-events] scan: now={}, namespaces_with_pending={}",
now.to_rfc3339(),
namespaces.len()
);
}
for ns_str in &namespaces {
let ns = match Namespace::parse(ns_str) {
Ok(n) => n,
Err(e) => {
if verbose {
eprintln!("[pending-events] skip invalid namespace {ns_str:?}: {e}");
}
continue;
}
};
let token = match rt.authorize(ns.clone()) {
Ok(t) => t,
Err(e) => {
if verbose {
eprintln!("[pending-events] authorize({ns_str}) failed: {e}");
}
continue;
}
};
let store = match rt.notes(&token) {
Ok(s) => s,
Err(e) => {
if verbose {
eprintln!("[pending-events] notes({ns_str}) failed: {e}");
}
continue;
}
};
let filter = NoteFilter {
kind: Some("scheduled_event".to_string()),
property_filters: vec![PropertyFilter {
json_path: "$.status".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("pending".to_string()),
}],
order_by: Some(("$.trigger_at".to_string(), SortDir::Asc)),
..Default::default()
};
const PAGE_SIZE: u32 = 200;
let mut offset: u64 = 0;
loop {
let page = store
.query_notes_filtered(
ns_str,
&filter,
PageRequest {
limit: PAGE_SIZE,
offset,
},
)
.await
.with_context(|| {
format!("pending-events: query_notes_filtered failed for ns={ns_str}")
})?;
let page_len = page.items.len() as u32;
for mut note in page.items {
summary.scanned += 1;
let trigger_at_str = note
.properties
.as_ref()
.and_then(|p| p.get("trigger_at"))
.and_then(Value::as_str)
.unwrap_or("");
let trigger_at = match trigger_at_str.parse::<DateTime<Utc>>() {
Ok(ts) => ts,
Err(_) => {
if verbose {
eprintln!(
"[pending-events] skip note {}: unparseable trigger_at {:?}",
note.id, trigger_at_str
);
}
summary.skipped_not_due += 1;
continue;
}
};
if trigger_at > now {
summary.skipped_not_due += 1;
continue;
}
let event_type = note
.properties
.as_ref()
.and_then(|p| p.get("event_type"))
.and_then(Value::as_str)
.unwrap_or("remind");
let action_dsl: Option<String> = if event_type == "schedule" {
note.properties
.as_ref()
.and_then(|p| p.get("payload"))
.and_then(Value::as_str)
.map(str::to_string)
} else {
None
};
if let Some(dsl) = &action_dsl {
let dispatch_result = dispatch_action(dsl, ns_str, &server, verbose).await;
if let Err(e) = dispatch_result {
if verbose {
eprintln!("[pending-events] dispatch failed for note {}: {e}", note.id);
}
summary.failed += 1;
}
}
let repeat = note
.properties
.as_ref()
.and_then(|p| p.get("repeat"))
.and_then(Value::as_str)
.map(str::to_string);
let fired_at_rfc = Utc::now().to_rfc3339();
let mut props = note.properties.clone().unwrap_or_else(|| json!({}));
match next_trigger_at(&repeat, trigger_at) {
Some(next_at) => {
props["trigger_at"] = json!(next_at.to_rfc3339());
props["status"] = json!("pending");
props["fired_at"] = json!(fired_at_rfc);
note.properties = Some(props);
note.updated_at = Utc::now().timestamp_micros();
summary.advanced += 1;
}
None => {
props["status"] = json!("fired");
props["fired_at"] = json!(fired_at_rfc);
note.properties = Some(props);
note.updated_at = Utc::now().timestamp_micros();
summary.fired += 1;
}
}
if let Err(e) = store.upsert_note(note.clone()).await {
if verbose {
eprintln!("[pending-events] upsert_note failed for {}: {e}", note.id);
}
summary.failed += 1;
if summary.fired > 0 {
summary.fired -= 1;
}
if summary.advanced > 0 {
summary.advanced -= 1;
}
}
}
if page_len < PAGE_SIZE {
break;
}
offset = offset
.checked_add(u64::from(PAGE_SIZE))
.ok_or_else(|| anyhow::anyhow!("pending-events: pagination offset overflow"))?;
}
}
Ok(summary)
}
fn next_trigger_at(repeat: &Option<String>, current: DateTime<Utc>) -> Option<DateTime<Utc>> {
match repeat.as_deref() {
Some("daily") => Some(current + Duration::days(1)),
Some("weekly") => Some(current + Duration::weeks(1)),
Some("monthly") => {
current.checked_add_months(Months::new(1))
}
Some(expr) if is_five_field_cron(expr) => {
tracing::warn!(
repeat = expr,
"pending-events: cron repeat expression cannot be advanced (not yet supported); \
event will be marked fired (one-shot)"
);
None
}
_ => None,
}
}
fn is_five_field_cron(expr: &str) -> bool {
expr.split_whitespace().count() == 5
}
async fn dispatch_action(
action_dsl: &str,
namespace: &str,
server: &KhiveMcpServer,
verbose: bool,
) -> Result<()> {
let parsed = khive_request::parse_request(action_dsl).map_err(|e| {
anyhow::anyhow!("pending-events: action DSL parse error ({e}): {action_dsl:?}")
})?;
let ops_json: Vec<Value> = parsed
.ops
.iter()
.map(|op| {
let mut args = serde_json::Map::new();
for (k, v) in &op.args {
if let khive_request::ArgValue::Value(val) = v {
args.insert(k.clone(), val.clone());
}
}
args.insert(
"namespace".to_string(),
Value::String(namespace.to_string()),
);
json!({ "tool": op.tool, "args": Value::Object(args) })
})
.collect();
let ops_str = serde_json::to_string(&ops_json)
.map_err(|e| anyhow::anyhow!("pending-events: serialize ops: {e}"))?;
if verbose {
eprintln!("[pending-events] dispatch ns={namespace}: {ops_str}");
}
let result = server
.dispatch_request_local(RequestParams {
ops: ops_str,
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
})
.await
.map_err(|e| anyhow::anyhow!("pending-events: dispatch error: {e}"))?;
let parsed_result: Value = serde_json::from_str(&result).unwrap_or(Value::Null);
if let Some(results) = parsed_result.get("results").and_then(Value::as_array) {
let failures: Vec<_> = results
.iter()
.filter(|r| r.get("ok").and_then(Value::as_bool) == Some(false))
.collect();
if !failures.is_empty() {
let errs: Vec<String> = failures
.iter()
.filter_map(|r| r.get("error").and_then(Value::as_str).map(str::to_string))
.collect();
return Err(anyhow::anyhow!(
"pending-events: action produced {} failure(s): {}",
failures.len(),
errs.join("; ")
));
}
}
Ok(())
}
async fn discover_pending_namespaces(rt: &KhiveRuntime, now: DateTime<Utc>) -> Result<Vec<String>> {
use khive_storage::types::{SqlStatement, SqlValue};
let sql_access = rt.sql();
let mut reader = sql_access
.reader()
.await
.context("pending-events: open SQL reader")?;
let now_rfc = now.to_rfc3339();
let rows = reader
.query_all(SqlStatement {
sql: "SELECT DISTINCT namespace \
FROM notes \
WHERE kind = 'scheduled_event' \
AND deleted_at IS NULL \
AND json_extract(properties, '$.status') = 'pending' \
AND json_extract(properties, '$.trigger_at') <= ?1"
.into(),
params: vec![SqlValue::Text(now_rfc)],
label: Some("pending_events_namespaces".into()),
})
.await
.context("pending-events: discover namespaces query")?;
let namespaces: Vec<String> = rows
.into_iter()
.filter_map(|row| {
row.get("namespace").and_then(|v| {
if let SqlValue::Text(s) = v {
Some(s.clone())
} else {
None
}
})
})
.collect();
Ok(namespaces)
}
pub fn print_summary(summary: &DrainSummary) {
let json = json!({
"scanned": summary.scanned,
"fired": summary.fired,
"advanced": summary.advanced,
"failed": summary.failed,
"skipped_not_due": summary.skipped_not_due,
});
println!(
"{}",
serde_json::to_string_pretty(&json).expect("serialize")
);
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::NamedTempFile;
fn tmp_db() -> (NamedTempFile, String) {
let f = NamedTempFile::new().expect("tempfile");
let path = f.path().to_str().expect("utf8 path").to_string();
(f, path)
}
async fn make_rt(db_path: &str) -> KhiveRuntime {
let cfg = RuntimeConfig {
db_path: Some(std::path::PathBuf::from(db_path)),
default_namespace: Namespace::parse("local").unwrap(),
embedding_model: None,
additional_embedding_models: vec![],
..Default::default()
};
KhiveRuntime::new(cfg).expect("runtime")
}
async fn create_scheduled_event(
rt: &KhiveRuntime,
namespace: &str,
trigger_at: &str,
action_dsl: Option<&str>,
repeat: Option<&str>,
event_type: &str,
) -> uuid::Uuid {
let props = json!({
"trigger_at": trigger_at,
"repeat": repeat,
"status": "pending",
"event_type": event_type,
"payload": action_dsl,
"fired_at": null,
"cancelled_at": null,
});
let ns = Namespace::parse(namespace).expect("ns");
let token = rt.authorize(ns).expect("authorize");
let content = action_dsl.unwrap_or("test reminder");
let note = rt
.create_note(
&token,
"scheduled_event",
None,
content,
None,
Some(props),
vec![],
)
.await
.expect("create_note");
note.id
}
async fn get_note_props(rt: &KhiveRuntime, id: uuid::Uuid) -> Value {
let ns = Namespace::parse("local").unwrap();
let token = rt.authorize(ns).expect("authorize");
let store = rt.notes(&token).expect("notes");
let note = store
.get_note(id)
.await
.expect("get_note")
.expect("note exists");
note.properties.unwrap_or(json!({}))
}
#[tokio::test]
async fn due_event_is_fired() {
let (_tmp, db_path) = tmp_db();
let rt = make_rt(&db_path).await;
let past = "2000-01-01T00:00:00Z";
let id =
create_scheduled_event(&rt, "local", past, Some("stats()"), None, "schedule").await;
let summary = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain");
assert!(summary.scanned >= 1, "must have scanned the due event");
assert!(
summary.fired >= 1 || summary.advanced >= 1,
"must fire or advance"
);
let props = get_note_props(&rt, id).await;
let status = props["status"].as_str().unwrap_or("");
assert!(
status == "fired" || status == "pending",
"status must be fired or pending (repeat), got {status:?}"
);
}
#[tokio::test]
async fn future_event_is_skipped() {
let (_tmp, db_path) = tmp_db();
let rt = make_rt(&db_path).await;
let future = "2099-01-01T00:00:00Z";
let id =
create_scheduled_event(&rt, "local", future, Some("stats()"), None, "schedule").await;
let summary = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain");
assert_eq!(summary.fired, 0, "future event must not be fired");
assert_eq!(summary.advanced, 0, "future event must not be advanced");
let props = get_note_props(&rt, id).await;
assert_eq!(
props["status"].as_str(),
Some("pending"),
"future event must remain pending"
);
}
#[tokio::test]
async fn fired_event_is_idempotent() {
let (_tmp, db_path) = tmp_db();
let rt = make_rt(&db_path).await;
let past = "2000-01-01T00:00:00Z";
let id =
create_scheduled_event(&rt, "local", past, Some("stats()"), None, "schedule").await;
let s1 = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain 1");
assert!(s1.scanned >= 1);
let s2 = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain 2");
assert_eq!(s2.scanned, 0, "no pending events on second drain");
assert_eq!(s2.fired, 0, "no new fires on second drain");
let props = get_note_props(&rt, id).await;
let fired_at_1 = props["fired_at"].as_str().unwrap_or("").to_string();
assert!(
!fired_at_1.is_empty(),
"fired_at must be set after first drain"
);
let props2 = get_note_props(&rt, id).await;
assert_eq!(
props2["fired_at"].as_str().unwrap_or(""),
fired_at_1.as_str(),
"fired_at must not change on second drain"
);
}
#[tokio::test]
async fn daily_repeat_advances() {
let (_tmp, db_path) = tmp_db();
let rt = make_rt(&db_path).await;
let past = "2000-06-01T09:00:00Z";
let id = create_scheduled_event(
&rt,
"local",
past,
Some("stats()"),
Some("daily"),
"schedule",
)
.await;
let summary = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain");
assert!(
summary.advanced >= 1,
"daily event must be advanced, not fired"
);
let props = get_note_props(&rt, id).await;
assert_eq!(
props["status"].as_str(),
Some("pending"),
"after advance, status must be pending"
);
let new_trigger = props["trigger_at"]
.as_str()
.expect("trigger_at must be set");
let new_ts: DateTime<Utc> = new_trigger.parse().expect("parseable ts");
let original: DateTime<Utc> = past.parse().unwrap();
assert_eq!(
new_ts,
original + Duration::days(1),
"daily advance must add 1 day"
);
}
#[tokio::test]
async fn namespace_isolation() {
let (_tmp, db_path) = tmp_db();
let rt = make_rt(&db_path).await;
let ns_a = "ns-a";
let ns_b = "ns-b";
let past = "2000-01-01T00:00:00Z";
let id_a = create_scheduled_event(&rt, ns_a, past, Some("stats()"), None, "schedule").await;
let _id_b = create_scheduled_event(
&rt,
ns_b,
"2099-01-01T00:00:00Z",
Some("stats()"),
None,
"schedule",
)
.await;
let summary = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain");
assert!(summary.scanned >= 1);
assert!(summary.fired >= 1 || summary.advanced >= 1);
let token_a = rt.authorize(Namespace::parse(ns_a).unwrap()).expect("auth");
let store_a = rt.notes(&token_a).expect("notes");
let note_a = store_a.get_note(id_a).await.expect("get").expect("exists");
let status_a = note_a
.properties
.as_ref()
.and_then(|p| p.get("status"))
.and_then(Value::as_str)
.unwrap_or("");
assert!(
status_a == "fired" || status_a == "pending",
"ns-a event must be fired or advanced, got {status_a:?}"
);
}
#[tokio::test]
async fn dispatch_failure_does_not_abort_drain() {
let (_tmp, db_path) = tmp_db();
let rt = make_rt(&db_path).await;
let past = "2000-01-01T00:00:00Z";
let _id_bad = create_scheduled_event(
&rt,
"local",
past,
Some("stats()"), None,
"schedule",
)
.await;
let id_bad2 = create_scheduled_event(
&rt,
"local",
past,
Some("this_verb_does_not_exist(foo=\"bar\")"),
None,
"schedule",
)
.await;
let summary = run_pending_events(Some(&db_path), "local", false)
.await
.expect("drain must not abort");
assert!(summary.scanned >= 2, "both events must be scanned");
assert!(
summary.failed >= 1 || summary.fired >= 1,
"at least one event processed (failed or fired)"
);
let props_bad2 = get_note_props(&rt, id_bad2).await;
let _ = props_bad2["status"].as_str(); }
#[test]
fn next_trigger_at_daily() {
let base: DateTime<Utc> = "2026-06-01T09:00:00Z".parse().unwrap();
let next = next_trigger_at(&Some("daily".to_string()), base).unwrap();
assert_eq!(next, base + Duration::days(1));
}
#[test]
fn next_trigger_at_weekly() {
let base: DateTime<Utc> = "2026-06-01T09:00:00Z".parse().unwrap();
let next = next_trigger_at(&Some("weekly".to_string()), base).unwrap();
assert_eq!(next, base + Duration::weeks(1));
}
#[test]
fn next_trigger_at_monthly() {
let base: DateTime<Utc> = "2026-06-01T09:00:00Z".parse().unwrap();
let next = next_trigger_at(&Some("monthly".to_string()), base).unwrap();
let expected: DateTime<Utc> = "2026-07-01T09:00:00Z".parse().unwrap();
assert_eq!(next, expected);
}
#[test]
fn next_trigger_at_none_repeat_returns_none() {
let base: DateTime<Utc> = "2026-06-01T09:00:00Z".parse().unwrap();
assert!(next_trigger_at(&None, base).is_none());
}
#[test]
fn next_trigger_at_cron_returns_none() {
let base: DateTime<Utc> = "2026-06-01T09:00:00Z".parse().unwrap();
assert!(next_trigger_at(&Some("0 9 * * 1".to_string()), base).is_none());
}
}