diff --git a/jest.config.cjs b/jest.config.cjs index be2cf47a..2217beb2 100644 --- a/jest.config.cjs +++ b/jest.config.cjs @@ -13,6 +13,7 @@ module.exports = { 'src/tests/dummy-connector/metadata-extraction.test.ts', 'src/http/axios-client-internal.test.ts', 'src/tests/event-data-size-limit/.*.test.ts', + 'src/tests/upload-failure/upload-failure.slow.test.ts', ], }, { @@ -33,6 +34,7 @@ module.exports = { '/src/tests/dummy-connector/metadata-extraction.test.ts', '/src/http/axios-client-internal.test.ts', '/src/tests/event-data-size-limit/size-limit-1.test.ts', + '/src/tests/upload-failure/upload-failure.slow.test.ts', ], }, ], diff --git a/src/mock-server/mock-server.interfaces.ts b/src/mock-server/mock-server.interfaces.ts index 0094407f..67685c19 100644 --- a/src/mock-server/mock-server.interfaces.ts +++ b/src/mock-server/mock-server.interfaces.ts @@ -56,6 +56,15 @@ export interface RouteConfig { headers?: Record; /** Optional retry configuration for simulating failures before success */ retry?: RetryConfig; + /** + * Respond successfully for the first N requests, then return an error. + * Useful for partial-upload scenarios (first batch succeeds, second fails). + */ + succeedThenFail?: { + successCount: number; + errorStatus?: number; + errorBody?: unknown; + }; /** Optional delay in milliseconds before sending the response */ delay?: number; } diff --git a/src/mock-server/mock-server.ts b/src/mock-server/mock-server.ts index 2ad6ee72..4d936fae 100644 --- a/src/mock-server/mock-server.ts +++ b/src/mock-server/mock-server.ts @@ -216,18 +216,51 @@ export class MockServer { * Configures a route to return a specific status code and optional response body. */ public setRoute(config: RouteConfig): void { - const { path, method, status, body, bodyBuffer, retry, headers, delay } = - config; + const { + path, + method, + status, + body, + bodyBuffer, + retry, + succeedThenFail, + headers, + delay, + } = config; const key = this.getRouteKey(method, path); - if (retry) { + if (retry || succeedThenFail) { this.requestCounts.set(key, 0); } this.routeHandlers.set(key, (req: ParsedRequest, res: MockResponse) => { const sendResponse = (responseDelay?: number) => { const send = () => { - if (retry) { + if (succeedThenFail) { + const currentCount = this.requestCounts.get(key) || 0; + this.requestCounts.set(key, currentCount + 1); + + if (currentCount < succeedThenFail.successCount) { + if (headers) { + res.set(headers); + } + + if (bodyBuffer !== undefined) { + res.status(status).buffer(bodyBuffer); + } else if (body !== undefined) { + res.status(status).json(body); + } else { + this.defaultRouteHandler(req, res); + } + } else { + const errorStatus = succeedThenFail.errorStatus ?? 400; + if (succeedThenFail.errorBody !== undefined) { + res.status(errorStatus).json(succeedThenFail.errorBody); + } else { + res.status(errorStatus).send(); + } + } + } else if (retry) { const currentCount = this.requestCounts.get(key) || 0; const failureCount = retry.failureCount ?? 4; const errorStatus = retry.errorStatus ?? 500; @@ -303,6 +336,16 @@ export class MockServer { this.requests = []; } + /** + * Returns all POST requests to the platform callback URL. + */ + public getCallbackRequests(): RequestInfo[] { + return this.requests.filter( + (req) => + req.method.toUpperCase() === 'POST' && req.url.includes('callback_url') + ); + } + /** * Returns the most recent request or undefined if no requests exist. */ diff --git a/src/multithreading/process-task.test.ts b/src/multithreading/process-task.test.ts index 015c5612..7739853f 100644 --- a/src/multithreading/process-task.test.ts +++ b/src/multithreading/process-task.test.ts @@ -62,6 +62,7 @@ jest.mock('./worker-adapter/worker-adapter', () => ({ WorkerAdapter: jest.fn().mockImplementation(() => ({ isTimeout: false, hasWorkerEmitted: false, + emitError: jest.fn(), })), })); @@ -150,7 +151,11 @@ describe(processTask.name, () => { // Arrange const event = makeEvent(); setWorkerData({ event, initialState: {}, options: {} }); - const mockAdapter = { isTimeout: false, hasWorkerEmitted: false }; + const mockAdapter = { + isTimeout: false, + hasWorkerEmitted: false, + emitError: jest.fn(), + }; (WorkerAdapter as jest.Mock).mockImplementation(() => mockAdapter); const task = jest.fn().mockResolvedValue(undefined); const onTimeout = jest.fn().mockResolvedValue(undefined); @@ -200,4 +205,31 @@ describe(processTask.name, () => { expect(processExitSpy).toHaveBeenCalledWith(1); expect(onTimeout).not.toHaveBeenCalled(); }); + + it('should emit a failure event through the adapter before falling back', async () => { + const event = makeEvent(); + setWorkerData({ event, initialState: {}, options: {} }); + const mockAdapter = { + isTimeout: false, + hasWorkerEmitted: false, + emitError: jest.fn().mockImplementation(() => { + mockAdapter.hasWorkerEmitted = true; + }), + }; + (WorkerAdapter as jest.Mock).mockImplementation(() => mockAdapter); + const taskError = new Error('task upload failed'); + const task = jest.fn().mockRejectedValue(taskError); + const onTimeout = jest.fn().mockResolvedValue(undefined); + + processTask({ task, onTimeout }); + await flush(); + + expect(mockAdapter.emitError).toHaveBeenCalledWith(taskError); + expect(mockParentPortPostMessage).not.toHaveBeenCalledWith( + expect.objectContaining({ + subject: WorkerMessageSubject.WorkerMessageFailed, + }) + ); + expect(processExitSpy).toHaveBeenCalledWith(1); + }); }); diff --git a/src/multithreading/process-task.ts b/src/multithreading/process-task.ts index 4f3e75a3..aa983863 100644 --- a/src/multithreading/process-task.ts +++ b/src/multithreading/process-task.ts @@ -23,6 +23,8 @@ export function processTask({ void (async () => { await runWithSdkLogContext(async () => { + let adapter: WorkerAdapter | undefined; + try { const event = workerData.event; @@ -44,27 +46,48 @@ export function processTask({ options, }); - const adapter = new WorkerAdapter({ + const workerAdapter = new WorkerAdapter({ event, adapterState, options, }); + adapter = workerAdapter; parentPort?.on(WorkerEvent.WorkerMessage, (message) => { if (message.subject !== WorkerMessageSubject.WorkerMessageExit) { return; } console.log('Timeout received. Waiting for the task to finish.'); - adapter.isTimeout = true; + workerAdapter.isTimeout = true; }); - await runWithUserLogContext(async () => task({ adapter })); - if (adapter.isTimeout && !adapter.hasWorkerEmitted) { - await runWithUserLogContext(async () => onTimeout({ adapter })); + await runWithUserLogContext(async () => + task({ adapter: workerAdapter }) + ); + if (workerAdapter.isTimeout && !workerAdapter.hasWorkerEmitted) { + await runWithUserLogContext(async () => + onTimeout({ adapter: workerAdapter }) + ); } process.exit(0); } catch (error) { - runWithUserLogContext(() => { + await runWithUserLogContext(async () => { + if (adapter && !adapter.hasWorkerEmitted) { + try { + await adapter.emitError(error); + } catch (failureError) { + console.error( + 'Error while emitting task failure.', + serializeError(failureError) + ); + } + } + + if (adapter?.hasWorkerEmitted) { + process.exit(1); + return; + } + const errorMessage = `Error while processing task. ${serializeError( error )}`; diff --git a/src/multithreading/worker-adapter/worker-adapter.emit.test.ts b/src/multithreading/worker-adapter/worker-adapter.emit.test.ts index 6812c369..bb092829 100644 --- a/src/multithreading/worker-adapter/worker-adapter.emit.test.ts +++ b/src/multithreading/worker-adapter/worker-adapter.emit.test.ts @@ -146,12 +146,11 @@ describe(`${WorkerAdapter.name}.emit`, () => { expect(mockPostMessage).toHaveBeenCalledTimes(1); }); - it('should correctly emit one event even if uploadAllRepos errors', async () => { + it('should log state size and current state before posting state', async () => { // Arrange + const logSpy = jest.spyOn(console, 'log').mockImplementation(); adapter['adapterState'].postState = jest.fn().mockResolvedValue(undefined); - adapter.uploadAllRepos = jest - .fn() - .mockRejectedValue(new Error('uploadAllRepos error')); + adapter.uploadAllRepos = jest.fn().mockResolvedValue(undefined); // Act await adapter.emit(ExtractorEventType.MetadataExtractionError, { @@ -160,7 +159,46 @@ describe(`${WorkerAdapter.name}.emit`, () => { }); // Assert - expect(mockPostMessage).toHaveBeenCalledTimes(1); + const stateLogCall = logSpy.mock.calls.find( + ([message]) => + typeof message === 'string' && + message.includes('Saving ') && + message.includes('Current state') + ); + expect(stateLogCall).toBeDefined(); + expect(stateLogCall?.[0]).toContain( + 'KB state before emitting event with event type: METADATA_EXTRACTION_ERROR. Current state' + ); + expect(stateLogCall?.[1]).toEqual( + expect.objectContaining({ attachments: { completed: false } }) + ); + }); + + it('should emit phase extraction error when uploadAllRepos fails', async () => { + const { emit: mockEmit } = require('../../common/control-protocol'); + adapter['adapterState'].postState = jest.fn().mockResolvedValue(undefined); + adapter.uploadAllRepos = jest + .fn() + .mockRejectedValue(new Error('uploadAllRepos error')); + + await adapter.emit(ExtractorEventType.DataExtractionDone); + + expect(mockEmit).toHaveBeenCalledWith( + expect.objectContaining({ + eventType: ExtractorEventType.DataExtractionError, + data: expect.objectContaining({ + error: expect.objectContaining({ + message: expect.stringContaining('uploadAllRepos error'), + }), + }), + }) + ); + expect(mockPostMessage).toHaveBeenCalledWith( + expect.objectContaining({ + subject: 'emit', + payload: { eventType: ExtractorEventType.DataExtractionError }, + }) + ); }); it('should include artifacts in data for extraction events', async () => { @@ -493,6 +531,35 @@ describe(`${WorkerAdapter.name}.emit — ExternalSyncUnitExtractionDone legacy p >; expect(emittedData).not.toHaveProperty('external_sync_units'); }); + + it('should emit an error when external sync unit upload fails', async () => { + const { adapter } = makeAdapter(EventType.StartExtractingExternalSyncUnits); + adapter['adapterState'].postState = jest.fn().mockResolvedValue(undefined); + adapter.uploadAllRepos = jest.fn().mockResolvedValue(undefined); + const pushMock = jest + .fn() + .mockRejectedValue(new Error('external sync unit upload failed')); + jest.spyOn(adapter, 'initializeRepos').mockImplementation(() => undefined); + jest.spyOn(adapter, 'getRepo').mockReturnValue({ push: pushMock } as never); + + await adapter.emit(ExtractorEventType.ExternalSyncUnitExtractionDone, { + external_sync_units: [{ id: 'esu-1' }] as never, + }); + + const { emit: mockEmit } = require('../../common/control-protocol'); + expect(mockEmit).toHaveBeenCalledWith( + expect.objectContaining({ + eventType: ExtractorEventType.ExternalSyncUnitExtractionError, + data: expect.objectContaining({ + error: expect.objectContaining({ + message: expect.stringContaining( + 'external sync unit upload failed' + ), + }), + }), + }) + ); + }); }); describe('WorkerAdapter — workersOldest / workersNewest boundary updates', () => { diff --git a/src/multithreading/worker-adapter/worker-adapter.ts b/src/multithreading/worker-adapter/worker-adapter.ts index aa5cee74..ceb59e8e 100644 --- a/src/multithreading/worker-adapter/worker-adapter.ts +++ b/src/multithreading/worker-adapter/worker-adapter.ts @@ -14,7 +14,7 @@ import { toRfc3339Timestamp, } from './worker-adapter.helpers'; import { ProgressData } from './worker-adapter.interfaces'; -import { serializeError } from '../../logger/logger'; +import { getPrintableState, serializeError } from '../../logger/logger'; import { runWithSdkLogContext, runWithUserLogContext, @@ -62,6 +62,7 @@ import { Uploader } from '../../uploader/uploader'; import { Artifact, SsorAttachment } from '../../uploader/uploader.interfaces'; import { translateOutgoingEventType } from '../../common/event-type-translation'; import { truncateMessage } from '../../common/helpers'; +import { getTimeoutErrorEventType } from '../spawn/spawn.helpers'; export function createWorkerAdapter({ event, @@ -273,127 +274,62 @@ export class WorkerAdapter { return; } - // If the event is ExternalSyncUnitExtractionDone, upload external sync units via a Repo before emitting - // TODO: Remove in v2.0.0 - if ( - newEventType === ExtractorEventType.ExternalSyncUnitExtractionDone && - data?.external_sync_units && - data.external_sync_units.length > 0 - ) { - console.log( - `Uploading ${data.external_sync_units.length} external sync units via repo before emitting event.` - ); - - this.initializeRepos([ - { - itemType: AirSyncDefaultItemTypes.EXTERNAL_SYNC_UNITS, - overridenOptions: { - batchSize: 25000, - skipConfirmation: true, - }, - }, - ]); - - await this.getRepo(AirSyncDefaultItemTypes.EXTERNAL_SYNC_UNITS)?.push( - data.external_sync_units - ); - - // Remove inline external_sync_units from data to avoid SQS size issues - delete data.external_sync_units; - } - - // Upload all repos before emitting the event - console.log( - `Uploading all repos before emitting event with event type: ${newEventType}.` - ); + let eventPayload: EventData; try { - await this.uploadAllRepos(); - } catch (error) { - console.error('Error while uploading repos', error); - parentPort?.postMessage(WorkerMessageSubject.WorkerMessageExit); - this.hasWorkerEmitted = true; - return; - } - - // If the extraction is done, we want to save the timestamp of the last successful sync - if (newEventType === ExtractorEventType.AttachmentExtractionDone) { - console.log( - `Overwriting lastSuccessfulSyncStarted with lastSyncStarted (${this.state.lastSyncStarted}).` - ); - - this.state.lastSuccessfulSyncStarted = this.state.lastSyncStarted; - this.state.lastSyncStarted = ''; - - // Clear pending extraction boundaries now that the cycle is complete - this.state.pendingWorkersOldest = ''; - this.state.pendingWorkersNewest = ''; - - // Update workersOldest and workersNewest boundaries from resolved extraction timestamps. - // Expand boundaries: workersOldest gets the earliest timestamp, workersNewest gets the latest. - const extractionStart = this.event.payload.event_context.extract_from; - const extractionEnd = this.event.payload.event_context.extract_to; - + // If the event is ExternalSyncUnitExtractionDone, upload external sync units via a Repo before emitting. + // TODO: Remove in v2.0.0 if ( - extractionStart && - (!this.state.workersOldest || - extractionStart < this.state.workersOldest) + newEventType === ExtractorEventType.ExternalSyncUnitExtractionDone && + data?.external_sync_units && + data.external_sync_units.length > 0 ) { console.log( - `Updating workersOldest from '${this.state.workersOldest}' to '${extractionStart}'.` + `Uploading ${data.external_sync_units.length} external sync units via repo before emitting event.` ); - this.state.workersOldest = extractionStart; - } - if ( - extractionEnd && - (!this.state.workersNewest || - extractionEnd > this.state.workersNewest) - ) { - console.log( - `Updating workersNewest from '${this.state.workersNewest}' to '${extractionEnd}'.` + this.initializeRepos([ + { + itemType: AirSyncDefaultItemTypes.EXTERNAL_SYNC_UNITS, + overridenOptions: { + batchSize: 25000, + skipConfirmation: true, + }, + }, + ]); + + await this.getRepo(AirSyncDefaultItemTypes.EXTERNAL_SYNC_UNITS)?.push( + data.external_sync_units ); - this.state.workersNewest = extractionEnd; + + // Remove inline external_sync_units from data to avoid SQS size issues + delete data.external_sync_units; } - } - // We want to save the state every time we emit an event, except for the start and delete events - if (!STATELESS_EVENT_TYPES.includes(this.event.payload.event_type)) { - console.log( - `Saving state before emitting event with event type: ${newEventType}.` + await this.beforeEmit(newEventType); + eventPayload = this.buildEmitPayload(newEventType, data); + } catch (error) { + console.error( + 'Error while preparing event for emission.', + serializeError(error) ); - - try { - await this.adapterState.postState(this.state); - } catch (error) { - console.error('Error while posting state', error); - parentPort?.postMessage(WorkerMessageSubject.WorkerMessageExit); - this.hasWorkerEmitted = true; - return; - } + await this.emitError(error); + return; } try { - // Always prune error messages to make them shorter before emit - if (data?.error?.message) { - data.error.message = truncateMessage(data.error.message); - } - const isExtractionEvent = Object.values(ExtractorEventType).includes( newEventType as ExtractorEventType ); - const isLoaderEvent = Object.values(LoaderEventType).includes( - newEventType as LoaderEventType - ); const progressData: ProgressData = {}; if ( isExtractionEvent && - (newEventType == ExtractorEventType.DataExtractionDone || - newEventType == ExtractorEventType.DataExtractionProgress || - newEventType == ExtractorEventType.AttachmentExtractionDone || - newEventType == ExtractorEventType.AttachmentExtractionProgress) + (newEventType === ExtractorEventType.DataExtractionDone || + newEventType === ExtractorEventType.DataExtractionProgress || + newEventType === ExtractorEventType.AttachmentExtractionDone || + newEventType === ExtractorEventType.AttachmentExtractionProgress) ) { const repo = this.lastExtractedItemType ? this.repos.find((r) => r.itemType === this.lastExtractedItemType) @@ -418,13 +354,7 @@ export class WorkerAdapter { await emit({ eventType: newEventType, event: this.event, - data: { - ...data, - ...(isExtractionEvent ? { artifacts: this.artifacts } : {}), - ...(isLoaderEvent - ? { reports: this.reports, processed_files: this.processedFiles } - : {}), - }, + data: eventPayload, worker_metadata: { ...progressData }, }); @@ -446,14 +376,162 @@ export class WorkerAdapter { }); } + /** + * Emits a phase-specific error after task or pre-emit execution fails. This + * path avoids flushing repos again, so artifacts uploaded before the failure + * are still reported without retrying the failed batch. + */ + async emitError(error: unknown): Promise { + return runWithSdkLogContext(async () => { + if (this.hasWorkerEmitted) { + return; + } + + // Include artifacts uploaded before a failure so they are not orphaned + // and can be associated with the emitted phase error. + for (const repo of this.repos) { + this.artifacts = [...this.artifacts, ...repo.uploadedArtifacts]; + } + + const { eventType } = getTimeoutErrorEventType( + this.event.payload.event_type + ); + + try { + await emit({ + eventType, + event: this.event, + data: this.buildEmitPayload(eventType, { + error: { message: serializeError(error) }, + }), + }); + + const message: WorkerMessageEmitted = { + subject: WorkerMessageSubject.WorkerMessageEmitted, + payload: { eventType }, + }; + this.artifacts = []; + parentPort?.postMessage(message); + this.hasWorkerEmitted = true; + } catch (emitError) { + console.error( + `Error while emitting error event with event type: ${eventType}.`, + serializeError(emitError) + ); + } + }); + } + async uploadAllRepos(): Promise { for (const repo of this.repos) { - const error = await repo.upload(); + await repo.upload(); this.artifacts.push(...repo.uploadedArtifacts); - if (error) { - throw error; + } + } + + /** + * Pre-emit hook: uploads all repos and updates extraction boundaries/state. + * Throws if uploading repos or persisting state fails. The caller emits a + * compensating error event with whatever artifacts did upload. + */ + private async beforeEmit( + eventType: ExtractorEventType | LoaderEventType + ): Promise { + // Upload all repos before emitting the event. + console.log( + `Uploading all repos before emitting event with event type: ${eventType}.` + ); + + await this.uploadAllRepos(); + + // If the extraction is done, save the timestamp of the last successful sync. + if (eventType === ExtractorEventType.AttachmentExtractionDone) { + console.log( + `Overwriting lastSuccessfulSyncStarted with lastSyncStarted (${this.state.lastSyncStarted}).` + ); + + this.state.lastSuccessfulSyncStarted = this.state.lastSyncStarted; + this.state.lastSyncStarted = ''; + + // Clear pending extraction boundaries now that the cycle is complete. + this.state.pendingWorkersOldest = ''; + this.state.pendingWorkersNewest = ''; + + // Update workersOldest and workersNewest boundaries from resolved extraction timestamps. + const extractionStart = this.event.payload.event_context.extract_from; + const extractionEnd = this.event.payload.event_context.extract_to; + + if ( + extractionStart && + (!this.state.workersOldest || + extractionStart < this.state.workersOldest) + ) { + console.log( + `Updating workersOldest from '${this.state.workersOldest}' to '${extractionStart}'.` + ); + this.state.workersOldest = extractionStart; + } + + if ( + extractionEnd && + (!this.state.workersNewest || extractionEnd > this.state.workersNewest) + ) { + console.log( + `Updating workersNewest from '${this.state.workersNewest}' to '${extractionEnd}'.` + ); + this.state.workersNewest = extractionEnd; } } + + if (STATELESS_EVENT_TYPES.includes(this.event.payload.event_type)) { + return; + } + + let stateSizeKb = 0; + try { + stateSizeKb = + Buffer.byteLength(JSON.stringify(this.state), 'utf8') / 1024; + } catch { + stateSizeKb = 0; + } + + console.log( + `Saving ${stateSizeKb.toFixed( + 2 + )} KB state before emitting event with event type: ${eventType}. Current state`, + getPrintableState(this.state) + ); + + try { + await this.adapterState.postState(this.state); + } catch (error) { + console.error('Error while posting state.', serializeError(error)); + throw error; + } + } + + private buildEmitPayload( + eventType: ExtractorEventType | LoaderEventType, + data?: EventData + ): EventData { + if (data?.error?.message) { + data.error.message = truncateMessage(data.error.message); + } + + const isExtractionEvent = Object.values(ExtractorEventType).includes( + eventType as ExtractorEventType + ); + const isLoaderEvent = Object.values(LoaderEventType).includes( + eventType as LoaderEventType + ); + + return { + ...data, + ...(isExtractionEvent ? { artifacts: this.artifacts } : {}), + ...(isLoaderEvent + ? { reports: this.reports, processed_files: this.processedFiles } + : {}), + }; } async loadItemTypes({ diff --git a/src/multithreading/worker-adapter/worker-adapter.upload-failure.test.ts b/src/multithreading/worker-adapter/worker-adapter.upload-failure.test.ts new file mode 100644 index 00000000..7f60202f --- /dev/null +++ b/src/multithreading/worker-adapter/worker-adapter.upload-failure.test.ts @@ -0,0 +1,168 @@ +import { State } from '../../state/state'; +import { mockServer } from '../../tests/jest.setup'; +import { createItems } from '../../tests/test-helpers'; +import { createMockEvent } from '../../common/test-utils'; +import { AdapterState, EventType, ExtractorEventType } from '../../types'; +import { WorkerAdapter } from './worker-adapter'; + +/* eslint-disable @typescript-eslint/no-require-imports */ + +jest.mock('../../common/control-protocol', () => ({ + emit: jest.fn().mockResolvedValue({}), +})); + +jest.mock('../../mappers/mappers'); +jest.mock('node:worker_threads', () => ({ + parentPort: { postMessage: jest.fn() }, +})); +jest.mock('../../attachments-streaming/attachments-streaming-pool', () => ({ + AttachmentsStreamingPool: jest.fn().mockImplementation(() => ({ + streamAll: jest.fn().mockResolvedValue(undefined), + })), +})); + +interface TestState { + attachments: { completed: boolean }; +} + +const UPLOAD_URL_PATH = '/internal/airdrop.artifacts.upload-url'; + +function makeAdapter(): { + adapter: WorkerAdapter; + mockPostMessage: jest.Mock; +} { + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + const initialState: AdapterState = { + attachments: { completed: false }, + lastSyncStarted: '', + lastSuccessfulSyncStarted: '', + snapInVersionId: '', + toDevRev: { + attachmentsMetadata: { + artifactIds: [], + lastProcessed: 0, + lastProcessedAttachmentsIdsList: [], + }, + }, + }; + const adapterState = new State({ event, initialState }); + const adapter = new WorkerAdapter({ event, adapterState }); + + const workerThreads = require('node:worker_threads'); + const mockPostMessage = jest.fn(); + if (workerThreads.parentPort) { + jest + .spyOn(workerThreads.parentPort, 'postMessage') + .mockImplementation(mockPostMessage); + } else { + workerThreads.parentPort = { postMessage: mockPostMessage }; + } + + return { adapter, mockPostMessage }; +} + +describe(`${WorkerAdapter.name} upload failure (near-integration)`, () => { + let adapter: WorkerAdapter; + let mockPostMessage: jest.Mock; + + beforeEach(() => { + jest.clearAllMocks(); + mockServer.resetRoutes(); + ({ adapter, mockPostMessage } = makeAdapter()); + adapter['adapterState'].postState = jest.fn().mockResolvedValue(undefined); + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + it('should throw from repo.push when upload-url returns 400 during batch upload', async () => { + mockServer.setRoute({ + path: UPLOAD_URL_PATH, + method: 'GET', + status: 400, + }); + + adapter.initializeRepos([ + { + itemType: 'tasks', + overridenOptions: { batchSize: 10 }, + }, + ]); + + await expect( + adapter.getRepo('tasks')?.push(createItems(20)) + ).rejects.toThrow('artifact upload URL'); + + expect(adapter.getRepo('tasks')?.getItems().length).toBe(20); + }); + + it('should emit DataExtractionError when uploadAllRepos fails', async () => { + const { emit: mockEmit } = require('../../common/control-protocol'); + + mockServer.setRoute({ + path: UPLOAD_URL_PATH, + method: 'GET', + status: 400, + }); + + adapter.initializeRepos([{ itemType: 'tasks' }]); + await adapter.getRepo('tasks')?.push(createItems(5)); + + await adapter.emit(ExtractorEventType.DataExtractionDone); + + expect(mockEmit).toHaveBeenCalledWith( + expect.objectContaining({ + eventType: ExtractorEventType.DataExtractionError, + data: expect.objectContaining({ + error: expect.objectContaining({ + message: expect.stringContaining('artifact upload URL'), + }), + }), + }) + ); + expect(mockPostMessage).toHaveBeenCalledWith( + expect.objectContaining({ + subject: 'emit', + payload: { eventType: ExtractorEventType.DataExtractionError }, + }) + ); + }); + + it('should include partial artifacts when push batch succeeded before uploadAllRepos fails', async () => { + const { emit: mockEmit } = require('../../common/control-protocol'); + + mockServer.setRoute({ + path: UPLOAD_URL_PATH, + method: 'GET', + status: 200, + succeedThenFail: { + successCount: 1, + errorStatus: 400, + }, + }); + + adapter.initializeRepos([ + { + itemType: 'tasks', + overridenOptions: { batchSize: 10 }, + }, + ]); + await adapter.getRepo('tasks')?.push(createItems(15)); + + await adapter.emit(ExtractorEventType.DataExtractionDone); + + expect(mockEmit).toHaveBeenCalledWith( + expect.objectContaining({ + eventType: ExtractorEventType.DataExtractionError, + data: expect.objectContaining({ + artifacts: expect.arrayContaining([ + expect.objectContaining({ item_count: 10, item_type: 'tasks' }), + ]), + }), + }) + ); + }); +}); diff --git a/src/repo/repo.test.ts b/src/repo/repo.test.ts index cfebef42..bb38811b 100644 --- a/src/repo/repo.test.ts +++ b/src/repo/repo.test.ts @@ -426,7 +426,9 @@ describe(Repo.name, () => { artifact: null, }); - await repo.upload([itemWithDate('1', '2022-01-01T00:00:00.000Z')]); + await expect( + repo.upload([itemWithDate('1', '2022-01-01T00:00:00.000Z')]) + ).rejects.toThrow('fail'); expect(repo.dateRanges.creationDate.oldest).toBe( ts('2022-01-01T00:00:00.000Z') @@ -493,4 +495,40 @@ describe(Repo.name, () => { ); }); }); + + it('should throw when upload fails', async () => { + mockUploadFn.mockResolvedValueOnce({ + error: { message: 'upload failed' }, + }); + + await expect(repo.upload(createItems(1))).rejects.toThrow('upload failed'); + }); + + it('should retain items in repo when batch upload fails during push', async () => { + repo = new Repo({ + event: createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.ExtractionDataStart }, + }), + itemType: 'test_item_type', + normalize, + onUpload: jest.fn(), + options: { batchSize: 10 }, + }); + + const items = createItems(20); + mockUploadFn + .mockResolvedValueOnce({ + artifact: { + id: 'artifact-1', + item_type: 'test_item_type', + item_count: 10, + }, + }) + .mockResolvedValueOnce({ + error: { message: 'second batch failed' }, + }); + + await expect(repo.push(items)).rejects.toThrow('second batch failed'); + expect(repo.getItems().length).toBe(10); + }); }); diff --git a/src/repo/repo.ts b/src/repo/repo.ts index 2ae1ee85..f477a9f6 100644 --- a/src/repo/repo.ts +++ b/src/repo/repo.ts @@ -4,7 +4,6 @@ import { SSOR_ATTACHMENT, } from '../common/constants'; import { Item } from '../repo/repo.interfaces'; -import { ErrorRecord } from '../types/common'; import { Uploader } from '../uploader/uploader'; import { Artifact } from '../uploader/uploader.interfaces'; @@ -55,7 +54,7 @@ export class Repo { async upload( batch?: (NormalizedItem | NormalizedAttachment | Item)[] - ): Promise { + ): Promise { const itemsToUpload = batch || this.items; if (itemsToUpload.length > 0) { @@ -86,8 +85,14 @@ export class Repo { ); if (error || !artifact) { - console.error('Error while uploading batch', error); - return error; + console.error( + `Error while uploading batch of ${itemsToUpload.length} items of type ${this.itemType}.`, + error + ); + throw new Error( + error?.message ?? + `Upload failed for item type "${this.itemType}" without artifact.` + ); } this.onUpload(artifact); @@ -136,16 +141,9 @@ export class Repo { // Upload in batches while the number of items exceeds the batch size const batchSize = this.options?.batchSize || ARTIFACT_BATCH_SIZE; while (this.items.length >= batchSize) { - // Slice out a batch of batchSize items to upload - const batch = this.items.splice(0, batchSize); - - try { - // Upload the batch - await this.upload(batch); - } catch (error) { - console.error('Error while uploading batch', error); - return false; - } + const batch = this.items.slice(0, batchSize); + await this.upload(batch); + this.items.splice(0, batchSize); } return true; diff --git a/src/tests/test-helpers.ts b/src/tests/test-helpers.ts index 74ff77ad..dcab2c70 100644 --- a/src/tests/test-helpers.ts +++ b/src/tests/test-helpers.ts @@ -179,3 +179,39 @@ export function spyOnPrivateMethod( // eslint-disable-next-line @typescript-eslint/no-explicit-any return jest.spyOn(instance as any, methodName as string); } + +export interface CallbackEventBody { + event_type: string; + event_data?: { + error?: { message: string }; + artifacts?: { id: string; item_type: string; item_count: number }[]; + }; +} + +export function getCallbackEventBodies( + requests: { body?: unknown }[] +): CallbackEventBody[] { + return requests.map((req) => req.body as CallbackEventBody); +} + +export function expectNoCallbackWithEventType( + bodies: CallbackEventBody[], + eventType: string +): void { + expect(bodies.some((b) => b.event_type === eventType)).toBe(false); +} + +export function expectLastCallbackError( + bodies: CallbackEventBody[], + expectedEventType: string, + messageSubstring?: string +): void { + expect(bodies.length).toBeGreaterThan(0); + const last = bodies[bodies.length - 1]; + expect(last.event_type).toBe(expectedEventType); + expect(last.event_data?.error?.message).toEqual(expect.any(String)); + expect(last.event_data?.error?.message?.length).toBeGreaterThan(0); + if (messageSubstring) { + expect(last.event_data?.error?.message).toContain(messageSubstring); + } +} diff --git a/src/tests/upload-failure/emit-partial-failure.ts b/src/tests/upload-failure/emit-partial-failure.ts new file mode 100644 index 00000000..3f64baf1 --- /dev/null +++ b/src/tests/upload-failure/emit-partial-failure.ts @@ -0,0 +1,39 @@ +import { + ExtractorEventType, + NormalizedItem, + processTask, + RepoInterface, +} from '../../index'; +import { Item } from '../../repo/repo.interfaces'; + +const repos: RepoInterface[] = [ + { + itemType: 'tasks', + overridenOptions: { batchSize: 10 }, + normalize: (task: Item): NormalizedItem => ({ + id: task.id, + created_date: task.created_at, + modified_date: task.updated_at, + data: { name: task.name }, + }), + }, +]; + +processTask({ + task: async ({ adapter }) => { + adapter.initializeRepos(repos); + + const tasks = Array.from({ length: 15 }, (_, i) => ({ + id: `task_${i}`, + name: `Task ${i}`, + created_at: new Date().toISOString(), + updated_at: new Date().toISOString(), + })); + + await adapter.getRepo('tasks')?.push(tasks); + await adapter.emit(ExtractorEventType.DataExtractionDone); + }, + onTimeout: async ({ adapter }) => { + await adapter.emit(ExtractorEventType.DataExtractionProgress); + }, +}); diff --git a/src/tests/upload-failure/emit-remainder-failure.ts b/src/tests/upload-failure/emit-remainder-failure.ts new file mode 100644 index 00000000..d7c268b5 --- /dev/null +++ b/src/tests/upload-failure/emit-remainder-failure.ts @@ -0,0 +1,38 @@ +import { + ExtractorEventType, + NormalizedItem, + processTask, + RepoInterface, +} from '../../index'; +import { Item } from '../../repo/repo.interfaces'; + +const repos: RepoInterface[] = [ + { + itemType: 'tasks', + normalize: (task: Item): NormalizedItem => ({ + id: task.id, + created_date: task.created_at, + modified_date: task.updated_at, + data: { name: task.name }, + }), + }, +]; + +processTask({ + task: async ({ adapter }) => { + adapter.initializeRepos(repos); + + const tasks = Array.from({ length: 5 }, (_, i) => ({ + id: `task_${i}`, + name: `Task ${i}`, + created_at: new Date().toISOString(), + updated_at: new Date().toISOString(), + })); + + await adapter.getRepo('tasks')?.push(tasks); + await adapter.emit(ExtractorEventType.DataExtractionDone); + }, + onTimeout: async ({ adapter }) => { + await adapter.emit(ExtractorEventType.DataExtractionProgress); + }, +}); diff --git a/src/tests/upload-failure/extraction.ts b/src/tests/upload-failure/extraction.ts new file mode 100644 index 00000000..4889fb6a --- /dev/null +++ b/src/tests/upload-failure/extraction.ts @@ -0,0 +1,24 @@ +import { AirdropEvent, spawn } from '../../index'; + +interface ExtractorState { + [key: string]: unknown; +} + +const initialState = {}; +const initialDomainMapping = {}; + +const run = async (events: AirdropEvent[], workerPath: string) => { + for (const event of events) { + await spawn({ + event, + initialState, + workerPath, + initialDomainMapping, + options: { + isLocalDevelopment: true, + }, + }); + } +}; + +export default run; diff --git a/src/tests/upload-failure/metadata-upload-failure.ts b/src/tests/upload-failure/metadata-upload-failure.ts new file mode 100644 index 00000000..1c00b06e --- /dev/null +++ b/src/tests/upload-failure/metadata-upload-failure.ts @@ -0,0 +1,24 @@ +import { ExtractorEventType, processTask } from '../../index'; + +const repos = [ + { + itemType: 'external_domain_metadata', + }, +]; + +processTask({ + task: async ({ adapter }) => { + adapter.initializeRepos(repos); + + await adapter + .getRepo('external_domain_metadata') + ?.push([{ id: 'metadata-1' }]); + + await adapter.emit(ExtractorEventType.MetadataExtractionDone); + }, + onTimeout: async ({ adapter }) => { + await adapter.emit(ExtractorEventType.MetadataExtractionError, { + error: { message: 'Failed to extract metadata. Lambda timeout.' }, + }); + }, +}); diff --git a/src/tests/upload-failure/mock-server.setup.ts b/src/tests/upload-failure/mock-server.setup.ts new file mode 100644 index 00000000..12936fad --- /dev/null +++ b/src/tests/upload-failure/mock-server.setup.ts @@ -0,0 +1,20 @@ +import { MockServer } from '../../mock-server/mock-server'; + +/** + * Dedicated mock server for upload-failure tests so they do not share the global + * jest.setup singleton with other suites running in parallel in the same worker. + * Port 0 assigns a unique port per instance. + */ +export const mockServer = new MockServer(0); + +beforeAll(async () => { + await mockServer.start(); +}); + +afterAll(async () => { + await mockServer.stop(); +}); + +beforeEach(() => { + mockServer.resetRoutes(); +}); diff --git a/src/tests/upload-failure/push-batch-failure.ts b/src/tests/upload-failure/push-batch-failure.ts new file mode 100644 index 00000000..94dab132 --- /dev/null +++ b/src/tests/upload-failure/push-batch-failure.ts @@ -0,0 +1,39 @@ +import { + ExtractorEventType, + NormalizedItem, + processTask, + RepoInterface, +} from '../../index'; +import { Item } from '../../repo/repo.interfaces'; + +const repos: RepoInterface[] = [ + { + itemType: 'tasks', + overridenOptions: { batchSize: 10 }, + normalize: (task: Item): NormalizedItem => ({ + id: task.id, + created_date: task.created_at, + modified_date: task.updated_at, + data: { name: task.name }, + }), + }, +]; + +processTask({ + task: async ({ adapter }) => { + adapter.initializeRepos(repos); + + const tasks = Array.from({ length: 20 }, (_, i) => ({ + id: `task_${i}`, + name: `Task ${i}`, + created_at: new Date().toISOString(), + updated_at: new Date().toISOString(), + })); + + await adapter.getRepo('tasks')?.push(tasks); + await adapter.emit(ExtractorEventType.DataExtractionDone); + }, + onTimeout: async ({ adapter }) => { + await adapter.emit(ExtractorEventType.DataExtractionProgress); + }, +}); diff --git a/src/tests/upload-failure/upload-failure.integration.test.ts b/src/tests/upload-failure/upload-failure.integration.test.ts new file mode 100644 index 00000000..c180e0e3 --- /dev/null +++ b/src/tests/upload-failure/upload-failure.integration.test.ts @@ -0,0 +1,242 @@ +import { EventType, ExtractorEventType } from '../../types/extraction'; +import { mockServer } from './mock-server.setup'; +import { createMockEvent } from '../../common/test-utils'; +import { + expectLastCallbackError, + expectNoCallbackWithEventType, + getCallbackEventBodies, +} from '../test-helpers'; + +import run from './extraction'; + +const UPLOAD_URL_PATH = '/internal/airdrop.artifacts.upload-url'; +const UPLOAD_URL_ERROR_SNIPPET = 'artifact upload URL'; + +function failUploadUrlRoute(status: number): void { + mockServer.setRoute({ + path: UPLOAD_URL_PATH, + method: 'GET', + status, + }); +} + +describe('Upload failure integration (fast)', () => { + beforeEach(() => { + mockServer.resetRoutes(); + }); + + describe('repo.push batch upload failure', () => { + it('should emit DataExtractionError when upload-url returns 400 during push', async () => { + failUploadUrlRoute(400); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/push-batch-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + UPLOAD_URL_ERROR_SNIPPET + ); + expect( + mockServer.getRequestCount('GET', UPLOAD_URL_PATH) + ).toBeGreaterThanOrEqual(1); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + UPLOAD_URL_ERROR_SNIPPET + ); + }); + + it('should emit DataExtractionError when presigned file upload returns 400 during push', async () => { + mockServer.setRoute({ + path: '/file-upload-url', + method: 'POST', + status: 400, + }); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/push-batch-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + 'uploading artifact' + ); + expect( + mockServer.getRequestCount('GET', UPLOAD_URL_PATH) + ).toBeGreaterThanOrEqual(1); + expect(mockServer.getRequestCount('POST', '/file-upload-url')).toBe(1); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + 'uploading artifact' + ); + }); + }); + + describe('uploadAllRepos failure at emit', () => { + it('should emit DataExtractionError when upload-url returns 400 during emit flush', async () => { + failUploadUrlRoute(400); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/emit-remainder-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + UPLOAD_URL_ERROR_SNIPPET + ); + expect( + callbacks[callbacks.length - 1].event_data?.error?.message + ).not.toContain('Worker exited without emitting'); + }); + + it('should emit error when presigned file upload returns 400', async () => { + mockServer.setRoute({ + path: '/file-upload-url', + method: 'POST', + status: 400, + }); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/emit-remainder-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + 'uploading artifact' + ); + }); + + it('should emit error when confirm-upload returns 400', async () => { + mockServer.setRoute({ + path: '/internal/airdrop.artifacts.confirm-upload', + method: 'POST', + status: 400, + }); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/emit-remainder-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + 'confirming artifact upload' + ); + }); + }); + + describe('partial upload then uploadAllRepos failure', () => { + it('should include one artifact when first batch succeeds and remainder upload fails', async () => { + mockServer.setRoute({ + path: UPLOAD_URL_PATH, + method: 'GET', + status: 200, + succeedThenFail: { + successCount: 1, + errorStatus: 400, + }, + }); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/emit-partial-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + UPLOAD_URL_ERROR_SNIPPET + ); + expect( + callbacks[callbacks.length - 1].event_data?.artifacts + ).toHaveLength(1); + expect( + callbacks[callbacks.length - 1].event_data?.artifacts?.[0].item_count + ).toBe(10); + }); + }); + + describe('metadata extraction phase', () => { + it('should emit MetadataExtractionError when upload-url returns 400', async () => { + failUploadUrlRoute(400); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingMetadata }, + }); + + await run([event], __dirname + '/metadata-upload-failure'); + + const callbacks = getCallbackEventBodies( + mockServer.getCallbackRequests() + ); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.MetadataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.MetadataExtractionError, + UPLOAD_URL_ERROR_SNIPPET + ); + }); + }); +}); diff --git a/src/tests/upload-failure/upload-failure.slow.test.ts b/src/tests/upload-failure/upload-failure.slow.test.ts new file mode 100644 index 00000000..a308b8a5 --- /dev/null +++ b/src/tests/upload-failure/upload-failure.slow.test.ts @@ -0,0 +1,83 @@ +/** + * Upload-failure retry tests (spawn integration). + * + * Uses a dedicated mock server (./mock-server.setup) for isolation when the slow + * project runs test files in parallel. These tests use the production retry + * count and are intentionally slower than the fast upload-failure suite. + */ +import { EventType, ExtractorEventType } from '../../types/extraction'; +import { createMockEvent } from '../../common/test-utils'; +import { + expectLastCallbackError, + expectNoCallbackWithEventType, + getCallbackEventBodies, +} from '../test-helpers'; +import { mockServer } from './mock-server.setup'; + +import run from './extraction'; + +jest.setTimeout(180000); + +const UPLOAD_URL_PATH = '/internal/airdrop.artifacts.upload-url'; +const TEST_HTTP_RETRIES = 5; + +function failUploadUrlPermanently(): void { + mockServer.setRoute({ + path: UPLOAD_URL_PATH, + method: 'GET', + status: 503, + }); +} + +describe('Upload failure integration (retry exhaustion)', () => { + it('should emit DataExtractionError after upload-url retries are exhausted during push', async () => { + failUploadUrlPermanently(); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/push-batch-failure'); + + const callbacks = getCallbackEventBodies(mockServer.getCallbackRequests()); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + 'artifact upload URL' + ); + expect(mockServer.getRequestCount('GET', UPLOAD_URL_PATH)).toBe( + TEST_HTTP_RETRIES + 1 + ); + }); + + it('should emit DataExtractionError after upload-url retries are exhausted during emit flush', async () => { + failUploadUrlPermanently(); + + const event = createMockEvent(mockServer.baseUrl, { + payload: { event_type: EventType.StartExtractingData }, + }); + + await run([event], __dirname + '/emit-remainder-failure'); + + const callbacks = getCallbackEventBodies(mockServer.getCallbackRequests()); + expectNoCallbackWithEventType( + callbacks, + ExtractorEventType.DataExtractionDone + ); + expectLastCallbackError( + callbacks, + ExtractorEventType.DataExtractionError, + 'artifact upload URL' + ); + expect( + callbacks[callbacks.length - 1].event_data?.error?.message + ).not.toContain('Worker exited without emitting'); + expect(mockServer.getRequestCount('GET', UPLOAD_URL_PATH)).toBe( + TEST_HTTP_RETRIES + 1 + ); + }); +});