-
Notifications
You must be signed in to change notification settings - Fork 40
/
AmqpSourceStage.cs
170 lines (163 loc) · 7.59 KB
/
AmqpSourceStage.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
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
using System;
using System.Collections.Generic;
using System.Linq;
using Akka.IO;
using Akka.Streams.Stage;
using RabbitMQ.Client;
namespace Akka.Streams.Amqp
{
/// <summary>
/// Connects to an AMQP server upon materialization and consumes messages from it emitting them
/// into the stream. Each materialized source will create one connection to the broker.
/// As soon as an <see cref="IncomingMessage"/> is sent downstream, an ack for it is sent to the broker.
/// </summary>
public sealed class AmqpSourceStage : GraphStage<SourceShape<IncomingMessage>>
{
public static readonly Attributes DefaultAttributes = Attributes.CreateName("AmqpSource");
public IAmqpSourceSettings Settings { get; }
public int BufferSize { get; }
/// <summary>
/// Constructor
/// </summary>
/// <param name="settings">The source settings</param>
/// <param name="bufferSize">The max number of elements to prefetch and buffer at any given time.</param>
public AmqpSourceStage(IAmqpSourceSettings settings, int bufferSize)
{
Settings = settings;
BufferSize = bufferSize;
}
public override SourceShape<IncomingMessage> Shape => new SourceShape<IncomingMessage>(Out);
public readonly Outlet<IncomingMessage> Out = new Outlet<IncomingMessage>("AmqpSource.out");
protected override Attributes InitialAttributes => DefaultAttributes;
protected override GraphStageLogic CreateLogic(Attributes inheritedAttributes)
{
return new AmqpSourceStageLogic(this);
}
public override string ToString()
{
return "AmqpSource";
}
private class AmqpSourceStageLogic : AmqpConnectorLogic
{
private readonly AmqpSourceStage _stage;
private readonly Queue<IncomingMessage> _queue = new Queue<IncomingMessage>();
private IBasicConsumer _amqpSourceConsumer;
public AmqpSourceStageLogic(AmqpSourceStage stage) : base(stage.Shape)
{
_stage = stage;
SetHandler(_stage.Out, () =>
{
if (_queue.Count > 0)
{
PushAndAckMessage(_queue.Dequeue());
}
});
}
public override IAmqpConnectorSettings Settings => _stage.Settings;
public override IConnectionFactory ConnectionFactoryFrom(IAmqpConnectionSettings settings) =>
AmqpConnector.ConnectionFactoryFrom(settings);
public override IConnection NewConnection(IConnectionFactory factory, IAmqpConnectionSettings settings) =>
AmqpConnector.NewConnection(factory, settings);
public override void WhenConnected()
{
// we have only one consumer per connection so global is ok
Channel.BasicQos(0, (ushort) _stage.BufferSize, false);
var consumerCallback = GetAsyncCallback<IncomingMessage>(HandleDelivery);
var shutdownCallback = GetAsyncCallback<ShutdownEventArgs>(args =>
{
if (args != null)
{
FailStage(ShutdownSignalException.FromArgs(args));
}
else
{
CompleteStage();
}
});
_amqpSourceConsumer = new DefaultConsumer(consumerCallback,shutdownCallback);
if (Settings is NamedQueueSourceSettings)
{
var namedSourceSettings = (NamedQueueSourceSettings) Settings;
SetupNamedQueue(namedSourceSettings);
}
else if (Settings is TemporaryQueueSourceSettings)
{
var tempSettings = (TemporaryQueueSourceSettings) Settings;
SetupTmeporaryQueue(tempSettings);
}
}
public override void OnFailure(Exception ex)
{
}
private void SetupNamedQueue(NamedQueueSourceSettings settings)
{
Channel.BasicConsume(settings.Queue,
false,// never auto-ack
settings.ConsumerTag,// consumer tag
settings.NoLocal,
settings.Exclusive,
settings.Arguments.ToDictionary(k=> k.Key, val=> val.Value), _amqpSourceConsumer);
}
private void SetupTmeporaryQueue(TemporaryQueueSourceSettings settings)
{
// this is a weird case that required dynamic declaration, the queue name is not known
// up front, it is only useful for sources, so that's why it's not placed in the AmqpConnectorLogic
var queueName = Channel.QueueDeclare().QueueName;
Channel.QueueBind(queueName, settings.Exchange, settings.RoutingKey ?? "");
Channel.BasicConsume(queueName, false, _amqpSourceConsumer);
}
private void HandleDelivery(IncomingMessage message)
{
if (IsAvailable(_stage.Out))
{
PushAndAckMessage(message);
}
else
{
if (_queue.Count + 1 > _stage.BufferSize)
{
FailStage(new ApplicationException($"Reached maximum buffer size {_stage.BufferSize}"));
}
else
{
_queue.Enqueue(message);
}
}
}
private void PushAndAckMessage(IncomingMessage message)
{
Push(_stage.Out, message);
// ack it as soon as we have passed it downstream
// TODO ack less often and do batch acks with multiple = true would probably be more performant
Channel.BasicAck(message.Envelope.DeliveryTag,
false);// just this single message
}
private class DefaultConsumer : DefaultBasicConsumer
{
private readonly Action<IncomingMessage> _consumerCallback;
private readonly Action<ShutdownEventArgs> _shutdownCallback;
public DefaultConsumer(Action<IncomingMessage> consumerCallback, Action<ShutdownEventArgs> shutdownCallback)
{
_consumerCallback = consumerCallback;
_shutdownCallback = shutdownCallback;
}
public override void HandleBasicDeliver(string consumerTag, ulong deliveryTag, bool redelivered, string exchange, string routingKey,
IBasicProperties properties, byte[] body)
{
var envelope = Envelope.Create(deliveryTag, redelivered, exchange, routingKey);
var incomingMessage = IncomingMessage.Create(ByteString.CopyFrom(body), envelope, properties);
_consumerCallback?.Invoke(incomingMessage);
}
public override void HandleBasicCancel(string consumerTag)
{
// non consumer initiated cancel, for example happens when the queue has been deleted.
_shutdownCallback?.Invoke(null);
}
public override void HandleModelShutdown(object model, ShutdownEventArgs reason)
{
_shutdownCallback?.Invoke(reason);
}
}
}
}
}