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


The following commit(s) were added to refs/heads/develop by this push:
     new 9205b257b7 feat(plc4go): implement BACnet/IP segmentation, write 
priority, and directed WhoIs
9205b257b7 is described below

commit 9205b257b72f8143798471801c6c40f2d7ad50b5
Author: Sebastian Rühl <[email protected]>
AuthorDate: Tue Jun 23 12:33:10 2026 +0200

    feat(plc4go): implement BACnet/IP segmentation, write priority, and 
directed WhoIs
    
    Three additive capabilities for the pure-Go BACnet/IP driver:
    
      - Reassemble segmented ComplexAck responses instead of returning 
UNSUPPORTED.
        A segmented read is handed to a goroutine (the codec holds its 
expectation
        lock during message handling, so the segment-ack/await loop can't run 
inline)
        that ACKs each segment, reassembles via inboundReassembler, and 
reparses the
        concatenated payload with BACnetServiceAckParse. decodeComplexAck is 
split so
        the reassembled service ack reuses the same decode path.
      - Support a write priority via a {N} tag suffix (1..16), e.g.
        ANALOG_OUTPUT,1/PRESENT_VALUE[2]{8}, encoded into the WriteProperty 
(context
        tag 4) and WritePropertyMultiple PropertyWriteDefinition (context tag 3)
        priority fields.
      - Add a remote-address protocol-specific discovery option that sends a 
directed
        unicast WhoIs to a specific host/subnet instead of only broadcasting.
    
    Adds unit tests for each (segment round-trip, priority parsing/range, 
address resolution and option parsing).
---
 plc4go/internal/bacnetip/Discoverer.go         |  60 ++++++++-
 plc4go/internal/bacnetip/Discoverer_test.go    |  31 +++++
 plc4go/internal/bacnetip/Reader.go             |  44 +++++--
 plc4go/internal/bacnetip/ReaderSegmentation.go | 174 +++++++++++++++++++++++++
 plc4go/internal/bacnetip/Reader_test.go        |  55 +++++++-
 plc4go/internal/bacnetip/Tag.go                |   6 +
 plc4go/internal/bacnetip/TagHandler.go         |  16 ++-
 plc4go/internal/bacnetip/TagHandler_test.go    |  29 +++++
 plc4go/internal/bacnetip/Writer.go             |  12 +-
 9 files changed, 405 insertions(+), 22 deletions(-)

