Nostr Client Architecture - Production Implementation¶
Verzió: 1.0
Mélység: System Architecture
Cél: Production-ready client design patterns
1. Client Architecture Overview¶
1.1 System Layers¶
┌─────────────────────────────────────────────────────────────┐
│ PRESENTATION LAYER │
│ ┌──────────────┐ ┌──────────────┐ ┌─────────────────────┐ │
│ │ UI │ │ Navigation │ │ View Models │ │
│ │ Components │ │ (Router) │ │ (State Binding) │ │
│ └──────────────┘ └──────────────┘ └─────────────────────┘ │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ APPLICATION LAYER │
│ ┌──────────────┐ ┌──────────────┐ ┌─────────────────────┐ │
│ │ Services │ │ Use Cases │ │ Coordinators │ │
│ │ (Business) │ │ (Actions) │ │ (Flows) │ │
│ └──────────────┘ └──────────────┘ └─────────────────────┘ │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ DOMAIN LAYER │
│ ┌──────────────┐ ┌──────────────┐ ┌─────────────────────┐ │
│ │ Models │ │ Repositories │ │ Domain Logic │ │
│ │ (Events) │ │ (Interface) │ │ (Validation) │ │
│ └──────────────┘ └──────────────┘ └─────────────────────┘ │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ INFRASTRUCTURE LAYER │
│ ┌──────────────┐ ┌──────────────┐ ┌─────────────────────┐ │
│ │ Network │ │ Storage │ │ Crypto │ │
│ │(Relay Pool) │ │ (SQLite/LMDB)│ │ (Signing) │ │
│ └──────────────┘ └──────────────┘ └─────────────────────┘ │
└─────────────────────────────────────────────────────────────┘
2. State Management¶
2.1 Event Store Architecture¶
// Core event store interface
pub trait EventStore: Send + Sync {
async fn save(&self, event: Event) -> Result<(), Error>;
async fn get_by_id(&self, id: EventId) -> Result<Option<Event>, Error>;
async fn query(&self, filter: Filter) -> Result<Vec<Event>, Error>;
async fn delete(&self, id: EventId) -> Result<(), Error>;
async fn observe(&self, filter: Filter) -> BoxStream<Event>;
}
// SQLite implementation with proper indexing
pub struct SQLiteEventStore {
pool: SqlitePool,
deduplicator: Arc<Mutex<HashSet<EventId>>>,
}
impl SQLiteEventStore {
pub async fn new(path: &Path) -> Result<Self, Error> {
let pool = SqlitePool::connect(path).await?;
// Create tables with proper indexes
sqlx::query(
r#"
CREATE TABLE IF NOT EXISTS events (
id BLOB PRIMARY KEY,
pubkey BLOB NOT NULL,
created_at INTEGER NOT NULL,
kind INTEGER NOT NULL,
content TEXT NOT NULL,
sig BLOB NOT NULL,
tags TEXT NOT NULL, -- JSON
raw BLOB NOT NULL -- Full serialized event
);
CREATE INDEX IF NOT EXISTS idx_events_pubkey
ON events(pubkey, created_at DESC);
CREATE INDEX IF NOT EXISTS idx_events_kind
ON events(kind, created_at DESC);
CREATE INDEX IF NOT EXISTS idx_events_time
ON events(created_at DESC);
CREATE VIRTUAL TABLE IF NOT EXISTS event_fts
USING fts5(content, content='events', content_rowid='rowid');
"#
).execute(&pool).await?;
Ok(Self {
pool,
deduplicator: Arc::new(Mutex::new(HashSet::new())),
})
}
}
2.2 Reactive State (Rust with Tokio/Rx)¶
use tokio::sync::{watch, broadcast};
pub struct ReactiveStore {
// Current state snapshot
state: watch::Sender<AppState>,
// Event bus for real-time updates
event_bus: broadcast::Sender<StoreEvent>,
// Separate caches for different data types
timeline_cache: Arc<DashMap<String, Vec<Event>>>,
profile_cache: Arc<DashMap<PublicKey, Metadata>>,
relay_cache: Arc<DashMap<PublicKey, Vec<RelayUrl>>>,
}
pub struct AppState {
current_user: Option<User>,
active_relays: Vec<RelayConnection>,
unread_counts: HashMap<String, usize>,
network_status: NetworkStatus,
}
impl ReactiveStore {
pub async fn subscribe_timeline(
&self,
filter: Filter,
) -> impl Stream<Item = Vec<Event>> {
let (tx, rx) = mpsc::channel(100);
let mut event_stream = self.event_bus.subscribe();
tokio::spawn(async move {
while let Ok(store_event) = event_stream.recv().await {
if let StoreEvent::NewEvents(events) = store_event {
let filtered: Vec<Event> = events
.into_iter()
.filter(|e| filter.matches(e))
.collect();
if !filtered.is_empty() {
let _ = tx.send(filtered).await;
}
}
}
});
ReceiverStream::new(rx)
}
}
3. Synchronization Strategy¶
3.1 Sync Algorithm¶
Initial Sync:
1. Load relay list from NIP-65 (kind 10002)
2. Connect to all write + read relays
3. Query metadata for followed users (kind 0)
4. Query contact lists (kind 3)
5. Query recent timeline (kind 1, last 24h)
6. Set up real-time subscriptions
Incremental Sync:
1. Maintain "last sync" timestamp per relay
2. Query: events since last_sync
3. Handle conflicts (latest wins for replaceable)
4. Merge new events into local store
5. Update last_sync timestamp
Conflict Resolution:
- Regular events: Keep all (unique by ID)
- Replaceable (0, 3, 10000-19999): Latest created_at wins
- Addressable (30000-39999): Latest per (kind, author, d-tag)
3.2 Offline Support¶
pub struct OfflineManager {
pending_events: Arc<Mutex<Vec<PendingEvent>>>,
sync_queue: Arc<Mutex<Vec<SyncJob>>>,
network_monitor: NetworkMonitor,
}
impl OfflineManager {
pub async fn publish_event(
&self,
event: Event,
) -> Result<PublishResult, Error> {
if self.network_monitor.is_online() {
// Try immediate publish
match self.relay_pool.broadcast(event.clone()).await {
Ok(results) => {
// Check acceptance rate
let accepted = results.iter()
.filter(|r| r.is_ok())
.count();
if accepted >= self.min_relays {
return Ok(PublishResult::Published);
}
}
Err(_) => {}
}
}
// Queue for retry
self.queue_event(event).await;
Ok(PublishResult::Queued)
}
async fn retry_queued(&self) {
let pending = self.pending_events.lock().await;
for event in pending.iter() {
match self.publish_event(event.clone()).await {
Ok(PublishResult::Published) => {
// Remove from queue
self.remove_from_queue(event.id).await;
}
_ => {
// Keep for next retry with backoff
event.increment_retry_count();
}
}
}
}
}
4. Relay Pool Management¶
4.1 Intelligent Routing¶
pub struct RelayRouter {
relays: Arc<DashMap<RelayUrl, RelayState>>,
outbox_cache: Arc<DashMap<PublicKey, Vec<RelayUrl>>>, // NIP-65 cache
metrics: Arc<RatelimiterMetrics>,
}
impl RelayRouter {
// Publish to user's write relays
pub async fn publish_to_outbox(
&self,
event: &Event,
) -> Vec<Result<(), Error>> {
let author = event.pubkey;
// Get author's write relays from NIP-65
let relays = self.outbox_cache
.get(&author)
.map(|e| e.value().clone())
.unwrap_or_else(|| self.get_default_relays());
let mut results = vec![];
for relay_url in relays {
if let Some(relay) = self.relays.get(&relay_url) {
let result = relay.send_event(event).await;
results.push(result);
}
}
results
}
// Read from user's read relays + their write relays (outbox model)
pub async fn subscribe_to_feed(
&self,
follows: Vec<PublicKey>,
) -> impl Stream<Item = Event> {
let mut all_relays = HashSet::new();
// My read relays
all_relays.extend(self.get_my_read_relays());
// Followed users' write relays (outbox)
for pubkey in follows {
if let Some(relays) = self.outbox_cache.get(&pubkey) {
all_relays.extend(relays.value().clone());
}
}
// Subscribe to all
let (tx, rx) = mpsc::channel(1000);
for relay_url in all_relays {
let tx = tx.clone();
tokio::spawn(async move {
// Subscribe logic
});
}
ReceiverStream::new(rx)
}
}
4.2 Quality Scoring¶
pub struct RelayQuality {
url: RelayUrl,
score: f64, // 0.0 to 1.0
metrics: RelayMetrics,
}
pub struct RelayMetrics {
latency_ms: u32,
success_rate: f64, // % of successful publishes
timeout_rate: f64, // % of timeouts
event_count: u64, // Total events received
last_seen: Instant,
connection_attempts: u32,
connection_failures: u32,
}
impl RelayQuality {
pub fn update_score(&mut self) {
self.score =
(1.0 - (self.metrics.latency_ms as f64 / 1000.0)).clamp(0.0, 1.0) * 0.3 +
self.metrics.success_rate * 0.4 +
(1.0 - self.metrics.timeout_rate) * 0.3;
}
}
5. Storage Optimization¶
5.1 Event Deduplication¶
pub struct Deduplicator {
bloom: BloomFilter, // Probabilistic check (fast)
cache: LruCache<EventId, ()>, // Exact check (bounded memory)
}
impl Deduplicator {
pub fn is_new(&mut self, event_id: &EventId) -> bool {
// Fast path: Bloom filter (may have false positives)
if !self.bloom.check(event_id) {
self.bloom.add(event_id);
self.cache.put(event_id.clone(), ());
return true;
}
// Slow path: LRU cache (exact)
if self.cache.get(event_id).is_some() {
return false;
}
self.cache.put(event_id.clone(), ());
true
}
}
5.2 Pruning Strategy¶
pub struct PruningStrategy {
// Keep all from followed users (last 90 days)
followed_retention: Duration,
// Keep others for shorter time
unfollowed_retention: Duration,
// Always keep replaceable events (metadata)
keep_replaceable: bool,
// Max database size
max_size_bytes: u64,
}
impl PruningStrategy {
pub async fn prune(&self, store: &dyn EventStore) -> Result<u64, Error> {
let cutoff_followed = Utc::now() - self.followed_retention;
let cutoff_unfollowed = Utc::now() - self.unfollowed_retention;
let removed = store.delete_where(
r#"(author IN (SELECT pubkey FROM follows) AND created_at < ?1)
OR (author NOT IN (SELECT pubkey FROM follows) AND created_at < ?2)
AND kind NOT IN (0, 3, 10000, 10001, 10002)"#,
params![cutoff_followed, cutoff_unfollowed],
).await?;
Ok(removed)
}
}
6. Multi-Account Support¶
6.1 Account Isolation¶
pub struct AccountManager {
accounts: HashMap<AccountId, Account>,
active_account: AccountId,
keyring: SecureKeyring,
}
pub struct Account {
id: AccountId,
name: String,
public_key: PublicKey,
// Key stored in keyring, not here
metadata: Metadata,
relay_preferences: Vec<RelayUrl>,
isolated_store: Box<dyn EventStore>,
}
impl AccountManager {
pub async fn switch_account(
&mut self,
account_id: AccountId,
) -> Result<(), Error> {
// Save current state
self.save_account_state(self.active_account).await?;
// Switch relay pools
let new_account = self.accounts.get(&account_id)
.ok_or(Error::AccountNotFound)?;
self.relay_pool.disconnect_all().await;
self.relay_pool.connect_all(&new_account.relay_preferences).await;
// Switch event store
self.store.switch_to(&new_account.isolated_store).await;
self.active_account = account_id;
Ok(())
}
pub async fn post_as(
&self,
account_id: AccountId,
content: String,
) -> Result<Event, Error> {
let account = self.accounts.get(&account_id)
.ok_or(Error::AccountNotFound)?;
// Get key from secure keyring
let secret_key = self.keyring.get_key(&account_id).await?;
// Build and sign event
let event = EventBuilder::text_note(content)
.to_event(&secret_key)?;
// Publish
self.relay_pool.broadcast(event.clone()).await?;
Ok(event)
}
}
7. Performance Patterns¶
7.1 Lazy Loading¶
pub struct LazyLoader {
cache: LruCache<EventId, Event>,
loading: Arc<Mutex<HashSet<EventId>>>,
}
impl LazyLoader {
pub async fn get_event(&self,
id: EventId,
) -> Result<Option<Event>, Error> {
// Check cache
if let Some(event) = self.cache.get(&id) {
return Ok(Some(event.clone()));
}
// Check if already loading
let mut loading = self.loading.lock().await;
if loading.contains(&id) {
// Wait for loading to complete
drop(loading);
self.wait_for_load(&id).await
} else {
// Start loading
loading.insert(id.clone());
drop(loading);
// Query relays
let event = self.query_relays(id.clone()).await?;
// Update cache
if let Some(ref e) = event {
self.cache.put(id, e.clone());
}
loading.lock().await.remove(&id);
Ok(event)
}
}
}
7.2 Pagination¶
pub struct PaginatedQuery {
filter: Filter,
limit: usize, // Events per page
cursor: Option<Cursor>, // Last event ID/timestamp
direction: Direction, // Newer or Older
}
impl PaginatedQuery {
pub fn next_page(&self,
last_event: &Event,
) -> PaginatedQuery {
let mut new_filter = self.filter.clone();
match self.direction {
Direction::Older => {
new_filter.until = Some(last_event.created_at - 1);
}
Direction::Newer => {
new_filter.since = Some(last_event.created_at + 1);
}
}
PaginatedQuery {
filter: new_filter,
limit: self.limit,
cursor: Some(Cursor::from_event(last_event)),
direction: self.direction,
}
}
}
8. Error Handling¶
8.1 Error Types¶
pub enum NostrError {
// Network
RelayDisconnected(RelayUrl),
Timeout(Duration),
RateLimited { relay: RelayUrl, retry_after: Duration },
// Protocol
InvalidEvent(String),
InvalidSignature,
RelayRejected { code: String, message: String },
// Storage
DatabaseError(String),
MigrationFailed,
// Crypto
InvalidKeyFormat,
SigningFailed,
DecryptionFailed,
}
impl NostrError {
pub fn is_retriable(&self) -> bool {
matches!(self,
NostrError::RelayDisconnected(_) |
NostrError::Timeout(_) |
NostrError::RateLimited { .. }
)
}
pub fn retry_delay(&self) -> Duration {
match self {
NostrError::RateLimited { retry_after, .. } => *retry_after,
NostrError::Timeout(_) => Duration::from_secs(2),
_ => Duration::from_secs(5),
}
}
}
Document Version: 1.0
Architecture patterns: Clean Architecture, CQRS, Event Sourcing
For: CreatorQuetzal