Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfa, Microsoft Agent Framework İş Akışı sistemindeki Denetim Noktalarına genel bir bakış sağlar.
Overview
Denetim noktaları, yürütme sırasında belirli noktalarda iş akışının durumunu kaydetmenize ve daha sonra bu noktalardan devam etmenizi sağlar. Bu özellik özellikle aşağıdaki senaryolar için kullanışlıdır:
- Uzun süre devam eden ve hata durumunda ilerlemeyi kaybetmek istemediğiniz iş akışları.
- Yürütmeyi daha sonra duraklatmak ve devam ettirmek istediğiniz uzun süreli iş akışları.
- Denetim veya uyumluluk amacıyla düzenli olarak durum kaydetme gerektiren iş akışları.
- Farklı ortamlar veya örnekler arasında taşınması gereken iş akışları.
Denetim Noktaları Ne Zaman Oluşturulur?
Unutmayın, iş akışları, iş akışı yürütme modelinde belgelendiği gibi, süper adımlar halinde yürütülür. Denetim noktaları, bu üst adımdaki tüm yürütücüler yürütmelerini tamamladıktan sonra her üst adımın sonunda oluşturulur. Denetim noktası iş akışının durumunun tamamını yakalar, örneğin:
- Tüm yürütücülerin geçerli durumu
- Sonraki üst adım için iş akışında bekleyen tüm iletiler
- Bekleyen istekler ve yanıtlar
- Paylaşılan durumlar
Uyarı
Python sürüm 1.13.0'dan başlayarak iş akışları, iş akışı girişini kaydetmek için ilk üst adımdan önce bir giriş denetim noktası ve istek olaylarına yanıtlar teslim edildiğinde başka bir giriş denetim noktası da oluşturur. Bu denetim noktaları, iş akışının tamamının yeniden oynatılabilir olmasını sağlar. Bu sürüm, yineleme sayılarına, ileti kaynağı kimliklerine veya denetim noktalarının sıralamasına bağlı uygulamalar için küçük çaplı uyumluluğu bozan değişiklikler içerir. Mevcut denetim noktaları desteklenmeye devam eder. Geçiş ayrıntıları için bkz. Python iş akışı denetim noktalarını 1.13.0'a yükseltme.
Denetim Noktalarını Yakalama
Denetim noktası oluşturmayı etkinleştirmek için iş akışı çalıştırılırken bir CheckpointManager sağlanması gerekir. Denetim noktasına SuperStepCompletedEvent aracılığıyla veya çalıştırmada Checkpoints özelliği aracılığıyla erişilebilir.
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;
İş akışı oluşturulurken denetim noktası oluşturmayı etkinleştirmek için bir CheckpointStorage sağlanması gerekir. Daha sonra bir denetim noktasına depolama alanı üzerinden erişilebilir. Agent Framework üç yerleşik uygulama sunar; dayanıklılık ve dağıtım gereksinimlerinize uyan uygulamayı seçin:
| Provider | Paket | Durability | En iyi kullanım alanları |
|---|---|---|---|
InMemoryCheckpointStorage |
agent-framework |
Yalnızca işlem içi | Testler, tanıtımlar, kısa süreli iş akışları |
FileCheckpointStorage |
agent-framework |
Yerel diske | Tek makineli iş akışları, yerel geliştirme |
CosmosCheckpointStorage |
agent-framework-azure-cosmos |
Azure Cosmos DB veritabanı | Üretim, dağıtılmış, işlemler arası iş akışları |
Her üçü de aynı CheckpointStorage protokolü uygular, böylece iş akışını veya yürütücü kodunu değiştirmeden sağlayıcıları değiştirebilirsiniz.
InMemoryCheckpointStorage denetim noktalarını işlem belleğinde tutar. Testler, tanıtımlar ve yeniden başlatmalar arasında dayanıklılığa ihtiyaç duymadığınız kısa süreli iş akışları için idealdir.
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)
Denetim noktası oluşturmayı etkinleştirmek için yürütme ortamını bir denetim noktası yöneticisiyle yapılandırın. Daha sonra bir denetim noktasına workflow.SuperStepCompletedEvent üzerinden veya çalıştırmanın denetim noktaları listesi üzerinden erişilebilir.
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()
Kontrol Noktalarından Devam Etme
Belirli bir denetim noktasından bir iş akışını aynı çalışmada doğrudan sürdürebilirsiniz.
// 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}");
}
}
Belirli bir denetim noktasından bir iş akışını doğrudan aynı iş akışı örneğinde sürdürebilirsiniz.
# 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):
...
Bir akış çalıştırmasını doğrudan aynı çalıştırmada belirli bir denetim noktasına geri yükleyebilirsiniz.
// 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)
}
}
Denetim Noktalarından Yeniden Doldurma
Yeniden oluşturulan bir iş akışı, kontrol noktasını oluşturan iş akışının topolojisini ve yürütücü kimliklerini korumalıdır. Yürütücü kimliğinin çözümlenme şekli SDK'ya ve yürütücü türüne bağlıdır.
Veya bir iş akışını bir denetim noktasından yeni bir çalıştırma örneğine yeniden yükleyebilirsiniz.
// 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);
Önemli
geçirilen ResumeStreamingAsync iş akışının, denetim noktasını oluşturan iş akışıyla aynı yapıya ve yürütücü kimliklerine sahip olması gerekir. İş akışı, istekler, bağımlılık enjeksiyonu kapsamları, işlemler veya dağıtımlar arasında yeniden oluşturulan yerel ChatClientAgent örnekleri içeriyorsa, her aracıya kararlı bir ChatClientAgentOptions.Id atayın. Bir aracı da bir Name ayarlarsa, Name öğesini de aynı şekilde bırakın.
Örneğin, ajanın mantıksal rolünü temsil eden bir kimlik tanımlayıcısı atayın:
// 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.
""",
},
});
Bu düzeni iş akışına katılan her aracıya uygulayın. Aracı kimlikleri iş akışı içinde benzersiz olmalı ve aynı mantıksal aracı yeniden oluştururken yeniden kullanılmalıdır. Konuşma kimliklerini, istek kimliklerini, kullanıcı kimliklerini, kişisel olarak tanımlanabilir bilgileri veya gizli bilgileri aracı kimliği olarak kullanmayın.
Bir aracı Name ayarlandığında, geçerli .NET iş akışı yürütücüsü kimliği hem Name hem de Id değerlerinden türetilir; bu nedenle bu değerlerden herhangi birinin değiştirilmesi, yeniden oluşturulan iş akışını denetim noktasıyla uyumsuz hale getirir. Kararlı değerlerin atanması, farklı veya rastgele oluşturulan kimliklerle oluşturulan denetim noktalarını onarmaz; bunun yerine yeni bir oturum başlatın ve denetim noktası kökenini kullanın.
İlgili senaryolar için bkz. Aracı olarak iş akışları ve İletim düzenlemesi.
Veya bir kontrol noktasından yeni bir iş akışı örneğini yeniden başlatabilirsiniz.
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,
):
...
Veya bir kontrol noktasından yeni bir iş akışı örneğini yeniden başlatabilirsiniz.
// 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)
}
}
Yürütücü Durumlarını Kaydet
Bir yürütücü durumunun bir denetim noktasında yakalanmasını sağlamak için, yürütücü OnCheckpointingAsync yöntemini geçersiz kılmalı ve durumunu iş akışı bağlamına kaydetmelidir.
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);
}
}
Ayrıca, denetim noktasından devam ederken durumun doğru şekilde geri yüklendiğinden emin olmak için yürütücü yöntemi geçersiz kılmalı OnCheckpointRestoredAsync ve iş akışı bağlamından durumunu yüklemelidir.
protected override async ValueTask OnCheckpointRestoredAsync(IWorkflowContext context, CancellationToken cancellation = default)
{
this.messages = await context.ReadStateAsync<List<string>>(StateKey).ConfigureAwait(false);
}
Bir yürütücünün durumunun bir kontrol noktasında yakalandığından emin olmak için, yürütücü on_checkpoint_save yöntemini geçersiz kılmalı ve durumunu bir sözlük olarak döndürmelidir.
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}
Ayrıca, denetim noktasından devam ederken durumun doğru şekilde geri yüklendiğinden emin olmak için yürütücü yöntemi on_checkpoint_restore geçersiz kılmalı ve sağlanan durum sözlüğünden durumunu geri yüklemelidir.
async def on_checkpoint_restore(self, state: dict[str, Any]) -> None:
self._messages = state.get("messages", [])
Yürütücü durumunun bir denetim noktasında yakalandığından emin olmak için denetim noktası kancalarını yürütücüye ekleyin ve iş akışı bağlamı aracılığıyla durumu depolayın.
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))
}
OnCheckpointRestoredFunc içindeki durumu geri yükleyin:
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()
Güvenlikle İlgili Dikkat Edilmesi Gerekenler
Önemli
Denetim noktası depolaması bir güven sınırıdır. İster yerleşik depolama uygulamalarını ister özel bir depolama uygulamasını kullanın, depolama arka ucu güvenilir, özel altyapı olarak ele alınmalıdır. Güvenilmeyen veya üzerinde oynanma olasılığı olan kaynaklardan hiçbir zaman denetim noktaları yüklemeyin.
Denetim noktaları için kullanılan depolama konumunun uygun şekilde güvenli olduğundan emin olun. Denetim noktası verilerine yalnızca yetkili hizmetlerin ve kullanıcıların okuma veya yazma erişimi olmalıdır.
Pickle modülünde serileştirme
Hem FileCheckpointStorage hem de CosmosCheckpointStorage, veri sınıfları, tarih zamanları ve özel nesneler gibi JSON yerel olmayan durumu seri hale getirmek için Python pickle modülünü kullanır. Seri durumdan çıkarma sırasında rastgele kod yürütme risklerini azaltmak için her iki sağlayıcı da varsayılan olarak kısıtlanmış bir unpickler kullanır. Seri durumdan çıkarma sırasında yalnızca yerleşik bir güvenli Python türü kümesine (temel öğeler, datetime, uuid, Decimal, ortak koleksiyonlar vb.) ve desteklenen Agent Framework veya OpenAI SDK türlerine izin verilir. Modül öneki izin listesi yalnızca türlerle sınırlıdır: yardımcı fonksiyonlar ve tür olmayan diğer global tanımlar reddedilir.
Agent Framework, pickle etmeden önce çerçeveye özgü nesnelerden geçici raw_representation değerlerini çıkarır ve bu alanları None olarak geri yükler. Desteklenmeyen başka herhangi bir tür, seri durumdan çıkarma işleminin WorkflowCheckpointException ile başarısız olmasına neden olur.
Uygulamaya özgü ek türlere izin vermek için, allowed_checkpoint_types bunları "module:qualname" biçimini kullanarak parametresi aracılığıyla iletin:
from agent_framework import FileCheckpointStorage
storage = FileCheckpointStorage(
"/tmp/checkpoints",
allowed_checkpoint_types=[
"my_app.models:SafeState",
"my_app.models:UserProfile",
],
)
Her allowed_checkpoint_types girişin bir türe çözümlenmesi gerekir. Modül düzeyinde bir fonksiyon veya tür olmayan başka bir global öğe eklemek, o global öğeyi seri durumdan çıkarılabilir hâle getirmez.
CosmosCheckpointStorage aynı parametreyi kabul eder:
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",
],
)
Mevcut Python işlemindeki her kısıtlanmış denetim noktası çözücüsü için bir tür kaydetmek üzere sınıfı register_checkpoint_type() öğesine iletin:
from agent_framework import register_checkpoint_type
from my_app.models import SafeState
register_checkpoint_type(SafeState)
Genel kayıt, çağrıdan önce oluşturulan depolama örnekleri için de geçerlidir.
Buna karşılık, allowed_checkpoint_types yalnızca FileCheckpointStorage bunu alan veya CosmosCheckpointStorage örneğine uygulanır.
Kodunuz depolamayı oluştururken depo başına parametresini tercih edin. Kodunuz depolama örneğini yapılandıramıyorsa, örneğin sizin için denetim noktası depolaması oluşturan barındırılan bir ortamda genel kayıt kullanın. Türü, onu içeren bir denetim noktasını yüklemeden önce kaydedin.
Her iki mekanizma da kısıtlı unpickler’ın izin listesini genişletir; ancak hiçbiri güvenilmeyen pickle verilerini güvenli kılmaz. İzin verilen bir sınıf özel turşu yeniden oluşturma davranışı tanımlayabildiği için yalnızca güvenilen uygulama sınıflarını kaydedin. Genel kayıt, işlemdeki her kısıtlı denetim noktası yükü için seri durumdan çıkarma yüzeyini artırır, bu nedenle genel kayıt defterini gerekli en düşük kümede tutun.
Tehdit modeliniz pickle tabanlı serileştirmeye hiç izin vermiyorsa, alternatif bir serileştirme stratejisiyle özel bir InMemoryCheckpointStorage kullanın veya CheckpointStorage uygulayın.
Depolama konumu sorumluluğu
FileCheckpointStorage açık storage_path bir parametre gerektirir; varsayılan dizin yoktur. Çerçeve yol geçişi saldırılarına karşı doğrulanırken, depolama dizininin güvenliğini sağlamak (dosya izinleri, bekleyen şifreleme, erişim denetimleri) geliştiricinin sorumluluğundadır. Yalnızca yetkili işlemlerin denetim noktası dizinine okuma veya yazma erişimi olmalıdır.
CosmosCheckpointStorage depolama için Azure Cosmos DB dayanır. Mümkün olduğunda yönetilen kimlik / RBAC kullanın, veritabanı ve kapsayıcının kapsamını iş akışı hizmetine alın ve anahtar tabanlı kimlik doğrulaması kullanıyorsanız hesap anahtarlarını döndürün. Dosya depolamada olduğu gibi, denetim noktası belgelerini barındıran Cosmos DB kapsayıcısına yalnızca yetkili sorumluların okuma veya yazma erişimi olmalıdır.
Git denetim noktası yöneticileri denetim noktası durumunu JSON olarak seri hale getirse de denetim noktası depolaması yine de güvenilir uygulama durumudur. kullanıyorsanız checkpoint.NewFileSystemJSONStore, denetim noktası dosyalarını korumalı bir dizinde depolayın ve okuma/yazma erişimini yalnızca yetkili işlemlerle kısıtlayın. Özel mağazalar kendi erişim denetimi, bütünlük ve dayanıklılık garantilerinden sorumludur.