Summary
Since 2.5.0 (#587), Session.WriteDelivery writes the transfer frame outside ThisLock and increments nextOutgoingId only afterwards, back under the lock:
// src/Session.cs (v2.5.4)
int len = this.connection.SendCommand(this.channel, transfer, ...); // L770, outside the lock
...
lock (this.ThisLock)
{
this.nextOutgoingId++; // L778
SendFlow reads nextOutgoingId under ThisLock (L296). A flow sent by another thread in between, e.g. a ReceiverLink credit refresh on the same session, carries next-outgoing-id = N although transfer N is already on the wire.
That breaks AMQP 1.0 Part 2, §2.5.6: "The next-outgoing-id is the transfer-id to assign to the next transfer frame". The spec also says that on receiving a flow the peer "MUST update the next-incoming-id directly from the next-outgoing-id of the frame", so the stale value moves the peer's next-incoming-id back by one.
In 2.4.11 the increment and the write happened under the same lock. The bug is still present in 2.5.4 and on master.
Reproduce
Use Apache Qpid Broker-J 10.1.0 or later. Since apache/qpid-broker-j@8b98449, Broker-J validates the next-outgoing-id of incoming flows and ends the session when it goes backwards. Earlier versions did not validate it and silently took the stale value, so the problem goes unnoticed there.
One session with a sender on a dedicated thread and a low-credit receiver, both on the same queue:
var connection = await Connection.Factory.CreateAsync(new Address("amqp://guest:guest@localhost:5672"));
var session = new Session(connection);
var receiver = new ReceiverLink(session, "receiver", "q");
receiver.Start(5, (link, message) => link.Accept(message)); // small credit -> frequent flows
var sender = new SenderLink(session, "sender", "q");
var slots = new SemaphoreSlim(100);
new Thread(() =>
{
for (int i = 0; ; i++)
{
slots.Wait();
sender.Send(new Message("payload " + i), (l, m, o, s) => slots.Release(), null);
}
}).Start();
With 2.5.1 against Broker-J 10.1.0 the session ends within the first few hundred messages, on every run:
SEND (ch=0) begin(next-outgoing-id:4294967293,...)
... 285 transfers, ids 4294967293..281
SEND (ch=0) transfer(handle:1,delivery-id:284,...)
SEND (ch=0) flow(next-in-id:14,in-window:2048,next-out-id:281,...,handle:0,link-credit:5) <- should be 282
...
RECV (ch=0) end(error:error(condition:amqp:session:window-violation,
description:Next outgoing id '281' is less than next incoming id '282'))
| Client |
Session layout |
Broker-J 10.1.0 |
| 2.5.1 |
sender and receiver on one session |
window-violation in < 1 s, 10/10 runs |
| 2.5.1 |
separate sessions |
no error, 3 × 20 s (~870k messages each) |
| 2.4.11 |
sender and receiver on one session |
no error, 3 × 5 s (~180k messages each) |
Suggested fix
This is a possible approach we came up with. You know the threading model of the library much better, so there may well be a more appropriate fix.
The value has to be exact, not just monotonic. The peer copies it into its next-incoming-id, so a value that is ahead is wrong too. Incrementing before the write would let a flow advertise N+1 before transfer N is on the wire; the peer would then count transfer N on top of it and the next flow would look stale again.
To keep the benefit of #587 without adding a new lock, SendFlow defers the flow while a delivery write is in progress (writingDelivery). WriteDelivery sends the pending flows right after incrementing nextOutgoingId, inside the lock block it already enters after every frame. The previous body of SendFlow moves unchanged to WriteFlow.
diff --git a/src/Session.cs b/src/Session.cs
index 88ab00b..1a102c0 100644
--- a/src/Session.cs
+++ b/src/Session.cs
@@ -95,6 +95,7 @@ namespace Amqp
SequenceNumber nextOutgoingId;
uint outgoingWindow;
bool writingDelivery;
+ System.Collections.Generic.List<Flow> pendingFlows;
/// <summary>
/// Initializes a session object.
@@ -291,6 +292,27 @@ namespace Amqp
internal void SendFlow(Flow flow)
{
lock (this.ThisLock)
+ {
+ if (this.writingDelivery)
+ {
+ // a transfer may be on the wire before nextOutgoingId is incremented;
+ // WriteDelivery sends this flow right after the increment
+ if (this.pendingFlows == null)
+ {
+ this.pendingFlows = new System.Collections.Generic.List<Flow>();
+ }
+
+ this.pendingFlows.Add(flow);
+ return;
+ }
+
+ this.WriteFlow(flow);
+ }
+ }
+
+ void WriteFlow(Flow flow)
+ {
+ // Must be called under lock
{
this.incomingWindow = defaultWindowSize;
flow.NextOutgoingId = this.nextOutgoingId;
@@ -781,6 +803,16 @@ namespace Amqp
this.outgoingWindow--;
}
+ if (this.pendingFlows != null && this.pendingFlows.Count > 0)
+ {
+ foreach (Flow f in this.pendingFlows)
+ {
+ this.WriteFlow(f);
+ }
+
+ this.pendingFlows.Clear();
+ }
+
if (delivery.Buffer.Length == 0 || delivery.Link.IsClosed)
{
delivery.InProgress = false;
- No new lock and no change in lock order. A flow is still written under
ThisLock, as SendFlow does today.
- All pending flows carry the same
next-outgoing-id: only transfer frames consume transfer-ids.
- A flow can be delayed by at most one transfer frame.
Summary
Since 2.5.0 (#587),
Session.WriteDeliverywrites thetransferframe outsideThisLockand incrementsnextOutgoingIdonly afterwards, back under the lock:SendFlowreadsnextOutgoingIdunderThisLock(L296). A flow sent by another thread in between, e.g. aReceiverLinkcredit refresh on the same session, carriesnext-outgoing-id = Nalthough transfer N is already on the wire.That breaks AMQP 1.0 Part 2, §2.5.6: "The next-outgoing-id is the transfer-id to assign to the next transfer frame". The spec also says that on receiving a flow the peer "MUST update the next-incoming-id directly from the next-outgoing-id of the frame", so the stale value moves the peer's
next-incoming-idback by one.In 2.4.11 the increment and the write happened under the same lock. The bug is still present in 2.5.4 and on
master.Reproduce
Use Apache Qpid Broker-J 10.1.0 or later. Since apache/qpid-broker-j@8b98449, Broker-J validates the
next-outgoing-idof incoming flows and ends the session when it goes backwards. Earlier versions did not validate it and silently took the stale value, so the problem goes unnoticed there.One session with a sender on a dedicated thread and a low-credit receiver, both on the same queue:
With 2.5.1 against Broker-J 10.1.0 the session ends within the first few hundred messages, on every run:
window-violationin < 1 s, 10/10 runsSuggested fix
This is a possible approach we came up with. You know the threading model of the library much better, so there may well be a more appropriate fix.
The value has to be exact, not just monotonic. The peer copies it into its
next-incoming-id, so a value that is ahead is wrong too. Incrementing before the write would let a flow advertise N+1 before transfer N is on the wire; the peer would then count transfer N on top of it and the next flow would look stale again.To keep the benefit of #587 without adding a new lock,
SendFlowdefers the flow while a delivery write is in progress (writingDelivery).WriteDeliverysends the pending flows right after incrementingnextOutgoingId, inside the lock block it already enters after every frame. The previous body ofSendFlowmoves unchanged toWriteFlow.ThisLock, asSendFlowdoes today.next-outgoing-id: only transfer frames consume transfer-ids.