LCOV - code coverage report
Current view: top level - libs/bsw/cpp2someip/src - BufferedEventSender.cpp (source / functions) Coverage Total Hit
Test: coverage.info Lines: 77.3 % 176 136
Test Date: 2026-09-11 12:05:06 Functions: 92.9 % 14 13

            Line data    Source code
       1              : /********************************************************************************
       2              :  * Copyright (c) 2026 Accenture
       3              :  *
       4              :  * This program and the accompanying materials are made available under the
       5              :  * terms of the Apache License Version 2.0 which is available at
       6              :  * https://www.apache.org/licenses/LICENSE-2.0
       7              :  *
       8              :  * SPDX-License-Identifier: Apache-2.0
       9              :  ********************************************************************************/
      10              : 
      11              : #include "someip/BufferedEventSender.h"
      12              : 
      13              : #include "bsp/timer/SystemTimer.h"
      14              : #include "someip/ISomeIpSerializable.h"
      15              : #include "someip/NetworkChannel.h"
      16              : #include "someip/SomeIpMessage.h"
      17              : #include "someip/SomeIpSerializer.h"
      18              : #include "someip/Statistics.h"
      19              : #include "someip/TpTransceiver.h"
      20              : #include "someip/logger.h"
      21              : 
      22              : #include <etl/algorithm.h>
      23              : 
      24              : // Logger API uses printf-style varargs for fixed diagnostic messages in this module.
      25              : // NOLINTBEGIN(cppcoreguidelines-pro-type-vararg)
      26              : 
      27              : namespace someip
      28              : {
      29              : using ::util::logger::SOMEIP;
      30              : 
      31           25 : BufferedEventSender::BufferedEventSender(
      32              :     INetwork& network,
      33              :     ::async::ContextType const ethernetContext,
      34              :     ITpTransceiver& tpTransceiver,
      35           25 :     ::etl::span<internal::EventMessage*> messages)
      36           25 : : _ethernetContext(ethernetContext)
      37           25 : , _sendNextMessageFunction(
      38              :       ::async::Function::CallType::
      39           25 :           create<BufferedEventSender, &BufferedEventSender::sendNextMessage>(*this))
      40              : , _sendNextMessageTimeout()
      41           25 : , _network(network)
      42           25 : , _tpTransceiver(tpTransceiver)
      43           25 : , _eventMessages(messages)
      44           25 : , _nextSendTime(0U)
      45           25 : {}
      46              : 
      47            2 : void BufferedEventSender::init() { _nextSendTime = 0U; }
      48              : 
      49            2 : void BufferedEventSender::shutdown()
      50              : {
      51            2 :     _sendNextMessageTimeout.cancel();
      52              : 
      53           11 :     for (internal::EventMessage* const message : _eventMessages)
      54              :     {
      55            9 :         if (message != nullptr)
      56              :         {
      57            9 :             message->clear();
      58              :         }
      59              :     }
      60            2 : }
      61              : 
      62           24 : IEventSender::ErrorCode BufferedEventSender::sendEvent(
      63              :     service_id::type const serviceId,
      64              :     major_version::type const majorVersion,
      65              :     uint16_t const eventId,
      66              :     uint32_t const maximumDelayTime,
      67              :     ISomeIpSerializable const* const payload,
      68              :     uint16_t const localPort,
      69              :     uint8_t const proto,
      70              :     ::ip::IPEndpoint const& destination,
      71              :     uint16_t const sessionId)
      72              : {
      73           24 :     Statistics::incCounter(Statistics::Counter::PDU_TX);
      74              : 
      75           24 :     auto channel = _network.getRpcChannel(localPort, destination, proto);
      76              : 
      77           24 :     if (!channel.has_value())
      78              :     {
      79            1 :         WARN_LOG(SOMEIP, "BufferedEventSender::sendSingleEvent() no channel");
      80            1 :         return IEventSender::ErrorCode::EVENT_SEND_ERROR;
      81              :     }
      82              : 
      83           23 :     size_t eventSize = 0U;
      84           23 :     if (!serializeEvent(
      85              :             serviceId,
      86              :             majorVersion,
      87              :             eventId,
      88              :             payload,
      89           46 :             channel->getOutputBuffer(),
      90              :             &eventSize,
      91              :             sessionId))
      92              :     {
      93            0 :         return IEventSender::ErrorCode::EVENT_SEND_ERROR;
      94              :     }
      95              : 
      96              :     // check if we should send this event immediately
      97           46 :     if ((_eventMessages.size() == 0U) || (maximumDelayTime == 0U)
      98           46 :         || (eventSize > internal::EventMessage::BUFFER_SIZE))
      99              :     {
     100            3 :         return sendSingleEvent(
     101            3 :             channel.value(), localPort, proto, destination, sessionId, eventSize);
     102              :     }
     103              : 
     104           20 :     uint32_t const pdu
     105           20 :         = ((static_cast<uint32_t>(serviceId) << 16U) | static_cast<uint32_t>(eventId));
     106           20 :     uint32_t const currentTime = static_cast<uint32_t>(getSystemTimeMs32Bit());
     107           20 :     uint32_t const sendTime
     108              :         = currentTime
     109           20 :           + ((maximumDelayTime > MAX_PACKET_DELAY) ? MAX_PACKET_DELAY : maximumDelayTime);
     110              : 
     111              :     // check if we already got a message for this endpoint
     112           20 :     internal::EventMessage* message = findMessage(destination, localPort, proto);
     113              : 
     114           20 :     if (message != nullptr)
     115              :     {
     116              :         // send message if not enough space left or already full
     117            3 :         if (((message->getBufferOffset() + eventSize) > internal::EventMessage::BUFFER_SIZE)
     118            3 :             || message->isFull())
     119              :         {
     120            0 :             Statistics::incCounter(Statistics::Counter::TRAIN_FULL);
     121            0 :             sendBufferedEvents(*message);
     122              :         }
     123            3 :         else if (message->hasPdu(pdu)) // send message if duplicate pdu
     124              :         {
     125            2 :             Statistics::incCounter(Statistics::Counter::TRAIN_DUPLICATE_PDU);
     126            2 :             sendBufferedEvents(*message);
     127              :         }
     128              :         else
     129              :         {
     130              :             // misra
     131              :         }
     132              :     }
     133              :     else
     134              :     {
     135              :         // check if we have a free message
     136           17 :         message = getEmptyMessage();
     137           17 :         if (message == nullptr)
     138              :         {
     139              :             // send next scheduled message
     140            1 :             message = getNextScheduledMessage();
     141            1 :             if (message != nullptr)
     142              :             {
     143            1 :                 sendBufferedEvents(*message);
     144              :             }
     145              :         }
     146              :     }
     147              : 
     148           20 :     if (message == nullptr) // no message available !
     149              :     {
     150            0 :         ERROR_LOG(SOMEIP, "BufferedEventSender::sendEvent: no free buffer");
     151            0 :         Statistics::incCounter(Statistics::Counter::TRAIN_NOT_AVAILABLE);
     152            0 :         return IEventSender::ErrorCode::EVENT_SEND_ERROR;
     153              :     }
     154              : 
     155           20 :     bool isNew = false;
     156           20 :     if (!message->isInitialized())
     157              :     {
     158           19 :         message->init(destination, localPort, proto);
     159           19 :         isNew = true;
     160              :     }
     161              : 
     162           20 :     if (bufferEvent(*message, eventSize, channel->getOutputBuffer()))
     163              :     {
     164           20 :         message->addPdu(pdu);
     165           20 :         message->adjustSendTime(sendTime);
     166              : 
     167           20 :         updateSchedule(currentTime);
     168           20 :         return IEventSender::ErrorCode::EVENT_SEND_OK;
     169              :     }
     170              : 
     171              :     // cleanup
     172            0 :     if (isNew)
     173              :     {
     174            0 :         message->clear();
     175              :     }
     176              : 
     177            0 :     return IEventSender::ErrorCode::EVENT_SEND_ERROR;
     178           24 : }
     179              : 
     180              : // virtual
     181            0 : void BufferedEventSender::sendNextMessage()
     182              : {
     183            0 :     uint32_t const currentTime = static_cast<uint32_t>(getSystemTimeMs32Bit());
     184            0 :     _nextSendTime              = 0U;
     185              : 
     186              :     // send expired messages
     187            0 :     for (internal::EventMessage* const message : _eventMessages)
     188              :     {
     189            0 :         if (message != nullptr)
     190              :         {
     191            0 :             if (!message->isInitialized())
     192              :             {
     193            0 :                 continue;
     194              :             }
     195              : 
     196            0 :             uint32_t const sendTime = message->getSendTime();
     197              : 
     198            0 :             if (currentTime >= sendTime)
     199              :             {
     200            0 :                 Statistics::incCounter(Statistics::Counter::TRAIN_TIMEOUT);
     201            0 :                 sendBufferedEvents(*message);
     202              :             }
     203            0 :             else if ((_nextSendTime == 0U) || (_nextSendTime > sendTime))
     204              :             {
     205            0 :                 _nextSendTime = sendTime;
     206              :             }
     207              :             else
     208              :             {
     209              :                 // misra
     210              :             }
     211              :         }
     212              :         else
     213              :         {
     214            0 :             ERROR_LOG(SOMEIP, "BufferedEventSender::expired() nullptr msg");
     215              :         }
     216              :     }
     217              : 
     218              :     // schedule next message
     219            0 :     if (_nextSendTime != 0U)
     220              :     {
     221            0 :         uint32_t const timeout = _nextSendTime - currentTime;
     222            0 :         async::schedule(
     223            0 :             _ethernetContext,
     224              :             _sendNextMessageFunction,
     225            0 :             _sendNextMessageTimeout,
     226              :             timeout,
     227              :             ::async::TimeUnit::MILLISECONDS);
     228              :     }
     229            0 : }
     230              : 
     231            3 : IEventSender::ErrorCode BufferedEventSender::sendSingleEvent(
     232              :     NetworkChannel& channel,
     233              :     uint16_t const /* localPort */,
     234              :     uint8_t const proto,
     235              :     ::ip::IPEndpoint const& /* destination */,
     236              :     uint16_t const sessionId,
     237              :     size_t const length) const
     238              : {
     239            3 :     IEventSender::ErrorCode errorCode = IEventSender::ErrorCode::EVENT_SEND_OK;
     240              : 
     241            3 :     if (ITpTransceiver::isOutgoingTpMessage(proto, static_cast<uint32_t>(length)))
     242              :     {
     243            1 :         SomeIpMessage message(channel.getOutputBuffer().first(length));
     244            1 :         message.setSessionId(sessionId);
     245              : 
     246            1 :         if (!_tpTransceiver.sendTpMessage(channel, message))
     247              :         {
     248            0 :             errorCode = IEventSender::ErrorCode::EVENT_SEND_ERROR;
     249              :         }
     250              :     }
     251              :     else
     252              :     {
     253            2 :         if (!channel.send(static_cast<uint32_t>(length)))
     254              :         {
     255            0 :             errorCode = IEventSender::ErrorCode::EVENT_SEND_ERROR;
     256              :         }
     257              :     }
     258              : 
     259            3 :     if (errorCode != IEventSender::ErrorCode::EVENT_SEND_OK)
     260              :     {
     261            0 :         WARN_LOG(SOMEIP, "BufferedEventSender::sendSingleEvent() send failed");
     262              :     }
     263              : 
     264            3 :     return errorCode;
     265              : }
     266              : 
     267            3 : void BufferedEventSender::sendBufferedEvents(internal::EventMessage& message)
     268              : {
     269            3 :     auto channel = _network.getRpcChannel(
     270            3 :         message.getLocalPort(), message.getDestination(), message.getProto());
     271              : 
     272            3 :     if (!channel.has_value())
     273              :     {
     274            0 :         WARN_LOG(SOMEIP, "BufferedEventSender::sendBufferedEvents() no channel");
     275              :     }
     276              :     else
     277              :     {
     278            3 :         size_t const length = message.getBufferOffset();
     279              : 
     280            3 :         if (!channel->send(static_cast<uint32_t>(length), message.getBuffer().first(length)))
     281              :         {
     282            0 :             WARN_LOG(SOMEIP, "BufferedEventSender::sendBufferedEvents() send failed");
     283              :         }
     284              :     }
     285              : 
     286            3 :     message.clear();
     287            3 : }
     288              : 
     289           20 : bool BufferedEventSender::bufferEvent(
     290              :     internal::EventMessage& message, size_t const length, ::etl::span<uint8_t> const& eventBuffer)
     291              : {
     292           20 :     ::etl::span<uint8_t> const buffer = message.getBuffer().subspan(message.getBufferOffset());
     293              : 
     294           20 :     etl::copy_n(eventBuffer.begin(), length, buffer.begin());
     295              : 
     296           20 :     message.incBufferOffset(static_cast<uint16_t>(length));
     297           20 :     return true;
     298              : }
     299              : 
     300           23 : bool BufferedEventSender::serializeEvent(
     301              :     service_id::type const serviceId,
     302              :     major_version::type const majorVersion,
     303              :     uint16_t const eventId,
     304              :     ISomeIpSerializable const* const payload,
     305              :     ::etl::span<uint8_t> const buffer,
     306              :     size_t* const length,
     307              :     uint16_t const sessionId)
     308              : {
     309           23 :     if (buffer.size() < SomeIpMessage::OFFSET_PAYLOAD)
     310              :     {
     311            0 :         WARN_LOG(SOMEIP, "BufferedEventSender::serializeEvent() buffer is too small");
     312            0 :         return false;
     313              :     }
     314              : 
     315           23 :     SomeIpMessage message(buffer);
     316              : 
     317           23 :     message.setServiceId(serviceId);
     318           23 :     message.setMethodId(eventId);
     319           23 :     message.setPayloadLength(0);
     320           23 :     message.setRequestId(0U);
     321           23 :     message.setMessageType(SomeIpMessage::MessageType::NOTIFICATION);
     322           23 :     message.setReturnCode(SomeIpMessage::ReturnCode::SOMEIP_E_OK);
     323           23 :     message.setProtocolVersion(1U);
     324           23 :     message.setInterfaceVersion(majorVersion);
     325           23 :     message.setSessionId(sessionId);
     326              : 
     327           23 :     if (payload != nullptr)
     328              :     {
     329              :         SomeIpSerializer serializer(
     330           21 :             ::etl::span<uint8_t>(message.getPayload(), message.getMaximumPayloadLength()));
     331              : 
     332           21 :         serializer << *payload;
     333              : 
     334           21 :         if (!serializer.isGood())
     335              :         {
     336            0 :             WARN_LOG(SOMEIP, "BufferedEventSender::serializeEvent() serialization failed");
     337            0 :             return false;
     338              :         }
     339           21 :         message.setPayloadLength(static_cast<uint32_t>(serializer.getCurrentPosition()));
     340              :     }
     341              : 
     342           23 :     if (length != nullptr)
     343              :     {
     344           23 :         *length = message.getTotalLength();
     345              :     }
     346              :     else
     347              :     {
     348            0 :         ERROR_LOG(SOMEIP, "BufferedEventSender::serializeEvent() invalid length");
     349              :     }
     350           23 :     return true;
     351              : }
     352              : 
     353           25 : size_t BufferedEventSender::countMessages() const
     354              : {
     355           25 :     size_t result = 0U;
     356              : 
     357          225 :     for (internal::EventMessage* const message : _eventMessages)
     358              :     {
     359          200 :         if ((message != nullptr) && (message->isInitialized()))
     360              :         {
     361           77 :             result++;
     362              :         }
     363              :     }
     364              : 
     365           25 :     return result;
     366              : }
     367              : 
     368           20 : internal::EventMessage* BufferedEventSender::findMessage(
     369              :     ::ip::IPEndpoint const& destination, port::type const localPort, proto::type const proto) const
     370              : {
     371          156 :     for (internal::EventMessage* const message : _eventMessages)
     372              :     {
     373          139 :         if ((message != nullptr) && (message->isMatching(destination, localPort, proto)))
     374              :         {
     375            3 :             return message;
     376              :         }
     377              :     }
     378              : 
     379           17 :     return nullptr;
     380              : }
     381              : 
     382           17 : internal::EventMessage* BufferedEventSender::getEmptyMessage() const
     383              : {
     384           54 :     for (internal::EventMessage* const message : _eventMessages)
     385              :     {
     386           53 :         if ((message != nullptr) && (!message->isInitialized()))
     387              :         {
     388           16 :             return message;
     389              :         }
     390              :     }
     391              : 
     392            1 :     return nullptr;
     393              : }
     394              : 
     395           21 : internal::EventMessage* BufferedEventSender::getNextScheduledMessage() const
     396              : {
     397           21 :     internal::EventMessage* result = nullptr;
     398              : 
     399          189 :     for (internal::EventMessage* const message : _eventMessages)
     400              :     {
     401          168 :         if ((message != nullptr) && (!message->isInitialized()))
     402              :         {
     403          104 :             continue;
     404              :         }
     405           64 :         if ((result == nullptr)
     406           64 :             || ((message != nullptr) && (result->getSendTime() > message->getSendTime())))
     407              :         {
     408           22 :             result = message;
     409              :         }
     410              :     }
     411              : 
     412           21 :     return result;
     413              : }
     414              : 
     415           20 : void BufferedEventSender::updateSchedule(uint32_t const currentTime)
     416              : {
     417           20 :     internal::EventMessage* const message = getNextScheduledMessage();
     418              : 
     419           20 :     if (message == nullptr)
     420              :     {
     421            0 :         return;
     422              :     }
     423              : 
     424           20 :     uint32_t const sendTime = message->getSendTime();
     425              : 
     426           20 :     if ((_nextSendTime == 0U) || (_nextSendTime > sendTime))
     427              :     {
     428            8 :         _nextSendTime          = sendTime;
     429            8 :         uint32_t const timeout = sendTime - currentTime;
     430              : 
     431            8 :         _sendNextMessageTimeout.cancel();
     432            8 :         async::schedule(
     433            8 :             _ethernetContext,
     434              :             _sendNextMessageFunction,
     435            8 :             _sendNextMessageTimeout,
     436              :             timeout,
     437              :             ::async::TimeUnit::MILLISECONDS);
     438              :     }
     439              : }
     440              : 
     441              : } // namespace someip
     442              : 
     443              : // NOLINTEND(cppcoreguidelines-pro-type-vararg)
        

Generated by: LCOV version 2.0-1