1use std::collections::{HashMap, HashSet};
17use std::sync::atomic::{AtomicUsize, Ordering};
18use std::sync::{Arc, Mutex};
19
20use arc_swap::ArcSwap;
21use lattice_core::BufferId;
22use lattice_mode::{
23 ElementContent, ElementId, ModelineElement, ModelineElementUpdate, ModelineKey, ModelineRole,
24 ModelineService, Zone, modeline::ROLE_MODE_ITEM,
25};
26use lattice_runtime::EventBus;
27use tokio::sync::Notify;
28
29use crate::mcp::server::ServerState;
30
31pub const STATUS_ELEMENT: &str = "claude-code";
34
35pub type IdeBuffers = Arc<Mutex<HashSet<BufferId>>>;
38
39pub fn register_status_descriptor(svc: &ModelineService) {
42 svc.register(ModelineElement::new(
43 ElementId::new(STATUS_ELEMENT),
44 Zone::Right,
45 6,
46 ));
47}
48
49pub fn status_content(
60 state: &ServerState,
61 conns: usize,
62 project: &str,
63 reviews: usize,
64 mention: bool,
65) -> ElementContent {
66 if !state.running {
67 return ElementContent::default();
68 }
69 let glyph = if conns > 0 { '●' } else { '○' };
71 let port = state.port.map(|p| p.to_string()).unwrap_or_default();
72 let proj = if project.is_empty() {
73 String::new()
74 } else {
75 format!(" · {project}")
76 };
77 let conn = match conns {
78 0 => String::new(),
79 1 => " · 1 conn".to_string(),
80 n => format!(" · {n} conns"),
81 };
82 let review = match reviews {
85 0 => String::new(),
86 1 => " · ◆ review".to_string(),
87 n => format!(" · ◆ {n} reviews"),
88 };
89 let mention = if mention { " · @sent" } else { "" };
91 let text = format!("{glyph} claude{proj} :{port}{conn}{review}{mention}");
92 ElementContent::text(text, ModelineRole::new(ROLE_MODE_ITEM))
93}
94
95pub struct MentionState {
102 until: Mutex<Option<std::time::Instant>>,
103 changed: Arc<Notify>,
104}
105
106pub type MentionHandle = Arc<MentionState>;
108
109const MENTION_ECHO: std::time::Duration = std::time::Duration::from_secs(3);
111
112impl MentionState {
113 pub fn new(changed: Arc<Notify>) -> MentionHandle {
114 Arc::new(Self {
115 until: Mutex::new(None),
116 changed,
117 })
118 }
119
120 pub fn ping(&self) {
122 *self.until.lock().unwrap_or_else(|e| e.into_inner()) =
123 Some(std::time::Instant::now() + MENTION_ECHO);
124 self.changed.notify_one();
125 }
126
127 pub fn remaining(&self) -> Option<std::time::Duration> {
131 let until = *self.until.lock().unwrap_or_else(|e| e.into_inner());
132 until.and_then(|d| d.checked_duration_since(std::time::Instant::now()))
133 }
134}
135
136pub struct ReviewState {
144 count: Arc<AtomicUsize>,
145 changed: Arc<Notify>,
146}
147
148pub type ReviewHandle = Arc<ReviewState>;
150
151impl ReviewState {
152 pub fn new(changed: Arc<Notify>) -> ReviewHandle {
155 Arc::new(Self {
156 count: Arc::new(AtomicUsize::new(0)),
157 changed,
158 })
159 }
160
161 pub fn count_handle(&self) -> Arc<AtomicUsize> {
163 self.count.clone()
164 }
165
166 pub fn begin(&self) -> ReviewGuard {
168 self.count.fetch_add(1, Ordering::Relaxed);
169 self.changed.notify_one();
170 ReviewGuard {
171 count: self.count.clone(),
172 changed: self.changed.clone(),
173 }
174 }
175}
176
177pub struct ReviewGuard {
181 count: Arc<AtomicUsize>,
182 changed: Arc<Notify>,
183}
184
185impl Drop for ReviewGuard {
186 fn drop(&mut self) {
187 self.count.fetch_sub(1, Ordering::Relaxed);
188 self.changed.notify_one();
189 }
190}
191
192pub fn project_name(workspace_folders: &[String]) -> String {
196 workspace_folders
197 .first()
198 .map(|f| f.trim_end_matches('/'))
199 .filter(|f| !f.is_empty())
200 .and_then(|f| std::path::Path::new(f).file_name())
201 .map(|n| n.to_string_lossy().into_owned())
202 .unwrap_or_default()
203}
204
205pub fn spawn_status_publisher(
210 bus: Arc<EventBus>,
211 state: Arc<ArcSwap<ServerState>>,
212 conn_count: Arc<AtomicUsize>,
213 project: String,
216 review_count: Arc<AtomicUsize>,
219 mention: MentionHandle,
221 ide_buffers: IdeBuffers,
222 changed: Arc<Notify>,
223 rt: &tokio::runtime::Handle,
224) {
225 rt.spawn(async move {
226 let id = ElementId::new(STATUS_ELEMENT);
227 let mut last: HashMap<BufferId, ElementContent> = HashMap::new();
228 loop {
229 let mention_clear = mention.remaining();
232 let content = status_content(
233 &state.load(),
234 conn_count.load(Ordering::Relaxed),
235 &project,
236 review_count.load(Ordering::Relaxed),
237 mention_clear.is_some(),
238 );
239 let bufs: Vec<BufferId> = {
240 let g = ide_buffers.lock().unwrap_or_else(|e| e.into_inner());
241 g.iter().copied().collect()
242 };
243
244 for buf in &bufs {
246 if last.get(buf) != Some(&content) {
247 publish(&bus, &id, *buf, content.clone());
248 last.insert(*buf, content.clone());
249 }
250 }
251 let removed: Vec<BufferId> =
253 last.keys().filter(|b| !bufs.contains(b)).copied().collect();
254 for buf in removed {
255 publish(&bus, &id, buf, ElementContent::default());
256 last.remove(&buf);
257 }
258
259 match mention_clear {
263 Some(rem) => {
264 tokio::select! {
265 _ = changed.notified() => {}
266 _ = tokio::time::sleep(rem) => {}
267 }
268 }
269 None => changed.notified().await,
270 }
271 }
272 });
273}
274
275fn publish(bus: &EventBus, id: &ElementId, buf: BufferId, content: ElementContent) {
276 bus.publish_typed(ModelineElementUpdate {
277 key: ModelineKey::Buffer(buf),
278 id: id.clone(),
279 content,
280 });
281}
282
283#[cfg(test)]
284mod tests {
285 #![allow(clippy::unwrap_used)]
286 use super::*;
287
288 fn state(running: bool, port: Option<u16>) -> ServerState {
289 ServerState { running, port }
290 }
291
292 #[test]
293 fn stopped_server_has_empty_content() {
294 assert!(status_content(&state(false, None), 0, "lattice", 0, false).is_empty());
295 }
296
297 #[test]
298 fn running_server_shows_glyph_project_port_and_conn_count() {
299 let c = status_content(&state(true, Some(8123)), 0, "lattice", 0, false);
301 assert_eq!(c.plain(), "○ claude · lattice :8123");
302 let c1 = status_content(&state(true, Some(8123)), 1, "lattice", 0, false);
303 assert_eq!(c1.plain(), "● claude · lattice :8123 · 1 conn");
304 let c2 = status_content(&state(true, Some(8123)), 3, "lattice", 0, false);
305 assert_eq!(c2.plain(), "● claude · lattice :8123 · 3 conns");
306 }
307
308 #[test]
309 fn empty_project_is_omitted() {
310 let c = status_content(&state(true, Some(8123)), 1, "", 0, false);
311 assert_eq!(c.plain(), "● claude :8123 · 1 conn");
312 }
313
314 #[test]
315 fn pending_review_shows_badge() {
316 let c1 = status_content(&state(true, Some(8123)), 1, "lattice", 1, false);
317 assert_eq!(c1.plain(), "● claude · lattice :8123 · 1 conn · ◆ review");
318 let c2 = status_content(&state(true, Some(8123)), 1, "lattice", 2, false);
319 assert_eq!(
320 c2.plain(),
321 "● claude · lattice :8123 · 1 conn · ◆ 2 reviews"
322 );
323 }
324
325 #[test]
326 fn at_mention_echo_shows_when_active() {
327 let c = status_content(&state(true, Some(8123)), 1, "lattice", 0, true);
330 assert_eq!(c.plain(), "● claude · lattice :8123 · 1 conn · @sent");
331 let c2 = status_content(&state(true, Some(8123)), 1, "lattice", 1, true);
332 assert_eq!(
333 c2.plain(),
334 "● claude · lattice :8123 · 1 conn · ◆ review · @sent"
335 );
336 }
337
338 #[test]
339 fn mention_remaining_is_some_after_ping_then_none_after_window() {
340 let m = MentionState::new(Arc::new(Notify::new()));
341 assert!(m.remaining().is_none(), "inactive before any ping");
342 m.ping();
343 assert!(m.remaining().is_some(), "active right after ping");
344 }
345
346 #[test]
347 fn review_guard_tracks_count_and_wakes() {
348 let changed = Arc::new(Notify::new());
349 let review = ReviewState::new(changed);
350 let count = review.count_handle();
351 assert_eq!(count.load(Ordering::Relaxed), 0);
352 {
353 let _g = review.begin();
354 assert_eq!(count.load(Ordering::Relaxed), 1, "begin increments");
355 let _g2 = review.begin();
356 assert_eq!(count.load(Ordering::Relaxed), 2, "concurrent reviews count");
357 }
358 assert_eq!(count.load(Ordering::Relaxed), 0, "guards decrement on drop");
359 }
360
361 #[test]
362 fn project_name_takes_workspace_basename() {
363 assert_eq!(
364 project_name(&["/home/me/src/lattice".to_string()]),
365 "lattice"
366 );
367 assert_eq!(
368 project_name(&["/home/me/src/lattice/".to_string()]),
369 "lattice"
370 );
371 assert_eq!(project_name(&[]), "");
372 assert_eq!(project_name(&["".to_string()]), "");
373 }
374
375 #[tokio::test]
376 async fn publisher_pushes_content_for_a_registered_buffer() {
377 use std::time::Duration;
378
379 let bus = Arc::new(EventBus::new());
380 let server_state = Arc::new(ArcSwap::from_pointee(state(true, Some(9001))));
381 let conn_count = Arc::new(AtomicUsize::new(0));
382 let ide_buffers: IdeBuffers = Arc::new(Mutex::new(HashSet::new()));
383 let changed = Arc::new(Notify::new());
384
385 let buf = BufferId(7);
388 ide_buffers.lock().unwrap().insert(buf);
389 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ModelineElementUpdate>();
390 bus.subscribe_typed(tx);
391
392 let review_count = Arc::new(AtomicUsize::new(0));
393 let mention = MentionState::new(Arc::new(Notify::new()));
394 spawn_status_publisher(
395 bus.clone(),
396 server_state,
397 conn_count,
398 "lattice".to_string(),
399 review_count,
400 mention,
401 ide_buffers,
402 changed,
403 &tokio::runtime::Handle::current(),
404 );
405
406 let update = tokio::time::timeout(Duration::from_secs(2), rx.recv())
407 .await
408 .expect("an update within the timeout")
409 .expect("the publisher pushed an update");
410 assert_eq!(update.key, ModelineKey::Buffer(buf));
411 assert_eq!(update.id, ElementId::new(STATUS_ELEMENT));
412 assert_eq!(update.content.plain(), "○ claude · lattice :9001");
414 }
415}