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
5 changes: 4 additions & 1 deletion src/InProcessTestHost/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,10 @@ Only the current execution is retained. `ContinueAsNew` replaces the previous ge
history when the new generation commits. The underlying in-memory service also accepts an
execution ID: null or empty selects the current execution, and a different execution ID
returns no history (gRPC `NotFound`). Worker history streaming continues to use the dispatched
episode's replay snapshot rather than this management snapshot.
episode's replay snapshot rather than this management snapshot. These temporary worker snapshots
remain available until the episode's final response and are released before the dispatcher commits
that episode or starts its next generation. A failed dispatch also releases its snapshot; committed
history is retained independently until purge or generation replacement.

### 5. Purge completed instances

Expand Down
22 changes: 16 additions & 6 deletions src/InProcessTestHost/Sidecar/Grpc/TaskHubGrpcServer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -822,6 +822,7 @@ async Task<GrpcOrchestratorExecutionResult> ITaskExecutor.ExecuteOrchestrator(
// This must be done before we start the orchestrator execution.
TaskCompletionSource<GrpcOrchestratorExecutionResult> tcs =
this.CreateTaskCompletionSourceForOrchestrator(instance.InstanceId);
List<P.HistoryEvent>? streamedPastEvents = null;

try
{
Expand All @@ -845,7 +846,8 @@ async Task<GrpcOrchestratorExecutionResult> ITaskExecutor.ExecuteOrchestrator(
if (this.supportsHistoryStreaming && totalBytes > HistoryStreamingThresholdBytes)
{
orkRequest.RequiresHistoryStreaming = true;
// Store past events to serve via StreamInstanceHistory
// Keep this episode's replay snapshot available until execution finishes.
streamedPastEvents = protoPastEvents;
this.streamingPastEvents[instance.InstanceId] = protoPastEvents;
}
else
Expand All @@ -858,18 +860,26 @@ await this.SendWorkItemToClientAsync(new P.WorkItem
{
OrchestratorRequest = orkRequest,
});

// The TCS will be completed on the message stream handler when it gets a response back from the remote process
// TODO: How should we handle timeouts if the remote process never sends a response?
// Probably need to have a static timeout (e.g. 5 minutes).
return await tcs.Task;
}
catch
{
// Remove the TaskCompletionSource that we just created
this.RemoveOrchestratorTaskCompletionSource(instance.InstanceId);
throw;
}

// The TCS will be completed on the message stream handler when it gets a response back from the remote process
// TODO: How should we handle timeouts if the remote process never sends a response?
// Probably need to have a static timeout (e.g. 5 minutes).
return await tcs.Task;
finally
{
if (streamedPastEvents is not null)
{
this.streamingPastEvents.TryRemove(
new KeyValuePair<string, List<P.HistoryEvent>>(instance.InstanceId, streamedPastEvents));
}
}
}

async Task<ActivityExecutionResult> ITaskExecutor.ExecuteActivity(OrchestrationInstance instance, TaskScheduledEvent activityEvent)
Expand Down
71 changes: 71 additions & 0 deletions test/InProcessTestHost.Tests/OrchestrationHistoryTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
using Xunit;
using Xunit.Abstractions;
using P = Microsoft.DurableTask.Protobuf;
using PurgeResult = Microsoft.DurableTask.Client.PurgeResult;

namespace InProcessTestHost.Tests;

Expand Down Expand Up @@ -132,6 +133,76 @@ public async Task GetHistoryAsync_RunningAndCompleted_ReturnsFreshSnapshots(int
Assert.Single(completed.OfType<ExecutionCompletedEvent>());
}

/// <summary>
/// Reuses one host without retaining completed episodes' replay snapshots or losing committed history.
/// </summary>
[Theory]
[InlineData(32, false)]
[InlineData(600 * 1024, false)]
[InlineData(600 * 1024, true)]
public async Task GetHistoryAsync_ReusedHost_ReleasesWorkerSnapshots(int payloadSize, bool continueAsNew)
{
// Arrange
using CancellationTokenSource timeout = new(TimeSpan.FromSeconds(30));
HistoryRequestInterceptor interceptor = new();
await using DurableTaskTestHost host = await DurableTaskTestHost.StartAsync(tasks =>
{
tasks.AddOrchestratorFunc<string, int>("PayloadLength", async (context, input) =>
{
int length = await context.CallActivityAsync<int>("Length", input);
if (continueAsNew && input[0] == 'x')
{
context.ContinueAsNew(new string('y', input.Length));
}

return length;
});
tasks.AddActivityFunc<string, int>("Length", (context, input) => input.Length);
}, new DurableTaskTestHostOptions
{
ConfigureServices = services => services.Configure<GrpcDurableTaskWorkerOptions>(
options => options.Interceptors.Add(interceptor)),
}, timeout.Token);
string[] instanceIds = new string[2];

for (int i = 0; i < instanceIds.Length; i++)
{
// Act
string instanceId = await host.Client.ScheduleNewOrchestrationInstanceAsync(
"PayloadLength", new string('x', payloadSize), cancellation: timeout.Token);
instanceIds[i] = instanceId;
OrchestrationMetadata metadata = await host.Client.WaitForInstanceCompletionAsync(
instanceId, getInputsAndOutputs: true, cancellation: timeout.Token);
IList<HistoryEvent> history = await host.Client.GetOrchestrationHistoryAsync(instanceId, timeout.Token);

// Assert
Assert.Equal(OrchestrationRuntimeStatus.Completed, metadata.RuntimeStatus);
Assert.Equal(payloadSize, metadata.ReadOutputAs<int>());
Assert.Equal(8, history.Count);
Assert.StartsWith(continueAsNew ? "\"y" : "\"x", Assert.Single(history.OfType<ExecutionStartedEvent>()).Input);
Assert.Single(history.OfType<TaskCompletedEvent>());
Assert.Single(history.OfType<ExecutionCompletedEvent>());
int pastEventBytes = history.Take(4).Sum(e => ProtobufUtils.ToHistoryEventProto(e).CalculateSize());
bool streamsHistory = payloadSize > 32;
int streamsPerInstance = streamsHistory ? (continueAsNew ? 2 : 1) : 0;
this.output.WriteLine(
$"Instance {i + 1}: past-event protobuf size {pastEventBytes} bytes; worker history requests {interceptor.HistoryRequestCount}");
Assert.Equal(streamsHistory, pastEventBytes > 1024 * 1024);
Assert.Equal(streamsPerInstance * (i + 1), interceptor.HistoryRequestCount);
Assert.Equal(0, WorkerHistorySnapshotTestHelpers.GetSnapshots(host).Count);

PurgeResult purge = await host.Client.PurgeInstanceAsync(instanceId, cancellation: timeout.Token);
Assert.Equal(1, purge.PurgedInstanceCount);
Assert.Null(await host.Client.GetInstanceAsync(instanceId, cancellation: timeout.Token));
ArgumentException missing = await Assert.ThrowsAsync<ArgumentException>(() =>
host.Client.GetOrchestrationHistoryAsync(instanceId, timeout.Token));
Assert.Equal(StatusCode.NotFound, Assert.IsType<RpcException>(missing.InnerException).StatusCode);
Assert.Equal(0, WorkerHistorySnapshotTestHelpers.GetSnapshots(host).Count);
}

Assert.NotEqual(instanceIds[0], instanceIds[1]);
}

[Fact]
public async Task GetHistoryAsync_ContinueAsNew_ReturnsOnlyCurrentGeneration()
{
Expand Down
Loading
Loading