Skip to content
This repository was archived by the owner on Mar 31, 2026. It is now read-only.
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ public class ChatResponseEvent : EventBase
[Id(2)] public string MemberName { get; set; }
[Id(3)] public ChatResponse ChatResponse { get; set; }
[Id(4)] public long Term { get; set; }
[Id(5)] public string FailureSummary { get; set; }
}

[GenerateSerializer]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,5 +4,6 @@ public enum WorkflowExecutionStatus
{
Pending,
Running,
Completed
Completed,
Failed
}
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
public class WorkUnitExecutionRecord
{
[Id(0)]
public string WorkUnitGrainId { get; set; }

Check warning on line 31 in src/Aevatar.GAgents.GroupChat.Core/States/WorkflowExecutionRecordState.cs

View workflow job for this annotation

GitHub Actions / publish (Aevatar.GAgents.GroupChat.GroupMember)

Non-nullable property 'WorkUnitGrainId' must contain a non-null value when exiting constructor. Consider adding the 'required' modifier or declaring the property as nullable.
[Id(1)]
public DateTime StartTime { get; set; }
[Id(2)]
Expand All @@ -39,4 +39,6 @@
public string InputData { get; set; }
[Id(5)]
public string OutputData { get; set; }

[Id(6)] public string FailureSummary { get; set; }
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,11 @@
using GroupChat.GAgent.Feature.Common;
using GroupChat.GAgent.GEvent;
using Aevatar.Core.Abstractions;
using Aevatar.GAgents.GroupChat.WorkflowCoordinator;
using GroupChat.GAgent.Dto;
using GroupChat.GAgent.Feature.Blackboard;
using GroupChat.GAgent.Feature.Coordinator.GEvent;
using Microsoft.Extensions.Logging;

namespace GroupChat.GAgent;

Expand Down Expand Up @@ -40,12 +42,24 @@ public async Task HandleEventAsync(ChatEvent @event)
}

// var history = await GetCareChatMessagesFromBlackboardAsync(@event.BlackboardId);
var talkResponse = await ChatAsync(@event.BlackboardId, @event.CoordinatorMessages);
await PublishAsync(new ChatResponseEvent()
try
{
BlackboardId = @event.BlackboardId, MemberId = this.GetPrimaryKey(), MemberName = State.MemberName,
ChatResponse = talkResponse, Term = @event.Term
});
var talkResponse = await ChatAsync(@event.BlackboardId, @event.CoordinatorMessages);
await PublishAsync(new ChatResponseEvent()
{
BlackboardId = @event.BlackboardId, MemberId = this.GetPrimaryKey(), MemberName = State.MemberName,
ChatResponse = talkResponse, Term = @event.Term
});
}
catch (Exception e)
{
Logger.LogError($"[MemberGAgentBase] Handler ChatEvent fail: {e.Message}");
await PublishAsync(new ChatResponseEvent()
{
BlackboardId = @event.BlackboardId, MemberId = this.GetPrimaryKey(), MemberName = State.MemberName,
FailureSummary = e.ToString(), Term = @event.Term
});
}
}

[EventHandler]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,8 @@ private async Task TrySaveWorkflowViewAsync(WorkflowViewConfigDto configuration)
AgentId = configuration.WorkflowCoordinatorGAgentId
});
}

await ConfirmEvents();
}

protected override void GAgentTransitionState(WorkflowViewState state,
Expand Down
33 changes: 24 additions & 9 deletions src/Aevatar.GAgents.GroupChat/GroupMemberGAgentBase.cs
Original file line number Diff line number Diff line change
Expand Up @@ -50,16 +50,31 @@ public async Task HandleEventAsync(ChatEvent @event)
return;
}

