diff --git a/Makefile b/Makefile index 92f044ed..192a16e6 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,4 @@ -.PHONY: dev backend frontend infra stop clean seed test lint setup docs sync pr ship merge branches worktree-clean branch-clean +.PHONY: dev backend frontend infra events stop clean seed test lint setup docs sync pr ship merge branches worktree-clean branch-clean # --- First time setup --- setup: @@ -11,6 +11,12 @@ infra: @until docker compose exec -T postgres pg_isready -U tasktime >/dev/null 2>&1; do sleep 1; done @echo "PostgreSQL ready on :5432, Redis ready on :6379" +events: + docker compose --profile events up -d postgres redis kafka backend-relay + @echo "Waiting for Kafka..." + @until docker compose --profile events exec -T kafka kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1; do sleep 1; done + @echo "Kafka ready on :9092, backend-relay started" + # --- Dev servers --- backend: infra cd backend && npm run dev diff --git a/backend/package-lock.json b/backend/package-lock.json index 880ff583..a1032311 100644 --- a/backend/package-lock.json +++ b/backend/package-lock.json @@ -21,6 +21,7 @@ "express": "^4.21.2", "helmet": "^8.0.0", "jsonwebtoken": "^9.0.2", + "kafkajs": "^2.2.4", "node-cron": "^4.2.1", "prom-client": "^15.1.3", "redis": "^5.11.0", @@ -6092,6 +6093,15 @@ "safe-buffer": "^5.0.1" } }, + "node_modules/kafkajs": { + "version": "2.2.4", + "resolved": "https://registry.npmjs.org/kafkajs/-/kafkajs-2.2.4.tgz", + "integrity": "sha512-j/YeapB1vfPT2iOIUn/vxdyKEuhuY2PxMBvf5JWux6iSaukAccrMtXEY/Lb7OvavDhOWME589bpLrEdnVHjfjA==", + "license": "MIT", + "engines": { + "node": ">=14.0.0" + } + }, "node_modules/keyv": { "version": "4.5.4", "resolved": "https://registry.npmjs.org/keyv/-/keyv-4.5.4.tgz", diff --git a/backend/package.json b/backend/package.json index 5ce059c6..8f470631 100644 --- a/backend/package.json +++ b/backend/package.json @@ -10,6 +10,10 @@ "build": "tsc", "typecheck": "tsc --noEmit", "start": "node dist/server.js", + "events:relay": "node dist/relay.js", + "events:relay:dev": "tsx src/relay.ts", + "events:tail": "node dist/scripts/events-tail.js", + "events:tail:dev": "tsx src/scripts/events-tail.ts", "user:promote-super-admin": "tsx src/scripts/promote-super-admin.ts", "user:rotate-password": "tsx src/scripts/rotate-password.ts", "db:bootstrap": "node dist/prisma/bootstrap.js", @@ -45,6 +49,7 @@ "express": "^4.21.2", "helmet": "^8.0.0", "jsonwebtoken": "^9.0.2", + "kafkajs": "^2.2.4", "node-cron": "^4.2.1", "prom-client": "^15.1.3", "redis": "^5.11.0", diff --git a/backend/src/config.ts b/backend/src/config.ts index 1ab0a4a9..16a24c6b 100644 --- a/backend/src/config.ts +++ b/backend/src/config.ts @@ -1,6 +1,8 @@ import 'dotenv/config'; import { z } from 'zod'; +const booleanFlag = z.enum(['true', 'false']).default('false').transform((value) => value === 'true'); + const envSchema = z.object({ DATABASE_URL: z.string(), JWT_SECRET: z.string().min(10), @@ -30,6 +32,16 @@ const envSchema = z.object({ BURNDOWN_RETENTION_CRON: z.string().default('0 3 * * 0'), BURNDOWN_RETENTION_DAYS_AFTER_DONE: z.coerce.number().min(7).max(3650).default(90), BURNDOWN_WEEKLY_AGG_AFTER_DAYS: z.coerce.number().min(30).max(3650).default(365), + + // TTBUS-0: Kafka event bus + transactional outbox relay. + NOTIFICATIONS_ENABLED: booleanFlag, + KAFKA_BROKERS: z.string().default('localhost:9092'), + KAFKA_CLIENT_ID: z.string().min(1).default('tasktime-backend'), + OUTBOX_RELAY_INTERVAL_MS: z.coerce.number().min(100).max(60000).default(500), + OUTBOX_RELAY_BATCH_SIZE: z.coerce.number().min(1).max(500).default(100), + OUTBOX_RELAY_MAX_ATTEMPTS: z.coerce.number().min(1).max(100).default(10), + OUTBOX_RELAY_TRANSACTION_TIMEOUT_MS: z.coerce.number().min(1000).max(300000).default(60000), + OUTBOX_CLEANUP_RETENTION_DAYS: z.coerce.number().min(1).max(365).default(7), }); export const config = envSchema.parse(process.env); diff --git a/backend/src/relay.ts b/backend/src/relay.ts new file mode 100644 index 00000000..75e6fcb1 --- /dev/null +++ b/backend/src/relay.ts @@ -0,0 +1,36 @@ +import { config } from './config.js'; +import { createEventBusProducer } from './shared/eventbus/producer.js'; +import { PrismaOutboxRelayStore } from './shared/outbox/relay-store.js'; +import { OutboxRelayWorker } from './shared/outbox/relay-worker.js'; + +if (!config.NOTIFICATIONS_ENABLED) { + console.log('Outbox relay disabled: NOTIFICATIONS_ENABLED=false'); + process.exit(0); +} + +const worker = new OutboxRelayWorker({ + store: new PrismaOutboxRelayStore(config.OUTBOX_RELAY_TRANSACTION_TIMEOUT_MS), + producer: createEventBusProducer(), + intervalMs: config.OUTBOX_RELAY_INTERVAL_MS, + batchSize: config.OUTBOX_RELAY_BATCH_SIZE, + maxAttempts: config.OUTBOX_RELAY_MAX_ATTEMPTS, + cleanupRetentionDays: config.OUTBOX_CLEANUP_RETENTION_DAYS, +}); + +async function shutdown(signal: string) { + console.log(`Outbox relay received ${signal}, shutting down`); + await worker.stop(); + process.exit(0); +} + +process.once('SIGINT', () => { + void shutdown('SIGINT'); +}); +process.once('SIGTERM', () => { + void shutdown('SIGTERM'); +}); + +worker.start().catch((err) => { + console.error('Outbox relay failed to start', err); + process.exit(1); +}); diff --git a/backend/src/scripts/events-tail.ts b/backend/src/scripts/events-tail.ts new file mode 100644 index 00000000..7c1ada8d --- /dev/null +++ b/backend/src/scripts/events-tail.ts @@ -0,0 +1,53 @@ +import { createKafkaClient } from '../shared/eventbus/client.js'; +import { eventEnvelopeSchema } from '../shared/eventbus/envelope.js'; + +const topics = process.argv.slice(2); + +if (topics.length === 0) { + console.error('Usage: npm run events:tail -- [topic...]'); + process.exit(1); +} + +const groupId = `tasktime-events-tail-${Date.now()}`; +const kafka = createKafkaClient({ clientId: process.env.KAFKA_CLIENT_ID ?? 'tasktime-events-tail' }); +const consumer = kafka.consumer({ groupId }); + +async function main() { + await consumer.connect(); + for (const topic of topics) { + await consumer.subscribe({ topic, fromBeginning: false }); + } + + console.log(`Tailing Kafka topics: ${topics.join(', ')}`); + await consumer.run({ + eachMessage: async ({ topic, partition, message }) => { + if (!message.value) { + console.log(JSON.stringify({ topic, partition, offset: message.offset, empty: true })); + return; + } + + const parsed = JSON.parse(message.value.toString()); + const envelope = eventEnvelopeSchema.parse(parsed); + console.log(JSON.stringify({ topic, partition, offset: message.offset, envelope })); + }, + }); +} + +async function shutdown(signal: string) { + console.error(`events:tail received ${signal}, shutting down`); + await consumer.disconnect(); + process.exit(0); +} + +process.once('SIGINT', () => { + void shutdown('SIGINT'); +}); +process.once('SIGTERM', () => { + void shutdown('SIGTERM'); +}); + +main().catch(async (err) => { + console.error(err); + await consumer.disconnect().catch(() => undefined); + process.exit(1); +}); diff --git a/backend/src/shared/eventbus/client.ts b/backend/src/shared/eventbus/client.ts new file mode 100644 index 00000000..9fc87faa --- /dev/null +++ b/backend/src/shared/eventbus/client.ts @@ -0,0 +1,22 @@ +import { Kafka, logLevel, type KafkaConfig } from 'kafkajs'; + +export function parseKafkaBrokers(value: string): string[] { + return value + .split(',') + .map((broker) => broker.trim()) + .filter(Boolean); +} + +export function createKafkaClient(overrides: Partial = {}): Kafka { + const brokers = overrides.brokers ?? parseKafkaBrokers(process.env.KAFKA_BROKERS ?? 'localhost:9092'); + if (brokers.length === 0) { + throw new Error('KAFKA_BROKERS must contain at least one broker'); + } + + return new Kafka({ + clientId: process.env.KAFKA_CLIENT_ID ?? 'tasktime-backend', + brokers, + logLevel: logLevel.INFO, + ...overrides, + }); +} diff --git a/backend/src/shared/eventbus/consumer.ts b/backend/src/shared/eventbus/consumer.ts new file mode 100644 index 00000000..58397a0e --- /dev/null +++ b/backend/src/shared/eventbus/consumer.ts @@ -0,0 +1,92 @@ +import type { Consumer, EachMessagePayload, Kafka, KafkaMessage } from 'kafkajs'; + +import { createKafkaClient } from './client.js'; +import { type EventEnvelope, eventEnvelopeSchema } from './envelope.js'; +import { hasProcessedMessage, markProcessedOnce } from '../outbox/processed-messages.service.js'; + +export type EventHandler> = ( + envelope: EventEnvelope, + payload: T, + kafka: Pick & { offset: string }, +) => Promise; + +export type ProcessKafkaMessageOptions> = { + consumerGroup: string; + message: Pick; + topic: string; + partition: number; + handler: EventHandler; +}; + +export async function processKafkaMessage>({ + consumerGroup, + message, + topic, + partition, + handler, +}: ProcessKafkaMessageOptions): Promise<'processed' | 'skipped'> { + if (!message.value) { + throw new Error(`Kafka message on ${topic}[${partition}] offset ${message.offset} has no value`); + } + + const envelope = eventEnvelopeSchema.parse(JSON.parse(message.value.toString())); + if (await hasProcessedMessage(consumerGroup, envelope.messageId)) { + return 'skipped'; + } + + await handler(envelope, envelope.payload as T, { topic, partition, offset: message.offset }); + await markProcessedOnce(consumerGroup, envelope.messageId); + return 'processed'; +} + +export type EventBusConsumerOptions> = { + consumerGroup: string; + topics: string[]; + handler: EventHandler; + kafka?: Kafka; + consumer?: Consumer; + fromBeginning?: boolean; +}; + +export class EventBusConsumer> { + private readonly consumer: Consumer; + private readonly consumerGroup: string; + private readonly topics: string[]; + private readonly handler: EventHandler; + private readonly fromBeginning: boolean; + private connected = false; + + constructor(options: EventBusConsumerOptions) { + this.consumerGroup = options.consumerGroup; + this.topics = options.topics; + this.handler = options.handler; + this.fromBeginning = options.fromBeginning ?? false; + this.consumer = options.consumer ?? (options.kafka ?? createKafkaClient()).consumer({ groupId: options.consumerGroup }); + } + + async connectAndRun(): Promise { + if (this.connected) return; + await this.consumer.connect(); + for (const topic of this.topics) { + await this.consumer.subscribe({ topic, fromBeginning: this.fromBeginning }); + } + await this.consumer.run({ + eachMessage: async ({ topic, partition, message }) => { + await processKafkaMessage({ + consumerGroup: this.consumerGroup, + message, + topic, + partition, + handler: this.handler, + }); + }, + }); + this.connected = true; + } + + async disconnect(): Promise { + if (!this.connected) return; + await this.consumer.disconnect(); + this.connected = false; + } +} diff --git a/backend/src/shared/eventbus/producer.ts b/backend/src/shared/eventbus/producer.ts new file mode 100644 index 00000000..b0a8bf5b --- /dev/null +++ b/backend/src/shared/eventbus/producer.ts @@ -0,0 +1,54 @@ +import type { Kafka, Message, Producer } from 'kafkajs'; + +import { createKafkaClient } from './client.js'; +import { type EventEnvelope, eventEnvelopeSchema } from './envelope.js'; + +export type EventBusProducerOptions = { + kafka?: Kafka; + producer?: Producer; +}; + +export class EventBusProducer { + private readonly producer: Producer; + private connected = false; + + constructor(options: EventBusProducerOptions = {}) { + this.producer = options.producer ?? (options.kafka ?? createKafkaClient()).producer(); + } + + async connect(): Promise { + if (this.connected) return; + await this.producer.connect(); + this.connected = true; + } + + async disconnect(): Promise { + if (!this.connected) return; + await this.producer.disconnect(); + this.connected = false; + } + + async sendEnvelope(topic: string, envelope: EventEnvelope): Promise { + const validEnvelope = eventEnvelopeSchema.parse(envelope); + await this.connect(); + + const message: Message = { + key: validEnvelope.messageId, + value: JSON.stringify(validEnvelope), + headers: { + type: validEnvelope.type, + v: String(validEnvelope.v), + occurredAt: validEnvelope.occurredAt, + }, + }; + + await this.producer.send({ + topic, + messages: [message], + }); + } +} + +export function createEventBusProducer(options?: EventBusProducerOptions): EventBusProducer { + return new EventBusProducer(options); +} diff --git a/backend/src/shared/outbox/processed-messages.service.ts b/backend/src/shared/outbox/processed-messages.service.ts index 75d96890..370789f4 100644 --- a/backend/src/shared/outbox/processed-messages.service.ts +++ b/backend/src/shared/outbox/processed-messages.service.ts @@ -26,3 +26,14 @@ export async function markProcessedOnce( throw err; } } + +export async function hasProcessedMessage( + consumerGroup: string, + messageId: string, + client: PrismaLike = prisma, +): Promise { + const count = await client.processedMessage.count({ + where: { consumerGroup, messageId }, + }); + return count > 0; +} diff --git a/backend/src/shared/outbox/relay-store.ts b/backend/src/shared/outbox/relay-store.ts new file mode 100644 index 00000000..3094daff --- /dev/null +++ b/backend/src/shared/outbox/relay-store.ts @@ -0,0 +1,87 @@ +import { prisma } from '../../prisma/client.js'; + +export type OutboxRelayMessage = { + id: string; + topic: string; + messageId: string; + envelope: unknown; + attempts: number; +}; + +export type OutboxRelayBatchActions = { + markSent(id: string): Promise; + markFailed(id: string, error: unknown): Promise; +}; + +export type OutboxRelayStore = { + withPendingBatch( + limit: number, + maxAttempts: number, + handler: (messages: OutboxRelayMessage[], actions: OutboxRelayBatchActions) => Promise, + ): Promise; + cleanupSentBefore(cutoff: Date): Promise; +}; + +function formatRelayError(error: unknown): string { + const message = error instanceof Error ? error.message : String(error); + return message.slice(0, 2000); +} + +export class PrismaOutboxRelayStore implements OutboxRelayStore { + constructor(private readonly transactionTimeoutMs: number) {} + + async withPendingBatch( + limit: number, + maxAttempts: number, + handler: (messages: OutboxRelayMessage[], actions: OutboxRelayBatchActions) => Promise, + ): Promise { + return prisma.$transaction( + async (tx) => { + const messages = await tx.$queryRaw` + SELECT + id, + topic, + message_id AS "messageId", + envelope, + attempts + FROM event_outbox + WHERE sent_at IS NULL + AND attempts < ${maxAttempts} + ORDER BY created_at ASC + LIMIT ${limit} + FOR UPDATE SKIP LOCKED + `; + + const actions: OutboxRelayBatchActions = { + markSent: async (id: string) => { + await tx.eventOutbox.update({ + where: { id }, + data: { sentAt: new Date(), lastError: null }, + }); + }, + markFailed: async (id: string, error: unknown) => { + await tx.eventOutbox.update({ + where: { id }, + data: { + attempts: { increment: 1 }, + lastError: formatRelayError(error), + }, + }); + }, + }; + + return handler(messages, actions); + }, + { timeout: this.transactionTimeoutMs }, + ); + } + + async cleanupSentBefore(cutoff: Date): Promise { + const result = await prisma.eventOutbox.deleteMany({ + where: { + sentAt: { lt: cutoff }, + }, + }); + return result.count; + } +} diff --git a/backend/src/shared/outbox/relay-worker.ts b/backend/src/shared/outbox/relay-worker.ts new file mode 100644 index 00000000..541ac3e2 --- /dev/null +++ b/backend/src/shared/outbox/relay-worker.ts @@ -0,0 +1,96 @@ +import { eventEnvelopeSchema, type EventEnvelope } from '../eventbus/envelope.js'; +import type { OutboxRelayStore } from './relay-store.js'; + +export type OutboxRelayProducer = { + connect?(): Promise; + disconnect?(): Promise; + sendEnvelope(topic: string, envelope: EventEnvelope): Promise; +}; + +export type OutboxRelayStats = { + scanned: number; + sent: number; + failed: number; +}; + +export type OutboxRelayWorkerOptions = { + store: OutboxRelayStore; + producer: OutboxRelayProducer; + intervalMs: number; + batchSize: number; + maxAttempts: number; + cleanupRetentionDays: number; + now?: () => Date; +}; + +export class OutboxRelayWorker { + private timer: NodeJS.Timeout | null = null; + private cleanupTimer: NodeJS.Timeout | null = null; + private running = false; + + constructor(private readonly options: OutboxRelayWorkerOptions) {} + + async runOnce(): Promise { + return this.options.store.withPendingBatch(this.options.batchSize, this.options.maxAttempts, async (messages, actions) => { + const stats: OutboxRelayStats = { scanned: messages.length, sent: 0, failed: 0 }; + + for (const message of messages) { + try { + const envelope = eventEnvelopeSchema.parse(message.envelope); + await this.options.producer.sendEnvelope(message.topic, envelope); + await actions.markSent(message.id); + stats.sent += 1; + } catch (err) { + await actions.markFailed(message.id, err); + stats.failed += 1; + } + } + + return stats; + }); + } + + async cleanupSent(): Promise { + const now = this.options.now?.() ?? new Date(); + const cutoff = new Date(now.getTime() - this.options.cleanupRetentionDays * 24 * 60 * 60 * 1000); + return this.options.store.cleanupSentBefore(cutoff); + } + + async start(): Promise { + if (this.timer) return; + await this.options.producer.connect?.(); + this.timer = setInterval(() => { + void this.tick().catch((err) => { + console.error('Outbox relay tick failed', err); + }); + }, this.options.intervalMs); + this.cleanupTimer = setInterval(() => { + void this.cleanupSent().catch((err) => { + console.error('Outbox relay cleanup failed', err); + }); + }, 60 * 60 * 1000); + await this.tick(); + } + + async stop(): Promise { + if (this.timer) { + clearInterval(this.timer); + this.timer = null; + } + if (this.cleanupTimer) { + clearInterval(this.cleanupTimer); + this.cleanupTimer = null; + } + await this.options.producer.disconnect?.(); + } + + private async tick(): Promise { + if (this.running) return; + this.running = true; + try { + await this.runOnce(); + } finally { + this.running = false; + } + } +} diff --git a/backend/tests/eventbus-consumer.unit.test.ts b/backend/tests/eventbus-consumer.unit.test.ts new file mode 100644 index 00000000..51ca65fd --- /dev/null +++ b/backend/tests/eventbus-consumer.unit.test.ts @@ -0,0 +1,66 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; + +import type { EventEnvelope } from '../src/shared/eventbus/envelope.js'; +import { processKafkaMessage } from '../src/shared/eventbus/consumer.js'; +import * as processedMessages from '../src/shared/outbox/processed-messages.service.js'; + +const envelope: EventEnvelope = { + v: 1, + messageId: '97d6ed68-8bec-43d7-bd71-eb855e11498b', + type: 'ISSUE_CREATED', + occurredAt: '2026-04-28T00:00:00.000Z', + actor: { userId: null }, + payload: { issueId: 'issue-1' }, + meta: { tenantId: null }, +}; + +function message(value: unknown) { + return { + value: Buffer.from(JSON.stringify(value)), + offset: '42', + }; +} + +describe('processKafkaMessage', () => { + beforeEach(() => { + vi.restoreAllMocks(); + }); + + it('handles a valid envelope and marks it processed', async () => { + const handler = vi.fn().mockResolvedValue(undefined); + vi.spyOn(processedMessages, 'hasProcessedMessage').mockResolvedValue(false); + const mark = vi.spyOn(processedMessages, 'markProcessedOnce').mockResolvedValue(true); + + await expect( + processKafkaMessage({ + consumerGroup: 'notifications-service', + topic: 'tt.issues', + partition: 0, + message: message(envelope), + handler, + }), + ).resolves.toBe('processed'); + + expect(handler).toHaveBeenCalledWith(envelope, envelope.payload, { topic: 'tt.issues', partition: 0, offset: '42' }); + expect(mark).toHaveBeenCalledWith('notifications-service', envelope.messageId); + }); + + it('skips an already processed envelope', async () => { + const handler = vi.fn().mockResolvedValue(undefined); + vi.spyOn(processedMessages, 'hasProcessedMessage').mockResolvedValue(true); + const mark = vi.spyOn(processedMessages, 'markProcessedOnce').mockResolvedValue(true); + + await expect( + processKafkaMessage({ + consumerGroup: 'notifications-service', + topic: 'tt.issues', + partition: 0, + message: message(envelope), + handler, + }), + ).resolves.toBe('skipped'); + + expect(handler).not.toHaveBeenCalled(); + expect(mark).not.toHaveBeenCalled(); + }); +}); diff --git a/backend/tests/outbox-relay.unit.test.ts b/backend/tests/outbox-relay.unit.test.ts new file mode 100644 index 00000000..595cb640 --- /dev/null +++ b/backend/tests/outbox-relay.unit.test.ts @@ -0,0 +1,121 @@ +import { describe, expect, it } from 'vitest'; + +import type { EventEnvelope } from '../src/shared/eventbus/envelope.js'; +import type { OutboxRelayProducer } from '../src/shared/outbox/relay-worker.js'; +import { OutboxRelayWorker } from '../src/shared/outbox/relay-worker.js'; +import type { OutboxRelayBatchActions, OutboxRelayMessage, OutboxRelayStore } from '../src/shared/outbox/relay-store.js'; + +const validEnvelope = (messageId: string): EventEnvelope => ({ + v: 1, + messageId, + type: 'ISSUE_CREATED', + occurredAt: '2026-04-28T00:00:00.000Z', + actor: { userId: null }, + payload: { issueId: 'issue-1' }, + meta: { tenantId: null }, +}); + +class MemoryRelayStore implements OutboxRelayStore { + sent = new Set(); + failed = new Map(); + cleanupCutoff: Date | null = null; + + constructor(private readonly messages: OutboxRelayMessage[]) {} + + async withPendingBatch( + limit: number, + maxAttempts: number, + handler: (messages: OutboxRelayMessage[], actions: OutboxRelayBatchActions) => Promise, + ): Promise { + const batch = this.messages + .filter((message) => !this.sent.has(message.id) && !this.failed.has(message.id) && message.attempts < maxAttempts) + .slice(0, limit); + + return handler(batch, { + markSent: async (id) => { + this.sent.add(id); + }, + markFailed: async (id, error) => { + this.failed.set(id, error); + }, + }); + } + + async cleanupSentBefore(cutoff: Date): Promise { + this.cleanupCutoff = cutoff; + return 3; + } +} + +class MemoryProducer implements OutboxRelayProducer { + sent: Array<{ topic: string; envelope: EventEnvelope }> = []; + + constructor(private readonly failingMessageId?: string) {} + + async sendEnvelope(topic: string, envelope: EventEnvelope): Promise { + if (envelope.messageId === this.failingMessageId) { + throw new Error('Kafka unavailable'); + } + this.sent.push({ topic, envelope }); + } +} + +function makeWorker(store: OutboxRelayStore, producer: OutboxRelayProducer, now = () => new Date('2026-04-28T00:00:00.000Z')) { + return new OutboxRelayWorker({ + store, + producer, + intervalMs: 500, + batchSize: 100, + maxAttempts: 10, + cleanupRetentionDays: 7, + now, + }); +} + +describe('OutboxRelayWorker', () => { + it('publishes valid pending messages and marks them sent', async () => { + const store = new MemoryRelayStore([ + { id: 'row-1', topic: 'tt.issues', messageId: 'm1', envelope: validEnvelope('475161a7-db66-4f70-8a22-4b737a6a6d28'), attempts: 0 }, + ]); + const producer = new MemoryProducer(); + + const stats = await makeWorker(store, producer).runOnce(); + + expect(stats).toEqual({ scanned: 1, sent: 1, failed: 0 }); + expect(producer.sent).toHaveLength(1); + expect(store.sent.has('row-1')).toBe(true); + }); + + it('marks a row failed when publishing throws', async () => { + const messageId = 'fc3fa374-7e6e-4fe1-a6cf-f6db38244bc5'; + const store = new MemoryRelayStore([ + { id: 'row-1', topic: 'tt.issues', messageId, envelope: validEnvelope(messageId), attempts: 0 }, + ]); + const producer = new MemoryProducer(messageId); + + const stats = await makeWorker(store, producer).runOnce(); + + expect(stats).toEqual({ scanned: 1, sent: 0, failed: 1 }); + expect(store.failed.get('row-1')).toBeInstanceOf(Error); + }); + + it('does not pick messages at max attempts', async () => { + const store = new MemoryRelayStore([ + { id: 'row-1', topic: 'tt.issues', messageId: 'm1', envelope: validEnvelope('475161a7-db66-4f70-8a22-4b737a6a6d28'), attempts: 10 }, + ]); + const producer = new MemoryProducer(); + + const stats = await makeWorker(store, producer).runOnce(); + + expect(stats).toEqual({ scanned: 0, sent: 0, failed: 0 }); + expect(producer.sent).toHaveLength(0); + }); + + it('calculates cleanup cutoff from retention days', async () => { + const store = new MemoryRelayStore([]); + + await expect(makeWorker(store, new MemoryProducer()).cleanupSent()).resolves.toBe(3); + + expect(store.cleanupCutoff?.toISOString()).toBe('2026-04-21T00:00:00.000Z'); + }); +}); diff --git a/deploy/docker-compose.production.yml b/deploy/docker-compose.production.yml index 7e9f43b5..8f858291 100644 --- a/deploy/docker-compose.production.yml +++ b/deploy/docker-compose.production.yml @@ -31,6 +31,61 @@ services: retries: 5 start_period: 20s + kafka: + image: bitnami/kafka:3.7 + profiles: ["events"] + restart: unless-stopped + environment: + KAFKA_CFG_NODE_ID: 0 + KAFKA_CFG_PROCESS_ROLES: controller,broker + KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 + KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 + KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT + KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093 + KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE: "false" + ALLOW_PLAINTEXT_LISTENER: "yes" + volumes: + - kafka-data:/bitnami/kafka + healthcheck: + test: ["CMD-SHELL", "kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1"] + interval: 10s + timeout: 5s + retries: 10 + + kafka-topics-init: + image: bitnami/kafka:3.7 + profiles: ["events"] + command: + - /bin/bash + - -lc + - | + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.issues --partitions 3 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.comments --partitions 3 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.releases --partitions 3 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.workflow --partitions 3 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.notifications.dlq --partitions 3 --replication-factor 1 --config retention.ms=2592000000 + depends_on: + kafka: + condition: service_healthy + + backend-relay: + image: ${BACKEND_IMAGE:?BACKEND_IMAGE is required}:${IMAGE_TAG:?IMAGE_TAG is required} + command: ["npm", "run", "events:relay"] + profiles: ["events"] + restart: unless-stopped + env_file: + - ./env/backend.production.env + environment: + NOTIFICATIONS_ENABLED: "true" + KAFKA_BROKERS: kafka:9092 + KAFKA_CLIENT_ID: tasktime-backend-relay + depends_on: + postgres: + condition: service_healthy + kafka-topics-init: + condition: service_completed_successfully + postgres: image: postgres:16-alpine restart: unless-stopped @@ -120,3 +175,4 @@ volumes: postgres-data: redis-data: pipeline-postgres-data: + kafka-data: diff --git a/deploy/docker-compose.staging.yml b/deploy/docker-compose.staging.yml index f3be4a0a..dbf67098 100644 --- a/deploy/docker-compose.staging.yml +++ b/deploy/docker-compose.staging.yml @@ -31,6 +31,61 @@ services: retries: 5 start_period: 20s + kafka: + image: bitnami/kafka:3.7 + profiles: ["events"] + restart: unless-stopped + environment: + KAFKA_CFG_NODE_ID: 0 + KAFKA_CFG_PROCESS_ROLES: controller,broker + KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 + KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 + KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT + KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093 + KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE: "true" + ALLOW_PLAINTEXT_LISTENER: "yes" + volumes: + - kafka-data:/bitnami/kafka + healthcheck: + test: ["CMD-SHELL", "kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1"] + interval: 10s + timeout: 5s + retries: 10 + + kafka-topics-init: + image: bitnami/kafka:3.7 + profiles: ["events"] + command: + - /bin/bash + - -lc + - | + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.issues --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.comments --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.releases --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.workflow --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.notifications.dlq --partitions 1 --replication-factor 1 --config retention.ms=2592000000 + depends_on: + kafka: + condition: service_healthy + + backend-relay: + image: ${BACKEND_IMAGE:?BACKEND_IMAGE is required}:${IMAGE_TAG:-staging} + command: ["npm", "run", "events:relay"] + profiles: ["events"] + restart: unless-stopped + env_file: + - ./env/backend.staging.env + environment: + NOTIFICATIONS_ENABLED: "true" + KAFKA_BROKERS: kafka:9092 + KAFKA_CLIENT_ID: tasktime-backend-relay + depends_on: + postgres: + condition: service_healthy + kafka-topics-init: + condition: service_completed_successfully + postgres: image: postgres:16-alpine restart: unless-stopped @@ -120,3 +175,4 @@ volumes: postgres-data: redis-data: pipeline-postgres-data: + kafka-data: diff --git a/deploy/env/backend.production.env.example b/deploy/env/backend.production.env.example index 64ddbe7b..2361344d 100644 --- a/deploy/env/backend.production.env.example +++ b/deploy/env/backend.production.env.example @@ -27,3 +27,13 @@ ANTHROPIC_API_KEY= MCP_SERVICE_TOKEN= # MCP agent account password (agent@flow-universe.internal — for write operations via backend API) MCP_AGENT_PASSWORD= + +# TTBUS-0 event bus / transactional outbox relay. +NOTIFICATIONS_ENABLED=false +KAFKA_BROKERS=kafka:9092 +KAFKA_CLIENT_ID=tasktime-backend +OUTBOX_RELAY_INTERVAL_MS=500 +OUTBOX_RELAY_BATCH_SIZE=100 +OUTBOX_RELAY_MAX_ATTEMPTS=10 +OUTBOX_RELAY_TRANSACTION_TIMEOUT_MS=60000 +OUTBOX_CLEANUP_RETENTION_DAYS=7 diff --git a/deploy/env/backend.staging.env.example b/deploy/env/backend.staging.env.example index 866beac3..bb1e5524 100644 --- a/deploy/env/backend.staging.env.example +++ b/deploy/env/backend.staging.env.example @@ -23,3 +23,13 @@ MCP_AGENT_PASSWORD= # Leave blank in the example; the workflow writes `=true` on staging. FEATURES_ADVANCED_SEARCH= FEATURES_CHECKPOINT_TTQL= + +# TTBUS-0 event bus / transactional outbox relay. +NOTIFICATIONS_ENABLED=false +KAFKA_BROKERS=kafka:9092 +KAFKA_CLIENT_ID=tasktime-backend +OUTBOX_RELAY_INTERVAL_MS=500 +OUTBOX_RELAY_BATCH_SIZE=100 +OUTBOX_RELAY_MAX_ATTEMPTS=10 +OUTBOX_RELAY_TRANSACTION_TIMEOUT_MS=60000 +OUTBOX_CLEANUP_RETENTION_DAYS=7 diff --git a/docker-compose.yml b/docker-compose.yml index 465b5958..e0dbd51b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -27,6 +27,65 @@ services: timeout: 3s retries: 5 + kafka: + image: bitnami/kafka:3.7 + profiles: ["events"] + environment: + KAFKA_CFG_NODE_ID: 0 + KAFKA_CFG_PROCESS_ROLES: controller,broker + KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 + KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 + KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT + KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093 + KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE: "true" + ALLOW_PLAINTEXT_LISTENER: "yes" + ports: + - "9092:9092" + volumes: + - kafka-data:/bitnami/kafka + healthcheck: + test: ["CMD-SHELL", "kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1"] + interval: 10s + timeout: 5s + retries: 10 + + kafka-topics-init: + image: bitnami/kafka:3.7 + profiles: ["events"] + command: + - /bin/bash + - -lc + - | + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.issues --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.comments --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.releases --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.workflow --partitions 1 --replication-factor 1 --config retention.ms=604800000 + kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic tt.notifications.dlq --partitions 1 --replication-factor 1 --config retention.ms=2592000000 + depends_on: + kafka: + condition: service_healthy + + backend-relay: + build: + context: . + dockerfile: backend/Dockerfile + profiles: ["events"] + command: ["npm", "run", "events:relay"] + environment: + DATABASE_URL: postgresql://tasktime:${POSTGRES_PASSWORD:-tasktime}@postgres:5432/tasktime?schema=public + JWT_SECRET: ${JWT_SECRET:-dev-jwt-secret-change-me} + JWT_REFRESH_SECRET: ${JWT_REFRESH_SECRET:-dev-refresh-secret-change-me} + NOTIFICATIONS_ENABLED: "true" + KAFKA_BROKERS: kafka:9092 + KAFKA_CLIENT_ID: tasktime-backend-relay + depends_on: + postgres: + condition: service_healthy + kafka-topics-init: + condition: service_completed_successfully + restart: unless-stopped + # --- Pipeline Service DB --- pipeline-postgres: image: postgres:16-alpine @@ -79,3 +138,4 @@ services: volumes: pgdata: pipeline-pgdata: + kafka-data: diff --git a/docs/architecture/event-bus.md b/docs/architecture/event-bus.md new file mode 100644 index 00000000..57af4a88 --- /dev/null +++ b/docs/architecture/event-bus.md @@ -0,0 +1,101 @@ +# Event Bus Architecture + +TTBUS-0 introduces a Kafka-backed event bus with a PostgreSQL transactional outbox. Business modules write domain events into `event_outbox` in the same database transaction as the business change. A separate relay process publishes pending rows to Kafka and marks rows as sent only after Kafka accepts the message. + +## Envelope + +All Kafka messages use the versioned envelope from `backend/src/shared/eventbus/envelope.ts`: + +```ts +{ + v: 1, + messageId: "uuid", + type: "ISSUE_CREATED", + occurredAt: "2026-04-28T00:00:00.000Z", + actor: { userId: "uuid-or-null", ip?: "...", userAgent?: "..." }, + payload: {}, + meta: { tenantId: null, correlationId?: "..." } +} +``` + +`messageId` is both the Kafka message key and the idempotency key for consumers. + +## Topics + +Domain topics are intentionally split by business area: + +- `tt.issues` +- `tt.comments` +- `tt.releases` +- `tt.workflow` +- `tt.notifications.dlq` + +Dev and staging use one partition per topic. Production compose creates three partitions per topic and keeps business events for 7 days. The DLQ topic keeps events for 30 days. + +## Producer Flow + +Use `publishInTx(tx, topic, type, payload, actor)` from inside the same `prisma.$transaction(...)` callback as the business mutation. + +Payloads must not include secrets. `publishInTx` rejects secret-like keys such as `password`, `password_hash`, `token`, `access_token`, `refresh_token`, `api_key`, and `secret`, including nested fields. + +## Relay Flow + +The relay entrypoint is `backend/src/relay.ts`. + +Run locally: + +```bash +make events +``` + +Run after a backend build: + +```bash +cd backend +npm run events:relay +``` + +The relay: + +- reads unsent rows with `FOR UPDATE SKIP LOCKED`; +- skips rows whose `attempts` reached `OUTBOX_RELAY_MAX_ATTEMPTS`; +- publishes one Kafka message per outbox row; +- sets `sent_at` only after successful publish; +- increments `attempts` and stores `last_error` after failed publish. + +## Consumer Gotchas + +Consumers should use `EventBusConsumer` or `processKafkaMessage` so every message goes through envelope validation and `processed_messages` deduplication. + +Handler functions must be idempotent. The event bus is at-least-once: retries can publish or deliver a duplicate message, and the dedup table only guarantees one successful handler run per `(consumerGroup, messageId)`. + +The current consumer helper marks a message as processed after the handler resolves. If the handler performs side effects and then throws, Kafka can redeliver the message. Consumer handlers should either make side effects idempotent or persist their own progress before external calls. + +## Operations + +Kafka and the relay are opt-in in Docker Compose via the `events` profile. Kafka is not exposed outside the Docker network in staging or production. + +Useful environment variables: + +- `NOTIFICATIONS_ENABLED` +- `KAFKA_BROKERS` +- `KAFKA_CLIENT_ID` +- `OUTBOX_RELAY_INTERVAL_MS` +- `OUTBOX_RELAY_BATCH_SIZE` +- `OUTBOX_RELAY_MAX_ATTEMPTS` +- `OUTBOX_RELAY_TRANSACTION_TIMEOUT_MS` +- `OUTBOX_CLEANUP_RETENTION_DAYS` + +Tail topics after a backend build: + +```bash +cd backend +npm run events:tail -- tt.issues +``` + +During development: + +```bash +cd backend +npm run events:tail:dev -- tt.issues +```