Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
168 changes: 138 additions & 30 deletions src/sse/sse.js
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,31 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions
/** @type {import("../htmx").HtmxInternalApi} */
var api

/**
* sseSources keeps a reference to the EventSource that was created for an element,
* together with the attribute values it was created from.
*
* htmx fires htmx:afterProcessNode again every time an element's attributes change,
* and it wipes that element's internal data right before doing so. Without this map
* the extension has no reference left to the connection it opened earlier, so it
* would open a second one and leave the first one open forever.
*
* @type {WeakMap<HTMLElement, {source: EventSource, url: string, closeAttribute: string}>}
*/
var sseSources = new WeakMap()

/**
* sseListeners keeps the listeners that were registered for an element, so that they
* can be removed before the element is registered again.
*
* This is needed for the same reason as sseSources: htmx wipes an element's internal
* data before it re-fires htmx:afterProcessNode, so sseEventListener no longer refers
* to what was registered earlier, and the element would end up listening twice.
*
* @type {WeakMap<HTMLElement, Array<{source: EventSource, name: string, listener: Function}>>}
*/
var sseListeners = new WeakMap()

htmx.defineExtension('sse', {

/**
Expand Down Expand Up @@ -42,16 +67,8 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions
var parent = evt.target || evt.detail.elt
switch (name) {
case 'htmx:beforeCleanupElement':
var internalData = api.getInternalData(parent)
// Try to remove remove an EventSource when elements are removed
var source = internalData.sseEventSource
if (source) {
api.triggerEvent(parent, 'htmx:sseClose', {
source,
type: 'nodeReplaced',
})
internalData.sseEventSource.close()
}
closeEventSource(parent, 'nodeReplaced')

return

Expand Down Expand Up @@ -85,6 +102,10 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions
* @param {HTMLElement} elt
*/
function registerSSE(elt) {
// Drop what this element registered earlier, so that processing it again does not
// leave it listening for the same event twice
removeSSEListeners(elt)

// Add message handlers for every `sse-swap` attribute
if (api.getAttributeValue(elt, 'sse-swap')) {
// Find closest existing event source
Expand Down Expand Up @@ -125,7 +146,7 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions

// Register the new listener
api.getInternalData(elt).sseEventListener = listener
source.addEventListener(sseEventName, listener)
addSSEListener(elt, source, sseEventName, listener)
}
}

Expand Down Expand Up @@ -162,7 +183,7 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions

// Register the new listener
api.getInternalData(elt).sseEventListener = listener
source.addEventListener(ts.trigger.slice(4), listener)
addSSEListener(elt, source, ts.trigger.slice(4), listener)
})
}
}
Expand All @@ -181,19 +202,49 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions
}

// handle extension source creation attribute
if (api.getAttributeValue(elt, 'sse-connect')) {
var sseURL = api.getAttributeValue(elt, 'sse-connect')
if (sseURL == null) {
return
}

var sseURL = api.getAttributeValue(elt, 'sse-connect')
if (sseURL) {
ensureEventSource(elt, sseURL, retryCount)
} else {
// the element does not ask for a connection anymore, close the one it owned
closeEventSource(elt, 'sourceReplaced')
}

registerSSE(elt)
}

/**
* ensureEventSource connects the given element to the given url.
*
* If the element already is connected to that url and the connection is still
* usable, then the existing EventSource is kept, and its reference is restored in
* the element's internal data, which htmx wipes every time it processes the element.
*
* @param {HTMLElement} elt
* @param {string} url
* @param {number} retryCount
*/
function ensureEventSource(elt, url, retryCount) {
var closeAttribute = api.getAttributeValue(elt, 'sse-close')
var previous = sseSources.get(elt)
var replacedSource = false

if (previous) {
if (previous.source.readyState !== EventSource.CLOSED) {
if (previous.url === url && previous.closeAttribute === closeAttribute) {
// Already connected, keep the stream instead of opening a second one
api.getInternalData(elt).sseEventSource = previous.source
return
}

// The element asks for a different connection now, close the previous one
closeEventSource(elt, 'sourceReplaced')
}

// Descendants are still listening on the source that is being replaced
replacedSource = true
}

var source = htmx.createEventSource(url)

source.onerror = function(err) {
Expand All @@ -219,20 +270,20 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions
source.onopen = function(evt) {
api.triggerEvent(elt, 'htmx:sseOpen', { source })

if (retryCount && retryCount > 0) {
if (replacedSource || (retryCount && retryCount > 0)) {
const childrenToFix = elt.querySelectorAll("[sse-swap], [data-sse-swap], [hx-trigger], [data-hx-trigger]")
for (let i = 0; i < childrenToFix.length; i++) {
registerSSE(childrenToFix[i])
}
replacedSource = false
// We want to increase the reconnection delay for consecutive failed attempts only
retryCount = 0
}
}

api.getInternalData(elt).sseEventSource = source
sseSources.set(elt, { source, url, closeAttribute })


var closeAttribute = api.getAttributeValue(elt, "sse-close");
if (closeAttribute) {
// close eventsource when this message is received
source.addEventListener(closeAttribute, function() {
Expand All @@ -254,20 +305,77 @@ This extension adds support for Server Sent Events to htmx. See /www/extensions
*/
function maybeCloseSSESource(elt) {
if (!api.bodyContains(elt)) {
var source = api.getInternalData(elt).sseEventSource
if (source != undefined) {
api.triggerEvent(elt, 'htmx:sseClose', {
source,
type: 'nodeMissing',
})
source.close()
// source = null
return true
}
return closeEventSource(elt, 'nodeMissing')
}
return false
}

/**
* addSSEListener starts listening for an event on the given source, and keeps track
* of the listener so that it can be removed again.
*
* @param {HTMLElement} elt
* @param {EventSource} source
* @param {string} name
* @param {Function} listener
*/
function addSSEListener(elt, source, name, listener) {
var listeners = sseListeners.get(elt)
if (!listeners) {
listeners = []
sseListeners.set(elt, listeners)
}

listeners.push({ source, name, listener })
source.addEventListener(name, listener)
}

/**
* removeSSEListeners stops listening for every event this element was registered for.
*
* @param {HTMLElement} elt
*/
function removeSSEListeners(elt) {
var listeners = sseListeners.get(elt)
if (!listeners) {
return
}

for (var i = 0; i < listeners.length; i++) {
listeners[i].source.removeEventListener(listeners[i].name, listeners[i].listener)
}
sseListeners.delete(elt)
}

/**
* closeEventSource closes the EventSource that was created for the given element,
* if there still is one, and reports the reason it was closed.
*
* The source is looked up in sseSources as well, because htmx wipes an element's
* internal data whenever it processes that element again.
*
* @param {HTMLElement} elt
* @param {string} type
* @returns boolean whether a source was closed
*/
function closeEventSource(elt, type) {
var internalData = api.getInternalData(elt)
var tracked = sseSources.get(elt)
var source = internalData.sseEventSource || (tracked && tracked.source)
if (!source) {
return false
}

api.triggerEvent(elt, 'htmx:sseClose', {
source,
type,
})
source.close()
sseSources.delete(elt)
delete internalData.sseEventSource
return true
}


/**
* @param {HTMLElement} elt
Expand Down
Loading