// var history = await GetCareChatMessagesFromBlackboardAsync(@event.BlackboardId);
var talkResponse = await ChatAsync(@event.BlackboardId, @event.CoordinatorMessages);
await PublishAsync(new ChatResponseEvent
try
{
// var history = await GetCareChatMessagesFromBlackboardAsync(@event.BlackboardId);
var talkResponse = await ChatAsync(@event.BlackboardId, @event.CoordinatorMessages);
await PublishAsync(new ChatResponseEvent
{
BlackboardId = @event.BlackboardId,
MemberId = this.GetPrimaryKey(),
MemberName = State.MemberName,
ChatResponse = talkResponse,
Term = @event.Term
});
}
catch (Exception e)
{
BlackboardId = @event.BlackboardId,
MemberId = this.GetPrimaryKey(),
MemberName = State.MemberName,
ChatResponse = talkResponse,
Term = @event.Term
});
Logger.LogError($"[GroupMemberGAgentBase] Handler ChatEvent fail: {e.Message}");
await PublishAsync(new ChatResponseEvent()
{
BlackboardId = @event.BlackboardId,
MemberId = this.GetPrimaryKey(),
MemberName = State.MemberName,
FailureSummary = e.ToString(),
Term = @event.Term
});
}
}

[EventHandler]
Expand Down
11 changes: 10 additions & 1 deletion src/Aevatar.GAgents.GroupChat/WorkflowCoordinatorGAgent.cs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,15 @@ public async Task HandleEventAsync(ChatResponseEvent @event)
return;
}

if (!@event.FailureSummary.IsNullOrEmpty())
{
Logger.LogError(
$"[WorkflowCoordinatorGAgent] ChatResponseEvent workUnit execute fail:{@event.FailureSummary}");
RaiseEvent(new WorkflowStartFailedLogEvent());
await ConfirmEvents();
return;
}

var blackboard = GrainFactory.GetGrain<IBlackboardGAgent>(State.BlackboardId);
await blackboard.SetMessageAsync(new CoordinatorConfirmChatResponse()
{
Expand Down Expand Up @@ -84,7 +93,7 @@ await blackboard.SetMessageAsync(new CoordinatorConfirmChatResponse()
public async Task HandleEventAsync(StartWorkflowCoordinatorEvent @event)
{
Logger.LogDebug("[WorkflowCoordinatorGAgent] handler StartWorkflowCoordinatorEvent start");
if (State.WorkflowStatus != WorkflowCoordinatorStatus.Pending)
if (State.WorkflowStatus != WorkflowCoordinatorStatus.Pending && State.WorkflowStatus != WorkflowCoordinatorStatus.Failed)
{
Logger.LogError("[WorkflowCoordinatorGAgent] The workflow is not ready to run.");
return;
Expand Down
38 changes: 34 additions & 4 deletions src/Aevatar.GAgents.GroupChat/WorkflowExecutionRecordGAgent.cs
Original file line number Diff line number Diff line change
Expand Up @@ -44,11 +44,23 @@ public async Task HandleEventAsync(StartExecuteWorkUnitEvent @event)
[EventHandler]
public async Task HandleEventAsync(ChatResponseEvent @event)
{
RaiseEvent(new FinishExecuteWorkUnitLogEvent
if (@event.FailureSummary.IsNullOrEmpty())
{
WorkUnitGrainId = @event.PublisherGrainId.ToString(),
OutputData = JsonConvert.SerializeObject(@event.ChatResponse?.Content)
});
RaiseEvent(new FinishExecuteWorkUnitLogEvent
{
WorkUnitGrainId = @event.PublisherGrainId.ToString(),
OutputData = JsonConvert.SerializeObject(@event.ChatResponse?.Content)
});
}
else
{
RaiseEvent(new FailExecuteWorkflowLogEvent()
{
WorkUnitGrainId = @event.PublisherGrainId.ToString(),
FailureSummary = @event.FailureSummary
});
}

await ConfirmEvents();
}

Expand Down Expand Up @@ -102,6 +114,15 @@ protected override void GAgentTransitionState(WorkflowExecutionRecordState state
workUnit.Status = WorkflowExecutionStatus.Completed;
workUnit.OutputData = finishExecuteWorkUnitLogEvent.OutputData;
break;
case FailExecuteWorkflowLogEvent failExecuteWorkflowLogEvent:
var failWorkUnit = state.WorkUnitRecords.First(o =>
o.WorkUnitGrainId == failExecuteWorkflowLogEvent.WorkUnitGrainId);
failWorkUnit.EndTime = DateTime.UtcNow;
failWorkUnit.Status = WorkflowExecutionStatus.Failed;
failWorkUnit.FailureSummary = failExecuteWorkflowLogEvent.FailureSummary;

state.Status = WorkflowExecutionStatus.Failed;
break;
}
}
}
Expand Down Expand Up @@ -143,4 +164,13 @@ public class FinishExecuteWorkUnitLogEvent : WorkflowExecutionRecordLogEvent
public string WorkUnitGrainId { get; set; }
[Id(1)]
public string OutputData { get; set; }
}

[GenerateSerializer]
public class FailExecuteWorkflowLogEvent : WorkflowExecutionRecordLogEvent
{
[Id(0)]
public string WorkUnitGrainId { get; set; }
[Id(1)]
public string FailureSummary { get; set; }
}
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
<ItemGroup>
<ProjectReference Include="..\..\src\Aevatar.GAgents.Basic\Aevatar.GAgents.Basic.csproj" />
<ProjectReference Include="..\..\src\Aevatar.GAgents.GroupChat\Aevatar.GAgents.GroupChat.csproj" />
<ProjectReference Include="..\..\src\Aevatar.GAgents.InputGAgent\Aevatar.GAgents.InputGAgent.csproj" />
<ProjectReference Include="..\Aevatar.GAgents.TestBase\Aevatar.GAgents.TestBase.csproj" />
</ItemGroup>

Expand Down
26 changes: 26 additions & 0 deletions test/Aevatar.GAgents.GroupChat.Test/GAgents/WorkerGAgent.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
using GroupChat.GAgent;
using GroupChat.GAgent.Feature.Common;
using GroupChat.GAgent.GEvent;
using Volo.Abp;

namespace Aevatar.GAgents.GroupChat.Test.GAgents;

Expand All @@ -20,6 +21,15 @@ public async Task SetDelayWorkAsync(int delaySeconds)
await ConfirmEvents();
}

