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 369c946d018b1ae7165e0474342bbf03da18b124 Author: Sebastian Rühl <[email protected]> AuthorDate: Wed Jul 15 09:50:24 2026 +0200 fix(plc4go): adopt DeliverResult in the remaining drivers Follow-up to 18ce26d7f9: converts the API-facing request/result channel sends in ads, cbus, eip, knxnetip, opcua and s7 (readers, writers, subscribers, ads browser) to utils.DeliverResult. All these channels have capacity 1 and a consumer that reads at most one result, so a surplus send - late timeout handler, duplicate error-handler invocation, abandoned caller - blocked its goroutine forever and, via the codec WaitGroup, wedged Disconnect. Notes: - knxnetip's ConnectionDriverSpecificOperations already delivered non-blocking through its sendResponse select/default helper and is left untouched. - The knxnetip Writer had no logger; it now takes the usual options and gets the connection's logger for the drop warning. - s7 Writer struct-literal results were normalized to the NewDefaultPlcWriteRequestResult constructor. - Connection-internal handshake channels (eip/s7/cbus Connection, knxnetip ConnectionInternalOperations, opcua SecureChannel and the deliberately oversized opcua subscribe channels) are out of scope: their consumers are driver-internal with bounded waits. --- plc4go/internal/ads/Browser.go | 5 ++-- plc4go/internal/ads/Reader.go | 42 +++++++++++++++++----------------- plc4go/internal/ads/Subscriber.go | 14 ++++++------ plc4go/internal/ads/Writer.go | 34 +++++++++++++-------------- plc4go/internal/cbus/Reader.go | 17 +++++++------- plc4go/internal/cbus/Subscriber.go | 9 ++++---- plc4go/internal/cbus/Writer.go | 19 +++++++-------- plc4go/internal/eip/Reader.go | 20 ++++++++-------- plc4go/internal/eip/Writer.go | 34 +++++++++++++-------------- plc4go/internal/knxnetip/Connection.go | 2 +- plc4go/internal/knxnetip/Reader.go | 6 ++--- plc4go/internal/knxnetip/Subscriber.go | 8 +++---- plc4go/internal/knxnetip/Writer.go | 14 +++++++++--- plc4go/internal/opcua/Reader.go | 14 ++++++------ plc4go/internal/opcua/Subscriber.go | 10 ++++---- plc4go/internal/opcua/Writer.go | 16 ++++++------- plc4go/internal/s7/Reader.go | 23 ++++++++++--------- plc4go/internal/s7/Writer.go | 24 +++++++------------ 18 files changed, 158 insertions(+), 153 deletions(-) diff --git a/plc4go/internal/ads/Browser.go b/plc4go/internal/ads/Browser.go index 3a5da78365..4ca3986791 100644 --- a/plc4go/internal/ads/Browser.go +++ b/plc4go/internal/ads/Browser.go @@ -29,6 +29,7 @@ import ( driverModel "github.com/apache/plc4x/plc4go/protocols/ads/readwrite/model" "github.com/apache/plc4x/plc4go/spi/errors" spiModel "github.com/apache/plc4x/plc4go/spi/model" + "github.com/apache/plc4x/plc4go/spi/utils" ) func (m *Connection) BrowseRequestBuilder() apiModel.PlcBrowseRequestBuilder { @@ -46,7 +47,7 @@ func (m *Connection) BrowseWithInterceptor(ctx context.Context, browseRequest ap m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcBrowseRequestResult(browseRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcBrowseRequestResult(browseRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() responseCodes := map[string]apiModel.PlcResponseCode{} @@ -56,7 +57,7 @@ func (m *Connection) BrowseWithInterceptor(ctx context.Context, browseRequest ap responseCodes[queryName], results[queryName] = m.BrowseQuery(ctx, interceptor, queryName, query) } browseResponse := spiModel.NewDefaultPlcBrowseResponse(browseRequest, results, responseCodes) - result <- spiModel.NewDefaultPlcBrowseRequestResult(browseRequest, browseResponse, nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcBrowseRequestResult(browseRequest, browseResponse, nil)) }) return result } diff --git a/plc4go/internal/ads/Reader.go b/plc4go/internal/ads/Reader.go index b0b2917f55..75d69174bb 100644 --- a/plc4go/internal/ads/Reader.go +++ b/plc4go/internal/ads/Reader.go @@ -46,7 +46,7 @@ func (m *Connection) Read(ctx context.Context, readRequest apiModel.PlcReadReque m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() if len(readRequest.GetTagNames()) <= 1 { @@ -60,7 +60,7 @@ func (m *Connection) Read(ctx context.Context, readRequest apiModel.PlcReadReque func (m *Connection) singleRead(ctx context.Context, readRequest apiModel.PlcReadRequest, result chan apiModel.PlcReadRequestResult) { if len(readRequest.GetTagNames()) != 1 { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.New("this part of the ads driver only supports single-item requests")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.New("this part of the ads driver only supports single-item requests"))) m.log.Debug().Int("nTags", len(readRequest.GetTagNames())).Msg("this part of the ads driver only supports single-item requests. Got nTags tags") return } @@ -71,25 +71,25 @@ func (m *Connection) singleRead(ctx context.Context, readRequest apiModel.PlcRea if model.NeedsResolving(tag) { adsField, err := model.CastToSymbolicPlcTagFromPlcTag(tag) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrap(err, "invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrap(err, "invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } // Replace the symbolic tag with a direct one tag, err = m.resolveSymbolicTag(ctx, adsField) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "invalid tag item type"), - ) + )) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } } directAdsTag, ok := tag.(*model.DirectPlcTag) if !ok { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.New("invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.New("invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } @@ -97,16 +97,16 @@ func (m *Connection) singleRead(ctx context.Context, readRequest apiModel.PlcRea m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() response, err := m.ExecuteAdsReadRequest(ctx, directAdsTag.IndexGroup, directAdsTag.IndexOffset, directAdsTag.DataType.GetSize()) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "got error executing the read request"), - ) + )) return } @@ -130,11 +130,11 @@ func (m *Connection) singleRead(ctx context.Context, readRequest apiModel.PlcRea } } // Return the response to the caller. - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, spiModel.NewDefaultPlcReadResponse(readRequest, responseCodes, plcValues), nil, - ) + )) }) } @@ -149,33 +149,33 @@ func (m *Connection) multiRead(ctx context.Context, readRequest apiModel.PlcRead if model.NeedsResolving(tag) { adsField, err := model.CastToSymbolicPlcTagFromPlcTag(tag) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "invalid tag item type"), - ) + )) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } // Replace the symbolic tag with a direct one tag, err = m.resolveSymbolicTag(ctx, adsField) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "invalid tag item type"), - ) + )) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } } directAdsTag, ok := tag.(*model.DirectPlcTag) if !ok { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.New("invalid tag item type"), - ) + )) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } @@ -202,11 +202,11 @@ func (m *Connection) multiRead(ctx context.Context, readRequest apiModel.PlcRead response, err := m.ExecuteAdsReadWriteRequest(ctx, uint32(driverModel.ReservedIndexGroups_ADSIGRP_MULTIPLE_READ), uint32(len(directAdsTags)), expectedResponseDataSize, requestItems, nil) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "error executing multi-item read request"), - ) + )) return } @@ -247,11 +247,11 @@ func (m *Connection) multiRead(ctx context.Context, readRequest apiModel.PlcRead } // Return the response to the caller. - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, spiModel.NewDefaultPlcReadResponse(readRequest, responseCodes, plcValues), nil, - ) + )) } func (m *Connection) parsePlcValue(dataType driverModel.AdsDataTypeTableEntry, arrayInfo []driverModel.AdsDataTypeArrayInfo, rb utils.ReadBufferByteBased) (apiValues.PlcValue, error) { diff --git a/plc4go/internal/ads/Subscriber.go b/plc4go/internal/ads/Subscriber.go index 0c77f919d2..0d7d9dc807 100644 --- a/plc4go/internal/ads/Subscriber.go +++ b/plc4go/internal/ads/Subscriber.go @@ -112,7 +112,7 @@ func (m *Connection) Subscribe(ctx context.Context, subscriptionRequest apiModel for _, subResultChannel := range subResultChannels { select { case <-ctx.Done(): - globalResultChannel <- spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, ctx.Err()) + utils.DeliverResult(m.log, globalResultChannel, spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, ctx.Err())) return case subResult := <-subResultChannel: // These are all single value requests ... so it's safe to assume this shortcut. @@ -123,7 +123,7 @@ func (m *Connection) Subscribe(ctx context.Context, subscriptionRequest apiModel // As soon as all are done, process the results result := m.processSubscriptionResponses(ctx, subscriptionRequest, subResults) // Return the final result - globalResultChannel <- result + utils.DeliverResult(m.log, globalResultChannel, result) }) return globalResultChannel @@ -134,7 +134,7 @@ func (m *Connection) subscribe(ctx context.Context, subscriptionRequest apiModel m.wg.Go(func() { defer func() { if err := recover(); err != nil { - responseChan <- spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, responseChan, spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() // At this point we are sure to only have single item direct tag requests. @@ -146,11 +146,11 @@ func (m *Connection) subscribe(ctx context.Context, subscriptionRequest apiModel response, err := m.ExecuteAdsAddDeviceNotificationRequest(ctx, directTag.IndexGroup, directTag.IndexOffset, directTag.DataType.GetSize(), model.AdsTransMode_ON_CHANGE, 0, 0) if err != nil { - responseChan <- spiModel.NewDefaultPlcSubscriptionRequestResult( + utils.DeliverResult(m.log, responseChan, spiModel.NewDefaultPlcSubscriptionRequestResult( subscriptionRequest, nil, err, - ) + )) } // Create a new subscription handle. subscriptionHandle := dirverModel.NewAdsSubscriptionHandle( @@ -159,7 +159,7 @@ func (m *Connection) subscribe(ctx context.Context, subscriptionRequest apiModel directTag, append(m._options, options.WithCustomLogger(m.log))..., ) - responseChan <- spiModel.NewDefaultPlcSubscriptionRequestResult( + utils.DeliverResult(m.log, responseChan, spiModel.NewDefaultPlcSubscriptionRequestResult( subscriptionRequest, spiModel.NewDefaultPlcSubscriptionResponse( subscriptionRequest, @@ -168,7 +168,7 @@ func (m *Connection) subscribe(ctx context.Context, subscriptionRequest apiModel append(m._options, options.WithCustomLogger(m.log))..., ), nil, - ) + )) // Store it together with the returned ADS handle. m.subscriptions[response.GetNotificationHandle()] = subscriptionHandle }) diff --git a/plc4go/internal/ads/Writer.go b/plc4go/internal/ads/Writer.go index b1d25f16f9..0bc7b0873f 100644 --- a/plc4go/internal/ads/Writer.go +++ b/plc4go/internal/ads/Writer.go @@ -45,7 +45,7 @@ func (m *Connection) Write(ctx context.Context, writeRequest apiModel.PlcWriteRe m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() if len(writeRequest.GetTagNames()) <= 1 { @@ -59,7 +59,7 @@ func (m *Connection) Write(ctx context.Context, writeRequest apiModel.PlcWriteRe func (m *Connection) singleWrite(ctx context.Context, writeRequest apiModel.PlcWriteRequest, result chan apiModel.PlcWriteRequestResult) { if len(writeRequest.GetTagNames()) != 1 { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("this part of the ads driver only supports single-item requests")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("this part of the ads driver only supports single-item requests"))) m.log.Debug().Int("nTags", len(writeRequest.GetTagNames())).Msg("this part of the ads driver only supports single-item requests. Got nTags tags") return } @@ -70,25 +70,25 @@ func (m *Connection) singleWrite(ctx context.Context, writeRequest apiModel.PlcW if model.NeedsResolving(tag) { adsField, err := model.CastToSymbolicPlcTagFromPlcTag(tag) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrap(err, "invalid tag item type"), - ) + )) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } // Replace the symbolic tag with a direct one tag, err = m.resolveSymbolicTag(ctx, adsField) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } } directAdsTag, ok := tag.(*model.DirectPlcTag) if !ok { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } @@ -98,7 +98,7 @@ func (m *Connection) singleWrite(ctx context.Context, writeRequest apiModel.PlcW io := utils.NewWriteBufferByteBased(utils.WithByteOrderForByteBasedBuffer(binary.LittleEndian)) err := m.serializePlcValue(directAdsTag.DataType, directAdsTag.GetArrayInfo(), value, io) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error serializing plc value")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error serializing plc value"))) return } data := io.GetBytes() @@ -106,12 +106,12 @@ func (m *Connection) singleWrite(ctx context.Context, writeRequest apiModel.PlcW m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() response, err := m.ExecuteAdsWriteRequest(ctx, directAdsTag.IndexGroup, directAdsTag.IndexOffset, data) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "got error executing the write request")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "got error executing the write request"))) return } @@ -127,7 +127,7 @@ func (m *Connection) singleWrite(ctx context.Context, writeRequest apiModel.PlcW responseCodes[tagName] = apiModel.PlcResponseCode_OK } // Return the response to the caller. - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes), nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes), nil)) }) } @@ -143,21 +143,21 @@ func (m *Connection) multiWrite(ctx context.Context, writeRequest apiModel.PlcWr if model.NeedsResolving(tag) { adsField, err := model.CastToSymbolicPlcTagFromPlcTag(tag) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } // Replace the symbolic tag with a direct one tag, err = m.resolveSymbolicTag(ctx, adsField) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } } directAdsTag, ok := tag.(*model.DirectPlcTag) if !ok { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("invalid tag item type"))) m.log.Debug().Type("tag", tag).Msg("Invalid tag item type") return } @@ -167,7 +167,7 @@ func (m *Connection) multiWrite(ctx context.Context, writeRequest apiModel.PlcWr // Serialize the plc value err := m.serializePlcValue(directAdsTag.DataType, directAdsTag.GetArrayInfo(), writeRequest.GetValue(tagName), io) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error serializing plc value")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error serializing plc value"))) return } @@ -193,12 +193,12 @@ func (m *Connection) multiWrite(ctx context.Context, writeRequest apiModel.PlcWr uint32(driverModel.ReservedIndexGroups_ADSIGRP_MULTIPLE_WRITE), uint32(len(directAdsTags)), expectedResponseDataSize, requestItems, io.GetBytes()) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error executing multi-item write request")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error executing multi-item write request"))) return } if response.GetResult() != driverModel.ReturnCode_OK { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, fmt.Errorf("got return result %s from remote", response.GetResult().String())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, fmt.Errorf("got return result %s from remote", response.GetResult().String()))) return } @@ -219,7 +219,7 @@ func (m *Connection) multiWrite(ctx context.Context, writeRequest apiModel.PlcWr } // Return the response to the caller. - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes), nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes), nil)) } func (m *Connection) serializePlcValue(dataType driverModel.AdsDataTypeTableEntry, arrayInfo []apiModel.ArrayInfo, plcValue apiValues.PlcValue, wb utils.WriteBufferByteBased) error { diff --git a/plc4go/internal/cbus/Reader.go b/plc4go/internal/cbus/Reader.go index 12f3e153e6..6e0aedbd47 100644 --- a/plc4go/internal/cbus/Reader.go +++ b/plc4go/internal/cbus/Reader.go @@ -35,6 +35,7 @@ import ( spiModel "github.com/apache/plc4x/plc4go/spi/model" "github.com/apache/plc4x/plc4go/spi/options" "github.com/apache/plc4x/plc4go/spi/transactions" + "github.com/apache/plc4x/plc4go/spi/utils" ) type Reader struct { @@ -66,12 +67,12 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadRequest, result chan apiModel.PlcReadRequestResult) { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() numTags := len(readRequest.GetTagNames()) if numTags > 20 { // letters g-z - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.New("Only 20 tags can be handled at once")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.New("Only 20 tags can be handled at once"))) return } messages := make(map[string]readWriteModel.CBusMessage) @@ -80,11 +81,11 @@ func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadReque message, supportsRead, _, _, err := TagToCBusMessage(tag, nil, m.alphaGenerator, m.messageCodec) switch { case err != nil: - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrapf(err, "Error encoding cbus message for tag %s", tagName), - ) + )) return case !supportsRead: // Note this should not be reachable panic("this should not be possible as we always should then get the error above") @@ -107,21 +108,21 @@ func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadReque } for tagName, messageToSend := range messages { if err := ctx.Err(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, err, - ) + )) return } m.createMessageTransactionAndWait(ctx, messageToSend, addResponseCode, tagName, addPlcValue) } readResponse := spiModel.NewDefaultPlcReadResponse(readRequest, responseCodes, plcValues) - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, readResponse, nil, - ) + )) } func (m *Reader) createMessageTransactionAndWait(ctx context.Context, messageToSend readWriteModel.CBusMessage, addResponseCode func(name string, responseCode apiModel.PlcResponseCode), tagName string, addPlcValue func(name string, plcValue apiValues.PlcValue)) { diff --git a/plc4go/internal/cbus/Subscriber.go b/plc4go/internal/cbus/Subscriber.go index 7cbc43ffe1..857fe91c15 100644 --- a/plc4go/internal/cbus/Subscriber.go +++ b/plc4go/internal/cbus/Subscriber.go @@ -35,6 +35,7 @@ import ( "github.com/apache/plc4x/plc4go/spi/errors" spiModel "github.com/apache/plc4x/plc4go/spi/model" "github.com/apache/plc4x/plc4go/spi/options" + "github.com/apache/plc4x/plc4go/spi/utils" spiValues "github.com/apache/plc4x/plc4go/spi/values" ) @@ -67,7 +68,7 @@ func (s *Subscriber) Subscribe(_ context.Context, subscriptionRequest apiModel.P s.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() internalPlcSubscriptionRequest := subscriptionRequest.(*spiModel.DefaultPlcSubscriptionRequest) @@ -94,7 +95,7 @@ func (s *Subscriber) Subscribe(_ context.Context, subscriptionRequest apiModel.P subscriptionValues[tagName] = handle } - result <- spiModel.NewDefaultPlcSubscriptionRequestResult( + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult( subscriptionRequest, spiModel.NewDefaultPlcSubscriptionResponse( subscriptionRequest, @@ -103,7 +104,7 @@ func (s *Subscriber) Subscribe(_ context.Context, subscriptionRequest apiModel.P append(s._options, options.WithCustomLogger(s.log))..., ), nil, - ) + )) }) return result } @@ -111,7 +112,7 @@ func (s *Subscriber) Subscribe(_ context.Context, subscriptionRequest apiModel.P func (s *Subscriber) Unsubscribe(ctx context.Context, unsubscriptionRequest apiModel.PlcUnsubscriptionRequest) <-chan apiModel.PlcUnsubscriptionRequestResult { // TODO: handle context result := make(chan apiModel.PlcUnsubscriptionRequestResult, 1) - result <- spiModel.NewDefaultPlcUnsubscriptionRequestResult(unsubscriptionRequest, nil, errors.New("Not Implemented")) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcUnsubscriptionRequestResult(unsubscriptionRequest, nil, errors.New("Not Implemented"))) // TODO: As soon as we establish a connection, we start getting data... // subscriptions are more a internal handling of which values to pass where. diff --git a/plc4go/internal/cbus/Writer.go b/plc4go/internal/cbus/Writer.go index b3cb504c5c..d0dafc37cb 100644 --- a/plc4go/internal/cbus/Writer.go +++ b/plc4go/internal/cbus/Writer.go @@ -33,6 +33,7 @@ import ( spiModel "github.com/apache/plc4x/plc4go/spi/model" "github.com/apache/plc4x/plc4go/spi/options" "github.com/apache/plc4x/plc4go/spi/transactions" + "github.com/apache/plc4x/plc4go/spi/utils" ) type Writer struct { @@ -62,16 +63,16 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() numTags := len(writeRequest.GetTagNames()) if numTags > 20 { // letters g-z - result <- spiModel.NewDefaultPlcWriteRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.New("Only 20 tags can be handled at once"), - ) + )) return } @@ -81,19 +82,19 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques plcValue := writeRequest.GetValue(tagName) message, _, supportsWrite, _, err := TagToCBusMessage(tag, plcValue, m.alphaGenerator, m.messageCodec) if !supportsWrite { - result <- spiModel.NewDefaultPlcWriteRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrapf(err, "Error encoding cbus message for tag %s. Tag is not meant to be written.", tagName), - ) + )) return } if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrapf(err, "Error encoding cbus message for tag %s", tagName), - ) + )) return } messages[tagName] = message @@ -107,7 +108,7 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques } for tagName, messageToSend := range messages { if err := ctx.Err(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, err) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, err)) return } tagNameCopy := tagName @@ -157,7 +158,7 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques }) } readResponse := spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes) - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, readResponse, nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, readResponse, nil)) }) return result } diff --git a/plc4go/internal/eip/Reader.go b/plc4go/internal/eip/Reader.go index 94ae2d2662..da9efee899 100644 --- a/plc4go/internal/eip/Reader.go +++ b/plc4go/internal/eip/Reader.go @@ -73,7 +73,7 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() classSegment := readWriteModel.NewLogicalSegment(readWriteModel.NewClassID(0, 6)) @@ -87,7 +87,7 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) } ansi, err := toAnsi(tag) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "Error encoding eip ansi for tag %s", tagName)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "Error encoding eip ansi for tag %s", tagName))) return } requestItem := readWriteModel.NewCipUnconnectedRequest(classSegment, instanceSegment, @@ -128,32 +128,32 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) m.log.Trace().Msg("convert response to PLC4X response") readResponse, err := m.ToPlc4xReadResponse(unconnectedDataItem.GetService(), readRequest) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "Error decoding response"), - ) + )) return transaction.EndRequest() } - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, readResponse, nil, - ) + )) return transaction.EndRequest() }, func(err error) error { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "got timeout while waiting for response"), - ) + )) return transaction.EndRequest() }); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "error sending message"), - ) + )) if err := transaction.FailRequest(errors.Errorf("timeout after %s", time.Second*1)); err != nil { m.log.Debug().Err(err).Msg("Error failing request") } diff --git a/plc4go/internal/eip/Writer.go b/plc4go/internal/eip/Writer.go index 38b735efd0..e10f4cec0e 100644 --- a/plc4go/internal/eip/Writer.go +++ b/plc4go/internal/eip/Writer.go @@ -69,7 +69,7 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() items := make([]readWriteModel.CipService, len(writeRequest.GetTagNames())) @@ -95,12 +95,12 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques } data, err := encodeValue(value, eipTag.GetType(), elements) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding value for eipTag %s", tagName)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding value for eipTag %s", tagName))) return } ansi, err := toAnsi(tag) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding eip ansi for eipTag %s", tagName)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding eip ansi for eipTag %s", tagName))) return } items[i] = readWriteModel.NewCipWriteRequest(ansi, eipTag.GetType(), elements, data) @@ -157,25 +157,25 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques readResponse, err := m.ToPlc4xWriteResponse(cipWriteResponse, writeRequest) if err != nil { - result <- &spiModel.DefaultPlcWriteRequestResult{ + utils.DeliverResult(m.log, result, &spiModel.DefaultPlcWriteRequestResult{ Request: writeRequest, Err: errors.Wrap(err, "Error decoding response"), - } + }) return transaction.EndRequest() } - result <- &spiModel.DefaultPlcWriteRequestResult{ + utils.DeliverResult(m.log, result, &spiModel.DefaultPlcWriteRequestResult{ Request: writeRequest, Response: readResponse, - } + }) return transaction.EndRequest() }, func(err error) error { - result <- &spiModel.DefaultPlcWriteRequestResult{ + utils.DeliverResult(m.log, result, &spiModel.DefaultPlcWriteRequestResult{ Request: writeRequest, Err: errors.New("got timeout while waiting for response"), - } + }) return transaction.EndRequest() }, time.Second*1); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrap(err, "error sending message")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrap(err, "error sending message"))) if err := transaction.FailRequest(errors.Errorf("timeout after %s", time.Second*1)); err != nil { m.log.Debug().Err(err).Msg("Error failing request") } @@ -254,27 +254,27 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques readResponse, err := m.ToPlc4xWriteResponse(multipleServiceResponse, writeRequest) if err != nil { - result <- &spiModel.DefaultPlcWriteRequestResult{ + utils.DeliverResult(m.log, result, &spiModel.DefaultPlcWriteRequestResult{ Request: writeRequest, Err: errors.Wrap(err, "Error decoding response"), - } + }) return transaction.EndRequest() } - result <- &spiModel.DefaultPlcWriteRequestResult{ + utils.DeliverResult(m.log, result, &spiModel.DefaultPlcWriteRequestResult{ Request: writeRequest, Response: readResponse, - } + }) return transaction.EndRequest() }, func(err error) error { - result <- &spiModel.DefaultPlcWriteRequestResult{ + utils.DeliverResult(m.log, result, &spiModel.DefaultPlcWriteRequestResult{ Request: writeRequest, Err: errors.New("got timeout while waiting for response"), - } + }) return transaction.EndRequest() }, time.Second*1); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrap(err, "error sending message")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult( writeRequest, nil, errors.Wrap(err, "error sending message"))) if err := transaction.FailRequest(errors.Errorf("timeout after %s", time.Second*1)); err != nil { m.log.Debug().Err(err).Msg("Error failing request") } diff --git a/plc4go/internal/knxnetip/Connection.go b/plc4go/internal/knxnetip/Connection.go index 4a26b3b491..a71470b70d 100644 --- a/plc4go/internal/knxnetip/Connection.go +++ b/plc4go/internal/knxnetip/Connection.go @@ -461,7 +461,7 @@ func (m *Connection) ReadRequestBuilder() apiModel.PlcReadRequestBuilder { func (m *Connection) WriteRequestBuilder() apiModel.PlcWriteRequestBuilder { return spiModel.NewDefaultPlcWriteRequestBuilder( - m.tagHandler, m.valueHandler, NewWriter(m.messageCodec)) + m.tagHandler, m.valueHandler, NewWriter(m.messageCodec, options.WithCustomLogger(m.log))) } func (m *Connection) SubscriptionRequestBuilder() apiModel.PlcSubscriptionRequestBuilder { diff --git a/plc4go/internal/knxnetip/Reader.go b/plc4go/internal/knxnetip/Reader.go index f7b7aa8b60..fdd585af20 100644 --- a/plc4go/internal/knxnetip/Reader.go +++ b/plc4go/internal/knxnetip/Reader.go @@ -60,7 +60,7 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) m.wg.Go(func() { defer func() { if err := recover(); err != nil { - resultChan <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, resultChan, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() responseCodes := map[string]apiModel.PlcResponseCode{} @@ -161,11 +161,11 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) // Assemble the results result := spiModel.NewDefaultPlcReadResponse(readRequest, responseCodes, plcValues) - resultChan <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, resultChan, spiModel.NewDefaultPlcReadRequestResult( readRequest, result, nil, - ) + )) }) return resultChan } diff --git a/plc4go/internal/knxnetip/Subscriber.go b/plc4go/internal/knxnetip/Subscriber.go index 24cbd761c1..2413b359e0 100644 --- a/plc4go/internal/knxnetip/Subscriber.go +++ b/plc4go/internal/knxnetip/Subscriber.go @@ -67,7 +67,7 @@ func (s *Subscriber) Subscribe(ctx context.Context, subscriptionRequest apiModel s.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() internalPlcSubscriptionRequest := subscriptionRequest.(*spiModel.DefaultPlcSubscriptionRequest) @@ -84,7 +84,7 @@ func (s *Subscriber) Subscribe(ctx context.Context, subscriptionRequest apiModel subscriptionValues[tagName] = NewSubscriptionHandle(s, tagName, internalPlcSubscriptionRequest.GetTag(tagName), tagType, internalPlcSubscriptionRequest.GetInterval(tagName)) } - result <- spiModel.NewDefaultPlcSubscriptionRequestResult( + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult( subscriptionRequest, spiModel.NewDefaultPlcSubscriptionResponse( subscriptionRequest, @@ -93,7 +93,7 @@ func (s *Subscriber) Subscribe(ctx context.Context, subscriptionRequest apiModel append(s._options, options.WithCustomLogger(s.log))..., ), nil, - ) + )) }) return result } @@ -101,7 +101,7 @@ func (s *Subscriber) Subscribe(ctx context.Context, subscriptionRequest apiModel func (s *Subscriber) Unsubscribe(ctx context.Context, unsubscriptionRequest apiModel.PlcUnsubscriptionRequest) <-chan apiModel.PlcUnsubscriptionRequestResult { // TODO: handle context result := make(chan apiModel.PlcUnsubscriptionRequestResult, 1) - result <- spiModel.NewDefaultPlcUnsubscriptionRequestResult(unsubscriptionRequest, nil, errors.New("Not Implemented")) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcUnsubscriptionRequestResult(unsubscriptionRequest, nil, errors.New("Not Implemented"))) // TODO: As soon as we establish a connection, we start getting data... // subscriptions are more an internal handling of which values to pass where. diff --git a/plc4go/internal/knxnetip/Writer.go b/plc4go/internal/knxnetip/Writer.go index 833554c336..99bec4f33e 100644 --- a/plc4go/internal/knxnetip/Writer.go +++ b/plc4go/internal/knxnetip/Writer.go @@ -22,20 +22,28 @@ package knxnetip import ( "context" + "github.com/rs/zerolog" + apiModel "github.com/apache/plc4x/plc4go/pkg/api/model" readWriteModel "github.com/apache/plc4x/plc4go/protocols/knxnetip/readwrite/model" "github.com/apache/plc4x/plc4go/spi" "github.com/apache/plc4x/plc4go/spi/errors" spiModel "github.com/apache/plc4x/plc4go/spi/model" + "github.com/apache/plc4x/plc4go/spi/options" + "github.com/apache/plc4x/plc4go/spi/utils" ) type Writer struct { messageCodec spi.MessageCodec + + log zerolog.Logger } -func NewWriter(messageCodec spi.MessageCodec) Writer { +func NewWriter(messageCodec spi.MessageCodec, _options ...options.WithOption) Writer { + customLogger := options.ExtractCustomLoggerOrDefaultToGlobal(_options...) return Writer{ messageCodec: messageCodec, + log: customLogger, } } @@ -50,7 +58,7 @@ func (m Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteRequest tag := writeRequest.GetTag(tagName) groupAddressTag, err := CastToGroupAddressTagFromPlcTag(tag) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("invalid tag item type")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("invalid tag item type"))) return result } @@ -59,7 +67,7 @@ func (m Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteRequest tagType := groupAddressTag.GetTagType() // TODO: why do we ignore the bytes here? if _, err := readWriteModel.KnxDatapointSerialize(value, *tagType); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("error serializing value: "+err.Error())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("error serializing value: "+err.Error()))) return result } } diff --git a/plc4go/internal/opcua/Reader.go b/plc4go/internal/opcua/Reader.go index f56f8ec9fc..7d640a63d3 100644 --- a/plc4go/internal/opcua/Reader.go +++ b/plc4go/internal/opcua/Reader.go @@ -59,7 +59,7 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadRequest, result chan apiModel.PlcReadRequestResult) { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() @@ -78,7 +78,7 @@ func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadReque nodeId, err := generateNodeId(tag) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "error generating node id from tag %s", tag)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "error generating node id from tag %s", tag))) return } @@ -106,19 +106,19 @@ func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadReque buffer := utils.NewWriteBufferByteBased(utils.WithByteOrderForByteBasedBuffer(binary.LittleEndian)) if err := extObject.SerializeWithWriteBuffer(ctx, buffer); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "Unable to serialise the ReadRequest")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "Unable to serialise the ReadRequest"))) return } consumer := func(opcuaResponse []byte) { reply, err := readWriteModel.ExtensionObjectParseWithBuffer[readWriteModel.ExtensionObject](ctx, utils.NewReadBufferByteBased(opcuaResponse, utils.WithByteOrderForReadBufferByteBased(binary.LittleEndian)), false) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "Unable to read the reply")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Wrapf(err, "Unable to read the reply"))) return } extensionObjectDefinition := reply.GetBody() if _readResponse, ok := extensionObjectDefinition.(readWriteModel.ReadResponse); ok { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, spiModel.NewDefaultPlcReadResponse(readResponse(m.log, readRequest, readRequest.GetTagNames(), _readResponse.GetResults())), nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, spiModel.NewDefaultPlcReadResponse(readResponse(m.log, readRequest, readRequest.GetTagNames(), _readResponse.GetResults())), nil)) return } else { if serviceFault, ok := extensionObjectDefinition.(readWriteModel.ServiceFault); ok { @@ -132,12 +132,12 @@ func (m *Reader) readSync(ctx context.Context, readRequest apiModel.PlcReadReque for _, tagName := range readRequest.GetTagNames() { responseCodes[tagName] = apiModel.PlcResponseCode_INTERNAL_ERROR } - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, spiModel.NewDefaultPlcReadResponse(readRequest, responseCodes, nil), nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, spiModel.NewDefaultPlcReadResponse(readRequest, responseCodes, nil), nil)) } } errorDispatcher := func(err error) { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, err) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, err)) } m.connection.channel.submit(ctx, m.connection.messageCodec, errorDispatcher, consumer, buffer) diff --git a/plc4go/internal/opcua/Subscriber.go b/plc4go/internal/opcua/Subscriber.go index 2385a4eefb..ada9054c21 100644 --- a/plc4go/internal/opcua/Subscriber.go +++ b/plc4go/internal/opcua/Subscriber.go @@ -71,7 +71,7 @@ func (s *Subscriber) Subscribe(ctx context.Context, subscriptionRequest apiModel func (s *Subscriber) subscribeSync(ctx context.Context, result chan apiModel.PlcSubscriptionRequestResult, subscriptionRequest apiModel.PlcSubscriptionRequest) { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() internalPlcSubscriptionRequest := subscriptionRequest.(*spiModel.DefaultPlcSubscriptionRequest) @@ -83,7 +83,7 @@ func (s *Subscriber) subscribeSync(ctx context.Context, result chan apiModel.Plc subscription, err := s.onSubscribeCreateSubscription(ctx, cycleTime) if err != nil { - result <- spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Wrap(err, "error create subscription")) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult(subscriptionRequest, nil, errors.Wrap(err, "error create subscription"))) return } subscriptionId := subscription.GetSubscriptionId() @@ -105,7 +105,7 @@ func (s *Subscriber) subscribeSync(ctx context.Context, result chan apiModel.Plc subscriptionValues[tagName] = handle } - result <- spiModel.NewDefaultPlcSubscriptionRequestResult( + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcSubscriptionRequestResult( subscriptionRequest, spiModel.NewDefaultPlcSubscriptionResponse( subscriptionRequest, @@ -114,7 +114,7 @@ func (s *Subscriber) subscribeSync(ctx context.Context, result chan apiModel.Plc append(s._options, options.WithCustomLogger(s.log))..., ), nil, - ) + )) } func (s *Subscriber) onSubscribeCreateSubscription(ctx context.Context, cycleTime time.Duration) (readWriteModel.CreateSubscriptionResponse, error) { @@ -200,7 +200,7 @@ func (s *Subscriber) onDisconnect() { func (s *Subscriber) Unsubscribe(ctx context.Context, unsubscriptionRequest apiModel.PlcUnsubscriptionRequest) <-chan apiModel.PlcUnsubscriptionRequestResult { result := make(chan apiModel.PlcUnsubscriptionRequestResult, 1) - result <- spiModel.NewDefaultPlcUnsubscriptionRequestResult(unsubscriptionRequest, nil, errors.New("Not Implemented")) + utils.DeliverResult(s.log, result, spiModel.NewDefaultPlcUnsubscriptionRequestResult(unsubscriptionRequest, nil, errors.New("Not Implemented"))) for _, handle := range unsubscriptionRequest.(*spiModel.DefaultPlcUnsubscriptionRequest).GetSubscriptionHandles() { handle.(*SubscriptionHandle).stopSubscriber() diff --git a/plc4go/internal/opcua/Writer.go b/plc4go/internal/opcua/Writer.go index 9b5eb92714..d94626f73c 100644 --- a/plc4go/internal/opcua/Writer.go +++ b/plc4go/internal/opcua/Writer.go @@ -61,7 +61,7 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques func (m *Writer) WriteSync(ctx context.Context, writeRequest apiModel.PlcWriteRequest, result chan apiModel.PlcWriteRequestResult) { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() @@ -80,13 +80,13 @@ func (m *Writer) WriteSync(ctx context.Context, writeRequest apiModel.PlcWriteRe nodeId, err := generateNodeId(tag) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "error generating node id from tag %s", tag)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "error generating node id from tag %s", tag))) return } plcValue, err := m.fromPlcValue(tagName, tag, writeRequest) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error getting plcValue")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error getting plcValue"))) return } writeValueArray[i] = readWriteModel.NewWriteValue(nodeId, @@ -130,18 +130,18 @@ func (m *Writer) WriteSync(ctx context.Context, writeRequest apiModel.PlcWriteRe ) buffer := utils.NewWriteBufferByteBased(utils.WithByteOrderForByteBasedBuffer(binary.LittleEndian)) if err := extObject.SerializeWithWriteBuffer(ctx, buffer); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Unable to serialise the ReadRequest")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Unable to serialise the ReadRequest"))) return } consumer := func(opcuaResponse []byte) { reply, err := readWriteModel.ExtensionObjectParseWithBuffer[readWriteModel.ExtensionObject](ctx, utils.NewReadBufferByteBased(opcuaResponse, utils.WithByteOrderForReadBufferByteBased(binary.LittleEndian)), false) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Unable to read the reply")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Unable to read the reply"))) return } if writeResponse, ok := reply.(readWriteModel.WriteResponse); ok { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(m.writeResponse(writeRequest, writeResponse.GetResults())), nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(m.writeResponse(writeRequest, writeResponse.GetResults())), nil)) return } else { if serviceFault, ok := reply.(readWriteModel.ServiceFault); ok { @@ -155,12 +155,12 @@ func (m *Writer) WriteSync(ctx context.Context, writeRequest apiModel.PlcWriteRe for _, tagName := range writeRequest.GetTagNames() { responseCodes[tagName] = apiModel.PlcResponseCode_INTERNAL_ERROR } - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes), nil) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, spiModel.NewDefaultPlcWriteResponse(writeRequest, responseCodes), nil)) } } errorDispatcher := func(err error) { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, err) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, err)) } m.connection.channel.submit(ctx, m.connection.messageCodec, errorDispatcher, consumer, buffer) diff --git a/plc4go/internal/s7/Reader.go b/plc4go/internal/s7/Reader.go index c99ee53497..7ddf5b158d 100644 --- a/plc4go/internal/s7/Reader.go +++ b/plc4go/internal/s7/Reader.go @@ -35,6 +35,7 @@ import ( spiModel "github.com/apache/plc4x/plc4go/spi/model" "github.com/apache/plc4x/plc4go/spi/options" "github.com/apache/plc4x/plc4go/spi/transactions" + "github.com/apache/plc4x/plc4go/spi/utils" spiValues "github.com/apache/plc4x/plc4go/spi/values" ) @@ -68,7 +69,7 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult(readRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() @@ -77,11 +78,11 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) tag := readRequest.GetTag(tagName) address, err := encodeS7Address(tag) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrapf(err, "Error encoding s7 address for tag %s", tagName), - ) + )) return } requestItems[i] = readWriteModel.NewS7VarRequestParameterItemAddress(address) @@ -145,32 +146,32 @@ func (m *Reader) Read(ctx context.Context, readRequest apiModel.PlcReadRequest) readResponse, err := m.ToPlc4xReadResponse(payload, readRequest) if err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "Error decoding response"), - ) + )) return transaction.EndRequest() } - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, readResponse, nil, - ) + )) return transaction.EndRequest() }, func(err error) error { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "got timeout while waiting for response"), - ) + )) return transaction.EndRequest() }); err != nil { - result <- spiModel.NewDefaultPlcReadRequestResult( + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcReadRequestResult( readRequest, nil, errors.Wrap(err, "error sending message"), - ) + )) if err := transaction.FailRequest(errors.Errorf("timeout after %s", 1*time.Second)); err != nil { m.log.Debug().Err(err).Msg("Error failing request") } diff --git a/plc4go/internal/s7/Writer.go b/plc4go/internal/s7/Writer.go index c086350b1a..8387fff491 100644 --- a/plc4go/internal/s7/Writer.go +++ b/plc4go/internal/s7/Writer.go @@ -35,6 +35,7 @@ import ( spiModel "github.com/apache/plc4x/plc4go/spi/model" "github.com/apache/plc4x/plc4go/spi/options" "github.com/apache/plc4x/plc4go/spi/transactions" + "github.com/apache/plc4x/plc4go/spi/utils" ) type Writer struct { @@ -63,7 +64,7 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques m.wg.Go(func() { defer func() { if err := recover(); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack())) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Errorf("panic-ed %v. Stack: %s", err, debug.Stack()))) } }() parameterItems := make([]readWriteModel.S7VarRequestParameterItem, len(writeRequest.GetTagNames())) @@ -73,13 +74,13 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques plcValue := writeRequest.GetValue(tagName) s7Address, err := encodeS7Address(tag) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding s7 address for tag %s", tagName)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding s7 address for tag %s", tagName))) return } parameterItems[i] = readWriteModel.NewS7VarRequestParameterItemAddress(s7Address) value, err := serializePlcValue(tag, plcValue) if err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding value for tag %s", tagName)) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrapf(err, "Error encoding value for tag %s", tagName))) return } payloadItems[i] = value @@ -136,25 +137,16 @@ func (m *Writer) Write(ctx context.Context, writeRequest apiModel.PlcWriteReques readResponse, err := m.ToPlc4xWriteResponse(payload, writeRequest) if err != nil { - result <- &spiModel.DefaultPlcWriteRequestResult{ - Request: writeRequest, - Err: errors.Wrap(err, "Error decoding response"), - } + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "Error decoding response"))) return transaction.EndRequest() } - result <- &spiModel.DefaultPlcWriteRequestResult{ - Request: writeRequest, - Response: readResponse, - } + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, readResponse, nil)) return transaction.EndRequest() }, func(err error) error { - result <- &spiModel.DefaultPlcWriteRequestResult{ - Request: writeRequest, - Err: errors.New("got timeout while waiting for response"), - } + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.New("got timeout while waiting for response"))) return transaction.EndRequest() }); err != nil { - result <- spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error sending message")) + utils.DeliverResult(m.log, result, spiModel.NewDefaultPlcWriteRequestResult(writeRequest, nil, errors.Wrap(err, "error sending message"))) if err := transaction.FailRequest(errors.Errorf("timeout after %s", 1*time.Second)); err != nil { m.log.Debug().Err(err).Msg("Error failing request") }
