diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2c471d5a7..dc24302de 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -709,6 +709,15 @@ jobs: --run-ignored ignored-only env: DATABASE_URL: postgres://buzz:${{ env.BUZZ_TEST_POSTGRES_PASSWORD }}@localhost:5432/buzz + - name: Tunnel fence Redis integration tests + run: | + cargo nextest run \ + --archive-file target/ci/backend-integration-tests.tar.zst \ + -E 'package(buzz-relay) and test(/tunnel::directory::tests/)' \ + --run-ignored ignored-only + env: + REDIS_URL: redis://localhost:6379 + - name: NIP-ER reminder e2e # Feature e2e for NIP-ER (Event Reminders, kind:30300): write-path # validation, author-only read filtering, and scheduler delivery against diff --git a/crates/buzz-relay/src/config.rs b/crates/buzz-relay/src/config.rs index 85a0ca2ef..c4821a130 100644 --- a/crates/buzz-relay/src/config.rs +++ b/crates/buzz-relay/src/config.rs @@ -152,6 +152,13 @@ pub struct Config { /// no-regression rollout. pub mesh: buzz_relay_mesh::MeshConfig, + /// Tunnel fence Redis key format used by the mesh session directory. + /// + /// Defaults to `legacy`; operators must explicitly select `tagged` only + /// after completing the drain phase documented in + /// `docs/tunnel-fence-key-migration.md`. + pub tunnel_fence_key_format: crate::tunnel::directory::SessionKeyFormat, + /// Testbed-only reliable-stream echo consumer (`BUZZ_MESH_DEMO_ECHO`). /// When `on`, the owner side of an inbound reliable mesh stream echoes /// every validated `Data` frame back to the sender — a transport/ @@ -279,6 +286,18 @@ pub struct Config { pub serve_git_web_gui: bool, } +fn parse_tunnel_fence_key_format( + value: Option<&str>, +) -> Result { + match value { + None | Some("legacy") => Ok(crate::tunnel::directory::SessionKeyFormat::Legacy), + Some("tagged") => Ok(crate::tunnel::directory::SessionKeyFormat::Tagged), + Some(value) => Err(ConfigError::InvalidValue(format!( + "BUZZ_TUNNEL_FENCE_KEY_FORMAT must be 'legacy' or 'tagged', got {value:?}" + ))), + } +} + fn parse_bind_addr(raw: &str) -> Result { raw.parse::() .map_err(|e| ConfigError::InvalidBindAddr(e.to_string())) @@ -560,6 +579,16 @@ impl Config { registry_refresh: std::time::Duration::from_secs(15), }; + let tunnel_fence_key_format = match std::env::var("BUZZ_TUNNEL_FENCE_KEY_FORMAT") { + Ok(value) => parse_tunnel_fence_key_format(Some(&value))?, + Err(std::env::VarError::NotPresent) => parse_tunnel_fence_key_format(None)?, + Err(std::env::VarError::NotUnicode(_)) => { + return Err(ConfigError::InvalidValue( + "BUZZ_TUNNEL_FENCE_KEY_FORMAT must be valid Unicode".to_string(), + )); + } + }; + // Demo echo opt-in: same strict pattern as BUZZ_MESH — explicit // `on`/`true`/`1` only, anything else (absent, `off`, typos) is off. let mesh_demo_echo = std::env::var("BUZZ_MESH_DEMO_ECHO") @@ -956,6 +985,7 @@ impl Config { require_relay_membership, huddle_audio_available, mesh, + tunnel_fence_key_format, mesh_demo_echo, relay_owner_pubkey, relay_operator_api_origin, @@ -997,6 +1027,25 @@ mod tests { // value set by `invalid_bind_addr_returns_error`, causing a flaky failure. static ENV_MUTEX: std::sync::Mutex<()> = std::sync::Mutex::new(()); + #[test] + fn tunnel_fence_key_format_is_default_off_and_strict() { + use crate::tunnel::directory::SessionKeyFormat; + + assert_eq!( + parse_tunnel_fence_key_format(None).unwrap(), + SessionKeyFormat::Legacy + ); + assert_eq!( + parse_tunnel_fence_key_format(Some("legacy")).unwrap(), + SessionKeyFormat::Legacy + ); + assert_eq!( + parse_tunnel_fence_key_format(Some("tagged")).unwrap(), + SessionKeyFormat::Tagged + ); + assert!(parse_tunnel_fence_key_format(Some("on")).is_err()); + } + #[test] fn defaults_are_valid() { let _guard = ENV_MUTEX.lock().unwrap(); diff --git a/crates/buzz-relay/src/mesh_boot.rs b/crates/buzz-relay/src/mesh_boot.rs index 20e550aa0..fc16e317f 100644 --- a/crates/buzz-relay/src/mesh_boot.rs +++ b/crates/buzz-relay/src/mesh_boot.rs @@ -508,7 +508,11 @@ pub async fn boot_mesh( transport.set_inbound(Box::new(dispatcher.clone())); Ok(Some(MeshHandle { - directory: SessionDirectory::new(redis_pool), + directory: SessionDirectory::with_key_format( + redis_pool, + std::time::Duration::from_secs(30), + config.tunnel_fence_key_format, + ), transport, membership: membership_arc, local_runtime_id: runtime_id, diff --git a/crates/buzz-relay/src/tunnel/directory.rs b/crates/buzz-relay/src/tunnel/directory.rs index 7cd6e22a2..4352418b5 100644 --- a/crates/buzz-relay/src/tunnel/directory.rs +++ b/crates/buzz-relay/src/tunnel/directory.rs @@ -14,18 +14,42 @@ use uuid::Uuid; const DEFAULT_LEASE_TTL: Duration = Duration::from_secs(30); +/// Redis key format used for tunnel fences during the staged cluster migration. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum SessionKeyFormat { + /// Original untagged keys. This remains the default during the drain deploy. + #[default] + Legacy, + /// Cluster-safe keys co-slotted by a first-position session hash tag. + Tagged, +} + const ACQUIRE_SCRIPT: &str = r#" local lease_key = KEYS[1] 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 +100,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. @@ -87,6 +116,7 @@ return {current, known_generation} pub struct SessionDirectory { pool: deadpool_redis::Pool, lease_ttl: Duration, + key_format: SessionKeyFormat, } /// Active session ownership lease read from Redis. @@ -184,12 +214,25 @@ pub enum DirectoryError { impl SessionDirectory { /// Create a directory backed by `pool` with the default lease TTL. pub fn new(pool: deadpool_redis::Pool) -> Self { - Self::with_lease_ttl(pool, DEFAULT_LEASE_TTL) + Self::with_key_format(pool, DEFAULT_LEASE_TTL, SessionKeyFormat::Legacy) } /// Create a directory backed by `pool` with an explicit lease TTL. pub fn with_lease_ttl(pool: deadpool_redis::Pool, lease_ttl: Duration) -> Self { - Self { pool, lease_ttl } + Self::with_key_format(pool, lease_ttl, SessionKeyFormat::Legacy) + } + + /// Create a directory using an explicitly selected migration key format. + pub fn with_key_format( + pool: deadpool_redis::Pool, + lease_ttl: Duration, + key_format: SessionKeyFormat, + ) -> Self { + Self { + pool, + lease_ttl, + key_format, + } } /// Attempt to create/take over the session lease. @@ -207,13 +250,27 @@ 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 (lease_key, generation_key, legacy) = match self.key_format { + SessionKeyFormat::Legacy => ( + &keys.legacy_lease, + &keys.legacy_generation, + LegacyValues::default(), + ), + SessionKeyFormat::Tagged => ( + &keys.lease, + &keys.generation, + read_legacy_keys(&mut conn, &keys).await?, + ), + }; let (status, value, _known_generation): (String, String, String) = Script::new(ACQUIRE_SCRIPT) - .key(&keys.lease) - .key(&keys.generation) + .key(lease_key) + .key(generation_key) .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)?; @@ -248,8 +305,8 @@ impl SessionDirectory { let ttl_ms = ttl_ms(self.lease_ttl)?; let mut conn = self.pool.get().await?; let (status, value, known_generation): (String, String, String) = Script::new(RENEW_SCRIPT) - .key(&keys.lease) - .key(&keys.generation) + .key(keys.lease(self.key_format)) + .key(keys.generation(self.key_format)) .arg(lease.owner_runtime_id.to_hex()) .arg(lease.generation) .arg(ttl_ms) @@ -279,8 +336,8 @@ impl SessionDirectory { let mut conn = self.pool.get().await?; let (status, value, known_generation): (String, String, String) = Script::new(RELEASE_SCRIPT) - .key(&keys.lease) - .key(&keys.generation) + .key(keys.lease(self.key_format)) + .key(keys.generation(self.key_format)) .arg(lease.owner_runtime_id.to_hex()) .arg(lease.generation) .invoke_async(&mut *conn) @@ -310,14 +367,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(&mut conn, &keys, self.key_format).await?; + parse_optional_lease(community_id, session_id, &value) } /// Read the non-expiring generation counter for a session, if it exists. @@ -328,14 +379,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(&mut conn, &keys, self.key_format).await?; + parse_optional_generation(community_id, session_id, &generation) } /// Validate a session-bearing mesh frame fence against Redis. @@ -356,11 +401,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(&mut conn, &keys, self.key_format) + .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 +496,120 @@ impl SessionLease { struct SessionKeys { lease: String, generation: String, + legacy_lease: String, + legacy_generation: String, } impl SessionKeys { + fn lease(&self, format: SessionKeyFormat) -> &str { + match format { + SessionKeyFormat::Legacy => &self.legacy_lease, + SessionKeyFormat::Tagged => &self.lease, + } + } + + fn generation(&self, format: SessionKeyFormat) -> &str { + match format { + SessionKeyFormat::Legacy => &self.legacy_generation, + SessionKeyFormat::Tagged => &self.generation, + } + } + 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"), } } } +#[derive(Default)] +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> { + read_current_keys_for(conn, keys, SessionKeyFormat::Tagged, legacy).await +} + +async fn read_keys( + conn: &mut deadpool_redis::Connection, + keys: &SessionKeys, + format: SessionKeyFormat, +) -> Result<(String, String), redis::RedisError> { + match format { + SessionKeyFormat::Legacy => { + let empty = LegacyValues::default(); + read_current_keys_for(conn, keys, format, &empty).await + } + SessionKeyFormat::Tagged => read_keys_with_legacy_fallback(conn, keys).await, + } +} + +async fn read_current_keys_for( + conn: &mut deadpool_redis::Connection, + keys: &SessionKeys, + format: SessionKeyFormat, + legacy: &LegacyValues, +) -> Result<(String, String), redis::RedisError> { + Script::new(VALIDATE_SCRIPT) + .key(keys.lease(format)) + .key(keys.generation(format)) + .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; } @@ -597,25 +742,36 @@ mod tests { .expect("create redis pool") } - async fn redis_directory_if_available() -> Option { + async fn redis_directory() -> SessionDirectory { let pool = pool(); - let mut conn = pool.get().await.ok()?; + let mut conn = pool + .get() + .await + .expect("REDIS_URL must be reachable for ignored tunnel directory tests"); redis::cmd("PING") .query_async::(&mut *conn) .await - .ok()?; - Some(SessionDirectory::with_lease_ttl( + .expect("REDIS_URL must answer PING for ignored tunnel directory tests"); + drop(conn); + SessionDirectory::with_key_format( pool, Duration::from_millis(150), - )) + SessionKeyFormat::Tagged, + ) } 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,20 +803,185 @@ 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 acquire_conflict_renew_release_and_monotonic_takeover() { - let Some(directory) = redis_directory_if_available().await else { - return; + #[ignore = "requires REDIS_URL; run by backend-integration"] + async fn legacy_generation_and_live_lease_migrate_without_fence_regression() { + let directory = redis_directory().await; + 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] + #[ignore = "requires REDIS_URL; run by backend-integration"] + async fn tagged_state_is_authoritative_when_both_formats_conflict() { + let directory = redis_directory().await; + 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 tagged = SessionLease { + community_id, + session_id, + owner_runtime_id: runtime(2), + generation: 42, + profile: Profile::ReliableStream, + }; + let legacy_value = format!("{}|99|reliable-stream", runtime(1).to_hex()); + let tagged_value = format!("{}|42|reliable-stream", runtime(2).to_hex()); + let mut conn = directory.pool.get().await.expect("redis conn"); + let _: () = redis::pipe() + .cmd("SET") + .arg(&keys.legacy_generation) + .arg(99_u64) + .cmd("SET") + .arg(&keys.legacy_lease) + .arg(&legacy_value) + .arg("PX") + .arg(5_000_u64) + .cmd("SET") + .arg(&keys.generation) + .arg(42_u64) + .cmd("SET") + .arg(&keys.lease) + .arg(&tagged_value) + .arg("PX") + .arg(5_000_u64) + .query_async(&mut *conn) + .await + .expect("seed conflicting key formats"); + drop(conn); + + assert_eq!( + directory.lookup(community_id, session_id).await.unwrap(), + Some(tagged.clone()) + ); + assert_eq!( + directory + .known_generation(community_id, session_id) + .await + .unwrap(), + Some(42) + ); + assert!(directory + .validate_fenced_header(community_id, &tagged.fenced_header()) + .await + .is_ok()); + assert!(matches!( + directory + .validate_fenced_header( + community_id, + &FencedHeader { + session_id, + generation: 99, + owner_runtime_id: runtime(1), + }, + ) + .await, + Err(MeshError::FutureGeneration { + known_generation: 42, + .. + }) + )); + assert!(matches!( + directory + .acquire(community_id, session_id, runtime(3), Profile::ReliableStream) + .await + .unwrap(), + AcquireResult::Exists(ref lease) if *lease == tagged + )); + } + + #[tokio::test] + #[ignore = "requires REDIS_URL; run by backend-integration"] + async fn acquire_conflict_renew_release_and_monotonic_takeover() { + let directory = redis_directory().await; let community_id = community(); let session_id = Uuid::new_v4(); clear_keys(&directory, community_id, session_id).await; @@ -733,10 +1054,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires REDIS_URL; run by backend-integration"] async fn takeover_after_ttl_expiry_increments_non_expiring_counter() { - let Some(directory) = redis_directory_if_available().await else { - return; - }; + let directory = redis_directory().await; let community_id = community(); let session_id = Uuid::new_v4(); clear_keys(&directory, community_id, session_id).await; @@ -785,10 +1105,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires REDIS_URL; run by backend-integration"] async fn validate_returns_typed_fence_rejections() { - let Some(directory) = redis_directory_if_available().await else { - return; - }; + let directory = redis_directory().await; let community_id = community(); let session_id = Uuid::new_v4(); clear_keys(&directory, community_id, session_id).await; @@ -880,10 +1199,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires REDIS_URL; run by backend-integration"] async fn validate_returns_no_active_lease_after_expiry_before_takeover() { - let Some(directory) = redis_directory_if_available().await else { - return; - }; + let directory = redis_directory().await; let community_id = community(); let session_id = Uuid::new_v4(); clear_keys(&directory, community_id, session_id).await; 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"); diff --git a/docs/tunnel-fence-key-migration.md b/docs/tunnel-fence-key-migration.md new file mode 100644 index 000000000..4ef517135 --- /dev/null +++ b/docs/tunnel-fence-key-migration.md @@ -0,0 +1,103 @@ +# Tunnel fence Redis key migration + +Tunnel lease and generation keys must share a Redis Cluster hash slot. The +cluster-safe tagged format is deliberately gated by +`BUZZ_TUNNEL_FENCE_KEY_FORMAT`; its default is `legacy` so an ordinary image +rollout cannot activate new-format writers while old writers still run. + +## Required two-phase deployment + +1. **Drain deploy:** deploy the new binary everywhere with + `BUZZ_TUNNEL_FENCE_KEY_FORMAT=legacy` (or unset). Keep the existing + standalone Redis endpoint. Wait until every old binary has terminated and + its maximum graceful-shutdown/lease window has elapsed. Verify no old pod + can acquire or renew a legacy lease. +2. **Enable deploy (non-rolling):** scale the relay deployment to zero and wait + until the maximum graceful-shutdown/lease window has elapsed. Confirm that no + relay process remains and no legacy lease can still be renewed. While writers + remain stopped, build a manifest from the source endpoint containing every + legacy `buzz:*:tunnel:*:generation` key and its exact string value. Transfer + those non-expiring generation keys, under the same legacy key names, to the + cluster-enforcing target endpoint. Do **not** copy lease keys: all expiring + leases must remain drained and must not be resurrected on the target. One + executable procedure is below; run it from a controlled host with + `redis-cli` access to both endpoints: + + ```bash + set -euo pipefail + : "${SOURCE_REDIS_URL:?set the drained source Redis URL}" + : "${TARGET_REDIS_URL:?set the target Redis URL}" + + work=$(mktemp -d) + trap 'rm -rf "$work"' EXIT + + manifest() { + endpoint=$1 + output=$2 + redis-cli -u "$endpoint" --scan \ + --pattern 'buzz:*:tunnel:*:generation' \ + | LC_ALL=C sort -u \ + | while IFS= read -r key; do + value=$(redis-cli -u "$endpoint" --raw GET "$key") + test -n "$value" + printf '%s\t%s\n' "$key" "$value" + done >"$output" + } + + # Never overwrite initialized tagged state without operator reconciliation. + test -z "$(redis-cli -u "$TARGET_REDIS_URL" --scan \ + --pattern 'buzz:{*}:*:tunnel:generation' | head -n 1)" + + manifest "$SOURCE_REDIS_URL" "$work/source-before" + while IFS=$'\t' read -r key value; do + redis-cli -u "$TARGET_REDIS_URL" SET "$key" "$value" >/dev/null + done <"$work/source-before" + + # Writers remain stopped while these final manifests are produced. + manifest "$SOURCE_REDIS_URL" "$work/source-after" + manifest "$TARGET_REDIS_URL" "$work/target-after" + cmp "$work/source-before" "$work/source-after" + cmp "$work/source-after" "$work/target-after" + + # Neither legacy nor tagged leases may exist on the target. + test -z "$(redis-cli -u "$TARGET_REDIS_URL" --scan \ + --pattern 'buzz:*:tunnel:*:lease' | head -n 1)" + test -z "$(redis-cli -u "$TARGET_REDIS_URL" --scan \ + --pattern 'buzz:{*}:*:tunnel:lease' | head -n 1)" + ``` + + Treat any non-zero exit as a failed transfer; do not continue to scale-up. + + Before starting any relay, rescan both endpoints and verify that the complete + source and target legacy-generation key sets are identical and that every + corresponding string value is byte-for-byte equal to the manifest. Also + verify that the source manifest did not change during transfer and that the + target has no legacy or tagged tunnel `*:lease` keys. If a generation key is + missing, extra, changed, unreadable, or cannot be written, **abort the + cutover**, keep all relays scaled to zero, and repair/repeat the transfer; + never start tagged writers from a partial or unverified copy. If the target + already contains tagged generation keys, stop and perform an + operator-reviewed reconciliation rather than overwriting or assuming that + the legacy manifest is authoritative. + + Only after that verification succeeds, set + `BUZZ_TUNNEL_FENCE_KEY_FORMAT=tagged` on the whole deployment, switch the + configured Redis endpoint to the verified target, and scale up. Do **not** + start the first tagged pod while any legacy-mode pod is alive. The first + tagged acquisition reads the copied legacy watermark from its currently + configured target pool, seeds its non-expiring tagged generation, and then + increments it. Tagged state is authoritative once initialized. + +Do not mix `legacy` and `tagged` writers. The two legacy keys occupy different +cluster slots, so Redis cannot atomically prove that an old writer did not race +the migration. + +## Rollback rule + +Before phase 2, roll back normally with the format left `legacy`. After any pod +has enabled `tagged`, **do not roll back to an old binary or set the format back +to `legacy`**. Roll back only to a tagged-capable binary with the format still +`tagged`; otherwise a legacy writer can reuse an already-issued generation and +break fencing. Restoring legacy mode requires a full mesh outage, drainage of +all tagged leases, and an operator-reviewed generation migration; it is not a +rolling rollback.