Skip to content
Open
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
140 changes: 109 additions & 31 deletions product-sdk/packages/individuality/src/watch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,11 @@
* best block, decoded by the same functions the pinned reads use.
*
* Each watch returns the function that stops it. A value that fails to decode goes
* to `onError` and the watch keeps running, while a failed subscription goes to
* `onError` and ends. Nothing arrives after the stop function is called.
* to `onError` and the watch keeps running. A failed subscription goes to `onError`
* too, and the watch subscribes again after a wait that doubles from one second to
* thirty, back to one second once a value arrives. PAPI ends a `watchValue` on its
* first failed query, and a host that drops one answer would otherwise silence the
* watch for good. Nothing arrives after the stop function is called.
*
* PAPI emits once per best block whether or not the value changed, measured against
* paseo on 2026-09-25, so a watch only delivers a decoded value that differs from the
Expand Down Expand Up @@ -112,7 +115,7 @@ export function watchCurrentGame(
onError: WatchErrorHandler,
): () => void {
return watch(
chain.individuality.query.Game.Game.watchValue({ at: "best" }),
() => chain.individuality.query.Game.Game.watchValue({ at: "best" }),
(raw) => (raw === undefined ? null : toCurrentGame(raw)),
onValue,
onError,
Expand All @@ -127,7 +130,7 @@ export function watchPlayer(
onError: WatchErrorHandler,
): () => void {
return watch(
chain.individuality.query.Game.Players.watchValue(options.player, { at: "best" }),
() => chain.individuality.query.Game.Players.watchValue(options.player, { at: "best" }),
(raw) => (raw === undefined ? null : { registered: raw.registered }),
onValue,
onError,
Expand All @@ -142,7 +145,8 @@ export function watchParticipant(
onError: WatchErrorHandler,
): () => void {
return watch(
chain.individuality.query.Score.Participants.watchValue(options.player, { at: "best" }),
() =>
chain.individuality.query.Score.Participants.watchValue(options.player, { at: "best" }),
(raw) => (raw === undefined ? null : toPersonhoodParticipant(raw)),
onValue,
onError,
Expand All @@ -152,41 +156,55 @@ export function watchParticipant(
/** Distinct from every storage value, `undefined` included, so the first emission is decoded. */
const UNSEEN = Symbol("unseen");

const RESUBSCRIBE_FIRST_MS = 1_000;
const RESUBSCRIBE_MAX_MS = 30_000;

function watch<Raw, Value>(
source: WatchedValue<Raw>,
source: () => WatchedValue<Raw>,
decode: (raw: Raw) => Value,
onValue: (value: Value, block: WatchedBlock) => void,
onError: WatchErrorHandler,
): () => void {
let stopped = false;
let lastRaw: Raw | typeof UNSEEN = UNSEEN;
let lastDelivered: string | undefined;
const subscription = source.subscribe({
next: ({ block, value }) => {
// PAPI hands back the same object while the stored bytes are unchanged.
if (stopped || value === lastRaw) return;
lastRaw = value;
let decoded: Value;
try {
decoded = decode(value);
} catch (cause) {
let subscription: { unsubscribe(): void } | undefined;
let resubscribe: ReturnType<typeof setTimeout> | undefined;
let wait = RESUBSCRIBE_FIRST_MS;
const subscribe = () => {
let lastRaw: Raw | typeof UNSEEN = UNSEEN;
subscription = source().subscribe({
next: ({ block, value }) => {
// PAPI hands back the same object while the stored bytes are unchanged.
if (stopped || value === lastRaw) return;
lastRaw = value;
wait = RESUBSCRIBE_FIRST_MS;
let decoded: Value;
try {
decoded = decode(value);
} catch (cause) {
onError(normalizeError(cause, ProductIndividualityError));
return;
}
// New bytes can decode to the same value when only a field the decoder
// drops has moved, such as an offchain-worker cursor.
const delivered = JSON.stringify(decoded, jsonSerialize);
if (delivered === lastDelivered) return;
lastDelivered = delivered;
onValue(decoded, { blockHash: block.hash, blockNumber: block.number });
},
error: (cause) => {
if (stopped) return;
onError(normalizeError(cause, ProductIndividualityError));
return;
}
// New bytes can decode to the same value when only a field the decoder
// drops has moved, such as an offchain-worker cursor.
const delivered = JSON.stringify(decoded, jsonSerialize);
if (delivered === lastDelivered) return;
lastDelivered = delivered;
onValue(decoded, { blockHash: block.hash, blockNumber: block.number });
},
error: (cause) => {
if (!stopped) onError(normalizeError(cause, ProductIndividualityError));
},
});
resubscribe = setTimeout(subscribe, wait);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We never unsubscribe the errored subscription before retrying, so this only holds if PAPI tears the stream down on error. I'd call subscription?.unsubscribe() here to be safe, and the fake can't catch a leak since it shares one unsubscribe across subscribes.

wait = Math.min(wait * 2, RESUBSCRIBE_MAX_MS);
},
});
};
subscribe();
return () => {
stopped = true;
subscription.unsubscribe();
clearTimeout(resubscribe);
subscription?.unsubscribe();
};
}

Expand Down Expand Up @@ -355,13 +373,73 @@ if (import.meta.vitest) {
test("sends a failed subscription to onError as a package error with its cause", () => {
const { chain, game } = fakeChain();
const onError = vi.fn();
watchCurrentGame(chain, vi.fn(), onError);
const stop = watchCurrentGame(chain, vi.fn(), onError);
game.fail(new Error("disconnected"));
stop();
const error = onError.mock.calls[0]?.[0] as ProductIndividualityError;
expect(error).toBeInstanceOf(ProductIndividualityError);
expect((error.cause as Error).message).toBe("disconnected");
});

test("subscribes again after a failed subscription, and delivers only a changed value", () => {
vi.useFakeTimers();
try {
const { chain, calls, game } = fakeChain();
const onValue = vi.fn();
const stop = watchCurrentGame(chain, onValue, vi.fn());
game.emit(rawGame());
game.fail(new Error("timed out"));
vi.advanceTimersByTime(999);
expect(calls).toHaveLength(1);
vi.advanceTimersByTime(1);
expect(calls).toHaveLength(2);
game.emit(rawGame(), block(11));
game.emit(rawGame({ state: { type: "PlayerProcess" } }), block(12));
stop();
expect(onValue).toHaveBeenCalledTimes(2);
expect(onValue.mock.calls[1]?.[0]).toMatchObject({ phase: "PlayerProcess" });
} finally {
vi.useRealTimers();
}
});

test("doubles the wait while subscriptions keep failing, up to thirty seconds", () => {
vi.useFakeTimers();
try {
const { chain, calls, game } = fakeChain();
const stop = watchCurrentGame(chain, vi.fn(), vi.fn());
const waits: number[] = [];
for (let failure = 0; failure < 7; failure++) {
game.fail(new Error("timed out"));
const before = calls.length;
let waited = 0;
while (calls.length === before && waited < 60_000) {
vi.advanceTimersByTime(1_000);
waited += 1_000;
}
waits.push(waited);
}
stop();
expect(waits).toEqual([1_000, 2_000, 4_000, 8_000, 16_000, 30_000, 30_000]);
} finally {
vi.useRealTimers();
}
});

test("stop cancels a subscription that is waiting to start again", () => {
vi.useFakeTimers();
try {
const { chain, calls, game } = fakeChain();
const stop = watchCurrentGame(chain, vi.fn(), vi.fn());
game.fail(new Error("timed out"));
stop();
vi.advanceTimersByTime(60_000);
expect(calls).toHaveLength(1);
} finally {
vi.useRealTimers();
}
});

test("stops the subscription, and delivers nothing after it", () => {
const { chain, game } = fakeChain();
const onValue = vi.fn();
Expand Down
12 changes: 12 additions & 0 deletions product-sdk/pending-changesets/individuality-watch-resubscribe.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
---
"@parity/product-sdk-individuality": minor
"@parity/product-sdk": minor
---

`watchCurrentGame`, `watchPlayer` and `watchParticipant` now subscribe again
after a failed subscription instead of ending. PAPI ends a `watchValue` on its
first failed query, so one lost or failed host answer used to silence the watch
for good. The failure still goes to `onError`, and the watch subscribes again
after a wait that doubles from one second to thirty and drops back to one second
once a value arrives. A value that matches the last one delivered is not
delivered again after the watch resubscribes.
Loading