From af0276410f537c46f6ad6912b8b027a10abbdd11 Mon Sep 17 00:00:00 2001 From: Guido X Jansen Date: Mon, 23 Feb 2026 21:17:22 +0100 Subject: [PATCH] fix(firehose): restore cursor, reconnection backoff, connected flag (#83) * fix(firehose): restore cursor, reconnection backoff, connected flag - Read saved cursor from CursorStore on startup and pass as query parameter to TapChannel URL (prevents event loss on restart) - Add exponential backoff reconnection (1s to 60s cap) when channel errors or closes unexpectedly - Defer connected=true until first event is successfully processed (was set prematurely before WebSocket established) - Graceful shutdown cancels pending reconnection timers - Guard against reconnection during shutdown with shuttingDown flag - Construct TapChannel directly to support cursor in WebSocket URL * fix(health): treat firehose startup as healthy in readiness check During startup (before the first event is processed), the firehose has connected=false and lastEventId=null. Treat this as healthy rather than unhealthy, since the service is still initializing. Only report unhealthy when the firehose was previously connected (lastEventId is set) but has since disconnected. --- src/firehose/service.ts | 146 +++++++-- src/routes/health.ts | 5 +- tests/unit/firehose/service.test.ts | 486 ++++++++++++++++++++++++++-- 3 files changed, 579 insertions(+), 58 deletions(-) diff --git a/src/firehose/service.ts b/src/firehose/service.ts index cdf1edf..a1d0fed 100644 --- a/src/firehose/service.ts +++ b/src/firehose/service.ts @@ -1,5 +1,4 @@ -import { Tap, SimpleIndexer } from '@atproto/tap' -import type { TapChannel } from '@atproto/tap' +import { Tap, SimpleIndexer, TapChannel } from '@atproto/tap' import type { RecordEvent as TapRecordEvent, IdentityEvent as TapIdentityEvent } from '@atproto/tap' import type { Database } from '../db/index.js' import type { Logger } from '../lib/logger.js' @@ -20,6 +19,9 @@ interface FirehoseStatus { lastEventId: number | null } +const MIN_BACKOFF_MS = 1000 +const MAX_BACKOFF_MS = 60_000 + export class FirehoseService { private tap: Tap private channel: TapChannel | null = null @@ -27,13 +29,17 @@ export class FirehoseService { private repoManager: RepoManager private recordHandler: RecordHandler private identityHandler: IdentityHandler + private indexer: SimpleIndexer | null = null private connected = false private lastEventId: number | null = null + private shuttingDown = false + private reconnectAttempts = 0 + private reconnectTimer: ReturnType | null = null constructor( db: Database, private logger: Logger, - env: Env + private env: Env ) { this.tap = new Tap(env.TAP_URL, { adminPassword: env.TAP_ADMIN_PASSWORD, @@ -60,11 +66,12 @@ export class FirehoseService { async start(): Promise { try { + this.shuttingDown = false await this.repoManager.restoreTrackedRepos() - const indexer = new SimpleIndexer() + this.indexer = new SimpleIndexer() - indexer.record(async (evt: TapRecordEvent) => { + this.indexer.record(async (evt: TapRecordEvent) => { const event: RecordEvent = { id: evt.id, action: evt.action, @@ -78,11 +85,10 @@ export class FirehoseService { } await this.recordHandler.handle(event) - this.lastEventId = evt.id - this.cursorStore.saveCursor(BigInt(evt.id)) + this.onEventProcessed(evt.id) }) - indexer.identity(async (evt: TapIdentityEvent) => { + this.indexer.identity(async (evt: TapIdentityEvent) => { const event: IdentityEvent = { id: evt.id, did: evt.did, @@ -92,26 +98,20 @@ export class FirehoseService { } await this.identityHandler.handle(event) - this.lastEventId = evt.id - this.cursorStore.saveCursor(BigInt(evt.id)) + this.onEventProcessed(evt.id) }) - indexer.error((err: Error) => { + this.indexer.error((err: Error) => { this.logger.error({ err }, 'Firehose indexer error') }) - this.channel = this.tap.channel(indexer) - - // channel.start() is a long-running loop over WebSocket messages that - // only resolves when the channel is destroyed. Run it as a background - // task so it does not block the Fastify onReady hook. - this.channel.start().catch((err: unknown) => { - this.logger.error({ err }, 'Firehose channel error') - this.connected = false - }) + const cursor = await this.cursorStore.getCursor() + this.startChannel(cursor) - this.connected = true - this.logger.info('Firehose subscription started') + this.logger.info( + { cursor: cursor !== null ? cursor.toString() : null }, + 'Firehose subscription started' + ) } catch (err) { this.logger.error({ err }, 'Failed to start firehose service') this.connected = false @@ -119,10 +119,18 @@ export class FirehoseService { } async stop(): Promise { + this.shuttingDown = true + + if (this.reconnectTimer !== null) { + clearTimeout(this.reconnectTimer) + this.reconnectTimer = null + } + if (this.channel) { await this.channel.destroy() this.channel = null } + await this.cursorStore.flush() this.connected = false this.logger.info('Firehose subscription stopped') @@ -138,4 +146,98 @@ export class FirehoseService { getRepoManager(): RepoManager { return this.repoManager } + + private onEventProcessed(id: number): void { + if (!this.connected) { + this.connected = true + this.logger.info('Firehose connection confirmed') + } + this.reconnectAttempts = 0 + this.lastEventId = id + this.cursorStore.saveCursor(BigInt(id)) + } + + private startChannel(cursor: bigint | null): void { + if (this.shuttingDown || !this.indexer) { + return + } + + this.channel = this.createChannel(cursor) + + this.channel + .start() + .then(() => { + this.connected = false + if (!this.shuttingDown) { + this.logger.warn('Firehose channel closed, scheduling reconnection') + this.scheduleReconnect() + } + }) + .catch((err: unknown) => { + this.connected = false + if (!this.shuttingDown) { + this.logger.error({ err }, 'Firehose channel error, scheduling reconnection') + this.scheduleReconnect() + } + }) + } + + private createChannel(cursor: bigint | null): TapChannel { + if (!this.indexer) { + throw new Error('Cannot create channel: indexer not initialized') + } + + const url = new URL(this.env.TAP_URL) + url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:' + url.pathname = '/channel' + if (cursor !== null) { + url.searchParams.set('cursor', cursor.toString()) + } + + return new TapChannel(url.toString(), this.indexer, { + adminPassword: this.env.TAP_ADMIN_PASSWORD, + }) + } + + private scheduleReconnect(): void { + if (this.shuttingDown) { + return + } + + this.reconnectAttempts++ + const backoffMs = Math.min( + MIN_BACKOFF_MS * Math.pow(2, this.reconnectAttempts - 1), + MAX_BACKOFF_MS + ) + + this.logger.info( + { attempt: this.reconnectAttempts, backoffMs }, + 'Scheduling firehose reconnection' + ) + + this.reconnectTimer = setTimeout(() => { + this.reconnectTimer = null + void this.attemptReconnect() + }, backoffMs) + } + + private async attemptReconnect(): Promise { + if (this.shuttingDown) { + return + } + + try { + const cursor = await this.cursorStore.getCursor() + + this.logger.info( + { attempt: this.reconnectAttempts, cursor: cursor?.toString() ?? null }, + 'Attempting firehose reconnection' + ) + + this.startChannel(cursor) + } catch (err) { + this.logger.error({ err, attempt: this.reconnectAttempts }, 'Firehose reconnection failed') + this.scheduleReconnect() + } + } } diff --git a/src/routes/health.ts b/src/routes/health.ts index a035bc7..a0a611f 100644 --- a/src/routes/health.ts +++ b/src/routes/health.ts @@ -38,9 +38,12 @@ export const healthRoutes: FastifyPluginCallback = (fastify, _opts, done) => { } // Check firehose + // During startup (no events processed yet), treat as healthy. + // Once events have been processed, require an active connection. const firehoseStatus = fastify.firehose.getStatus() + const firehoseHealthy = firehoseStatus.connected || firehoseStatus.lastEventId === null checks['firehose'] = { - status: firehoseStatus.connected ? 'healthy' : 'unhealthy', + status: firehoseHealthy ? 'healthy' : 'unhealthy', ...(firehoseStatus.lastEventId !== null ? { latency: firehoseStatus.lastEventId } : {}), } diff --git a/tests/unit/firehose/service.test.ts b/tests/unit/firehose/service.test.ts index 082a191..3866270 100644 --- a/tests/unit/firehose/service.test.ts +++ b/tests/unit/firehose/service.test.ts @@ -1,39 +1,84 @@ -import { describe, it, expect, vi, beforeEach } from 'vitest' +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest' import { FirehoseService } from '../../../src/firehose/service.js' import type { Env } from '../../../src/config/env.js' -// Mock the Tap and SimpleIndexer from @atproto/tap +// --- Hoisted mock state (available before vi.mock runs) --- + +const { mockTapChannelCtor, channelInstances, indexerInstances } = vi.hoisted(() => ({ + mockTapChannelCtor: vi.fn(), + channelInstances: [] as Array<{ + url: string + handler: unknown + opts: unknown + start: ReturnType + destroy: ReturnType + }>, + indexerInstances: [] as Array<{ + _recordHandler?: (evt: unknown) => Promise + _identityHandler?: (evt: unknown) => Promise + _errorHandler?: (err: Error) => void + }>, +})) + vi.mock('@atproto/tap', () => { - const mockChannel = { - start: vi.fn().mockResolvedValue(undefined), - destroy: vi.fn().mockResolvedValue(undefined), + class MockSimpleIndexer { + _recordHandler?: (evt: unknown) => Promise + _identityHandler?: (evt: unknown) => Promise + _errorHandler?: (err: Error) => void + + constructor() { + indexerInstances.push(this) + } + + record(fn: (evt: unknown) => Promise) { + this._recordHandler = fn + return this + } + identity(fn: (evt: unknown) => Promise) { + this._identityHandler = fn + return this + } + error(fn: (err: Error) => void) { + this._errorHandler = fn + return this + } + + onEvent = vi.fn() + onError = vi.fn() } class MockTap { addRepos = vi.fn().mockResolvedValue(undefined) removeRepos = vi.fn().mockResolvedValue(undefined) - channel = vi.fn().mockReturnValue(mockChannel) - } - - class MockSimpleIndexer { - identity = vi.fn().mockReturnThis() - record = vi.fn().mockReturnThis() - error = vi.fn().mockReturnThis() + channel = vi.fn() } return { Tap: MockTap, SimpleIndexer: MockSimpleIndexer, - _mockChannel: mockChannel, + TapChannel: mockTapChannelCtor, } }) +// --- Helpers --- + +/** + * Creates a mock Drizzle query chain that supports both: + * await db.select().from(table) -- thenable (restoreTrackedRepos) + * await db.select().from(table).where() -- chained (getCursor) + */ +function mockQueryChain(result: unknown[] = []) { + const promise = Promise.resolve(result) as Promise & { + where: ReturnType + } + promise.where = vi.fn().mockResolvedValue(result) + return promise +} + function createMockDb() { return { select: vi.fn().mockReturnValue({ - from: vi.fn().mockReturnValue({ - where: vi.fn().mockResolvedValue([]), - }), + from: vi.fn().mockReturnValue(mockQueryChain([])), }), insert: vi.fn().mockReturnValue({ values: vi.fn().mockReturnValue({ @@ -83,6 +128,63 @@ function createMinimalEnv(): Env { } } +/** Minimal record event with unsupported collection (handler returns immediately). */ +function createTestRecordEvent(id: number) { + return { + id, + type: 'record' as const, + action: 'create' as const, + did: 'did:plc:test', + rev: 'rev1', + collection: 'com.example.unsupported', + rkey: 'test', + record: undefined, + cid: undefined, + live: true, + } +} + +/** Default TapChannel mock: start() never resolves (channel stays alive). */ +function defaultChannelImpl( + this: Record, + url: string, + handler: unknown, + opts: unknown +) { + this.url = url + this.handler = handler + this.opts = opts + this.start = vi.fn().mockReturnValue(new Promise(() => {})) + this.destroy = vi.fn().mockResolvedValue(undefined) + channelInstances.push(this as (typeof channelInstances)[number]) +} + +/** Get the record handler from the first indexer instance, or throw. */ +function getRecordHandler() { + const handler = indexerInstances[0]._recordHandler + if (!handler) { + throw new Error('record handler not registered on indexer') + } + return handler +} + +/** TapChannel mock that rejects immediately on start(). */ +function failingChannelImpl( + this: Record, + url: string, + handler: unknown, + opts: unknown +) { + this.url = url + this.handler = handler + this.opts = opts + this.start = vi.fn().mockRejectedValue(new Error('Connection failed')) + this.destroy = vi.fn().mockResolvedValue(undefined) + channelInstances.push(this as (typeof channelInstances)[number]) +} + +// --- Tests --- + describe('FirehoseService', () => { let db: ReturnType let logger: ReturnType @@ -90,11 +192,19 @@ describe('FirehoseService', () => { beforeEach(() => { vi.clearAllMocks() + channelInstances.length = 0 + indexerInstances.length = 0 + mockTapChannelCtor.mockImplementation(defaultChannelImpl) + db = createMockDb() logger = createMockLogger() env = createMinimalEnv() }) + afterEach(() => { + vi.useRealTimers() + }) + describe('lifecycle', () => { it('creates a service instance', () => { const service = new FirehoseService(db as never, logger as never, env) @@ -102,55 +212,361 @@ describe('FirehoseService', () => { }) it('starts without throwing', async () => { - // Mock restoreTrackedRepos to find no repos - db.select.mockReturnValue({ - from: vi.fn().mockResolvedValue([]), - }) - const service = new FirehoseService(db as never, logger as never, env) await expect(service.start()).resolves.toBeUndefined() }) it('stops without throwing', async () => { - db.select.mockReturnValue({ - from: vi.fn().mockResolvedValue([]), + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + await expect(service.stop()).resolves.toBeUndefined() + }) + }) + + describe('cursor restoration', () => { + it('passes saved cursor as URL query parameter', async () => { + // restoreTrackedRepos → [], getCursor → [{ cursor: 42n }] + db.select + .mockReturnValueOnce({ from: vi.fn().mockReturnValue(mockQueryChain([])) }) + .mockReturnValueOnce({ from: vi.fn().mockReturnValue(mockQueryChain([{ cursor: 42n }])) }) + + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + expect(channelInstances).toHaveLength(1) + expect(channelInstances[0].url).toContain('cursor=42') + }) + + it('omits cursor when none is saved', async () => { + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + expect(channelInstances).toHaveLength(1) + expect(channelInstances[0].url).not.toContain('cursor') + }) + + it('constructs WebSocket URL from TAP_URL', async () => { + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + expect(channelInstances[0].url).toBe('ws://localhost:2480/channel') + }) + + it('converts https to wss in channel URL', async () => { + env.TAP_URL = 'https://tap.example.com' + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + expect(channelInstances[0].url).toBe('wss://tap.example.com/channel') + }) + }) + + describe('connected flag', () => { + it('reports disconnected before start', () => { + const service = new FirehoseService(db as never, logger as never, env) + expect(service.getStatus().connected).toBe(false) + }) + + it('reports disconnected immediately after start (before events)', async () => { + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + expect(service.getStatus().connected).toBe(false) + }) + + it('reports connected after first event is processed', async () => { + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + await getRecordHandler()(createTestRecordEvent(1)) + + expect(service.getStatus().connected).toBe(true) + expect(service.getStatus().lastEventId).toBe(1) + }) + + it('sets connected to false on channel error', async () => { + vi.useFakeTimers() + + let rejectChannel: ((err: Error) => void) | undefined + mockTapChannelCtor.mockImplementationOnce(function ( + this: Record, + url: string, + handler: unknown, + opts: unknown + ) { + this.url = url + this.handler = handler + this.opts = opts + this.start = vi.fn().mockReturnValue( + new Promise((_, reject) => { + rejectChannel = reject + }) + ) + this.destroy = vi.fn().mockResolvedValue(undefined) + channelInstances.push(this as (typeof channelInstances)[number]) }) const service = new FirehoseService(db as never, logger as never, env) await service.start() - await expect(service.stop()).resolves.toBeUndefined() + + await getRecordHandler()(createTestRecordEvent(1)) + expect(service.getStatus().connected).toBe(true) + + if (!rejectChannel) throw new Error('reject not captured') + rejectChannel(new Error('Connection lost')) + await vi.advanceTimersByTimeAsync(0) + + expect(service.getStatus().connected).toBe(false) + + await service.stop() + }) + + it('reports disconnected after stop', async () => { + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + await getRecordHandler()(createTestRecordEvent(1)) + expect(service.getStatus().connected).toBe(true) + + await service.stop() + expect(service.getStatus().connected).toBe(false) }) }) - describe('getStatus', () => { - it('returns status before start', () => { + describe('reconnection', () => { + it('retries after channel error with exponential backoff', async () => { + vi.useFakeTimers() + + let callCount = 0 + mockTapChannelCtor.mockImplementation(function ( + this: Record, + url: string, + handler: unknown, + opts: unknown + ) { + this.url = url + this.handler = handler + this.opts = opts + this.destroy = vi.fn().mockResolvedValue(undefined) + callCount++ + this.start = + callCount <= 2 + ? vi.fn().mockRejectedValue(new Error('Connection failed')) + : vi.fn().mockReturnValue(new Promise(() => {})) + channelInstances.push(this as (typeof channelInstances)[number]) + }) + const service = new FirehoseService(db as never, logger as never, env) - const status = service.getStatus() - expect(status.connected).toBe(false) - expect(status.lastEventId).toBeNull() + await service.start() + expect(channelInstances).toHaveLength(1) + + // First rejection → 1s backoff + await vi.advanceTimersByTimeAsync(0) + await vi.advanceTimersByTimeAsync(1000) + expect(channelInstances).toHaveLength(2) + + // Second rejection → 2s backoff + await vi.advanceTimersByTimeAsync(0) + await vi.advanceTimersByTimeAsync(2000) + expect(channelInstances).toHaveLength(3) + + await service.stop() }) - it('returns connected status after start', async () => { + it('resets backoff after successful event processing', async () => { + vi.useFakeTimers() + + let channel2Reject: ((err: Error) => void) | undefined + mockTapChannelCtor.mockImplementation(function ( + this: Record, + url: string, + handler: unknown, + opts: unknown + ) { + this.url = url + this.handler = handler + this.opts = opts + this.destroy = vi.fn().mockResolvedValue(undefined) + const idx = channelInstances.length + if (idx === 0) { + this.start = vi.fn().mockRejectedValue(new Error('fail')) + } else if (idx === 1) { + this.start = vi.fn().mockReturnValue( + new Promise((_, reject) => { + channel2Reject = reject + }) + ) + } else { + this.start = vi.fn().mockReturnValue(new Promise(() => {})) + } + channelInstances.push(this as (typeof channelInstances)[number]) + }) + + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + // Channel #1 fails → 1s backoff + await vi.advanceTimersByTimeAsync(0) + await vi.advanceTimersByTimeAsync(1000) + expect(channelInstances).toHaveLength(2) + + // Simulate event on channel #2 → resets backoff + await getRecordHandler()(createTestRecordEvent(100)) + + // Channel #2 fails + if (!channel2Reject) throw new Error('reject not captured') + channel2Reject(new Error('late failure')) + await vi.advanceTimersByTimeAsync(0) + + // Backoff should be 1s (reset), not 2s + const before = channelInstances.length + await vi.advanceTimersByTimeAsync(999) + expect(channelInstances).toHaveLength(before) + await vi.advanceTimersByTimeAsync(1) + expect(channelInstances).toHaveLength(before + 1) + + await service.stop() + }) + + it('reads latest cursor on reconnection', async () => { + vi.useFakeTimers() + + mockTapChannelCtor + .mockImplementationOnce(failingChannelImpl) + .mockImplementationOnce(defaultChannelImpl) + + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + expect(channelInstances[0].url).not.toContain('cursor') + + // Make getCursor return a value for the reconnection attempt db.select.mockReturnValue({ - from: vi.fn().mockResolvedValue([]), + from: vi.fn().mockReturnValue(mockQueryChain([{ cursor: 99n }])), }) + await vi.advanceTimersByTimeAsync(0) + await vi.advanceTimersByTimeAsync(1000) + + expect(channelInstances).toHaveLength(2) + expect(channelInstances[1].url).toContain('cursor=99') + + await service.stop() + }) + + it('caps backoff at 60 seconds', async () => { + vi.useFakeTimers() + mockTapChannelCtor.mockImplementation(failingChannelImpl) + const service = new FirehoseService(db as never, logger as never, env) await service.start() - const status = service.getStatus() - expect(status.connected).toBe(true) + + // Backoff sequence: 1s, 2s, 4s, 8s, 16s, 32s, 60s (capped) + const backoffs = [1000, 2000, 4000, 8000, 16000, 32000, 60000] + for (const ms of backoffs) { + await vi.advanceTimersByTimeAsync(0) + const before = channelInstances.length + await vi.advanceTimersByTimeAsync(ms) + expect(channelInstances.length).toBe(before + 1) + } + + // Next backoff should still be 60s (capped) + await vi.advanceTimersByTimeAsync(0) + const before = channelInstances.length + await vi.advanceTimersByTimeAsync(60000) + expect(channelInstances.length).toBe(before + 1) + + await service.stop() + }) + + it('does not reconnect after stop', async () => { + vi.useFakeTimers() + mockTapChannelCtor.mockImplementationOnce(failingChannelImpl) + + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + // Let rejection schedule reconnect timer + await vi.advanceTimersByTimeAsync(0) + + await service.stop() + + // Advance far past any backoff + await vi.advanceTimersByTimeAsync(120000) + + expect(channelInstances).toHaveLength(1) + }) + + it('does not reconnect when channel closes during shutdown', async () => { + vi.useFakeTimers() + + // Channel start resolves immediately (clean close) + mockTapChannelCtor.mockImplementationOnce(function ( + this: Record, + url: string, + handler: unknown, + opts: unknown + ) { + this.url = url + this.handler = handler + this.opts = opts + this.start = vi.fn().mockResolvedValue(undefined) + this.destroy = vi.fn().mockResolvedValue(undefined) + channelInstances.push(this as (typeof channelInstances)[number]) + }) + + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + await service.stop() + + // Let resolved promise propagate + await vi.advanceTimersByTimeAsync(0) + + // Should not reconnect + await vi.advanceTimersByTimeAsync(120000) + expect(channelInstances).toHaveLength(1) + }) + + it('logs reconnection attempts with attempt count and backoff', async () => { + vi.useFakeTimers() + + mockTapChannelCtor + .mockImplementationOnce(failingChannelImpl) + .mockImplementationOnce(defaultChannelImpl) + + const service = new FirehoseService(db as never, logger as never, env) + await service.start() + + await vi.advanceTimersByTimeAsync(0) + + const schedulingLog = logger.info.mock.calls.find( + (call: unknown[]) => + typeof call[1] === 'string' && (call[1]).includes('Scheduling') + ) + expect(schedulingLog).toBeDefined() + if (!schedulingLog) throw new Error('scheduling log not found') + expect(schedulingLog[0]).toEqual(expect.objectContaining({ attempt: 1, backoffMs: 1000 })) + + await vi.advanceTimersByTimeAsync(1000) + + const attemptLog = logger.info.mock.calls.find( + (call: unknown[]) => + typeof call[1] === 'string' && (call[1]).includes('Attempting') + ) + expect(attemptLog).toBeDefined() + if (!attemptLog) throw new Error('attempt log not found') + expect(attemptLog[0]).toEqual(expect.objectContaining({ attempt: 1 })) + + await service.stop() }) }) describe('error handling', () => { it('does not throw when start fails', async () => { - // Make restoreTrackedRepos fail db.select.mockReturnValue({ from: vi.fn().mockRejectedValue(new Error('DB down')), }) const service = new FirehoseService(db as never, logger as never, env) - // start() should catch errors internally await expect(service.start()).resolves.toBeUndefined() expect(logger.error).toHaveBeenCalled() }) -- 2.51.2