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
9 changes: 6 additions & 3 deletions include/async_io_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -121,8 +121,11 @@ class AsyncIoManager
std::span<FilePageId> page_ids,
std::vector<Page> &pages) = 0;

// Writes a single page. @p page is consumed (moved out) ONLY on success; on
// an error return it is left untouched so the caller still owns it and can
// release the cache page it holds (see WriteTask::WritePage).
virtual KvError WritePage(const TableIdent &tbl_id,
VarPage page,
VarPage &page,
FilePageId file_page_id) = 0;

virtual GlobalRegisteredMemory *GetGlobalRegisteredMemory() const
Expand Down Expand Up @@ -561,7 +564,7 @@ class IouringMgr : public AsyncIoManager
std::vector<Page> &pages) override;

KvError WritePage(const TableIdent &tbl_id,
VarPage page,
VarPage &page,
FilePageId file_page_id) override;

KvError ReadSegments(const TableIdent &tbl_id,
Expand Down Expand Up @@ -1590,7 +1593,7 @@ class MemStoreMgr : public AsyncIoManager
std::vector<Page> &pages) override;

KvError WritePage(const TableIdent &tbl_id,
VarPage page,
VarPage &page,
FilePageId file_page_id) override;
KvError SyncData(const TableIdent &tbl_id) override;
KvError AbortWrite(const TableIdent &tbl_id) override;
Expand Down
7 changes: 7 additions & 0 deletions include/tasks/write_task.h
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,13 @@ class WriteTask : public KvTask
KvError WritePage(MemCachedPage::Handle &page, FilePageId file_page_id);
KvError WritePage(VarPage page, FilePageId file_page_id);
KvError AppendWritePage(VarPage page, FilePageId file_page_id);
// Release a cache page still held when a write path bails out on error
// (AppendWritePage's early returns, BatchWriteTask::Pop's index build). A
// freshly promoted/allocated page is detached and pinned only by this
// handle, so it must be freed back to the buffer pool or its slot leaks;
// data/overflow variants and cache pages the caller still pins (index
// writes keep a second IO pin) own their storage and are left alone.
void ReleaseHeldPage(VarPage page);
void FlushAppendWrites();
// Build this task's CoW root and snapshot the branch-file-mapping tail in
// one step: forwards to PageManager::MakeCowRoot, and on success captures
Expand Down
13 changes: 11 additions & 2 deletions src/async_io_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -971,9 +971,13 @@ std::pair<ManifestFilePtr, KvError> IouringMgr::GetManifest(
}

KvError IouringMgr::WritePage(const TableIdent &tbl_id,
VarPage page,
VarPage &page,
FilePageId file_page_id)
{
// Test seam: fail before `page` is consumed (models an OpenOrCreateFD
// failure) so a test can drive the caller's release of the held page.
TEST_FAIL_POINT_RETURN("WritePageBeforeSubmit", KvError::Corrupted);

auto [file_id, offset] = ConvFilePageId(file_page_id);
uint64_t term = ProcessTerm();
std::string_view branch = GetActiveBranch();
Expand Down Expand Up @@ -1253,6 +1257,11 @@ KvError IouringMgr::SubmitMergedWrite(const TableIdent &tbl_id,
std::vector<uint16_t> &release_indices,
bool use_fixed)
{
// Test seam: force the append-mode merged write to fail so a test can drive
// FlushAppendWrites's error path (which must free the pages it holds) while
// AppendWritePage is still holding the current cache page.
TEST_FAIL_POINT_RETURN("SubmitMergedWrite", KvError::Corrupted);

const uint64_t term = ProcessTerm();
const std::string_view branch = GetActiveBranch();
// In append mode, offset 0 means this merged write targets a brand-new
Expand Down Expand Up @@ -7083,7 +7092,7 @@ std::pair<ManifestFilePtr, KvError> MemStoreMgr::GetManifest(
}

KvError MemStoreMgr::WritePage(const TableIdent &tbl_id,
VarPage page,
VarPage &page,
FilePageId file_page_id)
{
auto it = store_.find(tbl_id);
Expand Down
21 changes: 21 additions & 0 deletions src/tasks/batch_write_task.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
#include "async_io_manager.h"
#include "coding.h"
#include "compression.h"
#include "fail_point.h"
#include "storage/shard.h"
#include "tasks/task.h"
#include "utils.h"
Expand Down Expand Up @@ -896,6 +897,11 @@ std::pair<MemCachedPage::Handle, KvError> BatchWriteTask::Pop()
KvError err = add_to_page(new_key, new_page_id);
if (err != KvError::NoError)
{
// prev_handle is empty only when FinishIndexPage OutOfMem'd at
// its own AllocPage (never assigned the page); any other error
// leaves it holding the page to free.
assert(prev_handle || err == KvError::OutOfMem);
ReleaseHeldPage(VarPage(std::move(prev_handle)));
return {MemCachedPage::Handle(), err};
}
}
Expand All @@ -922,6 +928,8 @@ std::pair<MemCachedPage::Handle, KvError> BatchWriteTask::Pop()
KvError err = add_to_page(new_key, new_page_id);
if (err != KvError::NoError)
{
assert(prev_handle || err == KvError::OutOfMem);
ReleaseHeldPage(VarPage(std::move(prev_handle)));
return {MemCachedPage::Handle(), err};
}
AdvanceIndexPageIter(base_page_iter, is_base_iter_valid);
Expand All @@ -936,6 +944,8 @@ std::pair<MemCachedPage::Handle, KvError> BatchWriteTask::Pop()
KvError err = add_to_page(new_key, new_page);
if (err != KvError::NoError)
{
assert(prev_handle || err == KvError::OutOfMem);
ReleaseHeldPage(VarPage(std::move(prev_handle)));
return {MemCachedPage::Handle(), err};
}
}
Expand All @@ -961,12 +971,19 @@ std::pair<MemCachedPage::Handle, KvError> BatchWriteTask::Pop()
prev_handle, prev_key, prev_page_id, std::move(curr_page_key));
if (err != KvError::NoError)
{
// Empty only when FinishIndexPage OutOfMem'd at its AllocPage.
assert(prev_handle || err == KvError::OutOfMem);
ReleaseHeldPage(VarPage(std::move(prev_handle)));
return {MemCachedPage::Handle(), err};
}
err = FlushIndexPage(
prev_handle, std::move(prev_key), prev_page_id, splited);
if (err != KvError::NoError)
{
// FinishIndexPage above succeeded, so prev_handle always holds a
// page here (FlushIndexPage never clears it on failure).
assert(prev_handle);
ReleaseHeldPage(VarPage(std::move(prev_handle)));
return {MemCachedPage::Handle(), err};
}
if (!splited)
Expand Down Expand Up @@ -1023,6 +1040,10 @@ KvError BatchWriteTask::FlushIndexPage(MemCachedPage::Handle &idx_page,
PageId page_id,
bool split)
{
// Test seam: fail the flush while Pop still holds prev_handle, so a test
// can drive Pop's error returns (which must release the held index page).
TEST_FAIL_POINT_RETURN("FlushIndexPage", KvError::Corrupted);

// Flushes the built index page.
idx_page->SetPageId(page_id);
KvError err = WritePage(idx_page);
Expand Down
49 changes: 46 additions & 3 deletions src/tasks/write_task.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -299,8 +299,15 @@ KvError WriteTask::WritePage(VarPage page, FilePageId file_page_id)
return AppendWritePage(std::move(page), file_page_id);
}

