From 32a14fabb50419d8c623b4fda49f8c83fb97adf1 Mon Sep 17 00:00:00 2001 From: Kalle <38327916+Sendouc@users.noreply.github.com> Date: Tue, 18 Aug 2026 18:30:44 +0300 Subject: [PATCH] Chat without heartbeat --- .../src/lib/features/chat/chat-utils.test.ts | 52 ++++++++ apps/web/src/lib/features/chat/chat-utils.ts | 19 +++ apps/web/src/lib/features/chat/chat.remote.ts | 40 +++++- .../notifications/notifications.remote.ts | 8 +- .../web/src/lib/features/scrims/Scrim.test.ts | 34 ++++++ apps/web/src/lib/features/scrims/Scrim.ts | 25 +++- .../src/lib/features/scrims/scrims.remote.ts | 18 ++- apps/web/src/lib/server/events.test.ts | 115 ++++++++++++++---- apps/web/src/lib/server/events.ts | 60 +++++++-- 9 files changed, 317 insertions(+), 54 deletions(-) create mode 100644 apps/web/src/lib/features/chat/chat-utils.test.ts diff --git a/apps/web/src/lib/features/chat/chat-utils.test.ts b/apps/web/src/lib/features/chat/chat-utils.test.ts new file mode 100644 index 000000000..7c4e9c077 --- /dev/null +++ b/apps/web/src/lib/features/chat/chat-utils.test.ts @@ -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(); + }); +}); diff --git a/apps/web/src/lib/features/chat/chat-utils.ts b/apps/web/src/lib/features/chat/chat-utils.ts index 785ffc8b7..d1fee942a 100644 --- a/apps/web/src/lib/features/chat/chat-utils.ts +++ b/apps/web/src/lib/features/chat/chat-utils.ts @@ -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"; diff --git a/apps/web/src/lib/features/chat/chat.remote.ts b/apps/web/src/lib/features/chat/chat.remote.ts index e507d580e..eac89a1e2 100644 --- a/apps/web/src/lib/features/chat/chat.remote.ts +++ b/apps/web/src/lib/features/chat/chat.remote.ts @@ -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); diff --git a/apps/web/src/lib/features/notifications/notifications.remote.ts b/apps/web/src/lib/features/notifications/notifications.remote.ts index f314a6c20..49baa244e 100644 --- a/apps/web/src/lib/features/notifications/notifications.remote.ts +++ b/apps/web/src/lib/features/notifications/notifications.remote.ts @@ -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); } }); diff --git a/apps/web/src/lib/features/scrims/Scrim.test.ts b/apps/web/src/lib/features/scrims/Scrim.test.ts index 5e961adf7..43a0629e8 100644 --- a/apps/web/src/lib/features/scrims/Scrim.test.ts +++ b/apps/web/src/lib/features/scrims/Scrim.test.ts @@ -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(), + ); + }); +}); diff --git a/apps/web/src/lib/features/scrims/Scrim.ts b/apps/web/src/lib/features/scrims/Scrim.ts index cd423d68a..89cba3b70 100644 --- a/apps/web/src/lib/features/scrims/Scrim.ts +++ b/apps/web/src/lib/features/scrims/Scrim.ts @@ -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[] = [], 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[] = [], + mapLists: Pick[] = [], +): 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, + ); } /** diff --git a/apps/web/src/lib/features/scrims/scrims.remote.ts b/apps/web/src/lib/features/scrims/scrims.remote.ts index 0901cb126..5f80f4126 100644 --- a/apps/web/src/lib/features/scrims/scrims.remote.ts +++ b/apps/web/src/lib/features/scrims/scrims.remote.ts @@ -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; } }, ); diff --git a/apps/web/src/lib/server/events.test.ts b/apps/web/src/lib/server/events.test.ts index 0397f4ff9..952b4a501 100644 --- a/apps/web/src/lib/server/events.test.ts +++ b/apps/web/src/lib/server/events.test.ts @@ -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) { + let settled = false; + next.then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 20)); + + return !settled; +} diff --git a/apps/web/src/lib/server/events.ts b/apps/web/src/lib/server/events.ts index b37146012..7dc5fdca0 100644 --- a/apps/web/src/lib/server/events.ts +++ b/apps/web/src/lib/server/events.ts @@ -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 { 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 | undefined; + await new Promise((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;