• Home
  • Line#
  • Scopes#
  • Navigate#
  • Raw
  • Download
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