Microsoft-Agent-Framework-Workflows – Prüfpunkte

Diese Seite bietet eine Übersicht über Checkpoints im Microsoft Agent Framework-Workflowsystem.

Überblick

Prüfpunkte ermöglichen es Ihnen, den Status eines Workflows an bestimmten Punkten während der Ausführung zu speichern und von diesen Punkten später fortzusetzen. Dieses Feature eignet sich besonders für die folgenden Szenarien:

  • Lang andauernde Workflows, bei denen Sie den Fortschrittsverlust im Falle von Fehlern vermeiden möchten.
  • Lang andauernde Workflows, bei denen Sie die Ausführung zu einem späteren Zeitpunkt anhalten und fortsetzen möchten.
  • Workflows, die regelmäßige Zustandsspeicherung für Überwachungs- oder Compliancezwecke erfordern.
  • Workflows, die in verschiedenen Umgebungen oder Instanzen migriert werden müssen.

Wann werden Prüfpunkte erstellt?

Denken Sie daran, dass Workflows in Supersteps ausgeführt werden, wie im Workflowausführungsmodell dokumentiert. Prüfpunkte werden am Ende jedes Supersteps erstellt, nachdem alle Ausführenden in diesem Superstep ihre Ausführung abgeschlossen haben. Ein Prüfpunkt erfasst den gesamten Status des Workflows, einschließlich:

  • Der aktuelle Status aller Executors
  • Alle ausstehenden Nachrichten im Workflow für den nächsten Superstep
  • Ausstehende Anforderungen und Antworten
  • Freigegebene Zustände

Note

Ab Python Version 1.13.0 erstellen Workflows auch einen Eintragsprüfpunkt vor dem ersten Superstep zum Aufzeichnen der Workfloweingabe und einen weiteren Eintragsprüfpunkt, wenn Antworten auf Anforderungsereignisse übermittelt werden. Mit diesen Prüfpunkten kann der vollständige Workflow wieder wiedergegeben werden. Diese Version enthält geringfügige Änderungen für Anwendungen, die von Iterationsanzahlen, Nachrichtenquell-IDs oder Prüfpunkt-Sortierung abhängen. Vorhandene Prüfpunkte werden weiterhin unterstützt. Details zur Migration finden Sie unter Upgrade Python Workflowprüfpunkte auf 1.13.0.

Erfassen von Prüfpunkten

Zum Aktivieren der Prüfpunktausführung muss beim Ausführen des Workflows ein CheckpointManager bereitgestellt werden. Auf einen Prüfpunkt kann dann über ein SuperStepCompletedEvent oder über die Checkpoints-Eigenschaft während der Ausführung zugegriffen werden.

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;

Um das Checkpointing zu aktivieren, muss beim Erstellen eines Workflows ein CheckpointStorage bereitgestellt werden. Über den Speicher kann dann auf einen Prüfpunkt zugegriffen werden. Agent Framework enthält drei integrierte Implementierungen – wählen Sie die Implementierung aus, die Ihren Anforderungen an Haltbarkeit und Bereitstellung entspricht:

Provider Package Beständigkeit Am besten geeignet für:
InMemoryCheckpointStorage agent-framework Nur in Bearbeitung Tests, Demos, kurzlebige Workflows
FileCheckpointStorage agent-framework Lokaler Datenträger Workflows mit einem Computer, lokale Entwicklung
CosmosCheckpointStorage agent-framework-azure-cosmos Azure Cosmos DB Produktion, verteilte, prozessübergreifende Workflows

Alle drei implementieren dasselbe CheckpointStorage Protokoll, sodass Sie Anbieter austauschen können, ohne Workflow- oder Executorcode zu ändern.

InMemoryCheckpointStorage hält Prüfpunkte im Prozessspeicher. Am besten geeignet für Tests, Demos und kurzlebige Workflows, bei denen Sie keine Haltbarkeit über Neustarts hinweg benötigen.

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)

Um Checkpointing zu aktivieren, konfigurieren Sie die Ausführungsumgebung mit einem Checkpoint-Manager. Ein Checkpoint kann dann über workflow.SuperStepCompletedEvent oder über die Checkpoint-Liste des Durchlaufs aufgerufen werden.

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

Fortsetzen ab Prüfpunkten

