Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Android.bp
Original file line number Diff line number Diff line change
Expand Up @@ -19080,6 +19080,7 @@ filegroup {
"src/trace_processor/perfetto_sql/pipeline/catalog.cc",
"src/trace_processor/perfetto_sql/pipeline/column_pruning.cc",
"src/trace_processor/perfetto_sql/pipeline/compiler.cc",
"src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc",
"src/trace_processor/perfetto_sql/pipeline/plan_serialization.cc",
],
}
Expand Down
2 changes: 2 additions & 0 deletions BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -3879,6 +3879,8 @@ perfetto_filegroup(
"src/trace_processor/perfetto_sql/pipeline/compiler.cc",
"src/trace_processor/perfetto_sql/pipeline/compiler.h",
"src/trace_processor/perfetto_sql/pipeline/logical_plan.h",
"src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc",
"src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h",
"src/trace_processor/perfetto_sql/pipeline/plan_serialization.cc",
"src/trace_processor/perfetto_sql/pipeline/plan_serialization.h",
],
Expand Down
8 changes: 7 additions & 1 deletion src/perfetto_sql/analysis/relation.cc
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,7 @@ class RelationAnalyzer::Impl {
void Begin() {
preserves_rows_ = true;
views_.clear();
leaves_.clear();
}

// The columns of the relation `name`. Its hidden columns, if any, are added
Expand Down Expand Up @@ -191,6 +192,9 @@ class RelationAnalyzer::Impl {
// every OwnedView needs a stable address: growing a std::vector<OwnedView>
// would move the elements and moving `sql` can relocate its bytes (SSO).
std::vector<std::unique_ptr<OwnedView>> views_;
// Lineage string_views point into each leaf relation's strings, so they are
// kept at stable addresses for the same reason.
std::vector<std::unique_ptr<LeafRelation>> leaves_;
bool preserves_rows_ = true;
};

Expand Down Expand Up @@ -445,7 +449,9 @@ base::StatusOr<std::vector<ColumnLineage>> RelationAnalyzer::Impl::Relation(
std::string_view name,
int depth,
std::vector<std::string_view>& hidden) {
if (std::optional<LeafRelation> relation = catalog_.FindLeafRelation(name)) {
if (std::optional<LeafRelation> found = catalog_.FindLeafRelation(name)) {
leaves_.push_back(std::make_unique<LeafRelation>(std::move(*found)));
const LeafRelation* relation = leaves_.back().get();
std::vector<ColumnLineage> out;
out.reserve(relation->columns.size());
for (const LeafColumn& column : relation->columns) {
Expand Down
7 changes: 3 additions & 4 deletions src/perfetto_sql/analysis/relation.h
Original file line number Diff line number Diff line change
Expand Up @@ -50,20 +50,19 @@ using ColumnType = base::TypeSet<Id, Uint32, Int32, Int64, Double, String>;

// A leaf relation whose columns can be used as lineage origins.
struct LeafColumn {
std::string_view name;
std::string name;
// Nothing when the catalog does not know how the column is stored.
std::optional<ColumnType> type;
// Left out of `*` and `table.*`, as SQLite does for HIDDEN columns, but
// still found by name.
bool hidden = false;
};
struct LeafRelation {
std::string_view name;
std::string name;
std::vector<LeafColumn> columns;
};

// Supplies the schema objects referenced by parsed queries. Returned leaf
// strings only need to remain valid for the duration of an Analyze call.
// Supplies the schema objects referenced by parsed queries.
class Catalog {
public:
virtual ~Catalog();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,7 @@ namespace perfetto::trace_processor {

// Adapts a connection to semantic analysis and pipeline compilation. Each
// dataframe is served as a typed leaf relation and as a dataframe.
class ConnectionCatalog final : public perfetto_sql::analysis::Catalog,
public pipeline::Catalog {
class ConnectionCatalog final : public pipeline::Catalog {
public:
explicit ConnectionCatalog(PerfettoSqlConnection*);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@
#include "src/trace_processor/perfetto_sql/parser/perfetto_sql_parser.h"
#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h"
#include "src/trace_processor/perfetto_sql/pipeline/physical_plan.h"
#include "src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h"
#include "src/trace_processor/perfetto_sql/pipeline/plan_serialization.h"
#include "src/trace_processor/sqlite/bindings/sqlite_column.h"
#include "src/trace_processor/sqlite/bindings/sqlite_type.h"
Expand Down Expand Up @@ -361,7 +362,7 @@ PerfettoSqlConnection::PerfettoSqlConnection(
auto ctx = std::make_unique<PipelineModule::Context>();
ctx->pool = pool_;
ctx->connection = this;
RegisterVirtualTableModule<PipelineModule>(PipelineModule::kName,
RegisterVirtualTableModule<PipelineModule>(pipeline::kPipelineFunction,
std::move(ctx));
}
database_->InitializeSharedSchema(connection_.get());
Expand Down Expand Up @@ -1051,7 +1052,7 @@ base::Status PerfettoSqlConnection::ExecuteCreateTable(
base::StatusOr<SqliteConnection::PreparedStatement>
PerfettoSqlConnection::PreparePipeline(const pipeline::LogicalPlan& plan,
const SqlSource& source) {
auto sql = PipelineModule::SelectFrom(plan);
auto sql = pipeline::SelectPipelineSql(plan);
if (!sql.ok()) {
return base::ErrStatus("%s%s", source.AsTraceback(0).c_str(),
sql.status().c_message());
Expand Down
21 changes: 2 additions & 19 deletions src/trace_processor/perfetto_sql/engine/pipeline_module.cc
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,8 @@
#include "src/trace_processor/core/exec/variant.h"
#include "src/trace_processor/core/util/bit_vector.h"
#include "src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.h"
#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h"
#include "src/trace_processor/perfetto_sql/pipeline/physical_plan.h"
#include "src/trace_processor/perfetto_sql/pipeline/plan_serialization.h"
#include "src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h"
#include "src/trace_processor/sqlite/bindings/sqlite_result.h"
#include "src/trace_processor/sqlite/sqlite_utils.h"

Expand All @@ -56,7 +55,7 @@ constexpr int kFirstOutputColumn = 1;
std::string Schema() {
// Public names (which may repeat) are applied by the outer SELECT.
std::vector<std::string> columns{"pipeline HIDDEN"};
for (uint32_t i = 0; i < PipelineModule::kMaxColumns; ++i) {
for (uint32_t i = 0; i < pipeline::kMaxPipelineColumns; ++i) {
columns.push_back("c" + std::to_string(i));
}
return "CREATE TABLE x(" + base::Join(columns, ", ") + ")";
Expand Down Expand Up @@ -167,22 +166,6 @@ int CheckStatus(PipelineModule::Cursor* cursor) {

} // namespace

base::StatusOr<std::string> PipelineModule::SelectFrom(
const pipeline::LogicalPlan& plan) {
const std::vector<pipeline::NamedColumn>& output = plan.output;
if (output.size() > kMaxColumns) {
return base::ErrStatus("A pipeline can output at most %u columns, not %zu",
kMaxColumns, output.size());
}
std::vector<std::string> columns;
for (uint32_t i = 0; i < output.size(); ++i) {
columns.push_back("c" + std::to_string(i) + " AS \"" +
base::ReplaceAll(output[i].name, "\"", "\"\"") + "\"");
}
return "SELECT " + base::Join(columns, ", ") + " FROM " + kName + "(X'" +
base::ToHex(pipeline::SerializePlan(plan)) + "')";
}

int PipelineModule::Connect(sqlite3* db,
void* raw_ctx,
int,
Expand Down
7 changes: 0 additions & 7 deletions src/trace_processor/perfetto_sql/engine/pipeline_module.h
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,8 @@
#include <string>
#include <vector>

#include "perfetto/ext/base/status_or.h"
#include "src/trace_processor/containers/string_pool.h"
#include "src/trace_processor/core/exec/row_cursor.h"
#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h"
#include "src/trace_processor/perfetto_sql/pipeline/physical_plan.h"
#include "src/trace_processor/sqlite/bindings/sqlite_module.h"

Expand All @@ -44,8 +42,6 @@ struct PipelineModule : sqlite::Module<PipelineModule> {
static constexpr auto kType = kEponymousOnly;
static constexpr bool kSupportsWrites = false;
static constexpr bool kDoesOverloadFunctions = false;
static constexpr char kName[] = "__intrinsic_pipeline";
static constexpr uint32_t kMaxColumns = 256;

struct Context {
StringPool* pool;
Expand Down Expand Up @@ -76,9 +72,6 @@ struct PipelineModule : sqlite::Module<PipelineModule> {
std::optional<int64_t> target_rowid;
};

// SQL reading `plan`'s output under its own column names.
static base::StatusOr<std::string> SelectFrom(const pipeline::LogicalPlan&);

static int Connect(sqlite3*,
void*,
int,
Expand Down
3 changes: 3 additions & 0 deletions src/trace_processor/perfetto_sql/pipeline/BUILD.gn
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,15 @@ source_set("logical") {
"compiler.cc",
"compiler.h",
"logical_plan.h",
"pipeline_sql.cc",
"pipeline_sql.h",
"plan_serialization.cc",
"plan_serialization.h",
]
deps = [
"../../../../gn:default_deps",
"../../../base",
"../../../perfetto_sql/analysis",
"../../../perfetto_sql/syntaqlite",
"../../core/common",
"../../core/dataframe",
Expand Down
9 changes: 6 additions & 3 deletions src/trace_processor/perfetto_sql/pipeline/catalog.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,16 +20,19 @@
#include <string_view>

#include "perfetto/ext/base/status_or.h"
#include "src/perfetto_sql/analysis/relation.h"
#include "src/trace_processor/core/dataframe/dataframe.h"
#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h"
#include "src/trace_processor/sqlite/sql_source.h"

namespace perfetto::trace_processor::pipeline {

// Lookup interface the compiler uses to resolve what a pipeline reads.
class Catalog {
// Lookup interface the compiler uses to resolve what a pipeline reads: the
// relations semantic analysis can describe, and the dataframes a pipeline can
// read directly.
class Catalog : public perfetto_sql::analysis::Catalog {
public:
virtual ~Catalog();
~Catalog() override;

// Dataframe registered as `name`, or null.
virtual const dataframe::Dataframe* FindDataframe(
Expand Down
46 changes: 46 additions & 0 deletions src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
/*
* Copyright (C) 2026 The Android Open Source Project
*
* 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.
*/

#include "src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h"

#include <cstdint>
#include <string>
#include <vector>

#include "perfetto/base/status.h"
#include "perfetto/ext/base/status_or.h"
#include "perfetto/ext/base/string_utils.h"
#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h"
#include "src/trace_processor/perfetto_sql/pipeline/plan_serialization.h"

namespace perfetto::trace_processor::pipeline {

base::StatusOr<std::string> SelectPipelineSql(const LogicalPlan& plan) {
if (plan.output.size() > kMaxPipelineColumns) {
return base::ErrStatus("A pipeline can output at most %u columns, not %zu",
kMaxPipelineColumns, plan.output.size());
}
std::vector<std::string> columns;
for (uint32_t i = 0; i < plan.output.size(); ++i) {
columns.push_back("c" + std::to_string(i) + " AS \"" +
base::ReplaceAll(plan.output[i].name, "\"", "\"\"") +
"\"");
}
return "SELECT " + base::Join(columns, ", ") + " FROM " + kPipelineFunction +
"(X'" + base::ToHex(SerializePlan(plan)) + "')";
}

} // namespace perfetto::trace_processor::pipeline
40 changes: 40 additions & 0 deletions src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
/*
* Copyright (C) 2026 The Android Open Source Project
*
* 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.
*/

#ifndef SRC_TRACE_PROCESSOR_PERFETTO_SQL_PIPELINE_PIPELINE_SQL_H_
#define SRC_TRACE_PROCESSOR_PERFETTO_SQL_PIPELINE_PIPELINE_SQL_H_

#include <cstdint>
#include <string>

#include "perfetto/ext/base/status_or.h"
#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h"

namespace perfetto::trace_processor::pipeline {

// The table function which runs a serialized plan.
inline constexpr char kPipelineFunction[] = "__intrinsic_pipeline";

// The most columns the table function can output.
inline constexpr uint32_t kMaxPipelineColumns = 256;

// SQL reading `plan`'s output under its own column names. The plan is
// serialized into the SQL, so the SQL needs nothing else to run.
base::StatusOr<std::string> SelectPipelineSql(const LogicalPlan& plan);

} // namespace perfetto::trace_processor::pipeline

#endif // SRC_TRACE_PROCESSOR_PERFETTO_SQL_PIPELINE_PIPELINE_SQL_H_
2 changes: 1 addition & 1 deletion src/trace_processor/perfetto_sql/pipeline/test_catalog.h
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ namespace perfetto::trace_processor::pipeline {

// Catalog over dataframes built by the test. SQLite does not know about them.
// Anything else is described via `connection` (if given) and is untyped.
class TestCatalog : public Catalog, public perfetto_sql::analysis::Catalog {
class TestCatalog : public Catalog {
public:
explicit TestCatalog(StringPool* pool, SqliteConnection* connection = nullptr)
: pool_(pool), connection_(connection) {}
Expand Down
Loading