refactor: split node agent and Go control plane
This commit is contained in:
@@ -0,0 +1,113 @@
|
||||
namespace WxAgent.Core;
|
||||
|
||||
public enum MessageEventKind
|
||||
{
|
||||
MessageReceived,
|
||||
Reconnected
|
||||
}
|
||||
|
||||
public sealed record MessageEvent(
|
||||
string EventId,
|
||||
MessageEventKind Kind,
|
||||
string Session,
|
||||
ChatMessageSnapshot? Message,
|
||||
DateTimeOffset ObservedAt,
|
||||
bool Recovered);
|
||||
|
||||
public sealed record ListenerCheckpoint(
|
||||
int Version,
|
||||
string Session,
|
||||
DateTimeOffset UpdatedAt,
|
||||
IReadOnlyList<string> Fingerprints);
|
||||
|
||||
public sealed class MessageEventState
|
||||
{
|
||||
private readonly int _capacity;
|
||||
private readonly Queue<string> _order = new();
|
||||
private readonly HashSet<string> _seen = new(StringComparer.Ordinal);
|
||||
|
||||
public MessageEventState(int capacity = 4096, ListenerCheckpoint? checkpoint = null)
|
||||
{
|
||||
ArgumentOutOfRangeException.ThrowIfLessThan(capacity, 1);
|
||||
_capacity = capacity;
|
||||
if (checkpoint is null)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
foreach (var fingerprint in checkpoint.Fingerprints.TakeLast(capacity))
|
||||
{
|
||||
Add(fingerprint);
|
||||
}
|
||||
}
|
||||
|
||||
public MessageEvent? TryCreateMessage(
|
||||
string session,
|
||||
ChatMessageSnapshot message,
|
||||
DateTimeOffset observedAt,
|
||||
bool recovered)
|
||||
{
|
||||
ArgumentException.ThrowIfNullOrWhiteSpace(session);
|
||||
if (!Add(message.Fingerprint))
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
return new MessageEvent(
|
||||
message.Fingerprint,
|
||||
MessageEventKind.MessageReceived,
|
||||
session,
|
||||
message,
|
||||
observedAt,
|
||||
recovered);
|
||||
}
|
||||
|
||||
public ListenerCheckpoint CreateCheckpoint(string session, DateTimeOffset updatedAt) =>
|
||||
new(2, session, updatedAt, _order.ToArray());
|
||||
|
||||
private bool Add(string fingerprint)
|
||||
{
|
||||
if (!_seen.Add(fingerprint))
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
_order.Enqueue(fingerprint);
|
||||
if (_order.Count > _capacity)
|
||||
{
|
||||
_seen.Remove(_order.Dequeue());
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
public sealed record MessageCallbackFailure(int CallbackIndex, string ErrorType);
|
||||
|
||||
public static class MessageCallbackDispatcher
|
||||
{
|
||||
public static async Task<IReadOnlyList<MessageCallbackFailure>> DispatchAsync(
|
||||
IReadOnlyList<Func<MessageEvent, CancellationToken, Task>> callbacks,
|
||||
MessageEvent messageEvent,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
var failures = new List<MessageCallbackFailure>();
|
||||
for (var index = 0; index < callbacks.Count; index++)
|
||||
{
|
||||
try
|
||||
{
|
||||
await callbacks[index](messageEvent, cancellationToken);
|
||||
}
|
||||
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
failures.Add(new MessageCallbackFailure(index, exception.GetType().Name));
|
||||
}
|
||||
}
|
||||
|
||||
return failures;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user