From da08c389493066c1983beeb440709eebc483c8b3 Mon Sep 17 00:00:00 2001 From: oliver Date: Tue, 15 Sep 2026 22:37:18 +0800 Subject: [PATCH] Fix key-alarm IM delivery blocked by DSH session sink. Run session and WhatsApp/IM sinks in parallel so a failing or slow sticky session no longer skips proactive IM delivery (0.1.34). Co-authored-by: Cursor --- lib/index.js | 33 ++++++++++++++---- package.json | 2 +- src/index.ts | 16 ++++----- src/netx/alarm-dispatch.ts | 54 +++++++++++++++++++++++++++++ test/alarm-dispatch.test.mjs | 67 ++++++++++++++++++++++++++++++++++++ 5 files changed, 157 insertions(+), 15 deletions(-) create mode 100644 src/netx/alarm-dispatch.ts create mode 100644 test/alarm-dispatch.test.mjs diff --git a/lib/index.js b/lib/index.js index 9047449..bf912fd 100644 --- a/lib/index.js +++ b/lib/index.js @@ -546,6 +546,29 @@ function resetAlarmSession() { disposeSticky(previous); } +// src/netx/alarm-dispatch.ts +async function dispatchAlarmToSinks(ctx, payload, options, sinks = {}) { + const toSession = sinks.toSession ?? deliverAlarmToSession; + const toIm = sinks.toIm ?? deliverAlarmToIm; + const jobs = []; + if (options.deliverDsh) { + jobs.push(toSession(ctx, payload, options.lang)); + } + const imOptions = { + enabled: options.deliverIm, + targets: options.imTargets, + lang: options.lang + }; + jobs.push(toIm(ctx, payload, imOptions)); + const results = await Promise.allSettled(jobs); + for (const result of results) { + if (result.status === "rejected") { + ctx.logger.warn("netxops alarm-push: sink failed: %s", result.reason); + } + } + return results; +} + // src/netx/capability-groups.ts var DEFAULT_CAPABILITY_GROUPS = Object.freeze({ ops: Object.freeze({ inPreset: true, public: false }), @@ -4558,12 +4581,10 @@ function apply(ctx, config = Config({})) { token, logger: ctx.logger, onAlarm: async (payload) => { - if (deliverDsh) { - await deliverAlarmToSession(ctx, payload, lang); - } - await deliverAlarmToIm(ctx, payload, { - enabled: deliverIm, - targets: imTargets, + await dispatchAlarmToSinks(ctx, payload, { + deliverDsh, + deliverIm, + imTargets, lang }); } diff --git a/package.json b/package.json index 2b2c34a..e78849d 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "dsh-netxops", - "version": "0.1.33", + "version": "0.1.34", "description": "DeepSeek Harness Netx Ops: ops/topology + IM delivery + session export + knowledge-base MANIFEST/_skills + localSkills", "license": "MIT", "type": "module", diff --git a/src/index.ts b/src/index.ts index 5d18000..a91e974 100644 --- a/src/index.ts +++ b/src/index.ts @@ -35,8 +35,8 @@ import { publishAlarmPushStatus, resetAlarmPushStatus, } from './netx/alarm-push-status.ts' -import { deliverAlarmToIm } from './netx/alarm-im.ts' -import { deliverAlarmToSession, resetAlarmSession } from './netx/alarm-session.ts' +import { dispatchAlarmToSinks } from './netx/alarm-dispatch.ts' +import { resetAlarmSession } from './netx/alarm-session.ts' import { capabilityGroupsFromSettings, groupsForPlane, @@ -307,18 +307,18 @@ export function apply(ctx: Context, config: Config = Config({})): void { imBotId: current.imBotId ?? '', imTargetId: current.imTargetId ?? '', }) + // Prefer target list presence; UI also mirrors this into alarmDeliverIm. const deliverIm = imTargets.length > 0 stopAlarmPush = startAlarmPushClient({ apiUrl, token, logger: ctx.logger, onAlarm: async (payload) => { - if (deliverDsh) { - await deliverAlarmToSession(ctx, payload, lang) - } - await deliverAlarmToIm(ctx, payload, { - enabled: deliverIm, - targets: imTargets, + // Parallel sinks: a stuck/failing DSH sticky session must not block WhatsApp/IM. + await dispatchAlarmToSinks(ctx, payload, { + deliverDsh, + deliverIm, + imTargets, lang, }) }, diff --git a/src/netx/alarm-dispatch.ts b/src/netx/alarm-dispatch.ts new file mode 100644 index 0000000..888254d --- /dev/null +++ b/src/netx/alarm-dispatch.ts @@ -0,0 +1,54 @@ +/** + * Fan-out one key-alarm to optional DSH session + IM sinks. + * Sinks run in parallel so a failing/slow session path cannot block WhatsApp/IM. + */ + +import type { Context } from '@deepseek-ai/cordis' +import { deliverAlarmToIm, type AlarmImDeliveryOptions } from './alarm-im.ts' +import { deliverAlarmToSession } from './alarm-session.ts' +import type { KeyAlarmPayload } from './alarm-push.ts' +import type { ImDeliveryTarget } from './im-targets.ts' + +export interface AlarmSinkDispatchOptions { + deliverDsh: boolean + deliverIm: boolean + imTargets: readonly ImDeliveryTarget[] + lang: string +} + +export interface AlarmSinkFns { + toSession?: typeof deliverAlarmToSession + toIm?: typeof deliverAlarmToIm +} + +/** + * Deliver to every enabled sink; never let one sink's rejection cancel the others. + * @returns settled results in order: [session?, im] (session omitted when deliverDsh is false). + */ +export async function dispatchAlarmToSinks( + ctx: Context, + payload: KeyAlarmPayload, + options: AlarmSinkDispatchOptions, + sinks: AlarmSinkFns = {}, +): Promise[]> { + const toSession = sinks.toSession ?? deliverAlarmToSession + const toIm = sinks.toIm ?? deliverAlarmToIm + const jobs: Array> = [] + if (options.deliverDsh) { + jobs.push(toSession(ctx, payload, options.lang)) + } + const imOptions: AlarmImDeliveryOptions = { + enabled: options.deliverIm, + targets: options.imTargets, + lang: options.lang, + } + jobs.push(toIm(ctx, payload, imOptions)) + + const results = await Promise.allSettled(jobs) + for (const result of results) { + if (result.status === 'rejected') { + ctx.logger.warn('netxops alarm-push: sink failed: %s', result.reason) + } + } + return results +} diff --git a/test/alarm-dispatch.test.mjs b/test/alarm-dispatch.test.mjs new file mode 100644 index 0000000..155b802 --- /dev/null +++ b/test/alarm-dispatch.test.mjs @@ -0,0 +1,67 @@ +import assert from 'node:assert/strict' +import { test } from 'node:test' + +import { dispatchAlarmToSinks } from '../src/netx/alarm-dispatch.ts' + +test('dispatchAlarmToSinks still delivers IM when DSH session rejects', async () => { + let imCalls = 0 + const warnings = [] + const ctx = { + logger: { + warn: (...args) => { warnings.push(args.map(String).join(' ')) }, + info() {}, + }, + } + const results = await dispatchAlarmToSinks( + ctx, + { + action: 'inserted', + rule_label: 'Power Down', + alarm_key: 'k1', + ne: { host_name: 'PE1' }, + }, + { + deliverDsh: true, + deliverIm: true, + imTargets: [{ botId: 'bot', targetId: 'tgt' }], + lang: 'zh', + }, + { + toSession: async () => { + throw new Error('sticky session exploded') + }, + toIm: async () => { + imCalls += 1 + }, + }, + ) + assert.equal(imCalls, 1) + assert.equal(results.length, 2) + assert.equal(results[0].status, 'rejected') + assert.equal(results[1].status, 'fulfilled') + assert.ok(warnings.some((line) => line.includes('sink failed'))) +}) + +test('dispatchAlarmToSinks skips session job when deliverDsh is false', async () => { + let sessionCalls = 0 + let imCalls = 0 + const ctx = { logger: { warn() {}, info() {} } } + const results = await dispatchAlarmToSinks( + ctx, + { action: 'inserted', alarm_key: 'k1' }, + { + deliverDsh: false, + deliverIm: true, + imTargets: [{ botId: 'bot', targetId: 'tgt' }], + lang: 'zh', + }, + { + toSession: async () => { sessionCalls += 1 }, + toIm: async () => { imCalls += 1 }, + }, + ) + assert.equal(sessionCalls, 0) + assert.equal(imCalls, 1) + assert.equal(results.length, 1) + assert.equal(results[0].status, 'fulfilled') +})