use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime};
use async_trait::async_trait;
use serde_json::Value;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub enum WorkSchedule {
#[default]
Immediate,
At(SystemTime),
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub enum WakePolicy {
#[default]
Never,
OnCompletion,
}
#[derive(Clone, Debug, PartialEq)]
pub struct TaskRequest {
kind: String,
input: Value,
schedule: WorkSchedule,
idempotency_key: Option<String>,
wake_policy: WakePolicy,
}
impl TaskRequest {
pub fn new(kind: impl Into<String>, input: impl Into<Value>) -> Self {
Self {
kind: kind.into(),
input: input.into(),
schedule: WorkSchedule::Immediate,
idempotency_key: None,
wake_policy: WakePolicy::Never,
}
}
pub fn schedule(mut self, schedule: WorkSchedule) -> Self {
self.schedule = schedule;
self
}
pub fn idempotency_key(mut self, key: impl Into<String>) -> Self {
self.idempotency_key = Some(key.into());
self
}
pub fn wake_policy(mut self, policy: WakePolicy) -> Self {
self.wake_policy = policy;
self
}
pub fn kind(&self) -> &str {
&self.kind
}
pub fn input(&self) -> &Value {
&self.input
}
pub fn work_schedule(&self) -> WorkSchedule {
self.schedule
}
pub fn submission_key(&self) -> Option<&str> {
self.idempotency_key.as_deref()
}
pub fn completion_wake_policy(&self) -> WakePolicy {
self.wake_policy
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum TaskState {
Pending,
Running,
Succeeded,
Failed,
Canceled,
}
impl TaskState {
pub fn is_terminal(self) -> bool {
matches!(self, Self::Succeeded | Self::Failed | Self::Canceled)
}
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub enum TaskOutcome {
Succeeded {
summary: Option<String>,
output: Value,
},
Failed {
error: String,
},
Canceled,
}
impl TaskOutcome {
pub fn success(output: impl Into<Value>) -> Self {
Self::Succeeded {
summary: None,
output: output.into(),
}
}
pub fn success_with_summary(summary: impl Into<String>, output: impl Into<Value>) -> Self {
Self::Succeeded {
summary: Some(summary.into()),
output: output.into(),
}
}
pub fn failure(error: impl Into<String>) -> Self {
Self::Failed {
error: error.into(),
}
}
pub fn task_state(&self) -> TaskState {
match self {
Self::Succeeded { .. } => TaskState::Succeeded,
Self::Failed { .. } => TaskState::Failed,
Self::Canceled => TaskState::Canceled,
}
}
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub struct Task {
pub id: String,
pub session_id: String,
pub kind: String,
pub input: Value,
pub state: TaskState,
pub cancel_requested_at: Option<SystemTime>,
pub idempotency_key: Option<String>,
pub scheduled_at: SystemTime,
pub created_at: SystemTime,
pub attempts: u32,
pub finished_at: Option<SystemTime>,
pub outcome: Option<TaskOutcome>,
}
impl Task {
pub fn pending(
id: impl Into<String>,
session_id: impl Into<String>,
request: &TaskRequest,
created_at: SystemTime,
) -> Self {
let scheduled_at = match request.schedule {
WorkSchedule::Immediate => created_at,
WorkSchedule::At(at) => at,
};
Self {
id: id.into(),
session_id: session_id.into(),
kind: request.kind.clone(),
input: request.input.clone(),
state: TaskState::Pending,
cancel_requested_at: None,
idempotency_key: request.idempotency_key.clone(),
scheduled_at,
created_at,
attempts: 0,
finished_at: None,
outcome: None,
}
}
}
#[derive(Clone, PartialEq)]
pub struct TaskDelivery {
pub task: Task,
pub attempt: u32,
pub lease_expires_at: SystemTime,
lease_token: String,
}
impl std::fmt::Debug for TaskDelivery {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TaskDelivery")
.field("task", &self.task)
.field("attempt", &self.attempt)
.field("lease_expires_at", &self.lease_expires_at)
.field("lease_token", &"[redacted]")
.finish()
}
}
impl TaskDelivery {
pub fn from_claim(
task: Task,
attempt: u32,
lease_expires_at: SystemTime,
lease_token: impl Into<String>,
) -> Self {
Self {
task,
attempt,
lease_expires_at,
lease_token: lease_token.into(),
}
}
pub fn lease_token(&self) -> &str {
&self.lease_token
}
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub enum WakeReason {
Requested {
payload: Value,
},
TaskFinished {
task_id: String,
state: TaskState,
outcome: TaskOutcome,
},
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub struct SessionWake {
pub id: String,
pub session_id: String,
pub reason: WakeReason,
pub created_at: SystemTime,
pub idempotency_key: Option<String>,
}
impl SessionWake {
pub fn requested(
id: impl Into<String>,
session_id: impl Into<String>,
request: &WakeRequest,
created_at: SystemTime,
) -> Self {
Self {
id: id.into(),
session_id: session_id.into(),
reason: WakeReason::Requested {
payload: request.payload.clone(),
},
created_at,
idempotency_key: request.idempotency_key.clone(),
}
}
pub fn task_finished(
id: impl Into<String>,
task: &Task,
created_at: SystemTime,
) -> Option<Self> {
let outcome = task.outcome.clone()?;
if !task.state.is_terminal() || outcome.task_state() != task.state {
return None;
}
Some(Self {
id: id.into(),
session_id: task.session_id.clone(),
reason: WakeReason::TaskFinished {
task_id: task.id.clone(),
state: task.state,
outcome,
},
created_at,
idempotency_key: None,
})
}
}
#[derive(Clone, Debug, PartialEq)]
pub struct WakeRequest {
payload: Value,
idempotency_key: Option<String>,
}
impl WakeRequest {
pub fn new(payload: impl Into<Value>) -> Self {
Self {
payload: payload.into(),
idempotency_key: None,
}
}
pub fn idempotency_key(mut self, key: impl Into<String>) -> Self {
self.idempotency_key = Some(key.into());
self
}
pub fn payload(&self) -> &Value {
&self.payload
}
pub fn submission_key(&self) -> Option<&str> {
self.idempotency_key.as_deref()
}
}
#[derive(Clone, PartialEq)]
pub struct WakeDelivery {
pub wake: SessionWake,
pub attempt: u32,
pub lease_expires_at: SystemTime,
lease_token: String,
}
impl std::fmt::Debug for WakeDelivery {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("WakeDelivery")
.field("wake", &self.wake)
.field("attempt", &self.attempt)
.field("lease_expires_at", &self.lease_expires_at)
.field("lease_token", &"[redacted]")
.finish()
}
}
impl WakeDelivery {
pub fn from_claim(
wake: SessionWake,
attempt: u32,
lease_expires_at: SystemTime,
lease_token: impl Into<String>,
) -> Self {
Self {
wake,
attempt,
lease_expires_at,
lease_token: lease_token.into(),
}
}
pub fn lease_token(&self) -> &str {
&self.lease_token
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum WorkError {
InvalidRequest {
field: &'static str,
message: String,
},
TaskNotFound {
task_id: String,
},
IdempotencyConflict {
key: String,
},
StaleDelivery {
id: String,
},
InvalidTransition {
message: String,
},
Backend {
message: String,
},
}
impl WorkError {
pub fn backend(message: impl Into<String>) -> Self {
Self::Backend {
message: message.into(),
}
}
}
impl std::fmt::Display for WorkError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::InvalidRequest { field, message } => {
write!(f, "invalid work request field {field}: {message}")
}
Self::TaskNotFound { task_id } => write!(f, "task not found: {task_id}"),
Self::IdempotencyConflict { key } => {
write!(f, "idempotency key was reused for different work: {key}")
}
Self::StaleDelivery { id } => write!(f, "delivery no longer owns {id}"),
Self::InvalidTransition { message } => write!(f, "invalid work transition: {message}"),
Self::Backend { message } => write!(f, "work backend failed: {message}"),
}
}
}
impl std::error::Error for WorkError {}
#[async_trait]
pub trait WorkBackend: Send + Sync {
async fn submit(
&self,
session_id: &str,
request: TaskRequest,
now: SystemTime,
) -> Result<Task, WorkError>;
async fn task(&self, session_id: &str, task_id: &str) -> Result<Option<Task>, WorkError>;
async fn cancel(
&self,
session_id: &str,
task_id: &str,
now: SystemTime,
) -> Result<Task, WorkError>;
async fn claim_due(
&self,
now: SystemTime,
lease_for: Duration,
limit: usize,
) -> Result<Vec<TaskDelivery>, WorkError>;
async fn finish(
&self,
delivery: &TaskDelivery,
outcome: TaskOutcome,
now: SystemTime,
) -> Result<Task, WorkError>;
async fn request_wake(
&self,
session_id: &str,
request: WakeRequest,
now: SystemTime,
) -> Result<SessionWake, WorkError>;
async fn claim_wakes(
&self,
now: SystemTime,
lease_for: Duration,
limit: usize,
) -> Result<Vec<WakeDelivery>, WorkError>;
async fn acknowledge_wake(&self, delivery: &WakeDelivery) -> Result<(), WorkError>;
}
#[derive(Clone)]
pub struct WorkQueue {
backend: Arc<dyn WorkBackend>,
}
impl WorkQueue {
pub fn in_memory() -> Self {
Self::with_backend(Arc::new(InMemoryWorkBackend::new()))
}
pub fn with_backend(backend: Arc<dyn WorkBackend>) -> Self {
Self { backend }
}
pub fn for_session(&self, session_id: impl Into<String>) -> SessionWork {
SessionWork {
queue: self.clone(),
session_id: session_id.into(),
}
}
pub async fn claim_due(
&self,
lease_for: Duration,
limit: usize,
) -> Result<Vec<TaskDelivery>, WorkError> {
self.claim_due_at(SystemTime::now(), lease_for, limit).await
}
pub async fn claim_due_at(
&self,
now: SystemTime,
lease_for: Duration,
limit: usize,
) -> Result<Vec<TaskDelivery>, WorkError> {
validate_claim(lease_for, limit)?;
self.backend.claim_due(now, lease_for, limit).await
}
pub async fn finish(
&self,
delivery: &TaskDelivery,
outcome: TaskOutcome,
) -> Result<Task, WorkError> {
self.finish_at(delivery, outcome, SystemTime::now()).await
}
pub async fn finish_at(
&self,
delivery: &TaskDelivery,
outcome: TaskOutcome,
now: SystemTime,
) -> Result<Task, WorkError> {
self.backend.finish(delivery, outcome, now).await
}
pub async fn claim_wakes(
&self,
lease_for: Duration,
limit: usize,
) -> Result<Vec<WakeDelivery>, WorkError> {
self.claim_wakes_at(SystemTime::now(), lease_for, limit)
.await
}
pub async fn claim_wakes_at(
&self,
now: SystemTime,
lease_for: Duration,
limit: usize,
) -> Result<Vec<WakeDelivery>, WorkError> {
validate_claim(lease_for, limit)?;
self.backend.claim_wakes(now, lease_for, limit).await
}
pub async fn acknowledge_wake(&self, delivery: &WakeDelivery) -> Result<(), WorkError> {
self.backend.acknowledge_wake(delivery).await
}
}
impl Default for WorkQueue {
fn default() -> Self {
Self::in_memory()
}
}
#[derive(Clone)]
pub struct SessionWork {
queue: WorkQueue,
session_id: String,
}
impl SessionWork {
pub fn session_id(&self) -> &str {
&self.session_id
}
pub async fn submit(&self, request: TaskRequest) -> Result<Task, WorkError> {
self.submit_at(request, SystemTime::now()).await
}
pub async fn submit_at(
&self,
request: TaskRequest,
now: SystemTime,
) -> Result<Task, WorkError> {
validate_session_id(&self.session_id)?;
validate_task_request(&request)?;
self.queue
.backend
.submit(&self.session_id, request, now)
.await
}
pub async fn task(&self, task_id: &str) -> Result<Option<Task>, WorkError> {
validate_session_id(&self.session_id)?;
validate_task_id(task_id)?;
self.queue.backend.task(&self.session_id, task_id).await
}
pub async fn cancel(&self, task_id: &str) -> Result<Task, WorkError> {
self.cancel_at(task_id, SystemTime::now()).await
}
pub async fn cancel_at(&self, task_id: &str, now: SystemTime) -> Result<Task, WorkError> {
validate_session_id(&self.session_id)?;
validate_task_id(task_id)?;
self.queue
.backend
.cancel(&self.session_id, task_id, now)
.await
}
pub async fn wake(&self, request: WakeRequest) -> Result<SessionWake, WorkError> {
self.wake_at(request, SystemTime::now()).await
}
pub async fn wake_at(
&self,
request: WakeRequest,
now: SystemTime,
) -> Result<SessionWake, WorkError> {
validate_session_id(&self.session_id)?;
validate_wake_request(&request)?;
self.queue
.backend
.request_wake(&self.session_id, request, now)
.await
}
}
fn validate_session_id(session_id: &str) -> Result<(), WorkError> {
validate_non_blank("session_id", session_id)
}
fn validate_task_id(task_id: &str) -> Result<(), WorkError> {
validate_non_blank("task_id", task_id)
}
fn validate_non_blank(field: &'static str, value: &str) -> Result<(), WorkError> {
if value.trim().is_empty() {
return Err(WorkError::InvalidRequest {
field,
message: "must not be blank".to_string(),
});
}
Ok(())
}
fn validate_task_request(request: &TaskRequest) -> Result<(), WorkError> {
if request.kind.trim().is_empty() {
return Err(WorkError::InvalidRequest {
field: "kind",
message: "must not be blank".to_string(),
});
}
if request
.idempotency_key
.as_ref()
.is_some_and(|key| key.trim().is_empty())
{
return Err(WorkError::InvalidRequest {
field: "idempotency_key",
message: "must not be blank".to_string(),
});
}
Ok(())
}
fn validate_wake_request(request: &WakeRequest) -> Result<(), WorkError> {
if request
.idempotency_key
.as_ref()
.is_some_and(|key| key.trim().is_empty())
{
return Err(WorkError::InvalidRequest {
field: "idempotency_key",
message: "must not be blank".to_string(),
});
}
Ok(())
}
fn validate_claim(lease_for: Duration, limit: usize) -> Result<(), WorkError> {
if lease_for.is_zero() {
return Err(WorkError::InvalidRequest {
field: "lease_for",
message: "must be positive".to_string(),
});
}
if limit == 0 {
return Err(WorkError::InvalidRequest {
field: "limit",
message: "must be positive".to_string(),
});
}
Ok(())
}
pub struct InMemoryWorkBackend {
state: Mutex<MemoryState>,
}
impl InMemoryWorkBackend {
pub fn new() -> Self {
Self {
state: Mutex::new(MemoryState::default()),
}
}
fn lock(&self) -> Result<std::sync::MutexGuard<'_, MemoryState>, WorkError> {
self.state
.lock()
.map_err(|_| WorkError::backend("in-memory state lock was poisoned"))
}
}
impl Default for InMemoryWorkBackend {
fn default() -> Self {
Self::new()
}
}
#[derive(Default)]
struct MemoryState {
tasks: HashMap<String, StoredTask>,
task_order: Vec<String>,
task_keys: HashMap<(String, String), String>,
wakes: HashMap<String, StoredWake>,
wake_order: Vec<String>,
wake_keys: HashMap<(String, String), String>,
}
struct StoredTask {
task: Task,
request: TaskRequest,
wake_policy: WakePolicy,
claim: Option<Lease>,
last_settlement: Option<(String, TaskOutcome)>,
}
struct StoredWake {
wake: SessionWake,
request: Option<WakeRequest>,
claim: Option<Lease>,
acknowledged_by: Option<String>,
attempts: u32,
}
#[derive(Clone)]
struct Lease {
token: String,
expires_at: SystemTime,
}
#[async_trait]
impl WorkBackend for InMemoryWorkBackend {
async fn submit(
&self,
session_id: &str,
request: TaskRequest,
now: SystemTime,
) -> Result<Task, WorkError> {
let mut state = self.lock()?;
if let Some(key) = request.idempotency_key.as_ref() {
let scope = (session_id.to_string(), key.clone());
if let Some(task_id) = state.task_keys.get(&scope) {
let stored = state
.tasks
.get(task_id)
.expect("key references stored task");
if stored.request == request {
return Ok(stored.task.clone());
}
return Err(WorkError::IdempotencyConflict { key: key.clone() });
}
}
let task = Task::pending(new_id("task"), session_id, &request, now);
if let Some(key) = request.idempotency_key.as_ref() {
state
.task_keys
.insert((session_id.to_string(), key.clone()), task.id.clone());
}
state.task_order.push(task.id.clone());
state.tasks.insert(
task.id.clone(),
StoredTask {
task: task.clone(),
wake_policy: request.wake_policy,
request,
claim: None,
last_settlement: None,
},
);
Ok(task)
}
async fn task(&self, session_id: &str, task_id: &str) -> Result<Option<Task>, WorkError> {
let state = self.lock()?;
Ok(state
.tasks
.get(task_id)
.filter(|stored| stored.task.session_id == session_id)
.map(|stored| stored.task.clone()))
}
async fn cancel(
&self,
session_id: &str,
task_id: &str,
now: SystemTime,
) -> Result<Task, WorkError> {
let mut state = self.lock()?;
let (task, wake_policy, should_wake) = {
let stored = state
.tasks
.get_mut(task_id)
.filter(|stored| stored.task.session_id == session_id)
.ok_or_else(|| WorkError::TaskNotFound {
task_id: task_id.to_string(),
})?;
if stored.task.state.is_terminal() {
return Ok(stored.task.clone());
}
stored.task.cancel_requested_at.get_or_insert(now);
let mut should_wake = false;
if stored.task.state == TaskState::Pending {
stored.task.state = TaskState::Canceled;
stored.task.finished_at = Some(now);
stored.task.outcome = Some(TaskOutcome::Canceled);
should_wake = stored.wake_policy == WakePolicy::OnCompletion;
}
(stored.task.clone(), stored.wake_policy, should_wake)
};
if should_wake {
push_task_wake(&mut state, &task, wake_policy, now);
}
Ok(task)
}
async fn claim_due(
&self,
now: SystemTime,
lease_for: Duration,
limit: usize,
) -> Result<Vec<TaskDelivery>, WorkError> {
let mut state = self.lock()?;
let mut deliveries = Vec::new();
for task_id in state.task_order.clone() {
if deliveries.len() == limit {
break;
}
let stored = state.tasks.get_mut(&task_id).expect("ordered task exists");
let pending_due =
stored.task.state == TaskState::Pending && stored.task.scheduled_at <= now;
let lease_expired = stored.task.state == TaskState::Running
&& stored
.claim
.as_ref()
.is_none_or(|claim| claim.expires_at <= now);
if !pending_due && !lease_expired {
continue;
}
stored.task.state = TaskState::Running;
stored.task.attempts = stored.task.attempts.saturating_add(1);
let attempt = stored.task.attempts;
let expires_at =
now.checked_add(lease_for)
.ok_or_else(|| WorkError::InvalidRequest {
field: "lease_for",
message: "lease expiration is outside SystemTime range".to_string(),
})?;
let token = new_id("lease");
stored.claim = Some(Lease {
token: token.clone(),
expires_at,
});
deliveries.push(TaskDelivery::from_claim(
stored.task.clone(),
attempt,
expires_at,
token,
));
}
Ok(deliveries)
}
async fn finish(
&self,
delivery: &TaskDelivery,
outcome: TaskOutcome,
now: SystemTime,
) -> Result<Task, WorkError> {
let mut state = self.lock()?;
let (task, wake_policy, should_wake) =
{
let stored = state.tasks.get_mut(&delivery.task.id).ok_or_else(|| {
WorkError::TaskNotFound {
task_id: delivery.task.id.clone(),
}
})?;
if let Some((token, settled)) = stored.last_settlement.as_ref()
&& token == delivery.lease_token()
{
if settled == &outcome {
return Ok(stored.task.clone());
}
return Err(WorkError::InvalidTransition {
message: "the delivery was already settled with a different outcome"
.to_string(),
});
}
let owns_claim = stored.task.state == TaskState::Running
&& stored.task.attempts == delivery.attempt
&& stored
.claim
.as_ref()
.is_some_and(|claim| claim.token == delivery.lease_token());
if !owns_claim {
return Err(WorkError::StaleDelivery {
id: delivery.task.id.clone(),
});
}
stored.task.state = outcome.task_state();
stored.task.finished_at = Some(now);
stored.task.outcome = Some(outcome.clone());
stored.claim = None;
stored.last_settlement = Some((delivery.lease_token().to_string(), outcome));
let should_wake = stored.wake_policy == WakePolicy::OnCompletion;
(stored.task.clone(), stored.wake_policy, should_wake)
};
if should_wake {
push_task_wake(&mut state, &task, wake_policy, now);
}
Ok(task)
}
async fn request_wake(
&self,
session_id: &str,
request: WakeRequest,
now: SystemTime,
) -> Result<SessionWake, WorkError> {
let mut state = self.lock()?;
if let Some(key) = request.idempotency_key.as_ref() {
let scope = (session_id.to_string(), key.clone());
if let Some(wake_id) = state.wake_keys.get(&scope) {
let stored = state
.wakes
.get(wake_id)
.expect("key references stored wake");
if stored.request.as_ref() == Some(&request) {
return Ok(stored.wake.clone());
}
return Err(WorkError::IdempotencyConflict { key: key.clone() });
}
}
let wake = SessionWake::requested(new_id("wake"), session_id, &request, now);
if let Some(key) = request.idempotency_key.as_ref() {
state
.wake_keys
.insert((session_id.to_string(), key.clone()), wake.id.clone());
}
push_wake(&mut state, wake.clone(), Some(request));
Ok(wake)
}
async fn claim_wakes(
&self,
now: SystemTime,
lease_for: Duration,
limit: usize,
) -> Result<Vec<WakeDelivery>, WorkError> {
let mut state = self.lock()?;
let mut deliveries = Vec::new();
for wake_id in state.wake_order.clone() {
if deliveries.len() == limit {
break;
}
let stored = state.wakes.get_mut(&wake_id).expect("ordered wake exists");
if stored.acknowledged_by.is_some() {
continue;
}
let available = stored
.claim
.as_ref()
.is_none_or(|claim| claim.expires_at <= now);
if !available {
continue;
}
stored.attempts = stored.attempts.saturating_add(1);
let expires_at =
now.checked_add(lease_for)
.ok_or_else(|| WorkError::InvalidRequest {
field: "lease_for",
message: "lease expiration is outside SystemTime range".to_string(),
})?;
let token = new_id("wakelease");
stored.claim = Some(Lease {
token: token.clone(),
expires_at,
});
deliveries.push(WakeDelivery::from_claim(
stored.wake.clone(),
stored.attempts,
expires_at,
token,
));
}
Ok(deliveries)
}
async fn acknowledge_wake(&self, delivery: &WakeDelivery) -> Result<(), WorkError> {
let mut state = self.lock()?;
let stored =
state
.wakes
.get_mut(&delivery.wake.id)
.ok_or_else(|| WorkError::StaleDelivery {
id: delivery.wake.id.clone(),
})?;
if stored.acknowledged_by.as_deref() == Some(delivery.lease_token()) {
return Ok(());
}
let owns_claim = stored.attempts == delivery.attempt
&& stored
.claim
.as_ref()
.is_some_and(|claim| claim.token == delivery.lease_token());
if !owns_claim {
return Err(WorkError::StaleDelivery {
id: delivery.wake.id.clone(),
});
}
stored.acknowledged_by = Some(delivery.lease_token().to_string());
stored.claim = None;
Ok(())
}
}
fn push_task_wake(state: &mut MemoryState, task: &Task, wake_policy: WakePolicy, now: SystemTime) {
if wake_policy != WakePolicy::OnCompletion {
return;
}
let Some(wake) = SessionWake::task_finished(new_id("wake"), task, now) else {
return;
};
push_wake(state, wake, None);
}
fn push_wake(state: &mut MemoryState, wake: SessionWake, request: Option<WakeRequest>) {
state.wake_order.push(wake.id.clone());
state.wakes.insert(
wake.id.clone(),
StoredWake {
wake,
request,
claim: None,
acknowledged_by: None,
attempts: 0,
},
);
}
fn new_id(prefix: &str) -> String {
let generated = everruns_core::session_task::generate_task_id();
format!("{prefix}_{}", generated.trim_start_matches("task_"))
}