use std::{sync::Arc, time::Duration};
use fraiseql_core::security::SecurityContext;
use serde_json::Value;
use tokio::time::MissedTickBehavior;
use tracing::{debug, warn};
use super::{AsyncOperationsRuntime, store::ClaimedOperation};
use crate::routes::graphql::{AppState, tenant_dispatch};
pub async fn run(runtime: AsyncOperationsRuntime, state: AppState) {
let mut ticker = tokio::time::interval(Duration::from_millis(runtime.config.poll_interval_ms));
ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
loop {
ticker.tick().await;
let claimed = match runtime.store.claim(1, runtime.config.stuck_threshold_secs).await {
Ok(c) => c,
Err(e) => {
warn!(error = %e, "async-operations: claim failed");
continue;
},
};
for op in claimed {
execute_claimed(&runtime, &state, op).await;
}
match runtime.store.sweep_finished(runtime.config.result_ttl_secs).await {
Ok(0) => {},
Ok(n) => debug!(swept = n, "async-operations: swept expired finished operations"),
Err(e) => warn!(error = %e, "async-operations: sweep failed"),
}
}
}
async fn execute_claimed(
runtime: &AsyncOperationsRuntime,
state: &AppState,
claimed: ClaimedOperation,
) {
let ClaimedOperation { op, claim_token } = claimed;
let store = Arc::clone(&runtime.store);
if op.cancellation_requested {
match store.cancel_unstarted(op.op_id, claim_token).await {
Ok(_) => return,
Err(e) => {
warn!(error = %e, op_id = %op.op_id, "async-operations: cancel-unstarted failed");
return;
},
}
}
let ctx: SecurityContext =
match serde_json::from_value(op.security_context.clone()) {
Ok(ctx) => ctx,
Err(e) => {
record_failure(&store, op.op_id, claim_token, &format!(
"stored security context no longer deserializes ({e}) — refusing to execute \
with an unverifiable principal"
), None)
.await;
return;
},
};
if ctx.is_expired() {
record_failure(
&store,
op.op_id,
claim_token,
"the submitter's security context expired before execution started — resubmit with \
a live credential",
None,
)
.await;
return;
}
let dispatch = match tenant_dispatch::dispatch_to_tenant(state, op.tenant_key.as_deref()) {
Ok(d) => d,
Err(e) => {
record_failure(&store, op.op_id, claim_token, &format!("tenant dispatch: {e}"), None)
.await;
return;
},
};
let hb_store = Arc::clone(&store);
let hb_interval = Duration::from_secs((runtime.config.stuck_threshold_secs / 3).max(1));
let op_id = op.op_id;
let execution =
dispatch
.executor
.execute_with_security(&op.document, op.variables.as_ref(), &ctx);
tokio::pin!(execution);
let mut hb_ticker = tokio::time::interval(hb_interval);
hb_ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
hb_ticker.tick().await; let exec_result = loop {
tokio::select! {
result = &mut execution => break result,
_ = hb_ticker.tick() => {
match hb_store.heartbeat(op_id, claim_token).await {
Ok(true) => {},
Ok(false) => {
warn!(op_id = %op_id, "async-operations: claim lost mid-execution");
return;
},
Err(e) => warn!(error = %e, "async-operations: heartbeat failed"),
}
}
}
};
match exec_result {
Ok(result) => {
let has_errors =
result.get("errors").and_then(Value::as_array).is_some_and(|e| !e.is_empty());
if has_errors {
let rendered = result["errors"].to_string();
record_failure(&store, op_id, claim_token, &rendered, Some(&result)).await;
} else {
match store.complete(op_id, claim_token, &result).await {
Ok(true) => {},
Ok(false) => {
warn!(op_id = %op_id, "async-operations: completion superseded (claim lost)");
},
Err(e) => {
warn!(error = %e, op_id = %op_id, "async-operations: complete failed");
},
}
}
},
Err(e) => {
warn!(error = %e, op_id = %op_id, "async-operations: execution failed");
let rendered = state
.error_sanitizer
.sanitize(crate::error::GraphQLError::from_fraiseql_error(&e))
.message;
record_failure(&store, op_id, claim_token, &rendered, None).await;
},
}
}
async fn record_failure(
store: &super::AsyncOperationStore,
op_id: uuid::Uuid,
claim_token: uuid::Uuid,
error: &str,
partial: Option<&Value>,
) {
match store.fail(op_id, claim_token, error, partial).await {
Ok(true) => {},
Ok(false) => warn!(op_id = %op_id, "async-operations: failure record superseded"),
Err(e) => warn!(error = %e, op_id = %op_id, "async-operations: fail-record failed"),
}
}