KvError err = IoMgr()->WritePage(tbl_ident_, std::move(page), file_page_id);
CHECK_KV_ERR(err);
// WritePage consumes `page` only on success; on error it leaves it with us
// so a synchronous failure (e.g. OpenOrCreateFD) does not orphan the cache
// page this handle holds.
KvError err = IoMgr()->WritePage(tbl_ident_, page, file_page_id);
if (err != KvError::NoError)
{
ReleaseHeldPage(std::move(page));
return err;
}
if (inflight_io_ >= opts->max_write_batch_pages)
{
// Avoid long running WriteTask block ReadTask/ScanTask
Expand Down Expand Up @@ -345,6 +352,7 @@ KvError WriteTask::AppendWritePage(VarPage page, FilePageId file_page_id)
FlushAppendWrites();
if (write_err_ != KvError::NoError)
{
ReleaseHeldPage(std::move(page));
return write_err_;
}
// In cloud append mode, trigger immediate upload of sealed file
Expand All @@ -353,12 +361,17 @@ KvError WriteTask::AppendWritePage(VarPage page, FilePageId file_page_id)
{
KvError err = IoMgr()->OnDataFileSealed(
tbl_ident_, DataFileKey(sealed_file_id));
CHECK_KV_ERR(err);
if (err != KvError::NoError)
{
ReleaseHeldPage(std::move(page));
return err;
}
}
uint16_t buf_index = 0;
char *buf = IoMgr()->AcquireWriteBuffer(buf_index);
if (buf == nullptr)
{
ReleaseHeldPage(std::move(page));
return KvError::OutOfMem;
}
bool use_fixed = IoMgr()->WriteBufferUseFixed();
Expand All @@ -369,6 +382,7 @@ KvError WriteTask::AppendWritePage(VarPage page, FilePageId file_page_id)
char *dst = append_aggregator_.TryReserve(file_id, offset, page_size);
if (dst == nullptr)
{
ReleaseHeldPage(std::move(page));
return KvError::OutOfMem;
}
std::memcpy(dst, page_ptr, page_size);
Expand All @@ -385,6 +399,35 @@ KvError WriteTask::AppendWritePage(VarPage page, FilePageId file_page_id)
return KvError::NoError;
}

void WriteTask::ReleaseHeldPage(VarPage page)
{
if (VarPageType(page.index()) != VarPageType::MemCachedPage)
{
// Data/overflow pages carry their own buffer; destroying the VarPage
// releases it. Only promoted cache pages need explicit accounting.
return;
}
MemCachedPage::Handle &handle = std::get<MemCachedPage::Handle>(page);
MemCachedPage *cache_page = handle.Get();
if (cache_page == nullptr)
{
// Empty handle: nothing to release. This happens on an OutOfMem return
// where the allocation that would have populated the handle is exactly
// what failed (BatchWriteTask::Pop's index build via FinishIndexPage);
// the caller asserts the OutOfMem precondition.
return;
}
handle.Reset(); // drop this task's pin
// Free only when this handle was the sole owner (the promoted data-page
// case). A page still pinned by the caller -- index-page writes pass a
// second, temporary pin -- is the caller's to release, and a page already
// linked into the active/free list must not be freed here.
if (cache_page->IsDetached() && !cache_page->IsPinned())
{
shard->IndexManager()->FreePage(cache_page);
}
}

void WriteTask::FlushAppendWrites()
{
if (!append_aggregator_.HasData())
Expand Down
Loading
Loading