Skip to content
Closed
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
6 changes: 3 additions & 3 deletions Android.bp
Original file line number Diff line number Diff line change
Expand Up @@ -19005,15 +19005,15 @@ filegroup {
filegroup {
name: "perfetto_src_trace_processor_perfetto_sql_exec_exec",
srcs: [
"src/trace_processor/perfetto_sql/exec/sql_scan.cc",
"src/trace_processor/perfetto_sql/exec/collected_rows.cc",
],
}

// GN: //src/trace_processor/perfetto_sql/exec:unittests
filegroup {
name: "perfetto_src_trace_processor_perfetto_sql_exec_unittests",
srcs: [
"src/trace_processor/perfetto_sql/exec/sql_scan_unittest.cc",
"src/trace_processor/perfetto_sql/exec/collected_rows_unittest.cc",
],
}

Expand Down Expand Up @@ -19121,7 +19121,7 @@ filegroup {
filegroup {
name: "perfetto_src_trace_processor_perfetto_sql_schema_schema",
srcs: [
"src/trace_processor/perfetto_sql/schema/query_schema.cc",
"src/trace_processor/perfetto_sql/schema/sqlite_relations.cc",
],
}

Expand Down
8 changes: 4 additions & 4 deletions BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -3829,8 +3829,8 @@ perfetto_filegroup(
perfetto_filegroup(
name = "src_trace_processor_perfetto_sql_exec_exec",
srcs = [
"src/trace_processor/perfetto_sql/exec/sql_scan.cc",
"src/trace_processor/perfetto_sql/exec/sql_scan.h",
"src/trace_processor/perfetto_sql/exec/collected_rows.cc",
"src/trace_processor/perfetto_sql/exec/collected_rows.h",
],
)

Expand Down Expand Up @@ -3901,8 +3901,8 @@ perfetto_filegroup(
perfetto_filegroup(
name = "src_trace_processor_perfetto_sql_schema_schema",
srcs = [
"src/trace_processor/perfetto_sql/schema/query_schema.cc",
"src/trace_processor/perfetto_sql/schema/query_schema.h",
"src/trace_processor/perfetto_sql/schema/sqlite_relations.cc",
"src/trace_processor/perfetto_sql/schema/sqlite_relations.h",
"src/trace_processor/perfetto_sql/schema/type_mapping.h",
],
)
Expand Down
56 changes: 42 additions & 14 deletions src/perfetto_sql/analysis/relation.cc
Original file line number Diff line number Diff line change
Expand Up @@ -51,8 +51,14 @@ struct OwnedView {
uint32_t root = 0;
};

// The text SQLite sees for the name at `span`, after macro expansion. A name is
// one token, so its text is a slice of the one layer it is in, which lives as
// long as the statement.
std::string_view Text(SyntaqliteParser* p, SyntaqliteTextSpan span) {
return base::TrimWhitespace(SyntaqliteSpanText(p, span));
uint32_t len = 0;
const char* text = syntaqlite_parser_span_expanded_text(p, &span, &len);
return text ? base::TrimWhitespace(std::string_view(text, len))
: std::string_view();
}

const SyntaqliteNode* Node(SyntaqliteParser* p, uint32_t id) {
Expand Down Expand Up @@ -150,6 +156,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 @@ -185,15 +192,38 @@ class RelationAnalyzer::Impl {
static ColumnLineage Lookup(const Scope&,
std::string_view table,
std::string_view column);
// The columns of a leaf relation, kept for the rest of the analysis. Its
// hidden columns are added to `hidden`.
std::vector<ColumnLineage> LeafColumns(LeafRelation,
std::vector<std::string_view>& hidden);

const Catalog& catalog_;
// Lineage string_views point into each view's sql string and parse tree, so
// 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;
};

std::vector<ColumnLineage> RelationAnalyzer::Impl::LeafColumns(
LeafRelation found,
std::vector<std::string_view>& hidden) {
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) {
out.push_back({column.name, {{relation->name, column.name, column.type}}});
if (column.hidden) {
hidden.push_back(column.name);
}
}
return out;
}

ColumnLineage RelationAnalyzer::Impl::Lookup(const Scope& scope,
std::string_view table,
std::string_view column) {
Expand Down Expand Up @@ -436,26 +466,24 @@ RelationAnalyzer::Impl::Select(SyntaqliteParser* p, uint32_t id, int depth) {
}
return std::move(*left);
}
default:
return base::ErrStatus("relation analysis: not a select");
default: {
std::optional<LeafRelation> found = catalog_.FindNodeRelation({p, id});
if (!found) {
return base::ErrStatus("relation analysis: not a select");
}
preserves_rows_ = false;
std::vector<std::string_view> hidden;
return LeafColumns(std::move(*found), hidden);
}
}
}

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)) {
std::vector<ColumnLineage> out;
out.reserve(relation->columns.size());
for (const LeafColumn& column : relation->columns) {
out.push_back(
{column.name, {{relation->name, column.name, column.type}}});
if (column.hidden) {
hidden.push_back(column.name);
}
}
return out;
if (std::optional<LeafRelation> found = catalog_.FindLeafRelation(name)) {
return LeafColumns(std::move(*found), hidden);
}
if (depth >= kMaxDepth) {
return base::ErrStatus(
Expand Down
12 changes: 8 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 All @@ -72,6 +71,11 @@ class Catalog {
std::string_view name) const = 0;
virtual std::optional<std::string> FindViewSql(
std::string_view name) const = 0;
// The relation a node of the dialect's own stands for, such as a pipeline
// written as a subquery, when the host knows it.
virtual std::optional<LeafRelation> FindNodeRelation(SqlNode) const {
return std::nullopt;
}
};

struct ColumnOrigin {
Expand Down
17 changes: 8 additions & 9 deletions src/trace_processor/perfetto_sql/engine/connection_catalog.cc
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
#include "src/trace_processor/core/dataframe/dataframe.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/schema/query_schema.h"
#include "src/trace_processor/perfetto_sql/schema/sqlite_relations.h"
#include "src/trace_processor/perfetto_sql/schema/type_mapping.h"
#include "src/trace_processor/sqlite/bindings/sqlite_column.h"
#include "src/trace_processor/sqlite/sql_source.h"
Expand Down Expand Up @@ -77,10 +77,15 @@ std::optional<analysis::LeafRelation> ConnectionCatalog::FindLeafRelation(
std::string_view name) const {
const dataframe::Dataframe* dataframe = connection_->GetDataframeOrNull(name);
if (!dataframe) {
return std::nullopt;
// A view is expanded instead, so its columns are traced through it.
if (FindViewSql(name)) {
return std::nullopt;
}
return sql_schema::FindSqliteRelation(connection_->sqlite_connection(),
name);
}
analysis::LeafRelation relation;
relation.name = name;
relation.name = std::string(name);
const std::vector<std::string>& columns = dataframe->column_names();
relation.columns.reserve(columns.size());
for (uint32_t i = 0; i < columns.size(); ++i) {
Expand Down Expand Up @@ -109,10 +114,4 @@ const dataframe::Dataframe* ConnectionCatalog::FindDataframe(
return connection_->GetDataframeOrNull(name);
}

base::StatusOr<pipeline::Schema> ConnectionCatalog::DescribeQuery(
const SqlSource& sql) const {
return sql_schema::DescribeQuery(connection_->sqlite_connection(), sql,
*this);
}

} // namespace perfetto::trace_processor
5 changes: 1 addition & 4 deletions src/trace_processor/perfetto_sql/engine/connection_catalog.h
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 All @@ -44,8 +43,6 @@ class ConnectionCatalog final : public perfetto_sql::analysis::Catalog,

const dataframe::Dataframe* FindDataframe(
std::string_view name) const override;
base::StatusOr<pipeline::Schema> DescribeQuery(
const SqlSource& sql) const override;

private:
PerfettoSqlConnection* connection_;
Expand Down
15 changes: 11 additions & 4 deletions src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.cc
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,8 @@ PerfettoSqlConnection::PerfettoSqlConnection(
ctx->connection = this;
RegisterVirtualTableModule<PipelineModule>(pipeline::kPipelineFunction,
std::move(ctx));
base::Status status = RegisterAggregateFunction<CollectRows>(pool_);
PERFETTO_CHECK(status.ok());
}
database_->InitializeSharedSchema(connection_.get());

Expand Down Expand Up @@ -1057,18 +1059,23 @@ PerfettoSqlConnection::PreparePipeline(const pipeline::LogicalPlan& plan,
return base::ErrStatus("%s%s", source.AsTraceback(0).c_str(),
sql.status().c_message());
}
return connection_->PrepareStatement(source.RewriteAllIgnoreExisting(
SqlSource::FromTraceProcessorImplementation(std::move(*sql))));
SqliteConnection::PreparedStatement stmt =
connection_->PrepareStatement(source.RewriteAllIgnoreExisting(
SqlSource::FromTraceProcessorImplementation(std::move(*sql))));
RETURN_IF_ERROR(stmt.status());
return std::move(stmt);
}

base::StatusOr<std::unique_ptr<pipeline::PhysicalPlan>>
PerfettoSqlConnection::LoadPipeline(std::string_view serialized) {
PerfettoSqlConnection::LoadPipeline(
std::string_view serialized,
const exec::CollectedRowsScan::Inputs& inputs) {
PERFETTO_TP_TRACE(metatrace::Category::QUERY_TIMELINE, "PIPELINE_LOAD");
ASSIGN_OR_RETURN(pipeline::LogicalPlan plan,
pipeline::DeserializePlan(serialized, *catalog_));
pipeline::LowerEnvironment env;
env.connection = connection_.get();
env.pool = pool_;
env.inputs = &inputs;
return pipeline::Lower(plan, env);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -174,9 +174,11 @@ class PerfettoSqlConnection {
base::StatusOr<SqliteConnection::PreparedStatement> PrepareSqliteStatement(
SqlSource sql);

// Loads a plan written by pipeline::SerializePlan, ready to run.
// Loads a plan written by pipeline::SerializePlan, ready to run. Each run
// reads the rows of the plan's inputs from `inputs`, which must outlive it.
base::StatusOr<std::unique_ptr<pipeline::PhysicalPlan>> LoadPipeline(
std::string_view serialized);
std::string_view serialized,
const exec::CollectedRowsScan::Inputs& inputs);

// Registers a virtual table module with the given name.
//
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -856,7 +856,9 @@ class PerfettoSqlConnectionPipelineTest : public PerfettoSqlConnectionTest {
std::string PipelineSql(const std::string& pipeline) {
auto res = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(pipeline));
PERFETTO_CHECK(res.ok());
if (!res.ok()) {
PERFETTO_FATAL("%s", res.status().c_message());
}
return res->stmt.sql();
}

Expand Down Expand Up @@ -1176,17 +1178,25 @@ TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesInViewsAndFunctions) {
EXPECT_THAT(*rows, testing::ElementsAre("20"));
}

// A pipeline runs its SQL sources itself, so nothing binds a function's
// arguments there: reading one is refused when the pipeline is compiled, not
// silently NULL.
TEST_F(PerfettoSqlConnectionPipelineTest, PipelineSourcesCannotReadArguments) {
EXPECT_THAT(
Rows("CREATE PERFETTO FUNCTION scaled(k LONG) RETURNS TABLE(total LONG) "
"AS SELECT total FROM (FROM (SELECT id, parent_id, self * $k AS "
// A pipeline's SQL is evaluated where the pipeline is written, so it reads
// the arguments of the function it is in like any other SQL there.
TEST_F(PerfettoSqlConnectionPipelineTest, PipelineSourcesReadArguments) {
ASSERT_TRUE(
Rows("CREATE PERFETTO FUNCTION scaled(k LONG) "
"RETURNS TABLE(id LONG, total LONG) "
"AS SELECT id, total FROM (FROM (SELECT id, parent_id, self * $k AS "
"self FROM tree) |> TREE ACCUMULATE UP SUM(self) AS total)")
.status()
.message(),
testing::HasSubstr("Cannot read `$k`"));
.ok());
auto rows = Rows("SELECT id, total FROM scaled(2)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,200", "1,120", "2,60", "3,80"));

// Read again for each row of a join, with a different argument each time.
rows = Rows(
"SELECT k.v, s.total FROM (SELECT 1 AS v UNION ALL SELECT 2 "
"UNION ALL SELECT 1) k JOIN scaled(k.v) s WHERE s.id = 0");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::UnorderedElementsAre("1,100", "2,200", "1,100"));
}

// Macros and pipelines nest inside each other in every way: pipelines in
Expand Down Expand Up @@ -1239,9 +1249,9 @@ TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesAsIntersectionOperands) {
EXPECT_THAT(*rows, testing::ElementsAre("12,3", "20,2"));
}

// SQLite reads the inner side of a join once per outer row. A pipeline read
// again with the same plan is not run again each time: `random()` in its
// source would otherwise give each read different rows.
// SQLite reads the inner side of a join once per outer row. The relation a
// pipeline reads is collected once, so every read sees the same rows even when
// its SQL, like `random()`, would give different ones each time.
TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesReadAgainAreNotRunAgain) {
ASSERT_TRUE(Rows("CREATE TABLE picks(k); "
"INSERT INTO picks VALUES (1), (2), (3), (4), (5)")
Expand All @@ -1250,9 +1260,7 @@ TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesReadAgainAreNotRunAgain) {
"SELECT count(DISTINCT p.r) FROM picks CROSS JOIN "
"(FROM (SELECT 1 AS one, random() AS r)) p");
ASSERT_TRUE(rows.ok()) << rows.status().message();
// The second read runs the pipeline once more to keep its rows; every read
// after replays them.
EXPECT_THAT(*rows, testing::ElementsAre("2"));
EXPECT_THAT(*rows, testing::ElementsAre("1"));
}

TEST_F(PerfettoSqlConnectionPipelineTest, PipelineSubqueriesNeedPipelines) {
Expand Down
Loading
Loading