diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index de67df35..4b367034 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -5,11 +5,22 @@ 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 - name: Install Docker Compose @@ -21,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/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 fba37786..7d7443fe 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 = ( @@ -104,12 +106,27 @@ 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(); 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; } @@ -465,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 */ @@ -1980,7 +2015,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'); @@ -2570,7 +2608,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(); } @@ -3134,7 +3173,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. @@ -3156,10 +3195,22 @@ 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); + + /* 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)}`); @@ -3303,7 +3354,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/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/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/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}); 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/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 785dfb18..a2b87b38 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/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/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/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(); } 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(); +}); 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')); +}); 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)`); + } +}); 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'); +}); 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'] + ); +}); 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();