Page MenuHomePhorge

No OneTemporary

Size
21 KB
Referenced Files
None
Subscribers
None
diff --git a/lilybuild/lilybuild/podman_helper.py b/lilybuild/lilybuild/podman_helper.py
index 17a69e5..500f3ca 100755
--- a/lilybuild/lilybuild/podman_helper.py
+++ b/lilybuild/lilybuild/podman_helper.py
@@ -1,544 +1,547 @@
#!/usr/bin/env python3
# This file is part of lilybuild.
# SPDX-FileCopyrightText: 2025-2026 tusooa <tusooa@kazv.moe>
# SPDX-License-Identifier: GPL-2.0-only
import subprocess
import sys
import os
import json
import random
import traceback
import string
import time
import hashlib
import re
import tempfile
import lilybuild.safetar
col_info = '\x1b[1;34m[INFO]'
col_success = '\x1b[1;32m[SUCC]'
col_warn = '\x1b[1;33m<WARN>'
col_error = '\x1b[1;31m!ERROR!'
col_reset = '\x1b[0m'
any_spaces_re = re.compile(r'\s')
def pinfo(*args, **kwargs):
print(col_info, *args, col_reset, **kwargs)
sys.stdout.flush()
def perror(*args, **kwargs):
print(col_error, *args, col_reset, **kwargs)
sys.stdout.flush()
def pwarn(*args, **kwargs):
print(col_warn, *args, col_reset, **kwargs)
sys.stdout.flush()
def psuccess(*args, **kwargs):
print(col_success, *args, col_reset, **kwargs)
sys.stdout.flush()
def gen_random_id():
# https://stackoverflow.com/questions/2257441/random-string-generation-with-upper-case-letters-and-digits
return ''.join(random.SystemRandom().choice(string.ascii_lowercase + string.digits) for _ in range(10))
def image_to_podman_args(image):
name = image['name']
args = []
if 'entrypoint' in image:
# ci.json requires that the entrypoint is an array of strings
ep = json.dumps(image['entrypoint'])
args += ['--entrypoint', ep]
args += ['--', name]
return args
class PodmanHelper:
cache_storage_root_dir = '/cache'
cache_max_bytes = 10 * 1024 * 1024 * 1024
work_vol_mount_dir = '/build'
script_vol_mount_dir = '/script'
cache_tmp_dir = '/tmp/cache/'
script_name = script_vol_mount_dir + '/run.sh'
env_file_basename = 'env'
metadata_file_basename = 'metadata.json'
cur_cache_basename = 'cur'
volume_helper_image = os.environ.get('LILYBUILD_VOLUME_HELPER_IMAGE', 'r.lily-is.land/infra/lilybuild/volume-helper:servant')
key_file_pub = '/secrets/lilybuild-volume-helper-key.pub'
key_file_sub = '/secrets/lilybuild-volume-helper-key'
ssh_port = '2222'
ssh_command = f'ssh -p {ssh_port} -i {key_file_sub} -oStrictHostKeyChecking=no -oUserKnownHostsFile=/dev/null'
ssh_command_list = ['ssh', '-p', ssh_port, '-i', key_file_sub, '-oStrictHostKeyChecking=no', '-oUserKnownHostsFile=/dev/null']
worker_container_name = os.environ.get('HOSTNAME', '')
ssh_max_wait = 10
ssh_wait_interval_sec = 1
service_max_wait_sec = 60 * 5
service_wait_interval_sec = 10
container_run_timeout_sec = 60 * 60 * 2 # 2 hours by default
def __init__(self, **kwargs):
self.volumes_to_remove = []
self.helper_container_id = None
self.service_network_id = None
self.service_containers = []
self.metadata = {
'repo_id': None,
'caches': [],
'cache_last_invalidated_sec': 0,
'protected': False,
}
self.__dict__.update(kwargs)
# This calls podman which is hard to duplicate so we mock this function
# in tests instead
def verbose_run(self, *args, **kwargs):
print('run:', args, kwargs)
sys.stdout.flush()
return subprocess.run(*args, **kwargs)
def create_volume(self, t):
res = self.verbose_run([
'podman', 'volume', 'create',
'--label', 'lilybuild=' + t,
], check=True, capture_output=True, encoding='utf-8')
volname = res.stdout.strip()
self.volumes_to_remove.append(volname)
return volname
def clean_volumes(self):
self.verbose_run([
'podman', 'volume', 'rm', '-f', '--',
] + self.volumes_to_remove, capture_output=True)
def clean_helper_container(self):
self.verbose_run([
'podman', 'container', 'rm', '-f', '--', self.helper_container_id,
], capture_output=True)
def start_helper_service(self, work_volname, script_volname):
res = self.verbose_run([
'podman', 'container', 'inspect', '--', self.worker_container_name,
], check=True, capture_output=True, encoding='utf-8')
container_stat = json.loads(res.stdout)[0]
pod = container_stat.get('Pod')
networks = list(container_stat.get('NetworkSettings').get('Networks').keys())
alias = gen_random_id()
container_name = 'lilybuild-helper-' + alias
with open(self.key_file_pub) as f:
pub_key = f.readline().strip()
res = self.verbose_run([
'podman', 'run', '--rm', '-d', '--name', container_name,
+ '--pull=always',
f'--mount=type=volume,source={work_volname},destination={self.work_vol_mount_dir}',
f'--mount=type=volume,source={script_volname},destination={self.script_vol_mount_dir}',
f'--pod={pod}',
f'--net={networks[0]}',
f'--network-alias={alias}',
'--image-volume=ignore',
'--label', 'lilybuild=helper',
'-e', 'PUID=0',
'-e', 'PGID=0',
'-e', f'PUBLIC_KEY={pub_key}',
'-e', 'USER_NAME=helper',
'-e', 'SUDO_ACCESS=true',
'--',
self.volume_helper_image,
], check=True, capture_output=True, encoding='utf-8')
self.helper_container_id = container_name
self.helper_container_alias = alias
pinfo('Waiting for ssh service to be up...')
service_up = False
for i in range(self.ssh_max_wait):
chk = self.verbose_run(['nc', alias, self.ssh_port], input=b'', capture_output=True)
if chk.returncode == 0 and chk.stdout is not None and chk.stdout.startswith(b'SSH'):
service_up = True
break
else:
time.sleep(self.ssh_wait_interval_sec)
if not service_up:
raise RuntimeError('Service is still not up!')
psuccess('Service is up.')
return (container_name, alias)
def get_valid_caches(self, md):
using_caches = []
for c in md['caches']:
if c['policy'] != 'pull-push' and c['policy'] != 'pull':
continue
storage_dir = self.get_cache_storage_dir(md['repo_id'], c, protected=md['protected'])
pinfo('Cache storage dir is ', storage_dir)
if os.path.exists(storage_dir):
pinfo('Cache exists')
cur_cache_name = os.path.join(storage_dir, self.cur_cache_basename)
try:
stat_res = os.stat(cur_cache_name)
if stat_res.st_mtime > md['cache_last_invalidated_sec']:
pinfo('Cache not expired')
using_caches.append(cur_cache_name)
else:
pinfo('Cache expired')
except:
pinfo('Cache does not exist')
return using_caches
def import_caches(self, md, vol_mount_dir):
# Ensure we do not accidentally remove root
# although it should be pretty safe (it's a constant), but who knows
# Validate against spaces because anything passed after ssh is processed
# through a shell
if not (vol_mount_dir and
isinstance(vol_mount_dir, str) and
not any_spaces_re.search(vol_mount_dir)):
perror('vol_mount_dir cannot be empty and cannot contain spaces')
raise RuntimeError('vol_mount_dir cannot be empty and cannot contain spaces')
valid_caches = self.get_valid_caches(md)
cache_file = os.path.join(self.cache_tmp_dir, self.cur_cache_basename)
for c in valid_caches:
try:
pinfo('Uploading cache...')
self.verbose_run([
'rsync', '-a', '--delete',
'--rsh', self.ssh_command,
'--',
c,
f'helper@{self.helper_container_alias}:{self.cache_tmp_dir}',
], check=True)
pinfo('Extracting cache...')
self.verbose_run(self.ssh_command_list + [
f'helper@{self.helper_container_alias}',
'tar', '-xf', cache_file, '-C', vol_mount_dir,
], check=True)
except subprocess.CalledProcessError as e:
pwarn('Error when importing cache:', e)
# Cache is corrupt and should not be trusted
self.verbose_run(self.ssh_command_list + [
f'helper@{self.helper_container_alias}',
'rm', '-rf', '--', f'{vol_mount_dir}/*', f'{vol_mount_dir}/.*',
])
except:
pwarn('Other error occurred', sys.exception())
finally:
pinfo('Removing uploaded cache archive...')
self.verbose_run(self.ssh_command_list + [
f'helper@{self.helper_container_alias}',
'rm', '-f', '--', cache_file,
])
def save_caches(self, md, /, succeeded):
for c in md['caches']:
if c['policy'] != 'pull-push' and c['policy'] != 'push':
pinfo('Cache saving skipped because of policy')
continue
if (c['when'] == 'on_success' and not succeeded) or (c['when'] == 'on_failure' and succeeded):
pinfo('Cache saving skipped because mismatch in success status', c)
continue
storage_dir = self.get_cache_storage_dir(md['repo_id'], c, protected=md['protected'])
cache_file = os.path.join(storage_dir, self.cur_cache_basename)
try:
replaced = False
os.makedirs(storage_dir, exist_ok=True)
fd, fn = tempfile.mkstemp(dir=storage_dir)
os.close(fd)
lilybuild.safetar.create(
fn,
self.result_dir,
c['paths'],
self.cache_max_bytes,
items_to_exclude=None,
compression='gz',
)
os.replace(fn, cache_file)
replaced = True
except:
pwarn('Unable to create cache', sys.exception())
finally:
# Either the temp file is renamed, or it is not
if not replaced:
try:
os.remove(fn)
except:
pass
def import_volume(self, local_dir, vol_mount_dir):
# I'll just use the shell instead of pipe2+fork+exec+wait, much easier
self.verbose_run([
'rsync', '-a',
'--rsh', self.ssh_command,
f'{local_dir}/',
f'helper@{self.helper_container_alias}:{vol_mount_dir}',
], check=True)
def export_volume(self, local_dir, vol_mount_dir):
self.verbose_run([
'rsync', '-a', '--delete',
'--rsh', self.ssh_command,
f'helper@{self.helper_container_alias}:{vol_mount_dir}/',
local_dir,
], check=True)
def create_service_network(self):
res = self.verbose_run([
'podman', 'network', 'create', '--label', 'lilybuild=service-network'
], capture_output=True, check=True, encoding='utf-8')
self.service_network_id = res.stdout.strip()
return self.service_network_id
def maybe_clean_service_network(self):
if self.service_network_id is None:
return
res = self.verbose_run([
'podman', 'network', 'rm', '-f', '--', self.service_network_id
], capture_output=True, encoding='utf-8')
if res.returncode != 0:
perror('Cannot remove service network.')
def start_and_record_service_container(self, service):
image = service['name']
ep_args = []
if service['entrypoint']:
if isinstance(service['entrypoint'], str):
entrypoint = service['entrypoint']
else:
entrypoint = json.dumps(service['entrypoint'])
ep_args += [f'--entrypoint={entrypoint}']
cmd_args = []
if service['command']:
if isinstance(service['command'], str):
cmd_args += [service['command']]
else:
cmd_args += service['command']
res = self.verbose_run([
'podman', 'run', '-d', '--label', 'lilybuild=job-service',
+ '--pull=always',
f'--env-file={self.env_filename}',
f'--network={self.service_network_id}',
] + [
f'--network-alias={alias}' for alias in service['aliases']
] + ep_args + [
'--',
image,
] + cmd_args, check=True, capture_output=True, encoding='utf-8')
service_id = res.stdout.strip()
self.service_containers.append(service_id)
def ensure_service_containers_up(self):
waiting_container_ids = self.service_containers[:]
steady_deadline = time.monotonic() + self.service_max_wait_sec
pinfo('Waiting for service containers...')
while waiting_container_ids:
for cid in waiting_container_ids[:]:
res = self.verbose_run([
'podman', 'container', 'inspect', '--', cid
], check=True, capture_output=True, encoding='utf-8')
ins = json.loads(res.stdout)[0]
if ins.get('State', {}).get('Status') == 'running':
psuccess(f'Container {cid} is up')
waiting_container_ids.remove(cid)
if waiting_container_ids:
if time.monotonic() > steady_deadline:
perror('Containers are not yet up after deadline.')
raise TimeoutError('Service containers startup timeout')
pinfo('Some containers are not yet up. Waiting...')
time.sleep(self.service_wait_interval_sec)
psuccess('All service containers are up.')
def maybe_prune_service_containers(self):
container_ids = self.service_containers
if not container_ids:
return
stop_proc = self.verbose_run(['podman', 'container', 'stop', '--'] + container_ids)
if stop_proc.returncode != 0:
pwarn('Cannot stop container.')
# -v removes anonymous volumes associated with the container
rm_proc = self.verbose_run(['podman', 'container', 'rm', '-f', '-v', '--'] + container_ids)
def run_in_container(self, image, work_volname, script_volname):
timeout = self.container_run_timeout_sec
steady_deadline = time.monotonic() + timeout
network_args = []
if self.service_network_id:
network_args += [f'--network={self.service_network_id}']
start_process = self.verbose_run([
'podman', 'run', '-d',
+ '--pull=always',
f'--mount=type=volume,source={work_volname},destination={self.work_vol_mount_dir}',
f'--mount=type=volume,source={script_volname},destination={self.script_vol_mount_dir}',
f'--env-file={self.env_filename}',
] + network_args + image_to_podman_args(image) + [
self.script_name,
], capture_output=True, encoding='utf-8')
if start_process.returncode != 0:
perror('Cannot run container. Error message:')
print(start_process.stderr)
return start_process.returncode
container_id = start_process.stdout.strip()
steady_now = time.monotonic()
log_args = []
retcode = None
try:
while steady_deadline > steady_now:
log_process = self.verbose_run([
'podman', 'logs', '--follow'
] + log_args + ['--', container_id], timeout=steady_deadline - steady_now)
# Exited from `podman logs`: why? Is the container still running?
inspect_running = self.verbose_run([
'podman', 'container', 'inspect',
'--format', '{{.State.Status}}', '--', container_id,
], capture_output=True, encoding='utf-8', check=True)
if inspect_running.stdout.strip() == 'exited':
inspect_retcode = self.verbose_run([
'podman', 'container', 'inspect',
'--format', '{{.State.ExitCode}}', '--', container_id,
], capture_output=True, encoding='utf-8', check=True)
retcode = int(inspect_retcode.stdout.strip())
break
else:
pwarn('`podman logs` unexpectedly quits when the container is still running, resuming logs...')
log_args = ['--tail', '10']
steady_now = time.monotonic()
if retcode is None:
perror('Command timed out.')
retcode = 1
except subprocess.TimeoutExpired as e:
perror('Command timed out.')
retcode = 1
except subprocess.CalledProcessError as e:
perror('Cannot inspect container:', e)
except:
perror('Another exception happened:', sys.exception())
finally:
pinfo('Cleaning up container...')
stop_proc = self.verbose_run(['podman', 'container', 'stop', '--', container_id])
if stop_proc.returncode != 0:
pwarn('Cannot stop container.')
# -v removes anonymous volumes associated with the container
rm_proc = self.verbose_run(['podman', 'container', 'rm', '-f', '-v', '--', container_id])
pinfo('Cleaned.')
return retcode
def hash_cache_key(self, cache_key):
if not isinstance(cache_key, str):
cache_key = ''
m = hashlib.sha256()
m.update(cache_key.encode())
return m.hexdigest()
def get_cache_storage_dir(self, repo_id, cache_def, /, protected):
cache_key = cache_def.get('key', '')
hashed_key = self.hash_cache_key(cache_key)
return os.path.join(
self.cache_storage_root_dir,
'repos',
str(repo_id),
'protected' if protected else 'unprotected',
'cache-keys',
hashed_key,
)
def main(self, argv):
image = json.loads(argv[1])
self.work_dir = argv[2]
self.script_dir = argv[3]
self.result_dir = argv[4]
self.env_filename = os.path.join(self.script_dir, self.env_file_basename)
services = []
if len(argv) >= 6:
services = json.loads(argv[5])
metadata_filename = os.path.join(self.script_dir, self.metadata_file_basename)
if os.path.exists(metadata_filename):
pinfo('Parsing metadata...')
with open(metadata_filename) as f:
self.metadata = json.loads(f.read())
psuccess('Parsed.')
pinfo('Creating volumes...')
work_vol = self.create_volume('work')
script_vol = self.create_volume('script')
psuccess('Created.')
pinfo('Starting helper service...')
self.start_helper_service(work_vol, script_vol)
psuccess('Started...')
if services:
pinfo('Creating service network...')
self.create_service_network()
psuccess('Created.')
pinfo('Starting job-defined services...')
for service in services:
self.start_and_record_service_container(service)
pinfo('Waiting for job-defined services...')
self.ensure_service_containers_up()
if self.metadata['caches']:
pinfo('Importing caches...')
self.import_caches(self.metadata, self.work_vol_mount_dir)
psuccess('Imported.')
pinfo('Importing volumes...')
self.import_volume(self.work_dir, self.work_vol_mount_dir)
self.import_volume(self.script_dir, self.script_vol_mount_dir)
psuccess('Imported.')
pinfo('Running container...')
retcode = self.run_in_container(image, work_vol, script_vol)
succeeded = retcode == 0
pinfo(f'Returned {retcode}.')
if not succeeded:
perror('Job failed.')
else:
psuccess('Job succeeded.')
# We should collect the result regardless whether it succeeded
pinfo('Collecting build changes...')
self.export_volume(self.result_dir, self.work_vol_mount_dir)
psuccess('Collected.')
if self.metadata['caches']:
pinfo('Saving caches...')
self.save_caches(self.metadata, succeeded=succeeded)
psuccess('Saved.')
return retcode
def cleanup_all(self):
pinfo('Cleaning service containers...')
self.maybe_prune_service_containers()
psuccess('Cleaned.')
pinfo('Cleaning service network...')
self.maybe_clean_service_network()
psuccess('Cleaned.')
if self.helper_container_id:
pinfo('Cleaning helper container')
self.clean_helper_container()
psuccess('Cleaned.')
pinfo('Cleaning volumes...')
self.clean_volumes()
psuccess('Cleaned.')
if __name__ == '__main__':
retcode = 1
try:
ph = PodmanHelper()
retcode = ph.main(sys.argv)
except Exception as e:
perror('Error!', e)
print(traceback.format_exc())
raise
finally:
ph.cleanup_all()
sys.exit(retcode)

File Metadata

Mime Type
text/x-diff
Expires
Mon, Oct 12, 11:22 AM (1 d, 14 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1786162
Default Alt Text
(21 KB)

Event Timeline