Flujos de trabajo de Microsoft Agent Framework: puntos de control

En esta página se proporciona información general sobre los puntos de control en el sistema de flujo de trabajo de Microsoft Agent Framework.

Overview

Los puntos de control permiten guardar el estado de un flujo de trabajo en puntos específicos durante su ejecución y reanudarse desde esos puntos más adelante. Esta característica es especialmente útil para los siguientes escenarios:

  • Flujos de trabajo de larga duración en los que desea evitar perder el progreso en caso de errores.
  • Flujos de trabajo de larga duración en los que desea pausar y reanudar la ejecución más adelante.
  • Flujos de trabajo que requieren el guardado de estado periódico para fines de auditoría o cumplimiento.
  • Flujos de trabajo que deben migrarse en diferentes entornos o instancias.

¿Cuándo se crean los puntos de control?

Recuerde que los flujos de trabajo se ejecutan en superpasos, como se documenta en el modelo de ejecución de flujo de trabajo. Los puntos de control se crean al final de cada superpaso, después de que todos los ejecutores de ese superpaso hayan completado su ejecución. Un punto de control captura todo el estado del flujo de trabajo, entre los que se incluyen:

  • Estado actual de todos los ejecutores
  • Todos los mensajes pendientes del flujo de trabajo para el siguiente superpaso
  • Solicitudes y respuestas pendientes
  • Estados compartidos

Note

A partir de Python versión 1.13.0, los flujos de trabajo también crean un punto de control de entrada antes del primer superpaso para registrar la entrada del flujo de trabajo y otro punto de control de entrada cuando se entregan respuestas a eventos de solicitud. Estos puntos de control hacen que el flujo de trabajo completo se pueda reproducir. Esta versión incluye cambios importantes menores para las aplicaciones que dependen de recuentos de iteración, identificadores de origen de mensajes o ordenación de puntos de comprobación. Los puntos de control existentes siguen siendo compatibles. Para obtener más información sobre la migración, consulte Actualización de puntos de comprobación de flujo de trabajo Python a la versión 1.13.0.

Capturar puntos de control

Para habilitar los puntos de control, es necesario proporcionar un elemento CheckpointManager al ejecutar el flujo de trabajo. A continuación, se puede acceder a un punto de control a través de un SuperStepCompletedEvent, o a través de la propiedad Checkpoints en la ejecución.

using Microsoft.Agents.AI.Workflows;

// Create a checkpoint manager to manage checkpoints
CheckpointManager checkpointManager = CheckpointManager.CreateInMemory();

// Run the workflow with checkpointing enabled
StreamingRun run = await InProcessExecution
    .RunStreamingAsync(workflow, input, checkpointManager)
    .ConfigureAwait(false);
await foreach (WorkflowEvent evt in run.WatchStreamAsync().ConfigureAwait(false))
{
    if (evt is SuperStepCompletedEvent superStepCompletedEvt)
    {
        // Access the checkpoint
        CheckpointInfo? checkpoint = superStepCompletedEvt.CompletionInfo?.Checkpoint;
    }
}

// Checkpoints can also be accessed from the run directly
IReadOnlyList<CheckpointInfo> checkpoints = run.Checkpoints;

Para habilitar los puntos de control, es necesario proporcionar un CheckpointStorage elemento al crear un flujo de trabajo. A continuación, se puede acceder a un punto de control a través del almacenamiento. Agent Framework incluye tres implementaciones integradas: elija la que coincida con sus necesidades de durabilidad e implementación:

Provider Package Durability Más adecuado para
InMemoryCheckpointStorage agent-framework Solo en proceso Pruebas, demostraciones, flujos de trabajo de corta duración
FileCheckpointStorage agent-framework Disco local Flujos de trabajo de una sola máquina, desarrollo local
CosmosCheckpointStorage agent-framework-azure-cosmos Azure Cosmos DB (la base de datos de Azure Cosmos) Flujos de trabajo de producción, distribuidos y entre procesos

