3 #include "../../data_interfaces/source/ifdv2_synchronizer.hpp"
5 #include <launchdarkly/async/cancellation.hpp>
6 #include <launchdarkly/async/promise.hpp>
7 #include <launchdarkly/fdv2_protocol_handler.hpp>
8 #include <launchdarkly/logging/logger.hpp>
9 #include <launchdarkly/server_side/config/built/all_built.hpp>
10 #include <launchdarkly/sse/client.hpp>
12 #include <boost/asio/any_io_executor.hpp>
13 #include <boost/url/url.hpp>
22 namespace launchdarkly::server_side::data_systems {
24 class FDv2StreamingSynchronizerTestPeer;
39 friend class FDv2StreamingSynchronizerTestPeer;
47 boost::asio::any_io_executor
const& executor,
49 std::string streaming_base_url,
51 std::optional<std::string> filter_key,
52 std::chrono::milliseconds initial_reconnect_delay);
56 async::Future<data_interfaces::FDv2SourceResult>
Next(
57 data_model::Selector selector)
override;
59 void Close()
override;
61 [[nodiscard]] std::string
const&
Identity()
const override;
68 friend class FDv2StreamingSynchronizerTestPeer;
72 boost::asio::any_io_executor
const& executor,
73 std::string streaming_base_url,
75 std::optional<std::string> filter_key,
76 std::chrono::milliseconds initial_reconnect_delay);
88 async::Future<data_interfaces::FDv2SourceResult>
Next(
89 data_model::Selector
const& selector,
90 std::shared_ptr<State>
self);
97 void ClearPendingPromise();
108 boost::beast::http::request<boost::beast::http::string_body>;
109 using HttpResponseHeader = boost::beast::http::response_header<>;
116 void EnsureStarted(data_model::Selector
const& selector,
117 std::shared_ptr<State>
self);
126 void OnConnect(HttpRequest* req);
127 void OnResponse(HttpResponseHeader
const& headers);
128 void OnEvent(sse::Event
const& event);
129 void OnError(sse::Error
const& error);
135 std::string
const streaming_base_url_;
137 std::optional<std::string>
const filter_key_;
138 std::chrono::milliseconds
const initial_reconnect_delay_;
139 boost::asio::any_io_executor
const executor_;
143 FDv2ProtocolHandler protocol_handler_;
147 bool started_ =
false;
148 bool closed_ =
false;
150 std::optional<data_interfaces::FDv1FallbackDirective>
151 latest_fdv1_fallback_;
152 data_model::Selector latest_selector_;
153 std::optional<boost::urls::url> base_url_;
154 std::shared_ptr<sse::Client> sse_client_;
155 std::optional<async::Promise<data_interfaces::FDv2SourceResult>>
157 std::deque<data_interfaces::FDv2SourceResult> result_queue_;
161 async::Promise<std::monostate> close_promise_;
164 std::shared_ptr<State> state_;
Definition: http_properties.hpp:70
Definition: ifdv2_synchronizer.hpp:19
Definition: streaming_synchronizer.hpp:38
std::string const & Identity() const override
Definition: streaming_synchronizer.cpp:435
void Close() override
Definition: streaming_synchronizer.cpp:428
async::Future< data_interfaces::FDv2SourceResult > Next(data_model::Selector selector) override
Definition: streaming_synchronizer.cpp:402
FDv2StreamingSynchronizer(boost::asio::any_io_executor const &executor, Logger const &logger, std::string streaming_base_url, config::built::HttpProperties const &http_properties, std::optional< std::string > filter_key, std::chrono::milliseconds initial_reconnect_delay)
Definition: streaming_synchronizer.cpp:384
Definition: fdv2_source_result.hpp:56