mirror of
https://github.com/4gray/iptvnator.git
synced 2026-10-10 10:06:15 -08:00
feat(stalker): serialize base identity mutations
This commit is contained in:
1 parent
2a88060fae
commit
f3acc2ce39
2 files changed
+513
No files matched your search
+259
@@ -0,0 +1,259 @@
|
||||
import {
|
||||
StalkerBaseIdentityCoordinator,
|
||||
StalkerBaseIdentityCoordinatorPool,
|
||||
} from './stalker-base-identity-coordinator';
|
||||
|
||||
function deferred<T = void>() {
|
||||
let resolve!: (value: T | PromiseLike<T>) => void;
|
||||
let reject!: (reason?: unknown) => void;
|
||||
const promise = new Promise<T>((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',
|
||||
});
|
||||
});
|
||||
});
|
||||
+254
@@ -0,0 +1,254 @@
|
||||
import { normalizeStalkerMacAddress } from '@iptvnator/portal/stalker/protocol';
|
||||
|
||||
export interface StalkerBaseIdentitySnapshot {
|
||||
readonly epoch: number;
|
||||
readonly activePrincipal?: string;
|
||||
}
|
||||
|
||||
export type StalkerCoordinatedReadResult<T> =
|
||||
| {
|
||||
readonly kind: 'completed';
|
||||
readonly snapshot: StalkerBaseIdentitySnapshot;
|
||||
readonly value: T;
|
||||
}
|
||||
| {
|
||||
readonly kind: 'suspended';
|
||||
readonly snapshot: StalkerBaseIdentitySnapshot;
|
||||
};
|
||||
|
||||
export interface StalkerCoordinatedMutationResult<T> {
|
||||
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<T>(
|
||||
expectedPrincipal: string,
|
||||
expectedEpoch: number,
|
||||
operation: () => Promise<T>
|
||||
): Promise<StalkerCoordinatedReadResult<T>> {
|
||||
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<T>(
|
||||
nextPrincipal: string | undefined,
|
||||
operation: (mutationEpoch: number) => Promise<T>
|
||||
): Promise<StalkerCoordinatedMutationResult<T>> {
|
||||
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<void> {
|
||||
if (!this.#writerActive && this.#waiters.length === 0) {
|
||||
this.#activeReaders += 1;
|
||||
return;
|
||||
}
|
||||
await new Promise<void>((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<void> {
|
||||
if (
|
||||
!this.#writerActive &&
|
||||
this.#activeReaders === 0 &&
|
||||
this.#waiters.length === 0
|
||||
) {
|
||||
this.#writerActive = true;
|
||||
return;
|
||||
}
|
||||
await new Promise<void>((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,
|
||||
});
|
||||
}
|
||||
Reference in new issue
Block a user