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