Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 36 additions & 4 deletions src/Dapr.Workflow.Abstractions/WorkflowRuntimeOptions.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
// ------------------------------------------------------------------------
// ------------------------------------------------------------------------
// Copyright 2022 The Dapr Authors
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
Expand All @@ -25,6 +25,9 @@ public sealed class WorkflowRuntimeOptions
private readonly List<Action<IWorkflowsFactory>> _registrationActions = [];
private int _maxConcurrentWorkflows = 100;
private int _maxConcurrentActivities = 100;
private TimeSpan? _historyCacheTtl;
private int _historyCacheMaxInstances;
private long _historyCacheMaxBytes;

/// <summary>
/// Gets or sets the gRPC channel options used for connecting to the Dapr sidecar.
Expand All @@ -45,22 +48,51 @@ public sealed class WorkflowRuntimeOptions
/// when it has gone idle (no turn) for longer than this. <c>null</c> (the default) uses the built-in
/// default of one hour. Ignored when <see cref="DisableStatefulHistory"/> is <c>true</c>.
/// </summary>
public TimeSpan? HistoryCacheTtl { get; set; }
public TimeSpan? HistoryCacheTtl
{
get => _historyCacheTtl;
set
{
if (value is { } configured && configured <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(nameof(value), configured,
"The history cache TTL must be positive.");
}

_historyCacheTtl = value;
}
}

/// <summary>
/// Gets or sets the maximum number of per-instance histories retained on a single work-item stream;
/// least-recently-used entries are evicted beyond it. <c>0</c> (the default) uses the built-in default.
/// Ignored when <see cref="DisableStatefulHistory"/> is <c>true</c>.
/// </summary>
public int HistoryCacheMaxInstances { get; set; }
public int HistoryCacheMaxInstances
{
get => _historyCacheMaxInstances;
set
{
ArgumentOutOfRangeException.ThrowIfNegative(value);
_historyCacheMaxInstances = value;
}
}

/// <summary>
/// Gets or sets the byte budget for cached histories on a single work-item stream; least-recently-used
/// entries are evicted beyond it. <c>0</c> (the default) means unlimited (bounded only by
/// <see cref="HistoryCacheMaxInstances"/> and <see cref="HistoryCacheTtl"/>). Ignored when
/// <see cref="DisableStatefulHistory"/> is <c>true</c>.
/// </summary>
public long HistoryCacheMaxBytes { get; set; }
public long HistoryCacheMaxBytes
{
get => _historyCacheMaxBytes;
set
{
ArgumentOutOfRangeException.ThrowIfNegative(value);
_historyCacheMaxBytes = value;
}
}

/// <summary>
/// Gets the maximum number of concurrent workflow instances that can be executed at the same time.
Expand Down
79 changes: 45 additions & 34 deletions src/Dapr.Workflow/Worker/Grpc/WorkflowHistoryCache.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
// ------------------------------------------------------------------------
// ------------------------------------------------------------------------
// Copyright 2026 The Dapr Authors
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
Expand Down Expand Up @@ -38,11 +38,13 @@ private sealed class Entry
/// <summary>Matches the width of the running total and the configured budget, so a large
/// history cannot overflow and corrupt the accounting.</summary>
public required long Bytes { get; init; }
public required LinkedListNode<string> RecencyNode { get; init; }
public DateTime LastAccess { get; set; }
}

private readonly object _lock = new();
private readonly Dictionary<string, Entry> _entries = new();
private readonly LinkedList<string> _recency = new();
private readonly TimeSpan _ttl;
private readonly int _maxInstances;
private readonly long _maxBytes;
Expand All @@ -67,16 +69,24 @@ public int Generation
}
}

