feat(server): persist activity history on the server and expose it to the dashboard UI via /api/activity instead of relying on fragile localStorage
This commit is contained in:
1 parent
bf7bf48163
commit
ae0ea9ab22
4 files changed
+83
-41
No files matched your search
@@ -895,23 +895,17 @@
|
|||||||
// --- WebSocket Activity Feed ---
|
// --- WebSocket Activity Feed ---
|
||||||
const MAX_ACTIVITY_HISTORY = 100;
|
const MAX_ACTIVITY_HISTORY = 100;
|
||||||
|
|
||||||
function saveActivityHistory(time, htmlData) {
|
async function loadActivityHistory() {
|
||||||
try {
|
|
||||||
let history = JSON.parse(localStorage.getItem('activityHistory') || '[]');
|
|
||||||
history.push({ time, data: htmlData });
|
|
||||||
if (history.length > MAX_ACTIVITY_HISTORY) history = history.slice(-MAX_ACTIVITY_HISTORY);
|
|
||||||
localStorage.setItem('activityHistory', JSON.stringify(history));
|
|
||||||
} catch(e) {}
|
|
||||||
}
|
|
||||||
|
|
||||||
function loadActivityHistory() {
|
|
||||||
try {
|
try {
|
||||||
|
const response = await fetch('/api/activity');
|
||||||
|
const history = await response.json();
|
||||||
const feed = document.getElementById('activity-feed');
|
const feed = document.getElementById('activity-feed');
|
||||||
let history = JSON.parse(localStorage.getItem('activityHistory') || '[]');
|
feed.innerHTML = '';
|
||||||
history.forEach(item => {
|
history.forEach(item => {
|
||||||
const div = document.createElement('div');
|
const div = document.createElement('div');
|
||||||
div.className = 'feed-entry';
|
div.className = 'feed-entry';
|
||||||
div.innerHTML = `<span class="time">[${item.time}]</span> ${item.data}`;
|
const timeStr = new Date(item.time).toLocaleTimeString([], {hour: '2-digit', minute:'2-digit', second:'2-digit'});
|
||||||
|
div.innerHTML = `<span class="time">[${timeStr}]</span> ${item.message || item.data}`;
|
||||||
feed.appendChild(div);
|
feed.appendChild(div);
|
||||||
});
|
});
|
||||||
if (history.length > 0) {
|
if (history.length > 0) {
|
||||||
@@ -931,12 +925,10 @@
|
|||||||
if (data.type === 'activity') {
|
if (data.type === 'activity') {
|
||||||
const div = document.createElement('div');
|
const div = document.createElement('div');
|
||||||
div.className = 'feed-entry';
|
div.className = 'feed-entry';
|
||||||
const time = new Date().toLocaleTimeString();
|
const timeStr = new Date(data.data.time).toLocaleTimeString([], {hour: '2-digit', minute:'2-digit', second:'2-digit'});
|
||||||
div.innerHTML = `<span class="time">[${time}]</span> ${data.data}`;
|
div.innerHTML = `<span class="time">[${timeStr}]</span> ${data.data.message || data.data.data || data.data}`;
|
||||||
feed.appendChild(div);
|
feed.appendChild(div);
|
||||||
|
|
||||||
saveActivityHistory(time, data.data);
|
|
||||||
|
|
||||||
// Auto-scroll logic
|
// Auto-scroll logic
|
||||||
const isScrolledToBottom = feed.scrollHeight - feed.clientHeight <= feed.scrollTop + 20;
|
const isScrolledToBottom = feed.scrollHeight - feed.clientHeight <= feed.scrollTop + 20;
|
||||||
if (isScrolledToBottom) {
|
if (isScrolledToBottom) {
|
||||||
|
|||||||
+31
-24
@@ -225,6 +225,23 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
|||||||
next_id: AtomicUsize::new(1),
|
next_id: AtomicUsize::new(1),
|
||||||
});
|
});
|
||||||
|
|
||||||
|
let app_state_clone = Arc::clone(&app_state);
|
||||||
|
let mut rx = state.activity_tx.subscribe();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
while let Ok(msg) = rx.recv().await {
|
||||||
|
let senders: Vec<_> = app_state_clone
|
||||||
|
.clients
|
||||||
|
.read()
|
||||||
|
.unwrap_or_else(|e| e.into_inner())
|
||||||
|
.values()
|
||||||
|
.cloned()
|
||||||
|
.collect();
|
||||||
|
for client_tx in senders {
|
||||||
|
let _ = client_tx.try_send(msg.clone());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
let app = Router::new()
|
let app = Router::new()
|
||||||
.route(
|
.route(
|
||||||
"/api/version",
|
"/api/version",
|
||||||
@@ -350,6 +367,18 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
|||||||
}
|
}
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
.route(
|
||||||
|
"/api/activity",
|
||||||
|
get({
|
||||||
|
let state_clone = app_state.handler.state.clone();
|
||||||
|
move || async move {
|
||||||
|
let activities = state_clone.recent_activities.read_with(|a| {
|
||||||
|
a.iter().cloned().collect::<Vec<_>>()
|
||||||
|
});
|
||||||
|
axum::Json(serde_json::json!(activities))
|
||||||
|
}
|
||||||
|
}),
|
||||||
|
)
|
||||||
.route(
|
.route(
|
||||||
"/api/stats",
|
"/api/stats",
|
||||||
get({
|
get({
|
||||||
@@ -491,30 +520,7 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, client_type: Str
|
|||||||
.and_then(|p| p.get("name"))
|
.and_then(|p| p.get("name"))
|
||||||
.and_then(|n| n.as_str())
|
.and_then(|n| n.as_str())
|
||||||
.unwrap_or("unknown_tool");
|
.unwrap_or("unknown_tool");
|
||||||
let activity_msg = format!("Agent executed tool: {}", name);
|
handler.state.broadcast_activity(&format!("Agent executed tool: {}", name));
|
||||||
|
|
||||||
let event = serde_json::json!({
|
|
||||||
"type": "activity",
|
|
||||||
"data": activity_msg
|
|
||||||
});
|
|
||||||
|
|
||||||
let senders: Vec<_> = state_clone
|
|
||||||
.clients
|
|
||||||
.read()
|
|
||||||
.unwrap_or_else(|e| e.into_inner())
|
|
||||||
.iter()
|
|
||||||
.filter_map(|(id, tx)| {
|
|
||||||
if id != &session_id_clone {
|
|
||||||
Some(tx.clone())
|
|
||||||
} else {
|
|
||||||
None
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.collect();
|
|
||||||
|
|
||||||
for client_tx in senders {
|
|
||||||
let _ = client_tx.try_send(event.to_string());
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
} // End if proxy
|
} // End if proxy
|
||||||
|
|
||||||
@@ -840,6 +846,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
|
|||||||
tech_debts: Store::new("tech_debts", db.clone()),
|
tech_debts: Store::new("tech_debts", db.clone()),
|
||||||
gates: Store::new("gates", db.clone()),
|
gates: Store::new("gates", db.clone()),
|
||||||
context_workspaces: Store::new("context_workspaces", db.clone()),
|
context_workspaces: Store::new("context_workspaces", db.clone()),
|
||||||
|
recent_activities: Store::new("recent_activities", db.clone()),
|
||||||
activity_tx: tokio::sync::broadcast::channel(100).0,
|
activity_tx: tokio::sync::broadcast::channel(100).0,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
+19
-1
@@ -27,6 +27,7 @@ pub struct MemoryState {
|
|||||||
pub tech_debts: Store<Vec<TechDebt>>,
|
pub tech_debts: Store<Vec<TechDebt>>,
|
||||||
pub gates: Store<Vec<GateRecord>>,
|
pub gates: Store<Vec<GateRecord>>,
|
||||||
pub context_workspaces: Store<Vec<ContextWorkspace>>,
|
pub context_workspaces: Store<Vec<ContextWorkspace>>,
|
||||||
|
pub recent_activities: Store<std::collections::VecDeque<serde_json::Value>>,
|
||||||
pub activity_tx: tokio::sync::broadcast::Sender<String>,
|
pub activity_tx: tokio::sync::broadcast::Sender<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -37,9 +38,26 @@ impl MemoryState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn broadcast_activity(&self, message: &str) {
|
pub fn broadcast_activity(&self, message: &str) {
|
||||||
|
let time = std::time::SystemTime::now()
|
||||||
|
.duration_since(std::time::UNIX_EPOCH)
|
||||||
|
.unwrap_or_default()
|
||||||
|
.as_millis() as u64;
|
||||||
|
|
||||||
|
let item = serde_json::json!({
|
||||||
|
"time": time,
|
||||||
|
"message": message
|
||||||
|
});
|
||||||
|
|
||||||
|
self.recent_activities.modify(|activities| {
|
||||||
|
activities.push_back(item.clone());
|
||||||
|
if activities.len() > 100 {
|
||||||
|
activities.pop_front();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
let payload = serde_json::json!({
|
let payload = serde_json::json!({
|
||||||
"type": "activity",
|
"type": "activity",
|
||||||
"data": message
|
"data": item
|
||||||
})
|
})
|
||||||
.to_string();
|
.to_string();
|
||||||
let _ = self.activity_tx.send(payload);
|
let _ = self.activity_tx.send(payload);
|
||||||
|
|||||||
@@ -0,0 +1,25 @@
|
|||||||
|
--- server/src/main.rs
|
||||||
|
+++ server/src/main.rs
|
||||||
|
@@ -225,6 +225,20 @@
|
||||||
|
next_id: AtomicUsize::new(1),
|
||||||
|
});
|
||||||
|
|
||||||
|
+ let app_state_clone = Arc::clone(&app_state);
|
||||||
|
+ let mut rx = state.activity_tx.subscribe();
|
||||||
|
+ tokio::spawn(async move {
|
||||||
|
+ while let Ok(msg) = rx.recv().await {
|
||||||
|
+ let senders: Vec<_> = app_state_clone
|
||||||
|
+ .clients
|
||||||
|
+ .read()
|
||||||
|
+ .unwrap_or_else(|e| e.into_inner())
|
||||||
|
+ .values()
|
||||||
|
+ .cloned()
|
||||||
|
+ .collect();
|
||||||
|
+ for client_tx in senders {
|
||||||
|
+ let _ = client_tx.try_send(msg.clone());
|
||||||
|
+ }
|
||||||
|
+ }
|
||||||
|
+ });
|
||||||
|
+
|
||||||
|
let app = Router::new()
|
||||||
|
.route(
|
||||||
Reference in new issue
Block a user