Los tres implementan el mismo CheckpointStorage protocolo, por lo que puede intercambiar proveedores sin cambiar el código de flujo de trabajo o ejecutor.

InMemoryCheckpointStorage mantiene los puntos de control en la memoria del proceso. Ideal para pruebas, demostraciones y flujos de trabajo de corta duración en los que no necesita durabilidad en los reinicios.

from agent_framework import (
    InMemoryCheckpointStorage,
    WorkflowBuilder,
)

# Create a checkpoint storage to manage checkpoints
checkpoint_storage = InMemoryCheckpointStorage()

# Build a workflow with checkpointing enabled
builder = WorkflowBuilder(start_executor=start_executor, checkpoint_storage=checkpoint_storage)
builder.add_edge(start_executor, executor_b)
builder.add_edge(executor_b, executor_c)
builder.add_edge(executor_b, end_executor)
workflow = builder.build()

# Run the workflow
async for event in workflow.run(input, stream=True):
    ...

# Access checkpoints from the storage
checkpoints = await checkpoint_storage.list_checkpoints(workflow_name=workflow.name)

Para habilitar los puntos de control, configure el entorno de ejecución con un administrador de puntos de control. A continuación, se puede acceder a un punto de control desde workflow.SuperStepCompletedEvento a través de la lista de puntos de comprobación de la ejecución.

checkpointManager := checkpoint.NewInMemoryManager()

run, err := inproc.Default.
    WithCheckpointing(checkpointManager).
    RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer run.Close(ctx)

var checkpoints []workflow.CheckpointInfo
for evt, err := range run.WatchStream(ctx) {
    if err != nil {
        return err
    }
    if completed, ok := evt.(workflow.SuperStepCompletedEvent); ok && completed.CompletionInfo != nil {
        if completed.CompletionInfo.CheckpointInfo != nil {
            checkpoints = append(checkpoints, *completed.CompletionInfo.CheckpointInfo)
        }
    }
}

// Checkpoints can also be accessed from the run directly.
checkpoints = run.Checkpoints()

Reanudación desde puntos de control

Puede reanudar un flujo de trabajo desde un punto de control específico directamente en la misma ejecución.

// Assume we want to resume from the 6th checkpoint
CheckpointInfo savedCheckpoint = run.Checkpoints[5];
// Restore the state directly on the same run instance.
await run.RestoreCheckpointAsync(savedCheckpoint).ConfigureAwait(false);
await foreach (WorkflowEvent evt in run.WatchStreamAsync().ConfigureAwait(false))
{
    if (evt is WorkflowOutputEvent workflowOutputEvt)
    {
        Console.WriteLine($"Workflow completed with result: {workflowOutputEvt.Data}");
    }
}

Puede reanudar un flujo de trabajo desde un punto de control específico directamente en la misma instancia de flujo de trabajo.

# Assume we want to resume from the 6th checkpoint
saved_checkpoint = checkpoints[5]
async for event in workflow.run(checkpoint_id=saved_checkpoint.checkpoint_id, stream=True):
    ...

Puede restaurar una ejecución de streaming en un punto de control específico directamente en la misma ejecución.

// Assume we want to resume from the 6th checkpoint.
savedCheckpoint := checkpoints[5]
if err := run.RestoreCheckpoint(ctx, savedCheckpoint); err != nil {
    return err
}

for evt, err := range run.WatchStream(ctx) {
    if err != nil {
        return err
    }
    if outputEvent, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Workflow completed with result: %v\n", outputEvent.Output)
    }
}

Rehidratación desde puntos de control

Un flujo de trabajo rehidratado debe conservar la topología y las identidades de los ejecutores del flujo de trabajo que creó el punto de control. La forma en que se resuelve la identidad del ejecutor depende del SDK y del tipo de ejecutor.

O bien, puede rehidratar un flujo de trabajo desde un punto de control a una nueva instancia de ejecución.

