feat: flow add eval interval (#273)

* feat: create flow eval interval

Signed-off-by: discord9 <discord9@163.com>

* fix: eval interval

Signed-off-by: discord9 <discord9@163.com>

* chore: make all

Signed-off-by: discord9 <discord9@163.com>

* fix: flow create request

Signed-off-by: discord9 <discord9@163.com>

* chore: make all

Signed-off-by: discord9 <discord9@163.com>

* chore: fmt

Signed-off-by: discord9 <discord9@163.com>

---------

Signed-off-by: discord9 <discord9@163.com>
This commit is contained in:
discord9
2025-08-27 16:21:27 +08:00
committed by GitHub
parent df2bb74b59
commit 66eb089afa
15 changed files with 2513 additions and 802 deletions
+72 -22
View File
@@ -165,6 +165,7 @@ PROTOBUF_CONSTEXPR CreateRequest::CreateRequest(
, /*decltype(_impl_.flow_id_)*/nullptr
, /*decltype(_impl_.sink_table_name_)*/nullptr
, /*decltype(_impl_.expire_after_)*/nullptr
, /*decltype(_impl_.eval_interval_)*/nullptr
, /*decltype(_impl_.create_if_not_exists_)*/false
, /*decltype(_impl_.or_replace_)*/false
, /*decltype(_impl_._cached_size_)*/{}} {}
@@ -311,6 +312,7 @@ const uint32_t TableStruct_greptime_2fv1_2fflow_2fserver_2eproto::offsets[] PROT
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.sink_table_name_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.create_if_not_exists_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.expire_after_),
PROTOBUF_FIELD_OFFSET(::greptime::v1::flow::CreateRequest, _impl_.eval_interval_),
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_),
@@ -342,8 +344,8 @@ static const ::_pbi::MigrationSchema schemas[] PROTOBUF_SECTION_VARIABLE(protode
{ 69, -1, -1, sizeof(::greptime::v1::flow::FlowResponse)},
{ 79, 87, -1, sizeof(::greptime::v1::flow::CreateRequest_FlowOptionsEntry_DoNotUse)},
{ 89, -1, -1, sizeof(::greptime::v1::flow::CreateRequest)},
{ 104, -1, -1, sizeof(::greptime::v1::flow::DropRequest)},
{ 111, -1, -1, sizeof(::greptime::v1::flow::FlushFlow)},
{ 105, -1, -1, sizeof(::greptime::v1::flow::DropRequest)},
{ 112, -1, -1, sizeof(::greptime::v1::flow::FlushFlow)},
};
static const ::_pb::Message* const file_default_instances[] = {
@@ -389,30 +391,31 @@ const char descriptor_table_protodef_greptime_2fv1_2fflow_2fserver_2eproto[] PRO
"ected_flows\030\003 \003(\0132\023.greptime.v1.FlowId\022B"
"\n\nextensions\030\004 \003(\0132..greptime.v1.flow.Fl"
"owResponse.ExtensionsEntry\0321\n\017Extensions"
"Entry\022\013\n\003key\030\001 \001(\t\022\r\n\005value\030\002 \001(\014:\0028\001\"\222\003"
"Entry\022\013\n\003key\030\001 \001(\t\022\r\n\005value\030\002 \001(\014:\0028\001\"\304\003"
"\n\rCreateRequest\022$\n\007flow_id\030\001 \001(\0132\023.grept"
"ime.v1.FlowId\022.\n\020source_table_ids\030\002 \003(\0132"
"\024.greptime.v1.TableId\022/\n\017sink_table_name"
"\030\003 \001(\0132\026.greptime.v1.TableName\022\034\n\024create"
"_if_not_exists\030\004 \001(\010\022.\n\014expire_after\030\005 \001"
"(\0132\030.greptime.v1.ExpireAfter\022\017\n\007comment\030"
"\006 \001(\t\022\013\n\003sql\030\007 \001(\t\022F\n\014flow_options\030\010 \003(\013"
"20.greptime.v1.flow.CreateRequest.FlowOp"
"tionsEntry\022\022\n\nor_replace\030\t \001(\010\0322\n\020FlowOp"
"tionsEntry\022\013\n\003key\030\001 \001(\t\022\r\n\005value\030\002 \001(\t:\002"
"8\001\"3\n\013DropRequest\022$\n\007flow_id\030\001 \001(\0132\023.gre"
"ptime.v1.FlowId\"1\n\tFlushFlow\022$\n\007flow_id\030"
"\001 \001(\0132\023.greptime.v1.FlowId2\230\002\n\004Flow\022S\n\022H"
"andleCreateRemove\022\035.greptime.v1.flow.Flo"
"wRequest\032\036.greptime.v1.flow.FlowResponse"
"\022W\n\023HandleMirrorRequest\022 .greptime.v1.fl"
"ow.InsertRequests\032\036.greptime.v1.flow.Flo"
"wResponse\022b\n\031HandleMarkDirtyTimeWindow\022%"
".greptime.v1.flow.DirtyWindowRequests\032\036."
"greptime.v1.flow.FlowResponseBY\n\023io.grep"
"time.v1.flowB\006ServerZ:github.com/Greptim"
"eTeam/greptime-proto/go/greptime/v1/flow"
"b\006proto3"
"(\0132\030.greptime.v1.ExpireAfter\0220\n\reval_int"
"erval\030\n \001(\0132\031.greptime.v1.EvalInterval\022\017"
"\n\007comment\030\006 \001(\t\022\013\n\003sql\030\007 \001(\t\022F\n\014flow_opt"
"ions\030\010 \003(\01320.greptime.v1.flow.CreateRequ"
"est.FlowOptionsEntry\022\022\n\nor_replace\030\t \001(\010"
"\0322\n\020FlowOptionsEntry\022\013\n\003key\030\001 \001(\t\022\r\n\005val"
"ue\030\002 \001(\t:\0028\001\"3\n\013DropRequest\022$\n\007flow_id\030\001"
" \001(\0132\023.greptime.v1.FlowId\"1\n\tFlushFlow\022$"
"\n\007flow_id\030\001 \001(\0132\023.greptime.v1.FlowId2\230\002\n"
"\004Flow\022S\n\022HandleCreateRemove\022\035.greptime.v"
"1.flow.FlowRequest\032\036.greptime.v1.flow.Fl"
"owResponse\022W\n\023HandleMirrorRequest\022 .grep"
"time.v1.flow.InsertRequests\032\036.greptime.v"
"1.flow.FlowResponse\022b\n\031HandleMarkDirtyTi"
"meWindow\022%.greptime.v1.flow.DirtyWindowR"
"equests\032\036.greptime.v1.flow.FlowResponseB"
"Y\n\023io.greptime.v1.flowB\006ServerZ:github.c"
"om/GreptimeTeam/greptime-proto/go/grepti"
"me/v1/flowb\006proto3"
;
static const ::_pbi::DescriptorTable* const descriptor_table_greptime_2fv1_2fflow_2fserver_2eproto_deps[3] = {
&::descriptor_table_greptime_2fv1_2fcommon_2eproto,
@@ -421,7 +424,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, 1968, descriptor_table_protodef_greptime_2fv1_2fflow_2fserver_2eproto,
false, false, 2018, 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, 13,
schemas, file_default_instances, TableStruct_greptime_2fv1_2fflow_2fserver_2eproto::offsets,
@@ -2310,6 +2313,7 @@ class CreateRequest::_Internal {
static const ::greptime::v1::FlowId& flow_id(const CreateRequest* msg);
static const ::greptime::v1::TableName& sink_table_name(const CreateRequest* msg);
static const ::greptime::v1::ExpireAfter& expire_after(const CreateRequest* msg);
static const ::greptime::v1::EvalInterval& eval_interval(const CreateRequest* msg);
};
const ::greptime::v1::FlowId&
@@ -2324,6 +2328,10 @@ const ::greptime::v1::ExpireAfter&
CreateRequest::_Internal::expire_after(const CreateRequest* msg) {
return *msg->_impl_.expire_after_;
}
const ::greptime::v1::EvalInterval&
CreateRequest::_Internal::eval_interval(const CreateRequest* msg) {
return *msg->_impl_.eval_interval_;
}
void CreateRequest::clear_flow_id() {
if (GetArenaForAllocation() == nullptr && _impl_.flow_id_ != nullptr) {
delete _impl_.flow_id_;
@@ -2345,6 +2353,12 @@ void CreateRequest::clear_expire_after() {
}
_impl_.expire_after_ = nullptr;
}
void CreateRequest::clear_eval_interval() {
if (GetArenaForAllocation() == nullptr && _impl_.eval_interval_ != nullptr) {
delete _impl_.eval_interval_;
}
_impl_.eval_interval_ = nullptr;
}
CreateRequest::CreateRequest(::PROTOBUF_NAMESPACE_ID::Arena* arena,
bool is_message_owned)
: ::PROTOBUF_NAMESPACE_ID::Message(arena, is_message_owned) {
@@ -2365,6 +2379,7 @@ CreateRequest::CreateRequest(const CreateRequest& from)
, decltype(_impl_.flow_id_){nullptr}
, decltype(_impl_.sink_table_name_){nullptr}
, decltype(_impl_.expire_after_){nullptr}
, decltype(_impl_.eval_interval_){nullptr}
, decltype(_impl_.create_if_not_exists_){}
, decltype(_impl_.or_replace_){}
, /*decltype(_impl_._cached_size_)*/{}};
@@ -2396,6 +2411,9 @@ CreateRequest::CreateRequest(const CreateRequest& from)
if (from._internal_has_expire_after()) {
_this->_impl_.expire_after_ = new ::greptime::v1::ExpireAfter(*from._impl_.expire_after_);
}
if (from._internal_has_eval_interval()) {
_this->_impl_.eval_interval_ = new ::greptime::v1::EvalInterval(*from._impl_.eval_interval_);
}
::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_));
@@ -2414,6 +2432,7 @@ inline void CreateRequest::SharedCtor(
, decltype(_impl_.flow_id_){nullptr}
, decltype(_impl_.sink_table_name_){nullptr}
, decltype(_impl_.expire_after_){nullptr}
, decltype(_impl_.eval_interval_){nullptr}
, decltype(_impl_.create_if_not_exists_){false}
, decltype(_impl_.or_replace_){false}
, /*decltype(_impl_._cached_size_)*/{}
@@ -2448,6 +2467,7 @@ inline void CreateRequest::SharedDtor() {
if (this != internal_default_instance()) delete _impl_.flow_id_;
if (this != internal_default_instance()) delete _impl_.sink_table_name_;
if (this != internal_default_instance()) delete _impl_.expire_after_;
if (this != internal_default_instance()) delete _impl_.eval_interval_;
}
void CreateRequest::ArenaDtor(void* object) {
@@ -2480,6 +2500,10 @@ void CreateRequest::Clear() {
delete _impl_.expire_after_;
}
_impl_.expire_after_ = nullptr;
if (GetArenaForAllocation() == nullptr && _impl_.eval_interval_ != nullptr) {
delete _impl_.eval_interval_;
}
_impl_.eval_interval_ = nullptr;
::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_));
@@ -2578,6 +2602,14 @@ const char* CreateRequest::_InternalParse(const char* ptr, ::_pbi::ParseContext*
} else
goto handle_unusual;
continue;
// .greptime.v1.EvalInterval eval_interval = 10;
case 10:
if (PROTOBUF_PREDICT_TRUE(static_cast<uint8_t>(tag) == 82)) {
ptr = ctx->ParseMessage(_internal_mutable_eval_interval(), ptr);
CHK_(ptr);
} else
goto handle_unusual;
continue;
default:
goto handle_unusual;
} // switch
@@ -2698,6 +2730,13 @@ uint8_t* CreateRequest::_InternalSerialize(
target = ::_pbi::WireFormatLite::WriteBoolToArray(9, this->_internal_or_replace(), target);
}
// .greptime.v1.EvalInterval eval_interval = 10;
if (this->_internal_has_eval_interval()) {
target = ::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::
InternalWriteMessage(10, _Internal::eval_interval(this),
_Internal::eval_interval(this).GetCachedSize(), target, stream);
}
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);
@@ -2765,6 +2804,13 @@ size_t CreateRequest::ByteSizeLong() const {
*_impl_.expire_after_);
}
// .greptime.v1.EvalInterval eval_interval = 10;
if (this->_internal_has_eval_interval()) {
total_size += 1 +
::PROTOBUF_NAMESPACE_ID::internal::WireFormatLite::MessageSize(
*_impl_.eval_interval_);
}
// bool create_if_not_exists = 4;
if (this->_internal_create_if_not_exists() != 0) {
total_size += 1 + 1;
@@ -2813,6 +2859,10 @@ void CreateRequest::MergeImpl(::PROTOBUF_NAMESPACE_ID::Message& to_msg, const ::
_this->_internal_mutable_expire_after()->::greptime::v1::ExpireAfter::MergeFrom(
from._internal_expire_after());
}
if (from._internal_has_eval_interval()) {
_this->_internal_mutable_eval_interval()->::greptime::v1::EvalInterval::MergeFrom(
from._internal_eval_interval());
}
if (from._internal_create_if_not_exists() != 0) {
_this->_internal_set_create_if_not_exists(from._internal_create_if_not_exists());
}
+105
View File
@@ -1619,6 +1619,7 @@ class CreateRequest final :
kFlowIdFieldNumber = 1,
kSinkTableNameFieldNumber = 3,
kExpireAfterFieldNumber = 5,
kEvalIntervalFieldNumber = 10,
kCreateIfNotExistsFieldNumber = 4,
kOrReplaceFieldNumber = 9,
};
@@ -1739,6 +1740,24 @@ class CreateRequest final :
::greptime::v1::ExpireAfter* expire_after);
::greptime::v1::ExpireAfter* unsafe_arena_release_expire_after();
// .greptime.v1.EvalInterval eval_interval = 10;
bool has_eval_interval() const;
private:
bool _internal_has_eval_interval() const;
public:
void clear_eval_interval();
const ::greptime::v1::EvalInterval& eval_interval() const;
PROTOBUF_NODISCARD ::greptime::v1::EvalInterval* release_eval_interval();
::greptime::v1::EvalInterval* mutable_eval_interval();
void set_allocated_eval_interval(::greptime::v1::EvalInterval* eval_interval);
private:
const ::greptime::v1::EvalInterval& _internal_eval_interval() const;
::greptime::v1::EvalInterval* _internal_mutable_eval_interval();
public:
void unsafe_arena_set_allocated_eval_interval(
::greptime::v1::EvalInterval* eval_interval);
::greptime::v1::EvalInterval* unsafe_arena_release_eval_interval();
// bool create_if_not_exists = 4;
void clear_create_if_not_exists();
bool create_if_not_exists() const;
@@ -1776,6 +1795,7 @@ class CreateRequest final :
::greptime::v1::FlowId* flow_id_;
::greptime::v1::TableName* sink_table_name_;
::greptime::v1::ExpireAfter* expire_after_;
::greptime::v1::EvalInterval* eval_interval_;
bool create_if_not_exists_;
bool or_replace_;
mutable ::PROTOBUF_NAMESPACE_ID::internal::CachedSize _cached_size_;
@@ -3312,6 +3332,91 @@ inline void CreateRequest::set_allocated_expire_after(::greptime::v1::ExpireAfte
// @@protoc_insertion_point(field_set_allocated:greptime.v1.flow.CreateRequest.expire_after)
}
// .greptime.v1.EvalInterval eval_interval = 10;
inline bool CreateRequest::_internal_has_eval_interval() const {
return this != internal_default_instance() && _impl_.eval_interval_ != nullptr;
}
inline bool CreateRequest::has_eval_interval() const {
return _internal_has_eval_interval();
}
inline const ::greptime::v1::EvalInterval& CreateRequest::_internal_eval_interval() const {
const ::greptime::v1::EvalInterval* p = _impl_.eval_interval_;
return p != nullptr ? *p : reinterpret_cast<const ::greptime::v1::EvalInterval&>(
::greptime::v1::_EvalInterval_default_instance_);
}
inline const ::greptime::v1::EvalInterval& CreateRequest::eval_interval() const {
// @@protoc_insertion_point(field_get:greptime.v1.flow.CreateRequest.eval_interval)
return _internal_eval_interval();
}
inline void CreateRequest::unsafe_arena_set_allocated_eval_interval(
::greptime::v1::EvalInterval* eval_interval) {
if (GetArenaForAllocation() == nullptr) {
delete reinterpret_cast<::PROTOBUF_NAMESPACE_ID::MessageLite*>(_impl_.eval_interval_);
}
_impl_.eval_interval_ = eval_interval;
if (eval_interval) {
} else {
}
// @@protoc_insertion_point(field_unsafe_arena_set_allocated:greptime.v1.flow.CreateRequest.eval_interval)
}
inline ::greptime::v1::EvalInterval* CreateRequest::release_eval_interval() {
::greptime::v1::EvalInterval* temp = _impl_.eval_interval_;
_impl_.eval_interval_ = nullptr;
#ifdef PROTOBUF_FORCE_COPY_IN_RELEASE
auto* old = reinterpret_cast<::PROTOBUF_NAMESPACE_ID::MessageLite*>(temp);
temp = ::PROTOBUF_NAMESPACE_ID::internal::DuplicateIfNonNull(temp);
if (GetArenaForAllocation() == nullptr) { delete old; }
#else // PROTOBUF_FORCE_COPY_IN_RELEASE
if (GetArenaForAllocation() != nullptr) {
temp = ::PROTOBUF_NAMESPACE_ID::internal::DuplicateIfNonNull(temp);
}
#endif // !PROTOBUF_FORCE_COPY_IN_RELEASE
return temp;
}
inline ::greptime::v1::EvalInterval* CreateRequest::unsafe_arena_release_eval_interval() {
// @@protoc_insertion_point(field_release:greptime.v1.flow.CreateRequest.eval_interval)
::greptime::v1::EvalInterval* temp = _impl_.eval_interval_;
_impl_.eval_interval_ = nullptr;
return temp;
}
inline ::greptime::v1::EvalInterval* CreateRequest::_internal_mutable_eval_interval() {
if (_impl_.eval_interval_ == nullptr) {
auto* p = CreateMaybeMessage<::greptime::v1::EvalInterval>(GetArenaForAllocation());
_impl_.eval_interval_ = p;
}
return _impl_.eval_interval_;
}
inline ::greptime::v1::EvalInterval* CreateRequest::mutable_eval_interval() {
::greptime::v1::EvalInterval* _msg = _internal_mutable_eval_interval();
// @@protoc_insertion_point(field_mutable:greptime.v1.flow.CreateRequest.eval_interval)
return _msg;
}
inline void CreateRequest::set_allocated_eval_interval(::greptime::v1::EvalInterval* eval_interval) {
::PROTOBUF_NAMESPACE_ID::Arena* message_arena = GetArenaForAllocation();
if (message_arena == nullptr) {
delete reinterpret_cast< ::PROTOBUF_NAMESPACE_ID::MessageLite*>(_impl_.eval_interval_);
}
if (eval_interval) {
::PROTOBUF_NAMESPACE_ID::Arena* submessage_arena =
::PROTOBUF_NAMESPACE_ID::Arena::InternalGetOwningArena(
reinterpret_cast<::PROTOBUF_NAMESPACE_ID::MessageLite*>(eval_interval));
if (message_arena != submessage_arena) {
eval_interval = ::PROTOBUF_NAMESPACE_ID::internal::GetOwnedMessage(
message_arena, eval_interval, submessage_arena);
}
} else {
}
_impl_.eval_interval_ = eval_interval;
// @@protoc_insertion_point(field_set_allocated:greptime.v1.flow.CreateRequest.eval_interval)
}
// string comment = 6;
inline void CreateRequest::clear_comment() {
_impl_.comment_.ClearToEmpty();