Le notebook vérifie les barrières de readiness et de compatibilité, écrit six étapes puis simule un crash, rejoue le lot entier et constate six reçus replayed pour zéro doublon, provoque un IDEMPOTENCY_CONFLICT puis un IDEMPOTENCY_KEY_REQUIRED, ouvre une session et lit sa parité d’audit, et termine sur une mutation refusée par la route de lecture avec un HTTP 200 porteur d’un statut Rejected.
durable-runtime/notebook.ipynbEXTRAIT / 16 CELLULES
MARKDOWN [0]Reproduisez le scénario sur une copie de travail et conservez les identifiants de preuves, de session et de snapshot dans le résultat.
CODE [1] / PYTHONimport json, os, pathlib, time
import requests
import urllib3
# The quick-start server uses a self-signed certificate; disable verification
# for the local playbook only, never against a real deployment.
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
BASE_URL = os.environ.get('CORROBORE_URL', 'https://127.0.0.1:8080')
TOKEN = os.environ.get('CORROBORE_HTTP_AUTH_TOKEN', 'change-me')
http = requests.Session()
http.verify = False
http.headers.update({'Authorization': f'Bearer {TOKEN}', 'Content-Type': 'application/json'})
class MemoryError_(RuntimeError):
"""Carries the stable v1 error taxonomy instead of a bare HTTP status."""
def __init__(self, code, message, status):
super().__init__(f'{code}: {message}')
self.code, self.message, self.status = code, message, status
def memory_op(operation, payload, idempotency_key=None, expect_error=False):
"""POST /v1/memory/operations and unwrap the typed result."""
body = {'contract_version': 'v1', 'operation': operation, 'input': payload}
if idempotency_key is not None:
body['idempotency_key'] = idempotency_key
for attempt in range(12):
response = http.post(f'{BASE_URL}/v1/memory/operations', data=json.dumps(body))
# Protected routes share a global token bucket (50 rps sustained, 200 burst by
# default) that a bulk load will hit. Two details matter here: the 429 body is
# plain text, not the JSON error envelope, and `Retry-After` can be `0` — so
# honour it as a floor, never as the whole wait.
if response.status_code == 429:
hinted = float(response.headers.get('Retry-After', 0) or 0)
time.sleep(max(hinted, 0.2 * (attempt + 1)))
continue
break
if response.status_code != 200:
try:
error = response.json().get('error', {})
except ValueError: # 429 and other transport rejections are not JSON
error = {}
failure = MemoryError_(error.get('code', 'UNKNOWN'), error.get('message', response.text), response.status_code)
if expect_error:
return failure
raise failure
if expect_error:
raise AssertionError(f'{operation} unexpectedly succeeded')
return response.json()['result']['result']
ready = http.get(f'{BASE_URL}/health/ready').json()
version = http.get(f'{BASE_URL}/version').json()
print('ready :', json.dumps(ready)[:160])
print('version:', json.dumps(version)[:160])
CODE [2] / PYTHONlive = http.get(f'{BASE_URL}/health/live').json()
ready = http.get(f'{BASE_URL}/health/ready').json()
health = http.get(f'{BASE_URL}/health').json()
version = http.get(f'{BASE_URL}/version').json()
print('live :', live['live'], live['lifecycle_state'])
print('ready :', ready['ready'])
for check, value in ready['checks'].items():
print(f' {check:<24} {value}')
print('storage :', health['storage_mode'], '| uptime_ms', health['uptime_ms'])
print('version :', version['version'], version['build_target'])
print('storage compatibility:', json.dumps(version['storage_compatibility']))
# Refuse to write into a runtime that has not finished recovering.
assert ready['ready'] is True
assert ready['checks']['storage_recovered'] is True
CODE [3] / PYTHONDATASET = pathlib.Path('dataset.json')
dataset = json.loads(DATASET.read_text())
CRASH_AFTER = dataset['facts']['crash_after_step']
ids, receipts = {}, {}
def write_trail(records, stop_after=None):
"""Write the trail, optionally dying partway through."""
written = 0
for record in records:
key = record['identity_key']
result = memory_op('remember', dict(record), idempotency_key=f'runtime:{key}')
ids[key] = result['record']['id']
receipts.setdefault(key, []).append(result['receipt'])
written += 1
if stop_after is not None and written >= stop_after:
raise RuntimeError(f'process died after {written} writes')
return written
try:
write_trail(dataset['memories'], stop_after=CRASH_AFTER)
except RuntimeError as failure:
print('simulated crash:', failure)
print('committed before the crash:', len(ids))
CODE [4] / PYTHON# The restarted agent replays the whole batch. It has no idea what it already did.
written = write_trail(dataset['memories'])
print('write calls issued on resume:', written)
replayed = [key for key, history in receipts.items() if any(r['replayed'] for r in history)]
print('receipts marked replayed :', len(replayed))
for key in sorted(replayed)[:3]:
first, second = receipts[key][0], receipts[key][1]
print(f" {key}")
print(f" first : version {first['committed_version']} replayed={first['replayed']} audit={first['audit_correlation_id'][:8]}")
print(f" resume: version {second['committed_version']} replayed={second['replayed']} audit={second['audit_correlation_id'][:8]}")
# A replay returns the *original* receipt, including its audit correlation id.
for key in replayed:
assert receipts[key][0]['committed_id'] == receipts[key][1]['committed_id']
assert receipts[key][0]['committed_version'] == receipts[key][1]['committed_version']
assert len(replayed) == CRASH_AFTER
CODE [5] / PYTHONmutated = dict(dataset['memories'][1])
mutated['content'] = {'format': 'text_and_properties',
'value': {'text': 'rewritten after the fact', 'properties': {}}}
conflict = memory_op('remember', mutated,
idempotency_key=f"runtime:{mutated['identity_key']}", expect_error=True)
print('reused key, new payload:', conflict.status, conflict.code, '-', conflict.message)
assert conflict.code == 'IDEMPOTENCY_CONFLICT'
missing = memory_op('remember', dict(dataset['memories'][0]), expect_error=True)
print('no key at all :', missing.status, missing.code, '-', missing.message)
assert missing.code == 'IDEMPOTENCY_KEY_REQUIRED'
CODE [6] / PYTHONopened = http.post(f'{BASE_URL}/v1/sessions/start', data=json.dumps({
'workspace_id': 'workspace--standalone-default',
'actor_id': 'incident-bot',
'actor_kind': 'agent',
})).json()
session_id = opened['result']['session_id']
print('session:', session_id, opened['result']['status'])
http.post(f'{BASE_URL}/v1/cypher/write', data=json.dumps({
'query': "MERGE (n:IncidentTag {name:'INC-2291'}) RETURN n",
'session_id': session_id,
}))
logs = http.get(f'{BASE_URL}/v1/sessions/{session_id}/logs', params={'limit': 10}).json()['result']
print('log path :', logs['log_path'].rsplit('/', 1)[-1])
print('matched/total :', logs['matched_entries'], '/', logs['total_matched_entries'])
print('audit parity :', json.dumps(logs['audit_parity'], indent=2))
parity = logs['audit_parity']
assert parity['parity_ok'] is True
assert parity['missing_output_event_ids'] == []
assert parity['orphan_output_event_ids'] == []
CODE [7] / PYTHONresponse = http.post(f'{BASE_URL}/v1/cypher/read',
data=json.dumps({'query': 'CREATE (n:ShouldNotExist) RETURN n'}))
body = response.json()
print('HTTP status :', response.status_code)
print('execution status:', body['result']['status'])
for issue in body['result']['validation_errors']:
print(f" {issue['code']}: {issue['message']}")
assert response.status_code == 200
assert body['result']['status'] == 'Rejected'
assert any(e['code'] == 'WRITE_PERMISSION_REQUIRED' for e in body['result']['validation_errors'])