This is an automated email from the ASF dual-hosted git repository.

sruehl pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/plc4x.git

commit 5e4fb8e202a618eca3a491f50c825966f35d23b1
Author: Sebastian Rühl <[email protected]>
AuthorDate: Fri Jul 24 11:38:19 2026 +0200

    feat(plc4go/bacnetip): segmented request sending via peer-capability 
connection options
    
    A confirmed request whose APDU exceeds the target device's declared
    MaxApduLengthAccepted now goes out as a segmented request per ASHRAE 135
    clause 5.4: segment 0 alone, then windows of segments paced by the peer's
    SegmentAcks (actual window size honored, NAK-driven rewind, bounded
    retransmission on ack timeout). The final service response is matched by a
    pre-armed expectation that explicitly ignores SegmentAcks, since both wait
    on the same invoke id.
    
    The peer's capabilities arrive as two new connection options mirroring the
    I-Am vocabulary: PeerMaxApduLengthAccepted=<octets> and
    PeerSegmentationSupported=<segmented-both|segmented-receive|...>. With no
    declared ceiling nothing changes (requests are never segmented); with a
    ceiling but no segmentation support an oversized request fails fast with an
    actionable error instead of provoking a device-side abort. An unknown peer
    capability conservatively means no-segmentation — deliberately the opposite
    default of our OWN capability parsing.
    
    Reader and Writer share the transmit machinery (large ReadPropertyMultiple
    access lists segment too); the previously unconsumed outboundSegmenter now
    drives the wire format: segment 0 carries the service-choice byte as its
    first payload byte, follow-up segments repeat the choice in the discrete
    segmentServiceChoice field, matching the generated model's parse conditions.
    
    Roundtrip tests run a fake clause-5.4 receive-side device: window sizes 1
    and 3, NAK continuation, fail-fast without peer support, and a fitting
    request staying unsegmented.
---
 plc4go/internal/bacnetip/Configuration.go          |  15 +
 plc4go/internal/bacnetip/Configuration_plc4xgen.go |   8 +
 plc4go/internal/bacnetip/Connection.go             |   1 +
 plc4go/internal/bacnetip/DriverContext.go          |  37 ++-
 plc4go/internal/bacnetip/DriverContext_plc4xgen.go |   8 +
 plc4go/internal/bacnetip/Reader.go                 |  73 ++++-
 .../internal/bacnetip/ReaderResultDelivery_test.go |   2 +-
 plc4go/internal/bacnetip/Segmentation.go           |   9 +
 .../bacnetip/SegmentedWriteRoundtrip_test.go       | 355 +++++++++++++++++++++
 plc4go/internal/bacnetip/SenderSegmentation.go     | 259 +++++++++++++++
 plc4go/internal/bacnetip/Writer.go                 |  62 ++++
 11 files changed, 815 insertions(+), 14 deletions(-)

diff --git a/plc4go/internal/bacnetip/Configuration.go 
b/plc4go/internal/bacnetip/Configuration.go
index 424e212ef9..a5c9480445 100644
--- a/plc4go/internal/bacnetip/Configuration.go
+++ b/plc4go/internal/bacnetip/Configuration.go
@@ -103,6 +103,21 @@ type Configuration struct {
        // "<ip>:<port>" (encoded as the 6-byte B/IP DADR); for other datalinks 
a hex
        // string ("0x0C") supplies the raw MAC octets.
        RemoteAddress string
+
+       // PeerMaxApduLengthAccepted is the target device's 
MaxApduLengthAccepted in
+       // octets, as the device declared it in its I-Am. When a confirmed 
request's
+       // APDU would exceed this, the driver sends it as a segmented request
+       // (ASHRAE 135 clause 5.4) — provided PeerSegmentationSupported allows 
it.
+       // 0 (default) means unknown: requests are never segmented and an 
oversized
+       // request fails fast instead of provoking an abort from the device.
+       PeerMaxApduLengthAccepted uint16
+
+       // PeerSegmentationSupported is the target device's segmentation 
capability
+       // from its I-Am: "segmented-both", "segmented-transmit", 
"segmented-receive"
+       // or "no-segmentation". Segmented requests are only sent when the peer 
can
+       // RECEIVE segments ("segmented-both"/"segmented-receive"). Empty 
(default)
+       // means unknown and is treated as "no-segmentation".
+       PeerSegmentationSupported string
 }
 
 // ParseFromOptions populates a Configuration from the connection-URL query 
