-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathQueueListener.cs
More file actions
80 lines (74 loc) · 1.82 KB
/
Copy pathQueueListener.cs
File metadata and controls
80 lines (74 loc) · 1.82 KB
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
using System;
using System.Threading;
namespace SimpleSqlQueue
{
public class QueueListener<T>
{
bool _isRunning;
readonly IQueue<T> _queue;
readonly int _retriesBeforeFailure;
readonly int _visibilityTimeout;
public QueueListener(IQueue<T> queue, int retriesBeforeFailure, int visibilityTimeout)
{
_queue = queue;
_retriesBeforeFailure = retriesBeforeFailure;
_visibilityTimeout = visibilityTimeout;
}
public void StartListening(Action<QueueItem<T>> processingAction)
{
_isRunning = true;
ThreadPool.QueueUserWorkItem(delegate { MonitorQueue(processingAction); });
}
public void StopListening()
{
_isRunning = false;
}
void LogFailure(string source, string message)
{
}
void MonitorQueue(Action<QueueItem<T>> action)
{
while (_isRunning)
{
Thread.Sleep(100);
var item = _queue.Dequeue(_visibilityTimeout);
if (item != null)
{
try
{
action.Invoke(item);
_queue.DeleteQueueItem(item);
}
catch (Exception e)
{
if (item.FailedAttempts >= _retriesBeforeFailure)
{
LogFailure(
"Queue Processor",
string.Format(
"Failed to process Queue Item {0}, failure number {1}, item has been moved to failed item queue.{2}Exception Details:{2}{3}",
item.Id,
item.FailedAttempts,
Environment.NewLine,
e));
_queue.StoreFailedMessage(item);
_queue.DeleteQueueItem(item);
}
else
{
LogFailure(
"Queue Processor",
string.Format(
"Failed to process Queue Item {0}, failure number {1}, will retry {2} times. {3}Exception Details:{3}{4}",
item.Id,
item.FailedAttempts,
_retriesBeforeFailure - item.FailedAttempts,
Environment.NewLine,
e));
}
}
}
}
}
}
}