Files
wx-win-agent/node-agent/WxAgent.Core/MessageEvents.cs
T

114 lines
2.9 KiB
C#

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;
}
}