// A rehydrated workflow must preserve the topology and executor identities of the workflow that
// created the checkpoint. This executor-only workflow rebuilds identically because its executors
// use fixed ids. Agent-based workflows must recreate each local agent with the same
// ChatClientAgentOptions.Id (and, if set, the same Name), otherwise the executor ids no longer
// match the checkpoint and resume fails.
var newWorkflow = WorkflowFactory.BuildWorkflow();
const int CheckpointIndex = 5;
Console.WriteLine($"\n\nHydrating a new workflow instance from the {CheckpointIndex + 1}th checkpoint.");
CheckpointInfo savedCheckpoint = checkpoints[CheckpointIndex];

await using StreamingRun newCheckpointedRun =
    await InProcessExecution.ResumeStreamingAsync(newWorkflow, savedCheckpoint, checkpointManager);

Important

El flujo de trabajo pasado a ResumeStreamingAsync debe tener la misma estructura e identidades del ejecutor que el flujo de trabajo que creó el punto de control. Si el flujo de trabajo contiene instancias locales ChatClientAgent que se reconstruyen entre solicitudes, ámbitos de inserción de dependencias, procesos o implementaciones, asigne a cada agente un elemento estable ChatClientAgentOptions.Id. Si un agente también establece un Name, deje ese Name sin cambios también.

Por ejemplo, asigne un identificador que represente el rol lógico del agente:

// Give each agent a stable, unique Id so its workflow executor identity stays the same when the
// workflow is reconstructed (for example per request or dependency-injection scope), which keeps
// checkpoints resumable. If an agent also has a Name, keep that stable too, since the executor
// identity includes it. Use a fixed logical role here, not a conversation, request, or user id.
internal const string IntakeAgentName = "Assistant";
public AIAgent IntakeAgent { get; } = chatClient.AsAIAgent(new ChatClientAgentOptions
{
    Id = "intake-agent",
    Name = IntakeAgentName,
    ChatOptions = new()
    {
        Instructions =
            """
            You receive a user request and are responsible for routing to the correct initial expert agent.
            """,
    },
});

Aplique este patrón a todos los agentes que participan en el flujo de trabajo. Los identificadores de agente deben ser únicos dentro del flujo de trabajo y deben reutilizarse al reconstruir el mismo agente lógico. No use identificadores de conversación, identificadores de solicitud, identificadores de usuario, información de identificación personal o secretos como identificadores de agente.

Cuando se establece un agente Name, la identidad actual del ejecutor de flujos de trabajo de .NET se obtiene a partir de sus Name y Id, por lo que cambiar cualquiera de estos valores hace que el flujo de trabajo vuelto a compilar sea incompatible con el punto de control. La asignación de valores estables no repara los puntos de control creados con identificadores diferentes o generados aleatoriamente; inicie una nueva sesión y el linaje del punto de comprobación en su lugar.

Para casos relacionados, consulte Flujos de trabajo como agentes y orquestación de transferencias.

O bien, puede rehidratar una nueva instancia de flujo de trabajo desde un punto de control.

from agent_framework import WorkflowBuilder

builder = WorkflowBuilder(start_executor=start_executor)
builder.add_edge(start_executor, executor_b)
builder.add_edge(executor_b, executor_c)
builder.add_edge(executor_b, end_executor)
# This workflow instance doesn't require checkpointing enabled.
workflow = builder.build()

# Assume we want to resume from the 6th checkpoint
saved_checkpoint = checkpoints[5]
async for event in workflow.run(
    checkpoint_id=saved_checkpoint.checkpoint_id,
    checkpoint_storage=checkpoint_storage,
    stream=True,
):
    ...

O bien, puede rehidratar una nueva instancia de flujo de trabajo desde un punto de control.

// Assume we want to resume from the 6th checkpoint
savedCheckpoint := checkpoints[5]
newWorkflow := buildWorkflow()