options,
diff --git a/plc4go/internal/bacnetip/Configuration_plc4xgen.go 
b/plc4go/internal/bacnetip/Configuration_plc4xgen.go
index 4866962b39..e02fe59f41 100644
--- a/plc4go/internal/bacnetip/Configuration_plc4xgen.go
+++ b/plc4go/internal/bacnetip/Configuration_plc4xgen.go
@@ -114,6 +114,14 @@ func (d *Configuration) SerializeWithWriteBuffer(ctx 
context.Context, writeBuffe
        if err := writeBuffer.WriteString("remoteAddress", 
uint32(len(d.RemoteAddress)*8), d.RemoteAddress); err != nil {
                return err
        }
+
+       if err := writeBuffer.WriteUint16("peerMaxApduLengthAccepted", 16, 
d.PeerMaxApduLengthAccepted); err != nil {
+               return err
+       }
+
+       if err := writeBuffer.WriteString("peerSegmentationSupported", 
uint32(len(d.PeerSegmentationSupported)*8), d.PeerSegmentationSupported); err 
!= nil {
+               return err
+       }
        if err := writeBuffer.PopContext("Configuration"); err != nil {
                return err
        }
diff --git a/plc4go/internal/bacnetip/Connection.go 
b/plc4go/internal/bacnetip/Connection.go
index 881d3a00c2..5f41e394d3 100644
--- a/plc4go/internal/bacnetip/Connection.go
+++ b/plc4go/internal/bacnetip/Connection.go
@@ -224,6 +224,7 @@ func (c *Connection) ReadRequestBuilder() 
apiModel.PlcReadRequestBuilder {
                        &c.invokeIdGenerator,
                        c.messageCodec,
                        c.tm,
+                       c.driverContext,
                        c.routedDest,
                        append(c._options, options.WithCustomLogger(c.log))...,
                ),
diff --git a/plc4go/internal/bacnetip/DriverContext.go 
b/plc4go/internal/bacnetip/DriverContext.go
index ea2e69c7e5..dc9704b8ac 100644
--- a/plc4go/internal/bacnetip/DriverContext.go
+++ b/plc4go/internal/bacnetip/DriverContext.go
@@ -44,18 +44,43 @@ type DriverContext struct {
        // BACnet enum value (e.g. 16 → NUM_SEGMENTS_16).
        maxSegmentsAccepted model.MaxSegmentsAccepted `stringer:"true"`
 
+       // peerMaxApduBytes is Configuration.PeerMaxApduLengthAccepted: the 
target
+       // device's APDU ceiling in octets. 0 = unknown (requests are never 
segmented).
+       peerMaxApduBytes uint16
+
+       // peerAcceptsSegmentedRequests is derived from
+       // Configuration.PeerSegmentationSupported: true when the peer declared 
it can
+       // RECEIVE segmented requests ("segmented-both"/"segmented-receive").
+       peerAcceptsSegmentedRequests bool
+
        awaitSetupComplete      bool
        awaitDisconnectComplete bool
 }
 
 func NewDriverContext(configuration Configuration) DriverContext {
        return DriverContext{
-               configuration:           configuration,
-               maxApduLengthAccepted:   
bytesToMaxApduLength(configuration.MaxApduLengthAccepted),
-               segmentation:            
stringToSegmentation(configuration.SegmentationSupported),
-               maxSegmentsAccepted:     
numToMaxSegments(configuration.MaxSegmentsAccepted),
-               awaitSetupComplete:      true,
-               awaitDisconnectComplete: true,
+               configuration:                configuration,
+               maxApduLengthAccepted:        
bytesToMaxApduLength(configuration.MaxApduLengthAccepted),
+               segmentation:                 
stringToSegmentation(configuration.SegmentationSupported),
+               maxSegmentsAccepted:          
numToMaxSegments(configuration.MaxSegmentsAccepted),
+               peerMaxApduBytes:             
configuration.PeerMaxApduLengthAccepted,
+               peerAcceptsSegmentedRequests: 
segmentationAcceptsSegmentedRequests(configuration.PeerSegmentationSupported),
+               awaitSetupComplete:           true,
+               awaitDisconnectComplete:      true,
+       }
+}
+
+// segmentationAcceptsSegmentedRequests reports whether a peer declaring the
+// given segmentation-supported string can receive segmented confirmed 
requests.
+// Unlike stringToSegmentation (which defaults OUR capability to segmented-both
+// for unknown strings), an unknown or empty PEER capability must 
conservatively
+// mean "no".
+func segmentationAcceptsSegmentedRequests(s string) bool {
+       switch s {
+       case "segmented-both", "segmented-receive":
+               return true
+       default:
+               return false
        }
 }
 
diff --git a/plc4go/internal/bacnetip/DriverContext_plc4xgen.go 
b/plc4go/internal/bacnetip/DriverContext_plc4xgen.go
index 242653fe28..7854214943 100644
--- a/plc4go/internal/bacnetip/DriverContext_plc4xgen.go
+++ b/plc4go/internal/bacnetip/DriverContext_plc4xgen.go
@@ -68,6 +68,14 @@ func (d *DriverContext) SerializeWithWriteBuffer(ctx 
context.Context, writeBuffe
                return err
        }
 
+       if err := writeBuffer.WriteUint16("peerMaxApduBytes", 16, 
d.peerMaxApduBytes); err != nil {
+               return err
+       }
+
+       if err := writeBuffer.WriteBit("peerAcceptsSegmentedRequests", 
d.peerAcceptsSegmentedRequests); err != nil {
+               return err
+       }
+
        if err := writeBuffer.WriteBit("awaitSetupComplete", 
d.awaitSetupComplete); err != nil {
                return err
        }
