Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85806051
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
78 KB
Referenced Files
None
Subscribers
None
View Options
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
Details
Attached
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)
Attached To
Mode
rB lilybuild
Attached
Detach File
Event Timeline
Log In to Comment