Skip to content

Commit 878ca36

Browse files
committed
feat: add altimeter wal_started/wal_stored/wal_shipped tests
1 parent 942220f commit 878ca36

9 files changed

Lines changed: 651 additions & 39 deletions

File tree

docs/internal/20260119-altimeter-event-log-integration.md

Lines changed: 12 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -40,24 +40,13 @@ https://github.com/project-tsurugi/altimeter/blob/feature/multi_nodes/docs/ja/al
4040
### 出力位置
4141

4242
- wal_stored
43-
- `datastore::update_min_epoch_id` 関数内の次の行の直後に追加する
44-
```cpp
45-
TRACE_FINE << "epoch_id_record_finished_ updated to " << to_be_epoch;
46-
```
43+
- `datastore::persist_epoch_id` 内(成功/失敗の両方で出力)
4744
- wal_shipped
48-
- `datastore::persist_and_propagate_epoch_id` 関数内の次の行の直後に追加する。2箇所存在する。
49-
```cpp
50-
bool sent = impl_->propagate_group_commit(epoch_id);
51-
```
45+
- `datastore_impl::propagate_group_commit` 内(成功/失敗の両方で出力)
5246
- wal_received
53-
- `message_group_commit::post_receive` 関数内の次の行の直後に追加する
54-
```cpp
55-
TRACE_START << "epoch_number: " << epoch_number_;
56-
```
47+
- 未実装(`message_group_commit::post_receive` に追加予定)
5748
- wal_started
58-
- `datastore::ready() ` 関数内の次の行の直後に追加する
59-
```cpp
60-
blob_id_type max_blob_id = std::max(create_snapshot_and_get_max_blob_id(), compaction_catalog_->get_max_blob_id());
49+
- `datastore::create_snapshot_and_get_max_blob_id_with_wal_started_log` 内(成功/失敗の両方で出力)
6150
```
6251

6352
## ログ出力項目と取得方法(wal_*
@@ -78,20 +67,20 @@ Altimeterの仕様上、`wal_stored / wal_shipped / wal_received / wal_started`
7867
### イベント別メモ
7968

8069
- `wal_stored`
81-
- 位置: `datastore::update_min_epoch_id`
82-
- `result`: `write_epoch_callback_` が例外なく完了したら成功
70+
- 位置: `datastore::persist_epoch_id`
71+
- `result`: `persist_epoch_id` の処理が例外なく完了したら成功
8372

8473
- `wal_shipped`
85-
- 位置: `datastore::persist_and_propagate_epoch_id`
86-
- `result`: `impl_->propagate_group_commit(epoch_id)` の戻り値で判断
74+
- 位置: `datastore_impl::propagate_group_commit`
75+
- `result`: `propagate_group_commit` の戻り値で判断
8776

8877
- `wal_received`
89-
- 位置: `message_group_commit::post_receive`
90-
- `result`: 受信処理が例外なく完了したら成功
78+
- 位置: `message_group_commit::post_receive`(未実装)
79+
- `result`: 受信処理が例外なく完了したら成功(予定)
9180

9281
- `wal_started`
93-
- 位置: `datastore::ready`
94-
- `result`: 初期化処理が例外なく完了したら成功
82+
- 位置: `datastore::create_snapshot_and_get_max_blob_id_with_wal_started_log`
83+
- `result`: `create_snapshot_and_get_max_blob_id` が例外なく完了したら成功
9584

9685

9786
### 出力項目

