From db18b558ce7496f20a05dd0aab2b15536e12a86c Mon Sep 17 00:00:00 2001 From: Hoan Luu Huu <110280845+xquanluu@users.noreply.github.com> Date: Fri, 21 Aug 2026 18:35:38 +0700 Subject: [PATCH 1/8] add sip_reason_header to status Callback- #178 (#1579) * feat: surface the SIP Reason header on call status events Carriers fronting ISDN/E1 PRI trunks put the authoritative disconnect cause in an RFC 3326 Reason header rather than in the SIP status line, e.g. SIP/2.0 408 Request Timeout Reason: Q.850 ;cause=18 Q.850 cause 18 is "no user responding" - nobody answered, not a platform fault. Different causes also arrive under the same SIP status (503 with cause=38 network out of order, or cause=41 temporary failure), so the status code alone cannot classify the outcome of the call. drachtio relays the header intact and it is present on the response object, but the feature server only read the status line from it, so the cause was lost at the application boundary and never reached the call status webhook. Rather than extract the header at each emit site, carry the SIP message that caused the status change on the callStatusChange event and derive from it in one place, so provisional responses, the 200, final failures, BYE and CANCEL are all covered by the same code and future headers cost one line. Note this partly overlaps the existing _extractCustomHeaders/sip_headers passthrough: that already exposes a Reason header arriving on a BYE, but it only runs on the hangup path, so nothing covered the outbound INVITE failure responses where a Q.850 cause matters most. sip_reason_header is a dedicated, documented field that behaves the same on every status event. sip_reason keeps meaning the status line phrase, and no key is added when there is no Reason header, so existing consumers see an unchanged payload. Adds a test:unit script so the smoke test runs without the docker testbed. Co-Authored-By: Claude Opus 5 (1M context) * docs: note that the Reason header arrives re-serialized, not verbatim Verified live end-to-end against a cluster, capturing both the external leg and the leg into the feature server: 14.226.234.142 -> 10.0.197.31:5060 Reason: Q.850 ;cause=31 (as sent) 10.0.197.31:5060 -> :5070 Reason: Q.850;cause=31 (to fs) Proxying re-serializes the header, normalizing the optional whitespace RFC 3326 permits around ';'. The same normalization appears in customer captures from an unrelated deployment, so this is drachtio behaviour, not cluster-specific. sip_reason_header is therefore the header as this process received it, not the carrier's exact bytes - worth stating outright, since the spacing inconsistency is exactly what consumers ask about, and a consumer who string-matches on the spaced form would silently never match. Co-Authored-By: Claude Opus 5 (1M context) * fix: clear a stale Reason header in redis, pass it on the alloc path, run the tests in CI Three problems found reviewing the earlier commits. 1. The redis call record kept a stale header forever. updateCallStatus assigns unconditionally, but toJSON only emits truthy values, and the same object is written to the redis hash with hmset - a MERGE. An absent key therefore left the previous status change's header in place, readable via GET /Calls/:sid: a leg that got 183 + "Reason: Q.850;cause=31", then answered on a clean 200 and completed on a plain BYE, reported a Q.850 temporary-failure cause for a call that ended normally. The webhooks were right; only the call record was wrong. Storing absence as '' overwrites it. The existing "does not linger" test passed throughout because it only exercised toJSON in memory, so this adds one that pins what survives the redis filter (verified: it fails without the fix). 2. The endpoint-allocation failure path already extracted the Reason header and relayed it to the SBC, but called _notifyCallStatusChange without it - so a FreeSWITCH 488 with "Reason: Q.850;cause=88 INCOMPATIBLE_DESTINATION" told the SBC the cause and the application nothing, which is exactly the case this feature exists to expose. That error is an fsmrf object rather than a SipMessage, so it cannot go through msg (the guard in reasonHeaderFromSipMessage would silently return undefined); the event now takes an explicit sipReasonHeader for callers holding the value already. 3. The test:unit script added with these tests was never wired into CI - the workflow runs jslint and npm test, and npm test enumerates its files explicitly and skips test/unit entirely, so the suite would have rotted unnoticed. Co-Authored-By: Claude Opus 5 (1M context) * fix: clear the stale header on the CallSession path too, and keep the CANCEL Both from PR review. 1. The previous fix only worked for SingleDialer. Storing '' on the instance is enough for callers that write the instance itself, but CallSession writes toJSON(), and toJSON() drops falsy values on purpose so the key stays out of the webhook payload - so '' never reached hmset and the stale header survived. Reproduced: after 183, toJSON has: Q.850;cause=31 raw instance value: "" instance -> redis has key: true (SingleDialer, clears) toJSON() -> redis has key: false (CallSession, stale persists) The split is by session class, not call direction: place-outdial is required only by dial.js, so SingleDialer covers dial-verb child legs while inbound, REST-created and adulting all inherit CallSession._notifyCallStatusChange - including the REST outdial path this feature was written for. Adds CallInfo.toRedisJSON(), a named projection for the merge-semantics store, so the divergence from toJSON() lives in one documented place rather than being rediscovered at each call site. The unit test asserted against the raw instance, which is why it passed while the real path was broken; it now goes through both writers and fails if either regresses. 2. The caller-abandoned race dropped the CANCEL. middleware.js had it in hand and discarded it, so the constructor's _onCancel() passed nothing whenever the CANCEL beat the application fetch - the header landed or not depending on timing. Worth closing because the 487 and its 'Request Terminated' phrase are both ours, so every abandoned inbound call looks identical: a Reason: SIP;cause=200;text="Call completed elsewhere" is what separates a forked branch losing the race from a caller who gave up, and today both are just no-answer. Confirmed on a deployed srf (5.0.27) that this is not a no-op: copyUASHeaderToUACForOnlyCancel forwards a hardcoded ['Reason', 'X-Reason'] when proxying a CANCEL, so the header does reach us. Co-Authored-By: Claude Opus 5 (1M context) * refactor: keep the Reason header out of the redis call record Reversing the earlier approach after a closer look at who reads what. The header was reaching the redis call record simply because CallInfo feeds two sinks with opposite semantics: the status webhook is an EVENT (each POST is independent, an absent key means "not in this event") while the redis record is STATE written with hmset, a MERGE (an absent key means "keep what was there"). sipReasonHeader is the first field here that can legitimately go from set back to unset - sipStatus and sipReason are only ever overwritten - so it was the first to expose the mismatch, and two rounds of fixes had to chase it because the two writers project from different bases. Rather than keep managing that, exclude it: redis holds calls that are still live, where there is usually no interesting cause yet, and by the time there is one the call is over. Call history is served from RecentCalls (the CDR), which is what the webapp reads - not GET /Calls. So the field earned very little there while costing an invariant that has to be remembered forever. Worth being explicit that this is NOT simply a revert: dropping the '' would only have stopped the empty value being written. When the header is PRESENT it still reached redis through both writers - toJSON includes it, and SingleDialer writes the instance's own properties - and was then never cleared. Keeping it out takes the same machinery as clearing it, just inverted, so this is a choice about where the field belongs rather than a saving. WEBHOOK_ONLY_FIELDS names that intent in one place, with the reasoning, so the next field with the same property has somewhere obvious to go. The test asserts both writers, since they project from different bases and checking one would pass while the other still wrote the field (verified: it fails if the exclusion is removed). Co-Authored-By: Claude Opus 5 (1M context) * fix: apply the redis projection at the boundary, and drop dead adulting plumbing From PR review. 1. There was a THIRD call-record writer. Asking each call site to remember the projection was the wrong shape, and review found the proof: this class writes the record from the status change, the recording flag AND (private only) _persistConferenceState, and the last one still passed toJSON() straight through. A conferenced caller whose BYE carried Reason: Q.850;cause=16 would have that written into the call hash on the next conference-state persist, where hmset merges and nothing can ever clear it. Fixed by wrapping updateCallStatus once where it is bound, so every write is projected and a newly added writer cannot bypass it by forgetting to ask. The call sites go back to passing plain shapes. This is a class of miss that a unit test cannot catch - the contract test passed the whole time - so the fix is structural rather than another assertion. 2. The msg/byeReq parameters added to AdultingCallSession were dead code. The only path that reacts to the far-end BYE there is the inline sd.dlg.on('destroy') handler, which calls _callReleased() and discards the request; _hangup is reachable only from CallSession.hangup() (LCC), which passes nothing. The Reason header on that leg does reach the webhook, via SingleDialer's own destroy handler - so the plumbing was not just unused but misleading about which path carries it. Removed, with a comment at the handler recording where the event actually comes from so it does not get re-added. Co-Authored-By: Claude Opus 5 (1M context) * docs: correct the projection comment in SingleDialer Copied verbatim from CallSession, where the list of writers (status change, recording flag, conference state) is accurate. SingleDialer has one writer, so the comment described code that is not there. Co-Authored-By: Claude Opus 5 (1M context) --------- Co-authored-by: Claude Opus 5 (1M context) --- .github/workflows/build.yml | 1 + lib/http-routes/api/create-call.js | 8 +- lib/middleware.js | 5 + lib/session/adulting-call-session.js | 4 + lib/session/call-info.js | 46 +++++++- lib/session/call-session.js | 19 +++- lib/session/inbound-call-session.js | 11 +- lib/session/rest-call-session.js | 2 +- lib/utils/place-outdial.js | 26 +++-- lib/utils/sip-reason.js | 41 +++++++ package.json | 1 + test/unit/sip-reason-header.test.js | 163 +++++++++++++++++++++++++++ 12 files changed, 307 insertions(+), 20 deletions(-) create mode 100644 lib/utils/sip-reason.js create mode 100644 test/unit/sip-reason-header.test.js diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index de67df35..43d0fc86 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -12,6 +12,7 @@ jobs: node-version: 20 - run: npm ci - run: npm run jslint + - run: npm run test:unit - name: Install Docker Compose run: | sudo curl -L "https://github.com/docker/compose/releases/download/1.29.2/docker-compose-$(uname -s)-$(uname -m)" -o /usr/local/bin/docker-compose diff --git a/lib/http-routes/api/create-call.js b/lib/http-routes/api/create-call.js index e2f84142..ceb8e2da 100644 --- a/lib/http-routes/api/create-call.js +++ b/lib/http-routes/api/create-call.js @@ -339,7 +339,7 @@ router.post('/', cs.callInfo.sbcCallid = prov.get('X-CID'); if ([180, 183].includes(prov.status) && prov.body) connectStream(prov.body); restDial.emit('callStatus', prov.status, !!prov.body); - cs.emit('callStatusChange', {callStatus, sipStatus: prov.status}); + cs.emit('callStatusChange', {callStatus, sipStatus: prov.status, msg: prov}); } }); connectStream(dlg.remote.sdp); @@ -352,7 +352,8 @@ router.post('/', cs.emit('callStatusChange', { callStatus: CallStatus.InProgress, sipStatus: 200, - sipReason: 'OK' + sipReason: 'OK', + msg: dlg.res }); restDial.emit('callStatus', 200); restDial.emit('connect', dlg); @@ -367,7 +368,8 @@ router.post('/', if (cs) cs.emit('callStatusChange', { callStatus, sipStatus: err.status, - sipReason: err.reason + sipReason: err.reason, + msg: err.res }); cs.callGone = true; } diff --git a/lib/middleware.js b/lib/middleware.js index a8ef630c..d8ee0b7e 100644 --- a/lib/middleware.js +++ b/lib/middleware.js @@ -120,6 +120,11 @@ module.exports = function(srf, logger) { req.once('cancel', (sipMsg) => { logger.info(`${callId} got CANCEL request`); req.locals.canceled = true; + /* keep the CANCEL: it may carry an RFC 3326 Reason (e.g. SIP;cause=200 + ;text="Call completed elsewhere"), which is the only thing distinguishing a + forked branch losing the race from a caller who simply gave up - the 487 and + its reason phrase are ours, so both look identical without it */ + req.locals.cancelReq = sipMsg; }); next(); } diff --git a/lib/session/adulting-call-session.js b/lib/session/adulting-call-session.js index 03063d87..4302b48d 100644 --- a/lib/session/adulting-call-session.js +++ b/lib/session/adulting-call-session.js @@ -23,6 +23,10 @@ class AdultingCallSession extends CallSession { this.sd = singleDialer; this.req = callInfo.req; + /* The BYE is deliberately not threaded on from here: the Completed status event for + this leg - and with it any Reason header on the BYE - is emitted by SingleDialer's + own 'destroy' handler, which already carries the request. Emitting it here too would + double-notify. */ this.sd.dlg.on('destroy', () => { this.logger.info('AdultingCallSession: called party hung up'); this._callReleased(); diff --git a/lib/session/call-info.js b/lib/session/call-info.js index 2527dc99..9dd87450 100644 --- a/lib/session/call-info.js +++ b/lib/session/call-info.js @@ -2,6 +2,21 @@ const {CallDirection, CallStatus} = require('../utils/constants'); const parseUri = require('drachtio-srf').parseUri; const crypto = require('crypto'); const {JAMBONES_API_BASE_URL} = require('../config'); +/** + * Fields that belong on the status webhook but must never enter the redis call record. + * + * That record is written with hmset, which MERGES, so a field able to go from set back to + * unset - as sipReasonHeader does, unlike sipStatus/sipReason which are only ever + * overwritten - would strand a value from an earlier status change where GET /Calls/:sid + * reports it: a leg that saw "183 + Reason: Q.850;cause=31" then answered cleanly and ended + * on a plain BYE would report a temporary-failure cause for a call that completed normally. + * + * Excluding it removes that problem instead of managing it. Redis holds calls that are + * still live, where there is usually no interesting cause yet; by the time there is one the + * call is over and the history belongs in the CDR. Add any future webhook-only field here. + */ +const WEBHOOK_ONLY_FIELDS = ['sipReasonHeader']; + /** * @classdesc Represents the common information for all calls * that is provided in call status webhooks @@ -105,11 +120,26 @@ class CallInfo { * update the status of the call * @param {string} callStatus - current call status * @param {number} sipStatus - current sip status + * @param {string} [sipReason] - reason phrase from the SIP status line + * @param {string} [sipReasonHeader] - RFC 3326 Reason header of the SIP message that caused + * this status change, if it carried one. Unlike the fields above this is assigned + * unconditionally, so that it always describes the current change rather than lingering + * from an earlier one. + * + * This is a webhook-only field, deliberately kept OUT of the redis call record: that + * record is written with hmset, which MERGES, so a field that can legitimately go from + * set back to unset - as this one does, unlike sipStatus/sipReason - would strand a cause + * from an earlier status change where GET /Calls/:sid reports it. Keeping it undefined + * when absent means realtimedb-helpers filters it out and it is never written at all, + * which removes the problem rather than managing it. Redis holds calls that are still + * live, where there is usually no interesting cause yet; by the time there is one the + * call is over and the history belongs in the CDR. */ - updateCallStatus(callStatus, sipStatus, sipReason) { + updateCallStatus(callStatus, sipStatus, sipReason, sipReasonHeader) { this.callStatus = callStatus; if (sipStatus) this.sipStatus = sipStatus; if (sipReason) this.sipReason = sipReason; + this.sipReasonHeader = sipReasonHeader; } /** @@ -132,6 +162,17 @@ class CallInfo { return this._sipHeaders; } + /** + * Project a call-status shape into what may be written to the redis call record. + * Both writers go through this, from different bases: CallSession writes the webhook + * payload, SingleDialer writes the CallInfo instance itself. + */ + static toRedisRecord(obj) { + const record = Object.assign({}, obj); + WEBHOOK_ONLY_FIELDS.forEach((f) => delete record[f]); + return record; + } + toJSON() { const obj = { callSid: this.callSid, @@ -149,7 +190,8 @@ class CallInfo { applicationSid: this.applicationSid, fsSipAddress: this.localSipAddress }; - ['parentCallSid', 'originatingSipIp', 'originatingSipTrunkName', 'callTerminationBy'].forEach((prop) => { + ['parentCallSid', 'originatingSipIp', 'originatingSipTrunkName', 'callTerminationBy', + 'sipReasonHeader'].forEach((prop) => { if (this[prop]) obj[prop] = this[prop]; }); if (typeof this.duration === 'number') obj.duration = this.duration; diff --git a/lib/session/call-session.js b/lib/session/call-session.js index 6499caff..bca9c5b0 100644 --- a/lib/session/call-session.js +++ b/lib/session/call-session.js @@ -41,6 +41,8 @@ const { NonFatalTaskError} = require('../utils/error'); const { createMediaEndpoint } = require('../utils/media-endpoint'); const { isOnhold } = require('../utils/sdp-utils'); const SttLatencyCalculator = require('../utils/stt-latency-calculator'); +const {reasonHeaderFromSipMessage} = require('../utils/sip-reason'); +const CallInfo = require('./call-info'); const sqlRetrieveQueueEventHook = `SELECT * FROM webhooks WHERE webhook_sid = ( @@ -109,7 +111,12 @@ class CallSession extends Emitter { this.tmpFiles = new Set(); if (!this.isSmsCallSession) { - this.updateCallStatus = srf.locals.dbHelpers.updateCallStatus; + /* Route every call-record write through the redis projection here rather than at each + call site: this class writes the record from more than one place (status change, + recording flag, conference state) and a new one must not be able to leak a + webhook-only field into a store with merge semantics by forgetting to ask. */ + const {updateCallStatus} = srf.locals.dbHelpers; + this.updateCallStatus = (obj, serviceUrl) => updateCallStatus(CallInfo.toRedisRecord(obj), serviceUrl); this.serviceUrl = srf.locals.serviceUrl; } @@ -2541,7 +2548,8 @@ Duration=${duration} ` this._notifyCallStatusChange({ callStatus: CallStatus.Failed, sipStatus: err.status, - sipReason: err.reason || 'Endpoint Allocation Failed' + sipReason: err.reason || 'Endpoint Allocation Failed', + sipReasonHeader }); this._callReleased(); } @@ -3105,7 +3113,7 @@ Duration=${duration} ` * @param {number} sipStatus - current sip status * @param {number} [duration] - duration of a completed call, in seconds */ - async _notifyCallStatusChange({callStatus, sipStatus, sipReason, duration, headers}) { + async _notifyCallStatusChange({callStatus, sipStatus, sipReason, duration, headers, msg, sipReasonHeader}) { if (this.callMoved) return; // manage record all call. @@ -3127,7 +3135,10 @@ Duration=${duration} ` (!duration && callStatus !== CallStatus.Completed), 'duration MUST be supplied when call completed AND ONLY when call completed'); - this.callInfo.updateCallStatus(callStatus, sipStatus, sipReason); + /* a caller that already extracted the header (endpoint allocation failure, where the + error is an fsmrf object rather than a SipMessage) passes it directly */ + this.callInfo.updateCallStatus(callStatus, sipStatus, sipReason, + sipReasonHeader ?? reasonHeaderFromSipMessage(msg)); if (typeof duration === 'number') this.callInfo.duration = duration; if (headers) this.callInfo.sipHeaders = headers; this.executeStatusCallback(callStatus, sipStatus); diff --git a/lib/session/inbound-call-session.js b/lib/session/inbound-call-session.js index 4b1bb1b6..d0f60e97 100644 --- a/lib/session/inbound-call-session.js +++ b/lib/session/inbound-call-session.js @@ -25,7 +25,10 @@ class InboundCallSession extends CallSession { // if the call was canceled before we got here, handle it if (this.req.locals.canceled) { req.locals.logger.info('InboundCallSession: constructor - call was already canceled'); - this._onCancel(); + /* the CANCEL landed before we got here, so it came in via middleware rather than + the listener below; without this the header would survive only when the CANCEL + lost the race with the application fetch */ + this._onCancel(req.locals.cancelReq); } req.once('cancel', this._onCancel.bind(this)); @@ -38,13 +41,14 @@ class InboundCallSession extends CallSession { }); } - _onCancel() { + _onCancel(cancelReq) { this.rootSpan.setAttributes({'call.termination': 'caller abandoned'}); this.callInfo.callTerminationBy = 'caller'; this._notifyCallStatusChange({ callStatus: CallStatus.NoAnswer, sipStatus: 487, - sipReason: 'Request Terminated' + sipReason: 'Request Terminated', + msg: cancelReq }); this._callReleased(); } @@ -113,6 +117,7 @@ class InboundCallSession extends CallSession { this.emit('callStatusChange', { callStatus: CallStatus.Completed, duration, + msg: req, ...(headers && {headers}) }); this._callReleased(); diff --git a/lib/session/rest-call-session.js b/lib/session/rest-call-session.js index ba131be2..a63262a3 100644 --- a/lib/session/rest-call-session.js +++ b/lib/session/rest-call-session.js @@ -66,7 +66,7 @@ class RestCallSession extends CallSession { this.callInfo.callTerminationBy = terminatedBy; const duration = moment().diff(this.dlg.connectTime, 'seconds'); const headers = this._extractCustomHeaders(req); - this.emit('callStatusChange', {callStatus: CallStatus.Completed, duration, ...(headers && {headers})}); + this.emit('callStatusChange', {callStatus: CallStatus.Completed, duration, msg: req, ...(headers && {headers})}); this.logger.info(`RestCallSession: called party hung up by ${terminatedBy}`); this._callReleased(); } diff --git a/lib/utils/place-outdial.js b/lib/utils/place-outdial.js index 843c747b..da05301c 100644 --- a/lib/utils/place-outdial.js +++ b/lib/utils/place-outdial.js @@ -17,6 +17,7 @@ const HttpRequestor = require('./http-requestor'); const WsRequestor = require('./ws-requestor'); const {makeOpusFirst, removeVideoSdp} = require('./sdp-utils'); const { createMediaEndpoint } = require('./media-endpoint'); +const {reasonHeaderFromSipMessage} = require('./sip-reason'); class SingleDialer extends Emitter { constructor({logger, sbcAddress, target, opts, application, callInfo, accountInfo, rootSpan, startSpan, dialTask, @@ -136,7 +137,12 @@ class SingleDialer extends Emitter { assert(false, `invalid dial type ${this.target.type}: must be phone, user, or sip`); } - this.updateCallStatus = srf.locals.dbHelpers.updateCallStatus; + /* Route the call-record write through the redis projection here rather than at the + call site, matching CallSession: this class has only one writer today, but the + redis record has merge semantics and a second one must not be able to leak a + webhook-only field into it by forgetting to ask. */ + const {updateCallStatus} = srf.locals.dbHelpers; + this.updateCallStatus = (obj, serviceUrl) => updateCallStatus(CallInfo.toRedisRecord(obj), serviceUrl); this.serviceUrl = srf.locals.serviceUrl; this.ep = await this._createMediaEndpoint(); @@ -221,7 +227,7 @@ class SingleDialer extends Emitter { }); }, cbProvisional: (prov) => { - const status = {sipStatus: prov.status, sipReason: prov.reason}; + const status = {sipStatus: prov.status, sipReason: prov.reason, msg: prov}; // Update call-id for sbc outbound INVITE this.callInfo.sbcCallid = prov.get('X-CID'); if ([180, 183].includes(prov.status) && prov.body) { @@ -246,7 +252,8 @@ class SingleDialer extends Emitter { this.emit('callStatusChange', { sipStatus: 200, sipReason: 'OK', - callStatus: CallStatus.InProgress + callStatus: CallStatus.InProgress, + msg: this.dlg.res }); this.logger.debug(`SingleDialer:exec call connected: ${this.callSid}`); const connectTime = this.dlg.connectTime = moment(); @@ -273,7 +280,9 @@ class SingleDialer extends Emitter { const duration = moment().diff(connectTime, 'seconds'); const headers = this._extractCustomHeaders(req); this.logger.debug('SingleDialer:exec called party hung up'); - this.emit('callStatusChange', {callStatus: CallStatus.Completed, duration, ...(headers && {headers})}); + this.emit('callStatusChange', { + callStatus: CallStatus.Completed, duration, msg: req, ...(headers && {headers}) + }); this.ep && this.ep.destroy(); }) .on('refresh', () => this.logger.info('SingleDialer:exec - dialog refreshed by uas')) @@ -315,6 +324,7 @@ class SingleDialer extends Emitter { if (err instanceof SipError) { status.sipStatus = err.status; status.sipReason = err.reason; + status.msg = err.res; if (err.status === 487) status.callStatus = CallStatus.NoAnswer; else if ([486, 600].includes(err.status)) status.callStatus = CallStatus.Busy; this.logger.info(`SingleDialer:exec outdial failure ${err.status}`); @@ -547,13 +557,14 @@ class SingleDialer extends Emitter { return Object.keys(headers).length ? headers : null; } - _notifyCallStatusChange({callStatus, sipStatus, sipReason, duration, headers}) { + _notifyCallStatusChange({callStatus, sipStatus, sipReason, duration, headers, msg, sipReasonHeader}) { assert((typeof duration === 'number' && callStatus === CallStatus.Completed) || (!duration && callStatus !== CallStatus.Completed), 'duration MUST be supplied when call completed AND ONLY when call completed'); if (this.callInfo) { - this.callInfo.updateCallStatus(callStatus, sipStatus, sipReason); + this.callInfo.updateCallStatus(callStatus, sipStatus, sipReason, + sipReasonHeader ?? reasonHeaderFromSipMessage(msg)); if (typeof duration === 'number') this.callInfo.duration = duration; if (headers) this.callInfo.sipHeaders = headers; try { @@ -562,7 +573,8 @@ class SingleDialer extends Emitter { this.logger.info(err, `SingleDialer:_notifyCallStatusChange error sending ${callStatus} ${sipStatus}`); } // update calls db - this.updateCallStatus(this.callInfo, this.serviceUrl).catch((err) => this.logger.error(err, 'redis error')); + this.updateCallStatus(this.callInfo, this.serviceUrl) + .catch((err) => this.logger.error(err, 'redis error')); } else { this.logger.info('SingleDialer:_notifyCallStatusChange: call status change before sending the outbound INVITE!!'); diff --git a/lib/utils/sip-reason.js b/lib/utils/sip-reason.js new file mode 100644 index 00000000..34b45160 --- /dev/null +++ b/lib/utils/sip-reason.js @@ -0,0 +1,41 @@ +/** + * RFC 3326 Reason header support. + * + * Carriers fronting ISDN/E1 PRI trunks put the authoritative disconnect cause in + * a Reason header rather than in the SIP status line, e.g. + * + * SIP/2.0 408 Request Timeout + * Reason: Q.850 ;cause=18 + * + * (Q.850 cause 18 is "no user responding" - i.e. nobody answered, not a fault.) + * Different Q.850 causes can arrive under the same SIP status - 503 may carry + * cause=38 (network out of order) or cause=41 (temporary failure) - so the status + * code on its own is not enough to classify the outcome of the call. We surface + * the header verbatim on call status events and leave interpretation to the + * application. + */ + +/** + * Return the Reason header of a SIP message, or undefined if it has none. + * + * A message may legally carry more than one Reason header (RFC 3326), and trunks + * that report both a SIP and a Q.850 cause commonly do. The drachtio parser joins + * repeated headers into a single comma-separated string; we pass that through + * unchanged rather than picking one of them. + * + * What reaches us is the header as the SBC relayed it, NOT necessarily the + * carrier's exact bytes: proxying re-serializes the header, which normalizes the + * optional whitespace RFC 3326 permits around ';'. A carrier's + * "Q.850 ;cause=18" therefore arrives here as "Q.850;cause=18" - confirmed on + * the wire (external leg vs the leg into this process) and visible in customer + * captures too. Consumers should parse tolerantly rather than string-match. + * + * @param {object} [msg] - a drachtio SipMessage (request or response), if we have one + * @returns {string|undefined} the Reason header value, or undefined + */ +const reasonHeaderFromSipMessage = (msg) => { + if (!msg || typeof msg.get !== 'function') return; + return msg.get('Reason') || undefined; +}; + +module.exports = {reasonHeaderFromSipMessage}; diff --git a/package.json b/package.json index 4525639a..56404cc2 100644 --- a/package.json +++ b/package.json @@ -20,6 +20,7 @@ "scripts": { "start": "node app", "test": "NODE_ENV=test JAMBONES_HOSTING=1 HTTP_POOL=1 JAMBONES_TTS_TRIM_SILENCE=1 ENCRYPTION_SECRET=foobar DRACHTIO_HOST=127.0.0.1 DRACHTIO_PORT=9060 DRACHTIO_SECRET=cymru JAMBONES_MYSQL_HOST=127.0.0.1 JAMBONES_MYSQL_PORT=3360 JAMBONES_MYSQL_USER=jambones_test JAMBONES_MYSQL_PASSWORD=jambones_test JAMBONES_MYSQL_DATABASE=jambones_test JAMBONES_REDIS_HOST=127.0.0.1 JAMBONES_REDIS_PORT=16379 JAMBONES_LOGLEVEL=error ENABLE_METRICS=0 HTTP_PORT=3000 JAMBONES_SBCS=172.38.0.10 JAMBONES_FREESWITCH=127.0.0.1:8022:JambonzR0ck$:docker-host JAMBONES_TIME_SERIES_HOST=127.0.0.1 JAMBONES_NETWORK_CIDR=172.38.0.0/16 node test/ ", + "test:unit": "node --test test/unit/*.test.js", "coverage": "./node_modules/.bin/nyc --reporter html --report-dir ./coverage npm run test", "jslint": "eslint app.js tracer.js lib", "jslint:fix": "eslint app.js tracer.js lib --fix" diff --git a/test/unit/sip-reason-header.test.js b/test/unit/sip-reason-header.test.js new file mode 100644 index 00000000..79ac9ecc --- /dev/null +++ b/test/unit/sip-reason-header.test.js @@ -0,0 +1,163 @@ +const test = require('node:test'); +const assert = require('node:assert'); +const SipMessage = require('drachtio-srf/lib/sip-parser/message'); +const CallInfo = require('../../lib/session/call-info'); +const snakeCaseKeys = require('../../lib/utils/snakecase-keys'); +const {reasonHeaderFromSipMessage} = require('../../lib/utils/sip-reason'); +const {CallDirection, CallStatus} = require('../../lib/utils/constants'); + +/* build the outbound INVITE that CallInfo is constructed from */ +const makeReq = () => { + const req = new SipMessage([ + 'INVITE sip:+971555551234@example.com SIP/2.0', + 'Call-ID: daa1269b-0b91-1240-9db3-022758ab7fff', + 'From: ;tag=abc123', + 'To: ', + 'Content-Length: 0', + '', '' + ].join('\r\n')); + req.srf = {locals: {localSipAddress: '172.30.29.123:5060'}}; + return req; +}; + +const makeSipMessage = (startLine, headers = []) => new SipMessage([ + startLine, + 'Call-ID: daa24f5d-0b91-1240-14ab-0ec7040a32ad', + ...headers, + 'Content-Length: 0', + '', '' +].join('\r\n')); + +const makeCallInfo = () => new CallInfo({ + direction: CallDirection.Outbound, + req: makeReq(), + to: '+971555551234', + callSid: '9921be00-ced0-45cb-add1-e02f9ce555ab', + accountSid: 'e43117dc-4b91-430c-82ad-74d2725f3026', + applicationSid: '72c5c38f-9bba-40ce-aa83-aaa6be55e1b5', + traceId: '615e314ac26241863b905931d9aad440' +}); + + +/* mirrors filterNullsAndObjects in realtimedb-helpers, which decides what actually + reaches the redis call hash via hmset */ +const redisFields = (callInfo) => Object.keys(callInfo) + .filter((k) => callInfo[k] !== null && typeof callInfo[k] !== 'undefined' && typeof callInfo[k] !== 'object'); + +/* the payload a call status webhook consumer actually receives */ +const statusPayload = (callInfo) => snakeCaseKeys(callInfo.toJSON(), ['customerData', 'sip', 'env_vars', 'args']); + +test('Reason header on a final failure response is surfaced as sip_reason_header', () => { + const callInfo = makeCallInfo(); + const res = makeSipMessage('SIP/2.0 408 Request Timeout', ['Reason: Q.850 ;cause=18']); + + callInfo.updateCallStatus(CallStatus.Failed, 408, 'Request Timeout', reasonHeaderFromSipMessage(res)); + const payload = statusPayload(callInfo); + + assert.strictEqual(payload.sip_reason_header, 'Q.850 ;cause=18'); + /* sip_reason must keep meaning the status-line phrase - existing consumers depend on it */ + assert.strictEqual(payload.sip_reason, 'Request Timeout'); + assert.strictEqual(payload.sip_status, 408); +}); + +test('spacing variants of the Reason header are passed through verbatim', () => { + for (const raw of ['Q.850 ;cause=31', 'Q.850;cause=31', 'Q.850 ; cause=31']) { + const callInfo = makeCallInfo(); + const res = makeSipMessage('SIP/2.0 480 Temporarily Unavailable', [`Reason: ${raw}`]); + callInfo.updateCallStatus(CallStatus.Failed, 480, 'Temporarily Unavailable', reasonHeaderFromSipMessage(res)); + assert.strictEqual(statusPayload(callInfo).sip_reason_header, raw); + } +}); + +test('a response with no Reason header adds no key to the payload', () => { + const callInfo = makeCallInfo(); + const res = makeSipMessage('SIP/2.0 503 Service Unavailable'); + + callInfo.updateCallStatus(CallStatus.Failed, 503, 'Service Unavailable', reasonHeaderFromSipMessage(res)); + const payload = statusPayload(callInfo); + + assert.ok(!('sip_reason_header' in payload), 'payload must be unchanged for carriers that send no Reason'); +}); + +test('repeated Reason headers are preserved rather than one being dropped', () => { + const callInfo = makeCallInfo(); + const res = makeSipMessage('SIP/2.0 486 Busy Here', [ + 'Reason: SIP ;cause=486 ;text="busy"', + 'Reason: Q.850 ;cause=17' + ]); + + callInfo.updateCallStatus(CallStatus.Busy, 486, 'Busy Here', reasonHeaderFromSipMessage(res)); + const header = statusPayload(callInfo).sip_reason_header; + + assert.match(header, /SIP ;cause=486/); + assert.match(header, /Q\.850 ;cause=17/); +}); + +test('a Reason header on a BYE is surfaced on the completed event', () => { + const callInfo = makeCallInfo(); + const bye = new SipMessage([ + 'BYE sip:+971555551234@example.com SIP/2.0', + 'Call-ID: daa24f5d-0b91-1240-14ab-0ec7040a32ad', + 'Reason: Q.850 ;cause=16', + 'Content-Length: 0', + '', '' + ].join('\r\n')); + + callInfo.duration = 42; + callInfo.updateCallStatus(CallStatus.Completed, 200, 'OK', reasonHeaderFromSipMessage(bye)); + + assert.strictEqual(statusPayload(callInfo).sip_reason_header, 'Q.850 ;cause=16'); +}); + +test('a Reason header does not linger onto a later status change that has none', () => { + const callInfo = makeCallInfo(); + const prov = makeSipMessage('SIP/2.0 183 Session Progress', ['Reason: Q.850 ;cause=31']); + const ok = makeSipMessage('SIP/2.0 200 OK'); + + callInfo.updateCallStatus(CallStatus.EarlyMedia, 183, 'Session Progress', reasonHeaderFromSipMessage(prov)); + assert.strictEqual(statusPayload(callInfo).sip_reason_header, 'Q.850 ;cause=31'); + + callInfo.updateCallStatus(CallStatus.InProgress, 200, 'OK', reasonHeaderFromSipMessage(ok)); + assert.ok(!('sip_reason_header' in statusPayload(callInfo)), + 'each status event must report the Reason of the message that caused it'); +}); + +test('status changes with no SIP message at all are handled', () => { + /* e.g. jambonz hanging up the call itself, or a media timeout */ + assert.strictEqual(reasonHeaderFromSipMessage(undefined), undefined); + assert.strictEqual(reasonHeaderFromSipMessage(null), undefined); + assert.strictEqual(reasonHeaderFromSipMessage({}), undefined); + + const callInfo = makeCallInfo(); + callInfo.duration = 7; + callInfo.updateCallStatus(CallStatus.Completed, 200, 'OK', reasonHeaderFromSipMessage(undefined)); + assert.ok(!('sip_reason_header' in statusPayload(callInfo))); +}); + + +test('the Reason header never enters the redis call record', () => { + const callInfo = makeCallInfo(); + const res = makeSipMessage('SIP/2.0 480 Temporarily Unavailable', ['Reason: Q.850 ;cause=31']); + + callInfo.updateCallStatus(CallStatus.Failed, 480, 'Temporarily Unavailable', reasonHeaderFromSipMessage(res)); + + /* It belongs on the webhook... */ + assert.strictEqual(statusPayload(callInfo).sip_reason_header, 'Q.850 ;cause=31'); + + /* ...and must be kept out of the redis call record, which is written with hmset - a + MERGE. A field that can go from set back to unset would otherwise strand a cause from + an earlier status change where GET /Calls/:sid reports it. + There is more than one writer and they project from DIFFERENT bases - the status + change and recording-flag writes send the webhook payload, SingleDialer sends the + CallInfo instance - so assert the projection holds for both shapes. The sessions apply + it by wrapping updateCallStatus at the boundary rather than at each call site, so a + newly added writer cannot bypass it by forgetting to ask. */ + assert.ok(!redisFields(CallInfo.toRedisRecord(callInfo.toJSON())).includes('sipReasonHeader'), + 'CallSession must not write sipReasonHeader to the call record'); + assert.ok(!redisFields(CallInfo.toRedisRecord(callInfo)).includes('sipReasonHeader'), + 'SingleDialer must not write sipReasonHeader to the call record'); + + /* the exclusion must not take anything else with it */ + assert.ok(redisFields(CallInfo.toRedisRecord(callInfo.toJSON())).includes('sipReason')); + assert.ok(redisFields(CallInfo.toRedisRecord(callInfo.toJSON())).includes('callStatus')); +}); From a6663198dab2a2c34c9101c4a4cf033e62a1fa7b Mon Sep 17 00:00:00 2001 From: Dave Horton Date: Wed, 26 Aug 2026 13:53:14 -0400 Subject: [PATCH 2/8] ci: authenticate to AWS via GitHub OIDC using the role_arn credential path (#1582) This repo is public and held a long-lived AWS access key as repository secrets (set 2023-11-22). It is replaced with short-lived credentials from the GitHub OIDC provider; no AWS key and no account id remain in the repo. create-test-db.js wrote {access_key_id, secret_access_key, aws_region} into the test database as the aws speech credential, which sends speech-utils' getAwsAuthToken down its access-key branch and calls GetSessionToken -- rejected by AWS for session credentials. The role_arn branch calls AssumeRole instead, which accepts them, and is already plumbed through db-utils.js, call-session.js and stt-task.js. The pinned speech-utils 0.2.30 already supports it, so no dependency change is needed. Fork pull requests receive neither secrets nor an OIDC token, so the credentials step is guarded by a condition; the AWS tests then skip for forks exactly as they do today. Co-authored-by: Claude Opus 5 --- .github/workflows/build.yml | 21 ++++++++++++++++----- lib/config.js | 2 ++ test/create-test-db.js | 15 ++++++++++++++- test/gather-tests.js | 3 ++- test/transcribe-tests.js | 5 +++-- 5 files changed, 37 insertions(+), 9 deletions(-) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 43d0fc86..4b367034 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -5,14 +5,24 @@ on: [push, pull_request] jobs: build: runs-on: ubuntu-latest + permissions: + id-token: write # required to request the GitHub OIDC token for AWS + contents: read steps: - uses: actions/checkout@v4 - uses: actions/setup-node@v4 with: node-version: 20 + # Pull requests from forks receive neither secrets nor an OIDC token, so this + # step is skipped for them; the AWS tests then skip too, exactly as they do + # today. Without the guard the step would hard-fail every fork PR. + - uses: aws-actions/configure-aws-credentials@v4 + if: github.event_name == 'push' || github.event.pull_request.head.repo.full_name == github.repository + with: + role-to-assume: ${{ secrets.AWS_ROLE_ARN }} + aws-region: us-east-1 - run: npm ci - run: npm run jslint - - run: npm run test:unit - name: Install Docker Compose run: | sudo curl -L "https://github.com/docker/compose/releases/download/1.29.2/docker-compose-$(uname -s)-$(uname -m)" -o /usr/local/bin/docker-compose @@ -22,8 +32,9 @@ jobs: - run: npm test env: GCP_JSON_KEY: ${{ secrets.GCP_JSON_KEY }} - AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }} - AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }} - AWS_REGION: ${{ secrets.AWS_REGION }} + # No AWS keys are stored as repository secrets. configure-aws-credentials + # above provides short-lived OIDC credentials, and AWS_ROLE_ARN makes the + # test speech credential use the AssumeRole path, which accepts them. + AWS_ROLE_ARN: ${{ secrets.AWS_ROLE_ARN }} MICROSOFT_REGION: ${{ secrets.MICROSOFT_REGION }} - MICROSOFT_API_KEY: ${{ secrets.MICROSOFT_API_KEY }} \ No newline at end of file + MICROSOFT_API_KEY: ${{ secrets.MICROSOFT_API_KEY }} diff --git a/lib/config.js b/lib/config.js index 806d3834..8158de66 100644 --- a/lib/config.js +++ b/lib/config.js @@ -93,6 +93,7 @@ const getCleanupIntervalMins = () => { const AWS_REGION = process.env.AWS_REGION; const AWS_ACCESS_KEY_ID = process.env.AWS_ACCESS_KEY_ID; const AWS_SECRET_ACCESS_KEY = process.env.AWS_SECRET_ACCESS_KEY; +const AWS_ROLE_ARN = process.env.AWS_ROLE_ARN; const AWS_SNS_PORT = parseInt(process.env.AWS_SNS_PORT, 10) || 3001; const AWS_SNS_TOPIC_ARN = process.env.AWS_SNS_TOPIC_ARN; const AWS_SNS_PORT_MAX = parseInt(process.env.AWS_SNS_PORT_MAX, 10) || 3005; @@ -198,6 +199,7 @@ module.exports = { AWS_REGION, AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY, + AWS_ROLE_ARN, AWS_SNS_PORT, AWS_SNS_TOPIC_ARN, AWS_SNS_PORT_MAX, diff --git a/test/create-test-db.js b/test/create-test-db.js index 88613596..d5a73c08 100644 --- a/test/create-test-db.js +++ b/test/create-test-db.js @@ -5,6 +5,7 @@ const {encrypt} = require('../lib/utils/encrypt-decrypt'); const { GCP_JSON_KEY, AWS_ACCESS_KEY_ID, + AWS_ROLE_ARN, AWS_SECRET_ACCESS_KEY, AWS_REGION, MICROSOFT_REGION, @@ -32,7 +33,19 @@ test('creating schema', (t) => { t.pass('adding google credentials'); sql.push(`UPDATE speech_credentials SET credential='${google_credential}' WHERE vendor='google';`); } - if (AWS_ACCESS_KEY_ID && AWS_SECRET_ACCESS_KEY) { + // Prefer role_arn. Under GitHub OIDC the ambient credentials are temporary, and + // speech-utils' getAwsAuthToken calls GetSessionToken on its access-key branch -- + // which AWS rejects for session credentials. The role_arn branch calls AssumeRole + // instead, which works with temporary credentials. + if (AWS_ROLE_ARN) { + const aws_credential = encrypt(JSON.stringify({ + role_arn: AWS_ROLE_ARN, + aws_region: AWS_REGION + })); + t.pass('adding aws credentials (role_arn)'); + sql.push(`UPDATE speech_credentials SET credential='${aws_credential}' WHERE vendor='aws';`); + } + else if (AWS_ACCESS_KEY_ID && AWS_SECRET_ACCESS_KEY) { const aws_credential = encrypt(JSON.stringify({ access_key_id: AWS_ACCESS_KEY_ID, secret_access_key: AWS_SECRET_ACCESS_KEY, diff --git a/test/gather-tests.js b/test/gather-tests.js index 904201db..11cb7b0a 100644 --- a/test/gather-tests.js +++ b/test/gather-tests.js @@ -7,6 +7,7 @@ const {provisionCallHook} = require('./utils') const { GCP_JSON_KEY, AWS_ACCESS_KEY_ID, + AWS_ROLE_ARN, AWS_SECRET_ACCESS_KEY, SONIOX_API_KEY, DEEPGRAM_API_KEY, @@ -190,7 +191,7 @@ test('\'gather\' test - microsoft', async(t) => { }); test('\'gather\' test - aws', async(t) => { - if (!AWS_ACCESS_KEY_ID || !AWS_SECRET_ACCESS_KEY) { + if (!AWS_ROLE_ARN && (!AWS_ACCESS_KEY_ID || !AWS_SECRET_ACCESS_KEY)) { t.pass('skipping aws tests'); return t.end(); } diff --git a/test/transcribe-tests.js b/test/transcribe-tests.js index a1939ced..30273b45 100644 --- a/test/transcribe-tests.js +++ b/test/transcribe-tests.js @@ -6,7 +6,8 @@ const clearModule = require('clear-module'); const {provisionCallHook} = require('./utils') const { GCP_JSON_KEY, - AWS_ACCESS_KEY_ID, + AWS_ACCESS_KEY_ID, + AWS_ROLE_ARN, AWS_SECRET_ACCESS_KEY, MICROSOFT_REGION, MICROSOFT_API_KEY, @@ -103,7 +104,7 @@ test('\'transcribe\' test - microsoft', async(t) => { }); test('\'transcribe\' test - aws', async(t) => { - if (!AWS_ACCESS_KEY_ID || !AWS_SECRET_ACCESS_KEY) { + if (!AWS_ROLE_ARN && (!AWS_ACCESS_KEY_ID || !AWS_SECRET_ACCESS_KEY)) { t.pass('skipping aws tests'); return t.end(); } From f29aa0faef38ceae74b9ad05f1e90e5059ec05d5 Mon Sep 17 00:00:00 2001 From: Rehuz Date: Sat, 5 Sep 2026 16:02:26 +0530 Subject: [PATCH 3/8] fix(tts): send stream_resumed to the same hook path as the other stream events (#1585) _onTtsStreamingResume asked for "streaming-event" while stream_open, stream_paused, stream_closed and user_interruption all ask for "/streaming-event". That one missing slash decides whether the event reaches an HTTP application at all. HttpRequestor joins a hook path to the app's baseUrl only when it is relative, and relative means starting with a slash: _isRelativeUrl(u) { return typeof u === 'string' && u.startsWith('/'); } const absUrl = this._isRelativeUrl(url) ? `${this.baseUrl}${url}` : url; "streaming-event" is neither relative nor absolute, so it is passed through unjoined and never lands on the application. The call site catches and logs the failure rather than raising it, so an application simply sees stream_paused with no matching stream_resumed and nothing anywhere says why. Present since the verb was introduced in #994. Found by reading the file, so there is no issue to close. Tests: the two new cases fail against the unpatched line and pass with it. Co-authored-by: Claude Opus 5 --- lib/session/call-session.js | 2 +- test/unit/tts-streaming-events.test.js | 64 ++++++++++++++++++++++++++ 2 files changed, 65 insertions(+), 1 deletion(-) create mode 100644 test/unit/tts-streaming-events.test.js diff --git a/lib/session/call-session.js b/lib/session/call-session.js index bca9c5b0..752f9d42 100644 --- a/lib/session/call-session.js +++ b/lib/session/call-session.js @@ -3285,7 +3285,7 @@ Duration=${duration} ` } _onTtsStreamingResume() { - this.requestor?.request('tts:streaming-event', 'streaming-event', {event_type: 'stream_resumed'}) + this.requestor?.request('tts:streaming-event', '/streaming-event', {event_type: 'stream_resumed'}) .catch((err) => this.logger.info({err}, 'CallSession:_onTtsStreamingResume - Error sending')); } diff --git a/test/unit/tts-streaming-events.test.js b/test/unit/tts-streaming-events.test.js new file mode 100644 index 00000000..5de978f6 --- /dev/null +++ b/test/unit/tts-streaming-events.test.js @@ -0,0 +1,64 @@ +const test = require('node:test'); +const assert = require('node:assert'); + +/* call-session decrypts credentials at require time, so it needs a secret present */ +process.env.ENCRYPTION_SECRET = process.env.ENCRYPTION_SECRET || 'foobar'; +process.env.JAMBONES_LOGLEVEL = process.env.JAMBONES_LOGLEVEL || 'error'; + +const CallSession = require('../../lib/session/call-session'); + +/* Only the requestor and logger are touched by these handlers, so the rest of a + CallSession is deliberately left unbuilt. */ +const makeSession = () => { + const sent = []; + const session = Object.create(CallSession.prototype); + + Object.assign(session, { + application: { + requestor: { + request: async (type, hook, payload) => { + sent.push({type, hook, payload}); + } + } + }, + logger: {info: () => {}, debug: () => {}, error: () => {}} + }); + + return {session, sent}; +}; + +test('the stream_resumed event is sent to the same hook path as every other one', async () => { + /* A hook path only gets joined to the application's baseUrl when it starts with + a slash (HttpRequestor: `_isRelativeUrl(url) ? baseUrl + url : url`). Sending + "streaming-event" instead of "/streaming-event" is neither relative nor + absolute, so it never reaches an HTTP application at all — and the failure is + swallowed by the .catch on the call site, so nothing surfaces. */ + const {session, sent} = makeSession(); + + session._onTtsStreamingResume(); + await new Promise((resolve) => setImmediate(resolve)); + + assert.strictEqual(sent.length, 1); + assert.strictEqual(sent[0].hook, '/streaming-event', + 'stream_resumed must use an absolute-from-base hook path'); + assert.strictEqual(sent[0].payload.event_type, 'stream_resumed'); +}); + +test('pause and resume agree on the hook path', async () => { + /* These two are a pair; if they ever diverge again, an application would get + one half of the pause/resume cycle and silently lose the other. */ + const {session, sent} = makeSession(); + + session._onTtsStreamingPause(); + session._onTtsStreamingResume(); + await new Promise((resolve) => setImmediate(resolve)); + + assert.deepStrictEqual( + sent.map(({hook}) => hook), + ['/streaming-event', '/streaming-event'] + ); + assert.deepStrictEqual( + sent.map(({payload}) => payload.event_type), + ['stream_paused', 'stream_resumed'] + ); +}); From 6cc3dc45f2e58ea6c5cb46378bc4d20660bf08e7 Mon Sep 17 00:00:00 2001 From: Hoan Luu Huu <110280845+xquanluu@users.noreply.github.com> Date: Wed, 9 Sep 2026 13:45:30 +0700 Subject: [PATCH 4/8] Fix/siprec survives fs transfer (#1586) * fix: carry siprec recording state across a feature server transfer transferCallToFeatureServer() stored the application, callInfo and remaining tasks, but not the SIPREC recording state, so the receiving session started at RecordingOff while the SBC was still recording. From there every live call control on the recording was rejected locally before an INFO was ever sent: stop and pause fail the guards in notifyRecordOptions, and a fresh startCallRecording is refused by the SBC as a duplicate. Carry recordOptions and the record state through the transfer and restore both on the receiving session, which also keeps propagateAnswer from issuing a second startCallRecording for a call that is already being recorded. Co-Authored-By: Claude Opus 5 (1M context) * test: add standalone smoke test for siprec across a feature server transfer Drives the failing flow end to end - answer, start SIPREC, enqueue, create the agent call on a different feature server, dequeue by callSid - while acting as the SIPREC recorder, so it can tell whether the recording kept receiving media across the REFER and whether it was BYEd when the call ended. Needs two or more feature servers and a real inbound call, so it lives outside the mocha suite and is run by hand; npm test does not pick it up. Co-Authored-By: Claude Opus 5 (1M context) * test: let the siprec smoke test originate both call legs itself Waiting for a human to dial in with music playing made the test hard to run and easy to make meaningless (a muted caller sends no RTP, so "the recording is silent" proves nothing). Add a built-in SIP endpoint that answers and streams a PCMU tone, and have the test place both legs through createCall, so a run needs only an account_sid and no carrier, DID or softphone. That makes the default mode exercise the transfer path in sbc-outbound; the previous flow is kept as SMOKE_MODE=inbound for the sbc-inbound path. Co-Authored-By: Claude Opus 5 (1M context) * fix: hand the inherited recording state to one session only Review of the previous commit turned up two ways the carried state leaks, both because it was read off the shared application object and left there: - a call transferred twice kept a stale siprecRecording from the first hop. transferCallToFeatureServer copies cs.application, and the new guard only wrote the key, so a call recorded FS1->FS2, stopped there, then moved to FS3 arrived believing it was still recording: startCallRecording rejected as "already started", stopCallRecording sent an INFO for a session that was gone. - child legs got it too. dial passes cs.application straight to ConfirmCallSession and place-outdial spreads it for the adulting session, so those sessions came up with recordState=recording_on and the parent's SRS options while their own dialog had no siprec session at all. The transferred session now takes the state off the application as it adopts it, so exactly one session owns it. Also fixes the smoke test's silence measurement, which was computed per stream and reported a whole quiet window as a gap - a legitimately silent second stream failed the run - and splits the README's requirements by mode, since the default rest mode needs no portal application and no external caller. Co-Authored-By: Claude Opus 5 (1M context) * test: drop the standalone smoke script, the harness test covers it smoke-tests/siprec-fs-transfer was a self-contained driver for the cross-feature -server SIPREC flow, written before the same scenario landed in the smoke-tester harness as TestVerb_Siprec_SurvivesFeatureServerTransfer (plus a REST-leg variant for the sbc-outbound path). The harness version is the one that runs in CI, asserts the recorded audio with Deepgram rather than counting packets, and is what proved both fixes on a real cluster - keeping a second implementation of the same flow here only invites the two to drift. Co-Authored-By: Claude Opus 5 (1M context) --------- Co-authored-by: Claude Opus 5 (1M context) --- lib/session/call-session.js | 10 ++++++++++ lib/tasks/task.js | 8 +++++++- 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/lib/session/call-session.js b/lib/session/call-session.js index 752f9d42..400d2bc4 100644 --- a/lib/session/call-session.js +++ b/lib/session/call-session.js @@ -106,6 +106,16 @@ class CallSession extends Emitter { assert(rootSpan); this._recordState = RecordState.RecordingOff; + + /* A call transferred from another feature server is still being recorded by the SBC, so + this session inherits that state - and takes it off the application, which is handed + on to child legs (dial confirm, adulting) and to any later transfer. */ + const inheritedRecording = application?.siprecRecording; + if (inheritedRecording) { + delete application.siprecRecording; + this._recordState = inheritedRecording.state || RecordState.RecordingOff; + if (inheritedRecording.options) this.recordOptions = inheritedRecording.options; + } this._notifyEvents = false; this.tmpFiles = new Set(); diff --git a/lib/tasks/task.js b/lib/tasks/task.js index 9de60b89..5bb15051 100644 --- a/lib/tasks/task.js +++ b/lib/tasks/task.js @@ -3,7 +3,7 @@ const crypto = require('crypto'); const {TaskPreconditions} = require('../utils/constants'); const { normalizeJambones } = require('@jambonz/verb-specifications'); const WsRequestor = require('../utils/ws-requestor'); -const {TaskName} = require('../utils/constants'); +const {TaskName, RecordState} = require('../utils/constants'); const {trace} = require('@opentelemetry/api'); /** @@ -345,6 +345,12 @@ class Task extends Emitter { delete obj.notifier; obj.tasks = cs.getRemainingTaskData(); obj.callInfo = cs.callInfo.toJSON(); + + /* the SBC keeps the siprec session across the REFER, so the receiving session must + inherit the state or pause/resume/stop there will be rejected as out of sync */ + if (cs.recordState !== RecordState.RecordingOff) { + obj.siprecRecording = {state: cs.recordState, options: cs.recordOptions}; + } if (opts && obj.tasks.length > 0) { const key = Object.keys(obj.tasks[0])[0]; Object.assign(obj.tasks[0][key], {_: opts}); From a4e7afac8345051501bb772c2d5280ae9d9ca88d Mon Sep 17 00:00:00 2001 From: Rehan Sanjay Venkatesan Date: Wed, 9 Sep 2026 19:13:32 +0530 Subject: [PATCH 5/8] fix(status): don't re-send pre-transfer call statuses after a transfer (#1584) A call transferred between feature servers arrives at the receiving server as a fresh INVITE, so the session runs through trying, ringing, early-media and answered again as it accepts the REFER. The application was already told all of that by the server that first handled the call, so it sees a second round of events for a call it believes is long since answered. Skip only the application-facing status callback for those four statuses when the session is a transferred call. callInfo, the redis call record and record-all-calls still run exactly as before, so the call record stays accurate and answering a transferred call still starts recording. Fixes #1392 Co-authored-by: Claude Opus 5 --- lib/session/call-session.js | 29 +++++- test/unit/transferred-call-status.test.js | 103 ++++++++++++++++++++++ 2 files changed, 131 insertions(+), 1 deletion(-) create mode 100644 test/unit/transferred-call-status.test.js diff --git a/lib/session/call-session.js b/lib/session/call-session.js index 400d2bc4..ca743a75 100644 --- a/lib/session/call-session.js +++ b/lib/session/call-session.js @@ -482,6 +482,24 @@ class CallSession extends Emitter { return this.application.transferredCall === true; } + /** + * returns true if this status was already reported to the application by the feature + * server that handled the call before it was transferred here. + * + * A transferred call arrives as a fresh INVITE, so this session runs through trying, + * ringing and answered again as it accepts the REFER. The application has already seen + * all of those from the original server, so sending them a second time looks like + * duplicate events for a call it believes is long since answered. + */ + _isStatusAlreadyReportedBeforeTransfer(callStatus) { + return this.isTransferredCall && [ + CallStatus.Trying, + CallStatus.Ringing, + CallStatus.EarlyMedia, + CallStatus.InProgress + ].includes(callStatus); + } + /** * returns true if this session is an inbound call session */ @@ -3151,7 +3169,16 @@ Duration=${duration} ` sipReasonHeader ?? reasonHeaderFromSipMessage(msg)); if (typeof duration === 'number') this.callInfo.duration = duration; if (headers) this.callInfo.sipHeaders = headers; - this.executeStatusCallback(callStatus, sipStatus); + + /* callInfo and the redis call record are still updated below; only the + application-facing notification is skipped */ + if (this._isStatusAlreadyReportedBeforeTransfer(callStatus)) { + this.logger.debug({callStatus}, + 'CallSession:_notifyCallStatusChange - suppressing status already sent before transfer'); + } + else { + this.executeStatusCallback(callStatus, sipStatus); + } // update calls db //this.logger.debug(`updating redis with ${JSON.stringify(this.callInfo)}`); diff --git a/test/unit/transferred-call-status.test.js b/test/unit/transferred-call-status.test.js new file mode 100644 index 00000000..1e95385d --- /dev/null +++ b/test/unit/transferred-call-status.test.js @@ -0,0 +1,103 @@ +const test = require('node:test'); +const assert = require('node:assert'); + +/* call-session decrypts credentials at require time, so it needs a secret present */ +process.env.ENCRYPTION_SECRET = process.env.ENCRYPTION_SECRET || 'foobar'; +process.env.JAMBONES_LOGLEVEL = process.env.JAMBONES_LOGLEVEL || 'error'; + +const CallSession = require('../../lib/session/call-session'); +const {CallStatus} = require('../../lib/utils/constants'); + +/* a CallSession is far too entangled to construct here, and none of what + _notifyCallStatusChange touches needs the constructor to have run */ +const makeSession = ({transferredCall = false, recordAllCalls = false} = {}) => { + const calls = {statusCallback: [], redis: [], callInfo: [], recorderStarted: 0, recorderStopped: 0}; + const session = Object.create(CallSession.prototype); + + Object.assign(session, { + callMoved: false, + notifiedComplete: false, + serviceUrl: 'http://127.0.0.1:3000', + application: {transferredCall, record_all_calls: false}, + accountInfo: {account: {record_all_calls: recordAllCalls}}, + backgroundTaskManager: { + newTask: (name) => { + if (name === 'record') calls.recorderStarted++; + }, + stop: (name) => { + if (name === 'record') calls.recorderStopped++; + } + }, + callInfo: { + updateCallStatus: (...args) => calls.callInfo.push(args), + toJSON: () => ({callStatus: 'x'}) + }, + logger: {debug: () => {}, info: () => {}, error: () => {}}, + executeStatusCallback: (callStatus, sipStatus) => calls.statusCallback.push({callStatus, sipStatus}), + updateCallStatus: async (obj) => { + calls.redis.push(obj); + } + }); + + return {session, calls}; +}; + +const notified = (calls) => calls.statusCallback.map(({callStatus}) => callStatus); + +const EARLY_STATUSES = [ + [CallStatus.Trying, 100], + [CallStatus.Ringing, 180], + [CallStatus.EarlyMedia, 183], + [CallStatus.InProgress, 200] +]; + +test('a transferred call does not re-notify statuses the first server already sent', async () => { + const {session, calls} = makeSession({transferredCall: true}); + + for (const [callStatus, sipStatus] of EARLY_STATUSES) { + await session._notifyCallStatusChange({callStatus, sipStatus}); + } + + assert.deepStrictEqual(notified(calls), [], + 'no early status should reach the application for a transferred call'); +}); + +test('a normal call still notifies every one of those statuses', async () => { + const {session, calls} = makeSession({transferredCall: false}); + + for (const [callStatus, sipStatus] of EARLY_STATUSES) { + await session._notifyCallStatusChange({callStatus, sipStatus}); + } + + assert.deepStrictEqual(notified(calls), EARLY_STATUSES.map(([s]) => s), + 'suppression must apply only to transferred calls'); +}); + +test('a transferred call still notifies completion', async () => { + const {session, calls} = makeSession({transferredCall: true}); + + await session._notifyCallStatusChange({callStatus: CallStatus.Completed, sipStatus: 200, duration: 12}); + + assert.deepStrictEqual(notified(calls), [CallStatus.Completed], + 'only the statuses sent before the transfer are duplicates'); +}); + +test('suppressing the notification still updates callInfo and the redis call record', async () => { + const {session, calls} = makeSession({transferredCall: true}); + + await session._notifyCallStatusChange({callStatus: CallStatus.InProgress, sipStatus: 200}); + + assert.strictEqual(calls.statusCallback.length, 0); + assert.strictEqual(calls.callInfo.length, 1, 'callInfo must still track the status'); + assert.strictEqual(calls.redis.length, 1, 'the call record must still be written to redis'); +}); + +test('suppressing the notification still starts record-all-calls', async () => { + const {session, calls} = makeSession({transferredCall: true, recordAllCalls: true}); + + await session._notifyCallStatusChange({callStatus: CallStatus.InProgress, sipStatus: 200}); + + assert.strictEqual(calls.statusCallback.length, 0); + assert.strictEqual(calls.recorderStarted, 1, + 'answering a transferred call must still start recording when record_all_calls is set'); +}); From 83db903bc69ee8601e6417ab8cb4fdc255523e7f Mon Sep 17 00:00:00 2001 From: Ben Younes Date: Wed, 9 Sep 2026 16:14:51 +0200 Subject: [PATCH 6/8] fix: fire sip:refer actionHook when the far end BYEs before any NOTIFY (#1258) (#1580) * fix: fire sip:refer actionHook when the far end BYEs before any NOTIFY When a REFER is accepted with 202 the verb waits for a NOTIFY carrying the final status of the referred call, and only fired the actionHook from the 15s timeout callback. If the far end sends BYE straight after the 202 the call session kills the task, awaitTaskDone() resolves, and exec returned without ever performing the action. Perform the action once from exec on every exit path - before the session tears the requestor down - and de-duplicate it against the final-NOTIFY path. Fix #1258 Generated by Ora Studio Vibe coded by ousamabenyounes Co-Authored-By: Ora Agent * fix: keep final_referred_call_status when a BYE races the final NOTIFY (#1580) A final NOTIFY (status >= 200) whose eventHook round trip is still in flight when the far end BYEs lost final_referred_call_status on the actionHook: kill() woke exec(), which memoised the action via _performReferAction({refer_status}) without the final status, so the NOTIFY path's later call carrying it was deduped away. On main the app received final_referred_call_status here. Record the final status synchronously in _handleNotify before the eventHook await, and include it from exec() so the memoised action carries it regardless of which exit path wins the race. Extract the 200 threshold into a named constant while touching those lines. Addresses davehorton's review on #1580. Generated by Ora Studio Vibe coded by ousamabenyounes Co-Authored-By: Ora Agent --------- Co-authored-by: Ora Agent Co-authored-by: Ben Younes <2910651+ousamabenyounes@users.noreply.github.com> --- lib/tasks/sip_refer.js | 63 +++++-- test/unit/sip-refer-action-hook.test.js | 230 ++++++++++++++++++++++++ 2 files changed, 278 insertions(+), 15 deletions(-) create mode 100644 test/unit/sip-refer-action-hook.test.js diff --git a/lib/tasks/sip_refer.js b/lib/tasks/sip_refer.js index bac8e4c4..2521e5a8 100644 --- a/lib/tasks/sip_refer.js +++ b/lib/tasks/sip_refer.js @@ -1,7 +1,16 @@ const Task = require('./task'); -const {TaskName, TaskPreconditions} = require('../utils/constants'); +const {TaskName, TaskPreconditions, KillReason} = require('../utils/constants'); const {parseUri} = require('drachtio-srf'); +/* how long we wait for a NOTIFY carrying the final status of the referred call */ +const NOTIFY_TIMEOUT_MS = 15000; + +/* SIP status returned by a far end that accepted the REFER */ +const REFER_ACCEPTED = 202; + +/* lowest SIP status that counts as the final response of the referred call */ +const SIP_FINAL_RESPONSE_MIN = 200; + /** * sends a sip REFER to transfer the existing call */ @@ -15,6 +24,8 @@ class TaskSipRefer extends Task { this.referredByDisplayName = this.data.referredByDisplayName; this.headers = this.data.headers || {}; this.eventHook = this.data.eventHook; + this._actionPromise = null; + this._finalReferredCallStatus = null; } get name() { return TaskName.SipRefer; } @@ -47,30 +58,36 @@ class TaskSipRefer extends Task { this.logger.info(`TaskSipRefer:exec - received ${this.referStatus} to REFER`); /* if we fail, fall through to next verb. If success, we should get BYE from far end */ - if (this.referStatus === 202) { + if (this.referStatus === REFER_ACCEPTED) { this._notifyTimer = setTimeout(() => { - this.logger.info('TaskSipRefer:exec - no NOTIFY received in 15 secs, exiting'); - this.performAction({refer_status: this.referStatus}) - .catch((err) => this.logger.error(err, 'TaskSipRefer:exec - error performing action')); + this.logger.info(`TaskSipRefer:exec - no NOTIFY received in ${NOTIFY_TIMEOUT_MS} ms, exiting`); this.notifyTaskDone(); - }, 15000); + }, NOTIFY_TIMEOUT_MS); await this.awaitTaskDone(); if (this._notifyTimer) { clearTimeout(this._notifyTimer); this._notifyTimer = null; } } - else { - await this.performAction({refer_status: this.referStatus}); - } + /* the far end may send BYE before any NOTIFY arrives, which kills this task while we are + awaiting above; performing the action here means the actionHook runs on every exit path - + and while the requestor is still up - rather than only when the NOTIFY timer expires. + When a final NOTIFY was already parsed (even if its eventHook round trip had not returned + before the BYE woke us), _finalReferredCallStatus carries the status so the app still gets + final_referred_call_status here rather than losing it to the memoised action */ + await this._performReferAction({ + refer_status: this.referStatus, + ...(this._finalReferredCallStatus && {final_referred_call_status: this._finalReferredCallStatus}) + }); } catch (err) { this.logger.info({err}, 'TaskSipRefer:exec - error sending REFER'); } this.referSpan?.end(); } - async kill(cs) { + async kill(cs, reason) { super.kill(cs); + this.killReason = reason || KillReason.Hangup; const {dlg} = cs; dlg.off('notify', this.notifyHandler); this.notifyTaskDone(); @@ -87,24 +104,40 @@ class TaskSipRefer extends Task { if (arr) { const status = typeof arr[1] === 'string' ? parseInt(arr[1], 10) : arr[1]; this.logger.debug(`TaskSipRefer:_handleNotify: call got status ${status}`); + /* record the final status synchronously, before the eventHook round trip below can be + interrupted by a BYE, so exec() can include it even if it wins the race to the action */ + if (status >= SIP_FINAL_RESPONSE_MIN) this._finalReferredCallStatus = status; if (this.eventHook) { const b3 = this.getTracingPropagation(); const httpHeaders = b3 && {b3}; await cs.requestor.request('verb:hook', this.eventHook, {event: 'transfer-status', call_status: status}, httpHeaders); } - if (status >= 200) { + if (status >= SIP_FINAL_RESPONSE_MIN) { this.referSpan.setAttributes({'refer.finalNotify': status}); - await this.performAction({refer_status: 202, final_referred_call_status: status}) - .catch((err) => { - this.logger.error(err, 'TaskSipRefer:exec - error performing action finalNotify'); - }); + await this._performReferAction({refer_status: REFER_ACCEPTED, final_referred_call_status: status}); this.notifyTaskDone(); } } } } + /** + * fire the verb's actionHook exactly once, whichever exit path completes the task: + * a final NOTIFY, the NOTIFY timeout, or the task being killed by an early BYE + */ + _performReferAction(results) { + if (!this._actionPromise) { + /* when the app replaced this verb with new commands, still notify it, but do not let the + hook response replace the application a second time (same contract as the dial verb) */ + this._actionPromise = this.performAction(results, this.killReason !== KillReason.Replaced) + .catch((err) => this.logger.error({err}, 'TaskSipRefer:_performReferAction - error performing action')); + } + /* callers await the shared promise, so exec() cannot return - and let the session close the + requestor - while an actionHook started by another exit path is still in flight */ + return this._actionPromise; + } + _normalizeReferHeaders(cs, dlg) { let {referTo, referredBy, referredByDisplayName} = this; diff --git a/test/unit/sip-refer-action-hook.test.js b/test/unit/sip-refer-action-hook.test.js new file mode 100644 index 00000000..634474b4 --- /dev/null +++ b/test/unit/sip-refer-action-hook.test.js @@ -0,0 +1,230 @@ +const test = require('node:test'); +const assert = require('node:assert'); +const Emitter = require('events'); +const {context} = require('@opentelemetry/api'); +const {KillReason} = require('../../lib/utils/constants'); +const proxyquire = require('proxyquire').noCallThru(); + +const ACTION_HOOK = '/refer-action'; +const EVENT_HOOK = '/refer-event'; +const REFER_TO = '+15551234567'; +const REFER_ACCEPTED = 202; +const REFER_DECLINED = 603; +const FINAL_NOTIFY_STATUS = 200; +const NOTIFY_TIMEOUT_MS = 15000; +const FAKE_TRACE_ID = '0'.repeat(32); +const FAKE_SPAN_ID = '0'.repeat(16); + +/* Task#startSpan pulls the tracer off the app module singleton; stub it so the verb can be + exercised without booting the feature server */ +const fakeSpan = { + setAttributes: () => {}, + end: () => {}, + spanContext: () => ({traceId: FAKE_TRACE_ID, spanId: FAKE_SPAN_ID}) +}; +const TaskSipRefer = proxyquire('../../lib/tasks/sip_refer', { + '../..': { + srf: {locals: {otel: {tracer: {startSpan: () => fakeSpan}}}}, + '@global': true, + '@noCallThru': true + } +}); + +const noop = () => {}; +const logger = {info: noop, debug: noop, error: noop}; + +/* minimal CallSession stand-in: a dialog that answers the REFER, and a requestor that + records every hook it is asked to fire */ +const makeCallSession = (referStatus) => { + const hookCalls = []; + const dlg = new Emitter(); + dlg.local = {uri: 'sip:jambonz@example.com'}; + dlg.remote = {uri: 'sip:carrier@10.10.10.10'}; + dlg.request = async() => ({status: referStatus}); + + return { + hookCalls, + dlg, + replacedWith: [], + replaceApplication(tasks) { + this.replacedWith.push(tasks); + }, + req: {callingNumber: '+15550000000', callingName: 'jambonz'}, + callInfo: {toJSON: () => ({call_sid: 'call-sid-under-test'})}, + requestor: { + request: async(type, hook, params) => { + hookCalls.push({type, hook, params}); + } + } + }; +}; + +const makeTask = (opts = {}) => { + const task = new TaskSipRefer(logger, {referTo: REFER_TO, actionHook: ACTION_HOOK, ...opts}); + task.ctx = context.active(); + return task; +}; + +/* let the REFER request/response round trip settle before driving the next event */ +const settle = () => new Promise((resolve) => setImmediate(resolve)); + +const makeNotify = (status) => { + const req = { + get: (name) => ('Content-Type' === name ? 'message/sipfrag;version=2.0' : undefined), + body: `SIP/2.0 ${status} OK` + }; + return {req, res: {send: () => {}}}; +}; + +test('sip:refer actionHook fires when the far end sends BYE before any NOTIFY', async() => { + const cs = makeCallSession(REFER_ACCEPTED); + const task = makeTask(); + const execPromise = task.exec(cs); + await settle(); + + /* the BYE tears the call session down, which kills the running task */ + task.kill(cs); + await execPromise; + + assert.strictEqual(cs.hookCalls.length, 1, 'actionHook should have been called exactly once'); + assert.strictEqual(cs.hookCalls[0].hook, ACTION_HOOK); + assert.strictEqual(cs.hookCalls[0].params.refer_status, REFER_ACCEPTED); +}); + +test('sip:refer actionHook fires once, with the final status, when a NOTIFY arrives', async() => { + const cs = makeCallSession(REFER_ACCEPTED); + const task = makeTask(); + const execPromise = task.exec(cs); + await settle(); + + const {req, res} = makeNotify(FINAL_NOTIFY_STATUS); + cs.dlg.emit('notify', req, res); + await execPromise; + + assert.strictEqual(cs.hookCalls.length, 1, 'actionHook should have been called exactly once'); + assert.strictEqual(cs.hookCalls[0].params.refer_status, REFER_ACCEPTED); + assert.strictEqual(cs.hookCalls[0].params.final_referred_call_status, FINAL_NOTIFY_STATUS); +}); + +test('sip:refer actionHook fires when the far end rejects the REFER', async() => { + const cs = makeCallSession(REFER_DECLINED); + const task = makeTask(); + + await task.exec(cs); + + assert.strictEqual(cs.hookCalls.length, 1, 'actionHook should have been called exactly once'); + assert.strictEqual(cs.hookCalls[0].params.refer_status, REFER_DECLINED); +}); + +test('a failing actionHook is logged, not thrown out of the verb', async() => { + const cs = makeCallSession(REFER_ACCEPTED); + cs.requestor.request = async() => { + throw new Error('actionHook unreachable'); + }; + const task = makeTask(); + const execPromise = task.exec(cs); + await settle(); + + task.kill(cs); + await assert.doesNotReject(execPromise); +}); + +test('sip:refer actionHook fires when no NOTIFY arrives before the timeout', async(t) => { + t.mock.timers.enable({apis: ['setTimeout']}); + const cs = makeCallSession(REFER_ACCEPTED); + const task = makeTask(); + const execPromise = task.exec(cs); + await settle(); + + t.mock.timers.tick(NOTIFY_TIMEOUT_MS); + await execPromise; + + assert.strictEqual(cs.hookCalls.length, 1, 'actionHook should have been called exactly once'); + assert.strictEqual(cs.hookCalls[0].params.refer_status, REFER_ACCEPTED); +}); + +test('exec waits for an in-flight actionHook when a BYE races the final NOTIFY', async() => { + const cs = makeCallSession(REFER_ACCEPTED); + let releaseHook; + const hookGate = new Promise((resolve) => { + releaseHook = resolve; + }); + const recordHook = cs.requestor.request; + cs.requestor.request = async(...args) => { + await hookGate; + return recordHook(...args); + }; + + const task = makeTask(); + const execPromise = task.exec(cs); + await settle(); + + const {req, res} = makeNotify(FINAL_NOTIFY_STATUS); + cs.dlg.emit('notify', req, res); + await settle(); + + /* the BYE lands while the actionHook request started by the NOTIFY is still in flight */ + task.kill(cs); + let execResolved = false; + execPromise.then(() => { + execResolved = true; + }); + await settle(); + assert.strictEqual(execResolved, false, 'exec must not resolve while the actionHook is in flight'); + + releaseHook(); + await execPromise; + assert.strictEqual(cs.hookCalls.length, 1, 'actionHook should have been called exactly once'); +}); + +test('the final NOTIFY status still reaches the actionHook when a BYE races the eventHook round trip', async() => { + const cs = makeCallSession(REFER_ACCEPTED); + let releaseEvent; + const eventGate = new Promise((resolve) => { + releaseEvent = resolve; + }); + const recordHook = cs.requestor.request; + cs.requestor.request = async(type, hook, params, httpHeaders) => { + /* hold the transfer-status eventHook mid round trip so the BYE can race it */ + if (EVENT_HOOK === hook) await eventGate; + return recordHook(type, hook, params, httpHeaders); + }; + + const task = makeTask({eventHook: EVENT_HOOK}); + const execPromise = task.exec(cs); + await settle(); + + /* final NOTIFY arrives; _handleNotify parks awaiting the eventHook round trip */ + const {req, res} = makeNotify(FINAL_NOTIFY_STATUS); + cs.dlg.emit('notify', req, res); + await settle(); + + /* the BYE lands while that eventHook request is still in flight, killing the task */ + task.kill(cs); + await settle(); + + releaseEvent(); + await execPromise; + + const action = cs.hookCalls.find((c) => ACTION_HOOK === c.hook); + assert.ok(action, 'actionHook should have been called'); + assert.strictEqual(action.params.refer_status, REFER_ACCEPTED); + assert.strictEqual(action.params.final_referred_call_status, FINAL_NOTIFY_STATUS, + 'final_referred_call_status must survive a BYE racing the final NOTIFY eventHook'); +}); + +test('an actionHook response replaces the application only when the verb was not itself replaced', async() => { + for (const [killReason, expectedReplacements] of [[undefined, 1], [KillReason.Replaced, 0]]) { + const cs = makeCallSession(REFER_ACCEPTED); + cs.requestor.request = async() => [{verb: 'hangup'}]; + const task = makeTask(); + const execPromise = task.exec(cs); + await settle(); + + task.kill(cs, killReason); + await execPromise; + + assert.strictEqual(cs.replacedWith.length, expectedReplacements, + `killReason=${killReason} should produce ${expectedReplacements} application replacement(s)`); + } +}); From 1d1d16e221c7db519d31819c0efe564257414f4e Mon Sep 17 00:00:00 2001 From: Ben Younes Date: Wed, 9 Sep 2026 16:54:07 +0200 Subject: [PATCH 7/8] fix: return tts:tokens-result when tts:tokens command is missing id (#1548) (#1577) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _lccTtsTokens silently returned on a missing `id`, leaving the WS client with no feedback — no tts:tokens-result, no audio, only a server-side info log. Mirror the existing missing-`tokens` handling and reply with status 'failed' / reason 'missing id' so callers that omit `id` (e.g. the Python SDK) get an actionable result. --- lib/session/call-session.js | 5 ++++- test/index.js | 1 + test/tts-tokens-missing-id-test.js | 34 ++++++++++++++++++++++++++++++ 3 files changed, 39 insertions(+), 1 deletion(-) create mode 100644 test/tts-tokens-missing-id-test.js diff --git a/lib/session/call-session.js b/lib/session/call-session.js index ca743a75..b369e0ba 100644 --- a/lib/session/call-session.js +++ b/lib/session/call-session.js @@ -1986,7 +1986,10 @@ Duration=${duration} ` if (id === undefined) { this.logger.info({opts}, 'CallSession:_lccTtsTokens - invalid command since id is missing'); - return; + return this.requestor.request('tts:tokens-result', '/tokens-result', { + status: 'failed', + reason: 'missing id' + }).catch((err) => this.logger.debug({err}, 'CallSession:_notifyTaskStatus - Error sending')); } else if (tokens === undefined) { this.logger.info({opts}, 'CallSession:_lccTtsTokens - invalid command since tokens is missing'); diff --git a/test/index.js b/test/index.js index 46f525a4..1ec44bea 100644 --- a/test/index.js +++ b/test/index.js @@ -4,6 +4,7 @@ require('./ws-requestor-unit-test'); require('./http-requestor-retry-test'); require('./http-requestor-unit-test'); require('./unit-tests'); +require('./tts-tokens-missing-id-test'); require('./hold-unhold-test'); require('./docker_start'); require('./create-test-db'); diff --git a/test/tts-tokens-missing-id-test.js b/test/tts-tokens-missing-id-test.js new file mode 100644 index 00000000..89e0f0d8 --- /dev/null +++ b/test/tts-tokens-missing-id-test.js @@ -0,0 +1,34 @@ +const test = require('tape'); +const sinon = require('sinon'); +const CallSession = require('../lib/session/call-session'); + +// Regression test for #1548: _lccTtsTokens must not silently drop a tts:tokens +// command that is missing `id`. It should reply with a tts:tokens-result of +// status 'failed'/reason 'missing id', mirroring how a missing `tokens` field +// is already handled — otherwise the WS client (e.g. the Python SDK, which does +// not send `id`) gets no feedback at all. + +const buildFakeSession = () => { + const requestor = { request: sinon.stub().resolves({}) }; + const logger = { info: sinon.stub(), debug: sinon.stub() }; + return { requestor, logger, ttsStreamingBuffer: { bufferTokens: sinon.stub().resolves({}) } }; +}; + +test('_lccTtsTokens: missing id returns a failed tts:tokens-result instead of silently dropping', async (t) => { + const session = buildFakeSession(); + + // WHEN a tts:tokens command arrives without `id` + await CallSession.prototype._lccTtsTokens.call(session, {tokens: 'hello world'}); + + // THEN a tts:tokens-result with status 'failed'/reason 'missing id' is sent back + t.ok(session.requestor.request.calledOnce, 'requestor.request should be called once'); + const [type, path, payload] = session.requestor.request.firstCall.args; + t.equal(type, 'tts:tokens-result', 'response type is tts:tokens-result'); + t.equal(path, '/tokens-result', 'response path is /tokens-result'); + t.equal(payload.status, 'failed', 'status is failed'); + t.equal(payload.reason, 'missing id', 'reason is missing id'); + + // AND tokens are never buffered for an invalid command + t.ok(session.ttsStreamingBuffer.bufferTokens.notCalled, 'tokens are not buffered when id is missing'); + t.end(); +}); From 74510b0bef3e0a936a9d2e4fad07989ef948afd2 Mon Sep 17 00:00:00 2001 From: Ben Younes Date: Sun, 13 Sep 2026 12:49:40 +0200 Subject: [PATCH 8/8] fix: retry ws opening-handshake timeout under the default ct policy (#1565) (#1576) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The ws v8 library throws a plain Error('Opening handshake has timed out') with no .code and name 'Error' when the handshake timeout (JAMBONES_WS_HANDSHAKE_TIMEOUT_MS, default 1500ms) elapses. BaseRequestor._shouldRetry classified retryable errors only by .code/.name/.statusCode, so this error matched neither the ct nor rt bucket and was never retried under the default rp=ct policy — the call failed immediately on the first attempt, silently bypassing maxReconnects. Match the handshake-timeout message in the ct bucket so a transient slow WS upgrade (e.g. a serverless/edge backend cold start) is retried as intended. Adds tape coverage for both the default-ct retry path and the rp=4xx no-retry path. --- lib/utils/base-requestor.js | 11 +++- test/ws-requestor-retry-unit-test.js | 88 ++++++++++++++++++++++++++++ 2 files changed, 97 insertions(+), 2 deletions(-) diff --git a/lib/utils/base-requestor.js b/lib/utils/base-requestor.js index 989fb765..f7851b5f 100644 --- a/lib/utils/base-requestor.js +++ b/lib/utils/base-requestor.js @@ -6,6 +6,11 @@ const timeSeries = require('@jambonz/time-series'); const {NODE_ENV, JAMBONES_TIME_SERIES_HOST} = require('../config'); let alerter ; +// The ws (v8) library throws this plain Error on an opening-handshake timeout. It carries no +// .code and .name === 'Error', so it must be matched by message to be treated as a connection +// timeout (ct) for retry purposes (issue #1565). +const WS_HANDSHAKE_TIMEOUT_MESSAGE = 'Opening handshake has timed out'; + class BaseRequestor extends Emitter { constructor(logger, account_sid, hook, secret) { super(); @@ -99,11 +104,13 @@ class BaseRequestor extends Emitter { * @returns {boolean} True if the error should be retried */ _shouldRetry(err, rpValues) { - // ct = connection timeout (ECONNREFUSED, ETIMEDOUT, etc) + // ct = connection timeout (ECONNREFUSED, ETIMEDOUT, etc). The ws opening-handshake timeout + // has no .code, so match it by message so the default ct policy retries it (issue #1565). const isCt = err.code === 'ECONNREFUSED' || err.code === 'ETIMEDOUT' || err.code === 'ECONNRESET' || - err.code === 'ECONNABORTED'; + err.code === 'ECONNABORTED' || + err.message === WS_HANDSHAKE_TIMEOUT_MESSAGE; // rt = request timeout const isRt = err.name === 'TimeoutError'; // 4xx = client errors diff --git a/test/ws-requestor-retry-unit-test.js b/test/ws-requestor-retry-unit-test.js index 5af7ee59..666c87cf 100644 --- a/test/ws-requestor-retry-unit-test.js +++ b/test/ws-requestor-retry-unit-test.js @@ -8,6 +8,9 @@ const { } = require('../lib/config'); const logger = require('pino')({level: JAMBONES_LOGLEVEL}); +// The exact message the ws v8 library throws on an opening-handshake timeout (issue #1565). +const HANDSHAKE_TIMEOUT_MESSAGE = 'Opening handshake has timed out'; + // Mock WebSocket specifically for retry testing class RetryMockWebSocket { static retryScenarios = new Map(); @@ -124,6 +127,15 @@ class RetryMockWebSocket { this.eventListeners.get('error')(err); } }); + } else if (behavior.type === 'handshake-timeout') { + // Simulate the ws v8 opening-handshake timeout: a plain Error with no + // .code property and .name === 'Error' (see issue #1565). + setImmediate(() => { + console.log(`RetryMockWebSocket: triggering handshake timeout`); + if (this.eventListeners.has('error')) { + this.eventListeners.get('error')(new Error(HANDSHAKE_TIMEOUT_MESSAGE)); + } + }); } else if (behavior.type === 'success') { // Successful connection console.log(`RetryMockWebSocket: triggering success`); @@ -572,6 +584,82 @@ test('WS Retry - rp=ct (connection timeout) should retry network errors', async t.end(); }); +test('WS Retry - handshake timeout with default (ct) policy should retry and succeed', async (t) => { + // GIVEN a handshake-timeout error (ws v8 throws a code-less Error) on the first + // attempt, and the default retry policy (no hash params => ct). Regression for #1565: + // this error type was never classified as retryable, so the call failed on attempt 1. + RetryMockWebSocket.clearScenarios(); + + const retryScenario = { + attempts: [ + { type: 'handshake-timeout' }, + { type: 'success' } + ] + }; + RetryMockWebSocket.setRetryScenario('ws://localhost:3000', retryScenario); + + const hook = { + url: 'ws://localhost:3000', // No hash parameters - defaults to ct policy + username: 'username', + password: 'password' + }; + + const params = { + callSid: 'test_handshake_timeout_ct' + }; + + // WHEN + const requestor = new WsRequestor(logger, "account_sid", hook, "webhook_secret"); + const result = await requestor.request('session:new', hook, params, {}); + + // THEN + t.ok(result, 'ws retried the handshake timeout under the default ct policy and got a response'); + t.equal(RetryMockWebSocket.getConnectionAttempts('ws://localhost:3000'), 2, + 'should have made 2 connection attempts'); + t.end(); +}); + +test('WS Retry - handshake timeout with rp=4xx should not retry', async (t) => { + // GIVEN a handshake-timeout error but a retry policy that only covers 4xx. + // Guards that the #1565 fix buckets the error as ct, not as always-retryable. + RetryMockWebSocket.clearScenarios(); + + const originalUrl = 'ws://localhost:3000#rc=2&rp=4xx'; + const cleanUrl = 'ws://localhost:3000'; + RetryMockWebSocket.setUrlMapping(cleanUrl, originalUrl); + + const retryScenario = { + attempts: [ + { type: 'handshake-timeout' } + ] + }; + RetryMockWebSocket.setRetryScenario('rc=2&rp=4xx', retryScenario); + + const hook = { + url: originalUrl, + username: 'username', + password: 'password' + }; + + const params = { + callSid: 'test_handshake_timeout_no_retry' + }; + + // WHEN & THEN + const requestor = new WsRequestor(logger, "account_sid", hook, "webhook_secret"); + try { + await requestor.request('session:new', hook, params, {}); + t.fail('Should have thrown an error'); + } catch (err) { + const errorMessage = err.message || err.toString() || String(err); + t.ok(errorMessage.includes(HANDSHAKE_TIMEOUT_MESSAGE), + 'ws properly failed without retry for handshake timeout when rp=4xx'); + t.equal(RetryMockWebSocket.getConnectionAttempts('rc=2&rp=4xx'), 1, + 'should have made only 1 connection attempt'); + t.end(); + } +}); + test('WS Retry - default behavior (no hash params) should use ct policy', async (t) => { // GIVEN RetryMockWebSocket.clearScenarios();