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
20 changes: 20 additions & 0 deletions .changeset/olive-donkeys-jam.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
---
'@haverstack/adapter-api': patch
---

Stop the change-feed client reconnecting against a refusal that will repeat.

`isFatalFeedError` ended the reconnect loop only for an unrenewable
credential (401) and an authorization refusal (403). Every other refusal the
server faulted the request for — a malformed cursor or filter answered
`400 bad_request`, say — was treated as transient, so `subscribeChanges()`
retried it with backoff indefinitely, settling into an attempt roughly every
15 seconds and reporting the same error to `onError` each time. The
subscriber was never told to stop, and `onReset` never fired, so the
application had nothing to reconcile from either.

The predicate now decides on the wire status: a `4xx` ends the loop, since
the reconnect sends the same request and would be refused the same way. A
`5xx` still reconnects, which is what keeps `timeout` — the answer a server
gives while shedding query load — from turning a busy server into a
permanently dead subscription.
2 changes: 2 additions & 0 deletions docs/spec/wire-format.md
Original file line number Diff line number Diff line change
Expand Up @@ -539,6 +539,8 @@ data: {"reason":"cursor_expired"}

**Reconnection is the client's job, with exponential backoff and jitter.** A server restart otherwise produces a synchronized reconnect stampede from every client it dropped.

**A client stops reconnecting when the answer was `4xx`, and keeps reconnecting when it was `5xx`.** The reconnect sends the same request, so a status faulting that request — a malformed cursor or filter, an unrenewable credential, an authorization refusal — will be answered identically however long the client waits, and backing off only spins. A `5xx` says the server could not serve a request it did not fault, which is the case backoff exists for; `timeout` is the answer a server gives while [shedding query load](#bounding-query-cost) and a client that gave up on it would turn a busy server into a dead subscription. A client that stops reports the error to its subscriber first; the repair is to subscribe again.

### Permission scoping

**A connection delivers the events its token's session may read, and nothing else.** The predicate is literally `canRead` applied per event — no second vocabulary, no feed-specific ACL. A server subscribes **unscoped** at the storage owner and fans out per connection, filtering each through the `ScopedStack` its token's session names via `Stack.forSession()`, taking the `(principalId, subjectId)` pair whole. Delegated authority is then the ordinary [intersection](./access-control.md#delegation-principal-and-subject), inherited rather than reimplemented.
Expand Down
25 changes: 15 additions & 10 deletions packages/adapter-api/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@
* patchContent()/deleteRecord()/etc.'s expectedVersion option.
*/

import { StackQueryError, StackPermissionError } from '@haverstack/core';
import { StackError, StackQueryError } from '@haverstack/core';
import type {
StackAdapter,
StackRecord,
Expand Down Expand Up @@ -59,6 +59,7 @@ import {
isRetryableAuthError,
isValidSeq,
supportsChangeFeed,
WIRE_ERROR_STATUS,
supportsDidChallenge,
CHANGE_FRAME_READY,
CHANGE_FRAME_RECORD,
Expand Down Expand Up @@ -484,16 +485,20 @@ const reconnectDelay = (attempt: number): number => {
const sleep = (ms: number): Promise<void> => new Promise((resolve) => setTimeout(resolve, ms));

/**
* Whether a feed error is one reconnecting cannot recover. An
* authentication failure (a 401 whose re-auth did not renew, or no
* credential to renew with) and an authorization refusal (403) will reject
* the next connection identically, so retrying only spins. A connection or
* server error is transient and reconnects. The stream is closed by
* returning from the pump; the subscriber has already been told via
* onError.
* Whether a feed error is one reconnecting cannot recover. A 4xx faults the
* request, and the reconnect sends the same one, so retrying only spins. A
* 5xx is the server's own trouble and may clear — a shed-load `timeout`
* reconnects. The stream is closed by returning from the pump; the
* subscriber has already been told via onError.
*/
const isFatalFeedError = (err: unknown): boolean =>
err instanceof APIAdapterAuthError || err instanceof StackPermissionError;
const isFatalFeedError = (err: unknown): boolean => {
if (err instanceof APIAdapterAuthError) return true;
if (err instanceof StackError) {
const status = WIRE_ERROR_STATUS[err.code];
return status >= 400 && status < 500;
}
return false;
};

// -------------------------------------------------------
// Challenge–response handshake
Expand Down
63 changes: 63 additions & 0 deletions packages/adapter-api/tests/change-feed.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import {
APIAdapterAuthError,
} from '../src/index.js';
import { WIRE_PROTOCOL_VERSION } from '@haverstack/wire-types';
import { StackQueryError, StackTimeoutError } from '@haverstack/core';
import type { RecordChange } from '@haverstack/core';

const BASE_URL = 'https://stack.example.com';
Expand Down Expand Up @@ -486,6 +487,68 @@ describe('reconnection', () => {
stop();
});

// A request the server faulted is the same request on the next
// connection, so the loop ends rather than spinning against a verdict
// that will not change.
test('stops reconnecting after a 4xx the reconnect would only repeat', async () => {
vi.useFakeTimers();
const adapter = await openAdapter();
const first = feed();
mockFetch.mockResolvedValueOnce(first.response);

const onError = vi.fn();
const subscription = adapter.subscribeChanges({ onError }, () => {});
first.write(READY);
const stop = await subscription;

mockFetch.mockResolvedValueOnce(
jsonResponse({ error: { code: 'bad_request', message: 'Invalid cursor' } }, 400),
);
first.end();

await vi.advanceTimersByTimeAsync(60_000);
await vi.waitFor(() => expect(onError).toHaveBeenCalledOnce());
expect(onError.mock.calls[0]![0]).toBeInstanceOf(StackQueryError);

// No further reconnect: the fetch count holds at discovery + first +
// the refused reconnect, even after more time passes.
const calls = mockFetch.mock.calls.length;
await vi.advanceTimersByTimeAsync(120_000);
expect(mockFetch.mock.calls.length).toBe(calls);
stop();
});

// 503 is a server shedding load, not a verdict on the request — being
// reconnected against is the whole point of it.
test('keeps reconnecting after a 503 the server may recover from', async () => {
vi.useFakeTimers();
const adapter = await openAdapter();
const first = feed();
mockFetch.mockResolvedValueOnce(first.response);

const onError = vi.fn();
const subscription = adapter.subscribeChanges({ onError }, () => {});
first.write(READY);
const stop = await subscription;

mockFetch.mockResolvedValueOnce(
jsonResponse({ error: { code: 'timeout', message: 'Shedding load' } }, 503),
);
const recovered = feed();
mockFetch.mockResolvedValueOnce(recovered.response);
first.end();

await vi.advanceTimersByTimeAsync(60_000);
await vi.waitFor(() => expect(onError).toHaveBeenCalledOnce());
expect(onError.mock.calls[0]![0]).toBeInstanceOf(StackTimeoutError);

// discovery + first + the 503 + the attempt that reaches the server
// again: the subscription outlived the refusal.
await vi.advanceTimersByTimeAsync(120_000);
await vi.waitFor(() => expect(mockFetch.mock.calls.length).toBe(4));
stop();
});

test('stops reconnecting once unsubscribed, and aborts the open stream', async () => {
vi.useFakeTimers();
const adapter = await openAdapter();
Expand Down
Loading