Skip to content
Closed
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
49 changes: 49 additions & 0 deletions crates/ironclaw_gateway/static/app.js
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ let pairingPollInterval = null;
let unreadThreads = new Map(); // thread_id -> unread count
let _loadThreadsTimer = null;
const JOB_EVENTS_CAP = 500;
const JOB_EVENTS_MAX_JOBS = 50;
const MEMORY_SEARCH_QUERY_MAX_LENGTH = 100;
let stagedImages = [];
let authFlowPending = false;
Expand Down Expand Up @@ -236,6 +237,19 @@ const DONE_WITHOUT_RESPONSE_TIMEOUT_MS = 1500;
let _turnResponseReceived = false;
let _doneWithoutResponseTimer = null;

// Clean up connection-level timers and buffers.
// Called before creating a new connection, on tab hide, and on page unload
// to prevent leaked intervals/timeouts from accumulating across reconnects.
// Note: _doneWithoutResponseTimer is intentionally NOT cleared here — it is a
// turn-level concern managed by the onopen and response handlers (#2079).
function cleanupConnectionState() {
if (_streamDebounceTimer) { clearInterval(_streamDebounceTimer); _streamDebounceTimer = null; }
_streamBuffer = '';
if (_connectionLostTimer) { clearTimeout(_connectionLostTimer); _connectionLostTimer = null; }
Comment thread
henrypark133 marked this conversation as resolved.
_connectionLostAt = null;
if (jobListRefreshTimer) { clearTimeout(jobListRefreshTimer); jobListRefreshTimer = null; }
}

