96 std::ostringstream msg {};
99 msg <<
"value=" << *discarded_data.
GetValue();
100 if (oldest_retained_data && !oldest_retained_data->
GetValue()) {
102 msg <<
" (transferred to retained newer data)";
109 if (oldest_retained_data && !oldest_retained_data->
GetTimestamp()) {
111 msg <<
" (transferred to retained newer data)";
118 msg <<
"quality=" << quality_str;
122 LOG4CPLUS_DEBUG(this->logger,
"OLDB async (" << this->name <<
") discarded data: " << msg.str());
127 discarded_data.
GetPromise()->set_exception(boost::copy_exception(std::out_of_range{
"Dropped data."}));
129 std::lock_guard lck {this->m_sync_dtor_mutex};
130 this->m_sync_dtor_n_writing--;
148 std::unique_lock lck {m_sync_dtor_mutex};
154 m_sync_dtor_cv.wait(lck, [
this]{
155 LOG4CPLUS_TRACE(logger,
"~CiiOldbDataPointAsync (" << name <<
") m_sync_dtor_cv.wait predicate function: "
156 <<
"m_sync_dtor_n_writing=" << m_sync_dtor_n_writing);
157 return (m_sync_dtor_n_writing == 0);
160 LOG4CPLUS_TRACE(logger,
"~CiiOldbDataPointAsync (" << name <<
") is done.");
166 return delegate->ReadValue(check_bad_quality);
178 const T& value, int64_t timestamp, elt::oldb::CiiOldbDpQuality quality,
bool is_disable_publishing) {
183 std::lock_guard lck {m_sync_dtor_mutex};
184 m_sync_dtor_n_writing++;
188 return CiiOldbDataPointAsync<T>::WriteAsync(std::move(new_data), is_disable_publishing);
194 elt::oldb::CiiOldbDpQuality quality,
bool is_disable_publishing) {
197 std::lock_guard lck {m_sync_dtor_mutex};
198 m_sync_dtor_n_writing++;
201 return CiiOldbDataPointAsync<T>::WriteAsync(std::move(new_data), is_disable_publishing);
206boost::future<typename CiiOldbDataPointAsync<T>::OldbData> CiiOldbDataPointAsync<T>::WriteAsync(
207 OldbData&& new_data,
bool is_disable_publishing) {
213 std::chrono::steady_clock::time_point t_0 {std::chrono::steady_clock::now()};
223 boost::promise<OldbData> promise {};
224 boost::future<OldbData> future = promise.get_future();
228 std::chrono::steady_clock::time_point t_1 {std::chrono::steady_clock::now()};
232 new_data.GetValue(), new_data.GetTimestamp(), new_data.GetQuality(), std::move(promise)};
238 buffer.Push(std::move(new_data_with_promise));
240 std::chrono::steady_clock::time_point t_2 {std::chrono::steady_clock::now()};
242 if (new_data.GetValue()) {
243 LOG4CPLUS_TRACE(logger,
"OLDB async (" << name <<
") wrote to buffer: "
244 << *new_data_with_promise.GetValue() <<
", threadId=" << std::this_thread::get_id());
246 LOG4CPLUS_TRACE(logger,
"OLDB async (" << name <<
") wrote to buffer: <no value>"
247 <<
", threadId=" << std::this_thread::get_id());
251 bool expected =
false;
253 std::chrono::steady_clock::time_point t_3 {};
254 std::chrono::steady_clock::time_point t_4 {};
257 if (is_async_processing.compare_exchange_strong(expected,
true)) {
259 LOG4CPLUS_TRACE(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name << LOG4CPLUS_TEXT(
") will use a writer thread."));
289 t_3 = std::chrono::steady_clock::now();
292 boost::asio::post(async_exec, [
this](){ WriteBufferToOldb(); });
294 t_4 = std::chrono::steady_clock::now();
296 t_3 = t_4 = std::chrono::steady_clock::now();
299 LOG4CPLUS_TRACE(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name << LOG4CPLUS_TEXT(
") will rely on running writer thread."));
302 LOG4CPLUS_TRACE(logger,
"OLDB async (" << name <<
") Returning from WriteAsync");
304 int d1 = std::chrono::duration_cast<std::chrono::microseconds>(t_1 - t_0).count();
305 int d2 = std::chrono::duration_cast<std::chrono::microseconds>(t_2 - t_1).count();
306 int d3 = std::chrono::duration_cast<std::chrono::microseconds>(t_3 - t_2).count();
307 int d4 = std::chrono::duration_cast<std::chrono::microseconds>(t_4 - t_3).count();
309 if (d1 + d2 + d3 + d4 > 1000) {
310 LOG4CPLUS_DEBUG(logger,
"CiiOldbDataPointAsync " << name <<
" WriteAsync steps in micros: "
311 << d1 <<
", " << d2 <<
", " << d3 <<
", " << d4);;
318void CiiOldbDataPointAsync<T>::WriteBufferToOldb() {
319 LOG4CPLUS_TRACE(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name << LOG4CPLUS_TEXT(
") writer thread started. is_async_processing=") << is_async_processing <<
320 ", threadId=" << std::this_thread::get_id());
326 std::optional<OldbDataWithPromise> data_to_write_opt {buffer.Poll()};
328 if (!data_to_write_opt) {
333 is_async_processing =
false;
337 while (data_to_write_opt.has_value()) {
339 LOG4CPLUS_TRACE(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name <<
") thread got data and will call the sync writeValue.");
341 OldbDataWithPromise data_to_write {std::move(data_to_write_opt.value())};
344 auto start = std::chrono::steady_clock::now();
347 if (data_to_write.GetValue() && data_to_write.GetTimestamp() && data_to_write.GetQuality()) {
348 delegate->WriteValue(*data_to_write.GetValue(), *data_to_write.GetTimestamp(), *data_to_write.GetQuality(),
false);
349 LOG4CPLUS_TRACE(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name <<
350 LOG4CPLUS_TEXT(
") wrote DP to OLDB in ") << std::chrono::duration_cast<std::chrono::milliseconds>((std::chrono::steady_clock::now() - start)).count() <<
351 " ms, data=" << LOG4CPLUS_TEXT(*data_to_write.GetValue()) );
353 else if (data_to_write.GetQuality()) {
354 delegate->SetQuality(*data_to_write.GetQuality());
355 LOG4CPLUS_TRACE(logger,
"OLDB async (" << name <<
356 ") wrote DP quality to OLDB in " << std::chrono::duration_cast<std::chrono::milliseconds>((std::chrono::steady_clock::now() - start)).count() <<
360 LOG4CPLUS_DEBUG(logger, LOG4CPLUS_TEXT(
"Programming error: OLDB async (") << name <<
361 ") could not write DP because of missing data in OldbDataWithPromise");
364 data_to_write.GetPromise()->set_value(std::move(data_to_write));
366 }
catch (
const elt::oldb::CiiOldbException& ex) {
367 LOG4CPLUS_DEBUG(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name <<
") writing to OLDB failed: " << ex.what());
370 data_to_write.GetPromise()->set_exception(boost::copy_exception(ex));
374 LOG4CPLUS_WARN(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name <<
") unknown exception.");
381 auto lock = buffer.Lock();
382 data_to_write_opt = buffer.Poll();
383 if (!data_to_write_opt.has_value()) {
384 is_async_processing.store(
false);
389 LOG4CPLUS_TRACE(logger, LOG4CPLUS_TEXT(
"OLDB async (") << name <<
390 LOG4CPLUS_TEXT(
") performed ") << write_count << LOG4CPLUS_TEXT(
" DP writes in the background, now releasing writer thread."));
396 std::lock_guard lck {m_sync_dtor_mutex};
397 m_sync_dtor_n_writing -= write_count;
398 LOG4CPLUS_TRACE(logger,
"OLDB async (" << name <<
") WriteBufferToOldb done. m_sync_dtor_n_writing=" << m_sync_dtor_n_writing);
400 m_sync_dtor_cv.notify_all();