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
10 changes: 10 additions & 0 deletions include/lavik/storage/detail/grouped_hash.h
Original file line number Diff line number Diff line change
Expand Up @@ -294,6 +294,16 @@ class HashGroupMap {
}
return end();
}
// Exact metadata lookup without constructing the ancestor stack needed by
// an iterator. The returned pointer borrows this immutable tree version.
const RecoveredHashGroup* Get(Key key) const noexcept {
const auto* node = root_.get();
while (node) {
if (key == node->entry_.first) return &node->entry_.second;
node = key < node->entry_.first ? node->left_.get() : node->right_.get();
}
return nullptr;
}
const RecoveredHashGroup& at(Key key) const {
const auto it = find(key);
if (it == end()) throw std::out_of_range("group directory key");
Expand Down
67 changes: 48 additions & 19 deletions src/storage/engine/grouped_hash.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,12 @@
#include "lavik/storage/detail/grouped_hash.h"

#include <algorithm>
#include <cstring>
#include <limits>
#include <utility>

#include "absl/container/flat_hash_map.h"
#include "absl/container/inlined_vector.h"

namespace lavik::storage {
namespace {
Expand Down Expand Up @@ -397,9 +399,16 @@ absl::StatusOr<HashGroupEdit> ApplyHashGroupEdits(
Store(encoded, offset, view.field_.size(), 4);
Store(encoded, offset + 4, view.value_.size(), 4);
offset += 8;
encoded.replace(offset, view.field_.size(), view.field_);
// The destination was sized once above and cannot alias these views
// of the read lease/request. Copy bytes directly; Set values are empty
// and need no string mutation (or a zero-length copy from nullptr).
if (!view.field_.empty())
std::memcpy(encoded.data() + offset, view.field_.data(),
view.field_.size());
offset += view.field_.size();
encoded.replace(offset, view.value_.size(), view.value_);
if (!view.value_.empty())
std::memcpy(encoded.data() + offset, view.value_.data(),
view.value_.size());
offset += view.value_.size();
}
}
Expand Down Expand Up @@ -585,7 +594,7 @@ absl::StatusOr<HashGroupDirectory> HashGroupDirectory::Recover(
continue;
}
if (candidate.field_count_ > root.field_count_ - count ||
directory.groups_.find(id.prefix_) != directory.groups_.end()) {
directory.groups_.Get(id.prefix_) != nullptr) {
return absl::DataLossError("overlapping Hash groups or invalid count");
}
auto inserted = directory.groups_.Set(id.prefix_, candidate);
Expand Down Expand Up @@ -635,28 +644,43 @@ absl::StatusOr<HashGroupDirectory> HashGroupDirectory::Apply(
next.root_ = root;
next.sequence_ = root.revision_;
next.command_sequence_ = sequence;
std::map<HashGroupId, const RecoveredHashGroup*> writes;
// Most point writes replace one route. A sorted pointer list preserves
// duplicate detection and prefix order without allocating a map node per
// changed group; larger batches spill under the existing scratch admission.
absl::InlinedVector<const RecoveredHashGroup*, 4> writes;
writes.reserve(changes.size());
for (const auto& change : changes) writes.push_back(&change);
std::sort(writes.begin(), writes.end(),
[](const auto* a, const auto* b) { return a->id_ < b->id_; });
std::optional<HashGroupId> previous;
std::uint64_t count = root_.field_count_;
// Hash-prefix leaves partition a 64-bit domain. A wider scratch accumulator
// verifies total coverage after local replacements, without scanning every
// unchanged route. Pairwise overlap checks below make equal coverage imply
// there are no gaps either.
__uint128_t coverage = static_cast<__uint128_t>(1) << 64;
for (const auto& change : changes) {
for (const auto* changed : writes) {
const auto& change = *changed;
if (!change.id_.valid() || change.incarnation_ != root.incarnation_ ||
change.sequence_ != root.revision_ ||
change.field_count_ > std::numeric_limits<std::uint32_t>::max() ||
(change.retired_ && change.field_count_ != 0) ||
!writes.emplace(change.id_, &change).second) {
previous == change.id_) {
return absl::DataLossError("invalid grouped directory mutation record");
}
const auto current = groups_.find(change.id_.prefix_);
if (current != groups_.end() && current->second.id_ == change.id_) {
count -= current->second.field_count_;
next.total_group_bytes_ -= current->second.encoded_bytes_;
previous = change.id_;
const auto* current = groups_.Get(change.id_.prefix_);
if (current != nullptr && current->id_ == change.id_) {
count -= current->field_count_;
next.total_group_bytes_ -= current->encoded_bytes_;
coverage -= static_cast<__uint128_t>(1) << (64 - change.id_.bits_);
const auto erased = next.groups_.Erase(change.id_.prefix_);
if (!erased.ok()) return erased;
// An unchanged prefix identity only replaces metadata. Erasing it first
// would copy/rebalance an immutable AVL path that Set immediately copies
// again. Splits still remove their retired parent before adding children.
if (change.retired_) {
const auto erased = next.groups_.Erase(change.id_.prefix_);
if (!erased.ok()) return erased;
}
} else if (change.retired_) {
return absl::DataLossError("retiring a non-current Hash group");
}
Expand All @@ -665,22 +689,27 @@ absl::StatusOr<HashGroupDirectory> HashGroupDirectory::Apply(
if (!retired.ok()) return retired;
}
}
for (const auto& [id, changed] : writes) {
for (const auto* changed : writes) {
const auto& change = *changed;
const auto id = change.id_;
if (change.retired_) continue;
if (next.retired_.find(id) != next.retired_.end()) {
if (next.retired_.Get(id) != nullptr) {
return absl::DataLossError("resurrecting a retired Hash group");
}
const auto* floor = next.groups_.Floor(id.prefix_);
if (floor && floor->id_.last() >= id.prefix_) {
const bool replaces_route = floor && floor->id_ == id;
if (floor && !replaces_route && floor->id_.last() >= id.prefix_) {
return absl::DataLossError("group update overlaps an earlier route");
}
const auto inserted = next.groups_.Set(id.prefix_, change);
if (!inserted.ok()) return inserted;
auto after = next.groups_.find(id.prefix_);
++after;
if (after != next.groups_.end() && after->first <= id.last()) {
return absl::DataLossError("group update overlaps a later route");
// Replacing the exact interval cannot change either neighbour boundary.
if (!replaces_route) {
auto after = next.groups_.find(id.prefix_);
++after;
if (after != next.groups_.end() && after->first <= id.last()) {
return absl::DataLossError("group update overlaps a later route");
}
}
count += change.field_count_;
if (change.encoded_bytes_ >
Expand Down
23 changes: 13 additions & 10 deletions src/storage/engine/grouped_mutation.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

#include <charconv>

#include "absl/container/inlined_vector.h"
#include "impl.h"
#include "lavik/storage/detail/grouped_scratch.h"

Expand Down Expand Up @@ -276,8 +277,11 @@ Task<absl::Status> StorageEngine::Impl::CommitGroupedHashMutationLocked(
}
auto encoder_admission = encoder_budget.Reserve(1);
if (!encoder_admission.ok()) co_return encoder_admission.status();
std::vector<HashGroupEncoder> encoders;
std::vector<std::uint64_t> encoded_sizes;
// Point writes normally produce one page. Keep its preflight and publication
// metadata inside this coroutine frame; multi-page commands retain the same
// admitted spill capacity and lifetime through all asynchronous writes.
absl::InlinedVector<HashGroupEncoder, 1> encoders;
absl::InlinedVector<std::uint64_t, 1> encoded_sizes;
try {
// No current-command pages have been staged if preflight allocation fails.
LAVIK_FAULT_BAD_ALLOC("LAVIK_FAIL_GROUP_ENCODER_PREPARE_KEY", key);
Expand Down Expand Up @@ -306,12 +310,11 @@ Task<absl::Status> StorageEngine::Impl::CommitGroupedHashMutationLocked(
std::uint64_t total = previous->directory().total_group_bytes();
for (std::size_t i = 0; i < plan->writes_.size(); ++i) {
const auto& page = plan->writes_[i];
const auto old = previous->directory().groups().find(page.id_.prefix_);
if (old != previous->directory().groups().end() &&
old->second.id_ == page.id_) {
if (old->second.encoded_bytes_ > total)
const auto* old = previous->directory().groups().Get(page.id_.prefix_);
if (old != nullptr && old->id_ == page.id_) {
if (old->encoded_bytes_ > total)
co_return absl::DataLossError("invalid grouped byte total");
total -= old->second.encoded_bytes_;
total -= old->encoded_bytes_;
}
if (!page.retired_) {
if (encoded_sizes[i] > UINT64_MAX - total)
Expand Down Expand Up @@ -406,7 +409,7 @@ Task<absl::Status> StorageEngine::Impl::CommitGroupedHashMutationLocked(
if (!reserved.ok()) co_return reserved.status();
std::optional<GroupedObjectIndex::Publication> publication(
std::move(*reserved));
std::vector<HashGroupLocation> written;
absl::InlinedVector<HashGroupLocation, 1> written;
written.reserve(plan->writes_.size());
// A failed batch in an outer transaction retains its staged bytes until the
// outer commit retires them. Its existing fence must not point at a block we
Expand Down Expand Up @@ -461,8 +464,8 @@ Task<absl::Status> StorageEngine::Impl::CommitGroupedHashMutationLocked(
}
written.push_back(std::move(*group));
}
std::vector<RecoveredHashGroup> candidates;
std::vector<HashGroupId> written_ids;
absl::InlinedVector<RecoveredHashGroup, 1> candidates;
absl::InlinedVector<HashGroupId, 1> written_ids;
candidates.reserve(written.size());
written_ids.reserve(written.size());
for (std::size_t i = 0; i < written.size(); ++i) {
Expand Down
48 changes: 45 additions & 3 deletions src/storage/engine/grouped_object_index.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -438,6 +438,49 @@ absl::StatusOr<NodeHandle> UpdatePhysical(
if (!node) return BuildPhysical(changed, arena);
if (node->page_) {
const auto& page = *node->page_;
// A point replacement keeps the page's identities and cardinality. Copy
// the admitted compact arrays directly instead of expanding every old
// coordinate into a temporary RecordLocation and compacting it again.
// Pages with manifests retain the general ownership-preserving merge.
if (changed.size() == 1 && page.extents_.empty() &&
!changed.front().extents_ && !changed.front().location_.external()) {
const auto found = std::lower_bound(page.ids_.begin(), page.ids_.end(),
changed.front().id_);
if (found != page.ids_.end() && *found == changed.front().id_) {
auto replacement = AllocateLocalObject<GroupIndexNode>(arena);
if (!replacement.ok()) return replacement.status();
auto copied = AllocateLocalObject<GroupIndexPage>(arena);
if (!copied.ok()) return copied.status();
const auto arrays_bytes =
AllocatorUsableSizeForRequest(page.ids_.size() *
sizeof(HashGroupId)) +
AllocatorUsableSizeForRequest(page.locations_.size() *
sizeof(GroupedRecordIndexEntry));
auto reservation = TryReserveMemory(arrays_bytes);
if (!reservation) {
RecordMemoryRejection();
return absl::ResourceExhaustedError(
"OOM group page exceeds maxmemory");
}
(*copied)->arrays_charge_.Account(
arena->allocation_domain().owner_shard_, arrays_bytes);
reservation.reset();
(*copied)->ids_ = page.ids_;
(*copied)->locations_ = page.locations_;
(*copied)->extents_.SetEntryArena(arena);
const auto index = static_cast<std::size_t>(found - page.ids_.begin());
(*copied)->locations_[index] = {
RecordIndexValue(changed.front().location_)};
const auto mask = std::uint64_t{1} << index;
(*copied)->retired_ = changed.front().retired_ ? page.retired_ | mask
: page.retired_ & ~mask;
(*replacement)->size_ = node->size_;
(*replacement)->representative_ = node->representative_;
(*replacement)->branch_depth_ = node->branch_depth_;
(*replacement)->page_ = std::move(*copied);
return NodeHandle(std::move(*replacement));
}
}
// Both inputs have unique, sorted full identities. Merge directly rather
// than allocating a map node for every unchanged location on each write.
// Replaced locations need no old lookup or manifest reference at all.
Expand Down Expand Up @@ -1157,9 +1200,8 @@ const GroupedRecordIndexEntry* GroupedHashObject::FindGroup(
: nullptr;
}
if (is_ordered() && !has_member_index()) return nullptr;
const auto route = directory().groups().find(id.prefix_);
if (route == directory().groups().end() || route->second.id_ != id)
return nullptr;
const auto* route = directory().groups().Get(id.prefix_);
if (route == nullptr || route->id_ != id) return nullptr;
return FindRecord(id);
}

Expand Down
6 changes: 5 additions & 1 deletion src/storage/engine/grouped_read.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,11 @@ StorageEngine::Impl::LoadHashGroupPayload(
// hold graph pins; an ordinary reader captures current identity here and
// retries against a refreshed view if GC moves it during IO.
const RecordLocation location = MaterializeIndexLocation(*entry);
auto extents = object->ExtentsFor(id);
// The captured record says whether a manifest exists; inline pages need
// no second physical-index lookup. External reads retain the owned handle
// across suspension exactly as before.
auto extents =
location.external() ? object->ExtentsFor(id) : ExtentManifest{};
absl::StatusOr<LoadedValue> loaded;
if (location.external()) {
loaded =
Expand Down
76 changes: 60 additions & 16 deletions src/storage/engine/hash_tree.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@

#include "absl/container/flat_hash_map.h"
#include "absl/container/flat_hash_set.h"
#include "absl/container/inlined_vector.h"
#include "absl/hash/hash.h"
#include "impl.h"
#include "lavik/glob.h"
Expand Down Expand Up @@ -731,18 +732,27 @@ Task<absl::StatusOr<HashResult>> StorageEngine::Impl::ExecuteHashLikeLocked(
co_return absl::InvalidArgumentError("Hash field/value mismatch");
std::set<HashGroupId> selected;
std::optional<MemoryReservation> operand_scratch;
std::map<HashGroupId, std::vector<HashEntryView>> leaf_edits;
struct RoutedEdit {
HashGroupId id_;
HashEntryView view_;
std::size_t ordinal_;
};
// A single-field write needs no heap-allocated routing index. For batches,
// sort borrowed views by leaf and original position: duplicates retain
// command order, and one reusable page list replaces a vector per leaf.
absl::InlinedVector<RoutedEdit, 1> leaf_edits;
if (edit_leaves) {
GroupedScratchBudget budget;
if (operation.fields_.size() > SIZE_MAX / 256)
co_return absl::ResourceExhaustedError("Hash operand index overflow");
// Routing vectors, mutation views, map slots and growth rounding, even
// for empty operands. Request strings themselves remain client-owned.
// Routing and per-page views, including growth rounding and empty
// operands. Request strings themselves remain client-owned.
auto added = budget.AddBytes(operation.fields_.size() * 256);
if (!added.ok()) co_return added;
auto admitted = budget.Reserve(1);
if (!admitted.ok()) co_return admitted.status();
operand_scratch.emplace(std::move(*admitted));
leaf_edits.reserve(operation.fields_.size());
}
std::optional<MemoryReservation> grouped_scratch;
if (replace || unlocked_create) {
Expand Down Expand Up @@ -774,24 +784,50 @@ Task<absl::StatusOr<HashResult>> StorageEngine::Impl::ExecuteHashLikeLocked(
const auto* route = grouped->directory().Find(field);
if (route == nullptr)
co_return absl::DataLossError("missing Hash field route");
selected.insert(route->id_);
if (edit_leaves)
leaf_edits[route->id_].push_back(
{field, operation.kind_ == HashOperationKind::kDelete
? std::string_view{}
: operation.values_[i]});
if (edit_leaves) {
leaf_edits.push_back(
{route->id_,
{field, operation.kind_ == HashOperationKind::kDelete
? std::string_view{}
: operation.values_[i]},
i});
} else {
selected.insert(route->id_);
}
}
} else {
for (const auto& [prefix, metadata] : grouped->directory().groups())
selected.insert(metadata.id_);
}
for (const auto id : selected) {
const auto add_group = [&](HashGroupId id) -> absl::Status {
const auto* entry = grouped->FindGroup(id);
if (entry == nullptr)
co_return absl::DataLossError("missing Hash scratch page");
const auto added =
budget.AddGroup(entry->value_, grouped->ExtentsFor(id));
if (!added.ok()) co_return added;
return absl::DataLossError("missing Hash scratch page");
// Inline records cannot own a manifest. Avoid a second physical-index
// traversal merely to obtain the null handle used by admission.
return budget.AddGroup(entry->value_, entry->value_.external()
? grouped->ExtentsFor(id)
: ExtentManifest{});
};
if (edit_leaves) {
if (leaf_edits.size() > 1)
std::sort(leaf_edits.begin(), leaf_edits.end(),
[](const auto& a, const auto& b) {
return a.id_ != b.id_ ? a.id_ < b.id_
: a.ordinal_ < b.ordinal_;
});
std::optional<HashGroupId> previous;
for (const auto& edit : leaf_edits) {
if (previous == edit.id_) continue;
const auto added = add_group(edit.id_);
if (!added.ok()) co_return added;
previous = edit.id_;
}
} else {
for (const auto id : selected) {
const auto added = add_group(id);
if (!added.ok()) co_return added;
}
}
for (const auto field : operation.fields_) {
const auto added = budget.AddBytes(field.size());
Expand All @@ -818,8 +854,16 @@ Task<absl::StatusOr<HashResult>> StorageEngine::Impl::ExecuteHashLikeLocked(
: operation.kind_ == HashOperationKind::kSet
? HashGroupEditKind::kSet
: HashGroupEditKind::kSetIfAbsent;
for (const auto& [id, edits] : leaf_edits) {
co_await bycorf::Yield(*store.worker_);
absl::InlinedVector<HashEntryView, 1> edits;
for (std::size_t begin = 0; begin < leaf_edits.size();) {
// Yield between pages for batch fairness, without scheduling an extra
// turn before every single-field write. Actual IO still suspends.
if (begin != 0) co_await bycorf::Yield(*store.worker_);
const auto id = leaf_edits[begin].id_;
edits.clear();
do {
edits.push_back(leaf_edits[begin++].view_);
} while (begin < leaf_edits.size() && leaf_edits[begin].id_ == id);
auto loaded = co_await LoadHashGroupPayload(store, partition, db_id,
key, digest, grouped, id);
if (!loaded.ok()) co_return loaded.status();
Expand Down
Loading