Skip to content

Commit c10307c

Browse files
committed
fix(line): add chunk-idle timeout to inbound media download
1 parent 92fb79e commit c10307c

2 files changed

Lines changed: 222 additions & 7 deletions

File tree

extensions/line/src/download.test.ts

Lines changed: 109 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@ vi.mock("openclaw/plugin-sdk/media-store", () => ({
3232
}));
3333

3434
let downloadLineMedia: typeof import("./download.js").downloadLineMedia;
35+
let LineMediaDownloadTimeoutError: typeof import("./download.js").LineMediaDownloadTimeoutError;
36+
let LINE_DOWNLOAD_IDLE_TIMEOUT_MS: typeof import("./download.js").LINE_DOWNLOAD_IDLE_TIMEOUT_MS;
3537

3638
async function* chunks(parts: Buffer[]): AsyncGenerator<Buffer> {
3739
for (const part of parts) {
@@ -59,7 +61,8 @@ function detectMockContentType(buffer: Buffer, contentType?: string): string | u
5961

6062
describe("downloadLineMedia", () => {
6163
beforeAll(async () => {
62-
({ downloadLineMedia } = await import("./download.js"));
64+
({ downloadLineMedia, LineMediaDownloadTimeoutError, LINE_DOWNLOAD_IDLE_TIMEOUT_MS } =
65+
await import("./download.js"));
6366
});
6467

6568
afterAll(() => {
@@ -161,4 +164,109 @@ describe("downloadLineMedia", () => {
161164

162165
await expect(downloadLineMedia("mid-bad", "token")).rejects.toThrow(/Media exceeds/i);
163166
});
167+
168+
describe("chunk-idle timeout", () => {
169+
function neverYieldingStream(): AsyncIterable<Buffer> {
170+
return {
171+
[Symbol.asyncIterator]() {
172+
return {
173+
next(): Promise<IteratorResult<Buffer>> {
174+
return new Promise<IteratorResult<Buffer>>(() => {});
175+
},
176+
};
177+
},
178+
};
179+
}
180+
181+
function delayedStream(payload: Buffer, delayMs: number): AsyncIterable<Buffer> {
182+
let yielded = false;
183+
return {
184+
[Symbol.asyncIterator]() {
185+
return {
186+
async next(): Promise<IteratorResult<Buffer>> {
187+
if (yielded) {
188+
return { value: undefined as unknown as Buffer, done: true };
189+
}
190+
await new Promise((r) => setTimeout(r, delayMs));
191+
yielded = true;
192+
return { value: payload, done: false };
193+
},
194+
};
195+
},
196+
};
197+
}
198+
199+
it("rejects with LineMediaDownloadTimeoutError when the stream stalls past chunkTimeoutMs", async () => {
200+
getMessageContentMock.mockResolvedValueOnce(neverYieldingStream());
201+
const promise = downloadLineMedia("mid-stall", "token", 1024 * 1024, {
202+
chunkTimeoutMs: 50,
203+
});
204+
await expect(promise).rejects.toBeInstanceOf(LineMediaDownloadTimeoutError);
205+
await expect(promise.catch((e) => e.chunkTimeoutMs)).resolves.toBe(50);
206+
});
207+
208+
it("rejects when getMessageContent headers never arrive", async () => {
209+
getMessageContentMock.mockReturnValueOnce(new Promise(() => {}));
210+
const promise = downloadLineMedia("mid-stall-headers", "token", 1024 * 1024, {
211+
chunkTimeoutMs: 50,
212+
});
213+
await expect(promise).rejects.toBeInstanceOf(LineMediaDownloadTimeoutError);
214+
});
215+
216+
it("does not reject when chunks arrive within chunkTimeoutMs", async () => {
217+
const jpeg = Buffer.from([0xff, 0xd8, 0xff, 0x00]);
218+
getMessageContentMock.mockResolvedValueOnce(delayedStream(jpeg, 10));
219+
const result = await downloadLineMedia("mid-slow-but-progressing", "token", 1024 * 1024, {
220+
chunkTimeoutMs: 500,
221+
});
222+
expect(result.contentType).toBe("image/jpeg");
223+
});
224+
225+
it("exposes LINE_DOWNLOAD_IDLE_TIMEOUT_MS = 30s aligned with TELEGRAM_DOWNLOAD_IDLE_TIMEOUT_MS", () => {
226+
expect(LINE_DOWNLOAD_IDLE_TIMEOUT_MS).toBe(30_000);
227+
});
228+
229+
it("calls iterator.return() exactly once on timeout so the upstream Readable is destroyed", async () => {
230+
const returnSpy = vi.fn(async () => ({ value: undefined as unknown as Buffer, done: true }));
231+
const stream: AsyncIterable<Buffer> = {
232+
[Symbol.asyncIterator]() {
233+
return {
234+
next(): Promise<IteratorResult<Buffer>> {
235+
return new Promise<IteratorResult<Buffer>>(() => {});
236+
},
237+
return: returnSpy as () => Promise<IteratorResult<Buffer>>,
238+
};
239+
},
240+
};
241+
getMessageContentMock.mockResolvedValueOnce(stream);
242+
await expect(
243+
downloadLineMedia("mid-return-spy", "token", 1024 * 1024, { chunkTimeoutMs: 50 }),
244+
).rejects.toBeInstanceOf(LineMediaDownloadTimeoutError);
245+
expect(returnSpy).toHaveBeenCalledTimes(1);
246+
});
247+
248+
it("rejects when a second chunk stalls after a successful first chunk (partial-then-stall)", async () => {
249+
const jpegPart = Buffer.from([0xff, 0xd8, 0xff, 0x00]);
250+
let yieldedFirst = false;
251+
const stream: AsyncIterable<Buffer> = {
252+
[Symbol.asyncIterator]() {
253+
return {
254+
async next(): Promise<IteratorResult<Buffer>> {
255+
if (!yieldedFirst) {
256+
yieldedFirst = true;
257+
return { value: jpegPart, done: false };
258+
}
259+
return new Promise<IteratorResult<Buffer>>(() => {});
260+
},
261+
};
262+
},
263+
};
264+
getMessageContentMock.mockResolvedValueOnce(stream);
265+
const promise = downloadLineMedia("mid-partial-stall", "token", 1024 * 1024, {
266+
chunkTimeoutMs: 50,
267+
});
268+
await expect(promise).rejects.toBeInstanceOf(LineMediaDownloadTimeoutError);
269+
await expect(promise.catch((e) => e.chunkTimeoutMs)).resolves.toBe(50);
270+
});
271+
});
164272
});

extensions/line/src/download.ts

Lines changed: 113 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -8,22 +8,129 @@ interface DownloadResult {
88
size: number;
99
}
1010

11+
/**
12+
* Default per-chunk idle timeout for LINE inbound media downloads.
13+
*
14+
* Matches the Telegram `TELEGRAM_DOWNLOAD_IDLE_TIMEOUT_MS = 30_000` ceiling so
15+
* a stalled LINE content stream cannot block inbound dispatch indefinitely.
16+
* Idle (not overall) so legitimate slow-but-progressing transfers continue.
17+
*/
18+
export const LINE_DOWNLOAD_IDLE_TIMEOUT_MS = 30_000;
19+
20+
export class LineMediaDownloadTimeoutError extends Error {
21+
readonly chunkTimeoutMs: number;
22+
constructor(chunkTimeoutMs: number) {
23+
super(`LINE media download stalled: no data received for ${chunkTimeoutMs}ms`);
24+
this.name = "LineMediaDownloadTimeoutError";
25+
this.chunkTimeoutMs = chunkTimeoutMs;
26+
}
27+
}
28+
29+
// Wraps an AsyncIterable so each `next()` is bounded by `chunkTimeoutMs`. On
30+
// timeout we propagate `LineMediaDownloadTimeoutError` and call `return()` on
31+
// the upstream iterator. For Node `Readable` produced by `@line/bot-sdk`
32+
// `convertResponseToReadable`, that triggers `Readable.destroy()`, which stops
33+
// further `read()` calls. The underlying `fetch` itself is NOT cancelled —
34+
// `@line/bot-sdk`'s `HTTPFetchClient.get` does not accept `AbortSignal`, so the
35+
// in-flight HTTP body may continue draining the TCP buffer until the OS
36+
// keep-alive timer fires. The caller (and any concurrent inbound dispatch) is
37+
// unblocked immediately; the orphaned fetch is bounded by transport, not by
38+
// this code.
39+
function withChunkIdleTimeout<T>(
40+
source: AsyncIterable<T>,
41+
chunkTimeoutMs: number,
42+
): AsyncIterable<T> {
43+
return {
44+
async *[Symbol.asyncIterator]() {
45+
const iterator = source[Symbol.asyncIterator]();
46+
try {
47+
while (true) {
48+
const nextPromise = iterator.next();
49+
let timeoutHandle: ReturnType<typeof setTimeout> | undefined;
50+
const timeoutPromise = new Promise<never>((_, reject) => {
51+
timeoutHandle = setTimeout(
52+
() => reject(new LineMediaDownloadTimeoutError(chunkTimeoutMs)),
53+
chunkTimeoutMs,
54+
);
55+
});
56+
let result: IteratorResult<T>;
57+
try {
58+
result = await Promise.race([nextPromise, timeoutPromise]);
59+
} catch (err) {
60+
// The timeout won the race; `nextPromise` is still pending and may
61+
// later reject (e.g., the upstream Readable emits an error after
62+
// `destroy()`). Without this, that rejection escapes as
63+
// `unhandledRejection` and is fatal on Node 22+.
64+
nextPromise.then(
65+
() => undefined,
66+
() => undefined,
67+
);
68+
throw err;
69+
} finally {
70+
if (timeoutHandle !== undefined) {
71+
clearTimeout(timeoutHandle);
72+
}
73+
}
74+
if (result.done) {
75+
return;
76+
}
77+
yield result.value;
78+
}
79+
} finally {
80+
if (typeof iterator.return === "function") {
81+
await iterator.return().catch(() => undefined);
82+
}
83+
}
84+
},
85+
};
86+
}
87+
1188
export async function downloadLineMedia(
1289
messageId: string,
1390
channelAccessToken: string,
1491
maxBytes = 10 * 1024 * 1024,
92+
options?: { chunkTimeoutMs?: number },
1593
): Promise<DownloadResult> {
94+
const chunkTimeoutMs = options?.chunkTimeoutMs ?? LINE_DOWNLOAD_IDLE_TIMEOUT_MS;
1695
const client = new messagingApi.MessagingApiBlobClient({
1796
channelAccessToken,
1897
});
1998

20-
const response = await client.getMessageContent(messageId);
21-
const saved = await saveMediaStream(
22-
response as AsyncIterable<Buffer>,
23-
undefined,
24-
"inbound",
25-
maxBytes,
99+
// Race the headers-level fetch against the same idle ceiling. `@line/bot-sdk`
100+
// currently builds its fetch without `AbortSignal`, so a server that never
101+
// returns headers cannot otherwise be capped here. The fetch itself keeps
102+
// running on timeout (see `withChunkIdleTimeout` note above for the same
103+
// limitation on the body phase).
104+
let responseTimeoutHandle: ReturnType<typeof setTimeout> | undefined;
105+
const responseTimeoutPromise = new Promise<never>((_, reject) => {
106+
responseTimeoutHandle = setTimeout(
107+
() => reject(new LineMediaDownloadTimeoutError(chunkTimeoutMs)),
108+
chunkTimeoutMs,
109+
);
110+
});
111+
const responsePromise = client.getMessageContent(messageId);
112+
let response: Awaited<ReturnType<typeof client.getMessageContent>>;
113+
try {
114+
response = await Promise.race([responsePromise, responseTimeoutPromise]);
115+
} catch (err) {
116+
// Same as `withChunkIdleTimeout`: suppress any later rejection from the
117+
// orphaned `responsePromise` so it does not become an `unhandledRejection`.
118+
responsePromise.then(
119+
() => undefined,
120+
() => undefined,
121+
);
122+
throw err;
123+
} finally {
124+
if (responseTimeoutHandle !== undefined) {
125+
clearTimeout(responseTimeoutHandle);
126+
}
127+
}
128+
129+
const guardedStream = withChunkIdleTimeout(
130+
response as unknown as AsyncIterable<Buffer>,
131+
chunkTimeoutMs,
26132
);
133+
const saved = await saveMediaStream(guardedStream, undefined, "inbound", maxBytes);
27134
logVerbose(`line: persisted media ${messageId} to ${saved.path} (${saved.size} bytes)`);
28135

29136
return {

0 commit comments

Comments
 (0)