diff --git a/crates/taskito-core/src/error.rs b/crates/taskito-core/src/error.rs index 14f0bcab..41c7b43b 100644 --- a/crates/taskito-core/src/error.rs +++ b/crates/taskito-core/src/error.rs @@ -8,6 +8,13 @@ pub enum QueueError { #[error("connection pool error: {0}")] Pool(#[from] diesel::r2d2::PoolError), + #[cfg(feature = "redis")] + #[error("redis error: {0}")] + Redis(#[from] redis::RedisError), + + #[error("json error: {0}")] + Json(#[from] serde_json::Error), + #[error("job not found: {0}")] JobNotFound(String), diff --git a/crates/taskito-core/src/storage/redis_backend/archival.rs b/crates/taskito-core/src/storage/redis_backend/archival.rs index 91639833..8f8e5ff9 100644 --- a/crates/taskito-core/src/storage/redis_backend/archival.rs +++ b/crates/taskito-core/src/storage/redis_backend/archival.rs @@ -1,7 +1,7 @@ use redis::Commands; use super::{map_err, strip_list_blobs, RedisStorage}; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::job::{Job, JobStatus}; impl RedisStorage { @@ -73,8 +73,7 @@ impl RedisStorage { let archived_key = self.key(&["archived", id]); let data: Option = conn.get(&archived_key).map_err(map_err)?; if let Some(d) = data { - let mut job: Job = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let mut job: Job = serde_json::from_str(&d)?; strip_list_blobs(&mut job); jobs.push(job); } diff --git a/crates/taskito-core/src/storage/redis_backend/circuit_breakers.rs b/crates/taskito-core/src/storage/redis_backend/circuit_breakers.rs index d084ce50..1479737b 100644 --- a/crates/taskito-core/src/storage/redis_backend/circuit_breakers.rs +++ b/crates/taskito-core/src/storage/redis_backend/circuit_breakers.rs @@ -1,7 +1,7 @@ use redis::Commands; use super::{map_err, RedisStorage}; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::storage::records::CircuitBreakerState; impl RedisStorage { @@ -12,8 +12,7 @@ impl RedisStorage { let data: Option = conn.get(&cb_key).map_err(map_err)?; match data { Some(d) => { - let row: CircuitBreakerState = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let row: CircuitBreakerState = serde_json::from_str(&d)?; Ok(Some(row)) } None => Ok(None), @@ -25,7 +24,7 @@ impl RedisStorage { let cb_key = self.key(&["cb", &row.task_name]); let cb_all = self.key(&["cb", "all"]); - let json = serde_json::to_string(row).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(row)?; let pipe = &mut redis::pipe(); pipe.set(&cb_key, &json); @@ -45,8 +44,7 @@ impl RedisStorage { let cb_key = self.key(&["cb", &name]); let data: Option = conn.get(&cb_key).map_err(map_err)?; if let Some(d) = data { - let row: CircuitBreakerState = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let row: CircuitBreakerState = serde_json::from_str(&d)?; rows.push(row); } } diff --git a/crates/taskito-core/src/storage/redis_backend/dead_letter.rs b/crates/taskito-core/src/storage/redis_backend/dead_letter.rs index ac4953e7..ddef9096 100644 --- a/crates/taskito-core/src/storage/redis_backend/dead_letter.rs +++ b/crates/taskito-core/src/storage/redis_backend/dead_letter.rs @@ -88,7 +88,7 @@ impl RedisStorage { dlq_retry_count, }; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let dlq_key = self.key(&["dlq", &dlq_id]); let dlq_all = self.key(&["dlq", "all"]); @@ -98,8 +98,7 @@ impl RedisStorage { dead_job.status = JobStatus::Dead; dead_job.error = Some(error.to_string()); dead_job.completed_at = Some(now); - let dead_json = - serde_json::to_string(&dead_job).map_err(|e| QueueError::Other(e.to_string()))?; + let dead_json = serde_json::to_string(&dead_job)?; // Commit the DLQ entry and the live→archive move together, but only if the // job is still live in its expected state. A racing complete/fail or the @@ -148,8 +147,7 @@ impl RedisStorage { let dlq_key = self.key(&["dlq", &id]); let data: Option = conn.get(&dlq_key).map_err(map_err)?; if let Some(d) = data { - let entry: DeadJobEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: DeadJobEntry = serde_json::from_str(&d)?; let mut dead = DeadJob::from(entry); strip_dead_blob(&mut dead); results.push(dead); @@ -172,8 +170,7 @@ impl RedisStorage { let dlq_key = self.key(&["dlq", id]); let data: Option = conn.get(&dlq_key).map_err(map_err)?; if let Some(d) = data { - let entry: DeadJobEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: DeadJobEntry = serde_json::from_str(&d)?; let mut dead = DeadJob::from(entry); strip_dead_blob(&mut dead); results.push(dead); @@ -209,8 +206,7 @@ impl RedisStorage { let dlq_key = self.key(&["dlq", &id]); let data: Option = conn.get(&dlq_key).map_err(map_err)?; if let Some(d) = data { - let entry: DeadJobEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: DeadJobEntry = serde_json::from_str(&d)?; if entry.task_name == task_name { let mut dead = DeadJob::from(entry); strip_dead_blob(&mut dead); @@ -244,8 +240,7 @@ impl RedisStorage { if let Some(d) = data { // Propagate (don't skip) on a corrupt entry: silently ignoring it // would leave a task-owned row behind and under-report the count. - let entry: DeadJobEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: DeadJobEntry = serde_json::from_str(&d)?; if entry.task_name == task_name { to_delete.push((id, entry.notes, entry.original_job_id)); } @@ -277,8 +272,7 @@ impl RedisStorage { v.ok_or_else(|| QueueError::JobNotFound(dead_id.to_string())) })?; - let entry: DeadJobEntry = - serde_json::from_str(&data).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: DeadJobEntry = serde_json::from_str(&data)?; // Attribution + original job id for the sub:dead removal below, captured // before `entry`'s fields are moved into `new_job`. The re-enqueue below @@ -391,8 +385,7 @@ impl RedisStorage { let Some(d) = data else { return Ok(false); }; - let entry: DeadJobEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: DeadJobEntry = serde_json::from_str(&d)?; let pipe = &mut redis::pipe(); pipe.del(&dlq_key); diff --git a/crates/taskito-core/src/storage/redis_backend/jobs/dequeue.rs b/crates/taskito-core/src/storage/redis_backend/jobs/dequeue.rs index 6c64a0a8..a95174b4 100644 --- a/crates/taskito-core/src/storage/redis_backend/jobs/dequeue.rs +++ b/crates/taskito-core/src/storage/redis_backend/jobs/dequeue.rs @@ -2,7 +2,7 @@ use redis::Commands; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::job::{Job, JobStatus}; use crate::storage::redis_backend::{map_err, RedisStorage}; @@ -36,7 +36,7 @@ impl RedisStorage { job: &Job, queue_key: &str, ) -> Result { - let job_json = serde_json::to_string(job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(job)?; let job_key = self.key(&["job", &job.id]); let pending_status = self.key(&["jobs", "status", &(JobStatus::Pending as i32).to_string()]); @@ -88,8 +88,7 @@ impl RedisStorage { } }; - let mut job: Job = - serde_json::from_str(&data).map_err(|e| QueueError::Other(e.to_string()))?; + let mut job: Job = serde_json::from_str(&data)?; // Must be pending and scheduled_at <= now if job.status != JobStatus::Pending || job.scheduled_at > now { @@ -232,8 +231,7 @@ impl RedisStorage { } }; - let mut job: Job = - serde_json::from_str(&data).map_err(|e| QueueError::Other(e.to_string()))?; + let mut job: Job = serde_json::from_str(&data)?; // Must be pending and scheduled_at <= now if job.status != JobStatus::Pending || job.scheduled_at > now { diff --git a/crates/taskito-core/src/storage/redis_backend/jobs/enqueue.rs b/crates/taskito-core/src/storage/redis_backend/jobs/enqueue.rs index f1595335..1aeb7d49 100644 --- a/crates/taskito-core/src/storage/redis_backend/jobs/enqueue.rs +++ b/crates/taskito-core/src/storage/redis_backend/jobs/enqueue.rs @@ -30,9 +30,7 @@ impl RedisStorage { let dep_key = self.key(&["job", dep_id]); let data: Option = conn.get(&dep_key).map_err(map_err)?; let dep_job: Job = match data { - Some(d) => { - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))? - } + Some(d) => serde_json::from_str(&d)?, None => match self.load_archived_job(conn, dep_id)? { Some(archived) if archived.status == JobStatus::Complete => continue, _ => return Err(QueueError::DependencyNotFound(DEP_MISSING.to_string())), @@ -52,7 +50,7 @@ impl RedisStorage { self.validate_dep_ids(&mut conn, &depends_on, None)?; - let job_json = serde_json::to_string(&job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(&job)?; let job_key = self.key(&["job", &job.id]); let status_key = self.key(&["jobs", "status", &(job.status as i32).to_string()]); let queue_key = self.key(&["queue", &job.queue, "pending"]); @@ -102,8 +100,7 @@ impl RedisStorage { let pipe = &mut redis::pipe(); for (i, job) in jobs.iter().enumerate() { - let job_json = - serde_json::to_string(job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(job)?; let job_key = self.key(&["job", &job.id]); let status_key = self.key(&["jobs", "status", &(job.status as i32).to_string()]); let queue_key = self.key(&["queue", &job.queue, "pending"]); @@ -195,16 +192,14 @@ impl RedisStorage { .map_err(map_err)?; if let Some(job_data) = result { - let job: Job = serde_json::from_str(&job_data) - .map_err(|e| QueueError::Other(e.to_string()))?; + let job: Job = serde_json::from_str(&job_data)?; return Ok(job); } // No active duplicate — enqueue normally let depends_on = new_job.depends_on.clone(); let job = new_job.into_job(); - let job_json = - serde_json::to_string(&job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(&job)?; self.validate_dep_ids(&mut conn, &depends_on, None)?; @@ -331,8 +326,7 @@ impl RedisStorage { if let Some(existing_data) = result { // Lost the race — another caller created a job first - let existing_job: Job = serde_json::from_str(&existing_data) - .map_err(|e| QueueError::Other(e.to_string()))?; + let existing_job: Job = serde_json::from_str(&existing_data)?; return Ok(existing_job); } diff --git a/crates/taskito-core/src/storage/redis_backend/jobs/errors.rs b/crates/taskito-core/src/storage/redis_backend/jobs/errors.rs index 88d19f18..708ee8c5 100644 --- a/crates/taskito-core/src/storage/redis_backend/jobs/errors.rs +++ b/crates/taskito-core/src/storage/redis_backend/jobs/errors.rs @@ -2,7 +2,7 @@ use redis::Commands; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::job::now_millis; use crate::storage::records::JobError; use crate::storage::redis_backend::{map_err, RedisStorage}; @@ -20,7 +20,7 @@ impl RedisStorage { error: error.to_string(), failed_at: now, }; - let json = serde_json::to_string(&row).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&row)?; let errors_key = self.key(&["job_errors", job_id]); conn.rpush::<_, _, ()>(&errors_key, &json) @@ -36,8 +36,7 @@ impl RedisStorage { let mut rows = Vec::new(); for entry in entries { - let row: JobError = - serde_json::from_str(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let row: JobError = serde_json::from_str(&entry)?; rows.push(row); } rows.sort_by_key(|r| r.attempt); diff --git a/crates/taskito-core/src/storage/redis_backend/jobs/helpers.rs b/crates/taskito-core/src/storage/redis_backend/jobs/helpers.rs index 5074b5e2..1efc5f90 100644 --- a/crates/taskito-core/src/storage/redis_backend/jobs/helpers.rs +++ b/crates/taskito-core/src/storage/redis_backend/jobs/helpers.rs @@ -29,8 +29,7 @@ impl RedisStorage { let data: Option = conn.get(&job_key).map_err(map_err)?; match data { Some(d) => { - let job: Job = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let job: Job = serde_json::from_str(&d)?; Ok(Some(job)) } None => Ok(None), @@ -46,8 +45,7 @@ impl RedisStorage { let data: Option = conn.get(&archived_key).map_err(map_err)?; match data { Some(d) => { - let job: Job = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let job: Job = serde_json::from_str(&d)?; Ok(Some(job)) } None => Ok(None), @@ -71,7 +69,7 @@ impl RedisStorage { job: &Job, old_status: JobStatus, ) -> Result<()> { - let job_json = serde_json::to_string(job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(job)?; let job_key = self.key(&["job", &job.id]); let old_status_key = self.key(&["jobs", "status", &(old_status as i32).to_string()]); let new_status_key = self.key(&["jobs", "status", &(job.status as i32).to_string()]); @@ -104,7 +102,7 @@ impl RedisStorage { job: &Job, old_status: JobStatus, ) -> Result<()> { - let job_json = serde_json::to_string(job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(job)?; let job_key = self.key(&["job", &job.id]); let old_status_key = self.key(&["jobs", "status", &(old_status as i32).to_string()]); let pending_status_key = @@ -217,7 +215,7 @@ impl RedisStorage { job: &Job, old_status: JobStatus, ) -> Result<()> { - let job_json = serde_json::to_string(job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(job)?; let pipe = &mut redis::pipe(); pipe.atomic(); self.push_archive_ops(pipe, job, old_status, &job_json); diff --git a/crates/taskito-core/src/storage/redis_backend/jobs/state.rs b/crates/taskito-core/src/storage/redis_backend/jobs/state.rs index 2c7627a7..237ccc11 100644 --- a/crates/taskito-core/src/storage/redis_backend/jobs/state.rs +++ b/crates/taskito-core/src/storage/redis_backend/jobs/state.rs @@ -161,7 +161,7 @@ impl RedisStorage { job.error = None; job.cancel_requested = false; - let job_json = serde_json::to_string(&job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(&job)?; let job_key = self.key(&["job", id]); let running_key = self.key(&["jobs", "status", &(JobStatus::Running as i32).to_string()]); let pending_key = self.key(&["jobs", "status", &(JobStatus::Pending as i32).to_string()]); @@ -233,7 +233,7 @@ impl RedisStorage { } job.cancel_requested = true; - let job_json = serde_json::to_string(&job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(&job)?; let job_key = self.key(&["job", id]); let running_key = self.key(&["jobs", "status", &(JobStatus::Running as i32).to_string()]); let cancel_set = self.key(&["jobs", "cancel_requested"]); @@ -345,7 +345,7 @@ impl RedisStorage { let mut conn = self.conn()?; let mut job = self.get_job_required(id)?; job.progress = Some(progress); - let job_json = serde_json::to_string(&job).map_err(|e| QueueError::Other(e.to_string()))?; + let job_json = serde_json::to_string(&job)?; let job_key = self.key(&["job", id]); // Guarded write: only update if `job:` still exists. If the job was diff --git a/crates/taskito-core/src/storage/redis_backend/logs.rs b/crates/taskito-core/src/storage/redis_backend/logs.rs index 67811091..e7801511 100644 --- a/crates/taskito-core/src/storage/redis_backend/logs.rs +++ b/crates/taskito-core/src/storage/redis_backend/logs.rs @@ -2,7 +2,7 @@ use redis::Commands; use serde::{Deserialize, Serialize}; use super::{map_err, RedisStorage, SCAN_BATCH}; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::job::now_millis; use crate::storage::records::TaskLogEntry; @@ -54,7 +54,7 @@ impl RedisStorage { logged_at: now, }; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let log_key = self.key(&["log", &id]); let by_job_key = self.key(&["logs", "by_job", job_id]); @@ -82,8 +82,7 @@ impl RedisStorage { let log_key = self.key(&["log", &id]); let data: Option = conn.get(&log_key).map_err(map_err)?; if let Some(d) = data { - let entry: LogEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: LogEntry = serde_json::from_str(&d)?; rows.push(TaskLogEntry::from(entry)); } } @@ -116,8 +115,7 @@ impl RedisStorage { let log_key = self.key(&["log", &id]); let data: Option = conn.get(&log_key).map_err(map_err)?; if let Some(d) = data { - let entry: LogEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: LogEntry = serde_json::from_str(&d)?; rows.push(TaskLogEntry::from(entry)); } } @@ -154,8 +152,7 @@ impl RedisStorage { let log_key = self.key(&["log", &id]); let data: Option = conn.get(&log_key).map_err(map_err)?; if let Some(d) = data { - let entry: LogEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: LogEntry = serde_json::from_str(&d)?; rows.push(TaskLogEntry::from(entry)); } } @@ -187,8 +184,7 @@ impl RedisStorage { let entries: Vec> = pipe.query(&mut conn).map_err(map_err)?; for data in entries.into_iter().flatten() { - let entry: LogEntry = - serde_json::from_str(&data).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: LogEntry = serde_json::from_str(&data)?; if task_name.is_some_and(|n| entry.task_name != n) { continue; diff --git a/crates/taskito-core/src/storage/redis_backend/metrics.rs b/crates/taskito-core/src/storage/redis_backend/metrics.rs index bb1527cb..108659b0 100644 --- a/crates/taskito-core/src/storage/redis_backend/metrics.rs +++ b/crates/taskito-core/src/storage/redis_backend/metrics.rs @@ -2,7 +2,7 @@ use redis::Commands; use serde::{Deserialize, Serialize}; use super::{map_err, RedisStorage}; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::job::now_millis; use crate::storage::records::{ReplayEntry, TaskMetric}; @@ -81,7 +81,7 @@ impl RedisStorage { recorded_at: now, }; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let metric_key = self.key(&["metric", &id]); let all_key = self.key(&["metrics", "all"]); @@ -114,8 +114,7 @@ impl RedisStorage { let metric_key = self.key(&["metric", &id]); let data: Option = conn.get(&metric_key).map_err(map_err)?; if let Some(d) = data { - let entry: MetricEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: MetricEntry = serde_json::from_str(&d)?; rows.push(TaskMetric::from(entry)); } } @@ -180,7 +179,7 @@ impl RedisStorage { replay_error: replay_error.map(|s| s.to_string()), }; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let replay_key = self.key(&["replay", &id]); let by_original = self.key(&["replay", "by_original", original_job_id]); @@ -203,8 +202,7 @@ impl RedisStorage { let replay_key = self.key(&["replay", &id]); let data: Option = conn.get(&replay_key).map_err(map_err)?; if let Some(d) = data { - let entry: ReplayHistoryEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: ReplayHistoryEntry = serde_json::from_str(&d)?; rows.push(ReplayEntry::from(entry)); } } diff --git a/crates/taskito-core/src/storage/redis_backend/mod.rs b/crates/taskito-core/src/storage/redis_backend/mod.rs index 3f69f464..341ec6d8 100644 --- a/crates/taskito-core/src/storage/redis_backend/mod.rs +++ b/crates/taskito-core/src/storage/redis_backend/mod.rs @@ -118,7 +118,7 @@ impl crate::storage::notify::StorageNotifier for RedisStorage { } fn map_err(e: redis::RedisError) -> QueueError { - QueueError::Other(e.to_string()) + QueueError::Redis(e) } /// Batch size for the bounded history scans (SSCAN/ZSCAN COUNT hint and the diff --git a/crates/taskito-core/src/storage/redis_backend/periodic.rs b/crates/taskito-core/src/storage/redis_backend/periodic.rs index 1aa36836..5f54f602 100644 --- a/crates/taskito-core/src/storage/redis_backend/periodic.rs +++ b/crates/taskito-core/src/storage/redis_backend/periodic.rs @@ -2,7 +2,7 @@ use redis::Commands; use serde::{Deserialize, Serialize}; use super::{map_err, RedisStorage}; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::storage::records::{NewPeriodicTask, PeriodicTask}; #[derive(Serialize, Deserialize)] @@ -53,7 +53,7 @@ impl RedisStorage { timezone: task.timezone.clone(), }; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let pkey = self.key(&["periodic", &task.name]); let due_key = self.key(&["periodic", "due"]); @@ -81,8 +81,7 @@ impl RedisStorage { let pkey = self.key(&["periodic", &name]); let data: Option = conn.get(&pkey).map_err(map_err)?; if let Some(d) = data { - let entry: PeriodicEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: PeriodicEntry = serde_json::from_str(&d)?; if entry.enabled { rows.push(PeriodicTask::from(entry)); } @@ -98,13 +97,11 @@ impl RedisStorage { let data: Option = conn.get(&pkey).map_err(map_err)?; if let Some(d) = data { - let mut entry: PeriodicEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let mut entry: PeriodicEntry = serde_json::from_str(&d)?; entry.last_run = Some(last_run); entry.next_run = next_run; - let json = - serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let due_key = self.key(&["periodic", "due"]); let pipe = &mut redis::pipe(); @@ -145,8 +142,7 @@ impl RedisStorage { } let data: Option = conn.get(&key).map_err(map_err)?; if let Some(d) = data { - let entry: PeriodicEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: PeriodicEntry = serde_json::from_str(&d)?; rows.push(PeriodicTask::from(entry)); } } @@ -185,11 +181,10 @@ impl RedisStorage { return Ok(false); }; - let mut entry: PeriodicEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let mut entry: PeriodicEntry = serde_json::from_str(&d)?; entry.enabled = enabled; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let due_key = self.key(&["periodic", "due"]); let pipe = &mut redis::pipe(); diff --git a/crates/taskito-core/src/storage/redis_backend/pubsub.rs b/crates/taskito-core/src/storage/redis_backend/pubsub.rs index 66b05135..7d4e86f8 100644 --- a/crates/taskito-core/src/storage/redis_backend/pubsub.rs +++ b/crates/taskito-core/src/storage/redis_backend/pubsub.rs @@ -71,8 +71,7 @@ impl RedisStorage { let data: Option = conn.get(&blob_key).map_err(map_err)?; match data { Some(d) => { - let entry: SubEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: SubEntry = serde_json::from_str(&d)?; Ok(Some(entry)) } None => Ok(None), @@ -99,7 +98,7 @@ impl RedisStorage { blobs .into_iter() .flatten() - .map(|d| serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))) + .map(|d| serde_json::from_str(&d).map_err(QueueError::from)) .collect() } @@ -141,7 +140,7 @@ impl RedisStorage { entry.active = prior.active; entry.created_at = prior.created_at; } - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; let pipe = redis::pipe().atomic().to_owned(); let mut pipe = pipe; @@ -228,8 +227,7 @@ impl RedisStorage { let Some(d) = data else { return Ok(false); }; - let entry: SubEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let entry: SubEntry = serde_json::from_str(&d)?; let by_topic = self.key(&["subs", "by_topic", topic]); let all = self.key(&["subs", "all"]); @@ -263,11 +261,10 @@ impl RedisStorage { return Ok(false); }; - let mut entry: SubEntry = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let mut entry: SubEntry = serde_json::from_str(&d)?; entry.active = active; - let json = serde_json::to_string(&entry).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(&entry)?; conn.set::<_, _, ()>(&blob_key, &json).map_err(map_err)?; Ok(true) diff --git a/crates/taskito-core/src/storage/redis_backend/rate_limits.rs b/crates/taskito-core/src/storage/redis_backend/rate_limits.rs index 6bcbe549..8c0a984d 100644 --- a/crates/taskito-core/src/storage/redis_backend/rate_limits.rs +++ b/crates/taskito-core/src/storage/redis_backend/rate_limits.rs @@ -1,7 +1,7 @@ use redis::Commands; use super::{map_err, RedisStorage}; -use crate::error::{QueueError, Result}; +use crate::error::Result; use crate::job::now_millis; use crate::storage::records::RateLimitState; @@ -13,8 +13,7 @@ impl RedisStorage { let data: Option = conn.get(&rkey).map_err(map_err)?; match data { Some(d) => { - let row: RateLimitState = - serde_json::from_str(&d).map_err(|e| QueueError::Other(e.to_string()))?; + let row: RateLimitState = serde_json::from_str(&d)?; Ok(Some(row)) } None => Ok(None), @@ -24,7 +23,7 @@ impl RedisStorage { pub fn upsert_rate_limit(&self, row: &RateLimitState) -> Result<()> { let mut conn = self.conn()?; let rkey = self.key(&["rate_limit", &row.key]); - let json = serde_json::to_string(row).map_err(|e| QueueError::Other(e.to_string()))?; + let json = serde_json::to_string(row)?; conn.set::<_, _, ()>(&rkey, &json).map_err(map_err)?; Ok(()) }