using System;
using System.Collections.Concurrent;
using System.IO;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using TeamAAS.Feeder.Attributes;
using TeamAAS.Feeder.Interfaces;
using TeamAAS.Feeder.Enums;
using TeamAAS.Feeder.Models;
namespace TeamAAS.Feeder.Devices
{
///
/// TEAM 品牌供料器 TCP 协议实现
/// 协议格式:{命令}\r\n,ASCII 编码
/// 用 .NET 内置 System.Net.Sockets,无第三方 TCP 依赖
///
[Feeder("Team 供料器", FeederBrand.Team, Description = "内置 Team 品牌 TCP 协议实现")]
public class TeamFeeder : IFeeder
{
private const int MillisecondsTimeout = 5000;
private const int ReconnectIntervalMs = 3000;
private const string Terminator = "\r\n";
// 系统参数全局地址
private const int AddrBackLightFlash = 33; // 背光闪烁时间
private const int AddrLightBrightness = 34; // 光源亮度
private const int AddrLightTimeout = 44; // 背光超时
private const int AddrPlatShakeTimeout = 65; // 平台振动超时
private const int AddrHopperShakeTimeout = 66;// 料斗振动超时
private const int HopperParamStartAddress = 742; // 料斗输出参数起始地址(每组8个)
private const int InputTriggerStartAddress = 39; // 输入触发序列ID起始地址
private readonly object _ioLock = new object();
private TcpClient _tcp;
private NetworkStream _stream;
private Thread _recvThread;
private volatile bool _recvRunning;
private readonly ConcurrentQueue _rxQueue = new ConcurrentQueue();
private readonly StringBuilder _rxBuffer = new StringBuilder();
private bool _isBacklightOn;
private System.Threading.Timer _reconnectTimer;
private volatile bool _manualDisconnect;
public Guid Id { get; }
public int FeederNo { get; }
public string Name { get; }
public FeederBrand Brand => FeederBrand.Team;
public bool IsConnected => _tcp?.Connected ?? false;
public bool IsBacklightOn => _isBacklightOn;
public event Action ConnectionChanged;
public event Action MessageExchanged;
public TeamFeeder(Guid id, int feederNo, string name, string ip, int port)
{
Id = id;
FeederNo = feederNo;
Name = name;
IP = ip;
Port = port;
}
private string IP { get; }
private int Port { get; }
public async Task ConnectAsync()
{
_manualDisconnect = false;
StopReconnectTimer();
DisposeStream();
_tcp = new TcpClient();
await _tcp.ConnectAsync(IP, Port);
_stream = _tcp.GetStream();
_stream.ReadTimeout = MillisecondsTimeout;
_stream.WriteTimeout = MillisecondsTimeout;
lock (_ioLock)
{
while (_rxQueue.TryDequeue(out _)) { }
_rxBuffer.Clear();
}
StartReceiveLoop();
ConnectionChanged?.Invoke(this, true);
}
public async Task DisconnectAsync()
{
_manualDisconnect = true;
StopReconnectTimer();
_recvRunning = false;
DisposeStream();
_isBacklightOn = false;
ConnectionChanged?.Invoke(this, false);
await Task.CompletedTask;
}
public void Dispose()
{
_manualDisconnect = true;
StopReconnectTimer();
try { _recvRunning = false; DisposeStream(); } catch { }
}
private void DisposeStream()
{
try { _stream?.Close(); _stream?.Dispose(); } catch { }
_stream = null;
try { _tcp?.Close(); _tcp?.Dispose(); } catch { }
_tcp = null;
}
private void StartReceiveLoop()
{
_recvRunning = true;
_recvThread = new Thread(ReceiveLoop) { IsBackground = true, Name = $"FeederRecv-{Name}" };
_recvThread.Start();
}
private void StartReconnectTimer()
{
if (_manualDisconnect) return;
if (_reconnectTimer != null) return;
_reconnectTimer = new System.Threading.Timer(ReconnectCallback, null,
ReconnectIntervalMs, ReconnectIntervalMs);
}
private void StopReconnectTimer()
{
var t = _reconnectTimer;
_reconnectTimer = null;
try { t?.Dispose(); } catch { }
}
private void ReconnectCallback(object state)
{
if (_manualDisconnect) { StopReconnectTimer(); return; }
if (IsConnected) { StopReconnectTimer(); return; }
try
{
DisposeStream();
var c = new TcpClient();
c.Connect(IP, Port);
_tcp = c;
_stream = c.GetStream();
_stream.ReadTimeout = MillisecondsTimeout;
_stream.WriteTimeout = MillisecondsTimeout;
lock (_ioLock)
{
while (_rxQueue.TryDequeue(out _)) { }
_rxBuffer.Clear();
}
StartReceiveLoop();
StopReconnectTimer();
ConnectionChanged?.Invoke(this, true);
}
catch
{
try { DisposeStream(); } catch { }
}
}
private void ReceiveLoop()
{
var buf = new byte[4096];
try
{
while (_recvRunning && _stream != null && _tcp.Connected)
{
int n = _stream.Read(buf, 0, buf.Length);
if (n <= 0) break;
string chunk = Encoding.ASCII.GetString(buf, 0, n);
lock (_ioLock)
{
_rxBuffer.Append(chunk);
FlushMessages();
}
}
}
catch (IOException) { }
catch (ObjectDisposedException) { }
catch (Exception) { }
finally
{
if (_recvRunning)
{
_recvRunning = false;
ConnectionChanged?.Invoke(this, false);
if (!_manualDisconnect)
StartReconnectTimer();
}
}
}
private void FlushMessages()
{
string content = _rxBuffer.ToString();
int idx;
while ((idx = content.IndexOf(Terminator, StringComparison.Ordinal)) >= 0)
{
var msg = content.Substring(0, idx);
content = content.Substring(idx + Terminator.Length);
_rxQueue.Enqueue(msg);
}
_rxBuffer.Clear();
_rxBuffer.Append(content);
}
#region 背光
public async Task OpenBacklightAsync()
{
await SendAsync("{K1}");
_isBacklightOn = true;
}
public async Task CloseBacklightAsync()
{
await SendAsync("{K0}");
_isBacklightOn = false;
}
public async Task QueryBacklightAsync()
{
var resp = await SendAndRecvAsync("{K?}");
return resp.Contains("1");
}
#endregion
#region 方向执行
public async Task RunDirectionAsync(FeederAction action, int? durationMs = null)
{
char dirChar = (char)('A' + (int)action);
string cmd = durationMs.HasValue
? $"{{C{dirChar}{durationMs.Value}}}"
: $"{{C{dirChar}}}";
var resp = await SendAndRecvAsync(cmd);
var clean = resp.Trim()
.Replace("{", "").Replace("}", "")
.Replace("C" + dirChar, "");
if (string.IsNullOrEmpty(clean)) return 0;
int.TryParse(clean, out var ms);
return ms;
}
public async Task StopAsync()
{
await SendAsync("{HC}");
}
#endregion
#region 参数读写
public async Task GetActionParamAsync(int index)
{
if (index < 0 || index > 10) index = 0;
char ch = (char)('A' + index);
var resp = await SendAndRecvAsync($"{{LC{ch}}}");
var body = resp.Trim();
var start = body.IndexOf('(');
var end = body.IndexOf(')');
if (start < 0 || end < 0) throw new Exception($"Feeder LC{ch} 响应格式错误: {resp}");
var inner = body.Substring(start + 1, end - start - 1);
var parts = inner.Split(';');
if (parts.Length < 17) throw new Exception($"Feeder LC{ch} 参数数量不足: {parts.Length}");
var p = new SingleActionParam();
int i = 0;
p.Motor1 = ParseMotor(parts, ref i);
p.Motor2 = ParseMotor(parts, ref i);
p.Motor3 = ParseMotor(parts, ref i);
p.Motor4 = ParseMotor(parts, ref i);
p.DurationValue = int.Parse(parts[16]);
return p;
}
public async Task SetActionParamAsync(int index, SingleActionParam param)
{
if (index < 0 || index > 10) index = 0;
char ch = (char)('A' + index);
var sb = new StringBuilder();
sb.Append($"{param.Motor1.Amplitude};{param.Motor1.Frequency};{param.Motor1.Phase};{(int)param.Motor1.WaveShape};");
sb.Append($"{param.Motor2.Amplitude};{param.Motor2.Frequency};{param.Motor2.Phase};{(int)param.Motor2.WaveShape};");
sb.Append($"{param.Motor3.Amplitude};{param.Motor3.Frequency};{param.Motor3.Phase};{(int)param.Motor3.WaveShape};");
sb.Append($"{param.Motor4.Amplitude};{param.Motor4.Frequency};{param.Motor4.Phase};{(int)param.Motor4.WaveShape};");
sb.Append(param.DurationValue);
await SendAndRecvAsync($"{{SC{ch}=({sb})}}");
}
public async Task DownloadAllParamsAsync(FeederInfo info)
{
if (info?.ActionGroups == null) return;
for (int i = 0; i < info.ActionGroups.Count && i < 11; i++)
{
await SetActionParamAsync(i, info.ActionGroups[i]);
}
}
public async Task SaveToDeviceAsync()
{
await SendAndRecvAsync("{DV}");
}
#endregion
#region 系统参数(超时/背光)
public async Task GetSystemParamAsync()
{
var p = new FeederSystemParam
{
LightTimeout = (ushort)await ReadParamQueryAsync(AddrLightTimeout),
LightBrightness = (ushort)await ReadParamQueryAsync(AddrLightBrightness),
LightFlashTime = (ushort)await ReadParamQueryAsync(AddrBackLightFlash),
PlatShakeTimeout = (ushort)await ReadParamQueryAsync(AddrPlatShakeTimeout),
HopperShakeTimeout = (ushort)await ReadParamQueryAsync(AddrHopperShakeTimeout),
};
return p;
}
public async Task SetSystemParamAsync(FeederSystemParam param)
{
if (param == null) return;
await WriteParamAsync(AddrLightTimeout, param.LightTimeout);
await WriteParamAsync(AddrLightBrightness, param.LightBrightness);
await WriteParamAsync(AddrBackLightFlash, param.LightFlashTime);
await WriteParamAsync(AddrPlatShakeTimeout, param.PlatShakeTimeout);
await WriteParamAsync(AddrHopperShakeTimeout, param.HopperShakeTimeout);
// 全局参数存至非易失内存
await SendAndRecvAsync("{DG}");
}
#endregion
#region 料斗输入输出
public async Task GetHopperParamAsync(int group)
{
if (group < 0) group = 0;
int start = HopperParamStartAddress + group * 8;
var values = new int[8];
for (int i = 0; i < 8; i++)
values[i] = await ReadParamAsync(start + i);
return new HopperOutputParam
{
HopperDigitalOutput1 = values[0] == 1,
HopperAnalogOutput1 = values[1],
HopperDigitalOutput2 = values[2] == 1,
HopperAnalogOutput2 = values[3],
HopperVibrationAmpl = values[4],
HopperVibrationFreq = values[5],
HopperVibrationWaveform = values[6],
HopperDuration = values[7],
};
}
public async Task SetHopperParamAsync(int group, HopperOutputParam param)
{
if (param == null) return;
if (group < 0) group = 0;
int start = HopperParamStartAddress + group * 8;
await WriteParamAsync(start + 0, param.HopperDigitalOutput1 ? 1 : 0);
await WriteParamAsync(start + 1, param.HopperAnalogOutput1);
await WriteParamAsync(start + 2, param.HopperDigitalOutput2 ? 1 : 0);
await WriteParamAsync(start + 3, param.HopperAnalogOutput2);
await WriteParamAsync(start + 7, param.HopperDuration);
// 振动集参数存至非易失内存
await SendAndRecvAsync("{DV}");
}
public async Task RunHopperOutputAsync(int id, int? timespan = null)
{
if (id < 1 || id > 26) id = 1;
char ch = (char)('A' + (id - 1));
string output = "B" + ch;
string cmd = "{" + output + (timespan.HasValue ? timespan.Value.ToString() : "") + "}";
var resp = await SendAndRecvAsync(cmd);
var clean = resp.Trim().Replace("{", "").Replace("}", "").Replace(output, "");
if (string.IsNullOrWhiteSpace(clean)) return 0;
int.TryParse(clean, out var ms);
return ms;
}
public async Task StopHopperOutputAsync()
{
await SendAsync("{HB}");
}
#endregion
#region 输入触发
public async Task GetInputTriggerAsync(int index)
{
if (index < 0 || index > 2) index = 0;
return await ReadParamAsync(InputTriggerStartAddress + index);
}
public async Task SetInputTriggerAsync(int index, int sequenceId)
{
if (index < 0 || index > 2) index = 0;
await WriteParamAsync(InputTriggerStartAddress + index, sequenceId);
}
#endregion
#region 通用参数读写(RP/WP)
public async Task ReadParamAsync(int id)
{
var resp = await SendAndRecvAsync($"{{RP{id}}}");
return ParseParamValue(resp);
}
/// 查询形式 {RPid?},用于全局系统参数。
private async Task ReadParamQueryAsync(int id)
{
var resp = await SendAndRecvAsync($"{{RP{id}?}}");
return ParseParamValue(resp);
}
public async Task WriteParamAsync(int id, int value)
{
await SendAndRecvAsync($"{{WP{id}={value}}}");
}
/// 从 {RPxx:value} 形式响应中解析出整数值。
private static int ParseParamValue(string resp)
{
if (string.IsNullOrEmpty(resp)) return 0;
var body = resp.Trim().Replace("{", "").Replace("}", "");
var idx = body.IndexOf(':');
var token = idx >= 0 ? body.Substring(idx + 1) : body;
int.TryParse(token.Trim(), out var value);
return value;
}
#endregion
#region 底层通讯
private void RaiseTx(string cmd)
{
try { MessageExchanged?.Invoke(this, cmd, true); } catch { }
}
private void RaiseRx(string resp)
{
try { MessageExchanged?.Invoke(this, resp, false); } catch { }
}
private Task SendAsync(string cmd)
{
if (_stream == null) throw new InvalidOperationException("Feeder 未连接");
var data = Encoding.ASCII.GetBytes(cmd + Terminator);
lock (_ioLock)
{
_stream.Write(data, 0, data.Length);
}
RaiseTx(cmd);
return Task.CompletedTask;
}
private async Task SendAndRecvAsync(string cmd)
{
if (_stream == null) throw new InvalidOperationException("Feeder 未连接");
while (_rxQueue.TryDequeue(out _)) { }
var data = Encoding.ASCII.GetBytes(cmd + Terminator);
await _stream.WriteAsync(data, 0, data.Length);
RaiseTx(cmd);
var deadline = DateTime.UtcNow.AddMilliseconds(MillisecondsTimeout);
while (DateTime.UtcNow < deadline)
{
if (_rxQueue.TryDequeue(out var resp))
{
RaiseRx(resp);
return resp;
}
await Task.Delay(20);
}
throw new TimeoutException($"Feeder 命令超时: {cmd}");
}
private static VoiceMotorParam ParseMotor(string[] parts, ref int i)
{
return new VoiceMotorParam
{
Amplitude = ushort.Parse(parts[i++]),
Frequency = ushort.Parse(parts[i++]),
Phase = ushort.Parse(parts[i++]),
WaveShape = (WaveShape)ushort.Parse(parts[i++]),
};
}
#endregion
}
}