fix(relay): gate tunnel fence key migration

Co-authored-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz>
Signed-off-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@buzz.block.builderlab.xyz>
This commit is contained in:
npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr
2026-07-31 14:00:50 -04:00
parent 2e550ff139
commit ae4ce2e945
5 changed files with 286 additions and 35 deletions
+9
View File
@@ -692,6 +692,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
+49
View File
@@ -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<crate::tunnel::directory::SessionKeyFormat, ConfigError> {
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<SocketAddr, ConfigError> {
raw.parse::<SocketAddr>()
.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();
+5 -1
View File
@@ -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,
+187 -34
View File
@@ -14,6 +14,16 @@ 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]
@@ -106,6 +116,7 @@ return {legacy_lease, legacy_generation}
pub struct SessionDirectory {
pool: deadpool_redis::Pool,
lease_ttl: Duration,
key_format: SessionKeyFormat,
}
/// Active session ownership lease read from Redis.
@@ -203,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.
@@ -226,11 +250,22 @@ 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 (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)
@@ -270,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)
@@ -301,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)
@@ -332,7 +367,7 @@ impl SessionDirectory {
) -> Result<Option<SessionLease>, DirectoryError> {
let keys = SessionKeys::new(community_id, session_id);
let mut conn = self.pool.get().await?;
let (value, _known_generation) = read_keys_with_legacy_fallback(&mut conn, &keys).await?;
let (value, _known_generation) = read_keys(&mut conn, &keys, self.key_format).await?;
parse_optional_lease(community_id, session_id, &value)
}
@@ -344,7 +379,7 @@ impl SessionDirectory {
) -> Result<Option<u64>, DirectoryError> {
let keys = SessionKeys::new(community_id, session_id);
let mut conn = self.pool.get().await?;
let (_lease, generation) = read_keys_with_legacy_fallback(&mut conn, &keys).await?;
let (_lease, generation) = read_keys(&mut conn, &keys, self.key_format).await?;
parse_optional_generation(community_id, session_id, &generation)
}
@@ -366,7 +401,7 @@ impl SessionDirectory {
.get()
.await
.map_err(|e| MeshError::Transport(format!("redis pool: {e}")))?;
let (lease_value, known_generation) = read_keys_with_legacy_fallback(&mut conn, &keys)
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 =
@@ -466,6 +501,20 @@ struct SessionKeys {
}
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 {
// The session hash tag is deliberately the first `{...}` segment. Redis
// Cluster hashes both script keys to this session's slot, while standard
@@ -481,6 +530,7 @@ impl SessionKeys {
}
}
#[derive(Default)]
struct LegacyValues {
lease: String,
generation: String,
@@ -511,10 +561,33 @@ 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)
.key(&keys.generation)
.key(keys.lease(format))
.key(keys.generation(format))
.arg(&legacy.lease)
.arg(&legacy.generation)
.invoke_async(&mut **conn)
@@ -669,17 +742,22 @@ mod tests {
.expect("create redis pool")
}
async fn redis_directory_if_available() -> Option<SessionDirectory> {
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::<String>(&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) {
@@ -748,10 +826,9 @@ mod tests {
}
#[tokio::test]
#[ignore = "requires REDIS_URL; run by backend-integration"]
async fn legacy_generation_and_live_lease_migrate_without_fence_regression() {
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;
@@ -822,10 +899,89 @@ mod tests {
}
#[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 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;
@@ -898,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;
@@ -950,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;
@@ -1045,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;
+36
View File
@@ -0,0 +1,36 @@
# 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. Set
`BUZZ_TUNNEL_FENCE_KEY_FORMAT=tagged` on the whole deployment, move to the
cluster-enforcing endpoint, and then scale up. Do **not** start the first
tagged pod while any legacy-mode pod is alive. The first tagged acquisition
seeds its non-expiring generation from the drained legacy watermark before
incrementing 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.