Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion projects/kit/offline/src/lib/offline-request-policy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ export function shouldCommitOfflineCollection<T>(emission: OfflineReadEmission<r
}

/** Controls GET emission order between local replica and remote transport. */
export type OfflineReadStrategy = 'network-first' | 'local-first';
export type OfflineReadStrategy = 'network-first' | 'local-first' | 'fastest-first';

/** Product read policy backed by a local replica fallback for transport failures. */
export interface OfflineReadRequestPlan {
Expand All @@ -58,6 +58,9 @@ export interface OfflineReadRequestPlan {
* 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.
* `fastest-first` emits whichever usable local or remote response settles
* first. A winning local response is followed by remote revalidation; a
* winning remote response cancels and suppresses the slower local read.
*/
readStrategy?: OfflineReadStrategy;
/**
Expand Down
162 changes: 160 additions & 2 deletions projects/kit/offline/src/lib/offline.interceptor.spec.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,16 @@
import { HttpContext, HttpErrorResponse, HttpHeaders, HttpRequest, HttpResponse } from '@angular/common/http';
import {
HttpContext,
HttpErrorResponse,
type HttpEvent,
HttpEventType,
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 { finalize, firstValueFrom, of, Subject, throwError, type Observable } from 'rxjs';
import { finalize, firstValueFrom, Observable, of, Subject, throwError } from 'rxjs';
import { OfflineNetworkService } from './offline-network.service';
import { OFFLINE_MUTATION_PERSISTENCE_ENABLED } from './offline-mutation-persistence.service';
import { OfflineReplicaMutationCoordinator } from './offline-replica-mutation-coordinator';
Expand Down Expand Up @@ -559,6 +567,156 @@ describe('offlineInterceptor', () => {
});
});

describe('fastest-first read strategy', () => {
const fastestFirstPlan = (overrides: Partial<OfflineRequestPlan> = {}): 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<void>((resolve) => (releaseMutation = resolve));
let mutationStarted!: () => void;
const mutationReady = new Promise<void>((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<HttpResponse<unknown>>();
resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(async () => local) }));

const emissions: HttpResponse<unknown>[] = [];
let complete!: () => void;
const completed = new Promise<void>((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<HttpEvent<unknown>>();
resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(async () => local) }));

const emissions: HttpResponse<unknown>[] = [];
let complete!: () => void;
const completed = new Promise<void>((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<unknown>) => void;
const localReady = new Promise<HttpResponse<unknown>>((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<void>((resolve) => queueMicrotask(resolve));

expect(emissions).toEqual([remote]);
});

it('local miss時はremoteを待つ', async () => {
let emitRemote!: () => void;
const remote$ = new Observable<HttpResponse<unknown>>((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<void>((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<unknown>) => void;
const localReady = new Promise<HttpResponse<unknown>>((resolve) => (resolveLocal = resolve));
resolve.mockReturnValue(fastestFirstPlan({ readLocal: vi.fn(() => localReady) }));

const emissions: HttpResponse<unknown>[] = [];
let observedError: unknown;
const settled = new Promise<void>((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<void>((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 });
Expand Down
41 changes: 41 additions & 0 deletions projects/kit/offline/src/lib/offline.interceptor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,14 @@ import {
defer,
dematerialize,
EMPTY,
filter,
from,
map,
materialize,
NEVER,
of,
Observable,
race,
ReplaySubject,
take,
tap,
Expand Down Expand Up @@ -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)) {
Expand Down Expand Up @@ -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<unknown>,
plan: OfflineReadRequestPlan,
transport: () => Observable<HttpEvent<unknown>>,
fallback: OfflineRequestFallbackService,
errorHandler: ErrorHandler,
replicaMutations: OfflineReplicaMutationCoordinator,
): Observable<HttpEvent<unknown>> {
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<MaterializedTransport>() },
),
);
}

function resolveLocalAttempt(
plan: OfflineReadRequestPlan,
errorHandler: ErrorHandler,
Expand Down