From 0854ebc222e11b5e14d041a9cb1c52693664aed6 Mon Sep 17 00:00:00 2001 From: Dara Adedeji Date: Tue, 29 Sep 2026 09:18:31 -0400 Subject: [PATCH 1/4] Ask the server which images it still shows before dropping an upload The upload queue decided whether a pending upload was still wanted by reading the strategy's full snapshot. It now asks the new images:listReferencedAssetIds query: every image the server's content shows, read from the reference rows, or null while the reference backfill has not finished (the queue then keeps the upload and asks again). Nothing changes for users today. This ships before the page trash (#228). The trash hides deleted pages from the full snapshot while their images are still wanted, so a client still reading the snapshot could drop the bytes of an image on a page that can be restored. Co-Authored-By: Claude Opus 5.5 (1M context) --- convex/function_spec.json | 32 +++ convex/images.ts | 16 ++ convex/referencedAssetIds.test.ts | 169 +++++++++++++++ lib/collab/convex_strategy_repository.dart | 9 + lib/collab/generated/convex_models.dart | 18 ++ lib/collab/generated/icarus_convex_api.dart | 18 ++ .../cloud_media_upload_queue_provider.dart | 70 ++----- ...loud_media_upload_queue_provider_test.dart | 196 +++++++++++------- test/web_media_bytes_upload_test.dart | 2 +- 9 files changed, 406 insertions(+), 124 deletions(-) create mode 100644 convex/referencedAssetIds.test.ts diff --git a/convex/function_spec.json b/convex/function_spec.json index 5db7163d..6d8fd58b 100644 --- a/convex/function_spec.json +++ b/convex/function_spec.json @@ -40660,6 +40660,38 @@ "kind": "internal" } }, + { + "args": { + "type": "object", + "value": { + "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", diff --git a/convex/images.ts b/convex/images.ts index d967fd62..10c3d7ab 100644 --- a/convex/images.ts +++ b/convex/images.ts @@ -747,6 +747,22 @@ export const listForStrategy = query({ }, }); +/// Every image 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(), + }, + returns: v.union(v.array(v.string()), v.null()), + handler: async (ctx, args) => { + const strategy = await getStrategyByPublicId(ctx, args.strategyPublicId); + await assertStrategyRole(ctx, strategy, "viewer"); + if (!(await assetReferencesReady(ctx))) return null; + return [...(await collectLiveAssetIds(ctx, strategy._id))].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..be29245a --- /dev/null +++ b/convex/referencedAssetIds.test.ts @@ -0,0 +1,169 @@ +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 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"; + +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("lists every image 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, { strategyPublicId }), + ).toEqual(["shot", "shown"]); + }); + + 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, { strategyPublicId }), + ).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, { strategyPublicId }), + ).toEqual(["shot", "shown"]); + await expect( + stranger.query(listReferencedAssetIds, { strategyPublicId }), + ).rejects.toThrow(); + }); +}); diff --git a/lib/collab/convex_strategy_repository.dart b/lib/collab/convex_strategy_repository.dart index 96629b51..bf22ca58 100644 --- a/lib/collab/convex_strategy_repository.dart +++ b/lib/collab/convex_strategy_repository.dart @@ -157,6 +157,15 @@ class ConvexStrategyRepository { .map(_pageSnapshot); } + /// Every image the strategy's content shows on the server, or null while + /// the server cannot tell yet. + Future?> fetchReferencedAssetIds(String strategyPublicId) async { + final ids = await _api.images + .listReferencedAssetIds(strategyPublicId: strategyPublicId) + .fetch(); + return ids?.toSet(); + } + /// [shareToken]: see [fetchShell]. Future fetchFullSnapshot( String strategyPublicId, { diff --git a/lib/collab/generated/convex_models.dart b/lib/collab/generated/convex_models.dart index cbb2fcb7..f463cd39 100644 --- a/lib/collab/generated/convex_models.dart +++ b/lib/collab/generated/convex_models.dart @@ -5244,6 +5244,24 @@ List decodeImagesListForStrategyResult( ) .toList(growable: false); +ConvexObject encodeImagesListReferencedAssetIdsArgs({ + required String strategyPublicId, +}) => ConvexObject({'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(), diff --git a/lib/collab/generated/icarus_convex_api.dart b/lib/collab/generated/icarus_convex_api.dart index c7cea8fa..e7365156 100644 --- a/lib/collab/generated/icarus_convex_api.dart +++ b/lib/collab/generated/icarus_convex_api.dart @@ -417,6 +417,9 @@ abstract interface class ImagesModule { ConvexQuery> listForStrategy({ required String strategyPublicId, }); + ConvexQuery?> listReferencedAssetIds({ + required String strategyPublicId, + }); } final class _ImagesModule implements ImagesModule { @@ -537,6 +540,21 @@ final class _ImagesModule implements ImagesModule { decode: decodeImagesListForStrategyResult, ); } + + @override + ConvexQuery?> listReferencedAssetIds({ + required String strategyPublicId, + }) { + final args = encodeImagesListReferencedAssetIdsArgs( + strategyPublicId: strategyPublicId, + ); + return ConvexQuery( + transport: _transport, + name: 'images:listReferencedAssetIds', + args: args, + decode: decodeImagesListReferencedAssetIdsResult, + ); + } } abstract interface class InvitesModule { diff --git a/lib/providers/collab/cloud_media_upload_queue_provider.dart b/lib/providers/collab/cloud_media_upload_queue_provider.dart index c70fffc3..3a225cf1 100644 --- a/lib/providers/collab/cloud_media_upload_queue_provider.dart +++ b/lib/providers/collab/cloud_media_upload_queue_provider.dart @@ -97,12 +97,14 @@ final cloudMediaUploadQueueProvider = CloudMediaUploadQueueNotifier.new, ); -typedef CloudMediaReferenceSnapshotLoader = Future - Function(String strategyPublicId); - -final cloudMediaReferenceSnapshotLoaderProvider = - Provider( - (ref) => ref.watch(convexStrategyRepositoryProvider).fetchFullSnapshot, +/// Every image 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); + +final cloudMediaReferenceLoaderProvider = Provider( + (ref) => ref.watch(convexStrategyRepositoryProvider).fetchReferencedAssetIds, ); final cloudMediaAccountIdProvider = Provider( @@ -370,22 +372,20 @@ 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, ); } 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,9 +845,9 @@ class CloudMediaUploadQueueNotifier return null; } - late final RemoteFullStrategySnapshot snapshot; + final Set? serverReferences; try { - snapshot = await ref.read(cloudMediaReferenceSnapshotLoaderProvider)( + serverReferences = await ref.read(cloudMediaReferenceLoaderProvider)( job.strategyPublicId, ); } catch (error) { @@ -858,12 +858,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 +1137,17 @@ 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); } catch (error) { _logMedia( 'reference_reconcile.deferred strategy=${entry.key} error=$error', ); continue; } + if (serverReferences == null) continue; if (ref.read(cloudMediaAccountIdProvider) != accountId) return; @@ -1156,7 +1158,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 +1178,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..6ad6dc2d 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,62 @@ 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('keeps an upload and its bytes while the server cannot tell', + () async { + final mediaStore = MemoryDurableCloudMediaOutboxStore(); + await mediaStore.put(job('page-image')); + 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 +1335,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 +1358,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 +1385,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 +1443,37 @@ 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) async => + referenced; + + @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/web_media_bytes_upload_test.dart b/test/web_media_bytes_upload_test.dart index e2a5f177..e3216354 100644 --- a/test/web_media_bytes_upload_test.dart +++ b/test/web_media_bytes_upload_test.dart @@ -272,7 +272,7 @@ ProviderContainer _webSession({ convexConnectionProvider.overrideWith((ref) => Stream.value(online)), strategyProvider.overrideWith(_CloudStrategy.new), strategyOpQueueProvider.overrideWith(_OpQueue.new), - cloudMediaReferenceSnapshotLoaderProvider.overrideWithValue( + cloudMediaReferenceLoaderProvider.overrideWithValue( (_) async => throw StateError('no server reads in this test'), ), ], From 6a305869a4a74fbb5670d7beb2154a83f842dcb0 Mon Sep 17 00:00:00 2001 From: Dara Adedeji Date: Tue, 29 Sep 2026 09:30:46 -0400 Subject: [PATCH 2/4] Test the uncertain answer on a staged upload, the job it holds A durable upload may upload while the server cannot tell (nothing the user made is lost that way); a staged one keeps its check, its job and its bytes, and never uploads. The test now covers the staged case, where the answer decides. Co-Authored-By: Claude Opus 5.5 (1M context) --- test/cloud_media_upload_queue_provider_test.dart | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/test/cloud_media_upload_queue_provider_test.dart b/test/cloud_media_upload_queue_provider_test.dart index 6ad6dc2d..4b78d125 100644 --- a/test/cloud_media_upload_queue_provider_test.dart +++ b/test/cloud_media_upload_queue_provider_test.dart @@ -1198,10 +1198,13 @@ void main() { expect(bytes.bytesFor(key('page-image')), [1, 2, 3]); }); - test('keeps an upload and its bytes while the server cannot tell', - () async { + test( + 'a staged upload keeps its check, and its bytes, while the server ' + 'cannot tell', () async { final mediaStore = MemoryDurableCloudMediaOutboxStore(); - await mediaStore.put(job('page-image')); + // 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(), From d5f9209e7d4e1711ddba838515c8af204f4eecd4 Mon Sep 17 00:00:00 2001 From: Dara Adedeji Date: Tue, 29 Sep 2026 09:33:45 -0400 Subject: [PATCH 3/4] Ask only about the images in question, a bounded batch at a time listReferencedAssetIds read every live reference row of the strategy, so a strategy with tens of thousands of image references (thousands of lineups sharing screenshots) would pass Convex's read limit and never get an answer. It now takes the asset ids the client holds uploads for, at most 100 per call, and looks each up through the strategy/asset/ deleted index. The repository asks in batches of 100; if any batch cannot tell, the answer is null. Co-Authored-By: Claude Opus 5.5 (1M context) --- convex/function_spec.json | 9 +++ convex/images.ts | 31 ++++++++-- convex/referencedAssetIds.test.ts | 35 +++++++++-- lib/collab/convex_strategy_repository.dart | 34 ++++++++--- lib/collab/generated/convex_models.dart | 10 +++- lib/collab/generated/icarus_convex_api.dart | 3 + .../cloud_media_upload_queue_provider.dart | 18 ++++-- ...loud_media_upload_queue_provider_test.dart | 29 +++++----- test/collab/referenced_asset_ids_test.dart | 58 +++++++++++++++++++ test/web_media_bytes_upload_test.dart | 11 ++-- 10 files changed, 196 insertions(+), 42 deletions(-) create mode 100644 test/collab/referenced_asset_ids_test.dart diff --git a/convex/function_spec.json b/convex/function_spec.json index 6d8fd58b..228924bb 100644 --- a/convex/function_spec.json +++ b/convex/function_spec.json @@ -40664,6 +40664,15 @@ "args": { "type": "object", "value": { + "assetPublicIds": { + "fieldType": { + "type": "array", + "value": { + "type": "string" + } + }, + "optional": false + }, "strategyPublicId": { "fieldType": { "type": "string" diff --git a/convex/images.ts b/convex/images.ts index 10c3d7ab..71782b24 100644 --- a/convex/images.ts +++ b/convex/images.ts @@ -747,19 +747,42 @@ export const listForStrategy = query({ }, }); -/// Every image 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. +/// 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; - return [...(await collectLiveAssetIds(ctx, strategy._id))].sort(); + 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(); }, }); diff --git a/convex/referencedAssetIds.test.ts b/convex/referencedAssetIds.test.ts index be29245a..b1ff7eee 100644 --- a/convex/referencedAssetIds.test.ts +++ b/convex/referencedAssetIds.test.ts @@ -30,6 +30,11 @@ 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 { @@ -125,22 +130,42 @@ async function seed(t: RootHarness): Promise { } describe("images:listReferencedAssetIds", () => { - test("lists every image the strategy's content shows, deleted content left out", async () => { + 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, { strategyPublicId }), + 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, { strategyPublicId }), + await owner.query(listReferencedAssetIds, asked), ).toBeNull(); }); @@ -160,10 +185,10 @@ describe("images:listReferencedAssetIds", () => { const stranger = await user(t, "stranger"); expect( - await viewer.query(listReferencedAssetIds, { strategyPublicId }), + await viewer.query(listReferencedAssetIds, asked), ).toEqual(["shot", "shown"]); await expect( - stranger.query(listReferencedAssetIds, { strategyPublicId }), + stranger.query(listReferencedAssetIds, asked), ).rejects.toThrow(); }); }); diff --git a/lib/collab/convex_strategy_repository.dart b/lib/collab/convex_strategy_repository.dart index bf22ca58..5207f61d 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,15 +159,33 @@ class ConvexStrategyRepository { .map(_pageSnapshot); } - /// Every image the strategy's content shows on the server, or null while - /// the server cannot tell yet. - Future?> fetchReferencedAssetIds(String strategyPublicId) async { - final ids = await _api.images - .listReferencedAssetIds(strategyPublicId: strategyPublicId) - .fetch(); - return ids?.toSet(); + /// 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, { diff --git a/lib/collab/generated/convex_models.dart b/lib/collab/generated/convex_models.dart index f463cd39..b31e4b83 100644 --- a/lib/collab/generated/convex_models.dart +++ b/lib/collab/generated/convex_models.dart @@ -5245,8 +5245,16 @@ List decodeImagesListForStrategyResult( .toList(growable: false); ConvexObject encodeImagesListReferencedAssetIdsArgs({ + required List assetPublicIds, required String strategyPublicId, -}) => ConvexObject({'strategyPublicId': ConvexString(strategyPublicId)}); +}) => ConvexObject({ + 'assetPublicIds': ConvexArray( + assetPublicIds.indexed + .map((entry) => ConvexString(entry.$2)) + .toList(growable: false), + ), + 'strategyPublicId': ConvexString(strategyPublicId), +}); List? decodeImagesListReferencedAssetIdsResult(ConvexValue value) => (value) is ConvexNull diff --git a/lib/collab/generated/icarus_convex_api.dart b/lib/collab/generated/icarus_convex_api.dart index e7365156..5c9bb6a8 100644 --- a/lib/collab/generated/icarus_convex_api.dart +++ b/lib/collab/generated/icarus_convex_api.dart @@ -418,6 +418,7 @@ abstract interface class ImagesModule { required String strategyPublicId, }); ConvexQuery?> listReferencedAssetIds({ + required List assetPublicIds, required String strategyPublicId, }); } @@ -543,9 +544,11 @@ final class _ImagesModule implements ImagesModule { @override ConvexQuery?> listReferencedAssetIds({ + required List assetPublicIds, required String strategyPublicId, }) { final args = encodeImagesListReferencedAssetIdsArgs( + assetPublicIds: assetPublicIds, strategyPublicId: strategyPublicId, ); return ConvexQuery( diff --git a/lib/providers/collab/cloud_media_upload_queue_provider.dart b/lib/providers/collab/cloud_media_upload_queue_provider.dart index 3a225cf1..98d2993c 100644 --- a/lib/providers/collab/cloud_media_upload_queue_provider.dart +++ b/lib/providers/collab/cloud_media_upload_queue_provider.dart @@ -97,11 +97,13 @@ final cloudMediaUploadQueueProvider = CloudMediaUploadQueueNotifier.new, ); -/// Every image 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. +/// 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); + String strategyPublicId, + Iterable assetPublicIds, +); final cloudMediaReferenceLoaderProvider = Provider( (ref) => ref.watch(convexStrategyRepositoryProvider).fetchReferencedAssetIds, @@ -378,6 +380,7 @@ class CloudMediaUploadQueueNotifier try { serverReferences = await ref.read(cloudMediaReferenceLoaderProvider)( strategyPublicId, + [for (final job in missingReferences) job.assetPublicId], ); } catch (_) { serverReferences = null; @@ -849,6 +852,7 @@ class CloudMediaUploadQueueNotifier try { serverReferences = await ref.read(cloudMediaReferenceLoaderProvider)( job.strategyPublicId, + [job.assetPublicId], ); } catch (error) { _logMedia( @@ -1139,8 +1143,10 @@ class CloudMediaUploadQueueNotifier for (final entry in byStrategy.entries) { final Set? serverReferences; try { - serverReferences = - await ref.read(cloudMediaReferenceLoaderProvider)(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', diff --git a/test/cloud_media_upload_queue_provider_test.dart b/test/cloud_media_upload_queue_provider_test.dart index 4b78d125..d5dfba0d 100644 --- a/test/cloud_media_upload_queue_provider_test.dart +++ b/test/cloud_media_upload_queue_provider_test.dart @@ -455,7 +455,7 @@ void main() { store, strategyStore: ops, cloudReady: true, - referenceLoader: (_) { + referenceLoader: (_, __) { if (!started.isCompleted) started.complete(); return snapshot.future; }, @@ -522,7 +522,7 @@ void main() { final container = _container( store, cloudReady: true, - referenceLoader: (_) async { + referenceLoader: (_, __) async { if (!snapshotAvailable) throw StateError('offline'); return {}; }, @@ -587,7 +587,7 @@ void main() { store, accountId: 'account-b', cloudReady: true, - referenceLoader: (_) async { + referenceLoader: (_, __) async { snapshotReads += 1; return {}; }, @@ -715,7 +715,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), cloudReady: true, cloudEnabled: true, - referenceLoader: (_) async => {}, + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); @@ -868,7 +868,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), strategyOpen: false, cloudReady: true, - referenceLoader: (strategyId) async { + referenceLoader: (strategyId, _) async { snapshotReads += 1; expect(strategyId, 'strategy-a'); return snapshot; @@ -924,7 +924,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), strategyOpen: false, cloudReady: true, - referenceLoader: (_) async => {}, + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); @@ -960,7 +960,7 @@ void main() { strategyStore: MemoryDurableStrategyOutboxStore(), cloudReady: true, cloudEnabled: true, - referenceLoader: (_) async => {}, + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); @@ -980,7 +980,7 @@ void main() { mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), cloudReady: true, - referenceLoader: (_) async => {}, + referenceLoader: (_, __) async => {}, ); addTearDown(container.dispose); final queue = container.read(cloudMediaUploadQueueProvider.notifier); @@ -1116,7 +1116,7 @@ void main() { // With [serverRepository], its answers are the server's. if (serverRepository == null) cloudMediaReferenceLoaderProvider.overrideWithValue( - referenceLoader ?? (_) async => {'server-image'}, + referenceLoader ?? (_, __) async => {'server-image'}, ), ]); return ( @@ -1338,7 +1338,7 @@ void main() { final (:container, :repository, :goOnline) = setUp( mediaStore: mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), - referenceLoader: (_) async { + referenceLoader: (_, __) async { reads += 1; if (reads == 1) throw StateError('read failed'); return {}; @@ -1368,7 +1368,7 @@ void main() { final (:container, :repository, :goOnline) = setUp( mediaStore: mediaStore, strategyStore: MemoryDurableStrategyOutboxStore(), - referenceLoader: (_) { + referenceLoader: (_, __) { reads += 1; // A second, overlapping check could not tell and let it upload. if (reads > 1) throw StateError('read failed'); @@ -1455,8 +1455,11 @@ class _ReferencingRepository extends _UploadRecordingRepository { final Set? referenced; @override - Future?> fetchReferencedAssetIds(String strategyPublicId) async => - referenced; + Future?> fetchReferencedAssetIds( + String strategyPublicId, + Iterable assetPublicIds, + ) async => + referenced?.intersection(assetPublicIds.toSet()); @override Future fetchFullSnapshot( diff --git a/test/collab/referenced_asset_ids_test.dart b/test/collab/referenced_asset_ids_test.dart new file mode 100644 index 00000000..6297c309 --- /dev/null +++ b/test/collab/referenced_asset_ids_test.dart @@ -0,0 +1,58 @@ +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('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(); +} diff --git a/test/web_media_bytes_upload_test.dart b/test/web_media_bytes_upload_test.dart index e3216354..978ef07e 100644 --- a/test/web_media_bytes_upload_test.dart +++ b/test/web_media_bytes_upload_test.dart @@ -273,7 +273,7 @@ ProviderContainer _webSession({ strategyProvider.overrideWith(_CloudStrategy.new), strategyOpQueueProvider.overrideWith(_OpQueue.new), cloudMediaReferenceLoaderProvider.overrideWithValue( - (_) async => throw StateError('no server reads in this test'), + (_, __) 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(); From 79e444e0a441bbf79aaacb0328843acc554cc732 Mon Sep 17 00:00:00 2001 From: Dara Adedeji Date: Tue, 29 Sep 2026 09:52:06 -0400 Subject: [PATCH 4/4] Say, when reading a full snapshot, that the trash's pages may be left out getFullSnapshot takes an optional acceptsTrashedPagesLeftOut, ignored for now; this client sends it. Once the server keeps deleted pages in a trash (#228), it can refuse a snapshot to a client that does not send it while the strategy holds trashed pages: such a client decides from the snapshot whether an upload is still wanted, and would drop the bytes of an image on a page that can be restored. A refused read already keeps the bytes. Co-Authored-By: Claude Opus 5.5 (1M context) --- convex/function_spec.json | 6 +++++ convex/referencedAssetIds.test.ts | 17 ++++++++++++ convex/strategy.ts | 5 ++++ lib/collab/convex_strategy_repository.dart | 3 +++ lib/collab/generated/convex_models.dart | 6 +++++ lib/collab/generated/icarus_convex_api.dart | 5 ++++ test/collab/referenced_asset_ids_test.dart | 30 +++++++++++++++++++++ 7 files changed, 72 insertions(+) diff --git a/convex/function_spec.json b/convex/function_spec.json index 228924bb..84f4a898 100644 --- a/convex/function_spec.json +++ b/convex/function_spec.json @@ -182529,6 +182529,12 @@ "args": { "type": "object", "value": { + "acceptsTrashedPagesLeftOut": { + "fieldType": { + "type": "boolean" + }, + "optional": true + }, "shareToken": { "fieldType": { "type": "string" diff --git a/convex/referencedAssetIds.test.ts b/convex/referencedAssetIds.test.ts index b1ff7eee..6fa40ad1 100644 --- a/convex/referencedAssetIds.test.ts +++ b/convex/referencedAssetIds.test.ts @@ -20,6 +20,9 @@ const createStrategy = makeFunctionReference<"mutation">( 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", ); @@ -191,4 +194,18 @@ describe("images:listReferencedAssetIds", () => { 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 5207f61d..e7ccd5c4 100644 --- a/lib/collab/convex_strategy_repository.dart +++ b/lib/collab/convex_strategy_repository.dart @@ -196,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 b31e4b83..c9ca04dd 100644 --- a/lib/collab/generated/convex_models.dart +++ b/lib/collab/generated/convex_models.dart @@ -5868,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 5c9bb6a8..f2be0f19 100644 --- a/lib/collab/generated/icarus_convex_api.dart +++ b/lib/collab/generated/icarus_convex_api.dart @@ -1348,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, }); @@ -1362,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/test/collab/referenced_asset_ids_test.dart b/test/collab/referenced_asset_ids_test.dart index 6297c309..c9ebafd5 100644 --- a/test/collab/referenced_asset_ids_test.dart +++ b/test/collab/referenced_asset_ids_test.dart @@ -20,6 +20,22 @@ void main() { 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)); @@ -56,3 +72,17 @@ final class _ReferencesTransport implements ConvexTransport { @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(); +}