diff --git a/plc4go/internal/bacnetip/Discoverer.go 
b/plc4go/internal/bacnetip/Discoverer.go
index 441fb4dcf2..f7926a1eb7 100644
--- a/plc4go/internal/bacnetip/Discoverer.go
+++ b/plc4go/internal/bacnetip/Discoverer.go
@@ -136,8 +136,20 @@ func (d *Discoverer) broadcastAndDiscover(ctx 
context.Context, communicationChan
                        if err != nil {
                                return nil, err
                        }
-                       if _, err := 
communicationChannelInstance.broadcastConnection.WriteTo(theBytes, 
communicationChannelInstance.broadcastConnection.LocalAddr()); err != nil {
-                               d.log.Debug().Err(err).Msg("Error sending 
broadcast")
+                       // Directed (unicast) WhoIs when a remote address is 
supplied,
+                       // otherwise broadcast on the interface.
+                       sendConn := 
communicationChannelInstance.broadcastConnection
+                       target := 
communicationChannelInstance.broadcastConnection.LocalAddr()
+                       if specificOptions.remoteAddress != "" {
+                               if udpAddr, rerr := 
resolveBacnetUDPAddr(specificOptions.remoteAddress, 
specificOptions.bacNetPort); rerr == nil {
+                                       sendConn = 
communicationChannelInstance.unicastConnection
+                                       target = udpAddr
+                               } else {
+                                       
d.log.Warn().Err(rerr).Str("remoteAddress", 
specificOptions.remoteAddress).Msg("invalid remote-address; falling back to 
broadcast")
+                               }
+                       }
+                       if _, err := sendConn.WriteTo(theBytes, target); err != 
nil {
+                               d.log.Debug().Err(err).Msg("Error sending 
WhoIs")
                        }
                }
                if whoHasOptions := specificOptions.whoHasOptions; 
whoHasOptions != nil {
@@ -472,8 +484,13 @@ func extractInterfaces(discoveryOptions 
[]options.WithDiscoveryOption) ([]net.In
 }
 
 type protocolSpecificOptions struct {
-       bacNetPort   int
-       whoIsOptions *struct {
+       bacNetPort int
+       // remoteAddress, when set, sends the WhoIs as a directed unicast to 
this
+       // host (instead of the interface broadcast), enabling targeted 
discovery of
+       // a specific device or subnet. Host only or host:port; port defaults to
+       // bacNetPort.
+       remoteAddress string
+       whoIsOptions  *struct {
                limits *struct {
                        low  uint
                        high uint
@@ -501,6 +518,33 @@ func bacNetPort(port int) option {
        }
 }
 
+// resolveBacnetUDPAddr parses a "host" or "host:port" string into a UDP 
address,
+// defaulting the port to defaultPort when none is supplied.
+func resolveBacnetUDPAddr(addr string, defaultPort int) (*net.UDPAddr, error) {
+       host, portStr, err := net.SplitHostPort(addr)
+       if err != nil {
+               // No port present — treat the whole string as the host.
+               host = addr
+               portStr = strconv.Itoa(defaultPort)
+       }
+       ip := net.ParseIP(host)
+       if ip == nil {
+               return nil, errors.Errorf("invalid remote-address host %q", 
host)
+       }
+       port, err := strconv.Atoi(portStr)
+       if err != nil {
+               return nil, errors.Wrapf(err, "invalid remote-address port %q", 
portStr)
+       }
+       return &net.UDPAddr{IP: ip, Port: port}, nil
+}
+
+func remoteAddress(addr string) option {
+       return func(specificOptions *protocolSpecificOptions) error {
+               specificOptions.remoteAddress = addr
+               return nil
+       }
+}
+
 func whoIsLimits(whoIsLowLimit, whoIsHighLimit uint) option {
        return func(specificOptions *protocolSpecificOptions) error {
                specificOptions.whoIsOptions = &struct {
@@ -639,6 +683,14 @@ func extractProtocolSpecificOptions(discoveryOptions 
[]options.WithDiscoveryOpti
                collectedOptions = append(collectedOptions, bacNetPort(47808))
        }
 
+       if _, ok := filteredOptionMap["remote-address"]; ok {
+               addr, err := OneString(filteredOptionMap, "remote-address")
+               if err != nil {
+                       return nil, err
+               }
+               collectedOptions = append(collectedOptions, remoteAddress(addr))
+       }
+
        if whoIsLow, whoIsHigh, ok, err := func() (whoIsLowLimit uint, 
whoIsHighLimit uint, ok bool, err error) {
                if _, limitPresent := filteredOptionMap["who-is-low-limit"]; 
!limitPresent {
                        return
diff --git a/plc4go/internal/bacnetip/Discoverer_test.go 
b/plc4go/internal/bacnetip/Discoverer_test.go
index ed1895b39d..655f4ed97a 100644
--- a/plc4go/internal/bacnetip/Discoverer_test.go
+++ b/plc4go/internal/bacnetip/Discoverer_test.go
@@ -24,6 +24,9 @@ import (
        "time"
 
        "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+
+       "github.com/apache/plc4x/plc4go/spi/options"
 )
 
 func TestNewDiscoverer_DefaultTimeout(t *testing.T) {
@@ -48,3 +51,31 @@ func TestSetDiscoveryTimeout_NegativeFallsBackToDefault(t 
*testing.T) {
        d.SetDiscoveryTimeout(-1 * time.Second)
        assert.Equal(t, 5*time.Second, d.discoveryTimeout)
 }
+
+func TestResolveBacnetUDPAddr(t *testing.T) {
+       // Host only — default port applies.
+       addr, err := resolveBacnetUDPAddr("192.168.1.50", 47808)
+       require.NoError(t, err)
+       assert.Equal(t, "192.168.1.50", addr.IP.String())
+       assert.Equal(t, 47808, addr.Port)
+
+       // Host:port — explicit port wins.
+       addr, err = resolveBacnetUDPAddr("10.0.0.5:47809", 47808)
+       require.NoError(t, err)
+       assert.Equal(t, "10.0.0.5", addr.IP.String())
+       assert.Equal(t, 47809, addr.Port)
+
+       // Invalid host.
+       _, err = resolveBacnetUDPAddr("not-an-ip", 47808)
+       require.Error(t, err)
+}
+
+func TestExtractProtocolSpecificOptions_RemoteAddress(t *testing.T) {
+       opts := []options.WithDiscoveryOption{
+               options.WithDiscoveryOptionProtocolSpecific("remote-address", 
"192.168.1.50"),
+               options.WithDiscoveryOptionProtocolSpecific("bacnet-port", 
47808),
+       }
+       specific, err := extractProtocolSpecificOptions(opts)
+       require.NoError(t, err)
+       assert.Equal(t, "192.168.1.50", specific.remoteAddress)
+}
diff --git a/plc4go/internal/bacnetip/Reader.go 
b/plc4go/internal/bacnetip/Reader.go
index 32bf49537b..5ff38b1952 100644
--- a/plc4go/internal/bacnetip/Reader.go
+++ b/plc4go/internal/bacnetip/Reader.go
@@ -135,6 +135,10 @@ func (m *Reader) Read(ctx context.Context, readRequest 
apiModel.PlcReadRequest)
                // 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) {
+                       // readCtx is the caller's context; it outlives the 
transaction so a
+                       // multi-segment reassembly (driven after we end the 
transaction) can
+                       // keep awaiting segments without being cancelled by 
transaction teardown.
+                       readCtx := ctx
                        ctx, cancel := context.WithCancel(ctx)
                        context.AfterFunc(transactionContext, cancel)
                        // Send the  over the wire
@@ -166,7 +170,18 @@ func (m *Reader) Read(ctx context.Context, readRequest 
apiModel.PlcReadRequest)
                                m.log.Trace().Msg("convert response to ")
                                apdu := 
message.(readWriteModel.BVLC).(interface{ GetNpdu() readWriteModel.NPDU 
}).GetNpdu().GetApdu()
 
-                               // TODO: implement segment handling
+                               // Segmented response: the device split the 
service ack across
+                               // multiple APDUs. We can't drive the 
segment-ack/await loop here
+                               // because this callback runs while the codec 
holds its expectation
+                               // lock — instead hand off to a goroutine that 
uses the caller's
+                               // context, and free this transaction slot 
immediately.
+                               if complexAck, ok := 
apdu.(readWriteModel.APDUComplexAck); ok && complexAck.GetSegmentedMessage() {
+                                       m.log.Trace().Uint8("invokeId", 
complexAck.GetOriginalInvokeId()).Msg("segmented read response — starting 
reassembly")
+                                       m.wg.Go(func() {
+                                               
m.reassembleSegmentedRead(readCtx, readRequest, complexAck, result)
+                                       })
+                                       return transaction.EndRequest()
+                               }
 
                                // Convert the bacnet response into a PLC4X 
response
                                m.log.Trace().Msg("convert response to PLC4X 
response")
@@ -242,22 +257,27 @@ func (m *Reader) ToPlc4xReadResponse(apdu 
readWriteModel.APDU, readRequest apiMo
        }
 }
 
-// decodeComplexAck handles the happy-path APDUComplexAck for ReadProperty and
-// ReadPropertyMultiple. Segmented responses are flagged but not yet 
reassembled
-// — that happens in Phase 5.
+// decodeComplexAck handles a non-segmented APDUComplexAck for ReadProperty and
+// ReadPropertyMultiple. Segmented responses are reassembled separately (see
+// reassembleSegmentedRead) and then handed to decodeServiceAck directly.
 func (m *Reader) decodeComplexAck(apdu readWriteModel.APDUComplexAck, 
readRequest apiModel.PlcReadRequest) (apiModel.PlcReadResponse, error) {
+       if apdu.GetSegmentedMessage() {
+               // Should not happen: segmented responses are intercepted 
before this
+               // point. Surface UNSUPPORTED rather than panicking on a nil 
ServiceAck.
+               m.log.Warn().Uint8("invokeId", 
apdu.GetOriginalInvokeId()).Msg("segmented APDU reached decodeComplexAck 
unexpectedly")
+               return broadcastResponseCode(readRequest, 
apiModel.PlcResponseCode_UNSUPPORTED, map[string]apiModel.PlcResponseCode{}, 
map[string]values.PlcValue{}), nil
+       }
+       return m.decodeServiceAck(apdu.GetServiceAck(), readRequest)
+}
+
+// decodeServiceAck converts a decoded BACnetServiceAck (whether it arrived in 
a
+// single APDU or was reassembled from segments) into a PLC4X read response.
+func (m *Reader) decodeServiceAck(serviceAckMsg 
readWriteModel.BACnetServiceAck, readRequest apiModel.PlcReadRequest) 
(apiModel.PlcReadResponse, error) {
        tagNames := readRequest.GetTagNames()
        responseCodes := map[string]apiModel.PlcResponseCode{}
        plcValues := map[string]values.PlcValue{}
 
-       if apdu.GetSegmentedMessage() {
-               // Segmentation reassembly is implemented in Phase 5. Until 
then surface a
-               // well-formed response with UNSUPPORTED so callers handle it 
gracefully.
-               m.log.Warn().Uint8("invokeId", 
apdu.GetOriginalInvokeId()).Msg("segmented APDU response — reassembly not yet 
implemented")
-               return broadcastResponseCode(readRequest, 
apiModel.PlcResponseCode_UNSUPPORTED, responseCodes, plcValues), nil
-       }
-
-       switch serviceAck := apdu.GetServiceAck().(type) {
+       switch serviceAck := serviceAckMsg.(type) {
        case readWriteModel.BACnetServiceAckReadProperty:
                if len(tagNames) == 0 {
                        return nil, errors.New("ReadProperty response without a 
corresponding requested tag")
diff --git a/plc4go/internal/bacnetip/ReaderSegmentation.go 
b/plc4go/internal/bacnetip/ReaderSegmentation.go
new file mode 100644
index 0000000000..63f6b29173
--- /dev/null
+++ b/plc4go/internal/bacnetip/ReaderSegmentation.go
@@ -0,0 +1,174 @@
+/*
+ * 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"
+
+       apiModel "github.com/apache/plc4x/plc4go/pkg/api/model"
+       readWriteModel 
"github.com/apache/plc4x/plc4go/protocols/bacnetip/readwrite/model"
+       "github.com/apache/plc4x/plc4go/spi"
+       "github.com/apache/plc4x/plc4go/spi/errors"
+       spiModel "github.com/apache/plc4x/plc4go/spi/model"
+)
+
+// segmentWaitTimeout bounds how long we wait for each follow-up segment before
+// giving up on a segmented response.
+const segmentWaitTimeout = 5 * time.Second
+
+// reassembleSegmentedRead drives the BACnet segmented-response protocol for a
+// read: it acknowledges each segment and waits for the next until the device
+// signals no more follow, then reparses the concatenated payload into a 
service
+// ack and produces the final PLC4X read result.
+//
+// It runs in its own goroutine (NOT under the codec expectation lock) so it is
+// free to register fresh expectations and send segment acks.
+func (m *Reader) reassembleSegmentedRead(ctx context.Context, readRequest 
apiModel.PlcReadRequest, first readWriteModel.APDUComplexAck, result chan 
apiModel.PlcReadRequestResult) {
+       invokeId := first.GetOriginalInvokeId()
+
+       // Advertise a window size of 1 (ack every segment) — simplest and most
+       // broadly compatible; the reassembler echoes this in each SegmentAck.
+       reassembler := NewInboundReassembler(invokeId, 1)
+
+       ack, err := reassembler.AcceptSegment(first)
+       if err != nil {
+               m.failSegmentedRead(readRequest, result, errors.Wrap(err, 
"error accepting first segment"))
+               return
+       }
+
+       for {
+               if reassembler.Complete() {
+                       // Acknowledge the final segment (best-effort).
+                       if sendErr := m.sendSegmentAck(ctx, ack); sendErr != 
nil {
+                               m.log.Debug().Err(sendErr).Msg("error sending 
final segment ack")
+                       }
+                       break
+               }
+
+               // Register the expectation for the next segment BEFORE acking 
the
+               // current one, otherwise the device's next segment could race 
ahead of
+               // our expectation and be dropped.
+               segCh, errCh := m.expectSegment(ctx, invokeId)
+               if sendErr := m.sendSegmentAck(ctx, ack); sendErr != nil {
+                       m.failSegmentedRead(readRequest, result, 
errors.Wrap(sendErr, "error sending segment ack"))
+                       return
+               }
+
+               select {
+               case <-ctx.Done():
+                       m.failSegmentedRead(readRequest, result, 
errors.Wrap(ctx.Err(), "context cancelled awaiting segment"))
+                       return
+               case waitErr := <-errCh:
+                       m.failSegmentedRead(readRequest, result, 
errors.Wrap(waitErr, "error awaiting segment"))
+                       return
+               case seg := <-segCh:
+                       ack, err = reassembler.AcceptSegment(seg)
+                       if err != nil {
+                               m.failSegmentedRead(readRequest, result, 
errors.Wrap(err, "error accepting segment"))
+                               return
+                       }
+               }
+       }
+
+       // Reparse the concatenated payload (service-choice byte + data) into a
+       // service ack, then decode it the same way as a single-APDU response.
+       payload := reassembler.Bytes()
+       serviceAck, err := 
readWriteModel.BACnetServiceAckParse[readWriteModel.BACnetServiceAck](ctx, 
payload, uint32(len(payload)))
+       if err != nil {
+               m.failSegmentedRead(readRequest, result, errors.Wrap(err, 
"error parsing reassembled service ack"))
+               return
+       }
+
+       readResponse, err := m.decodeServiceAck(serviceAck, readRequest)
+       if err != nil {
+               m.failSegmentedRead(readRequest, result, errors.Wrap(err, 
"error decoding reassembled response"))
+               return
+       }
+       result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, 
readResponse, nil)
+}
+
+func (m *Reader) failSegmentedRead(readRequest apiModel.PlcReadRequest, result 
chan apiModel.PlcReadRequestResult, err error) {
+       m.log.Debug().Err(err).Msg("segmented read failed")
+       result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, err)
+}
+
+// expectSegment registers a one-shot expectation for the next segmented
+// APDUComplexAck with the given invoke id, returning channels for the segment 
or
+// an error/timeout.
+func (m *Reader) expectSegment(ctx context.Context, invokeId uint8) (<-chan 
readWriteModel.APDUComplexAck, <-chan error) {
+       segCh := make(chan readWriteModel.APDUComplexAck, 1)
+       errCh := make(chan error, 1)
+
+       expectCtx, cancel := context.WithTimeout(ctx, segmentWaitTimeout)
+
+       m.messageCodec.Expect(expectCtx, "readSegment",
+               func(message spi.Message) bool {
+                       apdu, ok := apduFromMessage(message)
+                       if !ok {
+                               return false
+                       }
+                       complexAck, ok := apdu.(readWriteModel.APDUComplexAck)
+                       return ok && complexAck.GetSegmentedMessage() && 
complexAck.GetOriginalInvokeId() == invokeId
+               },
+               func(message spi.Message) error {
+                       cancel()
+                       apdu, _ := apduFromMessage(message)
+                       segCh <- apdu.(readWriteModel.APDUComplexAck)
+                       return nil
+               },
+               func(err error) error {
+                       cancel()
+                       errCh <- err
+                       return nil
+               },
+       )
+       return segCh, errCh
+}
+
+// sendSegmentAck transmits a BACnet SegmentAck wrapped in NPDU+BVLC.
+func (m *Reader) sendSegmentAck(ctx context.Context, ack 
readWriteModel.APDUSegmentAck) error {
+       if ack == nil {
+               return nil
+       }
+       return m.messageCodec.Send(ctx, "segmentAck", wrapAPDU(ack, false))
+}
+
+// apduFromMessage safely extracts the APDU from a received BVLC message,
+// returning ok=false rather than panicking on unexpected shapes.
+func apduFromMessage(message spi.Message) (readWriteModel.APDU, bool) {
+       bvlc, ok := message.(readWriteModel.BVLC)
+       if !ok {
+               return nil, false
+       }
+       npduRetriever, ok := bvlc.(interface{ GetNpdu() readWriteModel.NPDU })
+       if !ok {
+               return nil, false
+       }
+       npdu := npduRetriever.GetNpdu()
+       if npdu == nil || npdu.GetControl().GetMessageTypeFieldPresent() {
+               return nil, false
+       }
+       apdu := npdu.GetApdu()
+       if apdu == nil {
+               return nil, false
+       }
+       return apdu, true
+}
diff --git a/plc4go/internal/bacnetip/Reader_test.go 
b/plc4go/internal/bacnetip/Reader_test.go
index 2550284a06..caecedd8cc 100644
--- a/plc4go/internal/bacnetip/Reader_test.go
+++ b/plc4go/internal/bacnetip/Reader_test.go
@@ -20,6 +20,7 @@
 package bacnetip
 
 import (
+       "context"
        "testing"
 
        "github.com/rs/zerolog"
@@ -252,7 +253,7 @@ func TestToPlc4xReadResponse_Reject(t *testing.T) {
        assert.Equal(t, apiModel.PlcResponseCode_INVALID_DATA, 
resp.GetResponseCode("r"))
 }
 
-func TestToPlc4xReadResponse_Segmented_ReturnsUnsupported(t *testing.T) {
+func TestToPlc4xReadResponse_Segmented_DefensiveFallback(t *testing.T) {
        reader := newTestReader(t)
        request := readRequestForTags("big")
        serviceAck := readWriteModel.NewBACnetServiceAckReadProperty(
@@ -262,13 +263,63 @@ func 
TestToPlc4xReadResponse_Segmented_ReturnsUnsupported(t *testing.T) {
                nil,
                
constructedDataFromTag(readWriteModel.CreateBACnetApplicationTagNull()),
        )
-       // segmentedMessage=true triggers the Phase-2 short-circuit until Phase 
5 lands.
+       // In the live read flow a segmented APDU is intercepted and reassembled
+       // before ToPlc4xReadResponse. Reaching decodeComplexAck with a 
segmented
+       // APDU is unexpected, so it falls back to UNSUPPORTED rather than 
panicking.
        apdu := readWriteModel.NewAPDUComplexAck(true, true, 1, nil, nil, 
serviceAck, nil, nil)
        resp, err := reader.ToPlc4xReadResponse(apdu, request)
        require.NoError(t, err)
        assert.Equal(t, apiModel.PlcResponseCode_UNSUPPORTED, 
resp.GetResponseCode("big"))
 }
 
+// TestReassembledSegments_RoundTrip verifies the core segmentation path: a 
real
+// service ack serialized and split across two APDUComplexAck segments is
+// reassembled, reparsed, and decoded back into the original value.
+func TestReassembledSegments_RoundTrip(t *testing.T) {
+       reader := newTestReader(t)
+       request := readRequestForTags("big")
+
+       serviceAck := readWriteModel.NewBACnetServiceAckReadProperty(
+               0,
+               readWriteModel.CreateBACnetContextTagObjectIdentifier(0, 
uint16(readWriteModel.BACnetObjectType_ANALOG_INPUT), 1),
+               readWriteModel.CreateBACnetPropertyIdentifierTagged(1, 
uint32(readWriteModel.BACnetPropertyIdentifier_PRESENT_VALUE)),
+               nil,
+               
constructedDataFromTag(readWriteModel.CreateBACnetApplicationTagReal(23.5)),
+       )
+       fullBytes, err := serviceAck.Serialize()
+       require.NoError(t, err)
+       require.Greater(t, len(fullBytes), 2)
+
+       // Split the serialized service ack across two segments. Segment 0 
carries
+       // the service-choice byte (as the wire format does); concatenation 
restores
+       // the original bytes.
+       split := len(fullBytes) / 2
+       seg0 := readWriteModel.NewAPDUComplexAck(true, true, 7, ptrU8(0), 
ptrU8(1), nil, nil, fullBytes[:split])
+       seg1 := readWriteModel.NewAPDUComplexAck(true, false, 7, ptrU8(1), 
ptrU8(1), nil, nil, fullBytes[split:])
+
+       r := NewInboundReassembler(7, 1)
+       ack0, err := r.AcceptSegment(seg0)
+       require.NoError(t, err)
+       require.NotNil(t, ack0)
+       require.False(t, r.Complete())
+       _, err = r.AcceptSegment(seg1)
+       require.NoError(t, err)
+       require.True(t, r.Complete())
+       require.Equal(t, fullBytes, r.Bytes())
+
+       parsed, err := 
readWriteModel.BACnetServiceAckParse[readWriteModel.BACnetServiceAck](context.Background(),
 r.Bytes(), uint32(len(r.Bytes())))
+       require.NoError(t, err)
+
+       resp, err := reader.decodeServiceAck(parsed, request)
+       require.NoError(t, err)
+       assert.Equal(t, apiModel.PlcResponseCode_OK, 
resp.GetResponseCode("big"))
+       val := resp.GetValue("big")
+       require.NotNil(t, val)
+       assert.InDelta(t, 23.5, val.GetFloat32(), 0.001)
+}
+
+func ptrU8(v uint8) *uint8 { return &v }
+
 // ── ValueDecoder unit tests ───────────────────────────────────────────────
 
 func TestAppTagToPlcValue_AllPrimitiveTypes(t *testing.T) {
diff --git a/plc4go/internal/bacnetip/Tag.go b/plc4go/internal/bacnetip/Tag.go
index 24ec6e58c3..4f6da0cf0c 100644
--- a/plc4go/internal/bacnetip/Tag.go
+++ b/plc4go/internal/bacnetip/Tag.go
@@ -95,6 +95,9 @@ type property struct {
        PropertyIdentifierProprietary *uint32
        // ArrayIndex Optional index of property
        ArrayIndex *uint
+       // WritePriority is the optional BACnet write priority (1..16) for 
commandable
+       // properties. Only meaningful on write requests; ignored on reads.
+       WritePriority *uint8
 }
 
 func (p property) getId() uint32 {
@@ -115,6 +118,9 @@ func (p property) String() string {
        if p.ArrayIndex != nil {
                result += fmt.Sprintf(":[%d]", p.ArrayIndex)
        }
+       if p.WritePriority != nil {
+               result += fmt.Sprintf(":{%d}", *p.WritePriority)
+       }
        return result
 }
 
diff --git a/plc4go/internal/bacnetip/TagHandler.go 
b/plc4go/internal/bacnetip/TagHandler.go
index 5608eb7e81..4a86af2503 100644
--- a/plc4go/internal/bacnetip/TagHandler.go
+++ b/plc4go/internal/bacnetip/TagHandler.go
@@ -38,8 +38,8 @@ type TagHandler struct {
 
 func NewTagHandler() TagHandler {
        return TagHandler{
-               addressPattern:       
regexp.MustCompile(`^(?P<objectType>[\d\w]+),(?P<objectInstance>\d+)/(?P<propertyIdentifiers>[\d\w]+(?:\[\d+])?(?:&[\d\w]+(?:\[\d+])?)*)`),
-               propertyFieldPattern: 
regexp.MustCompile(`^(?P<propertyIdentifier>[\d\w]+)(?:\[(?P<arrayIndex>\d+)])?$`),
+               addressPattern:       
regexp.MustCompile(`^(?P<objectType>[\d\w]+),(?P<objectInstance>\d+)/(?P<propertyIdentifiers>[\d\w]+(?:\[\d+])?(?:\{\d+})?(?:&[\d\w]+(?:\[\d+])?(?:\{\d+})?)*)`),
+               propertyFieldPattern: 
regexp.MustCompile(`^(?P<propertyIdentifier>[\d\w]+)(?:\[(?P<arrayIndex>\d+)])?(?:\{(?P<writePriority>\d+)})?$`),
        }
 }
 
@@ -49,6 +49,7 @@ const (
        PROPERTY_IDENTIFIERS = "propertyIdentifiers"
        PROPERTY_IDENTIFIER  = "propertyIdentifier"
        ARRAY_INDEX          = "arrayIndex"
+       WRITE_PRIORITY       = "writePriority"
 )
 
 func (m TagHandler) ParseTag(tagString string) (apiModel.PlcTag, error) {
@@ -78,6 +79,7 @@ func (m TagHandler) ParseTag(tagString string) 
(apiModel.PlcTag, error) {
                                PropertyIdentifier            
*readWriteModel.BACnetPropertyIdentifier
                                PropertyIdentifierProprietary *uint32
                                ArrayIndex                    *uint
+                               WritePriority                 *uint8
                        }
                        propertyMatch := 
utils.GetSubgroupMatches(m.propertyFieldPattern, propertyString)
                        propertyIdentifierMatch := 
propertyMatch[PROPERTY_IDENTIFIER]
@@ -99,6 +101,16 @@ func (m TagHandler) ParseTag(tagString string) 
(apiModel.PlcTag, error) {
                                        _property.ArrayIndex = &arrayIndex
                                }
                        }
+                       if writePriorityMatch := propertyMatch[WRITE_PRIORITY]; 
writePriorityMatch != "" {
+                               if parsedPriority, err := 
strconv.ParseUint(writePriorityMatch, 10, 8); err != nil {
+                                       return nil, errors.Wrap(err, "Error 
parsing write priority")
+                               } else if parsedPriority < 1 || parsedPriority 
> 16 {
+                                       return nil, errors.Errorf("write 
priority %d out of range (1..16)", parsedPriority)
+                               } else {
+                                       writePriority := uint8(parsedPriority)
+                                       _property.WritePriority = &writePriority
+                               }
+                       }
 
                        result.Properties = append(result.Properties, _property)
                }
diff --git a/plc4go/internal/bacnetip/TagHandler_test.go 
b/plc4go/internal/bacnetip/TagHandler_test.go
index 16432d668a..8bf1effb17 100644
--- a/plc4go/internal/bacnetip/TagHandler_test.go
+++ b/plc4go/internal/bacnetip/TagHandler_test.go
@@ -69,6 +69,35 @@ func TestTagHandler_ParseWithArrayIndex(t *testing.T) {
        assert.Equal(t, uint(3), *props[0].ArrayIndex)
 }
 
+func TestTagHandler_ParseWithWritePriority(t *testing.T) {
+       h := NewTagHandler()
+       tag, err := h.ParseTag("ANALOG_OUTPUT,1/PRESENT_VALUE{8}")
+       require.NoError(t, err)
+       props := tag.(BacNetPlcTag).GetProperties()
+       require.Len(t, props, 1)
+       require.Nil(t, props[0].ArrayIndex)
+       require.NotNil(t, props[0].WritePriority)
+       assert.Equal(t, uint8(8), *props[0].WritePriority)
+}
+
+func TestTagHandler_ParseWithArrayIndexAndWritePriority(t *testing.T) {
+       h := NewTagHandler()
+       tag, err := h.ParseTag("ANALOG_OUTPUT,5/PRESENT_VALUE[2]{16}")
+       require.NoError(t, err)
+       props := tag.(BacNetPlcTag).GetProperties()
+       require.Len(t, props, 1)
+       require.NotNil(t, props[0].ArrayIndex)
+       assert.Equal(t, uint(2), *props[0].ArrayIndex)
+       require.NotNil(t, props[0].WritePriority)
+       assert.Equal(t, uint8(16), *props[0].WritePriority)
+}
+
+func TestTagHandler_ParseWritePriorityOutOfRange(t *testing.T) {
+       h := NewTagHandler()
+       _, err := h.ParseTag("ANALOG_OUTPUT,1/PRESENT_VALUE{17}")
+       require.Error(t, err)
+}
+
 func TestTagHandler_ParseMultipleProperties(t *testing.T) {
        h := NewTagHandler()
        tag, err := h.ParseTag("ANALOG_INPUT,1/PRESENT_VALUE&OBJECT_NAME&UNITS")
diff --git a/plc4go/internal/bacnetip/Writer.go 
b/plc4go/internal/bacnetip/Writer.go
index 98c5cb701f..b0080e7ecb 100644
--- a/plc4go/internal/bacnetip/Writer.go
+++ b/plc4go/internal/bacnetip/Writer.go
@@ -177,7 +177,11 @@ func (m *Writer) buildServiceRequest(writeRequest 
apiModel.PlcWriteRequest) (rea
                        if prop.ArrayIndex != nil {
                                arrayIndex = 
readWriteModel.CreateBACnetContextTagUnsignedInteger(1, *prop.ArrayIndex)
                        }
-                       defs = append(defs, 
readWriteModel.NewBACnetPropertyWriteDefinition(propId, arrayIndex, cd, nil))
+                       var priority 
readWriteModel.BACnetContextTagUnsignedInteger
+                       if prop.WritePriority != nil {
+                               priority = 
readWriteModel.CreateBACnetContextTagUnsignedInteger(3, 
uint(*prop.WritePriority))
+                       }
+                       defs = append(defs, 
readWriteModel.NewBACnetPropertyWriteDefinition(propId, arrayIndex, cd, 
priority))
                }
                specs = append(specs, 
readWriteModel.NewBACnetWriteAccessSpecification(
                        objectIdTag,
@@ -204,8 +208,12 @@ func (m *Writer) buildSingleWriteProperty(tag 
BacNetPlcTag, plcValue apiValues.P
        if prop.ArrayIndex != nil {
                arrayIndex = 
readWriteModel.CreateBACnetContextTagUnsignedInteger(2, *prop.ArrayIndex)
        }
+       var priority readWriteModel.BACnetContextTagUnsignedInteger
+       if prop.WritePriority != nil {
+               priority = 
readWriteModel.CreateBACnetContextTagUnsignedInteger(4, 
uint(*prop.WritePriority))
+       }
        cd := constructedDataFromAppTag(appTag, 3)
-       return readWriteModel.NewBACnetConfirmedServiceRequestWriteProperty(0, 
objectIdTag, propId, arrayIndex, cd, nil), nil
+       return readWriteModel.NewBACnetConfirmedServiceRequestWriteProperty(0, 
objectIdTag, propId, arrayIndex, cd, priority), nil
 }
 
 // constructedDataFromAppTag wraps a single ApplicationTag into a generic

Reply via email to