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
14 changes: 9 additions & 5 deletions packages/presence/src/routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -175,19 +175,23 @@ export function createPresenceRoutes(
const surface = c.req.param("surface");
const key = { tenantId: tenant.id, surface };

registry.sweepStale(heartbeatTimeoutMs, now());

let patch: PresenceStatePatch = {};
if (body.cursor !== undefined) patch = { ...patch, cursor: body.cursor };
if (body.typing !== undefined) patch = { ...patch, typing: body.typing };
const states = registry.heartbeat(key, principal.id, patch, now());
if (states === undefined) {
// Refresh this principal's `lastSeenAt` *before* sweeping: this
// request arriving is itself proof of liveness, so the sweep below
// must judge staleness against the fresh timestamp, never the
// pre-request one — otherwise a heartbeat landing a moment past
// `heartbeatTimeoutMs` (ordinary jitter) would evict its own sender.
const heartbeatResult = registry.heartbeat(key, principal.id, patch, now());
if (heartbeatResult === undefined) {
return c.json(
errorEnvelope("not_joined", "principal has not joined this room"),
404,
);
}
return c.json({ members: states });
registry.sweepStale(heartbeatTimeoutMs, now());
return c.json({ members: registry.states(key) });
});

app.post(
Expand Down
92 changes: 92 additions & 0 deletions packages/presence/test/routes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,98 @@ describe("presence routes", () => {
expect(response.status).toBe(404);
});

test("a heartbeat arriving just past the timeout boundary does not evict its own sender", async () => {
let clock = 0;
const registry = createPresenceRoomRegistry();
const app = mountAs(
createPresenceRoutes({
registry,
requireGrant: allowAll,
now: () => clock,
}),
{
tenantId: "tnt_a",
principalId: "prn_alice",
},
);

await app.request(`/rooms/${SURFACE}/join`, {
method: "POST",
headers: { "content-type": "application/json" },
body: "{}",
});

// One second past the default 45s heartbeat timeout: normal jitter
// (a slow network tick, a throttled background tab), not a genuinely
// stale client — the heartbeat that arrives now is itself proof the
// sender is alive.
clock = 46_000;
const response = await app.request(`/rooms/${SURFACE}/heartbeat`, {
method: "POST",
headers: { "content-type": "application/json" },
body: "{}",
});

expect(response.status).toBe(200);
const body = (await response.json()) as {
members: JoinResponseBody["members"];
};
expect(body.members.map((m) => m.principalId)).toEqual(["prn_alice"]);
});

test("a heartbeat still sweeps a genuinely stale, different principal out of the response", async () => {
let clock = 0;
const registry = createPresenceRoomRegistry();
const alice = mountAs(
createPresenceRoutes({
registry,
requireGrant: allowAll,
now: () => clock,
}),
{
tenantId: "tnt_a",
principalId: "prn_alice",
},
);
const bob = mountAs(
createPresenceRoutes({
registry,
requireGrant: allowAll,
now: () => clock,
}),
{
tenantId: "tnt_a",
principalId: "prn_bob",
},
);

await alice.request(`/rooms/${SURFACE}/join`, {
method: "POST",
headers: { "content-type": "application/json" },
body: "{}",
});
await bob.request(`/rooms/${SURFACE}/join`, {
method: "POST",
headers: { "content-type": "application/json" },
body: "{}",
});

// Bob never heartbeats again; alice's next heartbeat lands well past
// the timeout for bob, but only 1ms past it for herself.
clock = 46_000;
const response = await alice.request(`/rooms/${SURFACE}/heartbeat`, {
method: "POST",
headers: { "content-type": "application/json" },
body: "{}",
});

expect(response.status).toBe(200);
const body = (await response.json()) as {
members: JoinResponseBody["members"];
};
expect(body.members.map((m) => m.principalId)).toEqual(["prn_alice"]);
});

test("an invalid join body is rejected with 400", async () => {
const app = mountAs(createPresenceRoutes({ requireGrant: allowAll }), {
tenantId: "tnt_a",
Expand Down
Loading