Skip to content

Commit a0beb5f

Browse files
committed
Add DumpMessages/AlterRole/UpdateUser interfaces
Signed-off-by: yhmo <yihua.mo@zilliz.com>
1 parent 673b0ea commit a0beb5f

34 files changed

Lines changed: 1357 additions & 17 deletions

cmake/MilvusProtoGen.cmake

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,15 +16,14 @@
1616

1717
include_guard(GLOBAL)
1818

19-
set(PROTO_VERSION v2.6.13)
20-
set(PROTO_URL https://github.com/milvus-io/milvus-proto/archive/refs/tags/${PROTO_VERSION}.tar.gz)
21-
19+
set(PROTO_COMMIT fd141a092113dd7a85ab6dc6c1b037b17a3bf96a)
2220

2321
include(FetchContent)
2422

2523
# download proto
2624
FetchContent_Declare(milvus_proto
27-
URL ${PROTO_URL}
25+
GIT_REPOSITORY https://github.com/milvus-io/milvus-proto.git
26+
GIT_TAG ${PROTO_COMMIT}
2827
)
2928
FetchContent_Populate(milvus_proto)
3029

examples/README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ Examples for MilvusClient:
3030

3131
Examples for MilvusClientV2:
3232
- `./cmake_build/examples/v2/sdk_array_v2`: example to show the usage of Array field.
33+
- `./cmake_build/examples/v2/sdk_cdc_v2`: example to show the usage of CDC interfaces including DumpMessages().
3334
- `./cmake_build/examples/v2/sdk_db_v2`: example to show the usage of databases.
3435
- `./cmake_build/examples/v2/sdk_default_value_v2`: example to show the usage of default value.
3536
- `./cmake_build/examples/v2/sdk_dml_v2`: example to show the usage of dml interfaces.

examples/src/v2/cdc.cpp

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
1+
// Licensed to the LF AI & Data foundation under one
2+
// or more contributor license agreements. See the NOTICE file
3+
// distributed with this work for additional information
4+
// regarding copyright ownership. The ASF licenses this file
5+
// to you under the Apache License, Version 2.0 (the
6+
// "License"); you may not use this file except in compliance
7+
// with the License. You may obtain a copy of the License at
8+
//
9+
// http://www.apache.org/licenses/LICENSE-2.0
10+
//
11+
// Unless required by applicable law or agreed to in writing, software
12+
// distributed under the License is distributed on an "AS IS" BASIS,
13+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
// See the License for the specific language governing permissions and
15+
// limitations under the License.
16+
17+
#include <cstdint>
18+
#include <iostream>
19+
#include <limits>
20+
#include <string>
21+
#include <vector>
22+
23+
#include "ExampleUtils.h"
24+
#include "milvus/MilvusClientV2.h"
25+
26+
namespace {
27+
28+
const std::string kClusterAUri = "http://localhost:19530";
29+
const std::string kClusterBUri = "http://localhost:29530";
30+
const std::string kClusterAId = "cdc-test-upstream";
31+
const std::string kClusterBId = "cdc-test-downstream";
32+
const int kPChannelNum = 16;
33+
34+
std::vector<std::string>
35+
GeneratePChannels(const std::string& cluster_id) {
36+
std::vector<std::string> pchannels;
37+
pchannels.reserve(kPChannelNum);
38+
for (int i = 0; i < kPChannelNum; ++i) {
39+
pchannels.push_back(cluster_id + "-rootcoord-dml_" + std::to_string(i));
40+
}
41+
return pchannels;
42+
}
43+
44+
void
45+
PrintDumpedMessage(const milvus::DumpedMessage& message) {
46+
std::cout << "message_id=" << message.MessageID().ID() << ", wal_name=" << message.MessageID().WalName()
47+
<< ", payload_size=" << message.Payload().size() << std::endl;
48+
if (!message.Properties().empty()) {
49+
std::cout << "properties=";
50+
util::PrintMap(message.Properties());
51+
}
52+
}
53+
54+
} // namespace
55+
56+
int
57+
main(int argc, char* argv[]) {
58+
printf("Example start...\n");
59+
60+
auto cluster_a_client = milvus::MilvusClientV2::Create();
61+
auto cluster_b_client = milvus::MilvusClientV2::Create();
62+
63+
auto status = cluster_a_client->Connect(milvus::ConnectParam{kClusterAUri});
64+
util::CheckStatus("connect cluster A", status);
65+
66+
status = cluster_b_client->Connect(milvus::ConnectParam{kClusterBUri});
67+
util::CheckStatus("connect cluster B", status);
68+
69+
milvus::MilvusCluster cluster_a;
70+
cluster_a.WithClusterID(kClusterAId).WithUri(kClusterAUri).WithPChannels(GeneratePChannels(kClusterAId));
71+
72+
milvus::MilvusCluster cluster_b;
73+
cluster_b.WithClusterID(kClusterBId).WithUri(kClusterBUri).WithPChannels(GeneratePChannels(kClusterBId));
74+
75+
milvus::CrossClusterTopology topology;
76+
topology.WithSourceClusterID(kClusterAId).WithTargetClusterID(kClusterBId);
77+
78+
milvus::ReplicateConfiguration configuration;
79+
configuration.AddCluster(std::move(cluster_a));
80+
configuration.AddCluster(std::move(cluster_b));
81+
configuration.AddCrossClusterTopology(std::move(topology));
82+
83+
milvus::UpdateReplicateConfigurationRequest update_request;
84+
update_request.WithConfiguration(std::move(configuration));
85+
86+
status = cluster_a_client->UpdateReplicateConfiguration(update_request);
87+
util::CheckStatus("update replicate configuration on cluster A", status);
88+
89+
status = cluster_b_client->UpdateReplicateConfiguration(update_request);
90+
util::CheckStatus("update replicate configuration on cluster B", status);
91+
92+
const std::string target_pchannel = GeneratePChannels(kClusterBId).front();
93+
94+
milvus::GetReplicateInfoRequest replicate_info_request;
95+
replicate_info_request.WithSourceClusterID(kClusterAId).WithTargetPChannel(target_pchannel);
96+
97+
milvus::GetReplicateInfoResponse replicate_info_response;
98+
status = cluster_b_client->GetReplicateInfo(replicate_info_request, replicate_info_response);
99+
util::CheckStatus("get replicate info from cluster B", status);
100+
101+
milvus::ReplicateMessageID start_message_id;
102+
start_message_id.WithID(replicate_info_response.Checkpoint().MessageID().ID())
103+
.WithWalName(replicate_info_response.Checkpoint().MessageID().WalName());
104+
105+
milvus::DumpMessagesRequest dump_request;
106+
dump_request.WithPChannel(target_pchannel)
107+
.WithStartMessageID(std::move(start_message_id))
108+
.WithStartTimeTick(replicate_info_response.Checkpoint().TimeTick())
109+
.WithEndTimeTick(std::numeric_limits<uint64_t>::max());
110+
111+
std::cout << "DumpMessages on pchannel: " << target_pchannel << std::endl;
112+
status = cluster_a_client->DumpMessages(dump_request, [](const milvus::DumpedMessage& message) {
113+
PrintDumpedMessage(message);
114+
return milvus::Status::OK();
115+
});
116+
util::CheckStatus("dump messages from cluster A", status);
117+
118+
cluster_a_client->Disconnect();
119+
cluster_b_client->Disconnect();
120+
return 0;
121+
}

src/impl/MilvusClientV2Impl.cpp

Lines changed: 85 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2507,6 +2507,41 @@ MilvusClientV2Impl::GetReplicateInfo(const GetReplicateInfoRequest& request, Get
25072507
pre, &MilvusConnection::GetReplicateInfo, post);
25082508
}
25092509

