mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-09-02 01:29:23 +00:00
* feat(ux): complete UX overhaul — design system, boot screen, onboarding, web polish Shared design system: CSS custom properties for spacing, typography, transitions, and color tokens used across web UI and boot screen. Boot screen: compact feature-tags line showing enabled subsystems (db, tools, routines, heartbeat, skills, sandbox, embeddings) at a glance. Downgrade startup info logs (libSQL, webhook, workspace seed) to debug level since the boot screen now covers this. Onboarding wizard: model picker with live API fetch, provider-aware auth flow, improved error recovery and progress display. Web UI: ARIA attributes, welcome card, streaming debounce, connection status banner, skeleton loaders, send cooldown. CLI: doctor command enhancements, status command cleanup, REPL banner consolidation, shared fmt module. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(ux): Apple-level design refinements — spring physics, glass morphism, chat polish Merge staging theme support (dark/light/system toggle) and layer UX polish on top: spring-physics motion, glass morphism depth, chat experience improvements, and responsive mobile refinements. Design system: - Restore and extend design token system (spacing, typography, timing, easing) with legacy aliases for theme compatibility - Add shadow tiers, accent glow, glass morphism, spring easing tokens - Tokens defined in both dark (:root) and light ([data-theme="light"]) Micro-interactions (Phase 2): - Spring-overshoot message entry animation (slideUp) - Spring-scale button press on all interactive buttons - Tab crossfade animation, tool card smooth accordion (max-height) - Modal scale(0.95) + blur(8px) entry, toast spring slide - Sidebar width crossfade, card hover lift Visual depth (Phase 3): - Tab bar glass morphism + surface highlight + sliding indicator - Active tab accent background pill - Assistant message accent left border, user message bubble tail - Floating input area (rounded + shadow + margin) Chat polish (Phase 4): - Smooth streaming cursor (cursorPulse), message hover timestamps - Time separators (Today/Yesterday/date) - Textarea smooth auto-expand, send button glow Settings & forms (Phase 5): - iOS-style toggle switches for boolean settings - Input focus glow, save feedback spring animation - Welcome card with gradient background + proper spacing - Sticky settings group headers with glass backdrop Accessibility & mobile (Phase 6): - Animated focus ring, prefers-reduced-motion global kill-switch - Touch target audit (44px min), mobile bottom-sheet modals - Mobile bottom tab bar, toast redesign (icon + border + countdown) - Thread hover translateX, badge in_progress pulse Bug fixes: - Gateway/TEE popover z-index (tab-bar z-index: 200, popovers 500) - Connection lost banner as fixed top bar instead of flex child - Sidebar collapse keeps toggle + new thread buttons visible - Downgrade noisy startup logs (db, webhook, vector) to debug - Remove green dot pulse animation on connected status - Deduplicate confirm-modal in HTML, add tab-indicator div Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(web): mobile layout improvements — sidebar toggle, settings drill-down, tab bar polish - Fix mobile sidebar toggle: use expanded-mobile class instead of collapsed, add backdrop overlay, auto-close on thread select, outside-click dismiss - Settings: replace cramped horizontal tabs with drill-down navigation (category list → detail view → back button) - Bottom tab bar: add glass morphism, hide theme toggle, flip tab indicator to top edge - Keep thread toggle button visible in collapsed 36px sidebar strip Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(repl): interactive approval selector and transient status lines - Replace ASCII-art approval box with clean horizontal rule card - Add inquire-based interactive selector for tool approvals (↑↓ + Enter) - Selector runs directly from send_status via spawn_blocking, with stdin_locked flag to prevent readline from competing for stdin - Transient thinking/tool-started lines: each replaces the previous, all erased before final output (no clutter left in scrollback) - Esc in selector sends denial so agent never gets stuck Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: widen TurnCost token fields to u64 and remove unused variable - Change input_tokens/output_tokens from u32 to u64 in StatusUpdate::TurnCost, SseEvent::TurnCost, and the thread_ops emit site to avoid truncation on large conversations - Remove unused _routine_engine_for_loop binding in agent_loop.rs Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * chore: reduce startup log noise — demote info to debug Demote routine startup messages (builder, WASM tools, tunnel, WASM channels) from info to debug so the default log output stays clean. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix(web): allow CDN scripts in CSP connect-src directive Add cdn.jsdelivr.net and cdnjs.cloudflare.com to connect-src so the browser can fetch marked.js and DOMPurify without CSP violations. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * style: fix cargo fmt in repl.rs Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix(web): gate turn_cost SSE handler on current thread Prevents cost badge from attaching to the wrong message when switching threads or receiving events from background threads. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * ci: retrigger CI * fix: add missing extension_manager to webhook EngineContext The webhook trigger path added in #736 was missing the extension_manager field introduced by #1453. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * chore: ignore RUSTSEC-2026-0049 rustls-webpki CRL advisory Low impact — requires compromised CA to exploit. Tracked for upstream rustls-webpki upgrade. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix(routines): use fields.join for cron normalization Use split_whitespace fields instead of re-trimming the original string to avoid preserving extra internal whitespace in cron expressions. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * feat(repl): Apple-style approval card — clean vertical flow - Drop verbose tool description (the command IS the decision surface) - Unified vertical pipe layout: ◆ header → │ params → │ selector - Selector options show keyboard shortcuts inline: Approve (y) - Compact help message, answered state uses └ to close the flow - No horizontal rules, no blank-line padding — just breathing room Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * refactor(repl): replace inquire with crossterm for approval selector Drop the inquire dependency (which pulled in crossterm 0.25, duplicating the existing 0.28). The 3-option approval selector is now built directly with crossterm raw mode — same UX, zero new dependencies. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * chore(deps): upgrade crossterm 0.28 → 0.29, eliminate duplication termimad (via crokey) uses crossterm 0.29. Upgrading our direct dependency from 0.28 to 0.29 collapses to a single crossterm version in the dependency tree. Also migrated termimad::crossterm:: references to the direct crossterm import. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address review comments — box_top off-by-one, smart_truncate overflow, mobile theme toggle - Fix box_top() fill calculation: was off-by-one, producing boxes 1 char too wide (fmt.rs) - Fix smart_truncate(): account for "..." in the budget so output never exceeds max_chars (repl.rs) - Move theme toggle to settings sidebar on mobile instead of display:none, so mobile users can still switch themes (style.css, index.html, app.js) Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * style: cargo fmt repl.rs Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> * fix: address review — retry duplication, CSP connect-src, deny color - Remove failed message before retry to prevent duplicate user messages - Revert connect-src to 'self' — CDN hosts only need script-src - Use red for Deny confirmation in REPL approval selector Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]> --------- Co-authored-by: Claude Opus 4.6 (1M context) <[email protected]>
298 lines
10 KiB
Rust
298 lines
10 KiB
Rust
//! SSE connection manager for broadcasting events to browser tabs.
|
|
|
|
use std::convert::Infallible;
|
|
use std::sync::Arc;
|
|
use std::sync::atomic::{AtomicU64, Ordering};
|
|
use std::time::Duration;
|
|
|
|
use axum::response::sse::{Event, KeepAlive, Sse};
|
|
use futures::Stream;
|
|
use tokio::sync::broadcast;
|
|
use tokio_stream::StreamExt;
|
|
use tokio_stream::wrappers::BroadcastStream;
|
|
|
|
use crate::channels::web::types::SseEvent;
|
|
|
|
/// Maximum number of concurrent SSE/WebSocket connections.
|
|
/// Prevents resource exhaustion from connection flooding.
|
|
const MAX_CONNECTIONS: u64 = 100;
|
|
|
|
/// Manages SSE broadcast to all connected browser tabs.
|
|
pub struct SseManager {
|
|
tx: broadcast::Sender<SseEvent>,
|
|
connection_count: Arc<AtomicU64>,
|
|
max_connections: u64,
|
|
}
|
|
|
|
impl SseManager {
|
|
/// Create a new SSE manager.
|
|
pub fn new() -> Self {
|
|
// Buffer 256 events; slow clients will miss events (acceptable for SSE with reconnect)
|
|
let (tx, _) = broadcast::channel(256);
|
|
Self {
|
|
tx,
|
|
connection_count: Arc::new(AtomicU64::new(0)),
|
|
max_connections: MAX_CONNECTIONS,
|
|
}
|
|
}
|
|
|
|
/// Create an SSE manager that reuses an existing broadcast sender.
|
|
///
|
|
/// This preserves the broadcast channel across `rebuild_state` calls so
|
|
/// that sender handles captured by other components remain valid.
|
|
///
|
|
/// **Important:** The connection counter is reset to zero. This method must
|
|
/// only be called before the server starts accepting connections (i.e.,
|
|
/// during startup wiring). Calling it after connections are established
|
|
/// will break connection tracking and allow exceeding `MAX_CONNECTIONS`.
|
|
pub fn from_sender(tx: broadcast::Sender<SseEvent>) -> Self {
|
|
Self {
|
|
tx,
|
|
connection_count: Arc::new(AtomicU64::new(0)),
|
|
max_connections: MAX_CONNECTIONS,
|
|
}
|
|
}
|
|
|
|
/// Broadcast an event to all connected clients.
|
|
pub fn broadcast(&self, event: SseEvent) {
|
|
// Ignore send errors (no receivers is fine)
|
|
let _ = self.tx.send(event);
|
|
}
|
|
|
|
/// Get a clone of the broadcast sender for use by other components.
|
|
pub fn sender(&self) -> broadcast::Sender<SseEvent> {
|
|
self.tx.clone()
|
|
}
|
|
|
|
/// Get current number of active connections.
|
|
pub fn connection_count(&self) -> u64 {
|
|
self.connection_count.load(Ordering::Relaxed)
|
|
}
|
|
|
|
/// Create a raw broadcast subscription for non-SSE consumers (e.g. WebSocket).
|
|
///
|
|
/// Returns a stream of `SseEvent` values and increments/decrements the
|
|
/// connection counter on creation/drop, just like `subscribe()` does for SSE.
|
|
///
|
|
/// Returns `None` if the maximum connection limit has been reached.
|
|
pub fn subscribe_raw(&self) -> Option<impl Stream<Item = SseEvent> + Send + 'static + use<>> {
|
|
// Atomically increment only if below the limit. This prevents
|
|
// concurrent callers from overshooting max_connections.
|
|
let counter = Arc::clone(&self.connection_count);
|
|
let max = self.max_connections;
|
|
counter
|
|
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
|
|
if current < max {
|
|
Some(current + 1)
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
.ok()?;
|
|
let rx = self.tx.subscribe();
|
|
|
|
let stream = BroadcastStream::new(rx).filter_map(|result| result.ok());
|
|
|
|
Some(CountedStream {
|
|
inner: stream,
|
|
counter,
|
|
})
|
|
}
|
|
|
|
/// Create a new SSE stream for a client connection.
|
|
///
|
|
/// Returns `None` if the maximum connection limit has been reached.
|
|
pub fn subscribe(
|
|
&self,
|
|
) -> Option<Sse<impl Stream<Item = Result<Event, Infallible>> + Send + 'static + use<>>> {
|
|
// Atomically increment only if below the limit.
|
|
let counter = Arc::clone(&self.connection_count);
|
|
let max = self.max_connections;
|
|
counter
|
|
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
|
|
if current < max {
|
|
Some(current + 1)
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
.ok()?;
|
|
let rx = self.tx.subscribe();
|
|
|
|
let stream = BroadcastStream::new(rx)
|
|
.filter_map(|result| result.ok())
|
|
.map(|event| {
|
|
let data = serde_json::to_string(&event).unwrap_or_default();
|
|
let event_type = match &event {
|
|
SseEvent::Response { .. } => "response",
|
|
SseEvent::Thinking { .. } => "thinking",
|
|
SseEvent::ToolStarted { .. } => "tool_started",
|
|
SseEvent::ToolCompleted { .. } => "tool_completed",
|
|
SseEvent::ToolResult { .. } => "tool_result",
|
|
SseEvent::StreamChunk { .. } => "stream_chunk",
|
|
SseEvent::Status { .. } => "status",
|
|
SseEvent::ApprovalNeeded { .. } => "approval_needed",
|
|
SseEvent::AuthRequired { .. } => "auth_required",
|
|
SseEvent::AuthCompleted { .. } => "auth_completed",
|
|
SseEvent::Error { .. } => "error",
|
|
SseEvent::JobStarted { .. } => "job_started",
|
|
SseEvent::JobMessage { .. } => "job_message",
|
|
SseEvent::JobToolUse { .. } => "job_tool_use",
|
|
SseEvent::JobToolResult { .. } => "job_tool_result",
|
|
SseEvent::JobStatus { .. } => "job_status",
|
|
SseEvent::JobResult { .. } => "job_result",
|
|
SseEvent::Heartbeat => "heartbeat",
|
|
SseEvent::ImageGenerated { .. } => "image_generated",
|
|
SseEvent::Suggestions { .. } => "suggestions",
|
|
SseEvent::TurnCost { .. } => "turn_cost",
|
|
SseEvent::ExtensionStatus { .. } => "extension_status",
|
|
};
|
|
Ok(Event::default().event(event_type).data(data))
|
|
});
|
|
|
|
// Wrap in a stream that decrements on drop
|
|
let counted_stream = CountedStream {
|
|
inner: stream,
|
|
counter,
|
|
};
|
|
|
|
Some(
|
|
Sse::new(counted_stream)
|
|
.keep_alive(KeepAlive::new().interval(Duration::from_secs(30)).text("")),
|
|
)
|
|
}
|
|
}
|
|
|
|
impl Default for SseManager {
|
|
fn default() -> Self {
|
|
Self::new()
|
|
}
|
|
}
|
|
|
|
/// Stream wrapper that decrements connection count on drop.
|
|
///
|
|
/// When the SSE client disconnects, this stream is dropped
|
|
/// and the counter is decremented.
|
|
struct CountedStream<S> {
|
|
inner: S,
|
|
counter: Arc<AtomicU64>,
|
|
}
|
|
|
|
impl<S: Stream + Unpin> Stream for CountedStream<S> {
|
|
type Item = S::Item;
|
|
|
|
fn poll_next(
|
|
mut self: std::pin::Pin<&mut Self>,
|
|
cx: &mut std::task::Context<'_>,
|
|
) -> std::task::Poll<Option<Self::Item>> {
|
|
std::pin::Pin::new(&mut self.inner).poll_next(cx)
|
|
}
|
|
}
|
|
|
|
impl<S> Drop for CountedStream<S> {
|
|
fn drop(&mut self) {
|
|
self.counter.fetch_sub(1, Ordering::Relaxed);
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn test_sse_manager_creation() {
|
|
let manager = SseManager::new();
|
|
assert_eq!(manager.connection_count(), 0);
|
|
}
|
|
|
|
#[test]
|
|
fn test_broadcast_without_receivers() {
|
|
let manager = SseManager::new();
|
|
// Should not panic even with no receivers
|
|
manager.broadcast(SseEvent::Heartbeat);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_broadcast_to_receiver() {
|
|
let manager = SseManager::new();
|
|
let mut rx = BroadcastStream::new(manager.tx.subscribe());
|
|
|
|
manager.broadcast(SseEvent::Status {
|
|
message: "test".to_string(),
|
|
thread_id: None,
|
|
});
|
|
|
|
let event = rx.next().await;
|
|
assert!(event.is_some());
|
|
let event = event.unwrap().unwrap();
|
|
match event {
|
|
SseEvent::Status { message, .. } => assert_eq!(message, "test"),
|
|
_ => panic!("unexpected event type"),
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_subscribe_raw_receives_events() {
|
|
let manager = SseManager::new();
|
|
let mut stream = Box::pin(manager.subscribe_raw().expect("should subscribe"));
|
|
|
|
assert_eq!(manager.connection_count(), 1);
|
|
|
|
manager.broadcast(SseEvent::Thinking {
|
|
message: "working".to_string(),
|
|
thread_id: None,
|
|
});
|
|
|
|
let event = stream.next().await.unwrap();
|
|
match event {
|
|
SseEvent::Thinking { message, .. } => assert_eq!(message, "working"),
|
|
_ => panic!("Expected Thinking event"),
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_subscribe_raw_decrements_on_drop() {
|
|
let manager = SseManager::new();
|
|
{
|
|
let _stream = Box::pin(manager.subscribe_raw().expect("should subscribe"));
|
|
assert_eq!(manager.connection_count(), 1);
|
|
}
|
|
// Stream dropped, counter should decrement
|
|
assert_eq!(manager.connection_count(), 0);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_subscribe_raw_multiple_subscribers() {
|
|
let manager = SseManager::new();
|
|
let mut s1 = Box::pin(manager.subscribe_raw().expect("should subscribe"));
|
|
let mut s2 = Box::pin(manager.subscribe_raw().expect("should subscribe"));
|
|
assert_eq!(manager.connection_count(), 2);
|
|
|
|
manager.broadcast(SseEvent::Heartbeat);
|
|
|
|
let e1 = s1.next().await.unwrap();
|
|
let e2 = s2.next().await.unwrap();
|
|
assert!(matches!(e1, SseEvent::Heartbeat));
|
|
assert!(matches!(e2, SseEvent::Heartbeat));
|
|
|
|
drop(s1);
|
|
assert_eq!(manager.connection_count(), 1);
|
|
drop(s2);
|
|
assert_eq!(manager.connection_count(), 0);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_subscribe_raw_rejects_over_limit() {
|
|
let mut manager = SseManager::new();
|
|
manager.max_connections = 2; // Low limit for testing
|
|
|
|
let _s1 = Box::pin(manager.subscribe_raw().expect("first should succeed"));
|
|
let _s2 = Box::pin(manager.subscribe_raw().expect("second should succeed"));
|
|
assert_eq!(manager.connection_count(), 2);
|
|
|
|
// Third should be rejected
|
|
assert!(manager.subscribe_raw().is_none());
|
|
assert!(manager.subscribe().is_none());
|
|
}
|
|
}
|