diff --git a/plc4go/internal/bacnetip/Reader.go 
b/plc4go/internal/bacnetip/Reader.go
index 654160301d..ee654801a1 100644
--- a/plc4go/internal/bacnetip/Reader.go
+++ b/plc4go/internal/bacnetip/Reader.go
@@ -42,6 +42,7 @@ type Reader struct {
        invokeIdGenerator *InvokeIdGenerator
        messageCodec      spi.MessageCodec
        tm                transactions.RequestTransactionManager
+       driverContext     DriverContext
        routedDest        *routedDestination // nil for local-segment 
connections
 
        maxSegmentsAccepted   readWriteModel.MaxSegmentsAccepted
@@ -52,12 +53,13 @@ type Reader struct {
        log zerolog.Logger
 }
 
-func NewReader(invokeIdGenerator *InvokeIdGenerator, messageCodec 
spi.MessageCodec, tm transactions.RequestTransactionManager, routedDest 
*routedDestination, _options ...options.WithOption) *Reader {
+func NewReader(invokeIdGenerator *InvokeIdGenerator, messageCodec 
spi.MessageCodec, tm transactions.RequestTransactionManager, driverContext 
DriverContext, routedDest *routedDestination, _options ...options.WithOption) 
*Reader {
        customLogger := 
options.ExtractCustomLoggerOrDefaultToGlobal(_options...)
        return &Reader{
                invokeIdGenerator: invokeIdGenerator,
                messageCodec:      messageCodec,
                tm:                tm,
+               driverContext:     driverContext,
                routedDest:        routedDest,
 
                maxSegmentsAccepted:   
readWriteModel.MaxSegmentsAccepted_MORE_THAN_64_SEGMENTS,
@@ -135,6 +137,27 @@ func (m *Reader) Read(ctx context.Context, readRequest 
apiModel.PlcReadRequest)
                        nil,
                )
 
+               // If the request exceeds the peer's declared APDU ceiling 
(large
+               // ReadPropertyMultiple access lists), it has to go out as a 
segmented
+               // request (ASHRAE 135 clause 5.4) — or fail fast when the peer 
can't
+               // receive segments, instead of provoking an abort.
+               var segmentedPayload []byte
+               if m.driverContext.peerMaxApduBytes > 0 {
+                       payload, serErr := serviceRequest.Serialize()
+                       if serErr != nil {
+                               utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrap(serErr, 
"Error serializing read request")))
+                               return
+                       }
+                       if m.driverContext.needsSegmentedRequest(len(payload)) {
+                               if 
!m.driverContext.peerAcceptsSegmentedRequests {
+                                       utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcReadRequestResult(readRequest, nil,
+                                               errors.Errorf("read request of 
%d bytes exceeds the peer's max APDU of %d and the peer does not support 
segmented requests", len(payload), m.driverContext.peerMaxApduBytes)))
+                                       return
+                               }
+                               segmentedPayload = payload
+                       }
+               }
+
                // Start a new request-transaction (Is ended in the 
response-handler)
                transaction := m.tm.StartTransaction("read")
                transaction.Submit("readOperation", func(transactionContext 
context.Context, transaction transactions.RequestTransaction) {
@@ -144,9 +167,8 @@ func (m *Reader) Read(ctx context.Context, readRequest 
apiModel.PlcReadRequest)
                        readCtx := ctx
                        ctx, cancel := context.WithCancel(ctx)
                        context.AfterFunc(transactionContext, cancel)
-                       // Send the  over the wire
-                       m.log.Trace().Msg("Send ")
-                       if err := m.messageCodec.SendRequest(ctx, "read", 
wrapAPDU(apdu, true, m.routedDest), func(message spi.Message) bool {
+
+                       acceptsReadResponse := func(message spi.Message) bool {
                                bvlc, ok := message.(readWriteModel.BVLC)
                                if !ok {
                                        m.log.Debug().Type("bvlc", 
bvlc).Msg("Received strange type")
@@ -168,7 +190,8 @@ func (m *Reader) Read(ctx context.Context, readRequest 
apiModel.PlcReadRequest)
                                } else {
                                        return invokeIdFromApdu == invokeId
                                }
-                       }, func(message spi.Message) error {
+                       }
+                       handleReadResponse := func(message spi.Message) error {
                                // Convert the response into an
                                m.log.Trace().Msg("convert response to ")
                                apdu := 
message.(readWriteModel.BVLC).(interface{ GetNpdu() readWriteModel.NPDU 
}).GetNpdu().GetApdu()
@@ -204,14 +227,50 @@ func (m *Reader) Read(ctx context.Context, readRequest 
apiModel.PlcReadRequest)
                                        nil,
                                ))
                                return transaction.EndRequest()
-                       }, func(err error) error {
+                       }
+                       handleReadError := func(err error) error {
                                utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcReadRequestResult(
                                        readRequest,
                                        nil,
                                        errors.Wrap(err, "got timeout while 
waiting for response"),
                                ))
                                return transaction.EndRequest()
-                       }); err != nil {
+                       }
+
+                       if segmentedPayload != nil {
+                               // Register the response expectation BEFORE 
driving the segments:
+                               // the peer may answer right after acking the 
final segment. The
+                               // matcher must not consume the peer's 
SegmentAcks — those belong
+                               // to the sender's own expectations.
+                               m.messageCodec.Expect(ctx, 
"segmentedReadResponse",
+                                       
responseMatcherExcludingSegmentAcks(acceptsReadResponse),
+                                       handleReadResponse,
+                                       handleReadError,
+                               )
+
+                               sender := &segmentedRequestSender{
+                                       messageCodec:  m.messageCodec,
+                                       routedDest:    m.routedDest,
+                                       driverContext: m.driverContext,
+                                       log:           m.log,
+                               }
+                               if err := sender.send(ctx, invokeId, 
segmentedPayload); err != nil {
+                                       // Deliver the real cause first (the 
buffered result channel
+                                       // takes it), then cancel so the armed 
response expectation
+                                       // resolves without delivering a 
second, less specific error.
+                                       utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcReadRequestResult(
+                                               readRequest,
+                                               nil,
+                                               errors.Wrap(err, "error sending 
segmented read request"),
+                                       ))
+                                       cancel()
+                               }
+                               return
+                       }
+
+                       // Send the  over the wire
+                       m.log.Trace().Msg("Send ")
+                       if err := m.messageCodec.SendRequest(ctx, "read", 
wrapAPDU(apdu, true, m.routedDest), acceptsReadResponse, handleReadResponse, 
handleReadError); err != nil {
                                utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcReadRequestResult(
                                        readRequest,
                                        nil,
diff --git a/plc4go/internal/bacnetip/ReaderResultDelivery_test.go 
b/plc4go/internal/bacnetip/ReaderResultDelivery_test.go
index a6a74918dc..b592d12570 100644
--- a/plc4go/internal/bacnetip/ReaderResultDelivery_test.go
+++ b/plc4go/internal/bacnetip/ReaderResultDelivery_test.go
@@ -76,7 +76,7 @@ func (c *captureCodec) GetDefaultIncomingMessageChannel() 
chan spi.Message { ret
 func TestReader_lateErrorHandlerAfterFailedSendMustNotBlock(t *testing.T) {
        codec := newCaptureCodec(errors.New("send failed: broken pipe"))
        tm := transactions.NewRequestTransactionManager(1)
-       reader := NewReader(&InvokeIdGenerator{}, codec, tm, nil)
+       reader := NewReader(&InvokeIdGenerator{}, codec, tm, 
NewDriverContext(createDefaultConfiguration()), nil)
 
        objType := readWriteModel.BACnetObjectType_ANALOG_INPUT
        propId := readWriteModel.BACnetPropertyIdentifier_PRESENT_VALUE
diff --git a/plc4go/internal/bacnetip/Segmentation.go 
b/plc4go/internal/bacnetip/Segmentation.go
index c6bcc4575a..de964d22ea 100644
--- a/plc4go/internal/bacnetip/Segmentation.go
+++ b/plc4go/internal/bacnetip/Segmentation.go
@@ -135,6 +135,15 @@ func NewOutboundSegmenter(invokeId uint8, payload []byte, 
maxApdu uint16, window
 // HasMore reports whether more segments remain to be sent.
 func (s *outboundSegmenter) HasMore() bool { return !s.done }
 
+// TotalSegments returns how many segments the payload spans at the configured
+// per-segment budget. An empty payload still occupies one (empty) segment.
+func (s *outboundSegmenter) TotalSegments() int {
+       if len(s.payload) == 0 {
+               return 1
+       }
+       return (len(s.payload) + int(s.maxSegment) - 1) / int(s.maxSegment)
+}
+
 // NextSegment returns the next chunk of payload bytes along with a flag
 // indicating whether more segments follow. Advances the internal cursor.
 func (s *outboundSegmenter) NextSegment() (seq uint8, segment []byte, 
moreFollows bool) {
diff --git a/plc4go/internal/bacnetip/SegmentedWriteRoundtrip_test.go 
b/plc4go/internal/bacnetip/SegmentedWriteRoundtrip_test.go
new file mode 100644
index 0000000000..67907268e9
--- /dev/null
+++ b/plc4go/internal/bacnetip/SegmentedWriteRoundtrip_test.go
@@ -0,0 +1,355 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package bacnetip
+
+import (
+       "context"
+       "fmt"
+       "net"
+       "strings"
+       "sync"
+       "testing"
+       "time"
+
+       "github.com/rs/zerolog"
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+
+       plc4go "github.com/apache/plc4x/plc4go/pkg/api"
+       apiModel "github.com/apache/plc4x/plc4go/pkg/api/model"
+       apiTransports "github.com/apache/plc4x/plc4go/pkg/api/transports"
+       model 
"github.com/apache/plc4x/plc4go/protocols/bacnetip/readwrite/model"
+       "github.com/apache/plc4x/plc4go/spi/options"
+       "github.com/apache/plc4x/plc4go/spi/testutils"
+)
+
+// fakeSegmentedDevice is a UDP BACnet/IP device that understands SEGMENTED
+// confirmed requests (ASHRAE 135 clause 5.4 receive side): it collects the
+// segments, acknowledges per its actual window size, reassembles the service
+// request, and answers with a SimpleAck once the final segment arrived.
+// Unsegmented confirmed requests get an immediate SimpleAck.
+type fakeSegmentedDevice struct {
+       conn *net.UDPConn
+       wg   sync.WaitGroup
+       log  zerolog.Logger
+       ctx  context.Context
+
+       // actualWindowSize is what the device answers in its SegmentAcks (and 
how
+       // many segments it lets pass between acks).
+       actualWindowSize uint8
+       // nakFirstSegment, when true, answers segment 0 with a negative ack
+       // (still confirming seq 0) — the sender must treat it as "received up 
to
+       // 0" and continue with segment 1.
+       nakFirstSegment bool
+
+       mu               sync.Mutex
+       buf              []byte
+       sinceAck         uint8
+       segmentsReceived int
+       unsegmented      int
+       reassembled      model.BACnetConfirmedServiceRequest
+       nakSent          bool
+}
+
+func startFakeSegmentedDevice(t *testing.T, log zerolog.Logger, 
actualWindowSize uint8) *fakeSegmentedDevice {
+       t.Helper()
+       conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 
1), Port: 0})
+       require.NoError(t, err)
+       d := &fakeSegmentedDevice{conn: conn, log: log, ctx: t.Context(), 
actualWindowSize: actualWindowSize}
+       d.wg.Add(1)
+       go d.serve()
+       return d
+}
+
+func (d *fakeSegmentedDevice) port() int { return 
d.conn.LocalAddr().(*net.UDPAddr).Port }
+
+func (d *fakeSegmentedDevice) stop() {
+       _ = d.conn.Close()
+       d.wg.Wait()
+}
+
+func (d *fakeSegmentedDevice) stats() (segments, unsegmented int, reassembled 
model.BACnetConfirmedServiceRequest) {
+       d.mu.Lock()
+       defer d.mu.Unlock()
+       return d.segmentsReceived, d.unsegmented, d.reassembled
+}
+
+func (d *fakeSegmentedDevice) serve() {
+       defer d.wg.Done()
+       buf := make([]byte, 4096)
+       for {
+               n, src, err := d.conn.ReadFromUDP(buf)
+               if err != nil {
+                       return // socket closed
+               }
+               data := make([]byte, n)
+               copy(data, buf[:n])
+
+               bvlc, err := model.BVLCParse[model.BVLC](d.ctx, data)
+               if err != nil {
+                       d.log.Error().Err(err).Msg("fake segmented device: 
parse BVLC")
+                       continue
+               }
+               npduRetriever, ok := bvlc.(interface{ GetNpdu() model.NPDU })
+               if !ok {
+                       continue
+               }
+               cr, ok := 
npduRetriever.GetNpdu().GetApdu().(model.APDUConfirmedRequest)
+               if !ok {
+                       continue
+               }
+
+               for _, reply := range d.handleConfirmedRequest(cr) {
+                       theBytes, err := reply.Serialize()
+                       if err != nil {
+                               d.log.Error().Err(err).Msg("fake segmented 
device: serialize reply")
+                               continue
+                       }
+                       if _, err := d.conn.WriteToUDP(theBytes, src); err != 
nil {
+                               d.log.Error().Err(err).Msg("fake segmented 
device: send reply")
+                       }
+               }
+       }
+}
+
+// handleConfirmedRequest implements the device-side segment collection and
+// returns the BVLC replies to put on the wire (segment acks and/or the final
+// SimpleAck).
+func (d *fakeSegmentedDevice) handleConfirmedRequest(cr 
model.APDUConfirmedRequest) []model.BVLC {
+       d.mu.Lock()
+       defer d.mu.Unlock()
+
+       invokeId := cr.GetInvokeId()
+
+       if !cr.GetSegmentedMessage() {
+               d.unsegmented++
+               return []model.BVLC{wrapAPDU(model.NewAPDUSimpleAck(invokeId, 
cr.GetServiceRequest().GetServiceChoice()), false, nil)}
+       }
+
+       seqPtr := cr.GetSequenceNumber()
+       if seqPtr == nil {
+               d.log.Error().Msg("fake segmented device: segmented request 
without sequence number")
+               return nil
+       }
+       seq := *seqPtr
+       d.segmentsReceived++
+
+       if seq == 0 {
+               // First segment: raw bytes start with the service-choice octet.
+               d.buf = append([]byte{}, cr.GetSegment()...)
+       } else {
+               if cr.GetSegmentServiceChoice() == nil {
+                       d.log.Error().Uint8("seq", seq).Msg("fake segmented 
device: follow-up segment without service choice")
+                       return nil
+               }
+               d.buf = append(d.buf, cr.GetSegment()...)
+       }
+
+       var replies []model.BVLC
+
+       if d.nakFirstSegment && seq == 0 && !d.nakSent {
+               // Negative ack still confirming segment 0 — sender must 
continue at 1.
+               d.nakSent = true
+               d.sinceAck = 0
+               replies = append(replies, 
wrapAPDU(model.NewAPDUSegmentAck(true, true, invokeId, 0, d.actualWindowSize), 
false, nil))
+               return replies
+       }
+
+       d.sinceAck++
+       final := !cr.GetMoreFollows()
+       // Clause 5.4.5: segment 0 is always acknowledged on its own; 
afterwards the
+       // device acks every actual-window-size segments and the final segment.
+       if seq == 0 || d.sinceAck >= d.actualWindowSize || final {
+               d.sinceAck = 0
+               replies = append(replies, 
wrapAPDU(model.NewAPDUSegmentAck(false, true, invokeId, seq, 
d.actualWindowSize), false, nil))
+       }
+
+       if final {
+               reassembled, err := 
model.BACnetConfirmedServiceRequestParse[model.BACnetConfirmedServiceRequest](d.ctx,
 d.buf, uint32(len(d.buf)))
+               if err != nil {
+                       d.log.Error().Err(err).Int("bytes", 
len(d.buf)).Msg("fake segmented device: reassembled payload does not parse")
+                       return replies
+               }
+               d.reassembled = reassembled
+               replies = append(replies, 
wrapAPDU(model.NewAPDUSimpleAck(invokeId, reassembled.GetServiceChoice()), 
false, nil))
+       }
+       return replies
+}
+
+// newSegmentedTestConnection opens a driver connection declaring the peer's
+// APDU ceiling and segmentation capability via the new connection options.
+func newSegmentedTestConnection(t *testing.T, log zerolog.Logger, port int, 
peerOptions string) plc4go.PlcConnection {
+       t.Helper()
+       dm := plc4go.NewPlcDriverManager(options.WithCustomLogger(log))
+       dm.RegisterDriver(NewDriver(options.WithCustomLogger(log)))
+       apiTransports.RegisterUdpTransport(dm)
+
+       connStr := 
fmt.Sprintf("bacnet-ip:udp://127.0.0.1:%d?local-port=0&ApduTimeoutMs=3000&%s", 
port, peerOptions)
+       ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second)
+       t.Cleanup(cancel)
+       conn, err := dm.GetConnection(ctx, connStr)
+       require.NoError(t, err)
+       t.Cleanup(func() { assert.NoError(t, conn.Close()) })
+       return conn
+}
+
+// executeLargeWrite issues a WritePropertyMultiple across enough tags that the
+// serialized request exceeds a 206-byte peer APDU ceiling.
+func executeLargeWrite(t *testing.T, conn plc4go.PlcConnection, tags int) 
(apiModel.PlcWriteResponse, error) {
+       t.Helper()
+       builder := conn.WriteRequestBuilder()
+       for i := 0; i < tags; i++ {
+               builder.AddTagAddress(fmt.Sprintf("pv%d", i), 
fmt.Sprintf("ANALOG_VALUE,%d/PRESENT_VALUE", i+1), float32(i)+0.5)
+       }
+       wr, err := builder.Build()
+       require.NoError(t, err)
+       ctx, cancel := context.WithTimeout(t.Context(), 15*time.Second)
+       t.Cleanup(cancel)
+       select {
+       case <-ctx.Done():
+               return nil, ctx.Err()
+       case res := <-wr.Execute(ctx):
+               return res.GetResponse(), res.GetErr()
+       }
+}
+
+// TestSegmentedWrite_RoundtripWindow1 pins the core clause 5.4 transmit flow
+// against a device acking every segment (actual window size 1): the oversized
+// WritePropertyMultiple must arrive as multiple segments whose reassembly
+// parses back to the original request, and the write must resolve OK.
+func TestSegmentedWrite_RoundtripWindow1(t *testing.T) {
+       log := testutils.ProduceTestingLogger(t)
+       device := startFakeSegmentedDevice(t, log, 1)
+       t.Cleanup(device.stop)
+
+       conn := newSegmentedTestConnection(t, log, device.port(),
+               
"PeerMaxApduLengthAccepted=206&PeerSegmentationSupported=segmented-both")
+
+       response, err := executeLargeWrite(t, conn, 24)
+       require.NoError(t, err)
+
+       segments, unsegmented, reassembled := device.stats()
+       t.Logf("device saw %d segments (%d unsegmented requests)", segments, 
unsegmented)
+       require.GreaterOrEqual(t, segments, 2, "request must have been 
segmented")
+       assert.Zero(t, unsegmented, "the oversized request must not go out 
unsegmented")
+       require.NotNil(t, reassembled, "device must have reassembled the 
request")
+       wpm, ok := 
reassembled.(model.BACnetConfirmedServiceRequestWritePropertyMultiple)
+       require.True(t, ok, "reassembled request must be a 
WritePropertyMultiple, got %T", reassembled)
+       assert.Len(t, wpm.GetData(), 24, "all write-access specs must survive 
reassembly")
+
+       for i := 0; i < 12; i++ {
+               assert.Equal(t, apiModel.PlcResponseCode_OK, 
response.GetResponseCode(fmt.Sprintf("pv%d", i)))
+       }
+}
+
+// TestSegmentedWrite_RoundtripWindow3 exercises the windowed burst path: the
+// device only acks every 3rd segment, so the sender must honor the actual
+// window size from the first ack and keep multiple segments in flight.
+func TestSegmentedWrite_RoundtripWindow3(t *testing.T) {
+       log := testutils.ProduceTestingLogger(t)
+       device := startFakeSegmentedDevice(t, log, 3)
+       t.Cleanup(device.stop)
+
+       conn := newSegmentedTestConnection(t, log, device.port(),
+               
"PeerMaxApduLengthAccepted=206&PeerSegmentationSupported=segmented-both")
+
+       response, err := executeLargeWrite(t, conn, 50)
+       require.NoError(t, err)
+
+       segments, _, reassembled := device.stats()
+       t.Logf("device saw %d segments", segments)
+       require.GreaterOrEqual(t, segments, 5, "a 50-tag request at a 206-byte 
ceiling must span several segments")
+       require.NotNil(t, reassembled)
+       wpm, ok := 
reassembled.(model.BACnetConfirmedServiceRequestWritePropertyMultiple)
+       require.True(t, ok)
+       assert.Len(t, wpm.GetData(), 50)
+
+       assert.Equal(t, apiModel.PlcResponseCode_OK, 
response.GetResponseCode("pv0"))
+       assert.Equal(t, apiModel.PlcResponseCode_OK, 
response.GetResponseCode("pv49"))
+}
+
+// TestSegmentedWrite_NegativeAckContinues pins the NAK handling: a negative
+// SegmentAck confirming segment 0 must not abort the transfer — the sender
+// resumes with segment 1 and the write still completes.
+func TestSegmentedWrite_NegativeAckContinues(t *testing.T) {
+       log := testutils.ProduceTestingLogger(t)
+       device := startFakeSegmentedDevice(t, log, 1)
+       device.nakFirstSegment = true
+       t.Cleanup(device.stop)
+
+       conn := newSegmentedTestConnection(t, log, device.port(),
+               
"PeerMaxApduLengthAccepted=206&PeerSegmentationSupported=segmented-both")
+
+       response, err := executeLargeWrite(t, conn, 24)
+       require.NoError(t, err)
+       _, _, reassembled := device.stats()
+       require.NotNil(t, reassembled, "transfer must complete despite the NAK 
on segment 0")
+       assert.Equal(t, apiModel.PlcResponseCode_OK, 
response.GetResponseCode("pv0"))
+}
+
+// TestSegmentedWrite_PeerWithoutSegmentationFailsFast: with a declared APDU
+// ceiling but no segmentation support, an oversized request must fail fast
+// with an actionable error instead of provoking a device-side abort.
+func TestSegmentedWrite_PeerWithoutSegmentationFailsFast(t *testing.T) {
+       log := testutils.ProduceTestingLogger(t)
+       device := startFakeSegmentedDevice(t, log, 1)
+       t.Cleanup(device.stop)
+
+       conn := newSegmentedTestConnection(t, log, device.port(),
+               
"PeerMaxApduLengthAccepted=206&PeerSegmentationSupported=no-segmentation")
+
+       _, err := executeLargeWrite(t, conn, 24)
+       require.Error(t, err)
+       assert.Contains(t, err.Error(), "does not support segmented requests")
+       segments, unsegmented, _ := device.stats()
+       assert.Zero(t, segments, "nothing may reach the wire")
+       assert.Zero(t, unsegmented, "nothing may reach the wire")
+}
+
+// TestSegmentedWrite_SmallRequestStaysUnsegmented guards the common path: with
+// peer capabilities declared, a request that FITS the ceiling must go out as a
+// plain unsegmented confirmed request.
+func TestSegmentedWrite_SmallRequestStaysUnsegmented(t *testing.T) {
+       log := testutils.ProduceTestingLogger(t)
+       device := startFakeSegmentedDevice(t, log, 1)
+       t.Cleanup(device.stop)
+
+       conn := newSegmentedTestConnection(t, log, device.port(),
+               
"PeerMaxApduLengthAccepted=480&PeerSegmentationSupported=segmented-both")
+
+       response, err := executeLargeWrite(t, conn, 2)
+       require.NoError(t, err)
+       segments, unsegmented, _ := device.stats()
+       assert.Zero(t, segments, "a fitting request must not be segmented")
+       assert.Equal(t, 1, unsegmented)
+       assert.Equal(t, apiModel.PlcResponseCode_OK, 
response.GetResponseCode("pv0"))
+}
+
+// TestSegmentationOptionParsing_Unknown guards the conservative default: an
+// unknown or absent peer segmentation string must NOT enable segmented sends.
+func TestSegmentationOptionParsing_Unknown(t *testing.T) {
+       for _, s := range []string{"", "bogus", "SEGMENTED-BOTH", 
strings.ToUpper("segmented-receive")} {
+               assert.False(t, segmentationAcceptsSegmentedRequests(s), "%q 
must not enable segmented requests", s)
+       }
+       assert.True(t, segmentationAcceptsSegmentedRequests("segmented-both"))
+       assert.True(t, 
segmentationAcceptsSegmentedRequests("segmented-receive"))
+       assert.False(t, 
segmentationAcceptsSegmentedRequests("segmented-transmit"),
+               "a peer that can only TRANSMIT segments cannot receive ours")
+}
diff --git a/plc4go/internal/bacnetip/SenderSegmentation.go 
b/plc4go/internal/bacnetip/SenderSegmentation.go
new file mode 100644
index 0000000000..85b7efd8fe
--- /dev/null
+++ b/plc4go/internal/bacnetip/SenderSegmentation.go
@@ -0,0 +1,259 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   https://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package bacnetip
+
+import (
+       "context"
+       "time"
+
+       "github.com/rs/zerolog"
+
+       readWriteModel 
"github.com/apache/plc4x/plc4go/protocols/bacnetip/readwrite/model"
+       "github.com/apache/plc4x/plc4go/spi"
+       "github.com/apache/plc4x/plc4go/spi/errors"
+)
+
+// segmentAckWaitTimeout bounds how long we wait for the peer's SegmentAck 
after
+// transmitting a window of request segments (ASHRAE 135 T-seg, default 5s).
+const segmentAckWaitTimeout = 5 * time.Second
+
+// proposedSegmentWindowSize is the window size we propose in every request
+// segment. The peer's SegmentAck answers with the ACTUAL window size, which
+// then governs how many segments we put on the wire per ack.
+const proposedSegmentWindowSize = 16
+
+// maxOutboundRequestSegments caps a segmented request at the uint8 sequence
+// space so we never have to disambiguate wrapped sequence numbers. A request
+// this large (~250 * peer-max-APDU bytes) exceeds any real device's
+// max-segments-accepted anyway.
+const maxOutboundRequestSegments = 255
+
+// unsegmentedConfirmedRequestHeaderBytes is the fixed APDU header of an
+// UNsegmented confirmed request: PDU-type/flags octet, max-segments/max-APDU
+// octet, and the invoke id. Everything after it is the serialized service
+// request, so `header + len(serviceRequest)` is the wire APDU size used to
+// decide whether a request fits the peer's MaxApduLengthAccepted.
+const unsegmentedConfirmedRequestHeaderBytes = 3
+
+// segmentedRequestSender transmits an oversized confirmed request as a
+// sequence of APDUConfirmedRequest segments per ASHRAE 135 clause 5.4
+// (segment 0 alone, then windows of segments each answered by a SegmentAck,
+// with NAK-driven rewind and bounded retransmission on ack timeout).
+//
+// The caller must register its expectation for the FINAL service response
+// (excluding APDUSegmentAck — see responseMatcherExcludingSegmentAcks) BEFORE
+// calling send: the peer may answer immediately after acking the last segment.
+type segmentedRequestSender struct {
+       messageCodec  spi.MessageCodec
+       routedDest    *routedDestination
+       driverContext DriverContext
+       log           zerolog.Logger
+}
+
+// needsSegmentedRequest reports whether the serialized service-request payload
+// exceeds the peer's declared APDU ceiling. With an unknown ceiling (0) the
+// answer is always false — we then send unsegmented exactly as before.
+func (d DriverContext) needsSegmentedRequest(payloadBytes int) bool {
+       return d.peerMaxApduBytes > 0 && 
payloadBytes+unsegmentedConfirmedRequestHeaderBytes > int(d.peerMaxApduBytes)
+}
+
+// send drives the segmented transmission of payload (the serialized
+// BACnetConfirmedServiceRequest, starting with its service-choice byte) under
+// the given invoke id. It returns once the peer has acknowledged the final
+// segment; the service response itself is delivered through the caller's
+// pre-registered expectation.
+func (s *segmentedRequestSender) send(ctx context.Context, invokeId uint8, 
payload []byte) error {
+       if len(payload) == 0 {
+               return errors.New("empty confirmed-request payload")
+       }
+       segmenter := NewOutboundSegmenter(invokeId, payload, 
s.driverContext.peerMaxApduBytes, proposedSegmentWindowSize)
+       total := segmenter.TotalSegments()
+       if total > maxOutboundRequestSegments {
+               return errors.Errorf("confirmed request of %d bytes needs %d 
segments, exceeding the %d-segment ceiling", len(payload), total, 
maxOutboundRequestSegments)
+       }
+       serviceChoice := readWriteModel.BACnetConfirmedServiceChoice(payload[0])
+
+       s.log.Debug().
+               Uint8("invokeId", invokeId).
+               Int("payloadBytes", len(payload)).
+               Int("segments", total).
+               Uint16("peerMaxApdu", s.driverContext.peerMaxApduBytes).
+               Msg("sending segmented confirmed request")
+
+       // Window is 1 until the first SegmentAck reveals the peer's actual 
window
+       // size (clause 5.4.4: segment 0 is sent alone).
+       window := uint8(1)
+       lastAcked := -1
+       retries := uint8(0)
+
+       for {
+               // Register the ack expectation BEFORE putting segments on the 
wire,
+               // otherwise the peer's ack could race ahead of the expectation.
+               ackCh, errCh := s.expectSegmentAck(ctx, invokeId)
+
+               sent := 0
+               for sent < int(window) && segmenter.HasMore() {
+                       seq, chunk, moreFollows := segmenter.NextSegment()
+                       apdu := s.segmentApdu(invokeId, seq, moreFollows, 
chunk, serviceChoice)
+                       if err := s.messageCodec.Send(ctx, "requestSegment", 
wrapAPDU(apdu, true, s.routedDest)); err != nil {
+                               return errors.Wrapf(err, "error sending request 
segment %d", seq)
+                       }
+                       sent++
+               }
+
+               ack, err := s.awaitMeaningfulAck(ctx, invokeId, lastAcked, 
ackCh, errCh)
+               if err != nil {
+                       if ctx.Err() != nil {
+                               return errors.Wrap(ctx.Err(), "context done 
awaiting segment ack")
+                       }
+                       // Ack timeout: retransmit from the last acknowledged 
segment,
+                       // bounded by the configured APDU retry budget.
+                       retries++
+                       if retries > s.driverContext.configuration.ApduRetries {
+                               return errors.Wrapf(err, "no segment ack after 
%d retries", s.driverContext.configuration.ApduRetries)
+                       }
+                       segmenter.Rewind(uint8(lastAcked + 1))
+                       continue
+               }
+               retries = 0
+
+               // The actual window size in the peer's ack is authoritative 
for the
+               // remainder of the transmission (capped at our own proposal).
+               if actual := ack.GetActualWindowSize(); actual > 0 {
+                       window = min(actual, proposedSegmentWindowSize)
+               }
+
+               if ackSeq := int(ack.GetSequenceNumber()); ackSeq > lastAcked {
+                       lastAcked = ackSeq
+               }
+
+               if lastAcked >= total-1 {
+                       if !ack.GetNegativeAck() {
+                               return nil
+                       }
+                       // A NAK naming the final segment is a protocol oddity; 
resend just
+                       // the final segment rather than looping forever.
+                       segmenter.Rewind(uint8(total - 1))
+                       lastAcked = total - 2
+                       continue
+               }
+
+               // Resume right after the last segment the peer confirmed in 
order.
+               // After an in-order positive ack this is a no-op rewind; after 
a NAK
+               // (or an ack for a partial window) it retransmits the missing 
tail.
+               segmenter.Rewind(uint8(lastAcked + 1))
+       }
+}
+
+// awaitMeaningfulAck waits on the registered expectation and swallows stale
+// duplicate acks (positive acks for segments at or below lastAcked), re-arming
+// the expectation until a meaningful ack, an error, or ctx cancellation.
+func (s *segmentedRequestSender) awaitMeaningfulAck(ctx context.Context, 
invokeId uint8, lastAcked int, ackCh <-chan readWriteModel.APDUSegmentAck, 
errCh <-chan error) (readWriteModel.APDUSegmentAck, error) {
+       for {
+               select {
+               case <-ctx.Done():
+                       return nil, ctx.Err()
+               case err := <-errCh:
+                       return nil, err
+               case ack := <-ackCh:
+                       if !ack.GetNegativeAck() && 
int(ack.GetSequenceNumber()) <= lastAcked {
+                               s.log.Debug().Uint8("ackSeq", 
ack.GetSequenceNumber()).Msg("ignoring stale duplicate segment ack")
+                               ackCh, errCh = s.expectSegmentAck(ctx, invokeId)
+                               continue
+                       }
+                       return ack, nil
+               }
+       }
+}
+
+// segmentApdu builds one APDUConfirmedRequest segment. Segment 0 carries the
+// service-choice byte as the first payload byte (the generated model only
+// parses/serializes a discrete segmentServiceChoice for sequence numbers > 0),
+// so chunks slice the serialized service request contiguously and the peer's
+// reassembly concatenates them back to the original payload.
+func (s *segmentedRequestSender) segmentApdu(invokeId uint8, seq uint8, 
moreFollows bool, chunk []byte, serviceChoice 
readWriteModel.BACnetConfirmedServiceChoice) 
readWriteModel.APDUConfirmedRequest {
+       seqV := seq
+       windowV := uint8(proposedSegmentWindowSize)
+       var choicePtr *readWriteModel.BACnetConfirmedServiceChoice
+       if seq != 0 {
+               choicePtr = &serviceChoice
+       }
+       return readWriteModel.NewAPDUConfirmedRequest(
+               true,
+               moreFollows,
+               true,
+               s.driverContext.maxSegmentsAccepted,
+               s.driverContext.maxApduLengthAccepted,
+               invokeId,
+               &seqV,
+               &windowV,
+               nil,
+               choicePtr,
+               chunk,
+       )
+}
+
+// expectSegmentAck registers a one-shot expectation for the next SegmentAck
+// with the given invoke id, returning channels for the ack or an 
error/timeout.
+func (s *segmentedRequestSender) expectSegmentAck(ctx context.Context, 
invokeId uint8) (<-chan readWriteModel.APDUSegmentAck, <-chan error) {
+       ackCh := make(chan readWriteModel.APDUSegmentAck, 1)
+       errCh := make(chan error, 1)
+
+       expectCtx, cancel := context.WithTimeout(ctx, segmentAckWaitTimeout)
+
+       m := s.messageCodec
+       m.Expect(expectCtx, "requestSegmentAck",
+               func(message spi.Message) bool {
+                       apdu, ok := apduFromMessage(message)
+                       if !ok {
+                               return false
+                       }
+                       ack, ok := apdu.(readWriteModel.APDUSegmentAck)
+                       return ok && ack.GetOriginalInvokeId() == invokeId
+               },
+               func(message spi.Message) error {
+                       cancel()
+                       apdu, _ := apduFromMessage(message)
+                       ackCh <- apdu.(readWriteModel.APDUSegmentAck)
+                       return nil
+               },
+               func(err error) error {
+                       cancel()
+                       errCh <- err
+                       return nil
+               },
+       )
+       return ackCh, errCh
+}
+
+// responseMatcherExcludingSegmentAcks wraps a response matcher so it never
+// consumes the peer's SegmentAcks for our own request segments — those belong
+// to the segmentedRequestSender's expectations, which are armed concurrently
+// on the same invoke id.
+func responseMatcherExcludingSegmentAcks(matcher func(spi.Message) bool) 
func(spi.Message) bool {
+       return func(message spi.Message) bool {
+               if apdu, ok := apduFromMessage(message); ok {
+                       if _, isSegmentAck := 
apdu.(readWriteModel.APDUSegmentAck); isSegmentAck {
+                               return false
+                       }
+               }
+               return matcher(message)
+       }
+}
diff --git a/plc4go/internal/bacnetip/Writer.go 
b/plc4go/internal/bacnetip/Writer.go
index a054e66aa7..1f019d6508 100644
--- a/plc4go/internal/bacnetip/Writer.go
+++ b/plc4go/internal/bacnetip/Writer.go
@@ -106,7 +106,69 @@ func (m *Writer) Write(ctx context.Context, writeRequest 
apiModel.PlcWriteReques
                        nil,
                )
 
+               // If the request exceeds the peer's declared APDU ceiling, it 
has to go
+               // out as a segmented request (ASHRAE 135 clause 5.4) — or fail 
fast when
+               // the peer can't receive segments, instead of provoking an 
abort.
+               var segmentedPayload []byte
+               if m.driverContext.peerMaxApduBytes > 0 {
+                       payload, serErr := serviceRequest.Serialize()
+                       if serErr != nil {
+                               utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(serErr, 
"Error serializing WriteProperty request")))
+                               return
+                       }
+                       if m.driverContext.needsSegmentedRequest(len(payload)) {
+                               if 
!m.driverContext.peerAcceptsSegmentedRequests {
+                                       utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil,
+                                               errors.Errorf("write request of 
%d bytes exceeds the peer's max APDU of %d and the peer does not support 
segmented requests", len(payload), m.driverContext.peerMaxApduBytes)))
+                                       return
+                               }
+                               segmentedPayload = payload
+                       }
+               }
+
                transaction := m.tm.StartTransaction("write")
+               if segmentedPayload != nil {
+                       transaction.Submit("segmentedWriteOperation", 
func(transactionContext context.Context, transaction 
transactions.RequestTransaction) {
+                               ctx, cancel := context.WithCancel(ctx)
+                               context.AfterFunc(transactionContext, cancel)
+
+                               // Register the response expectation BEFORE 
driving the segments:
+                               // the peer may answer right after acking the 
final segment. The
+                               // matcher must not consume the peer's 
SegmentAcks — those belong
+                               // to the sender's own expectations.
+                               m.messageCodec.Expect(ctx, 
"segmentedWriteResponse",
+                                       
responseMatcherExcludingSegmentAcks(func(message spi.Message) bool {
+                                               return 
m.acceptsResponse(message, invokeId)
+                                       }),
+                                       func(message spi.Message) error {
+                                               bvlc := 
message.(readWriteModel.BVLC)
+                                               responseApdu := 
bvlc.(interface{ GetNpdu() readWriteModel.NPDU }).GetNpdu().GetApdu()
+                                               writeResponse := 
m.toPlcWriteResponse(responseApdu, writeRequest)
+                                               utils.DeliverResult(m.log, 
result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, writeResponse, 
nil))
+                                               return transaction.EndRequest()
+                                       },
+                                       func(err error) error {
+                                               utils.DeliverResult(m.log, 
result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, 
errors.Wrap(err, "got timeout while waiting for segmented write response")))
+                                               return transaction.EndRequest()
+                                       },
+                               )
+
+                               sender := &segmentedRequestSender{
+                                       messageCodec:  m.messageCodec,
+                                       routedDest:    m.routedDest,
+                                       driverContext: m.driverContext,
+                                       log:           m.log,
+                               }
+                               if sendErr := sender.send(ctx, invokeId, 
segmentedPayload); sendErr != nil {
+                                       // Deliver the real cause first (the 
buffered result channel
+                                       // takes it), then cancel so the armed 
response expectation
+                                       // resolves without delivering a 
second, less specific error.
+                                       utils.DeliverResult(m.log, result, 
spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, 
errors.Wrap(sendErr, "error sending segmented write request")))
+                                       cancel()
+                               }
+                       })
+                       return
+               }
                transaction.Submit("writeOperation", func(transactionContext 
context.Context, transaction transactions.RequestTransaction) {
                        ctx, cancel := context.WithCancel(ctx)
                        context.AfterFunc(transactionContext, cancel)


Reply via email to