Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion conanfile.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@

class HomeObjectConan(ConanFile):
name = "homeobject"
version = "4.3.5"
version = "4.3.6"

homepage = "https://github.com/eBay/HomeObject"
description = "Blob Store built on HomeStore"
Expand Down
4 changes: 2 additions & 2 deletions src/lib/homestore_backend/gc_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,8 @@ SISL_LOGGING_DECL(gcmgr)
#define RECOVERD_GC_TASK_ID 0

#define GCLOG(level, gc_task_id, pg_id, shard_id, msg, ...) \
LOG##level##MOD(gcmgr, "[gc_task_id={}, pg_id={}, shard_id=0x{:x}] " msg, gc_task_id, pg_id, shard_id, \
##__VA_ARGS__)
LOG##level##MOD(gcmgr, "[gc_task_id={}, pg_id={}, shard_id=0x{:x}] " msg, gc_task_id, pg_id, \
shard_id, ##__VA_ARGS__)

#define GCLOGT(gc_task_id, pg_id, shard_id, msg, ...) GCLOG(TRACE, gc_task_id, pg_id, shard_id, msg, ##__VA_ARGS__)
#define GCLOGD(gc_task_id, pg_id, shard_id, msg, ...) GCLOG(DEBUG, gc_task_id, pg_id, shard_id, msg, ##__VA_ARGS__)
Expand Down
44 changes: 16 additions & 28 deletions src/lib/homestore_backend/gc_manager.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -155,14 +155,10 @@ class GCManager {

// Backlog / pressure snapshot gauges. Values refreshed once per scan cycle by
// GCManager::scan_chunks_for_gc; worst-case staleness = gc_scan_interval_sec.
REGISTER_GAUGE(pending_gc_bytes,
"Total reclaimable garbage bytes in PG-owned chunks on this pdev");
REGISTER_GAUGE(eligible_gc_bytes,
"Reclaimable bytes currently eligible for normal GC on this pdev");
REGISTER_GAUGE(eligible_gc_chunk_count,
"Chunks currently eligible for normal GC on this pdev");
REGISTER_GAUGE(pending_normal_gc_task_count,
"Normal-priority GC tasks queued or running on this pdev");
REGISTER_GAUGE(pending_gc_bytes, "Total reclaimable garbage bytes in PG-owned chunks on this pdev");
REGISTER_GAUGE(eligible_gc_bytes, "Reclaimable bytes currently eligible for normal GC on this pdev");
REGISTER_GAUGE(eligible_gc_chunk_count, "Chunks currently eligible for normal GC on this pdev");
REGISTER_GAUGE(pending_normal_gc_task_count, "Normal-priority GC tasks queued or running on this pdev");

// Distribution of the pending backlog by garbage-ratio bucket. We register 10
// gauges under a single Prometheus metric name (`pending_gc_chunks_ratio`),
Expand All @@ -179,18 +175,16 @@ class GCManager {
// shape. Prometheus text output is unaffected — it uses HELP (unchanged across
// registrations) and disambiguates series by labels.
static constexpr std::array< const char*, 10 > kRatioBucketLabels = {
"00-10", "10-20", "20-30", "30-40", "40-50",
"50-60", "60-70", "70-80", "80-90", "90-100"};
"00-10", "10-20", "20-30", "30-40", "40-50", "50-60", "60-70", "70-80", "80-90", "90-100"};
for (size_t i = 0; i < kRatioBucketLabels.size(); ++i) {
const auto lo = i * 10;
const auto hi = (i + 1) * 10;
const auto desc = fmt::format(
"Snapshot count of pending chunks with garbage ratio in ({}, {}]% "
"(bucket={})",
lo, hi, kRatioBucketLabels[i]);
ratio_bucket_indices_[i] = m_impl_ptr->register_gauge(
"pending_gc_chunks_ratio", desc, "" /* report_name */,
sisl::metric_label{"bucket", kRatioBucketLabels[i]});
const auto desc = fmt::format("Snapshot count of pending chunks with garbage ratio in ({}, {}]% "
"(bucket={})",
lo, hi, kRatioBucketLabels[i]);
ratio_bucket_indices_[i] =
m_impl_ptr->register_gauge("pending_gc_chunks_ratio", desc, "" /* report_name */,
sisl::metric_label{"bucket", kRatioBucketLabels[i]});
}

register_me_to_farm();
Expand Down Expand Up @@ -230,9 +224,8 @@ class GCManager {
// Bypass GAUGE_UPDATE for the same reason as bucket registration: we need to
// address 10 distinct gauge indices that share one metric name.
for (size_t i = 0; i < ratio_bucket_indices_.size(); ++i) {
m_impl_ptr->gauge_update(
ratio_bucket_indices_[i],
static_cast< int64_t >(gc_actor_.get_pending_ratio_bucket(i)));
m_impl_ptr->gauge_update(ratio_bucket_indices_[i],
static_cast< int64_t >(gc_actor_.get_pending_ratio_bucket(i)));
}
}

Expand Down Expand Up @@ -315,12 +308,8 @@ class GCManager {

// Snapshot readers used by pdev_gc_metrics::on_gather. Return the last value published
// by GCManager::scan_chunks_for_gc for this pdev; 0 before the first scan completes.
uint64_t get_pending_gc_bytes() const {
return m_pending_gc_bytes.load(std::memory_order_relaxed);
}
uint64_t get_eligible_gc_bytes() const {
return m_eligible_gc_bytes.load(std::memory_order_relaxed);
}
uint64_t get_pending_gc_bytes() const { return m_pending_gc_bytes.load(std::memory_order_relaxed); }
uint64_t get_eligible_gc_bytes() const { return m_eligible_gc_bytes.load(std::memory_order_relaxed); }
uint32_t get_eligible_gc_chunk_count() const {
return m_eligible_gc_chunk_count.load(std::memory_order_relaxed);
}
Expand All @@ -332,8 +321,7 @@ class GCManager {
// the totals locally over all chunks on this pdev, then hand them in via one call so the
// metrics stay internally consistent within a scan cycle. Between-gauge drift is bounded
// by one scan interval; individual scalars are aligned and therefore torn-read safe.
void publish_scan_snapshot(uint64_t pending_bytes, uint64_t eligible_bytes,
uint32_t eligible_chunks,
void publish_scan_snapshot(uint64_t pending_bytes, uint64_t eligible_bytes, uint32_t eligible_chunks,
const std::array< uint32_t, 10 >& ratio_buckets) {
m_pending_gc_bytes.store(pending_bytes, std::memory_order_relaxed);
m_eligible_gc_bytes.store(eligible_bytes, std::memory_order_relaxed);
Expand Down
4 changes: 2 additions & 2 deletions src/lib/homestore_backend/hs_blob_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,8 @@ SISL_LOGGING_DECL(blobmgr)

#define BLOG(level, trace_id, shard_id, blob_id, msg, ...) \
LOG##level##MOD(blobmgr, "[traceID={},shardID=0x{:x},pg={},shard=0x{:x},blob={}] " msg, trace_id, shard_id, \
(shard_id >> homeobject::shard_width), (shard_id & homeobject::shard_mask), blob_id, \
##__VA_ARGS__)
(shard_id >> homeobject::shard_width), (shard_id & homeobject::shard_mask), \
blob_id, ##__VA_ARGS__)

#define BLOGT(trace_id, shard_id, blob_id, msg, ...) BLOG(TRACE, trace_id, shard_id, blob_id, msg, ##__VA_ARGS__)
#define BLOGD(trace_id, shard_id, blob_id, msg, ...) BLOG(DEBUG, trace_id, shard_id, blob_id, msg, ##__VA_ARGS__)
Expand Down
4 changes: 4 additions & 0 deletions src/lib/homestore_backend/hs_homeobject.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -439,9 +439,11 @@ class HSHomeObject : public HomeObjectImpl {
* RPC handlers for scrub:
* 1. on_scrub_req_received: receive the scrub req from leader
* 2. on_scrub_result_received: receive the scrub result from followers
* 3. on_scrub_timestamp_push_received: receive the pushed scrub superblk timestamps from the leader
*/
void on_scrub_req_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data);
void on_scrub_result_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data);
void on_scrub_timestamp_push_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data);

/**
* Register data RPC handlers for this PG
Expand Down Expand Up @@ -572,6 +574,8 @@ class HSHomeObject : public HomeObjectImpl {
inline const static std::string PUSH_SCRUB_REQ{"PUSH_SCRUB_REQ"};
// return scrub result to leader
inline const static std::string PUSH_SCRUB_RESULT{"PUSH_SCRUB_RESULT"};
// sync last_deep_scrub_timestamp/last_shallow_scrub_timestamp between leader and followers
inline const static std::string PUSH_SCRUB_TIMESTAMP{"PUSH_SCRUB_TIMESTAMP"};

class PGBlobIterator {
public:
Expand Down
39 changes: 39 additions & 0 deletions src/lib/homestore_backend/hs_pg_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1182,6 +1182,14 @@ void HSHomeObject::HS_PG::register_data_rpc_handlers() {
} else {
LOGW("PUSH_SCRUB_RESULT RPC handler already registered for pg={}", pg_id);
}

success =
repl_dev_->add_data_rpc_service(PUSH_SCRUB_TIMESTAMP, bind_this(HS_PG::on_scrub_timestamp_push_received, 1));
if (success) {
LOGI("Successfully registered PUSH_SCRUB_TIMESTAMP RPC handler for pg={}", pg_id);
} else {
LOGW("PUSH_SCRUB_TIMESTAMP RPC handler already registered for pg={}", pg_id);
}
}

void HSHomeObject::HS_PG::on_scrub_req_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data) {
Expand Down Expand Up @@ -1329,6 +1337,37 @@ void HSHomeObject::HS_PG::on_scrub_result_received(boost::intrusive_ptr< sisl::G
scrub_mgr->handle_scrub_req_resp(pg_id, scrub_result);
}

void HSHomeObject::HS_PG::on_scrub_timestamp_push_received(boost::intrusive_ptr< sisl::GenericRpcData >& rpc_data) {
const auto pg_id = pg_info_.id;
LOGD("Received scrub superblk push for pg={}", pg_id);

auto const& incoming_buf = rpc_data->request_blob();
if (incoming_buf.size() != sizeof(ScrubManager::scrub_timestamp_info)) {
LOGW("scrub superblk push received with invalid buffer size for pg={}, size={}", pg_id, incoming_buf.size());
rpc_data->send_response();
return;
}

ScrubManager::scrub_timestamp_info incoming;
std::memcpy(&incoming, incoming_buf.cbytes(), sizeof(incoming));

auto scrub_mgr = home_obj_.scrub_manager();
if (!scrub_mgr) {
LOGW("ScrubManager is not initialized in HS_PG::on_scrub_timestamp_push_received for pg={}", pg_id);
rpc_data->send_response();
return;
}

auto merged =
std::make_shared< ScrubManager::scrub_timestamp_info >(scrub_mgr->sync_scrub_timestamp(pg_id, incoming));

sisl::io_blob_list_t blob_list;
blob_list.emplace_back(reinterpret_cast< uint8_t* >(merged.get()),
static_cast< uint32_t >(sizeof(ScrubManager::scrub_timestamp_info)), false);
rpc_data->set_comp_cb([merged](boost::intrusive_ptr< sisl::GenericRpcData >&) {});
rpc_data->send_response(blob_list);
}

// NOTE: caller should hold the _pg_lock
const HSHomeObject::HS_PG* HSHomeObject::_get_hs_pg_unlocked(pg_id_t pg_id) const {
auto iter = _pg_map.find(pg_id);
Expand Down
2 changes: 1 addition & 1 deletion src/lib/homestore_backend/replication_state_machine.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -277,7 +277,7 @@ class ReplicationStateMachine : public homestore::ReplDevListener {

std::pair< homestore::repl_lsn_t, homestore::chunk_num_t > get_no_space_left_error_info() const;

void handle_no_space_left(homestore ::repl_lsn_t lsn, homestore ::chunk_num_t chunk_id);
void handle_no_space_left(homestore::repl_lsn_t lsn, homestore::chunk_num_t chunk_id);
};

} // namespace homeobject
Loading
Loading