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
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
39 changes: 39 additions & 0 deletions src/perfetto_sql/syntaqlite/perfetto.y
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,12 @@ perfetto_pipe_source(A) ::= LP select(S) RP as(Z). {
A = synq_parse_perfetto_pipe_source(pCtx, SYNQ_NO_SPAN, SYNQ_NO_SPAN,
S, Z.name, Z.has_as ? SYNTAQLITE_BOOL_TRUE : SYNTAQLITE_BOOL_FALSE);
}
// A pipeline read by a pipeline is compiled into the same plan, so it is left
// as it is written rather than expanded into SQL.
perfetto_pipe_source(A) ::= LP perfetto_pipeline(P) RP as(Z). {
A = synq_parse_perfetto_pipe_source(pCtx, SYNQ_NO_SPAN, SYNQ_NO_SPAN,
P, Z.name, Z.has_as ? SYNTAQLITE_BOOL_TRUE : SYNTAQLITE_BOOL_FALSE);
}

%type perfetto_tree_direction {int}
perfetto_tree_direction(A) ::= UP. { A = SYNTAQLITE_PERFETTO_TREE_DIRECTION_UP; }
Expand Down Expand Up @@ -424,6 +430,39 @@ perfetto_pipeline(A) ::= INTERVAL INTERSECTION OF LP

cmd(A) ::= perfetto_pipeline(P). { A = P; }

// A pipeline in parentheses reads like any other subquery, in a FROM clause or
// as a CTE. Once parsed, it is handed to the engine, which replaces it with
// SQL reading it.
%type perfetto_subquery_pipeline {uint32_t}
perfetto_subquery_pipeline(A) ::= perfetto_pipeline(P). {
A = P;
synq_parser_expand_node(pCtx, P, "pipeline", 8);
}

seltablist(A) ::= stl_prefix(A) LP perfetto_subquery_pipeline(P) RP as(Z)
on_using(N). {
pCtx->saw_subquery = 1;
uint32_t sub = synq_parse_subquery_table_source(
pCtx, P, Z.name,
Z.has_as ? SYNTAQLITE_BOOL_TRUE : SYNTAQLITE_BOOL_FALSE);
if (A == SYNTAQLITE_NULL_NODE) {
synq_reject_dangling_on_using(pCtx, N);
A = sub;
} else {
SyntaqliteNode *pfx = AST_NODE(&pCtx->ast, A);
A = synq_parse_join_clause(pCtx,
pfx->join_prefix.join_type,
pfx->join_prefix.modifiers,
pfx->join_prefix.source,
sub, N.on_expr, N.using_cols);
}
}
wqitem(A) ::= withnm(X) eidlist_opt(Y) wqas(M) LP perfetto_subquery_pipeline(P)
RP. {
A = synq_parse_cte_definition(pCtx, synq_span_dequote(pCtx, X),
(SyntaqliteMaterialized)M, Y, P);
}

// ---------- PERFETTO PRAGMA ----------

// A setting of the engine, rather than of SQLite. Any expression parses;
Expand Down
3,788 changes: 1,920 additions & 1,868 deletions src/perfetto_sql/syntaqlite/syntaqlite_perfetto.c

Large diffs are not rendered by default.

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
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,10 @@

#include "src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.h"

#include <sqlite3.h>

#include <algorithm>
#include <cstdint>
#include <memory>
#include <string>
#include <utility>
Expand Down Expand Up @@ -1133,6 +1136,265 @@ TEST_F(PerfettoSqlConnectionPipelineTest,
EXPECT_TRUE(first->stmt.status().ok());
}

TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesAsSubqueries) {
auto rows = Rows(
"SELECT t.id, p.total * 2 FROM tree t "
"JOIN (FROM tree |> TREE ACCUMULATE UP SUM(self) AS total) p USING (id)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,200", "1,120", "2,60", "3,80"));

rows = Rows(
"WITH totals AS (FROM tree |> TREE ACCUMULATE UP SUM(self) AS total) "
"SELECT id, total FROM totals");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,100", "1,60", "2,30", "3,40"));
}

TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesInViewsAndFunctions) {
ASSERT_TRUE(Rows("CREATE PERFETTO VIEW totals AS SELECT id, total FROM "
"(FROM tree |> TREE ACCUMULATE UP SUM(self) AS total)")
.ok());
ASSERT_TRUE(Rows("CREATE PERFETTO FUNCTION subtree(root LONG) "
"RETURNS TABLE(total LONG) AS SELECT total FROM "
"(FROM tree |> TREE ACCUMULATE UP SUM(self) AS total) "
"WHERE id = $root")
.ok());
auto rows = Rows("SELECT id, total FROM totals");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,100", "1,60", "2,30", "3,40"));
rows = Rows("SELECT total FROM subtree(1)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("60"));

// Both read the table as it is when they run.
ASSERT_TRUE(Rows("DELETE FROM tree WHERE id = 3").ok());
rows = Rows("SELECT id, total FROM totals");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,60", "1,20", "2,30"));
rows = Rows("SELECT total FROM subtree(1)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
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 "
"self FROM tree) |> TREE ACCUMULATE UP SUM(self) AS total)")
.status()
.message(),
testing::HasSubstr("Cannot read `$k`"));
}