/// <summary>Initializes the cache. Non-positive ttl/maxInstances use defaults; maxBytes &lt;= 0 means unlimited.</summary>
/// <summary>Initializes the cache. Null ttl and zero maxInstances use defaults; zero maxBytes means unlimited.</summary>
public WorkflowHistoryCache(
TimeSpan? ttl = null,
int maxInstances = 0,
long maxBytes = 0,
Func<DateTime>? clock = null)
{
_ttl = ttl is { } configured && configured > TimeSpan.Zero ? configured : DefaultTtl;
_maxInstances = maxInstances > 0 ? maxInstances : DefaultMaxInstances;
_maxBytes = maxBytes > 0 ? maxBytes : 0;
if (ttl is { } configured && configured <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(nameof(ttl), configured, "The history cache TTL must be positive.");
}

ArgumentOutOfRangeException.ThrowIfNegative(maxInstances);
ArgumentOutOfRangeException.ThrowIfNegative(maxBytes);

_ttl = ttl ?? DefaultTtl;
_maxInstances = maxInstances == 0 ? DefaultMaxInstances : maxInstances;
_maxBytes = maxBytes;
_clock = clock ?? (() => DateTime.UtcNow);
}

Expand All @@ -90,7 +100,7 @@ public WorkflowHistoryCache(
return null;
}

entry.LastAccess = _clock();
TouchLocked(entry);
return entry.Events;
}
}
Expand Down Expand Up @@ -125,12 +135,20 @@ public void Put(string instanceId, IEnumerable<HistoryEvent> events, int generat

if (_entries.TryGetValue(instanceId, out var existing))
{
_totalBytes -= existing.Bytes;
_entries.Remove(instanceId);
RemoveEntryLocked(existing);
}

_entries[instanceId] = new Entry { Events = snapshot, Bytes = bytes, LastAccess = _clock() };
var recencyNode = _recency.AddFirst(instanceId);
_entries[instanceId] = new Entry
{
Events = snapshot,
Bytes = bytes,
LastAccess = _clock(),
RecencyNode = recencyNode
};
_totalBytes += bytes;
EvictToFit(instanceId);
EvictToFit();
}
}

Expand Down Expand Up @@ -161,6 +179,7 @@ public void Reset()
lock (_lock)
{
_entries.Clear();
_recency.Clear();
_totalBytes = 0;
_generation++;
}
Expand Down Expand Up @@ -210,20 +229,33 @@ internal long TotalBytes
}
}

private void TouchLocked(Entry entry)
{
entry.LastAccess = _clock();
_recency.Remove(entry.RecencyNode);
_recency.AddFirst(entry.RecencyNode);
}

private void RemoveLocked(string instanceId)
{
if (_entries.Remove(instanceId, out var entry))
{
_totalBytes -= entry.Bytes;
RemoveEntryLocked(entry);
}
}

private void RemoveEntryLocked(Entry entry)
{
_recency.Remove(entry.RecencyNode);
_totalBytes -= entry.Bytes;
}

/// <summary>
/// Evicts least-recently-used entries until within the count and byte bounds. Always keeps the
/// just-touched entry so the active working set is never evicted; a lone entry over the byte
/// budget is kept (a soft overage) rather than thrashing.
/// </summary>
private void EvictToFit(string keep)
private void EvictToFit()
{
while (_entries.Count > 1)
{
Expand All @@ -234,34 +266,13 @@ private void EvictToFit(string keep)
return;
}

var victim = LeastRecentlyUsedExcept(keep);
var victim = _recency.Last;
if (victim is null)
{
return;
}

RemoveLocked(victim);
RemoveLocked(victim.Value);
}
}

private string? LeastRecentlyUsedExcept(string keep)
{
string? oldestId = null;
var oldestAccess = DateTime.MaxValue;
foreach (var (instanceId, entry) in _entries)
{
if (instanceId == keep)
{
continue;
}

if (oldestId is null || entry.LastAccess < oldestAccess)
{
oldestId = instanceId;
oldestAccess = entry.LastAccess;
}
}

return oldestId;
}
}

This file was deleted.

Loading
Loading