Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85805970
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
21 KB
Referenced Files
None
Subscribers
None
View Options
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
Details
Attached
Mime Type
text/x-diff
Expires
Mon, Oct 12, 11:22 AM (1 d, 13 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1786162
Default Alt Text
(21 KB)
Attached To
Mode
rB lilybuild
Attached
Detach File
Event Timeline
Log In to Comment