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 <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-15 22:37:18 +08:00
parent 3a731389f0
commit da08c38949
5 changed files with 157 additions and 15 deletions

View file

@ -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
});
}

View file

@ -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",

View file

@ -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,
})
},

View file

@ -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<PromiseSettledResult<void>[]> {
const toSession = sinks.toSession ?? deliverAlarmToSession
const toIm = sinks.toIm ?? deliverAlarmToIm
const jobs: Array<Promise<void>> = []
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
}

View file

@ -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')
})