From f3acc2ce397730ffbba5f60d30a036212adc2b41 Mon Sep 17 00:00:00 2001 From: 4gray Date: Mon, 27 Jul 2026 09:52:24 +0200 Subject: [PATCH] feat(stalker): serialize base identity mutations --- .../stalker-base-identity-coordinator.spec.ts | 259 ++++++++++++++++++ .../stalker-base-identity-coordinator.ts | 254 +++++++++++++++++ 2 files changed, 513 insertions(+) create mode 100644 apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.spec.ts create mode 100644 apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.ts diff --git a/apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.spec.ts b/apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.spec.ts new file mode 100644 index 000000000..9938ea391 --- /dev/null +++ b/apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.spec.ts @@ -0,0 +1,259 @@ +import { + StalkerBaseIdentityCoordinator, + StalkerBaseIdentityCoordinatorPool, +} from './stalker-base-identity-coordinator'; + +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, reject, resolve }; +} + +describe('StalkerBaseIdentityCoordinatorPool', () => { + it('keys coordinators by approved origin and normalized MAC', () => { + const pool = new StalkerBaseIdentityCoordinatorPool(); + + const first = pool.get( + 'https://portal.example/customer/c/', + '00-1a-79-00-00-01' + ); + const sameBaseIdentity = pool.get( + 'https://portal.example/other/path', + '00:1A:79:00:00:01' + ); + const otherOrigin = pool.get( + 'https://other.example/customer/c/', + '00:1A:79:00:00:01' + ); + const otherMac = pool.get( + 'https://portal.example/customer/c/', + '00:1A:79:00:00:02' + ); + + expect(sameBaseIdentity).toBe(first); + expect(otherOrigin).not.toBe(first); + expect(otherMac).not.toBe(first); + expect(pool.size).toBe(3); + }); +}); + +describe('StalkerBaseIdentityCoordinator', () => { + it('allows concurrent readers for the active principal and epoch', async () => { + const coordinator = new StalkerBaseIdentityCoordinator(); + const ready = await coordinator.runMutation('principal-a', async () => { + return 'authenticated'; + }); + const firstGate = deferred(); + const secondGate = deferred(); + const started: string[] = []; + + const first = coordinator.runRead( + 'principal-a', + ready.snapshot.epoch, + async () => { + started.push('first'); + await firstGate.promise; + return 1; + } + ); + const second = coordinator.runRead( + 'principal-a', + ready.snapshot.epoch, + async () => { + started.push('second'); + await secondGate.promise; + return 2; + } + ); + await Promise.resolve(); + + expect(started).toEqual(['first', 'second']); + firstGate.resolve(); + secondGate.resolve(); + await expect(Promise.all([first, second])).resolves.toEqual([ + { + kind: 'completed', + snapshot: { activePrincipal: 'principal-a', epoch: 1 }, + value: 1, + }, + { + kind: 'completed', + snapshot: { activePrincipal: 'principal-a', epoch: 1 }, + value: 2, + }, + ]); + }); + + it('runs mutations exclusively and does not admit later readers ahead of a writer', async () => { + const coordinator = new StalkerBaseIdentityCoordinator(); + const ready = await coordinator.runMutation('principal-a', async () => { + return undefined; + }); + const readerGate = deferred(); + const mutationGate = deferred(); + const order: string[] = []; + + const reader = coordinator.runRead( + 'principal-a', + ready.snapshot.epoch, + async () => { + order.push('reader-start'); + await readerGate.promise; + order.push('reader-end'); + } + ); + await Promise.resolve(); + const mutation = coordinator.runMutation( + 'principal-a', + async () => { + order.push('mutation-start'); + await mutationGate.promise; + order.push('mutation-end'); + } + ); + const laterReader = coordinator.runRead( + 'principal-a', + ready.snapshot.epoch + 1, + async () => { + order.push('later-reader'); + } + ); + + await Promise.resolve(); + expect(order).toEqual(['reader-start']); + readerGate.resolve(); + await reader; + await Promise.resolve(); + expect(order).toEqual([ + 'reader-start', + 'reader-end', + 'mutation-start', + ]); + mutationGate.resolve(); + await mutation; + await laterReader; + expect(order).toEqual([ + 'reader-start', + 'reader-end', + 'mutation-start', + 'mutation-end', + 'later-reader', + ]); + }); + + it('suspends stale epochs and inactive principals without running their operation', async () => { + const coordinator = new StalkerBaseIdentityCoordinator(); + const first = await coordinator.runMutation( + 'principal-a', + async () => undefined + ); + const second = await coordinator.runMutation( + 'principal-b', + async () => undefined + ); + const operation = jest.fn(); + + await expect( + coordinator.runRead( + 'principal-a', + first.snapshot.epoch, + operation + ) + ).resolves.toEqual({ + kind: 'suspended', + snapshot: { activePrincipal: 'principal-b', epoch: 2 }, + }); + await expect( + coordinator.runRead( + 'principal-b', + first.snapshot.epoch, + operation + ) + ).resolves.toEqual({ + kind: 'suspended', + snapshot: { activePrincipal: 'principal-b', epoch: 2 }, + }); + expect(second.snapshot.epoch).toBe(2); + expect(operation).not.toHaveBeenCalled(); + }); + + it('keeps alternating principals isolated behind monotonic mutation epochs', async () => { + const coordinator = new StalkerBaseIdentityCoordinator(); + const events: string[] = []; + + const first = await coordinator.runMutation( + 'principal-a', + async (epoch) => { + events.push(`auth-a-${epoch}`); + } + ); + const second = await coordinator.runMutation( + 'principal-b', + async (epoch) => { + events.push(`auth-b-${epoch}`); + } + ); + const third = await coordinator.runMutation( + 'principal-a', + async (epoch) => { + events.push(`reauth-a-${epoch}`); + } + ); + + expect(events).toEqual(['auth-a-1', 'auth-b-2', 'reauth-a-3']); + expect(first.snapshot.epoch).toBe(1); + expect(second.snapshot.epoch).toBe(2); + expect(third.snapshot).toEqual({ + activePrincipal: 'principal-a', + epoch: 3, + }); + }); + + it('advances the epoch and suspends all principals after a failed mutation', async () => { + const coordinator = new StalkerBaseIdentityCoordinator(); + const ready = await coordinator.runMutation( + 'principal-a', + async () => undefined + ); + + await expect( + coordinator.runMutation('principal-b', async () => { + throw new Error('authentication-failed'); + }) + ).rejects.toThrow('authentication-failed'); + + expect(coordinator.snapshot).toEqual({ epoch: 2 }); + await expect( + coordinator.runRead( + 'principal-a', + ready.snapshot.epoch, + jest.fn() + ) + ).resolves.toEqual({ + kind: 'suspended', + snapshot: { epoch: 2 }, + }); + }); + + it('can explicitly enter a no-active-principal transition epoch', async () => { + const coordinator = new StalkerBaseIdentityCoordinator(); + await coordinator.runMutation( + 'principal-a', + async () => undefined + ); + + const transition = await coordinator.runMutation( + undefined, + async () => 'transition' + ); + + expect(transition).toEqual({ + snapshot: { epoch: 2 }, + value: 'transition', + }); + }); +}); diff --git a/apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.ts b/apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.ts new file mode 100644 index 000000000..f85a64be7 --- /dev/null +++ b/apps/electron-backend/src/app/services/stalker-session/stalker-base-identity-coordinator.ts @@ -0,0 +1,254 @@ +import { normalizeStalkerMacAddress } from '@iptvnator/portal/stalker/protocol'; + +export interface StalkerBaseIdentitySnapshot { + readonly epoch: number; + readonly activePrincipal?: string; +} + +export type StalkerCoordinatedReadResult = + | { + readonly kind: 'completed'; + readonly snapshot: StalkerBaseIdentitySnapshot; + readonly value: T; + } + | { + readonly kind: 'suspended'; + readonly snapshot: StalkerBaseIdentitySnapshot; + }; + +export interface StalkerCoordinatedMutationResult { + readonly snapshot: StalkerBaseIdentitySnapshot; + readonly value: T; +} + +interface ReadWaiter { + readonly kind: 'read'; + readonly resolve: () => void; +} + +interface WriteWaiter { + readonly kind: 'write'; + readonly resolve: () => void; +} + +type LockWaiter = ReadWaiter | WriteWaiter; + +/** + * A fair asynchronous read/write gate for one approved-origin + MAC pair. + * + * Catalog operations may share a read epoch. Every wire mutation advances the + * epoch, runs alone, and either installs exactly one active principal or leaves + * the base identity suspended when it fails. + */ +export class StalkerBaseIdentityCoordinator { + #epoch = 0; + #activePrincipal: string | undefined; + #activeReaders = 0; + #writerActive = false; + readonly #waiters: LockWaiter[] = []; + + get snapshot(): StalkerBaseIdentitySnapshot { + return createSnapshot(this.#epoch, this.#activePrincipal); + } + + async runRead( + expectedPrincipal: string, + expectedEpoch: number, + operation: () => Promise + ): Promise> { + validatePrincipal(expectedPrincipal); + validateEpoch(expectedEpoch); + await this.#acquireRead(); + try { + if ( + this.#epoch !== expectedEpoch || + this.#activePrincipal !== expectedPrincipal + ) { + return { + kind: 'suspended', + snapshot: this.snapshot, + }; + } + + const value = await operation(); + return { + kind: 'completed', + snapshot: this.snapshot, + value, + }; + } finally { + this.#releaseRead(); + } + } + + async runMutation( + nextPrincipal: string | undefined, + operation: (mutationEpoch: number) => Promise + ): Promise> { + if (nextPrincipal !== undefined) { + validatePrincipal(nextPrincipal); + } + await this.#acquireWrite(); + try { + if (this.#epoch >= Number.MAX_SAFE_INTEGER) { + throw new Error('stalker-coordinator-epoch-exhausted'); + } + this.#epoch += 1; + this.#activePrincipal = undefined; + + const value = await operation(this.#epoch); + this.#activePrincipal = nextPrincipal; + return { + snapshot: this.snapshot, + value, + }; + } finally { + this.#releaseWrite(); + } + } + + async #acquireRead(): Promise { + if (!this.#writerActive && this.#waiters.length === 0) { + this.#activeReaders += 1; + return; + } + await new Promise((resolve) => { + this.#waiters.push({ kind: 'read', resolve }); + }); + } + + #releaseRead(): void { + this.#activeReaders -= 1; + if (this.#activeReaders < 0) { + this.#activeReaders = 0; + throw new Error('stalker-coordinator-read-underflow'); + } + this.#drainWaiters(); + } + + async #acquireWrite(): Promise { + if ( + !this.#writerActive && + this.#activeReaders === 0 && + this.#waiters.length === 0 + ) { + this.#writerActive = true; + return; + } + await new Promise((resolve) => { + this.#waiters.push({ kind: 'write', resolve }); + }); + } + + #releaseWrite(): void { + if (!this.#writerActive) { + throw new Error('stalker-coordinator-write-underflow'); + } + this.#writerActive = false; + this.#drainWaiters(); + } + + #drainWaiters(): void { + if (this.#writerActive || this.#activeReaders > 0) { + return; + } + + const first = this.#waiters[0]; + if (!first) { + return; + } + if (first.kind === 'write') { + this.#waiters.shift(); + this.#writerActive = true; + first.resolve(); + return; + } + + while (this.#waiters[0]?.kind === 'read') { + const reader = this.#waiters.shift() as ReadWaiter; + this.#activeReaders += 1; + reader.resolve(); + } + } +} + +/** + * Process-local registry. The key deliberately excludes endpoint path, + * identity revision, playlist, and principal because portals may rotate a + * single server-side token for the whole origin/MAC pair. + */ +export class StalkerBaseIdentityCoordinatorPool { + readonly #coordinators = new Map< + string, + StalkerBaseIdentityCoordinator + >(); + + get size(): number { + return this.#coordinators.size; + } + + get( + approvedPortalUrl: string, + macAddress: string + ): StalkerBaseIdentityCoordinator { + const key = createBaseIdentityKey(approvedPortalUrl, macAddress); + const existing = this.#coordinators.get(key); + if (existing) { + return existing; + } + const coordinator = new StalkerBaseIdentityCoordinator(); + this.#coordinators.set(key, coordinator); + return coordinator; + } + + clear(): void { + this.#coordinators.clear(); + } +} + +function createBaseIdentityKey( + approvedPortalUrl: string, + macAddress: string +): string { + let url: URL; + try { + url = new URL(approvedPortalUrl); + } catch { + throw new Error('invalid-url'); + } + if ( + (url.protocol !== 'http:' && url.protocol !== 'https:') || + url.username !== '' || + url.password !== '' + ) { + throw new Error('invalid-url'); + } + const normalizedMac = normalizeStalkerMacAddress(macAddress); + return JSON.stringify([url.origin, normalizedMac]); +} + +function validatePrincipal(principal: string): void { + if ( + principal.length === 0 || + principal.length > 256 || + /[\u0000-\u001f\u007f]/.test(principal) + ) { + throw new Error('invalid-stalker-principal'); + } +} + +function validateEpoch(epoch: number): void { + if (!Number.isSafeInteger(epoch) || epoch < 0) { + throw new Error('invalid-stalker-coordinator-epoch'); + } +} + +function createSnapshot( + epoch: number, + activePrincipal: string | undefined +): StalkerBaseIdentitySnapshot { + return Object.freeze({ + ...(activePrincipal === undefined ? {} : { activePrincipal }), + epoch, + }); +}