blob: 2e3532bfed058c71e0a01c3a363a6b0a05aae937 [file]
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
#include "io/cache/async_cache_write_manager.h"
#include <gtest/gtest.h>
#include <algorithm>
#include <atomic>
#include <barrier>
#include <condition_variable>
#include <cstring>
#include <future>
#include <limits>
#include <memory>
#include <mutex>
#include <new>
#include <string>
#include <thread>
#include <vector>
#include "common/config.h"
#include "common/exception.h"
#include "cpp/sync_point.h"
#include "io/cache/async_cache_write_manager_metrics.h"
#include "io/cache/block_file_cache_test_common.h"
#include "io/cache/inflight_write_buffer_index.h"
#include "util/defer_op.h"
#include "util/mem_info.h"
#include "util/time.h"
namespace doris::io {
namespace {
FileCacheSettings async_write_cache_settings() {
FileCacheSettings settings;
settings.query_queue_size = 4_mb;
settings.query_queue_elements = 1024;
settings.index_queue_size = 1_mb;
settings.index_queue_elements = 256;
settings.disposable_queue_size = 1_mb;
settings.disposable_queue_elements = 256;
settings.capacity = 8_mb;
settings.max_file_block_size = 4096;
settings.max_query_cache_size = 0;
return settings;
}
AsyncCacheWriteTask make_async_write_task(
AsyncCacheWriteManager* manager, const std::string& key, char fill,
std::function<void(const AsyncCacheWriteTask&)> on_finalized = nullptr,
int64_t submit_ts_us = MonotonicMicros()) {
DORIS_CHECK(manager != nullptr);
AsyncCacheWriteBufferPtr buffer;
EXPECT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
DORIS_CHECK(buffer != nullptr);
memset(buffer->data(), fill, buffer->size());
const auto cache_hash = BlockFileCache::hash(key);
return AsyncCacheWriteTask {
.cache_hash = cache_hash,
.file_offset = 0,
.write_size = buffer->size(),
.buffer = std::move(buffer),
.admission_ctx = {},
.submit_ts_us = submit_ts_us,
.write_epoch = manager->current_write_epoch(cache_hash),
.on_finalized = std::move(on_finalized),
};
}
bool is_cache_range_downloaded(BlockFileCache* cache, const UInt128Wrapper& hash, size_t offset = 0,
size_t size = 4096) {
DORIS_CHECK(cache != nullptr);
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
FileBlocks blocks;
bool fully_covered = false;
DORIS_CHECK(cache->get_downloaded_blocks_if_fully_covered(hash, offset, size, context, &blocks,
&fully_covered)
.ok());
return fully_covered;
}
class AsyncCacheWriteManagerTest : public BlockFileCacheTest {
protected:
std::unique_ptr<BlockFileCache> create_cache(const std::string& name) {
auto path = caches_dir / name;
std::error_code error;
fs::remove_all(path, error);
fs::create_directories(path);
_paths.emplace_back(path);
auto cache = std::make_unique<BlockFileCache>(path.string(), async_write_cache_settings());
EXPECT_TRUE(cache->initialize().ok());
wait_until_cache_ready(*cache);
EXPECT_TRUE(cache->async_write_manager()->start().ok());
return cache;
}
void TearDown() override {
for (const auto& path : _paths) {
std::error_code error;
fs::remove_all(path, error);
}
}
private:
std::vector<fs::path> _paths;
};
TEST_F(AsyncCacheWriteManagerTest, TrackedBufferAllocationCatchesAllocatorFailure) {
auto cache = create_cache("async_write_manager_buffer_allocation_failure");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
const uint64_t baseline_failures = manager->_metrics->snapshot().buffer_alloc_fail;
const double old_fault_probability = config::mem_alloc_fault_probability;
Defer restore_fault_probability {
[&]() { config::mem_alloc_fault_probability = old_fault_probability; }};
config::mem_alloc_fault_probability = 1.0;
AsyncCacheWriteBufferPtr buffer;
const Status status = manager->allocate_tracked_buffer(4096, &buffer);
config::mem_alloc_fault_probability = old_fault_probability;
EXPECT_TRUE(status.is<ErrorCode::MEM_LIMIT_EXCEEDED>());
EXPECT_EQ(buffer, nullptr);
EXPECT_EQ(manager->_metrics->snapshot().buffer_alloc_fail, baseline_failures + 1);
}
class OneShotSyncPointGate {
public:
void arrive_and_wait() {
std::unique_lock lock(_mutex);
if (_arrived) {
return;
}
_arrived = true;
_cv.notify_all();
_cv.wait(lock, [&]() { return _released; });
}
bool wait_until_arrived() {
std::unique_lock lock(_mutex);
return _cv.wait_for(lock, std::chrono::seconds(5), [&]() { return _arrived; });
}
void release() {
{
std::lock_guard lock(_mutex);
_released = true;
}
_cv.notify_all();
}
private:
std::mutex _mutex;
std::condition_variable _cv;
bool _arrived {false};
bool _released {false};
};
TEST(AsyncCacheWriteConfigTest, ResolveMaxPendingBytes) {
constexpr int64_t kMiB = 1024 * 1024;
constexpr int64_t kGiB = 1024 * kMiB;
size_t resolved_bytes = 0;
ASSERT_TRUE(resolve_async_file_cache_write_max_pending_bytes(123 * kMiB, 100 * kGiB,
&resolved_bytes)
.ok());
EXPECT_EQ(resolved_bytes, 123 * kMiB);
ASSERT_TRUE(
resolve_async_file_cache_write_max_pending_bytes(-1, 32 * kGiB, &resolved_bytes).ok());
EXPECT_EQ(resolved_bytes, 1 * kGiB);
ASSERT_TRUE(
resolve_async_file_cache_write_max_pending_bytes(-1, 200 * kGiB, &resolved_bytes).ok());
EXPECT_EQ(resolved_bytes, 2 * kGiB);
EXPECT_FALSE(
resolve_async_file_cache_write_max_pending_bytes(0, 100 * kGiB, &resolved_bytes).ok());
EXPECT_FALSE(
resolve_async_file_cache_write_max_pending_bytes(-2, 100 * kGiB, &resolved_bytes).ok());
}
TEST_F(AsyncCacheWriteManagerTest, WriteEpochRegistryReusesAndReclaimsLiveKey) {
auto cache = create_cache("async_write_epoch_registry_lifecycle");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
EXPECT_EQ(manager->active_write_epoch_key_count(), 0);
const auto hash = BlockFileCache::hash("epoch_registry_lifecycle");
{
const auto first = manager->current_write_epoch(hash);
EXPECT_TRUE(manager->is_current_write_epoch(first));
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
{
const auto second = manager->current_write_epoch(hash);
EXPECT_EQ(second.key_token, first.key_token);
EXPECT_EQ(second.cache_epoch, first.cache_epoch);
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
}
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
}
EXPECT_EQ(manager->active_write_epoch_key_count(), 0);
const auto old_epoch = manager->current_write_epoch(hash);
manager->invalidate_pending_writes(hash);
EXPECT_FALSE(manager->is_current_write_epoch(old_epoch));
EXPECT_EQ(manager->active_write_epoch_key_count(), 0);
const auto new_epoch = manager->current_write_epoch(hash);
EXPECT_TRUE(manager->is_current_write_epoch(new_epoch));
EXPECT_EQ(new_epoch.cache_epoch, old_epoch.cache_epoch);
EXPECT_LT(old_epoch.key_token->generation(), new_epoch.key_token->generation());
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
}
TEST_F(AsyncCacheWriteManagerTest, WriteEpochRegistryUnwindsPublicationFailures) {
auto cache = create_cache("async_write_epoch_registry_exception");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
const auto hash = BlockFileCache::hash("epoch_registry_exception");
auto* sync_point = SyncPoint::get_instance();
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto expect_capture_failure = [&](const std::string& point) {
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
point, [](auto&&) { throw std::bad_alloc(); }, &guard);
EXPECT_THROW(static_cast<void>(manager->current_write_epoch(hash)), std::bad_alloc);
EXPECT_EQ(manager->active_write_epoch_key_count(), 0);
};
expect_capture_failure("AsyncCacheWriteEpochRegistry::capture:after_candidate_object_created");
expect_capture_failure("AsyncCacheWriteEpochRegistry::capture:before_candidate_publish");
{
const auto epoch = manager->current_write_epoch(hash);
EXPECT_TRUE(manager->is_current_write_epoch(epoch));
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
}
EXPECT_EQ(manager->active_write_epoch_key_count(), 0);
}
TEST_F(AsyncCacheWriteManagerTest, ConcurrentWriteEpochCapturePublishesSingleToken) {
auto cache = create_cache("async_write_epoch_registry_concurrent");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
const auto hash = BlockFileCache::hash("epoch_registry_concurrent");
constexpr size_t thread_count = 16;
std::barrier start(static_cast<std::ptrdiff_t>(thread_count));
std::vector<AsyncCacheWriteEpoch> epochs(thread_count);
std::vector<std::thread> threads;
threads.reserve(thread_count);
for (size_t index = 0; index < thread_count; ++index) {
threads.emplace_back([&, index]() {
start.arrive_and_wait();
epochs[index] = manager->current_write_epoch(hash);
});
}
for (auto& thread : threads) {
thread.join();
}
ASSERT_NE(epochs.front().key_token, nullptr);
for (const auto& epoch : epochs) {
EXPECT_EQ(epoch.key_token, epochs.front().key_token);
EXPECT_EQ(epoch.cache_epoch, epochs.front().cache_epoch);
}
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
epochs.clear();
EXPECT_EQ(manager->active_write_epoch_key_count(), 0);
}
TEST_F(AsyncCacheWriteManagerTest, PerKeyRemoveDoesNotInvalidateUnrelatedWriteEpoch) {
auto cache = create_cache("async_write_epoch_scope");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
const auto target_hash = BlockFileCache::hash("epoch_scope_target");
const auto target_epoch = manager->current_write_epoch(target_hash);
const uint64_t cache_epoch = manager->current_cache_epoch();
const uint64_t baseline_key_invalidations = manager->_metrics->snapshot().key_epoch_invalidate;
for (size_t index = 0; index < 128; ++index) {
cache->remove_if_cached_async(
BlockFileCache::hash("unrelated_removed_" + std::to_string(index)));
}
EXPECT_EQ(manager->current_cache_epoch(), cache_epoch);
EXPECT_TRUE(manager->is_current_write_epoch(target_epoch));
EXPECT_EQ(manager->active_write_epoch_key_count(), 1);
const auto removed_hash = BlockFileCache::hash("epoch_scope_removed");
const auto removed_epoch = manager->current_write_epoch(removed_hash);
cache->remove_if_cached_async(removed_hash);
EXPECT_FALSE(manager->is_current_write_epoch(removed_epoch));
EXPECT_TRUE(manager->is_current_write_epoch(target_epoch));
EXPECT_EQ(manager->_metrics->snapshot().key_epoch_invalidate - baseline_key_invalidations, 129);
const auto replacement_epoch = manager->current_write_epoch(removed_hash);
const uint64_t baseline_cache_invalidations =
manager->_metrics->snapshot().cache_epoch_invalidate;
EXPECT_EQ(manager->invalidate_all_pending_writes(), cache_epoch + 1);
EXPECT_FALSE(manager->is_current_write_epoch(target_epoch));
EXPECT_FALSE(manager->is_current_write_epoch(replacement_epoch));
EXPECT_EQ(manager->_metrics->snapshot().cache_epoch_invalidate - baseline_cache_invalidations,
1);
}
TEST_F(AsyncCacheWriteManagerTest, TaskWritesDownloadedBlockAndCleansInflightEntry) {
auto cache = create_cache("async_write_manager_single_task");
auto* manager = cache->async_write_manager();
auto* index = cache->inflight_write_buffer_index();
ASSERT_NE(manager, nullptr);
ASSERT_NE(index, nullptr);
const uint64_t baseline_submitted = manager->_metrics->snapshot().submitted;
const uint64_t baseline_submitted_bytes = manager->_metrics->snapshot().submitted_bytes;
const uint64_t baseline_finished = manager->_metrics->snapshot().finished;
const uint64_t baseline_finished_bytes = manager->_metrics->snapshot().finished_bytes;
const uint64_t baseline_worker_finished_bytes =
manager->_metrics->snapshot().worker_finished_bytes;
const int64_t baseline_submit_latency_count =
manager->_metrics->snapshot().submit_latency_count;
const int64_t baseline_buffer_alloc_latency_count =
manager->_metrics->snapshot().buffer_alloc_latency_count;
const int64_t baseline_queue_wait_latency_count =
manager->_metrics->snapshot().queue_wait_latency_count;
const int64_t baseline_worker_task_latency_count =
manager->_metrics->snapshot().worker_task_latency_count;
const int64_t baseline_get_or_set_latency_count =
manager->_metrics->snapshot().get_or_set_latency_count;
const int64_t baseline_append_latency_count =
manager->_metrics->snapshot().append_latency_count;
const int64_t baseline_finalize_latency_count =
manager->_metrics->snapshot().finalize_latency_count;
constexpr size_t block_size = 4096;
const auto hash = BlockFileCache::hash("async_single_task");
const std::string payload(block_size, 'a');
const int64_t baseline_memory = manager->buffer_memory_bytes();
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(block_size, &buffer).ok());
memcpy(buffer->data(), payload.data(), payload.size());
EXPECT_GE(manager->buffer_memory_bytes(), baseline_memory + static_cast<int64_t>(block_size));
const auto epoch = manager->current_write_epoch(hash);
auto entry =
std::make_shared<InflightWriteBufferEntry>(buffer, 0, block_size, MonotonicMicros());
ASSERT_EQ(index->insert_if_absent(hash, 0, entry), nullptr);
std::promise<void> finished;
auto finished_future = finished.get_future();
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = block_size,
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = epoch,
.on_finalized =
[index, hash, entry, &finished](const AsyncCacheWriteTask&) {
index->remove_if(hash, 0, entry);
finished.set_value();
},
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->pending_bytes(), 0);
EXPECT_EQ(manager->queued_count(), 0);
EXPECT_EQ(manager->queued_bytes(), 0);
EXPECT_EQ(manager->active_task_count(), 0);
EXPECT_EQ(manager->active_bytes(), 0);
EXPECT_EQ(index->lookup(hash, 0), nullptr);
EXPECT_EQ(manager->_metrics->snapshot().submitted - baseline_submitted, 1);
EXPECT_EQ(manager->_metrics->snapshot().submitted_bytes - baseline_submitted_bytes, block_size);
EXPECT_EQ(manager->_metrics->snapshot().finished - baseline_finished, 1);
EXPECT_EQ(manager->_metrics->snapshot().finished_bytes - baseline_finished_bytes, block_size);
EXPECT_EQ(manager->_metrics->snapshot().worker_finished_bytes - baseline_worker_finished_bytes,
block_size);
EXPECT_EQ(manager->_metrics->snapshot().submit_latency_count - baseline_submit_latency_count,
1);
EXPECT_EQ(manager->_metrics->snapshot().buffer_alloc_latency_count -
baseline_buffer_alloc_latency_count,
1);
EXPECT_EQ(manager->_metrics->snapshot().queue_wait_latency_count -
baseline_queue_wait_latency_count,
1);
EXPECT_EQ(manager->_metrics->snapshot().worker_task_latency_count -
baseline_worker_task_latency_count,
1);
EXPECT_EQ(manager->_metrics->snapshot().get_or_set_latency_count -
baseline_get_or_set_latency_count,
1);
EXPECT_EQ(manager->_metrics->snapshot().append_latency_count - baseline_append_latency_count,
1);
EXPECT_EQ(
manager->_metrics->snapshot().finalize_latency_count - baseline_finalize_latency_count,
1);
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
FileBlocks blocks;
bool fully_covered = false;
ASSERT_TRUE(cache->get_downloaded_blocks_if_fully_covered(hash, 0, block_size, context, &blocks,
&fully_covered)
.ok());
ASSERT_TRUE(fully_covered);
ASSERT_EQ(blocks.size(), 1);
std::string actual(block_size, '\0');
ASSERT_TRUE(blocks.front()->read(Slice(actual.data(), actual.size()), 0).ok());
EXPECT_EQ(actual, payload);
blocks.clear();
entry.reset();
buffer.reset();
for (int attempt = 0; attempt < 100 && manager->buffer_memory_bytes() != baseline_memory;
++attempt) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
EXPECT_EQ(manager->buffer_memory_bytes(), baseline_memory);
}
TEST_F(AsyncCacheWriteManagerTest, LockedQueuePreservesFifoWithSingleWorker) {
auto cache = create_cache("async_write_manager_locked_fifo");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 4 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const std::vector<UInt128Wrapper> hashes {BlockFileCache::hash("locked_fifo_a"),
BlockFileCache::hash("locked_fifo_b"),
BlockFileCache::hash("locked_fifo_c")};
std::mutex mutex;
std::condition_variable cv;
std::vector<size_t> take_order;
size_t finalized = 0;
bool release_worker = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&& args) {
const auto* task = try_any_cast<const AsyncCacheWriteTask*>(args[0]);
const auto iterator = std::find(hashes.begin(), hashes.end(), task->cache_hash);
DORIS_CHECK(iterator != hashes.end());
std::unique_lock lock(mutex);
take_order.emplace_back(static_cast<size_t>(iterator - hashes.begin()));
cv.notify_all();
cv.wait(lock, [&]() { return release_worker; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto on_finalized = [&](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finalized;
cv.notify_all();
};
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "locked_fifo_a", 'a', on_finalized)));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5),
[&]() { return take_order.size() == 1; }));
}
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "locked_fifo_b", 'b', on_finalized)));
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "locked_fifo_c", 'c', on_finalized)));
EXPECT_EQ(manager->queued_count(), 2);
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return finalized == 3; }));
}
EXPECT_EQ(take_order, (std::vector<size_t> {0, 1, 2}));
EXPECT_EQ(manager->pending_count(), 0);
}
TEST_F(AsyncCacheWriteManagerTest, WorkerContinuesAfterTaskException) {
auto cache = create_cache("async_write_manager_worker_exception");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 2 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const auto failed_hash = BlockFileCache::hash("worker_exception_failed");
const auto later_hash = BlockFileCache::hash("worker_exception_later");
std::atomic<bool> exception_injected {false};
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&& args) {
const auto* task = try_any_cast<const AsyncCacheWriteTask*>(args[0]);
if (task->cache_hash == failed_hash && !exception_injected.exchange(true)) {
throw Exception(ErrorCode::INTERNAL_ERROR, "injected worker task exception");
}
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
std::mutex mutex;
std::condition_variable cv;
size_t finalized_count = 0;
const auto on_finalized = [&](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finalized_count;
cv.notify_all();
};
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "worker_exception_failed", 'f', on_finalized)));
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "worker_exception_later", 'l', on_finalized)));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(
cv.wait_for(lock, std::chrono::seconds(5), [&]() { return finalized_count == 2; }));
}
EXPECT_TRUE(exception_injected.load());
EXPECT_FALSE(is_cache_range_downloaded(cache.get(), failed_hash));
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), later_hash));
EXPECT_EQ(manager->running_worker_count(), 1);
auto shutdown_future = std::async(std::launch::async, [manager]() { manager->shutdown(); });
ASSERT_EQ(shutdown_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
shutdown_future.get();
}
TEST_F(AsyncCacheWriteManagerTest, RejectsWhenAllPendingTasksAreActive) {
auto cache = create_cache("async_write_manager_backpressure");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.max_pending_bytes = 4096;
ASSERT_TRUE(manager->update_options(options).ok());
std::mutex mutex;
std::condition_variable cv;
bool worker_entered = false;
bool release_worker = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&&) {
std::unique_lock lock(mutex);
worker_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_worker; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
AsyncCacheWriteBufferPtr first_buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &first_buffer).ok());
memset(first_buffer->data(), 'b', first_buffer->size());
std::promise<void> first_finished;
auto first_future = first_finished.get_future();
const auto first_hash = BlockFileCache::hash("backpressure_first");
AsyncCacheWriteTask first_task {
.cache_hash = first_hash,
.file_offset = 0,
.write_size = first_buffer->size(),
.buffer = first_buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(first_hash),
.on_finalized =
[&first_finished](const AsyncCacheWriteTask&) { first_finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(first_task)));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entered; }));
}
AsyncCacheWriteBufferPtr rejected_buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &rejected_buffer).ok());
const auto rejected_hash = BlockFileCache::hash("backpressure_rejected");
AsyncCacheWriteTask rejected_task {
.cache_hash = rejected_hash,
.file_offset = 0,
.write_size = rejected_buffer->size(),
.buffer = rejected_buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(rejected_hash),
.on_finalized = nullptr,
};
EXPECT_FALSE(manager->try_submit(std::move(rejected_task)));
EXPECT_EQ(manager->pending_count(), 1);
EXPECT_EQ(manager->pending_bytes(), 4096);
EXPECT_EQ(manager->queued_count(), 0);
EXPECT_EQ(manager->queued_bytes(), 0);
EXPECT_EQ(manager->active_task_count(), 1);
EXPECT_EQ(manager->active_bytes(), 4096);
EXPECT_EQ(manager->_active_get_or_set_count.load(std::memory_order_relaxed), 1);
EXPECT_GE(manager->_metrics->snapshot().reject_backpressure, 1);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest, 0);
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
ASSERT_EQ(first_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->pending_bytes(), 0);
}
TEST_F(AsyncCacheWriteManagerTest, RejectsTaskLargerThanPendingMemoryLimit) {
auto cache = create_cache("async_write_manager_task_too_large");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.max_pending_bytes = 4095;
ASSERT_TRUE(manager->update_options(options).ok());
const uint64_t baseline_rejected = manager->_metrics->snapshot().rejected;
const uint64_t baseline_backpressure = manager->_metrics->snapshot().reject_backpressure;
size_t finalized = 0;
EXPECT_FALSE(manager->try_submit(make_async_write_task(
manager, "task_too_large", 'l', [&](const AsyncCacheWriteTask&) { ++finalized; })));
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->pending_bytes(), 0);
EXPECT_EQ(manager->_metrics->snapshot().rejected - baseline_rejected, 1);
EXPECT_EQ(manager->_metrics->snapshot().reject_backpressure - baseline_backpressure, 1);
EXPECT_EQ(finalized, 0);
}
TEST_F(AsyncCacheWriteManagerTest,
DropOldestReplacesOnlyOldestQueuedTaskAndKeepsInflightReaderAlive) {
auto cache = create_cache("async_write_manager_drop_oldest");
auto* manager = cache->async_write_manager();
auto* index = cache->inflight_write_buffer_index();
ASSERT_NE(manager, nullptr);
ASSERT_NE(index, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 3 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const int64_t baseline_memory = manager->buffer_memory_bytes();
const std::vector<std::string> keys {"drop_oldest_a", "drop_oldest_b", "drop_oldest_c",
"drop_oldest_d"};
std::vector<UInt128Wrapper> hashes;
hashes.reserve(keys.size());
for (const auto& key : keys) {
hashes.emplace_back(BlockFileCache::hash(key));
}
std::mutex mutex;
std::condition_variable cv;
bool worker_entered = false;
bool release_worker = false;
std::vector<size_t> finalized(keys.size(), 0);
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&&) {
std::unique_lock lock(mutex);
worker_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_worker; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto record_finalized = [&](size_t task_id) {
return [&, task_id](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finalized[task_id];
cv.notify_all();
};
};
ASSERT_TRUE(
manager->try_submit(make_async_write_task(manager, keys[0], 'a', record_finalized(0))));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entered; }));
}
std::vector<std::shared_ptr<InflightWriteBufferEntry>> entries(keys.size());
const auto make_indexed_task = [&](size_t task_id, char fill) {
AsyncCacheWriteBufferPtr buffer;
EXPECT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
DORIS_CHECK(buffer != nullptr);
memset(buffer->data(), fill, buffer->size());
const int64_t submit_ts_us = MonotonicMicros();
auto write_epoch = manager->current_write_epoch(hashes[task_id]);
entries[task_id] =
std::make_shared<InflightWriteBufferEntry>(buffer, 0, buffer->size(), submit_ts_us);
DORIS_CHECK(index->insert_if_absent(hashes[task_id], 0, entries[task_id]) == nullptr);
return AsyncCacheWriteTask {
.cache_hash = hashes[task_id],
.file_offset = 0,
.write_size = buffer->size(),
.buffer = std::move(buffer),
.admission_ctx = {},
.submit_ts_us = submit_ts_us,
.write_epoch = std::move(write_epoch),
.on_finalized =
[&, task_id](const AsyncCacheWriteTask&) {
index->remove_if(hashes[task_id], 0, entries[task_id]);
std::lock_guard lock(mutex);
++finalized[task_id];
cv.notify_all();
},
};
};
ASSERT_TRUE(manager->try_submit(make_indexed_task(1, 'b')));
ASSERT_TRUE(manager->try_submit(make_indexed_task(2, 'c')));
ASSERT_EQ(manager->pending_count(), 3);
ASSERT_EQ(manager->pending_bytes(), 3 * 4096);
ASSERT_EQ(manager->queued_count(), 2);
ASSERT_EQ(manager->queued_bytes(), 2 * 4096);
ASSERT_EQ(manager->active_task_count(), 1);
ASSERT_EQ(manager->active_bytes(), 4096);
auto concurrent_reader = index->lookup(hashes[1], 0);
ASSERT_NE(concurrent_reader, nullptr);
const uint64_t baseline_evicted = manager->_metrics->snapshot().evicted_oldest;
const uint64_t baseline_evicted_bytes = manager->_metrics->snapshot().evicted_oldest_bytes;
const int64_t baseline_evicted_age_count =
manager->_metrics->snapshot().evicted_oldest_age_count;
ASSERT_TRUE(manager->try_submit(make_indexed_task(3, 'd')));
EXPECT_EQ(manager->pending_count(), 3);
EXPECT_EQ(manager->pending_bytes(), 3 * 4096);
EXPECT_EQ(manager->queued_count(), 2);
EXPECT_EQ(manager->queued_bytes(), 2 * 4096);
EXPECT_EQ(manager->active_task_count(), 1);
EXPECT_EQ(manager->active_bytes(), 4096);
EXPECT_EQ(finalized[0], 0);
EXPECT_EQ(finalized[1], 1);
EXPECT_EQ(finalized[2], 0);
EXPECT_EQ(finalized[3], 0);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest - baseline_evicted, 1);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest_bytes - baseline_evicted_bytes, 4096);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest_age_count - baseline_evicted_age_count,
1);
EXPECT_EQ(index->lookup(hashes[1], 0), nullptr);
EXPECT_NE(index->lookup(hashes[2], 0), nullptr);
EXPECT_NE(index->lookup(hashes[3], 0), nullptr);
ASSERT_NE(concurrent_reader->buffer, nullptr);
EXPECT_EQ(concurrent_reader->buffer->data()[0], 'b');
EXPECT_EQ(concurrent_reader->buffer->data()[4095], 'b');
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() {
return finalized[0] == 1 && finalized[1] == 1 && finalized[2] == 1 && finalized[3] == 1;
}));
}
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), hashes[0]));
EXPECT_FALSE(is_cache_range_downloaded(cache.get(), hashes[1]));
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), hashes[2]));
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), hashes[3]));
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->pending_bytes(), 0);
EXPECT_EQ(index->count(), 0);
concurrent_reader.reset();
entries.clear();
for (int attempt = 0; attempt < 100 && manager->buffer_memory_bytes() != baseline_memory;
++attempt) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
EXPECT_EQ(manager->buffer_memory_bytes(), baseline_memory);
}
TEST_F(AsyncCacheWriteManagerTest, InvalidatedVictimBufferRemainsReadableUntilFinalized) {
auto cache = create_cache("async_write_manager_old_epoch_victim");
auto* manager = cache->async_write_manager();
auto* index = cache->inflight_write_buffer_index();
ASSERT_NE(manager, nullptr);
ASSERT_NE(index, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 2 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
std::mutex mutex;
std::condition_variable cv;
bool worker_entered = false;
bool release_worker = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&&) {
std::unique_lock lock(mutex);
worker_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_worker; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
ASSERT_TRUE(manager->try_submit(make_async_write_task(manager, "old_epoch_active", 'a')));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entered; }));
}
const auto hash = BlockFileCache::hash("old_epoch_same_key");
const auto old_epoch = manager->current_write_epoch(hash);
AsyncCacheWriteBufferPtr old_buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &old_buffer).ok());
memset(old_buffer->data(), 'o', old_buffer->size());
auto old_entry = std::make_shared<InflightWriteBufferEntry>(old_buffer, 0, old_buffer->size(),
MonotonicMicros());
ASSERT_EQ(index->insert_if_absent(hash, 0, old_entry), nullptr);
size_t old_callback_count = 0;
AsyncCacheWriteTask old_task {
.cache_hash = hash,
.file_offset = 0,
.write_size = old_buffer->size(),
.buffer = old_buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = old_epoch,
.on_finalized =
[&, old_entry](const AsyncCacheWriteTask&) {
index->remove_if(hash, 0, old_entry);
++old_callback_count;
},
};
ASSERT_TRUE(manager->try_submit(std::move(old_task)));
manager->invalidate_pending_writes(hash);
EXPECT_FALSE(manager->is_current_write_epoch(old_epoch));
const auto new_epoch = manager->current_write_epoch(hash);
EXPECT_TRUE(manager->is_current_write_epoch(new_epoch));
EXPECT_EQ(new_epoch.cache_epoch, old_epoch.cache_epoch);
EXPECT_LT(old_epoch.key_token->generation(), new_epoch.key_token->generation());
AsyncCacheWriteBufferPtr replacement_buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &replacement_buffer).ok());
memset(replacement_buffer->data(), 'n', replacement_buffer->size());
auto replacement_entry = std::make_shared<InflightWriteBufferEntry>(
replacement_buffer, 0, replacement_buffer->size(), MonotonicMicros());
EXPECT_EQ(index->lookup(hash, 0), old_entry);
EXPECT_EQ(index->insert_if_absent(hash, 0, replacement_entry), old_entry);
const uint64_t baseline_evicted = manager->_metrics->snapshot().evicted_oldest;
ASSERT_TRUE(manager->try_submit(make_async_write_task(manager, "old_epoch_replacer", 'r')));
EXPECT_EQ(old_callback_count, 1);
EXPECT_EQ(index->lookup(hash, 0), nullptr);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest - baseline_evicted, 1);
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
for (int attempt = 0; attempt < 5000 && manager->pending_count() != 0; ++attempt) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
ASSERT_EQ(manager->pending_count(), 0);
EXPECT_EQ(old_callback_count, 1);
EXPECT_EQ(index->count(), 0);
}
TEST_F(AsyncCacheWriteManagerTest, EvictedCallbackRunsOutsideQueueMutex) {
auto cache = create_cache("async_write_manager_callback_outside_lock");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 2 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const auto active_hash = BlockFileCache::hash("callback_outside_active");
const auto replacement_hash = BlockFileCache::hash("callback_outside_replacement");
std::mutex mutex;
std::condition_variable cv;
bool active_entered = false;
bool release_active = false;
bool replacement_entered = false;
bool victim_callback_entered = false;
bool release_victim_callback = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&& args) {
const auto* task = try_any_cast<const AsyncCacheWriteTask*>(args[0]);
std::unique_lock lock(mutex);
if (task->cache_hash == active_hash) {
active_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_active; });
} else if (task->cache_hash == replacement_hash) {
replacement_entered = true;
cv.notify_all();
}
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_active = true;
release_victim_callback = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
ASSERT_TRUE(
manager->try_submit(make_async_write_task(manager, "callback_outside_active", 'a')));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return active_entered; }));
}
ASSERT_TRUE(manager->try_submit(make_async_write_task(
manager, "callback_outside_victim", 'v', [&](const AsyncCacheWriteTask&) {
std::unique_lock lock(mutex);
victim_callback_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_victim_callback; });
})));
auto replacement_future = std::async(std::launch::async, [&]() {
SCOPED_ATTACH_TASK(ExecEnv::GetInstance()->orphan_mem_tracker());
return manager->try_submit(
make_async_write_task(manager, "callback_outside_replacement", 'r'));
});
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5),
[&]() { return victim_callback_entered; }));
release_active = true;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(
cv.wait_for(lock, std::chrono::seconds(5), [&]() { return replacement_entered; }));
}
EXPECT_EQ(replacement_future.wait_for(std::chrono::milliseconds(0)),
std::future_status::timeout);
{
std::lock_guard lock(mutex);
release_victim_callback = true;
}
cv.notify_all();
ASSERT_TRUE(replacement_future.get());
}
TEST_F(AsyncCacheWriteManagerTest, OldQueuedTaskIsStillWrittenAndCleansInflightEntry) {
auto cache = create_cache("async_write_manager_old_queued_task");
auto* manager = cache->async_write_manager();
auto* index = cache->inflight_write_buffer_index();
ASSERT_NE(manager, nullptr);
ASSERT_NE(index, nullptr);
const auto hash = BlockFileCache::hash("old_queued_task");
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
memset(buffer->data(), 'w', buffer->size());
const auto epoch = manager->current_write_epoch(hash);
auto entry = std::make_shared<InflightWriteBufferEntry>(buffer, 0, buffer->size(),
MonotonicMicros());
ASSERT_EQ(index->insert_if_absent(hash, 0, entry), nullptr);
std::promise<void> finished;
auto finished_future = finished.get_future();
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = buffer->size(),
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros() - 60LL * 60 * 1000 * 1000,
.write_epoch = epoch,
.on_finalized =
[index, hash, entry, &finished](const AsyncCacheWriteTask&) {
index->remove_if(hash, 0, entry);
finished.set_value();
},
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(index->lookup(hash, 0), nullptr);
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), hash));
}
TEST_F(AsyncCacheWriteManagerTest, ShutdownDrainsAcceptedTask) {
auto cache = create_cache("async_write_manager_shutdown_drain");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
const auto hash = BlockFileCache::hash("shutdown_drain");
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
memset(buffer->data(), 's', buffer->size());
std::promise<void> finished;
auto finished_future = finished.get_future();
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = buffer->size(),
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(hash),
.on_finalized = [&finished](const AsyncCacheWriteTask&) { finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
manager->shutdown();
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(0)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->running_worker_count(), 0);
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
FileBlocks blocks;
bool fully_covered = false;
ASSERT_TRUE(cache->get_downloaded_blocks_if_fully_covered(hash, 0, 4096, context, &blocks,
&fully_covered)
.ok());
EXPECT_TRUE(fully_covered);
ASSERT_EQ(blocks.size(), 1);
std::string actual(4096, '\0');
ASSERT_TRUE(blocks.front()->read(Slice(actual.data(), actual.size()), 0).ok());
EXPECT_EQ(actual, std::string(4096, 's'));
}
TEST_F(AsyncCacheWriteManagerTest, OneTaskWritesMultipleContainedCells) {
auto cache = create_cache("async_write_manager_multiple_cells");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
constexpr size_t cell_size = 4096;
constexpr size_t task_size = cell_size * 2;
const auto hash = BlockFileCache::hash("multiple_cells");
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(task_size, &buffer).ok());
memset(buffer->data(), 'm', cell_size);
memset(buffer->data() + cell_size, 'n', cell_size);
std::promise<void> finished;
auto finished_future = finished.get_future();
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = task_size,
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(hash),
.on_finalized = [&finished](const AsyncCacheWriteTask&) { finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
FileBlocks blocks;
bool fully_covered = false;
ASSERT_TRUE(cache->get_downloaded_blocks_if_fully_covered(hash, 0, task_size, context, &blocks,
&fully_covered)
.ok());
EXPECT_TRUE(fully_covered);
ASSERT_EQ(blocks.size(), 2);
auto iterator = blocks.begin();
std::string first(cell_size, '\0');
ASSERT_TRUE((*iterator)->read(Slice(first.data(), first.size()), 0).ok());
EXPECT_EQ(first, std::string(cell_size, 'm'));
++iterator;
std::string second(cell_size, '\0');
ASSERT_TRUE((*iterator)->read(Slice(second.data(), second.size()), 0).ok());
EXPECT_EQ(second, std::string(cell_size, 'n'));
}
TEST_F(AsyncCacheWriteManagerTest, ExistingAndDeletingCellsKeepTheirOwners) {
auto cache = create_cache("async_write_manager_existing_states");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
constexpr size_t cell_size = 4096;
constexpr size_t task_size = cell_size * 3;
const auto hash = BlockFileCache::hash("existing_states");
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
auto holder = cache->get_or_set(hash, 0, task_size, context);
ASSERT_EQ(holder.file_blocks.size(), 3);
auto iterator = holder.file_blocks.begin();
const auto downloaded_block = *iterator;
++iterator;
const auto downloading_block = *iterator;
++iterator;
const auto deleting_block = *iterator;
ASSERT_EQ(downloaded_block->get_or_set_downloader(), FileBlock::get_caller_id());
const std::string original(cell_size, 'x');
ASSERT_TRUE(downloaded_block->append(Slice(original.data(), original.size())).ok());
ASSERT_TRUE(downloaded_block->finalize().ok());
ASSERT_EQ(downloading_block->get_or_set_downloader(), FileBlock::get_caller_id());
deleting_block->set_deleting();
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(task_size, &buffer).ok());
memset(buffer->data(), 'y', buffer->size());
std::promise<void> finished;
auto finished_future = finished.get_future();
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = task_size,
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(hash),
.on_finalized = [&finished](const AsyncCacheWriteTask&) { finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_GE(manager->_metrics->snapshot().skip_downloaded, 1);
EXPECT_GE(manager->_metrics->snapshot().skip_downloading, 1);
EXPECT_GE(manager->_metrics->snapshot().skip_deleting, 1);
std::string actual(cell_size, '\0');
ASSERT_TRUE(downloaded_block->read(Slice(actual.data(), actual.size()), 0).ok());
EXPECT_EQ(actual, original);
EXPECT_EQ(downloading_block->state(), FileBlock::State::DOWNLOADING);
EXPECT_EQ(downloading_block->get_downloader(), FileBlock::get_caller_id());
EXPECT_TRUE(cache->is_block_deleting(deleting_block));
}
TEST_F(AsyncCacheWriteManagerTest, RemoveDuringAppendDoesNotLeaveResurrectedCacheData) {
auto cache = create_cache("async_write_manager_remove_during_append");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
std::mutex mutex;
std::condition_variable cv;
bool before_append = false;
bool release_worker = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_append",
[&](auto&&) {
std::unique_lock lock(mutex);
before_append = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_worker; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto hash = BlockFileCache::hash("remove_during_append");
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
memset(buffer->data(), 'r', buffer->size());
std::promise<void> finished;
auto finished_future = finished.get_future();
const auto write_epoch = manager->current_write_epoch(hash);
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = buffer->size(),
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = write_epoch,
.on_finalized = [&finished](const AsyncCacheWriteTask&) { finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return before_append; }));
}
EXPECT_EQ(manager->_active_append_count.load(std::memory_order_relaxed), 1);
fs::path cache_file;
{
ReadStatistics probe_stats;
CacheContext probe_context;
probe_context.stats = &probe_stats;
auto probe_result = cache->probe(hash, 0, 4096, probe_context);
ASSERT_EQ(probe_result.file_blocks.size(), 1);
ASSERT_NE(probe_result.file_blocks[0], nullptr);
EXPECT_EQ(probe_result.file_blocks[0]->state(), FileBlock::State::DOWNLOADING);
cache_file = probe_result.file_blocks[0]->get_cache_file();
}
const uint64_t old_cache_epoch = manager->current_cache_epoch();
cache->remove_if_cached_async(hash);
EXPECT_EQ(manager->current_cache_epoch(), old_cache_epoch);
EXPECT_FALSE(manager->is_current_write_epoch(write_epoch));
const auto new_write_epoch = manager->current_write_epoch(hash);
EXPECT_TRUE(manager->is_current_write_epoch(new_write_epoch));
EXPECT_LT(write_epoch.key_token->generation(), new_write_epoch.key_token->generation());
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->_active_append_count.load(std::memory_order_relaxed), 0);
bool removed = false;
for (int attempt = 0; attempt < 5000; ++attempt) {
ReadStatistics probe_stats;
CacheContext probe_context;
probe_context.stats = &probe_stats;
bool metadata_removed = false;
{
auto probe_result = cache->probe(hash, 0, 4096, probe_context);
metadata_removed =
probe_result.file_blocks.size() == 1 && probe_result.file_blocks[0] == nullptr;
}
if (metadata_removed && !fs::exists(cache_file)) {
removed = true;
break;
}
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
EXPECT_TRUE(removed);
}
TEST_F(AsyncCacheWriteManagerTest, PartialOverlapIsSkippedWithoutOutOfBoundsWrite) {
auto cache = create_cache("async_write_manager_partial_overlap");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
constexpr size_t existing_size = 4096;
const auto hash = BlockFileCache::hash("partial_overlap");
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
{
auto holder = cache->get_or_set(hash, 0, existing_size, context);
ASSERT_EQ(holder.file_blocks.size(), 1);
const auto& block = holder.file_blocks.front();
ASSERT_EQ(block->get_or_set_downloader(), FileBlock::get_caller_id());
const std::string payload(existing_size, 'x');
ASSERT_TRUE(block->append(Slice(payload.data(), payload.size())).ok());
ASSERT_TRUE(block->finalize().ok());
}
constexpr size_t task_offset = 1024;
constexpr size_t task_size = 4096;
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(task_size, &buffer).ok());
memset(buffer->data(), 'y', buffer->size());
std::promise<void> finished;
auto finished_future = finished.get_future();
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = task_offset,
.write_size = task_size,
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(hash),
.on_finalized = [&finished](const AsyncCacheWriteTask&) { finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
ASSERT_EQ(finished_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_GE(manager->_metrics->snapshot().skip_partial_overlap, 1);
FileBlocks blocks;
bool fully_covered = false;
ASSERT_TRUE(cache->get_downloaded_blocks_if_fully_covered(hash, 0, task_offset + task_size,
context, &blocks, &fully_covered)
.ok());
EXPECT_TRUE(fully_covered);
ASSERT_EQ(blocks.size(), 2);
auto iterator = blocks.begin();
std::string existing(existing_size, '\0');
ASSERT_TRUE((*iterator)->read(Slice(existing.data(), existing.size()), 0).ok());
EXPECT_EQ(existing, std::string(existing_size, 'x'));
++iterator;
std::string tail(task_offset, '\0');
ASSERT_TRUE((*iterator)->read(Slice(tail.data(), tail.size()), 0).ok());
EXPECT_EQ(tail, std::string(task_offset, 'y'));
}
TEST_F(AsyncCacheWriteManagerTest, RemoveInvalidatesOnlyMatchingTaskAndCleansEmptyCell) {
auto cache = create_cache("async_write_manager_epoch");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.worker_count = 1;
ASSERT_TRUE(manager->update_options(options).ok());
std::mutex mutex;
std::condition_variable cv;
bool worker_entered = false;
bool release_worker = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&&) {
std::unique_lock lock(mutex);
worker_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_worker; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto active_hash = BlockFileCache::hash("epoch_drop_active");
AsyncCacheWriteBufferPtr active_buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &active_buffer).ok());
memset(active_buffer->data(), 'a', active_buffer->size());
std::promise<void> active_finished;
auto active_future = active_finished.get_future();
const auto active_epoch = manager->current_write_epoch(active_hash);
AsyncCacheWriteTask active_task {
.cache_hash = active_hash,
.file_offset = 0,
.write_size = active_buffer->size(),
.buffer = active_buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = active_epoch,
.on_finalized =
[&active_finished](const AsyncCacheWriteTask&) { active_finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(active_task)));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entered; }));
}
const auto queued_hash = BlockFileCache::hash("epoch_drop_queued");
AsyncCacheWriteBufferPtr queued_buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &queued_buffer).ok());
memset(queued_buffer->data(), 'q', queued_buffer->size());
std::promise<void> queued_finished;
auto queued_future = queued_finished.get_future();
const auto queued_epoch = manager->current_write_epoch(queued_hash);
AsyncCacheWriteTask queued_task {
.cache_hash = queued_hash,
.file_offset = 0,
.write_size = queued_buffer->size(),
.buffer = queued_buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = queued_epoch,
.on_finalized =
[&queued_finished](const AsyncCacheWriteTask&) { queued_finished.set_value(); },
};
ASSERT_TRUE(manager->try_submit(std::move(queued_task)));
ASSERT_EQ(manager->pending_count(), 2);
const uint64_t cache_epoch = manager->current_cache_epoch();
const uint64_t baseline_stale_key = manager->_metrics->snapshot().drop_stale_key_epoch;
cache->remove_if_cached_async(active_hash);
EXPECT_EQ(manager->current_cache_epoch(), cache_epoch);
EXPECT_FALSE(manager->is_current_write_epoch(active_epoch));
EXPECT_TRUE(manager->is_current_write_epoch(queued_epoch));
{
std::lock_guard lock(mutex);
release_worker = true;
}
cv.notify_all();
ASSERT_EQ(active_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
ASSERT_EQ(queued_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_GE(manager->_metrics->snapshot().drop_stale_key_epoch - baseline_stale_key, 1);
ReadStatistics read_stats;
CacheContext context;
context.stats = &read_stats;
const auto expect_cache_gap = [&](const UInt128Wrapper& hash) {
auto probe_result = cache->probe(hash, 0, 4096, context);
ASSERT_EQ(probe_result.file_blocks.size(), 1);
EXPECT_EQ(probe_result.file_blocks[0], nullptr);
};
expect_cache_gap(active_hash);
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), queued_hash));
}
TEST_F(AsyncCacheWriteManagerTest, PendingLimitDecreaseKeepsReplacingOldestQueuedTask) {
auto cache = create_cache("async_write_manager_limit_decrease");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 4 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
std::mutex mutex;
std::condition_variable cv;
size_t worker_entries = 0;
size_t released_entries = 0;
std::vector<size_t> finalized(7, 0);
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&&) {
std::unique_lock lock(mutex);
const size_t current_entry = ++worker_entries;
cv.notify_all();
cv.wait(lock, [&]() { return released_entries >= current_entry; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
released_entries = std::numeric_limits<size_t>::max();
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto finalizer = [&](size_t task_id) {
return [&, task_id](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finalized[task_id];
cv.notify_all();
};
};
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_a", 'a', finalizer(0))));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(
cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entries == 1; }));
}
const auto first_evicted_hash = BlockFileCache::hash("limit_decrease_b");
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_b", 'b', finalizer(1))));
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_c", 'c', finalizer(2))));
const auto second_evicted_hash = BlockFileCache::hash("limit_decrease_d");
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_d", 'd', finalizer(3))));
ASSERT_EQ(manager->pending_count(), 4);
ASSERT_EQ(manager->pending_bytes(), 4 * 4096);
ASSERT_EQ(manager->queued_count(), 3);
ASSERT_EQ(manager->queued_bytes(), 3 * 4096);
options.max_pending_bytes = 2 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const uint64_t baseline_evicted = manager->_metrics->snapshot().evicted_oldest;
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_e", 'e', finalizer(4))));
EXPECT_EQ(finalized[1], 1);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest - baseline_evicted, 1);
EXPECT_EQ(manager->pending_count(), 4);
EXPECT_EQ(manager->pending_bytes(), 4 * 4096);
EXPECT_EQ(manager->queued_count(), 3);
EXPECT_EQ(manager->queued_bytes(), 3 * 4096);
{
std::lock_guard lock(mutex);
released_entries = 1;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(
cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entries == 2; }));
}
EXPECT_EQ(manager->pending_count(), 3);
EXPECT_EQ(manager->pending_bytes(), 3 * 4096);
const auto third_evicted_hash = BlockFileCache::hash("limit_decrease_g");
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_g", 'g', finalizer(5))));
EXPECT_EQ(finalized[3], 1);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest - baseline_evicted, 2);
EXPECT_EQ(manager->pending_count(), 3);
EXPECT_EQ(manager->pending_bytes(), 3 * 4096);
EXPECT_EQ(manager->queued_count(), 2);
EXPECT_EQ(manager->queued_bytes(), 2 * 4096);
{
std::lock_guard lock(mutex);
released_entries = 2;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(
cv.wait_for(lock, std::chrono::seconds(5), [&]() { return worker_entries == 3; }));
}
ASSERT_EQ(manager->pending_count(), 2);
ASSERT_EQ(manager->pending_bytes(), 2 * 4096);
ASSERT_EQ(manager->queued_count(), 1);
ASSERT_EQ(manager->queued_bytes(), 4096);
ASSERT_EQ(manager->active_task_count(), 1);
ASSERT_EQ(manager->active_bytes(), 4096);
const auto replacement_hash = BlockFileCache::hash("limit_decrease_f");
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "limit_decrease_f", 'f', finalizer(6))));
EXPECT_EQ(manager->pending_count(), 2);
EXPECT_EQ(manager->pending_bytes(), 2 * 4096);
EXPECT_EQ(manager->queued_count(), 1);
EXPECT_EQ(manager->queued_bytes(), 4096);
EXPECT_EQ(finalized[5], 1);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest - baseline_evicted, 3);
{
std::lock_guard lock(mutex);
released_entries = std::numeric_limits<size_t>::max();
}
cv.notify_all();
for (int attempt = 0; attempt < 5000 && manager->pending_count() != 0; ++attempt) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
ASSERT_EQ(manager->pending_count(), 0);
ASSERT_EQ(manager->pending_bytes(), 0);
for (size_t finalized_count : finalized) {
EXPECT_EQ(finalized_count, 1);
}
EXPECT_FALSE(is_cache_range_downloaded(cache.get(), first_evicted_hash));
EXPECT_FALSE(is_cache_range_downloaded(cache.get(), second_evicted_hash));
EXPECT_FALSE(is_cache_range_downloaded(cache.get(), third_evicted_hash));
EXPECT_TRUE(is_cache_range_downloaded(cache.get(), replacement_hash));
}
TEST_F(AsyncCacheWriteManagerTest, UpdateOptionsValidatesAndAppliesAtRuntime) {
auto cache = create_cache("async_write_manager_resize");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
auto invalid_options = options;
invalid_options.worker_count = 0;
EXPECT_TRUE(manager->update_options(invalid_options).is<ErrorCode::INVALID_ARGUMENT>());
invalid_options = options;
invalid_options.max_pending_bytes = 0;
EXPECT_TRUE(manager->update_options(invalid_options).is<ErrorCode::INVALID_ARGUMENT>());
options.worker_count = 3;
options.max_pending_bytes = 7 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const auto updated = manager->options();
EXPECT_EQ(updated.worker_count, 3);
EXPECT_EQ(updated.max_pending_bytes, 7 * 4096);
EXPECT_EQ(manager->_worker_pool->min_threads(), 3);
EXPECT_EQ(manager->_worker_pool->max_threads(), 3);
options.worker_count = 1;
ASSERT_TRUE(manager->update_options(options).ok());
EXPECT_EQ(manager->options().worker_count, 1);
EXPECT_EQ(manager->_worker_pool->min_threads(), 1);
EXPECT_EQ(manager->_worker_pool->max_threads(), 1);
}
TEST_F(AsyncCacheWriteManagerTest, WorkerStopPublishesUnderQueueMutex) {
OneShotSyncPointGate worker_before_wait;
OneShotSyncPointGate stop_before_request;
std::promise<void> stop_requested;
auto stop_requested_future = stop_requested.get_future();
std::unique_ptr<BlockFileCache> cache;
std::future<Status> resize_future;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard worker_guard;
SyncPoint::CallbackGuard resize_before_guard;
SyncPoint::CallbackGuard resize_after_guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::Worker::_run:before_wait",
[&](auto&&) { worker_before_wait.arrive_and_wait(); }, &worker_guard);
sync_point->set_call_back(
"AsyncCacheWriteManager::_stop_workers_locked:before_request_stop",
[&](auto&&) { stop_before_request.arrive_and_wait(); }, &resize_before_guard);
sync_point->set_call_back(
"AsyncCacheWriteManager::_stop_workers_locked:after_request_stop",
[&](auto&&) { stop_requested.set_value(); }, &resize_after_guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
worker_before_wait.release();
stop_before_request.release();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
cache = create_cache("async_write_manager_resize_wait_wakeup");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
ASSERT_GT(manager->options().worker_count, 1);
ASSERT_TRUE(worker_before_wait.wait_until_arrived());
resize_future =
std::async(std::launch::async, [manager]() { return manager->resize_workers(1); });
ASSERT_TRUE(stop_before_request.wait_until_arrived());
stop_before_request.release();
EXPECT_EQ(stop_requested_future.wait_for(std::chrono::milliseconds(500)),
std::future_status::timeout);
worker_before_wait.release();
ASSERT_EQ(stop_requested_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
ASSERT_EQ(resize_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
ASSERT_TRUE(resize_future.get().ok());
EXPECT_EQ(manager->options().worker_count, 1);
}
TEST_F(AsyncCacheWriteManagerTest, ResizeWorkersPreservesActiveTaskOwnership) {
auto cache = create_cache("async_write_manager_resize_workers");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
constexpr size_t worker_count = 8;
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = worker_count * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
options.worker_count = worker_count;
ASSERT_TRUE(manager->update_options(options).ok());
std::mutex mutex;
std::condition_variable cv;
size_t entered_workers = 0;
size_t finished_tasks = 0;
bool release_workers = false;
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&&) {
std::unique_lock lock(mutex);
++entered_workers;
cv.notify_all();
cv.wait(lock, [&]() { return release_workers; });
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_workers = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
for (size_t task_id = 0; task_id < worker_count; ++task_id) {
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
memset(buffer->data(), static_cast<int>('a' + task_id), buffer->size());
const auto hash = BlockFileCache::hash("resize_worker_" + std::to_string(task_id));
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = buffer->size(),
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(hash),
.on_finalized =
[&](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finished_tasks;
cv.notify_all();
},
};
ASSERT_TRUE(manager->try_submit(std::move(task)));
}
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5),
[&]() { return entered_workers == worker_count; }));
EXPECT_EQ(manager->pending_count(), worker_count);
}
auto shrink_options = manager->options();
shrink_options.worker_count = 1;
auto shrink_future = std::async(std::launch::async, [manager, shrink_options]() {
return manager->update_options(shrink_options);
});
for (int attempt = 0; attempt < 5000 && manager->options().worker_count != 1; ++attempt) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
EXPECT_EQ(manager->options().worker_count, 1);
EXPECT_EQ(shrink_future.wait_for(std::chrono::milliseconds(0)), std::future_status::timeout);
{
std::lock_guard lock(mutex);
release_workers = true;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5),
[&]() { return finished_tasks == worker_count; }));
}
ASSERT_TRUE(shrink_future.get().ok());
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->options().worker_count, 1);
}
TEST_F(AsyncCacheWriteManagerTest, ConcurrentDropOldestMaintainsCounterConservation) {
auto cache = create_cache("async_write_manager_concurrent_drop_oldest");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
constexpr size_t producer_count = 4;
constexpr size_t tasks_per_producer = 64;
constexpr size_t producer_tasks = producer_count * tasks_per_producer;
constexpr size_t total_tasks = producer_tasks + 1;
constexpr size_t max_pending_blocks = 16;
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = max_pending_blocks * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const auto active_hash = BlockFileCache::hash("concurrent_drop_oldest_0");
std::mutex mutex;
std::condition_variable cv;
bool active_entered = false;
bool release_active = false;
size_t worker_tasks = 0;
std::vector<size_t> finalized(total_tasks, 0);
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&& args) {
const auto* task = try_any_cast<const AsyncCacheWriteTask*>(args[0]);
std::unique_lock lock(mutex);
++worker_tasks;
if (task->cache_hash == active_hash) {
active_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_active; });
}
},
&guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_active = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
std::vector<AsyncCacheWriteTask> tasks;
tasks.reserve(total_tasks);
for (size_t task_id = 0; task_id < total_tasks; ++task_id) {
tasks.emplace_back(make_async_write_task(
manager, "concurrent_drop_oldest_" + std::to_string(task_id),
static_cast<char>('a' + task_id % 26), [&, task_id](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finalized[task_id];
cv.notify_all();
}));
}
const uint64_t baseline_submitted = manager->_metrics->snapshot().submitted;
const uint64_t baseline_submitted_bytes = manager->_metrics->snapshot().submitted_bytes;
const uint64_t baseline_finished = manager->_metrics->snapshot().finished;
const uint64_t baseline_finished_bytes = manager->_metrics->snapshot().finished_bytes;
const uint64_t baseline_worker_finished = manager->_metrics->snapshot().worker_finished;
const uint64_t baseline_worker_finished_bytes =
manager->_metrics->snapshot().worker_finished_bytes;
const uint64_t baseline_evicted = manager->_metrics->snapshot().evicted_oldest;
const uint64_t baseline_evicted_bytes = manager->_metrics->snapshot().evicted_oldest_bytes;
const uint64_t baseline_rejected = manager->_metrics->snapshot().rejected;
ASSERT_TRUE(manager->try_submit(std::move(tasks[0])));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return active_entered; }));
}
std::atomic<size_t> accepted {0};
std::vector<std::thread> producers;
producers.reserve(producer_count);
std::barrier start_barrier(producer_count);
for (size_t producer_id = 0; producer_id < producer_count; ++producer_id) {
producers.emplace_back([&, producer_id]() {
start_barrier.arrive_and_wait();
const size_t first_task = 1 + producer_id * tasks_per_producer;
for (size_t offset = 0; offset < tasks_per_producer; ++offset) {
if (manager->try_submit(std::move(tasks[first_task + offset]))) {
accepted.fetch_add(1, std::memory_order_relaxed);
}
}
});
}
for (auto& producer : producers) {
producer.join();
}
EXPECT_EQ(accepted.load(std::memory_order_relaxed), producer_tasks);
EXPECT_EQ(manager->pending_count(), max_pending_blocks);
EXPECT_EQ(manager->pending_bytes(), max_pending_blocks * 4096);
EXPECT_EQ(manager->queued_count(), max_pending_blocks - 1);
EXPECT_EQ(manager->queued_bytes(), (max_pending_blocks - 1) * 4096);
EXPECT_EQ(manager->active_task_count(), 1);
EXPECT_EQ(manager->active_bytes(), 4096);
EXPECT_EQ(manager->_metrics->snapshot().submitted - baseline_submitted, total_tasks);
EXPECT_EQ(manager->_metrics->snapshot().submitted_bytes - baseline_submitted_bytes,
total_tasks * 4096);
EXPECT_EQ(manager->_metrics->snapshot().rejected - baseline_rejected, 0);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest - baseline_evicted,
total_tasks - max_pending_blocks);
EXPECT_EQ(manager->_metrics->snapshot().evicted_oldest_bytes - baseline_evicted_bytes,
(total_tasks - max_pending_blocks) * 4096);
{
std::lock_guard lock(mutex);
release_active = true;
}
cv.notify_all();
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(10), [&]() {
return std::all_of(finalized.begin(), finalized.end(),
[](size_t count) { return count == 1; });
}));
}
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->pending_bytes(), 0);
EXPECT_EQ(manager->queued_count(), 0);
EXPECT_EQ(manager->queued_bytes(), 0);
EXPECT_EQ(manager->active_task_count(), 0);
EXPECT_EQ(manager->active_bytes(), 0);
EXPECT_EQ(worker_tasks, max_pending_blocks);
EXPECT_EQ(manager->_metrics->snapshot().finished - baseline_finished, total_tasks);
EXPECT_EQ(manager->_metrics->snapshot().finished_bytes - baseline_finished_bytes,
total_tasks * 4096);
EXPECT_EQ(manager->_metrics->snapshot().worker_finished - baseline_worker_finished,
max_pending_blocks);
EXPECT_EQ(manager->_metrics->snapshot().worker_finished_bytes - baseline_worker_finished_bytes,
max_pending_blocks * 4096);
EXPECT_EQ((manager->_metrics->snapshot().worker_finished - baseline_worker_finished) +
(manager->_metrics->snapshot().evicted_oldest - baseline_evicted),
manager->_metrics->snapshot().submitted - baseline_submitted);
}
TEST_F(AsyncCacheWriteManagerTest, MutableConfigUpdatesManagersExplicitly) {
auto* factory = FileCacheFactory::instance();
factory->_caches.clear();
factory->_path_to_cache.clear();
factory->_capacity = 0;
const bool old_enable = config::enable_async_file_cache_write;
const int32_t old_workers = config::async_file_cache_write_workers_per_disk;
const int64_t old_max_pending_bytes = config::async_file_cache_write_max_pending_bytes;
const auto path1 = caches_dir / "async_write_manager_config_update_1";
const auto path2 = caches_dir / "async_write_manager_config_update_2";
std::error_code error;
Defer restore {[&]() {
EXPECT_TRUE(config::set_config("async_file_cache_write_workers_per_disk",
std::to_string(old_workers))
.ok());
EXPECT_TRUE(config::set_config("async_file_cache_write_max_pending_bytes",
std::to_string(old_max_pending_bytes))
.ok());
EXPECT_TRUE(
config::set_config("enable_async_file_cache_write", old_enable ? "true" : "false")
.ok());
factory->_caches.clear();
factory->_path_to_cache.clear();
factory->_capacity = 0;
fs::remove_all(path1, error);
fs::remove_all(path2, error);
}};
ASSERT_TRUE(config::set_config("enable_async_file_cache_write", "false").ok());
ASSERT_TRUE(config::set_config("async_file_cache_write_max_pending_bytes", "-1").ok());
const int32_t valid_workers = config::async_file_cache_write_workers_per_disk;
EXPECT_FALSE(config::set_config("async_file_cache_write_workers_per_disk", "0").ok());
EXPECT_EQ(config::async_file_cache_write_workers_per_disk, valid_workers);
EXPECT_FALSE(config::set_config("async_file_cache_write_max_pending_bytes", "0").ok());
EXPECT_EQ(config::async_file_cache_write_max_pending_bytes, -1);
fs::remove_all(path1, error);
fs::remove_all(path2, error);
fs::create_directories(path1);
fs::create_directories(path2);
ASSERT_TRUE(factory->create_file_cache(path1.string(), async_write_cache_settings()).ok());
auto* cache1 = factory->get_by_path(path1.string());
ASSERT_NE(cache1, nullptr);
wait_until_cache_ready(*cache1);
size_t auto_total_max_pending_bytes = 0;
ASSERT_TRUE(resolve_async_file_cache_write_max_pending_bytes(-1, MemInfo::mem_limit(),
&auto_total_max_pending_bytes)
.ok());
EXPECT_EQ(cache1->async_write_manager()->options().max_pending_bytes,
auto_total_max_pending_bytes);
ASSERT_TRUE(factory->create_file_cache(path2.string(), async_write_cache_settings()).ok());
auto* cache2 = factory->get_by_path(path2.string());
ASSERT_NE(cache2, nullptr);
wait_until_cache_ready(*cache2);
const size_t auto_max_pending_bytes_per_instance = auto_total_max_pending_bytes / 2;
for (auto* cache : {cache1, cache2}) {
ASSERT_FALSE(cache->async_write_manager()->_started.load(std::memory_order_acquire));
EXPECT_EQ(cache->async_write_manager()->options().max_pending_bytes,
auto_max_pending_bytes_per_instance);
}
AsyncCacheWriteBufferPtr disabled_buffer;
ASSERT_TRUE(
cache1->async_write_manager()->allocate_tracked_buffer(4096, &disabled_buffer).ok());
const auto disabled_hash = BlockFileCache::hash("disabled_async_write_manager");
AsyncCacheWriteTask disabled_task {
.cache_hash = disabled_hash,
.file_offset = 0,
.write_size = disabled_buffer->size(),
.buffer = disabled_buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = cache1->async_write_manager()->current_write_epoch(disabled_hash),
.on_finalized = nullptr,
};
EXPECT_FALSE(cache1->async_write_manager()->try_submit(std::move(disabled_task)));
EXPECT_EQ(cache1->async_write_manager()->pending_count(), 0);
const int32_t new_workers = old_workers == 1 ? 2 : 1;
constexpr int64_t new_total_max_pending_bytes = 64 * 1024 * 1024 + 8192;
ASSERT_TRUE(config::set_config("async_file_cache_write_workers_per_disk",
std::to_string(new_workers))
.ok());
ASSERT_TRUE(config::set_config("async_file_cache_write_max_pending_bytes",
std::to_string(new_total_max_pending_bytes))
.ok());
ASSERT_TRUE(config::set_config("enable_async_file_cache_write", "true").ok());
for (auto* cache : {cache1, cache2}) {
const auto updated = cache->async_write_manager()->options();
EXPECT_TRUE(cache->async_write_manager()->_started.load(std::memory_order_acquire));
EXPECT_EQ(updated.worker_count, new_workers);
EXPECT_EQ(updated.max_pending_bytes, new_total_max_pending_bytes / 2);
}
ASSERT_TRUE(config::set_config("async_file_cache_write_max_pending_bytes", "-1").ok());
for (auto* cache : {cache1, cache2}) {
EXPECT_EQ(cache->async_write_manager()->options().max_pending_bytes,
auto_max_pending_bytes_per_instance);
}
}
TEST_F(AsyncCacheWriteManagerTest, ShutdownWaitsForConcurrentReplacementAndDrains) {
auto cache = create_cache("async_write_manager_shutdown_replacement");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
auto options = manager->options();
options.worker_count = 1;
options.max_pending_bytes = 2 * 4096;
ASSERT_TRUE(manager->update_options(options).ok());
const auto active_hash = BlockFileCache::hash("shutdown_replacement_active");
std::mutex mutex;
std::condition_variable cv;
bool active_entered = false;
bool release_active = false;
bool victim_callback_entered = false;
bool release_victim_callback = false;
std::vector<size_t> finalized(3, 0);
std::promise<void> shutdown_stopped_accepting;
auto shutdown_stopped_accepting_future = shutdown_stopped_accepting.get_future();
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard worker_guard;
SyncPoint::CallbackGuard shutdown_guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::_persist_task:before_get_or_set",
[&](auto&& args) {
const auto* task = try_any_cast<const AsyncCacheWriteTask*>(args[0]);
if (task->cache_hash != active_hash) {
return;
}
std::unique_lock lock(mutex);
active_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_active; });
},
&worker_guard);
sync_point->set_call_back(
"AsyncCacheWriteManager::shutdown:after_stop_accepting",
[&](auto&&) { shutdown_stopped_accepting.set_value(); }, &shutdown_guard);
sync_point->enable_processing();
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_active = true;
release_victim_callback = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
const auto finalizer = [&](size_t task_id) {
return [&, task_id](const AsyncCacheWriteTask&) {
std::lock_guard lock(mutex);
++finalized[task_id];
cv.notify_all();
};
};
ASSERT_TRUE(manager->try_submit(
make_async_write_task(manager, "shutdown_replacement_active", 'a', finalizer(0))));
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5), [&]() { return active_entered; }));
}
ASSERT_TRUE(manager->try_submit(make_async_write_task(
manager, "shutdown_replacement_victim", 'v', [&](const AsyncCacheWriteTask&) {
std::unique_lock lock(mutex);
++finalized[1];
victim_callback_entered = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_victim_callback; });
})));
auto replacement_future = std::async(std::launch::async, [&]() {
SCOPED_ATTACH_TASK(ExecEnv::GetInstance()->orphan_mem_tracker());
return manager->try_submit(
make_async_write_task(manager, "shutdown_replacement_new", 'n', finalizer(2)));
});
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(5),
[&]() { return victim_callback_entered; }));
}
auto shutdown_future = std::async(std::launch::async, [manager]() { manager->shutdown(); });
ASSERT_EQ(shutdown_stopped_accepting_future.wait_for(std::chrono::seconds(5)),
std::future_status::ready);
EXPECT_EQ(shutdown_future.wait_for(std::chrono::milliseconds(0)), std::future_status::timeout);
{
std::lock_guard lock(mutex);
release_victim_callback = true;
}
cv.notify_all();
ASSERT_TRUE(replacement_future.get());
EXPECT_EQ(shutdown_future.wait_for(std::chrono::milliseconds(0)), std::future_status::timeout);
{
std::lock_guard lock(mutex);
release_active = true;
}
cv.notify_all();
ASSERT_EQ(shutdown_future.wait_for(std::chrono::seconds(5)), std::future_status::ready);
shutdown_future.get();
EXPECT_EQ(finalized, (std::vector<size_t> {1, 1, 1}));
EXPECT_EQ(manager->pending_count(), 0);
EXPECT_EQ(manager->queued_count(), 0);
EXPECT_EQ(manager->active_task_count(), 0);
}
TEST_F(AsyncCacheWriteManagerTest, ShutdownRejectsSubmitterWhenShutdownWinsAdmissionRace) {
auto cache = create_cache("async_write_manager_shutdown");
auto* manager = cache->async_write_manager();
ASSERT_NE(manager, nullptr);
std::mutex mutex;
std::condition_variable cv;
bool submitter_waiting_to_register = false;
bool release_submitter = false;
std::promise<void> shutdown_stopped_accepting;
auto shutdown_stopped_accepting_future = shutdown_stopped_accepting.get_future();
auto* sync_point = SyncPoint::get_instance();
SyncPoint::CallbackGuard submit_guard;
SyncPoint::CallbackGuard shutdown_guard;
sync_point->set_call_back(
"AsyncCacheWriteManager::try_submit:before_register",
[&](auto&&) {
std::unique_lock lock(mutex);
submitter_waiting_to_register = true;
cv.notify_all();
cv.wait(lock, [&]() { return release_submitter; });
},
&submit_guard);
sync_point->set_call_back(
"AsyncCacheWriteManager::shutdown:after_stop_accepting",
[&](auto&&) { shutdown_stopped_accepting.set_value(); }, &shutdown_guard);
sync_point->enable_processing();
std::future<bool> submit_future;
std::future<void> shutdown_future;
Defer clear_sync_point {[&]() {
{
std::lock_guard lock(mutex);
release_submitter = true;
}
cv.notify_all();
sync_point->disable_processing();
sync_point->clear_all_call_backs();
}};
AsyncCacheWriteBufferPtr buffer;
ASSERT_TRUE(manager->allocate_tracked_buffer(4096, &buffer).ok());
const auto hash = BlockFileCache::hash("shutdown_submitter");
AsyncCacheWriteTask task {
.cache_hash = hash,
.file_offset = 0,
.write_size = buffer->size(),
.buffer = buffer,
.admission_ctx = {},
.submit_ts_us = MonotonicMicros(),
.write_epoch = manager->current_write_epoch(hash),
.on_finalized = nullptr,
};
submit_future = std::async(std::launch::async, [manager, task = std::move(task)]() mutable {
return manager->try_submit(std::move(task));
});
{
std::unique_lock lock(mutex);
ASSERT_TRUE(cv.wait_for(lock, std::chrono::seconds(1),
[&]() { return submitter_waiting_to_register; }));
}
shutdown_future = std::async(std::launch::async, [manager]() { manager->shutdown(); });
ASSERT_EQ(shutdown_stopped_accepting_future.wait_for(std::chrono::seconds(1)),
std::future_status::ready);
ASSERT_EQ(shutdown_future.wait_for(std::chrono::seconds(1)), std::future_status::ready);
shutdown_future.get();
{
std::lock_guard lock(mutex);
release_submitter = true;
}
cv.notify_all();
EXPECT_FALSE(submit_future.get());
EXPECT_EQ(manager->pending_count(), 0);
}
} // namespace
} // namespace doris::io