import 'dart:async'; import 'package:commons_core/commons_core.dart'; import 'package:connectivity_plus/connectivity_plus.dart'; import 'package:flutter/foundation.dart'; import 'message_store.dart'; import 'profile_cache.dart'; import 'social_service.dart'; import 'social_settings.dart'; /// App-wide inbox listener for private messages (NIP-17). /// /// Without this, a message only arrived while its specific chat screen was /// open — so a first message from a new peer was invisible (nothing in local /// storage, so nothing in the inbox list, so you never opened the chat, so you /// never subscribed). This keeps ONE long-lived subscription for the whole app /// while online: it persists every incoming message to the [MessageStore] and /// fires [changes] so the inbox list refreshes. Reconnects when the network /// returns; degrades to nothing offline (local-first). /// /// Foreground only — background/push delivery is a later concern. class InboxService { InboxService({ required SocialService social, required SocialSettings settings, required MessageStore store, ProfileCache? profileCache, }) : _social = social, _settings = settings, _store = store, _profileCache = profileCache; final SocialService _social; final SocialSettings _settings; final MessageStore _store; final ProfileCache? _profileCache; final _changes = StreamController.broadcast(); SocialSession? _session; StreamSubscription? _inboxSub; StreamSubscription>? _connSub; bool _started = false; bool _connecting = false; /// Fires (with no payload) after a NEW message is persisted — the inbox list /// listens to this to reload. Broadcast, so several screens can listen. Stream get changes => _changes.stream; /// Begins listening. Idempotent. Connects now if online, and (re)connects /// whenever the network returns. Safe to call once at app start. Future start() async { if (_started) return; _started = true; try { _connSub = Connectivity().onConnectivityChanged.listen((results) { final offline = results.isEmpty || results.every((r) => r == ConnectivityResult.none); if (offline) { _dropSession(); } else { unawaited(_connect()); } }); } catch (_) { // Platform without connectivity support — just try once below. } await _connect(); } /// Opens a session and subscribes to the inbox, if not already connected. /// Guarded so overlapping connectivity events can't open two sessions. Future _connect() async { if (_session != null || _connecting) return; // already (being) connected _connecting = true; try { final relays = await _settings.relayUrls(); if (relays.isEmpty) return; // offline / unconfigured final session = await _social.openSession(relays); if (!_started) { await session.close(); // stopped while connecting return; } _session = session; _inboxSub = session.messages.inbox().listen( ingest, onError: (_) => _dropSession(), // relay dropped — retry on next change ); } catch (_) { // No relay reachable — a later connectivity change retries. } finally { _connecting = false; } } /// Persists one incoming [message] (deduped) and, when it's new, fires /// [changes] and best-effort caches the sender's display name. Separated from /// the transport so it can be unit-tested without a relay. @visibleForTesting Future ingest(PrivateMessage message) async { final stored = await _store.append(message.fromPubkey, message); if (!stored) return; // a re-delivered duplicate if (!_changes.isClosed) _changes.add(null); await _cacheName(message.fromPubkey); } Future _cacheName(String peerPubkey) async { final cache = _profileCache; final session = _session; if (cache == null || session == null) return; if (await cache.name(peerPubkey) != null) return; // already known try { final profile = await session.profile.fetch(peerPubkey); if (profile != null && profile.name.isNotEmpty) { await cache.setName(peerPubkey, profile.name); if (!_changes.isClosed) _changes.add(null); // re-render with the name } } catch (_) { // best effort — the short key shows meanwhile } } void _dropSession() { unawaited(_inboxSub?.cancel()); _inboxSub = null; unawaited(_session?.close()); _session = null; } Future stop() async { _started = false; await _connSub?.cancel(); _connSub = null; await _inboxSub?.cancel(); _inboxSub = null; await _session?.close(); _session = null; if (!_changes.isClosed) await _changes.close(); } }