diff --git a/docs/decisions-log.md b/docs/decisions-log.md index c62d615..6aecc55 100644 --- a/docs/decisions-log.md +++ b/docs/decisions-log.md @@ -44,7 +44,7 @@ item lands, strike it here. | **Webhook drain concurrency** — drain is serialized in-process | Serialized drain lags under event volume | ADR-003 | | **ChirpStack ingress enqueue-then-ack** — vendor HTTP posts once and does not retry; v1 awaits `handle` before 204 for local Redis durability only | Designing durable raw-event enqueue → 204 → async process | — | | **Thorough cleanup suite** — every exit path (success, final failure, PULL age cap, cancel) must leave zero Redis references, PUSH and PULL. Sweep list is derived from the stage table (ADR-008 §7); remaining work is the dedicated integration suite | Dedicated integration suite | ADR-006 **D2**; ADR-008 §7 | -| **PUSH ingress `retryOrFail` uses the device queue unconditionally** — `incoming.handle` passes `STAGES.device.key()` into `processEvent`. A ChirpStack nack (`event=ack`, `acknowledged: false`) while the message is still in `queue_in_flight_to_relay_node` calls `enterRetry` on `queue_in_flight_to_device`, ZSCORE misses, WARN `enter retry lost the claim; another writer already moved this message` (wording is wrong: nothing moved it). Retry does not run; the message stays on the GW wait. Later txack/up can still complete (seen 2026-09-02 LoRaWAN cutover, e.g. `mi:6296686`). A nack with no later success waits until relay-node timeout (`getRemoteStatus`, 15 min) | Nacks at GW leave `PROCESSING` past a few seconds, or we change ingress to pass the real stage key | ADR-008; `src/engine/incoming.ts` `handle` → `processEvent(..., STAGES.device.key())` | +| ~~**PUSH ingress `retryOrFail` uses the device queue unconditionally**~~ — landed 2026-09-04. `processEvent` claims the stage from the hash (`stageForStatus` / `stageKeyFor`) and no longer takes a queue key. Smoke: nack while still on the relay-node wait in `test/integration/incoming-ingress.smoke.spec.ts`. ADR-008. | — | — | ### Product / ops trigger diff --git a/src/engine/incoming.ts b/src/engine/incoming.ts index 953eb17..e18efd6 100644 --- a/src/engine/incoming.ts +++ b/src/engine/incoming.ts @@ -22,7 +22,7 @@ import type { } from '../plugins/plugin.interface.js'; import type { BaseService } from './base.js'; import type { StageMoves } from './lifecycle/moves.js'; -import { STAGES } from './lifecycle/stages.js'; +import { STAGES, stageForStatus, stageKeyFor } from './lifecycle/stages.js'; import type { StageOutcome } from './lifecycle/types.js'; /** @@ -48,12 +48,10 @@ export type IncomingService = { * member either a new score or a removal. * * @param parsedEvent - Normalized event from the plugin - * @param currentQueueKey - Queue the message sits in (the caller's stage) * @param plugin - Owning delivery plugin */ processEvent( parsedEvent: ParsedIncomingEvent, - currentQueueKey: string, plugin: DeliveryPlugin, ): Promise; }; @@ -87,12 +85,10 @@ export function createIncomingService(options: CreateIncomingServiceOptions): In * stage on every tick (A2). * * @param parsedEvent - Normalized event from the plugin - * @param currentQueueKey - Queue for retry/fail (PUSH handle uses the device queue) * @param plugin - Owning delivery plugin (tuning for stage moves) */ async function processEvent( parsedEvent: ParsedIncomingEvent, - currentQueueKey: string, plugin: DeliveryPlugin, ): Promise { const { deliveryQueueId, deliveryStatus, device, commandType, response, unsolicited, failureContext } = parsedEvent; @@ -117,7 +113,7 @@ export function createIncomingService(options: CreateIncomingServiceOptions): In const messageId = await messageStore.getMessageIdFromDeliveryQueueId(deliveryQueueId); if (!messageId) { - logger.warn({ module: 'incoming', deliveryQueueId, parsedEvent }, 'message not found for deliveryQueueId'); + logger.debug({ module: 'incoming', deliveryQueueId, parsedEvent }, 'message not found for deliveryQueueId'); return 'orphaned'; } @@ -130,7 +126,33 @@ export function createIncomingService(options: CreateIncomingServiceOptions): In if (deliveryStatus === 'DELIVERY_FAILED') { const context = failureContext ?? { reason: 'Unable to deliver message after negative remote response' }; - return baseService.retryOrFail(messageId, currentQueueKey, context, plugin); + const storedMessage = await messageStore.getMessageById(messageId); + if (!storedMessage) { + logger.warn({ module: 'incoming', messageId }, 'message not found (already cleaned up?)'); + return 'orphaned'; + } + + // Claim the stage the hash is in. A nack can arrive while the member is + // still on the relay-node wait. + const stage = stageForStatus(storedMessage.deliveryStatus, plugin.deliveryPattern); + if (!stage) { + logger.warn( + { + module: 'incoming', + messageId, + deliveryStatus: storedMessage.deliveryStatus, + }, + 'cannot retry; message is not in a stage', + ); + return 'orphaned'; + } + + return baseService.retryOrFail( + messageId, + stageKeyFor(STAGES[stage], plugin.id), + context, + plugin, + ); } if (deliveryStatus !== 'DELIVERY_SUCCESSFUL') { @@ -191,7 +213,7 @@ export function createIncomingService(options: CreateIncomingServiceOptions): In return; } - await processEvent(parsedEvent, STAGES.device.key(), plugin); + await processEvent(parsedEvent, plugin); } return { handle, processEvent }; diff --git a/src/engine/lifecycle/actions.ts b/src/engine/lifecycle/actions.ts index 4cbdc92..90421df 100644 --- a/src/engine/lifecycle/actions.ts +++ b/src/engine/lifecycle/actions.ts @@ -141,7 +141,7 @@ export function createStageActions(options: CreateStageActionsOptions): StageAct * vendor how the task is doing, and either resolve the message or wait again on the * poll ladder. */ - async awaitingTask({ message, plugin, queueKey, messageAgeMs }) { + async awaitingTask({ message, plugin, messageAgeMs }) { if (messageAgeMs >= PULL_MAX_MESSAGE_AGE_MS) { await _failPermanently(message, plugin); return 'removed'; @@ -165,7 +165,7 @@ export function createStageActions(options: CreateStageActionsOptions): StageAct const parsedEvent = await fetchStatus(message); if (!parsedEvent) return 'rescheduled'; - return incomingService.processEvent(parsedEvent, queueKey, plugin); + return incomingService.processEvent(parsedEvent, plugin); }, /** diff --git a/src/engine/lifecycle/moves.ts b/src/engine/lifecycle/moves.ts index 6347652..2f973c2 100644 --- a/src/engine/lifecycle/moves.ts +++ b/src/engine/lifecycle/moves.ts @@ -245,7 +245,7 @@ export function createStageMoves(options: CreateStageMovesOptions): StageMoves { // legitimate (cancel, or a deadline that fired first), it is the rate that matters. if (!claimed) { metrics.recordStageClaimMiss(from); - logger.warn( + logger.debug( { module: 'lifecycle', messageId, from, to: next, pluginId: plugin.id }, 'stage advance lost the claim; another writer already moved this message', ); diff --git a/src/lib/redis-repository/admission-store.ts b/src/lib/redis-repository/admission-store.ts index 72a5523..8dd3144 100644 --- a/src/lib/redis-repository/admission-store.ts +++ b/src/lib/redis-repository/admission-store.ts @@ -80,7 +80,7 @@ export function createAdmissionStore( if (deadMembers.length > 0) { await client.srem(concurrencyRateLimitKey, ...deadMembers); - logger.warn({ + logger.debug({ module: 'redis', deadCount: deadMembers.length, concurrencyRateLimitKey, diff --git a/src/plugins/calin-api-v1/incoming.ts b/src/plugins/calin-api-v1/incoming.ts index 7049be9..3ef649b 100644 --- a/src/plugins/calin-api-v1/incoming.ts +++ b/src/plugins/calin-api-v1/incoming.ts @@ -220,7 +220,7 @@ export function createCalinApiV1Incoming( // A numeric code means CALIN answered over HTTP — that is a real failure. // Returning null keeps the TaskNo and the poll ladder (ADR-008 awaitingTask). if (typeof errCode !== 'number') { - logger.warn({ + logger.debug({ module: 'calin-api-v1.incoming', err, messageId: id, diff --git a/src/plugins/calin-api-v1/lib/repo.ts b/src/plugins/calin-api-v1/lib/repo.ts index 66a93fa..ca91f60 100644 --- a/src/plugins/calin-api-v1/lib/repo.ts +++ b/src/plugins/calin-api-v1/lib/repo.ts @@ -109,6 +109,21 @@ type DownResponseBody = { readonly Message?: unknown; }; +/** + * Node / undici errnos that mean we never got an HTTP response from CALIN. + * `UND_ERR_*` is matched by prefix as well (connect / socket / header timeouts). + */ +const TRANSPORT_CODES = new Set([ + 'ECONNREFUSED', + 'ECONNRESET', + 'ETIMEDOUT', + 'EAI_AGAIN', + 'ENOTFOUND', + 'EHOSTUNREACH', + 'ENETUNREACH', + 'EPIPE', +]); + /** * Read a Node-style errno from a fetch failure (`cause.code` or top-level `code`). */ @@ -129,6 +144,41 @@ function getNetworkErrorCode(err: unknown): string | undefined { return undefined; } +function isAbortTimeout(err: unknown): boolean { + let current: unknown = err; + for (let i = 0; i < 4; i++) { + if (typeof current !== 'object' || current === null) return false; + const name = (current as { name?: unknown }).name; + if (name === 'TimeoutError' || name === 'AbortError') return true; + current = (current as { cause?: unknown }).cause; + } + return false; +} + +function isTransportCode(code: string): boolean { + return TRANSPORT_CODES.has(code) || code.startsWith('UND_ERR_'); +} + +/** + * String errno for a transport failure, or `undefined` when this is not transport. + */ +function transportCodeOf(err: unknown): string | undefined { + const code = getNetworkErrorCode(err); + if (code !== undefined && isTransportCode(code)) return code; + if (isAbortTimeout(err)) return code ?? 'TimeoutError'; + return undefined; +} + +function messageForTransport(code: string): string { + if (code === 'ECONNREFUSED') { + return '[CALIN API-V1] could not be reached, connection was refused'; + } + if (code === 'ECONNRESET') { + return '[CALIN API-V1] abruptly closed its end of the connection'; + } + return '[CALIN API-V1] could not be reached'; +} + /** * Build a CALIN API V1 client closed over `apiBaseUrl`. * @@ -204,41 +254,21 @@ export function createCalinApiV1Client(deps: { readonly apiBaseUrl: string }) { throw err; } - let message: string; - let code: number | string | null | undefined; - const networkCode = getNetworkErrorCode(err); - - if (networkCode === 'ECONNREFUSED') { - logger.error({ module: 'calin-api-v1.repo', path, err }, 'ECONNREFUSED'); - message = '[CALIN API-V1] could not be reached, connection was refused'; - code = networkCode; - } - else if (networkCode === 'ECONNRESET') { - logger.error({ module: 'calin-api-v1.repo', path, err }, 'ECONNRESET'); - message = '[CALIN API-V1] abruptly closed its end of the connection'; - code = networkCode; - } - else if ( - typeof err === 'object' - && err !== null - && 'cause' in err - && (err as { cause?: unknown }).cause - ) { - logger.error({ module: 'calin-api-v1.repo', path, err }, 'unhandled cause'); - message = '[CALIN API-V1] is down'; - code = getNetworkErrorCode((err as { cause: unknown }).cause) - ?? getNetworkErrorCode(err); - } - else if (err instanceof Error && err.message) { - logger.error({ module: 'calin-api-v1.repo', path, err }, 'fetch failed'); - message = '[CALIN API-V1] is down'; - } - else { - logger.error({ module: 'calin-api-v1.repo', path, err }, 'fetch failed'); - message = '[CALIN API-V1] is down'; + const transportCode = transportCodeOf(err); + if (transportCode !== undefined) { + logger.warn( + { module: 'calin-api-v1.repo', path, code: transportCode }, + 'unreachable', + ); + throw new CalinApiV1Error(messageForTransport(transportCode), { + code: transportCode, + }); } - throw new CalinApiV1Error(message, { code }); + logger.error({ module: 'calin-api-v1.repo', path, err }, 'fetch failed'); + throw new CalinApiV1Error('[CALIN API-V1] is down', { + code: getNetworkErrorCode(err), + }); } }; diff --git a/test/helpers/in-memory-incoming.ts b/test/helpers/in-memory-incoming.ts index 683e378..62273a3 100644 --- a/test/helpers/in-memory-incoming.ts +++ b/test/helpers/in-memory-incoming.ts @@ -36,7 +36,6 @@ export function createInMemoryIncomingService( }, async processEvent( _parsedEvent: ParsedIncomingEvent, - _queueKey: string, _plugin: DeliveryPlugin, ): Promise { // HTTP unit tests do not exercise poll / processEvent. diff --git a/test/integration/incoming-ingress.smoke.spec.ts b/test/integration/incoming-ingress.smoke.spec.ts index afa4f48..f5fa931 100644 --- a/test/integration/incoming-ingress.smoke.spec.ts +++ b/test/integration/incoming-ingress.smoke.spec.ts @@ -7,7 +7,7 @@ * docker compose up -d valkey * pnpm exec vitest run test/integration/incoming-ingress.smoke.spec.ts */ -import { afterAll, describe, expect, it } from 'vitest'; +import { afterAll, beforeAll, describe, expect, it } from 'vitest'; import { buildApp } from '#src/app.js'; import { deviceMessagingConfigSchema } from '#src/config/schema.js'; @@ -16,7 +16,7 @@ import { createInFlightSends } from '#src/engine/in-flight-sends.js'; import { createIncomingService } from '#src/engine/incoming.js'; import { createStageMoves } from '#src/engine/lifecycle/moves.js'; import { STAGES } from '#src/engine/lifecycle/stages.js'; -import { createOutgoingService } from '#src/engine/outgoing.js'; +import { createOutgoingService, type OutgoingService } from '#src/engine/outgoing.js'; import { createAdmissionStore } from '#src/lib/redis-repository/admission-store.js'; import { createMessageStore } from '#src/lib/redis-repository/message-store.js'; import { createStageStore } from '#src/lib/redis-repository/stage-store.js'; @@ -28,20 +28,27 @@ import { waitForPostSend } from '../helpers/wait-for-post-send.js'; const delivery = deviceMessagingConfigSchema.parse({ $schemaVersion: '1' }).delivery; +type IngressStack = { + readonly app: Awaited>; + readonly outgoing: OutgoingService; +}; + describe('incoming PUSH ingress', () => { let redis: typeof import('../../src/lib/redis-repository/client.js').redis; let redisKeys: typeof import('../../src/lib/redis-repository/keys.js').redisKeys; + beforeAll(async () => { + ({ redis } = await import('../../src/lib/redis-repository/client.js')); + ({ redisKeys } = await import('../../src/lib/redis-repository/keys.js')); + }); + afterAll(async () => { if (redis) { await redis.quit(); } }); - it('stub-push: GW → ingress success → message cleaned up', async () => { - ({ redis } = await import('../../src/lib/redis-repository/client.js')); - ({ redisKeys } = await import('../../src/lib/redis-repository/keys.js')); - + async function startStack(): Promise { const registry = createPluginRegistry([ { id: STUB_PUSH_ID } ]); const metrics = noopMetrics; const messageStore = createMessageStore({ client: redis }); @@ -49,7 +56,7 @@ describe('incoming PUSH ingress', () => { const moves = createStageMoves({ delivery, metrics, stageStore }); const baseService = createBaseService({ delivery, metrics, messageStore, moves }); const inFlightSends = createInFlightSends(); - const outgoingService = createOutgoingService({ + const outgoing = createOutgoingService({ registry, delivery, baseService, @@ -67,13 +74,17 @@ describe('incoming PUSH ingress', () => { metrics, }); const app = await buildApp({ metrics, incomingService, registry }); + return { app, outgoing }; + } + it('stub-push: GW → ingress success → message cleaned up', async () => { + const { app, outgoing } = await startStack(); const correlationId = `ingress-push-${ Date.now() }`; const networkId = 92; const queueKey = `queue:stub-push:network:${ networkId }`; try { - const enqueued = await outgoingService.enqueue({ + const enqueued = await outgoing.enqueue({ commandType: 'READ_CREDIT', priority: 1, pluginId: STUB_PUSH_ID, @@ -85,8 +96,8 @@ describe('incoming PUSH ingress', () => { }, }); - await outgoingService.distributeToNetworkServers(); - const afterSend = await waitForPostSend(outgoingService, correlationId); + await outgoing.distributeToNetworkServers(); + const afterSend = await waitForPostSend(outgoing, correlationId); expect(afterSend.deliveryQueueId).toMatch(/^stub-ext-/); const deliveryQueueId = afterSend.deliveryQueueId; @@ -102,7 +113,7 @@ describe('incoming PUSH ingress', () => { }); expect(ack.statusCode).toBe(204); - const afterAck = await outgoingService.getByCorrelationId(correlationId); + const afterAck = await outgoing.getByCorrelationId(correlationId); expect(afterAck?.deliveryStatus).toBe('SENT_TO_DEVICE'); expect(await redis.zscore(STAGES.device.key(), enqueued.id)).not.toBeNull(); @@ -119,12 +130,74 @@ describe('incoming PUSH ingress', () => { }); expect(success.statusCode).toBe(204); - const afterSuccess = await outgoingService.getByCorrelationId(correlationId); + const afterSuccess = await outgoing.getByCorrelationId(correlationId); expect(afterSuccess).toBeNull(); expect(await redis.zscore(STAGES.device.key(), enqueued.id)).toBeNull(); } finally { - const leftover = await outgoingService.getByCorrelationId(correlationId); + const leftover = await outgoing.getByCorrelationId(correlationId); + if (leftover) { + await purgeMessageReferences(leftover.id, { correlationId }); + } + await redis.srem( + redisKeys.listOfInitialQueuesToDistributeFrom(), + queueKey, + ); + await redis.del(redisKeys.lockForQueue(queueKey)); + await app.close(); + } + }); + + it('stub-push: nack while still on the relay-node wait enters retry', async () => { + const { app, outgoing } = await startStack(); + const correlationId = `ingress-push-nack-${ Date.now() }`; + const networkId = 93; + const queueKey = `queue:stub-push:network:${ networkId }`; + + try { + const enqueued = await outgoing.enqueue({ + commandType: 'READ_CREDIT', + priority: 1, + pluginId: STUB_PUSH_ID, + networkId, + correlationId, + device: { + type: 'ELECTRICITY_METER', + externalReference: 'ingress-push-nack-meter', + }, + }); + + await outgoing.distributeToNetworkServers(); + const afterSend = await waitForPostSend(outgoing, correlationId); + expect(afterSend.deliveryQueueId).toMatch(/^stub-ext-/); + expect(afterSend.deliveryStatus).toBe('DELIVERED_TO_NS'); + expect(await redis.zscore(STAGES.relayNode.key(), enqueued.id)).not.toBeNull(); + const deliveryQueueId = afterSend.deliveryQueueId; + + const nack = await app.inject({ + method: 'POST', + url: `/ingress/${ STUB_PUSH_ID }`, + headers: { 'content-type': 'application/json' }, + payload: { + deliveryQueueId, + deliveryStatus: 'DELIVERY_FAILED', + device: enqueued.device, + failureContext: { reason: 'Downlink not acknowledged by device' }, + }, + }); + expect(nack.statusCode).toBe(204); + + const afterNack = await outgoing.getByCorrelationId(correlationId); + expect(afterNack?.deliveryStatus).toBe('TO_RETRY'); + expect(await redis.zscore(STAGES.relayNode.key(), enqueued.id)).toBeNull(); + expect(await redis.zscore(STAGES.device.key(), enqueued.id)).toBeNull(); + expect(await redis.zscore(STAGES.retry.key(), enqueued.id)).not.toBeNull(); + expect( + await redis.exists(redisKeys.indexExternalDeliveryId(deliveryQueueId)), + ).toBe(0); + } + finally { + const leftover = await outgoing.getByCorrelationId(correlationId); if (leftover) { await purgeMessageReferences(leftover.id, { correlationId }); } diff --git a/test/unit/plugins/calin-api-v1-repo.spec.ts b/test/unit/plugins/calin-api-v1-repo.spec.ts index 9a6f6e8..0807046 100644 --- a/test/unit/plugins/calin-api-v1-repo.spec.ts +++ b/test/unit/plugins/calin-api-v1-repo.spec.ts @@ -135,6 +135,40 @@ describe('createCalinApiV1Client', () => { }); }); + it('maps UND_ERR_CONNECT_TIMEOUT (via cause.code) to CalinApiV1Error', async () => { + const timedOut = Object.assign(new Error('fetch failed'), { + cause: Object.assign(new Error('Connect Timeout Error'), { + code: 'UND_ERR_CONNECT_TIMEOUT', + }), + }); + vi.stubGlobal('fetch', vi.fn().mockRejectedValue(timedOut)); + + const client = createCalinApiV1Client({ apiBaseUrl: API_BASE }); + await expect( + client.sendRequest('/COMM_RemoteTokenTask', { TaskNo: 't-1' }), + ).rejects.toMatchObject({ + name: 'CalinApiV1Error', + message: '[CALIN API-V1] could not be reached', + code: 'UND_ERR_CONNECT_TIMEOUT', + }); + }); + + it('maps AbortSignal timeout (TimeoutError) to CalinApiV1Error', async () => { + const aborted = Object.assign(new Error('The operation was aborted due to timeout'), { + name: 'TimeoutError', + }); + vi.stubGlobal('fetch', vi.fn().mockRejectedValue(aborted)); + + const client = createCalinApiV1Client({ apiBaseUrl: API_BASE }); + await expect( + client.sendRequest('/COMM_RemoteReading', { MeterNo: 'm-1' }), + ).rejects.toMatchObject({ + name: 'CalinApiV1Error', + message: '[CALIN API-V1] could not be reached', + code: 'TimeoutError', + }); + }); + it('maps generic fetch failures to CalinApiV1Error', async () => { vi.stubGlobal( 'fetch',