From 08873462b5067267a554f4145e540ffbf7f8dc4f Mon Sep 17 00:00:00 2001 From: rdlabo Date: Tue, 11 Aug 2026 08:56:39 +0900 Subject: [PATCH] feat: add local-first offline reads --- .../offline/src/lib/offline-request-policy.ts | 15 + .../src/lib/offline.interceptor.spec.ts | 260 +++++++++++++++++- .../offline/src/lib/offline.interceptor.ts | 141 +++++++++- 3 files changed, 408 insertions(+), 8 deletions(-) diff --git a/projects/kit/offline/src/lib/offline-request-policy.ts b/projects/kit/offline/src/lib/offline-request-policy.ts index 9c92878..27114ad 100644 --- a/projects/kit/offline/src/lib/offline-request-policy.ts +++ b/projects/kit/offline/src/lib/offline-request-policy.ts @@ -13,9 +13,24 @@ export type OfflineResponseSource = 'local' | 'optimistic'; /** Source passed to a read policy's shared response projector. */ export type OfflineReadResponseSource = 'remote' | 'local'; +/** Controls GET emission order between local replica and remote transport. */ +export type OfflineReadStrategy = 'network-first' | 'local-first'; + /** Product read policy backed by a local replica fallback for transport failures. */ export interface OfflineReadRequestPlan { kind: 'read'; + /** + * Controls whether the interceptor hits the network first or emits a local + * replica immediately and revalidates in the background. Defaults to + * `network-first` when omitted so existing policies stay compatible. + * + * @remarks + * `local-first` starts transport and `readLocal()` concurrently, emits + * projected local first on hit, then drains the buffered remote response. + * Callers must keep the HTTP observable subscribed through revalidation; + * `firstValueFrom` and `take(1)` cancel in-flight transport. + */ + readStrategy?: OfflineReadStrategy; /** * Persists and projects a remote response, or projects a local fallback. * diff --git a/projects/kit/offline/src/lib/offline.interceptor.spec.ts b/projects/kit/offline/src/lib/offline.interceptor.spec.ts index 81e3bb2..c844061 100644 --- a/projects/kit/offline/src/lib/offline.interceptor.spec.ts +++ b/projects/kit/offline/src/lib/offline.interceptor.spec.ts @@ -1,7 +1,8 @@ import { HttpContext, HttpErrorResponse, HttpHeaders, HttpRequest, HttpResponse } from '@angular/common/http'; +import { ErrorHandler } from '@angular/core'; import { TestBed } from '@angular/core/testing'; import { beforeEach, describe, expect, it, vi } from 'vitest'; -import { firstValueFrom, of, throwError } from 'rxjs'; +import { finalize, firstValueFrom, of, Subject, throwError, type Observable } from 'rxjs'; import { OfflineNetworkService } from './offline-network.service'; import { offlineInterceptor } from './offline.interceptor'; import { @@ -19,17 +20,20 @@ describe('offlineInterceptor', () => { let resolveMutation: ReturnType) => OfflineMutationRequestPlan | null>>; let markApiSuccess: ReturnType; let markApiFailure: ReturnType; + let handleError: ReturnType; beforeEach(() => { resolve = vi.fn(() => null); resolveMutation = vi.fn(() => null); markApiSuccess = vi.fn(); markApiFailure = vi.fn(); + handleError = vi.fn(); TestBed.configureTestingModule({ providers: [ { provide: OfflineRequestPolicyRegistry, useValue: { resolve } }, { provide: OfflineMutationRequestPolicyRegistry, useValue: { resolve: resolveMutation } }, { provide: OfflineNetworkService, useValue: { markApiSuccess, markApiFailure } }, + { provide: ErrorHandler, useValue: { handleError } }, ], }); }); @@ -194,8 +198,262 @@ describe('offlineInterceptor', () => { expect(resolve).not.toHaveBeenCalled(); expect(markApiFailure).toHaveBeenCalledOnce(); }); + + describe('local-first read strategy', () => { + const localFirstPlan = (overrides: Partial = {}): OfflineRequestPlan => ({ + kind: 'read', + readStrategy: 'local-first', + readLocal: vi.fn(async () => null), + ...overrides, + }); + + it('local hitのあとremoteを順にemitしreachabilityを更新する', async () => { + const rawLocal = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + const projectedLocal = new HttpResponse({ body: { value: 'local-projected' }, status: 200 }); + const transportResponse = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const projectedRemote = new HttpResponse({ body: { value: 'remote-projected' }, status: 200 }); + const projectResponse = vi.fn(async (_response: HttpResponse, source: OfflineReadResponseSource) => + source === 'local' ? projectedLocal : projectedRemote, + ); + resolve.mockReturnValue( + localFirstPlan({ + readLocal: vi.fn(async () => rawLocal), + projectResponse, + }), + ); + + const emissions = await collect(run(new HttpRequest('GET', '/bootstrap'), () => of(transportResponse))); + + expect(emissions).toHaveLength(2); + expect(emissions[0] instanceof HttpResponse && emissions[0].body).toEqual({ value: 'local-projected' }); + expect(emissions[0] instanceof HttpResponse && emissions[0].headers.get(OFFLINE_RESPONSE_HEADER)).toBe('local'); + expect(emissions[1] instanceof HttpResponse && emissions[1].body).toEqual({ value: 'remote-projected' }); + expect(emissions[1] instanceof HttpResponse && emissions[1].headers.has(OFFLINE_RESPONSE_HEADER)).toBe(false); + expect(projectResponse.mock.calls.map(([, source]) => source)).toEqual(['local', 'remote']); + expect(markApiSuccess).toHaveBeenCalledOnce(); + expect(markApiFailure).not.toHaveBeenCalled(); + }); + + it('deferred localでもsubscribe直後にtransportを開始する', async () => { + let transportSubscribed = false; + let resolveLocal!: (value: HttpResponse) => void; + const localReady = new Promise>((resolve) => { + resolveLocal = resolve; + }); + const next = vi.fn(() => { + transportSubscribed = true; + return of(new HttpResponse({ body: { value: 'remote' }, status: 200 })); + }); + resolve.mockReturnValue(localFirstPlan({ readLocal: vi.fn(() => localReady) })); + + const pending = collect(run(new HttpRequest('GET', '/bootstrap'), next)); + await vi.waitFor(() => expect(transportSubscribed).toBe(true)); + expect(next).toHaveBeenCalledOnce(); + resolveLocal(new HttpResponse({ body: { value: 'cached' }, status: 200 })); + await pending; + }); + + it('remoteが先に完了してもemit順はlocal→remote', async () => { + const transportSubject = new Subject>(); + let resolveLocal!: (value: HttpResponse) => void; + const localReady = new Promise>((resolve) => { + resolveLocal = resolve; + }); + const next = vi.fn(() => transportSubject.asObservable()); + resolve.mockReturnValue(localFirstPlan({ readLocal: vi.fn(() => localReady) })); + + const pending = collect(run(new HttpRequest('GET', '/bootstrap'), next)); + await vi.waitFor(() => expect(next).toHaveBeenCalledOnce()); + transportSubject.next(new HttpResponse({ body: { value: 'remote' }, status: 200 })); + transportSubject.complete(); + resolveLocal(new HttpResponse({ body: { value: 'cached' }, status: 200 })); + + const emissions = await pending; + expect(emissions).toHaveLength(2); + expect(emissions[0] instanceof HttpResponse && emissions[0].body).toEqual({ value: 'cached' }); + expect(emissions[0] instanceof HttpResponse && emissions[0].headers.get(OFFLINE_RESPONSE_HEADER)).toBe('local'); + expect(emissions[1] instanceof HttpResponse && emissions[1].body).toEqual({ value: 'remote' }); + }); + + it('local emit後のstatus=0はerrorにせずreachability failureだけ記録する', async () => { + const local = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + resolve.mockReturnValue(localFirstPlan({ readLocal: vi.fn(async () => local) })); + const error = new HttpErrorResponse({ status: 0, error: new Error('offline') }); + + const emissions = await collect(run(new HttpRequest('GET', '/bootstrap'), () => throwError(() => error))); + + expect(emissions).toHaveLength(1); + expect(emissions[0] instanceof HttpResponse && emissions[0].headers.get(OFFLINE_RESPONSE_HEADER)).toBe('local'); + expect(markApiFailure).toHaveBeenCalledOnce(); + expect(markApiSuccess).not.toHaveBeenCalled(); + }); + + it('local emit後の401/403/500はerrorのまま', async () => { + const local = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + resolve.mockReturnValue(localFirstPlan({ readLocal: vi.fn(async () => local) })); + + for (const status of [401, 403, 500]) { + const error = new HttpErrorResponse({ status }); + await expect( + collect(run(new HttpRequest('GET', '/bootstrap'), () => throwError(() => error))), + ).rejects.toBe(error); + } + }); + + it('local missはnetwork-firstと同じremote→fallback動作', async () => { + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const readLocal = vi.fn(async () => null); + const next = vi.fn(() => of(remote)); + resolve.mockReturnValue(localFirstPlan({ readLocal })); + + await expect(firstValueFrom(run(new HttpRequest('GET', '/bootstrap'), next))).resolves.toBe(remote); + + expect(readLocal).toHaveBeenCalledOnce(); + expect(next).toHaveBeenCalledOnce(); + expect(markApiSuccess).toHaveBeenCalledOnce(); + }); + + it('local read失敗はErrorHandlerへ報告しremoteを継続する', async () => { + const localError = new Error('sqlite locked'); + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const readLocal = vi.fn(async () => { + throw localError; + }); + const next = vi.fn(() => of(remote)); + resolve.mockReturnValue(localFirstPlan({ readLocal })); + + await expect(firstValueFrom(run(new HttpRequest('GET', '/bootstrap'), next))).resolves.toBe(remote); + + expect(handleError).toHaveBeenCalledWith(localError); + expect(next).toHaveBeenCalledOnce(); + expect(markApiSuccess).toHaveBeenCalledOnce(); + }); + + it('remote projection失敗はlocal emit後もerrorのまま', async () => { + const local = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + const projectionError = new HttpErrorResponse({ status: 500, error: new Error('projection failed') }); + resolve.mockReturnValue( + localFirstPlan({ + readLocal: vi.fn(async () => local), + projectResponse: vi.fn(async (_response, source) => { + if (source === 'remote') throw projectionError; + return local; + }), + }), + ); + + await expect( + collect(run(new HttpRequest('GET', '/bootstrap'), () => of(new HttpResponse({ status: 200 })))), + ).rejects.toBe(projectionError); + }); + + it('remote projectResponseのstatus=0はtransport fallbackとして握りつぶさない', async () => { + const local = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + const projectionError = new HttpErrorResponse({ status: 0, error: new Error('local persistence failed') }); + resolve.mockReturnValue( + localFirstPlan({ + readLocal: vi.fn(async () => local), + projectResponse: vi.fn(async (_response, source) => { + if (source === 'remote') throw projectionError; + return local; + }), + }), + ); + + await expect( + collect(run(new HttpRequest('GET', '/bootstrap'), () => of(new HttpResponse({ status: 200 })))), + ).rejects.toBe(projectionError); + }); + + it('local projectResponse失敗はErrorHandlerへ報告しnetwork-firstへ継続する', async () => { + const local = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + const projectionError = new Error('corrupt cache'); + const remote = new HttpResponse({ body: { value: 'remote' }, status: 200 }); + const next = vi.fn(() => of(remote)); + resolve.mockReturnValue( + localFirstPlan({ + readLocal: vi.fn(async () => local), + projectResponse: vi.fn(async (_response, source) => { + if (source === 'local') throw projectionError; + return remote; + }), + }), + ); + + const emissions = await collect(run(new HttpRequest('GET', '/bootstrap'), next)); + + expect(emissions).toEqual([remote]); + expect(handleError).toHaveBeenCalledWith(projectionError); + expect(next).toHaveBeenCalledOnce(); + expect(markApiSuccess).toHaveBeenCalledOnce(); + }); + + it('local pending中のunsubscribeはtransportをcancelしemit/reportしない', async () => { + let transportUnsubscribed = false; + let resolveLocal!: (value: HttpResponse) => void; + const localReady = new Promise>((resolve) => { + resolveLocal = resolve; + }); + const next = vi.fn(() => + of(new HttpResponse({ status: 200 })).pipe( + finalize(() => { + transportUnsubscribed = true; + }), + ), + ); + resolve.mockReturnValue(localFirstPlan({ readLocal: vi.fn(() => localReady) })); + + const subscription = run(new HttpRequest('GET', '/bootstrap'), next).subscribe(); + await vi.waitFor(() => expect(next).toHaveBeenCalledOnce()); + subscription.unsubscribe(); + expect(transportUnsubscribed).toBe(true); + + resolveLocal(new HttpResponse({ body: { value: 'cached' }, status: 200 })); + await Promise.resolve(); + expect(handleError).not.toHaveBeenCalled(); + }); + + it('remote完了前のunsubscribeはtransportをcancelする', async () => { + const local = new HttpResponse({ body: { value: 'cached' }, status: 200 }); + resolve.mockReturnValue(localFirstPlan({ readLocal: vi.fn(async () => local) })); + let transportUnsubscribed = false; + const transportSubject = new Subject>(); + const transport$ = transportSubject.asObservable().pipe(finalize(() => { + transportUnsubscribed = true; + })); + + const emissions: HttpResponse[] = []; + const subscription = run(new HttpRequest('GET', '/bootstrap'), () => transport$).subscribe({ + next: (event) => { + if (event instanceof HttpResponse) emissions.push(event); + }, + }); + + await vi.waitFor(() => expect(emissions).toHaveLength(1)); + subscription.unsubscribe(); + expect(transportUnsubscribed).toBe(true); + }); + }); + + it('readStrategy未指定はnetwork-firstのまま', async () => { + resolve.mockReturnValue({ kind: 'read', readLocal: vi.fn() }); + const response = new HttpResponse({ body: { userId: 1 }, status: 200 }); + await expect(firstValueFrom(run(new HttpRequest('GET', '/bootstrap'), () => of(response)))).resolves.toBe(response); + expect(markApiSuccess).toHaveBeenCalledOnce(); + }); }); function run(request: HttpRequest, next: Parameters[1]) { return TestBed.runInInjectionContext(() => offlineInterceptor(request, next)); } + +function collect(source: Observable): Promise { + return new Promise((resolvePromise, reject) => { + const values: T[] = []; + source.subscribe({ + next: (value) => values.push(value), + complete: () => resolvePromise(values), + error: reject, + }); + }); +} diff --git a/projects/kit/offline/src/lib/offline.interceptor.ts b/projects/kit/offline/src/lib/offline.interceptor.ts index f941ecb..39bb09d 100644 --- a/projects/kit/offline/src/lib/offline.interceptor.ts +++ b/projects/kit/offline/src/lib/offline.interceptor.ts @@ -1,8 +1,24 @@ import type { HttpEvent, HttpInterceptorFn, HttpRequest } from '@angular/common/http'; import { HttpResponse as AngularHttpResponse } from '@angular/common/http'; import { ErrorHandler, inject, Injectable } from '@angular/core'; -import type { Observable } from 'rxjs'; -import { catchError, concatMap, defer, from, map, of, tap, throwError } from 'rxjs'; +import type { Notification, Observable, ObservableNotification } from 'rxjs'; +import { + catchError, + concat, + concatMap, + connect, + defer, + dematerialize, + EMPTY, + from, + map, + materialize, + of, + ReplaySubject, + take, + tap, + throwError, +} from 'rxjs'; import { isOfflineFallbackError, OfflineNetworkService } from './offline-network.service'; import { OFFLINE_BYPASS, @@ -14,6 +30,8 @@ import { const LOCAL_FIRST_MUTATION_METHODS = new Set(['POST', 'PUT', 'PATCH', 'DELETE']); +type MaterializedTransport = Notification> & ObservableNotification>; + /** Applies product read and local-first mutation policies while observing real API reachability. */ export const offlineInterceptor: HttpInterceptorFn = (request, next) => { const network = inject(OfflineNetworkService); @@ -24,10 +42,10 @@ export const offlineInterceptor: HttpInterceptorFn = (request, next) => { const fallback = inject(OfflineRequestFallbackService); const plan = registry.resolve(request); if (!plan) return transport(); - return defer(transport).pipe( - catchError((error: unknown) => fallback.handle(request, error, plan) ?? throwError(() => error)), - concatMap((event) => projectRemoteResponse(event, plan)), - ); + if (plan.readStrategy === 'local-first') { + return readLocalFirst(request, plan, transport, fallback, inject(ErrorHandler)); + } + return readNetworkFirst(request, plan, transport, fallback); } if (LOCAL_FIRST_MUTATION_METHODS.has(request.method)) { const plan = inject(OfflineMutationRequestPolicyRegistry).resolve(request); @@ -40,7 +58,116 @@ export const offlineInterceptor: HttpInterceptorFn = (request, next) => { return transport(); }; -function projectRemoteResponse(event: HttpEvent, plan: OfflineReadRequestPlan): Observable> { +function readNetworkFirst( + request: HttpRequest, + plan: OfflineReadRequestPlan, + transport: () => Observable>, + fallback: OfflineRequestFallbackService, +): Observable> { + return defer(transport).pipe( + catchError((error: unknown) => fallback.handle(request, error, plan) ?? throwError(() => error)), + concatMap((event) => projectReadResponse(event, plan)), + ); +} + +/** + * Stale-while-revalidate local-first GET handling. + * + * @remarks + * Starts raw transport and `readLocal()` concurrently at outer subscription, + * buffers materialized remote notifications until the local attempt settles, + * emits a projected local response first on hit, then drains/projects the + * buffered transport. Remote projection begins only after the local decision. + * + * Consumers must keep the returned observable subscribed through revalidation; + * `firstValueFrom` and `take(1)` cancel in-flight transport and suppress + * further emissions. + */ +function readLocalFirst( + request: HttpRequest, + plan: OfflineReadRequestPlan, + transport: () => Observable>, + fallback: OfflineRequestFallbackService, + errorHandler: ErrorHandler, +): Observable> { + return defer(transport).pipe( + materialize(), + connect( + (bufferedTransport$) => + resolveLocalAttempt(plan, errorHandler).pipe( + concatMap((localResponse) => + localResponse + ? concat(of(localResponse), drainRemoteAfterLocal(bufferedTransport$, plan)) + : drainRemoteNetworkFirst(bufferedTransport$, request, plan, fallback), + ), + ), + { connector: () => new ReplaySubject() }, + ), + ); +} + +function resolveLocalAttempt( + plan: OfflineReadRequestPlan, + errorHandler: ErrorHandler, +): Observable | null> { + return defer(() => + from(plan.readLocal()).pipe( + catchError((localError: unknown) => { + errorHandler.handleError(localError); + return of(null); + }), + ), + ).pipe( + concatMap((local) => (local ? tryProjectLocal(local, plan, errorHandler) : of(null))), + take(1), + ); +} + +function tryProjectLocal( + cached: AngularHttpResponse, + plan: OfflineReadRequestPlan, + errorHandler: ErrorHandler, +): Observable | null> { + return emitTaggedLocalResponse(cached, plan).pipe( + catchError((error: unknown) => { + errorHandler.handleError(error); + return of(null); + }), + ); +} + +function emitTaggedLocalResponse( + cached: AngularHttpResponse, + plan: OfflineReadRequestPlan, +): Observable> { + return projectReadResponse(cached.clone({ headers: cached.headers.set(OFFLINE_RESPONSE_HEADER, 'local') }), plan); +} + +function drainRemoteAfterLocal( + bufferedTransport$: Observable, + plan: OfflineReadRequestPlan, +): Observable> { + return bufferedTransport$.pipe( + dematerialize(), + catchError((error: unknown) => (isOfflineFallbackError(error) ? EMPTY : throwError(() => error))), + concatMap((event) => projectReadResponse(event, plan)), + ); +} + +function drainRemoteNetworkFirst( + bufferedTransport$: Observable, + request: HttpRequest, + plan: OfflineReadRequestPlan, + fallback: OfflineRequestFallbackService, +): Observable> { + return bufferedTransport$.pipe( + dematerialize(), + catchError((error: unknown) => fallback.handle(request, error, plan) ?? throwError(() => error)), + concatMap((event) => projectReadResponse(event, plan)), + ); +} + +function projectReadResponse(event: HttpEvent, plan: OfflineReadRequestPlan): Observable> { if (!(event instanceof AngularHttpResponse) || !plan.projectResponse) return of(event); const source = event.headers.get(OFFLINE_RESPONSE_HEADER) === 'local' ? 'local' : 'remote'; return from(plan.projectResponse(event, source)).pipe(