From a32019794c9b6e7b81d21706628d80001c5caee0 Mon Sep 17 00:00:00 2001 From: rdlabo Date: Sun, 30 Aug 2026 22:55:30 +0900 Subject: [PATCH] feat(offline): add fastest available read strategy --- .../offline/src/lib/offline-request-policy.ts | 5 +- .../src/lib/offline.interceptor.spec.ts | 162 +++++++++++++++++- .../offline/src/lib/offline.interceptor.ts | 41 +++++ 3 files changed, 205 insertions(+), 3 deletions(-) diff --git a/projects/kit/offline/src/lib/offline-request-policy.ts b/projects/kit/offline/src/lib/offline-request-policy.ts index a020c47..e0d93fb 100644 --- a/projects/kit/offline/src/lib/offline-request-policy.ts +++ b/projects/kit/offline/src/lib/offline-request-policy.ts @@ -43,7 +43,7 @@ export function shouldCommitOfflineCollection(emission: OfflineReadEmission { }); }); + describe('fastest-first read strategy', () => { + const fastestFirstPlan = (overrides: Partial = {}): OfflineRequestPlan => ({ + kind: 'read', + readStrategy: 'fastest-first', + readLocal: vi.fn(async () => null), + ...overrides, + }); + + it('remote先着時はreplica lane待ちのlocal readを抑止する', async () => { + const coordinator = TestBed.inject(OfflineReplicaMutationCoordinator); + let releaseMutation!: () => void; + const mutationGate = new Promise((resolve) => (releaseMutation = resolve)); + let mutationStarted!: () => void; + const mutationReady = new Promise((resolve) => (mutationStarted = resolve)); + const mutation = coordinator.run(async () => { + mutationStarted(); + await mutationGate; + }); + await mutationReady; + + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const readLocal = vi.fn(async () => new HttpResponse({ body: { value: 'stale local' }, status: 200 })); + resolve.mockReturnValue(fastestFirstPlan({ readLocal })); + + await expect(firstValueFrom(run(new HttpRequest('GET', '/bootstrap'), () => of(remote)))).resolves.toBe(remote); + releaseMutation(); + await mutation; + await coordinator.drain(); + + expect(readLocal).not.toHaveBeenCalled(); + expect(markApiSuccess).toHaveBeenCalledOnce(); + }); + + it('local先着時はlocalをemitしてからremoteで再検証する', async () => { + const local = new HttpResponse({ body: { value: 'local' }, status: 200 }); + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const transport$ = new Subject>(); + resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(async () => local) })); + + const emissions: HttpResponse[] = []; + let complete!: () => void; + const completed = new Promise((resolve) => (complete = resolve)); + run(new HttpRequest('GET', '/bootstrap'), () => transport$).subscribe({ + next: (event) => { + if (event instanceof HttpResponse) emissions.push(event); + }, + complete, + }); + await vi.waitFor(() => expect(emissions).toHaveLength(1)); + transport$.next(remote); + transport$.complete(); + await completed; + + expect(emissions).toHaveLength(2); + expect(emissions[0] instanceof HttpResponse && emissions[0].body).toEqual({ value: 'local' }); + expect(emissions[0] instanceof HttpResponse && emissions[0].headers.get(OFFLINE_RESPONSE_HEADER)).toBe('local'); + expect(emissions[1]).toBe(remote); + }); + + it('remoteのSent eventは勝敗に使わず、即時localの後にremote responseをemitする', async () => { + const local = new HttpResponse({ body: { value: 'local' }, status: 200 }); + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const transport$ = new Subject>(); + resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(async () => local) })); + + const emissions: HttpResponse[] = []; + let complete!: () => void; + const completed = new Promise((resolve) => (complete = resolve)); + run(new HttpRequest('GET', '/bootstrap'), () => transport$).subscribe({ + next: (event) => { + if (event instanceof HttpResponse) emissions.push(event); + }, + complete, + }); + transport$.next({ type: HttpEventType.Sent }); + await vi.waitFor(() => expect(emissions).toHaveLength(1)); + transport$.next(remote); + transport$.complete(); + await completed; + + expect(emissions.map((event) => event.body)).toEqual([{ value: 'local' }, { value: 'remote' }]); + }); + + it('remote先着後に遅いlocalをemitしない', async () => { + let resolveLocal!: (value: HttpResponse) => void; + const localReady = new Promise>((resolve) => (resolveLocal = resolve)); + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(() => localReady) })); + + const emissions = await collect(run(new HttpRequest('GET', '/bootstrap'), () => of(remote))); + resolveLocal(new HttpResponse({ body: { value: 'stale local' }, status: 200 })); + await new Promise((resolve) => queueMicrotask(resolve)); + + expect(emissions).toEqual([remote]); + }); + + it('local miss時はremoteを待つ', async () => { + let emitRemote!: () => void; + const remote$ = new Observable>((subscriber) => { + emitRemote = () => { + subscriber.next(new HttpResponse({ body: { value: 'remote' }, status: 200 })); + subscriber.complete(); + }; + }); + resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(async () => null) })); + + const pending = firstValueFrom(run(new HttpRequest('GET', '/bootstrap'), () => remote$)); + await new Promise((resolve) => queueMicrotask(resolve)); + emitRemote(); + await expect(pending).resolves.toMatchObject({ body: { value: 'remote' } }); + }); + + it('remote error先着でもlocal hitを先にemitする', async () => { + const remoteError = new HttpErrorResponse({ status: 500, error: 'server error' }); + let resolveLocal!: (value: HttpResponse) => void; + const localReady = new Promise>((resolve) => (resolveLocal = resolve)); + resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(() => localReady) })); + + const emissions: HttpResponse[] = []; + let observedError: unknown; + const settled = new Promise((resolve) => { + run(new HttpRequest('GET', '/bootstrap'), () => throwError(() => remoteError)).subscribe({ + next: (event) => { + if (event instanceof HttpResponse) emissions.push(event); + }, + error: (error: unknown) => { + observedError = error; + resolve(); + }, + }); + }); + await new Promise((resolve) => queueMicrotask(resolve)); + resolveLocal(new HttpResponse({ body: { value: 'local' }, status: 200 })); + await settled; + + expect(emissions).toHaveLength(1); + expect(emissions[0]?.body).toEqual({ value: 'local' }); + expect(observedError).toBe(remoteError); + }); + + it('remote error先着かつlocal missならremote errorを返す', async () => { + const remoteError = new HttpErrorResponse({ status: 500, error: 'server error' }); + resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(async () => null) })); + + await expect( + firstValueFrom(run(new HttpRequest('GET', '/bootstrap'), () => throwError(() => remoteError))), + ).rejects.toBe(remoteError); + }); + }); + it('readStrategy未指定はnetwork-firstのまま', async () => { resolve.mockReturnValue({ kind: 'read', readLocal: vi.fn() }); const response = new HttpResponse({ body: { userId: 1 }, status: 200 }); diff --git a/projects/kit/offline/src/lib/offline.interceptor.ts b/projects/kit/offline/src/lib/offline.interceptor.ts index e5853b3..3204eae 100644 --- a/projects/kit/offline/src/lib/offline.interceptor.ts +++ b/projects/kit/offline/src/lib/offline.interceptor.ts @@ -10,11 +10,14 @@ import { defer, dematerialize, EMPTY, + filter, from, map, materialize, + NEVER, of, Observable, + race, ReplaySubject, take, tap, @@ -48,6 +51,9 @@ export const offlineInterceptor: HttpInterceptorFn = (request, next) => { if (plan.readStrategy === 'local-first') { return readLocalFirst(request, plan, transport, fallback, inject(ErrorHandler), inject(OfflineReplicaMutationCoordinator)); } + if (plan.readStrategy === 'fastest-first') { + return readFastestFirst(request, plan, transport, fallback, inject(ErrorHandler), inject(OfflineReplicaMutationCoordinator)); + } return readNetworkFirst(request, plan, transport, fallback); } if (LOCAL_FIRST_MUTATION_METHODS.has(request.method)) { @@ -111,6 +117,41 @@ function readLocalFirst( ); } +/** Races a serialized local replica read against remote transport without allowing stale local data to follow remote. */ +function readFastestFirst( + request: HttpRequest, + plan: OfflineReadRequestPlan, + transport: () => Observable>, + fallback: OfflineRequestFallbackService, + errorHandler: ErrorHandler, + replicaMutations: OfflineReplicaMutationCoordinator, +): Observable> { + return defer(transport).pipe( + materialize(), + connect( + (bufferedTransport$) => { + const localDecision$ = resolveLocalAttempt(plan, errorHandler, replicaMutations).pipe( + concatMap((localResponse) => + localResponse + ? concat(of(localResponse), drainRemoteAfterLocal(bufferedTransport$, plan)) + : drainRemoteNetworkFirst(bufferedTransport$, request, plan, fallback), + ), + ); + // Angular transport emits Sent/progress events before the response. + // They are not usable read results and therefore must not win the race. + const remoteWinner$ = drainRemoteNetworkFirst(bufferedTransport$, request, plan, fallback).pipe( + filter((event) => event instanceof AngularHttpResponse), + // A remote error is not a usable response. Keep it buffered until the + // local attempt decides whether it can satisfy the read. + catchError(() => NEVER), + ); + return race(localDecision$, remoteWinner$); + }, + { connector: () => new ReplaySubject() }, + ), + ); +} + function resolveLocalAttempt( plan: OfflineReadRequestPlan, errorHandler: ErrorHandler,