eventLog.ts 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146
  1. import { appendFile, mkdir, readFile, rename, rm, stat } from 'node:fs/promises';
  2. import { join } from 'node:path';
  3. import type { AgentEvent } from './types';
  4. /** 移植自 Noobi.ai src/main/eventLog.ts,维度由 projectId 改为 threadId */
  5. const MAX_EVENT_BYTES = 64 * 1024;
  6. const MAX_TAIL_EVENTS = 400;
  7. const MAX_LOG_BYTES = 8 * 1024 * 1024;
  8. export class EventLog {
  9. readonly #directory: string;
  10. readonly #queues = new Map<string, Promise<void>>();
  11. constructor(directory: string) {
  12. this.#directory = directory;
  13. }
  14. async init(): Promise<void> {
  15. await mkdir(this.#directory, { recursive: true });
  16. }
  17. async append(event: AgentEvent): Promise<void> {
  18. validateThreadId(event.threadId);
  19. const serialized = serializeEvent(event);
  20. const previous = this.#queues.get(event.threadId) ?? Promise.resolve();
  21. const current = previous
  22. .catch(() => undefined)
  23. .then(async () => {
  24. await rotateIfNeeded(this.#path(event.threadId));
  25. await appendFile(this.#path(event.threadId), `${serialized}\n`, {
  26. encoding: 'utf8',
  27. mode: 0o600,
  28. });
  29. });
  30. this.#queues.set(event.threadId, current);
  31. const release = (): void => {
  32. if (this.#queues.get(event.threadId) === current) this.#queues.delete(event.threadId);
  33. };
  34. void current.then(release, release);
  35. await current;
  36. }
  37. /** 只取尾部若干条,并把同一 id 的增量拼回完整消息 */
  38. async read(threadId: string, limit = MAX_TAIL_EVENTS): Promise<AgentEvent[]> {
  39. validateThreadId(threadId);
  40. await this.#queues.get(threadId)?.catch(() => undefined);
  41. let source: string;
  42. try {
  43. source = await readFile(this.#path(threadId), 'utf8');
  44. } catch (error) {
  45. if (asNodeError(error).code === 'ENOENT') return [];
  46. throw error;
  47. }
  48. const events = new Map<string, AgentEvent>();
  49. const lines = source.split(/\r?\n/u).filter(Boolean).slice(-limit);
  50. for (const line of lines) {
  51. try {
  52. const event = JSON.parse(line) as AgentEvent;
  53. if (event.threadId === threadId && typeof event.id === 'string') {
  54. const previous = events.get(event.id);
  55. events.set(
  56. event.id,
  57. previous && event.isDelta
  58. ? {
  59. ...previous,
  60. ...event,
  61. message: `${previous.message}${event.message}`.slice(-120_000),
  62. }
  63. : event,
  64. );
  65. }
  66. } catch {
  67. // 最后一行可能只写了一半,或是不兼容的旧记录,忽略
  68. }
  69. }
  70. return [...events.values()].sort((left, right) => left.timestamp.localeCompare(right.timestamp));
  71. }
  72. async flush(): Promise<void> {
  73. await Promise.allSettled([...this.#queues.values()]);
  74. }
  75. async remove(threadId: string): Promise<void> {
  76. validateThreadId(threadId);
  77. await this.#queues.get(threadId)?.catch(() => undefined);
  78. await Promise.all([
  79. rm(this.#path(threadId), { force: true }),
  80. rm(`${this.#path(threadId)}.previous`, { force: true }),
  81. ]);
  82. }
  83. #path(threadId: string): string {
  84. return join(this.#directory, `${threadId}.jsonl`);
  85. }
  86. }
  87. /** 超过 64 KiB 的事件用二分法截断 message,保证元数据仍能落盘 */
  88. function serializeEvent(event: AgentEvent): string {
  89. const serialized = JSON.stringify(event);
  90. if (Buffer.byteLength(serialized, 'utf8') <= MAX_EVENT_BYTES) return serialized;
  91. const marker = '\n…(事件日志已按 64 KiB 截断)';
  92. const codePoints = Array.from(event.message);
  93. let lower = 0;
  94. let upper = codePoints.length;
  95. let best = JSON.stringify({ ...event, message: marker });
  96. if (Buffer.byteLength(best, 'utf8') > MAX_EVENT_BYTES) {
  97. throw new Error('事件元数据超过 64 KiB 落盘上限');
  98. }
  99. while (lower <= upper) {
  100. const middle = Math.floor((lower + upper) / 2);
  101. const candidate = JSON.stringify({
  102. ...event,
  103. message: `${codePoints.slice(0, middle).join('')}${marker}`,
  104. });
  105. if (Buffer.byteLength(candidate, 'utf8') <= MAX_EVENT_BYTES) {
  106. best = candidate;
  107. lower = middle + 1;
  108. } else {
  109. upper = middle - 1;
  110. }
  111. }
  112. return best;
  113. }
  114. async function rotateIfNeeded(path: string): Promise<void> {
  115. try {
  116. if ((await stat(path)).size <= MAX_LOG_BYTES) return;
  117. const previous = `${path}.previous`;
  118. await rm(previous, { force: true });
  119. await rename(path, previous);
  120. } catch (error) {
  121. if (asNodeError(error).code !== 'ENOENT') throw error;
  122. }
  123. }
  124. function validateThreadId(value: string): void {
  125. if (!/^[a-zA-Z0-9_-]{1,128}$/u.test(value)) throw new Error('非法的 threadId');
  126. }
  127. function asNodeError(value: unknown): NodeJS.ErrnoException {
  128. return value as NodeJS.ErrnoException;
  129. }