Files
agent-framework/dotnet/src/Microsoft.Agents.AI.Workflows/Checkpointing/FileSystemJsonCheckpointStore.cs
Jacob Alber 0086d38f58 .NET: [BREAKING] Workflows API Review Naming Changes (Part 1?) (#4090)
* refactor: Normalize Run/RunStreaming with AIAgent

* refactor: Clarify Session vs. Run -level concepts

* Rename RunId to SessionId to better match Run/Session terminology in AIAgent
* [BREAKING]: Will break existing checkpointed sessions in CosmosDb due to field rename

* refactor: Rename and simplify interface around getting typed data out of ExternalRequest/Response

* Also adds hints around using value types in PortableValue

* refactor: Rename AddFanInEdge to AddFanInBarrierEdge

This will prevent a breaking change later when we introduce a programmable FanIn edge, analogous to the FanOut edge's EdgeSelector.

The goal, in the long run is to support a number of different FanIn scenarios, with naive FanIn (no barrier) by default, similar to FanOut.

* refactor: AsAgent(this Workflow, ...) => AsAIAgent(...)

* misc - part1: SwitchBuilder internal

---------

Co-authored-by: Dmytro Struk <13853051+dmytrostruk@users.noreply.github.com>
2026-02-20 02:05:18 +00:00

173 lines
6.7 KiB
C#

// Copyright (c) Microsoft. All rights reserved.
using System;
using System.Collections.Generic;
using System.IO;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
namespace Microsoft.Agents.AI.Workflows.Checkpointing;
/// <summary>
/// Provides a file system-based implementation of a JSON checkpoint store that persists checkpoint data and index
/// information to disk using JSON files.
/// </summary>
/// <remarks>This class manages checkpoint storage by writing JSON files to a specified directory and maintaining
/// an index file for efficient retrieval. It is intended for scenarios where durable, process-exclusive checkpoint
/// persistence is required. Instances of this class are not thread-safe and should not be shared across multiple
/// threads without external synchronization. The class implements IDisposable; callers should ensure Dispose is called
/// to release file handles and system resources when the store is no longer needed.</remarks>
public sealed class FileSystemJsonCheckpointStore : JsonCheckpointStore, IDisposable
{
[System.Diagnostics.CodeAnalysis.SuppressMessage("Usage", "CA2213:Disposable fields should be disposed",
Justification = "It is disposed, the analyzer is just not picking it up properly")]
private FileStream? _indexFile;
internal DirectoryInfo Directory { get; }
internal HashSet<CheckpointInfo> CheckpointIndex { get; }
/// <summary>
/// Initializes a new instance of the <see cref="FileSystemJsonCheckpointStore"/> class that uses the specified directory
/// </summary>
/// <param name="directory"></param>
/// <exception cref="ArgumentNullException"></exception>
/// <exception cref="InvalidOperationException"></exception>
public FileSystemJsonCheckpointStore(DirectoryInfo directory)
{
this.Directory = directory ?? throw new ArgumentNullException(nameof(directory));
if (!directory.Exists)
{
directory.Create();
}
try
{
this._indexFile = File.Open(Path.Combine(directory.FullName, "index.jsonl"), FileMode.OpenOrCreate, FileAccess.ReadWrite, FileShare.None);
}
catch
{
throw new InvalidOperationException($"The store at '{directory.FullName}' is already in use by another process.");
}
try
{
// read the lines of indexfile and parse them as CheckpointInfos
this.CheckpointIndex = [];
#if NET
const int BufferSize = -1;
#else
const int BufferSize = 1024;
#endif
using StreamReader reader = new(this._indexFile, encoding: Encoding.UTF8, detectEncodingFromByteOrderMarks: false, BufferSize, leaveOpen: true);
while (reader.ReadLine() is string line)
{
if (JsonSerializer.Deserialize(line, KeyTypeInfo) is { } info)
{
this.CheckpointIndex.Add(info);
}
}
}
catch (Exception exception)
{
throw new InvalidOperationException($"Could not load store at '{directory.FullName}'. Index corrupted.", exception);
}
}
/// <inheritdoc/>
public void Dispose()
{
FileStream? indexFileLocal = Interlocked.Exchange(ref this._indexFile, null);
indexFileLocal?.Dispose();
}
[System.Diagnostics.CodeAnalysis.SuppressMessage("Maintainability", "CA1513:Use ObjectDisposedException throw helper",
Justification = "Throw helper does not exist in NetFx 4.7.2")]
private void CheckDisposed()
{
if (this._indexFile is null)
{
throw new ObjectDisposedException($"{nameof(FileSystemJsonCheckpointStore)}({this.Directory.FullName})");
}
}
private string GetFileNameForCheckpoint(string sessionId, CheckpointInfo key)
=> Path.Combine(this.Directory.FullName, $"{sessionId}_{key.CheckpointId}.json");
private CheckpointInfo GetUnusedCheckpointInfo(string sessionId)
{
CheckpointInfo key;
do
{
key = new(sessionId);
} while (!this.CheckpointIndex.Add(key));
return key;
}
/// <inheritdoc/>
[System.Diagnostics.CodeAnalysis.SuppressMessage("Performance", "CA1835:Prefer the 'Memory'-based overloads for 'ReadAsync' and 'WriteAsync'",
Justification = "Memory-based overload is missing for 4.7.2")]
public override async ValueTask<CheckpointInfo> CreateCheckpointAsync(string sessionId, JsonElement value, CheckpointInfo? parent = null)
{
this.CheckDisposed();
CheckpointInfo key = this.GetUnusedCheckpointInfo(sessionId);
string fileName = this.GetFileNameForCheckpoint(sessionId, key);
try
{
using Stream checkpointStream = File.Open(fileName, FileMode.Create, FileAccess.Write, FileShare.None);
using Utf8JsonWriter jsonWriter = new(checkpointStream, new JsonWriterOptions() { Indented = false });
value.WriteTo(jsonWriter);
JsonSerializer.Serialize(this._indexFile!, key, KeyTypeInfo);
byte[] bytes = Encoding.UTF8.GetBytes(Environment.NewLine);
await this._indexFile!.WriteAsync(bytes, 0, bytes.Length, CancellationToken.None).ConfigureAwait(false);
await this._indexFile!.FlushAsync(CancellationToken.None).ConfigureAwait(false);
return key;
}
catch (Exception ex)
{
this.CheckpointIndex.Remove(key);
try
{
// try to clean up after ourselves
File.Delete(fileName);
}
catch { }
throw new InvalidOperationException($"Could not create checkpoint in store at '{this.Directory.FullName}'.", ex);
}
}
/// <inheritdoc/>
public override async ValueTask<JsonElement> RetrieveCheckpointAsync(string sessionId, CheckpointInfo key)
{
this.CheckDisposed();
string fileName = this.GetFileNameForCheckpoint(sessionId, key);
if (!this.CheckpointIndex.Contains(key) ||
!File.Exists(fileName))
{
throw new KeyNotFoundException($"Checkpoint '{key.CheckpointId}' not found in store at '{this.Directory.FullName}'.");
}
using FileStream checkpointFileStream = File.Open(fileName, FileMode.Open, FileAccess.Read, FileShare.Read);
using JsonDocument document = await JsonDocument.ParseAsync(checkpointFileStream).ConfigureAwait(false);
return document.RootElement.Clone();
}
/// <inheritdoc/>
public override ValueTask<IEnumerable<CheckpointInfo>> RetrieveIndexAsync(string sessionId, CheckpointInfo? withParent = null)
{
this.CheckDisposed();
return new(this.CheckpointIndex);
}
}