use crate::activity::activity::{
Activity, ActivityFuture, ActivityHandler, ActivityHandlerRegistry, ActivityOption,
ActivityPriority, OnDuplicate,
};
use crate::config::WorkerConfig;
use crate::queue::queue::{ActivityQueueTrait, ActivityResult, ResultState};
use crate::runner::error::WorkerError;
use crate::storage::{FailureKind, IdempotencyBehavior, QueuedActivity, Storage};
use crate::{ActivityContext, ActivityError};
use chrono::Utc;
#[cfg(feature = "postgres")]
use crate::observability::QueueInspector;
use futures::FutureExt;
use serde_json::json;
use std::panic::AssertUnwindSafe;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{watch, RwLock};
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
pub trait MetricsSink: Send + Sync + 'static {
fn inc_counter(&self, name: &str, value: u64);
fn observe_duration(&self, _name: &str, _dur: Duration) {
let _ = (_name, _dur);
}
}
pub struct NoopMetrics;
impl MetricsSink for NoopMetrics {
fn inc_counter(&self, _name: &str, _value: u64) {}
}
struct Backoff {
current: Duration,
base: Duration,
max: Duration,
}
impl Backoff {
fn new(base: Duration, max: Duration) -> Self {
Self {
current: base,
base,
max,
}
}
fn reset(&mut self) {
self.current = self.base;
}
fn next(&mut self) -> Duration {
let next = self.current;
self.current = (self.current.mul_f32(2.0)).min(self.max);
next
}
}
pub struct WorkerEngine {
activity_queue: Arc<dyn ActivityQueueTrait>,
backend: Arc<dyn Storage>,
activity_handlers: ActivityHandlerRegistry, config: WorkerConfig,
running: Arc<RwLock<bool>>, shutdown_tx: watch::Sender<bool>,
cancel_token: CancellationToken,
metrics: Arc<dyn MetricsSink>,
}
impl WorkerEngine {
pub fn new_with_backend(backend: Arc<dyn Storage>, config: WorkerConfig) -> Self {
let (shutdown_tx, _shutdown_rx) = watch::channel(false);
let adapter = Arc::new(BackendQueueAdapter::new(
backend.clone(),
config.activity_types.clone(),
));
Self {
activity_queue: adapter,
backend,
activity_handlers: ActivityHandlerRegistry::new(),
config,
running: Arc::new(RwLock::new(false)),
shutdown_tx,
cancel_token: CancellationToken::new(),
metrics: Arc::new(NoopMetrics),
}
}
pub fn with_metrics(&mut self, sink: Arc<dyn MetricsSink>) {
self.metrics = sink;
}
#[cfg(feature = "postgres")]
pub fn inspector(&self) -> QueueInspector {
QueueInspector::new(self.backend.clone())
.with_max_workers(self.config.max_concurrent_activities)
}
pub fn builder() -> WorkerEngineBuilder {
WorkerEngineBuilder::new()
}
pub async fn start(&self) -> Result<(), WorkerError> {
{
let mut running = self.running.write().await;
if *running {
return Err(WorkerError::AlreadyRunning);
}
*running = true;
}
if let Some(ref types) = self.config.activity_types {
let missing: Vec<_> = types
.iter()
.filter(|t| !self.activity_handlers.contains_key(t.as_str()))
.collect();
if !missing.is_empty() {
panic!(
"activity_types filter contains types with no registered handler: {:?}",
missing
);
}
}
info!(
max_concurrent_activities = self.config.max_concurrent_activities,
"Starting worker engine"
);
let mut join_handles = Vec::new();
if !self.activity_queue.schedules_natively() {
let scheduled_handle = self.start_scheduled_activities_processor().await;
join_handles.push(scheduled_handle);
}
let reaper_handle = self.start_reaper_processor().await;
join_handles.push(reaper_handle);
for worker_id in 0..self.config.max_concurrent_activities {
let handle = self.start_worker_loop(worker_id).await;
join_handles.push(handle);
}
tokio::select! {
_ = self.wait_for_shutdown() => {
info!("Shutdown signal received, stopping worker engine");
}
result = futures::future::try_join_all(join_handles) => {
match result {
Ok(_) => info!("All worker loops completed"),
Err(e) => error!(error = %e, "A worker task failed"),
}
}
}
self.stop().await;
info!("Worker engine stopped");
Ok(())
}
pub async fn stop(&self) {
info!("Stopping worker engine");
let mut running = self.running.write().await;
*running = false;
let _ = self.shutdown_tx.send(true);
self.cancel_token.cancel();
}
async fn start_worker_loop(
&self,
worker_id: usize,
) -> tokio::task::JoinHandle<Result<(), WorkerError>> {
let running = self.running.clone();
let activity_queue = self.activity_queue.clone();
let activity_handlers = self.activity_handlers.clone();
let activity_queue_for_context = self.activity_queue.clone();
let mut shutdown_rx = self.shutdown_tx.subscribe();
let cancel_token = self.cancel_token.clone();
let metrics = self.metrics.clone();
tokio::spawn(async move {
debug!(%worker_id, "Starting worker loop");
let worker_label = format!("worker-{}", worker_id);
let mut backoff = Backoff::new(Duration::from_millis(100), Duration::from_secs(5));
while *running.read().await {
if *shutdown_rx.borrow() {
break;
}
let dequeue_fut = activity_queue.dequeue(Duration::from_secs(1), &worker_label);
let activity_opt = tokio::select! {
_ = shutdown_rx.changed() => { break; }
res = dequeue_fut => res
};
let activity = match activity_opt {
Ok(Some(a)) => {
backoff.reset();
a
}
Ok(None) => {
let sleep_for = backoff.next();
tokio::select! {
_ = tokio::time::sleep(sleep_for) => {},
_ = shutdown_rx.changed() => break,
}
continue;
}
Err(e) => {
error!(%worker_id, error = %e, "Failed to dequeue activity");
tokio::select! {
_ = tokio::time::sleep(Duration::from_secs(1)) => {},
_ = shutdown_rx.changed() => break,
}
continue;
}
};
let activity_id = activity.id;
let activity_type = activity.activity_type.clone();
debug!(%worker_id, activity_id = %activity_id, activity_type = ?activity_type, "Worker processing activity");
let handler = match activity_handlers.get(&activity.activity_type) {
Some(h) => h.clone(),
None => {
error!(%worker_id, activity_id = %activity_id, activity_type = ?activity_type, "No handler found for activity type");
if let Err(e) = activity_queue
.mark_failed(
activity,
"handler_not_found".to_string(),
false,
worker_label.as_str(),
)
.await
{
error!(%worker_id, activity_id = %activity_id, error = %e, "Failed to mark activity as failed");
}
continue;
}
};
let context = ActivityContext {
activity_id,
activity_type: activity_type.clone(),
retry_count: activity.retry_count,
metadata: activity.metadata.clone(),
cancel_token: cancel_token.child_token(),
activity_executor: Arc::new(WorkerEngineWrapper::new(
activity_queue_for_context.clone(),
)),
};
let payload_for_dead_letter = activity.payload.clone();
let activity_timeout = Duration::from_secs(activity.timeout_seconds);
let handle_fut = handler.handle(activity.payload.clone(), context);
let timed = tokio::select! {
_ = shutdown_rx.changed() => {
break;
}
res = tokio::time::timeout(activity_timeout, async {
match AssertUnwindSafe(handle_fut).catch_unwind().await {
Ok(value) => value,
Err(e) => {
let err_msg = e.downcast_ref::<String>()
.map(|s| format!("panic: {}", s))
.or_else(|| e.downcast_ref::<&str>().map(|s| format!("panic: {}", s)))
.unwrap_or_else(|| "panic (unknown)".to_string());
Err(ActivityError::Retry(err_msg))
},
}
}) => res,
};
match timed {
Ok(Ok(value)) => {
metrics.inc_counter("activity_completed", 1);
if let Err(e) = activity_queue
.mark_completed(&activity, worker_label.as_str())
.await
{
error!(%worker_id, activity_id = %activity_id, error = %e, "Failed to mark activity as completed");
}
info!(%worker_id, activity_id = %activity_id, activity_type = ?activity_type, "Activity completed successfully");
let aq = activity_queue.clone();
let result_to_store = ActivityResult {
data: value,
state: ResultState::Ok,
};
tokio::spawn(async move {
if let Err(e) = aq.store_result(activity_id, result_to_store).await {
error!(activity_id = %activity_id, error = %e, "Failed to store activity result");
}
});
}
Ok(Err(ActivityError::Retry(reason))) => {
metrics.inc_counter("activity_retry", 1);
warn!(%worker_id, activity_id = %activity_id, activity_type = ?activity_type, reason = %reason, "Activity requesting retry");
match activity_queue
.mark_failed(activity, reason.clone(), true, worker_label.as_str())
.await
{
Ok(true) => {
let dead_letter_context = ActivityContext {
activity_id,
activity_type: activity_type.clone(),
retry_count: 0, metadata: Default::default(),
cancel_token: cancel_token.child_token(),
activity_executor: Arc::new(WorkerEngineWrapper::new(
activity_queue_for_context.clone(),
)),
};
handler
.on_dead_letter(
payload_for_dead_letter.clone(),
dead_letter_context,
reason,
)
.await;
}
Ok(false) => {
}
Err(e) => {
error!(%worker_id, activity_id = %activity_id, error = %e, "Failed to mark activity for retry");
}
}
}
Ok(Err(ActivityError::NonRetry(reason))) => {
metrics.inc_counter("activity_failed_non_retry", 1);
error!(%worker_id, activity_id = %activity_id, activity_type = ?activity_type, reason = %reason, "Activity failed");
if let Err(e) = activity_queue
.mark_failed(activity, reason.clone(), false, worker_label.as_str())
.await
{
error!(%worker_id, activity_id = %activity_id, error = %e, "Failed to mark activity as failed");
}
let aq = activity_queue.clone();
tokio::spawn(async move {
let activity_result = ActivityResult {
data: Some(json!({
"error": reason,
"type": "non_retryable",
"failed_at": Utc::now().to_rfc3339()
})),
state: ResultState::Err,
};
if let Err(e) = aq.store_result(activity_id, activity_result).await {
error!(activity_id = %activity_id, error = %e, "Failed to store activity result");
}
});
}
Err(_elapsed) => {
metrics.inc_counter("activity_timeout", 1);
let error_msg = "Activity execution timed out".to_string();
error!(%worker_id, activity_id = %activity_id, activity_type = ?activity_type, timeout = ?activity_timeout, "Activity timed out");
match activity_queue
.mark_failed(activity, error_msg.clone(), true, worker_label.as_str())
.await
{
Ok(true) => {
let dead_letter_context = ActivityContext {
activity_id,
activity_type: activity_type.clone(),
retry_count: 0,
metadata: Default::default(),
cancel_token: cancel_token.child_token(),
activity_executor: Arc::new(WorkerEngineWrapper::new(
activity_queue_for_context.clone(),
)),
};
handler
.on_dead_letter(
payload_for_dead_letter.clone(),
dead_letter_context,
error_msg,
)
.await;
}
Ok(false) => {
}
Err(e) => {
error!(%worker_id, activity_id = %activity_id, error = %e, "Failed to mark activity as failed");
}
}
}
}
}
debug!(%worker_id, "Worker loop stopped");
Ok(())
})
}
async fn start_scheduled_activities_processor(
&self,
) -> tokio::task::JoinHandle<Result<(), WorkerError>> {
let activity_queue = self.activity_queue.clone();
let running = self.running.clone();
let mut shutdown_rx = self.shutdown_tx.subscribe();
let poll_interval = {
let secs = self
.config
.schedule_poll_interval_seconds
.unwrap_or(5)
.max(1);
Duration::from_secs(secs)
};
tokio::spawn(async move {
debug!("Starting scheduled activities processor");
while *running.read().await {
if *shutdown_rx.borrow() {
break;
}
if let Err(e) = activity_queue.process_scheduled_activities().await {
error!(error = %e, "Failed to process scheduled activities");
}
tokio::select! {
_ = tokio::time::sleep(poll_interval) => {},
_ = shutdown_rx.changed() => break,
}
}
debug!("Scheduled activities processor stopped");
Ok(())
})
}
async fn start_reaper_processor(&self) -> tokio::task::JoinHandle<Result<(), WorkerError>> {
let activity_queue = self.activity_queue.clone();
let running = self.running.clone();
let mut shutdown_rx = self.shutdown_tx.subscribe();
let interval = Duration::from_secs(self.config.reaper_interval_seconds.unwrap_or(5).max(1));
let batch_size = self.config.reaper_batch_size.unwrap_or(100);
tokio::spawn(async move {
debug!("Starting reaper processor");
while *running.read().await {
if *shutdown_rx.borrow() {
break;
}
if let Err(e) = activity_queue.requeue_expired(batch_size).await {
error!(error = %e, "Reaper failed to requeue expired items");
}
tokio::select! {
_ = tokio::time::sleep(interval) => {},
_ = shutdown_rx.changed() => break,
}
}
debug!("Reaper processor stopped");
Ok(())
})
}
async fn wait_for_shutdown(&self) {
let ctrl_c = async {
tokio::signal::ctrl_c()
.await
.expect("failed to install Ctrl+C handler");
};
#[cfg(unix)]
let terminate = async {
use tokio::signal::unix::{signal, SignalKind};
let mut sigterm =
signal(SignalKind::terminate()).expect("failed to install SIGTERM handler");
sigterm.recv().await;
};
#[cfg(not(unix))]
let terminate = std::future::pending::<()>();
tokio::select! {
_ = ctrl_c => { info!("Received Ctrl+C signal"); },
_ = terminate => { info!("Received SIGTERM signal"); },
}
self.stop().await;
}
}
impl WorkerEngine {
pub fn register_activity(&mut self, activity_type: String, activity: Arc<dyn ActivityHandler>) {
self.activity_handlers.insert(activity_type, activity);
}
pub fn get_activity_executor(&self) -> Arc<dyn ActivityExecutor> {
Arc::new(WorkerEngineWrapper::new(self.activity_queue.clone()))
}
}
pub struct WorkerEngineBuilder {
queue_name: Option<String>,
max_workers: Option<usize>,
schedule_poll_interval: Option<Duration>,
metrics: Option<Arc<dyn MetricsSink>>,
backend: Option<Arc<dyn Storage>>,
activity_types: Option<Vec<String>>,
}
impl WorkerEngineBuilder {
pub fn new() -> Self {
Self {
queue_name: None,
max_workers: None,
schedule_poll_interval: None,
metrics: None,
backend: None,
activity_types: None,
}
}
pub fn queue_name(mut self, name: &str) -> Self {
self.queue_name = Some(name.to_string());
self
}
pub fn max_workers(mut self, max: usize) -> Self {
self.max_workers = Some(max);
self
}
pub fn schedule_poll_interval(mut self, interval: Duration) -> Self {
self.schedule_poll_interval = Some(interval);
self
}
pub fn activity_types(mut self, types: &[&str]) -> Self {
self.activity_types = Some(types.iter().map(|s| s.to_string()).collect());
self
}
pub fn metrics(mut self, sink: Arc<dyn MetricsSink>) -> Self {
self.metrics = Some(sink);
self
}
pub fn backend(mut self, backend: Arc<dyn Storage>) -> Self {
self.backend = Some(backend);
self
}
pub async fn build(self) -> Result<WorkerEngine, WorkerError> {
let max_concurrent_activities = self.max_workers.unwrap_or(10);
let schedule_poll_interval_seconds = self.schedule_poll_interval.map_or(5, |d| d.as_secs());
if let Some(backend) = self.backend {
let queue_name = self.queue_name.unwrap_or_else(|| "default".to_string());
let config = WorkerConfig {
queue_name,
max_concurrent_activities,
schedule_poll_interval_seconds: Some(schedule_poll_interval_seconds),
lease_ms: Some(60_000),
reaper_interval_seconds: Some(5),
reaper_batch_size: Some(100),
activity_types: self.activity_types.clone(),
};
let mut worker_engine = WorkerEngine::new_with_backend(backend, config);
if let Some(metrics) = self.metrics {
worker_engine.with_metrics(metrics);
}
return Ok(worker_engine);
}
Err(WorkerError::Configuration(
"No backend configured. Call .backend(Arc::new(your_backend)) before .build(). \
Use PostgresBackend::new(...) for PostgreSQL, or the runner_q_redis crate for Redis.".to_string()
))
}
}
impl Default for WorkerEngineBuilder {
fn default() -> Self {
Self::new()
}
}
pub struct ActivityBuilder<'a> {
engine: &'a WorkerEngineWrapper,
activity_type: String,
payload: Option<serde_json::Value>,
priority: Option<ActivityPriority>,
max_retries: Option<u32>,
timeout: Option<Duration>,
delay: Option<Duration>,
idempotency_key: Option<(String, OnDuplicate)>,
}
impl<'a> ActivityBuilder<'a> {
pub fn new(engine: &'a WorkerEngineWrapper, activity_type: String) -> Self {
Self {
engine,
activity_type,
payload: None,
priority: None,
max_retries: None,
timeout: None,
delay: None,
idempotency_key: None,
}
}
pub fn payload(mut self, payload: serde_json::Value) -> Self {
self.payload = Some(payload);
self
}
pub fn priority(mut self, priority: ActivityPriority) -> Self {
self.priority = Some(priority);
self
}
pub fn max_retries(mut self, retries: u32) -> Self {
self.max_retries = Some(retries);
self
}
pub fn timeout(mut self, timeout: Duration) -> Self {
self.timeout = Some(timeout);
self
}
pub fn delay(mut self, delay: Duration) -> Self {
self.delay = Some(delay);
self
}
pub fn idempotency_key(mut self, key: impl Into<String>, behavior: OnDuplicate) -> Self {
self.idempotency_key = Some((key.into(), behavior));
self
}
pub async fn execute(self) -> Result<ActivityFuture, WorkerError> {
let payload = self
.payload
.ok_or_else(|| WorkerError::QueueError("Activity payload is required".to_string()))?;
let option = if self.priority.is_some()
|| self.max_retries.is_some()
|| self.timeout.is_some()
|| self.delay.is_some()
|| self.idempotency_key.is_some()
{
let idempotency_key = self
.idempotency_key
.map(|(key, behavior)| (format!("{}-{}", key, self.activity_type), behavior));
Some(ActivityOption {
priority: self.priority,
max_retries: self.max_retries.unwrap_or(3),
timeout_seconds: self.timeout.map(|d| d.as_secs()).unwrap_or(300),
delay_seconds: self.delay.map(|d| d.as_secs()),
idempotency_key,
})
} else {
None
};
self.engine
.execute_activity(self.activity_type, payload, option)
.await
}
}
#[async_trait::async_trait]
pub trait ActivityExecutor: Send + Sync {
fn activity(&self, activity_type: &str) -> ActivityBuilder<'_>;
}
#[derive(Clone)]
pub struct WorkerEngineWrapper {
activity_queue: Arc<dyn ActivityQueueTrait>,
}
impl WorkerEngineWrapper {
pub(crate) fn new(activity_queue: Arc<dyn ActivityQueueTrait>) -> Self {
Self { activity_queue }
}
async fn execute_activity(
&self,
activity_type: String,
payload: serde_json::Value,
option: Option<ActivityOption>,
) -> Result<ActivityFuture, WorkerError> {
let activity = Activity::new(activity_type, payload, option);
let activity_id = activity.id;
if let Some(existing_id) = self
.activity_queue
.evaluate_idempotency_rule(&activity)
.await?
{
return Ok(ActivityFuture::new(
self.activity_queue.clone(),
existing_id,
));
}
match activity.scheduled_at {
None => self.activity_queue.enqueue(activity).await?,
Some(_) => self.activity_queue.schedule_activity(activity).await?,
}
Ok(ActivityFuture::new(
self.activity_queue.clone(),
activity_id,
))
}
}
#[async_trait::async_trait]
impl ActivityExecutor for WorkerEngineWrapper {
fn activity(&self, activity_type: &str) -> ActivityBuilder<'_> {
ActivityBuilder::new(self, activity_type.to_string())
}
}
struct BackendQueueAdapter {
backend: Arc<dyn Storage>,
activity_types: Option<Vec<String>>,
}
impl BackendQueueAdapter {
fn new(backend: Arc<dyn Storage>, activity_types: Option<Vec<String>>) -> Self {
Self {
backend,
activity_types,
}
}
fn activity_to_queued(activity: &Activity) -> QueuedActivity {
QueuedActivity {
id: activity.id,
activity_type: activity.activity_type.clone(),
payload: activity.payload.clone(),
priority: activity.priority.clone(),
max_retries: activity.max_retries,
retry_count: activity.retry_count,
timeout_seconds: activity.timeout_seconds,
retry_delay_seconds: activity.retry_delay_seconds,
scheduled_at: activity.scheduled_at,
metadata: activity.metadata.clone(),
idempotency_key: activity.idempotency_key.as_ref().map(|(k, b)| {
let behavior = match b {
OnDuplicate::AllowReuse => IdempotencyBehavior::AllowReuse,
OnDuplicate::ReturnExisting => IdempotencyBehavior::ReturnExisting,
OnDuplicate::AllowReuseOnFailure => IdempotencyBehavior::AllowReuseOnFailure,
OnDuplicate::NoReuse => IdempotencyBehavior::NoReuse,
};
(k.clone(), behavior)
}),
created_at: activity.created_at,
}
}
fn queued_to_activity(queued: &QueuedActivity) -> Activity {
Activity {
id: queued.id,
activity_type: queued.activity_type.clone(),
payload: queued.payload.clone(),
priority: queued.priority.clone(),
status: crate::ActivityStatus::Pending,
created_at: queued.created_at,
scheduled_at: queued.scheduled_at,
retry_count: queued.retry_count,
max_retries: queued.max_retries,
timeout_seconds: queued.timeout_seconds,
retry_delay_seconds: queued.retry_delay_seconds,
metadata: queued.metadata.clone(),
idempotency_key: queued.idempotency_key.as_ref().map(|(k, b)| {
let behavior = match b {
IdempotencyBehavior::AllowReuse => OnDuplicate::AllowReuse,
IdempotencyBehavior::ReturnExisting => OnDuplicate::ReturnExisting,
IdempotencyBehavior::AllowReuseOnFailure => OnDuplicate::AllowReuseOnFailure,
IdempotencyBehavior::NoReuse => OnDuplicate::NoReuse,
};
(k.clone(), behavior)
}),
}
}
}
#[async_trait::async_trait]
impl ActivityQueueTrait for BackendQueueAdapter {
async fn enqueue(&self, activity: Activity) -> Result<(), WorkerError> {
let queued = Self::activity_to_queued(&activity);
self.backend.enqueue(queued).await.map_err(Into::into)
}
async fn dequeue(
&self,
timeout: Duration,
worker_id: &str,
) -> Result<Option<Activity>, WorkerError> {
let types_ref = self.activity_types.as_deref();
match self.backend.dequeue(worker_id, timeout, types_ref).await? {
Some(activity) => Ok(Some(Self::queued_to_activity(&activity))),
None => Ok(None),
}
}
async fn mark_completed(
&self,
activity: &Activity,
worker_id: &str,
) -> Result<(), WorkerError> {
self.backend
.ack_success(activity.id, None, worker_id)
.await
.map_err(Into::into)
}
async fn mark_failed(
&self,
activity: Activity,
error_message: String,
retryable: bool,
worker_id: &str,
) -> Result<bool, WorkerError> {
let failure = if retryable {
FailureKind::Retryable {
reason: error_message,
}
} else {
FailureKind::NonRetryable {
reason: error_message,
}
};
self.backend
.ack_failure(activity.id, failure, worker_id)
.await
.map_err(Into::into)
}
async fn schedule_activity(&self, activity: Activity) -> Result<(), WorkerError> {
let mut queued = Self::activity_to_queued(&activity);
if queued.scheduled_at.is_none() {
queued.scheduled_at = Some(Utc::now());
}
self.backend.enqueue(queued).await.map_err(Into::into)
}
async fn process_scheduled_activities(&self) -> Result<Vec<Activity>, WorkerError> {
let _count = self.backend.process_scheduled().await?;
Ok(vec![])
}
async fn requeue_expired(&self, max_to_process: usize) -> Result<u64, WorkerError> {
self.backend
.requeue_expired(max_to_process)
.await
.map_err(Into::into)
}
async fn evaluate_idempotency_rule(
&self,
activity: &Activity,
) -> Result<Option<uuid::Uuid>, WorkerError> {
let queued = Self::activity_to_queued(activity);
self.backend
.check_idempotency(&queued)
.await
.map_err(Into::into)
}
async fn extend_lease(
&self,
activity_id: uuid::Uuid,
extend_by: Duration,
) -> Result<bool, WorkerError> {
self.backend
.extend_lease(activity_id, extend_by)
.await
.map_err(Into::into)
}
async fn store_result(
&self,
activity_id: uuid::Uuid,
result: ActivityResult,
) -> Result<(), WorkerError> {
let backend_result = crate::storage::ActivityResult {
data: result.data,
state: match result.state {
ResultState::Ok => crate::storage::ResultState::Ok,
ResultState::Err => crate::storage::ResultState::Err,
},
};
self.backend
.store_result(activity_id, backend_result)
.await
.map_err(Into::into)
}
async fn get_result(
&self,
activity_id: uuid::Uuid,
) -> Result<Option<ActivityResult>, WorkerError> {
match self.backend.get_result(activity_id).await? {
Some(backend_result) => Ok(Some(ActivityResult {
data: backend_result.data,
state: match backend_result.state {
crate::storage::ResultState::Ok => ResultState::Ok,
crate::storage::ResultState::Err => ResultState::Err,
},
})),
None => Ok(None),
}
}
fn schedules_natively(&self) -> bool {
self.backend.schedules_natively()
}
}