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)
|