Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Ta strona zawiera omówienie Checkpoints w systemie przepływu pracy Microsoft Agent Framework.
Przegląd
Punkty kontrolne umożliwiają zapisywanie stanu procesu w określonych punktach podczas jego wykonywania i kontynuowanie od tych punktów później. Ta funkcja jest szczególnie przydatna w następujących scenariuszach:
- Długotrwałe przepływy pracy, w których chcesz uniknąć utraty postępu w przypadku awarii.
- Długotrwałe przepływy pracy, w których chcesz wstrzymać i wznowić wykonywanie w późniejszym czasie.
- Przepływy pracy, które wymagają okresowego zapisywania stanu na potrzeby inspekcji lub zgodności.
- Przepływy pracy, które należy migrować w różnych środowiskach lub instancjach.
Kiedy są tworzone punkty kontrolne?
Pamiętaj, że przepływy pracy są wykonywane w superkrokach, zgodnie z dokumentacją w modelu wykonywania przepływu pracy. Punkty kontrolne są tworzone na końcu każdego superkroku, po zakończeniu wykonywania wszystkich funkcji wykonawczych w tym superkroku. Punkt kontrolny przechwytuje cały stan przepływu pracy, w tym:
- Bieżący stan wszystkich funkcji wykonawczych
- Wszystkie oczekujące komunikaty w przepływie pracy dla następnego superkroku
- Oczekujące żądania i odpowiedzi
- Stany udostępnione
Note
Począwszy od Python wersji 1.13.0, przepływy pracy tworzą również punkt kontrolny wejścia przed pierwszym superkrokem w celu zarejestrowania danych wejściowych przepływu pracy i innego punktu kontrolnego wejścia, gdy są dostarczane odpowiedzi na żądania zdarzeń. Te punkty kontrolne sprawiają, że pełny przepływ pracy można odtworzyć. Ta wersja zawiera drobne zmiany powodujące niezgodność dla aplikacji, które zależą od liczby iteracji, identyfikatorów źródła komunikatów lub porządkowania punktów kontrolnych. Istniejące punkty kontrolne pozostają obsługiwane. Aby uzyskać szczegółowe informacje na temat migracji, zobacz Uaktualnianie punktów kontrolnych przepływu pracy Python do wersji 1.13.0.
Przechwytywanie punktów kontrolnych
Aby włączyć tworzenie punktów kontrolnych, należy podać element CheckpointManager podczas uruchamiania przepływu pracy. Następnie można uzyskać dostęp do punktu kontrolnego za pośrednictwem SuperStepCompletedEvent, lub poprzez właściwość Checkpoints w ramach uruchomienia.
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;
Aby włączyć tworzenie punktów kontrolnych, należy podać element CheckpointStorage podczas tworzenia przepływu pracy. Następnie można uzyskać dostęp do punktu kontrolnego za pośrednictwem magazynu. Struktura agenta dostarcza trzy wbudowane implementacje — wybierz jedną zgodną z potrzebami dotyczącymi trwałości i wdrażania:
| Dostawca | Pakiet | Durability | Najlepsze dla |
|---|---|---|---|
InMemoryCheckpointStorage |
agent-framework |
Tylko w trakcie procesu | Testy, pokazy, krótkotrwałe przepływy pracy |
FileCheckpointStorage |
agent-framework |
Dysk lokalny | Przepływy pracy z jedną maszyną, programowanie lokalne |
CosmosCheckpointStorage |
agent-framework-azure-cosmos |
Azure Cosmos DB | Przepływy pracy w środowisku produkcyjnym, rozproszonym i międzyprocesowym |
Wszystkie trzy implementują ten sam CheckpointStorage protokół, dzięki czemu można zamienić dostawców bez zmiany przepływu pracy lub kodu wykonawczego.
InMemoryCheckpointStorage program przechowuje punkty kontrolne w pamięci procesu. Idealne do testów, demonstracji i krótkich procesów roboczych, w których nie potrzebujesz trwałości w przypadku ponownych uruchomień.
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)
Aby włączyć tworzenie punktów kontrolnych, skonfiguruj środowisko wykonywania za pomocą menedżera punktów kontrolnych. Następnie można uzyskać dostęp do punktu kontrolnego z poziomu workflow.SuperStepCompletedEvent lub za pośrednictwem listy punktów kontrolnych uruchomienia.
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()
Wznawianie z punktów kontrolnych
Przepływ pracy można wznowić bezpośrednio z konkretnego punktu kontrolnego w tym samym przebiegu.
// 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}");
}
}
Możesz wznowić przepływ pracy z określonego punktu kontrolnego bezpośrednio w tym samym wystąpieniu.
# 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):
...
Możesz przywrócić przebieg przesyłania strumieniowego do określonego punktu kontrolnego bezpośrednio w tym samym przebiegu.
// 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)
}
}
Przywracanie z punktów kontrolnych
Przepływ pracy z ponownym wypełnianiem musi zachować tożsamość topologii i funkcji wykonawczej przepływu pracy, który utworzył punkt kontrolny. Sposób rozpoznawania tożsamości funkcji wykonawczej zależy od zestawu SDK i typu funkcji wykonawczej.
Możesz też przywrócić przepływ pracy z punktu kontrolnego do nowego wystąpienia uruchomienia.
// 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);
Ważna
Przekazany przepływ pracy ResumeStreamingAsync musi mieć taką samą strukturę i tożsamości funkcji wykonawczej jak przepływ pracy, który utworzył punkt kontrolny. Jeśli przepływ pracy zawiera wystąpienia lokalne ChatClientAgent , które są odtwarzane między żądaniami, zakresami iniekcji zależności, procesami lub wdrożeniami, przypisz każdemu agentowi stabilną wartość ChatClientAgentOptions.Id. Jeśli agent ustawia Namerównież wartość , zachowaj to Name bez zmian.
Na przykład przypisz identyfikator reprezentujący rolę logiczną agenta:
// 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.
""",
},
});
Zastosuj ten wzorzec do każdego agenta, który uczestniczy w przepływie pracy. Identyfikatory agentów muszą być unikatowe w przepływie pracy i muszą być ponownie używane podczas odbudowy tego samego agenta logicznego. Nie używaj identyfikatorów konwersacji, identyfikatorów żądań, identyfikatorów użytkowników, danych osobowych lub wpisów tajnych jako identyfikatorów agentów.
Po ustawieniu agenta Name bieżąca tożsamość funkcji wykonawczej przepływu pracy .NET pochodzi zarówno z jejName, jak i Id, dlatego zmiana jednej z wartości powoduje, że przebudowany przepływ pracy jest niezgodny z punktem kontrolnym. Przypisywanie stabilnych wartości nie naprawia punktów kontrolnych utworzonych przy użyciu różnych lub losowo wygenerowanych identyfikatorów; Zamiast tego uruchom nową sesję i pochodzenie punktów kontrolnych.
W przypadku powiązanych scenariuszy zobacz Przepływy pracy jako agenci i aranżacja handoff.
Możesz też przywrócić nowe wystąpienie przepływu pracy z punktu kontrolnego.
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,
):
...
Możesz też przywrócić nowe wystąpienie przepływu pracy z punktu kontrolnego.
// 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)
}
}
Zapisz stany wykonawcze
Aby upewnić się, że stan funkcji wykonawczej jest przechwytywany w punkcie OnCheckpointingAsync kontrolnym, funkcja wykonawcza musi zastąpić metodę i zapisać jej stan w kontekście przepływu pracy.
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);
}
}
Ponadto, aby zapewnić prawidłowe przywrócenie stanu podczas wznawiania z punktu kontrolnego, egzekutor musi zastąpić OnCheckpointRestoredAsync metodę i załadować jego stan z kontekstu przepływu pracy.
protected override async ValueTask OnCheckpointRestoredAsync(IWorkflowContext context, CancellationToken cancellation = default)
{
this.messages = await context.ReadStateAsync<List<string>>(StateKey).ConfigureAwait(false);
}
Aby upewnić się, że stan funkcji wykonawczej jest przechwytywany w punkcie kontrolnym, funkcja wykonawcza musi zastąpić on_checkpoint_save metodę i zwrócić jej stan jako słownik.
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}
Ponadto, aby upewnić się, że stan jest poprawnie przywracany podczas wznawiania z punktu kontrolnego, funkcja wykonawcza musi zastąpić metodę on_checkpoint_restore i przywrócić jej stan z dostarczonego słownika stanu.
async def on_checkpoint_restore(self, state: dict[str, Any]) -> None:
self._messages = state.get("messages", [])
Aby upewnić się, że stan egzekutora zostanie zapisany w punkcie kontrolnym, podłącz hooki punktu kontrolnego do egzekutora i zapisz stan za pomocą kontekstu przepływu pracy.
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))
}
Przywróć stan w pliku 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()
Zagadnienia związane z zabezpieczeniami
Ważna
Przechowywanie punktów kontrolnych jest granicą zaufania. Niezależnie od tego, czy używasz wbudowanych implementacji magazynu, czy niestandardowej, zaplecze magazynu musi być traktowane jako zaufana, prywatna infrastruktura. Nigdy nie ładuj punktów kontrolnych z niezaufanych lub potencjalnie naruszonych źródeł.
Upewnij się, że lokalizacja magazynu używana dla punktów kontrolnych jest odpowiednio zabezpieczona. Tylko autoryzowane usługi i użytkownicy powinni mieć dostęp do odczytu lub zapisu do danych punktu kontrolnego.
Serializacja Pickle
Zarówno FileCheckpointStorage, jak i CosmosCheckpointStorage używają modułu pickle Python do serializacji stanu natywnego innego niż JSON, takich jak klasy danych, daty/godziny i obiekty niestandardowe. Aby ograniczyć ryzyko dowolnego wykonania kodu podczas deserializacji, obydwaj dostawcy domyślnie używają ograniczonego unpicklera. Podczas deserializacji dozwolony jest tylko wbudowany zestaw bezpiecznych typów Python (pierwotnych, datetime, , uuidDecimal, wspólnych kolekcji itp.) i obsługiwanych typów platformy Agent Framework lub zestawu OPENAI SDK. Lista dozwolonych prefiksów modułów jest tylko typem: funkcje pomocnicze i inne nietypowe globalne są odrzucane. Każdy nieobsługiwany typ powoduje niepowodzenie deserializacji z elementem WorkflowCheckpointException.
Aby zezwolić na dodatkowe typy specyficzne dla aplikacji, przekaż je za pośrednictwem parametru allowed_checkpoint_types przy użyciu "module:qualname" formatu:
from agent_framework import FileCheckpointStorage
storage = FileCheckpointStorage(
"/tmp/checkpoints",
allowed_checkpoint_types=[
"my_app.models:SafeState",
"my_app.models:UserProfile",
],
)
Każdy allowed_checkpoint_types wpis musi być rozpoznawany jako typ. Dodanie funkcji na poziomie modułu lub innej nietypowej globalnej nie sprawia, że ta globalna deserializowana.
CosmosCheckpointStorage akceptuje ten sam parametr:
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",
],
)
Jeśli model zagrożeń w ogóle nie zezwala na serializację opartą na module pickle, użyj InMemoryCheckpointStorage lub zaimplementuj niestandardowy CheckpointStorage z alternatywną strategią serializacji.
Odpowiedzialność za lokalizację magazynu
FileCheckpointStorage wymaga jawnego storage_path parametru — nie ma katalogu domyślnego. Chociaż struktura chroni przed atakami typu path traversal, zabezpieczenie samego katalogu magazynu (uprawnienia do plików, szyfrowanie w spoczynku, kontrolę dostępu) jest odpowiedzialnością dewelopera. Tylko autoryzowane procesy powinny mieć dostęp do odczytu lub zapisu do katalogu punktu kontrolnego.
CosmosCheckpointStorage opiera się na Azure Cosmos DB do przechowywania. Użyj tożsamości zarządzanej/kontroli dostępu opartej na rolach, jeśli to możliwe, określ zakres bazy danych i kontenera dla usługi przepływu pracy oraz rotuj klucze konta, jeśli używasz uwierzytelniania opartego na kluczach. Podobnie jak w przypadku magazynu plików, tylko autoryzowane podmioty powinny mieć dostęp do odczytu lub zapisu do kontenera Cosmos DB, który przechowuje dokumenty punktów kontrolnych.
Menedżerowie punktów kontrolnych języka Go serializują stan punktu kontrolnego jako JSON, ale magazyn punktów kontrolnych jest nadal zaufanym stanem aplikacji. Jeśli używasz programu checkpoint.NewFileSystemJSONStore, zapisz pliki punktu kontrolnego w chronionym katalogu i ogranicz dostęp do odczytu/zapisu tylko do autoryzowanych procesów. Magazyny niestandardowe są odpowiedzialne za własną kontrolę dostępu, integralność i trwałość gwarancji.