Skip to content

Session sends a flow with a stale next-outgoing-id while a transfer is being written (regression from #587) #653

Description

@mgeri

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions