diff --git a/src/Aevatar.GAgents.GroupChat.Core/Events/CoordinatorEvent.cs b/src/Aevatar.GAgents.GroupChat.Core/Events/CoordinatorEvent.cs index c08f858d..426b02bc 100644 --- a/src/Aevatar.GAgents.GroupChat.Core/Events/CoordinatorEvent.cs +++ b/src/Aevatar.GAgents.GroupChat.Core/Events/CoordinatorEvent.cs @@ -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] diff --git a/src/Aevatar.GAgents.GroupChat.Core/Models/WorkUnitExecutionStatus.cs b/src/Aevatar.GAgents.GroupChat.Core/Models/WorkUnitExecutionStatus.cs index 5a9eacd8..3c2099b7 100644 --- a/src/Aevatar.GAgents.GroupChat.Core/Models/WorkUnitExecutionStatus.cs +++ b/src/Aevatar.GAgents.GroupChat.Core/Models/WorkUnitExecutionStatus.cs @@ -4,5 +4,6 @@ public enum WorkflowExecutionStatus { Pending, Running, - Completed + Completed, + Failed } \ No newline at end of file diff --git a/src/Aevatar.GAgents.GroupChat.Core/States/WorkflowExecutionRecordState.cs b/src/Aevatar.GAgents.GroupChat.Core/States/WorkflowExecutionRecordState.cs index ab29b2b8..b94736d5 100644 --- a/src/Aevatar.GAgents.GroupChat.Core/States/WorkflowExecutionRecordState.cs +++ b/src/Aevatar.GAgents.GroupChat.Core/States/WorkflowExecutionRecordState.cs @@ -39,4 +39,6 @@ public class WorkUnitExecutionRecord public string InputData { get; set; } [Id(5)] public string OutputData { get; set; } + + [Id(6)] public string FailureSummary { get; set; } } \ No newline at end of file diff --git a/src/Aevatar.GAgents.GroupChat.GroupMember/GAgent/MemberGAgentBase.cs b/src/Aevatar.GAgents.GroupChat.GroupMember/GAgent/MemberGAgentBase.cs index b502ed89..6cd8a42d 100644 --- a/src/Aevatar.GAgents.GroupChat.GroupMember/GAgent/MemberGAgentBase.cs +++ b/src/Aevatar.GAgents.GroupChat.GroupMember/GAgent/MemberGAgentBase.cs @@ -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; @@ -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] diff --git a/src/Aevatar.GAgents.GroupChat/GAgent/Coordinator/WorkflowView/WorkflowViewGAgent.cs b/src/Aevatar.GAgents.GroupChat/GAgent/Coordinator/WorkflowView/WorkflowViewGAgent.cs index cc040fca..94e02c22 100644 --- a/src/Aevatar.GAgents.GroupChat/GAgent/Coordinator/WorkflowView/WorkflowViewGAgent.cs +++ b/src/Aevatar.GAgents.GroupChat/GAgent/Coordinator/WorkflowView/WorkflowViewGAgent.cs @@ -98,6 +98,8 @@ private async Task TrySaveWorkflowViewAsync(WorkflowViewConfigDto configuration) AgentId = configuration.WorkflowCoordinatorGAgentId }); } + + await ConfirmEvents(); } protected override void GAgentTransitionState(WorkflowViewState state, diff --git a/src/Aevatar.GAgents.GroupChat/GroupMemberGAgentBase.cs b/src/Aevatar.GAgents.GroupChat/GroupMemberGAgentBase.cs index f0ca0f9a..eed8565e 100644 --- a/src/Aevatar.GAgents.GroupChat/GroupMemberGAgentBase.cs +++ b/src/Aevatar.GAgents.GroupChat/GroupMemberGAgentBase.cs @@ -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] diff --git a/src/Aevatar.GAgents.GroupChat/WorkflowCoordinatorGAgent.cs b/src/Aevatar.GAgents.GroupChat/WorkflowCoordinatorGAgent.cs index 6fda7a35..f792b549 100644 --- a/src/Aevatar.GAgents.GroupChat/WorkflowCoordinatorGAgent.cs +++ b/src/Aevatar.GAgents.GroupChat/WorkflowCoordinatorGAgent.cs @@ -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(State.BlackboardId); await blackboard.SetMessageAsync(new CoordinatorConfirmChatResponse() { @@ -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; diff --git a/src/Aevatar.GAgents.GroupChat/WorkflowExecutionRecordGAgent.cs b/src/Aevatar.GAgents.GroupChat/WorkflowExecutionRecordGAgent.cs index 2a32e0c0..f2c15466 100644 --- a/src/Aevatar.GAgents.GroupChat/WorkflowExecutionRecordGAgent.cs +++ b/src/Aevatar.GAgents.GroupChat/WorkflowExecutionRecordGAgent.cs @@ -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(); } @@ -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; } } } @@ -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; } } \ No newline at end of file diff --git a/test/Aevatar.GAgents.GroupChat.Test/Aevatar.GAgents.GroupChat.Test.csproj b/test/Aevatar.GAgents.GroupChat.Test/Aevatar.GAgents.GroupChat.Test.csproj index 73ca4eaa..e90a2f4e 100644 --- a/test/Aevatar.GAgents.GroupChat.Test/Aevatar.GAgents.GroupChat.Test.csproj +++ b/test/Aevatar.GAgents.GroupChat.Test/Aevatar.GAgents.GroupChat.Test.csproj @@ -28,6 +28,7 @@ + diff --git a/test/Aevatar.GAgents.GroupChat.Test/GAgents/WorkerGAgent.cs b/test/Aevatar.GAgents.GroupChat.Test/GAgents/WorkerGAgent.cs index de224a27..664b211f 100644 --- a/test/Aevatar.GAgents.GroupChat.Test/GAgents/WorkerGAgent.cs +++ b/test/Aevatar.GAgents.GroupChat.Test/GAgents/WorkerGAgent.cs @@ -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; @@ -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 GetState() { return Task.FromResult(State); @@ -39,6 +49,11 @@ protected override async Task ChatAsync(Guid blackboardId, Lists.AgentName).ToList()}); @@ -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; } } } @@ -78,10 +96,17 @@ public class WorkerDelayLogEvent : WorkerEventLog [Id(0)] public int DelaySeconds; } +[GenerateSerializer] +public class WorkerFailureLogEvent : WorkerEventLog +{ + [Id(0)] public string FailureSummary; +} + public interface IWorkerGAgent : IStateGAgent { Task SetDelayWorkAsync(int delaySeconds); + Task SetFailureSummary(string failureSummary); } [GenerateSerializer] @@ -89,4 +114,5 @@ public class WorkerState : GroupMemberState { [Id(0)] public List PreWorkUnits = new List(); [Id(1)] public int DelaySeconds { get; set; } = 0; + [Id(2)] public string FailureSummary { get; set; } } \ No newline at end of file diff --git a/test/Aevatar.GAgents.GroupChat.Test/Tests/GroupChatWorkflowTest.cs b/test/Aevatar.GAgents.GroupChat.Test/Tests/GroupChatWorkflowTest.cs index 2bb81ccd..ce4ed6b2 100644 --- a/test/Aevatar.GAgents.GroupChat.Test/Tests/GroupChatWorkflowTest.cs +++ b/test/Aevatar.GAgents.GroupChat.Test/Tests/GroupChatWorkflowTest.cs @@ -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; @@ -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(Guid.NewGuid()); + await toni.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Toni" }); + var leader = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await leader.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Leader" }); + + var groupAgent = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var workflows = new List() + { + new WorkflowUnitDto() { GrainId = toni.GetGrainId().ToString(), NextGrainId = leader.GetGrainId().ToString() }, + new WorkflowUnitDto() { GrainId = leader.GetGrainId().ToString(), NextGrainId = "" } + }; + + var coordinator = await _agentFactory.GetGAgentAsync(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(Guid.NewGuid()); + await toni.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Toni" }); + await toni.SetFailureSummary("Exception"); + var leader = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await leader.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Leader" }); + + var groupAgent = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var workflows = new List() + { + new WorkflowUnitDto() { GrainId = toni.GetGrainId().ToString(), NextGrainId = leader.GetGrainId().ToString() }, + new WorkflowUnitDto() { GrainId = leader.GetGrainId().ToString(), NextGrainId = "" } + }; + + var coordinator = await _agentFactory.GetGAgentAsync(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); + } } \ No newline at end of file diff --git a/test/Aevatar.GAgents.GroupChat.Test/Tests/WorkflowExecutionRecordGAgentTests.cs b/test/Aevatar.GAgents.GroupChat.Test/Tests/WorkflowExecutionRecordGAgentTests.cs index fb833ebf..22dd16b2 100644 --- a/test/Aevatar.GAgents.GroupChat.Test/Tests/WorkflowExecutionRecordGAgentTests.cs +++ b/test/Aevatar.GAgents.GroupChat.Test/Tests/WorkflowExecutionRecordGAgentTests.cs @@ -1,13 +1,21 @@ using Aevatar.Core.Abstractions; +using Aevatar.Core; using Aevatar.GAgents.Basic.BasicGAgents.GroupGAgent; using Aevatar.GAgents.GroupChat.Core; +using Aevatar.GAgents.GroupChat.Core.Dto; +using Aevatar.GAgents.GroupChat.Core.States; using Aevatar.GAgents.GroupChat.Test.GAgents; using Aevatar.GAgents.GroupChat.WorkflowCoordinator; +using Aevatar.GAgents.GroupChat.WorkflowCoordinator.Dto; using Aevatar.GAgents.GroupChat.WorkflowCoordinator.GEvent; +using GroupChat.GAgent; using GroupChat.GAgent.Feature.Common; using GroupChat.GAgent.Feature.Coordinator.GEvent; using Newtonsoft.Json; using Shouldly; +using Aevatar.GAgents.InputGAgent.GAgent; +using Aevatar.GAgents.InputGAgent.Dto; +using Volo.Abp; namespace Aevatar.GAgents.GroupChat.Test.Tests; @@ -141,6 +149,40 @@ public async Task Handle_ChatResponseEvent_Test() grainARecord.OutputData.ShouldBe(JsonConvert.SerializeObject(finishExecuteGrainA.ChatResponse.Content)); } + [Fact] + public async Task Handle_ChatResponseEvent_Failure_Test() + { + var groupAgent = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var recordAgent = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await groupAgent.RegisterAsync(recordAgent); + var workerGrainId = groupAgent.GetGrainId(); + + await StartExecuteWorkflowAsync(groupAgent, workerGrainId); + + var startExecuteGrain = new StartExecuteWorkUnitEvent + { + WorkUnitGrainId = workerGrainId.ToString(), + CoordinatorMessages = new List { new ChatMessage { Content = "Input A" } } + }; + await groupAgent.PublishEventAsync(startExecuteGrain); + await Task.Delay(500); + + var failure = new ChatResponseEvent + { + PublisherGrainId = workerGrainId, + FailureSummary = "unit crashed" + }; + await groupAgent.PublishEventAsync(failure); + await Task.Delay(1000); + + var state = await recordAgent.GetStateAsync(); + state.Status.ShouldBe(WorkflowExecutionStatus.Failed); + var unit = state.WorkUnitRecords.First(o => o.WorkUnitGrainId == workerGrainId.ToString()); + unit.Status.ShouldBe(WorkflowExecutionStatus.Failed); + unit.FailureSummary.ShouldBe("unit crashed"); + unit.EndTime.ShouldNotBeNull(); + } + [Fact] public async Task IncorrectSequence_Test() { @@ -180,7 +222,7 @@ public async Task IncorrectSequence_Test() grainARecord.Status.ShouldBe(WorkflowExecutionStatus.Completed); grainARecord.InputData.ShouldBe(JsonConvert.SerializeObject(startExecuteGrain.CoordinatorMessages)); } - + private async Task StartExecuteWorkflowAsync(IGroupGAgent groupAgent, GrainId workerGrainId) { var startExecuteWorkflowEvent = new StartExecuteWorkflowEvent @@ -200,4 +242,362 @@ private async Task StartExecuteWorkflowAsync(IGroupGAgent groupAgent, GrainId wo await groupAgent.PublishEventAsync(startExecuteWorkflowEvent); await Task.Delay(1000); } + + [Fact] + public async Task Member_GetDescription_ShouldIncludeMemberName() + { + var member = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await member.ConfigAsync(new GroupMemberConfigDto { MemberName = "Alice" }); + + var desc = await member.DescribeAsync(); + desc.ShouldContain("Member Name: Alice"); + } + + [Fact] + public async Task Member_Ping_ShouldPublishPong() + { + var group = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var member = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await member.ConfigAsync(new GroupMemberConfigDto { MemberName = "Pingy" }); + var collector = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + + await group.RegisterAsync(member); + await group.RegisterAsync(collector); + + var blackboardId = Guid.NewGuid(); + await group.PublishEventAsync(new CoordinatorPingEvent { BlackboardId = blackboardId }); + await Task.Delay(500); + + var s = await collector.GetStateAsync(); + s.LastPongBlackboardId.ShouldBe(blackboardId); + s.LastPongMemberName.ShouldBe("Pingy"); + s.LastPongMemberId.ShouldBe(member.GetGrainId().GetGuidKey()); + } + + [Fact] + public async Task Member_EvaluationInterest_ShouldPublishResponse() + { + var group = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var member = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await member.ConfigAsync(new GroupMemberConfigDto { MemberName = "Eva" }); + var collector = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + + await group.RegisterAsync(member); + await group.RegisterAsync(collector); + + var blackboardId = Guid.NewGuid(); + await group.PublishEventAsync(new EvaluationInterestEvent { BlackboardId = blackboardId, ChatTerm = 123 }); + await Task.Delay(500); + + var s = await collector.GetStateAsync(); + s.LastInterestBlackboardId.ShouldBe(blackboardId); + s.LastInterestMemberId.ShouldBe(member.GetGrainId().GetGuidKey()); + s.LastInterestValue.ShouldBe(77); + s.LastInterestChatTerm.ShouldBe(123); + } + + [Fact] + public async Task Member_GetMessageFromBlackboard_ShouldReturnContent() + { + var member = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await member.ConfigAsync(new GroupMemberConfigDto { MemberName = "Reader" }); + + var blackboardId = Guid.NewGuid(); + var blackboard = await _agentFactory.GetGAgentAsync(blackboardId); + await blackboard.SetTopic("topic-x"); + + var msgs = await member.FetchBlackboardMessages(blackboardId); + msgs.ShouldNotBeNull(); + msgs.Any(m => m.MessageType == MessageType.BlackboardTopic && m.Content == "topic-x").ShouldBeTrue(); + } + + [Fact] + public async Task InputGAgent_ChatEvent_SpeakerMatch_ShouldPublishChatResponse() + { + var group = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var input = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var collector = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + + await group.RegisterAsync(input); + await group.RegisterAsync(collector); + + await input.ConfigAsync(new InputConfigDto { MemberName = "Inny", Input = "Hello" }); + + var blackboardId = Guid.NewGuid(); + await group.PublishEventAsync(new ChatEvent + { + BlackboardId = blackboardId, + Speaker = input.GetGrainId().GetGuidKey(), + Term = 1, + CoordinatorMessages = new List { new ChatMessage { Content = "start" } } + }); + + await Task.Delay(500); + + var s = await collector.GetStateAsync(); + s.LastChatBlackboardId.ShouldBe(blackboardId); + s.LastChatMemberId.ShouldBe(input.GetGrainId().GetGuidKey()); + s.LastChatMemberName.ShouldBe("Inny"); + s.LastChatContent.ShouldBe("Hello"); + s.LastChatTerm.ShouldBe(1); + s.LastChatFailure.ShouldBeNull(); + } + + [Fact] + public async Task InputGAgent_ChatEvent_SpeakerMismatch_ShouldBeIgnored() + { + var group = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var input = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var collector = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + + await group.RegisterAsync(input); + await group.RegisterAsync(collector); + + await input.ConfigAsync(new InputConfigDto { MemberName = "Inny", Input = "Hello" }); + + var blackboardId = Guid.NewGuid(); + await group.PublishEventAsync(new ChatEvent + { + BlackboardId = blackboardId, + Speaker = Guid.NewGuid(), + Term = 2, + CoordinatorMessages = new List { new ChatMessage { Content = "start" } } + }); + + await Task.Delay(500); + + var s = await collector.GetStateAsync(); + s.LastChatBlackboardId.ShouldNotBe(blackboardId); + s.LastChatContent.ShouldBeNull(); + } + + [Fact] + public async Task InputGAgent_EvaluationInterest_ShouldPublish100() + { + var group = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var input = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var collector = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + + await group.RegisterAsync(input); + await group.RegisterAsync(collector); + + await input.ConfigAsync(new InputConfigDto { MemberName = "Scorer", Input = "whatever" }); + + var blackboardId = Guid.NewGuid(); + await group.PublishEventAsync(new EvaluationInterestEvent { BlackboardId = blackboardId, ChatTerm = 9 }); + await Task.Delay(500); + + var s = await collector.GetStateAsync(); + s.LastInterestBlackboardId.ShouldBe(blackboardId); + s.LastInterestMemberId.ShouldBe(input.GetGrainId().GetGuidKey()); + s.LastInterestValue.ShouldBe(100); + s.LastInterestChatTerm.ShouldBe(9); + } + + [Fact] + public async Task Workflow_WithExecutionRecord_Failure_ShouldMarkFailed() + { + var toni = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await toni.ConfigAsync(new InputConfigDto { MemberName = "Scorer", Input = "whatever" }); + var leader = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + await leader.ConfigAsync(new GroupMemberConfigDto() { MemberName = "Leader" }); + + var groupAgent = await _agentFactory.GetGAgentAsync(Guid.NewGuid()); + var workflows = new List() + { + new WorkflowUnitDto() { GrainId = toni.GetGrainId().ToString(), NextGrainId = leader.GetGrainId().ToString() }, + new WorkflowUnitDto() { GrainId = leader.GetGrainId().ToString(), NextGrainId = "" } + }; + + var coordinator = await _agentFactory.GetGAgentAsync(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(1500); + + var cstate = await coordinator.GetStateAsync(); + cstate.WorkflowStatus.ShouldBe(WorkflowCoordinatorStatus.Failed); + } + + // Test helper agents for coverage + [GAgent(nameof(EventCollectorGAgent))] + public class EventCollectorGAgent : GroupMemberGAgentBase, IEventCollectorGAgent + { + public override Task GetDescriptionAsync() => Task.FromResult("collector"); + + protected override Task GetInterestValueAsync(Guid blackboardId) => Task.FromResult(0); + + protected override Task ChatAsync(Guid blackboardId, List? coordinatorMessages) + { + return Task.FromResult(new ChatResponse { Content = "noop" }); + } + + [EventHandler] + public async Task HandleEventAsync(CoordinatorPongEvent @event) + { + RaiseEvent(new SetPongLogEvent + { + BlackboardId = @event.BlackboardId, + MemberId = @event.MemberId, + MemberName = @event.MemberName + }); + await ConfirmEvents(); + } + + [EventHandler] + public async Task HandleEventAsync(EvaluationInterestResponseEvent @event) + { + RaiseEvent(new SetInterestLogEvent + { + BlackboardId = @event.BlackboardId, + MemberId = @event.MemberId, + InterestValue = @event.InterestValue, + ChatTerm = @event.ChatTerm + }); + await ConfirmEvents(); + } + + [EventHandler] + public async Task HandleEventAsync(ChatResponseEvent @event) + { + RaiseEvent(new SetChatLogEvent + { + BlackboardId = @event.BlackboardId, + MemberId = @event.MemberId, + MemberName = @event.MemberName, + Content = @event.ChatResponse?.Content, + Term = @event.Term, + Failure = @event.FailureSummary + }); + await ConfirmEvents(); + } + + protected override void GroupMemberTransitionState(CollectorState state, StateLogEventBase @event) + { + switch (@event) + { + case SetPongLogEvent e: + state.LastPongBlackboardId = e.BlackboardId; + state.LastPongMemberId = e.MemberId; + state.LastPongMemberName = e.MemberName; + return; + case SetInterestLogEvent e2: + state.LastInterestBlackboardId = e2.BlackboardId; + state.LastInterestMemberId = e2.MemberId; + state.LastInterestValue = e2.InterestValue; + state.LastInterestChatTerm = e2.ChatTerm; + return; + case SetChatLogEvent e3: + state.LastChatBlackboardId = e3.BlackboardId; + state.LastChatMemberId = e3.MemberId; + state.LastChatMemberName = e3.MemberName; + state.LastChatContent = e3.Content; + state.LastChatTerm = e3.Term; + state.LastChatFailure = e3.Failure; + return; + } + } + } + + public interface IEventCollectorGAgent : IStateGAgent + { + } + + [GenerateSerializer] + public class CollectorLogEvent : StateLogEventBase + { + } + + [GenerateSerializer] + public class SetPongLogEvent : CollectorLogEvent + { + [Id(0)] public Guid BlackboardId { get; set; } + [Id(1)] public Guid MemberId { get; set; } + [Id(2)] public string MemberName { get; set; } + } + + [GenerateSerializer] + public class SetInterestLogEvent : CollectorLogEvent + { + [Id(0)] public Guid BlackboardId { get; set; } + [Id(1)] public Guid MemberId { get; set; } + [Id(2)] public int InterestValue { get; set; } + [Id(3)] public long ChatTerm { get; set; } + } + + [GenerateSerializer] + public class SetChatLogEvent : CollectorLogEvent + { + [Id(0)] public Guid BlackboardId { get; set; } + [Id(1)] public Guid MemberId { get; set; } + [Id(2)] public string MemberName { get; set; } + [Id(3)] public string Content { get; set; } + [Id(4)] public long Term { get; set; } + [Id(5)] public string Failure { get; set; } + } + + [GenerateSerializer] + public class CollectorState : WorkerState + { + [Id(10)] public Guid LastPongBlackboardId { get; set; } + [Id(11)] public Guid LastPongMemberId { get; set; } + [Id(12)] public string LastPongMemberName { get; set; } + [Id(13)] public Guid LastInterestBlackboardId { get; set; } + [Id(14)] public Guid LastInterestMemberId { get; set; } + [Id(15)] public int LastInterestValue { get; set; } + [Id(16)] public long LastInterestChatTerm { get; set; } + [Id(17)] public Guid LastChatBlackboardId { get; set; } + [Id(18)] public Guid LastChatMemberId { get; set; } + [Id(19)] public string LastChatMemberName { get; set; } + [Id(20)] public string LastChatContent { get; set; } + [Id(21)] public long LastChatTerm { get; set; } + [Id(22)] public string LastChatFailure { get; set; } + } + + [GAgent(nameof(TestMemberHelperGAgent))] + public class TestMemberHelperGAgent : GroupMemberGAgentBase, ITestMemberHelperGAgent + { + protected override Task GetInterestValueAsync(Guid blackboardId) => Task.FromResult(77); + + protected override Task ChatAsync(Guid blackboardId, List? coordinatorMessages) + { + return Task.FromResult(new ChatResponse { Content = "ok" }); + } + + public Task> FetchBlackboardMessages(Guid blackboardId) => GetMessageFromBlackboardAsync(blackboardId); + + public Task DescribeAsync() => GetDescriptionAsync(); + } + + public interface ITestMemberHelperGAgent : IStateGAgent + { + Task> FetchBlackboardMessages(Guid blackboardId); + Task DescribeAsync(); + } + + [GenerateSerializer] + public class TestMemberEventLog : StateLogEventBase + { + } + + [GAgent(nameof(FailInputGAgent))] + public class FailInputGAgent : InputGAgent.GAgent.InputGAgent, IFailInputGAgent + { + protected override Task ChatAsync(Guid blackboardId, List? messages) + { + throw new UserFriendlyException("InputGAgent fail"); + } + } + + public interface IFailInputGAgent : IInputGAgent + { + } } \ No newline at end of file