Workflows
Cette page présente les schémas pratiques pour concevoir des workflows. Pour mieux comprendre la logique sous-jacente (durabilité, relecture, historique, déterminisme), consultez Concepts clés > Workflows.
Workflow versus activité
Lorsque vous concevez l’orchestration de votre application, la première décision à prendre est de déterminer ce qui relève du workflow et ce qui relève d’une activité.
Un workflow agit comme le cerveau de l’orchestration. Il coordonne les étapes, prend des décisions, conserve l’état sur de longues durées (de quelques secondes à plusieurs années) et attend des événements externes. Le code du workflow doit être déterministe pour permettre à la plateforme de le rejouer depuis l’historique en cas de redémarrage d’un worker.
Une activité est l’unité de travail effectif — tout ce qui interagit avec l’extérieur : appels LLM, requêtes HTTP, écritures en base de données, lectures de fichiers, exécution d’outils, etc. Les activités n’ont pas besoin d’être déterministes, s’exécutent sur des périodes courtes (de quelques secondes à quelques minutes) et sont automatiquement retentées avec un backoff configurable en cas d’échec.
Par exemple, un workflow qui traite des factures entrantes peut appeler une activité fetch_invoice_pdf (téléchargement HTTP), puis une activité extract_fields (appel LLM avec sortie structurée), effectuer une branche sur le montant extrait pour déterminer si une validation humaine est requise, attendre un signal approve, puis appeler une activité post_to_accounting_system. L’ordre, la logique de branchement et l’attente s’expriment dans le workflow ; chaque effet de bord vit dans une activité.
Pour la contrepartie « activité » de cette page, voir Activités.
Définir un workflow
Pour définir un workflow, utilisez le décorateur @workflow.define sur une classe et marquez la fonction d’entrée avec @workflow.entrypoint :
import mistralai.workflows as workflows
@workflows.workflow.define(name="report_workflow")
class ReportWorkflow:
@workflows.workflow.entrypoint
async def run(self, report_type: str, include_details: bool = False) -> dict:
"""Workflow implementation."""
# Orchestrer les activités ici
...Entrée du workflow
La méthode run accepte tout type sérialisable en JSON — modèles Pydantic, dictionnaires natifs ou types primitifs. La plateforme valide les entrées et génère les schémas d’UI Console dans tous les cas. Pour une vue conceptuelle, voir Concepts clés > Entrée du workflow.
Entrées primitives et multi-paramètres
@workflow.define(name="report_workflow")
class ReportWorkflow:
@workflow.entrypoint
async def run(self, report_type: str, include_details: bool = False) -> dict:
...Exemple d’entrée :
{"report_type": "daily", "include_details": false}Modèle Pydantic unique
Quand le point d’entrée prend un seul BaseModel Pydantic, ses champs deviennent les clés d’entrée de premier niveau ; il n’y a pas d’objet englobant :
from pydantic import BaseModel
class ReportParams(BaseModel):
report_type: str
include_details: bool = False
@workflow.define(name="report_workflow")
class ReportWorkflow:
@workflow.entrypoint
async def run(self, params: ReportParams) -> dict:
...Exemple d’entrée :
{"report_type": "daily", "include_details": true}Union de modèles Pydantic
Quand un workflow peut être déclenché avec différentes structures d’entrée, utilisez une union de sous-classes de BaseModel :
from pydantic import BaseModel, ConfigDict
class PromptInput(BaseModel):
model_config = ConfigDict(extra="forbid")
prompt: str
class CountInput(BaseModel):
model_config = ConfigDict(extra="forbid")
count: int
@workflow.define(name="flexible_workflow")
class FlexibleWorkflow:
@workflow.entrypoint
async def run(self, params: PromptInput | CountInput) -> str:
if isinstance(params, PromptInput):
return f"prompt: {params.prompt}"
return f"count: {params.count}"Le SDK valide l’entrée selon les membres de l’union dans l’ordre et transmet le premier correspondant au handler :
{"prompt": "hello"} → PromptInput
{"count": 42} → CountInput
{"count": 42, "extra": "field"} → ValidationErrorUtilisez extra="forbid" sur chaque modèle membre. Cela permet une discrimination précise et des messages d’erreur explicites quand l’entrée ne correspond à aucune forme attendue.
| None)Union optionnelle (| None)
Ajoutez | None (ou utilisez Optional depuis typing) pour rendre l’entrée optionnelle :
@workflow.define(name="optional_workflow")
class OptionalWorkflow:
@workflow.entrypoint
async def run(self, params: PromptInput | CountInput | None) -> str:
if params is None:
return "no input"
if isinstance(params, PromptInput):
return f"prompt: {params.prompt}"
return f"count: {params.count}"Combiner BaseModel et types primitifs
Les unions qui combinent un BaseModel et un type non-modèle (comme str ou int) ne sont pas prises en charge et provoqueront une erreur lors de l’application du décorateur @workflow.entrypoint :
# ❌ Lève une erreur lors de la définition de la classe
@workflow.define(name="my_workflow")
class MyWorkflow:
@workflow.entrypoint
async def run(self, params: PromptInput | str) -> str:
...Le paramètre 'params' de 'run' comporte des membres union non pris en charge : str.
Les unions d’entrées n’acceptent que des sous-classes de Pydantic BaseModel et None.
Utilisez un BaseModel dédié à la place.Si vous devez vraiment traiter une union avec un type non-modèle, définissez une sous-classe nommée de RootModel et utilisez-la à la place :
from pydantic import BaseModel, RootModel
class PromptInput(BaseModel):
prompt: str
class PromptOrRaw(RootModel[PromptInput | str]):
pass
@workflow.define(name="my_workflow")
class MyWorkflow:
@workflow.entrypoint
async def run(self, params: PromptOrRaw) -> str:
if isinstance(params.root, PromptInput):
return f"structured: {params.root.prompt}"
return f"raw: {params.root}"Cette sous-classe nommée permet aussi de générer un schéma JSON plus clair.
Timeout d’exécution
Pour le rôle et la valeur par défaut de execution_timeout, voir Concepts clés > Timeout d’exécution. Pour modifier la valeur par défaut (1h), transmettez un timedelta au décorateur :
from datetime import timedelta
from mistralai.workflows import workflow
@workflow.define(name="long_running_workflow", execution_timeout=timedelta(days=7))
class LongRunningWorkflow:
@workflow.entrypoint
async def run(self, params: MyParams) -> MyResult:
...Un workflow encore en cours après expiration du execution_timeout est annulé avec une erreur WORKFLOW_EXECUTION_TIMED_OUT.
execution_timeout est une limite totale en temps réel, pas un timeout d’activité. Les timeouts d’activité se gèrent avec start_to_close_timeout ou schedule_to_close_timeout sur chaque appel d’activité.
Fonctionnalités clés du workflow
Signaux
Un signal est un message fire-and-forget livré à un workflow en cours d’exécution. L’émetteur n’attend pas de réponse : le workflow traite le signal de façon asynchrone et met à jour son propre état. Utilisez les signaux pour des validations humaines, des callbacks de systèmes externes, ou tout événement destiné à réveiller un workflow sans attendre de valeur de retour.
Par exemple, un workflow de traitement de factures peut se mettre en pause jusqu’à ce qu’un validateur envoie un signal d’approbation :
@workflows.workflow.signal(name="approve", description="Approval signal")
async def approve_signal(self, data: ApprovalData) -> None:
self.approved = TruePlus d’informations sur les signaux
Requêtes
Une requête lit l’état du workflow sans le modifier. Utilisez des requêtes pour exposer l’avancement, l’état actuel ou des valeurs calculées à des clients externes — par exemple, un dashboard qui consulte le pourcentage d’éléments traités.
@workflows.workflow.query(name="get_status", description="Get current status")
def get_status(self) -> WorkflowStatus:
return WorkflowStatus(
progress=self.progress,
status=self.current_status
)Plus d’informations sur les requêtes
Mises à jour
Une mise à jour fonctionne comme un signal, mais l’émetteur attend que le workflow traite le message et réponde. Utilisez une mise à jour si vous avez besoin d’une confirmation ou d’un résultat calculé — par exemple, pour ajuster une configuration de workflow en cours d’exécution et récupérer le nouvel état, ou pour soumettre une valeur et recevoir la version validée et persistée.
@workflows.workflow.update(name="update_config", description="Update workflow configuration")
async def update_config(self, config: UpdateConfig) -> UpdateResult:
# Traitement de la mise à jour
return UpdateResult(success=True)Plus d’informations sur les mises à jour
Plannings
Définissez des exécutions planifiées de workflow à l’aide d’expressions cron :
from mistralai.workflows.models import ScheduleDefinition
schedule = ScheduleDefinition(
input={"report_type": "daily", "include_details": False}, # Paramètres d’entrée pour les exécutions planifiées
cron_expressions=["0 0 * * *"] # Exécution quotidienne à minuit
)
@workflows.workflow.define(name="report_workflow", schedules=[schedule])
class ReportWorkflow:
@workflows.workflow.entrypoint
async def run(self, report_type: str = "daily", include_details: bool = False) -> None:
# Générer le rapport
passFonctionnalités principales du planning :
- Planification basée sur des expressions cron standard
- Paramètres d’entrée pour les exécutions planifiées
- Possibilité de spécifier plusieurs expressions cron
Plus d’informations sur la planification
Workflows enfants
Exécutez d’autres workflows en tant qu’enfants :
from datetime import timedelta
from mistralai.workflows import workflow
result = await workflow.execute_workflow(
ChildWorkflow,
params=child_params,
execution_timeout=timedelta(hours=1)
)Vous pouvez également lancer un workflow enfant sans attendre son résultat (fire-and-forget) en passant wait=False :
handle = await workflow.execute_workflow(
ChildWorkflow,
params=child_params,
execution_timeout=timedelta(hours=1),
wait=False,
)
# Le parent continue immédiatement ; l’enfant s’exécute indépendamment
# Vous pouvez attendre plus tard : result = await handlePar défaut, wait=False applique la politique de fermeture parent à ABANDON, ce qui signifie que l’enfant continue à s’exécuter même si le parent termine. Vous pouvez remplacer ce comportement avec le paramètre parent_close_policy :
from mistralai.workflows import ParentClosePolicy
handle = await workflow.execute_workflow(
ChildWorkflow,
params=child_params,
execution_timeout=timedelta(hours=1),
wait=False,
parent_close_policy=ParentClosePolicy.TERMINATE,
)Attendre des conditions
Mettez le workflow en pause jusqu’à ce qu’une condition soit remplie (en général via un signal ou une mise à jour reçue). Pendant ce temps, le workflow « dort » sans consommer de ressources. Si la condition n’est pas remplie avant expiration du timeout, une TimeoutError est levée.
from mistralai.workflows import workflow
await workflow.wait_condition(
lambda: self.ready,
timeout=timedelta(minutes=5)
)Continue-as-new
Réinitialisez l’historique d’un workflow tout en le maintenant logique en vie — utile pour des workflows à itérations infinies ou des exécutions proches de la limite de 50 000 événements historiques.
Pour le schéma complet (paramètres d’état, déclencheur, implications sur le déterminisme), voir Continue-as-new.
Fonctionnalités avancées
Mécanisme de réinitialisation
Vous pouvez « réinitialiser » une exécution de workflow pour la redémarrer à partir d’un point précis de son historique. Cela est utile si un workflow est bloqué (par exemple, suite à une erreur non déterministe) ou ne parvient pas à se terminer. Lorsque vous réinitialisez, le workflow termine son exécution en cours puis redémarre depuis l’événement de votre choix dans l’historique.
Tous les événements du workflow jusqu’au point de réinitialisation sont copiés dans la nouvelle exécution. Le workflow reprend alors à partir de là avec la version courante de votre code. Toute avancée après l’événement de réinitialisation est perdue.
Ne réinitialisez qu’après avoir corrigé le problème d’origine. Nous vous recommandons d’indiquer une raison pour la réinitialisation ; celle-ci sera enregistrée dans l’historique des événements du workflow.
from mistralai.client import Mistral
client = Mistral(api_key="your_api_key")
client.workflows.executions.reset_workflow(
execution_id="your-execution-id",
event_id=42, # Doit être un événement WORKFLOW_TASK_COMPLETED
reason="Bug corrigé dans la logique de l’activité",
exclude_signals=True, # Optionnel : ne pas rejouer les signaux après ce point
exclude_updates=True # Optionnel : ne pas rejouer les mises à jour après ce point
)Contexte d’exécution du workflow
Récupérez à l’exécution l’identifiant unique courant. Utile pour journaliser, relier des événements entre services ou transmettre cet ID à des systèmes externes afin qu’ils puissent émettre des signaux en retour.
import mistralai.workflows as workflows
execution_id = workflows.get_execution_id()