C#. Прослушка данных через TcpClient. Поток данных с буфера чтения в IObservable
есть задача клиент подключается к Tcp socket сервера и слушает его, получая пакеты данных.
Также нужно отправлять запросы на сервер и получать в ответ запрашиваемые данные.
Т.к. основной режим это слушать данные, я решил сделать отдельный поток ОТВЕТОВ через Observable.Create
При успешном подключении к серверу выдается поток ответов (IObservable<object> - object приведен для макета).
Сервисы уровнем выше вешают обработчики (where фильтры, select преобразования и т.д.) на этот поток, и в конце сервис подписывается на IObservable и получает уже обработанные данные.
Вся цепочка строится Lazy, т.е. материализуется только после подписки - что мне очень удобно, т.к. на это будет отдельная команда.
Использование Observable.Create - дает преимущество:
- при отписке на стороне подписчика - срабатывает токен отмены и завершается задача генерации последовательности.
- последующая подписка запускает генерацию заново.
КРАСОТА!
К это системе прослушки надо добавить систему запрос-ожидание-ответ.
Для ответов уже есть поток, запрос пошлем тоже независимо.
Осталось связать это в связку запрос-ожидание-ответ, т.е. сделали запрос, из потока ответов нашли нужный ответ за отведенное время ожидания, если ответ есть, true если не нашли ответ то false.
Т.е. метод отправки запроса должен вернуть флаг валидности ответа на этот запрос, а сами данные ответа летят дальше в потоке ответов.
ответ идеально было бы дожидаться через ResponseDataFlow.FirstOrDefaultAsync(...), но любая подписка на IObservable запускает новую задачу!! Т.к. мы используем Observable.Create и в параллельном потоке опрашивается буфер чтения сокета - что не верно.
public async Task<Result<OutData, ErrorAll>> SendCommandAndWaitAnswer()
{
// Отправить Запрос, и ожидать ответ на шине ResponseDataFlow.
await _tcpIpTransport.SendCommandExec();
var response = await ResponseDataFlow.FirstOrDefaultAsync(data => data.ToString() == "05"); // ожидаем ответ "05"
var outData = new OutData(response);
outData.Validate();// Примем решение о валидности ответа.
return ResultExt.Success(outData);
}
Мне нужно подписаться или как-то по другому получить данные из ТОГО же потока данных, не генерируя новую последовательность подпиской.
Я уже думал отказаться от Observable.Create и просто создать вручную Task опроса и дергать Rx событие создавая последовательность, но прелести запуска/останова задачи опроса на другом конце слоя сервисов большой плюс.
Еще вариант вручную вмешаться в поток ответов и создавать задачу TaskCompletionSource по ожиданию нужного ответа из потока- но это колхозно выглядит, через общие поля будут связанны методы SendCommandAndWaitAnswer() и ReConnect().
public class SystemTcpIpClient : BaseTcpIpClient
{
protected override async Task<Result<IObservable<object>, ErrorAll>> ReOpen(CancellationToken ct)
{
DisposeTransport();
try
{
_client = new TcpClient { NoDelay = false }; //true - пакет будет отправлен мгновенно (при NetworkStream.Write). false - пока не собранно значительное кол-во данных отправки не будет.
await _client.ConnectAsync(ipAddress, Address.Port);
_netStream = _client.GetStream();
return CreateNewDataFlow();
}
catch (Exception ex)
{
DisposeTransport();
return Result.Failure<IObservable<object>, ErrorAll>(ErrorPresets.TcpIpTransport.ReConnect(Address, StatusString));
}
}
protected override Result<IObservable<object>, ErrorAll> CreateNewDataFlow()
{
var dataFlow = GetResponseDataFlow();
return Result.Success<IObservable<object>, ErrorAll>(dataFlow);
}
private IObservable<object> GetResponseDataFlow()
{
return Observable.Create<object>(
(obs, ct) =>
{
return Task.Run(async () =>
{
while (!ct.IsCancellationRequested)
{
try
{
var (_, isFailure, data) = await ReadDataPolling();
if (isFailure)
{
obs.OnCompleted();
}
else
{
obs.OnNext(data);
}
}
catch (Exception ex)
{
obs.OnError(ex);
}
}
_logger?.Information("GetResponseDataFlow Loop terminated");
}, ct);
});
}
private async Task<Result<string>> ReadDataPolling()
{
await Task.Delay(100); //Время поллинга
int buferSize = 1024;// читаем всегда весь буфер
byte[] bDataTemp = new byte[buferSize];
var numberOfBytes = await _netStream!.ReadAsync(bDataTemp, 0, buferSize);
if(numberOfBytes == 0)
return Result.Failure<string>("Сервер закрыл подключение. Получен конец потока чтения данных с сокета.");
var realData= bDataTemp.Take(numberOfBytes);
var resStr= realData.ArrayByteToString("X2");//DEBUG
return resStr;
}
}
public class DataProvider4TcpIpService : IDataProvider
{
private readonly ITcpIpTransport _tcpIpTransport;
public DataProvider4TcpIpService(ITcpIpTransport tcpIpTransport)
{
_tcpIpTransport = tcpIpTransport;
}
public IObservable<object> ResponseDataFlow { get; private set; }
/// <summary>
/// После успешного коннекта получить поток данных.
/// </summary>
public async Task<Result<bool, ErrorAll>> ReConnect(CancellationToken ct)
{
var res= await _tcpIpTransport.CycleReOpenedExec(); //Циклический реконнект
if (res.IsSuccess)
{
ResponseDataFlow = res.Value;
}
return ResultExt.Success();
}
public async Task<Result<OutData, ErrorAll>> SendCommandAndWaitAnswer()
{
// Отправить Запрос, и ожидать ответ на шине ResponseDataFlow.
await _tcpIpTransport.SendCommandExec();
var response = await ResponseDataFlow.FirstOrDefaultAsync(data => data.ToString() == "05");
var outData = new OutData();
return ResultExt.Success(outData);
}
}
ВОПРОСЫ:
- Можно ли как-то реализовать эту задачу (Один издатель-много подписчиков на т.е. же данные) в классе DataProvider4TcpIpService, имея поток ответов?
- Есть ли готовые реализации прослушки веб сокетов реализующие IObservable, чтобы не использовать polling для опроса?
Дополнение: Основной сервис который использует поток ответов
public class ConfigDataFlowBg
{
private readonly DataProviderControllService _dataProviderControllService;
private readonly ILogger _logger;
private IDisposable _dataProviderFlowDataLifeTime; //контроль за потоком генерации данных
public ConfigDataFlowBg(
DataProviderControllService dataProviderControllService,
ILogger logger
)
{
_dataProviderControllService = dataProviderControllService;
_logger = logger;
}
/// <summary>
/// Подписка на поток получения данных.
/// </summary>
public async Task StartAsync(CancellationToken ct)
{
var dataProviderRes= await _dataProviderControllService.BuildDataProvider(new DataProviderId(TransportKey.TcpIp, "LocalTcpServer"));
var dp = dataProviderRes.Value;
var connectRes= await dp.ReConnect(ct);
if (connectRes.IsSuccess)
{
_dataProviderFlowDataLifeTime= dp.ResponseDataFlow
.TakeWhile(v => !ct.IsCancellationRequested)
.Subscribe( //подписка только в самом последнем модуле
o =>
{
//получение данных последним модулем.
_logger.Information(o.ToString());
},
exception =>
{
}, async () =>
{
await StartAsync(ct); //DEBUG для коннекта при обрыве связи.
//Разорвать всю цепочку обработки IObservable
//Выполнить ReConnect, получить новый IObservable и заново создать цепочку.
});
}
}
}