mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-27 08:00:17 +00:00
Wait for one-shot event routines before exit
This commit is contained in:
@@ -1150,9 +1150,14 @@ impl Agent {
|
||||
&& let Submission::UserInput { ref content } = submission
|
||||
&& let Some(engine) = self.routine_engine().await
|
||||
{
|
||||
let single_message_repl = is_single_message_repl(message);
|
||||
// Use post-hook content so that BeforeInbound hooks that rewrite
|
||||
// input are respected by event trigger matching.
|
||||
let fired = engine.check_event_triggers(message, content).await;
|
||||
let fired = if single_message_repl {
|
||||
engine.check_event_triggers_and_wait(message, content).await
|
||||
} else {
|
||||
engine.check_event_triggers(message, content).await
|
||||
};
|
||||
if fired > 0 {
|
||||
tracing::debug!(
|
||||
channel = %message.channel,
|
||||
@@ -1160,7 +1165,7 @@ impl Agent {
|
||||
fired,
|
||||
"Consumed inbound user message with matching event-triggered routine(s)"
|
||||
);
|
||||
return if is_single_message_repl(message) {
|
||||
return if single_message_repl {
|
||||
Ok(None)
|
||||
} else {
|
||||
Ok(Some(String::new()))
|
||||
|
||||
+60
-10
@@ -18,6 +18,7 @@ use std::time::Duration;
|
||||
use chrono::Utc;
|
||||
use regex::Regex;
|
||||
use tokio::sync::{RwLock, mpsc};
|
||||
use tokio::task::JoinHandle;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::agent::Scheduler;
|
||||
@@ -45,6 +46,11 @@ enum EventMatcher {
|
||||
System { routine: Routine },
|
||||
}
|
||||
|
||||
struct TriggeredRoutine {
|
||||
routine: Routine,
|
||||
detail: String,
|
||||
}
|
||||
|
||||
/// Distinguishes why sandbox is unavailable so error messages are accurate.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum SandboxReadiness {
|
||||
@@ -202,6 +208,44 @@ impl RoutineEngine {
|
||||
|
||||
/// Check incoming message against event triggers. Returns number of routines fired.
|
||||
pub async fn check_event_triggers(&self, message: &IncomingMessage, content: &str) -> usize {
|
||||
let triggered = self.matching_event_triggers(message, content).await;
|
||||
let fired = triggered.len();
|
||||
for triggered in triggered {
|
||||
std::mem::drop(self.spawn_fire(triggered.routine, "event", Some(triggered.detail)));
|
||||
}
|
||||
fired
|
||||
}
|
||||
|
||||
/// Fire matching event-triggered routines and wait for them to complete.
|
||||
///
|
||||
/// Used by single-message REPL mode so the process does not exit before
|
||||
/// background event-triggered routines finish.
|
||||
pub async fn check_event_triggers_and_wait(
|
||||
&self,
|
||||
message: &IncomingMessage,
|
||||
content: &str,
|
||||
) -> usize {
|
||||
let triggered = self.matching_event_triggers(message, content).await;
|
||||
let fired = triggered.len();
|
||||
let handles: Vec<JoinHandle<()>> = triggered
|
||||
.into_iter()
|
||||
.map(|triggered| self.spawn_fire(triggered.routine, "event", Some(triggered.detail)))
|
||||
.collect();
|
||||
|
||||
for handle in handles {
|
||||
if let Err(e) = handle.await {
|
||||
tracing::warn!(error = %e, "Event-triggered routine task failed");
|
||||
}
|
||||
}
|
||||
|
||||
fired
|
||||
}
|
||||
|
||||
async fn matching_event_triggers(
|
||||
&self,
|
||||
message: &IncomingMessage,
|
||||
content: &str,
|
||||
) -> Vec<TriggeredRoutine> {
|
||||
let cache = self.event_cache.read().await;
|
||||
|
||||
// Early return if there are no message matchers at all.
|
||||
@@ -209,10 +253,9 @@ impl RoutineEngine {
|
||||
.iter()
|
||||
.any(|m| matches!(m, EventMatcher::Message { .. }))
|
||||
{
|
||||
return 0;
|
||||
return Vec::new();
|
||||
}
|
||||
|
||||
let mut fired = 0;
|
||||
let mut triggered = Vec::new();
|
||||
|
||||
// Collect routine IDs for batch query
|
||||
let routine_ids: Vec<Uuid> = cache
|
||||
@@ -224,13 +267,13 @@ impl RoutineEngine {
|
||||
.collect();
|
||||
|
||||
if routine_ids.is_empty() {
|
||||
return 0;
|
||||
return Vec::new();
|
||||
}
|
||||
|
||||
// Single batch query instead of N queries
|
||||
let concurrent_counts = match self.batch_concurrent_counts(&routine_ids).await {
|
||||
Some(counts) => counts,
|
||||
None => return 0,
|
||||
None => return Vec::new(),
|
||||
};
|
||||
|
||||
for matcher in cache.iter() {
|
||||
@@ -285,11 +328,13 @@ impl RoutineEngine {
|
||||
}
|
||||
|
||||
let detail = truncate(content, 200);
|
||||
self.spawn_fire(routine.clone(), "event", Some(detail));
|
||||
fired += 1;
|
||||
triggered.push(TriggeredRoutine {
|
||||
routine: routine.clone(),
|
||||
detail,
|
||||
});
|
||||
}
|
||||
|
||||
fired
|
||||
triggered
|
||||
}
|
||||
|
||||
/// Emit a structured event to system-event routines.
|
||||
@@ -845,7 +890,12 @@ impl RoutineEngine {
|
||||
}
|
||||
|
||||
/// Spawn a fire in a background task.
|
||||
fn spawn_fire(&self, routine: Routine, trigger_type: &str, trigger_detail: Option<String>) {
|
||||
fn spawn_fire(
|
||||
&self,
|
||||
routine: Routine,
|
||||
trigger_type: &str,
|
||||
trigger_detail: Option<String>,
|
||||
) -> JoinHandle<()> {
|
||||
let run = RoutineRun {
|
||||
id: Uuid::new_v4(),
|
||||
routine_id: routine.id,
|
||||
@@ -882,7 +932,7 @@ impl RoutineEngine {
|
||||
return;
|
||||
}
|
||||
execute_routine(engine, routine, run).await;
|
||||
});
|
||||
})
|
||||
}
|
||||
|
||||
fn check_cooldown(&self, routine: &Routine) -> bool {
|
||||
|
||||
Reference in New Issue
Block a user