Files
inbuxa-server/crates/store/src/backend/redis/lookup.rs
T
jcoffey-dev 6e50ba25a9
ci / fork-checks (pull_request) Successful in 47s
ci / build (pull_request) Successful in 4m58s
SQL pools time out; task locks are a renewed five-minute lease
A 3-node rehearsal (PostgreSQL + NATS + Garage) found two ways a crash
leaves work stuck:

Pool hangs. The PostgreSQL pool (deadpool) was built with no timeouts,
so a request waited for a free connection, and for one to be opened or
recycled, for as long as it took: forever when the server stopped
answering. MySQL's pool (mysql_async) has no wait timeout at all.

- PostgreSQL: wait 30 s (or the store's timeout if longer), create the
  store's timeout or 15 s (it bounds the whole handshake, where
  tokio-postgres's connect_timeout covers only the TCP connect), recycle
  10 s. The pool config is now always set, not only with
  poolMaxConnections.
- MySQL: every connection is taken through MysqlStore::conn(), which
  gives up after 30 s.
- Both: TCP keepalive after 60 s idle, so a server that vanished
  without closing the connection is noticed in minutes rather than the
  two-hour system default.

The DataStore schema has no pool timeout settings, so these are fixed
defaults; the store's own timeout bounds connecting on PostgreSQL.

Task locks. A task lock lasted an hour, so after a hard crash the dead
node's tasks waited up to an hour and five minutes. The lock is now a
five-minute lease: while this node runs a task, the task manager renews
its lock every third of the lifetime (InMemoryStore::renew_lock, a
compare-and-set on the store backends and SET XX EX on Redis, which
leaves a lock that already expired alone). A killed node's tasks run
elsewhere within about five minutes plus the claim recheck. A task this
node holds isn't handed to a worker again by the scan.

store::pool_timeout (new): a local listener that accepts connections
and never answers plays a hung server; a PostgreSQL store with a 2 s
timeout returns an error in about 4 s, and a MySQL store in 30 s.
Without the timeouts both wait for good. store::task_locks gains a
task held for 1.5 lock lifetimes: its lease is still held, and released
when the task ends.
2026-09-24 13:11:26 -07:00

334 lines
10 KiB
Rust

/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::{RedisPool, RedisStore, into_error};
use crate::{Deserialize, write::now};
use deadpool::managed::{Manager, Object, Pool};
use redis::{AsyncCommands, RedisError, RedisResult, RetryMethod, Script};
use std::sync::LazyLock;
static INCR_EXPIRE: LazyLock<Script> = LazyLock::new(|| {
Script::new(
"redis.call('INCRBY', KEYS[1], ARGV[1])
redis.call('EXPIRE', KEYS[1], ARGV[2])
return redis.call('GET', KEYS[1])",
)
});
impl RedisStore {
pub async fn key_set(&self, key: &[u8], value: &[u8], expires: Option<u64>) -> trc::Result<()> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| {
Self::key_set_(conn, key, value, expires).await
})
.await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| {
Self::key_set_(conn, key, value, expires).await
})
.await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| {
Self::key_set_(conn, key, value, expires).await
})
.await
}
}
}
pub async fn key_incr(&self, key: &[u8], value: i64, expires: Option<u64>) -> trc::Result<i64> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| {
Self::key_incr_(conn, key, value, expires).await
})
.await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| {
Self::key_incr_(conn, key, value, expires).await
})
.await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| {
Self::key_incr_(conn, key, value, expires).await
})
.await
}
}
}
pub async fn try_lock(&self, key: &[u8], expires: u64) -> trc::Result<bool> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| Self::try_lock_(conn, key, expires).await).await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| Self::try_lock_(conn, key, expires).await).await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| Self::try_lock_(conn, key, expires).await).await
}
}
}
// inbuxa: see InMemoryStore::renew_lock
pub async fn renew_lock(&self, key: &[u8], expires: u64) -> trc::Result<bool> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| {
Self::renew_lock_(conn, key, expires).await
})
.await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| {
Self::renew_lock_(conn, key, expires).await
})
.await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| {
Self::renew_lock_(conn, key, expires).await
})
.await
}
}
}
pub async fn key_delete(&self, key: &[u8]) -> trc::Result<()> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| Self::key_delete_(conn, key).await).await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| Self::key_delete_(conn, key).await).await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| Self::key_delete_(conn, key).await).await
}
}
}
pub async fn key_delete_prefix(&self, prefix: &[u8]) -> trc::Result<()> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| {
Self::key_delete_prefix_(conn, prefix).await
})
.await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| {
Self::key_delete_prefix_(conn, prefix).await
})
.await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| {
Self::key_delete_prefix_(conn, prefix).await
})
.await
}
}
}
pub async fn key_get<T: Deserialize + std::fmt::Debug + 'static>(
&self,
key: &[u8],
) -> trc::Result<Option<T>> {
let value = match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| Self::key_get_(conn, key).await).await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| Self::key_get_(conn, key).await).await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| Self::key_get_(conn, key).await).await
}
}?;
value.map(T::deserialize_owned).transpose()
}
pub async fn counter_get(&self, key: &[u8]) -> trc::Result<i64> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| Self::counter_get_(conn, key).await).await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| Self::counter_get_(conn, key).await).await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| Self::counter_get_(conn, key).await).await
}
}
}
pub async fn key_exists(&self, key: &[u8]) -> trc::Result<bool> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| Self::key_exists_(conn, key).await).await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| Self::key_exists_(conn, key).await).await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| Self::key_exists_(conn, key).await).await
}
}
}
async fn key_get_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<Option<Vec<u8>>> {
redis::cmd("GET").arg(key).query_async(conn).await
}
async fn counter_get_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<i64> {
redis::cmd("GET")
.arg(key)
.query_async::<Option<i64>>(conn)
.await
.map(|value| value.unwrap_or(0))
}
async fn key_exists_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<bool> {
conn.exists(key).await
}
async fn key_set_(
conn: &mut impl AsyncCommands,
key: &[u8],
value: &[u8],
expires: Option<u64>,
) -> RedisResult<()> {
if let Some(expires) = expires {
conn.set_ex(key, value, expires).await
} else {
conn.set(key, value).await
}
}
async fn key_incr_(
conn: &mut impl AsyncCommands,
key: &[u8],
value: i64,
expires: Option<u64>,
) -> RedisResult<i64> {
if let Some(expires) = expires {
INCR_EXPIRE
.key(key)
.arg(value)
.arg(expires as i64)
.invoke_async(conn)
.await
} else {
conn.incr(key, value).await
}
}
async fn try_lock_(
conn: &mut impl AsyncCommands,
key: &[u8],
expires: u64,
) -> RedisResult<bool> {
redis::cmd("SET")
.arg(key)
.arg(now() + expires)
.arg("NX")
.arg("EX")
.arg(expires as i64)
.query_async::<Option<String>>(conn)
.await
.map(|reply| reply.is_some())
}
async fn renew_lock_(
conn: &mut impl AsyncCommands,
key: &[u8],
expires: u64,
) -> RedisResult<bool> {
redis::cmd("SET")
.arg(key)
.arg(now() + expires)
.arg("XX")
.arg("EX")
.arg(expires as i64)
.query_async::<Option<String>>(conn)
.await
.map(|reply| reply.is_some())
}
async fn key_delete_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<()> {
conn.del(key).await
}
async fn key_delete_prefix_(conn: &mut impl AsyncCommands, prefix: &[u8]) -> RedisResult<()> {
let mut pattern = Vec::with_capacity(prefix.len() + 1);
pattern.extend_from_slice(prefix);
pattern.push(b'*');
let mut cursor = 0;
loop {
let (new_cursor, keys): (u64, Vec<Vec<u8>>) = redis::cmd("SCAN")
.cursor_arg(cursor)
.arg("MATCH")
.arg(&pattern)
.arg("COUNT")
.arg(100)
.query_async(conn)
.await?;
if !keys.is_empty() {
conn.del::<_, ()>(&keys).await?;
}
if new_cursor != 0 {
cursor = new_cursor;
} else {
return Ok(());
}
}
}
}
async fn with_conn<M, T>(
pool: &Pool<M>,
operation: impl AsyncFnOnce(&mut M::Type) -> RedisResult<T>,
) -> trc::Result<T>
where
M: Manager<Error = trc::Error>,
{
let mut conn = pool.get().await.map_err(into_error)?;
match operation(conn.as_mut()).await {
Ok(value) => Ok(value),
Err(err) => {
if is_stale_connection(&err) {
drop(Object::take(conn));
}
Err(into_error(err))
}
}
}
fn is_stale_connection(err: &RedisError) -> bool {
matches!(
err.retry_method(),
RetryMethod::Reconnect
| RetryMethod::ReconnectFromInitialConnections
| RetryMethod::RefreshSlotsAndRetry
| RetryMethod::MovedRedirect
| RetryMethod::AskRedirect
)
}