twenty / packages /twenty-server /src /engine /subscriptions /utils /wrap-async-iterator-with-lifecycle.ts
Jules
Initial clean CRM deployment with prebuilt assets
d9494a5
Raw
History Blame Contribute Delete
2.26 kB
import { isDefined } from 'twenty-shared/utils';
type AsyncIteratorLifecycleOptions<T> = {
initialValue?: T;
onHeartbeat?: () => Promise<boolean>;
heartbeatIntervalMs?: number;
onCleanup?: () => Promise<void>;
onCleanupError?: (error: unknown) => void;
};
export function wrapAsyncIteratorWithLifecycle<T>(
iterator: AsyncIterableIterator<T>,
options: AsyncIteratorLifecycleOptions<T>,
): AsyncIterableIterator<T> {
const {
initialValue,
onHeartbeat,
heartbeatIntervalMs,
onCleanup,
onCleanupError,
} = options;
let heartbeatInterval: NodeJS.Timeout | null = null;
let hasYieldedInitialValue = false;
const startHeartbeat = () => {
if (onHeartbeat && heartbeatIntervalMs) {
heartbeatInterval = setInterval(() => {
try {
void onHeartbeat().catch(() => {});
} catch {
// Heartbeat failure shouldn't crash the stream
}
}, heartbeatIntervalMs);
}
};
const cleanup = async () => {
if (heartbeatInterval) {
clearInterval(heartbeatInterval);
heartbeatInterval = null;
}
if (onCleanup) {
try {
await onCleanup();
} catch (error) {
onCleanupError?.(error);
}
}
};
return {
next: async () => {
if (!isDefined(heartbeatInterval)) {
startHeartbeat();
}
if (isDefined(initialValue) && !hasYieldedInitialValue) {
hasYieldedInitialValue = true;
return { done: false, value: initialValue };
}
let result: IteratorResult<T>;
try {
result = await iterator.next();
} catch (error) {
await cleanup();
throw error;
}
if (result.done) {
await cleanup();
}
return result;
},
return: async () => {
let result: IteratorResult<T>;
try {
await cleanup();
} finally {
result = (await iterator.return?.()) ?? {
done: true,
value: undefined,
};
}
return result;
},
throw: async (error) => {
await cleanup();
if (iterator.throw) {
return iterator.throw(error);
}
throw error;
},
[Symbol.asyncIterator]() {
return this;
},
};
}