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
2 changes: 2 additions & 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/plan_serialization.cc",
],
}

Expand Down Expand Up @@ -19109,6 +19110,7 @@ filegroup {
name: "perfetto_src_trace_processor_perfetto_sql_pipeline_unittests",
srcs: [
"src/trace_processor/perfetto_sql/pipeline/physical_plan_unittest.cc",
"src/trace_processor/perfetto_sql/pipeline/plan_serialization_unittest.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/plan_serialization.cc",
"src/trace_processor/perfetto_sql/pipeline/plan_serialization.h",
],
)

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/plan_serialization.h"
#include "src/trace_processor/sqlite/bindings/sqlite_column.h"
#include "src/trace_processor/sqlite/bindings/sqlite_type.h"
#include "src/trace_processor/sqlite/bindings/sqlite_value.h"
Expand Down Expand Up @@ -359,7 +360,7 @@ PerfettoSqlConnection::PerfettoSqlConnection(
{
auto ctx = std::make_unique<PipelineModule::Context>();
ctx->pool = pool_;
pipeline_context_ = ctx.get();
ctx->connection = this;
RegisterVirtualTableModule<PipelineModule>(PipelineModule::kName,
std::move(ctx));
}
Expand Down Expand Up @@ -701,9 +702,8 @@ PerfettoSqlConnection::ProcessFrame(size_t frame_idx) {
std::holds_alternative<PerfettoSqlParser::SqliteSql>(stmt))) {
source_to_prepare = parser->TakeStatementSql();
} else if (std::holds_alternative<PerfettoSqlParser::Pipeline>(stmt)) {
auto pipeline =
std::get<PerfettoSqlParser::Pipeline>(parser->TakeStatement());
pipeline_plan = std::move(pipeline.plan);
pipeline_plan = std::move(
std::get<PerfettoSqlParser::Pipeline>(parser->TakeStatement()).plan);
source_to_prepare = parser->TakeStatementSql();
} else {
is_dummy = true;
Expand All @@ -714,8 +714,8 @@ PerfettoSqlConnection::ProcessFrame(size_t frame_idx) {
{
PERFETTO_TP_TRACE(metatrace::Category::QUERY_TIMELINE, "QUERY_PREPARE");
if (pipeline_plan) {
ASSIGN_OR_RETURN(next_stmt, PreparePipeline(std::move(*pipeline_plan),
*source_to_prepare));
ASSIGN_OR_RETURN(next_stmt,
PreparePipeline(*pipeline_plan, *source_to_prepare));
} else {
auto stmt_result =
connection_->PrepareStatement(std::move(*source_to_prepare));
Expand Down Expand Up @@ -979,9 +979,9 @@ base::Status PerfettoSqlConnection::ExecuteCreateTable(
[&create_table](metatrace::Record* record) {
record->AddArg("table_name", create_table.name);
});
auto* logical = std::get_if<pipeline::LogicalPlan>(&create_table.body);
const auto* logical = std::get_if<pipeline::LogicalPlan>(&create_table.body);
base::StatusOr<SqliteConnection::PreparedStatement> stmt_or =
logical ? PreparePipeline(std::move(*logical), statement_sql)
logical ? PreparePipeline(*logical, statement_sql)
: connection_->PrepareStatement(
std::move(std::get<SqlSource>(create_table.body)));
ASSIGN_OR_RETURN(auto stmt, std::move(stmt_or));
Expand Down Expand Up @@ -1049,14 +1049,26 @@ base::Status PerfettoSqlConnection::ExecuteCreateTable(
}

base::StatusOr<SqliteConnection::PreparedStatement>
PerfettoSqlConnection::PreparePipeline(pipeline::LogicalPlan logical,
PerfettoSqlConnection::PreparePipeline(const pipeline::LogicalPlan& plan,
const SqlSource& source) {
PERFETTO_TP_TRACE(metatrace::Category::QUERY_TIMELINE, "PIPELINE_PLAN");
auto sql = PipelineModule::SelectFrom(plan);
if (!sql.ok()) {
return base::ErrStatus("%s%s", source.AsTraceback(0).c_str(),
sql.status().c_message());
}
return connection_->PrepareStatement(source.RewriteAllIgnoreExisting(
SqlSource::FromTraceProcessorImplementation(std::move(*sql))));
}

base::StatusOr<std::unique_ptr<pipeline::PhysicalPlan>>
PerfettoSqlConnection::LoadPipeline(std::string_view serialized) {
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_;
return PipelineModule::Prepare(connection_.get(), pipeline_context_,
pipeline::Lower(logical, env), source);
return pipeline::Lower(plan, env);
}

base::Status PerfettoSqlConnection::ExecuteCreateView(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,12 +39,12 @@
#include "src/trace_processor/core/plugin/registration.h"
#include "src/trace_processor/perfetto_sql/engine/dataframe_module.h"
#include "src/trace_processor/perfetto_sql/engine/perfetto_sql_database.h"
#include "src/trace_processor/perfetto_sql/engine/pipeline_module.h"
#include "src/trace_processor/perfetto_sql/engine/runtime_table_function.h"
#include "src/trace_processor/perfetto_sql/engine/static_table_function_module.h"
#include "src/trace_processor/perfetto_sql/parser/function_util.h"
#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/sqlite/bindings/sqlite_module.h"
#include "src/trace_processor/sqlite/bindings/sqlite_result.h"
#include "src/trace_processor/sqlite/bindings/sqlite_window_function.h"
Expand Down Expand Up @@ -174,6 +174,10 @@ class PerfettoSqlConnection {
base::StatusOr<SqliteConnection::PreparedStatement> PrepareSqliteStatement(
SqlSource sql);

// Loads a plan written by pipeline::SerializePlan, ready to run.
base::StatusOr<std::unique_ptr<pipeline::PhysicalPlan>> LoadPipeline(
std::string_view serialized);

// Registers a virtual table module with the given name.
//
// |name|: name of the module in SQL.
Expand Down Expand Up @@ -466,7 +470,7 @@ class PerfettoSqlConnection {
base::Status ExecuteCreateMacro(const PerfettoSqlParser::CreateMacro&);

base::StatusOr<SqliteConnection::PreparedStatement> PreparePipeline(
pipeline::LogicalPlan,
const pipeline::LogicalPlan&,
const SqlSource&);

base::Status ExecuteCreateIndex(const PerfettoSqlParser::CreateIndex&);
Expand Down Expand Up @@ -591,7 +595,6 @@ class PerfettoSqlConnection {
// context class of the module inherits from ModuleStateManagerBase.
std::vector<sqlite::ModuleStateManagerBase*> virtual_module_state_managers_;

PipelineModule::Context* pipeline_context_ = nullptr;
RuntimeTableFunctionModule::Context* runtime_table_fn_context_ = nullptr;
StaticTableFunctionModule::Context* static_table_fn_context_ = nullptr;
DataframeModule::Context* dataframe_context_ = nullptr;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -849,6 +849,19 @@ class PerfettoSqlConnectionPipelineTest : public PerfettoSqlConnectionTest {
return rows;
}

// The SQL a pipeline statement is prepared as.
std::string PipelineSql(const std::string& pipeline) {
auto res = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(pipeline));
PERFETTO_CHECK(res.ok());
return res->stmt.sql();
}

// The FROM clause of `sql`, from its leading space.
static std::string From(const std::string& sql) {
return sql.substr(sql.find(" FROM "));
}

std::vector<std::string> ColumnNames(const std::string& sql) {
auto res = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery(sql));
Expand Down Expand Up @@ -1011,21 +1024,6 @@ TEST_F(PerfettoSqlConnectionPipelineTest, AReadTableCanBeReplaced) {
EXPECT_EQ(rows, 4u);
}

// Finalization retires the TEMP table; the next execution drops it safely.
TEST_F(PerfettoSqlConnectionPipelineTest, CompletedPipelineTablesAreDropped) {
ASSERT_TRUE(Rows("FROM tree").ok());
auto tables = Rows(
"SELECT name FROM sqlite_temp_schema WHERE name GLOB "
"'__intrinsic_pipeline_*'");
ASSERT_TRUE(tables.ok()) << tables.status().message();
EXPECT_TRUE(tables->empty());
auto modules = Rows(
"SELECT name FROM pragma_module_list WHERE name GLOB "
"'__intrinsic_pipeline*'");
ASSERT_TRUE(modules.ok()) << modules.status().message();
EXPECT_THAT(*modules, testing::ElementsAre("__intrinsic_pipeline"));
}

TEST_F(PerfettoSqlConnectionPipelineTest,
DifferentSchemasAndConcurrentStatements) {
auto first =
Expand All @@ -1048,7 +1046,7 @@ TEST_F(PerfettoSqlConnectionPipelineTest,
EXPECT_EQ(rows, 4u);
EXPECT_TRUE(second->stmt.status().ok());
// The schema changed after preparing the first statement. Resetting and
// rerunning it must retain the bound plan through SQLite's reprepare path.
// rerunning it loads the pipeline again from its plan.
ASSERT_EQ(sqlite3_reset(first->stmt.sqlite_stmt()), SQLITE_OK);
rows = 0;
while (first->stmt.Step())
Expand All @@ -1057,54 +1055,32 @@ TEST_F(PerfettoSqlConnectionPipelineTest,
EXPECT_TRUE(first->stmt.status().ok());
}

TEST_F(PerfettoSqlConnectionPipelineTest,
CleanupRetriesWhileAnotherStatementIsActive) {
{
auto active = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery("FROM tree"));
ASSERT_TRUE(active.ok()) << active.status().message();
ASSERT_TRUE(Rows("FROM (SELECT 123 AS value)").ok());
// The finished pipeline can be retired even though the active statement
// prevents schema changes. Retrying cleanup must not interrupt either.
ASSERT_TRUE(Rows("SELECT 1").ok());
while (active->stmt.Step()) {
}
EXPECT_TRUE(active->stmt.status().ok());
}
auto tables = Rows(
"SELECT name FROM sqlite_temp_schema WHERE name GLOB "
"'__intrinsic_pipeline_*'");
ASSERT_TRUE(tables.ok()) << tables.status().message();
EXPECT_TRUE(tables->empty());
}
// SQL holding a pipeline's plan can be stored and run later, against the
// tables as they are by then.
TEST_F(PerfettoSqlConnectionPipelineTest, StoredPipelinesRunLater) {
std::string sql = PipelineSql(
"FROM tree |> TREE ACCUMULATE UP SUM(self) AS total |> SELECT id, total");
ASSERT_TRUE(Rows("CREATE VIEW totals AS " + sql).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"));

TEST_F(PerfettoSqlConnectionPipelineTest,
CleanupAcrossRollbackAndExecutionFailure) {
ASSERT_TRUE(Rows("FROM tree").ok());
ASSERT_TRUE(Rows("BEGIN; FROM tree").ok());
ASSERT_TRUE(Rows("ROLLBACK").ok());
EXPECT_FALSE(Rows("FROM (SELECT 0 AS id, NULL AS parent_id, 'bad' AS value) "
"|> TREE ACCUMULATE UP SUM(value) AS total")
.ok());
auto tables = Rows(
"SELECT name FROM sqlite_temp_schema WHERE name GLOB "
"'__intrinsic_pipeline_*'");
ASSERT_TRUE(tables.ok()) << tables.status().message();
EXPECT_TRUE(tables->empty());
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"));
}

TEST_F(PerfettoSqlConnectionPipelineTest,
FailedCreateDoesNotRetireAnExistingTable) {
ASSERT_TRUE(Rows("CREATE TEMP TABLE __intrinsic_pipeline_0(value); "
"INSERT INTO __intrinsic_pipeline_0 VALUES(123)")
.ok());
EXPECT_FALSE(Rows("FROM tree").ok());
auto existing = Rows("SELECT value FROM temp.__intrinsic_pipeline_0");
ASSERT_TRUE(existing.ok()) << existing.status().message();
EXPECT_THAT(*existing, testing::ElementsAre("123"));
TEST_F(PerfettoSqlConnectionPipelineTest, BadPlans) {
auto rows = Rows("SELECT c0 FROM __intrinsic_pipeline('FROM tree')");
EXPECT_THAT(rows.status().message(), testing::HasSubstr("expected a plan"));
rows = Rows("SELECT c0 FROM __intrinsic_pipeline(X'00')");
EXPECT_THAT(rows.status().message(), testing::HasSubstr("malformed plan"));
rows = Rows("SELECT c1" + From(PipelineSql("FROM (SELECT 1 AS x)")));
EXPECT_THAT(rows.status().message(), testing::HasSubstr("no column c1"));
}

TEST_F(PerfettoSqlConnectionPipelineTest, TemporaryTablesAreConnectionLocal) {
TEST_F(PerfettoSqlConnectionPipelineTest, ForksRunPipelinesIndependently) {
auto fork = connection_->Fork();
auto first = connection_->ExecuteUntilLastStatement(
SqlSource::FromExecuteQuery("FROM (SELECT 123 AS value)"));
Expand All @@ -1120,52 +1096,16 @@ TEST_F(PerfettoSqlConnectionPipelineTest, TemporaryTablesAreConnectionLocal) {

TEST_F(PerfettoSqlConnectionPipelineTest,
OutputConstraintsApplyAfterAccumulation) {
auto context = std::make_unique<PipelineModule::Context>();
context->pool = &pool_;
auto* ctx = context.get();
connection_->RegisterVirtualTableModule<PipelineModule>("test_pipeline",
std::move(context));
ASSERT_TRUE(
Rows("CREATE VIRTUAL TABLE temp.test_output USING test_pipeline(4)")
.ok());

pipeline::LogicalPlan logical;
for (const char* name : {"id", "parent_id", "self"}) {
auto id = logical.AddColumn(name, core::Int64{});
logical.output.push_back({name, id});
}
pipeline::PlanNodeId scan = logical.AddNode(pipeline::op::Scan{
SqlSource::FromExecuteQuery("SELECT id, parent_id, self FROM tree"),
logical.output});
auto total = logical.AddColumn("total", core::Int64{});
pipeline::op::TreeAccumulate fold;
fold.direction = pipeline::op::TreeDirection::kUp;
fold.node_column = 0;
fold.parent_column = 1;
fold.aggregates.push_back(
{pipeline::op::TreeAccumulate::Function::kSum, 2, total});
logical.AddNode(std::move(fold), {scan});
logical.output.push_back({"total", total});
pipeline::LowerEnvironment env{connection_->sqlite_connection(), &pool_};
std::string from =
From(PipelineSql("FROM (SELECT id, parent_id, self FROM tree) "
"|> TREE ACCUMULATE UP SUM(self) AS total"));
// The root is last in child-first output. Filtering by its output rowid
// must retain all descendants while calculating its total.
for (const char* rhs : {"3", "3.0", "'3'"}) {
auto stmt = connection_->sqlite_connection()->PrepareStatement(
SqlSource::FromExecuteQuery(
"SELECT c3 FROM temp.test_output(?) WHERE rowid = " +
std::string(rhs) + " AND c3 > 50"));
ASSERT_TRUE(stmt.status().ok()) << stmt.status().message();
PipelineModule::Invocation invocation{ctx, "unused",
pipeline::Lower(logical, env)};
ASSERT_EQ(sqlite3_bind_pointer(stmt.sqlite_stmt(), 1, &invocation,
PipelineModule::kPlanPointerType, nullptr),
SQLITE_OK);
ASSERT_TRUE(stmt.Step()) << stmt.status().message();
EXPECT_EQ(sqlite3_column_int64(stmt.sqlite_stmt(), 0), 100);
EXPECT_FALSE(stmt.Step());
EXPECT_TRUE(stmt.status().ok());
// Explicitly close cursors before the borrowed plan goes out of scope.
sqlite3_reset(stmt.sqlite_stmt());
auto rows = Rows("SELECT c3" + from + " WHERE rowid = " + std::string(rhs) +
" AND c3 > 50");
ASSERT_TRUE(rows.ok()) << rows.status().message();
EXPECT_THAT(*rows, testing::ElementsAre("100"));
}
}

Expand Down
Loading
Loading