pub struct RedisQueueLayer { /* private fields */ }Expand description
Redis ZSET layer for queued runs ({prefix}:ready:{pool} or hash-tagged for cluster).
Used only as the claim hot path inside crate::PostgresRedisSchedulerStore — not a
standalone SchedulerStore. Connect with Self::connect:
| Knob | Effect |
|---|---|
url / CHRONON_REDIS_URL | Standalone Redis |
CHRONON_REDIS_CLUSTER_URLS | Comma-separated cluster nodes |
CHRONON_REDIS_HASH_TAGS=1 | Hash-tagged pool keys for slot affinity |
key_prefix (default chronon) | Key namespace |
§Examples
use chronon_backend_redis::RedisQueueLayer;
let redis = RedisQueueLayer::connect("redis://127.0.0.1:6379", Some("myapp")).await?;Implementations§
Source§impl RedisQueueLayer
impl RedisQueueLayer
Sourcepub async fn connect(
url: &str,
key_prefix: Option<&str>,
) -> Result<RedisQueueLayer, ChrononError>
pub async fn connect( url: &str, key_prefix: Option<&str>, ) -> Result<RedisQueueLayer, ChrononError>
Connect to Redis at url with optional key prefix (default chronon).
When CHRONON_REDIS_CLUSTER_URLS is set (comma-separated), uses Redis Cluster.
When CHRONON_REDIS_HASH_TAGS=1, pool keys use hash tags for cluster slot affinity.
§Errors
Returns a storage error when the connection cannot be established.
Sourcepub fn test_url() -> String
pub fn test_url() -> String
Redis URL for tests (CHRONON_TEST_REDIS_URL, then CHRONON_REDIS_URL, or local default).
Sourcepub async fn enqueue_run(
&self,
pool_id: &str,
run_id: &str,
scheduled_for: DateTime<Utc>,
) -> Result<(), ChrononError>
pub async fn enqueue_run( &self, pool_id: &str, run_id: &str, scheduled_for: DateTime<Utc>, ) -> Result<(), ChrononError>
Enqueue a run id ordered by scheduled_for (ZADD score = epoch millis).
§Errors
Returns a storage error when Redis commands fail.
Sourcepub async fn claim_next_run_id(
&self,
pool_id: &str,
now: DateTime<Utc>,
) -> Result<Option<String>, ChrononError>
pub async fn claim_next_run_id( &self, pool_id: &str, now: DateTime<Utc>, ) -> Result<Option<String>, ChrononError>
Atomically pop the earliest due run id (score <= now) from the pool queue.
§Errors
Returns a storage error when Redis commands fail.
Sourcepub async fn claim_next_run_ids(
&self,
pool_id: &str,
count: usize,
now: DateTime<Utc>,
) -> Result<Vec<String>, ChrononError>
pub async fn claim_next_run_ids( &self, pool_id: &str, count: usize, now: DateTime<Utc>, ) -> Result<Vec<String>, ChrononError>
Pop up to count earliest due run ids (scheduled_for score ≤ now).
§Errors
Returns a storage error when Redis commands fail.
Sourcepub async fn remove_run(
&self,
pool_id: &str,
run_id: &str,
) -> Result<(), ChrononError>
pub async fn remove_run( &self, pool_id: &str, run_id: &str, ) -> Result<(), ChrononError>
Remove a run from the ready queue (e.g. after cancellation).
§Errors
Returns a storage error when Redis commands fail.
Sourcepub async fn flush_keys(&self) -> Result<(), ChrononError>
pub async fn flush_keys(&self) -> Result<(), ChrononError>
Delete all keys with this layer’s prefix (test isolation).
§Errors
Returns a storage error when Redis commands fail.
Trait Implementations§
Auto Trait Implementations§
impl !RefUnwindSafe for RedisQueueLayer
impl !UnwindSafe for RedisQueueLayer
impl Freeze for RedisQueueLayer
impl Send for RedisQueueLayer
impl Sync for RedisQueueLayer
impl Unpin for RedisQueueLayer
impl UnsafeUnpin for RedisQueueLayer
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
impl<A, B, T> HttpServerConnExec<A, B> for Twhere
B: Body,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more