Move app to apps/web-react in pnpm workspace layout

This commit is contained in:
Kalle
2026-08-16 09:20:33 +03:00
parent 7e365ccfcf
commit bbc8ea57af
2913 changed files with 227 additions and 250 deletions

View File

@@ -0,0 +1,320 @@
/**
* AnalyzerWorker: owns OpenCV.js (WASM), the detector registry and a
* DetectorScheduler. Two entry points: "frame" — the main thread posts one
* ImageBitmap/VideoFrame at a time (live capture, screenshot harness, VoD
* seek fallback); results come back per detector, then a "done" carrying the
* scheduler's calm signal and telemetry. "scanChunk" — a VoD time slice is
* demuxed/decoded entirely in the worker with mediabunny: the worker owns a
* contiguous slice so scheduling is exact, undue frames skip canvas
* readback, and calm stretches skim by keyframe hops instead of decoding
* every frame — the big VoD speedup, since sequential decode bounds scan
* wall-clock time.
*/
import {
ALL_FORMATS,
BlobSource,
EncodedPacketSink,
Input,
type VideoSample,
VideoSampleSink,
} from "mediabunny";
import { loadOpenCV } from "../core/cv";
import { MAP_START_EVENT_TYPE } from "../core/detectors/map-start/index";
import {
createAllDetectors,
SCOREBOARD_EVENT_TYPES,
} from "../core/detectors/registry";
import { DetectorScheduler } from "../core/detectors/scheduler";
import {
createScanTelemetry,
detectorTelemetry,
type ScanTelemetry,
} from "../core/detectors/telemetry";
import type { Detector } from "../core/detectors/types";
import { normalizeFrame, toMat } from "../core/image";
import type {
AnalyzeRequest,
InitRequest,
ScanChunkRequest,
WorkerRequest,
WorkerResponse,
} from "./protocol";
import { fetchScoreboardResources } from "./resources";
/**
* Widest skim hop: calm footage is sampled at the keyframe cadence, capped
* here so long-GOP recordings still cannot slip a results screen (~10s) or
* a match intro (~7s) between two samples.
*/
const MAX_SKIM_STRIDE_S = 2.5;
const PROGRESS_POST_INTERVAL_MS = 400;
const PREVIEW_POST_INTERVAL_MS = 600;
const PREVIEW_WIDTH = 480;
const PREVIEW_HEIGHT = 270;
let detectors: Detector<unknown>[] = [];
let scheduler: DetectorScheduler | null = null;
/** null unless the init message asked for telemetry */
let telemetry: ScanTelemetry | null = null;
let collectTelemetry = false;
let chunkAborted = false;
/** last per-frame t, to reset telemetry when a new session rewinds the clock */
let lastFrameT = Number.NEGATIVE_INFINITY;
function post(message: WorkerResponse, transfer: Transferable[] = []): void {
self.postMessage(message, { transfer });
}
async function init({
assetsBaseUrl,
suppressSteadyFrames = true,
collectTelemetry: collect = false,
}: InitRequest): Promise<void> {
try {
await loadOpenCV();
const resources = await fetchScoreboardResources(assetsBaseUrl);
detectors = createAllDetectors(resources);
scheduler = new DetectorScheduler(detectors, {
suppressSteadyFrames,
matchOpeningTypes: [MAP_START_EVENT_TYPE],
matchClosingTypes: SCOREBOARD_EVENT_TYPES,
});
collectTelemetry = collect;
telemetry = freshTelemetry();
post({ kind: "ready" });
} catch (error) {
post({ kind: "error", message: `init failed: ${String(error)}` });
}
}
/**
* Run the due detectors over one frame; closes `bitmap`. When the scheduler
* has no detector due, the canvas readback and normalize are skipped too.
*/
async function analyzeFrame(
bitmap: ImageBitmap | VideoFrame,
t: number,
): Promise<void> {
const due = scheduler!.dueDetectors(t);
if (due.length === 0) {
bitmap.close();
return;
}
const width = "displayWidth" in bitmap ? bitmap.displayWidth : bitmap.width;
const height =
"displayHeight" in bitmap ? bitmap.displayHeight : bitmap.height;
const canvas = new OffscreenCanvas(width, height);
const ctx = canvas.getContext("2d", { willReadFrequently: true })!;
ctx.drawImage(bitmap, 0, 0);
bitmap.close();
const imageData = ctx.getImageData(0, 0, canvas.width, canvas.height);
const src = toMat({
width: imageData.width,
height: imageData.height,
data: imageData.data,
});
let frame: ReturnType<typeof normalizeFrame>;
try {
frame = normalizeFrame(src);
} finally {
src.delete();
}
if (telemetry) telemetry.analyzedFrames++;
// On detection, ship back the exact analyzed pixels (lossless, at capture
// resolution) so the UI never has to re-grab a later frame — encoded at
// most once per frame, however many detectors fire on it.
let encoded: Promise<Blob> | null = null;
const frameBlob = () =>
(encoded ??= canvas.convertToBlob({ type: "image/png" }));
try {
for (const detector of detectors) {
if (!due.includes(detector.id)) continue;
const counters = telemetry
? detectorTelemetry(telemetry, detector.id)
: null;
const gateStart = counters ? performance.now() : 0;
const gate = detector.gate(frame);
if (counters) {
counters.checks++;
counters.gateMs += performance.now() - gateStart;
}
scheduler!.recordGate(detector.id, t, gate.pass, gate.signature);
if (counters && gate.pass) counters.gatePasses++;
const runParse = gate.pass && scheduler!.shouldParse(detector.id, t);
if (counters && gate.pass && !runParse) counters.suppressedParses++;
let events: ReturnType<typeof detector.parse> = [];
if (runParse) {
const parseStart = counters ? performance.now() : 0;
events = detector.parse(frame, t, gate);
if (counters) {
counters.parses++;
counters.parseMs += performance.now() - parseStart;
}
scheduler!.recordParse(detector.id, t, events);
}
const blob =
events.length > 0 && detector.attachFrame !== false
? await frameBlob()
: undefined;
post({
kind: "result",
detector: detector.id,
t,
gate,
events,
frame: blob,
});
}
} finally {
frame.delete();
}
}
async function analyze({ bitmap, t }: AnalyzeRequest): Promise<void> {
if (t + 5 < lastFrameT) telemetry = freshTelemetry();
lastFrameT = t;
try {
await analyzeFrame(bitmap, t);
} catch (error) {
post({ kind: "error", message: `analyze failed: ${String(error)}` });
}
post({ kind: "done", t, calm: scheduler!.calm(t), telemetry });
}
async function scanChunk({
file,
chunkIndex,
tStart,
tEnd,
}: ScanChunkRequest): Promise<void> {
chunkAborted = false;
scheduler!.reset(tStart);
telemetry = freshTelemetry();
const wallStart = performance.now();
let lastProgressAt = 0;
let lastPreviewAt = 0;
let cursor = tStart;
let mode: "active" | "skim" = "active";
const input = new Input({
formats: ALL_FORMATS,
source: new BlobSource(file),
});
try {
const track = await input.getPrimaryVideoTrack();
if (!track || !(await track.canDecode())) {
throw new Error("worker cannot decode this file");
}
const samples = new VideoSampleSink(track);
const packets = new EncodedPacketSink(track);
const handleSample = async (sample: VideoSample): Promise<void> => {
const t = sample.timestamp;
if (telemetry) {
telemetry.decodedFrames++;
const span = Math.max(0, t - cursor);
if (mode === "active") telemetry.activeVideoS += span;
else telemetry.skimVideoS += span;
}
cursor = Math.max(cursor, t);
const frame = sample.toVideoFrame();
sample.close();
const now = performance.now();
let preview: ImageBitmap | undefined;
if (now - lastPreviewAt >= PREVIEW_POST_INTERVAL_MS) {
lastPreviewAt = now;
preview = await createImageBitmap(frame, {
resizeWidth: PREVIEW_WIDTH,
resizeHeight: PREVIEW_HEIGHT,
});
}
if (t >= scheduler!.nextDueT()) {
await analyzeFrame(frame, t);
} else {
frame.close();
}
if (preview || now - lastProgressAt >= PROGRESS_POST_INTERVAL_MS) {
lastProgressAt = now;
if (telemetry) telemetry.wallMs = performance.now() - wallStart;
post(
{
kind: "chunkProgress",
chunkIndex,
t: cursor,
mode,
telemetry,
preview,
},
preview ? [preview] : [],
);
}
};
scan: while (!chunkAborted && cursor < tEnd) {
if (mode === "active") {
// dense sequential decode: every frame is seen, the scheduler
// decides which are worth analyzing
for await (const sample of samples.samples(cursor)) {
if (!sample) continue;
if (chunkAborted || sample.timestamp >= tEnd) {
sample.close();
break scan;
}
await handleSample(sample);
if (scheduler!.calm(cursor)) {
mode = "skim";
break;
}
}
if (mode === "active") break; // media ended before tEnd
} else {
// skim: hop keyframe to keyframe (single-frame decodes) while
// calm, capped so long GOPs cannot hide a short screen
const key = await packets.getKeyPacket(cursor + MAX_SKIM_STRIDE_S, {
verifyKeyPackets: true,
});
const target =
key && key.timestamp > cursor
? key.timestamp
: cursor + MAX_SKIM_STRIDE_S;
if (target >= tEnd) {
cursor = tEnd;
break;
}
const sample = await samples.getSample(target);
if (!sample) {
cursor = target;
continue;
}
await handleSample(sample);
cursor = Math.max(cursor, target);
if (!scheduler!.calm(cursor)) mode = "active";
}
}
if (telemetry) telemetry.wallMs = performance.now() - wallStart;
post({ kind: "chunkDone", chunkIndex, telemetry });
} catch (error) {
post({
kind: "error",
message: `chunk ${chunkIndex} scan failed: ${String(error)}`,
});
} finally {
input.dispose();
}
}
function freshTelemetry(): ScanTelemetry | null {
return collectTelemetry ? createScanTelemetry() : null;
}
self.onmessage = (e: MessageEvent) => {
const msg = e.data as WorkerRequest;
if (msg.kind === "init") void init(msg);
else if (msg.kind === "frame") void analyze(msg);
else if (msg.kind === "scanChunk") void scanChunk(msg);
else if (msg.kind === "abortChunk") chunkAborted = true;
};

