Cleaned up request acks. Added internal service bus for internal messaging.

This commit is contained in:
Tom
2024-11-08 15:32:42 +00:00
parent fe2eb86a08
commit 66f2bf7ec6
33 changed files with 1326 additions and 415 deletions

View File

@@ -0,0 +1,41 @@
using System.Reactive;
namespace TwitchChatTTS.Bus
{
public class ServiceBusObservable : ObservableBase<ServiceBusData>
{
private readonly string _topic;
private readonly ServiceBusCentral _central;
public ServiceBusObservable(string topic, ServiceBusCentral central)
{
_topic = topic;
_central = central;
}
protected override IDisposable SubscribeCore(IObserver<ServiceBusData> observer)
{
_central.Add(_topic, observer);
return new ServiceBusUnsubscriber(_topic, _central, observer);
}
private sealed class ServiceBusUnsubscriber : IDisposable
{
private readonly string _topic;
private readonly ServiceBusCentral _central;
private readonly IObserver<ServiceBusData> _receiver;
public ServiceBusUnsubscriber(string topic, ServiceBusCentral central, IObserver<ServiceBusData> receiver)
{
_topic = topic;
_central = central;
_receiver = receiver;
}
public void Dispose()
{
_central.RemoveObserver(_topic, _receiver);
}
}
}
}