feat: add wait and timeout field (#303)

* feat: add `wait` and `timeout` field

* chore: remove `wait` from `Repartition`

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fmt

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fmt

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fmt

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-01-20 16:41:47 +08:00
committed by GitHub
parent 1353b0ada9
commit 7755821c27
10 changed files with 1191 additions and 931 deletions
+101 -40
View File
@@ -251,6 +251,8 @@ PROTOBUF_CONSTEXPR DdlTaskRequest::DdlTaskRequest(
::_pbi::ConstantInitialized): _impl_{
/*decltype(_impl_.header_)*/nullptr
, /*decltype(_impl_.query_context_)*/nullptr
, /*decltype(_impl_.wait_)*/false
, /*decltype(_impl_.timeout_secs_)*/0u
, /*decltype(_impl_.task_)*/{}
, /*decltype(_impl_._cached_size_)*/{}
, /*decltype(_impl_._oneof_case_)*/{}} {}
@@ -416,6 +418,8 @@ const uint32_t TableStruct_greptime_2fv1_2fmeta_2fddl_2eproto::offsets[] PROTOBU
~0u, // no _inlined_string_donated_
PROTOBUF_FIELD_OFFSET(::greptime::v1::meta::DdlTaskRequest, _impl_.header_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::meta::DdlTaskRequest, _impl_.query_context_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::meta::DdlTaskRequest, _impl_.wait_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::meta::DdlTaskRequest, _impl_.timeout_secs_),
::_pbi::kInvalidFieldOffsetTag,
::_pbi::kInvalidFieldOffsetTag,
::_pbi::kInvalidFieldOffsetTag,
@@ -463,7 +467,7 @@ static const ::_pbi::MigrationSchema schemas[] PROTOBUF_SECTION_VARIABLE(protode
{ 108, -1, -1, sizeof(::greptime::v1::meta::DropTriggerTask)},
{ 115, -1, -1, sizeof(::greptime::v1::meta::CommentOnTask)},
{ 122, -1, -1, sizeof(::greptime::v1::meta::DdlTaskRequest)},
{ 148, -1, -1, sizeof(::greptime::v1::meta::DdlTaskResponse)},
{ 150, -1, -1, sizeof(::greptime::v1::meta::DdlTaskResponse)},
};
static const ::_pb::Message* const file_default_instances[] = {
@@ -524,44 +528,45 @@ const char descriptor_table_protodef_greptime_2fv1_2fmeta_2fddl_2eproto[] PROTOB
"ggerExpr\"E\n\017DropTriggerTask\0222\n\014drop_trig"
"ger\030\001 \001(\0132\034.greptime.v1.DropTriggerExpr\""
"\?\n\rCommentOnTask\022.\n\ncomment_on\030\001 \001(\0132\032.g"
"reptime.v1.CommentOnExpr\"\265\t\n\016DdlTaskRequ"
"reptime.v1.CommentOnExpr\"\331\t\n\016DdlTaskRequ"
"est\022/\n\006header\030\001 \001(\0132\037.greptime.v1.meta.R"
"equestHeader\0220\n\rquery_context\030@ \001(\0132\031.gr"
"eptime.v1.QueryContext\022>\n\021create_table_t"
"ask\030\002 \001(\0132!.greptime.v1.meta.CreateTable"
"TaskH\000\022:\n\017drop_table_task\030\003 \001(\0132\037.grepti"
"me.v1.meta.DropTableTaskH\000\022<\n\020alter_tabl"
"e_task\030\004 \001(\0132 .greptime.v1.meta.AlterTab"
"leTaskH\000\022B\n\023truncate_table_task\030\005 \001(\0132#."
"greptime.v1.meta.TruncateTableTaskH\000\022@\n\022"
"create_table_tasks\030\006 \001(\0132\".greptime.v1.m"
"eta.CreateTableTasksH\000\022<\n\020drop_table_tas"
"ks\030\007 \001(\0132 .greptime.v1.meta.DropTableTas"
"ksH\000\022>\n\021alter_table_tasks\030\010 \001(\0132!.grepti"
"me.v1.meta.AlterTableTasksH\000\022@\n\022drop_dat"
"abase_task\030\t \001(\0132\".greptime.v1.meta.Drop"
"DatabaseTaskH\000\022D\n\024create_database_task\030\n"
" \001(\0132$.greptime.v1.meta.CreateDatabaseTa"
"skH\000\022<\n\020create_flow_task\030\013 \001(\0132 .greptim"
"e.v1.meta.CreateFlowTaskH\000\0228\n\016drop_flow_"
"task\030\014 \001(\0132\036.greptime.v1.meta.DropFlowTa"
"skH\000\022<\n\020create_view_task\030\r \001(\0132 .greptim"
"e.v1.meta.CreateViewTaskH\000\0228\n\016drop_view_"
"task\030\016 \001(\0132\036.greptime.v1.meta.DropViewTa"
"skH\000\022B\n\023alter_database_task\030\017 \001(\0132#.grep"
"time.v1.meta.AlterDatabaseTaskH\000\022B\n\023crea"
"te_trigger_task\030\020 \001(\0132#.greptime.v1.meta"
".CreateTriggerTaskH\000\022>\n\021drop_trigger_tas"
"k\030\021 \001(\0132!.greptime.v1.meta.DropTriggerTa"
"skH\000\022:\n\017comment_on_task\030\022 \001(\0132\037.greptime"
".v1.meta.CommentOnTaskH\000B\006\n\004task\"\230\001\n\017Ddl"
"TaskResponse\0220\n\006header\030\001 \001(\0132 .greptime."
"v1.meta.ResponseHeader\022*\n\003pid\030\002 \001(\0132\035.gr"
"eptime.v1.meta.ProcedureId\022\'\n\ttable_ids\030"
"\005 \003(\0132\024.greptime.v1.TableId*#\n\013DdlTaskTy"
"pe\022\n\n\006Create\020\000\022\010\n\004Drop\020\001B<Z:github.com/G"
"reptimeTeam/greptime-proto/go/greptime/v"
"1/metab\006proto3"
"eptime.v1.QueryContext\022\014\n\004wait\030A \001(\010\022\024\n\014"
"timeout_secs\030B \001(\r\022>\n\021create_table_task\030"
"\002 \001(\0132!.greptime.v1.meta.CreateTableTask"
"H\000\022:\n\017drop_table_task\030\003 \001(\0132\037.greptime.v"
"1.meta.DropTableTaskH\000\022<\n\020alter_table_ta"
"sk\030\004 \001(\0132 .greptime.v1.meta.AlterTableTa"
"skH\000\022B\n\023truncate_table_task\030\005 \001(\0132#.grep"
"time.v1.meta.TruncateTableTaskH\000\022@\n\022crea"
"te_table_tasks\030\006 \001(\0132\".greptime.v1.meta."
"CreateTableTasksH\000\022<\n\020drop_table_tasks\030\007"
" \001(\0132 .greptime.v1.meta.DropTableTasksH\000"
"\022>\n\021alter_table_tasks\030\010 \001(\0132!.greptime.v"
"1.meta.AlterTableTasksH\000\022@\n\022drop_databas"
"e_task\030\t \001(\0132\".greptime.v1.meta.DropData"
"baseTaskH\000\022D\n\024create_database_task\030\n \001(\013"
"2$.greptime.v1.meta.CreateDatabaseTaskH\000"
"\022<\n\020create_flow_task\030\013 \001(\0132 .greptime.v1"
".meta.CreateFlowTaskH\000\0228\n\016drop_flow_task"
"\030\014 \001(\0132\036.greptime.v1.meta.DropFlowTaskH\000"
"\022<\n\020create_view_task\030\r \001(\0132 .greptime.v1"
".meta.CreateViewTaskH\000\0228\n\016drop_view_task"
"\030\016 \001(\0132\036.greptime.v1.meta.DropViewTaskH\000"
"\022B\n\023alter_database_task\030\017 \001(\0132#.greptime"
".v1.meta.AlterDatabaseTaskH\000\022B\n\023create_t"
"rigger_task\030\020 \001(\0132#.greptime.v1.meta.Cre"
"ateTriggerTaskH\000\022>\n\021drop_trigger_task\030\021 "
"\001(\0132!.greptime.v1.meta.DropTriggerTaskH\000"
"\022:\n\017comment_on_task\030\022 \001(\0132\037.greptime.v1."
"meta.CommentOnTaskH\000B\006\n\004task\"\230\001\n\017DdlTask"
"Response\0220\n\006header\030\001 \001(\0132 .greptime.v1.m"
"eta.ResponseHeader\022*\n\003pid\030\002 \001(\0132\035.grepti"
"me.v1.meta.ProcedureId\022\'\n\ttable_ids\030\005 \003("
"\0132\024.greptime.v1.TableId*#\n\013DdlTaskType\022\n"
"\n\006Create\020\000\022\010\n\004Drop\020\001B<Z:github.com/Grept"
"imeTeam/greptime-proto/go/greptime/v1/me"
"tab\006proto3"
;
static const ::_pbi::DescriptorTable* const descriptor_table_greptime_2fv1_2fmeta_2fddl_2eproto_deps[4] = {
&::descriptor_table_greptime_2fv1_2fcommon_2eproto,
@@ -571,7 +576,7 @@ static const ::_pbi::DescriptorTable* const descriptor_table_greptime_2fv1_2fmet
};
static ::_pbi::once_flag descriptor_table_greptime_2fv1_2fmeta_2fddl_2eproto_once;
const ::_pbi::DescriptorTable descriptor_table_greptime_2fv1_2fmeta_2fddl_2eproto = {
false, false, 2894, descriptor_table_protodef_greptime_2fv1_2fmeta_2fddl_2eproto,
false, false, 2930, descriptor_table_protodef_greptime_2fv1_2fmeta_2fddl_2eproto,
"greptime/v1/meta/ddl.proto",
&descriptor_table_greptime_2fv1_2fmeta_2fddl_2eproto_once, descriptor_table_greptime_2fv1_2fmeta_2fddl_2eproto_deps, 4, 19,
schemas, file_default_instances, TableStruct_greptime_2fv1_2fmeta_2fddl_2eproto::offsets,
@@ -4454,6 +4459,8 @@ DdlTaskRequest::DdlTaskRequest(const DdlTaskRequest& from)
new (&_impl_) Impl_{
decltype(_impl_.header_){nullptr}
, decltype(_impl_.query_context_){nullptr}
, decltype(_impl_.wait_){}
, decltype(_impl_.timeout_secs_){}
, decltype(_impl_.task_){}
, /*decltype(_impl_._cached_size_)*/{}
, /*decltype(_impl_._oneof_case_)*/{}};
@@ -4465,6 +4472,9 @@ DdlTaskRequest::DdlTaskRequest(const DdlTaskRequest& from)
if (from._internal_has_query_context()) {
_this->_impl_.query_context_ = new ::greptime::v1::QueryContext(*from._impl_.query_context_);
}
::memcpy(&_impl_.wait_, &from._impl_.wait_,
static_cast<size_t>(reinterpret_cast<char*>(&_impl_.timeout_secs_) -
reinterpret_cast<char*>(&_impl_.wait_)) + sizeof(_impl_.timeout_secs_));
clear_has_task();
switch (from.task_case()) {
case kCreateTableTask: {
@@ -4566,6 +4576,8 @@ inline void DdlTaskRequest::SharedCtor(
new (&_impl_) Impl_{
decltype(_impl_.header_){nullptr}
, decltype(_impl_.query_context_){nullptr}
, decltype(_impl_.wait_){false}
, decltype(_impl_.timeout_secs_){0u}
, decltype(_impl_.task_){}
, /*decltype(_impl_._cached_size_)*/{}
, /*decltype(_impl_._oneof_case_)*/{}
@@ -4722,6 +4734,9 @@ void DdlTaskRequest::Clear() {
delete _impl_.query_context_;
}
_impl_.query_context_ = nullptr;
::memset(&_impl_.wait_, 0, static_cast<size_t>(
reinterpret_cast<char*>(&_impl_.timeout_secs_) -
reinterpret_cast<char*>(&_impl_.wait_)) + sizeof(_impl_.timeout_secs_));
clear_task();
_internal_metadata_.Clear<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>();
}
@@ -4884,6 +4899,22 @@ const char* DdlTaskRequest::_InternalParse(const char* ptr, ::_pbi::ParseContext
} else
goto handle_unusual;
continue;
// bool wait = 65;
case 65:
if (PROTOBUF_PREDICT_TRUE(static_cast<uint8_t>(tag) == 8)) {
_impl_.wait_ = ::PROTOBUF_NAMESPACE_ID::internal::ReadVarint64(&ptr);
CHK_(ptr);
} else
goto handle_unusual;
continue;
// uint32 timeout_secs = 66;
case 66:
if (PROTOBUF_PREDICT_TRUE(static_cast<uint8_t>(tag) == 16)) {
_impl_.timeout_secs_ = ::PROTOBUF_NAMESPACE_ID::internal::ReadVarint32(&ptr);
CHK_(ptr);
} else
goto handle_unusual;
continue;
default:
goto handle_unusual;
} // switch
@@ -5046,6 +5077,18 @@ uint8_t* DdlTaskRequest::_InternalSerialize(
_Internal::query_context(this).GetCachedSize(), target, stream);
}
// bool wait = 65;
if (this->_internal_wait() != 0) {
target = stream->EnsureSpace(target);
target = ::_pbi::WireFormatLite::WriteBoolToArray(65, this->_internal_wait(), target);
}
// uint32 timeout_secs = 66;
if (this->_internal_timeout_secs() != 0) {
target = stream->EnsureSpace(target);
target = ::_pbi::WireFormatLite::WriteUInt32ToArray(66, this->_internal_timeout_secs(), target);
}
if (PROTOBUF_PREDICT_FALSE(_internal_metadata_.have_unknown_fields())) {
target = ::_pbi::WireFormat::InternalSerializeUnknownFieldsToArray(
_internal_metadata_.unknown_fields<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(::PROTOBUF_NAMESPACE_ID::UnknownFieldSet::default_instance), target, stream);
@@ -5076,6 +5119,18 @@ size_t DdlTaskRequest::ByteSizeLong() const {
*_impl_.query_context_);
}
// bool wait = 65;
if (this->_internal_wait() != 0) {
total_size += 2 + 1;
}
// uint32 timeout_secs = 66;
if (this->_internal_timeout_secs() != 0) {
total_size += 2 +
::_pbi::WireFormatLite::UInt32Size(
this->_internal_timeout_secs());
}
switch (task_case()) {
// .greptime.v1.meta.CreateTableTask create_table_task = 2;
case kCreateTableTask: {
@@ -5226,6 +5281,12 @@ void DdlTaskRequest::MergeImpl(::PROTOBUF_NAMESPACE_ID::Message& to_msg, const :
_this->_internal_mutable_query_context()->::greptime::v1::QueryContext::MergeFrom(
from._internal_query_context());
}
if (from._internal_wait() != 0) {
_this->_internal_set_wait(from._internal_wait());
}
if (from._internal_timeout_secs() != 0) {
_this->_internal_set_timeout_secs(from._internal_timeout_secs());
}
switch (from.task_case()) {
case kCreateTableTask: {
_this->_internal_mutable_create_table_task()->::greptime::v1::meta::CreateTableTask::MergeFrom(
@@ -5334,8 +5395,8 @@ void DdlTaskRequest::InternalSwap(DdlTaskRequest* other) {
using std::swap;
_internal_metadata_.InternalSwap(&other->_internal_metadata_);
::PROTOBUF_NAMESPACE_ID::internal::memswap<
PROTOBUF_FIELD_OFFSET(DdlTaskRequest, _impl_.query_context_)
+ sizeof(DdlTaskRequest::_impl_.query_context_)
PROTOBUF_FIELD_OFFSET(DdlTaskRequest, _impl_.timeout_secs_)
+ sizeof(DdlTaskRequest::_impl_.timeout_secs_)
- PROTOBUF_FIELD_OFFSET(DdlTaskRequest, _impl_.header_)>(
reinterpret_cast<char*>(&_impl_.header_),
reinterpret_cast<char*>(&other->_impl_.header_));
+62
View File
@@ -3029,6 +3029,8 @@ class DdlTaskRequest final :
enum : int {
kHeaderFieldNumber = 1,
kQueryContextFieldNumber = 64,
kWaitFieldNumber = 65,
kTimeoutSecsFieldNumber = 66,
kCreateTableTaskFieldNumber = 2,
kDropTableTaskFieldNumber = 3,
kAlterTableTaskFieldNumber = 4,
@@ -3083,6 +3085,24 @@ class DdlTaskRequest final :
::greptime::v1::QueryContext* query_context);
::greptime::v1::QueryContext* unsafe_arena_release_query_context();
// bool wait = 65;
void clear_wait();
bool wait() const;
void set_wait(bool value);
private:
bool _internal_wait() const;
void _internal_set_wait(bool value);
public:
// uint32 timeout_secs = 66;
void clear_timeout_secs();
uint32_t timeout_secs() const;
void set_timeout_secs(uint32_t value);
private:
uint32_t _internal_timeout_secs() const;
void _internal_set_timeout_secs(uint32_t value);
public:
// .greptime.v1.meta.CreateTableTask create_table_task = 2;
bool has_create_table_task() const;
private:
@@ -3421,6 +3441,8 @@ class DdlTaskRequest final :
struct Impl_ {
::greptime::v1::meta::RequestHeader* header_;
::greptime::v1::QueryContext* query_context_;
bool wait_;
uint32_t timeout_secs_;
union TaskUnion {
constexpr TaskUnion() : _constinit_{} {}
::PROTOBUF_NAMESPACE_ID::internal::ConstantInitialized _constinit_;
@@ -5342,6 +5364,46 @@ inline void DdlTaskRequest::set_allocated_query_context(::greptime::v1::QueryCon
// @@protoc_insertion_point(field_set_allocated:greptime.v1.meta.DdlTaskRequest.query_context)
}
// bool wait = 65;
inline void DdlTaskRequest::clear_wait() {
_impl_.wait_ = false;
}
inline bool DdlTaskRequest::_internal_wait() const {
return _impl_.wait_;
}
inline bool DdlTaskRequest::wait() const {
// @@protoc_insertion_point(field_get:greptime.v1.meta.DdlTaskRequest.wait)
return _internal_wait();
}
inline void DdlTaskRequest::_internal_set_wait(bool value) {
_impl_.wait_ = value;
}
inline void DdlTaskRequest::set_wait(bool value) {
_internal_set_wait(value);
// @@protoc_insertion_point(field_set:greptime.v1.meta.DdlTaskRequest.wait)
}
// uint32 timeout_secs = 66;
inline void DdlTaskRequest::clear_timeout_secs() {
_impl_.timeout_secs_ = 0u;
}
inline uint32_t DdlTaskRequest::_internal_timeout_secs() const {
return _impl_.timeout_secs_;
}
inline uint32_t DdlTaskRequest::timeout_secs() const {
// @@protoc_insertion_point(field_get:greptime.v1.meta.DdlTaskRequest.timeout_secs)
return _internal_timeout_secs();
}
inline void DdlTaskRequest::_internal_set_timeout_secs(uint32_t value) {
_impl_.timeout_secs_ = value;
}
inline void DdlTaskRequest::set_timeout_secs(uint32_t value) {
_internal_set_timeout_secs(value);
// @@protoc_insertion_point(field_set:greptime.v1.meta.DdlTaskRequest.timeout_secs)
}
// .greptime.v1.meta.CreateTableTask create_table_task = 2;
inline bool DdlTaskRequest::_internal_has_create_table_task() const {
return task_case() == kCreateTableTask;