score/message_passing/client_connection.h
Line | Count | Source |
1 | | /******************************************************************************** |
2 | | * Copyright (c) 2025 Contributors to the Eclipse Foundation |
3 | | * |
4 | | * See the NOTICE file(s) distributed with this work for additional |
5 | | * information regarding copyright ownership. |
6 | | * |
7 | | * This program and the accompanying materials are made available under the |
8 | | * terms of the Apache License Version 2.0 which is available at |
9 | | * https://www.apache.org/licenses/LICENSE-2.0 |
10 | | * |
11 | | * SPDX-License-Identifier: Apache-2.0 |
12 | | ********************************************************************************/ |
13 | | #ifndef SCORE_LIB_MESSAGE_PASSING_CLIENT_CONNECTION_H |
14 | | #define SCORE_LIB_MESSAGE_PASSING_CLIENT_CONNECTION_H |
15 | | |
16 | | #include "score/message_passing/i_client_connection.h" |
17 | | #include "score/message_passing/i_client_factory.h" |
18 | | #include "score/message_passing/i_shared_resource_engine.h" |
19 | | |
20 | | #include <score/string.hpp> |
21 | | #include <score/vector.hpp> |
22 | | |
23 | | #include <atomic> |
24 | | #include <condition_variable> |
25 | | #include <mutex> |
26 | | #include <optional> |
27 | | |
28 | | namespace score::message_passing::detail |
29 | | { |
30 | | |
31 | | class ClientConnection final : public IClientConnection |
32 | | { |
33 | | public: |
34 | | ClientConnection(std::shared_ptr<ISharedResourceEngine> engine, |
35 | | const ServiceProtocolConfig& protocol_config, |
36 | | const IClientFactory::ClientConfig& client_config) noexcept; |
37 | | ~ClientConnection() noexcept override; |
38 | | |
39 | | ClientConnection(const ClientConnection&) = delete; |
40 | | ClientConnection(ClientConnection&&) = delete; |
41 | | ClientConnection& operator=(const ClientConnection&) = delete; |
42 | | ClientConnection& operator=(ClientConnection&&) = delete; |
43 | | |
44 | | score::cpp::expected_blank<score::os::Error> Send(score::cpp::span<const std::uint8_t> message) noexcept override; |
45 | | |
46 | | score::cpp::expected<score::cpp::span<const std::uint8_t>, score::os::Error> SendWaitReply( |
47 | | score::cpp::span<const std::uint8_t> message, |
48 | | score::cpp::span<std::uint8_t> reply) noexcept override; |
49 | | |
50 | | score::cpp::expected_blank<score::os::Error> SendWithCallback(score::cpp::span<const std::uint8_t> message, |
51 | | ReplyCallback callback) noexcept override; |
52 | | |
53 | | State GetState() const noexcept override; |
54 | | |
55 | | StopReason GetStopReason() const noexcept override; |
56 | | |
57 | | void Start(StateCallback state_callback, NotifyCallback notify_callback) noexcept override; |
58 | | |
59 | | void Stop() noexcept override; |
60 | | |
61 | | void Restart() noexcept override; |
62 | | |
63 | | private: |
64 | | void TryConnect() noexcept; |
65 | | bool TryQueueMessage(score::cpp::span<const std::uint8_t> message, ReplyCallback callback) noexcept; |
66 | | StopReason ProcessInputEvent() noexcept; |
67 | | |
68 | | // The lock shall be already taken. |
69 | | // The function may release it, call a user callback, and then lock it again. |
70 | | void ProcessSendQueueUnderLock(std::unique_lock<std::mutex>& lock) noexcept; |
71 | | void ArmSendQueueUnderLock() noexcept; |
72 | | |
73 | | bool TrySetStopReason(const StopReason stop_reason) noexcept; |
74 | | |
75 | | void ProcessStateChangeToStopped() noexcept; |
76 | | void ProcessStateChange(const State state) noexcept; |
77 | | void SwitchToStopState() noexcept; |
78 | | |
79 | | bool IsInCallback() const noexcept |
80 | 109 | { |
81 | 109 | return engine_->IsOnCallbackThread(); |
82 | 109 | } |
83 | | |
84 | | void DoRestart() noexcept; |
85 | | |
86 | | const std::shared_ptr<ISharedResourceEngine> engine_; |
87 | | const score::cpp::pmr::string identifier_; |
88 | | const std::uint32_t max_send_size_; |
89 | | const std::uint32_t max_receive_size_; |
90 | | const IClientFactory::ClientConfig client_config_; |
91 | | |
92 | | std::int32_t client_fd_; |
93 | | std::atomic<State> state_; |
94 | | std::atomic<StopReason> stop_reason_; |
95 | | |
96 | | // to detach and survive the destructor, if needed for stopping |
97 | | struct CallbackContext |
98 | | { |
99 | | StateCallback state_callback; |
100 | | std::recursive_mutex finalize_mutex; |
101 | | }; |
102 | | std::shared_ptr<CallbackContext> callback_context_; |
103 | | |
104 | | NotifyCallback notify_callback_; |
105 | | |
106 | | std::int32_t connect_retry_ms_; |
107 | | |
108 | | std::mutex send_mutex_; |
109 | | std::condition_variable send_condition_; |
110 | | |
111 | | // at the time of construction of the connection object, we preallocate the storage for the amount of messages |
112 | | // requested in the client_config to be sent asynchronously. The send_storage_ container is responsible for managing |
113 | | // the lifetime of these message objects. In addition, we arrange a list of the currently unused message objects |
114 | | // in send_pool_, relying on the allocation-free nature of intrusive lists. When we have a message to send |
115 | | // asynchronously, we borrow an element from the send_pool and put it into the send_queue_, which is also an |
116 | | // intrusive list container. When we send this message later, we return the freed element back to the send_pool_. |
117 | | // Thus, we have no extra memory allocation after creation of a ClientConnection object. |
118 | | class SendCommand : public score::containers::intrusive_list_element<> |
119 | | { |
120 | | public: |
121 | | using allocator_type = score::cpp::pmr::polymorphic_allocator<SendCommand>; |
122 | | explicit SendCommand(const allocator_type& allocator) |
123 | 74 | : score::containers::intrusive_list_element<>{}, message(allocator), callback{} |
124 | 74 | { |
125 | 74 | } |
126 | | |
127 | | score::cpp::pmr::vector<std::uint8_t> message; |
128 | | ReplyCallback callback; |
129 | | }; |
130 | | score::cpp::pmr::vector<SendCommand> send_storage_; |
131 | | score::containers::intrusive_list<SendCommand> send_pool_; |
132 | | score::containers::intrusive_list<SendCommand> send_queue_; |
133 | | |
134 | | std::optional<ReplyCallback> waiting_for_reply_; |
135 | | |
136 | | ISharedResourceEngine::CommandQueueEntry connection_timer_; |
137 | | ISharedResourceEngine::CommandQueueEntry disconnection_command_; |
138 | | ISharedResourceEngine::CommandQueueEntry async_send_command_; |
139 | | ISharedResourceEngine::PosixEndpointEntry posix_endpoint_; |
140 | | }; |
141 | | |
142 | | } // namespace score::message_passing::detail |
143 | | |
144 | | #endif // SCORE_LIB_MESSAGE_PASSING_CLIENT_CONNECTION_H |