285 lines
13 KiB
JavaScript
285 lines
13 KiB
JavaScript
import assert from 'node:assert/strict';
|
|
import { readFileSync } from 'node:fs';
|
|
import test from 'node:test';
|
|
import ts from 'typescript';
|
|
import { parse } from 'vue/compiler-sfc';
|
|
|
|
function clock() {
|
|
let now = 0, sequence = 0;
|
|
const jobs = new Map();
|
|
const add = (fn, ms, repeat = false) => {
|
|
const id = ++sequence;
|
|
jobs.set(id, { fn, at: now + ms, ms, repeat });
|
|
return id;
|
|
};
|
|
const timers = {
|
|
setTimeout: (fn, ms) => add(fn, ms), clearTimeout: (id) => jobs.delete(id),
|
|
setInterval: (fn, ms) => add(fn, ms, true), clearInterval: (id) => jobs.delete(id),
|
|
};
|
|
const flush = async () => { for (let i = 0; i < 15; i++) await Promise.resolve(); };
|
|
return {
|
|
timers, flush, pending: () => jobs.size,
|
|
async tick(ms) {
|
|
await flush();
|
|
const end = now + ms;
|
|
for (;;) {
|
|
const next = [...jobs].filter(([, job]) => job.at <= end).sort((a, b) => a[1].at - b[1].at)[0];
|
|
if (!next) break;
|
|
const [id, job] = next;
|
|
now = job.at;
|
|
if (job.repeat) job.at += job.ms; else jobs.delete(id);
|
|
job.fn(); await flush();
|
|
}
|
|
now = end; await flush();
|
|
},
|
|
};
|
|
}
|
|
|
|
function evaluate(source, require, timer, suffix = '') {
|
|
const js = ts.transpileModule(source.replaceAll('import.meta.env', '{}'), {
|
|
compilerOptions: { module: ts.ModuleKind.CommonJS, target: ts.ScriptTarget.ES2020 },
|
|
}).outputText;
|
|
const exports = {};
|
|
return new Function('exports', 'require', ...Object.keys(timer.timers), js + suffix)(
|
|
exports, require, ...Object.values(timer.timers),
|
|
) || exports;
|
|
}
|
|
|
|
function socketFixture() {
|
|
const timer = clock(), tasks = [], frames = [];
|
|
let token = 'session-a', network, refreshes = 0;
|
|
globalThis.uni = {
|
|
onNetworkStatusChange: (fn) => { network = fn; },
|
|
connectSocket: (options) => {
|
|
const callbacks = {};
|
|
const task = {
|
|
options, sent: [], closed: false,
|
|
onOpen: (fn) => { callbacks.open = fn; }, onMessage: (fn) => { callbacks.message = fn; },
|
|
onError: (fn) => { callbacks.error = fn; }, onClose: (fn) => { callbacks.close = fn; },
|
|
send(data) { this.sent.push(JSON.parse(data.data)); if (this.failSend) data.fail?.({}); },
|
|
close() { this.closed = true; callbacks.close?.({ code: 1000 }); },
|
|
open: () => callbacks.open(), error: () => callbacks.error({}),
|
|
message: (frame) => callbacks.message({ data: JSON.stringify(frame) }),
|
|
auth() { this.open(); this.message({ command: 'AUTH_ACK' }); },
|
|
};
|
|
tasks.push(task); return task;
|
|
},
|
|
};
|
|
const im = evaluate(readFileSync(new URL('../src/utils/im.ts', import.meta.url), 'utf8'), () => ({
|
|
API_ORIGIN: 'https://test.invalid', getToken: () => token,
|
|
ensureAccessToken: async () => { refreshes++; token = 'refreshed-session'; return true; },
|
|
}), timer);
|
|
const off = im.connectIM((frame) => frames.push(frame));
|
|
return { ...im, timer, tasks, frames, off, network: (state) => network(state),
|
|
token: (value) => { token = value; }, refreshes: () => refreshes };
|
|
}
|
|
|
|
test('half-open sockets without a PONG are closed and reconnected', async () => {
|
|
const f = socketFixture();
|
|
try {
|
|
f.tasks[0].auth();
|
|
await f.timer.tick(40_000);
|
|
assert.equal(f.tasks[0].closed, true, 'a stale SocketTask must not block reconnection forever');
|
|
assert.ok(f.tasks.length >= 2);
|
|
} finally { f.off(); }
|
|
});
|
|
|
|
test('handshake watchdog, send failures and network changes recover without duplicate sockets', async () => {
|
|
const f = socketFixture();
|
|
try {
|
|
await f.timer.tick(15_000);
|
|
assert.ok(f.tasks.length >= 2, 'a connectSocket call with no callbacks must time out');
|
|
const current = f.tasks.at(-1); current.auth(); current.failSend = true;
|
|
await f.timer.tick(27_000);
|
|
assert.equal(current.closed, true);
|
|
f.tasks.at(-1).auth();
|
|
f.network({ isConnected: false });
|
|
assert.equal(f.tasks.at(-1).closed, true);
|
|
const count = f.tasks.length;
|
|
await f.timer.tick(35_000); assert.equal(f.tasks.length, count);
|
|
f.network({ isConnected: true }); assert.equal(f.tasks.length, count + 1);
|
|
f.network({ isConnected: true }); assert.equal(f.tasks.length, count + 1);
|
|
} finally { f.off(); }
|
|
});
|
|
|
|
test('foreground resume and account changes replace stale sockets and ignore late frames', async () => {
|
|
const f = socketFixture();
|
|
try {
|
|
const old = f.tasks[0]; old.auth();
|
|
f.suspendIM(); await f.timer.tick(60_000);
|
|
assert.equal(f.tasks.length, 1);
|
|
f.resumeIM(); assert.equal(f.tasks.length, 2); f.tasks[1].auth();
|
|
f.token('session-b'); f.resumeIM();
|
|
assert.deepEqual(f.tasks[2].options.protocols, ['xingyu.jwt.session-b']);
|
|
const before = f.frames.length;
|
|
old.message({ command: 'MESSAGE_PUSH', data: { secret: 'old account' } });
|
|
assert.equal(f.frames.length, before);
|
|
f.token(''); f.resumeIM(); await f.timer.tick(60_000);
|
|
assert.equal(f.tasks.length, 3);
|
|
} finally { f.off(); }
|
|
assert.equal(f.timer.pending(), 0);
|
|
});
|
|
|
|
test('only AUTH_ACK establishes IM readiness; listener exceptions cannot swallow other deliveries', async () => {
|
|
const f = socketFixture();
|
|
const offBad = f.connectIM(() => { throw new Error('broken page'); });
|
|
try {
|
|
f.tasks[0].open();
|
|
assert.equal(f.frames.some((frame) => frame.data?.status === 'connected'), false);
|
|
f.tasks[0].message({ command: 'AUTH_ACK' });
|
|
const seen = [];
|
|
const offPage = f.connectIM((frame) => seen.push(frame.command));
|
|
f.tasks[0].message({ command: 'MESSAGE_PUSH', data: { id: 1 } });
|
|
assert.ok(seen.includes('MESSAGE_PUSH'));
|
|
offPage();
|
|
await f.timer.tick(25_000); f.tasks[0].message({ command: 'PONG' });
|
|
await f.timer.tick(11_000); assert.equal(f.tasks.length, 1);
|
|
} finally { offBad(); f.off(); }
|
|
});
|
|
|
|
test('failed authentication handshake refreshes credentials even without close code 1008', async () => {
|
|
const f = socketFixture();
|
|
try {
|
|
f.tasks[0].error(); await f.timer.tick(2_000);
|
|
assert.equal(f.refreshes(), 1);
|
|
assert.deepEqual(f.tasks.at(-1).options.protocols, ['xingyu.jwt.refreshed-session']);
|
|
} finally { f.off(); }
|
|
});
|
|
|
|
function messagesFixture() {
|
|
const timer = clock(), hooks = {}, listeners = new Set(), replies = [], calls = [], counts = [];
|
|
let result = { items: [] };
|
|
globalThis.uni = { stopPullDownRefresh() {}, navigateTo() {} };
|
|
const source = parse(readFileSync(new URL('../src/pages/messages/index.vue', import.meta.url), 'utf8')).descriptor.scriptSetup.content;
|
|
const page = evaluate(source, (name) => {
|
|
if (name === 'vue') return { ref: (value) => ({ value }) };
|
|
if (name === '@dcloudio/uni-app') return Object.fromEntries(['onLoad', 'onShow', 'onHide', 'onUnload', 'onPullDownRefresh'].map((key) => [key, (fn) => { hooks[key] = fn; }]));
|
|
if (name.endsWith('.vue')) return {};
|
|
if (name.endsWith('/request')) return {
|
|
getToken: () => 'session-a',
|
|
api: { conversations: async (options) => { calls.push(options); return replies.length ? replies.shift() : result; } },
|
|
};
|
|
if (name.endsWith('/im')) return { resumeIM() {}, connectIM: (fn) => { listeners.add(fn); return () => listeners.delete(fn); } };
|
|
if (name.endsWith('/app-state')) return { syncConversationUnread: (items) => { counts.push(items.reduce((sum, item) => sum + item.unread, 0)); } };
|
|
if (name.endsWith('/message-display')) return { formatConversationTime: (value) => value };
|
|
throw new Error(name);
|
|
}, timer, '\nreturn {items,loading};');
|
|
return { ...page, timer, hooks, calls, replies, counts,
|
|
setResult: (value) => { result = { items: value }; },
|
|
emit: (frame) => listeners.forEach((fn) => fn(frame)),
|
|
show: async () => { hooks.onLoad?.(); hooks.onShow(); await timer.flush(); },
|
|
close: () => hooks.onUnload(),
|
|
};
|
|
}
|
|
|
|
test('visible message list reconciles missed pushes and unread badge without navigation', async () => {
|
|
const f = messagesFixture();
|
|
try {
|
|
await f.show();
|
|
const latest = [{ id: 2, lastMessage: 'new message', unread: 1 }, { id: 1, lastMessage: 'old', unread: 0 }];
|
|
f.setResult(latest);
|
|
await f.timer.tick(5_500);
|
|
assert.deepEqual(f.items.value, latest, 'remaining on the page must not require a push or a second onShow');
|
|
assert.equal(f.counts.at(-1), 1);
|
|
assert.equal(f.calls.at(-1).silent, true);
|
|
f.hooks.onHide(); const count = f.calls.length;
|
|
await f.timer.tick(60_000); assert.equal(f.calls.length, count);
|
|
f.hooks.onShow(); await f.timer.flush(); assert.ok(f.calls.length > count);
|
|
} finally { f.close(); }
|
|
assert.equal(f.timer.pending(), 0);
|
|
});
|
|
|
|
test('push bursts refresh previews/order/counts without skeleton flashes or starvation', async () => {
|
|
const f = messagesFixture();
|
|
try {
|
|
await f.show();
|
|
f.setResult([{ id: 1, lastMessage: 'incoming', unread: 1 }]);
|
|
for (let i = 0; i < 15; i++) {
|
|
f.emit({ command: 'MESSAGE_PUSH' });
|
|
await f.timer.tick(100);
|
|
assert.equal(f.loading.value, false);
|
|
}
|
|
assert.equal(f.items.value[0]?.lastMessage, 'incoming');
|
|
assert.ok(f.calls.length > 1 && f.calls.length < 15);
|
|
f.setResult([{ id: 1, lastMessage: 'incoming', unread: 0 }]);
|
|
f.emit({ command: 'READ_ACK' }); await f.timer.tick(200);
|
|
assert.equal(f.items.value[0].unread, 0);
|
|
} finally { f.close(); }
|
|
});
|
|
|
|
test('push during a pending fetch is replayed once; hidden-page responses cannot overwrite resumed state', async () => {
|
|
const f = messagesFixture();
|
|
try {
|
|
await f.show();
|
|
let resolve;
|
|
f.replies.push(new Promise((done) => { resolve = done; }));
|
|
f.emit({ command: 'MESSAGE_PUSH' }); await f.timer.tick(200);
|
|
const before = f.calls.length;
|
|
for (let i = 0; i < 10; i++) f.emit({ command: 'MESSAGE_PUSH' });
|
|
await f.timer.tick(200); assert.equal(f.calls.length, before);
|
|
f.setResult([{ id: 2, lastMessage: 'latest', unread: 2 }]);
|
|
resolve({ items: [] }); await f.timer.flush(); await f.timer.tick(200);
|
|
assert.equal(f.calls.length, before + 1); assert.equal(f.items.value[0].id, 2);
|
|
f.replies.push(new Promise((done) => { resolve = done; }));
|
|
f.emit({ command: 'MESSAGE_RECALLED' }); await f.timer.tick(200);
|
|
f.hooks.onHide(); f.setResult([{ id: 3, lastMessage: 'after resume', unread: 1 }]);
|
|
f.hooks.onShow(); await f.timer.flush();
|
|
resolve({ items: [{ id: 99, unread: 99 }] }); await f.timer.flush();
|
|
assert.equal(f.items.value[0].id, 3);
|
|
} finally { f.close(); }
|
|
});
|
|
|
|
test('connection-status flapping cannot postpone reconciliation; failed refresh preserves the list and retries', async () => {
|
|
const f = messagesFixture();
|
|
try {
|
|
f.setResult([{ id: 1, lastMessage: 'retained', unread: 1 }]);
|
|
await f.show();
|
|
f.setResult([{ id: 2, lastMessage: 'caught up', unread: 2 }]);
|
|
for (let i = 0; i < 6; i++) {
|
|
f.emit({ command: 'IM_STATUS', data: { status: 'reconnecting' } });
|
|
await f.timer.tick(1000);
|
|
}
|
|
assert.equal(f.items.value[0].id, 2);
|
|
// The deferred rejection happens after the request has installed its catch handler.
|
|
let reject;
|
|
f.replies.push(new Promise((_, fail) => { reject = fail; }));
|
|
f.emit({ command: 'MESSAGE_PUSH' }); await f.timer.tick(200);
|
|
reject(new Error('offline')); await f.timer.flush();
|
|
assert.equal(f.items.value[0].id, 2);
|
|
f.setResult([{ id: 3, lastMessage: 'retry success', unread: 3 }]);
|
|
await f.timer.tick(5_100);
|
|
assert.equal(f.items.value[0].id, 3);
|
|
f.emit({ command: 'IM_STATUS', data: { status: 'connected' } }); await f.timer.tick(200);
|
|
f.setResult([{ id: 4, lastMessage: 'lost push with healthy socket', unread: 4 }]);
|
|
await f.timer.tick(15_100);
|
|
assert.equal(f.items.value[0].id, 4);
|
|
} finally { f.close(); }
|
|
});
|
|
|
|
test('silent list requests do not toast on each retry and keep normal error feedback', async () => {
|
|
const timer = clock(), notices = [], requests = [];
|
|
let reply = 'network';
|
|
globalThis.uni = {
|
|
getStorageSync: () => '', showToast: (value) => notices.push(value),
|
|
request(options) {
|
|
requests.push(options);
|
|
if (reply === 'network') options.fail({ errMsg: 'offline' });
|
|
else options.success({ statusCode: 503, data: { message: 'retry later' } });
|
|
},
|
|
};
|
|
const source = readFileSync(new URL('../src/utils/request.ts', import.meta.url), 'utf8');
|
|
const request = evaluate(source, (name) => {
|
|
if (name === './device') return { getDeviceHeaders: () => ({}) };
|
|
if (name === './network-errors') return { isLoopbackServer: () => false, networkFailureMessage: () => 'offline' };
|
|
throw new Error(name);
|
|
}, timer);
|
|
await assert.rejects(request.api.conversations({ silent: true }));
|
|
assert.equal(notices.length, 0);
|
|
assert.equal('silent' in requests[0], false, 'internal option must not leak to uni.request');
|
|
reply = 'server';
|
|
await assert.rejects(request.api.conversations({ silent: true }));
|
|
assert.equal(notices.length, 0);
|
|
await assert.rejects(request.api.conversations());
|
|
assert.equal(notices.length, 1);
|
|
});
|