Dreamine.Communication.Core
1.0.2
Dreamine.Communication.Core 통신 기능과 관련 API를 제공합니다.
Toggle main menu visibility
로딩중...
검색중...
일치하는것 없음
InMemoryOutboundMessageQueue.cs
이 파일의 문서화 페이지로 가기
1
using
Dreamine.Communication.Abstractions.Options;
2
3
namespace
Dreamine.Communication.Core.Queues
;
4
13
public
sealed
class
InMemoryOutboundMessageQueue
:
IOutboundMessageQueue
14
{
23
private
readonly
object
_sync
=
new
();
32
private
readonly LinkedList<QueuedOutboundMessage>
_messages
=
new
();
41
private
readonly OutboundQueueOptions
_options
;
42
67
public
InMemoryOutboundMessageQueue
(OutboundQueueOptions? options =
null
)
68
{
69
_options
= options ??
new
OutboundQueueOptions();
70
71
if
(
_options
.MaxQueueSize <= 0)
72
{
73
throw
new
ArgumentOutOfRangeException(
74
nameof(options),
75
"MaxQueueSize must be greater than zero."
);
76
}
77
}
78
87
public
int
Count
88
{
89
get
90
{
91
lock (
_sync
)
92
{
93
return
_messages
.Count;
94
}
95
}
96
}
97
154
public
ValueTask
EnqueueAsync
(
155
QueuedOutboundMessage
message,
156
CancellationToken cancellationToken =
default
)
157
{
158
ArgumentNullException.ThrowIfNull(message);
159
cancellationToken.ThrowIfCancellationRequested();
160
161
lock (
_sync
)
162
{
163
EnsureCapacityForAddToBack
();
164
_messages
.AddLast(message);
165
}
166
167
return
ValueTask.CompletedTask;
168
}
169
226
public
ValueTask
EnqueueFrontAsync
(
227
QueuedOutboundMessage
message,
228
CancellationToken cancellationToken =
default
)
229
{
230
ArgumentNullException.ThrowIfNull(message);
231
cancellationToken.ThrowIfCancellationRequested();
232
233
lock (
_sync
)
234
{
235
EnsureCapacityForAddToFront
();
236
_messages
.AddFirst(message);
237
}
238
239
return
ValueTask.CompletedTask;
240
}
241
274
public
ValueTask<QueuedOutboundMessage?>
TryDequeueAsync
(
275
CancellationToken cancellationToken =
default
)
276
{
277
cancellationToken.ThrowIfCancellationRequested();
278
279
lock (
_sync
)
280
{
281
if
(
_messages
.First is
null
)
282
{
283
return
ValueTask.FromResult<
QueuedOutboundMessage
?>(
null
);
284
}
285
286
var message =
_messages
.First.Value;
287
_messages
.RemoveFirst();
288
289
return
ValueTask.FromResult<
QueuedOutboundMessage
?>(message);
290
}
291
}
292
301
public
void
Clear
()
302
{
303
lock (
_sync
)
304
{
305
_messages
.Clear();
306
}
307
}
308
325
private
void
EnsureCapacityForAddToBack
()
326
{
327
if
(
_messages
.Count <
_options
.MaxQueueSize)
328
{
329
return
;
330
}
331
332
if
(!
_options
.DropOldestWhenFull)
333
{
334
throw
new
InvalidOperationException(
"Outbound message queue is full."
);
335
}
336
337
_messages
.RemoveFirst();
338
}
339
356
private
void
EnsureCapacityForAddToFront
()
357
{
358
if
(
_messages
.Count <
_options
.MaxQueueSize)
359
{
360
return
;
361
}
362
363
if
(!
_options
.DropOldestWhenFull)
364
{
365
throw
new
InvalidOperationException(
"Outbound message queue is full."
);
366
}
367
368
_messages
.RemoveLast();
369
}
370
}
Dreamine.Communication.Core.Queues
Definition
InMemoryOutboundMessageQueue.cs:3
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue._sync
readonly object _sync
Definition
InMemoryOutboundMessageQueue.cs:23
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.EnqueueFrontAsync
ValueTask EnqueueFrontAsync(QueuedOutboundMessage message, CancellationToken cancellationToken=default)
Definition
InMemoryOutboundMessageQueue.cs:226
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.Clear
void Clear()
Definition
InMemoryOutboundMessageQueue.cs:301
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue._options
readonly OutboundQueueOptions _options
Definition
InMemoryOutboundMessageQueue.cs:41
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.Count
int Count
Definition
InMemoryOutboundMessageQueue.cs:88
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue._messages
readonly LinkedList< QueuedOutboundMessage > _messages
Definition
InMemoryOutboundMessageQueue.cs:32
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.TryDequeueAsync
ValueTask< QueuedOutboundMessage?> TryDequeueAsync(CancellationToken cancellationToken=default)
Definition
InMemoryOutboundMessageQueue.cs:274
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.EnqueueAsync
ValueTask EnqueueAsync(QueuedOutboundMessage message, CancellationToken cancellationToken=default)
Definition
InMemoryOutboundMessageQueue.cs:154
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.InMemoryOutboundMessageQueue
InMemoryOutboundMessageQueue(OutboundQueueOptions? options=null)
Definition
InMemoryOutboundMessageQueue.cs:67
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.EnsureCapacityForAddToFront
void EnsureCapacityForAddToFront()
Definition
InMemoryOutboundMessageQueue.cs:356
Dreamine.Communication.Core.Queues.InMemoryOutboundMessageQueue.EnsureCapacityForAddToBack
void EnsureCapacityForAddToBack()
Definition
InMemoryOutboundMessageQueue.cs:325
Dreamine.Communication.Core.Queues.IOutboundMessageQueue
Definition
IOutboundMessageQueue.cs:12
Dreamine.Communication.Core.Queues.QueuedOutboundMessage
Definition
QueuedOutboundMessage.cs:14
Queues
InMemoryOutboundMessageQueue.cs
다음에 의해 생성됨 :
1.17.0