newRun, err := inproc.Default.
    WithCheckpointing(checkpointManager).
    ResumeStreaming(ctx, newWorkflow, savedCheckpoint)
if err != nil {
    return err
}
defer newRun.Close(ctx)

for evt, err := range newRun.WatchStream(ctx) {
    if err != nil {
        return err
    }
    if outputEvent, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Workflow completed with result: %v\n", outputEvent.Output)
    }
}

Guardar estados de ejecución

Para asegurarse de que el estado de un ejecutor se captura en un punto de control, el ejecutor debe invalidar el OnCheckpointingAsync método y guardar su estado en el contexto de flujo de trabajo.

using Microsoft.Agents.AI.Workflows;

internal sealed partial class CustomExecutor() : Executor("CustomExecutor")
{
    private const string StateKey = "CustomExecutorState";

    private List<string> messages = new();

    [MessageHandler]
    private async ValueTask HandleAsync(string message, IWorkflowContext context)
    {
        this.messages.Add(message);
        // Executor logic...
    }

    protected override ValueTask OnCheckpointingAsync(IWorkflowContext context, CancellationToken cancellation = default)
    {
        return context.QueueStateUpdateAsync(StateKey, this.messages);
    }
}

Además, para asegurarse de que el estado se restaura correctamente al reanudar desde un punto de control, el ejecutor debe invalidar el OnCheckpointRestoredAsync método y cargar su estado desde el contexto de flujo de trabajo.

protected override async ValueTask OnCheckpointRestoredAsync(IWorkflowContext context, CancellationToken cancellation = default)
{
    this.messages = await context.ReadStateAsync<List<string>>(StateKey).ConfigureAwait(false);
}

Para asegurarse de que el estado de un ejecutor se captura en un punto de control, el ejecutor debe invalidar el on_checkpoint_save método y devolver su estado como diccionario.

class CustomExecutor(Executor):
    def __init__(self, id: str) -> None:
        super().__init__(id=id)
        self._messages: list[str] = []

    @handler
    async def handle(self, message: str, ctx: WorkflowContext):
        self._messages.append(message)
        # Executor logic...

    async def on_checkpoint_save(self) -> dict[str, Any]:
        return {"messages": self._messages}

Además, para asegurarse de que el estado se restaura correctamente al reanudarse desde un punto de control, el ejecutor debe invalidar el on_checkpoint_restore método y restaurar su estado desde el diccionario de estado proporcionado.

async def on_checkpoint_restore(self, state: dict[str, Any]) -> None:
    self._messages = state.get("messages", [])

Para asegurarse de que el estado del ejecutor se captura en un punto de control, adjunte los enlaces de punto de control al ejecutor y almacene el estado a través del contexto de flujo de trabajo.

type customExecutor struct {
    messages []string
}

func (e *customExecutor) Handle(message string) {
    e.messages = append(e.messages, message)
}

func (e *customExecutor) OnCheckpoint(ctx *workflow.Context) error {
    return ctx.QueueStateUpdate("CustomExecutorState", "", slices.Clone(e.messages))
}

Restaure el estado en OnCheckpointRestoredFunc:

func (e *customExecutor) OnCheckpointRestored(ctx *workflow.Context) error {
    value, err := ctx.ReadState("CustomExecutorState", "")
    if err != nil {
        return err
    }
    if value == nil {
        e.messages = nil
        return nil
    }

    messages, ok := value.([]string)
    if !ok {
        return fmt.Errorf("unexpected custom executor state type %T", value)
    }
    e.messages = slices.Clone(messages)
    return nil
}

executorState := &customExecutor{}
custom := workflow.NewExecutor("CustomExecutor", executorState).Extend(&workflow.Executor{
    OnCheckpointFunc:         executorState.OnCheckpoint,
    OnCheckpointRestoredFunc: executorState.OnCheckpointRestored,
}).Bind()

Consideraciones de seguridad

Important

