From 2e550ff139069af7e2770a74f7841fa66adbed9c Mon Sep 17 00:00:00 2001 From: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz> Date: Fri, 31 Jul 2026 13:17:42 -0400 Subject: [PATCH] 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> --- crates/buzz-relay/src/tunnel/directory.rs | 217 +++++++++++++++++++--- crates/buzz-relay/src/tunnel/reliable.rs | 11 +- 2 files changed, 200 insertions(+), 28 deletions(-) diff --git a/crates/buzz-relay/src/tunnel/directory.rs b/crates/buzz-relay/src/tunnel/directory.rs index 7cd6e22a2..6ccd40600 100644 --- a/crates/buzz-relay/src/tunnel/directory.rs +++ b/crates/buzz-relay/src/tunnel/directory.rs @@ -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, DirectoryError> { let keys = SessionKeys::new(community_id, session_id); let mut conn = self.pool.get().await?; - let value: Option = 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, DirectoryError> { let keys = SessionKeys::new(community_id, session_id); let mut conn = self.pool.get().await?; - let value: Option = 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 { + // 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, Option) = 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 { diff --git a/crates/buzz-relay/src/tunnel/reliable.rs b/crates/buzz-relay/src/tunnel/reliable.rs index 5153e104d..fae957264 100644 --- a/crates/buzz-relay/src/tunnel/reliable.rs +++ b/crates/buzz-relay/src/tunnel/reliable.rs @@ -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");