Files
laatzen/Common/NamedPipes/NamedPipe.cs
T
2023-11-09 08:59:09 +01:00

151 lines
3.8 KiB
C#

using System.ComponentModel;
using System.IO.Pipes;
using System;
using System.Threading;
namespace NamedPipes
{
public abstract class NamedPipe
{
protected const string PIPE_NAME = "NamedPipes";
protected const byte FLOW_REQUEST = 1;
protected const byte DATA_REQUEST = 2;
protected const byte MSG_LENGTH = 9;
protected PipeStream pipe;
private readonly BackgroundWorker worker;
public NamedPipe()
{
this.worker = new BackgroundWorker();
this.worker.DoWork += this.Run;
}
public event Action Connected;
public event Action Disconnected;
public event Action<double> DataRequest;
public event Action<double> FlowRequest;
public void Disconnect()
{
this.pipe?.Dispose();
this.worker?.Dispose();
}
public void Connect()
=> this.worker.RunWorkerAsync();
public double RequestReference()
{
var completed = false;
var reference = 0D;
var timeout = 60_000;
var delay = 1000;
void DataRequested(double val)
{
completed = true;
reference = val;
}
this.DataRequest += DataRequested;
this.SendMessage(reference, DATA_REQUEST);
while (!completed && timeout > 0)
{
Thread.Sleep(delay);
timeout -= delay;
}
this.FlowRequest -= DataRequested;
return reference;
}
public void RequestFlowRate(double value)
{
var completed = false;
var timeout = 60_000;
var delay = 1000;
void FlowRequested(double _) => completed = true;
this.FlowRequest += FlowRequested;
this.SendMessage(value, FLOW_REQUEST);
while (!completed && timeout > 0)
{
Thread.Sleep(delay);
timeout -= delay;
}
this.FlowRequest -= FlowRequested;
}
protected void SendMessage(double value, byte endpoint)
{
if (this.pipe.IsConnected)
{
var message = new byte[MSG_LENGTH];
message[0] = endpoint;
BitConverter
.GetBytes(value)
.CopyTo(message, 1);
this.pipe.Write(message, 0, message.Length);
}
}
protected virtual void Run(object sender, DoWorkEventArgs args)
{
this.Connected?.Invoke();
var buffer = new byte[byte.MaxValue];
while (this.pipe.IsConnected)
{
var asyncResult = this.pipe.BeginRead(buffer, 0, byte.MaxValue, this.EndRead, buffer);
if (asyncResult.AsyncWaitHandle.WaitOne())
{
continue;
}
else
{
break;
}
}
if (!this.pipe.IsConnected)
{
this.Disconnected?.Invoke();
}
}
private void EndRead(IAsyncResult asyncResult)
{
if (this.pipe.IsConnected && asyncResult.AsyncState is byte[] buffer)
{
var readLength = this.pipe.EndRead(asyncResult);
if (readLength > 0 && buffer.Length >= readLength)
{
var value = BitConverter.ToDouble(buffer, 1);
switch (buffer[0])
{
case FLOW_REQUEST: this.FlowRequest?.Invoke(value); break;
case DATA_REQUEST: this.DataRequest?.Invoke(value); break;
}
}
}
}
}
}