feat: import people
Signed-off-by: Sebastian Krupinski <krupinski01@gmail.com>
This commit is contained in:
+31
-34
@@ -53,32 +53,55 @@ export async function transceivePost<TRequest, TResponse>(
|
||||
* Stream an NDJSON API response, unwrapping data frames for the caller.
|
||||
*
|
||||
* The server emits one JSON object per line with a transport-level `type`
|
||||
* discriminant. This helper consumes control and error frames, forwards only
|
||||
* unwrapped `data` payloads to the caller, and returns the final stream total.
|
||||
* discriminant. This consumes the chunked body, splits it into lines, forwards
|
||||
* only unwrapped `data` payloads to the caller, and returns the final total.
|
||||
*
|
||||
* @param operation - Operation name, e.g. 'entity.listStream'
|
||||
* @param operation - Operation name, e.g. 'entity.listStream', 'entity.import'
|
||||
* @param data - Operation-specific request data
|
||||
* @param onData - Synchronous callback invoked for every unwrapped data payload.
|
||||
* May throw to abort the stream.
|
||||
* @param user - Optional user identifier override
|
||||
* @param options - Optional `user` override and an `onStart` hook, invoked once
|
||||
* with the expected total (the progress denominator) when the
|
||||
* server declares one on the start frame.
|
||||
* @returns Promise resolving to the final stream total from the control/end frame
|
||||
*/
|
||||
export async function transceiveStream<TRequest, TData>(
|
||||
operation: string,
|
||||
data: TRequest,
|
||||
onData: (data: TData) => void,
|
||||
user?: string
|
||||
options?: { user?: string; onStart?: (expected?: number) => void }
|
||||
): Promise<{ total: number }> {
|
||||
const request: ApiRequest<TRequest> = {
|
||||
version: API_VERSION,
|
||||
transaction: generateTransactionId(),
|
||||
operation,
|
||||
data,
|
||||
user,
|
||||
user: options?.user,
|
||||
};
|
||||
|
||||
let total = 0;
|
||||
|
||||
// Interpret one NDJSON line: control frames carry start/end metadata, error
|
||||
// frames abort, data frames are unwrapped to the caller.
|
||||
const dispatch = (line: string): void => {
|
||||
const message = JSON.parse(line) as ApiStreamResponse<TData>;
|
||||
|
||||
if (message.type === 'control') {
|
||||
if (message.status === 'start') {
|
||||
options?.onStart?.(message.total);
|
||||
} else if (message.status === 'end') {
|
||||
total = message.total;
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === 'error') {
|
||||
throw new Error(`[${operation}] ${message.message}`);
|
||||
}
|
||||
|
||||
onData(message.data);
|
||||
};
|
||||
|
||||
await fetchWrapper.post(API_URL, request, {
|
||||
headers: { 'Accept': 'application/json' },
|
||||
onStream: async (response: Response) => {
|
||||
@@ -100,38 +123,12 @@ export async function transceiveStream<TRequest, TData>(
|
||||
buffer = lines.pop()!; // retain any incomplete trailing chunk
|
||||
|
||||
for (const line of lines) {
|
||||
if (!line.trim()) continue;
|
||||
const message = JSON.parse(line) as ApiStreamResponse<TData>;
|
||||
|
||||
if (message.type === 'control') {
|
||||
if (message.status === 'end') {
|
||||
total = message.total;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
if (message.type === 'error') {
|
||||
throw new Error(`[${operation}] ${message.message}`);
|
||||
}
|
||||
|
||||
onData(message.data);
|
||||
if (line.trim()) dispatch(line);
|
||||
}
|
||||
}
|
||||
|
||||
// flush any remaining bytes still in the buffer
|
||||
if (buffer.trim()) {
|
||||
const message = JSON.parse(buffer) as ApiStreamResponse<TData>;
|
||||
|
||||
if (message.type === 'control') {
|
||||
if (message.status === 'end') {
|
||||
total = message.total;
|
||||
}
|
||||
} else if (message.type === 'error') {
|
||||
throw new Error(`[${operation}] ${message.message}`);
|
||||
} else {
|
||||
onData(message.data);
|
||||
}
|
||||
}
|
||||
if (buffer.trim()) dispatch(buffer);
|
||||
} finally {
|
||||
reader.releaseLock();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user