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 punktów kontrolnych w systemie przepływu pracy programu Microsoft Agent Framework.
Overview
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 superetapach, jak opisano w modelu wykonywania przepływów 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
Uwaga / Notatka
Począwszy od wersji 1.13.0 języka Python przepływy pracy tworzą również początkowy punkt kontrolny przed pierwszym superkrokiem, aby zarejestrować dane wejściowe przepływu pracy, a także kolejny początkowy punkt kontrolny, gdy dostarczane są odpowiedzi na zdarzenia żądań. Te punkty kontrolne sprawiają, że pełny przepływ pracy można odtworzyć. Ta wersja zawiera drobne niekompatybilne zmiany w przypadku aplikacji, które zależą od liczby iteracji, identyfikatorów źródeł komunikatów lub kolejności punktów kontrolnych. Istniejące punkty kontrolne są nadal 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:
| Provider | Package | 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
Odtworzony przepływ pracy musi zachować topologię i tożsamości egzekutorów 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);
Important
Przepływ pracy przekazany do ResumeStreamingAsync musi mieć taką samą strukturę i identyfikatory wykonawców 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 również Name, pozostaw również 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 ani poufnych informacji jako identyfikatorów agentów.
Gdy agent Name jest ustawiony, tożsamość bieżącego wykonawcy przepływu pracy platformy .NET jest wyprowadzana zarówno z jego Name, jak i Id, więc zmiana którejkolwiek z tych wartości powoduje, że ponownie skompilowany 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.
Informacje o powiązanych scenariuszach znajdują się w artykułach Przepływy pracy jako agenci i Orkiestracja przekazywania.
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
Important
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 dozwolone są wyłącznie wbudowany zestaw bezpiecznych typów Pythona (typów prostych, datetime, uuid, Decimal, typowych kolekcji itp.) oraz typy obsługiwane przez Agent Framework lub OpenAI SDK. Lista dozwolonych prefiksów modułu dotyczy wyłącznie typów: funkcje pomocnicze i inne globalne elementy niebędące typami są odrzucane.
Struktura agenta pomija wartości przejściowe raw_representation z obiektów natywnych dla platformy przed wybraniem i przywróceniem tych pól jako None. Każdy inny nieobsługiwany typ powoduje niepowodzenie deserializacji z parametrem 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 wpis allowed_checkpoint_types musi wskazywać na typ. Dodanie funkcji na poziomie modułu lub innego globalnego obiektu, który nie jest typem, nie sprawia, że ten obiekt globalny można deserializować.
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",
],
)
Aby zarejestrować typ dla każdego dekodera ograniczonego punktu kontrolnego w bieżącym procesie Python, przekaż klasę do register_checkpoint_type():
from agent_framework import register_checkpoint_type
from my_app.models import SafeState
register_checkpoint_type(SafeState)
Rejestracja globalna ma również zastosowanie do instancji magazynu utworzonych przed wywołaniem.
Natomiast allowed_checkpoint_types ma zastosowanie tylko do instancji FileCheckpointStorage lub CosmosCheckpointStorage, która go otrzymuje.
Preferuj parametr dla magazynu, gdy kod tworzy magazyn. Użyj rejestracji globalnej, gdy kod nie może skonfigurować instancji magazynu, na przykład w środowisku hostowanym, które tworzy za Ciebie magazyn punktów kontrolnych. Zarejestruj typ przed załadowaniem punktu kontrolnego zawierającego go.
Oba mechanizmy rozszerzają listę dozwolonych elementów ograniczonego unpicklera; żaden z nich nie sprawia, że niezaufane dane pickle są bezpieczne. Rejestruj tylko zaufane klasy aplikacji, ponieważ dopuszczona klasa może definiować niestandardowy sposób odtwarzania obiektu przy użyciu mechanizmu pickle. Globalna rejestracja zwiększa powierzchnię deserializacji dla każdego ograniczonego obciążenia punktu kontrolnego w procesie, więc zachowaj rejestr globalny do minimalnego wymaganego zestawu.
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.