include/limestone/api/datastore.h

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -443,6 +443,17 @@ class datastore {
443443

444444
private:
445445
void persist_epoch_id(epoch_id_type epoch_id);
446+
/**
447+
* @brief Log that the write-ahead log (WAL) has been started.
448+
* @param wal_version the WAL version (epoch) associated with the start event.
449+
* @param success whether starting the WAL succeeded (true) or failed (false).
450+
*/
451+
void log_wal_started(epoch_id_type wal_version, bool success) const;
452+
/**
453+
* @brief Create a snapshot and record WAL started log, then return the maximum blob ID.
454+
* @return The maximum blob ID observed while creating the snapshot.
455+
*/
456+
blob_id_type create_snapshot_and_get_max_blob_id_with_wal_started_log();
446457

447458
std::function<void(epoch_id_type)> write_epoch_callback_{
448459
[this](epoch_id_type epoch) { this->persist_and_propagate_epoch_id(epoch); }

src/limestone/datastore.cpp

Lines changed: 93 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,11 @@
2424

2525
#include <glog/logging.h>
2626
#include <limestone/logging.h>
27+
#ifdef ENABLE_ALTIMETER
28+
#include <altimeter/event/constants.h>
29+
#include <altimeter/log_item.h>
30+
#include <altimeter/logger.h>
31+
#endif
2732
#include "logging_helper.h"
2833
#include "limestone_exception_helper.h"
2934

@@ -200,26 +205,99 @@ void datastore::recover() const noexcept {
200205

201206
void datastore::persist_epoch_id(epoch_id_type epoch_id) {
202207
TRACE_START << "epoch_id=" << epoch_id;
203-
if (++epoch_write_counter >= max_entries_in_epoch_file) {
204-
write_epoch_to_file_internal(tmp_epoch_file_path_.string(), epoch_id, file_write_mode::overwrite);
208+
try {
209+
if (++epoch_write_counter >= max_entries_in_epoch_file) {
210+
write_epoch_to_file_internal(tmp_epoch_file_path_.string(), epoch_id, file_write_mode::overwrite);
205211

206-
boost::system::error_code ec;
207-
if (::rename(tmp_epoch_file_path_.c_str(), epoch_file_path_.c_str()) != 0) {
208-
TRACE_ABORT;
209-
LOG_AND_THROW_IO_EXCEPTION("Failed to rename temp file: " + tmp_epoch_file_path_.string() + " to " + epoch_file_path_.string(), errno);
212+
boost::system::error_code ec;
213+
if (::rename(tmp_epoch_file_path_.c_str(), epoch_file_path_.c_str()) != 0) {
214+
TRACE_ABORT;
215+
LOG_AND_THROW_IO_EXCEPTION("Failed to rename temp file: " + tmp_epoch_file_path_.string() + " to " + epoch_file_path_.string(), errno);
216+
}
217+
boost::filesystem::remove(tmp_epoch_file_path_, ec);
218+
if (ec) {
219+
TRACE_ABORT;
220+
LOG_AND_THROW_IO_EXCEPTION("Failed to remove temp file: " + tmp_epoch_file_path_.string(), ec);
221+
}
222+
epoch_write_counter = 0;
223+
} else {
224+
write_epoch_to_file_internal(epoch_file_path_.string(), epoch_id, file_write_mode::append);
210225
}
211-
boost::filesystem::remove(tmp_epoch_file_path_, ec);
212-
if (ec) {
213-
TRACE_ABORT;
214-
LOG_AND_THROW_IO_EXCEPTION("Failed to remove temp file: " + tmp_epoch_file_path_.string(), ec);
226+
#ifdef ENABLE_ALTIMETER
227+
if (::altimeter::logger::is_log_on(::altimeter::event::category,
228+
::altimeter::event::level::log_data_store)) {
229+
::altimeter::log_item log_item;
230+
log_item.category(::altimeter::event::category);
231+
log_item.type(::altimeter::event::type::wal_stored);
232+
log_item.level(::altimeter::event::level::log_data_store);
233+
log_item.add(::altimeter::event::item::instance_id, impl_->instance_id());
234+
log_item.add(::altimeter::event::item::dbname, impl_->db_name());
235+
log_item.add(::altimeter::event::item::pid, static_cast<std::int64_t>(impl_->pid()));
236+
std::string wal_version = std::to_string(epoch_id);
237+
log_item.add(::altimeter::event::item::wal_version, wal_version);
238+
log_item.add(::altimeter::event::item::result, ::altimeter::event::result::success);
239+
::altimeter::logger::log(log_item);
215240
}
216-
epoch_write_counter = 0;
217-
} else {
218-
write_epoch_to_file_internal(epoch_file_path_.string(), epoch_id, file_write_mode::append);
241+
#endif
242+
} catch (...) {
243+
#ifdef ENABLE_ALTIMETER
244+
if (::altimeter::logger::is_log_on(::altimeter::event::category,
245+
::altimeter::event::level::log_data_store)) {
246+
::altimeter::log_item log_item;
247+
log_item.category(::altimeter::event::category);
248+
log_item.type(::altimeter::event::type::wal_stored);
249+
log_item.level(::altimeter::event::level::log_data_store);
250+
log_item.add(::altimeter::event::item::instance_id, impl_->instance_id());
251+
log_item.add(::altimeter::event::item::dbname, impl_->db_name());
252+
log_item.add(::altimeter::event::item::pid, static_cast<std::int64_t>(impl_->pid()));
253+
std::string wal_version = std::to_string(epoch_id);
254+
log_item.add(::altimeter::event::item::wal_version, wal_version);
255+
log_item.add(::altimeter::event::item::result, ::altimeter::event::result::failure);
256+
::altimeter::logger::log(log_item);
257+
}
258+
#endif
259+
throw;
219260
}
220261
TRACE_END;
221262
}
222263

264+
void datastore::log_wal_started(epoch_id_type wal_version, bool success) const {
265+
#ifdef ENABLE_ALTIMETER
266+
if (!::altimeter::logger::is_log_on(::altimeter::event::category,
267+
::altimeter::event::level::log_data_store)) {
268+
return;
269+
}
270+
::altimeter::log_item log_item;
271+
log_item.category(::altimeter::event::category);
272+
log_item.type(::altimeter::event::type::wal_started);
273+
log_item.level(::altimeter::event::level::log_data_store);
274+
log_item.add(::altimeter::event::item::instance_id, impl_->instance_id());
275+
log_item.add(::altimeter::event::item::dbname, impl_->db_name());
276+
log_item.add(::altimeter::event::item::pid, static_cast<std::int64_t>(impl_->pid()));
277+
std::string wal_version_str = std::to_string(wal_version);
278+
log_item.add(::altimeter::event::item::wal_version, wal_version_str);
279+
log_item.add(::altimeter::event::item::result,
280+
success ? ::altimeter::event::result::success : ::altimeter::event::result::failure);
281+
::altimeter::logger::log(log_item);
282+
#else
283+
(void)wal_version;
284+
(void)success;
285+
#endif
286+
}
287+
288+
blob_id_type datastore::create_snapshot_and_get_max_blob_id_with_wal_started_log() {
289+
try {
290+
auto max_blob_id = create_snapshot_and_get_max_blob_id();
291+
log_wal_started(static_cast<epoch_id_type>(epoch_id_informed_.load()), true);
292+
return max_blob_id;
293+
} catch (...) {
294+
// NOTE: This path may not appear in coverage because death tests run in a separate process,
295+
// and their coverage data is not merged into the parent process report.
296+
log_wal_started(static_cast<epoch_id_type>(epoch_id_informed_.load()), false);
297+
throw;
298+
}
299+
}
300+
223301
void datastore::persist_and_propagate_epoch_id(epoch_id_type epoch_id) {
224302
TRACE_START << "epoch_id=" << epoch_id;
225303
if (impl_->is_async_group_commit_enabled()) {
@@ -247,7 +325,8 @@ blob_reference_tag_type datastore::generate_reference_tag(
247325
void datastore::ready() {
248326
TRACE_START;
249327
try {
250-
blob_id_type max_blob_id = std::max(create_snapshot_and_get_max_blob_id(), compaction_catalog_->get_max_blob_id());
328+
blob_id_type max_blob_id =
329+
std::max(create_snapshot_and_get_max_blob_id_with_wal_started_log(), compaction_catalog_->get_max_blob_id());
251330
blob_file_garbage_collector_ = std::make_unique<blob_file_garbage_collector>(*blob_file_resolver_);
252331
blob_file_garbage_collector_->scan_blob_files(max_blob_id);
253332

src/limestone/datastore_impl.cpp

Lines changed: 56 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,10 +20,16 @@
2020
#include <iostream>
2121
#include <cstring>
2222
#include <stdexcept>
23+
#include <string>
2324
#include <openssl/hmac.h>
2425
#include <openssl/rand.h>
2526
#include <openssl/evp.h>
2627
#include <openssl/err.h>
28+
#ifdef ENABLE_ALTIMETER
29+
#include <altimeter/event/constants.h>
30+
#include <altimeter/log_item.h>
31+
#include <altimeter/logger.h>
32+
#endif
2733

2834
#include <replication/replica_connector.h>
2935
#include <limestone_exception_helper.h>
@@ -139,13 +145,57 @@ bool datastore_impl::propagate_group_commit(uint64_t epoch_id) {
139145
}
140146
if (replica_exists_.load(std::memory_order_acquire)) {
141147
TRACE_START << "epoch_id=" << epoch_id;
142-
message_group_commit message{epoch_id};
143-
if (!control_channel_->send_message(message)) {
148+
bool sent = false;
149+
if (group_commit_sender_for_tests_) {
150+
sent = group_commit_sender_for_tests_(epoch_id);
151+
} else {
152+
if (!control_channel_) {
153+
LOG_LP(ERROR) << "Control channel is not initialized.";
154+
TRACE_END << "Failed to send group commit message.";
155+
sent = false;
156+
} else {
157+
message_group_commit message{epoch_id};
158+
sent = control_channel_->send_message(message);
159+
}
160+
}
161+
if (!sent) {
144162
LOG_LP(ERROR) << "Failed to send group commit message to replica.";
145163
TRACE_END << "Failed to send group commit message.";
164+
#ifdef ENABLE_ALTIMETER
165+
if (::altimeter::logger::is_log_on(::altimeter::event::category,
166+
::altimeter::event::level::log_data_store)) {
167+
::altimeter::log_item log_item;
168+
log_item.category(::altimeter::event::category);
169+
log_item.type(::altimeter::event::type::wal_shipped);
170+
log_item.level(::altimeter::event::level::log_data_store);
171+
log_item.add(::altimeter::event::item::instance_id, instance_id_);
172+
log_item.add(::altimeter::event::item::dbname, db_name_);
173+
log_item.add(::altimeter::event::item::pid, static_cast<std::int64_t>(pid_));
174+
std::string wal_version = std::to_string(epoch_id);
175+
log_item.add(::altimeter::event::item::wal_version, wal_version);
176+
log_item.add(::altimeter::event::item::result, ::altimeter::event::result::failure);
177+
::altimeter::logger::log(log_item);
178+
}
179+
#endif
146180
return false;
147181
}
148182
TRACE_END;
183+
#ifdef ENABLE_ALTIMETER
184+
if (::altimeter::logger::is_log_on(::altimeter::event::category,
185+
::altimeter::event::level::log_data_store)) {
186+
::altimeter::log_item log_item;
187+
log_item.category(::altimeter::event::category);
188+
log_item.type(::altimeter::event::type::wal_shipped);
189+
log_item.level(::altimeter::event::level::log_data_store);
190+
log_item.add(::altimeter::event::item::instance_id, instance_id_);
191+
log_item.add(::altimeter::event::item::dbname, db_name_);
192+
log_item.add(::altimeter::event::item::pid, static_cast<std::int64_t>(pid_));
193+
std::string wal_version = std::to_string(epoch_id);
194+
log_item.add(::altimeter::event::item::wal_version, wal_version);
195+
log_item.add(::altimeter::event::item::result, ::altimeter::event::result::success);
196+
::altimeter::logger::log(log_item);
197+
}
198+
#endif
149199
return true;
150200
}
151201
return false;
@@ -226,6 +276,10 @@ bool datastore_impl::is_async_group_commit_enabled() const noexcept {
226276
return async_group_commit_enabled_;
227277
}
228278

279+
void datastore_impl::set_group_commit_sender_for_tests(std::function<bool(uint64_t)> const& sender) {
280+
group_commit_sender_for_tests_ = sender;
281+
}
282+
229283
void datastore_impl::set_instance_id(std::string_view instance_id) {
230284
instance_id_ = instance_id;
231285
}

src/limestone/datastore_impl.h

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
#include <array>
2323
#include <memory>
2424
#include <optional>
25+
#include <functional>
2526
#include <string>
2627
#include <sys/types.h>
2728

@@ -160,6 +161,13 @@ class datastore_impl {
160161
*/
161162
[[nodiscard]] pid_t pid() const noexcept;
162163

164+
/**
165+
* @brief Sets a custom group commit sender for tests.
166+
* @param sender The sender function(epoch_id) used to simulate group commit sending.
167+
* The function must return true on success and false on failure.
168+
*/
169+
void set_group_commit_sender_for_tests(std::function<bool(uint64_t)> const& sender);
170+
163171
private:
164172
// Atomic counter for tracking active backup operations.
165173
std::atomic<int> backup_counter_;
@@ -187,6 +195,7 @@ class datastore_impl {
187195
std::string instance_id_{"instance_id_not_set"};
188196
std::string db_name_{"db_name_not_set"};
189197
pid_t pid_{0};
198+
std::function<bool(uint64_t)> group_commit_sender_for_tests_{};
190199

191200
/**
192201
* @brief generates HMAC secret key for BLOB reference tag generation.

test/CMakeLists.txt

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,12 @@ target_link_libraries(${test_target}
1717
PRIVATE gflags::gflags
1818
)
1919

20+
if (ENABLE_ALTIMETER)
21+
target_link_libraries(${test_target}
22+
PRIVATE altimeter
23+
)
24+
endif()
25+
2026
function (add_test_executable source_file)
2127
get_filename_component(test_name "${source_file}" NAME_WE)
2228
target_sources(${test_target}

0 commit comments

Comments
 (0)