Dreamine.Communication.Core 1.0.2
Dreamine.Communication.Core 통신 기능과 관련 API를 제공합니다.
로딩중...
검색중...
일치하는것 없음
InMemoryOutboundMessageQueue.cs
이 파일의 문서화 페이지로 가기
1using Dreamine.Communication.Abstractions.Options;
2
4
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 {
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 {
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
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
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}
ValueTask EnqueueFrontAsync(QueuedOutboundMessage message, CancellationToken cancellationToken=default)
ValueTask< QueuedOutboundMessage?> TryDequeueAsync(CancellationToken cancellationToken=default)
ValueTask EnqueueAsync(QueuedOutboundMessage message, CancellationToken cancellationToken=default)