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
9 changes: 9 additions & 0 deletions src/db/sqlite/connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,8 @@ function createSchema(d: DatabaseSync): void {
install_id TEXT NOT NULL,
capability TEXT NOT NULL,
muted TEXT NOT NULL DEFAULT '[]',
loud TEXT NOT NULL DEFAULT '[]',
everyone INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (server_user_id, install_id)
Expand Down Expand Up @@ -715,6 +717,13 @@ function runMigrations(d: DatabaseSync): void {
if (!hasColumn(d, "push_devices", "muted")) {
d.exec("ALTER TABLE push_devices ADD COLUMN muted TEXT NOT NULL DEFAULT '[]'");
}
// Conversations at "All messages", and whether @everyone gets through (GRYT-1696). Old phones send neither.
if (!hasColumn(d, "push_devices", "loud")) {
d.exec("ALTER TABLE push_devices ADD COLUMN loud TEXT NOT NULL DEFAULT '[]'");
}
if (!hasColumn(d, "push_devices", "everyone")) {
d.exec("ALTER TABLE push_devices ADD COLUMN everyone INTEGER NOT NULL DEFAULT 0");
}

const cols = d.prepare("PRAGMA table_info(users)").all() as { name: string }[];
const colNames = new Set(cols.map((c) => c.name));
Expand Down
2 changes: 1 addition & 1 deletion src/db/sqlite/pushDevices.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ describe("push devices", () => {
it("keeps one row per install and replaces its capability", () => {
savePushDevice("u1", "install-1", cap("a"));
savePushDevice("u1", "install-1", cap("b"));
assert.deepEqual(listPushDevices("u1"), [{ installId: "install-1", capability: cap("b"), muted: new Set() }]);
assert.deepEqual(listPushDevices("u1"), [{ installId: "install-1", capability: cap("b"), muted: new Set(), loud: new Set(), everyone: false }]);
removePushDevice("u1", "install-1");
assert.deepEqual(listPushDevices("u1"), []);
});
Expand Down
46 changes: 38 additions & 8 deletions src/db/sqlite/pushDevices.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,17 @@ export interface PushDevice {
capability: string;
/** Conversations muted on that phone, which never reach the relay (GRYT-1689). */
muted: ReadonlySet<string>;
/** Conversations at "All messages" on that phone: every message there wakes it (GRYT-1696). */
loud: ReadonlySet<string>;
/** Whether @everyone and @here wake it, which "Suppress @everyone" turns off. */
everyone: boolean;
}

/** What the phone said about its notification settings when it last checked in. */
export interface PushSettings {
muted?: readonly string[];
loud?: readonly string[];
everyone?: boolean;
}

/** More than this and the oldest goes. Nobody has ten phones; a reinstall loop might. */
Expand All @@ -16,7 +27,7 @@ export const PUSH_DEVICE_STALE_DAYS = 30;

const DAY_MS = 24 * 60 * 60 * 1000;

function parseMuted(raw: string): ReadonlySet<string> {
function parseIds(raw: string): ReadonlySet<string> {
try {
const value = JSON.parse(raw) as unknown;
return new Set(Array.isArray(value) ? value.filter((v): v is string => typeof v === "string") : []);
Expand All @@ -29,15 +40,20 @@ export function savePushDevice(
serverUserId: string,
installId: string,
capability: string,
muted: readonly string[] = [],
settings: PushSettings = {},
now = new Date(),
): void {
const db = getSqliteDb();
const at = toIso(now);
db.prepare(
`INSERT INTO push_devices (server_user_id, install_id, capability, muted, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT(server_user_id, install_id) DO UPDATE SET capability = excluded.capability, muted = excluded.muted, updated_at = excluded.updated_at`,
).run(serverUserId, installId, capability, JSON.stringify(muted), at, at);
`INSERT INTO push_devices (server_user_id, install_id, capability, muted, loud, everyone, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(server_user_id, install_id) DO UPDATE SET capability = excluded.capability, muted = excluded.muted,
loud = excluded.loud, everyone = excluded.everyone, updated_at = excluded.updated_at`,
).run(
serverUserId, installId, capability,
JSON.stringify(settings.muted ?? []), JSON.stringify(settings.loud ?? []), settings.everyone ? 1 : 0,
at, at,
);
db.prepare(
`DELETE FROM push_devices WHERE server_user_id = ? AND install_id NOT IN (
SELECT install_id FROM push_devices WHERE server_user_id = ? ORDER BY updated_at DESC, install_id LIMIT ?)`,
Expand All @@ -50,9 +66,23 @@ export function listPushDevices(serverUserId: string, now = new Date()): PushDev
const cutoff = toIso(new Date(now.getTime() - PUSH_DEVICE_STALE_DAYS * DAY_MS));
db.prepare(`DELETE FROM push_devices WHERE server_user_id = ? AND updated_at < ?`).run(serverUserId, cutoff);
const rows = db
.prepare(`SELECT install_id, capability, muted FROM push_devices WHERE server_user_id = ?`)
.all(serverUserId) as { install_id: string; capability: string; muted: string }[];
return rows.map((r) => ({ installId: r.install_id, capability: r.capability, muted: parseMuted(r.muted) }));
.prepare(`SELECT install_id, capability, muted, loud, everyone FROM push_devices WHERE server_user_id = ?`)
.all(serverUserId) as { install_id: string; capability: string; muted: string; loud: string; everyone: number }[];
return rows.map((r) => ({
installId: r.install_id,
capability: r.capability,
muted: parseIds(r.muted),
loud: parseIds(r.loud),
everyone: r.everyone === 1,
}));
}

/** Accounts with a phone that wants every message in this conversation. Stale ones are dropped later, per account. */
export function loudPushAccounts(conversationId: string): string[] {
const rows = getSqliteDb()
.prepare(`SELECT DISTINCT p.server_user_id AS id FROM push_devices p, json_each(p.loud) j WHERE j.value = ?`)
.all(conversationId) as { id: string }[];
return rows.map((r) => r.id);
}

export function removePushDevice(serverUserId: string, installId: string): void {
Expand Down
10 changes: 5 additions & 5 deletions src/services/push.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,14 +57,14 @@ describe("who is at a screen", () => {

describe("pushing", () => {
it("sends the capability in the header and only the kind in the body", async () => {
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>() }] });
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>(), loud: new Set<string>(), everyone: false }] });
h.pusher.notify({}, ["u1"], "dm", "conv");
await flush();
assert.deepEqual(h.calls, [{ url: "https://push.test/v1/push", auth: `Bearer ${CAP_A}`, body: JSON.stringify({ kind: "dm" }) }]);
});

it("skips somebody who is at a screen, and wakes them once they put the phone away", async () => {
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>() }] });
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>(), loud: new Set<string>(), everyone: false }] });
h.pusher.notify({ s: client("u1") }, ["u1"], "mention", "conv");
await flush();
assert.equal(h.calls.length, 0);
Expand All @@ -74,7 +74,7 @@ describe("pushing", () => {
});

