Merge pull request #423 from Vaso73/fix/oci-stack-recreate

fix(oci): recreate coordinated stacks from saved state
This commit is contained in:
MacRimi
2026-10-06 19:12:28 +02:00
committed by GitHub
19 changed files with 396 additions and 59 deletions
+3 -1
View File
@@ -166,6 +166,7 @@ def resolve_candidate(reference, architecture):
if '@' in reference:
transport = repo + '@' + reference.split('@', 1)[1]
raw = command('skopeo', 'inspect', '--raw', 'docker://' + transport)
registry_digest = 'sha256:' + sha(raw)
manifest = json.loads(raw)
if 'manifests' in manifest:
matches = [m for m in manifest['manifests'] if m.get('platform', {}).get('architecture') == architecture
@@ -192,7 +193,8 @@ def resolve_candidate(reference, architecture):
# The build date identifies the image as the publisher released it: it is
# what changes when an image is rebuilt, whether or not the application
# version inside it moved.
return {'manifest_digest': digest, 'layers': [layer['digest'] for layer in json.loads(raw).get('layers', [])],
return {'manifest_digest': digest, 'registry_digest': registry_digest,
'layers': [layer['digest'] for layer in json.loads(raw).get('layers', [])],
'defaults': defaults, 'version': version, 'created': config.get('created')}
+38 -20
View File
@@ -32,7 +32,7 @@ def validate_database_transition(previous, candidate):
raise ValueError(translate('The new image changes the PostgreSQL major version; the data must be migrated before updating'))
def nextcloud_plan(primary, records, inventory, lifecycle):
def nextcloud_plan(primary, records, inventory, lifecycle, operation):
"""Translate only the known three-member stack and retain rollback contracts."""
intent = primary.get('native_stack_intent', {})
if intent.get('adapter', {}).get('name') != 'install_nextcloud_stack.sh':
@@ -60,13 +60,13 @@ def nextcloud_plan(primary, records, inventory, lifecycle):
snapshot = copy.deepcopy(translated[service['vmid']])
snapshot.pop('stack', None)
parent['stack']['members'].append(snapshot)
plan = oci_stack_plan.build(parent, translated, inventory, 'update')
plan = oci_stack_plan.build(parent, translated, inventory, operation)
plan['original_members'] = copy.deepcopy(list(records.values()))
plan['nextcloud_replay'] = True
return plan
def paperless_plan(primary, records, inventory, lifecycle):
def paperless_plan(primary, records, inventory, lifecycle, operation):
"""Prepare the known Paperless stack without publishing translated recipes."""
if primary.get('native_stack_intent', {}).get('adapter', {}).get('name') != 'install_paperless_stack.sh':
raise ValueError(f"{translate('Unrecognized stack adapter:')} Paperless")
@@ -89,13 +89,13 @@ def paperless_plan(primary, records, inventory, lifecycle):
snapshot = copy.deepcopy(translated[service['vmid']])
snapshot.pop('stack', None)
parent['stack']['members'].append(snapshot)
plan = oci_stack_plan.build(parent, translated, inventory, 'update')
plan = oci_stack_plan.build(parent, translated, inventory, operation)
plan['original_members'] = copy.deepcopy(list(records.values()))
plan['paperless_replay'] = True
return plan
def tandoor_plan(primary, records, inventory, lifecycle):
def tandoor_plan(primary, records, inventory, lifecycle, operation):
"""Prepare exactly the application and PostgreSQL without publishing state."""
if primary.get('native_stack_intent', {}).get('adapter', {}).get('name') != 'install_tandoor_stack.sh':
raise ValueError(f"{translate('Unrecognized stack adapter:')} Tandoor")
@@ -118,13 +118,13 @@ def tandoor_plan(primary, records, inventory, lifecycle):
snapshot = copy.deepcopy(translated[service['vmid']])
snapshot.pop('stack', None)
parent['stack']['members'].append(snapshot)
plan = oci_stack_plan.build(parent, translated, inventory, 'update')
plan = oci_stack_plan.build(parent, translated, inventory, operation)
plan['original_members'] = copy.deepcopy(list(records.values()))
plan['tandoor_replay'] = True
return plan
def immich_plan(primary, records, inventory, lifecycle):
def immich_plan(primary, records, inventory, lifecycle, operation):
if primary.get('native_stack_intent', {}).get('adapter', {}).get('name') != 'install_immich_stack.sh':
raise ValueError(f"{translate('Unrecognized stack adapter:')} Immich")
translated = {vmid: replay.immich_record(record) for vmid, record in records.items()}
@@ -145,7 +145,7 @@ def immich_plan(primary, records, inventory, lifecycle):
snapshot = copy.deepcopy(translated[service['vmid']])
snapshot.pop('stack', None)
parent['stack']['members'].append(snapshot)
plan = oci_stack_plan.build(parent, translated, inventory, 'update')
plan = oci_stack_plan.build(parent, translated, inventory, operation)
plan['original_members'] = copy.deepcopy(list(records.values()))
plan['immich_replay'] = True
return plan
@@ -262,7 +262,11 @@ class NativeAdapter:
current['pending_stack_transaction'] = str(self.journal)
instances.write(instances.location(self.root, vmid), current)
config = instances.command('pct', 'config', str(record['vmid']))
archive, digest = resolve_archive(record, config)
if operation == 'recreate':
digest = record.get('observed', {}).get('resolved_registry_digest')
archive, digest = resolve_archive(record, config, required_digest=digest)
else:
archive, digest = resolve_archive(record, config)
image = image_from_archive(str(archive))
old = record['observed']['image']
validate_database_transition(old, image)
@@ -279,7 +283,13 @@ class NativeAdapter:
msg_info(f"{translate('Checking the new image without starting it:')} {self.describe(record['vmid'])}")
self.probe_nextcloud_image(record, archive, image)
msg_ok(f"{translate('New image compatible:')} {self.describe(record['vmid'])}")
return {'archive': str(archive), 'digest': digest}
proposal = None
if operation == 'recreate':
proposal = {'operation': 'recreate', 'candidate': copy.deepcopy(record)}
config_hash = record.get('observed', {}).get('config_sha256')
if config_hash:
proposal['base_config_sha256'] = config_hash
return {'archive': str(archive), 'digest': digest, 'proposal': proposal}
def probe_nextcloud_image(self, record, archive, image):
"""Import but never start a disposable rootfs before stopping the stack."""
@@ -346,9 +356,12 @@ class NativeAdapter:
context.update(tandoor_replay=True, effective_record=self.records[vmid])
if self.plan.get('immich_replay'):
context.update(immich_replay=True, effective_record=self.records[vmid])
member_tx.apply(self.root, vmid, Path(prepared['archive']), 'update',
registry_digest=prepared['digest'], acknowledge_external_data=self.acknowledge,
coordinated=context, progress=f"{translate('Updating')} {self.describe(vmid)}:")
operation = self.plan['operation']
member_tx.apply(self.root, vmid, Path(prepared['archive']), operation,
proposal=prepared.get('proposal'), registry_digest=prepared['digest'],
acknowledge_external_data=self.acknowledge, coordinated=context,
progress=f"{translate('Recreating') if operation == 'recreate' else translate('Updating')} "
f"{self.describe(vmid)}:")
def start(self, vmid):
self.validate(self.plan)
@@ -617,10 +630,13 @@ class NativeAdapter:
_current = {'journal': None, 'primary': None}
def run(vmid, recover=False, acknowledge_external_data=False, keep_backup=None):
def run(vmid, recover=False, acknowledge_external_data=False, keep_backup=None, operation='update'):
if operation not in ('update', 'recreate'):
raise ValueError(translate('Invalid stack operation'))
root = instances.ROOT
msg_info(translate('Checking the interrupted stack operation...') if recover
else translate('Checking the stack before the update...'))
else (translate('Checking the stack before recreating...') if operation == 'recreate'
else translate('Checking the stack before the update...')))
with instances.locked(root):
selected = instances.read(root, vmid)
primary_id = selected.get('stack_member', {}).get('primary_vmid', vmid)
@@ -665,9 +681,9 @@ def run(vmid, recover=False, acknowledge_external_data=False, keep_backup=None):
if not stat.S_ISREG(info.st_mode) or info.st_uid != 0 or info.st_mode & 0o077:
raise ValueError(translate('Unsafe dependency contract'))
builder = builders[adapter_name]
plan = builder(primary, records, inventory, json.loads(lifecycle.read_text()))
plan = builder(primary, records, inventory, json.loads(lifecycle.read_text()), operation)
else:
plan = oci_stack_plan.build(primary, records, inventory, 'update')
plan = oci_stack_plan.build(primary, records, inventory, operation)
stack_tx.validate_plan(plan)
directory = instances.location(root, primary_id).parent / 'stack-transactions' / uuid.uuid4().hex
private_directory(directory)
@@ -689,10 +705,11 @@ def run(vmid, recover=False, acknowledge_external_data=False, keep_backup=None):
adapter.keep_backup = keep_backup
import oci_operation_notice
import oci_update_current
with oci_operation_notice.operation([member['vmid'] for member in plan['members']], 'update',
with oci_operation_notice.operation([member['vmid'] for member in plan['members']], operation,
oci_update_current.application_name(primary, primary_id), primary_id):
result = stack_tx.execute(journal, adapter, plan)
msg_ok(translate('Stack update completed. Data kept.'))
msg_ok(translate('Stack recreation completed. Data kept.') if operation == 'recreate'
else translate('Stack update completed. Data kept.'))
return result
@@ -728,13 +745,14 @@ def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('vmid', type=int)
parser.add_argument('--recover', action='store_true')
parser.add_argument('--operation', choices=['update', 'recreate'], default='update')
parser.add_argument('--acknowledge-external-data', action='store_true')
parser.add_argument('--keep-backup', metavar='STORAGE')
args = parser.parse_args()
if os.geteuid() != 0:
parser.error(translate('Root privileges on the Proxmox node are required'))
try:
run(args.vmid, args.recover, args.acknowledge_external_data, args.keep_backup)
run(args.vmid, args.recover, args.acknowledge_external_data, args.keep_backup, args.operation)
return 0
except BlockingIOError:
msg_error(translate('Another OCI operation is using the registry. This operation was not started.'))
+10 -2
View File
@@ -50,20 +50,28 @@ def run_quiet(args, error, capture=False):
process.wait()
def resolve_archive(desired, config, current=None, check=None):
def resolve_archive(desired, config, current=None, check=None, required_digest=None):
# Shared by individual and coordinated operations; no guest mutation here.
# When the registry still serves the current digest nothing is downloaded.
reference = desired['template']['container_contract']['image']['reference']
architecture = transaction.parse_config(config)['arch']
msg_info(translate('Checking the image in the registry...'))
transaction.log(f'image: {reference} ({architecture})')
if required_digest is not None and not re.fullmatch(r'sha256:[a-f0-9]{64}', required_digest):
raise ValueError(translate('Invalid saved image digest'))
lookup = repository(reference) + '@' + required_digest if required_digest else reference
code = ('import json,sys; from oci_installation_state import resolve_candidate; '
'print(json.dumps(resolve_candidate(sys.argv[1],sys.argv[2])))')
candidate = json.loads(run_quiet([sys.executable, '-c', code, reference, architecture],
candidate = json.loads(run_quiet([sys.executable, '-c', code, lookup, architecture],
translate('Could not query the image registry'), capture=True))
digest = candidate['manifest_digest']
if not re.fullmatch(r'sha256:[a-f0-9]{64}', digest):
raise ValueError(translate('Invalid registry digest'))
registry_digest = candidate.get('registry_digest')
if registry_digest is not None and not re.fullmatch(r'sha256:[a-f0-9]{64}', registry_digest):
raise ValueError(translate('Invalid registry digest'))
if required_digest is not None and registry_digest != required_digest:
raise ValueError(translate('The registry did not return the saved image digest'))
msg_ok(f"{translate('Image:')} {reference} ({candidate.get('version') or digest[7:19]})")
if current and same_image(candidate, current):
return None, digest
+2 -2
View File
@@ -595,9 +595,9 @@ def build_parser() -> argparse.ArgumentParser:
rclone_parser.add_argument("--host", default="auto")
rclone_parser.add_argument("--dry-run", action="store_true")
subparsers.add_parser("recover", help="Register again the OCI applications restored from a backup")
manage_parser = subparsers.add_parser("manage", help="Update or recreate one installed OCI instance")
manage_parser = subparsers.add_parser("manage", help="Update, modify or recreate one installed OCI instance")
manage_parser.add_argument("vmid", type=int)
manage_parser.add_argument("--action", choices=("update", "recreate"), required=True)
manage_parser.add_argument("--action", choices=("update", "modify", "recreate"), required=True)
manage_parser.add_argument("--keep-backup", metavar="STORAGE",
help="Keep the backup taken before the update in this Proxmox storage")
manage_parser.add_argument("--unattended", action="store_true",
+26 -12
View File
@@ -285,6 +285,8 @@ def manage_instance(project, ui, row, action=None, lifecycle_args=()):
return False
if row['stack']:
return _manage_stack(project, ui, row, action, lifecycle_args)
if action == 'modify':
return False
if not row['pending']:
if row['status'] != 'installed' or row['reason'] != 'matched':
ui.message(translate('The instance identity or status must be reviewed before updating.'), translate('OCI management'))
@@ -483,24 +485,32 @@ def _manage_stack(project, ui, row, action=None, lifecycle_args=()):
if action is None:
# A stack that cannot be updated can still be removed.
options = [('update', translate('Update every container of the application'))] if updatable else []
options.append(('recreate', translate('Recreate: add or remove extra paths and devices')))
options.append(('modify', translate('Modify extra paths and devices')))
if updatable:
options.append(('recreate', translate('Recreate every container with its saved configuration')))
options.append(('remove', translate('Remove: delete the application and its containers')))
action = ui.choose(translate('Manage OCI stack'), options, options[0][0])
if action is None:
return False
if action == 'remove':
return _remove(project, ui, primary_id)
if action == 'modify':
from .stack_recreation import modify_stack
return modify_stack(project, ui, primary, _run_lifecycle)
if action == 'recreate':
# Nothing is rebuilt: only the extra paths and devices of the
# application container change.
from .stack_recreation import recreate_stack
return recreate_stack(project, ui, primary, _run_lifecycle)
if not ui.review(translate('All {count} containers of the application are updated together (main CT: {vmid}). '
'If there are new versions, all images are downloaded and verified, the application '
'is stopped, each container is backed up and replaced with its new image. If anything '
'fails, the backups are restored.').format(count=len(members), vmid=primary_id),
translate('Update OCI stack'), question=translate('Update the whole stack?'), default=True):
return False
if not ui.review(translate('All {count} containers of the application will be recreated from their saved '
'image digests (main CT: {vmid}). The stack is stopped, every container is '
'backed up and replaced, then checked. If anything fails, the backups are restored.')
.format(count=len(members), vmid=primary_id), translate('Recreate OCI stack'),
question=translate('Recreate the whole stack?'), default=False):
return False
else:
if not ui.review(translate('All {count} containers of the application are updated together (main CT: {vmid}). '
'If there are new versions, all images are downloaded and verified, the application '
'is stopped, each container is backed up and replaced with its new image. If anything '
'fails, the backups are restored.').format(count=len(members), vmid=primary_id),
translate('Update OCI stack'), question=translate('Update the whole stack?'), default=True):
return False
else:
if not ui.review(translate('A coordinated operation has a saved journal. Continuing attempts to recover the previous stack where needed, or finish cleanup for a completed operation. Recovery or cleanup can fail.'), translate('Recover OCI stack'),
question=translate('Recover or complete the operation?'), default=True):
@@ -512,9 +522,13 @@ def _manage_stack(project, ui, row, action=None, lifecycle_args=()):
command.append('--recover')
else:
command.extend(lifecycle_args)
if action == 'recreate':
command.extend(['--operation', 'recreate'])
if '--acknowledge-external-data' not in command:
command.append('--acknowledge-external-data')
completed = _run_lifecycle(command, translate('Recover OCI stack') if pending else translate('Update OCI stack'))
title = (translate('Recover OCI stack') if pending else translate('Recreate OCI stack')
if action == 'recreate' else translate('Update OCI stack'))
completed = _run_lifecycle(command, title)
if completed and not pending and not getattr(ui, 'unattended', False):
images.offer_removal(ui, [int(member['vmid']) for member in members])
return completed
+9 -8
View File
@@ -1,6 +1,7 @@
"""Recreate for a multi-container application: add or remove the extra paths
and devices of its application container. Its own data, its database and the
other containers are never part of it."""
"""Modify the extra paths and devices of a multi-container application.
Only its application container is restarted. Its own data, its database and
the other containers are never part of this operation."""
from __future__ import annotations
import json
@@ -206,7 +207,7 @@ def change_recognition(project, ui, primary, member, run_lifecycle):
return run_lifecycle(command, translate('Recreate OCI'))
def recreate_stack(project, ui, primary, run_lifecycle):
def modify_stack(project, ui, primary, run_lifecycle):
learning = immich_learning(primary)
if learning is not None:
what = ui.choose(translate('What to recreate'),
@@ -230,13 +231,13 @@ def recreate_stack(project, ui, primary, run_lifecycle):
member = next(m for m in members if str(m['vmid']) == selected)
changes = plan_changes(ui, member)
if changes is None:
ui.message(translate('Nothing was changed.'), translate('Recreate OCI'))
ui.message(translate('Nothing was changed.'), translate('Modify OCI stack'))
return False
if not ui.review(summary(member, changes), translate('Recreate OCI'),
question=translate('Recreate with these options?'), default=True):
if not ui.review(summary(member, changes), translate('Modify OCI stack'),
question=translate('Apply these changes?'), default=True):
return False
with tempfile.NamedTemporaryFile(mode='w', suffix='.json') as file:
json.dump(changes, file)
file.flush()
return run_lifecycle([sys.executable, str(project / 'remote/oci_stack_modify.py'),
str(member['vmid']), '--changes', file.name], translate('Recreate OCI'))
str(member['vmid']), '--changes', file.name], translate('Modify OCI stack'))
@@ -0,0 +1,58 @@
"""A Recreate pins the registry/index digest while downloading a platform manifest."""
import json
from pathlib import Path
import sys
import unittest
from unittest.mock import patch
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT / "remote"))
import oci_installation_state as state
import oci_update_current as update
def digest(data):
return "sha256:" + state.sha(data)
class RecreateRegistryDigestTests(unittest.TestCase):
def test_candidate_keeps_the_multi_arch_registry_digest_separate_from_platform_manifest(self):
platform = json.dumps({
"schemaVersion": 2,
"config": {"digest": "sha256:" + "c" * 64},
"layers": [],
}).encode()
index = json.dumps({
"schemaVersion": 2,
"manifests": [{"digest": digest(platform),
"platform": {"architecture": "amd64", "os": "linux"}}],
}).encode()
config = json.dumps({"architecture": "amd64", "os": "linux",
"config": {"Labels": {}, "Env": []}}).encode()
with patch.object(state, "command", side_effect=[index, platform, config]):
candidate = state.resolve_candidate("postgres@" + digest(index), "amd64")
self.assertEqual(candidate["registry_digest"], digest(index))
self.assertEqual(candidate["manifest_digest"], digest(platform))
def test_recreate_accepts_saved_index_digest_and_returns_platform_manifest(self):
index_digest = "sha256:" + "a" * 64
platform_digest = "sha256:" + "b" * 64
candidate = {"registry_digest": index_digest, "manifest_digest": platform_digest,
"layers": [], "created": "2026-10-05T00:00:00Z"}
desired = {"template": {"container_contract": {"image": {"reference": "postgres:latest"}}},
"deployment": {"template_storage": "local"}}
with patch.object(update, "run_quiet", return_value=json.dumps(candidate)) as query:
archive, digest_value = update.resolve_archive(
desired, b"arch: amd64\n", current={"manifest_digest": platform_digest},
required_digest=index_digest)
self.assertIsNone(archive)
self.assertEqual(digest_value, platform_digest)
self.assertIn("postgres@" + index_digest, query.call_args.args[0])
if __name__ == "__main__":
unittest.main()
+70
View File
@@ -0,0 +1,70 @@
"""The coordinated Recreate path must remain distinct from Update."""
from pathlib import Path
import sys
import unittest
from unittest.mock import patch
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT / "remote"))
import oci_stack_native as native
DIGEST = "sha256:" + "a" * 64
MEMBER = {
"vmid": 138,
"installation_id": "blinko-installation",
"template": {"id": "blinko", "container_contract": {"image": {"reference": "blinkospace/blinko:latest"}}},
"deployment": {"mounts": []},
"observed": {"config_sha256": "saved-config", "resolved_registry_digest": DIGEST,
"image": {"architecture": "amd64", "os": "linux", "defaults": {}}},
"stack": {"deployment": {"services": [{"vmid": 138, "name": "Blinko",
"healthcheck": {"type": "running", "timeout_seconds": 120}}]}},
}
class StackRecreateV1Tests(unittest.TestCase):
def adapter(self):
plan = {"operation": "recreate", "primary_vmid": 138, "members": [MEMBER]}
return native.NativeAdapter(Path("/registry"), Path("/journal"), plan)
def test_prepare_resolves_the_saved_digest_and_builds_an_unchanged_proposal(self):
adapter = self.adapter()
image = {"architecture": "amd64", "os": "linux", "defaults": {}}
with patch.object(native.instances, "read", return_value={}), \
patch.object(native.instances, "write"), \
patch.object(native.instances, "location", return_value=Path("/registry/138.json")), \
patch.object(native.instances, "command", return_value=b"arch: amd64\n"), \
patch.object(native, "resolve_archive", return_value=(Path("/cache/blinko.tar"), DIGEST)) as resolve, \
patch.object(native, "image_from_archive", return_value=image):
prepared = adapter.prepare(MEMBER, "recreate")
resolve.assert_called_once_with(MEMBER, b"arch: amd64\n", required_digest=DIGEST)
self.assertEqual(prepared["digest"], DIGEST)
self.assertEqual(prepared["proposal"]["operation"], "recreate")
self.assertEqual(prepared["proposal"]["candidate"], MEMBER)
self.assertEqual(prepared["proposal"]["base_config_sha256"], "saved-config")
def test_replace_passes_the_recreate_operation_and_proposal_to_every_member(self):
adapter = self.adapter()
prepared = {"archive": "/cache/blinko.tar", "digest": DIGEST,
"proposal": {"operation": "recreate", "candidate": MEMBER,
"base_config_sha256": "saved-config"}}
with patch.object(adapter, "validate"), \
patch.object(adapter, "state", return_value={"backups": {"138": {"archive": "/backup"}}}), \
patch.object(native.member_tx, "apply") as apply:
adapter.replace(138, prepared, "transaction-id")
self.assertEqual(apply.call_args.args[3], "recreate")
self.assertEqual(apply.call_args.kwargs["proposal"], prepared["proposal"])
self.assertEqual(apply.call_args.kwargs["registry_digest"], DIGEST)
self.assertIn("Recreating", apply.call_args.kwargs["progress"])
def test_an_unknown_stack_operation_is_rejected_before_the_stack_is_touched(self):
with self.assertRaises(ValueError):
native.run(138, operation="not-a-lifecycle-operation")
if __name__ == "__main__":
unittest.main()