Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/decisions-log.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
38 changes: 30 additions & 8 deletions src/engine/incoming.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

/**
Expand All @@ -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<StageOutcome>;
};
Expand Down Expand Up @@ -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<StageOutcome> {
const { deliveryQueueId, deliveryStatus, device, commandType, response, unsolicited, failureContext } = parsedEvent;
Expand All @@ -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';
}

Expand All @@ -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') {
Expand Down Expand Up @@ -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 };
Expand Down
4 changes: 2 additions & 2 deletions src/engine/lifecycle/actions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -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);
},

/**
Expand Down
2 changes: 1 addition & 1 deletion src/engine/lifecycle/moves.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
);
Expand Down
2 changes: 1 addition & 1 deletion src/lib/redis-repository/admission-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion src/plugins/calin-api-v1/incoming.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
96 changes: 63 additions & 33 deletions src/plugins/calin-api-v1/lib/repo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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`).
*/
Expand All @@ -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`.
*
Expand Down Expand Up @@ -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),
});
}
};

Expand Down
1 change: 0 additions & 1 deletion test/helpers/in-memory-incoming.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@ export function createInMemoryIncomingService(
},
async processEvent(
_parsedEvent: ParsedIncomingEvent,
_queueKey: string,
_plugin: DeliveryPlugin,
): Promise<StageOutcome> {
// HTTP unit tests do not exercise poll / processEvent.
Expand Down
Loading
Loading