Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,15 @@

## Unreleased

- New `@maple-dev/effect-orm/tinybird` entry: `defineDatasource`, `defineMaterializedView`,
`column`, `t`, `engine`, `node` and `InferRow` with the call shapes of `@tinybirdco/sdk`, and
`buildProject` writing the same datafiles byte for byte. A datasource is a ClickHouse
`table`, so the query builder takes it and `renderSchema` writes its DDL. See
[docs/tinybird.md](./docs/tinybird.md).
- `CH.lowCardinality(t)`, `CH.simpleAggregateFunction(fn, t)` and `CH.precision(t, digits)`:
storage wrappers that change the DDL and not how the column reads. `CH.nullable` of a
`LowCardinality` type renders `LowCardinality(Nullable(T))`, which ClickHouse requires.

- Fixes found migrating a large consumer onto the strict params and brands:
- A param inside `when` / `whenTrue`, or one whose name is a union, is optional in the params
`compile` asks for: whether it renders is only known at runtime. It used to be required,
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,7 @@ regressions live in [`src/docs-examples.test.ts`](./src/docs-examples.test.ts).
| `@maple-dev/effect-orm/postgres` | The same for Postgres: the builder, Postgres column types (`text`, `int8`, `timestamptz`, …) and functions, `table` with keys, indexes and foreign keys, and a `compile` for Postgres. |
| `@maple-dev/effect-orm/expr` | Kitchen-sink namespace: every expression helper plus all ClickHouse functions under their raw names (`min_`, `toString_`, `toStartOfInterval`, `dynamicColumn`, …). |
| `@maple-dev/effect-orm/sql` | The low-level `SqlFragment` AST (`raw`, `ident`, `compile`, …) for hand-rolling fragments. |
| `@maple-dev/effect-orm/tinybird` | Tinybird datasources and materialized views, SDK-compatible, that are also query tables; `buildProject` writes the datafiles. See [Tinybird](./docs/tinybird.md). |
| `@maple-dev/effect-orm/schema` | Migration tooling over `table` values: DDL rendering, snapshots, the schema diff. Pure. |
| `@maple-dev/effect-orm/kit`, `/migrate` | `effect-orm generate` and `check`; applying migrations through a driver you provide. See [Schema and migrations](./docs/migrations.md). |
| `@maple-dev/effect-orm/database` | `Database` over your `SqlClient`: `run`, `execute`, `transaction` with retry. |
Expand Down
1 change: 1 addition & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ Roughly in reading order.
| `@maple-dev/effect-orm/benchmark` | Driver-free suite definitions, runner, report schemas, and comparisons |
| `@maple-dev/effect-orm/benchmark/http` | ClickHouse HTTP transport, environment configuration, and query-log collection |
| `@maple-dev/effect-orm/benchmark/cli` | `runCli(args)` for embedding the bundled `ch-bench` commands |
| `@maple-dev/effect-orm/tinybird` | Tinybird datasources and views with `@tinybirdco/sdk` call shapes; each datasource is a query table. See [Tinybird](./tinybird.md) |
| `@maple-dev/effect-orm/schema` | Tooling over `table` values: `renderSchema` / `renderPgSchema`, `entitiesOf` / `pgEntitiesOf`, snapshots, and the schema diff. Pure |
| `@maple-dev/effect-orm/kit` | `generate` and `check` over a migrations folder, `defineConfig`, and `runCli` for the bundled `effect-orm` command. Node or Bun |
| `@maple-dev/effect-orm/migrate` | `run`, `status`, `verify`, `baseline`, and `MigrationDriver`: applies ClickHouse or Postgres migrations through a driver you provide |
Expand Down
3 changes: 2 additions & 1 deletion docs/reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -300,7 +300,8 @@ that cannot render throws `SchemaDefinitionDefect` while the module loads.

