diff --git a/src/utils/s3/object-stream.ts b/src/utils/s3/object-stream.ts new file mode 100644 index 0000000..52af667 --- /dev/null +++ b/src/utils/s3/object-stream.ts @@ -0,0 +1,110 @@ +import { contentRange, type RangeParseResult } from './range'; + +export interface ObjectPartSource { + telegramFileId: string; + telegramUrl: string; + sizeBytes: number; + partNumber: number; +} + +export interface ObjectResponseInput { + reqId: string; + contentType: string; + etag: string; + lastModified: Date; + totalSize: number; + parts: ObjectPartSource[]; + range: RangeParseResult; +} + +interface PlannedPart { + part: ObjectPartSource; + relativeStart: number; + relativeEnd: number; +} + +const baseHeaders = (input: ObjectResponseInput, contentLength: number): Headers => { + const headers = new Headers({ + 'content-type': input.contentType, + 'content-length': String(contentLength), + etag: `"${input.etag}"`, + 'last-modified': input.lastModified.toUTCString(), + 'x-amz-request-id': input.reqId, + 'accept-ranges': 'bytes', + 'cache-control': 'public, max-age=31536000', + }); + return headers; +}; + +const planParts = (parts: ObjectPartSource[], start: number, end: number): PlannedPart[] => { + const planned: PlannedPart[] = []; + let offset = 0; + for (const part of parts) { + const partStart = offset; + const partEnd = offset + part.sizeBytes - 1; + offset += part.sizeBytes; + if (end < partStart || start > partEnd) continue; + planned.push({ + part, + relativeStart: Math.max(start, partStart) - partStart, + relativeEnd: Math.min(end, partEnd) - partStart, + }); + } + return planned; +}; + +const fetchPartBody = async (planned: PlannedPart): Promise> => { + const rangeHeader = `bytes=${planned.relativeStart}-${planned.relativeEnd}`; + const wantsWholePart = + planned.relativeStart === 0 && planned.relativeEnd === planned.part.sizeBytes - 1; + const res = await fetch( + planned.part.telegramUrl, + wantsWholePart ? undefined : { headers: { range: rangeHeader } }, + ); + if (!res.ok) throw new Error(`Telegram fetch failed: ${res.status}`); + if (wantsWholePart || res.status === 206) return res.body!; + + const bytes = new Uint8Array(await res.arrayBuffer()); + return new Response(bytes.slice(planned.relativeStart, planned.relativeEnd + 1)).body!; +}; + +const concatPartStreams = (plannedParts: PlannedPart[]): ReadableStream => + new ReadableStream({ + async start(controller) { + try { + for (const planned of plannedParts) { + const stream = await fetchPartBody(planned); + const reader = stream.getReader(); + while (true) { + const { value, done } = await reader.read(); + if (done) break; + if (value) controller.enqueue(value); + } + } + controller.close(); + } catch (error) { + controller.error(error); + } + }, + }); + +export const createGetObjectResponse = async (input: ObjectResponseInput): Promise => { + if (input.range.type === 'invalid') { + throw new Error('createGetObjectResponse received invalid range'); + } + + const start = input.range.type === 'valid' ? input.range.start : 0; + const end = input.range.type === 'valid' ? input.range.end : input.totalSize - 1; + const plannedParts = planParts(input.parts, start, end); + const contentLength = end >= start ? end - start + 1 : 0; + const headers = baseHeaders(input, contentLength); + + if (input.range.type === 'valid') { + headers.set('content-range', contentRange(start, end, input.totalSize)); + } + + return new Response(concatPartStreams(plannedParts), { + status: input.range.type === 'valid' ? 206 : 200, + headers, + }); +}; diff --git a/test/s3-object-stream.test.ts b/test/s3-object-stream.test.ts new file mode 100644 index 0000000..0ce58c3 --- /dev/null +++ b/test/s3-object-stream.test.ts @@ -0,0 +1,85 @@ +import { afterEach, describe, expect, it } from 'bun:test'; +import { createGetObjectResponse } from '../src/utils/s3/object-stream'; + +const originalFetch = globalThis.fetch; + +const streamText = (text: string) => new Response(text).body!; + +const installFetch = () => { + globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { + const range = new Headers(init?.headers).get('range'); + const url = String(_url); + const text = url.includes('part-1') ? 'hello ' : 'world'; + if (range === 'bytes=1-3') { + return new Response(text.slice(1, 4), { + status: 206, + headers: { 'content-range': `bytes 1-3/${text.length}`, 'content-length': '3' }, + }); + } + return new Response(streamText(text), { + status: 200, + headers: { 'content-length': String(text.length) }, + }); + }) as typeof fetch; +}; + +afterEach(() => { + globalThis.fetch = originalFetch; +}); + +describe('S3 object stream response builder', () => { + it('concatenates multiple Telegram part streams', async () => { + installFetch(); + const res = await createGetObjectResponse({ + reqId: 'req-1', + contentType: 'text/plain', + etag: 'etag123', + lastModified: new Date('2026-07-07T00:00:00Z'), + totalSize: 11, + parts: [ + { + telegramFileId: 'part-1', + telegramUrl: 'https://telegram.test/part-1', + sizeBytes: 6, + partNumber: 1, + }, + { + telegramFileId: 'part-2', + telegramUrl: 'https://telegram.test/part-2', + sizeBytes: 5, + partNumber: 2, + }, + ], + range: { type: 'none' }, + }); + + expect(res.status).toBe(200); + expect(res.headers.get('content-length')).toBe('11'); + expect(await res.text()).toBe('hello world'); + }); + + it('returns 206 with content-range for a single-part byte range', async () => { + installFetch(); + const res = await createGetObjectResponse({ + reqId: 'req-2', + contentType: 'text/plain', + etag: 'etag123', + lastModified: new Date('2026-07-07T00:00:00Z'), + totalSize: 6, + parts: [ + { + telegramFileId: 'part-1', + telegramUrl: 'https://telegram.test/part-1', + sizeBytes: 6, + partNumber: 1, + }, + ], + range: { type: 'valid', start: 1, end: 3 }, + }); + + expect(res.status).toBe(206); + expect(res.headers.get('content-range')).toBe('bytes 1-3/6'); + expect(res.headers.get('content-length')).toBe('3'); + expect(await res.text()).toBe('ell'); + }); +});