diff --git a/src/audit/AuditEventStream.ts b/src/audit/AuditEventStream.ts new file mode 100644 index 0000000..1015ec7 --- /dev/null +++ b/src/audit/AuditEventStream.ts @@ -0,0 +1,108 @@ +import type { AuditEvent } from '../types/audit.js'; + +export type AuditEventHandler = (event: AuditEvent) => void; + +export interface AuditStreamOptions { + /** Buffer size for event replay on new subscriber (default: 0 = no replay) */ + bufferSize?: number; + /** Filter function - only events passing this are emitted */ + filter?: (event: AuditEvent) => boolean; +} + +/** + * AuditEventStream provides a pub/sub stream for SDK audit events. + * Supports multiple subscribers, optional event buffering for replay, + * and per-subscriber filtering. + */ +export class AuditEventStream { + private readonly handlers = new Set(); + private readonly buffer: AuditEvent[] = []; + private readonly bufferSize: number; + private readonly filter?: (event: AuditEvent) => boolean; + private closed = false; + + constructor(options: AuditStreamOptions = {}) { + this.bufferSize = options.bufferSize ?? 0; + this.filter = options.filter; + } + + /** + * Publish an audit event to all subscribers. + * If the stream is closed, the event is silently dropped. + */ + publish(event: AuditEvent): void { + if (this.closed) return; + if (this.filter && !this.filter(event)) return; + + if (this.bufferSize > 0) { + this.buffer.push(event); + if (this.buffer.length > this.bufferSize) { + this.buffer.shift(); + } + } + + for (const handler of this.handlers) { + try { + handler(event); + } catch { + // subscriber errors must not break the stream + } + } + } + + /** + * Subscribe to audit events. Replays buffered events immediately if any. + * Returns an unsubscribe function. + */ + subscribe(handler: AuditEventHandler): () => void { + this.handlers.add(handler); + // Replay buffered events + for (const event of this.buffer) { + try { + handler(event); + } catch { + // ignore replay errors + } + } + return () => { + this.handlers.delete(handler); + }; + } + + /** + * Returns the number of active subscribers. + */ + get subscriberCount(): number { + return this.handlers.size; + } + + /** + * Returns a copy of the current replay buffer. + */ + getBuffer(): AuditEvent[] { + return [...this.buffer]; + } + + /** + * Clears all buffered events. + */ + clearBuffer(): void { + this.buffer.length = 0; + } + + /** + * Close the stream, preventing any further event publishing. + * All existing subscribers are removed. + */ + close(): void { + this.closed = true; + this.handlers.clear(); + } + + /** + * Whether the stream has been closed. + */ + get isClosed(): boolean { + return this.closed; + } +} diff --git a/test/auditEventStream.test.ts b/test/auditEventStream.test.ts new file mode 100644 index 0000000..eb48d26 --- /dev/null +++ b/test/auditEventStream.test.ts @@ -0,0 +1,132 @@ +import { describe, it, expect, vi } from 'vitest'; +import { AuditEventStream } from '../src/audit/AuditEventStream'; +import type { AuditEvent } from '../src/types/audit'; + +const makeEvent = (action: string, id = 'inv-1'): AuditEvent => ({ + invoiceId: id, + actorId: 'actor-1', + action, + payload: { value: 42 }, + timestamp: Date.now(), +}); + +describe('AuditEventStream', () => { + it('delivers published events to subscribers', () => { + const stream = new AuditEventStream(); + const received: AuditEvent[] = []; + stream.subscribe((e) => received.push(e)); + + const evt = makeEvent('CREATE'); + stream.publish(evt); + + expect(received).toHaveLength(1); + expect(received[0]).toEqual(evt); + }); + + it('supports multiple subscribers', () => { + const stream = new AuditEventStream(); + const a: AuditEvent[] = []; + const b: AuditEvent[] = []; + stream.subscribe((e) => a.push(e)); + stream.subscribe((e) => b.push(e)); + + stream.publish(makeEvent('UPDATE_STATUS')); + + expect(a).toHaveLength(1); + expect(b).toHaveLength(1); + }); + + it('unsubscribe stops delivery', () => { + const stream = new AuditEventStream(); + const received: AuditEvent[] = []; + const unsub = stream.subscribe((e) => received.push(e)); + + stream.publish(makeEvent('CREATE')); + unsub(); + stream.publish(makeEvent('REFUND')); + + expect(received).toHaveLength(1); + }); + + it('buffers events and replays to new subscribers', () => { + const stream = new AuditEventStream({ bufferSize: 5 }); + stream.publish(makeEvent('CREATE')); + stream.publish(makeEvent('UPDATE_STATUS')); + + const replayed: AuditEvent[] = []; + stream.subscribe((e) => replayed.push(e)); + + expect(replayed).toHaveLength(2); + }); + + it('respects bufferSize cap - oldest events are evicted', () => { + const stream = new AuditEventStream({ bufferSize: 2 }); + stream.publish(makeEvent('A')); + stream.publish(makeEvent('B')); + stream.publish(makeEvent('C')); + + expect(stream.getBuffer()).toHaveLength(2); + expect(stream.getBuffer()[0]?.action).toBe('B'); + expect(stream.getBuffer()[1]?.action).toBe('C'); + }); + + it('filter option prevents non-matching events', () => { + const stream = new AuditEventStream({ + filter: (e) => e.action === 'CREATE', + }); + const received: AuditEvent[] = []; + stream.subscribe((e) => received.push(e)); + + stream.publish(makeEvent('CREATE')); + stream.publish(makeEvent('REFUND')); + + expect(received).toHaveLength(1); + expect(received[0]?.action).toBe('CREATE'); + }); + + it('closed stream drops events', () => { + const stream = new AuditEventStream(); + const received: AuditEvent[] = []; + stream.subscribe((e) => received.push(e)); + + stream.close(); + stream.publish(makeEvent('CREATE')); + + expect(received).toHaveLength(0); + expect(stream.isClosed).toBe(true); + }); + + it('subscriber count tracks correctly', () => { + const stream = new AuditEventStream(); + expect(stream.subscriberCount).toBe(0); + + const unsub1 = stream.subscribe(vi.fn()); + const unsub2 = stream.subscribe(vi.fn()); + expect(stream.subscriberCount).toBe(2); + + unsub1(); + expect(stream.subscriberCount).toBe(1); + + unsub2(); + expect(stream.subscriberCount).toBe(0); + }); + + it('clearBuffer empties the replay buffer', () => { + const stream = new AuditEventStream({ bufferSize: 10 }); + stream.publish(makeEvent('CREATE')); + expect(stream.getBuffer()).toHaveLength(1); + + stream.clearBuffer(); + expect(stream.getBuffer()).toHaveLength(0); + }); + + it('subscriber errors do not crash the stream', () => { + const stream = new AuditEventStream(); + const good: AuditEvent[] = []; + stream.subscribe(() => { throw new Error('boom'); }); + stream.subscribe((e) => good.push(e)); + + stream.publish(makeEvent('CREATE')); + expect(good).toHaveLength(1); + }); +});