diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index 451e2d09..76149a7d 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -216,8 +216,9 @@ model JobRun { jobType String @map("job_type") status String parentJobId Int? @map("parent_job_id") - jobTrigger String @map("job_trigger") - retryCount Int @default(0) @map("retry_count") + jobTrigger String @map("job_trigger") + triggeredByUser String? @map("triggered_by_user") + retryCount Int @default(0) @map("retry_count") error String? metadata Json @default("{}") createdAt DateTime @default(now()) @map("created_at") @db.Timestamptz() @@ -225,10 +226,28 @@ model JobRun { completedAt DateTime? @map("completed_at") @db.Timestamptz() parentJob JobRun? @relation("ChildJobs", fields: [parentJobId], references: [id]) childJobs JobRun[] @relation("ChildJobs") + activities JobActivity[] @@index([status]) @@index([parentJobId]) @@index([jobType, status]) + @@index([createdAt(sort: Desc)]) @@map("job_runs") @@schema("csa") } + +model JobActivity { + id Int @id @default(autoincrement()) + jobRunId Int? @map("job_run_id") + when DateTime @default(now()) @db.Timestamptz() + severity String + type String + related String? + jobRun JobRun? @relation(fields: [jobRunId], references: [id], onDelete: SetNull) + + @@index([jobRunId, when(sort: Desc)]) + @@index([severity]) + @@index([type]) + @@map("job_activities") + @@schema("csa") +} diff --git a/backend/src/api/batches/batches.service.spec.ts b/backend/src/api/batches/batches.service.spec.ts index deaa340a..c4402339 100644 --- a/backend/src/api/batches/batches.service.spec.ts +++ b/backend/src/api/batches/batches.service.spec.ts @@ -877,6 +877,37 @@ describe('BatchesService', () => { expectIncomplete(result, 104, ['Cancellation End Date', 'Cancellation Reason Code']) }) + + it('should tag manual add incomplete validation with BATCH warn metadata', async () => { + const contactMissingProvince = makeContact({ + id: 102, + caseNumber: 'CASE-102', + birthProvince: null, + }) + + setupCommonMocks([contactMissingProvince]) + mockPrisma.batch.update.mockResolvedValue({ + id: 1, + batchNumber: 1, + status: 'pending', + recordCount: 0, + batchDate: null, + createdAt: new Date(), + systemComments: null, + }) + const warnSpy = vi.spyOn(service['logger'], 'warn').mockImplementation(() => {}) + + await service.addContactsToPendingBatch([102], 'jsmith') + + expect(warnSpy).toHaveBeenCalledWith( + expect.stringContaining('Manual add to batch by jsmith'), + expect.objectContaining({ + activityType: 'BATCH', + related: expect.stringContaining('skipped due to missing CRA mandatory fields'), + }), + ) + warnSpy.mockRestore() + }) }) describe('User Story 40101 - S2: Auto-batch with CRA validation & auto-hold', () => { diff --git a/backend/src/api/batches/batches.service.ts b/backend/src/api/batches/batches.service.ts index b4da9ae3..f3949a7d 100644 --- a/backend/src/api/batches/batches.service.ts +++ b/backend/src/api/batches/batches.service.ts @@ -1,4 +1,6 @@ -import { BadRequestException, Injectable, Logger, NotFoundException } from '@nestjs/common' +import { BadRequestException, Injectable, NotFoundException } from '@nestjs/common' +import { AppLogger } from 'src/common/logger/app-logger' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { Prisma } from '@prisma/client' import { PrismaService } from 'src/common/database/prisma.service' import { @@ -87,7 +89,7 @@ export interface UpdateBatchStatusOptions { @Injectable() export class BatchesService { - private readonly logger = new Logger(BatchesService.name) + private readonly logger = new AppLogger(BatchesService.name) constructor( private prisma: PrismaService, @@ -104,6 +106,44 @@ export class BatchesService { }) } + private logBatchOperationIssues( + operation: 'add' | 'remove', + userId: string, + batchId: number, + result: Pick, + ): void { + const batchLabel = `batch ${batchId}` + const isAutoBatch = userId === 'SYSTEM' + const trigger = isAutoBatch ? 'Auto-batch' : `Manual ${operation} to batch by ${userId}` + + if (result.incomplete.length > 0) { + const count = result.incomplete.length + const detail = isAutoBatch + ? `${count} contacts auto-held due to missing CRA mandatory fields` + : `${count} contacts skipped due to missing CRA mandatory fields` + this.logger.warn(`${trigger}: ${detail} (${batchLabel})`, { + activityType: JobActivityType.BATCH, + related: `${detail} (${batchLabel})`, + }) + } + + const errors = result.skipped.filter((entry) => entry.reason === 'error') + if (errors.length > 0) { + this.logger.error(`${trigger}: ${errors.length} contacts failed (${batchLabel})`, { + activityType: JobActivityType.BATCH, + related: `${errors.length} contacts failed during ${operation} (${batchLabel})`, + }) + } + + const otherSkipped = result.skipped.filter((entry) => entry.reason !== 'error') + if (otherSkipped.length > 0) { + this.logger.warn(`${trigger}: ${otherSkipped.length} contacts skipped (${batchLabel})`, { + activityType: JobActivityType.BATCH, + related: `${otherSkipped.length} contacts skipped during ${operation} (${batchLabel})`, + }) + } + } + private async nextBatchNumber(tx: Prisma.TransactionClient): Promise { await tx.$executeRaw( Prisma.sql`SELECT pg_advisory_xact_lock(${BATCH_ADVISORY_LOCK_CLASS}, ${BATCH_NUMBER_ADVISORY_LOCK_OBJECT})`, @@ -566,6 +606,8 @@ export class BatchesService { }), ) + this.logBatchOperationIssues('add', userId, result.batch.id, result) + return result } @@ -712,6 +754,13 @@ export class BatchesService { ) if (!transition.success) { + this.logger.error( + `Manual remove from batch failed for contact ${contactId}: ${transition.reason}`, + { + activityType: JobActivityType.BATCH, + related: `Manual remove from batch contact ${contactId} by ${userId ?? 'unknown'}: ${transition.reason}`, + }, + ) throw new BadRequestException( `Failed to transition contact ${contactId} on REMOVE_FROM_BATCH: ${transition.reason}`, ) @@ -805,6 +854,8 @@ export class BatchesService { }), ) + this.logBatchOperationIssues('remove', userId, result.batch.id, result) + return result } diff --git a/backend/src/api/contacts/contacts.controller.spec.ts b/backend/src/api/contacts/contacts.controller.spec.ts index eeada64d..76de39f6 100644 --- a/backend/src/api/contacts/contacts.controller.spec.ts +++ b/backend/src/api/contacts/contacts.controller.spec.ts @@ -399,7 +399,19 @@ describe('ContactsController', () => { const res = await request(app.getHttpServer()).post('/contacts/1/run-eligibility').expect(200) expect(res.body).toEqual(result) - expect(service.runContactEligibility).toHaveBeenCalledWith(1) + expect(service.runContactEligibility).toHaveBeenCalledWith(1, 'SYSTEM') + }) + + it('should pass username from @CurrentUser when guard sets it', async () => { + const result = { previousStatus: 'eligible', newStatus: 'in_pay' } + vi.spyOn(service, 'runContactEligibility').mockResolvedValue(result) + + await request(app.getHttpServer()) + .post('/contacts/1/run-eligibility') + .set('x-test-username', 'jsmith') + .expect(200) + + expect(service.runContactEligibility).toHaveBeenCalledWith(1, 'jsmith') }) it('should return 404 when contact not found', async () => { diff --git a/backend/src/api/contacts/contacts.controller.ts b/backend/src/api/contacts/contacts.controller.ts index 9b931be3..8f049265 100644 --- a/backend/src/api/contacts/contacts.controller.ts +++ b/backend/src/api/contacts/contacts.controller.ts @@ -227,8 +227,8 @@ export class ContactsController { @ApiResponse({ status: 200, description: 'Eligibility result with previous and new status' }) @ApiResponse({ status: 404, description: 'Contact not found' }) @ApiResponse({ status: 422, description: 'Contact not found in staging tables' }) - async runEligibility(@Param('id', ParseIntPipe) id: number) { - return this.contactsService.runContactEligibility(id) + async runEligibility(@Param('id', ParseIntPipe) id: number, @CurrentUser() userId: string) { + return this.contactsService.runContactEligibility(id, userId) } @Patch(':id/review-flag') diff --git a/backend/src/api/contacts/contacts.service.spec.ts b/backend/src/api/contacts/contacts.service.spec.ts index dc5c2fec..e4aa7456 100644 --- a/backend/src/api/contacts/contacts.service.spec.ts +++ b/backend/src/api/contacts/contacts.service.spec.ts @@ -1714,8 +1714,10 @@ describe('ContactsService', () => { it('should throw NotFoundException when contact does not exist', async () => { vi.spyOn(prisma.contact, 'findUnique').mockResolvedValue(null) - await expect(service.runContactEligibility(999)).rejects.toThrow(NotFoundException) - await expect(service.runContactEligibility(999)).rejects.toThrow('Contact 999 not found') + await expect(service.runContactEligibility(999, 'JSMITH')).rejects.toThrow(NotFoundException) + await expect(service.runContactEligibility(999, 'JSMITH')).rejects.toThrow( + 'Contact 999 not found', + ) }) it('should map EligibilityInputError to UnprocessableEntityException', async () => { @@ -1724,11 +1726,23 @@ describe('ContactsService', () => { eligibility.runForContact = vi .fn() .mockRejectedValue(new EligibilityInputError('Contact ICM-1 not found in staging tables')) + const errorSpy = vi.spyOn(service['logger'], 'error').mockImplementation(() => {}) - await expect(service.runContactEligibility(1)).rejects.toThrow(UnprocessableEntityException) - await expect(service.runContactEligibility(1)).rejects.toThrow( + await expect(service.runContactEligibility(1, 'JSMITH')).rejects.toThrow( + UnprocessableEntityException, + ) + await expect(service.runContactEligibility(1, 'JSMITH')).rejects.toThrow( 'Contact ICM-1 not found in staging tables', ) + expect(errorSpy).toHaveBeenCalledWith( + 'Manual eligibility failed for contact 1: Contact ICM-1 not found in staging tables', + { + activityType: 'DATA_QUALITY', + related: + 'Manual eligibility contact 1 (ICM-1) by JSMITH: Contact ICM-1 not found in staging tables', + }, + ) + errorSpy.mockRestore() }) it('should propagate generic Errors without wrapping (becomes 500 at HTTP layer)', async () => { @@ -1737,9 +1751,9 @@ describe('ContactsService', () => { const dbError = new Error('connection terminated unexpectedly') eligibility.runForContact = vi.fn().mockRejectedValue(dbError) - await expect(service.runContactEligibility(1)).rejects.toBe(dbError) + await expect(service.runContactEligibility(1, 'JSMITH')).rejects.toBe(dbError) // Specifically must NOT have been rewrapped as a 422 - await expect(service.runContactEligibility(1)).rejects.not.toBeInstanceOf( + await expect(service.runContactEligibility(1, 'JSMITH')).rejects.not.toBeInstanceOf( UnprocessableEntityException, ) }) @@ -1751,7 +1765,7 @@ describe('ContactsService', () => { .fn() .mockResolvedValue({ previousStatus: 'eligible', newStatus: 'in_pay' }) - await expect(service.runContactEligibility(1)).resolves.toEqual({ + await expect(service.runContactEligibility(1, 'JSMITH')).resolves.toEqual({ previousStatus: 'eligible', newStatus: 'in_pay', }) @@ -1766,7 +1780,7 @@ describe('ContactsService', () => { .mockResolvedValue({ previousStatus: 'eligible', newStatus: 'in_pay' }) const icmSync = vi.spyOn(service['icmSyncBackService'], 'syncSingleContact') - await service.runContactEligibility(1) + await service.runContactEligibility(1, 'JSMITH') expect(icmSync).toHaveBeenCalledWith(1) }) @@ -1779,7 +1793,7 @@ describe('ContactsService', () => { .mockResolvedValue({ previousStatus: 'eligible', newStatus: 'eligible' }) const icmSync = vi.spyOn(service['icmSyncBackService'], 'syncSingleContact') - await service.runContactEligibility(1) + await service.runContactEligibility(1, 'JSMITH') expect(icmSync).not.toHaveBeenCalled() }) @@ -1793,11 +1807,20 @@ describe('ContactsService', () => { vi.spyOn(service['icmSyncBackService'], 'syncSingleContact').mockRejectedValue( new Error('ICM down'), ) + const warnSpy = vi.spyOn(service['logger'], 'warn').mockImplementation(() => {}) - await expect(service.runContactEligibility(1)).resolves.toEqual({ + await expect(service.runContactEligibility(1, 'JSMITH')).resolves.toEqual({ previousStatus: 'eligible', newStatus: 'in_pay', }) + + await vi.waitFor(() => { + expect(warnSpy).toHaveBeenCalledWith('Immediate ICM sync failed for contact 1: ICM down', { + activityType: 'ICM', + related: 'ICM sync failed after manual eligibility contact 1 by JSMITH', + }) + }) + warnSpy.mockRestore() }) }) }) diff --git a/backend/src/api/contacts/contacts.service.ts b/backend/src/api/contacts/contacts.service.ts index 5de579c2..c535f0a2 100644 --- a/backend/src/api/contacts/contacts.service.ts +++ b/backend/src/api/contacts/contacts.service.ts @@ -1,10 +1,11 @@ import { BadRequestException, Injectable, - Logger, NotFoundException, UnprocessableEntityException, } from '@nestjs/common' +import { AppLogger } from 'src/common/logger/app-logger' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { PaginatedResponse } from 'src/api/common/dto/paginated-response.dto' import { PrismaService } from 'src/common/database/prisma.service' import { @@ -36,7 +37,7 @@ import type { @Injectable() export class ContactsService { - private readonly logger = new Logger(ContactsService.name) + private readonly logger = new AppLogger(ContactsService.name) constructor( private prisma: PrismaService, @@ -829,6 +830,7 @@ export class ContactsService { async runContactEligibility( contactId: number, + triggeredByUser: string, ): Promise<{ previousStatus: string | null; newStatus: string }> { const contact = await this.prisma.contact.findUnique({ where: { id: contactId }, @@ -844,6 +846,10 @@ export class ContactsService { result = await this.eligibilityService.runForContact(contact.personIdIcm) } catch (err) { if (err instanceof EligibilityInputError) { + this.logger.error(`Manual eligibility failed for contact ${contactId}: ${err.message}`, { + activityType: JobActivityType.DATA_QUALITY, + related: `Manual eligibility contact ${contactId} (${contact.personIdIcm}) by ${triggeredByUser}: ${err.message}`, + }) throw new UnprocessableEntityException(err.message) } throw err @@ -866,6 +872,10 @@ export class ContactsService { this.icmSyncBackService.syncSingleContact(contactId).catch((err) => { this.logger.warn( `Immediate ICM sync failed for contact ${contactId}: ${(err as Error).message}`, + { + activityType: JobActivityType.ICM, + related: `ICM sync failed after manual eligibility contact ${contactId} by ${triggeredByUser}`, + }, ) }) } diff --git a/backend/src/api/jobs/jobs.controller.ts b/backend/src/api/jobs/jobs.controller.ts index bed6593f..abdcd1e9 100644 --- a/backend/src/api/jobs/jobs.controller.ts +++ b/backend/src/api/jobs/jobs.controller.ts @@ -16,13 +16,26 @@ import { ConfigService } from '@nestjs/config' import { ApiQuery, ApiResponse, ApiTags } from '@nestjs/swagger' import { Prisma } from '@prisma/client' import type { DeployEnv } from 'src/config/app.config' +import { JobActivitySeverity } from 'src/jobs/enums/job-activity-severity.enum' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobStatus } from 'src/jobs/enums/job-status.enum' import { JobTrigger } from 'src/jobs/enums/job-trigger.enum' import { JobType } from 'src/jobs/enums/job-type.enum' import { JobRunner } from 'src/jobs/job-runner.service' -import { JobsService } from 'src/jobs/jobs.service' +import { + JobsService, + type MonitoringActivityFilters, + type MonitoringHistoryFilters, +} from 'src/jobs/jobs.service' import { OpenshiftJobLauncher } from 'src/jobs/openshift-job-launcher.service' +import { + formatJobDisplayName, + formatJobSummary, + formatMonitoringStatus, + formatTriggeredBy, +} from 'src/jobs/job-monitoring.utils' import { CSAGuard } from '../common/guards/csa.guard' +import { CurrentUser } from '../common/decorators/current-user.decorator' import { canRunBulkJobInApiProcess } from './bulk-job-deploy-env' import { getJobRunWarning } from './job-openshift-advisory' @@ -46,6 +59,18 @@ interface JobRunResponse { warning?: string } +interface MonitoringJobResponse { + id: number + jobId: number + jobName: string + status: string + triggeredBy: string + started: Date | null + finished: Date | null + summary: string | null + warning: string | null +} + function toUserFacingJobError(jobType: string, error: string | null): string | null { return error ? `${jobType} ${GENERIC_JOB_FAILURE_SUFFIX}` : null } @@ -93,16 +118,12 @@ export class JobsController { private readonly openshiftJobLauncher: OpenshiftJobLauncher, ) {} - // Concurrency: csa.job_runs has a partial unique index on (job_type) WHERE status='RUNNING' - // so the second concurrent createJob for the same type raises P2002. We translate that to 409. - private async startFireAndForgetJob( - jobType: JobType, - ): Promise<{ jobRunId: number; message: string }> { - let jobRun + private async createEndUserJobRun(jobType: JobType, triggeredByUser: string) { try { - jobRun = await this.jobsService.createJob({ + return await this.jobsService.createJob({ jobType, jobTrigger: JobTrigger.END_USER, + triggeredByUser, }) } catch (err) { if (err instanceof Prisma.PrismaClientKnownRequestError && err.code === 'P2002') { @@ -110,6 +131,15 @@ export class JobsController { } throw err } + } + + // Concurrency: csa.job_runs has a partial unique index on (job_type) WHERE status='RUNNING' + // so the second concurrent createJob for the same type raises P2002. We translate that to 409. + private async startFireAndForgetJob( + jobType: JobType, + triggeredByUser: string, + ): Promise<{ jobRunId: number; message: string }> { + const jobRun = await this.createEndUserJobRun(jobType, triggeredByUser) this.jobRunner.executeJob(jobRun.id).catch((err) => { this.logger.error( @@ -126,30 +156,19 @@ export class JobsController { private async launchOpenShiftJob( jobType: JobType, + triggeredByUser: string, ): Promise<{ jobRunId: number; message: string; openshiftJobName?: string }> { if (!this.openshiftJobLauncher.isEnabled()) { const deployEnv = this.configService.get('app.deployEnv', 'local') if (canRunBulkJobInApiProcess(deployEnv)) { - return this.startFireAndForgetJob(jobType) + return this.startFireAndForgetJob(jobType, triggeredByUser) } throw new ServiceUnavailableException( `Bulk ${jobType} jobs must run in OpenShift when DEPLOY_ENV is ${deployEnv}. The job launcher is not available.`, ) } - - let jobRun - try { - jobRun = await this.jobsService.createJob({ - jobType, - jobTrigger: JobTrigger.END_USER, - }) - } catch (err) { - if (err instanceof Prisma.PrismaClientKnownRequestError && err.code === 'P2002') { - throw new ConflictException(`${jobType} is already running`) - } - throw err - } + const jobRun = await this.createEndUserJobRun(jobType, triggeredByUser) this.logger.log(`Created job_run ${jobRun.id} for ${jobType}, launching OpenShift Job...`) @@ -229,24 +248,114 @@ export class JobsController { @ApiResponse({ status: 201, description: 'RUN_ELIGIBILITY job started' }) @ApiResponse({ status: 409, description: 'RUN_ELIGIBILITY is already running' }) @ApiResponse({ status: 503, description: 'Failed to launch OpenShift Job' }) - async runEligibility() { - return this.launchOpenShiftJob(JobType.RUN_ELIGIBILITY) + async runEligibility(@CurrentUser() userId: string) { + return this.launchOpenShiftJob(JobType.RUN_ELIGIBILITY, userId) } @Post('auto-batch') @ApiResponse({ status: 201, description: 'AUTO_BATCH job started' }) @ApiResponse({ status: 409, description: 'AUTO_BATCH is already running' }) @ApiResponse({ status: 503, description: 'Failed to launch OpenShift Job' }) - async autoBatch() { - return this.launchOpenShiftJob(JobType.AUTO_BATCH) + async autoBatch(@CurrentUser() userId: string) { + return this.launchOpenShiftJob(JobType.AUTO_BATCH, userId) } @Post('send-cra-file') @ApiResponse({ status: 201, description: 'SEND_CRA_FILE job started' }) @ApiResponse({ status: 409, description: 'SEND_CRA_FILE is already running' }) @ApiResponse({ status: 503, description: 'Failed to launch OpenShift Job' }) - async sendCraFile() { - return this.launchOpenShiftJob(JobType.SEND_CRA_FILE) + async sendCraFile(@CurrentUser() userId: string) { + return this.launchOpenShiftJob(JobType.SEND_CRA_FILE, userId) + } + + @Get('monitoring/latest') + @ApiResponse({ status: 200, description: 'Latest monitored job run per job type' }) + async getLatestJobs() { + const jobs = await this.jobsService.getLatestJobsPerType() + return Promise.all( + jobs.map(async (job) => { + const advisory = await getJobRunWarning(job, this.openshiftJobLauncher) + return this.toMonitoringResponse(job, advisory) + }), + ) + } + + @Get('monitoring/history') + @ApiResponse({ status: 200, description: 'Monitored job history (last month)' }) + async getJobHistory( + @Query('jobType', new ParseEnumPipe(JobType, { optional: true })) jobType?: JobType, + @Query('status', new ParseEnumPipe(JobStatus, { optional: true })) status?: JobStatus, + @Query('jobId', new ParseIntPipe({ optional: true })) jobId?: number, + @Query('page', new ParseIntPipe({ optional: true })) page = 1, + @Query('limit', new ParseIntPipe({ optional: true })) limit = 10, + @Query('triggeredBy') triggeredBy?: string, + @Query('sortBy') sortBy?: MonitoringHistoryFilters['sortBy'], + @Query('sortOrder') sortOrder?: MonitoringHistoryFilters['sortOrder'], + ) { + const result = await this.jobsService.getJobHistory({ + jobType, + status, + jobId, + page, + limit, + triggeredBy, + sortBy, + sortOrder, + }) + + const data = await Promise.all( + result.data.map(async (job) => { + const advisory = await getJobRunWarning(job, this.openshiftJobLauncher) + return this.toMonitoringResponse(job, advisory) + }), + ) + + return { + ...result, + data, + } + } + + @Get('monitoring/activities') + @ApiResponse({ status: 200, description: 'Recent monitoring activities' }) + async getRecentActivities( + @Query('page', new ParseIntPipe({ optional: true })) page = 1, + @Query('limit', new ParseIntPipe({ optional: true })) limit = 10, + @Query('severity', new ParseEnumPipe(JobActivitySeverity, { optional: true })) + severity?: JobActivitySeverity, + @Query('type', new ParseEnumPipe(JobActivityType, { optional: true })) type?: JobActivityType, + @Query('sortBy') sortBy?: MonitoringActivityFilters['sortBy'], + @Query('sortOrder') sortOrder?: MonitoringActivityFilters['sortOrder'], + ) { + return this.jobsService.getRecentActivities(page, limit, { + severity, + type, + sortBy, + sortOrder, + }) + } + + @Get(':id/activities') + @ApiResponse({ status: 200, description: 'Monitoring activities for a specific job run' }) + async getJobActivities( + @Param('id', ParseIntPipe) jobId: number, + @Query('page', new ParseIntPipe({ optional: true })) page = 1, + @Query('limit', new ParseIntPipe({ optional: true })) limit = 10, + @Query('severity', new ParseEnumPipe(JobActivitySeverity, { optional: true })) + severity?: JobActivitySeverity, + @Query('type', new ParseEnumPipe(JobActivityType, { optional: true })) type?: JobActivityType, + @Query('sortBy') sortBy?: MonitoringActivityFilters['sortBy'], + @Query('sortOrder') sortOrder?: MonitoringActivityFilters['sortOrder'], + ) { + return this.jobsService.getActivities({ + jobRunId: jobId, + page, + limit, + severity, + type, + sortBy, + sortOrder, + }) } @Get(':id') @@ -261,4 +370,30 @@ export class JobsController { const warning = await getJobRunWarning(job, this.openshiftJobLauncher) return toJobRunResponse(job, warning) } + + private toMonitoringResponse( + job: { + id: number + jobType: string + status: string + jobTrigger: string + triggeredByUser?: string | null + startedAt: Date | null + completedAt: Date | null + metadata?: unknown + }, + advisoryWarning?: string, + ): MonitoringJobResponse { + return { + id: job.id, + jobId: job.id, + jobName: formatJobDisplayName(job.jobType), + status: formatMonitoringStatus(job.status), + triggeredBy: formatTriggeredBy(job), + started: job.startedAt, + finished: job.completedAt, + summary: formatJobSummary(job), + warning: advisoryWarning ?? null, + } + } } diff --git a/backend/src/api/jobs/jobs.controllers.spec.ts b/backend/src/api/jobs/jobs.controllers.spec.ts index 8ad19342..cac72601 100644 --- a/backend/src/api/jobs/jobs.controllers.spec.ts +++ b/backend/src/api/jobs/jobs.controllers.spec.ts @@ -11,7 +11,13 @@ import { CSAGuard } from '../common/guards/csa.guard' import { clearOpenshiftStatusCacheForTests } from './job-openshift-advisory' import { JobsController } from './jobs.controller' -const mockCSAGuard = { canActivate: () => true } +const mockCSAGuard = { + canActivate: (context: { switchToHttp: () => { getRequest: () => any } }) => { + const req = context.switchToHttp().getRequest() + req.username = 'JSMITH' + return true + }, +} describe('JobsController', () => { let app: INestApplication @@ -25,6 +31,10 @@ describe('JobsController', () => { createJob: vi.fn(), getJob: vi.fn(), getJobs: vi.fn(), + getLatestJobsPerType: vi.fn(), + getJobHistory: vi.fn(), + getRecentActivities: vi.fn(), + getActivities: vi.fn(), markFailed: vi.fn(), } @@ -136,6 +146,7 @@ describe('JobsController', () => { expect(mockJobsService.createJob).toHaveBeenCalledWith({ jobType: 'RUN_ELIGIBILITY', jobTrigger: 'END_USER', + triggeredByUser: 'JSMITH', }) expect(mockOpenshiftJobLauncher.launchJob).toHaveBeenCalledWith('RUN_ELIGIBILITY', 42) }) @@ -312,6 +323,7 @@ describe('JobsController', () => { expect(mockJobsService.createJob).toHaveBeenCalledWith({ jobType: 'SEND_CRA_FILE', jobTrigger: 'END_USER', + triggeredByUser: 'JSMITH', }) expect(mockOpenshiftJobLauncher.launchJob).toHaveBeenCalledWith('SEND_CRA_FILE', 789) }) @@ -381,4 +393,133 @@ describe('JobsController', () => { expect(mockJobRunner.executeJob).not.toHaveBeenCalled() }) }) + + describe('GET /jobs/monitoring/latest', () => { + it('should return latest monitored jobs with mapped display name and triggeredBy', async () => { + const now = new Date('2026-07-01T12:00:00Z') + mockJobsService.getLatestJobsPerType.mockResolvedValue([ + { + id: 41, + jobType: 'RUN_ELIGIBILITY', + status: 'SUCCESS', + jobTrigger: 'END_USER', + triggeredByUser: 'JSMITH', + startedAt: now, + completedAt: now, + metadata: { processed: 100, statusChanges: 2, newContacts: 1, skipped: 3 }, + }, + ]) + + const res = await request(app.getHttpServer()).get('/jobs/monitoring/latest').expect(200) + + expect(res.body).toEqual([ + { + id: 41, + jobId: 41, + jobName: 'Eligibility', + status: 'Success', + triggeredBy: 'JSMITH', + started: now.toISOString(), + finished: now.toISOString(), + summary: '100 processed, 2 updated, 1 new, 3 skipped', + warning: null, + }, + ]) + }) + }) + + describe('GET /jobs/monitoring/history', () => { + it('should forward filters and return paginated mapped rows', async () => { + const now = new Date('2026-07-02T10:00:00Z') + mockJobsService.getJobHistory.mockResolvedValue({ + data: [ + { + id: 7, + jobType: 'SEND_CRA_FILE', + status: 'FAILED', + jobTrigger: 'CRON', + startedAt: now, + completedAt: now, + metadata: null, + }, + ], + total: 1, + page: 1, + limit: 10, + }) + + const res = await request(app.getHttpServer()) + .get('/jobs/monitoring/history') + .query({ status: 'FAILED', triggeredBy: 'SYSTEM', page: 1, limit: 10 }) + .expect(200) + + expect(mockJobsService.getJobHistory).toHaveBeenCalledWith( + expect.objectContaining({ + status: 'FAILED', + triggeredBy: 'SYSTEM', + page: 1, + limit: 10, + }), + ) + expect(res.body.data[0]).toMatchObject({ + id: 7, + jobName: 'Send CRA File', + triggeredBy: 'SYSTEM', + summary: 'Job failed', + }) + }) + }) + + describe('GET /jobs/monitoring/activities', () => { + it('should return recent activities with pagination and filters', async () => { + mockJobsService.getRecentActivities.mockResolvedValue({ + data: [ + { id: 1, jobRunId: 5, severity: 'WARNING', type: 'CRA', related: 'Invalid file format' }, + ], + total: 1, + page: 1, + limit: 10, + }) + + const res = await request(app.getHttpServer()) + .get('/jobs/monitoring/activities') + .query({ severity: 'WARNING', type: 'CRA', page: 1, limit: 10 }) + .expect(200) + + expect(mockJobsService.getRecentActivities).toHaveBeenCalledWith(1, 10, { + severity: 'WARNING', + type: 'CRA', + sortBy: undefined, + sortOrder: undefined, + }) + expect(res.body.total).toBe(1) + }) + }) + + describe('GET /jobs/:id/activities', () => { + it('should return activities for selected job run', async () => { + mockJobsService.getActivities.mockResolvedValue({ + data: [{ id: 11, jobRunId: 99, severity: 'ERROR', type: 'JOB', related: 'boom' }], + total: 1, + page: 1, + limit: 10, + }) + + const res = await request(app.getHttpServer()) + .get('/jobs/99/activities') + .query({ page: 1, limit: 10, severity: 'ERROR', type: 'JOB' }) + .expect(200) + + expect(mockJobsService.getActivities).toHaveBeenCalledWith({ + jobRunId: 99, + page: 1, + limit: 10, + severity: 'ERROR', + type: 'JOB', + sortBy: undefined, + sortOrder: undefined, + }) + expect(res.body.data[0].jobRunId).toBe(99) + }) + }) }) diff --git a/backend/src/api/mock/mock.service.ts b/backend/src/api/mock/mock.service.ts index 83477039..bbca5b30 100644 --- a/backend/src/api/mock/mock.service.ts +++ b/backend/src/api/mock/mock.service.ts @@ -12,14 +12,14 @@ export class MockService { getFile(filename: string): unknown { if (/[^a-zA-Z0-9_-]/.test(filename)) { - this.logger.warn(`Invalid mock filename rejected: ${filename}`) + this.logger.log(`Invalid mock filename rejected: ${filename}`) return null } const filePath = join(this.dataDir, `${filename}.json`) if (!existsSync(filePath)) { - this.logger.warn(`Mock file not found: ${filePath}`) + this.logger.log(`Mock file not found: ${filePath}`) return null } diff --git a/backend/src/api/weekly-files/weekly-files.service.ts b/backend/src/api/weekly-files/weekly-files.service.ts index 731cafce..c50d5ace 100644 --- a/backend/src/api/weekly-files/weekly-files.service.ts +++ b/backend/src/api/weekly-files/weekly-files.service.ts @@ -1,12 +1,14 @@ -import { BadRequestException, Injectable, Logger, NotFoundException } from '@nestjs/common' +import { BadRequestException, Injectable, NotFoundException } from '@nestjs/common' import { Prisma } from '@prisma/client' import { PaginatedResponse } from 'src/api/common/dto/paginated-response.dto' import { PrismaService } from 'src/common/database/prisma.service' +import { AppLogger } from 'src/common/logger/app-logger' import { csaProcessingBatchDate } from 'src/common/utils' import { CRA_DATA_HANDLING_CONSTANT } from 'src/cra/cra.constant' import type { DetailRecord04, HeaderRecord } from 'src/cra/inbound/inbound-weekly.interface' import { RecordTypeCode, TranCode } from 'src/cra/inbound/inbound-weekly.interface' import { WklAssociatedRecordProcessorService } from 'src/cra/inbound/wkl-associated-record-processor.service' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { IcmSyncBackService } from 'src/sync/icm/icm-sync-back.service' import { BatchesService } from '../batches/batches.service' import type { ReprocessWeeklyFileResultDto } from './dto/associate-wkl-record.dto' @@ -268,7 +270,7 @@ const wklRecordDtoInclude = { @Injectable() export class WeeklyFilesService { - private readonly logger = new Logger(WeeklyFilesService.name) + private readonly logger = new AppLogger(WeeklyFilesService.name) constructor( private readonly prisma: PrismaService, @@ -538,7 +540,13 @@ export class WeeklyFilesService { try { await this.icmSyncBackService.syncFlaggedWithRetry() } catch (err) { - this.logger.warn(`ICM sync-back failed after WKL record reprocess: ${(err as Error).message}`) + this.logger.warn( + `ICM sync-back failed after WKL record reprocess: ${(err as Error).message}`, + { + activityType: JobActivityType.ICM, + related: `ICM sync-back failed after WKL record reprocess: ${(err as Error).message}`, + }, + ) } const updated = await this.prisma.wklFileRecord.findFirst({ @@ -627,7 +635,10 @@ export class WeeklyFilesService { try { await this.icmSyncBackService.syncFlaggedWithRetry() } catch (err) { - this.logger.warn(`ICM sync-back failed after WKL reprocess: ${(err as Error).message}`) + this.logger.warn(`ICM sync-back failed after WKL reprocess: ${(err as Error).message}`, { + activityType: JobActivityType.ICM, + related: `ICM sync-back failed after WKL reprocess: ${(err as Error).message}`, + }) } return { processedRecordIds, skippedRecords } diff --git a/backend/src/common/logger/app-logger.spec.ts b/backend/src/common/logger/app-logger.spec.ts index bdfc6fd4..f98393fb 100644 --- a/backend/src/common/logger/app-logger.spec.ts +++ b/backend/src/common/logger/app-logger.spec.ts @@ -1,17 +1,28 @@ import { winstonInstance } from './logger.config' import { AppLogger } from './app-logger' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' +import { JobActivitySeverity } from 'src/jobs/enums/job-activity-severity.enum' +import { + registerJobActivityRecorder, + resetJobActivityRecorder, + runWithJobExecutionScope, +} from 'src/jobs/job-monitoring-log' describe('AppLogger', () => { let logger: AppLogger let spy: ReturnType + const recordActivity = vi.fn().mockResolvedValue(undefined) beforeEach(() => { logger = new AppLogger('TestService') spy = vi.spyOn(winstonInstance, 'log').mockReturnValue(winstonInstance) + recordActivity.mockClear() + registerJobActivityRecorder(recordActivity) }) afterEach(() => { spy.mockRestore() + resetJobActivityRecorder() }) describe('alert()', () => { @@ -48,4 +59,73 @@ describe('AppLogger', () => { }) }) }) + + describe('warn() monitoring dual-write', () => { + it('should persist a job activity when activityType is provided', async () => { + await runWithJobExecutionScope(12, async () => { + logger.warn('Manual add to batch: 2 contacts skipped (batch 5)', { + activityType: JobActivityType.BATCH, + related: '2 contacts skipped during add (batch 5)', + }) + }) + + expect(recordActivity).toHaveBeenCalledWith({ + jobRunId: 12, + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.BATCH, + related: '2 contacts skipped during add (batch 5)', + }) + }) + + it('should not persist when warn is untagged', async () => { + await runWithJobExecutionScope(12, async () => { + logger.warn('Job already running, skipping') + }) + + expect(recordActivity).not.toHaveBeenCalled() + }) + }) + + describe('error() monitoring dual-write', () => { + it('should persist a job activity when activityType is provided', async () => { + await runWithJobExecutionScope(7, async () => { + logger.error('Manual remove from batch failed for contact 3', { + activityType: JobActivityType.BATCH, + related: 'Manual remove from batch contact 3: invalid_transition', + }) + }) + + expect(recordActivity).toHaveBeenCalledWith({ + jobRunId: 7, + severity: JobActivitySeverity.ERROR, + activityType: JobActivityType.BATCH, + related: 'Manual remove from batch contact 3: invalid_transition', + }) + }) + }) + + describe('crit() monitoring dual-write', () => { + it('should not persist when crit is untagged', async () => { + await runWithJobExecutionScope(5, async () => { + logger.crit('Job failed after all retries') + }) + + expect(recordActivity).not.toHaveBeenCalled() + }) + + it('should aggregate DATA_QUALITY crit logs tagged by category', async () => { + await runWithJobExecutionScope(5, async () => { + logger.crit('Skipping contact: empty/null in required fields [dob]', { + category: 'DATA_QUALITY', + }) + }) + + expect(recordActivity).toHaveBeenCalledWith({ + jobRunId: 5, + severity: JobActivitySeverity.CRITICAL, + activityType: JobActivityType.DATA_QUALITY, + related: 'Skipping contact: empty/null in required fields [dob]', + }) + }) + }) }) diff --git a/backend/src/common/logger/app-logger.ts b/backend/src/common/logger/app-logger.ts index e48a8f94..95ea1695 100644 --- a/backend/src/common/logger/app-logger.ts +++ b/backend/src/common/logger/app-logger.ts @@ -1,6 +1,23 @@ import { Logger } from '@nestjs/common' +import { JobMonitoringLogMeta, persistJobMonitoringLog } from 'src/jobs/job-monitoring-log' import { winstonInstance } from './logger.config' +function asMonitoringMeta(value: unknown): JobMonitoringLogMeta | undefined { + if (!value || typeof value !== 'object' || Array.isArray(value)) { + return undefined + } + return value as JobMonitoringLogMeta +} + +/** + * App logger with syslog levels and optional Monitoring dual-write. + * + * Convention: + * - `log` / `debug` / `verbose` — engineering narrative (Splunk only) + * - `warn` / `error` / `crit` — operator-relevant when tagged with `activityType` or `category` + * - Untagged `warn` / `error` / `crit` — Splunk only (auth, integration failures, routine guards) + * - Demote to `log` only mocks, config fallbacks, and internal engineering detail + */ export class AppLogger extends Logger { alert(message: string, metadata?: Record): void { winstonInstance.log('alert', message, { context: this.context, ...metadata }) @@ -8,5 +25,34 @@ export class AppLogger extends Logger { crit(message: string, metadata?: Record): void { winstonInstance.log('crit', message, { context: this.context, ...metadata }) + this.persistTaggedMonitoringLog('crit', message, metadata) + } + + warn(message: string, ...optionalParams: unknown[]): void { + super.warn(message, ...optionalParams) + this.persistTaggedMonitoringLog('warn', message, optionalParams[0]) + } + + error(message: string, ...optionalParams: unknown[]): void { + super.error(message, ...optionalParams) + this.persistTaggedMonitoringLog('error', message, optionalParams[0]) + } + + private persistTaggedMonitoringLog( + level: 'warn' | 'error' | 'crit', + message: string, + metadata: unknown, + ): void { + const meta = asMonitoringMeta(metadata) + if (!meta?.activityType && !meta?.category) { + return + } + + const aggregateDefault = level === 'crit' && !!meta.category ? undefined : false + + void persistJobMonitoringLog(level, message, { + ...meta, + aggregate: meta.aggregate ?? aggregateDefault, + }) } } diff --git a/backend/src/cra/cra.module.ts b/backend/src/cra/cra.module.ts index 80736b52..bf0e8e29 100644 --- a/backend/src/cra/cra.module.ts +++ b/backend/src/cra/cra.module.ts @@ -63,7 +63,7 @@ import { S3CraTransferService } from './transfer/s3-cra-transfer.service' return new S3CraTransferService(configService) } if (transferMode !== 'http') { - logger.warn(`Unknown CRA_TRANSFER_MODE "${transferMode}", falling back to http`) + logger.log(`Unknown CRA_TRANSFER_MODE "${transferMode}", falling back to http`) } return new HttpCraTransferService(httpService, configService) }, diff --git a/backend/src/cra/handlers/poll-cra-response.handler.spec.ts b/backend/src/cra/handlers/poll-cra-response.handler.spec.ts index 38a2a0a7..eaee5c36 100644 --- a/backend/src/cra/handlers/poll-cra-response.handler.spec.ts +++ b/backend/src/cra/handlers/poll-cra-response.handler.spec.ts @@ -61,6 +61,7 @@ describe('PollCraResponseHandler', () => { let mockPrisma: any let mockBatchesService: any let mockContactsService: any + let mockJobsService: any let mockIcmSyncBackService: any let mockWeeklyContactMatcher: any let mockWklFileRecordService: any @@ -149,6 +150,10 @@ describe('PollCraResponseHandler', () => { updateCsaStatus: vi.fn().mockResolvedValue({ success: true }), } + mockJobsService = { + addActivity: vi.fn().mockResolvedValue(undefined), + } + mockIcmSyncBackService = { syncFlaggedWithRetry: vi.fn().mockResolvedValue({ totalFlagged: 0, @@ -183,6 +188,7 @@ describe('PollCraResponseHandler', () => { mockPrisma, mockBatchesService, mockContactsService, + mockJobsService, mockIcmSyncBackService as any, mockWeeklyContactMatcher, mockWklFileRecordService, diff --git a/backend/src/cra/handlers/poll-cra-response.handler.ts b/backend/src/cra/handlers/poll-cra-response.handler.ts index 778b9983..395bfdc8 100644 --- a/backend/src/cra/handlers/poll-cra-response.handler.ts +++ b/backend/src/cra/handlers/poll-cra-response.handler.ts @@ -10,8 +10,10 @@ import { BATCH_DETAIL_EVENT, CSA_EVENT } from 'src/common/state-machine/constant import { csaProcessingBatchDate, parseWklDate } from 'src/common/utils' import { BaseJob } from 'src/jobs/base-job' import { JobType } from 'src/jobs/enums/job-type.enum' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobResult } from 'src/jobs/interfaces/job-result.interface' import { JobContext } from 'src/jobs/interfaces/job.interface' +import { JobsService } from 'src/jobs/jobs.service' import { IcmSyncBackService, SyncBackResult } from 'src/sync/icm/icm-sync-back.service' import { CRA_DATA_HANDLING_CONSTANT } from '../cra.constant' import { InboundFileService } from '../inbound/inbound-file.service' @@ -19,10 +21,10 @@ import { InboundResponseService } from '../inbound/inbound-response.service' import { InboundWeeklyResponseService } from '../inbound/inbound-weekly-response.service' import type { DetailRecord04, HeaderRecord } from '../inbound/inbound-weekly.interface' import { DETAIL_OUTCOME, type CraResDetail } from '../inbound/inbound.interface' -import { WklAssociatedRecordProcessorService } from '../inbound/wkl-associated-record-processor.service' -import { buildWklUpdatePayloads } from '../inbound/wkl-snapshot-data' import { WeeklyContactMatcherService } from '../inbound/weekly-contact-matcher.service' +import { WklAssociatedRecordProcessorService } from '../inbound/wkl-associated-record-processor.service' import { WklFileRecordService } from '../inbound/wkl-file-record.service' +import { buildWklUpdatePayloads } from '../inbound/wkl-snapshot-data' import { CraTransferService } from '../transfer/cra-transfer.service' const { DESTINATION_ID, @@ -57,6 +59,7 @@ export class PollCraResponseHandler extends BaseJob { private recordsWklUnmatchedApproved!: number private recordsWklUnmatchedRefused!: number private recordsWklUnmatchedSkipped!: number + private invalidFileFormatCount!: number private newCraRecordsInWkl: DetailRecord04[] = [] private unmatchedWklBatchId: number | null = null @@ -68,6 +71,7 @@ export class PollCraResponseHandler extends BaseJob { private readonly prisma: PrismaService, private readonly batchesService: BatchesService, private readonly contactsService: ContactsService, + private readonly jobsService: JobsService, private readonly icmSyncBackService: IcmSyncBackService, private readonly weeklyContactMatcher: WeeklyContactMatcherService, private readonly wklFileRecordService: WklFileRecordService, @@ -87,6 +91,7 @@ export class PollCraResponseHandler extends BaseJob { this.recordsWklUnmatchedApproved = 0 this.recordsWklUnmatchedRefused = 0 this.recordsWklUnmatchedSkipped = 0 + this.invalidFileFormatCount = 0 this.unmatchedWklBatchId = null this.newCraRecordsInWkl = [] @@ -120,7 +125,10 @@ export class PollCraResponseHandler extends BaseJob { try { syncResult = await this.icmSyncBackService.syncFlaggedWithRetry() } catch (err) { - this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`) + this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`, { + activityType: JobActivityType.ICM, + related: `ICM sync-back failed: ${(err as Error).message}`, + }) } const totalUpdated = @@ -131,11 +139,13 @@ export class PollCraResponseHandler extends BaseJob { this.recordsWklUnmatchedApproved + this.recordsWklUnmatchedRefused + const fileNames = sortedFiles.map((f) => f.fileName).join(', ') return { success: true, - message: `Processed ${totalRecordsProcessed} CRA response records from ${sortedFiles.length} file(s)`, + message: `Processed ${totalRecordsProcessed} records from ${sortedFiles.length} file(s): ${fileNames}`, metadata: { files_processed: sortedFiles.length, + file_names: sortedFiles.map((f) => f.fileName), records_updated: totalUpdated, records_accepted: this.recordsAccepted, records_rejected: this.recordsRejected, @@ -146,6 +156,7 @@ export class PollCraResponseHandler extends BaseJob { records_wkl_unmatched_approved: this.recordsWklUnmatchedApproved, records_wkl_unmatched_refused: this.recordsWklUnmatchedRefused, records_wkl_unmatched_skipped: this.recordsWklUnmatchedSkipped, + invalid_file_format_count: this.invalidFileFormatCount, batch_ids: [...this.processedBatchIds], syncResult, craNewRecordsInWkl: { @@ -207,7 +218,13 @@ export class PollCraResponseHandler extends BaseJob { const valid = this.inboundFileService.isValidResponseFile(file.fileName) if (!valid) { - this.logger.warn(`Invalid response file format: ${file.fileName}`) + this.invalidFileFormatCount += 1 + this.logger.warn(`Invalid response file format: ${file.fileName}`, { + activityType: JobActivityType.CRA, + aggregate: true, + aggregateKey: 'invalid-response-file-format', + related: `Invalid response file format (example: ${file.fileName})`, + }) } await this.prisma.transferFile.create({ @@ -250,7 +267,10 @@ export class PollCraResponseHandler extends BaseJob { parsed = this.inboundResponseService.parseFile(localFilePath) } } catch (error) { - this.logger.error(`Failed to parse response file ${responseFile.fileName}: ${error}`) + this.logger.error(`Failed to parse response file ${responseFile.fileName}: ${error}`, { + activityType: JobActivityType.CRA, + related: `Failed to parse response file ${responseFile.fileName}`, + }) await this.prisma.transferFile.update({ where: { id: responseFile.id }, data: { isValid: false, isDetailsProcessed: true }, @@ -318,7 +338,12 @@ export class PollCraResponseHandler extends BaseJob { }) if (!batchDetail) { - this.logger.warn(`Batch detail not found for referenceNum ${detail.referenceNum}`) + this.logger.warn(`Batch detail not found for referenceNum ${detail.referenceNum}`, { + activityType: JobActivityType.CRA, + aggregate: true, + aggregateKey: 'batch-detail-not-found', + related: `Batch detail not found for reference number (example: ${detail.referenceNum})`, + }) return } @@ -390,7 +415,12 @@ export class PollCraResponseHandler extends BaseJob { const wklType = TRANSACTION_TYPE_MAP[detail.transactionType] if (!wklType || !TRANSACTION_TYPES.includes(wklType)) { - this.logger.warn(`WKL: unexpected transaction type ${detail.transactionType}, skipping`) + this.logger.warn(`WKL: unexpected transaction type ${detail.transactionType}, skipping`, { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-unexpected-transaction', + related: `Unexpected WKL transaction type (example: ${detail.transactionType})`, + }) await this.persistWklRecord(ctx, detail, { matchStatus: WKL_MATCH_STATUS.SKIPPED }) this.recordsWklSkipped++ return @@ -401,12 +431,24 @@ export class PollCraResponseHandler extends BaseJob { this.logger.warn( `WKL: no matching batch detail for ${detail.childGivenName.trim()} ${detail.childSurName.trim()} ` + `(DIN: ${detail.childDin?.trim() || 'none'})`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-no-batch-detail', + related: 'No matching batch detail for WKL record', + }, ) const contacts = await this.weeklyContactMatcher.findMatchingContact(detail) if (!contacts) { this.logger.warn( `WKL: no matching contacts for ${detail.childGivenName.trim()} ${detail.childSurName.trim()} ` + `(DIN: ${detail.childDin?.trim() || 'none'})`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-no-matching-contacts', + related: 'No matching contacts for WKL record', + }, ) this.newCraRecordsInWkl.push(detail) await this.persistWklRecord(ctx, detail, { matchStatus: WKL_MATCH_STATUS.UNMATCHED }) @@ -441,6 +483,12 @@ export class PollCraResponseHandler extends BaseJob { this.logger.warn( `WKL: transaction type mismatch for contact ${batchDetail.contactId} — ` + `WKL says ${wklType}, batch detail says ${batchDetail.transactionType}`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-transaction-type-mismatch', + related: 'WKL transaction type mismatch with batch detail', + }, ) } @@ -505,6 +553,12 @@ export class PollCraResponseHandler extends BaseJob { } else { this.logger.warn( `WKL: unexpected status '${detail.status}' for contact ${batchDetail.contactId}, skipping`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-unexpected-status', + related: `Unexpected WKL status (example: ${detail.status})`, + }, ) await this.persistWklRecord(ctx, detail, { matchStatus: WKL_MATCH_STATUS.SKIPPED }) this.recordsWklSkipped++ diff --git a/backend/src/cra/handlers/send-cra-file.handler.spec.ts b/backend/src/cra/handlers/send-cra-file.handler.spec.ts index 908ca9f6..363f6294 100644 --- a/backend/src/cra/handlers/send-cra-file.handler.spec.ts +++ b/backend/src/cra/handlers/send-cra-file.handler.spec.ts @@ -74,6 +74,7 @@ describe('SendCraFileHandler', () => { let mockPrisma: any let mockBatchesService: any let mockContactsService: any + let mockJobsService: any let mockOutboundDataService: any let mockOutboundFileService: any let mockCraTransferService: any @@ -105,6 +106,10 @@ describe('SendCraFileHandler', () => { updateCsaStatus: vi.fn().mockResolvedValue({ success: true }), } + mockJobsService = { + addActivity: vi.fn().mockResolvedValue(undefined), + } + mockOutboundDataService = { buildCraFileData: vi.fn().mockReturnValue({ header: { tranCode: 6133 }, @@ -158,6 +163,7 @@ describe('SendCraFileHandler', () => { mockOutboundDataService, mockOutboundFileService, mockCraTransferService, + mockJobsService, mockIcmSyncBackService, ) }) @@ -175,6 +181,7 @@ describe('SendCraFileHandler', () => { expect(result.success).toBe(true) expect(result.message).toContain('No batch to process') + expect(result.metadata).toEqual({ no_batch: true }) expect(mockBatchesService.updateBatchStatus).not.toHaveBeenCalled() expect(mockOutboundFileService.createFile).not.toHaveBeenCalled() }) @@ -190,6 +197,7 @@ describe('SendCraFileHandler', () => { expect(result.success).toBe(true) expect(result.message).toContain('No batch to process') + expect(result.metadata).toEqual({ no_batch: true }) expect(mockBatchesService.updateBatchStatus).not.toHaveBeenCalled() expect(mockOutboundFileService.createFile).not.toHaveBeenCalled() }) @@ -329,6 +337,7 @@ describe('SendCraFileHandler', () => { expect(result.metadata).toEqual({ batch_id: 10, file_path: '/tmp/cra/testfile.txt', + file_name: 'testfile.txt', record_count: 3, contacts_count: 2, }) @@ -493,6 +502,9 @@ describe('SendCraFileHandler', () => { expect(warnSpy).toHaveBeenCalledWith( 'Batch 10: SEND_FAILED transition failed (Invalid transition); persisted systemComments and status via direct update', + expect.objectContaining({ + activityType: 'BATCH', + }), ) expect(mockPrisma.batch.update).toHaveBeenCalledWith({ where: { id: 10 }, @@ -524,6 +536,9 @@ describe('SendCraFileHandler', () => { expect(warnSpy).toHaveBeenCalledWith( 'Batch 10: SEND_FAILED transition failed (Invalid transition); persisted systemComments only', + expect.objectContaining({ + activityType: 'BATCH', + }), ) expect(mockPrisma.batch.update).toHaveBeenCalledWith({ where: { id: 10 }, diff --git a/backend/src/cra/handlers/send-cra-file.handler.ts b/backend/src/cra/handlers/send-cra-file.handler.ts index 4d30d9b1..75403827 100644 --- a/backend/src/cra/handlers/send-cra-file.handler.ts +++ b/backend/src/cra/handlers/send-cra-file.handler.ts @@ -15,8 +15,10 @@ import { import { appendSystemComment, pacificToday } from 'src/common/utils' import { BaseJob } from 'src/jobs/base-job' import { JobType } from 'src/jobs/enums/job-type.enum' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobResult } from 'src/jobs/interfaces/job-result.interface' import { JobContext } from 'src/jobs/interfaces/job.interface' +import { JobsService } from 'src/jobs/jobs.service' import { IcmSyncBackService } from 'src/sync/icm/icm-sync-back.service' import { CRA_DATA_HANDLING_CONSTANT } from '../cra.constant' import { OutboundDataService } from '../outbound/outbound-data.service' @@ -41,6 +43,7 @@ export class SendCraFileHandler extends BaseJob { private readonly outboundDataService: OutboundDataService, private readonly outboundFileService: OutboundFileService, private readonly craTransferService: CraTransferService, + private readonly jobsService: JobsService, private readonly icmSyncBackService: IcmSyncBackService, ) { super() @@ -91,7 +94,7 @@ export class SendCraFileHandler extends BaseJob { async execute(_context: JobContext): Promise { if (!this.batch || this.batchDetails.length === 0) { - return { success: true, message: 'No batch to process' } + return { success: true, message: 'No batch to process', metadata: { no_batch: true } } } const { header, details, trailer } = this.outboundDataService.buildCraFileData( @@ -148,10 +151,11 @@ export class SendCraFileHandler extends BaseJob { return { success: true, - message: `Batch ${this.batch.id} sent to CRA`, + message: `Batch ${this.batch.id} sent to CRA (file: ${fileName}, ${recordCount} records)`, metadata: { batch_id: this.batch.id, file_path: filePath, + file_name: fileName, record_count: recordCount, contacts_count: this.batchDetails.length, }, @@ -182,14 +186,17 @@ export class SendCraFileHandler extends BaseJob { try { await this.icmSyncBackService.syncFlaggedWithRetry() } catch (err) { - this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`) + this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`, { + activityType: JobActivityType.ICM, + related: `ICM sync-back failed: ${(err as Error).message}`, + }) } } private async ensureBatchDetailsReady(): Promise { const missingRefDetails = this.batchDetails.filter((detail) => !detail.referenceNumber) if (missingRefDetails.length > 0) { - this.logger.warn( + this.logger.log( `Batch ${this.batch!.id}: ${missingRefDetails.length} details missing referenceNumber, backfilling`, ) for (const detail of missingRefDetails) { @@ -209,7 +216,10 @@ export class SendCraFileHandler extends BaseJob { if (!this.batch) return - this.logger.error(`File transfer failed for batch ${this.batch.id}`, error) + this.logger.error(`File transfer failed for batch ${this.batch.id}`, { + activityType: JobActivityType.CRA, + related: error.message || 'File transfer failed', + }) const errorMessage = error.message || 'File transfer failed' const batchRecord = await this.prisma.batch.findUnique({ where: { id: this.batch.id }, @@ -235,10 +245,18 @@ export class SendCraFileHandler extends BaseJob { if (data.status) { this.logger.warn( `Batch ${this.batch.id}: SEND_FAILED transition failed (${result.reason}); persisted systemComments and status via direct update`, + { + activityType: JobActivityType.BATCH, + related: `Batch ${this.batch.id}: send failed transition error`, + }, ) } else { this.logger.warn( `Batch ${this.batch.id}: SEND_FAILED transition failed (${result.reason}); persisted systemComments only`, + { + activityType: JobActivityType.BATCH, + related: `Batch ${this.batch.id}: send failed transition error`, + }, ) } await this.prisma.batch.update({ diff --git a/backend/src/cra/inbound/weekly-contact-matcher.service.ts b/backend/src/cra/inbound/weekly-contact-matcher.service.ts index 1622ba8c..d0c40116 100644 --- a/backend/src/cra/inbound/weekly-contact-matcher.service.ts +++ b/backend/src/cra/inbound/weekly-contact-matcher.service.ts @@ -4,6 +4,7 @@ import { PrismaService } from 'src/common/database/prisma.service' import { BATCH_DETAIL_STATUS } from 'src/common/state-machine/constants/batch-detail-status.constants' import { BATCH_STATUS } from 'src/common/state-machine/constants/batch-status.constants' import { AppLogger } from 'src/common/logger/app-logger' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { normalize, parseWklDate } from 'src/common/utils' import { CRA_DATA_HANDLING_CONSTANT } from '../cra.constant' import { CraMatchingSnapshot } from './cra-matching-snapshot.interface' @@ -165,6 +166,12 @@ export class WeeklyContactMatcherService { this.logger.warn( `WKL: multiple batch detail matches (${matches.length}) for ` + `${wklDetail.childGivenName.trim()} ${wklDetail.childSurName.trim()}, skipping`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-multiple-batch-detail-matches', + related: 'Multiple batch detail matches for WKL record', + }, ) return null } @@ -226,6 +233,12 @@ export class WeeklyContactMatcherService { if (matches.length > 1) { this.logger.warn( `WKL backfill: multiple CRA batch details for contact ${contactId} on ${weeklyFileDate.toISOString().slice(0, 10)}`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-multiple-cra-batch-details', + related: 'Multiple CRA batch details for WKL backfill match', + }, ) } return null @@ -268,7 +281,12 @@ export class WeeklyContactMatcherService { }) if (dinMatches.length === 1) return dinMatches[0] if (dinMatches.length > 1) { - this.logger.warn(`WKL contact match: multiple contacts with DIN ${din}, skipping`) + this.logger.warn(`WKL contact match: multiple contacts with DIN ${din}, skipping`, { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-multiple-contacts-din', + related: `Multiple contacts matched for DIN (example: ${din})`, + }) return null } } @@ -328,6 +346,12 @@ export class WeeklyContactMatcherService { this.logger.warn( `WKL contact match: multiple contacts (${detailMatches.length}) for ` + `${wklDetail.childGivenName.trim()} ${wklDetail.childSurName.trim()}, skipping`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-multiple-contacts-details', + related: 'Multiple contacts matched for WKL child details', + }, ) return null } diff --git a/backend/src/cra/inbound/wkl-associated-record-processor.service.ts b/backend/src/cra/inbound/wkl-associated-record-processor.service.ts index 8e58f975..78411dfd 100644 --- a/backend/src/cra/inbound/wkl-associated-record-processor.service.ts +++ b/backend/src/cra/inbound/wkl-associated-record-processor.service.ts @@ -1,7 +1,9 @@ -import { Injectable, Logger } from '@nestjs/common' +import { Injectable } from '@nestjs/common' import { BatchesService } from 'src/api/batches/batches.service' import { ContactsService } from 'src/api/contacts/contacts.service' +import { AppLogger } from 'src/common/logger/app-logger' import { BATCH_DETAIL_EVENT, CSA_STATUS } from 'src/common/state-machine/constants' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { buildWklUpdatePayloads } from './wkl-snapshot-data' import { CRA_DATA_HANDLING_CONSTANT } from '../cra.constant' import type { DetailRecord04, HeaderRecord } from './inbound-weekly.interface' @@ -29,7 +31,7 @@ export interface WklUnmatchedProcessCounters { @Injectable() export class WklAssociatedRecordProcessorService { - private readonly logger = new Logger(WklAssociatedRecordProcessorService.name) + private readonly logger = new AppLogger(WklAssociatedRecordProcessorService.name) constructor( private readonly batchesService: BatchesService, @@ -48,6 +50,12 @@ export class WklAssociatedRecordProcessorService { if (!wklType || !TRANSACTION_TYPES.includes(wklType)) { this.logger.warn( `WKL: unexpected transaction type ${detail.transactionType}, skipping [origin: ${ctx.origin}]`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-unexpected-transaction', + related: `Unexpected WKL transaction type (example: ${detail.transactionType})`, + }, ) counters.skipped++ return null @@ -70,6 +78,12 @@ export class WklAssociatedRecordProcessorService { `WKL: transaction type mismatch for contact ${contactId} — ` + `WKL says ${wklType}, batch detail says ${batchDetail.transactionType} ` + `[origin: ${ctx.origin}]`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-transaction-type-mismatch', + related: 'WKL transaction type mismatch with batch detail', + }, ) counters.skipped++ return null @@ -144,6 +158,12 @@ export class WklAssociatedRecordProcessorService { this.logger.warn( `WKL: unexpected status '${detail.status}' for contact ${batchDetail.contactId}, skipping ` + `[origin: ${ctx.origin}]`, + { + activityType: JobActivityType.WKL, + aggregate: true, + aggregateKey: 'wkl-unexpected-status', + related: `Unexpected WKL status (example: ${detail.status})`, + }, ) counters.skipped++ return null diff --git a/backend/src/jobs/enums/job-activity-severity.enum.ts b/backend/src/jobs/enums/job-activity-severity.enum.ts new file mode 100644 index 00000000..c64e05d5 --- /dev/null +++ b/backend/src/jobs/enums/job-activity-severity.enum.ts @@ -0,0 +1,6 @@ +/** FDD APL-12 severity levels (OS-CSA-JA-02). */ +export enum JobActivitySeverity { + WARNING = 'WARNING', + ERROR = 'ERROR', + CRITICAL = 'CRITICAL', +} diff --git a/backend/src/jobs/enums/job-activity-type.enum.ts b/backend/src/jobs/enums/job-activity-type.enum.ts new file mode 100644 index 00000000..da71e27d --- /dev/null +++ b/backend/src/jobs/enums/job-activity-type.enum.ts @@ -0,0 +1,9 @@ +/** FDD APL-12 activity categories (OS-CSA-JA-03). */ +export enum JobActivityType { + DATA_QUALITY = 'DATA_QUALITY', + JOB = 'JOB', + CRA = 'CRA', + WKL = 'WKL', + ICM = 'ICM', + BATCH = 'BATCH', +} diff --git a/backend/src/jobs/handlers/retry-failed.handler.ts b/backend/src/jobs/handlers/retry-failed.handler.ts index 4970ad31..db827f2d 100644 --- a/backend/src/jobs/handlers/retry-failed.handler.ts +++ b/backend/src/jobs/handlers/retry-failed.handler.ts @@ -2,6 +2,7 @@ import { Injectable } from '@nestjs/common' import { IcmSyncBackService, SyncBackResult } from 'src/sync/icm/icm-sync-back.service' import { BaseJob } from '../base-job' import { JobType } from '../enums/job-type.enum' +import { JobActivityType } from '../enums/job-activity-type.enum' import { JobResult } from '../interfaces/job-result.interface' import { JobContext } from '../interfaces/job.interface' import { JobRunner } from '../job-runner.service' @@ -35,6 +36,10 @@ export class RetryFailedHandler extends BaseJob { if (syncResult.failed > 0) { this.logger.warn( `ICM sweep partial failure: ${syncResult.synced} synced, ${syncResult.failed} failed`, + { + activityType: JobActivityType.ICM, + related: `ICM sweep partial failure (${syncResult.synced} synced, ${syncResult.failed} failed)`, + }, ) } } @@ -45,7 +50,10 @@ export class RetryFailedHandler extends BaseJob { metadata: { syncResult }, } } catch (error) { - this.logger.error(`Error processing failed jobs: ${error.message}`, error.stack) + this.logger.error(`Error processing failed jobs: ${error.message}`, { + activityType: JobActivityType.JOB, + related: error.message, + }) return { success: false, message: error.message, diff --git a/backend/src/jobs/job-activity-aggregator.spec.ts b/backend/src/jobs/job-activity-aggregator.spec.ts new file mode 100644 index 00000000..70fdc319 --- /dev/null +++ b/backend/src/jobs/job-activity-aggregator.spec.ts @@ -0,0 +1,53 @@ +import { describe, it, expect, vi } from 'vitest' +import { JobActivityAggregator } from './job-activity-aggregator' +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' + +describe('JobActivityAggregator', () => { + it('should flush a single occurrence unchanged', async () => { + const aggregator = new JobActivityAggregator() + const recordFn = vi.fn().mockResolvedValue(undefined) + + aggregator.note({ + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.CRA, + related: 'Invalid response file format (example: bad.rsp)', + aggregateKey: 'invalid-response-file-format', + }) + + await aggregator.flush(recordFn) + + expect(recordFn).toHaveBeenCalledWith({ + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.CRA, + related: 'Invalid response file format (example: bad.rsp)', + }) + }) + + it('should aggregate multiple notes into one row', async () => { + const aggregator = new JobActivityAggregator() + const recordFn = vi.fn().mockResolvedValue(undefined) + + aggregator.note({ + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.WKL, + related: 'Unexpected WKL transaction type (example: 99)', + aggregateKey: 'wkl-unexpected-transaction', + }) + aggregator.note({ + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.WKL, + related: 'Unexpected WKL transaction type (example: 99)', + aggregateKey: 'wkl-unexpected-transaction', + }) + + await aggregator.flush(recordFn) + + expect(recordFn).toHaveBeenCalledOnce() + expect(recordFn).toHaveBeenCalledWith({ + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.WKL, + related: '2 occurrences — Unexpected WKL transaction type (example: 99)', + }) + }) +}) diff --git a/backend/src/jobs/job-activity-aggregator.ts b/backend/src/jobs/job-activity-aggregator.ts new file mode 100644 index 00000000..a1485dc1 --- /dev/null +++ b/backend/src/jobs/job-activity-aggregator.ts @@ -0,0 +1,60 @@ +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' + +type ActivityBucket = { + severity: JobActivitySeverity + activityType: JobActivityType + count: number + sampleRelated: string +} + +export class JobActivityAggregator { + private readonly buckets = new Map() + + note(params: { + severity: JobActivitySeverity + activityType: JobActivityType + related: string + aggregateKey: string + }): void { + const existing = this.buckets.get(params.aggregateKey) + if (existing) { + existing.count += 1 + return + } + + this.buckets.set(params.aggregateKey, { + severity: params.severity, + activityType: params.activityType, + count: 1, + sampleRelated: params.related, + }) + } + + async flush( + recordFn: (params: { + severity: JobActivitySeverity + activityType: JobActivityType + related: string + }) => Promise, + ): Promise { + for (const bucket of this.buckets.values()) { + const related = + bucket.count === 1 + ? bucket.sampleRelated + : `${bucket.count} occurrences — ${bucket.sampleRelated}` + + await recordFn({ + severity: bucket.severity, + activityType: bucket.activityType, + related: related.slice(0, 512), + }) + } + + this.buckets.clear() + } + + get size(): number { + return this.buckets.size + } +} diff --git a/backend/src/jobs/job-activity.service.spec.ts b/backend/src/jobs/job-activity.service.spec.ts new file mode 100644 index 00000000..af12fe24 --- /dev/null +++ b/backend/src/jobs/job-activity.service.spec.ts @@ -0,0 +1,59 @@ +import type { TestingModule } from '@nestjs/testing' +import { Test } from '@nestjs/testing' +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' +import { JobActivityService } from './job-activity.service' +import { JobsService } from './jobs.service' + +describe('JobActivityService', () => { + let service: JobActivityService + let jobsService: JobsService + + beforeEach(async () => { + const module: TestingModule = await Test.createTestingModule({ + providers: [ + JobActivityService, + { + provide: JobsService, + useValue: { addActivity: vi.fn().mockResolvedValue({ id: 1 }) }, + }, + ], + }).compile() + + service = module.get(JobActivityService) + jobsService = module.get(JobsService) + }) + + it('records activity via JobsService', async () => { + await service.recordActivity({ + jobRunId: 5, + severity: JobActivitySeverity.CRITICAL, + activityType: JobActivityType.DATA_QUALITY, + related: '3 contacts skipped', + }) + + expect(jobsService.addActivity).toHaveBeenCalledWith(5, { + severity: JobActivitySeverity.CRITICAL, + type: JobActivityType.DATA_QUALITY, + related: '3 contacts skipped', + }) + }) + + it('truncates long related text', async () => { + const longText = 'x'.repeat(600) + + await service.recordActivity({ + jobRunId: 1, + severity: JobActivitySeverity.ERROR, + activityType: JobActivityType.JOB, + related: longText, + }) + + expect(jobsService.addActivity).toHaveBeenCalledWith( + 1, + expect.objectContaining({ + related: 'x'.repeat(512), + }), + ) + }) +}) diff --git a/backend/src/jobs/job-activity.service.ts b/backend/src/jobs/job-activity.service.ts new file mode 100644 index 00000000..24b7a940 --- /dev/null +++ b/backend/src/jobs/job-activity.service.ts @@ -0,0 +1,24 @@ +import { Injectable } from '@nestjs/common' +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' +import { JobsService } from './jobs.service' + +const MAX_RELATED_LENGTH = 512 + +@Injectable() +export class JobActivityService { + constructor(private readonly jobsService: JobsService) {} + + async recordActivity(params: { + jobRunId: number | null + severity: JobActivitySeverity + activityType: JobActivityType + related?: string + }): Promise { + await this.jobsService.addActivity(params.jobRunId, { + severity: params.severity, + type: params.activityType, + related: params.related?.slice(0, MAX_RELATED_LENGTH), + }) + } +} diff --git a/backend/src/jobs/job-execution.scope.ts b/backend/src/jobs/job-execution.scope.ts new file mode 100644 index 00000000..3fe4ea39 --- /dev/null +++ b/backend/src/jobs/job-execution.scope.ts @@ -0,0 +1,41 @@ +import { AsyncLocalStorage } from 'async_hooks' +import { JobActivityAggregator } from './job-activity-aggregator' + +type JobExecutionStore = { + jobRunId: number + aggregator: JobActivityAggregator + pendingWrites: Promise[] +} + +const storage = new AsyncLocalStorage() + +export function getCurrentJobRunId(): number | undefined { + return storage.getStore()?.jobRunId +} + +export function getJobActivityAggregator(): JobActivityAggregator | undefined { + return storage.getStore()?.aggregator +} + +export function trackPendingActivityWrite(write: Promise): void { + const store = storage.getStore() + if (store) { + store.pendingWrites.push(write) + } +} + +export async function flushPendingActivityWrites(): Promise { + const store = storage.getStore() + if (!store || store.pendingWrites.length === 0) { + return + } + + const pending = store.pendingWrites + store.pendingWrites = [] + await Promise.all(pending) +} + +export async function runWithJobScope(jobRunId: number, fn: () => Promise): Promise { + const aggregator = new JobActivityAggregator() + return storage.run({ jobRunId, aggregator, pendingWrites: [] }, fn) +} diff --git a/backend/src/jobs/job-monitoring-log.spec.ts b/backend/src/jobs/job-monitoring-log.spec.ts new file mode 100644 index 00000000..dfa79f30 --- /dev/null +++ b/backend/src/jobs/job-monitoring-log.spec.ts @@ -0,0 +1,89 @@ +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest' +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' +import { + persistJobMonitoringLog, + registerJobActivityRecorder, + resetJobActivityRecorder, + runWithJobExecutionScope, +} from './job-monitoring-log' + +describe('job-monitoring-log', () => { + const recordActivity = vi.fn().mockResolvedValue(undefined) + + beforeEach(() => { + recordActivity.mockClear() + registerJobActivityRecorder(recordActivity) + }) + + afterEach(() => { + resetJobActivityRecorder() + }) + + it('should write immediately when activityType is set', async () => { + await runWithJobExecutionScope(42, async () => { + await persistJobMonitoringLog('warn', 'ICM sync-back failed', { + activityType: JobActivityType.ICM, + related: 'ICM sync-back failed: timeout', + }) + }) + + expect(recordActivity).toHaveBeenCalledWith({ + jobRunId: 42, + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.ICM, + related: 'ICM sync-back failed: timeout', + }) + }) + + it('should aggregate category-tagged crit logs and flush at job end', async () => { + await runWithJobExecutionScope(7, async () => { + await persistJobMonitoringLog( + 'crit', + 'Skipping contact: empty/null in required fields [dob]', + { + category: 'DATA_QUALITY', + }, + ) + await persistJobMonitoringLog( + 'crit', + 'Skipping contact: empty/null in required fields [name]', + { + category: 'DATA_QUALITY', + }, + ) + }) + + expect(recordActivity).toHaveBeenCalledOnce() + expect(recordActivity).toHaveBeenCalledWith({ + jobRunId: 7, + severity: JobActivitySeverity.CRITICAL, + activityType: JobActivityType.DATA_QUALITY, + related: '2 occurrences — Skipping contact: empty/null in required fields [dob]', + }) + }) + + it('should no-op when no recorder is registered', async () => { + resetJobActivityRecorder() + + await persistJobMonitoringLog('error', 'Job failed', { + activityType: JobActivityType.JOB, + }) + + expect(recordActivity).not.toHaveBeenCalled() + }) + + it('should persist without a job run when no scope is active', async () => { + await persistJobMonitoringLog('warn', 'ICM sync-back failed after WKL reprocess', { + activityType: JobActivityType.ICM, + related: 'ICM sync-back failed after WKL reprocess: timeout', + }) + + expect(recordActivity).toHaveBeenCalledWith({ + jobRunId: null, + severity: JobActivitySeverity.WARNING, + activityType: JobActivityType.ICM, + related: 'ICM sync-back failed after WKL reprocess: timeout', + }) + }) +}) diff --git a/backend/src/jobs/job-monitoring-log.ts b/backend/src/jobs/job-monitoring-log.ts new file mode 100644 index 00000000..5379a69c --- /dev/null +++ b/backend/src/jobs/job-monitoring-log.ts @@ -0,0 +1,166 @@ +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' +import { + getCurrentJobRunId, + getJobActivityAggregator, + runWithJobScope, + trackPendingActivityWrite, + flushPendingActivityWrites, +} from './job-execution.scope' + +export type JobMonitoringLogMeta = { + activityType?: JobActivityType + category?: string + related?: string + aggregate?: boolean + aggregateKey?: string +} + +const CATEGORY_TO_ACTIVITY_TYPE: Record = { + DATA_QUALITY: JobActivityType.DATA_QUALITY, + JOB: JobActivityType.JOB, + CRA: JobActivityType.CRA, + WKL: JobActivityType.WKL, + ICM: JobActivityType.ICM, + BATCH: JobActivityType.BATCH, +} + +type RecordActivityFn = (params: { + jobRunId: number | null + severity: JobActivitySeverity + activityType: JobActivityType + related?: string +}) => Promise + +let recordActivityFn: RecordActivityFn | null = null + +export function registerJobActivityRecorder(fn: RecordActivityFn): void { + recordActivityFn = fn +} + +/** @internal test helper */ +export function resetJobActivityRecorder(): void { + recordActivityFn = null +} + +export function resolveActivityType(meta: JobMonitoringLogMeta): JobActivityType | undefined { + if (meta.activityType) { + return meta.activityType + } + + if (meta.category) { + return CATEGORY_TO_ACTIVITY_TYPE[meta.category] + } + + return undefined +} + +export function mapLogLevelToSeverity( + level: 'warn' | 'error' | 'crit' | 'alert', +): JobActivitySeverity { + switch (level) { + case 'warn': + return JobActivitySeverity.WARNING + case 'error': + return JobActivitySeverity.ERROR + case 'crit': + case 'alert': + return JobActivitySeverity.CRITICAL + } +} + +async function writeActivity(params: { + jobRunId: number | null + severity: JobActivitySeverity + activityType: JobActivityType + related?: string +}): Promise { + if (!recordActivityFn) { + return + } + + try { + await recordActivityFn(params) + } catch { + // Activity logging must not fail job execution. + } +} + +function queueImmediateActivity(params: { + jobRunId: number | null + severity: JobActivitySeverity + activityType: JobActivityType + related: string +}): void { + trackPendingActivityWrite(writeActivity(params)) +} + +export async function persistJobMonitoringLog( + level: 'warn' | 'error' | 'crit' | 'alert', + message: string, + meta?: JobMonitoringLogMeta, + explicitJobRunId?: number | null, +): Promise { + if (!meta) { + return + } + + const activityType = resolveActivityType(meta) + if (!activityType) { + return + } + + const jobRunId = + explicitJobRunId !== undefined ? explicitJobRunId : (getCurrentJobRunId() ?? null) + + const severity = mapLogLevelToSeverity(level) + const related = (meta.related ?? message).slice(0, 512) + const shouldAggregate = meta.aggregate ?? (level === 'crit' && !!meta.category) + + if (shouldAggregate) { + const aggregator = getJobActivityAggregator() + if (aggregator) { + aggregator.note({ + severity, + activityType, + related, + aggregateKey: meta.aggregateKey ?? `${severity}:${activityType}`, + }) + return + } + } + + queueImmediateActivity({ jobRunId, severity, activityType, related }) +} + +export async function flushJobMonitoringAggregates(): Promise { + const aggregator = getJobActivityAggregator() + const jobRunId = getCurrentJobRunId() ?? null + + if (!aggregator || aggregator.size === 0) { + return + } + + await aggregator.flush(async (bucket) => { + await writeActivity({ + jobRunId, + severity: bucket.severity, + activityType: bucket.activityType, + related: bucket.related, + }) + }) +} + +export async function runWithJobExecutionScope( + jobRunId: number, + fn: () => Promise, +): Promise { + return runWithJobScope(jobRunId, async () => { + try { + return await fn() + } finally { + await flushPendingActivityWrites() + await flushJobMonitoringAggregates() + } + }) +} diff --git a/backend/src/jobs/job-monitoring.utils.spec.ts b/backend/src/jobs/job-monitoring.utils.spec.ts new file mode 100644 index 00000000..c3b24b0c --- /dev/null +++ b/backend/src/jobs/job-monitoring.utils.spec.ts @@ -0,0 +1,136 @@ +import { JobStatus } from './enums/job-status.enum' +import { JobTrigger } from './enums/job-trigger.enum' +import { JobType } from './enums/job-type.enum' +import { + formatJobDisplayName, + formatJobSummary, + formatMonitoringStatus, + formatTriggeredBy, + MONITORED_JOB_HISTORY_TYPES, + MONITORED_JOB_LIST_TYPES, +} from './job-monitoring.utils' + +describe('job-monitoring.utils', () => { + it('lists six FDD job types for the job list', () => { + expect(MONITORED_JOB_HISTORY_TYPES).toHaveLength(13) + expect(MONITORED_JOB_LIST_TYPES).toEqual([ + JobType.INGEST_ICM, + JobType.INGEST_MIS, + JobType.RUN_ELIGIBILITY, + JobType.AUTO_BATCH, + JobType.SEND_CRA_FILE, + JobType.POLL_CRA_RESPONSE, + ]) + }) + + it('maps job types to display names', () => { + expect(formatJobDisplayName(JobType.INGEST_ICM)).toBe('Data Fetch - ICM') + expect(formatJobDisplayName(JobType.RUN_ELIGIBILITY)).toBe('Eligibility') + }) + + it('maps job status to FDD labels', () => { + expect(formatMonitoringStatus(JobStatus.SUCCESS)).toBe('Success') + expect(formatMonitoringStatus(JobStatus.RUNNING)).toBe('Running') + expect(formatMonitoringStatus(JobStatus.FAILED)).toBe('Failed') + }) + + it('returns IDIR for end-user jobs and SYSTEM otherwise', () => { + expect(formatTriggeredBy({ jobTrigger: JobTrigger.END_USER, triggeredByUser: 'jsmith' })).toBe( + 'JSMITH', + ) + expect(formatTriggeredBy({ jobTrigger: JobTrigger.CRON, triggeredByUser: null })).toBe('SYSTEM') + }) + + it('returns Job failed for failed runs', () => { + expect( + formatJobSummary({ + jobType: JobType.SEND_CRA_FILE, + status: JobStatus.FAILED, + metadata: null, + }), + ).toBe('Job failed') + }) + + it('returns null for running jobs', () => { + expect( + formatJobSummary({ + jobType: JobType.RUN_ELIGIBILITY, + status: JobStatus.RUNNING, + metadata: { processed: 1 }, + }), + ).toBeNull() + }) + + it('formats eligibility summary from metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.RUN_ELIGIBILITY, + status: JobStatus.SUCCESS, + metadata: { processed: 10, statusChanges: 2, newContacts: 1, skipped: 3 }, + }), + ).toBe('10 processed, 2 updated, 1 new, 3 skipped') + }) + + it('formats ICM ingestion summary from metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.INGEST_ICM, + status: JobStatus.SUCCESS, + metadata: { totalFetched: 100, totalUpserted: 95 }, + }), + ).toBe('100 fetched, 95 upserted') + }) + + it('formats MIS ingestion summary from metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.INGEST_MIS, + status: JobStatus.SUCCESS, + metadata: { totalRows: 500, results: [{}, {}, {}] }, + }), + ).toBe('500 rows loaded across 3 files') + }) + + it('formats auto-batch summary from metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.AUTO_BATCH, + status: JobStatus.SUCCESS, + metadata: { application: 4, cancellation: 2 }, + }), + ).toBe('4 application, 2 cancellation') + }) + + it('formats weekly response summary from metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.POLL_CRA_RESPONSE, + status: JobStatus.SUCCESS, + metadata: { files_processed: 2, records_accepted: 10, records_rejected: 1 }, + }), + ).toBe('2 files processed, 10 accepted, 1 rejected') + }) + + it('formats send CRA no-batch summary from metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.SEND_CRA_FILE, + status: JobStatus.SUCCESS, + metadata: { no_batch: true }, + }), + ).toBe('No batch to process') + }) + + it('formats ingest data summary from child job metadata', () => { + expect( + formatJobSummary({ + jobType: JobType.INGEST_DATA, + status: JobStatus.SUCCESS, + metadata: { + icmResult: { success: true, metadata: { totalFetched: 50, totalUpserted: 48 } }, + misResult: { success: true, metadata: { totalRows: 200, results: [{}, {}] } }, + }, + }), + ).toBe('ICM: 50 fetched, 48 upserted; MIS: 200 rows, 2 files') + }) +}) diff --git a/backend/src/jobs/job-monitoring.utils.ts b/backend/src/jobs/job-monitoring.utils.ts new file mode 100644 index 00000000..2de3b5e2 --- /dev/null +++ b/backend/src/jobs/job-monitoring.utils.ts @@ -0,0 +1,203 @@ +import { JobStatus } from './enums/job-status.enum' +import { JobTrigger } from './enums/job-trigger.enum' +import { JobType } from './enums/job-type.enum' + +/** FDD Job List (APL-10) — one row per type; ICM/MIS are child ingestion runs. */ +export const MONITORED_JOB_LIST_TYPES: JobType[] = [ + JobType.INGEST_ICM, + JobType.INGEST_MIS, + JobType.RUN_ELIGIBILITY, + JobType.AUTO_BATCH, + JobType.SEND_CRA_FILE, + JobType.POLL_CRA_RESPONSE, +] + +/** Ingestion child jobs — latest run is not filtered by parentJobId. */ +export const MONITORED_CHILD_JOB_TYPES: JobType[] = [JobType.INGEST_ICM, JobType.INGEST_MIS] + +/** FDD Job History (APL-11) — list types plus orchestrator, sync, retry, and backfills. */ +export const MONITORED_JOB_HISTORY_TYPES: JobType[] = [ + ...MONITORED_JOB_LIST_TYPES, + JobType.INGEST_DATA, + JobType.SYNC_ICM, + JobType.RETRY_FAILED, + JobType.BACKFILL_OOC_AGREEMENT_LINES, + JobType.BACKFILL_ICM_CASE_CLOSE_DATES, + JobType.BACKFILL_WKL_FILE_RECORDS, + JobType.BACKFILL_BATCH_EFFECTIVE_DATE_REASON, +] + +export const JOB_DISPLAY_NAMES: Record = { + [JobType.INGEST_DATA]: 'Data Fetch', + [JobType.INGEST_MIS]: 'Data Fetch - MIS', + [JobType.INGEST_ICM]: 'Data Fetch - ICM', + [JobType.RUN_ELIGIBILITY]: 'Eligibility', + [JobType.AUTO_BATCH]: 'Auto Batch', + [JobType.SEND_CRA_FILE]: 'Send CRA File', + [JobType.POLL_CRA_RESPONSE]: 'Weekly Response', + [JobType.SYNC_ICM]: 'ICM Sync-Back', + [JobType.RETRY_FAILED]: 'Retry Failed Jobs', +} + +const MONITORING_STATUS_LABELS: Record = { + [JobStatus.SUCCESS]: 'Success', + [JobStatus.RUNNING]: 'Running', + [JobStatus.FAILED]: 'Failed', +} + +export type JobRunForMonitoring = { + jobType: string + jobTrigger: string + triggeredByUser?: string | null + status: string + metadata?: unknown +} + +export function formatJobDisplayName(jobType: string): string { + return JOB_DISPLAY_NAMES[jobType] ?? jobType +} + +export function formatMonitoringStatus(status: string): string { + return MONITORING_STATUS_LABELS[status] ?? status +} + +export function formatTriggeredBy( + job: Pick, +): string { + if (job.jobTrigger === JobTrigger.END_USER && job.triggeredByUser) { + return job.triggeredByUser.toUpperCase() + } + return 'SYSTEM' +} + +export function formatJobSummary( + job: Pick, +): string | null { + if (job.status === JobStatus.FAILED) { + return 'Job failed' + } + + if (job.status === JobStatus.RUNNING) { + return null + } + + const metadata = (job.metadata as Record | null) ?? null + if (!metadata) { + return null + } + + switch (job.jobType) { + case JobType.INGEST_ICM: { + const fetched = metadata.totalFetched ?? metadata.fetched + const upserted = metadata.totalUpserted ?? metadata.upserted + if (fetched !== undefined && upserted !== undefined) { + return `${fetched} fetched, ${upserted} upserted` + } + break + } + case JobType.INGEST_MIS: { + const totalRows = metadata.totalRows + const fileCount = Array.isArray(metadata.results) ? metadata.results.length : undefined + if (totalRows !== undefined && fileCount !== undefined) { + return `${totalRows} rows loaded across ${fileCount} files` + } + break + } + case JobType.RUN_ELIGIBILITY: { + const processed = metadata.processed + const updated = metadata.statusChanges ?? metadata.updated + const newContacts = metadata.newContacts ?? metadata.new + const skipped = metadata.skipped + if (processed !== undefined) { + return `${processed} processed, ${updated ?? 0} updated, ${newContacts ?? 0} new, ${skipped ?? 0} skipped` + } + break + } + case JobType.AUTO_BATCH: { + const application = metadata.application + const cancellation = metadata.cancellation + if (application !== undefined && cancellation !== undefined) { + return `${application} application, ${cancellation} cancellation` + } + break + } + case JobType.SEND_CRA_FILE: { + if (metadata.no_batch === true) { + return 'No batch to process' + } + const batchId = metadata.batch_id + const recordCount = metadata.record_count + const contactsCount = metadata.contacts_count + if (batchId !== undefined) { + return `Batch ${batchId}, ${recordCount ?? 0} records, ${contactsCount ?? 0} contacts` + } + break + } + case JobType.POLL_CRA_RESPONSE: { + const filesProcessed = metadata.files_processed + if (filesProcessed !== undefined) { + if (filesProcessed === 0) { + return 'No new CRA response files to process' + } + const accepted = metadata.records_accepted ?? 0 + const rejected = metadata.records_rejected ?? 0 + return `${filesProcessed} files processed, ${accepted} accepted, ${rejected} rejected` + } + break + } + case JobType.SYNC_ICM: { + if (metadata.totalFlagged === 0) { + return 'No contacts flagged for ICM sync' + } + const synced = metadata.synced + const failed = metadata.failed + if (synced !== undefined || failed !== undefined) { + return `${synced ?? 0} synced, ${failed ?? 0} failed` + } + break + } + case JobType.RETRY_FAILED: { + const syncResult = metadata.syncResult as Record | null | undefined + if (syncResult && (syncResult.synced !== undefined || syncResult.failed !== undefined)) { + return `${syncResult.synced ?? 0} synced, ${syncResult.failed ?? 0} failed` + } + return 'Failed job processing completed' + } + case JobType.INGEST_DATA: { + const icm = metadata.icmResult as { metadata?: Record } | undefined + const mis = metadata.misResult as { metadata?: Record } | undefined + const parts: string[] = [] + + const icmFetched = icm?.metadata?.totalFetched + const icmUpserted = icm?.metadata?.totalUpserted + if (icmFetched !== undefined && icmUpserted !== undefined) { + parts.push(`ICM: ${icmFetched} fetched, ${icmUpserted} upserted`) + } + + const misRows = mis?.metadata?.totalRows + const misFileCount = Array.isArray(mis?.metadata?.results) + ? mis.metadata.results.length + : undefined + if (misRows !== undefined && misFileCount !== undefined) { + parts.push(`MIS: ${misRows} rows, ${misFileCount} files`) + } + + if (parts.length > 0) { + return parts.join('; ') + } + + if (metadata.icmResult && metadata.misResult) { + return 'Data ingestion completed successfully' + } + break + } + default: + break + } + + return null +} + +export function isMonitoredChildJobType(jobType: JobType): boolean { + return MONITORED_CHILD_JOB_TYPES.includes(jobType) +} diff --git a/backend/src/jobs/job-runner.service.ts b/backend/src/jobs/job-runner.service.ts index 08051b61..bbbe3242 100644 --- a/backend/src/jobs/job-runner.service.ts +++ b/backend/src/jobs/job-runner.service.ts @@ -5,6 +5,7 @@ import { JobType } from './enums/job-type.enum' import { JobResult } from './interfaces/job-result.interface' import { JobContext } from './interfaces/job.interface' import { JobRegistry } from './job-registry.service' +import { runWithJobExecutionScope } from './job-monitoring-log' import { JobsService } from './jobs.service' import { OpenshiftJobLauncher } from './openshift-job-launcher.service' @@ -55,6 +56,10 @@ export class JobRunner { // Execute a job with inline retry async executeJob(jobId: number): Promise { + return runWithJobExecutionScope(jobId, () => this.executeJobScoped(jobId)) + } + + private async executeJobScoped(jobId: number): Promise { const job = await this.jobsService.getJob(jobId) if (!job) { throw new Error(`Job ${jobId} not found`) diff --git a/backend/src/jobs/jobs.module.ts b/backend/src/jobs/jobs.module.ts index d1d8b336..f6b72cd8 100644 --- a/backend/src/jobs/jobs.module.ts +++ b/backend/src/jobs/jobs.module.ts @@ -2,8 +2,10 @@ import { Module, OnModuleInit } from '@nestjs/common' import { PrismaModule } from 'src/common/database/prisma.module' import { IcmSyncBackModule } from 'src/sync/icm/icm-sync-back.module' import { RetryFailedHandler } from './handlers/retry-failed.handler' +import { JobActivityService } from './job-activity.service' import { JobRegistry } from './job-registry.service' import { JobRunner } from './job-runner.service' +import { registerJobActivityRecorder } from './job-monitoring-log' import { JobsService } from './jobs.service' import { OpenshiftJobLauncher } from './openshift-job-launcher.service' @@ -19,16 +21,32 @@ import { OpenshiftJobLauncher } from './openshift-job-launcher.service' */ @Module({ imports: [PrismaModule, IcmSyncBackModule], - providers: [JobsService, JobRunner, JobRegistry, RetryFailedHandler, OpenshiftJobLauncher], - exports: [JobsService, JobRunner, JobRegistry, RetryFailedHandler, OpenshiftJobLauncher], + providers: [ + JobsService, + JobActivityService, + JobRunner, + JobRegistry, + RetryFailedHandler, + OpenshiftJobLauncher, + ], + exports: [ + JobsService, + JobActivityService, + JobRunner, + JobRegistry, + RetryFailedHandler, + OpenshiftJobLauncher, + ], }) export class JobsModule implements OnModuleInit { constructor( private readonly registry: JobRegistry, private readonly retryFailedHandler: RetryFailedHandler, + private readonly jobActivityService: JobActivityService, ) {} onModuleInit() { + registerJobActivityRecorder((params) => this.jobActivityService.recordActivity(params)) this.registry.register(this.retryFailedHandler.jobType, this.retryFailedHandler) } } diff --git a/backend/src/jobs/jobs.service.spec.ts b/backend/src/jobs/jobs.service.spec.ts index b48b33a3..8ede2075 100644 --- a/backend/src/jobs/jobs.service.spec.ts +++ b/backend/src/jobs/jobs.service.spec.ts @@ -1,6 +1,8 @@ import type { TestingModule } from '@nestjs/testing' import { Test } from '@nestjs/testing' import { PrismaService } from 'src/common/database/prisma.service' +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' import { JobStatus } from './enums/job-status.enum' import { JobTrigger } from './enums/job-trigger.enum' import { JobType } from './enums/job-type.enum' @@ -40,6 +42,11 @@ describe('JobsService', () => { update: vi.fn().mockResolvedValue(mockJobRun), updateMany: vi.fn().mockResolvedValue({ count: 1 }), }, + jobActivity: { + create: vi.fn().mockResolvedValue({ id: 1 }), + findMany: vi.fn().mockResolvedValue([{ id: 1 }]), + count: vi.fn().mockResolvedValue(1), + }, }, }, ], @@ -221,6 +228,132 @@ describe('JobsService', () => { completedAt: expect.any(Date), }), }) + expect(prisma.jobActivity.create).toHaveBeenCalledWith({ + data: expect.objectContaining({ + jobRunId: 1, + severity: JobActivitySeverity.ERROR, + type: JobActivityType.JOB, + related: 'Connection timeout', + }), + }) + }) + }) + + describe('monitoring', () => { + it('should return latest run per monitored job type', async () => { + const latestIcm = { + ...mockJobRun, + id: 3, + jobType: JobType.INGEST_ICM, + parentJobId: 99, + startedAt: new Date('2026-01-02T00:00:00Z'), + } + const latestEligibility = { + ...mockJobRun, + id: 4, + jobType: JobType.RUN_ELIGIBILITY, + parentJobId: null, + } + vi.spyOn(prisma.jobRun, 'findFirst') + .mockResolvedValueOnce(latestIcm as any) + .mockResolvedValueOnce(null) + .mockResolvedValueOnce(latestEligibility as any) + .mockResolvedValueOnce(null) + .mockResolvedValueOnce(null) + .mockResolvedValueOnce(null) + + const result = await service.getLatestJobsPerType() + + expect(result).toHaveLength(2) + expect(result.find((x) => x.jobType === JobType.INGEST_ICM)?.id).toBe(3) + expect(prisma.jobRun.findFirst).toHaveBeenCalledWith({ + where: { jobType: JobType.INGEST_ICM }, + orderBy: { startedAt: 'desc' }, + }) + expect(prisma.jobRun.findFirst).toHaveBeenCalledWith({ + where: { jobType: JobType.RUN_ELIGIBILITY, parentJobId: null }, + orderBy: { startedAt: 'desc' }, + }) + }) + + it('should store triggeredByUser when provided', async () => { + await service.createJob({ + jobType: JobType.RUN_ELIGIBILITY, + jobTrigger: JobTrigger.END_USER, + triggeredByUser: 'JSMITH', + }) + + expect(prisma.jobRun.create).toHaveBeenCalledWith({ + data: expect.objectContaining({ + triggeredByUser: 'JSMITH', + }), + }) + }) + + it('should filter monitoring history by SYSTEM using jobTrigger', async () => { + await service.getJobHistory({ triggeredBy: 'SYSTEM' }) + + expect(prisma.jobRun.findMany).toHaveBeenCalledWith( + expect.objectContaining({ + where: expect.objectContaining({ + jobTrigger: { in: [JobTrigger.CRON, JobTrigger.SYSTEM] }, + }), + }), + ) + }) + + it('should filter monitoring history by USER using jobTrigger', async () => { + await service.getJobHistory({ triggeredBy: 'USER' }) + + expect(prisma.jobRun.findMany).toHaveBeenCalledWith( + expect.objectContaining({ + where: expect.objectContaining({ + jobTrigger: JobTrigger.END_USER, + }), + }), + ) + }) + + it('should query history for monitored types including ICM/MIS child runs', async () => { + await service.getJobHistory() + + expect(prisma.jobRun.findMany).toHaveBeenCalledWith( + expect.objectContaining({ + where: expect.objectContaining({ + startedAt: { gte: expect.any(Date) }, + OR: [ + { jobType: { in: [JobType.INGEST_ICM, JobType.INGEST_MIS] } }, + { parentJobId: null }, + ], + }), + orderBy: { startedAt: 'desc' }, + }), + ) + }) + + it('should filter monitoring history by user IDIR', async () => { + await service.getJobHistory({ triggeredBy: 'jsmith' }) + + expect(prisma.jobRun.findMany).toHaveBeenCalledWith( + expect.objectContaining({ + where: expect.objectContaining({ + jobTrigger: JobTrigger.END_USER, + triggeredByUser: { equals: 'jsmith', mode: 'insensitive' }, + }), + }), + ) + }) + + it('should apply recent activity time window', async () => { + await service.getRecentActivities(1, 10) + + expect(prisma.jobActivity.findMany).toHaveBeenCalledWith( + expect.objectContaining({ + where: expect.objectContaining({ + when: { gte: expect.any(Date) }, + }), + }), + ) }) }) diff --git a/backend/src/jobs/jobs.service.ts b/backend/src/jobs/jobs.service.ts index 3f7a766a..2085fed1 100644 --- a/backend/src/jobs/jobs.service.ts +++ b/backend/src/jobs/jobs.service.ts @@ -1,15 +1,47 @@ import { Injectable } from '@nestjs/common' +import { Prisma } from '@prisma/client' import { PrismaService } from 'src/common/database/prisma.service' +import { JobActivitySeverity } from './enums/job-activity-severity.enum' +import { JobActivityType } from './enums/job-activity-type.enum' import { JobStatus } from './enums/job-status.enum' import { JobTrigger } from './enums/job-trigger.enum' import { JobType } from './enums/job-type.enum' +import { + isMonitoredChildJobType, + MONITORED_CHILD_JOB_TYPES, + MONITORED_JOB_HISTORY_TYPES, + MONITORED_JOB_LIST_TYPES, +} from './job-monitoring.utils' const RETRYABLE_END_USER_JOB_TYPES: JobType[] = [JobType.SEND_CRA_FILE] +export interface MonitoringHistoryFilters { + page?: number + limit?: number + jobType?: JobType + status?: JobStatus + triggeredBy?: string + jobId?: number + sortBy?: 'id' | 'jobType' | 'status' | 'jobTrigger' | 'startedAt' | 'completedAt' | 'createdAt' + sortOrder?: 'asc' | 'desc' +} + +export interface MonitoringActivityFilters { + page?: number + limit?: number + jobRunId?: number + severity?: JobActivitySeverity + type?: JobActivityType + sortBy?: 'when' | 'severity' | 'type' | 'jobRunId' + sortOrder?: 'asc' | 'desc' + fromWhen?: Date +} + export interface CreateJobDto { jobType: JobType jobTrigger: JobTrigger parentJobId?: number + triggeredByUser?: string metadata?: Record } @@ -19,11 +51,12 @@ export class JobsService { async createJob(dto: CreateJobDto) { const now = new Date() - return this.prisma.jobRun.create({ + const job = await this.prisma.jobRun.create({ data: { jobType: dto.jobType, jobTrigger: dto.jobTrigger, parentJobId: dto.parentJobId, + triggeredByUser: dto.triggeredByUser, status: JobStatus.RUNNING, retryCount: 0, metadata: (dto.metadata ?? {}) as any, @@ -31,6 +64,8 @@ export class JobsService { startedAt: now, }, }) + + return job } async getJobs(filters: { jobType?: JobType; status?: JobStatus; page?: number; limit?: number }) { @@ -62,7 +97,7 @@ export class JobsService { } async markSuccess(id: number, metadata?: Record) { - return this.prisma.jobRun.updateMany({ + const result = await this.prisma.jobRun.updateMany({ where: { id, status: JobStatus.RUNNING }, data: { status: JobStatus.SUCCESS, @@ -70,10 +105,12 @@ export class JobsService { ...(metadata && { metadata: metadata as any }), }, }) + + return result } async markFailed(id: number, error: string) { - return this.prisma.jobRun.updateMany({ + const result = await this.prisma.jobRun.updateMany({ where: { id, status: JobStatus.RUNNING }, data: { status: JobStatus.FAILED, @@ -82,6 +119,16 @@ export class JobsService { completedAt: new Date(), }, }) + + if (result.count > 0) { + await this.addActivity(id, { + severity: JobActivitySeverity.ERROR, + type: JobActivityType.JOB, + related: error.slice(0, 512), + }) + } + + return result } async resetToRunning(id: number) { @@ -138,12 +185,12 @@ export class JobsService { } async markStuckJobsAsFailed(stuckThresholdMinutes: number = 60) { - const threshold = new Date(Date.now() - stuckThresholdMinutes * 60 * 1000) - return this.prisma.jobRun.updateMany({ + const stuckJobs = await this.getStuckRunningJobs(stuckThresholdMinutes) + const result = await this.prisma.jobRun.updateMany({ where: { status: JobStatus.RUNNING, startedAt: { - lt: threshold, + lt: new Date(Date.now() - stuckThresholdMinutes * 60 * 1000), }, }, data: { @@ -152,10 +199,22 @@ export class JobsService { completedAt: new Date(), }, }) + + if (result.count > 0) { + for (const job of stuckJobs) { + await this.addActivity(job.id, { + severity: JobActivitySeverity.WARNING, + type: JobActivityType.JOB, + related: 'Job timed out (stuck)', + }) + } + } + + return result } async markStuckJobAsFailed(id: number, error: string = 'Job timed out (stuck)') { - return this.prisma.jobRun.updateMany({ + const result = await this.prisma.jobRun.updateMany({ where: { id, status: JobStatus.RUNNING }, data: { status: JobStatus.FAILED, @@ -163,6 +222,16 @@ export class JobsService { completedAt: new Date(), }, }) + + if (result.count > 0) { + await this.addActivity(id, { + severity: JobActivitySeverity.WARNING, + type: JobActivityType.JOB, + related: error.slice(0, 512), + }) + } + + return result } async getChildJobs(parentJobId: number) { @@ -215,4 +284,128 @@ export class JobsService { }) return failed !== null } + + async getLatestJobsPerType() { + const runs = await Promise.all( + MONITORED_JOB_LIST_TYPES.map((jobType) => + this.prisma.jobRun.findFirst({ + where: { + jobType, + ...(isMonitoredChildJobType(jobType) ? {} : { parentJobId: null }), + }, + orderBy: { startedAt: 'desc' }, + }), + ), + ) + + return runs.filter((run): run is NonNullable => run !== null) + } + + async getJobHistory(filters: MonitoringHistoryFilters = {}) { + const page = filters.page ?? 1 + const limit = filters.limit ?? 10 + const oneMonthAgo = new Date() + oneMonthAgo.setMonth(oneMonthAgo.getMonth() - 1) + + const monitoredTypes = filters.jobType ? [filters.jobType] : MONITORED_JOB_HISTORY_TYPES + + const where: Prisma.JobRunWhereInput = { + startedAt: { gte: oneMonthAgo }, + jobType: { in: monitoredTypes }, + OR: [{ jobType: { in: MONITORED_CHILD_JOB_TYPES } }, { parentJobId: null }], + ...(filters.status && { status: filters.status }), + ...(filters.jobId && { id: filters.jobId }), + } + + if (filters.triggeredBy) { + if (filters.triggeredBy === 'SYSTEM') { + where.jobTrigger = { in: [JobTrigger.CRON, JobTrigger.SYSTEM] } + } else if (filters.triggeredBy === 'USER' || filters.triggeredBy === JobTrigger.END_USER) { + where.jobTrigger = JobTrigger.END_USER + } else if ( + filters.triggeredBy === JobTrigger.CRON || + filters.triggeredBy === JobTrigger.SYSTEM + ) { + where.jobTrigger = filters.triggeredBy + } else { + where.jobTrigger = JobTrigger.END_USER + where.triggeredByUser = { equals: filters.triggeredBy, mode: 'insensitive' } + } + } + + const sortBy = filters.sortBy ?? 'startedAt' + const sortOrder = filters.sortOrder ?? 'desc' + + const [data, total] = await Promise.all([ + this.prisma.jobRun.findMany({ + where, + orderBy: { [sortBy]: sortOrder }, + skip: (page - 1) * limit, + take: limit, + }), + this.prisma.jobRun.count({ where }), + ]) + + return { data, total, page, limit } + } + + async addActivity( + jobRunId: number | null, + activity: { severity: JobActivitySeverity; type: JobActivityType; related?: string }, + ) { + return this.prisma.jobActivity.create({ + data: { + jobRunId, + severity: activity.severity, + type: activity.type, + related: activity.related, + when: new Date(), + }, + }) + } + + async getActivities(filters: MonitoringActivityFilters = {}) { + const page = filters.page ?? 1 + const limit = filters.limit ?? 10 + const sortBy = filters.sortBy ?? 'when' + const sortOrder = filters.sortOrder ?? 'desc' + + const where = { + ...(filters.jobRunId && { jobRunId: filters.jobRunId }), + ...(filters.severity && { severity: filters.severity }), + ...(filters.type && { type: filters.type }), + ...(filters.fromWhen && { when: { gte: filters.fromWhen } }), + } + + const [data, total] = await Promise.all([ + this.prisma.jobActivity.findMany({ + where, + orderBy: { [sortBy]: sortOrder }, + skip: (page - 1) * limit, + take: limit, + }), + this.prisma.jobActivity.count({ where }), + ]) + + return { data, total, page, limit } + } + + async getRecentActivities( + page: number = 1, + limit: number = 10, + filters: Omit = {}, + ) { + const oneMonthAgo = new Date() + oneMonthAgo.setMonth(oneMonthAgo.getMonth() - 1) + + return this.getActivities({ + page, + limit, + severity: filters.severity, + type: filters.type, + sortBy: filters.sortBy ?? 'when', + sortOrder: filters.sortOrder ?? 'desc', + fromWhen: oneMonthAgo, + }) + } } diff --git a/backend/src/jobs/openshift-job-launcher.service.ts b/backend/src/jobs/openshift-job-launcher.service.ts index d1795e7c..ecf45e41 100644 --- a/backend/src/jobs/openshift-job-launcher.service.ts +++ b/backend/src/jobs/openshift-job-launcher.service.ts @@ -45,7 +45,7 @@ export class OpenshiftJobLauncher { if (!namespace) { this.enabled = false this.namespace = 'local' - this.logger.warn( + this.logger.log( 'OpenShift job launcher disabled (in-cluster config loaded but namespace could not be read)', ) return @@ -89,7 +89,7 @@ export class OpenshiftJobLauncher { return readFileSync(namespacePath, 'utf8').trim() } } catch (error) { - this.logger.warn(`Could not read namespace from service account: ${(error as Error).message}`) + this.logger.log(`Could not read namespace from service account: ${(error as Error).message}`) } return null } diff --git a/backend/src/sync/eligibility/auto-batch.service.ts b/backend/src/sync/eligibility/auto-batch.service.ts index 953d30dd..f525e327 100644 --- a/backend/src/sync/eligibility/auto-batch.service.ts +++ b/backend/src/sync/eligibility/auto-batch.service.ts @@ -70,18 +70,6 @@ export class AutoBatchService { else cancellation++ } - if (result.skipped.length > 0) { - this.logger.log( - `Auto-batch skipped ${result.skipped.length} contacts (batch ${result.batch.id})`, - ) - } - - if (result.incomplete.length > 0) { - this.logger.log( - `Auto-batch Records Validation: ${result.incomplete.length} contacts auto-held due to missing CRA mandatory fields (batch ${result.batch.id})`, - ) - } - this.logger.log( `Auto-batched ${application} application + ${cancellation} cancellation contacts into batch ${result.batch.id}`, ) diff --git a/backend/src/sync/eligibility/eligibility.service.spec.ts b/backend/src/sync/eligibility/eligibility.service.spec.ts index d645340e..171e4fec 100644 --- a/backend/src/sync/eligibility/eligibility.service.spec.ts +++ b/backend/src/sync/eligibility/eligibility.service.spec.ts @@ -835,8 +835,17 @@ describe('EligibilityService', () => { await service.runForContact('ICM-ELIG') - expect(warnSpy).toHaveBeenCalledWith(expect.stringContaining('ICM-ELIG')) - expect(warnSpy).toHaveBeenCalledWith(expect.stringContaining('BL-14B/14C skip')) + expect(warnSpy).toHaveBeenCalledWith( + expect.stringContaining('ICM-ELIG'), + expect.objectContaining({ + activityType: 'DATA_QUALITY', + aggregateKey: 'user-set-missing-effective-date', + }), + ) + expect(warnSpy).toHaveBeenCalledWith( + expect.stringContaining('BL-14B/14C skip'), + expect.any(Object), + ) warnSpy.mockRestore() }) diff --git a/backend/src/sync/eligibility/eligibility.service.ts b/backend/src/sync/eligibility/eligibility.service.ts index 5a6e3187..59f9b54b 100644 --- a/backend/src/sync/eligibility/eligibility.service.ts +++ b/backend/src/sync/eligibility/eligibility.service.ts @@ -1,6 +1,7 @@ import { Injectable } from '@nestjs/common' import { PrismaService } from 'src/common/database/prisma.service' import { AppLogger } from 'src/common/logger/app-logger' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { getAgeCutoffDate, isEligibleAge, @@ -684,7 +685,12 @@ export class EligibilityService { for (const profile of profiles) { if (!profile.dateOfBirth) { - this.logger.warn(`Skipping contact ${profile.personIdIcm}: missing date of birth`) + this.logger.warn(`Skipping contact ${profile.personIdIcm}: missing date of birth`, { + activityType: JobActivityType.DATA_QUALITY, + aggregate: true, + aggregateKey: 'missing-dob', + related: 'Contact skipped: missing date of birth', + }) stats.skipped++ continue } @@ -775,6 +781,12 @@ export class EligibilityService { private warnUserSetWithoutEffectiveDate(profile: ContactProfile): void { this.logger.warn( `User-set CSA status for ${profile.personIdIcm} (last_updated_by=${profile.lastUpdatedBy}) but no csa_status_effective_date on master or ICM; running eligibility without BL-14B/14C skip`, + { + activityType: JobActivityType.DATA_QUALITY, + aggregate: true, + aggregateKey: 'user-set-missing-effective-date', + related: 'User-set CSA status without effective date', + }, ) } diff --git a/backend/src/sync/handlers/auto-batch.handler.ts b/backend/src/sync/handlers/auto-batch.handler.ts index a3fe794e..7d4d879b 100644 --- a/backend/src/sync/handlers/auto-batch.handler.ts +++ b/backend/src/sync/handlers/auto-batch.handler.ts @@ -1,6 +1,7 @@ import { Injectable } from '@nestjs/common' import { BaseJob } from 'src/jobs/base-job' import { JobType } from 'src/jobs/enums/job-type.enum' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobResult } from 'src/jobs/interfaces/job-result.interface' import { JobContext } from 'src/jobs/interfaces/job.interface' import { IcmSyncBackService, SyncBackResult } from '../icm/icm-sync-back.service' @@ -29,7 +30,10 @@ export class AutoBatchHandler extends BaseJob { try { syncResult = await this.icmSyncBackService.syncFlaggedWithRetry() } catch (err) { - this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`) + this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`, { + activityType: JobActivityType.ICM, + related: `ICM sync-back failed: ${(err as Error).message}`, + }) } } diff --git a/backend/src/sync/handlers/ingest-data.handler.spec.ts b/backend/src/sync/handlers/ingest-data.handler.spec.ts index e381bb1e..2f55d82b 100644 --- a/backend/src/sync/handlers/ingest-data.handler.spec.ts +++ b/backend/src/sync/handlers/ingest-data.handler.spec.ts @@ -45,6 +45,10 @@ describe('IngestDataHandler', () => { expect(result.success).toBe(true) expect(result.message).toBe('Data ingestion completed successfully') + expect(result.metadata).toEqual({ + icmResult: { success: true }, + misResult: { success: true }, + }) expect(runJobTypeSpy).toHaveBeenCalledTimes(2) expect(runJobTypeSpy).toHaveBeenCalledWith(JobType.INGEST_ICM, JobTrigger.CRON, { parentJobId: 1, diff --git a/backend/src/sync/handlers/ingest-data.handler.ts b/backend/src/sync/handlers/ingest-data.handler.ts index 72abb71a..1a57db07 100644 --- a/backend/src/sync/handlers/ingest-data.handler.ts +++ b/backend/src/sync/handlers/ingest-data.handler.ts @@ -4,6 +4,7 @@ import { JobTrigger } from 'src/jobs/enums/job-trigger.enum' import { JobType } from 'src/jobs/enums/job-type.enum' import { JobContext } from 'src/jobs/interfaces/job.interface' import { JobResult } from 'src/jobs/interfaces/job-result.interface' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobRunner } from 'src/jobs/job-runner.service' /* @@ -29,6 +30,10 @@ export class IngestDataHandler extends BaseJob { ]) if (!icmResult.success || !misResult.success) { + this.logger.error('Ingestion failed', { + activityType: JobActivityType.JOB, + related: `Data ingestion failed (ICM success=${icmResult.success}, MIS success=${misResult.success})`, + }) return { success: false, message: 'Ingestion failed', @@ -39,6 +44,7 @@ export class IngestDataHandler extends BaseJob { return { success: true, message: 'Data ingestion completed successfully', + metadata: { icmResult, misResult }, } } } diff --git a/backend/src/sync/handlers/run-eligibility.handler.ts b/backend/src/sync/handlers/run-eligibility.handler.ts index 974cc084..6796ad18 100644 --- a/backend/src/sync/handlers/run-eligibility.handler.ts +++ b/backend/src/sync/handlers/run-eligibility.handler.ts @@ -1,12 +1,13 @@ import { Injectable } from '@nestjs/common' import { ConfigService } from '@nestjs/config' import { BaseJob } from 'src/jobs/base-job' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobType } from 'src/jobs/enums/job-type.enum' import { JobResult } from 'src/jobs/interfaces/job-result.interface' import { JobContext } from 'src/jobs/interfaces/job.interface' import { JobsService } from 'src/jobs/jobs.service' -import { IcmSyncBackService, SyncBackResult } from '../icm/icm-sync-back.service' import { EligibilityService } from '../eligibility/eligibility.service' +import { IcmSyncBackService, SyncBackResult } from '../icm/icm-sync-back.service' /* * Runs eligibility rules against staged data, then syncs flagged contacts back to ICM. @@ -29,11 +30,31 @@ export class RunEligibilityHandler extends BaseJob { const threshold = await this.computeThreshold() const result = await this.eligibilityService.run(threshold) + if (result.skipped > 0) { + this.logger.crit(`${result.skipped} contacts skipped (missing required fields)`, { + activityType: JobActivityType.DATA_QUALITY, + related: `${result.skipped} contacts skipped (missing required fields)`, + }) + } + let syncResult: SyncBackResult | null = null try { syncResult = await this.icmSyncBackService.syncFlaggedWithRetry() } catch (err) { - this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`) + this.logger.warn(`ICM sync-back failed: ${(err as Error).message}`, { + activityType: JobActivityType.ICM, + related: `ICM sync-back failed: ${(err as Error).message}`, + }) + } + + if (syncResult && syncResult.failed > 0) { + this.logger.warn( + `ICM sync-back partial failure (${syncResult.synced} synced, ${syncResult.failed} failed)`, + { + activityType: JobActivityType.ICM, + related: `ICM sync-back partial failure (${syncResult.synced} synced, ${syncResult.failed} failed)`, + }, + ) } return { diff --git a/backend/src/sync/handlers/sync-icm.handler.ts b/backend/src/sync/handlers/sync-icm.handler.ts index c3283ef8..d84546a0 100644 --- a/backend/src/sync/handlers/sync-icm.handler.ts +++ b/backend/src/sync/handlers/sync-icm.handler.ts @@ -1,6 +1,7 @@ import { Injectable } from '@nestjs/common' import { BaseJob } from 'src/jobs/base-job' import { JobType } from 'src/jobs/enums/job-type.enum' +import { JobActivityType } from 'src/jobs/enums/job-activity-type.enum' import { JobResult } from 'src/jobs/interfaces/job-result.interface' import { JobContext } from 'src/jobs/interfaces/job.interface' import { IcmSyncBackService } from '../icm/icm-sync-back.service' @@ -31,6 +32,10 @@ export class SyncIcmHandler extends BaseJob { const result = await this.icmSyncBackService.syncFlaggedContacts() if (result.failed > 0 && result.synced === 0) { + this.logger.error(`ICM sync failed: all ${result.failed} contacts failed`, { + activityType: JobActivityType.ICM, + related: `ICM sync failed: all ${result.failed} contacts failed`, + }) return { success: false, message: `ICM sync failed: all ${result.failed} contacts failed`, @@ -38,6 +43,16 @@ export class SyncIcmHandler extends BaseJob { } } + if (result.failed > 0) { + this.logger.warn( + `ICM sync partial failure: ${result.synced} synced, ${result.failed} failed`, + { + activityType: JobActivityType.ICM, + related: `ICM sync partial failure (${result.synced} synced, ${result.failed} failed)`, + }, + ) + } + return { success: true, message: `ICM sync complete: ${result.synced} synced, ${result.failed} failed`, diff --git a/backend/src/sync/icm/data-source/mock-icm-data-source.ts b/backend/src/sync/icm/data-source/mock-icm-data-source.ts index 79bdbc85..298f2500 100644 --- a/backend/src/sync/icm/data-source/mock-icm-data-source.ts +++ b/backend/src/sync/icm/data-source/mock-icm-data-source.ts @@ -13,7 +13,7 @@ export class MockIcmDataSource extends IcmDataSource { const mockFile = path.join(mockDir, `${config.name}.json`) if (!fs.existsSync(mockFile)) { - this.logger.warn(`Mock file not found for ${config.name}: ${mockFile}`) + this.logger.log(`Mock file not found for ${config.name}: ${mockFile}`) return [] } diff --git a/frontend/src/components/JobMonitoringTab.tsx b/frontend/src/components/JobMonitoringTab.tsx index 1a89fd09..5ce53e2b 100644 --- a/frontend/src/components/JobMonitoringTab.tsx +++ b/frontend/src/components/JobMonitoringTab.tsx @@ -1,431 +1,1116 @@ import AccessTimeIcon from '@mui/icons-material/AccessTime' +import ArrowDownwardIcon from '@mui/icons-material/ArrowDownward' +import ArrowUpwardIcon from '@mui/icons-material/ArrowUpward' import CheckCircleIcon from '@mui/icons-material/CheckCircle' +import ClearIcon from '@mui/icons-material/Clear' import ErrorOutlineIcon from '@mui/icons-material/ErrorOutline' +import FilterListIcon from '@mui/icons-material/FilterList' +import UnfoldMoreIcon from '@mui/icons-material/UnfoldMore' import WarningAmberIcon from '@mui/icons-material/WarningAmber' import { + Alert, Box, + Button, + CircularProgress, + IconButton, + Menu, + MenuItem, Pagination, Paper, + Select, Table, TableBody, TableCell, TableContainer, TableHead, TableRow, + TextField, Tooltip, Typography, } from '@mui/material' -import { useState } from 'react' - -interface JobListRow { - id: number - jobName: string - status: 'Success' | 'Running' | 'Failed' - triggerBy?: string - started: string - finished?: string - summary: string - warning?: string +import { useCallback, useEffect, useState, type MouseEvent } from 'react' +import { + getJobActivities, + getJobHistory, + getLatestJobs, + getRecentActivities, + type ActivityParams, + type JobActivityRow, + type JobHistoryParams, + type MonitoringJobRow, +} from '../service/jobs-service' + +const ITEMS_PER_PAGE = 10 +const RUNNING_JOB_POLL_MS = 30_000 + +// Map display job names to backend JobType enum values for server-side filtering +const JOB_NAME_TO_TYPE: Record = { + 'Data Fetch - MIS': 'INGEST_MIS', + 'Data Fetch - ICM': 'INGEST_ICM', + Eligibility: 'RUN_ELIGIBILITY', + 'Auto Batch': 'AUTO_BATCH', + 'Send CRA File': 'SEND_CRA_FILE', + 'Weekly Response': 'POLL_CRA_RESPONSE', } -interface JobHistoryRow { - id: number - jobName: string - status: 'Success' | 'Failed' | 'Running' - triggerBy?: string - started: string - finished?: string - summary: string - warning?: string +const MONITORED_JOB_NAMES = Object.keys(JOB_NAME_TO_TYPE) +const STATUSES = ['Success', 'Running', 'Failed'] +const STATUS_TO_API: Record = { + Success: 'SUCCESS', + Running: 'RUNNING', + Failed: 'FAILED', +} +const TRIGGER_OPTIONS = ['SYSTEM', 'USER'] +const ACTIVITY_SEVERITIES = ['ERROR', 'WARNING', 'CRITICAL'] +const ACTIVITY_TYPES = ['DATA_QUALITY', 'JOB', 'CRA', 'WKL', 'ICM', 'BATCH'] +const ACTIVITY_TYPE_LABELS: Record = { + DATA_QUALITY: 'Data quality', + JOB: 'Job', + CRA: 'CRA', + WKL: 'Weekly file (WKL)', + ICM: 'ICM', + BATCH: 'Batch', } -interface ActivityRow { - id: number - when: string - severity: 'Error' | 'Warning' - type: string - related: string - jobId?: number +/** Format a date string to PT timezone: yyyy-Mmm-dd HH:mm:ss */ +const formatDatePT = (dateStr: string | null | undefined): string => { + if (!dateStr) return '—' + try { + const date = new Date(dateStr) + if (isNaN(date.getTime())) return '—' + const tz = 'America/Vancouver' + const year = new Intl.DateTimeFormat('en-US', { timeZone: tz, year: 'numeric' }).format(date) + const month = new Intl.DateTimeFormat('en-US', { timeZone: tz, month: 'short' }).format(date) + const day = new Intl.DateTimeFormat('en-US', { timeZone: tz, day: '2-digit' }).format(date) + const time = new Intl.DateTimeFormat('en-US', { + timeZone: tz, + hour: '2-digit', + minute: '2-digit', + second: '2-digit', + hour12: false, + }).format(date) + return `${year}-${month}-${day} ${time}` + } catch { + return '—' + } } -const ITEMS_PER_PAGE = 10 +const normalizeStatus = (status: string): string => { + const map: Record = { + SUCCESS: 'Success', + FAILED: 'Failed', + RUNNING: 'Running', + Success: 'Success', + Failed: 'Failed', + Running: 'Running', + } + return map[status] ?? status +} -// Mock data for Job List -const mockJobListData: JobListRow[] = [ - { - id: 1, - jobName: 'Data fetch', - status: 'Success', - triggerBy: 'System', - started: 'Jul 12, 11:00 PM', - finished: 'Jul 12, 11:04 PM', - summary: '500 processed · 12 updated · 2 skipped', - }, - { - id: 2, - jobName: 'Eligibility', - status: 'Running', - triggerBy: 'User', - started: 'Jul 13, 12:01 AM', - summary: 'In progress', - }, - { - id: 3, - jobName: 'Send CRA file', - status: 'Failed', - started: 'Jul 11, 08:00 AM', - finished: 'Jul 11, 08:02 AM', - summary: '2140 processed · 31 updated · 6 skipped', - warning: 'Possible stuck run', - }, - { - id: 4, - jobName: 'Weekly response', - status: 'Success', - started: 'Jul 11, 09:00 AM', - finished: 'Jul 11, 09:03 AM', - summary: '1902 processed · 25 updated · 3 skipped', - }, - { - id: 5, - jobName: 'Auto Batch', - status: 'Success', - started: 'Jul 11, 09:00 AM', - finished: 'Jul 11, 09:03 AM', - summary: '1903 processed · 25 updated · 3 skipped', - }, -] - -// Mock data for Job History -const mockJobHistoryData: JobHistoryRow[] = [ - { - id: 1, - jobName: 'Data Fetch', - status: 'Success', - started: 'Jul 08, 08:00 AM', - finished: 'Jul 08, 08:03 AM', - summary: '2098 processed · 20 updated · 1 skipped', - }, - { - id: 2, - jobName: 'Eligibility', - status: 'Success', - started: 'Jul 09, 08:00 AM', - finished: 'Jul 09, 08:02 AM', - summary: '2110 processed · 27 updated · 2 skipped', - }, - { - id: 3, - jobName: 'Auto Batch', - status: 'Failed', - started: 'Jul 11, 08:00 AM', - finished: 'Jul 11, 08:02 AM', - summary: '2140 processed · 31 updated · 6 skipped', - warning: 'Possible stuck run', - }, - { - id: 4, - jobName: 'Send CRA file', - status: 'Running', - started: 'Jul 12, 08:00 AM', - summary: 'In progress', - }, - { - id: 5, - jobName: 'Data Fetch', - status: 'Success', - started: 'Jul 08, 08:00 AM', - finished: 'Jul 08, 08:03 AM', - summary: '2098 processed · 20 updated · 1 skipped', - }, - { - id: 6, - jobName: 'Eligibility', - status: 'Success', - started: 'Jul 09, 08:00 AM', - finished: 'Jul 09, 08:02 AM', - summary: '2110 processed · 27 updated · 2 skipped', - }, - { - id: 7, - jobName: 'Auto Batch', - status: 'Failed', - started: 'Jul 11, 08:00 AM', - finished: 'Jul 11, 08:02 AM', - summary: '2140 processed · 31 updated · 6 skipped', - warning: 'Possible stuck run', - }, - { - id: 8, - jobName: 'Send CRA file', - status: 'Running', - started: 'Jul 12, 08:00 AM', - summary: 'In progress', - }, - { - id: 9, - jobName: 'Data Fetch', - status: 'Success', - started: 'Jul 08, 08:00 AM', - finished: 'Jul 08, 08:03 AM', - summary: '2098 processed · 20 updated · 1 skipped', - }, - { - id: 10, - jobName: 'Eligibility', - status: 'Success', - started: 'Jul 09, 08:00 AM', - finished: 'Jul 09, 08:02 AM', - summary: '2110 processed · 27 updated · 2 skipped', - }, -] - -// Mock data for Activities/Audit Log -const mockActivitiesData: ActivityRow[] = [ - { - id: 1, - when: 'Jul 11, 08:01 AM', - severity: 'Error', - type: 'CRA', - related: 'Send CRA file run @ Jul 11 08:00', - jobId: 5, - }, - { - id: 2, - when: 'Jul 11, 08:01 AM', - severity: 'Warning', - type: 'Data quality', - related: 'Send CRA file run @ Jul 11 08:00', - jobId: 3, - }, - { - id: 3, - when: 'Jul 11, 08:02 AM', - severity: 'Error', - type: 'Job', - related: 'Send CRA file run @ Jul 11 08:00', - }, - { - id: 4, - when: 'Jul 10, 09:15 AM', - severity: 'Warning', - type: 'File Processing', - related: 'Weekly response run @ Jul 10 09:00', - jobId: 4, - }, - { - id: 5, - when: 'Jul 09, 08:30 AM', - severity: 'Error', - type: 'Data quality', - related: 'Auto Batch run @ Jul 09 08:00', - jobId: 7, - }, - { - id: 6, - when: 'Jul 08, 10:45 AM', - severity: 'Warning', - type: 'CRA', - related: 'Data Fetch run @ Jul 08 08:00', - jobId: 1, - }, - { - id: 7, - when: 'Jul 12, 01:00 AM', - severity: 'Error', - type: 'Job', - related: 'Eligibility run @ Jul 13 12:00', - jobId: 2, - }, - { - id: 8, - when: 'Jul 11, 09:10 AM', - severity: 'Warning', - type: 'Data quality', - related: 'Weekly response run @ Jul 11 09:00', - jobId: 4, - }, - { - id: 9, - when: 'Jul 07, 02:20 PM', - severity: 'Error', - type: 'File Processing', - related: 'Data Fetch run @ Jul 07 14:00', - jobId: 1, - }, - { - id: 10, - when: 'Jul 06, 11:40 AM', - severity: 'Warning', - type: 'CRA', - related: 'Send CRA file run @ Jul 06 11:00', - jobId: 3, - }, -] +const matchesTriggerFilter = (triggeredBy: string, filter: string): boolean => { + if (!filter) return true + if (filter === 'USER') return triggeredBy !== 'SYSTEM' + return triggeredBy === filter +} -const getStatusIcon = (status: string) => { - switch (status) { - case 'Success': - return - case 'Failed': - return - case 'Running': - return - default: - return null +const normalizeSeverity = (severity: string): string => { + const map: Record = { + ERROR: 'Error', + WARNING: 'Warning', + CRITICAL: 'Critical', } + return map[severity] ?? severity +} + +const getStatusIcon = (status: string) => { + const s = status.toUpperCase() + if (s === 'SUCCESS') return + if (s === 'FAILED') return + if (s === 'RUNNING') return + return null } const getSeverityIcon = (severity: string) => { - switch (severity) { - case 'Error': - return - case 'Warning': - return - default: - return null + const s = severity.toUpperCase() + if (s === 'ERROR' || s === 'CRITICAL') { + return ( + + ) + } + if (s === 'WARNING') return + return null +} + +const warningChip = (text: string) => ( + + + {text} + + +) + +interface SortableHeaderCellProps { + label: string + field: string + currentSortField: string + currentSortOrder: 'asc' | 'desc' + onSort: (field: string) => void + sortable?: boolean + onFilterClick?: (event: MouseEvent) => void + filterActive?: boolean +} + +function SortableHeaderCell({ + label, + field, + currentSortField, + currentSortOrder, + onSort, + sortable = true, + onFilterClick, + filterActive = false, +}: SortableHeaderCellProps) { + const isActive = currentSortField === field + return ( + onSort(field) : undefined} + > + + + {label} + {sortable && + (isActive ? ( + currentSortOrder === 'asc' ? ( + + ) : ( + + ) + ) : ( + + ))} + + {onFilterClick && ( + { + event.stopPropagation() + onFilterClick(event) + }} + sx={{ p: 0.5, color: filterActive ? '#1976d2' : '#666' }} + > + + + )} + + + ) +} + +const cellSx = { fontSize: '0.875rem' } +const nonSortHeaderSx = { fontWeight: 600, fontSize: '0.875rem', backgroundColor: '#f5f5f5' } + +type JobListFilterColumn = 'id' | 'jobName' | 'status' | 'triggeredBy' +type JobHistoryFilterColumn = 'id' | 'jobType' | 'status' | 'jobTrigger' +type ActivityFilterColumn = 'severity' | 'type' | 'jobRunId' + +const JOB_LIST_FILTER_LABELS: Record = { + id: 'Job ID', + jobName: 'Job Name', + status: 'Status', + triggeredBy: 'Trigger By', +} + +const JOB_HISTORY_FILTER_LABELS: Record = { + id: 'Job ID', + jobType: 'Job Name', + status: 'Status', + jobTrigger: 'Trigger By', +} + +const ACTIVITY_FILTER_LABELS: Record = { + severity: 'Severity', + type: 'Type', + jobRunId: 'Job ID', +} + +const JOB_LIST_DATE_FIELDS = new Set(['started', 'finished']) +const JOB_HISTORY_DATE_FIELDS = new Set(['startedAt', 'completedAt', 'createdAt']) +const ACTIVITIES_DATE_FIELDS = new Set(['when']) + +const toTimeOrNull = (value: unknown): number | null => { + if (!value) return null + const time = new Date(String(value)).getTime() + return Number.isNaN(time) ? null : time +} + +const compareWithNullsLast = ( + left: number | string | null, + right: number | string | null, +): number => { + if (left === null && right === null) return 0 + if (left === null) return 1 + if (right === null) return -1 + if (typeof left === 'number' && typeof right === 'number') return left - right + return String(left).localeCompare(String(right), undefined, { + numeric: true, + sensitivity: 'base', + }) +} + +const compareJobListRows = (a: MonitoringJobRow, b: MonitoringJobRow, field: string): number => { + const getSortValue = (row: MonitoringJobRow, sortField: string): unknown => { + switch (sortField) { + case 'id': + return row.id + case 'jobName': + return row.jobName + case 'status': + return row.status + case 'triggeredBy': + return row.triggeredBy + case 'started': + return row.started + case 'finished': + return row.finished + case 'summary': + return row.summary + case 'warning': + return row.warning + default: + return null + } + } + + if (field === 'id') { + return compareWithNullsLast(a.id ?? null, b.id ?? null) } + + if (JOB_LIST_DATE_FIELDS.has(field)) { + const aTime = toTimeOrNull(getSortValue(a, field)) + const bTime = toTimeOrNull(getSortValue(b, field)) + return compareWithNullsLast(aTime, bTime) + } + + const valA = getSortValue(a, field) + const valB = getSortValue(b, field) + return compareWithNullsLast( + valA == null ? null : String(valA), + valB == null ? null : String(valB), + ) } export default function JobMonitoringTab() { + // ── Job List ──────────────────────────────────────────────────────────── + const [jobListData, setJobListData] = useState([]) + const [jobListLoading, setJobListLoading] = useState(true) + const [jobListError, setJobListError] = useState(null) + // Job List: client-side filter state + const [jlFilterId, setJlFilterId] = useState('') + const [jlFilterName, setJlFilterName] = useState('') + const [jlFilterStatus, setJlFilterStatus] = useState('') + const [jlFilterTrigger, setJlFilterTrigger] = useState('') + const [jlSortField, setJlSortField] = useState('id') + const [jlSortOrder, setJlSortOrder] = useState<'asc' | 'desc'>('asc') + const [jlFilterAnchor, setJlFilterAnchor] = useState<{ + element: HTMLElement | null + column: JobListFilterColumn | '' + }>({ element: null, column: '' }) + + // ── Job History ───────────────────────────────────────────────────────── + const [jobHistoryData, setJobHistoryData] = useState([]) + const [jobHistoryTotal, setJobHistoryTotal] = useState(0) const [jobHistoryPage, setJobHistoryPage] = useState(1) + const [jobHistoryLoading, setJobHistoryLoading] = useState(true) + const [jobHistoryError, setJobHistoryError] = useState(null) + // Job History: server-side filter state (text inputs apply on Enter) + const [jhFilterId, setJhFilterId] = useState('') + const [jhAppliedFilterId, setJhAppliedFilterId] = useState('') + const [jhFilterJobName, setJhFilterJobName] = useState('') + const [jhFilterStatus, setJhFilterStatus] = useState('') + const [jhFilterTrigger, setJhFilterTrigger] = useState('') + const [jhSortField, setJhSortField] = useState('startedAt') + const [jhSortOrder, setJhSortOrder] = useState<'asc' | 'desc'>('desc') + const [selectedJobHistoryId, setSelectedJobHistoryId] = useState(null) + const [jhFilterAnchor, setJhFilterAnchor] = useState<{ + element: HTMLElement | null + column: JobHistoryFilterColumn | '' + }>({ element: null, column: '' }) + + // ── Activities ────────────────────────────────────────────────────────── + const [activitiesData, setActivitiesData] = useState([]) + const [activitiesTotal, setActivitiesTotal] = useState(0) const [activitiesPage, setActivitiesPage] = useState(1) + const [activitiesLoading, setActivitiesLoading] = useState(true) + const [activitiesError, setActivitiesError] = useState(null) + // Activities: server-side filter state (Job ID applies on Enter) + const [actFilterSeverity, setActFilterSeverity] = useState('') + const [actFilterType, setActFilterType] = useState('') + const [actFilterJobId, setActFilterJobId] = useState('') + const [actAppliedFilterJobId, setActAppliedFilterJobId] = useState('') + const [actSortField, setActSortField] = useState('when') + const [actSortOrder, setActSortOrder] = useState<'asc' | 'desc'>('desc') + const [actFilterAnchor, setActFilterAnchor] = useState<{ + element: HTMLElement | null + column: ActivityFilterColumn | '' + }>({ element: null, column: '' }) - // Pagination calculations - const jobHistoryTotalPages = Math.ceil(mockJobHistoryData.length / ITEMS_PER_PAGE) - const activitiesTotalPages = Math.ceil(mockActivitiesData.length / ITEMS_PER_PAGE) + // ── Data fetching ─────────────────────────────────────────────────────── + const fetchJobList = useCallback(async () => { + setJobListLoading(true) + setJobListError(null) + try { + const data = await getLatestJobs() + setJobListData(data) + } catch { + setJobListError('Failed to load job list. Please try again.') + } finally { + setJobListLoading(false) + } + }, []) - const jobHistoryPaginatedData = mockJobHistoryData.slice( - (jobHistoryPage - 1) * ITEMS_PER_PAGE, - jobHistoryPage * ITEMS_PER_PAGE, - ) + const fetchJobHistory = useCallback(async () => { + setJobHistoryLoading(true) + setJobHistoryError(null) + try { + const params: JobHistoryParams = { + page: jobHistoryPage, + limit: ITEMS_PER_PAGE, + sortBy: jhSortField, + sortOrder: jhSortOrder, + } + if (jhAppliedFilterId) params.jobId = Number(jhAppliedFilterId) + if (jhFilterJobName) params.jobType = JOB_NAME_TO_TYPE[jhFilterJobName] + if (jhFilterStatus) params.status = STATUS_TO_API[jhFilterStatus] ?? jhFilterStatus + if (jhFilterTrigger) params.triggeredBy = jhFilterTrigger + const result = await getJobHistory(params) + setJobHistoryData(result.data) + setJobHistoryTotal(result.total) + } catch { + setJobHistoryError('Failed to load job history. Please try again.') + } finally { + setJobHistoryLoading(false) + } + }, [ + jobHistoryPage, + jhAppliedFilterId, + jhFilterJobName, + jhFilterStatus, + jhFilterTrigger, + jhSortField, + jhSortOrder, + ]) - const activitiesPaginatedData = mockActivitiesData.slice( - (activitiesPage - 1) * ITEMS_PER_PAGE, - activitiesPage * ITEMS_PER_PAGE, - ) + const fetchActivities = useCallback(async () => { + setActivitiesLoading(true) + setActivitiesError(null) + try { + const params: ActivityParams = { + page: activitiesPage, + limit: ITEMS_PER_PAGE, + sortBy: actSortField, + sortOrder: actSortOrder, + } + if (actFilterSeverity) params.severity = actFilterSeverity + if (actFilterType) params.type = actFilterType + + // When a Job History row is selected, show activities for that specific job. + // When a Job ID filter is manually entered, show activities for that job. + // Otherwise show recent monitoring activities. + let result + if (selectedJobHistoryId) { + result = await getJobActivities(selectedJobHistoryId, params) + } else if (actAppliedFilterJobId) { + result = await getJobActivities(Number(actAppliedFilterJobId), params) + } else { + result = await getRecentActivities(params) + } + + setActivitiesData(result.data) + setActivitiesTotal(result.total) + } catch { + setActivitiesError('Failed to load activities. Please try again.') + } finally { + setActivitiesLoading(false) + } + }, [ + activitiesPage, + selectedJobHistoryId, + actAppliedFilterJobId, + actFilterSeverity, + actFilterType, + actSortField, + actSortOrder, + ]) + + useEffect(() => { + const timeoutId = window.setTimeout(() => { + void fetchJobList() + }, 0) + + return () => window.clearTimeout(timeoutId) + }, [fetchJobList]) + + useEffect(() => { + const timeoutId = window.setTimeout(() => { + void fetchJobHistory() + }, 0) + + return () => window.clearTimeout(timeoutId) + }, [fetchJobHistory]) + + useEffect(() => { + const timeoutId = window.setTimeout(() => { + void fetchActivities() + }, 0) + + return () => window.clearTimeout(timeoutId) + }, [fetchActivities]) + + const hasRunningJobs = jobListData.some((row) => row.status.toUpperCase() === 'RUNNING') + + useEffect(() => { + if (!hasRunningJobs) { + return + } + + const interval = setInterval(() => { + void fetchJobList() + void fetchJobHistory() + void fetchActivities() + }, RUNNING_JOB_POLL_MS) + + return () => clearInterval(interval) + }, [hasRunningJobs, fetchJobList, fetchJobHistory, fetchActivities]) + + // ── Sort handlers ──────────────────────────────────────────────────────── + const handleJlSort = (field: string) => { + if (jlSortField === field) { + setJlSortOrder((prev) => (prev === 'asc' ? 'desc' : 'asc')) + } else { + setJlSortField(field) + setJlSortOrder(JOB_LIST_DATE_FIELDS.has(field) ? 'desc' : 'asc') + } + } + + const handleJhSort = (field: string) => { + setJobHistoryPage(1) + if (jhSortField === field) { + setJhSortOrder((prev) => (prev === 'asc' ? 'desc' : 'asc')) + } else { + setJhSortField(field) + setJhSortOrder(JOB_HISTORY_DATE_FIELDS.has(field) ? 'desc' : 'asc') + } + } + + const handleActSort = (field: string) => { + setActivitiesPage(1) + if (actSortField === field) { + setActSortOrder((prev) => (prev === 'asc' ? 'desc' : 'asc')) + } else { + setActSortField(field) + setActSortOrder(ACTIVITIES_DATE_FIELDS.has(field) ? 'desc' : 'asc') + } + } + + const openJlFilter = (event: MouseEvent, column: JobListFilterColumn) => { + setJlFilterAnchor({ element: event.currentTarget, column }) + } + + const closeJlFilter = () => { + setJlFilterAnchor({ element: null, column: '' }) + } + + const openJhFilter = (event: MouseEvent, column: JobHistoryFilterColumn) => { + setJhFilterAnchor({ element: event.currentTarget, column }) + } + + const closeJhFilter = () => { + setJhFilterAnchor({ element: null, column: '' }) + } + + const openActFilter = (event: MouseEvent, column: ActivityFilterColumn) => { + setActFilterAnchor({ element: event.currentTarget, column }) + } + + const closeActFilter = () => { + setActFilterAnchor({ element: null, column: '' }) + } + + // ── Job History row selection ───────────────────────────────────────────── + const handleJobHistoryRowClick = (jobId: number) => { + setSelectedJobHistoryId((prev) => (prev === jobId ? null : jobId)) + setActivitiesPage(1) + } + + // ── Clear all filters ───────────────────────────────────────────────────── + const clearJobListFilters = () => { + setJlFilterId('') + setJlFilterName('') + setJlFilterStatus('') + setJlFilterTrigger('') + setJlSortField('id') + setJlSortOrder('asc') + } + + const clearJobHistoryFilters = () => { + setJhFilterId('') + setJhAppliedFilterId('') + setJhFilterJobName('') + setJhFilterStatus('') + setJhFilterTrigger('') + setJhSortField('startedAt') + setJhSortOrder('desc') + setJobHistoryPage(1) + setSelectedJobHistoryId(null) + } + + const clearActivitiesFilters = () => { + setActFilterSeverity('') + setActFilterType('') + setActFilterJobId('') + setActAppliedFilterJobId('') + setActSortField('when') + setActSortOrder('desc') + setActivitiesPage(1) + setSelectedJobHistoryId(null) + } + + // ── Job List: client-side filter + sort ────────────────────────────────── + const filteredJobList = jobListData + .filter((row) => { + if (jlFilterId && !String(row.id).includes(jlFilterId)) return false + if (jlFilterName && !row.jobName.toLowerCase().includes(jlFilterName.toLowerCase())) + return false + if (jlFilterStatus && row.status !== jlFilterStatus) return false + if (!matchesTriggerFilter(row.triggeredBy, jlFilterTrigger)) return false + return true + }) + .sort((a, b) => { + const cmp = compareJobListRows(a, b, jlSortField) + return jlSortOrder === 'asc' ? cmp : -cmp + }) + + const jhTotalPages = Math.ceil(jobHistoryTotal / ITEMS_PER_PAGE) + const actTotalPages = Math.ceil(activitiesTotal / ITEMS_PER_PAGE) + + const filterTextFieldProps = { + size: 'small' as const, + placeholder: 'Filter...', + inputProps: { style: { fontSize: '0.75rem', padding: '2px 6px' } }, + } + + const filterSelectSx = { fontSize: '0.75rem', minWidth: 110 } return ( - {/* Job List Table */} + {/* ── Job List ── */} - - Job List - + + + Job List + + + + + {jobListError && ( + + {jobListError} + + )} + - - Job ID - Job Name - Status - Trigger By - Started - Finished - Summary - Warning + + openJlFilter(event, 'id')} + filterActive={jlFilterId.length > 0} + /> + openJlFilter(event, 'jobName')} + filterActive={jlFilterName.length > 0} + /> + openJlFilter(event, 'status')} + filterActive={jlFilterStatus.length > 0} + /> + openJlFilter(event, 'triggeredBy')} + filterActive={jlFilterTrigger.length > 0} + /> + + + Summary + Warning - {mockJobListData.map((row) => ( - - {row.id} - {row.jobName} - - - {getStatusIcon(row.status)} - {row.status} - + {jobListLoading ? ( + + + - {row.triggerBy || '—'} - {row.started} - {row.finished || '—'} - {row.summary} - - {row.warning ? ( - - - {row.warning} - - - ) : ( - '—' - )} + + ) : filteredJobList.length === 0 ? ( + + + No jobs found - ))} + ) : ( + filteredJobList.map((row) => ( + + {row.id} + {row.jobName} + + + {getStatusIcon(row.status)} + {normalizeStatus(row.status)} + + + {row.triggeredBy || '—'} + {formatDatePT(row.started)} + {formatDatePT(row.finished)} + {row.summary || '—'} + + {row.warning ? warningChip(row.warning) : '—'} + + + )) + )}
+ + + + + + Filter by{' '} + {jlFilterAnchor.column ? JOB_LIST_FILTER_LABELS[jlFilterAnchor.column] : ''} + + + + {jlFilterAnchor.column === 'id' && ( + setJlFilterId(e.target.value)} + /> + )} + {jlFilterAnchor.column === 'jobName' && ( + setJlFilterName(e.target.value)} + /> + )} + {jlFilterAnchor.column === 'status' && ( + + )} + {jlFilterAnchor.column === 'triggeredBy' && ( + + )} + +
- {/* Job History Table */} + {/* ── Job History ── */} - - Job History - + + + Job History + + + + + {jobHistoryError && ( + + {jobHistoryError} + + )} + - - Job ID - Job Name - Status - Trigger By - Started (PT) - Finished (PT) - Summary - Warning + + openJhFilter(event, 'id')} + filterActive={jhAppliedFilterId.length > 0} + /> + openJhFilter(event, 'jobType')} + filterActive={jhFilterJobName.length > 0} + /> + openJhFilter(event, 'status')} + filterActive={jhFilterStatus.length > 0} + /> + openJhFilter(event, 'jobTrigger')} + filterActive={jhFilterTrigger.length > 0} + /> + + + Summary + Warning - {jobHistoryPaginatedData.map((row) => ( - - {row.id} - {row.jobName} - - - {getStatusIcon(row.status)} - {row.status} - + {jobHistoryLoading ? ( + + + - {row.triggerBy || '—'} - {row.started} - {row.finished || '—'} - {row.summary} - - {row.warning ? ( - - - {row.warning} - - - ) : ( - '—' - )} + + ) : jobHistoryData.length === 0 ? ( + + + No history found - ))} + ) : ( + jobHistoryData.map((row) => ( + handleJobHistoryRowClick(row.id)} + sx={{ + cursor: 'pointer', + '&.Mui-selected': { backgroundColor: 'rgba(25, 118, 210, 0.08)' }, + '&.Mui-selected:hover': { backgroundColor: 'rgba(25, 118, 210, 0.12)' }, + }} + > + {row.id} + {row.jobName} + + + {getStatusIcon(row.status)} + {normalizeStatus(row.status)} + + + {row.triggeredBy || '—'} + {formatDatePT(row.started)} + {formatDatePT(row.finished)} + {row.summary || '—'} + + {row.warning ? warningChip(row.warning) : '—'} + + + )) + )}
+ + + + + Filter by{' '} + {jhFilterAnchor.column ? JOB_HISTORY_FILTER_LABELS[jhFilterAnchor.column] : ''} + + + + {jhFilterAnchor.column === 'id' && ( + setJhFilterId(e.target.value)} + onKeyDown={(e) => { + if (e.key === 'Enter') { + setJhAppliedFilterId(jhFilterId) + setJobHistoryPage(1) + closeJhFilter() + } + }} + /> + )} + {jhFilterAnchor.column === 'jobType' && ( + + )} + {jhFilterAnchor.column === 'status' && ( + + )} + {jhFilterAnchor.column === 'jobTrigger' && ( + + )} + + + - Showing {jobHistoryPaginatedData.length} of {mockJobHistoryData.length} records + Showing {jobHistoryData.length} of {jobHistoryTotal} records - {jobHistoryTotalPages > 1 && ( + {jhTotalPages > 1 && ( setJobHistoryPage(page)} + onChange={(_, p) => setJobHistoryPage(p)} color="primary" showFirstButton showLastButton @@ -452,47 +1137,228 @@ export default function JobMonitoringTab() {
- {/* Activities/Audit Log Table */} + {/* ── Activities ── */} - - Activities - + + + Activities + {selectedJobHistoryId && ( + + — Filtered by Job #{selectedJobHistoryId} + + )} + + + + + {activitiesError && ( + + {activitiesError} + + )} + - - Header - When (PT) - Severity - Type - Related - Job ID + + + openActFilter(event, 'severity')} + filterActive={actFilterSeverity.length > 0} + /> + openActFilter(event, 'type')} + filterActive={actFilterType.length > 0} + /> + {/* Related is not sortable per FDD */} + Related + openActFilter(event, 'jobRunId')} + filterActive={actAppliedFilterJobId.length > 0} + /> - {activitiesPaginatedData.map((row) => ( - - {row.id} - {row.when} - - - {getSeverityIcon(row.severity)} - {row.severity} - + {activitiesLoading ? ( + + + - {row.type} - - - {row.related} - + + ) : activitiesData.length === 0 ? ( + + + No activities found - {row.jobId || '—'} - ))} + ) : ( + activitiesData.map((row) => ( + + {formatDatePT(row.when)} + + + {getSeverityIcon(row.severity)} + {normalizeSeverity(row.severity)} + + + {ACTIVITY_TYPE_LABELS[row.type] ?? row.type} + + {row.related ? ( + + {row.related} + + ) : ( + '—' + )} + + {row.jobRunId ?? '—'} + + )) + )}
+ + + + + Filter by{' '} + {actFilterAnchor.column ? ACTIVITY_FILTER_LABELS[actFilterAnchor.column] : ''} + + + + {actFilterAnchor.column === 'severity' && ( + + )} + {actFilterAnchor.column === 'type' && ( + + )} + {actFilterAnchor.column === 'jobRunId' && ( + setActFilterJobId(e.target.value)} + onKeyDown={(e) => { + if (e.key === 'Enter') { + setActAppliedFilterJobId(actFilterJobId) + setActivitiesPage(1) + closeActFilter() + } + }} + /> + )} + + + - Showing {activitiesPaginatedData.length} of {mockActivitiesData.length} records + Showing {activitiesData.length} of {activitiesTotal} records - {activitiesTotalPages > 1 && ( + {actTotalPages > 1 && ( setActivitiesPage(page)} + onChange={(_, p) => setActivitiesPage(p)} color="primary" showFirstButton showLastButton diff --git a/frontend/src/service/jobs-service.ts b/frontend/src/service/jobs-service.ts new file mode 100644 index 00000000..a0effe1d --- /dev/null +++ b/frontend/src/service/jobs-service.ts @@ -0,0 +1,119 @@ +import APIService from './api-service' + +export interface MonitoringJobRow { + id: number + jobName: string + status: string + triggeredBy: string + started: string | null + finished: string | null + summary: string | null + warning: string | null +} + +export interface JobActivityRow { + id: number + jobRunId: number | null + when: string + severity: string + type: string + related: string | null +} + +export interface PaginatedResponse { + data: T[] + total: number + page: number + limit: number +} + +export interface JobHistoryParams { + page?: number + limit?: number + jobType?: string + status?: string + jobId?: number + triggeredBy?: string + sortBy?: string + sortOrder?: 'asc' | 'desc' +} + +export interface ActivityParams { + page?: number + limit?: number + severity?: string + type?: string + sortBy?: string + sortOrder?: 'asc' | 'desc' +} + +/** + * Get the latest job run per monitored job type (Job List view) + */ +export const getLatestJobs = async (): Promise => { + const response = await APIService.getAxiosInstance().get('/jobs/monitoring/latest') + return response.data +} + +/** + * Get paginated job history for the last month (Job History view) + */ +export const getJobHistory = async ( + params: JobHistoryParams = {}, +): Promise> => { + const query: Record = {} + if (params.page !== undefined) query.page = params.page + if (params.limit !== undefined) query.limit = params.limit + if (params.jobType) query.jobType = params.jobType + if (params.status) query.status = params.status + if (params.jobId) query.jobId = params.jobId + if (params.triggeredBy) query.triggeredBy = params.triggeredBy + if (params.sortBy) query.sortBy = params.sortBy + if (params.sortOrder) query.sortOrder = params.sortOrder + + const response = await APIService.getAxiosInstance().get('/jobs/monitoring/history', { + params: query, + }) + return response.data +} + +/** + * Get paginated recent monitoring activities (default Activities view) + */ +export const getRecentActivities = async ( + params: ActivityParams = {}, +): Promise> => { + const query: Record = {} + if (params.page !== undefined) query.page = params.page + if (params.limit !== undefined) query.limit = params.limit + if (params.severity) query.severity = params.severity + if (params.type) query.type = params.type + if (params.sortBy) query.sortBy = params.sortBy + if (params.sortOrder) query.sortOrder = params.sortOrder + + const response = await APIService.getAxiosInstance().get('/jobs/monitoring/activities', { + params: query, + }) + return response.data +} + +/** + * Get paginated activities for a specific job run (when a Job History row is selected) + */ +export const getJobActivities = async ( + jobId: number, + params: ActivityParams = {}, +): Promise> => { + const query: Record = {} + if (params.page !== undefined) query.page = params.page + if (params.limit !== undefined) query.limit = params.limit + if (params.severity) query.severity = params.severity + if (params.type) query.type = params.type + if (params.sortBy) query.sortBy = params.sortBy + if (params.sortOrder) query.sortOrder = params.sortOrder + + const response = await APIService.getAxiosInstance().get(`/jobs/${jobId}/activities`, { + params: query, + }) + return response.data +} diff --git a/migrations/sql/V26__add_job_monitoring_activity_log.sql b/migrations/sql/V26__add_job_monitoring_activity_log.sql new file mode 100644 index 00000000..c1437cda --- /dev/null +++ b/migrations/sql/V26__add_job_monitoring_activity_log.sql @@ -0,0 +1,32 @@ +-- VW-02 Job Monitoring: user IDIR and activity log. + +ALTER TABLE csa.job_runs + ADD COLUMN IF NOT EXISTS triggered_by_user TEXT; + +COMMENT ON COLUMN csa.job_runs.triggered_by_user IS + 'IDIR username when job_trigger is END_USER; NULL for CRON/SYSTEM jobs'; + +-- Activity log table for job-level monitoring details. +CREATE TABLE IF NOT EXISTS csa.job_activities ( + id SERIAL PRIMARY KEY, + job_run_id INTEGER REFERENCES csa.job_runs(id) ON DELETE SET NULL, + "when" TIMESTAMPTZ NOT NULL DEFAULT NOW(), + severity TEXT NOT NULL, + type TEXT NOT NULL, + related TEXT +); + +COMMENT ON COLUMN csa.job_activities.job_run_id IS + 'Associated job run when the activity occurred during a job; NULL for standalone operator actions'; + +CREATE INDEX IF NOT EXISTS idx_job_runs_created_at_desc + ON csa.job_runs (created_at DESC); + +CREATE INDEX IF NOT EXISTS idx_job_activities_job_run_id_when + ON csa.job_activities (job_run_id, "when" DESC); + +CREATE INDEX IF NOT EXISTS idx_job_activities_severity + ON csa.job_activities (severity); + +CREATE INDEX IF NOT EXISTS idx_job_activities_type + ON csa.job_activities (type);