View File

@@ -0,0 +1,186 @@
/**
* Main-thread wrapper around the AnalyzerWorker: init handshake, then either
* one in-flight frame at a time (live capture / screenshot / seek fallback —
* each frame yields one result per due detector, then a "done" carrying the
* scheduler's calm signal and telemetry) or one in-flight chunk scan (the
* worker decodes and analyzes a VoD time slice by itself, streaming results
* and progress until "chunkDone").
*/
import { Config } from "../../../config";
import type { ScanTelemetry } from "../core/detectors/telemetry";
import type { WorkerResponse } from "./protocol";
export type ResultHandler = (
result: Extract<WorkerResponse, { kind: "result" }>,
) => void;
export type ErrorHandler = (message: string) => void;
export interface DoneInfo {
calm: boolean;
/** null unless the client was created with collectTelemetry */
telemetry: ScanTelemetry | null;
}
export type DoneHandler = (t: number, info: DoneInfo) => void;
export type ChunkProgress = Extract<WorkerResponse, { kind: "chunkProgress" }>;
export type ChunkProgressHandler = (progress: ChunkProgress) => void;
/** each chunk-scanning worker decodes and analyzes on its own; leave a core
* for the main thread and one for the browser's media stack */
export function defaultScanWorkerCount(): number {
return Math.min(4, Math.max(1, (navigator.hardwareConcurrency || 4) - 2));
}
interface PendingChunk {
resolve(telemetry: ScanTelemetry | null): void;
reject(error: Error): void;
onProgress?: ChunkProgressHandler;
}
export class AnalyzerClient {
#worker: Worker;
#ready = false;
#busy = false;
#onResult: ResultHandler;
#onError: ErrorHandler;
#onDone: DoneHandler | undefined;
#readyPromise: Promise<void>;
#rejectReady: ((error: Error) => void) | undefined;
#idleWaiters: (() => void)[] = [];
#chunk: PendingChunk | null = null;
constructor(
onResult: ResultHandler,
// biome-ignore lint/suspicious/noConsole: default sink for worker errors when no handler is passed
onError: ErrorHandler = console.error,
onDone?: DoneHandler,
options: {
suppressSteadyFrames?: boolean;
collectTelemetry?: boolean;
} = {},
) {
this.#onResult = onResult;
this.#onError = onError;
this.#onDone = onDone;
this.#worker = new Worker(
new URL("./analyzer.worker.ts", import.meta.url),
{
type: "module",
},
);
let resolveReady!: () => void;
this.#readyPromise = new Promise((resolve, reject) => {
resolveReady = resolve;
this.#rejectReady = reject;
});
// init failure also surfaces via onError; don't let an un-awaited
// whenReady() turn it into an unhandled rejection as well
this.#readyPromise.catch(() => {});
this.#worker.onmessage = (e: MessageEvent<WorkerResponse>) => {
const msg = e.data;
if (msg.kind === "ready") {
this.#ready = true;
resolveReady();
} else if (msg.kind === "result") {
this.#onResult(msg);
} else if (msg.kind === "done") {
this.#settle();
this.#onDone?.(msg.t, { calm: msg.calm, telemetry: msg.telemetry });
} else if (msg.kind === "chunkProgress") {
this.#chunk?.onProgress?.(msg);
} else if (msg.kind === "chunkDone") {
const chunk = this.#chunk;
this.#chunk = null;
this.#settle();
chunk?.resolve(msg.telemetry);
} else if (msg.kind === "error") {
this.#fail(msg.message);
}
};
// A throw outside the worker's own try/catch posts neither "error" nor
// "done"; without these handlers `busy` would stay true forever and the
// sampler / VoD scan would silently freeze.
this.#worker.onerror = (e: ErrorEvent) => {
this.#fail(`worker error: ${e.message || String(e)}`);
};
this.#worker.onmessageerror = () => {
this.#fail("worker message deserialization failed");
};
this.#worker.postMessage({
kind: "init",
assetsBaseUrl: Config.staticAssetsUrl,
suppressSteadyFrames: options.suppressSteadyFrames ?? true,
collectTelemetry: options.collectTelemetry ?? false,
});
}
whenReady(): Promise<void> {
return this.#readyPromise;
}
get busy(): boolean {
return this.#busy || !this.#ready;
}
/** Resolves once no frame or chunk scan is in flight. Call whenReady() first. */
async whenIdle(): Promise<void> {
while (this.#busy) {
await new Promise<void>((resolve) => this.#idleWaiters.push(resolve));
}
}
/** Returns false (and closes the bitmap) if the worker is still busy. */
analyze(bitmap: ImageBitmap | VideoFrame, t: number): boolean {
if (this.busy) {
bitmap.close();
return false;
}
this.#busy = true;
this.#worker.postMessage({ kind: "frame", bitmap, t }, [bitmap]);
return true;
}
/**
* Scan [tStart, tEnd) of `file` inside the worker. Results stream to the
* shared result handler; resolves with the chunk's telemetry once done
* (an aborted chunk resolves too — abort is not an error).
*/
scanChunk(
request: { file: File; chunkIndex: number; tStart: number; tEnd: number },
onProgress?: ChunkProgressHandler,
): Promise<ScanTelemetry | null> {
if (this.busy) {
return Promise.reject(new Error("analyzer is busy"));
}
this.#busy = true;
return new Promise((resolve, reject) => {
this.#chunk = { resolve, reject, onProgress };
this.#worker.postMessage({ kind: "scanChunk", ...request });
});
}
/** Ask a running chunk scan to stop; it resolves after the current frame. */
abortChunk(): void {
if (this.#chunk) this.#worker.postMessage({ kind: "abortChunk" });
}
dispose(): void {
this.#worker.terminate();
}
#settle(): void {
this.#busy = false;
const waiters = this.#idleWaiters;
this.#idleWaiters = [];
for (const waiter of waiters) waiter();
}
#fail(message: string): void {
// an error before "ready" means init failed — reject whenReady() so
// callers don't hang on a client that will never become usable
if (!this.#ready) this.#rejectReady?.(new Error(message));
const chunk = this.#chunk;
this.#chunk = null;
this.#settle();
if (chunk) chunk.reject(new Error(message));
else this.#onError(message);
}
}