2510+
Status
2511+
MilvusClientV2Impl::DumpMessages(const DumpMessagesRequest& request,
2512+
const std::function<Status(const DumpedMessage&)>& on_message) {
2513+
if (!on_message) {
2514+
return {StatusCode::INVALID_ARGUMENT, "DumpMessages callback cannot be empty"};
2515+
}
2516+
2517+
auto connection = connection_.GetConnection();
2518+
if (connection == nullptr) {
2519+
return {StatusCode::NOT_CONNECTED, "Connection is not created!"};
2520+
}
2521+
2522+
proto::milvus::DumpMessagesRequest rpc_request;
2523+
rpc_request.set_pchannel(request.PChannel());
2524+
rpc_request.mutable_start_message_id()->set_id(request.StartMessageID().ID());
2525+
const auto& wal_name_str = request.StartMessageID().WalName();
2526+
if (!wal_name_str.empty()) {
2527+
proto::common::WALName wal_name;
2528+
if (!proto::common::WALName_Parse(wal_name_str, &wal_name)) {
2529+
return {StatusCode::INVALID_ARGUMENT, "Unknown WAL name: " + wal_name_str};
2530+
}
2531+
rpc_request.mutable_start_message_id()->set_wal_name(wal_name);
2532+
}
2533+
rpc_request.set_start_timetick(request.StartTimeTick());
2534+
rpc_request.set_end_timetick(request.EndTimeTick());
2535+
2536+
auto callback = [&on_message](const proto::common::ImmutableMessage& rpc_message) {
2537+
DumpedMessage message;
2538+
ConvertImmutableMessage(rpc_message, message);
2539+
return on_message(message);
2540+
};
2541+
2542+
return connection->DumpMessages(rpc_request, GrpcOpts{}, callback);
2543+
}
2544+
25102545
Status
25112546
MilvusClientV2Impl::CreateResourceGroup(const CreateResourceGroupRequest& request) {
25122547
auto pre = [&request](proto::milvus::CreateResourceGroupRequest& rpc_request) {
@@ -2656,6 +2691,18 @@ MilvusClientV2Impl::UpdatePassword(const UpdatePasswordRequest& request) {
26562691
pre, &MilvusConnection::UpdateCredential, nullptr);
26572692
}
26582693

2694+
Status
2695+
MilvusClientV2Impl::UpdateUser(const UpdateUserRequest& request) {
2696+
auto pre = [&request](proto::milvus::UpdateCredentialRequest& rpc_request) {
2697+
rpc_request.set_username(request.UserName());
2698+
rpc_request.set_description(request.Description());
2699+
return Status::OK();
2700+
};
2701+
2702+
return connection_.Invoke<proto::milvus::UpdateCredentialRequest, proto::common::Status>(
2703+
pre, &MilvusConnection::UpdateCredential, nullptr);
2704+
}
2705+
26592706
Status
26602707
MilvusClientV2Impl::DropUser(const DropUserRequest& request) {
26612708
auto pre = [&request](proto::milvus::DeleteCredentialRequest& rpc_request) {
@@ -2680,6 +2727,8 @@ MilvusClientV2Impl::DescribeUser(const DescribeUserRequest& request, DescribeUse
26802727
desc.SetName(request.UserName());
26812728
if (rpc_response.results().size() > 0) {
26822729
auto result = rpc_response.results().at(0);
2730+
desc.SetName(result.user().name());
2731+
desc.SetDescription(result.description());
26832732
for (const auto& role : result.roles()) {
26842733
desc.AddRole(role.name());
26852734
}
@@ -2719,6 +2768,18 @@ MilvusClientV2Impl::CreateRole(const CreateRoleRequest& request) {
27192768
pre, &MilvusConnection::CreateRole, nullptr);
27202769
}
27212770

2771+
Status
2772+
MilvusClientV2Impl::AlterRole(const AlterRoleRequest& request) {
2773+
auto pre = [&request](proto::milvus::AlterRoleRequest& rpc_request) {
2774+
rpc_request.set_role_name(request.RoleName());
2775+
rpc_request.set_description(request.Description());
2776+
return Status::OK();
2777+
};
2778+
2779+
return connection_.Invoke<proto::milvus::AlterRoleRequest, proto::common::Status>(pre, &MilvusConnection::AlterRole,
2780+
nullptr);
2781+
}
2782+
27222783
Status
27232784
MilvusClientV2Impl::DropRole(const DropRoleRequest& request) {
27242785
auto pre = [&request](proto::milvus::DropRoleRequest& rpc_request) {
@@ -2740,13 +2801,36 @@ MilvusClientV2Impl::DescribeRole(const DescribeRoleRequest& request, DescribeRol
27402801
return Status::OK();
27412802
};
27422803

2743-
auto post = [&request, &response](const proto::milvus::SelectGrantResponse& rpc_response) {
2804+
auto post = [this, &request, &response](const proto::milvus::SelectGrantResponse& rpc_response) {
27442805
RoleDesc desc;
27452806
desc.SetName(request.RoleName());
27462807
for (const auto& entity : rpc_response.entities()) {
27472808
desc.AddGrantItem({entity.object().name(), entity.object_name(), entity.db_name(), entity.role().name(),
27482809
entity.grantor().user().name(), entity.grantor().privilege().name()});
27492810
}
2811+
2812+
proto::milvus::SelectRoleRequest role_request;
2813+
role_request.set_include_user_info(false);
2814+
role_request.mutable_role()->set_name(request.RoleName());
2815+
proto::milvus::SelectRoleResponse role_response;
2816+
2817+
auto connection = connection_.GetConnection();
2818+
if (connection == nullptr) {
2819+
return Status{StatusCode::NOT_CONNECTED, "Connection is not created!"};
2820+
}
2821+
auto retry_param = connection_.GetRetryParam();
2822+
auto caller = [&]() {
2823+
return connection->SelectRole(role_request, role_response, GrpcOpts{connection_.GetRpcDeadlineMs()});
2824+
};
2825+
auto status = Retry(caller, retry_param);
2826+
if (!status.IsOk()) {
2827+
return status;
2828+
}
2829+
if (role_response.results_size() > 0) {
2830+
desc.SetName(role_response.results(0).role().name());
2831+
desc.SetDescription(role_response.results(0).role().description());
2832+
}
2833+
27502834
response.SetDesc(std::move(desc));
27512835
return Status::OK();
27522836
};

src/impl/MilvusClientV2Impl.h

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -268,6 +268,10 @@ class MilvusClientV2Impl : public MilvusClientV2, public std::enable_shared_from
268268
Status
269269
GetReplicateInfo(const GetReplicateInfoRequest& request, GetReplicateInfoResponse& response) final;
270270

271+
Status
272+
DumpMessages(const DumpMessagesRequest& request,
273+
const std::function<Status(const DumpedMessage&)>& on_message) final;
274+
271275
Status
272276
CreateResourceGroup(const CreateResourceGroupRequest& request) final;
273277

@@ -295,6 +299,9 @@ class MilvusClientV2Impl : public MilvusClientV2, public std::enable_shared_from
295299
Status
296300
UpdatePassword(const UpdatePasswordRequest& request) final;
297301

302+
Status
303+
UpdateUser(const UpdateUserRequest& request) final;
304+
298305
Status
299306
DropUser(const DropUserRequest& request) final;
300307

@@ -307,6 +314,9 @@ class MilvusClientV2Impl : public MilvusClientV2, public std::enable_shared_from
307314
Status
308315
CreateRole(const CreateRoleRequest& request) final;
309316

317+
Status
318+
AlterRole(const AlterRoleRequest& request) final;
319+
310320
Status
311321
DropRole(const DropRoleRequest& request) final;
312322

src/impl/MilvusConnection.cpp

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -683,6 +683,55 @@ MilvusConnection::GetReplicateInfo(const proto::milvus::GetReplicateInfoRequest&
683683
return grpcCall("GetReplicateInfo", &Stub::GetReplicateInfo, request, response, options);
684684
}
685685

686+
Status
687+
MilvusConnection::DumpMessages(const proto::milvus::DumpMessagesRequest& request, const GrpcContextOptions& options,
688+
const std::function<Status(const proto::common::ImmutableMessage&)>& on_message) {
689+
std::shared_ptr<proto::milvus::MilvusService::Stub> stub;
690+
{
691+
std::lock_guard<std::mutex> lock(stub_mtx_);
692+
stub = stub_;
693+
}
694+
if (stub == nullptr) {
695+
return {StatusCode::NOT_CONNECTED, "Connection is not ready!"};
696+
}
697+
698+
::grpc::ClientContext context;
699+
if (options.timeout > 0) {
700+
auto deadline = std::chrono::system_clock::now() + std::chrono::milliseconds{options.timeout};
701+
context.set_deadline(deadline);
702+
}
703+
704+
auto reader = stub->DumpMessages(&context, request);
705+
proto::milvus::DumpMessagesResponse response;
706+
while (reader->Read(&response)) {
707+
switch (response.response_case()) {
708+
case proto::milvus::DumpMessagesResponse::kStatus: {
709+
auto status = StatusByProtoResponse(response.status());
710+
if (!status.IsOk()) {
711+
return status;
712+
}
713+
break;
714+
}
715+
case proto::milvus::DumpMessagesResponse::kMessage: {
716+
auto status = on_message(response.message());
717+
if (!status.IsOk()) {
718+
return status;
719+
}
720+
break;
721+
}
722+
case proto::milvus::DumpMessagesResponse::RESPONSE_NOT_SET:
723+
default:
724+
return {StatusCode::SERVER_FAILED, "Unexpected empty dumpMessages response"};
725+
}
726+
}
727+
728+
auto grpc_status = reader->Finish();
729+
if (!grpc_status.ok()) {
730+
return StatusCodeFromGrpcStatus(grpc_status);
731+
}
732+
return Status::OK();
733+
}
734+
686735
Status
687736
MilvusConnection::SelectUser(const proto::milvus::SelectUserRequest& request,
688737
proto::milvus::SelectUserResponse& response, const GrpcContextOptions& options) {
@@ -707,6 +756,12 @@ MilvusConnection::CreateRole(const proto::milvus::CreateRoleRequest& request, pr
707756
return grpcCall("CreateRole", &Stub::CreateRole, request, response, options);
708757
}
709758

759+
Status
760+
MilvusConnection::AlterRole(const proto::milvus::AlterRoleRequest& request, proto::common::Status& response,
761+
const GrpcContextOptions& options) {
762+
return grpcCall("AlterRole", &Stub::AlterRole, request, response, options);
763+
}
764+
710765
Status
711766
MilvusConnection::DropRole(const proto::milvus::DropRoleRequest& request, proto::common::Status& response,
712767
const GrpcContextOptions& options) {

0 commit comments

Comments
 (0)