From 491b1b9282e9f0c2f273c3bf223fe7587004bb02 Mon Sep 17 00:00:00 2001 From: Bill Guowei Yang Date: Wed, 29 Jul 2026 10:05:18 -0400 Subject: [PATCH] perf: add paired query catalog expansion --- tests/perf/README.md | 56 +++++++ tests/perf/core/catalog.go | 258 +++++++++++++++++++++++++++++++- tests/perf/core/catalog_test.go | 213 +++++++++++++++++++++++++- tests/perf/core/runner_test.go | 43 +++++- tests/perf/core/sink_test.go | 56 +++++++ tests/perf/core/types.go | 23 ++- 6 files changed, 639 insertions(+), 10 deletions(-) diff --git a/tests/perf/README.md b/tests/perf/README.md index 51edf11d..53230904 100644 --- a/tests/perf/README.md +++ b/tests/perf/README.md @@ -2,6 +2,62 @@ This package contains the golden-query performance harness. +## Paired Query Catalogs + +Existing catalogs continue to use `queries:` unchanged. A catalog may contain +legacy `queries:`, `paired_queries:`, or both. Paired definitions let one +semantic SQL template run against the frozen raw Parquet views and the +production-shaped DuckLake tables without changing the runner or artifact +contracts: + +```yaml +relation_variants: + raw_view: + events: frozen_v1.events_file_view + persons: frozen_v1.persons_file_view + managed_table: + events: posthog.events + persons: posthog.persons + +paired_queries: + - query_id_base: q_events_daily + intent_id: ph.events.daily.v1 + tags: [posthog, events, time-series] + params: {} + sql_template: | + SELECT date_trunc('day', "timestamp") AS day, COUNT(*) AS events + FROM {{ relation "events" }} + WHERE "timestamp" >= TIMESTAMPTZ '2026-03-01 00:00:00+00' + AND "timestamp" < TIMESTAMPTZ '2026-03-18 00:00:00+00' + GROUP BY 1 + ORDER BY 1 +``` + +Paired catalogs must declare exactly the `raw_view` and `managed_table` +variants. A template expands in declaration order, with `raw_view` before +`managed_table`, into `q_events_daily__raw_view` and +`q_events_daily__managed_table`. Generated queries retain the same +`intent_id`, tags, parameters, and semantic template; only declared relation +placeholders differ. They carry in-memory storage-target metadata, so later +code does not need to infer the target from the generated ID. Legacy queries +remain unpaired. The v1 artifact and publisher schemas remain unchanged, so +artifact rows distinguish paired targets only by these generated query IDs; +they do not include a storage-target column. + +Templating is intentionally limited to `{{ relation "" }}`. Each role +must have a binding in both variants, and multiple roles may be used in one +template. Bindings are unquoted, dot-separated identifiers such as +`posthog.events`; the loader validates every identifier segment and emits it +as a safely quoted relation. SQL expressions, comments, semicolons, +whitespace, quoted identifiers, and malformed names are rejected in bindings; +all template actions other than the relation placeholder are rejected. The +rendered SQL must be a single read-only `SELECT` statement and is copied into +both current protocol SQL fields. + +This is catalog abstraction only. Fair scheduling, migration of the real +frozen PostHog query catalog, paired artifacts, dashboards, and Grafana work +are deliberately deferred to later PRs. + ## Local Smoke Run ```bash diff --git a/tests/perf/core/catalog.go b/tests/perf/core/catalog.go index 022fc3b8..59ef70bb 100644 --- a/tests/perf/core/catalog.go +++ b/tests/perf/core/catalog.go @@ -3,11 +3,41 @@ package core import ( "fmt" "os" + "regexp" "strings" "gopkg.in/yaml.v3" ) +var ( + relationPlaceholderRE = regexp.MustCompile(`\{\{\s*relation\s+"([A-Za-z_][A-Za-z0-9_]*)"\s*\}\}`) + identifierPartRE = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`) +) + +type catalogFile struct { + Name string `yaml:"name"` + Description string `yaml:"description"` + Seed int64 `yaml:"seed"` + DatasetScale int `yaml:"dataset_scale"` + Targets []Protocol `yaml:"targets"` + WarmupIterations int `yaml:"warmup_iterations"` + MeasureIterations int `yaml:"measure_iterations"` + RelationVariants map[StorageTarget]map[string]string `yaml:"relation_variants"` +} + +type pairedQueryDefinition struct { + QueryIDBase string `yaml:"query_id_base"` + IntentID string `yaml:"intent_id"` + Tags []string `yaml:"tags"` + Params map[string]any `yaml:"params"` + SQLTemplate string `yaml:"sql_template"` +} + +type catalogEntry struct { + legacy *Query + paired *pairedQueryDefinition +} + func LoadCatalog(path string) (Catalog, error) { b, err := os.ReadFile(path) if err != nil { @@ -17,16 +47,179 @@ func LoadCatalog(path string) (Catalog, error) { } func ParseCatalog(raw []byte) (Catalog, error) { - var c Catalog - if err := yaml.Unmarshal(raw, &c); err != nil { + var file catalogFile + if err := yaml.Unmarshal(raw, &file); err != nil { return Catalog{}, fmt.Errorf("parse catalog: %w", err) } + c := Catalog{ + Name: file.Name, + Description: file.Description, + Seed: file.Seed, + DatasetScale: file.DatasetScale, + Targets: file.Targets, + WarmupIterations: file.WarmupIterations, + MeasureIterations: file.MeasureIterations, + } + entries, err := catalogEntries(raw) + if err != nil { + return Catalog{}, err + } + if len(entries) > 0 { + if err := validateRelationVariants(file.RelationVariants, entries); err != nil { + return Catalog{}, err + } + } + for _, entry := range entries { + switch { + case entry.legacy != nil: + c.Queries = append(c.Queries, *entry.legacy) + case entry.paired != nil: + queries, err := expandPairedQuery(*entry.paired, file.RelationVariants) + if err != nil { + return Catalog{}, err + } + c.Queries = append(c.Queries, queries...) + } + } if err := validateCatalog(c); err != nil { return Catalog{}, err } return c, nil } +// catalogEntries preserves the declaration order of legacy and paired lists +// when they are mixed in a YAML mapping. The runtime still receives only the +// expanded Catalog.Queries slice. +func catalogEntries(raw []byte) ([]catalogEntry, error) { + var document yaml.Node + if err := yaml.Unmarshal(raw, &document); err != nil { + return nil, fmt.Errorf("parse catalog: %w", err) + } + if len(document.Content) != 1 || document.Content[0].Kind != yaml.MappingNode { + return nil, fmt.Errorf("parse catalog: expected a mapping") + } + mapping := document.Content[0] + var entries []catalogEntry + for index := 0; index < len(mapping.Content); index += 2 { + key, value := mapping.Content[index], mapping.Content[index+1] + switch key.Value { + case "queries": + var queries []Query + if err := value.Decode(&queries); err != nil { + return nil, fmt.Errorf("parse legacy queries: %w", err) + } + for i := range queries { + entries = append(entries, catalogEntry{legacy: &queries[i]}) + } + case "paired_queries": + var paired []pairedQueryDefinition + if err := value.Decode(&paired); err != nil { + return nil, fmt.Errorf("parse paired queries: %w", err) + } + for i := range paired { + entries = append(entries, catalogEntry{paired: &paired[i]}) + } + } + } + return entries, nil +} + +func validateRelationVariants(variants map[StorageTarget]map[string]string, entries []catalogEntry) error { + hasPairedQueries := false + for _, entry := range entries { + if entry.paired != nil { + hasPairedQueries = true + break + } + } + if !hasPairedQueries { + return nil + } + if len(variants) != 2 { + return fmt.Errorf("paired catalogs must declare exactly the raw_view and managed_table storage variants") + } + for _, target := range []StorageTarget{StorageTargetRawView, StorageTargetManagedTable} { + if _, ok := variants[target]; !ok { + return fmt.Errorf("paired catalogs must declare exactly the raw_view and managed_table storage variants") + } + } + return nil +} + +func expandPairedQuery(def pairedQueryDefinition, variants map[StorageTarget]map[string]string) ([]Query, error) { + if def.QueryIDBase == "" { + return nil, fmt.Errorf("paired query missing query_id_base") + } + if def.IntentID == "" { + return nil, fmt.Errorf("paired query %s missing intent_id", def.QueryIDBase) + } + matches := relationPlaceholderRE.FindAllStringSubmatch(def.SQLTemplate, -1) + remaining := relationPlaceholderRE.ReplaceAllString(def.SQLTemplate, "") + if strings.Contains(remaining, "{{") || strings.Contains(remaining, "}}") { + return nil, fmt.Errorf("paired query %s has unsupported template action", def.QueryIDBase) + } + if len(matches) == 0 { + return nil, fmt.Errorf("paired query %s must contain at least one relation placeholder", def.QueryIDBase) + } + + queries := make([]Query, 0, 2) + for _, target := range []StorageTarget{StorageTargetRawView, StorageTargetManagedTable} { + rendered, err := renderRelationTemplate(def.QueryIDBase, def.SQLTemplate, matches, variants[target], target) + if err != nil { + return nil, err + } + queryID := def.QueryIDBase + "__" + string(target) + if err := validateSelectOnlySQL("sql_template", queryID, rendered); err != nil { + return nil, err + } + queries = append(queries, Query{ + QueryID: queryID, + IntentID: def.IntentID, + Tags: def.Tags, + Params: def.Params, + PGWireSQL: rendered, + DuckhogSQL: rendered, + StorageTarget: target, + }) + } + return queries, nil +} + +func renderRelationTemplate(queryID, template string, matches [][]string, bindings map[string]string, target StorageTarget) (string, error) { + replacements := make(map[string]string, len(matches)) + for _, match := range matches { + role := match[1] + binding, ok := bindings[role] + if !ok || binding == "" { + return "", fmt.Errorf("paired query %s missing relation binding for role %q in storage target %q", queryID, role, target) + } + quoted, err := quoteRelationIdentifier(binding) + if err != nil { + return "", fmt.Errorf("paired query %s has invalid relation identifier for role %q in storage target %q: %w", queryID, role, target, err) + } + replacements[role] = quoted + } + return relationPlaceholderRE.ReplaceAllStringFunc(template, func(placeholder string) string { + role := relationPlaceholderRE.FindStringSubmatch(placeholder)[1] + return replacements[role] + }), nil +} + +func quoteRelationIdentifier(identifier string) (string, error) { + parts := strings.Split(identifier, ".") + if len(parts) == 0 { + return "", fmt.Errorf("empty identifier") + } + quoted := make([]string, 0, len(parts)) + for _, part := range parts { + if !identifierPartRE.MatchString(part) { + return "", fmt.Errorf("%q is not a dot-separated identifier", identifier) + } + quoted = append(quoted, `"`+part+`"`) + } + return strings.Join(quoted, "."), nil +} + func validateCatalog(c Catalog) error { if c.Name == "" { return fmt.Errorf("catalog name is required") @@ -108,9 +301,70 @@ func validateSelectOnlySQL(field, queryID, sql string) error { if strings.Contains(trimmed, ";") { return fmt.Errorf("query %s %s must contain a single SELECT statement in frozen mode", queryID, field) } + if containsSQLKeyword(trimmed, "INTO") { + return fmt.Errorf("query %s %s must be SELECT-only in frozen mode", queryID, field) + } return nil } +func containsSQLKeyword(sql, keyword string) bool { + for i := 0; i < len(sql); { + switch { + case sql[i] == '\'': + i = skipQuotedSQLString(sql, i, '\'') + case sql[i] == '"': + i = skipQuotedSQLString(sql, i, '"') + case i+1 < len(sql) && sql[i] == '-' && sql[i+1] == '-': + i += 2 + for i < len(sql) && sql[i] != '\n' { + i++ + } + case i+1 < len(sql) && sql[i] == '/' && sql[i+1] == '*': + i += 2 + for i+1 < len(sql) && (sql[i] != '*' || sql[i+1] != '/') { + i++ + } + if i+1 < len(sql) { + i += 2 + } + case isSQLIdentifierStart(sql[i]): + start := i + i++ + for i < len(sql) && isSQLIdentifierPart(sql[i]) { + i++ + } + if strings.EqualFold(sql[start:i], keyword) { + return true + } + default: + i++ + } + } + return false +} + +func skipQuotedSQLString(sql string, start int, quote byte) int { + for i := start + 1; i < len(sql); i++ { + if sql[i] != quote { + continue + } + if i+1 < len(sql) && sql[i+1] == quote { + i++ + continue + } + return i + 1 + } + return len(sql) +} + +func isSQLIdentifierStart(ch byte) bool { + return ch == '_' || (ch >= 'A' && ch <= 'Z') || (ch >= 'a' && ch <= 'z') +} + +func isSQLIdentifierPart(ch byte) bool { + return isSQLIdentifierStart(ch) || (ch >= '0' && ch <= '9') || ch == '$' +} + func trimLeadingSQLComments(sql string) string { remaining := strings.TrimSpace(sql) for { diff --git a/tests/perf/core/catalog_test.go b/tests/perf/core/catalog_test.go index a4c8005b..fffa4d1d 100644 --- a/tests/perf/core/catalog_test.go +++ b/tests/perf/core/catalog_test.go @@ -1,6 +1,7 @@ package core import ( + "reflect" "strings" "testing" ) @@ -36,6 +37,166 @@ queries: if catalog.Queries[0].QueryID != "q1" || catalog.Queries[0].IntentID != "i1" { t.Fatalf("unexpected query identity: %+v", catalog.Queries[0]) } + if catalog.Queries[0].StorageTarget != "" { + t.Fatalf("legacy query unexpectedly has storage target %q", catalog.Queries[0].StorageTarget) + } +} + +func TestParseCatalogExpandsPairedQueriesIntoStorageTargets(t *testing.T) { + catalog, err := ParseCatalog([]byte(pairedCatalogYAML(` +paired_queries: + - query_id_base: q_events_daily + intent_id: ph.events.daily.v1 + tags: [posthog, events, time-series] + params: + tenant_id: 42 + sql_template: | + SELECT date_trunc('day', "timestamp") AS day, COUNT(*) AS events + FROM {{ relation "events" }} + WHERE "timestamp" >= TIMESTAMPTZ '2026-03-01 00:00:00+00' + GROUP BY 1 + ORDER BY 1 +`))) + if err != nil { + t.Fatalf("ParseCatalog returned error: %v", err) + } + if got, want := queryIDs(catalog), []string{"q_events_daily__raw_view", "q_events_daily__managed_table"}; !reflect.DeepEqual(got, want) { + t.Fatalf("unexpected generated query order: got %v want %v", got, want) + } + for _, query := range catalog.Queries { + if query.IntentID != "ph.events.daily.v1" { + t.Fatalf("generated query did not retain intent_id: %+v", query) + } + if !reflect.DeepEqual(query.Tags, []string{"posthog", "events", "time-series"}) || !reflect.DeepEqual(query.Params, map[string]any{"tenant_id": 42}) { + t.Fatalf("generated query did not retain shared metadata: %+v", query) + } + if query.PGWireSQL != query.DuckhogSQL { + t.Fatalf("generated SQL must be copied into both protocol fields: %+v", query) + } + } + if got, want := catalog.Queries[0].StorageTarget, StorageTargetRawView; got != want { + t.Fatalf("raw query target: got %q want %q", got, want) + } + if got, want := catalog.Queries[1].StorageTarget, StorageTargetManagedTable; got != want { + t.Fatalf("managed query target: got %q want %q", got, want) + } + if got, want := catalog.Queries[0].PGWireSQL, "SELECT date_trunc('day', \"timestamp\") AS day, COUNT(*) AS events\nFROM \"frozen_v1\".\"events_file_view\"\nWHERE \"timestamp\" >= TIMESTAMPTZ '2026-03-01 00:00:00+00'\nGROUP BY 1\nORDER BY 1\n"; got != want { + t.Fatalf("raw query SQL: got %q want %q", got, want) + } + if got, want := catalog.Queries[1].PGWireSQL, "SELECT date_trunc('day', \"timestamp\") AS day, COUNT(*) AS events\nFROM \"posthog\".\"events\"\nWHERE \"timestamp\" >= TIMESTAMPTZ '2026-03-01 00:00:00+00'\nGROUP BY 1\nORDER BY 1\n"; got != want { + t.Fatalf("managed query SQL: got %q want %q", got, want) + } +} + +func TestParseCatalogExpandsMultipleRelationsInDeclarationOrder(t *testing.T) { + catalog, err := ParseCatalog([]byte(pairedCatalogYAML(` +paired_queries: + - query_id_base: q_join + intent_id: ph.join.v1 + tags: [posthog] + params: {} + sql_template: SELECT COUNT(*) FROM {{ relation "events" }} e JOIN {{ relation "persons" }} p ON e.person_id = p.id + - query_id_base: q_events + intent_id: ph.events.v1 + tags: [posthog] + params: {} + sql_template: SELECT COUNT(*) FROM {{ relation "events" }} +queries: + - query_id: legacy_after + intent_id: legacy.intent + pgwire_sql: SELECT 1 + duckhog_sql: SELECT 1 +`))) + if err != nil { + t.Fatalf("ParseCatalog returned error: %v", err) + } + if got, want := queryIDs(catalog), []string{"q_join__raw_view", "q_join__managed_table", "q_events__raw_view", "q_events__managed_table", "legacy_after"}; !reflect.DeepEqual(got, want) { + t.Fatalf("unexpected mixed catalog order: got %v want %v", got, want) + } + if got, want := catalog.Queries[0].PGWireSQL, "SELECT COUNT(*) FROM \"frozen_v1\".\"events_file_view\" e JOIN \"frozen_v1\".\"persons_file_view\" p ON e.person_id = p.id"; got != want { + t.Fatalf("raw multi-relation SQL: got %q want %q", got, want) + } + if got, want := catalog.Queries[1].PGWireSQL, "SELECT COUNT(*) FROM \"posthog\".\"events\" e JOIN \"posthog\".\"persons\" p ON e.person_id = p.id"; got != want { + t.Fatalf("managed multi-relation SQL: got %q want %q", got, want) + } +} + +func TestParseCatalogRejectsInvalidPairedDefinitions(t *testing.T) { + tests := []struct { + name string + yaml string + want string + }{ + {name: "missing variants", yaml: catalogYAML("relation_variants: {}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "storage variants"}, + {name: "invalid variant", yaml: catalogYAML("relation_variants:\n raw_view: {events: frozen_v1.events_file_view}\n archive: {events: posthog.events}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "storage variants"}, + {name: "missing base id", yaml: pairedCatalogYAML("paired_queries:\n - intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "query_id_base"}, + {name: "missing intent", yaml: pairedCatalogYAML("paired_queries:\n - query_id_base: q\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "intent_id"}, + {name: "missing binding", yaml: catalogYAML("relation_variants:\n raw_view: {events: frozen_v1.events_file_view}\n managed_table: {persons: posthog.persons}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "missing relation binding"}, + {name: "unknown binding", yaml: pairedCatalogYAML("paired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"orders\" }}\n"), want: "missing relation binding"}, + {name: "malicious identifier", yaml: catalogYAML("relation_variants:\n raw_view: {events: frozen_v1.events_file_view}\n managed_table: {events: 'posthog.events; DROP TABLE posthog.events'}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "invalid relation identifier"}, + {name: "whitespace identifier", yaml: catalogYAML("relation_variants:\n raw_view: {events: frozen_v1.events_file_view}\n managed_table: {events: 'posthog. events'}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "invalid relation identifier"}, + {name: "comment identifier", yaml: catalogYAML("relation_variants:\n raw_view: {events: frozen_v1.events_file_view}\n managed_table: {events: 'posthog.events -- managed table'}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "invalid relation identifier"}, + {name: "expression identifier", yaml: catalogYAML("relation_variants:\n raw_view: {events: frozen_v1.events_file_view}\n managed_table: {events: 'lower(posthog.events)'}\npaired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ relation \"events\" }}\n"), want: "invalid relation identifier"}, + {name: "unsupported action", yaml: pairedCatalogYAML("paired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * FROM {{ .Events }}\n"), want: "unsupported template action"}, + {name: "no placeholder", yaml: pairedCatalogYAML("paired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT 1\n"), want: "relation placeholder"}, + {name: "rendered write", yaml: pairedCatalogYAML("paired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: INSERT INTO {{ relation \"events\" }} VALUES (1)\n"), want: "SELECT-only"}, + {name: "rendered select into", yaml: pairedCatalogYAML("paired_queries:\n - query_id_base: q\n intent_id: i\n sql_template: SELECT * INTO derived_events FROM {{ relation \"events\" }}\n"), want: "SELECT-only"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, err := ParseCatalog([]byte(tt.yaml)) + if err == nil || !strings.Contains(err.Error(), tt.want) { + t.Fatalf("ParseCatalog error = %v, want substring %q", err, tt.want) + } + }) + } +} + +func TestParseCatalogRejectsGeneratedQueryIDCollisions(t *testing.T) { + tests := []struct { + name string + yaml string + }{ + {name: "explicit legacy id", yaml: pairedCatalogYAML(` +queries: + - query_id: q_events__raw_view + intent_id: legacy.intent + pgwire_sql: SELECT 1 + duckhog_sql: SELECT 1 +paired_queries: + - query_id_base: q_events + intent_id: paired.intent + sql_template: SELECT * FROM {{ relation "events" }} +`)}, + {name: "two paired bases", yaml: pairedCatalogYAML(` +paired_queries: + - query_id_base: q_events + intent_id: one + sql_template: SELECT * FROM {{ relation "events" }} + - query_id_base: q_events + intent_id: two + sql_template: SELECT * FROM {{ relation "events" }} +`)}, + {name: "legacy duplicate", yaml: pairedCatalogYAML(` +queries: + - query_id: q_legacy + intent_id: one + pgwire_sql: SELECT 1 + duckhog_sql: SELECT 1 + - query_id: q_legacy + intent_id: two + pgwire_sql: SELECT 2 + duckhog_sql: SELECT 2 +`)}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, err := ParseCatalog([]byte(tt.yaml)) + if err == nil || !strings.Contains(err.Error(), "duplicate query_id") { + t.Fatalf("ParseCatalog error = %v, want duplicate query_id", err) + } + }) + } } func TestParseCatalogRejectsDuplicateQueryIDs(t *testing.T) { @@ -73,7 +234,7 @@ func TestValidateReadOnlyCatalogAcceptsSelectOnlyQueries(t *testing.T) { { QueryID: "q1", IntentID: "i1", - PGWireSQL: "SELECT 1;", + PGWireSQL: "SELECT $$into$$ AS label;", DuckhogSQL: "/* comment */ SELECT 1", }, }, @@ -103,3 +264,53 @@ func TestValidateReadOnlyCatalogRejectsNonSelectQueries(t *testing.T) { t.Fatalf("expected SELECT-only error, got %v", err) } } + +func TestArtifactSinkContractRemainsUnchangedForPairedQueries(t *testing.T) { + catalog, err := ParseCatalog([]byte(pairedCatalogYAML(` +paired_queries: + - query_id_base: q_events + intent_id: ph.events.v1 + sql_template: SELECT COUNT(*) FROM {{ relation "events" }} +`))) + if err != nil { + t.Fatalf("ParseCatalog returned error: %v", err) + } + if got, want := queryIDs(catalog), []string{"q_events__raw_view", "q_events__managed_table"}; !reflect.DeepEqual(got, want) { + t.Fatalf("unexpected runtime query representation: got %v want %v", got, want) + } + if got := catalog.Queries[0].StorageTarget; got != StorageTargetRawView { + t.Fatalf("unexpected storage target %q", got) + } +} + +func pairedCatalogYAML(body string) string { + return catalogYAML(strings.TrimSuffix(` +relation_variants: + raw_view: + events: frozen_v1.events_file_view + persons: frozen_v1.persons_file_view + managed_table: + events: posthog.events + persons: posthog.persons + `, "\t") + body) +} + +func catalogYAML(body string) string { + return ` +name: paired +description: paired query catalog +seed: 1 +dataset_scale: 1 +targets: [pgwire] +warmup_iterations: 0 +measure_iterations: 1 +` + body +} + +func queryIDs(c Catalog) []string { + ids := make([]string, 0, len(c.Queries)) + for _, query := range c.Queries { + ids = append(ids, query.QueryID) + } + return ids +} diff --git a/tests/perf/core/runner_test.go b/tests/perf/core/runner_test.go index 1296b174..8ffdb88c 100644 --- a/tests/perf/core/runner_test.go +++ b/tests/perf/core/runner_test.go @@ -2,6 +2,7 @@ package core import ( "context" + "reflect" "testing" "time" ) @@ -9,15 +10,55 @@ import ( type testDriver struct { protocol Protocol calls int + queryIDs []string } func (d *testDriver) Protocol() Protocol { return d.protocol } -func (d *testDriver) Execute(_ context.Context, _ Query, _ []any) (ExecutionResult, error) { +func (d *testDriver) Execute(_ context.Context, query Query, _ []any) (ExecutionResult, error) { d.calls++ + d.queryIDs = append(d.queryIDs, query.QueryID) return ExecutionResult{Rows: 1}, nil } +func TestRunnerExecutesPairedQueriesThroughExistingRuntimeContract(t *testing.T) { + catalog, err := ParseCatalog([]byte(pairedCatalogYAML(` +paired_queries: + - query_id_base: q_events + intent_id: ph.events.v1 + tags: [posthog] + params: {} + sql_template: SELECT COUNT(*) FROM {{ relation "events" }} +`))) + if err != nil { + t.Fatalf("ParseCatalog returned error: %v", err) + } + driver := &testDriver{protocol: ProtocolPGWire} + sink := &inMemorySink{} + runner := NewQueryRunner(RunnerConfig{ + Catalog: catalog, + Drivers: map[Protocol]ProtocolDriver{ + ProtocolPGWire: driver, + }, + Sink: sink, + Now: func() time.Time { return time.Unix(1700000000, 0) }, + }) + if _, err := runner.Run(context.Background()); err != nil { + t.Fatalf("Run returned error: %v", err) + } + wantIDs := []string{"q_events__raw_view", "q_events__managed_table"} + if !reflect.DeepEqual(driver.queryIDs, wantIDs) { + t.Fatalf("driver query order: got %v want %v", driver.queryIDs, wantIDs) + } + gotIDs := make([]string, 0, len(sink.results)) + for _, result := range sink.results { + gotIDs = append(gotIDs, result.QueryID) + } + if !reflect.DeepEqual(gotIDs, wantIDs) { + t.Fatalf("result query IDs: got %v want %v", gotIDs, wantIDs) + } +} + func (d *testDriver) Close() error { return nil } type inMemorySink struct { diff --git a/tests/perf/core/sink_test.go b/tests/perf/core/sink_test.go index f5625288..65b29eb6 100644 --- a/tests/perf/core/sink_test.go +++ b/tests/perf/core/sink_test.go @@ -1,9 +1,11 @@ package core import ( + "encoding/csv" "encoding/json" "os" "path/filepath" + "reflect" "strings" "testing" "time" @@ -86,3 +88,57 @@ func TestArtifactSinkWritesSummaryCSVAndMetrics(t *testing.T) { t.Fatalf("csv rows missing measure_iteration values: %q", csvText) } } + +func TestPairedQueriesPreserveArtifactCSVContract(t *testing.T) { + catalog, err := ParseCatalog([]byte(pairedCatalogYAML(` +paired_queries: + - query_id_base: q_events + intent_id: ph.events.v1 + sql_template: SELECT COUNT(*) FROM {{ relation "events" }} +`))) + if err != nil { + t.Fatalf("ParseCatalog returned error: %v", err) + } + dir := t.TempDir() + sink, err := NewArtifactSink(dir) + if err != nil { + t.Fatalf("NewArtifactSink returned error: %v", err) + } + for _, query := range catalog.Queries { + if err := sink.Record(QueryResult{ + QueryID: query.QueryID, + IntentID: query.IntentID, + MeasureIteration: 1, + Protocol: ProtocolPGWire, + Status: "ok", + Rows: 1, + Duration: time.Millisecond, + StartedAt: time.Unix(1700000000, 0), + }); err != nil { + t.Fatalf("Record returned error: %v", err) + } + } + if err := sink.Close(RunSummary{}, ""); err != nil { + t.Fatalf("Close returned error: %v", err) + } + file, err := os.Open(filepath.Join(dir, "query_results.csv")) + if err != nil { + t.Fatalf("open query_results.csv: %v", err) + } + defer func() { + if err := file.Close(); err != nil { + t.Errorf("close query_results.csv: %v", err) + } + }() + records, err := csv.NewReader(file).ReadAll() + if err != nil { + t.Fatalf("read query_results.csv: %v", err) + } + wantHeader := []string{"query_id", "intent_id", "measure_iteration", "protocol", "status", "error", "error_class", "rows", "duration_ms", "started_at"} + if !reflect.DeepEqual(records[0], wantHeader) { + t.Fatalf("CSV header: got %v want %v", records[0], wantHeader) + } + if got, want := []string{records[1][0], records[2][0]}, []string{"q_events__raw_view", "q_events__managed_table"}; !reflect.DeepEqual(got, want) { + t.Fatalf("CSV query IDs: got %v want %v", got, want) + } +} diff --git a/tests/perf/core/types.go b/tests/perf/core/types.go index b0d97e22..bfc22b6c 100644 --- a/tests/perf/core/types.go +++ b/tests/perf/core/types.go @@ -9,6 +9,16 @@ const ( ProtocolFlight Protocol = "flight" ) +// StorageTarget identifies the physical relation family selected for a paired +// catalog query. It is runtime-only metadata; artifacts continue to use the +// existing query ID and intent ID fields. +type StorageTarget string + +const ( + StorageTargetRawView StorageTarget = "raw_view" + StorageTargetManagedTable StorageTarget = "managed_table" +) + type Catalog struct { Name string `yaml:"name"` Description string `yaml:"description"` @@ -21,12 +31,13 @@ type Catalog struct { } type Query struct { - QueryID string `yaml:"query_id"` - IntentID string `yaml:"intent_id"` - Tags []string `yaml:"tags"` - Params map[string]any `yaml:"params"` - PGWireSQL string `yaml:"pgwire_sql"` - DuckhogSQL string `yaml:"duckhog_sql"` + QueryID string `yaml:"query_id"` + IntentID string `yaml:"intent_id"` + Tags []string `yaml:"tags"` + Params map[string]any `yaml:"params"` + PGWireSQL string `yaml:"pgwire_sql"` + DuckhogSQL string `yaml:"duckhog_sql"` + StorageTarget StorageTarget `yaml:"-" json:"-"` } type ExecutionResult struct {