/** * wj-pi-auto-compact 与其他有状态扩展之间的通用协调协议。 * * 协议故意只描述压缩屏障,不携带任何宿主或业务插件的领域身份。 */ export const COORDINATION_PROTOCOL_VERSION = "wj-pi-auto-compact/coordination/v1" as const; export const COORDINATION_CHANNELS = Object.freeze({ discover: "wj-pi-auto-compact/coordination/v1/discover", discovered: "wj-pi-auto-compact/coordination/v1/discovered", prepare: "wj-pi-auto-compact/coordination/v1/prepare", prepared: "wj-pi-auto-compact/coordination/v1/prepared", complete: "wj-pi-auto-compact/coordination/v1/complete", completed: "wj-pi-auto-compact/coordination/v1/completed", } as const); export type CompactionOutcome = "succeeded" | "failed" | "cancelled" | "not_started"; export interface CoordinationDiscoverRequest { readonly protocolVersion: typeof COORDINATION_PROTOCOL_VERSION; readonly requestId: string; } /** 发现响应的字段是协议的完整公开面;participantId 必须保持不透明。 */ export interface CoordinationDiscovered { readonly protocolVersion: typeof COORDINATION_PROTOCOL_VERSION; readonly requestId: string; readonly participantId: string; readonly requiresBarrier: boolean; } export interface CoordinationPrepareRequest { readonly protocolVersion: typeof COORDINATION_PROTOCOL_VERSION; readonly requestId: string; readonly participantId: string; } export interface CoordinationPrepared { readonly protocolVersion: typeof COORDINATION_PROTOCOL_VERSION; readonly requestId: string; readonly participantId: string; readonly prepared: boolean; } export interface CoordinationCompleteRequest { readonly protocolVersion: typeof COORDINATION_PROTOCOL_VERSION; readonly requestId: string; readonly participantId: string; readonly outcome: CompactionOutcome; readonly continuationExpected?: boolean; } export interface CoordinationCompleted { readonly protocolVersion: typeof COORDINATION_PROTOCOL_VERSION; readonly requestId: string; readonly participantId: string; readonly completed: boolean; } export type CoordinationEventBus = { emit(channel: string, data: unknown): void; on(channel: string, handler: (data: unknown) => void): () => void; }; export const DISCOVERY_WINDOW_MS = 100; // 必须覆盖参与者内部正向确认与独立 not_started 补偿两个有界阶段。 export const PREPARE_ACK_TIMEOUT_MS = 12_000; export const COMPLETE_ACK_TIMEOUT_MS = 12_000; const OUTCOMES: ReadonlySet = new Set([ "succeeded", "failed", "cancelled", "not_started", ]); export function isCompactionOutcome(value: unknown): value is CompactionOutcome { return typeof value === "string" && OUTCOMES.has(value); } export function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function hasOnlyKeys(value: Record, keys: readonly string[]): boolean { return Object.keys(value).every((key) => keys.includes(key)); } function nonEmptyId(value: unknown): value is string { return typeof value === "string" && value.length > 0 && value.length <= 256; } export function parseDiscoverRequest(value: unknown): CoordinationDiscoverRequest | undefined { if (!isRecord(value) || !hasOnlyKeys(value, ["protocolVersion", "requestId"])) return undefined; if (value.protocolVersion !== COORDINATION_PROTOCOL_VERSION || !nonEmptyId(value.requestId)) return undefined; return Object.freeze({ protocolVersion: COORDINATION_PROTOCOL_VERSION, requestId: value.requestId, }); } export function parseDiscovered(value: unknown): CoordinationDiscovered | undefined { if (!isRecord(value) || !hasOnlyKeys(value, ["protocolVersion", "requestId", "participantId", "requiresBarrier"])) { return undefined; } if ( value.protocolVersion !== COORDINATION_PROTOCOL_VERSION || !nonEmptyId(value.requestId) || !nonEmptyId(value.participantId) || typeof value.requiresBarrier !== "boolean" ) return undefined; return Object.freeze({ protocolVersion: COORDINATION_PROTOCOL_VERSION, requestId: value.requestId, participantId: value.participantId, requiresBarrier: value.requiresBarrier, }); } export function parsePrepareRequest(value: unknown): CoordinationPrepareRequest | undefined { if (!isRecord(value) || !hasOnlyKeys(value, ["protocolVersion", "requestId", "participantId"])) return undefined; if ( value.protocolVersion !== COORDINATION_PROTOCOL_VERSION || !nonEmptyId(value.requestId) || !nonEmptyId(value.participantId) ) return undefined; return Object.freeze({ protocolVersion: COORDINATION_PROTOCOL_VERSION, requestId: value.requestId, participantId: value.participantId, }); } export function parsePrepared(value: unknown): CoordinationPrepared | undefined { if (!isRecord(value) || !hasOnlyKeys(value, ["protocolVersion", "requestId", "participantId", "prepared"])) { return undefined; } if ( value.protocolVersion !== COORDINATION_PROTOCOL_VERSION || !nonEmptyId(value.requestId) || !nonEmptyId(value.participantId) || typeof value.prepared !== "boolean" ) return undefined; return Object.freeze({ protocolVersion: COORDINATION_PROTOCOL_VERSION, requestId: value.requestId, participantId: value.participantId, prepared: value.prepared, }); } export function parseCompleteRequest(value: unknown): CoordinationCompleteRequest | undefined { if (!isRecord(value) || !hasOnlyKeys(value, [ "protocolVersion", "requestId", "participantId", "outcome", "continuationExpected", ])) { return undefined; } const continuationExpected = Object.prototype.hasOwnProperty.call(value, "continuationExpected") ? value.continuationExpected : false; if ( value.protocolVersion !== COORDINATION_PROTOCOL_VERSION || !nonEmptyId(value.requestId) || !nonEmptyId(value.participantId) || !isCompactionOutcome(value.outcome) || typeof continuationExpected !== "boolean" || (continuationExpected && value.outcome !== "succeeded") ) return undefined; return Object.freeze({ protocolVersion: COORDINATION_PROTOCOL_VERSION, requestId: value.requestId, participantId: value.participantId, outcome: value.outcome, continuationExpected, }); } export function parseCompleted(value: unknown): CoordinationCompleted | undefined { if (!isRecord(value) || !hasOnlyKeys(value, ["protocolVersion", "requestId", "participantId", "completed"])) { return undefined; } if ( value.protocolVersion !== COORDINATION_PROTOCOL_VERSION || !nonEmptyId(value.requestId) || !nonEmptyId(value.participantId) || typeof value.completed !== "boolean" ) return undefined; return Object.freeze({ protocolVersion: COORDINATION_PROTOCOL_VERSION, requestId: value.requestId, participantId: value.participantId, completed: value.completed, }); } export function freezeEvent(value: T): Readonly { return Object.freeze({ ...value }); }