Spaces:
Build error
Build error
File size: 2,256 Bytes
d9494a5 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 | 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;
},
};
}
|