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
16 changes: 14 additions & 2 deletions BlogWorkflow.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,16 +22,24 @@ public class BlogWorkflow(
// Emits the root span for a workflow run. Activated by the ActivityListener
// registered in Program.cs (or an OpenTelemetry TracerProvider).
private static readonly ActivitySource s_activitySource = new("BlogWriter.Workflow");
private static long s_operationVersion;

public async Task<ResearchState> RunAsync(ResearchState state, CancellationToken cancellationToken = default)
public async Task<ResearchState> RunAsync(
ResearchState state,
CancellationToken cancellationToken = default,
IProgress<WorkflowOutputUpdate>? output = null)
{
long operationVersion = Interlocked.Increment(ref s_operationVersion);
var publisher = new WorkflowOutputPublisher(output, operationVersion);
publisher.PublishLifecycle(WorkflowOutputOutcome.Progress, "Writing workflow started.");

using Activity? activity = s_activitySource.StartActivity("Workflow.Run");
activity?.SetTag("blog.topic", state.MainTask);

var bloggerExecutor = new BloggerExecutor(blogger);
var researcherExecutor = new ResearcherExecutor(researcher);
var authorExecutor = new AuthorExecutor(author);
var reviewerExecutor = new ReviewerExecutor(reviewer);
var reviewerExecutor = new ReviewerExecutor(reviewer, publisher);

Workflow workflow = new WorkflowBuilder(bloggerExecutor)
.AddEdge(bloggerExecutor, researcherExecutor)
Expand Down Expand Up @@ -59,14 +67,17 @@ public async Task<ResearchState> RunAsync(ResearchState state, CancellationToken
{
case ExecutorInvokedEvent invoked:
logger.LogInformation("[workflow] -> {ExecutorId} started", invoked.ExecutorId);
publisher.PublishLifecycle(WorkflowOutputOutcome.Progress, $"{invoked.ExecutorId} started.");
break;

case ExecutorCompletedEvent completed:
logger.LogInformation("[workflow] {ExecutorId} completed", completed.ExecutorId);
publisher.PublishLifecycle(WorkflowOutputOutcome.Progress, $"{completed.ExecutorId} completed.");
break;

case ExecutorFailedEvent failed:
logger.LogError(failed.Data as Exception, "[workflow] {ExecutorId} failed", failed.ExecutorId);
publisher.PublishLifecycle(WorkflowOutputOutcome.Failure, $"{failed.ExecutorId} failed.");

// A token-cap breach must abort the whole run, not just the
// node. Re-throw it so it unwinds to the application entry point.
Expand All @@ -81,6 +92,7 @@ public async Task<ResearchState> RunAsync(ResearchState state, CancellationToken
case WorkflowOutputEvent { Data: ResearchState finalState }:
// The reviewer yielded the final, approved (or revision-capped) state.
result = finalState;
publisher.PublishLifecycle(WorkflowOutputOutcome.Success, "Writing workflow completed.");
break;
}
}
Expand Down
110 changes: 110 additions & 0 deletions BlogWriter.Tests/BlogWorkflowTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
using Microsoft.Extensions.Logging.Abstractions;
using Xunit;

namespace BlogWriter.Tests;

public sealed class BlogWorkflowTests
{
[Fact]
public async Task RunAsync_EmitsLifecycleAndReviewerUpdatesWithoutChangingFinalState()
{
var workflow = new BlogWorkflow(
new TestBlogger(),
new TestResearcher(),
new TestAuthor(),
new TestReviewer(),
NullLogger<BlogWorkflow>.Instance);
var output = new WorkflowOutputCollector();

var service = new BlogWriterSessionService(workflow, new RecordingStore());
BlogSession session = await service.StartAsync("topic", output: output);
ResearchState result = session.State;

Assert.Equal("draft", result.Draft);
Assert.Equal("APPROVED", result.ReviewNotes);
Assert.Contains(output.Updates, update =>
update.Kind == WorkflowOutputKind.Lifecycle &&
update.Outcome == WorkflowOutputOutcome.Progress);
Assert.Contains(output.Updates, update =>
update.Kind == WorkflowOutputKind.ReviewerFeedback &&
update.Message == "APPROVED");
Assert.Contains(output.Updates, update =>
update.Kind == WorkflowOutputKind.Lifecycle &&
update.Outcome == WorkflowOutputOutcome.Success);
Assert.Equal(output.Updates.Count, output.Updates.Select(update => update.Sequence).Distinct().Count());
}

private sealed class TestBlogger : IBloggerAgent
{
public Task<BloggerDecision> InvokeAsync(ResearchState state, CancellationToken cancellationToken = default) =>
Task.FromResult(new BloggerDecision("research", state.MainTask));

public Task<ResearchState> BloggerNodeAsync(ResearchState state, CancellationToken cancellationToken = default)
{
state.NextStep = "research";
state.CurrentSubTask = state.MainTask;
return Task.FromResult(state);
}
}

private sealed class TestResearcher : IResearcherAgent
{
public Task<string> InvokeAsync(string query, CancellationToken cancellationToken = default) =>
Task.FromResult("finding");

public Task<ResearchState> ResearchNodeAsync(ResearchState state, CancellationToken cancellationToken = default)
{
state.ResearchFindings.Add("finding");
return Task.FromResult(state);
}
}

private sealed class TestAuthor : IAuthorAgent
{
public Task<string?> InvokeAsync(ResearchState state, CancellationToken cancellationToken = default) =>
Task.FromResult<string?>("draft");

public Task<ResearchState> AuthorNodeAsync(ResearchState state, CancellationToken cancellationToken = default)
{
state.Draft = "draft";
return Task.FromResult(state);
}
}

private sealed class TestReviewer : IReviewerAgent
{
public Task<string> InvokeAsync(ResearchState state, CancellationToken cancellationToken = default) =>
Task.FromResult("APPROVED");

public Task<ResearchState> ReviewerNodeAsync(ResearchState state, CancellationToken cancellationToken = default)
{
state.ReviewNotes = "APPROVED";
return Task.FromResult(state);
}
}

private sealed class RecordingStore : IBlogSessionStore
{
public Task<BlogSession> CreateAsync(ResearchState state, CancellationToken cancellationToken = default) =>
Task.FromResult(new BlogSession
{
Id = Guid.NewGuid().ToString("N"),
OwnerId = "owner",
CreatedAt = DateTimeOffset.UtcNow,
UpdatedAt = DateTimeOffset.UtcNow,
State = state,
});

public Task<BlogSession?> GetAsync(string sessionId, CancellationToken cancellationToken = default) =>
Task.FromResult<BlogSession?>(null);

public Task<IReadOnlyList<BlogSessionSummary>> ListAsync(CancellationToken cancellationToken = default) =>
Task.FromResult<IReadOnlyList<BlogSessionSummary>>([]);

public Task SaveAsync(BlogSession session, CancellationToken cancellationToken = default) =>
Task.CompletedTask;

public Task DeleteOwnerSessionsAsync(string ownerId, CancellationToken cancellationToken = default) =>
Task.CompletedTask;
}
}
22 changes: 21 additions & 1 deletion BlogWriter.Tests/BlogWriterSessionServiceTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,21 @@ public async Task ListAndLoadAsync_DelegateToOwnerScopedStore()
Assert.Same(existing, loaded);
}

[Fact]
public async Task StartAndReviseAsync_ForwardOutputObserverToWorkflow()
{
var workflow = new StubWorkflow(state => state);
var service = new BlogWriterSessionService(workflow, new RecordingStore());
var output = new Progress<WorkflowOutputUpdate>();

await service.StartAsync("topic", output: output);
Assert.Same(output, workflow.LastOutput);

BlogSession session = CreateSession("draft", "review");
await service.ReviseAsync(session, "change it", 500, 900, output: output);
Assert.Same(output, workflow.LastOutput);
}

private static BlogSession CreateSession(string draft, string review) => new()
{
Id = Guid.NewGuid().ToString("N"),
Expand All @@ -115,11 +130,16 @@ public async Task ListAndLoadAsync_DelegateToOwnerScopedStore()
private sealed class StubWorkflow(Func<ResearchState, ResearchState> run) : IBlogWorkflow
{
public int CallCount { get; private set; }
public IProgress<WorkflowOutputUpdate>? LastOutput { get; private set; }

public Task<ResearchState> RunAsync(ResearchState state, CancellationToken cancellationToken = default)
public Task<ResearchState> RunAsync(
ResearchState state,
CancellationToken cancellationToken = default,
IProgress<WorkflowOutputUpdate>? output = null)
{
cancellationToken.ThrowIfCancellationRequested();
CallCount++;
LastOutput = output;
return Task.FromResult(run(state));
}
}
Expand Down
8 changes: 8 additions & 0 deletions BlogWriter.Tests/WorkflowOutputTestDoubles.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
namespace BlogWriter.Tests;

internal sealed class WorkflowOutputCollector : IProgress<WorkflowOutputUpdate>
{
public List<WorkflowOutputUpdate> Updates { get; } = [];

public void Report(WorkflowOutputUpdate value) => Updates.Add(value);
}
51 changes: 51 additions & 0 deletions BlogWriter.Tests/WorkflowOutputUpdateTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
using Xunit;

namespace BlogWriter.Tests;

public sealed class WorkflowOutputUpdateTests
{
[Fact]
public void Create_RequiresUserVisibleMessage()
{
Assert.Throws<ArgumentException>(() => WorkflowOutputUpdate.Create(
WorkflowOutputKind.Lifecycle,
WorkflowOutputOutcome.Progress,
"",
operationVersion: 1,
sequence: 1,
updateKey: "op-1"));
}

[Fact]
public void Create_PreservesRoutingAndOrderingMetadata()
{
WorkflowOutputUpdate update = WorkflowOutputUpdate.Create(
WorkflowOutputKind.ReviewerFeedback,
WorkflowOutputOutcome.Review,
"Needs a stronger conclusion.",
operationVersion: 4,
sequence: 7,
updateKey: "op-4-review-1",
revisionNumber: 1);

Assert.Equal(WorkflowOutputKind.ReviewerFeedback, update.Kind);
Assert.Equal(WorkflowOutputOutcome.Review, update.Outcome);
Assert.Equal("Needs a stronger conclusion.", update.Message);
Assert.Equal(4, update.OperationVersion);
Assert.Equal(7, update.Sequence);
Assert.Equal("op-4-review-1", update.UpdateKey);
Assert.Equal(1, update.RevisionNumber);
}

[Fact]
public void Create_RequiresStableUpdateKey()
{
Assert.Throws<ArgumentException>(() => WorkflowOutputUpdate.Create(
WorkflowOutputKind.Lifecycle,
WorkflowOutputOutcome.Success,
"Complete",
operationVersion: 1,
sequence: 1,
updateKey: " "));
}
}
34 changes: 34 additions & 0 deletions BlogWriter.Web.Tests/BlogWorkspaceOutputTestHelpers.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
using BlogWriter.Web.Services;

namespace BlogWriter.Web.Tests;

internal static class BlogWorkspaceOutputTestHelpers
{
public static WorkflowOutputUpdate Lifecycle(
string message,
long operationVersion = 1,
long sequence = 1,
WorkflowOutputOutcome outcome = WorkflowOutputOutcome.Progress) =>
WorkflowOutputUpdate.Create(
WorkflowOutputKind.Lifecycle,
outcome,
message,
operationVersion,
sequence,
$"lifecycle-{operationVersion}-{sequence}");

public static WorkflowOutputUpdate Review(
string message,
string updateKey,
long operationVersion = 1,
long sequence = 1,
int revisionNumber = 0) =>
WorkflowOutputUpdate.Create(
WorkflowOutputKind.ReviewerFeedback,
WorkflowOutputOutcome.Review,
message,
operationVersion,
sequence,
updateKey,
revisionNumber);
}
Loading
Loading