El almacenamiento de los puntos de control es un límite de confianza. Independientemente de si usa las implementaciones de almacenamiento integradas o una personalizada, el back-end de almacenamiento debe tratarse como infraestructura privada de confianza. Nunca cargue puntos de control de orígenes que no sean de confianza o potencialmente alterados.

Asegúrese de que la ubicación de almacenamiento usada para los puntos de control está protegida correctamente. Solo los servicios y usuarios autorizados deben tener acceso de lectura o escritura a los datos de punto de control.

Serialización pickle de Python

Tanto FileCheckpointStorage como CosmosCheckpointStorage usan el módulo pickle de Python para serializar el estado no nativo de JSON, como dataclasses, datetimes y objetos personalizados. Para mitigar los riesgos de ejecución arbitraria de código durante la deserialización, ambos proveedores usan un unpickler restringido de forma predeterminada. Solo se permite un conjunto integrado de tipos de Python seguros (primitivos, datetime, uuid, colecciones Decimalcomunes, etc.) y los tipos de SDK de Agent Framework o OpenAI compatibles durante la deserialización. La lista de prefijos de módulo permitidos solo admite tipos: se rechazan las funciones auxiliares y otros elementos globales que no sean tipos. Cualquier tipo no admitido hace que la deserialización falle con un WorkflowCheckpointException.

Para permitir tipos adicionales específicos de la aplicación, páselos mediante el parámetro allowed_checkpoint_types usando el formato "module:qualname".

from agent_framework import FileCheckpointStorage

storage = FileCheckpointStorage(
    "/tmp/checkpoints",
    allowed_checkpoint_types=[
        "my_app.models:SafeState",
        "my_app.models:UserProfile",
    ],
)

Cada entrada allowed_checkpoint_types debe corresponder a un tipo. Añadir una función a nivel de módulo u otro elemento global que no sea un tipo no hace que ese elemento global se pueda deserializar.

CosmosCheckpointStorage acepta el mismo parámetro:

from azure.identity.aio import DefaultAzureCredential
from agent_framework_azure_cosmos import CosmosCheckpointStorage

storage = CosmosCheckpointStorage(
    endpoint="https://my-account.documents.azure.com:443/",
    credential=DefaultAzureCredential(),
    database_name="agent-db",
    container_name="checkpoints",
    allowed_checkpoint_types=[
        "my_app.models:SafeState",
        "my_app.models:UserProfile",
    ],
)

Si su modelo de amenazas no permite en absoluto la serialización basada en pickle, use InMemoryCheckpointStorage o implemente una versión personalizada de CheckpointStorage con una estrategia de serialización alternativa.

Responsabilidad de ubicación de almacenamiento

FileCheckpointStorage requiere un parámetro explícito storage_path : no hay ningún directorio predeterminado. Aunque el marco se valida contra ataques de recorrido de rutas, proteger el propio directorio de almacenamiento (permisos de archivo, cifrado de datos en reposo, controles de acceso) es responsabilidad del desarrollador. Solo los procesos autorizados deben tener acceso de lectura o escritura al directorio de punto de control.

CosmosCheckpointStorage se basa en Azure Cosmos DB para el almacenamiento. Use la identidad administrada o RBAC siempre que sea posible, limite la base de datos y el contenedor al servicio de workflows y rote las claves de cuenta si usa la autenticación basada en claves. Como con el almacenamiento de archivos, solo los roles autorizados deben tener acceso de lectura o escritura al contenedor de Cosmos DB que contiene documentos de punto de control.

Los gestores de puntos de control de Go serializan el estado del punto de control como JSON, pero el almacenamiento de los puntos de control sigue siendo un estado de la aplicación de confianza. Si usa checkpoint.NewFileSystemJSONStore, almacene los archivos de punto de control en un directorio protegido y restrinja el acceso de lectura y escritura solo a los procesos autorizados. Los almacenes personalizados son responsables de sus propias garantías de control de acceso, integridad y durabilidad.

Pasos siguientes