From 5c9810fb301dbdafb2d09b3650c791942b42f6a9 Mon Sep 17 00:00:00 2001 From: Jeff Repanich Date: Sun, 13 Sep 2026 16:48:02 -0400 Subject: [PATCH] Support latest server correlation contract --- compose.yml | 2 +- fitz/client.go | 7 + internal/core/client/client.go | 17 ++ internal/core/connection/connection.go | 153 +++++++++++++++--- internal/core/connection/mux.go | 102 +++++++++++- internal/core/connection/mux_internal_test.go | 21 +++ internal/protocol/frame.go | 60 +++++++ internal/protocol/frame_test.go | 17 ++ test/queue_test.go | 46 ++++++ 9 files changed, 394 insertions(+), 31 deletions(-) diff --git a/compose.yml b/compose.yml index 478193b..0252e37 100644 --- a/compose.yml +++ b/compose.yml @@ -16,7 +16,7 @@ # - Anon TCP: tcp://127.0.0.1:${FITZ_ANON_HOST_TCP_PORT:-4191} x-broker-common: &broker-common - image: ghcr.io/cntryl/fitz:latest + image: ghcr.io/cntryl/fitz@sha256:976b4baa57b91841021e43563112cf0b686a3e79af89c1242a4f5cbf362e7afa restart: unless-stopped stop_grace_period: 15s healthcheck: diff --git a/fitz/client.go b/fitz/client.go index 5a14e3d..d4680dd 100644 --- a/fitz/client.go +++ b/fitz/client.go @@ -66,6 +66,13 @@ func (c *Client) State() ConnectionState { return fromCoreConnectionState(c.inner.State()) } +// CorrelationEnabled reports whether the current broker session supports +// frame-level request correlation. +func (c *Client) CorrelationEnabled() bool { return c.inner.CorrelationEnabled() } + +// ServerCapabilities returns the advertised protocol version and capability bits. +func (c *Client) ServerCapabilities() (uint16, uint32) { return c.inner.ServerCapabilities() } + // Notice returns the Notice domain client for publish/subscribe messaging. func (c *Client) Notice() NoticeClient { return ¬iceClient{inner: c.inner.Notice()} diff --git a/internal/core/client/client.go b/internal/core/client/client.go index fdbb9a8..608959c 100644 --- a/internal/core/client/client.go +++ b/internal/core/client/client.go @@ -594,6 +594,23 @@ func (c *Client) Metrics() connection.MultiplexerMetrics { return connection.MultiplexerMetrics{} } +// CorrelationEnabled reports whether the current broker session advertised +// frame-level request correlation. +func (c *Client) CorrelationEnabled() bool { + if conn := c.currentConnection(); conn != nil { + return conn.CorrelationEnabled() + } + return false +} + +// ServerCapabilities returns the current protocol version and capability bits. +func (c *Client) ServerCapabilities() (uint16, uint32) { + if conn := c.currentConnection(); conn != nil { + return conn.ServerCapabilities() + } + return 0, 0 +} + // Domain client accessors. // KV returns the KV domain client. diff --git a/internal/core/connection/connection.go b/internal/core/connection/connection.go index ce77a95..761e7ef 100644 --- a/internal/core/connection/connection.go +++ b/internal/core/connection/connection.go @@ -3,6 +3,7 @@ package connection import ( "bytes" "context" + "encoding/binary" "errors" "fmt" "log/slog" @@ -146,6 +147,10 @@ type Config struct { Meter metric.Meter // When nil, otel.Meter(module) is used. } +func (c *Connection) CorrelationEnabled() bool { return c.mux.CorrelationEnabled() } + +func (c *Connection) ServerCapabilities() (uint16, uint32) { return c.mux.Capabilities() } + // DefaultConfig returns default configuration. func DefaultConfig() Config { return Config{ @@ -582,8 +587,7 @@ func (c *Connection) dispatchLoop() { } c.recordActivity() - // Decode frame (MessageType + payload) - msgType, payload, err := protocol.DecodeFrame(frame) + hasResponse, err := c.dispatchTransportFrame(frame) if err != nil { if c.logger != nil { c.logger.Error("decode frame failed", "error", err) @@ -591,18 +595,13 @@ func (c *Connection) dispatchLoop() { c.setConnError(fmt.Errorf("decode frame: %w", err)) return } - if c.logger != nil { - c.logger.Debug("frame received", "msg_type", msgType) - } // First valid response confirms authentication - if firstResponse { + if firstResponse && hasResponse { c.confirmAuthentication() firstResponse = false } - // Route to multiplexer (non-blocking dispatch) - c.mux.Dispatch(msgType, payload) continue } @@ -613,8 +612,7 @@ func (c *Connection) dispatchLoop() { } c.recordActivity() - // Decode frame (MessageType + payload) - msgType, payload, err := protocol.DecodeFrame(frame) + hasResponse, err := c.dispatchTransportFrame(frame) if err != nil { if c.logger != nil { c.logger.Error("decode frame failed", "error", err) @@ -622,19 +620,55 @@ func (c *Connection) dispatchLoop() { c.setConnError(fmt.Errorf("decode frame: %w", err)) return } - if c.logger != nil { - c.logger.Debug("frame received", "msg_type", msgType) - } - // First valid response confirms authentication - if firstResponse { + if firstResponse && hasResponse { c.confirmAuthentication() firstResponse = false } + } +} - // Route to multiplexer (non-blocking dispatch) - c.mux.Dispatch(msgType, payload) +func (c *Connection) dispatchTransportFrame(data []byte) (bool, error) { + frames, err := protocol.DecodeFrames(data) + if err != nil { + return false, err } + var correlationID uint64 + hasResponse := false + for _, frame := range frames { + if c.logger != nil { + c.logger.Debug("frame received", "msg_type", frame.MessageType) + } + switch frame.MessageType { + case protocol.MessageTypeServerHello: + if correlationID != 0 { + return false, errors.New("CORRELATED record cannot label SERVER_HELLO") + } + if len(frame.Payload) >= 6 { + c.mux.SetCapabilities(binary.BigEndian.Uint16(frame.Payload[:2]), binary.BigEndian.Uint32(frame.Payload[2:6])) + } + case protocol.MessageTypeCorrelated: + if len(frame.Payload) != 8 || correlationID != 0 { + return false, errors.New("malformed CORRELATED record") + } + correlationID = binary.BigEndian.Uint64(frame.Payload) + if correlationID == 0 { + return false, errors.New("zero CORRELATED identifier") + } + default: + hasResponse = true + if correlationID != 0 { + c.mux.DispatchCorrelated(correlationID, frame.MessageType, frame.Payload) + correlationID = 0 + } else { + c.mux.Dispatch(frame.MessageType, frame.Payload) + } + } + } + if correlationID != 0 { + return false, errors.New("CORRELATED record did not label a response") + } + return hasResponse, nil } // handleReadError processes transport read errors. @@ -698,7 +732,20 @@ func (c *Connection) SendRequest(ctx context.Context, msgType uint16, payload [] } defer c.ReleaseRequestSlot() - frame := protocol.EncodeFrameOwned(msgType, payload) + correlationID := uint64(0) + if c.mux.CorrelationEnabled() && msgType != protocol.MessageTypeRpcRequest && msgType != protocol.MessageTypeRpcResponse { + correlationID = c.mux.NextCorrelationID() + } else { + lane := c.mux.LegacyLane(msgType) + lane.Lock() + defer lane.Unlock() + } + var frame *protocol.FrameBuffer + if correlationID == 0 { + frame = protocol.EncodeFrameOwned(msgType, payload) + } else { + frame = protocol.EncodeCorrelatedFrameOwned(correlationID, msgType, payload) + } if frame == nil { err := errors.New("encode frame") span.RecordError(err) @@ -726,9 +773,20 @@ func (c *Connection) SendRequest(ctx context.Context, msgType uint16, payload [] } c.writeMu.Lock() - c.mux.RegisterRequestWaiter(msgType, waiter, nil) + if correlationID == 0 { + c.mux.RegisterRequestWaiter(msgType, waiter, nil) + } else if !c.mux.RegisterCorrelatedRequest(correlationID, waiter) { + c.writeMu.Unlock() + return nil, errors.New("register correlated request") + } defer func() { - if c.mux.UnregisterRequestWaiter(msgType, waiter) { + var removed bool + if correlationID == 0 { + removed = c.mux.UnregisterRequestWaiter(msgType, waiter) + } else { + removed = c.mux.UnregisterCorrelatedRequest(correlationID, waiter) + } + if removed { releaseWaiter = true } }() @@ -756,7 +814,13 @@ func (c *Connection) SendRequest(ctx context.Context, msgType uint16, payload [] } return waiter.response, nil case <-ctx.Done(): - if c.mux.AbandonRequestWaiter(msgType, waiter) { + var removed bool + if correlationID == 0 { + removed = c.mux.AbandonRequestWaiter(msgType, waiter) + } else { + removed = c.mux.UnregisterCorrelatedRequest(correlationID, waiter) + } + if removed { releaseWaiter = true } span.RecordError(ctx.Err()) @@ -822,14 +886,34 @@ func (c *Connection) SendRequestWithWriter(ctx context.Context, msgType uint16, } defer c.ReleaseRequestSlot() - frame, err := protocol.EncodeFrameWithPayloadWriter(msgType, writePayload) + baseFrame, err := protocol.EncodeFrameWithPayloadWriter(msgType, writePayload) if err != nil { wrapped := fmt.Errorf("encode frame: %w", err) span.RecordError(wrapped) span.SetStatus(codes.Error, wrapped.Error()) return nil, wrapped } - defer frame.Release() + defer baseFrame.Release() + correlationID := uint64(0) + if c.mux.CorrelationEnabled() && msgType != protocol.MessageTypeRpcRequest && msgType != protocol.MessageTypeRpcResponse { + correlationID = c.mux.NextCorrelationID() + } else { + lane := c.mux.LegacyLane(msgType) + lane.Lock() + defer lane.Unlock() + } + frame := baseFrame + if correlationID != 0 { + decodedType, decodedPayload, decodeErr := protocol.DecodeFrame(baseFrame.Bytes()) + if decodeErr != nil { + return nil, decodeErr + } + frame = protocol.EncodeCorrelatedFrameOwned(correlationID, decodedType, decodedPayload) + if frame == nil { + return nil, errors.New("encode correlated frame") + } + defer frame.Release() + } waiter := acquireRequestWaiter() releaseWaiter := false @@ -850,9 +934,20 @@ func (c *Connection) SendRequestWithWriter(ctx context.Context, msgType uint16, } c.writeMu.Lock() - c.mux.RegisterRequestWaiter(msgType, waiter, nil) + if correlationID == 0 { + c.mux.RegisterRequestWaiter(msgType, waiter, nil) + } else if !c.mux.RegisterCorrelatedRequest(correlationID, waiter) { + c.writeMu.Unlock() + return nil, errors.New("register correlated request") + } defer func() { - if c.mux.UnregisterRequestWaiter(msgType, waiter) { + var removed bool + if correlationID == 0 { + removed = c.mux.UnregisterRequestWaiter(msgType, waiter) + } else { + removed = c.mux.UnregisterCorrelatedRequest(correlationID, waiter) + } + if removed { releaseWaiter = true } }() @@ -880,7 +975,13 @@ func (c *Connection) SendRequestWithWriter(ctx context.Context, msgType uint16, } return waiter.response, nil case <-ctx.Done(): - if c.mux.AbandonRequestWaiter(msgType, waiter) { + var removed bool + if correlationID == 0 { + removed = c.mux.AbandonRequestWaiter(msgType, waiter) + } else { + removed = c.mux.UnregisterCorrelatedRequest(correlationID, waiter) + } + if removed { releaseWaiter = true } span.RecordError(ctx.Err()) diff --git a/internal/core/connection/mux.go b/internal/core/connection/mux.go index 051b19b..9de4c21 100644 --- a/internal/core/connection/mux.go +++ b/internal/core/connection/mux.go @@ -7,6 +7,8 @@ import ( "log/slog" "sync" "sync/atomic" + + "github.com/cntryl/fitz-go/v2/internal/protocol" ) // pendingRequest represents one in-flight request awaiting response. @@ -225,9 +227,11 @@ type Multiplexer struct { // FIFO queue of pending requests per MessageType // Key = MessageType (100-199 for KV, 200-299 for Queue, etc.) // Value = queue of pendingRequest (oldest at front) - pending map[uint16]*requestQueue - mu sync.Mutex - handlerMu sync.RWMutex + pending map[uint16]*requestQueue + correlated map[uint64]pendingRequest + legacyLanes map[uint16]*sync.Mutex + mu sync.Mutex + handlerMu sync.RWMutex // Async delivery handlers (Notice NOTIFY, Schedule NOTIFY, RPC REQUEST to worker, RPC RESPONSE per CLIENT_SPEC.md) // notifyHandlers are keyed by message type so every subscription-capable domain can register independently. @@ -243,18 +247,56 @@ type Multiplexer struct { responsesDropped atomic.Uint64 logger *slog.Logger - closed atomic.Bool + closed atomic.Bool + protocolVersion atomic.Uint32 + capabilities atomic.Uint32 + nextCorrelationID atomic.Uint64 } // NewMultiplexer creates a new multiplexer. func NewMultiplexer() *Multiplexer { return &Multiplexer{ pending: make(map[uint16]*requestQueue), + correlated: make(map[uint64]pendingRequest), + legacyLanes: make(map[uint16]*sync.Mutex), notifyHandlers: make(map[uint16]func(subID uint64, route string, payload []byte)), rawPushHandlers: make(map[uint16]func(payload []byte)), } } +func (m *Multiplexer) SetCapabilities(protocolVersion uint16, capabilities uint32) { + m.protocolVersion.Store(uint32(protocolVersion)) + m.capabilities.Store(capabilities) +} + +func (m *Multiplexer) CorrelationEnabled() bool { + return m.capabilities.Load()&protocol.CapabilityCorrelation != 0 +} + +func (m *Multiplexer) Capabilities() (uint16, uint32) { + return uint16(m.protocolVersion.Load()), m.capabilities.Load() +} + +func (m *Multiplexer) NextCorrelationID() uint64 { + for { + id := m.nextCorrelationID.Add(1) + if id != 0 { + return id + } + } +} + +func (m *Multiplexer) LegacyLane(msgType uint16) *sync.Mutex { + m.mu.Lock() + defer m.mu.Unlock() + lane := m.legacyLanes[msgType] + if lane == nil { + lane = &sync.Mutex{} + m.legacyLanes[msgType] = lane + } + return lane +} + func (m *Multiplexer) setLogger(logger *slog.Logger) { m.logger = logger } @@ -302,6 +344,34 @@ func (m *Multiplexer) RegisterRequestWaiter(msgType uint16, waiter *requestWaite m.requestsTotal.Add(1) } +func (m *Multiplexer) RegisterCorrelatedRequest(id uint64, waiter *requestWaiter) bool { + m.mu.Lock() + defer m.mu.Unlock() + if m.closed.Load() { + waiter.fail() + return false + } + if _, exists := m.correlated[id]; exists { + return false + } + m.correlated[id] = pendingRequest{waiter: waiter} + m.requestsInFlight.Add(1) + m.requestsTotal.Add(1) + return true +} + +func (m *Multiplexer) UnregisterCorrelatedRequest(id uint64, waiter *requestWaiter) bool { + m.mu.Lock() + defer m.mu.Unlock() + req, exists := m.correlated[id] + if !exists || req.waiter != waiter { + return false + } + delete(m.correlated, id) + m.requestsInFlight.Add(-1) + return true +} + func (m *Multiplexer) UnregisterRequestWaiter(msgType uint16, waiter *requestWaiter) bool { m.mu.Lock() defer m.mu.Unlock() @@ -411,6 +481,22 @@ func (m *Multiplexer) Dispatch(msgType uint16, payload []byte) { req.waiter.deliver(payload) } +func (m *Multiplexer) DispatchCorrelated(id uint64, msgType uint16, payload []byte) { + m.mu.Lock() + req, exists := m.correlated[id] + if exists { + delete(m.correlated, id) + } + m.mu.Unlock() + if !exists { + m.Dispatch(msgType, payload) + return + } + m.requestsInFlight.Add(-1) + m.responsesTotal.Add(1) + req.waiter.deliver(payload) +} + // handleKVNotify processes KV NOTIFY messages (111): // [u64 subscription_id][string exact_route][u64 mutation_count]. func (m *Multiplexer) handleKVNotify(msgType uint16, payload []byte, handler func(subID uint64, route string, payload []byte)) { @@ -635,9 +721,17 @@ func (m *Multiplexer) Close() error { req.waiter.fail() } } + for _, req := range m.correlated { + if req.waiter != nil { + req.waiter.fail() + } + } // Clear pending requests m.pending = make(map[uint16]*requestQueue) + m.correlated = make(map[uint64]pendingRequest) + m.capabilities.Store(0) + m.protocolVersion.Store(0) return nil } diff --git a/internal/core/connection/mux_internal_test.go b/internal/core/connection/mux_internal_test.go index 6dfe5b8..3f834b6 100644 --- a/internal/core/connection/mux_internal_test.go +++ b/internal/core/connection/mux_internal_test.go @@ -11,6 +11,27 @@ import ( "github.com/stretchr/testify/require" ) +func TestShouldDispatchCorrelatedRequestsOutOfOrderGivenSameMessageType(t *testing.T) { + mux := NewMultiplexer() + first := acquireRequestWaiter() + second := acquireRequestWaiter() + t.Cleanup(func() { + releaseRequestWaiter(first) + releaseRequestWaiter(second) + require.NoError(t, mux.Close()) + }) + require.True(t, mux.RegisterCorrelatedRequest(1, first)) + require.True(t, mux.RegisterCorrelatedRequest(2, second)) + + mux.DispatchCorrelated(2, protocol.MessageTypeQueueReserve, []byte("second")) + mux.DispatchCorrelated(1, protocol.MessageTypeQueueReserve, []byte("first")) + + <-first.ready + <-second.ready + require.Equal(t, []byte("first"), first.response) + require.Equal(t, []byte("second"), second.response) +} + func TestShouldRouteRpcWorkerRequestGivenPendingRpcRequestWhenDispatchCalled(t *testing.T) { mux := NewMultiplexer() t.Cleanup(func() { diff --git a/internal/protocol/frame.go b/internal/protocol/frame.go index 0f226d0..ee446ef 100644 --- a/internal/protocol/frame.go +++ b/internal/protocol/frame.go @@ -94,6 +94,18 @@ const ( MessageTypeEscape = 0xFF ) +const ( + MessageTypeCorrelate uint16 = 2 + MessageTypeCorrelated uint16 = 3 + MessageTypeServerHello uint16 = 4 + CapabilityCorrelation uint32 = 1 << 0 +) + +type Frame struct { + MessageType uint16 + Payload []byte +} + // EncodeMessageType encodes MessageType using variable-length encoding // Per CLIENT_SPEC.md: types 0-254 = 1 byte, types 255+ = [0xFF][u16 BE] func EncodeMessageType(msgType uint16) []byte { @@ -177,6 +189,32 @@ func EncodeFrameOwned(msgType uint16, payload []byte) *FrameBuffer { return frame } +// EncodeCorrelatedFrameOwned encodes CORRELATE immediately before the labeled +// request in the same transport frame. +func EncodeCorrelatedFrameOwned(correlationID uint64, msgType uint16, payload []byte) *FrameBuffer { + if correlationID == 0 || len(payload) > MaxPayloadSize { + return nil + } + buf := getBuffer() + buf.Grow(11 + 3 + 2 + len(payload)) + buf.WriteByte(byte(MessageTypeCorrelate)) + writeU16BE(buf, 8) + var id [8]byte + binary.BigEndian.PutUint64(id[:], correlationID) + buf.Write(id[:]) + if msgType <= 254 { + buf.WriteByte(byte(msgType)) + } else { + buf.WriteByte(MessageTypeEscape) + writeU16BE(buf, msgType) + } + writeU16BE(buf, uint16(len(payload))) + buf.Write(payload) + frame := frameBufferPool.Get().(*FrameBuffer) + frame.buf = buf + return frame +} + // EncodeFrameWithPayloadWriter encodes a frame using a payload writer callback. // This avoids intermediate payload allocations by writing directly into the frame buffer. func EncodeFrameWithPayloadWriter(msgType uint16, writePayload func(*bytes.Buffer)) (*FrameBuffer, error) { @@ -251,6 +289,28 @@ func DecodeFrame(data []byte) (msgType uint16, payload []byte, err error) { return msgType, payload, nil } +// DecodeFrames decodes every TLV record in one message-bounded transport frame. +func DecodeFrames(data []byte) ([]Frame, error) { + frames := make([]Frame, 0, 2) + for len(data) > 0 { + msgType, typeLen, err := DecodeMessageType(data) + if err != nil { + return nil, fmt.Errorf("decode message type: %w", err) + } + if len(data) < typeLen+2 { + return nil, errors.New("insufficient data for length field") + } + length := int(binary.BigEndian.Uint16(data[typeLen : typeLen+2])) + end := typeLen + 2 + length + if len(data) < end { + return nil, fmt.Errorf("insufficient data for payload: need %d, have %d", length, len(data)-typeLen-2) + } + frames = append(frames, Frame{MessageType: msgType, Payload: data[typeLen+2 : end]}) + data = data[end:] + } + return frames, nil +} + // EncodeTCPFrame encodes a frame for TCP transport using buffer pools // Format: [Frame Length (u32 BE)][MessageType][Length][Payload] // where Frame Length = total size of [MessageType][Length][Payload] diff --git a/internal/protocol/frame_test.go b/internal/protocol/frame_test.go index a4fd75b..852b126 100644 --- a/internal/protocol/frame_test.go +++ b/internal/protocol/frame_test.go @@ -235,6 +235,23 @@ func TestShouldRejectTrailingBytesGivenExtraFrameDataWhenDecodeFrameCalled(t *te require.EqualError(t, err, "unexpected trailing bytes after frame payload") } +func TestShouldDecodeMultipleRecordsGivenCorrelatedTransportFrame(t *testing.T) { + label := EncodeFrame(MessageTypeCorrelated, []byte{0, 0, 0, 0, 0, 0, 0, 7}) + response := EncodeFrame(202, []byte("reserved")) + + frames, err := DecodeFrames(append(label, response...)) + + if err != nil { + t.Fatal(err) + } + if len(frames) != 2 { + t.Fatalf("got %d frames", len(frames)) + } + if frames[0].MessageType != MessageTypeCorrelated || frames[1].MessageType != 202 { + t.Fatalf("unexpected message types: %d, %d", frames[0].MessageType, frames[1].MessageType) + } +} + func TestShouldHandleMaxSizePayloadGivenExactLimitWhenEncodeFrameCalled(t *testing.T) { // Arrange // Create exactly 65535 byte payload (max allowed) diff --git a/test/queue_test.go b/test/queue_test.go index 2392912..b9ca3c0 100644 --- a/test/queue_test.go +++ b/test/queue_test.go @@ -186,6 +186,52 @@ func TestShouldLongPollGivenReserveWithOptionsWhenMessageArrivesLater(t *testing }) } +func TestShouldCorrelateSameTypeReservesGivenOutOfOrderBrokerResponses(t *testing.T) { + fixture.RunWithBothTransports(t, func(t *testing.T, transport fixture.TransportType) { + f := fixture.NewTestFixture(t, transport) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + require.NoError(t, f.Connect(ctx)) + require.Eventually(t, f.Client().CorrelationEnabled, time.Second, 10*time.Millisecond) + version, capabilities := f.Client().ServerCapabilities() + require.Equal(t, uint16(1), version) + require.Equal(t, uint32(1), capabilities) + parkedRoute := f.UniqueRoute("queue") + readyRoute := f.UniqueRoute("queue") + _, err := f.Client().Queue().Enqueue(ctx, readyRoute, []byte("second")) + require.NoError(t, err) + + type reserveResult struct { + items []*fitz.QueueItem + err error + } + parked := make(chan reserveResult, 1) + go func() { + items, reserveErr := f.Client().Queue().ReserveWithOptions( + ctx, parkedRoute, 30, + fitz.WithQueueReserveBatchSize(1), + fitz.WithQueueReserveWaitSeconds(5), + ) + parked <- reserveResult{items: items, err: reserveErr} + }() + time.Sleep(100 * time.Millisecond) + second, err := f.Client().Queue().ReserveWithOptions( + ctx, readyRoute, 30, + fitz.WithQueueReserveBatchSize(1), + fitz.WithQueueReserveWaitSeconds(0), + ) + require.NoError(t, err) + require.Len(t, second, 1) + require.Equal(t, []byte("second"), second[0].Body) + _, err = f.Client().Queue().Enqueue(ctx, parkedRoute, []byte("first")) + require.NoError(t, err) + first := <-parked + require.NoError(t, first.err) + require.Len(t, first.items, 1) + require.Equal(t, []byte("first"), first.items[0].Body) + }) +} + func TestShouldDistributeMessagesGivenMultipleConsumersWhenConcurrentReserve(t *testing.T) { fixture.RunWithBothTransports(t, func(t *testing.T, transport fixture.TransportType) { f1 := fixture.NewTestFixture(t, transport)