diff --git a/@vates/generator-toolbox/src/synchronized.mts b/@vates/generator-toolbox/src/synchronized.mts index e612db1025..e9802d6806 100644 --- a/@vates/generator-toolbox/src/synchronized.mts +++ b/@vates/generator-toolbox/src/synchronized.mts @@ -144,14 +144,7 @@ class Forked implements AsyncGenerator { return this.#parent.remove(this.#uid, e) } - async *[Symbol.asyncIterator](): AsyncGenerator { - while (true) { - const res = await this.next() - if (res.done) { - break - } - - yield res.value - } + [Symbol.asyncIterator](): AsyncGenerator { + return this } } diff --git a/@vates/generator-toolbox/src/synchronized.test.mts b/@vates/generator-toolbox/src/synchronized.test.mts index c0c183f4e4..3609fe61a7 100644 --- a/@vates/generator-toolbox/src/synchronized.test.mts +++ b/@vates/generator-toolbox/src/synchronized.test.mts @@ -2,8 +2,7 @@ import assert from 'node:assert' import { suite, test } from 'node:test' import { Synchronized } from './synchronized.mjs' -// eslint-disable-next-line @typescript-eslint/no-unused-vars -async function* makeRangeGenerator(end = Infinity, progress = { yielded: 0 }, onYielded = (val: unknown) => {}) { +async function* makeRangeGenerator(end = Infinity, progress = { yielded: 0 }, onYielded = (_val: unknown) => {}) { for (let i = 0; i < end; i++) { await new Promise(resolve => setTimeout(resolve, 10)) yield i @@ -17,7 +16,7 @@ async function consume( iterable: AsyncGenerator, delay = 2500, // eslint-disable-next-line @typescript-eslint/no-unused-vars - onConsumed = (val: unknown, iterable: AsyncGenerator) => Promise.resolve(false) + onConsumed = (_val: unknown, _iterable: AsyncGenerator) => Promise.resolve(false) ) { for await (const val of iterable) { await new Promise(resolve => setTimeout(resolve, delay)) @@ -28,6 +27,30 @@ async function consume( } suite('success', () => { + test('for-await-of break stops the source generator', async () => { + let sourceClosed = false + const source = (async function* () { + try { + for (let i = 0; i < 10; i++) { + await new Promise(resolve => setTimeout(resolve, 10)) + yield i + } + } finally { + sourceClosed = true + } + })() + + const forker = new Synchronized(source) + const fork = forker.fork('first') + + // for-await-of must call fork.return() on break, which propagates to the source + for await (const val of fork) { + if (val === 2) break + } + + assert.strictEqual(sourceClosed, true) + }) + test('if works with multiple consumer', async () => { const progress = { yielded: 0 } const generator = makeRangeGenerator(3, progress) diff --git a/@vates/generator-toolbox/src/timeout.mts b/@vates/generator-toolbox/src/timeout.mts index f41f4e9ee8..e31b31d092 100644 --- a/@vates/generator-toolbox/src/timeout.mts +++ b/@vates/generator-toolbox/src/timeout.mts @@ -1,4 +1,7 @@ import assert from 'node:assert' + +const TIMEOUT = Symbol('timeout') + export class Timeout implements AsyncGenerator { #source: AsyncGenerator #timeout: number @@ -8,21 +11,30 @@ export class Timeout implements AsyncGenerator { this.#timeout = timeout } async next(): Promise> { - let timeout: ReturnType - const promiseTimeout = new Promise((_, reject) => { - timeout = setTimeout(() => { - reject(new Error('Timeout reached ')) - }, this.#timeout) - }) - const promiseNext = new Promise((resolve, reject) => { - this.#source.next().then(res => { - // ensure timetout won't fire later - clearTimeout(timeout) - resolve(res) - }, reject) - }) - // promiseTimeout will never resolve - return Promise.race([promiseNext, promiseTimeout]) as Promise> + let timeoutHandle: ReturnType + + const sourceNext = this.#source.next() + const result = await Promise.race([ + sourceNext.then(res => { + clearTimeout(timeoutHandle) + return res + }), + new Promise(resolve => { + timeoutHandle = setTimeout(() => resolve(TIMEOUT), this.#timeout) + }), + ]) + + if (result === TIMEOUT) { + // stop the source once the in-flight next settles to avoid data loss on the next call + // we don't want to await for the end of timeouted code + sourceNext.then( + () => this.#source.return(undefined).catch(() => {}), + () => {} // source already errored, nothing to clean up + ) + throw new Error('Timeout reached') + } + + return result } return(): Promise> { return this.#source.return(undefined) diff --git a/@vates/generator-toolbox/src/timeout.test.mts b/@vates/generator-toolbox/src/timeout.test.mts index b600486d87..5f0e21ee5f 100644 --- a/@vates/generator-toolbox/src/timeout.test.mts +++ b/@vates/generator-toolbox/src/timeout.test.mts @@ -64,4 +64,24 @@ describe('Timeout class', () => { await assert.rejects(timeout.throw(new Error('Test error')), 'Expected error was not thrown') }) + + it('should close the source generator after a timeout', async () => { + let sourceClosed = false + const slowSource = (async function* () { + try { + // resolves after 100ms so the in-flight next() eventually settles + await new Promise(resolve => setTimeout(resolve, 100)) + yield 1 + } finally { + sourceClosed = true + } + })() + + const timeout = new Timeout(slowSource, 10) // fires before source yields + await assert.rejects(() => timeout.next(), /Timeout reached/) + + // wait for the in-flight next() to settle so source.return() can be called + await new Promise(resolve => setTimeout(resolve, 200)) + assert.strictEqual(sourceClosed, true) + }) }) diff --git a/@xen-orchestra/disk-transform/src/Disk.mts b/@xen-orchestra/disk-transform/src/Disk.mts index 4e442a8ce8..cbbaf49c8d 100644 --- a/@xen-orchestra/disk-transform/src/Disk.mts +++ b/@xen-orchestra/disk-transform/src/Disk.mts @@ -48,17 +48,10 @@ export abstract class Disk { abstract buildDiskBlockGenerator(): Promise> | AsyncGenerator async *diskBlocks(uid?: string): AsyncGenerator { try { - // compute next block while the destination is consuming the current block const blockGenerator = await this.buildDiskBlockGenerator() - let next = blockGenerator.next() - while (true) { - const res = await next - next = blockGenerator.next() - if (res.done) { - break - } + for await (const block of blockGenerator) { this.#generatedDiskBlocks++ - yield res.value + yield block } } finally { await this.close() diff --git a/@xen-orchestra/disk-transform/src/Disk.test.mts b/@xen-orchestra/disk-transform/src/Disk.test.mts index 1fc91421ff..ddfcd176a5 100644 --- a/@xen-orchestra/disk-transform/src/Disk.test.mts +++ b/@xen-orchestra/disk-transform/src/Disk.test.mts @@ -28,7 +28,7 @@ test('Disk class', async t => { disk.blocks[key] = true } const foundKeys = [] - const generator = disk.diskBlocks() + const generator = disk.diskBlocks() for await (const block of generator) { foundKeys.push(block.index) assert.strictEqual(block.data.length, blockSize) @@ -39,4 +39,24 @@ test('Disk class', async t => { foundKeys.sort() assert.deepStrictEqual(keys, foundKeys) }) + + await t.test('diskBlocks calls progressHandler.done when consumer exits early', async () => { + let doneCalled = false + disk.progressHandler = { + setProgress: async () => {}, + done: async () => { + doneCalled = true + }, + } + disk.blocks[4] = true + disk.blocks[5] = true + disk.blocks[6] = true + + // break after the first block — the inner generator's finally must still run + for await (const _ of disk.diskBlocks()) { + break + } + + assert.strictEqual(doneCalled, true) + }) }) diff --git a/@xen-orchestra/disk-transform/src/ReadAhead.mts b/@xen-orchestra/disk-transform/src/ReadAhead.mts index 5192a5434e..b8ae3e6d7b 100644 --- a/@xen-orchestra/disk-transform/src/ReadAhead.mts +++ b/@xen-orchestra/disk-transform/src/ReadAhead.mts @@ -24,12 +24,14 @@ export class ReadAhead extends RandomDiskPassthrough { const PRELOAD_SIZE = Math.max(Math.min(Math.floor((blockIndexes.length * maxPercent) / 100), maxNumber), 1) let counter = 0 - try{ + try { for (const index of blockIndexes) { counter++ - await self.progressHandler?.setProgress(counter/ blockIndexes.length) + await self.progressHandler?.setProgress(counter / blockIndexes.length) if (preloaded.length < PRELOAD_SIZE) { - preloaded.push(self.source.readBlock(index)) + const p = self.source.readBlock(index) + p.catch(() => {}) // suppress unhandled rejection; error still surfaces when awaited via shift() + preloaded.push(p) } if (preloaded.length === PRELOAD_SIZE) { const next = (await preloaded.shift())! @@ -40,10 +42,11 @@ export class ReadAhead extends RandomDiskPassthrough { const next = (await preloaded.shift())! yield next } - }finally{ + } finally { + await Promise.allSettled(preloaded) // wait for in-flight I/O & release resources + preloaded.length = 0 await self.progressHandler?.done() } - } return generator() } diff --git a/CHANGELOG.unreleased.md b/CHANGELOG.unreleased.md index 8b1ffc99f3..7a4ad5b729 100644 --- a/CHANGELOG.unreleased.md +++ b/CHANGELOG.unreleased.md @@ -15,6 +15,8 @@ > Users must be able to say: “I had this issue, happy to know it's fixed” +- [Backup] Reduce backup memory consumption (PR [#9557](https://github.com/vatesfr/xen-orchestra/pull/9557)) + ### Packages to release > When modifying a package, add it here with its release type. @@ -31,4 +33,8 @@ +- @vates/generator-toolbox patch +- @xen-orchestra/disk-transform patch +- xo-server minor + diff --git a/packages/xo-server/src/xo-mixins/logs/index.mjs b/packages/xo-server/src/xo-mixins/logs/index.mjs index 3d9961cabe..00b75413b1 100644 --- a/packages/xo-server/src/xo-mixins/logs/index.mjs +++ b/packages/xo-server/src/xo-mixins/logs/index.mjs @@ -1,5 +1,5 @@ import { configure } from '@xen-orchestra/log/configure' -import { createCaptureTransport } from '@xen-orchestra/log/capture' +import { createLogger, createCaptureTransport } from '@xen-orchestra/log/capture' import { dedupe } from '@xen-orchestra/log/dedupe' import { defer, fromEvent } from 'promise-toolbox' @@ -7,12 +7,20 @@ import LevelDbLogger from './loggers/leveldb.mjs' const { DEBUG } = process.env +const { warn } = createLogger('xo:mixins:logs') export default class Logs { constructor(app) { this._app = app app.hooks.on('clean', () => this._gc()) + setInterval( + () => { + this._gc().catch(error => warn('error while interval log cleaning', error)) + }, + 6 * 60 * 60 * 1000 + ).unref() + app.config.watch('logs', ({ filter, level, transport: transportsObject }) => { const transports = [] for (const id of Object.keys(transportsObject)) {