From ec4e59f93f5d7250d4d63bdbb2d17fe4eea01568 Mon Sep 17 00:00:00 2001 From: Lalit Maganti Date: Sun, 27 Sep 2026 20:50:11 +0100 Subject: [PATCH] tp: pass the SQL a pipeline reads to it as dataframes A pipeline ran its SQL sources itself, as statements of their own. So they could see nothing of the statement the pipeline was written in. Now SQLite evaluates each relation a pipeline reads where the pipeline is written, and builds it into a dataframe, as a PERFETTO TABLE is built. __intrinsic_dataframes gathers them into one list, passed to __intrinsic_pipeline after the plan, and the pipeline reads each like any other dataframe. A pipeline never runs SQL itself, and a plan never holds SQL. - __intrinsic_dataframe_agg builds the dataframe. A plan reads its i-th as dataframe argument i, which is bound to the dataframe when the plan is loaded, before it is lowered, so lowering only ever sees dataframes. - The columns of the relation come from semantic analysis rather than from preparing its SQL. Relations analysis cannot describe yet, such as VALUES, fail clearly. Tables SQLite knows, including table functions, are described from its schema. - As in a PERFETTO TABLE, a column must hold one type throughout. The dataframe is built without analysis, keeping integers as Int64 and estimating no statistics, as a pipeline only scans it. This removes SqlScan, DescribeQuery and positional pruning of SQL sources, and lowering no longer needs a string pool. --- Android.bp | 25 - BUILD | 15 - src/perfetto_sql/analysis/relation.cc | 8 +- src/perfetto_sql/syntaqlite/BUILD.gn | 1 - src/trace_processor/BUILD.gn | 1 - .../core/dataframe/adhoc_dataframe_builder.cc | 40 +- .../core/dataframe/adhoc_dataframe_builder.h | 7 + .../adhoc_dataframe_builder_unittest.cc | 23 + .../core/dataframe/dataframe.cc | 12 +- .../core/dataframe/dataframe.h | 8 + .../perfetto_sql/engine/connection_catalog.cc | 48 +- .../perfetto_sql/engine/connection_catalog.h | 5 - .../engine/perfetto_sql_connection.cc | 23 +- .../engine/perfetto_sql_connection.h | 6 +- .../perfetto_sql_connection_unittest.cc | 103 +++- .../perfetto_sql/engine/pipeline_module.cc | 149 ++++- .../perfetto_sql/engine/pipeline_module.h | 37 ++ .../engine/sqlite_dataframe_builder.cc | 43 ++ .../engine/sqlite_dataframe_builder.h | 8 + .../perfetto_sql/exec/BUILD.gn | 51 -- .../perfetto_sql/exec/sql_scan.cc | 283 ---------- .../perfetto_sql/exec/sql_scan.h | 102 ---- .../perfetto_sql/exec/sql_scan_unittest.cc | 532 ------------------ .../parser/perfetto_sql_parser_unittest.cc | 55 +- .../perfetto_sql/pipeline/BUILD.gn | 6 +- .../perfetto_sql/pipeline/catalog.h | 7 - .../perfetto_sql/pipeline/column_pruning.cc | 38 +- .../perfetto_sql/pipeline/compiler.cc | 76 ++- .../perfetto_sql/pipeline/logical_plan.h | 13 +- .../pipeline/logical_plan_test_utils.cc | 6 + .../perfetto_sql/pipeline/physical_plan.cc | 24 +- .../perfetto_sql/pipeline/physical_plan.h | 16 +- .../pipeline/physical_plan_unittest.cc | 33 +- .../perfetto_sql/pipeline/pipeline_sql.cc | 60 +- .../perfetto_sql/pipeline/pipeline_sql.h | 22 +- .../pipeline/plan_serialization.cc | 91 ++- .../pipeline/plan_serialization.h | 14 + .../pipeline/plan_serialization_unittest.cc | 19 +- .../perfetto_sql/pipeline/test_catalog.h | 92 ++- .../perfetto_sql/schema/BUILD.gn | 6 +- .../perfetto_sql/schema/query_schema.cc | 104 ---- .../perfetto_sql/schema/query_schema.h | 38 -- .../span_join_operator/span_join_operator.cc | 45 +- src/trace_processor/sqlite/sqlite_utils.cc | 86 +-- src/trace_processor/sqlite/sqlite_utils.h | 14 +- .../sqlite/sqlite_utils_unittest.cc | 33 +- 46 files changed, 915 insertions(+), 1513 deletions(-) delete mode 100644 src/trace_processor/perfetto_sql/exec/BUILD.gn delete mode 100644 src/trace_processor/perfetto_sql/exec/sql_scan.cc delete mode 100644 src/trace_processor/perfetto_sql/exec/sql_scan.h delete mode 100644 src/trace_processor/perfetto_sql/exec/sql_scan_unittest.cc delete mode 100644 src/trace_processor/perfetto_sql/schema/query_schema.cc delete mode 100644 src/trace_processor/perfetto_sql/schema/query_schema.h diff --git a/Android.bp b/Android.bp index 741a27cb5b6..8f070f2dc3d 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", @@ -24531,8 +24510,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", @@ -25691,7 +25668,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", @@ -26409,7 +26385,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 91e24c40529..6740c5841b2 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", ], ) @@ -11857,7 +11844,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", @@ -12213,7 +12199,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..73ddf17c46a 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) { 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..7702ebb8799 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,15 @@ 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(/*finalized=*/false, std::move(column_names_), + std::move(columns), static_cast(row_count), + string_pool_); + if (analyze_) { + dataframe.Finalize(); + } else { + 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..1d7d8ed8e7f 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. @@ -451,6 +456,9 @@ class Dataframe { uint32_t row_count, StringPool* string_pool); + // Finalize, estimating distinct counts only if `estimate_distinct`. + void FinalizeColumns(bool estimate_distinct); + template PERFETTO_ALWAYS_INLINE void InsertUncheckedInternal( std::index_sequence, diff --git a/src/trace_processor/perfetto_sql/engine/connection_catalog.cc b/src/trace_processor/perfetto_sql/engine/connection_catalog.cc index 8cf1357d0d5..5867b82a4f3 100644 --- a/src/trace_processor/perfetto_sql/engine/connection_catalog.cc +++ b/src/trace_processor/perfetto_sql/engine/connection_catalog.cc @@ -22,17 +22,16 @@ #include #include -#include "perfetto/ext/base/status_or.h" -#include "perfetto/ext/base/string_utils.h" +#include + #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 { @@ -55,19 +54,6 @@ constexpr char kFindViewSql[] = R"( LIMIT 1 )"; -// Escapes |name| as a single-quoted SQL string literal. -std::string Quoted(std::string_view name) { - std::string out = "'"; - for (char c : name) { - if (c == '\'') { - out.push_back('\''); - } - out.push_back(c); - } - out.push_back('\''); - return out; -} - } // namespace ConnectionCatalog::ConnectionCatalog(PerfettoSqlConnection* connection) @@ -77,7 +63,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; @@ -95,8 +93,12 @@ std::optional ConnectionCatalog::FindViewSql( std::string_view name) const { SqliteConnection::PreparedStatement stmt = connection_->sqlite_connection()->PrepareStatement( - SqlSource::FromTraceProcessorImplementation( - base::ReplaceAll(kFindViewSql, "$name", Quoted(name)))); + SqlSource::FromTraceProcessorImplementation(kFindViewSql)); + // A named parameter binds every occurrence of $name in the query. + sqlite3_stmt* raw = stmt.sqlite_stmt(); + sqlite3_bind_text(raw, sqlite3_bind_parameter_index(raw, "$name"), + name.data(), static_cast(name.size()), + sqlite::utils::kSqliteTransient); if (!stmt.Step()) { return std::nullopt; } @@ -109,10 +111,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..1af1c4b518b 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,34 @@ 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"); } + // Only alive while Load runs: once bound, the plan shares ownership of each + // column it reads, so nothing needs to keep the dataframes alive. + std::vector inputs; + 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"); + } + for (const SharedDataframe& input : list->dataframes) { + 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 +250,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::utils::SetError(ctx, kName, dataframe.status()); + } + 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 +335,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 +382,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 f327019220f..675a25eed89 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; 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..9b7fced40ce 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,39 @@ 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); + if (IsPresent(n->schema)) { + return Err(from, Error::kUnsupported, + "reading a relation whose columns cannot be worked out", + " (a schema-qualified table)"); + } + 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()) { + return Err(from, Error::kUnsupported, + "reading a relation whose columns cannot be worked out", + " (" + lineage.status().message() + ")"); } 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 aa4555d3ef0..3537fbe5e48 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" @@ -30,8 +31,27 @@ namespace perfetto::trace_processor::pipeline { // kMaxPipelineColumns columns. 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"; + +// 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 fb5eacd4cb6..a92d60ed8d6 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" @@ -197,10 +199,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()); @@ -304,24 +307,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; @@ -380,6 +377,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) { @@ -397,21 +405,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(); } @@ -424,6 +427,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 28db0710962..6485649d26a 100644 --- a/src/trace_processor/perfetto_sql/pipeline/plan_serialization.h +++ b/src/trace_processor/perfetto_sql/pipeline/plan_serialization.h @@ -20,8 +20,12 @@ #include #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" @@ -40,6 +44,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 564ecb45062..cb747da3d78 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