diff --git a/convex/function_spec.json b/convex/function_spec.json index 5db7163d..84f4a898 100644 --- a/convex/function_spec.json +++ b/convex/function_spec.json @@ -40660,6 +40660,47 @@ "kind": "internal" } }, + { + "args": { + "type": "object", + "value": { + "assetPublicIds": { + "fieldType": { + "type": "array", + "value": { + "type": "string" + } + }, + "optional": false + }, + "strategyPublicId": { + "fieldType": { + "type": "string" + }, + "optional": false + } + } + }, + "functionType": "Query", + "identifier": "images.js:listReferencedAssetIds", + "returns": { + "type": "union", + "value": [ + { + "type": "array", + "value": { + "type": "string" + } + }, + { + "type": "null" + } + ] + }, + "visibility": { + "kind": "public" + } + }, { "args": { "type": "object", @@ -182488,6 +182529,12 @@ "args": { "type": "object", "value": { + "acceptsTrashedPagesLeftOut": { + "fieldType": { + "type": "boolean" + }, + "optional": true + }, "shareToken": { "fieldType": { "type": "string" diff --git a/convex/images.ts b/convex/images.ts index d967fd62..71782b24 100644 --- a/convex/images.ts +++ b/convex/images.ts @@ -747,6 +747,45 @@ export const listForStrategy = query({ }, }); +/// At most this many images are asked about at once, so the answer reads a +/// bounded number of rows however much content the strategy has. +export const MAX_REFERENCED_ASSET_IDS_PER_QUERY = 100; + +/// Which of [assetPublicIds] the strategy's content shows, deleted content +/// left out. A client asks before it drops an upload it holds. Null until +/// the reference backfill has finished, when it cannot be told. +export const listReferencedAssetIds = query({ + args: { + strategyPublicId: v.string(), + assetPublicIds: v.array(v.string()), + }, + returns: v.union(v.array(v.string()), v.null()), + handler: async (ctx, args) => { + if (args.assetPublicIds.length > MAX_REFERENCED_ASSET_IDS_PER_QUERY) { + throw invalidPayloadError( + `Ask about at most ${MAX_REFERENCED_ASSET_IDS_PER_QUERY} images at once.`, + ); + } + const strategy = await getStrategyByPublicId(ctx, args.strategyPublicId); + await assertStrategyRole(ctx, strategy, "viewer"); + if (!(await assetReferencesReady(ctx))) return null; + const referenced: string[] = []; + for (const assetPublicId of new Set(args.assetPublicIds)) { + const reference = await ctx.db + .query("assetReferences") + .withIndex("by_strategyId_and_assetPublicId_and_deleted", (q) => + q + .eq("strategyId", strategy._id) + .eq("assetPublicId", assetPublicId) + .eq("deleted", false), + ) + .first(); + if (reference !== null) referenced.push(assetPublicId); + } + return referenced.sort(); + }, +}); + export const getAssetUrl = query({ args: { strategyPublicId: v.string(), diff --git a/convex/referencedAssetIds.test.ts b/convex/referencedAssetIds.test.ts new file mode 100644 index 00000000..6fa40ad1 --- /dev/null +++ b/convex/referencedAssetIds.test.ts @@ -0,0 +1,211 @@ +import { + convexTest, + type TestConvexForDataModel, + type TestConvexForDataModelAndIdentity, +} from "convex-test"; +import { makeFunctionReference } from "convex/server"; +import { describe, expect, test } from "vitest"; +import type { DataModel } from "./_generated/dataModel"; +import { markAssetReferencesReady } from "./lib/assetReferences"; +import { CURRENT_CLOUD_PROTOCOL_VERSION } from "./lib/cloudProtocol"; +import schema from "./schema"; +import { modules } from "./test.setup"; + +const ensureCurrentUser = makeFunctionReference<"mutation">( + "users:ensureCurrentUser", +); +const createStrategy = makeFunctionReference<"mutation">( + "strategies:createWithInitialPage", +); +const applyBatch = makeFunctionReference<"mutation">("ops:applyBatch"); +const createShare = makeFunctionReference<"mutation">("shares:create"); +const redeemShare = makeFunctionReference<"mutation">("shares:redeem"); +const getFullSnapshot = makeFunctionReference<"query">( + "strategy:getFullSnapshot", +); +const listReferencedAssetIds = makeFunctionReference<"query">( + "images:listReferencedAssetIds", +); + +type Harness = TestConvexForDataModel; +type RootHarness = TestConvexForDataModelAndIdentity; + +const protocol = { clientProtocolVersion: CURRENT_CLOUD_PROTOCOL_VERSION }; +const strategyPublicId = "referenced-assets-strategy"; +const pagePublicId = "referenced-assets-page"; +// Every image the tests' content ever showed, and one it never did. +const asked = { + strategyPublicId, + assetPublicIds: ["removed", "shown", "shot", "never-placed"], +}; + +function identity(subject: string) { + return { + issuer: "https://referenced-assets.test", + subject, + tokenIdentifier: `referenced-assets|${subject}`, + name: subject, + }; +} + +function imagePayload(assetPublicId: string) { + return { + kind: "image" as const, + payloadVersion: 1, + data: { id: assetPublicId, elementType: "image" }, + }; +} + +async function user(t: RootHarness, subject: string): Promise { + const harness = t.withIdentity(identity(subject)); + await harness.mutation(ensureCurrentUser, protocol); + return harness; +} + +let opCounter = 0; + +async function apply(owner: Harness, ops: Array>) { + const result = (await owner.mutation(applyBatch, { + ...protocol, + strategyPublicId, + clientId: "client", + ops: ops.map((op) => ({ opId: `op-${++opCounter}`, ...op })), + })) as { results: Array<{ status: string }> }; + expect(result.results.map((entry) => entry.status)).toEqual( + ops.map(() => "applied"), + ); +} + +/// A strategy whose page shows 'shown' as an image and 'shot' in a lineup, +/// and once showed 'removed'. +async function seed(t: RootHarness): Promise { + const owner = await user(t, "owner"); + await owner.mutation(createStrategy, { + ...protocol, + publicId: strategyPublicId, + name: "Referenced assets", + mapData: "ascent", + initialPagePublicId: pagePublicId, + initialPageName: "Page 1", + initialPageIsAttack: true, + }); + await apply(owner, [ + { + type: "element.add", + elementPublicId: "shown", + pagePublicId, + payload: imagePayload("shown"), + sortIndex: 0, + }, + { + type: "element.add", + elementPublicId: "removed", + pagePublicId, + payload: imagePayload("removed"), + sortIndex: 1, + }, + { + type: "lineup.add", + lineupPublicId: "lineupLink:k", + pagePublicId, + payload: { + kind: "lineupLink" as const, + payloadVersion: 1, + data: { + id: "k", + originId: "o", + landingId: "l", + images: [{ id: "shot" }], + }, + }, + sortIndex: 0, + }, + ]); + await apply(owner, [ + { + type: "element.delete", + elementPublicId: "removed", + pagePublicId, + expectedElementRevision: 1, + }, + ]); + return owner; +} + +describe("images:listReferencedAssetIds", () => { + test("answers which of the images asked about the strategy's content shows, deleted content left out", async () => { + const t = convexTest(schema, modules); + await t.run(markAssetReferencesReady); + const owner = await seed(t); + + expect( + await owner.query(listReferencedAssetIds, asked), + ).toEqual(["shot", "shown"]); + }); + + test("reads a bounded number of rows: it answers at most 100 images at once", async () => { + const t = convexTest(schema, modules); + await t.run(markAssetReferencesReady); + const owner = await seed(t); + const many = Array.from({ length: 100 }, (_, index) => `image-${index}`); + + expect( + await owner.query(listReferencedAssetIds, { + strategyPublicId, + assetPublicIds: [...many.slice(1), "shown"], + }), + ).toEqual(["shown"]); + await expect( + owner.query(listReferencedAssetIds, { + strategyPublicId, + assetPublicIds: [...many, "shown"], + }), + ).rejects.toThrow(/at most 100/); + }); + + test("is null until the reference backfill has finished", async () => { + const t = convexTest(schema, modules); + const owner = await seed(t); + + expect( + await owner.query(listReferencedAssetIds, asked), + ).toBeNull(); + }); + + test("answers anyone who can view the strategy, and no one else", async () => { + const t = convexTest(schema, modules); + await t.run(markAssetReferencesReady); + const owner = await seed(t); + const viewer = await user(t, "viewer"); + await owner.mutation(createShare, { + ...protocol, + targetType: "strategy", + targetPublicId: strategyPublicId, + token: "viewer-token", + role: "viewer", + }); + await viewer.mutation(redeemShare, { ...protocol, token: "viewer-token" }); + const stranger = await user(t, "stranger"); + + expect( + await viewer.query(listReferencedAssetIds, asked), + ).toEqual(["shot", "shown"]); + await expect( + stranger.query(listReferencedAssetIds, asked), + ).rejects.toThrow(); + }); + + test("a client that checks references apart says so when it reads a full snapshot", async () => { + const t = convexTest(schema, modules); + await t.run(markAssetReferencesReady); + const owner = await seed(t); + + const snapshot = (await owner.query(getFullSnapshot, { + strategyPublicId, + acceptsTrashedPagesLeftOut: true, + })) as { pages: Array<{ publicId: string }> }; + expect(snapshot.pages.map((page) => page.publicId)).toEqual([ + pagePublicId, + ]); + }); +}); diff --git a/convex/strategy.ts b/convex/strategy.ts index 39039281..32df99f7 100644 --- a/convex/strategy.ts +++ b/convex/strategy.ts @@ -64,6 +64,11 @@ export const getFullSnapshot = query({ args: { strategyPublicId: v.string(), shareToken: v.optional(v.string()), + // Set by clients that ask images:listReferencedAssetIds, not this + // snapshot, whether an upload is still wanted: a snapshot that leaves + // deleted pages out cannot make them drop one. Ignored until pages can + // be deleted into a trash. + acceptsTrashedPagesLeftOut: v.optional(v.boolean()), }, returns: fullStrategySnapshotValidator, handler: async (ctx, args) => { diff --git a/lib/collab/convex_strategy_repository.dart b/lib/collab/convex_strategy_repository.dart index 96629b51..e7ccd5c4 100644 --- a/lib/collab/convex_strategy_repository.dart +++ b/lib/collab/convex_strategy_repository.dart @@ -1,3 +1,5 @@ +import 'dart:math' show min; + import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:icarus/collab/cloud_media_models.dart'; import 'package:icarus/collab/cloud_library_models.dart'; @@ -157,6 +159,33 @@ class ConvexStrategyRepository { .map(_pageSnapshot); } + /// Which of [assetPublicIds] the strategy's content shows on the server, + /// or null while the server cannot tell yet. Asked a batch at a time. + Future?> fetchReferencedAssetIds( + String strategyPublicId, + Iterable assetPublicIds, + ) async { + final asked = assetPublicIds.toSet().toList(growable: false); + final referenced = {}; + for (var start = 0; start < asked.length; start += _referencedIdsBatch) { + final ids = await _api.images + .listReferencedAssetIds( + strategyPublicId: strategyPublicId, + assetPublicIds: asked.sublist( + start, + min(start + _referencedIdsBatch, asked.length), + ), + ) + .fetch(); + if (ids == null) return null; + referenced.addAll(ids); + } + return referenced; + } + + /// The server's MAX_REFERENCED_ASSET_IDS_PER_QUERY. + static const _referencedIdsBatch = 100; + /// [shareToken]: see [fetchShell]. Future fetchFullSnapshot( String strategyPublicId, { @@ -167,6 +196,9 @@ class ConvexStrategyRepository { .getFullSnapshot( strategyPublicId: strategyPublicId, shareToken: _optional(shareToken), + // This client checks image references apart, so it can take a + // snapshot without the pages in the server's trash. + acceptsTrashedPagesLeftOut: const ConvexOptional.present(true), ) .fetch(), ); diff --git a/lib/collab/generated/convex_models.dart b/lib/collab/generated/convex_models.dart index cbb2fcb7..c9ca04dd 100644 --- a/lib/collab/generated/convex_models.dart +++ b/lib/collab/generated/convex_models.dart @@ -5244,6 +5244,32 @@ List decodeImagesListForStrategyResult( ) .toList(growable: false); +ConvexObject encodeImagesListReferencedAssetIdsArgs({ + required List assetPublicIds, + required String strategyPublicId, +}) => ConvexObject({ + 'assetPublicIds': ConvexArray( + assetPublicIds.indexed + .map((entry) => ConvexString(entry.$2)) + .toList(growable: false), + ), + 'strategyPublicId': ConvexString(strategyPublicId), +}); + +List? decodeImagesListReferencedAssetIdsResult(ConvexValue value) => + (value) is ConvexNull + ? null + : _decodeArray(value, 'images.js:listReferencedAssetIds.returns') + .value + .indexed + .map( + (entry) => _decodeString( + entry.$2, + _indexPath('images.js:listReferencedAssetIds.returns', entry.$1), + ), + ) + .toList(growable: false); + ConvexObject encodeInvitesCreateArgs({ required double clientProtocolVersion, ConvexOptional expiresAt = const ConvexOptional.absent(), @@ -5842,9 +5868,15 @@ ConvexValue decodeStrategiesUpdateResult(ConvexValue value) => _decodeRaw(value, 'strategies.js:update.returns', _validatePagesAddResult); ConvexObject encodeStrategyGetFullSnapshotArgs({ + ConvexOptional acceptsTrashedPagesLeftOut = + const ConvexOptional.absent(), ConvexOptional shareToken = const ConvexOptional.absent(), required String strategyPublicId, }) => ConvexObject({ + if (acceptsTrashedPagesLeftOut.isPresent) + 'acceptsTrashedPagesLeftOut': ConvexBoolean( + acceptsTrashedPagesLeftOut.value, + ), if (shareToken.isPresent) 'shareToken': ConvexString(shareToken.value), 'strategyPublicId': ConvexString(strategyPublicId), }); diff --git a/lib/collab/generated/icarus_convex_api.dart b/lib/collab/generated/icarus_convex_api.dart index c7cea8fa..f2be0f19 100644 --- a/lib/collab/generated/icarus_convex_api.dart +++ b/lib/collab/generated/icarus_convex_api.dart @@ -417,6 +417,10 @@ abstract interface class ImagesModule { ConvexQuery> listForStrategy({ required String strategyPublicId, }); + ConvexQuery?> listReferencedAssetIds({ + required List assetPublicIds, + required String strategyPublicId, + }); } final class _ImagesModule implements ImagesModule { @@ -537,6 +541,23 @@ final class _ImagesModule implements ImagesModule { decode: decodeImagesListForStrategyResult, ); } + + @override + ConvexQuery?> listReferencedAssetIds({ + required List assetPublicIds, + required String strategyPublicId, + }) { + final args = encodeImagesListReferencedAssetIdsArgs( + assetPublicIds: assetPublicIds, + strategyPublicId: strategyPublicId, + ); + return ConvexQuery( + transport: _transport, + name: 'images:listReferencedAssetIds', + args: args, + decode: decodeImagesListReferencedAssetIdsResult, + ); + } } abstract interface class InvitesModule { @@ -1327,6 +1348,8 @@ final class _StrategiesModule implements StrategiesModule { abstract interface class StrategyModule { ConvexQuery getFullSnapshot({ + ConvexOptional acceptsTrashedPagesLeftOut = + const ConvexOptional.absent(), ConvexOptional shareToken = const ConvexOptional.absent(), required String strategyPublicId, }); @@ -1341,10 +1364,13 @@ final class _StrategyModule implements StrategyModule { final ConvexTransport _transport; @override ConvexQuery getFullSnapshot({ + ConvexOptional acceptsTrashedPagesLeftOut = + const ConvexOptional.absent(), ConvexOptional shareToken = const ConvexOptional.absent(), required String strategyPublicId, }) { final args = encodeStrategyGetFullSnapshotArgs( + acceptsTrashedPagesLeftOut: acceptsTrashedPagesLeftOut, shareToken: shareToken, strategyPublicId: strategyPublicId, ); diff --git a/lib/providers/collab/cloud_media_upload_queue_provider.dart b/lib/providers/collab/cloud_media_upload_queue_provider.dart index c70fffc3..98d2993c 100644 --- a/lib/providers/collab/cloud_media_upload_queue_provider.dart +++ b/lib/providers/collab/cloud_media_upload_queue_provider.dart @@ -97,12 +97,16 @@ final cloudMediaUploadQueueProvider = CloudMediaUploadQueueNotifier.new, ); -typedef CloudMediaReferenceSnapshotLoader = Future - Function(String strategyPublicId); +/// Which of the images asked about a strategy's content shows on the +/// server, on any of its pages, or null while the server cannot tell. Asked +/// before a pending upload is dropped as no longer wanted. +typedef CloudMediaReferenceLoader = Future?> Function( + String strategyPublicId, + Iterable assetPublicIds, +); -final cloudMediaReferenceSnapshotLoaderProvider = - Provider( - (ref) => ref.watch(convexStrategyRepositoryProvider).fetchFullSnapshot, +final cloudMediaReferenceLoaderProvider = Provider( + (ref) => ref.watch(convexStrategyRepositoryProvider).fetchReferencedAssetIds, ); final cloudMediaAccountIdProvider = Provider( @@ -370,22 +374,21 @@ class CloudMediaUploadQueueNotifier ) .toList(growable: false); if (missingReferences.isNotEmpty) { - RemoteFullStrategySnapshot? serverSnapshot; + Set? serverReferences; if (ref.read(authProvider).isConvexUserReady && ref.read(convexConnectionSnapshotProvider)) { try { - serverSnapshot = - await ref.read(cloudMediaReferenceSnapshotLoaderProvider)( + serverReferences = await ref.read(cloudMediaReferenceLoaderProvider)( strategyPublicId, + [for (final job in missingReferences) job.assetPublicId], ); } catch (_) { - serverSnapshot = null; + serverReferences = null; } } - if (serverSnapshot == null || + if (serverReferences == null || missingReferences.any( - (job) => - !_snapshotReferencesAsset(serverSnapshot!, job.assetPublicId), + (job) => !serverReferences!.contains(job.assetPublicId), )) { _scheduleRetryForNextEligibleJob( minimumDelay: _blockedRetryDelay, @@ -845,10 +848,11 @@ class CloudMediaUploadQueueNotifier return null; } - late final RemoteFullStrategySnapshot snapshot; + final Set? serverReferences; try { - snapshot = await ref.read(cloudMediaReferenceSnapshotLoaderProvider)( + serverReferences = await ref.read(cloudMediaReferenceLoaderProvider)( job.strategyPublicId, + [job.assetPublicId], ); } catch (error) { _logMedia( @@ -858,12 +862,13 @@ class CloudMediaUploadQueueNotifier return null; } final current = _getJob(job.jobId); - if (!_belongsToActiveAccount(job) || + if (serverReferences == null || + !_belongsToActiveAccount(job) || current == null || current.updatedAt != job.updatedAt) { return null; } - if (_snapshotReferencesAsset(snapshot, job.assetPublicId)) return false; + if (serverReferences.contains(job.assetPublicId)) return false; final deleted = await _deleteJob(job, onlyIfUnreferenced: true); if (!deleted) return null; @@ -1136,16 +1141,19 @@ class CloudMediaUploadQueueNotifier (byStrategy[job.strategyPublicId] ??= []).add(job); } for (final entry in byStrategy.entries) { - late final RemoteFullStrategySnapshot snapshot; + final Set? serverReferences; try { - snapshot = await ref - .read(cloudMediaReferenceSnapshotLoaderProvider)(entry.key); + serverReferences = await ref.read(cloudMediaReferenceLoaderProvider)( + entry.key, + [for (final job in entry.value) job.assetPublicId], + ); } catch (error) { _logMedia( 'reference_reconcile.deferred strategy=${entry.key} error=$error', ); continue; } + if (serverReferences == null) continue; if (ref.read(cloudMediaAccountIdProvider) != accountId) return; @@ -1156,7 +1164,7 @@ class CloudMediaUploadQueueNotifier if (!identical(_getJob(job.jobId), job)) continue; final key = durableCloudMediaOutboxStorageKey(job); final localReference = _hasLocalReference(job); - if (_snapshotReferencesAsset(snapshot, job.assetPublicId) || + if (serverReferences.contains(job.assetPublicId) || localReference == true) { final promoted = job.copyWith( referenceDurable: true, @@ -1176,32 +1184,6 @@ class CloudMediaUploadQueueNotifier } } - bool _snapshotReferencesAsset( - RemoteFullStrategySnapshot snapshot, - String assetPublicId, - ) { - for (final elements in snapshot.elementsByPage.values) { - if (elements.any( - (element) => - !element.deleted && - element.elementType == 'image' && - element.publicId == assetPublicId, - )) { - return true; - } - } - for (final lineups in snapshot.lineupsByPage.values) { - if (lineups.any( - (lineup) => - !lineup.deleted && - _jsonContainsAssetId(lineup.payload, assetPublicId), - )) { - return true; - } - } - return false; - } - bool _opReferencesAsset(StrategyOp op, String assetPublicId) { if (op is ElementAddOp) { return op.elementPublicId == assetPublicId; diff --git a/test/cloud_media_upload_queue_provider_test.dart b/test/cloud_media_upload_queue_provider_test.dart index 334e0b08..d5dfba0d 100644 --- a/test/cloud_media_upload_queue_provider_test.dart +++ b/test/cloud_media_upload_queue_provider_test.dart @@ -207,7 +207,7 @@ ProviderContainer _container( bool cloudReady = false, bool cloudEnabled = false, String? accountId = 'account-a', - CloudMediaReferenceSnapshotLoader? referenceSnapshotLoader, + CloudMediaReferenceLoader? referenceLoader, ConvexStrategyRepository? repository, }) { final resolvedOpQueueState = opQueueState ?? @@ -246,9 +246,9 @@ ProviderContainer _container( strategyOpQueueProvider.overrideWith( () => _FixedOpQueue(resolvedOpQueueState), ), - if (referenceSnapshotLoader != null) - cloudMediaReferenceSnapshotLoaderProvider.overrideWithValue( - referenceSnapshotLoader, + if (referenceLoader != null) + cloudMediaReferenceLoaderProvider.overrideWithValue( + referenceLoader, ), ], ); @@ -267,27 +267,6 @@ DurableOutboxRecord _durableRecord(StrategyOp op) { ); } -RemoteFullStrategySnapshot _fullSnapshot({ - List elements = const [], - List lineups = const [], -}) { - final now = DateTime.utc(2026, 9, 3); - return RemoteFullStrategySnapshot( - header: RemoteStrategyHeader( - publicId: 'strategy-a', - name: 'Strategy A', - mapData: 'ascent', - revision: 1, - createdAt: now, - updatedAt: now, - ), - pages: const [], - elementsByPage: {'page-a': elements}, - lineupsByPage: {'page-a': lineups}, - assetsById: const {}, - ); -} - void main() { TestWidgetsFlutterBinding.ensureInitialized(); @@ -470,13 +449,13 @@ void main() { updatedAt: DateTime.utc(2026, 9, 4), ); await store.put(job); - final snapshot = Completer(); + final snapshot = Completer?>(); final started = Completer(); final container = _container( store, strategyStore: ops, cloudReady: true, - referenceSnapshotLoader: (_) { + referenceLoader: (_, __) { if (!started.isCompleted) started.complete(); return snapshot.future; }, @@ -510,7 +489,7 @@ void main() { createdAt: DateTime.now(), updatedAt: DateTime.now(), )); - snapshot.complete(_fullSnapshot()); + snapshot.complete({}); await queue.retryNow(); expect(store.load().jobs.single.referenceDurable, isTrue); expect(store.load().jobs.single.width, restage ? 42 : null); @@ -543,9 +522,9 @@ void main() { final container = _container( store, cloudReady: true, - referenceSnapshotLoader: (_) async { + referenceLoader: (_, __) async { if (!snapshotAvailable) throw StateError('offline'); - return _fullSnapshot(); + return {}; }, ); addTearDown(container.dispose); @@ -608,9 +587,9 @@ void main() { store, accountId: 'account-b', cloudReady: true, - referenceSnapshotLoader: (_) async { + referenceLoader: (_, __) async { snapshotReads += 1; - return _fullSnapshot(); + return {}; }, ); addTearDown(accountB.dispose); @@ -736,7 +715,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), cloudReady: true, cloudEnabled: true, - referenceSnapshotLoader: (_) async => _fullSnapshot(), + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); @@ -883,29 +862,13 @@ void main() { ), ); var snapshotReads = 0; - final snapshot = _fullSnapshot( - elements: [ - RemoteElement( - publicId: 'acked-image', - strategyPublicId: 'strategy-a', - pagePublicId: 'page-a', - elementType: 'image', - payload: cloudElementPayload( - kind: 'image', - data: const {'id': 'acked-image', 'elementType': 'image'}, - ), - sortIndex: 0, - revision: 1, - deleted: false, - ), - ], - ); + final snapshot = {'acked-image'}; final container = _container( mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), strategyOpen: false, cloudReady: true, - referenceSnapshotLoader: (strategyId) async { + referenceLoader: (strategyId, _) async { snapshotReads += 1; expect(strategyId, 'strategy-a'); return snapshot; @@ -961,7 +924,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), strategyOpen: false, cloudReady: true, - referenceSnapshotLoader: (_) async => _fullSnapshot(), + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); @@ -997,7 +960,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), cloudReady: true, cloudEnabled: true, - referenceSnapshotLoader: (_) async => _fullSnapshot(), + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); @@ -1017,7 +980,7 @@ void main() { mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), cloudReady: true, - referenceSnapshotLoader: (_) async => _fullSnapshot(), + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); final queue = container.read(cloudMediaUploadQueueProvider.notifier); @@ -1131,10 +1094,11 @@ void main() { required MemoryDurableCloudMediaOutboxStore mediaStore, required MemoryDurableStrategyOutboxStore strategyStore, PendingMediaBytesStore? bytesStore, - CloudMediaReferenceSnapshotLoader? referenceSnapshotLoader, + CloudMediaReferenceLoader? referenceLoader, + _UploadRecordingRepository? serverRepository, }) { var online = false; - final repository = _UploadRecordingRepository(); + final repository = serverRepository ?? _UploadRecordingRepository(); final container = ProviderContainer(overrides: [ durableCloudMediaOutboxStoreProvider.overrideWithValue(mediaStore), convexStrategyRepositoryProvider.overrideWithValue(repository), @@ -1149,21 +1113,11 @@ void main() { imageFilesOnDeviceProvider.overrideWithValue(false), pendingMediaBytesStoreProvider .overrideWithValue(bytesStore ?? MemoryPendingMediaBytesStore()), - cloudMediaReferenceSnapshotLoaderProvider.overrideWithValue( - referenceSnapshotLoader ?? - (_) async => _fullSnapshot(elements: [ - const RemoteElement( - publicId: 'server-image', - strategyPublicId: 'strategy-a', - pagePublicId: 'page-a', - elementType: 'image', - payload: {'id': 'server-image'}, - sortIndex: 0, - revision: 1, - deleted: false, - ), - ]), - ), + // With [serverRepository], its answers are the server's. + if (serverRepository == null) + cloudMediaReferenceLoaderProvider.overrideWithValue( + referenceLoader ?? (_, __) async => {'server-image'}, + ), ]); return ( container: container, @@ -1215,6 +1169,65 @@ void main() { expect(repository.uploadedAssetIds, isEmpty); }); + test( + 'keeps an upload whose image the server still names, whatever the ' + 'full snapshot shows', () async { + final mediaStore = MemoryDurableCloudMediaOutboxStore(); + await mediaStore.put(job('page-image')); + final (:container, :repository, :goOnline) = setUp( + mediaStore: mediaStore, + strategyStore: MemoryDurableStrategyOutboxStore(), + // The full snapshot shows only the pages on screen; the reference + // list counts every image the server's content still shows. + serverRepository: _ReferencingRepository({'page-image'}), + ); + addTearDown(container.dispose); + final bytes = container.read(pendingMediaBytesProvider.notifier); + await bytes.put(key('page-image'), Uint8List.fromList([1, 2, 3])); + final queue = container.read(cloudMediaUploadQueueProvider.notifier); + + await queue.recheckAfterDiscardedWork('strategy-a'); + goOnline(); + await queue.retryNow(ignoreBackoff: true); + await Future.delayed(const Duration(milliseconds: 50)); + + expect( + mediaStore.load().jobs.map((job) => job.assetPublicId), + ['page-image'], + ); + expect(bytes.bytesFor(key('page-image')), [1, 2, 3]); + }); + + test( + 'a staged upload keeps its check, and its bytes, while the server ' + 'cannot tell', () async { + final mediaStore = MemoryDurableCloudMediaOutboxStore(); + // Staged: the change placing its image has not landed, so only the + // server can say whether anything still shows it. + await mediaStore.put(job('page-image', referenceDurable: false)); + final (:container, :repository, :goOnline) = setUp( + mediaStore: mediaStore, + strategyStore: MemoryDurableStrategyOutboxStore(), + serverRepository: _ReferencingRepository(null), + ); + addTearDown(container.dispose); + final bytes = container.read(pendingMediaBytesProvider.notifier); + await bytes.put(key('page-image'), Uint8List.fromList([1, 2, 3])); + final queue = container.read(cloudMediaUploadQueueProvider.notifier); + + await queue.recheckAfterDiscardedWork('strategy-a'); + goOnline(); + await queue.retryNow(ignoreBackoff: true); + await Future.delayed(const Duration(milliseconds: 50)); + + expect( + mediaStore.load().jobs.map((job) => job.assetPublicId), + ['page-image'], + ); + expect(bytes.bytesFor(key('page-image')), [1, 2, 3]); + expect(repository.uploadedAssetIds, isEmpty); + }); + test('a restart before the server answers still checks', () async { final mediaStore = MemoryDurableCloudMediaOutboxStore(); await mediaStore.put(job('page-image')); @@ -1325,10 +1338,10 @@ void main() { final (:container, :repository, :goOnline) = setUp( mediaStore: mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), - referenceSnapshotLoader: (_) async { + referenceLoader: (_, __) async { reads += 1; if (reads == 1) throw StateError('read failed'); - return _fullSnapshot(); + return {}; }, ); addTearDown(container.dispose); @@ -1348,14 +1361,14 @@ void main() { }); test('a retry during a check waits for it, uploading nothing', () async { - final snapshot = Completer(); + final snapshot = Completer?>(); var reads = 0; final mediaStore = MemoryDurableCloudMediaOutboxStore(); await mediaStore.put(job('page-image')); final (:container, :repository, :goOnline) = setUp( mediaStore: mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), - referenceSnapshotLoader: (_) { + referenceLoader: (_, __) { reads += 1; // A second, overlapping check could not tell and let it upload. if (reads > 1) throw StateError('read failed'); @@ -1375,7 +1388,7 @@ void main() { await Future.delayed(const Duration(milliseconds: 20)); expect(repository.uploadedAssetIds, isEmpty); - snapshot.complete(_fullSnapshot()); + snapshot.complete({}); await Future.wait([first, second]); await Future.delayed(const Duration(milliseconds: 20)); @@ -1433,3 +1446,40 @@ class _UploadRecordingRepository implements ConvexStrategyRepository { @override dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); } + +/// A server whose content shows [referenced] images (null: it cannot tell +/// yet) and whose full snapshot shows none. +class _ReferencingRepository extends _UploadRecordingRepository { + _ReferencingRepository(this.referenced); + + final Set? referenced; + + @override + Future?> fetchReferencedAssetIds( + String strategyPublicId, + Iterable assetPublicIds, + ) async => + referenced?.intersection(assetPublicIds.toSet()); + + @override + Future fetchFullSnapshot( + String strategyPublicId, { + String? shareToken, + }) async { + final now = DateTime.utc(2026, 9, 3); + return RemoteFullStrategySnapshot( + header: RemoteStrategyHeader( + publicId: strategyPublicId, + name: 'Strategy A', + mapData: 'ascent', + revision: 1, + createdAt: now, + updatedAt: now, + ), + pages: const [], + elementsByPage: const {}, + lineupsByPage: const {}, + assetsById: const {}, + ); + } +} diff --git a/test/collab/referenced_asset_ids_test.dart b/test/collab/referenced_asset_ids_test.dart new file mode 100644 index 00000000..c9ebafd5 --- /dev/null +++ b/test/collab/referenced_asset_ids_test.dart @@ -0,0 +1,88 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:icarus/collab/convex_strategy_repository.dart'; +import 'package:icarus/collab/generated/generated.dart'; +import 'package:icarus/collab/transport/convex_transport.dart'; + +void main() { + final asked = [for (var i = 0; i < 250; i += 1) 'image-$i']; + + test('asks about images a bounded batch at a time', () async { + final transport = + _ReferencesTransport({'image-3', 'image-120', 'image-249'}); + final repository = ConvexStrategyRepository(IcarusConvexApi(transport)); + + final referenced = await repository.fetchReferencedAssetIds( + 'strategy-a', + [...asked, 'image-3'], + ); + + expect(referenced, {'image-3', 'image-120', 'image-249'}); + expect(transport.batchSizes, [100, 100, 50]); + }); + + test( + 'full snapshots are read as a client that can take one without the ' + "trash's pages", () async { + final transport = _RecordingTransport(); + final repository = ConvexStrategyRepository(IcarusConvexApi(transport)); + + await expectLater( + repository.fetchFullSnapshot('strategy-a'), throwsStateError); + + final (name, args) = transport.calls.single; + expect(name, 'strategy:getFullSnapshot'); + expect(args.value['acceptsTrashedPagesLeftOut'], isA()); + expect((args.value['acceptsTrashedPagesLeftOut'] as ConvexBoolean).value, + isTrue); + }); + + test('cannot tell if any batch cannot', () async { + final transport = _ReferencesTransport({'image-3'}, nullBatch: 1); + final repository = ConvexStrategyRepository(IcarusConvexApi(transport)); + + expect( + await repository.fetchReferencedAssetIds('strategy-a', asked), isNull); + }); +} + +/// A server whose content shows [referenced]; the batch numbered [nullBatch] +/// is answered as before the reference backfill. +final class _ReferencesTransport implements ConvexTransport { + _ReferencesTransport(this.referenced, {this.nullBatch}); + + final Set referenced; + final int? nullBatch; + final batchSizes = []; + + @override + Future query(String name, ConvexObject args) async { + expect(name, 'images:listReferencedAssetIds'); + final ids = [ + for (final id in (args.value['assetPublicIds'] as ConvexArray).value) + (id as ConvexString).value, + ]; + batchSizes.add(ids.length); + if (batchSizes.length - 1 == nullBatch) return const ConvexNull(); + return ConvexArray([ + for (final id in ids) + if (referenced.contains(id)) ConvexString(id), + ]); + } + + @override + dynamic noSuchMethod(Invocation invocation) => throw UnimplementedError(); +} + +/// Records the arguments of each query, then fails it. +final class _RecordingTransport implements ConvexTransport { + final calls = <(String, ConvexObject)>[]; + + @override + Future query(String name, ConvexObject args) async { + calls.add((name, args)); + throw StateError('not answered in this test'); + } + + @override + dynamic noSuchMethod(Invocation invocation) => throw UnimplementedError(); +} diff --git a/test/web_media_bytes_upload_test.dart b/test/web_media_bytes_upload_test.dart index e2a5f177..978ef07e 100644 --- a/test/web_media_bytes_upload_test.dart +++ b/test/web_media_bytes_upload_test.dart @@ -272,8 +272,8 @@ ProviderContainer _webSession({ convexConnectionProvider.overrideWith((ref) => Stream.value(online)), strategyProvider.overrideWith(_CloudStrategy.new), strategyOpQueueProvider.overrideWith(_OpQueue.new), - cloudMediaReferenceSnapshotLoaderProvider.overrideWithValue( - (_) async => throw StateError('no server reads in this test'), + cloudMediaReferenceLoaderProvider.overrideWithValue( + (_, __) async => throw StateError('no server reads in this test'), ), ], ); @@ -313,7 +313,8 @@ RemoteImageAsset _activeAsset(String id) => RemoteImageAsset( uploadStatus: 'active', ); -Future _settle() => Future.delayed(const Duration(milliseconds: 50)); +Future _settle() => + Future.delayed(const Duration(milliseconds: 50)); Future _until(bool Function() done) async { for (var i = 0; i < 60 && !done(); i++) { @@ -359,8 +360,7 @@ void main() { // Offline: the picked image is saved and queued, nothing is sent. await tester.runAsync(() async { - final first = - _webSession(mediaStore: mediaStore, bytesStore: bytesStore); + final first = _webSession(mediaStore: mediaStore, bytesStore: bytesStore); await first.read(placedImageProvider.notifier).saveSecureImage( _imageBytes, _imageId, @@ -730,8 +730,7 @@ void main() { expect(bytesStore.values, hasLength(1)); }); - test('drafts whose queuing failed go when the dialog has closed', - () async { + test('drafts whose queuing failed go when the dialog has closed', () async { await drafts.add(_imageBytes, '.png'); final handedOver = drafts.handOver(); await drafts.dismissed();