fix(stores): skip the notify read-back; read OIDC link conflicts from the primary

NotificationStore gains `insert`, which writes and returns the uid;
`create` is now `insert` plus the primary read-back. NotificationService
.notify ignores the row, so it uses `insert` and no longer pays a primary
round trip per recipient. NotificationDriver.create still reads the row
back from the primary.

OIDCStore.link reads the conflicting row from the primary after a unique
violation, so a replica that has not seen another account's link can no
longer make the call report success.
This commit is contained in:
Daniel Salazar committed 2026-10-10 10:25:09 -04:00
1 parent c5e37aba73
commit 005b3255cf
5 files changed
+119 -23

No files matched your search

@@ -160,12 +160,34 @@ describe('NotificationService.notify', () => {
server.clients.event.off?.('outer.gui.notif.message', handler);
});
it('does not read the row back after writing it', async () => {
const user = await makeUser();
const db = server.clients.db;
const read = vi.spyOn(db, 'read');
const pread = vi.spyOn(db, 'pread');
try {
const uid = await notifications.notify(
[user.id],
{ title: 'no read-back' },
{ type: 'share.received' },
);
expect(uid).toEqual(expect.any(String));
const readsOfRow = [...read.mock.calls, ...pread.mock.calls].filter(
([sql]) => /FROM `notification`/u.test(String(sql)),
);
expect(readsOfRow).toEqual([]);
} finally {
read.mockRestore();
pread.mockRestore();
}
});
it('pushes nothing for a recipient whose insert failed', async () => {
const good = await makeUser();
const pushed = collect('outer.gui.notif.message');
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {});
// userId 0 is falsy — NotificationStore.create rejects it. The
// userId 0 is falsy — the store's insert rejects it. The
// remaining user must still be written and still be pushed to.
await notifications.notify(
[0, good.id],
@@ -149,7 +149,7 @@ export class NotificationService extends PuterService {
userIds.map(async (userId) => {
const uid = uuidv4();
try {
await this.stores.notification.create({
await this.stores.notification.insert({
userId,
value: payload,
uid,
@@ -20,6 +20,17 @@
import { v4 as uuidv4 } from 'uuid';
import { PuterStore } from '../types';
/**
* @typedef {{
* userId: number;
* value: unknown;
* uid?: string;
* type?: string;
* audience?: string;
* appUid?: string | null;
* }} NotificationInsert
*/
export class NotificationStore extends PuterStore {
// -- Reads --------------------------------------------------------
@@ -156,24 +167,30 @@ export class NotificationStore extends PuterStore {
// -- Writes -------------------------------------------------------
/**
* Write a notification and return the stored row, read from the primary.
*
* @param {NotificationInsert} args
*/
async create(args) {
const uid = await this.insert(args);
return this.getByUid(uid, { userId: args.userId, primary: true });
}
/**
* `create` without the read-back; returns the uid only.
*
* `uid` lets the caller name the row up front — NotificationService pushes
* the uid to the socket before this insert lands, and the ack / mark-shown
* the uid to the socket after this insert lands, and the ack / mark-shown
* round trip comes back keyed on it.
*
* `type`, `audience` and `appUid` are the scope tuple; the registry in the
* notification service is what decides which combinations are legal, and
* the empty `type` written by a caller that names none reads as legacy.
*
* @param {{
* userId: number;
* value: unknown;
* uid?: string;
* type?: string;
* audience?: string;
* appUid?: string | null;
* }} args
* @param {NotificationInsert} args
* @returns {Promise<string>}
*/
async create({
async insert({
userId,
value,
uid = uuidv4(),
@@ -181,14 +198,14 @@ export class NotificationStore extends PuterStore {
audience = 'account',
appUid = null,
}) {
if (!userId) throw new Error('create: userId is required');
if (!userId) throw new Error('insert: userId is required');
const serialized =
typeof value === 'string' ? value : JSON.stringify(value ?? {});
await this.clients.db.write(
'INSERT INTO `notification` (`uid`, `user_id`, `value`, `type`, `audience`, `app_uid`) VALUES (?, ?, ?, ?, ?, ?)',
[uid, userId, serialized, type, audience, appUid],
);
return this.getByUid(uid, { userId, primary: true });
return uid;
}
/**
+12 -6
View File
@@ -38,14 +38,17 @@ const isUniqueConstraintError = (e) => {
export class OIDCStore extends PuterStore {
// -- Reads --------------------------------------------------------
async getByProviderSub(provider, providerSub) {
/** Pass `primary` when the row may be newer than the replica. */
async getByProviderSub(provider, providerSub, { primary = false } = {}) {
// Ordered so a sub that predates the UNIQUE index (two callbacks for the
// same new identity used to be able to both insert) always resolves to
// the same link, instead of bouncing the user between two accounts.
const rows = await this.clients.db.read(
'SELECT * FROM `user_oidc_providers` WHERE `provider` = ? AND `provider_sub` = ? ORDER BY `id` ASC LIMIT 1',
[provider, providerSub],
);
const sql =
'SELECT * FROM `user_oidc_providers` WHERE `provider` = ? AND `provider_sub` = ? ORDER BY `id` ASC LIMIT 1';
const params = [provider, providerSub];
const rows = primary
? await this.clients.db.pread(sql, params)
: await this.clients.db.read(sql, params);
return rows[0] ?? null;
}
@@ -73,7 +76,10 @@ export class OIDCStore extends PuterStore {
// the same (user, provider, sub) triple (idempotent no-op), or the
// sub already belongs to a DIFFERENT user. The latter must fail loudly
// so callers don't assume success and act on an unrelated account.
const existing = await this.getByProviderSub(provider, providerSub);
// The primary: it rejected the insert, so it has the row.
const existing = await this.getByProviderSub(provider, providerSub, {
primary: true,
});
if (!existing) return;
if (existing.user_id !== userId) {
// The caller is told there is a conflict, not whose.
+54 -3
View File
@@ -1,4 +1,7 @@
import { describe, expect, it, vi } from 'vitest';
import { v4 as uuidv4 } from 'uuid';
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest';
import { PuterServer } from '../../server.js';
import { setupTestServer } from '../../testUtil.js';
import { OIDCStore } from './OIDCStore.js';
type OidcRow = {
@@ -26,7 +29,7 @@ const createStore = (rows: readonly OidcRow[]) => {
throw createPostgresUniqueError();
},
),
read: vi.fn(
pread: vi.fn(
async (
_sql: string,
_params: readonly unknown[],
@@ -52,7 +55,7 @@ describe('OIDCStore', () => {
store.link(123, 'test-provider', 'subject-1'),
).resolves.toBeUndefined();
expect(db.read).toHaveBeenCalledWith(
expect(db.pread).toHaveBeenCalledWith(
expect.stringContaining('user_oidc_providers'),
['test-provider', 'subject-1'],
);
@@ -82,3 +85,51 @@ describe('OIDCStore', () => {
expect(error.message).not.toContain('subject-1');
});
});
describe('OIDCStore.link behind a lagging replica', () => {
let server: PuterServer;
beforeAll(async () => {
server = await setupTestServer();
});
afterAll(async () => {
await server?.shutdown();
});
const makeUser = async () => {
const username = `oidc-${Math.random().toString(36).slice(2, 10)}`;
return server.stores.user.create({
username,
uuid: uuidv4(),
password: null,
email: `${username}@test.local`,
});
};
it('refuses an identity linked to another account the replica has not seen yet', async () => {
const owner = await makeUser();
const other = await makeUser();
const sub = `sub-${uuidv4()}`;
await server.stores.oidc.link(owner.id, 'test-provider', sub);
// sqlite's pread delegates to read, so the primary is pinned to the
// real one while replica reads of the link table come back empty.
const db = server.clients.db;
const realRead = db.read.bind(db);
const pread = vi.spyOn(db, 'pread').mockImplementation(realRead);
const read = vi
.spyOn(db, 'read')
.mockImplementation(async (q, p) =>
/FROM `user_oidc_providers`/u.test(q) ? [] : realRead(q, p),
);
try {
await expect(
server.stores.oidc.link(other.id, 'test-provider', sub),
).rejects.toMatchObject({ statusCode: 409 });
} finally {
read.mockRestore();
pread.mockRestore();
}
});
});