Compare commits

..

2 Commits

Author SHA1 Message Date
Florent BEAUCHAMP
d4ca0da8b9 fix(backups/cleanVm): vhd size computation 2023-07-05 16:15:02 +02:00
Florent BEAUCHAMP
486d50f2f1 feat(@xen-orchestra/backups): store more detailled sizes 2023-07-05 14:38:17 +02:00
539 changed files with 9062 additions and 18385 deletions

View File

@@ -1,11 +1,8 @@
'use strict'
module.exports = { module.exports = {
arrowParens: 'avoid', arrowParens: 'avoid',
jsxSingleQuote: true, jsxSingleQuote: true,
semi: false, semi: false,
singleQuote: true, singleQuote: true,
trailingComma: 'es5',
// 2020-11-24: Requested by nraynaud and approved by the rest of the team // 2020-11-24: Requested by nraynaud and approved by the rest of the team
// //

View File

@@ -33,7 +33,7 @@
"test": "node--test" "test": "node--test"
}, },
"devDependencies": { "devDependencies": {
"sinon": "^16.0.0", "sinon": "^15.0.1",
"tap": "^16.3.0", "tap": "^16.3.0",
"test": "^3.2.1" "test": "^3.2.1"
} }

View File

@@ -13,15 +13,12 @@ describe('decorateWith', () => {
const expectedFn = Function.prototype const expectedFn = Function.prototype
const newFn = () => {} const newFn = () => {}
const decorator = decorateWith( const decorator = decorateWith(function wrapper(fn, ...args) {
function wrapper(fn, ...args) {
assert.deepStrictEqual(fn, expectedFn) assert.deepStrictEqual(fn, expectedFn)
assert.deepStrictEqual(args, expectedArgs) assert.deepStrictEqual(args, expectedArgs)
return newFn return newFn
}, }, ...expectedArgs)
...expectedArgs
)
const descriptor = { const descriptor = {
configurable: true, configurable: true,

View File

@@ -29,7 +29,7 @@
"ensure-array": "^1.0.0" "ensure-array": "^1.0.0"
}, },
"devDependencies": { "devDependencies": {
"sinon": "^16.0.0", "sinon": "^15.0.1",
"test": "^3.2.1" "test": "^3.2.1"
} }
} }

View File