public async Task SetFailureSummary(string failureSummary)
{
RaiseEvent(new WorkerFailureLogEvent()
{
FailureSummary = failureSummary
});
await ConfirmEvents();
}

public Task<WorkerState> GetState()
{
return Task.FromResult(State);
Expand All @@ -39,6 +49,11 @@ protected override async Task<ChatResponse> ChatAsync(Guid blackboardId, List<Ch
await Task.Delay(TimeSpan.FromSeconds(State.DelaySeconds));
}

if (!State.FailureSummary.IsNullOrEmpty())
{
throw new UserFriendlyException(State.FailureSummary);
}

var response = new ChatResponse();
response.Content = $"{State.MemberName} Send the message";
RaiseEvent(new WorkHandleMessageLogEvent(){ PreWorkUnits = coordinatorMessages!.Select(s=>s.AgentName).ToList()});
Expand All @@ -56,6 +71,9 @@ protected override void GroupMemberTransitionState(WorkerState state, StateLogEv
case WorkerDelayLogEvent workerDelayLogEvent:
state.DelaySeconds = workerDelayLogEvent.DelaySeconds;
return;
case WorkerFailureLogEvent workerFailureLogEvent:
state.FailureSummary = workerFailureLogEvent.FailureSummary;
return;
}
}
}
Expand All @@ -78,15 +96,23 @@ public class WorkerDelayLogEvent : WorkerEventLog
[Id(0)] public int DelaySeconds;
}

[GenerateSerializer]
public class WorkerFailureLogEvent : WorkerEventLog
{
[Id(0)] public string FailureSummary;
}


public interface IWorkerGAgent : IStateGAgent<WorkerState>
{
Task SetDelayWorkAsync(int delaySeconds);
Task SetFailureSummary(string failureSummary);
}

[GenerateSerializer]
public class WorkerState : GroupMemberState
{
[Id(0)] public List<string> PreWorkUnits = new List<string>();
[Id(1)] public int DelaySeconds { get; set; } = 0;
[Id(2)] public string FailureSummary { get; set; }
}
91 changes: 91 additions & 0 deletions test/Aevatar.GAgents.GroupChat.Test/Tests/GroupChatWorkflowTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
using Aevatar.GAgents.GroupChat.WorkflowCoordinator.GEvent;
using Aevatar.GAgents.GroupChat.Test.GAgents;
using Aevatar.GAgents.GroupChat.WorkflowCoordinator;
using GroupChat.GAgent.Feature.Coordinator.GEvent;
using Shouldly;

