diff --git a/.golangci.yml b/.golangci.yml index 7b8b3d46..8327c6e9 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -211,7 +211,7 @@ linters: path: _test\.go - linters: - mnd - path: ^internal/workloads/ + path: ^workloads/(execute_sql|simple|tpcb|tpcc|tpcds|tpch)/ - linters: - unused path: bind.go diff --git a/AGENTS.md b/AGENTS.md index bc268e42..6d0edbaa 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -58,9 +58,10 @@ with `go generate ./pkg/config`. There is no application protobuf generation ste and protobuf remains only as an indirect dependency of external SDKs. Load-time generation is plain Go under `pkg/gen/` (see `docs/parallelism.md`). -**Embedded FS rebuild rule:** `workloads/` is `//go:embed *` (SQL/JSON/MD only). +**Embedded FS rebuild rule:** each asset-bearing package under `workloads//` +embeds its own SQL/JSON/README files and registers them with the shared catalog. If you pass a workload by short name (`tpcc/tx`), the binary serves its SQL from -the embedded snapshot. Edits to `workloads/*.sql` on disk have **no effect** until +the embedded snapshot. Edits to `workloads//*.sql` on disk have **no effect** until `make build` reruns. **Local path bypass:** If you pass an explicit local `.sql` path @@ -79,7 +80,8 @@ Resolution order for SQL files: **cwd → `~/.stroppy/` → embedded**. | `cmd/stroppy/` | entrypoint (`main.go` blank-imports the drivers) + cobra subcommands: run, probe, version | | `cmd/stroppy/commands/run/` | arg parsing, driver/env/step resolution, dispatch to `bench.Run` | | `pkg/bench/` | Go-native engine: `Workload` interface, `Run`, scenario executor, VU/Bench SDK, metrics sink + summary | -| `internal/workloads/` | the Go workloads (simple, tpcb, tpcc, tpch, tpcds, execute_sql); aggregated by blank import | +| `workloads//` | one package per built-in workload, containing Go implementation/tests plus owned SQL/JSON/README assets | +| `workloads/all/` | explicit blank-import aggregation for built-in workload registration | | `pkg/driver/dispatcher.go` | driver registry: `RegisterDriver()` + `Dispatch()` | | `pkg/driver/{postgres,mysql,picodata,ydb,noop,csv}/` | driver implementations | | `pkg/driver/sqldriver/` | shared sql.DB-backed base (mysql, ydb use this) | @@ -87,7 +89,7 @@ Resolution order for SQL files: **cwd → `~/.stroppy/` → embedded**. | `pkg/datagen/` | row-production seam: `source` (Partitionable/RowSource) + canonical TPC-DS/TPC-H generator adapters (`tpcdsgen`, `tpchgen`) | | `internal/runner/` | run-config merge, env override parsing, driver presets, config-file load | | `pkg/config/` | plain-Go application config types + strict recursive JSON normalizer; schema source for `docs/jsonschema/run.schema.json` | -| `workloads/` | embedded SQL/JSON workloads: tpcb, tpcc, tpch, tpcds | +| `workloads/` | shared embedded-asset registry and catalog | | `docs/parallelism.md` | InsertRequest parallelism contract and tuning | ## Drivers @@ -220,8 +222,11 @@ Section layout (must be identical across dialects): ``` Each Go workload implements the `bench.Workload` interface (`Setup`, `Iterate`, -`Teardown`) in `internal/workloads//`. TPC-B and TPC-C each ship two -registered variants: +`Teardown`) in `workloads//`, beside its tests and owned assets. Asset-bearing +packages embed only their local SQL/JSON/README files, register the filesystem with +`workloads.Register`, and add a blank import to `workloads/all/import.go`. Contract +tests in each package pin required files, sections, and named queries. TPC-B and +TPC-C each ship two registered variants: - `procs` — calls stored procs via the `workload_procs` section; pg + mysql only - `tx` — runs ordered DML steps inside `driver.beginTx()`; all SQL drivers (pg/mysql/pico/ydb) diff --git a/cmd/stroppy/commands/help/topic_datagen.go b/cmd/stroppy/commands/help/topic_datagen.go index 0d6dec36..e73ae071 100644 --- a/cmd/stroppy/commands/help/topic_datagen.go +++ b/cmd/stroppy/commands/help/topic_datagen.go @@ -58,9 +58,9 @@ CSV OUTPUT REFERENCES docs/parallelism.md InsertRequest parallelism contract and tuning - internal/workloads/simple/ Minimal Go workload - internal/workloads/tpcb/ Small relational workload - internal/workloads/tpch/ Stateful canonical generator + workloads/simple/ Minimal Go workload + workloads/tpcb/ Small relational workload + workloads/tpch/ Stateful canonical generator pkg/gen/ Deterministic generation primitives pkg/bench/query.go b.Insert entrypoint `, diff --git a/cmd/stroppy/commands/probe/probe_test.go b/cmd/stroppy/commands/probe/probe_test.go index 5df169da..fbb8fd00 100644 --- a/cmd/stroppy/commands/probe/probe_test.go +++ b/cmd/stroppy/commands/probe/probe_test.go @@ -7,8 +7,8 @@ import ( "strings" "testing" - _ "github.com/stroppy-io/stroppy/internal/workloads" "github.com/stroppy-io/stroppy/pkg/bench" + _ "github.com/stroppy-io/stroppy/workloads/all" ) func TestJSONCatalogIncludesWorkloadSchemas(t *testing.T) { diff --git a/cmd/stroppy/commands/root.go b/cmd/stroppy/commands/root.go index a79d6a94..7538565d 100644 --- a/cmd/stroppy/commands/root.go +++ b/cmd/stroppy/commands/root.go @@ -15,8 +15,8 @@ import ( "github.com/stroppy-io/stroppy/cmd/stroppy/commands/probe" "github.com/stroppy-io/stroppy/cmd/stroppy/commands/run" "github.com/stroppy-io/stroppy/internal/version" - _ "github.com/stroppy-io/stroppy/internal/workloads" "github.com/stroppy-io/stroppy/pkg/common/shutdown" + _ "github.com/stroppy-io/stroppy/workloads/all" ) // appName is the binary / command name. diff --git a/cmd/stroppy/commands/run/run_test.go b/cmd/stroppy/commands/run/run_test.go index ea1a20e3..57359466 100644 --- a/cmd/stroppy/commands/run/run_test.go +++ b/cmd/stroppy/commands/run/run_test.go @@ -14,10 +14,10 @@ import ( "github.com/stretchr/testify/require" "github.com/stroppy-io/stroppy/internal/runner" - _ "github.com/stroppy-io/stroppy/internal/workloads/simple" "github.com/stroppy-io/stroppy/pkg/bench" "github.com/stroppy-io/stroppy/pkg/config" _ "github.com/stroppy-io/stroppy/pkg/driver/noop" + _ "github.com/stroppy-io/stroppy/workloads/simple" ) func unsetLoggerEnv(t *testing.T) { diff --git a/internal/workloads/import.go b/internal/workloads/import.go deleted file mode 100644 index ee2d345c..00000000 --- a/internal/workloads/import.go +++ /dev/null @@ -1,13 +0,0 @@ -// Package workloads aggregates every Go-native workload via blank import so the -// binary pulls them in and their init() registrations run. Add a line per ported -// workload. -package workloads - -import ( - _ "github.com/stroppy-io/stroppy/internal/workloads/execute_sql" - _ "github.com/stroppy-io/stroppy/internal/workloads/simple" - _ "github.com/stroppy-io/stroppy/internal/workloads/tpcb" - _ "github.com/stroppy-io/stroppy/internal/workloads/tpcc" - _ "github.com/stroppy-io/stroppy/internal/workloads/tpcds" - _ "github.com/stroppy-io/stroppy/internal/workloads/tpch" -) diff --git a/pkg/bench/sql_file_test.go b/pkg/bench/sql_file_test.go new file mode 100644 index 00000000..de42d296 --- /dev/null +++ b/pkg/bench/sql_file_test.go @@ -0,0 +1,62 @@ +package bench + +import ( + "os" + "path/filepath" + "sync" + "testing" + "testing/fstest" + + "github.com/stroppy-io/stroppy/workloads" +) + +var registerSQLResolutionPreset sync.Once + +func TestLoadSQLUsesLocalOverrideBeforeEmbeddedFallback(t *testing.T) { + const ( + preset = "bench-sql-resolution-test" + fileName = "dialect.sql" + embeddedBody = "--+ query\n--= body\nSELECT 'embedded';\n" + localBody = "--+ query\n--= body\nSELECT 'local';\n" + ) + + registerSQLResolutionPreset.Do(func() { + workloads.Register(preset, fstest.MapFS{ + fileName: &fstest.MapFile{Data: []byte(embeddedBody)}, + }) + }) + + t.Run("embedded fallback", func(t *testing.T) { + t.Chdir(t.TempDir()) + + assertSQLBody(t, preset, fileName, "SELECT 'embedded';") + }) + + t.Run("local workload file", func(t *testing.T) { + directory := t.TempDir() + t.Chdir(directory) + + if err := os.MkdirAll(filepath.Join("workloads", preset), 0o755); err != nil { + t.Fatal(err) + } + + if err := os.WriteFile(filepath.Join("workloads", preset, fileName), []byte(localBody), 0o600); err != nil { + t.Fatal(err) + } + + assertSQLBody(t, preset, fileName, "SELECT 'local';") + }) +} + +func assertSQLBody(t *testing.T, preset, fileName, want string) { + t.Helper() + + sql, err := LoadSQL(preset, fileName) + if err != nil { + t.Fatal(err) + } + + if got, ok := sql.Query("query", "body"); !ok || got != want { + t.Fatalf("query body = %q, %v; want %q, true", got, ok, want) + } +} diff --git a/workloads/all/import.go b/workloads/all/import.go new file mode 100644 index 00000000..744373ac --- /dev/null +++ b/workloads/all/import.go @@ -0,0 +1,11 @@ +// Package all registers every built-in workload through blank imports. +package all + +import ( + _ "github.com/stroppy-io/stroppy/workloads/execute_sql" + _ "github.com/stroppy-io/stroppy/workloads/simple" + _ "github.com/stroppy-io/stroppy/workloads/tpcb" + _ "github.com/stroppy-io/stroppy/workloads/tpcc" + _ "github.com/stroppy-io/stroppy/workloads/tpcds" + _ "github.com/stroppy-io/stroppy/workloads/tpch" +) diff --git a/workloads/catalog.go b/workloads/catalog.go index ef2a99cd..0afd57e0 100644 --- a/workloads/catalog.go +++ b/workloads/catalog.go @@ -2,6 +2,7 @@ package workloads import ( "fmt" + "io/fs" "sort" "strings" ) @@ -22,9 +23,14 @@ func Catalog() ([]PresetInfo, error) { out := make([]PresetInfo, 0, len(presets)) for _, name := range presets { - entries, err := Content.ReadDir(name) + files, err := presetFiles(Preset(name)) if err != nil { - return nil, fmt.Errorf("%w: %s", ErrUnknownPreset, name) + return nil, err + } + + entries, err := fs.ReadDir(files, ".") + if err != nil { + return nil, fmt.Errorf("read preset %q: %w", name, err) } info := PresetInfo{Name: name} diff --git a/workloads/catalog_test.go b/workloads/catalog_test.go index 91b26204..5bb24d08 100644 --- a/workloads/catalog_test.go +++ b/workloads/catalog_test.go @@ -1,27 +1,45 @@ -package workloads +package workloads_test -import "testing" +import ( + "testing" + + "github.com/stroppy-io/stroppy/workloads" + _ "github.com/stroppy-io/stroppy/workloads/all" +) func TestCatalog(t *testing.T) { - catalog, err := Catalog() + catalog, err := workloads.Catalog() if err != nil { t.Fatalf("Catalog() error: %v", err) } - if len(catalog) != len(AvailablePresets()) { - t.Fatalf("got %d presets, want %d", len(catalog), len(AvailablePresets())) + expected := []workloads.Preset{ + workloads.PresetTPCB, + workloads.PresetTPCC, + workloads.PresetTPCDS, + workloads.PresetTPCH, + } + if len(catalog) != len(expected) { + t.Fatalf("got %d presets, want %d", len(catalog), len(expected)) + } + + byName := make(map[string]workloads.PresetInfo, len(catalog)) + for _, preset := range catalog { + byName[preset.Name] = preset } - byName := make(map[string]PresetInfo, len(catalog)) - for _, p := range catalog { - byName[p.Name] = p + for _, preset := range expected { + if _, ok := byName[string(preset)]; !ok { + t.Errorf("missing preset %q", preset) + } } - if len(byName["tpcc"].SQL) == 0 { + tpcc := byName[string(workloads.PresetTPCC)] + if len(tpcc.SQL) == 0 { t.Error("tpcc: expected SQL dialects, got none") } - if len(byName["tpcc"].Docs) == 0 { + if len(tpcc.Docs) == 0 { t.Error("tpcc: expected docs, got none") } } diff --git a/workloads/embed.go b/workloads/embed.go index 63f8cb90..8771d116 100644 --- a/workloads/embed.go +++ b/workloads/embed.go @@ -1,16 +1,17 @@ -// Package workloads provides embedded SQL/JSON workload files. +// Package workloads provides workload-owned embedded SQL, JSON, and documentation files. package workloads import ( - "embed" "errors" "fmt" "io" + "io/fs" "os" - "path" + "path/filepath" + "sort" ) -// Preset represents an available example preset. +// Preset identifies an embedded workload asset set. type Preset string const ( @@ -23,24 +24,43 @@ const ( // ErrUnknownPreset is returned when an unknown preset name is requested. var ErrUnknownPreset = errors.New("unknown preset") -//go:embed * -var Content embed.FS +var registry = make(map[Preset]fs.FS) -// AvailablePresets returns list of available preset names. +// Register associates a preset with assets embedded during package initialization. +func Register(preset Preset, files fs.FS) { + if files == nil { + panic(fmt.Sprintf("workloads: register %q with nil filesystem", preset)) + } + + if _, exists := registry[preset]; exists { + panic(fmt.Sprintf("workloads: preset %q already registered", preset)) + } + + registry[preset] = files +} + +// AvailablePresets returns registered preset names in sorted order. func AvailablePresets() []string { - return []string{ - string(PresetTPCC), - string(PresetTPCB), - string(PresetTPCDS), - string(PresetTPCH), + presets := make([]string, 0, len(registry)) + for preset := range registry { + presets = append(presets, string(preset)) } + + sort.Strings(presets) + + return presets } -// CopyPresetToPath copies preset files to the target directory. +// CopyPresetToPath copies preset files to target directory. func CopyPresetToPath(targetPath string, preset Preset, perm os.FileMode) error { - entries, err := Content.ReadDir(string(preset)) + files, err := presetFiles(preset) if err != nil { - return fmt.Errorf("%w: %s", ErrUnknownPreset, preset) + return err + } + + entries, err := fs.ReadDir(files, ".") + if err != nil { + return fmt.Errorf("read preset %q: %w", preset, err) } for _, entry := range entries { @@ -48,42 +68,52 @@ func CopyPresetToPath(targetPath string, preset Preset, perm os.FileMode) error continue } - err = copyFileToPath(targetPath, string(preset), entry.Name(), perm) - if err != nil { - return fmt.Errorf("preset '%s' file copy error: %w", preset, err) + if err := copyFileToPath(files, targetPath, entry.Name(), perm); err != nil { + return fmt.Errorf("preset %q file copy: %w", preset, err) } } return nil } -// ReadPresetFile reads a single file from an embedded preset. -// presetName is the preset directory (e.g., "tpcc"), fileName is the file within it (e.g., "tpcc.ts"). +// ReadPresetFile reads one file embedded by a workload package. func ReadPresetFile(presetName, fileName string) ([]byte, error) { - return Content.ReadFile(path.Join(presetName, fileName)) + files, err := presetFiles(Preset(presetName)) + if err != nil { + return nil, err + } + + return fs.ReadFile(files, fileName) } -// copyFileToPath copies a single file from examples to the target directory. -func copyFileToPath(targetPath, preset, fileName string, perm os.FileMode) error { - file, err := Content.Open(path.Join(preset, fileName)) +func presetFiles(preset Preset) (fs.FS, error) { + files, exists := registry[preset] + if !exists { + return nil, fmt.Errorf("%w: %s", ErrUnknownPreset, preset) + } + + return files, nil +} + +func copyFileToPath(files fs.FS, targetPath, fileName string, perm os.FileMode) error { + source, err := files.Open(fileName) if err != nil { - return fmt.Errorf("failed to open file %s: %w", fileName, err) + return fmt.Errorf("open %s: %w", fileName, err) } - defer file.Close() + defer source.Close() - destFile, err := os.OpenFile( - path.Join(targetPath, fileName), + destination, err := os.OpenFile( + filepath.Join(targetPath, fileName), os.O_WRONLY|os.O_CREATE|os.O_TRUNC, perm, ) if err != nil { - return fmt.Errorf("failed to open file %s: %w", fileName, err) + return fmt.Errorf("open %s: %w", fileName, err) } - defer destFile.Close() + defer destination.Close() - _, err = io.Copy(destFile, file) - if err != nil { - return fmt.Errorf("failed to copy file %s to %s: %w", fileName, targetPath, err) + if _, err := io.Copy(destination, source); err != nil { + return fmt.Errorf("copy %s to %s: %w", fileName, targetPath, err) } return nil diff --git a/internal/workloads/execute_sql/execute_sql.go b/workloads/execute_sql/execute_sql.go similarity index 85% rename from internal/workloads/execute_sql/execute_sql.go rename to workloads/execute_sql/execute_sql.go index 779ecfac..2fcc1129 100644 --- a/internal/workloads/execute_sql/execute_sql.go +++ b/workloads/execute_sql/execute_sql.go @@ -1,8 +1,6 @@ -// Package execute_sql is the Go-native port of workloads/execute_sql/execute_sql.ts: -// a generic runner that executes every query in a SQL file (or inline SQL string) once. -// The SQL source is one of two typed workload parameters: --sql-file (a path resolved -// cwd → workloads/execute_sql/ → embedded) or --sql-body (inline SQL text). Queries are -// delimited by `--= name` markers, matching parse_sql.ts — a markerless source yields none. +// Package execute_sql runs every named query from an inline or file SQL source. +// Files resolve from cwd before embedded workload assets; markerless sources contain +// no executable named queries. package execute_sql import ( diff --git a/internal/workloads/execute_sql/execute_sql_test.go b/workloads/execute_sql/execute_sql_test.go similarity index 100% rename from internal/workloads/execute_sql/execute_sql_test.go rename to workloads/execute_sql/execute_sql_test.go diff --git a/workloads/internal/workloadtest/contract.go b/workloads/internal/workloadtest/contract.go new file mode 100644 index 00000000..338a0d4f --- /dev/null +++ b/workloads/internal/workloadtest/contract.go @@ -0,0 +1,49 @@ +// Package workloadtest provides contract assertions for embedded workload assets. +package workloadtest + +import ( + "io/fs" + "testing" + + "github.com/stroppy-io/stroppy/pkg/bench" +) + +// Query identifies one required named query. +type Query struct { + Section string + Name string +} + +// Files asserts that every named asset exists. +func Files(t *testing.T, files fs.FS, names ...string) { + t.Helper() + + for _, name := range names { + if _, err := fs.Stat(files, name); err != nil { + t.Errorf("required asset %q: %v", name, err) + } + } +} + +// SQL asserts required sections and named queries in one embedded SQL asset. +func SQL(t *testing.T, files fs.FS, name string, sections []string, queries []Query) { + t.Helper() + + data, err := fs.ReadFile(files, name) + if err != nil { + t.Fatalf("read %q: %v", name, err) + } + + parsed := bench.ParseSQL(string(data)) + for _, section := range sections { + if len(parsed.Section(section)) == 0 { + t.Errorf("%s: missing or empty section %q", name, section) + } + } + + for _, query := range queries { + if body, ok := parsed.Query(query.Section, query.Name); !ok || body == "" { + t.Errorf("%s: missing or empty query %s/%s", name, query.Section, query.Name) + } + } +} diff --git a/internal/workloads/params_test.go b/workloads/params_test.go similarity index 98% rename from internal/workloads/params_test.go rename to workloads/params_test.go index 6852b62b..56f05a5b 100644 --- a/internal/workloads/params_test.go +++ b/workloads/params_test.go @@ -1,4 +1,4 @@ -package workloads +package workloads_test import ( "go/ast" @@ -12,6 +12,7 @@ import ( "testing" "github.com/stroppy-io/stroppy/pkg/bench" + _ "github.com/stroppy-io/stroppy/workloads/all" ) func TestBuiltInWorkloadParameterSchemas(t *testing.T) { diff --git a/internal/workloads/simple/simple.go b/workloads/simple/simple.go similarity index 91% rename from internal/workloads/simple/simple.go rename to workloads/simple/simple.go index 465eee85..22f9bfc9 100644 --- a/internal/workloads/simple/simple.go +++ b/workloads/simple/simple.go @@ -1,8 +1,5 @@ -// Package simple is the Go-native port of workloads/simple/simple.ts — the -// minimal stroppy demo. Loads a small table via a typed InsertRequest, runs an -// aggregate count plus per-row lookups, and tears down. First Go workload on -// the typed insert path; proves Setup/Iterate/Teardown + Step + Driver.Insert -// with a plain-Go row formula end to end. +// Package simple provides Stroppy's minimal workload example. It owns a small +// typed load, aggregate count, per-row lookups, and teardown lifecycle. package simple import ( diff --git a/internal/workloads/simple/simple_source_test.go b/workloads/simple/simple_source_test.go similarity index 100% rename from internal/workloads/simple/simple_source_test.go rename to workloads/simple/simple_source_test.go diff --git a/workloads/tpcb/assets.go b/workloads/tpcb/assets.go new file mode 100644 index 00000000..3ac2b16a --- /dev/null +++ b/workloads/tpcb/assets.go @@ -0,0 +1,14 @@ +package tpcb + +import ( + "embed" + + "github.com/stroppy-io/stroppy/workloads" +) + +//go:embed *.sql README.md +var files embed.FS + +func init() { + workloads.Register(workloads.PresetTPCB, files) +} diff --git a/workloads/tpcb/assets_test.go b/workloads/tpcb/assets_test.go new file mode 100644 index 00000000..74a039b0 --- /dev/null +++ b/workloads/tpcb/assets_test.go @@ -0,0 +1,30 @@ +package tpcb + +import ( + "testing" + + "github.com/stroppy-io/stroppy/workloads/internal/workloadtest" +) + +func TestEmbeddedAssetContract(t *testing.T) { + workloadtest.Files(t, files, "README.md", "pg.sql", "mysql.sql", "pico.sql", "ydb.sql") + + txQueries := make([]workloadtest.Query, 0, len(requiredTxQueries)) + for _, query := range requiredTxQueries { + txQueries = append(txQueries, workloadtest.Query{Section: query.section, Name: query.query}) + } + + for _, dialect := range []string{"pg.sql", "mysql.sql", "pico.sql", "ydb.sql"} { + t.Run(dialect, func(t *testing.T) { + queries := txQueries + if dialect == "pg.sql" || dialect == "mysql.sql" { + queries = append(append([]workloadtest.Query(nil), txQueries...), workloadtest.Query{ + Section: "workload_procs", + Name: "tpcb_transaction", + }) + } + + workloadtest.SQL(t, files, dialect, requiredSetupSections, queries) + }) + } +} diff --git a/internal/workloads/tpcb/params_test.go b/workloads/tpcb/params_test.go similarity index 100% rename from internal/workloads/tpcb/params_test.go rename to workloads/tpcb/params_test.go diff --git a/internal/workloads/tpcb/procs_test.go b/workloads/tpcb/procs_test.go similarity index 100% rename from internal/workloads/tpcb/procs_test.go rename to workloads/tpcb/procs_test.go diff --git a/internal/workloads/tpcb/sql_validation_test.go b/workloads/tpcb/sql_validation_test.go similarity index 100% rename from internal/workloads/tpcb/sql_validation_test.go rename to workloads/tpcb/sql_validation_test.go diff --git a/internal/workloads/tpcb/tpcb.go b/workloads/tpcb/tpcb.go similarity index 99% rename from internal/workloads/tpcb/tpcb.go rename to workloads/tpcb/tpcb.go index 4a806f7c..43b76e70 100644 --- a/internal/workloads/tpcb/tpcb.go +++ b/workloads/tpcb/tpcb.go @@ -1,5 +1,5 @@ -// Package tpcb is the Go-native port of pgbench's canonical 5-statement TPC-B -// transaction, shipping two registered variants: +// Package tpcb owns Stroppy's TPC-B implementation, tests, and dialect SQL. +// It ships two registered variants: // // - tpcb/tx runs the five DML steps inline under one client-side transaction // per iteration; supports pg/mysql/picodata/ydb. diff --git a/internal/workloads/tpcb/tpcb_source_test.go b/workloads/tpcb/tpcb_source_test.go similarity index 100% rename from internal/workloads/tpcb/tpcb_source_test.go rename to workloads/tpcb/tpcb_source_test.go diff --git a/workloads/tpcc/assets.go b/workloads/tpcc/assets.go new file mode 100644 index 00000000..3f327c7b --- /dev/null +++ b/workloads/tpcc/assets.go @@ -0,0 +1,14 @@ +package tpcc + +import ( + "embed" + + "github.com/stroppy-io/stroppy/workloads" +) + +//go:embed *.sql README.md +var files embed.FS + +func init() { + workloads.Register(workloads.PresetTPCC, files) +} diff --git a/workloads/tpcc/assets_test.go b/workloads/tpcc/assets_test.go new file mode 100644 index 00000000..62caebaa --- /dev/null +++ b/workloads/tpcc/assets_test.go @@ -0,0 +1,87 @@ +package tpcc + +import ( + "testing" + + "github.com/stroppy-io/stroppy/workloads/internal/workloadtest" +) + +func TestEmbeddedAssetContract(t *testing.T) { + dialects := []string{"pg.sql", "mysql.sql", "pico.sql", "ydb.sql", "ydb_no_indexes.sql"} + workloadtest.Files(t, files, append([]string{"README.md"}, dialects...)...) + + sections := []string{ + "drop_schema", + "create_schema", + "workload_tx_new_order", + "workload_tx_payment", + "workload_tx_order_status", + "workload_tx_delivery", + "workload_tx_stock_level", + } + queries := []workloadtest.Query{ + {Section: "workload_tx_new_order", Name: "get_customer"}, + {Section: "workload_tx_new_order", Name: "get_warehouse"}, + {Section: "workload_tx_new_order", Name: "get_district"}, + {Section: "workload_tx_new_order", Name: "update_district"}, + {Section: "workload_tx_new_order", Name: "insert_order"}, + {Section: "workload_tx_new_order", Name: "insert_new_order"}, + {Section: "workload_tx_new_order", Name: "get_items_batch"}, + {Section: "workload_tx_new_order", Name: "get_stocks_batch"}, + {Section: "workload_tx_new_order", Name: "update_stock"}, + {Section: "workload_tx_new_order", Name: "insert_order_line"}, + {Section: "workload_tx_payment", Name: "count_customers_by_name"}, + {Section: "workload_tx_payment", Name: "get_customer_by_name"}, + {Section: "workload_tx_payment", Name: "get_customer_by_id"}, + {Section: "workload_tx_payment", Name: "update_customer_bc"}, + {Section: "workload_tx_payment", Name: "update_customer"}, + {Section: "workload_tx_payment", Name: "insert_history"}, + {Section: "workload_tx_order_status", Name: "count_customers_by_name"}, + {Section: "workload_tx_order_status", Name: "get_customer_by_name"}, + {Section: "workload_tx_order_status", Name: "get_customer_by_id"}, + {Section: "workload_tx_order_status", Name: "get_last_order"}, + {Section: "workload_tx_order_status", Name: "get_order_lines"}, + {Section: "workload_tx_delivery", Name: "get_min_new_order"}, + {Section: "workload_tx_delivery", Name: "delete_new_order"}, + {Section: "workload_tx_delivery", Name: "get_order"}, + {Section: "workload_tx_delivery", Name: "update_order"}, + {Section: "workload_tx_delivery", Name: "update_order_line"}, + {Section: "workload_tx_delivery", Name: "get_order_line_amount"}, + {Section: "workload_tx_delivery", Name: "update_customer"}, + {Section: "workload_tx_stock_level", Name: "get_district"}, + {Section: "workload_tx_stock_level", Name: "get_window_items"}, + {Section: "workload_tx_stock_level", Name: "stock_count_in"}, + } + returningQueries := []workloadtest.Query{ + {Section: "workload_tx_payment", Name: "update_get_warehouse"}, + {Section: "workload_tx_payment", Name: "update_get_district"}, + } + nonReturningQueries := []workloadtest.Query{ + {Section: "workload_tx_payment", Name: "update_warehouse"}, + {Section: "workload_tx_payment", Name: "get_warehouse"}, + {Section: "workload_tx_payment", Name: "update_district"}, + {Section: "workload_tx_payment", Name: "get_district"}, + } + + for _, dialect := range dialects { + t.Run(dialect, func(t *testing.T) { + dialectSections := sections + + dialectQueries := append([]workloadtest.Query(nil), queries...) + if dialect == "pg.sql" || dialect == "ydb.sql" || dialect == "ydb_no_indexes.sql" { + dialectQueries = append(dialectQueries, returningQueries...) + } else { + dialectQueries = append(dialectQueries, nonReturningQueries...) + } + + if dialect == "pg.sql" || dialect == "mysql.sql" { + dialectSections = append(append([]string(nil), sections...), "create_procedures", "workload_procs") + for _, name := range txNames { + dialectQueries = append(dialectQueries, workloadtest.Query{Section: "workload_procs", Name: name}) + } + } + + workloadtest.SQL(t, files, dialect, dialectSections, dialectQueries) + }) + } +} diff --git a/internal/workloads/tpcc/config.go b/workloads/tpcc/config.go similarity index 100% rename from internal/workloads/tpcc/config.go rename to workloads/tpcc/config.go diff --git a/internal/workloads/tpcc/helpers.go b/workloads/tpcc/helpers.go similarity index 100% rename from internal/workloads/tpcc/helpers.go rename to workloads/tpcc/helpers.go diff --git a/internal/workloads/tpcc/load.go b/workloads/tpcc/load.go similarity index 100% rename from internal/workloads/tpcc/load.go rename to workloads/tpcc/load.go diff --git a/internal/workloads/tpcc/neworder_test.go b/workloads/tpcc/neworder_test.go similarity index 100% rename from internal/workloads/tpcc/neworder_test.go rename to workloads/tpcc/neworder_test.go diff --git a/internal/workloads/tpcc/params_test.go b/workloads/tpcc/params_test.go similarity index 100% rename from internal/workloads/tpcc/params_test.go rename to workloads/tpcc/params_test.go diff --git a/internal/workloads/tpcc/procs.go b/workloads/tpcc/procs.go similarity index 99% rename from internal/workloads/tpcc/procs.go rename to workloads/tpcc/procs.go index 4c8238e0..1ec216ba 100644 --- a/internal/workloads/tpcc/procs.go +++ b/workloads/tpcc/procs.go @@ -191,7 +191,7 @@ func (w *workload) procDelivery(ctx context.Context, b *bench.Bench, vs *vuState start := time.Now() defer func() { w.m.deliveryDur.Add(float64(time.Since(start).Milliseconds())) }() - carrierID := vs.ri(vs.dCarrier, 1, districtsPerWarehouse) + carrierID := vs.ri(vs.dCarrier, 1, 10) args := map[string]any{"d_w_id": vs.homeWID, "d_o_carrier_id": carrierID} return bench.Retry0(ctx, w.retryPolicy, func() error { diff --git a/internal/workloads/tpcc/report.go b/workloads/tpcc/report.go similarity index 100% rename from internal/workloads/tpcc/report.go rename to workloads/tpcc/report.go diff --git a/internal/workloads/tpcc/report_test.go b/workloads/tpcc/report_test.go similarity index 100% rename from internal/workloads/tpcc/report_test.go rename to workloads/tpcc/report_test.go diff --git a/internal/workloads/tpcc/tpcc.go b/workloads/tpcc/tpcc.go similarity index 99% rename from internal/workloads/tpcc/tpcc.go rename to workloads/tpcc/tpcc.go index d11f7d87..0bd84092 100644 --- a/internal/workloads/tpcc/tpcc.go +++ b/workloads/tpcc/tpcc.go @@ -1,5 +1,5 @@ -// Package tpcc is the Go-native port of workloads/tpcc/tx.ts: the five TPC-C -// transactions as ordered DML steps inside driver transactions, with the standard +// Package tpcc owns Stroppy's TPC-C implementation, tests, and dialect SQL. It +// runs five transactions as ordered DML steps inside driver transactions, with // 45/43/4/4/4 mix, full population, and §1.3.1 validation. Load/config/prepare are // shared structure ported from tpcc_common.ts. Covers pg + mysql; picodata (no // OFFSET) and ydb (bound IN-list) dialect branches are ported faithfully so all diff --git a/internal/workloads/tpcc/tpcc_complex_source_test.go b/workloads/tpcc/tpcc_complex_source_test.go similarity index 100% rename from internal/workloads/tpcc/tpcc_complex_source_test.go rename to workloads/tpcc/tpcc_complex_source_test.go diff --git a/internal/workloads/tpcc/tpcc_source_test.go b/workloads/tpcc/tpcc_source_test.go similarity index 100% rename from internal/workloads/tpcc/tpcc_source_test.go rename to workloads/tpcc/tpcc_source_test.go diff --git a/internal/workloads/tpcc/validate.go b/workloads/tpcc/validate.go similarity index 100% rename from internal/workloads/tpcc/validate.go rename to workloads/tpcc/validate.go diff --git a/internal/workloads/tpcc/validate_test.go b/workloads/tpcc/validate_test.go similarity index 100% rename from internal/workloads/tpcc/validate_test.go rename to workloads/tpcc/validate_test.go diff --git a/workloads/tpcds/assets.go b/workloads/tpcds/assets.go new file mode 100644 index 00000000..b644abcc --- /dev/null +++ b/workloads/tpcds/assets.go @@ -0,0 +1,14 @@ +package tpcds + +import ( + "embed" + + "github.com/stroppy-io/stroppy/workloads" +) + +//go:embed *.sql *.json README.md +var files embed.FS + +func init() { + workloads.Register(workloads.PresetTPCDS, files) +} diff --git a/workloads/tpcds/assets_test.go b/workloads/tpcds/assets_test.go new file mode 100644 index 00000000..324534bc --- /dev/null +++ b/workloads/tpcds/assets_test.go @@ -0,0 +1,79 @@ +package tpcds + +import ( + "strconv" + "testing" + + "github.com/stroppy-io/stroppy/workloads/internal/workloadtest" +) + +func TestEmbeddedAssetContract(t *testing.T) { + assets := []string{ + "README.md", + "answers_sf1.json", + "pg.sql", + "mysql.sql", + "pico.sql", + "ydb.sql", + "schema.pg.sql", + "schema.mysql.sql", + "schema.pico.sql", + "schema.ydb.sql", + } + workloadtest.Files(t, files, assets...) + + schemas := map[string][]string{ + "schema.pg.sql": {"create_schema", "create_indexes"}, + "schema.mysql.sql": {"create_schema", "create_indexes"}, + "schema.pico.sql": {"drop_schema", "create_schema", "create_indexes"}, + "schema.ydb.sql": {"drop_schema", "create_schema", "create_schema_column", "create_indexes"}, + } + for name, sections := range schemas { + t.Run(name, func(t *testing.T) { + workloadtest.SQL(t, files, name, sections, nil) + }) + } + + queries := tpcdsContractQueries(nil) + picoQueries := tpcdsContractQueries(map[int]bool{ + 36: true, + 44: true, + 47: true, + 49: true, + 57: true, + 67: true, + 70: true, + 86: true, + }) + + for _, name := range []string{"pg.sql", "mysql.sql", "pico.sql", "ydb.sql"} { + t.Run(name, func(t *testing.T) { + required := queries + if name == "pico.sql" { + required = picoQueries + } + + workloadtest.SQL(t, files, name, nil, required) + }) + } +} + +func tpcdsContractQueries(skip map[int]bool) []workloadtest.Query { + queries := make([]workloadtest.Query, 0, 103) + + for number := 1; number <= 99; number++ { + if skip[number] { + continue + } + + name := "query_" + strconv.Itoa(number) + switch number { + case 14, 23, 24, 39: + queries = append(queries, workloadtest.Query{Name: name + "_a"}, workloadtest.Query{Name: name + "_b"}) + default: + queries = append(queries, workloadtest.Query{Name: name}) + } + } + + return queries +} diff --git a/internal/workloads/tpcds/config.go b/workloads/tpcds/config.go similarity index 90% rename from internal/workloads/tpcds/config.go rename to workloads/tpcds/config.go index 6e38784b..dce4707e 100644 --- a/internal/workloads/tpcds/config.go +++ b/workloads/tpcds/config.go @@ -1,5 +1,5 @@ -// Package tpcds is the Go-native port of workloads/tpcds/tpcds.ts: the relational load -// of the 24 TPC-DS tables via the ported dsdgen generator (bench.InsertTpcds) plus the +// Package tpcds owns Stroppy's TPC-DS implementation, tests, dialect SQL, and +// answer data. It loads 24 tables through the canonical generator and runs the // 99 business queries, run either from the baked canonical qualification set or from an // in-process generated stream (throughput test). SF=1 answer validation (pg/mysql) is // ported from tpcds_validate.ts as a multiset comparison against answers_sf1.json. diff --git a/internal/workloads/tpcds/params_test.go b/workloads/tpcds/params_test.go similarity index 100% rename from internal/workloads/tpcds/params_test.go rename to workloads/tpcds/params_test.go diff --git a/internal/workloads/tpcds/tpcds.go b/workloads/tpcds/tpcds.go similarity index 100% rename from internal/workloads/tpcds/tpcds.go rename to workloads/tpcds/tpcds.go diff --git a/internal/workloads/tpcds/validate.go b/workloads/tpcds/validate.go similarity index 100% rename from internal/workloads/tpcds/validate.go rename to workloads/tpcds/validate.go diff --git a/workloads/tpch/assets.go b/workloads/tpch/assets.go new file mode 100644 index 00000000..5f78de20 --- /dev/null +++ b/workloads/tpch/assets.go @@ -0,0 +1,14 @@ +package tpch + +import ( + "embed" + + "github.com/stroppy-io/stroppy/workloads" +) + +//go:embed *.sql *.json README.md +var files embed.FS + +func init() { + workloads.Register(workloads.PresetTPCH, files) +} diff --git a/workloads/tpch/assets_test.go b/workloads/tpch/assets_test.go new file mode 100644 index 00000000..4dc4780a --- /dev/null +++ b/workloads/tpch/assets_test.go @@ -0,0 +1,34 @@ +package tpch + +import ( + "testing" + + "github.com/stroppy-io/stroppy/workloads/internal/workloadtest" +) + +func TestEmbeddedAssetContract(t *testing.T) { + dialects := []string{"pg.sql", "mysql.sql", "pico.sql", "ydb.sql"} + workloadtest.Files( + t, + files, + append([]string{"README.md", "answers_sf1.json", "distributions.json"}, dialects...)..., + ) + + sections := map[string][]string{ + "pg.sql": {"drop_schema", "create_schema", "set_unlogged", "create_indexes", "set_logged", "analyze"}, + "mysql.sql": {"drop_schema", "create_schema", "create_indexes", "analyze"}, + "pico.sql": {"drop_schema", "create_schema", "create_indexes"}, + "ydb.sql": {"drop_schema", "create_schema", "create_schema_column", "create_indexes"}, + } + queries := make([]workloadtest.Query, 0, len(queryNames)) + + for _, section := range queryNames { + queries = append(queries, workloadtest.Query{Section: section, Name: "body"}) + } + + for _, dialect := range dialects { + t.Run(dialect, func(t *testing.T) { + workloadtest.SQL(t, files, dialect, sections[dialect], queries) + }) + } +} diff --git a/internal/workloads/tpch/config.go b/workloads/tpch/config.go similarity index 100% rename from internal/workloads/tpch/config.go rename to workloads/tpch/config.go diff --git a/internal/workloads/tpch/params_test.go b/workloads/tpch/params_test.go similarity index 100% rename from internal/workloads/tpch/params_test.go rename to workloads/tpch/params_test.go diff --git a/internal/workloads/tpch/tpch.go b/workloads/tpch/tpch.go similarity index 95% rename from internal/workloads/tpch/tpch.go rename to workloads/tpch/tpch.go index 0142e8ce..49eddad6 100644 --- a/internal/workloads/tpch/tpch.go +++ b/workloads/tpch/tpch.go @@ -1,6 +1,6 @@ -// Package tpch is the Go-native port of workloads/tpch/tx.ts: the relational load of -// the 8 TPC-H tables via the ported dbgen generator (bench.InsertTpch) plus the q1–q22 -// business queries run once with §2.4 pinned defaults, and SF=1 answer validation +// Package tpch owns Stroppy's TPC-H implementation, tests, dialect SQL, and +// answer data. It loads eight tables through the canonical generator and runs +// q1–q22 once with §2.4 pinned defaults, and SF=1 answer validation // (postgres only). Supports pg/mysql/pico/ydb dialect files; date shifts for pico/ydb // are precomputed client-side. package tpch diff --git a/internal/workloads/tpch/validate.go b/workloads/tpch/validate.go similarity index 100% rename from internal/workloads/tpch/validate.go rename to workloads/tpch/validate.go