← Toutes les capacités

10 / OPERATIONS

Reprise durable, audit et exécution sûre

La mémoire d’un agent doit survivre au processus, refuser les écritures ambiguës et expliquer ce qui s’est passé avant comme après une reprise.

HTTP + RUST CORE
01

Pourquoi cette capacité

Faire de la fiabilité une capacité produit

Le runtime combine stockage WAL, snapshots, promotion atomique, récupération stricte, lectures et écritures séparées, authentification, sessions persistantes, logs JSONL et signaux de readiness.

02

Scénario

FIELD NOTE / SCENARIO

Reprendre après une interruption pendant un commit

Le service s’arrête entre la préparation d’un snapshot et sa promotion. Au redémarrage, la récupération sélectionne la dernière génération durable, vérifie les journaux et n’annonce ready qu’après validation.

  1. 01

    Écrire par transaction avec journal durable

  2. 02

    Restaurer la dernière génération cohérente

  3. 03

    Vérifier readiness, requête protégée et parité d’audit

Résultat attendu

L’agent reprend sur un état explicable sans accepter un snapshot partiel ni masquer le défaut opérationnel.

03

Jeu de données synthétique

Une piste d’incident de treize étapes, interrompue à la sixième

Le jeu de données n’est pas un graphe de connaissance mais une charge de travail : la trace ordonnée qu’un agent écrit pendant la remédiation d’un incident de pipeline. Quatorze mémoires, vingt-cinq relations, une clé d’idempotence déterministe par écriture. Le notebook tue le processus après la sixième étape puis rejoue le lot entier depuis le début, exactement comme le ferait un agent redémarré sans mémoire de sa progression.

14Mémoires
25Relations
13Points de contrôle
6Crash après l’étape
6Phases
0Doublons tolérés
Télécharger le jeu de donnéesdataset.json · JSON · 23 kB
04

Pourquoi Corrobore

Avoir des journaux et pouvoir affirmer qu’ils sont complets sont deux choses différentes.

Ce qu’il faut construire soi-même

  • Un healthcheck qui répond 200 dès que le processus écoute ne dit rien de la reprise du stockage. On écrit dans un moteur qui n’a pas fini de récupérer, et on l’apprend plus tard.
  • La reprise après interruption demande d’écrire sa propre table de progression, de la garder cohérente avec les écritures, et d’espérer que les deux n’ont pas divergé pendant la panne.
  • Un fichier de journal ne prouve rien à lui seul : rien ne réconcilie ce que le moteur a accepté avec ce qu’il a émis, donc rien ne signale un événement manquant.

Ce que le contrat garantit

  • /health/ready expose storage_recovered séparément de la vivacité, et /version publie l’enveloppe de compatibilité du stockage : deux barrières avant la première écriture.
  • Le rejeu est gratuit si les clés sont déterministes. Une mutation rejouée renvoie le reçu d’origine, identifiant de corrélation d’audit compris, et ne crée aucun doublon.
  • Les journaux de session arrivent avec leur propre réconciliation audit_parity, qui nomme les événements manquants et orphelins au lieu de vous laisser faire confiance au fichier.
05

Notebook exécutable

Quinze cellules qui tuent un agent en plein vol et le ramènent

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] / PYTHON
import 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] / PYTHON
live = 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] / PYTHON
DATASET = 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] / PYTHON
mutated = 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] / PYTHON
opened = 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] / PYTHON
response = 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'])
Télécharger le notebooknotebook.ipynb · IPYNB · 16 cellules
06

Playbook à télécharger

Trois minutes pour vérifier qu’une reprise ne duplique rien

Le playbook démarre un Corrobore éphémère, exécute le notebook et propose de déplacer l’étape du crash, ou de tuer réellement le conteneur en stockage persistant pour voir la récupération s’exécuter.

  1. 01

    Démarrer un Corrobore jetable en stockage éphémère

  2. 02

    Installer requests

  3. 03

    Exécuter le notebook et lire les reçus rejoués

  4. 04

    Déplacer l’étape du crash, relancer, recompter

Télécharger le playbookplaybook.md · MARKDOWN · requests
07

Ce qui est disponible

Serveur HTTP authentifié, stockage durable, reprise, sessions, logs, métriques et health checks sont publiés. Les primitives avancées de raisonnement des autres pages peuvent rester Rust-only.