PI / src /client /agentClient.ts
noodcon's picture
Upload 64 files
9b9eafc verified
Raw
History Blame Contribute Delete
2.71 kB
// Cầu nối main thread ↔ worker: gọi lệnh dạng Promise và phát sự kiện cho giao diện.
import type { WorkerAction, WorkerEvent } from '../shared/protocol';
export class WorkerError extends Error {
code?: string;
constructor(message: string, name?: string, code?: string) {
super(message);
this.name = name ?? 'Error';
this.code = code;
}
}
type Listener = (event: WorkerEvent) => void;
export class AgentClient {
private worker: Worker;
private nextId = 1;
private pending = new Map<number, { resolve: (v: any) => void; reject: (e: Error) => void }>();
private listeners = new Set<Listener>();
/** Kết quả `initialized` đầu tiên — giao diện chờ nó trước khi cho phép thao tác. */
readonly ready: Promise<Extract<WorkerEvent, { type: 'initialized' }>>;
constructor() {
this.worker = new Worker(new URL('../worker/agent.worker.ts', import.meta.url), { type: 'module' });
let resolveReady!: (e: Extract<WorkerEvent, { type: 'initialized' }>) => void;
this.ready = new Promise((resolve) => (resolveReady = resolve));
this.worker.onmessage = ({ data }: MessageEvent<WorkerEvent>) => {
if (data.type === 'reply') {
const p = this.pending.get(data.id);
if (!p) return;
this.pending.delete(data.id);
if (data.error !== undefined) p.reject(new WorkerError(data.error, data.errorName, data.errorCode));
else p.resolve(data.value);
return;
}
if (data.type === 'initialized') resolveReady(data);
this.listeners.forEach((l) => l(data));
};
const fail = (message: string) => {
const error = new WorkerError(message);
this.pending.forEach((p) => p.reject(error));
this.pending.clear();
this.listeners.forEach((l) => l({ type: 'fatal', error: message }));
};
this.worker.onerror = (e) => fail(e.message || 'The background worker stopped unexpectedly. Reload the page.');
this.worker.onmessageerror = () => fail('The background worker sent an unreadable message.');
}
call<T = unknown>(action: WorkerAction['action'], args?: Record<string, unknown>): Promise<T> {
const id = this.nextId++;
return new Promise<T>((resolve, reject) => {
this.pending.set(id, { resolve, reject });
this.worker.postMessage({ id, action, args });
});
}
subscribe(listener: Listener): () => void {
this.listeners.add(listener);
return () => this.listeners.delete(listener);
}
}
let singleton: AgentClient | undefined;
/** Một worker duy nhất cho cả ứng dụng (an toàn với React StrictMode). */
export function getClient(): AgentClient {
return (singleton ??= new AgentClient());
}