| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146 |
- import { appendFile, mkdir, readFile, rename, rm, stat } from 'node:fs/promises';
- import { join } from 'node:path';
- import type { AgentEvent } from './types';
- /** 移植自 Noobi.ai src/main/eventLog.ts,维度由 projectId 改为 threadId */
- const MAX_EVENT_BYTES = 64 * 1024;
- const MAX_TAIL_EVENTS = 400;
- const MAX_LOG_BYTES = 8 * 1024 * 1024;
- export class EventLog {
- readonly #directory: string;
- readonly #queues = new Map<string, Promise<void>>();
- constructor(directory: string) {
- this.#directory = directory;
- }
- async init(): Promise<void> {
- await mkdir(this.#directory, { recursive: true });
- }
- async append(event: AgentEvent): Promise<void> {
- validateThreadId(event.threadId);
- const serialized = serializeEvent(event);
- const previous = this.#queues.get(event.threadId) ?? Promise.resolve();
- const current = previous
- .catch(() => undefined)
- .then(async () => {
- await rotateIfNeeded(this.#path(event.threadId));
- await appendFile(this.#path(event.threadId), `${serialized}\n`, {
- encoding: 'utf8',
- mode: 0o600,
- });
- });
- this.#queues.set(event.threadId, current);
- const release = (): void => {
- if (this.#queues.get(event.threadId) === current) this.#queues.delete(event.threadId);
- };
- void current.then(release, release);
- await current;
- }
- /** 只取尾部若干条,并把同一 id 的增量拼回完整消息 */
- async read(threadId: string, limit = MAX_TAIL_EVENTS): Promise<AgentEvent[]> {
- validateThreadId(threadId);
- await this.#queues.get(threadId)?.catch(() => undefined);
- let source: string;
- try {
- source = await readFile(this.#path(threadId), 'utf8');
- } catch (error) {
- if (asNodeError(error).code === 'ENOENT') return [];
- throw error;
- }
- const events = new Map<string, AgentEvent>();
- const lines = source.split(/\r?\n/u).filter(Boolean).slice(-limit);
- for (const line of lines) {
- try {
- const event = JSON.parse(line) as AgentEvent;
- if (event.threadId === threadId && typeof event.id === 'string') {
- const previous = events.get(event.id);
- events.set(
- event.id,
- previous && event.isDelta
- ? {
- ...previous,
- ...event,
- message: `${previous.message}${event.message}`.slice(-120_000),
- }
- : event,
- );
- }
- } catch {
- // 最后一行可能只写了一半,或是不兼容的旧记录,忽略
- }
- }
- return [...events.values()].sort((left, right) => left.timestamp.localeCompare(right.timestamp));
- }
- async flush(): Promise<void> {
- await Promise.allSettled([...this.#queues.values()]);
- }
- async remove(threadId: string): Promise<void> {
- validateThreadId(threadId);
- await this.#queues.get(threadId)?.catch(() => undefined);
- await Promise.all([
- rm(this.#path(threadId), { force: true }),
- rm(`${this.#path(threadId)}.previous`, { force: true }),
- ]);
- }
- #path(threadId: string): string {
- return join(this.#directory, `${threadId}.jsonl`);
- }
- }
- /** 超过 64 KiB 的事件用二分法截断 message,保证元数据仍能落盘 */
- function serializeEvent(event: AgentEvent): string {
- const serialized = JSON.stringify(event);
- if (Buffer.byteLength(serialized, 'utf8') <= MAX_EVENT_BYTES) return serialized;
- const marker = '\n…(事件日志已按 64 KiB 截断)';
- const codePoints = Array.from(event.message);
- let lower = 0;
- let upper = codePoints.length;
- let best = JSON.stringify({ ...event, message: marker });
- if (Buffer.byteLength(best, 'utf8') > MAX_EVENT_BYTES) {
- throw new Error('事件元数据超过 64 KiB 落盘上限');
- }
- while (lower <= upper) {
- const middle = Math.floor((lower + upper) / 2);
- const candidate = JSON.stringify({
- ...event,
- message: `${codePoints.slice(0, middle).join('')}${marker}`,
- });
- if (Buffer.byteLength(candidate, 'utf8') <= MAX_EVENT_BYTES) {
- best = candidate;
- lower = middle + 1;
- } else {
- upper = middle - 1;
- }
- }
- return best;
- }
- async function rotateIfNeeded(path: string): Promise<void> {
- try {
- if ((await stat(path)).size <= MAX_LOG_BYTES) return;
- const previous = `${path}.previous`;
- await rm(previous, { force: true });
- await rename(path, previous);
- } catch (error) {
- if (asNodeError(error).code !== 'ENOENT') throw error;
- }
- }
- function validateThreadId(value: string): void {
- if (!/^[a-zA-Z0-9_-]{1,128}$/u.test(value)) throw new Error('非法的 threadId');
- }
- function asNodeError(value: unknown): NodeJS.ErrnoException {
- return value as NodeJS.ErrnoException;
- }
|