it("buzzes once per phone per conversation in fifteen seconds", async () => {
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>() }, { installId: "tablet-12", capability: CAP_B, muted: new Set<string>() }] });
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>(), loud: new Set<string>(), everyone: false }, { installId: "tablet-12", capability: CAP_B, muted: new Set<string>(), loud: new Set<string>(), everyone: false }] });
h.pusher.notify({}, ["u1"], "dm", "conv");
h.pusher.notify({}, ["u1", "u1"], "dm", "conv");
h.pusher.notify({}, ["u1"], "dm", "other");
Expand All @@ -86,7 +86,7 @@ describe("pushing", () => {

it("forgets a capability the relay calls gone or unknown", async () => {
for (const status of [404, 410]) {
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>() }] }, status);
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>(), loud: new Set<string>(), everyone: false }] }, status);
h.pusher.notify({}, ["u1"], "dm", "conv");
await flush();
await flush();
Expand All @@ -95,7 +95,7 @@ describe("pushing", () => {
});

it("keeps it when the relay is only busy", async () => {
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>() }] }, 502);
const h = harness({ u1: [{ installId: "phone-1234", capability: CAP_A, muted: new Set<string>(), loud: new Set<string>(), everyone: false }] }, 502);
h.pusher.notify({}, ["u1"], "dm", "conv");
await flush();
await flush();
Expand Down
41 changes: 34 additions & 7 deletions src/services/push.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,10 @@ import type { Clients } from "../types";
* a kind, nothing else: it writes the notification text itself.
*/