@@ -1,7 +1,9 @@
import LRU from 'lru-cache' 'use strict'
import Fuse from 'fuse-native'
import { VhdSynthetic } from 'vhd-lib' const LRU = require('lru-cache')
import { Disposable, fromCallback } from 'promise-toolbox' const Fuse = require('fuse-native')
const { VhdSynthetic } = require('vhd-lib')
const { Disposable, fromCallback } = require('promise-toolbox')
// build a s stat object from https://github.com/fuse-friends/fuse-native/blob/master/test/fixtures/stat.js // build a s stat object from https://github.com/fuse-friends/fuse-native/blob/master/test/fixtures/stat.js
const stat = st => ({ const stat = st => ({
@@ -14,7 +16,7 @@ const stat = st => ({
gid: st.gid !== undefined ? st.gid : process.getgid(), gid: st.gid !== undefined ? st.gid : process.getgid(),
}) })
export const mount = Disposable.factory(async function* mount(handler, diskPath, mountDir) { exports.mount = Disposable.factory(async function* mount(handler, diskPath, mountDir) {
const vhd = yield VhdSynthetic.fromVhdChain(handler, diskPath) const vhd = yield VhdSynthetic.fromVhdChain(handler, diskPath)
const cache = new LRU({ const cache = new LRU({

View File

@@ -1,6 +1,6 @@
{ {
"name": "@vates/fuse-vhd", "name": "@vates/fuse-vhd",
"version": "2.0.0", "version": "1.0.0",
"license": "ISC", "license": "ISC",
"private": false, "private": false,
"homepage": "https://github.com/vatesfr/xen-orchestra/tree/master/@vates/fuse-vhd", "homepage": "https://github.com/vatesfr/xen-orchestra/tree/master/@vates/fuse-vhd",
@@ -15,14 +15,13 @@
"url": "https://vates.fr" "url": "https://vates.fr"
}, },
"engines": { "engines": {
"node": ">=14" "node": ">=10.0"
}, },
"main": "./index.mjs",
"dependencies": { "dependencies": {
"fuse-native": "^2.2.6", "fuse-native": "^2.2.6",
"lru-cache": "^7.14.0", "lru-cache": "^7.14.0",
"promise-toolbox": "^0.21.0", "promise-toolbox": "^0.21.0",
"vhd-lib": "^4.6.1" "vhd-lib": "^4.5.0"
}, },
"scripts": { "scripts": {
"postversion": "npm publish --access public" "postversion": "npm publish --access public"

View File

@@ -0,0 +1,42 @@
'use strict'
exports.INIT_PASSWD = Buffer.from('NBDMAGIC') // "NBDMAGIC" ensure we're connected to a nbd server
exports.OPTS_MAGIC = Buffer.from('IHAVEOPT') // "IHAVEOPT" start an option block
exports.NBD_OPT_REPLY_MAGIC = 1100100111001001n // magic received during negociation
exports.NBD_OPT_EXPORT_NAME = 1
exports.NBD_OPT_ABORT = 2
exports.NBD_OPT_LIST = 3
exports.NBD_OPT_STARTTLS = 5
exports.NBD_OPT_INFO = 6
exports.NBD_OPT_GO = 7
exports.NBD_FLAG_HAS_FLAGS = 1 << 0
exports.NBD_FLAG_READ_ONLY = 1 << 1
exports.NBD_FLAG_SEND_FLUSH = 1 << 2
exports.NBD_FLAG_SEND_FUA = 1 << 3
exports.NBD_FLAG_ROTATIONAL = 1 << 4
exports.NBD_FLAG_SEND_TRIM = 1 << 5
exports.NBD_FLAG_FIXED_NEWSTYLE = 1 << 0
exports.NBD_CMD_FLAG_FUA = 1 << 0
exports.NBD_CMD_FLAG_NO_HOLE = 1 << 1
exports.NBD_CMD_FLAG_DF = 1 << 2
exports.NBD_CMD_FLAG_REQ_ONE = 1 << 3
exports.NBD_CMD_FLAG_FAST_ZERO = 1 << 4
exports.NBD_CMD_READ = 0
exports.NBD_CMD_WRITE = 1
exports.NBD_CMD_DISC = 2
exports.NBD_CMD_FLUSH = 3
exports.NBD_CMD_TRIM = 4
exports.NBD_CMD_CACHE = 5
exports.NBD_CMD_WRITE_ZEROES = 6
exports.NBD_CMD_BLOCK_STATUS = 7
exports.NBD_CMD_RESIZE = 8
exports.NBD_REQUEST_MAGIC = 0x25609513 // magic number to create a new NBD request to send to the server
exports.NBD_REPLY_MAGIC = 0x67446698 // magic number received from the server when reading response to a nbd request
exports.NBD_REPLY_ACK = 1
exports.NBD_DEFAULT_PORT = 10809
exports.NBD_DEFAULT_BLOCK_SIZE = 64 * 1024

View File

@@ -1,41 +0,0 @@
export const INIT_PASSWD = Buffer.from('NBDMAGIC') // "NBDMAGIC" ensure we're connected to a nbd server
export const OPTS_MAGIC = Buffer.from('IHAVEOPT') // "IHAVEOPT" start an option block
export const NBD_OPT_REPLY_MAGIC = 1100100111001001n // magic received during negociation
export const NBD_OPT_EXPORT_NAME = 1
export const NBD_OPT_ABORT = 2
export const NBD_OPT_LIST = 3
export const NBD_OPT_STARTTLS = 5
export const NBD_OPT_INFO = 6
export const NBD_OPT_GO = 7
export const NBD_FLAG_HAS_FLAGS = 1 << 0
export const NBD_FLAG_READ_ONLY = 1 << 1
export const NBD_FLAG_SEND_FLUSH = 1 << 2
export const NBD_FLAG_SEND_FUA = 1 << 3
export const NBD_FLAG_ROTATIONAL = 1 << 4
export const NBD_FLAG_SEND_TRIM = 1 << 5
export const NBD_FLAG_FIXED_NEWSTYLE = 1 << 0
export const NBD_CMD_FLAG_FUA = 1 << 0
export const NBD_CMD_FLAG_NO_HOLE = 1 << 1
export const NBD_CMD_FLAG_DF = 1 << 2
export const NBD_CMD_FLAG_REQ_ONE = 1 << 3
export const NBD_CMD_FLAG_FAST_ZERO = 1 << 4
export const NBD_CMD_READ = 0
export const NBD_CMD_WRITE = 1
export const NBD_CMD_DISC = 2
export const NBD_CMD_FLUSH = 3
export const NBD_CMD_TRIM = 4
export const NBD_CMD_CACHE = 5
export const NBD_CMD_WRITE_ZEROES = 6
export const NBD_CMD_BLOCK_STATUS = 7
export const NBD_CMD_RESIZE = 8
export const NBD_REQUEST_MAGIC = 0x25609513 // magic number to create a new NBD request to send to the server
export const NBD_REPLY_MAGIC = 0x67446698 // magic number received from the server when reading response to a nbd request
export const NBD_REPLY_ACK = 1
export const NBD_DEFAULT_PORT = 10809
export const NBD_DEFAULT_BLOCK_SIZE = 64 * 1024

View File

@@ -1,11 +1,8 @@
import assert from 'node:assert' 'use strict'
import { Socket } from 'node:net' const assert = require('node:assert')
import { connect } from 'node:tls' const { Socket } = require('node:net')
import { fromCallback, pRetry, pDelay, pTimeout } from 'promise-toolbox' const { connect } = require('node:tls')
import { readChunkStrict } from '@vates/read-chunk' const {
import { createLogger } from '@xen-orchestra/log'
import {
INIT_PASSWD, INIT_PASSWD,
NBD_CMD_READ, NBD_CMD_READ,
NBD_DEFAULT_BLOCK_SIZE, NBD_DEFAULT_BLOCK_SIZE,
@@ -20,13 +17,16 @@ import {
NBD_REQUEST_MAGIC, NBD_REQUEST_MAGIC,
OPTS_MAGIC, OPTS_MAGIC,
NBD_CMD_DISC, NBD_CMD_DISC,
} from './constants.mjs' } = require('./constants.js')
const { fromCallback, pRetry, pDelay, pTimeout } = require('promise-toolbox')
const { readChunkStrict } = require('@vates/read-chunk')
const { createLogger } = require('@xen-orchestra/log')
const { warn } = createLogger('vates:nbd-client') const { warn } = createLogger('vates:nbd-client')
// documentation is here : https://github.com/NetworkBlockDevice/nbd/blob/master/doc/proto.md // documentation is here : https://github.com/NetworkBlockDevice/nbd/blob/master/doc/proto.md
export default class NbdClient { module.exports = class NbdClient {
#serverAddress #serverAddress
#serverCert #serverCert
#serverPort #serverPort

View File

@@ -13,18 +13,17 @@
"url": "https://vates.fr" "url": "https://vates.fr"
}, },
"license": "ISC", "license": "ISC",
"version": "2.0.0", "version": "1.2.1",
"engines": { "engines": {
"node": ">=14.0" "node": ">=14.0"
}, },
"main": "./index.mjs",
"dependencies": { "dependencies": {
"@vates/async-each": "^1.0.0", "@vates/async-each": "^1.0.0",
"@vates/read-chunk": "^1.2.0", "@vates/read-chunk": "^1.1.1",
"@xen-orchestra/async-map": "^0.1.2", "@xen-orchestra/async-map": "^0.1.2",
"@xen-orchestra/log": "^0.6.0", "@xen-orchestra/log": "^0.6.0",
"promise-toolbox": "^0.21.0", "promise-toolbox": "^0.21.0",
"xen-api": "^1.3.6" "xen-api": "^1.3.3"
}, },
"devDependencies": { "devDependencies": {
"tap": "^16.3.0", "tap": "^16.3.0",
@@ -32,6 +31,6 @@
}, },
"scripts": { "scripts": {
"postversion": "npm publish --access public", "postversion": "npm publish --access public",
"test-integration": "tap --lines 97 --functions 95 --branches 74 --statements 97 tests/*.integ.mjs" "test-integration": "tap --lines 97 --functions 95 --branches 74 --statements 97 tests/*.integ.js"
} }
} }

View File

@@ -1,12 +1,13 @@
import NbdClient from '../index.mjs' 'use strict'
import { spawn, exec } from 'node:child_process' const NbdClient = require('../index.js')
import fs from 'node:fs/promises' const { spawn, exec } = require('node:child_process')
import { test } from 'tap' const fs = require('node:fs/promises')
import tmp from 'tmp' const { test } = require('tap')
import { pFromCallback } from 'promise-toolbox' const tmp = require('tmp')
import { Socket } from 'node:net' const { pFromCallback } = require('promise-toolbox')
import { NBD_DEFAULT_PORT } from '../constants.mjs' const { Socket } = require('node:net')
import assert from 'node:assert' const { NBD_DEFAULT_PORT } = require('../constants.js')
const assert = require('node:assert')
const FILE_SIZE = 10 * 1024 * 1024 const FILE_SIZE = 10 * 1024 * 1024

View File

@@ -1,3 +1,4 @@
'use strict'
/* /*
node-vsphere-soap node-vsphere-soap
@@ -11,18 +12,17 @@
*/ */
import { EventEmitter } from 'events' const EventEmitter = require('events').EventEmitter
import axios from 'axios' const axios = require('axios')
import https from 'node:https' const https = require('node:https')
import util from 'util' const util = require('util')
import soap from 'soap' const soap = require('soap')
import Cookie from 'soap-cookie' // required for session persistence const Cookie = require('soap-cookie') // required for session persistence
// Client class // Client class
// inherits from EventEmitter // inherits from EventEmitter
// possible events: connect, error, ready // possible events: connect, error, ready
export function Client(vCenterHostname, username, password, sslVerify) { function Client(vCenterHostname, username, password, sslVerify) {
this.status = 'disconnected' this.status = 'disconnected'
this.reconnectCount = 0 this.reconnectCount = 0
@@ -228,3 +228,4 @@ function _soapErrorHandler(self, emitter, command, args, err) {
} }
// end // end
exports.Client = Client

View File

@@ -1,8 +1,8 @@
{ {
"name": "@vates/node-vsphere-soap", "name": "@vates/node-vsphere-soap",
"version": "2.0.0", "version": "1.0.0",
"description": "interface to vSphere SOAP/WSDL from node for interfacing with vCenter or ESXi, forked from node-vsphere-soap", "description": "interface to vSphere SOAP/WSDL from node for interfacing with vCenter or ESXi, forked from node-vsphere-soap",
"main": "lib/client.mjs", "main": "lib/client.js",
"author": "reedog117", "author": "reedog117",
"repository": { "repository": {
"directory": "@vates/node-vsphere-soap", "directory": "@vates/node-vsphere-soap",
@@ -30,7 +30,7 @@
"private": false, "private": false,
"homepage": "https://github.com/vatesfr/xen-orchestra/tree/master/@vates/node-vsphere-soap", "homepage": "https://github.com/vatesfr/xen-orchestra/tree/master/@vates/node-vsphere-soap",
"engines": { "engines": {
"node": ">=14" "node": ">=8.10"
}, },
"scripts": { "scripts": {
"postversion": "npm publish --access public" "postversion": "npm publish --access public"

View File

@@ -1,11 +1,15 @@
'use strict'
// place your own credentials here for a vCenter or ESXi server // place your own credentials here for a vCenter or ESXi server
// this information will be used for connecting to a vCenter instance // this information will be used for connecting to a vCenter instance
// for module testing // for module testing
// name the file config-test.js // name the file config-test.js
export const vCenterTestCreds = { const vCenterTestCreds = {
vCenterIP: 'vcsa', vCenterIP: 'vcsa',
vCenterUser: 'vcuser', vCenterUser: 'vcuser',
vCenterPassword: 'vcpw', vCenterPassword: 'vcpw',
vCenter: true, vCenter: true,
} }
exports.vCenterTestCreds = vCenterTestCreds

View File

@@ -1,16 +1,18 @@
'use strict'
/* /*
vsphere-soap.test.js vsphere-soap.test.js
tests for the vCenterConnectionInstance class tests for the vCenterConnectionInstance class
*/ */
import assert from 'assert' const assert = require('assert')
import { describe, it } from 'test' const { describe, it } = require('test')
import * as vc from '../lib/client.mjs' const vc = require('../lib/client')
// eslint-disable-next-line n/no-missing-import // eslint-disable-next-line n/no-missing-require
import { vCenterTestCreds as TestCreds } from '../config-test.mjs' const TestCreds = require('../config-test.js').vCenterTestCreds
const VItest = new vc.Client(TestCreds.vCenterIP, TestCreds.vCenterUser, TestCreds.vCenterPassword, false) const VItest = new vc.Client(TestCreds.vCenterIP, TestCreds.vCenterUser, TestCreds.vCenterPassword, false)

View File

@@ -1,7 +1,6 @@
'use strict' 'use strict'
const assert = require('assert') const assert = require('assert')
const isUtf8 = require('isutf8')
/** /**
* Read a chunk of data from a stream. * Read a chunk of data from a stream.
@@ -82,13 +81,6 @@ exports.readChunkStrict = async function readChunkStrict(stream, size) {
if (size !== undefined && chunk.length !== size) { if (size !== undefined && chunk.length !== size) {
const error = new Error(`stream has ended with not enough data (actual: ${chunk.length}, expected: ${size})`) const error = new Error(`stream has ended with not enough data (actual: ${chunk.length}, expected: ${size})`)
// Buffer.isUtf8 is too recent for now
// @todo : replace external package by Buffer.isUtf8 when the supported version of node reach 18
if (chunk.length < 1024 && isUtf8(chunk)) {
error.text = chunk.toString('utf8')
}
Object.defineProperties(error, { Object.defineProperties(error, {
chunk: { chunk: {
value: chunk, value: chunk,

View File

@@ -102,37 +102,12 @@ describe('readChunkStrict', function () {
assert.strictEqual(error.chunk, undefined) assert.strictEqual(error.chunk, undefined)
}) })
it('throws if stream ends with not enough data, utf8', async () => { it('throws if stream ends with not enough data', async () => {
const error = await rejectionOf(readChunkStrict(makeStream(['foo', 'bar']), 10)) const error = await rejectionOf(readChunkStrict(makeStream(['foo', 'bar']), 10))
assert(error instanceof Error) assert(error instanceof Error)
assert.strictEqual(error.message, 'stream has ended with not enough data (actual: 6, expected: 10)') assert.strictEqual(error.message, 'stream has ended with not enough data (actual: 6, expected: 10)')
assert.strictEqual(error.text, 'foobar')
assert.deepEqual(error.chunk, Buffer.from('foobar')) assert.deepEqual(error.chunk, Buffer.from('foobar'))
}) })
it('throws if stream ends with not enough data, non utf8 ', async () => {
const source = [Buffer.alloc(10, 128), Buffer.alloc(10, 128)]
const error = await rejectionOf(readChunkStrict(makeStream(source), 30))
assert(error instanceof Error)
assert.strictEqual(error.message, 'stream has ended with not enough data (actual: 20, expected: 30)')
assert.strictEqual(error.text, undefined)
assert.deepEqual(error.chunk, Buffer.concat(source))
})
it('throws if stream ends with not enough data, utf8 , long data', async () => {
const source = Buffer.from('a'.repeat(1500))
const error = await rejectionOf(readChunkStrict(makeStream([source]), 2000))
assert(error instanceof Error)
assert.strictEqual(error.message, `stream has ended with not enough data (actual: 1500, expected: 2000)`)
assert.strictEqual(error.text, undefined)
assert.deepEqual(error.chunk, source)
})
it('succeed', async () => {
const source = Buffer.from('a'.repeat(20))
const chunk = await readChunkStrict(makeStream([source]), 10)
assert.deepEqual(source.subarray(10), chunk)
})
}) })
describe('skip', function () { describe('skip', function () {
@@ -159,16 +134,6 @@ describe('skip', function () {
it('returns less size if stream ends', async () => { it('returns less size if stream ends', async () => {
assert.deepEqual(await skip(makeStream('foo bar'), 10), 7) assert.deepEqual(await skip(makeStream('foo bar'), 10), 7)
}) })
it('put back if it read too much', async () => {
let source = makeStream(['foo', 'bar'])
await skip(source, 1) // read part of data chunk
const chunk = (await readChunkStrict(source, 2)).toString('utf-8')
assert.strictEqual(chunk, 'oo')
source = makeStream(['foo', 'bar'])
assert.strictEqual(await skip(source, 3), 3) // read aligned with data chunk
})
}) })
describe('skipStrict', function () { describe('skipStrict', function () {
@@ -179,9 +144,4 @@ describe('skipStrict', function () {
assert.strictEqual(error.message, 'stream has ended with not enough data (actual: 7, expected: 10)') assert.strictEqual(error.message, 'stream has ended with not enough data (actual: 7, expected: 10)')
assert.deepEqual(error.bytesSkipped, 7) assert.deepEqual(error.bytesSkipped, 7)
}) })
it('succeed', async () => {
const source = makeStream(['foo', 'bar', 'baz'])
const res = await skipStrict(source, 4)
assert.strictEqual(res, undefined)
})
}) })

View File

@@ -19,7 +19,7 @@
"type": "git", "type": "git",
"url": "https://github.com/vatesfr/xen-orchestra.git" "url": "https://github.com/vatesfr/xen-orchestra.git"
}, },
"version": "1.2.0", "version": "1.1.1",
"engines": { "engines": {
"node": ">=8.10" "node": ">=8.10"
}, },
@@ -33,8 +33,5 @@
}, },
"devDependencies": { "devDependencies": {
"test": "^3.2.1" "test": "^3.2.1"
},
"dependencies": {
"isutf8": "^4.0.0"
} }
} }

View File

@@ -27,7 +27,7 @@
"license": "ISC", "license": "ISC",
"version": "0.1.0", "version": "0.1.0",
"engines": { "engines": {
"node": ">=12.3" "node": ">=10"
}, },
"scripts": { "scripts": {
"postversion": "npm publish --access public", "postversion": "npm publish --access public",

View File

@@ -35,7 +35,7 @@
"test": "node--test" "test": "node--test"
}, },
"devDependencies": { "devDependencies": {
"sinon": "^16.0.0", "sinon": "^15.0.1",
"test": "^3.2.1" "test": "^3.2.1"
} }
} }

View File

@@ -1,5 +1,5 @@
import { asyncMap } from '@xen-orchestra/async-map' import { asyncMap } from '@xen-orchestra/async-map'
import { RemoteAdapter } from '@xen-orchestra/backups/RemoteAdapter.mjs' import { RemoteAdapter } from '@xen-orchestra/backups/RemoteAdapter.js'
import { getSyncedHandler } from '@xen-orchestra/fs' import { getSyncedHandler } from '@xen-orchestra/fs'
import getopts from 'getopts' import getopts from 'getopts'
import { basename, dirname } from 'path' import { basename, dirname } from 'path'

View File

@@ -7,9 +7,9 @@
"bugs": "https://github.com/vatesfr/xen-orchestra/issues", "bugs": "https://github.com/vatesfr/xen-orchestra/issues",
"dependencies": { "dependencies": {
"@xen-orchestra/async-map": "^0.1.2", "@xen-orchestra/async-map": "^0.1.2",
"@xen-orchestra/backups": "^0.43.0", "@xen-orchestra/backups": "^0.39.0",
"@xen-orchestra/fs": "^4.1.0", "@xen-orchestra/fs": "^4.0.1",
"filenamify": "^6.0.0", "filenamify": "^4.1.0",
"getopts": "^2.2.5", "getopts": "^2.2.5",
"lodash": "^4.17.15", "lodash": "^4.17.15",
"promise-toolbox": "^0.21.0" "promise-toolbox": "^0.21.0"
@@ -27,7 +27,7 @@
"scripts": { "scripts": {
"postversion": "npm publish --access public" "postversion": "npm publish --access public"
}, },
"version": "1.0.13", "version": "1.0.9",
"license": "AGPL-3.0-or-later", "license": "AGPL-3.0-or-later",
"author": { "author": {
"name": "Vates SAS", "name": "Vates SAS",

View File

@@ -1,8 +1,10 @@
import { Metadata } from './_runners/Metadata.mjs' 'use strict'
import { VmsRemote } from './_runners/VmsRemote.mjs'
import { VmsXapi } from './_runners/VmsXapi.mjs'
export function createRunner(opts) { const { Metadata } = require('./_runners/Metadata.js')
const { VmsRemote } = require('./_runners/VmsRemote.js')
const { VmsXapi } = require('./_runners/VmsXapi.js')
exports.createRunner = function createRunner(opts) {
const { type } = opts.job const { type } = opts.job
switch (type) { switch (type) {
case 'backup': case 'backup':

View File

@@ -1,6 +1,8 @@
import { asyncMap } from '@xen-orchestra/async-map' 'use strict'
export class DurablePartition { const { asyncMap } = require('@xen-orchestra/async-map')
exports.DurablePartition = class DurablePartition {
// private resource API is used exceptionally to be able to separate resource creation and release // private resource API is used exceptionally to be able to separate resource creation and release
#partitionDisposers = {} #partitionDisposers = {}

View File

@@ -1,6 +1,8 @@
import { Task } from './Task.mjs' 'use strict'
export class HealthCheckVmBackup { const { Task } = require('./Task')
exports.HealthCheckVmBackup = class HealthCheckVmBackup {
#restoredVm #restoredVm
#timeout #timeout
#xapi #xapi

View File

@@ -1,11 +1,13 @@
import assert from 'node:assert' 'use strict'
import { formatFilenameDate } from './_filenameDate.mjs' const assert = require('assert')
import { importIncrementalVm } from './_incrementalVm.mjs'
import { Task } from './Task.mjs'
import { watchStreamSize } from './_watchStreamSize.mjs'
export class ImportVmBackup { const { formatFilenameDate } = require('./_filenameDate.js')
const { importIncrementalVm } = require('./_incrementalVm.js')
const { Task } = require('./Task.js')
const { watchStreamSize } = require('./_watchStreamSize.js')
exports.ImportVmBackup = class ImportVmBackup {
constructor({ adapter, metadata, srUuid, xapi, settings: { newMacAddresses, mapVdisSrs = {} } = {} }) { constructor({ adapter, metadata, srUuid, xapi, settings: { newMacAddresses, mapVdisSrs = {} } = {} }) {
this._adapter = adapter this._adapter = adapter
this._importIncrementalVmSettings = { newMacAddresses, mapVdisSrs } this._importIncrementalVmSettings = { newMacAddresses, mapVdisSrs }

View File

@@ -1,39 +1,43 @@
import { asyncEach } from '@vates/async-each' 'use strict'
import { asyncMap, asyncMapSettled } from '@xen-orchestra/async-map'
import { compose } from '@vates/compose'
import { createLogger } from '@xen-orchestra/log'
import { createVhdDirectoryFromStream, openVhd, VhdAbstract, VhdDirectory, VhdSynthetic } from 'vhd-lib'
import { decorateMethodsWith } from '@vates/decorate-with'
import { deduped } from '@vates/disposable/deduped.js'
import { dirname, join, resolve } from 'node:path'
import { execFile } from 'child_process'
import { mount } from '@vates/fuse-vhd'
import { readdir, lstat } from 'node:fs/promises'
import { synchronized } from 'decorator-synchronized'
import { v4 as uuidv4 } from 'uuid'
import { ZipFile } from 'yazl'
import Disposable from 'promise-toolbox/Disposable'
import fromCallback from 'promise-toolbox/fromCallback'
import fromEvent from 'promise-toolbox/fromEvent'
import groupBy from 'lodash/groupBy.js'
import pDefer from 'promise-toolbox/defer'
import pickBy from 'lodash/pickBy.js'
import tar from 'tar'
import zlib from 'zlib'
import { BACKUP_DIR } from './_getVmBackupDir.mjs' const { asyncMap, asyncMapSettled } = require('@xen-orchestra/async-map')
import { cleanVm } from './_cleanVm.mjs' const { synchronized } = require('decorator-synchronized')
import { formatFilenameDate } from './_filenameDate.mjs' const Disposable = require('promise-toolbox/Disposable')
import { getTmpDir } from './_getTmpDir.mjs' const fromCallback = require('promise-toolbox/fromCallback')
import { isMetadataFile } from './_backupType.mjs' const fromEvent = require('promise-toolbox/fromEvent')
import { isValidXva } from './_isValidXva.mjs' const pDefer = require('promise-toolbox/defer')
import { listPartitions, LVM_PARTITION_TYPE } from './_listPartitions.mjs' const groupBy = require('lodash/groupBy.js')
import { lvs, pvs } from './_lvm.mjs' const pickBy = require('lodash/pickBy.js')
import { watchStreamSize } from './_watchStreamSize.mjs' const { dirname, join, normalize, resolve } = require('path')
const { createLogger } = require('@xen-orchestra/log')
const { createVhdDirectoryFromStream, openVhd, VhdAbstract, VhdDirectory, VhdSynthetic } = require('vhd-lib')
const { deduped } = require('@vates/disposable/deduped.js')
const { decorateMethodsWith } = require('@vates/decorate-with')
const { compose } = require('@vates/compose')
const { execFile } = require('child_process')
const { readdir, lstat } = require('fs-extra')
const { v4: uuidv4 } = require('uuid')
const { ZipFile } = require('yazl')
const zlib = require('zlib')
export const DIR_XO_CONFIG_BACKUPS = 'xo-config-backups' const { BACKUP_DIR } = require('./_getVmBackupDir.js')
const { cleanVm } = require('./_cleanVm.js')
const { formatFilenameDate } = require('./_filenameDate.js')
const { getTmpDir } = require('./_getTmpDir.js')
const { isMetadataFile } = require('./_backupType.js')
const { isValidXva } = require('./_isValidXva.js')
const { listPartitions, LVM_PARTITION_TYPE } = require('./_listPartitions.js')
const { lvs, pvs } = require('./_lvm.js')
const { watchStreamSize } = require('./_watchStreamSize')
// @todo : this import is marked extraneous , sould be fixed when lib is published
const { mount } = require('@vates/fuse-vhd')
const { asyncEach } = require('@vates/async-each')
export const DIR_XO_POOL_METADATA_BACKUPS = 'xo-pool-metadata-backups' const DIR_XO_CONFIG_BACKUPS = 'xo-config-backups'
exports.DIR_XO_CONFIG_BACKUPS = DIR_XO_CONFIG_BACKUPS
const DIR_XO_POOL_METADATA_BACKUPS = 'xo-pool-metadata-backups'
exports.DIR_XO_POOL_METADATA_BACKUPS = DIR_XO_POOL_METADATA_BACKUPS
const { debug, warn } = createLogger('xo:backups:RemoteAdapter') const { debug, warn } = createLogger('xo:backups:RemoteAdapter')
@@ -42,23 +46,20 @@ const compareTimestamp = (a, b) => a.timestamp - b.timestamp
const noop = Function.prototype const noop = Function.prototype
const resolveRelativeFromFile = (file, path) => resolve('/', dirname(file), path).slice(1) const resolveRelativeFromFile = (file, path) => resolve('/', dirname(file), path).slice(1)
const makeRelative = path => resolve('/', path).slice(1)
const resolveSubpath = (root, path) => resolve(root, makeRelative(path))
async function addZipEntries(zip, realBasePath, virtualBasePath, relativePaths) { const resolveSubpath = (root, path) => resolve(root, `.${resolve('/', path)}`)
for (const relativePath of relativePaths) {
const realPath = join(realBasePath, relativePath)
const virtualPath = join(virtualBasePath, relativePath)
async function addDirectory(files, realPath, metadataPath) {
const stats = await lstat(realPath) const stats = await lstat(realPath)
const { mode, mtime } = stats
const opts = { mode, mtime }
if (stats.isDirectory()) { if (stats.isDirectory()) {
zip.addEmptyDirectory(virtualPath, opts) await asyncMap(await readdir(realPath), file =>
await addZipEntries(zip, realPath, virtualPath, await readdir(realPath)) addDirectory(files, realPath + '/' + file, metadataPath + '/' + file)
)
} else if (stats.isFile()) { } else if (stats.isFile()) {
zip.addFile(realPath, virtualPath, opts) files.push({
} realPath,
metadataPath,
})
} }
} }
@@ -75,7 +76,7 @@ const debounceResourceFactory = factory =>
return this._debounceResource(factory.apply(this, arguments)) return this._debounceResource(factory.apply(this, arguments))
} }
export class RemoteAdapter { class RemoteAdapter {
constructor( constructor(
handler, handler,
{ debounceResource = res => res, dirMode, vhdDirectoryCompression, useGetDiskLegacy = false } = {} { debounceResource = res => res, dirMode, vhdDirectoryCompression, useGetDiskLegacy = false } = {}
@@ -186,6 +187,17 @@ export class RemoteAdapter {
}) })
} }
async *_usePartitionFiles(diskId, partitionId, paths) {
const path = yield this.getPartition(diskId, partitionId)
const files = []
await asyncMap(paths, file =>
addDirectory(files, resolveSubpath(path, file), normalize('./' + file).replace(/\/+$/, ''))
)
return files
}
// check if we will be allowed to merge a a vhd created in this adapter // check if we will be allowed to merge a a vhd created in this adapter
// with the vhd at path `path` // with the vhd at path `path`
async isMergeableParent(packedParentUid, path) { async isMergeableParent(packedParentUid, path) {
@@ -202,24 +214,15 @@ export class RemoteAdapter {
}) })
} }
fetchPartitionFiles(diskId, partitionId, paths, format) { fetchPartitionFiles(diskId, partitionId, paths) {
const { promise, reject, resolve } = pDefer() const { promise, reject, resolve } = pDefer()
Disposable.use( Disposable.use(
async function* () { async function* () {
const path = yield this.getPartition(diskId, partitionId) const files = yield this._usePartitionFiles(diskId, partitionId, paths)
let outputStream
if (format === 'tgz') {
outputStream = tar.c({ cwd: path, gzip: true }, paths.map(makeRelative))
} else if (format === 'zip') {
const zip = new ZipFile() const zip = new ZipFile()
await addZipEntries(zip, path, '', paths.map(makeRelative)) files.forEach(({ realPath, metadataPath }) => zip.addFile(realPath, metadataPath))
zip.end() zip.end()
;({ outputStream } = zip) const { outputStream } = zip
} else {
throw new Error('unsupported format ' + format)
}
resolve(outputStream) resolve(outputStream)
await fromEvent(outputStream, 'end') await fromEvent(outputStream, 'end')
}.bind(this) }.bind(this)
@@ -666,7 +669,7 @@ export class RemoteAdapter {
const handler = this._handler const handler = this._handler
if (this.useVhdDirectory()) { if (this.useVhdDirectory()) {
const dataPath = `${dirname(path)}/data/${uuidv4()}.vhd` const dataPath = `${dirname(path)}/data/${uuidv4()}.vhd`
const size = await createVhdDirectoryFromStream(handler, dataPath, input, { const sizes = await createVhdDirectoryFromStream(handler, dataPath, input, {
concurrency: writeBlockConcurrency, concurrency: writeBlockConcurrency,
compression: this.#getCompressionType(), compression: this.#getCompressionType(),
async validator() { async validator() {
@@ -675,19 +678,22 @@ export class RemoteAdapter {
}, },
}) })
await VhdAbstract.createAlias(handler, path, dataPath) await VhdAbstract.createAlias(handler, path, dataPath)
return size return sizes
} else { } else {
return this.outputStream(path, input, { checksum, validator }) const size = this.outputStream(path, input, { checksum, validator })
return {
compressedSize: size,
sourceSize: size,
writtenSize: size,
}
} }
} }
async outputStream(path, input, { checksum = true, maxStreamLength, streamLength, validator = noop } = {}) { async outputStream(path, input, { checksum = true, validator = noop } = {}) {
const container = watchStreamSize(input) const container = watchStreamSize(input)
await this._handler.outputStream(path, input, { await this._handler.outputStream(path, input, {
checksum, checksum,
dirMode: this._dirMode, dirMode: this._dirMode,
maxStreamLength,
streamLength,
async validator() { async validator() {
await input.task await input.task
return validator.apply(this, arguments) return validator.apply(this, arguments)
@@ -828,7 +834,11 @@ decorateMethodsWith(RemoteAdapter, {
debounceResourceFactory, debounceResourceFactory,
]), ]),
_usePartitionFiles: Disposable.factory,
getDisk: compose([Disposable.factory, [deduped, diskId => [diskId]], debounceResourceFactory]), getDisk: compose([Disposable.factory, [deduped, diskId => [diskId]], debounceResourceFactory]),
getPartition: Disposable.factory, getPartition: Disposable.factory,
}) })
exports.RemoteAdapter = RemoteAdapter

View File

@@ -0,0 +1,29 @@
'use strict'
const { join, resolve } = require('node:path/posix')
const { DIR_XO_POOL_METADATA_BACKUPS } = require('./RemoteAdapter.js')
const { PATH_DB_DUMP } = require('./_runners/_PoolMetadataBackup.js')
exports.RestoreMetadataBackup = class RestoreMetadataBackup {
constructor({ backupId, handler, xapi }) {
this._backupId = backupId
this._handler = handler
this._xapi = xapi
}
async run() {
const backupId = this._backupId
const handler = this._handler
const xapi = this._xapi
if (backupId.split('/')[0] === DIR_XO_POOL_METADATA_BACKUPS) {
return xapi.putResource(await handler.createReadStream(`${backupId}/data`), PATH_DB_DUMP, {
task: xapi.task_create('Import pool metadata'),
})
} else {
const metadata = JSON.parse(await handler.readFile(join(backupId, 'metadata.json')))
return String(await handler.readFile(resolve(backupId, metadata.data ?? 'data.json')))
}
}
}

View File

@@ -1,32 +0,0 @@
import { join, resolve } from 'node:path/posix'
import { DIR_XO_POOL_METADATA_BACKUPS } from './RemoteAdapter.mjs'
import { PATH_DB_DUMP } from './_runners/_PoolMetadataBackup.mjs'
export class RestoreMetadataBackup {
constructor({ backupId, handler, xapi }) {
this._backupId = backupId
this._handler = handler
this._xapi = xapi
}
async run() {
const backupId = this._backupId
const handler = this._handler
const xapi = this._xapi
if (backupId.split('/')[0] === DIR_XO_POOL_METADATA_BACKUPS) {
return xapi.putResource(await handler.createReadStream(`${backupId}/data`), PATH_DB_DUMP, {
task: xapi.task_create('Import pool metadata'),
})
} else {
const metadata = JSON.parse(await handler.readFile(join(backupId, 'metadata.json')))
const dataFileName = resolve(backupId, metadata.data ?? 'data.json')
const data = await handler.readFile(dataFileName)
// if data is JSON, sent it as a plain string, otherwise, consider the data as binary and encode it
const isJson = dataFileName.endsWith('.json')
return isJson ? data.toString() : { encoding: 'base64', data: data.toString('base64') }
}
}
}

View File

@@ -1,5 +1,7 @@
import CancelToken from 'promise-toolbox/CancelToken' 'use strict'
import Zone from 'node-zone'
const CancelToken = require('promise-toolbox/CancelToken')
const Zone = require('node-zone')
const logAfterEnd = log => { const logAfterEnd = log => {
const error = new Error('task has already ended') const error = new Error('task has already ended')
@@ -28,7 +30,7 @@ const serializeError = error =>
const $$task = Symbol('@xen-orchestra/backups/Task') const $$task = Symbol('@xen-orchestra/backups/Task')
export class Task { class Task {
static get cancelToken() { static get cancelToken() {
const task = Zone.current.data[$$task] const task = Zone.current.data[$$task]
return task !== undefined ? task.#cancelToken : CancelToken.none return task !== undefined ? task.#cancelToken : CancelToken.none
@@ -149,6 +151,7 @@ export class Task {
}) })
} }
} }
exports.Task = Task
for (const method of ['info', 'warning']) { for (const method of ['info', 'warning']) {
Task[method] = (...args) => Zone.current.data[$$task]?.[method](...args) Task[method] = (...args) => Zone.current.data[$$task]?.[method](...args)

View File

@@ -0,0 +1,6 @@
'use strict'
exports.isMetadataFile = filename => filename.endsWith('.json')
exports.isVhdFile = filename => filename.endsWith('.vhd')
exports.isXvaFile = filename => filename.endsWith('.xva')
exports.isXvaSumFile = filename => filename.endsWith('.xva.checksum')

View File

@@ -1,4 +0,0 @@
export const isMetadataFile = filename => filename.endsWith('.json')
export const isVhdFile = filename => filename.endsWith('.vhd')
export const isXvaFile = filename => filename.endsWith('.xva')
export const isXvaSumFile = filename => filename.endsWith('.xva.checksum')

View File

@@ -1,25 +1,25 @@
import { createLogger } from '@xen-orchestra/log' 'use strict'
import { catchGlobalErrors } from '@xen-orchestra/log/configure'
import Disposable from 'promise-toolbox/Disposable' const logger = require('@xen-orchestra/log').createLogger('xo:backups:worker')
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { compose } from '@vates/compose'
import { createCachedLookup } from '@vates/cached-dns.lookup'
import { createDebounceResource } from '@vates/disposable/debounceResource.js'
import { createRunner } from './Backup.mjs'
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 { Xapi } from '@xen-orchestra/xapi'
import { RemoteAdapter } from './RemoteAdapter.mjs' require('@xen-orchestra/log/configure').catchGlobalErrors(logger)
import { Task } from './Task.mjs'
createCachedLookup().patchGlobal() require('@vates/cached-dns.lookup').createCachedLookup().patchGlobal()
const Disposable = require('promise-toolbox/Disposable')
const ignoreErrors = require('promise-toolbox/ignoreErrors')
const { compose } = require('@vates/compose')
const { createDebounceResource } = require('@vates/disposable/debounceResource.js')
const { decorateMethodsWith } = require('@vates/decorate-with')
const { deduped } = require('@vates/disposable/deduped.js')
const { getHandler } = require('@xen-orchestra/fs')
const { createRunner } = require('./Backup.js')
const { parseDuration } = require('@vates/parse-duration')
const { Xapi } = require('@xen-orchestra/xapi')
const { RemoteAdapter } = require('./RemoteAdapter.js')
const { Task } = require('./Task.js')
const logger = createLogger('xo:backups:worker')
catchGlobalErrors(logger)
const { debug } = logger const { debug } = logger
class BackupWorker { class BackupWorker {

View File

@@ -1,11 +1,13 @@
import cancelable from 'promise-toolbox/cancelable' 'use strict'
import CancelToken from 'promise-toolbox/CancelToken'
const cancelable = require('promise-toolbox/cancelable')
const CancelToken = require('promise-toolbox/CancelToken')
// Similar to `Promise.all` + `map` but pass a cancel token to the callback // Similar to `Promise.all` + `map` but pass a cancel token to the callback
// //
// If any of the executions fails, the cancel token will be triggered and the // If any of the executions fails, the cancel token will be triggered and the
// first reason will be rejected. // first reason will be rejected.
export const cancelableMap = cancelable(async function cancelableMap($cancelToken, iterable, callback) { exports.cancelableMap = cancelable(async function cancelableMap($cancelToken, iterable, callback) {
const { cancel, token } = CancelToken.source([$cancelToken]) const { cancel, token } = CancelToken.source([$cancelToken])
try { try {
return await Promise.all( return await Promise.all(

View File

@@ -1,19 +1,19 @@
import test from 'test' 'use strict'
import { strict as assert } from 'node:assert'
import tmp from 'tmp' const { beforeEach, afterEach, test, describe } = require('test')
import fs from 'fs-extra' const assert = require('assert').strict
import * as uuid from 'uuid'
import { getHandler } from '@xen-orchestra/fs'
import { pFromCallback } from 'promise-toolbox'
import { RemoteAdapter } from './RemoteAdapter.mjs'
import { VHDFOOTER, VHDHEADER } from './tests.fixtures.mjs'
import { VhdFile, Constants, VhdDirectory, VhdAbstract } from 'vhd-lib'
import { checkAliases } from './_cleanVm.mjs'
import { dirname, basename } from 'node:path'
import { rimraf } from 'rimraf'
const { beforeEach, afterEach, describe } = test const tmp = require('tmp')
const fs = require('fs-extra')
const uuid = require('uuid')
const { getHandler } = require('@xen-orchestra/fs')
const { pFromCallback } = require('promise-toolbox')
const { RemoteAdapter } = require('./RemoteAdapter')
const { VHDFOOTER, VHDHEADER } = require('./tests.fixtures.js')
const { VhdFile, Constants, VhdDirectory, VhdAbstract } = require('vhd-lib')
const { checkAliases } = require('./_cleanVm')
const { dirname, basename } = require('path')
const { rimraf } = require('rimraf')
let tempDir, adapter, handler, jobId, vdiId, basePath, relativePath let tempDir, adapter, handler, jobId, vdiId, basePath, relativePath
const rootPath = 'xo-vm-backups/VMUUID/' const rootPath = 'xo-vm-backups/VMUUID/'

View File

@@ -1,38 +1,29 @@
import * as UUID from 'uuid' 'use strict'
import sum from 'lodash/sum.js'
import { asyncMap } from '@xen-orchestra/async-map'
import { Constants, openVhd, VhdAbstract, VhdFile } from 'vhd-lib'
import { isVhdAlias, resolveVhdAlias } from 'vhd-lib/aliases.js'
import { dirname, resolve } from 'node:path'
import { isMetadataFile, isVhdFile, isXvaFile, isXvaSumFile } from './_backupType.mjs'
import { limitConcurrency } from 'limit-concurrency-decorator'
import { mergeVhdChain } from 'vhd-lib/merge.js'
import { Task } from './Task.mjs'
import { Disposable } from 'promise-toolbox'
import handlerPath from '@xen-orchestra/fs/path'
const sum = require('lodash/sum')
const UUID = require('uuid')
const { asyncMap } = require('@xen-orchestra/async-map')
const { Constants, openVhd, VhdAbstract } = require('vhd-lib')
const { isVhdAlias, resolveVhdAlias } = require('vhd-lib/aliases')
const { dirname, resolve } = require('path')
const { DISK_TYPES } = Constants const { DISK_TYPES } = Constants
const { isMetadataFile, isVhdFile, isXvaFile, isXvaSumFile } = require('./_backupType.js')
const { limitConcurrency } = require('limit-concurrency-decorator')
const { mergeVhdChain } = require('vhd-lib/merge')
// checking the size of a vhd directory is costly const { Task } = require('./Task.js')
// 1 Http Query per 1000 blocks const { Disposable } = require('promise-toolbox')
// we only check size of all the vhd are VhdFiles const handlerPath = require('@xen-orchestra/fs/path')
function shouldComputeVhdsSize(handler, vhds) {
if (handler.isEncrypted) {
return false
}
return vhds.every(vhd => vhd instanceof VhdFile)
}
const computeVhdsSize = (handler, vhdPaths) => const computeVhdsSize = (handler, vhdPaths) =>
Disposable.use( Disposable.use(
vhdPaths.map(vhdPath => openVhd(handler, vhdPath)), vhdPaths.map(vhdPath => openVhd(handler, vhdPath)),
async vhds => { async vhds => {
if (shouldComputeVhdsSize(handler, vhds)) { await Promise.all(vhds.map(vhd => vhd.readBlockAllocationTable()))
const sizes = await asyncMap(vhds, vhd => vhd.getSize()) // get file size for vhdfile, computed size from bat for vhd directory
const sizes = await asyncMap(vhds, vhd => vhd.streamSize())
return sum(sizes) return sum(sizes)
} }
}
) )
// chain is [ ancestor, child_1, ..., child_n ] // chain is [ ancestor, child_1, ..., child_n ]
@@ -116,7 +107,7 @@ const listVhds = async (handler, vmDir, logWarn) => {
return { vhds, interruptedVhds, aliases } return { vhds, interruptedVhds, aliases }
} }
export async function checkAliases( async function checkAliases(
aliasPaths, aliasPaths,
targetDataRepository, targetDataRepository,
{ handler, logInfo = noop, logWarn = console.warn, remove = false } { handler, logInfo = noop, logWarn = console.warn, remove = false }
@@ -175,9 +166,11 @@ export async function checkAliases(
}) })
} }
exports.checkAliases = checkAliases
const defaultMergeLimiter = limitConcurrency(1) const defaultMergeLimiter = limitConcurrency(1)
export async function cleanVm( exports.cleanVm = async function cleanVm(
vmDir, vmDir,
{ {
fixMetadata, fixMetadata,
@@ -531,11 +524,6 @@ export async function cleanVm(
const linkedVhds = Object.keys(vhds).map(key => resolve('/', vmDir, vhds[key])) const linkedVhds = Object.keys(vhds).map(key => resolve('/', vmDir, vhds[key]))
fileSystemSize = await computeVhdsSize(handler, linkedVhds) fileSystemSize = await computeVhdsSize(handler, linkedVhds)
// the size is not computed in some cases (e.g. VhdDirectory)
if (fileSystemSize === undefined) {
return
}
// don't warn if the size has changed after a merge // don't warn if the size has changed after a merge
if (!merged && fileSystemSize !== size) { if (!merged && fileSystemSize !== size) {
// FIXME: figure out why it occurs so often and, once fixed, log the real problems with `logWarn` // FIXME: figure out why it occurs so often and, once fixed, log the real problems with `logWarn`
@@ -553,6 +541,8 @@ export async function cleanVm(
// systematically update size after a merge // systematically update size after a merge
if ((merged || fixMetadata) && size !== fileSystemSize) { if ((merged || fixMetadata) && size !== fileSystemSize) {
// @todo add a cumulatedTransferSize property ?
// @todo update writtenSize, compressedSize
metadata.size = fileSystemSize metadata.size = fileSystemSize
mustRegenerateCache = true mustRegenerateCache = true
try { try {

View File

@@ -0,0 +1,8 @@
'use strict'
const { utcFormat, utcParse } = require('d3-time-format')
// Format a date in ISO 8601 in a safe way to be used in filenames
// (even on Windows).
exports.formatFilenameDate = utcFormat('%Y%m%dT%H%M%SZ')
exports.parseFilenameDate = utcParse('%Y%m%dT%H%M%SZ')

View File

@@ -1,6 +0,0 @@
import { utcFormat, utcParse } from 'd3-time-format'
// Format a date in ISO 8601 in a safe way to be used in filenames
// (even on Windows).
export const formatFilenameDate = utcFormat('%Y%m%dT%H%M%SZ')
export const parseFilenameDate = utcParse('%Y%m%dT%H%M%SZ')

View File

@@ -1,4 +1,6 @@
'use strict'
// returns all entries but the last retention-th // returns all entries but the last retention-th
export function getOldEntries(retention, entries) { exports.getOldEntries = function getOldEntries(retention, entries) {
return entries === undefined ? [] : retention > 0 ? entries.slice(0, -retention) : entries return entries === undefined ? [] : retention > 0 ? entries.slice(0, -retention) : entries
} }

View File

@@ -1,11 +1,13 @@
import Disposable from 'promise-toolbox/Disposable' 'use strict'
import { join } from 'node:path'
import { mkdir, rmdir } from 'node:fs/promises' const Disposable = require('promise-toolbox/Disposable')
import { tmpdir } from 'os' const { join } = require('path')
const { mkdir, rmdir } = require('fs-extra')
const { tmpdir } = require('os')
const MAX_ATTEMPTS = 3 const MAX_ATTEMPTS = 3
export async function getTmpDir() { exports.getTmpDir = async function getTmpDir() {
for (let i = 0; true; ++i) { for (let i = 0; true; ++i) {
const path = join(tmpdir(), Math.random().toString(36).slice(2)) const path = join(tmpdir(), Math.random().toString(36).slice(2))
try { try {

View File

@@ -0,0 +1,8 @@
'use strict'
const BACKUP_DIR = 'xo-vm-backups'
exports.BACKUP_DIR = BACKUP_DIR
exports.getVmBackupDir = function getVmBackupDir(uuid) {
return `${BACKUP_DIR}/${uuid}`
}

View File

@@ -1,5 +0,0 @@
export const BACKUP_DIR = 'xo-vm-backups'
export function getVmBackupDir(uuid) {
return `${BACKUP_DIR}/${uuid}`
}

View File

@@ -1,22 +1,24 @@
import find from 'lodash/find.js' 'use strict'
import groupBy from 'lodash/groupBy.js'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import omit from 'lodash/omit.js'
import { asyncMap } from '@xen-orchestra/async-map'
import { CancelToken } from 'promise-toolbox'
import { compareVersions } from 'compare-versions'
import { createVhdStreamWithLength } from 'vhd-lib'
import { defer } from 'golike-defer'
import { cancelableMap } from './_cancelableMap.mjs' const find = require('lodash/find.js')
import { Task } from './Task.mjs' const groupBy = require('lodash/groupBy.js')
import pick from 'lodash/pick.js' const ignoreErrors = require('promise-toolbox/ignoreErrors')
const omit = require('lodash/omit.js')
const { asyncMap } = require('@xen-orchestra/async-map')
const { CancelToken } = require('promise-toolbox')
const { compareVersions } = require('compare-versions')
const { createVhdStreamWithLength } = require('vhd-lib')
const { defer } = require('golike-defer')
export const TAG_BASE_DELTA = 'xo:base_delta' const { cancelableMap } = require('./_cancelableMap.js')
const { Task } = require('./Task.js')
const pick = require('lodash/pick.js')
export const TAG_COPY_SRC = 'xo:copy_of' const TAG_BASE_DELTA = 'xo:base_delta'
exports.TAG_BASE_DELTA = TAG_BASE_DELTA
const TAG_BACKUP_SR = 'xo:backup:sr' const TAG_COPY_SRC = 'xo:copy_of'
exports.TAG_COPY_SRC = TAG_COPY_SRC
const ensureArray = value => (value === undefined ? [] : Array.isArray(value) ? value : [value]) const ensureArray = value => (value === undefined ? [] : Array.isArray(value) ? value : [value])
const resolveUuid = async (xapi, cache, uuid, type) => { const resolveUuid = async (xapi, cache, uuid, type) => {
@@ -31,7 +33,7 @@ const resolveUuid = async (xapi, cache, uuid, type) => {
return ref return ref
} }
export async function exportIncrementalVm( exports.exportIncrementalVm = async function exportIncrementalVm(
vm, vm,
baseVm, baseVm,
{ {
@@ -41,7 +43,6 @@ export async function exportIncrementalVm(
fullVdisRequired = new Set(), fullVdisRequired = new Set(),
disableBaseTags = false, disableBaseTags = false,
preferNbd,
} = {} } = {}
) { ) {
// refs of VM's VDIs → base's VDIs. // refs of VM's VDIs → base's VDIs.
@@ -89,7 +90,6 @@ export async function exportIncrementalVm(
baseRef: baseVdi?.$ref, baseRef: baseVdi?.$ref,
cancelToken, cancelToken,
format: 'vhd', format: 'vhd',
preferNbd,
}) })
}) })
@@ -143,7 +143,7 @@ export async function exportIncrementalVm(
) )
} }
export const importIncrementalVm = defer(async function importIncrementalVm( exports.importIncrementalVm = defer(async function importIncrementalVm(
$defer, $defer,
incrementalVm, incrementalVm,
sr, sr,
@@ -161,10 +161,7 @@ export const importIncrementalVm = defer(async function importIncrementalVm(
if (detectBase) { if (detectBase) {
const remoteBaseVmUuid = vmRecord.other_config[TAG_BASE_DELTA] const remoteBaseVmUuid = vmRecord.other_config[TAG_BASE_DELTA]
if (remoteBaseVmUuid) { if (remoteBaseVmUuid) {
baseVm = find( baseVm = find(xapi.objects.all, obj => (obj = obj.other_config) && obj[TAG_COPY_SRC] === remoteBaseVmUuid)
xapi.objects.all,
obj => (obj = obj.other_config) && obj[TAG_COPY_SRC] === remoteBaseVmUuid && obj[TAG_BACKUP_SR] === sr.$id
)
if (!baseVm) { if (!baseVm) {
throw new Error(`could not find the base VM (copy of ${remoteBaseVmUuid})`) throw new Error(`could not find the base VM (copy of ${remoteBaseVmUuid})`)

View File

@@ -1,4 +1,6 @@
import assert from 'node:assert' 'use strict'
const assert = require('assert')
const COMPRESSED_MAGIC_NUMBERS = [ const COMPRESSED_MAGIC_NUMBERS = [
// https://tools.ietf.org/html/rfc1952.html#page-5 // https://tools.ietf.org/html/rfc1952.html#page-5
@@ -45,7 +47,7 @@ const isValidTar = async (handler, size, fd) => {
} }
// TODO: find an heuristic for compressed files // TODO: find an heuristic for compressed files
export async function isValidXva(path) { async function isValidXva(path) {
const handler = this._handler const handler = this._handler
// size is longer when encrypted + reading part of an encrypted file is not implemented // size is longer when encrypted + reading part of an encrypted file is not implemented
@@ -72,5 +74,6 @@ export async function isValidXva(path) {
return true return true
} }
} }
exports.isValidXva = isValidXva
const noop = Function.prototype const noop = Function.prototype

View File

@@ -1,7 +1,9 @@
import fromCallback from 'promise-toolbox/fromCallback' 'use strict'
import { createLogger } from '@xen-orchestra/log'
import { createParser } from 'parse-pairs' const fromCallback = require('promise-toolbox/fromCallback')
import { execFile } from 'child_process' const { createLogger } = require('@xen-orchestra/log')
const { createParser } = require('parse-pairs')
const { execFile } = require('child_process')
const { debug } = createLogger('xo:backups:listPartitions') const { debug } = createLogger('xo:backups:listPartitions')
@@ -22,7 +24,8 @@ const IGNORED_PARTITION_TYPES = {
0x82: true, // swap 0x82: true, // swap
} }
export const LVM_PARTITION_TYPE = 0x8e const LVM_PARTITION_TYPE = 0x8e
exports.LVM_PARTITION_TYPE = LVM_PARTITION_TYPE
const parsePartxLine = createParser({ const parsePartxLine = createParser({
keyTransform: key => (key === 'UUID' ? 'id' : key.toLowerCase()), keyTransform: key => (key === 'UUID' ? 'id' : key.toLowerCase()),
@@ -30,7 +33,7 @@ const parsePartxLine = createParser({
}) })
// returns an empty array in case of a non-partitioned disk // returns an empty array in case of a non-partitioned disk
export async function listPartitions(devicePath) { exports.listPartitions = async function listPartitions(devicePath) {
const parts = await fromCallback(execFile, 'partx', [ const parts = await fromCallback(execFile, 'partx', [
'--bytes', '--bytes',
'--output=NR,START,SIZE,NAME,UUID,TYPE', '--output=NR,START,SIZE,NAME,UUID,TYPE',

View File

@@ -1,6 +1,8 @@
import fromCallback from 'promise-toolbox/fromCallback' 'use strict'
import { createParser } from 'parse-pairs'
import { execFile } from 'child_process' const fromCallback = require('promise-toolbox/fromCallback')
const { createParser } = require('parse-pairs')
const { execFile } = require('child_process')
// =================================================================== // ===================================================================
@@ -27,5 +29,5 @@ const makeFunction =
.map(Array.isArray(fields) ? parse : line => parse(line)[fields]) .map(Array.isArray(fields) ? parse : line => parse(line)[fields])
} }
export const lvs = makeFunction('lvs') exports.lvs = makeFunction('lvs')
export const pvs = makeFunction('pvs') exports.pvs = makeFunction('pvs')

View File

@@ -1,20 +1,22 @@
import { asyncMap } from '@xen-orchestra/async-map' 'use strict'
import Disposable from 'promise-toolbox/Disposable'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { extractIdsFromSimplePattern } from '../extractIdsFromSimplePattern.mjs' const { asyncMap } = require('@xen-orchestra/async-map')
import { PoolMetadataBackup } from './_PoolMetadataBackup.mjs' const Disposable = require('promise-toolbox/Disposable')
import { XoMetadataBackup } from './_XoMetadataBackup.mjs' const ignoreErrors = require('promise-toolbox/ignoreErrors')
import { DEFAULT_SETTINGS, Abstract } from './_Abstract.mjs'
import { runTask } from './_runTask.mjs' const { extractIdsFromSimplePattern } = require('../extractIdsFromSimplePattern.js')
import { getAdaptersByRemote } from './_getAdaptersByRemote.mjs' const { PoolMetadataBackup } = require('./_PoolMetadataBackup.js')
const { XoMetadataBackup } = require('./_XoMetadataBackup.js')
const { DEFAULT_SETTINGS, Abstract } = require('./_Abstract.js')
const { runTask } = require('./_runTask.js')
const { getAdaptersByRemote } = require('./_getAdaptersByRemote.js')
const DEFAULT_METADATA_SETTINGS = { const DEFAULT_METADATA_SETTINGS = {
retentionPoolMetadata: 0, retentionPoolMetadata: 0,
retentionXoMetadata: 0, retentionXoMetadata: 0,
} }
export const Metadata = class MetadataBackupRunner extends Abstract { exports.Metadata = class MetadataBackupRunner extends Abstract {
_computeBaseSettings(config, job) { _computeBaseSettings(config, job) {
const baseSettings = { ...DEFAULT_SETTINGS } const baseSettings = { ...DEFAULT_SETTINGS }
Object.assign(baseSettings, DEFAULT_METADATA_SETTINGS, config.defaultSettings, config.metadata?.defaultSettings) Object.assign(baseSettings, DEFAULT_METADATA_SETTINGS, config.defaultSettings, config.metadata?.defaultSettings)

View File

@@ -1,15 +1,17 @@
import { asyncMapSettled } from '@xen-orchestra/async-map' 'use strict'
import Disposable from 'promise-toolbox/Disposable'
import { limitConcurrency } from 'limit-concurrency-decorator'
import { extractIdsFromSimplePattern } from '../extractIdsFromSimplePattern.mjs' const { asyncMapSettled } = require('@xen-orchestra/async-map')
import { Task } from '../Task.mjs' const Disposable = require('promise-toolbox/Disposable')
import createStreamThrottle from './_createStreamThrottle.mjs' const { limitConcurrency } = require('limit-concurrency-decorator')
import { DEFAULT_SETTINGS, Abstract } from './_Abstract.mjs'
import { runTask } from './_runTask.mjs' const { extractIdsFromSimplePattern } = require('../extractIdsFromSimplePattern.js')
import { getAdaptersByRemote } from './_getAdaptersByRemote.mjs' const { Task } = require('../Task.js')
import { FullRemote } from './_vmRunners/FullRemote.mjs' const createStreamThrottle = require('./_createStreamThrottle.js')
import { IncrementalRemote } from './_vmRunners/IncrementalRemote.mjs' const { DEFAULT_SETTINGS, Abstract } = require('./_Abstract.js')
const { runTask } = require('./_runTask.js')
const { getAdaptersByRemote } = require('./_getAdaptersByRemote.js')
const { FullRemote } = require('./_vmRunners/FullRemote.js')
const { IncrementalRemote } = require('./_vmRunners/IncrementalRemote.js')
const DEFAULT_REMOTE_VM_SETTINGS = { const DEFAULT_REMOTE_VM_SETTINGS = {
concurrency: 2, concurrency: 2,
@@ -25,7 +27,7 @@ const DEFAULT_REMOTE_VM_SETTINGS = {
vmTimeout: 0, vmTimeout: 0,
} }
export const VmsRemote = class RemoteVmsBackupRunner extends Abstract { exports.VmsRemote = class RemoteVmsBackupRunner extends Abstract {
_computeBaseSettings(config, job) { _computeBaseSettings(config, job) {
const baseSettings = { ...DEFAULT_SETTINGS } const baseSettings = { ...DEFAULT_SETTINGS }
Object.assign(baseSettings, DEFAULT_REMOTE_VM_SETTINGS, config.defaultSettings, config.vm?.defaultSettings) Object.assign(baseSettings, DEFAULT_REMOTE_VM_SETTINGS, config.defaultSettings, config.vm?.defaultSettings)

View File

@@ -1,15 +1,17 @@
import { asyncMapSettled } from '@xen-orchestra/async-map' 'use strict'
import Disposable from 'promise-toolbox/Disposable'
import { limitConcurrency } from 'limit-concurrency-decorator'
import { extractIdsFromSimplePattern } from '../extractIdsFromSimplePattern.mjs' const { asyncMapSettled } = require('@xen-orchestra/async-map')
import { Task } from '../Task.mjs' const Disposable = require('promise-toolbox/Disposable')
import createStreamThrottle from './_createStreamThrottle.mjs' const { limitConcurrency } = require('limit-concurrency-decorator')
import { DEFAULT_SETTINGS, Abstract } from './_Abstract.mjs'
import { runTask } from './_runTask.mjs' const { extractIdsFromSimplePattern } = require('../extractIdsFromSimplePattern.js')
import { getAdaptersByRemote } from './_getAdaptersByRemote.mjs' const { Task } = require('../Task.js')
import { IncrementalXapi } from './_vmRunners/IncrementalXapi.mjs' const createStreamThrottle = require('./_createStreamThrottle.js')
import { FullXapi } from './_vmRunners/FullXapi.mjs' const { DEFAULT_SETTINGS, Abstract } = require('./_Abstract.js')
const { runTask } = require('./_runTask.js')
const { getAdaptersByRemote } = require('./_getAdaptersByRemote.js')
const { IncrementalXapi } = require('./_vmRunners/IncrementalXapi.js')
const { FullXapi } = require('./_vmRunners/FullXapi.js')
const DEFAULT_XAPI_VM_SETTINGS = { const DEFAULT_XAPI_VM_SETTINGS = {
bypassVdiChainsCheck: false, bypassVdiChainsCheck: false,
@@ -34,7 +36,7 @@ const DEFAULT_XAPI_VM_SETTINGS = {
vmTimeout: 0, vmTimeout: 0,
} }
export const VmsXapi = class VmsXapiBackupRunner extends Abstract { exports.VmsXapi = class VmsXapiBackupRunner extends Abstract {
_computeBaseSettings(config, job) { _computeBaseSettings(config, job) {
const baseSettings = { ...DEFAULT_SETTINGS } const baseSettings = { ...DEFAULT_SETTINGS }
Object.assign(baseSettings, DEFAULT_XAPI_VM_SETTINGS, config.defaultSettings, config.vm?.defaultSettings) Object.assign(baseSettings, DEFAULT_XAPI_VM_SETTINGS, config.defaultSettings, config.vm?.defaultSettings)

View File

@@ -1,15 +1,17 @@
import Disposable from 'promise-toolbox/Disposable' 'use strict'
import pTimeout from 'promise-toolbox/timeout'
import { compileTemplate } from '@xen-orchestra/template'
import { runTask } from './_runTask.mjs'
import { RemoteTimeoutError } from './_RemoteTimeoutError.mjs'
export const DEFAULT_SETTINGS = { const Disposable = require('promise-toolbox/Disposable')
const pTimeout = require('promise-toolbox/timeout')
const { compileTemplate } = require('@xen-orchestra/template')
const { runTask } = require('./_runTask.js')
const { RemoteTimeoutError } = require('./_RemoteTimeoutError.js')
exports.DEFAULT_SETTINGS = {
getRemoteTimeout: 300e3, getRemoteTimeout: 300e3,
reportWhen: 'failure', reportWhen: 'failure',
} }
export const Abstract = class AbstractRunner { exports.Abstract = class AbstractRunner {
constructor({ config, getAdapter, getConnectedRecord, job, schedule }) { constructor({ config, getAdapter, getConnectedRecord, job, schedule }) {
this._config = config this._config = config
this._getRecord = getConnectedRecord this._getRecord = getConnectedRecord

View File

@@ -1,13 +1,16 @@
import { asyncMap } from '@xen-orchestra/async-map' 'use strict'
import { DIR_XO_POOL_METADATA_BACKUPS } from '../RemoteAdapter.mjs' const { asyncMap } = require('@xen-orchestra/async-map')
import { forkStreamUnpipe } from './_forkStreamUnpipe.mjs'
import { formatFilenameDate } from '../_filenameDate.mjs'
import { Task } from '../Task.mjs'
export const PATH_DB_DUMP = '/pool/xmldbdump' const { DIR_XO_POOL_METADATA_BACKUPS } = require('../RemoteAdapter.js')
const { forkStreamUnpipe } = require('./_forkStreamUnpipe.js')
const { formatFilenameDate } = require('../_filenameDate.js')
const { Task } = require('../Task.js')
export class PoolMetadataBackup { const PATH_DB_DUMP = '/pool/xmldbdump'
exports.PATH_DB_DUMP = PATH_DB_DUMP
exports.PoolMetadataBackup = class PoolMetadataBackup {
constructor({ config, job, pool, remoteAdapters, schedule, settings }) { constructor({ config, job, pool, remoteAdapters, schedule, settings }) {
this._config = config this._config = config
this._job = job this._job = job

View File

@@ -1,6 +1,8 @@
export class RemoteTimeoutError extends Error { 'use strict'
class RemoteTimeoutError extends Error {
constructor(remoteId) { constructor(remoteId) {
super('timeout while getting the remote ' + remoteId) super('timeout while getting the remote ' + remoteId)
this.remoteId = remoteId this.remoteId = remoteId
} }
} }
exports.RemoteTimeoutError = RemoteTimeoutError

View File

@@ -1,11 +1,13 @@
import { asyncMap } from '@xen-orchestra/async-map' 'use strict'
import { join } from '@xen-orchestra/fs/path'
import { DIR_XO_CONFIG_BACKUPS } from '../RemoteAdapter.mjs' const { asyncMap } = require('@xen-orchestra/async-map')
import { formatFilenameDate } from '../_filenameDate.mjs' const { join } = require('@xen-orchestra/fs/path')
import { Task } from '../Task.mjs'
export class XoMetadataBackup { const { DIR_XO_CONFIG_BACKUPS } = require('../RemoteAdapter.js')
const { formatFilenameDate } = require('../_filenameDate.js')
const { Task } = require('../Task.js')
exports.XoMetadataBackup = class XoMetadataBackup {
constructor({ config, job, remoteAdapters, schedule, settings }) { constructor({ config, job, remoteAdapters, schedule, settings }) {
this._config = config this._config = config
this._job = job this._job = job
@@ -22,13 +24,7 @@ export class XoMetadataBackup {
const dir = `${scheduleDir}/${formatFilenameDate(timestamp)}` const dir = `${scheduleDir}/${formatFilenameDate(timestamp)}`
const data = job.xoMetadata const data = job.xoMetadata
let dataBaseName = './data' const dataBaseName = './data.json'
// JSON data is sent as plain string, binary data is sent as an object with `data` and `encoding properties
const isJson = typeof data === 'string'
if (isJson) {
dataBaseName += '.json'
}
const metadata = JSON.stringify( const metadata = JSON.stringify(
{ {
@@ -60,7 +56,7 @@ export class XoMetadataBackup {
async () => { async () => {
const handler = adapter.handler const handler = adapter.handler
const dirMode = this._config.dirMode const dirMode = this._config.dirMode
await handler.outputFile(dataFileName, isJson ? data : Buffer.from(data.data, data.encoding), { dirMode }) await handler.outputFile(dataFileName, data, { dirMode })
await handler.outputFile(metaDataFileName, metadata, { await handler.outputFile(metaDataFileName, metadata, {
dirMode, dirMode,
}) })

View File

@@ -1,10 +1,12 @@
import { pipeline } from 'node:stream' 'use strict'
import { ThrottleGroup } from '@kldzj/stream-throttle'
import identity from 'lodash/identity.js' const { pipeline } = require('node:stream')
const { ThrottleGroup } = require('@kldzj/stream-throttle')
const identity = require('lodash/identity.js')
const noop = Function.prototype const noop = Function.prototype
export default function createStreamThrottle(rate) { module.exports = function createStreamThrottle(rate) {
if (rate === 0) { if (rate === 0) {
return identity return identity
} }

View File

@@ -1,13 +1,14 @@
import { createLogger } from '@xen-orchestra/log' 'use strict'
import { finished, PassThrough } from 'node:stream'
const { debug } = createLogger('xo:backups:forkStreamUnpipe') const { finished, PassThrough } = require('node:stream')
const { debug } = require('@xen-orchestra/log').createLogger('xo:backups:forkStreamUnpipe')
// create a new readable stream from an existing one which may be piped later // create a new readable stream from an existing one which may be piped later
// //
// in case of error in the new readable stream, it will simply be unpiped // in case of error in the new readable stream, it will simply be unpiped
// from the original one // from the original one
export function forkStreamUnpipe(source) { exports.forkStreamUnpipe = function forkStreamUnpipe(source) {
const { forks = 0 } = source const { forks = 0 } = source
source.forks = forks + 1 source.forks = forks + 1

View File

@@ -1,7 +1,9 @@
export function getAdaptersByRemote(adapters) { 'use strict'
const getAdaptersByRemote = adapters => {
const adaptersByRemote = {} const adaptersByRemote = {}
adapters.forEach(({ adapter, remoteId }) => { adapters.forEach(({ adapter, remoteId }) => {
adaptersByRemote[remoteId] = adapter adaptersByRemote[remoteId] = adapter
}) })
return adaptersByRemote return adaptersByRemote
} }
exports.getAdaptersByRemote = getAdaptersByRemote

View File

@@ -0,0 +1,6 @@
'use strict'
const { Task } = require('../Task.js')
const noop = Function.prototype
const runTask = (...args) => Task.run(...args).catch(noop) // errors are handled by logs
exports.runTask = runTask

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,12 +1,14 @@
import { decorateMethodsWith } from '@vates/decorate-with' 'use strict'
import { defer } from 'golike-defer'
import { AbstractRemote } from './_AbstractRemote.mjs'
import { FullRemoteWriter } from '../_writers/FullRemoteWriter.mjs'
import { forkStreamUnpipe } from '../_forkStreamUnpipe.mjs'
import { watchStreamSize } from '../../_watchStreamSize.mjs'
import { Task } from '../../Task.mjs'
export const FullRemote = class FullRemoteVmBackupRunner extends AbstractRemote { const { decorateMethodsWith } = require('@vates/decorate-with')
const { defer } = require('golike-defer')
const { AbstractRemote } = require('./_AbstractRemote')
const { FullRemoteWriter } = require('../_writers/FullRemoteWriter')
const { forkStreamUnpipe } = require('../_forkStreamUnpipe')
const { watchStreamSize } = require('../../_watchStreamSize')
const { Task } = require('../../Task')
class FullRemoteVmBackupRunner extends AbstractRemote {
_getRemoteWriter() { _getRemoteWriter() {
return FullRemoteWriter return FullRemoteWriter
} }
@@ -29,8 +31,6 @@ export const FullRemote = class FullRemoteVmBackupRunner extends AbstractRemote
writer => writer =>
writer.run({ writer.run({
stream: forkStreamUnpipe(stream), stream: forkStreamUnpipe(stream),
// stream will be forked and transformed, it's not safe to attach additionnal properties to it
streamLength: stream.length,
timestamp: metadata.timestamp, timestamp: metadata.timestamp,
vm: metadata.vm, vm: metadata.vm,
vmSnapshot: metadata.vmSnapshot, vmSnapshot: metadata.vmSnapshot,
@@ -47,6 +47,7 @@ export const FullRemote = class FullRemoteVmBackupRunner extends AbstractRemote
} }
} }
decorateMethodsWith(FullRemote, { exports.FullRemote = FullRemoteVmBackupRunner
decorateMethodsWith(FullRemoteVmBackupRunner, {
_run: defer, _run: defer,
}) })

View File

@@ -1,14 +1,16 @@
import { createLogger } from '@xen-orchestra/log' 'use strict'
import { forkStreamUnpipe } from '../_forkStreamUnpipe.mjs' const { createLogger } = require('@xen-orchestra/log')
import { FullRemoteWriter } from '../_writers/FullRemoteWriter.mjs'
import { FullXapiWriter } from '../_writers/FullXapiWriter.mjs' const { forkStreamUnpipe } = require('../_forkStreamUnpipe.js')
import { watchStreamSize } from '../../_watchStreamSize.mjs' const { FullRemoteWriter } = require('../_writers/FullRemoteWriter.js')
import { AbstractXapi } from './_AbstractXapi.mjs' const { FullXapiWriter } = require('../_writers/FullXapiWriter.js')
const { watchStreamSize } = require('../../_watchStreamSize.js')
const { AbstractXapi } = require('./_AbstractXapi.js')
const { debug } = createLogger('xo:backups:FullXapiVmBackup') const { debug } = createLogger('xo:backups:FullXapiVmBackup')
export const FullXapi = class FullXapiVmBackupRunner extends AbstractXapi { exports.FullXapi = class FullXapiVmBackupRunner extends AbstractXapi {
_getWriters() { _getWriters() {
return [FullRemoteWriter, FullXapiWriter] return [FullRemoteWriter, FullXapiWriter]
} }
@@ -35,23 +37,13 @@ export const FullXapi = class FullXapiVmBackupRunner extends AbstractXapi {
useSnapshot: false, useSnapshot: false,
}) })
) )
const vdis = await exportedVm.$getDisks()
let maxStreamLength = 1024 * 1024 // Ovf file and tar headers are a few KB, let's stay safe
for (const vdiRef of vdis) {
const vdi = await this._xapi.getRecord(vdiRef)
// at most the xva will take the physical usage of the disk
// the resulting stream can be smaller due to the smaller block size for xva than vhd, and compression of xcp-ng
maxStreamLength += vdi.physical_utilisation
}
const sizeContainer = watchStreamSize(stream) const sizeContainer = watchStreamSize(stream)
const timestamp = Date.now() const timestamp = Date.now()
await this._callWriters( await this._callWriters(
writer => writer =>
writer.run({ writer.run({
maxStreamLength,
sizeContainer, sizeContainer,
stream: forkStreamUnpipe(stream), stream: forkStreamUnpipe(stream),
timestamp, timestamp,

View File

@@ -1,14 +1,15 @@
import { asyncEach } from '@vates/async-each' 'use strict'
import { decorateMethodsWith } from '@vates/decorate-with' const assert = require('node:assert')
import { defer } from 'golike-defer'
import assert from 'node:assert'
import isVhdDifferencingDisk from 'vhd-lib/isVhdDifferencingDisk.js'
import mapValues from 'lodash/mapValues.js'
import { AbstractRemote } from './_AbstractRemote.mjs' const { decorateMethodsWith } = require('@vates/decorate-with')
import { forkDeltaExport } from './_forkDeltaExport.mjs' const { defer } = require('golike-defer')
import { IncrementalRemoteWriter } from '../_writers/IncrementalRemoteWriter.mjs' const { mapValues } = require('lodash')
import { Task } from '../../Task.mjs' const { Task } = require('../../Task')
const { AbstractRemote } = require('./_AbstractRemote')
const { IncrementalRemoteWriter } = require('../_writers/IncrementalRemoteWriter')
const { forkDeltaExport } = require('./_forkDeltaExport')
const isVhdDifferencingDisk = require('vhd-lib/isVhdDifferencingDisk')
const { asyncEach } = require('@vates/async-each')
class IncrementalRemoteVmBackupRunner extends AbstractRemote { class IncrementalRemoteVmBackupRunner extends AbstractRemote {
_getRemoteWriter() { _getRemoteWriter() {
@@ -60,7 +61,7 @@ class IncrementalRemoteVmBackupRunner extends AbstractRemote {
} }
} }
export const IncrementalRemote = IncrementalRemoteVmBackupRunner exports.IncrementalRemote = IncrementalRemoteVmBackupRunner
decorateMethodsWith(IncrementalRemoteVmBackupRunner, { decorateMethodsWith(IncrementalRemoteVmBackupRunner, {
_run: defer, _run: defer,
}) })

View File

@@ -1,26 +1,28 @@
import { asyncEach } from '@vates/async-each' 'use strict'
import { asyncMap } from '@xen-orchestra/async-map'
import { createLogger } from '@xen-orchestra/log'
import { pipeline } from 'node:stream'
import findLast from 'lodash/findLast.js'
import isVhdDifferencingDisk from 'vhd-lib/isVhdDifferencingDisk.js'
import keyBy from 'lodash/keyBy.js'
import mapValues from 'lodash/mapValues.js'
import vhdStreamValidator from 'vhd-lib/vhdStreamValidator.js'
import { AbstractXapi } from './_AbstractXapi.mjs' const findLast = require('lodash/findLast.js')
import { exportIncrementalVm } from '../../_incrementalVm.mjs' const keyBy = require('lodash/keyBy.js')
import { forkDeltaExport } from './_forkDeltaExport.mjs' const mapValues = require('lodash/mapValues.js')
import { IncrementalRemoteWriter } from '../_writers/IncrementalRemoteWriter.mjs' const vhdStreamValidator = require('vhd-lib/vhdStreamValidator.js')
import { IncrementalXapiWriter } from '../_writers/IncrementalXapiWriter.mjs' const { asyncMap } = require('@xen-orchestra/async-map')
import { Task } from '../../Task.mjs' const { createLogger } = require('@xen-orchestra/log')
import { watchStreamSize } from '../../_watchStreamSize.mjs' const { pipeline } = require('node:stream')
const { IncrementalRemoteWriter } = require('../_writers/IncrementalRemoteWriter.js')
const { IncrementalXapiWriter } = require('../_writers/IncrementalXapiWriter.js')
const { exportIncrementalVm } = require('../../_incrementalVm.js')
const { Task } = require('../../Task.js')
const { watchStreamSize } = require('../../_watchStreamSize.js')
const { AbstractXapi } = require('./_AbstractXapi.js')
const { forkDeltaExport } = require('./_forkDeltaExport.js')
const isVhdDifferencingDisk = require('vhd-lib/isVhdDifferencingDisk')
const { asyncEach } = require('@vates/async-each')
const { debug } = createLogger('xo:backups:IncrementalXapiVmBackup') const { debug } = createLogger('xo:backups:IncrementalXapiVmBackup')
const noop = Function.prototype const noop = Function.prototype
export const IncrementalXapi = class IncrementalXapiVmBackupRunner extends AbstractXapi { exports.IncrementalXapi = class IncrementalXapiVmBackupRunner extends AbstractXapi {
_getWriters() { _getWriters() {
return [IncrementalRemoteWriter, IncrementalXapiWriter] return [IncrementalRemoteWriter, IncrementalXapiWriter]
} }
@@ -41,7 +43,6 @@ export const IncrementalXapi = class IncrementalXapiVmBackupRunner extends Abstr
const deltaExport = await exportIncrementalVm(exportedVm, baseVm, { const deltaExport = await exportIncrementalVm(exportedVm, baseVm, {
fullVdisRequired, fullVdisRequired,
preferNbd: this._settings.preferNbd,
}) })
// since NBD is network based, if one disk use nbd , all the disk use them // since NBD is network based, if one disk use nbd , all the disk use them
// except the suspended VDI // except the suspended VDI

View File

@@ -1,6 +1,8 @@
import { asyncMap } from '@xen-orchestra/async-map' 'use strict'
import { createLogger } from '@xen-orchestra/log'
import { Task } from '../../Task.mjs' const { asyncMap } = require('@xen-orchestra/async-map')
const { createLogger } = require('@xen-orchestra/log')
const { Task } = require('../../Task.js')
const { debug, warn } = createLogger('xo:backups:AbstractVmRunner') const { debug, warn } = createLogger('xo:backups:AbstractVmRunner')
@@ -17,7 +19,7 @@ const asyncEach = async (iterable, fn, thisArg = iterable) => {
} }
} }
export const Abstract = class AbstractVmBackupRunner { exports.Abstract = class AbstractVmBackupRunner {
// calls fn for each function, warns of any errors, and throws only if there are no writers left // calls fn for each function, warns of any errors, and throws only if there are no writers left
async _callWriters(fn, step, parallel = true) { async _callWriters(fn, step, parallel = true) {
const writers = this._writers const writers = this._writers

View File

@@ -1,12 +1,11 @@
import { asyncEach } from '@vates/async-each' 'use strict'
import { Disposable } from 'promise-toolbox' const { Abstract } = require('./_Abstract')
import { getVmBackupDir } from '../../_getVmBackupDir.mjs' const { getVmBackupDir } = require('../../_getVmBackupDir')
const { asyncEach } = require('@vates/async-each')
const { Disposable } = require('promise-toolbox')
import { Abstract } from './_Abstract.mjs' exports.AbstractRemote = class AbstractRemoteVmBackupRunner extends Abstract {
import { extractIdsFromSimplePattern } from '../../extractIdsFromSimplePattern.mjs'
export const AbstractRemote = class AbstractRemoteVmBackupRunner extends Abstract {
constructor({ constructor({
config, config,
job, job,
@@ -35,8 +34,7 @@ export const AbstractRemote = class AbstractRemoteVmBackupRunner extends Abstrac
this._writers = writers this._writers = writers
const RemoteWriter = this._getRemoteWriter() const RemoteWriter = this._getRemoteWriter()
extractIdsFromSimplePattern(job.remotes).forEach(remoteId => { Object.entries(remoteAdapters).forEach(([remoteId, adapter]) => {
const adapter = remoteAdapters[remoteId]
const targetSettings = { const targetSettings = {
...settings, ...settings,
...allSettings[remoteId], ...allSettings[remoteId],

View File

@@ -1,16 +1,18 @@
import assert from 'node:assert' 'use strict'
import groupBy from 'lodash/groupBy.js'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { asyncMap } from '@xen-orchestra/async-map'
import { decorateMethodsWith } from '@vates/decorate-with'
import { defer } from 'golike-defer'
import { formatDateTime } from '@xen-orchestra/xapi'
import { getOldEntries } from '../../_getOldEntries.mjs' const assert = require('assert')
import { Task } from '../../Task.mjs' const groupBy = require('lodash/groupBy.js')
import { Abstract } from './_Abstract.mjs' const ignoreErrors = require('promise-toolbox/ignoreErrors')
const { asyncMap } = require('@xen-orchestra/async-map')
const { decorateMethodsWith } = require('@vates/decorate-with')
const { defer } = require('golike-defer')
const { formatDateTime } = require('@xen-orchestra/xapi')
export const AbstractXapi = class AbstractXapiVmBackupRunner extends Abstract { const { getOldEntries } = require('../../_getOldEntries.js')
const { Task } = require('../../Task.js')
const { Abstract } = require('./_Abstract.js')
class AbstractXapiVmBackupRunner extends Abstract {
constructor({ constructor({
config, config,
getSnapshotNameLabel, getSnapshotNameLabel,
@@ -31,11 +33,6 @@ export const AbstractXapi = class AbstractXapiVmBackupRunner extends Abstract {
throw new Error('cannot backup a VM created by this very job') throw new Error('cannot backup a VM created by this very job')
} }
const currentOperations = Object.values(vm.current_operations)
if (currentOperations.some(_ => _ === 'migrate_send' || _ === 'pool_migrate')) {
throw new Error('cannot backup a VM currently being migrated')
}
this.config = config this.config = config
this.job = job this.job = job
this.remoteAdapters = remoteAdapters this.remoteAdapters = remoteAdapters
@@ -261,15 +258,7 @@ export const AbstractXapi = class AbstractXapiVmBackupRunner extends Abstract {
} }
if (this._writers.size !== 0) { if (this._writers.size !== 0) {
const { pool_migrate = null, migrate_send = null } = this._exportedVm.blocked_operations
const reason = 'VM migration is blocked during backup'
await this._exportedVm.update_blocked_operations({ pool_migrate: reason, migrate_send: reason })
try {
await this._copy() await this._copy()
} finally {
await this._exportedVm.update_blocked_operations({ pool_migrate, migrate_send })
}
} }
} finally { } finally {
if (startAfter) { if (startAfter) {
@@ -282,7 +271,8 @@ export const AbstractXapi = class AbstractXapiVmBackupRunner extends Abstract {
await this._healthCheck() await this._healthCheck()
} }
} }
exports.AbstractXapi = AbstractXapiVmBackupRunner
decorateMethodsWith(AbstractXapi, { decorateMethodsWith(AbstractXapiVmBackupRunner, {
run: defer, run: defer,
}) })

View File

@@ -0,0 +1,12 @@
'use strict'
const { mapValues } = require('lodash')
const { forkStreamUnpipe } = require('../_forkStreamUnpipe')
exports.forkDeltaExport = function forkDeltaExport(deltaExport) {
return Object.create(deltaExport, {
streams: {
value: mapValues(deltaExport.streams, forkStreamUnpipe),
},
})
}

View File

@@ -1,11 +0,0 @@
import mapValues from 'lodash/mapValues.js'
import { forkStreamUnpipe } from '../_forkStreamUnpipe.mjs'
export function forkDeltaExport(deltaExport) {
return Object.create(deltaExport, {
streams: {
value: mapValues(deltaExport.streams, forkStreamUnpipe),
},
})
}

View File

@@ -1,11 +1,13 @@
import { formatFilenameDate } from '../../_filenameDate.mjs' 'use strict'
import { getOldEntries } from '../../_getOldEntries.mjs'
import { Task } from '../../Task.mjs'
import { MixinRemoteWriter } from './_MixinRemoteWriter.mjs' const { formatFilenameDate } = require('../../_filenameDate.js')
import { AbstractFullWriter } from './_AbstractFullWriter.mjs' const { getOldEntries } = require('../../_getOldEntries.js')
const { Task } = require('../../Task.js')
export class FullRemoteWriter extends MixinRemoteWriter(AbstractFullWriter) { const { MixinRemoteWriter } = require('./_MixinRemoteWriter.js')
const { AbstractFullWriter } = require('./_AbstractFullWriter.js')
exports.FullRemoteWriter = class FullRemoteWriter extends MixinRemoteWriter(AbstractFullWriter) {
constructor(props) { constructor(props) {
super(props) super(props)
@@ -24,7 +26,7 @@ export class FullRemoteWriter extends MixinRemoteWriter(AbstractFullWriter) {
) )
} }
async _run({ maxStreamLength, timestamp, sizeContainer, stream, streamLength, vm, vmSnapshot }) { async _run({ timestamp, sizeContainer, stream, vm, vmSnapshot }) {
const settings = this._settings const settings = this._settings
const job = this._job const job = this._job
const scheduleId = this._scheduleId const scheduleId = this._scheduleId
@@ -65,8 +67,6 @@ export class FullRemoteWriter extends MixinRemoteWriter(AbstractFullWriter) {
await Task.run({ name: 'transfer' }, async () => { await Task.run({ name: 'transfer' }, async () => {
await adapter.outputStream(dataFilename, stream, { await adapter.outputStream(dataFilename, stream, {
maxStreamLength,
streamLength,
validator: tmpPath => adapter.isValidXva(tmpPath), validator: tmpPath => adapter.isValidXva(tmpPath),
}) })
return { size: sizeContainer.size } return { size: sizeContainer.size }

View File

@@ -1,16 +1,18 @@
import ignoreErrors from 'promise-toolbox/ignoreErrors' 'use strict'
import { asyncMap, asyncMapSettled } from '@xen-orchestra/async-map'
import { formatDateTime } from '@xen-orchestra/xapi'
import { formatFilenameDate } from '../../_filenameDate.mjs' const ignoreErrors = require('promise-toolbox/ignoreErrors')
import { getOldEntries } from '../../_getOldEntries.mjs' const { asyncMap, asyncMapSettled } = require('@xen-orchestra/async-map')
import { Task } from '../../Task.mjs' const { formatDateTime } = require('@xen-orchestra/xapi')
import { AbstractFullWriter } from './_AbstractFullWriter.mjs' const { formatFilenameDate } = require('../../_filenameDate.js')
import { MixinXapiWriter } from './_MixinXapiWriter.mjs' const { getOldEntries } = require('../../_getOldEntries.js')
import { listReplicatedVms } from './_listReplicatedVms.mjs' const { Task } = require('../../Task.js')
export class FullXapiWriter extends MixinXapiWriter(AbstractFullWriter) { const { AbstractFullWriter } = require('./_AbstractFullWriter.js')
const { MixinXapiWriter } = require('./_MixinXapiWriter.js')
const { listReplicatedVms } = require('./_listReplicatedVms.js')
exports.FullXapiWriter = class FullXapiWriter extends MixinXapiWriter(AbstractFullWriter) {
constructor(props) { constructor(props) {
super(props) super(props)

View File

@@ -1,27 +1,29 @@
import assert from 'node:assert' 'use strict'
import mapValues from 'lodash/mapValues.js'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { asyncEach } from '@vates/async-each'
import { asyncMap } from '@xen-orchestra/async-map'
import { chainVhd, checkVhdChain, openVhd, VhdAbstract } from 'vhd-lib'
import { createLogger } from '@xen-orchestra/log'
import { decorateClass } from '@vates/decorate-with'
import { defer } from 'golike-defer'
import { dirname } from 'node:path'
import { formatFilenameDate } from '../../_filenameDate.mjs' const assert = require('assert')
import { getOldEntries } from '../../_getOldEntries.mjs' const mapValues = require('lodash/mapValues.js')
import { Task } from '../../Task.mjs' const ignoreErrors = require('promise-toolbox/ignoreErrors')
const { asyncEach } = require('@vates/async-each')
const { asyncMap } = require('@xen-orchestra/async-map')
const { chainVhd, checkVhdChain, openVhd, VhdAbstract } = require('vhd-lib')
const { createLogger } = require('@xen-orchestra/log')
const { decorateClass } = require('@vates/decorate-with')
const { defer } = require('golike-defer')
const { dirname } = require('path')
import { MixinRemoteWriter } from './_MixinRemoteWriter.mjs' const { formatFilenameDate } = require('../../_filenameDate.js')
import { AbstractIncrementalWriter } from './_AbstractIncrementalWriter.mjs' const { getOldEntries } = require('../../_getOldEntries.js')
import { checkVhd } from './_checkVhd.mjs' const { Task } = require('../../Task.js')
import { packUuid } from './_packUuid.mjs'
import { Disposable } from 'promise-toolbox' const { MixinRemoteWriter } = require('./_MixinRemoteWriter.js')
const { AbstractIncrementalWriter } = require('./_AbstractIncrementalWriter.js')
const { checkVhd } = require('./_checkVhd.js')
const { packUuid } = require('./_packUuid.js')
const { Disposable } = require('promise-toolbox')
const { warn } = createLogger('xo:backups:DeltaBackupWriter') const { warn } = createLogger('xo:backups:DeltaBackupWriter')
export class IncrementalRemoteWriter extends MixinRemoteWriter(AbstractIncrementalWriter) { class IncrementalRemoteWriter extends MixinRemoteWriter(AbstractIncrementalWriter) {
async checkBaseVdis(baseUuidToSrcVdi) { async checkBaseVdis(baseUuidToSrcVdi) {
const { handler } = this._adapter const { handler } = this._adapter
const adapter = this._adapter const adapter = this._adapter
@@ -203,7 +205,7 @@ export class IncrementalRemoteWriter extends MixinRemoteWriter(AbstractIncrement
// TODO remove when this has been done before the export // TODO remove when this has been done before the export
await checkVhd(handler, parentPath) await checkVhd(handler, parentPath)
} }
// @todo : sum per property
transferSize += await adapter.writeVhd(path, deltaExport.streams[`${id}.vhd`], { transferSize += await adapter.writeVhd(path, deltaExport.streams[`${id}.vhd`], {
// no checksum for VHDs, because they will be invalidated by // no checksum for VHDs, because they will be invalidated by
// merges and chainings // merges and chainings
@@ -230,12 +232,12 @@ export class IncrementalRemoteWriter extends MixinRemoteWriter(AbstractIncrement
return { size: transferSize } return { size: transferSize }
}) })
metadataContent.size = size metadataContent.size = size // @todo: transferSize
this._metadataFileName = await adapter.writeVmBackupMetadata(vm.uuid, metadataContent) this._metadataFileName = await adapter.writeVmBackupMetadata(vm.uuid, metadataContent)
// TODO: run cleanup? // TODO: run cleanup?
} }
} }
decorateClass(IncrementalRemoteWriter, { exports.IncrementalRemoteWriter = decorateClass(IncrementalRemoteWriter, {
_transfer: defer, _transfer: defer,
}) })

View File

@@ -1,17 +1,19 @@
import { asyncMap, asyncMapSettled } from '@xen-orchestra/async-map' 'use strict'
import ignoreErrors from 'promise-toolbox/ignoreErrors'
import { formatDateTime } from '@xen-orchestra/xapi'
import { formatFilenameDate } from '../../_filenameDate.mjs' const { asyncMap, asyncMapSettled } = require('@xen-orchestra/async-map')
import { getOldEntries } from '../../_getOldEntries.mjs' const ignoreErrors = require('promise-toolbox/ignoreErrors')
import { importIncrementalVm, TAG_COPY_SRC } from '../../_incrementalVm.mjs' const { formatDateTime } = require('@xen-orchestra/xapi')
import { Task } from '../../Task.mjs'
import { AbstractIncrementalWriter } from './_AbstractIncrementalWriter.mjs' const { formatFilenameDate } = require('../../_filenameDate.js')
import { MixinXapiWriter } from './_MixinXapiWriter.mjs' const { getOldEntries } = require('../../_getOldEntries.js')
import { listReplicatedVms } from './_listReplicatedVms.mjs' const { importIncrementalVm, TAG_COPY_SRC } = require('../../_incrementalVm.js')
const { Task } = require('../../Task.js')
export class IncrementalXapiWriter extends MixinXapiWriter(AbstractIncrementalWriter) { const { AbstractIncrementalWriter } = require('./_AbstractIncrementalWriter.js')
const { MixinXapiWriter } = require('./_MixinXapiWriter.js')
const { listReplicatedVms } = require('./_listReplicatedVms.js')
exports.IncrementalXapiWriter = class IncrementalXapiWriter extends MixinXapiWriter(AbstractIncrementalWriter) {
async checkBaseVdis(baseUuidToSrcVdi, baseVm) { async checkBaseVdis(baseUuidToSrcVdi, baseVm) {
const sr = this._sr const sr = this._sr
const replicatedVm = listReplicatedVms(sr.$xapi, this._job.id, sr.uuid, this._vmUuid).find( const replicatedVm = listReplicatedVms(sr.$xapi, this._job.id, sr.uuid, this._vmUuid).find(

View File

@@ -0,0 +1,14 @@
'use strict'
const { AbstractWriter } = require('./_AbstractWriter.js')
exports.AbstractFullWriter = class AbstractFullWriter extends AbstractWriter {
async run({ timestamp, sizeContainer, stream, vm, vmSnapshot }) {
try {
return await this._run({ timestamp, sizeContainer, stream, vm, vmSnapshot })
} finally {
// ensure stream is properly closed
stream.destroy()
}
}
}

View File

@@ -1,12 +0,0 @@
import { AbstractWriter } from './_AbstractWriter.mjs'
export class AbstractFullWriter extends AbstractWriter {
async run({ maxStreamLength, timestamp, sizeContainer, stream, streamLength, vm, vmSnapshot }) {
try {
return await this._run({ maxStreamLength, timestamp, sizeContainer, stream, streamLength, vm, vmSnapshot })
} finally {
// ensure stream is properly closed
stream.destroy()
}
}
}

View File

@@ -1,6 +1,8 @@
import { AbstractWriter } from './_AbstractWriter.mjs' 'use strict'
export class AbstractIncrementalWriter extends AbstractWriter { const { AbstractWriter } = require('./_AbstractWriter.js')
exports.AbstractIncrementalWriter = class AbstractIncrementalWriter extends AbstractWriter {
checkBaseVdis(baseUuidToSrcVdi, baseVm) { checkBaseVdis(baseUuidToSrcVdi, baseVm) {
throw new Error('Not implemented') throw new Error('Not implemented')
} }

View File

@@ -1,7 +1,9 @@
import { formatFilenameDate } from '../../_filenameDate.mjs' 'use strict'
import { getVmBackupDir } from '../../_getVmBackupDir.mjs'
export class AbstractWriter { const { formatFilenameDate } = require('../../_filenameDate')
const { getVmBackupDir } = require('../../_getVmBackupDir')
exports.AbstractWriter = class AbstractWriter {
constructor({ config, healthCheckSr, job, vmUuid, scheduleId, settings }) { constructor({ config, healthCheckSr, job, vmUuid, scheduleId, settings }) {
this._config = config this._config = config
this._healthCheckSr = healthCheckSr this._healthCheckSr = healthCheckSr

View File

@@ -1,17 +1,19 @@
import { createLogger } from '@xen-orchestra/log' 'use strict'
import { join } from 'node:path'
import assert from 'node:assert'
import { formatFilenameDate } from '../../_filenameDate.mjs' const { createLogger } = require('@xen-orchestra/log')
import { getVmBackupDir } from '../../_getVmBackupDir.mjs' const { join } = require('path')
import { HealthCheckVmBackup } from '../../HealthCheckVmBackup.mjs'
import { ImportVmBackup } from '../../ImportVmBackup.mjs' const assert = require('assert')
import { Task } from '../../Task.mjs' const { formatFilenameDate } = require('../../_filenameDate.js')
import * as MergeWorker from '../../merge-worker/index.mjs' const { getVmBackupDir } = require('../../_getVmBackupDir.js')
const { HealthCheckVmBackup } = require('../../HealthCheckVmBackup.js')
const { ImportVmBackup } = require('../../ImportVmBackup.js')
const { Task } = require('../../Task.js')
const MergeWorker = require('../../merge-worker/index.js')
const { info, warn } = createLogger('xo:backups:MixinBackupWriter') const { info, warn } = createLogger('xo:backups:MixinBackupWriter')
export const MixinRemoteWriter = (BaseClass = Object) => exports.MixinRemoteWriter = (BaseClass = Object) =>
class MixinRemoteWriter extends BaseClass { class MixinRemoteWriter extends BaseClass {
#lock #lock

View File

@@ -1,10 +1,12 @@
import { extractOpaqueRef } from '@xen-orchestra/xapi' 'use strict'
import assert from 'node:assert/strict'
import { HealthCheckVmBackup } from '../../HealthCheckVmBackup.mjs' const { extractOpaqueRef } = require('@xen-orchestra/xapi')
import { Task } from '../../Task.mjs'
export const MixinXapiWriter = (BaseClass = Object) => const { Task } = require('../../Task')
const assert = require('node:assert/strict')
const { HealthCheckVmBackup } = require('../../HealthCheckVmBackup')
exports.MixinXapiWriter = (BaseClass = Object) =>
class MixinXapiWriter extends BaseClass { class MixinXapiWriter extends BaseClass {
constructor({ sr, ...rest }) { constructor({ sr, ...rest }) {
super(rest) super(rest)
@@ -18,7 +20,7 @@ export const MixinXapiWriter = (BaseClass = Object) =>
const vdiRefs = await xapi.VM_getDisks(baseVm.$ref) const vdiRefs = await xapi.VM_getDisks(baseVm.$ref)
for (const vdiRef of vdiRefs) { for (const vdiRef of vdiRefs) {
const vdi = xapi.getObject(vdiRef) const vdi = xapi.getObject(vdiRef)
if (vdi.$SR.uuid !== this._healthCheckSr.uuid) { if (vdi.$SR.uuid !== this._heathCheckSr.uuid) {
return false return false
} }
} }

View File

@@ -0,0 +1,8 @@
'use strict'
const openVhd = require('vhd-lib').openVhd
const Disposable = require('promise-toolbox/Disposable')
exports.checkVhd = async function checkVhd(handler, path) {
await Disposable.use(openVhd(handler, path), () => {})
}

View File

@@ -1,6 +0,0 @@
import { openVhd } from 'vhd-lib'
import Disposable from 'promise-toolbox/Disposable'
export async function checkVhd(handler, path) {
await Disposable.use(openVhd(handler, path), () => {})
}

View File

@@ -1,3 +1,5 @@
'use strict'
const getReplicatedVmDatetime = vm => { const getReplicatedVmDatetime = vm => {
const { 'xo:backup:datetime': datetime = vm.name_label.slice(-17, -1) } = vm.other_config const { 'xo:backup:datetime': datetime = vm.name_label.slice(-17, -1) } = vm.other_config
return datetime return datetime
@@ -5,7 +7,7 @@ const getReplicatedVmDatetime = vm => {
const compareReplicatedVmDatetime = (a, b) => (getReplicatedVmDatetime(a) < getReplicatedVmDatetime(b) ? -1 : 1) const compareReplicatedVmDatetime = (a, b) => (getReplicatedVmDatetime(a) < getReplicatedVmDatetime(b) ? -1 : 1)
export function listReplicatedVms(xapi, scheduleOrJobId, srUuid, vmUuid) { exports.listReplicatedVms = function listReplicatedVms(xapi, scheduleOrJobId, srUuid, vmUuid) {
const { all } = xapi.objects const { all } = xapi.objects
const vms = {} const vms = {}
for (const key in all) { for (const key in all) {

View File

@@ -1,5 +1,7 @@
'use strict'
const PARSE_UUID_RE = /-/g const PARSE_UUID_RE = /-/g
export function packUuid(uuid) { exports.packUuid = function packUuid(uuid) {
return Buffer.from(uuid.replace(PARSE_UUID_RE, ''), 'hex') return Buffer.from(uuid.replace(PARSE_UUID_RE, ''), 'hex')
} }

View File

@@ -1,4 +1,6 @@
export function watchStreamSize(stream, container = { size: 0 }) { 'use strict'
exports.watchStreamSize = function watchStreamSize(stream, container = { size: 0 }) {
stream.on('data', data => { stream.on('data', data => {
container.size += data.length container.size += data.length
}) })

View File

@@ -1,4 +1,6 @@
export function extractIdsFromSimplePattern(pattern) { 'use strict'
exports.extractIdsFromSimplePattern = function extractIdsFromSimplePattern(pattern) {
if (pattern === undefined) { if (pattern === undefined) {
return [] return []
} }

View File

@@ -1,5 +1,7 @@
import mapValues from 'lodash/mapValues.js' 'use strict'
import { dirname } from 'node:path'
const mapValues = require('lodash/mapValues.js')
const { dirname } = require('path')
function formatVmBackup(backup) { function formatVmBackup(backup) {
return { return {
@@ -29,6 +31,6 @@ function formatVmBackup(backup) {
} }
// format all backups as returned by RemoteAdapter#listAllVmBackups() // format all backups as returned by RemoteAdapter#listAllVmBackups()
export function formatVmBackups(backupsByVM) { exports.formatVmBackups = function formatVmBackups(backupsByVM) {
return mapValues(backupsByVM, backups => backups.map(formatVmBackup)) return mapValues(backupsByVM, backups => backups.map(formatVmBackup))
} }

View File

@@ -0,0 +1,92 @@
#!/usr/bin/env node
// eslint-disable-next-line eslint-comments/disable-enable-pair
/* eslint-disable n/shebang */
'use strict'
const { catchGlobalErrors } = require('@xen-orchestra/log/configure')
const { createLogger } = require('@xen-orchestra/log')
const { getSyncedHandler } = require('@xen-orchestra/fs')
const { join } = require('path')
const Disposable = require('promise-toolbox/Disposable')
const min = require('lodash/min')
const { getVmBackupDir } = require('../_getVmBackupDir.js')
const { RemoteAdapter } = require('../RemoteAdapter.js')
const { CLEAN_VM_QUEUE } = require('./index.js')
// -------------------------------------------------------------------
catchGlobalErrors(createLogger('xo:backups:mergeWorker'))
const { fatal, info, warn } = createLogger('xo:backups:mergeWorker')
// -------------------------------------------------------------------
const main = Disposable.wrap(async function* main(args) {
const handler = yield getSyncedHandler({ url: 'file://' + process.cwd() })
yield handler.lock(CLEAN_VM_QUEUE)
const adapter = new RemoteAdapter(handler)
const listRetry = async () => {
const timeoutResolver = resolve => setTimeout(resolve, 10e3)
for (let i = 0; i < 10; ++i) {
const entries = await handler.list(CLEAN_VM_QUEUE)
if (entries.length !== 0) {
return entries
}
await new Promise(timeoutResolver)
}
}
let taskFiles
while ((taskFiles = await listRetry()) !== undefined) {
const taskFileBasename = min(taskFiles)
const previousTaskFile = join(CLEAN_VM_QUEUE, taskFileBasename)
const taskFile = join(CLEAN_VM_QUEUE, '_' + taskFileBasename)
// move this task to the end
try {
await handler.rename(previousTaskFile, taskFile)
} catch (error) {
// this error occurs if the task failed too many times (i.e. too many `_` prefixes)
// there is nothing more that can be done
if (error.code === 'ENAMETOOLONG') {
await handler.unlink(previousTaskFile)
}
throw error
}
try {
const vmDir = getVmBackupDir(String(await handler.readFile(taskFile)))
try {
await adapter.cleanVm(vmDir, { merge: true, logInfo: info, logWarn: warn, remove: true })
} catch (error) {
// consider the clean successful if the VM dir is missing
if (error.code !== 'ENOENT') {
throw error
}
}
handler.unlink(taskFile).catch(error => warn('deleting task failure', { error }))
} catch (error) {
warn('failure handling task', { error })
}
}
})
info('starting')
main(process.argv.slice(2)).then(
() => {
info('bye :-)')
},
error => {
fatal(error)
process.exit(1)
}
)

View File

@@ -1,103 +0,0 @@
#!/usr/bin/env node
// eslint-disable-next-line eslint-comments/disable-enable-pair
/* eslint-disable n/shebang */
import { asyncEach } from '@vates/async-each'
import { catchGlobalErrors } from '@xen-orchestra/log/configure'
import { createLogger } from '@xen-orchestra/log'
import { getSyncedHandler } from '@xen-orchestra/fs'
import { join } from 'node:path'
import { load as loadConfig } from 'app-conf'
import Disposable from 'promise-toolbox/Disposable'
import { getVmBackupDir } from '../_getVmBackupDir.mjs'
import { RemoteAdapter } from '../RemoteAdapter.mjs'
import { CLEAN_VM_QUEUE } from './index.mjs'
const APP_NAME = 'xo-merge-worker'
const APP_DIR = new URL('.', import.meta.url).pathname
// -------------------------------------------------------------------
catchGlobalErrors(createLogger('xo:backups:mergeWorker'))
const { fatal, info, warn } = createLogger('xo:backups:mergeWorker')
// -------------------------------------------------------------------
const main = Disposable.wrap(async function* main(args) {
const handler = yield getSyncedHandler({ url: 'file://' + process.cwd() })
yield handler.lock(CLEAN_VM_QUEUE)
const adapter = new RemoteAdapter(handler)
const listRetry = async () => {
const timeoutResolver = resolve => setTimeout(resolve, 10e3)
for (let i = 0; i < 10; ++i) {
const entries = await handler.list(CLEAN_VM_QUEUE)
if (entries.length !== 0) {
entries.sort()
return entries
}
await new Promise(timeoutResolver)
}
}
let taskFiles
while ((taskFiles = await listRetry()) !== undefined) {
const { concurrency } = await loadConfig(APP_NAME, {
appDir: APP_DIR,
ignoreUnknownFormats: true,
})
await asyncEach(
taskFiles,
async taskFileBasename => {
const previousTaskFile = join(CLEAN_VM_QUEUE, taskFileBasename)
const taskFile = join(CLEAN_VM_QUEUE, '_' + taskFileBasename)
// move this task to the end
try {
await handler.rename(previousTaskFile, taskFile)
} catch (error) {
// this error occurs if the task failed too many times (i.e. too many `_` prefixes)
// there is nothing more that can be done
if (error.code === 'ENAMETOOLONG') {
await handler.unlink(previousTaskFile)
}
throw error
}
try {
const vmDir = getVmBackupDir(String(await handler.readFile(taskFile)))
try {
await adapter.cleanVm(vmDir, { merge: true, logInfo: info, logWarn: warn, remove: true })
} catch (error) {
// consider the clean successful if the VM dir is missing
if (error.code !== 'ENOENT') {
throw error
}
}
handler.unlink(taskFile).catch(error => warn('deleting task failure', { error }))
} catch (error) {
warn('failure handling task', { error })
}
},
{ concurrency }
)
}
})
info('starting')
main(process.argv.slice(2)).then(
() => {
info('bye :-)')
},
error => {
fatal(error)
process.exit(1)
}
)

View File

@@ -1 +0,0 @@
concurrency = 1

View File

@@ -1,12 +1,13 @@
import { join } from 'node:path' 'use strict'
import { spawn } from 'child_process'
import { check } from 'proper-lockfile'
export const CLEAN_VM_QUEUE = '/xo-vm-backups/.queue/clean-vm/' const { join, resolve } = require('path')
const { spawn } = require('child_process')
const { check } = require('proper-lockfile')
const CLI_PATH = new URL('cli.mjs', import.meta.url).pathname const CLEAN_VM_QUEUE = (exports.CLEAN_VM_QUEUE = '/xo-vm-backups/.queue/clean-vm/')
export const run = async function runMergeWorker(remotePath) { const CLI_PATH = resolve(__dirname, 'cli.js')
exports.run = async function runMergeWorker(remotePath) {
try { try {
// TODO: find a way to pass the acquire the lock and then pass it down the worker // TODO: find a way to pass the acquire the lock and then pass it down the worker
if (await check(join(remotePath, CLEAN_VM_QUEUE))) { if (await check(join(remotePath, CLEAN_VM_QUEUE))) {

View File

@@ -8,33 +8,32 @@
"type": "git", "type": "git",
"url": "https://github.com/vatesfr/xen-orchestra.git" "url": "https://github.com/vatesfr/xen-orchestra.git"
}, },
"version": "0.43.0", "version": "0.39.0",
"engines": { "engines": {
"node": ">=14.18" "node": ">=14.18"
}, },
"scripts": { "scripts": {
"postversion": "npm publish --access public", "postversion": "npm publish --access public",
"test-integration": "node--test *.integ.mjs" "test-integration": "node--test *.integ.js"
}, },
"dependencies": { "dependencies": {
"@iarna/toml": "^2.2.5",
"@kldzj/stream-throttle": "^1.1.1", "@kldzj/stream-throttle": "^1.1.1",
"@vates/async-each": "^1.0.0", "@vates/async-each": "^1.0.0",
"@vates/cached-dns.lookup": "^1.0.0", "@vates/cached-dns.lookup": "^1.0.0",
"@vates/compose": "^2.1.0", "@vates/compose": "^2.1.0",
"@vates/decorate-with": "^2.0.0", "@vates/decorate-with": "^2.0.0",
"@vates/disposable": "^0.1.4", "@vates/disposable": "^0.1.4",
"@vates/fuse-vhd": "^2.0.0", "@vates/fuse-vhd": "^1.0.0",
"@vates/nbd-client": "^2.0.0", "@vates/nbd-client": "^1.2.1",
"@vates/parse-duration": "^0.1.1", "@vates/parse-duration": "^0.1.1",
"@xen-orchestra/async-map": "^0.1.2", "@xen-orchestra/async-map": "^0.1.2",
"@xen-orchestra/fs": "^4.1.0", "@xen-orchestra/fs": "^4.0.1",
"@xen-orchestra/log": "^0.6.0", "@xen-orchestra/log": "^0.6.0",
"@xen-orchestra/template": "^0.1.0", "@xen-orchestra/template": "^0.1.0",
"app-conf": "^2.3.0", "compare-versions": "^5.0.1",
"compare-versions": "^6.0.0", "d3-time-format": "^3.0.0",
"d3-time-format": "^4.1.0",
"decorator-synchronized": "^0.6.0", "decorator-synchronized": "^0.6.0",
"fs-extra": "^11.1.0",
"golike-defer": "^0.5.1", "golike-defer": "^0.5.1",
"limit-concurrency-decorator": "^0.5.0", "limit-concurrency-decorator": "^0.5.0",
"lodash": "^4.17.20", "lodash": "^4.17.20",
@@ -42,21 +41,19 @@
"parse-pairs": "^2.0.0", "parse-pairs": "^2.0.0",
"promise-toolbox": "^0.21.0", "promise-toolbox": "^0.21.0",
"proper-lockfile": "^4.1.2", "proper-lockfile": "^4.1.2",
"tar": "^6.1.15",
"uuid": "^9.0.0", "uuid": "^9.0.0",
"vhd-lib": "^4.6.1", "vhd-lib": "^4.5.0",
"xen-api": "^1.3.6", "xen-api": "^1.3.3",
"yazl": "^2.5.1" "yazl": "^2.5.1"
}, },
"devDependencies": { "devDependencies": {
"fs-extra": "^11.1.0",
"rimraf": "^5.0.1", "rimraf": "^5.0.1",
"sinon": "^16.0.0", "sinon": "^15.0.1",
"test": "^3.2.1", "test": "^3.2.1",
"tmp": "^0.2.1" "tmp": "^0.2.1"
}, },
"peerDependencies": { "peerDependencies": {
"@xen-orchestra/xapi": "^3.2.0" "@xen-orchestra/xapi": "^2.2.1"
}, },
"license": "AGPL-3.0-or-later", "license": "AGPL-3.0-or-later",
"author": { "author": {

View File

@@ -1,6 +1,8 @@
import { DIR_XO_CONFIG_BACKUPS, DIR_XO_POOL_METADATA_BACKUPS } from './RemoteAdapter.mjs' 'use strict'
export function parseMetadataBackupId(backupId) { const { DIR_XO_CONFIG_BACKUPS, DIR_XO_POOL_METADATA_BACKUPS } = require('./RemoteAdapter.js')
exports.parseMetadataBackupId = function parseMetadataBackupId(backupId) {
const [dir, ...rest] = backupId.split('/') const [dir, ...rest] = backupId.split('/')
if (dir === DIR_XO_CONFIG_BACKUPS) { if (dir === DIR_XO_CONFIG_BACKUPS) {
const [scheduleId, timestamp] = rest const [scheduleId, timestamp] = rest

View File

@@ -1,11 +1,14 @@
import { createLogger } from '@xen-orchestra/log' 'use strict'
import { fork } from 'child_process'
const path = require('path')
const { createLogger } = require('@xen-orchestra/log')
const { fork } = require('child_process')
const { warn } = createLogger('xo:backups:backupWorker') const { warn } = createLogger('xo:backups:backupWorker')
const PATH = new URL('_backupWorker.mjs', import.meta.url).pathname const PATH = path.resolve(__dirname, '_backupWorker.js')
export function runBackupWorker(params, onLog) { exports.runBackupWorker = function runBackupWorker(params, onLog) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
const worker = fork(PATH) const worker = fork(PATH)

View File

@@ -1,5 +1,7 @@
'use strict'
// a valid footer of a 2 // a valid footer of a 2
export const VHDFOOTER = { exports.VHDFOOTER = {
cookie: 'conectix', cookie: 'conectix',
features: 2, features: 2,
fileFormatVersion: 65536, fileFormatVersion: 65536,
@@ -18,7 +20,7 @@ export const VHDFOOTER = {
hidden: '', hidden: '',
reserved: '', reserved: '',
} }
export const VHDHEADER = { exports.VHDHEADER = {
cookie: 'cxsparse', cookie: 'cxsparse',
dataOffset: undefined, dataOffset: undefined,
tableOffset: 2048, tableOffset: 2048,

View File

@@ -18,7 +18,7 @@
"preferGlobal": true, "preferGlobal": true,
"dependencies": { "dependencies": {
"golike-defer": "^0.5.1", "golike-defer": "^0.5.1",
"xen-api": "^1.3.6" "xen-api": "^1.3.3"
}, },
"scripts": { "scripts": {
"postversion": "npm publish" "postversion": "npm publish"

View File

@@ -42,7 +42,7 @@
"test": "node--test" "test": "node--test"
}, },
"devDependencies": { "devDependencies": {
"sinon": "^16.0.0", "sinon": "^15.0.1",
"test": "^3.2.1" "test": "^3.2.1"
} }
} }

View File

@@ -1,7 +1,7 @@
{ {
"private": false, "private": false,
"name": "@xen-orchestra/fs", "name": "@xen-orchestra/fs",
"version": "4.1.0", "version": "4.0.1",
"license": "AGPL-3.0-or-later", "license": "AGPL-3.0-or-later",
"description": "The File System for Xen Orchestra backups.", "description": "The File System for Xen Orchestra backups.",
"homepage": "https://github.com/vatesfr/xen-orchestra/tree/master/@xen-orchestra/fs", "homepage": "https://github.com/vatesfr/xen-orchestra/tree/master/@xen-orchestra/fs",
@@ -29,7 +29,7 @@
"@vates/async-each": "^1.0.0", "@vates/async-each": "^1.0.0",
"@vates/coalesce-calls": "^0.1.0", "@vates/coalesce-calls": "^0.1.0",
"@vates/decorate-with": "^2.0.0", "@vates/decorate-with": "^2.0.0",
"@vates/read-chunk": "^1.2.0", "@vates/read-chunk": "^1.1.1",
"@xen-orchestra/log": "^0.6.0", "@xen-orchestra/log": "^0.6.0",
"bind-property-descriptor": "^2.0.0", "bind-property-descriptor": "^2.0.0",
"decorator-synchronized": "^0.6.0", "decorator-synchronized": "^0.6.0",
@@ -53,7 +53,7 @@
"cross-env": "^7.0.2", "cross-env": "^7.0.2",
"dotenv": "^16.0.0", "dotenv": "^16.0.0",
"rimraf": "^5.0.1", "rimraf": "^5.0.1",
"sinon": "^16.0.0", "sinon": "^15.0.4",
"test": "^3.3.0", "test": "^3.3.0",
"tmp": "^0.2.1" "tmp": "^0.2.1"
}, },

View File

@@ -189,7 +189,7 @@ export default class RemoteHandlerAbstract {
* @param {number} [options.dirMode] * @param {number} [options.dirMode]
* @param {(this: RemoteHandlerAbstract, path: string) => Promise<undefined>} [options.validator] Function that will be called before the data is commited to the remote, if it fails, file should not exist * @param {(this: RemoteHandlerAbstract, path: string) => Promise<undefined>} [options.validator] Function that will be called before the data is commited to the remote, if it fails, file should not exist
*/ */
async outputStream(path, input, { checksum = true, dirMode, maxStreamLength, streamLength, validator } = {}) { async outputStream(path, input, { checksum = true, dirMode, validator } = {}) {
path = normalizePath(path) path = normalizePath(path)
let checksumStream let checksumStream
@@ -201,8 +201,6 @@ export default class RemoteHandlerAbstract {
} }
await this._outputStream(path, input, { await this._outputStream(path, input, {
dirMode, dirMode,
maxStreamLength,
streamLength,
validator, validator,
}) })
if (checksum) { if (checksum) {
@@ -626,18 +624,14 @@ export default class RemoteHandlerAbstract {
const files = await this._list(dir) const files = await this._list(dir)
await asyncEach(files, file => await asyncEach(files, file =>
this._unlink(`${dir}/${file}`).catch( this._unlink(`${dir}/${file}`).catch(error => {
error => {
// Unlink dir behavior is not consistent across platforms // Unlink dir behavior is not consistent across platforms
// https://github.com/nodejs/node-v0.x-archive/issues/5791 // https://github.com/nodejs/node-v0.x-archive/issues/5791
if (error.code === 'EISDIR' || error.code === 'EPERM') { if (error.code === 'EISDIR' || error.code === 'EPERM') {
return this._rmtree(`${dir}/${file}`) return this._rmtree(`${dir}/${file}`)
} }
throw error throw error
}, })
// real unlink concurrency will be 2**max directory depth
{ concurrency: 2 }
)
) )
return this._rmtree(dir) return this._rmtree(dir)
} }

View File

@@ -209,7 +209,7 @@ describe('encryption', () => {
// encrypt with a non default algorithm // encrypt with a non default algorithm
const encryptor = _getEncryptor('aes-256-cbc', '73c1838d7d8a6088ca2317fb5f29cd91') const encryptor = _getEncryptor('aes-256-cbc', '73c1838d7d8a6088ca2317fb5f29cd91')
await fs.writeFile(`${dir}/encryption.json`, `{"algorithm": "aes-256-gcm"}`) await fs.writeFile(`${dir}/encryption.json`, `{"algorithm": "aes-256-gmc"}`)
await fs.writeFile(`${dir}/metadata.json`, encryptor.encryptData(`{"random": "NOTSORANDOM"}`)) await fs.writeFile(`${dir}/metadata.json`, encryptor.encryptData(`{"random": "NOTSORANDOM"}`))
// remote is now non empty : can't modify key anymore // remote is now non empty : can't modify key anymore

View File

@@ -5,7 +5,6 @@ import {
CreateMultipartUploadCommand, CreateMultipartUploadCommand,
DeleteObjectCommand, DeleteObjectCommand,
GetObjectCommand, GetObjectCommand,
GetObjectLockConfigurationCommand,
HeadObjectCommand, HeadObjectCommand,
ListObjectsV2Command, ListObjectsV2Command,
PutObjectCommand, PutObjectCommand,
@@ -17,22 +16,21 @@ import { NodeHttpHandler } from '@aws-sdk/node-http-handler'
import { getApplyMd5BodyChecksumPlugin } from '@aws-sdk/middleware-apply-body-checksum' import { getApplyMd5BodyChecksumPlugin } from '@aws-sdk/middleware-apply-body-checksum'
import { Agent as HttpAgent } from 'http' import { Agent as HttpAgent } from 'http'
import { Agent as HttpsAgent } from 'https' import { Agent as HttpsAgent } from 'https'
import pRetry from 'promise-toolbox/retry'
import { createLogger } from '@xen-orchestra/log' import { createLogger } from '@xen-orchestra/log'
import { PassThrough, Transform, pipeline } from 'stream' import { decorateWith } from '@vates/decorate-with'
import { PassThrough, pipeline } from 'stream'
import { parse } from 'xo-remote-parser' import { parse } from 'xo-remote-parser'
import copyStreamToBuffer from './_copyStreamToBuffer.js' import copyStreamToBuffer from './_copyStreamToBuffer.js'
import guessAwsRegion from './_guessAwsRegion.js' import guessAwsRegion from './_guessAwsRegion.js'
import RemoteHandlerAbstract from './abstract' import RemoteHandlerAbstract from './abstract'
import { basename, join, split } from './path' import { basename, join, split } from './path'
import { asyncEach } from '@vates/async-each' import { asyncEach } from '@vates/async-each'
import { pRetry } from 'promise-toolbox'
// endpoints https://docs.aws.amazon.com/general/latest/gr/s3.html // endpoints https://docs.aws.amazon.com/general/latest/gr/s3.html
// limits: https://docs.aws.amazon.com/AmazonS3/latest/dev/qfacts.html // limits: https://docs.aws.amazon.com/AmazonS3/latest/dev/qfacts.html
const MAX_PART_SIZE = 1024 * 1024 * 1024 * 5 // 5GB const MAX_PART_SIZE = 1024 * 1024 * 1024 * 5 // 5GB
const MAX_PART_NUMBER = 10000
const MIN_PART_SIZE = 5 * 1024 * 1024
const { warn } = createLogger('xo:fs:s3') const { warn } = createLogger('xo:fs:s3')
export default class S3Handler extends RemoteHandlerAbstract { export default class S3Handler extends RemoteHandlerAbstract {
@@ -74,47 +72,12 @@ export default class S3Handler extends RemoteHandlerAbstract {
}), }),
}) })
// Workaround for https://github.com/aws/aws-sdk-js-v3/issues/2673
this.#s3.middlewareStack.use(getApplyMd5BodyChecksumPlugin(this.#s3.config))
const parts = split(path) const parts = split(path)
this.#bucket = parts.shift() this.#bucket = parts.shift()
this.#dir = join(...parts) this.#dir = join(...parts)
const WITH_RETRY = [
'_closeFile',
'_copy',
'_getInfo',
'_getSize',
'_list',
'_mkdir',
'_openFile',
'_outputFile',
'_read',
'_readFile',
'_rename',
'_rmdir',
'_truncate',
'_unlink',
'_write',
'_writeFile',
]
WITH_RETRY.forEach(functionName => {
if (this[functionName] !== undefined) {
// adding the retry on the top level mtehod won't
// cover when _functionName are called internally
this[functionName] = pRetry.wrap(this[functionName], {
delays: [100, 200, 500, 1000, 2000],
// these errors should not change on retry
when: err => !['EEXIST', 'EISDIR', 'ENOTEMPTY', 'ENOENT', 'ENOTDIR', 'EISDIR'].includes(err?.code),
onRetry(error) {
warn('retrying method on fs ', {
method: functionName,
attemptNumber: this.attemptNumber,
delay: this.delay,
error,
file: this.arguments?.[0],
})
},
})
}
})
} }
get type() { get type() {
@@ -223,35 +186,11 @@ export default class S3Handler extends RemoteHandlerAbstract {
} }
} }
async _outputStream(path, input, { streamLength, maxStreamLength = streamLength, validator }) { async _outputStream(path, input, { validator }) {
// S3 storage is limited to 10K part, each part is limited to 5GB. And the total upload must be smaller than 5TB
// a bigger partSize increase the memory consumption of aws/lib-storage exponentially
let partSize
if (maxStreamLength === undefined) {
warn(`Writing ${path} to a S3 remote without a max size set will cut it to 50GB`, { path })
partSize = MIN_PART_SIZE // min size for S3
} else {
partSize = Math.min(Math.max(Math.ceil(maxStreamLength / MAX_PART_NUMBER), MIN_PART_SIZE), MAX_PART_SIZE)
}
// ensure we don't try to upload a stream to big for this partSize
let readCounter = 0
const MAX_SIZE = MAX_PART_NUMBER * partSize
const streamCutter = new Transform({
transform(chunk, encoding, callback) {
readCounter += chunk.length
if (readCounter > MAX_SIZE) {
callback(new Error(`read ${readCounter} bytes, maximum size allowed is ${MAX_SIZE} `))
} else {
callback(null, chunk)
}
},
})
// Workaround for "ReferenceError: ReadableStream is not defined" // Workaround for "ReferenceError: ReadableStream is not defined"
// https://github.com/aws/aws-sdk-js-v3/issues/2522 // https://github.com/aws/aws-sdk-js-v3/issues/2522
const Body = new PassThrough() const Body = new PassThrough()
pipeline(input, streamCutter, Body, () => {}) pipeline(input, Body, () => {})
const upload = new Upload({ const upload = new Upload({
client: this.#s3, client: this.#s3,
@@ -259,8 +198,6 @@ export default class S3Handler extends RemoteHandlerAbstract {
...this.#createParams(path), ...this.#createParams(path),
Body, Body,
}, },
partSize,
leavePartsOnError: false,
}) })
await upload.done() await upload.done()
@@ -275,6 +212,21 @@ export default class S3Handler extends RemoteHandlerAbstract {
} }
} }
// some objectstorage provider like backblaze, can answer a 500/503 routinely
// in this case we should retry, and let their load balancing do its magic
// https://www.backblaze.com/b2/docs/calling.html#error_handling
@decorateWith(pRetry.wrap, {
delays: [100, 200, 500, 1000, 2000],
when: e => e.$metadata?.httpStatusCode === 500,
onRetry(error) {
warn('retrying writing file', {
attemptNumber: this.attemptNumber,
delay: this.delay,
error,
file: this.arguments[0],
})
},
})
async _writeFile(file, data, options) { async _writeFile(file, data, options) {
return this.#s3.send( return this.#s3.send(
new PutObjectCommand({ new PutObjectCommand({
@@ -444,24 +396,6 @@ export default class S3Handler extends RemoteHandlerAbstract {
async _closeFile(fd) {} async _closeFile(fd) {}
async _sync() {
await super._sync()
try {
// if Object Lock is enabled, each upload must come with a contentMD5 header
// the computation of this md5 is memory-intensive, especially when uploading a stream
const res = await this.#s3.send(new GetObjectLockConfigurationCommand({ Bucket: this.#bucket }))
if (res.ObjectLockConfiguration?.ObjectLockEnabled === 'Enabled') {
// Workaround for https://github.com/aws/aws-sdk-js-v3/issues/2673
// will automatically add the contentMD5 header to any upload to S3
this.#s3.middlewareStack.use(getApplyMd5BodyChecksumPlugin(this.#s3.config))
}
} catch (error) {
if (error.Code !== 'ObjectLockConfigurationNotFoundError') {
throw error
}
}
}
useVhdDirectory() { useVhdDirectory() {
return true return true
} }

View File

@@ -1,4 +1,2 @@
// Keeping this file to prevent applying the global monorepo config for now // Keeping this file to prevent applying the global monorepo config for now
module.exports = { module.exports = {};
trailingComma: "es5",
};

View File

@@ -1,29 +1,6 @@
# ChangeLog # ChangeLog
## **next** ## **0.2.0**
- Ability to snapshot/copy a VM from its view (PR [#7087](https://github.com/vatesfr/xen-orchestra/pull/7087))
- [Header] Replace logo with "XO LITE" (PR [#7118](https://github.com/vatesfr/xen-orchestra/pull/7118))
- New VM console toolbar + Ability to send Ctrl+Alt+Del (PR [#7088](https://github.com/vatesfr/xen-orchestra/pull/7088))
## **0.1.4** (2023-10-03)
- Ability to migrate selected VMs to another host (PR [#7040](https://github.com/vatesfr/xen-orchestra/pull/7040))
- Ability to snapshot selected VMs (PR [#7021](https://github.com/vatesfr/xen-orchestra/pull/7021))
- Add Patches to Pool Dashboard (PR [#6709](https://github.com/vatesfr/xen-orchestra/pull/6709))
- Add remember me checkbox on the login page (PR [#7030](https://github.com/vatesfr/xen-orchestra/pull/7030))
## **0.1.3** (2023-09-01)
- Add Alarms to Pool Dashboard (PR [#6976](https://github.com/vatesfr/xen-orchestra/pull/6976))
## **0.1.2** (2023-07-28)
- Ability to export selected VMs as CSV file (PR [#6915](https://github.com/vatesfr/xen-orchestra/pull/6915))
- [Pool/VMs] Ability to export selected VMs as JSON file (PR [#6911](https://github.com/vatesfr/xen-orchestra/pull/6911))
- Add Tasks to Pool Dashboard (PR [#6713](https://github.com/vatesfr/xen-orchestra/pull/6713))
## **0.1.1** (2023-07-03)
- Invalidate sessionId token after logout (PR [#6480](https://github.com/vatesfr/xen-orchestra/pull/6480)) - Invalidate sessionId token after logout (PR [#6480](https://github.com/vatesfr/xen-orchestra/pull/6480))
- Settings page (PR [#6418](https://github.com/vatesfr/xen-orchestra/pull/6418)) - Settings page (PR [#6418](https://github.com/vatesfr/xen-orchestra/pull/6418))

Some files were not shown because too many files have changed in this diff Show More