Skip to content
Draft
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
14 changes: 14 additions & 0 deletions .changeset/fix-cluster-mssql-for-update.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
---
"effect": patch
---

Fix SQL Server support in cluster SqlMessageStorage.

The mssql code paths contained failures across core storage paths:

- The unprocessed-message reads emitted a trailing `FOR UPDATE` clause, which T-SQL only allows on cursors, so the reads failed with a syntax error. The mssql dialect now omits the clause (as sqlite does) without adding an `UPDLOCK` hint, so concurrent duplicate observation stays within the existing at-least-once delivery contract.
- The messages and replies tables declared plain `UNIQUE` constraints on `message_id`, `(request_id, kind)` and `(request_id, sequence)`. SQL Server treats NULLs as equal in unique constraints, so the second message without a primary key — and the second chunk reply of any stream — failed with a unique violation. The constraints are replaced with filtered unique indexes (`WHERE ... IS NOT NULL`), matching the NULL-exempt semantics of the other dialects.
- The `insertEnvelope` MERGE statement used subqueries in its `OUTPUT` clause, which SQL Server rejects (error 10705), so saving a request with a primary key always failed — and its duplicate-detection logic was wrong besides, since a MERGE with only `WHEN NOT MATCHED` emits no OUTPUT row for an existing key.
- The `payload` and `headers` columns were declared as the deprecated `TEXT` type. The mssql client binds NULL parameters as `BIT`, which has no implicit conversion to `TEXT`, so inserting envelopes with NULL payloads (e.g. `AckChunk`) failed with an operand type clash. The columns are now `NVARCHAR(MAX)`.

