move protobuf from GreptimeDB (#1)
* move protobuf from GreptimeDB * add license check * add other github actions
This commit is contained in:
@@ -0,0 +1,10 @@
|
|||||||
|
{
|
||||||
|
"LABEL": {
|
||||||
|
"name": "Invalid PR Title",
|
||||||
|
"color": "B60205"
|
||||||
|
},
|
||||||
|
"CHECKS": {
|
||||||
|
"regexp": "^(feat|fix|test|refactor|chore|style|docs|perf|build|ci|revert)(\\(.*\\))?:.*",
|
||||||
|
"ignoreLabels" : ["ignore-title"]
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,19 @@
|
|||||||
|
I hereby agree to the terms of the [GreptimeDB CLA](https://gist.github.com/xtang/6378857777706e568c1949c7578592cc)
|
||||||
|
|
||||||
|
## What's changed and what's your intention?
|
||||||
|
|
||||||
|
_PLEASE DO NOT LEAVE THIS EMPTY !!!_
|
||||||
|
|
||||||
|
Please explain IN DETAIL what the changes are in this PR and why they are needed:
|
||||||
|
|
||||||
|
- Summarize your change (**mandatory**)
|
||||||
|
- How does this PR work? Need a brief introduction for the changed logic (optional)
|
||||||
|
- Describe clearly one logical change and avoid lazy messages (optional)
|
||||||
|
- Describe any limitations of the current code (optional)
|
||||||
|
|
||||||
|
## Checklist
|
||||||
|
|
||||||
|
- [ ] I have written the necessary comments.
|
||||||
|
- [ ] I have added the necessary unit tests and integration tests.
|
||||||
|
|
||||||
|
## Refer to a related PR or issue link (optional)
|
||||||
@@ -0,0 +1,16 @@
|
|||||||
|
name: License checker
|
||||||
|
|
||||||
|
on:
|
||||||
|
push:
|
||||||
|
branches:
|
||||||
|
- main
|
||||||
|
pull_request:
|
||||||
|
types: [opened, synchronize, reopened, ready_for_review]
|
||||||
|
jobs:
|
||||||
|
license-header-check:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
name: license-header-check
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v2
|
||||||
|
- name: Check License Header
|
||||||
|
uses: apache/skywalking-eyes/header@main
|
||||||
@@ -0,0 +1,76 @@
|
|||||||
|
on:
|
||||||
|
pull_request:
|
||||||
|
types: [opened, synchronize, reopened, ready_for_review]
|
||||||
|
push:
|
||||||
|
branches:
|
||||||
|
- main
|
||||||
|
workflow_dispatch:
|
||||||
|
|
||||||
|
name: CI
|
||||||
|
|
||||||
|
env:
|
||||||
|
RUST_TOOLCHAIN: stable
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
typos:
|
||||||
|
name: Spell Check with Typos
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v2
|
||||||
|
- uses: crate-ci/typos@v1.0.4
|
||||||
|
|
||||||
|
check:
|
||||||
|
name: Check
|
||||||
|
if: github.event.pull_request.draft == false
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 60
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v3
|
||||||
|
- uses: arduino/setup-protoc@v1
|
||||||
|
with:
|
||||||
|
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||||
|
- uses: dtolnay/rust-toolchain@master
|
||||||
|
with:
|
||||||
|
toolchain: ${{ env.RUST_TOOLCHAIN }}
|
||||||
|
- name: Rust Cache
|
||||||
|
uses: Swatinem/rust-cache@v2
|
||||||
|
- name: Run cargo check
|
||||||
|
run: cargo check --workspace --all-targets
|
||||||
|
|
||||||
|
fmt:
|
||||||
|
name: Rustfmt
|
||||||
|
if: github.event.pull_request.draft == false
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 60
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v3
|
||||||
|
- uses: arduino/setup-protoc@v1
|
||||||
|
with:
|
||||||
|
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||||
|
- uses: dtolnay/rust-toolchain@master
|
||||||
|
with:
|
||||||
|
toolchain: ${{ env.RUST_TOOLCHAIN }}
|
||||||
|
components: rustfmt
|
||||||
|
- name: Rust Cache
|
||||||
|
uses: Swatinem/rust-cache@v2
|
||||||
|
- name: Run cargo fmt
|
||||||
|
run: cargo fmt --all -- --check
|
||||||
|
|
||||||
|
clippy:
|
||||||
|
name: Clippy
|
||||||
|
if: github.event.pull_request.draft == false
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 60
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v3
|
||||||
|
- uses: arduino/setup-protoc@v1
|
||||||
|
with:
|
||||||
|
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||||
|
- uses: dtolnay/rust-toolchain@master
|
||||||
|
with:
|
||||||
|
toolchain: ${{ env.RUST_TOOLCHAIN }}
|
||||||
|
components: clippy
|
||||||
|
- name: Rust Cache
|
||||||
|
uses: Swatinem/rust-cache@v2
|
||||||
|
- name: Run cargo clippy
|
||||||
|
run: cargo clippy --workspace --all-targets -- -D warnings -D clippy::print_stdout -D clippy::print_stderr
|
||||||
@@ -0,0 +1,20 @@
|
|||||||
|
name: "PR Title Checker"
|
||||||
|
on:
|
||||||
|
pull_request_target:
|
||||||
|
types:
|
||||||
|
- opened
|
||||||
|
- edited
|
||||||
|
- synchronize
|
||||||
|
- labeled
|
||||||
|
- unlabeled
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
check:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 10
|
||||||
|
steps:
|
||||||
|
- uses: thehanimo/pr-title-checker@v1.3.4
|
||||||
|
with:
|
||||||
|
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||||
|
pass_on_octokit_error: false
|
||||||
|
configuration_path: ".github/pr-title-checker-config.json"
|
||||||
@@ -0,0 +1,14 @@
|
|||||||
|
header:
|
||||||
|
license:
|
||||||
|
spdx-id: Apache-2.0
|
||||||
|
copyright-owner: Greptime Team
|
||||||
|
|
||||||
|
paths:
|
||||||
|
- "**/*.rs"
|
||||||
|
- "**/*.py"
|
||||||
|
|
||||||
|
comment: on-failure
|
||||||
|
|
||||||
|
dependency:
|
||||||
|
files:
|
||||||
|
- Cargo.toml
|
||||||
+12
@@ -0,0 +1,12 @@
|
|||||||
|
[package]
|
||||||
|
name = "greptime-proto"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2021"
|
||||||
|
license = "Apache-2.0"
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
prost = "0.11"
|
||||||
|
tonic = "0.8"
|
||||||
|
|
||||||
|
[build-dependencies]
|
||||||
|
tonic-build = "0.8"
|
||||||
@@ -0,0 +1,30 @@
|
|||||||
|
// Copyright 2023 Greptime Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
fn main() {
|
||||||
|
tonic_build::configure()
|
||||||
|
.compile(
|
||||||
|
&[
|
||||||
|
"proto/greptime/v1/database.proto",
|
||||||
|
"proto/greptime/v1/meta/common.proto",
|
||||||
|
"proto/greptime/v1/meta/heartbeat.proto",
|
||||||
|
"proto/greptime/v1/meta/route.proto",
|
||||||
|
"proto/greptime/v1/meta/store.proto",
|
||||||
|
"proto/greptime/v1/meta/cluster.proto",
|
||||||
|
"proto/prometheus/remote/remote.proto",
|
||||||
|
],
|
||||||
|
&["proto"],
|
||||||
|
)
|
||||||
|
.expect("compile proto");
|
||||||
|
}
|
||||||
@@ -0,0 +1,85 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1;
|
||||||
|
|
||||||
|
message Column {
|
||||||
|
string column_name = 1;
|
||||||
|
|
||||||
|
enum SemanticType {
|
||||||
|
TAG = 0;
|
||||||
|
FIELD = 1;
|
||||||
|
TIMESTAMP = 2;
|
||||||
|
}
|
||||||
|
SemanticType semantic_type = 2;
|
||||||
|
|
||||||
|
message Values {
|
||||||
|
repeated int32 i8_values = 1;
|
||||||
|
repeated int32 i16_values = 2;
|
||||||
|
repeated int32 i32_values = 3;
|
||||||
|
repeated int64 i64_values = 4;
|
||||||
|
|
||||||
|
repeated uint32 u8_values = 5;
|
||||||
|
repeated uint32 u16_values = 6;
|
||||||
|
repeated uint32 u32_values = 7;
|
||||||
|
repeated uint64 u64_values = 8;
|
||||||
|
|
||||||
|
repeated float f32_values = 9;
|
||||||
|
repeated double f64_values = 10;
|
||||||
|
|
||||||
|
repeated bool bool_values = 11;
|
||||||
|
repeated bytes binary_values = 12;
|
||||||
|
repeated string string_values = 13;
|
||||||
|
|
||||||
|
repeated int32 date_values = 14;
|
||||||
|
repeated int64 datetime_values = 15;
|
||||||
|
repeated int64 ts_second_values = 16;
|
||||||
|
repeated int64 ts_millisecond_values = 17;
|
||||||
|
repeated int64 ts_microsecond_values = 18;
|
||||||
|
repeated int64 ts_nanosecond_values = 19;
|
||||||
|
}
|
||||||
|
// The array of non-null values in this column.
|
||||||
|
//
|
||||||
|
// For example: suppose there is a column "foo" that contains some int32 values (1, 2, 3, 4, 5, null, 7, 8, 9, null);
|
||||||
|
// column:
|
||||||
|
// column_name: foo
|
||||||
|
// semantic_type: Tag
|
||||||
|
// values: 1, 2, 3, 4, 5, 7, 8, 9
|
||||||
|
// null_masks: 00100000 00000010
|
||||||
|
Values values = 3;
|
||||||
|
|
||||||
|
// Mask maps the positions of null values.
|
||||||
|
// If a bit in null_mask is 1, it indicates that the column value at that position is null.
|
||||||
|
bytes null_mask = 4;
|
||||||
|
|
||||||
|
// Helpful in creating vector from column.
|
||||||
|
ColumnDataType datatype = 5;
|
||||||
|
}
|
||||||
|
|
||||||
|
message ColumnDef {
|
||||||
|
string name = 1;
|
||||||
|
ColumnDataType datatype = 2;
|
||||||
|
bool is_nullable = 3;
|
||||||
|
bytes default_constraint = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
enum ColumnDataType {
|
||||||
|
BOOLEAN = 0;
|
||||||
|
INT8 = 1;
|
||||||
|
INT16 = 2;
|
||||||
|
INT32 = 3;
|
||||||
|
INT64 = 4;
|
||||||
|
UINT8 = 5;
|
||||||
|
UINT16 = 6;
|
||||||
|
UINT32 = 7;
|
||||||
|
UINT64 = 8;
|
||||||
|
FLOAT32 = 9;
|
||||||
|
FLOAT64 = 10;
|
||||||
|
BINARY = 11;
|
||||||
|
STRING = 12;
|
||||||
|
DATE = 13;
|
||||||
|
DATETIME = 14;
|
||||||
|
TIMESTAMP_SECOND = 15;
|
||||||
|
TIMESTAMP_MILLISECOND = 16;
|
||||||
|
TIMESTAMP_MICROSECOND = 17;
|
||||||
|
TIMESTAMP_NANOSECOND = 18;
|
||||||
|
}
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1;
|
||||||
|
|
||||||
|
import "greptime/v1/ddl.proto";
|
||||||
|
import "greptime/v1/column.proto";
|
||||||
|
|
||||||
|
message RequestHeader {
|
||||||
|
// The `catalog` that is selected to be used in this request.
|
||||||
|
string catalog = 1;
|
||||||
|
// The `schema` that is selected to be used in this request.
|
||||||
|
string schema = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message GreptimeRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
oneof request {
|
||||||
|
InsertRequest insert = 2;
|
||||||
|
QueryRequest query = 3;
|
||||||
|
DdlRequest ddl = 4;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
message QueryRequest {
|
||||||
|
oneof query {
|
||||||
|
string sql = 1;
|
||||||
|
bytes logical_plan = 2;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
message InsertRequest {
|
||||||
|
string table_name = 1;
|
||||||
|
|
||||||
|
// Data is represented here.
|
||||||
|
repeated Column columns = 3;
|
||||||
|
|
||||||
|
// The row_count of all columns, which include null and non-null values.
|
||||||
|
//
|
||||||
|
// Note: the row_count of all columns in a InsertRequest must be same.
|
||||||
|
uint32 row_count = 4;
|
||||||
|
|
||||||
|
// The region number of current insert request.
|
||||||
|
uint32 region_number = 5;
|
||||||
|
}
|
||||||
|
|
||||||
|
message AffectedRows {
|
||||||
|
uint32 value = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message FlightMetadata {
|
||||||
|
AffectedRows affected_rows = 1;
|
||||||
|
}
|
||||||
@@ -0,0 +1,79 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1;
|
||||||
|
|
||||||
|
import "greptime/v1/column.proto";
|
||||||
|
|
||||||
|
// "Data Definition Language" requests, that create, modify or delete the database structures but not the data.
|
||||||
|
// `DdlRequest` could carry more information than plain SQL, for example, the "table_id" in `CreateTableExpr`.
|
||||||
|
// So create a new DDL expr if you need it.
|
||||||
|
message DdlRequest {
|
||||||
|
oneof expr {
|
||||||
|
CreateDatabaseExpr create_database = 1;
|
||||||
|
CreateTableExpr create_table = 2;
|
||||||
|
AlterExpr alter = 3;
|
||||||
|
DropTableExpr drop_table = 4;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
message CreateTableExpr {
|
||||||
|
string catalog_name = 1;
|
||||||
|
string schema_name = 2;
|
||||||
|
string table_name = 3;
|
||||||
|
string desc = 4;
|
||||||
|
repeated ColumnDef column_defs = 5;
|
||||||
|
string time_index = 6;
|
||||||
|
repeated string primary_keys = 7;
|
||||||
|
bool create_if_not_exists = 8;
|
||||||
|
map<string, string> table_options = 9;
|
||||||
|
TableId table_id = 10;
|
||||||
|
repeated uint32 region_ids = 11;
|
||||||
|
}
|
||||||
|
|
||||||
|
message AlterExpr {
|
||||||
|
string catalog_name = 1;
|
||||||
|
string schema_name = 2;
|
||||||
|
string table_name = 3;
|
||||||
|
oneof kind {
|
||||||
|
AddColumns add_columns = 4;
|
||||||
|
DropColumns drop_columns = 5;
|
||||||
|
RenameTable rename_table = 6;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
message DropTableExpr {
|
||||||
|
string catalog_name = 1;
|
||||||
|
string schema_name = 2;
|
||||||
|
string table_name = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message CreateDatabaseExpr {
|
||||||
|
//TODO(hl): maybe rename to schema_name?
|
||||||
|
string database_name = 1;
|
||||||
|
bool create_if_not_exists = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message AddColumns {
|
||||||
|
repeated AddColumn add_columns = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message DropColumns {
|
||||||
|
repeated DropColumn drop_columns = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message RenameTable {
|
||||||
|
string new_table_name = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message AddColumn {
|
||||||
|
ColumnDef column_def = 1;
|
||||||
|
bool is_key = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message DropColumn {
|
||||||
|
string name = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message TableId {
|
||||||
|
uint32 id = 1;
|
||||||
|
}
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1.meta;
|
||||||
|
|
||||||
|
import "greptime/v1/meta/common.proto";
|
||||||
|
import "greptime/v1/meta/store.proto";
|
||||||
|
|
||||||
|
// Cluster service is used for communication between meta nodes.
|
||||||
|
service Cluster {
|
||||||
|
// Batch get kvs by input keys from leader's in_memory kv store.
|
||||||
|
rpc BatchGet(BatchGetRequest) returns (BatchGetResponse);
|
||||||
|
|
||||||
|
// Range get the kvs from leader's in_memory kv store.
|
||||||
|
rpc Range(RangeRequest) returns (RangeResponse);
|
||||||
|
}
|
||||||
|
|
||||||
|
message BatchGetRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
repeated bytes keys = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message BatchGetResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
repeated KeyValue kvs = 2;
|
||||||
|
}
|
||||||
@@ -0,0 +1,48 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1.meta;
|
||||||
|
|
||||||
|
message RequestHeader {
|
||||||
|
uint64 protocol_version = 1;
|
||||||
|
// cluster_id is the ID of the cluster which be sent to.
|
||||||
|
uint64 cluster_id = 2;
|
||||||
|
// member_id is the ID of the sender server.
|
||||||
|
uint64 member_id = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message ResponseHeader {
|
||||||
|
uint64 protocol_version = 1;
|
||||||
|
// cluster_id is the ID of the cluster which sent the response.
|
||||||
|
uint64 cluster_id = 2;
|
||||||
|
Error error = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Error {
|
||||||
|
int32 code = 1;
|
||||||
|
string err_msg = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Peer {
|
||||||
|
uint64 id = 1;
|
||||||
|
string addr = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message TableName {
|
||||||
|
string catalog_name = 1;
|
||||||
|
string schema_name = 2;
|
||||||
|
string table_name = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message TimeInterval {
|
||||||
|
// The unix timestamp in millis of the start of this period.
|
||||||
|
uint64 start_timestamp_millis = 1;
|
||||||
|
// The unix timestamp in millis of the end of this period.
|
||||||
|
uint64 end_timestamp_millis = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message KeyValue {
|
||||||
|
// key is the key in bytes. An empty key is not allowed.
|
||||||
|
bytes key = 1;
|
||||||
|
// value is the value held by the key, in bytes.
|
||||||
|
bytes value = 2;
|
||||||
|
}
|
||||||
@@ -0,0 +1,92 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1.meta;
|
||||||
|
|
||||||
|
import "greptime/v1/meta/common.proto";
|
||||||
|
|
||||||
|
service Heartbeat {
|
||||||
|
// Heartbeat, there may be many contents of the heartbeat, such as:
|
||||||
|
// 1. Metadata to be registered to meta server and discoverable by other nodes.
|
||||||
|
// 2. Some performance metrics, such as Load, CPU usage, etc.
|
||||||
|
// 3. The number of computing tasks being executed.
|
||||||
|
rpc Heartbeat(stream HeartbeatRequest) returns (stream HeartbeatResponse) {}
|
||||||
|
|
||||||
|
// Ask leader's endpoint.
|
||||||
|
rpc AskLeader(AskLeaderRequest) returns (AskLeaderResponse) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
message HeartbeatRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
// Self peer
|
||||||
|
Peer peer = 2;
|
||||||
|
// Leader node
|
||||||
|
bool is_leader = 3;
|
||||||
|
// Actually reported time interval
|
||||||
|
TimeInterval report_interval = 4;
|
||||||
|
// Node stat
|
||||||
|
NodeStat node_stat = 5;
|
||||||
|
// Region stats on this node
|
||||||
|
repeated RegionStat region_stats = 6;
|
||||||
|
// Follower nodes and stats, empty on follower nodes
|
||||||
|
repeated ReplicaStat replica_stats = 7;
|
||||||
|
}
|
||||||
|
|
||||||
|
message NodeStat {
|
||||||
|
// The read capacity units during this period
|
||||||
|
int64 rcus = 1;
|
||||||
|
// The write capacity units during this period
|
||||||
|
int64 wcus = 2;
|
||||||
|
// How many tables on this node
|
||||||
|
int64 table_num = 3;
|
||||||
|
// How many regions on this node
|
||||||
|
int64 region_num = 4;
|
||||||
|
|
||||||
|
double cpu_usage = 5;
|
||||||
|
double load = 6;
|
||||||
|
// Read disk IO on this node
|
||||||
|
double read_io_rate = 7;
|
||||||
|
// Write disk IO on this node
|
||||||
|
double write_io_rate = 8;
|
||||||
|
|
||||||
|
// Others
|
||||||
|
map<string, string> attrs = 100;
|
||||||
|
}
|
||||||
|
|
||||||
|
message RegionStat {
|
||||||
|
uint64 region_id = 1;
|
||||||
|
TableName table_name = 2;
|
||||||
|
// The read capacity units during this period
|
||||||
|
int64 rcus = 3;
|
||||||
|
// The write capacity units during this period
|
||||||
|
int64 wcus = 4;
|
||||||
|
// Approximate bytes of this region
|
||||||
|
int64 approximate_bytes = 5;
|
||||||
|
// Approximate number of rows in this region
|
||||||
|
int64 approximate_rows = 6;
|
||||||
|
|
||||||
|
// Others
|
||||||
|
map<string, string> attrs = 100;
|
||||||
|
}
|
||||||
|
|
||||||
|
message ReplicaStat {
|
||||||
|
Peer peer = 1;
|
||||||
|
bool in_sync = 2;
|
||||||
|
bool is_learner = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message HeartbeatResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
repeated bytes payload = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message AskLeaderRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message AskLeaderResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
Peer leader = 2;
|
||||||
|
}
|
||||||
@@ -0,0 +1,99 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1.meta;
|
||||||
|
|
||||||
|
import "greptime/v1/meta/common.proto";
|
||||||
|
|
||||||
|
service Router {
|
||||||
|
rpc Create(CreateRequest) returns (RouteResponse) {}
|
||||||
|
|
||||||
|
// Fetch routing information for tables. The smallest unit is the complete
|
||||||
|
// routing information(all regions) of a table.
|
||||||
|
//
|
||||||
|
// ```text
|
||||||
|
// table_1
|
||||||
|
// table_name
|
||||||
|
// table_schema
|
||||||
|
// regions
|
||||||
|
// region_1
|
||||||
|
// leader_peer
|
||||||
|
// follower_peer_1, follower_peer_2
|
||||||
|
// region_2
|
||||||
|
// leader_peer
|
||||||
|
// follower_peer_1, follower_peer_2, follower_peer_3
|
||||||
|
// region_xxx
|
||||||
|
// table_2
|
||||||
|
// ...
|
||||||
|
// ```
|
||||||
|
//
|
||||||
|
rpc Route(RouteRequest) returns (RouteResponse) {}
|
||||||
|
|
||||||
|
rpc Delete(DeleteRequest) returns (RouteResponse) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
message CreateRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
TableName table_name = 2;
|
||||||
|
repeated Partition partitions = 3;
|
||||||
|
bytes table_info = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
message RouteRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
repeated TableName table_names = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message DeleteRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
TableName table_name = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message RouteResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
repeated Peer peers = 2;
|
||||||
|
repeated TableRoute table_routes = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message TableRoute {
|
||||||
|
Table table = 1;
|
||||||
|
repeated RegionRoute region_routes = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message RegionRoute {
|
||||||
|
Region region = 1;
|
||||||
|
// single leader node for write task
|
||||||
|
uint64 leader_peer_index = 2;
|
||||||
|
// multiple follower nodes for read task
|
||||||
|
repeated uint64 follower_peer_indexes = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Table {
|
||||||
|
uint64 id = 1;
|
||||||
|
TableName table_name = 2;
|
||||||
|
bytes table_schema = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Region {
|
||||||
|
// TODO(LFC): Maybe use message RegionNumber?
|
||||||
|
uint64 id = 1;
|
||||||
|
string name = 2;
|
||||||
|
Partition partition = 3;
|
||||||
|
|
||||||
|
map<string, string> attrs = 100;
|
||||||
|
}
|
||||||
|
|
||||||
|
// PARTITION `region_name` VALUES LESS THAN (value_list)
|
||||||
|
message Partition {
|
||||||
|
repeated bytes column_list = 1;
|
||||||
|
repeated bytes value_list = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
// This message is only for saving into store.
|
||||||
|
message TableRouteValue {
|
||||||
|
repeated Peer peers = 1;
|
||||||
|
TableRoute table_route = 2;
|
||||||
|
}
|
||||||
@@ -0,0 +1,159 @@
|
|||||||
|
syntax = "proto3";
|
||||||
|
|
||||||
|
package greptime.v1.meta;
|
||||||
|
|
||||||
|
import "greptime/v1/meta/common.proto";
|
||||||
|
|
||||||
|
service Store {
|
||||||
|
// Range gets the keys in the range from the key-value store.
|
||||||
|
rpc Range(RangeRequest) returns (RangeResponse);
|
||||||
|
|
||||||
|
// Put puts the given key into the key-value store.
|
||||||
|
rpc Put(PutRequest) returns (PutResponse);
|
||||||
|
|
||||||
|
// BatchPut atomically puts the given keys into the key-value store.
|
||||||
|
rpc BatchPut(BatchPutRequest) returns (BatchPutResponse);
|
||||||
|
|
||||||
|
// CompareAndPut atomically puts the value to the given updated
|
||||||
|
// value if the current value == the expected value.
|
||||||
|
rpc CompareAndPut(CompareAndPutRequest) returns (CompareAndPutResponse);
|
||||||
|
|
||||||
|
// DeleteRange deletes the given range from the key-value store.
|
||||||
|
rpc DeleteRange(DeleteRangeRequest) returns (DeleteRangeResponse);
|
||||||
|
|
||||||
|
// MoveValue atomically renames the key to the given updated key.
|
||||||
|
rpc MoveValue(MoveValueRequest) returns (MoveValueResponse);
|
||||||
|
}
|
||||||
|
|
||||||
|
message RangeRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
// key is the first key for the range, If range_end is not given, the
|
||||||
|
// request only looks up key.
|
||||||
|
bytes key = 2;
|
||||||
|
// range_end is the upper bound on the requested range [key, range_end).
|
||||||
|
// If range_end is '\0', the range is all keys >= key.
|
||||||
|
// If range_end is key plus one (e.g., "aa"+1 == "ab", "a\xff"+1 == "b"),
|
||||||
|
// then the range request gets all keys prefixed with key.
|
||||||
|
// If both key and range_end are '\0', then the range request returns all
|
||||||
|
// keys.
|
||||||
|
bytes range_end = 3;
|
||||||
|
// limit is a limit on the number of keys returned for the request. When
|
||||||
|
// limit is set to 0, it is treated as no limit.
|
||||||
|
int64 limit = 4;
|
||||||
|
// keys_only when set returns only the keys and not the values.
|
||||||
|
bool keys_only = 5;
|
||||||
|
}
|
||||||
|
|
||||||
|
message RangeResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
// kvs is the list of key-value pairs matched by the range request.
|
||||||
|
repeated KeyValue kvs = 2;
|
||||||
|
// more indicates if there are more keys to return in the requested range.
|
||||||
|
bool more = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message PutRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
// key is the key, in bytes, to put into the key-value store.
|
||||||
|
bytes key = 2;
|
||||||
|
// value is the value, in bytes, to associate with the key in the
|
||||||
|
// key-value store.
|
||||||
|
bytes value = 3;
|
||||||
|
// If prev_kv is set, gets the previous key-value pair before changing it.
|
||||||
|
// The previous key-value pair will be returned in the put response.
|
||||||
|
bool prev_kv = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
message PutResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
// If prev_kv is set in the request, the previous key-value pair will be
|
||||||
|
// returned.
|
||||||
|
KeyValue prev_kv = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message BatchPutRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
repeated KeyValue kvs = 2;
|
||||||
|
// If prev_kv is set, gets the previous key-value pairs before changing it.
|
||||||
|
// The previous key-value pairs will be returned in the batch put response.
|
||||||
|
bool prev_kv = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message BatchPutResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
// If prev_kv is set in the request, the previous key-value pairs will be
|
||||||
|
// returned.
|
||||||
|
repeated KeyValue prev_kvs = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message CompareAndPutRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
// key is the key, in bytes, to put into the key-value store.
|
||||||
|
bytes key = 2;
|
||||||
|
// expect is the previous value, in bytes
|
||||||
|
bytes expect = 3;
|
||||||
|
// value is the value, in bytes, to associate with the key in the
|
||||||
|
// key-value store.
|
||||||
|
bytes value = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
message CompareAndPutResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
bool success = 2;
|
||||||
|
KeyValue prev_kv = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message DeleteRangeRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
// key is the first key to delete in the range.
|
||||||
|
bytes key = 2;
|
||||||
|
// range_end is the key following the last key to delete for the range
|
||||||
|
// [key, range_end).
|
||||||
|
// If range_end is not given, the range is defined to contain only the key
|
||||||
|
// argument.
|
||||||
|
// If range_end is one bit larger than the given key, then the range is all
|
||||||
|
// the keys with the prefix (the given key).
|
||||||
|
// If range_end is '\0', the range is all keys greater than or equal to the
|
||||||
|
// key argument.
|
||||||
|
bytes range_end = 3;
|
||||||
|
// If prev_kv is set, gets the previous key-value pairs before deleting it.
|
||||||
|
// The previous key-value pairs will be returned in the delete response.
|
||||||
|
bool prev_kv = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
message DeleteRangeResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
// deleted is the number of keys deleted by the delete range request.
|
||||||
|
int64 deleted = 2;
|
||||||
|
// If prev_kv is set in the request, the previous key-value pairs will be
|
||||||
|
// returned.
|
||||||
|
repeated KeyValue prev_kvs = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message MoveValueRequest {
|
||||||
|
RequestHeader header = 1;
|
||||||
|
|
||||||
|
// If from_key dose not exist, return the value of to_key (if it exists).
|
||||||
|
// If from_key exists, move the value of from_key to to_key (i.e. rename),
|
||||||
|
// and return the value.
|
||||||
|
bytes from_key = 2;
|
||||||
|
bytes to_key = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message MoveValueResponse {
|
||||||
|
ResponseHeader header = 1;
|
||||||
|
|
||||||
|
// If from_key dose not exist, return the value of to_key (if it exists).
|
||||||
|
// If from_key exists, return the value of from_key.
|
||||||
|
KeyValue kv = 2;
|
||||||
|
}
|
||||||
@@ -0,0 +1,85 @@
|
|||||||
|
// Copyright 2016 Prometheus Team
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
syntax = "proto3";
|
||||||
|
package prometheus;
|
||||||
|
|
||||||
|
option go_package = "prompb";
|
||||||
|
|
||||||
|
import "prometheus/remote/types.proto";
|
||||||
|
|
||||||
|
message WriteRequest {
|
||||||
|
repeated prometheus.TimeSeries timeseries = 1;
|
||||||
|
// Cortex uses this field to determine the source of the write request.
|
||||||
|
// We reserve it to avoid any compatibility issues.
|
||||||
|
reserved 2;
|
||||||
|
repeated prometheus.MetricMetadata metadata = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
// ReadRequest represents a remote read request.
|
||||||
|
message ReadRequest {
|
||||||
|
repeated Query queries = 1;
|
||||||
|
|
||||||
|
enum ResponseType {
|
||||||
|
// Server will return a single ReadResponse message with matched series that includes list of raw samples.
|
||||||
|
// It's recommended to use streamed response types instead.
|
||||||
|
//
|
||||||
|
// Response headers:
|
||||||
|
// Content-Type: "application/x-protobuf"
|
||||||
|
// Content-Encoding: "snappy"
|
||||||
|
SAMPLES = 0;
|
||||||
|
// Server will stream a delimited ChunkedReadResponse message that contains XOR encoded chunks for a single series.
|
||||||
|
// Each message is following varint size and fixed size bigendian uint32 for CRC32 Castagnoli checksum.
|
||||||
|
//
|
||||||
|
// Response headers:
|
||||||
|
// Content-Type: "application/x-streamed-protobuf; proto=prometheus.ChunkedReadResponse"
|
||||||
|
// Content-Encoding: ""
|
||||||
|
STREAMED_XOR_CHUNKS = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
// accepted_response_types allows negotiating the content type of the response.
|
||||||
|
//
|
||||||
|
// Response types are taken from the list in the FIFO order. If no response type in `accepted_response_types` is
|
||||||
|
// implemented by server, error is returned.
|
||||||
|
// For request that do not contain `accepted_response_types` field the SAMPLES response type will be used.
|
||||||
|
repeated ResponseType accepted_response_types = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
// ReadResponse is a response when response_type equals SAMPLES.
|
||||||
|
message ReadResponse {
|
||||||
|
// In same order as the request's queries.
|
||||||
|
repeated QueryResult results = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Query {
|
||||||
|
int64 start_timestamp_ms = 1;
|
||||||
|
int64 end_timestamp_ms = 2;
|
||||||
|
repeated prometheus.LabelMatcher matchers = 3;
|
||||||
|
prometheus.ReadHints hints = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
message QueryResult {
|
||||||
|
// Samples within a time series must be ordered by time.
|
||||||
|
repeated prometheus.TimeSeries timeseries = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
// ChunkedReadResponse is a response when response_type equals STREAMED_XOR_CHUNKS.
|
||||||
|
// We strictly stream full series after series, optionally split by time. This means that a single frame can contain
|
||||||
|
// partition of the single series, but once a new series is started to be streamed it means that no more chunks will
|
||||||
|
// be sent for previous one. Series are returned sorted in the same way TSDB block are internally.
|
||||||
|
message ChunkedReadResponse {
|
||||||
|
repeated prometheus.ChunkedSeries chunked_series = 1;
|
||||||
|
|
||||||
|
// query_index represents an index of the query from ReadRequest.queries these chunks relates to.
|
||||||
|
int64 query_index = 2;
|
||||||
|
}
|
||||||
@@ -0,0 +1,117 @@
|
|||||||
|
// Copyright 2017 Prometheus Team
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
syntax = "proto3";
|
||||||
|
package prometheus;
|
||||||
|
|
||||||
|
option go_package = "prompb";
|
||||||
|
|
||||||
|
message MetricMetadata {
|
||||||
|
enum MetricType {
|
||||||
|
UNKNOWN = 0;
|
||||||
|
COUNTER = 1;
|
||||||
|
GAUGE = 2;
|
||||||
|
HISTOGRAM = 3;
|
||||||
|
GAUGEHISTOGRAM = 4;
|
||||||
|
SUMMARY = 5;
|
||||||
|
INFO = 6;
|
||||||
|
STATESET = 7;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Represents the metric type, these match the set from Prometheus.
|
||||||
|
// Refer to model/textparse/interface.go for details.
|
||||||
|
MetricType type = 1;
|
||||||
|
string metric_family_name = 2;
|
||||||
|
string help = 4;
|
||||||
|
string unit = 5;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Sample {
|
||||||
|
double value = 1;
|
||||||
|
// timestamp is in ms format, see model/timestamp/timestamp.go for
|
||||||
|
// conversion from time.Time to Prometheus timestamp.
|
||||||
|
int64 timestamp = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Exemplar {
|
||||||
|
// Optional, can be empty.
|
||||||
|
repeated Label labels = 1;
|
||||||
|
double value = 2;
|
||||||
|
// timestamp is in ms format, see model/timestamp/timestamp.go for
|
||||||
|
// conversion from time.Time to Prometheus timestamp.
|
||||||
|
int64 timestamp = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
// TimeSeries represents samples and labels for a single time series.
|
||||||
|
message TimeSeries {
|
||||||
|
// For a timeseries to be valid, and for the samples and exemplars
|
||||||
|
// to be ingested by the remote system properly, the labels field is required.
|
||||||
|
repeated Label labels = 1;
|
||||||
|
repeated Sample samples = 2;
|
||||||
|
repeated Exemplar exemplars = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Label {
|
||||||
|
string name = 1;
|
||||||
|
string value = 2;
|
||||||
|
}
|
||||||
|
|
||||||
|
message Labels {
|
||||||
|
repeated Label labels = 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Matcher specifies a rule, which can match or set of labels or not.
|
||||||
|
message LabelMatcher {
|
||||||
|
enum Type {
|
||||||
|
EQ = 0;
|
||||||
|
NEQ = 1;
|
||||||
|
RE = 2;
|
||||||
|
NRE = 3;
|
||||||
|
}
|
||||||
|
Type type = 1;
|
||||||
|
string name = 2;
|
||||||
|
string value = 3;
|
||||||
|
}
|
||||||
|
|
||||||
|
message ReadHints {
|
||||||
|
int64 step_ms = 1; // Query step size in milliseconds.
|
||||||
|
string func = 2; // String representation of surrounding function or aggregation.
|
||||||
|
int64 start_ms = 3; // Start time in milliseconds.
|
||||||
|
int64 end_ms = 4; // End time in milliseconds.
|
||||||
|
repeated string grouping = 5; // List of label names used in aggregation.
|
||||||
|
bool by = 6; // Indicate whether it is without or by.
|
||||||
|
int64 range_ms = 7; // Range vector selector range in milliseconds.
|
||||||
|
}
|
||||||
|
|
||||||
|
// Chunk represents a TSDB chunk.
|
||||||
|
// Time range [min, max] is inclusive.
|
||||||
|
message Chunk {
|
||||||
|
int64 min_time_ms = 1;
|
||||||
|
int64 max_time_ms = 2;
|
||||||
|
|
||||||
|
// We require this to match chunkenc.Encoding.
|
||||||
|
enum Encoding {
|
||||||
|
UNKNOWN = 0;
|
||||||
|
XOR = 1;
|
||||||
|
}
|
||||||
|
Encoding type = 3;
|
||||||
|
bytes data = 4;
|
||||||
|
}
|
||||||
|
|
||||||
|
// ChunkedSeries represents single, encoded time series.
|
||||||
|
message ChunkedSeries {
|
||||||
|
// Labels should be sorted.
|
||||||
|
repeated Label labels = 1;
|
||||||
|
// Chunks will be in start time order and may overlap.
|
||||||
|
repeated Chunk chunks = 2;
|
||||||
|
}
|
||||||
+22
@@ -0,0 +1,22 @@
|
|||||||
|
// Copyright 2023 Greptime Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
mod serde;
|
||||||
|
pub mod v1;
|
||||||
|
|
||||||
|
pub mod prometheus {
|
||||||
|
pub mod remote {
|
||||||
|
tonic::include_proto!("prometheus");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,38 @@
|
|||||||
|
// Copyright 2023 Greptime Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
use prost::DecodeError;
|
||||||
|
use prost::Message;
|
||||||
|
|
||||||
|
use crate::v1::meta::TableRouteValue;
|
||||||
|
|
||||||
|
macro_rules! impl_convert_with_bytes {
|
||||||
|
($data_type: ty) => {
|
||||||
|
impl From<$data_type> for Vec<u8> {
|
||||||
|
fn from(entity: $data_type) -> Self {
|
||||||
|
entity.encode_to_vec()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TryFrom<&[u8]> for $data_type {
|
||||||
|
type Error = DecodeError;
|
||||||
|
|
||||||
|
fn try_from(value: &[u8]) -> Result<Self, Self::Error> {
|
||||||
|
<$data_type>::decode(value.as_ref())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
impl_convert_with_bytes!(TableRouteValue);
|
||||||
@@ -0,0 +1,17 @@
|
|||||||
|
// Copyright 2023 Greptime Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
tonic::include_proto!("greptime.v1");
|
||||||
|
|
||||||
|
pub mod meta;
|
||||||
+209
@@ -0,0 +1,209 @@
|
|||||||
|
// Copyright 2023 Greptime Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
tonic::include_proto!("greptime.v1.meta");
|
||||||
|
|
||||||
|
use std::collections::HashMap;
|
||||||
|
use std::hash::{Hash, Hasher};
|
||||||
|
|
||||||
|
pub const PROTOCOL_VERSION: u64 = 1;
|
||||||
|
|
||||||
|
#[derive(Default)]
|
||||||
|
pub struct PeerDict {
|
||||||
|
peers: HashMap<Peer, usize>,
|
||||||
|
index: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PeerDict {
|
||||||
|
pub fn get_or_insert(&mut self, peer: Peer) -> usize {
|
||||||
|
let index = self.peers.entry(peer).or_insert_with(|| {
|
||||||
|
let v = self.index;
|
||||||
|
self.index += 1;
|
||||||
|
v
|
||||||
|
});
|
||||||
|
|
||||||
|
*index
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn into_peers(self) -> Vec<Peer> {
|
||||||
|
let mut array = vec![Peer::default(); self.index];
|
||||||
|
for (p, i) in self.peers {
|
||||||
|
array[i] = p;
|
||||||
|
}
|
||||||
|
array
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[allow(clippy::derive_hash_xor_eq)]
|
||||||
|
impl Hash for Peer {
|
||||||
|
fn hash<H: Hasher>(&self, state: &mut H) {
|
||||||
|
self.id.hash(state);
|
||||||
|
self.addr.hash(state);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Eq for Peer {}
|
||||||
|
|
||||||
|
impl RequestHeader {
|
||||||
|
#[inline]
|
||||||
|
pub fn new((cluster_id, member_id): (u64, u64)) -> Self {
|
||||||
|
Self {
|
||||||
|
protocol_version: PROTOCOL_VERSION,
|
||||||
|
cluster_id,
|
||||||
|
member_id,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ResponseHeader {
|
||||||
|
#[inline]
|
||||||
|
pub fn success(cluster_id: u64) -> Self {
|
||||||
|
Self {
|
||||||
|
protocol_version: PROTOCOL_VERSION,
|
||||||
|
cluster_id,
|
||||||
|
..Default::default()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
pub fn failed(cluster_id: u64, error: Error) -> Self {
|
||||||
|
Self {
|
||||||
|
protocol_version: PROTOCOL_VERSION,
|
||||||
|
cluster_id,
|
||||||
|
error: Some(error),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
pub fn is_not_leader(&self) -> bool {
|
||||||
|
if let Some(error) = &self.error {
|
||||||
|
if error.code == ErrorCode::NotLeader as i32 {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum ErrorCode {
|
||||||
|
NoActiveDatanodes = 1,
|
||||||
|
NotLeader = 2,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Error {
|
||||||
|
#[inline]
|
||||||
|
pub fn no_active_datanodes() -> Self {
|
||||||
|
Self {
|
||||||
|
code: ErrorCode::NoActiveDatanodes as i32,
|
||||||
|
err_msg: "No active datanodes".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
pub fn is_not_leader() -> Self {
|
||||||
|
Self {
|
||||||
|
code: ErrorCode::NotLeader as i32,
|
||||||
|
err_msg: "Current server is not leader".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl HeartbeatResponse {
|
||||||
|
#[inline]
|
||||||
|
pub fn is_not_leader(&self) -> bool {
|
||||||
|
if let Some(header) = &self.header {
|
||||||
|
return header.is_not_leader();
|
||||||
|
}
|
||||||
|
false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
macro_rules! gen_set_header {
|
||||||
|
($req: ty) => {
|
||||||
|
impl $req {
|
||||||
|
#[inline]
|
||||||
|
pub fn set_header(&mut self, (cluster_id, member_id): (u64, u64)) {
|
||||||
|
self.header = Some(RequestHeader::new((cluster_id, member_id)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
gen_set_header!(HeartbeatRequest);
|
||||||
|
gen_set_header!(RouteRequest);
|
||||||
|
gen_set_header!(CreateRequest);
|
||||||
|
gen_set_header!(RangeRequest);
|
||||||
|
gen_set_header!(DeleteRequest);
|
||||||
|
gen_set_header!(PutRequest);
|
||||||
|
gen_set_header!(BatchPutRequest);
|
||||||
|
gen_set_header!(CompareAndPutRequest);
|
||||||
|
gen_set_header!(DeleteRangeRequest);
|
||||||
|
gen_set_header!(MoveValueRequest);
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use std::vec;
|
||||||
|
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_peer_dict() {
|
||||||
|
let mut dict = PeerDict::default();
|
||||||
|
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 1,
|
||||||
|
addr: "111".to_string(),
|
||||||
|
});
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 2,
|
||||||
|
addr: "222".to_string(),
|
||||||
|
});
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 1,
|
||||||
|
addr: "111".to_string(),
|
||||||
|
});
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 1,
|
||||||
|
addr: "111".to_string(),
|
||||||
|
});
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 1,
|
||||||
|
addr: "111".to_string(),
|
||||||
|
});
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 1,
|
||||||
|
addr: "111".to_string(),
|
||||||
|
});
|
||||||
|
dict.get_or_insert(Peer {
|
||||||
|
id: 2,
|
||||||
|
addr: "222".to_string(),
|
||||||
|
});
|
||||||
|
|
||||||
|
assert_eq!(2, dict.index);
|
||||||
|
assert_eq!(
|
||||||
|
vec![
|
||||||
|
Peer {
|
||||||
|
id: 1,
|
||||||
|
addr: "111".to_string(),
|
||||||
|
},
|
||||||
|
Peer {
|
||||||
|
id: 2,
|
||||||
|
addr: "222".to_string(),
|
||||||
|
}
|
||||||
|
],
|
||||||
|
dict.into_peers()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user