feat: flow replace option (#198)

* feat: flow replace

* chore: make

* refactor: rename

* chore: make all

* refactor: more rename

* chore: make all
This commit is contained in:
discord9
2024-11-14 14:45:09 +08:00
committed by GitHub
parent 75c5fb5691
commit e1070ad3e7
8 changed files with 1695 additions and 169 deletions
+51 -20
View File
@@ -138,6 +138,7 @@ PROTOBUF_CONSTEXPR CreateRequest::CreateRequest(
, /*decltype(_impl_.sink_table_name_)*/nullptr
, /*decltype(_impl_.expire_after_)*/nullptr
, /*decltype(_impl_.create_if_not_exists_)*/false
, /*decltype(_impl_.or_replace_)*/false
, /*decltype(_impl_._cached_size_)*/{}} {}
struct CreateRequestDefaultTypeInternal {
PROTOBUF_CONSTEXPR CreateRequestDefaultTypeInternal()
@@ -270,6 +271,7 @@ const uint32_t TableStruct_greptime_2fv1_2fflow_2fserver_2eproto::offsets[] PROT
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.comment_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.sql_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.flow_options_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.or_replace_),
~0u, // no _has_bits_
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::DropRequest, _internal_metadata_),
~0u, // no _extensions_
@@ -295,8 +297,8 @@ static const ::_pbi::MigrationSchema schemas[] PROTOBUF_SECTION_VARIABLE(protode
{ 54, -1, -1, sizeof(::greptime::v1::flow::FlowResponse)},
{ 64, 72, -1, sizeof(::greptime::v1::flow::CreateRequest_FlowOptionsEntry_DoNotUse)},
{ 74, -1, -1, sizeof(::greptime::v1::flow::CreateRequest)},
{ 88, -1, -1, sizeof(::greptime::v1::flow::DropRequest)},
{ 95, -1, -1, sizeof(::greptime::v1::flow::FlushFlow)},
{ 89, -1, -1, sizeof(::greptime::v1::flow::DropRequest)},
{ 96, -1, -1, sizeof(::greptime::v1::flow::FlushFlow)},
};
static const ::_pb::Message* const file_default_instances[] = {
@@ -337,7 +339,7 @@ const char descriptor_table_protodef_greptime_2fv1_2fflow_2fserver_2eproto[] PRO
".greptime.v1.FlowId\022B\n\nextensions\030\004 \003(\0132"
"..greptime.v1.flow.FlowResponse.Extensio"
"nsEntry\0321\n\017ExtensionsEntry\022\013\n\003key\030\001 \001(\t\022"
"\r\n\005value\030\002 \001(\014:\0028\001\"\376\002\n\rCreateRequest\022$\n\007"
"\r\n\005value\030\002 \001(\014:\0028\001\"\222\003\n\rCreateRequest\022$\n\007"
"flow_id\030\001 \001(\0132\023.greptime.v1.FlowId\022.\n\020so"
"urce_table_ids\030\002 \003(\0132\024.greptime.v1.Table"
"Id\022/\n\017sink_table_name\030\003 \001(\0132\026.greptime.v"
@@ -345,18 +347,19 @@ const char descriptor_table_protodef_greptime_2fv1_2fflow_2fserver_2eproto[] PRO
"\010\022.\n\014expire_after\030\005 \001(\0132\030.greptime.v1.Ex"
"pireAfter\022\017\n\007comment\030\006 \001(\t\022\013\n\003sql\030\007 \001(\t\022"
"F\n\014flow_options\030\010 \003(\01320.greptime.v1.flow"
".CreateRequest.FlowOptionsEntry\0322\n\020FlowO"
"ptionsEntry\022\013\n\003key\030\001 \001(\t\022\r\n\005value\030\002 \001(\t:"
"\0028\001\"3\n\013DropRequest\022$\n\007flow_id\030\001 \001(\0132\023.gr"
"eptime.v1.FlowId\"1\n\tFlushFlow\022$\n\007flow_id"
"\030\001 \001(\0132\023.greptime.v1.FlowId2\264\001\n\004Flow\022S\n\022"
"HandleCreateRemove\022\035.greptime.v1.flow.Fl"
"owRequest\032\036.greptime.v1.flow.FlowRespons"
"e\022W\n\023HandleMirrorRequest\022 .greptime.v1.f"
"low.InsertRequests\032\036.greptime.v1.flow.Fl"
"owResponseBY\n\023io.greptime.v1.flowB\006Serve"
"rZ:github.com/GreptimeTeam/greptime-prot"
"o/go/greptime/v1/flowb\006proto3"
".CreateRequest.FlowOptionsEntry\022\022\n\nor_re"
"place\030\t \001(\010\0322\n\020FlowOptionsEntry\022\013\n\003key\030\001"
" \001(\t\022\r\n\005value\030\002 \001(\t:\0028\001\"3\n\013DropRequest\022$"
"\n\007flow_id\030\001 \001(\0132\023.greptime.v1.FlowId\"1\n\t"
"FlushFlow\022$\n\007flow_id\030\001 \001(\0132\023.greptime.v1"
".FlowId2\264\001\n\004Flow\022S\n\022HandleCreateRemove\022\035"
".greptime.v1.flow.FlowRequest\032\036.greptime"
".v1.flow.FlowResponse\022W\n\023HandleMirrorReq"
"uest\022 .greptime.v1.flow.InsertRequests\032\036"
".greptime.v1.flow.FlowResponseBY\n\023io.gre"
"ptime.v1.flowB\006ServerZ:github.com/Grepti"
"meTeam/greptime-proto/go/greptime/v1/flo"
"wb\006proto3"
;
static const ::_pbi::DescriptorTable* const descriptor_table_greptime_2fv1_2fflow_2fserver_2eproto_deps[3] = {
&::descriptor_table_greptime_2fv1_2fcommon_2eproto,
@@ -365,7 +368,7 @@ static const ::_pbi::DescriptorTable* const descriptor_table_greptime_2fv1_2fflo
};
static ::_pbi::once_flag descriptor_table_greptime_2fv1_2fflow_2fserver_2eproto_once;
const ::_pbi::DescriptorTable descriptor_table_greptime_2fv1_2fflow_2fserver_2eproto = {
false, false, 1709, descriptor_table_protodef_greptime_2fv1_2fflow_2fserver_2eproto,
false, false, 1729, descriptor_table_protodef_greptime_2fv1_2fflow_2fserver_2eproto,
"greptime/v1/flow/server.proto",
&descriptor_table_greptime_2fv1_2fflow_2fserver_2eproto_once, descriptor_table_greptime_2fv1_2fflow_2fserver_2eproto_deps, 3, 11,
schemas, file_default_instances, TableStruct_greptime_2fv1_2fflow_2fserver_2eproto::offsets,
@@ -1905,6 +1908,7 @@ CreateRequest::CreateRequest(const CreateRequest& from)
, decltype(_impl_.sink_table_name_){nullptr}
, decltype(_impl_.expire_after_){nullptr}
, decltype(_impl_.create_if_not_exists_){}
, decltype(_impl_.or_replace_){}
, /*decltype(_impl_._cached_size_)*/{}};
_internal_metadata_.MergeFrom<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(from._internal_metadata_);
@@ -1934,7 +1938,9 @@ CreateRequest::CreateRequest(const CreateRequest& from)
if (from._internal_has_expire_after()) {
_this->_impl_.expire_after_ = new ::greptime::v1::ExpireAfter(*from._impl_.expire_after_);
}
_this->_impl_.create_if_not_exists_ = from._impl_.create_if_not_exists_;
::memcpy(&_impl_.create_if_not_exists_, &from._impl_.create_if_not_exists_,
static_cast<size_t>(reinterpret_cast<char*>(&_impl_.or_replace_) -
reinterpret_cast<char*>(&_impl_.create_if_not_exists_)) + sizeof(_impl_.or_replace_));
// @@protoc_insertion_point(copy_constructor:greptime.v1.flow.CreateRequest)
}
@@ -1951,6 +1957,7 @@ inline void CreateRequest::SharedCtor(
, decltype(_impl_.sink_table_name_){nullptr}
, decltype(_impl_.expire_after_){nullptr}
, decltype(_impl_.create_if_not_exists_){false}
, decltype(_impl_.or_replace_){false}
, /*decltype(_impl_._cached_size_)*/{}
};
_impl_.comment_.InitDefault();
@@ -2015,7 +2022,9 @@ void CreateRequest::Clear() {
delete _impl_.expire_after_;
}
_impl_.expire_after_ = nullptr;
_impl_.create_if_not_exists_ = false;
::memset(&_impl_.create_if_not_exists_, 0, static_cast<size_t>(
reinterpret_cast<char*>(&_impl_.or_replace_) -
reinterpret_cast<char*>(&_impl_.create_if_not_exists_)) + sizeof(_impl_.or_replace_));
_internal_metadata_.Clear<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>();
}
@@ -2103,6 +2112,14 @@ const char* CreateRequest::_InternalParse(const char* ptr, ::_pbi::ParseContext*
} else
goto handle_unusual;
continue;
// bool or_replace = 9;
case 9:
if (PROTOBUF_PREDICT_TRUE(static_cast<uint8_t>(tag) == 72)) {
_impl_.or_replace_ = ::PROTOBUF_NAMESPACE_ID::internal::ReadVarint64(&ptr);
CHK_(ptr);
} else
goto handle_unusual;
continue;
default:
goto handle_unusual;
} // switch
@@ -2217,6 +2234,12 @@ uint8_t* CreateRequest::_InternalSerialize(
}
}
// bool or_replace = 9;
if (this->_internal_or_replace() != 0) {
target = stream->EnsureSpace(target);
target = ::_pbi::WireFormatLite::WriteBoolToArray(9, this->_internal_or_replace(), 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);
@@ -2289,6 +2312,11 @@ size_t CreateRequest::ByteSizeLong() const {
total_size += 1 + 1;
}
// bool or_replace = 9;
if (this->_internal_or_replace() != 0) {
total_size += 1 + 1;
}
return MaybeComputeUnknownFieldsSize(total_size, &_impl_._cached_size_);
}
@@ -2330,6 +2358,9 @@ void CreateRequest::MergeImpl(::PROTOBUF_NAMESPACE_ID::Message& to_msg, const ::
if (from._internal_create_if_not_exists() != 0) {
_this->_internal_set_create_if_not_exists(from._internal_create_if_not_exists());
}
if (from._internal_or_replace() != 0) {
_this->_internal_set_or_replace(from._internal_or_replace());
}
_this->_internal_metadata_.MergeFrom<::PROTOBUF_NAMESPACE_ID::UnknownFieldSet>(from._internal_metadata_);
}
@@ -2360,8 +2391,8 @@ void CreateRequest::InternalSwap(CreateRequest* other) {
&other->_impl_.sql_, rhs_arena
);
::PROTOBUF_NAMESPACE_ID::internal::memswap<
PROTOBUF_FIELD_OFFSET(CreateRequest, _impl_.create_if_not_exists_)
+ sizeof(CreateRequest::_impl_.create_if_not_exists_)
PROTOBUF_FIELD_OFFSET(CreateRequest, _impl_.or_replace_)
+ sizeof(CreateRequest::_impl_.or_replace_)
- PROTOBUF_FIELD_OFFSET(CreateRequest, _impl_.flow_id_)>(
reinterpret_cast<char*>(&_impl_.flow_id_),
reinterpret_cast<char*>(&other->_impl_.flow_id_));
+31
View File
@@ -1282,6 +1282,7 @@ class CreateRequest final :
kSinkTableNameFieldNumber = 3,
kExpireAfterFieldNumber = 5,
kCreateIfNotExistsFieldNumber = 4,
kOrReplaceFieldNumber = 9,
};
// repeated .greptime.v1.TableId source_table_ids = 2;
int source_table_ids_size() const;
@@ -1409,6 +1410,15 @@ class CreateRequest final :
void _internal_set_create_if_not_exists(bool value);
public:
// bool or_replace = 9;
void clear_or_replace();
bool or_replace() const;
void set_or_replace(bool value);
private:
bool _internal_or_replace() const;
void _internal_set_or_replace(bool value);
public:
// @@protoc_insertion_point(class_scope:greptime.v1.flow.CreateRequest)
private:
class _Internal;
@@ -1429,6 +1439,7 @@ class CreateRequest final :
::greptime::v1::TableName* sink_table_name_;
::greptime::v1::ExpireAfter* expire_after_;
bool create_if_not_exists_;
bool or_replace_;
mutable ::PROTOBUF_NAMESPACE_ID::internal::CachedSize _cached_size_;
};
union { Impl_ _impl_; };
@@ -2977,6 +2988,26 @@ CreateRequest::mutable_flow_options() {
return _internal_mutable_flow_options();
}
// bool or_replace = 9;
inline void CreateRequest::clear_or_replace() {
_impl_.or_replace_ = false;
}
inline bool CreateRequest::_internal_or_replace() const {
return _impl_.or_replace_;
}
inline bool CreateRequest::or_replace() const {
// @@protoc_insertion_point(field_get:greptime.v1.flow.CreateRequest.or_replace)
return _internal_or_replace();
}
inline void CreateRequest::_internal_set_or_replace(bool value) {
_impl_.or_replace_ = value;
}
inline void CreateRequest::set_or_replace(bool value) {
_internal_set_or_replace(value);
// @@protoc_insertion_point(field_set:greptime.v1.flow.CreateRequest.or_replace)
}
// -------------------------------------------------------------------
// DropRequest