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

255 lines
8.0 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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,
},
});
}