mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(relay): co-slot tunnel fencing keys
Co-authored-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz> Signed-off-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz>
This commit is contained in:
parent
052174a148
commit
2e550ff139
@@ -20,12 +20,26 @@ local generation_key = KEYS[2]
|
||||
local owner = ARGV[1]
|
||||
local profile = ARGV[2]
|
||||
local ttl_ms = tonumber(ARGV[3])
|
||||
local legacy_lease = ARGV[4]
|
||||
local legacy_generation = ARGV[5]
|
||||
|
||||
local current = redis.call('GET', lease_key)
|
||||
if current then
|
||||
return {'exists', current, redis.call('GET', generation_key) or ''}
|
||||
end
|
||||
|
||||
-- A live legacy lease is allowed to drain before this key format takes over.
|
||||
-- Deployments must drain legacy writers before enabling the tagged format; the
|
||||
-- separate legacy reads cannot be made atomic across Redis Cluster slots.
|
||||
if not redis.call('GET', generation_key) and legacy_lease ~= '' then
|
||||
return {'exists', legacy_lease, legacy_generation}
|
||||
end
|
||||
|
||||
-- Preserve the old non-expiring fence watermark on first use. SETNX accepts the
|
||||
-- integer as an exact string, avoiding Lua-number precision loss.
|
||||
if legacy_generation ~= '' then
|
||||
redis.call('SETNX', generation_key, legacy_generation)
|
||||
end
|
||||
local generation = redis.call('INCR', generation_key)
|
||||
local value = owner .. '|' .. tostring(generation) .. '|' .. profile
|
||||
redis.call('SET', lease_key, value, 'PX', ttl_ms)
|
||||
@@ -76,10 +90,15 @@ return {'lost', current, redis.call('GET', generation_key) or current_generation
|
||||
const VALIDATE_SCRIPT: &str = r#"
|
||||
local lease_key = KEYS[1]
|
||||
local generation_key = KEYS[2]
|
||||
local legacy_lease = ARGV[1]
|
||||
local legacy_generation = ARGV[2]
|
||||
|
||||
local current = redis.call('GET', lease_key) or ''
|
||||
local known_generation = redis.call('GET', generation_key) or ''
|
||||
return {current, known_generation}
|
||||
local known_generation = redis.call('GET', generation_key)
|
||||
if known_generation then
|
||||
return {current, known_generation}
|
||||
end
|
||||
return {legacy_lease, legacy_generation}
|
||||
"#;
|
||||
|
||||
/// Redis-backed owner directory for mesh tunnel sessions.
|
||||
@@ -207,6 +226,7 @@ impl SessionDirectory {
|
||||
let keys = SessionKeys::new(community_id, session_id);
|
||||
let ttl_ms = ttl_ms(self.lease_ttl)?;
|
||||
let mut conn = self.pool.get().await?;
|
||||
let legacy = read_legacy_keys(&mut conn, &keys).await?;
|
||||
let (status, value, _known_generation): (String, String, String) =
|
||||
Script::new(ACQUIRE_SCRIPT)
|
||||
.key(&keys.lease)
|
||||
@@ -214,6 +234,8 @@ impl SessionDirectory {
|
||||
.arg(owner_runtime_id.to_hex())
|
||||
.arg(profile.as_wire_str())
|
||||
.arg(ttl_ms)
|
||||
.arg(&legacy.lease)
|
||||
.arg(&legacy.generation)
|
||||
.invoke_async(&mut *conn)
|
||||
.await?;
|
||||
let lease = parse_lease(community_id, session_id, &value)?;
|
||||
@@ -310,14 +332,8 @@ impl SessionDirectory {
|
||||
) -> Result<Option<SessionLease>, DirectoryError> {
|
||||
let keys = SessionKeys::new(community_id, session_id);
|
||||
let mut conn = self.pool.get().await?;
|
||||
let value: Option<String> = redis::cmd("GET")
|
||||
.arg(&keys.lease)
|
||||
.query_async(&mut *conn)
|
||||
.await?;
|
||||
value
|
||||
.as_deref()
|
||||
.map(|v| parse_lease(community_id, session_id, v))
|
||||
.transpose()
|
||||
let (value, _known_generation) = read_keys_with_legacy_fallback(&mut conn, &keys).await?;
|
||||
parse_optional_lease(community_id, session_id, &value)
|
||||
}
|
||||
|
||||
/// Read the non-expiring generation counter for a session, if it exists.
|
||||
@@ -328,14 +344,8 @@ impl SessionDirectory {
|
||||
) -> Result<Option<u64>, DirectoryError> {
|
||||
let keys = SessionKeys::new(community_id, session_id);
|
||||
let mut conn = self.pool.get().await?;
|
||||
let value: Option<String> = redis::cmd("GET")
|
||||
.arg(&keys.generation)
|
||||
.query_async(&mut *conn)
|
||||
.await?;
|
||||
match value.as_deref() {
|
||||
Some(value) => parse_optional_generation(community_id, session_id, value),
|
||||
None => Ok(None),
|
||||
}
|
||||
let (_lease, generation) = read_keys_with_legacy_fallback(&mut conn, &keys).await?;
|
||||
parse_optional_generation(community_id, session_id, &generation)
|
||||
}
|
||||
|
||||
/// Validate a session-bearing mesh frame fence against Redis.
|
||||
@@ -356,11 +366,9 @@ impl SessionDirectory {
|
||||
.get()
|
||||
.await
|
||||
.map_err(|e| MeshError::Transport(format!("redis pool: {e}")))?;
|
||||
let (lease_value, known_generation): (String, String) = Script::new(VALIDATE_SCRIPT)
|
||||
.key(&keys.lease)
|
||||
.key(&keys.generation)
|
||||
.invoke_async(&mut *conn)
|
||||
.await?;
|
||||
let (lease_value, known_generation) = read_keys_with_legacy_fallback(&mut conn, &keys)
|
||||
.await
|
||||
.map_err(|e| MeshError::Transport(e.to_string()))?;
|
||||
let known_from_counter =
|
||||
parse_optional_generation(community_id, fenced.session_id, &known_generation)
|
||||
.map_err(|e| MeshError::Transport(e.to_string()))?
|
||||
@@ -453,18 +461,82 @@ impl SessionLease {
|
||||
struct SessionKeys {
|
||||
lease: String,
|
||||
generation: String,
|
||||
legacy_lease: String,
|
||||
legacy_generation: String,
|
||||
}
|
||||
|
||||
impl SessionKeys {
|
||||
fn new(community_id: CommunityId, session_id: Uuid) -> Self {
|
||||
let base = format!("buzz:{}:tunnel:{}", community_id, session_id);
|
||||
// The session hash tag is deliberately the first `{...}` segment. Redis
|
||||
// Cluster hashes both script keys to this session's slot, while standard
|
||||
// Redis treats the braces as ordinary key bytes.
|
||||
let base = format!("buzz:{{{session_id}}}:{community_id}:tunnel");
|
||||
let legacy_base = format!("buzz:{community_id}:tunnel:{session_id}");
|
||||
Self {
|
||||
lease: format!("{base}:lease"),
|
||||
generation: format!("{base}:generation"),
|
||||
legacy_lease: format!("{legacy_base}:lease"),
|
||||
legacy_generation: format!("{legacy_base}:generation"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
struct LegacyValues {
|
||||
lease: String,
|
||||
generation: String,
|
||||
}
|
||||
|
||||
async fn read_legacy_keys(
|
||||
conn: &mut deadpool_redis::Connection,
|
||||
keys: &SessionKeys,
|
||||
) -> Result<LegacyValues, redis::RedisError> {
|
||||
// This pipeline is intentionally non-atomic: MULTI/EXEC would put the two
|
||||
// legacy keys in one cross-slot transaction on Redis Cluster. These reads
|
||||
// bridge the drain-aware key migration; tagged state always wins once its
|
||||
// non-expiring generation key has been initialized.
|
||||
let (lease, generation): (Option<String>, Option<String>) = redis::pipe()
|
||||
.cmd("GET")
|
||||
.arg(&keys.legacy_lease)
|
||||
.cmd("GET")
|
||||
.arg(&keys.legacy_generation)
|
||||
.query_async(&mut **conn)
|
||||
.await?;
|
||||
Ok(LegacyValues {
|
||||
lease: lease.unwrap_or_default(),
|
||||
generation: generation.unwrap_or_default(),
|
||||
})
|
||||
}
|
||||
|
||||
async fn read_current_keys(
|
||||
conn: &mut deadpool_redis::Connection,
|
||||
keys: &SessionKeys,
|
||||
legacy: &LegacyValues,
|
||||
) -> Result<(String, String), redis::RedisError> {
|
||||
Script::new(VALIDATE_SCRIPT)
|
||||
.key(&keys.lease)
|
||||
.key(&keys.generation)
|
||||
.arg(&legacy.lease)
|
||||
.arg(&legacy.generation)
|
||||
.invoke_async(&mut **conn)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn read_keys_with_legacy_fallback(
|
||||
conn: &mut deadpool_redis::Connection,
|
||||
keys: &SessionKeys,
|
||||
) -> Result<(String, String), redis::RedisError> {
|
||||
let empty = LegacyValues {
|
||||
lease: String::new(),
|
||||
generation: String::new(),
|
||||
};
|
||||
let current = read_current_keys(conn, keys, &empty).await?;
|
||||
if !current.1.is_empty() {
|
||||
return Ok(current);
|
||||
}
|
||||
let legacy = read_legacy_keys(conn, keys).await?;
|
||||
read_current_keys(conn, keys, &legacy).await
|
||||
}
|
||||
|
||||
trait ProfileWireExt {
|
||||
fn as_wire_str(&self) -> &'static str;
|
||||
}
|
||||
@@ -613,9 +685,15 @@ mod tests {
|
||||
async fn clear_keys(directory: &SessionDirectory, community_id: CommunityId, session_id: Uuid) {
|
||||
let keys = SessionKeys::new(community_id, session_id);
|
||||
let mut conn = directory.pool.get().await.expect("redis conn");
|
||||
let _: () = redis::cmd("DEL")
|
||||
let _: () = redis::pipe()
|
||||
.cmd("DEL")
|
||||
.arg(keys.lease)
|
||||
.cmd("DEL")
|
||||
.arg(keys.generation)
|
||||
.cmd("DEL")
|
||||
.arg(keys.legacy_lease)
|
||||
.cmd("DEL")
|
||||
.arg(keys.legacy_generation)
|
||||
.query_async(&mut *conn)
|
||||
.await
|
||||
.expect("clear keys");
|
||||
@@ -647,15 +725,102 @@ mod tests {
|
||||
let keys = SessionKeys::new(community(), session());
|
||||
assert_eq!(
|
||||
keys.lease,
|
||||
format!("buzz:{}:tunnel:{}:lease", community(), session())
|
||||
format!("buzz:{{{}}}:{}:tunnel:lease", session(), community())
|
||||
);
|
||||
assert_eq!(
|
||||
keys.generation,
|
||||
format!("buzz:{{{}}}:{}:tunnel:generation", session(), community())
|
||||
);
|
||||
assert_eq!(
|
||||
keys.legacy_lease,
|
||||
format!("buzz:{}:tunnel:{}:lease", community(), session())
|
||||
);
|
||||
assert_eq!(
|
||||
keys.legacy_generation,
|
||||
format!("buzz:{}:tunnel:{}:generation", community(), session())
|
||||
);
|
||||
assert_eq!(keys.lease.find('{'), Some(5));
|
||||
assert_eq!(keys.generation.find('{'), Some(5));
|
||||
let tag = format!("{{{}}}", session());
|
||||
assert!(keys.lease.starts_with(&format!("buzz:{tag}:")));
|
||||
assert!(keys.generation.starts_with(&format!("buzz:{tag}:")));
|
||||
assert_ne!(keys.lease, keys.generation);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn legacy_generation_and_live_lease_migrate_without_fence_regression() {
|
||||
let Some(directory) = redis_directory_if_available().await else {
|
||||
return;
|
||||
};
|
||||
let community_id = community();
|
||||
let session_id = Uuid::new_v4();
|
||||
clear_keys(&directory, community_id, session_id).await;
|
||||
let keys = SessionKeys::new(community_id, session_id);
|
||||
let legacy_lease = format!("{}|41|reliable-stream", runtime(1).to_hex());
|
||||
let mut conn = directory.pool.get().await.expect("redis conn");
|
||||
let _: () = redis::pipe()
|
||||
.cmd("SET")
|
||||
.arg(&keys.legacy_generation)
|
||||
.arg(41_u64)
|
||||
.cmd("SET")
|
||||
.arg(&keys.legacy_lease)
|
||||
.arg(&legacy_lease)
|
||||
.arg("PX")
|
||||
.arg(5_000_u64)
|
||||
.query_async(&mut *conn)
|
||||
.await
|
||||
.expect("seed legacy keys");
|
||||
drop(conn);
|
||||
|
||||
let existing = directory
|
||||
.acquire(
|
||||
community_id,
|
||||
session_id,
|
||||
runtime(2),
|
||||
Profile::ReliableStream,
|
||||
)
|
||||
.await
|
||||
.expect("legacy lease remains authoritative");
|
||||
assert!(matches!(existing, AcquireResult::Exists(ref lease) if lease.generation == 41));
|
||||
let legacy = match existing {
|
||||
AcquireResult::Exists(lease) => lease,
|
||||
AcquireResult::Acquired(_) => unreachable!(),
|
||||
};
|
||||
assert_eq!(
|
||||
directory.lookup(community_id, session_id).await.unwrap(),
|
||||
Some(legacy)
|
||||
);
|
||||
let mut conn = directory.pool.get().await.expect("redis conn");
|
||||
let _: () = redis::cmd("DEL")
|
||||
.arg(&keys.legacy_lease)
|
||||
.query_async(&mut *conn)
|
||||
.await
|
||||
.expect("drain legacy lease");
|
||||
drop(conn);
|
||||
|
||||
let migrated = match directory
|
||||
.acquire(
|
||||
community_id,
|
||||
session_id,
|
||||
runtime(2),
|
||||
Profile::ReliableStream,
|
||||
)
|
||||
.await
|
||||
.expect("acquire tagged lease")
|
||||
{
|
||||
AcquireResult::Acquired(lease) => lease,
|
||||
AcquireResult::Exists(_) => panic!("released legacy lease must drain"),
|
||||
};
|
||||
assert_eq!(migrated.generation, 42);
|
||||
assert_eq!(
|
||||
directory
|
||||
.known_generation(community_id, session_id)
|
||||
.await
|
||||
.unwrap(),
|
||||
Some(42)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn acquire_conflict_renew_release_and_monotonic_takeover() {
|
||||
let Some(directory) = redis_directory_if_available().await else {
|
||||
|
||||
@@ -694,7 +694,8 @@ mod tests {
|
||||
}
|
||||
|
||||
async fn clear_keys(directory: &SessionDirectory, community_id: CommunityId, session_id: Uuid) {
|
||||
let base = format!("buzz:{}:tunnel:{}", community_id, session_id);
|
||||
let base = format!("buzz:{{{session_id}}}:{community_id}:tunnel");
|
||||
let legacy_base = format!("buzz:{community_id}:tunnel:{session_id}");
|
||||
let _ = directory
|
||||
.release(&SessionLease {
|
||||
community_id,
|
||||
@@ -705,9 +706,15 @@ mod tests {
|
||||
})
|
||||
.await;
|
||||
let mut conn = pool().get().await.expect("redis conn");
|
||||
let _: () = redis::cmd("DEL")
|
||||
let _: () = redis::pipe()
|
||||
.cmd("DEL")
|
||||
.arg(format!("{base}:lease"))
|
||||
.cmd("DEL")
|
||||
.arg(format!("{base}:generation"))
|
||||
.cmd("DEL")
|
||||
.arg(format!("{legacy_base}:lease"))
|
||||
.cmd("DEL")
|
||||
.arg(format!("{legacy_base}:generation"))
|
||||
.query_async(&mut *conn)
|
||||
.await
|
||||
.expect("clear keys");
|
||||
|
||||
Reference in New Issue
Block a user