fix: let a revocation be seen by the reads and the nodes that outlive it

Four places where a revoke lands on the primary and the next reader can
still be told the old answer.

The u2a cache key was dropped but not re-filled, so the next miss read a
replica that may still carry the row and cached it for the TTL -- and a
revoke's own settling pass is usually that miss. It re-fills from the
primary, as the u2u side already does.

A revoked session dropped its keys the same way. It now leaves the
revoked row in the cache, read from the primary, so there is nothing for
a lagging replica to re-introduce.

A link share is its own authority rather than an index over a
permission, so a replica still carrying a withdrawn row answers with it.
For a short window after one is withdrawn, those reads go to the
primary.

Socket eviction rode a local event, so it only reached the node that did
the revoking. It is announced on the fan-out channel too, which is where
the node actually holding the connection hears it.
This commit is contained in:
Juan Castro committed 2026-10-09 16:24:01 -04:00
1 parent ff35fbd66e
commit 155341ee6c
10 files changed
+192 -10

No files matched your search

+5
View File
@@ -708,6 +708,11 @@ export type EventMap = {
* the next handshake.
*/
'auth.sessions.revoked': { user_id: number; session_uids: string[] };
/** The same revocation, carried to sibling nodes and peer clusters. */
'outer.pubsub.auth.sessions.revoked': {
user_id: number;
session_uids: string[];
};
/** One access token was revoked; only its own connections should drop. */
'auth.access-token.revoked': { token_uid: string };
+8 -2
View File
@@ -728,9 +728,15 @@ export class AuthService extends PuterService {
revoked: { userId: number; uuids: string[] } | null | undefined,
): void {
if (!revoked?.userId || revoked.uuids.length === 0) return;
const payload = {
user_id: revoked.userId,
session_uids: revoked.uuids,
};
this.clients.event?.emit('auth.sessions.revoked', payload, {});
// Whoever holds the socket may not be whoever revoked the session.
this.clients.event?.emit(
'auth.sessions.revoked',
{ user_id: revoked.userId, session_uids: revoked.uuids },
'outer.pubsub.auth.sessions.revoked',
payload,
{},
);
}
@@ -662,6 +662,28 @@ describe('SocketService (live socket.io)', () => {
});
});
it('drops the connection on a revocation that reached it from elsewhere', async () => {
const evicted = await createTestUser(server, {
username: 'sock-fanout',
password: 'sock-fanout-password',
});
const row = await server.stores.user.getByUsername(evicted.username);
const socket = await connect({ auth_token: `Bearer ${evicted.token}` });
expect(socket.connected).toBe(true);
// As the fan-out delivers it: the node holding this socket is not
// the one that revoked the session.
server.clients.event.emit(
'outer.pubsub.auth.sessions.revoked',
{ user_id: row!.id, session_uids: ['whatever'] },
{},
);
await vi.waitFor(() => expect(socket.connected).toBe(false), {
timeout: 5_000,
});
});
it('leaves other accounts alone when one is revoked', async () => {
const kept = await createTestUser(server, {
username: 'sock-kept',
+8 -4
View File
@@ -834,15 +834,19 @@ export class SocketService extends PuterService {
},
);
this.clients.event.on(
// Both channels: the local one for the node that did the revoking,
// the fanned-out one for every other node holding its sockets.
for (const key of [
'auth.sessions.revoked',
(_key: string, data: unknown) => {
'outer.pubsub.auth.sessions.revoked',
] as const) {
this.clients.event.on(key, (_key: string, data: unknown) => {
const { user_id } = data as { user_id: number };
this.#evictUserSockets(user_id).catch((err: unknown) => {
console.error('[socket] session eviction failed', err);
});
},
);
});
}
this.clients.event.on(
'auth.access-token.revoked',
@@ -152,6 +152,28 @@ describe('PermissionStore', () => {
expect(rows[0].extra).toEqual({ v: 2 });
});
it('re-fills the cache from the primary when a grant is revoked', async () => {
const user = await makeUser();
const app = await makeApp(user.id);
await store.upsertUserAppPerm(user.id, app.id, 'driver:kv', {});
// Warm it, so a revoke has something to leave behind.
await store.hasUserAppPerm(user.id, app.id, 'driver:kv');
const read = vi.spyOn(server.clients.db, 'read');
try {
await store.deleteUserAppPerm(user.id, app.id, 'driver:kv');
read.mockClear();
// Answered from the re-filled cache, not from a replica that
// may still carry the row that was just removed.
expect(
await store.hasUserAppPerm(user.id, app.id, 'driver:kv'),
).toBe(false);
expect(read).not.toHaveBeenCalled();
} finally {
read.mockRestore();
}
});
it('deletes a single grant and every grant for an app', async () => {
const user = await makeUser();
const app = await makeApp(user.id);
@@ -798,6 +798,25 @@ export class PermissionStore extends PuterStore {
return rows.map((row) => this.#decodeExtra<LinkedUserAppPermRow>(row));
}
/** Re-fill from the primary; a miss on a replica re-caches a stale row. */
async warmUserAppPermsFromPrimary(
userId: number,
appId: number,
): Promise<void> {
try {
const rows = await this.listUserAppPermsFromPrimary(userId, appId);
await this.clients.redis.set(
this.#u2aCacheKey(userId, appId),
JSON.stringify(rows),
'EX',
U2A_CACHE_TTL_SECONDS,
);
} catch (e) {
// A key left absent is the safe miss: the next read goes to the DB.
console.warn('[permission] u2a cache re-fill failed:', e);
}
}
async upsertUserAppPerm(
userId: number,
appId: number,
@@ -838,6 +857,7 @@ export class PermissionStore extends PuterStore {
keys: [this.#u2aCacheKey(userId, appId)],
broadcast: true,
});
await this.warmUserAppPermsFromPrimary(userId, appId);
}
async deleteUserAppAll(userId: number, appId: number): Promise<void> {
@@ -849,6 +869,7 @@ export class PermissionStore extends PuterStore {
keys: [this.#u2aCacheKey(userId, appId)],
broadcast: true,
});
await this.warmUserAppPermsFromPrimary(userId, appId);
}
/**
@@ -388,6 +388,27 @@ export class SessionStore extends PuterStore {
[now, uuid],
);
await this.publishCacheKeys({ keys, broadcast: true });
await this.#cacheRevokedFromPrimary(uuid);
}
/**
* Seed the cache with the revoked row, read from the primary. Dropping the
* key alone leaves the next miss to a replica that may still show the row
* active, which would then be cached for a full TTL.
*/
async #cacheRevokedFromPrimary(uuid) {
try {
const rows = await this.clients.db.pread(
'SELECT * FROM `sessions` WHERE `uuid` = ? LIMIT 1',
[uuid],
);
const normalized = this.#normalizeRow(rows[0]);
if (normalized?.revoked_at != null) {
await this.#writeCache(normalized);
}
} catch {
// No entry is the safe miss; the next read goes to the database.
}
}
/**
@@ -205,6 +205,23 @@ describe('SessionStore', () => {
expect(row.revoked_at).toBeGreaterThan(0);
});
it('leaves the revoked row cached, so no replica can re-introduce it', async () => {
const user = await makeUser();
const session = await target.create(user.id);
await target.getByUuid(session.uuid);
await target.removeByUuid(session.uuid);
const read = vi.spyOn(server.clients.db, 'read');
try {
// Answered from the cache the revoke left, not from a read.
expect(await target.getByUuid(session.uuid)).toBeNull();
expect(read).not.toHaveBeenCalled();
} finally {
read.mockRestore();
}
});
it('is idempotent — second call does not overwrite revoked_at', async () => {
const user = await makeUser();
const session = await target.create(user.id);
+40 -4
View File
@@ -1017,8 +1017,28 @@ export class ShareStore extends PuterStore {
*
* @param {number} fsentryId
*/
/** How long after a link is withdrawn these reads go to the primary. */
static REVOKED_LINK_PRIMARY_WINDOW_SECONDS = 30;
#revokedLinkKey() {
return 'share:anyone:revoked-recently';
}
/** Past the replica while a link has just been withdrawn. */
async #readAnyone(sql, params) {
let recent = false;
try {
recent = !!(await this.clients.redis.get(this.#revokedLinkKey()));
} catch {
// Unreadable marker reads as no revoke; the row is still checked.
}
return recent
? this.clients.db.pread(sql, params)
: this.clients.db.read(sql, params);
}
async getAnyone(fsentryId) {
const rows = await this.clients.db.read(
const rows = await this.#readAnyone(
'SELECT * FROM `share` WHERE `fsentry_id` = ? AND `anyone` = 1 LIMIT 1',
[fsentryId],
);
@@ -1033,7 +1053,7 @@ export class ShareStore extends PuterStore {
async listAnyoneOnFsentries(fsentryIds) {
if (fsentryIds.length === 0) return [];
const placeholders = fsentryIds.map(() => '?').join(', ');
const rows = await this.clients.db.read(
const rows = await this.#readAnyone(
`SELECT * FROM \`share\` WHERE \`fsentry_id\` IN (${placeholders}) ` +
'AND `anyone` = 1 ORDER BY `id`',
fsentryIds,
@@ -1054,7 +1074,7 @@ export class ShareStore extends PuterStore {
async listAnyoneReaching(uuids) {
if (uuids.length === 0) return [];
const placeholders = uuids.map(() => '?').join(', ');
const rows = await this.clients.db.read(
const rows = await this.#readAnyone(
'SELECT `share`.`mode`, `fsentries`.`uuid` AS `entry_uuid`, ' +
'`fsentries`.`user_id` AS `owner_user_id` FROM `share` ' +
'JOIN `fsentries` ON `fsentries`.`id` = `share`.`fsentry_id` ' +
@@ -1133,7 +1153,23 @@ export class ShareStore extends PuterStore {
'DELETE FROM `share` WHERE `fsentry_id` = ? AND `anyone` = 1',
[fsentryId],
);
return (result?.affectedRows ?? result?.changes ?? 0) > 0;
const removed = (result?.affectedRows ?? result?.changes ?? 0) > 0;
if (removed) await this.#markLinkRevoked();
return removed;
}
/** Open the window in which link reads go to the primary. */
async #markLinkRevoked() {
try {
await this.clients.redis.set(
this.#revokedLinkKey(),
'1',
'EX',
ShareStore.REVOKED_LINK_PRIMARY_WINDOW_SECONDS,
);
} catch (e) {
console.warn('[share] link-revoke marker not set:', e);
}
}
// -- Daily quota --------------------------------------------------
@@ -332,6 +332,34 @@ describe('ShareStore', () => {
expect(page.items.map((r) => r.uid)).toContain(created.uid);
});
it('reads a link past the replica once one has been withdrawn', async () => {
const entry = await makeEntry(issuer);
await store.upsertAnyone({
issuerUserId: issuer.id,
fsentryId: entry.id,
mode: 'read',
});
const read = vi.spyOn(server.clients.db, 'read');
const pread = vi.spyOn(server.clients.db, 'pread');
try {
// Nothing withdrawn: the ordinary read serves it.
await store.getAnyone(entry.id);
expect(read).toHaveBeenCalled();
expect(pread).not.toHaveBeenCalled();
read.mockClear();
pread.mockClear();
await store.deleteAnyone(entry.id);
// Withdrawn just now, so the row is read where it is gone.
expect(await store.getAnyone(entry.id)).toBeNull();
expect(pread).toHaveBeenCalled();
} finally {
read.mockRestore();
pread.mockRestore();
}
});
it('lets a link share lose its app, and never take another one', async () => {
const entry = await makeEntry(issuer);
const appOf = (row) => {