← Files ModRetro Chromatic PluginARCHIVED FILE
dist/web-preview-recording.js
19.9 KB · Oct 2, 2026 · 00:37 UTC
import { randomUUID } from "node:crypto";
import { RecordingJournal } from "./web-preview-recording-journal.js";
import { parseRecordingState, recordingDuration, recordingId, savedRecordingKey, RECORDING_DRAIN_MS } from "./web-annotations/recording-protocol.js";
const errorText = (error) => (error instanceof Error ? error.message : String(error)).slice(0, 512);
/** One async recording, with independent durable ownership and fresh close acknowledgments. */
export class PreviewRecordingBroker {
now;
cartridgeType;
mono;
views = new Map();
retired = new Map();
journal;
job;
admission = true;
pending = Promise.resolve();
timer;
closeState;
cancelledEpoch;
reserving = false;
completion;
completed = false;
journalFailed = false;
startDeadline = 0;
stopDeadline = 0;
finalizationDeadline;
constructor(projectRoot, now = () => Date.now(), cartridgeType = 0, mono = () => performance.now()) {
this.now = now;
this.cartridgeType = cartridgeType;
this.mono = mono;
this.journal = new RecordingJournal(projectRoot);
}
get busy() { return this.reserving || !!this.job && (!this.completed || this.job.phase === "unknown"); }
get remainingDrainMs() { return Math.max(0, (this.closeState?.deadline ?? this.mono() + RECORDING_DRAIN_MS) - this.mono()); }
get closeUncertain() { return !!this.closeState && !this.closeState.human; }
async unresolved() {
if (this.job)
return this.state();
const entries = await this.journal.entries();
const entry = entries.find(e => e.phase !== "finished" && e.phase !== "failed");
if (entry)
return { ...entry, phase: "unknown" };
const active = await this.journal.activeId();
return active ? { id: active, phase: "unknown", error: "Recording reservation is incomplete; preserve the original journal." } : undefined;
}
get active() { return this.busy || !!this.closeState; }
state() {
this.tick();
const result = this.job ? structuredClone(this.job) : undefined;
if (result && !this.completed && (result.phase === "finished" || result.phase === "failed"))
result.phase = "stopping";
return result;
}
async status(id) {
recordingId(id);
if (this.job?.id === id)
return this.state();
return this.read(id);
}
view(id) { const s = this.views.get(id)?.status; return s ? structuredClone(s) : undefined; }
isRetired(id) { return this.retired.has(id); }
persist() {
const job = this.job;
if (!job)
return Promise.resolve();
const snapshot = structuredClone(job);
const task = this.pending.then(async () => {
// An earlier queued write may have failed after this snapshot was made.
// Once uncertain, no later terminal snapshot is allowed to release ownership.
if (this.journalFailed || job.phase === "unknown") {
snapshot.phase = "unknown";
snapshot.error = job.error;
}
await this.journal.write(job.id, snapshot);
if (this.journalFailed)
throw new Error("Earlier recording journal commit remains UNKNOWN.");
});
this.pending = task.catch(error => { this.journalFailed = true; job.phase = "unknown"; job.error = `Journal commit unconfirmed: ${errorText(error)}`; });
return task;
}
arm() { this.timer ??= setInterval(() => this.tick(), 100); this.timer.unref?.(); }
tick() {
const j = this.job;
if (!j || j.phase === "finished" || j.phase === "failed")
return;
const now = this.mono();
if ((j.phase === "queued" || j.phase === "starting") && now >= this.startDeadline) {
if (!j.delivered) {
j.phase = "failed";
j.stopped = true;
j.released = true;
j.error = "Recording expired before delivery; no recorder was started.";
void this.completeWithoutCapture();
}
else
this.stop(j.id);
}
if (j.delivered && ["starting", "recording"].includes(j.phase) && now - (this.views.get(j.binding.viewId)?.lastSeen ?? 0) >= 2500)
this.stop(j.id);
if (now >= this.stopDeadline && !j.finalizationDeadline)
this.stop(j.id);
if (this.finalizationDeadline !== undefined && now >= this.finalizationDeadline && j.phase !== "unknown") {
j.phase = "unknown";
j.error = "Recording finalization was not acknowledged within 30 seconds. Preserve the original browser and capture.";
void this.persist().catch(() => { });
}
}
async start(binding, duration) {
if (!this.admission || this.busy)
throw new Error("A recording or close operation is already pending or UNKNOWN.");
const v = this.views.get(binding.viewId);
if (!v || this.mono() - v.lastSeen > 2500 || !v.status.available || !v.status.visible || v.status.pending || v.status.phase !== "idle" ||
v.status.generation !== binding.recordingGeneration || v.status.owner !== null)
throw new Error("Choose a current visible idle recording view from status.");
if ([...this.views.values()].some(view => view.status.owner === "human" || view.status.pending || view.status.phase === "unknown"))
throw new Error("A person is recording or saving in another preview. Finish that capture first.");
const durationMs = recordingDuration(duration), now = this.now();
const job = { version: 1, id: randomUUID(), cartridgeType: this.cartridgeType, binding: { ...binding }, phase: "queued", durationMs, admittedAt: now, startBy: now + 5000, stopBy: now + 5000 + durationMs,
delivered: false, uploadAttempted: false, stopped: false, released: false, browserAcknowledgedSave: false };
this.reserving = true;
try {
await this.journal.reserve(job.id, job);
this.job = job;
this.completion = undefined;
this.completed = false;
this.journalFailed = false;
this.startDeadline = this.mono() + 5000;
this.stopDeadline = this.startDeadline + durationMs;
this.finalizationDeadline = undefined;
this.arm();
if (!this.admission)
this.stop(job.id);
return structuredClone(job);
}
finally {
this.reserving = false;
}
}
async read(id) {
recordingId(id);
if (this.job?.id === id) {
this.tick();
await this.pending;
await this.completion;
return structuredClone(this.job);
}
const old = await this.journal.read(id);
// Restart never turns an unacknowledged job into success, even if a media file exists.
if (old.phase !== "finished" && old.phase !== "failed") {
old.phase = "unknown";
old.error = "The original recording owner is unavailable. Preserve its original receipt; new conflicting recordings are blocked.";
}
return old;
}
stop(id) {
if (!this.job || this.job.id !== recordingId(id))
throw new Error("This recording belongs to another or unavailable owner; read its original status.");
const j = this.job;
if (j.phase === "queued" && !j.delivered) {
j.phase = "failed";
j.stopped = true;
j.released = true;
j.error = "Stopped before delivery; no recorder was started.";
void this.completeWithoutCapture();
return structuredClone(j);
}
if (j.phase !== "finished" && j.phase !== "failed" && !j.finalizationDeadline) {
j.finalizationDeadline = this.now() + RECORDING_DRAIN_MS;
this.finalizationDeadline = this.mono() + RECORDING_DRAIN_MS;
if (j.phase !== "unknown")
j.phase = "stopping";
void this.persist().catch(() => { });
}
return structuredClone(j);
}
poll(viewId, value) {
if (value === undefined)
return undefined;
const status = parseRecordingState(value), now = this.mono();
const retired = this.retired.get(viewId);
const idle = status.owner === null && status.phase === "idle" && !status.pending && !status.recordingId && !status.closeEpoch;
if (retired) {
if (status.retireId === retired && idle)
return { retired };
throw new Error("This recording view has retired; stale activity is not accepted.");
}
if (status.retireId && this.cancelledEpoch && status.closeEpoch === this.cancelledEpoch)
return { cancelClose: this.cancelledEpoch };
if (status.retireId && !this.closeState) {
const previous = this.views.get(viewId);
if (!previous || status.generation !== previous.generation || !idle || this.reserving || this.job?.binding.viewId === viewId && this.busy)
throw new Error("Recording retirement requires a known idle view with no pending owner.");
if (this.retired.size >= 256)
throw new Error("Recording view retirement limit reached. Close this preview before opening another.");
this.retired.set(viewId, status.retireId);
this.views.delete(viewId);
return { retired: status.retireId };
}
if (!this.views.has(viewId) && this.views.size >= 8)
throw new Error("This recording listener has reached its eight-view limit.");
const v = this.views.get(viewId) ?? { status, lastSeen: now, generation: status.generation };
if (status.generation >= v.generation) {
v.status = status;
v.generation = status.generation;
}
v.lastSeen = now;
this.views.set(viewId, v);
const reply = {};
this.tick();
const close = this.closeState;
if (close) {
reply.close = close.id;
if (close.terminal)
reply.terminalClose = true;
if (status.closeEpoch === close.id && status.generation >= v.generation) {
v.ack = close.id;
if (status.owner === "human" && (status.pending || status.phase !== "idle"))
close.human = true;
}
}
else if (this.cancelledEpoch)
reply.cancelClose = this.cancelledEpoch;
const j = this.job;
if (j && j.binding.viewId === viewId) {
if (!j.delivered && j.phase === "queued" && !close) {
j.delivered = true;
j.phase = "starting";
void this.persist().catch(() => { });
reply.start = { id: j.id, generation: j.binding.recordingGeneration, durationMs: j.durationMs, startBy: j.startBy, stopBy: j.stopBy };
}
if (status.recordingId === j.id) {
if (status.owner === "mcp" && status.generation === j.binding.recordingGeneration) {
const previous = JSON.stringify([j.diagnostics, j.audioIncluded, j.audioUnavailableReason]);
if (status.diagnostics)
j.diagnostics = status.diagnostics;
if (status.audioIncluded !== undefined) {
j.audioIncluded = status.audioIncluded;
j.audioUnavailableReason = status.audioUnavailableReason;
}
// Late evidence can improve diagnostics, never clear UNKNOWN. Save
// it even when no phase transition remains to trigger a journal write.
if (previous !== JSON.stringify([j.diagnostics, j.audioIncluded, j.audioUnavailableReason]))
void this.persist().catch(() => { });
}
if (status.phase === "failed" && status.stopped && status.released && !j.uploadAttempted && j.phase !== "failed" && j.phase !== "finished" && j.phase !== "unknown") {
j.phase = "failed";
j.stopped = true;
j.released = true;
j.error = status.error;
void this.completeWithoutCapture();
}
else if (status.generation !== j.binding.recordingGeneration && j.phase !== "failed" && j.phase !== "finished") {
j.phase = "unknown";
j.error = "Recording generation changed unexpectedly.";
void this.persist().catch(() => { });
}
else {
if (status.phase === "recording" && j.phase === "starting") {
j.phase = "recording";
void this.persist().catch(() => { });
}
if (status.phase === "stopping" || status.phase === "saving" || status.phase === "unknown")
this.stop(j.id);
j.stopped ||= status.stopped === true;
j.released ||= status.released === true;
if (j.saved && status.saved && savedRecordingKey(j.saved) === savedRecordingKey(status.saved) && j.stopped && j.released && !j.browserAcknowledgedSave) {
j.browserAcknowledgedSave = true;
if (j.phase !== "unknown") {
j.phase = "finished";
void this.completeWithoutCapture();
}
else
void this.persist().catch(() => { });
}
}
}
if (j.phase === "starting" || j.phase === "recording")
reply.lease = j.id;
if (j.finalizationDeadline && j.phase !== "finished" && j.phase !== "failed")
reply.stop = j.id;
if (this.completed && (j.phase === "finished" || j.phase === "failed"))
reply.settled = j.id;
if (j.saved && !j.browserAcknowledgedSave)
reply.saved = { recordingId: j.id, capture: { ...j.saved } };
}
return reply;
}
completeWithoutCapture() {
return this.completion ??= this.finish();
}
async finish() {
try {
await this.persist();
if (this.journalFailed || this.job?.phase === "unknown")
throw new Error("Recording journal remains UNKNOWN.");
await this.journal.release(this.job.id);
this.completed = true;
}
catch (error) {
if (this.job) {
this.job.phase = "unknown";
this.job.error = errorText(error);
}
}
if (!this.busy && this.timer) {
clearInterval(this.timer);
this.timer = undefined;
}
}
/** Called before consuming any multipart bytes, including attempts to omit ownership fields. */
admitUpload(id, viewId, generation) {
if (typeof viewId === "string" && this.retired.has(viewId))
throw new Error("A retired recording view cannot upload a capture.");
const j = this.job;
if (id === undefined) {
if (typeof viewId !== "string" || typeof generation !== "string" || !/^(0|[1-9][0-9]*)$/.test(generation))
throw new Error("Capture upload requires its original view ownership fields.");
const v = this.views.get(viewId), g = Number(generation);
if (!v || !Number.isSafeInteger(g) || j?.binding.viewId === viewId && this.busy && j.phase !== "failed" || v.status.owner === "mcp" || v.uploadedGeneration === g || g < v.generation || g > v.generation + 1 ||
g === v.generation && v.status.owner !== "human" || this.closeState && v.ack === this.closeState.id && g > v.generation)
throw new Error("Human capture upload ownership is stale, repeated, or blocked by close.");
v.generation = g;
v.uploadedGeneration = g;
return;
}
if (!j || j.id !== recordingId(id) || j.binding.viewId !== viewId || String(j.binding.recordingGeneration) !== generation || !j.delivered ||
j.uploadAttempted || ["finished", "failed", "unknown"].includes(j.phase))
throw new Error("Recording upload ownership is stale, repeated, or unresolved.");
j.uploadAttempted = true;
this.stop(j.id);
j.phase = "saving";
void this.persist().catch(() => { });
return j;
}
async published(id, saved, outcome = "complete") {
const j = this.job;
if (!j || j.id !== id)
throw new Error("Recording owner changed during publication.");
j.saved = { ...saved };
j.outcome = outcome;
try {
await this.persist();
}
catch (error) {
j.phase = "unknown";
j.error = "Media published but journal commit is unconfirmed. Preserve the original publication; do not upload again.";
throw error;
}
}
async uploadFailed(id, error) {
const j = this.job;
if (!j || j.id !== id)
return;
j.phase = "unknown";
j.error = errorText(error);
const partial = error?.partial;
if (partial)
j.partial = partial;
await this.persist().catch(() => { });
}
markUnknown(message) {
if (this.job && this.busy) {
this.job.phase = "unknown";
this.job.error = message.slice(0, 512);
void this.persist().catch(() => { });
}
}
async drain(terminal = false) {
this.admission = false;
const close = this.closeState ??= { id: randomUUID(), deadline: this.mono() + RECORDING_DRAIN_MS, human: false, terminal };
close.terminal ||= terminal;
this.arm();
if (this.job && this.busy)
this.stop(this.job.id);
let expired = false, timer;
const work = async () => {
for (;;) {
this.tick();
await this.pending;
await this.completion;
if (expired)
throw new Error("Recording drain expired.");
if (close.human && !terminal) {
this.cancelledEpoch = close.id;
this.closeState = undefined;
this.admission = true;
if (!this.busy && this.timer) {
clearInterval(this.timer);
this.timer = undefined;
}
throw new Error("A person is recording or saving in this preview. Finish that capture before closing.");
}
const acknowledged = [...this.views.values()].every(v => v.ack === close.id && !v.status.pending && (v.status.phase === "idle" || v.status.owner === "mcp" && (v.status.phase === "finished" || v.status.phase === "failed")));
if (acknowledged && !this.busy)
return;
if (this.mono() >= close.deadline)
throw new Error("Browser recording close acknowledgement is UNKNOWN.");
await new Promise(resolve => setTimeout(resolve, 50));
}
};
try {
await Promise.race([work(), new Promise((_, reject) => { timer = setTimeout(() => { expired = true; reject(new Error("Browser recording close acknowledgement is UNKNOWN.")); }, this.remainingDrainMs); })]);
}
catch (error) {
if (this.closeState === close) {
expired = true;
this.markUnknown("Browser close or journal finalization is UNKNOWN. Preserve the original capture.");
}
throw error;
}
finally {
if (timer)
clearTimeout(timer);
}
}
dispose() { if (this.timer)
clearInterval(this.timer); this.timer = undefined; }
}
//# sourceMappingURL=web-preview-recording.js.mapSHA-256: e4327cd2bf005dae608f110fc8750498124ecbdcac7a95c1a78e910804286d36