Skip to content

Commit 90bd3e2

Browse files
committed
chore: bump version to 1.2.34 in Cargo.toml and package.json; enhance createWorkerClient with Web Worker support and update conversation merging logic
1 parent 03aa1ab commit 90bd3e2

9 files changed

Lines changed: 329 additions & 16 deletions

File tree

crates/restsend-wasm/Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "restsend-wasm"
3-
version = "1.2.33"
3+
version = "1.2.34"
44
edition = "2021"
55
description = "Restsend Instant Messaging Javascript/Wasm SDK"
66
authors = ["Restsend Team <kui@fourz.cn>"]

crates/restsend-wasm/test/worker.spec.js

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import { endpoint } from './common.js'
66
describe('Worker quick path', async function () {
77
it('#sync chat logs in worker', async () => {
88
const info = await signin(endpoint, 'vitalik', 'vitalik:demo')
9-
const workerClient = await createWorkerClient(info, '')
9+
const workerClient = await createWorkerClient(info, '', { enableWorker: true })
1010

1111
const result = await workerClient.syncChatLogs('vitalik:guido', undefined, { limit: 5 })
1212

crates/restsend-wasm/worker/client.d.ts

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,14 @@
11
export interface WorkerClientOptions {
2+
/** Set to true to offload syncChatLogs to a Web Worker. Default: false */
3+
enableWorker?: boolean
4+
/** @deprecated Use enableWorker instead */
5+
forceFallback?: boolean
6+
/** Custom worker entry URL */
27
workerUrl?: string | URL
8+
/** Worker RPC call timeout (ms). Default: 8000 */
39
rpcTimeoutMs?: number
10+
/** Worker init timeout (ms). Default: 5000 */
411
initTimeoutMs?: number
5-
forceFallback?: boolean
612
}
713

814
export declare function createWorkerClient(

crates/restsend-wasm/worker/client.mjs

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -73,23 +73,32 @@ function buildProxy(facade, client) {
7373
}
7474

7575
/**
76-
* Create a worker-backed client wrapper.
76+
* Create a client wrapper with optional Web Worker support for syncChatLogs.
7777
*
78-
* The returned object is API-compatible with `new Client(info, dbName)` and forwards
79-
* all methods/properties to the underlying client except `syncChatLogs`, which is routed
80-
* to worker by default for better UI responsiveness.
78+
* By default the worker is **disabled**. Set `options.enableWorker = true` to offload
79+
* `syncChatLogs` calls to a Web Worker, keeping the main thread responsive during
80+
* heavy log fetching.
81+
*
82+
* @param {Object} info - Auth info (same as `new Client(info, dbName)`)
83+
* @param {string} [dbName=''] - IndexedDB database name
84+
* @param {Object} [options={}]
85+
* @param {boolean} [options.enableWorker=false] - Set to true to use a Web Worker for syncChatLogs
86+
* @param {number} [options.rpcTimeoutMs=8000] - Worker RPC call timeout (ms)
87+
* @param {number} [options.initTimeoutMs=5000] - Worker init timeout (ms)
88+
* @param {string} [options.workerUrl] - Custom worker entry URL
89+
* @returns {Promise<Object>} A proxy object that delegates to the underlying client
8190
*/
8291
export async function createWorkerClient(info, dbName = '', options = {}) {
8392
const rpcTimeoutMs = options.rpcTimeoutMs ?? 8000;
8493
const initTimeoutMs = options.initTimeoutMs ?? 5000;
85-
const forceFallback = options.forceFallback ?? false;
94+
const enableWorker = options.enableWorker === true;
8695

8796
const client = new Client(info, dbName);
8897

8998
let workerEnabled = false;
9099
let rpc = null;
91100

92-
if (!forceFallback) {
101+
if (enableWorker) {
93102
try {
94103
const workerUrl = options.workerUrl ?? new URL('./entry.mjs', import.meta.url);
95104
const worker = new Worker(workerUrl, { type: 'module' });

crates/restsend/src/client/store/conversations.rs

Lines changed: 276 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -395,6 +395,41 @@ impl ClientStore {
395395
}
396396

397397
conversation.merge_local_read_state(&old_conversation);
398+
399+
// Prefer local last_message when it is newer or server's is unreadable,
400+
// mirroring the single merge_conversation logic.
401+
let server_last_message_seq = conversation.last_message_seq;
402+
let local_last_message_readable = old_conversation
403+
.last_message
404+
.as_ref()
405+
.map(|c| !c.unreadable)
406+
.unwrap_or(false);
407+
let server_last_message_unreadable = conversation
408+
.last_message
409+
.as_ref()
410+
.map(|c| c.unreadable)
411+
.unwrap_or(false);
412+
413+
if let Some(local_seq) = old_conversation.last_message_seq {
414+
let should_use_local = match server_last_message_seq {
415+
Some(server_seq) if server_last_message_unreadable => {
416+
local_last_message_readable
417+
}
418+
Some(server_seq) => {
419+
local_last_message_readable && local_seq >= server_seq
420+
}
421+
None => local_last_message_readable,
422+
};
423+
if should_use_local {
424+
conversation.last_message = old_conversation.last_message.clone();
425+
conversation.last_message_at =
426+
old_conversation.last_message_at.clone();
427+
conversation.last_sender_id =
428+
old_conversation.last_sender_id.clone();
429+
conversation.last_message_seq =
430+
old_conversation.last_message_seq;
431+
}
432+
}
398433
}
399434

400435
if needs_last_readable_refresh(&conversation) {
@@ -762,6 +797,10 @@ impl ClientStore {
762797
new_conversation.last_read_at = conversation.last_read_at;
763798
new_conversation.last_read_seq = conversation.last_read_seq;
764799
new_conversation.unread = conversation.unread;
800+
new_conversation.last_message = conversation.last_message.clone();
801+
new_conversation.last_message_at = conversation.last_message_at.clone();
802+
new_conversation.last_sender_id = conversation.last_sender_id.clone();
803+
new_conversation.last_message_seq = conversation.last_message_seq;
765804
}
766805
self.ensure_topic_owner_id(&mut new_conversation).await;
767806
if let Ok(t) = self.message_storage.table::<Conversation>().await {
@@ -1579,4 +1618,241 @@ mod tests {
15791618
assert_eq!(second.items[0].id, "chat-3");
15801619
assert_eq!(second.items[2].id, "chat-1");
15811620
}
1621+
1622+
#[tokio::test]
1623+
async fn merge_conversations_prefers_local_newer_last_message_over_server_stale() {
1624+
let store = ClientStore::new("", ":memory:", "http://test", "token", "user1");
1625+
1626+
let agent_reply = ChatLog {
1627+
id: "log_101".to_string(),
1628+
topic_id: "topic-1".to_string(),
1629+
seq: 101,
1630+
sender_id: "user1".to_string(),
1631+
created_at: "2026-05-21T10:00:01Z".to_string(),
1632+
content: Content {
1633+
content_type: "text".to_string(),
1634+
text: "agent reply".to_string(),
1635+
unreadable: false,
1636+
..Default::default()
1637+
},
1638+
..Default::default()
1639+
};
1640+
let log_table = store.message_storage.table::<ChatLog>().await.unwrap();
1641+
log_table
1642+
.set(&agent_reply.topic_id, &agent_reply.id, Some(&agent_reply))
1643+
.await
1644+
.unwrap();
1645+
1646+
let local_conversation = Conversation {
1647+
topic_id: "topic-1".to_string(),
1648+
last_seq: 101,
1649+
last_read_seq: 101,
1650+
unread: 0,
1651+
last_message_seq: Some(101),
1652+
last_message: Some(agent_reply.content.clone()),
1653+
last_message_at: agent_reply.created_at.clone(),
1654+
last_sender_id: agent_reply.sender_id.clone(),
1655+
..Default::default()
1656+
};
1657+
let conv_table = store.message_storage.table::<Conversation>().await.unwrap();
1658+
conv_table
1659+
.set("", &local_conversation.topic_id, Some(&local_conversation))
1660+
.await
1661+
.unwrap();
1662+
1663+
let server_conversations = vec![Conversation {
1664+
topic_id: "topic-1".to_string(),
1665+
last_seq: 100,
1666+
last_read_seq: 101,
1667+
unread: 0,
1668+
last_message_seq: Some(100),
1669+
last_message: Some(Content {
1670+
content_type: "text".to_string(),
1671+
text: "customer message".to_string(),
1672+
unreadable: false,
1673+
..Default::default()
1674+
}),
1675+
last_message_at: "2026-05-21T10:00:00Z".to_string(),
1676+
last_sender_id: "user2".to_string(),
1677+
..Default::default()
1678+
}];
1679+
1680+
let merged = store.merge_conversations(server_conversations).await;
1681+
assert_eq!(merged.len(), 1);
1682+
assert_eq!(merged[0].last_message_seq, Some(101));
1683+
assert_eq!(
1684+
merged[0].last_message.as_ref().unwrap().text,
1685+
"agent reply"
1686+
);
1687+
assert_eq!(merged[0].unread, 0);
1688+
}
1689+
1690+
#[tokio::test]
1691+
async fn merge_conversations_keeps_server_last_message_when_server_is_newer() {
1692+
let store = ClientStore::new("", ":memory:", "http://test", "token", "user1");
1693+
1694+
let old_local = Conversation {
1695+
topic_id: "topic-1".to_string(),
1696+
last_seq: 100,
1697+
last_read_seq: 100,
1698+
unread: 0,
1699+
last_message_seq: Some(100),
1700+
last_message: Some(Content {
1701+
content_type: "text".to_string(),
1702+
text: "old message".to_string(),
1703+
unreadable: false,
1704+
..Default::default()
1705+
}),
1706+
last_message_at: "2026-05-21T10:00:00Z".to_string(),
1707+
last_sender_id: "user2".to_string(),
1708+
..Default::default()
1709+
};
1710+
let conv_table = store.message_storage.table::<Conversation>().await.unwrap();
1711+
conv_table
1712+
.set("", &old_local.topic_id, Some(&old_local))
1713+
.await
1714+
.unwrap();
1715+
1716+
let server_conversations = vec![Conversation {
1717+
topic_id: "topic-1".to_string(),
1718+
last_seq: 102,
1719+
last_read_seq: 101,
1720+
unread: 1,
1721+
last_message_seq: Some(102),
1722+
last_message: Some(Content {
1723+
content_type: "text".to_string(),
1724+
text: "newer server message".to_string(),
1725+
unreadable: false,
1726+
..Default::default()
1727+
}),
1728+
last_message_at: "2026-05-21T10:00:02Z".to_string(),
1729+
last_sender_id: "user3".to_string(),
1730+
..Default::default()
1731+
}];
1732+
1733+
let merged = store.merge_conversations(server_conversations).await;
1734+
assert_eq!(merged.len(), 1);
1735+
assert_eq!(merged[0].last_message_seq, Some(102));
1736+
assert_eq!(
1737+
merged[0].last_message.as_ref().unwrap().text,
1738+
"newer server message"
1739+
);
1740+
}
1741+
1742+
#[tokio::test]
1743+
async fn merge_conversations_preserves_local_read_state_with_zero_unread() {
1744+
let store = ClientStore::new("", ":memory:", "http://test", "token", "user1");
1745+
1746+
// Scenario: user read the conversation locally, unread=0
1747+
// Server still has unread>0 (read not synced, or server counts differently)
1748+
let local = Conversation {
1749+
topic_id: "topic-1".to_string(),
1750+
last_seq: 105,
1751+
last_read_seq: 100,
1752+
unread: 0,
1753+
last_message_seq: Some(105),
1754+
last_message: Some(Content {
1755+
content_type: "text".to_string(),
1756+
text: "readable".to_string(),
1757+
unreadable: false,
1758+
..Default::default()
1759+
}),
1760+
last_message_at: "2026-05-21T10:00:05Z".to_string(),
1761+
last_sender_id: "user1".to_string(),
1762+
..Default::default()
1763+
};
1764+
let conv_table = store.message_storage.table::<Conversation>().await.unwrap();
1765+
conv_table
1766+
.set("", &local.topic_id, Some(&local))
1767+
.await
1768+
.unwrap();
1769+
1770+
let server = Conversation {
1771+
topic_id: "topic-1".to_string(),
1772+
last_seq: 105,
1773+
last_read_seq: 100,
1774+
unread: 3,
1775+
last_message_seq: Some(105),
1776+
last_message: Some(Content {
1777+
content_type: "text".to_string(),
1778+
text: "same readable".to_string(),
1779+
unreadable: false,
1780+
..Default::default()
1781+
}),
1782+
last_message_at: "2026-05-21T10:00:05Z".to_string(),
1783+
last_sender_id: "user2".to_string(),
1784+
..Default::default()
1785+
};
1786+
1787+
let merged = store.merge_conversations(vec![server]).await;
1788+
assert_eq!(merged.len(), 1);
1789+
assert_eq!(merged[0].unread, 0);
1790+
}
1791+
1792+
#[tokio::test]
1793+
async fn merge_conversations_prefers_local_last_message_when_server_is_unreadable() {
1794+
let store = ClientStore::new("", ":memory:", "http://test", "token", "user1");
1795+
1796+
let readable_log = ChatLog {
1797+
id: "log_100".to_string(),
1798+
topic_id: "topic-1".to_string(),
1799+
seq: 100,
1800+
sender_id: "user2".to_string(),
1801+
created_at: "2026-05-21T10:00:00Z".to_string(),
1802+
content: Content {
1803+
content_type: "text".to_string(),
1804+
text: "readable summary".to_string(),
1805+
unreadable: false,
1806+
..Default::default()
1807+
},
1808+
..Default::default()
1809+
};
1810+
let log_table = store.message_storage.table::<ChatLog>().await.unwrap();
1811+
log_table
1812+
.set(&readable_log.topic_id, &readable_log.id, Some(&readable_log))
1813+
.await
1814+
.unwrap();
1815+
1816+
let local = Conversation {
1817+
topic_id: "topic-1".to_string(),
1818+
last_seq: 100,
1819+
last_read_seq: 100,
1820+
unread: 0,
1821+
last_message_seq: Some(100),
1822+
last_message: Some(readable_log.content.clone()),
1823+
last_message_at: readable_log.created_at.clone(),
1824+
last_sender_id: "user2".to_string(),
1825+
..Default::default()
1826+
};
1827+
let conv_table = store.message_storage.table::<Conversation>().await.unwrap();
1828+
conv_table
1829+
.set("", &local.topic_id, Some(&local))
1830+
.await
1831+
.unwrap();
1832+
1833+
let server = Conversation {
1834+
topic_id: "topic-1".to_string(),
1835+
last_seq: 101,
1836+
last_read_seq: 100,
1837+
unread: 1,
1838+
last_message_seq: Some(101),
1839+
last_message: Some(Content {
1840+
content_type: "text".to_string(),
1841+
text: "hidden system message".to_string(),
1842+
unreadable: true,
1843+
..Default::default()
1844+
}),
1845+
last_message_at: "2026-05-21T10:00:01Z".to_string(),
1846+
last_sender_id: "system".to_string(),
1847+
..Default::default()
1848+
};
1849+
1850+
let merged = store.merge_conversations(vec![server]).await;
1851+
assert_eq!(merged.len(), 1);
1852+
assert_eq!(merged[0].last_message_seq, Some(100));
1853+
assert_eq!(
1854+
merged[0].last_message.as_ref().unwrap().text,
1855+
"readable summary"
1856+
);
1857+
}
15821858
}

0 commit comments

Comments
 (0)