test(ql3): gate failed plugin upgrades before activation

This commit is contained in:
whyour
2026-08-14 13:57:16 +08:00
parent 056eebef94
commit 2485d61e8a
6 changed files with 736 additions and 256 deletions
@@ -6,14 +6,9 @@ const fs = require('node:fs');
const https = require('node:https');
const path = require('node:path');
const { createRequire } = require('node:module');
const {
createHash,
generateKeyPairSync,
sign,
} = require('node:crypto');
const { createHash, generateKeyPairSync, sign } = require('node:crypto');
const FIXTURE_SCHEMA =
'qinglong/plugin-package-recovery-e2e-fixture@v1';
const FIXTURE_SCHEMA = 'qinglong/plugin-package-recovery-e2e-fixture@v1';
const REGISTRY_EVENT_SCHEMA =
'qinglong/plugin-package-recovery-e2e-registry-event@v1';
const OCI_MANIFEST = 'application/vnd.oci.image.manifest.v1+json';
@@ -74,10 +69,7 @@ function tarHeader(entryPath, size) {
Buffer.from('ustar\0').copy(header, 257);
Buffer.from('00').copy(header, 263);
const checksum = header.reduce((total, byte) => total + byte, 0);
Buffer.from(`${checksum.toString(8).padStart(6, '0')}\0 `).copy(
header,
148,
);
Buffer.from(`${checksum.toString(8).padStart(6, '0')}\0 `).copy(header, 148);
return header;
}
@@ -92,18 +84,17 @@ function canonicalTar(entries) {
return Buffer.concat(parts);
}
function pluginManifest(architecture) {
const {
PLUGIN_PACKAGE_API_VERSION,
PLUGIN_PACKAGE_KIND,
} = ql3Require('@qinglong/runtime-core/plugin-package');
function pluginManifest(architecture, version, invalidUpgrade = false) {
const { PLUGIN_PACKAGE_API_VERSION, PLUGIN_PACKAGE_KIND } = ql3Require(
'@qinglong/runtime-core/plugin-package',
);
return {
apiVersion: PLUGIN_PACKAGE_API_VERSION,
kind: PLUGIN_PACKAGE_KIND,
metadata: {
name: 'e2e-monitor',
displayName: 'E2E Monitor',
version: '1.0.0',
version,
description: 'One bounded end-to-end recovery package',
license: 'Apache-2.0',
},
@@ -121,13 +112,55 @@ function pluginManifest(architecture) {
permissions: {
network: { allowedHosts: [] },
secrets: [],
tools: [],
tools: invalidUpgrade ? ['system.command'] : [],
},
contents: { tasks: [], workflows: [], prompts: [], tools: [] },
contents: invalidUpgrade
? {
tasks: ['tasks/noop.json'],
workflows: ['workflows/cycle.json'],
prompts: [],
tools: [],
}
: { tasks: [], workflows: [], prompts: [], tools: [] },
},
};
}
function invalidUpgradeResources() {
return Object.freeze({
'tasks/noop.json': Object.freeze({
schema: 'qinglong/plugin-package-task-resource@v1',
id: 'noop',
name: 'No-op',
labels: Object.freeze({}),
enabled: true,
kind: 'command',
spec: Object.freeze({
schema: 'qinglong/command@v1',
config: Object.freeze({
command: Object.freeze({
kind: 'argv',
file: '/usr/bin/printf',
args: Object.freeze(['ok']),
}),
environment: Object.freeze([]),
timeoutMs: 30_000,
}),
}),
}),
'workflows/cycle.json': Object.freeze({
schema: 'qinglong/plugin-package-workflow-resource@v1',
id: 'cycle',
name: 'Rejected cyclic workflow',
enabled: true,
steps: Object.freeze([
Object.freeze({ id: 'first', task: 'noop', needs: ['second'] }),
Object.freeze({ id: 'second', task: 'noop', needs: ['first'] }),
]),
}),
});
}
function route(path, mediaType, body) {
return Object.freeze({
path,
@@ -137,44 +170,38 @@ function route(path, mediaType, body) {
});
}
function createFixture({ registry, architecture, createdAtMs = Date.now() }) {
if (
typeof registry !== 'string' ||
!/^[a-z0-9](?:[-a-z0-9.]{0,251}[a-z0-9])?$/.test(registry) ||
!['amd64', 'arm64'].includes(architecture) ||
!Number.isSafeInteger(createdAtMs) ||
createdAtMs < 1
) {
throw new TypeError('Plugin Package E2E fixture options are invalid');
}
const {
PluginPackagePublisherTrustRegistry,
PLUGIN_PACKAGE_SIGNATURE_SCHEMA,
pluginPackageContentTreeDigest,
pluginPackagePublisherSignaturePayload,
} = ql3Require('@qinglong/runtime-core/plugin-package-bundle');
const {
planPluginPackageInstall,
} = ql3Require('@qinglong/runtime-core/plugin-package');
const {
createPluginPackageLock,
pluginPackageInstallActionDigest,
pluginPackageInstallPlanDigest,
serializePluginPackageManifest,
} = ql3Require('@qinglong/runtime-core/plugin-package-install');
function packageMaterial(registry, manifest, resourceValues) {
const { pluginPackageContentTreeDigest } = ql3Require(
'@qinglong/runtime-core/plugin-package-bundle',
);
const { serializePluginPackageManifest } = ql3Require(
'@qinglong/runtime-core/plugin-package-install',
);
const {
PLUGIN_PACKAGE_OCI_ARTIFACT_TYPE,
PLUGIN_PACKAGE_OCI_CONFIG_MEDIA_TYPE,
PLUGIN_PACKAGE_OCI_SIGNATURE_ARTIFACT_TYPE,
PLUGIN_PACKAGE_OCI_SIGNATURE_CONFIG_MEDIA_TYPE,
} = ql3Require('@qinglong/cluster-admin/plugin-package-oci-stage');
const manifest = pluginManifest(architecture);
const resourceEntries = Object.entries(resourceValues)
.map(([entryPath, value]) => ({
path: entryPath,
body: jsonBytes(value),
}))
.sort((left, right) => left.path.localeCompare(right.path));
const contentDigest = pluginPackageContentTreeDigest(
resourceEntries.map((entry) => ({
path: entry.path,
bytes: entry.body.byteLength,
digest: sha256(entry.body),
})),
);
const artifact = canonicalTar([
{
path: 'package.json',
body: Buffer.from(serializePluginPackageManifest(manifest), 'utf8'),
},
...resourceEntries,
]);
const packageConfig = jsonBytes({
schema: 'qinglong/plugin-package-oci-config@v1',
@@ -201,6 +228,142 @@ function createFixture({ registry, architecture, createdAtMs = Date.now() }) {
};
const packageManifestBytes = jsonBytes(packageManifestValue);
const packageManifestDigest = sha256(packageManifestBytes);
return Object.freeze({
manifest,
source: Object.freeze({
kind: 'oci',
locator: `oci://${registry}/${REPOSITORY}@sha256:${packageManifestDigest}`,
artifactDigest,
artifactBytes: artifact.byteLength,
contentDigest,
}),
packageConfig,
packageConfigDigest,
packageManifestBytes,
packageManifestDigest,
artifact,
artifactDigest,
mediaTypes: Object.freeze({
packageConfig: PLUGIN_PACKAGE_OCI_CONFIG_MEDIA_TYPE,
signatureArtifact: PLUGIN_PACKAGE_OCI_SIGNATURE_ARTIFACT_TYPE,
signatureConfig: PLUGIN_PACKAGE_OCI_SIGNATURE_CONFIG_MEDIA_TYPE,
}),
});
}
function signedPackageRoutes(material, lock, privateKey) {
const {
PLUGIN_PACKAGE_SIGNATURE_SCHEMA,
pluginPackagePublisherSignaturePayload,
} = ql3Require('@qinglong/runtime-core/plugin-package-bundle');
const signature = {
schema: PLUGIN_PACKAGE_SIGNATURE_SCHEMA,
publisher: PUBLISHER,
keyId: KEY_ID,
signature: sign(
null,
pluginPackagePublisherSignaturePayload(lock, PUBLISHER, KEY_ID),
privateKey,
).toString('base64url'),
};
const signatureConfig = jsonBytes(signature);
const signatureConfigDigest = sha256(signatureConfig);
const signatureManifestValue = {
schemaVersion: 2,
mediaType: OCI_MANIFEST,
artifactType: material.mediaTypes.signatureArtifact,
config: {
mediaType: material.mediaTypes.signatureConfig,
digest: `sha256:${signatureConfigDigest}`,
size: signatureConfig.byteLength,
},
layers: [],
subject: {
mediaType: OCI_MANIFEST,
digest: `sha256:${material.packageManifestDigest}`,
size: material.packageManifestBytes.byteLength,
},
};
const signatureManifestBytes = jsonBytes(signatureManifestValue);
const signatureManifestDigest = sha256(signatureManifestBytes);
const referrers = jsonBytes({
schemaVersion: 2,
mediaType: OCI_INDEX,
manifests: [
{
mediaType: OCI_MANIFEST,
digest: `sha256:${signatureManifestDigest}`,
size: signatureManifestBytes.byteLength,
artifactType: material.mediaTypes.signatureArtifact,
annotations: {
'qinglong.io/plugin-package-lock-digest': lock.lockDigest,
},
},
],
});
const prefix = `/v2/${REPOSITORY}`;
return Object.freeze([
route(
`${prefix}/manifests/sha256:${material.packageManifestDigest}`,
OCI_MANIFEST,
material.packageManifestBytes,
),
route(
`${prefix}/blobs/sha256:${material.packageConfigDigest}`,
material.mediaTypes.packageConfig,
material.packageConfig,
),
route(
`${prefix}/referrers/sha256:${
material.packageManifestDigest
}?artifactType=${encodeURIComponent(
material.mediaTypes.signatureArtifact,
)}`,
OCI_INDEX,
referrers,
),
route(
`${prefix}/manifests/sha256:${signatureManifestDigest}`,
OCI_MANIFEST,
signatureManifestBytes,
),
route(
`${prefix}/blobs/sha256:${signatureConfigDigest}`,
material.mediaTypes.signatureConfig,
signatureConfig,
),
route(
`${prefix}/blobs/sha256:${material.artifactDigest}`,
BUNDLE,
material.artifact,
),
]);
}
function createFixture({ registry, architecture, createdAtMs = Date.now() }) {
if (
typeof registry !== 'string' ||
!/^[a-z0-9](?:[-a-z0-9.]{0,251}[a-z0-9])?$/.test(registry) ||
!['amd64', 'arm64'].includes(architecture) ||
!Number.isSafeInteger(createdAtMs) ||
createdAtMs < 1
) {
throw new TypeError('Plugin Package E2E fixture options are invalid');
}
const { PluginPackagePublisherTrustRegistry } = ql3Require(
'@qinglong/runtime-core/plugin-package-bundle',
);
const { planPluginPackageInstall } = ql3Require(
'@qinglong/runtime-core/plugin-package',
);
const {
createPluginPackageLock,
pluginPackageInstallActionDigest,
pluginPackageInstallPlanDigest,
} = ql3Require('@qinglong/runtime-core/plugin-package-install');
const { createPluginPackageResourceGenerationFromReferences } = ql3Require(
'@qinglong/runtime-core/plugin-package-resource-generation',
);
const environment = {
qinglongVersion: '3.0.0-alpha.0',
architecture,
@@ -209,33 +372,28 @@ function createFixture({ registry, architecture, createdAtMs = Date.now() }) {
availableMemoryBytes: 128 * 1024 * 1024,
availableDiskBytes: 256 * 1024 * 1024,
};
const plan = planPluginPackageInstall(manifest, environment);
const source = {
kind: 'oci',
locator: `oci://${registry}/${REPOSITORY}@sha256:${packageManifestDigest}`,
artifactDigest,
artifactBytes: artifact.byteLength,
contentDigest: pluginPackageContentTreeDigest([]),
};
const action = {
lockId: 'lock-plugin-recovery-e2e',
const initialManifest = pluginManifest(architecture, '1.0.0');
const initialMaterial = packageMaterial(registry, initialManifest, {});
const initialPlan = planPluginPackageInstall(initialManifest, environment);
const initialAction = {
lockId: 'lock-plugin-recovery-e2e-initial',
projectId: 'default',
manifest,
plan,
manifest: initialManifest,
plan: initialPlan,
environment,
source,
source: initialMaterial.source,
architecture,
deploymentProfile: 'cluster-control',
targetGeneration: 1,
};
const lock = createPluginPackageLock({
...action,
const initialLock = createPluginPackageLock({
...initialAction,
approval: {
requestId: 'approval-plugin-recovery-e2e',
requestId: 'approval-plugin-recovery-e2e-initial',
requestVersion: 1,
dispatchId: 'dispatch-plugin-recovery-e2e',
actionDigest: pluginPackageInstallActionDigest(action),
previewDigest: pluginPackageInstallPlanDigest(plan),
dispatchId: 'dispatch-plugin-recovery-e2e-initial',
actionDigest: pluginPackageInstallActionDigest(initialAction),
previewDigest: pluginPackageInstallPlanDigest(initialPlan),
approvedBy: { type: 'user', id: 'e2e-owner' },
approvedAtMs: createdAtMs - 1,
expiresAtMs: createdAtMs + 60 * 60 * 1000,
@@ -243,6 +401,46 @@ function createFixture({ registry, architecture, createdAtMs = Date.now() }) {
},
createdAtMs,
});
const upgradeCreatedAtMs = createdAtMs + 10;
const upgradeManifest = pluginManifest(architecture, '2.0.0', true);
const upgradeMaterial = packageMaterial(
registry,
upgradeManifest,
invalidUpgradeResources(),
);
const upgradePlan = planPluginPackageInstall(
upgradeManifest,
environment,
initialManifest,
);
const upgradeAction = {
lockId: 'lock-plugin-recovery-e2e-upgrade',
projectId: 'default',
manifest: upgradeManifest,
plan: upgradePlan,
environment,
previousManifest: initialManifest,
source: upgradeMaterial.source,
architecture,
deploymentProfile: 'cluster-control',
targetGeneration: 2,
previousLockDigest: initialLock.lockDigest,
};
const upgradeLock = createPluginPackageLock({
...upgradeAction,
approval: {
requestId: 'approval-plugin-recovery-e2e-upgrade',
requestVersion: 1,
dispatchId: 'dispatch-plugin-recovery-e2e-upgrade',
actionDigest: pluginPackageInstallActionDigest(upgradeAction),
previewDigest: pluginPackageInstallPlanDigest(upgradePlan),
approvedBy: { type: 'user', id: 'e2e-owner' },
approvedAtMs: upgradeCreatedAtMs - 1,
expiresAtMs: upgradeCreatedAtMs + 60 * 60 * 1000,
fence: { projectVersion: 1, bindingVersion: 1 },
},
createdAtMs: upgradeCreatedAtMs,
});
const { publicKey, privateKey } = generateKeyPairSync('ed25519');
const publicKeyPem = publicKey.export({ format: 'pem', type: 'spki' });
const trust = {
@@ -258,94 +456,58 @@ function createFixture({ registry, architecture, createdAtMs = Date.now() }) {
],
};
new PluginPackagePublisherTrustRegistry(trust.keys);
const signature = {
schema: PLUGIN_PACKAGE_SIGNATURE_SCHEMA,
publisher: PUBLISHER,
keyId: KEY_ID,
signature: sign(
null,
pluginPackagePublisherSignaturePayload(lock, PUBLISHER, KEY_ID),
privateKey,
).toString('base64url'),
};
const signatureConfig = jsonBytes(signature);
const signatureConfigDigest = sha256(signatureConfig);
const signatureManifestValue = {
schemaVersion: 2,
mediaType: OCI_MANIFEST,
artifactType: PLUGIN_PACKAGE_OCI_SIGNATURE_ARTIFACT_TYPE,
config: {
mediaType: PLUGIN_PACKAGE_OCI_SIGNATURE_CONFIG_MEDIA_TYPE,
digest: `sha256:${signatureConfigDigest}`,
size: signatureConfig.byteLength,
},
layers: [],
subject: {
mediaType: OCI_MANIFEST,
digest: `sha256:${packageManifestDigest}`,
size: packageManifestBytes.byteLength,
},
};
const signatureManifestBytes = jsonBytes(signatureManifestValue);
const signatureManifestDigest = sha256(signatureManifestBytes);
const referrers = jsonBytes({
schemaVersion: 2,
mediaType: OCI_INDEX,
manifests: [
{
mediaType: OCI_MANIFEST,
digest: `sha256:${signatureManifestDigest}`,
size: signatureManifestBytes.byteLength,
artifactType: PLUGIN_PACKAGE_OCI_SIGNATURE_ARTIFACT_TYPE,
annotations: {
'qinglong.io/plugin-package-lock-digest': lock.lockDigest,
},
},
],
const initialRoutes = signedPackageRoutes(
initialMaterial,
initialLock,
privateKey,
);
const upgradeRoutes = signedPackageRoutes(
upgradeMaterial,
upgradeLock,
privateKey,
);
const initial = Object.freeze({
installationId: 'install-plugin-recovery-e2e-initial',
manifest: initialManifest,
lock: initialLock,
generation: createPluginPackageResourceGenerationFromReferences({
installationId: 'install-plugin-recovery-e2e-initial',
projectId: initialLock.projectId,
packageName: initialLock.packageName,
lockDigest: initialLock.lockDigest,
generation: initialLock.targetGeneration,
previousActiveLockDigest: null,
contentDigest: initialLock.source.contentDigest,
resources: initialLock.resources,
}),
routes: initialRoutes,
});
const upgrade = Object.freeze({
installationId: 'install-plugin-recovery-e2e-upgrade',
manifest: upgradeManifest,
lock: upgradeLock,
generation: createPluginPackageResourceGenerationFromReferences({
installationId: 'install-plugin-recovery-e2e-upgrade',
projectId: upgradeLock.projectId,
packageName: upgradeLock.packageName,
lockDigest: upgradeLock.lockDigest,
generation: upgradeLock.targetGeneration,
previousActiveLockDigest: initialLock.lockDigest,
contentDigest: upgradeLock.source.contentDigest,
resources: upgradeLock.resources,
}),
routes: upgradeRoutes,
});
const prefix = `/v2/${REPOSITORY}`;
const routes = [
route(
`${prefix}/manifests/sha256:${packageManifestDigest}`,
OCI_MANIFEST,
packageManifestBytes,
),
route(
`${prefix}/blobs/sha256:${packageConfigDigest}`,
PLUGIN_PACKAGE_OCI_CONFIG_MEDIA_TYPE,
packageConfig,
),
route(
`${prefix}/referrers/sha256:${packageManifestDigest}?artifactType=${encodeURIComponent(
PLUGIN_PACKAGE_OCI_SIGNATURE_ARTIFACT_TYPE,
)}`,
OCI_INDEX,
referrers,
),
route(
`${prefix}/manifests/sha256:${signatureManifestDigest}`,
OCI_MANIFEST,
signatureManifestBytes,
),
route(
`${prefix}/blobs/sha256:${signatureConfigDigest}`,
PLUGIN_PACKAGE_OCI_SIGNATURE_CONFIG_MEDIA_TYPE,
signatureConfig,
),
route(
`${prefix}/blobs/sha256:${artifactDigest}`,
BUNDLE,
artifact,
),
];
return Object.freeze({
schema: FIXTURE_SCHEMA,
registry,
repository: REPOSITORY,
architecture,
lock,
lock: initialLock,
initial,
upgrade,
trust,
routes,
routes: Object.freeze([...initialRoutes, ...upgradeRoutes]),
});
}
@@ -355,7 +517,10 @@ function readFixture(filePath) {
!value ||
value.schema !== FIXTURE_SCHEMA ||
!Array.isArray(value.routes) ||
!value.lock ||
!value.initial?.lock ||
!value.initial?.generation ||
!value.upgrade?.lock ||
!value.upgrade?.generation ||
!value.trust
) {
throw new Error('Plugin Package E2E fixture is invalid');
@@ -459,17 +624,27 @@ async function runSeed() {
assertPostgresPackageExecutorSchemaReady,
createPostgresDatabaseOpener,
loadPostgresConnectionEnvironment,
PostgresPluginPackageSecretBindingTransitionRepository,
} = ql3Require('@qinglong/cluster-postgres/package-executor');
const {
PostgresPluginPackageInstallRepository,
} = ql3Require('@qinglong/cluster-postgres/plugin-package-install');
const { PostgresPluginPackageInstallRepository } = ql3Require(
'@qinglong/cluster-postgres/plugin-package-install',
);
const {
createPluginPackageInstall,
normalizePluginPackageLock,
pluginPackageInstallCreate,
} = ql3Require('@qinglong/runtime-core/plugin-package-install');
const { createPluginPackageSecretBindingTarget } = ql3Require(
'@qinglong/runtime-core/plugin-package-secret-binding',
);
const { createPluginPackageSecretBindingTransitionPlan } = ql3Require(
'@qinglong/runtime-core/plugin-package-secret-binding-transition-plan',
);
const fixture = readFixture(process.env.QL3_E2E_FIXTURE_FILE);
const lock = normalizePluginPackageLock(fixture.lock);
const mode = process.env.QL3_E2E_MODE;
if (!['seed-initial', 'seed-upgrade', 'commit-transition'].includes(mode)) {
throw new Error('Plugin Package E2E seed mode is invalid');
}
const connection = loadPostgresConnectionEnvironment(process.env, {
host: 'QL3_E2E_POSTGRES_HOST',
port: 'QL3_E2E_POSTGRES_PORT',
@@ -489,21 +664,72 @@ async function runSeed() {
})();
try {
await assertPostgresPackageExecutorSchemaReady(database.pool);
if (mode === 'commit-transition') {
const plannedAtMs = Date.now();
const transitionPlan = createPluginPackageSecretBindingTransitionPlan({
previousTarget: createPluginPackageSecretBindingTarget(
fixture.initial.generation,
fixture.initial.manifest,
),
previousBinding: null,
previousAttemptGeneration: 1,
nextGeneration: fixture.upgrade.generation,
nextManifest: fixture.upgrade.manifest,
assignments: [],
plannedAtMs,
});
const result =
await new PostgresPluginPackageSecretBindingTransitionRepository(
database.pool,
).apply({
transitionPlan,
evidenceDigest: transitionPlan.transitionDigest,
committedAtMs: plannedAtMs + 1,
});
process.stdout.write(
`${JSON.stringify({
schema: 'qinglong/plugin-package-recovery-e2e-transition-result@v1',
event: 'transition_completed',
status: result.status,
generationDigest: transitionPlan.nextTarget.generationDigest,
transitionDigest: transitionPlan.transitionDigest,
bindingDigest: result.receipt.bindingDigest,
receiptDigest: result.receipt.receiptDigest,
})}\n`,
);
return;
}
const selected =
mode === 'seed-initial' ? fixture.initial : fixture.upgrade;
const lock = normalizePluginPackageLock(selected.lock);
const repository = new PostgresPluginPackageInstallRepository(
database.pool,
);
const previous = await repository.find(lock.projectId, lock.packageName);
if (
(mode === 'seed-initial' && previous !== null) ||
(mode === 'seed-upgrade' &&
(previous?.state !== 'active' ||
previous.installationId !== fixture.initial.installationId ||
previous.lockDigest !== fixture.initial.lock.lockDigest))
) {
throw new Error('Plugin Package E2E previous install head is invalid');
}
const record = createPluginPackageInstall(lock, {
installationId: 'install-plugin-recovery-e2e',
mutationId: 'mutation-plugin-recovery-e2e-create',
installationId: selected.installationId,
mutationId: `mutation-plugin-recovery-e2e-${
mode === 'seed-initial' ? 'initial' : 'upgrade'
}-create`,
occurredAtMs: lock.createdAtMs + 1,
});
const result = await repository.create(
pluginPackageInstallCreate(lock, record, null),
pluginPackageInstallCreate(lock, record, previous),
);
process.stdout.write(
`${JSON.stringify({
schema: 'qinglong/plugin-package-recovery-e2e-seed-result@v1',
event: 'seed_completed',
phase: mode === 'seed-initial' ? 'initial' : 'upgrade',
status: result.status,
state: result.record.state,
installationId: result.record.installationId,
@@ -521,11 +747,17 @@ async function main() {
await runRegistry();
return;
}
if (process.env.QL3_E2E_MODE === 'seed') {
if (
['seed-initial', 'seed-upgrade', 'commit-transition'].includes(
process.env.QL3_E2E_MODE,
)
) {
await runSeed();
return;
}
throw new Error('QL3_E2E_MODE must be registry or seed');
throw new Error(
'QL3_E2E_MODE must be registry, seed-initial, seed-upgrade or commit-transition',
);
}
if (require.main === module) {
@@ -35,7 +35,13 @@ const POSTGRES_REPOSITORY_DIGEST = `postgres@${POSTGRES_INDEX_DIGEST}`;
const DEFAULT_ADMIN_IMAGE = 'qinglong3-cluster-admin:ql3-plugin-recovery-e2e';
const DEFAULT_CONTROL_IMAGE =
'qinglong3-cluster-control:ql3-plugin-recovery-e2e';
const REPORT_SCHEMA = 'qinglong/plugin-package-recovery-e2e-live-contract@v1';
const REPORT_SCHEMA = 'qinglong/plugin-package-recovery-e2e-live-contract@v2';
const INITIAL_SEED_JOB = 'ql3-plugin-package-e2e-seed-initial';
const INITIAL_RECOVERY_JOB = 'ql3-plugin-package-recovery-initial';
const UPGRADE_SEED_JOB = 'ql3-plugin-package-e2e-seed-upgrade';
const UPGRADE_STAGE_JOB = 'ql3-plugin-package-recovery-stage-upgrade';
const TRANSITION_JOB = 'ql3-plugin-package-e2e-transition';
const UPGRADE_REJECTION_JOB = 'ql3-plugin-package-recovery-reject-upgrade';
const SAFE_CLUSTER =
/^ql3-plugin-recovery-e2e(?:-[a-z0-9](?:[-a-z0-9]{0,24}[a-z0-9])?)?$/;
@@ -727,11 +733,11 @@ function migrationJob() {
return job;
}
function seedJob() {
function seedJob(name, mode) {
return {
apiVersion: 'batch/v1',
kind: 'Job',
metadata: { name: 'ql3-plugin-package-e2e-seed', namespace: NAMESPACE },
metadata: { name, namespace: NAMESPACE },
spec: {
backoffLimit: 0,
activeDeadlineSeconds: 300,
@@ -754,7 +760,7 @@ function seedJob() {
imagePullPolicy: 'Never',
command: ['node', '/opt/ql3-e2e/fixture.cjs'],
env: [
{ name: 'QL3_E2E_MODE', value: 'seed' },
{ name: 'QL3_E2E_MODE', value: mode },
{
name: 'QL3_E2E_FIXTURE_FILE',
value: '/opt/ql3-e2e/fixture.json',
@@ -811,7 +817,7 @@ function seedJob() {
};
}
function recoveryResources(fixture) {
function recoveryResources(fixture, jobName) {
const rbac = readYamlDocuments(
path.join(
ROOT,
@@ -830,6 +836,7 @@ function recoveryResources(fixture) {
'deploy/kubernetes/ql3-cluster/operations/plugin-package-recovery/base/recover-job.yaml',
),
);
job.metadata.name = jobName;
job.metadata.namespace = NAMESPACE;
job.spec.activeDeadlineSeconds = 300;
const container = job.spec.template.spec.containers[0];
@@ -956,7 +963,14 @@ function recoveryResources(fixture) {
];
}
function waitForJob(name, timeoutMs = 5 * 60 * 1000) {
function waitForJob(
name,
expectedStatus = 'complete',
timeoutMs = 5 * 60 * 1000,
) {
if (!['complete', 'failed'].includes(expectedStatus)) {
fail('expected Job status is invalid');
}
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
const job = kubectlJson(['-n', NAMESPACE, 'get', 'job', name]);
@@ -964,11 +978,17 @@ function waitForJob(name, timeoutMs = 5 * 60 * 1000) {
(condition) =>
condition.type === 'Complete' && condition.status === 'True',
);
if (complete) return job;
if (complete) {
if (expectedStatus !== 'complete') {
fail(`${name} completed but failure was required`);
}
return job;
}
const failed = job.status?.conditions?.find(
(condition) => condition.type === 'Failed' && condition.status === 'True',
);
if (failed) {
if (expectedStatus === 'failed') return job;
const logs = kubectl(
['-n', NAMESPACE, 'logs', `job/${name}`, '--all-containers=true'],
{ capture: true, quiet: true, allowFailure: true },
@@ -1039,7 +1059,72 @@ function canI(verb, resource) {
return result.stdout === 'yes';
}
function databaseEvidence(lockDigest) {
function upgradeStageEvidence(fixture, transitionReceiptCount) {
const generationDigest = fixture.upgrade.generation.generationDigest;
assert.match(generationDigest, /^[0-9a-f]{64}$/);
const sql = `
SELECT json_build_object(
'state', (
SELECT state FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'previousActiveLockDigest', (
SELECT previous_active_lock_digest FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'activeLockDigest', (
SELECT active_lock_digest FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'mutationCount', (
SELECT count(*) FROM ql3.plugin_package_install_mutations
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'transitionReceiptCount', (
SELECT count(*)
FROM ql3.plugin_package_secret_binding_transition_receipts
WHERE generation_digest = '${generationDigest}'
),
'candidateRevisionCount', (
SELECT count(*) FROM ql3.plugin_package_materialized_revisions
WHERE generation_digest = '${generationDigest}'
)
)::text;
`.trim();
const output = kubectl(
[
'-n',
NAMESPACE,
'exec',
POSTGRES_NAME,
'--',
'psql',
'--username',
'postgres',
'--dbname',
'qinglong',
'--tuples-only',
'--no-align',
'--command',
sql,
],
{ capture: true, quiet: true },
).stdout;
const value = JSON.parse(output);
assert.equal(value.state, 'staged');
assert.equal(value.previousActiveLockDigest, fixture.initial.lock.lockDigest);
assert.equal(value.activeLockDigest, fixture.initial.lock.lockDigest);
assert.equal(value.mutationCount, 2);
assert.equal(value.transitionReceiptCount, transitionReceiptCount);
assert.equal(value.candidateRevisionCount, 0);
return value;
}
function databaseEvidence(fixture) {
const initialGenerationDigest = fixture.initial.generation.generationDigest;
const upgradeGenerationDigest = fixture.upgrade.generation.generationDigest;
assert.match(initialGenerationDigest, /^[0-9a-f]{64}$/);
assert.match(upgradeGenerationDigest, /^[0-9a-f]{64}$/);
const sql = `
SELECT json_build_object(
'migrationCount', (SELECT count(*) FROM ql3.schema_migrations),
@@ -1048,25 +1133,63 @@ SELECT json_build_object(
FROM ql3.schema_capabilities
WHERE contract_name = 'control-core'
),
'state', (
'initialState', (
SELECT state
FROM ql3.plugin_package_installs
WHERE installation_id = 'install-plugin-recovery-e2e'
WHERE installation_id = '${fixture.initial.installationId}'
),
'lockDigest', (
SELECT lock_digest
FROM ql3.plugin_package_installs
WHERE installation_id = 'install-plugin-recovery-e2e'
),
'activeLockDigest', (
'initialActiveLockDigest', (
SELECT active_lock_digest
FROM ql3.plugin_package_installs
WHERE installation_id = 'install-plugin-recovery-e2e'
WHERE installation_id = '${fixture.initial.installationId}'
),
'mutationCount', (
'upgradeState', (
SELECT state
FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'upgradePreviousActiveLockDigest', (
SELECT previous_active_lock_digest
FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'upgradeActiveLockDigest', (
SELECT active_lock_digest
FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'upgradeFailureReason', (
SELECT record_json #>> '{failure,reason}'
FROM ql3.plugin_package_installs
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'initialMutationCount', (
SELECT count(*)
FROM ql3.plugin_package_install_mutations
WHERE installation_id = 'install-plugin-recovery-e2e'
WHERE installation_id = '${fixture.initial.installationId}'
),
'upgradeMutationCount', (
SELECT count(*)
FROM ql3.plugin_package_install_mutations
WHERE installation_id = '${fixture.upgrade.installationId}'
),
'headInstallationId', (
SELECT installation_id
FROM ql3.plugin_package_install_heads
WHERE project_id = 'default' AND package_name = 'e2e-monitor'
),
'transitionReceiptCount', (
SELECT count(*)
FROM ql3.plugin_package_secret_binding_transition_receipts
WHERE generation_digest = '${upgradeGenerationDigest}'
),
'initialRevisionCount', (
SELECT count(*) FROM ql3.plugin_package_materialized_revisions
WHERE generation_digest = '${initialGenerationDigest}'
),
'upgradeRevisionCount', (
SELECT count(*) FROM ql3.plugin_package_materialized_revisions
WHERE generation_digest = '${upgradeGenerationDigest}'
),
'recoverableCount', (
SELECT count(*)
@@ -1095,17 +1218,28 @@ SELECT json_build_object(
{ capture: true, quiet: true },
).stdout;
const value = JSON.parse(output);
assert.equal(value.migrationCount, 22);
assert.equal(value.capabilityVersion, 21);
assert.equal(value.state, 'active');
assert.equal(value.lockDigest, lockDigest);
assert.equal(value.activeLockDigest, lockDigest);
assert.equal(value.mutationCount, 4);
assert.equal(value.migrationCount, 65);
assert.equal(value.capabilityVersion, 64);
assert.equal(value.initialState, 'active');
assert.equal(value.initialActiveLockDigest, fixture.initial.lock.lockDigest);
assert.equal(value.upgradeState, 'failed');
assert.equal(
value.upgradePreviousActiveLockDigest,
fixture.initial.lock.lockDigest,
);
assert.equal(value.upgradeActiveLockDigest, fixture.initial.lock.lockDigest);
assert.equal(value.upgradeFailureReason, 'activation_fact_conflict');
assert.equal(value.initialMutationCount, 4);
assert.equal(value.upgradeMutationCount, 3);
assert.equal(value.headInstallationId, fixture.upgrade.installationId);
assert.equal(value.transitionReceiptCount, 1);
assert.equal(value.initialRevisionCount, 1);
assert.equal(value.upgradeRevisionCount, 0);
assert.equal(value.recoverableCount, 0);
return value;
}
function activePointerEvidence(lockDigest) {
function activePointerEvidence(fixture) {
const values = kubectlJson([
'-n',
NAMESPACE,
@@ -1117,13 +1251,14 @@ function activePointerEvidence(lockDigest) {
assert.equal(values.items.length, 1);
const configMap = values.items[0];
const pointer = JSON.parse(configMap.data['active.json']);
assert.equal(pointer.intent.installationId, 'install-plugin-recovery-e2e');
assert.equal(pointer.intent.lockDigest, lockDigest);
assert.equal(pointer.intent.installationId, fixture.initial.installationId);
assert.equal(pointer.intent.lockDigest, fixture.initial.lock.lockDigest);
assert.equal(pointer.receipt.generation, 1);
return Object.freeze({
name: configMap.metadata.name,
uid: configMap.metadata.uid,
resourceVersion: configMap.metadata.resourceVersion,
activeJson: configMap.data['active.json'],
intentDigest: pointer.intent.intentDigest,
activationRef: pointer.receipt.activationRef,
});
@@ -1141,10 +1276,17 @@ function registryEvidence(fixture) {
.map((line) => JSON.parse(line))
.filter((value) => value.schema === REGISTRY_EVENT_SCHEMA);
const packageRequests = events.filter((event) => event.path !== '/v2/');
assert.equal(packageRequests.length, fixture.routes.length);
const expectedPaths = [
...fixture.initial.routes.map((routeValue) => routeValue.path),
...fixture.upgrade.routes.flatMap((routeValue) => [
routeValue.path,
routeValue.path,
]),
].sort();
assert.equal(packageRequests.length, expectedPaths.length);
assert.deepEqual(
packageRequests.map((event) => event.path).sort(),
fixture.routes.map((routeValue) => routeValue.path).sort(),
expectedPaths,
);
assert.ok(packageRequests.every((event) => event.status === 200));
assert.ok(packageRequests.every((event) => event.authenticated === true));
@@ -1154,6 +1296,8 @@ function registryEvidence(fixture) {
authenticatedRequestCount: packageRequests.length,
requestCount: packageRequests.length,
uniquePaths: new Set(packageRequests.map((event) => event.path)).size,
initialRequestCount: fixture.initial.routes.length,
upgradeRequestCount: fixture.upgrade.routes.length * 2,
redirects: 0,
});
}
@@ -1572,36 +1716,109 @@ async function main() {
(value) => value.event === 'migration_completed',
);
apply(seedJob(), 'persist one durable queued Plugin Package installation');
waitForJob('ql3-plugin-package-e2e-seed');
const seed = lastJsonLine(
jobLog('ql3-plugin-package-e2e-seed'),
apply(
seedJob(INITIAL_SEED_JOB, 'seed-initial'),
'persist initial durable queued Plugin Package installation',
);
waitForJob(INITIAL_SEED_JOB);
const initialSeed = lastJsonLine(
jobLog(INITIAL_SEED_JOB),
(value) => value.event === 'seed_completed',
);
assert.equal(seed.status, 'created');
assert.equal(seed.state, 'queued');
assert.equal(seed.lockDigest, fixture.lock.lockDigest);
assert.equal(initialSeed.phase, 'initial');
assert.equal(initialSeed.status, 'created');
assert.equal(initialSeed.state, 'queued');
assert.equal(initialSeed.lockDigest, fixture.initial.lock.lockDigest);
for (const resource of recoveryResources(fixture)) {
apply(resource, `apply recovery ${resource.kind}`);
for (const resource of recoveryResources(fixture, INITIAL_RECOVERY_JOB)) {
apply(resource, `apply initial recovery ${resource.kind}`);
}
const recovered = waitForJob('ql3-plugin-package-recovery');
const recoveryLog = jobLog('ql3-plugin-package-recovery');
const completed = lastJsonLine(
recoveryLog,
const initialRecovered = waitForJob(INITIAL_RECOVERY_JOB);
const initialCompleted = lastJsonLine(
jobLog(INITIAL_RECOVERY_JOB),
(value) => value.event === 'recovery_completed',
);
assert.equal(completed.recovery.safeToAdmit, true);
assert.equal(completed.recovery.remaining, false);
assert.equal(completed.recovery.manualRequired, 0);
assert.equal(initialCompleted.recovery.safeToAdmit, true);
assert.equal(initialCompleted.recovery.remaining, false);
assert.equal(initialCompleted.recovery.manualRequired, 0);
const pointerBeforeUpgrade = activePointerEvidence(fixture);
const database = databaseEvidence(fixture.lock.lockDigest);
const pointer = activePointerEvidence(fixture.lock.lockDigest);
apply(
seedJob(UPGRADE_SEED_JOB, 'seed-upgrade'),
'persist invalid durable queued Plugin Package upgrade',
);
waitForJob(UPGRADE_SEED_JOB);
const upgradeSeed = lastJsonLine(
jobLog(UPGRADE_SEED_JOB),
(value) => value.event === 'seed_completed',
);
assert.equal(upgradeSeed.phase, 'upgrade');
assert.equal(upgradeSeed.status, 'created');
assert.equal(upgradeSeed.state, 'queued');
assert.equal(upgradeSeed.lockDigest, fixture.upgrade.lock.lockDigest);
for (const resource of recoveryResources(fixture, UPGRADE_STAGE_JOB)) {
apply(resource, `apply upgrade staging recovery ${resource.kind}`);
}
const stagedUpgrade = waitForJob(UPGRADE_STAGE_JOB, 'failed');
const stagedFailureCondition = stagedUpgrade.status.conditions.find(
(condition) => condition.type === 'Failed' && condition.status === 'True',
);
assert.ok(stagedFailureCondition?.lastTransitionTime);
const stageFailure = lastJsonLine(
jobLog(UPGRADE_STAGE_JOB),
(value) => value.event === 'recovery_failed',
);
assert.equal(
stageFailure.name,
'ClusterPluginPackageRecoveryRequiredError',
);
const stagedDatabase = upgradeStageEvidence(fixture, 0);
assert.deepEqual(activePointerEvidence(fixture), pointerBeforeUpgrade);
apply(
seedJob(TRANSITION_JOB, 'commit-transition'),
'commit durable no-secret binding transition receipt',
);
const committedTransition = waitForJob(TRANSITION_JOB);
const transition = lastJsonLine(
jobLog(TRANSITION_JOB),
(value) => value.event === 'transition_completed',
);
assert.equal(transition.status, 'created');
assert.equal(
transition.generationDigest,
fixture.upgrade.generation.generationDigest,
);
assert.equal(transition.bindingDigest, null);
upgradeStageEvidence(fixture, 1);
for (const resource of recoveryResources(fixture, UPGRADE_REJECTION_JOB)) {
apply(resource, `apply upgrade rejection recovery ${resource.kind}`);
}
const rejectedUpgrade = waitForJob(UPGRADE_REJECTION_JOB);
const rejectionCompleted = lastJsonLine(
jobLog(UPGRADE_REJECTION_JOB),
(value) => value.event === 'recovery_completed',
);
assert.equal(rejectionCompleted.recovery.safeToAdmit, true);
assert.equal(rejectionCompleted.recovery.remaining, false);
assert.equal(rejectionCompleted.recovery.manualRequired, 0);
const database = databaseEvidence(fixture);
const pointerAfterRejection = activePointerEvidence(fixture);
assert.deepEqual(pointerAfterRejection, pointerBeforeUpgrade);
const oci = registryEvidence(fixture);
const rbac = recoveryRbacEvidence();
assertRuntimeCannotReadPluginAuthority();
const runtime = applyRuntimeAfterRecovery(recovered, migrated, secrets);
const recoveryImageId = jobImageId('ql3-plugin-package-recovery');
const runtime = applyRuntimeAfterRecovery(
rejectedUpgrade,
migrated,
secrets,
);
const initialRecoveryImageId = jobImageId(INITIAL_RECOVERY_JOB);
const stageRecoveryImageId = jobImageId(UPGRADE_STAGE_JOB);
const rejectionRecoveryImageId = jobImageId(UPGRADE_REJECTION_JOB);
const migrationImageId = jobImageId('ql3-cluster-migration');
const postgresPod = kubectlJson([
'-n',
@@ -1621,21 +1838,40 @@ async function main() {
controlBuildId: imageId(CONTROL_IMAGE),
postgresRepositoryDigest: POSTGRES_REPOSITORY_DIGEST,
migrationImageId,
recoveryImageId,
initialRecoveryImageId,
stageRecoveryImageId,
rejectionRecoveryImageId,
postgresImageId: postgresPod.status.containerStatuses[0].imageID,
}),
ordering: Object.freeze({
migrationJobUid: migrated.metadata.uid,
migrationCompletedAt: migrated.status.completionTime,
recoveryJobUid: recovered.metadata.uid,
recoveryCompletedAt: recovered.status.completionTime,
initialRecoveryJobUid: initialRecovered.metadata.uid,
initialRecoveryCompletedAt: initialRecovered.status.completionTime,
upgradeStageJobUid: stagedUpgrade.metadata.uid,
upgradeStageFailedAt: stagedFailureCondition.lastTransitionTime,
transitionJobUid: committedTransition.metadata.uid,
transitionCompletedAt: committedTransition.status.completionTime,
rejectionRecoveryJobUid: rejectedUpgrade.metadata.uid,
rejectionRecoveryCompletedAt: rejectedUpgrade.status.completionTime,
runtimeCreatedAt: runtime.creationTimestamp,
runtimeBoundRecoveryJobUid: runtime.recoveryJobUid,
}),
failedUpgrade: Object.freeze({
stageFailure: Object.freeze({
jobUid: stagedUpgrade.metadata.uid,
reason: stageFailure.name,
durableState: stagedDatabase.state,
}),
transitionReceiptDigest: transition.receiptDigest,
rejectionReason: database.upgradeFailureReason,
candidateRevisionCount: database.upgradeRevisionCount,
activePointerUnchanged: true,
}),
database,
oci,
kubernetes: Object.freeze({
activePointer: pointer,
activePointer: pointerAfterRejection,
rbac,
}),
runtime,