Files
Jyotisha/frontend/src/lib/stream-text-response.ts
T

255 lines
8.0 KiB
TypeScript

type StreamHooks = {
readonly onFirstOutput?: () => Promise<void>;
readonly onComplete?: (output: string) => Promise<void>;
readonly onError?: (error: unknown, emitted: boolean, output: string) => Promise<void>;
readonly onCancel?: (emitted: boolean) => Promise<void>;
};
type StreamTextResponseOptions = StreamHooks & {
readonly mode: "engine" | "mastra";
readonly requestId: string;
readonly headers?: Record<string, string>;
readonly transformText?: (text: string) => string;
readonly continueAfterDisconnect?: boolean;
};
const hiddenBlockOpeners = [
"<!--AYANAM_SUGGESTIONS:",
"<!--AYANAM_TITLE:",
] as const;
function longestOpenerPrefixSuffix(value: string) {
const maximum = Math.min(
value.length,
Math.max(...hiddenBlockOpeners.map((opener) => opener.length - 1)),
);
for (let length = maximum; length > 0; length -= 1) {
const suffix = value.slice(-length);
if (hiddenBlockOpeners.some((opener) => opener.startsWith(suffix))) return length;
}
return 0;
}
/** Sends only visible prose through the output guard and preserves metadata bytes. */
export function createVisibleTextTransformer(transform: (text: string) => string) {
let rawBuffer = "";
let visibleBuffer = "";
let hiddenBuffer = "";
let hiddenBlocks: Array<{ readonly offset: number; readonly text: string }> = [];
let hidden = false;
function parse(value: string, final: boolean) {
rawBuffer += value;
while (rawBuffer) {
if (hidden) {
const closeIndex = rawBuffer.indexOf("-->");
if (closeIndex < 0) {
if (final) {
hiddenBlocks.push({
offset: visibleBuffer.length,
text: hiddenBuffer + rawBuffer,
});
hiddenBuffer = "";
rawBuffer = "";
hidden = false;
} else {
const retainedLength = Math.min(2, rawBuffer.length);
hiddenBuffer += rawBuffer.slice(0, rawBuffer.length - retainedLength);
rawBuffer = rawBuffer.slice(rawBuffer.length - retainedLength);
}
break;
}
hiddenBuffer += rawBuffer.slice(0, closeIndex + 3);
rawBuffer = rawBuffer.slice(closeIndex + 3);
hiddenBlocks.push({ offset: visibleBuffer.length, text: hiddenBuffer });
hiddenBuffer = "";
hidden = false;
continue;
}
const openerIndex = hiddenBlockOpeners.reduce<number>((earliest, opener) => {
const index = rawBuffer.indexOf(opener);
return index >= 0 && (earliest < 0 || index < earliest) ? index : earliest;
}, -1);
if (openerIndex >= 0) {
visibleBuffer += rawBuffer.slice(0, openerIndex);
rawBuffer = rawBuffer.slice(openerIndex);
hiddenBuffer = "";
hidden = true;
continue;
}
if (final) {
visibleBuffer += rawBuffer;
rawBuffer = "";
break;
}
const retainedLength = longestOpenerPrefixSuffix(rawBuffer);
const visibleLength = rawBuffer.length - retainedLength;
if (visibleLength > 0) visibleBuffer += rawBuffer.slice(0, visibleLength);
rawBuffer = rawBuffer.slice(visibleLength);
break;
}
}
function renderVisiblePrefix(length: number) {
if (length === 0) return "";
const visible = visibleBuffer.slice(0, length);
const included = hiddenBlocks.filter((block) => block.offset <= length);
const remaining = hiddenBlocks
.filter((block) => block.offset > length)
.map((block) => ({ ...block, offset: block.offset - length }));
const transformed = transform(visible);
let output = "";
if (transformed === visible) {
let start = 0;
for (const block of included) {
output += visible.slice(start, block.offset) + block.text;
start = block.offset;
}
output += visible.slice(start);
} else {
// A refusal may replace the whole sentence, so an in-sentence byte offset
// no longer has meaning. Keep metadata exact and in order after the safe
// visible replacement; the frontend parser accepts metadata at any point.
output = transformed + included.map((block) => block.text).join("");
}
visibleBuffer = visibleBuffer.slice(length);
hiddenBlocks = remaining;
return output;
}
function lastCompleteClauseBoundary() {
let boundary = 0;
for (const match of visibleBuffer.matchAll(/[。!?.!?\n]+/gu)) {
boundary = (match.index ?? 0) + match[0].length;
}
return boundary;
}
function consume(value: string, final: boolean) {
parse(value, final);
if (final) {
const output = renderVisiblePrefix(visibleBuffer.length);
if (hiddenBlocks.length === 0) return output;
const metadata = hiddenBlocks.map((block) => block.text).join("");
hiddenBlocks = [];
return output + metadata;
}
return renderVisiblePrefix(lastCompleteClauseBoundary());
}
return Object.freeze({
push: (value: string) => consume(value, false),
finish: (value: string) => consume(value, true),
});
}
export function streamTextResponse(
stream: AsyncIterable<string>,
options: StreamTextResponseOptions,
) {
const iterator = stream[Symbol.asyncIterator]();
const encoder = new TextEncoder();
const visibleTransformer = options.transformText
? createVisibleTextTransformer(options.transformText)
: null;
let settled = false;
let cancellationStarted = false;
let disconnected = false;
let emitted = false;
let fullOutput = "";
let firstOutputSettlementStarted = false;
function startFirstOutputSettlement(value: string) {
if (!/\S/.test(value) || firstOutputSettlementStarted) return undefined;
firstOutputSettlementStarted = true;
return options.onFirstOutput?.();
}
async function output(
controller: ReadableStreamDefaultController<Uint8Array> | undefined,
value: string,
) {
if (!value) return;
if (disconnected || !controller) {
fullOutput += value;
return;
}
const firstOutputSettlement = startFirstOutputSettlement(value);
controller.enqueue(encoder.encode(value));
fullOutput += value;
if (/\S/.test(value)) emitted = true;
await firstOutputSettlement;
}
async function consume(
controller: ReadableStreamDefaultController<Uint8Array> | undefined,
) {
try {
while (true) {
const { done, value } = await iterator.next();
if (settled) return;
if (done) {
await output(controller, visibleTransformer ? visibleTransformer.finish("") : "");
if (settled) return;
settled = true;
if (!/\S/.test(fullOutput)) {
const error = new Error("empty_stream");
await options.onError?.(error, false, fullOutput);
if (!disconnected) controller?.error(error);
return;
}
await options.onComplete?.(fullOutput);
if (!disconnected) controller?.close();
return;
}
const transformed = visibleTransformer
? visibleTransformer.push(value)
: value;
await output(controller, transformed);
if (settled) return;
}
} catch (error) {
if (cancellationStarted) return;
if (!settled) {
settled = true;
await options.onError?.(error, emitted, fullOutput);
}
if (!disconnected) controller?.error(error);
}
}
const body = new ReadableStream<Uint8Array>({
start(controller) {
void consume(controller).catch(() => {});
},
async cancel() {
if (settled) return;
if (options.continueAfterDisconnect) {
disconnected = true;
return;
}
settled = true;
cancellationStarted = true;
try {
await iterator.return?.();
} finally {
await options.onCancel?.(emitted);
}
},
});
return new Response(body, {
headers: {
"cache-control": "no-cache, no-transform",
"content-type": "text/plain; charset=utf-8",
"x-accel-buffering": "no",
"x-ayanam-mode": options.mode,
"x-ayanam-request-id": options.requestId,
...options.headers,
},
});
}