Sie können einen Workflow von einem bestimmten Checkpoint direkt im selben Ausführungslauf fortsetzen.

// 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}");
    }
}

Sie können einen Workflow aus einem bestimmten Prüfpunkt direkt in derselben Workflowinstanz fortsetzen.

# 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):
    ...

Sie können eine Streamingausführung direkt im selben Lauf auf einen bestimmten Prüfpunkt zurücksetzen.

// 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)
    }
}

Aktivierung aus Prüfpunkten

Ein rehydratierter Workflow muss die Topologie- und Executoridentitäten des Workflows beibehalten, der den Prüfpunkt erstellt hat. Die Auflösung der Executoridentität hängt vom SDK- und Executortyp ab.

Sie können einen Workflow auch aus einem Prüfpunkt heraus in eine neue Ausführungsinstanz aktivieren.

// 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);

Von Bedeutung

Der an den Workflow übergebene ResumeStreamingAsync Workflow muss die gleiche Struktur und Vollstreckungsidentität aufweisen wie der Workflow, der den Prüfpunkt erstellt hat. Wenn der Workflow lokale ChatClientAgent Instanzen enthält, die über Anforderungen, Abhängigkeitseinfügungsbereiche, Prozesse oder Bereitstellungen hinweg rekonstruiert werden, weisen Sie jedem Agent einen stabilen ChatClientAgentOptions.IdZu. Wenn ein Agent auch einen NameWert festlegt, behalten Sie dies Name ebenfalls unverändert.

Weisen Sie beispielsweise eine ID zu, die die logische Rolle des Agents darstellt:

