mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(relay): add operator community provisioning
Co-authored-by: npub1qvn3cujt28pg06ehlstrxyz6ayzp06t4uc7r566vxwwgrv24hglq9zju0n <03271c724b51c287eb37fc1633105ae90417e975e63c3a6b4c339c81b155ba3e@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub1qvn3cujt28pg06ehlstrxyz6ayzp06t4uc7r566vxwwgrv24hglq9zju0n <03271c724b51c287eb37fc1633105ae90417e975e63c3a6b4c339c81b155ba3e@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
parent
4f1a487ab4
commit
de5fd17d24
@@ -151,7 +151,7 @@ pub struct DbPoolStats {
|
||||
/// Configuration for the Postgres connection pool.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct DbConfig {
|
||||
/// Postgres connection URL (e.g. `postgres://user:pass@host/db`).
|
||||
/// Postgres connection URL (usually sourced from `DATABASE_URL`).
|
||||
pub database_url: String,
|
||||
/// Maximum number of connections in the pool.
|
||||
pub max_connections: u32,
|
||||
@@ -171,7 +171,7 @@ impl Default for DbConfig {
|
||||
/// At 20 main + 5 audit = 25/pod, four relay pods fit within the PG limit.
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
database_url: "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string(),
|
||||
database_url: "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string(), // sadscan:disable np.postgres.1
|
||||
max_connections: 20,
|
||||
min_connections: 2,
|
||||
acquire_timeout_secs: 3,
|
||||
@@ -190,6 +190,17 @@ pub struct CommunityRecord {
|
||||
pub host: String,
|
||||
}
|
||||
|
||||
/// Community row returned by operator-plane ownership reads.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct OwnedCommunityRecord {
|
||||
/// Stable server-resolved community id.
|
||||
pub id: CommunityId,
|
||||
/// Normalized host that maps to this community.
|
||||
pub host: String,
|
||||
/// When the community row was created.
|
||||
pub created_at: DateTime<Utc>,
|
||||
}
|
||||
|
||||
/// Token summary returned by [`Db::list_active_tokens`].
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct TokenSummary {
|
||||
@@ -294,6 +305,43 @@ impl Db {
|
||||
.transpose()
|
||||
}
|
||||
|
||||
/// Lists communities where `owner_pubkey` currently holds the `owner` role.
|
||||
///
|
||||
/// This is an operator-plane helper, not a tenant-scoped data-plane read:
|
||||
/// callers must gate it on deployment-level operator auth before exposing it.
|
||||
pub async fn list_communities_owned_by(
|
||||
&self,
|
||||
owner_pubkey: &str,
|
||||
) -> Result<Vec<OwnedCommunityRecord>> {
|
||||
let owner_pubkey = owner_pubkey.to_ascii_lowercase();
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT c.id, c.host, c.created_at
|
||||
FROM communities c
|
||||
JOIN relay_members rm ON rm.community_id = c.id
|
||||
WHERE rm.pubkey = $1
|
||||
AND rm.role = 'owner'
|
||||
ORDER BY c.created_at ASC, c.host ASC
|
||||
"#,
|
||||
)
|
||||
.bind(owner_pubkey)
|
||||
.fetch_all(&self.pool)
|
||||
.await?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|row| {
|
||||
let id: Uuid = row.try_get("id")?;
|
||||
let host: String = row.try_get("host")?;
|
||||
let created_at: DateTime<Utc> = row.try_get("created_at")?;
|
||||
Ok(OwnedCommunityRecord {
|
||||
id: CommunityId::from_uuid(id),
|
||||
host,
|
||||
created_at,
|
||||
})
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Returns the normalized host mapped to a community id, if the community
|
||||
/// exists.
|
||||
///
|
||||
@@ -2864,6 +2912,35 @@ mod tests {
|
||||
assert_eq!(found.id, CommunityId::from_uuid(id));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[ignore = "requires Postgres"]
|
||||
async fn list_communities_owned_by_returns_only_owner_rows() {
|
||||
let db = setup_db().await;
|
||||
let community_a = CommunityId::from_uuid(make_community(&db.pool).await);
|
||||
let community_b = CommunityId::from_uuid(make_community(&db.pool).await);
|
||||
let community_c = CommunityId::from_uuid(make_community(&db.pool).await);
|
||||
let owner = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
|
||||
let other = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
|
||||
|
||||
db.bootstrap_owner(community_a, owner)
|
||||
.await
|
||||
.expect("owner A");
|
||||
db.bootstrap_owner(community_b, other)
|
||||
.await
|
||||
.expect("other owner B");
|
||||
db.add_relay_member(community_c, owner, "admin", None)
|
||||
.await
|
||||
.expect("admin C");
|
||||
|
||||
let owned = db
|
||||
.list_communities_owned_by(owner)
|
||||
.await
|
||||
.expect("list owned communities");
|
||||
|
||||
assert_eq!(owned.len(), 1);
|
||||
assert_eq!(owned[0].id, community_a);
|
||||
}
|
||||
|
||||
async fn insert_channel(pool: &PgPool, community_id: Uuid, channel_id: Uuid) {
|
||||
let creator: Vec<u8> = vec![0u8; 32];
|
||||
sqlx::query(
|
||||
|
||||
@@ -25,7 +25,7 @@ use super::{api_error, internal_error, not_found};
|
||||
///
|
||||
/// Returns the authenticated public key and an event ID for replay detection.
|
||||
/// For X-Pubkey dev mode, the event ID is a zero hash (no replay concern).
|
||||
fn verify_bridge_auth(
|
||||
pub(crate) fn verify_bridge_auth(
|
||||
headers: &HeaderMap,
|
||||
method: &str,
|
||||
url: &str,
|
||||
@@ -76,7 +76,7 @@ fn verify_bridge_auth(
|
||||
/// `AppState`, not process-local memory. Any Redis/guard error fails closed:
|
||||
/// without the shared `SET NX EX` proof, a stateless worker cannot admit the
|
||||
/// NIP-98 request safely.
|
||||
async fn check_nip98_replay(
|
||||
pub(crate) async fn check_nip98_replay(
|
||||
state: &AppState,
|
||||
tenant: &TenantContext,
|
||||
event_id_bytes: [u8; 32],
|
||||
@@ -135,7 +135,11 @@ async fn check_nip98_replay_with_guard(
|
||||
/// pass and the relay would proceed against the wrong tenant's auth context),
|
||||
/// and (b) reject every legitimate request whose community host isn't the
|
||||
/// single configured one. Substituting `tenant.host()` closes both directions.
|
||||
fn nip98_expected_url(config_relay_url: &str, tenant: &TenantContext, path: &str) -> String {
|
||||
pub(crate) fn nip98_expected_url(
|
||||
config_relay_url: &str,
|
||||
tenant: &TenantContext,
|
||||
path: &str,
|
||||
) -> String {
|
||||
let scheme = if config_relay_url.trim_start().starts_with("wss://") {
|
||||
"https"
|
||||
} else {
|
||||
@@ -560,6 +564,10 @@ pub async fn submit_event(
|
||||
check_nip98_replay(&state, &tenant, event_id_bytes).await?;
|
||||
let pubkey_bytes = pubkey.to_bytes().to_vec();
|
||||
|
||||
let event: nostr::Event = serde_json::from_slice(&body)
|
||||
.map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid event JSON: {e}")))?;
|
||||
let kind_u32 = buzz_core::kind::event_kind_u32(&event);
|
||||
|
||||
// Enforce relay membership (with NIP-OA fallback via x-auth-tag header).
|
||||
let auth_tag = headers.get("x-auth-tag").and_then(|v| v.to_str().ok());
|
||||
super::relay_members::enforce_relay_membership(
|
||||
@@ -570,16 +578,12 @@ pub async fn submit_event(
|
||||
)
|
||||
.await?;
|
||||
|
||||
let event: nostr::Event = serde_json::from_slice(&body)
|
||||
.map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid event JSON: {e}")))?;
|
||||
|
||||
// Mesh signaling kinds (24620 status report, 24621 connect request) are
|
||||
// ephemeral and deliberately absent from ingest_event's per-kind allowlist.
|
||||
// The desktop's Rust coordinator publishes them via this bridge, so route
|
||||
// them to the mesh handlers — the HTTP twin of the WS door's special-casing
|
||||
// in handlers::event. Membership was enforced above; the handlers re-check
|
||||
// it fail-closed.
|
||||
let kind_u32 = buzz_core::kind::event_kind_u32(&event);
|
||||
if kind_u32 == buzz_core::kind::KIND_MESH_STATUS_REPORT
|
||||
|| kind_u32 == buzz_core::kind::KIND_MESH_CONNECT_REQUEST
|
||||
{
|
||||
|
||||
@@ -5,6 +5,7 @@ pub mod events;
|
||||
pub mod git;
|
||||
pub mod media;
|
||||
pub mod nip05;
|
||||
pub mod operator;
|
||||
|
||||
// Re-export imeta helpers used by ingest pipeline.
|
||||
pub use crate::handlers::imeta::{validate_imeta_tags, verify_imeta_blobs};
|
||||
|
||||
@@ -0,0 +1,210 @@
|
||||
//! Deployment-operator HTTP APIs.
|
||||
//!
|
||||
//! These routes are outside the Nostr event data plane. They still use NIP-98
|
||||
//! request signing and replay protection, but they do not run through event
|
||||
//! ingest, relay membership, channel scoping, storage, or fan-out.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::{
|
||||
extract::{Query, RawQuery, State},
|
||||
http::{HeaderMap, StatusCode},
|
||||
response::Json,
|
||||
};
|
||||
use serde::Deserialize;
|
||||
use serde_json::Value;
|
||||
|
||||
use crate::handlers::community_provisioning::{
|
||||
normalize_candidate_host, validate_pubkey_hex, ProvisionCommunityRequest,
|
||||
};
|
||||
use crate::state::AppState;
|
||||
|
||||
use super::{api_error, bridge, internal_error};
|
||||
|
||||
/// Query parameters for `GET /operator/communities`.
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct ListCommunitiesQuery {
|
||||
owner_pubkey: String,
|
||||
}
|
||||
|
||||
/// Query parameters for `GET /operator/communities/availability`.
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct CommunityAvailabilityQuery {
|
||||
host: String,
|
||||
}
|
||||
|
||||
/// Shared operator auth prelude: bind an ingress host for NIP-98 URL/replay
|
||||
/// scoping, verify the signed request, then gate on `RELAY_OPERATOR_PUBKEYS`.
|
||||
async fn authorize_operator_request(
|
||||
state: &Arc<AppState>,
|
||||
headers: &HeaderMap,
|
||||
method: &str,
|
||||
path: &str,
|
||||
raw_query: Option<&str>,
|
||||
body: Option<&[u8]>,
|
||||
) -> Result<nostr::PublicKey, (StatusCode, Json<Value>)> {
|
||||
// Bind to an existing ingress community only for NIP-98 URL/replay scoping.
|
||||
// This is not a tenant data-plane operation, so do not run relay-membership
|
||||
// checks and do not route through event ingest.
|
||||
let raw_host = headers
|
||||
.get(axum::http::header::HOST)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.unwrap_or("");
|
||||
let tenant = crate::tenant::bind_community(&state.db, raw_host)
|
||||
.await
|
||||
.map_err(|_| {
|
||||
api_error(
|
||||
StatusCode::NOT_FOUND,
|
||||
"relay: no community is configured for this host",
|
||||
)
|
||||
})?;
|
||||
|
||||
let path_with_query = match raw_query {
|
||||
Some(q) if !q.is_empty() => format!("{path}?{q}"),
|
||||
_ => path.to_string(),
|
||||
};
|
||||
let url = bridge::nip98_expected_url(&state.config.relay_url, &tenant, &path_with_query);
|
||||
let (pubkey, event_id_bytes) = bridge::verify_bridge_auth(
|
||||
headers, method, &url, body,
|
||||
true, // operator endpoints always require NIP-98; no X-Pubkey dev fallback
|
||||
)?;
|
||||
bridge::check_nip98_replay(state, &tenant, event_id_bytes).await?;
|
||||
|
||||
let pubkey_hex = pubkey.to_hex();
|
||||
if !state
|
||||
.config
|
||||
.relay_operator_pubkeys
|
||||
.iter()
|
||||
.any(|pk| pk == &pubkey_hex)
|
||||
{
|
||||
return Err(api_error(
|
||||
StatusCode::FORBIDDEN,
|
||||
"actor not authorized: not a relay operator",
|
||||
));
|
||||
}
|
||||
|
||||
Ok(pubkey)
|
||||
}
|
||||
|
||||
/// Provision or converge a community host.
|
||||
///
|
||||
/// `POST /operator/communities`, NIP-98 signed by a pubkey in
|
||||
/// `RELAY_OPERATOR_PUBKEYS`, body:
|
||||
///
|
||||
/// ```json
|
||||
/// { "host": "acme.communities.buzz.xyz", "initial_owner_pubkey": "<hex>" }
|
||||
/// ```
|
||||
///
|
||||
/// The request is authenticated against the host it arrives on (so NIP-98 `u`
|
||||
/// still binds to the request authority) but it intentionally does not require
|
||||
/// relay membership in that host's community. The operator allowlist is the
|
||||
/// authority for this deployment-root control-plane surface.
|
||||
pub async fn provision_community(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
body: axum::body::Bytes,
|
||||
) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
|
||||
let pubkey = authorize_operator_request(
|
||||
&state,
|
||||
&headers,
|
||||
"POST",
|
||||
"/operator/communities",
|
||||
None,
|
||||
Some(&body),
|
||||
)
|
||||
.await?;
|
||||
|
||||
let request: ProvisionCommunityRequest = serde_json::from_slice(&body).map_err(|e| {
|
||||
api_error(
|
||||
StatusCode::BAD_REQUEST,
|
||||
&format!("invalid provision-community JSON: {e}"),
|
||||
)
|
||||
})?;
|
||||
|
||||
match crate::handlers::community_provisioning::provision_community(&state, &pubkey, request)
|
||||
.await
|
||||
{
|
||||
Ok(response) => Ok(Json(serde_json::to_value(response).map_err(|e| {
|
||||
tracing::error!("failed to serialize provision-community response: {e}");
|
||||
internal_error("operator provision response serialization failed")
|
||||
})?)),
|
||||
Err(msg) if msg.starts_with("actor not authorized") => {
|
||||
Err(api_error(StatusCode::FORBIDDEN, &msg))
|
||||
}
|
||||
Err(msg) => Err(api_error(StatusCode::BAD_REQUEST, &msg)),
|
||||
}
|
||||
}
|
||||
|
||||
/// List communities where a pubkey currently holds the `owner` role.
|
||||
pub async fn list_owned_communities(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
RawQuery(raw_query): RawQuery,
|
||||
Query(query): Query<ListCommunitiesQuery>,
|
||||
) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
|
||||
authorize_operator_request(
|
||||
&state,
|
||||
&headers,
|
||||
"GET",
|
||||
"/operator/communities",
|
||||
raw_query.as_deref(),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let owner_pubkey = validate_pubkey_hex(&query.owner_pubkey).ok_or_else(|| {
|
||||
api_error(
|
||||
StatusCode::BAD_REQUEST,
|
||||
"invalid owner_pubkey: expected 64-char hex pubkey",
|
||||
)
|
||||
})?;
|
||||
|
||||
let rows = state
|
||||
.db
|
||||
.list_communities_owned_by(&owner_pubkey)
|
||||
.await
|
||||
.map_err(|e| internal_error(&format!("list owned communities: {e}")))?;
|
||||
|
||||
Ok(Json(serde_json::json!({
|
||||
"owner_pubkey": owner_pubkey,
|
||||
"communities": rows.into_iter().map(|row| serde_json::json!({
|
||||
"community_id": row.id.to_string(),
|
||||
"host": row.host,
|
||||
"created_at": row.created_at,
|
||||
})).collect::<Vec<_>>(),
|
||||
})))
|
||||
}
|
||||
|
||||
/// Check whether a community host is available, returning the relay-canonical
|
||||
/// normalized authority used by create.
|
||||
pub async fn community_availability(
|
||||
State(state): State<Arc<AppState>>,
|
||||
headers: HeaderMap,
|
||||
RawQuery(raw_query): RawQuery,
|
||||
Query(query): Query<CommunityAvailabilityQuery>,
|
||||
) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
|
||||
authorize_operator_request(
|
||||
&state,
|
||||
&headers,
|
||||
"GET",
|
||||
"/operator/communities/availability",
|
||||
raw_query.as_deref(),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let normalized_host = normalize_candidate_host(&query.host)
|
||||
.map_err(|msg| api_error(StatusCode::BAD_REQUEST, &msg))?;
|
||||
let existing = state
|
||||
.db
|
||||
.lookup_community_by_host(&normalized_host)
|
||||
.await
|
||||
.map_err(|e| internal_error(&format!("check community availability: {e}")))?;
|
||||
|
||||
Ok(Json(serde_json::json!({
|
||||
"host": query.host,
|
||||
"normalized_host": normalized_host,
|
||||
"available": existing.is_none(),
|
||||
"community_id": existing.map(|record| record.id.to_string()),
|
||||
})))
|
||||
}
|
||||
@@ -98,6 +98,19 @@ pub struct Config {
|
||||
/// with the `owner` role on first startup.
|
||||
pub relay_owner_pubkey: Option<String>,
|
||||
|
||||
/// Deployment-level relay operator pubkeys allowed to use the
|
||||
/// `POST /operator/communities` provisioning endpoint.
|
||||
///
|
||||
/// Unlike `relay_owner_pubkey` (a role *within* the deployment community),
|
||||
/// operators span tenants: they may create new communities and rotate owners
|
||||
/// via the operator endpoint, but hold no implicit tenant membership row.
|
||||
/// Empty (the default) disables community provisioning entirely — fail closed.
|
||||
///
|
||||
/// Set via `RELAY_OPERATOR_PUBKEYS` as a comma-separated list of 64-char
|
||||
/// hex pubkeys. Invalid entries are rejected at startup (config error), not
|
||||
/// skipped — a typo must not silently disable an operator.
|
||||
pub relay_operator_pubkeys: Vec<String>,
|
||||
|
||||
/// Allow NIP-OA owner attestation for relay membership.
|
||||
///
|
||||
/// When `true` and `require_relay_membership` is also `true`, agents
|
||||
@@ -184,7 +197,7 @@ impl Config {
|
||||
let bind_addr = parse_bind_addr(&bind_addr_raw)?;
|
||||
|
||||
let database_url = std::env::var("DATABASE_URL")
|
||||
.unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string());
|
||||
.unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string()); // sadscan:disable np.postgres.1
|
||||
|
||||
let redis_url =
|
||||
std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://localhost:6379".to_string());
|
||||
@@ -260,6 +273,34 @@ impl Config {
|
||||
}
|
||||
});
|
||||
|
||||
// Note: intentionally not prefixed with BUZZ_ — same relay-identity
|
||||
// config family as RELAY_OWNER_PUBKEY. Comma-separated 64-char hex
|
||||
// pubkeys. Unlike RELAY_OWNER_PUBKEY (warn-and-ignore), an invalid
|
||||
// entry here is a hard config error: silently dropping an operator
|
||||
// pubkey would silently disable provisioning for that operator.
|
||||
let relay_operator_pubkeys = match std::env::var("RELAY_OPERATOR_PUBKEYS") {
|
||||
Ok(raw) => {
|
||||
let mut pubkeys = Vec::new();
|
||||
for entry in raw.split(',') {
|
||||
let entry = entry.trim().to_lowercase();
|
||||
if entry.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let valid = entry.len() == 64 && entry.chars().all(|c| c.is_ascii_hexdigit());
|
||||
if !valid {
|
||||
return Err(ConfigError::InvalidValue(format!(
|
||||
"RELAY_OPERATOR_PUBKEYS entry is not a valid 64-char hex pubkey: {entry:?}"
|
||||
)));
|
||||
}
|
||||
if !pubkeys.contains(&entry) {
|
||||
pubkeys.push(entry);
|
||||
}
|
||||
}
|
||||
pubkeys
|
||||
}
|
||||
Err(_) => Vec::new(),
|
||||
};
|
||||
|
||||
let auth = buzz_auth::AuthConfig::default();
|
||||
|
||||
if !require_auth_token {
|
||||
@@ -428,6 +469,7 @@ impl Config {
|
||||
require_relay_membership,
|
||||
huddle_audio_available,
|
||||
relay_owner_pubkey,
|
||||
relay_operator_pubkeys,
|
||||
allow_nip_oa_auth,
|
||||
media,
|
||||
media_max_concurrent_uploads,
|
||||
@@ -477,6 +519,10 @@ mod tests {
|
||||
config.relay_owner_pubkey.is_none(),
|
||||
"relay_owner_pubkey should default to None"
|
||||
);
|
||||
assert!(
|
||||
config.relay_operator_pubkeys.is_empty(),
|
||||
"relay_operator_pubkeys should default empty (provisioning disabled)"
|
||||
);
|
||||
assert!(
|
||||
!config.allow_nip_oa_auth,
|
||||
"allow_nip_oa_auth should default to false"
|
||||
@@ -487,6 +533,38 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn relay_operator_pubkeys_parse_dedupe_and_normalize() {
|
||||
let _guard = ENV_MUTEX.lock().unwrap();
|
||||
std::env::set_var(
|
||||
"RELAY_OPERATOR_PUBKEYS",
|
||||
"AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA,bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb,aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
|
||||
);
|
||||
let config = Config::from_env().expect("config");
|
||||
std::env::remove_var("RELAY_OPERATOR_PUBKEYS");
|
||||
|
||||
assert_eq!(
|
||||
config.relay_operator_pubkeys,
|
||||
vec![
|
||||
"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string(),
|
||||
"bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".to_string(),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn relay_operator_pubkeys_invalid_entry_is_error() {
|
||||
let _guard = ENV_MUTEX.lock().unwrap();
|
||||
std::env::set_var("RELAY_OPERATOR_PUBKEYS", "not-a-pubkey");
|
||||
let result = Config::from_env();
|
||||
std::env::remove_var("RELAY_OPERATOR_PUBKEYS");
|
||||
|
||||
assert!(matches!(
|
||||
result,
|
||||
Err(ConfigError::InvalidValue(ref msg)) if msg.contains("RELAY_OPERATOR_PUBKEYS")
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn huddle_audio_available_can_be_disabled_for_horizontal_scaling() {
|
||||
let _guard = ENV_MUTEX.lock().unwrap();
|
||||
|
||||
@@ -0,0 +1,304 @@
|
||||
//! Relay-operator community provisioning HTTP handler support.
|
||||
//!
|
||||
//! ## Authorization: operator, not owner
|
||||
//!
|
||||
//! Every other admin surface in this relay is community-scoped — the sender's
|
||||
//! role is looked up in `relay_members (community_id, pubkey)` for the
|
||||
//! host-resolved tenant. Community *creation* cannot work that way: its effect
|
||||
//! is the creation of tenancy itself, so the authorizing identity must sit
|
||||
//! above tenants. The gate here is the deployment-level
|
||||
//! `RELAY_OPERATOR_PUBKEYS` allowlist (see `Config::relay_operator_pubkeys`).
|
||||
//! An empty allowlist (the default) disables provisioning entirely.
|
||||
//!
|
||||
//! The public surface is `POST /operator/communities`, authenticated by NIP-98
|
||||
//! and gated by the deployment-level `RELAY_OPERATOR_PUBKEYS` allowlist. The
|
||||
//! endpoint is intentionally outside the Nostr event ingest data plane: no
|
||||
//! relay-membership bypass, no special event kind, no storage or fan-out.
|
||||
//!
|
||||
//! ## Request shape
|
||||
//!
|
||||
//! ```json
|
||||
//! { "host": "acme.communities.buzz.xyz", "initial_owner_pubkey": "<hex>" }
|
||||
//! ```
|
||||
//!
|
||||
//! `initial_owner_pubkey` is optional. When present for an existing community,
|
||||
//! it rotates that community owner through the same bootstrap path used by
|
||||
//! `RELAY_OWNER_PUBKEY`; relay operators are deployment-root authorities.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tracing::info;
|
||||
|
||||
use buzz_core::tenant::normalize_host;
|
||||
|
||||
use crate::state::AppState;
|
||||
|
||||
/// Maximum accepted host length. Hostnames cap at 253 octets; leave
|
||||
/// headroom for a `:port` suffix.
|
||||
const MAX_HOST_LEN: usize = 260;
|
||||
|
||||
/// JSON body for `POST /operator/communities`.
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct ProvisionCommunityRequest {
|
||||
/// Normalized authority for the community to ensure/create.
|
||||
pub host: String,
|
||||
/// Optional initial owner pubkey. When set on an existing community this
|
||||
/// rotates the owner through the same bootstrap path used at startup.
|
||||
#[serde(default)]
|
||||
pub initial_owner_pubkey: Option<String>,
|
||||
}
|
||||
|
||||
/// JSON response from `POST /operator/communities`.
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct ProvisionCommunityResponse {
|
||||
/// UUID of the ensured/created community.
|
||||
pub community_id: String,
|
||||
/// Canonical host stored on the community row.
|
||||
pub host: String,
|
||||
/// `created` when the host row was inserted, `existed` when it was already
|
||||
/// present and the request converged idempotently.
|
||||
pub status: &'static str,
|
||||
/// Echoes the validated owner pubkey when an owner bootstrap/rotation ran.
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub owner_pubkey: Option<String>,
|
||||
}
|
||||
|
||||
pub(crate) fn validate_pubkey_hex(value: &str) -> Option<String> {
|
||||
let normalized = value.to_ascii_lowercase();
|
||||
(normalized.len() == 64 && normalized.chars().all(|c| c.is_ascii_hexdigit()))
|
||||
.then_some(normalized)
|
||||
}
|
||||
|
||||
/// Validate a normalized host value for a community.
|
||||
///
|
||||
/// The host must already be in normalized shape (`normalize_host` is a
|
||||
/// no-op on it): lowercase, no default port, no trailing dot. Requiring the
|
||||
/// caller to send the normalized form keeps the stored `communities.host`
|
||||
/// value byte-identical to what request-time host resolution will look up.
|
||||
fn validate_host(host: &str) -> Result<(), String> {
|
||||
if host.is_empty() {
|
||||
return Err("host is empty".to_string());
|
||||
}
|
||||
if host.len() > MAX_HOST_LEN {
|
||||
return Err(format!(
|
||||
"host too long: {} bytes (max {MAX_HOST_LEN})",
|
||||
host.len()
|
||||
));
|
||||
}
|
||||
if host.chars().any(|c| c.is_control() || c.is_whitespace()) {
|
||||
return Err("host contains invalid characters".to_string());
|
||||
}
|
||||
if host.contains('/') || host.contains('?') || host.contains('#') || host.contains('@') {
|
||||
return Err(
|
||||
"host must be a bare authority (no scheme, path, query, or userinfo)".to_string(),
|
||||
);
|
||||
}
|
||||
if normalize_host(host) != host {
|
||||
return Err(format!(
|
||||
"host is not normalized: expected {:?}",
|
||||
normalize_host(host)
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Normalize and validate a host supplied to read-only operator endpoints.
|
||||
///
|
||||
/// Unlike create, availability checks may accept non-canonical but normalizable
|
||||
/// authority values (uppercase host, trailing dot, default port) so kgoose can
|
||||
/// ask the relay for the canonical spelling before creating. Schemes, paths,
|
||||
/// userinfo, whitespace/control characters, and oversized values are still
|
||||
/// rejected.
|
||||
pub(crate) fn normalize_candidate_host(host: &str) -> Result<String, String> {
|
||||
if host.is_empty() {
|
||||
return Err("host is empty".to_string());
|
||||
}
|
||||
if host.len() > MAX_HOST_LEN {
|
||||
return Err(format!(
|
||||
"host too long: {} bytes (max {MAX_HOST_LEN})",
|
||||
host.len()
|
||||
));
|
||||
}
|
||||
if host.chars().any(|c| c.is_control() || c.is_whitespace()) {
|
||||
return Err("host contains invalid characters".to_string());
|
||||
}
|
||||
if host.contains('/') || host.contains('?') || host.contains('#') || host.contains('@') {
|
||||
return Err(
|
||||
"host must be a bare authority (no scheme, path, query, or userinfo)".to_string(),
|
||||
);
|
||||
}
|
||||
|
||||
let normalized = normalize_host(host);
|
||||
validate_host(&normalized)?;
|
||||
Ok(normalized)
|
||||
}
|
||||
|
||||
/// Validate and execute a relay-operator community provisioning request.
|
||||
///
|
||||
/// The caller is an HTTP operator endpoint, not the Nostr event ingest path.
|
||||
/// That keeps the tenant data-plane fences unchanged: no relay-membership
|
||||
/// bypass, no special event kind, no command routed ahead of moderation/write
|
||||
/// blocks. The endpoint authenticates its NIP-98 signer first, then passes the
|
||||
/// signer here for the deployment-level `RELAY_OPERATOR_PUBKEYS` allowlist.
|
||||
///
|
||||
/// Idempotency and owner semantics: the request is idempotent on the host row
|
||||
/// (re-sending it never duplicates a community). When `initial_owner_pubkey` is
|
||||
/// present, the owner is (re)bootstrapped via [`buzz_db::Db::bootstrap_owner`]
|
||||
/// even if the community already existed — any previous owner is demoted to
|
||||
/// admin, exactly like rotating `RELAY_OWNER_PUBKEY` for the deployment
|
||||
/// community. This makes a retry after a partial failure (row created, owner
|
||||
/// bootstrap crashed) converge, at the cost that an operator-signed request can
|
||||
/// rotate an existing community's owner. The operator allowlist is therefore
|
||||
/// documented as deployment-root authority, not create-only authority.
|
||||
pub async fn provision_community(
|
||||
state: &Arc<AppState>,
|
||||
operator_pubkey: &nostr::PublicKey,
|
||||
request: ProvisionCommunityRequest,
|
||||
) -> Result<ProvisionCommunityResponse, String> {
|
||||
let operator_hex = operator_pubkey.to_hex();
|
||||
|
||||
// Operator gate. Deliberately NOT a relay_members lookup: provisioning
|
||||
// authority spans tenants and lives in deployment config only. Empty
|
||||
// allowlist → everyone is rejected (fail closed).
|
||||
if !state
|
||||
.config
|
||||
.relay_operator_pubkeys
|
||||
.iter()
|
||||
.any(|pk| pk == &operator_hex)
|
||||
{
|
||||
return Err("actor not authorized: not a relay operator".to_string());
|
||||
}
|
||||
|
||||
validate_host(&request.host)?;
|
||||
|
||||
let initial_owner = request
|
||||
.initial_owner_pubkey
|
||||
.as_deref()
|
||||
.map(|value| {
|
||||
validate_pubkey_hex(value).ok_or_else(|| {
|
||||
"invalid initial_owner_pubkey: expected 64-char hex pubkey".to_string()
|
||||
})
|
||||
})
|
||||
.transpose()?;
|
||||
|
||||
let existed = state
|
||||
.db
|
||||
.lookup_community_by_host(&request.host)
|
||||
.await
|
||||
.map_err(|e| format!("database error: {e}"))?
|
||||
.is_some();
|
||||
|
||||
// Same idempotent upsert as the startup seed — creating a community is an
|
||||
// INSERT, never DDL (docs/multi-tenant-relay.md §System Model).
|
||||
let record = state
|
||||
.db
|
||||
.ensure_configured_community(&request.host)
|
||||
.await
|
||||
.map_err(|e| format!("failed to create community: {e}"))?;
|
||||
|
||||
if let Some(owner_hex) = &initial_owner {
|
||||
state
|
||||
.db
|
||||
.bootstrap_owner(record.id, owner_hex)
|
||||
.await
|
||||
.map_err(|e| format!("community provisioned but owner bootstrap failed: {e}"))?;
|
||||
}
|
||||
|
||||
info!(
|
||||
operator = %operator_hex,
|
||||
community = %record.id,
|
||||
host = %record.host,
|
||||
owner = initial_owner.as_deref().unwrap_or("<none>"),
|
||||
existed,
|
||||
"community provisioned via operator endpoint"
|
||||
);
|
||||
|
||||
Ok(ProvisionCommunityResponse {
|
||||
community_id: record.id.to_string(),
|
||||
host: record.host,
|
||||
status: if existed { "existed" } else { "created" },
|
||||
owner_pubkey: initial_owner,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn host_valid_bare_domain() {
|
||||
assert!(validate_host("acme.communities.buzz.xyz").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_valid_with_port() {
|
||||
assert!(validate_host("localhost:3000").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_empty() {
|
||||
assert!(validate_host("").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_uppercase() {
|
||||
assert!(validate_host("Acme.example").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_default_port() {
|
||||
assert!(validate_host("acme.example:443").is_err());
|
||||
assert!(validate_host("acme.example:80").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_trailing_dot() {
|
||||
assert!(validate_host("acme.example.").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_scheme_path_userinfo() {
|
||||
assert!(validate_host("wss://acme.example").is_err());
|
||||
assert!(validate_host("acme.example/path").is_err());
|
||||
assert!(validate_host("user@acme.example").is_err());
|
||||
assert!(validate_host("acme.example?x=1").is_err());
|
||||
assert!(validate_host("acme.example#frag").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_whitespace_and_control() {
|
||||
assert!(validate_host("acme .example").is_err());
|
||||
assert!(validate_host("acme\n.example").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_rejects_oversized() {
|
||||
let long = format!("{}.example", "a".repeat(260));
|
||||
assert!(validate_host(&long).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn host_accepts_ipv6_bracket_literal() {
|
||||
assert!(validate_host("[::1]:3000").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn candidate_host_normalizes_safe_variants() {
|
||||
assert_eq!(
|
||||
normalize_candidate_host("Acme.Example:443").unwrap(),
|
||||
"acme.example"
|
||||
);
|
||||
assert_eq!(
|
||||
normalize_candidate_host("acme.example.").unwrap(),
|
||||
"acme.example"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn candidate_host_rejects_non_authorities() {
|
||||
assert!(normalize_candidate_host("https://acme.example").is_err());
|
||||
assert!(normalize_candidate_host("acme.example/path").is_err());
|
||||
assert!(normalize_candidate_host("acme .example").is_err());
|
||||
}
|
||||
}
|
||||
@@ -200,7 +200,7 @@ fn required_scope_for_kind(kind: u32, event: &Event) -> Result<Scope, &'static s
|
||||
Ok(Scope::AdminChannels)
|
||||
}
|
||||
// NIP-43: relay membership admin commands (9030–9032) + Buzz
|
||||
// workspace-profile command (9033)
|
||||
// workspace-profile command (9033).
|
||||
k if k == RELAY_ADMIN_ADD_MEMBER
|
||||
|| k == RELAY_ADMIN_REMOVE_MEMBER
|
||||
|| k == RELAY_ADMIN_CHANGE_ROLE
|
||||
|
||||
@@ -4,6 +4,8 @@ pub mod auth;
|
||||
pub mod close;
|
||||
/// Command executor — transactional processing for command kinds.
|
||||
pub mod command_executor;
|
||||
/// Relay-operator community provisioning HTTP support.
|
||||
pub mod community_provisioning;
|
||||
/// NIP-45 COUNT handler.
|
||||
pub mod count;
|
||||
/// EVENT handler — WS dispatcher → ingest pipeline → fan-out.
|
||||
|
||||
@@ -60,6 +60,14 @@ pub fn build_router(state: Arc<AppState>) -> Router {
|
||||
.route("/events", post(api::bridge::submit_event))
|
||||
.route("/query", post(api::bridge::query_events))
|
||||
.route("/count", post(api::bridge::count_events))
|
||||
.route(
|
||||
"/operator/communities",
|
||||
get(api::operator::list_owned_communities).post(api::operator::provision_community),
|
||||
)
|
||||
.route(
|
||||
"/operator/communities/availability",
|
||||
get(api::operator::community_availability),
|
||||
)
|
||||
// Moderation queue reads (NIP-98 auth + mod-authz gate, L6)
|
||||
.route("/moderation/reports", get(api::bridge::moderation_reports))
|
||||
.route("/moderation/audit", get(api::bridge::moderation_audit))
|
||||
@@ -96,6 +104,7 @@ pub fn build_router(state: Arc<AppState>) -> Router {
|
||||
// Reserved API prefixes must 404 normally, not serve index.html.
|
||||
let reserved = path.starts_with("/api/")
|
||||
|| path.starts_with("/media/")
|
||||
|| path.starts_with("/operator/")
|
||||
|| path.starts_with("/git/")
|
||||
|| path.starts_with("/internal/")
|
||||
|| path.starts_with("/.well-known/")
|
||||
|
||||
Reference in New Issue
Block a user