**Constructors** — `string`, `bool`, `uint8`, `uint16`, `uint32`, `uint64`, `int32`,
`int64`, `float64`, `dateTime`, `dateTime64`, `dateTimeString`, `dateTime64String`, `map`,
`array`, `nullable`, `aggregateState(fn, ...args)`, `custom(sql, schema, literalSchema?)`,
`array`, `nullable`, `lowCardinality(type)`, `simpleAggregateFunction(fn, type)`,
`precision(dateTime64Type, digits)`, `aggregateState(fn, ...args)`, `custom(sql, schema, literalSchema?)`,
`brand(type, schema)` (the type narrowed by `schema`: a branded id, a literal union; see
[Branded columns](./tables-and-types.md#branded-columns)), and `untyped(sql)` for a wire value
passed through unvalidated. See [Tables and column types](./tables-and-types.md).
Expand Down
5 changes: 5 additions & 0 deletions docs/tables-and-types.md
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,11 @@ The constructors are values, not calls (except the parameterised ones):
| `CH.nullable(t)` | `Nullable(T)` | `T \| null` | value or `null` |
| `CH.untyped(sql)` | whatever you name | `unknown` | unvalidated |

Three wrappers change only the DDL, never how a column reads back:
`CH.lowCardinality(t)` (`LowCardinality(T)`; inside `nullable` it renders
`LowCardinality(Nullable(T))`), `CH.simpleAggregateFunction("sum", t)`
(`SimpleAggregateFunction(sum, T)`), and `CH.precision(CH.dateTime64String, 9)` (`DateTime64(9)`).

Column types, functions, and the query builder share the one `CH` namespace, so a schema module
and a query module import the same thing.

Expand Down
40 changes: 40 additions & 0 deletions docs/tinybird.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
# Tinybird datasources

`@maple-dev/effect-orm/tinybird` defines Tinybird datasources and materialized views with the
call shapes of `@tinybirdco/sdk` (`defineDatasource`, `column`, `t`, `engine`,
`defineMaterializedView`, `node`, `InferRow`), so a project moves over by changing its import.

What it adds over the SDK: a datasource **is** a ClickHouse `table` with its DDL. The query
builder takes it directly, `renderSchema` writes it for a plain ClickHouse server, and the
datafiles Tinybird deploys come from the same definition.

```ts
import * as CH from "@maple-dev/effect-orm/clickhouse"
import { buildProject, column, defineDatasource, engine, t } from "@maple-dev/effect-orm/tinybird"

export const events = defineDatasource("events", {
schema: {
OrgId: column(t.string().lowCardinality().brand(OrgId), { jsonPath: "$.org_id" }),
Timestamp: t.dateTime64(9),
Kind: t.string().lowCardinality().default("click"),
},
engine: engine.mergeTree({ sortingKey: ["OrgId", "Timestamp"], ttl: "toDate(Timestamp) + INTERVAL 30 DAY" }),
tenantColumn: "OrgId",
})

CH.from(events).select("Kind").where(($) => [$.OrgId.eq(CH.param.of(events.columns.OrgId, "orgId"))])

const { datasources, pipes } = buildProject(await import("./datasources"), await import("./views"))
```

- `t.*` builds a ClickHouse column type plus its datafile modifiers. `.lowCardinality()`,
`.nullable()`, `.default(v)`, `.defaultExpr(sql)`, `.codec(c)` and `.jsonPath(p)` behave as in
the SDK; `.brand(schema)` narrows the query column and leaves the ingested row alone.
- `DateTime` and `DateTime64` columns read back as the string ClickHouse sends, which is also
what `InferRow` types them as. `InferRow` types a `Map` as a `Record`, the JSON Tinybird ingests.
- `buildProject(...modules)` writes every datasource and view a module exports, in export order.
Its output matches `@tinybirdco/sdk` 0.0.84 byte for byte for the features here: schemas with
json paths, defaults and codecs, the MergeTree family and `Null` engines, indexes, forward
queries, and materialized views. Kafka, S3, tokens, endpoints and copy pipes are not ported.
- Materialized view SQL is a string. For a view checked against its target's columns, use
`CH.materializedView` from `/clickhouse`.
4 changes: 4 additions & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,10 @@
"types": "./dist/schema.d.mts",
"import": "./dist/schema.mjs"
},
"./tinybird": {
"types": "./dist/tinybird.d.mts",
"import": "./dist/tinybird.mjs"
},
"./migrate": {
"types": "./dist/migrate.d.mts",
"import": "./dist/migrate.mjs"
Expand Down
27 changes: 26 additions & 1 deletion src/ch/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -325,15 +325,40 @@ export const map = <K extends CHType<string, string, any>, V extends CHType<stri
export const array = <E extends CHType<string, any, any>>(e: E): CHArray<E> =>
chType("Array", `Array(${e.sql})`, Schema.Array(e.schema), undefined, e) as CHArray<E>

const LOW_CARDINALITY = /^LowCardinality\((.+)\)$/s

export const nullable = <T extends CHType<string, any, any>>(t: T): CHNullable<T> =>
chType(
"Nullable",
`Nullable(${t.sql})`,
// ClickHouse rejects `Nullable(LowCardinality(T))`; the wrapper goes outside.
LOW_CARDINALITY.test(t.sql) ? t.sql.replace(LOW_CARDINALITY, "LowCardinality(Nullable($1))") : `Nullable(${t.sql})`,
Schema.NullOr(t.schema),
Schema.NullOr(t.literalSchema),
t,
) as CHNullable<T>

/**
* `LowCardinality(T)`: dictionary-encoded storage. Queries see `T` unchanged,
* so only the DDL differs.
*/
export const lowCardinality = <T extends CHType<string, any, any>>(t: T): T =>
LOW_CARDINALITY.test(t.sql) ? t : { ...t, sql: `LowCardinality(${t.sql})` }

/**
* `SimpleAggregateFunction(fn, T)`: stores the aggregated value itself, so it
* reads back as `T`. Only the DDL differs.
*/
export const simpleAggregateFunction = <T extends CHType<string, any, any>>(fn: string, t: T): T => ({
...t,
sql: `SimpleAggregateFunction(${fn}, ${t.sql})`,
})

/** A `DateTime64` with its sub-second digits: `precision(dateTime64, 9)` is `DateTime64(9)`. */
export const precision = <T extends CHType<"DateTime64", any, any>>(
t: T,
digits: 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9,
): T => ({ ...t, sql: `DateTime64(${digits})` })

/**
* A column type of your own: a ClickHouse type name and the schema its wire
* value decodes with.
Expand Down
3 changes: 3 additions & 0 deletions src/clickhouse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,9 @@ export {
map,
array,
nullable,
lowCardinality,
simpleAggregateFunction,
precision,
custom,
brand,
} from "./ch/types"
Expand Down
20 changes: 20 additions & 0 deletions src/schema/schema.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,26 @@ describe("defineTable", () => {
)
})

it("renders storage-only wrappers without changing the column's query type", () => {
const Rollup = CH.table("rollup", {
columns: {
OrgId: CH.lowCardinality(CH.string),
Env: CH.nullable(CH.lowCardinality(CH.string)),
At: CH.precision(CH.dateTime64String, 9),
Calls: CH.simpleAggregateFunction("sum", CH.uint64),
},
engine: CH.engine.aggregatingMergeTree(),
orderBy: ["OrgId"],
})
const [table] = S.renderSchema(S.entitiesOf([Rollup]))
expect(table).toContain("\tOrgId LowCardinality(String),")
expect(table).toContain("\tEnv LowCardinality(Nullable(String)),")
expect(table).toContain("\tAt DateTime64(9),")
expect(table).toContain("\tCalls SimpleAggregateFunction(sum, UInt64)")
expect(Rollup.columns.OrgId._tag).toBe("String")
expect(Rollup.columns.Calls._tag).toBe("UInt64")
})

it("rejects a MergeTree without a sorting key", () => {
expect(() =>
CH.table("bad", { columns: { a: CH.string }, engine: CH.engine.mergeTree() }),
Expand Down
50 changes: 50 additions & 0 deletions src/tinybird.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
// @maple-dev/effect-orm/tinybird
//
// Tinybird datasources and materialized views with the call shapes of
// `@tinybirdco/sdk`. A datasource is a ClickHouse schema table, so the query
// builder takes it as-is, and `buildProject` writes its datafiles. See
// docs/tinybird.md.

export {
t,
getTinybirdType,
getModifiers,
isTinybirdType,
type TinybirdType,
type AnyTinybirdType,
type TypeModifiers,
type RowOf,
} from "./tinybird/types"

export {
column,
defineDatasource,
defineMaterializedView,
engine,
node,
getColumnType,
getColumnJsonPath,
isDatasourceDefinition,
isPipeDefinition,
formatDefaultValue,
type AnyDatasource,
type ColumnDefinition,
type ColumnsOfSchema,
type Datasource,
type DatasourceIndex,
type DatasourceOptions,
type DefaultedColumnsOfSchema,
type EngineConfig,
type InferRow,
type MaterializedViewDefinition,
type NodeDefinition,
type SchemaDefinition,
} from "./tinybird/datasource"

export {
buildProject,
generateDatasource,
generatePipe,
type Datafile,
type TinybirdProject,
} from "./tinybird/datafile"
143 changes: 143 additions & 0 deletions src/tinybird/datafile.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
// Tinybird datafiles (`.datasource`, `.pipe`), byte-compatible with the
// generator of `@tinybirdco/sdk` 0.0.84 (MIT), for the features ported here.

import {
defaultSqlOf,
getColumnJsonPath,
getColumnType,
isDatasourceDefinition,
isPipeDefinition,
splitKey,
type AnyDatasource,
type EngineConfig,
type MaterializedViewDefinition,
type NodeDefinition,
} from "./datasource"
import { getTinybirdType } from "./types"

export interface Datafile {
readonly name: string
readonly content: string
}

export interface TinybirdProject {
readonly datasources: ReadonlyArray<Datafile>
readonly pipes: ReadonlyArray<Datafile>
}

const indent = (text: string): ReadonlyArray<string> => text.split(/\r?\n/).map((line) => ` ${line}`)

const engineClause = (config: EngineConfig | undefined): string => {
if (config === undefined) return 'ENGINE "MergeTree"'
if (config.type === "Null") return "ENGINE Null"
const lines = [`ENGINE "${config.type}"`]
if (config.partitionKey) lines.push(`ENGINE_PARTITION_KEY "${config.partitionKey}"`)
lines.push(`ENGINE_SORTING_KEY "${normalizeKey(config.sortingKey)}"`)
if (config.primaryKey) lines.push(`ENGINE_PRIMARY_KEY "${normalizeKey(config.primaryKey)}"`)
if (config.ttl) lines.push(`ENGINE_TTL "${config.ttl}"`)
if (config.type === "ReplacingMergeTree" && config.ver) lines.push(`ENGINE_VER "${config.ver}"`)
if (config.type === "ReplacingMergeTree" && config.isDeleted) lines.push(`ENGINE_IS_DELETED "${config.isDeleted}"`)
if (config.type === "CollapsingMergeTree" || config.type === "VersionedCollapsingMergeTree") {
lines.push(`ENGINE_SIGN "${config.sign}"`)
}
if (config.type === "VersionedCollapsingMergeTree") lines.push(`ENGINE_VERSION "${config.version}"`)
if (config.type === "SummingMergeTree" && config.columns && config.columns.length > 0) {
lines.push(`ENGINE_SUMMING_COLUMNS "${config.columns.join(", ")}"`)
}
if (config.settings && Object.keys(config.settings).length > 0) {
const settings = Object.entries(config.settings)
.map(([key, value]) => (typeof value === "string" ? `${key}='${value.replace(/'/g, "\\'")}'` : `${key}=${value}`))
.join(", ")
lines.push(`ENGINE_SETTINGS "${settings}"`)
}
return lines.join("\n")
}

// The SDK joins an array key with ", " and keeps a string key as written.
const normalizeKey = (key: string | ReadonlyArray<string>): string => (typeof key === "string" ? key : key.join(", "))

export const generateDatasource = (datasource: AnyDatasource): Datafile => {
const { options } = datasource
const parts: Array<string> = []
if (options.description) parts.push(`DESCRIPTION >\n ${options.description}`, "")

const includeJsonPaths = options.jsonPaths !== false
const names = Object.keys(options.schema)
const columns = names.map((name, index) => {
const input = options.schema[name]!
const type = getColumnType(input)
const line = [` ${name} ${getTinybirdType(type)}`]
if (includeJsonPaths) line.push(`\`json:${getColumnJsonPath(input) ?? `$.${name}`}\``)
const defaultSql = defaultSqlOf(type)
if (defaultSql !== undefined) line.push(`DEFAULT ${defaultSql}`)
if (type.modifiers.codec) line.push(`CODEC(${type.modifiers.codec})`)
return line.join(" ") + (index < names.length - 1 ? "," : "")
})
parts.push(["SCHEMA >", ...columns].join("\n"), "")
parts.push(engineClause(options.engine))

if (options.indexes && options.indexes.length > 0) {
parts.push(
"",
[
"INDEXES >",
...options.indexes.map(
(index) => ` ${index.name} ${index.expr} TYPE ${index.type} GRANULARITY ${index.granularity}`,
),
].join("\n"),
)
}

const forwardQuery = options.forwardQuery?.trim()
if (forwardQuery) parts.push("", ["FORWARD_QUERY >", ...indent(forwardQuery)].join("\n"))

return { name: datasource._name, content: parts.join("\n") }
}

const TEMPLATE = /\{\{[^}]+\}\}|\{%[^%]+%\}/

const generateNode = (node: NodeDefinition): string => {
const lines = [`NODE ${node._name}`]
if (node.description) lines.push("DESCRIPTION >", ` ${node.description}`)
lines.push("SQL >")
if (TEMPLATE.test(node.sql)) lines.push(" %")
for (const line of node.sql.trim().split("\n")) lines.push(` ${line}`)
return lines.join("\n")
}

export const generatePipe = (pipe: MaterializedViewDefinition): Datafile => {
const { options } = pipe
const parts: Array<string> = []
if (options.description) parts.push(`DESCRIPTION >\n ${options.description}`, "")
options.nodes.forEach((node, index) => {
parts.push(generateNode(node))
if (index < options.nodes.length - 1) parts.push("")
})
const materialized = ["TYPE MATERIALIZED", `DATASOURCE ${options.materialized.datasource._name}`]
if (options.materialized.deploymentMethod === "alter") materialized.push("DEPLOYMENT_METHOD alter")
parts.push("", materialized.join("\n"))
return { name: pipe._name, content: parts.join("\n") }
}

/**
* The datafiles of every datasource and materialized view exported by the
* given modules, in export order (a module namespace lists exports by name).
*/
export const buildProject = (...modules: ReadonlyArray<Readonly<Record<string, unknown>>>): TinybirdProject => {
const seen = new Set<unknown>()
const datasources: Array<Datafile> = []
const pipes: Array<Datafile> = []
for (const module of modules) {
for (const value of Object.values(module)) {
if (seen.has(value)) continue
if (isDatasourceDefinition(value)) {
seen.add(value)
datasources.push(generateDatasource(value))
} else if (isPipeDefinition(value)) {
seen.add(value)
pipes.push(generatePipe(value))
}
}
}
return { datasources, pipes }
}
Loading
Loading