// Macros and pipelines nest inside each other in every way: pipelines in
// macros, macros in those pipelines' sources and stages, and macro arguments
// which are themselves macros expanding to pipelines.
TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesAndMacrosNest) {
ASSERT_TRUE(Rows(R"(
CREATE PERFETTO MACRO own() RETURNS ColumnName AS self;
CREATE PERFETTO MACRO src() RETURNS TableOrSubquery
AS (SELECT id, parent_id, self FROM tree);
CREATE PERFETTO MACRO piped_src() RETURNS TableOrSubquery
AS (FROM src!() |> SELECT id, parent_id, own!());
CREATE PERFETTO MACRO fold(t TableOrSubquery) RETURNS TableOrSubquery
AS (FROM (SELECT * FROM $t) |> TREE ACCUMULATE UP SUM(own!()) AS total);
)")
.ok());

auto rows = Rows("SELECT id, total FROM fold!(piped_src!())");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,100", "1,60", "2,30", "3,40"));

// Again, read by a pipeline through SQL, and passed to a macro as a
// pipeline written out in the call.
rows = Rows(
"SELECT id, total FROM (FROM (SELECT id, parent_id, total AS self FROM "
"fold!((FROM piped_src!() |> SELECT id, parent_id, self))) "
"|> TREE ACCUMULATE UP SUM(self) AS total)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,230", "1,100", "2,30", "3,40"));
}

TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesAsPipelineSources) {
auto rows = Rows(
"FROM (FROM tree |> TREE ACCUMULATE UP SUM(self) AS total) AS t "
"|> TREE ACCUMULATE UP SUM(total) AS again |> SELECT t.id, again");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,230", "1,100", "2,30", "3,40"));
}

TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesAsIntersectionOperands) {
ASSERT_TRUE(Rows("CREATE PERFETTO TABLE a AS SELECT 10 AS ts, 5 AS dur "
"UNION ALL SELECT 20, 5")
.ok());
ASSERT_TRUE(
Rows("CREATE PERFETTO TABLE b AS SELECT 12 AS ts, 10 AS dur").ok());
auto rows = Rows(
"INTERVAL INTERSECTION OF ((FROM a |> SELECT ts, dur) AS x, b AS y) "
"|> SELECT ts, dur");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("12,3", "20,2"));
}

TEST_F(PerfettoSqlConnectionPipelineTest, PipelineSubqueriesNeedPipelines) {
ASSERT_TRUE(Rows("PERFETTO PRAGMA pipelines = 0").ok());
EXPECT_THAT(Rows("SELECT * FROM (FROM tree)").status().message(),
testing::HasSubstr("Pipelines are not enabled"));
}

// A pipeline is replaced where it was written, so one in a macro works as
// one written out.
TEST_F(PerfettoSqlConnectionPipelineTest, PipelinesInMacros) {
ASSERT_TRUE(Rows("CREATE PERFETTO MACRO totals(t TableOrSubquery) "
"RETURNS TableOrSubquery "
"AS (FROM $t |> TREE ACCUMULATE UP SUM(self) AS total)")
.ok());
auto rows = Rows("SELECT id, total FROM totals!(tree)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,100", "1,60", "2,30", "3,40"));

// In a macro which a pipeline's SQL source calls, with another macro
// called inside it.
ASSERT_TRUE(Rows("CREATE PERFETTO MACRO just_tree() RETURNS TableOrSubquery "
"AS tree")
.ok());
rows = Rows(
"SELECT id, total FROM (FROM (SELECT id, parent_id, total AS self "
"FROM totals!(just_tree!())) "
"|> TREE ACCUMULATE UP SUM(self) AS total)");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("0,230", "1,100", "2,30", "3,40"));

// An error in one traces back through the call.
ASSERT_TRUE(Rows("CREATE PERFETTO MACRO bad() RETURNS TableOrSubquery "
"AS (FROM tree |> SELECT nope)")
.ok());
std::string error = Rows("SELECT * FROM bad!()").status().message();
EXPECT_THAT(error, testing::HasSubstr("bad!()"));
EXPECT_THAT(error, testing::HasSubstr("no such column: 'nope'"));
}

