| |
| |
|
|
| import type { WorkerEvent, WorkerRequest } from '../shared/protocol'; |
| import type { WorkspaceState } from '../shared/types'; |
| import { Agent } from './agent'; |
| import { MiniCpmEngine } from './engine'; |
| import { isModelCached } from './modelCache'; |
| import { Workspace } from './workspace'; |
|
|
| const post = (event: WorkerEvent) => (self as DedicatedWorkerGlobalScope).postMessage(event); |
|
|
| |
|
|
| interface Store { |
| load(): Promise<WorkspaceState | undefined>; |
| save(state: WorkspaceState): Promise<void>; |
| } |
|
|
| async function openStore(): Promise<Store> { |
| const db = await new Promise<IDBDatabase>((resolve, reject) => { |
| const req = indexedDB.open('minicpm5-pi-react-workspace', 1); |
| req.onupgradeneeded = () => req.result.createObjectStore('state'); |
| req.onsuccess = () => resolve(req.result); |
| req.onerror = () => reject(req.error); |
| }); |
| return { |
| load: () => |
| new Promise((resolve, reject) => { |
| const req = db.transaction('state').objectStore('state').get('current'); |
| req.onsuccess = () => resolve(req.result as WorkspaceState | undefined); |
| req.onerror = () => reject(req.error); |
| }), |
| save: (state) => |
| new Promise((resolve, reject) => { |
| const tx = db.transaction('state', 'readwrite'); |
| tx.objectStore('state').put(state, 'current'); |
| tx.oncomplete = () => resolve(); |
| tx.onerror = () => reject(tx.error); |
| tx.onabort = () => reject(tx.error); |
| }), |
| }; |
| } |
|
|
| |
|
|
| const engine = new MiniCpmEngine(); |
| let store: Store | undefined; |
| let workspace: Workspace; |
| let agent: Agent; |
| let busy = false; |
| let controller: AbortController | undefined; |
|
|
| async function currentState(): Promise<WorkspaceState> { |
| return { |
| version: 1, |
| files: await workspace.snapshot(), |
| entries: await workspace.serialize(), |
| messages: agent.messages, |
| }; |
| } |
|
|
| async function persist(): Promise<WorkspaceState> { |
| const state = await currentState(); |
| try { |
| await store?.save(state); |
| } catch (err) { |
| post({ type: 'persistence_error', error: 'Changes are in memory, but browser storage failed: ' + (err as Error).message }); |
| } |
| post({ type: 'workspace', files: state.files }); |
| return state; |
| } |
|
|
| const initialized = (async () => { |
| let saved: WorkspaceState | undefined; |
| try { |
| store = await openStore(); |
| saved = await store.load(); |
| } catch (err) { |
| post({ type: 'persistence_error', error: 'Workspace saving is unavailable: ' + (err as Error).message }); |
| } |
| workspace = new Workspace(saved?.entries, saved?.files); |
| await workspace.ready; |
| agent = new Agent({ engine, workspace, emit: post, persist: async () => void (await persist()) }, saved?.messages ?? []); |
|
|
| const state = await currentState(); |
| if (!saved) { |
| try { |
| await store?.save(state); |
| } catch (err) { |
| post({ type: 'persistence_error', error: 'Changes are in memory, but browser storage failed: ' + (err as Error).message }); |
| } |
| } |
| post({ type: 'initialized', ...state }); |
| })(); |
|
|
| |
|
|
| self.onmessage = async ({ data }: MessageEvent<WorkerRequest>) => { |
| const { id, action } = data; |
| const args = (data.args ?? {}) as Record<string, any>; |
| const reply = (value?: unknown) => post({ type: 'reply', id, value }); |
|
|
| try { |
| await initialized; |
|
|
| if (action === 'stop') { |
| controller?.abort(); |
| engine.stop(); |
| return reply(true); |
| } |
| if (action === 'snapshot') return reply(await currentState()); |
| if (action === 'cache_status') return reply(await isModelCached()); |
|
|
| if (busy) throw new Error('Wait for the current operation or stop it first.'); |
| busy = true; |
| controller = new AbortController(); |
| post({ type: 'busy', action, cachedOnly: args.cachedOnly === true }); |
|
|
| try { |
| let value: unknown; |
| switch (action) { |
| case 'load': |
| value = await engine.load({ cachedOnly: args.cachedOnly === true }, controller.signal, post); |
| break; |
| case 'prompt': |
| await agent.prompt(args.text, controller.signal); |
| value = await currentState(); |
| break; |
| case 'shell': |
| value = await workspace.exec(args.command, controller.signal); |
| await persist(); |
| break; |
| case 'write': |
| await workspace.write(args.path, args.content); |
| value = await persist(); |
| break; |
| case 'new_chat': |
| agent.reset(); |
| value = await persist(); |
| break; |
| case 'read': |
| value = await workspace.read(args.path); |
| break; |
| default: |
| throw new Error('Unknown action: ' + action); |
| } |
| reply(value); |
| } finally { |
| busy = false; |
| controller = undefined; |
| post({ type: 'idle' }); |
| } |
| } catch (err) { |
| const e = err as Error & { code?: string }; |
| post({ |
| type: 'reply', |
| id, |
| error: e?.name === 'AbortError' ? 'Stopped.' : String(e?.message ?? e), |
| errorName: e?.name, |
| errorCode: e?.code, |
| }); |
| } |
| }; |
|
|