fix(e2e): make live relay checks topology-correct

Give unarchive membership signals unique event ids, derive live-test authorities and media URLs from the configured relay, use metadata-free PNG bytes, respect community-scoped Git pointer keys, and let the subscription-cap proof coexist with the independent admission quota.

Co-authored-by: npub122y0pqkertljmedu303rl0aqrj3w8pvu43t6jxm6875lzg6f2pwqegc3xc <5288f082d91aff2de5bc8be23fbfa01ca2e3859cac57a91b7a3fa9f12349505c@buzz.block.builderlab.xyz>
Signed-off-by: npub122y0pqkertljmedu303rl0aqrj3w8pvu43t6jxm6875lzg6f2pwqegc3xc <5288f082d91aff2de5bc8be23fbfa01ca2e3859cac57a91b7a3fa9f12349505c@buzz.block.builderlab.xyz>
This commit is contained in:
npub122y0pqkertljmedu303rl0aqrj3w8pvu43t6jxm6875lzg6f2pwqegc3xc
2026-08-02 13:55:26 -07:00
parent e120231d70
commit 8f5947d76c
5 changed files with 186 additions and 89 deletions
+35 -14
View File
@@ -904,6 +904,27 @@ pub async fn emit_membership_notification(
target_pubkey: &[u8],
actor_pubkey: &[u8],
notification_kind: u32,
) -> anyhow::Result<()> {
emit_membership_notification_with_nonce(
tenant,
state,
channel_id,
target_pubkey,
actor_pubkey,
notification_kind,
None,
)
.await
}
async fn emit_membership_notification_with_nonce(
tenant: &TenantContext,
state: &Arc<AppState>,
channel_id: Uuid,
target_pubkey: &[u8],
actor_pubkey: &[u8],
notification_kind: u32,
nonce: Option<&str>,
) -> anyhow::Result<()> {
let target_hex = hex::encode(target_pubkey);
let actor_hex = hex::encode(actor_pubkey);
@@ -931,8 +952,15 @@ pub async fn emit_membership_notification(
})
.to_string();
let mut tags = vec![p_tag, h_tag];
if let Some(nonce) = nonce {
tags.push(
Tag::parse(["nonce", nonce])
.map_err(|e| anyhow::anyhow!("failed to build nonce tag: {e}"))?,
);
}
let event = EventBuilder::new(Kind::Custom(notification_kind as u16), content)
.tags([p_tag, h_tag])
.tags(tags)
.sign_with_keys(&state.relay_keypair)
.map_err(|e| anyhow::anyhow!("failed to sign membership notification: {e}"))?;
@@ -1605,28 +1633,21 @@ async fn handle_edit_metadata(
// Resubscribe connected agents after restore: archiving evicts their
// live subscriptions (CLOSED "channel access revoked") and unarchive
// otherwise emits no signal that makes a connected agent resubscribe.
// We reuse the member_added notification (44100) purely as a resubscribe
// trigger — no membership actually changed here — because it flows on the
// agent's always-live global membership subscription, the same path
// remove/re-add uses to recover. Humans self-heal via the re-emitted
// kind:39000 discovery, so this is intentionally agent-scoped.
//
// Known limitation: emit_membership_notification builds a created_at=now
// event with no nonce, and insert_event skips fan-out on a duplicate id.
// Four sub-second toggles (archive->unarchive->archive->unarchive) on the
// same channel by the same actor could collide ids and skip a fan-out.
// Not reachable in practice — unarchive has a single human-driven caller;
// the reaper only auto-archives — so we don't engineer around it.
// We reuse member_added (44100) as the existing global resubscribe
// trigger. Each unarchive carries the accepted edit event id as a
// nonce so it cannot collide with the original member-added event (or
// another unarchive) within Nostr's one-second timestamp precision.
for member in
state.db.get_members(tenant.community(), channel_id).await?
{
if let Err(e) = emit_membership_notification(
if let Err(e) = emit_membership_notification_with_nonce(
tenant,
state,
channel_id,
&member.pubkey,
&actor_bytes,
KIND_MEMBER_ADDED_NOTIFICATION,
Some(&event.id.to_hex()),
)
.await
{
+54 -9
View File
@@ -189,16 +189,62 @@ impl GitS3Probe {
Self { bucket }
}
fn pointer_key(owner: &str, repo: &str) -> String {
fn configured_pointer_key(owner: &str, repo: &str) -> Option<String> {
let repo = repo.strip_suffix(".git").unwrap_or(repo);
if let Ok(community) = std::env::var("BUZZ_E2E_GIT_COMMUNITY_ID") {
return format!("repos/{community}/{owner}/{repo}/pointer");
std::env::var("BUZZ_E2E_GIT_COMMUNITY_ID")
.ok()
.map(|community| format!("repos/{community}/{owner}/{repo}/pointer"))
}
async fn find_pointer_key(&self, owner: &str, repo: &str) -> Option<String> {
if let Some(key) = Self::configured_pointer_key(owner, repo) {
return Some(key);
}
format!("repos/{owner}/{repo}/pointer")
// Git pointers are community-scoped. Discover the exact pointer created
// for this random owner/repo rather than probing the obsolete unscoped
// `repos/{owner}/{repo}/pointer` layout. The optional env override above
// remains useful when a harness already knows the tenant UUID.
let repo = repo.strip_suffix(".git").unwrap_or(repo);
let suffix = format!("/{owner}/{repo}/pointer");
let mut continuation = None;
let mut matches = Vec::new();
loop {
let (page, _) = self
.bucket
.list_page(
"repos/".to_string(),
None,
continuation.take(),
None,
Some(1000),
)
.await
.expect("list Git pointer namespace");
matches.extend(
page.contents
.into_iter()
.map(|object| object.key)
.filter(|key| key.ends_with(&suffix)),
);
if !page.is_truncated {
break;
}
continuation = page.next_continuation_token;
assert!(
continuation.is_some(),
"truncated Git pointer listing must include a continuation token"
);
}
assert!(
matches.len() <= 1,
"random E2E owner/repo resolved multiple community pointers: {matches:?}"
);
matches.pop()
}
async fn pointer(&self, owner: &str, repo: &str) -> Option<PointerSnapshot> {
let key = Self::pointer_key(owner, repo);
let key = self.find_pointer_key(owner, repo).await?;
match self.bucket.get_object(&key).await {
Ok(resp) => {
let etag = resp
@@ -225,10 +271,9 @@ impl GitS3Probe {
}
tokio::time::sleep(Duration::from_millis(250)).await;
}
panic!(
"S3 manifest pointer {} never appeared; git may have fallen back to disk",
Self::pointer_key(owner, repo)
);
let expected = Self::configured_pointer_key(owner, repo)
.unwrap_or_else(|| format!("repos/<community>/{owner}/{repo}/pointer"));
panic!("S3 manifest pointer {expected} never appeared")
}
async fn assert_manifest_exists(&self, digest: &str) {
@@ -19,6 +19,11 @@ fn relay_ws_url() -> String {
.replace("https://", "wss://")
}
fn relay_authority() -> String {
let url = url::Url::parse(&relay_http_url()).expect("relay HTTP URL");
url[url::Position::BeforeHost..url::Position::AfterPort].to_string()
}
fn http_client() -> Client {
Client::builder()
.timeout(Duration::from_secs(15))
@@ -98,15 +103,13 @@ fn tiny_jpeg() -> Vec<u8> {
}
fn tiny_png() -> Vec<u8> {
// Valid 2x2 red PNG generated by ffmpeg
// Valid, metadata-free 2x2 red PNG.
vec![
0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a, 0x00, 0x00, 0x00, 0x0d, 0x49, 0x48, 0x44,
0x52, 0x00, 0x00, 0x00, 0x02, 0x00, 0x00, 0x00, 0x02, 0x08, 0x02, 0x00, 0x00, 0x00, 0xfd,
0xd4, 0x9a, 0x73, 0x00, 0x00, 0x00, 0x09, 0x70, 0x48, 0x59, 0x73, 0x00, 0x00, 0x00, 0x01,
0x00, 0x00, 0x00, 0x01, 0x00, 0x4f, 0x25, 0xc4, 0xd6, 0x00, 0x00, 0x00, 0x10, 0x49, 0x44,
0x41, 0x54, 0x78, 0x9c, 0x63, 0xfc, 0xc3, 0x00, 0x02, 0x2c, 0x60, 0x92, 0x01, 0x00, 0x0d,
0x04, 0x01, 0x02, 0xbf, 0x50, 0x15, 0xb3, 0x00, 0x00, 0x00, 0x00, 0x49, 0x45, 0x4e, 0x44,
0xae, 0x42, 0x60, 0x82,
0xd4, 0x9a, 0x73, 0x00, 0x00, 0x00, 0x10, 0x49, 0x44, 0x41, 0x54, 0x78, 0x9c, 0x63, 0xf8,
0xcf, 0xc0, 0x00, 0x44, 0x0c, 0x10, 0x0a, 0x00, 0x1f, 0xee, 0x03, 0xfd, 0x8b, 0x5f, 0x14,
0xd4, 0x00, 0x00, 0x00, 0x00, 0x49, 0x45, 0x4e, 0x44, 0xae, 0x42, 0x60, 0x82,
]
}
@@ -371,7 +374,7 @@ async fn test_auth_server_tag_correct() {
Tag::parse(["t", "upload"]).unwrap(),
Tag::parse(["x", &sha256]).unwrap(),
Tag::parse(["expiration", &(now + 300).to_string()]).unwrap(),
Tag::parse(["server", "localhost:3000"]).unwrap(),
Tag::parse(["server", &relay_authority()]).unwrap(),
],
);
let resp = upload_with_auth(&client, &auth, &sha256, &jpeg).await;
@@ -558,6 +561,10 @@ async fn test_ws_valid_imeta() {
let resp = upload(&http, &keys, &jpeg).await;
assert_eq!(resp.status(), 200);
let upload_desc: serde_json::Value = resp.json().await.expect("upload descriptor");
let media_url = upload_desc["url"].as_str().expect("uploaded media URL");
let media_size = upload_desc["size"].as_u64().expect("uploaded media size");
// Connect via WebSocket
let mut client = BuzzTestClient::connect(&relay_ws_url(), &keys)
.await
@@ -569,10 +576,10 @@ async fn test_ws_valid_imeta() {
Tag::parse(["h", &channel_id]).unwrap(),
Tag::parse([
"imeta",
&format!("url http://localhost:3000/media/{sha256}.jpg"),
&format!("url {media_url}"),
"m image/jpeg",
&format!("x {sha256}"),
"size 347",
&format!("size {media_size}"),
])
.unwrap(),
])
+67 -47
View File
@@ -120,7 +120,12 @@ async fn seed_relay_member(host: &str, keys: &Keys, role: &str) {
}
async fn seed_relay_owner(keys: &Keys) {
seed_relay_member("localhost:3000", keys, "owner").await;
seed_relay_member(&relay_authority(), keys, "owner").await;
}
fn relay_authority() -> String {
let url = url::Url::parse(&relay_http_url()).expect("relay HTTP URL");
url[url::Position::BeforeHost..url::Position::AfterPort].to_string()
}
fn http_origin_for_host(host: &str) -> String {
@@ -315,7 +320,7 @@ async fn test_invite_claim_rejects_invalid_code() {
#[ignore]
async fn test_invite_mint_requires_owner_or_admin() {
let member = Keys::generate();
seed_relay_member("localhost:3000", &member, "member").await;
seed_relay_member(&relay_authority(), &member, "member").await;
let response = invite_post(&member, "/api/invites", "{}").await;
assert_eq!(response.status(), reqwest::StatusCode::FORBIDDEN);
@@ -791,10 +796,10 @@ async fn test_auth_event_kind_rejected() {
/// NIP-11 max_subscriptions must be enforced; (limit+1)th REQ gets CLOSED.
///
/// The relay's MAX_SUBSCRIPTIONS is 1024. Opening 1024 subs in a test is slow,
/// so we open a smaller batch and verify the NIP-11 advertised limit matches
/// the actual enforcement constant. The full-limit test is covered by the
/// NIP-11 assertion below (which verifies the advertised value is 1024).
/// This is a protocol-cap test, not an admission-throughput test. Open one REQ
/// at a time and wait out any shared fixed-window quota before retrying a REQ
/// rejected specifically as `rate-limited`, so production admission remains
/// enabled while the test deterministically reaches the independent 1024 cap.
#[tokio::test]
#[ignore]
async fn test_subscription_limit_enforced() {
@@ -802,60 +807,75 @@ async fn test_subscription_limit_enforced() {
let keys = Keys::generate();
let mut client = BuzzTestClient::connect(&url, &keys).await.expect("connect");
// Open 1024 subscriptions (the relay's MAX_SUBSCRIPTIONS).
for i in 0..1024 {
let sid = format!("limit-sub-{i}");
let filter = Filter::new().kind(Kind::Custom(9));
client
.subscribe(&sid, vec![filter])
.await
.expect("subscribe");
// Drain EOSE to avoid buffer buildup.
client
.collect_until_eose(&sid, Duration::from_secs(5))
.await
.expect("EOSE");
let filter = Filter::new().kind(Kind::Custom(49_999));
subscribe_until_eose(&mut client, &sid, filter).await;
}
let overflow_sid = sub_id("overflow");
// Use a kind that no other test writes, so we don't receive stale events.
let filter = Filter::new().kind(Kind::Custom(49999));
client
.subscribe(&overflow_sid, vec![filter])
.await
.expect("send REQ");
// Drain EOSE and stale events from the 100 earlier subscriptions
// until we receive the CLOSED for the overflow subscription.
let msg = loop {
let m = client
.recv_event(Duration::from_secs(5))
let filter = Filter::new().kind(Kind::Custom(49_999));
loop {
client
.subscribe(&overflow_sid, vec![filter.clone()])
.await
.expect("recv CLOSED (or timeout)");
match &m {
RelayMessage::Eose { .. } => continue,
RelayMessage::Event { .. } => continue, // stale event from earlier subs
_ => break m,
}
};
.expect("send overflow REQ");
match msg {
RelayMessage::Closed {
subscription_id,
message,
} => {
assert_eq!(subscription_id, overflow_sid);
assert!(
message.to_lowercase().contains("too many"),
"Expected 'too many' in CLOSED message, got: {message}"
);
match client
.recv_event(Duration::from_secs(6))
.await
.expect("recv overflow CLOSED")
{
RelayMessage::Closed {
subscription_id,
message,
} if subscription_id == overflow_sid && message.starts_with("rate-limited:") => {
tokio::time::sleep(Duration::from_secs(5)).await;
}
RelayMessage::Closed {
subscription_id,
message,
} => {
assert_eq!(subscription_id, overflow_sid);
assert!(
message.to_lowercase().contains("too many"),
"Expected 'too many' in CLOSED message, got: {message}"
);
break;
}
other => panic!("Expected CLOSED for overflow subscription, got {other:?}"),
}
other => panic!("Expected CLOSED for overflow subscription, got {other:?}"),
}
client.disconnect().await.expect("disconnect");
}
async fn subscribe_until_eose(client: &mut BuzzTestClient, sid: &str, filter: Filter) {
loop {
client
.subscribe(sid, vec![filter.clone()])
.await
.expect("subscribe");
match client
.recv_event(Duration::from_secs(6))
.await
.expect("EOSE or rate-limit CLOSED")
{
RelayMessage::Eose { subscription_id } => {
assert_eq!(subscription_id, sid);
return;
}
RelayMessage::Closed {
subscription_id,
message,
} if subscription_id == sid && message.starts_with("rate-limited:") => {
tokio::time::sleep(Duration::from_secs(5)).await;
}
other => panic!("unexpected response while opening {sid}: {other:?}"),
}
}
}
#[tokio::test]
#[ignore]
async fn test_nip11_relay_info() {
+14 -10
View File
@@ -103,14 +103,17 @@ if [ "$helm_rc" != 0 ]; then
fi
# ── 4. verify 3/3 Ready ───────────────────────────────────────────────────────
# Select the relay by its explicit component label. Optional worker/sidecar
# Deployments share the release instance and must never become the rollout target.
DEPLOYMENTS=$(kubectl -n "$NS" get deploy \
-l "app.kubernetes.io/instance=$RELEASE,app.kubernetes.io/component=relay" \
-o jsonpath='{range .items[*]}{.metadata.name}{"\n"}{end}')
[ "$(printf '%s\n' "$DEPLOYMENTS" | sed '/^$/d' | wc -l | tr -d ' ')" = "1" ] || \
die "expected exactly one relay Deployment, got: ${DEPLOYMENTS:-<none>}"
DEPLOY=$(printf '%s\n' "$DEPLOYMENTS" | sed '/^$/d')
# Find the relay Deployment: everything under this release named "buzz" except
# the bundled "*-minio" Deployment. (The chart fullname collapses
# "<release>-<chart>" to "<release>" when the release name already contains the
# chart name, so the name isn't always "<release>-buzz".)
DEPLOY=""
for d in $(kubectl -n "$NS" get deploy -l "app.kubernetes.io/instance=$RELEASE" \
-o jsonpath='{range .items[*]}{.metadata.name}{"\n"}{end}'); do
case "$d" in *-minio) continue;; esac
DEPLOY="$d"; break
done
[ -n "$DEPLOY" ] || die "could not locate the relay Deployment"
log "waiting for $REPLICAS relay pods Ready (deployment: $DEPLOY)"
kubectl -n "$NS" rollout status deployment/"$DEPLOY" --timeout=4m | tee "$EVID/rollout.txt"
kubectl -n "$NS" get pods -o wide | tee "$EVID/pods.txt"
@@ -126,8 +129,9 @@ log "probing /_readiness on each relay pod individually"
: > "$EVID/readiness.txt"
FAIL=0
RELAY_PODS=$(kubectl -n "$NS" get pods \
-l "app.kubernetes.io/instance=$RELEASE,app.kubernetes.io/component=relay" \
-o jsonpath='{range .items[*]}{.metadata.name}{"\n"}{end}')
-l "app.kubernetes.io/name=buzz,app.kubernetes.io/instance=$RELEASE" \
-o jsonpath='{range .items[*]}{.metadata.name}{" "}{.metadata.labels.app\.kubernetes\.io/component}{"\n"}{end}' \
| awk '$2 != "minio" && $2 != "minio-init" {print $1}')
for pod in $RELAY_PODS; do
body=$(kubectl -n "$NS" exec "$pod" -- \
sh -c 'curl -sS --max-time 5 http://127.0.0.1:8080/_readiness' 2>/dev/null || echo '<curl-failed>')