diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 3fe9aca05..0d1924161 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -2340,17 +2340,34 @@ impl Db { created_by: &[u8], ttl_seconds: Option, ) -> Result { - channel::create_channel( - self.pg_pool()?, - community_id, - name, - channel_type, - visibility, - description, - created_by, - ttl_seconds, - ) - .await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::create_channel( + pool, + community_id, + name, + channel_type, + visibility, + description, + created_by, + ttl_seconds, + ) + .await + } + DbBackend::Postgres => { + channel::create_channel( + self.pg_pool()?, + community_id, + name, + channel_type, + visibility, + description, + created_by, + ttl_seconds, + ) + .await + } + } } /// Creates a channel with a client-supplied UUID. @@ -2416,7 +2433,12 @@ impl Db { community_id: CommunityId, channel_id: Uuid, ) -> Result> { - channel::get_canvas(self.pg_pool()?, community_id, channel_id).await + match &self.backend { + DbBackend::SQLite(pool) => sqlite::get_canvas(pool, community_id, channel_id).await, + DbBackend::Postgres => { + channel::get_canvas(self.pg_pool()?, community_id, channel_id).await + } + } } /// Sets or clears the canvas content for a channel. @@ -2426,7 +2448,14 @@ impl Db { channel_id: Uuid, canvas: Option<&str>, ) -> Result<()> { - channel::set_canvas(self.pg_pool()?, community_id, channel_id, canvas).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::set_canvas(pool, community_id, channel_id, canvas).await + } + DbBackend::Postgres => { + channel::set_canvas(self.pg_pool()?, community_id, channel_id, canvas).await + } + } } /// Adds a member to a channel. @@ -2461,14 +2490,21 @@ impl Db { pubkey: &[u8], actor_pubkey: &[u8], ) -> Result<()> { - channel::remove_member( - self.pg_pool()?, - community_id, - channel_id, - pubkey, - actor_pubkey, - ) - .await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::remove_member(pool, community_id, channel_id, pubkey, actor_pubkey).await + } + DbBackend::Postgres => { + channel::remove_member( + self.pg_pool()?, + community_id, + channel_id, + pubkey, + actor_pubkey, + ) + .await + } + } } /// Returns `true` if the pubkey is an active member. @@ -2496,7 +2532,14 @@ impl Db { channel_ids: &[Uuid], pubkeys: &[Vec], ) -> Result)>> { - channel::membership_pairs(self.pg_pool()?, community_id, channel_ids, pubkeys).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::membership_pairs(pool, community_id, channel_ids, pubkeys).await + } + DbBackend::Postgres => { + channel::membership_pairs(self.pg_pool()?, community_id, channel_ids, pubkeys).await + } + } } /// Returns all active members of a channel. @@ -2519,7 +2562,14 @@ impl Db { community_id: CommunityId, channel_ids: &[Uuid], ) -> Result> { - channel::get_members_bulk(self.pg_pool()?, community_id, channel_ids).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::get_members_bulk(pool, community_id, channel_ids).await + } + DbBackend::Postgres => { + channel::get_members_bulk(self.pg_pool()?, community_id, channel_ids).await + } + } } /// Get all channel IDs accessible to a pubkey. @@ -2544,7 +2594,12 @@ impl Db { community_id: CommunityId, visibility: Option<&str>, ) -> Result> { - channel::list_channels(self.pg_pool()?, community_id, visibility).await + match &self.backend { + DbBackend::SQLite(pool) => sqlite::list_channels(pool, community_id, visibility).await, + DbBackend::Postgres => { + channel::list_channels(self.pg_pool()?, community_id, visibility).await + } + } } /// Returns full channel records for all channels a user can access. @@ -2555,14 +2610,28 @@ impl Db { visibility_filter: Option<&str>, member_only: Option, ) -> Result> { - channel::get_accessible_channels( - self.pg_pool()?, - community_id, - pubkey, - visibility_filter, - member_only, - ) - .await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::get_accessible_channels( + pool, + community_id, + pubkey, + visibility_filter, + member_only, + ) + .await + } + DbBackend::Postgres => { + channel::get_accessible_channels( + self.pg_pool()?, + community_id, + pubkey, + visibility_filter, + member_only, + ) + .await + } + } } /// Returns all bot-role members with their aggregated channel names in one community. @@ -2570,7 +2639,10 @@ impl Db { &self, community_id: CommunityId, ) -> Result> { - channel::get_bot_members(self.pg_pool()?, community_id).await + match &self.backend { + DbBackend::SQLite(pool) => sqlite::get_bot_members(pool, community_id).await, + DbBackend::Postgres => channel::get_bot_members(self.pg_pool()?, community_id).await, + } } /// Bulk-fetch user records by pubkey. @@ -2579,7 +2651,12 @@ impl Db { community_id: CommunityId, pubkeys: &[Vec], ) -> Result> { - channel::get_users_bulk(self.pg_pool()?, community_id, pubkeys).await + match &self.backend { + DbBackend::SQLite(pool) => sqlite::get_users_bulk(pool, community_id, pubkeys).await, + DbBackend::Postgres => { + channel::get_users_bulk(self.pg_pool()?, community_id, pubkeys).await + } + } } /// Updates a channel's name and/or description. @@ -2589,7 +2666,14 @@ impl Db { channel_id: Uuid, updates: channel::ChannelUpdate, ) -> Result { - channel::update_channel(self.pg_pool()?, community_id, channel_id, updates).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::update_channel(pool, community_id, channel_id, updates).await + } + DbBackend::Postgres => { + channel::update_channel(self.pg_pool()?, community_id, channel_id, updates).await + } + } } /// Sets the topic for a channel. @@ -2600,7 +2684,14 @@ impl Db { topic: &str, set_by: &[u8], ) -> Result<()> { - channel::set_topic(self.pg_pool()?, community_id, channel_id, topic, set_by).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::set_topic(pool, community_id, channel_id, topic, set_by).await + } + DbBackend::Postgres => { + channel::set_topic(self.pg_pool()?, community_id, channel_id, topic, set_by).await + } + } } /// Sets the purpose for a channel. @@ -2611,12 +2702,27 @@ impl Db { purpose: &str, set_by: &[u8], ) -> Result<()> { - channel::set_purpose(self.pg_pool()?, community_id, channel_id, purpose, set_by).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::set_purpose(pool, community_id, channel_id, purpose, set_by).await + } + DbBackend::Postgres => { + channel::set_purpose(self.pg_pool()?, community_id, channel_id, purpose, set_by) + .await + } + } } /// Archives a channel. pub async fn archive_channel(&self, community_id: CommunityId, channel_id: Uuid) -> Result<()> { - channel::archive_channel(self.pg_pool()?, community_id, channel_id).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::archive_channel(pool, community_id, channel_id).await + } + DbBackend::Postgres => { + channel::archive_channel(self.pg_pool()?, community_id, channel_id).await + } + } } /// Unarchives a channel. @@ -2625,7 +2731,14 @@ impl Db { community_id: CommunityId, channel_id: Uuid, ) -> Result<()> { - channel::unarchive_channel(self.pg_pool()?, community_id, channel_id).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::unarchive_channel(pool, community_id, channel_id).await + } + DbBackend::Postgres => { + channel::unarchive_channel(self.pg_pool()?, community_id, channel_id).await + } + } } /// Soft-delete a channel. @@ -2634,7 +2747,14 @@ impl Db { community_id: CommunityId, channel_id: Uuid, ) -> Result { - channel::soft_delete_channel(self.pg_pool()?, community_id, channel_id).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::soft_delete_channel(pool, community_id, channel_id).await + } + DbBackend::Postgres => { + channel::soft_delete_channel(self.pg_pool()?, community_id, channel_id).await + } + } } /// Returns the count of active members in a channel. @@ -2643,7 +2763,14 @@ impl Db { community_id: CommunityId, channel_id: Uuid, ) -> Result { - channel::get_member_count(self.pg_pool()?, community_id, channel_id).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::get_member_count(pool, community_id, channel_id).await + } + DbBackend::Postgres => { + channel::get_member_count(self.pg_pool()?, community_id, channel_id).await + } + } } /// Bulk-fetch member counts for a set of channel IDs. @@ -2652,7 +2779,14 @@ impl Db { community_id: CommunityId, channel_ids: &[Uuid], ) -> Result> { - channel::get_member_counts_bulk(self.pg_pool()?, community_id, channel_ids).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::get_member_counts_bulk(pool, community_id, channel_ids).await + } + DbBackend::Postgres => { + channel::get_member_counts_bulk(self.pg_pool()?, community_id, channel_ids).await + } + } } /// Get the active role of a pubkey in a channel. @@ -2676,7 +2810,10 @@ impl Db { pub async fn reap_expired_ephemeral_channels( &self, ) -> Result> { - channel::reap_expired_ephemeral_channels(self.pg_pool()?).await + match &self.backend { + DbBackend::SQLite(pool) => sqlite::reap_expired_ephemeral_channels(pool).await, + DbBackend::Postgres => channel::reap_expired_ephemeral_channels(self.pg_pool()?).await, + } } /// Query due reminders ready for delivery. @@ -5040,14 +5177,21 @@ impl Db { role: &str, policy_version: Option<&str>, ) -> Result { - relay_members::claim_relay_membership( - self.pg_pool()?, - community, - pubkey, - role, - policy_version, - ) - .await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::claim_relay_membership(pool, community, pubkey, role, policy_version).await + } + DbBackend::Postgres => { + relay_members::claim_relay_membership( + self.pg_pool()?, + community, + pubkey, + role, + policy_version, + ) + .await + } + } } /// Returns whether a member has persisted acceptance evidence for a policy version. @@ -5057,13 +5201,20 @@ impl Db { pubkey: &str, policy_version: &str, ) -> Result { - relay_members::has_join_policy_acceptance( - self.pg_pool()?, - community, - pubkey, - policy_version, - ) - .await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::has_join_policy_acceptance(pool, community, pubkey, policy_version).await + } + DbBackend::Postgres => { + relay_members::has_join_policy_acceptance( + self.pg_pool()?, + community, + pubkey, + policy_version, + ) + .await + } + } } /// Removes a relay member from `community` atomically, refusing to delete the owner. @@ -5072,7 +5223,14 @@ impl Db { community: CommunityId, pubkey: &str, ) -> Result { - relay_members::remove_relay_member(self.pg_pool()?, community, pubkey).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::remove_relay_member(pool, community, pubkey, None).await + } + DbBackend::Postgres => { + relay_members::remove_relay_member(self.pg_pool()?, community, pubkey).await + } + } } /// Removes a relay member from `community` only if their current role matches `expected_role`. @@ -5085,13 +5243,20 @@ impl Db { pubkey: &str, expected_role: &str, ) -> Result { - relay_members::remove_relay_member_if_role( - self.pg_pool()?, - community, - pubkey, - expected_role, - ) - .await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::remove_relay_member(pool, community, pubkey, Some(expected_role)).await + } + DbBackend::Postgres => { + relay_members::remove_relay_member_if_role( + self.pg_pool()?, + community, + pubkey, + expected_role, + ) + .await + } + } } /// Updates the role of an existing relay member in `community`. Returns `true` if updated. @@ -5101,7 +5266,20 @@ impl Db { pubkey: &str, new_role: &str, ) -> Result { - relay_members::update_relay_member_role(self.pg_pool()?, community, pubkey, new_role).await + match &self.backend { + DbBackend::SQLite(pool) => { + sqlite::update_relay_member_role(pool, community, pubkey, new_role).await + } + DbBackend::Postgres => { + relay_members::update_relay_member_role( + self.pg_pool()?, + community, + pubkey, + new_role, + ) + .await + } + } } /// Ensures the owner pubkey exists with role `"owner"` in `community`. Called at startup. diff --git a/crates/buzz-db/src/sqlite.rs b/crates/buzz-db/src/sqlite.rs index 7d3658049..4d3c173f4 100644 --- a/crates/buzz-db/src/sqlite.rs +++ b/crates/buzz-db/src/sqlite.rs @@ -2949,6 +2949,464 @@ pub(crate) async fn get_member_role( .bind(community.as_uuid().to_string()).bind(channel_id.to_string()).bind(pubkey).fetch_optional(pool).await?) } +pub(crate) async fn create_channel( + pool: &SqlitePool, + community: CommunityId, + name: &str, + channel_type: crate::channel::ChannelType, + visibility: crate::channel::ChannelVisibility, + description: Option<&str>, + created_by: &[u8], + ttl_seconds: Option, +) -> Result { + let id = Uuid::new_v4(); + create_channel_with_id( + pool, + community, + id, + name, + channel_type, + visibility, + description, + created_by, + ttl_seconds, + ) + .await + .map(|(r, _)| r) +} + +pub(crate) async fn get_canvas( + pool: &SqlitePool, + community: CommunityId, + channel_id: Uuid, +) -> Result> { + Ok(sqlx::query_scalar( + "SELECT canvas FROM channels WHERE community_id = ?1 AND id = ?2 AND deleted_at IS NULL", + ) + .bind(community.as_uuid().to_string()) + .bind(channel_id.to_string()) + .fetch_optional(pool) + .await? + .flatten()) +} + +pub(crate) async fn set_canvas( + pool: &SqlitePool, + community: CommunityId, + channel_id: Uuid, + canvas: Option<&str>, +) -> Result<()> { + let n = sqlx::query("UPDATE channels SET canvas = ?1, updated_at = unixepoch() WHERE community_id = ?2 AND id = ?3 AND deleted_at IS NULL") + .bind(canvas).bind(community.as_uuid().to_string()).bind(channel_id.to_string()).execute(pool).await?.rows_affected(); + if n == 0 { + return Err(crate::DbError::ChannelNotFound(channel_id)); + } + Ok(()) +} + +pub(crate) async fn remove_member( + pool: &SqlitePool, + community: CommunityId, + channel_id: Uuid, + pubkey: &[u8], + actor: &[u8], +) -> Result<()> { + let mut tx = pool.begin().await?; + let cid = community.as_uuid().to_string(); + let chid = channel_id.to_string(); + let is_self = pubkey == actor; + if !is_self { + let actor_role: Option = sqlx::query_scalar("SELECT cm.role FROM channel_members cm JOIN channels c ON c.id=cm.channel_id AND c.community_id=?1 WHERE cm.channel_id=?2 AND cm.pubkey=?3 AND cm.removed_at IS NULL AND c.deleted_at IS NULL") + .bind(&cid).bind(&chid).bind(actor).fetch_optional(&mut *tx).await?; + let agent_owner: i64 = sqlx::query_scalar("SELECT count(*) FROM users WHERE community_id=?1 AND pubkey=?2 AND agent_owner_pubkey=?3") + .bind(&cid).bind(pubkey).bind(actor).fetch_one(&mut *tx).await?; + let elevated = matches!(actor_role.as_deref(), Some("owner" | "admin")); + if !elevated && agent_owner == 0 { + return Err(crate::DbError::AccessDenied( + "only owners/admins or the agent's owner may remove other members".into(), + )); + } + } + let target_role: Option = sqlx::query_scalar( + "SELECT role FROM channel_members WHERE channel_id=?1 AND pubkey=?2 AND removed_at IS NULL", + ) + .bind(&chid) + .bind(pubkey) + .fetch_optional(&mut *tx) + .await?; + if target_role.is_none() { + return Err(crate::DbError::MemberNotFound(channel_id)); + } + if target_role.as_deref() == Some("owner") { + let count: i64 = sqlx::query_scalar("SELECT count(*) FROM channel_members WHERE channel_id=?1 AND role='owner' AND removed_at IS NULL").bind(&chid).fetch_one(&mut *tx).await?; + if count <= 1 { + return Err(crate::DbError::AccessDenied( + "cannot remove the last owner — transfer ownership first".into(), + )); + } + } + let n = sqlx::query("UPDATE channel_members SET removed_at=unixepoch() WHERE channel_id=?1 AND pubkey=?2 AND removed_at IS NULL").bind(&chid).bind(pubkey).execute(&mut *tx).await?.rows_affected(); + if n == 0 { + return Err(crate::DbError::MemberNotFound(channel_id)); + } + tx.commit().await?; + Ok(()) +} + +pub(crate) async fn membership_pairs( + pool: &SqlitePool, + community: CommunityId, + channel_ids: &[Uuid], + pubkeys: &[Vec], +) -> Result)>> { + if channel_ids.is_empty() || pubkeys.is_empty() { + return Ok(Vec::new()); + } + let ids = channel_ids + .iter() + .map(|_| "?") + .collect::>() + .join(","); + let keys = pubkeys.iter().map(|_| "?").collect::>().join(","); + let sql = format!("SELECT cm.channel_id,cm.pubkey FROM channel_members cm JOIN channels c ON c.id=cm.channel_id AND c.community_id=? AND c.deleted_at IS NULL WHERE cm.channel_id IN ({ids}) AND cm.pubkey IN ({keys}) AND cm.removed_at IS NULL"); + let mut q = sqlx::query(sqlx::AssertSqlSafe(sql)).bind(community.as_uuid().to_string()); + for id in channel_ids { + q = q.bind(id.to_string()); + } + for pk in pubkeys { + q = q.bind(pk); + } + let mut out = Vec::new(); + for r in q.fetch_all(pool).await? { + let id: String = r.try_get("channel_id")?; + out.push(( + Uuid::parse_str(&id).map_err(|e| crate::DbError::InvalidData(e.to_string()))?, + r.try_get("pubkey")?, + )); + } + Ok(out) +} + +pub(crate) async fn get_members_bulk( + pool: &SqlitePool, + community: CommunityId, + channel_ids: &[Uuid], +) -> Result> { + if channel_ids.is_empty() { + return Ok(Vec::new()); + } + let ids = channel_ids + .iter() + .map(|_| "?") + .collect::>() + .join(","); + let sql=format!("SELECT cm.channel_id,cm.pubkey,cm.role,cm.joined_at,cm.invited_by,cm.removed_at FROM channel_members cm JOIN channels c ON c.id=cm.channel_id AND c.community_id=? AND c.deleted_at IS NULL WHERE cm.channel_id IN ({ids}) AND cm.removed_at IS NULL ORDER BY cm.joined_at ASC"); + let mut q = sqlx::query(sqlx::AssertSqlSafe(sql)).bind(community.as_uuid().to_string()); + for id in channel_ids { + q = q.bind(id.to_string()); + } + q.fetch_all(pool) + .await? + .into_iter() + .map(member_record) + .collect() +} + +pub(crate) async fn list_channels( + pool: &SqlitePool, + community: CommunityId, + visibility: Option<&str>, +) -> Result> { + let mut sql = "SELECT * FROM channels WHERE community_id=? AND deleted_at IS NULL".to_string(); + if visibility.is_some() { + sql.push_str(" AND visibility=?"); + } + sql.push_str(" ORDER BY created_at DESC LIMIT 1000"); + let mut q = sqlx::query(sqlx::AssertSqlSafe(sql)).bind(community.as_uuid().to_string()); + if let Some(v) = visibility { + q = q.bind(v); + } + q.fetch_all(pool) + .await? + .into_iter() + .map(channel_record) + .collect() +} + +pub(crate) async fn get_accessible_channels( + pool: &SqlitePool, + community: CommunityId, + pubkey: &[u8], + visibility: Option<&str>, + member_only: Option, +) -> Result> { + let mut sql="SELECT c.*, (cm.channel_id IS NOT NULL) AS is_member FROM channels c LEFT JOIN channel_members cm ON cm.channel_id=c.id AND cm.pubkey=? AND cm.removed_at IS NULL WHERE c.community_id=? AND c.deleted_at IS NULL AND (c.channel_type != 'dm' OR cm.hidden_at IS NULL)".to_string(); + if member_only == Some(true) { + sql.push_str(" AND cm.channel_id IS NOT NULL"); + } else { + sql.push_str(" AND (c.visibility='open' OR cm.channel_id IS NOT NULL)"); + } + if visibility.is_some() { + sql.push_str(" AND c.visibility=?"); + } + sql.push_str(" ORDER BY CASE c.channel_type WHEN 'stream' THEN 0 WHEN 'forum' THEN 1 WHEN 'dm' THEN 2 ELSE 3 END,c.name LIMIT 1000"); + let mut q = sqlx::query(sqlx::AssertSqlSafe(sql)) + .bind(pubkey) + .bind(community.as_uuid().to_string()); + if let Some(v) = visibility { + q = q.bind(v); + } + let mut out = Vec::new(); + for r in q.fetch_all(pool).await? { + let m = r.try_get::("is_member")? != 0; + out.push(crate::channel::AccessibleChannel { + channel: channel_record(r)?, + is_member: m, + }); + } + Ok(out) +} + +pub(crate) async fn get_bot_members( + pool: &SqlitePool, + community: CommunityId, +) -> Result> { + let rows=sqlx::query("SELECT cm.pubkey,u.display_name FROM channel_members cm LEFT JOIN users u ON u.community_id=? AND u.pubkey=cm.pubkey JOIN channels c ON c.id=cm.channel_id AND c.deleted_at IS NULL WHERE c.community_id=? AND cm.role='bot' AND cm.removed_at IS NULL GROUP BY cm.pubkey,u.display_name LIMIT 1000").bind(community.as_uuid().to_string()).bind(community.as_uuid().to_string()).fetch_all(pool).await?; + let mut out = Vec::new(); + for r in rows { + let pk = r.try_get("pubkey")?; + let chans=sqlx::query("SELECT c.name,c.id FROM channel_members cm JOIN channels c ON c.id=cm.channel_id WHERE c.community_id=? AND cm.pubkey=? AND cm.role='bot' AND cm.removed_at IS NULL AND c.deleted_at IS NULL ORDER BY c.name").bind(community.as_uuid().to_string()).bind(&pk).fetch_all(pool).await?; + let channels = chans + .into_iter() + .map(|x| { + Ok(crate::channel::BotChannelEntry { + name: x.try_get("name")?, + id: x.try_get("id")?, + }) + }) + .collect::>>()?; + out.push(crate::channel::BotMemberRecord { + pubkey: pk, + display_name: r.try_get("display_name")?, + agent_type: None, + capabilities: None, + channels, + }); + } + Ok(out) +} + +pub(crate) async fn get_users_bulk( + pool: &SqlitePool, + community: CommunityId, + pubkeys: &[Vec], +) -> Result> { + if pubkeys.is_empty() { + return Ok(Vec::new()); + } + let keys = pubkeys.iter().map(|_| "?").collect::>().join(","); + let mut q=sqlx::query(sqlx::AssertSqlSafe(format!("SELECT pubkey,display_name,avatar_url,nip05_handle FROM users WHERE community_id=? AND pubkey IN ({keys})"))).bind(community.as_uuid().to_string()); + for pk in pubkeys { + q = q.bind(pk); + } + q.fetch_all(pool) + .await? + .into_iter() + .map(|r| { + Ok(crate::channel::UserRecord { + pubkey: r.try_get("pubkey")?, + display_name: r.try_get("display_name")?, + avatar_url: r.try_get("avatar_url")?, + nip05_handle: r.try_get("nip05_handle")?, + }) + }) + .collect() +} + +fn channel_not_found_if_zero(n: u64, id: Uuid) -> Result<()> { + if n == 0 { + Err(crate::DbError::ChannelNotFound(id)) + } else { + Ok(()) + } +} + +pub(crate) async fn update_channel( + pool: &SqlitePool, + community: CommunityId, + id: Uuid, + mut u: crate::channel::ChannelUpdate, +) -> Result { + if u.name.is_none() + && u.description.is_none() + && u.visibility.is_none() + && u.ttl_seconds.is_none() + { + return Err(crate::DbError::InvalidData( + "at least one field must be provided for update".into(), + )); + } + if let Some(n) = u.name.as_mut() { + *n = buzz_core::channel::canonical_channel_name(n).to_owned(); + if n.trim().is_empty() { + return Err(crate::DbError::InvalidData( + "channel name is required".into(), + )); + } + } + let mut sets = Vec::new(); + if u.name.is_some() { + sets.push("name=?"); + } + if u.description.is_some() { + sets.push("description=?"); + } + if u.visibility.is_some() { + sets.push("visibility=?"); + } + if u.ttl_seconds.is_some() { + sets.push("ttl_seconds=?"); + sets.push("ttl_deadline=CASE WHEN ? IS NULL THEN NULL ELSE unixepoch()+? END"); + } + sets.push("updated_at=unixepoch()"); + let mut q = sqlx::query(sqlx::AssertSqlSafe(format!( + "UPDATE channels SET {} WHERE community_id=? AND id=? AND deleted_at IS NULL", + sets.join(",") + ))); + if let Some(v) = &u.name { + q = q.bind(v); + } + if let Some(v) = &u.description { + q = q.bind(v); + } + if let Some(v) = &u.visibility { + q = q.bind(v); + } + if let Some(v) = u.ttl_seconds { + q = q.bind(v).bind(v).bind(v); + } + q = q.bind(community.as_uuid().to_string()).bind(id.to_string()); + channel_not_found_if_zero(q.execute(pool).await?.rows_affected(), id)?; + get_channel(pool, community, id).await +} + +pub(crate) async fn set_topic( + pool: &SqlitePool, + c: CommunityId, + id: Uuid, + v: &str, + by: &[u8], +) -> Result<()> { + let n=sqlx::query("UPDATE channels SET topic=?,topic_set_by=?,topic_set_at=unixepoch(),updated_at=unixepoch() WHERE community_id=? AND id=? AND deleted_at IS NULL").bind(v).bind(by).bind(c.as_uuid().to_string()).bind(id.to_string()).execute(pool).await?.rows_affected(); + channel_not_found_if_zero(n, id) +} +pub(crate) async fn set_purpose( + pool: &SqlitePool, + c: CommunityId, + id: Uuid, + v: &str, + by: &[u8], +) -> Result<()> { + let n=sqlx::query("UPDATE channels SET purpose=?,purpose_set_by=?,purpose_set_at=unixepoch(),updated_at=unixepoch() WHERE community_id=? AND id=? AND deleted_at IS NULL").bind(v).bind(by).bind(c.as_uuid().to_string()).bind(id.to_string()).execute(pool).await?.rows_affected(); + channel_not_found_if_zero(n, id) +} +pub(crate) async fn archive_channel(pool: &SqlitePool, c: CommunityId, id: Uuid) -> Result<()> { + let row = sqlx::query_scalar::<_, Option>( + "SELECT archived_at FROM channels WHERE community_id=? AND id=? AND deleted_at IS NULL", + ) + .bind(c.as_uuid().to_string()) + .bind(id.to_string()) + .fetch_optional(pool) + .await?; + match row { + None => Err(crate::DbError::ChannelNotFound(id)), + Some(Some(_)) => Err(crate::DbError::AccessDenied( + "channel is already archived".into(), + )), + Some(None) => { + sqlx::query("UPDATE channels SET archived_at=unixepoch() WHERE community_id=? AND id=? AND archived_at IS NULL AND deleted_at IS NULL").bind(c.as_uuid().to_string()).bind(id.to_string()).execute(pool).await?; + Ok(()) + } + } +} +pub(crate) async fn unarchive_channel(pool: &SqlitePool, c: CommunityId, id: Uuid) -> Result<()> { + let row = sqlx::query_scalar::<_, Option>( + "SELECT archived_at FROM channels WHERE community_id=? AND id=? AND deleted_at IS NULL", + ) + .bind(c.as_uuid().to_string()) + .bind(id.to_string()) + .fetch_optional(pool) + .await?; + match row { + None => Err(crate::DbError::ChannelNotFound(id)), + Some(None) => Err(crate::DbError::AccessDenied( + "channel is not archived".into(), + )), + Some(Some(_)) => { + sqlx::query("UPDATE channels SET archived_at=NULL,ttl_deadline=CASE WHEN ttl_seconds IS NULL THEN ttl_deadline ELSE unixepoch()+ttl_seconds END WHERE community_id=? AND id=? AND deleted_at IS NULL").bind(c.as_uuid().to_string()).bind(id.to_string()).execute(pool).await?; + Ok(()) + } + } +} +pub(crate) async fn soft_delete_channel( + pool: &SqlitePool, + c: CommunityId, + id: Uuid, +) -> Result { + Ok(sqlx::query("UPDATE channels SET deleted_at=unixepoch() WHERE community_id=? AND id=? AND deleted_at IS NULL").bind(c.as_uuid().to_string()).bind(id.to_string()).execute(pool).await?.rows_affected()>0) +} +pub(crate) async fn get_member_count(pool: &SqlitePool, c: CommunityId, id: Uuid) -> Result { + Ok(sqlx::query_scalar("SELECT count(*) FROM channel_members cm JOIN channels ch ON ch.id=cm.channel_id AND ch.community_id=? WHERE cm.channel_id=? AND cm.removed_at IS NULL").bind(c.as_uuid().to_string()).bind(id.to_string()).fetch_one(pool).await?) +} +pub(crate) async fn get_member_counts_bulk( + pool: &SqlitePool, + c: CommunityId, + ids: &[Uuid], +) -> Result> { + if ids.is_empty() { + return Ok(std::collections::HashMap::new()); + } + let ph = ids.iter().map(|_| "?").collect::>().join(","); + let mut q=sqlx::query(sqlx::AssertSqlSafe(format!("SELECT cm.channel_id,count(*) cnt FROM channel_members cm JOIN channels ch ON ch.id=cm.channel_id AND ch.community_id=? WHERE cm.channel_id IN ({ph}) AND cm.removed_at IS NULL GROUP BY cm.channel_id"))).bind(c.as_uuid().to_string()); + for id in ids { + q = q.bind(id.to_string()); + } + let mut out = std::collections::HashMap::new(); + for r in q.fetch_all(pool).await? { + let id: String = r.try_get("channel_id")?; + out.insert( + Uuid::parse_str(&id).map_err(|e| crate::DbError::InvalidData(e.to_string()))?, + r.try_get("cnt")?, + ); + } + Ok(out) +} +pub(crate) async fn reap_expired_ephemeral_channels( + pool: &SqlitePool, +) -> Result> { + let rows=sqlx::query("SELECT ch.community_id,c.host,ch.id FROM channels ch JOIN communities c ON c.id=ch.community_id WHERE ch.ttl_seconds IS NOT NULL AND ch.ttl_deadline Result> { chrono::DateTime::from_timestamp(value, 0).ok_or(crate::DbError::InvalidTimestamp(value)) } @@ -3170,6 +3628,107 @@ pub(crate) async fn add_relay_member( .rows_affected() == 1) } +pub(crate) async fn claim_relay_membership( + pool: &SqlitePool, + community: CommunityId, + pubkey: &str, + role: &str, + policy_version: Option<&str>, +) -> Result { + let mut tx = pool.begin().await?; + let inserted = sqlx::query( + "INSERT INTO relay_members (community_id, pubkey, role, added_by) VALUES (?1, lower(?2), ?3, 'invite') ON CONFLICT DO NOTHING", + ) + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .bind(role) + .execute(&mut *tx) + .await? + .rows_affected() + == 1; + if let Some(version) = policy_version { + sqlx::query("INSERT INTO join_policy_acceptances (community_id, pubkey, policy_version) VALUES (?1, lower(?2), ?3) ON CONFLICT DO NOTHING") + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .bind(version) + .execute(&mut *tx) + .await?; + } + tx.commit().await?; + Ok(inserted) +} + +pub(crate) async fn has_join_policy_acceptance( + pool: &SqlitePool, + community: CommunityId, + pubkey: &str, + policy_version: &str, +) -> Result { + Ok(sqlx::query_scalar::<_, i64>( + "SELECT count(*) FROM join_policy_acceptances WHERE community_id = ?1 AND pubkey = lower(?2) AND policy_version = ?3", + ) + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .bind(policy_version) + .fetch_one(pool) + .await? + != 0) +} + +pub(crate) async fn remove_relay_member( + pool: &SqlitePool, + community: CommunityId, + pubkey: &str, + expected_role: Option<&str>, +) -> Result { + use crate::relay_members::RemoveResult; + let result = if let Some(role) = expected_role { + sqlx::query("DELETE FROM relay_members WHERE community_id = ?1 AND pubkey = lower(?2) AND role = ?3") + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .bind(role) + .execute(pool) + .await? + } else { + sqlx::query("DELETE FROM relay_members WHERE community_id = ?1 AND pubkey = lower(?2) AND role <> 'owner'") + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .execute(pool) + .await? + }; + if result.rows_affected() == 1 { + return Ok(RemoveResult::Removed); + } + let role = sqlx::query_scalar::<_, String>( + "SELECT role FROM relay_members WHERE community_id = ?1 AND pubkey = lower(?2)", + ) + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .fetch_optional(pool) + .await?; + Ok(match role.as_deref() { + None => RemoveResult::NotFound, + Some("owner") => RemoveResult::IsOwner, + Some(_) if expected_role.is_some() => RemoveResult::RoleMismatch, + Some(_) => RemoveResult::NotFound, + }) +} + +pub(crate) async fn update_relay_member_role( + pool: &SqlitePool, + community: CommunityId, + pubkey: &str, + new_role: &str, +) -> Result { + Ok(sqlx::query("UPDATE relay_members SET role = ?3, updated_at = unixepoch() WHERE community_id = ?1 AND pubkey = lower(?2) AND role <> 'owner'") + .bind(community.as_uuid().to_string()) + .bind(pubkey) + .bind(new_role) + .execute(pool) + .await? + .rows_affected() == 1) +} + pub(crate) async fn bootstrap_owner( pool: &SqlitePool, community: CommunityId, @@ -4333,6 +4892,70 @@ mod tests { assert!(get_workflow_run(&pool, community, run_id).await.is_err()); } + #[tokio::test] + async fn relay_member_mutations_preserve_policy_owner_and_role_contracts() { + use crate::relay_members::RemoveResult; + + let pool = connect(":memory:").await.unwrap(); + let community = ensure_configured_community(&pool, "members.example") + .await + .unwrap() + .id; + bootstrap_owner(&pool, community, "OWNER").await.unwrap(); + + assert!( + claim_relay_membership(&pool, community, "ALICE", "member", Some("v1")) + .await + .unwrap() + ); + assert!( + !claim_relay_membership(&pool, community, "alice", "member", Some("v1")) + .await + .unwrap() + ); + assert!(has_join_policy_acceptance(&pool, community, "ALICE", "v1") + .await + .unwrap()); + assert_eq!( + remove_relay_member(&pool, community, "alice", Some("admin")) + .await + .unwrap(), + RemoveResult::RoleMismatch + ); + assert!(update_relay_member_role(&pool, community, "alice", "admin") + .await + .unwrap()); + assert_eq!( + remove_relay_member(&pool, community, "alice", Some("member")) + .await + .unwrap(), + RemoveResult::RoleMismatch + ); + assert_eq!( + remove_relay_member(&pool, community, "alice", Some("admin")) + .await + .unwrap(), + RemoveResult::Removed + ); + assert_eq!( + remove_relay_member(&pool, community, "owner", None) + .await + .unwrap(), + RemoveResult::IsOwner + ); + assert!( + !update_relay_member_role(&pool, community, "owner", "member") + .await + .unwrap() + ); + assert_eq!( + remove_relay_member(&pool, community, "missing", None) + .await + .unwrap(), + RemoveResult::NotFound + ); + } + #[tokio::test] async fn core_slice_survives_temporary_file_reopen() { let path = std::env::temp_dir().join(format!("buzz-db-{}.sqlite", Uuid::new_v4()));