From ea5ee0c46f31b71435925c128bf7d6cba242641c Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Fri, 11 Jul 2025 02:20:05 +0300 Subject: [PATCH] Free space on overwrites correctly --- src/blockstore/blockstore_heap.cpp | 214 +++++++++++++++++++---------- src/blockstore/blockstore_heap.h | 19 ++- 2 files changed, 158 insertions(+), 75 deletions(-) diff --git a/src/blockstore/blockstore_heap.cpp b/src/blockstore/blockstore_heap.cpp index 80f74140..8867bca8 100644 --- a/src/blockstore/blockstore_heap.cpp +++ b/src/blockstore/blockstore_heap.cpp @@ -371,7 +371,7 @@ skip_object: { fprintf(stderr, "Warning: Object in metadata block %u at %u is a newer duplicate, overriding\n", block_num, block_offset); - erase_object(dup_block, dup_obj); + erase_object(dup_block, dup_obj, 0, false); } } // Verify checksums @@ -660,7 +660,7 @@ bool blockstore_heap_t::recheck_small_writes(std::functioninode, obj->stripe); - erase_object(block_num, obj); + erase_object(block_num, obj, 0, false); } else { @@ -1159,22 +1159,39 @@ int blockstore_heap_t::add_object(object_id oid, heap_write_t *wr, uint32_t *mod return 0; } -heap_object_t *blockstore_heap_t::mvcc_save_copy(heap_object_t *obj) +bool blockstore_heap_t::mvcc_check_tracking(object_id oid) +{ + auto mvcc_it = object_mvcc.lower_bound({ .oid = { .inode = oid.inode, .stripe = oid.stripe }, .lsn = 0 }); + if (mvcc_it == object_mvcc.end()) + { + // no copies + return false; + } + return (mvcc_it->first.oid == oid && mvcc_it->second.entry_copy); +} + +// returns tracking_active, i.e. true if there exists at least one copied MVCC version of the object +// it's used for reference tracking because tracking_active=false means that there is 1 implicit reference +// for heap_writes of the current version of the object and true means that there isn't +bool blockstore_heap_t::mvcc_save_copy(heap_object_t *obj) { auto oid = (object_id){ .inode = obj->inode, .stripe = obj->stripe }; auto mvcc_it = object_mvcc.lower_bound({ .oid = { .inode = oid.inode, .stripe = oid.stripe+1 }, .lsn = 0 }); if (mvcc_it == object_mvcc.begin()) { - return NULL; + // no copies + return false; } mvcc_it--; if (mvcc_it->first.oid != oid) { - return NULL; + // no copies :-) + return false; } if (mvcc_it->second.entry_copy) { - return mvcc_it->second.entry_copy; + // active current copy :-) + return true; } assert(obj->size == sizeof(heap_object_t)); uint32_t total_size = obj->size; @@ -1223,7 +1240,30 @@ heap_object_t *blockstore_heap_t::mvcc_save_copy(heap_object_t *obj) mvcc_buffer_refs[wr->location] += add_ref; } } - return mvcc_it->second.entry_copy; + // copied :-) + return true; +} + +void blockstore_heap_t::mark_overwritten(uint64_t over_lsn, uint64_t inode, heap_write_t *wr, heap_write_t *end_wr, bool tracking_active) +{ + while (wr && wr != end_wr) + { + if (wr->needs_compact(0)) + { + mark_lsn_compacted(wr->lsn); + } + if ((wr->flags & BS_HEAP_TYPE) == BS_HEAP_BIG_WRITE) + { + overwrite_ref_queue.push_back((heap_refqi_t){ .lsn = over_lsn, .inode = inode, .location = wr->location, .len = 0, .is_data = true }); + mvcc_data_refs[wr->location] += !tracking_active; + } + else if ((wr->flags & BS_HEAP_TYPE) == BS_HEAP_SMALL_WRITE && wr->size > 0) + { + overwrite_ref_queue.push_back((heap_refqi_t){ .lsn = over_lsn, .inode = inode, .location = wr->location, .len = wr->len, .is_data = false }); + mvcc_buffer_refs[wr->location] += !tracking_active; + } + wr = wr->next(); + } } int blockstore_heap_t::update_object(uint32_t block_num, heap_object_t *obj, heap_write_t *wr, uint32_t *modified_block) @@ -1245,7 +1285,7 @@ int blockstore_heap_t::update_object(uint32_t block_num, heap_object_t *obj, hea } if (!(obj->get_writes()->flags & BS_HEAP_STABLE) && (wr->flags & BS_HEAP_STABLE)) { - // Stable overwrites are not allowed over unstable + // Stable overwrites are not allowed over unstable - FIXME maybe they are? :-) return EINVAL; } if (wr->version <= obj->get_writes()->version) @@ -1265,13 +1305,7 @@ int blockstore_heap_t::update_object(uint32_t block_num, heap_object_t *obj, hea *modified_block = block_num; } // Save a copy of the object - only when overwriting - heap_object_t *obj_copy = is_overwrite ? mvcc_save_copy(obj) : NULL; - bool tracking_active = !!obj_copy; - if (!tracking_active) - { - auto mvcc_it = object_mvcc.lower_bound((heap_object_lsn_t){ .oid = oid, .lsn = 0 }); - tracking_active = (mvcc_it != object_mvcc.end() && mvcc_it->first.oid == oid && mvcc_it->second.entry_copy); - } + bool tracking_active = is_overwrite ? mvcc_save_copy(obj) : mvcc_check_tracking(oid); if (tracking_active) { // MVCC reference tracking is in action for the object, increase the refcount @@ -1293,15 +1327,12 @@ int blockstore_heap_t::update_object(uint32_t block_num, heap_object_t *obj, hea assert(offset != UINT32_MAX); memcpy(inf.data + offset, wr, wr_size); heap_write_t *new_wr = (heap_write_t*)(inf.data + offset); + new_wr->size = wr_size; + new_wr->lsn = ++next_lsn; int32_t used_delta = wr_size; if (is_overwrite) { - for (auto wr = obj->get_writes(); wr; wr = wr->next()) - { - if (wr->needs_compact(0)) - mark_lsn_compacted(wr->lsn); - } - free_object_space(obj->inode, obj->get_writes(), NULL); + mark_overwritten(new_wr->lsn, obj->inode, obj->get_writes(), NULL, tracking_active); // Free old write entries used_delta -= free_writes(obj->get_writes(), NULL); new_wr->next_pos = 0; @@ -1311,7 +1342,6 @@ int blockstore_heap_t::update_object(uint32_t block_num, heap_object_t *obj, hea { assert(wr->flags == (BS_HEAP_INTENT_WRITE|BS_HEAP_STABLE)); auto second_wr = obj->get_writes()->next(); - free_object_space(obj->inode, obj->get_writes(), second_wr); used_delta -= free_writes(obj->get_writes(), second_wr); new_wr->next_pos = (uint8_t*)second_wr - (uint8_t*)new_wr; } @@ -1319,8 +1349,6 @@ int blockstore_heap_t::update_object(uint32_t block_num, heap_object_t *obj, hea { new_wr->next_pos = ((uint8_t*)obj + obj->write_pos) - (uint8_t*)new_wr; } - new_wr->size = wr_size; - new_wr->lsn = ++next_lsn; wr->lsn = new_wr->lsn; push_inflight_lsn(oid, new_wr->lsn, new_wr->needs_compact(0) ? HEAP_INFLIGHT_COMPACTABLE : 0); if ((wr->flags & BS_HEAP_TYPE) == BS_HEAP_BIG_WRITE) @@ -1398,8 +1426,8 @@ int blockstore_heap_t::post_stabilize(object_id oid, uint64_t version, uint32_t if (unstable_big_wr && unstable_big_wr->next()) { // Remove previous stable entry series - mvcc_save_copy(obj); - free_object_space(obj->inode, unstable_big_wr->next(), NULL); + bool tracking_active = mvcc_save_copy(obj); + mark_overwritten(next_lsn+1, obj->inode, unstable_big_wr->next(), NULL, tracking_active); add_used_space(block_num, -free_writes(unstable_big_wr->next(), NULL)); unstable_big_wr->next_pos = 0; } @@ -1461,16 +1489,18 @@ int blockstore_heap_t::post_rollback(object_id oid, uint64_t version, uint32_t * { *modified_block = block_num; } - mvcc_save_copy(obj); + bool tracking_active = mvcc_save_copy(obj); + ++next_lsn; + push_inflight_lsn(oid, next_lsn, 0); if (!wr) { - erase_object(block_num, obj); + erase_object(block_num, obj, next_lsn, tracking_active); } else { // Erase head versions heap_write_t *first_wr = obj->get_writes(); - free_object_space(obj->inode, first_wr, wr); + mark_overwritten(next_lsn, obj->inode, first_wr, wr, tracking_active); add_used_space(block_num, -free_writes(first_wr, wr)); obj->write_pos = ((uint8_t*)wr - (uint8_t*)obj); obj->crc32c = obj->calc_crc32c(); @@ -1491,10 +1521,12 @@ int blockstore_heap_t::post_delete(object_id oid, uint32_t *modified_block) { *modified_block = block_num; } - mvcc_save_copy(obj); + bool tracking_active = mvcc_save_copy(obj); auto & inf = block_info.at(block_num); assert(inf.data); - erase_object(block_num, obj); + ++next_lsn; + push_inflight_lsn(oid, next_lsn, 0); + erase_object(block_num, obj, next_lsn, tracking_active); return 0; } @@ -1547,34 +1579,74 @@ int blockstore_heap_t::compact_object(object_id oid, uint64_t compact_lsn, uint8 return res; } +void blockstore_heap_t::deref_data(uint64_t inode, uint64_t location, bool free_at_0) +{ + auto ref_it = mvcc_data_refs.find(location); + if (ref_it != mvcc_data_refs.end()) + { + assert(ref_it->second > 0); + ref_it->second--; + if (!ref_it->second) + { + mvcc_data_refs.erase(ref_it); + ref_it = mvcc_data_refs.end(); + } + } + if (ref_it == mvcc_data_refs.end() && free_at_0) + { + assert(data_alloc->get(location >> dsk->block_order)); + data_alloc->set(location >> dsk->block_order, false); + auto & space = inode_space_stats[inode]; + assert(space >= dsk->data_block_size); + space -= dsk->data_block_size; + data_used_space -= dsk->data_block_size; + if (!space) + inode_space_stats.erase(inode); + } +} + +void blockstore_heap_t::deref_buffer(uint64_t inode, uint64_t location, uint32_t len, bool free_at_0) +{ + auto ref_it = mvcc_buffer_refs.find(location); + if (ref_it != mvcc_buffer_refs.end()) + { + assert(ref_it->second > 0); + ref_it->second--; + if (!ref_it->second) + { + mvcc_buffer_refs.erase(ref_it); + ref_it = mvcc_buffer_refs.end(); + } + } + if (ref_it == mvcc_buffer_refs.end() && free_at_0) + { + free_buffer_area(inode, location, len); + } +} + +void blockstore_heap_t::deref_overwrites(uint64_t lsn) +{ + while (overwrite_ref_queue.size() > 0) + { + // Dereference item + auto & el = overwrite_ref_queue.front(); + if (el.lsn > lsn) + break; + if (el.is_data) + deref_data(el.inode, el.location, true); + else + deref_buffer(el.inode, el.location, el.len, true); + overwrite_ref_queue.pop_front(); + } +} + void blockstore_heap_t::free_object_space(inode_t inode, heap_write_t *from, heap_write_t *to, int mode) { for (heap_write_t *wr = from; wr && wr != to; wr = wr->next()) { if ((wr->flags & BS_HEAP_TYPE) == BS_HEAP_BIG_WRITE) { - auto ref_it = mvcc_data_refs.find(wr->location); - if (ref_it != mvcc_data_refs.end()) - { - assert(ref_it->second > 0); - ref_it->second--; - if (!ref_it->second) - { - mvcc_data_refs.erase(ref_it); - ref_it = mvcc_data_refs.end(); - } - } - if (ref_it == mvcc_data_refs.end() && mode != BS_HEAP_FREE_MAIN) - { - assert(data_alloc->get(wr->location >> dsk->block_order)); - data_alloc->set(wr->location >> dsk->block_order, false); - auto & space = inode_space_stats[inode]; - assert(space >= dsk->data_block_size); - space -= dsk->data_block_size; - data_used_space -= dsk->data_block_size; - if (!space) - inode_space_stats.erase(inode); - } + deref_data(inode, wr->location, mode != BS_HEAP_FREE_MAIN); if (mode == BS_HEAP_FREE_MVCC && (wr->flags & BS_HEAP_STABLE)) { // Stop at the last visible version @@ -1583,21 +1655,7 @@ void blockstore_heap_t::free_object_space(inode_t inode, heap_write_t *from, hea } else if ((wr->flags & BS_HEAP_TYPE) == BS_HEAP_SMALL_WRITE) { - auto ref_it = mvcc_buffer_refs.find(wr->location); - if (ref_it != mvcc_buffer_refs.end()) - { - assert(ref_it->second > 0); - ref_it->second--; - if (!ref_it->second) - { - mvcc_buffer_refs.erase(ref_it); - ref_it = mvcc_buffer_refs.end(); - } - } - if (ref_it == mvcc_buffer_refs.end() && mode != BS_HEAP_FREE_MAIN) - { - free_buffer_area(inode, wr->location, wr->len); - } + deref_buffer(inode, wr->location, wr->len, mode != BS_HEAP_FREE_MAIN); } } } @@ -1613,16 +1671,21 @@ void blockstore_heap_t::erase_block_index(inode_t inode, uint64_t stripe) } } -void blockstore_heap_t::erase_object(uint32_t block_num, heap_object_t *obj) +void blockstore_heap_t::erase_object(uint32_t block_num, heap_object_t *obj, uint64_t lsn, bool tracking_active) { // Erase object - for (auto wr = obj->get_writes(); wr; wr = wr->next()) + if (!lsn) { - if (wr->needs_compact(0)) - mark_lsn_compacted(wr->lsn); + for (auto wr = obj->get_writes(); wr; wr = wr->next()) + if (wr->needs_compact(0)) + mark_lsn_compacted(wr->lsn); + free_object_space(obj->inode, obj->get_writes(), NULL); + } + else + { + mark_overwritten(lsn, obj->inode, obj->get_writes(), NULL, tracking_active); } erase_block_index(obj->inode, obj->stripe); - free_object_space(obj->inode, obj->get_writes(), NULL); auto freed = free_writes(obj->get_writes(), NULL); auto obj_size = obj->size; memset((uint8_t*)obj, 0, obj_size); @@ -1896,6 +1959,10 @@ void blockstore_heap_t::mark_lsn_completed(uint64_t lsn) { completed_lsn++; } + if (dsk->disable_meta_fsync && dsk->disable_journal_fsync) + { + deref_overwrites(completed_lsn); + } } } @@ -1905,6 +1972,7 @@ void blockstore_heap_t::mark_lsn_fsynced(uint64_t lsn) { assert(lsn >= first_inflight_lsn && lsn <= completed_lsn); fsynced_lsn = lsn; + deref_overwrites(lsn); } } diff --git a/src/blockstore/blockstore_heap.h b/src/blockstore/blockstore_heap.h index cc8113b3..9d74edf7 100644 --- a/src/blockstore/blockstore_heap.h +++ b/src/blockstore/blockstore_heap.h @@ -111,6 +111,15 @@ struct heap_inflight_lsn_t uint64_t flags; }; +struct heap_refqi_t +{ + uint64_t lsn; + uint64_t inode; + uint64_t location; + uint32_t len; + bool is_data; +}; + class blockstore_heap_t { friend class heap_write_t; @@ -150,6 +159,7 @@ class blockstore_heap_t uint64_t fsynced_lsn = 0; uint64_t compacted_lsn = 0; uint64_t next_compact_lsn = 0; + std::deque overwrite_ref_queue; std::vector tmp_compact_queue; std::deque recheck_queue; @@ -165,12 +175,17 @@ class blockstore_heap_t uint32_t find_block_run(heap_block_info_t & block, uint32_t space); uint32_t find_block_space(uint32_t block_num, uint32_t space); uint32_t compact_object_to(heap_object_t *obj, uint64_t lsn, uint8_t *new_csums); - heap_object_t *mvcc_save_copy(heap_object_t *obj); + bool mvcc_save_copy(heap_object_t *obj); + bool mvcc_check_tracking(object_id oid); int add_object(object_id oid, heap_write_t *wr, uint32_t *modified_block); + void mark_overwritten(uint64_t over_lsn, uint64_t inode, heap_write_t *wr, heap_write_t *end_wr, bool tracking_active); int update_object(uint32_t block_num, heap_object_t *obj, heap_write_t *wr, uint32_t *modified_block); - void erase_object(uint32_t block_num, heap_object_t *obj); + void erase_object(uint32_t block_num, heap_object_t *obj, uint64_t lsn, bool tracking_active); void reindex_block(uint32_t block_num, heap_object_t *from_obj); void erase_block_index(inode_t inode, uint64_t stripe); + void deref_data(uint64_t inode, uint64_t location, bool free_at_0); + void deref_buffer(uint64_t inode, uint64_t location, uint32_t len, bool free_at_0); + void deref_overwrites(uint64_t lsn); void free_object_space(inode_t inode, heap_write_t *from, heap_write_t *to, int mode = 0); void add_used_space(uint32_t block_num, int32_t used_delta); void push_inflight_lsn(object_id oid, uint64_t lsn, uint64_t flags);