-
Notifications
You must be signed in to change notification settings - Fork 37
/
MqttConsumerEndpoint.cs
153 lines (127 loc) · 5.67 KB
/
MqttConsumerEndpoint.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
// Copyright (c) 2020 Sergio Aquilini
// This code is licensed under MIT license (see LICENSE file for details)
using System;
using System.Collections.Generic;
using System.Diagnostics.CodeAnalysis;
using MQTTnet.Client;
using MQTTnet.Protocol;
using Silverback.Messaging.Configuration.Mqtt;
using Silverback.Messaging.Inbound.ErrorHandling;
using Silverback.Util;
namespace Silverback.Messaging
{
/// <summary>
/// Represents a topic to consume from.
/// </summary>
public sealed class MqttConsumerEndpoint : ConsumerEndpoint, IEquatable<MqttConsumerEndpoint>
{
/// <summary>
/// Initializes a new instance of the <see cref="MqttConsumerEndpoint" /> class.
/// </summary>
/// <param name="topics">
/// The name of the topics or the topic filter strings.
/// </param>
public MqttConsumerEndpoint(params string[] topics)
: base(string.Empty)
{
Topics = topics;
if (topics == null || topics.Length == 0)
return;
Name = topics.Length > 1 ? "[" + string.Join(",", topics) + "]" : topics[0];
}
/// <summary>
/// Gets the name of the topics or the topic filter strings.
/// </summary>
public IReadOnlyCollection<string> Topics { get; }
/// <summary>
/// Gets or sets the MQTT client configuration. This is actually a wrapper around the
/// <see cref="MqttClientOptions" /> from the MQTTnet library.
/// </summary>
public MqttClientConfig Configuration { get; set; } = new();
/// <summary>
/// Gets or sets the quality of service level (at most once, at least once or exactly once).
/// The default is <see cref="MqttQualityOfServiceLevel.AtMostOnce" />.
/// </summary>
public MqttQualityOfServiceLevel QualityOfServiceLevel { get; set; }
/// <summary>
/// Gets or sets the maximum number of incoming message that can be processed concurrently.
/// The default is 1.
/// </summary>
public int MaxDegreeOfParallelism { get; set; } = 1;
/// <summary>
/// Gets or sets the maximum number of messages to be consumed and enqueued waiting to be processed.
/// The default is 1.
/// </summary>
public int BackpressureLimit { get; set; } = 1;
/// <inheritdoc cref="Endpoint.Validate" />
public override void Validate()
{
base.Validate();
if (Configuration == null)
throw new EndpointConfigurationException("Configuration cannot be null.");
if (!Configuration.AreHeadersSupported)
{
if (Serializer.RequireHeaders)
{
throw new EndpointConfigurationException(
"Wrong serializer configuration. Since headers (user properties) are not " +
"supported by MQTT prior to version 5, the serializer must be configured with an " +
"hardcoded message type.");
}
switch (ErrorPolicy)
{
case ErrorPolicyBase errorPolicyBase:
CheckErrorPolicyMaxFailedAttemptsNotSet(errorPolicyBase);
break;
case ErrorPolicyChain errorPolicyChain:
errorPolicyChain.Policies.ForEach(CheckErrorPolicyMaxFailedAttemptsNotSet);
break;
}
}
Configuration.Validate();
if (MaxDegreeOfParallelism < 1)
throw new EndpointConfigurationException("MaxDegreeOfParallelism must be greater or equal to 1.");
if (BackpressureLimit < 1)
throw new EndpointConfigurationException("BackpressureLimit must be greater or equal to 1.");
}
/// <inheritdoc cref="ConsumerEndpoint.GetUniqueConsumerGroupName" />
public override string GetUniqueConsumerGroupName() => Name;
/// <inheritdoc cref="IEquatable{T}.Equals(T)" />
public bool Equals(MqttConsumerEndpoint? other)
{
if (other is null)
return false;
if (ReferenceEquals(this, other))
return true;
return BaseEquals(other) &&
Equals(Configuration, other.Configuration) &&
Equals(QualityOfServiceLevel, other.QualityOfServiceLevel);
}
/// <inheritdoc cref="object.Equals(object)" />
public override bool Equals(object? obj)
{
if (obj is null)
return false;
if (ReferenceEquals(this, obj))
return true;
if (obj.GetType() != GetType())
return false;
return Equals((MqttConsumerEndpoint)obj);
}
/// <inheritdoc cref="object.GetHashCode" />
[SuppressMessage(
"ReSharper",
"NonReadonlyMemberInGetHashCode",
Justification = "Protected set is not abused")]
public override int GetHashCode() => Name.GetHashCode(StringComparison.Ordinal);
private static void CheckErrorPolicyMaxFailedAttemptsNotSet(ErrorPolicyBase errorPolicyBase)
{
if (errorPolicyBase is not RetryErrorPolicy && errorPolicyBase.MaxFailedAttemptsCount > 1)
{
throw new EndpointConfigurationException(
"Cannot set MaxFailedAttempts on the error policies (except for the RetryPolicy) " +
"because headers (user properties) are not supported by MQTT prior to version 5.");
}
}
}
}