diff --git a/extensions/imessage/src/client.ts b/extensions/imessage/src/client.ts index 7d13b25d8698..ff55c822a3e2 100644 --- a/extensions/imessage/src/client.ts +++ b/extensions/imessage/src/client.ts @@ -1,6 +1,6 @@ // Imessage plugin module implements client behavior. import { type ChildProcessWithoutNullStreams, spawn } from "node:child_process"; -import { createInterface, type Interface } from "node:readline"; +import { StringDecoder } from "node:string_decoder"; import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime"; import type { RuntimeEnv } from "openclaw/plugin-sdk/runtime-env"; import { normalizeLowercaseStringOrEmpty } from "openclaw/plugin-sdk/string-coerce-runtime"; @@ -68,7 +68,8 @@ export class IMessageRpcClient { private readonly closed: Promise; private closedResolve: (() => void) | null = null; private child: ChildProcessWithoutNullStreams | null = null; - private reader: Interface | null = null; + private stdoutBuffer = ""; + private readonly stdoutDecoder = new StringDecoder("utf8"); private nextId = 1; private publicProcessError: string | null = null; @@ -97,14 +98,12 @@ export class IMessageRpcClient { stdio: ["pipe", "pipe", "pipe"], }); this.child = child; - this.reader = createInterface({ input: child.stdout }); - this.reader.on("line", (line) => { - const trimmed = line.trim(); - if (!trimmed) { + child.stdout.on("data", (chunk) => { + if (this.child !== child) { return; } - this.handleLine(trimmed); + this.handleStdoutChunk(chunk); }); child.stderr?.on("data", (chunk) => { @@ -131,6 +130,9 @@ export class IMessageRpcClient { }); child.on("close", (code, signal) => { + if (this.child === child) { + this.flushStdoutBuffer(); + } this.failAll(this.buildCloseError(code, signal)); this.closedResolve?.(); }); @@ -140,8 +142,8 @@ export class IMessageRpcClient { if (!this.child) { return; } - this.reader?.close(); - this.reader = null; + this.stdoutBuffer = ""; + this.stdoutDecoder.end(); this.child.stdin?.end(); const child = this.child; this.child = null; @@ -215,6 +217,40 @@ export class IMessageRpcClient { return await response; } + private handleStdoutChunk(chunk: Buffer | string) { + const text = typeof chunk === "string" ? chunk : this.stdoutDecoder.write(chunk); + this.stdoutBuffer += text; + + let newlineIndex = this.stdoutBuffer.indexOf("\n"); + while (newlineIndex !== -1) { + const line = this.stdoutBuffer.slice(0, newlineIndex); + this.stdoutBuffer = this.stdoutBuffer.slice(newlineIndex + 1); + this.handleStdoutLine(line); + newlineIndex = this.stdoutBuffer.indexOf("\n"); + } + } + + private flushStdoutBuffer() { + const tail = this.stdoutDecoder.end(); + if (tail) { + this.stdoutBuffer += tail; + } + if (!this.stdoutBuffer) { + return; + } + const line = this.stdoutBuffer; + this.stdoutBuffer = ""; + this.handleStdoutLine(line); + } + + private handleStdoutLine(line: string) { + const trimmed = line.trim(); + if (!trimmed) { + return; + } + this.handleLine(trimmed); + } + private handleLine(line: string) { let parsed: IMessageRpcResponse; try { diff --git a/extensions/imessage/src/status.test.ts b/extensions/imessage/src/status.test.ts index 013cbb5e917d..330e9977c1c6 100644 --- a/extensions/imessage/src/status.test.ts +++ b/extensions/imessage/src/status.test.ts @@ -1,4 +1,6 @@ // Imessage tests cover status plugin behavior. +import { EventEmitter } from "node:events"; +import { PassThrough } from "node:stream"; import { createPluginSetupWizardStatus } from "openclaw/plugin-sdk/plugin-test-runtime"; import * as processRuntime from "openclaw/plugin-sdk/process-runtime"; import * as setupRuntime from "openclaw/plugin-sdk/setup"; @@ -20,6 +22,26 @@ const getIMessageSetupStatus = createPluginSetupWizardStatus({ const spawnMock = vi.hoisted(() => vi.fn()); +function createMockChildProcess() { + const child = new EventEmitter() as EventEmitter & { + stdin: PassThrough; + stdout: PassThrough; + stderr: PassThrough; + killed: boolean; + kill: (signal?: string) => boolean; + }; + child.stdin = new PassThrough(); + child.stdout = new PassThrough(); + child.stderr = new PassThrough(); + child.killed = false; + child.kill = (signal?: string) => { + child.killed = true; + child.emit("close", 0, signal ?? null); + return true; + }; + return child; +} + vi.mock("node:child_process", async () => { const actual = await vi.importActual("node:child_process"); return { @@ -67,6 +89,96 @@ describe("createIMessageRpcClient", () => { expect(internals.buildCloseError(1, null).message).toBe(PUBLIC_IMESSAGE_FULL_DISK_ACCESS_ERROR); }); + + it.each([ + ["U+2028", "\u2028"], + ["U+2029", "\u2029"], + ])( + "frames stdout on LF only so raw %s inside JSON strings stays intact", + async (_, separator) => { + const { IMessageRpcClient } = await import("./client.js"); + const client = new IMessageRpcClient(); + const internals = client as unknown as { + handleStdoutChunk: (chunk: Buffer | string) => void; + pending: Map< + string, + { + resolve: (value: unknown) => void; + reject: (error: Error) => void; + } + >; + }; + const result = new Promise((resolve, reject) => { + internals.pending.set("1", { resolve, reject }); + }); + const text = `line one${separator}line two`; + const payload = `${JSON.stringify({ + jsonrpc: "2.0", + id: 1, + result: { messages: [{ text }] }, + })}\n`; + const bytes = Buffer.from(payload, "utf8"); + const separatorIndex = bytes.indexOf(Buffer.from(separator, "utf8")); + + internals.handleStdoutChunk(bytes.subarray(0, separatorIndex + 1)); + internals.handleStdoutChunk(bytes.subarray(separatorIndex + 1)); + + await expect(result).resolves.toEqual({ + messages: [{ text }], + }); + }, + ); + + it("handles multiple LF-delimited stdout responses in one chunk", async () => { + const { IMessageRpcClient } = await import("./client.js"); + const client = new IMessageRpcClient(); + const internals = client as unknown as { + handleStdoutChunk: (chunk: Buffer | string) => void; + pending: Map< + string, + { + resolve: (value: unknown) => void; + reject: (error: Error) => void; + } + >; + }; + const first = new Promise((resolve, reject) => { + internals.pending.set("1", { resolve, reject }); + }); + const second = new Promise((resolve, reject) => { + internals.pending.set("2", { resolve, reject }); + }); + + internals.handleStdoutChunk( + `${JSON.stringify({ jsonrpc: "2.0", id: 1, result: { ok: "first" } })}\n${JSON.stringify({ + jsonrpc: "2.0", + id: 2, + result: { ok: "second" }, + })}\n`, + ); + + await expect(first).resolves.toEqual({ ok: "first" }); + await expect(second).resolves.toEqual({ ok: "second" }); + }); + + it("ignores stdout from a stale child after stop so late notifications cannot leak (#89830)", async () => { + vi.stubEnv("VITEST", ""); + vi.stubEnv("NODE_ENV", ""); + const child = createMockChildProcess(); + spawnMock.mockReturnValue(child); + const onNotification = vi.fn(); + const { IMessageRpcClient } = await import("./client.js"); + const client = new IMessageRpcClient({ onNotification }); + + await client.start(); + await client.stop(); + + // A not-yet-exited imsg child emits a complete notification after stop(). + // The `this.child !== child` guard must drop it before handleStdoutChunk. + child.stdout.write('{"jsonrpc":"2.0","method":"messages.changed","params":{}}\n'); + + expect(onNotification).not.toHaveBeenCalled(); + }); }); describe("imessage setup status", () => {