A new `0003_mssql_schema_fixes` migration converges existing SQL Server databases (a no-op on other dialects): it drops the anonymous and named unique constraints, creates the filtered unique indexes, and alters the `TEXT` columns to `NVARCHAR(MAX)` in place, preserving existing data.
142 changes: 103 additions & 39 deletions packages/effect/src/unstable/cluster/SqlMessageStorage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -279,37 +279,17 @@ export const make: (options?: {
ON target.message_id = source.message_id
WHEN NOT MATCHED THEN
INSERT ${sql.insert(row)}
OUTPUT
inserted.id,
CASE
WHEN inserted.id IS NULL THEN (
SELECT r.id, r.kind, r.payload
FROM ${repliesTableSql} r
WHERE r.id = target.last_reply_id
)
END as reply_id,
CASE
WHEN inserted.id IS NULL THEN (
SELECT r.kind
FROM ${repliesTableSql} r
WHERE r.id = target.last_reply_id
)
END as reply_kind,
CASE
WHEN inserted.id IS NULL THEN (
SELECT r.payload
FROM ${repliesTableSql} r
WHERE r.id = target.last_reply_id
)
END as reply_payload,
CASE
WHEN inserted.id IS NULL THEN (
SELECT r.sequence
FROM ${repliesTableSql} r
WHERE r.id = target.last_reply_id
)
END as reply_sequence;
`,
OUTPUT inserted.id;
`.pipe(Effect.flatMap((rows) => {
// inserted a new row
if (rows.length > 0) return Effect.succeed([])
return sql`
SELECT m.id, r.id as reply_id, r.kind as reply_kind, r.payload as reply_payload, r.sequence as reply_sequence
FROM ${messagesTableSql} m
LEFT JOIN ${repliesTableSql} r ON r.id = m.last_reply_id
WHERE m.message_id = ${message_id}
`
})),
orElse: () => (row, message_id) =>
sql`
SELECT m.id, r.id as reply_id, r.kind as reply_kind, r.payload as reply_payload, r.sequence as reply_sequence
Expand Down Expand Up @@ -343,6 +323,8 @@ export const make: (options?: {
})
const forUpdate = sql.onDialectOrElse({
sqlite: () => sql.literal(""),
// SQL Server uses table hints instead of a trailing FOR UPDATE clause.
mssql: () => sql.literal(""),
orElse: () => sql.literal("FOR UPDATE")
})

Expand Down Expand Up @@ -711,8 +693,8 @@ const migrations = (options?: {
entity_id VARCHAR(255) NOT NULL,
kind INT NOT NULL,
tag VARCHAR(50),
payload TEXT,
headers TEXT,
payload NVARCHAR(MAX),
headers NVARCHAR(MAX),
trace_id VARCHAR(32),
span_id VARCHAR(16),
sampled BIT,
Expand All @@ -721,8 +703,7 @@ const migrations = (options?: {
reply_id BIGINT,
last_reply_id BIGINT,
last_read DATETIME,
deliver_at BIGINT,
UNIQUE (message_id)
deliver_at BIGINT
)
`,
mysql: () =>
Expand Down Expand Up @@ -866,11 +847,9 @@ const migrations = (options?: {
rowid BIGINT IDENTITY(1,1),
kind INT,
request_id BIGINT NOT NULL,
payload TEXT NOT NULL,
payload NVARCHAR(MAX) 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)
acked BIT NOT NULL DEFAULT 0
)
`,
mysql: () =>
Expand Down Expand Up @@ -977,6 +956,91 @@ const migrations = (options?: {
// sqlite
Effect.void
})
}),
"0003_mssql_schema_fixes": Effect.gen(function*() {
const sql = (yield* SqlClient.SqlClient).withoutTransforms()
const messagesTableSql = sql(messagesTable)
const repliesTableSql = sql(repliesTable)
const messageIdUniqueIndex = `${messagesTable}_message_id_idx`
const oneExitUniqueIndex = `${repliesTable}_one_exit`
const sequenceUniqueIndex = `${repliesTable}_sequence`

// SQL Server unique constraints treat NULLs as equal, unlike the other
// dialects, so the inline constraints created by 0001 reject rows that
// are legal everywhere else: the second message without a primary key
// (message_id NULL) and the second chunk reply of a request (kind
// NULL). Replace them with filtered unique indexes covering only
// non-NULL keys; the sequence key receives the same filtered treatment
// to match the other dialects' nullable-unique semantics. The
// message_id constraint from 0001 was anonymous, so it is looked up by
// shape before being dropped.
//
// 0001 also declared the payload and headers columns as TEXT, which is
// deprecated and has no conversion from BIT, the type the client binds
// NULL parameters as. Convert them to NVARCHAR(MAX).
yield* sql.onDialectOrElse({
mssql: () =>
sql`
DECLARE @constraint_name sysname;
DECLARE @drop_constraint_sql nvarchar(max);

SELECT @constraint_name = kc.name
FROM sys.key_constraints AS kc
WHERE kc.parent_object_id = OBJECT_ID(N'${messagesTableSql}')
AND kc.type = 'UQ'
AND EXISTS (
SELECT 1
FROM sys.index_columns AS ic
INNER JOIN sys.columns AS c
ON c.object_id = ic.object_id
AND c.column_id = ic.column_id
WHERE ic.object_id = kc.parent_object_id
AND ic.index_id = kc.unique_index_id
AND ic.is_included_column = 0
AND ic.key_ordinal > 0
GROUP BY ic.object_id, ic.index_id
HAVING COUNT(*) = 1
AND MAX(CASE WHEN c.name = N'message_id' THEN 1 ELSE 0 END) = 1
);

IF @constraint_name IS NOT NULL
BEGIN
SET @drop_constraint_sql = N'ALTER TABLE ${messagesTableSql} DROP CONSTRAINT ' + QUOTENAME(@constraint_name);
EXEC sp_executesql @drop_constraint_sql;
END;

IF EXISTS (SELECT * FROM sys.key_constraints WHERE name = ${oneExitUniqueIndex} AND parent_object_id = OBJECT_ID(N'${repliesTableSql}'))
ALTER TABLE ${repliesTableSql} DROP CONSTRAINT ${sql(oneExitUniqueIndex)};

IF EXISTS (SELECT * FROM sys.key_constraints WHERE name = ${sequenceUniqueIndex} AND parent_object_id = OBJECT_ID(N'${repliesTableSql}'))
ALTER TABLE ${repliesTableSql} DROP CONSTRAINT ${sql(sequenceUniqueIndex)};

IF NOT EXISTS (SELECT * FROM sys.indexes WHERE name = ${messageIdUniqueIndex} AND object_id = OBJECT_ID(N'${messagesTableSql}'))
CREATE UNIQUE INDEX ${sql(messageIdUniqueIndex)}
ON ${messagesTableSql} (message_id)
WHERE message_id IS NOT NULL;

IF NOT EXISTS (SELECT * FROM sys.indexes WHERE name = ${oneExitUniqueIndex} AND object_id = OBJECT_ID(N'${repliesTableSql}'))
CREATE UNIQUE INDEX ${sql(oneExitUniqueIndex)}
ON ${repliesTableSql} (request_id, kind)
WHERE kind IS NOT NULL;

IF NOT EXISTS (SELECT * FROM sys.indexes WHERE name = ${sequenceUniqueIndex} AND object_id = OBJECT_ID(N'${repliesTableSql}'))
CREATE UNIQUE INDEX ${sql(sequenceUniqueIndex)}
ON ${repliesTableSql} (request_id, sequence)
WHERE sequence IS NOT NULL;

IF EXISTS (SELECT * FROM sys.columns WHERE object_id = OBJECT_ID(N'${messagesTableSql}') AND name = N'payload' AND system_type_id = TYPE_ID(N'text'))
ALTER TABLE ${messagesTableSql} ALTER COLUMN payload NVARCHAR(MAX);

IF EXISTS (SELECT * FROM sys.columns WHERE object_id = OBJECT_ID(N'${messagesTableSql}') AND name = N'headers' AND system_type_id = TYPE_ID(N'text'))
ALTER TABLE ${messagesTableSql} ALTER COLUMN headers NVARCHAR(MAX);

IF EXISTS (SELECT * FROM sys.columns WHERE object_id = OBJECT_ID(N'${repliesTableSql}') AND name = N'payload' AND system_type_id = TYPE_ID(N'text'))
ALTER TABLE ${repliesTableSql} ALTER COLUMN payload NVARCHAR(MAX) NOT NULL;
`,
orElse: () => Effect.void
})
})
})
}
Expand Down
1 change: 1 addition & 0 deletions packages/platform-node/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@
"ioredis": "^5.7.0"
},
"devDependencies": {
"@testcontainers/mssqlserver": "^11.14.0",
"@testcontainers/mysql": "^11.14.0",
"@testcontainers/postgresql": "^11.14.0",
"@testcontainers/redis": "^11.14.0",
Expand Down
49 changes: 19 additions & 30 deletions packages/platform-node/test/cluster/SqlMessageStorage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ 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 { SqlClient } from "effect/unstable/sql"
import { MssqlContainer } from "../fixtures/mssql-utils.ts"
import { MysqlContainer } from "../fixtures/mysql2-utils.ts"
import { PgContainer } from "../fixtures/pg-utils.ts"
import {
Expand All @@ -16,24 +17,22 @@ import {
StreamRpc
} from "./MessageStorageTest.ts"

const StorageLive = SqlMessageStorage.layer.pipe(
Layer.provideMerge(Snowflake.layerGenerator),
Layer.provide(ShardingConfig.layerDefaults)
)
// Each test provides its own storage layer with a distinct table prefix, so
// concurrently running tests never share tables.
const storageLive = (prefix: string) =>
SqlMessageStorage.layerWith({ prefix }).pipe(
Layer.provideMerge(Snowflake.layerGenerator),
Layer.provide(ShardingConfig.layerDefaults)
)

const truncate = Effect.gen(function*() {
const sql = yield* SqlClient.SqlClient
yield* sql`DELETE FROM cluster_replies`
yield* sql`DELETE FROM cluster_messages`
})

describe("SqlMessageStorage", () => {
describe("SqlMessageStorage", { timeout: 30_000 }, () => {
;([
["pg", Layer.orDie(PgContainer.layerClient)],
["mysql", Layer.orDie(MysqlContainer.layerClient)],
["mssql", Layer.orDie(MssqlContainer.layerClient)],
["sqlite", Layer.orDie(SqliteLayer)]
] as const).forEach(([label, layer]) => {
it.layer(StorageLive.pipe(Layer.provideMerge(layer)), {
it.layer(layer, {
timeout: 120000
})(label, (it) => {
it.effect("saveRequest", () =>
Expand All @@ -59,7 +58,7 @@ describe("SqlMessageStorage", () => {
messages = yield* storage.unprocessedMessages([request.envelope.address.shardId])
expect(messages).toHaveLength(5)
expect(messages.map((m: any) => m.envelope.payload.id)).toEqual([6, 7, 8, 9, 10])
}))
}).pipe(Effect.provide(storageLive("save_request"))))

it.effect("saveReply + saveRequest duplicate", () =>
Effect.gen(function*() {
Expand Down Expand Up @@ -113,12 +112,10 @@ describe("SqlMessageStorage", () => {
}
const error = yield* Effect.flip(Fiber.join(fiber))
expect(error._tag).toEqual("PersistenceError")
}))
}).pipe(Effect.provide(storageLive("stream_replies"))))

it.effect("detects duplicates", () =>
Effect.gen(function*() {
yield* truncate

const storage = yield* MessageStorage.MessageStorage
yield* storage.saveRequest(
yield* makeRequest({
Expand All @@ -133,12 +130,10 @@ describe("SqlMessageStorage", () => {
})
)
expect(result._tag).toEqual("Duplicate")
}))
}).pipe(Effect.provide(storageLive("detects_duplicates"))))

it.effect("unprocessedMessages", () =>
Effect.gen(function*() {
yield* truncate

const storage = yield* MessageStorage.MessageStorage
const request = yield* makeRequest()
yield* storage.saveRequest(request)
Expand All @@ -149,24 +144,20 @@ describe("SqlMessageStorage", () => {
yield* storage.saveRequest(yield* makeRequest())
messages = yield* storage.unprocessedMessages([request.envelope.address.shardId])
expect(messages).toHaveLength(1)
}))
}).pipe(Effect.provide(storageLive("unprocessed_messages"))))

it.effect("unprocessedMessages excludes complete requests", () =>
Effect.gen(function*() {
yield* truncate

const storage = yield* MessageStorage.MessageStorage
const request = yield* makeRequest()
yield* storage.saveRequest(request)
yield* storage.saveReply(yield* makeReply(request))
const messages = yield* storage.unprocessedMessages([request.envelope.address.shardId])
expect(messages).toHaveLength(0)
}))
}).pipe(Effect.provide(storageLive("unprocessed_complete"))))

it.effect("repliesFor", () =>
Effect.gen(function*() {
yield* truncate

const storage = yield* MessageStorage.MessageStorage
const request = yield* makeRequest()
yield* storage.saveRequest(request)
Expand All @@ -176,7 +167,7 @@ describe("SqlMessageStorage", () => {
replies = yield* storage.repliesFor([request])
expect(replies).toHaveLength(1)
expect(replies[0].requestId).toEqual(request.envelope.requestId)
}))
}).pipe(Effect.provide(storageLive("replies_for"))))

it.effect("registerReplyHandler", () =>
Effect.gen(function*() {
Expand All @@ -194,12 +185,10 @@ describe("SqlMessageStorage", () => {
yield* storage.saveReply(yield* makeReply(request))
yield* latch.await
yield* Fiber.await(fiber)
}))
}).pipe(Effect.provide(storageLive("register_reply_handler"))))

it.effect("unprocessedMessagesById", () =>
Effect.gen(function*() {
yield* truncate

const storage = yield* MessageStorage.MessageStorage
const request = yield* makeRequest()
yield* storage.saveRequest(request)
Expand All @@ -208,7 +197,7 @@ describe("SqlMessageStorage", () => {
yield* storage.saveReply(yield* makeReply(request))
messages = yield* storage.unprocessedMessagesById([request.envelope.requestId])
expect(messages).toHaveLength(0)
}))
}).pipe(Effect.provide(storageLive("unprocessed_by_id"))))
})
})
})
Expand Down
Loading