From 0e3543bb0133b6e1d145dc8e90a996690277f40f Mon Sep 17 00:00:00 2001 From: saqib mir Date: Wed, 9 Sep 2026 17:17:09 +0530 Subject: [PATCH] =?UTF-8?q?fix(016-load-concurrency-testing):=20US2=20?= =?UTF-8?q?=E2=80=94=20version-guard=20SLA=20pause/resume/sweep?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit SlaRunRepository.updateWithVersion replaces the old unguarded update(), mirroring TicketsRepository.updateStatus's exact atomic-updateMany pattern. pause/resume/complete now retry (bounded) against fresh state on a version conflict; the breach sweep skips a run that lost the race to a concurrent pause/resume/complete rather than clobbering it, deferring to the next scheduled pass. Verified against real Postgres: concurrent pause/resume/sweep activity against the same run now always leaves it in one internally-consistent state, across repeated runs. Co-Authored-By: Claude Sonnet 5 --- .../sla/repository/sla-run.repository.ts | 22 +- .../orchestration/sla/service/sla.service.ts | 97 +++++--- tests/concurrency/sla-race.test.ts | 218 ++++++++++++++++++ 3 files changed, 306 insertions(+), 31 deletions(-) create mode 100644 tests/concurrency/sla-race.test.ts diff --git a/src/modules/orchestration/sla/repository/sla-run.repository.ts b/src/modules/orchestration/sla/repository/sla-run.repository.ts index 9a84d3b..f7a4ca1 100644 --- a/src/modules/orchestration/sla/repository/sla-run.repository.ts +++ b/src/modules/orchestration/sla/repository/sla-run.repository.ts @@ -21,8 +21,26 @@ export class SlaRunRepository { return this.prisma.sLARun.findUnique({ where: { ticketId } }); } - async update(id: string, data: Prisma.SLARunUpdateInput): Promise { - return this.prisma.sLARun.update({ where: { id }, data }); + /** + * 016-load-concurrency-testing research.md §2: optimistic concurrency, identical shape to + * TicketsRepository.updateStatus (003-ticketing) — an atomic single-statement `updateMany` + * conditioned on the row's `version` still matching `expectedVersion`. Replaces the old plain + * `update(id, data)`, which let two racing callers (e.g. `resume` and the breach sweep + * evaluating the same run at once) silently clobber each other's writes. Returns `null` on a + * stale-version mismatch — the caller re-reads and retries, exactly like `TicketsService` + * already does for ticket-status conflicts. + */ + async updateWithVersion( + id: string, + expectedVersion: number, + data: Prisma.SLARunUpdateInput, + ): Promise { + const result = await this.prisma.sLARun.updateMany({ + where: { id, version: expectedVersion }, + data: { ...data, version: { increment: 1 } } as Prisma.SLARunUncheckedUpdateManyInput, + }); + if (result.count === 0) return null; + return this.prisma.sLARun.findUnique({ where: { id } }); } /** research.md "Breach detection — one repeatable BullMQ job": every running run whose diff --git a/src/modules/orchestration/sla/service/sla.service.ts b/src/modules/orchestration/sla/service/sla.service.ts index 2eaeaec..29c57d8 100644 --- a/src/modules/orchestration/sla/service/sla.service.ts +++ b/src/modules/orchestration/sla/service/sla.service.ts @@ -98,13 +98,28 @@ export class SlaService { }); } + // 016-load-concurrency-testing research.md §2: pause/resume/complete each read-then-write a + // run with no guard, so two of them racing (or one racing the sweep below) could silently + // clobber each other — e.g. a resume reading a run just before the sweep marks it breached, + // then overwriting that breach back to 'running' moments later. Bounded to a handful of + // attempts, re-reading fresh state each time (mirroring TicketsService's own version-conflict + // handling), so a losing attempt still applies correctly against the winner's result instead + // of being silently dropped or corrupting state. + private static readonly MAX_UPDATE_ATTEMPTS = 5; + /** FR-007: pausing on WAITING_FOR_CUSTOMER records pausedAt and flips status — no-ops if * there's no run or it isn't currently running. */ async pause(ticketId: string): Promise { - const run = await this.runs.findByTicketId(ticketId); - if (!run || run.status !== 'running') return; + for (let attempt = 1; attempt <= SlaService.MAX_UPDATE_ATTEMPTS; attempt++) { + const run = await this.runs.findByTicketId(ticketId); + if (!run || run.status !== 'running') return; - await this.runs.update(run.id, { status: 'paused', pausedAt: new Date() }); + const updated = await this.runs.updateWithVersion(run.id, run.version, { + status: 'paused', + pausedAt: new Date(), + }); + if (updated) return; + } } /** @@ -113,38 +128,49 @@ export class SlaService { * durability mechanism (no separate remaining-minutes bookkeeping, no in-memory state). */ async resume(ticketId: string): Promise { - const run = await this.runs.findByTicketId(ticketId); - if (!run || run.status !== 'paused' || !run.pausedAt) return; + for (let attempt = 1; attempt <= SlaService.MAX_UPDATE_ATTEMPTS; attempt++) { + const run = await this.runs.findByTicketId(ticketId); + if (!run || run.status !== 'paused' || !run.pausedAt) return; - const pausedMs = Date.now() - run.pausedAt.getTime(); - await this.runs.update(run.id, { - status: 'running', - pausedAt: null, - resumedAt: new Date(), - firstResponseDueAt: run.firstResponseDueAt - ? new Date(run.firstResponseDueAt.getTime() + pausedMs) - : null, - resolutionDueAt: run.resolutionDueAt - ? new Date(run.resolutionDueAt.getTime() + pausedMs) - : null, - }); + const pausedMs = Date.now() - run.pausedAt.getTime(); + const updated = await this.runs.updateWithVersion(run.id, run.version, { + status: 'running', + pausedAt: null, + resumedAt: new Date(), + firstResponseDueAt: run.firstResponseDueAt + ? new Date(run.firstResponseDueAt.getTime() + pausedMs) + : null, + resolutionDueAt: run.resolutionDueAt + ? new Date(run.resolutionDueAt.getTime() + pausedMs) + : null, + }); + if (updated) return; + } } /** FR-010: a run that resolves before its due date is marked completed and is never later * flagged breached (the breach sweep only ever looks at status: 'running' runs). */ async complete(ticketId: string): Promise { - const run = await this.runs.findByTicketId(ticketId); - if (!run || run.status === 'completed') return; + for (let attempt = 1; attempt <= SlaService.MAX_UPDATE_ATTEMPTS; attempt++) { + const run = await this.runs.findByTicketId(ticketId); + if (!run || run.status === 'completed') return; - // 014-full-observability data-model.md #6: read BEFORE the update below — a run already - // 'breached' by the time it resolves was already counted breached by the sweep and must - // never also be counted 'met' here, even though this update still (pre-existing behavior, - // unrelated to this feature — see research.md §5) overwrites its status to 'completed'. - if (run.status !== 'breached') { - slaRunOutcomesCounter.inc({ outcome: 'met' }); + const updated = await this.runs.updateWithVersion(run.id, run.version, { + status: 'completed', + completedAt: new Date(), + }); + if (updated) { + // 014-full-observability data-model.md #6: gated on this same successful transition's + // pre-update status (not re-read afterward) — a run already 'breached' by the time it + // resolves was already counted breached by the sweep and must never also be counted + // 'met' here, even though this update still (pre-existing behavior, unrelated to this + // feature — see research.md §5) overwrites its status to 'completed'. + if (run.status !== 'breached') { + slaRunOutcomesCounter.inc({ outcome: 'met' }); + } + return; + } } - - await this.runs.update(run.id, { status: 'completed', completedAt: new Date() }); } /** @@ -153,13 +179,23 @@ export class SlaService { * worker process needed to invoke it, tests call this directly). Marks resolution breaches * (status -> breached) and first-response breaches (firstResponseBreachedAt, status * unchanged), then fires escalation for each newly-detected breach. + * + * 016-load-concurrency-testing research.md §2: each run's update is now version-guarded. A + * lost race here (a concurrent pause/resume/complete changed the run first) means this run's + * status is no longer what the sweep's own query assumed — skipped for this pass rather than + * retried, since the next scheduled sweep re-evaluates every run fresh against real + * conditions anyway; this only ever defers, never drops, a genuine breach. */ async runBreachDetectionSweep(): Promise { const now = new Date(); const resolutionBreaches = await this.runs.findRunningPastResolutionDueAt(now); for (const run of resolutionBreaches) { - await this.runs.update(run.id, { status: 'breached', breachedAt: now }); + const updated = await this.runs.updateWithVersion(run.id, run.version, { + status: 'breached', + breachedAt: now, + }); + if (!updated) continue; slaRunOutcomesCounter.inc({ outcome: 'breached' }); await this.escalation.handleBreach(run.ticketId, 'resolution_breach'); } @@ -170,7 +206,10 @@ export class SlaService { const hasAgentResponse = messages.some((m) => m.type === 'AGENT_MESSAGE'); if (hasAgentResponse) continue; - await this.runs.update(run.id, { firstResponseBreachedAt: now }); + const updated = await this.runs.updateWithVersion(run.id, run.version, { + firstResponseBreachedAt: now, + }); + if (!updated) continue; await this.escalation.handleBreach(run.ticketId, 'first_response_breach'); } } diff --git a/tests/concurrency/sla-race.test.ts b/tests/concurrency/sla-race.test.ts new file mode 100644 index 0000000..695528f --- /dev/null +++ b/tests/concurrency/sla-race.test.ts @@ -0,0 +1,218 @@ +import { describe, it, expect, beforeAll, afterAll } from 'vitest'; +import { buildApp } from '@/app'; +import { prismaClient } from '@/infrastructure/database'; +import { FastifyInstance } from 'fastify'; +import { loginAs, authHeader } from '../helpers/auth'; +import { + encryptCredential, + generateCredentialSecret, + issueIntegrationToken, +} from '@/modules/catalog/products'; +import { slaService } from '@/modules/orchestration/sla'; + +/** + * specs/016-load-concurrency-testing User Story 2 / FR-002 / SC-002: concurrent pause, resume, + * and breach-sweep activity against the same SLA run must always leave it in one + * internally-consistent, auditable state — never a state with contradictory fields (paused with + * no pause timestamp, or a legitimately breached run silently reverted to running by a racing + * resume). Before research.md §2's fix (SLARun.version + SlaRunRepository.updateWithVersion), + * SlaService.pause/resume/complete/runBreachDetectionSweep each did a plain read-then-write with + * no guard, so two racing calls could clobber each other's writes. + */ +describe('SLA pause/resume/sweep race (User Story 2)', () => { + let app: FastifyInstance; + let authToken: string; + const externalProductId = `TEST_SLA_RACE_PROD_${Date.now()}`; + const skillTag = `sla_race_skill_${Date.now()}`; + let productId: string; + let teamId: string; + let agentId: string; + let nodeId: string; + let policyId: string; + let secret: string; + const createdTicketIds: string[] = []; + + async function createTicketAndAssign(): Promise { + const token = issueIntegrationToken(secret, { + externalProductId, + tenantId: 'tenant-1', + userId: 'user-1', + }); + const created = await app.inject({ + method: 'POST', + url: '/v1/support/requests', + headers: { authorization: `Bearer ${token}` }, + payload: { + productId: externalProductId, + tenantId: 'tenant-1', + userId: 'user-1', + source: 'test', + problem: `SLA race test ${Date.now()}-${Math.random()}`, + }, + }); + expect(created.statusCode).toBe(202); + const ticketId = created.json().data.ticketId as string; + createdTicketIds.push(ticketId); + + const ticket = await prismaClient.ticket.findUniqueOrThrow({ where: { id: ticketId } }); + await app.inject({ + method: 'PATCH', + url: `/tickets/${ticketId}/status`, + headers: authHeader(authToken), + payload: { status: 'HUMAN_ESCALATION', expectedVersion: ticket.version }, + }); + return ticketId; + } + + beforeAll(async () => { + app = await buildApp(); + authToken = await loginAs(app, 'ADMIN'); + + const product = await prismaClient.product.create({ + data: { externalProductId, name: 'SLA Race Test Product', status: 'active' }, + }); + productId = product.id; + secret = generateCredentialSecret(); + await prismaClient.productIntegration.create({ + data: { + productId, + credentialRef: encryptCredential(secret), + authMechanism: 'signed_token', + allowedScope: { tenantIds: ['tenant-1'] }, + status: 'active', + rateLimitPerMinute: 1000, + rateLimitPerUserPerMinute: 1000, + }, + }); + + const team = await app.inject({ + method: 'POST', + url: '/admin/teams', + headers: authHeader(authToken), + payload: { name: `SLA Race Team ${Date.now()}` }, + }); + teamId = team.json().data.id; + + const agent = await app.inject({ + method: 'POST', + url: `/admin/teams/${teamId}/agents`, + headers: authHeader(authToken), + payload: { name: 'SLA Race Agent' }, + }); + agentId = agent.json().data.id; + await app.inject({ + method: 'PUT', + url: `/admin/agents/${agentId}/skills/${skillTag}`, + headers: authHeader(authToken), + payload: { level: 3 }, + }); + + const node = await app.inject({ + method: 'POST', + url: '/admin/hierarchy-nodes', + headers: authHeader(authToken), + payload: { + name: 'SLA Race Node', + order: 0, + productScope: [externalProductId], + skills: [skillTag], + assignmentStrategy: 'ROUND_ROBIN', + }, + }); + nodeId = node.json().data.id; + + const policy = await app.inject({ + method: 'POST', + url: '/admin/sla-policies', + headers: authHeader(authToken), + payload: { name: 'SLA Race Policy', productId, firstResponseMinutes: 30, resolutionMinutes: 240 }, + }); + policyId = policy.json().data.id; + }); + + afterAll(async () => { + const ticketFilter = { ticketId: { in: createdTicketIds } }; + await prismaClient.escalationEvent.deleteMany({ where: ticketFilter }); + await prismaClient.sLARun.deleteMany({ where: ticketFilter }); + await prismaClient.sLAPolicy.deleteMany({ where: { id: policyId } }); + await prismaClient.assignmentHistory.deleteMany({ where: ticketFilter }); + await prismaClient.assignment.deleteMany({ where: ticketFilter }); + await prismaClient.hierarchyNode.deleteMany({ where: { id: nodeId } }); + await prismaClient.agentSkill.deleteMany({ where: { agentId } }); + await prismaClient.agent.deleteMany({ where: { teamId } }); + await prismaClient.team.deleteMany({ where: { id: teamId } }); + await prismaClient.ticketMessage.deleteMany({ where: ticketFilter }); + await prismaClient.ticket.deleteMany({ where: { id: { in: createdTicketIds } } }); + await prismaClient.problem.deleteMany({ where: { productId } }); + await prismaClient.productIntegration.deleteMany({ where: { productId } }); + await prismaClient.product.deleteMany({ where: { id: productId } }); + await app.close(); + }); + + function assertInternallyConsistent(run: { + status: string; + pausedAt: Date | null; + resumedAt: Date | null; + breachedAt: Date | null; + }) { + if (run.status === 'paused') { + expect(run.pausedAt).not.toBeNull(); + } else { + expect(run.pausedAt).toBeNull(); + } + // A run the sweep has genuinely marked breached must never be silently reverted to running + // by a racing resume — status and breachedAt must agree with each other. + if (run.status === 'breached') { + expect(run.breachedAt).not.toBeNull(); + } + } + + it( + 'leaves an internally-consistent final state under concurrent pause/resume/sweep', + async () => { + const ticketId = await createTicketAndAssign(); + + // Force the run's resolution due date into the past so the breach sweep genuinely has + // something real to detect concurrently with pause/resume, not a no-op query. + await prismaClient.sLARun.update({ + where: { ticketId }, + data: { resolutionDueAt: new Date(Date.now() - 60_000) }, + }); + + await Promise.all([ + slaService.pause(ticketId), + slaService.resume(ticketId), + slaService.runBreachDetectionSweep(), + slaService.pause(ticketId), + slaService.resume(ticketId), + ]); + + const run = await prismaClient.sLARun.findUniqueOrThrow({ where: { ticketId } }); + assertInternallyConsistent(run); + }, + 30000, + ); + + it( + 'holds consistently across 10 repeated runs (SC-002: zero contradictory-state outcomes)', + async () => { + for (let run = 0; run < 10; run++) { + const ticketId = await createTicketAndAssign(); + await prismaClient.sLARun.update({ + where: { ticketId }, + data: { resolutionDueAt: new Date(Date.now() - 60_000) }, + }); + + await Promise.all([ + slaService.pause(ticketId), + slaService.resume(ticketId), + slaService.runBreachDetectionSweep(), + ]); + + const finalRun = await prismaClient.sLARun.findUniqueOrThrow({ where: { ticketId } }); + assertInternallyConsistent(finalRun); + } + }, + 60000, + ); +});