// 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.
            """,
    },
});

Wenden Sie dieses Muster auf jeden Agent an, der am Workflow teilnimmt. Agent-IDs müssen innerhalb des Workflows eindeutig sein und beim Rekonstruieren desselben logischen Agents wiederverwendet werden. Verwenden Sie keine Unterhaltungs-IDs, Anforderungs-IDs, Benutzer-IDs, persönlich identifizierbare Informationen oder geheime Schlüssel als Agent-IDs.

Wenn ein Agent Name festgelegt ist, wird die aktuelle .NET Workflowausführeridentität sowohl von seinem als Idauch von ihm Name abgeleitet, sodass durch ändern eines der Werte der neu erstellte Workflow mit dem Prüfpunkt nicht kompatibel ist. Das Zuweisen stabiler Werte repariert keine Prüfpunkte, die mit unterschiedlichen oder zufällig generierten IDs erstellt wurden; starten Sie stattdessen eine neue Sitzungs- und Prüfpunktlinie.

Verwandte Szenarien finden Sie unter Workflows als Agents und Handoff-Orchestrierung.

Alternativ können Sie eine neue Workflowinstanz aus einem Prüfpunkt heraus aktivieren.

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,
):
    ...

Alternativ können Sie eine neue Workflowinstanz aus einem Prüfpunkt heraus aktivieren.

// 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)
    }
}

Executorstatus speichern

Um sicherzustellen, dass der Status eines Executors in einem Prüfpunkt erfasst wird, muss der Executor die OnCheckpointingAsync Methode überschreiben und den Status im Workflowkontext speichern.

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);
    }
}

Um sicherzustellen, dass der Zustand beim Fortsetzen aus einem Prüfpunkt ordnungsgemäß wiederhergestellt wird, muss der Executor die OnCheckpointRestoredAsync-Methode überschreiben und seinen Zustand aus dem Workflowkontext laden.

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

Um sicherzustellen, dass der Zustand eines Executors in einem Prüfpunkt erfasst wird, muss der Executor die on_checkpoint_save-Methode überschreiben und seinen Zustand als Dictionary zurückgeben.

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}

Um sicherzustellen, dass der Zustand beim Fortsetzen aus einem Prüfpunkt ordnungsgemäß wiederhergestellt wird, muss der Executor die on_checkpoint_restore-Methode überschreiben und seinen Zustand aus dem bereitgestellten Zustandsdictionary wiederherstellen.

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

Um sicherzustellen, dass der Status des Executors in einem Checkpoint erfasst wird, fügen Sie dem Executor Checkpoint-Hooks hinzu und speichern Sie den Status über den Workflow-Kontext.

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

Wiederherstellen des Zustands in 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()

Sicherheitsüberlegungen

Von Bedeutung

Der Prüfpunktspeicher ist eine Vertrauensgrenze. Unabhängig davon, ob Sie die integrierten Speicherimplementierungen oder eine benutzerdefinierte Implementierung verwenden, muss das Speicher-Back-End als vertrauenswürdige private Infrastruktur behandelt werden. Laden Sie niemals Prüfpunkte aus nicht vertrauenswürdigen oder potenziell manipulierten Quellen.

Stellen Sie sicher, dass der für Prüfpunkte verwendete Speicherort entsprechend gesichert ist. Nur autorisierte Dienste und Benutzer sollten Lese- oder Schreibzugriff auf Prüfpunktdaten haben.

Pickle-Serialisierung

Sowohl FileCheckpointStorage als auch CosmosCheckpointStorage verwenden das modul pickle Python, um nicht-JSON-nativen Zustand wie Datenklassen, Datetimes und benutzerdefinierte Objekte zu serialisieren. Um die Risiken einer willkürlichen Codeausführung während der Deserialisierung zu minimieren, verwenden beide Anbieter standardmäßig einen eingeschränkten Unpickler. Während der Deserialisierung sind nur eine integrierte Gruppe sicherer Python Typen (Grundtypen, datetimeuuid, , Decimalallgemeine Auflistungen usw.) und unterstützte Agent Framework- oder OpenAI SDK-Typen zulässig. Die Zulassungsliste mit Modulpräfix ist schreibgeschützt: Hilfsfunktionen und andere Nicht-Typ-Globalen werden abgelehnt. Jeder nicht unterstützte Typ führt dazu, dass die Deserialisierung mit einem WorkflowCheckpointExceptionFehler auftritt.

Um zusätzliche anwendungsspezifische Typen zuzulassen, übergeben Sie diese über den allowed_checkpoint_types-Parameter mithilfe des "module:qualname"-Formats.

from agent_framework import FileCheckpointStorage

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

Jeder allowed_checkpoint_types Eintrag muss in einen Typ aufgelöst werden. Das Hinzufügen einer Funktion auf Modulebene oder einer anderen globalen Nichttypfunktion macht diese globale Deserialisierbarkeit nicht möglich.

CosmosCheckpointStorage akzeptiert denselben Parameter:

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",
    ],
)

Wenn Ihr Bedrohungsmodell überhaupt keine pickle-basierte Serialisierung zulässt, verwenden Sie InMemoryCheckpointStorage oder implementieren Sie eine benutzerdefinierte CheckpointStorage mit einer alternativen Serialisierungsstrategie.

Verantwortung für den Speicherort

FileCheckpointStorage erfordert einen expliziten storage_path Parameter – es gibt kein Standardverzeichnis. Während das Framework gegen Pfad-Traversalangriffe validiert, liegt es in der Verantwortung des Entwicklers, das Speicherverzeichnis selbst (Dateiberechtigungen, Verschlüsselung im Ruhezustand, Zugriffssteuerungen) zu sichern. Nur autorisierte Prozesse sollten Lese- oder Schreibzugriff auf das Prüfpunktverzeichnis haben.

CosmosCheckpointStorage basiert auf Azure Cosmos DB für den Speicher. Verwenden Sie, wo möglich, verwaltete Identität/RBAC, begrenzen Sie den Geltungsbereich der Datenbank und des Containers auf den Workflow-Dienst und wechseln Sie die Kontoschlüssel, wenn Sie die schlüsselbasierte Authentifizierung verwenden. Wie bei der Dateispeicherung sollten nur autorisierte Entitäten Lese- oder Schreibzugriff auf den Cosmos DB-Container haben, der Checkpoint-Dokumente enthält.

Checkpoint-Manager in Go serialisieren den Checkpoint-Zustand als JSON, aber der Checkpoint-Speicher gilt weiterhin als vertrauenswürdiger Anwendungszustand. Wenn Sie checkpoint.NewFileSystemJSONStore verwenden, speichern Sie Prüfpunktdateien in einem geschützten Verzeichnis und beschränken Sie den Lese-/Schreibzugriff nur auf autorisierte Prozesse. Kundenspezifische Stores sind selbst für die Garantien hinsichtlich Zugriffssteuerung, Integrität und Dauerhaftigkeit verantwortlich.

Nächste Schritte