Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
de49272
feat(relay): draft private read-state accessory API
Sep 26, 2026
b6a9a9a
Merge origin/main into meli/buzz-v1-read-state
Sep 26, 2026
85ccf27
feat(relay): complete bounded private read-state projection
Sep 28, 2026
f4a54a4
Merge origin/main into meli/buzz-v1-read-state
Sep 28, 2026
beb9922
fix(db): follow renumbered personal read migration in parity test
Sep 28, 2026
7451d79
docs(relay): separate extension design from current API contract
Sep 28, 2026
4288ba5
test(db): separate participation root cap from SQL deadline
Sep 28, 2026
a81e8fb
feat(relay): whole-channel reads, thread summaries and targeted refresh
Sep 28, 2026
64de580
Merge origin/main (4ef23609b) into meli/buzz-v1-read-state
Sep 30, 2026
49dfc89
fix(relay): make /buzz/v1 invisible when disabled and drop its events…
Sep 30, 2026
0f01e45
docs(buzz-v1): scope capability absence; witness that diffs are not r…
Sep 30, 2026
abc7dd8
Merge origin/main (0ee609379) into meli/buzz-v1-read-state
Sep 30, 2026
bb3aaa1
Merge origin/main (965e1997b) into meli/buzz-v1-read-state
Sep 30, 2026
6904051
Merge origin/main (53a12100b) into meli/buzz-v1-read-state
Sep 30, 2026
9cfe40d
test(relay): cover signed sidebar deletion counts
Sep 30, 2026
b8c1eb0
refactor(db): share sidebar evidence completeness flag
Sep 30, 2026
2ffa6bf
Merge origin/main (5fdb2e536) into meli/buzz-v1-read-state
Oct 1, 2026
1f0333a
Merge origin/main (fad4ff637) into meli/buzz-v1-read-state
Oct 1, 2026
c489742
Merge origin/main (83aab8cb5) into meli/buzz-v1-read-state
Oct 1, 2026
644dbbd
docs(buzz-v1): state the rollout order and scope the NIP-RS sentence
Oct 1, 2026
81e6657
Merge origin/main (04b040ef6) into meli/buzz-v1-read-state
Oct 1, 2026
d21507f
Merge origin/main (16839a077) into meli/buzz-v1-read-state
Oct 1, 2026
4be404e
fix(relay): keep the NIP-FI denial wire shape on /buzz/v1
Oct 1, 2026
533de68
feat(buzz-v1): count a reply as unread only when it is directed at th…
Oct 2, 2026
a09bbdf
Merge origin/main (639593bba) into meli/buzz-v1-read-state
Oct 2, 2026
9894641
fix(buzz-v1): let a thread on a never-unread root be read
Oct 2, 2026
136e247
Merge origin/main (8af2d91f3) into meli/buzz-v1-read-state
Oct 2, 2026
ce9481a
fix(buzz-v1): answer as Off in NIP-FI shadow mode, and record its ver…
Oct 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 3 additions & 3 deletions crates/buzz-db/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,9 +64,9 @@ pub(crate) use runtime::{
pub use store::{
admin_moderation, allowlist, api_token, archived_identities, artifact, channel,
channel_members, community, deletion, dm, event, feed, git_repo, moderation, operator_listener,
partition, product_feedback, push, reaction, read_state, relay_admin_actions, relay_invite,
relay_members, relay_operators, reminder, replaceable, storage_accounting, thread,
thread_window, usage, user, workflow,
partition, personal_read, product_feedback, push, reaction, read_state, relay_admin_actions,
relay_invite, relay_members, relay_operators, reminder, replaceable, storage_accounting,
thread, thread_window, usage, user, workflow,
};

pub use allowlist::AllowlistEntry;
Expand Down
23 changes: 22 additions & 1 deletion crates/buzz-db/src/runtime/migration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -705,7 +705,12 @@ mod postgres_tests {
let mut migrations: Vec<_> = MIGRATOR.iter().collect();
migrations.sort_by_key(|migration| migration.version);

assert_eq!(migrations.len(), 55);
assert_eq!(migrations.len(), 56);
assert_eq!(migrations[55].version, 56);
assert!(migrations[55]
.sql
.as_str()
.contains("CREATE TABLE personal_read_accounts"));
assert_eq!(migrations[48].version, 49);
assert_eq!(migrations[49].version, 50);
assert_eq!(migrations[50].version, 51);
Expand Down Expand Up @@ -2048,6 +2053,22 @@ mod postgres_tests {
let mut expected_fences = migration.fence_attachments.clone();
expected_fences.remove("product_feedback");
expected_fences.remove("rate_limit_violations");
let personal = surface(
MIGRATOR
.iter()
.find(|m| m.version == 56)
.expect("personal read migration")
.sql
.as_ref(),
);
for (table, definition) in personal.tables {
assert_eq!(
schema.tables.get(&table),
Some(&definition),
"personal read table {table} differs"
);
}
expected_fences.extend(personal.fence_attachments);
expected_fences.extend(["artifact_heads", "artifact_revisions"].map(str::to_owned));
assert_eq!(
expected_fences, schema.fence_attachments,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -311,7 +311,7 @@ async fn migration_schema_thread_window_prebuild_does_not_queue_behind_writer()
production_result.is_ok(),
"production migrator must preserve ingestion progress: {production_result:?}"
);
assert_eq!(version, 55);
assert_eq!(version, 56);
assert_eq!(final_oid, oid, "prebuild must not be replaced");
assert_eq!(count, 4, "all writer witnesses must persist");
}
4 changes: 4 additions & 0 deletions crates/buzz-db/src/store/deletion.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,8 @@ pub const EXPECTED_SCOPED_TABLES: &[&str] = &[
"moderation_actions",
"moderation_reports",
"parameterized_event_watermarks",
"personal_read_accounts",
"personal_read_frontiers",
"pubkey_allowlist",
"push_leases",
"push_match_queue",
Expand All @@ -110,6 +112,8 @@ pub const EXPECTED_SCOPED_TABLES: &[&str] = &[

/// Foreign-key-safe child-before-parent order for the PostgreSQL purge.
pub const PURGE_SCOPED_TABLES: &[&str] = &[
"personal_read_frontiers",
"personal_read_accounts",
"workflow_approvals",
"scheduled_workflow_fires",
"workflow_runs",
Expand Down
2 changes: 2 additions & 0 deletions crates/buzz-db/src/store/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ pub mod moderation;
pub mod operator_listener;
/// Monthly table partition management.
pub mod partition;
/// Private signer-owned accessory read progress.
pub mod personal_read;
/// Buzz product-feedback sidecar persistence.
pub mod product_feedback;
/// Community-scoped push lease and durable wake-outbox persistence.
Expand Down
30 changes: 30 additions & 0 deletions crates/buzz-db/src/store/personal_read/classification.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
//! Selector eligibility and shared directed-reason rules. Aggregate SQL applies
//! the same eligibility before grouping, covered by the PostgreSQL parity test.
use super::model::{Reason, ELIGIBLE_KINDS};

pub(super) fn eligible(
kind: i32,
own: bool,
deleted: bool,
created_ms: i64,
cutoff_ms: i64,
) -> bool {
ELIGIBLE_KINDS.contains(&kind) && !own && !deleted && created_ms >= cutoff_ms
}

/// Why a message is directed, before conversation membership is known.
pub(super) fn reason(channel_type: &str, actor_hex: &str, tags: &[Vec<String>]) -> Option<Reason> {
let tagged = |name: &str, matches: &dyn Fn(&str) -> bool| {
tags.iter()
.any(|tag| tag.len() >= 2 && tag[0] == name && matches(&tag[1]))
};
if channel_type == "dm" {
Some(Reason::Direct)
} else if tagged("p", &|value| value.eq_ignore_ascii_case(actor_hex)) {
Some(Reason::Mention)
} else if tagged("broadcast", &|value| value == "1") {
Some(Reason::Broadcast)
} else {
None
}
}
259 changes: 259 additions & 0 deletions crates/buzz-db/src/store/personal_read/context.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,259 @@
//! Explicit selectors over the same private frontier authority. No history API.
use super::{classification, model::*, participation, projection::read_account, writes};
use crate::{observability, Db, DbError, Result};
use buzz_core::CommunityId;
use chrono::{DateTime, Utc};
use sqlx::{Acquire, Row};
use std::collections::HashMap;
use uuid::Uuid;

impl Db {
/// Resolve bounded explicit contexts/messages in a single read-only snapshot.
/// Callers must recheck admission and resource access outside this snapshot.
pub async fn personal_read_contexts(
&self,
community: CommunityId,
actor: &nostr::PublicKey,
retention_seconds: u32,
queries: &[ContextQuery],
) -> Result<ContextPage> {
if queries.is_empty()
|| queries.len() > MAX_CONTEXTS
|| queries.iter().map(|q| q.message_ids.len()).sum::<usize>() > MAX_CONTEXT_MESSAGES
|| queries.iter().any(|q| {
q.message_ids
.iter()
.any(|id| writes::event_id(id).is_none())
})
{
return Err(DbError::InvalidData("invalid context selectors".into()));
}
let mut conn = observability::acquire_writer(
&self.pool,
observability::WriterOperation::SubscriptionHistory,
)
.await?;
let mut tx = conn.begin().await?;
sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ READ ONLY")
.execute(&mut *tx)
.await?;
writes::deadlines(&mut tx).await?;
let actor_bytes = actor.to_bytes();
let account = read_account(&mut tx, retention_seconds).await?;
let mut contexts = Vec::with_capacity(queries.len());
// Replies whose state turns on conversation membership, by position.
let mut pending = Vec::new();
for query in queries {
let root =
match writes::valid_target(&mut tx, community, &actor_bytes, &query.target).await {
Ok(Some(root)) => root,
Ok(None) => {
contexts.push(ContextState::Unavailable);
continue;
}
Err(DbError::InvalidData(_)) => {
contexts.push(ContextState::Unknown);
continue;
}
Err(error) => return Err(error),
};
// A thread's effective prefix includes the channel's whole-channel cut,
// exactly as the sidebar projection counts it.
let prefix: Option<i64> = sqlx::query_scalar(
"SELECT GREATEST(
(SELECT through_timestamp FROM personal_read_frontiers
WHERE community_id=$1 AND actor=$2 AND channel_id=$3 AND root_id=$4),
(SELECT threads_through_timestamp FROM personal_read_frontiers
WHERE community_id=$1 AND actor=$2 AND channel_id=$3 AND root_id=''::bytea
AND $4<>''::bytea))",
)
.bind(community.as_uuid())
.bind(actor_bytes.as_slice())
.bind(query.target.channel_id)
.bind(&root)
.fetch_one(&mut *tx)
.await?;
let ids: Vec<Vec<u8>> = query
.message_ids
.iter()
.filter_map(|id| writes::event_id(id))
.collect();
let rows = sqlx::query(
"SELECT encode(e.id,'hex') AS id,e.kind,e.created_at,
e.deleted_at IS NOT NULL AS deleted,e.pubkey=$3 AS own,
CASE WHEN octet_length(e.tags::text)<=8192 THEN e.tags ELSE NULL END AS tags,
tm.root_event_id,tm.parent_event_id,c.channel_type::text AS channel_type
FROM events e JOIN channels c ON c.community_id=e.community_id AND c.id=e.channel_id
LEFT JOIN thread_metadata tm ON tm.community_id=e.community_id
AND tm.event_id=e.id AND tm.event_created_at=e.created_at AND tm.channel_id=e.channel_id
WHERE e.community_id=$1 AND e.channel_id=$2 AND e.id=ANY($4)",
).bind(community.as_uuid()).bind(query.target.channel_id)
.bind(actor_bytes.as_slice()).bind(&ids).fetch_all(&mut *tx).await?;
let by_id: HashMap<String, _> = rows
.into_iter()
.map(|row| Ok((row.try_get::<String, _>("id")?, row)))
.collect::<Result<_>>()?;
let mut messages = Vec::with_capacity(ids.len());
for id in &query.message_ids {
let state = if let Some(row) = by_id.get(&id.to_ascii_lowercase()) {
let tags: Option<serde_json::Value> = row.try_get("tags")?;
let parsed =
tags.and_then(|v| serde_json::from_value::<Vec<Vec<String>>>(v).ok());
if let Some(tags) = parsed {
let canonical: Option<Vec<u8>> = row.try_get("root_event_id")?;
let marked_reply = buzz_core::nip10::parse_thread_markers_from_parts(
tags.iter().map(Vec::as_slice),
)
.resolve()
.is_some();
let message_id = writes::event_id(id).unwrap_or_default();
let is_reply = canonical.as_ref().is_some_and(|r| r != &message_id);
if marked_reply && canonical.is_none() {
MessageReadState::Unknown
} else if (root.is_empty() && is_reply)
|| (!root.is_empty()
&& (!is_reply || canonical.as_ref() != Some(&root)))
{
// The root's own timeline state is never the thread's state.
MessageReadState::Unavailable
} else {
let created: DateTime<Utc> = row.try_get("created_at")?;
let kind: i32 = row.try_get("kind")?;
if !classification::eligible(
kind,
row.try_get("own")?,
row.try_get("deleted")?,
created.timestamp_millis(),
account.cutoff_ms,
) {
MessageReadState::NotCounted
} else if prefix.is_some_and(|p| created.timestamp() <= p) {
MessageReadState::Read
} else {
let reason = classification::reason(
&row.try_get::<String, _>("channel_type")?,
&actor.to_hex(),
&tags,
);
// Membership outranks a broadcast and decides a
// plain reply. A DM or mention needs no lookup.
if is_reply
&& !matches!(reason, Some(Reason::Direct | Reason::Mention))
{
let parent: Option<Vec<u8>> = row.try_get("parent_event_id")?;
if let Some(parent) = parent {
pending.push((
contexts.len(),
messages.len(),
(query.target.channel_id, parent),
));
}
}
// Provisional for a pending reply: see `settle`.
if is_reply && reason.is_none() {
MessageReadState::Unknown
} else {
MessageReadState::Unread { reason }
}
}
}
} else {
MessageReadState::Unknown
}
} else {
MessageReadState::Unavailable
};
messages.push(ContextMessage {
message_id: id.clone(),
state,
});
}
contexts.push(ContextState::Available {
through_timestamp: prefix,
messages,
});
}
let targets: Vec<_> = pending.iter().map(|(.., key)| key.clone()).collect();
let members = participation::resolve(&mut tx, community, &actor_bytes, &targets).await?;
for (context, message, key) in pending {
if let ContextState::Available { messages, .. } = &mut contexts[context] {
settle(&mut messages[message].state, members.get(&key).copied());
}
}
tx.commit().await?;
Ok(ContextPage { account, contexts })
}

/// Final bounded access check on the writer, outside a projection snapshot.
/// Open-channel access is independent of joined-sidebar membership.
pub async fn personal_read_accessible_contexts(
&self,
community: CommunityId,
actor: &nostr::PublicKey,
channels: &[Uuid],
) -> Result<Vec<Uuid>> {
if channels.len() > MAX_CONTEXTS {
return Err(DbError::InvalidData("too many context channels".into()));
}
let mut conn = observability::acquire_writer(
&self.pool,
observability::WriterOperation::Authorization,
)
.await?;
Ok(sqlx::query_scalar(
"SELECT c.id FROM channels c WHERE c.community_id=$1 AND c.id=ANY($2)
AND c.deleted_at IS NULL AND (c.visibility='open' OR EXISTS (
SELECT 1 FROM channel_members cm WHERE cm.community_id=$1
AND cm.channel_id=c.id AND cm.pubkey=$3 AND cm.removed_at IS NULL))",
)
.bind(community.as_uuid())
.bind(channels)
.bind(actor.to_bytes().as_slice())
.fetch_all(&mut *conn)
.await?)
}
}

/// Apply conversation membership to a pending reply's provisional state:
/// `unknown` for a plain reply, `unread` with reason `broadcast` for a
/// broadcast. An undecided lookup (`None`) leaves it as it is.
fn settle(state: &mut MessageReadState, member: Option<bool>) {
match member {
Some(true) => {
*state = MessageReadState::Unread {
reason: Some(Reason::Conversation),
}
}
// A broadcast counts outside the actor's conversations too.
Some(false) if matches!(state, MessageReadState::Unknown) => {
*state = MessageReadState::NotCounted
}
_ => {}
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn membership_settles_a_pending_reply_and_an_undecided_lookup_fabricates_nothing() {
let broadcast = || MessageReadState::Unread {
reason: Some(Reason::Broadcast),
};
for (provisional, member, expected) in [
(MessageReadState::Unknown, Some(true), "conversation"),
(MessageReadState::Unknown, Some(false), "not_counted"),
(MessageReadState::Unknown, None, "unknown"),
(broadcast(), Some(true), "conversation"),
(broadcast(), Some(false), "broadcast"),
(broadcast(), None, "broadcast"),
] {
let mut state = provisional;
settle(&mut state, member);
let wire = serde_json::to_value(&state).unwrap();
let got = wire["reason"].as_str().or(wire["status"].as_str());
assert_eq!(got, Some(expected), "{member:?}");
}
}
}
Loading
Loading