Page MenuHomePhorge

No OneTemporary

Size
78 KB
Referenced Files
None
Subscribers
None
diff --git a/lilybuild/lilybuild/ci_steps.py b/lilybuild/lilybuild/ci_steps.py
index 2e1e8b3..0c65af3 100644
--- a/lilybuild/lilybuild/ci_steps.py
+++ b/lilybuild/lilybuild/ci_steps.py
@@ -1,588 +1,608 @@
from buildbot.plugins import *
from buildbot.process import buildstep, logobserver
from buildbot.interfaces import IRenderable
from twisted.internet import defer
from .ci_syntax import ci_file
from .ci_syntax import rules as ci_rules
-from .helpers import rsync_rules_from_artifacts, get_job_script, normalize_image, normalize_services, ci_vars_to_env_file
+from .helpers import (
+ rsync_rules_from_artifacts,
+ get_job_script,
+ normalize_image,
+ normalize_services,
+ ci_vars_to_env_file,
+ generate_metadata_from_job,
+)
from .phorge import SendCoverageToPhorge
import re
import sys
import json
SAFETAR_EXEC = '/lilybuild/lilybuild/safetar.py'
COVERAGE_EXEC = '/lilybuild/lilybuild/coverage.py'
def on_success(step):
return step.build.results == util.SUCCESS
def on_always(_step):
return True
def fill_list(*args):
return list(args)
class RunCIJobStep(steps.BuildStep):
# 200 MiB
artifact_max_size = 200 * 1024 * 1024
default_image = 'alpine'
master_job_artifact_dir_pattern = '%(kw:st)s/repos/%(prop:lilybuild_repo_id)s/builds/%(prop:lilybuild_root_build_id)s/jobs/%(kw:job)s/artifacts'
artifact_file_name = 'artifacts.tar'
master_job_artifact_file_name_pattern = master_job_artifact_dir_pattern + '/' + artifact_file_name
reports_file_name = 'reports.tar'
master_reports_file_name_pattern = master_job_artifact_dir_pattern + '/' + reports_file_name
master_pages_dir_pattern = '%(kw:st)s/repos/%(prop:lilybuild_repo_id)s/pages'
phorge_coverage_file_name = 'coverage-phorge.json'
def __init__(
self,
lbc,
src_relative=None,
src_dir=None,
storage_dir=None,
repo_id=None,
result_relative=None,
result_dir=None,
artifact_stage_relative=None,
artifact_stage_dir=None,
job_prop=None,
artifact_link_base=None,
**kwargs):
self.lbc = lbc
self.src_relative = src_relative
self.src_dir = src_dir
self.work_root_dir = kwargs['workdir']
self.script_dir = 'script'
self.storage_dir = storage_dir
self.repo_id = repo_id
self.artifact_stage_relative = artifact_stage_relative
self.artifact_stage_dir = artifact_stage_dir
self.result_relative = result_relative
self.result_dir = result_dir
self.artifact_link_base = artifact_link_base
super().__init__(name='Run step', **kwargs)
def get_cur_repo_config(self):
return self.lbc.repos[self.getProperty('lilybuild_repo_id')]
@defer.inlineCallbacks
def run(self):
job_prop = self.getProperty('lilybuild_job_prop')
job = ci_file.CIJob.from_prop(job_prop)
job_index = self.getProperty('lilybuild_job_index')
variables = yield self.get_ci_variables(job)
should_run = True
if len(job.rules):
should_run = False
default_when = job.struct_raw.get('when', 'on_success')
for r in job.rules:
try:
when = r.get('when', default_when)
if when == 'never' or when == 'manual':
should_run_to_set = False
else:
should_run_to_set = True
rule_str = r.get('if')
if not rule_str:
# No condition == always true
# TODO: `changes` rule
should_run = should_run_to_set
break
res = ci_rules.evaluate_rule(ci_rules.parse_rule(rule_str), variables)
if not res:
continue
should_run = should_run_to_set
break
except SyntaxError:
self.addCompleteLog('error', f'Rule "{rule_str}" has syntax errors')
except:
pass
if should_run:
next_steps = self.job_to_steps(job, job_index, variables)
self.build.addStepsAfterCurrentStep(next_steps)
return util.SUCCESS
else:
self.addCompleteLog('info', 'Job skipped by a rule')
return util.SKIPPED
@defer.inlineCallbacks
def get_ci_variables(self, job):
res = {}
res.update(job.get_predefined_ci_variables())
res.update(self.getProperty('lilybuild_pipeline_vars'))
res['CI_JOB_IMAGE'] = job.image or self.default_image
res['CI_JOB_URL'] = yield self.build.getUrl()
res['CI_JOB_ID'] = self.build.buildid
res['CI_PROJECT_DIR'] = '/build'
try:
repo = self.lbc.repos[self.getProperty('lilybuild_repo_id')]
variables = yield repo['variables_getter'](self.build)
res_vars = {}
for var in variables:
value = variables[var]
if IRenderable.providedBy(value):
value = yield self.build.render(value)
res_vars[var] = value
res.update(res_vars)
except Exception as e:
self.addCompleteLog('exception', f'{e}')
return res
def get_upload_artifacts_jobs(self, short_name, artifact_type, artifact_name, base_dir, paths, exclude, master_pattern, job_index, doStepIf=on_success, has_pages=False):
archive_artifact_step = steps.ShellCommand(
name=f'Archive artifacts: {short_name}',
command=[
SAFETAR_EXEC,
],
initialStdin=json.dumps({
'op': 'create',
'archive_file': artifact_name,
'base_dir': base_dir,
'content': paths,
'items_to_exclude': exclude,
'compression': 'gz',
'limit_bytes': self.get_cur_repo_config()['artifact_uncompressed_limit'],
}),
workdir=self.work_root_dir,
doStepIf=doStepIf,
)
masterdest = util.Interpolate(
master_pattern,
st=self.storage_dir,
job=job_index,
doStepIf=doStepIf,
)
parent_build_id = self.getProperty('lilybuild_pipeline_vars')['CI_PIPELINE_ID']
artifact_url = f'{self.artifact_link_base}/plugins/lilybuild_artifacts/builds/{parent_build_id}/jobs/{job_index}/artifacts/{artifact_type}' if self.artifact_link_base else None
upload_artifact_step = steps.FileUpload(
workersrc=artifact_name,
maxsize=self.get_cur_repo_config()['artifact_compressed_limit'],
name=f'Upload artifacts: {short_name}',
masterdest=masterdest,
workdir=self.work_root_dir,
url=artifact_url,
doStepIf=doStepIf,
)
r = [archive_artifact_step, upload_artifact_step]
if has_pages:
r.append(steps.MasterShellCommand(
command=util.Transform(fill_list,
sys.executable,
'-m', 'lilybuild.pages',
util.Interpolate(
self.master_pages_dir_pattern,
st=self.storage_dir,
),
masterdest,
),
name='Deploy pages',
logEnviron=False,
doStepIf=doStepIf,
))
return r
def job_to_steps(self, job, job_index, variables):
script_name = self.script_dir + '/run.sh'
env_filename = self.script_dir + '/env'
+ metadata_filename = self.script_dir + '/metadata.json'
source_step = self.lbc.create_source_step()
script_step = steps.StringDownload(
get_job_script(job),
name='Set up script',
workerdest=script_name,
workdir=self.work_root_dir,
doStepIf=on_success,
)
env_step = steps.StringDownload(
ci_vars_to_env_file(variables),
name='Set up env file',
workerdest=env_filename,
workdir=self.work_root_dir,
doStepIf=on_success,
)
+ metadata_step = steps.StringDownload(
+ generate_metadata_from_job(
+ self.getProperty('lilybuild_repo_id'),
+ job,
+ variables,
+ ),
+ name='Set up metadata file',
+ workerdest=metadata_filename,
+ workdir=self.work_root_dir,
+ doStepIf=on_success,
+ )
+
chmod_step = steps.ShellCommand(
name='Make script executable',
command=['chmod', '+x', script_name],
workdir=self.work_root_dir,
doStepIf=on_success,
)
artifact_steps = []
dep_job_indices = self.getProperty('lilybuild_dependency_job_indices')
if dep_job_indices:
for i in dep_job_indices:
# The steps may not run or may not have an artifact even if it runs
download_job = steps.FileDownload(
mastersrc=util.Interpolate(
self.master_job_artifact_file_name_pattern,
st=self.storage_dir,
job=i,
),
maxsize=self.get_cur_repo_config()['artifact_compressed_limit'],
name=f'Download artifacts from job #{i}',
workerdest=self.artifact_file_name,
workdir=self.work_root_dir,
doStepIf=on_success,
haltOnFailure=False,
flunkOnFailure=False,
# https://github.com/buildbot/buildbot/issues/3709
blocksize=256 * 1024,
)
unarchive_job = steps.ShellCommand(
name=f'Unarchive artifacts from job #{i}',
command=[
SAFETAR_EXEC,
],
initialStdin=json.dumps({
'op': 'extract',
'archive_file': self.artifact_file_name,
'target_dir': self.src_relative,
}),
workdir=self.work_root_dir,
doStepIf=on_success,
haltOnFailure=False,
flunkOnFailure=False,
)
artifact_steps += [download_job, unarchive_job]
run_step = steps.ShellCommand(
name='Run script in container',
command=[
'/lilybuild/podman-helper',
normalize_image(job.image or self.default_image),
self.src_relative,
self.script_dir,
self.result_relative,
normalize_services(job.services),
],
workdir=self.work_root_dir,
doStepIf=on_success,
# 2h timeout by default
# TODO support timeout by each job
timeout=60 * 60 * 2,
)
clean_script_step = steps.ShellCommand(
name='Clean script dir',
command=[
'rm',
'-rf',
self.script_dir,
],
workdir=self.work_root_dir,
alwaysRun=True,
)
- steps_to_run = [source_step, script_step, chmod_step, env_step] + artifact_steps + [run_step, clean_script_step]
+ steps_to_run = [source_step, script_step, chmod_step, env_step, metadata_step] + artifact_steps + [run_step, clean_script_step]
if 'paths' in job.artifacts:
steps_to_run += self.get_upload_artifacts_jobs(
'files',
'archive',
self.artifact_file_name,
self.result_relative,
job.artifacts.get('paths', []),
job.artifacts.get('exclude', []),
self.master_job_artifact_file_name_pattern,
job_index,
has_pages=job.is_pages()
)
if job.has_supported_coverage_report():
steps_to_run += [steps.ShellCommand(
name='Process reports',
command=[
COVERAGE_EXEC,
],
initialStdin=json.dumps({
'source_dir': self.src_relative,
'result_dir': self.result_relative,
'untrusted_coverage_file': job.artifacts['reports']['coverage_report']['path'],
'output_dir': self.artifact_stage_relative,
}),
workdir=self.work_root_dir,
flunkOnFailure=False,
doStepIf=on_always,
)] + self.get_upload_artifacts_jobs(
'reports',
'reports',
self.reports_file_name,
self.artifact_stage_relative,
['*'],
[],
self.master_reports_file_name_pattern,
job_index,
doStepIf=on_always
) + [SendCoverageToPhorge(
self.lbc,
self.artifact_stage_relative + '/' + self.phorge_coverage_file_name,
workdir=self.work_root_dir,
)]
clean_stage_dir_again_step = steps.ShellCommand(
name='Clean stage, result and artifact',
command=[
'rm',
'-rf',
self.result_relative,
self.artifact_file_name,
self.artifact_stage_relative,
self.reports_file_name,
],
workdir=self.work_root_dir,
alwaysRun=True,
)
steps_to_run.append(clean_stage_dir_again_step)
return steps_to_run
class TriggerMultipleJobsStep(steps.Trigger):
properties_to_keep = [
'branch',
'revision',
'repository',
'harbormaster_build_target_phid',
'harbormaster_variable_buildable.diff',
'harbormaster_variable_repository.staging.ref',
'harbormaster_variable_repository.staging.uri',
'harbormaster_variable_repository.uri',
'lilybuild_repo',
'lilybuild_repo_id',
'lilybuild_pipeline_vars',
]
def __init__(self, lbc, jobs_with_data, **kwargs):
self.lbc = lbc
self.jobs = jobs_with_data
super().__init__(schedulerNames=[self.lbc.triggerable_scheduler_name], **kwargs)
def getSchedulersAndProperties(self):
ret = []
common_properties = {
'lilybuild_root_build_id': self.build.buildid,
}
for prop in self.properties_to_keep:
if self.hasProperty(prop):
common_properties[prop] = self.getProperty(prop)
for (job, i, dep_job_indices) in self.jobs:
properties = common_properties.copy()
properties['lilybuild_job_prop'] = job.to_prop()
properties['lilybuild_job_index'] = i
properties['virtual_builder_name'] = 'lilybuild-job - ' + common_properties['lilybuild_repo'] + ' - ' + job.name
properties['lilybuild_dependency_job_indices'] = dep_job_indices
ret.append({
'sched_name': self.lbc.triggerable_scheduler_name,
'props_to_set': properties,
'unimportant': False,
})
return ret
class LatestMixin:
latest_build_dir_pattern = '%(kw:st)s/repos/%(prop:lilybuild_repo_id)s/latest'
latest_good_build_dir_pattern = '%(kw:st)s/repos/%(prop:lilybuild_repo_id)s/latest-good'
latest_name_pattern = '%(kw:st)s/repos/%(prop:lilybuild_repo_id)s/latest/%(kw:name)s'
latest_good_name_pattern = '%(kw:st)s/repos/%(prop:lilybuild_repo_id)s/latest-good/%(kw:name)s'
class EnsureLatestDirs(steps.MasterShellCommand, LatestMixin):
def __init__(self, lbc, storage_dir, **kwargs):
self.lbc = lbc
self.storage_dir = storage_dir
command = util.Transform(
fill_list,
'mkdir',
'-pv',
util.Interpolate(self.latest_build_dir_pattern, st=self.storage_dir),
util.Interpolate(self.latest_good_build_dir_pattern, st=self.storage_dir))
super().__init__(
name='Ensure latest dirs',
command=command,
logEnviron=False,
doStepIf=on_always,
**kwargs
)
class MarkLatest(steps.MasterShellCommand, LatestMixin):
def __init__(self, lbc, storage_dir, buildid, ref_name, is_good, **kwargs):
self.lbc = lbc
self.storage_dir = storage_dir
name_pattern = self.latest_good_name_pattern if is_good else self.latest_name_pattern
command = util.Transform(
fill_list,
'ln',
'-sfvn',
util.Interpolate('../builds/%(kw:buildid)s', buildid=buildid),
util.Interpolate(name_pattern, st=self.storage_dir, name=ref_name),
)
super().__init__(
name='Mark build as latest-good' if is_good else 'Mark build as latest',
command=command,
logEnviron=False,
doStepIf=on_success if is_good else on_always,
**kwargs
)
class AnalyzeCIFileCommand(buildstep.ShellMixin, steps.BuildStep):
ci_def_file = '.gitlab-ci.yml'
build_target_prop_name = 'harbormaster_build_target_phid'
def __init__(
self,
lbc,
src_relative=None,
src_dir=None,
storage_dir=None,
repo_id=None,
result_relative=None,
result_dir=None,
artifact_stage_relative=None,
artifact_stage_dir=None,
**kwargs):
kwargs['name'] = 'Analyze CI file'
kwargs['command'] = ['cat', self.ci_def_file]
self.lbc = lbc
self.src_relative = src_relative
self.src_dir = src_dir
self.work_root_dir = kwargs['workdir']
self.script_dir = 'script'
self.storage_dir = storage_dir
self.repo_id = repo_id
self.artifact_stage_relative = artifact_stage_relative
self.artifact_stage_dir = artifact_stage_dir
self.result_relative = result_relative
self.result_dir = result_dir
kwargs['workdir'] = self.src_dir
kwargs = self.setupShellMixin(kwargs)
super().__init__(**kwargs)
self.observer = logobserver.BufferLogObserver()
self.addLogObserver('stdio', self.observer)
def stage_to_step(self, stage_name, stage_jobs, job_name_to_index_map, ci_file):
jobs_with_data = []
for job in stage_jobs:
dep_job_names = [
jn
for jn in ci_file.get_jobs_to_pull_artifacts_from(job.name)
if ci_file.jobs[jn].has_artifacts_archive()
]
dep_job_indices = [job_name_to_index_map[jn] for jn in dep_job_names]
jobs_with_data.append((job, job_name_to_index_map[job.name], dep_job_indices))
trigger = TriggerMultipleJobsStep(
name=stage_name,
lbc=self.lbc,
jobs_with_data=jobs_with_data,
waitForFinish=True,
doStepIf=on_success,
)
return trigger
def get_steps_and_job_map(self, stdout):
f = ci_file.CIFile(stdout)
stages = f.get_grouped_jobs()
jobs = [job for (stage, js) in f.get_grouped_jobs() for job in js]
job_names = [job.name for job in jobs]
job_name_to_index_map = {}
for (i, j) in enumerate(jobs):
job_name_to_index_map[j.name] = i
steps = [self.stage_to_step(stage_name, stage_jobs, job_name_to_index_map, f) for (stage_name, stage_jobs) in stages]
print('steps:', steps)
return (steps, job_name_to_index_map)
def get_is_phorge(self):
return not not self.getProperty(self.build_target_prop_name)
def get_ref_and_type(self):
ref_type = 'branch'
ref = self.getProperty('branch')
if self.getProperty('category') == 'tag':
ref_type = 'tag'
if ref is not None:
m = re.match(r'^refs/(heads|tags)/(.+)$', ref)
if m:
ref = m.group(2)
return (ref, ref_type)
def get_cur_repo_config(self):
return self.lbc.repos[self.getProperty('lilybuild_repo_id')]
@defer.inlineCallbacks
def get_pipeline_ci_vars(self):
url = yield self.build.getUrl()
res = {
'CI_PIPELINE_ID': self.build.buildid,
'CI_PIPELINE_IID': self.build.buildid,
'CI_PIPELINE_URL': url,
'CI_PROJECT_ID': self.getProperty('lilybuild_repo_id'),
'CI_CONFIG_PATH': self.ci_def_file,
}
if not self.get_is_phorge():
res['CI_COMMIT_SHA'] = self.getProperty('got_revision')
res['CI_COMMIT_SHORT_SHA'] = res['CI_COMMIT_SHA'][:8]
(ref, ref_type) = self.get_ref_and_type()
res['CI_COMMIT_REF_NAME'] = ref
res['CI_COMMIT_REF_SLUG'] = ci_file.ci_slugify(ref)
res['CI_COMMIT_REF_PROTECTED'] = 'false'
if ref_type == 'tag':
res['CI_COMMIT_TAG'] = ref
elif ref_type == 'branch':
res['CI_COMMIT_BRANCH'] = ref
return res
@defer.inlineCallbacks
def run(self):
# run './build.sh --list-stages' to generate the list of stages
cmd = yield self.makeRemoteShellCommand()
yield self.runCommand(cmd)
# if the command passes extract the list of stages
result = cmd.results()
if result == util.SUCCESS:
pipeline_vars = yield self.get_pipeline_ci_vars()
self.setProperty('lilybuild_pipeline_vars', pipeline_vars, self.__class__.__name__)
# create a ShellCommand for each stage and add them to the build
(steps, job_map) = self.get_steps_and_job_map(self.observer.getStdout())
self.setProperty('lilybuild_job_map', job_map, self.__class__.__name__)
self.build.addStepsAfterCurrentStep(steps)
latest_branch_map = self.get_cur_repo_config()['artifact_latest_branch_map']
ref, _ref_type = self.get_ref_and_type()
if ref in latest_branch_map:
latest_name = latest_branch_map[ref]
self.build.addStepsAfterLastStep([
EnsureLatestDirs(lbc=self.lbc, storage_dir=self.storage_dir),
MarkLatest(
lbc=self.lbc,
storage_dir=self.storage_dir,
buildid=self.build.buildid,
ref_name=latest_name,
is_good=True,
),
MarkLatest(
lbc=self.lbc,
storage_dir=self.storage_dir,
buildid=self.build.buildid,
ref_name=latest_name,
is_good=False,
),
])
return result
diff --git a/lilybuild/lilybuild/helpers.py b/lilybuild/lilybuild/helpers.py
index 35e1cae..9c093c6 100644
--- a/lilybuild/lilybuild/helpers.py
+++ b/lilybuild/lilybuild/helpers.py
@@ -1,105 +1,139 @@
import json
import shlex
import re
def normalize_path_for_rsync(path):
n = path
if n.startswith('./'):
n = n[2:]
if n.endswith('/'):
n = n[:-1]
return '/' + n
def rsync_rules_from_artifacts(artifacts):
paths = artifacts.get('paths', [])
# Include all dirs
rules = ['--include', '*/']
for p in paths:
normalized_path = normalize_path_for_rsync(p)
rules += [
# If path already has /** at the end, the second will actually do nothing,
# but it's fine to add it anyway. The directory itself will still
# be visited because of the --include */ option we add at the beginning.
'--include', normalized_path,
'--include', normalized_path + '/**',
]
# Exclude everything else
rules += ['--exclude', '*']
return rules
def normalize_base_url(base_url):
return base_url.rstrip('/') if base_url else None
def phorge_token_to_arcrc(normalized_base_url, token):
return json.dumps({
'hosts': {
normalized_base_url + '/api/': {
'token': token,
},
},
})
def ci_vars_to_env_file(v):
res = []
for name in v:
value = v[name]
if not isinstance(value, str):
value = str(value)
if '\n' not in value:
res.append(f'{name}={value}')
# Otherwise, ignore multiline variables because podman cannot pass it in env file
return '\n'.join(res)
DEFAULT_SCRIPT_HEADER = '''\
#!/bin/sh
set -e -x
cd /build
'''
def get_job_script(job):
return (
DEFAULT_SCRIPT_HEADER +
'\n\n'.join(job.before_script) + '\n\n' +
'\n\n'.join(job.script) +
'\n\nset +e\n\n' +
'\n\n'.join(job.after_script) +
'\n\nexit 0'
)
def normalize_image(image):
if isinstance(image, str):
return json.dumps({'name': image})
else:
return json.dumps(image)
def get_service_aliases_from_name(name):
# https://docs.gitlab.com/ci/services/#accessing-the-services
pos = name.find(':')
if pos != -1:
name = name[:pos]
primary = name.replace('/', '__')
secondary = name.replace('/', '-')
if primary == secondary:
return [primary]
else:
return [primary, secondary]
SERVICE_ALIAS_SEPARATOR = re.compile(r'[ ,]+')
def normalize_services(services):
res = []
for s in services:
so = s if isinstance(s, dict) else {'name': s}
normalized_service = {
'name': so['name'],
'aliases': SERVICE_ALIAS_SEPARATOR.split(so.get('alias')) if so.get('alias') else get_service_aliases_from_name(so['name']),
'entrypoint': so.get('entrypoint'),
'command': so.get('command'),
}
res.append(normalized_service)
return json.dumps(res)
+
+VAR_REGEX = re.compile(r'\$([A-Za-z0-9_]+|\{[A-Za-z0-9_]+\})')
+def expand_in_vars(value, vs):
+ def replacement(match):
+ varname = match.group(1)
+ if varname.startswith('{'):
+ varname = varname[1:-1]
+ return vs.get(varname, '')
+ return VAR_REGEX.sub(replacement, value)
+
+def generate_metadata_from_job(repo_id, job, vs):
+ res = {
+ 'repo_id': repo_id,
+ 'caches': [],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ caches = job.struct_raw.get('cache') or []
+ if not isinstance(caches, list):
+ caches = [caches]
+ for cache_def in caches:
+ cache_key = cache_def.get('key')
+ if isinstance(cache_key, str):
+ cache_key = expand_in_vars(cache_key, vs)
+
+ paths = cache_def.get('paths') or []
+ res['caches'].append({
+ 'key': cache_key,
+ 'paths': paths,
+ 'when': cache_def.get('when') or 'on_success',
+ 'policy': cache_def.get('policy') or 'pull-push',
+ })
+
+ return json.dumps(res)
diff --git a/lilybuild/lilybuild/podman_helper.py b/lilybuild/lilybuild/podman_helper.py
index 5ad73e7..ba12a47 100755
--- a/lilybuild/lilybuild/podman_helper.py
+++ b/lilybuild/lilybuild/podman_helper.py
@@ -1,405 +1,541 @@
#!/usr/bin/env python3
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
+ 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,
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', '--delete',
+ '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',
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',
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 retcode != 0:
+ 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)
diff --git a/lilybuild/lilybuild/tests/ci_syntax/res/cache.yaml b/lilybuild/lilybuild/tests/ci_syntax/res/cache.yaml
new file mode 100644
index 0000000..7d76c7d
--- /dev/null
+++ b/lilybuild/lilybuild/tests/ci_syntax/res/cache.yaml
@@ -0,0 +1,23 @@
+
+a:
+ cache:
+ - key: test
+ paths:
+ - abc
+ when: always
+ policy: pull
+
+b:
+ cache:
+ key: test
+ paths:
+ - abc
+
+c:
+ cache:
+ - key: $CI_JOB_NAME
+ paths:
+ - abc
+ - key: xx$CI_JOB_NAME
+ paths:
+ - def
diff --git a/lilybuild/lilybuild/tests/helpers_test.py b/lilybuild/lilybuild/tests/helpers_test.py
index 6cbba29..b666942 100644
--- a/lilybuild/lilybuild/tests/helpers_test.py
+++ b/lilybuild/lilybuild/tests/helpers_test.py
@@ -1,191 +1,281 @@
import unittest
import json
from lilybuild.ci_syntax.ci_file import CIFile
from lilybuild.helpers import (
rsync_rules_from_artifacts,
normalize_base_url,
phorge_token_to_arcrc,
ci_vars_to_env_file,
get_job_script,
DEFAULT_SCRIPT_HEADER,
normalize_image,
get_service_aliases_from_name,
normalize_services,
+ expand_in_vars,
+ generate_metadata_from_job,
)
from lilybuild.tests.resources import get_res
class RsyncRulesTest(unittest.TestCase):
def test_empty(self):
self.assertEqual(
rsync_rules_from_artifacts({}),
['--include', '*/', '--exclude', '*']
)
def test_simple(self):
self.assertEqual(
rsync_rules_from_artifacts({'paths': ['public']}),
['--include', '*/',
'--include', '/public',
'--include', '/public/**',
'--exclude', '*']
)
def test_dotslash(self):
self.assertEqual(
rsync_rules_from_artifacts({'paths': ['./public/']}),
['--include', '*/',
'--include', '/public',
'--include', '/public/**',
'--exclude', '*']
)
def test_doublestar(self):
self.assertEqual(
rsync_rules_from_artifacts({'paths': ['./public/**']}),
['--include', '*/',
'--include', '/public/**',
'--include', '/public/**/**',
'--exclude', '*']
)
def test_doublestar_middle(self):
self.assertEqual(
rsync_rules_from_artifacts({'paths': ['./public/**/*.html']}),
['--include', '*/',
'--include', '/public/**/*.html',
'--include', '/public/**/*.html/**',
'--exclude', '*']
)
def test_dotdotslash(self):
# No exploit possible because rsync will not visit beyond the source root
self.assertEqual(
rsync_rules_from_artifacts({'paths': ['../etc/passwd']}),
['--include', '*/',
'--include', '/../etc/passwd',
'--include', '/../etc/passwd/**',
'--exclude', '*']
)
class PhorgeUtilsTest(unittest.TestCase):
def test_normalize_base_url(self):
self.assertEqual(normalize_base_url('https://iron.lily-is.land/'), 'https://iron.lily-is.land')
self.assertEqual(normalize_base_url('https://iron.lily-is.land'), 'https://iron.lily-is.land')
self.assertEqual(normalize_base_url(''), None)
def test_phorge_token_to_arcrc(self):
self.assertEqual(
json.loads(phorge_token_to_arcrc('https://iron.lily-is.land', 'some-token')),
{
'hosts': {
'https://iron.lily-is.land/api/': {
'token': 'some-token',
},
},
},
)
class CiVarsTest(unittest.TestCase):
def test_simple(self):
self.assertEqual(ci_vars_to_env_file({}), '')
self.assertEqual(ci_vars_to_env_file({'VAR': 'val'}), 'VAR=val')
self.assertEqual(ci_vars_to_env_file({'VAR': 'foo bar'}), "VAR=foo bar")
self.assertEqual(ci_vars_to_env_file({'VAR': '\nbar', 'MEW': 'abc def'}), "MEW=abc def")
self.assertEqual(ci_vars_to_env_file({'VAR': 12345}), "VAR=12345")
class GetJobScriptTest(unittest.TestCase):
def test_only_script(self):
r = CIFile(get_res('pages_attr'))
job_script = get_job_script(r.jobs['is-pages'])
self.assertEqual(job_script, f'''\
{DEFAULT_SCRIPT_HEADER}
make docs
set +e
exit 0''')
def test_before_and_after(self):
r = CIFile(get_res('defaults'))
job_script = get_job_script(r.jobs['build-a'])
self.assertEqual(job_script, f'''\
{DEFAULT_SCRIPT_HEADER}ls
make
make install
set +e
find
echo ok
exit 0''')
class NormalizeImageTest(unittest.TestCase):
def test_str(self):
res = normalize_image('alpine')
self.assertEqual(json.loads(res), {'name': 'alpine'})
def test_object(self):
orig = {'name': 'alpine', 'entrypoint': ['/docker-run', '/bin/bb']}
res = normalize_image(orig)
self.assertEqual(json.loads(res), orig)
class GetServiceAliasesFromNameTest(unittest.TestCase):
def test_simple(self):
self.assertEqual(
get_service_aliases_from_name('mewmew:abcdefg'),
['mewmew'])
self.assertEqual(
get_service_aliases_from_name('mewmew/a:abcdefg'),
['mewmew__a', 'mewmew-a'])
self.assertEqual(
get_service_aliases_from_name('mew-mew/a:abc-defg'),
['mew-mew__a', 'mew-mew-a'])
self.assertEqual(
get_service_aliases_from_name('a.example/mew-mew/a:abc-defg'),
['a.example__mew-mew__a', 'a.example-mew-mew-a'])
class NormalizeServicesTest(unittest.TestCase):
def test_simple(self):
self.assertEqual(
json.loads(normalize_services([
'mysql:latest',
'mysql:latest',
])),
[{ 'name': 'mysql:latest', 'aliases': ['mysql'], 'entrypoint': None, 'command': None },
{ 'name': 'mysql:latest', 'aliases': ['mysql'], 'entrypoint': None, 'command': None }],
)
def test_own_alias(self):
self.assertEqual(
json.loads(normalize_services([
{'name': 'mysql:latest', 'alias': 'a, b c'},
'mysql:latest',
])),
[{ 'name': 'mysql:latest', 'aliases': ['a', 'b', 'c'], 'entrypoint': None, 'command': None },
{ 'name': 'mysql:latest', 'aliases': ['mysql'], 'entrypoint': None, 'command': None }],
)
def test_entrypoint_command(self):
self.assertEqual(
json.loads(normalize_services([
{'name': 'mysql:latest', 'entrypoint': 'a', 'command': 'b c'},
])),
[{ 'name': 'mysql:latest', 'aliases': ['mysql'], 'entrypoint': 'a', 'command': 'b c' }]
)
self.assertEqual(
json.loads(normalize_services([
{'name': 'mysql:latest', 'entrypoint': ['a', 'b'], 'command': ['b c', 'c d']},
])),
[{ 'name': 'mysql:latest', 'aliases': ['mysql'], 'entrypoint': ['a', 'b'], 'command': ['b c', 'c d'] }]
)
+class ExpandInVarsTest(unittest.TestCase):
+ def test_simple(self):
+ self.assertEqual(
+ expand_in_vars('a', {}),
+ 'a'
+ )
+ self.assertEqual(
+ expand_in_vars('$abc', {'a': '1', 'abc': '2'}),
+ '2'
+ )
+ self.assertEqual(
+ expand_in_vars('$abc_$def-$g', {'abc': '1', 'def': '2'}),
+ '2-'
+ )
+ self.assertEqual(
+ expand_in_vars('$abc$def-$g', {'abc': '$def', 'def': '2'}),
+ '$def2-'
+ )
+ self.assertEqual(
+ expand_in_vars('${abc}def-$g', {'abc': '$def', 'def': '2'}),
+ '$defdef-'
+ )
+ self.assertEqual(
+ expand_in_vars('${ab$defc}', {'abc': '$def', 'def': '2'}),
+ '${ab}'
+ )
+ self.assertEqual(
+ expand_in_vars('${ab${def}c}', {'abc': '$def', 'def': '2'}),
+ '${ab2c}'
+ )
+ self.assertEqual(
+ expand_in_vars('${ab${def}c', {'abc': '$def', 'def': '2'}),
+ '${ab2c'
+ )
+
+class GenerateMetadataFromJobTest(unittest.TestCase):
+ def test_simple(self):
+ r = CIFile(get_res('cache'))
+ self.assertEqual(
+ json.loads(generate_metadata_from_job(2, r.jobs['a'], {'CI_JOB_NAME': 'a'})),
+ {
+ 'repo_id': 2,
+ 'caches': [{
+ 'key': 'test',
+ 'paths': ['abc'],
+ 'when': 'always',
+ 'policy': 'pull',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ )
+
+ self.assertEqual(
+ json.loads(generate_metadata_from_job(2, r.jobs['b'], {'CI_JOB_NAME': 'b'})),
+ {
+ 'repo_id': 2,
+ 'caches': [{
+ 'key': 'test',
+ 'paths': ['abc'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ )
+
+ self.assertEqual(
+ json.loads(generate_metadata_from_job(2, r.jobs['c'], {'CI_JOB_NAME': 'c'})),
+ {
+ 'repo_id': 2,
+ 'caches': [{
+ 'key': 'c',
+ 'paths': ['abc'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }, {
+ 'key': 'xxc',
+ 'paths': ['def'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ )
+
if __name__ == '__main__':
unittest.main()
diff --git a/lilybuild/lilybuild/tests/podman_helper_test_worker.py b/lilybuild/lilybuild/tests/podman_helper_test_worker.py
index 08734ff..f4b29e7 100644
--- a/lilybuild/lilybuild/tests/podman_helper_test_worker.py
+++ b/lilybuild/lilybuild/tests/podman_helper_test_worker.py
@@ -1,273 +1,481 @@
import unittest
from unittest.mock import Mock
import tempfile
import time
import os
import json
import subprocess
from dataclasses import dataclass
from contextlib import contextmanager
from lilybuild.podman_helper import PodmanHelper, image_to_podman_args
+try:
+ import tarfile
+ tarfile.FilterError
+except AttributeError:
+ import backports.tarfile as tarfile
@dataclass
class MockedCompletedProcess:
stdout: str | bytes | None = None
stderr: str | bytes | None = None
returncode: int = 0
def mocked(ph, mock=None):
ph.verbose_run = mock or Mock()
return ph
+def make_cache_file(cache_file_name):
+ os.makedirs(os.path.dirname(cache_file_name), exist_ok=True)
+ with tempfile.TemporaryDirectory() as dir_name:
+ os.makedirs(os.path.join(dir_name, 'a'))
+ with open(os.path.join(dir_name, 'a', 'b'), 'w') as f:
+ print('bbb', file=f)
+ with tarfile.open(cache_file_name, 'w:gz') as f:
+ f.add(os.path.join(dir_name, 'a'), 'a')
+ return cache_file_name
+
class PodmanHelperTest(unittest.TestCase):
+ def test_get_cache_storage_dir(self):
+ ph = PodmanHelper(cache_storage_root_dir='/foo/cache')
+ self.assertTrue(
+ ph.get_cache_storage_dir(1, {
+ 'key': 'foo',
+ }, protected=False)
+ .startswith('/foo/cache/repos/1/unprotected/cache-keys/')
+ )
+ self.assertTrue(
+ ph.get_cache_storage_dir(1, {
+ 'key': 'bar',
+ }, protected=True)
+ .startswith('/foo/cache/repos/1/protected/cache-keys/')
+ )
+
+ def test_get_valid_caches(self):
+ md = {
+ 'repo_id': 1,
+ 'caches': [{
+ 'key': 'bar',
+ 'paths': ['a'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ # No cache directory
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = PodmanHelper(cache_storage_root_dir=dir_name)
+ self.assertEqual(ph.get_valid_caches(md), [])
+
+ # With cache directory, no cache file
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = PodmanHelper(cache_storage_root_dir=dir_name)
+ d1 = ph.get_cache_storage_dir(1, md['caches'][0], protected=False)
+ os.makedirs(d1)
+ self.assertEqual(ph.get_valid_caches(md), [])
+
+ # With cache directory, with cache file
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = PodmanHelper(cache_storage_root_dir=dir_name)
+ d1 = ph.get_cache_storage_dir(1, md['caches'][0], protected=False)
+ os.makedirs(d1)
+ cache_file = os.path.join(d1, 'cur')
+ with open(cache_file, 'w') as f:
+ print('', file=f)
+ self.assertEqual(ph.get_valid_caches(md), [cache_file])
+
+ # With cache directory, with expired cache file
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = PodmanHelper(cache_storage_root_dir=dir_name)
+ d1 = ph.get_cache_storage_dir(1, md['caches'][0], protected=False)
+ os.makedirs(d1)
+ cache_file = os.path.join(d1, 'cur')
+ with open(cache_file, 'w') as f:
+ print('', file=f)
+ md2 = md.copy()
+ md2['cache_last_invalidated_sec'] = time.time() + 1
+ self.assertEqual(ph.get_valid_caches(md2), [])
+
+ # With cache directory and cache file, but policy does not contain pull
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = PodmanHelper(cache_storage_root_dir=dir_name)
+ d1 = ph.get_cache_storage_dir(1, md['caches'][0], protected=False)
+ os.makedirs(d1)
+ cache_file = os.path.join(d1, 'cur')
+ with open(cache_file, 'w') as f:
+ print('', file=f)
+ md2 = md.copy()
+ md2['caches'] = [md['caches'][0].copy()]
+ md2['caches'][0]['policy'] = 'push'
+ self.assertEqual(ph.get_valid_caches(md2), [])
+
+ def test_import_caches(self):
+ md = {
+ 'repo_id': 1,
+ 'caches': [{
+ 'key': 'bar',
+ 'paths': ['a'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }, {
+ 'key': 'mew',
+ 'paths': ['b'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = mocked(
+ PodmanHelper(cache_storage_root_dir=dir_name)
+ )
+ ph.helper_container_alias = 'helper-xxx'
+ with self.assertRaises(RuntimeError) as m:
+ ph.import_caches(md, '')
+
+ with tempfile.TemporaryDirectory() as dir_name:
+ ph = mocked(
+ PodmanHelper(cache_storage_root_dir=dir_name)
+ )
+ ph.helper_container_alias = 'helper-xxx'
+ d = ph.get_cache_storage_dir(md['repo_id'], md['caches'][0], protected=md['protected'])
+ cache_name = os.path.join(d, ph.cur_cache_basename)
+ make_cache_file(cache_name)
+ ph.import_caches(md, ph.work_vol_mount_dir)
+ # rsync -> tar -> rm archive
+ self.assertEqual(ph.verbose_run.call_count, 3)
+ self.assertTrue('rsync' in ph.verbose_run.call_args_list[0].args[0])
+ self.assertTrue('tar' in ph.verbose_run.call_args_list[1].args[0])
+ self.assertTrue('rm' in ph.verbose_run.call_args_list[2].args[0])
+
+ with tempfile.TemporaryDirectory() as dir_name:
+ def handle(run_args, **kwargs):
+ if 'tar' in run_args:
+ raise subprocess.CalledProcessError(returncode=1, cmd=run_args)
+ return MockedCompletedProcess()
+ ph = mocked(
+ PodmanHelper(cache_storage_root_dir=dir_name),
+ Mock(side_effect=handle),
+ )
+ ph.helper_container_alias = 'helper-xxx'
+ d = ph.get_cache_storage_dir(md['repo_id'], md['caches'][0], protected=md['protected'])
+ cache_name = os.path.join(d, ph.cur_cache_basename)
+ make_cache_file(cache_name)
+ ph.import_caches(md, ph.work_vol_mount_dir)
+ # rsync -> tar -> clean up extracted dir -> rm archive
+ self.assertEqual(ph.verbose_run.call_count, 4)
+ self.assertTrue('rsync' in ph.verbose_run.call_args_list[0].args[0])
+ self.assertTrue('tar' in ph.verbose_run.call_args_list[1].args[0])
+ self.assertTrue('rm' in ph.verbose_run.call_args_list[2].args[0])
+ self.assertTrue('rm' in ph.verbose_run.call_args_list[3].args[0])
+
+ def test_save_caches(self):
+ md = {
+ 'repo_id': 1,
+ 'caches': [{
+ 'key': 'bar',
+ 'paths': ['a'],
+ 'when': 'on_success',
+ 'policy': 'pull-push',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+ with tempfile.TemporaryDirectory() as dir_name:
+ cache_root = os.path.join(dir_name, 'cache')
+ result_dir = os.path.join(dir_name, 'result')
+ os.makedirs(os.path.join(result_dir, 'a'))
+ with open(os.path.join(result_dir, 'a', 'b'), 'w') as f:
+ print('mewmew', file=f)
+ ph = PodmanHelper(cache_storage_root_dir=cache_root, result_dir=result_dir)
+ ph.save_caches(md, succeeded=True)
+ cache_dir = ph.get_cache_storage_dir(md['repo_id'], md['caches'][0], protected=md['protected'])
+ cache_file = os.path.join(cache_dir, ph.cur_cache_basename)
+ self.assertTrue(os.path.exists(cache_file))
+ with tarfile.open(cache_file) as f:
+ f.getmember('a/b')
+
+ # did not succeed
+ with tempfile.TemporaryDirectory() as dir_name:
+ cache_root = os.path.join(dir_name, 'cache')
+ result_dir = os.path.join(dir_name, 'result')
+ ph = PodmanHelper(cache_storage_root_dir=cache_root, result_dir=result_dir)
+ ph.save_caches(md, succeeded=False)
+ cache_dir = ph.get_cache_storage_dir(md['repo_id'], md['caches'][0], protected=md['protected'])
+ cache_file = os.path.join(cache_dir, ph.cur_cache_basename)
+ self.assertFalse(os.path.exists(cache_file))
+
+ md2 = {
+ 'repo_id': 1,
+ 'caches': [{
+ 'key': 'bar',
+ 'paths': ['a'],
+ 'when': 'on_success',
+ 'policy': 'pull',
+ }],
+ 'cache_last_invalidated_sec': 0,
+ 'protected': False,
+ }
+
+ # policy does not contain push
+ with tempfile.TemporaryDirectory() as dir_name:
+ cache_root = os.path.join(dir_name, 'cache')
+ result_dir = os.path.join(dir_name, 'result')
+ ph = PodmanHelper(cache_storage_root_dir=cache_root, result_dir=result_dir)
+ ph.save_caches(md2, succeeded=True)
+ cache_dir = ph.get_cache_storage_dir(md2['repo_id'], md2['caches'][0], protected=md2['protected'])
+ cache_file = os.path.join(cache_dir, ph.cur_cache_basename)
+ self.assertFalse(os.path.exists(cache_file))
+
def test_create_and_clean_volume(self):
ph = mocked(
PodmanHelper(),
Mock(side_effect=[
MockedCompletedProcess(stdout='volume1\n'),
MockedCompletedProcess(stdout='volume2\n'),
MockedCompletedProcess(stdout='volume1\nvolume2\n'),
])
)
res = ph.create_volume('work')
self.assertEqual(res, 'volume1')
res = ph.create_volume('script')
self.assertEqual(res, 'volume2')
self.assertEqual(ph.volumes_to_remove, ['volume1', 'volume2'])
ph.clean_volumes()
self.assertEqual(
ph.verbose_run.call_args.args[0][-3:],
['--', 'volume1', 'volume2'],
)
def test_helper_service(self):
key_content = 'ssh-ed25519 somethingsomething a@example.com'
with tempfile.NamedTemporaryFile(mode='w+', encoding='utf-8') as key_pub:
print(key_content, file=key_pub)
key_pub.flush()
ph = mocked(
PodmanHelper(
worker_container_name='workerhostname',
key_file_pub=key_pub.name,
ssh_wait_interval_sec=0.001,
volume_helper_image='lilybuild-volume-helper',
),
Mock(side_effect=[
# podman inspect
MockedCompletedProcess(stdout=json.dumps([{
'Pod': 'pod0',
'NetworkSettings': {
'Networks': {
'network0': {
},
}
}
}])),
# podman run
MockedCompletedProcess(stdout='container0\n'),
# nc
MockedCompletedProcess(returncode=1),
# nc
MockedCompletedProcess(stdout=b'SSH 1.1.1\n'),
# podman container rm
MockedCompletedProcess(),
])
)
(cont_name, alias) = ph.start_helper_service('volume1', 'volume2')
self.assertEqual(
ph.verbose_run.call_args_list[0].args[0][-3:],
['inspect', '--', 'workerhostname'],
)
self.assertEqual(
ph.verbose_run.call_args_list[1].args[0][-2:],
['--', 'lilybuild-volume-helper'],
)
self.assertTrue(
'run' in ph.verbose_run.call_args_list[1].args[0]
)
self.assertTrue(ph.helper_container_id is not None)
ph.clean_helper_container()
self.assertEqual(
ph.verbose_run.call_args.args[0][-2:],
['--', cont_name],
)
def test_create_prune_service_containers(self):
env = '/path/to/script/env'
ph = mocked(
PodmanHelper(env_filename=env),
Mock(side_effect=
# network create
[MockedCompletedProcess(stdout='network0')]
# container run
+ [MockedCompletedProcess(stdout=f'container{i}') for i in range(4)]
# stop & rm
+ [MockedCompletedProcess(), MockedCompletedProcess()]
),
)
ph.create_service_network()
self.assertEqual(ph.service_network_id, 'network0')
# no entrypoint, no command
ph.start_and_record_service_container({
'name': 'service:latest',
'aliases': ['foo', 'bar'],
'command': None,
'entrypoint': None,
})
self.assertTrue(f'--env-file={env}' in ph.verbose_run.call_args.args[0])
self.assertTrue('--network=network0' in ph.verbose_run.call_args.args[0])
self.assertTrue('--network-alias=foo' in ph.verbose_run.call_args.args[0])
self.assertTrue('--network-alias=bar' in ph.verbose_run.call_args.args[0])
self.assertEqual(
ph.verbose_run.call_args.args[0][-2:],
['--', 'service:latest'],
)
# no entrypoint, with command
ph.start_and_record_service_container({
'name': 'service:latest',
'aliases': ['foo', 'bar'],
'command': ['abc', 'def'],
'entrypoint': None,
})
self.assertEqual(
ph.verbose_run.call_args.args[0][-4:],
['--', 'service:latest', 'abc', 'def'],
)
# entrypoint str
ph.start_and_record_service_container({
'name': 'service:latest',
'aliases': ['foo', 'bar'],
'command': ['abc', 'def'],
'entrypoint': '/bin/sh',
})
self.assertTrue('--entrypoint=/bin/sh', ph.verbose_run.call_args.args[0])
# entrypoint list
ph.start_and_record_service_container({
'name': 'service:latest',
'aliases': ['foo', 'bar'],
'command': ['abc', 'def'],
'entrypoint': ['/bin/sh', '-c'],
})
self.assertTrue('--entrypoint=["/bin/sh", "-c"]', ph.verbose_run.call_args.args[0])
self.assertEqual(ph.service_containers, ['container0', 'container1', 'container2', 'container3'])
ph.maybe_prune_service_containers()
self.assertEqual(
ph.verbose_run.call_args_list[-2].args[0][-5:],
['--', 'container0', 'container1', 'container2', 'container3'],
)
self.assertEqual(
ph.verbose_run.call_args_list[-1].args[0][-5:],
['--', 'container0', 'container1', 'container2', 'container3'],
)
def test_ensure_service_containers_up(self):
count = 0
def inspect_func(run_args, **kwargs):
nonlocal count
count += 1
cont_id = run_args[-1]
res = [{
'State': {
'Status': 'running',
}
}]
if cont_id == 'container0' and count < 5:
res[0]['State']['Status'] = 'starting'
return MockedCompletedProcess(stdout=json.dumps(res))
ph = mocked(
PodmanHelper(service_wait_interval_sec=0.001, service_max_wait_sec=1),
Mock(side_effect=inspect_func),
)
ph.service_containers += ['container0', 'container1']
ph.ensure_service_containers_up()
def test_ensure_service_containers_up_failed(self):
def inspect_func(run_args, **kwargs):
res = [{
'State': {
'Status': 'starting',
}
}]
return MockedCompletedProcess(stdout=json.dumps(res))
ph = mocked(
PodmanHelper(service_wait_interval_sec=0.5, service_max_wait_sec=1),
Mock(side_effect=inspect_func),
)
ph.service_containers += ['container0', 'container1']
with self.assertRaises(TimeoutError) as m:
ph.ensure_service_containers_up()
def test_image_to_podman_args(self):
self.assertEqual(
image_to_podman_args({'name': 'aaa'}),
['--', 'aaa'],
)
self.assertEqual(
image_to_podman_args({'name': 'aaa', 'entrypoint': ['mew', 'abc']}),
['--entrypoint', json.dumps(['mew', 'abc']), '--', 'aaa'],
)
def test_run_in_container(self):
def make_handle(rc):
def handle(run_args, **kwargs):
if run_args[1] == 'run':
return MockedCompletedProcess(stdout='container0\n')
if run_args[1] == 'logs':
return MockedCompletedProcess()
if run_args[2] == 'inspect' and '{{.State.Status}}' in run_args:
return MockedCompletedProcess(stdout='exited\n')
if run_args[2] == 'inspect' and '{{.State.ExitCode}}' in run_args:
return MockedCompletedProcess(stdout=f'{rc}\n')
if run_args[2] == 'stop' or run_args[2] == 'rm':
return MockedCompletedProcess()
raise RuntimeError(f'Unexpected command called: {run_args}')
return handle
ph = mocked(
PodmanHelper(container_run_timeout_sec=10, env_filename='/env'),
Mock(side_effect=make_handle(0)),
)
rc = ph.run_in_container({'name': 'foo:latest'}, 'work_vol', 'script_vol')
self.assertEqual(rc, 0)
ph = mocked(
PodmanHelper(container_run_timeout_sec=10, env_filename='/env'),
Mock(side_effect=make_handle(1)),
)
rc = ph.run_in_container({'name': 'foo:latest'}, 'work_vol', 'script_vol')
self.assertEqual(rc, 1)
def test_run_in_container_timeout(self):
def handle(run_args, **kwargs):
if run_args[1] == 'run':
return MockedCompletedProcess(stdout='container0\n')
if run_args[1] == 'logs':
raise subprocess.TimeoutExpired(run_args, 10)
if run_args[2] == 'stop' or run_args[2] == 'rm':
return MockedCompletedProcess()
raise RuntimeError(f'Unexpected command called: {run_args}')
ph = mocked(
PodmanHelper(container_run_timeout_sec=10, env_filename='/env'),
Mock(side_effect=handle),
)
rc = ph.run_in_container({'name': 'foo:latest'}, 'work_vol', 'script_vol')
self.assertEqual(rc, 1)
# check the container is cleaned up
self.assertTrue('stop' in ph.verbose_run.call_args_list[-2].args[0])
self.assertTrue('rm' in ph.verbose_run.call_args_list[-1].args[0])
def test_cleanup_all(self):
ph = mocked(PodmanHelper())
ph.cleanup_all()
if __name__ == '__main__':
unittest.main()

File Metadata

Mime Type
text/x-diff
Expires
Mon, Oct 12, 1:07 PM (1 d, 15 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1786225
Default Alt Text
(78 KB)

Event Timeline