View File

@@ -0,0 +1,93 @@
import type { ScanTelemetry } from "../core/detectors/telemetry";
import type { DetectedEvent, GateResult } from "../core/detectors/types";
export interface InitRequest {
kind: "init";
/**
* static assets CDN root the worker fetches icons and atlases from
* (`Config.staticAssetsUrl`); passed in the init message so the worker
* bundle stays free of the app config graph
*/
assetsBaseUrl: string;
/**
* skip parse() for a detector whose gate keeps firing without confidence
* improving (static screen), and let the scheduler thin out checks;
* default true — one-shot consumers like the screenshot harness turn it
* off to get every detector on every frame
*/
suppressSteadyFrames?: boolean;
/**
* accumulate scan telemetry counters (and time the detectors) so they can
* be reported back with progress and done messages; default false — the
* VoD tab only asks for them when the telemetry panel is opted into
*/
collectTelemetry?: boolean;
}
export interface AnalyzeRequest {
kind: "frame";
/** VideoFrame is what VoD decode produces; transferring it directly skips
* a main-thread ImageBitmap conversion */
bitmap: ImageBitmap | VideoFrame;
/** seconds into the stream */
t: number;
}
/**
* Scan a time slice of a VoD entirely inside the worker: demux + decode with
* mediabunny, schedule detectors, post results as they fire. Decoding in the
* worker removes the per-frame main-thread hop and lets each worker own a
* contiguous slice, so scheduler state (cadence, suppression, calm) is exact
* instead of split across a pool.
*/
export interface ScanChunkRequest {
kind: "scanChunk";
file: File;
chunkIndex: number;
/** seconds; the chunk scans [tStart, tEnd) */
tStart: number;
tEnd: number;
}
export interface AbortChunkRequest {
kind: "abortChunk";
}
export type WorkerRequest =
| InitRequest
| AnalyzeRequest
| ScanChunkRequest
| AbortChunkRequest;
export type WorkerResponse =
| { kind: "ready" }
| {
kind: "result";
detector: string;
t: number;
gate: GateResult;
events: DetectedEvent<unknown>[];
/** lossless PNG of the exact frame that was analyzed; present when events fired */
frame?: Blob;
}
/** all due detectors have reported for frame t (per-frame path only) */
| {
kind: "done";
t: number;
/** scheduler sees dead air — the caller may widen its sampling stride */
calm: boolean;
/** null when the worker was not asked to collect telemetry */
telemetry: ScanTelemetry | null;
}
| {
kind: "chunkProgress";
chunkIndex: number;
/** seconds of video the chunk scan has reached */
t: number;
mode: "active" | "skim";
telemetry: ScanTelemetry | null;
/** small bitmap of the latest decoded frame, for the preview canvas */
preview?: ImageBitmap;
}
| { kind: "chunkDone"; chunkIndex: number; telemetry: ScanTelemetry | null }
| { kind: "error"; message: string };

