Compare commits

..
Author SHA1 Message Date
HeinandClaude Sonnet 5 ce3b615b0a feat(dbml): @postgres/@sqlite dialect directives (#19)
Add parseable `@<namespace>[(<target>)]: <args>` directives embedded in DBML.
They are stored losslessly on each object's Metadata, round-trip unchanged
through the DBML writer, and are translated to SQL only by the writer for the
matching dialect.

- models: Directive type + catalog; Metadata map added to Column and Index
- dbml reader: parse and attach directives at database/table/column/index
  level; line-numbered errors; repeatable by default with singleton duplicate
  detection. Fixes a preexisting bug where an `indexes {}` closing brace ended
  the table early, dropping trailing Note: and directive lines.
- dbml writer: re-emit directives at their location; idempotent output
- pgsql writer: PARTITION BY / INHERITS / WITH / TABLESPACE (table),
  STORAGE / COMPRESSION / identity (column), WITH / TABLESPACE (index)
- sqlite writer: WITHOUT ROWID / STRICT (table), COLLATE (column)
- --strict-directives flag on ReaderOptions and WriterOptions
- docs/DBML_DIRECTIVES.md + reader/writer READMEs

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Ss2MY5J11cRGwEz86ZXk7d
2026-09-08 16:17:37 +02:00
34 changed files with 1736 additions and 1525 deletions
+9 -317
View File
@@ -1,7 +1,6 @@
package main
import (
"bytes"
"fmt"
"io"
"os"
@@ -12,8 +11,6 @@ import (
"github.com/spf13/cobra"
"git.warky.dev/wdevs/relspecgo/pkg/diff"
"git.warky.dev/wdevs/relspecgo/pkg/inspector"
"git.warky.dev/wdevs/relspecgo/pkg/jobs"
"git.warky.dev/wdevs/relspecgo/pkg/merge"
"git.warky.dev/wdevs/relspecgo/pkg/models"
@@ -21,7 +18,6 @@ import (
"git.warky.dev/wdevs/relspecgo/pkg/readers/sqldir"
"git.warky.dev/wdevs/relspecgo/pkg/writers"
wpgsql "git.warky.dev/wdevs/relspecgo/pkg/writers/pgsql"
"git.warky.dev/wdevs/relspecgo/pkg/writers/sqlexec"
wtemplate "git.warky.dev/wdevs/relspecgo/pkg/writers/template"
)
@@ -126,9 +122,6 @@ func loadJobSet() (*jobs.Set, error) {
if err := set.Validate(); err != nil {
return nil, err
}
for _, w := range set.Warnings {
fmt.Fprintf(os.Stderr, "warning: %s\n", w)
}
return set, nil
}
@@ -187,14 +180,12 @@ func executeJobPlan(set *jobs.Set, name string, dryRun, noDeps bool, out io.Writ
// Pre-flight: resolve and check paths, output policy and env vars for the
// whole plan before anything runs. A failure here means no job executes.
resolved := make([]*resolvedJob, len(plan))
byName := make(map[string]*resolvedJob, len(plan))
for i, j := range plan {
rj, perr := preflightJob(j, byName)
rj, perr := preflightJob(j)
if perr != nil {
return fmt.Errorf("job %q: %w", j.Name, perr)
}
resolved[i] = rj
byName[j.Name] = rj
}
if dryRun {
@@ -226,12 +217,7 @@ type resolvedJob struct {
outputConn string // resolved connection string (secret)
outputConnEnv string
logPath string
logPolicy jobs.LogPolicy
templatePath string
reportPath string // "" for a diff summary written to the log
reportFormat string
rulesPath string // "" means inspector defaults
selection *splitSelection
secrets []string // resolved secret values to redact from logs
}
@@ -240,12 +226,11 @@ type resolvedInput struct {
path string // "" when the input is a database
conn string // resolved connection string (secret)
connEnv string
fromJob string // producer job name when this input came from from_job
}
func preflightJob(j *jobs.Job, resolvedByName map[string]*resolvedJob) (*resolvedJob, error) {
func preflightJob(j *jobs.Job) (*resolvedJob, error) {
root := j.Dir()
rj := &resolvedJob{job: j, root: root, logPolicy: j.ResolvedLogPolicy()}
rj := &resolvedJob{job: j, root: root}
if j.Logfile != "" {
p, err := jobs.SafeJoin(root, j.Logfile)
@@ -268,20 +253,6 @@ func preflightJob(j *jobs.Job, resolvedByName map[string]*resolvedJob) (*resolve
for i, in := range j.Inputs {
ri := resolvedInput{format: strings.ToLower(in.Format)}
if in.FromJob != "" {
producer, ok := resolvedByName[in.FromJob]
if !ok {
return nil, fmt.Errorf("input[%d]: from_job %q is not in this plan (do not use --no-deps with from_job inputs)", i, in.FromJob)
}
if producer.outputPath == "" {
return nil, fmt.Errorf("input[%d]: from_job %q does not write a file output", i, in.FromJob)
}
ri.path = producer.outputPath
ri.format = strings.ToLower(producer.job.Output.Format)
ri.fromJob = in.FromJob
rj.inputs = append(rj.inputs, ri)
continue
}
if in.ConnEnv != "" {
v, ok := os.LookupEnv(in.ConnEnv)
if !ok || v == "" {
@@ -342,43 +313,6 @@ func preflightJob(j *jobs.Job, resolvedByName map[string]*resolvedJob) (*resolve
rj.outputPath = p
}
}
if j.Rules != "" {
p, err := jobs.SafeJoin(root, j.Rules)
if err != nil {
return nil, fmt.Errorf("rules: %w", err)
}
info, err := os.Stat(p)
if err != nil || info.IsDir() {
return nil, fmt.Errorf("rules %q: not found or is a directory", j.Rules)
}
rj.rulesPath = p
}
if j.Report != nil {
rj.reportFormat = strings.ToLower(j.Report.Format)
if j.Report.Path != "" {
p, err := jobs.SafeJoin(root, j.Report.Path)
if err != nil {
return nil, fmt.Errorf("report: %w", err)
}
if _, err := os.Stat(p); err == nil && !j.Report.Overwrite {
return nil, fmt.Errorf("report %s already exists (set report.overwrite: true to replace it)", j.Report.Path)
}
rj.reportPath = p
}
}
if j.Select != nil {
rj.selection = &splitSelection{
Schemas: j.Select.Schemas,
Tables: j.Select.Tables,
ExcludeSchemas: j.Select.ExcludeSchemas,
ExcludeTables: j.Select.ExcludeTables,
DatabaseName: j.Select.DatabaseName,
}
}
return rj, nil
}
@@ -391,12 +325,9 @@ func printResolvedJob(out io.Writer, n, total int, rj *resolvedJob) {
}
fmt.Fprintf(out, " job file: %s\n", j.SourceFile)
for _, ri := range rj.inputs {
switch {
case ri.fromJob != "":
fmt.Fprintf(out, " input: %s (%s) from job %q\n", ri.path, ri.format, ri.fromJob)
case ri.path != "":
if ri.path != "" {
fmt.Fprintf(out, " input: %s (%s)\n", ri.path, ri.format)
default:
} else {
fmt.Fprintf(out, " input: env:%s (%s)\n", ri.connEnv, ri.format)
}
}
@@ -408,31 +339,15 @@ func printResolvedJob(out io.Writer, n, total int, rj *resolvedJob) {
} else if rj.outputConnEnv != "" {
fmt.Fprintf(out, " output: env:%s (%s)\n", rj.outputConnEnv, j.Output.Format)
}
if j.Report != nil {
format := valueOr(rj.reportFormat, "default")
if rj.reportPath != "" {
fmt.Fprintf(out, " report: %s (%s)\n", rj.reportPath, format)
} else {
fmt.Fprintf(out, " report: (log) (%s)\n", format)
}
}
if rj.rulesPath != "" {
fmt.Fprintf(out, " rules: %s\n", rj.rulesPath)
} else if j.Command == jobs.CommandInspect {
fmt.Fprintf(out, " rules: (built-in defaults)\n")
}
if rj.selection != nil {
fmt.Fprintf(out, " select: %s\n", rj.selection.summary())
}
if rj.logPath != "" {
fmt.Fprintf(out, " logfile: %s (rotate >= %d bytes, keep %d)\n", rj.logPath, rj.logPolicy.MaxSizeBytes, rj.logPolicy.Keep)
fmt.Fprintf(out, " logfile: %s\n", rj.logPath)
}
fmt.Fprintln(out)
}
// executeResolvedJob runs a single already-validated job.
func executeResolvedJob(rj *resolvedJob) (err error) {
lg, closeLog, lerr := newJobLogger(rj.logPath, rj.logPolicy, rj.secrets)
lg, closeLog, lerr := newJobLogger(rj.logPath, rj.secrets)
if lerr != nil {
return lerr
}
@@ -449,14 +364,6 @@ func executeResolvedJob(rj *resolvedJob) (err error) {
err = runScriptsListJob(rj, lg)
case jobs.CommandTempl:
err = runTemplJob(rj, lg)
case jobs.CommandSplit:
err = runSplitJob(rj, lg)
case jobs.CommandInspect:
err = runInspectJob(rj, lg)
case jobs.CommandDiff:
err = runDiffJob(rj, lg)
case jobs.CommandScriptsExec:
err = runScriptsExecJob(rj, lg)
default:
err = fmt.Errorf("unsupported command %q", rj.job.Command)
}
@@ -598,168 +505,6 @@ func runScriptsListJob(rj *resolvedJob, lg *jobLogger) error {
return nil
}
func runSplitJob(rj *resolvedJob, lg *jobLogger) error {
db, err := readJobInputs(rj, lg)
if err != nil {
return err
}
sel := splitSelection{}
if rj.selection != nil {
sel = *rj.selection
}
filtered, err := filterDatabaseSelection(db, sel)
if err != nil {
return fmt.Errorf("split selection: %w", err)
}
if sel.DatabaseName != "" {
filtered.Name = sel.DatabaseName
}
tables := 0
for _, s := range filtered.Schemas {
tables += len(s.Tables)
}
lg.logf("split: selected %d schema(s), %d table(s)", len(filtered.Schemas), tables)
return writeJobOutput(rj, filtered, lg)
}
func runInspectJob(rj *resolvedJob, lg *jobLogger) error {
db, err := readJobInputs(rj, lg)
if err != nil {
return err
}
config, err := inspector.LoadConfig(rj.rulesPath) // "" -> built-in defaults
if err != nil {
return fmt.Errorf("load rules: %w", err)
}
report, err := inspector.NewInspector(db, config).Inspect()
if err != nil {
return fmt.Errorf("inspection failed: %w", err)
}
var formatted string
switch valueOr(rj.reportFormat, "markdown") {
case "json":
formatted, err = inspector.NewJSONFormatter().Format(report)
default:
formatted, err = inspector.NewMarkdownFormatter(io.Discard).Format(report)
}
if err != nil {
return fmt.Errorf("format report: %w", err)
}
if werr := atomicWrite(rj.reportPath, func(tmp string) error {
return os.WriteFile(tmp, []byte(formatted), 0o644)
}); werr != nil {
return werr
}
lg.logf("inspect: %d error(s), %d warning(s) -> %s",
report.Summary.ErrorCount, report.Summary.WarningCount, rj.reportPath)
if report.HasErrors() {
return fmt.Errorf("inspection found %d error(s)", report.Summary.ErrorCount)
}
return nil
}
func runDiffJob(rj *resolvedJob, lg *jobLogger) error {
if len(rj.inputs) != 2 {
return fmt.Errorf("diff requires exactly 2 inputs, got %d", len(rj.inputs))
}
source, err := readOneJobInput(rj.inputs[0])
if err != nil {
return fmt.Errorf("input[0]: %w", err)
}
lg.logf("diff source: %s", inputLabel(rj.inputs[0]))
target, err := readOneJobInput(rj.inputs[1])
if err != nil {
return fmt.Errorf("input[1]: %w", err)
}
lg.logf("diff target: %s", inputLabel(rj.inputs[1]))
result := diff.CompareDatabases(source, target)
s := diff.ComputeSummary(result)
lg.logf("diff: schemas %d/%d/%d, tables %d/%d/%d, columns %d/%d/%d (missing/extra/modified)",
s.Schemas.Missing, s.Schemas.Extra, s.Schemas.Modified,
s.Tables.Missing, s.Tables.Extra, s.Tables.Modified,
s.Columns.Missing, s.Columns.Extra, s.Columns.Modified)
format := diff.FormatSummary
switch rj.reportFormat {
case "json":
format = diff.FormatJSON
case "html":
format = diff.FormatHTML
}
if rj.reportPath == "" {
var buf bytes.Buffer
if err := diff.FormatDiff(result, format, &buf); err != nil {
return fmt.Errorf("format diff: %w", err)
}
for _, line := range strings.Split(strings.TrimRight(buf.String(), "\n"), "\n") {
lg.logf("%s", line)
}
return nil
}
if werr := atomicWrite(rj.reportPath, func(tmp string) error {
f, err := os.Create(tmp)
if err != nil {
return err
}
defer f.Close()
return diff.FormatDiff(result, format, f)
}); werr != nil {
return werr
}
lg.logf("diff report written: %s", rj.reportPath)
return nil
}
func runScriptsExecJob(rj *resolvedJob, lg *jobLogger) error {
schemaName := valueOr(rj.job.Options.Schema, "public")
combined := &models.Schema{Name: schemaName}
for _, dir := range rj.scriptDirs {
reader := sqldir.NewReader(&readers.ReaderOptions{
FilePath: dir,
Metadata: map[string]any{
"schema_name": schemaName,
"database_name": "database",
},
})
db, err := reader.ReadDatabase()
if err != nil {
return fmt.Errorf("%s: %w", dir, err)
}
if len(db.Schemas) == 0 {
continue
}
combined.Scripts = append(combined.Scripts, db.Schemas[0].Scripts...)
}
if len(combined.Scripts) == 0 {
lg.logf("no scripts found; nothing to execute")
return nil
}
lg.logf("executing %d script(s) against database env:%s", len(combined.Scripts), rj.outputConnEnv)
writer := sqlexec.NewWriter(&writers.WriterOptions{
Metadata: map[string]any{
"connection_string": rj.outputConn,
"ignore_errors": rj.job.Options.ContinueOnError,
},
})
if err := writer.WriteSchema(combined); err != nil {
return fmt.Errorf("script execution failed: %w", err)
}
opts := writer.Options()
total, _ := opts.Metadata["execution_total"].(int)
success, _ := opts.Metadata["execution_success"].(int)
failed, _ := opts.Metadata["execution_failed"].(int)
lg.logf("executed %d script(s): %d succeeded, %d failed", total, success, failed)
if failed > 0 && !rj.job.Options.ContinueOnError {
return fmt.Errorf("%d script(s) failed", failed)
}
return nil
}
// readJobInputs reads every input and additively merges them into one model.
func readJobInputs(rj *resolvedJob, lg *jobLogger) (*models.Database, error) {
var base *models.Database
@@ -814,37 +559,7 @@ func writeJobOutput(rj *resolvedJob, db *models.Database, lg *jobLogger) error {
return fmt.Errorf("failed to create output directory: %w", err)
}
lg.logf("writing output: %s (%s)", rj.outputPath, format)
write := func(target string) error {
return writeDatabase(db, format, target, o.Package, o.Schema, o.FlattenSchema, "", "", o.ContinueOnError, "")
}
// Single-file formats are written to a temp file and renamed into place so
// a failure never leaves a partial or truncated output. Directory-emitting
// formats (gorm/bun/drizzle/typeorm/prisma) write in place.
if jobs.SingleFileOutputFormat(format) {
return atomicWrite(rj.outputPath, write)
}
return write(rj.outputPath)
}
// atomicWrite calls produce with a temp path in the same directory as
// finalPath, then renames it over finalPath. The temp file is removed on any
// error so the destination is only ever replaced by a complete file.
func atomicWrite(finalPath string, produce func(tmpPath string) error) error {
dir := filepath.Dir(finalPath)
if err := os.MkdirAll(dir, 0o755); err != nil {
return fmt.Errorf("failed to create output directory: %w", err)
}
tmp := filepath.Join(dir, fmt.Sprintf(".%s.relspec-tmp-%d", filepath.Base(finalPath), os.Getpid()))
if err := produce(tmp); err != nil {
_ = os.Remove(tmp)
return err
}
if err := os.Rename(tmp, finalPath); err != nil {
_ = os.Remove(tmp)
return fmt.Errorf("failed to finalize %s: %w", finalPath, err)
}
return nil
return writeDatabase(db, format, rj.outputPath, o.Package, o.Schema, o.FlattenSchema, "", "", o.ContinueOnError, "")
}
// --- logging + redaction ---------------------------------------------------
@@ -857,7 +572,7 @@ type jobLogger struct {
// newJobLogger returns a logger that mirrors to stderr and, when path is set,
// to a job logfile. Connection strings and known secret values are redacted
// from everything it writes.
func newJobLogger(path string, policy jobs.LogPolicy, secrets []string) (*jobLogger, func(err error), error) {
func newJobLogger(path string, secrets []string) (*jobLogger, func(err error), error) {
lg := &jobLogger{secrets: secrets}
if path == "" {
return lg, func(error) {}, nil
@@ -865,7 +580,6 @@ func newJobLogger(path string, policy jobs.LogPolicy, secrets []string) (*jobLog
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return nil, nil, fmt.Errorf("failed to create log directory: %w", err)
}
rotateLogIfNeeded(path, policy)
f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644)
if err != nil {
return nil, nil, fmt.Errorf("failed to open logfile %q: %w", path, err)
@@ -879,28 +593,6 @@ func newJobLogger(path string, policy jobs.LogPolicy, secrets []string) (*jobLog
}, nil
}
// rotateLogIfNeeded renames path -> path.1 -> path.2 ... up to policy.Keep
// when path has grown to policy.MaxSizeBytes or more. The oldest file beyond
// Keep is deleted. A zero/negative MaxSizeBytes disables rotation.
func rotateLogIfNeeded(path string, policy jobs.LogPolicy) {
if policy.MaxSizeBytes <= 0 {
return
}
info, err := os.Stat(path)
if err != nil || info.Size() < policy.MaxSizeBytes {
return
}
if policy.Keep < 1 {
_ = os.Remove(path)
return
}
_ = os.Remove(fmt.Sprintf("%s.%d", path, policy.Keep))
for i := policy.Keep - 1; i >= 1; i-- {
_ = os.Rename(fmt.Sprintf("%s.%d", path, i), fmt.Sprintf("%s.%d", path, i+1))
}
_ = os.Rename(path, path+".1")
}
func (l *jobLogger) logf(format string, args ...interface{}) {
line := l.redact(fmt.Sprintf(format, args...))
fmt.Fprintf(os.Stderr, " %s\n", line)
-280
View File
@@ -376,286 +376,6 @@ jobs:
}
}
func TestJobRun_SplitJob(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
extract:
command: split
inputs:
- path: schema/core.dbml
format: dbml
- path: schema/tenant.dbml
format: dbml
select:
tables: [users]
output:
format: json
path: build/subset.json
overwrite: true
`)
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
if err := executeJobPlan(set, "extract", false, false, &bytes.Buffer{}); err != nil {
t.Fatalf("execute split job: %v", err)
}
out, err := os.ReadFile(filepath.Join(dir, "build", "subset.json"))
if err != nil {
t.Fatalf("read split output: %v", err)
}
s := string(out)
if !strings.Contains(s, "users") {
t.Fatalf("split output missing selected table:\n%s", s)
}
if strings.Contains(s, "posts") {
t.Fatalf("split output should have excluded posts:\n%s", s)
}
}
func TestJobRun_InspectJob(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
lint:
command: inspect
inputs:
- path: schema/core.dbml
format: dbml
report:
format: json
path: build/report.json
overwrite: true
logfile: .relspec/lint.log
`)
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
// Default rules only warn, so the job succeeds.
if err := executeJobPlan(set, "lint", false, false, &bytes.Buffer{}); err != nil {
t.Fatalf("execute inspect job: %v", err)
}
if _, err := os.ReadFile(filepath.Join(dir, "build", "report.json")); err != nil {
t.Fatalf("expected report file: %v", err)
}
logData, _ := os.ReadFile(filepath.Join(dir, ".relspec", "lint.log"))
if !strings.Contains(string(logData), "inspect:") {
t.Fatalf("logfile missing inspect summary:\n%s", logData)
}
}
func TestJobRun_InspectJobFailsOnRuleError(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
lint:
command: inspect
inputs:
- path: schema/core.dbml
format: dbml
rules: rules.yaml
report:
format: json
path: build/report.json
overwrite: true
logfile: .relspec/lint.log
`)
// A rule set to "error" level for a violation the fixture triggers.
writeFile(t, filepath.Join(dir, "rules.yaml"), `version: "1.0"
rules:
primary_key_naming:
enabled: enforce
function: primary_key_naming
pattern: "^id_"
message: "Primary key columns should start with id_"
`)
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
err := executeJobPlan(set, "lint", false, false, &bytes.Buffer{})
if err == nil || !strings.Contains(err.Error(), "error(s)") {
t.Fatalf("expected inspect job to fail on rule error, got %v", err)
}
logData, _ := os.ReadFile(filepath.Join(dir, ".relspec", "lint.log"))
if !strings.Contains(string(logData), "FAILED") {
t.Fatalf("failed inspect job should log FAILED:\n%s", logData)
}
}
func TestJobRun_DiffJob(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
compare:
command: diff
inputs:
- path: schema/core.dbml
format: dbml
- path: schema/tenant.dbml
format: dbml
report:
format: json
path: build/diff.json
overwrite: true
`)
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
if err := executeJobPlan(set, "compare", false, false, &bytes.Buffer{}); err != nil {
t.Fatalf("execute diff job: %v", err)
}
out, err := os.ReadFile(filepath.Join(dir, "build", "diff.json"))
if err != nil {
t.Fatalf("read diff report: %v", err)
}
if len(out) == 0 {
t.Fatal("diff report is empty")
}
}
func TestJobRun_FromJobWiring(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
a:
command: convert
inputs:
- path: schema/core.dbml
format: dbml
output:
format: json
path: build/a.json
overwrite: true
b:
command: convert
inputs:
- from_job: a
output:
format: yaml
path: build/b.yaml
overwrite: true
`)
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
if err := executeJobPlan(set, "b", false, false, &bytes.Buffer{}); err != nil {
t.Fatalf("execute from_job chain: %v", err)
}
if _, err := os.Stat(filepath.Join(dir, "build", "a.json")); err != nil {
t.Fatalf("producer output missing: %v", err)
}
out, err := os.ReadFile(filepath.Join(dir, "build", "b.yaml"))
if err != nil {
t.Fatalf("consumer output missing: %v", err)
}
if !strings.Contains(string(out), "users") {
t.Fatalf("consumer did not consume producer output:\n%s", out)
}
}
func TestJobRun_LogRotation(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
build:
command: scripts-list
script_dirs: [migrations]
log_max_size: "150B"
log_keep: 2
logfile: .relspec/build.log
`)
writeFile(t, filepath.Join(dir, "migrations", "1_001_a.sql"), "CREATE TABLE a();\n")
logPath := filepath.Join(dir, ".relspec", "build.log")
writeFile(t, logPath, strings.Repeat("x", 300)+"\n")
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
if err := executeJobPlan(set, "build", false, false, &bytes.Buffer{}); err != nil {
t.Fatalf("execute job: %v", err)
}
rotated, err := os.ReadFile(logPath + ".1")
if err != nil {
t.Fatalf("expected rotated logfile build.log.1: %v", err)
}
if !strings.Contains(string(rotated), strings.Repeat("x", 300)) {
t.Fatalf("rotated logfile should hold the old content")
}
fresh, err := os.ReadFile(logPath)
if err != nil {
t.Fatalf("expected fresh logfile: %v", err)
}
if strings.Contains(string(fresh), strings.Repeat("x", 300)) {
t.Fatalf("fresh logfile should not contain the rotated-out content:\n%s", fresh)
}
if !strings.Contains(string(fresh), "OK") {
t.Fatalf("fresh logfile should hold the new run:\n%s", fresh)
}
}
func TestJobRun_AtomicOutputLeavesOriginalOnFailure(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
x:
command: convert
inputs:
- path: schema/core.dbml
format: dbml
output:
format: json
path: build/out.json
overwrite: true
`)
// Seed the destination, then make its parent directory read-only so the
// rename step fails. The seeded file must survive intact.
seeded := filepath.Join(dir, "build", "out.json")
writeFile(t, seeded, `{"seeded":true}`)
if err := os.Chmod(filepath.Join(dir, "build"), 0o500); err != nil {
t.Skipf("cannot chmod: %v", err)
}
t.Cleanup(func() { _ = os.Chmod(filepath.Join(dir, "build"), 0o755) })
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
if err := executeJobPlan(set, "x", false, false, &bytes.Buffer{}); err == nil {
t.Skip("write unexpectedly succeeded (running as root?)")
}
if err := os.Chmod(filepath.Join(dir, "build"), 0o755); err != nil {
t.Fatal(err)
}
data, err := os.ReadFile(seeded)
if err != nil {
t.Fatalf("seeded file gone: %v", err)
}
if !strings.Contains(string(data), "seeded") {
t.Fatalf("seeded file was corrupted: %s", data)
}
}
func TestJobRun_ScriptsExecMissingConnEnv(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
migrate:
command: scripts-exec
script_dirs: [migrations]
output:
conn_env: RELSPEC_TEST_EXEC_MISSING
logfile: .relspec/migrate.log
`)
writeFile(t, filepath.Join(dir, "migrations", "1_001_a.sql"), "CREATE TABLE a();\n")
os.Unsetenv("RELSPEC_TEST_EXEC_MISSING")
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
err := executeJobPlan(set, "migrate", false, false, &bytes.Buffer{})
if err == nil || !strings.Contains(err.Error(), "conn_env") {
t.Fatalf("expected missing conn_env error, got %v", err)
}
}
func TestJobRun_ScriptsExecDryRun(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
migrate:
command: scripts-exec
script_dirs: [migrations]
output:
conn_env: RELSPEC_TEST_EXEC_CONN
`)
writeFile(t, filepath.Join(dir, "migrations", "1_001_a.sql"), "CREATE TABLE a();\n")
t.Setenv("RELSPEC_TEST_EXEC_CONN", "postgres://u:secretpw@h/db")
set := mustLoadSet(t, filepath.Join(dir, "relspec.yml"))
var buf bytes.Buffer
if err := executeJobPlan(set, "migrate", true, false, &buf); err != nil {
t.Fatalf("dry run: %v", err)
}
if strings.Contains(buf.String(), "secretpw") {
t.Fatalf("plan leaked secret:\n%s", buf.String())
}
if !strings.Contains(buf.String(), "env:RELSPEC_TEST_EXEC_CONN") {
t.Fatalf("plan should name the env var:\n%s", buf.String())
}
}
func TestJobRun_TemplDatabaseMode(t *testing.T) {
dir := jobFixture(t, `version: 1
jobs:
+2
View File
@@ -10,6 +10,7 @@ func newReaderOptions(filePath, connString string) *readers.ReaderOptions {
FilePath: filePath,
ConnectionString: connString,
Prisma7: prisma7,
StrictDirectives: strictDirectives,
}
}
@@ -22,5 +23,6 @@ func newWriterOptions(outputPath, packageName string, flattenSchema bool, nullab
NullableArrays: nullableArrays,
Prisma7: prisma7,
ContinueOnError: continueOnError,
StrictDirectives: strictDirectives,
}
}
+2
View File
@@ -14,6 +14,7 @@ var (
buildDate = "unknown"
prisma7 bool
noVersion bool
strictDirectives bool
)
func init() {
@@ -72,6 +73,7 @@ func init() {
rootCmd.AddCommand(reportCmd)
rootCmd.PersistentFlags().BoolVar(&prisma7, "prisma7", false, "Use Prisma 7 generator conventions when reading/writing Prisma schemas")
rootCmd.PersistentFlags().BoolVar(&noVersion, "no-version", false, "Suppress the RelSpec version header")
rootCmd.PersistentFlags().BoolVar(&strictDirectives, "strict-directives", false, "Fail on unknown or untranslatable DBML dialect directives (@postgres:, @sqlite:, …)")
}
// printVersionHeader prints the "RelSpec <version> (built: <date>)" banner
+6 -50
View File
@@ -205,52 +205,8 @@ func runSplit(cmd *cobra.Command, args []string) error {
return nil
}
// splitSelection is the schema/table selection for a split, independent of the
// CLI flag globals so the job runner can build one directly.
type splitSelection struct {
Schemas []string
Tables []string
ExcludeSchemas []string
ExcludeTables []string
DatabaseName string
}
// summary renders a one-line human description of the selection.
func (s splitSelection) summary() string {
var parts []string
if len(s.Schemas) > 0 {
parts = append(parts, "schemas="+strings.Join(s.Schemas, ","))
}
if len(s.Tables) > 0 {
parts = append(parts, "tables="+strings.Join(s.Tables, ","))
}
if len(s.ExcludeSchemas) > 0 {
parts = append(parts, "exclude_schemas="+strings.Join(s.ExcludeSchemas, ","))
}
if len(s.ExcludeTables) > 0 {
parts = append(parts, "exclude_tables="+strings.Join(s.ExcludeTables, ","))
}
if s.DatabaseName != "" {
parts = append(parts, "database_name="+s.DatabaseName)
}
if len(parts) == 0 {
return "(all schemas/tables)"
}
return strings.Join(parts, " ")
}
// filterDatabase filters the database based on the CLI split flags.
// filterDatabase filters the database based on provided criteria
func filterDatabase(db *models.Database) (*models.Database, error) {
return filterDatabaseSelection(db, splitSelection{
Schemas: parseCommaSeparated(splitSchemas),
Tables: parseCommaSeparated(splitTables),
ExcludeSchemas: parseCommaSeparated(splitExcludeSchema),
ExcludeTables: parseCommaSeparated(splitExcludeTables),
})
}
// filterDatabaseSelection filters db down to the schemas/tables named by sel.
func filterDatabaseSelection(db *models.Database, sel splitSelection) (*models.Database, error) {
filteredDB := &models.Database{
Name: db.Name,
Description: db.Description,
@@ -264,11 +220,11 @@ func filterDatabaseSelection(db *models.Database, sel splitSelection) (*models.D
Domains: db.Domains, // Keep domains for now
}
// Selection criteria
includeSchemas := sel.Schemas
includeTables := sel.Tables
excludeSchemas := sel.ExcludeSchemas
excludeTables := sel.ExcludeTables
// Parse filter flags
includeSchemas := parseCommaSeparated(splitSchemas)
includeTables := parseCommaSeparated(splitTables)
excludeSchemas := parseCommaSeparated(splitExcludeSchema)
excludeTables := parseCommaSeparated(splitExcludeTables)
// Convert table names to lowercase for case-insensitive matching
includeTablesLower := make(map[string]bool)
+115
View File
@@ -0,0 +1,115 @@
# DBML Dialect Directives
DBML has no dialect-neutral way to express database-specific features such as
PostgreSQL table partitioning or SQLite `WITHOUT ROWID`. RelSpec adds **dialect
directives** — explicit, parseable lines embedded in a `.dbml` file that are:
- stored losslessly in the intermediate model (under each object's `Metadata`),
- preserved unchanged through a `DBML → model → DBML` round-trip,
- translated to SQL **only** by the writer for the matching dialect
(`@postgres:` clauses appear in PostgreSQL output, never in SQLite output, and
vice-versa).
## Grammar
A directive is a single line, matched on its trimmed content:
```
@<namespace>[(<target>)]: <args>
```
| Part | Rules |
|------|-------|
| `namespace` | `^[a-z][a-z0-9_]*$` — e.g. `postgres`, `sqlite`. Future dialects allowed. |
| `(target)` | Optional. A **column name** only, valid only on a directive line inside a table body. Bare or single/double quoted. |
| `args` | Everything after the first `:`, trimmed. Otherwise preserved **verbatim**. Must be non-empty. |
The **key** of a directive is derived: the lowercased first whitespace-delimited
token of `args` (`partition by RANGE (created_at)``partition`). It drives
duplicate detection and writer dispatch.
## Location
Where the line appears determines which object it attaches to:
| Position in the file | Attaches to |
|----------------------|-------------|
| Before the first `Table {` | database (`db.Metadata`) |
| Table body, no `(target)` | that table |
| Table body, `(col)` target | column `col` of that table (error if `col` is unknown) |
| Inside an `indexes { }` block | the **most recently listed** index entry in that block; `(target)` is not allowed |
```dbml
@postgres: search_path myapp -- database
Table myapp.events {
id bigint [pk]
created_at timestamp [not null]
@postgres(id): identity always -- column "id"
@postgres: partition by RANGE (created_at) -- table
@postgres: tablespace fast_data -- table
@sqlite: without rowid -- table
indexes {
(created_at) [name: 'idx_events_created']
@postgres: with (fillfactor=90) -- index "idx_events_created"
@postgres: tablespace idx_space -- index "idx_events_created"
}
}
```
## Duplicate policy
- **Repeatable by default** — every directive with the same `(namespace, key)` at
one location is kept, in source order.
- **Singletons** raise a line-numbered error on a second occurrence at the same
location. Current singletons: `postgres` `partition`, `tablespace`, `inherits`,
`storage`, `compression`, `identity`; `sqlite` `without`, `strict`, `collate`.
## Strict mode
CLI flag `--strict-directives` (also `ReaderOptions.StrictDirectives` /
`WriterOptions.StrictDirectives`):
- **Reader**: an unknown namespace or key is a hard error. Without strict mode it
is stored and preserved silently, and round-trips unchanged.
- **PostgreSQL / SQLite writer**: a directive for **that** writer's own dialect
whose key it cannot translate is a hard error. Without strict mode, translatable
keys are emitted and the rest are skipped. Directives for other dialects are
always ignored, never emitted.
## Errors
All are line-numbered (`dbml: line N: …`):
- no colon, or empty `args`
- namespace empty or not matching `[a-z][a-z0-9_]*`
- `(target)` naming an unknown column, or used at the top level / in an `indexes` block
- a directive in the catalog used at a location it is not valid for
- duplicate singleton at the same location
- (strict mode) unknown `(namespace, key)`
## Supported directive matrix
### `@postgres`
| Key | Locations | SQL emitted | Notes |
|-----|-----------|-------------|-------|
| `partition` | table | `PARTITION BY <args>` appended to `CREATE TABLE` | e.g. `@postgres: partition by RANGE (created_at)` |
| `inherits` | table | `INHERITS (<args>)` — args verbatim | |
| `with` | table, index | `WITH (<params>)` | On an index, wins over `WITH` derived from the index comment. `@postgres: with (fillfactor=90)` |
| `tablespace` | table, index | `TABLESPACE <name>` | Emitted after `WITH`, before `WHERE` on indexes |
| `storage` | column | `STORAGE <mode>` in the column definition | e.g. `@postgres(blob): storage external` |
| `compression` | column | `COMPRESSION <method>` | |
| `identity` | column | `identity always``GENERATED ALWAYS AS IDENTITY`; `identity default` / `identity by default``GENERATED BY DEFAULT AS IDENTITY` | |
### `@sqlite`
| Key | Locations | SQL emitted | Notes |
|-----|-----------|-------------|-------|
| `without` | table | `WITHOUT ROWID` table option | `@sqlite: without rowid` |
| `strict` | table | `STRICT` table option | `WITHOUT ROWID` is emitted before `STRICT` |
| `collate` | column | ` COLLATE <name>` in the column definition | e.g. `@sqlite(name): collate NOCASE` |
Unknown namespaces and keys not in these tables are still preserved losslessly
(and round-trip through the DBML writer) whenever strict mode is off.
+28 -161
View File
@@ -9,11 +9,10 @@ relspec job run build-schema --plan # validate + print plan, execute nothing
relspec job run build-schema # run the job (and its dependencies)
```
## Design contract
## Design contract (first release)
This is a deliberately small, safe contract. Every capability is offline-testable
except live database execution (`scripts-exec`), which is validated and planned
offline and only connects at run time.
This is the smallest coherent contract that is safe and useful end to end.
Anything not listed under "Supported" is intentionally deferred.
### Not a shell
@@ -25,15 +24,14 @@ means adding a vetted adapter in the RelSpec source.
|----------------|--------------------------------------------------------------------|
| `convert` | read one or more input schemas, additively merge them, write one output |
| `merge` | like `convert` but requires ≥2 inputs and exposes `skip_*` merge options |
| `split` | read one or more schemas, keep the selected schemas/tables, write one output |
| `scripts-list` | deterministically list SQL scripts across one or more directories |
| `scripts-exec` | execute SQL scripts across one or more directories against a live PostgreSQL database |
| `templ` | apply a custom Go text template to one or more input schemas |
| `inspect` | validate one or more schemas against rules and write a report |
| `diff` | compare exactly two schemas and write a differences report |
`convert`, `merge` and `split` are **producers**: their file output can be fed
directly into another job with `from_job` (see below).
Deferred (documented, not implemented here): `scripts` execution against a live
database, `split`, `inspect`, `diff`, job-to-job output wiring,
log rotation/retention. Live SQL execution already exists as
`relspec scripts execute`; wiring it into the job runner is a follow-up because
it needs live database credentials and cannot be covered by offline tests.
### Discovery and precedence
@@ -51,16 +49,12 @@ already forbid duplicate keys within a single file.
### Paths
* Every path (`inputs[].path`, `output.path`, `report.path`, `rules`,
`script_dirs[]`, `template`, `logfile`) is **relative to the directory
containing the job file that declared the job**, not the process working
directory.
* Every path (`inputs[].path`, `output.path`, `script_dirs[]`, `logfile`) is
**relative to the directory containing the job file that declared the job**,
not the process working directory.
* Absolute paths, `~`-relative paths and any path that resolves outside the job
file directory (`../`, `a/../../b`, …) are **rejected during validation**
before anything runs.
* At run time each path is additionally resolved through its symlinks: a symlink
inside the job-file directory that points outside it is rejected before the
path is opened.
### Credentials
@@ -82,104 +76,58 @@ already forbid duplicate keys within a single file.
first. Nothing is read, written, connected to, or executed if validation fails.
Checks include:
* schema `version` **forward-permissive**: any version `>= 1` is accepted.
An omitted `version` is treated as the current one. A version newer than this
build understands loads best-effort (unknown YAML fields are ignored and a
warning is printed); at the current version unknown YAML fields are still
rejected.
* schema `version` (must be `1`), unknown YAML fields rejected
* duplicate job names across files
* unknown / missing `command`
* per-command input/output shape:
* `convert` needs ≥1 input + output; `merge` needs ≥2 inputs + output
* `split` needs ≥1 input + a file output, plus an optional `select:` block
* `scripts-list` needs `script_dirs` and forbids inputs/output
* `scripts-exec` needs `script_dirs` and `output.conn_env` (pgsql only)
* `inspect` needs ≥1 input + `report:` (format `markdown`|`json`)
* `diff` needs **exactly 2** inputs + `report:` (format `summary`|`json`|`html`)
* per-command input/output shape (`convert`/`merge` need inputs + output;
`scripts-list` needs `script_dirs` and forbids inputs/output)
* unknown input/output `format`
* `from_job` targets exist, are producers (`convert`/`merge`/`split`) and write a
single-file output
* path traversal / absolute / home-relative paths
* `depends_on` and `from_job` targets exist
* dependency cycles over the combined `depends_on` + `from_job` graph
(reported as `a -> b -> c -> a`)
* `depends_on` targets exist
* dependency cycles (reported as `a -> b -> c -> a`)
Then, immediately before running, per-job pre-flight resolves paths and checks:
* every input file exists and is a file (a `from_job` input is exempt — its
producer runs earlier in the same plan)
* every input file exists and is a file
* every `script_dir` exists and is a directory
* every `conn_env` variable is set
* `output.path` / `report.path` does not already exist unless the matching
`overwrite: true` is set
* `rules` (inspect), when given, exists and is a file
* symlinks in every resolved path stay inside the job-file directory
* `output.path` does not already exist unless `output.overwrite: true`
If any pre-flight check fails for **any** job in the plan, **no** job runs.
### Execution and exit codes
* `relspec job run <name>` runs the job's dependency closure first
(`depends_on` plus any `from_job` producers), in topological order
(deterministic), then the job. `--no-deps` runs only the named job and is
incompatible with `from_job` inputs.
* `relspec job run <name>` runs the job's `depends_on` closure first, in
topological order (deterministic), then the job. `--no-deps` runs only the
named job.
* `--dry-run` (alias `--plan`) prints the resolved plan and exits 0 without
touching inputs, outputs or databases.
* A failing job returns the underlying non-zero status (the process exits 1)
and the error names the job. The logfile records `FAILED: <error>`; a
successful job records `OK`. No separate success-marker file is written, so a
failure can never leave a stale "success".
* `inspect` fails the job when the report contains rule **errors** (enforced
rules); warnings do not fail it. `diff` never fails on differences.
* Single-file outputs and reports are written to a temporary file in the target
directory and atomically renamed into place, so an interrupted run never
leaves a partial file. Directory-emitting formats (`gorm`, `bun`, `drizzle`,
`typeorm`, `prisma`) are written in place.
### Logfile rotation
When a job has a `logfile`, it is size-rotated before each run. Defaults are
**5 MB** with **3** rotated files kept (`build.log``build.log.1` → …). Override
per job with `log_max_size` / `log_keep`, or for a whole file with a top-level
`defaults:` block. `log_max_size` accepts `B`/`KB`/`MB`/`GB` suffixes (e.g.
`"512KB"`, `"5MB"`).
## Schema reference
```yaml
version: 1 # optional; any value >= 1 is accepted
defaults: # optional, file-wide
log_max_size: 5MB # B / KB / MB / GB
log_keep: 3
version: 1 # required, must be 1
jobs:
<job-name>:
command: convert | merge | split | scripts-list | scripts-exec | templ | inspect | diff
command: convert | merge | scripts-list # required
description: "free text" # optional, shown by `job list`
depends_on: [other-job, ...] # optional
inputs: # convert (≥1) / merge (≥2) / split (≥1) / inspect (≥1) / diff (exactly 2)
inputs: # convert (≥1) / merge (≥2)
- path: relative/file.dbml # file inputs
format: dbml
- format: pgsql # live-connection inputs
conn_env: SOURCE_DB_URL # env var NAME
- from_job: build-schema # consume another job's file output
script_dirs: # scripts-list / scripts-exec (≥1)
script_dirs: # scripts-list (≥1)
- migrations/core
- migrations/tenant
template: templates/schema.tmpl # templ (required)
mode: table # templ: database/schema/script/table
filename_pattern: "{{.Name}}.go" # templ multi-output modes
select: # split (optional; default = keep everything)
schemas: [public]
tables: [users, orders]
exclude_schemas: []
exclude_tables: []
database_name: SubsetDB # optional rename of the output database
rules: .relspec-rules.yaml # inspect (optional; built-in defaults if omitted)
report: # inspect (required) / diff (required)
format: json # inspect: markdown|json ; diff: summary|json|html
path: build/report.json # required, except a diff "summary" (goes to the log)
overwrite: false
output: # convert / merge / split (required); scripts-exec (required, conn_env)
output: # convert / merge (required)
format: pgsql
path: build/schema.sql # file output, OR:
conn_env: TARGET_DB_URL # execute against DB (pgsql only)
@@ -188,15 +136,13 @@ jobs:
flatten_schema: false
schema: public
package: models # for gorm/bun output
continue_on_error: false # pgsql / scripts-exec output
continue_on_error: false # pgsql output
skip_relations: false # merge only
skip_enums: false
skip_views: false
skip_domains: false
skip_sequences: false
logfile: .relspec/log/<job-name>.log # optional; appended to, size-rotated
log_max_size: 5MB # optional per-job override
log_keep: 3 # optional per-job override
logfile: .relspec/log/<job-name>.log # optional; appended to
```
For `templ`, `inputs` use the same file or `pgsql`/`conn_env` source forms as
@@ -204,11 +150,6 @@ schema conversion. `output` is optional (empty means stdout); when present it
contains only `path` and `overwrite`, because templates do not select a schema
writer format.
A `from_job` input takes no `path`, `format` or `conn_env`: it resolves to the
named job's `output.path` and inherits its format, and implies a dependency on
that job. The producer must be a `convert`, `merge` or `split` job writing a
single-file output.
### Supported input formats
`dbml`, `dctx`, `drawdb`, `graphql`, `json`, `yaml`, `gorm`, `bun`, `drizzle`,
@@ -289,77 +230,3 @@ jobs:
path: snapshots/prod.dbml
overwrite: true
```
### Chain jobs with `from_job`, then lint the result
```yaml
version: 1
jobs:
build-json:
command: convert
inputs:
- { path: schema/core.dbml, format: dbml }
- { path: schema/tenant.dbml, format: dbml }
output: { format: json, path: build/schema.json, overwrite: true }
lint-schema:
command: inspect
inputs:
- from_job: build-json # implies depends_on: [build-json]
rules: .relspec-rules.yaml # optional; built-in rules if omitted
report:
format: markdown
path: build/lint-report.md
overwrite: true
```
`relspec job run lint-schema` runs `build-json` first, then inspects its output.
The job fails (exit 1) if any enforced rule is violated.
### Split a subset out of a larger schema
```yaml
version: 1
jobs:
posts-only:
command: split
inputs:
- { path: schema/core.dbml, format: dbml }
- { path: schema/tenant.dbml, format: dbml }
select:
tables: [posts]
output: { format: dbml, path: build/posts.dbml, overwrite: true }
```
### Diff two schemas
```yaml
version: 1
jobs:
drift:
command: diff
inputs: # exactly two
- { path: build/schema.json, format: json }
- format: pgsql
conn_env: PROD_DB_URL
report:
format: summary # summary → logfile; json/html need a path
```
`diff` reports differences and always exits 0.
### Execute migration scripts against a live database
```yaml
version: 1
jobs:
apply-migrations:
command: scripts-exec
script_dirs:
- migrations/core
- migrations/tenant
output:
conn_env: TARGET_DB_URL # pgsql only; no path
options:
continue_on_error: false
logfile: .relspec/log/apply-migrations.log
```
-34
View File
@@ -4,14 +4,7 @@
# relspec job list
# relspec job run build-schema --plan
# relspec job run build-schema
# relspec job run lint-schema # inspect, consuming build-json's output
version: 1
# File-wide defaults. Individual jobs may override log_max_size / log_keep.
defaults:
log_max_size: 2MB
log_keep: 5
jobs:
build-schema:
command: convert
@@ -50,30 +43,3 @@ jobs:
- migrations/core
- migrations/tenant
logfile: .relspec/log/migration-order.log
lint-schema:
command: inspect
description: Validate build-json's output against the built-in rules
# No depends_on needed: the from_job input implies a dependency on build-json.
inputs:
- from_job: build-json
report:
format: markdown
path: build/lint-report.md
overwrite: true
logfile: .relspec/log/lint-schema.log
posts-only:
command: split
description: Extract just the posts table into its own DBML file
inputs:
- path: schema/core.dbml
format: dbml
- path: schema/tenant.dbml
format: dbml
select:
tables: [posts]
output:
format: dbml
path: build/posts.dbml
overwrite: true
+12 -387
View File
@@ -20,52 +20,24 @@ import (
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"gopkg.in/yaml.v3"
)
// CurrentSchemaVersion is the highest job-file schema version this build was
// written for. MinSchemaVersion is the oldest it still accepts. A file that
// declares a version in between loads normally; a newer version loads
// best-effort with a warning (see Load); an older-than-minimum version is a
// hard error.
const (
CurrentSchemaVersion = 1
MinSchemaVersion = 1
)
// Built-in logfile rotation policy, used when neither the job nor its file's
// defaults block sets one.
const (
defaultLogMaxSizeBytes int64 = 5 << 20 // 5 MiB
defaultLogKeep = 3
)
// SchemaVersion is the only job-file schema version this build understands.
const SchemaVersion = 1
// Command names are a closed allow-list. Arbitrary strings are rejected.
const (
CommandConvert = "convert" // read one or more schema files, optionally merge, write one output
CommandMerge = "merge" // additive merge of two or more schema files into one output
CommandScriptsList = "scripts-list" // deterministically list SQL scripts across one or more directories
CommandScriptsExec = "scripts-exec" // execute SQL scripts across one or more directories against a live database
CommandTempl = "templ" // apply a custom Go text template to one or more schemas
CommandSplit = "split" // extract selected schemas/tables into a separate output
CommandInspect = "inspect" // validate one or more schemas against rules and write a report
CommandDiff = "diff" // compare exactly two schemas and write a differences report
)
// SupportedCommands lists every accepted command, in help order.
var SupportedCommands = []string{
CommandConvert, CommandMerge, CommandScriptsList, CommandScriptsExec,
CommandTempl, CommandSplit, CommandInspect, CommandDiff,
}
// producerCommands are commands whose output is a schema file that another job
// may consume via from_job.
var producerCommands = map[string]bool{
CommandConvert: true, CommandMerge: true, CommandSplit: true,
}
var SupportedCommands = []string{CommandConvert, CommandMerge, CommandScriptsList, CommandTempl}
// readerFormats are the file-based input formats a job may declare (path).
var readerFormats = map[string]bool{
@@ -89,51 +61,17 @@ var writerFormats = map[string]bool{
// live database) is supported instead of writing a file.
var execOutputFormats = map[string]bool{"pgsql": true}
// singleFileFormats are output formats that emit exactly one file (as opposed
// to a directory of files). Only these are eligible for atomic temp+rename
// writes and for being consumed by another job via from_job.
var singleFileFormats = map[string]bool{
"json": true, "yaml": true, "dbml": true, "dctx": true, "drawdb": true,
"graphql": true, "pgsql": true, "mssql": true, "sqlite": true,
}
// SingleFileOutputFormat reports whether format writes exactly one file.
func SingleFileOutputFormat(format string) bool {
return singleFileFormats[strings.ToLower(format)]
}
// diffReportFormats and inspectReportFormats are the report.format values
// accepted by the diff and inspect commands respectively.
var (
diffReportFormats = map[string]bool{"summary": true, "json": true, "html": true}
inspectReportFormats = map[string]bool{"markdown": true, "json": true}
)
// File is the on-disk shape of a single job file.
type File struct {
Version int `yaml:"version"`
Defaults *Defaults `yaml:"defaults"`
Jobs map[string]*Job `yaml:"jobs"`
}
// Defaults carries file-wide settings that individual jobs may override.
type Defaults struct {
// LogMaxSize is a human-readable size ("5MB", "512KB", "1GB"). Empty
// means "use the built-in default".
LogMaxSize string `yaml:"log_max_size"`
// LogKeep is how many rotated logfiles to retain. Zero means "use the
// built-in default".
LogKeep int `yaml:"log_keep"`
}
// Job is one named job within a job file.
type Job struct {
// Name and SourceFile are populated by Load, not parsed from YAML.
Name string `yaml:"-"`
SourceFile string `yaml:"-"`
// fileDefaults is the Defaults block of the file that declared this job,
// captured by Load. nil when the file had none.
fileDefaults *Defaults `yaml:"-"`
Command string `yaml:"command"`
Description string `yaml:"description"`
@@ -144,13 +82,8 @@ type Job struct {
Mode string `yaml:"mode"`
FilenamePattern string `yaml:"filename_pattern"`
Output *Output `yaml:"output"`
Rules string `yaml:"rules"`
Report *Report `yaml:"report"`
Select *Select `yaml:"select"`
Options Options `yaml:"options"`
Logfile string `yaml:"logfile"`
LogMaxSize string `yaml:"log_max_size"`
LogKeep *int `yaml:"log_keep"`
}
// Input is one declared input schema.
@@ -161,10 +94,6 @@ type Input struct {
// ConnEnv is the NAME of an environment variable holding a connection
// string, used with database formats. The value is never stored here.
ConnEnv string `yaml:"conn_env"`
// FromJob names another job in the set whose file output is used as this
// input. It implies a dependency on that job. Path/Format/ConnEnv must be
// empty when FromJob is set; the format is inherited from the producer.
FromJob string `yaml:"from_job"`
}
// Output is the declared output target.
@@ -175,105 +104,6 @@ type Output struct {
Overwrite bool `yaml:"overwrite"`
}
// Report is the output target for the inspect and diff commands.
type Report struct {
// Format is the report format: diff accepts summary|json|html, inspect
// accepts markdown|json. Empty means the command's default.
Format string `yaml:"format"`
Path string `yaml:"path"`
Overwrite bool `yaml:"overwrite"`
}
// Select carries the schema/table selection for the split command.
type Select struct {
Schemas []string `yaml:"schemas"`
Tables []string `yaml:"tables"`
ExcludeSchemas []string `yaml:"exclude_schemas"`
ExcludeTables []string `yaml:"exclude_tables"`
DatabaseName string `yaml:"database_name"`
}
// LogPolicy is the resolved logfile rotation policy for a job.
type LogPolicy struct {
MaxSizeBytes int64
Keep int
}
// ResolvedLogPolicy returns the effective rotation policy: the job's own
// overrides win, then its file's defaults block, then the built-in default.
func (j *Job) ResolvedLogPolicy() LogPolicy {
p := LogPolicy{MaxSizeBytes: defaultLogMaxSizeBytes, Keep: defaultLogKeep}
if j.fileDefaults != nil {
if n, err := parseHumanSize(j.fileDefaults.LogMaxSize); err == nil && n > 0 {
p.MaxSizeBytes = n
}
if j.fileDefaults.LogKeep > 0 {
p.Keep = j.fileDefaults.LogKeep
}
}
if n, err := parseHumanSize(j.LogMaxSize); err == nil && n > 0 {
p.MaxSizeBytes = n
}
if j.LogKeep != nil && *j.LogKeep >= 0 {
p.Keep = *j.LogKeep
}
return p
}
// effectiveDeps returns the union of explicit depends_on entries and the jobs
// referenced by from_job inputs, deduplicated in stable order.
func (j *Job) effectiveDeps() []string {
seen := map[string]bool{}
var deps []string
add := func(name string) {
if name == "" || name == j.Name || seen[name] {
return
}
seen[name] = true
deps = append(deps, name)
}
for _, d := range j.DependsOn {
add(d)
}
for _, in := range j.Inputs {
add(in.FromJob)
}
return deps
}
// parseHumanSize parses a byte size such as "5MB", "512 KB", "1gb" or a bare
// byte count. An empty string returns (0, nil) so callers can fall back.
func parseHumanSize(s string) (int64, error) {
s = strings.TrimSpace(s)
if s == "" {
return 0, nil
}
upper := strings.ToUpper(s)
mult := int64(1)
// Check multi-character suffixes before the bare "B".
for _, u := range []struct {
suffix string
m int64
}{
{"KB", 1 << 10}, {"MB", 1 << 20}, {"GB", 1 << 30}, {"B", 1},
} {
if strings.HasSuffix(upper, u.suffix) {
mult = u.m
upper = strings.TrimSpace(strings.TrimSuffix(upper, u.suffix))
break
}
}
n, err := strconv.ParseFloat(upper, 64)
if err != nil {
return 0, fmt.Errorf("invalid size %q", s)
}
if n < 0 {
return 0, fmt.Errorf("negative size %q", s)
}
return int64(n * float64(mult)), nil
}
// Options carries the subset of command flags a job file may set.
type Options struct {
FlattenSchema bool `yaml:"flatten_schema"`
@@ -297,9 +127,6 @@ type Set struct {
Files []string
// Jobs is keyed by job name.
Jobs map[string]*Job
// Warnings holds non-fatal load-time messages (e.g. a newer-than-known
// schema version). Callers should surface these to the user.
Warnings []string
}
// Names returns all job names in deterministic (sorted) order.
@@ -372,35 +199,15 @@ func Load(paths []string) (*Set, error) {
if err != nil {
return nil, fmt.Errorf("failed to read job file %q: %w", path, err)
}
// Peek at the version first so a newer file can be parsed leniently
// (unknown fields ignored) instead of failing outright.
var probe struct {
Version int `yaml:"version"`
}
if err := yaml.Unmarshal(data, &probe); err != nil {
return nil, fmt.Errorf("invalid job file %q: %w", path, err)
}
version := probe.Version
if version == 0 {
version = CurrentSchemaVersion
}
if version < MinSchemaVersion {
return nil, fmt.Errorf("job file %q: unsupported version %d (this build accepts %d or newer)", path, version, MinSchemaVersion)
}
strict := version <= CurrentSchemaVersion
if !strict {
set.Warnings = append(set.Warnings, fmt.Sprintf(
"job file %q declares version %d, newer than this build understands (%d); loading best-effort and ignoring unknown fields",
path, version, CurrentSchemaVersion))
}
dec := yaml.NewDecoder(strings.NewReader(string(data)))
dec.KnownFields(strict)
dec.KnownFields(true)
var f File
if err := dec.Decode(&f); err != nil {
return nil, fmt.Errorf("invalid job file %q: %w", path, err)
}
if f.Version != SchemaVersion {
return nil, fmt.Errorf("job file %q: unsupported version %d (expected %d)", path, f.Version, SchemaVersion)
}
if len(f.Jobs) == 0 {
return nil, fmt.Errorf("job file %q: no jobs defined", path)
}
@@ -413,7 +220,6 @@ func Load(paths []string) (*Set, error) {
}
job.Name = name
job.SourceFile = path
job.fileDefaults = f.Defaults
origin[name] = path
set.Jobs[name] = job
}
@@ -434,30 +240,13 @@ func (s *Set) Validate() error {
errs = append(errs, fmt.Sprintf("job %q: %s", name, msg))
}
}
// Dependency references + cycles + from_job wiring.
// Dependency references + cycles.
for _, name := range s.Names() {
j := s.Jobs[name]
for _, dep := range j.DependsOn {
for _, dep := range s.Jobs[name].DependsOn {
if _, ok := s.Jobs[dep]; !ok {
errs = append(errs, fmt.Sprintf("job %q: depends_on unknown job %q", name, dep))
}
}
for i, in := range j.Inputs {
if in.FromJob == "" {
continue
}
producer, ok := s.Jobs[in.FromJob]
if !ok {
errs = append(errs, fmt.Sprintf("job %q: input[%d] from_job references unknown job %q", name, i, in.FromJob))
continue
}
if !producerCommands[producer.Command] || producer.Output == nil ||
producer.Output.Path == "" || !SingleFileOutputFormat(producer.Output.Format) {
errs = append(errs, fmt.Sprintf(
"job %q: input[%d] from_job %q must name a convert/merge/split job that writes a single-file output",
name, i, in.FromJob))
}
}
}
if cycle := s.findCycle(); cycle != "" {
errs = append(errs, fmt.Sprintf("dependency cycle detected: %s", cycle))
@@ -473,8 +262,7 @@ func (j *Job) validate() []string {
var e []string
switch j.Command {
case CommandConvert, CommandMerge, CommandScriptsList, CommandScriptsExec,
CommandTempl, CommandSplit, CommandInspect, CommandDiff:
case CommandConvert, CommandMerge, CommandScriptsList, CommandTempl:
case "":
e = append(e, "missing command")
return e
@@ -494,7 +282,6 @@ func (j *Job) validate() []string {
}
checkPath("logfile", j.Logfile)
checkPath("template", j.Template)
checkPath("rules", j.Rules)
for _, in := range j.Inputs {
checkPath("input path", in.Path)
}
@@ -504,13 +291,6 @@ func (j *Job) validate() []string {
if j.Output != nil {
checkPath("output path", j.Output.Path)
}
if j.Report != nil {
checkPath("report path", j.Report.Path)
}
if _, err := parseHumanSize(j.LogMaxSize); err != nil {
e = append(e, fmt.Sprintf("log_max_size: %v", err))
}
switch j.Command {
case CommandConvert, CommandMerge:
@@ -570,113 +350,11 @@ func (j *Job) validate() []string {
if j.Output != nil && j.Output.Format != "" {
e = append(e, "output.format is not valid for command \"templ\"")
}
case CommandSplit:
if len(j.Inputs) < 1 {
e = append(e, "command \"split\" requires at least 1 input")
}
for i, in := range j.Inputs {
e = append(e, validateInput(i, in)...)
}
if len(j.ScriptDirs) > 0 {
e = append(e, "script_dirs is not valid for command \"split\"")
}
if j.Report != nil {
e = append(e, "report is not valid for command \"split\" (use output)")
}
if j.Output == nil {
e = append(e, "missing output")
} else {
if j.Output.ConnEnv != "" {
e = append(e, "command \"split\" writes a file; output.conn_env is not supported")
}
e = append(e, validateOutput(*j.Output)...)
}
case CommandInspect:
if len(j.Inputs) < 1 {
e = append(e, "command \"inspect\" requires at least 1 input")
}
for i, in := range j.Inputs {
e = append(e, validateInput(i, in)...)
}
if len(j.ScriptDirs) > 0 {
e = append(e, "script_dirs is not valid for command \"inspect\"")
}
if j.Output != nil {
e = append(e, "output is not valid for command \"inspect\" (use report)")
}
e = append(e, validateReport(j.Report, "inspect", inspectReportFormats, "markdown")...)
case CommandDiff:
if len(j.Inputs) != 2 {
e = append(e, "command \"diff\" requires exactly 2 inputs (source, target)")
}
for i, in := range j.Inputs {
e = append(e, validateInput(i, in)...)
}
if len(j.ScriptDirs) > 0 {
e = append(e, "script_dirs is not valid for command \"diff\"")
}
if j.Output != nil {
e = append(e, "output is not valid for command \"diff\" (use report)")
}
e = append(e, validateReport(j.Report, "diff", diffReportFormats, "summary")...)
case CommandScriptsExec:
if len(j.ScriptDirs) == 0 {
e = append(e, "command \"scripts-exec\" requires at least one script_dir")
}
if len(j.Inputs) > 0 {
e = append(e, "inputs is not valid for command \"scripts-exec\"")
}
if j.Report != nil {
e = append(e, "report is not valid for command \"scripts-exec\"")
}
if j.Output == nil || j.Output.ConnEnv == "" {
e = append(e, "command \"scripts-exec\" requires output.conn_env (an environment variable name holding a connection string)")
} else {
if j.Output.Path != "" {
e = append(e, "command \"scripts-exec\" executes against a database; output.path is not supported")
}
f := strings.ToLower(j.Output.Format)
if f != "" && f != "pgsql" {
e = append(e, fmt.Sprintf("command \"scripts-exec\" only supports pgsql databases (got %q)", j.Output.Format))
}
if looksLikeSecret(j.Output.ConnEnv) {
e = append(e, "output: conn_env must be an environment variable name, not a connection string")
}
}
}
return e
}
// validateReport checks a Report block for the inspect/diff commands.
func validateReport(r *Report, cmd string, allowed map[string]bool, defFmt string) []string {
if r == nil {
return []string{fmt.Sprintf("command %q requires a report block", cmd)}
}
var e []string
f := strings.ToLower(r.Format)
if f == "" {
f = defFmt
}
if !allowed[f] {
names := make([]string, 0, len(allowed))
for k := range allowed {
names = append(names, k)
}
sort.Strings(names)
e = append(e, fmt.Sprintf("command %q report.format %q is not supported (use: %s)", cmd, r.Format, strings.Join(names, ", ")))
}
// A diff summary may be written to the log; everything else needs a path.
summaryToLog := cmd == "diff" && f == "summary"
if r.Path == "" && !summaryToLog {
e = append(e, fmt.Sprintf("command %q requires report.path", cmd))
}
return e
}
func validateTemplInput(i int, in Input) []string {
if in.FromJob != "" {
return fromJobInputShape(i, in)
}
var e []string
if in.Format == "" {
return []string{fmt.Sprintf("input[%d]: missing format", i)}
@@ -705,27 +383,7 @@ func validateTemplInput(i int, in Input) []string {
return e
}
// fromJobInputShape checks the structural rules for an input that pulls its
// schema from another job's output. The referenced job's existence and kind
// are checked in Set.Validate, which can see the whole set.
func fromJobInputShape(i int, in Input) []string {
var e []string
if in.Path != "" {
e = append(e, fmt.Sprintf("input[%d]: from_job takes no path", i))
}
if in.Format != "" {
e = append(e, fmt.Sprintf("input[%d]: from_job inherits the producer's format; drop format", i))
}
if in.ConnEnv != "" {
e = append(e, fmt.Sprintf("input[%d]: from_job takes no conn_env", i))
}
return e
}
func validateInput(i int, in Input) []string {
if in.FromJob != "" {
return fromJobInputShape(i, in)
}
var e []string
if in.Format == "" {
e = append(e, fmt.Sprintf("input[%d]: missing format", i))
@@ -829,42 +487,9 @@ func SafeJoin(root, rel string) (string, error) {
if rp == ".." || strings.HasPrefix(rp, ".."+string(filepath.Separator)) {
return "", fmt.Errorf("path %q escapes the job file directory", rel)
}
// Symlink hardening: resolve symlinks on the root and on the deepest
// existing ancestor of the target, and require the target to still live
// inside the resolved root. This catches a symlink inside the job-file
// directory that points outside it.
realRoot, err := filepath.EvalSymlinks(absRoot)
if err != nil {
return "", fmt.Errorf("cannot resolve job file directory: %w", err)
}
realAnc, err := filepath.EvalSymlinks(deepestExistingAncestor(joined))
if err != nil {
return "", fmt.Errorf("cannot resolve path %q: %w", rel, err)
}
if realAnc != realRoot {
if r, err := filepath.Rel(realRoot, realAnc); err != nil ||
r == ".." || strings.HasPrefix(r, ".."+string(filepath.Separator)) {
return "", fmt.Errorf("path %q resolves outside the job file directory via a symlink", rel)
}
}
return joined, nil
}
// deepestExistingAncestor returns p itself if it exists, otherwise the nearest
// existing parent directory (falling back to the filesystem root).
func deepestExistingAncestor(p string) string {
for {
if _, err := os.Lstat(p); err == nil {
return p
}
parent := filepath.Dir(p)
if parent == p {
return p
}
p = parent
}
}
// Plan returns the jobs to execute for name in dependency order. When
// includeDeps is false only the named job is returned (its declared
// dependencies are still validated to exist and be acyclic by Validate).
@@ -889,7 +514,7 @@ func (s *Set) Plan(name string, includeDeps bool) ([]*Job, error) {
}
inProgress[n] = true
j := s.Jobs[n]
deps := j.effectiveDeps()
deps := append([]string(nil), j.DependsOn...)
sort.Strings(deps)
for _, d := range deps {
if _, ok := s.Jobs[d]; !ok {
@@ -918,7 +543,7 @@ func (s *Set) findCycle() string {
dfs = func(n string) []string {
color[n] = 1
stack = append(stack, n)
deps := s.Jobs[n].effectiveDeps()
deps := append([]string(nil), s.Jobs[n].DependsOn...)
sort.Strings(deps)
for _, d := range deps {
if _, ok := s.Jobs[d]; !ok {
+5 -203
View File
@@ -56,82 +56,13 @@ func TestLoadRejectsUnknownFields(t *testing.T) {
}
}
func TestLoadWarnsOnNewerVersion(t *testing.T) {
func TestLoadRejectsBadVersion(t *testing.T) {
dir := t.TempDir()
p := filepath.Join(dir, "relspec.yml")
// A newer version loads best-effort with a warning, and unknown fields
// from the newer schema are ignored rather than rejected.
write(t, p, "version: 99\njobs:\n a:\n command: convert\n"+
" inputs:\n - path: a.dbml\n format: dbml\n"+
" output:\n format: json\n path: out.json\n"+
" future_field: whatever\n")
set, err := Load([]string{p})
if err != nil {
t.Fatalf("newer version should load, got %v", err)
}
if len(set.Warnings) == 0 {
t.Fatal("expected a warning about the newer version")
}
if err := set.Validate(); err != nil {
t.Fatalf("validate: %v", err)
}
}
func TestLoadAcceptsOmittedVersion(t *testing.T) {
dir := t.TempDir()
p := filepath.Join(dir, "relspec.yml")
write(t, p, "jobs:\n a:\n command: convert\n"+
" inputs:\n - path: a.dbml\n format: dbml\n"+
" output:\n format: json\n path: out.json\n")
set, err := Load([]string{p})
if err != nil {
t.Fatalf("omitted version should load, got %v", err)
}
if len(set.Warnings) != 0 {
t.Fatalf("omitted version should not warn, got %v", set.Warnings)
}
}
func TestLoadStillRejectsUnknownFieldsAtCurrentVersion(t *testing.T) {
dir := t.TempDir()
p := filepath.Join(dir, "relspec.yml")
write(t, p, "version: 1\njobs:\n a:\n command: convert\n bogus: true\n")
if _, err := Load([]string{p}); err == nil {
t.Fatal("expected unknown-field rejection at the current version")
}
}
func TestParseHumanSize(t *testing.T) {
cases := []struct {
in string
want int64
bad bool
}{
{"", 0, false},
{"512", 512, false},
{"512B", 512, false},
{"1KB", 1 << 10, false},
{"5MB", 5 << 20, false},
{"1gb", 1 << 30, false},
{" 2 MB ", 2 << 20, false},
{"nonsense", 0, true},
{"-1MB", 0, true},
}
for _, c := range cases {
got, err := parseHumanSize(c.in)
if c.bad {
if err == nil {
t.Errorf("parseHumanSize(%q): expected error", c.in)
}
continue
}
if err != nil {
t.Errorf("parseHumanSize(%q): %v", c.in, err)
continue
}
if got != c.want {
t.Errorf("parseHumanSize(%q) = %d, want %d", c.in, got, c.want)
}
write(t, p, "version: 2\njobs:\n a:\n command: convert\n")
_, err := Load([]string{p})
if err == nil || !strings.Contains(err.Error(), "unsupported version") {
t.Fatalf("expected unsupported version error, got %v", err)
}
}
@@ -308,135 +239,6 @@ func TestShippedExampleIsValid(t *testing.T) {
}
}
func TestFromJobWiring(t *testing.T) {
content := "version: 1\njobs:\n" +
" producer:\n command: convert\n" +
" inputs:\n - path: a.dbml\n format: dbml\n" +
" output:\n format: json\n path: build/schema.json\n" +
" consumer:\n command: convert\n" +
" inputs:\n - from_job: producer\n" +
" output:\n format: yaml\n path: build/schema.yaml\n"
set := loadOne(t, content)
if err := set.Validate(); err != nil {
t.Fatalf("validate: %v", err)
}
plan, err := set.Plan("consumer", true)
if err != nil {
t.Fatal(err)
}
if len(plan) != 2 || plan[0].Name != "producer" || plan[1].Name != "consumer" {
t.Fatalf("plan = %v, want [producer consumer]", plan)
}
}
func TestFromJobRejectsNonProducer(t *testing.T) {
content := "version: 1\njobs:\n" +
" lister:\n command: scripts-list\n script_dirs: [migrations]\n" +
" consumer:\n command: convert\n" +
" inputs:\n - from_job: lister\n" +
" output:\n format: yaml\n path: out.yaml\n"
set := loadOne(t, content)
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "from_job") {
t.Fatalf("want from_job producer error, got %v", err)
}
}
func TestFromJobRejectsUnknownJob(t *testing.T) {
content := "version: 1\njobs:\n" +
" consumer:\n command: convert\n" +
" inputs:\n - from_job: ghost\n" +
" output:\n format: yaml\n path: out.yaml\n"
set := loadOne(t, content)
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "unknown job") {
t.Fatalf("want unknown job error, got %v", err)
}
}
func TestFromJobCycleDetected(t *testing.T) {
content := "version: 1\njobs:\n" +
" a:\n command: convert\n" +
" inputs:\n - from_job: b\n" +
" output:\n format: json\n path: a.json\n" +
" b:\n command: convert\n" +
" inputs:\n - from_job: a\n" +
" output:\n format: json\n path: b.json\n"
set := loadOne(t, content)
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "cycle") {
t.Fatalf("want cycle error, got %v", err)
}
}
func TestSplitJobValidation(t *testing.T) {
set := loadOne(t, "version: 1\njobs:\n s:\n command: split\n"+
" inputs:\n - path: a.dbml\n format: dbml\n")
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "missing output") {
t.Fatalf("want missing output, got %v", err)
}
set = loadOne(t, "version: 1\njobs:\n s:\n command: split\n"+
" inputs:\n - path: a.dbml\n format: dbml\n"+
" select:\n tables: [users]\n"+
" output:\n format: json\n path: out.json\n")
if err := set.Validate(); err != nil {
t.Fatalf("expected valid split job, got %v", err)
}
}
func TestInspectJobValidation(t *testing.T) {
set := loadOne(t, "version: 1\njobs:\n i:\n command: inspect\n"+
" inputs:\n - path: a.dbml\n format: dbml\n")
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "report") {
t.Fatalf("want report required, got %v", err)
}
set = loadOne(t, "version: 1\njobs:\n i:\n command: inspect\n"+
" inputs:\n - path: a.dbml\n format: dbml\n"+
" report:\n format: json\n path: build/report.json\n")
if err := set.Validate(); err != nil {
t.Fatalf("expected valid inspect job, got %v", err)
}
}
func TestDiffJobValidation(t *testing.T) {
set := loadOne(t, "version: 1\njobs:\n d:\n command: diff\n"+
" inputs:\n - path: a.dbml\n format: dbml\n"+
" report:\n format: summary\n")
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "exactly 2 inputs") {
t.Fatalf("want exactly 2 inputs, got %v", err)
}
set = loadOne(t, "version: 1\njobs:\n d:\n command: diff\n"+
" inputs:\n - path: a.dbml\n format: dbml\n"+
" - path: b.dbml\n format: dbml\n"+
" report:\n format: summary\n")
if err := set.Validate(); err != nil {
t.Fatalf("expected valid diff job, got %v", err)
}
}
func TestScriptsExecValidation(t *testing.T) {
set := loadOne(t, "version: 1\njobs:\n x:\n command: scripts-exec\n"+
" script_dirs: [migrations]\n")
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "conn_env") {
t.Fatalf("want output.conn_env required, got %v", err)
}
set = loadOne(t, "version: 1\njobs:\n x:\n command: scripts-exec\n"+
" script_dirs: [migrations]\n"+
" output:\n conn_env: TARGET_DB_URL\n")
if err := set.Validate(); err != nil {
t.Fatalf("expected valid scripts-exec job, got %v", err)
}
}
func TestSafeJoinRejectsSymlinkEscape(t *testing.T) {
root := t.TempDir()
outside := t.TempDir()
link := filepath.Join(root, "link")
if err := os.Symlink(outside, link); err != nil {
t.Skipf("symlink not supported: %v", err)
}
if _, err := SafeJoin(root, "link/x.sql"); err == nil {
t.Fatal("expected rejection of a path escaping via a symlink")
}
}
func TestScriptsListValidation(t *testing.T) {
set := loadOne(t, "version: 1\njobs:\n s:\n command: scripts-list\n")
if err := set.Validate(); err == nil || !strings.Contains(err.Error(), "script_dir") {
+237
View File
@@ -0,0 +1,237 @@
package models
import (
"fmt"
"sort"
"strings"
)
// Directive is a dialect-specific instruction embedded in a source schema
// (currently DBML) that is preserved losslessly in the intermediate model and
// consumed only by the writer for its namespace. Directives are stored in the
// Metadata map of the object they apply to, under DirectivesMetadataKey.
//
// Example DBML: `@postgres: partition by RANGE (created_at)` parses to
// Directive{Namespace: "postgres", Key: "partition", Args: "partition by RANGE (created_at)"}.
type Directive struct {
// Namespace is the dialect the directive targets, e.g. "postgres" or "sqlite".
Namespace string `json:"namespace" yaml:"namespace"`
// Key is the lowercased first token of Args, used for duplicate detection
// and writer dispatch.
Key string `json:"key,omitempty" yaml:"key,omitempty"`
// Args is the verbatim argument text following the "@namespace:" prefix.
Args string `json:"args" yaml:"args"`
// Line is the 1-based source line the directive was read from, when known.
Line int `json:"line,omitempty" yaml:"line,omitempty"`
}
// DirectivesMetadataKey is the Metadata map key under which the ordered list of
// dialect directives for an object is stored.
const DirectivesMetadataKey = "directives"
// DirectiveKey derives the Key for a directive from its argument text: the
// lowercased first whitespace-delimited token.
func DirectiveKey(args string) string {
fields := strings.Fields(args)
if len(fields) == 0 {
return ""
}
return strings.ToLower(fields[0])
}
// AddDirective appends d to the directive list stored in meta. The caller is
// responsible for ensuring meta is non-nil (all Init* constructors allocate it).
// If d.Key is empty it is derived from d.Args.
func AddDirective(meta map[string]any, d Directive) {
if meta == nil {
return
}
if d.Key == "" {
d.Key = DirectiveKey(d.Args)
}
existing := GetDirectives(meta)
existing = append(existing, d)
meta[DirectivesMetadataKey] = existing
}
// GetDirectives returns the directives stored in meta, sorted deterministically
// by (Namespace, Line, Args). It tolerates both a freshly built []Directive and
// the []any of map[string]any produced by a JSON/YAML round-trip.
func GetDirectives(meta map[string]any) []Directive {
if meta == nil {
return nil
}
raw, ok := meta[DirectivesMetadataKey]
if !ok || raw == nil {
return nil
}
var out []Directive
switch v := raw.(type) {
case []Directive:
out = append(out, v...)
case []any:
for _, item := range v {
if d, ok := directiveFromAny(item); ok {
out = append(out, d)
}
}
}
sort.SliceStable(out, func(i, j int) bool {
if out[i].Namespace != out[j].Namespace {
return out[i].Namespace < out[j].Namespace
}
if out[i].Line != out[j].Line {
return out[i].Line < out[j].Line
}
return out[i].Args < out[j].Args
})
return out
}
// directiveFromAny decodes a single directive from the loosely typed forms that
// survive a JSON or YAML round-trip (map[string]any / map[any]any).
func directiveFromAny(item any) (Directive, bool) {
switch m := item.(type) {
case Directive:
return m, true
case map[string]any:
return directiveFromStringMap(m), true
case map[any]any:
sm := make(map[string]any, len(m))
for k, val := range m {
if ks, ok := k.(string); ok {
sm[ks] = val
}
}
return directiveFromStringMap(sm), true
}
return Directive{}, false
}
func directiveFromStringMap(m map[string]any) Directive {
d := Directive{}
if s, ok := m["namespace"].(string); ok {
d.Namespace = s
}
if s, ok := m["key"].(string); ok {
d.Key = s
}
if s, ok := m["args"].(string); ok {
d.Args = s
}
switch n := m["line"].(type) {
case int:
d.Line = n
case int64:
d.Line = int(n)
case float64:
d.Line = int(n)
}
if d.Key == "" {
d.Key = DirectiveKey(d.Args)
}
return d
}
// DirectivesForNamespace returns the directives in meta that target ns, in the
// deterministic order of GetDirectives.
func DirectivesForNamespace(meta map[string]any, ns string) []Directive {
all := GetDirectives(meta)
if len(all) == 0 {
return nil
}
out := make([]Directive, 0, len(all))
for _, d := range all {
if d.Namespace == ns {
out = append(out, d)
}
}
return out
}
// HasDirective reports whether meta contains a directive with the given
// namespace and key.
func HasDirective(meta map[string]any, ns, key string) bool {
for _, d := range GetDirectives(meta) {
if d.Namespace == ns && d.Key == key {
return true
}
}
return false
}
// DirectiveSpec describes a documented directive in the catalog.
type DirectiveSpec struct {
// Singleton means only one directive with this namespace/key may appear at
// a single location; a second one is a parse error.
Singleton bool
// Locations lists the location kinds the directive is valid at
// ("database", "table", "column", "index").
Locations []string
}
// Location kinds a directive may attach to.
const (
DirectiveLocationDatabase = "database"
DirectiveLocationTable = "table"
DirectiveLocationColumn = "column"
DirectiveLocationIndex = "index"
)
// DirectiveCatalog is the set of documented directives per namespace. It is used
// for strict-mode validation in readers and writers; unknown namespaces/keys are
// still preserved losslessly when strict mode is off.
var DirectiveCatalog = map[string]map[string]DirectiveSpec{
"postgres": {
"partition": {Singleton: true, Locations: []string{DirectiveLocationTable}},
"tablespace": {Singleton: true, Locations: []string{DirectiveLocationTable, DirectiveLocationIndex}},
"inherits": {Singleton: true, Locations: []string{DirectiveLocationTable}},
"with": {Singleton: false, Locations: []string{DirectiveLocationTable, DirectiveLocationIndex}},
"storage": {Singleton: true, Locations: []string{DirectiveLocationColumn}},
"compression": {Singleton: true, Locations: []string{DirectiveLocationColumn}},
"identity": {Singleton: true, Locations: []string{DirectiveLocationColumn}},
},
"sqlite": {
"without": {Singleton: true, Locations: []string{DirectiveLocationTable}},
"strict": {Singleton: true, Locations: []string{DirectiveLocationTable}},
"collate": {Singleton: true, Locations: []string{DirectiveLocationColumn}},
},
}
// LookupDirectiveSpec returns the catalog spec for a namespace/key and whether
// it is documented.
func LookupDirectiveSpec(ns, key string) (DirectiveSpec, bool) {
keys, ok := DirectiveCatalog[ns]
if !ok {
return DirectiveSpec{}, false
}
spec, ok := keys[key]
return spec, ok
}
// DirectiveLocationAllowed reports whether a documented directive may appear at
// the given location. Unknown directives (not in the catalog) are allowed
// everywhere so they can be preserved.
func DirectiveLocationAllowed(ns, key, location string) bool {
spec, ok := LookupDirectiveSpec(ns, key)
if !ok {
return true
}
for _, l := range spec.Locations {
if l == location {
return true
}
}
return false
}
// FormatDirectiveLine renders a directive back to its DBML source form, e.g.
// "@postgres: partition by RANGE (created_at)" or "@postgres(id): identity always".
func FormatDirectiveLine(d Directive, target string) string {
if target != "" {
return fmt.Sprintf("@%s(%s): %s", d.Namespace, target, d.Args)
}
return fmt.Sprintf("@%s: %s", d.Namespace, d.Args)
}
+133
View File
@@ -0,0 +1,133 @@
package models
import (
"encoding/json"
"testing"
)
func TestDirectiveKey(t *testing.T) {
cases := map[string]string{
"partition by RANGE (created_at)": "partition",
"WITHOUT ROWID": "without",
" strict ": "strict",
"": "",
}
for args, want := range cases {
if got := DirectiveKey(args); got != want {
t.Errorf("DirectiveKey(%q) = %q, want %q", args, got, want)
}
}
}
func TestAddDirectiveDerivesKey(t *testing.T) {
meta := map[string]any{}
AddDirective(meta, Directive{Namespace: "postgres", Args: "partition by RANGE (x)", Line: 2})
AddDirective(meta, Directive{Namespace: "postgres", Key: "tablespace", Args: "tablespace fast", Line: 3})
got := GetDirectives(meta)
if len(got) != 2 {
t.Fatalf("got %d directives, want 2", len(got))
}
if got[0].Key != "partition" {
t.Errorf("derived key = %q, want %q", got[0].Key, "partition")
}
if got[1].Key != "tablespace" {
t.Errorf("explicit key = %q, want %q", got[1].Key, "tablespace")
}
}
func TestAddDirectiveNilMeta(t *testing.T) {
// Must not panic.
AddDirective(nil, Directive{Namespace: "postgres", Args: "strict"})
}
func TestGetDirectivesOrdering(t *testing.T) {
meta := map[string]any{}
AddDirective(meta, Directive{Namespace: "sqlite", Args: "strict", Line: 9})
AddDirective(meta, Directive{Namespace: "postgres", Args: "with (b)", Line: 5})
AddDirective(meta, Directive{Namespace: "postgres", Args: "with (a)", Line: 5})
AddDirective(meta, Directive{Namespace: "postgres", Args: "partition by x", Line: 2})
got := GetDirectives(meta)
wantArgs := []string{"partition by x", "with (a)", "with (b)", "strict"}
if len(got) != len(wantArgs) {
t.Fatalf("got %d directives, want %d", len(got), len(wantArgs))
}
for i, w := range wantArgs {
if got[i].Args != w {
t.Errorf("directive[%d].Args = %q, want %q", i, got[i].Args, w)
}
}
}
func TestGetDirectivesTolerantDecodeAfterJSON(t *testing.T) {
meta := map[string]any{}
AddDirective(meta, Directive{Namespace: "postgres", Args: "partition by RANGE (created_at)", Line: 4})
AddDirective(meta, Directive{Namespace: "sqlite", Args: "without rowid", Line: 6})
blob, err := json.Marshal(meta)
if err != nil {
t.Fatalf("marshal: %v", err)
}
var round map[string]any
if err := json.Unmarshal(blob, &round); err != nil {
t.Fatalf("unmarshal: %v", err)
}
got := GetDirectives(round)
if len(got) != 2 {
t.Fatalf("got %d directives after JSON round-trip, want 2", len(got))
}
if got[0].Namespace != "postgres" || got[0].Key != "partition" || got[0].Line != 4 {
t.Errorf("post-JSON directive[0] = %+v", got[0])
}
if got[0].Args != "partition by RANGE (created_at)" {
t.Errorf("post-JSON args not verbatim: %q", got[0].Args)
}
if got[1].Namespace != "sqlite" || got[1].Key != "without" {
t.Errorf("post-JSON directive[1] = %+v", got[1])
}
}
func TestDirectivesForNamespaceAndHasDirective(t *testing.T) {
meta := map[string]any{}
AddDirective(meta, Directive{Namespace: "postgres", Args: "partition by x", Line: 1})
AddDirective(meta, Directive{Namespace: "sqlite", Args: "strict", Line: 2})
pg := DirectivesForNamespace(meta, "postgres")
if len(pg) != 1 || pg[0].Key != "partition" {
t.Errorf("DirectivesForNamespace(postgres) = %+v", pg)
}
if !HasDirective(meta, "sqlite", "strict") {
t.Error("HasDirective(sqlite, strict) = false, want true")
}
if HasDirective(meta, "postgres", "tablespace") {
t.Error("HasDirective(postgres, tablespace) = true, want false")
}
}
func TestDirectiveLocationAllowed(t *testing.T) {
if !DirectiveLocationAllowed("postgres", "partition", DirectiveLocationTable) {
t.Error("partition should be allowed at table level")
}
if DirectiveLocationAllowed("postgres", "partition", DirectiveLocationColumn) {
t.Error("partition should not be allowed at column level")
}
// Unknown directives are allowed everywhere so they can be preserved.
if !DirectiveLocationAllowed("postgres", "bogus", DirectiveLocationDatabase) {
t.Error("unknown key should be allowed everywhere")
}
if !DirectiveLocationAllowed("madeup", "x", DirectiveLocationTable) {
t.Error("unknown namespace should be allowed everywhere")
}
}
func TestFormatDirectiveLine(t *testing.T) {
d := Directive{Namespace: "postgres", Key: "identity", Args: "identity always"}
if got := FormatDirectiveLine(d, ""); got != "@postgres: identity always" {
t.Errorf("FormatDirectiveLine no target = %q", got)
}
if got := FormatDirectiveLine(d, "id"); got != "@postgres(id): identity always" {
t.Errorf("FormatDirectiveLine with target = %q", got)
}
}
+6
View File
@@ -32,6 +32,7 @@ type Database struct {
DatabaseType DatabaseType `json:"database_type,omitempty" yaml:"database_type,omitempty" xml:"database_type,omitempty"`
DatabaseVersion string `json:"database_version,omitempty" yaml:"database_version,omitempty" xml:"database_version,omitempty"`
SourceFormat string `json:"source_format,omitempty" yaml:"source_format,omitempty" xml:"source_format,omitempty"` // Source Format of the database.
Metadata map[string]any `json:"metadata,omitempty" yaml:"metadata,omitempty" xml:"-"`
UpdatedAt string `json:"updatedat,omitempty" yaml:"updatedat,omitempty" xml:"updatedat,omitempty"`
GUID string `json:"guid" yaml:"guid" xml:"guid"`
}
@@ -240,6 +241,7 @@ type Column struct {
IsPrimaryKey bool `json:"is_primary_key" yaml:"is_primary_key" xml:"is_primary_key"`
Comment string `json:"comment,omitempty" yaml:"comment,omitempty" xml:"comment,omitempty"`
Collation string `json:"collation,omitempty" yaml:"collation,omitempty" xml:"collation,omitempty"`
Metadata map[string]any `json:"metadata,omitempty" yaml:"metadata,omitempty" xml:"-"`
Sequence uint `json:"sequence,omitempty" yaml:"sequence,omitempty" xml:"sequence,omitempty"`
GUID string `json:"guid" yaml:"guid" xml:"guid"`
}
@@ -263,6 +265,7 @@ type Index struct {
Concurrent bool `json:"concurrent,omitempty" yaml:"concurrent,omitempty" xml:"concurrent,omitempty"`
Include []string `json:"include,omitempty" yaml:"include,omitempty" xml:"include,omitempty"` // INCLUDE columns
Comment string `json:"comment,omitempty" yaml:"comment,omitempty" xml:"comment,omitempty"`
Metadata map[string]any `json:"metadata,omitempty" yaml:"metadata,omitempty" xml:"-"`
Sequence uint `json:"sequence,omitempty" yaml:"sequence,omitempty" xml:"sequence,omitempty"`
GUID string `json:"guid" yaml:"guid" xml:"guid"`
}
@@ -396,6 +399,7 @@ func InitDatabase(name string) *Database {
Name: name,
Schemas: make([]*Schema, 0),
Domains: make([]*Domain, 0),
Metadata: make(map[string]any),
GUID: uuid.New().String(),
}
}
@@ -434,6 +438,7 @@ func InitColumn(name, table, schema string) *Column {
Name: name,
Table: table,
Schema: schema,
Metadata: make(map[string]any),
GUID: uuid.New().String(),
}
}
@@ -446,6 +451,7 @@ func InitIndex(name, table, schema string) *Index {
Schema: schema,
Columns: make([]string, 0),
Include: make([]string, 0),
Metadata: make(map[string]any),
GUID: uuid.New().String(),
}
}
+44
View File
@@ -93,6 +93,50 @@ Ref: posts.user_id > users.id [delete: cascade]
- Indexes and composite indexes
- Table notes and column notes
- Enums
- Dialect directives (`@postgres:` / `@sqlite:` — see below)
## Dialect directives
Lines of the form `@<namespace>[(<column>)]: <args>` embed database-specific
features that plain DBML cannot express (partitioning, `WITHOUT ROWID`,
tablespaces, index storage parameters, …). They are stored losslessly on the
relevant object's `Metadata` and round-trip unchanged through the DBML writer;
the PostgreSQL and SQLite writers translate the ones they understand to SQL.
```dbml
@postgres: search_path myapp
Table myapp.events {
id bigint [pk]
created_at timestamp [not null]
@postgres(id): identity always
@postgres: partition by RANGE (created_at)
@sqlite: without rowid
indexes {
(created_at) [name: 'idx_events_created']
@postgres: with (fillfactor=90)
}
}
```
| Position | Attaches to |
|----------|-------------|
| Before the first `Table {` | database |
| Table body, no `(target)` | that table |
| Table body, `(col)` target | column `col` (error if unknown) |
| Inside `indexes { }` | the most recently listed index entry |
`args` is preserved verbatim; the **key** (lowercased first token) drives
duplicate detection. Repeated directives are kept in order; catalog "singleton"
keys error on a second occurrence at the same location. All errors are
line-numbered.
`ReaderOptions.StrictDirectives` (CLI `--strict-directives`) turns an unknown
namespace or key into an error instead of preserving it silently.
See [`docs/DBML_DIRECTIVES.md`](../../../docs/DBML_DIRECTIVES.md) for the full
grammar and the supported-directive matrix.
## Notes
+137
View File
@@ -0,0 +1,137 @@
package dbml
import (
"fmt"
"regexp"
"strings"
"git.warky.dev/wdevs/relspecgo/pkg/models"
)
// directiveLineRegex matches a dialect directive line:
//
// @postgres: partition by RANGE (created_at)
// @postgres(id): identity always
//
// Group 1 is the namespace, group 2 the optional (column) target, group 3 the
// raw argument text (validated separately so error messages can be specific).
var directiveLineRegex = regexp.MustCompile(`^@([^():]*)(?:\(([^()]*)\))?\s*:(.*)$`)
// namespaceRegex is the grammar for a directive namespace.
var namespaceRegex = regexp.MustCompile(`^[a-z][a-z0-9_]*$`)
// parsedDirective is a directive line that has been parsed but not yet attached
// to a model object.
type parsedDirective struct {
namespace string
target string // column name; "" when absent
args string
line int
}
// parseDirectiveLine parses a single "@namespace[(target)]: args" line.
func parseDirectiveLine(line string, lineNo int) (parsedDirective, error) {
m := directiveLineRegex.FindStringSubmatch(line)
if m == nil {
return parsedDirective{}, fmt.Errorf(
"dbml: line %d: malformed directive %q (expected \"@namespace: args\")", lineNo, line)
}
ns := strings.TrimSpace(m[1])
target := strings.TrimSpace(m[2])
args := strings.TrimSpace(m[3])
if !namespaceRegex.MatchString(ns) {
return parsedDirective{}, fmt.Errorf(
"dbml: line %d: invalid directive namespace %q (must match [a-z][a-z0-9_]*)", lineNo, ns)
}
if args == "" {
return parsedDirective{}, fmt.Errorf("dbml: line %d: directive @%s has no arguments", lineNo, ns)
}
if target != "" {
target = stripQuotes(target)
}
return parsedDirective{namespace: ns, target: target, args: args, line: lineNo}, nil
}
// attachDirective resolves the target model object from the current parser state
// and stores the directive in its Metadata, enforcing location, duplicate and
// strict-mode rules.
func (r *Reader) attachDirective(
pd parsedDirective,
db *models.Database,
table *models.Table,
inTable, inIndexes bool,
lastIndex *models.Index,
) error {
strict := r.options != nil && r.options.StrictDirectives
key := models.DirectiveKey(pd.args)
var meta map[string]any
var location string
switch {
case inIndexes:
if pd.target != "" {
return fmt.Errorf("dbml: line %d: directive target (%s) is not allowed inside an indexes block", pd.line, pd.target)
}
if lastIndex == nil {
return fmt.Errorf("dbml: line %d: directive @%s must follow an index definition", pd.line, pd.namespace)
}
if lastIndex.Metadata == nil {
lastIndex.Metadata = make(map[string]any)
}
meta = lastIndex.Metadata
location = models.DirectiveLocationIndex
case inTable && table != nil:
if pd.target != "" {
col, ok := table.Columns[pd.target]
if !ok {
return fmt.Errorf("dbml: line %d: directive target column %q not found in table %q", pd.line, pd.target, table.Name)
}
if col.Metadata == nil {
col.Metadata = make(map[string]any)
}
meta = col.Metadata
location = models.DirectiveLocationColumn
} else {
if table.Metadata == nil {
table.Metadata = make(map[string]any)
}
meta = table.Metadata
location = models.DirectiveLocationTable
}
default:
if pd.target != "" {
return fmt.Errorf("dbml: line %d: directive target (%s) is only valid inside a table", pd.line, pd.target)
}
if db.Metadata == nil {
db.Metadata = make(map[string]any)
}
meta = db.Metadata
location = models.DirectiveLocationDatabase
}
spec, documented := models.LookupDirectiveSpec(pd.namespace, key)
if strict && !documented {
return fmt.Errorf("dbml: line %d: unknown directive @%s: %s (strict mode)", pd.line, pd.namespace, key)
}
if documented && !models.DirectiveLocationAllowed(pd.namespace, key, location) {
return fmt.Errorf("dbml: line %d: directive @%s: %s is not valid at %s level", pd.line, pd.namespace, key, location)
}
if documented && spec.Singleton && models.HasDirective(meta, pd.namespace, key) {
return fmt.Errorf("dbml: line %d: duplicate @%s directive %q at %s level", pd.line, pd.namespace, key, location)
}
models.AddDirective(meta, models.Directive{
Namespace: pd.namespace,
Key: key,
Args: pd.args,
Line: pd.line,
})
return nil
}
+181
View File
@@ -0,0 +1,181 @@
package dbml
import (
"strings"
"testing"
"git.warky.dev/wdevs/relspecgo/pkg/models"
"git.warky.dev/wdevs/relspecgo/pkg/readers"
)
func parse(t *testing.T, strict bool, src string) (*models.Database, error) {
t.Helper()
r := NewReader(&readers.ReaderOptions{StrictDirectives: strict})
return r.parseDBML(src)
}
func firstTable(t *testing.T, db *models.Database) *models.Table {
t.Helper()
if len(db.Schemas) == 0 || len(db.Schemas[0].Tables) == 0 {
t.Fatal("no table parsed")
}
return db.Schemas[0].Tables[0]
}
func TestDirectives_AttachAtEachLocation(t *testing.T) {
src := `@postgres: search_path myapp
Table myapp.events {
id bigint [pk]
created_at timestamp [not null]
@postgres(id): identity always
@postgres: partition by RANGE (created_at)
indexes {
(created_at) [name: 'idx_events_created']
@postgres: with (fillfactor=90)
}
}
`
db, err := parse(t, false, src)
if err != nil {
t.Fatalf("parse: %v", err)
}
if !models.HasDirective(db.Metadata, "postgres", "search_path") {
t.Errorf("database-level directive missing: %+v", db.Metadata)
}
tbl := firstTable(t, db)
if !models.HasDirective(tbl.Metadata, "postgres", "partition") {
t.Errorf("table-level directive missing: %+v", tbl.Metadata)
}
col := tbl.Columns["id"]
if col == nil || !models.HasDirective(col.Metadata, "postgres", "identity") {
t.Errorf("column-level directive missing")
}
// Verbatim args preserved.
if d := models.DirectivesForNamespace(col.Metadata, "postgres"); len(d) != 1 || d[0].Args != "identity always" {
t.Errorf("column directive args = %+v", d)
}
var idx *models.Index
for _, i := range tbl.Indexes {
idx = i
}
if idx == nil || !models.HasDirective(idx.Metadata, "postgres", "with") {
t.Errorf("index-level directive missing: %+v", idx)
}
}
func TestDirectives_RepeatablePreservedAndOrdered(t *testing.T) {
src := `Table s.t {
id int [pk]
@postgres: with (fillfactor=90)
@postgres: with (autovacuum_enabled=off)
}
`
db, err := parse(t, false, src)
if err != nil {
t.Fatalf("parse: %v", err)
}
tbl := firstTable(t, db)
got := models.DirectivesForNamespace(tbl.Metadata, "postgres")
if len(got) != 2 {
t.Fatalf("got %d directives, want 2", len(got))
}
if got[0].Args != "with (fillfactor=90)" || got[1].Args != "with (autovacuum_enabled=off)" {
t.Errorf("repeatable directives out of order: %+v", got)
}
}
func TestDirectives_SingletonDuplicateErrors(t *testing.T) {
src := `Table s.t {
id int [pk]
@postgres: partition by RANGE (a)
@postgres: partition by LIST (b)
}
`
_, err := parse(t, false, src)
if err == nil || !strings.Contains(err.Error(), "duplicate") {
t.Fatalf("want duplicate error, got %v", err)
}
if !strings.Contains(err.Error(), "line 4") {
t.Errorf("error not line-numbered: %v", err)
}
}
func TestDirectives_MalformedErrors(t *testing.T) {
cases := map[string]string{
"no colon": "@postgres partition by x",
"empty args": "@postgres:",
"bad namespace": "@Postgres: partition by x",
"numeric prefix": "@1x: foo",
}
for name, line := range cases {
t.Run(name, func(t *testing.T) {
src := "Table s.t {\n id int [pk]\n " + line + "\n}\n"
_, err := parse(t, false, src)
if err == nil {
t.Fatalf("want error for %q", line)
}
if !strings.Contains(err.Error(), "line 3") {
t.Errorf("error not line-numbered: %v", err)
}
})
}
}
func TestDirectives_UnknownPreservedNonStrict(t *testing.T) {
src := `Table s.t {
id int [pk]
@postgres: frobnicate all the things
@clickhouse: engine MergeTree
}
`
db, err := parse(t, false, src)
if err != nil {
t.Fatalf("parse: %v", err)
}
tbl := firstTable(t, db)
if !models.HasDirective(tbl.Metadata, "postgres", "frobnicate") {
t.Error("unknown postgres key not preserved")
}
if !models.HasDirective(tbl.Metadata, "clickhouse", "engine") {
t.Error("unknown namespace not preserved")
}
}
func TestDirectives_StrictErrors(t *testing.T) {
src := `Table s.t {
id int [pk]
@postgres: frobnicate x
}
`
_, err := parse(t, true, src)
if err == nil || !strings.Contains(err.Error(), "strict mode") {
t.Fatalf("want strict-mode error, got %v", err)
}
}
func TestDirectives_UnknownColumnTargetErrors(t *testing.T) {
src := `Table s.t {
id int [pk]
@postgres(missing): identity always
}
`
_, err := parse(t, false, src)
if err == nil || !strings.Contains(err.Error(), "not found") {
t.Fatalf("want unknown-column error, got %v", err)
}
}
func TestDirectives_WrongLocationErrors(t *testing.T) {
// partition is table-only.
src := "@postgres: partition by RANGE (x)\n\nTable s.t {\n id int [pk]\n}\n"
_, err := parse(t, false, src)
if err == nil || !strings.Contains(err.Error(), "not valid at database level") {
t.Fatalf("want location error, got %v", err)
}
}
+24 -2
View File
@@ -435,11 +435,14 @@ func (r *Reader) parseDBML(content string) (*models.Database, error) {
var inIndexes bool
var inTable bool
var columnSeq uint
var lastIndex *models.Index // most recent index in the current Indexes block
lineNo := 0
tableRegex := regexp.MustCompile(`^Table\s+(.+?)\s*{`)
refRegex := regexp.MustCompile(`^Ref:\s+(.+)`)
for scanner.Scan() {
lineNo++
line := strings.TrimSpace(scanner.Text())
// Skip empty lines and comments
@@ -447,6 +450,20 @@ func (r *Reader) parseDBML(content string) (*models.Database, error) {
continue
}
// Parse a dialect directive (@postgres:, @sqlite:, …). Handled before
// table/column/index parsing so directive lines are never mistaken for
// columns.
if strings.HasPrefix(line, "@") {
pd, err := parseDirectiveLine(line, lineNo)
if err != nil {
return nil, err
}
if err := r.attachDirective(pd, db, currentTable, inTable, inIndexes, lastIndex); err != nil {
return nil, err
}
continue
}
// Parse Table definition
if matches := tableRegex.FindStringSubmatch(line); matches != nil {
tableName := matches[1]
@@ -474,8 +491,10 @@ func (r *Reader) parseDBML(content string) (*models.Database, error) {
continue
}
// End of table definition
if inTable && line == "}" {
// End of table definition. Guarded by !inIndexes so the closing brace
// of an `indexes { }` block is not mistaken for the end of the table
// (which would drop any table-level content that follows it).
if inTable && !inIndexes && line == "}" {
if currentTable != nil && currentSchema != "" {
schemaMap[currentSchema].Tables = append(schemaMap[currentSchema].Tables, currentTable)
currentTable = nil
@@ -488,12 +507,14 @@ func (r *Reader) parseDBML(content string) (*models.Database, error) {
// Parse indexes section
if inTable && (strings.HasPrefix(line, "Indexes {") || strings.HasPrefix(line, "indexes {")) {
inIndexes = true
lastIndex = nil
continue
}
// End of indexes section
if inIndexes && line == "}" {
inIndexes = false
lastIndex = nil
continue
}
@@ -513,6 +534,7 @@ func (r *Reader) parseDBML(content string) (*models.Database, error) {
index := r.parseIndex(line, currentTable.Name, currentSchema)
if index != nil {
currentTable.Indexes[index.Name] = index
lastIndex = index
}
continue
}
+4
View File
@@ -28,6 +28,10 @@ type ReaderOptions struct {
// Prisma7 enables Prisma 7-specific handling for Prisma schemas.
Prisma7 bool
// StrictDirectives makes DBML dialect directives (@postgres:, @sqlite:, …)
// fail on an unknown namespace or key instead of preserving them silently.
StrictDirectives bool
// Additional options can be added here as needed
Metadata map[string]interface{}
}
+35
View File
@@ -137,6 +137,41 @@ indexes {
}
```
### Dialect directives
Dialect directives stored on a model object's `Metadata` (namespace `postgres`,
`sqlite`, …) are re-emitted verbatim, one line per directive, at the location
they belong to:
```dbml
@postgres: search_path myapp
Table myapp.events {
id bigint [pk]
created_at timestamp [not null]
@postgres(id): identity always
@postgres: partition by RANGE (created_at)
@sqlite: without rowid
indexes {
(created_at) [name: 'idx_events_created']
@postgres: with (fillfactor=90)
}
}
```
| Emitted at | From |
|------------|------|
| Before the first table | `Database.Metadata` |
| After a column line, as `@ns(col): …` | `Column.Metadata` |
| After an index line, inside `indexes { }` | `Index.Metadata` |
| After the `indexes` block, before `Note:` | `Table.Metadata` |
Output is deterministic (ordered by namespace, then source line, then args), so a
`DBML → model → DBML` round-trip is idempotent. See
[`docs/DBML_DIRECTIVES.md`](../../../docs/DBML_DIRECTIVES.md) for the grammar and
the list of directives the PostgreSQL and SQLite writers translate to SQL.
## Type Mapping
| SQL Type | DBML Type |
+23
View File
@@ -0,0 +1,23 @@
package dbml
import (
"git.warky.dev/wdevs/relspecgo/pkg/models"
)
// directiveLines renders every dialect directive stored in meta back to its DBML
// source form, one line per directive, each prefixed with indent. When target is
// non-empty it is emitted as the "(column)" target, e.g.
// " @postgres(id): identity always". Order is deterministic (see
// models.GetDirectives).
func directiveLines(meta map[string]any, indent, target string) []string {
directives := models.GetDirectives(meta)
if len(directives) == 0 {
return nil
}
lines := make([]string, 0, len(directives))
for _, d := range directives {
lines = append(lines, indent+models.FormatDirectiveLine(d, target))
}
return lines
}
+97
View File
@@ -0,0 +1,97 @@
package dbml
import (
"os"
"path/filepath"
"testing"
dbmlreader "git.warky.dev/wdevs/relspecgo/pkg/readers/dbml"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.warky.dev/wdevs/relspecgo/pkg/models"
"git.warky.dev/wdevs/relspecgo/pkg/readers"
"git.warky.dev/wdevs/relspecgo/pkg/writers"
)
const directiveSrc = `@postgres: search_path myapp
Table myapp.events {
id bigint [pk]
created_at timestamp [not null]
@postgres(id): identity always
@postgres: partition by RANGE (created_at)
@postgres: tablespace fast_data
@sqlite: without rowid
indexes {
(created_at) [name: 'idx_events_created']
@postgres: with (fillfactor=90)
}
}
`
func writeDBML(t *testing.T, db *models.Database) string {
t.Helper()
out := filepath.Join(t.TempDir(), "out.dbml")
require.NoError(t, NewWriter(&writers.WriterOptions{OutputPath: out}).WriteDatabase(db))
b, err := os.ReadFile(out)
require.NoError(t, err)
return string(b)
}
func readDBML(t *testing.T, src string) *models.Database {
t.Helper()
f := filepath.Join(t.TempDir(), "in.dbml")
require.NoError(t, os.WriteFile(f, []byte(src), 0o644))
db, err := dbmlreader.NewReader(&readers.ReaderOptions{FilePath: f}).ReadDatabase()
require.NoError(t, err)
return db
}
func collectDirectives(db *models.Database) map[string][]string {
got := map[string][]string{}
add := func(loc string, meta map[string]any) {
for _, d := range models.GetDirectives(meta) {
got[loc] = append(got[loc], models.FormatDirectiveLine(d, ""))
}
}
add("database", db.Metadata)
for _, s := range db.Schemas {
for _, tbl := range s.Tables {
add("table:"+tbl.Name, tbl.Metadata)
for _, c := range tbl.Columns {
add("column:"+c.Name, c.Metadata)
}
for _, i := range tbl.Indexes {
add("index:"+i.Name, i.Metadata)
}
}
}
return got
}
func TestDirectives_RoundTrip(t *testing.T) {
db1 := readDBML(t, directiveSrc)
out1 := writeDBML(t, db1)
db2 := readDBML(t, out1)
out2 := writeDBML(t, db2)
assert.Equal(t, out1, out2, "DBML directive output should be idempotent")
assert.Equal(t, collectDirectives(db1), collectDirectives(db2), "directives preserved through round-trip")
// Spot-check each location survived.
d := collectDirectives(db2)
assert.Contains(t, d["database"], "@postgres: search_path myapp")
assert.Contains(t, d["table:events"], "@postgres: partition by RANGE (created_at)")
assert.Contains(t, d["table:events"], "@sqlite: without rowid")
assert.Contains(t, d["column:id"], "@postgres: identity always")
assert.Contains(t, d["index:idx_events_created"], "@postgres: with (fillfactor=90)")
}
func TestDirectives_WriterEmitsColumnTarget(t *testing.T) {
db := readDBML(t, directiveSrc)
out := writeDBML(t, db)
assert.Contains(t, out, "@postgres(id): identity always")
}
+23
View File
@@ -72,6 +72,14 @@ func (w *Writer) databaseToDBML(d *models.Database) string {
sb.WriteString("\n")
}
if dirLines := directiveLines(d.Metadata, "", ""); len(dirLines) > 0 {
for _, line := range dirLines {
sb.WriteString(line)
sb.WriteString("\n")
}
sb.WriteString("\n")
}
for _, schema := range d.Schemas {
sb.WriteString(w.schemaToDBML(schema))
}
@@ -146,6 +154,11 @@ func (w *Writer) tableToDBML(t *models.Table) string {
fmt.Fprintf(&sb, " // %s", column.Comment)
}
sb.WriteString("\n")
for _, line := range directiveLines(column.Metadata, " ", column.Name) {
sb.WriteString(line)
sb.WriteString("\n")
}
}
if len(t.Indexes) > 0 {
@@ -167,10 +180,20 @@ func (w *Writer) tableToDBML(t *models.Table) string {
fmt.Fprintf(&sb, " [%s]", strings.Join(indexAttrs, ", "))
}
sb.WriteString("\n")
for _, line := range directiveLines(index.Metadata, " ", "") {
sb.WriteString(line)
sb.WriteString("\n")
}
}
sb.WriteString(" }\n")
}
for _, line := range directiveLines(t.Metadata, " ", "") {
sb.WriteString(line)
sb.WriteString("\n")
}
note := strings.TrimSpace(t.Description + " " + t.Comment)
if note != "" {
fmt.Fprintf(&sb, "\n Note: '%s'\n", note)
+20
View File
@@ -172,6 +172,26 @@ When `include_audit` is enabled, adds:
- Concurrent index creation (`CREATE INDEX CONCURRENTLY`) via `Index.Concurrent`
- Check constraints with expressions
- Extension types and indexes: PostGIS, pgvector, citext, hstore, ltree (see below)
- DBML dialect directives (`@postgres:` — see below)
### DBML dialect directives
`@postgres:` directives carried on a model object's `Metadata` (typically from a
DBML source file) are translated to SQL:
| Directive | Location | Emitted |
|-----------|----------|---------|
| `@postgres: partition by …` | table | `PARTITION BY …` on `CREATE TABLE` |
| `@postgres: inherits …` | table | `INHERITS (…)` |
| `@postgres: with (…)` | table, index | `WITH (…)` (on an index, overrides the comment-derived `WITH`) |
| `@postgres: tablespace …` | table, index | `TABLESPACE …` |
| `@postgres(col): storage …` | column | `STORAGE …` |
| `@postgres(col): compression …` | column | `COMPRESSION …` |
| `@postgres(col): identity always` / `identity by default` | column | `GENERATED ALWAYS/BY DEFAULT AS IDENTITY` |
Directives for other dialects (`@sqlite:` …) are ignored. With
`WriterOptions.StrictDirectives` (CLI `--strict-directives`) an untranslatable
`@postgres:` key is an error. Full reference: [`docs/DBML_DIRECTIVES.md`](../../../docs/DBML_DIRECTIVES.md).
## Data Types
+184
View File
@@ -0,0 +1,184 @@
package pgsql
import (
"fmt"
"strings"
"git.warky.dev/wdevs/relspecgo/pkg/models"
)
// directiveNamespace is the dialect namespace this writer consumes. Directives
// for other namespaces (e.g. "sqlite") are ignored and never emitted as SQL.
const directiveNamespace = "postgres"
// pgHandledDirectives maps a directive location to the set of postgres keys this
// writer knows how to translate. In strict mode an unknown key for this
// namespace at a supported location is a hard error.
var pgHandledDirectives = map[string]map[string]bool{
models.DirectiveLocationTable: {"partition": true, "inherits": true, "with": true, "tablespace": true},
models.DirectiveLocationColumn: {"storage": true, "compression": true, "identity": true},
models.DirectiveLocationIndex: {"with": true, "tablespace": true},
}
// checkDirectives validates postgres directives across a schema when strict mode
// is enabled. It returns an error for any postgres directive whose key this
// writer cannot translate. With strict mode off it is a no-op.
func (w *Writer) checkDirectives(schema *models.Schema) error {
if w.options == nil || !w.options.StrictDirectives {
return nil
}
for _, table := range schema.Tables {
if err := checkObjectDirectives(table.Metadata, models.DirectiveLocationTable, table.Name); err != nil {
return err
}
for _, col := range table.Columns {
if err := checkObjectDirectives(col.Metadata, models.DirectiveLocationColumn, table.Name+"."+col.Name); err != nil {
return err
}
}
for _, idx := range table.Indexes {
if err := checkObjectDirectives(idx.Metadata, models.DirectiveLocationIndex, idx.Name); err != nil {
return err
}
}
}
return nil
}
func checkObjectDirectives(meta map[string]any, location, owner string) error {
for _, d := range models.DirectivesForNamespace(meta, directiveNamespace) {
if !pgHandledDirectives[location][d.Key] {
return fmt.Errorf("pgsql: %s: unsupported @postgres directive %q at %s level (strict mode)", owner, d.Key, location)
}
}
return nil
}
// upperLeadingClause upcases a known leading keyword phrase in a directive
// argument so the emitted SQL reads conventionally. Identifiers that follow are
// left untouched.
func upperLeadingClause(args, lowerPrefix, upperPrefix string) string {
args = strings.TrimSpace(args)
if strings.HasPrefix(strings.ToLower(args), lowerPrefix) {
return upperPrefix + args[len(lowerPrefix):]
}
return args
}
// pgTableDirectiveSuffix returns the clause appended after the closing ")" of a
// CREATE TABLE statement, e.g. " PARTITION BY RANGE (created_at) TABLESPACE fast".
func pgTableDirectiveSuffix(table *models.Table) string {
directives := models.DirectivesForNamespace(table.Metadata, directiveNamespace)
if len(directives) == 0 {
return ""
}
byKey := firstByKey(directives)
var parts []string
if d, ok := byKey["partition"]; ok {
parts = append(parts, upperLeadingClause(d.Args, "partition by", "PARTITION BY"))
}
if d, ok := byKey["inherits"]; ok {
parts = append(parts, upperLeadingClause(d.Args, "inherits", "INHERITS"))
}
if d, ok := byKey["with"]; ok {
parts = append(parts, upperLeadingClause(d.Args, "with", "WITH"))
}
if d, ok := byKey["tablespace"]; ok {
parts = append(parts, upperLeadingClause(d.Args, "tablespace", "TABLESPACE"))
}
if len(parts) == 0 {
return ""
}
return " " + strings.Join(parts, " ")
}
// pgColumnDirectiveSuffix returns the clause appended to a column definition,
// e.g. " STORAGE PLAIN" or " GENERATED ALWAYS AS IDENTITY".
func pgColumnDirectiveSuffix(col *models.Column) string {
directives := models.DirectivesForNamespace(col.Metadata, directiveNamespace)
if len(directives) == 0 {
return ""
}
byKey := firstByKey(directives)
var parts []string
if d, ok := byKey["storage"]; ok {
parts = append(parts, upperLeadingClause(d.Args, "storage", "STORAGE"))
}
if d, ok := byKey["compression"]; ok {
parts = append(parts, upperLeadingClause(d.Args, "compression", "COMPRESSION"))
}
if d, ok := byKey["identity"]; ok {
parts = append(parts, identityClause(d.Args))
}
if len(parts) == 0 {
return ""
}
return " " + strings.Join(parts, " ")
}
// identityClause maps the two documented identity forms to standard SQL,
// falling back to a verbatim (upcased-keyword) rendering.
func identityClause(args string) string {
switch strings.ToLower(strings.Join(strings.Fields(args), " ")) {
case "identity always":
return "GENERATED ALWAYS AS IDENTITY"
case "identity default", "identity by default":
return "GENERATED BY DEFAULT AS IDENTITY"
default:
return upperLeadingClause(args, "identity", "IDENTITY")
}
}
// pgIndexDirectiveWith returns the parenthesised storage-parameter list from an
// @postgres: with (...) index directive, e.g. "fillfactor=90", or "".
func pgIndexDirectiveWith(index *models.Index) string {
for _, d := range models.DirectivesForNamespace(index.Metadata, directiveNamespace) {
if d.Key != "with" {
continue
}
inner := d.Args
if i := strings.Index(inner, "("); i >= 0 {
if j := strings.LastIndex(inner, ")"); j > i {
return strings.TrimSpace(inner[i+1 : j])
}
}
return strings.TrimSpace(strings.TrimPrefix(strings.ToLower(inner), "with"))
}
return ""
}
// pgIndexWithParams returns the storage-parameter list to use for an index,
// preferring an @postgres: with (...) directive over the given fallback (e.g.
// one derived from the index comment).
func pgIndexWithParams(index *models.Index, fallback string) string {
if p := pgIndexDirectiveWith(index); p != "" {
return p
}
return fallback
}
// pgIndexDirectiveTablespace returns the tablespace name from an
// @postgres: tablespace <name> index directive, or "".
func pgIndexDirectiveTablespace(index *models.Index) string {
for _, d := range models.DirectivesForNamespace(index.Metadata, directiveNamespace) {
if d.Key == "tablespace" {
return strings.TrimSpace(strings.TrimPrefix(strings.ToLower(d.Args), "tablespace"))
}
}
return ""
}
// firstByKey indexes directives by key, keeping the first occurrence (the
// documented postgres keys used here are all singletons).
func firstByKey(directives []models.Directive) map[string]models.Directive {
byKey := make(map[string]models.Directive, len(directives))
for _, d := range directives {
if _, exists := byKey[d.Key]; !exists {
byKey[d.Key] = d
}
}
return byKey
}
+107
View File
@@ -0,0 +1,107 @@
package pgsql
import (
"bytes"
"strings"
"testing"
"git.warky.dev/wdevs/relspecgo/pkg/models"
"git.warky.dev/wdevs/relspecgo/pkg/writers"
)
func directiveTestDB(t *testing.T) *models.Database {
t.Helper()
db := models.InitDatabase("testdb")
schema := models.InitSchema("public")
table := models.InitTable("events", "public")
id := models.InitColumn("id", "events", "public")
id.Type = "bigint"
id.IsPrimaryKey = true
id.NotNull = true
models.AddDirective(id.Metadata, models.Directive{Namespace: "postgres", Args: "identity always"})
models.AddDirective(id.Metadata, models.Directive{Namespace: "sqlite", Args: "collate NOCASE"})
table.Columns["id"] = id
created := models.InitColumn("created_at", "events", "public")
created.Type = "timestamp"
created.NotNull = true
table.Columns["created_at"] = created
models.AddDirective(table.Metadata, models.Directive{Namespace: "postgres", Args: "partition by RANGE (created_at)"})
models.AddDirective(table.Metadata, models.Directive{Namespace: "postgres", Args: "tablespace fast_data"})
models.AddDirective(table.Metadata, models.Directive{Namespace: "sqlite", Args: "without rowid"})
idx := models.InitIndex("idx_events_created", "events", "public")
idx.Columns = []string{"created_at"}
models.AddDirective(idx.Metadata, models.Directive{Namespace: "postgres", Args: "with (fillfactor=90)"})
models.AddDirective(idx.Metadata, models.Directive{Namespace: "postgres", Args: "tablespace idx_space"})
table.Indexes["idx_events_created"] = idx
schema.Tables = append(schema.Tables, table)
db.Schemas = append(db.Schemas, schema)
return db
}
func TestPgDirectives_WriteDatabasePath(t *testing.T) {
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{})
w.writer = &buf
if err := w.WriteDatabase(directiveTestDB(t)); err != nil {
t.Fatalf("WriteDatabase: %v", err)
}
out := buf.String()
for _, want := range []string{
") PARTITION BY RANGE (created_at) TABLESPACE fast_data",
"GENERATED ALWAYS AS IDENTITY",
"WITH (fillfactor=90) TABLESPACE idx_space",
} {
if !strings.Contains(out, want) {
t.Errorf("missing %q in:\n%s", want, out)
}
}
// sqlite directives must never reach PG output.
if strings.Contains(out, "WITHOUT ROWID") || strings.Contains(strings.ToUpper(out), "COLLATE NOCASE") {
t.Errorf("sqlite directive leaked into PG output:\n%s", out)
}
}
func TestPgDirectives_WriteSchemaPath(t *testing.T) {
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{})
w.writer = &buf
if err := w.WriteSchema(directiveTestDB(t).Schemas[0]); err != nil {
t.Fatalf("WriteSchema: %v", err)
}
out := buf.String()
if !strings.Contains(out, ") PARTITION BY RANGE (created_at) TABLESPACE fast_data;") {
t.Errorf("table suffix missing from WriteSchema path:\n%s", out)
}
if !strings.Contains(out, "WITH (fillfactor=90) TABLESPACE idx_space") {
t.Errorf("index clauses missing from WriteSchema path:\n%s", out)
}
}
func TestPgDirectives_StrictUnknownKeyErrors(t *testing.T) {
db := directiveTestDB(t)
models.AddDirective(db.Schemas[0].Tables[0].Metadata, models.Directive{Namespace: "postgres", Args: "frobnicate x"})
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{StrictDirectives: true})
w.writer = &buf
err := w.WriteSchema(db.Schemas[0])
if err == nil || !strings.Contains(err.Error(), "frobnicate") {
t.Fatalf("want strict error for unknown postgres key, got %v", err)
}
}
func TestPgDirectives_StrictIgnoresOtherNamespaces(t *testing.T) {
// sqlite directives are present but must not trip PG strict mode.
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{StrictDirectives: true})
w.writer = &buf
if err := w.WriteSchema(directiveTestDB(t).Schemas[0]); err != nil {
t.Fatalf("strict mode should ignore non-postgres directives, got %v", err)
}
}
+29 -10
View File
@@ -143,6 +143,10 @@ func (w *Writer) GenerateDatabaseStatements(db *models.Database) ([]string, erro
func (w *Writer) GenerateSchemaStatements(schema *models.Schema) ([]string, error) {
statements := []string{}
if err := w.checkDirectives(schema); err != nil {
return nil, err
}
// Phase 1: Create schema (skip entirely when flattening)
if schema.Name != "public" && !w.options.FlattenSchema {
statements = append(statements, fmt.Sprintf("-- Schema: %s", schema.Name))
@@ -277,17 +281,22 @@ func (w *Writer) GenerateSchemaStatements(schema *models.Schema) ([]string, erro
columnExprs := buildIndexColumnExpressions(table, index, indexType)
withClause := ""
if params := indexStorageParameters(index.Comment); params != "" {
if params := pgIndexWithParams(index, indexStorageParameters(index.Comment)); params != "" {
withClause = fmt.Sprintf(" WITH (%s)", params)
}
tablespaceClause := ""
if ts := pgIndexDirectiveTablespace(index); ts != "" {
tablespaceClause = fmt.Sprintf(" TABLESPACE %s", ts)
}
whereClause := ""
if index.Where != "" {
whereClause = fmt.Sprintf(" WHERE %s", index.Where)
}
stmt := fmt.Sprintf("CREATE %sINDEX IF NOT EXISTS %s ON %s USING %s (%s)%s%s",
uniqueStr, quoteIdentifier(index.Name), w.qualTable(schema.SQLName(), table.SQLName()), indexType, strings.Join(columnExprs, ", "), withClause, whereClause)
stmt := fmt.Sprintf("CREATE %sINDEX IF NOT EXISTS %s ON %s USING %s (%s)%s%s%s",
uniqueStr, quoteIdentifier(index.Name), w.qualTable(schema.SQLName(), table.SQLName()), indexType, strings.Join(columnExprs, ", "), withClause, tablespaceClause, whereClause)
statements = append(statements, stmt)
}
}
@@ -581,8 +590,9 @@ func (w *Writer) generateCreateTableStatement(schema *models.Schema, table *mode
columnDefs = append(columnDefs, " "+def)
}
stmt := fmt.Sprintf("CREATE TABLE IF NOT EXISTS %s (\n%s\n)",
w.qualTable(schema.SQLName(), table.SQLName()), strings.Join(columnDefs, ",\n"))
stmt := fmt.Sprintf("CREATE TABLE IF NOT EXISTS %s (\n%s\n)%s",
w.qualTable(schema.SQLName(), table.SQLName()), strings.Join(columnDefs, ",\n"),
pgTableDirectiveSuffix(table))
statements = append(statements, stmt)
return statements, nil
@@ -611,7 +621,7 @@ func (w *Writer) generateColumnDefinition(col *models.Column) string {
}
}
return strings.Join(parts, " ")
return strings.Join(parts, " ") + pgColumnDirectiveSuffix(col)
}
func effectiveColumnSQLType(col *models.Column) string {
@@ -678,6 +688,10 @@ func (w *Writer) WriteSchema(schema *models.Schema) error {
w.writer = os.Stdout
}
if err := w.checkDirectives(schema); err != nil {
return err
}
// Phase 1: Create schema (priority 1)
if err := w.writeCreateSchema(schema); err != nil {
return err
@@ -884,7 +898,7 @@ func (w *Writer) writeCreateTables(schema *models.Schema) error {
}
fmt.Fprintf(w.writer, "%s\n", strings.Join(columnDefs, ",\n"))
fmt.Fprintf(w.writer, ");\n\n")
fmt.Fprintf(w.writer, ")%s;\n\n", pgTableDirectiveSuffix(table))
}
return nil
@@ -1079,10 +1093,15 @@ func (w *Writer) writeIndexes(schema *models.Schema) error {
}
withClause := ""
if params := indexStorageParameters(index.Comment); params != "" {
if params := pgIndexWithParams(index, indexStorageParameters(index.Comment)); params != "" {
withClause = fmt.Sprintf(" WITH (%s)", params)
}
tablespaceClause := ""
if ts := pgIndexDirectiveTablespace(index); ts != "" {
tablespaceClause = fmt.Sprintf(" TABLESPACE %s", ts)
}
whereClause := ""
if index.Where != "" {
whereClause = fmt.Sprintf(" WHERE %s", index.Where)
@@ -1095,8 +1114,8 @@ func (w *Writer) writeIndexes(schema *models.Schema) error {
fmt.Fprintf(w.writer, "CREATE %sINDEX %sIF NOT EXISTS %s\n",
unique, concurrently, indexName)
fmt.Fprintf(w.writer, " ON %s USING %s (%s)%s%s;\n\n",
w.qualTable(schema.SQLName(), table.SQLName()), indexType, strings.Join(columnExprs, ", "), withClause, whereClause)
fmt.Fprintf(w.writer, " ON %s USING %s (%s)%s%s%s;\n\n",
w.qualTable(schema.SQLName(), table.SQLName()), indexType, strings.Join(columnExprs, ", "), withClause, tablespaceClause, whereClause)
}
}
+15
View File
@@ -118,6 +118,21 @@ CREATE TABLE "posts" (
- **Check Constraints**: Generated as comments (should be added to CREATE TABLE manually)
- **Indexes**: Generated without PostgreSQL-specific features (no GIN, GiST, operator classes)
## DBML dialect directives
`@sqlite:` directives carried on a model object's `Metadata` (typically from a
DBML source file) are translated to SQL:
| Directive | Location | Emitted |
|-----------|----------|---------|
| `@sqlite: without rowid` | table | `WITHOUT ROWID` table option |
| `@sqlite: strict` | table | `STRICT` table option (after `WITHOUT ROWID`) |
| `@sqlite(col): collate …` | column | ` COLLATE …` in the column definition |
Directives for other dialects (`@postgres:` …) are ignored. With
`WriterOptions.StrictDirectives` (CLI `--strict-directives`) an untranslatable
`@sqlite:` key is an error. Full reference: [`docs/DBML_DIRECTIVES.md`](../../../docs/DBML_DIRECTIVES.md).
## Output Structure
Generated SQL follows this order:
+84
View File
@@ -0,0 +1,84 @@
package sqlite
import (
"fmt"
"strings"
"git.warky.dev/wdevs/relspecgo/pkg/models"
)
// directiveNamespace is the dialect namespace this writer consumes. Directives
// for other namespaces (e.g. "postgres") are ignored and never emitted as SQL.
const directiveNamespace = "sqlite"
// sqliteHandledDirectives maps a directive location to the set of sqlite keys
// this writer knows how to translate. In strict mode an unknown key for this
// namespace at a supported location is a hard error.
var sqliteHandledDirectives = map[string]map[string]bool{
models.DirectiveLocationTable: {"without": true, "strict": true},
models.DirectiveLocationColumn: {"collate": true},
}
// checkDirectives validates sqlite directives across a schema when strict mode
// is enabled. With strict mode off it is a no-op.
func (w *Writer) checkDirectives(schema *models.Schema) error {
if w.options == nil || !w.options.StrictDirectives {
return nil
}
for _, table := range schema.Tables {
if err := checkObjectDirectives(table.Metadata, models.DirectiveLocationTable, table.Name); err != nil {
return err
}
for _, col := range table.Columns {
if err := checkObjectDirectives(col.Metadata, models.DirectiveLocationColumn, table.Name+"."+col.Name); err != nil {
return err
}
}
for _, idx := range table.Indexes {
if err := checkObjectDirectives(idx.Metadata, models.DirectiveLocationIndex, idx.Name); err != nil {
return err
}
}
}
return nil
}
func checkObjectDirectives(meta map[string]any, location, owner string) error {
for _, d := range models.DirectivesForNamespace(meta, directiveNamespace) {
if !sqliteHandledDirectives[location][d.Key] {
return fmt.Errorf("sqlite: %s: unsupported @sqlite directive %q at %s level (strict mode)", owner, d.Key, location)
}
}
return nil
}
// sqliteTableOptions returns the trailing table-option clause for a CREATE TABLE
// statement, e.g. "WITHOUT ROWID, STRICT". WITHOUT ROWID is emitted before
// STRICT, matching SQLite's own grammar ordering.
func sqliteTableOptions(table *models.Table) string {
var opts []string
if models.HasDirective(table.Metadata, directiveNamespace, "without") {
opts = append(opts, "WITHOUT ROWID")
}
if models.HasDirective(table.Metadata, directiveNamespace, "strict") {
opts = append(opts, "STRICT")
}
return strings.Join(opts, ", ")
}
// sqliteColumnCollate returns a " COLLATE <name>" clause for a column carrying an
// @sqlite(col): collate <name> directive, or "".
func sqliteColumnCollate(col *models.Column) string {
for _, d := range models.DirectivesForNamespace(col.Metadata, directiveNamespace) {
if d.Key != "collate" {
continue
}
name := strings.TrimSpace(strings.TrimPrefix(strings.TrimSpace(d.Args), "collate"))
name = strings.TrimSpace(name)
if name == "" {
return ""
}
return " COLLATE " + name
}
return ""
}
+82
View File
@@ -0,0 +1,82 @@
package sqlite
import (
"bytes"
"strings"
"testing"
"git.warky.dev/wdevs/relspecgo/pkg/models"
"git.warky.dev/wdevs/relspecgo/pkg/writers"
)
func sqliteDirectiveDB(t *testing.T) *models.Database {
t.Helper()
db := models.InitDatabase("testdb")
schema := models.InitSchema("public")
table := models.InitTable("events", "public")
id := models.InitColumn("id", "events", "public")
id.Type = "bigint"
id.IsPrimaryKey = true
id.NotNull = true
table.Columns["id"] = id
name := models.InitColumn("name", "events", "public")
name.Type = "varchar(200)"
name.NotNull = true
models.AddDirective(name.Metadata, models.Directive{Namespace: "sqlite", Args: "collate NOCASE"})
// A postgres directive on the same column must be ignored by the sqlite writer.
models.AddDirective(name.Metadata, models.Directive{Namespace: "postgres", Args: "storage plain"})
table.Columns["name"] = name
models.AddDirective(table.Metadata, models.Directive{Namespace: "sqlite", Args: "without rowid"})
models.AddDirective(table.Metadata, models.Directive{Namespace: "sqlite", Args: "strict"})
models.AddDirective(table.Metadata, models.Directive{Namespace: "postgres", Args: "partition by RANGE (id)"})
schema.Tables = append(schema.Tables, table)
db.Schemas = append(db.Schemas, schema)
return db
}
func TestSqliteDirectives_TableOptionsAndCollate(t *testing.T) {
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{})
w.writer = &buf
if err := w.WriteDatabase(sqliteDirectiveDB(t)); err != nil {
t.Fatalf("WriteDatabase: %v", err)
}
out := buf.String()
if !strings.Contains(out, ") WITHOUT ROWID, STRICT;") {
t.Errorf("missing table options clause:\n%s", out)
}
if !strings.Contains(out, `"name" TEXT COLLATE NOCASE NOT NULL`) {
t.Errorf("missing column COLLATE clause:\n%s", out)
}
// postgres directives must never reach sqlite output.
if strings.Contains(strings.ToUpper(out), "PARTITION BY") || strings.Contains(strings.ToUpper(out), "STORAGE PLAIN") {
t.Errorf("postgres directive leaked into sqlite output:\n%s", out)
}
}
func TestSqliteDirectives_StrictUnknownKeyErrors(t *testing.T) {
db := sqliteDirectiveDB(t)
models.AddDirective(db.Schemas[0].Tables[0].Metadata, models.Directive{Namespace: "sqlite", Args: "frobnicate x"})
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{StrictDirectives: true})
w.writer = &buf
err := w.WriteDatabase(db)
if err == nil || !strings.Contains(err.Error(), "frobnicate") {
t.Fatalf("want strict error for unknown sqlite key, got %v", err)
}
}
func TestSqliteDirectives_StrictIgnoresPostgres(t *testing.T) {
var buf bytes.Buffer
w := NewWriter(&writers.WriterOptions{StrictDirectives: true})
w.writer = &buf
if err := w.WriteDatabase(sqliteDirectiveDB(t)); err != nil {
t.Fatalf("strict mode should ignore postgres directives, got %v", err)
}
}
+1
View File
@@ -25,6 +25,7 @@ func GetTemplateFuncs(opts *writers.WriterOptions) template.FuncMap {
"join": strings.Join,
"lower": strings.ToLower,
"upper": strings.ToUpper,
"column_collate": sqliteColumnCollate,
}
}
+2
View File
@@ -45,6 +45,7 @@ type TableTemplateData struct {
Columns []*models.Column
PrimaryKey *models.Constraint
ForeignKeys []ForeignKeyTemplateData
TableOptions string
}
// ForeignKeyTemplateData contains data for an inline FOREIGN KEY clause
@@ -193,6 +194,7 @@ func BuildTableTemplateData(schema string, table *models.Table) TableTemplateDat
Columns: columns,
PrimaryKey: pk,
ForeignKeys: fks,
TableOptions: sqliteTableOptions(table),
}
}
@@ -1,7 +1,7 @@
CREATE TABLE {{quote_ident (qualified_table_name .Schema .Name)}} (
{{- $hasAutoIncrement := false}}
{{- range $i, $col := .Columns}}{{if $i}},{{end}}
{{quote_ident $col.Name}} {{map_type $col.Type}}{{if is_autoincrement $col}}{{$hasAutoIncrement = true}} PRIMARY KEY AUTOINCREMENT{{else}}{{if $col.NotNull}} NOT NULL{{end}}{{if ne (format_default $col) ""}} DEFAULT {{format_default $col}}{{end}}{{end}}
{{quote_ident $col.Name}} {{map_type $col.Type}}{{column_collate $col}}{{if is_autoincrement $col}}{{$hasAutoIncrement = true}} PRIMARY KEY AUTOINCREMENT{{else}}{{if $col.NotNull}} NOT NULL{{end}}{{if ne (format_default $col) ""}} DEFAULT {{format_default $col}}{{end}}{{end}}
{{- end}}
{{- if and .PrimaryKey (not $hasAutoIncrement)}}{{if gt (len .Columns) 0}},{{end}}
PRIMARY KEY ({{range $i, $colName := .PrimaryKey.Columns}}{{if $i}}, {{end}}{{quote_ident $colName}}{{end}})
@@ -9,4 +9,4 @@ CREATE TABLE {{quote_ident (qualified_table_name .Schema .Name)}} (
{{- range .ForeignKeys}},
FOREIGN KEY ({{range $i, $col := .Columns}}{{if $i}}, {{end}}{{quote_ident $col}}{{end}}) REFERENCES {{quote_ident (qualified_table_name .ForeignSchema .ForeignTable)}} ({{range $i, $col := .ForeignColumns}}{{if $i}}, {{end}}{{quote_ident $col}}{{end}}){{if .OnDelete}} ON DELETE {{.OnDelete}}{{end}}{{if .OnUpdate}} ON UPDATE {{.OnUpdate}}{{end}}
{{- end}}
);
){{if .TableOptions}} {{.TableOptions}}{{end}};
+4
View File
@@ -186,6 +186,10 @@ func tableSchemaName(schema string) string {
func (w *Writer) WriteSchema(schema *models.Schema) error {
tableSchema := tableSchemaName(schema.Name)
if err := w.checkDirectives(schema); err != nil {
return err
}
// SQLite doesn't have schemas, so we just write a comment (skip for the
// default schema, since its tables aren't actually being prefixed)
if tableSchema != "" {
+4
View File
@@ -86,6 +86,10 @@ type WriterOptions struct {
// Prisma7 enables Prisma 7-specific output for Prisma writers.
Prisma7 bool
// StrictDirectives makes dialect directive translation fail on an
// unsupported key for the writer's own namespace instead of skipping it.
StrictDirectives bool
// ContinueOnError instructs SQL writers to prepend `\set ON_ERROR_STOP off`
// to their output so that psql continues past errors instead of stopping.
ContinueOnError bool