99using System . Formats . Cbor ;
1010using System . Threading ;
1111using System . Threading . Tasks ;
12+ using Microsoft . Extensions . Logging ;
1213using ThingSet . Common . Protocols . Binary ;
1314
1415namespace ThingSet . Common . Transports ;
@@ -22,24 +23,31 @@ public abstract class ClientTransportBase<TEndpoint> : IClientTransport
2223{
2324 protected readonly ConcurrentDictionary < TEndpoint , ReceiveBuffer > _buffersBySender = new ConcurrentDictionary < TEndpoint , ReceiveBuffer > ( ) ;
2425
26+ private readonly ILogger ? _logger ;
27+ private readonly int _bufferSize ;
28+
2529 private Action < ulong ? , CborReader > ? _callback ;
2630
2731 private readonly Thread _subscriptionThread ;
2832 private bool _runSubscriptionThread = true ;
2933
30- protected ClientTransportBase ( )
34+ protected ClientTransportBase ( int bufferSize = 16384 , ILogger ? logger = null )
3135 {
36+ _logger = logger ;
37+ _bufferSize = bufferSize ;
38+
3239 _subscriptionThread = new Thread ( RunSubscriptionThread )
3340 {
3441 IsBackground = true ,
35- Name = $ "Subscription { Address } ",
42+ Name = $ "Subscription { PeerAddress } ",
3643 } ;
3744 }
3845
3946 /// <summary>
40- /// String representation of a network identifier.
47+ /// String representation of a network identifier for the peer device
48+ /// to which this client is connected.
4149 /// </summary>
42- protected abstract string Address { get ; }
50+ public abstract string PeerAddress { get ; }
4351
4452 /// <summary>
4553 /// Connects this transport.
@@ -90,6 +98,12 @@ protected async void RunSubscriptionThread()
9098
9199 protected abstract ValueTask HandleIncomingPublicationsAsync ( ) ;
92100
101+ protected ReceiveBuffer GetOrCreateBuffer ( TEndpoint endpoint )
102+ {
103+ byte [ ] buffer = new byte [ _bufferSize ] ;
104+ return _buffersBySender . GetOrAdd ( endpoint , _ => new ReceiveBuffer ( buffer ) ) ;
105+ }
106+
93107 protected void NotifyReport ( ulong ? eui , CborReader reader )
94108 {
95109 reader . ReadUInt32 ( ) ; // subset ID
@@ -104,6 +118,13 @@ protected void NotifyReport(ulong? eui, CborReader reader)
104118 protected abstract class ReportParser < TMessageType >
105119 where TMessageType : Enum
106120 {
121+ private readonly ILogger ? _logger ;
122+
123+ protected ReportParser ( ILogger ? logger )
124+ {
125+ _logger = logger ;
126+ }
127+
107128 /// <returns>True if a complete message has been assembled.</returns>
108129 public bool TryParse ( byte sequenceNumber , byte messageNumber , TMessageType messageType ,
109130 ReceiveBuffer buffer , byte [ ] data , [ MaybeNullWhen ( true ) ] out ulong ? eui ,
@@ -116,14 +137,17 @@ public bool TryParse(byte sequenceNumber, byte messageNumber, TMessageType messa
116137 {
117138 buffer . Started = true ;
118139 buffer . MessageNumber = messageNumber ;
140+ _logger ? . LogDebug ( $ "Message { messageNumber } started") ;
119141 }
120142 else if ( buffer . MessageNumber != messageNumber )
121143 {
144+ _logger ? . LogDebug ( $ "Message { messageNumber } mismatch; expected { buffer . MessageNumber } ") ;
122145 buffer . Reset ( ) ;
123146 return false ;
124147 }
125148 else if ( ! buffer . Started )
126149 {
150+ _logger ? . LogDebug ( $ "Message unexpected") ;
127151 buffer . Reset ( ) ;
128152 return false ;
129153 }
@@ -135,6 +159,7 @@ public bool TryParse(byte sequenceNumber, byte messageNumber, TMessageType messa
135159 }
136160 if ( IsLast ( messageType ) )
137161 {
162+ _logger ? . LogDebug ( $ "Dispatching message { messageNumber } of length { buffer . Position } ") ;
138163 ReadOnlyMemory < byte > memory = buffer . Buffer ;
139164 reader = new CborReader ( memory . Slice ( 1 ) , CborConformanceMode . Lax , allowMultipleRootLevelValues : true ) ;
140165 if ( buffer . Buffer [ 0 ] == ( byte ) ThingSetRequest . ReportEnhanced )
@@ -156,7 +181,12 @@ public bool TryParse(byte sequenceNumber, byte messageNumber, TMessageType messa
156181
157182 protected class ReceiveBuffer
158183 {
159- public byte [ ] Buffer = new byte [ 32768 ] ;
184+ public ReceiveBuffer ( byte [ ] buffer )
185+ {
186+ Buffer = buffer ;
187+ }
188+
189+ public byte [ ] Buffer ;
160190 public int Position ;
161191 public byte Sequence ;
162192 public byte MessageNumber ;
0 commit comments