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é

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

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

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

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

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

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"}  → ValidationError
Astuce

Utilisez 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.

Union optionnelle (| 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

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

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.

Note

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

Fonctionnalités clés du workflow

Signaux

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 = True

Plus d’informations sur les signaux

Requêtes

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

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

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
        pass

Fonctionnalité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

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 handle

Par 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

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

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

Fonctionnalités avancées

Mécanisme de réinitialisation

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

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()