소스 검색

Replace a manual lock with the lock class.

release/3.x.x
Christian Kratky 6 년 전
부모
커밋
4838da48ef
1개의 변경된 파일6개의 추가작업 그리고 20개의 파일을 삭제
  1. +6
    -20
      Frameworks/MQTTnet.NetStandard/ManagedClient/ManagedMqttClient.cs

+ 6
- 20
Frameworks/MQTTnet.NetStandard/ManagedClient/ManagedMqttClient.cs 파일 보기

@@ -7,6 +7,7 @@ using System.Threading.Tasks;
using MQTTnet.Client;
using MQTTnet.Diagnostics;
using MQTTnet.Exceptions;
using MQTTnet.Internal;
using MQTTnet.Protocol;

namespace MQTTnet.ManagedClient
@@ -15,7 +16,7 @@ namespace MQTTnet.ManagedClient
{
private readonly BlockingCollection<MqttApplicationMessage> _messageQueue = new BlockingCollection<MqttApplicationMessage>();
private readonly Dictionary<string, MqttQualityOfServiceLevel> _subscriptions = new Dictionary<string, MqttQualityOfServiceLevel>();
private readonly SemaphoreSlim _subscriptionsSemaphore = new SemaphoreSlim(1, 1);
private readonly AsyncLock _subscriptionsLock = new AsyncLock();
private readonly List<string> _unsubscriptions = new List<string>();

private readonly IMqttClient _mqttClient;
@@ -117,8 +118,7 @@ namespace MQTTnet.ManagedClient
{
if (topicFilters == null) throw new ArgumentNullException(nameof(topicFilters));

await _subscriptionsSemaphore.WaitAsync().ConfigureAwait(false);
try
using (await _subscriptionsLock.LockAsync(CancellationToken.None).ConfigureAwait(false))
{
foreach (var topicFilter in topicFilters)
{
@@ -126,16 +126,11 @@ namespace MQTTnet.ManagedClient
_subscriptionsNotPushed = true;
}
}
finally
{
_subscriptionsSemaphore.Release();
}
}

public async Task UnsubscribeAsync(IEnumerable<string> topics)
{
await _subscriptionsSemaphore.WaitAsync().ConfigureAwait(false);
try
using (await _subscriptionsLock.LockAsync(CancellationToken.None).ConfigureAwait(false))
{
foreach (var topic in topics)
{
@@ -146,16 +141,12 @@ namespace MQTTnet.ManagedClient
}
}
}
finally
{
_subscriptionsSemaphore.Release();
}
}

public void Dispose()
{
_messageQueue?.Dispose();
_subscriptionsSemaphore?.Dispose();
_subscriptionsLock?.Dispose();
_connectionCancellationToken?.Dispose();
_publishingCancellationToken?.Dispose();
}
@@ -294,8 +285,7 @@ namespace MQTTnet.ManagedClient
List<TopicFilter> subscriptions;
List<string> unsubscriptions;

await _subscriptionsSemaphore.WaitAsync().ConfigureAwait(false);
try
using (await _subscriptionsLock.LockAsync(CancellationToken.None).ConfigureAwait(false))
{
subscriptions = _subscriptions.Select(i => new TopicFilter(i.Key, i.Value)).ToList();

@@ -304,10 +294,6 @@ namespace MQTTnet.ManagedClient

_subscriptionsNotPushed = false;
}
finally
{
_subscriptionsSemaphore.Release();
}
if (!subscriptions.Any() && !unsubscriptions.Any())
{


불러오는 중...
취소
저장