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/batch-persistence-expiration-cleanup.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"effect": patch
---

Run shared-table SQL persistence expiration cleanup in indexed, bounded background batches.
2 changes: 2 additions & 0 deletions packages/effect/src/internal/persistence.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
/** @internal */
export const sqlCleanupBatchSize = 1000
129 changes: 124 additions & 5 deletions packages/effect/src/unstable/persistence/Persistence.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,11 @@ import * as Duration from "../../Duration.ts"
import * as Effect from "../../Effect.ts"
import * as Exit from "../../Exit.ts"
import { identity } from "../../Function.ts"
import { sqlCleanupBatchSize } from "../../internal/persistence.ts"
import * as Layer from "../../Layer.ts"
import * as PrimaryKey from "../../PrimaryKey.ts"
import * as Result from "../../Result.ts"
import * as Schedule from "../../Schedule.ts"
import * as Schema from "../../Schema.ts"
import type * as Scope from "../../Scope.ts"
import * as SqlClient from "../sql/SqlClient.ts"
Expand Down Expand Up @@ -560,6 +563,127 @@ export const layerBackingSql: Layer.Layer<
`
}).pipe(Effect.orDie)

yield* sql.onDialectOrElse({
pg: () =>
sql`CREATE INDEX IF NOT EXISTS effect_persistence_expires_idx ON ${table} (expires) WHERE expires IS NOT NULL`
.pipe(Effect.orDie, Effect.asVoid),
mysql: () =>
Effect.gen(function*() {
const indexExists = sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM information_schema.statistics
WHERE table_schema = DATABASE()
AND table_name = 'effect_persistence'
AND index_name = 'effect_persistence_expires_idx'
`.pipe(
Effect.map((rows) => Number(rows[0].count) > 0),
Effect.orDie
)
if (yield* indexExists) return

const createIndexResult = yield* sql`CREATE INDEX effect_persistence_expires_idx ON ${table} (expires)`
.pipe(Effect.result)
if (!(yield* indexExists)) {
if (Result.isFailure(createIndexResult)) {
return yield* Effect.die(createIndexResult.failure)
}
return yield* Effect.die(new Error("Failed to create effect_persistence_expires_idx"))
}
}),
Comment thread
AVtheking marked this conversation as resolved.
mssql: () =>
sql`
IF NOT EXISTS (
SELECT * FROM sys.indexes
WHERE name = N'effect_persistence_expires_idx'
AND object_id = OBJECT_ID(N'effect_persistence')
)
CREATE INDEX effect_persistence_expires_idx ON ${table} (expires) WHERE expires IS NOT NULL
`.pipe(Effect.orDie, Effect.asVoid),
// sqlite
orElse: () =>
sql`CREATE INDEX IF NOT EXISTS effect_persistence_expires_idx ON ${table} (expires) WHERE expires IS NOT NULL`
.pipe(Effect.orDie, Effect.asVoid)
})

const cleanupBatchDelay = Duration.millis(10)
const cleanupInterval = Duration.minutes(5)

const deleteExpiredBatch = sql.onDialectOrElse({
pg: () => (expiresAtOrBefore: number) =>
sql<{ readonly tupleId: string }>`
WITH expired_entries AS (
SELECT ctid FROM ${table}
WHERE expires IS NOT NULL AND expires <= ${expiresAtOrBefore}
ORDER BY expires
LIMIT ${sql.literal(String(sqlCleanupBatchSize))}
)
DELETE FROM ${table}
WHERE ctid IN (SELECT ctid FROM expired_entries)
RETURNING ctid::text AS "tupleId"
`.pipe(Effect.map((deletedEntries) => deletedEntries.length)),
mysql: () =>
Effect.fnUntraced(
function*(expiresAtOrBefore: number) {
yield* sql`
DELETE FROM ${table}
WHERE expires IS NOT NULL AND expires <= ${expiresAtOrBefore}
ORDER BY expires
LIMIT ${sql.literal(String(sqlCleanupBatchSize))}
`
const rows = yield* sql<{ readonly count: number }>`SELECT ROW_COUNT() AS count`
return Number(rows[0].count)
},
(effect) => effect.pipe(sql.withTransaction)
),
mssql: () => (expiresAtOrBefore: number) =>
sql<{ readonly store_id: string }>`
WITH expired_entries AS (
SELECT TOP ${sql.literal(String(sqlCleanupBatchSize))} store_id, id FROM ${table}
WITH (UPDLOCK, READPAST, READCOMMITTEDLOCK)
WHERE expires IS NOT NULL AND expires <= ${expiresAtOrBefore}
ORDER BY expires
)
DELETE persistence
OUTPUT DELETED.store_id
FROM ${table} AS persistence
INNER JOIN expired_entries
ON persistence.store_id = expired_entries.store_id
AND persistence.id = expired_entries.id
`.pipe(Effect.map((deletedEntries) => deletedEntries.length)),
// Some sqlite clients do not support interactive transactions, so use one bounded statement.
orElse: () => (expiresAtOrBefore: number) =>
sql<{ readonly store_id: string }>`
DELETE FROM ${table}
WHERE rowid IN (
SELECT rowid FROM ${table}
WHERE expires IS NOT NULL AND expires <= ${expiresAtOrBefore}
ORDER BY expires
LIMIT ${sql.literal(String(sqlCleanupBatchSize))}
)
RETURNING store_id
`.pipe(Effect.map((deletedEntries) => deletedEntries.length))
})

