1 /*
2 * Copyright (C) 2022 The Android Open Source Project
3 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16
17 #include "src/cloud_trace_processor/worker_impl.h"
18
19 #include <memory>
20
21 #include "perfetto/base/status.h"
22 #include "perfetto/ext/base/status_or.h"
23 #include "perfetto/ext/base/threading/stream.h"
24 #include "perfetto/ext/base/uuid.h"
25 #include "protos/perfetto/cloud_trace_processor/common.pb.h"
26 #include "protos/perfetto/cloud_trace_processor/orchestrator.pb.h"
27 #include "protos/perfetto/cloud_trace_processor/worker.pb.h"
28 #include "src/cloud_trace_processor/trace_processor_wrapper.h"
29 #include "src/trace_processor/util/status_macros.h"
30
31 namespace perfetto {
32 namespace cloud_trace_processor {
33
34 Worker::~Worker() = default;
35
CreateInProcesss(CtpEnvironment * environment,base::ThreadPool * pool)36 std::unique_ptr<Worker> Worker::CreateInProcesss(CtpEnvironment* environment,
37 base::ThreadPool* pool) {
38 return std::make_unique<WorkerImpl>(environment, pool);
39 }
40
WorkerImpl(CtpEnvironment * environment,base::ThreadPool * pool)41 WorkerImpl::WorkerImpl(CtpEnvironment* environment, base::ThreadPool* pool)
42 : environment_(environment), thread_pool_(pool) {}
43
44 base::StatusOrFuture<protos::TracePoolShardCreateResponse>
TracePoolShardCreate(const protos::TracePoolShardCreateArgs & args)45 WorkerImpl::TracePoolShardCreate(const protos::TracePoolShardCreateArgs& args) {
46 if (args.pool_type() == protos::TracePoolType::DEDICATED) {
47 return base::ErrStatus("Dedicated pools are not currently supported");
48 }
49 auto it_and_inserted = shards_.Insert(args.pool_id(), TracePoolShard());
50 if (!it_and_inserted.second) {
51 return base::ErrStatus("Shard for pool %s already exists",
52 args.pool_id().c_str());
53 }
54 return base::StatusOr(protos::TracePoolShardCreateResponse());
55 }
56
57 base::StatusOrStream<protos::TracePoolShardSetTracesResponse>
TracePoolShardSetTraces(const protos::TracePoolShardSetTracesArgs & args)58 WorkerImpl::TracePoolShardSetTraces(
59 const protos::TracePoolShardSetTracesArgs& args) {
60 using Response = protos::TracePoolShardSetTracesResponse;
61 using StatusOrResponse = base::StatusOr<Response>;
62
63 TracePoolShard* shard = shards_.Find(args.pool_id());
64 if (!shard) {
65 return base::StreamOf<StatusOrResponse>(base::ErrStatus(
66 "Unable to find shard for pool %s", args.pool_id().c_str()));
67 }
68
69 std::vector<base::StatusOrStream<Response>> streams;
70 for (const std::string& trace : args.traces()) {
71 // TODO(lalitm): add support for stateful trace processor in dedicated
72 // pools.
73 auto tp = std::make_unique<TraceProcessorWrapper>(
74 trace, thread_pool_, TraceProcessorWrapper::Statefulness::kStateless);
75 auto load_trace_future =
76 tp->LoadTrace(environment_->ReadFile(trace))
77 .ContinueWith(
78 [trace](base::Status status) -> base::Future<StatusOrResponse> {
79 RETURN_IF_ERROR(status);
80 protos::TracePoolShardSetTracesResponse resp;
81 *resp.mutable_trace() = trace;
82 return resp;
83 });
84 streams.emplace_back(base::StreamFromFuture(std::move(load_trace_future)));
85 shard->tps.emplace_back(std::move(tp));
86 }
87 return base::FlattenStreams(std::move(streams));
88 }
89
90 base::StatusOrStream<protos::TracePoolShardQueryResponse>
TracePoolShardQuery(const protos::TracePoolShardQueryArgs & args)91 WorkerImpl::TracePoolShardQuery(const protos::TracePoolShardQueryArgs& args) {
92 using Response = protos::TracePoolShardQueryResponse;
93 using StatusOrResponse = base::StatusOr<Response>;
94 TracePoolShard* shard = shards_.Find(args.pool_id());
95 if (!shard) {
96 return base::StreamOf<StatusOrResponse>(base::ErrStatus(
97 "Unable to find shard for pool %s", args.pool_id().c_str()));
98 }
99 std::vector<base::StatusOrStream<Response>> streams;
100 streams.reserve(shard->tps.size());
101 for (std::unique_ptr<TraceProcessorWrapper>& tp : shard->tps) {
102 streams.emplace_back(tp->Query(args.sql_query()));
103 }
104 return base::FlattenStreams(std::move(streams));
105 }
106
107 base::StatusOrFuture<protos::TracePoolShardDestroyResponse>
TracePoolShardDestroy(const protos::TracePoolShardDestroyArgs & args)108 WorkerImpl::TracePoolShardDestroy(
109 const protos::TracePoolShardDestroyArgs& args) {
110 if (!shards_.Erase(args.pool_id())) {
111 return base::ErrStatus("Unable to find shard for pool %s",
112 args.pool_id().c_str());
113 }
114 return base::StatusOr(protos::TracePoolShardDestroyResponse());
115 }
116
117 } // namespace cloud_trace_processor
118 } // namespace perfetto
119