feat(backups): use XO Tasks (#9734)

Change the backups so they stop using their own task system, and use the same system as XO Tasks.

For the moment, it shouldn't affect the users, but when the previous backup logs will be old enough to be erased, this will allow us to have shorter loading times for backup logs.

XO-52
This commit is contained in:
Bastien Nollet
2026-04-17 17:49:18 +02:00
committed by GitHub
parent 9dc52f2e5b
commit fbba05bbde
37 changed files with 1692 additions and 1735 deletions

View File

@@ -23,7 +23,7 @@ exports.makeOnProgress = function ({ onRootTaskEnd = noop, onRootTaskStart = noo
Object.defineProperty(taskLog, '$root', { value: taskLog })
// start of a root task
onRootTaskStart(taskLog)
onRootTaskStart(taskLog, event)
} else {
// start of a subtask
const parent = taskLogs.get(parentId)
@@ -54,10 +54,10 @@ exports.makeOnProgress = function ({ onRootTaskEnd = noop, onRootTaskStart = noo
}
if (type === 'end' && taskLog.$root === taskLog) {
onRootTaskEnd(taskLog)
onRootTaskEnd(taskLog, event)
}
}
onTaskUpdate(taskLog)
onTaskUpdate(taskLog, event)
}
}

View File

@@ -7,6 +7,27 @@ function alreadyEnded() {
throw new Error('task has already ended')
}
const serializeErrors = errors => (Array.isArray(errors) ? errors.map(serializeError) : errors)
// Create a serializable object from an error.
//
// Otherwise some fields might be non-enumerable and missing from logs.
const serializeError = error => {
return error instanceof Error
? {
...error, // Copy enumerable properties.
cause: error.cause,
code: error.code,
errors: serializeErrors(error.errors), // supports AggregateError
message: error.message,
name: error.name,
stack: error.stack,
}
: error
}
exports.serializeError = serializeError
// define a read-only, non-enumerable, non-configurable property
function define(object, property, value) {
Object.defineProperty(object, property, { value })
@@ -72,7 +93,11 @@ class Task {
_onProgress
get id() {
return (this.id = Math.random().toString(36).slice(2))
// Use a compact, sortable, string representation of the creation date,
// and add a random part to differentiate tasks created at the same time
//
// Due to the padding, dates are sortable up to 5188-04-22T11:04:28.415Z
return (this.id = Date.now().toString(36).padStart(9, '0') + '-' + Math.random().toString(36).slice(2))
}
set id(value) {
define(this, 'id', value)
@@ -136,6 +161,10 @@ class Task {
assert.equal(this.#status, PENDING)
assert.equal(this.#running, false)
if (result instanceof Error && result.toJSON === undefined) {
result = serializeError(result)
}
this.#emit('end', { status, result })
this._onProgress = alreadyEnded
this.#status = status

View File

@@ -4,7 +4,7 @@ const assert = require('node:assert').strict
const { describe, it } = require('node:test')
const { makeOnProgress } = require('./combineEvents.js')
const { Task } = require('./index.js')
const { Task, serializeError } = require('./index.js')
const noop = Function.prototype
@@ -28,6 +28,36 @@ function createTask(opts) {
return task
}
describe('serializeError', function () {
it('serializes simple error', function () {
const err = new Error('test Error')
err.code = 404
const serialized = serializeError(err)
assert.equal(serialized.message, 'test Error')
assert.equal(serialized.name, 'Error')
assert.equal(serialized.code, 404)
assert.notEqual(serialized.stack, undefined)
})
it('serializes aggregate error', function () {
const err = new AggregateError([new Error('test error 1'), new Error('test error 2')], 'Error message', {
cause: 'broken',
})
const serialized = serializeError(err)
assert.equal(serialized.message, 'Error message')
assert.equal(serialized.name, 'AggregateError')
assert.equal(serialized.cause, 'broken')
assert.notEqual(serialized.stack, undefined)
assert.equal(serialized.errors.length, 2)
assert.equal(serialized.errors[0].name, 'Error')
assert.equal(serialized.errors[1].name, 'Error')
assert.equal(serialized.errors[0].message, 'test error 1')
assert.equal(serialized.errors[1].message, 'test error 2')
assert.notEqual(serialized.errors[0].stack, undefined)
assert.notEqual(serialized.errors[1].stack, undefined)
})
})
describe('Task', function () {
describe('constructor', function () {
it('data properties are passed to the start event', async function () {
@@ -106,7 +136,7 @@ describe('Task', function () {
assert.equal(task.$events.length, 3)
assertEvent(task, { type: 'start' }, 0)
assertEvent(task, { type: 'abortionRequested', reason }, 1)
assertEvent(task, { type: 'end', status: 'failure', result }, 2)
assertEvent(task, { type: 'end', status: 'failure', result: serializeError(result) }, 2)
})
it('does not abort if the task succeed', async function () {
@@ -176,7 +206,7 @@ describe('Task', function () {
assertEvent(task, {
status: 'failure',
result: error,
result: serializeError(error),
type: 'end',
})
})
@@ -475,7 +505,7 @@ describe('Task', function () {
assert.equal(task.status, 'failure')
assertEvent(task, {
status: 'failure',
result: e,
result: serializeError(e),
type: 'end',
})
})
@@ -511,7 +541,7 @@ describe('Task', function () {
assert.equal(task.status, 'failure')
assertEvent(task, {
status: 'failure',
result: e,
result: serializeError(e),
type: 'end',
})
})

View File

@@ -1,4 +1,4 @@
import { Task } from './Task.mjs'
import { Task } from '@vates/task'
export class HealthCheckVmBackup {
#restoredVm
@@ -14,7 +14,7 @@ export class HealthCheckVmBackup {
async run() {
return Task.run(
{
name: 'vmstart',
properties: { name: 'vmstart' },
},
async () => {
let restoredVm = this.#restoredVm

View File

@@ -2,9 +2,9 @@ import assert from 'node:assert'
import { formatFilenameDate } from './_filenameDate.mjs'
import { importIncrementalVm } from './_incrementalVm.mjs'
import { Task } from './Task.mjs'
import { watchStreamSize } from './_watchStreamSize.mjs'
import { decorateClass } from '@vates/decorate-with'
import { Task } from '@vates/task'
import { createLogger } from '@xen-orchestra/log'
import { dirname, join } from 'node:path'
import pickBy from 'lodash/pickBy.js'
@@ -244,7 +244,7 @@ export class ImportVmBackup {
return Task.run(
{
name: 'transfer',
properties: { name: 'transfer' },
},
async () => {
const xapi = this._xapi

View File

@@ -1,155 +0,0 @@
import CancelToken from 'promise-toolbox/CancelToken'
import Zone from 'node-zone'
const logAfterEnd = log => {
const error = new Error('task has already ended')
error.log = log
throw error
}
const noop = Function.prototype
const serializeErrors = errors => (Array.isArray(errors) ? errors.map(serializeError) : errors)
// Create a serializable object from an error.
//
// Otherwise some fields might be non-enumerable and missing from logs.
const serializeError = error =>
error instanceof Error
? {
...error, // Copy enumerable properties.
code: error.code,
errors: serializeErrors(error.errors), // supports AggregateError
message: error.message,
name: error.name,
stack: error.stack,
}
: error
const $$task = Symbol('@xen-orchestra/backups/Task')
export class Task {
static get cancelToken() {
const task = Zone.current.data[$$task]
return task !== undefined ? task.#cancelToken : CancelToken.none
}
static run(opts, fn) {
return new this(opts).run(fn, true)
}
static wrapFn(opts, fn) {
// compatibility with @decorateWith
if (typeof fn !== 'function') {
;[fn, opts] = [opts, fn]
}
return function () {
return Task.run(typeof opts === 'function' ? opts.apply(this, arguments) : opts, () => fn.apply(this, arguments))
}
}
#cancelToken
#id = Math.random().toString(36).slice(2)
#onLog
#zone
constructor({ name, data, onLog }) {
let parentCancelToken, parentId
if (onLog === undefined) {
const parent = Zone.current.data[$$task]
if (parent === undefined) {
onLog = noop
} else {
onLog = log => parent.#onLog(log)
parentCancelToken = parent.#cancelToken
parentId = parent.#id
}
}
const zone = Zone.current.fork('@xen-orchestra/backups/Task')
zone.data[$$task] = this
this.#zone = zone
const { cancel, token } = CancelToken.source(parentCancelToken && [parentCancelToken])
this.#cancelToken = token
this.cancel = cancel
this.#onLog = onLog
this.#log('start', {
data,
message: name,
parentId,
})
}
failure(error) {
this.#end('failure', serializeError(error))
}
info(message, data) {
this.#log('info', { data, message })
}
/**
* Run a function in the context of this task
*
* In case of error, the task will be failed.
*
* @typedef Result
* @param {() => Result} fn
* @param {boolean} last - Whether the task should succeed if there is no error
* @returns Result
*/
run(fn, last = false) {
return this.#zone.run(() => {
try {
const result = fn()
let then
if (result != null && typeof (then = result.then) === 'function') {
then.call(result, last && (value => this.success(value)), error => this.failure(error))
} else if (last) {
this.success(result)
}
return result
} catch (error) {
this.failure(error)
throw error
}
})
}
success(value) {
this.#end('success', value)
}
warning(message, data) {
this.#log('warning', { data, message })
}
wrapFn(fn, last) {
const task = this
return function () {
return task.run(() => fn.apply(this, arguments), last)
}
}
#end(status, result) {
this.#log('end', { result, status })
this.#onLog = logAfterEnd
}
#log(event, props) {
this.#onLog({
...props,
event,
taskId: this.#id,
timestamp: Date.now(),
})
}
}
for (const method of ['info', 'warning']) {
Task[method] = (...args) => Zone.current.data[$$task]?.[method](...args)
}

View File

@@ -14,10 +14,10 @@ import { decorateMethodsWith } from '@vates/decorate-with'
import { deduped } from '@vates/disposable/deduped.js'
import { getHandler } from '@xen-orchestra/fs'
import { parseDuration } from '@vates/parse-duration'
import { Task } from '@vates/task'
import { Xapi } from '@xen-orchestra/xapi'
import { RemoteAdapter } from './RemoteAdapter.mjs'
import { Task } from './Task.mjs'
createCachedLookup().patchGlobal()
@@ -178,8 +178,8 @@ process.on('message', async message => {
const result = message.runWithLogs
? await Task.run(
{
name: 'backup run',
onLog: data =>
properties: { name: 'backup run', ...message.data.jobData },
onProgress: data =>
emitMessage({
data,
type: 'log',

View File

@@ -9,8 +9,8 @@ import { RemoteVhdDisk } from './disks/RemoteVhdDisk.mjs'
import { RemoteVhdDiskChain } from './disks/RemoteVhdDiskChain.mjs'
import { MergeRemoteDisk } from './disks/MergeRemoteDisk.mjs'
import { Task } from './Task.mjs'
import { Disposable } from 'promise-toolbox'
import { Task } from '@vates/task'
import handlerPath from '@xen-orchestra/fs/path'
const { DISK_TYPES } = Constants
@@ -534,7 +534,7 @@ export async function cleanVm(
await Promise.all([
...unusedVhdsDeletion,
toMerge.length !== 0 && (merge ? Task.run({ name: 'merge' }, doMerge) : () => Promise.resolve()),
toMerge.length !== 0 && (merge ? Task.run({ properties: { name: 'merge' } }, doMerge) : () => Promise.resolve()),
asyncMap(unusedXvas, path => {
logWarn('unused XVA', { path })
if (remove) {

View File

@@ -4,9 +4,9 @@ import { asyncMap } from '@xen-orchestra/async-map'
import { CancelToken } from 'promise-toolbox'
import { compareVersions } from 'compare-versions'
import { defer } from 'golike-defer'
import { Task } from '@vates/task'
import { cancelableMap } from './_cancelableMap.mjs'
import { Task } from './Task.mjs'
import pick from 'lodash/pick.js'
import { BASE_DELTA_VDI, CONTENT_KEY, COPY_OF, VM_UUID } from './_otherConfig.mjs'

View File

@@ -1,4 +1,5 @@
import { asyncMap } from '@xen-orchestra/async-map'
import { Task } from '@vates/task'
import Disposable from 'promise-toolbox/Disposable'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
@@ -6,9 +7,10 @@ import { extractIdsFromSimplePattern } from '../extractIdsFromSimplePattern.mjs'
import { PoolMetadataBackup } from './_PoolMetadataBackup.mjs'
import { XoMetadataBackup } from './_XoMetadataBackup.mjs'
import { DEFAULT_SETTINGS, Abstract } from './_Abstract.mjs'
import { runTask } from './_runTask.mjs'
import { getAdaptersByRemote } from './_getAdaptersByRemote.mjs'
const noop = Function.prototype
const DEFAULT_METADATA_SETTINGS = {
retentionPoolMetadata: 0,
retentionXoMetadata: 0,
@@ -55,13 +57,16 @@ export const Metadata = class MetadataBackupRunner extends Abstract {
poolIds.map(id =>
this._getRecord('pool', id).catch(error => {
// See https://github.com/vatesfr/xen-orchestra/commit/6aa6cfba8ec939c0288f0fa740f6dfad98c43cbb
runTask(
Task.run(
{
name: 'get pool record',
data: { type: 'pool', id },
properties: {
id,
name: 'get pool record',
type: 'pool',
},
},
() => Promise.reject(error)
)
).catch(noop)
})
)
),
@@ -81,11 +86,11 @@ export const Metadata = class MetadataBackupRunner extends Abstract {
if (pools.length !== 0 && settings.retentionPoolMetadata !== 0) {
promises.push(
asyncMap(pools, async pool =>
runTask(
Task.run(
{
name: `Starting metadata backup for the pool (${pool.$id}). (${job.id})`,
data: {
properties: {
id: pool.$id,
name: `Starting metadata backup for the pool (${pool.$id}). (${job.id})`,
pool,
poolMaster: await ignoreErrors.call(pool.$xapi.getRecord('host', pool.master)),
type: 'pool',
@@ -100,17 +105,17 @@ export const Metadata = class MetadataBackupRunner extends Abstract {
schedule,
settings,
}).run()
)
).catch(noop)
)
)
}
if (job.xoMetadata !== undefined && settings.retentionXoMetadata !== 0) {
promises.push(
runTask(
Task.run(
{
name: `Starting XO metadata backup. (${job.id})`,
data: {
properties: {
name: `Starting XO metadata backup. (${job.id})`,
type: 'xo',
},
},
@@ -122,7 +127,7 @@ export const Metadata = class MetadataBackupRunner extends Abstract {
schedule,
settings,
}).run()
)
).catch(noop)
)
}
await Promise.all(promises)

View File

@@ -1,9 +1,9 @@
import { asyncMapSettled } from '@xen-orchestra/async-map'
import Disposable from 'promise-toolbox/Disposable'
import { limitConcurrency } from 'limit-concurrency-decorator'
import { Task } from '@vates/task'
import { extractIdsFromSimplePattern } from '../extractIdsFromSimplePattern.mjs'
import { Task } from '../Task.mjs'
import { DEFAULT_SETTINGS, Abstract } from './_Abstract.mjs'
import { getAdaptersByRemote } from './_getAdaptersByRemote.mjs'
import { FullRemote } from './_vmRunners/FullRemote.mjs'
@@ -78,7 +78,7 @@ export const VmsRemote = class RemoteVmsBackupRunner extends Abstract {
}
nTriesByVmId[vmUuid]++
const taskStart = { name: 'backup VM', data: { type: 'VM', id: vmUuid } }
const taskStart = { properties: { id: vmUuid, name: 'backup VM', type: 'VM' } }
const vmSettings = { ...settings, ...allSettings[vmUuid] }
const isLastRun = nTriesByVmId[vmUuid] === vmSettings.nRetriesVmBackupFailures + 1
@@ -113,23 +113,27 @@ export const VmsRemote = class RemoteVmsBackupRunner extends Abstract {
taskByVmId[vmUuid] = new Task(taskStart)
}
const task = taskByVmId[vmUuid]
// error has to be caught in the task to prevent its failure, but handled outside the task to execute another task.run()
let taskError
return task
.run(async () => {
try {
const result = await vmBackup.run()
task.success(result)
return result
} catch (error) {
if (isLastRun) {
throw error
} else {
Task.warning(`Retry the VM mirror backup due to an error`, {
attempt: nTriesByVmId[vmUuid],
error: error.message,
})
queue.add(vmUuid)
}
.runInside(async () =>
vmBackup.run().catch(error => {
taskError = error
})
)
.then(result => {
if (taskError === undefined) {
return task.success(result)
}
if (isLastRun) {
return task.failure(taskError)
}
// don't end the task
task.warning(`Retry the VM mirror backup due to an error`, {
attempt: nTriesByVmId[vmUuid],
error: taskError.message,
})
queue.add(vmUuid)
})
.catch(noop)
}

View File

@@ -1,11 +1,10 @@
import { asyncMapSettled } from '@xen-orchestra/async-map'
import Disposable from 'promise-toolbox/Disposable'
import { limitConcurrency } from 'limit-concurrency-decorator'
import { Task } from '@vates/task'
import { extractIdsFromSimplePattern } from '../extractIdsFromSimplePattern.mjs'
import { Task } from '../Task.mjs'
import { DEFAULT_SETTINGS, Abstract } from './_Abstract.mjs'
import { runTask } from './_runTask.mjs'
import { getAdaptersByRemote } from './_getAdaptersByRemote.mjs'
import { IncrementalXapi } from './_vmRunners/IncrementalXapi.mjs'
import { FullXapi } from './_vmRunners/FullXapi.mjs'
@@ -65,13 +64,12 @@ export const VmsXapi = class VmsXapiBackupRunner extends Abstract {
Disposable.all(
extractIdsFromSimplePattern(job.srs).map(id =>
this._getRecord('SR', id).catch(error => {
runTask(
Task.run(
{
name: 'get SR record',
data: { type: 'SR', id },
properties: { id, name: 'get SR record', type: 'SR' },
},
() => Promise.reject(error)
)
).catch(noop)
})
)
),
@@ -108,11 +106,12 @@ export const VmsXapi = class VmsXapiBackupRunner extends Abstract {
}
return taskByVmId[vmUuid]
}
const vmBackupFailed = error => {
const vmBackupFailed = async (error, task) => {
if (isLastRun) {
throw error
return task.failure(error)
} else {
Task.warning(`Retry the VM backup due to an error`, {
// don't end the task
task.warning(`Retry the VM backup due to an error`, {
attempt: nTriesByVmId[vmUuid],
error: error.message,
})
@@ -126,19 +125,21 @@ export const VmsXapi = class VmsXapiBackupRunner extends Abstract {
nTriesByVmId[vmUuid]++
const vmSettings = { ...settings, ...allSettings[vmUuid] }
const taskStart = { name: 'backup VM', data: { type: 'VM', id: vmUuid } }
const taskStart = { properties: { id: vmUuid, name: 'backup VM', type: 'VM' } }
const isLastRun = nTriesByVmId[vmUuid] === vmSettings.nRetriesVmBackupFailures + 1
return this._getRecord('VM', vmUuid).then(
disposableVm =>
Disposable.use(disposableVm, async vm => {
if (taskStart.data.name_label === undefined) {
taskStart.data.name_label = vm.name_label
if (taskStart.properties.name_label === undefined) {
taskStart.properties.name_label = vm.name_label
}
const task = getVmTask()
// error has to be caught in the task to prevent its failure, but handled outside the task to execute another task.run()
let taskError
return task
.run(async () => {
.runInside(async () => {
const opts = {
baseSettings,
config,
@@ -164,21 +165,21 @@ export const VmsXapi = class VmsXapiBackupRunner extends Abstract {
throw new Error(`Job mode ${job.mode} not implemented`)
}
}
try {
const result = await vmBackup.run()
return vmBackup.run().catch(error => {
taskError = error
})
})
.then(result => {
if (taskError === undefined) {
task.success(result)
return result
} catch (error) {
vmBackupFailed(error)
} else {
// ending the task with error or not ending the task
vmBackupFailed(taskError, task)
}
})
.catch(noop) // errors are handled by logs
}),
error =>
getVmTask().run(() => {
vmBackupFailed(error)
})
error => vmBackupFailed(error, getVmTask())
)
}
const { concurrency } = settings

View File

@@ -1,8 +1,10 @@
import Disposable from 'promise-toolbox/Disposable'
import pTimeout from 'promise-toolbox/timeout'
import { compileTemplate } from '@xen-orchestra/template'
import { runTask } from './_runTask.mjs'
import { RemoteTimeoutError } from './_RemoteTimeoutError.mjs'
import { Task } from '@vates/task'
const noop = Function.prototype
export const DEFAULT_SETTINGS = {
getRemoteTimeout: 300e3,
@@ -36,13 +38,16 @@ export const Abstract = class AbstractRunner {
})
} catch (error) {
// See https://github.com/vatesfr/xen-orchestra/commit/6aa6cfba8ec939c0288f0fa740f6dfad98c43cbb
runTask(
Task.run(
{
name: 'get remote adapter',
data: { type: 'remote', id: remoteId },
properties: {
id: remoteId,
name: 'get remote adapter',
type: 'remote',
},
},
() => Promise.reject(error)
)
).catch(noop)
}
}
}

View File

@@ -1,9 +1,9 @@
import { asyncMap } from '@xen-orchestra/async-map'
import { Task } from '@vates/task'
import { DIR_XO_POOL_METADATA_BACKUPS } from '../RemoteAdapter.mjs'
import { forkStreamUnpipe } from './_forkStreamUnpipe.mjs'
import { formatFilenameDate } from '../_filenameDate.mjs'
import { Task } from '../Task.mjs'
export const PATH_DB_DUMP = '/pool/xmldbdump'
@@ -54,8 +54,8 @@ export class PoolMetadataBackup {
([remoteId, adapter]) =>
Task.run(
{
name: `Starting metadata backup for the pool (${pool.$id}) for the remote (${remoteId}). (${job.id})`,
data: {
properties: {
name: `Starting metadata backup for the pool (${pool.$id}) for the remote (${remoteId}). (${job.id})`,
id: remoteId,
type: 'remote',
},

View File

@@ -1,9 +1,9 @@
import { asyncMap } from '@xen-orchestra/async-map'
import { join } from '@xen-orchestra/fs/path'
import { Task } from '@vates/task'
import { DIR_XO_CONFIG_BACKUPS } from '../RemoteAdapter.mjs'
import { formatFilenameDate } from '../_filenameDate.mjs'
import { Task } from '../Task.mjs'
export class XoMetadataBackup {
constructor({ config, job, remoteAdapters, schedule, settings }) {
@@ -51,8 +51,8 @@ export class XoMetadataBackup {
([remoteId, adapter]) =>
Task.run(
{
name: `Starting XO metadata backup for the remote (${remoteId}). (${job.id})`,
data: {
properties: {
name: `Starting XO metadata backup for the remote (${remoteId}). (${job.id})`,
id: remoteId,
type: 'remote',
},

View File

@@ -1,5 +0,0 @@
import { Task } from '../Task.mjs'
const noop = Function.prototype
export const runTask = (...args) => Task.run(...args).catch(noop) // errors are handled by logs

View File

@@ -1,4 +1,5 @@
import { createLogger } from '@xen-orchestra/log'
import { Task } from '@vates/task'
import keyBy from 'lodash/keyBy.js'
import { AbstractXapi } from './_AbstractXapi.mjs'
@@ -6,7 +7,6 @@ import { forkDeltaExport } from './_forkDeltaExport.mjs'
import { exportIncrementalVm } from '../../_incrementalVm.mjs'
import { IncrementalRemoteWriter } from '../_writers/IncrementalRemoteWriter.mjs'
import { IncrementalXapiWriter } from '../_writers/IncrementalXapiWriter.mjs'
import { Task } from '../../Task.mjs'
import {
DATETIME,
DELTA_CHAIN_LENGTH,

View File

@@ -1,6 +1,6 @@
import { asyncMap } from '@xen-orchestra/async-map'
import { createLogger } from '@xen-orchestra/log'
import { Task } from '../../Task.mjs'
import { Task } from '@vates/task'
const { debug, warn } = createLogger('xo:backups:AbstractVmRunner')
@@ -83,7 +83,7 @@ export const Abstract = class AbstractVmBackupRunner {
// create a task to have an info in the logs and reports
return Task.run(
{
name: 'health check',
properties: { name: 'health check' },
},
() => {
Task.info(`This VM doesn't match the health check's tags for this schedule`)

View File

@@ -4,12 +4,12 @@ import { decorateMethodsWith } from '@vates/decorate-with'
import { defer } from 'golike-defer'
import { Disposable } from 'promise-toolbox'
import { createPredicate } from 'value-matcher'
import { Task } from '@vates/task'
import { getVmBackupDir } from '../../_getVmBackupDir.mjs'
import { Abstract } from './_Abstract.mjs'
import { extractIdsFromSimplePattern } from '../../extractIdsFromSimplePattern.mjs'
import { Task } from '../../Task.mjs'
export const AbstractRemote = class AbstractRemoteVmBackupRunner extends Abstract {
_filterPredicate

View File

@@ -6,9 +6,9 @@ import { asyncMap } from '@xen-orchestra/async-map'
import { asyncEach } from '@vates/async-each'
import { decorateMethodsWith } from '@vates/decorate-with'
import { defer } from 'golike-defer'
import { Task } from '@vates/task'
import { getOldEntries } from '../../_getOldEntries.mjs'
import { Task } from '../../Task.mjs'
import { Abstract } from './_Abstract.mjs'
import {
COPY_OF,
@@ -186,7 +186,7 @@ export const AbstractXapi = class AbstractXapiVmBackupRunner extends Abstract {
const settings = this._settings
if (await this._mustDoSnapshot()) {
await Task.run({ name: 'snapshot' }, async () => {
await Task.run({ properties: { name: 'snapshot' } }, async () => {
if (!settings.bypassVdiChainsCheck) {
await vm.$assertHealthyVdiChains()
}

View File

@@ -1,6 +1,7 @@
import { Task } from '@vates/task'
import { formatFilenameDate } from '../../_filenameDate.mjs'
import { getOldEntries } from '../../_getOldEntries.mjs'
import { Task } from '../../Task.mjs'
import { MixinRemoteWriter } from './_MixinRemoteWriter.mjs'
import { AbstractFullWriter } from './_AbstractFullWriter.mjs'
@@ -9,11 +10,11 @@ export class FullRemoteWriter extends MixinRemoteWriter(AbstractFullWriter) {
constructor(props) {
super(props)
this.run = Task.wrapFn(
this.run = Task.wrap(
{
name: 'export',
data: {
properties: {
id: props.remoteId,
name: 'export',
type: 'remote',
// necessary?
@@ -64,7 +65,7 @@ export class FullRemoteWriter extends MixinRemoteWriter(AbstractFullWriter) {
await deleteOldBackups()
}
await Task.run({ name: 'transfer' }, async () => {
await Task.run({ properties: { name: 'transfer' } }, async () => {
await adapter.outputStream(dataFilename, stream, {
maxStreamLength,
streamLength,

View File

@@ -1,9 +1,9 @@
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { asyncMap, asyncMapSettled } from '@xen-orchestra/async-map'
import { Task } from '@vates/task'
import { formatFilenameDate } from '../../_filenameDate.mjs'
import { getOldEntries } from '../../_getOldEntries.mjs'
import { Task } from '../../Task.mjs'
import { AbstractFullWriter } from './_AbstractFullWriter.mjs'
import { MixinXapiWriter } from './_MixinXapiWriter.mjs'
@@ -14,11 +14,11 @@ export class FullXapiWriter extends MixinXapiWriter(AbstractFullWriter) {
constructor(props) {
super(props)
this.run = Task.wrapFn(
this.run = Task.wrap(
{
name: 'export',
data: {
properties: {
id: props.sr.uuid,
name: 'export',
name_label: this._sr.name_label,
type: 'SR',
@@ -52,7 +52,7 @@ export class FullXapiWriter extends MixinXapiWriter(AbstractFullWriter) {
}
let targetVmRef
await Task.run({ name: 'transfer' }, async () => {
await Task.run({ properties: { name: 'transfer' } }, async () => {
targetVmRef = await xapi.VM_import(stream, sr.$ref, vm =>
Promise.all([
!_warmMigration && vm.add_tags('Disaster Recovery'),

View File

@@ -7,10 +7,10 @@ import { createLogger } from '@xen-orchestra/log'
import { decorateClass } from '@vates/decorate-with'
import { defer } from 'golike-defer'
import { dirname, basename } from 'node:path'
import { Task } from '@vates/task'
import { formatFilenameDate } from '../../_filenameDate.mjs'
import { getOldEntries } from '../../_getOldEntries.mjs'
import { Task } from '../../Task.mjs'
import { MixinRemoteWriter } from './_MixinRemoteWriter.mjs'
import { AbstractIncrementalWriter } from './_AbstractIncrementalWriter.mjs'
@@ -72,19 +72,20 @@ export class IncrementalRemoteWriter extends MixinRemoteWriter(AbstractIncrement
prepare({ isFull }) {
// create the task related to this export and ensure all methods are called in this context
const task = new Task({
name: 'export',
data: {
properties: {
id: this._remoteId,
isFull,
name: 'export',
type: 'remote',
},
})
this.transfer = task.wrapFn(this.transfer)
this.healthCheck = task.wrapFn(this.healthCheck)
this.cleanup = task.wrapFn(this.cleanup)
this.afterBackup = task.wrapFn(this.afterBackup, true)
this._prepare = task.wrapInside(this._prepare)
this.transfer = task.wrapInside(this.transfer)
this.healthCheck = task.wrapInside(this.healthCheck)
this.cleanup = task.wrapInside(this.cleanup)
this.afterBackup = task.wrap(this.afterBackup)
return task.run(() => this._prepare())
return this._prepare()
}
async _prepare() {
@@ -223,7 +224,7 @@ export class IncrementalRemoteWriter extends MixinRemoteWriter(AbstractIncrement
vtpms: deltaExport.vtpms,
}
let size = 0
await Task.run({ name: 'transfer' }, async () => {
await Task.run({ properties: { name: 'transfer' } }, async () => {
await asyncEach(
Object.entries(deltaExport.disks),
async ([diskRef, disk]) => {

View File

@@ -1,11 +1,11 @@
import humanFormat from 'human-format'
import { asyncMapSettled } from '@xen-orchestra/async-map'
import { Task } from '@vates/task'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { getOldEntries } from '../../_getOldEntries.mjs'
import { importIncrementalVm } from '../../_incrementalVm.mjs'
import { Task } from '../../Task.mjs'
import { AbstractIncrementalWriter } from './_AbstractIncrementalWriter.mjs'
import { MixinXapiWriter } from './_MixinXapiWriter.mjs'
@@ -151,20 +151,24 @@ export class IncrementalXapiWriter extends MixinXapiWriter(AbstractIncrementalWr
prepare({ isFull }) {
// create the task related to this export and ensure all methods are called in this context
const task = new Task({
name: 'export',
data: {
properties: {
id: this._sr.uuid,
isFull,
name: 'export',
name_label: this._sr.name_label,
type: 'SR',
},
})
const hasHealthCheckSr = this._healthCheckSr !== undefined
this.transfer = task.wrapFn(this.transfer)
this.cleanup = task.wrapFn(this.cleanup, !hasHealthCheckSr)
this.healthCheck = task.wrapFn(this.healthCheck, hasHealthCheckSr)
this._prepare = task.wrapInside(this._prepare)
this.transfer = task.wrapInside(this.transfer)
if (this._healthCheckSr !== undefined) {
this.cleanup = task.wrapInside(this.cleanup)
this.healthCheck = task.wrap(this.healthCheck)
} else {
this.cleanup = task.wrap(this.cleanup)
}
return task.run(() => this._prepare(isFull))
return this._prepare(isFull)
}
async _prepare(isFull) {
@@ -265,7 +269,7 @@ export class IncrementalXapiWriter extends MixinXapiWriter(AbstractIncrementalWr
const { uuid: srUuid, $xapi: xapi } = sr
let targetVmRef
await Task.run({ name: 'transfer' }, async () => {
await Task.run({ properties: { name: 'transfer' } }, async () => {
targetVmRef = await importIncrementalVm(this.#decorateVmMetadata(deltaExport, timestamp), sr, {
targetRef: this._targetVmRef,
})
@@ -289,7 +293,7 @@ export class IncrementalXapiWriter extends MixinXapiWriter(AbstractIncrementalWr
` -- last replication: ${formatFilenameDate(timestamp)} ${humanFormat.bytes(size)} read`
)
// take a snapshot to ensure these data are not modified until next snapshot
await Task.run({ name: 'target snapshot' }, async () => {
await Task.run({ properties: { name: 'target snapshot' } }, async () => {
await xapi.VM_snapshot(targetVmRef, {
name_label: `${vm.name_label} - ${job.name} / ${schedule.name} ${formatFilenameDate(timestamp)}`,
})

View File

@@ -1,12 +1,12 @@
import { createLogger } from '@xen-orchestra/log'
import { join } from 'node:path'
import { Task } from '@vates/task'
import assert from 'node:assert'
import { formatFilenameDate } from '../../_filenameDate.mjs'
import { getVmBackupDir } from '../../_getVmBackupDir.mjs'
import { HealthCheckVmBackup } from '../../HealthCheckVmBackup.mjs'
import { ImportVmBackup } from '../../ImportVmBackup.mjs'
import { Task } from '../../Task.mjs'
import * as MergeWorker from '../../merge-worker/index.mjs'
import ms from 'ms'
import { getEntryStatus } from '../../_getOldEntries.mjs'
@@ -28,7 +28,7 @@ export const MixinRemoteWriter = (BaseClass = Object) =>
async _cleanVm(options) {
try {
return await Task.run({ name: 'clean-vm' }, () => {
return await Task.run({ properties: { name: 'clean-vm' } }, () => {
return this._adapter.cleanVm(this._vmBackupDir, {
...options,
fixMetadata: true,
@@ -100,7 +100,7 @@ export const MixinRemoteWriter = (BaseClass = Object) =>
}
return Task.run(
{
name: 'health check',
properties: { name: 'health check' },
},
async () => {
const xapi = sr.$xapi

View File

@@ -1,7 +1,7 @@
import assert from 'node:assert/strict'
import { Task } from '@vates/task'
import { HealthCheckVmBackup } from '../../HealthCheckVmBackup.mjs'
import { Task } from '../../Task.mjs'
import ms from 'ms'
export const MixinXapiWriter = (BaseClass = Object) =>
@@ -33,7 +33,7 @@ export const MixinXapiWriter = (BaseClass = Object) =>
// copy VM
return Task.run(
{
name: 'health check',
properties: { name: 'health check' },
},
async () => {
const { $xapi: xapi } = sr
@@ -47,12 +47,12 @@ export const MixinXapiWriter = (BaseClass = Object) =>
}
if (await this.#isAlreadyOnHealthCheckSr(baseVm)) {
healthCheckVmRef = await Task.run(
{ name: 'cloning-vm' },
{ properties: { name: 'cloning-vm' } },
async () => await xapi.callAsync('VM.clone', this._targetVmRef, `Health Check - ${baseVm.name_label}`)
)
} else {
healthCheckVmRef = await Task.run(
{ name: 'copying-vm' },
{ properties: { name: 'copying-vm' } },
async () =>
await xapi.callAsync(
'VM.copy',

View File

@@ -29,6 +29,7 @@
"@vates/generator-toolbox": "^1.1.1",
"@vates/nbd-client": "^3.3.0",
"@vates/parse-duration": "^0.1.1",
"@vates/task": "^0.6.0",
"@xen-orchestra/async-map": "^0.1.2",
"@xen-orchestra/disk-transform": "^1.2.2",
"@xen-orchestra/fs": "^4.7.0",

View File

@@ -2,7 +2,7 @@ import { createLogger } from '@xen-orchestra/log'
import { EventEmitter } from 'node:events'
import { makeOnProgress } from '@vates/task/combineEvents'
import { noSuchObject } from 'xo-common/api-errors.js'
import { Task } from '@vates/task'
import { Task, serializeError } from '@vates/task'
import iteratee from 'lodash/iteratee.js'
import stubTrue from 'lodash/stubTrue.js'
@@ -10,18 +10,20 @@ export { Task }
const { warn } = createLogger('xo:mixins:Tasks')
const formatId = timestamp => timestamp.toString(36).padStart(9, '0')
const noop = Function.prototype
// Create a serializable object from an error.
const serializeError = error => ({
...error, // Copy enumerable properties.
code: error.code,
message: error.message,
name: error.name,
stack: error.stack,
})
const DEFAULT_BACKUP_LOG_KEEP_DURATION = 31 * 24 * 60 * 60 * 1000
const getLogAge = (log, now) => {
if (log.end !== undefined) {
return now - log.end
}
if (log.start !== undefined) {
return now - log.start
}
// if we can't evaluate the log's age, assume it's too old
return now
}
export default class Tasks extends EventEmitter {
#logsToClearOnSuccess = new Set()
@@ -110,7 +112,13 @@ export default class Tasks extends EventEmitter {
this.#app = app
app.hooks
.on('clean', () => this.#gc(app.config.getOptional('tasks.gc.keep') ?? 1e3))
.on('clean', () =>
this.#gc({
keepRegularLogs: app.config.getOptional('tasks.gc.keep') ?? 1e3,
backupKeepDuration:
app.config.getOptionalDuration('tasks.gc.backupKeepDuration') ?? DEFAULT_BACKUP_LOG_KEEP_DURATION,
})
)
.on('start', async () => {
this.#store = await app.getStore('tasks')
@@ -124,10 +132,15 @@ export default class Tasks extends EventEmitter {
})
}
#gc(keep) {
/**
* keepRegularLogs: number of non-backup log entries to keep.
* backupKeepDuration: Max duration to keep backup logs (only backup logs)
*/
#gc({ keepRegularLogs, backupKeepDuration }) {
return new Promise((resolve, reject) => {
const db = this.#store
// used to prevent resolving the promise before all undesired entries are deleted
let count = 1
const cb = () => {
@@ -135,25 +148,31 @@ export default class Tasks extends EventEmitter {
resolve()
}
}
const stream = db.createKeyStream({
const deleteEntry = data => {
++count
this.deleteLog(data.key, cb)
}
const stream = db.createReadStream({
reverse: true,
})
const deleteEntry = key => {
++count
const now = Date.now()
this.deleteLog(key, cb)
const onData = data => {
if (data.value?.properties?.name === 'backup run') {
if (getLogAge(data.value, now) > backupKeepDuration) {
return deleteEntry(data)
}
} else {
if (--keepRegularLogs < 0) {
return deleteEntry(data)
}
}
}
const onData =
keep !== 0
? () => {
if (--keep === 0) {
stream.on('data', deleteEntry)
stream.removeListener('data', onData)
}
}
: deleteEntry
stream.on('data', onData)
stream.on('end', cb).on('error', reject)
@@ -169,7 +188,7 @@ export default class Tasks extends EventEmitter {
}
async clearLogs() {
await this.#gc(0)
await this.#gc({ keepRegularLogs: 0, backupKeepDuration: 0 })
}
/**
@@ -190,20 +209,9 @@ export default class Tasks extends EventEmitter {
const task = new Task({ properties: { ...props, name, objectId, userId, type }, onProgress: this.#onProgress })
// Use a compact, sortable, string representation of the creation date
//
// Due to the padding, dates are sortable up to 5188-04-22T11:04:28.415Z
let now = Date.now()
let id
while (tasks.has((id = formatId(now)))) {
// if the current id is already taken, use the next millisecond
++now
}
task.id = id
tasks.set(id, task)
tasks.set(task.id, task)
if (clearLogOnSuccess) {
this.#logsToClearOnSuccess.add(id)
this.#logsToClearOnSuccess.add(task.id)
}
return task

View File

@@ -15,8 +15,8 @@ import { Readable } from 'stream'
import { RemoteAdapter } from '@xen-orchestra/backups/RemoteAdapter.mjs'
import { RestoreMetadataBackup } from '@xen-orchestra/backups/RestoreMetadataBackup.mjs'
import { runBackupWorker } from '@xen-orchestra/backups/runBackupWorker.mjs'
import { Task } from '@xen-orchestra/backups/Task.mjs'
import { Xapi } from '@xen-orchestra/xapi'
import { Task } from '@vates/task'
const noop = Function.prototype
@@ -113,6 +113,7 @@ export default class Backups {
}
return run.apply(this, arguments)
})(run)
run = (run => async (params, onLog) => {
if (onLog === undefined || !app.config.get('backups').disableWorkers) {
return run(params, onLog)
@@ -122,15 +123,15 @@ export default class Backups {
try {
await Task.run(
{
name: 'backup run',
data: {
properties: {
jobId: job.id,
jobName: job.name,
mode: job.mode,
name: 'backup run',
reportWhen: job.settings['']?.reportWhen,
scheduleId: schedule.id,
},
onLog,
onProgress: onLog,
},
() => run(params)
)
@@ -205,14 +206,14 @@ export default class Backups {
async (args, onLog) =>
Task.run(
{
data: {
properties: {
backupId,
jobId: metadata.jobId,
name: 'restore',
srId: srUuid,
time: metadata.timestamp,
},
name: 'restore',
onLog,
onProgress: onLog,
},
run
).catch(() => {}), // errors are handled by logs,
@@ -354,9 +355,11 @@ export default class Backups {
async (args, onLog) =>
Task.run(
{
name: 'metadataRestore',
data: JSON.parse(String(await handler.readFile(`${backupId}/metadata.json`))),
onLog,
properties: {
name: 'metadataRestore',
metadata: JSON.parse(String(await handler.readFile(`${backupId}/metadata.json`))),
},
onProgress: onLog,
},
() =>
new RestoreMetadataBackup({
@@ -385,6 +388,7 @@ export default class Backups {
remotes: { type: 'object' },
schedule: { type: 'object' },
xapis: { type: 'object', optional: true },
jobData: { type: 'object', optional: true },
recordToXapi: { type: 'object', optional: true },
streamLogs: { type: 'boolean', optional: true },
},

View File

@@ -34,7 +34,8 @@
- [OpenMetrics] Add per-VDI disk size metrics: `xcp_vdi_virtual_size_bytes` and `xcp_vdi_physical_usage_bytes` (PR [#9680](https://github.com/vatesfr/xen-orchestra/pull/9680))
- [i18n] Update Chinese (Simplified Han script), Czech, Danish, Dutch, Finnish, German, Italian, Korean, Norwegian, Persian, Polish, Portuguese, Portuguese (Brasil), Russian, Slovak and Spanish translations (PR [#9649](https://github.com/vatesfr/xen-orchestra/pull/9649))
- [OpenMetrics] Add 9 missing host RRD metrics: `hostload`, `memory_reclaimed`, `memory_reclaimed_max`, `running_vcpus`, `pif_aggr_rx`, `pif_aggr_tx`, `iops_total`, `io_throughput_total`, `latency` per SR (PR [#9696](https://github.com/vatesfr/xen-orchestra/pull/9696))
- [Netbox] Use platform hierarchy to assign versioned OS names (e.g. "Debian 12" instead of "Debian") when the major version is known (requires Netbox >= 4.4) #7773 (PR [#9644](https://github.com/vatesfr/xen-orchestra/pull/9644))
- [Netbox] Use platform hierarchy to assign versioned OS names (e.g. "Debian 12" instead of "Debian") when the major version is known (requires Netbox >= 4.4) [#7773](https://github.com/vatesfr/xen-orchestra/issues/7773) (PR [#9644](https://github.com/vatesfr/xen-orchestra/pull/9644))
- [Backups] Backups no longer use their own task system, but instead use the same system as XO Task. This will help improve loading times in the future (PR [#9734](https://github.com/vatesfr/xen-orchestra/pull/9734))
### Bug fixes
@@ -67,6 +68,7 @@
- @vates/generator-toolbox patch
- @vates/nbd-client minor
- @vates/node-vsphere-soap patch
- @vates/task minor
- @vates/types minor
- @xen-orchestra/async-map patch
- @xen-orchestra/backups minor
@@ -79,8 +81,8 @@
- @xen-orchestra/log patch
- @xen-orchestra/mcp minor
- @xen-orchestra/mixin patch
- @xen-orchestra/mixins patch
- @xen-orchestra/proxy patch
- @xen-orchestra/mixins minor
- @xen-orchestra/proxy minor
- @xen-orchestra/qcow2 minor
- @xen-orchestra/rest-api minor
- @xen-orchestra/template patch

View File

@@ -44,6 +44,7 @@
"@vates/parse-duration": "^0.1.1",
"@vates/predicates": "^1.1.0",
"@vates/read-chunk": "^1.2.1",
"@vates/task": "^0.6.0",
"@vates/types": "^1.22.0",
"@vates/xml": "^2.0.0",
"@vates/xml-rpc": "^1.0.0",

View File

@@ -1,5 +1,6 @@
import humanFormat from 'human-format'
import ms from 'ms'
import { computeStatusAndSortTasks, getStatus } from './xo-mixins/backups-ng-logs.mjs'
import { createLogger } from '@xen-orchestra/log'
const { warn } = createLogger('xo:server:handleBackupLog')
@@ -31,85 +32,41 @@ async function sendToNagios(app, jobName, vmBackupInfo) {
}
}
function forwardResult(log) {
export function forwardResult(log) {
if (log.status === 'failure') {
throw log.result
}
return log.result
}
// it records logs generated by `@xen-orchestra/backups/Task#run`
export const handleBackupLog = (
log,
{ vmBackupInfo, app, jobName, logger, localTaskIds, rootTaskId, runJobId, handleRootTaskId }
) => {
const { event, message, parentId, taskId } = log
// it records logs generated by backup task logs
export const handleBackupLog = (taskLog, event, { app, jobName, store }) => {
if (event.type === 'end') {
taskLog.status = computeStatusAndSortTasks(getStatus(taskLog.result, taskLog.status), taskLog.tasks)
}
// sending data to Nagios
if (app !== undefined && jobName !== undefined) {
if (event === 'start') {
if (log.data?.type === 'VM') {
vmBackupInfo.set('vm-' + taskId, {
id: log.data.id,
start: log.timestamp,
})
} else if (vmBackupInfo.has('vm-' + parentId) && log.message === 'export') {
vmBackupInfo.set('export-' + taskId, {
parentId: 'vm-' + parentId,
})
} else if (vmBackupInfo.has('export-' + parentId) && log.message === 'transfer') {
vmBackupInfo.set('transfer-' + taskId, {
parentId: 'export-' + parentId,
})
}
} else if (event === 'end') {
if (vmBackupInfo.has('vm-' + taskId)) {
const data = vmBackupInfo.get('vm-' + taskId)
data.result = log.status
data.end = log.timestamp
sendToNagios(app, jobName, data)
} else if (vmBackupInfo.has('transfer-' + taskId)) {
vmBackupInfo.get(vmBackupInfo.get(vmBackupInfo.get('transfer-' + taskId).parentId).parentId).size =
log.result.size
if (event.type === 'end' && taskLog.properties?.type === 'VM') {
// we arbitrary pick one transfer to get the size
const exportTask = taskLog.tasks?.find(task => task.properties?.name === 'export')
const transferTask =
exportTask === undefined ? undefined : exportTask.tasks.find(task => task.properties?.name === 'transfer')
const vmBackupInfo = {
start: taskLog.start,
id: taskLog.properties?.id,
result: taskLog.status,
end: taskLog.end,
size: transferTask?.result?.size,
}
sendToNagios(app, jobName, vmBackupInfo)
}
}
// If `runJobId` is defined, it means that the root task is already handled by `runJob`
if (runJobId !== undefined) {
// Ignore the start of the root task
if (event === 'start' && log.parentId === undefined) {
localTaskIds[taskId] = runJobId
return
}
store.put(taskLog.$root.id, taskLog.$root)
// Return/throw the result of the root task
if (event === 'end' && localTaskIds[taskId] === runJobId) {
return forwardResult(log)
}
}
const common = {
data: log.data,
event: 'task.' + event,
result: log.result,
status: log.status,
}
if (event === 'start') {
const { parentId } = log
if (parentId === undefined) {
handleRootTaskId((localTaskIds[taskId] = logger.notice(message, common)))
} else {
common.parentId = localTaskIds[parentId]
localTaskIds[taskId] = logger.notice(message, common)
}
} else {
common.taskId = localTaskIds[taskId]
logger.notice(message, common)
}
// special case for the end of the root task: return/throw the result
if (event === 'end' && localTaskIds[taskId] === rootTaskId) {
return forwardResult(log)
// end of the root task: return/throw the result
if (event.type === 'end' && taskLog.$root === taskLog) {
return forwardResult(taskLog)
}
}

View File

@@ -1,4 +1,3 @@
import forEach from 'lodash/forEach.js'
import isEmpty from 'lodash/isEmpty.js'
import iteratee from 'lodash/iteratee.js'
import ms from 'ms'
@@ -13,10 +12,10 @@ const isSkippedError = error =>
error.message === 'no VMs match this pattern' ||
error.message === 'unhealthy VDI chain')
const getStatus = (error, status = error === undefined ? 'success' : 'failure') =>
export const getStatus = (error, status = error === undefined ? 'success' : 'failure') =>
status === 'failure' && isSkippedError(error) ? 'skipped' : status
const computeStatusAndSortTasks = (status, tasks) => {
export const computeStatusAndSortTasks = (status, tasks) => {
if (status === 'failure' || tasks === undefined) {
return status
}
@@ -52,6 +51,30 @@ const taskTimeComparator = ({ start: s1, end: e1 }, { start: s2, end: e2 }) => {
return 1
}
function adaptTask(task) {
const { name, ...data } = task.properties
task.message = name
if (Object.keys(data).length > 0) {
task.data = data
}
delete task.properties
task.tasks?.forEach(adaptTask)
}
function taskFormatAdapter(log) {
const isXoTask = log.tasks?.length > 0 && !!log.tasks[0].properties
if (isXoTask) {
adaptTask(log.tasks[0])
if (log.tasks[0].infos !== undefined && log.tasks[0].infos.length > 0) {
log.infos = [...(log.infos ?? []), ...log.tasks[0].infos]
}
if (log.tasks[0].warnings !== undefined && log.tasks[0].warnings.length !== 0) {
log.warnings = [...(log.warnings ?? []), ...log.tasks[0].warnings]
}
log.tasks = log.tasks[0].tasks
}
}
// type Task = {
// data: any,
// end?: number,
@@ -73,12 +96,13 @@ export default {
this.getLogs('restore'),
this.getLogs('metadataRestore'),
])
const taskStore = await this.getStore('tasks')
const { runningJobs, runningRestores, runningMetadataRestores } = this
const consolidated = {}
const started = {}
const handleLog = ({ data, time, message }, id) => {
const handleLog = async ({ data, time, message }, id) => {
const { event } = data
if (event === 'job.start') {
if ((data.type === 'backup' || data.key === undefined) && (runId === undefined || runId === id)) {
@@ -103,6 +127,13 @@ export default {
log.end = time
log.status = computeStatusAndSortTasks(getStatus((log.result = data.error)), log.tasks)
}
} else if (event === 'job.backupTaskStart') {
// happens once, only for backups using XO Tasks
const { runJobId } = data
const log = started[runJobId]
if (log !== undefined && (await taskStore.has(data.backupTaskId))) {
log.tasks = [await taskStore.get(data.backupTaskId)]
}
} else if (event === 'task.start') {
const task = {
data: data.data,
@@ -175,9 +206,19 @@ export default {
}
}
forEach(jobLogs, handleLog)
forEach(restoreLogs, handleLog)
forEach(restoreMetadataLogs, handleLog)
for (const [logId, log] of Object.entries(jobLogs)) {
await handleLog(log, logId)
}
for (const [logId, log] of Object.entries(restoreLogs)) {
await handleLog(log, logId)
}
for (const [logId, log] of Object.entries(restoreMetadataLogs)) {
await handleLog(log, logId)
}
for (const log of Object.values(consolidated)) {
taskFormatAdapter(log)
}
if (runId !== undefined) {
if (consolidated[runId] === undefined) {

View File

@@ -3,6 +3,7 @@ import Disposable from 'promise-toolbox/Disposable'
import forOwn from 'lodash/forOwn.js'
import groupBy from 'lodash/groupBy.js'
import merge from 'lodash/merge.js'
import { makeOnProgress } from '@vates/task/combineEvents'
import { createLogger } from '@xen-orchestra/log'
import { createPredicate } from 'value-matcher'
import { decorateWith } from '@vates/decorate-with'
@@ -12,14 +13,14 @@ import { ImportVmBackup } from '@xen-orchestra/backups/ImportVmBackup.mjs'
import { createRunner } from '@xen-orchestra/backups/Backup.mjs'
import { invalidParameters, noMatchingVm } from 'xo-common/api-errors.js'
import { runBackupWorker } from '@xen-orchestra/backups/runBackupWorker.mjs'
import { Task } from '@xen-orchestra/backups/Task.mjs'
import { Task } from '@vates/task'
import { debounceWithKey, REMOVE_CACHE_ENTRY } from '../../_pDebounceWithKey.mjs'
import { handleBackupLog } from '../../_handleBackupLog.mjs'
import { forwardResult, handleBackupLog } from '../../_handleBackupLog.mjs'
import { serializeError, unboxIdsFromPattern } from '../../utils.mjs'
import { waitAll } from '../../_waitAll.mjs'
const log = createLogger('xo:xo-mixins:backups-ng')
const logger = createLogger('xo:xo-mixins:backups-ng')
const parseVmBackupId = id => {
const i = id.indexOf('/')
@@ -65,16 +66,24 @@ export default class BackupNg {
constructor(app) {
this._app = app
this._logger = undefined
this._store = undefined
this._runningRestores = new Set()
app.hooks.on('start', async () => {
this._logger = await app.getLogger('restore')
this._store = await app.getStore('tasks')
const executor = async ({ cancelToken, data, job: job_, logger, runJobId, schedule }) => {
const executor = async ({
cancelToken,
data,
job,
jobData,
logger: jobsLogger,
runJobId,
schedule,
jobUpdateFct,
}) => {
const backupsConfig = app.config.get('backups')
let job = job_
let vmIds
if (job.type === 'backup') {
@@ -99,7 +108,7 @@ export default class BackupNg {
try {
app.getObject(id)
} catch (error) {
const taskId = logger.notice('missing pool', {
const taskId = jobsLogger.notice('missing pool', {
data: {
type: 'pool',
id,
@@ -107,7 +116,7 @@ export default class BackupNg {
event: 'task.start',
parentId: runJobId,
})
logger.error('missing pool', {
jobsLogger.error('missing pool', {
event: 'task.end',
result: serializeError(error),
status: 'failure',
@@ -153,20 +162,15 @@ export default class BackupNg {
const targetRemoteIds = unboxIdsFromPattern(job.remotes)
try {
if (!useXoProxy && backupsConfig.disableWorkers) {
const localTaskIds = { __proto__: null }
const vmBackupInfo = new Map()
const onLogFct = makeOnProgress({
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { app: this._app, jobName: job.name, store: this._store })
},
})
return await Task.run(
{
name: 'backup run',
onLog: log =>
handleBackupLog(log, {
vmBackupInfo,
app: this._app,
jobName: job.name,
localTaskIds,
logger,
runJobId,
}),
properties: { name: 'backup run', ...jobData },
onProgress: onLogFct,
},
() =>
createRunner({
@@ -190,7 +194,7 @@ export default class BackupNg {
recordToXapi[uuid] = serverId
servers.add(serverId)
} catch (error) {
log.warn(error)
logger.warn(error)
}
}
// can be empty for mirror backup job
@@ -215,7 +219,7 @@ export default class BackupNg {
try {
remote = await app.getRemoteWithCredentials(id)
} catch (error) {
log.warn('Error while instantiating remote', { error, remoteId: id })
logger.warn('Error while instantiating remote', { error, remoteId: id })
remoteErrors[id] = error
return
}
@@ -280,6 +284,7 @@ export default class BackupNg {
const params = {
job,
jobData,
recordToXapi,
remotes,
schedule,
@@ -300,18 +305,18 @@ export default class BackupNg {
}
)
const localTaskIds = { __proto__: null }
let result
const vmBackupInfo = new Map()
const onLogFct = makeOnProgress({
onRootTaskEnd: log => {
result = forwardResult(log)
},
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { app: this._app, jobName: job.name, store: this._store })
},
})
for await (const log of logsStream) {
result = handleBackupLog(log, {
vmBackupInfo,
app: this._app,
jobName: job.name,
logger,
localTaskIds,
runJobId,
})
onLogFct(log)
}
return result
} catch (error) {
@@ -322,26 +327,31 @@ export default class BackupNg {
throw error
}
} else {
const localTaskIds = { __proto__: null }
const vmBackupInfo = new Map()
return await runBackupWorker(
let result
const onLogFct = makeOnProgress({
onRootTaskStart: log => {
jobUpdateFct(log.id).catch(logger.warn) // is async, but makeOnProgress doesn't await onRootTaskXXX functions
},
onRootTaskEnd: log => {
result = forwardResult(log)
},
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { app: this._app, jobName: job.name, store: this._store })
},
})
await runBackupWorker(
{
config: backupsConfig,
jobData,
remoteOptions: app.config.get('remoteOptions'),
resourceCacheDelay: app.config.getDuration('resourceCacheDelay'),
xapiOptions: app.config.get('xapiOptions'),
...params,
},
log =>
handleBackupLog(log, {
vmBackupInfo,
app: this._app,
jobName: job.name,
logger,
localTaskIds,
runJobId,
})
onLogFct
)
return result
}
} finally {
targetRemoteIds.forEach(id => this._listVmBackupsOnRemote(REMOVE_CACHE_ENTRY, id))
@@ -349,6 +359,7 @@ export default class BackupNg {
}
app.registerJobExecutor('backup', executor)
app.registerJobExecutor('mirrorBackup', executor)
return () => this._store.close()
})
}
@@ -466,7 +477,6 @@ export default class BackupNg {
const remote = await app.getRemoteWithCredentials(remoteId)
let rootTaskId
const logger = this._logger
try {
let result
if (remote.proxy !== undefined) {
@@ -499,17 +509,21 @@ export default class BackupNg {
assertType: 'iterator',
})
const localTaskIds = { __proto__: null }
const onLogFct = makeOnProgress({
onRootTaskStart: log => {
this._runningRestores.add(log.id)
rootTaskId = log.id
},
onRootTaskEnd: log => {
result = forwardResult(log)
},
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { store: this._store })
},
})
for await (const log of logsStream) {
result = handleBackupLog(log, {
logger,
localTaskIds,
handleRootTaskId: id => {
this._runningRestores.add(id)
rootTaskId = id
},
rootTaskId,
})
onLogFct(log)
}
} catch (error) {
if (invalidParameters.is(error)) {
@@ -521,25 +535,27 @@ export default class BackupNg {
} else {
result = await Disposable.use(app.getBackupsRemoteAdapter(remote), async adapter => {
const metadata = await adapter.readVmBackupMetadata(metadataFilename)
const localTaskIds = { __proto__: null }
const onLogFct = makeOnProgress({
onRootTaskStart: log => {
this._runningRestores.add(log.id)
rootTaskId = log.id
},
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { store: this._store })
},
})
return Task.run(
{
data: {
properties: {
name: 'restore',
backupId: id,
jobId: metadata.jobId,
srId,
time: metadata.timestamp,
},
name: 'restore',
onLog: log =>
handleBackupLog(log, {
logger,
localTaskIds,
handleRootTaskId: id => {
this._runningRestores.add(id)
rootTaskId = id
},
}),
onProgress: onLogFct,
},
async () =>
new ImportVmBackup({
@@ -604,7 +620,7 @@ export default class BackupNg {
)
return backupsByVm
} catch (error) {
log.warn(`listVmBackups for remote ${remoteId}:`, { error })
logger.warn(`listVmBackups for remote ${remoteId}:`, { error })
}
}

View File

@@ -169,7 +169,7 @@ export default class Jobs {
const logger = this._logger
const { id, type } = job
const runJobId = logger.notice(`Starting execution of ${id}.`, {
const jobData = {
data:
type === 'backup' || type === 'metadataBackup' || type === 'mirrorBackup'
? {
@@ -187,7 +187,11 @@ export default class Jobs {
scheduleId: schedule?.id,
key: job.key,
type,
})
}
const runJobId = logger.notice(`Starting execution of ${id}.`, jobData)
// Links the backup log to the job run
// We keep the jobs for this because of some mechanism related to jobs, like preventing double execution.
jobData.runJobId = runJobId
const app = this._app
try {
@@ -204,6 +208,7 @@ export default class Jobs {
// runId is a temporary property used to check if the report is sent after the server interruption
this.updateJob({ id, runId: runJobId })::ignoreErrors()
runningJobs[id] = runJobId
$defer(() => {
@@ -277,12 +282,23 @@ export default class Jobs {
runs[runJobId] = { cancel }
$defer(() => delete runs[runJobId])
// Links the job run to its backup log
const jobUpdateFct = async backupTaskId => {
await logger.notice(`Adding backupTaskId to job run ${runJobId}`, {
backupTaskId,
event: 'job.backupTaskStart',
runJobId,
})
}
await executor({
app,
cancelToken: token,
connection,
data: data_,
job,
jobUpdateFct,
jobData,
logger,
runJobId,
schedule,

View File

@@ -3,16 +3,17 @@ import cloneDeep from 'lodash/cloneDeep.js'
import Disposable from 'promise-toolbox/Disposable'
import { createLogger } from '@xen-orchestra/log'
import { createRunner } from '@xen-orchestra/backups/Backup.mjs'
import { makeOnProgress } from '@vates/task/combineEvents'
import { parseMetadataBackupId } from '@xen-orchestra/backups/parseMetadataBackupId.mjs'
import { RestoreMetadataBackup } from '@xen-orchestra/backups/RestoreMetadataBackup.mjs'
import { Task } from '@xen-orchestra/backups/Task.mjs'
import { Task } from '@vates/task'
import { debounceWithKey, REMOVE_CACHE_ENTRY } from '../_pDebounceWithKey.mjs'
import { handleBackupLog } from '../_handleBackupLog.mjs'
import { forwardResult, handleBackupLog } from '../_handleBackupLog.mjs'
import { waitAll } from '../_waitAll.mjs'
import { serializeError, unboxIdsFromPattern } from '../utils.mjs'
const log = createLogger('xo:xo-mixins:metadata-backups')
const logger = createLogger('xo:xo-mixins:metadata-backups')
const METADATA_BACKUP_JOB_TYPE = 'metadataBackup'
@@ -23,7 +24,7 @@ export default class metadataBackup {
constructor(app) {
this._app = app
this._logger = undefined
this._store = undefined
this._runningMetadataRestores = new Set()
const debounceDelay = app.config.getDuration('backups.listingDebounce')
@@ -31,13 +32,14 @@ export default class metadataBackup {
this._listPoolMetadataBackups = debounceWithKey(this._listPoolMetadataBackups, debounceDelay, remoteId => remoteId)
app.hooks.on('start', async () => {
this._logger = await app.getLogger('metadataRestore')
this._store = await app.getStore('tasks')
app.registerJobExecutor(METADATA_BACKUP_JOB_TYPE, this._executor.bind(this))
return () => this._store.close()
})
}
async _executor({ cancelToken, job: job_, logger, runJobId, schedule }) {
async _executor({ cancelToken, job: job_, jobData, logger: jobsLogger, runJobId, schedule }) {
const job = cloneDeep(job_)
const scheduleSettings = job.settings[schedule.id]
@@ -96,6 +98,7 @@ export default class metadataBackup {
const params = {
job,
jobData,
recordToXapi,
remotes,
schedule,
@@ -107,30 +110,32 @@ export default class metadataBackup {
assertType: 'iterator',
})
const localTaskIds = { __proto__: null }
let result
const onLogFct = makeOnProgress({
onRootTaskEnd: log => {
result = forwardResult(log)
},
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { store: this._store })
},
})
for await (const log of logsStream) {
result = handleBackupLog(log, {
localTaskIds,
logger,
runJobId,
})
onLogFct(log)
}
return result
} else {
cancelToken.throwIfRequested()
const localTaskIds = { __proto__: null }
const onLogFct = makeOnProgress({
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { store: this._store })
},
})
return Task.run(
{
name: 'backup run',
onLog: log =>
handleBackupLog(log, {
localTaskIds,
logger,
runJobId,
}),
properties: { name: 'backup run', ...jobData },
onProgress: onLogFct,
},
() =>
createRunner({
@@ -313,7 +318,7 @@ export default class metadataBackup {
pool[remoteId] = poolList
}
} catch (error) {
log.warn(`listMetadataBackups for remote ${remoteId}`, { error })
logger.warn(`listMetadataBackups for remote ${remoteId}`, { error })
}
})
)
@@ -330,7 +335,6 @@ export default class metadataBackup {
// └─ task.end
async restoreMetadataBackup({ id, poolUuid }) {
const app = this._app
const logger = this._logger
const [remoteId, ...path] = id.split('/')
const backupId = path.join('/')
@@ -343,29 +347,30 @@ export default class metadataBackup {
}
let rootTaskId
const localTaskIds = { __proto__: null }
const onLog = async log => {
if (type === 'xoConfig' && localTaskIds[log.taskId] === rootTaskId && log.status === 'success') {
const onProgressFct = makeOnProgress({
onRootTaskStart: log => {
this._runningMetadataRestores.add(log.id)
rootTaskId = log.id
},
onTaskUpdate: (log, event) => {
handleBackupLog(log, event, { store: this._store })
},
})
const onLogFct = async event => {
if (type === 'xoConfig' && event.status === 'success' && event.parentId === undefined) {
try {
const { result } = log
const { result } = event
await app.importConfig(typeof result === 'string' ? result : Buffer.from(result.data, result.encoding))
// don't log the XO config
log.result = undefined
event.result = undefined
} catch (error) {
log.result = serializeError(error)
log.status = 'failure'
event.result = serializeError(error)
event.status = 'failure'
}
}
handleBackupLog(log, {
logger,
localTaskIds,
handleRootTaskId: id => {
this._runningMetadataRestores.add(id)
rootTaskId = id
},
})
onProgressFct(event)
}
try {
@@ -398,15 +403,19 @@ export default class metadataBackup {
}
)
for await (const log of logsStream) {
onLog(log)
await onLogFct(log)
}
} else {
const handler = await app.getRemoteHandler(remoteId)
await Task.run(
{
name: 'metadataRestore',
data: JSON.parse(String(await handler.readFile(`${backupId}/metadata.json`))),
onLog,
properties: {
name: 'metadataRestore',
metadata: JSON.parse(String(await handler.readFile(`${backupId}/metadata.json`))),
},
onProgress: event => {
onLogFct(event).catch(logger.warn)
},
},
async () =>
new RestoreMetadataBackup({

2400
yarn.lock

File diff suppressed because it is too large Load Diff