Files
AFFiNE/packages/backend/server/src/plugins/copilot/utils.ts
DarkSky 5a49d5cd24 fix(server): abort behavior in sse stream (#12211)
fix AI-121
fix AI-118

<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## Summary by CodeRabbit

- **Bug Fixes**
- Improved handling of connection closures and request abortion for
streaming and non-streaming chat endpoints, ensuring session data is
saved appropriately even if the connection is interrupted.
- **Refactor**
- Streamlined internal logic for managing request signals and connection
events, resulting in more robust and explicit session management during
streaming interactions.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
2025-07-04 06:07:45 +00:00

49 lines
1.2 KiB
TypeScript

import { Readable } from 'node:stream';
import type { Request } from 'express';
import { readBufferWithLimit } from '../../base';
import { MAX_EMBEDDABLE_SIZE } from './types';
export function readStream(
readable: Readable,
maxSize = MAX_EMBEDDABLE_SIZE
): Promise<Buffer> {
return readBufferWithLimit(readable, maxSize);
}
type RequestClosedCallback = (isAborted: boolean) => void;
type SignalReturnType = {
signal: AbortSignal;
onConnectionClosed: (cb: RequestClosedCallback) => void;
};
export function getSignal(req: Request): SignalReturnType {
const controller = new AbortController();
let isAborted = true;
let callback: ((isAborted: boolean) => void) | undefined = undefined;
const onSocketEnd = () => {
isAborted = false;
};
const onSocketClose = (hadError: boolean) => {
req.socket.off('end', onSocketEnd);
req.socket.off('close', onSocketClose);
const aborted = hadError || isAborted;
if (aborted) {
controller.abort();
}
callback?.(aborted);
};
req.socket.on('end', onSocketEnd);
req.socket.on('close', onSocketClose);
return {
signal: controller.signal,
onConnectionClosed: cb => (callback = cb),
};
}