From df674eb7160435695ae98deafc76e1a60e12262c Mon Sep 17 00:00:00 2001 From: cyq <15000851237@163.com> Date: Fri, 21 Aug 2026 11:56:21 +0800 Subject: [PATCH] fix(channel): preserve UTF-8 across incremental event reads --- .../cli/src/commands/channel/store/watch.ts | 12 ++- .../cli/test/commands/channel-watch.test.ts | 80 +++++++++++++++++++ .../core/src/channel/internal/store/watch.ts | 12 ++- packages/core/test/channel/watch.test.ts | 64 +++++++++++++++ 4 files changed, 164 insertions(+), 4 deletions(-) create mode 100644 packages/cli/test/commands/channel-watch.test.ts create mode 100644 packages/core/test/channel/watch.test.ts diff --git a/packages/cli/src/commands/channel/store/watch.ts b/packages/cli/src/commands/channel/store/watch.ts index 62f4f3fb2..25e44e374 100644 --- a/packages/cli/src/commands/channel/store/watch.ts +++ b/packages/cli/src/commands/channel/store/watch.ts @@ -1,4 +1,5 @@ import fs from "node:fs"; +import { StringDecoder } from "node:string_decoder"; import { eventsPath, channelDir } from "./paths.js"; import type { ChannelEvent } from "./events.js"; @@ -9,6 +10,7 @@ export type WatchFilter = ChannelEventFilter; interface ReadProgress { byteOffset: number; carry: string; + decoder: StringDecoder; } async function readNewEvents( @@ -20,6 +22,7 @@ async function readNewEvents( // the file reappears we'll re-scan it from byte 0. state.byteOffset = 0; state.carry = ""; + state.decoder = new StringDecoder("utf8"); return []; } const stat = await fs.promises.stat(filePath); @@ -28,6 +31,7 @@ async function readNewEvents( // post-truncate events aren't lost forever. state.byteOffset = 0; state.carry = ""; + state.decoder = new StringDecoder("utf8"); } if (stat.size <= state.byteOffset) return []; @@ -37,7 +41,7 @@ async function readNewEvents( const buf = Buffer.alloc(length); await fh.read(buf, 0, length, state.byteOffset); state.byteOffset = stat.size; - const text = state.carry + buf.toString("utf-8"); + const text = state.carry + state.decoder.write(buf); const lines = text.split("\n"); state.carry = lines.pop() ?? ""; const events: ChannelEvent[] = []; @@ -102,7 +106,11 @@ export async function* watchEvents( initialOffset = 0; } } - const state: ReadProgress = { byteOffset: initialOffset, carry: "" }; + const state: ReadProgress = { + byteOffset: initialOffset, + carry: "", + decoder: new StringDecoder("utf8"), + }; const sinceSeq = opts.sinceSeq; let resolveNext: (() => void) | null = null; diff --git a/packages/cli/test/commands/channel-watch.test.ts b/packages/cli/test/commands/channel-watch.test.ts new file mode 100644 index 000000000..af656696a --- /dev/null +++ b/packages/cli/test/commands/channel-watch.test.ts @@ -0,0 +1,80 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +import { createChannel } from "../../src/commands/channel/create.js"; +import { eventsPath } from "../../src/commands/channel/store/paths.js"; +import { watchEvents } from "../../src/commands/channel/store/watch.js"; + +describe("watchEvents", () => { + let tmpDir: string; + let projectDir: string; + let oldRoot: string | undefined; + let oldProject: string | undefined; + + beforeEach(() => { + tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), "trellis-watch-test-")); + projectDir = path.join(tmpDir, "project"); + fs.mkdirSync(projectDir); + oldRoot = process.env.TRELLIS_CHANNEL_ROOT; + oldProject = process.env.TRELLIS_CHANNEL_PROJECT; + process.env.TRELLIS_CHANNEL_ROOT = path.join(tmpDir, "channels"); + delete process.env.TRELLIS_CHANNEL_PROJECT; + vi.spyOn(process, "cwd").mockReturnValue(projectDir); + vi.spyOn(console, "log").mockImplementation(() => undefined); + }); + + afterEach(() => { + vi.restoreAllMocks(); + if (oldRoot === undefined) delete process.env.TRELLIS_CHANNEL_ROOT; + else process.env.TRELLIS_CHANNEL_ROOT = oldRoot; + if (oldProject === undefined) delete process.env.TRELLIS_CHANNEL_PROJECT; + else process.env.TRELLIS_CHANNEL_PROJECT = oldProject; + fs.rmSync(tmpDir, { recursive: true, force: true }); + }); + + it("preserves UTF-8 characters split across incremental reads", async () => { + const channel = "utf8-boundary"; + await createChannel(channel, { by: "main" }); + + const messageLine = Buffer.from( + `${JSON.stringify({ + seq: 2, + ts: "2026-08-18T00:00:00.000Z", + kind: "message", + by: "worker", + text: "中", + })}\n`, + "utf8", + ); + const characterStart = messageLine.indexOf(Buffer.from("中", "utf8")); + expect(characterStart).toBeGreaterThanOrEqual(0); + + const file = eventsPath(channel); + fs.appendFileSync(file, messageLine.subarray(0, characterStart + 1)); + + const abortController = new AbortController(); + const events = watchEvents( + channel, + {}, + { fromStart: true, signal: abortController.signal }, + ); + + try { + const first = await events.next(); + expect(first.value).toMatchObject({ kind: "create" }); + + const secondEvent = events.next(); + fs.appendFileSync(file, messageLine.subarray(characterStart + 1)); + + await expect(secondEvent).resolves.toMatchObject({ + done: false, + value: { kind: "message", text: "中" }, + }); + } finally { + abortController.abort(); + await events.return(undefined); + } + }); +}); diff --git a/packages/core/src/channel/internal/store/watch.ts b/packages/core/src/channel/internal/store/watch.ts index 851845669..84ac91ce6 100644 --- a/packages/core/src/channel/internal/store/watch.ts +++ b/packages/core/src/channel/internal/store/watch.ts @@ -1,4 +1,5 @@ import fs from "node:fs"; +import { StringDecoder } from "node:string_decoder"; import type { ChannelEvent } from "./events.js"; import { matchesEventFilter, type ChannelEventFilter } from "./filter.js"; @@ -9,6 +10,7 @@ export type WatchFilter = ChannelEventFilter; interface ReadProgress { byteOffset: number; carry: string; + decoder: StringDecoder; } async function readNewEvents( @@ -18,12 +20,14 @@ async function readNewEvents( if (!fs.existsSync(filePath)) { state.byteOffset = 0; state.carry = ""; + state.decoder = new StringDecoder("utf8"); return []; } const stat = await fs.promises.stat(filePath); if (stat.size < state.byteOffset) { state.byteOffset = 0; state.carry = ""; + state.decoder = new StringDecoder("utf8"); } if (stat.size <= state.byteOffset) return []; @@ -33,7 +37,7 @@ async function readNewEvents( const buf = Buffer.alloc(length); await fh.read(buf, 0, length, state.byteOffset); state.byteOffset = stat.size; - const text = state.carry + buf.toString("utf-8"); + const text = state.carry + state.decoder.write(buf); const lines = text.split("\n"); state.carry = lines.pop() ?? ""; const events: ChannelEvent[] = []; @@ -79,7 +83,11 @@ export async function* watchEvents( initialOffset = 0; } } - const state: ReadProgress = { byteOffset: initialOffset, carry: "" }; + const state: ReadProgress = { + byteOffset: initialOffset, + carry: "", + decoder: new StringDecoder("utf8"), + }; const sinceSeq = opts.sinceSeq; let resolveNext: (() => void) | null = null; diff --git a/packages/core/test/channel/watch.test.ts b/packages/core/test/channel/watch.test.ts new file mode 100644 index 000000000..fbff0e5cf --- /dev/null +++ b/packages/core/test/channel/watch.test.ts @@ -0,0 +1,64 @@ +import fs from "node:fs"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +import { createChannel, watchChannelEvents } from "../../src/channel/index.js"; +import { eventsPath } from "../../src/channel/internal/store/paths.js"; +import { setupChannelTmp, type TmpEnv } from "./setup.js"; + +describe("watchChannelEvents", () => { + let env: TmpEnv; + + beforeEach(() => { + env = setupChannelTmp(); + vi.spyOn(process, "cwd").mockReturnValue(env.projectDir); + }); + + afterEach(() => { + vi.restoreAllMocks(); + env.cleanup(); + }); + + it("preserves UTF-8 characters split across incremental reads", async () => { + const channel = "utf8-boundary"; + await createChannel({ channel, by: "main" }); + + const messageLine = Buffer.from( + `${JSON.stringify({ + seq: 2, + ts: "2026-08-18T00:00:00.000Z", + kind: "message", + by: "worker", + text: "中", + })}\n`, + "utf8", + ); + const characterStart = messageLine.indexOf(Buffer.from("中", "utf8")); + expect(characterStart).toBeGreaterThanOrEqual(0); + + const file = eventsPath(channel); + fs.appendFileSync(file, messageLine.subarray(0, characterStart + 1)); + + const abortController = new AbortController(); + const events = watchChannelEvents({ + channel, + fromStart: true, + signal: abortController.signal, + }); + + try { + const first = await events.next(); + expect(first.value).toMatchObject({ kind: "create" }); + + const secondEvent = events.next(); + fs.appendFileSync(file, messageLine.subarray(characterStart + 1)); + + await expect(secondEvent).resolves.toMatchObject({ + done: false, + value: { kind: "message", text: "中" }, + }); + } finally { + abortController.abort(); + await events.return(undefined); + } + }); +});