Aller au contenu
English

Enregistrer et reprendre un workflow

Utilisez des checkpoints pour reprendre un workflow et autorisez explicitement les nouvelles tentatives après une interruption.

Passez un checkpoint à la méthode start() du workflow si vous devez poursuivre dans un autre processus. Outpost enregistre les changements d’état des tâches et leurs résultats sous le runId du checkpoint.

import { reportValue } from "./reporter.ts";
import {
  createLocalTransport,
  defineTask,
  defineWorkflow,
  createWorkflowCheckpointStore,
} from "@elie-laloum/outpost";

const store = createWorkflowCheckpointStore({
  transporter: createLocalTransport({ directory: ".outpost/storage" }),
});
const scan = defineTask({ key: "scan", perform: () => ({ files: 12 }) });
const result = await defineWorkflow("scan", [scan]).start({
  checkpoint: { store, runId: "scan-2026-09", version: "1" },
});
result.unwrap();
reportValue(result.value(scan));
// Example output: { files: 12 }

Le script affiche { files: 12 } et enregistre le checkpoint sous .outpost/storage. Relancez-le : scan ne s’exécute pas, sa valeur vient du checkpoint.

  • Enregistrements des tâchesStatut, tentatives, erreurs, demandes et décisions des étapes d’approbation de chaque tâche.
  • SortiesLa valeur de chaque tâche done, restaurée au lieu d’exécuter la tâche à nouveau.
  • ConsommationTentatives et tokens cumulés, pour qu’un budget couvre toutes les reprises.

Une exécution reprise garde son executionId : context.idempotencyKey reste donc identique pour chaque tâche. Les sandboxes, leurs fichiers et le code des tâches ne sont pas enregistrés : reprenez en appelant start() sur la même définition de workflow.

Une tâche avec checkpoint doit renvoyer undefined ou une valeur qui peut être enregistrée et relue en JSON sans perdre d’information. Sinon, la tentative échoue. Convertissez les dates en chaînes et ne gardez que les champs utiles.

Un résultat de dispatch porte des méthodes comme resume(). Projetez-le dans une defineTask, comme dans l’exemple ci-dessous :

import { defineIsolatedTask, defineTask } from "@elie-laloum/outpost";
import { coder, repository, sandboxProvider } from "./outpost.config.ts";

const agent = defineIsolatedTask({
  key: "fix-agent",
  request: () => ({
    repository,
    sandboxProvider,
    agent: coder,
    brief: { text: "Fix the failing tests and commit the fix." },
  }),
});
export const fix = defineTask({
  key: "fix",
  perform: async (context) => {
    const { branch, commits } = await agent.perform(context);
    return { branch, commits, finishedAt: new Date().toISOString() };
  },
});

Un checkpoint est limité à 16 Mio. Stockez les contenus volumineux comme artefacts et renvoyez leur référence.

Un checkpoint ne reprend que le workflow qui l’a écrit. start() refuse un checkpoint dont l’identité diffère.

Élément de l’identitéOù vous le définissez
Nom du workflowdefineWorkflow(name, tasks)
Versioncheckpoint.version
GrapheClés des tâches et leurs dépendances after
Réglages d’exécutiontimeoutMs, réglages de retry, présence de condition ou retry.accepts
GatesType, prompt, actors et authentication de chaque approbation ou pause
Tâches en bouclemaxRounds
Tâches interactivesactors, agent, modèle, brief, dépôt, maxTurns et fournisseur de sandbox

Le reste du code des tâches, les briefs et les entrées du workflow n’en font pas partie : changez version quand vous les modifiez. Une exécution enregistrée ne peut pas changer d’identité, pas même de version : relancez-la sous un nouveau runId.

Le budget du workflow ne fait pas non plus partie de l’identité.

Pour réutiliser des résultats entre exécutions différentes, utilisez plutôt le cache de résultats.

Une exécution terminée sur une tâche en échec, annulée ou interrompue ne reprend qu’avec resume: "retry-incomplete". Cette option autorise à exécuter ces tâches à nouveau, avec leurs effets de bord.

import {
  createWorkflowCheckpointStore,
  createLocalTransport,
} from "@elie-laloum/outpost";

export const store = createWorkflowCheckpointStore({
  transporter: createLocalTransport({ directory: ".outpost/storage" }),
});
import { defineTask } from "@elie-laloum/outpost";

export let calls = 0;
export const upload = defineTask({
  key: "upload",
  perform: () => {
    calls += 1;
    if (calls === 1) throw new Error("Network unavailable");
    return { uploaded: true };
  },
});
export function uploadCount() {
  return calls;
}
import { reportValue } from "./reporter.ts";
import { defineWorkflow } from "@elie-laloum/outpost";
import { upload } from "./upload.ts";
import { store } from "./upload-store.ts";

