Вы не можете выбрать более 25 тем Темы должны начинаться с буквы или цифры, могут содержать дефисы(-) и должны содержать не более 35 символов.
 
 
 
 

65 строки
1.8 KiB

  1. using System;
  2. using System.Collections.Concurrent;
  3. using System.Collections.Generic;
  4. using System.Threading;
  5. using System.Threading.Tasks;
  6. using MQTTnet.Core.Adapter;
  7. using MQTTnet.Core.Packets;
  8. using MQTTnet.Core.Serializer;
  9. namespace MQTTnet.Core.Tests
  10. {
  11. public class TestMqttCommunicationAdapter : IMqttCommunicationAdapter
  12. {
  13. private readonly BlockingCollection<MqttBasePacket> _incomingPackets = new BlockingCollection<MqttBasePacket>();
  14. public TestMqttCommunicationAdapter Partner { get; set; }
  15. public IMqttPacketSerializer PacketSerializer { get; } = new MqttPacketSerializer();
  16. public Task ConnectAsync(TimeSpan timeout)
  17. {
  18. return Task.FromResult(0);
  19. }
  20. public Task DisconnectAsync(TimeSpan timeout)
  21. {
  22. return Task.FromResult(0);
  23. }
  24. public Task SendPacketsAsync(TimeSpan timeout, CancellationToken cancellationToken, IEnumerable<MqttBasePacket> packets)
  25. {
  26. ThrowIfPartnerIsNull();
  27. foreach (var packet in packets)
  28. {
  29. Partner.SendPacketInternal(packet);
  30. }
  31. return Task.FromResult(0);
  32. }
  33. public Task<MqttBasePacket> ReceivePacketAsync(TimeSpan timeout, CancellationToken cancellationToken)
  34. {
  35. ThrowIfPartnerIsNull();
  36. return Task.Run(() => _incomingPackets.Take(), cancellationToken);
  37. }
  38. private void SendPacketInternal(MqttBasePacket packet)
  39. {
  40. if (packet == null) throw new ArgumentNullException(nameof(packet));
  41. _incomingPackets.Add(packet);
  42. }
  43. private void ThrowIfPartnerIsNull()
  44. {
  45. if (Partner == null)
  46. {
  47. throw new InvalidOperationException("Partner is not set.");
  48. }
  49. }
  50. }
  51. }