From 00817fffc4794f5f555d4ece3d1abdca21405a5c Mon Sep 17 00:00:00 2001 From: Emmanuel BUU Date: Thu, 1 Oct 2026 17:40:57 +0200 Subject: [PATCH] Publisher: event state publication (RFC 3903, PUBLISH) JsSIP can SUBSCRIBE to an event package and NOTIFY one, but cannot PUBLISH its own state to a presence server (or any RFC 3903 state agent). Applications had to build the request on SIPMessage and RequestSender internals. `ua.publish(target, eventName, contentType, options)` returns a Publisher, modelled on Subscriber: - publish(body): the initial publication, or a modification. The SIP-ETag of each 2xx names the publication from then on. - Refresh before expiry, with no body and SIP-If-Match, scheduled as Subscriber does to survive Chrome's intensive timer throttling. - terminate(): removal, Expires: 0 and SIP-If-Match; nothing is sent if nothing was published. - At most one PUBLISH in flight: a body given meanwhile is sent when the answer comes back, so the entity-tag always names the last accepted publication. - 412 (the server lost the entity-tag: restart, or expiry while the page was frozen): published again from scratch, once. - 423: published again with the Min-Expires given. - Any other failure terminates the Publisher; `terminated` carries the code and the final response, so that 489/405/501 (no PUBLISH here) can be told from a passing failure. - One Call-ID and an increasing CSeq for all requests, as for REGISTER. - UA.stop() removes every publication before closing, so that watchers do not see a stale state until it expires. Co-Authored-By: Claude Opus 5.5 --- src/Constants.d.ts | 1 + src/Constants.js | 1 + src/Publisher.d.ts | 60 ++++ src/Publisher.js | 447 ++++++++++++++++++++++++++++ src/UA.d.ts | 8 + src/UA.js | 37 +++ src/test/include/ResponderSocket.ts | 160 ++++++++++ src/test/test-Publisher.ts | 369 +++++++++++++++++++++++ 8 files changed, 1083 insertions(+) create mode 100644 src/Publisher.d.ts create mode 100644 src/Publisher.js create mode 100644 src/test/include/ResponderSocket.ts create mode 100644 src/test/test-Publisher.ts diff --git a/src/Constants.d.ts b/src/Constants.d.ts index b12d171a0..f0a0ca2f4 100644 --- a/src/Constants.d.ts +++ b/src/Constants.d.ts @@ -46,6 +46,7 @@ export const INVITE = 'INVITE'; export const MESSAGE = 'MESSAGE'; export const NOTIFY = 'NOTIFY'; export const OPTIONS = 'OPTIONS'; +export const PUBLISH = 'PUBLISH'; export const REGISTER = 'REGISTER'; export const REFER = 'REFER'; export const UPDATE = 'UPDATE'; diff --git a/src/Constants.js b/src/Constants.js index 0212f5626..1f1a5644f 100644 --- a/src/Constants.js +++ b/src/Constants.js @@ -59,6 +59,7 @@ module.exports = { MESSAGE: 'MESSAGE', NOTIFY: 'NOTIFY', OPTIONS: 'OPTIONS', + PUBLISH: 'PUBLISH', REGISTER: 'REGISTER', REFER: 'REFER', UPDATE: 'UPDATE', diff --git a/src/Publisher.d.ts b/src/Publisher.d.ts new file mode 100644 index 000000000..e0b69010b --- /dev/null +++ b/src/Publisher.d.ts @@ -0,0 +1,60 @@ +import { EventEmitter } from 'events'; +import { IncomingResponse } from './SIPMessage'; +import { UA } from './UA'; +import { URI } from './URI'; + +declare enum PublisherTerminatedCode { + PUBLISH_RESPONSE_TIMEOUT = 0, + PUBLISH_TRANSPORT_ERROR = 1, + PUBLISH_NON_OK_RESPONSE = 2, + PUBLISH_WRONG_OK_RESPONSE = 3, + PUBLISH_AUTHENTICATION_FAILED = 4, + UNPUBLISHED = 5, +} + +declare enum PublisherState { + STATE_INIT = 0, + STATE_PUBLISHED = 1, + STATE_TERMINATED = 2, +} + +export interface PublisherEventMap { + published: [response: IncomingResponse]; + terminated: [ + terminationCode: PublisherTerminatedCode, + response: IncomingResponse | undefined, + ]; +} + +export interface PublisherParams { + from_uri?: URI; + from_display_name?: string; + to_uri?: URI; + to_display_name?: string; +} + +export interface PublisherOptions { + expires?: number; + params?: PublisherParams; + extraHeaders?: string[]; +} + +export class Publisher extends EventEmitter { + constructor( + ua: UA, + target: string | URI, + eventName: string, + contentType: string, + options?: PublisherOptions + ); + publish(body: string): void; + terminate(): void; + get state(): PublisherState; + get etag(): string | null; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + set data(_data: any); + // eslint-disable-next-line @typescript-eslint/no-explicit-any + get data(): any; + static get C(): typeof PublisherTerminatedCode & typeof PublisherState; + get C(): typeof PublisherTerminatedCode & typeof PublisherState; +} diff --git a/src/Publisher.js b/src/Publisher.js new file mode 100644 index 000000000..de4169a00 --- /dev/null +++ b/src/Publisher.js @@ -0,0 +1,447 @@ +const EventEmitter = require('events').EventEmitter; +const Exceptions = require('./Exceptions'); +const Logger = require('./Logger'); +const JsSIP_C = require('./Constants'); +const Utils = require('./Utils'); +const Grammar = require('./Grammar'); +const SIPMessage = require('./SIPMessage'); +const RequestSender = require('./RequestSender'); + +const logger = new Logger('Publisher'); + +/** + * Termination codes. + */ +const C = { + // Termination codes. + PUBLISH_RESPONSE_TIMEOUT: 0, + PUBLISH_TRANSPORT_ERROR: 1, + PUBLISH_NON_OK_RESPONSE: 2, + PUBLISH_WRONG_OK_RESPONSE: 3, + PUBLISH_AUTHENTICATION_FAILED: 4, + UNPUBLISHED: 5, + + // Publisher states. + STATE_INIT: 0, + STATE_PUBLISHED: 1, + STATE_TERMINATED: 2, + + // Default expires value, in seconds. + DEFAULT_EXPIRES_SEC: 3600, +}; + +/** + * RFC 3903 event state publication. + * + * - publish(body) sends the initial publication; the 2xx carries the + * SIP-ETag that names it from then on. + * - Before it expires, the publication is refreshed: no body, SIP-If-Match. + * - publish(body) again modifies it: SIP-If-Match and the new body. + * - terminate() removes it: Expires: 0 and SIP-If-Match. + * + * At most one PUBLISH is in flight: a body given meanwhile is sent when the + * answer comes back, so the entity-tag always names the last publication + * the server accepted. + * + * The server may lose the entity-tag (it restarted, or the publication + * expired): it answers 412, and the body is published again from scratch, + * once. A 423 is answered with the Min-Expires it gives. Any other failure + * terminates the Publisher. + */ +module.exports = class Publisher extends EventEmitter { + /** + * Expose C object. + */ + static get C() { + return C; + } + + /** + * @param {UA} ua - reference to JsSIP.UA + * @param {string} target - Request-URI and To, usually our own address. + * @param {string} eventName - Event header value. + * @param {string} contentType - Content-Type of the published bodies. + * + * @param {PublisherOptions} options - optional parameters. + * @param {number} expires - Expires header value. Default is 3600. + * @param {RequestParams} params - Will have priority over ua.configuration. + * @param {Array} extraHeaders - Additional SIP headers. + */ + constructor( + ua, + target, + eventName, + contentType, + { expires, params, extraHeaders } = {} + ) { + logger.debug('new'); + + super(); + + if (!target) { + throw new TypeError('Not enough arguments: Missing target'); + } + + if (!eventName) { + throw new TypeError('Not enough arguments: Missing eventName'); + } + + if (!contentType) { + throw new TypeError('Not enough arguments: Missing contentType'); + } + + if (Grammar.parse(eventName, 'Event') === -1) { + throw new TypeError(`Invalid eventName: ${eventName}`); + } + + this._ua = ua; + this._target = this._ua.normalizeTarget(target); + + if (!this._target) { + throw new TypeError(`Invalid target: ${target}`); + } + + if (!Utils.isDecimal(expires) || expires <= 0) { + expires = C.DEFAULT_EXPIRES_SEC; + } + + this._expires = expires; + this._event_name = eventName; + this._content_type = contentType; + this._extraHeaders = Utils.cloneArray(extraHeaders); + + // One Call-ID for all the PUBLISH requests and a CSeq that only goes + // up, as for REGISTER (RFC 3261 10.2). + this._params = Utils.cloneObject(params); + + if (!this._params.to_uri) { + this._params.to_uri = this._target; + } + + this._params.call_id = Utils.createRandomToken(22); + this._params.cseq = 0; + + this._state = C.STATE_INIT; + + // Entity-tag of our publication, from the last 2xx. + this._etag = null; + + // The last body we were given, and the one the server holds. + this._body = null; + this._published_body = null; + + // A PUBLISH is waiting for its final response. + this._sending = false; + + // terminate() was called: the next request is the removal. + this._removing = false; + + // The last 412 was answered by publishing from scratch. + this._recovering = false; + + this._refresh_timer = null; + + // Custom publisher empty object for high level use. + this._data = {}; + } + + // Expose Publisher constants as a property of the Publisher instance. + get C() { + return C; + } + + get state() { + return this._state; + } + + /** + * Entity-tag of the publication, null until the first 2xx. + */ + get etag() { + return this._etag; + } + + get data() { + return this._data; + } + + set data(_data) { + this._data = _data; + } + + /** + * User API + */ + + /** + * Publish this body: the initial publication, or a modification. + * @param {string} body - publication body. + */ + publish(body) { + logger.debug('publish()'); + + if (this._state === C.STATE_TERMINATED || this._removing) { + throw new Exceptions.InvalidStateError(this._state); + } + + if (!body) { + throw new TypeError('Not enough arguments: Missing body'); + } + + this._ua.newPublisher(this); + + this._body = body; + + clearTimeout(this._refresh_timer); + this._refresh_timer = null; + + this._send(); + } + + /** + * Remove the publication (Expires: 0), and terminate. + */ + terminate() { + logger.debug('terminate()'); + + if (this._state === C.STATE_TERMINATED || this._removing) { + return; + } + + this._removing = true; + + clearTimeout(this._refresh_timer); + this._refresh_timer = null; + + // Nothing published and nothing on its way: nothing to remove. + if (this._etag === null && !this._sending) { + this._terminate(C.UNPUBLISHED); + + return; + } + + this._send(); + } + + /** + * Private API. + */ + + _send() { + if (this._sending) { + return; + } + + const headers = Utils.cloneArray(this._extraHeaders); + const removal = this._removing; + let body = null; + + headers.push(`Event: ${this._event_name}`); + + if (removal) { + headers.push('Expires: 0'); + } else { + headers.push(`Expires: ${this._expires}`); + } + + if (this._etag !== null) { + headers.push(`SIP-If-Match: ${this._etag}`); + } + + // The initial publication and a modification carry the body; a + // refresh and a removal do not. + if ( + !removal && + (this._etag === null || this._body !== this._published_body) + ) { + body = this._body; + headers.push(`Content-Type: ${this._content_type}`); + } + + this._params.cseq += 1; + + const request = new SIPMessage.OutgoingRequest( + JsSIP_C.PUBLISH, + this._target, + this._ua, + this._params, + headers, + body || undefined + ); + + this._sending = true; + + const request_sender = new RequestSender(this._ua, request, { + onRequestTimeout: () => { + this._sending = false; + this._terminate(C.PUBLISH_RESPONSE_TIMEOUT); + }, + onTransportError: () => { + this._sending = false; + this._terminate(C.PUBLISH_TRANSPORT_ERROR); + }, + // RequestSender increased the CSeq of its copy of the request. + onAuthenticated: () => { + this._params.cseq += 1; + }, + onReceiveResponse: response => { + if (response.status_code < 200) { + return; + } + + this._sending = false; + this._receiveResponse(response, body, removal); + }, + }); + + request_sender.send(); + } + + _receiveResponse(response, sentBody, removal) { + if (this._state === C.STATE_TERMINATED) { + return; + } + + const status_code = response.status_code; + + // Answer to the removal. + if (removal) { + if ( + (status_code >= 200 && status_code < 300) || + // The server had already forgotten the publication. + status_code === 412 + ) { + this._terminate(C.UNPUBLISHED, response); + } else { + this._terminate(this._failureCode(status_code), response); + } + + return; + } + + if (status_code >= 200 && status_code < 300) { + const etag = response.getHeader('SIP-ETag'); + + // RFC 3903 11.3: the 2xx carries the entity-tag. + if (!etag) { + logger.warn('2xx response without SIP-ETag'); + this._terminate(C.PUBLISH_WRONG_OK_RESPONSE, response); + + return; + } + + this._etag = etag.trim(); + this._recovering = false; + + if (sentBody !== null) { + this._published_body = sentBody; + } + + this._state = C.STATE_PUBLISHED; + + logger.debug('emit "published"'); + + this.emit('published', response); + + // terminate() was called while this PUBLISH was in flight. + if (this._removing) { + this._send(); + + return; + } + + // publish() was called while this PUBLISH was in flight. + if (this._body !== this._published_body) { + this._send(); + + return; + } + + let expires = parseInt(response.getHeader('expires')); + + if (!Utils.isDecimal(expires) || expires <= 0) { + expires = this._expires; + } + + this._scheduleRefresh(expires); + + return; + } + + // RFC 3903 6: the server does not know our entity-tag any more. + // Publish from scratch, once. + if (status_code === 412 && !this._recovering) { + logger.debug('412, publishing again from scratch'); + + this._recovering = true; + this._etag = null; + this._published_body = null; + + if (this._removing) { + this._terminate(C.UNPUBLISHED, response); + } else { + this._send(); + } + + return; + } + + // RFC 3903 5: interval too brief. + if (status_code === 423 && response.hasHeader('min-expires')) { + const min_expires = parseInt(response.getHeader('min-expires')); + + if (Utils.isDecimal(min_expires) && min_expires > this._expires) { + logger.debug(`423, publishing again with Expires: ${min_expires}`); + + this._expires = min_expires; + this._send(); + + return; + } + } + + this._terminate(this._failureCode(status_code), response); + } + + _failureCode(status_code) { + return status_code === 401 || status_code === 407 + ? C.PUBLISH_AUTHENTICATION_FAILED + : C.PUBLISH_NON_OK_RESPONSE; + } + + _scheduleRefresh(expires) { + // Same margin as Subscriber, for Chrome intensive timer throttling. + const timeout = + expires >= 140 + ? (expires * 1000) / 2 + + Math.floor((expires / 2 - 70) * 1000 * Math.random()) + : expires * 1000 - 5000; + + logger.debug( + `next PUBLISH will be sent in ${Math.floor(timeout / 1000)} sec` + ); + + clearTimeout(this._refresh_timer); + this._refresh_timer = setTimeout(() => { + this._refresh_timer = null; + this._send(); + }, timeout); + } + + _terminate(terminationCode, response = undefined) { + if (this._state === C.STATE_TERMINATED) { + return; + } + + this._state = C.STATE_TERMINATED; + + clearTimeout(this._refresh_timer); + this._refresh_timer = null; + + this._ua.destroyPublisher(this); + + logger.debug(`emit "terminated" code=${terminationCode}`); + + this.emit('terminated', terminationCode, response); + } +}; diff --git a/src/UA.d.ts b/src/UA.d.ts index ae562cc6f..c5dec5571 100644 --- a/src/UA.d.ts +++ b/src/UA.d.ts @@ -16,6 +16,7 @@ import { import { Message, SendMessageOptions } from './Message'; import { Registrator } from './Registrator'; import { Notifier } from './Notifier'; +import { Publisher, PublisherOptions } from './Publisher'; import { Subscriber } from './Subscriber'; import { URI } from './URI'; import { causes } from './Constants'; @@ -268,6 +269,13 @@ export class UA extends EventEmitter { options?: NotifierOptions ): Notifier; + publish( + target: string | URI, + eventName: string, + contentType: string, + options?: PublisherOptions + ): Publisher; + terminateSessions(options?: TerminateOptions): void; isRegistered(): boolean; diff --git a/src/UA.js b/src/UA.js index 687ff37a2..41a3fbd6e 100644 --- a/src/UA.js +++ b/src/UA.js @@ -7,6 +7,7 @@ const Registrator = require('./Registrator'); const RTCSession = require('./RTCSession'); const Subscriber = require('./Subscriber'); const Notifier = require('./Notifier'); +const Publisher = require('./Publisher'); const Message = require('./Message'); const Options = require('./Options'); const Transactions = require('./Transactions'); @@ -75,6 +76,9 @@ module.exports = class UA extends EventEmitter { // User actions outside any session/dialog (MESSAGE/OPTIONS). this._applicants = {}; + // Publisher instances with a publication sent or established. + this._publishers = new Set(); + this._sessions = {}; this._transport = null; this._contact = null; @@ -261,6 +265,15 @@ module.exports = class UA extends EventEmitter { return new Notifier(this, subscribe, contentType, options); } + /** + * Create publisher instance + */ + publish(target, eventName, contentType, options) { + logger.debug('publish()'); + + return new Publisher(this, target, eventName, contentType, options); + } + /** * Send a SIP OPTIONS. * @@ -310,6 +323,16 @@ module.exports = class UA extends EventEmitter { return; } + // Remove every publication (RFC 3903 4.5), before the un-REGISTER. + // The PUBLISH must leave before the UA is closed: others would + // otherwise see our state until the publication expires. + for (const publisher of this._publishers) { + logger.debug('closing publisher'); + try { + publisher.terminate(); + } catch (error) {} + } + // Close registrator. this._registrator.close(); @@ -482,6 +505,20 @@ module.exports = class UA extends EventEmitter { delete this._dialogs[dialog.id]; } + /** + * Publisher sent its initial PUBLISH. + */ + newPublisher(publisher) { + this._publishers.add(publisher); + } + + /** + * Publisher terminated. + */ + destroyPublisher(publisher) { + this._publishers.delete(publisher); + } + /** * new Message */ diff --git a/src/test/include/ResponderSocket.ts b/src/test/include/ResponderSocket.ts new file mode 100644 index 000000000..370547dd4 --- /dev/null +++ b/src/test/include/ResponderSocket.ts @@ -0,0 +1,160 @@ +// ResponderSocket answers every request the UA sends with the response +// built by the given handler. Used to play the server side of a client +// transaction (SUBSCRIBE, PUBLISH...) without a network. + +import type { Socket } from '../../Socket'; + +export type Responder = (request: SentRequest) => string | string[] | null; + +export interface SentRequest { + method: string; + raw: string; + header(name: string): string | undefined; + body: string; +} + +export default class ResponderSocket implements Socket { + url = 'ws://localhost:12345'; + via_transport = 'WS'; + sip_uri = 'sip:localhost:12345;transport=ws'; + + // Every message sent by the UA, requests and responses. + readonly sent: string[] = []; + + // The requests among them. + readonly requests: SentRequest[] = []; + + constructor(private readonly responder: Responder) {} + + connect(): void { + setTimeout(() => { + this.onconnect(); + }, 0); + } + + disconnect(): void { + setTimeout(() => { + this.ondisconnect(); + }, 0); + } + + send(message: string): boolean { + this.sent.push(message); + + // Only requests are answered. + if (message.startsWith('SIP/2.0')) { + return true; + } + + const request = parse(message); + + this.requests.push(request); + + const answer = this.responder(request); + const messages = answer === null ? [] : ([] as string[]).concat(answer); + + for (const data of messages) { + this.inject(data); + } + + return true; + } + + // Deliver a message to the UA, as if the server had sent it. + inject(message: string): void { + setTimeout(() => { + this.ondata(message); + }, 0); + } + + isConnected(): boolean { + return true; + } + + isConnecting(): boolean { + return false; + } + + onconnect(): void {} + + ondisconnect(): void {} + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + ondata(_event: T): void {} +} + +function parse(message: string): SentRequest { + const [head, body = ''] = message.split('\r\n\r\n'); + const lines = head.split('\r\n'); + + return { + method: lines[0].split(' ')[0], + raw: message, + header(name: string): string | undefined { + const prefix = `${name.toLowerCase()}:`; + const line = lines.find(l => l.toLowerCase().startsWith(prefix)); + + return line?.substring(prefix.length).trim(); + }, + body, + }; +} + +/** + * Build the response to `request`: Via, From, To (tagged), Call-ID and CSeq + * copied from it, then `headers`. + */ +export function reply( + request: SentRequest, + status: number, + reason: string, + headers: string[] = [] +): string { + const lines = request.raw.split('\r\n\r\n')[0].split('\r\n'); + const copied = lines.filter(line => /^(via|from|call-id|cseq):/i.test(line)); + let to = lines.find(line => /^to:/i.test(line))!; + + if (!/;tag=/i.test(to)) { + to += ';tag=responder'; + } + + return [ + `SIP/2.0 ${status} ${reason}`, + ...copied, + to, + ...headers, + 'Content-Length: 0', + '', + '', + ].join('\r\n'); +} + +/** + * Build the final NOTIFY of the subscription `subscribe` (any request of + * its dialog), as the notifier of `reply` would send it. + */ +export function finalNotify(subscribe: SentRequest, event: string): string { + const contact = subscribe.header('Contact')!.replace(/^<([^>]*)>.*$/, '$1'); + const from = subscribe.header('From')!; + let to = subscribe.header('To')!; + + if (!/;tag=/i.test(to)) { + to += ';tag=responder'; + } + + return [ + `NOTIFY ${contact} SIP/2.0`, + 'Via: SIP/2.0/WS example.com;branch=z9hG4bKresponder', + 'Max-Forwards: 70', + `From: ${to}`, + `To: ${from}`, + `Call-ID: ${subscribe.header('Call-ID')}`, + 'CSeq: 1 NOTIFY', + `Event: ${event}`, + 'Subscription-State: terminated', + 'Contact: ', + 'Content-Length: 0', + '', + '', + ].join('\r\n'); +} diff --git a/src/test/test-Publisher.ts b/src/test/test-Publisher.ts new file mode 100644 index 000000000..73e8b7ea9 --- /dev/null +++ b/src/test/test-Publisher.ts @@ -0,0 +1,369 @@ +import './include/common'; +import ResponderSocket, { reply } from './include/ResponderSocket'; +import type { SentRequest } from './include/ResponderSocket'; + +// eslint-disable-next-line @typescript-eslint/no-require-imports +const JsSIP = require('../JsSIP.js'); +const { UA } = JsSIP; + +const TARGET = 'sip:alice@example.com'; +const CONTENT_TYPE = 'application/pidf+xml'; + +interface TestUA { + on(event: string, listener: (...args: never[]) => void): void; + stop(): void; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + publish(...args: any[]): any; +} + +/** + * A presence server that keeps one entity-tag. Each PUBLISH gets the + * response returned by the next function in `script`, if any, else the + * normal one: 200 with a new SIP-ETag, or 412 when SIP-If-Match is not the + * current entity-tag. It grants `granted` seconds, or what was asked. + */ +function startUA( + script: ((r: SentRequest) => string)[] = [], + granted?: number +): { + ua: TestUA; + socket: ResponderSocket; +} { + let etag: string | null = null; + let n = 0; + + const socket = new ResponderSocket(request => { + if (request.method !== 'PUBLISH') { + return null; + } + + const scripted = script.shift(); + + if (scripted) { + return scripted(request); + } + + const ifMatch = request.header('SIP-If-Match'); + + if (ifMatch !== undefined && ifMatch !== etag) { + return reply(request, 412, 'Conditional Request Failed'); + } + + const expires = request.header('Expires'); + + if (expires === '0') { + etag = null; + + return reply(request, 200, 'OK', ['Expires: 0']); + } + + etag = `e${++n}`; + + return reply(request, 200, 'OK', [ + `SIP-ETag: ${etag}`, + `Expires: ${granted ?? expires}`, + ]); + }); + const ua = new UA({ + sockets: socket, + uri: TARGET, + register: false, + }); + + ua.start(); + + return { ua, socket }; +} + +function connected(ua: TestUA): Promise { + return new Promise(resolve => ua.on('connected', () => resolve())); +} + +function stopAndWait(ua: TestUA): Promise { + return new Promise(resolve => { + ua.on('disconnected', () => resolve()); + ua.stop(); + }); +} + +// eslint-disable-next-line @typescript-eslint/no-explicit-any +function next(publisher: any, event: string): Promise { + return new Promise(resolve => + publisher.once(event, (...args: unknown[]) => resolve(args)) + ); +} + +/** What each PUBLISH carried: Expires, SIP-If-Match, body. */ +function sent(socket: ResponderSocket): (string | undefined)[][] { + return socket.requests + .filter(r => r.method === 'PUBLISH') + .map(r => [r.header('Expires'), r.header('SIP-If-Match'), r.body]); +} + +describe('Publisher', () => { + test('publishes, modifies, and removes', async () => { + const { ua, socket } = startUA(); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + + expect(publisher.etag).toBe('e1'); + expect(publisher.state).toBe(publisher.C.STATE_PUBLISHED); + + const [initial] = socket.requests; + + expect(initial.header('Event')).toBe('presence'); + expect(initial.header('Content-Type')).toBe(CONTENT_TYPE); + + publisher.publish('closed'); + await next(publisher, 'published'); + + expect(socket.requests[1].header('CSeq')).toBe( + `${Number(initial.header('CSeq')?.split(' ')[0]) + 1} PUBLISH` + ); + expect(socket.requests[1].header('Call-ID')).toBe( + initial.header('Call-ID') + ); + + publisher.terminate(); + + const [code] = await next(publisher, 'terminated'); + + expect(code).toBe(publisher.C.UNPUBLISHED); + expect(sent(socket)).toEqual([ + ['3600', undefined, 'open'], + ['3600', 'e1', 'closed'], + ['0', 'e2', ''], + ]); + + await stopAndWait(ua); + }); + + test('keeps one PUBLISH in flight, and sends the latest body', async () => { + const { ua, socket } = startUA(); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('one'); + publisher.publish('two'); + publisher.publish('three'); + await next(publisher, 'published'); + await next(publisher, 'published'); + + expect(sent(socket)).toEqual([ + ['3600', undefined, 'one'], + ['3600', 'e1', 'three'], + ]); + + await stopAndWait(ua); + }); + + test('refreshes without a body before the publication expires', async () => { + // The server grants 6 s: the refresh leaves 1 s later. + const { ua, socket } = startUA([], 6); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + await next(publisher, 'published'); + + expect(sent(socket)).toEqual([ + ['3600', undefined, 'open'], + ['3600', 'e1', ''], + ]); + expect(socket.requests[1].header('Content-Type')).toBeUndefined(); + + await stopAndWait(ua); + }); + + test('publishes from scratch once after a 412', async () => { + const { ua, socket } = startUA([ + r => reply(r, 200, 'OK', ['SIP-ETag: lost', 'Expires: 3600']), + ]); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + + // The server does not know "lost": 412, then a new publication. + publisher.publish('closed'); + await next(publisher, 'published'); + + expect(sent(socket)).toEqual([ + ['3600', undefined, 'open'], + ['3600', 'lost', 'closed'], + ['3600', undefined, 'closed'], + ]); + expect(publisher.etag).toBe('e1'); + + await stopAndWait(ua); + }); + + test('terminates on a second 412 in a row', async () => { + const { ua } = startUA([ + r => reply(r, 200, 'OK', ['SIP-ETag: lost', 'Expires: 3600']), + r => reply(r, 412, 'Conditional Request Failed'), + r => reply(r, 412, 'Conditional Request Failed'), + ]); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + publisher.publish('closed'); + + const [code, response] = await next(publisher, 'terminated'); + + expect(code).toBe(publisher.C.PUBLISH_NON_OK_RESPONSE); + expect(response.status_code).toBe(412); + + await stopAndWait(ua); + }); + + test('publishes again with the Min-Expires of a 423', async () => { + const { ua, socket } = startUA([ + r => reply(r, 423, 'Interval Too Brief', ['Min-Expires: 7200']), + ]); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + + expect(sent(socket)).toEqual([ + ['3600', undefined, 'open'], + ['7200', undefined, 'open'], + ]); + + await stopAndWait(ua); + }); + + test('passes the response of a refusal with "terminated"', async () => { + const { ua } = startUA([r => reply(r, 489, 'Bad Event')]); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + + const [code, response] = await next(publisher, 'terminated'); + + expect(code).toBe(publisher.C.PUBLISH_NON_OK_RESPONSE); + expect(response.status_code).toBe(489); + expect(() => publisher.publish('closed')).toThrow( + JsSIP.Exceptions.InvalidStateError + ); + + await stopAndWait(ua); + }); + + test('terminates on a 2xx without SIP-ETag', async () => { + const { ua } = startUA([r => reply(r, 200, 'OK', ['Expires: 3600'])]); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + + const [code] = await next(publisher, 'terminated'); + + expect(code).toBe(publisher.C.PUBLISH_WRONG_OK_RESPONSE); + + await stopAndWait(ua); + }); + + test('terminate() before any publication sends nothing', async () => { + const { ua, socket } = startUA(); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + const ended = next(publisher, 'terminated'); + + publisher.terminate(); + + const [code] = await ended; + + expect(code).toBe(publisher.C.UNPUBLISHED); + expect(socket.requests).toHaveLength(0); + + await stopAndWait(ua); + }); + + test('terminate() during a modification removes after it', async () => { + const { ua, socket } = startUA(); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + publisher.publish('closed'); + publisher.terminate(); + + const [code] = await next(publisher, 'terminated'); + + expect(code).toBe(publisher.C.UNPUBLISHED); + expect(sent(socket)).toEqual([ + ['3600', undefined, 'open'], + ['3600', 'e1', 'closed'], + ['0', 'e2', ''], + ]); + + await stopAndWait(ua); + }); + + test('UA.stop() removes the publication', async () => { + const { ua, socket } = startUA(); + + await connected(ua); + + const publisher = ua.publish(TARGET, 'presence', CONTENT_TYPE); + + publisher.publish('open'); + await next(publisher, 'published'); + + const ended = next(publisher, 'terminated'); + + await stopAndWait(ua); + + const [code] = await ended; + + expect(code).toBe(publisher.C.UNPUBLISHED); + expect(sent(socket)).toEqual([ + ['3600', undefined, 'open'], + ['0', 'e1', ''], + ]); + }); + + test('refuses an invalid target', async () => { + const { ua } = startUA(); + + await connected(ua); + + expect(() => ua.publish('sip:', 'presence', CONTENT_TYPE)).toThrow( + TypeError + ); + + await stopAndWait(ua); + }); +});