export const workflow = defineWorkflow("upload", [upload]);
export const checkpoint = { store, runId: "upload-1", version: "1" };
reportValue((await workflow.start({ checkpoint })).status);
// Example output: failed
export const resumed = await workflow.start({
  checkpoint: { ...checkpoint, resume: "retry-incomplete" },
});
reportValue(resumed.status);
// Example output: done

Le script affiche failed, puis done. Sans resume, le second start() est refusé.

État enregistréSans resumeAvec resume: "retry-incomplete"
Toutes les tâches done ou skippedRenvoie le résultat enregistré, n’exécute rienIdentique
En pause sur une étape d’approbation ou en attente d’une réponseContinue avec vos décisions ou réponsesIdentique
En pause sur un quotaRelance la tâche après la réinitialisationIdentique
Une tâche failed, cancelled ou interrompuestart() est refuséLa relance, ainsi que les tâches qu’elle a sautées

Une tâche relancée repart pour une nouvelle série de tentatives retry. Une tâche done ne s’exécute jamais à nouveau. Une exécution arrêtée par son budget reprend de la même façon ; passez un budget plus large, car la consommation continue de s’additionner.

Une exécution possède son checkpoint pendant start() et le libère quand start() se termine. Si le processus meurt, la propriété reste : tout start() suivant pour ce runId est refusé jusqu’à ce que vous la libériez.

import { createHash } from "node:crypto";
import {
  createLocalTransport,
  recoverWorkflowCheckpoint,
} from "@elie-laloum/outpost";

const transporter = createLocalTransport({ directory: ".outpost/storage" });
const runId = "scan-2026-09";
const digest = createHash("sha256").update(runId).digest("hex");
const saved = await transporter.read(`checkpoints/${digest}.json`);
if (saved)
  await recoverWorkflowCheckpoint({
    transporter,
    runId,
    revision: saved.revision,
  });
Glissez pour vous déplacer · Ctrl + molette pour zoomer
100 %
  • ArrêterAssurez-vous que l’ancien processus n’écrit plus.
    1. Arrêter le processusConfirmez qu’il s’est terminé. Un PID ne prouve pas qu’un processus distant s’est arrêté.
    (Étapes)
    • → Déverrouiller : puis
  • DéverrouillerLibérez la propriété, gardez la progression.
    1. Lire la révisionLisez l’objet du checkpoint de l’exécution depuis le transport. Transport
    2. Libérer la propriétéL’appel est refusé si l’objet a changé depuis votre lecture. recoverWorkflowCheckpoint()
    (Étapes)
    • → Reprendre : puis
  • ReprendreRelancez le même workflow avec le même checkpoint.
    1. Autoriser le rejeuLa tâche interrompue est relancée avec resume: "retry-incomplete". start()
    (Étapes)

La clé du checkpoint est checkpoints/ suivi du SHA-256 du runId. La progression reste intacte ; relancez l’exécution avec resume: "retry-incomplete".

defineWorkflowJob() exécute chaque job sous son runId. Une file renvoie le job existant pour un identifiant qu’elle connaît déjà : un job terminé ne s’exécute donc jamais à nouveau.

Pour poursuivre l’exécution, publiez un nouvel identifiant de job avec le même runId et la même input :

import { createSqliteTaskQueue } from "@elie-laloum/outpost";

const queue = await createSqliteTaskQueue(".outpost/jobs.sqlite");
try {
  await queue.enqueue({
    id: "fix-42-resume-1",
    handler: "fix",
    input: { runId: "fix-42", input: { issue: 42 } },
  });
} finally {
  queue.close();
}

Une autre input change la version du checkpoint et le job échoue. Pour rejouer des tâches en échec ou interrompues, le traitement doit recevoir checkpoint: { store, version, resume: "retry-incomplete" }.

Le stockage accepte n’importe quel Transport. Utilisez un transport S3 ou R2 pour que des workers sur plusieurs machines partagent les exécutions : voir Où vivent les données.

  • Un seul start() à la fois par runId. Un second est refusé tant que le premier s’exécute.
  • La propriété n’expire jamais d’elle-même. Libérez-la avec recoverWorkflowCheckpoint() après avoir arrêté l’ancien processus.
  • Le rejeu répète les effets de bord qu’une tâche interrompue a déjà produits. Dédupliquez-les avec context.idempotencyKey : voir Files de jobs et workers.

API : createWorkflowCheckpointStore · WorkflowCheckpointOptions · recoverWorkflowCheckpoint · WorkflowCheckpoint · defineWorkflowJob.