diff --git a/Android.bp b/Android.bp index 1d0e9450a5a..5143c1a5321 100644 --- a/Android.bp +++ b/Android.bp @@ -3157,7 +3157,6 @@ cc_test { ":perfetto_src_trace_processor_metatrace", ":perfetto_src_trace_processor_metrics_metrics", ":perfetto_src_trace_processor_perfetto_sql_engine_engine", - ":perfetto_src_trace_processor_perfetto_sql_exec_exec", ":perfetto_src_trace_processor_perfetto_sql_generator_generator", ":perfetto_src_trace_processor_perfetto_sql_intrinsics_types_types", ":perfetto_src_trace_processor_perfetto_sql_parser_parser", @@ -19000,22 +18999,6 @@ filegroup { ], } -// GN: //src/trace_processor/perfetto_sql/exec:exec -filegroup { - name: "perfetto_src_trace_processor_perfetto_sql_exec_exec", - srcs: [ - "src/trace_processor/perfetto_sql/exec/sql_scan.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", - ], -} - // GN: //src/trace_processor/perfetto_sql/generator:gen_cc_perfetto_sql_descriptor genrule { name: "perfetto_src_trace_processor_perfetto_sql_generator_gen_cc_perfetto_sql_descriptor", @@ -19119,9 +19102,6 @@ filegroup { // GN: //src/trace_processor/perfetto_sql/schema:schema filegroup { name: "perfetto_src_trace_processor_perfetto_sql_schema_schema", - srcs: [ - "src/trace_processor/perfetto_sql/schema/query_schema.cc", - ], } // GN: //src/trace_processor/perfetto_sql/stdlib/android:android @@ -21513,7 +21493,6 @@ cc_library_static { ":perfetto_src_trace_processor_metatrace", ":perfetto_src_trace_processor_metrics_metrics", ":perfetto_src_trace_processor_perfetto_sql_engine_engine", - ":perfetto_src_trace_processor_perfetto_sql_exec_exec", ":perfetto_src_trace_processor_perfetto_sql_generator_generator", ":perfetto_src_trace_processor_perfetto_sql_intrinsics_types_types", ":perfetto_src_trace_processor_perfetto_sql_parser_parser", @@ -24529,8 +24508,6 @@ cc_test { ":perfetto_src_trace_processor_metrics_unittests", ":perfetto_src_trace_processor_perfetto_sql_engine_engine", ":perfetto_src_trace_processor_perfetto_sql_engine_unittests", - ":perfetto_src_trace_processor_perfetto_sql_exec_exec", - ":perfetto_src_trace_processor_perfetto_sql_exec_unittests", ":perfetto_src_trace_processor_perfetto_sql_generator_generator", ":perfetto_src_trace_processor_perfetto_sql_generator_unittests", ":perfetto_src_trace_processor_perfetto_sql_intrinsics_types_types", @@ -25689,7 +25666,6 @@ cc_library_static { ":perfetto_src_trace_processor_metatrace", ":perfetto_src_trace_processor_metrics_metrics", ":perfetto_src_trace_processor_perfetto_sql_engine_engine", - ":perfetto_src_trace_processor_perfetto_sql_exec_exec", ":perfetto_src_trace_processor_perfetto_sql_generator_generator", ":perfetto_src_trace_processor_perfetto_sql_intrinsics_types_types", ":perfetto_src_trace_processor_perfetto_sql_parser_parser", @@ -26407,7 +26383,6 @@ cc_binary_host { ":perfetto_src_trace_processor_metatrace", ":perfetto_src_trace_processor_metrics_metrics", ":perfetto_src_trace_processor_perfetto_sql_engine_engine", - ":perfetto_src_trace_processor_perfetto_sql_exec_exec", ":perfetto_src_trace_processor_perfetto_sql_generator_generator", ":perfetto_src_trace_processor_perfetto_sql_intrinsics_types_types", ":perfetto_src_trace_processor_perfetto_sql_parser_parser", diff --git a/BUILD b/BUILD index a8bbfafd0ff..e28547d2081 100644 --- a/BUILD +++ b/BUILD @@ -470,7 +470,6 @@ perfetto_cc_library( ":src_trace_processor_metatrace", ":src_trace_processor_metrics_metrics", ":src_trace_processor_perfetto_sql_engine_engine", - ":src_trace_processor_perfetto_sql_exec_exec", ":src_trace_processor_perfetto_sql_generator_generator", ":src_trace_processor_perfetto_sql_intrinsics_types_types", ":src_trace_processor_perfetto_sql_parser_parser", @@ -795,7 +794,6 @@ perfetto_cc_library( ":src_trace_processor_metatrace", ":src_trace_processor_metrics_metrics", ":src_trace_processor_perfetto_sql_engine_engine", - ":src_trace_processor_perfetto_sql_exec_exec", ":src_trace_processor_perfetto_sql_generator_generator", ":src_trace_processor_perfetto_sql_intrinsics_types_types", ":src_trace_processor_perfetto_sql_parser_parser", @@ -3825,15 +3823,6 @@ perfetto_filegroup( ], ) -# GN target: //src/trace_processor/perfetto_sql/exec:exec -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", - ], -) - # GN target: //src/trace_processor/perfetto_sql/generator:generator perfetto_filegroup( name = "src_trace_processor_perfetto_sql_generator_generator", @@ -3901,8 +3890,6 @@ 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/type_mapping.h", ], ) @@ -11856,7 +11843,6 @@ perfetto_cc_library( ":src_trace_processor_metatrace", ":src_trace_processor_metrics_metrics", ":src_trace_processor_perfetto_sql_engine_engine", - ":src_trace_processor_perfetto_sql_exec_exec", ":src_trace_processor_perfetto_sql_generator_generator", ":src_trace_processor_perfetto_sql_intrinsics_types_types", ":src_trace_processor_perfetto_sql_parser_parser", @@ -12212,7 +12198,6 @@ perfetto_cc_binary( ":src_trace_processor_metatrace", ":src_trace_processor_metrics_metrics", ":src_trace_processor_perfetto_sql_engine_engine", - ":src_trace_processor_perfetto_sql_exec_exec", ":src_trace_processor_perfetto_sql_generator_generator", ":src_trace_processor_perfetto_sql_intrinsics_types_types", ":src_trace_processor_perfetto_sql_parser_parser", diff --git a/src/perfetto_sql/analysis/relation.cc b/src/perfetto_sql/analysis/relation.cc index 468afcdc815..2e4fa27fec1 100644 --- a/src/perfetto_sql/analysis/relation.cc +++ b/src/perfetto_sql/analysis/relation.cc @@ -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) { @@ -186,6 +192,10 @@ 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 LeafColumns(LeafRelation, + std::vector& hidden); const Catalog& catalog_; // Lineage string_views point into each view's sql string and parse tree, so @@ -198,6 +208,22 @@ class RelationAnalyzer::Impl { bool preserves_rows_ = true; }; +std::vector RelationAnalyzer::Impl::LeafColumns( + LeafRelation found, + std::vector& hidden) { + leaves_.push_back(std::make_unique(std::move(found))); + const LeafRelation* relation = leaves_.back().get(); + std::vector 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) { @@ -450,18 +476,7 @@ base::StatusOr> RelationAnalyzer::Impl::Relation( int depth, std::vector& hidden) { if (std::optional found = catalog_.FindLeafRelation(name)) { - leaves_.push_back(std::make_unique(std::move(*found))); - const LeafRelation* relation = leaves_.back().get(); - std::vector 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; + return LeafColumns(std::move(*found), hidden); } if (depth >= kMaxDepth) { return base::ErrStatus( diff --git a/src/perfetto_sql/syntaqlite/BUILD.gn b/src/perfetto_sql/syntaqlite/BUILD.gn index 53c6c6208dd..56948403207 100644 --- a/src/perfetto_sql/syntaqlite/BUILD.gn +++ b/src/perfetto_sql/syntaqlite/BUILD.gn @@ -37,7 +37,6 @@ source_set("syntaqlite") { visibility = [ "..:intrinsic_macro_expansion", "../../trace_processor/perfetto_sql/engine:unittests", - "../../trace_processor/perfetto_sql/exec:*", "../../trace_processor/perfetto_sql/parser:*", "../../trace_processor/perfetto_sql/pipeline:*", "../../trace_processor/perfetto_sql/schema:*", diff --git a/src/trace_processor/BUILD.gn b/src/trace_processor/BUILD.gn index fe409d043f6..931ebebc9c3 100644 --- a/src/trace_processor/BUILD.gn +++ b/src/trace_processor/BUILD.gn @@ -473,7 +473,6 @@ perfetto_unittest_source_set("unittests") { deps += [ "../perfetto_sql/analysis:unittests", "perfetto_sql/engine:unittests", - "perfetto_sql/exec:unittests", "perfetto_sql/parser:unittests", "perfetto_sql/pipeline:unittests", "perfetto_sql/tokenizer:unittests", diff --git a/src/trace_processor/core/dataframe/adhoc_dataframe_builder.cc b/src/trace_processor/core/dataframe/adhoc_dataframe_builder.cc index 3280469a8bc..a17e06bad65 100644 --- a/src/trace_processor/core/dataframe/adhoc_dataframe_builder.cc +++ b/src/trace_processor/core/dataframe/adhoc_dataframe_builder.cc @@ -42,13 +42,31 @@ #include "src/trace_processor/core/util/flex_vector.h" namespace perfetto::trace_processor::core::dataframe { +namespace { + +// The number of values collected in `storage`. +size_t CollectedSize(const Storage& storage) { + switch (storage.type().index()) { + case StorageType::GetTypeIndex(): + return storage.unchecked_get().size(); + case StorageType::GetTypeIndex(): + return storage.unchecked_get().size(); + case StorageType::GetTypeIndex(): + return storage.unchecked_get().size(); + default: + PERFETTO_FATAL("Unexpected storage type"); + } +} + +} // namespace AdhocDataframeBuilder::AdhocDataframeBuilder(std::vector names, StringPool* pool, const Options& options) : string_pool_(pool), did_declare_types_(!options.types.empty()), - emit_auto_id_(options.emit_auto_id) { + emit_auto_id_(options.emit_auto_id), + analyze_(options.analyze) { PERFETTO_DCHECK(options.types.empty() || options.types.size() == names.size()); for (uint32_t i = 0; i < names.size(); ++i) { @@ -125,6 +143,15 @@ base::StatusOr AdhocDataframeBuilder::Build() && { Unsorted{}, HasDuplicates{}, })); + } else if (!analyze_) { + non_null_row_count = CollectedSize(*state.storage); + columns.emplace_back(std::make_shared(Column{ + std::move(*state.storage), + CreateNullStorageFromBitvector(std::move(state.null_overlay), + state.nullability_type), + Unsorted{}, + HasDuplicates{}, + })); } else if (state.storage->type().Is()) { auto& data = state.storage->unchecked_get(); non_null_row_count = data.size(); @@ -227,8 +254,12 @@ base::StatusOr AdhocDataframeBuilder::Build() && { Column{Storage{Storage::Id{static_cast(row_count)}}, NullStorage::NonNull{}, IdSorted{}, NoDuplicates{}})); } - return Dataframe(true, std::move(column_names_), std::move(columns), - static_cast(row_count), string_pool_); + Dataframe dataframe(analyze_, std::move(column_names_), std::move(columns), + static_cast(row_count), string_pool_); + if (!analyze_) { + dataframe.FinalizeWithoutStatistics(); + } + return dataframe; } Storage AdhocDataframeBuilder::CreateIntegerStorage( diff --git a/src/trace_processor/core/dataframe/adhoc_dataframe_builder.h b/src/trace_processor/core/dataframe/adhoc_dataframe_builder.h index 98c3fabe284..353dc51d457 100644 --- a/src/trace_processor/core/dataframe/adhoc_dataframe_builder.h +++ b/src/trace_processor/core/dataframe/adhoc_dataframe_builder.h @@ -117,6 +117,12 @@ struct AdhocDataframeBuilderOptions { // resulting dataframe is consumed somewhere that supplies its own primary // key (e.g. `StaticTableFunctionModule` adds a HIDDEN `_auto_id`). bool emit_auto_id = true; + + // If false, `Build()` keeps columns as they were collected: integers stay + // Int64, nothing is marked sorted or free of duplicates, and the dataframe + // is finalized without statistics for query planning. Cheaper, for a + // dataframe which is only ever scanned. + bool analyze = true; }; class AdhocDataframeBuilder { @@ -474,6 +480,7 @@ class AdhocDataframeBuilder { std::vector column_states_; bool did_declare_types_ = false; bool emit_auto_id_ = true; + bool analyze_ = true; base::Status current_status_ = base::OkStatus(); core::BitVector duplicate_bit_vector_; }; diff --git a/src/trace_processor/core/dataframe/adhoc_dataframe_builder_unittest.cc b/src/trace_processor/core/dataframe/adhoc_dataframe_builder_unittest.cc index f793b6ceacf..9a1da354aab 100644 --- a/src/trace_processor/core/dataframe/adhoc_dataframe_builder_unittest.cc +++ b/src/trace_processor/core/dataframe/adhoc_dataframe_builder_unittest.cc @@ -76,6 +76,29 @@ TEST_F(AdhocDataframeBuilderTest, StringColumnWithNullId) { ColumnSpec{Id{}, NonNull{}, IdSorted{}, NoDuplicates{}})); } +// Without analysis, columns are kept as collected: integers stay Int64 and +// nothing is claimed about their order or duplicates. +TEST_F(AdhocDataframeBuilderTest, WithoutAnalysisColumnsStayAsCollected) { + AdhocDataframeBuilder::Options options; + options.emit_auto_id = false; + options.analyze = false; + AdhocDataframeBuilder builder({"id", "value"}, &pool_, options); + for (int64_t i = 0; i < 3; ++i) { + builder.PushNonNull(0, i); + builder.PushNonNull(1, 2.5); + } + base::StatusOr df = std::move(builder).Build(); + ASSERT_OK(df.status()); + EXPECT_EQ(df->row_count(), 3u); + // Finalized all the same, so its columns can be shared. + EXPECT_TRUE(df->finalized()); + EXPECT_THAT( + df->CreateSpec().column_specs, + ElementsAre( + ColumnSpec{Int64{}, NonNull{}, Unsorted{}, HasDuplicates{}}, + ColumnSpec{Double{}, NonNull{}, Unsorted{}, HasDuplicates{}})); +} + // Callback for reading cell values in tests. struct TestCellCallback : CellCallback { void OnCell(int64_t v) { diff --git a/src/trace_processor/core/dataframe/dataframe.cc b/src/trace_processor/core/dataframe/dataframe.cc index c141f15350c..f8f289608de 100644 --- a/src/trace_processor/core/dataframe/dataframe.cc +++ b/src/trace_processor/core/dataframe/dataframe.cc @@ -218,6 +218,14 @@ Dataframe Dataframe::RemoveIndexAt(uint32_t pos) const { } void Dataframe::Finalize() { + FinalizeColumns(/*estimate_distinct=*/true); +} + +void Dataframe::FinalizeWithoutStatistics() { + FinalizeColumns(/*estimate_distinct=*/false); +} + +void Dataframe::FinalizeColumns(bool estimate_distinct) { if (finalized_) { return; } @@ -273,7 +281,9 @@ void Dataframe::Finalize() { default: PERFETTO_FATAL("Invalid nullability type"); } - c->estimated_distinct = EstimateDistinct(distinct_counts, *c); + if (estimate_distinct) { + c->estimated_distinct = EstimateDistinct(distinct_counts, *c); + } } // Bump the mutation counter so that any cursors with cached pointers // know to refresh them: shrink_to_fit() may have reallocated the internal diff --git a/src/trace_processor/core/dataframe/dataframe.h b/src/trace_processor/core/dataframe/dataframe.h index 333db8db602..c25eb94bd31 100644 --- a/src/trace_processor/core/dataframe/dataframe.h +++ b/src/trace_processor/core/dataframe/dataframe.h @@ -264,6 +264,11 @@ class Dataframe { // If the dataframe is already finalized, this function does nothing. void Finalize(); + // As Finalize, but without estimating how many distinct values each column + // holds, which only query planning uses. Cheaper, for a dataframe which is + // only ever scanned. + void FinalizeWithoutStatistics(); + // Makes a copy of the dataframe which has been finalized. Unfinalized // dataframes *cannot* be copied, so this function will assert if not // finalized. @@ -752,6 +757,9 @@ class Dataframe { uint32_t non_column_mutations_ = 0; // Whether the dataframe is "finalized". See `Finalize()`. + // Finalize, estimating distinct counts only if `estimate_distinct`. + void FinalizeColumns(bool estimate_distinct); + bool finalized_ = false; }; diff --git a/src/trace_processor/perfetto_sql/engine/connection_catalog.cc b/src/trace_processor/perfetto_sql/engine/connection_catalog.cc index 8cf1357d0d5..d20677f6ef1 100644 --- a/src/trace_processor/perfetto_sql/engine/connection_catalog.cc +++ b/src/trace_processor/perfetto_sql/engine/connection_catalog.cc @@ -22,17 +22,15 @@ #include #include -#include "perfetto/ext/base/status_or.h" #include "perfetto/ext/base/string_utils.h" #include "src/perfetto_sql/analysis/relation.h" #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/type_mapping.h" #include "src/trace_processor/sqlite/bindings/sqlite_column.h" #include "src/trace_processor/sqlite/sql_source.h" #include "src/trace_processor/sqlite/sqlite_connection.h" +#include "src/trace_processor/sqlite/sqlite_utils.h" namespace perfetto::trace_processor { namespace { @@ -77,7 +75,19 @@ std::optional 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; + } + analysis::LeafRelation relation{std::string(name), {}}; + for (const auto& column : sqlite::utils::GetColumns( + connection_->sqlite_connection()->db(), relation.name)) { + relation.columns.push_back({column.name, std::nullopt, column.hidden}); + } + if (relation.columns.empty()) { + return std::nullopt; + } + return relation; } analysis::LeafRelation relation; relation.name = name; @@ -109,10 +119,4 @@ const dataframe::Dataframe* ConnectionCatalog::FindDataframe( return connection_->GetDataframeOrNull(name); } -base::StatusOr ConnectionCatalog::DescribeQuery( - const SqlSource& sql) const { - return sql_schema::DescribeQuery(connection_->sqlite_connection(), sql, - *this); -} - } // namespace perfetto::trace_processor diff --git a/src/trace_processor/perfetto_sql/engine/connection_catalog.h b/src/trace_processor/perfetto_sql/engine/connection_catalog.h index d343a7a3e71..fb841e20bcc 100644 --- a/src/trace_processor/perfetto_sql/engine/connection_catalog.h +++ b/src/trace_processor/perfetto_sql/engine/connection_catalog.h @@ -21,13 +21,10 @@ #include #include -#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/engine/perfetto_sql_connection.h" #include "src/trace_processor/perfetto_sql/pipeline/catalog.h" -#include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" -#include "src/trace_processor/sqlite/sql_source.h" namespace perfetto::trace_processor { @@ -43,8 +40,6 @@ class ConnectionCatalog final : public pipeline::Catalog { const dataframe::Dataframe* FindDataframe( std::string_view name) const override; - base::StatusOr DescribeQuery( - const SqlSource& sql) const override; private: PerfettoSqlConnection* connection_; diff --git a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.cc b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.cc index cfa7433413a..4e752261c02 100644 --- a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.cc +++ b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.cc @@ -364,6 +364,12 @@ PerfettoSqlConnection::PerfettoSqlConnection( ctx->connection = this; RegisterVirtualTableModule(pipeline::kPipelineFunction, std::move(ctx)); + base::Status status = RegisterAggregateFunction(pool_); + PERFETTO_CHECK(status.ok()); + // Not deterministic: each call makes a new list. + status = RegisterFunction( + nullptr, RegisterFunctionArgs(nullptr, /*deterministic=*/false)); + PERFETTO_CHECK(status.ok()); } database_->InitializeSharedSchema(connection_.get()); @@ -1057,19 +1063,22 @@ 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> -PerfettoSqlConnection::LoadPipeline(std::string_view serialized) { +PerfettoSqlConnection::LoadPipeline( + std::string_view serialized, + const std::vector& args) { 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 pipeline::Lower(plan, env); + RETURN_IF_ERROR(pipeline::BindDataframeArgs(plan, args, pool_)); + return pipeline::Lower(plan); } base::Status PerfettoSqlConnection::ExecuteCreateView( diff --git a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.h b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.h index 97f7ee0f04d..ca1927bd725 100644 --- a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.h +++ b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection.h @@ -174,9 +174,11 @@ class PerfettoSqlConnection { base::StatusOr PrepareSqliteStatement( SqlSource sql); - // Loads a plan written by pipeline::SerializePlan, ready to run. + // Loads a plan written by pipeline::SerializePlan, ready to run, reading + // `args` as its dataframe arguments. The plan keeps what it reads of them. base::StatusOr> LoadPipeline( - std::string_view serialized); + std::string_view serialized, + const std::vector& args); // Registers a virtual table module with the given name. // diff --git a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc index dbfeeaa637f..536309a3c42 100644 --- a/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc +++ b/src/trace_processor/perfetto_sql/engine/perfetto_sql_connection_unittest.cc @@ -853,7 +853,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(); } @@ -999,6 +1001,10 @@ TEST_F(PerfettoSqlConnectionPipelineTest, Errors) { .status() .message(), testing::HasSubstr("syntax error near 'WHERE'")); + // Semantic analysis cannot yet describe every relation. + EXPECT_THAT(Rows("FROM (VALUES (1, 2))").status().message(), + testing::HasSubstr( + "reading a relation whose columns cannot be worked out")); } // Replacing a table mid-read must not affect a running pipeline. @@ -1080,6 +1086,82 @@ TEST_F(PerfettoSqlConnectionPipelineTest, BadPlans) { EXPECT_THAT(rows.status().message(), testing::HasSubstr("no such column")); } +// SQLite runs the SQL a pipeline reads as part of the statement the pipeline +// is written in, and passes the pipeline a dataframe of its rows. +TEST_F(PerfettoSqlConnectionPipelineTest, SqlSourcesAreDataframeArgs) { + std::string sql = PipelineSql(R"( + FROM (SELECT id, parent_id, self FROM tree) + |> TREE ACCUMULATE UP SUM(self) AS total + )"); + EXPECT_THAT( + sql, + testing::HasSubstr( + R"((SELECT __intrinsic_dataframes((SELECT __intrinsic_dataframe_agg('id,parent_id,self', "id", "parent_id", "self") FROM (SELECT id, parent_id, self FROM tree)))))")); + + // A relation with no rows too. + auto rows = Rows(R"( + FROM (SELECT id, parent_id, self FROM tree WHERE id > 10) + |> TREE ACCUMULATE UP SUM(self) AS total + )"); + ASSERT_TRUE(rows.ok()) << rows.status().message(); + EXPECT_TRUE(rows->empty()); +} + +// Anyone can write SQL running a plan, so the dataframes it is passed are +// checked. +TEST_F(PerfettoSqlConnectionPipelineTest, BadDataframeArgs) { + std::string sql = PipelineSql("FROM (SELECT 1 AS x)"); + size_t start = sql.find("X'"); + std::string plan = sql.substr(start, sql.find('\'', start + 2) + 1 - start); + auto run = [&](const std::string& args) { + return Rows("SELECT c0 FROM __intrinsic_pipeline(" + plan + args + ")") + .status(); + }; + EXPECT_THAT(run("").message(), testing::HasSubstr("no dataframe argument 0")); + EXPECT_THAT(run(", 1").message(), testing::HasSubstr("expected dataframes")); + EXPECT_THAT(run(", __intrinsic_dataframes(1)").message(), + testing::HasSubstr("expected dataframes")); + EXPECT_THAT(run(", __intrinsic_dataframes((SELECT " + "__intrinsic_dataframe_agg('y', 1)))") + .message(), + testing::HasSubstr("dataframe argument 0 has no column 'x'")); + EXPECT_THAT(run(", __intrinsic_dataframes(" + "(SELECT __intrinsic_dataframe_agg('x,y', 1)))") + .message(), + testing::HasSubstr("expected 2 values, not 1")); + EXPECT_THAT(run(R"(, __intrinsic_dataframes( + (SELECT __intrinsic_dataframe_agg('x', v) + FROM (SELECT 1 AS v UNION ALL SELECT 'a'))))") + .message(), + testing::HasSubstr("column 'x' was inferred to be")); + EXPECT_THAT(run(", __intrinsic_dataframes(" + "(SELECT __intrinsic_dataframe_agg('x', X'00')))") + .message(), + testing::HasSubstr("holds a blob")); + auto rows = Rows( + "SELECT c0 FROM __intrinsic_pipeline(" + plan + + ", __intrinsic_dataframes((SELECT __intrinsic_dataframe_agg('x', 7))))"); + ASSERT_TRUE(rows.ok()) << rows.status().message(); + EXPECT_THAT(*rows, testing::ElementsAre("7")); + // A relation with no rows is passed as NULL. + EXPECT_TRUE(run(", __intrinsic_dataframes(NULL)").ok()); +} + +// A cursor keeps the dataframes it loaded, so SQL passing it others later +// cannot free them from under it. +TEST_F(PerfettoSqlConnectionPipelineTest, NewDataframeArgsAreSafe) { + std::string sql = PipelineSql("FROM (SELECT 1 AS x)"); + size_t start = sql.find("X'"); + std::string plan = sql.substr(start, sql.find('\'', start + 2) + 1 - start); + auto rows = + Rows("SELECT v, (SELECT c0 FROM __intrinsic_pipeline(" + plan + + ", __intrinsic_dataframes((SELECT __intrinsic_dataframe_agg('x', y)" + " FROM (SELECT v AS y)))))" + " FROM (SELECT 1 AS v UNION ALL SELECT 2 UNION ALL SELECT 3)"); + ASSERT_TRUE(rows.ok()) << rows.status().message(); + EXPECT_EQ(rows->size(), 3u); +} + TEST_F(PerfettoSqlConnectionPipelineTest, ForksRunPipelinesIndependently) { auto fork = connection_->Fork(); auto first = connection_->ExecuteUntilLastStatement( @@ -1140,14 +1222,17 @@ TEST_F(PerfettoSqlConnectionPipelineTest, ColumnReadersRefreshAcrossBatches) { } EXPECT_TRUE(res->stmt.status().ok()); EXPECT_EQ(count, 5000u); - // A dynamically typed column must continue dispatching per value, including - // a type change in a later batch. - auto mixed = Rows( - "FROM (SELECT CASE WHEN id < 3000 THEN id ELSE 'last' END AS value FROM " - "values_table)"); - ASSERT_TRUE(mixed.ok()) << mixed.status().message(); - EXPECT_EQ(mixed->size(), 5000u); - EXPECT_EQ(std::count(mixed->begin(), mixed->end(), "last"), 2000); + // A pipeline reads SQL as a dataframe, whose columns each hold one type, so + // a column changing type part way through is refused. + EXPECT_THAT(Rows(R"( + FROM ( + SELECT CASE WHEN id < 3000 THEN id ELSE 'last' END AS value + FROM values_table + ) + )") + .status() + .message(), + testing::HasSubstr("column 'value' was inferred to be")); } } // namespace diff --git a/src/trace_processor/perfetto_sql/engine/pipeline_module.cc b/src/trace_processor/perfetto_sql/engine/pipeline_module.cc index fd3f04852b9..1f99beab853 100644 --- a/src/trace_processor/perfetto_sql/engine/pipeline_module.cc +++ b/src/trace_processor/perfetto_sql/engine/pipeline_module.cc @@ -18,6 +18,7 @@ #include +#include #include #include #include @@ -29,17 +30,22 @@ #include "perfetto/base/compiler.h" #include "perfetto/base/logging.h" #include "perfetto/base/status.h" +#include "perfetto/ext/base/status_or.h" #include "perfetto/ext/base/string_utils.h" #include "src/trace_processor/containers/string_pool.h" #include "src/trace_processor/core/common/storage_types.h" +#include "src/trace_processor/core/dataframe/dataframe.h" +#include "src/trace_processor/core/dataframe/runtime_dataframe_builder.h" #include "src/trace_processor/core/exec/column_view.h" #include "src/trace_processor/core/exec/row_cursor.h" #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/engine/sqlite_dataframe_builder.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/sqlite/bindings/sqlite_result.h" +#include "src/trace_processor/sqlite/bindings/sqlite_value.h" #include "src/trace_processor/sqlite/sqlite_utils.h" namespace perfetto::trace_processor { @@ -48,14 +54,33 @@ namespace { using core::exec::ColumnView; using core::exec::Variant; -// The plan comes before the outputs: SQLite only says which of a table's +// The pointer types of what __intrinsic_dataframe_agg and +// __intrinsic_dataframes give. +constexpr char kDataframePointerType[] = "PIPELINE_DATAFRAME"; +constexpr char kDataframesPointerType[] = "PIPELINE_DATAFRAMES"; + +// Shared, so a cursor keeps the dataframes it loaded alive. +using SharedDataframe = std::shared_ptr; + +// What __intrinsic_dataframes gives: the dataframes, in the order the plan +// numbers them, null for a relation with no rows. +struct DataframeList { + std::vector dataframes; +}; + +// The table function's arguments: the plan, then the list of dataframes it +// reads. They come before the outputs: SQLite only says which of a table's // first 63 columns a query reads, and those are best spent on the outputs. constexpr int kPlanColumn = 0; -constexpr int kFirstOutputColumn = 1; +constexpr int kDataframesColumn = 1; +constexpr int kFirstOutputColumn = 2; + +// idxNum: whether dataframes are given. +constexpr int kHasDataframes = 1; std::string Schema() { // Public names (which may repeat) are applied by the outer SELECT. - std::vector columns{"pipeline HIDDEN"}; + std::vector columns{"pipeline HIDDEN", "dataframes HIDDEN"}; for (uint32_t i = 0; i < pipeline::kMaxPipelineColumns; ++i) { columns.push_back("c" + std::to_string(i)); } @@ -179,16 +204,33 @@ int CheckStatus(PipelineModule::Cursor* cursor) { : sqlite::utils::SetError(cursor->pVtab, status); } -// The slow path of Filter: loads the plan in `value` into `c`. -PERFETTO_NO_INLINE int Load(PipelineModule::Cursor* c, sqlite3_value* value) { - if (sqlite3_value_type(value) != SQLITE_BLOB) { +// The slow path of Filter: loads the plan in `value`, reading `args`, into +// `c`. +PERFETTO_NO_INLINE int Load(PipelineModule::Cursor* c, + int idx_num, + sqlite3_value** argv) { + if (sqlite3_value_type(argv[0]) != SQLITE_BLOB) { return sqlite::utils::SetError(c->pVtab, "__intrinsic_pipeline: expected a plan"); } + if (idx_num & kHasDataframes) { + const auto* list = + sqlite::value::Pointer(argv[1], kDataframesPointerType); + if (!list) { + return sqlite::utils::SetError( + c->pVtab, "__intrinsic_pipeline: expected dataframes"); + } + c->inputs = list->dataframes; + } + std::vector inputs; + for (const SharedDataframe& input : c->inputs) { + inputs.push_back(input.get()); + } PipelineModule::Context* context = PipelineModule::GetVtab(c->pVtab)->context; auto plan = context->connection->LoadPipeline( - std::string_view(static_cast(sqlite3_value_blob(value)), - static_cast(sqlite3_value_bytes(value)))); + std::string_view(static_cast(sqlite3_value_blob(argv[0])), + static_cast(sqlite3_value_bytes(argv[0]))), + inputs); if (!plan.ok()) { return sqlite::utils::SetError(c->pVtab, plan.status()); } @@ -207,6 +249,68 @@ PERFETTO_NO_INLINE int Load(PipelineModule::Cursor* c, sqlite3_value* value) { } // namespace +void DataframesFunction::Step(sqlite3_context* ctx, + int argc, + sqlite3_value** argv) { + auto list = std::make_unique(); + list->dataframes.reserve(static_cast(argc)); + for (int i = 0; i < argc; ++i) { + // A pointer reads as NULL, so only a NULL without one is no rows. + const auto* dataframe = + sqlite::value::Pointer(argv[i], kDataframePointerType); + if (!dataframe && !sqlite::value::IsNull(argv[i])) { + return sqlite::utils::SetError( + ctx, base::ErrStatus("%s: expected dataframes", kName)); + } + list->dataframes.push_back(dataframe ? *dataframe : nullptr); + } + return sqlite::result::UniquePointer(ctx, std::move(list), + kDataframesPointerType); +} + +void DataframeAgg::Step(sqlite3_context* ctx, int argc, sqlite3_value** argv) { + if (argc < 1) { + return sqlite::utils::SetError( + ctx, base::ErrStatus("%s: expected column names", kName)); + } + AggCtx& agg = AggCtx::GetOrCreateContextForStep(ctx); + auto count = static_cast(argc - 1); + if (!agg.builder) { + const char* text = sqlite::value::Text(argv[0]); + std::vector names = base::SplitString(text ? text : "", ","); + if (names.size() != count) { + return sqlite::utils::SetError( + ctx, base::ErrStatus("%s: expected %zu values, not %u", kName, + names.size(), count)); + } + dataframe::RuntimeDataframeBuilder::Options options; + options.emit_auto_id = false; + options.analyze = false; + agg.builder.emplace(std::move(names), GetUserData(ctx), options); + } + base::Status status = AddSqliteValuesRow(*agg.builder, argv + 1, count); + if (!status.ok()) { + return sqlite::utils::SetError(ctx, kName, status); + } +} + +void DataframeAgg::Final(sqlite3_context* ctx) { + auto agg = AggCtx::GetContextOrNullForFinal(ctx); + if (!agg.get() || !agg.get()->builder) { + return sqlite::result::Null(ctx); + } + base::StatusOr dataframe = + std::move(*agg.get()->builder).Build(); + if (!dataframe.ok()) { + return sqlite::result::Error(ctx, dataframe.status().c_message()); + } + return sqlite::result::UniquePointer( + ctx, + std::make_unique( + std::make_shared(std::move(*dataframe))), + kDataframePointerType); +} + int PipelineModule::Connect(sqlite3* db, void* raw_ctx, int, @@ -230,24 +334,38 @@ int PipelineModule::Disconnect(sqlite3_vtab* vtab) { int PipelineModule::BestIndex(sqlite3_vtab*, sqlite3_index_info* info) { int plan = -1; + int dataframes = -1; for (int i = 0; i < info->nConstraint; ++i) { const auto& constraint = info->aConstraint[i]; if (constraint.op != SQLITE_INDEX_CONSTRAINT_EQ) { continue; } + // Without the plan and its arguments there is nothing to run. if (constraint.iColumn == kPlanColumn) { - // Without the plan there is nothing to run. if (!constraint.usable) { return SQLITE_CONSTRAINT; } plan = i; } + if (constraint.iColumn == kDataframesColumn) { + if (!constraint.usable) { + return SQLITE_CONSTRAINT; + } + dataframes = i; + } } if (plan == -1) { return SQLITE_CONSTRAINT; } - info->aConstraintUsage[plan].argvIndex = 1; + int argc = 0; + info->aConstraintUsage[plan].argvIndex = ++argc; info->aConstraintUsage[plan].omit = true; + info->idxNum = 0; + if (dataframes != -1) { + info->aConstraintUsage[dataframes].argvIndex = ++argc; + info->aConstraintUsage[dataframes].omit = true; + info->idxNum |= kHasDataframes; + } info->estimatedCost = 1e9; return SQLITE_OK; } @@ -263,15 +381,15 @@ int PipelineModule::Close(sqlite3_vtab_cursor* cursor) { } int PipelineModule::Filter(sqlite3_vtab_cursor* cursor, - int, + int idx_num, const char*, - int argc, + int, sqlite3_value** argv) { Cursor* c = GetCursor(cursor); - PERFETTO_DCHECK(argc == 1); - // The plan is a constant, so it is loaded once per cursor. + // The plan and its dataframes are the same on every filter, so it is loaded + // once per cursor. if (PERFETTO_UNLIKELY(!c->plan)) { - if (int rc = Load(c, argv[0]); rc != SQLITE_OK) { + if (int rc = Load(c, idx_num, argv); rc != SQLITE_OK) { return rc; } } diff --git a/src/trace_processor/perfetto_sql/engine/pipeline_module.h b/src/trace_processor/perfetto_sql/engine/pipeline_module.h index f162779b144..c581ed8ff5e 100644 --- a/src/trace_processor/perfetto_sql/engine/pipeline_module.h +++ b/src/trace_processor/perfetto_sql/engine/pipeline_module.h @@ -21,22 +21,59 @@ #include #include +#include #include #include #include "src/trace_processor/containers/string_pool.h" +#include "src/trace_processor/core/dataframe/dataframe.h" +#include "src/trace_processor/core/dataframe/runtime_dataframe_builder.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/perfetto_sql/pipeline/pipeline_sql.h" +#include "src/trace_processor/sqlite/bindings/sqlite_aggregate_function.h" +#include "src/trace_processor/sqlite/bindings/sqlite_function.h" #include "src/trace_processor/sqlite/bindings/sqlite_module.h" namespace perfetto::trace_processor { class PerfettoSqlConnection; +// `__intrinsic_dataframe_agg('a,b', a, b)`: builds a dataframe of a relation, +// as a PERFETTO TABLE is built, for a pipeline to read. Gives NULL for a +// relation with no rows, as the aggregate then never sees its column names. +struct DataframeAgg : sqlite::AggregateFunction { + static constexpr const char* kName = pipeline::kDataframeAggFunction; + static constexpr int kArgCount = -1; + using UserData = StringPool; + + struct AggCtx : sqlite::AggregateContext { + std::optional builder; + }; + + static void Step(sqlite3_context*, int argc, sqlite3_value** argv); + static void Final(sqlite3_context*); +}; + +// `__intrinsic_dataframes(df, ...)`: gathers the dataframes a pipeline reads +// into one list, so the table function takes them all as one argument. The +// list only points at them: it is read while the arguments are still alive. +struct DataframesFunction : sqlite::Function { + static constexpr const char* kName = pipeline::kDataframesFunction; + static constexpr int kArgCount = -1; + + static void Step(sqlite3_context*, int argc, sqlite3_value** argv); +}; + // Runs a pipeline from its plan, serialized into the SQL which reads it: // `__intrinsic_pipeline(X'...')`. The plan is all a pipeline needs, so the SQL // can be stored, in a view say, and run later. Pipelines output into fixed // columns `c0`, `c1`, ..., which the SQL reading them renames. +// +// The relations a pipeline reads from SQL are passed as a list of dataframes +// after the plan: `__intrinsic_pipeline(X'...', __intrinsic_dataframes((SELECT +// __intrinsic_dataframe_agg(...) FROM ...), ...))`. struct PipelineModule : sqlite::Module { static constexpr auto kType = kEponymousOnly; static constexpr bool kSupportsWrites = false; @@ -60,6 +97,8 @@ struct PipelineModule : sqlite::Module { ResultFn result; }; std::unique_ptr plan; + // The dataframes `plan` reads, kept alive for it. + std::vector> inputs; StringPool* pool = nullptr; std::unique_ptr rows; // By declared column. diff --git a/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.cc b/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.cc index 684070a0faa..b6902cca412 100644 --- a/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.cc +++ b/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.cc @@ -24,6 +24,7 @@ #include "src/trace_processor/core/dataframe/runtime_dataframe_builder.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" namespace perfetto::trace_processor { namespace { @@ -57,6 +58,32 @@ struct SqliteValueFetcher : public dataframe::ValueFetcher { bool blobs_as_null; }; +struct SqliteValuesFetcher : public dataframe::ValueFetcher { + explicit SqliteValuesFetcher(sqlite3_value** row) : values(row) {} + + using Type = sqlite::Type; + static constexpr Type kInt64 = sqlite::Type::kInteger; + static constexpr Type kDouble = sqlite::Type::kFloat; + static constexpr Type kString = sqlite::Type::kText; + static constexpr Type kNull = sqlite::Type::kNull; + static constexpr Type kBytes = sqlite::Type::kBlob; + + int64_t GetInt64Value(uint32_t column) const { + return sqlite::value::Int64(values[column]); + } + double GetDoubleValue(uint32_t column) const { + return sqlite::value::Double(values[column]); + } + const char* GetStringValue(uint32_t column) const { + return sqlite::value::Text(values[column]); + } + Type GetValueType(uint32_t column) const { + return sqlite::value::Type(values[column]); + } + + sqlite3_value** values; +}; + } // namespace base::StatusOr @@ -82,4 +109,20 @@ BuildRuntimeDataframeFromSqliteStatement( return std::move(builder); } +base::Status AddSqliteValuesRow(dataframe::RuntimeDataframeBuilder& builder, + sqlite3_value** values, + uint32_t count) { + for (uint32_t i = 0; i < count; ++i) { + if (sqlite::value::Type(values[i]) == sqlite::Type::kBlob) { + return base::ErrStatus("column %u holds a blob, which cannot be read", + i + 1); + } + } + SqliteValuesFetcher fetcher(values); + if (!builder.AddRow(&fetcher)) { + return builder.status(); + } + return base::OkStatus(); +} + } // namespace perfetto::trace_processor diff --git a/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.h b/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.h index 51e3bebe790..0e442eaa154 100644 --- a/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.h +++ b/src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.h @@ -23,6 +23,7 @@ #include #include +#include "perfetto/base/status.h" #include "perfetto/ext/base/status_or.h" #include "src/trace_processor/containers/string_pool.h" #include "src/trace_processor/core/dataframe/adhoc_dataframe_builder.h" @@ -48,6 +49,13 @@ BuildRuntimeDataframeFromSqliteStatement( std::string_view error_context, SqliteDataframeBuilderOptions options = {}); +// Adds a row of `values`, one for each column, to `builder`. Fails on a blob, +// which a dataframe cannot hold, and on a value of a different type to the +// rest of its column. +base::Status AddSqliteValuesRow(dataframe::RuntimeDataframeBuilder& builder, + sqlite3_value** values, + uint32_t count); + } // namespace perfetto::trace_processor #endif // SRC_TRACE_PROCESSOR_PERFETTO_SQL_ENGINE_SQLITE_DATAFRAME_BUILDER_H_ diff --git a/src/trace_processor/perfetto_sql/exec/BUILD.gn b/src/trace_processor/perfetto_sql/exec/BUILD.gn deleted file mode 100644 index d559857494d..00000000000 --- a/src/trace_processor/perfetto_sql/exec/BUILD.gn +++ /dev/null @@ -1,51 +0,0 @@ -# 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. - -import("../../../../gn/test.gni") - -source_set("exec") { - sources = [ - "sql_scan.cc", - "sql_scan.h", - ] - deps = [ - "../../../../gn:default_deps", - "../../../../gn:sqlite", - "../../../base", - "../../containers", - "../../core/common", - "../../core/exec", - "../../core/util", - "../../sqlite", - ] -} - -perfetto_unittest_source_set("unittests") { - testonly = true - sources = [ "sql_scan_unittest.cc" ] - deps = [ - ":exec", - "../../../../gn:default_deps", - "../../../../gn:gtest_and_gmock", - "../../../base", - "../../../perfetto_sql/analysis", - "../../containers", - "../../core/common", - "../../core/exec", - "../../core/exec:test_utils", - "../../core/util", - "../../sqlite", - "../schema", - ] -} diff --git a/src/trace_processor/perfetto_sql/exec/sql_scan.cc b/src/trace_processor/perfetto_sql/exec/sql_scan.cc deleted file mode 100644 index 7f1bef06570..00000000000 --- a/src/trace_processor/perfetto_sql/exec/sql_scan.cc +++ /dev/null @@ -1,283 +0,0 @@ -/* - * 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/exec/sql_scan.h" - -#include - -#include -#include -#include -#include -#include -#include -#include - -#include "perfetto/base/compiler.h" -#include "perfetto/base/logging.h" -#include "perfetto/base/status.h" -#include "src/trace_processor/containers/string_pool.h" -#include "src/trace_processor/core/common/storage_types.h" -#include "src/trace_processor/core/exec/column_chunk.h" -#include "src/trace_processor/core/exec/column_view.h" -#include "src/trace_processor/core/exec/operator.h" -#include "src/trace_processor/core/exec/row_batch.h" -#include "src/trace_processor/core/exec/row_selection.h" -#include "src/trace_processor/core/exec/variant.h" -#include "src/trace_processor/core/util/bit_vector.h" -#include "src/trace_processor/sqlite/bindings/sqlite_column.h" -#include "src/trace_processor/sqlite/bindings/sqlite_type.h" -#include "src/trace_processor/sqlite/sql_source.h" -#include "src/trace_processor/sqlite/sqlite_connection.h" - -namespace perfetto::trace_processor::exec { -namespace { - -using core::Double; -using core::Int64; -using core::StorageType; -using core::String; -using core::exec::ColumnChunk; -using core::exec::ColumnView; -using core::exec::kMaxBatchRows; -using core::exec::RowBatch; -using core::exec::RowSelection; -using core::exec::Variant; -} // namespace - -SqlScan::SqlScan(SqliteConnection* connection, - SqlSource sql, - core::Schema columns, - StringPool* pool) - : connection_(connection), - sql_(std::move(sql)), - columns_(std::move(columns)), - pool_(pool) {} - -SqlScan::~SqlScan() = default; -SqlScan::State::~State() = default; - -std::unique_ptr SqlScan::MakeState() const { - auto state = std::make_unique(); - Prepare(*state); - return state; -} - -void SqlScan::PrepareColumns(State& state) const { - state.columns.clear(); - state.data.clear(); - state.buffers.resize(columns_.size()); - state.columns.reserve(columns_.size()); - state.data.reserve(columns_.size()); - for (uint32_t i = 0; i < columns_.size(); ++i) { - auto column = state.buffers[i].Acquire(); - void* data = nullptr; - if (!columns_[i].type) { - data = column->Values().data(); - } else { - switch (columns_[i].type->index()) { - case StorageType::GetTypeIndex(): - data = column->Values().data(); - break; - case StorageType::GetTypeIndex(): - data = column->Values().data(); - break; - case StorageType::GetTypeIndex(): - data = column->Values().data(); - break; - case StorageType::GetTypeIndex(): - data = column->Values().data(); - break; - case StorageType::GetTypeIndex(): - data = column->Values().data(); - break; - default: - // An Id was already materialised as a Uint32 by ResolveTypes. - PERFETTO_FATAL("Unreachable"); - } - column->validity.resize(kMaxBatchRows); - } - state.columns.push_back(std::move(column)); - state.data.push_back(data); - } -} - -void SqlScan::Prepare(State& state) const { - state.statement.emplace(connection_->PrepareStatement(sql_)); - state.status = state.statement->status(); - state.done = false; - if (!state.status.ok()) { - return; - } - sqlite3_stmt* stmt = state.statement->sqlite_stmt(); - uint32_t count = sqlite::column::Count(stmt); - if (count != columns_.size()) { - state.status = - base::ErrStatus("SQL source: result shape changed between executions"); - return; - } - for (uint32_t i = 0; i < count; ++i) { - const char* name = sqlite::column::Name(stmt, i); - if (columns_[i].name != (name ? name : "")) { - state.status = base::ErrStatus( - "SQL source: result shape changed between executions"); - return; - } - } -} - -base::Status SqlScan::status(const core::exec::OperatorState& state) const { - return state.Cast().status; -} - -void SqlScan::Rewind(core::exec::OperatorState& state) const { - Prepare(state.Cast()); -} - -bool SqlScan::ReadValue(State& s, - sqlite3_stmt* stmt, - uint32_t index, - uint32_t row) const { - if (!columns_[index].type) { - auto* data = static_cast(s.data[index]); - switch (sqlite::column::Type(stmt, index)) { - case sqlite::Type::kInteger: - data[row] = Variant::Int64(sqlite::column::Int64(stmt, index)); - return true; - case sqlite::Type::kFloat: - data[row] = Variant::Double(sqlite::column::Double(stmt, index)); - return true; - case sqlite::Type::kText: - data[row] = Variant::String( - pool_->InternString(sqlite::column::Text(stmt, index))); - return true; - case sqlite::Type::kNull: - data[row] = Variant::Null(); - return true; - case sqlite::Type::kBlob: - s.status = base::ErrStatus( - "SQL source: column '%s' holds a blob, which a pipeline cannot " - "carry", - columns_[index].name.c_str()); - return false; - } - PERFETTO_FATAL("For GCC"); - } - switch (columns_[index].type->index()) { - case StorageType::GetTypeIndex(): - return ReadTypedValue(s, stmt, index, - row); - case StorageType::GetTypeIndex(): - return ReadTypedValue(s, stmt, index, - row); - case StorageType::GetTypeIndex(): - return ReadTypedValue(s, stmt, index, - row); - case StorageType::GetTypeIndex(): - return ReadTypedValue(s, stmt, index, row); - case StorageType::GetTypeIndex(): - return ReadTypedValue(s, stmt, index, - row); - default: - // An Id was already materialised as a Uint32 by ResolveTypes. - PERFETTO_FATAL("Unreachable"); - } -} - -template -bool SqlScan::ReadTypedValue(State& s, - sqlite3_stmt* stmt, - uint32_t index, - uint32_t row) const { - auto* data = static_cast(s.data[index]); - sqlite::Type type = sqlite::column::Type(stmt, index); - if (PERFETTO_LIKELY(type == SqliteType)) { - if constexpr (std::is_same_v) { - data[row] = pool_->InternString(sqlite::column::Text(stmt, index)); - } else if constexpr (std::is_same_v) { - data[row] = sqlite::column::Double(stmt, index); - } else { - data[row] = static_cast(sqlite::column::Int64(stmt, index)); - } - s.columns[index]->validity.set(row); - return true; - } - if (type == sqlite::Type::kNull) { - // The row is null, but write the slot anyway. A flat column's storage is - // readable at every row, so a reader summing it needs no per-row branch and - // never sees a value left over from the previous batch. - if constexpr (std::is_same_v) { - data[row] = StringPool::Id::Null(); - } else { - data[row] = T{}; - } - return true; - } - // Only reachable if the type lineage established turned out to be wrong. - s.status = base::ErrStatus( - "SQL source: column '%s' does not hold what it was traced back to", - columns_[index].name.c_str()); - return false; -} - -bool SqlScan::GetData(RowBatch& out, core::exec::OperatorState& state) const { - State& s = state.Cast(); - if (s.done || !s.status.ok()) { - return false; - } - out.Reset(); - PrepareColumns(s); - for (const std::shared_ptr& column : s.columns) { - if (column->validity.size() != 0) { - column->validity.ClearAllBits(); - } - } - sqlite3_stmt* stmt = s.statement->sqlite_stmt(); - uint32_t count = 0; - while (count < kMaxBatchRows && s.statement->Step()) { - for (uint32_t i = 0; i < s.columns.size(); ++i) { - if (!ReadValue(s, stmt, i, count)) { - return false; - } - } - ++count; - } - if (!s.statement->status().ok()) { - s.status = s.statement->status(); - return false; - } - s.done = count < kMaxBatchRows; - if (count == 0) { - return false; - } - - for (uint32_t i = 0; i < s.columns.size(); ++i) { - const std::shared_ptr& column = s.columns[i]; - if (!columns_[i].type) { - out.AddColumn( - ColumnView::Variants(static_cast(s.data[i])), column); - } else { - out.AddColumn(ColumnView::Reference(*columns_[i].type, s.data[i], - &column->validity), - column); - } - } - out.Compose(RowSelection::Range(0), count); - out.SetCardinality(count); - return true; -} - -} // namespace perfetto::trace_processor::exec diff --git a/src/trace_processor/perfetto_sql/exec/sql_scan.h b/src/trace_processor/perfetto_sql/exec/sql_scan.h deleted file mode 100644 index 1cd11d95750..00000000000 --- a/src/trace_processor/perfetto_sql/exec/sql_scan.h +++ /dev/null @@ -1,102 +0,0 @@ -/* - * 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_EXEC_SQL_SCAN_H_ -#define SRC_TRACE_PROCESSOR_PERFETTO_SQL_EXEC_SQL_SCAN_H_ - -#include -#include -#include -#include -#include - -#include "perfetto/base/status.h" -#include "src/trace_processor/containers/string_pool.h" -#include "src/trace_processor/core/common/schema.h" -#include "src/trace_processor/core/exec/buffer_pool.h" -#include "src/trace_processor/core/exec/column_chunk.h" -#include "src/trace_processor/core/exec/operator.h" -#include "src/trace_processor/core/exec/row_batch.h" -#include "src/trace_processor/sqlite/bindings/sqlite_type.h" -#include "src/trace_processor/sqlite/sql_source.h" -#include "src/trace_processor/sqlite/sqlite_connection.h" - -struct sqlite3_stmt; - -namespace perfetto::trace_processor::exec { - -// Reads a pipeline's rows from a SQL query. -// -// Promises nothing about the order the rows arrive in, because SQLite does -// not. -// -// Each column carries its type per row unless the query can be traced back to -// a dataframe column, which is the only way to establish a type. SQLite's -// declared types establish nothing: an INTEGER column holds text if something -// puts text in it. -class SqlScan : public core::exec::Source { - public: - // A scan over `sql` with its previously resolved result columns. - SqlScan(SqliteConnection*, SqlSource, core::Schema, StringPool*); - ~SqlScan() override; - - // The query's columns, in the order a batch carries them. - const core::Schema& columns() const { return columns_; } - - // The type of column `i`, or nothing when the column carries a type per - // row. - std::optional column_type(uint32_t i) const { - return columns_[i].type; - } - - std::unique_ptr MakeState() const override; - bool GetData(core::exec::RowBatch& out, - core::exec::OperatorState&) const override; - void Rewind(core::exec::OperatorState&) const override; - base::Status status(const core::exec::OperatorState&) const override; - - private: - struct State : core::exec::OperatorState { - ~State() override; - std::optional statement; - // Shared so a batch can keep the values alive. - std::vector> columns; - std::vector> buffers; - // Each column's value buffer, resolved out of its chunk once. - std::vector data; - bool done = false; - base::Status status = base::OkStatus(); - }; - - void Prepare(State&) const; - void PrepareColumns(State&) const; - - bool ReadValue(State&, sqlite3_stmt*, uint32_t index, uint32_t row) const; - template - bool ReadTypedValue(State&, - sqlite3_stmt*, - uint32_t index, - uint32_t row) const; - - SqliteConnection* connection_; - SqlSource sql_; - core::Schema columns_; - StringPool* pool_; -}; - -} // namespace perfetto::trace_processor::exec - -#endif // SRC_TRACE_PROCESSOR_PERFETTO_SQL_EXEC_SQL_SCAN_H_ diff --git a/src/trace_processor/perfetto_sql/exec/sql_scan_unittest.cc b/src/trace_processor/perfetto_sql/exec/sql_scan_unittest.cc deleted file mode 100644 index 7f6c2b6770e..00000000000 --- a/src/trace_processor/perfetto_sql/exec/sql_scan_unittest.cc +++ /dev/null @@ -1,532 +0,0 @@ -/* - * 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/exec/sql_scan.h" - -#include -#include -#include -#include -#include -#include -#include -#include - -#include "perfetto/ext/base/status_macros.h" -#include "src/trace_processor/containers/string_pool.h" -#include "src/trace_processor/core/common/storage_types.h" -#include "src/trace_processor/core/exec/assert_type.h" -#include "src/trace_processor/core/exec/column_view.h" -#include "src/trace_processor/core/exec/operator.h" -#include "src/trace_processor/core/exec/pipeline.h" -#include "src/trace_processor/core/exec/row_batch.h" -#include "src/trace_processor/core/exec/row_cursor.h" -#include "src/trace_processor/core/exec/row_selection.h" -#include "src/trace_processor/core/exec/test_utils.h" -#include "src/trace_processor/core/exec/tree_accumulate.h" -#include "src/trace_processor/core/exec/tree_number_nodes.h" -#include "src/trace_processor/core/exec/tree_order.h" -#include "src/trace_processor/core/exec/variant.h" -#include "src/trace_processor/core/util/bit_vector.h" -#include "src/trace_processor/perfetto_sql/schema/query_schema.h" -#include "src/trace_processor/perfetto_sql/schema/type_mapping.h" -#include "src/trace_processor/sqlite/sql_source.h" -#include "src/trace_processor/sqlite/sqlite_connection.h" -#include "test/gtest_and_gmock.h" - -namespace perfetto::trace_processor::exec { -namespace { - -namespace analysis = ::perfetto::perfetto_sql::analysis; - -using core::BitVector; -using core::Double; -using core::Int64; -using core::StorageType; -using core::String; -using core::exec::ColumnView; -using core::exec::kMaxBatchRows; -using core::exec::RowBatch; -using core::exec::RowCursor; -using core::exec::Variant; - -// Drives a plan the way an executor does: creates the state, owns the batch. -class Execution { - public: - explicit Execution(const core::exec::Source& source) - : source_(source), state_(source.MakeState()) {} - - RowBatch* Next() { - return source_.GetData(batch_, *state_) ? &batch_ : nullptr; - } - void Rewind() { source_.Rewind(*state_); } - base::Status status() const { return source_.status(*state_); } - - private: - const core::exec::Source& source_; - std::unique_ptr state_; - RowBatch batch_; -}; - -std::vector ReadInts(const RowBatch& batch, uint32_t index) { - std::vector out; - for (const Variant& cell : - core::exec::test::ReadColumn(batch, index)) { - out.push_back(cell.AsInt64()); - } - return out; -} - -struct TestColumn { - std::string name; - core::StorageType type; -}; - -TestColumn Typed(std::string name, core::StorageType type) { - return {std::move(name), type}; -} - -class TestCatalog : public analysis::Catalog { - public: - void Add(std::string name, std::vector columns) { - dataframes_[std::move(name)] = std::move(columns); - } - - std::optional FindLeafRelation( - std::string_view name) const override { - auto dataframe = dataframes_.find(std::string(name)); - if (dataframe == dataframes_.end()) { - return std::nullopt; - } - analysis::LeafRelation relation; - relation.name = name; - for (const TestColumn& column : dataframe->second) { - relation.columns.push_back( - {column.name, sql_schema::ToAnalysisType(column.type)}); - } - return relation; - } - - std::optional FindViewSql(std::string_view) const override { - return std::nullopt; - } - - private: - std::map> dataframes_; -}; - -class SqlScanTest : public ::testing::Test { - protected: - SqlScanTest() - : connection_(SqliteConnection::CreateConnectionToNewDatabase()) {} - - void Exec(const std::string& sql) { - auto statement = - connection_->PrepareStatement(SqlSource::FromExecuteQuery(sql)); - ASSERT_TRUE(statement.status().ok()) << statement.status().c_message(); - while (statement.Step()) { - } - ASSERT_TRUE(statement.status().ok()) << statement.status().c_message(); - } - - base::StatusOr> Scan( - const std::string& sql, - const analysis::Catalog& catalog) { - auto source = SqlSource::FromExecuteQuery(sql); - ASSIGN_OR_RETURN(auto columns, sql_schema::DescribeQuery(connection_.get(), - source, catalog)); - return std::make_unique(connection_.get(), std::move(source), - std::move(columns), &pool_); - } - - // An empty catalog traces nothing, so every column is a variant. - base::StatusOr> Scan(const std::string& sql) { - return Scan(sql, empty_catalog_); - } - - StringPool pool_; - std::unique_ptr connection_; - TestCatalog empty_catalog_; -}; - -TEST_F(SqlScanTest, AQuerysColumnsAreKnownBeforeItsRows) { - auto scan = Scan("SELECT 1 AS a, 'x' AS b"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - ASSERT_EQ((*scan)->columns().size(), 2u); - EXPECT_EQ((*scan)->columns()[0].name, "a"); - EXPECT_EQ((*scan)->columns()[1].name, "b"); -} - -TEST_F(SqlScanTest, AQuerysRowsArriveAsABatch) { - auto scan = Scan( - "SELECT a FROM (" - "SELECT 2 AS ord, 20 AS a UNION ALL " - "SELECT 3, 30 UNION ALL SELECT 1, 10) ORDER BY ord"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - EXPECT_EQ(batch->size(), 3u); - EXPECT_THAT(ReadInts(*batch, 0), testing::ElementsAre(10, 20, 30)); - EXPECT_EQ(run.Next(), nullptr); - EXPECT_TRUE(run.status().ok()); -} - -// One column holding three different types and a null, which SQLite allows. -TEST_F(SqlScanTest, OneColumnCanHoldMoreThanOneType) { - auto scan = Scan( - "SELECT value FROM (" - "SELECT 1 AS ord, 7 AS value UNION ALL SELECT 2, 1.5 " - "UNION ALL SELECT 3, 'hello' UNION ALL SELECT 4, NULL) ORDER BY ord"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - std::vector cells = core::exec::test::ReadColumn(*batch, 0); - ASSERT_EQ(cells.size(), 4u); - EXPECT_EQ(cells[0].AsInt64(), 7); - EXPECT_EQ(cells[1].AsDouble(), 1.5); - EXPECT_EQ(pool_.Get(cells[2].AsString()).ToStdString(), "hello"); - EXPECT_EQ(cells[3].type, Variant::Type::kNull); - EXPECT_TRUE(run.status().ok()); -} - -// A declared type is not binding in SQLite, so it is not trusted here. -TEST_F(SqlScanTest, ADeclaredTypeIsNotBelieved) { - Exec("CREATE TABLE t(i INTEGER)"); - Exec("INSERT INTO t VALUES(1), ('not a number')"); - auto scan = Scan("SELECT i FROM t"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - std::vector cells = core::exec::test::ReadColumn(*batch, 0); - ASSERT_EQ(cells.size(), 2u); - EXPECT_EQ(cells[0].AsInt64(), 1); - EXPECT_EQ(pool_.Get(cells[1].AsString()).ToStdString(), "not a number"); - EXPECT_TRUE(run.status().ok()); -} - -TEST_F(SqlScanTest, AColumnWhichIsNeverAnythingIsAColumnOfNulls) { - auto scan = Scan("SELECT NULL AS a UNION ALL SELECT NULL"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - for (const Variant& cell : core::exec::test::ReadColumn(*batch, 0)) { - EXPECT_EQ(cell.type, Variant::Type::kNull); - } -} - -TEST_F(SqlScanTest, MoreRowsThanFitInABatchArriveInSeveral) { - Exec( - "CREATE TABLE t AS WITH RECURSIVE r(x) AS (" - " SELECT 0 UNION ALL SELECT x + 1 FROM r WHERE x < 4999" - ") SELECT x FROM r"); - auto scan = Scan("SELECT x FROM t ORDER BY x"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - std::vector seen; - std::vector sizes; - while (RowBatch* batch = run.Next()) { - sizes.push_back(batch->size()); - for (int64_t value : ReadInts(*batch, 0)) { - seen.push_back(value); - } - } - ASSERT_TRUE(run.status().ok()) << run.status().c_message(); - EXPECT_THAT(sizes, testing::ElementsAre(kMaxBatchRows, kMaxBatchRows, 904u)); - ASSERT_EQ(seen.size(), 5000u); - for (uint32_t i = 0; i < seen.size(); ++i) { - ASSERT_EQ(seen[i], int64_t{i}); - } -} - -TEST_F(SqlScanTest, AQueryWhichDoesNotRunIsReported) { - auto scan = Scan("SELECT * FROM not_a_table"); - EXPECT_FALSE(scan.ok()); -} - -TEST_F(SqlScanTest, AQueryWhichFailsPartWayThroughIsReported) { - auto scan = - Scan("SELECT 1 AS a UNION ALL SELECT abs(-9223372036854775807 - 1)"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - while (run.Next()) { - } - EXPECT_FALSE(run.status().ok()); -} - -TEST_F(SqlScanTest, ABlobIsReportedRatherThanCarried) { - auto scan = Scan("SELECT x'0102' AS a"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - EXPECT_EQ(run.Next(), nullptr); - EXPECT_FALSE(run.status().ok()); - EXPECT_THAT(run.status().message(), testing::HasSubstr("blob")); -} - -TEST_F(SqlScanTest, RewindClearsThePreviousExecutionError) { - auto scan = Scan("SELECT x'0102' AS a"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - ASSERT_EQ(run.Next(), nullptr); - ASSERT_FALSE(run.status().ok()); - - run.Rewind(); - EXPECT_TRUE(run.status().ok()); - EXPECT_EQ(run.Next(), nullptr); - EXPECT_FALSE(run.status().ok()); -} - -TEST_F(SqlScanTest, EachExecutionValidatesItsResultShape) { - Exec("CREATE TABLE t(a INTEGER)"); - Exec("INSERT INTO t VALUES(1)"); - auto scan = Scan("SELECT * FROM t"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Exec("ALTER TABLE t ADD COLUMN b TEXT"); - - Execution run(**scan); - EXPECT_EQ(run.Next(), nullptr); - EXPECT_FALSE(run.status().ok()); - EXPECT_THAT(run.status().message(), testing::HasSubstr("shape changed")); -} - -TEST_F(SqlScanTest, AScanCanBeRunAgain) { - auto scan = Scan("SELECT 1 AS a UNION ALL SELECT 2"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - - auto drain = [&] { - std::vector out; - while (RowBatch* batch = run.Next()) { - for (int64_t value : ReadInts(*batch, 0)) { - out.push_back(value); - } - } - return out; - }; - std::vector first = drain(); - run.Rewind(); - EXPECT_EQ(drain(), first); -} - -TEST_F(SqlScanTest, AQueryReachesARowCursor) { - auto scan = Scan( - "SELECT a FROM (SELECT 2 AS ord, 6 AS a UNION ALL " - "SELECT 3, 7 UNION ALL SELECT 1, 5) ORDER BY ord"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - - RowCursor cursor(**scan); - std::vector values; - for (bool more = cursor.Open(); more; more = cursor.Next()) { - values.push_back(cursor.Value(0).AsInt64()); - } - EXPECT_THAT(values, testing::ElementsAre(5, 6, 7)); -} - -// The whole pipeline: a query, its columns asserted to be integers, put into a -// tree order and folded up the tree. -TEST_F(SqlScanTest, AQueryReachesTheTreeOperators) { - Exec( - "CREATE TABLE t AS " - "SELECT 0 AS id, NULL AS parent_id, 10 AS self " - "UNION ALL SELECT 1, 0, 20 " - "UNION ALL SELECT 2, 0, 30 " - "UNION ALL SELECT 3, 1, 40"); - auto scan = Scan("SELECT id, parent_id, self FROM t ORDER BY id"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - - std::vector> ops; - ops.push_back(std::make_unique( - 0, core::exec::AssertTypeTarget{core::Int64{}}, "id")); - ops.push_back(std::make_unique( - 1, core::exec::AssertTypeTarget{core::Int64{}}, "parent_id")); - ops.push_back(std::make_unique( - 2, core::exec::AssertTypeTarget{core::Int64{}}, "self")); - ops.push_back(std::make_unique(0, 1)); - ops.push_back(std::make_unique(3, 4)); - core::exec::TreeAccumulateSpec spec{3, 4, 2}; - ops.push_back(std::make_unique(spec)); - core::exec::Pipeline folded(**scan, std::move(ops), {}); - - std::unique_ptr state = folded.MakeState(); - RowBatch batch; - std::vector totals(4, 0); - while (folded.GetData(batch, *state)) { - std::vector ids = core::exec::test::ReadColumn(batch, 0); - std::vector values = - core::exec::test::ReadColumn(batch, 5); - for (uint32_t row = 0; row < batch.size(); ++row) { - totals[static_cast(ids[row])] = values[row]; - } - } - ASSERT_TRUE(folded.status(*state).ok()) << folded.status(*state).message(); - EXPECT_THAT(totals, testing::ElementsAre(100, 60, 30, 40)); -} - -// A column which is not the type claimed for it fails the pipeline. -TEST_F(SqlScanTest, AColumnWhichIsNotWhatWasAssertedIsReported) { - Exec("CREATE TABLE t(i INTEGER)"); - Exec("INSERT INTO t VALUES(1), ('not a number')"); - auto scan = Scan("SELECT i FROM t"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - - std::vector> ops; - ops.push_back(std::make_unique( - 0, core::exec::AssertTypeTarget{core::Int64{}}, "i")); - core::exec::Pipeline typed(**scan, std::move(ops), {}); - - Execution run(typed); - while (run.Next()) { - } - EXPECT_FALSE(run.status().ok()); - EXPECT_THAT(run.status().message(), testing::HasSubstr("'i'")); -} - -// A column which can be traced back to a dataframe needs neither a variant nor -// an assertion: it comes out flat. -TEST_F(SqlScanTest, AColumnFollowedBackToADataframeComesOutFlat) { - Exec("CREATE TABLE df(id INTEGER, name TEXT)"); - Exec("INSERT INTO df VALUES(7, 'hello'), (8, NULL)"); - TestCatalog catalog; - catalog.Add("df", {Typed("id", core::StorageType{core::Int64{}}), - Typed("name", core::StorageType{core::String{}})}); - - auto scan = Scan("SELECT id, name FROM df ORDER BY id", catalog); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - ASSERT_TRUE((*scan)->column_type(0).has_value()); - EXPECT_TRUE((*scan)->column_type(0)->Is()); - EXPECT_TRUE((*scan)->column_type(1)->Is()); - - Execution run(**scan); - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - EXPECT_EQ(batch->column(0).kind(), ColumnView::Kind::kFlat); - const auto* ids = static_cast(batch->column(0).data()); - EXPECT_EQ(ids[0], 7); - EXPECT_EQ(ids[1], 8); - const BitVector* validity = batch->column(1).validity(); - ASSERT_NE(validity, nullptr); - EXPECT_TRUE(validity->is_set(0)); - EXPECT_FALSE(validity->is_set(1)); - std::vector names = - core::exec::test::ReadColumn(*batch, 1); - EXPECT_EQ(pool_.Get(names[0]).ToStdString(), "hello"); -} - -TEST_F(SqlScanTest, MixedNumericCompoundResultsStayVariants) { - Exec("CREATE TABLE ints(value INTEGER)"); - Exec("INSERT INTO ints VALUES(7)"); - Exec("CREATE TABLE doubles(value REAL)"); - Exec("INSERT INTO doubles VALUES(1.5)"); - TestCatalog catalog; - catalog.Add("ints", {Typed("value", StorageType{Int64{}})}); - catalog.Add("doubles", {Typed("value", StorageType{Double{}})}); - - auto scan = Scan( - "SELECT value FROM (" - "SELECT 1 AS ord, value FROM ints UNION ALL " - "SELECT 2, value FROM doubles) ORDER BY ord", - catalog); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - EXPECT_FALSE((*scan)->column_type(0).has_value()); - Execution run(**scan); - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - std::vector values = - core::exec::test::ReadColumn(*batch, 0); - ASSERT_EQ(values.size(), 2u); - EXPECT_EQ(values[0].AsInt64(), 7); - EXPECT_EQ(values[1].AsDouble(), 1.5); -} - -// An expression cannot be traced back, so it stays a variant even when the -// column beside it does not. -TEST_F(SqlScanTest, OnlyTheColumnsWhichCanBeFollowedComeOutFlat) { - Exec("CREATE TABLE df(id INTEGER)"); - Exec("INSERT INTO df VALUES(7)"); - TestCatalog catalog; - catalog.Add("df", {Typed("id", core::StorageType{core::Int64{}})}); - - auto scan = Scan("SELECT id, id * 2 AS doubled FROM df", catalog); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - EXPECT_TRUE((*scan)->column_type(0).has_value()); - EXPECT_FALSE((*scan)->column_type(1).has_value()); - - Execution run(**scan); - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - EXPECT_EQ(batch->column(0).kind(), ColumnView::Kind::kFlat); - EXPECT_EQ(batch->column(1).kind(), ColumnView::Kind::kVariant); -} - -// Without a catalog nothing can be traced back. -TEST_F(SqlScanTest, WithoutACatalogEveryColumnIsAVariant) { - Exec("CREATE TABLE df(id INTEGER)"); - auto scan = Scan("SELECT id FROM df"); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - EXPECT_FALSE((*scan)->column_type(0).has_value()); -} - -// A flat column's storage is readable at every row, so a reader which sums it -// without checking validity gets zero rather than a value left over from the -// previous batch. -TEST_F(SqlScanTest, ANullSlotOfAFlatColumnHoldsZero) { - Exec("CREATE TABLE df(id INTEGER)"); - Exec("INSERT INTO df VALUES(7), (NULL)"); - TestCatalog catalog; - catalog.Add("df", {Typed("id", core::StorageType{core::Int64{}})}); - - auto scan = Scan("SELECT id FROM df ORDER BY id IS NULL", catalog); - ASSERT_TRUE(scan.ok()) << scan.status().c_message(); - Execution run(**scan); - RowBatch* batch = run.Next(); - ASSERT_NE(batch, nullptr); - ASSERT_EQ(batch->size(), 2u); - - const auto* ids = static_cast(batch->column(0).data()); - EXPECT_EQ(ids[0], 7); - EXPECT_EQ(ids[1], 0); - EXPECT_FALSE(batch->column(0).validity()->is_set(1)); -} - -TEST_F(SqlScanTest, RetainedBatchSurvivesAdvanceAndRewind) { - auto scan = Scan( - "WITH RECURSIVE n(x) AS (SELECT 0 UNION ALL " - "SELECT x+1 FROM n WHERE x<4096) SELECT x FROM n"); - ASSERT_TRUE(scan.ok()) << scan.status().message(); - auto state = (*scan)->MakeState(); - RowBatch output, retained; - ASSERT_TRUE((*scan)->GetData(output, *state)); - retained.CopyFrom(output); - ASSERT_TRUE((*scan)->GetData(output, *state)); - (*scan)->Rewind(*state); - ASSERT_TRUE((*scan)->GetData(output, *state)); - auto values = ReadInts(retained, 0); - ASSERT_EQ(values.size(), kMaxBatchRows); - for (uint32_t i = 0; i < values.size(); ++i) - EXPECT_EQ(values[i], i); -} - -} // namespace -} // namespace perfetto::trace_processor::exec diff --git a/src/trace_processor/perfetto_sql/parser/perfetto_sql_parser_unittest.cc b/src/trace_processor/perfetto_sql/parser/perfetto_sql_parser_unittest.cc index 715a1971124..100c07b361e 100644 --- a/src/trace_processor/perfetto_sql/parser/perfetto_sql_parser_unittest.cc +++ b/src/trace_processor/perfetto_sql/parser/perfetto_sql_parser_unittest.cc @@ -31,6 +31,7 @@ #include "src/trace_processor/perfetto_sql/parser/function_util.h" #include "src/trace_processor/perfetto_sql/parser/perfetto_sql_test_utils.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" +#include "src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h" #include "src/trace_processor/perfetto_sql/pipeline/test_catalog.h" #include "src/trace_processor/sqlite/sql_source.h" #include "src/trace_processor/sqlite/sqlite_connection.h" @@ -615,7 +616,7 @@ TEST_F(PerfettoSqlParserTest, PipelineExpandsMacros) { "FROM (SELECT * FROM tree) |> TREE ACCUMULATE UP SUM(self) AS " "total"); EXPECT_THAT(pipeline::LogicalPlanToString(pipeline->plan), - HasSubstr("Scan(sql SELECT * FROM (SELECT * FROM tree))")); + HasSubstr("Scan(sql (SELECT * FROM tree))")); } TEST_F(PerfettoSqlParserTest, CreatePerfettoTableAsPipeline) { @@ -655,7 +656,7 @@ TEST_F(PerfettoSqlParserTest, PipelineSyntaxErrors) { TEST_F(PerfettoSqlParserTest, PipelineCompileErrors) { EXPECT_THAT(ParsePipeline("FROM nope").status().message(), - HasSubstr("no such table: nope")); + HasSubstr("'nope' is not known")); EXPECT_THAT(ParsePipeline("FROM (SELECT id, self FROM tree) |> TREE " "ACCUMULATE UP SUM(self) AS total") .status() @@ -848,27 +849,23 @@ TEST_F(PerfettoSqlParserTest, PipelineSkipsUnusedTreeAggregates) { } TEST_F(PerfettoSqlParserTest, PipelinePushesPruningIntoSql) { - // SQL sources only ask SQLite for the columns that are used. - auto plan = ParsePipeline( + // The dataframe built for each SQL source holds only the columns used. + PerfettoSqlParser parser(macros_, catalog_, /*pipelines_allowed=*/true); + parser.Reset(SqlSource::FromExecuteQuery( "INTERVAL INTERSECTION OF (sql_spans AS a, (SELECT * FROM sql_spans) AS " - "b) |> SELECT ts, dur, b.utid"); - ASSERT_TRUE(plan.ok()) << plan.status().message(); - EXPECT_THAT(*plan, - HasSubstr("Scan(sql WITH __pipeline_source(c0, c1, c2, c3) AS " - "(SELECT * FROM sql_spans AS a) SELECT c0 AS \"ts\", " - "c1 AS \"dur\" FROM __pipeline_source) " - "[#2 AS ts, #3 AS dur]")); + "b) |> SELECT ts, dur, b.utid")); + ASSERT_TRUE(parser.Next()) << parser.status().message(); + auto sql = + pipeline::SelectPipelineSql(std::get(parser.statement()).plan); + ASSERT_TRUE(sql.ok()) << sql.status().message(); EXPECT_THAT( - *plan, - HasSubstr("Scan(sql WITH __pipeline_source(c0, c1, c2, c3) AS " - "(SELECT * FROM (SELECT * FROM sql_spans) AS b) SELECT c0 AS " - "\"ts\", c1 AS \"dur\", c3 AS \"utid\" FROM __pipeline_source) " - "[#6 AS ts, #7 AS dur, #9 AS utid]")); - - // A source that uses every column is left alone. - plan = ParsePipeline("FROM tree |> TREE ACCUMULATE UP SUM(self) AS total"); - ASSERT_TRUE(plan.ok()) << plan.status().message(); - EXPECT_THAT(*plan, HasSubstr("Scan(sql SELECT * FROM tree)")); + *sql, + HasSubstr( + R"((SELECT __intrinsic_dataframe_agg('ts,dur', "ts", "dur") FROM sql_spans AS a))")); + EXPECT_THAT( + *sql, + HasSubstr( + R"((SELECT __intrinsic_dataframe_agg('ts,dur,utid', "ts", "dur", "utid") FROM (SELECT * FROM sql_spans) AS b))")); } // The relational operators which reshape a pipeline's row: EXTEND, DROP, @@ -1082,29 +1079,27 @@ TEST_F(PerfettoSqlParserSelectLikeTest, SourceNames) { Check({ {"FROM (SELECT 1 AS x, count(*) AS n FROM tree) AS t |> SELECT *", "Output(#0 AS x, #1 AS n)"}, - // SQLite names an unaliased expression after its text. + // An expression is only named by its alias. {"FROM (SELECT 1 AS x, 1 + 1 FROM tree) AS t |> SELECT x", - "expected every column to have a valid name, but '1 + 1' is not one: " - "give it one with AS"}, + "expected every column to have a name, but column 2 has none: give it " + "one with AS"}, {"FROM (SELECT 1 AS \"my col\") AS t", kInvalid}, {"FROM (SELECT 1 AS \"1x\") AS t", kInvalid}, // Without an earlier x, x:1 is simply not a valid name. {"FROM (SELECT 1 AS \"x:1\") AS t", kInvalid}, {"INTERVAL INTERSECTION OF ((SELECT 0 AS ts, 10 AS dur, 1 + 1) AS a, " "(SELECT 5 AS ts, 10 AS dur) AS b)", - kInvalid}, + "expected every column to have a name, but column 3 has none"}, }); } -// SQLite renames the second of two columns sharing a name `x` to `x:1`. That -// is reported as the duplicate it is, before the names are checked. +// Names are compared as SQL compares them, ignoring case. TEST_F(PerfettoSqlParserSelectLikeTest, SourceDuplicateNames) { Check({ {"FROM (SELECT 1 AS x, 2 AS x) AS t", - "expected distinct column names, but there are two named 'x', which " - "SQLite renamed to 'x' and 'x:1'"}, + "expected distinct column names, but there are two named 'x'"}, {"FROM (SELECT 1 AS x, 2 AS X) AS t", - "two named 'x', which SQLite renamed to 'x' and 'X:1'"}, + "expected distinct column names, but there are two named 'x'"}, {"INTERVAL INTERSECTION OF (" "(SELECT 0 AS ts, 10 AS dur, 1 AS x, 2 AS x) AS a, " "(SELECT 5 AS ts, 10 AS dur) AS b)", diff --git a/src/trace_processor/perfetto_sql/pipeline/BUILD.gn b/src/trace_processor/perfetto_sql/pipeline/BUILD.gn index 84da7184395..8d7edc96f48 100644 --- a/src/trace_processor/perfetto_sql/pipeline/BUILD.gn +++ b/src/trace_processor/perfetto_sql/pipeline/BUILD.gn @@ -36,10 +36,12 @@ source_set("logical") { "../../../base", "../../../perfetto_sql/analysis", "../../../perfetto_sql/syntaqlite", + "../../containers", "../../core/common", "../../core/dataframe", "../../sqlite", "../../util:sql_argument", + "../schema", ] } @@ -52,12 +54,9 @@ source_set("plan") { deps = [ "../../../../gn:default_deps", "../../../base", - "../../containers", "../../core/common", "../../core/dataframe", "../../core/exec", - "../../sqlite", - "../exec", ] } @@ -72,6 +71,7 @@ source_set("test_catalog") { "../../containers", "../../core/dataframe", "../../sqlite", + "../engine", "../schema", ] } diff --git a/src/trace_processor/perfetto_sql/pipeline/catalog.h b/src/trace_processor/perfetto_sql/pipeline/catalog.h index 7dee828b010..a128a68b780 100644 --- a/src/trace_processor/perfetto_sql/pipeline/catalog.h +++ b/src/trace_processor/perfetto_sql/pipeline/catalog.h @@ -19,11 +19,8 @@ #include -#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 { @@ -37,10 +34,6 @@ class Catalog : public perfetto_sql::analysis::Catalog { // Dataframe registered as `name`, or null. virtual const dataframe::Dataframe* FindDataframe( std::string_view name) const = 0; - - // Columns of the query `sql`, typed where they trace back to a dataframe - // column. Fails if SQLite cannot prepare the query. - virtual base::StatusOr DescribeQuery(const SqlSource& sql) const = 0; }; } // namespace perfetto::trace_processor::pipeline diff --git a/src/trace_processor/perfetto_sql/pipeline/column_pruning.cc b/src/trace_processor/perfetto_sql/pipeline/column_pruning.cc index d63e3563468..24a71844ea0 100644 --- a/src/trace_processor/perfetto_sql/pipeline/column_pruning.cc +++ b/src/trace_processor/perfetto_sql/pipeline/column_pruning.cc @@ -25,7 +25,6 @@ #include #include "perfetto/base/logging.h" -#include "perfetto/ext/base/string_utils.h" #include "perfetto/ext/base/variant.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" #include "src/trace_processor/sqlite/sql_source.h" @@ -36,35 +35,6 @@ namespace { // The columns something downstream uses, indexed by column ID. using Needed = std::vector; -// Narrows `sql` down to the columns at the positions in `keep`, which keep -// their names from `kept`. We pick by position rather than name because a -// query can return two columns with the same name. The scan checks the names -// it reads back, so each column is given its name again. -SqlSource SelectPositions(const SqlSource& sql, - uint32_t count, - const std::vector& keep, - const std::vector& kept) { - std::string names; - for (uint32_t i = 0; i < count; ++i) { - if (i) { - names += ", "; - } - names += "c" + std::to_string(i); - } - std::string selected; - for (uint32_t i = 0; i < keep.size(); ++i) { - if (i) { - selected += ", "; - } - selected += "c" + std::to_string(keep[i]) + " AS \"" + - base::ReplaceAll(kept[i].name, "\"", "\"\"") + "\""; - } - return sql.RewriteAllIgnoreExisting( - SqlSource::FromTraceProcessorImplementation( - "WITH __pipeline_source(" + names + ") AS (" + sql.sql() + - ") SELECT " + selected + " FROM __pipeline_source")); -} - void PruneScan(op::Scan& scan, const Needed& needed) { std::vector keep; for (uint32_t i = 0; i < scan.columns.size(); ++i) { @@ -79,7 +49,6 @@ void PruneScan(op::Scan& scan, const Needed& needed) { if (keep.empty()) { keep.push_back(0); } - auto count = static_cast(scan.columns.size()); std::vector columns; for (uint32_t i : keep) { columns.push_back(std::move(scan.columns[i])); @@ -96,11 +65,10 @@ void PruneScan(op::Scan& scan, const Needed& needed) { dataframe.columns = std::move(kept); return; } - case base::variant_index(): { - auto& sql = base::unchecked_get(scan.source); - sql = SelectPositions(sql, count, keep, scan.columns); + case base::variant_index(): + // The SQL building a dataframe of the relation reads only the columns + // the scan keeps. return; - } default: PERFETTO_FATAL("Unknown scan source"); } diff --git a/src/trace_processor/perfetto_sql/pipeline/compiler.cc b/src/trace_processor/perfetto_sql/pipeline/compiler.cc index f953589a1ce..4b7c02c6658 100644 --- a/src/trace_processor/perfetto_sql/pipeline/compiler.cc +++ b/src/trace_processor/perfetto_sql/pipeline/compiler.cc @@ -31,18 +31,23 @@ #include "perfetto/ext/base/status_or.h" #include "perfetto/ext/base/string_utils.h" #include "perfetto/ext/base/string_view.h" +#include "src/perfetto_sql/analysis/relation.h" #include "src/perfetto_sql/syntaqlite/syntaqlite_perfetto.h" #include "src/trace_processor/core/common/storage_types.h" #include "src/trace_processor/core/dataframe/dataframe.h" #include "src/trace_processor/perfetto_sql/pipeline/catalog.h" #include "src/trace_processor/perfetto_sql/pipeline/column_pruning.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" +#include "src/trace_processor/perfetto_sql/schema/type_mapping.h" #include "src/trace_processor/sqlite/sql_source.h" #include "src/trace_processor/util/sql_argument.h" namespace perfetto::trace_processor::pipeline { namespace { +namespace analysis = ::perfetto::perfetto_sql::analysis; +using core::StorageType; + // The name `span` spells. A quoted name escapes its closing quote by doubling // it, which the span, pointing into the source, still contains. std::string SpanText(SyntaqliteParser* p, SyntaqliteTextSpan span) { @@ -327,28 +332,25 @@ std::optional Compiler::SourceQualifier( base::Status Compiler::CheckSourceNames(const op::Scan& scan, uint32_t at) const { - // SQLite renames the second of two columns sharing a name `x` to `x:1`. - // Say so, rather than only that `x:1` is not a valid name. for (size_t i = 0; i < scan.columns.size(); ++i) { - const std::string& name = scan.columns[i].name; - size_t colon = name.rfind(':'); - if (colon == std::string::npos || colon + 1 == name.size() || - name.find_first_not_of("0123456789", colon + 1) != std::string::npos) { - continue; - } for (size_t j = 0; j < i; ++j) { - const std::string& first = scan.columns[j].name; - if (base::CaseInsensitiveEqual(first, name.substr(0, colon))) { + if (base::CaseInsensitiveEqual(scan.columns[j].name, + scan.columns[i].name)) { return Expected(at, "distinct column names, but there are two named '" + - first + "', which SQLite renamed to '" + first + - "' and '" + name + "'"); + scan.columns[j].name + "'"); } } } - for (const NamedColumn& column : scan.columns) { - if (!sql_argument::IsValidColumnName(base::StringView(column.name))) { - return Expected(at, "every column to have a valid name, but '" + - column.name + + for (size_t i = 0; i < scan.columns.size(); ++i) { + const std::string& name = scan.columns[i].name; + // An expression is only named by its alias. + if (name.empty()) { + return Expected(at, "every column to have a name, but column " + + std::to_string(i + 1) + + " has none: give it one with AS"); + } + if (!sql_argument::IsValidColumnName(base::StringView(name))) { + return Expected(at, "every column to have a valid name, but '" + name + "' is not one: give it one with AS"); } } @@ -375,19 +377,36 @@ op::Scan Compiler::CompileDataframeSource(const dataframe::Dataframe& dataframe, } base::StatusOr Compiler::CompileSqlSource(uint32_t from) { - SqlSource sql = source_(from); - sql = - sql.RewriteAllIgnoreExisting(SqlSource::FromTraceProcessorImplementation( - "SELECT * FROM " + sql.sql())); - auto described = catalog_.DescribeQuery(sql); - if (!described.ok()) { - return base::ErrStatus("%s%s", Traceback(from).c_str(), - described.status().c_message()); + // The columns come from semantic analysis alone: the SQL is only run where + // the pipeline is written, which is the only place everything it reads is + // in scope. + const auto* n = Node(p_, from); + analysis::RelationAnalyzer analyzer(catalog_); + base::StatusOr lineage = + syntaqlite_node_is_present(n->select) + ? analyzer.AnalyzeQuery({p_, n->select}) + : analyzer.AnalyzeRelation(SpanText(p_, n->table_name)); + if (!lineage.ok() || IsPresent(n->schema)) { + std::string reason = + lineage.ok() ? "a schema-qualified table" : lineage.status().message(); + return Err(from, Error::kUnsupported, + "reading a relation whose columns cannot be worked out", + " (" + reason + ")"); } op::Scan scan; - scan.source = std::move(sql); - for (ColumnSchema& column : *described) { - AddScanColumn(scan, std::move(column)); + // The relation as written, which the SQL building its dataframe reads. + scan.source = source_(from); + for (const analysis::ColumnLineage& column : lineage->columns()) { + std::optional type; + if (std::optional traced = column.type()) { + type = sql_schema::ToStorageType(*traced); + // An Id only means something in the table it indexes: read through + // SQL, it is a plain number. + if (type->Is()) { + type = StorageType{core::Uint32{}}; + } + } + AddScanColumn(scan, {std::string(column.output_name), type}); } return scan; } diff --git a/src/trace_processor/perfetto_sql/pipeline/logical_plan.h b/src/trace_processor/perfetto_sql/pipeline/logical_plan.h index 94062b3b597..7d3817b7c0d 100644 --- a/src/trace_processor/perfetto_sql/pipeline/logical_plan.h +++ b/src/trace_processor/perfetto_sql/pipeline/logical_plan.h @@ -55,11 +55,20 @@ struct ScanDataframe { uint32_t row_count = 0; }; +// The dataframe the plan is passed as its `index`-th argument when it runs. +// SQLite builds it from a relation the plan read from SQL. +struct ScanDataframeArg { + uint32_t index = 0; +}; + // Reads all rows of a source. Always the first op. struct Scan { using Dataframe = ScanDataframe; - // Where a scan reads from: a dataframe directly, or a query run by SQLite. - using Source = std::variant; + using DataframeArg = ScanDataframeArg; + // Where a scan reads from: a dataframe, SQL, or a dataframe argument. SQL + // sources become dataframe arguments when the plan is written into SQL, and + // those are bound to dataframes when it is loaded. + using Source = std::variant; Source source; // Bindings in source column order. std::vector columns; diff --git a/src/trace_processor/perfetto_sql/pipeline/logical_plan_test_utils.cc b/src/trace_processor/perfetto_sql/pipeline/logical_plan_test_utils.cc index 1527535dac8..f2508500185 100644 --- a/src/trace_processor/perfetto_sql/pipeline/logical_plan_test_utils.cc +++ b/src/trace_processor/perfetto_sql/pipeline/logical_plan_test_utils.cc @@ -66,6 +66,12 @@ std::string ScanString(const LogicalPlan& plan, const op::Scan& scan) { case base::variant_index(): out += "sql " + base::unchecked_get(scan.source).sql(); break; + case base::variant_index(): + out += + "argument " + + std::to_string( + base::unchecked_get(scan.source).index); + break; default: PERFETTO_FATAL("Unknown scan source"); } diff --git a/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc b/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc index 08a232b1083..8b9e835bf0a 100644 --- a/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc +++ b/src/trace_processor/perfetto_sql/pipeline/physical_plan.cc @@ -37,7 +37,6 @@ #include "src/trace_processor/core/exec/tree_accumulate.h" #include "src/trace_processor/core/exec/tree_number_nodes.h" #include "src/trace_processor/core/exec/tree_order.h" -#include "src/trace_processor/perfetto_sql/exec/sql_scan.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" namespace perfetto::trace_processor::pipeline { @@ -46,9 +45,8 @@ namespace ex = core::exec; // Builds one pipeline of operators, including any blocking ordering stages. class Lowering { public: - explicit Lowering(const LogicalPlan& plan, const LowerEnvironment& env) + explicit Lowering(const LogicalPlan& plan) : plan_(plan), - env_(env), out_(std::make_unique()), positions_(plan.columns.size(), std::numeric_limits::max()), int64_columns_(plan.columns.size(), false) {} @@ -80,9 +78,8 @@ class Lowering { } void Define(ColumnId id) { positions_[id] = column_count_++; } - // Inputs, borrowed for the duration of lowering. + // Borrowed for the duration of lowering. const LogicalPlan& plan_; - const LowerEnvironment& env_; // Execution graph under construction. std::unique_ptr out_; @@ -131,17 +128,9 @@ std::unique_ptr Lowering::MakeSource(const op::Scan& scan) const { return std::make_unique(source.columns, source.row_count); } - case base::variant_index(): { - Schema columns; - columns.reserve(scan.columns.size()); - for (const NamedColumn& column : scan.columns) { - columns.push_back({column.name, plan_.columns[column.id].type}); - } - return std::make_unique( - env_.connection, base::unchecked_get(scan.source), - std::move(columns), env_.pool); - } default: + // SQL is moved out into dataframe arguments, and those are bound to + // dataframes, before a plan is run. PERFETTO_FATAL("Unknown scan source"); } } @@ -279,9 +268,8 @@ std::unique_ptr Lowering::Finish() { PhysicalPlan::PhysicalPlan() = default; PhysicalPlan::~PhysicalPlan() = default; -std::unique_ptr Lower(const LogicalPlan& plan, - const LowerEnvironment& env) { - Lowering lowering(plan, env); +std::unique_ptr Lower(const LogicalPlan& plan) { + Lowering lowering(plan); lowering.LowerNode(plan.root); return lowering.Finish(); } diff --git a/src/trace_processor/perfetto_sql/pipeline/physical_plan.h b/src/trace_processor/perfetto_sql/pipeline/physical_plan.h index d9af94566dd..5d1eb80c52e 100644 --- a/src/trace_processor/perfetto_sql/pipeline/physical_plan.h +++ b/src/trace_processor/perfetto_sql/pipeline/physical_plan.h @@ -26,20 +26,8 @@ #include "src/trace_processor/core/exec/pipeline.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" -namespace perfetto::trace_processor { -class SqliteConnection; -class StringPool; -} // namespace perfetto::trace_processor - namespace perfetto::trace_processor::pipeline { -// Connection state needed by lowering. The connection and pool must outlive -// the plan. Dataframe columns have already been resolved by logical planning. -struct LowerEnvironment { - SqliteConnection* connection = nullptr; - StringPool* pool = nullptr; -}; - // A pipeline ready to run: the executor nodes plus which batch columns are // the output. Const once built, so it can be run any number of times with one // state per run. @@ -73,8 +61,8 @@ class PhysicalPlan { // Builds executor nodes from a logical plan. Establishes tree numbering, row // ordering, and column types as needed, reusing them across consecutive folds. -std::unique_ptr Lower(const LogicalPlan&, - const LowerEnvironment&); +// Every dataframe the plan reads must already be resolved. +std::unique_ptr Lower(const LogicalPlan&); } // namespace perfetto::trace_processor::pipeline diff --git a/src/trace_processor/perfetto_sql/pipeline/physical_plan_unittest.cc b/src/trace_processor/perfetto_sql/pipeline/physical_plan_unittest.cc index 566c4711003..e4d2711eb7f 100644 --- a/src/trace_processor/perfetto_sql/pipeline/physical_plan_unittest.cc +++ b/src/trace_processor/perfetto_sql/pipeline/physical_plan_unittest.cc @@ -37,6 +37,7 @@ #include "src/trace_processor/core/exec/variant.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/pipeline_sql.h" #include "src/trace_processor/perfetto_sql/pipeline/plan_serialization.h" #include "src/trace_processor/perfetto_sql/pipeline/test_catalog.h" #include "src/trace_processor/sqlite/sql_source.h" @@ -82,10 +83,7 @@ class PhysicalPlanTest : public ::testing::Test { protected: PhysicalPlanTest() : connection_(SqliteConnection::CreateConnectionToNewDatabase()), - catalog_(&pool_, connection_.get()) { - env_.connection = connection_.get(); - env_.pool = &pool_; - } + catalog_(&pool_, connection_.get()) {} void Exec(const std::string& sql) { auto statement = @@ -127,7 +125,19 @@ class PhysicalPlanTest : public ::testing::Test { base::StatusOr> Plan(const std::string& sql) { ASSIGN_OR_RETURN(LogicalPlan plan, Compile(sql)); - return Lower(plan, env_); + return LowerReadingSql(std::move(plan)); + } + + // Lowers `plan` to read each of its SQL sources as a dataframe, built up + // front as SQLite builds it where the pipeline is written. + base::StatusOr> LowerReadingSql( + LogicalPlan plan) { + ASSIGN_OR_RETURN(auto dataframes, + BuildSqlSources(connection_.get(), &pool_, plan)); + LogicalPlan moved = MoveSqlSourcesToDataframeArgs(std::move(plan)).plan; + RETURN_IF_ERROR( + BindDataframeArgs(moved, DataframeArgs(dataframes), &pool_)); + return Lower(moved); } std::vector Names(const PhysicalPlan& plan) { @@ -165,7 +175,6 @@ class PhysicalPlanTest : public ::testing::Test { StringPool pool_; std::unique_ptr connection_; TestCatalog catalog_; - LowerEnvironment env_; base::FlatHashMap macros_; }; @@ -296,7 +305,7 @@ TEST_F(PhysicalPlanTest, OutputBindingsUseIdsRatherThanBatchPositions) { ColumnId id = logical.output.front().id; // Project and alias the same value twice, independently of the source names. logical.output = {{"total", total}, {"id", id}, {"again", total}}; - auto plan = Lower(logical, env_); + auto plan = std::move(*LowerReadingSql(std::move(logical))); EXPECT_THAT(Names(*plan), ElementsAre("total", "id", "again")); EXPECT_NE(plan->columns()[0].index, total); EXPECT_EQ(plan->columns()[0].index, plan->columns()[2].index); @@ -342,7 +351,7 @@ TEST_F(PhysicalPlanTest, LogicalPlanRetainsColumnsBeforeLowering) { LogicalPlan plan = std::get(parser.statement()).plan; catalog_.RemoveTable("df"); - auto physical = Lower(plan, env_); + auto physical = Lower(plan); auto rows = Run(*physical, "total"); ASSERT_TRUE(rows.ok()) << rows.status().message(); EXPECT_THAT(*rows, @@ -351,7 +360,7 @@ TEST_F(PhysicalPlanTest, LogicalPlanRetainsColumnsBeforeLowering) { TEST_F(PhysicalPlanTest, ValuesWhichAreNotIntegersFailTheRun) { CreateTree(); - Exec("INSERT INTO tree VALUES (4, 0, 'many')"); + Exec("UPDATE tree SET self = 'many'"); auto plan = Plan("FROM tree |> TREE ACCUMULATE UP SUM(self) AS total"); ASSERT_TRUE(plan.ok()) << plan.status().message(); EXPECT_THAT(Run(**plan, "total").status().message(), HasSubstr("'self'")); @@ -359,7 +368,7 @@ TEST_F(PhysicalPlanTest, ValuesWhichAreNotIntegersFailTheRun) { TEST_F(PhysicalPlanTest, DiagnosticsUseDefiningNamesAfterProjection) { CreateTree(); - Exec("INSERT INTO tree VALUES (4, 0, 'many')"); + Exec("UPDATE tree SET self = 'many'"); PerfettoSqlParser parser(macros_, catalog_, /*pipelines_allowed=*/true); parser.Reset(SqlSource::FromExecuteQuery( @@ -369,7 +378,7 @@ TEST_F(PhysicalPlanTest, DiagnosticsUseDefiningNamesAfterProjection) { // Neither the original input name nor its value is exposed in the result. logical.output = {{"id", logical.output.front().id}, {"renamed_total", logical.output.back().id}}; - auto physical = Lower(logical, env_); + auto physical = std::move(*LowerReadingSql(std::move(logical))); EXPECT_THAT(Run(*physical, "renamed_total").status().message(), HasSubstr("'self'")); } @@ -405,7 +414,7 @@ TEST_F(PhysicalPlanTest, APlanReadBackReadsTablesAsTheyAreNow) { {{1, 0, 2}, {0, std::nullopt, 1}}); auto read = DeserializePlan(bytes, catalog_); ASSERT_TRUE(read.ok()) << read.status().message(); - auto rows = Run(*Lower(*read, env_), "total"); + auto rows = Run(*Lower(*read), "total"); ASSERT_TRUE(rows.ok()) << rows.status().message(); EXPECT_THAT(*rows, ElementsAre(Pair(0, 3), Pair(1, 2))); } diff --git a/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc b/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc index 74aebb5e1c4..1ea1edcecaf 100644 --- a/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc +++ b/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.cc @@ -18,29 +18,81 @@ #include #include +#include +#include #include #include "perfetto/base/status.h" #include "perfetto/ext/base/status_or.h" #include "perfetto/ext/base/string_utils.h" +#include "perfetto/ext/base/variant.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" #include "src/trace_processor/perfetto_sql/pipeline/plan_serialization.h" +#include "src/trace_processor/sqlite/sql_source.h" namespace perfetto::trace_processor::pipeline { +namespace { + +std::string QuoteIdentifier(const std::string& name) { + return "\"" + base::ReplaceAll(name, "\"", "\"\"") + "\""; +} + +std::string QuoteString(const std::string& text) { + return "'" + base::ReplaceAll(text, "'", "''") + "'"; +} + +} // namespace + +PlanWithDataframeArgs MoveSqlSourcesToDataframeArgs(LogicalPlan plan) { + PlanWithDataframeArgs out; + for (PlanNode& node : plan.nodes) { + if (!node.Is()) { + continue; + } + auto& scan = node.Cast(); + if (!std::holds_alternative(scan.source)) { + continue; + } + std::vector names; + std::vector references; + for (const NamedColumn& column : scan.columns) { + names.push_back(column.name); + references.push_back(QuoteIdentifier(column.name)); + } + const std::string& from = base::unchecked_get(scan.source).sql(); + out.args.push_back("SELECT " + std::string(kDataframeAggFunction) + "(" + + QuoteString(base::Join(names, ",")) + ", " + + base::Join(references, ", ") + ") FROM " + from); + scan.source = + op::Scan::DataframeArg{static_cast(out.args.size() - 1)}; + } + out.plan = std::move(plan); + return out; +} base::StatusOr 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()); } + PlanWithDataframeArgs moved = MoveSqlSourcesToDataframeArgs(plan); std::vector 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, "\"", "\"\"") + - "\""); + columns.push_back("c" + std::to_string(i) + " AS " + + QuoteIdentifier(plan.output[i].name)); + } + std::string arguments = "X'" + base::ToHex(SerializePlan(moved.plan)) + "'"; + if (!moved.args.empty()) { + std::vector dataframes; + for (const std::string& arg : moved.args) { + dataframes.push_back("(" + arg + ")"); + } + // A subquery, so SQLite builds the list once rather than on every filter. + arguments += ", (SELECT " + std::string(kDataframesFunction) + "(" + + base::Join(dataframes, ", ") + "))"; } return "SELECT " + base::Join(columns, ", ") + " FROM " + kPipelineFunction + - "(X'" + base::ToHex(SerializePlan(plan)) + "')"; + "(" + arguments + ")"; } } // namespace perfetto::trace_processor::pipeline diff --git a/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h b/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h index a2dbf43b27e..3b512f6ae95 100644 --- a/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h +++ b/src/trace_processor/perfetto_sql/pipeline/pipeline_sql.h @@ -19,6 +19,7 @@ #include #include +#include #include "perfetto/ext/base/status_or.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" @@ -28,11 +29,30 @@ namespace perfetto::trace_processor::pipeline { // The table function which runs a serialized plan. inline constexpr char kPipelineFunction[] = "__intrinsic_pipeline"; +// The aggregate which builds a dataframe of a relation's rows, for a pipeline +// to read: `__intrinsic_dataframe_agg('a,b', a, b)`. +inline constexpr char kDataframeAggFunction[] = "__intrinsic_dataframe_agg"; + +// The function which gathers the dataframes a pipeline reads into the one +// argument the table function takes them in: +// `__intrinsic_dataframes((SELECT ...), (SELECT ...))`. +inline constexpr char kDataframesFunction[] = "__intrinsic_dataframes"; + // The most columns the table function can output. inline constexpr uint32_t kMaxPipelineColumns = 256; +// A plan whose SQL sources have been moved out into dataframe arguments: the +// plan reads its i-th SQL source as dataframe argument i, and `args[i]` is the +// SQL building it. +struct PlanWithDataframeArgs { + LogicalPlan plan; + std::vector args; +}; +PlanWithDataframeArgs MoveSqlSourcesToDataframeArgs(LogicalPlan plan); + // 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. +// serialized into the SQL, and each relation it reads from SQL is built into a +// dataframe where that SQL runs, so it sees whatever is in scope there. base::StatusOr SelectPipelineSql(const LogicalPlan& plan); } // namespace perfetto::trace_processor::pipeline diff --git a/src/trace_processor/perfetto_sql/pipeline/plan_serialization.cc b/src/trace_processor/perfetto_sql/pipeline/plan_serialization.cc index d480f624a9b..b7c0803fb98 100644 --- a/src/trace_processor/perfetto_sql/pipeline/plan_serialization.cc +++ b/src/trace_processor/perfetto_sql/pipeline/plan_serialization.cc @@ -31,7 +31,9 @@ #include "perfetto/base/status.h" #include "perfetto/ext/base/status_macros.h" #include "perfetto/ext/base/status_or.h" +#include "src/trace_processor/containers/string_pool.h" #include "src/trace_processor/core/common/storage_types.h" +#include "src/trace_processor/core/dataframe/adhoc_dataframe_builder.h" #include "src/trace_processor/core/dataframe/dataframe.h" #include "src/trace_processor/core/dataframe/specs.h" #include "src/trace_processor/perfetto_sql/pipeline/catalog.h" @@ -196,10 +198,11 @@ class PlanWriter { case base::variant_index(): w_.Str(base::unchecked_get(scan.source).name); break; - case base::variant_index(): - w_.Str(base::unchecked_get(scan.source).sql()); + case base::variant_index(): + w_.U32(base::unchecked_get(scan.source).index); break; default: + // SQL is moved out into dataframe arguments before a plan is written. PERFETTO_FATAL("Unknown scan source"); } w_.Size(scan.columns.size()); @@ -297,24 +300,18 @@ class PlanReader { scan.source = std::move(source); break; } - case base::variant_index(): - scan.source = SqlSource::FromTraceProcessorImplementation(r_.Str()); + case base::variant_index(): + scan.source = op::Scan::DataframeArg{r_.U32()}; break; default: r_.Fail(); return {}; } scan.columns.resize(r_.Count()); - bool from_sql = std::holds_alternative(scan.source); Available available; for (NamedColumn& column : scan.columns) { column.name = r_.Str(); - std::optional type = ReadType(r_); - // SqlScan cannot produce Ids. - if (from_sql && type && type->Is()) { - r_.Fail(); - } - column.id = plan_.AddColumn(column.name, type); + column.id = plan_.AddColumn(column.name, ReadType(r_)); available.push_back(column.id); } return available; @@ -373,6 +370,17 @@ class PlanReader { LogicalPlan plan_; }; +std::optional FindScanColumn(const dataframe::Dataframe& dataframe, + const std::string& name) { + const std::vector& names = dataframe.column_names(); + for (uint32_t i = 0; i < names.size(); ++i) { + if (names[i] == name && !dataframe::IsHiddenColumn(names[i])) { + return i; + } + } + return std::nullopt; +} + // Points each dataframe scan at the dataframe now registered under its name, // which must still have every column the plan reads, with the same type. base::Status ResolveDataframes(LogicalPlan& plan, const Catalog& catalog) { @@ -390,21 +398,16 @@ base::Status ResolveDataframes(LogicalPlan& plan, const Catalog& catalog) { return base::ErrStatus("Pipeline: table '%s' no longer exists", source->name.c_str()); } - const std::vector& names = dataframe->column_names(); for (const NamedColumn& column : scan.columns) { - uint32_t i = 0; - while (i < names.size() && - (names[i] != column.name || dataframe::IsHiddenColumn(names[i]))) { - ++i; - } + std::optional i = FindScanColumn(*dataframe, column.name); const std::optional& type = plan.columns[column.id].type; - if (i == names.size() || !type || !(*type == dataframe->column_type(i))) { + if (!i || !type || !(*type == dataframe->column_type(*i))) { return base::ErrStatus( "Pipeline: table '%s' has changed since the pipeline was written", source->name.c_str()); } - source->columns.push_back(dataframe->shared_column(i)); + source->columns.push_back(dataframe->shared_column(*i)); } source->row_count = dataframe->row_count(); } @@ -417,6 +420,56 @@ std::string SerializePlan(const LogicalPlan& plan) { return PlanWriter(plan).Write(); } +base::Status BindDataframeArgs( + LogicalPlan& plan, + const std::vector& args, + StringPool* pool) { + for (PlanNode& node : plan.nodes) { + if (!node.Is()) { + continue; + } + auto& scan = node.Cast(); + const auto* arg = std::get_if(&scan.source); + if (!arg) { + continue; + } + if (arg->index >= args.size()) { + return base::ErrStatus("Pipeline: no dataframe argument %u", arg->index); + } + // A relation with no rows is passed as null: read it as empty columns. + std::optional empty; + const dataframe::Dataframe* dataframe = args[arg->index]; + if (!dataframe) { + std::vector names; + for (const NamedColumn& column : scan.columns) { + names.push_back(column.name); + } + dataframe::AdhocDataframeBuilder::Options options; + options.emit_auto_id = false; + ASSIGN_OR_RETURN(empty, dataframe::AdhocDataframeBuilder(std::move(names), + pool, options) + .Build()); + dataframe = &*empty; + } + op::Scan::Dataframe source; + source.name = "dataframe argument " + std::to_string(arg->index); + for (const NamedColumn& column : scan.columns) { + std::optional i = FindScanColumn(*dataframe, column.name); + if (!i) { + return base::ErrStatus("Pipeline: %s has no column '%s'", + source.name.c_str(), column.name.c_str()); + } + // The dataframe was built after the plan was written, so it decides + // what each column holds. + plan.columns[column.id].type = dataframe->column_type(*i); + source.columns.push_back(dataframe->shared_column(*i)); + } + source.row_count = dataframe->row_count(); + scan.source = std::move(source); + } + return base::OkStatus(); +} + base::StatusOr DeserializePlan(std::string_view bytes, const Catalog& catalog) { Reader r(bytes); diff --git a/src/trace_processor/perfetto_sql/pipeline/plan_serialization.h b/src/trace_processor/perfetto_sql/pipeline/plan_serialization.h index b7ccb5f7441..dfb362f4a5d 100644 --- a/src/trace_processor/perfetto_sql/pipeline/plan_serialization.h +++ b/src/trace_processor/perfetto_sql/pipeline/plan_serialization.h @@ -19,8 +19,12 @@ #include #include +#include +#include "perfetto/base/status.h" #include "perfetto/ext/base/status_or.h" +#include "src/trace_processor/containers/string_pool.h" +#include "src/trace_processor/core/dataframe/dataframe.h" #include "src/trace_processor/perfetto_sql/pipeline/catalog.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" @@ -36,6 +40,16 @@ std::string SerializePlan(const LogicalPlan&); // changed. base::StatusOr DeserializePlan(std::string_view, const Catalog&); +// Points each scan of a dataframe argument at `args[i]`, the dataframe the +// plan is passed as its i-th argument when it runs, and types the scan's +// columns as the dataframe does. A null argument is a relation with no rows. +// Anyone can pass dataframes to a plan in SQL, so each must have the columns +// the plan reads from it. +base::Status BindDataframeArgs( + LogicalPlan& plan, + const std::vector& args, + StringPool* pool); + } // namespace perfetto::trace_processor::pipeline #endif // SRC_TRACE_PROCESSOR_PERFETTO_SQL_PIPELINE_PLAN_SERIALIZATION_H_ diff --git a/src/trace_processor/perfetto_sql/pipeline/plan_serialization_unittest.cc b/src/trace_processor/perfetto_sql/pipeline/plan_serialization_unittest.cc index 0be4f97c94e..d858c6690df 100644 --- a/src/trace_processor/perfetto_sql/pipeline/plan_serialization_unittest.cc +++ b/src/trace_processor/perfetto_sql/pipeline/plan_serialization_unittest.cc @@ -33,6 +33,7 @@ #include "src/trace_processor/perfetto_sql/pipeline/logical_plan.h" #include "src/trace_processor/perfetto_sql/pipeline/logical_plan_test_utils.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/test_catalog.h" #include "src/trace_processor/sqlite/sql_source.h" #include "src/trace_processor/sqlite/sqlite_connection.h" @@ -48,8 +49,6 @@ class PlanSerializationTest : public ::testing::Test { PlanSerializationTest() : connection_(SqliteConnection::CreateConnectionToNewDatabase()), catalog_(&pool_, connection_.get()) { - env_.connection = connection_.get(); - env_.pool = &pool_; // 0 (10) -> 1 (20) -> 3 (40) // -> 2 (30) catalog_.AddTable( @@ -67,12 +66,18 @@ class PlanSerializationTest : public ::testing::Test { ASSERT_TRUE(statement.status().ok()) << statement.status().c_message(); } + // The plan as it is written into SQL, reading its SQL sources as dataframe + // arguments, whose dataframes are built into `dataframes_`. LogicalPlan Compile(const std::string& sql) { PerfettoSqlParser parser(macros_, catalog_, /*pipelines_allowed=*/true); parser.Reset(SqlSource::FromExecuteQuery(sql)); PERFETTO_CHECK(parser.Next()); - return std::move( + LogicalPlan plan = std::move( std::get(parser.TakeStatement()).plan); + auto dataframes = BuildSqlSources(connection_.get(), &pool_, plan); + PERFETTO_CHECK(dataframes.ok()); + dataframes_ = std::move(*dataframes); + return MoveSqlSourcesToDataframeArgs(std::move(plan)).plan; } base::StatusOr RoundTrip(const LogicalPlan& plan) { @@ -82,7 +87,7 @@ class PlanSerializationTest : public ::testing::Test { StringPool pool_; std::unique_ptr connection_; TestCatalog catalog_; - LowerEnvironment env_; + std::vector> dataframes_; base::FlatHashMap macros_; }; @@ -152,7 +157,11 @@ TEST_F(PlanSerializationTest, MalformedPlansAreRefusedOrRun) { if (!read.ok()) { continue; } - auto physical = Lower(*read, env_); + if (!BindDataframeArgs(*read, DataframeArgs(dataframes_), &pool_) + .ok()) { + continue; + } + auto physical = Lower(*read); core::exec::RowCursor cursor(physical->source()); for (bool row = cursor.Open(); row; row = cursor.Next()) { } diff --git a/src/trace_processor/perfetto_sql/pipeline/test_catalog.h b/src/trace_processor/perfetto_sql/pipeline/test_catalog.h index 20a5e89eb68..f57ea647bd8 100644 --- a/src/trace_processor/perfetto_sql/pipeline/test_catalog.h +++ b/src/trace_processor/perfetto_sql/pipeline/test_catalog.h @@ -23,26 +23,33 @@ #include #include #include +#include #include #include "perfetto/base/logging.h" #include "perfetto/base/status.h" #include "perfetto/ext/base/flat_hash_map.h" +#include "perfetto/ext/base/status_macros.h" #include "perfetto/ext/base/status_or.h" +#include "perfetto/ext/base/string_utils.h" +#include "perfetto/ext/base/variant.h" #include "src/perfetto_sql/analysis/relation.h" #include "src/trace_processor/containers/string_pool.h" #include "src/trace_processor/core/dataframe/adhoc_dataframe_builder.h" #include "src/trace_processor/core/dataframe/dataframe.h" +#include "src/trace_processor/core/dataframe/runtime_dataframe_builder.h" +#include "src/trace_processor/perfetto_sql/engine/sqlite_dataframe_builder.h" #include "src/trace_processor/perfetto_sql/pipeline/catalog.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/type_mapping.h" #include "src/trace_processor/sqlite/sql_source.h" #include "src/trace_processor/sqlite/sqlite_connection.h" +#include "src/trace_processor/sqlite/sqlite_utils.h" 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. +// Anything else is looked up in `connection` (if given) and is untyped. class TestCatalog : public Catalog { public: explicit TestCatalog(StringPool* pool, SqliteConnection* connection = nullptr) @@ -80,16 +87,32 @@ class TestCatalog : public Catalog { return dataframe ? dataframe->get() : nullptr; } - base::StatusOr DescribeQuery(const SqlSource& sql) const override { - if (!connection_) { - return base::ErrStatus("no such table"); - } - return sql_schema::DescribeQuery(connection_, sql, *this); - } - std::optional FindLeafRelation( - std::string_view) const override { - return std::nullopt; + std::string_view name) const override { + const dataframe::Dataframe* dataframe = FindDataframe(name); + if (!dataframe) { + if (!connection_) { + return std::nullopt; + } + perfetto_sql::analysis::LeafRelation relation{std::string(name), {}}; + for (const auto& column : + sqlite::utils::GetColumns(connection_->db(), relation.name)) { + relation.columns.push_back({column.name, std::nullopt, column.hidden}); + } + if (relation.columns.empty()) { + return std::nullopt; + } + return relation; + } + perfetto_sql::analysis::LeafRelation relation; + relation.name = std::string(name); + const std::vector& columns = dataframe->column_names(); + for (uint32_t i = 0; i < columns.size(); ++i) { + relation.columns.push_back( + {columns[i], sql_schema::ToAnalysisType(dataframe->column_type(i)), + dataframe::IsHiddenColumn(columns[i])}); + } + return relation; } std::optional FindViewSql(std::string_view) const override { return std::nullopt; @@ -102,6 +125,53 @@ class TestCatalog : public Catalog { dataframes_; }; +// A dataframe of each of `plan`'s SQL sources, built from `connection` as +// SQLite builds them where the pipeline is written. Once the SQL sources are +// moved out, the plan reads the i-th as dataframe argument i. +inline base::StatusOr>> +BuildSqlSources(SqliteConnection* connection, + StringPool* pool, + const LogicalPlan& plan) { + std::vector> out; + for (const PlanNode& node : plan.nodes) { + if (!node.Is()) { + continue; + } + const auto& scan = node.Cast(); + if (!std::holds_alternative(scan.source)) { + continue; + } + std::vector names; + std::vector references; + for (const NamedColumn& column : scan.columns) { + names.push_back(column.name); + references.push_back("\"" + column.name + "\""); + } + auto statement = connection->PrepareStatement(SqlSource::FromExecuteQuery( + "SELECT " + base::Join(references, ", ") + " FROM " + + base::unchecked_get(scan.source).sql())); + statement.Step(); + RETURN_IF_ERROR(statement.status()); + ASSIGN_OR_RETURN(dataframe::RuntimeDataframeBuilder builder, + BuildRuntimeDataframeFromSqliteStatement( + pool, std::move(names), &statement, "SQL source")); + ASSIGN_OR_RETURN(dataframe::Dataframe dataframe, + std::move(builder).Build()); + out.push_back(std::make_unique(std::move(dataframe))); + } + return std::move(out); +} + +// `dataframes`, as the arguments to pass a plan. +inline std::vector DataframeArgs( + const std::vector>& dataframes) { + std::vector args; + for (const auto& dataframe : dataframes) { + args.push_back(dataframe.get()); + } + return args; +} + } // namespace perfetto::trace_processor::pipeline #endif // SRC_TRACE_PROCESSOR_PERFETTO_SQL_PIPELINE_TEST_CATALOG_H_ diff --git a/src/trace_processor/perfetto_sql/schema/BUILD.gn b/src/trace_processor/perfetto_sql/schema/BUILD.gn index ec7dd41e623..025901360cf 100644 --- a/src/trace_processor/perfetto_sql/schema/BUILD.gn +++ b/src/trace_processor/perfetto_sql/schema/BUILD.gn @@ -13,11 +13,7 @@ # limitations under the License. source_set("schema") { - sources = [ - "query_schema.cc", - "query_schema.h", - "type_mapping.h", - ] + sources = [ "type_mapping.h" ] deps = [ "../../../../gn:default_deps", "../../../../gn:sqlite", diff --git a/src/trace_processor/perfetto_sql/schema/query_schema.cc b/src/trace_processor/perfetto_sql/schema/query_schema.cc deleted file mode 100644 index 2f56e53163b..00000000000 --- a/src/trace_processor/perfetto_sql/schema/query_schema.cc +++ /dev/null @@ -1,104 +0,0 @@ -/* - * 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/schema/query_schema.h" - -#include - -#include -#include -#include -#include - -#include "perfetto/ext/base/status_macros.h" -#include "src/perfetto_sql/syntaqlite/syntaqlite_perfetto.h" -#include "src/trace_processor/perfetto_sql/schema/type_mapping.h" -#include "src/trace_processor/sqlite/bindings/sqlite_column.h" - -namespace perfetto::trace_processor::sql_schema { -namespace { - -using core::StorageType; - -struct ParserDeleter { - void operator()(SyntaqliteParser* parser) const { - syntaqlite_parser_destroy(parser); - } -}; -using ScopedParser = std::unique_ptr; - -// The types lineage established, lined up with the query's columns. If the two -// disagree on the number of columns they are not describing the same query, so -// no type is claimed for any of them. -std::vector> ResolveTypes( - const SqlSource& sql, - uint32_t count, - const analysis::Catalog& catalog) { - std::vector> types(count); - ScopedParser parser(syntaqlite_parser_create_perfetto(nullptr)); - syntaqlite_parser_reset(parser.get(), sql.sql().data(), - static_cast(sql.sql().size())); - if (syntaqlite_parser_next(parser.get()) != SYNTAQLITE_PARSE_OK) { - return types; - } - analysis::RelationAnalyzer analyzer(catalog); - auto resolved = analyzer.AnalyzeQuery( - {parser.get(), syntaqlite_result_root(parser.get())}); - if (!resolved.ok() || resolved->columns().size() != count) { - return types; - } - for (uint32_t i = 0; i < count; ++i) { - std::optional type = resolved->columns()[i].type(); - if (!type) { - continue; - } - StorageType storage = ToStorageType(*type); - // An Id has no storage of its own: its value is the row it sits at. A - // query result has no such rows to point at, so materialise it at the - // narrowest width which holds one. - if (storage.Is()) { - storage = StorageType{core::Uint32{}}; - } - types[i] = storage; - } - return types; -} - -} // namespace - -base::StatusOr DescribeQuery(SqliteConnection* connection, - const SqlSource& sql, - const analysis::Catalog& catalog) { - // Prepared here only to read the column names, then discarded: a statement - // belongs to one execution, but the columns belong to the query. - SqliteConnection::PreparedStatement statement = - connection->PrepareStatement(sql); - RETURN_IF_ERROR(statement.status()); - - sqlite3_stmt* stmt = statement.sqlite_stmt(); - uint32_t count = sqlite::column::Count(stmt); - std::vector> types = - ResolveTypes(sql, count, catalog); - core::Schema columns; - columns.reserve(count); - for (uint32_t i = 0; i < count; ++i) { - const char* name = sqlite::column::Name(stmt, i); - columns.push_back({name ? name : "", types[i]}); - } - return columns; -} - -} // namespace perfetto::trace_processor::sql_schema diff --git a/src/trace_processor/perfetto_sql/schema/query_schema.h b/src/trace_processor/perfetto_sql/schema/query_schema.h deleted file mode 100644 index d4063a83f47..00000000000 --- a/src/trace_processor/perfetto_sql/schema/query_schema.h +++ /dev/null @@ -1,38 +0,0 @@ -/* - * 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_SCHEMA_QUERY_SCHEMA_H_ -#define SRC_TRACE_PROCESSOR_PERFETTO_SQL_SCHEMA_QUERY_SCHEMA_H_ - -#include "perfetto/ext/base/status_or.h" -#include "src/perfetto_sql/analysis/relation.h" -#include "src/trace_processor/core/common/schema.h" -#include "src/trace_processor/sqlite/sql_source.h" -#include "src/trace_processor/sqlite/sqlite_connection.h" - -namespace perfetto::trace_processor::sql_schema { - -// SQLite supplies result names; semantic analysis supplies types where known. -// Fails if SQLite cannot prepare the query. Unknown types remain per-row -// variants. -base::StatusOr DescribeQuery( - SqliteConnection*, - const SqlSource&, - const perfetto_sql::analysis::Catalog&); - -} // namespace perfetto::trace_processor::sql_schema - -#endif // SRC_TRACE_PROCESSOR_PERFETTO_SQL_SCHEMA_QUERY_SCHEMA_H_ diff --git a/src/trace_processor/plugins/span_join_operator/span_join_operator.cc b/src/trace_processor/plugins/span_join_operator/span_join_operator.cc index 1d565418a5b..d31f8786ac0 100644 --- a/src/trace_processor/plugins/span_join_operator/span_join_operator.cc +++ b/src/trace_processor/plugins/span_join_operator/span_join_operator.cc @@ -33,6 +33,7 @@ #include "perfetto/base/logging.h" #include "perfetto/base/status.h" #include "perfetto/ext/base/status_macros.h" +#include "perfetto/ext/base/status_or.h" #include "perfetto/ext/base/string_splitter.h" #include "perfetto/ext/base/string_utils.h" #include "perfetto/trace_processor/basic_types.h" @@ -138,6 +139,35 @@ std::string EscapedSqliteValueAsString(sqlite3_value* value) { } } +base::StatusOr ColumnType( + const sqlite::utils::SqliteColumn& column, + const std::string& table) { + const char* type = column.type.c_str(); + if (base::CaseInsensitiveEqual(type, "STRING") || + base::CaseInsensitiveEqual(type, "TEXT")) { + return SqlValue::Type::kString; + } + if (base::CaseInsensitiveEqual(type, "DOUBLE")) { + return SqlValue::Type::kDouble; + } + if (base::CaseInsensitiveEqual(type, "BIG INT") || + base::CaseInsensitiveEqual(type, "BIGINT") || + base::CaseInsensitiveEqual(type, "UNSIGNED INT") || + base::CaseInsensitiveEqual(type, "INT") || + base::CaseInsensitiveEqual(type, "BOOLEAN") || + base::CaseInsensitiveEqual(type, "INTEGER")) { + return SqlValue::Type::kLong; + } + if (base::CaseInsensitiveEqual(type, "BLOB")) { + return SqlValue::Type::kBytes; + } + if (column.type.empty()) { + return SqlValue::Type::kNull; + } + return base::ErrStatus("Unknown column type '%s' on table %s", type, + table.c_str()); +} + } // namespace void SpanJoinOperatorModule::Vtab::PopulateColumnLocatorMap(uint32_t offset) { @@ -227,8 +257,19 @@ base::Status SpanJoinOperatorModule::TableDefinition::Create( } std::vector> cols; - RETURN_IF_ERROR(sqlite::utils::GetColumnsForTable( - connection->sqlite_connection()->db(), desc.name, cols)); + // Table functions are named with their arguments. + std::string table = desc.name.substr(0, desc.name.find('(')); + for (const auto& column : sqlite::utils::GetColumns( + connection->sqlite_connection()->db(), table)) { + if (!column.hidden) { + ASSIGN_OR_RETURN(SqlValue::Type type, ColumnType(column, desc.name)); + cols.emplace_back(type, column.name); + } + } + if (cols.empty()) { + return base::ErrStatus("Unknown table or view name '%s'", + desc.name.c_str()); + } uint32_t required_columns_found = 0; uint32_t ts_idx = std::numeric_limits::max(); diff --git a/src/trace_processor/sqlite/sqlite_utils.cc b/src/trace_processor/sqlite/sqlite_utils.cc index 1d429bb601d..a5d637c29fb 100644 --- a/src/trace_processor/sqlite/sqlite_utils.cc +++ b/src/trace_processor/sqlite/sqlite_utils.cc @@ -96,82 +96,28 @@ base::StatusOr ExtractArgument(size_t argc, } } // namespace internal -base::Status GetColumnsForTable( - sqlite3* db, - const std::string& raw_table_name, - std::vector>& columns) { - PERFETTO_DCHECK(columns.empty()); - char sql[1024]; - const char kRawSql[] = "SELECT name, type from pragma_table_info(\"%s\")"; - - // Support names which are table valued functions with arguments. - std::string table_name = raw_table_name.substr(0, raw_table_name.find('(')); - size_t n = base::SprintfTrunc(sql, sizeof(sql), kRawSql, table_name.c_str()); - PERFETTO_DCHECK(n > 0); - +std::vector GetColumns(sqlite3* db, const std::string& name) { sqlite3_stmt* raw_stmt = nullptr; - int err = - sqlite3_prepare_v2(db, sql, static_cast(n), &raw_stmt, nullptr); - if (err != SQLITE_OK) { - return base::ErrStatus("Preparing database failed"); + if (sqlite3_prepare_v2(db, + "SELECT name, type, hidden FROM pragma_table_xinfo(?)", + -1, &raw_stmt, nullptr) != SQLITE_OK) { + return {}; } ScopedStmt stmt(raw_stmt); - PERFETTO_DCHECK(sqlite3_column_count(*stmt) == 2); - - for (;;) { - err = sqlite3_step(raw_stmt); - if (err == SQLITE_DONE) - break; - if (err != SQLITE_ROW) { - return base::ErrStatus("Querying schema of table %s failed", - raw_table_name.c_str()); - } - - const char* name = + sqlite3_bind_text(*stmt, 1, name.c_str(), static_cast(name.size()), + kSqliteStatic); + std::vector columns; + int rc; + while ((rc = sqlite3_step(*stmt)) == SQLITE_ROW) { + const auto* column = reinterpret_cast(sqlite3_column_text(*stmt, 0)); - const char* raw_type = + const auto* type = reinterpret_cast(sqlite3_column_text(*stmt, 1)); - if (!name || !raw_type || !*name) { - return base::ErrStatus("Schema for %s has invalid column values", - raw_table_name.c_str()); - } - - SqlValue::Type type; - if (base::CaseInsensitiveEqual(raw_type, "STRING") || - base::CaseInsensitiveEqual(raw_type, "TEXT")) { - type = SqlValue::Type::kString; - } else if (base::CaseInsensitiveEqual(raw_type, "DOUBLE")) { - type = SqlValue::Type::kDouble; - } else if (base::CaseInsensitiveEqual(raw_type, "BIG INT") || - base::CaseInsensitiveEqual(raw_type, "BIGINT") || - base::CaseInsensitiveEqual(raw_type, "UNSIGNED INT") || - base::CaseInsensitiveEqual(raw_type, "INT") || - base::CaseInsensitiveEqual(raw_type, "BOOLEAN") || - base::CaseInsensitiveEqual(raw_type, "INTEGER")) { - type = SqlValue::Type::kLong; - } else if (base::CaseInsensitiveEqual(raw_type, "BLOB")) { - type = SqlValue::Type::kBytes; - } else if (!*raw_type) { - PERFETTO_DLOG("Unknown column type for %s %s", raw_table_name.c_str(), - name); - type = SqlValue::Type::kNull; - } else { - return base::ErrStatus("Unknown column type '%s' on table %s", raw_type, - raw_table_name.c_str()); - } - columns.emplace_back(type, name); - } - - // Catch mis-spelt table names. - // - // A SELECT on pragma_table_info() returns no rows if the - // table that was queried is not present. - if (columns.empty()) { - return base::ErrStatus("Unknown table or view name '%s'", - raw_table_name.c_str()); + // 1 is a hidden column; 2 and 3 are generated ones, which `*` includes. + columns.push_back({column ? column : "", type ? type : "", + sqlite3_column_int(*stmt, 2) == 1}); } - - return base::OkStatus(); + return rc == SQLITE_DONE ? columns : std::vector(); } const char* SqliteTypeToFriendlyString(SqlValue::Type type) { diff --git a/src/trace_processor/sqlite/sqlite_utils.h b/src/trace_processor/sqlite/sqlite_utils.h index c051e8abb10..689d8d37ec7 100644 --- a/src/trace_processor/sqlite/sqlite_utils.h +++ b/src/trace_processor/sqlite/sqlite_utils.h @@ -277,11 +277,15 @@ inline std::string SqlValueTypeToSqliteTypeName(SqlValue::Type type) { PERFETTO_FATAL("Not reached"); // For gcc } -// Returns the column names for the table named by |raw_table_name|. -base::Status GetColumnsForTable( - sqlite3* db, - const std::string& raw_table_name, - std::vector>& columns); +struct SqliteColumn { + std::string name; + // As declared, which SQLite does not enforce. + std::string type; + bool hidden = false; +}; + +// The columns of the table, view or table function SQLite knows as `name`. +std::vector GetColumns(sqlite3* db, const std::string& name); // Given an SqlValue::Type, converts it to a human-readable string. // This should really only be used for debugging messages. diff --git a/src/trace_processor/sqlite/sqlite_utils_unittest.cc b/src/trace_processor/sqlite/sqlite_utils_unittest.cc index 52ab54a50bc..f5cfeb0b50a 100644 --- a/src/trace_processor/sqlite/sqlite_utils_unittest.cc +++ b/src/trace_processor/sqlite/sqlite_utils_unittest.cc @@ -25,18 +25,16 @@ #include "perfetto/base/logging.h" #include "perfetto/trace_processor/basic_types.h" -#include "src/base/test/status_matchers.h" #include "src/trace_processor/sqlite/scoped_db.h" #include "test/gtest_and_gmock.h" namespace perfetto::trace_processor::sqlite::utils { namespace { -using base::gtest_matchers::IsError; -class GetColumnsForTableTest : public ::testing::Test { +class GetColumnsTest : public ::testing::Test { public: - GetColumnsForTableTest() { + GetColumnsTest() { sqlite3* db = nullptr; PERFETTO_CHECK(sqlite3_initialize() == SQLITE_OK); PERFETTO_CHECK(sqlite3_open(":memory:", &db) == SQLITE_OK); @@ -61,26 +59,17 @@ class GetColumnsForTableTest : public ::testing::Test { ScopedStmt stmt_; }; -TEST_F(GetColumnsForTableTest, ValidInput) { - RunStatement("CREATE TABLE foo (name STRING, ts INT, dur INT);"); - std::vector> columns; - ASSERT_OK(sqlite::utils::GetColumnsForTable(*db_, "foo", columns)); +TEST_F(GetColumnsTest, Columns) { + RunStatement("CREATE TABLE foo (name STRING, ts INT, dur);"); + std::vector columns = GetColumns(*db_, "foo"); + ASSERT_EQ(columns.size(), 3u); + EXPECT_EQ(columns[0].name, "name"); + EXPECT_EQ(columns[0].type, "STRING"); + EXPECT_EQ(columns[2].type, ""); } -TEST_F(GetColumnsForTableTest, UnknownType) { - // Currently GetColumnsForTable does not work with tables containing types it - // doesn't recognise. This just ensures that the query fails rather than - // crashing. - RunStatement("CREATE TABLE foo (name NUM, ts INT, dur INT);"); - std::vector> columns; - ASSERT_THAT(sqlite::utils::GetColumnsForTable(*db_, "foo", columns), - IsError()); -} - -TEST_F(GetColumnsForTableTest, UnknownTableName) { - std::vector> columns; - ASSERT_THAT(sqlite::utils::GetColumnsForTable(*db_, "unknowntable", columns), - IsError()); +TEST_F(GetColumnsTest, UnknownTableName) { + EXPECT_TRUE(GetColumns(*db_, "unknowntable").empty()); } } // namespace