View File

@@ -0,0 +1,91 @@
/**
* Worker/browser IO for ScoreboardResources: fetches over HTTP. Game icons
* come from the CDN's shared `img/<dir>/<id>.avif` sets the rest of the app
* already uses; the scanner-specific atlases (glyphs, planner signatures) are
* served same-origin from this repo's `public/scanner/v1/**`. What the bundle
* contains — every key, template option set, and atlas name — lives in
* core/resources.ts, shared with the Node loader. The base URL arrives via
* the worker init message (see worker/protocol.ts) so this module never
* imports the app config.
*/
import {
loadPlannerStages,
type PlannerManifest,
type PlannerStage,
} from "../core/detectors/minimap/stage";
import type { ScoreboardResources } from "../core/detectors/scoreboard/index";
import { type AtlasMeta, type GlyphSet, loadGlyphSet } from "../core/glyphs";
import type { FrameData } from "../core/image";
import { assembleScoreboardResources } from "../core/resources";
/** Scanner parser atlases; the version segment guards against CDN cache skew —
* bump it together with breaking atlas format changes (must match the
* Node-side SCANNER_ASSETS_DIR default in node/assets-dir.ts).
* xxx: temporarily served same-origin from this repo's public/ while the
* feature is in development; move to the assets repo CDN
* (`${base}/scanner/v1`) later */
const ATLAS_BASE = "/scanner/v1";
async function fetchImage(url: string): Promise<FrameData> {
const res = await fetch(url);
if (!res.ok) throw new Error(`fetch ${url}: ${res.status}`);
const bitmap = await createImageBitmap(await res.blob());
const canvas = new OffscreenCanvas(bitmap.width, bitmap.height);
const ctx = canvas.getContext("2d")!;
ctx.drawImage(bitmap, 0, 0);
bitmap.close();
const data = ctx.getImageData(0, 0, canvas.width, canvas.height);
return { width: data.width, height: data.height, data: data.data };
}
function makeFetchAtlas(base: string) {
return async function fetchAtlas(
name: string,
): Promise<() => GlyphSet | null> {
try {
const [meta, image] = await Promise.all([
fetch(`${base}/glyphs/${name}.json`).then((r) => {
if (!r.ok) throw new Error(String(r.status));
return r.json() as Promise<AtlasMeta>;
}),
fetchImage(`${base}/glyphs/${name}.png`),
]);
const set = loadGlyphSet(image, meta);
return () => set;
} catch {
return () => null;
}
};
}
function makeFetchPlannerStages(base: string) {
return async function fetchPlannerStages(): Promise<
() => PlannerStage[] | null
> {
try {
const [manifest, atlas] = await Promise.all([
fetch(`${base}/planner/manifest.json`).then((r) => {
if (!r.ok) throw new Error(String(r.status));
return r.json() as Promise<PlannerManifest>;
}),
fetchImage(`${base}/planner/signatures.png`),
]);
const stages = loadPlannerStages(atlas, manifest);
return () => stages;
} catch {
return () => null;
}
};
}
/** Requires loadOpenCV() to have resolved. */
export function fetchScoreboardResources(
base: string,
): Promise<ScoreboardResources> {
return assembleScoreboardResources({
readIcon: (dir, id) => fetchImage(`${base}/img/${dir}/${id}.avif`),
loadAtlas: makeFetchAtlas(ATLAS_BASE),
loadPlannerStages: makeFetchPlannerStages(ATLAS_BASE),
});
}