Kihagyás

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

Vissza a tetejére