export type PushKind = "mention" | "dm";
export type PushKind = "mention" | "dm" | "message";

/** Which of an account's phones want this push, by their own settings. Muted is checked on top. */
export type PushAccept = (device: PushDevice) => boolean;

export const DEFAULT_PUSH_RELAY = "https://push.gryt.chat";

Expand Down Expand Up @@ -53,6 +56,7 @@ export interface PusherDeps {

export function createPusher(deps: PusherDeps) {
const lastSent = new Map<string, number>();
const inFlight = new Set<Promise<void>>();

async function send(device: PushDevice, kind: PushKind): Promise<void> {
try {
Expand All @@ -71,7 +75,13 @@ export function createPusher(deps: PusherDeps) {
}

/** Fire and forget: called after delivery, and never awaited by a send. */
function notify(clientsInfo: Clients, serverUserIds: Iterable<string>, kind: PushKind, conversationId: string): void {
function notify(
clientsInfo: Clients,
serverUserIds: Iterable<string>,
kind: PushKind,
conversationId: string,
accept: PushAccept = () => true,
): void {
if (!deps.relay) return;
const now = deps.now();
for (const serverUserId of new Set(serverUserIds)) {
Expand All @@ -84,19 +94,25 @@ export function createPusher(deps: PusherDeps) {
continue;
}
for (const device of devices) {
if (device.muted.has(conversationId)) continue;
if (device.muted.has(conversationId) || !accept(device)) continue;
const key = `${device.capability}:${conversationId}`;
if (now - (lastSent.get(key) ?? 0) < QUIET_MS) continue;
lastSent.set(key, now);
void send(device, kind);
const sending = send(device, kind).finally(() => inFlight.delete(sending));
inFlight.add(sending);
}
}
if (lastSent.size > 5_000) {
for (const [key, at] of lastSent) if (now - at >= QUIET_MS) lastSent.delete(key);
}
}

return { notify };
/** Tests only: resolves once every push sent so far has had its answer. */
async function settled(): Promise<void> {
while (inFlight.size > 0) await Promise.all([...inFlight]);
}

return { notify, settled };
}

let shared: ReturnType<typeof createPusher> | null = null;
Expand All @@ -108,15 +124,26 @@ function configuredRelay(): string | null {
return relay;
}

export function pushNotify(clientsInfo: Clients, serverUserIds: Iterable<string>, kind: PushKind, conversationId: string): void {
export function pushNotify(
clientsInfo: Clients,
serverUserIds: Iterable<string>,
kind: PushKind,
conversationId: string,
accept?: PushAccept,
): void {
shared ??= createPusher({
relay: configuredRelay(),
listDevices: listPushDevices,
forget: removePushCapability,
fetch,
now: Date.now,
});
shared.notify(clientsInfo, serverUserIds, kind, conversationId);
shared.notify(clientsInfo, serverUserIds, kind, conversationId, accept);
}

/** Tests only: every push so far has reached the relay and been answered. */
export function pushesSettled(): Promise<void> {
return shared?.settled() ?? Promise.resolve();
}

/** Tests only: forget the throttle and read the relay setting again. */
Expand Down
30 changes: 26 additions & 4 deletions src/socket/handlers/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ import {
setThreadStatus,
setThreadTags,
type ThreadRecord,
loudPushAccounts,
} from "../../db";
import { isQuarantined } from "../../services/quarantineUpload";
import { processProfanity, type CensorStyle, type ProfanityMode } from "../../utils/profanityFilter";
Expand Down Expand Up @@ -822,23 +823,44 @@ export function registerChatHandlers(ctx: HandlerContext): EventHandlerMap {
});
}

/* Direct and role mentions wake a phone. @everyone and @here would wake
the whole server, and a DM already pushed as a DM. */
/* Direct and role mentions wake a phone. @everyone and @here only wake one that
doesn't suppress them, like the desktop (GRYT-1696). A DM already pushed as a DM. */
if (access.kind !== "dm") {
const blockers = await blockersOfSender(auth.tokenPayload.serverUserId);
const woken: string[] = [];
const crowd: string[] = [];
for (const [id, kind] of kindOf) {
if ((kind !== "user" && kind !== "role") || blockers.has(id)) continue;
if (blockers.has(id)) continue;
if (kind === "user" && !(await mayViewChannel(created.conversation_id, id))) continue;
woken.push(id);
(kind === "user" || kind === "role" ? woken : crowd).push(id);
}
pushNotify(clientsInfo, woken, "mention", created.conversation_id);
pushNotify(clientsInfo, crowd, "mention", created.conversation_id, (device) => device.everyone);
}
} catch (err) {
consola.warn("recording mentions failed", created.message_id, err);
}
}

/* A phone at "All messages" here wakes for every message (GRYT-1696). After the mention
push, so a mention takes the conversation's quiet window and reads as one. */
if (access.kind !== "dm") {
try {
const sender = auth.tokenPayload.serverUserId;
const blockers = await blockersOfSender(sender);
const loud: string[] = [];
for (const id of loudPushAccounts(created.conversation_id)) {
if (id === sender || blockers.has(id)) continue;
if (!(await mayViewChannel(created.conversation_id, id))) continue;
loud.push(id);
}
const conversationId = created.conversation_id;
pushNotify(clientsInfo, loud, "message", conversationId, (device) => device.loud.has(conversationId));
} catch (err) {
consola.warn("message push failed", created.message_id, err);
}
}

// After the message is out, so a promotion can neither slow one down
// nor stop one. One roles read on a server using none of this.
const promoted = await applyAutoRoles(
Expand Down
25 changes: 17 additions & 8 deletions src/socket/handlers/push.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,12 @@ type Reply = { ok: true } | { ok: false; error: string };
type Ack = (reply: Reply) => void;

const INSTALL_ID = /^[A-Za-z0-9_-]{8,64}$/;
const MAX_MUTED = 1000;
const MAX_IDS = 2000;

/** Conversation ids the phone muted. Null when it is not a short list of short strings. */
function mutedFrom(raw: unknown): string[] | null {
/** Conversation ids from the phone, muted or at All. Null when it is not a short list of short strings. */
function idsFrom(raw: unknown): string[] | null {
if (raw === undefined) return [];
if (!Array.isArray(raw) || raw.length > MAX_MUTED) return null;
if (!Array.isArray(raw) || raw.length > MAX_IDS) return null;
if (!raw.every((id) => typeof id === "string" && id.length > 0 && id.length <= 128)) return null;
return [...new Set(raw as string[])];
}
Expand All @@ -35,7 +35,7 @@ export function registerPushHandlers(ctx: HandlerContext): EventHandlerMap {

return {
"push:register": async (
payload: { accessToken: string; installId?: unknown; capability?: unknown; muted?: unknown },
payload: { accessToken: string; installId?: unknown; capability?: unknown; muted?: unknown; all?: unknown; everyone?: unknown },
ack: Ack,
) => {
ack = typeof ack === "function" ? ack : () => {};
Expand All @@ -46,11 +46,20 @@ export function registerPushHandlers(ctx: HandlerContext): EventHandlerMap {
|| typeof payload.capability !== "string" || !CAPABILITY_SHAPE.test(payload.capability)) {
return ack({ ok: false, error: "invalid_payload" });
}
const muted = mutedFrom(payload.muted);
if (!muted) return ack({ ok: false, error: "invalid_payload" });
const muted = idsFrom(payload.muted);
const loud = idsFrom(payload.all);
if (!muted || !loud || (payload.everyone !== undefined && typeof payload.everyone !== "boolean")) {
return ack({ ok: false, error: "invalid_payload" });
}
const auth = await requireAuth(socket, payload);
if (!auth) return ack({ ok: false, error: "unauthorized" });
savePushDevice(auth.tokenPayload.serverUserId, payload.installId, payload.capability, muted);
// Muted wins over All, so a list that names one conversation twice stays quiet.
const quiet = new Set(muted);
savePushDevice(auth.tokenPayload.serverUserId, payload.installId, payload.capability, {
muted,
loud: loud.filter((id) => !quiet.has(id)),
everyone: payload.everyone === true,
});
ack({ ok: true });
} catch (err) {
consola.error("push:register failed", err);
Expand Down
Loading
Loading