// SQL -> pipeline -> SQL -> pipeline -> SQL -> pipeline, SQLite and the
// executor taking turns at every level.
class PipelineNestingTest : public PerfettoSqlConnectionPipelineTest {
protected:
// A binary tree: node i's parent is (i - 1) / 2.
static constexpr uint32_t kNodes = 5000;

// Each node's subtree size, tripled, summed over its subtree, plus one.
static constexpr char kThreeLevels[] =
"SELECT id, total + 1 AS z FROM ("
" FROM (SELECT id, parent_id, total * 3 AS y FROM ("
" FROM big |> TREE ACCUMULATE UP SUM(x) AS total))"
" |> TREE ACCUMULATE UP SUM(y) AS total)";

void SetUp() override {
PerfettoSqlConnectionPipelineTest::SetUp();
ASSERT_TRUE(Rows("CREATE TABLE big AS WITH RECURSIVE n(i) AS "
"(SELECT 0 UNION ALL SELECT i + 1 FROM n WHERE i < " +
std::to_string(kNodes - 1) +
") SELECT i AS id, CASE WHEN i = 0 THEN NULL "
"ELSE (i - 1) / 2 END AS parent_id, 1 AS x FROM n")
.ok());
}

static std::vector<int64_t> Expected() {
std::vector<int64_t> size(kNodes, 0);
std::vector<int64_t> total(kNodes, 0);
for (uint32_t i = kNodes; i-- > 0;) {
size[i] += 1;
total[i] += 3 * size[i];
if (i != 0) {
size[(i - 1) / 2] += size[i];
total[(i - 1) / 2] += total[i];
}
}
for (int64_t& t : total) {
t += 1;
}
return total;
}

// Checks each row of `stmt` against Expected(), returning how many there
// were.
static uint32_t CheckRows(SqliteConnection::PreparedStatement& stmt) {
std::vector<int64_t> expected = Expected();
uint32_t rows = 0;
for (bool more = !stmt.IsDone(); more; more = stmt.Step()) {
auto id =
static_cast<uint32_t>(sqlite3_column_int64(stmt.sqlite_stmt(), 0));
EXPECT_EQ(sqlite3_column_int64(stmt.sqlite_stmt(), 1), expected[id]);
++rows;
}
EXPECT_TRUE(stmt.status().ok()) << stmt.status().message();
return rows;
}
};

TEST_F(PipelineNestingTest, EveryLevelStreams) {
auto res = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(kThreeLevels));
ASSERT_TRUE(res.ok()) << res.status().message();
EXPECT_EQ(CheckRows(res->stmt), kNodes);
}

TEST_F(PipelineNestingTest, ReadingAgainRunsEveryLevelAgain) {
ASSERT_TRUE(Rows("CREATE TABLE picks(k); "
"INSERT INTO picks VALUES (0), (1), (4999)")
.ok());
auto res = connection_->ExecuteUntilLastStatement(SqlSource::FromExecuteQuery(
std::string("SELECT p.id, p.z FROM picks CROSS JOIN (") + kThreeLevels +
") p WHERE p.id = picks.k"));
ASSERT_TRUE(res.ok()) << res.status().message();
EXPECT_EQ(CheckRows(res->stmt), 3u);
}

TEST_F(PipelineNestingTest, ErrorsSurfaceThroughEveryLevel) {
ASSERT_TRUE(Rows("UPDATE big SET x = 'bad' WHERE id = 4000").ok());
EXPECT_THAT(Rows(kThreeLevels).status().message(),
testing::HasSubstr("column 'x' holds a string"));
}

TEST_F(PipelineNestingTest, StoppingEarlyReleasesEveryLevel) {
auto res = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(kThreeLevels));
ASSERT_TRUE(res.ok()) << res.status().message();
for (uint32_t i = 0; i < 10; ++i) {
ASSERT_TRUE(res->stmt.Step()) << res->stmt.status().message();
}
}

TEST_F(PipelineNestingTest, InterleavedRuns) {
auto a = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(kThreeLevels));
auto b = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(kThreeLevels));
ASSERT_TRUE(a.ok()) << a.status().message();
ASSERT_TRUE(b.ok()) << b.status().message();
std::vector<int64_t> expected = Expected();
uint32_t rows = 1;
for (;;) {
for (auto* res : {&*a, &*b}) {
sqlite3_stmt* stmt = res->stmt.sqlite_stmt();
auto id = static_cast<uint32_t>(sqlite3_column_int64(stmt, 0));
EXPECT_EQ(sqlite3_column_int64(stmt, 1), expected[id]);
}
bool more_a = a->stmt.Step();
bool more_b = b->stmt.Step();
ASSERT_EQ(more_a, more_b);
if (!more_a) {
break;
}
++rows;
}
EXPECT_TRUE(a->stmt.status().ok()) << a->stmt.status().message();
EXPECT_TRUE(b->stmt.status().ok()) << b->stmt.status().message();
EXPECT_EQ(rows, kNodes);
}

TEST_F(PerfettoSqlConnectionPipelineTest, ColumnReadersRefreshAcrossBatches) {
ASSERT_TRUE(connection_
->Execute(SqlSource::FromExecuteQuery(R"(
Expand Down
Loading
Loading