Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Tato stránka obsahuje přehled Checkpoints v systému pracovních postupů Microsoft Agent Framework.
Přehled
Kontrolní body umožňují uložit stav pracovního postupu v konkrétních bodech během jeho provádění a pokračovat z těchto bodů později. Tato funkce je užitečná zejména pro následující scénáře:
- Dlouhotrvající pracovní postupy, ve kterých chcete zabránit ztrátě průběhu v případě selhání.
- Dlouhotrvající pracovní postupy, ve kterých chcete pozastavit a obnovit spouštění později.
- Pracovní postupy, které vyžadují pravidelné ukládání stavu pro účely auditování nebo dodržování předpisů
- Pracovní postupy, které je potřeba migrovat napříč různými prostředími nebo instancemi
Kdy jsou vytvořeny kontrolní body?
Mějte na paměti, že pracovní postupy se provádějí v superkrocích, jak je uvedeno v modelu provádění pracovního postupu. Kontrolní body se vytvoří na konci každého superkroku poté, co všechny exekutory v daném superkroku dokončily provádění. Kontrolní bod zachycuje celý stav pracovního postupu, včetně:
- Aktuální stav všech exekutorů
- Všechny čekající zprávy ve workflow pro další superstep
- Čekající žádosti a odpovědi
- Sdílené stavy
Note
Počínaje Python verze 1.13.0 pracovní postupy také vytvoří vstupní kontrolní bod před prvním superkrokem pro zaznamenání vstupu pracovního postupu a další vstupní kontrolní bod při doručení odpovědí na události požadavku. Díky těmto kontrolním bodům se celý pracovní postup dá přehrát znovu. Tato verze obsahuje menší zásadní změny pro aplikace, které závisí na počtu iterací, ID zdroje zpráv nebo řazení kontrolních bodů. Stávající kontrolní body zůstanou podporovány. Podrobnosti o migraci najdete v tématu Upgrade Python kontrolních bodů pracovního postupu na verzi 1.13.0.
Zachytávání kontrolních bodů
Pokud chcete povolit vytváření kontrolních bodů, je potřeba poskytnout CheckpointManager při spuštění pracovního postupu. K bodu obnovení je pak možné přistupovat prostřednictvím SuperStepCompletedEvent nebo vlastnosti Checkpoints při spuště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;
Pro povolení vytváření kontrolních bodů je třeba poskytnout CheckpointStorage při vytváření pracovního postupu. K kontrolnímu bodu se pak dostanete přes úložiště. Agent Framework dodává tři předdefinované implementace – vyberte si ty, které odpovídají vašim potřebám odolnosti a nasazení:
| Poskytovatel | Balíček | Durability | Nejvhodnější pro |
|---|---|---|---|
InMemoryCheckpointStorage |
agent-framework |
Pouze během procesu | Testy, ukázky, krátkodobé pracovní postupy |
FileCheckpointStorage |
agent-framework |
Místní disk | Pracovní postupy s jedním počítačem, místní vývoj |
CosmosCheckpointStorage |
agent-framework-azure-cosmos |
Azure Cosmos DB | Produkční, distribuované pracovní postupy napříč procesy |
Všechny tři implementují stejný CheckpointStorage protokol, takže můžete prohodit poskytovatele beze změny kódu pracovního postupu nebo exekutoru.
InMemoryCheckpointStorage uchovává kontrolní body v paměti procesu. Nejvhodnější pro testy, ukázky a krátkodobé pracovní postupy, ve kterých nepotřebujete odolnost při restartování.
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)
Pokud chcete povolit vytváření kontrolních bodů, nakonfigurujte spouštěcí prostředí pomocí správce kontrolních bodů. Ke kontrolnímu bodu pak lze přistupovat z workflow.SuperStepCompletedEvent nebo prostřednictvím seznamu kontrolních bodů běhu.
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()
Obnovení z kontrolních bodů
Pracovní postup můžete obnovit z konkrétního kontrolního bodu přímo v tom samém běhu.
// 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}");
}
}
Pracovní postup můžete obnovit z konkrétního kontrolního bodu přímo ve stejné instanci pracovního postupu.
# 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):
...
Běh streamování můžete obnovit z konkrétního kontrolního bodu přímo v rámci stejného běhu.
// 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)
}
}
Obnovení z kontrolních bodů
Rehydrovaný pracovní postup musí zachovat topologii a identity exekutoru pracovního postupu, který vytvořil kontrolní bod. Způsob řešení identity exekutoru závisí na typu sady SDK a exekutoru.
Nebo můžete pracovní postup obnovit z kontrolního bodu do nové spustitelné instance.
// 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);
Důležité
Pracovní postup předaný ResumeStreamingAsync musí mít stejnou strukturu a identity exekutoru jako pracovní postup, který vytvořil kontrolní bod. Pokud pracovní postup obsahuje místní ChatClientAgent instance, které jsou rekonstruovány napříč požadavky, obory injektáže závislostí, procesy nebo nasazení, přiřaďte každému agentu stabilní ChatClientAgentOptions.Id. Pokud agent také nastaví Name, ponechte to Name beze změny.
Například přiřaďte ID, které představuje logickou roli 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.
""",
},
});
Tento vzor použijte u každého agenta, který se účastní pracovního postupu. ID agenta musí být v rámci pracovního postupu jedinečná a musí se znovu použít při rekonstrukci stejného logického agenta. Nepoužívejte ID konverzací, ID žádostí, ID uživatelů, identifikovatelné osobní údaje ani tajné kódy jako ID agenta.
Když je agent Name nastaven, aktuální .NET identita exekutoru pracovního postupu je odvozena z jeho Name a Id, takže změna obou hodnot způsobí, že znovu sestavený pracovní postup není kompatibilní s kontrolním bodem. Přiřazení stabilních hodnot neopravuje kontrolní body vytvořené s různými nebo náhodně generovanými ID; místo toho spusťte novou relaci a rodokmen kontrolních bodů.
Související scénáře najdete v tématu Pracovní postupy jako agenti a orchestrace předání.
Nebo můžete obnovit novou instanci pracovního postupu z kontrolního bodu.
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,
):
...
Nebo můžete obnovit novou instanci pracovního postupu z kontrolního bodu.
// 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)
}
}
Stavy exekutoru
Aby se zajistilo, že se stav exekutoru zachytí do kontrolního bodu, musí exekutor přepsat OnCheckpointingAsync metodu a uložit svůj stav do kontextu pracovního postupu.
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);
}
}
Aby se zajistilo, že se stav při obnovení z kontrolního bodu správně obnoví, musí exekutor přepsat OnCheckpointRestoredAsync metodu a načíst její stav z kontextu pracovního postupu.
protected override async ValueTask OnCheckpointRestoredAsync(IWorkflowContext context, CancellationToken cancellation = default)
{
this.messages = await context.ReadStateAsync<List<string>>(StateKey).ConfigureAwait(false);
}
Aby se zajistilo, že se stav exekutoru zachytí do kontrolního bodu, musí exekutor přepsat on_checkpoint_save metodu a vrátit její stav jako slovník.
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}
K zajištění správného obnovení stavu při obnově z kontrolního bodu musí exekutor přepsat on_checkpoint_restore metodu a obnovit stav ze zadaného slovníku stavu.
async def on_checkpoint_restore(self, state: dict[str, Any]) -> None:
self._messages = state.get("messages", [])
Chcete-li zajistit, aby byl stav executoru zachycen v checkpointu, připojte k executoru checkpoint hooky a stav ukládejte prostřednictvím kontextu workflow.
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))
}
Obnovení stavu v 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()
Aspekty zabezpečení
Důležité
Úložiště kontrolních bodů je hranice důvěryhodnosti. Bez ohledu na to, jestli používáte integrované implementace úložiště nebo vlastní, musí být back-end úložiště považován za důvěryhodnou privátní infrastrukturu. Nikdy nenačítejte kontrolní body z nedůvěryhodných nebo potenciálně manipulovaných zdrojů.
Ujistěte se, že je umístění úložiště používané pro kontrolní body správně zabezpečené. Pouze autorizované služby a uživatelé by měli mít přístup ke čtení nebo zápisu k datům kontrolního bodu.
Serializace Pickle
Modul FileCheckpointStorage i CosmosCheckpointStorage používá modul Python pickle k serializaci jiného než nativního stavu JSON, jako jsou datové třídy, data a časy a vlastní objekty. Aby se zmírnila rizika spuštění libovolného kódu během deserializace, oba poskytovatelé ve výchozím nastavení používají omezený unpickler. Při deserializaci jsou povoleny pouze předdefinované sady bezpečných typů Python (primitiv, datetime, uuid, Decimalběžné kolekce atd.) a podporované typy sady Agent Framework nebo OpenAI SDK. Seznam povolených předpon modulů je pouze typ: pomocné funkce a jiné globální globální hodnoty jiného typu než typ jsou odmítnuty. Jakýkoli nepodporovaný typ způsobí, že deserializace selže s chybou WorkflowCheckpointException.
Pokud chcete povolit další typy specifické pro aplikaci, předejte je prostřednictvím parametru allowed_checkpoint_types pomocí "module:qualname" formátu:
from agent_framework import FileCheckpointStorage
storage = FileCheckpointStorage(
"/tmp/checkpoints",
allowed_checkpoint_types=[
"my_app.models:SafeState",
"my_app.models:UserProfile",
],
)
Každá allowed_checkpoint_types položka se musí přeložit na typ. Přidání funkce na úrovni modulu nebo jiného netypového globálního globálního objektu neumožňuje globální deserializovatelnost.
CosmosCheckpointStorage přijímá stejný 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",
],
)
Pokud váš model hrozeb nepovoluje serializaci založenou na pickle, použijte InMemoryCheckpointStorage nebo implementujte vlastní CheckpointStorage s alternativní strategií serializace.
Odpovědnost za umístění úložiště
FileCheckpointStorage vyžaduje explicitní storage_path parametr – neexistuje žádný výchozí adresář. I když rámec ověřuje útoky proti procházení cest, zabezpečení samotného adresáře úložiště (oprávnění k souborům, šifrování dat v klidu, řízení přístupu) je zodpovědností vývojáře. K adresáři kontrolních bodů by měly mít přístup jen autorizované procesy pro čtení nebo zápis.
CosmosCheckpointStorage spoléhá na úložiště Azure Cosmos DB. Pokud je to možné, použijte spravovanou identitu nebo RBAC, omezte rozsah databáze a kontejneru na službu pracovního postupu a měňte klíče účtu pravidelně, pokud používáte ověřování založené na klíčích. Stejně jako u úložiště souborů by k kontejneru Cosmos DB, který obsahuje dokumenty kontrolních bodů, měly mít přístup jen autorizované subjekty.
Správci kontrolních bodů v jazyce Go serializují stav kontrolních bodů do formátu JSON, ale úložiště kontrolních bodů je stále považováno za důvěryhodný stav aplikace. Pokud používáte checkpoint.NewFileSystemJSONStore, uložte soubory kontrolních bodů do chráněného adresáře a omezte přístup pro čtení a zápis pouze na autorizované procesy. Vlastní datová úložiště si sama zajišťují řízení přístupu, integritu a záruky trvalosti dat.