// Delay between batches and stop once a batch deletes no rows.
const deleteExpiredSchedule = Schedule.forever.pipe(
Schedule.setInputType<number>(),
Schedule.passthrough,
Schedule.while(({ input: deletedCount }) => deletedCount !== 0),
Schedule.addDelay(() => Effect.succeed(cleanupBatchDelay))
)

const deleteExpired = Effect.fnUntraced(function*() {
const expiresAtOrBefore = yield* Clock.currentTimeMillis
return yield* deleteExpiredBatch(expiresAtOrBefore).pipe(
Effect.repeat(deleteExpiredSchedule)
)
})

yield* deleteExpired().pipe(
Effect.catch((cause) => Effect.logWarning("Failed to clean up expired persistence entries", cause)),
Effect.repeat(Schedule.spaced(cleanupInterval)),
Effect.forkScoped
)

type UpsertFn = (
entries: Array<{ store_id: string; id: string; value: string; expires: number | null }>
) => Effect.Effect<unknown, SqlError>
Expand Down Expand Up @@ -606,11 +730,6 @@ export const layerBackingSql: Layer.Layer<
make: Effect.fnUntraced(function*(storeId) {
const clock = yield* Clock.Clock

// Cleanup expired entries on startup
yield* Effect.ignore(
sql`DELETE FROM ${table} WHERE store_id = ${storeId} AND expires IS NOT NULL AND expires <= ${clock.currentTimeMillisUnsafe()}`
)

return identity<BackingPersistenceStore>({
get: (key) =>
sql<
Expand Down
20 changes: 20 additions & 0 deletions packages/effect/test/unstable/persistence/SqlCleanupTest.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
import { Duration, Effect, Schedule } from "effect"
import { sqlCleanupBatchSize } from "effect/internal/persistence"
import { TestClock } from "effect/testing"

export const expiredEntryCount = sqlCleanupBatchSize + 1
export const cleanupBatchDelay = Duration.millis(10)
export const testTimeout = 30_000
Comment thread
coderabbitai[bot] marked this conversation as resolved.

export const waitForCount = <E, R>(
effect: Effect.Effect<number, E, R>,
predicate: (count: number) => boolean
) =>
effect.pipe(
Effect.repeat({
until: predicate,
schedule: Schedule.spaced(cleanupBatchDelay)
}),
Effect.timeout(Duration.millis(testTimeout)),
TestClock.withLive
)
65 changes: 64 additions & 1 deletion packages/sql/mysql2/test/Persistence.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
import { Layer } from "effect"
import { assert, it } from "@effect/vitest"
import { Effect, Layer } from "effect"
import * as PersistedCacheTest from "effect-test/unstable/persistence/PersistedCacheTest"
import * as PersistedQueueTest from "effect-test/unstable/persistence/PersistedQueueTest"
import * as SqlCleanupTest from "effect-test/unstable/persistence/SqlCleanupTest"
import { TestClock } from "effect/testing"
import { PersistedQueue, Persistence } from "effect/unstable/persistence"
import { SqlClient } from "effect/unstable/sql"
import { MysqlContainer } from "./utils.ts"

PersistedCacheTest.suite(
Expand All @@ -20,3 +24,62 @@ PersistedQueueTest.suite(
Layer.provide(MysqlContainer.layerClient)
)
)

it.layer(MysqlContainer.layerClient, { timeout: "30 seconds" })("Persistence SQL cleanup", (it) => {
it.effect("deletes expired entries in batches", () =>
Effect.gen(function*() {
const sql = (yield* SqlClient.SqlClient).withoutTransforms()
const table = sql("effect_persistence")
const expiredCount = sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM ${table} WHERE store_id = 'expired'
`.pipe(Effect.map((rows) => Number(rows[0].count)))
yield* sql`
CREATE TABLE ${table} (
store_id VARCHAR(191) NOT NULL,
id VARCHAR(191) NOT NULL,
value TEXT NOT NULL,
expires BIGINT,
PRIMARY KEY (store_id, id)
)
`

const entries = Array.from({ length: SqlCleanupTest.expiredEntryCount }, (_, i) => ({
store_id: "expired",
id: String(i),
value: "{}",
expires: 0
}))
yield* sql`INSERT INTO ${table} ${sql.insert(entries)}`.unprepared
yield* sql`
INSERT INTO ${table} (store_id, id, value, expires)
VALUES ('live', 'live', '{}', NULL), ('live', 'future', '{}', 1000000)
`

yield* Layer.build(Persistence.layerBackingSql)

let expired = yield* SqlCleanupTest.waitForCount(
expiredCount,
(count) => count < SqlCleanupTest.expiredEntryCount
)
assert.strictEqual(expired, 1)

while (expired > 0) {
const previous = expired
yield* TestClock.adjust(SqlCleanupTest.cleanupBatchDelay)
expired = yield* SqlCleanupTest.waitForCount(expiredCount, (count) => count < previous)
}
assert.strictEqual(expired, 0)
const live = yield* sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM ${table} WHERE store_id = 'live'
`
assert.strictEqual(Number(live[0].count), 2)

const indexes = yield* sql<{ readonly count: number }>`
SELECT COUNT(*) AS count FROM information_schema.statistics
WHERE table_schema = DATABASE()
AND table_name = 'effect_persistence'
AND index_name = 'effect_persistence_expires_idx'
`
assert.strictEqual(Number(indexes[0].count), 1)
}), { timeout: SqlCleanupTest.testTimeout })
})
59 changes: 59 additions & 0 deletions packages/sql/pg/test/Persistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { assert, it } from "@effect/vitest"
import { Effect, Exit, Fiber, Latch, Layer, Schema } from "effect"
import * as PersistedCacheTest from "effect-test/unstable/persistence/PersistedCacheTest"
import * as PersistedQueueTest from "effect-test/unstable/persistence/PersistedQueueTest"
import * as SqlCleanupTest from "effect-test/unstable/persistence/SqlCleanupTest"
import { TestClock } from "effect/testing"
import { PersistedQueue, Persistence } from "effect/unstable/persistence"
import { SqlClient } from "effect/unstable/sql"
Expand Down Expand Up @@ -104,3 +105,61 @@ it.layer(PgContainer.layerClient, { timeout: "30 seconds" })("PersistedQueue SQL
assert.strictEqual(value, "valid")
}).pipe(TestClock.withLive))
})

it.layer(PgContainer.layerClient, { timeout: "30 seconds" })("Persistence SQL cleanup", (it) => {
it.effect("deletes expired entries in batches", () =>
Effect.gen(function*() {
const sql = (yield* SqlClient.SqlClient).withoutTransforms()
const table = sql("effect_persistence")
const expiredCount = sql<{ readonly count: number }>`
SELECT COUNT(*)::INT AS count FROM ${table} WHERE store_id = 'expired'
`.pipe(Effect.map((rows) => rows[0].count))
yield* sql`
CREATE TABLE ${table} (
store_id TEXT NOT NULL,
id TEXT NOT NULL,
value TEXT NOT NULL,
expires BIGINT,
PRIMARY KEY (store_id, id)
)
`

const entries = Array.from({ length: SqlCleanupTest.expiredEntryCount }, (_, i) => ({
store_id: "expired",
id: String(i),
value: "{}",
expires: 0
}))
yield* sql`INSERT INTO ${table} ${sql.insert(entries)}`.unprepared
yield* sql`
INSERT INTO ${table} (store_id, id, value, expires)
VALUES ('live', 'live', '{}', NULL), ('live', 'future', '{}', 1000000)
`

yield* Layer.build(Persistence.layerBackingSql)

let expired = yield* SqlCleanupTest.waitForCount(
expiredCount,
(count) => count < SqlCleanupTest.expiredEntryCount
)
assert.strictEqual(expired, 1)

while (expired > 0) {
const previous = expired
yield* TestClock.adjust(SqlCleanupTest.cleanupBatchDelay)
expired = yield* SqlCleanupTest.waitForCount(expiredCount, (count) => count < previous)
}
assert.strictEqual(expired, 0)
const live = yield* sql<{ readonly count: number }>`
SELECT COUNT(*)::INT AS count FROM ${table} WHERE store_id = 'live'
`
assert.strictEqual(live[0].count, 2)

const indexes = yield* sql<{ readonly count: number }>`
SELECT COUNT(*)::INT AS count FROM pg_indexes
WHERE tablename = 'effect_persistence'
AND indexname = 'effect_persistence_expires_idx'
`
assert.strictEqual(indexes[0].count, 1)
}), { timeout: SqlCleanupTest.testTimeout })
})
Loading