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/ServiceAnnouncer.h"
12 :
13 : #include "bsp/timer/SystemTimer.h"
14 : #include "someip/INetwork.h"
15 : #include "someip/IServiceRegistry.h"
16 : #include "someip/NetworkChannel.h"
17 : #include "someip/ProvidedService.h"
18 : #include "someip/ServiceDescription.h"
19 : #include "someip/SomeIpConstants.h"
20 : #include "someip/SomeIpMessage.h"
21 : #include "someip/Statistics.h"
22 : #include "someip/logger.h"
23 :
24 : #include <etl/algorithm.h>
25 :
26 : // Logger API uses printf-style varargs for fixed diagnostic messages in this module.
27 : // NOLINTBEGIN(cppcoreguidelines-pro-type-vararg)
28 :
29 : namespace someip
30 : {
31 : using ::util::logger::SOMEIP;
32 :
33 24 : ServiceAnnouncer::ServiceAnnouncer(
34 : INetwork& network,
35 : ServiceManager& serviceManager,
36 : IServiceRegistry& serviceRegistry,
37 : ::async::ContextType const ethernetContext,
38 : QueryManager& queryManager,
39 24 : SessionManager& sessionManager)
40 24 : : _network(network)
41 24 : , _serviceManager(serviceManager)
42 24 : , _serviceRegistry(serviceRegistry)
43 24 : , _ethernetContext(ethernetContext)
44 24 : , _cyclicFunction(
45 48 : ::async::Function::CallType::create<ServiceAnnouncer, &ServiceAnnouncer::cyclic>(*this))
46 : , _cyclicTimeout()
47 24 : , _eventFunction(
48 : ::async::Function::CallType::
49 24 : create<ServiceAnnouncer, &ServiceAnnouncer::checkPendingTasksAndSendDueMessages>(*this))
50 : , _eventTimeout()
51 24 : , _messageBuilder()
52 24 : , _messageBuffer()
53 24 : , _isStarted(false)
54 24 : , _taskPool()
55 24 : , _pendingBrowseRequests()
56 24 : , _pendingTxMessages()
57 24 : , _queryManager(queryManager)
58 48 : , _sessionManager(sessionManager)
59 24 : {}
60 :
61 21 : void ServiceAnnouncer::init()
62 : {
63 21 : _sessionManager.init();
64 21 : etl::fill(_messageBuffer.begin(), _messageBuffer.end(), 0);
65 :
66 21 : _taskPool.release_all();
67 21 : _pendingBrowseRequests.clear();
68 21 : _pendingTxMessages.clear();
69 21 : }
70 :
71 1 : void ServiceAnnouncer::shutdown() { stop(); }
72 :
73 22 : void ServiceAnnouncer::start()
74 : {
75 22 : if (!_isStarted)
76 : {
77 21 : async::scheduleAtFixedRate(
78 21 : _ethernetContext,
79 : _cyclicFunction,
80 21 : _cyclicTimeout,
81 : CYCLE_TIME_MS,
82 : ::async::TimeUnit::MILLISECONDS);
83 21 : _isStarted = true;
84 : }
85 22 : }
86 :
87 1 : void ServiceAnnouncer::stop()
88 : {
89 1 : if (_isStarted)
90 : {
91 1 : checkPendingTasks(getSystemTimeMs32Bit());
92 1 : sendStopOffers();
93 1 : _cyclicTimeout.cancel();
94 1 : _eventTimeout.cancel();
95 1 : _isStarted = false;
96 : }
97 1 : }
98 :
99 5 : void ServiceAnnouncer::respondToFindService(
100 : service_id::type serviceId,
101 : instance_id::type instanceId,
102 : major_version::type majorVersion,
103 : minor_version::type minorVersion,
104 : ttl::type,
105 : ::ip::IPAddress const& sourceIpAddress,
106 : bool unicast)
107 : {
108 5 : if ((instanceId == instance_id::ANY) || (majorVersion == major_version::ANY)
109 2 : || (minorVersion == minor_version::ANY))
110 : {
111 4 : if (!_taskPool.full())
112 : {
113 4 : ServiceAnnouncerTask& task = *_taskPool.create();
114 4 : task.init(
115 : serviceId,
116 : instance_id::ANY,
117 : eventgroup_id::ALL,
118 : ttl::INVALID,
119 : majorVersion,
120 : minorVersion);
121 :
122 4 : uint32_t const delay = REQ_RES_MAX_DELAY_MS; // TODO: determine delay randomly
123 4 : task.setDestinationAddress(sourceIpAddress);
124 4 : task.setUnicast(unicast);
125 4 : task.setTimestamp(getSystemTimeMs32Bit() + delay);
126 4 : _pendingBrowseRequests.push_front(task);
127 :
128 4 : if (_taskPool.full())
129 : {
130 0 : checkPendingTasks(getSystemTimeMs32Bit());
131 : }
132 : }
133 : else
134 : {
135 0 : WARN_LOG(SOMEIP, "ServiceAnnouncer::respondToFindService() task pool is full");
136 : }
137 4 : }
138 : else
139 : {
140 1 : auto service = ::someip::make<ServiceDescription>();
141 1 : service.serviceId = serviceId;
142 1 : service.majorVersion = majorVersion;
143 1 : service.minorVersion = minorVersion; // don't care
144 1 : service.instanceId = instanceId;
145 :
146 1 : ProvidedService const* const providedService = _serviceManager.getService(service);
147 :
148 : // only answer if service is in its MAIN phase
149 1 : if ((providedService != nullptr)
150 1 : && (ProvidedService::ProvidedServiceState::MAIN_PHASE == providedService->getState()))
151 : {
152 1 : if (!_taskPool.full())
153 : {
154 1 : ServiceAnnouncerTask& task = *_taskPool.create();
155 :
156 1 : task.initFrom(
157 1 : providedService->description,
158 : sourceIpAddress,
159 : unicast,
160 1 : getSystemTimeMs32Bit() + REQ_RES_MIN_DELAY_MS,
161 : ServiceAnnouncerTask::TaskType::TASK_ANNOUNCE);
162 :
163 1 : _pendingTxMessages.push_front(task);
164 :
165 1 : if (_taskPool.full())
166 : {
167 0 : checkPendingTasks(getSystemTimeMs32Bit());
168 : }
169 : else
170 : {
171 1 : triggerEventTimeout();
172 : }
173 : }
174 : else
175 : {
176 0 : WARN_LOG(SOMEIP, "ServiceAnnouncer::respondToFindService() task pool is full");
177 : }
178 : }
179 : }
180 5 : }
181 :
182 4 : void ServiceAnnouncer::respondToSubscribe(
183 : service_id::type const serviceId,
184 : instance_id::type const instanceId,
185 : major_version::type const majorVersion,
186 : uint16_t const reserved,
187 : eventgroup_id::type const eventgroup,
188 : ttl::type const ttl,
189 : ::ip::IPAddress const& sourceIpAddress,
190 : ::ip::IPAddress const& endpointIpAddress,
191 : uint16_t const endpointPort,
192 : uint8_t const endpointProto)
193 : {
194 : IServiceRegistry::SubscriptionResult result;
195 4 : result = _serviceRegistry.subscribeReceived(
196 : serviceId,
197 : instanceId,
198 : majorVersion,
199 : eventgroup,
200 : ttl,
201 : endpointIpAddress,
202 : endpointPort,
203 : endpointProto);
204 :
205 4 : if (IServiceRegistry::SubscriptionResult::SUBSCRIBE_OK == result)
206 : {
207 0 : sendSubscribeAck(
208 : serviceId, instanceId, eventgroup, majorVersion, reserved, ttl, sourceIpAddress);
209 : }
210 4 : else if (IServiceRegistry::SubscriptionResult::SUBSCRIBE_OK_MULTICAST == result)
211 : {
212 2 : auto service = ::someip::make<ServiceDescription>();
213 2 : service.serviceId = serviceId;
214 2 : service.instanceId = instanceId;
215 2 : service.majorVersion = majorVersion;
216 2 : service.eventGroup = eventgroup;
217 :
218 2 : ProvidedService const* const providedService = _serviceManager.getEventGroup(service);
219 :
220 2 : if (providedService != nullptr)
221 : {
222 1 : sendSubscribeAckMulticast(
223 : serviceId,
224 : instanceId,
225 : eventgroup,
226 : majorVersion,
227 : reserved,
228 : ttl,
229 1 : providedService->description.ipAddress,
230 1 : providedService->description.port,
231 : sourceIpAddress);
232 : }
233 : }
234 2 : else if (IServiceRegistry::SubscriptionResult::SUBSCRIBE_ERROR == result)
235 : {
236 1 : sendSubscribeNack(
237 : serviceId, instanceId, eventgroup, majorVersion, reserved, sourceIpAddress);
238 : }
239 : else
240 : {
241 : // misra
242 : }
243 4 : }
244 :
245 1 : void ServiceAnnouncer::sendSubscribeAck(
246 : service_id::type const serviceId,
247 : instance_id::type const instanceId,
248 : eventgroup_id::type const eventgroup,
249 : major_version::type const majorVersion,
250 : uint16_t const reserved,
251 : ttl::type const ttl,
252 : ::ip::IPAddress const& sourceIpAddress)
253 : {
254 1 : ServiceAnnouncerTask task;
255 1 : task.init(
256 : serviceId, instanceId, eventgroup, ttl, majorVersion, static_cast<uint32_t>(reserved));
257 1 : task.setDestinationAddress(sourceIpAddress);
258 1 : task.setTimestamp(getSystemTimeMs32Bit() + REQ_RES_MIN_DELAY_MS);
259 1 : task.setType(ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE_ACK);
260 1 : enqueueTxMessage(task);
261 1 : }
262 :
263 2 : void ServiceAnnouncer::sendSubscribeNack(
264 : service_id::type const serviceId,
265 : instance_id::type const instanceId,
266 : eventgroup_id::type const eventgroup,
267 : major_version::type const majorVersion,
268 : uint16_t const reserved,
269 : ::ip::IPAddress const& sourceIpAddress)
270 : {
271 2 : ServiceAnnouncerTask task;
272 2 : task.init(serviceId, instanceId, eventgroup, 0U, majorVersion, static_cast<uint32_t>(reserved));
273 2 : task.setDestinationAddress(sourceIpAddress);
274 2 : task.setTimestamp(getSystemTimeMs32Bit() + REQ_RES_MIN_DELAY_MS);
275 2 : task.setType(ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE_NACK);
276 2 : enqueueTxMessage(task);
277 2 : }
278 :
279 3 : void ServiceAnnouncer::sendSubscribeAckMulticast(
280 : service_id::type const serviceId,
281 : instance_id::type const instanceId,
282 : eventgroup_id::type const eventgroup,
283 : major_version::type const majorVersion,
284 : uint16_t const reserved,
285 : ttl::type const ttl,
286 : ::ip::IPAddress const& endpointAddress,
287 : uint16_t const endpointPort,
288 : ::ip::IPAddress const& sourceIpAddress)
289 : {
290 3 : ServiceAnnouncerTask task;
291 3 : task.init(
292 : serviceId, instanceId, eventgroup, ttl, majorVersion, static_cast<uint32_t>(reserved));
293 3 : task.setProto(proto::SD_L4_PROTO_UDP);
294 3 : task.setEndpoint(endpointAddress, endpointPort);
295 :
296 3 : task.setDestinationAddress(sourceIpAddress);
297 3 : task.setTimestamp(getSystemTimeMs32Bit() + REQ_RES_MIN_DELAY_MS);
298 3 : task.setType(ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE_ACK_MULTICAST);
299 3 : enqueueTxMessage(task);
300 3 : }
301 :
302 2 : void ServiceAnnouncer::subscribe(
303 : ServiceDescription const& service, ::ip::IPAddress const& sourceAddress)
304 : {
305 2 : if (!containsEventGroup(service))
306 : {
307 1 : return; // only event groups can be subscribed
308 : }
309 :
310 1 : ServiceAnnouncerTask task;
311 :
312 1 : task.initFrom(
313 : service,
314 : sourceAddress,
315 : false,
316 1 : getSystemTimeMs32Bit() + REQ_RES_MIN_DELAY_MS,
317 : ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE);
318 :
319 1 : task.setMinorVersion(minor_version::INVALID);
320 1 : task.setEndpoint(_network.getLocalIp(), service.port);
321 1 : enqueueTxMessage(task);
322 : }
323 :
324 3 : void ServiceAnnouncer::unsubscribe(
325 : ServiceDescription const& service, ::ip::IPAddress const& sourceAddress)
326 : {
327 3 : if (!containsEventGroup(service))
328 : {
329 1 : return; // only event groups can be unsubscribed
330 : }
331 :
332 2 : ServiceAnnouncerTask task;
333 :
334 2 : task.initFrom(
335 : service,
336 : sourceAddress,
337 : false,
338 2 : getSystemTimeMs32Bit() + REQ_RES_MIN_DELAY_MS,
339 : ServiceAnnouncerTask::TaskType::TASK_UNSUBSCRIBE);
340 :
341 2 : task.setMinorVersion(minor_version::INVALID);
342 2 : task.setEndpoint(_network.getLocalIp(), service.port);
343 2 : enqueueTxMessage(task);
344 : }
345 :
346 1 : void ServiceAnnouncer::sendStopOffers()
347 : {
348 1 : initializeMulticastMessage();
349 1 : _serviceManager.triggerStopOffers(); // callback stopOffer
350 1 : finalizeMessage();
351 1 : }
352 :
353 0 : void ServiceAnnouncer::cyclic() { checkPendingTasksAndSendDueMessages(); }
354 :
355 0 : void ServiceAnnouncer::sendDueMessages()
356 : {
357 0 : initializeMulticastMessage();
358 0 : addProvidedServices(getSystemTimeMs32Bit());
359 0 : addQueries();
360 0 : finalizeMessage();
361 0 : }
362 :
363 9 : void ServiceAnnouncer::enqueueTxMessage(ServiceAnnouncerTask const& txMessage)
364 : {
365 9 : if (!_taskPool.full())
366 : {
367 9 : ServiceAnnouncerTask& task = *_taskPool.create();
368 9 : task = txMessage;
369 9 : _pendingTxMessages.push_front(task);
370 :
371 9 : if (_taskPool.full())
372 : {
373 0 : checkPendingTasks(getSystemTimeMs32Bit());
374 : }
375 : else
376 : {
377 9 : triggerEventTimeout();
378 : }
379 : }
380 : else
381 : {
382 0 : WARN_LOG(SOMEIP, "ServiceAnnouncer::enqueueTxMessage() task pool is empty");
383 : }
384 9 : }
385 :
386 4 : void ServiceAnnouncer::checkPendingTasks(uint64_t const now)
387 : {
388 4 : tPendingTaskList::iterator itr = _pendingBrowseRequests.begin();
389 4 : tPendingTaskList::iterator prev = _pendingBrowseRequests.before_begin();
390 4 : while (itr != _pendingBrowseRequests.end())
391 : {
392 0 : if (itr->getTimestamp() <= now)
393 : {
394 0 : initializeUnicastMessage(itr->getDestinationAddress());
395 0 : ::ip::IPAddress const currentDestination = itr->getDestinationAddress();
396 0 : processTaskBrowseResults(*itr);
397 0 : releaseTask(*itr);
398 0 : itr = _pendingBrowseRequests.erase_after(prev);
399 :
400 0 : tPendingTaskList::iterator followUpItr = itr;
401 0 : tPendingTaskList::iterator prevFollowUpItr = prev;
402 0 : while ((followUpItr != _pendingBrowseRequests.end())
403 0 : && (followUpItr->getDestinationAddress() == currentDestination))
404 : {
405 0 : if (followUpItr->getTimestamp() <= now)
406 : {
407 0 : processTaskBrowseResults(*followUpItr);
408 0 : releaseTask(*followUpItr);
409 0 : followUpItr = _pendingBrowseRequests.erase_after(prevFollowUpItr);
410 : }
411 : else
412 : {
413 0 : prevFollowUpItr = followUpItr;
414 0 : ++followUpItr;
415 : }
416 0 : prev = prevFollowUpItr;
417 0 : itr = followUpItr;
418 : }
419 0 : finalizeMessage();
420 : }
421 : else
422 : {
423 0 : prev = itr;
424 0 : ++itr;
425 : }
426 : }
427 :
428 7 : while (!_pendingTxMessages.empty())
429 : {
430 3 : ServiceAnnouncerTask& task = _pendingTxMessages.front();
431 3 : initializeUnicastMessage(task.getDestinationAddress());
432 3 : ::ip::IPAddress const currentDestination = task.getDestinationAddress();
433 3 : executeTask(task);
434 3 : _pendingTxMessages.pop_front();
435 3 : releaseTask(task);
436 :
437 3 : while ((!_pendingTxMessages.empty())
438 3 : && (_pendingTxMessages.front().getDestinationAddress() == currentDestination))
439 : {
440 0 : ServiceAnnouncerTask& followUpTask = _pendingTxMessages.front();
441 0 : executeTask(followUpTask);
442 0 : _pendingTxMessages.pop_front();
443 0 : releaseTask(followUpTask);
444 : }
445 3 : finalizeMessage();
446 : }
447 4 : }
448 :
449 0 : void ServiceAnnouncer::checkPendingTasksAndSendDueMessages()
450 : {
451 0 : checkPendingTasks(getSystemTimeMs32Bit());
452 0 : sendDueMessages();
453 0 : }
454 :
455 10 : void ServiceAnnouncer::triggerEventTimeout()
456 : {
457 10 : if (!_isStarted)
458 : {
459 0 : return;
460 : }
461 :
462 10 : _eventTimeout.cancel();
463 10 : async::schedule(
464 10 : _ethernetContext, _eventFunction, _eventTimeout, 1U, ::async::TimeUnit::MILLISECONDS);
465 : }
466 :
467 3 : void ServiceAnnouncer::executeTask(ServiceAnnouncerTask const& task)
468 : {
469 3 : switch (task.getType())
470 : {
471 0 : case ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE:
472 : {
473 0 : processTaskSubscribe(task);
474 0 : break;
475 : }
476 0 : case ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE_ACK:
477 : {
478 0 : processTaskSubscribeAck(task);
479 0 : break;
480 : }
481 0 : case ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE_NACK:
482 : {
483 0 : processTaskSubscribeNack(task);
484 0 : break;
485 : }
486 2 : case ServiceAnnouncerTask::TaskType::TASK_SUBSCRIBE_ACK_MULTICAST:
487 : {
488 2 : processTaskSubscribeAckMulticast(task);
489 2 : break;
490 : }
491 1 : case ServiceAnnouncerTask::TaskType::TASK_UNSUBSCRIBE:
492 : {
493 1 : processTaskUnsubscribe(task);
494 1 : break;
495 : }
496 0 : case ServiceAnnouncerTask::TaskType::TASK_ANNOUNCE:
497 : {
498 0 : processTaskOffer(task);
499 0 : break;
500 : }
501 0 : default:
502 : {
503 0 : WARN_LOG(SOMEIP, "ServiceAnnouncer::executeTask() invalid task %d", task.getType());
504 0 : break;
505 : }
506 : }
507 3 : }
508 :
509 1 : void ServiceAnnouncer::initializeMulticastMessage()
510 : {
511 : /*
512 : * No need to check the result of 'startMessage' if we know that provided message buffer
513 : * is enough to contain at least header at ServiceAnnouncer construction stage.
514 : * Message buffer is declared as a static array.
515 : */
516 1 : (void)_messageBuilder.startMessage(_messageBuffer);
517 1 : _destinationAddress = _network.getMulticastIp();
518 1 : }
519 :
520 3 : void ServiceAnnouncer::initializeUnicastMessage(::ip::IPAddress const& destinationAddress)
521 : {
522 : /*
523 : * No need to check the result of 'startMessage' if we know that provided message buffer
524 : * is enough to contain at least header at ServiceAnnouncer construction stage.
525 : * Message buffer is declared as a static array.
526 : */
527 3 : (void)_messageBuilder.startMessage(_messageBuffer);
528 3 : _destinationAddress = destinationAddress;
529 3 : }
530 :
531 4 : void ServiceAnnouncer::finalizeMessage()
532 : {
533 4 : if (_messageBuilder.isEmpty())
534 : {
535 2 : _messageBuilder.discardMessage();
536 : }
537 : else
538 : {
539 2 : uint16_t sessionId = 0U;
540 2 : bool rebootFlag = false;
541 2 : getSessionInfoForNextMessage(sessionId, rebootFlag);
542 2 : (void)_messageBuilder.finishMessage(sessionId, rebootFlag);
543 :
544 2 : sendMessage();
545 : }
546 4 : etl::fill(_messageBuffer.begin(), _messageBuffer.end(), 0);
547 4 : }
548 :
549 2 : void ServiceAnnouncer::getSessionInfoForNextMessage(uint16_t& sessionId, bool& rebootFlag)
550 : {
551 2 : if (_destinationAddress == _network.getMulticastIp())
552 : {
553 0 : _sessionManager.getSessionInfoForNextMulticastMessage(sessionId, rebootFlag);
554 : }
555 : else
556 : {
557 2 : _sessionManager.getSessionInfoForNextUnicastMessage(
558 2 : _destinationAddress, sessionId, rebootFlag);
559 : }
560 2 : }
561 :
562 0 : void ServiceAnnouncer::processTaskOffer(ServiceAnnouncerTask const& task)
563 : {
564 0 : auto service = ::someip::make<ServiceDescription>();
565 0 : task.copyTo(service);
566 :
567 0 : addOffer(service);
568 0 : }
569 :
570 0 : void ServiceAnnouncer::processTaskBrowseResults(ServiceAnnouncerTask const& task) const
571 : {
572 0 : _serviceManager.triggerOffers(task); // callback offer
573 0 : }
574 :
575 0 : void ServiceAnnouncer::processTaskSubscribe(ServiceAnnouncerTask const& task)
576 : {
577 0 : if (!task.containsEventGroup())
578 : {
579 0 : return;
580 : }
581 :
582 0 : auto service = ::someip::make<ServiceDescription>();
583 0 : task.copyTo(service);
584 :
585 0 : addSubscribe(service);
586 : }
587 :
588 0 : void ServiceAnnouncer::processTaskSubscribeAck(ServiceAnnouncerTask const& task)
589 : {
590 0 : auto service = ::someip::make<ServiceDescription>();
591 0 : task.copyTo(service);
592 :
593 0 : addSubscribeAck(service);
594 0 : }
595 :
596 0 : void ServiceAnnouncer::processTaskSubscribeNack(ServiceAnnouncerTask const& task)
597 : {
598 0 : auto service = ::someip::make<ServiceDescription>();
599 0 : task.copyTo(service);
600 :
601 0 : addSubscribeNack(service);
602 0 : }
603 :
604 2 : void ServiceAnnouncer::processTaskSubscribeAckMulticast(ServiceAnnouncerTask const& task)
605 : {
606 2 : if (!task.containsEventGroup())
607 : {
608 1 : return;
609 : }
610 :
611 1 : auto service = ::someip::make<ServiceDescription>();
612 1 : task.copyTo(service);
613 :
614 1 : addSubscribeAckMulticast(service);
615 : }
616 :
617 1 : void ServiceAnnouncer::processTaskUnsubscribe(ServiceAnnouncerTask const& task)
618 : {
619 1 : if (!task.containsEventGroup())
620 : {
621 0 : return;
622 : }
623 :
624 1 : auto service = ::someip::make<ServiceDescription>();
625 1 : task.copyTo(service);
626 :
627 1 : addUnsubscribe(service);
628 : }
629 :
630 3 : void ServiceAnnouncer::releaseTask(ServiceAnnouncerTask& task)
631 : {
632 3 : task.clear();
633 3 : _taskPool.release(&task);
634 3 : }
635 :
636 0 : void ServiceAnnouncer::addProvidedServices(uint64_t const now)
637 : {
638 0 : _serviceManager.updateServices(now); // callback offer / stopOffer
639 0 : }
640 :
641 0 : void ServiceAnnouncer::addQueries() const
642 : {
643 0 : uint64_t const now = getSystemTimeMs32Bit();
644 0 : _queryManager.updateQueries(now); // callback to ServiceAnnouncer::find()
645 0 : }
646 :
647 0 : void ServiceAnnouncer::find(ServiceDescription const& service)
648 : {
649 : // find service
650 0 : resetIfFull(_messageBuilder.addFind(
651 0 : service.serviceId,
652 0 : service.instanceId,
653 0 : service.majorVersion,
654 0 : service.minorVersion,
655 0 : service.ttl));
656 0 : }
657 :
658 0 : void ServiceAnnouncer::offer(ServiceDescription const& service)
659 : {
660 0 : ServiceDescription temp(service);
661 0 : addOffer(temp);
662 0 : }
663 :
664 0 : void ServiceAnnouncer::addOffer(ServiceDescription& service)
665 : {
666 0 : if (containsEventGroup(service))
667 : {
668 0 : return; // eventgroups are not announced, only services
669 : }
670 :
671 0 : service.ipAddress = _network.getLocalIp();
672 0 : resetIfFull(_messageBuilder.addOffer(
673 0 : service.serviceId,
674 0 : service.instanceId,
675 0 : service.majorVersion,
676 : service.minorVersion,
677 : service.ttl,
678 0 : service.ipAddress,
679 0 : service.port,
680 0 : service.proto));
681 : }
682 :
683 0 : void ServiceAnnouncer::stopOffer(ServiceDescription const& service)
684 : {
685 0 : ServiceDescription temp(service);
686 0 : addStopOffer(temp);
687 0 : }
688 :
689 0 : void ServiceAnnouncer::addStopOffer(ServiceDescription& service)
690 : {
691 0 : if (containsEventGroup(service))
692 : {
693 0 : return; // eventgroups are not announced, only services
694 : }
695 :
696 0 : service.ipAddress = _network.getLocalIp();
697 0 : resetIfFull(_messageBuilder.addDenounce(
698 0 : service.serviceId,
699 0 : service.instanceId,
700 0 : service.majorVersion,
701 : service.minorVersion,
702 : service.ttl,
703 0 : service.ipAddress,
704 0 : service.port,
705 0 : service.proto));
706 : }
707 :
708 0 : void ServiceAnnouncer::addSubscribe(ServiceDescription& service)
709 : {
710 0 : service.ipAddress = _network.getLocalIp();
711 0 : resetIfFull(_messageBuilder.addSubscribe(
712 0 : service.serviceId,
713 0 : service.instanceId,
714 0 : service.eventGroup,
715 0 : service.majorVersion,
716 : service.ttl,
717 0 : service.ipAddress,
718 0 : service.port,
719 0 : service.proto));
720 0 : }
721 :
722 1 : void ServiceAnnouncer::addUnsubscribe(ServiceDescription& service)
723 : {
724 1 : service.ipAddress = _network.getLocalIp();
725 1 : resetIfFull(_messageBuilder.addUnsubscribe(
726 1 : service.serviceId,
727 1 : service.instanceId,
728 1 : service.eventGroup,
729 1 : service.majorVersion,
730 : service.ttl,
731 1 : service.ipAddress,
732 1 : service.port,
733 1 : service.proto));
734 1 : }
735 :
736 0 : void ServiceAnnouncer::addSubscribeAck(ServiceDescription const& service)
737 : {
738 0 : resetIfFull(_messageBuilder.addSubscribeAck(
739 0 : service.serviceId,
740 0 : service.instanceId,
741 0 : service.eventGroup,
742 0 : service.majorVersion,
743 0 : service.minorVersion,
744 0 : service.ttl));
745 0 : }
746 :
747 1 : void ServiceAnnouncer::addSubscribeAckMulticast(ServiceDescription const& service)
748 : {
749 1 : resetIfFull(_messageBuilder.addSubscribeAckMulticast(
750 1 : service.serviceId,
751 1 : service.instanceId,
752 1 : service.eventGroup,
753 1 : service.majorVersion,
754 1 : service.minorVersion,
755 1 : service.ttl,
756 1 : service.ipAddress,
757 1 : service.port,
758 1 : service.proto));
759 1 : }
760 :
761 0 : void ServiceAnnouncer::addSubscribeNack(ServiceDescription const& service)
762 : {
763 0 : resetIfFull(_messageBuilder.addSubscribeNack(
764 0 : service.serviceId,
765 0 : service.instanceId,
766 0 : service.eventGroup,
767 0 : service.majorVersion,
768 0 : service.minorVersion,
769 0 : service.ttl));
770 0 : }
771 :
772 2 : void ServiceAnnouncer::resetIfFull(SdMessageReturnCode const returnCode)
773 : {
774 2 : if (returnCode == SdMessageReturnCode::SD_MESSAGE_IS_FULL)
775 : {
776 0 : finalizeMessage();
777 : /*
778 : * No need to check the result of 'startMessage' if we know that provided message buffer
779 : * is enough to contain at least header at ServiceAnnouncer construction stage.
780 : * Message buffer is declared as a static array.
781 : */
782 0 : (void)_messageBuilder.startMessage(_messageBuffer);
783 : }
784 2 : }
785 :
786 2 : void ServiceAnnouncer::sendMessage()
787 : {
788 2 : auto const portResult = _network.getSdPort();
789 :
790 2 : if (!portResult.has_value())
791 : {
792 0 : WARN_LOG(SOMEIP, "ServiceAnnouncer::sendMessage() no SD port available");
793 0 : return;
794 : }
795 :
796 2 : uint16_t const port = portResult.value();
797 :
798 2 : auto channel = _network.getSdChannel(port, ::ip::IPEndpoint(_destinationAddress, port));
799 :
800 2 : if (!channel.has_value())
801 : {
802 2 : WARN_LOG(SOMEIP, "ServiceAnnouncer::sendMessage() no channel");
803 2 : return;
804 : }
805 :
806 0 : uint32_t const length = readTotalLength(_messageBuffer);
807 0 : memcpy(channel->getOutputBuffer().data(), _messageBuffer.data(), length);
808 :
809 0 : if (!channel->send(length))
810 : {
811 0 : WARN_LOG(SOMEIP, "ServiceAnnouncer::sendMessage() send failed");
812 0 : return;
813 : }
814 :
815 0 : Statistics::incCounter(Statistics::Counter::SD_FRAME_TX);
816 2 : }
817 :
818 : } // namespace someip
819 :
820 : // NOLINTEND(cppcoreguidelines-pro-type-vararg)
|