Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
51 changes: 49 additions & 2 deletions mooncake-store/include/master_service.h
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,53 @@ class MasterService {
std::vector<tl::expected<void, ErrorCode>> BatchPutRevoke(
const UUID& client_id, const std::vector<std::string>& keys);

/**
* @brief Start an upsert operation. If the key does not exist, behaves
* like PutStart. If the key exists with the same size, performs in-place
* update (reuses existing buffers). If the key exists with a different
* size, deletes old replicas and allocates new ones.
* @return Replica descriptors on success, or error code on failure.
* Possible errors: OBJECT_HAS_REPLICATION_TASK (Copy/Move/Offload in
* progress), OBJECT_REPLICA_BUSY (replicas have non-zero refcnt).
*/
auto UpsertStart(const UUID& client_id, const std::string& key,
const uint64_t slice_length, const ReplicateConfig& config)
-> tl::expected<std::vector<Replica::Descriptor>, ErrorCode>;

/**
* @brief Complete an upsert operation. Delegates to PutEnd.
*/
auto UpsertEnd(const UUID& client_id, const std::string& key,
ReplicaType replica_type) -> tl::expected<void, ErrorCode>;

/**
* @brief Revoke an upsert operation. Delegates to PutRevoke.
*/
auto UpsertRevoke(const UUID& client_id, const std::string& key,
ReplicaType replica_type)
-> tl::expected<void, ErrorCode>;

/**
* @brief Start a batch of upsert operations.
*/
std::vector<tl::expected<std::vector<Replica::Descriptor>, ErrorCode>>
BatchUpsertStart(const UUID& client_id,
const std::vector<std::string>& keys,
const std::vector<uint64_t>& slice_lengths,
const ReplicateConfig& config);

/**
* @brief Complete a batch of upsert operations. Delegates to BatchPutEnd.
*/
std::vector<tl::expected<void, ErrorCode>> BatchUpsertEnd(
const UUID& client_id, const std::vector<std::string>& keys);

/**
* @brief Revoke a batch of upsert operations. Delegates to BatchPutRevoke.
*/
std::vector<tl::expected<void, ErrorCode>> BatchUpsertRevoke(
const UUID& client_id, const std::vector<std::string>& keys);

/**
* @brief Evict a disk replica for a key (triggered by client-side disk
* eviction).
Expand Down Expand Up @@ -493,8 +540,8 @@ class MasterService {
ObjectMetadata(ObjectMetadata&&) = delete;
ObjectMetadata& operator=(ObjectMetadata&&) = delete;

const UUID client_id;
const std::chrono::system_clock::time_point put_start_time;
UUID client_id;
std::chrono::system_clock::time_point put_start_time;
const size_t size;

mutable SpinLock lock;
Expand Down
14 changes: 14 additions & 0 deletions mooncake-store/include/replica.h
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,12 @@ class Replica {
return replica.is_processing();
}

[[nodiscard]] bool is_busy() const { return refcnt_.load() > 0; }

[[nodiscard]] static bool fn_is_busy(const Replica& replica) {
return replica.is_busy();
}

[[nodiscard]] ReplicaType type() const {
return std::visit(ReplicaTypeVisitor{}, data_);
}
Expand Down Expand Up @@ -297,6 +303,14 @@ class Replica {
}
}

void mark_processing() {
if (status_ == ReplicaStatus::COMPLETE) {
status_ = ReplicaStatus::PROCESSING;
} else {
LOG(ERROR) << "Cannot mark_processing from status: " << status_;
}
}

void inc_refcnt() { refcnt_.fetch_add(1); }

void dec_refcnt() { refcnt_.fetch_sub(1); }
Expand Down
2 changes: 2 additions & 0 deletions mooncake-store/include/types.h
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,8 @@ enum class ErrorCode : int32_t {
REPLICA_IS_GONE = -712, ///< Replica existed once, but is gone now.
REPLICA_NOT_IN_LOCAL_MEMORY =
-713, ///< Replica does not reside in current node memory.
OBJECT_REPLICA_BUSY =
-714, ///< Object replicas have non-zero refcnt.

// Transfer errors (Range: -800 to -899)
TRANSFER_FAIL = -800, ///< Transfer operation failed.
Expand Down
269 changes: 269 additions & 0 deletions mooncake-store/src/master_service.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -942,6 +942,275 @@ std::vector<tl::expected<void, ErrorCode>> MasterService::BatchPutRevoke(
return results;
}

auto MasterService::UpsertStart(const UUID& client_id, const std::string& key,
const uint64_t slice_length,
const ReplicateConfig& config)
-> tl::expected<std::vector<Replica::Descriptor>, ErrorCode> {
if (config.replica_num == 0 || key.empty() || slice_length == 0) {
LOG(ERROR) << "key=" << key << ", replica_num=" << config.replica_num
<< ", slice_length=" << slice_length
<< ", key_size=" << key.size() << ", error=invalid_params";
return tl::make_unexpected(ErrorCode::INVALID_PARAMS);
}

// Validate slice lengths
uint64_t total_length = 0;
if ((memory_allocator_type_ == BufferAllocatorType::CACHELIB) &&
(slice_length > kMaxSliceSize)) {
LOG(ERROR) << "key=" << key << ", slice_length=" << slice_length
<< ", max_size=" << kMaxSliceSize
<< ", error=invalid_slice_size";
return tl::make_unexpected(ErrorCode::INVALID_PARAMS);
}
total_length += slice_length;

VLOG(1) << "key=" << key << ", value_length=" << total_length
<< ", slice_length=" << slice_length << ", config=" << config
<< ", action=upsert_start_begin";

std::shared_lock<std::shared_mutex> shared_lock(snapshot_mutex_);
MetadataShardAccessorRW shard(this, getShardIndex(key));

const auto now = std::chrono::system_clock::now();
auto it = shard->metadata.find(key);

// If key found but all handles are stale, erase and treat as non-existent
if (it != shard->metadata.end() && CleanupStaleHandles(it->second)) {
shard->processing_keys.erase(key);
shard->metadata.erase(it);
it = shard->metadata.end();
}

// If key exists and is valid, handle Case B/C or preemption
if (it != shard->metadata.end()) {
auto& metadata = it->second;

// Safety check: reject if Copy/Move in progress
if (shard->replication_tasks.count(key) > 0) {
LOG(INFO) << "key=" << key
<< ", error=object_has_replication_task";
return tl::make_unexpected(ErrorCode::OBJECT_HAS_REPLICATION_TASK);
}

// Safety check: reject if offloading in progress
if (shard->offloading_tasks.count(key) > 0) {
LOG(INFO) << "key=" << key
<< ", error=object_has_offloading_task";
return tl::make_unexpected(ErrorCode::OBJECT_HAS_REPLICATION_TASK);
}

// Preemption: if another Put/Upsert is in progress, preempt it
if (shard->processing_keys.count(key) > 0) {
auto processing_replicas =
metadata.PopReplicas(&Replica::fn_is_processing);
if (!processing_replicas.empty()) {
std::lock_guard lock(discarded_replicas_mutex_);
discarded_replicas_.emplace_back(
std::move(processing_replicas),
now + put_start_release_timeout_sec_);
}
shard->processing_keys.erase(key);

// If no COMPLETE replicas remain after preemption, erase metadata
// and fall through to Case A
if (!metadata.HasReplica(&Replica::fn_is_completed)) {
shard->metadata.erase(it);
it = shard->metadata.end();
}
}
}

// If metadata was erased (stale, preempted with no COMPLETE replicas,
// or never existed), fall through to Case A
if (it == shard->metadata.end()) {
// Case A: key does not exist — same as PutStart
std::vector<Replica> replicas;
{
ScopedAllocatorAccess allocator_access =
segment_manager_.getAllocatorAccess();
const auto& allocator_manager =
allocator_access.getAllocatorManager();

std::vector<std::string> preferred_segments;
if (!config.preferred_segment.empty()) {
preferred_segments.push_back(config.preferred_segment);
} else if (!config.preferred_segments.empty()) {
preferred_segments = config.preferred_segments;
}

auto allocation_result = allocation_strategy_->Allocate(
allocator_manager, slice_length, config.replica_num,
preferred_segments);

if (!allocation_result.has_value()) {
VLOG(1) << "Failed to allocate replicas for key=" << key
<< ", error: " << allocation_result.error();
if (allocation_result.error() == ErrorCode::INVALID_PARAMS) {
return tl::make_unexpected(ErrorCode::INVALID_PARAMS);
}
need_eviction_ = true;
return tl::make_unexpected(ErrorCode::NO_AVAILABLE_HANDLE);
}

replicas = std::move(allocation_result.value());
}

if (use_disk_replica_) {
std::string file_path =
ResolvePathFromKey(key, root_fs_dir_, cluster_id_);
replicas.emplace_back(file_path, total_length,
ReplicaStatus::PROCESSING);
}

std::vector<Replica::Descriptor> replica_list;
replica_list.reserve(replicas.size());
for (const auto& replica : replicas) {
replica_list.emplace_back(replica.get_descriptor());
}

shard->metadata.emplace(
std::piecewise_construct, std::forward_as_tuple(key),
std::forward_as_tuple(client_id, now, total_length,
std::move(replicas), config.with_soft_pin));
shard->processing_keys.insert(key);

VLOG(1) << "key=" << key << ", action=upsert_start_case_a";
return replica_list;
Comment thread
00fish0 marked this conversation as resolved.
Outdated
}

// Key exists with COMPLETE replicas — Case B or Case C
auto& metadata = it->second;

// Safety check: reject if any replica has non-zero refcnt
if (metadata.HasReplica(&Replica::fn_is_busy)) {
LOG(INFO) << "key=" << key << ", error=object_replica_busy";
return tl::make_unexpected(ErrorCode::OBJECT_REPLICA_BUSY);
}

if (metadata.size == total_length) {
// Case B: same size — in-place update
metadata.client_id = client_id;
metadata.put_start_time = now;

metadata.VisitReplicas(
&Replica::fn_is_completed,
[](Replica& replica) { replica.mark_processing(); });

shard->processing_keys.insert(key);

std::vector<Replica::Descriptor> replica_list;
const auto& all_replicas = metadata.GetAllReplicas();
replica_list.reserve(all_replicas.size());
for (const auto& replica : all_replicas) {
replica_list.emplace_back(replica.get_descriptor());
}

VLOG(1) << "key=" << key << ", action=upsert_start_case_b_inplace";
return replica_list;
}

// Case C: different size — delete and reallocate
auto old_replicas = metadata.PopReplicas();
if (!old_replicas.empty()) {
std::lock_guard lock(discarded_replicas_mutex_);
discarded_replicas_.emplace_back(
std::move(old_replicas),
now + put_start_release_timeout_sec_);
}
shard->metadata.erase(it);

// Allocate new replicas (same as Case A)
std::vector<Replica> replicas;
{
ScopedAllocatorAccess allocator_access =
segment_manager_.getAllocatorAccess();
const auto& allocator_manager =
allocator_access.getAllocatorManager();

std::vector<std::string> preferred_segments;
if (!config.preferred_segment.empty()) {
preferred_segments.push_back(config.preferred_segment);
} else if (!config.preferred_segments.empty()) {
preferred_segments = config.preferred_segments;
}

auto allocation_result = allocation_strategy_->Allocate(
allocator_manager, slice_length, config.replica_num,
preferred_segments);

if (!allocation_result.has_value()) {
VLOG(1) << "Failed to allocate replicas for key=" << key
<< ", error: " << allocation_result.error();
if (allocation_result.error() == ErrorCode::INVALID_PARAMS) {
return tl::make_unexpected(ErrorCode::INVALID_PARAMS);
}
need_eviction_ = true;
return tl::make_unexpected(ErrorCode::NO_AVAILABLE_HANDLE);
}

replicas = std::move(allocation_result.value());
}

if (use_disk_replica_) {
std::string file_path =
ResolvePathFromKey(key, root_fs_dir_, cluster_id_);
replicas.emplace_back(file_path, total_length,
ReplicaStatus::PROCESSING);
}

std::vector<Replica::Descriptor> replica_list;
replica_list.reserve(replicas.size());
for (const auto& replica : replicas) {
replica_list.emplace_back(replica.get_descriptor());
}

shard->metadata.emplace(
std::piecewise_construct, std::forward_as_tuple(key),
std::forward_as_tuple(client_id, now, total_length,
std::move(replicas), config.with_soft_pin));
Comment thread
00fish0 marked this conversation as resolved.
Outdated
shard->processing_keys.insert(key);

VLOG(1) << "key=" << key << ", action=upsert_start_case_c_reallocate";
return replica_list;
}

auto MasterService::UpsertEnd(const UUID& client_id, const std::string& key,
ReplicaType replica_type)
-> tl::expected<void, ErrorCode> {
return PutEnd(client_id, key, replica_type);
}

auto MasterService::UpsertRevoke(const UUID& client_id, const std::string& key,
ReplicaType replica_type)
-> tl::expected<void, ErrorCode> {
return PutRevoke(client_id, key, replica_type);
}

std::vector<tl::expected<std::vector<Replica::Descriptor>, ErrorCode>>
MasterService::BatchUpsertStart(
const UUID& client_id, const std::vector<std::string>& keys,
const std::vector<uint64_t>& slice_lengths,
const ReplicateConfig& config) {
std::vector<tl::expected<std::vector<Replica::Descriptor>, ErrorCode>>
results;
results.reserve(keys.size());
for (size_t i = 0; i < keys.size(); ++i) {
results.emplace_back(
UpsertStart(client_id, keys[i], slice_lengths[i], config));
}
Comment on lines +1211 to +1229
return results;
}

std::vector<tl::expected<void, ErrorCode>> MasterService::BatchUpsertEnd(
const UUID& client_id, const std::vector<std::string>& keys) {
return BatchPutEnd(client_id, keys);
}

std::vector<tl::expected<void, ErrorCode>> MasterService::BatchUpsertRevoke(
const UUID& client_id, const std::vector<std::string>& keys) {
return BatchPutRevoke(client_id, keys);
}

auto MasterService::EvictDiskReplica(const UUID& client_id,
const std::string& key,
ReplicaType replica_type)
Expand Down
1 change: 1 addition & 0 deletions mooncake-store/src/types.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ const std::string& toString(ErrorCode errorCode) noexcept {
{ErrorCode::REPLICA_NOT_FOUND, "REPLICA_NOT_FOUND"},
{ErrorCode::REPLICA_ALREADY_EXISTS, "REPLICA_ALREADY_EXISTS"},
{ErrorCode::REPLICA_IS_GONE, "REPLICA_IS_GONE"},
{ErrorCode::OBJECT_REPLICA_BUSY, "OBJECT_REPLICA_BUSY"},
{ErrorCode::TRANSFER_FAIL, "TRANSFER_FAIL"},
{ErrorCode::RPC_FAIL, "RPC_FAIL"},
{ErrorCode::ETCD_OPERATION_ERROR, "ETCD_OPERATION_ERROR"},
Expand Down
Loading
Loading