add basic concord backend

This commit is contained in:
2026-09-16 16:57:04 +07:00
parent 5f2a5d7a37
commit 66f75ad105
7 changed files with 891 additions and 27 deletions
+156
View File
@@ -0,0 +1,156 @@
use std::collections::BTreeMap;
use std::sync::LazyLock;
use anyhow::{Result, anyhow};
use nostr_sdk::prelude::*;
use crate::ChannelId;
use crate::stream::OpenedStream;
static LOCAL_KEYS: LazyLock<Keys> = LazyLock::new(Keys::generate);
const CHANNEL_TAG: SingleLetterTag = SingleLetterTag::LOWERCASE_C;
const MARK_TAG: SingleLetterTag = SingleLetterTag::LOWERCASE_T;
const MARK_VALUE: &str = "concord";
const WRAP_TAG: &str = "e";
const KIND_TAG: &str = "k";
pub async fn cache_rumor(
database: &dyn NostrDatabase,
channel: &ChannelId,
opened: &OpenedStream,
) -> Result<()> {
let tags = vec![
Tag::identifier(opened.rumor_id),
Tag::custom(KIND_TAG, [opened.rumor.kind.to_string()]),
Tag::custom(WRAP_TAG, [opened.wrapper_id.to_string()]),
Tag::custom(MARK_TAG.as_str(), [MARK_VALUE]),
Tag::custom(CHANNEL_TAG.as_str(), [channel.to_hex()]),
Tag::public_key(opened.author),
];
let at = Timestamp::from_secs(opened.at_ms / 1000);
let event = EventBuilder::new(Kind::ApplicationSpecificData, opened.rumor.as_json())
.tags(tags)
.custom_created_at(at)
.finalize_async(&*LOCAL_KEYS)
.await?;
database.save_event(&event).await?;
Ok(())
}
/// Read a channel's cached rumors.
pub async fn query_rumors(
database: &dyn NostrDatabase,
channel: &ChannelId,
until: Option<Timestamp>,
limit: usize,
) -> Result<Vec<UnsignedEvent>> {
let mut filter = Filter::new()
.kind(Kind::ApplicationSpecificData)
.custom_tag(MARK_TAG, MARK_VALUE)
.custom_tag(CHANNEL_TAG, channel.to_hex());
if let Some(until) = until {
filter = filter.until(until);
}
let mut newest: BTreeMap<String, Event> = BTreeMap::new();
for event in database.query(filter).await? {
let Some(rumor_id) = event.tags.identifier() else {
continue;
};
match newest.get(&rumor_id) {
Some(existing) if existing.created_at >= event.created_at => {}
_ => {
newest.insert(rumor_id, event);
}
}
}
let mut events: Vec<Event> = newest.into_values().collect();
events.sort_by_key(|event| std::cmp::Reverse(event.created_at));
events.truncate(limit);
let mut rumors = Vec::with_capacity(events.len());
for event in events {
let rumor = UnsignedEvent::from_json(event.content)
.map_err(|error| anyhow!("cached rumor is not a valid event: {error}"))?;
rumors.push(rumor);
}
Ok(rumors)
}
#[cfg(test)]
mod tests {
use nostr_memory::MemoryDatabase;
use super::*;
use crate::Epoch;
use crate::derive::channel_group_key;
use crate::stream::{
KIND_WRAP, SealForm, build_rumor_ms, build_seal, channel_binding_tags, open_wrap, wrap_seal,
};
const SECRET: [u8; 32] = [0x07u8; 32];
#[test]
fn rumors_read_back_after_a_restart() {
let database = MemoryDatabase::unbounded();
let channel = ChannelId::from_bytes([0xabu8; 32]);
let author = Keys::generate();
smol::block_on(async {
let group = channel_group_key(&SECRET, &channel, Epoch(0)).expect("derives");
for (content, at_ms) in [("first", 1_000_000u64), ("second", 2_000_000)] {
let rumor = build_rumor_ms(
9,
author.public_key(),
content,
channel_binding_tags(&channel, Epoch(0)),
at_ms,
);
let seal = build_seal(&rumor, SealForm::Encrypted, &group, &author).expect("seals");
let (wrap, _) = wrap_seal(
&seal,
&group,
KIND_WRAP,
Timestamp::from_secs(at_ms / 1000),
&[],
)
.expect("wraps");
let opened = open_wrap(&wrap, &group).expect("opens");
cache_rumor(&database, &channel, &opened)
.await
.expect("caches");
}
// The group key is gone; only the local cache stands in for it.
let rumors = query_rumors(&database, &channel, None, 10)
.await
.expect("queries");
assert_eq!(rumors.len(), 2, "both messages come back");
assert_eq!(rumors[0].content, "second", "newest first");
assert_eq!(rumors[1].content, "first");
// A page boundary in message time, not in cache time.
let until = Timestamp::from_secs(1_500);
let page = query_rumors(&database, &channel, Some(until), 10)
.await
.expect("queries");
assert_eq!(page.len(), 1);
assert_eq!(page[0].content, "first");
let capped = query_rumors(&database, &channel, None, 1)
.await
.expect("queries");
assert_eq!(capped.len(), 1);
assert_eq!(capped[0].content, "second");
});
}
}