Initial LisgloSIPS V2 implementation
This commit is contained in:
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"name": "@lisglosips/redis",
|
||||
"version": "0.2.0",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
"types": "dist/index.d.ts",
|
||||
"scripts": {
|
||||
"build": "tsc -p tsconfig.json"
|
||||
},
|
||||
"dependencies": {
|
||||
"ioredis": "5.8.2"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,273 @@
|
||||
import { describe, expect, it } from 'vitest';
|
||||
|
||||
import {
|
||||
CDR_CONSUMER_GROUP,
|
||||
CDR_DEADLETTER_STREAM,
|
||||
CDR_STREAM,
|
||||
ensureCdrConsumerGroup,
|
||||
parseCdrStreamEvent,
|
||||
processCdrBatch,
|
||||
processPendingCdrBatch,
|
||||
publishCdrEvent,
|
||||
CdrRetryableError,
|
||||
type CdrStreamPublishInput,
|
||||
type CdrRedisCommands,
|
||||
type RedisAutoClaimResponse,
|
||||
type RedisStreamEntry,
|
||||
type RedisStreamReadResponse
|
||||
} from './index.js';
|
||||
|
||||
class FakeRedis implements CdrRedisCommands {
|
||||
readonly streams = new Map<string, RedisStreamEntry[]>();
|
||||
readonly groups = new Map<string, { delivered: number; pending: Map<string, RedisStreamEntry> }>();
|
||||
readonly locks = new Set<string>();
|
||||
|
||||
async xadd(stream: string, id: string, ...fieldValues: string[]): Promise<string> {
|
||||
const entries = this.streams.get(stream) ?? [];
|
||||
const redisId = id === '*' ? `${entries.length + 1}-0` : id;
|
||||
entries.push([redisId, fieldValues]);
|
||||
this.streams.set(stream, entries);
|
||||
return redisId;
|
||||
}
|
||||
|
||||
async xack(stream: string, group: string, ...ids: string[]): Promise<number> {
|
||||
const state = this.group(stream, group);
|
||||
let count = 0;
|
||||
for (const id of ids) {
|
||||
if (state.pending.delete(id)) {
|
||||
count += 1;
|
||||
}
|
||||
}
|
||||
return count;
|
||||
}
|
||||
|
||||
async xgroup(...args: string[]): Promise<string> {
|
||||
const [command, stream, group] = args;
|
||||
if (command !== 'CREATE' || !stream || !group) {
|
||||
throw new Error(`Unsupported xgroup call: ${args.join(' ')}`);
|
||||
}
|
||||
const key = this.groupKey(stream, group);
|
||||
if (this.groups.has(key)) {
|
||||
throw new Error('BUSYGROUP Consumer Group name already exists');
|
||||
}
|
||||
this.groups.set(key, { delivered: 0, pending: new Map() });
|
||||
if (!this.streams.has(stream)) {
|
||||
this.streams.set(stream, []);
|
||||
}
|
||||
return 'OK';
|
||||
}
|
||||
|
||||
async xreadgroup(...args: Array<string | number>): Promise<RedisStreamReadResponse | null> {
|
||||
const group = String(args[1]);
|
||||
const consumer = String(args[2]);
|
||||
const countIndex = args.indexOf('COUNT');
|
||||
const streamIndex = args.indexOf('STREAMS');
|
||||
const count = countIndex === -1 ? 10 : Number(args[countIndex + 1]);
|
||||
const stream = String(args[streamIndex + 1]);
|
||||
const state = this.group(stream, group);
|
||||
const entries = this.streams.get(stream) ?? [];
|
||||
const selected = entries.slice(state.delivered, state.delivered + count);
|
||||
if (selected.length === 0) {
|
||||
return null;
|
||||
}
|
||||
state.delivered += selected.length;
|
||||
for (const entry of selected) {
|
||||
state.pending.set(entry[0], entry);
|
||||
}
|
||||
void consumer;
|
||||
return [[stream, selected]];
|
||||
}
|
||||
|
||||
async xautoclaim(...args: Array<string | number>): Promise<RedisAutoClaimResponse> {
|
||||
const stream = String(args[0]);
|
||||
const group = String(args[1]);
|
||||
const consumer = String(args[2]);
|
||||
const countIndex = args.indexOf('COUNT');
|
||||
const count = countIndex === -1 ? 10 : Number(args[countIndex + 1]);
|
||||
const state = this.group(stream, group);
|
||||
const entries = Array.from(state.pending.values()).slice(0, count);
|
||||
for (const entry of entries) {
|
||||
state.pending.set(entry[0], entry);
|
||||
}
|
||||
void consumer;
|
||||
return ['0-0', entries];
|
||||
}
|
||||
|
||||
async set(key: string, _value: string, mode: 'NX', _expireMode: 'EX', _ttlSeconds: number): Promise<'OK' | null> {
|
||||
if (mode !== 'NX') {
|
||||
throw new Error('FakeRedis only supports SET NX');
|
||||
}
|
||||
if (this.locks.has(key)) {
|
||||
return null;
|
||||
}
|
||||
this.locks.add(key);
|
||||
return 'OK';
|
||||
}
|
||||
|
||||
async del(...keys: string[]): Promise<number> {
|
||||
let deleted = 0;
|
||||
for (const key of keys) {
|
||||
if (this.locks.delete(key)) {
|
||||
deleted += 1;
|
||||
}
|
||||
}
|
||||
return deleted;
|
||||
}
|
||||
|
||||
pendingSize(stream = CDR_STREAM, group = CDR_CONSUMER_GROUP): number {
|
||||
return this.group(stream, group).pending.size;
|
||||
}
|
||||
|
||||
private group(stream: string, group: string): { delivered: number; pending: Map<string, RedisStreamEntry> } {
|
||||
const state = this.groups.get(this.groupKey(stream, group));
|
||||
if (!state) {
|
||||
throw new Error('NOGROUP No such key or consumer group');
|
||||
}
|
||||
return state;
|
||||
}
|
||||
|
||||
private groupKey(stream: string, group: string): string {
|
||||
return `${stream}:${group}`;
|
||||
}
|
||||
}
|
||||
|
||||
function sampleEvent(overrides: Partial<CdrStreamPublishInput> = {}): CdrStreamPublishInput {
|
||||
return {
|
||||
eventId: 'evt-001',
|
||||
callId: 'call-001',
|
||||
nodeId: 'a1',
|
||||
opensipsInstance: 'opensips-a1',
|
||||
ingressAIp: '100.90.90.90',
|
||||
rtpengineNode: 'a1',
|
||||
sourceIp: '100.93.185.30',
|
||||
caller: 's21-ip-1001',
|
||||
callee: '13800138000',
|
||||
sipCode: 503,
|
||||
hangupReason: 'CONFIG_MISSING',
|
||||
configVersion: '0',
|
||||
createdAt: '2026-06-21T07:30:00.000Z',
|
||||
...overrides
|
||||
};
|
||||
}
|
||||
|
||||
describe('CDR Redis Stream contract', () => {
|
||||
it('publishes multi-A-aware CDR fields and parses typed values', async () => {
|
||||
const redis = new FakeRedis();
|
||||
const redisId = await publishCdrEvent(redis, sampleEvent());
|
||||
const entry = redis.streams.get(CDR_STREAM)?.find(([id]) => id === redisId);
|
||||
|
||||
expect(entry).toBeDefined();
|
||||
const parsed = parseCdrStreamEvent(entry?.[1] ?? []);
|
||||
expect(parsed.event_id).toBe('evt-001');
|
||||
expect(parsed.node_id).toBe('a1');
|
||||
expect(parsed.opensips_instance).toBe('opensips-a1');
|
||||
expect(parsed.rtpengine_node).toBe('a1');
|
||||
expect(parsed.sipCode).toBe(503);
|
||||
});
|
||||
|
||||
it('creates the consumer group idempotently', async () => {
|
||||
const redis = new FakeRedis();
|
||||
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
|
||||
expect(redis.groups.size).toBe(1);
|
||||
});
|
||||
|
||||
it('reads with XREADGROUP and ACKs processed entries', async () => {
|
||||
const redis = new FakeRedis();
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
await publishCdrEvent(redis, sampleEvent());
|
||||
|
||||
const summary = await processCdrBatch(redis, async () => 'processed', {
|
||||
consumer: 'test-consumer',
|
||||
blockMs: 1
|
||||
});
|
||||
|
||||
expect(summary.processed).toBe(1);
|
||||
expect(summary.ackedIds).toEqual(['1-0']);
|
||||
expect(redis.pendingSize()).toBe(0);
|
||||
});
|
||||
|
||||
it('ACKs duplicate event_id without invoking downstream processing again', async () => {
|
||||
const redis = new FakeRedis();
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
await publishCdrEvent(redis, sampleEvent({ callId: 'call-001' }));
|
||||
await publishCdrEvent(redis, sampleEvent({ callId: 'call-duplicate' }));
|
||||
let handled = 0;
|
||||
|
||||
const summary = await processCdrBatch(
|
||||
redis,
|
||||
async () => {
|
||||
handled += 1;
|
||||
return 'processed';
|
||||
},
|
||||
{ consumer: 'test-consumer', count: 2, blockMs: 1 }
|
||||
);
|
||||
|
||||
expect(handled).toBe(1);
|
||||
expect(summary.processed).toBe(1);
|
||||
expect(summary.duplicates).toBe(1);
|
||||
expect(summary.ackedIds).toEqual(['1-0', '2-0']);
|
||||
});
|
||||
|
||||
it('moves invalid payloads to the deadletter stream and ACKs them', async () => {
|
||||
const redis = new FakeRedis();
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
await redis.xadd(CDR_STREAM, '*', 'event_id', 'evt-bad');
|
||||
|
||||
const summary = await processCdrBatch(redis, async () => 'processed', {
|
||||
consumer: 'test-consumer',
|
||||
blockMs: 1
|
||||
});
|
||||
|
||||
expect(summary.deadlettered).toBe(1);
|
||||
expect(redis.pendingSize()).toBe(0);
|
||||
expect(redis.streams.get(CDR_DEADLETTER_STREAM)?.[0]?.[1]).toContain('evt-bad');
|
||||
});
|
||||
|
||||
it('leaves retryable handler failures pending and releases the idempotency lock', async () => {
|
||||
const redis = new FakeRedis();
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
await publishCdrEvent(redis, sampleEvent());
|
||||
|
||||
const summary = await processCdrBatch(
|
||||
redis,
|
||||
async () => {
|
||||
throw new CdrRetryableError('database is temporarily unavailable');
|
||||
},
|
||||
{ consumer: 'test-consumer', blockMs: 1 }
|
||||
);
|
||||
|
||||
expect(summary).toMatchObject({ processed: 0, duplicates: 0, deadlettered: 0, pendingLeft: 1 });
|
||||
expect(summary.ackedIds).toEqual([]);
|
||||
expect(redis.pendingSize()).toBe(1);
|
||||
expect(redis.streams.get(CDR_DEADLETTER_STREAM)).toBeUndefined();
|
||||
|
||||
const retry = await processPendingCdrBatch(redis, async () => 'processed', {
|
||||
consumer: 'retry-consumer',
|
||||
minIdleMs: 1
|
||||
});
|
||||
|
||||
expect(retry.processed).toBe(1);
|
||||
expect(retry.ackedIds).toEqual(['1-0']);
|
||||
expect(redis.pendingSize()).toBe(0);
|
||||
});
|
||||
|
||||
it('reclaims pending entries with XAUTOCLAIM and ACKs after retry', async () => {
|
||||
const redis = new FakeRedis();
|
||||
await ensureCdrConsumerGroup(redis);
|
||||
await publishCdrEvent(redis, sampleEvent());
|
||||
await redis.xreadgroup('GROUP', CDR_CONSUMER_GROUP, 'stalled-consumer', 'COUNT', 1, 'STREAMS', CDR_STREAM, '>');
|
||||
|
||||
expect(redis.pendingSize()).toBe(1);
|
||||
const summary = await processPendingCdrBatch(redis, async () => 'processed', {
|
||||
consumer: 'retry-consumer',
|
||||
minIdleMs: 1
|
||||
});
|
||||
|
||||
expect(summary.processed).toBe(1);
|
||||
expect(summary.ackedIds).toEqual(['1-0']);
|
||||
expect(redis.pendingSize()).toBe(0);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,408 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { CDR_CONSUMER_GROUP, CDR_DEADLETTER_STREAM, CDR_STREAM } from './index.js';
|
||||
|
||||
export const CDR_EVENT_SCHEMA_VERSION = '1';
|
||||
export const CDR_IDEMPOTENCY_KEY_PREFIX = 'lock:cdr:';
|
||||
|
||||
export type CdrStreamFieldMap = Record<string, string>;
|
||||
|
||||
export interface CdrRedisCommands {
|
||||
xadd(stream: string, id: string, ...fieldValues: string[]): Promise<string>;
|
||||
xack(stream: string, group: string, ...ids: string[]): Promise<number>;
|
||||
xgroup(...args: string[]): Promise<string>;
|
||||
xreadgroup(...args: Array<string | number>): Promise<RedisStreamReadResponse | null>;
|
||||
xautoclaim(...args: Array<string | number>): Promise<RedisAutoClaimResponse>;
|
||||
set(key: string, value: string, mode: 'NX', expireMode: 'EX', ttlSeconds: number): Promise<'OK' | null>;
|
||||
del(...keys: string[]): Promise<number>;
|
||||
}
|
||||
|
||||
export type RedisStreamEntry = [id: string, fields: string[]];
|
||||
export type RedisStreamReadResponse = Array<[stream: string, entries: RedisStreamEntry[]]>;
|
||||
export type RedisAutoClaimResponse = [nextStartId: string, entries: RedisStreamEntry[], deletedIds?: string[]];
|
||||
|
||||
export interface CdrStreamEvent {
|
||||
schema_version: string;
|
||||
event_id: string;
|
||||
idempotency_key: string;
|
||||
call_id: string;
|
||||
node_id: string;
|
||||
opensips_instance: string;
|
||||
ingress_a_ip: string;
|
||||
rtpengine_node: string;
|
||||
customer_id: string;
|
||||
customer_gateway_id: string;
|
||||
customer_gateway_policy_id: string;
|
||||
source_ip: string;
|
||||
caller: string;
|
||||
callee: string;
|
||||
vendor_id: string;
|
||||
vendor_gateway_id: string;
|
||||
line_group_id: string;
|
||||
started_at: string;
|
||||
answered_at: string;
|
||||
ended_at: string;
|
||||
duration_sec: string;
|
||||
sip_code: string;
|
||||
hangup_reason: string;
|
||||
recording_key: string;
|
||||
config_version: string;
|
||||
created_at: string;
|
||||
}
|
||||
|
||||
export interface ParsedCdrStreamEvent extends CdrStreamEvent {
|
||||
durationSeconds: number;
|
||||
sipCode: number;
|
||||
}
|
||||
|
||||
export interface CdrStreamPublishInput {
|
||||
eventId?: string;
|
||||
idempotencyKey?: string;
|
||||
callId: string;
|
||||
nodeId?: string;
|
||||
opensipsInstance?: string;
|
||||
ingressAIp?: string;
|
||||
rtpengineNode?: string;
|
||||
customerId?: string;
|
||||
customerGatewayId?: string;
|
||||
customerGatewayPolicyId?: string;
|
||||
sourceIp: string;
|
||||
caller?: string;
|
||||
callee?: string;
|
||||
vendorId?: string;
|
||||
vendorGatewayId?: string;
|
||||
lineGroupId?: string;
|
||||
startedAt?: string;
|
||||
answeredAt?: string;
|
||||
endedAt?: string;
|
||||
durationSec?: number;
|
||||
sipCode: number;
|
||||
hangupReason: string;
|
||||
recordingKey?: string;
|
||||
configVersion?: string;
|
||||
createdAt?: string;
|
||||
}
|
||||
|
||||
export type CdrHandlerResult = 'processed' | 'duplicate';
|
||||
export type CdrStreamHandler = (event: ParsedCdrStreamEvent, redisId: string) => Promise<CdrHandlerResult>;
|
||||
|
||||
export class CdrRetryableError extends Error {
|
||||
constructor(message: string) {
|
||||
super(message);
|
||||
this.name = 'CdrRetryableError';
|
||||
}
|
||||
}
|
||||
|
||||
export interface ProcessCdrBatchOptions {
|
||||
stream?: string;
|
||||
group?: string;
|
||||
consumer: string;
|
||||
count?: number;
|
||||
blockMs?: number;
|
||||
idempotencyTtlSeconds?: number;
|
||||
}
|
||||
|
||||
export interface ProcessPendingOptions {
|
||||
stream?: string;
|
||||
group?: string;
|
||||
consumer: string;
|
||||
minIdleMs?: number;
|
||||
startId?: string;
|
||||
count?: number;
|
||||
idempotencyTtlSeconds?: number;
|
||||
}
|
||||
|
||||
export interface CdrProcessSummary {
|
||||
processed: number;
|
||||
duplicates: number;
|
||||
deadlettered: number;
|
||||
pendingLeft: number;
|
||||
ackedIds: string[];
|
||||
}
|
||||
|
||||
const requiredFields = [
|
||||
'schema_version',
|
||||
'event_id',
|
||||
'idempotency_key',
|
||||
'call_id',
|
||||
'node_id',
|
||||
'opensips_instance',
|
||||
'ingress_a_ip',
|
||||
'rtpengine_node',
|
||||
'source_ip',
|
||||
'sip_code',
|
||||
'hangup_reason',
|
||||
'created_at'
|
||||
] as const;
|
||||
|
||||
export function buildCdrStreamEvent(input: CdrStreamPublishInput): CdrStreamEvent {
|
||||
const eventId = input.eventId ?? randomUUID();
|
||||
const endedAt = input.endedAt ?? '';
|
||||
|
||||
return {
|
||||
schema_version: CDR_EVENT_SCHEMA_VERSION,
|
||||
event_id: eventId,
|
||||
idempotency_key: input.idempotencyKey ?? `${input.callId}:${endedAt || eventId}`,
|
||||
call_id: input.callId,
|
||||
node_id: input.nodeId ?? 'a1',
|
||||
opensips_instance: input.opensipsInstance ?? 'opensips-a1',
|
||||
ingress_a_ip: input.ingressAIp ?? '',
|
||||
rtpengine_node: input.rtpengineNode ?? input.nodeId ?? 'a1',
|
||||
customer_id: input.customerId ?? '',
|
||||
customer_gateway_id: input.customerGatewayId ?? '',
|
||||
customer_gateway_policy_id: input.customerGatewayPolicyId ?? '',
|
||||
source_ip: input.sourceIp,
|
||||
caller: input.caller ?? '',
|
||||
callee: input.callee ?? '',
|
||||
vendor_id: input.vendorId ?? '',
|
||||
vendor_gateway_id: input.vendorGatewayId ?? '',
|
||||
line_group_id: input.lineGroupId ?? '',
|
||||
started_at: input.startedAt ?? '',
|
||||
answered_at: input.answeredAt ?? '',
|
||||
ended_at: endedAt,
|
||||
duration_sec: String(input.durationSec ?? 0),
|
||||
sip_code: String(input.sipCode),
|
||||
hangup_reason: input.hangupReason,
|
||||
recording_key: input.recordingKey ?? '',
|
||||
config_version: input.configVersion ?? '',
|
||||
created_at: input.createdAt ?? new Date().toISOString()
|
||||
};
|
||||
}
|
||||
|
||||
export function serializeCdrEvent(event: CdrStreamEvent): string[] {
|
||||
return Object.entries(event).flatMap(([field, value]) => [field, value]);
|
||||
}
|
||||
|
||||
export function parseStreamFields(fields: string[]): CdrStreamFieldMap {
|
||||
const parsed: CdrStreamFieldMap = {};
|
||||
for (let index = 0; index < fields.length; index += 2) {
|
||||
const key = fields[index];
|
||||
const value = fields[index + 1];
|
||||
if (key !== undefined && value !== undefined) {
|
||||
parsed[key] = value;
|
||||
}
|
||||
}
|
||||
return parsed;
|
||||
}
|
||||
|
||||
export function parseCdrStreamEvent(fields: string[]): ParsedCdrStreamEvent {
|
||||
const parsed = parseStreamFields(fields);
|
||||
for (const field of requiredFields) {
|
||||
if (!parsed[field]) {
|
||||
throw new Error(`CDR stream event missing required field: ${field}`);
|
||||
}
|
||||
}
|
||||
if (parsed.schema_version !== CDR_EVENT_SCHEMA_VERSION) {
|
||||
throw new Error(`Unsupported CDR schema version: ${parsed.schema_version}`);
|
||||
}
|
||||
|
||||
const durationSeconds = Number.parseInt(parsed.duration_sec ?? '0', 10);
|
||||
const sipCode = Number.parseInt(parsed.sip_code, 10);
|
||||
if (!Number.isInteger(durationSeconds) || durationSeconds < 0) {
|
||||
throw new Error(`Invalid CDR duration_sec: ${parsed.duration_sec}`);
|
||||
}
|
||||
if (!Number.isInteger(sipCode) || sipCode < 100 || sipCode > 699) {
|
||||
throw new Error(`Invalid CDR sip_code: ${parsed.sip_code}`);
|
||||
}
|
||||
|
||||
return {
|
||||
schema_version: parsed.schema_version,
|
||||
event_id: parsed.event_id,
|
||||
idempotency_key: parsed.idempotency_key,
|
||||
call_id: parsed.call_id,
|
||||
node_id: parsed.node_id,
|
||||
opensips_instance: parsed.opensips_instance,
|
||||
ingress_a_ip: parsed.ingress_a_ip,
|
||||
rtpengine_node: parsed.rtpengine_node,
|
||||
customer_id: parsed.customer_id ?? '',
|
||||
customer_gateway_id: parsed.customer_gateway_id ?? '',
|
||||
customer_gateway_policy_id: parsed.customer_gateway_policy_id ?? '',
|
||||
source_ip: parsed.source_ip,
|
||||
caller: parsed.caller ?? '',
|
||||
callee: parsed.callee ?? '',
|
||||
vendor_id: parsed.vendor_id ?? '',
|
||||
vendor_gateway_id: parsed.vendor_gateway_id ?? '',
|
||||
line_group_id: parsed.line_group_id ?? '',
|
||||
started_at: parsed.started_at ?? '',
|
||||
answered_at: parsed.answered_at ?? '',
|
||||
ended_at: parsed.ended_at ?? '',
|
||||
duration_sec: parsed.duration_sec ?? '0',
|
||||
sip_code: parsed.sip_code,
|
||||
hangup_reason: parsed.hangup_reason,
|
||||
recording_key: parsed.recording_key ?? '',
|
||||
config_version: parsed.config_version ?? '',
|
||||
created_at: parsed.created_at,
|
||||
durationSeconds,
|
||||
sipCode
|
||||
};
|
||||
}
|
||||
|
||||
export async function ensureCdrConsumerGroup(
|
||||
redis: CdrRedisCommands,
|
||||
stream = CDR_STREAM,
|
||||
group = CDR_CONSUMER_GROUP
|
||||
): Promise<void> {
|
||||
try {
|
||||
await redis.xgroup('CREATE', stream, group, '0', 'MKSTREAM');
|
||||
} catch (error) {
|
||||
if (!(error instanceof Error) || !error.message.includes('BUSYGROUP')) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export async function publishCdrEvent(
|
||||
redis: CdrRedisCommands,
|
||||
input: CdrStreamPublishInput,
|
||||
stream = CDR_STREAM
|
||||
): Promise<string> {
|
||||
const event = buildCdrStreamEvent(input);
|
||||
return redis.xadd(stream, '*', ...serializeCdrEvent(event));
|
||||
}
|
||||
|
||||
export async function processCdrBatch(
|
||||
redis: CdrRedisCommands,
|
||||
handler: CdrStreamHandler,
|
||||
options: ProcessCdrBatchOptions
|
||||
): Promise<CdrProcessSummary> {
|
||||
const stream = options.stream ?? CDR_STREAM;
|
||||
const group = options.group ?? CDR_CONSUMER_GROUP;
|
||||
const count = options.count ?? 10;
|
||||
const blockMs = options.blockMs ?? 1000;
|
||||
const response = await redis.xreadgroup(
|
||||
'GROUP',
|
||||
group,
|
||||
options.consumer,
|
||||
'COUNT',
|
||||
count,
|
||||
'BLOCK',
|
||||
blockMs,
|
||||
'STREAMS',
|
||||
stream,
|
||||
'>'
|
||||
);
|
||||
return processReadResponse(redis, handler, response, {
|
||||
stream,
|
||||
group,
|
||||
idempotencyTtlSeconds: options.idempotencyTtlSeconds ?? 86400
|
||||
});
|
||||
}
|
||||
|
||||
export async function processPendingCdrBatch(
|
||||
redis: CdrRedisCommands,
|
||||
handler: CdrStreamHandler,
|
||||
options: ProcessPendingOptions
|
||||
): Promise<CdrProcessSummary> {
|
||||
const stream = options.stream ?? CDR_STREAM;
|
||||
const group = options.group ?? CDR_CONSUMER_GROUP;
|
||||
const response = await redis.xautoclaim(
|
||||
stream,
|
||||
group,
|
||||
options.consumer,
|
||||
options.minIdleMs ?? 60000,
|
||||
options.startId ?? '0-0',
|
||||
'COUNT',
|
||||
options.count ?? 10
|
||||
);
|
||||
return processEntries(redis, handler, response[1], {
|
||||
stream,
|
||||
group,
|
||||
idempotencyTtlSeconds: options.idempotencyTtlSeconds ?? 86400
|
||||
});
|
||||
}
|
||||
|
||||
async function processReadResponse(
|
||||
redis: CdrRedisCommands,
|
||||
handler: CdrStreamHandler,
|
||||
response: RedisStreamReadResponse | null,
|
||||
options: { stream: string; group: string; idempotencyTtlSeconds: number }
|
||||
): Promise<CdrProcessSummary> {
|
||||
const entries = response?.flatMap(([, streamEntries]) => streamEntries) ?? [];
|
||||
return processEntries(redis, handler, entries, options);
|
||||
}
|
||||
|
||||
async function processEntries(
|
||||
redis: CdrRedisCommands,
|
||||
handler: CdrStreamHandler,
|
||||
entries: RedisStreamEntry[],
|
||||
options: { stream: string; group: string; idempotencyTtlSeconds: number }
|
||||
): Promise<CdrProcessSummary> {
|
||||
const summary: CdrProcessSummary = {
|
||||
processed: 0,
|
||||
duplicates: 0,
|
||||
deadlettered: 0,
|
||||
pendingLeft: 0,
|
||||
ackedIds: []
|
||||
};
|
||||
|
||||
for (const [redisId, fields] of entries) {
|
||||
try {
|
||||
const event = parseCdrStreamEvent(fields);
|
||||
const lock = await redis.set(
|
||||
`${CDR_IDEMPOTENCY_KEY_PREFIX}${event.event_id}`,
|
||||
redisId,
|
||||
'NX',
|
||||
'EX',
|
||||
options.idempotencyTtlSeconds
|
||||
);
|
||||
if (lock === null) {
|
||||
await ack(redis, options.stream, options.group, redisId, summary);
|
||||
summary.duplicates += 1;
|
||||
continue;
|
||||
}
|
||||
|
||||
const result = await handler(event, redisId);
|
||||
await ack(redis, options.stream, options.group, redisId, summary);
|
||||
if (result === 'duplicate') {
|
||||
summary.duplicates += 1;
|
||||
} else {
|
||||
summary.processed += 1;
|
||||
}
|
||||
} catch (error) {
|
||||
if (error instanceof CdrRetryableError) {
|
||||
const fieldMap = parseStreamFields(fields);
|
||||
if (fieldMap.event_id) {
|
||||
await redis.del(`${CDR_IDEMPOTENCY_KEY_PREFIX}${fieldMap.event_id}`);
|
||||
}
|
||||
summary.pendingLeft += 1;
|
||||
continue;
|
||||
}
|
||||
await moveToDeadletter(redis, redisId, fields, error);
|
||||
await ack(redis, options.stream, options.group, redisId, summary);
|
||||
summary.deadlettered += 1;
|
||||
}
|
||||
}
|
||||
|
||||
return summary;
|
||||
}
|
||||
|
||||
async function ack(
|
||||
redis: CdrRedisCommands,
|
||||
stream: string,
|
||||
group: string,
|
||||
redisId: string,
|
||||
summary: CdrProcessSummary
|
||||
): Promise<void> {
|
||||
await redis.xack(stream, group, redisId);
|
||||
summary.ackedIds.push(redisId);
|
||||
}
|
||||
|
||||
async function moveToDeadletter(redis: CdrRedisCommands, redisId: string, fields: string[], error: unknown): Promise<void> {
|
||||
const fieldMap = parseStreamFields(fields);
|
||||
await redis.xadd(
|
||||
CDR_DEADLETTER_STREAM,
|
||||
'*',
|
||||
'original_redis_id',
|
||||
redisId,
|
||||
'event_id',
|
||||
fieldMap.event_id ?? '',
|
||||
'call_id',
|
||||
fieldMap.call_id ?? '',
|
||||
'error',
|
||||
error instanceof Error ? error.message : 'Unknown CDR processing error',
|
||||
'payload',
|
||||
JSON.stringify(fieldMap),
|
||||
'deadlettered_at',
|
||||
new Date().toISOString()
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
import { Redis } from 'ioredis';
|
||||
|
||||
export const CDR_STREAM = 'stream:cdr_payload';
|
||||
export const CDR_DEADLETTER_STREAM = 'stream:cdr_deadletter';
|
||||
export const CDR_CONSUMER_GROUP = 'billing-workers';
|
||||
export const CONFIG_ACTIVE_VERSION_KEY = 'cfg:active_version';
|
||||
export const CONFIG_PREVIOUS_VERSION_KEY = 'cfg:previous_version';
|
||||
|
||||
export type RedisClient = Redis;
|
||||
|
||||
export function configVersionPrefix(version: string): string {
|
||||
return `cfg:v:${version}`;
|
||||
}
|
||||
|
||||
export function configVersionManifestKey(version: string): string {
|
||||
return `${configVersionPrefix(version)}:manifest`;
|
||||
}
|
||||
|
||||
export function createRedisClient(redisUrl: string): Redis {
|
||||
return new Redis(redisUrl, {
|
||||
lazyConnect: true,
|
||||
maxRetriesPerRequest: 3,
|
||||
enableReadyCheck: true
|
||||
});
|
||||
}
|
||||
|
||||
export * from './cdr-stream.js';
|
||||
@@ -0,0 +1,10 @@
|
||||
{
|
||||
"extends": "../../tsconfig.base.json",
|
||||
"compilerOptions": {
|
||||
"composite": true,
|
||||
"rootDir": "src",
|
||||
"outDir": "dist",
|
||||
"tsBuildInfoFile": "dist/tsconfig.tsbuildinfo"
|
||||
},
|
||||
"include": ["src/**/*.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user