Sounder approach for catalog fetcher

This commit is contained in:
franzap
2026-04-06 13:31:58 -03:00
parent 5b0d5c61fd
commit bf32208f57
2 changed files with 356 additions and 170 deletions
+297 -108
View File
@@ -1,9 +1,9 @@
import 'dart:async';
import 'package:models/models.dart';
/// Max identifiers per relay subscription to stay under the relay's 100-event
/// hard cap. Each app may have several versions on the relay, so a small batch
/// keeps every app represented in the response.
const kCatalogBatchSize = 60;
/// Relay hard limit for filters in a single REQ message.
const kMaxFiltersPerReq = 50;
/// Raw catalog data returned by [fetchCatalog]. Callers decide what to do with
/// it (update UI state, compare for updates, show notifications, etc.).
@@ -27,144 +27,333 @@ class CatalogResult {
/// Fetch the latest catalog data for a set of installed app identifiers.
///
/// Used by both the foreground update poller and the background WorkManager
/// task. Pass [source] to override the default [RemoteSource] (e.g. use
/// [LocalSource] for a cheap local-only re-derivation).
/// 1. Queries local DB for all known 3063/1063 events of [installedIds].
/// 2. Builds one incremental filter per app (`since` = most recent local
/// timestamp + 1 s) and sends them in batches of [kMaxFiltersPerReq].
/// 3. Loads related Apps and Profiles in the background after new events arrive.
///
/// Pass [localOnly] = true to skip the remote phase (cheap local re-derivation).
Future<CatalogResult> fetchCatalog({
required StorageNotifier storage,
required Set<String> installedIds,
required String platform,
required String subscriptionPrefix,
Source source = const RemoteSource(relays: 'AppCatalog', stream: false),
bool localOnly = false,
}) async {
if (installedIds.isEmpty) return CatalogResult.empty;
final installableByApp = <String, Installable>{};
final allApps = <String, App>{};
// ── Phase 1: Asset-first (3063 → 32267) ────────────────────────────
final assets = await batchedQuery<SoftwareAsset>(
storage: storage,
allIds: installedIds,
tagKey: '#i',
extraTags: {'#f': {platform}},
subscriptionPrefix: '$subscriptionPrefix-assets',
source: source,
// ── Phase 1: local baseline ──────────────────────────────────────────
final localAssets = await storage.query(
RequestFilter<SoftwareAsset>(
tags: {'#i': installedIds, '#f': {platform}},
).toRequest(),
source: const LocalSource(),
subscriptionPrefix: '$subscriptionPrefix-local-assets',
);
for (final a in assets) {
final installableByApp = <String, Installable>{};
final newestTimestamp = <String, DateTime>{};
for (final a in localAssets) {
final id = a.appIdentifier;
if (id.isEmpty) continue;
final existing = installableByApp[id];
if (existing == null ||
(a.versionCode ?? 0) > (existing.versionCode ?? 0)) {
installableByApp[id] = a;
_mergeInstallable(installableByApp, id, a);
final prev = newestTimestamp[id];
if (prev == null || a.createdAt.isAfter(prev)) {
newestTimestamp[id] = a.createdAt;
}
}
final assetCoveredIds = installableByApp.keys.toSet();
final assetCoveredIds = <String>{...installableByApp.keys};
if (assetCoveredIds.isNotEmpty) {
final apps = await storage.query(
RequestFilter<App>(
tags: {'#d': assetCoveredIds, '#f': {platform}},
).toRequest(),
source: source,
subscriptionPrefix: '$subscriptionPrefix-apps',
);
for (final app in apps) {
allApps[app.identifier] = app;
}
}
// ── Phase 2: Legacy (32267 → 30063 → 1063) ────────────────────────
// TODO(cleanup): Remove when all apps are migrated to 3063.
// Legacy 1063: local baseline for apps not covered by 3063
final uncoveredIds = installedIds.difference(assetCoveredIds);
if (uncoveredIds.isNotEmpty) {
final legacyApps = await storage.query(
RequestFilter<App>(
tags: {'#d': uncoveredIds, '#f': {platform}},
).toRequest(),
source: source,
subscriptionPrefix: '$subscriptionPrefix-legacy-apps',
final legacyResult = await _localLegacyBaseline(
storage: storage,
uncoveredIds: uncoveredIds,
platform: platform,
subscriptionPrefix: '$subscriptionPrefix-local-legacy',
);
if (legacyApps.isNotEmpty) {
final releaseFilters = legacyApps
.map((a) => a.latestRelease.req?.filters.firstOrNull)
.nonNulls
.toList();
var releases = const <Release>[];
if (releaseFilters.isNotEmpty) {
releases = await storage.query(
Request<Release>(releaseFilters),
source: source,
subscriptionPrefix: '$subscriptionPrefix-legacy-releases',
);
}
if (releases.isNotEmpty) {
final metadataFilters = releases
.map((r) => r.latestMetadata.req?.filters.firstOrNull)
.nonNulls
.toList();
if (metadataFilters.isNotEmpty) {
final metadatas = await storage.query(
Request<FileMetadata>(metadataFilters),
source: source,
subscriptionPrefix: '$subscriptionPrefix-legacy-meta',
);
for (final m in metadatas) {
final id = m.appIdentifier;
if (id.isEmpty) continue;
final existing = installableByApp[id];
if (existing == null ||
(m.versionCode ?? 0) > (existing.versionCode ?? 0)) {
installableByApp[id] = m;
}
}
}
}
for (final app in legacyApps) {
allApps.putIfAbsent(app.identifier, () => app);
}
}
installableByApp.addAll(legacyResult.installableByApp);
newestTimestamp.addAll(legacyResult.timestamps);
}
return CatalogResult(
apps: allApps.values.toList(),
installableByApp: installableByApp,
catalogedIds: installableByApp.keys.toSet(),
if (localOnly) {
return _buildResult(
storage: storage,
installableByApp: installableByApp,
platform: platform,
subscriptionPrefix: subscriptionPrefix,
localOnly: true,
);
}
// ── Phase 2: incremental remote fetch (3063) ──────────────────────────
final newAssets = await _incrementalRemoteFetch<SoftwareAsset>(
storage: storage,
installedIds: installedIds,
newestTimestamp: newestTimestamp,
platform: platform,
subscriptionPrefix: '$subscriptionPrefix-assets',
);
for (final a in newAssets) {
final id = a.appIdentifier;
if (id.isEmpty) continue;
_mergeInstallable(installableByApp, id, a);
// Track 3063 coverage so we know which apps are asset-covered
assetCoveredIds.add(id);
}
// Legacy 1063: remote App→Release→FileMetadata chain for apps without
// any 3063 coverage. Uses the proper relationship chain because old 1063
// events may lack an #i tag.
final legacyIds = installedIds.difference(assetCoveredIds);
if (legacyIds.isNotEmpty) {
final legacyResult = await _remoteLegacyChain(
storage: storage,
legacyIds: legacyIds,
platform: platform,
subscriptionPrefix: '$subscriptionPrefix-legacy',
);
legacyResult.installableByApp.forEach((id, candidate) {
_mergeInstallable(installableByApp, id, candidate);
});
}
// ── Phase 3: load App events (awaited) + profiles (background) ───────
final result = await _buildResult(
storage: storage,
installableByApp: installableByApp,
platform: platform,
subscriptionPrefix: subscriptionPrefix,
);
final pubkeys = result.apps.map((a) => a.event.pubkey).toSet();
if (pubkeys.isNotEmpty) {
unawaited(
storage.query(
RequestFilter<Profile>(authors: pubkeys).toRequest(),
source: const LocalAndRemoteSource(
relays: {'social', 'vertex'},
cachedFor: Duration(hours: 2),
stream: false,
),
subscriptionPrefix: '$subscriptionPrefix-profiles-bg',
),
);
}
return result;
}
/// Query a model type in batches of [kCatalogBatchSize] identifiers to avoid
/// hitting the relay's per-subscription event cap.
Future<List<T>> batchedQuery<T extends Model<T>>({
// ═══════════════════════════════════════════════════════════════════════════════
// INTERNALS
// ═══════════════════════════════════════════════════════════════════════════════
/// Build one filter per app with `since` = newest local + 1 s, batched into
/// groups of [kMaxFiltersPerReq]. Returns only genuinely new events.
Future<List<T>> _incrementalRemoteFetch<T extends Model<T>>({
required StorageNotifier storage,
required Set<String> allIds,
required String tagKey,
Map<String, Set<String>> extraTags = const {},
required Set<String> installedIds,
required Map<String, DateTime> newestTimestamp,
required String platform,
required String subscriptionPrefix,
Source source = const RemoteSource(relays: 'AppCatalog', stream: false),
}) async {
if (allIds.isEmpty) return const [];
final idList = allIds.toList();
final filters = <RequestFilter<T>>[];
for (final appId in installedIds) {
final local = newestTimestamp[appId];
filters.add(RequestFilter<T>(
tags: {'#i': {appId}, '#f': {platform}},
since: local?.add(const Duration(seconds: 1)),
limit: 1,
));
}
final results = <T>[];
for (var i = 0; i < idList.length; i += kCatalogBatchSize) {
final batch = idList.sublist(
for (var i = 0; i < filters.length; i += kMaxFiltersPerReq) {
final batch = filters.sublist(
i,
(i + kCatalogBatchSize).clamp(0, idList.length),
(i + kMaxFiltersPerReq).clamp(0, filters.length),
);
final batchResults = await storage.query(
RequestFilter<T>(
tags: {tagKey: batch.toSet(), ...extraTags},
).toRequest(),
source: source,
Request<T>(batch),
source: const RemoteSource(relays: 'AppCatalog', stream: false),
subscriptionPrefix: '$subscriptionPrefix-$i',
);
results.addAll(batchResults);
}
return results;
}
/// Local-only legacy pass: App → Release → FileMetadata chain for apps
/// without any 3063 coverage.
Future<_LegacyBaseline> _localLegacyBaseline({
required StorageNotifier storage,
required Set<String> uncoveredIds,
required String platform,
required String subscriptionPrefix,
}) async {
final installableByApp = <String, Installable>{};
final timestamps = <String, DateTime>{};
final legacyApps = await storage.query(
RequestFilter<App>(
tags: {'#d': uncoveredIds, '#f': {platform}},
).toRequest(),
source: const LocalSource(),
subscriptionPrefix: '$subscriptionPrefix-apps',
);
if (legacyApps.isEmpty) return _LegacyBaseline(installableByApp, timestamps);
final releaseFilters = legacyApps
.map((a) => a.latestRelease.req?.filters.firstOrNull)
.nonNulls
.toList();
if (releaseFilters.isEmpty) {
return _LegacyBaseline(installableByApp, timestamps);
}
final releases = await storage.query(
Request<Release>(releaseFilters),
source: const LocalSource(),
subscriptionPrefix: '$subscriptionPrefix-releases',
);
if (releases.isEmpty) return _LegacyBaseline(installableByApp, timestamps);
final metadataFilters = releases
.map((r) => r.latestMetadata.req?.filters.firstOrNull)
.nonNulls
.toList();
if (metadataFilters.isEmpty) {
return _LegacyBaseline(installableByApp, timestamps);
}
final metadatas = await storage.query(
Request<FileMetadata>(metadataFilters),
source: const LocalSource(),
subscriptionPrefix: '$subscriptionPrefix-meta',
);
for (final m in metadatas) {
final id = m.appIdentifier;
if (id.isEmpty) continue;
_mergeInstallable(installableByApp, id, m);
final prev = timestamps[id];
if (prev == null || m.createdAt.isAfter(prev)) {
timestamps[id] = m.createdAt;
}
}
return _LegacyBaseline(installableByApp, timestamps);
}
/// Remote legacy pass: App → Release → FileMetadata chain for apps without
/// 3063 coverage. Old 1063 events may lack #i tags, so we must traverse
/// the relationship chain rather than filtering by tag directly.
Future<_LegacyBaseline> _remoteLegacyChain({
required StorageNotifier storage,
required Set<String> legacyIds,
required String platform,
required String subscriptionPrefix,
}) async {
const source = RemoteSource(relays: 'AppCatalog', stream: false);
final installableByApp = <String, Installable>{};
final legacyApps = await storage.query(
RequestFilter<App>(
tags: {'#d': legacyIds, '#f': {platform}},
).toRequest(),
source: source,
subscriptionPrefix: '$subscriptionPrefix-apps',
);
if (legacyApps.isEmpty) return _LegacyBaseline(installableByApp, const {});
final releaseFilters = legacyApps
.map((a) => a.latestRelease.req?.filters.firstOrNull)
.nonNulls
.toList();
if (releaseFilters.isEmpty) {
return _LegacyBaseline(installableByApp, const {});
}
final releases = await storage.query(
Request<Release>(releaseFilters),
source: source,
subscriptionPrefix: '$subscriptionPrefix-releases',
);
if (releases.isEmpty) return _LegacyBaseline(installableByApp, const {});
final metadataFilters = releases
.map((r) => r.latestMetadata.req?.filters.firstOrNull)
.nonNulls
.toList();
if (metadataFilters.isEmpty) {
return _LegacyBaseline(installableByApp, const {});
}
final metadatas = await storage.query(
Request<FileMetadata>(metadataFilters),
source: source,
subscriptionPrefix: '$subscriptionPrefix-meta',
);
for (final m in metadatas) {
final id = m.appIdentifier;
if (id.isEmpty) continue;
_mergeInstallable(installableByApp, id, m);
}
return _LegacyBaseline(installableByApp, const {});
}
/// Resolve App objects for all cataloged identifiers.
///
/// Uses [LocalAndRemoteSource] so Apps not yet in the local DB are fetched
/// from the relay. Pass [localOnly] to skip the network round-trip.
Future<CatalogResult> _buildResult({
required StorageNotifier storage,
required Map<String, Installable> installableByApp,
required String platform,
required String subscriptionPrefix,
bool localOnly = false,
}) async {
if (installableByApp.isEmpty) return CatalogResult.empty;
final source = localOnly
? const LocalSource() as Source
: const LocalAndRemoteSource(relays: 'AppCatalog', stream: false);
final apps = await storage.query(
RequestFilter<App>(
tags: {'#d': installableByApp.keys.toSet(), '#f': {platform}},
).toRequest(),
source: source,
subscriptionPrefix: '$subscriptionPrefix-resolve-apps',
);
return CatalogResult(
apps: apps,
installableByApp: installableByApp,
catalogedIds: installableByApp.keys.toSet(),
);
}
void _mergeInstallable(
Map<String, Installable> map,
String id,
Installable candidate,
) {
final existing = map[id];
if (existing == null ||
(candidate.versionCode ?? 0) > (existing.versionCode ?? 0)) {
map[id] = candidate;
}
}
class _LegacyBaseline {
const _LegacyBaseline(this.installableByApp, this.timestamps);
final Map<String, Installable> installableByApp;
final Map<String, DateTime> timestamps;
}
+59 -62
View File
@@ -6,6 +6,7 @@ import 'package:models/models.dart';
import 'package:zapstore/main.dart';
import 'package:zapstore/services/catalog_fetcher.dart';
import 'package:zapstore/services/package_manager/package_manager.dart';
import 'package:zapstore/utils/extensions.dart';
/// How often to poll for updates from remote relays
const _pollInterval = Duration(minutes: 5);
@@ -50,8 +51,6 @@ class UpdatePollerState {
this.isChecking = false,
this.lastCheckTime,
this.lastError,
this.apps = const [],
this.installableByApp = const {},
this.catalogedIds = const {},
});
@@ -59,14 +58,8 @@ class UpdatePollerState {
final DateTime? lastCheckTime;
final String? lastError;
/// App objects fetched from relay, used for display (name, icon, author)
final List<App> apps;
/// appIdentifier → Installable (3063 or 1063), used for update comparison
final Map<String, Installable> installableByApp;
/// All app identifiers found in relay catalog (superset of apps list,
/// since an installable may exist without a corresponding App object)
/// App identifiers found in the relay catalog. The categorizer uses this
/// to know which installed apps to query (with relationships) from local DB.
final Set<String> catalogedIds;
UpdatePollerState copyWith({
@@ -74,16 +67,12 @@ class UpdatePollerState {
DateTime? lastCheckTime,
String? lastError,
bool clearError = false,
List<App>? apps,
Map<String, Installable>? installableByApp,
Set<String>? catalogedIds,
}) {
return UpdatePollerState(
isChecking: isChecking ?? this.isChecking,
lastCheckTime: lastCheckTime ?? this.lastCheckTime,
lastError: clearError ? null : (lastError ?? this.lastError),
apps: apps ?? this.apps,
installableByApp: installableByApp ?? this.installableByApp,
catalogedIds: catalogedIds ?? this.catalogedIds,
);
}
@@ -144,50 +133,27 @@ class UpdatePollerNotifier extends StateNotifier<UpdatePollerState> {
}
}
/// Fetch catalog data from relays and store in state.
/// Categorization is done by [categorizedUpdatesProvider].
/// Fetch catalog data from relays and store in local DB.
/// The poller only retains [catalogedIds]; the categorizer reactively
/// queries Apps with relationships from local cache.
Future<void> _fetchCatalog() async {
final pmState = ref.read(packageManagerProvider);
if (pmState.installed.isEmpty) {
state = state.copyWith(
apps: const [],
installableByApp: const {},
catalogedIds: const {},
);
state = state.copyWith(catalogedIds: const {});
return;
}
final storage = ref.read(storageNotifierProvider.notifier);
final result = await fetchCatalog(
storage: storage,
storage: ref.read(storageNotifierProvider.notifier),
installedIds: pmState.installed.keys.toSet(),
platform: ref.read(packageManagerProvider.notifier).platform,
subscriptionPrefix: 'app-updates-poll',
);
final authorPubkeys = result.apps.map((a) => a.event.pubkey).toSet();
if (authorPubkeys.isNotEmpty) {
unawaited(
storage.query(
RequestFilter<Profile>(authors: authorPubkeys).toRequest(),
source: const LocalAndRemoteSource(
relays: {'social', 'vertex'},
cachedFor: Duration(hours: 2),
stream: false,
),
subscriptionPrefix: 'app-updates-profiles',
),
);
}
state = state.copyWith(
apps: result.apps,
installableByApp: result.installableByApp,
catalogedIds: result.catalogedIds,
);
state = state.copyWith(catalogedIds: result.catalogedIds);
}
/// Re-derive catalog from local DB without hitting relays.
/// Re-derive catalog IDs from local DB without hitting relays.
/// Call when returning to the updates screen so that data written by
/// other code paths (detail screen, background service) is picked up
/// without waiting for the next poll cycle.
@@ -195,20 +161,15 @@ class UpdatePollerNotifier extends StateNotifier<UpdatePollerState> {
final pmState = ref.read(packageManagerProvider);
if (pmState.installed.isEmpty) return;
final storage = ref.read(storageNotifierProvider.notifier);
final result = await fetchCatalog(
storage: storage,
storage: ref.read(storageNotifierProvider.notifier),
installedIds: pmState.installed.keys.toSet(),
platform: ref.read(packageManagerProvider.notifier).platform,
subscriptionPrefix: 'app-updates-local',
source: const LocalSource(),
localOnly: true,
);
state = state.copyWith(
apps: result.apps,
installableByApp: result.installableByApp,
catalogedIds: result.catalogedIds,
);
state = state.copyWith(catalogedIds: result.catalogedIds);
}
@override
@@ -227,8 +188,9 @@ final updatePollerProvider =
// CATEGORIZED UPDATES PROVIDER
// ═══════════════════════════════════════════════════════════════════════════════
/// Pure synchronous derivation from poller catalog + installed packages.
/// Rebuilds when either changes.
/// Derives update categories from a reactive local query that loads each App
/// with its installable relationships. This is the same data path used by
/// the detail screen, install button, and version pill — one source of truth.
final categorizedUpdatesProvider = Provider<CategorizedUpdates>((ref) {
final pollerState = ref.watch(updatePollerProvider);
final installed = ref.watch(
@@ -248,23 +210,58 @@ final categorizedUpdatesProvider = Provider<CategorizedUpdates>((ref) {
);
}
final pm = ref.read(packageManagerProvider.notifier);
final installableByApp = pollerState.installableByApp;
final catalogedIds = pollerState.catalogedIds;
final catalogedInstalledIds =
installed.keys.where(catalogedIds.contains).toSet();
if (catalogedInstalledIds.isEmpty) {
return CategorizedUpdates(
automaticUpdates: const [],
manualUpdates: const [],
upToDateApps: const [],
uncatalogedApps: installed.values.toList()
..sort(
(a, b) => (a.name ?? a.appId)
.toLowerCase()
.compareTo((b.name ?? b.appId).toLowerCase()),
),
);
}
final platform = ref.read(packageManagerProvider.notifier).platform;
// Reactive local query: loads Apps with their installable relationships.
// The catalog fetcher already wrote SoftwareAssets / FileMetadatas to the
// local DB; this query picks them up via the model relationship graph —
// the same path used by the detail screen and install button.
final appsState = ref.watch(
query<App>(
tags: {'#d': catalogedInstalledIds, '#f': {platform}},
and: (app) => {
app.latestAsset.query(source: const LocalSource()),
app.latestRelease.query(
source: const LocalSource(),
and: (release) => {
release.latestMetadata.query(source: const LocalSource()),
},
),
},
source: const LocalSource(),
subscriptionPrefix: 'app-updates-categorize',
),
);
final apps = appsState.models;
final automaticUpdates = <App>[];
final manualUpdates = <App>[];
final upToDateApps = <App>[];
for (final app in pollerState.apps) {
for (final app in apps) {
final pkg = installed[app.identifier];
if (pkg == null) continue;
final installable = installableByApp[app.identifier];
final hasUpdate =
installable != null && pm.hasUpdate(app.identifier, installable);
if (hasUpdate) {
if (app.hasUpdate) {
if (pkg.canInstallSilently) {
automaticUpdates.add(app);
} else {