mirror of
https://github.com/Sendouc/sendou.ink.git
synced 2026-09-30 15:17:43 -05:00
Chat without heartbeat
This commit is contained in:
52
apps/web/src/lib/features/chat/chat-utils.test.ts
Normal file
52
apps/web/src/lib/features/chat/chat-utils.test.ts
Normal file
@@ -0,0 +1,52 @@
|
||||
import { addHours, subHours } from "date-fns";
|
||||
import { describe, expect, test } from "vitest";
|
||||
import { dateToDatabaseTimestamp } from "#lib/utils/dates.ts";
|
||||
import { CHAT } from "./chat-constants.ts";
|
||||
import { nextLifecycleChangeAt, roomLifecycle } from "./chat-utils.ts";
|
||||
|
||||
const NOW = new Date("2026-01-01T12:00:00Z");
|
||||
const timestamp = (date: Date) => dateToDatabaseTimestamp(date);
|
||||
|
||||
describe("roomLifecycle", () => {
|
||||
test("is active when nothing schedules inactivity", () => {
|
||||
expect(roomLifecycle(null, NOW)).toBe("ACTIVE");
|
||||
});
|
||||
|
||||
test("is active while inactivity is still in the future", () => {
|
||||
expect(roomLifecycle(timestamp(addHours(NOW, 1)), NOW)).toBe("ACTIVE");
|
||||
});
|
||||
|
||||
test("is inactive between the inactivity and archival boundaries", () => {
|
||||
expect(roomLifecycle(timestamp(subHours(NOW, 1)), NOW)).toBe("INACTIVE");
|
||||
});
|
||||
|
||||
test("is archived once the archival window elapsed", () => {
|
||||
const inactiveAt = subHours(NOW, CHAT.INACTIVE_TO_ARCHIVED_HOURS + 1);
|
||||
expect(roomLifecycle(timestamp(inactiveAt), NOW)).toBe("ARCHIVED");
|
||||
});
|
||||
});
|
||||
|
||||
describe("nextLifecycleChangeAt", () => {
|
||||
test("returns null when nothing schedules inactivity", () => {
|
||||
expect(nextLifecycleChangeAt(null, NOW)).toBeNull();
|
||||
});
|
||||
|
||||
test("returns the inactivity boundary while the room is active", () => {
|
||||
const inactiveAt = addHours(NOW, 1);
|
||||
expect(nextLifecycleChangeAt(timestamp(inactiveAt), NOW)).toEqual(
|
||||
inactiveAt,
|
||||
);
|
||||
});
|
||||
|
||||
test("returns the archival boundary while the room is inactive", () => {
|
||||
const inactiveAt = subHours(NOW, 1);
|
||||
expect(nextLifecycleChangeAt(timestamp(inactiveAt), NOW)).toEqual(
|
||||
addHours(inactiveAt, CHAT.INACTIVE_TO_ARCHIVED_HOURS),
|
||||
);
|
||||
});
|
||||
|
||||
test("returns null once the room is archived", () => {
|
||||
const inactiveAt = subHours(NOW, CHAT.INACTIVE_TO_ARCHIVED_HOURS + 1);
|
||||
expect(nextLifecycleChangeAt(timestamp(inactiveAt), NOW)).toBeNull();
|
||||
});
|
||||
});
|
||||
@@ -22,6 +22,25 @@ export function roomLifecycle(
|
||||
return "ARCHIVED";
|
||||
}
|
||||
|
||||
/**
|
||||
* When the room's lifecycle next advances without anything happening to it, or
|
||||
* null once it is archived (the final state) or has no scheduled inactivity.
|
||||
*/
|
||||
export function nextLifecycleChangeAt(
|
||||
inactiveAt: number | null,
|
||||
now = new Date(),
|
||||
): Date | null {
|
||||
if (inactiveAt === null) return null;
|
||||
|
||||
const inactiveDate = databaseTimestampToDate(inactiveAt);
|
||||
if (now < inactiveDate) return inactiveDate;
|
||||
|
||||
const archivedDate = addHours(inactiveDate, CHAT.INACTIVE_TO_ARCHIVED_HOURS);
|
||||
if (now < archivedDate) return archivedDate;
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
/** Messages can be sent while the room is active or inactive, never archived. */
|
||||
export function canSendToRoom(lifecycle: ChatRoomLifecycle) {
|
||||
return lifecycle !== "ARCHIVED";
|
||||
|
||||
@@ -1,12 +1,17 @@
|
||||
import { error } from "@sveltejs/kit";
|
||||
import * as R from "remeda";
|
||||
import * as v from "valibot";
|
||||
import { getUser, requireUser } from "#lib/features/auth/user.server.ts";
|
||||
import * as Events from "#lib/server/events.ts";
|
||||
import { id } from "#lib/utils/schemas.ts";
|
||||
import { command, query } from "$app/server";
|
||||
import { command, getRequestEvent, query } from "$app/server";
|
||||
import { CHAT } from "./chat-constants.ts";
|
||||
import { publishChatRoom } from "./chat.server.ts";
|
||||
import { canSendToRoom, roomLifecycle } from "./chat-utils.ts";
|
||||
import {
|
||||
canSendToRoom,
|
||||
nextLifecycleChangeAt,
|
||||
roomLifecycle,
|
||||
} from "./chat-utils.ts";
|
||||
import * as ChatRepository from "./ChatRepository.server.ts";
|
||||
|
||||
/**
|
||||
@@ -18,16 +23,24 @@ export const getChatRoom = query.live(
|
||||
v.object({ chatRoomId: id }),
|
||||
async function* ({ chatRoomId }) {
|
||||
const user = requireUser();
|
||||
await requireRoomMember(chatRoomId, user);
|
||||
let room = await requireRoomMember(chatRoomId, user);
|
||||
|
||||
yield await roomSnapshot(chatRoomId);
|
||||
|
||||
for await (const _ of Events.subscribe(
|
||||
Events.chatRoomChannel(chatRoomId),
|
||||
{
|
||||
signal: getRequestEvent().request.signal,
|
||||
// the room turning inactive/archived is not published by anyone
|
||||
wakeAt: () => nextLifecycleChangeAt(room.inactiveAt),
|
||||
},
|
||||
)) {
|
||||
if (!(await ChatRepository.findRoomById(chatRoomId))) {
|
||||
const currentRoom = await ChatRepository.findRoomById(chatRoomId);
|
||||
if (!currentRoom) {
|
||||
return;
|
||||
}
|
||||
room = currentRoom;
|
||||
|
||||
yield await roomSnapshot(chatRoomId);
|
||||
}
|
||||
},
|
||||
@@ -63,15 +76,30 @@ export const getChatRooms = query.live(async function* () {
|
||||
return;
|
||||
}
|
||||
|
||||
yield await roomsSnapshot(user.id);
|
||||
let snapshot = await roomsSnapshot(user.id);
|
||||
yield snapshot;
|
||||
|
||||
for await (const _ of Events.subscribe(
|
||||
Events.chatRoomsOfUserChannel(user.id),
|
||||
{
|
||||
signal: getRequestEvent().request.signal,
|
||||
// rooms turn inactive, then drop off the list once archived, unpublished
|
||||
wakeAt: () => earliestLifecycleChangeAt(snapshot.rooms),
|
||||
},
|
||||
)) {
|
||||
yield await roomsSnapshot(user.id);
|
||||
snapshot = await roomsSnapshot(user.id);
|
||||
yield snapshot;
|
||||
}
|
||||
});
|
||||
|
||||
function earliestLifecycleChangeAt(rooms: { inactiveAt: number | null }[]) {
|
||||
const changes = rooms
|
||||
.map((room) => nextLifecycleChangeAt(room.inactiveAt))
|
||||
.filter((changeAt) => changeAt !== null);
|
||||
|
||||
return R.firstBy(changes, (changeAt) => changeAt.getTime()) ?? null;
|
||||
}
|
||||
|
||||
async function roomsSnapshot(userId: number) {
|
||||
const rooms = await ChatRepository.findRoomsOfUser(userId);
|
||||
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import * as v from "valibot";
|
||||
import { getUser, requireUser } from "#lib/features/auth/user.server.ts";
|
||||
import * as Events from "#lib/server/events.ts";
|
||||
import { command, query } from "$app/server";
|
||||
import { command, getRequestEvent, query } from "$app/server";
|
||||
import * as NotificationRepository from "./NotificationRepository.server.ts";
|
||||
import { notifyNotificationsChanged } from "./core/notify.server.ts";
|
||||
import { NOTIFICATIONS } from "./notifications-constants.ts";
|
||||
@@ -22,9 +22,9 @@ export const getNotifications = query.live(async function* () {
|
||||
|
||||
yield await peek(user.id);
|
||||
|
||||
for await (const _ of Events.subscribe(
|
||||
Events.notificationsChannel(user.id),
|
||||
)) {
|
||||
for await (const _ of Events.subscribe(Events.notificationsChannel(user.id), {
|
||||
signal: getRequestEvent().request.signal,
|
||||
})) {
|
||||
yield await peek(user.id);
|
||||
}
|
||||
});
|
||||
|
||||
@@ -11,6 +11,7 @@ import {
|
||||
participantIdsListFromAccepted,
|
||||
sideDisplayName,
|
||||
sideOfUser,
|
||||
trackingLocksAt,
|
||||
} from "./Scrim.ts";
|
||||
|
||||
type MockUser = { id: number };
|
||||
@@ -578,3 +579,36 @@ describe("isTrackingLocked", () => {
|
||||
).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe("trackingLocksAt", () => {
|
||||
const ONE_HOUR_MS = 60 * 60 * 1000;
|
||||
const lockWindowMs = SCRIM_TRACKING_AUTO_LOCK_HOURS * ONE_HOUR_MS;
|
||||
|
||||
test("returns null when no map list submitted yet", () => {
|
||||
expect(trackingLocksAt([], [])).toBeNull();
|
||||
});
|
||||
|
||||
test("returns the auto-lock window past the list submission", () => {
|
||||
const updatedAt = 1_000_000;
|
||||
expect(trackingLocksAt([], [{ updatedAt }])).toEqual(
|
||||
new Date(updatedAt * 1000 + lockWindowMs),
|
||||
);
|
||||
});
|
||||
|
||||
test("counts from the most recent reported map when there is one", () => {
|
||||
const updatedAt = 1_000_000;
|
||||
const reportedAt = updatedAt + 60 * 60;
|
||||
expect(trackingLocksAt([{ reportedAt }], [{ updatedAt }])).toEqual(
|
||||
new Date(reportedAt * 1000 + lockWindowMs),
|
||||
);
|
||||
});
|
||||
|
||||
test("returns a moment in the past once tracking is locked", () => {
|
||||
const updatedAt = dateToDatabaseTimestamp(
|
||||
new Date(Date.now() - lockWindowMs * 2),
|
||||
);
|
||||
expect(trackingLocksAt([], [{ updatedAt }])?.getTime()).toBeLessThan(
|
||||
Date.now(),
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { logger } from "@sendou/utils/logger";
|
||||
import { format, isWeekend } from "date-fns";
|
||||
import { addHours, format, isWeekend } from "date-fns";
|
||||
import * as R from "remeda";
|
||||
import type { Tables } from "#lib/server/db/tables.ts";
|
||||
import { databaseTimestampToDate } from "#lib/utils/dates.ts";
|
||||
@@ -168,6 +168,20 @@ export function isTrackingLocked(
|
||||
mapLists: Pick<Tables["ScrimMapList"], "updatedAt">[] = [],
|
||||
now: number = Date.now(),
|
||||
): boolean {
|
||||
const locksAt = trackingLocksAt(maps, mapLists);
|
||||
|
||||
return locksAt !== null && now > locksAt.getTime();
|
||||
}
|
||||
|
||||
/**
|
||||
* When map-by-map tracking auto-locks given the current activity, or null when
|
||||
* tracking is not active. The moment can be in the past (tracking is locked
|
||||
* already).
|
||||
*/
|
||||
export function trackingLocksAt(
|
||||
maps: Pick<Tables["ScrimMap"], "reportedAt">[] = [],
|
||||
mapLists: Pick<Tables["ScrimMapList"], "updatedAt">[] = [],
|
||||
): Date | null {
|
||||
const latestReported = R.firstBy(
|
||||
maps.filter((m) => m.reportedAt !== null),
|
||||
[(m) => m.reportedAt!, "desc"],
|
||||
@@ -176,11 +190,12 @@ export function isTrackingLocked(
|
||||
|
||||
const referenceSeconds =
|
||||
latestReported?.reportedAt ?? latestList?.updatedAt ?? null;
|
||||
if (referenceSeconds === null) return false;
|
||||
if (referenceSeconds === null) return null;
|
||||
|
||||
const elapsedHours = (now - referenceSeconds * 1000) / (60 * 60 * 1000);
|
||||
|
||||
return elapsedHours > SCRIM_TRACKING_AUTO_LOCK_HOURS;
|
||||
return addHours(
|
||||
databaseTimestampToDate(referenceSeconds),
|
||||
SCRIM_TRACKING_AUTO_LOCK_HOURS,
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -38,7 +38,7 @@ import {
|
||||
} from "#lib/utils/respond.server.ts";
|
||||
import { id } from "#lib/utils/schemas.ts";
|
||||
import { toDBBoolean } from "#lib/utils/sql.ts";
|
||||
import { command, query, requested } from "$app/server";
|
||||
import { command, getRequestEvent, query, requested } from "$app/server";
|
||||
import * as Scrim from "./Scrim.ts";
|
||||
import * as ScrimMapByMap from "./ScrimMapByMap.ts";
|
||||
import * as ScrimMapListRepository from "./ScrimMapListRepository.server.ts";
|
||||
@@ -159,10 +159,20 @@ export const getScrimsNewData = query(async () => {
|
||||
export const getScrim = query.live(
|
||||
v.object({ scrimPostId: id }),
|
||||
async function* ({ scrimPostId }) {
|
||||
yield await scrimSnapshot(scrimPostId);
|
||||
let snapshot = await scrimSnapshot(scrimPostId);
|
||||
yield snapshot;
|
||||
|
||||
for await (const _ of Events.subscribe(Events.scrimChannel(scrimPostId))) {
|
||||
yield await scrimSnapshot(scrimPostId);
|
||||
for await (const _ of Events.subscribe(Events.scrimChannel(scrimPostId), {
|
||||
signal: getRequestEvent().request.signal,
|
||||
// tracking auto-locking is the one change nobody publishes
|
||||
wakeAt: () =>
|
||||
Scrim.trackingLocksAt(
|
||||
snapshot.mapByMap.maps,
|
||||
snapshot.mapByMap.mapLists,
|
||||
),
|
||||
})) {
|
||||
snapshot = await scrimSnapshot(scrimPostId);
|
||||
yield snapshot;
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
import { describe, expect, test } from "vitest";
|
||||
import * as Events from "./events.ts";
|
||||
|
||||
const neverAborts = () => new AbortController().signal;
|
||||
|
||||
describe("Events.subscribe", () => {
|
||||
test("yields once per publish", async () => {
|
||||
const iterator = Events.subscribe("test-1", {
|
||||
heartbeatMs: Number.POSITIVE_INFINITY,
|
||||
});
|
||||
const iterator = Events.subscribe("test-1", { signal: neverAborts() });
|
||||
|
||||
const first = iterator.next();
|
||||
await waitForSubscriber("test-1");
|
||||
@@ -20,9 +20,7 @@ describe("Events.subscribe", () => {
|
||||
});
|
||||
|
||||
test("coalesces publishes that land before the consumer resumes", async () => {
|
||||
const iterator = Events.subscribe("test-2", {
|
||||
heartbeatMs: Number.POSITIVE_INFINITY,
|
||||
});
|
||||
const iterator = Events.subscribe("test-2", { signal: neverAborts() });
|
||||
|
||||
const first = iterator.next();
|
||||
await waitForSubscriber("test-2");
|
||||
@@ -33,46 +31,103 @@ describe("Events.subscribe", () => {
|
||||
|
||||
// all three publishes collapsed into that one yield, so the next
|
||||
// next() must block until a fresh publish
|
||||
let resolved = false;
|
||||
const blocked = iterator.next().then((result) => {
|
||||
resolved = true;
|
||||
return result;
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
expect(resolved).toBe(false);
|
||||
expect(await blocksFor(iterator.next())).toBe(true);
|
||||
|
||||
Events.publish("test-2");
|
||||
expect((await blocked).done).toBe(false);
|
||||
|
||||
await iterator.return(undefined);
|
||||
});
|
||||
|
||||
test("wakes on the heartbeat interval without a publish", async () => {
|
||||
const iterator = Events.subscribe("test-3", { heartbeatMs: 5 });
|
||||
test("never wakes on its own without a wakeAt deadline", async () => {
|
||||
const iterator = Events.subscribe("test-3", { signal: neverAborts() });
|
||||
|
||||
expect(await blocksFor(iterator.next())).toBe(true);
|
||||
|
||||
Events.publish("test-3");
|
||||
await iterator.return(undefined);
|
||||
});
|
||||
|
||||
test("wakes at the wakeAt deadline without a publish", async () => {
|
||||
const iterator = Events.subscribe("test-4", {
|
||||
signal: neverAborts(),
|
||||
wakeAt: () => new Date(Date.now() + 5),
|
||||
});
|
||||
|
||||
expect((await iterator.next()).done).toBe(false);
|
||||
|
||||
await iterator.return(undefined);
|
||||
});
|
||||
|
||||
test("unsubscribes when the consumer stops iterating", async () => {
|
||||
const iterator = Events.subscribe("test-4", {
|
||||
heartbeatMs: Number.POSITIVE_INFINITY,
|
||||
test("ignores a wakeAt deadline that already passed", async () => {
|
||||
const iterator = Events.subscribe("test-5", {
|
||||
signal: neverAborts(),
|
||||
wakeAt: () => new Date(Date.now() - 60_000),
|
||||
});
|
||||
|
||||
expect(await blocksFor(iterator.next())).toBe(true);
|
||||
|
||||
Events.publish("test-5");
|
||||
await iterator.return(undefined);
|
||||
});
|
||||
|
||||
test("re-reads wakeAt before every sleep", async () => {
|
||||
const deadlines: (Date | null)[] = [null, new Date(Date.now() + 5)];
|
||||
let call = 0;
|
||||
|
||||
const iterator = Events.subscribe("test-6", {
|
||||
signal: neverAborts(),
|
||||
wakeAt: () => deadlines[call++] ?? null,
|
||||
});
|
||||
|
||||
const first = iterator.next();
|
||||
await waitForSubscriber("test-4");
|
||||
expect(Events.subscriberCount("test-4")).toBe(1);
|
||||
await waitForSubscriber("test-6");
|
||||
Events.publish("test-6");
|
||||
await first;
|
||||
|
||||
Events.publish("test-4");
|
||||
// the second sleep gets the deadline, so no publish is needed
|
||||
expect((await iterator.next()).done).toBe(false);
|
||||
expect(call).toBe(2);
|
||||
|
||||
await iterator.return(undefined);
|
||||
});
|
||||
|
||||
test("unsubscribes when the consumer stops iterating", async () => {
|
||||
const iterator = Events.subscribe("test-7", { signal: neverAborts() });
|
||||
|
||||
const first = iterator.next();
|
||||
await waitForSubscriber("test-7");
|
||||
expect(Events.subscriberCount("test-7")).toBe(1);
|
||||
|
||||
Events.publish("test-7");
|
||||
await first;
|
||||
await iterator.return(undefined);
|
||||
|
||||
expect(Events.subscriberCount("test-4")).toBe(0);
|
||||
expect(Events.subscriberCount("test-7")).toBe(0);
|
||||
});
|
||||
|
||||
test("ends and unsubscribes when the signal aborts mid-sleep", async () => {
|
||||
const controller = new AbortController();
|
||||
const iterator = Events.subscribe("test-8", { signal: controller.signal });
|
||||
|
||||
const first = iterator.next();
|
||||
await waitForSubscriber("test-8");
|
||||
controller.abort();
|
||||
|
||||
expect((await first).done).toBe(true);
|
||||
expect(Events.subscriberCount("test-8")).toBe(0);
|
||||
});
|
||||
|
||||
test("never starts when the signal aborted beforehand", async () => {
|
||||
const controller = new AbortController();
|
||||
controller.abort();
|
||||
|
||||
const iterator = Events.subscribe("test-9", { signal: controller.signal });
|
||||
|
||||
expect((await iterator.next()).done).toBe(true);
|
||||
expect(Events.subscriberCount("test-9")).toBe(0);
|
||||
});
|
||||
|
||||
test("publishing to a channel with no subscribers is a no-op", () => {
|
||||
expect(() => Events.publish("test-5")).not.toThrow();
|
||||
expect(() => Events.publish("test-10")).not.toThrow();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -81,3 +136,13 @@ async function waitForSubscriber(channel: string) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 1));
|
||||
}
|
||||
}
|
||||
|
||||
async function blocksFor(next: Promise<unknown>) {
|
||||
let settled = false;
|
||||
next.then(() => {
|
||||
settled = true;
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||
|
||||
return !settled;
|
||||
}
|
||||
|
||||
@@ -7,9 +7,17 @@
|
||||
* Subscribers coalesce: publishes that land while the subscriber is busy
|
||||
* producing a snapshot collapse into one pending wake-up, so generators always
|
||||
* yield snapshots, never an event log.
|
||||
*
|
||||
* There is deliberately no periodic wake-up: SvelteKit's live-query transport
|
||||
* writes its own `: keep-alive` SSE comment on an idle timer, so a quiet
|
||||
* subscription does not need to build (and throw away) snapshots to hold the
|
||||
* connection open. Snapshots that go stale purely with the passage of time —
|
||||
* a chat room turning archived, scrim tracking auto-locking — instead pass
|
||||
* `wakeAt` and get woken once, at the boundary.
|
||||
*/
|
||||
|
||||
const DEFAULT_HEARTBEAT_MS = 30_000;
|
||||
/** Node clamps longer delays to 1ms, which would spin; waking early is harmless. */
|
||||
const MAX_TIMEOUT_MS = 2_147_483_647;
|
||||
|
||||
type Wake = () => void;
|
||||
|
||||
@@ -26,14 +34,28 @@ export function publish(channel: string) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Yields once per (coalesced) publish on the channel, and additionally on a
|
||||
* heartbeat interval so long-lived streams keep writing through proxies that
|
||||
* idle-timeout quiet connections. Unsubscribes when the consumer stops
|
||||
* iterating (client disconnect unwinds the generator via `finally`).
|
||||
* Yields once per (coalesced) publish on the channel, until `signal` aborts.
|
||||
*
|
||||
* `signal` is required rather than optional because a generator parked on an
|
||||
* `await` cannot be unwound: `return()` on it is queued until it next reaches a
|
||||
* `yield`, so without the abort waking the sleep, a disconnected client's
|
||||
* subscriber would stay in the channel forever.
|
||||
*
|
||||
* `wakeAt` is consulted before every sleep and may name the moment the
|
||||
* consumer's latest snapshot goes stale on its own, waking it then as well; a
|
||||
* deadline already in the past is ignored, since that transition is part of the
|
||||
* snapshot just produced.
|
||||
*/
|
||||
export async function* subscribe(
|
||||
channel: string,
|
||||
{ heartbeatMs = DEFAULT_HEARTBEAT_MS }: { heartbeatMs?: number } = {},
|
||||
{
|
||||
signal,
|
||||
wakeAt,
|
||||
}: {
|
||||
/** Request abort signal, i.e. `getRequestEvent().request.signal`. */
|
||||
signal: AbortSignal;
|
||||
wakeAt?: () => Date | null;
|
||||
},
|
||||
): AsyncGenerator<void> {
|
||||
let pending = false;
|
||||
let wake: Wake | null = null;
|
||||
@@ -51,15 +73,24 @@ export async function* subscribe(
|
||||
wakes.add(listener);
|
||||
|
||||
try {
|
||||
while (true) {
|
||||
while (!signal.aborted) {
|
||||
if (!pending) {
|
||||
const staleInMs = msUntilStale(wakeAt?.() ?? null);
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
|
||||
await new Promise<void>((resolve) => {
|
||||
wake = resolve;
|
||||
if (heartbeatMs !== Number.POSITIVE_INFINITY) {
|
||||
setTimeout(resolve, heartbeatMs).unref();
|
||||
wake = () => resolve();
|
||||
signal.addEventListener("abort", wake, { once: true });
|
||||
if (staleInMs !== null) {
|
||||
timer = setTimeout(wake, staleInMs).unref();
|
||||
}
|
||||
});
|
||||
|
||||
clearTimeout(timer);
|
||||
if (wake) signal.removeEventListener("abort", wake);
|
||||
wake = null;
|
||||
|
||||
if (signal.aborted) return;
|
||||
}
|
||||
pending = false;
|
||||
yield;
|
||||
@@ -72,6 +103,15 @@ export async function* subscribe(
|
||||
}
|
||||
}
|
||||
|
||||
function msUntilStale(deadline: Date | null) {
|
||||
if (deadline === null) return null;
|
||||
|
||||
const delay = deadline.getTime() - Date.now();
|
||||
if (delay <= 0) return null;
|
||||
|
||||
return Math.min(delay, MAX_TIMEOUT_MS);
|
||||
}
|
||||
|
||||
/** Number of active subscribers on a channel (for tests & diagnostics). */
|
||||
export function subscriberCount(channel: string) {
|
||||
return channels.get(channel)?.size ?? 0;
|
||||
|
||||
Reference in New Issue
Block a user