mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
* feat: complete multi-tenant isolation — per-user budgets, model selection, heartbeat cycling Finishes the remaining isolation work from phases 2–4 of #59: Phase 2 (DB scoping): Fix /status and /list commands to use _for_user DB variants instead of global queries that leaked cross-user job data. Phase 3 (Runtime isolation): Per-user workspace in routine engine's spawn_fire so lightweight routines run in the correct user context. Per-user daily cost tracking in CostGuard with configurable budget via MAX_COST_PER_USER_PER_DAY_CENTS. Multi-user heartbeat that cycles through all users with routines, auto-detected from GATEWAY_USER_TOKENS. Phase 4 (Provider/tools): Per-user model selection via preferred_model setting — looked up from SettingsStore on first iteration, threaded through ReasoningContext.model_override to CompletionRequest. Works with providers that support per-request model overrides (NearAI). Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: use selected_model setting key to match /model command persistence The dispatcher was reading "preferred_model" but the /model command (merged from staging) persists to "selected_model". Since set_setting is already per-user scoped, using the same key makes /model work as the per-user model override in multi-tenant mode. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: heartbeat hygiene, /model multi-tenant guard, RigAdapter model override Three follow-up fixes for multi-tenant isolation: 1. Multi-user heartbeat now runs memory hygiene per user before each heartbeat check, matching single-user heartbeat behavior. 2. /model command in multi-tenant mode only persists to per-user settings (selected_model) without calling set_model() on the shared LlmProvider. The per-request model_override in the dispatcher reads from the same setting. Added multi_tenant flag to AgentConfig (auto-detected from GATEWAY_USER_TOKENS). 3. RigAdapter now supports per-request model overrides by injecting the model name into rig-core's additional_params. OpenAI/Anthropic/Ollama API servers use last-key-wins for duplicate JSON keys, so the override takes effect via serde's flatten serialization order. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address PR review — cost model attribution, heartbeat concurrency, pruning Fixes from review comments on #1614: - Cost tracking now uses the override model name (not active_model_name) when a per-user model override is active, for accurate attribution. - Multi-user heartbeat runs per-user checks concurrently via JoinSet instead of sequentially, preventing one slow user from blocking others. - Per-user failure counts tracked independently; users exceeding max_failures are skipped (matching single-user semantics). - per_user_daily_cost HashMap pruned on day rollover to prevent unbounded growth in long-lived deployments. - Doc comment fixed: says "routines" not "active routines". Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: /status ownership, model persistence scoping, heartbeat robustness Addresses second round of PR review on #1614: - /status <job_id> DB path now validates job.user_id == requesting user before returning data (was missing ownership check, security fix). - persist_selected_model takes user_id param instead of owner_id, and skips .env/TOML writes in multi-tenant mode (these are shared global files). handle_system_command now receives user_id from caller. - JoinSet collection handles Err(JoinError) explicitly instead of silently dropping panicked tasks. - Notification forwarder extracts owner_id from response metadata in multi-tenant mode for per-user routing instead of broadcasting to the agent owner. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: cost pricing, fire_manual workspace, heartbeat concurrency cap Round 3 review fixes: - Cost tracking passes None for cost_per_token when model override is active, letting CostGuard look up pricing by model name instead of using the default provider's rates (serrrfirat). - fire_manual() now uses per-user workspace, matching spawn_fire() pattern (serrrfirat). - Removed MULTI_TENANT env var — multi-tenant mode is auto-detected solely from GATEWAY_USER_TOKENS presence (serrrfirat + Copilot). - Multi-user heartbeat capped at 8 concurrent tasks to avoid flooding the LLM provider (serrrfirat + Copilot). - Fixed inject_model_override doc comment accuracy (Copilot). - Added comment explaining multi-tenant notification routing priority (Copilot). Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: user-scoped webhook endpoint for multi-tenant isolation Adds POST /api/webhooks/u/{user_id}/{path} — a user-scoped webhook endpoint that filters the routine lookup by user_id, preventing cross-user webhook triggering when paths collide. The existing /api/webhooks/{path} endpoint remains unchanged for backward compatibility in single-user deployments. Changes: - get_webhook_routine_by_path gains user_id: Option<&str> param - Both postgres and libsql implementations add AND user_id = ? filter when user_id is provided - New webhook_trigger_user_scoped_handler extracts (user_id, path) from URL and passes to shared fire_webhook_inner logic - Route registered on public router (webhooks are called by external services that can't send bearer tokens) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat: add TenantCtx for compile-time tenant isolation Implements zmanian's architectural proposal from #1614 review: two-tier scoped database access (TenantScope/AdminScope) so handler code cannot accidentally bypass tenant scoping. TenantScope (default): wraps user_id + Arc<dyn Database>, auto-binds user_id on every operation. ID-based lookups return None for cross- tenant resources. No escape hatch — forgetting to scope is a compile error. AdminScope (explicit opt-in): cross-tenant access for system-level components (heartbeat, routine engine, self-repair, scheduler, worker). TenantCtx bundles TenantScope + workspace + cost guard + per-user rate limiting. Constructed once per request in handle_message, threaded through all command handlers and ChatDelegate. Key changes: - New src/tenant.rs (~920 lines): TenantScope, AdminScope, TenantCtx, TenantRateState, TenantRateRegistry - All command handlers: user_id: &str → ctx: &TenantCtx - ChatDelegate: cost check/record/settings via self.tenant - System components: store field changed to AdminScope - Config: TENANT_MAX_LLM_CONCURRENT, TENANT_MAX_JOBS_CONCURRENT env vars - Fixes bug: /status <job_id> cross-tenant leak (now auto-filtered) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> --------- Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
597 lines
20 KiB
Rust
597 lines
20 KiB
Rust
//! Routine-related RoutineStore implementation for LibSqlBackend.
|
|
|
|
use std::collections::{HashMap, HashSet};
|
|
|
|
use async_trait::async_trait;
|
|
use chrono::{DateTime, Utc};
|
|
use libsql::params;
|
|
use uuid::Uuid;
|
|
|
|
use super::{
|
|
LibSqlBackend, ROUTINE_COLUMNS, ROUTINE_RUN_COLUMNS, fmt_opt_ts, fmt_ts, get_i64, get_text,
|
|
opt_text, opt_text_owned, row_to_routine_libsql, row_to_routine_run_libsql,
|
|
};
|
|
use crate::agent::routine::{Routine, RoutineRun, RunStatus};
|
|
use crate::db::RoutineStore;
|
|
use crate::error::DatabaseError;
|
|
|
|
#[async_trait]
|
|
impl RoutineStore for LibSqlBackend {
|
|
async fn create_routine(&self, routine: &Routine) -> Result<(), DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let trigger_type = routine.trigger.type_tag();
|
|
let trigger_config = routine.trigger.to_config_json();
|
|
let action_type = routine.action.type_tag();
|
|
let action_config = routine.action.to_config_json();
|
|
let cooldown_secs = routine.guardrails.cooldown.as_secs() as i64;
|
|
let max_concurrent = routine.guardrails.max_concurrent as i64;
|
|
let dedup_window_secs = routine.guardrails.dedup_window.map(|d| d.as_secs() as i64);
|
|
|
|
conn.execute(
|
|
r#"
|
|
INSERT INTO routines (
|
|
id, name, description, user_id, enabled,
|
|
trigger_type, trigger_config, action_type, action_config,
|
|
cooldown_secs, max_concurrent, dedup_window_secs,
|
|
notify_channel, notify_user, notify_on_success, notify_on_failure, notify_on_attention,
|
|
state, next_fire_at, created_at, updated_at
|
|
) VALUES (
|
|
?1, ?2, ?3, ?4, ?5,
|
|
?6, ?7, ?8, ?9,
|
|
?10, ?11, ?12,
|
|
?13, ?14, ?15, ?16, ?17,
|
|
?18, ?19, ?20, ?21
|
|
)
|
|
"#,
|
|
params![
|
|
routine.id.to_string(),
|
|
routine.name.as_str(),
|
|
routine.description.as_str(),
|
|
routine.user_id.as_str(),
|
|
routine.enabled as i64,
|
|
trigger_type,
|
|
trigger_config.to_string(),
|
|
action_type,
|
|
action_config.to_string(),
|
|
cooldown_secs,
|
|
max_concurrent,
|
|
dedup_window_secs,
|
|
opt_text(routine.notify.channel.as_deref()),
|
|
opt_text(routine.notify.user.as_deref()),
|
|
routine.notify.on_success as i64,
|
|
routine.notify.on_failure as i64,
|
|
routine.notify.on_attention as i64,
|
|
routine.state.to_string(),
|
|
fmt_opt_ts(&routine.next_fire_at),
|
|
fmt_ts(&routine.created_at),
|
|
fmt_ts(&routine.updated_at),
|
|
],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(())
|
|
}
|
|
|
|
async fn get_routine(&self, id: Uuid) -> Result<Option<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!("SELECT {} FROM routines WHERE id = ?1", ROUTINE_COLUMNS),
|
|
params![id.to_string()],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
match rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
Some(row) => Ok(Some(row_to_routine_libsql(&row)?)),
|
|
None => Ok(None),
|
|
}
|
|
}
|
|
|
|
async fn get_routine_by_name(
|
|
&self,
|
|
user_id: &str,
|
|
name: &str,
|
|
) -> Result<Option<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!(
|
|
"SELECT {} FROM routines WHERE user_id = ?1 AND name = ?2",
|
|
ROUTINE_COLUMNS
|
|
),
|
|
params![user_id, name],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
match rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
Some(row) => Ok(Some(row_to_routine_libsql(&row)?)),
|
|
None => Ok(None),
|
|
}
|
|
}
|
|
|
|
async fn list_routines(&self, user_id: &str) -> Result<Vec<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!(
|
|
"SELECT {} FROM routines WHERE user_id = ?1 ORDER BY name",
|
|
ROUTINE_COLUMNS
|
|
),
|
|
params![user_id],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
let mut routines = Vec::new();
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
routines.push(row_to_routine_libsql(&row)?);
|
|
}
|
|
Ok(routines)
|
|
}
|
|
|
|
async fn list_all_routines(&self) -> Result<Vec<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!("SELECT {} FROM routines ORDER BY name", ROUTINE_COLUMNS),
|
|
(),
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
let mut routines = Vec::new();
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
routines.push(row_to_routine_libsql(&row)?);
|
|
}
|
|
Ok(routines)
|
|
}
|
|
|
|
async fn list_event_routines(&self) -> Result<Vec<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!(
|
|
"SELECT {} FROM routines WHERE enabled = 1 AND trigger_type IN ('event', 'system_event')",
|
|
ROUTINE_COLUMNS
|
|
),
|
|
(),
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
let mut routines = Vec::new();
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
routines.push(row_to_routine_libsql(&row)?);
|
|
}
|
|
Ok(routines)
|
|
}
|
|
|
|
async fn list_due_cron_routines(&self) -> Result<Vec<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let now = fmt_ts(&Utc::now());
|
|
let mut rows = conn
|
|
.query(
|
|
&format!(
|
|
"SELECT {} FROM routines WHERE enabled = 1 AND trigger_type = 'cron' AND next_fire_at IS NOT NULL AND next_fire_at <= ?1",
|
|
ROUTINE_COLUMNS
|
|
),
|
|
params![now],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
let mut routines = Vec::new();
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
routines.push(row_to_routine_libsql(&row)?);
|
|
}
|
|
Ok(routines)
|
|
}
|
|
|
|
async fn update_routine(&self, routine: &Routine) -> Result<(), DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let trigger_type = routine.trigger.type_tag();
|
|
let trigger_config = routine.trigger.to_config_json();
|
|
let action_type = routine.action.type_tag();
|
|
let action_config = routine.action.to_config_json();
|
|
let cooldown_secs = routine.guardrails.cooldown.as_secs() as i64;
|
|
let max_concurrent = routine.guardrails.max_concurrent as i64;
|
|
let dedup_window_secs = routine.guardrails.dedup_window.map(|d| d.as_secs() as i64);
|
|
let now = fmt_ts(&Utc::now());
|
|
|
|
conn.execute(
|
|
r#"
|
|
UPDATE routines SET
|
|
name = ?2, description = ?3, enabled = ?4,
|
|
trigger_type = ?5, trigger_config = ?6,
|
|
action_type = ?7, action_config = ?8,
|
|
cooldown_secs = ?9, max_concurrent = ?10, dedup_window_secs = ?11,
|
|
notify_channel = ?12, notify_user = ?13,
|
|
notify_on_success = ?14, notify_on_failure = ?15, notify_on_attention = ?16,
|
|
state = ?17, next_fire_at = ?18,
|
|
updated_at = ?19
|
|
WHERE id = ?1
|
|
"#,
|
|
params![
|
|
routine.id.to_string(),
|
|
routine.name.as_str(),
|
|
routine.description.as_str(),
|
|
routine.enabled as i64,
|
|
trigger_type,
|
|
trigger_config.to_string(),
|
|
action_type,
|
|
action_config.to_string(),
|
|
cooldown_secs,
|
|
max_concurrent,
|
|
dedup_window_secs,
|
|
opt_text(routine.notify.channel.as_deref()),
|
|
opt_text(routine.notify.user.as_deref()),
|
|
routine.notify.on_success as i64,
|
|
routine.notify.on_failure as i64,
|
|
routine.notify.on_attention as i64,
|
|
routine.state.to_string(),
|
|
fmt_opt_ts(&routine.next_fire_at),
|
|
now,
|
|
],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(())
|
|
}
|
|
|
|
async fn update_routine_runtime(
|
|
&self,
|
|
id: Uuid,
|
|
last_run_at: DateTime<Utc>,
|
|
next_fire_at: Option<DateTime<Utc>>,
|
|
run_count: u64,
|
|
consecutive_failures: u32,
|
|
state: &serde_json::Value,
|
|
) -> Result<(), DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let now = fmt_ts(&Utc::now());
|
|
conn.execute(
|
|
r#"
|
|
UPDATE routines SET
|
|
last_run_at = ?2, next_fire_at = ?3,
|
|
run_count = ?4, consecutive_failures = ?5,
|
|
state = ?6, updated_at = ?7
|
|
WHERE id = ?1
|
|
"#,
|
|
params![
|
|
id.to_string(),
|
|
fmt_ts(&last_run_at),
|
|
fmt_opt_ts(&next_fire_at),
|
|
run_count as i64,
|
|
consecutive_failures as i64,
|
|
state.to_string(),
|
|
now,
|
|
],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(())
|
|
}
|
|
|
|
async fn delete_routine(&self, id: Uuid) -> Result<bool, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let count = conn
|
|
.execute(
|
|
"DELETE FROM routines WHERE id = ?1",
|
|
params![id.to_string()],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(count > 0)
|
|
}
|
|
|
|
async fn create_routine_run(&self, run: &RoutineRun) -> Result<(), DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
conn.execute(
|
|
r#"
|
|
INSERT INTO routine_runs (
|
|
id, routine_id, trigger_type, trigger_detail,
|
|
started_at, status, job_id
|
|
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
|
|
"#,
|
|
params![
|
|
run.id.to_string(),
|
|
run.routine_id.to_string(),
|
|
run.trigger_type.as_str(),
|
|
opt_text(run.trigger_detail.as_deref()),
|
|
fmt_ts(&run.started_at),
|
|
run.status.to_string(),
|
|
opt_text_owned(run.job_id.map(|id| id.to_string())),
|
|
],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(())
|
|
}
|
|
|
|
async fn complete_routine_run(
|
|
&self,
|
|
id: Uuid,
|
|
status: RunStatus,
|
|
result_summary: Option<&str>,
|
|
tokens_used: Option<i32>,
|
|
) -> Result<(), DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let now = fmt_ts(&Utc::now());
|
|
conn.execute(
|
|
r#"
|
|
UPDATE routine_runs SET
|
|
completed_at = ?5, status = ?2,
|
|
result_summary = ?3, tokens_used = ?4
|
|
WHERE id = ?1
|
|
"#,
|
|
params![
|
|
id.to_string(),
|
|
status.to_string(),
|
|
opt_text(result_summary),
|
|
tokens_used.map(|t| t as i64),
|
|
now,
|
|
],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(())
|
|
}
|
|
|
|
async fn list_routine_runs(
|
|
&self,
|
|
routine_id: Uuid,
|
|
limit: i64,
|
|
) -> Result<Vec<RoutineRun>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!(
|
|
"SELECT {} FROM routine_runs WHERE routine_id = ?1 ORDER BY started_at DESC LIMIT ?2",
|
|
ROUTINE_RUN_COLUMNS
|
|
),
|
|
params![routine_id.to_string(), limit],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
let mut runs = Vec::new();
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
runs.push(row_to_routine_run_libsql(&row)?);
|
|
}
|
|
Ok(runs)
|
|
}
|
|
|
|
async fn count_running_routine_runs(&self, routine_id: Uuid) -> Result<i64, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
"SELECT COUNT(*) as cnt FROM routine_runs WHERE routine_id = ?1 AND status = 'running'",
|
|
params![routine_id.to_string()],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
match rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
Some(row) => Ok(get_i64(&row, 0)),
|
|
None => Ok(0),
|
|
}
|
|
}
|
|
|
|
async fn count_running_routine_runs_batch(
|
|
&self,
|
|
routine_ids: &[Uuid],
|
|
) -> Result<HashMap<Uuid, i64>, DatabaseError> {
|
|
if routine_ids.is_empty() {
|
|
return Ok(HashMap::new());
|
|
}
|
|
|
|
let mut counts = HashMap::new();
|
|
let conn = self.connect().await?;
|
|
|
|
// Query all running routines and filter in memory
|
|
// This is simpler for libSQL than building dynamic parameter lists
|
|
let mut rows = conn
|
|
.query(
|
|
"SELECT routine_id, COUNT(*) as cnt FROM routine_runs
|
|
WHERE status = 'running'
|
|
GROUP BY routine_id",
|
|
params![],
|
|
)
|
|
.await
|
|
.map_err(|e| {
|
|
DatabaseError::Query(format!("Failed to batch count running routines: {}", e))
|
|
})?;
|
|
|
|
let routine_id_set: HashSet<Uuid> = routine_ids.iter().copied().collect();
|
|
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
let id_str: String = get_text(&row, 0);
|
|
let id = Uuid::parse_str(&id_str)
|
|
.map_err(|e| DatabaseError::Query(format!("Invalid routine UUID: {}", e)))?;
|
|
|
|
// Only include if this routine ID was requested
|
|
if routine_id_set.contains(&id) {
|
|
let cnt: i64 = get_i64(&row, 1);
|
|
counts.insert(id, cnt);
|
|
}
|
|
}
|
|
|
|
// Ensure all requested IDs are in the map (defaults to 0 for no running runs)
|
|
for id in routine_ids {
|
|
counts.entry(*id).or_insert(0);
|
|
}
|
|
|
|
Ok(counts)
|
|
}
|
|
|
|
async fn batch_get_last_run_status(
|
|
&self,
|
|
routine_ids: &[Uuid],
|
|
) -> Result<HashMap<Uuid, RunStatus>, DatabaseError> {
|
|
if routine_ids.is_empty() {
|
|
return Ok(HashMap::new());
|
|
}
|
|
|
|
let conn = self.connect().await?;
|
|
|
|
// SQLite doesn't support ANY($1), so we query all latest runs and filter in memory.
|
|
// Uses a subquery to pick only the most recent run per routine.
|
|
let mut rows = conn
|
|
.query(
|
|
"SELECT routine_id, status FROM routine_runs r1
|
|
WHERE started_at = (
|
|
SELECT MAX(started_at) FROM routine_runs r2
|
|
WHERE r2.routine_id = r1.routine_id
|
|
)
|
|
GROUP BY routine_id",
|
|
params![],
|
|
)
|
|
.await
|
|
.map_err(|e| {
|
|
DatabaseError::Query(format!("Failed to batch get last run status: {}", e))
|
|
})?;
|
|
|
|
let routine_id_set: HashSet<Uuid> = routine_ids.iter().copied().collect();
|
|
let mut statuses = HashMap::new();
|
|
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
let id_str: String = get_text(&row, 0);
|
|
let id = Uuid::parse_str(&id_str)
|
|
.map_err(|e| DatabaseError::Query(format!("Invalid routine UUID: {}", e)))?;
|
|
|
|
if routine_id_set.contains(&id) {
|
|
let status_str: String = get_text(&row, 1);
|
|
if let std::result::Result::Ok(status) = status_str.parse::<RunStatus>() {
|
|
statuses.insert(id, status);
|
|
}
|
|
}
|
|
}
|
|
|
|
Ok(statuses)
|
|
}
|
|
|
|
async fn link_routine_run_to_job(
|
|
&self,
|
|
run_id: Uuid,
|
|
job_id: Uuid,
|
|
) -> Result<(), DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
conn.execute(
|
|
"UPDATE routine_runs SET job_id = ?1 WHERE id = ?2",
|
|
params![job_id.to_string(), run_id.to_string()],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
Ok(())
|
|
}
|
|
|
|
async fn get_webhook_routine_by_path(
|
|
&self,
|
|
path: &str,
|
|
user_id: Option<&str>,
|
|
) -> Result<Option<Routine>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = if let Some(uid) = user_id {
|
|
conn.query(
|
|
&format!(
|
|
"SELECT {} FROM routines WHERE enabled = 1 AND trigger_type = 'webhook' \
|
|
AND user_id = ?2 \
|
|
AND (json_extract(trigger_config, '$.path') = ?1 \
|
|
OR (json_extract(trigger_config, '$.path') IS NULL AND CAST(id AS TEXT) = ?1))",
|
|
ROUTINE_COLUMNS
|
|
),
|
|
params![path, uid],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
} else {
|
|
conn.query(
|
|
&format!(
|
|
"SELECT {} FROM routines WHERE enabled = 1 AND trigger_type = 'webhook' \
|
|
AND (json_extract(trigger_config, '$.path') = ?1 \
|
|
OR (json_extract(trigger_config, '$.path') IS NULL AND CAST(id AS TEXT) = ?1))",
|
|
ROUTINE_COLUMNS
|
|
),
|
|
params![path],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
};
|
|
|
|
match rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
Some(row) => Ok(Some(row_to_routine_libsql(&row)?)),
|
|
None => Ok(None),
|
|
}
|
|
}
|
|
|
|
async fn list_dispatched_routine_runs(&self) -> Result<Vec<RoutineRun>, DatabaseError> {
|
|
let conn = self.connect().await?;
|
|
let mut rows = conn
|
|
.query(
|
|
&format!(
|
|
"SELECT {} FROM routine_runs WHERE status = 'running' AND job_id IS NOT NULL",
|
|
ROUTINE_RUN_COLUMNS
|
|
),
|
|
params![],
|
|
)
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?;
|
|
|
|
let mut runs = Vec::new();
|
|
while let Some(row) = rows
|
|
.next()
|
|
.await
|
|
.map_err(|e| DatabaseError::Query(e.to_string()))?
|
|
{
|
|
runs.push(row_to_routine_run_libsql(&row)?);
|
|
}
|
|
Ok(runs)
|
|
}
|
|
}
|