// --- Send Cooldown State ---
let _sendCooldown = false;
let _recentLocalPairingApprovals = new Map();
Expand Down Expand Up @@ -390,6 +404,7 @@ document.getElementById('token-input').addEventListener('keydown', (e) => {
// Without this, stale SSE connections from prior page loads linger and exhaust
// the HTTP/1.1 per-origin connection limit (6), blocking API fetch calls.
window.addEventListener('beforeunload', () => {
cleanupConnectionState();
if (eventSource) { eventSource.close(); eventSource = null; }
if (logEventSource) { logEventSource.close(); logEventSource = null; }
});
Expand All @@ -400,6 +415,7 @@ window.addEventListener('beforeunload', () => {
// the 3rd tab exhausts the browser's per-origin limit.
document.addEventListener('visibilitychange', () => {
if (document.hidden) {
cleanupConnectionState();
if (eventSource) { eventSource.close(); eventSource = null; }
if (logEventSource) { logEventSource.close(); logEventSource = null; }
} else if (token) {
Expand Down Expand Up @@ -693,6 +709,7 @@ function rememberSseEventId(event) {

function connectSSE(lastEventIdOverride) {
if (eventSource) eventSource.close();
cleanupConnectionState();

// In OIDC mode the reverse proxy provides auth; no query token needed.
let chatSseUrl = (token && !oidcProxyAuth)
Expand Down Expand Up @@ -846,6 +863,7 @@ function connectSSE(lastEventIdOverride) {
}
finalizeActivityGroup();
addMessage('assistant', data.content);
pruneOldMessages();
enableChatInput();
// Refresh thread list so new titles appear after first message
loadThreads();
Expand Down Expand Up @@ -1047,6 +1065,16 @@ function connectSSE(lastEventIdOverride) {
events.push({ type: evtType, data: data, ts: Date.now() });
// Cap per-job events to prevent memory leak
while (events.length > JOB_EVENTS_CAP) events.shift();
// Cap total tracked jobs — evict the one with the oldest last event
if (jobEvents.size > JOB_EVENTS_MAX_JOBS) {
let oldestKey = null, oldestTs = Infinity;
for (const [k, v] of jobEvents) {
if (k === jobId || k === currentJobId) continue; // never evict the active or viewed job
const lastTs = v.length > 0 ? v[v.length - 1].ts : 0;
if (lastTs < oldestTs) { oldestTs = lastTs; oldestKey = k; }
}
if (oldestKey) jobEvents.delete(oldestKey);
}
// If the Activity tab is currently visible for this job, refresh it
refreshActivityTab(jobId);
// Auto-refresh job list when on jobs tab (debounced)
Expand Down Expand Up @@ -1826,6 +1854,26 @@ function maybeInsertTimeSeparator(container, timestamp) {
container.appendChild(sep);
}

const MAX_DOM_MESSAGES = 200;

// Remove oldest messages/activity groups from the DOM when the chat container
// exceeds MAX_DOM_MESSAGES elements. Users can scroll up to trigger
// loadHistory() for older content. This prevents unbounded DOM growth during
// long sessions. Elements with data-streaming="true" are preserved to avoid
// breaking mid-stream responses.
function pruneOldMessages() {
const container = document.getElementById('chat-messages');
const items = container.querySelectorAll('.message, .activity-group, .time-separator');
if (items.length <= MAX_DOM_MESSAGES) return;
let removed = 0;
const target = items.length - MAX_DOM_MESSAGES;
for (let i = 0; i < items.length && removed < target; i++) {
if (items[i].getAttribute('data-streaming') === 'true') continue;
items[i].remove();
removed++;
}
}

function addMessage(role, content) {
const container = document.getElementById('chat-messages');
maybeInsertTimeSeparator(container);
Expand Down Expand Up @@ -2991,6 +3039,7 @@ function loadHistory(before) {

hasMore = data.has_more || false;
oldestTimestamp = data.oldest_timestamp || null;
if (!isPaginating) pruneOldMessages();
}).catch(() => {
// No history or no active thread
}).finally(() => {
Expand Down
18 changes: 16 additions & 2 deletions src/channels/web/sse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,15 @@ use crate::channels::web::types::AppEvent;
/// Prevents resource exhaustion from connection flooding.
pub const DEFAULT_MAX_CONNECTIONS: u64 = 100;

/// Default broadcast buffer size. Under heavy tool use each tool call generates
/// 5-10 SSE events; a small buffer causes slow clients to lag and reconnect in
/// a cascade. Configurable via `SSE_BROADCAST_BUFFER` env var.
const DEFAULT_BROADCAST_BUFFER: usize = 1024;

/// Upper bound for the broadcast buffer to prevent accidental OOM from
/// misconfigured env vars. 65 536 events * ~1 KB each ≈ 64 MB worst case.
const MAX_BROADCAST_BUFFER: usize = 65_536;

/// Envelope for broadcast events: carries an optional user scope.
///
/// `user_id = None` means the event is global (e.g. Heartbeat) and delivered
Expand Down Expand Up @@ -52,8 +61,13 @@ impl SseManager {

/// Create a new SSE manager with a custom connection limit.
pub fn with_max_connections(max_connections: u64) -> Self {
// Buffer 256 events; slow clients will miss events (acceptable for SSE with reconnect)
let (tx, _) = broadcast::channel(256);
let buffer_size = std::env::var("SSE_BROADCAST_BUFFER")
.ok()
.and_then(|v| v.parse::<usize>().ok())
.filter(|&n| n > 0)
.unwrap_or(DEFAULT_BROADCAST_BUFFER)
.min(MAX_BROADCAST_BUFFER);
let (tx, _) = broadcast::channel(buffer_size);
Self {
Comment thread
henrypark133 marked this conversation as resolved.
tx,
connection_count: Arc::new(AtomicU64::new(0)),
Expand Down
102 changes: 102 additions & 0 deletions tests/e2e/scenarios/test_dom_resource_limits.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
"""DOM and timer resource limit tests for issue #2406.

Verifies that the web UI does not exhaust browser resources during extended
sessions: DOM node count stays bounded, timers are cleaned up on reconnect,
and streaming messages survive pruning.
"""

from helpers import SEL


async def _wait_for_connected(page, *, timeout: int = 10000) -> None:
await page.wait_for_function(
"() => typeof sseHasConnectedBefore !== 'undefined' && sseHasConnectedBefore === true",
timeout=timeout,
)


async def test_dom_pruned_after_many_messages(page):
"""DOM stays bounded at MAX_DOM_MESSAGES after many insertions (#2406)."""
# Inject 250 messages directly (faster than round-tripping through LLM)
await page.evaluate("""() => {
for (let i = 0; i < 250; i++) {
addMessage(i % 2 === 0 ? 'user' : 'assistant', 'Message ' + i);
}
pruneOldMessages();
}""")

count = await page.locator(f"{SEL['chat_messages']} .message").count()
assert count <= 200, f"Expected <= 200 DOM messages after pruning, got {count}"
assert count >= 100, f"Expected at least 100 DOM messages (not over-pruned), got {count}"


async def test_no_timer_leak_across_reconnects(page):
"""Reconnect cycles do not accumulate leaked setInterval timers (#2406)."""
await _wait_for_connected(page, timeout=10000)

# Instrument: track net interval count via monkeypatching
await page.evaluate("""() => {
window.__testIntervalCount = 0;
const origSet = window.setInterval;
const origClear = window.clearInterval;
window.setInterval = function(...args) {
const id = origSet.apply(this, args);
window.__testIntervalCount++;
return id;
};
window.clearInterval = function(id) {
if (id != null) window.__testIntervalCount--;
return origClear.call(this, id);
};
}""")

baseline = await page.evaluate("window.__testIntervalCount")

# Force 5 reconnect cycles, triggering _streamDebounceTimer each time
# by dispatching a synthetic stream_chunk event before disconnecting.
for _ in range(5):
await page.evaluate("""() => {
// Trigger the stream_chunk handler which creates _streamDebounceTimer
const evt = new MessageEvent('stream_chunk', {
data: JSON.stringify({ content: 'test', thread_id: currentThreadId })
});
if (eventSource) eventSource.dispatchEvent(evt);
}""")
await page.evaluate("if (eventSource) eventSource.close()")
await page.evaluate("sseHasConnectedBefore = false; connectSSE()")
await _wait_for_connected(page, timeout=10000)

after = await page.evaluate("window.__testIntervalCount")
# Allow at most 1 net new interval (the gateway status poller or similar)
assert after <= baseline + 1, (
f"Interval leak detected: baseline={baseline}, after 5 reconnects={after}"
)
Comment thread
henrypark133 marked this conversation as resolved.


async def test_prune_preserves_streaming_message(page):
"""pruneOldMessages must not remove a message with data-streaming=true (#2406)."""
# Fill the DOM to just under the cap
await page.evaluate("""() => {
for (let i = 0; i < 199; i++) {
addMessage('assistant', 'msg ' + i);
}
}""")

# Mark the last assistant message as actively streaming
await page.evaluate("""() => {
const msgs = document.querySelectorAll('#chat-messages .message.assistant');
msgs[msgs.length - 1].setAttribute('data-streaming', 'true');
}""")

# Push over the cap and prune
await page.evaluate("""() => {
for (let i = 0; i < 10; i++) {
addMessage('user', 'overflow ' + i);
}
pruneOldMessages();
}""")

streaming_count = await page.locator('[data-streaming="true"]').count()
assert streaming_count == 1, (
f"Streaming message was pruned: expected 1 element with data-streaming, got {streaming_count}"
)
Loading