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)
