diff --git a/src/editors/jetbrains/htmx.web-types.json b/src/editors/jetbrains/htmx.web-types.json
index 64af559bd..cfefde7fc 100644
--- a/src/editors/jetbrains/htmx.web-types.json
+++ b/src/editors/jetbrains/htmx.web-types.json
@@ -566,18 +566,18 @@
"doc-url": "https://four.htmx.org/extensions/hx-sse#htmxafterssemessage"
},
{
- "name": "before:ws:connection",
+ "name": "ws:before:connection",
"description": "Fires before a WebSocket connection attempt. `detail.connection` can be modified; set `cancelled` or cancel the event to stop connecting.",
- "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxbeforewsconnection"
+ "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsbeforeconnection"
},
{
- "name": "after:ws:connection",
+ "name": "ws:after:connection",
"description": "Fires after a successful WebSocket connection. `detail.connection` describes the connection.",
- "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxafterwsconnection"
+ "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsafterconnection"
},
{
"name": "ws:close",
- "description": "Fires when a WebSocket connection closes. `detail.connection`, `detail.reason`, and `detail.code` describe the close.",
+ "description": "Fires when a WebSocket connection closes. `detail.connection`, `detail.reason`, and `detail.code` describe the close. Codes in `detail.connection.config.reconnectCodes` reconnect.",
"doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsclose"
},
{
@@ -586,24 +586,24 @@
"doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwserror"
},
{
- "name": "before:ws:request",
- "description": "Fires before sending a WebSocket message. `detail.headers` and `detail.body` are modifiable. Cancel to skip sending.",
- "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxbeforewsrequest"
+ "name": "ws:before:message:outgoing",
+ "description": "Fires before sending a WebSocket message. Modify `detail.message`, use `detail.waitUntil()` to delay sending, or set `detail.cancelled` to cancel.",
+ "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsbeforemessageoutgoing"
},
{
- "name": "after:ws:request",
- "description": "Fires after a WebSocket message is sent. `detail.headers` and `detail.body` contain the sent payload.",
- "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxafterwsrequest"
+ "name": "ws:after:message:outgoing",
+ "description": "Fires after a WebSocket message is sent. `detail.message.data` is the value passed to `WebSocket.send()`.",
+ "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsaftermessageoutgoing"
},
{
- "name": "before:ws:message",
- "description": "Fires before a WebSocket message is processed. `detail.message.text`, `detail.message.json`, and `detail.message.cancelled` are available.",
- "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxbeforewsmessage"
+ "name": "ws:before:message:incoming",
+ "description": "Fires before an incoming WebSocket message is processed. Convert `detail.message`, use `detail.waitUntil()` to delay processing, or set `detail.cancelled` to cancel. Associated messages fire from their sending element.",
+ "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsbeforemessageincoming"
},
{
- "name": "after:ws:message",
- "description": "Fires after a WebSocket message is processed. `detail.message.text` and `detail.message.json` are available.",
- "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxafterwsmessage"
+ "name": "ws:after:message:incoming",
+ "description": "Fires after an incoming WebSocket message is processed. `detail.message.data` preserves the native data; `text()`, `json()`, `blob()`, and `arrayBuffer()` convert it. Associated messages fire from their sending element.",
+ "doc-url": "https://four.htmx.org/extensions/hx-ws#htmxwsaftermessageincoming"
},
{
"name": "download:start",
diff --git a/src/ext/hx-ws.js b/src/ext/hx-ws.js
index 3138c6bad..67ed6194f 100644
--- a/src/ext/hx-ws.js
+++ b/src/ext/hx-ws.js
@@ -1,4 +1,5 @@
(() => {
+ const MESSAGE_ID_MAX_AGE = 30000;
let api;
// Build a CSS selector for querySelectorAll, respecting prefix + metaCharacter
@@ -15,29 +16,28 @@
// ========================================
function getConfig(element) {
- const defaults = {
+ let hxConfig = api.HCON.parse(api.attributeValue(element, 'hx-config')).ws || {};
+
+ return {
reconnect: true,
+ reconnectCodes: [
+ 1001, // Going Away
+ 1005, // No Status Received
+ 1006, // Abnormal Closure
+ 1011, // Internal Error
+ 1012, // Service Restart
+ 1013, // Try Again Later
+ 1014 // Bad Gateway
+ ],
reconnectDelay: 500,
reconnectMaxDelay: 60000,
reconnectMaxAttempts: Infinity,
reconnectJitter: 0.3,
pauseOnBackground: true,
- pendingRequestTTL: 30000
+ maxOutgoingMessagesQueueSize: 100,
+ ...htmx.config.ws, // global defaults
+ ...hxConfig // hx-config overrides
};
- let global = htmx.config.ws || {};
- let perElement = {};
- if (element) {
- let ctx = api.createRequestContext(element, new CustomEvent('_'));
- perElement = ctx.request.ws || {};
- }
- let merged = { ...defaults, ...global, ...perElement };
-
- // Backwards compat: boolean reconnectJitter (old API used true/false)
- if (typeof merged.reconnectJitter === 'boolean') {
- merged.reconnectJitter = merged.reconnectJitter ? 0.3 : 0;
- }
-
- return merged;
}
// ========================================
@@ -96,13 +96,16 @@
socket: null,
attempt: 0,
timer: null,
- pendingRequests: new Map(),
+ pendingMessages: new Map(),
+ queue: [],
+ receiving: Promise.resolve(),
+ sending: Promise.resolve(),
abortController: null,
visibilityHandler: null,
cancelled: false
};
- if (!api.triggerHtmxEvent(element, 'htmx:before:ws:connection', {connection}) || connection.cancelled) {
+ if (!api.triggerHtmxEvent(element, 'htmx:ws:before:connection', {connection}) || connection.cancelled) {
api.triggerHtmxEvent(element, 'htmx:ws:close', {
connection, reason: 'cancelled', code: null
});
@@ -148,7 +151,8 @@
if (connection.abortController) {
connection.abortController.abort();
}
- connection.pendingRequests.clear();
+ connection.pendingMessages.clear();
+ connection.queue.length = 0;
if (connection.socket) {
try {
if (connection.socket.readyState === WebSocket.OPEN || connection.socket.readyState === WebSocket.CONNECTING) {
@@ -187,17 +191,23 @@
connection.socket.addEventListener('open', () => {
let elt = findConnectedElement(url);
if (elt) {
- api.triggerHtmxEvent(elt, 'htmx:after:ws:connection', {connection});
+ api.triggerHtmxEvent(elt, 'htmx:ws:after:connection', {connection});
} else {
// Element was removed while connecting — orphaned socket
cleanupOrphanedConnection(url, connection);
return;
}
connection.attempt = 0;
+ flushQueue(connection);
}, opts);
connection.socket.addEventListener('message', (event) => {
- handleMessage(connection, event);
+ connection.receiving = connection.receiving
+ .then(() => handleMessage(connection, event))
+ .catch(error => {
+ let elt = findConnectedElement(connection.url);
+ if (elt) api.triggerHtmxEvent(elt, 'htmx:ws:error', { url: connection.url, error });
+ });
}, opts);
connection.socket.addEventListener('close', (event) => {
@@ -213,7 +223,7 @@
let config = connection.config;
if (config.pauseOnBackground && document.hidden) return;
- if (config.reconnect && findConnectedElement(url)) {
+ if (config.reconnect && config.reconnectCodes.includes(event.code) && findConnectedElement(url)) {
scheduleReconnect(url, connection);
} else {
// No element or reconnect disabled — full cleanup
@@ -259,7 +269,7 @@
let elt = findConnectedElement(url);
if (elt) {
connection.cancelled = false;
- if (!api.triggerHtmxEvent(elt, 'htmx:before:ws:connection', {connection}) || connection.cancelled) {
+ if (!api.triggerHtmxEvent(elt, 'htmx:ws:before:connection', {connection}) || connection.cancelled) {
api.triggerHtmxEvent(elt, 'htmx:ws:close', {
connection, reason: 'cancelled', code: null
});
@@ -292,7 +302,8 @@
if (connection.abortController) {
connection.abortController.abort();
}
- connection.pendingRequests.clear();
+ connection.pendingMessages.clear();
+ connection.queue.length = 0;
api.triggerHtmxEvent(element, 'htmx:ws:close', {
connection, reason: 'removed', code: null
});
@@ -303,25 +314,42 @@
}
// ========================================
- // PENDING REQUEST MANAGEMENT // ========================================
+ // PENDING MESSAGE MANAGEMENT
+ // ========================================
- function cleanupExpiredRequests(connection) {
- let config = connection.config;
+ function cleanupExpiredMessages(connection) {
let now = Date.now();
- let timeout = config.pendingRequestTTL || 30000;
- for (let [requestId, pending] of connection.pendingRequests) {
- if (now - pending.timestamp > timeout) {
- connection.pendingRequests.delete(requestId);
+ for (let [messageId, pending] of connection.pendingMessages) {
+ if (now - pending.timestamp > MESSAGE_ID_MAX_AGE) {
+ connection.pendingMessages.delete(messageId);
}
}
}
// ========================================
- // REQUESTS
+ // MESSAGES
// ========================================
- async function sendRequest(element, event) {
+ function transmitMessage(connection, element, message) {
+ try {
+ connection.socket.send(message.data);
+ let messageId = message.headers['HX-Message-ID'];
+ if (messageId) connection.pendingMessages.set(messageId, { element, timestamp: Date.now() });
+ api.triggerHtmxEvent(element, 'htmx:ws:after:message:outgoing', {message});
+ } catch (error) {
+ api.triggerHtmxEvent(element, 'htmx:ws:error', { url: connection.url, error });
+ }
+ }
+
+ function flushQueue(connection) {
+ while (connection.queue.length && connection.socket?.readyState === WebSocket.OPEN) {
+ let queuedMessage = connection.queue.shift();
+ transmitMessage(connection, queuedMessage.element, queuedMessage.message);
+ }
+ }
+
+ async function sendMessage(element, event) {
// hx-ws:send="/url" creates its own connection; hx-ws:send (no value) uses ancestor's
let sendAttr = api.attributeValue(element, 'hx-ws:send');
let url = (sendAttr && sendAttr !== 'true') ? sendAttr : null;
@@ -342,143 +370,197 @@
let normalizedUrl = normalizeWebSocketUrl(url);
let connection = connections.get(normalizedUrl);
- // Wait for socket to open if still connecting
- if (connection && connection.socket && connection.socket.readyState === WebSocket.CONNECTING) {
- await new Promise(resolve => {
- connection.socket.addEventListener('open', resolve, { once: true });
- connection.socket.addEventListener('close', resolve, { once: true });
- connection.socket.addEventListener('error', resolve, { once: true });
- });
- }
-
- if (!connection || !connection.socket || connection.socket.readyState !== WebSocket.OPEN) {
+ if (!connection) {
api.triggerHtmxEvent(element, 'htmx:ws:error', { url: normalizedUrl, error: 'Connection not open' });
return;
}
- // [Correlation] Cleanup expired pending requests periodically
- cleanupExpiredRequests(connection);
+ // [Correlation] Cleanup expired pending messages periodically
+ cleanupExpiredMessages(connection);
// Build headers using core's request context (same as HTTP requests)
let ctx = api.createRequestContext(element, event);
let headers = {...ctx.request.headers};
delete headers['Accept'];
- // [Correlation] Add request ID as a header
- let requestId = crypto.randomUUID();
- headers['HX-Request-ID'] = requestId;
+ // [Correlation] Add message ID as a header
+ headers['HX-Message-ID'] = crypto.randomUUID();
- // Build body from form data
+ // Build outgoing values from form data.
let form = element.form || element.closest('form');
let formData = api.collectFormData(element, form, event.submitter);
// Preserve multi-value form fields (checkboxes, multi-selects)
- let body = {};
+ let values = {};
for (let [key, value] of formData) {
- if (key in body) {
- body[key] = [].concat(body[key], value);
+ if (key in values) {
+ values[key] = [].concat(values[key], value);
} else {
- body[key] = value;
+ values[key] = value;
}
}
// Merge hx-vals after serialization to preserve JS types (numbers, booleans)
- let valsResult = api.getAttributeObject(element, 'hx-vals', obj => Object.assign(body, obj));
- if (valsResult) await valsResult;
+ let hxValsResult = api.getAttributeObject(element, 'hx-vals', obj => Object.assign(values, obj));
- let detail = { headers, body };
- if (!api.triggerHtmxEvent(element, 'htmx:before:ws:request', detail)) {
- return;
- }
+ let outgoingMessage = connection.sending.then(async () => {
+ if (hxValsResult) await hxValsResult;
+ delete values.headers;
- try {
- connection.socket.send(JSON.stringify(detail));
+ let pendingWork = [];
+ let message = {
+ headers,
+ values,
+ data: undefined
+ };
+ let detail = {
+ message,
+ cancelled: false,
+ waitUntil(promise) {
+ pendingWork.push(Promise.resolve(promise));
+ }
+ };
+ let shouldSend = api.triggerHtmxEvent(element, 'htmx:ws:before:message:outgoing', detail);
- // [Correlation] Store pending request for response matching
- connection.pendingRequests.set(requestId, { element, timestamp: Date.now() });
+ try {
+ await Promise.all(pendingWork);
+ if (!shouldSend || detail.cancelled) return;
- api.triggerHtmxEvent(element, 'htmx:after:ws:request', detail);
- } catch (error) {
- api.triggerHtmxEvent(element, 'htmx:ws:error', { url: normalizedUrl, error });
- }
+ message.data ??= JSON.stringify({ ...message.values, headers: message.headers });
+ if (connections.get(normalizedUrl) !== connection) {
+ api.triggerHtmxEvent(element, 'htmx:ws:error', { url: normalizedUrl, error: 'Connection closed' });
+ return;
+ }
+
+ if (connection.socket?.readyState === WebSocket.OPEN) {
+ transmitMessage(connection, element, message);
+ } else if (connection.queue.length >= connection.config.maxOutgoingMessagesQueueSize) {
+ api.triggerHtmxEvent(element, 'htmx:ws:error', {
+ url: normalizedUrl,
+ error: 'Outgoing messages queue is full'
+ });
+ } else {
+ connection.queue.push({element, message});
+ }
+ } catch (error) {
+ api.triggerHtmxEvent(element, 'htmx:ws:error', { url: normalizedUrl, error });
+ }
+ });
+ connection.sending = outgoingMessage.catch(() => {});
+ await outgoingMessage;
}
// ========================================
// MESSAGE RECEIVING & ROUTING
// ========================================
- function handleMessage(connection, event) {
+ async function handleMessage(connection, event) {
+ let data = event.data;
+ let textResult;
+ let jsonResult;
+ let arrayBufferResult;
+ let blobResult;
+ let pendingWork = [];
+ let message = {
+ data,
+ type: typeof data === 'string' ? 'text' : 'binary',
+ text() {
+ return textResult ??= typeof data === 'string'
+ ? Promise.resolve(data)
+ : data instanceof Blob
+ ? data.text()
+ : Promise.resolve(new TextDecoder().decode(data));
+ },
+ json() {
+ return jsonResult ??= message.text().then(JSON.parse);
+ },
+ arrayBuffer() {
+ return arrayBufferResult ??= data instanceof ArrayBuffer
+ ? Promise.resolve(data)
+ : data instanceof Blob
+ ? data.arrayBuffer()
+ : Promise.resolve(new TextEncoder().encode(data).buffer);
+ },
+ blob() {
+ return blobResult ??= data instanceof Blob
+ ? Promise.resolve(data)
+ : Promise.resolve(new Blob([data]));
+ }
+ };
+
let json = null;
- try {
- json = JSON.parse(event.data);
- } catch (e) {
- // Not JSON - will be treated as raw HTML below
- }
-
- // [Correlation] Cleanup expired pending requests on every message
- cleanupExpiredRequests(connection);
-
- // [Correlation] Match response to originating element, or fall back to first subscriber
- let connectionElement = null;
- let requestId = json?.['HX-Request-ID'] || json?.request_id;
- if (requestId && connection.pendingRequests.has(requestId)) {
- connectionElement = connection.pendingRequests.get(requestId).element;
- connection.pendingRequests.delete(requestId);
- // If the correlated element has been removed from the DOM, fall back
- if (!connectionElement.isConnected) {
- connectionElement = findConnectedElement(connection.url);
+ if (message.type === 'text') {
+ try {
+ json = await message.json();
+ } catch (e) {
+ // Non-JSON text is treated as raw HTML.
}
- } else {
- connectionElement = findConnectedElement(connection.url);
}
- if (!connectionElement) {
+ // [Correlation] Cleanup expired pending messages on every message
+ cleanupExpiredMessages(connection);
+
+ let messageId = json?.headers?.['HX-Message-ID'];
+ let pending = connection.pendingMessages.get(messageId);
+ if (pending) connection.pendingMessages.delete(messageId);
+
+ // Route associated incoming messages through their sender.
+ let element = pending?.element;
+ if (!element?.isConnected) element = findConnectedElement(connection.url);
+
+ if (!element) {
// No element in DOM for this connection — orphan cleanup
cleanupOrphanedConnection(connection.url, connection);
return;
}
let detail = {
- message: { text: event.data, json, cancelled: false }
+ message,
+ cancelled: false,
+ waitUntil(promise) {
+ pendingWork.push(Promise.resolve(promise));
+ }
};
+ let shouldProcess = api.triggerHtmxEvent(element, 'htmx:ws:before:message:incoming', detail);
- if (!api.triggerHtmxEvent(connectionElement, 'htmx:before:ws:message', detail) || detail.message.cancelled) {
- return;
- }
+ await Promise.all(pendingWork);
+ if (!shouldProcess || detail.cancelled) return;
// JSON with 'content' or 'payload' field: swap the HTML
// Raw (non-JSON) string: swap the entire string as HTML
// JSON without 'content'/'payload': data-only message, no swap (handle via events)
let html;
- if (detail.message.json) {
- if (detail.message.json.content !== undefined) {
- html = detail.message.json.content;
- } else if (detail.message.json.payload !== undefined) {
- html = detail.message.json.payload; // backwards compat
+ if (json) {
+ if (json.content !== undefined) {
+ html = json.content;
+ } else if (json.payload !== undefined) {
+ html = json.payload; // backwards compat
// Warn once per connection (not on every message)
if (!connection._payloadWarnFired) {
console.warn('htmx: [hx-ws] json.payload is deprecated; use json.content instead');
connection._payloadWarnFired = true;
}
}
- } else {
- html = detail.message.text;
+ } else if (message.type === 'text') {
+ html = await message.text();
}
if (html != null) {
- let target = detail.message.json?.target || api.attributeValue(connectionElement, 'hx-target');
- let swap = detail.message.json?.swap || api.attributeValue(connectionElement, 'hx-swap');
-
- htmx.swap({
- sourceElement: connectionElement,
- target: target || connectionElement,
- swap: swap || (target ? htmx.config.defaultSwap : 'none'),
+ let target = json?.target || api.attributeValue(element, 'hx-target');
+ let swap = json?.swap || api.attributeValue(element, 'hx-swap') || htmx.config.defaultSwap;
+ if (!/(?:^|\s)swapEmpty(?::(?:true|false))?(?=\s|$)/.test(swap)) swap += ' swapEmpty:false';
+
+ await htmx.swap({
+ sourceElement: element,
+ target: target || element,
+ swap,
+ select: json?.select ?? api.attributeValue(element, 'hx-select'),
+ selectOOB: api.attributeValue(element, 'hx-select-oob'),
text: html,
transition: false
});
}
- delete detail.message.cancelled;
- api.triggerHtmxEvent(connectionElement, 'htmx:after:ws:message', detail);
+ api.triggerHtmxEvent(element, 'htmx:ws:after:message:incoming', {message});
}
// ========================================
@@ -526,7 +608,7 @@
element._htmx.ws.url = connection.url;
}
}
- await sendRequest(element, evt);
+ await sendMessage(element, evt);
});
element._htmx.ws.sendInitialized = true;
}
@@ -636,7 +718,8 @@
if (connection.socket) {
connection.socket.close();
}
- connection.pendingRequests.clear();
+ connection.pendingMessages.clear();
+ connection.queue.length = 0;
});
},
get: (key) => connections.get(normalizeWebSocketUrl(key)),
diff --git a/src/scripts/upgrade-check.py b/src/scripts/upgrade-check.py
index a91f939c8..33eae8612 100755
--- a/src/scripts/upgrade-check.py
+++ b/src/scripts/upgrade-check.py
@@ -90,13 +90,13 @@
}
WS_EVENT_RENAMES = {
- "htmx:wsOpen": "htmx:after:ws:connection",
+ "htmx:wsOpen": "htmx:ws:after:connection",
"htmx:wsClose": "htmx:ws:close",
- "htmx:wsConfigSend": "htmx:before:ws:request",
- "htmx:wsBeforeSend": "htmx:before:ws:request",
- "htmx:wsAfterSend": "htmx:after:ws:request",
- "htmx:wsBeforeMessage": "htmx:before:ws:message",
- "htmx:wsAfterMessage": "htmx:after:ws:message",
+ "htmx:wsConfigSend": "htmx:ws:before:message:outgoing",
+ "htmx:wsBeforeSend": "htmx:ws:before:message:outgoing",
+ "htmx:wsAfterSend": "htmx:ws:after:message:outgoing",
+ "htmx:wsBeforeMessage": "htmx:ws:before:message:incoming",
+ "htmx:wsAfterMessage": "htmx:ws:after:message:incoming",
}
# Extension attribute renames
diff --git a/test/manual/WS_README.md b/test/manual/WS_README.md
index d10bb32d5..d217e3f09 100644
--- a/test/manual/WS_README.md
+++ b/test/manual/WS_README.md
@@ -27,7 +27,7 @@ A beautiful, comprehensive demonstration of the `hx-ws` extension showcasing rea
### 2. **Live Notifications**
- Receive random notifications every 5-8 seconds
- Shows real-time server push
-- Uses `beforeend` swap to prepend new notifications
+- Uses `afterbegin` to prepend new notifications
### 3. **Shared Counter**
- Multiple clients share the same counter state
@@ -64,32 +64,40 @@ A beautiful, comprehensive demonstration of the `hx-ws` extension showcasing rea
## 🎨 Key Concepts
-### HTML Partial Format
+### HTML Message Format
-Server messages use this format:
+Server messages use `content` for HTML and may specify a target and serialized swap specification:
```json
{
- "channel": "ui",
- "format": "html",
- "payload": "Content"
+ "content": "
Content
",
+ "target": "#target-id",
+ "swap": "beforeend settle:10ms"
}
```
-### Request/Response Pattern
+### Message Flow
Client sends:
```json
{
- "type": "request",
- "request_id": "uuid-here",
- "values": {
- "message": "Hello!"
- }
+ "headers": {
+ "HX-Message-ID": "uuid-here"
+ },
+ "message": "Hello!"
}
```
-Server responds with matching `request_id` to target the originating element.
+The server copies `HX-Message-ID` into the incoming message:
+
+```json
+{
+ "headers": {
+ "HX-Message-ID": "uuid-here"
+ },
+ "content": "Saved
"
+}
+```
### Multiple Partials
@@ -112,7 +120,8 @@ htmx.config.ws = {
reconnectMaxDelay: 60000, // Max delay (ms)
reconnectMaxAttempts: Infinity,// Max reconnect attempts
reconnectJitter: 0.3, // Jitter factor (0-1)
- pauseOnBackground: true // Pause connection when tab is backgrounded
+ pauseOnBackground: true, // Pause connection when tab is backgrounded
+ protocols: null // Optional WebSocket subprotocols
};
```
@@ -143,7 +152,7 @@ htmx.config.ws = {
### Button Actions
```html
-