Implement scoped Git Argo transport and gap-sensitive health observer
Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e241-8285-7a63-8c0c-51c9cb824dc3
This commit is contained in:
parent
22935955cc
commit
0ae5e84e7b
6 changed files with 512 additions and 11 deletions
|
|
@ -14,7 +14,7 @@ COPY .forgejo/workflows/image.yaml ./.forgejo/workflows/image.yaml
|
|||
COPY schemas/ ./schemas/
|
||||
COPY k8s/ ./k8s/
|
||||
COPY scripts/render_gitops.py scripts/check_gitops_promotion.py ./scripts/
|
||||
RUN uv sync --frozen --extra dev && uv run --no-sync python scripts/render_gitops.py --check --verify-frontend --verify-platform && uv run --no-sync pytest -q -p no:cacheprovider tests/test_frontend_patterns.py tests/test_instruction_evaluation.py tests/test_admin_sync_api.py tests/test_gitops_release.py tests/test_release_broker.py
|
||||
RUN uv sync --frozen --extra dev && uv run --no-sync python scripts/render_gitops.py --check --verify-frontend --verify-platform && uv run --no-sync pytest -q -p no:cacheprovider tests/test_frontend_patterns.py tests/test_instruction_evaluation.py tests/test_admin_sync_api.py tests/test_gitops_release.py tests/test_release_broker.py tests/test_release_transport.py
|
||||
|
||||
# Stage 2 — runtime image
|
||||
FROM python:3.12-slim@sha256:78387bc3881b8273120a12ebe6c1ab22b018ccc2c9adf565ae1ac9b536e184ea AS runtime
|
||||
|
|
|
|||
|
|
@ -42,19 +42,18 @@ Failed rollback retains the active slot. Retried publication/sync must be
|
|||
idempotent, using compare-and-swap and accepting already-at-target as success.
|
||||
The adapter must refuse any unexpected third revision rather than overwriting it.
|
||||
|
||||
Tests use generated fixture keys and an in-memory adapter. They prove policy and
|
||||
state-machine behavior, including restart and failed-health rollback; they are
|
||||
**not** proof of production Git/ArgoCD rollback or admitted authority.
|
||||
Tests use generated fixture keys, disposable real Git repositories and a simulated
|
||||
Kubernetes transport. They prove signed admission, real Git publication/rollback,
|
||||
concurrent-writer rejection and lost-push-response recovery. Kubernetes authorization,
|
||||
production credentials and live Git/ArgoCD rollback are **not** proven by these fixtures.
|
||||
|
||||
## Still required before activation
|
||||
|
||||
1. Implement and verify an authenticated transport adapter that resolves exact
|
||||
source commits, compares rendered manifests to the signed hashes, publishes
|
||||
only the platform child revision field, and runs the fixed selective syncs.
|
||||
Every network operation needs a bounded timeout; health must wait within that
|
||||
bound for the exact revision and verify deployments, report sink and schedules.
|
||||
It must persist/recover the platform commit across lost responses and serialize
|
||||
with other writers. A generic repository-write token is not path enforcement.
|
||||
1. Admit and exercise the implemented transport adapter against an isolated
|
||||
authenticated Git/Kubernetes environment. Supply the actual bounded report/schedule
|
||||
invariant probe, pinned cluster UID and immutable credential/context configuration.
|
||||
Verify the real authorization denials and restart behavior; simulated Kubernetes
|
||||
responses do not prove those. A generic repository-write token is not path enforcement.
|
||||
2. Admit the dedicated principal and credential custody through the owner lane.
|
||||
ArgoCD Core has no API-server token lane. Kubernetes Application patch RBAC
|
||||
alone cannot restrict fields: keep it behind the reviewed broker boundary.
|
||||
|
|
@ -70,3 +69,40 @@ state-machine behavior, including restart and failed-health rollback; they are
|
|||
|
||||
Platform enforcement contract: `railiance-platform/docs/activity-core-release-admission.md`.
|
||||
These requirements remain live work in ACTIVITY-WP-0041-T03 and RPF-WP-0048-T02.
|
||||
|
||||
|
||||
## Transport and observer implementation
|
||||
|
||||
`release_transport.GitArgoBackend` retrieves both full source commits from the fixed
|
||||
activity-core repository and compares parsed manifest hashes with the admitted
|
||||
binding. It accepts only the audited static `runtime.yaml` Kustomization; proposed
|
||||
source cannot execute a renderer/plugin. Git publication creates a commit changing
|
||||
only the platform child revision field. Non-force pushes reject concurrent branch
|
||||
advances. After a lost response, remote Git is the durable source of truth: a retry
|
||||
accepts already-at-target and selective sync uses the current verified platform tip.
|
||||
Unexpected third revisions and unrelated changed paths stop the operation.
|
||||
|
||||
The adapter verifies the operator-pinned kube-system namespace UID before Git
|
||||
publication or Kubernetes access. Context/credentials must remain immutable during
|
||||
its lifetime. Every subprocess is non-shell, with timeouts and sanitized failures.
|
||||
Kubernetes JSON patches compare resourceVersion and spec before setting the fixed
|
||||
sync operation. Existing foreign operations are never overwritten. Root sync selects
|
||||
only activity-core; child sync never prunes. Health checks require exact revision,
|
||||
Synced/Healthy and all three deployments ready at their observed generations, plus
|
||||
a mandatory report/schedule invariant probe. That trusted callback must itself use
|
||||
bounded I/O; no production probe is supplied or implied by the fixture implementation.
|
||||
Polling windows are bounded, with individual command timeout overhead possible.
|
||||
|
||||
`release_observer.HealthObserver` records samples durably. It requires consecutive
|
||||
healthy samples of one revision with gaps no greater than 90 seconds. Failed samples,
|
||||
revision changes and missed samples reset the interval; clock regression durably
|
||||
invalidates it. Attestation requires a complete 24-hour sampled interval and a fresh
|
||||
last sample, including after restart. This is evidence at a bounded sampling cadence,
|
||||
not a claim that every instant between samples was observed. Two distant healthy
|
||||
snapshots cannot establish the interval. New observers start from their first actual
|
||||
sample; they do not backfill the earlier deployment timestamp.
|
||||
|
||||
Neither module is connected to a production schedule, endpoint or credential source.
|
||||
Trusted signer custody, real invariant probes, Temporal dispatch and isolated
|
||||
Kubernetes authorization/rollback tests remain required before activation. Production
|
||||
revision and the deployment observation clock are unchanged by this source-only work.
|
||||
|
|
|
|||
69
src/activity_core/release_observer.py
Normal file
69
src/activity_core/release_observer.py
Normal file
|
|
@ -0,0 +1,69 @@
|
|||
"""Persistent, gap-sensitive health observations; no scheduler or signing key here."""
|
||||
from datetime import datetime, timedelta, timezone
|
||||
import sqlite3
|
||||
|
||||
from .release_operations import revision
|
||||
|
||||
|
||||
class HealthObserver:
|
||||
def __init__(self, database, *, max_gap_seconds=90):
|
||||
if not 1 <= max_gap_seconds <= 90:
|
||||
raise ValueError('health gap must be at most 90 seconds')
|
||||
self.database, self.max_gap = database, timedelta(seconds=max_gap_seconds)
|
||||
with self.connect() as db:
|
||||
db.execute('CREATE TABLE IF NOT EXISTS observations (at TEXT PRIMARY KEY, revision TEXT NOT NULL, healthy INTEGER NOT NULL, since TEXT)')
|
||||
db.execute('CREATE TABLE IF NOT EXISTS observer_guard (singleton INTEGER PRIMARY KEY, invalidated INTEGER NOT NULL)')
|
||||
db.execute('INSERT OR IGNORE INTO observer_guard VALUES (1,0)')
|
||||
|
||||
def connect(self):
|
||||
db=sqlite3.connect(self.database,timeout=5,isolation_level=None)
|
||||
db.execute('PRAGMA synchronous=FULL')
|
||||
return db
|
||||
|
||||
def observe(self, commit, healthy, now=None):
|
||||
"""Called by a trusted sampler after exact-revision readiness/invariant checks.
|
||||
|
||||
Failed/unavailable probes must be recorded False. A missed sample, clock
|
||||
regression or revision change resets the interval; no timestamp backfill.
|
||||
"""
|
||||
revision(commit)
|
||||
if type(healthy) is not bool:raise ValueError('explicit boolean health required')
|
||||
now=now or datetime.now(timezone.utc)
|
||||
if now.tzinfo is None:raise ValueError('timezone required')
|
||||
db=self.connect()
|
||||
try:
|
||||
db.execute('BEGIN IMMEDIATE')
|
||||
prior=db.execute('SELECT at,revision,healthy,since FROM observations ORDER BY rowid DESC LIMIT 1').fetchone()
|
||||
invalidated=db.execute('SELECT invalidated FROM observer_guard WHERE singleton=1').fetchone()[0]
|
||||
since=now.isoformat() if healthy else None
|
||||
if prior:
|
||||
last=datetime.fromisoformat(prior[0])
|
||||
if now<=last:
|
||||
db.execute('UPDATE observer_guard SET invalidated=1 WHERE singleton=1')
|
||||
db.commit()
|
||||
raise ValueError('non-increasing observation time')
|
||||
if not invalidated and healthy and prior[2] and prior[1]==commit and now-last<=self.max_gap:
|
||||
since=prior[3]
|
||||
db.execute('INSERT INTO observations VALUES (?,?,?,?)',(now.isoformat(),commit,int(healthy),since))
|
||||
db.execute('UPDATE observer_guard SET invalidated=0 WHERE singleton=1')
|
||||
db.commit()
|
||||
return since
|
||||
except Exception:
|
||||
db.rollback();raise
|
||||
finally:db.close()
|
||||
|
||||
def attestation(self, commit, now=None):
|
||||
"""Return facts for a trusted health issuer, never fabricate a signature."""
|
||||
revision(commit);now=now or datetime.now(timezone.utc)
|
||||
db=self.connect()
|
||||
try:
|
||||
row=db.execute('SELECT at,revision,healthy,since FROM observations ORDER BY rowid DESC LIMIT 1').fetchone()
|
||||
if db.execute('SELECT invalidated FROM observer_guard WHERE singleton=1').fetchone()[0]:
|
||||
raise ValueError('observation interval invalidated by clock regression')
|
||||
finally:db.close()
|
||||
if not row or row[1]!=commit or not row[2] or not row[3]:raise ValueError('no healthy revision interval')
|
||||
measured=datetime.fromisoformat(row[0]);since=datetime.fromisoformat(row[3])
|
||||
if not timedelta(0)<=now-measured<=self.max_gap:raise ValueError('health observations stale or clock regressed')
|
||||
if measured-since<timedelta(hours=24):raise ValueError('24-hour observed interval incomplete')
|
||||
return {'revision':commit,'continuous':True,'healthy':True,'synced':True,
|
||||
'healthy_since':since.isoformat(),'observed_at':measured.isoformat()}
|
||||
183
src/activity_core/release_transport.py
Normal file
183
src/activity_core/release_transport.py
Normal file
|
|
@ -0,0 +1,183 @@
|
|||
"""Fixed Git and Kubernetes transport for an admitted release broker.
|
||||
|
||||
Credentials and context are operator-provided process configuration. No bearer
|
||||
values, URLs, arbitrary commands or repository paths are accepted from a release.
|
||||
There is no service endpoint or automatic production activation in this module.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import subprocess
|
||||
import time
|
||||
|
||||
import yaml
|
||||
|
||||
from .release_broker import fingerprint
|
||||
from .release_operations import APPLICATION_PATH, revision, sync_operations, update_child
|
||||
|
||||
PLATFORM_URL='https://forgejo.coulomb.social/coulomb/railiance-platform.git'
|
||||
ACTIVITY_URL='https://forgejo.coulomb.social/coulomb/activity-core.git'
|
||||
|
||||
|
||||
class TransportError(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
def run(argv, *, input=None, env=None, timeout=30):
|
||||
"""Bounded, non-shell invocation; suppress outputs that could contain credentials."""
|
||||
try:
|
||||
return subprocess.run(argv,input=input,text=True,capture_output=True,check=True,
|
||||
timeout=timeout,env=env).stdout
|
||||
except (subprocess.SubprocessError,OSError):
|
||||
raise TransportError('release transport command failed or timed out') from None
|
||||
|
||||
|
||||
class GitArgoBackend:
|
||||
def __init__(self, plan, workspace, invariants, *, cluster_uid, admitted=False, runner=run,
|
||||
clock=time.monotonic, pause=time.sleep, wait_seconds=90):
|
||||
if not admitted:raise ValueError('transport identity not admitted')
|
||||
if not callable(invariants):raise ValueError('report and schedule invariant probe required')
|
||||
if not 1<=wait_seconds<=90:raise ValueError('bounded wait required')
|
||||
if not isinstance(cluster_uid,str) or not cluster_uid:raise ValueError('pinned cluster UID required')
|
||||
self.cluster_uid=cluster_uid;self.cluster_verified=False
|
||||
self.plan=dict(plan);self.workspace=Path(workspace);self.invariants=invariants
|
||||
self.runner,self.clock,self.pause,self.wait_seconds=runner,clock,pause,wait_seconds
|
||||
if plan.get('repository')!='coulomb/activity-core' or plan.get('application')!='activity-core':
|
||||
raise ValueError('transport plan outside scope')
|
||||
for name in ('candidate','rollback'):revision(plan[name])
|
||||
self.workspace.mkdir(parents=True,exist_ok=True,mode=0o700)
|
||||
# Environment is trusted service configuration, not producer input. In
|
||||
# particular kubeconfig and Git's credential helper belong to this identity.
|
||||
self.env={**os.environ,'GIT_TERMINAL_PROMPT':'0','GIT_AUTHOR_NAME':'activity-core release broker',
|
||||
'GIT_AUTHOR_EMAIL':'release-broker@activity-core.invalid',
|
||||
'GIT_COMMITTER_NAME':'activity-core release broker',
|
||||
'GIT_COMMITTER_EMAIL':'release-broker@activity-core.invalid'}
|
||||
self.platform=self.workspace/'platform.git';self.activity=self.workspace/'activity.git'
|
||||
for path in (self.platform,self.activity):
|
||||
if not path.exists():self.runner(['git','init','--bare',str(path)],env=self.env,timeout=30)
|
||||
|
||||
def git(self, repo, *args, input=None):
|
||||
result=self.runner(['git','-c','core.hooksPath=/dev/null','-c','protocol.file.allow=never',
|
||||
'-C',str(repo),*args],input=input,env=self.env,timeout=30)
|
||||
return result if args[0]=='show' else result.strip()
|
||||
|
||||
def head(self):
|
||||
self.git(self.platform,'fetch','--no-tags',PLATFORM_URL,'main')
|
||||
return revision(self.git(self.platform,'rev-parse','FETCH_HEAD'))
|
||||
|
||||
def validate_sources(self):
|
||||
for key,hash_key in [('rollback','before_sha256'),('candidate','after_sha256')]:
|
||||
commit=self.plan[key]
|
||||
self.git(self.activity,'fetch','--no-tags',ACTIVITY_URL,commit)
|
||||
body=self.git(self.activity,'show',commit+':k8s/gitops/runtime.yaml')
|
||||
docs=list(yaml.safe_load_all(body))
|
||||
if fingerprint(docs)!=self.plan[hash_key]:raise ValueError('source commit manifest binding mismatch')
|
||||
# Never execute kustomize/plugins from a proposed source. Only the
|
||||
# audited single static resource file is accepted by this lane.
|
||||
kustomization=yaml.safe_load(self.git(self.activity,'show',commit+':k8s/gitops/kustomization.yaml'))
|
||||
if kustomization!={'apiVersion':'kustomize.config.k8s.io/v1beta1','kind':'Kustomization','resources':['runtime.yaml']}:
|
||||
raise ValueError('source rendering outside static projection contract')
|
||||
|
||||
def publish_revision(self, expected, target):
|
||||
self.verify_cluster()
|
||||
if {expected,target}!={self.plan['candidate'],self.plan['rollback']}:
|
||||
raise ValueError('revision outside admitted release')
|
||||
self.validate_sources()
|
||||
base=self.head()
|
||||
document=self.git(self.platform,'show',base+':'+APPLICATION_PATH)
|
||||
updated=update_child(document,expected,target)
|
||||
if updated==document:return # Recovery after a lost push response.
|
||||
blob=self.git(self.platform,'hash-object','-w','--stdin',input=updated)
|
||||
# The trusted per-broker workspace is serialized by Broker's ledger lock.
|
||||
self.git(self.platform,'read-tree',base)
|
||||
self.git(self.platform,'update-index','--add','--cacheinfo','100644',blob,APPLICATION_PATH)
|
||||
tree=self.git(self.platform,'write-tree')
|
||||
commit=self.git(self.platform,'commit-tree',tree,'-p',base,input='Promote activity-core to '+target+'\n')
|
||||
changed=self.git(self.platform,'diff-tree','--no-commit-id','--name-only','-r',commit)
|
||||
if changed!=APPLICATION_PATH:raise ValueError('publication contains unrelated paths')
|
||||
# Non-force push rejects concurrent branch advances. A retry fetches the
|
||||
# new tip and repeats the child-revision compare-and-swap.
|
||||
self.git(self.platform,'push',PLATFORM_URL,revision(commit)+':refs/heads/main')
|
||||
|
||||
def verify_cluster(self):
|
||||
if not self.cluster_verified:
|
||||
uid=self.runner(['kubectl','--request-timeout=20s','get','namespace','kube-system',
|
||||
'-o','jsonpath={.metadata.uid}'],env=self.env,timeout=25).strip()
|
||||
if uid!=self.cluster_uid:raise ValueError('cluster identity mismatch')
|
||||
self.cluster_verified=True
|
||||
|
||||
def kube(self, *args, input=None):
|
||||
self.verify_cluster()
|
||||
return self.runner(['kubectl','--request-timeout=20s',*args],input=input,env=self.env,timeout=25)
|
||||
|
||||
def app(self, name):
|
||||
if name not in {'railiance-apps-root','activity-core'}:raise ValueError('foreign application')
|
||||
obj=json.loads(self.kube('-n','argocd','get','application',name,'-o','json'))
|
||||
if obj.get('metadata',{}).get('name')!=name or obj['metadata'].get('namespace')!='argocd':
|
||||
raise ValueError('application identity mismatch')
|
||||
if name=='railiance-apps-root':
|
||||
source=obj.get('spec',{}).get('source',{})
|
||||
if (source.get('repoURL')!=PLATFORM_URL or source.get('path')!='argocd/railiance01/applications'
|
||||
or source.get('targetRevision')!='main' or obj['spec'].get('syncPolicy',{}).get('automated') is not None):
|
||||
raise ValueError('root outside release contract')
|
||||
else:
|
||||
current=obj.get('spec',{}).get('source',{}).get('targetRevision')
|
||||
# Reuse validation without changing the live object.
|
||||
update_child(yaml.safe_dump(obj),revision(current),current)
|
||||
return obj
|
||||
|
||||
def _sync(self, name, body):
|
||||
desired=body['operation']['sync'];obj=self.app(name)
|
||||
operation=obj.get('operation')
|
||||
if operation is not None:
|
||||
active=dict(operation.get('sync',{}));active.setdefault('prune',False)
|
||||
if active!=desired:raise ValueError('another ArgoCD operation is active')
|
||||
done=obj.get('status',{}).get('operationState',{})
|
||||
completed=(done.get('phase')=='Succeeded' and done.get('syncResult',{}).get('revision')==desired['revision'])
|
||||
if operation is None and not completed:
|
||||
patch=[{'op':'test','path':'/metadata/resourceVersion','value':obj['metadata']['resourceVersion']},
|
||||
{'op':'test','path':'/spec','value':obj['spec']},
|
||||
{'op':'add','path':'/operation','value':body['operation']}]
|
||||
self.kube('-n','argocd','patch','application',name,'--type=json','--patch-file=/dev/stdin',input=json.dumps(patch))
|
||||
until=self.clock()+self.wait_seconds
|
||||
while self.clock()<until:
|
||||
state=self.app(name).get('status',{}).get('operationState',{})
|
||||
result=state.get('syncResult',{}).get('revision')
|
||||
if result==desired['revision']:
|
||||
if state.get('phase')=='Succeeded':return
|
||||
if state.get('phase') in {'Failed','Error'}:raise TransportError('ArgoCD sync failed')
|
||||
self.pause(1)
|
||||
raise TransportError('ArgoCD sync wait timed out')
|
||||
|
||||
def sync_revision(self, target):
|
||||
if target not in {self.plan['candidate'],self.plan['rollback']}:raise ValueError('unbound revision')
|
||||
base=self.head()
|
||||
document=self.git(self.platform,'show',base+':'+APPLICATION_PATH)
|
||||
# Publication must remain the current source tip; never sync a stale pin.
|
||||
update_child(document,target,target)
|
||||
root,child=sync_operations(base,target)
|
||||
self._sync(*root)
|
||||
if self.app('activity-core')['spec']['source']['targetRevision']!=target:
|
||||
raise ValueError('selective root sync did not set expected child revision')
|
||||
self._sync(*child)
|
||||
|
||||
def healthy(self, target):
|
||||
if target not in {self.plan['candidate'],self.plan['rollback']}:raise ValueError('unbound revision')
|
||||
until=self.clock()+self.wait_seconds
|
||||
while self.clock()<until:
|
||||
obj=self.app('activity-core');status=obj.get('status',{})
|
||||
synced=(obj['spec']['source']['targetRevision']==target
|
||||
and status.get('sync',{}).get('revision')==target
|
||||
and status.get('sync',{}).get('status')=='Synced'
|
||||
and status.get('health',{}).get('status')=='Healthy')
|
||||
if synced:
|
||||
ready=True
|
||||
for name in ('actcore-api','actcore-worker','actcore-event-router'):
|
||||
d=json.loads(self.kube('-n','activity-core','get','deployment',name,'-o','json'))
|
||||
n=d['spec'].get('replicas',1);s=d.get('status',{})
|
||||
ready &= n>0 and s.get('observedGeneration')==d['metadata']['generation'] and all(s.get(k,0)==n for k in ('readyReplicas','updatedReplicas','availableReplicas'))
|
||||
if ready and self.invariants(target) is True:return True
|
||||
self.pause(1)
|
||||
return False
|
||||
191
tests/test_release_transport.py
Normal file
191
tests/test_release_transport.py
Normal file
|
|
@ -0,0 +1,191 @@
|
|||
"""Real disposable Git repositories; simulated Kubernetes boundary."""
|
||||
import copy
|
||||
import json
|
||||
import subprocess
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
import yaml
|
||||
|
||||
from activity_core.release_broker import fingerprint
|
||||
from activity_core.release_observer import HealthObserver
|
||||
from activity_core.release_operations import APPLICATION_PATH
|
||||
from activity_core.release_transport import ACTIVITY_URL, PLATFORM_URL, GitArgoBackend, TransportError, run
|
||||
NOW=datetime(2026,9,28,16,tzinfo=timezone.utc)
|
||||
|
||||
ROOT=Path(__file__).resolve().parents[1]
|
||||
|
||||
|
||||
def git(path,*args):
|
||||
return subprocess.check_output(['git','-C',str(path),*args],text=True,stderr=subprocess.DEVNULL).strip()
|
||||
|
||||
|
||||
def init(path):
|
||||
path.mkdir();git(path,'init','-b','main');git(path,'config','user.name','Fixture');git(path,'config','user.email','fixture@example.invalid')
|
||||
|
||||
|
||||
def commit(path):
|
||||
git(path,'add','.');git(path,'commit','-m','fixture');return git(path,'rev-parse','HEAD')
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def transport(tmp_path):
|
||||
activity=tmp_path/'activity';init(activity)
|
||||
k=activity/'k8s/gitops';k.mkdir(parents=True)
|
||||
before=list(yaml.safe_load_all((ROOT/'k8s/gitops/runtime.yaml').read_text()))
|
||||
(k/'runtime.yaml').write_text(yaml.safe_dump_all(before));(k/'kustomization.yaml').write_text((ROOT/'k8s/gitops/kustomization.yaml').read_text())
|
||||
rollback=commit(activity)
|
||||
after=copy.deepcopy(before)
|
||||
next(d for d in after if d['metadata']['name']=='actcore-worker')['spec']['template']['spec']['containers'][0]['image']='forgejo.coulomb.social/coulomb/activity-core@sha256:'+'f'*64
|
||||
(k/'runtime.yaml').write_text(yaml.safe_dump_all(after));candidate=commit(activity)
|
||||
seed=tmp_path/'seed';init(seed);p=seed/APPLICATION_PATH;p.parent.mkdir(parents=True)
|
||||
child={'apiVersion':'argoproj.io/v1alpha1','kind':'Application','metadata':{'name':'activity-core','namespace':'argocd'},'spec':{'project':'activity-core','source':{'repoURL':ACTIVITY_URL,'path':'k8s/gitops','targetRevision':rollback},'destination':{'namespace':'activity-core','server':'https://kubernetes.default.svc'},'syncPolicy':{'syncOptions':['CreateNamespace=false','ApplyOutOfSyncOnly=true','PruneLast=true','FailOnSharedResource=true']}}}
|
||||
p.write_text(yaml.safe_dump(child));commit(seed)
|
||||
platform=tmp_path/'platform.git';git(tmp_path,'clone','--bare',str(seed),str(platform))
|
||||
plan={'repository':'coulomb/activity-core','application':'activity-core','candidate':candidate,'rollback':rollback,'before_sha256':fingerprint(before),'after_sha256':fingerprint(after)}
|
||||
state={'tick':0,'calls':[],'lose_push':False,'fail_health':False}
|
||||
root={'metadata':{'name':'railiance-apps-root','namespace':'argocd','resourceVersion':'1'},'spec':{'source':{'repoURL':PLATFORM_URL,'path':'argocd/railiance01/applications','targetRevision':'main'}}}
|
||||
child['metadata']['resourceVersion']='1';apps={'railiance-apps-root':root,'activity-core':child}
|
||||
def runner(argv,**kwargs):
|
||||
if argv[0]=='git':
|
||||
argv=[str(activity) if a==ACTIVITY_URL else str(platform) if a==PLATFORM_URL else 'protocol.file.allow=always' if a=='protocol.file.allow=never' else a for a in argv]
|
||||
result=run(argv,**kwargs)
|
||||
if 'push' in argv and state['lose_push']:
|
||||
state['lose_push']=False;raise TransportError('fixture lost push response')
|
||||
return result
|
||||
state['calls'].append(argv)
|
||||
if 'get' in argv:
|
||||
i=argv.index('get');kind,name=argv[i+1:i+3]
|
||||
if kind=='namespace':return 'fixture-cluster'
|
||||
if kind=='application':return json.dumps(apps[name])
|
||||
return json.dumps({'metadata':{'generation':1},'spec':{'replicas':1},'status':{'observedGeneration':1,'readyReplicas':1,'updatedReplicas':1,'availableReplicas':1}})
|
||||
i=argv.index('patch');name=argv[i+2];patch=json.loads(kwargs['input']);obj=apps[name]
|
||||
assert patch[0]=={'op':'test','path':'/metadata/resourceVersion','value':obj['metadata']['resourceVersion']}
|
||||
assert patch[1]=={'op':'test','path':'/spec','value':obj['spec']}
|
||||
operation=patch[2]['value'];sync=operation['sync'];assert sync['prune'] is False
|
||||
if name=='railiance-apps-root':
|
||||
assert sync['resources']==[{'group':'argoproj.io','kind':'Application','name':'activity-core','namespace':'argocd'}]
|
||||
source=yaml.safe_load(git(platform,'show',sync['revision']+':'+APPLICATION_PATH))
|
||||
apps['activity-core']['spec']=source['spec']
|
||||
obj['status']={'operationState':{'phase':'Succeeded','syncResult':{'revision':sync['revision']}},'sync':{'status':'Synced','revision':sync['revision']},'health':{'status':'Healthy'}}
|
||||
return '{}'
|
||||
def pause(seconds):state['tick']+=seconds
|
||||
backend=GitArgoBackend(plan,tmp_path/'broker',lambda target:not(state['fail_health'] and target==candidate),cluster_uid='fixture-cluster',admitted=True,runner=runner,clock=lambda:state['tick'],pause=pause,wait_seconds=2)
|
||||
return backend,plan,state,platform,apps,before,after
|
||||
|
||||
|
||||
def test_real_git_publish_recover_sync_and_rollback(transport):
|
||||
b,p,state,platform,apps,*_=transport
|
||||
before=git(platform,'rev-parse','main');state['lose_push']=True
|
||||
with pytest.raises(TransportError):b.publish_revision(p['rollback'],p['candidate'])
|
||||
published=git(platform,'rev-parse','main');assert published!=before
|
||||
b.publish_revision(p['rollback'],p['candidate'])
|
||||
assert git(platform,'rev-parse','main')==published
|
||||
assert git(platform,'diff-tree','--no-commit-id','--name-only','-r',published)==APPLICATION_PATH
|
||||
b.sync_revision(p['candidate']);assert b.healthy(p['candidate'])
|
||||
state['fail_health']=True;assert not b.healthy(p['candidate'])
|
||||
b.publish_revision(p['candidate'],p['rollback']);b.sync_revision(p['rollback']);assert b.healthy(p['rollback'])
|
||||
assert apps['activity-core']['spec']['source']['targetRevision']==p['rollback']
|
||||
|
||||
|
||||
def test_manifest_tampering_refused_before_push(transport):
|
||||
b,p,state,platform,*_=transport;before=git(platform,'rev-parse','main');b.plan['after_sha256']='0'*64
|
||||
with pytest.raises(ValueError,match='manifest binding'):b.publish_revision(p['rollback'],p['candidate'])
|
||||
assert git(platform,'rev-parse','main')==before
|
||||
assert not any('patch' in c for c in state['calls'])
|
||||
|
||||
|
||||
def test_no_unbound_target_or_unadmitted_transport(transport,tmp_path):
|
||||
b,p,*_=transport
|
||||
with pytest.raises(ValueError):b.publish_revision(p['rollback'],'c'*40)
|
||||
with pytest.raises(ValueError):b.sync_revision('c'*40)
|
||||
with pytest.raises(ValueError):GitArgoBackend(p,tmp_path/'denied',lambda _:True,cluster_uid='fixture-cluster')
|
||||
|
||||
|
||||
def test_active_foreign_operation_not_overwritten(transport):
|
||||
b,p,state,platform,apps,*_=transport
|
||||
b.publish_revision(p['rollback'],p['candidate'])
|
||||
apps['railiance-apps-root']['operation']={'sync':{'revision':'d'*40,'prune':True}}
|
||||
with pytest.raises(ValueError,match='active'):b.sync_revision(p['candidate'])
|
||||
assert not any('patch' in c for c in state['calls'])
|
||||
|
||||
|
||||
def test_health_observer_requires_samples_and_survives_restart(tmp_path):
|
||||
o=HealthObserver(tmp_path/'health.sqlite');sha='a'*40
|
||||
o.observe(sha,True,NOW)
|
||||
o.observe(sha,True,NOW+timedelta(hours=24))
|
||||
with pytest.raises(ValueError,match='incomplete'):o.attestation(sha,NOW+timedelta(hours=24))
|
||||
start=NOW+timedelta(hours=24)
|
||||
for n in range(1,1441):o.observe(sha,True,start+timedelta(minutes=n))
|
||||
restarted=HealthObserver(tmp_path/'health.sqlite')
|
||||
result=restarted.attestation(sha,start+timedelta(hours=24))
|
||||
assert result['healthy_since']==start.isoformat()
|
||||
with pytest.raises(ValueError,match='stale'):restarted.attestation(sha,start+timedelta(hours=24,minutes=2))
|
||||
restarted.observe(sha,False,start+timedelta(hours=24,minutes=1))
|
||||
with pytest.raises(ValueError,match='no healthy'):restarted.attestation(sha,start+timedelta(hours=24,minutes=1))
|
||||
|
||||
|
||||
def test_observer_revision_change_and_clock_regression(tmp_path):
|
||||
o=HealthObserver(tmp_path/'health.sqlite');o.observe('a'*40,True,NOW)
|
||||
with pytest.raises(ValueError,match='non-increasing'):o.observe('a'*40,True,NOW)
|
||||
t=NOW+timedelta(seconds=60)
|
||||
assert o.observe('b'*40,True,t)==t.isoformat()
|
||||
with pytest.raises(ValueError):o.attestation('a'*40,t)
|
||||
|
||||
|
||||
def test_concurrent_git_writer_is_not_overwritten(transport):
|
||||
b,p,state,platform,apps,*_=transport
|
||||
original=b.runner
|
||||
def racing(argv,**kwargs):
|
||||
if argv[0]=='git' and 'push' in argv:
|
||||
# Inject another committed platform update after broker's CAS read.
|
||||
tree=git(platform,'rev-parse','main^{tree}')
|
||||
parent=git(platform,'rev-parse','main')
|
||||
env=dict(b.env)
|
||||
new=run(['git','-C',str(platform),'commit-tree',tree,'-p',parent],input='concurrent writer',env=env).strip()
|
||||
git(platform,'update-ref','refs/heads/main',new,parent)
|
||||
return original(argv,**kwargs)
|
||||
b.runner=racing
|
||||
with pytest.raises(TransportError):b.publish_revision(p['rollback'],p['candidate'])
|
||||
current=yaml.safe_load(git(platform,'show','main:'+APPLICATION_PATH))
|
||||
assert current['spec']['source']['targetRevision']==p['rollback']
|
||||
|
||||
|
||||
def test_broker_rolls_back_using_real_git_transport(transport,tmp_path):
|
||||
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
|
||||
from cryptography.hazmat.primitives.serialization import Encoding,PublicFormat
|
||||
from activity_core.release_broker import Broker,Receipts,ROLES
|
||||
from tests.test_release_broker import envelopes
|
||||
b,p,state,platform,apps,before,after=transport
|
||||
keys={role:Ed25519PrivateKey.generate() for role in ROLES}
|
||||
trust={r:{'public_key':k.public_key().public_bytes(Encoding.Raw,PublicFormat.Raw),'principal':r,'roles':{r}} for r,k in keys.items()}
|
||||
broker=Broker(tmp_path/'coordinator.sqlite',Receipts(trust),admitted=True)
|
||||
signed=envelopes(before,after,keys,candidate=p['candidate'],rollback=p['rollback'])
|
||||
rid=broker.admit(before,after,p['candidate'],p['rollback'],signed,NOW)
|
||||
state['fail_health']=True
|
||||
phases=[]
|
||||
for _ in range(10):
|
||||
phases.append(broker.advance(rid,b,NOW))
|
||||
if phases[-1]=='rolled_back':break
|
||||
assert phases==['publish_pending','published','synced','rollback_planned','rollback_published','rollback_synced','rolled_back']
|
||||
assert apps['activity-core']['spec']['source']['targetRevision']==p['rollback']
|
||||
assert yaml.safe_load(git(platform,'show','main:'+APPLICATION_PATH))['spec']['source']['targetRevision']==p['rollback']
|
||||
|
||||
|
||||
def test_clock_regression_invalidates_interval_until_new_sample(tmp_path):
|
||||
o=HealthObserver(tmp_path/'health.sqlite');sha='a'*40
|
||||
o.observe(sha,True,NOW)
|
||||
with pytest.raises(ValueError):o.observe(sha,True,NOW-timedelta(seconds=1))
|
||||
with pytest.raises(ValueError,match='invalidated'):HealthObserver(tmp_path/'health.sqlite').attestation(sha,NOW)
|
||||
later=NOW+timedelta(seconds=30)
|
||||
assert o.observe(sha,True,later)==later.isoformat()
|
||||
|
||||
|
||||
def test_wrong_cluster_refused_before_git_publication(transport):
|
||||
b,p,state,platform,*_=transport
|
||||
before=git(platform,'rev-parse','main');b.cluster_uid='wrong-cluster'
|
||||
with pytest.raises(ValueError,match='cluster identity'):
|
||||
b.publish_revision(p['rollback'],p['candidate'])
|
||||
assert git(platform,'rev-parse','main')==before
|
||||
assert not any('patch' in c for c in state['calls'])
|
||||
|
|
@ -169,3 +169,25 @@ continuous observer, Temporal dispatch and isolated transport/rollback proof.
|
|||
The broker is disabled by default and not connected to live credentials or a
|
||||
production schedule. Current deployed revision and healthy-soak clock are unchanged.
|
||||
T03 stays progress; the full unattended acceptance is not complete.
|
||||
|
||||
## Transport adapter and sampled health observer — 2026-09-27
|
||||
|
||||
Implemented the fixed Git/ArgoCD adapter: exact source/hash verification, static
|
||||
Kustomization restriction, one-field child revision commits, non-force publication,
|
||||
remote-tip recovery after lost responses, cluster UID binding, resourceVersion/spec
|
||||
CAS, selective sync and deployment health checks. A bounded report/schedule probe
|
||||
remains a required trusted integration. Implemented persistent health sampling with
|
||||
90-second maximum gaps, revision/failure reset, clock-regression invalidation and
|
||||
24-hour attestation refusal unless the complete sample interval exists.
|
||||
|
||||
Tests use disposable real Git repositories and simulated Kubernetes responses.
|
||||
The broker completes failed-health rollback through the real Git adapter; tests
|
||||
cover lost push responses, concurrent publication, source tampering, wrong cluster,
|
||||
foreign Argo operations and observer restart/gaps. Latest focused suite: 48 passed;
|
||||
full affected suite before the final cluster-binding test: 75 passed. CI includes
|
||||
all these tests. No production credentials or activation were introduced.
|
||||
|
||||
T03 remains progress for admitted identity/custody, authenticated isolated Kubernetes
|
||||
verification, real build/review/retention issuers, production invariant probes,
|
||||
Temporal observer/dispatcher registration and evidence sinks. The implemented
|
||||
observer cannot retroactively assert the earlier soak interval. See docs/release-broker.md.
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue