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
4 changes: 4 additions & 0 deletions xstreamroll-processing/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@ POLL_INTERVAL_MS=5000
# Maximum number of stream sessions processed concurrently. Defaults to 32.
MAX_CONCURRENT_SESSIONS=32

# Durable local dead-letter file. Mount this path in production if records
# must survive worker replacement. GET /dead-letters exposes its contents.
DEAD_LETTER_STORE_PATH=./data/dead-letters.json

# ──────────────────────────────────────────────────────────────────────────────
# Worker identification (issue #347)
#
Expand Down
1 change: 1 addition & 0 deletions xstreamroll-processing/.gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -9,3 +9,4 @@ dist/
*.swp
*.swo
coverage/
data/dead-letters.json
62 changes: 62 additions & 0 deletions xstreamroll-processing/__tests__/dead-letter-store.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
import { mkdtemp, readFile, rm } from "fs/promises"
import { tmpdir } from "os"
import { join } from "path"

import { FileDeadLetterStore } from "../src/dead-letter-store"
import { StreamEvent } from "../src/session"

describe("FileDeadLetterStore", () => {
let directory: string
let store: FileDeadLetterStore

beforeEach(async () => {
directory = await mkdtemp(join(tmpdir(), "xstreamroll-dlq-"))
store = new FileDeadLetterStore(join(directory, "dead-letters.json"))
})

afterEach(async () => {
await rm(directory, { recursive: true, force: true })
})

it("persists a dead-letter and deduplicates the same event", async () => {
const event: StreamEvent = {
id: "event-1",
streamId: "stream-1",
data: { value: 1 },
timestamp: "2026-08-28T00:00:00.000Z",
}

await store.record(event, new Error("first failure"), 3)
await store.record(event, new Error("second failure"), 4)

const records = await store.list()
expect(records).toHaveLength(1)
expect(records[0]).toMatchObject({
key: "stream-1:event-1",
streamId: "stream-1",
event,
error: "second failure",
attemptCount: 4,
})
expect(JSON.parse(await readFile(join(directory, "dead-letters.json"), "utf8")))
.toHaveProperty("stream-1:event-1")
})

it("deduplicates events without an id using stable event content", async () => {
const event: StreamEvent = {
streamId: "stream-1",
data: { z: 2, a: 1 },
timestamp: "2026-08-28T00:00:00.000Z",
}
const sameEventWithDifferentKeyOrder: StreamEvent = {
streamId: "stream-1",
data: { a: 1, z: 2 },
timestamp: "2026-08-28T00:00:00.000Z",
}

await store.record(event, new Error("failure"), 1)
await store.record(sameEventWithDifferentKeyOrder, new Error("failure"), 1)

expect(await store.list()).toHaveLength(1)
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ test("filtered events are not published (integration)", async () => {
let sent = false
nock("http://mock-api")
.get("/streams/pending")
.query(true)
.times(100)
.reply(() => {
if (!sent) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ test("single event: polled -> session -> published", async () => {
let sent = false
nock("http://mock-api")
.get("/streams/pending")
.query(true)
.times(100)
.reply(() => {
if (!sent) {
Expand Down Expand Up @@ -86,6 +87,7 @@ test("multiple events same stream -> routed to same session", async () => {
let sentOnce = false
nock("http://mock-api")
.get("/streams/pending")
.query(true)
.times(100)
.reply(() => {
if (!sentOnce) {
Expand Down Expand Up @@ -127,6 +129,7 @@ test("capacity exceeded -> event dropped, not published", async () => {
let once = false
nock("http://mock-api")
.get("/streams/pending")
.query(true)
.times(100)
.reply(() => {
if (!once) {
Expand Down Expand Up @@ -164,6 +167,7 @@ test("graceful shutdown flushes pending publishes", async () => {
let sent = false
nock("http://mock-api")
.get("/streams/pending")
.query(true)
.times(100)
.reply(() => {
if (!sent) {
Expand All @@ -183,7 +187,7 @@ test("graceful shutdown flushes pending publishes", async () => {
return "ok"
})

const workerMod = await import("../../src/worker")
workerMod = await import("../../src/worker")

// give worker a moment to pick up the event
await new Promise((r) => setTimeout(r, 100))
Expand All @@ -202,6 +206,7 @@ test("api error then recovery -> worker retries next poll", async () => {
let calls = 0
nock("http://mock-api")
.get("/streams/pending")
.query(true)
.times(100)
.reply(() => {
calls++
Expand All @@ -219,7 +224,7 @@ test("api error then recovery -> worker retries next poll", async () => {
})
})

const workerMod = await import("../../src/worker")
workerMod = await import("../../src/worker")
await awaitWithTimeout(publishedPromise, 5000, "publish timeout")
await workerMod.shutdown("test")
})
5 changes: 5 additions & 0 deletions xstreamroll-processing/__tests__/session.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -267,6 +267,7 @@ describe("StreamSession — publish retry and dead-letter (issue #343)", () => {
it("dead-letters after exhausting retries and continues processing remaining queue", async () => {
const published: ProcessedStreamEvent[] = []
const deadLettered: StreamEvent[] = []
const deadLetterAttempts: number[] = []
let callsForFirst = 0

// First event always fails; second event always succeeds
Expand All @@ -281,6 +282,9 @@ describe("StreamSession — publish retry and dead-letter (issue #343)", () => {
}
published.push(e)
},
deadLetter: async (_event, _error, attempts) => {
deadLetterAttempts.push(attempts)
},
},
1000,
2,
Expand All @@ -306,6 +310,7 @@ describe("StreamSession — publish retry and dead-letter (issue #343)", () => {
expect(callsForFirst).toBe(3)
expect(deadLettered).toHaveLength(1)
expect(deadLettered[0].data).toEqual({ seq: 1 })
expect(deadLetterAttempts).toEqual([3])

// Second event was still published successfully
expect(published).toHaveLength(1)
Expand Down
1 change: 1 addition & 0 deletions xstreamroll-processing/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ const envSchema = z.object({
.default("3")
.transform((s) => Number(s))
.pipe(z.number().int().min(0)),
DEAD_LETTER_STORE_PATH: z.string().default("./data/dead-letters.json"),
/**
* Backend for the per-stream {@link EventFilter} config store
* (issue #351). `memory` keeps every config in-process and matches
Expand Down
103 changes: 103 additions & 0 deletions xstreamroll-processing/src/dead-letter-store.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
import { createHash } from "crypto"
import { mkdir, readFile, rename, writeFile } from "fs/promises"
import { dirname } from "path"

import { StreamEvent } from "./session"

export interface DeadLetterRecord {
key: string
streamId: string
event: StreamEvent
error: string
attemptCount: number
timestamp: string
}

export interface DeadLetterStore {
record(event: StreamEvent, error: Error, attempts: number): Promise<void>
list(): Promise<DeadLetterRecord[]>
}

/** File-backed dead-letter store used by the dependency-free worker. */
export class FileDeadLetterStore implements DeadLetterStore {
private pendingWrite: Promise<void> = Promise.resolve()

constructor(private readonly filePath: string) {}

async record(
event: StreamEvent,
error: Error,
attempts: number,
): Promise<void> {
this.pendingWrite = this.pendingWrite.catch(() => undefined).then(async () => {
const records = await this.readRecords()
const key = eventKey(event)
const existing = records[key]
records[key] = existing
? {
...existing,
error: error.message,
attemptCount: attempts,
timestamp: new Date().toISOString(),
}
: {
key,
streamId: event.streamId,
event,
error: error.message,
attemptCount: attempts,
timestamp: new Date().toISOString(),
}
await this.writeRecords(records)
})
return this.pendingWrite
}

async list(): Promise<DeadLetterRecord[]> {
await this.pendingWrite
const records = await this.readRecords()
return Object.values(records).sort((left, right) =>
left.timestamp.localeCompare(right.timestamp),
)
}

private async readRecords(): Promise<Record<string, DeadLetterRecord>> {
try {
const raw = await readFile(this.filePath, "utf8")
return JSON.parse(raw) as Record<string, DeadLetterRecord>
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return {}
throw error
}
}

private async writeRecords(
records: Record<string, DeadLetterRecord>,
): Promise<void> {
await mkdir(dirname(this.filePath), { recursive: true })
const temporaryPath = `${this.filePath}.${process.pid}.tmp`
await writeFile(temporaryPath, JSON.stringify(records, null, 2), "utf8")
await rename(temporaryPath, this.filePath)
}
}

function eventKey(event: StreamEvent): string {
if (event.id) return `${event.streamId}:${event.id}`
return createHash("sha256")
.update(stableStringify(event))
.digest("hex")
}

function stableStringify(value: unknown): string {
if (Array.isArray(value)) {
return `[${value.map((item) => stableStringify(item)).join(",")}]`
}
if (value !== null && typeof value === "object") {
const object = value as Record<string, unknown>
return `{${Object.keys(object)
.sort()
.map((key) => `${JSON.stringify(key)}:${stableStringify(object[key])}`)
.join(",")}}`
}
return JSON.stringify(value) ?? "null"
}
7 changes: 6 additions & 1 deletion xstreamroll-processing/src/lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,12 @@
*/

export type ShutdownReason =
"SIGINT" | "SIGTERM" | "uncaughtException" | "unhandledRejection" | "manual"
| "SIGINT"
| "SIGTERM"
| "uncaughtException"
| "unhandledRejection"
| "manual"
| "test"

export interface ShutdownHook {
/** Human-readable name for logging. */
Expand Down
26 changes: 25 additions & 1 deletion xstreamroll-processing/src/metrics.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { createServer, IncomingMessage, ServerResponse } from "http"

import { DeadLetterStore } from "./dead-letter-store"

export interface Metrics {
messagesProcessed: number
errors: number
Expand Down Expand Up @@ -104,7 +106,10 @@ export function markReady(): void {
live = true
}

export function startMetricsServer(port = 3002): ReturnType<typeof createServer> {
export function startMetricsServer(
port = 3002,
deadLetterStore?: DeadLetterStore,
): ReturnType<typeof createServer> {
const server = createServer((req: IncomingMessage, res: ServerResponse) => {
const url = req.url || ""
const accept = (req.headers.accept || "").toLowerCase()
Expand Down Expand Up @@ -147,6 +152,25 @@ export function startMetricsServer(port = 3002): ReturnType<typeof createServer>
return
}

if (url === "/dead-letters" || url === "/dead-letters/") {
if (!deadLetterStore) {
res.writeHead(503, { "Content-Type": "application/json" })
res.end(JSON.stringify({ error: "dead-letter store unavailable" }))
return
}
void deadLetterStore.list().then(
(records) => {
res.writeHead(200, { "Content-Type": "application/json" })
res.end(JSON.stringify(records))
},
() => {
res.writeHead(500, { "Content-Type": "application/json" })
res.end(JSON.stringify({ error: "dead-letter store unavailable" }))
},
)
return
}

if (url === "/metrics" || url === "/metrics/" || url === "/metrics/prometheus") {
const wantsJson =
accept.includes("application/json") &&
Expand Down
Loading
Loading