Skip to content
Open
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
5 changes: 5 additions & 0 deletions .changeset/cluster-sql-identifier-limit.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"effect": minor
---

Add opt-in `limitIdentifiers` to `SqlMessageStorage` so long table/index/constraint names can be capped at the Postgres 63-character identifier limit with a stable hash suffix. Default naming is unchanged.
28 changes: 21 additions & 7 deletions packages/effect/src/unstable/cluster/SqlMessageStorage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import type { Row } from "../sql/SqlConnection.ts"
import { isSqlError, type SqlError } from "../sql/SqlError.ts"
import { PersistenceError } from "./ClusterError.ts"
import type * as Envelope from "./Envelope.ts"
import { storageObjectName, storageTableName } from "./internal/sqlIdentifier.ts"
import * as MessageStorage from "./MessageStorage.ts"
import { SaveResultEncoded } from "./MessageStorage.ts"
import type * as Reply from "./Reply.ts"
Expand All @@ -45,6 +46,12 @@ const withTracerDisabled = Effect.withTracerEnabled(false)
* The optional `prefix` controls the table names for messages, replies, and
* migrations; when omitted, `cluster` is used.
*
* Set `limitIdentifiers` when using long table prefixes. Without it, names are
* used as-is (including lengths that Postgres may silently truncate). Enabling
* the flag hash-suffixes identifiers longer than 63 characters so table and
* index names stay unique — restart / migrate carefully if existing truncated
* names are already in the database.
*
* **Gotchas**
*
* Changing `prefix` after deployment points the runtime at a different set of
Expand All @@ -58,14 +65,16 @@ const withTracerDisabled = Effect.withTracerEnabled(false)
*/
export const make: (options?: {
readonly prefix?: string | undefined
readonly limitIdentifiers?: boolean | undefined
}) => Effect.Effect<
MessageStorage.MessageStorage["Service"],
never,
SqlClient.SqlClient | Snowflake.Generator
> = Effect.fnUntraced(function*(options) {
const sql = (yield* SqlClient.SqlClient).withoutTransforms()
const prefix = options?.prefix ?? "cluster"
const table = (name: string) => `${prefix}_${name}`
const limitIdentifiers = options?.limitIdentifiers === true
const table = (name: string) => storageTableName(prefix, name, limitIdentifiers)

yield* Effect.orDie(
Migrator.make({})({
Expand Down Expand Up @@ -675,6 +684,7 @@ export const layer: Layer.Layer<
*/
export const layerWith = (options: {
readonly prefix?: string | undefined
readonly limitIdentifiers?: boolean | undefined
}): Layer.Layer<MessageStorage.MessageStorage, never, SqlClient.SqlClient | ShardingConfig> =>
Layer.effect(MessageStorage.MessageStorage, make(options)).pipe(
Layer.provide(Snowflake.layerGenerator)
Expand All @@ -686,9 +696,13 @@ export const layerWith = (options: {

const migrations = (options?: {
readonly prefix?: string | undefined
readonly limitIdentifiers?: boolean | undefined
}) => {
const prefix = options?.prefix ?? "cluster"
const table = (name: string) => `${prefix}_${name}`
const limitIdentifiers = options?.limitIdentifiers === true
const table = (name: string) => storageTableName(prefix, name, limitIdentifiers)
const objectName = (tableName: string, suffix: string) =>
storageObjectName(prefix, tableName, suffix, limitIdentifiers)
const messagesTable = table("messages")
const repliesTable = table("replies")

Expand Down Expand Up @@ -804,8 +818,8 @@ const migrations = (options?: {
})

// Add message indexes optimized for the specific query patterns
const shardLookupIndex = `${messagesTable}_shard_idx`
const requestIdLookupIndex = `${messagesTable}_request_id_idx`
const shardLookupIndex = objectName("messages", "shard_idx")
const requestIdLookupIndex = objectName("messages", "request_id_idx")
yield* sql.onDialectOrElse({
mssql: () =>
sql`
Expand Down Expand Up @@ -869,8 +883,8 @@ const migrations = (options?: {
payload TEXT NOT NULL,
sequence INT,
acked BIT NOT NULL DEFAULT 0,
CONSTRAINT ${sql(repliesTable + "_one_exit")} UNIQUE (request_id, kind),
CONSTRAINT ${sql(repliesTable + "_sequence")} UNIQUE (request_id, sequence)
CONSTRAINT ${sql(objectName("replies", "one_exit"))} UNIQUE (request_id, kind),
CONSTRAINT ${sql(objectName("replies", "sequence"))} UNIQUE (request_id, sequence)
)
`,
mysql: () =>
Expand Down Expand Up @@ -919,7 +933,7 @@ const migrations = (options?: {
})

// Add reply indexes optimized for request_id lookups
const replyLookupIndex = `${repliesTable}_request_lookup_idx`
const replyLookupIndex = objectName("replies", "request_lookup_idx")
yield* sql.onDialectOrElse({
mssql: () =>
sql`
Expand Down
54 changes: 54 additions & 0 deletions packages/effect/src/unstable/cluster/internal/sqlIdentifier.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
/**
* Fits SQL identifiers into PostgreSQL's `NAMEDATALEN - 1` limit (63) without
* silent truncation collisions. Long names keep a stable hash suffix.
*
* @internal
*/
import { hashString } from "./hash.ts"

/** PostgreSQL `NAMEDATALEN - 1` identifier limit (also safe for MySQL). */
/** @internal */
export const SQL_IDENTIFIER_MAX_LENGTH = 63

/**
* Fits an identifier into {@link SQL_IDENTIFIER_MAX_LENGTH}.
*
* @internal
*/
export const sqlIdentifier = (name: string): string => {
if (name.length <= SQL_IDENTIFIER_MAX_LENGTH) {
return name
}
const digest = (hashString(name) >>> 0).toString(36)
const keep = SQL_IDENTIFIER_MAX_LENGTH - digest.length - 1
return `${name.slice(0, keep)}_${digest}`
}

/**
* Table name from prefix + logical table name.
*
* When `limit` is true, long names are capped at {@link SQL_IDENTIFIER_MAX_LENGTH}.
*
* @internal
*/
export const storageTableName = (prefix: string, name: string, limit = false): string => {
const raw = `${prefix}_${name}`
return limit ? sqlIdentifier(raw) : raw
}

/**
* Index / constraint name from prefix, logical table name, and suffix.
*
* When `limit` is true, long names are capped at {@link SQL_IDENTIFIER_MAX_LENGTH}.
*
* @internal
*/
export const storageObjectName = (
prefix: string,
tableName: string,
suffix: string,
limit = false
): string => {
const raw = `${prefix}_${tableName}_${suffix}`
return limit ? sqlIdentifier(raw) : raw
}
57 changes: 57 additions & 0 deletions packages/effect/test/cluster/sqlIdentifier.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
import { assert, describe, it } from "@effect/vitest"
import {
SQL_IDENTIFIER_MAX_LENGTH,
sqlIdentifier,
storageObjectName,
storageTableName
} from "effect/unstable/cluster/internal/sqlIdentifier"

describe("sqlIdentifier", () => {
it("leaves short identifiers unchanged", () => {
assert.strictEqual(sqlIdentifier("cluster_messages_shard_idx"), "cluster_messages_shard_idx")
assert.strictEqual(storageTableName("cluster", "messages"), "cluster_messages")
assert.strictEqual(storageObjectName("cluster", "messages", "shard_idx"), "cluster_messages_shard_idx")
assert.strictEqual(storageObjectName("cluster", "messages", "request_id_idx"), "cluster_messages_request_id_idx")
assert.strictEqual(
storageObjectName("cluster", "replies", "request_lookup_idx"),
"cluster_replies_request_lookup_idx"
)
})

it("does not limit long names unless opted in", () => {
const longPrefix = "p".repeat(80)
const unlimited = storageObjectName(longPrefix, "messages", "request_id_idx")
assert.isTrue(unlimited.length > SQL_IDENTIFIER_MAX_LENGTH)
assert.strictEqual(unlimited, `${longPrefix}_messages_request_id_idx`)
})

it("keeps long table and index names within the SQL identifier limit when enabled", () => {
const longPrefix = "p".repeat(80)
const logicalIndex = `${longPrefix}_messages_request_id_idx`
assert.isTrue(logicalIndex.length > SQL_IDENTIFIER_MAX_LENGTH)

const names = [
storageTableName(longPrefix, "messages", true),
storageObjectName(longPrefix, "messages", "shard_idx", true),
storageObjectName(longPrefix, "messages", "request_id_idx", true),
storageObjectName(longPrefix, "replies", "request_lookup_idx", true),
storageObjectName(longPrefix, "replies", "one_exit", true),
storageObjectName(longPrefix, "replies", "sequence", true)
]
for (const name of names) {
assert.isTrue(name.length <= SQL_IDENTIFIER_MAX_LENGTH, name)
}
assert.strictEqual(
storageObjectName(longPrefix, "messages", "shard_idx", true),
storageObjectName(longPrefix, "messages", "shard_idx", true)
)
assert.notStrictEqual(
storageObjectName(longPrefix, "messages", "shard_idx", true),
storageObjectName(longPrefix + "x", "messages", "shard_idx", true)
)
assert.notStrictEqual(
storageObjectName(longPrefix, "messages", "shard_idx", true),
storageObjectName(longPrefix, "messages", "request_id_idx", true)
)
})
})
71 changes: 71 additions & 0 deletions packages/platform-node/test/cluster/SqlMessageStorage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@ import { assert, describe, expect, it } from "@effect/vitest"
import { Effect, Fiber, FileSystem, Latch, Layer, Option } from "effect"
import { TestClock } from "effect/testing"
import { Message, MessageStorage, ShardingConfig, Snowflake, SqlMessageStorage } from "effect/unstable/cluster"
import {
SQL_IDENTIFIER_MAX_LENGTH,
storageObjectName,
storageTableName
} from "effect/unstable/cluster/internal/sqlIdentifier"
import { SqlClient } from "effect/unstable/sql"
import { MysqlContainer } from "../fixtures/mysql2-utils.ts"
import { PgContainer } from "../fixtures/pg-utils.ts"
Expand Down Expand Up @@ -211,6 +216,72 @@ describe("SqlMessageStorage", () => {
}))
})
})

it.layer(
SqlMessageStorage.layerWith({
prefix: "p".repeat(80),
limitIdentifiers: true
}).pipe(
Layer.provideMerge(Snowflake.layerGenerator),
Layer.provide(ShardingConfig.layerDefaults),
Layer.provideMerge(Layer.orDie(PgContainer.layerClient))
),
{ timeout: 120000 }
)("pg long identifiers", (it) => {
it.effect("creates truncated table and index names and stores messages", () =>
Effect.gen(function*() {
const sql = yield* SqlClient.SqlClient
const prefix = "p".repeat(80)
const messagesTable = storageTableName(prefix, "messages", true)
const indexes = [
storageObjectName(prefix, "messages", "shard_idx", true),
storageObjectName(prefix, "messages", "request_id_idx", true),
storageObjectName(prefix, "replies", "request_lookup_idx", true)
]

assert.isTrue(messagesTable.length <= SQL_IDENTIFIER_MAX_LENGTH)
for (const name of indexes) {
assert.isTrue(name.length <= SQL_IDENTIFIER_MAX_LENGTH)
}

const tables = yield* sql<{ tablename: string }>`
SELECT tablename FROM pg_tables
WHERE schemaname = 'public'
AND tablename IN ${
sql.in([
messagesTable,
storageTableName(prefix, "migrations", true),
storageTableName(prefix, "replies", true)
])
}
ORDER BY tablename
`
assert.deepStrictEqual(
tables.map((row) => row.tablename).sort(),
[
messagesTable,
storageTableName(prefix, "migrations", true),
storageTableName(prefix, "replies", true)
].sort()
)

const indexRows = yield* sql<{ indexname: string }>`
SELECT indexname FROM pg_indexes
WHERE schemaname = 'public'
AND indexname IN ${sql.in(indexes)}
`
assert.deepStrictEqual(indexRows.map((row) => row.indexname).sort(), [...indexes].sort())

const storage = yield* MessageStorage.MessageStorage
const request = yield* makeRequest({ payload: { id: 1 } })
const result = yield* storage.saveRequest(request)
assert.strictEqual(result._tag, "Success")
assert.strictEqual(
(yield* storage.unprocessedMessages([request.envelope.address.shardId])).length,
1
)
}))
})
})

const SqliteLayer = Effect.gen(function*() {
Expand Down