Skip to content

Commit 32143a4

Browse files
committed
[config] add bounded async CopyGA HTTP timeouts
1 parent c3ec31b commit 32143a4

10 files changed

Lines changed: 718 additions & 1252 deletions

File tree

docs/design/async_copyga.md

Lines changed: 512 additions & 1252 deletions
Large diffs are not rendered by default.

kv_cache_manager/config/cache_config.h

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,9 @@ class CacheConfig : public Jsonizable {
8282
int64_t migration_copy_poll_max_interval_ms() const {
8383
return migration_config_.copy_poll_max_interval_ms();
8484
}
85+
int64_t migration_copy_connect_timeout_ms() const { return migration_config_.copy_connect_timeout_ms(); }
86+
int64_t migration_copy_submit_timeout_ms() const { return migration_config_.copy_submit_timeout_ms(); }
87+
int64_t migration_copy_query_timeout_ms() const { return migration_config_.copy_query_timeout_ms(); }
8588
// Setters
8689
void set_reclaim_strategy(const std::shared_ptr<CacheReclaimStrategy> &reclaim_strategy) {
8790
reclaim_strategy_ = reclaim_strategy;
@@ -122,6 +125,15 @@ class CacheConfig : public Jsonizable {
122125
void set_migration_copy_poll_max_interval_ms(int64_t value) {
123126
migration_config_.set_copy_poll_max_interval_ms(value);
124127
}
128+
void set_migration_copy_connect_timeout_ms(int64_t value) {
129+
migration_config_.set_copy_connect_timeout_ms(value);
130+
}
131+
void set_migration_copy_submit_timeout_ms(int64_t value) {
132+
migration_config_.set_copy_submit_timeout_ms(value);
133+
}
134+
void set_migration_copy_query_timeout_ms(int64_t value) {
135+
migration_config_.set_copy_query_timeout_ms(value);
136+
}
125137
void set_migration_config(const MigrationConfig &migration_config) {
126138
migration_config_ = migration_config;
127139
}

kv_cache_manager/config/migration_strategy.cc

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,18 @@ bool MigrationConfig::FromRapidValue(const rapidjson::Value &rapid_value) {
144144
"copy_poll_max_interval_ms",
145145
copy_poll_max_interval_ms_,
146146
MigrationConfig::kDefaultCopyPollMaxIntervalMs);
147+
KVCM_JSON_GET_DEFAULT_MACRO(rapid_value,
148+
"copy_connect_timeout_ms",
149+
copy_connect_timeout_ms_,
150+
MigrationConfig::kDefaultCopyConnectTimeoutMs);
151+
KVCM_JSON_GET_DEFAULT_MACRO(rapid_value,
152+
"copy_submit_timeout_ms",
153+
copy_submit_timeout_ms_,
154+
MigrationConfig::kDefaultCopySubmitTimeoutMs);
155+
KVCM_JSON_GET_DEFAULT_MACRO(rapid_value,
156+
"copy_query_timeout_ms",
157+
copy_query_timeout_ms_,
158+
MigrationConfig::kDefaultCopyQueryTimeoutMs);
147159
KVCM_JSON_GET_MACRO(rapid_value, "strategies", strategies_);
148160
return true;
149161
}
@@ -158,6 +170,9 @@ void MigrationConfig::ToRapidWriter(rapidjson::Writer<rapidjson::StringBuffer> &
158170
Put(writer, "copy_operation_deadline_ms", copy_operation_deadline_ms_);
159171
Put(writer, "copy_poll_initial_interval_ms", copy_poll_initial_interval_ms_);
160172
Put(writer, "copy_poll_max_interval_ms", copy_poll_max_interval_ms_);
173+
Put(writer, "copy_connect_timeout_ms", copy_connect_timeout_ms_);
174+
Put(writer, "copy_submit_timeout_ms", copy_submit_timeout_ms_);
175+
Put(writer, "copy_query_timeout_ms", copy_query_timeout_ms_);
161176
Put(writer, "strategies", strategies_);
162177
}
163178

@@ -185,6 +200,17 @@ bool MigrationConfig::ValidateRequiredFields(std::string &invalid_fields) const
185200
valid = false;
186201
local_invalid_fields += "{copy_async_timing}";
187202
}
203+
if (copy_connect_timeout_ms_ <= 0 || copy_submit_timeout_ms_ < copy_connect_timeout_ms_ ||
204+
copy_query_timeout_ms_ < copy_connect_timeout_ms_ ||
205+
copy_submit_timeout_ms_ >= copy_operation_deadline_ms_ ||
206+
copy_query_timeout_ms_ >= copy_operation_deadline_ms_ ||
207+
copy_submit_timeout_ms_ >
208+
copy_operation_deadline_ms_ / MigrationConfig::kMinCopyHttpTimeoutWindowsPerDeadline ||
209+
copy_query_timeout_ms_ >
210+
copy_operation_deadline_ms_ / MigrationConfig::kMinCopyHttpTimeoutWindowsPerDeadline) {
211+
valid = false;
212+
local_invalid_fields += "{copy_async_http_timing}";
213+
}
188214
if (copy_execution_mode_ == MigrationCopyExecutionMode::ASYNC_REQUIRED &&
189215
(copy_max_inflight_bytes_ == 0 || copy_max_quarantine_operations_ <= 0 ||
190216
copy_max_quarantine_bytes_ == 0)) {

kv_cache_manager/config/migration_strategy.h

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -131,6 +131,13 @@ class MigrationConfig : public Jsonizable {
131131
static constexpr int64_t kDefaultCopyOperationDeadlineMs = 10 * 60 * 1000;
132132
static constexpr int64_t kDefaultCopyPollInitialIntervalMs = 20;
133133
static constexpr int64_t kDefaultCopyPollMaxIntervalMs = 1000;
134+
static constexpr int64_t kDefaultCopyConnectTimeoutMs = 1000;
135+
static constexpr int64_t kDefaultCopySubmitTimeoutMs = 3000;
136+
static constexpr int64_t kDefaultCopyQueryTimeoutMs = 3000;
137+
// Keep one control-plane request from consuming most of the operation
138+
// deadline. This guarantees at least sixteen timeout windows for a C1
139+
// coordinator; shared-backend concurrency still needs a deployment gate.
140+
static constexpr int64_t kMinCopyHttpTimeoutWindowsPerDeadline = 16;
134141

135142
MigrationConfig() = default;
136143
~MigrationConfig() override;
@@ -146,6 +153,9 @@ class MigrationConfig : public Jsonizable {
146153
int64_t copy_operation_deadline_ms() const { return copy_operation_deadline_ms_; }
147154
int64_t copy_poll_initial_interval_ms() const { return copy_poll_initial_interval_ms_; }
148155
int64_t copy_poll_max_interval_ms() const { return copy_poll_max_interval_ms_; }
156+
int64_t copy_connect_timeout_ms() const { return copy_connect_timeout_ms_; }
157+
int64_t copy_submit_timeout_ms() const { return copy_submit_timeout_ms_; }
158+
int64_t copy_query_timeout_ms() const { return copy_query_timeout_ms_; }
149159

150160
void set_strategies(const std::vector<std::shared_ptr<MigrationStrategy>> &strategies) {
151161
strategies_ = strategies;
@@ -163,6 +173,9 @@ class MigrationConfig : public Jsonizable {
163173
void set_copy_operation_deadline_ms(int64_t value) { copy_operation_deadline_ms_ = value; }
164174
void set_copy_poll_initial_interval_ms(int64_t value) { copy_poll_initial_interval_ms_ = value; }
165175
void set_copy_poll_max_interval_ms(int64_t value) { copy_poll_max_interval_ms_ = value; }
176+
void set_copy_connect_timeout_ms(int64_t value) { copy_connect_timeout_ms_ = value; }
177+
void set_copy_submit_timeout_ms(int64_t value) { copy_submit_timeout_ms_ = value; }
178+
void set_copy_query_timeout_ms(int64_t value) { copy_query_timeout_ms_ = value; }
166179

167180
bool FromRapidValue(const rapidjson::Value &rapid_value) override;
168181
void ToRapidWriter(rapidjson::Writer<rapidjson::StringBuffer> &writer) const noexcept override;
@@ -179,6 +192,9 @@ class MigrationConfig : public Jsonizable {
179192
int64_t copy_operation_deadline_ms_ = kDefaultCopyOperationDeadlineMs;
180193
int64_t copy_poll_initial_interval_ms_ = kDefaultCopyPollInitialIntervalMs;
181194
int64_t copy_poll_max_interval_ms_ = kDefaultCopyPollMaxIntervalMs;
195+
int64_t copy_connect_timeout_ms_ = kDefaultCopyConnectTimeoutMs;
196+
int64_t copy_submit_timeout_ms_ = kDefaultCopySubmitTimeoutMs;
197+
int64_t copy_query_timeout_ms_ = kDefaultCopyQueryTimeoutMs;
182198
};
183199

184200
} // namespace kv_cache_manager

kv_cache_manager/config/test/instance_group_test.cc

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -193,6 +193,54 @@ TEST_F(InstanceGroupTest, ProtoRoundTripWithoutBuckets) {
193193
EXPECT_TRUE(restored.revisit_interval_buckets().empty());
194194
}
195195

196+
TEST_F(InstanceGroupTest, ProtoRoundTripPreservesAsyncCopyHttpTimeouts) {
197+
InstanceGroup original;
198+
original.set_name("test_group");
199+
original.set_storage_candidates({"local"});
200+
original.set_global_quota_group_name("default");
201+
original.set_max_instance_count(10);
202+
original.set_version(1);
203+
204+
auto cache_config = std::make_shared<CacheConfig>();
205+
cache_config->set_reclaim_strategy(std::make_shared<CacheReclaimStrategy>());
206+
cache_config->set_migration_copy_connect_timeout_ms(750);
207+
cache_config->set_migration_copy_submit_timeout_ms(2500);
208+
cache_config->set_migration_copy_query_timeout_ms(2800);
209+
original.set_cache_config(cache_config);
210+
211+
proto::admin::InstanceGroup proto_msg;
212+
ProtoConvert::InstanceGroupToProto(original, &proto_msg);
213+
ASSERT_TRUE(proto_msg.has_cache_config());
214+
ASSERT_TRUE(proto_msg.cache_config().has_migration_config());
215+
const auto &migration_config = proto_msg.cache_config().migration_config();
216+
ASSERT_TRUE(migration_config.has_copy_connect_timeout_ms());
217+
ASSERT_TRUE(migration_config.has_copy_submit_timeout_ms());
218+
ASSERT_TRUE(migration_config.has_copy_query_timeout_ms());
219+
EXPECT_EQ(750, migration_config.copy_connect_timeout_ms().value());
220+
EXPECT_EQ(2500, migration_config.copy_submit_timeout_ms().value());
221+
EXPECT_EQ(2800, migration_config.copy_query_timeout_ms().value());
222+
223+
InstanceGroup restored;
224+
ProtoConvert::InstanceGroupFromProto(&proto_msg, restored);
225+
ASSERT_NE(nullptr, restored.cache_config());
226+
EXPECT_EQ(750, restored.cache_config()->migration_copy_connect_timeout_ms());
227+
EXPECT_EQ(2500, restored.cache_config()->migration_copy_submit_timeout_ms());
228+
EXPECT_EQ(2800, restored.cache_config()->migration_copy_query_timeout_ms());
229+
230+
proto_msg.mutable_cache_config()->mutable_migration_config()->clear_copy_connect_timeout_ms();
231+
proto_msg.mutable_cache_config()->mutable_migration_config()->clear_copy_submit_timeout_ms();
232+
proto_msg.mutable_cache_config()->mutable_migration_config()->clear_copy_query_timeout_ms();
233+
InstanceGroup restored_legacy;
234+
ProtoConvert::InstanceGroupFromProto(&proto_msg, restored_legacy);
235+
ASSERT_NE(nullptr, restored_legacy.cache_config());
236+
EXPECT_EQ(MigrationConfig::kDefaultCopyConnectTimeoutMs,
237+
restored_legacy.cache_config()->migration_copy_connect_timeout_ms());
238+
EXPECT_EQ(MigrationConfig::kDefaultCopySubmitTimeoutMs,
239+
restored_legacy.cache_config()->migration_copy_submit_timeout_ms());
240+
EXPECT_EQ(MigrationConfig::kDefaultCopyQueryTimeoutMs,
241+
restored_legacy.cache_config()->migration_copy_query_timeout_ms());
242+
}
243+
196244
TEST_F(InstanceGroupTest, EventReportStorageSpecProtoRoundTripPreservesSnapshotSettings) {
197245
proto::admin::StorageConfig legacy_proto_config;
198246
legacy_proto_config.set_global_unique_name("legacy_event_report");

kv_cache_manager/config/test/migration_strategy_test.cc

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -183,6 +183,9 @@ TEST_F(MigrationStrategyTest, TestInCacheConfig) {
183183
"migration_config": {
184184
"copy_max_concurrency": 6,
185185
"mark_clear_policy": 1,
186+
"copy_connect_timeout_ms": 750,
187+
"copy_submit_timeout_ms": 2500,
188+
"copy_query_timeout_ms": 2800,
186189
"strategies": [
187190
{
188191
"source_storage_name": "pace_mempool_01",
@@ -204,6 +207,9 @@ TEST_F(MigrationStrategyTest, TestInCacheConfig) {
204207
ASSERT_TRUE(cache_config.FromJsonString(json));
205208
ASSERT_EQ(6, cache_config.migration_copy_max_concurrency());
206209
ASSERT_EQ(MigrationMarkClearPolicy::CLEAR_ON_FULL_BLOCK_COVERED, cache_config.migration_mark_clear_policy());
210+
ASSERT_EQ(750, cache_config.migration_copy_connect_timeout_ms());
211+
ASSERT_EQ(2500, cache_config.migration_copy_submit_timeout_ms());
212+
ASSERT_EQ(2800, cache_config.migration_copy_query_timeout_ms());
207213
ASSERT_EQ(2u, cache_config.migration_strategies().size());
208214
ASSERT_EQ("pace_mempool_01", cache_config.migration_strategies()[0]->source_storage_name());
209215
ASSERT_EQ("pace_ssd_02", cache_config.migration_strategies()[1]->target_storage_name());
@@ -214,6 +220,9 @@ TEST_F(MigrationStrategyTest, TestInCacheConfig) {
214220
ASSERT_TRUE(parsed.FromJsonString(cache_config.ToJsonString()));
215221
ASSERT_EQ(6, parsed.migration_copy_max_concurrency());
216222
ASSERT_EQ(MigrationMarkClearPolicy::CLEAR_ON_FULL_BLOCK_COVERED, parsed.migration_mark_clear_policy());
223+
ASSERT_EQ(750, parsed.migration_copy_connect_timeout_ms());
224+
ASSERT_EQ(2500, parsed.migration_copy_submit_timeout_ms());
225+
ASSERT_EQ(2800, parsed.migration_copy_query_timeout_ms());
217226
ASSERT_EQ(2u, parsed.migration_strategies().size());
218227
ASSERT_DOUBLE_EQ(0.70, parsed.migration_strategies()[0]->trigger_threshold());
219228
ASSERT_EQ(120000, parsed.migration_strategies()[1]->methods().mark().timeout_ms());
@@ -228,9 +237,66 @@ TEST_F(MigrationStrategyTest, TestInCacheConfig) {
228237
ASSERT_TRUE(no_migration.FromJsonString(json2));
229238
ASSERT_EQ(CacheConfig::kDefaultMigrationCopyMaxConcurrency, no_migration.migration_copy_max_concurrency());
230239
ASSERT_EQ(MigrationMarkClearPolicy::CLEAR_ON_NEXT_WRITE_SUCCESS, no_migration.migration_mark_clear_policy());
240+
ASSERT_EQ(MigrationConfig::kDefaultCopyConnectTimeoutMs,
241+
no_migration.migration_copy_connect_timeout_ms());
242+
ASSERT_EQ(MigrationConfig::kDefaultCopySubmitTimeoutMs,
243+
no_migration.migration_copy_submit_timeout_ms());
244+
ASSERT_EQ(MigrationConfig::kDefaultCopyQueryTimeoutMs, no_migration.migration_copy_query_timeout_ms());
231245
ASSERT_TRUE(no_migration.migration_strategies().empty());
232246
}
233247

248+
TEST_F(MigrationStrategyTest, TestMigrationConfigRejectsInvalidAsyncHttpTiming) {
249+
const auto expect_invalid = [](const MigrationConfig &config) {
250+
std::string invalid_fields;
251+
EXPECT_FALSE(config.ValidateRequiredFields(invalid_fields));
252+
EXPECT_NE(std::string::npos, invalid_fields.find("copy_async_http_timing"));
253+
};
254+
255+
MigrationConfig config;
256+
std::string invalid_fields;
257+
ASSERT_TRUE(config.ValidateRequiredFields(invalid_fields)) << invalid_fields;
258+
259+
config.set_copy_connect_timeout_ms(0);
260+
expect_invalid(config);
261+
262+
config = MigrationConfig();
263+
config.set_copy_submit_timeout_ms(config.copy_connect_timeout_ms() - 1);
264+
expect_invalid(config);
265+
266+
config = MigrationConfig();
267+
config.set_copy_query_timeout_ms(config.copy_connect_timeout_ms() - 1);
268+
expect_invalid(config);
269+
270+
config = MigrationConfig();
271+
config.set_copy_operation_deadline_ms(config.copy_submit_timeout_ms());
272+
expect_invalid(config);
273+
274+
config = MigrationConfig();
275+
config.set_copy_operation_deadline_ms(config.copy_query_timeout_ms());
276+
expect_invalid(config);
277+
278+
config = MigrationConfig();
279+
config.set_copy_submit_timeout_ms(
280+
config.copy_operation_deadline_ms() /
281+
MigrationConfig::kMinCopyHttpTimeoutWindowsPerDeadline +
282+
1);
283+
expect_invalid(config);
284+
285+
config = MigrationConfig();
286+
config.set_copy_query_timeout_ms(
287+
config.copy_operation_deadline_ms() /
288+
MigrationConfig::kMinCopyHttpTimeoutWindowsPerDeadline +
289+
1);
290+
expect_invalid(config);
291+
292+
config = MigrationConfig();
293+
config.set_copy_operation_deadline_ms(
294+
config.copy_query_timeout_ms() *
295+
MigrationConfig::kMinCopyHttpTimeoutWindowsPerDeadline);
296+
invalid_fields.clear();
297+
EXPECT_TRUE(config.ValidateRequiredFields(invalid_fields)) << invalid_fields;
298+
}
299+
234300
TEST_F(MigrationStrategyTest, TestCacheConfigRejectsInvalidMigrationCopyConcurrency) {
235301
CacheConfig cache_config;
236302
cache_config.set_cache_prefer_strategy(CachePreferStrategy::CPS_ALWAYS_TAIR_MEMPOOL);

kv_cache_manager/data_storage/data_storage_backend.h

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,9 +53,17 @@ struct AsyncCopyBatchResult {
5353
};
5454

5555
struct AsyncCopyOptions {
56+
static constexpr int64_t kMinHttpTimeoutWindowsPerDeadline = 16;
57+
5658
int64_t operation_deadline_ms = 10 * 60 * 1000;
5759
int64_t initial_poll_interval_ms = 20;
5860
int64_t max_poll_interval_ms = 1000;
61+
// Per-request control-plane budgets. They are deliberately independent
62+
// from the end-to-end operation deadline: an HTTP timeout never proves
63+
// that a remote Copy task is terminal or that its destination is safe.
64+
int64_t connect_timeout_ms = 1000;
65+
int64_t submit_timeout_ms = 3000;
66+
int64_t query_timeout_ms = 3000;
5967
};
6068

6169
using AsyncCopyCompletion = std::function<void(AsyncCopyBatchResult)>;

kv_cache_manager/manager/migration_manager.cc

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -668,6 +668,9 @@ bool MigrationManager::RecoverAsyncCopyGuards() {
668668
options.operation_deadline_ms = group->cache_config()->migration_copy_operation_deadline_ms();
669669
options.initial_poll_interval_ms = group->cache_config()->migration_copy_poll_initial_interval_ms();
670670
options.max_poll_interval_ms = group->cache_config()->migration_copy_poll_max_interval_ms();
671+
options.connect_timeout_ms = group->cache_config()->migration_copy_connect_timeout_ms();
672+
options.submit_timeout_ms = group->cache_config()->migration_copy_submit_timeout_ms();
673+
options.query_timeout_ms = group->cache_config()->migration_copy_query_timeout_ms();
671674
}
672675
for (const auto &instance : instances) {
673676
if (!instance) {
@@ -3927,6 +3930,9 @@ void MigrationManager::RunAsyncMigrationPrepare(AsyncMigrationPrepareJob job, st
39273930
params.async_copy_options.operation_deadline_ms = cache_config->migration_copy_operation_deadline_ms();
39283931
params.async_copy_options.initial_poll_interval_ms = cache_config->migration_copy_poll_initial_interval_ms();
39293932
params.async_copy_options.max_poll_interval_ms = cache_config->migration_copy_poll_max_interval_ms();
3933+
params.async_copy_options.connect_timeout_ms = cache_config->migration_copy_connect_timeout_ms();
3934+
params.async_copy_options.submit_timeout_ms = cache_config->migration_copy_submit_timeout_ms();
3935+
params.async_copy_options.query_timeout_ms = cache_config->migration_copy_query_timeout_ms();
39303936
params.mark_timeout_ms = current_strategy->methods().mark().timeout_ms();
39313937
params.dedup_marks = true;
39323938
return DispatchMigrationBatchWithLifecycleLockHeld(job.trace_id,
@@ -4018,6 +4024,9 @@ MigrationManager::MigrateResult MigrationManager::MigrateCache(RequestContext *r
40184024
params.async_copy_options.initial_poll_interval_ms =
40194025
cache_config->migration_copy_poll_initial_interval_ms();
40204026
params.async_copy_options.max_poll_interval_ms = cache_config->migration_copy_poll_max_interval_ms();
4027+
params.async_copy_options.connect_timeout_ms = cache_config->migration_copy_connect_timeout_ms();
4028+
params.async_copy_options.submit_timeout_ms = cache_config->migration_copy_submit_timeout_ms();
4029+
params.async_copy_options.query_timeout_ms = cache_config->migration_copy_query_timeout_ms();
40214030
}
40224031
}
40234032
const auto dispatch =

kv_cache_manager/protocol/protobuf/admin_service.proto

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -270,6 +270,9 @@ message MigrationConfig {
270270
google.protobuf.Int64Value copy_operation_deadline_ms = 8;
271271
google.protobuf.Int64Value copy_poll_initial_interval_ms = 9;
272272
google.protobuf.Int64Value copy_poll_max_interval_ms = 10;
273+
google.protobuf.Int64Value copy_connect_timeout_ms = 11;
274+
google.protobuf.Int64Value copy_submit_timeout_ms = 12;
275+
google.protobuf.Int64Value copy_query_timeout_ms = 13;
273276
}
274277

275278
enum CachePreferStrategy {

0 commit comments

Comments
 (0)