diff --git a/README.md b/README.md index b81632d1..695fb0e4 100644 --- a/README.md +++ b/README.md @@ -173,7 +173,7 @@ export default class OrdersController extends Controller { | `table-context` | Selection state for the active table | | `abilities` | Permission checks | | `notifications` | Toast notifications | -| `socket` | Real-time channels | +| `socket` | Real-time channels, authenticated with session socket tokens | | `chat` | Chat channels and messages | | `events` | Application event tracking | | `theme` | Light and dark themes and route body classes | diff --git a/addon/services/socket.js b/addon/services/socket.js index bd379244..1d210324 100644 --- a/addon/services/socket.js +++ b/addon/services/socket.js @@ -1,16 +1,153 @@ -import Service from '@ember/service'; +import Service, { inject as service } from '@ember/service'; import { tracked } from '@glimmer/tracking'; import { isBlank } from '@ember/utils'; import { later } from '@ember/runloop'; +import { debug } from '@ember/debug'; import toBoolean from '../utils/to-boolean'; import config from 'ember-get-config'; +/** + * Console mint route for socket tokens, always on the internal namespace + * regardless of what the fetch service's namespace has been set to. + */ +export const TOKEN_PATH = 'socket/token'; +export const TOKEN_NAMESPACE = 'int/v1'; + +/** + * A cached token is reused only while it has more than this many seconds left, + * and the refresh timer fires this many seconds before expiry. + */ +export const REFRESH_LEEWAY_SECONDS = 60; + +/** + * Floor for the refresh timer, so a server handing out very short-lived tokens + * can never make the client spin. + */ +export const MIN_REFRESH_DELAY_MS = 5000; + +/** + * At most this many recoveries (fetch a token, authenticate, resubscribe) are + * attempted within the window; beyond that the service stops trying until the + * window has passed, so a server that keeps rejecting tokens cannot cause a loop. + */ +export const MAX_RECOVERIES_PER_WINDOW = 3; +export const RECOVERY_WINDOW_MS = 60000; + +/** + * Reasons the socket server gives (as a kickOut message or an AuthError's + * `reason`) when a channel was lost because of the token rather than because + * the user may not see the channel. Only these are retried with a fresh token. + */ +export const RETRYABLE_REASONS = ['no_token', 'token_expired', 'token_changed', 'deauthenticated', 'identity_changed']; + +/** + * Client events the service reacts to. + */ +export const CLIENT_EVENTS = ['connect', 'deauthenticate', 'kickOut', 'subscribeFail']; + +/** + * Universe events (fired by the session and current-user services) the service + * reacts to, mapped to the handler that runs. + */ +export const SESSION_EVENTS = { + 'session.authenticated': 'handleLogin', + 'user.loaded': 'handleUserLoaded', + 'user.organization_switched': 'handleOrganizationSwitch', + 'user.deauthenticated': 'handleLogout', +}; + +/** + * The `client` tag sent in the handshake query, used by the socket server to + * attribute authorization decisions to an app and version. Not a secret. + * + * @param {Object} appConfig + * @return {String} + */ +export function clientTag(appConfig) { + return `console/${appConfig.version ?? 'unknown'}`; +} + +/** + * The server-side reason a subscription was refused, or null when the failure + * was not an authorization failure. + * + * @param {Error} error + * @return {String|null} + */ +export function subscribeFailReason(error) { + return error.name === 'AuthError' ? error.reason : null; +} + +/** + * Validates a mint-route response. Returns null for anything that is not a + * token with a positive lifetime, which the caller treats as "no token". + * + * @param {Object|null} response `{ token, expires_in, expires_at }` + * @return {{token: String, expiresIn: Number}|null} + */ +export function parseTokenResponse(response) { + if (!response || typeof response.token !== 'string') { + return null; + } + + const expiresIn = Number(response.expires_in); + return expiresIn > 0 ? { token: response.token, expiresIn } : null; +} + +/** + * SocketService + * + * Owns the console's single SocketCluster client and keeps it authenticated + * with short-lived socket tokens minted from the console session. + * + * - The token is delivered in the handshake through an in-memory auth engine, + * so every connect and reconnect (and the subscriptions queued before it) is + * authenticated. Tokens are never written to localStorage or the URL. + * - When there is no session, or the server cannot mint tokens (older server or + * socket auth disabled), the client connects anonymously as it always has. + * - Tokens are refreshed shortly before they expire; when the server drops the + * token or kicks the client out of channels for a token reason, a fresh token + * is fetched and the lost channels are resubscribed by name, which also + * restores delivery to callers iterating channels from `instance().subscribe`. + * - Login, organization switches and logout re-key or tear down the socket. + */ export default class SocketService extends Service { + @service session; + @service fetch; + @service universe; + @tracked channels = []; + /** @type {{token: String, expiresAt: Number}|null} */ + tokenState = null; + tokenRequest = null; + refreshTimer = null; + recovery = null; + authenticating = false; + generation = 0; + recoveryLog = []; + pendingChannels = new Set(); + retriedChannels = new Map(); + constructor() { super(...arguments); + this.authEngine = this.createAuthEngine(); this.socket = this.createSocketClusterClient(); + this.listenToClient(); + this.listenToSession(); + } + + willDestroy() { + super.willDestroy(...arguments); + this.cancelRefresh(); + + for (const eventName of CLIENT_EVENTS) { + this.socket.closeListener(eventName); + } + + for (const [eventName, handler] of Object.entries(SESSION_EVENTS)) { + this.sessionBus.off(eventName, this, this[handler]); + } } instance() { @@ -25,10 +162,411 @@ export default class SocketService extends Service { } socketConfig.secure = toBoolean(socketConfig.secure); + socketConfig.query = { ...socketConfig.query, client: clientTag(config) }; + socketConfig.authEngine = this.authEngine; return socketClusterClient.create(socketConfig); } + /** + * The auth engine handed to the SocketCluster client. It replaces the default + * engine, which would persist the token in localStorage. + * + * @return {Object} + */ + createAuthEngine() { + return { + // The only tokens the client is ever given are the ones this service + // fetched and already holds in memory, so there is nothing to persist. + saveToken: (name, token) => Promise.resolve(token), + removeToken: () => { + const previous = this.tokenState ? this.tokenState.token : null; + this.forgetToken(); + return Promise.resolve(previous); + }, + loadToken: () => this.loadToken(), + }; + } + + /** + * Resolves the token for a handshake: a fresh socket token while the session + * is authenticated, otherwise null so the client connects anonymously. + * Never rejects. + * + * @return {Promise} + */ + async loadToken() { + if (!this.session.isAuthenticated) { + return null; + } + + return this.getToken(); + } + + /** + * Returns the cached token while it has more than the leeway left, otherwise + * fetches one. Concurrent callers share a single request. + * + * @param {Object} [options] + * @param {Boolean} [options.force=false] ignore the cached token + * @return {Promise} + */ + getToken({ force = false } = {}) { + if (!force && this.hasFreshToken()) { + return Promise.resolve(this.tokenState.token); + } + + if (!this.tokenRequest) { + const request = this.requestToken().finally(() => { + if (this.tokenRequest === request) { + this.tokenRequest = null; + } + }); + this.tokenRequest = request; + } + + return this.tokenRequest; + } + + hasFreshToken() { + return this.tokenState !== null && this.tokenState.expiresAt - Date.now() > REFRESH_LEEWAY_SECONDS * 1000; + } + + /** + * Mints a token from the console session. Any failure, including the 404 an + * older server or one with socket auth disabled returns, resolves to null so + * the socket falls back to an anonymous connection. + * + * @return {Promise} + */ + async requestToken() { + const generation = this.generation; + let response; + + try { + response = await this.fetch.post(TOKEN_PATH, {}, { namespace: TOKEN_NAMESPACE }); + } catch (error) { + debug(`[socket] socket token unavailable, connecting anonymously: ${error.message}`); + return null; + } + + // The session ended while the request was in flight. + if (generation !== this.generation) { + return null; + } + + const parsed = parseTokenResponse(response); + if (parsed === null) { + debug('[socket] socket token response was malformed, connecting anonymously'); + return null; + } + + this.tokenState = { token: parsed.token, expiresAt: Date.now() + parsed.expiresIn * 1000 }; + this.scheduleRefresh(parsed.expiresIn); + + return parsed.token; + } + + /** + * Arms the refresh timer for `expiresIn - leeway` seconds from now. A native + * timer is used deliberately: a run-loop timer would hold the test waiters + * (and `settled()`) open for the token's whole lifetime. + * + * @param {Number} expiresIn seconds until the token expires + */ + scheduleRefresh(expiresIn) { + this.cancelRefresh(); + + const delay = Math.max((expiresIn - REFRESH_LEEWAY_SECONDS) * 1000, MIN_REFRESH_DELAY_MS); + this.refreshTimer = setTimeout(() => { + this.refreshTimer = null; + this.refreshToken(); + }, delay); + } + + cancelRefresh() { + clearTimeout(this.refreshTimer); + this.refreshTimer = null; + } + + forgetToken() { + this.tokenState = null; + this.cancelRefresh(); + } + + /** + * Fetches a new token and, when connected, authenticates the open socket with + * it. When the socket is not connected the new token simply waits in memory + * for the next handshake. + * + * @return {Promise} + */ + async refreshToken() { + if (!this.session.isAuthenticated) { + return null; + } + + const token = await this.getToken({ force: true }); + if (token !== null && this.isOpen()) { + await this.authenticateWith(token); + } + + return token; + } + + /** + * Re-keys the socket for the current session: used on login and when the + * user switches organization (the old token names the old organization). + * + * @param {Object} [options] + * @param {Boolean} [options.force=true] discard the cached token first + * @return {Promise} whether the socket is now authenticated + */ + async reauthenticate({ force = true } = {}) { + if (!this.session.isAuthenticated) { + return false; + } + + if (force) { + this.forgetToken(); + } + + // Not connected (logged out earlier, or between reconnect attempts): the + // handshake loads the token itself. + if (!this.isOpen()) { + this.socket.connect(); + return false; + } + + const token = await this.getToken(); + return token !== null && this.authenticateWith(token); + } + + /** + * @param {String} token + * @return {Promise} whether the server accepted the token + */ + async authenticateWith(token) { + this.authenticating = true; + + try { + const status = await this.socket.authenticate(token); + if (status.isAuthenticated) { + return true; + } + + debug('[socket] socket token was rejected'); + } catch (error) { + debug(`[socket] socket authentication failed: ${error.message}`); + } finally { + this.authenticating = false; + } + + // A rejected token must not be offered again on the next handshake. + this.forgetToken(); + return false; + } + + isOpen() { + return this.socket.state === this.socket.OPEN; + } + + isAuthenticatedWith(token) { + return this.socket.authState === this.socket.AUTHENTICATED && this.socket.signedAuthToken === token; + } + + // ------------------------------------------------------------------------- + // Client events + // ------------------------------------------------------------------------- + + listenToClient() { + this.consume('connect', () => this.handleConnect()); + this.consume('deauthenticate', () => this.handleDeauthenticate()); + this.consume('kickOut', (event) => this.handleChannelLoss(event.channel, event.message)); + this.consume('subscribeFail', (event) => this.handleChannelLoss(event.channel, subscribeFailReason(event.error))); + } + + consume(eventName, handler) { + const stream = this.socket.listener(eventName); + + (async () => { + for await (const event of stream) { + handler(event); + } + })(); + } + + /** + * A handshake completed anonymously although there is a session: the socket + * connected before the session was restored, or the token request failed. + * Try once more now, within the recovery budget. + */ + handleConnect() { + if (this.socket.authState === this.socket.AUTHENTICATED || !this.session.isAuthenticated) { + return null; + } + + return this.recover(); + } + + /** + * The server dropped the socket's token (it expired, or the server removed + * it). Deauthentication caused by this service's own authenticate call is + * ignored: that path already handles the rejected token. + */ + handleDeauthenticate() { + if (this.authenticating || !this.session.isAuthenticated) { + return null; + } + + this.forgetToken(); + return this.recover(); + } + + /** + * A channel was kicked out or refused. Token-related losses are queued for + * resubscription after re-authenticating; anything else (the user may not see + * the channel) is left for the subscriber's own `subscribeFail` handling. + * + * @param {String} channelName + * @param {String|null} reason + */ + handleChannelLoss(channelName, reason) { + if (!RETRYABLE_REASONS.includes(reason) || !this.session.isAuthenticated) { + return null; + } + + this.pendingChannels.add(channelName); + return this.recover(); + } + + /** + * Single-flight recovery. Losses reported while a recovery is running are + * picked up by it, or by a follow-up run if they arrive after it resubscribed. + * + * @return {Promise} + */ + recover() { + if (!this.recovery) { + this.recovery = this.runRecovery().finally(() => { + this.recovery = null; + + if (this.pendingChannels.size > 0) { + this.recover(); + } + }); + } + + return this.recovery; + } + + async runRecovery() { + if (!this.takeRecoveryBudget()) { + debug('[socket] too many socket re-authentications, waiting before trying again'); + this.pendingChannels.clear(); + return false; + } + + const token = await this.getToken(); + if (token === null) { + this.pendingChannels.clear(); + return false; + } + + if (!this.isAuthenticatedWith(token) && !(await this.authenticateWith(token))) { + this.pendingChannels.clear(); + return false; + } + + this.resubscribePending(token); + return true; + } + + takeRecoveryBudget() { + const now = Date.now(); + this.recoveryLog = this.recoveryLog.filter((at) => now - at < RECOVERY_WINDOW_MS); + + if (this.recoveryLog.length >= MAX_RECOVERIES_PER_WINDOW) { + return false; + } + + this.recoveryLog.push(now); + return true; + } + + /** + * Resubscribes lost channels by name. Each channel is retried at most once + * per token, so a channel the server keeps refusing is not retried forever. + * + * @param {String} token the token the socket is now authenticated with + */ + resubscribePending(token) { + const channelNames = [...this.pendingChannels]; + this.pendingChannels.clear(); + + for (const channelName of channelNames) { + if (this.retriedChannels.get(channelName) === token) { + debug(`[socket] not retrying channel ${channelName} again with the same token`); + continue; + } + + this.retriedChannels.set(channelName, token); + this.socket.subscribe(channelName); + } + } + + // ------------------------------------------------------------------------- + // Session events + // ------------------------------------------------------------------------- + + listenToSession() { + // Held directly so teardown does not have to look the service up again + // while the owner is being destroyed. + this.sessionBus = this.universe; + + for (const [eventName, handler] of Object.entries(SESSION_EVENTS)) { + this.sessionBus.on(eventName, this, this[handler]); + } + } + + handleLogin() { + return this.reauthenticate({ force: true }); + } + + /** + * A restored session loads the user without a login event; if the socket + * connected anonymously before the session was restored, authenticate it now. + */ + handleUserLoaded() { + if (this.socket.authState === this.socket.AUTHENTICATED) { + return null; + } + + return this.reauthenticate({ force: false }); + } + + handleOrganizationSwitch() { + return this.reauthenticate({ force: true }); + } + + /** + * Logout: drop the token and every pending retry, and close the connection so + * nothing keeps flowing to a signed-out console. The next login reconnects. + */ + handleLogout() { + this.generation++; + this.tokenRequest = null; + this.forgetToken(); + this.pendingChannels.clear(); + this.retriedChannels.clear(); + this.recoveryLog = []; + this.socket.disconnect(); + } + + // ------------------------------------------------------------------------- + // Channel helpers + // ------------------------------------------------------------------------- + async listen(channelId, callback) { later( this, diff --git a/tests/helpers/stub-socketcluster.js b/tests/helpers/stub-socketcluster.js index 174d4423..e11f88b5 100644 --- a/tests/helpers/stub-socketcluster.js +++ b/tests/helpers/stub-socketcluster.js @@ -12,8 +12,8 @@ * so no production code has to change; and * 2. the global itself, which is replaced with an inert fake. * - * Tests that exercise socket behaviour should register their own fake on the owner - * rather than relying on the shape of this one. + * Tests that exercise socket behaviour swap in their own fake (see + * createRecordingClient) rather than relying on the shape of the inert one. */ const MARKER_SELECTOR = 'script[data-socketcluster-client]'; @@ -33,23 +33,148 @@ function createFakeChannel(name) { }; } +/** + * A client-level event listener that is both awaitable once and async-iterable; + * the iterator finishes immediately so the socket service's event loops end. + */ +function inertListener() { + return { + once: () => Promise.resolve(), + [Symbol.asyncIterator]() { + return { next: () => Promise.resolve({ done: true, value: undefined }) }; + }, + }; +} + +/** + * The client methods the socket service calls on every instance, as no-ops. + * Test fakes spread this in so they only spell out what they observe. + */ +export function inertClientMethods() { + return { + OPEN: 'open', + CLOSED: 'closed', + AUTHENTICATED: 'authenticated', + UNAUTHENTICATED: 'unauthenticated', + state: 'closed', + authState: 'unauthenticated', + signedAuthToken: null, + listener: inertListener, + closeListener() {}, + authenticate: () => Promise.resolve({ isAuthenticated: false }), + connect() {}, + disconnect() {}, + }; +} + export function createFakeSocketClusterClient() { return { create() { return { + ...inertClientMethods(), subscribe: (channelId) => createFakeChannel(channelId), transmit() {}, invoke: () => Promise.resolve(), closeAllChannels() {}, - disconnect() {}, - listener() { - return { once: () => Promise.resolve() }; + }; + }, + }; +} + +/** + * An async-iterable event stream that stays open until closed, so a test can + * push client events (deauthenticate, kickOut, subscribeFail, connect) into the + * socket service's listeners. + */ +export function createEventStream() { + const queue = []; + const waiting = []; + let closed = false; + + return { + push(value) { + if (waiting.length) { + waiting.shift()({ done: false, value }); + } else { + queue.push(value); + } + }, + close() { + closed = true; + while (waiting.length) { + waiting.shift()({ done: true, value: undefined }); + } + }, + once: () => Promise.resolve(), + [Symbol.asyncIterator]() { + return { + next() { + if (queue.length) { + return Promise.resolve({ done: false, value: queue.shift() }); + } + if (closed) { + return Promise.resolve({ done: true, value: undefined }); + } + return new Promise((resolve) => waiting.push(resolve)); }, }; }, }; } +/** + * A SocketCluster client fake that records what the socket service does with it + * and lets a test emit client events. `authenticate` accepts every token unless + * `client.authenticateImpl` is replaced. + */ +export function createRecordingClient() { + const streams = {}; + const client = { + ...inertClientMethods(), + state: 'open', + subscribed: [], + authenticated: [], + closedListeners: [], + connects: 0, + disconnects: 0, + listener(eventName) { + if (!streams[eventName]) { + streams[eventName] = createEventStream(); + } + return streams[eventName]; + }, + closeListener(eventName) { + client.closedListeners.push(eventName); + client.listener(eventName).close(); + }, + emit(eventName, data) { + client.listener(eventName).push(data); + }, + subscribe(channelName) { + client.subscribed.push(channelName); + return createFakeChannel(channelName); + }, + authenticateImpl(token) { + client.authState = client.AUTHENTICATED; + client.signedAuthToken = token; + return Promise.resolve({ isAuthenticated: true, authError: null }); + }, + authenticate(token) { + client.authenticated.push(token); + return client.authenticateImpl(token); + }, + connect() { + client.connects++; + }, + disconnect() { + client.disconnects++; + client.state = client.CLOSED; + }, + }; + + return client; +} + export default function stubSocketCluster() { if (!document.querySelector(MARKER_SELECTOR)) { const marker = document.createElement('script'); diff --git a/tests/unit/services/socket-auth-test.js b/tests/unit/services/socket-auth-test.js new file mode 100644 index 00000000..74df375f --- /dev/null +++ b/tests/unit/services/socket-auth-test.js @@ -0,0 +1,627 @@ +import { module, test } from 'qunit'; +import { setupTest } from 'dummy/tests/helpers'; +import Service from '@ember/service'; +import Evented from '@ember/object/evented'; +import { run } from '@ember/runloop'; +import { settled } from '@ember/test-helpers'; +import config from 'dummy/config/environment'; +import { createRecordingClient } from 'dummy/tests/helpers/stub-socketcluster'; +import { + TOKEN_PATH, + TOKEN_NAMESPACE, + REFRESH_LEEWAY_SECONDS, + MIN_REFRESH_DELAY_MS, + MAX_RECOVERIES_PER_WINDOW, + RECOVERY_WINDOW_MS, + CLIENT_EVENTS, + clientTag, + subscribeFailReason, + parseTokenResponse, +} from '@fleetbase/ember-core/services/socket'; + +/** + * Socket authentication: the in-memory auth engine, token refresh, recovery + * after the server drops the token or kicks the client out of channels, and the + * session lifecycle (login, restore, organization switch, logout). + * + * The SocketCluster global is swapped for a recording fake (see + * tests/helpers/stub-socketcluster), and the session, fetch and universe + * services are replaced with small stubs so each test controls whether there is + * a session and what the mint route answers. + */ + +// Lets the service's async event loops and promise chains run. +async function flush() { + for (let i = 0; i < 10; i++) { + await new Promise((resolve) => setTimeout(resolve, 0)); + } +} + +function authError(reason) { + const error = new Error(`Subscription to channel denied: ${reason}`); + error.name = 'AuthError'; + error.reason = reason; + return error; +} + +function deferred() { + let resolve; + let reject; + const promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} + +class SessionStub extends Service { + isAuthenticated = true; +} + +// Records mint requests; by default every call mints a new token valid for 15 minutes. +class FetchStub extends Service { + calls = []; + responder = (call) => Promise.resolve({ token: `token-${call}`, expires_in: 900, expires_at: '2026-01-01T00:15:00Z' }); + + post(path, data, options) { + this.calls.push({ path, data, options }); + return this.responder(this.calls.length); + } +} + +class UniverseStub extends Service.extend(Evented) {} + +module('Unit | Service | socket (authentication)', function (hooks) { + setupTest(hooks); + + hooks.beforeEach(function () { + const testContext = this; + this.created = []; + + this.originalClient = window.socketClusterClient; + window.socketClusterClient = { + create(socketConfig) { + const client = createRecordingClient(); + testContext.created.push(socketConfig); + testContext.client = client; + return client; + }, + }; + + this.owner.register('service:session', SessionStub); + this.owner.register('service:fetch', FetchStub); + this.owner.register('service:universe', UniverseStub); + + this.session = this.owner.lookup('service:session'); + this.fetch = this.owner.lookup('service:fetch'); + this.universe = this.owner.lookup('service:universe'); + this.service = this.owner.lookup('service:socket'); + this.authEngine = this.created[0].authEngine; + }); + + hooks.afterEach(function () { + window.socketClusterClient = this.originalClient; + }); + + // ------------------------------------------------------------------------- + // Pure helpers + // ------------------------------------------------------------------------- + + test('clientTag names the console and its version', function (assert) { + assert.strictEqual(clientTag({ version: '0.7.69' }), 'console/0.7.69'); + assert.strictEqual(clientTag({}), 'console/unknown'); + }); + + test('subscribeFailReason reads the reason off authorization errors only', function (assert) { + assert.strictEqual(subscribeFailReason(authError('no_token')), 'no_token'); + + const other = new Error('Socket hung up'); + other.name = 'BadConnectionError'; + other.reason = 'no_token'; + assert.strictEqual(subscribeFailReason(other), null); + }); + + test('parseTokenResponse accepts only a token with a positive lifetime', function (assert) { + assert.strictEqual(parseTokenResponse(null), null); + assert.strictEqual(parseTokenResponse({ expires_in: 900 }), null); + assert.strictEqual(parseTokenResponse({ token: 'abc', expires_in: 0 }), null); + assert.strictEqual(parseTokenResponse({ token: 'abc', expires_in: 'soon' }), null); + assert.deepEqual(parseTokenResponse({ token: 'abc', expires_in: '900' }), { token: 'abc', expiresIn: 900 }); + }); + + // ------------------------------------------------------------------------- + // Client construction + // ------------------------------------------------------------------------- + + test('the client gets the in-memory auth engine and a client tag in the query', function (assert) { + const socketConfig = this.created[0]; + + assert.strictEqual(socketConfig.authEngine, this.service.authEngine); + assert.strictEqual(typeof socketConfig.authEngine.loadToken, 'function'); + assert.strictEqual(socketConfig.query.client, clientTag(config)); + assert.notOk('token' in socketConfig.query, 'the token never travels in the URL'); + }); + + test('an existing query in the socket config is kept', function (assert) { + const original = config.socket.query; + config.socket.query = { region: 'a' }; + + try { + const service = this.owner.factoryFor('service:socket').create(); + const socketConfig = this.created[this.created.length - 1]; + + assert.deepEqual(socketConfig.query, { region: 'a', client: clientTag(config) }); + run(() => service.destroy()); + } finally { + config.socket.query = original; + } + }); + + // ------------------------------------------------------------------------- + // Auth engine + // ------------------------------------------------------------------------- + + test('loadToken connects anonymously without a session', async function (assert) { + this.session.isAuthenticated = false; + + assert.strictEqual(await this.authEngine.loadToken('socketcluster.authToken'), null); + assert.strictEqual(this.fetch.calls.length, 0, 'no token is requested'); + }); + + test('loadToken mints a token from the console session route', async function (assert) { + const token = await this.authEngine.loadToken('socketcluster.authToken'); + + assert.strictEqual(token, 'token-1'); + assert.strictEqual(this.fetch.calls.length, 1); + assert.strictEqual(this.fetch.calls[0].path, TOKEN_PATH); + assert.strictEqual(this.fetch.calls[0].options.namespace, TOKEN_NAMESPACE); + assert.strictEqual(TOKEN_NAMESPACE, 'int/v1'); + }); + + test('a cached token is reused while it is comfortably valid', async function (assert) { + await this.authEngine.loadToken(); + const second = await this.authEngine.loadToken(); + + assert.strictEqual(second, 'token-1'); + assert.strictEqual(this.fetch.calls.length, 1, 'only one request was made'); + }); + + test('a token inside the refresh leeway is replaced', async function (assert) { + this.fetch.responder = (call) => Promise.resolve({ token: `token-${call}`, expires_in: REFRESH_LEEWAY_SECONDS - 1 }); + + assert.strictEqual(await this.authEngine.loadToken(), 'token-1'); + assert.strictEqual(await this.authEngine.loadToken(), 'token-2'); + }); + + test('concurrent loads share one request', async function (assert) { + const pending = deferred(); + this.fetch.responder = () => pending.promise; + + const first = this.authEngine.loadToken(); + const second = this.authEngine.loadToken(); + pending.resolve({ token: 'shared', expires_in: 900 }); + + assert.deepEqual(await Promise.all([first, second]), ['shared', 'shared']); + assert.strictEqual(this.fetch.calls.length, 1); + assert.strictEqual(this.service.tokenRequest, null, 'the finished request is released'); + }); + + test('a failed mint request (e.g. 404 from an older server) falls back to anonymous', async function (assert) { + this.fetch.responder = () => Promise.reject(new Error('Not Found')); + + assert.strictEqual(await this.authEngine.loadToken(), null); + assert.strictEqual(this.service.tokenState, null); + assert.strictEqual(this.service.refreshTimer, null, 'nothing is scheduled'); + }); + + test('a malformed mint response falls back to anonymous', async function (assert) { + this.fetch.responder = () => Promise.resolve({ token: 'abc' }); + + assert.strictEqual(await this.authEngine.loadToken(), null); + assert.strictEqual(this.service.tokenState, null); + }); + + test('saveToken keeps nothing outside memory', async function (assert) { + assert.strictEqual(await this.authEngine.saveToken('socketcluster.authToken', 'abc', {}), 'abc'); + assert.strictEqual(window.localStorage.getItem('socketcluster.authToken'), null); + }); + + test('removeToken forgets the token and returns it', async function (assert) { + await this.authEngine.loadToken(); + + assert.strictEqual(await this.authEngine.removeToken('socketcluster.authToken'), 'token-1'); + assert.strictEqual(this.service.tokenState, null); + assert.strictEqual(this.service.refreshTimer, null); + assert.strictEqual(await this.authEngine.removeToken('socketcluster.authToken'), null, 'nothing left to remove'); + }); + + test('a token that arrives after logout is discarded', async function (assert) { + const stale = deferred(); + this.fetch.responder = () => stale.promise; + + const load = this.authEngine.loadToken(); + this.service.handleLogout(); + + // The next session's request is separate from the stale one. + this.fetch.responder = () => Promise.resolve({ token: 'fresh', expires_in: 900 }); + const next = this.service.getToken(); + assert.notStrictEqual(next, load); + + stale.resolve({ token: 'stale', expires_in: 900 }); + + assert.strictEqual(await load, null); + assert.strictEqual(await next, 'fresh'); + assert.strictEqual(this.service.tokenState.token, 'fresh', 'the stale token never landed'); + assert.strictEqual(this.service.tokenRequest, null); + }); + + // ------------------------------------------------------------------------- + // Refresh + // ------------------------------------------------------------------------- + + test('the refresh timer fires the leeway before expiry, with a floor', function (assert) { + const scheduled = []; + const originalSetTimeout = window.setTimeout; + window.setTimeout = (callback, delay) => { + scheduled.push({ callback, delay }); + // Negative ids can never clear a real timer when the service cancels. + return -scheduled.length; + }; + + try { + this.service.scheduleRefresh(900); + this.service.scheduleRefresh(10); + } finally { + window.setTimeout = originalSetTimeout; + } + + assert.strictEqual(scheduled[0].delay, (900 - REFRESH_LEEWAY_SECONDS) * 1000); + assert.strictEqual(scheduled[1].delay, MIN_REFRESH_DELAY_MS); + + let refreshed = 0; + this.service.refreshToken = () => refreshed++; + scheduled[1].callback(); + + assert.strictEqual(refreshed, 1, 'the timer refreshes the token'); + assert.strictEqual(this.service.refreshTimer, null); + }); + + test('minting a token arms the refresh timer', async function (assert) { + await this.authEngine.loadToken(); + + assert.notStrictEqual(this.service.refreshTimer, null); + }); + + test('refreshToken does nothing without a session', async function (assert) { + this.session.isAuthenticated = false; + + assert.strictEqual(await this.service.refreshToken(), null); + assert.strictEqual(this.fetch.calls.length, 0); + }); + + test('refreshToken re-authenticates the open socket with a new token', async function (assert) { + await this.authEngine.loadToken(); + + assert.strictEqual(await this.service.refreshToken(), 'token-2', 'the cached token is not reused'); + assert.deepEqual(this.client.authenticated, ['token-2']); + }); + + test('refreshToken only caches the token while disconnected', async function (assert) { + this.client.state = this.client.CLOSED; + + assert.strictEqual(await this.service.refreshToken(), 'token-1'); + assert.deepEqual(this.client.authenticated, []); + assert.strictEqual(await this.authEngine.loadToken(), 'token-1', 'the next handshake uses it'); + }); + + test('refreshToken skips authentication when no token could be minted', async function (assert) { + this.fetch.responder = () => Promise.reject(new Error('Not Found')); + + assert.strictEqual(await this.service.refreshToken(), null); + assert.deepEqual(this.client.authenticated, []); + }); + + test('a rejected token is forgotten', async function (assert) { + await this.authEngine.loadToken(); + this.client.authenticateImpl = () => Promise.resolve({ isAuthenticated: false, authError: new Error('bad') }); + + assert.false(await this.service.authenticateWith('token-1')); + assert.strictEqual(this.service.tokenState, null); + assert.false(this.service.authenticating); + }); + + test('a failed authenticate call is treated as a rejection', async function (assert) { + await this.authEngine.loadToken(); + this.client.authenticateImpl = () => Promise.reject(new Error('TimeoutError')); + + assert.false(await this.service.authenticateWith('token-1')); + assert.strictEqual(this.service.tokenState, null); + assert.false(this.service.authenticating); + }); + + // ------------------------------------------------------------------------- + // reauthenticate (login, restore, organization switch) + // ------------------------------------------------------------------------- + + test('reauthenticate does nothing without a session', async function (assert) { + this.session.isAuthenticated = false; + + assert.false(await this.service.reauthenticate()); + assert.strictEqual(this.fetch.calls.length, 0); + }); + + test('reauthenticate reconnects a closed socket and lets the handshake load the token', async function (assert) { + await this.authEngine.loadToken(); + this.client.state = this.client.CLOSED; + + assert.false(await this.service.reauthenticate()); + assert.strictEqual(this.client.connects, 1); + assert.strictEqual(this.service.tokenState, null, 'the old token is discarded'); + }); + + test('reauthenticate with force mints a new token for the open socket', async function (assert) { + await this.authEngine.loadToken(); + + assert.true(await this.service.reauthenticate({ force: true })); + assert.deepEqual(this.client.authenticated, ['token-2']); + }); + + test('reauthenticate without force reuses the cached token', async function (assert) { + await this.authEngine.loadToken(); + + assert.true(await this.service.reauthenticate({ force: false })); + assert.deepEqual(this.client.authenticated, ['token-1']); + assert.strictEqual(this.fetch.calls.length, 1); + }); + + test('reauthenticate reports failure when no token can be minted', async function (assert) { + this.fetch.responder = () => Promise.reject(new Error('Not Found')); + + assert.false(await this.service.reauthenticate()); + assert.deepEqual(this.client.authenticated, []); + }); + + test('login authenticates the socket', async function (assert) { + this.universe.trigger('session.authenticated'); + await flush(); + + assert.deepEqual(this.client.authenticated, ['token-1']); + }); + + test('a restored session authenticates an anonymous socket', async function (assert) { + this.universe.trigger('user.loaded'); + await flush(); + + assert.deepEqual(this.client.authenticated, ['token-1']); + }); + + test('a restored session leaves an authenticated socket alone', async function (assert) { + this.client.authState = this.client.AUTHENTICATED; + + assert.strictEqual(this.service.handleUserLoaded(), null); + assert.strictEqual(this.fetch.calls.length, 0); + }); + + test('an organization switch re-keys the socket with a new token', async function (assert) { + await this.authEngine.loadToken(); + + this.universe.trigger('user.organization_switched', { id: 'org-2' }); + await flush(); + + assert.deepEqual(this.client.authenticated, ['token-2']); + }); + + test('logout clears the token, timers and retries and disconnects', async function (assert) { + await this.authEngine.loadToken(); + this.service.pendingChannels.add('order.1'); + this.service.retriedChannels.set('order.2', 'token-1'); + this.service.recoveryLog = [Date.now()]; + + this.universe.trigger('user.deauthenticated'); + + assert.strictEqual(this.service.tokenState, null); + assert.strictEqual(this.service.refreshTimer, null); + assert.strictEqual(this.service.pendingChannels.size, 0); + assert.strictEqual(this.service.retriedChannels.size, 0); + assert.deepEqual(this.service.recoveryLog, []); + assert.strictEqual(this.client.disconnects, 1); + }); + + // ------------------------------------------------------------------------- + // Recovery + // ------------------------------------------------------------------------- + + test('an anonymous handshake with a session is authenticated on connect', async function (assert) { + assert.true(await this.service.handleConnect()); + assert.deepEqual(this.client.authenticated, ['token-1']); + }); + + test('connect is ignored when already authenticated or signed out', function (assert) { + this.client.authState = this.client.AUTHENTICATED; + assert.strictEqual(this.service.handleConnect(), null); + + this.client.authState = this.client.UNAUTHENTICATED; + this.session.isAuthenticated = false; + assert.strictEqual(this.service.handleConnect(), null); + assert.strictEqual(this.fetch.calls.length, 0); + }); + + test('deauthentication mints a new token and re-authenticates', async function (assert) { + await this.authEngine.loadToken(); + + assert.true(await this.service.handleDeauthenticate()); + assert.deepEqual(this.client.authenticated, ['token-2'], 'the dropped token is not reused'); + }); + + test('deauthentication caused by our own authenticate call is ignored', function (assert) { + this.service.authenticating = true; + + assert.strictEqual(this.service.handleDeauthenticate(), null); + assert.strictEqual(this.fetch.calls.length, 0); + }); + + test('deauthentication after sign-out is ignored', function (assert) { + this.session.isAuthenticated = false; + + assert.strictEqual(this.service.handleDeauthenticate(), null); + }); + + test('a channel lost for a token reason is resubscribed after re-authenticating', async function (assert) { + assert.true(await this.service.handleChannelLoss('order.1', 'token_expired')); + assert.deepEqual(this.client.authenticated, ['token-1']); + assert.deepEqual(this.client.subscribed, ['order.1']); + }); + + test('a channel the user may not see is not retried', function (assert) { + assert.strictEqual(this.service.handleChannelLoss('order.1', 'forbidden'), null); + assert.strictEqual(this.service.handleChannelLoss('order.1', null), null); + assert.strictEqual(this.fetch.calls.length, 0); + }); + + test('channel losses after sign-out are not retried', function (assert) { + this.session.isAuthenticated = false; + + assert.strictEqual(this.service.handleChannelLoss('order.1', 'no_token'), null); + }); + + test('recovery skips authenticate when the socket already holds the token', async function (assert) { + await this.authEngine.loadToken(); + this.client.authState = this.client.AUTHENTICATED; + this.client.signedAuthToken = 'token-1'; + + assert.true(await this.service.handleChannelLoss('order.1', 'identity_changed')); + assert.deepEqual(this.client.authenticated, []); + assert.deepEqual(this.client.subscribed, ['order.1']); + }); + + test('recovery re-authenticates a socket holding a different token', async function (assert) { + this.client.authState = this.client.AUTHENTICATED; + this.client.signedAuthToken = 'previous-organization'; + + assert.true(await this.service.handleChannelLoss('order.1', 'identity_changed')); + assert.deepEqual(this.client.authenticated, ['token-1']); + assert.deepEqual(this.client.subscribed, ['order.1']); + }); + + test('recovery gives up when no token can be minted', async function (assert) { + this.fetch.responder = () => Promise.reject(new Error('Not Found')); + + assert.false(await this.service.handleChannelLoss('order.1', 'no_token')); + assert.deepEqual(this.client.subscribed, []); + assert.strictEqual(this.service.pendingChannels.size, 0); + }); + + test('recovery gives up when the token is rejected', async function (assert) { + this.client.authenticateImpl = () => Promise.resolve({ isAuthenticated: false }); + + assert.false(await this.service.handleChannelLoss('order.1', 'no_token')); + assert.deepEqual(this.client.subscribed, []); + assert.strictEqual(this.service.pendingChannels.size, 0); + }); + + test('a channel is retried at most once per token', async function (assert) { + await this.service.handleChannelLoss('order.1', 'token_changed'); + await this.service.handleChannelLoss('order.1', 'token_changed'); + + assert.deepEqual(this.client.subscribed, ['order.1'], 'the second loss with the same token is not retried'); + + // A new token earns the channel another attempt. + await this.service.reauthenticate({ force: true }); + await this.service.handleChannelLoss('order.1', 'token_changed'); + + assert.deepEqual(this.client.subscribed, ['order.1', 'order.1']); + }); + + test('losses reported during a recovery are batched into it', async function (assert) { + const pending = deferred(); + this.fetch.responder = () => pending.promise; + + const first = this.service.handleChannelLoss('order.1', 'deauthenticated'); + const second = this.service.handleChannelLoss('order.2', 'deauthenticated'); + assert.strictEqual(first, second, 'one recovery is shared'); + + pending.resolve({ token: 'shared', expires_in: 900 }); + await first; + + assert.deepEqual(this.client.authenticated, ['shared']); + assert.deepEqual(this.client.subscribed, ['order.1', 'order.2']); + }); + + test('a loss reported while resubscribing gets a follow-up recovery', async function (assert) { + const service = this.service; + const client = this.client; + const originalSubscribe = client.subscribe; + client.subscribe = (channelName) => { + if (channelName === 'order.1') { + // The server refuses the next channel while we are still resubscribing. + service.handleChannelLoss('order.2', 'token_expired'); + } + return originalSubscribe(channelName); + }; + + await service.handleChannelLoss('order.1', 'token_expired'); + await flush(); + + assert.deepEqual(client.subscribed, ['order.1', 'order.2']); + assert.strictEqual(service.recovery, null); + }); + + test('recoveries are capped per window', async function (assert) { + for (let i = 0; i < MAX_RECOVERIES_PER_WINDOW; i++) { + assert.true(await this.service.handleChannelLoss(`order.${i}`, 'no_token')); + } + + assert.false(await this.service.handleChannelLoss('order.extra', 'no_token'), 'the next one is refused'); + assert.notOk(this.client.subscribed.includes('order.extra')); + assert.strictEqual(this.service.pendingChannels.size, 0); + }); + + test('old recoveries fall out of the window', async function (assert) { + const longAgo = Date.now() - RECOVERY_WINDOW_MS - 1; + this.service.recoveryLog = Array(MAX_RECOVERIES_PER_WINDOW).fill(longAgo); + + assert.true(await this.service.handleChannelLoss('order.1', 'no_token')); + assert.strictEqual(this.service.recoveryLog.length, 1); + }); + + // ------------------------------------------------------------------------- + // Client event wiring and teardown + // ------------------------------------------------------------------------- + + test('client events reach their handlers', async function (assert) { + this.client.emit('subscribeFail', { channel: 'order.1', error: authError('no_token') }); + await flush(); + assert.deepEqual(this.client.subscribed, ['order.1'], 'subscribeFail with an auth reason resubscribes'); + + this.client.emit('kickOut', { channel: 'order.2', message: 'token_expired' }); + await flush(); + assert.deepEqual(this.client.subscribed, ['order.1', 'order.2'], 'kickOut with an auth reason resubscribes'); + + this.client.authState = this.client.UNAUTHENTICATED; + this.client.emit('deauthenticate', {}); + await flush(); + assert.deepEqual(this.client.authenticated, ['token-1', 'token-2'], 'deauthenticate mints and authenticates'); + + this.service.recoveryLog = []; + this.client.authState = this.client.UNAUTHENTICATED; + this.client.emit('connect', {}); + await flush(); + assert.strictEqual(this.client.authenticated.length, 3, 'an anonymous connect authenticates'); + }); + + test('teardown stops the timer and every listener', async function (assert) { + await this.authEngine.loadToken(); + const client = this.client; + const universe = this.universe; + const service = this.service; + + run(() => service.destroy()); + await settled(); + + assert.deepEqual(client.closedListeners, CLIENT_EVENTS); + assert.strictEqual(service.refreshTimer, null); + + universe.trigger('user.deauthenticated'); + assert.strictEqual(client.disconnects, 0, 'session events no longer reach the service'); + }); +}); diff --git a/tests/unit/services/socket-listen-test.js b/tests/unit/services/socket-listen-test.js index 63dcb53f..5a5348a4 100644 --- a/tests/unit/services/socket-listen-test.js +++ b/tests/unit/services/socket-listen-test.js @@ -1,6 +1,7 @@ import { module, test } from 'qunit'; import { setupTest } from 'dummy/tests/helpers'; import { settled } from '@ember/test-helpers'; +import { inertClientMethods } from 'dummy/tests/helpers/stub-socketcluster'; /** * The body of `listen`'s async-iteration loop. @@ -43,6 +44,7 @@ module('Unit | Service | socket (listening)', function (hooks) { window.socketClusterClient = { create() { return { + ...inertClientMethods(), subscribe(channelId) { return fakeChannel(channelId, testContext.messages); }, diff --git a/tests/unit/services/socket-test.js b/tests/unit/services/socket-test.js index 0ccfc10e..31fcbe56 100644 --- a/tests/unit/services/socket-test.js +++ b/tests/unit/services/socket-test.js @@ -2,6 +2,7 @@ import { module, test } from 'qunit'; import { setupTest } from 'dummy/tests/helpers'; import { settled } from '@ember/test-helpers'; import config from 'dummy/config/environment'; +import { inertClientMethods } from 'dummy/tests/helpers/stub-socketcluster'; // The global SocketCluster client is replaced for the whole suite by // tests/helpers/stub-socketcluster, so no real connection is ever opened. These @@ -41,6 +42,7 @@ module('Unit | Service | socket', function (hooks) { create(socketConfig) { testContext.created.push(socketConfig); return { + ...inertClientMethods(), subscribe(channelId) { const channel = fakeChannel(channelId); testContext.subscribed.push(channel);