namespace Aevatar.GAgents.GroupChat.Test.Tests;
Expand Down Expand Up @@ -585,4 +586,94 @@ await workflowCoordinator.ConfigAsync(new WorkflowCoordinatorConfigDto()
workflowState = await workflowCoordinator.GetStateAsync();
workflowState.WorkflowStatus.ShouldBe(WorkflowCoordinatorStatus.Pending);
}

[Fact]
public async Task Workflow_WithExecutionRecord_Succeeds_ShouldCreateAndCompleteRecord()
{
var toni = await _agentFactory.GetGAgentAsync<IWorkerGAgent>(Guid.NewGuid());
await toni.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Toni" });
var leader = await _agentFactory.GetGAgentAsync<ILeaderGAgent>(Guid.NewGuid());
await leader.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Leader" });

var groupAgent = await _agentFactory.GetGAgentAsync<IGroupGAgent>(Guid.NewGuid());
var workflows = new List<WorkflowUnitDto>()
{
new WorkflowUnitDto() { GrainId = toni.GetGrainId().ToString(), NextGrainId = leader.GetGrainId().ToString() },
new WorkflowUnitDto() { GrainId = leader.GetGrainId().ToString(), NextGrainId = "" }
};

var coordinator = await _agentFactory.GetGAgentAsync<IWorkflowCoordinatorGAgent>(Guid.NewGuid());
await coordinator.ConfigAsync(new WorkflowCoordinatorConfigDto
{
WorkflowUnitList = workflows,
InitContent = "init",
EnableExecutionRecord = true
});

await groupAgent.RegisterAsync(coordinator);
await groupAgent.PublishEventAsync(new StartWorkflowCoordinatorEvent());

// Wait until execution record id is assigned
Guid recordId = Guid.Empty;
for (int i = 0; i < 20; i++)
{
var s = await coordinator.GetStateAsync();
recordId = s.CurrentExecutionRecordId;
if (recordId != Guid.Empty) break;
await Task.Delay(100);
}
recordId.ShouldNotBe(Guid.Empty);

// Wait for workflow to finish and record to be unregistered
await Task.Delay(TimeSpan.FromSeconds(2));
var state = await coordinator.GetStateAsync();
state.WorkflowStatus.ShouldBe(WorkflowCoordinatorStatus.Pending);
state.CurrentExecutionRecordId.ShouldBe(Guid.Empty);
}

[Fact]
public async Task Workflow_WithExecutionRecord_Failure_ShouldMarkFailed()
{
var toni = await _agentFactory.GetGAgentAsync<IWorkerGAgent>(Guid.NewGuid());
await toni.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Toni" });
await toni.SetFailureSummary("Exception");
var leader = await _agentFactory.GetGAgentAsync<ILeaderGAgent>(Guid.NewGuid());
await leader.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Leader" });

var groupAgent = await _agentFactory.GetGAgentAsync<IGroupGAgent>(Guid.NewGuid());
var workflows = new List<WorkflowUnitDto>()
{
new WorkflowUnitDto() { GrainId = toni.GetGrainId().ToString(), NextGrainId = leader.GetGrainId().ToString() },
new WorkflowUnitDto() { GrainId = leader.GetGrainId().ToString(), NextGrainId = "" }
};

var coordinator = await _agentFactory.GetGAgentAsync<IWorkflowCoordinatorGAgent>(Guid.NewGuid());
await coordinator.ConfigAsync(new WorkflowCoordinatorConfigDto
{
WorkflowUnitList = workflows,
InitContent = "init",
EnableExecutionRecord = true
});

await groupAgent.RegisterAsync(coordinator);
await groupAgent.PublishEventAsync(new StartWorkflowCoordinatorEvent());

// Wait a bit for first unit to be activated and term to be set
await Task.Delay(500);
// var cstateBefore = await coordinator.GetStateAsync();
// var currentTerm = cstateBefore.Term;
//
// await groupAgent.PublishEventAsync(new ChatResponseEvent
// {
// BlackboardId = cstateBefore.BlackboardId,
// MemberId = toni.GetPrimaryKey(),
// MemberName = "Toni",
// FailureSummary = "boom",
// Term = currentTerm
// });
await Task.Delay(1000);

var cstate = await coordinator.GetStateAsync();
cstate.WorkflowStatus.ShouldBe(WorkflowCoordinatorStatus.Failed);
}
}
Loading
Loading