Skip to content

Commit ec3722a

Browse files
authored
feat(web): session list status indicators (attention + scheduled) (#699)
1 parent 7457a8f commit ec3722a

33 files changed

Lines changed: 910 additions & 37 deletions

hub/src/sse/sseManager.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -155,7 +155,7 @@ export class SSEManager {
155155
}
156156
}
157157

158-
if (event.type === 'message-received') {
158+
if (event.type === 'message-received' || event.type === 'scheduled-matured') {
159159
return connection.all || connection.sessionId === event.sessionId
160160
}
161161

hub/src/store/messageStore.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import type { Database } from 'bun:sqlite'
22

33
import type { StoredMessage } from './types'
4-
import { addMessage, cancelQueuedMessage, deleteQueuedMessageById, lookupQueuedMessage, getMessages, getFirstMessages, getDeliverableMessagesAfter, getMessagesByPosition, getUninvokedLocalMessages, getMatureScheduledMessages, getImmediateQueuedLocalMessages, markMessagesInvoked, mergeSessionMessages, type CancelQueuedMessageResult, type LookupQueuedMessageResult } from './messages'
4+
import { addMessage, cancelQueuedMessage, deleteQueuedMessageById, lookupQueuedMessage, getMessages, getFirstMessages, getDeliverableMessagesAfter, getMessagesByPosition, getUninvokedLocalMessages, getMatureScheduledMessages, getImmediateQueuedLocalMessages, countFutureScheduledBySessionIds, countFutureScheduledLocalMessages, markMessagesInvoked, mergeSessionMessages, type CancelQueuedMessageResult, type LookupQueuedMessageResult } from './messages'
55

66
export class MessageStore {
77
private readonly db: Database
@@ -42,6 +42,14 @@ export class MessageStore {
4242
return getImmediateQueuedLocalMessages(this.db, sessionId)
4343
}
4444

45+
countFutureScheduledLocalMessages(sessionId: string, now: number = Date.now()): number {
46+
return countFutureScheduledLocalMessages(this.db, sessionId, now)
47+
}
48+
49+
countFutureScheduledBySessionIds(sessionIds: string[], now: number = Date.now()): Map<string, number> {
50+
return countFutureScheduledBySessionIds(this.db, sessionIds, now)
51+
}
52+
4553
cancelQueuedMessage(sessionId: string, messageId: string): CancelQueuedMessageResult {
4654
return cancelQueuedMessage(this.db, sessionId, messageId)
4755
}

hub/src/store/messages.test.ts

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -287,3 +287,60 @@ describe('getDeliverableMessagesAfter: CLI backfill excludes future-scheduled ro
287287
expect(empty).toHaveLength(0)
288288
})
289289
})
290+
291+
describe('countFutureScheduledLocalMessages', () => {
292+
it('counts only future scheduled uninvoked local messages', () => {
293+
const store = makeStore()
294+
const session = makeSession(store, 'sched-count')
295+
const now = Date.now()
296+
297+
store.messages.addMessage(
298+
session.id,
299+
{ role: 'user', content: { type: 'text', text: 'immediate queued' } },
300+
'local-immediate'
301+
)
302+
store.messages.addMessage(
303+
session.id,
304+
{ role: 'user', content: { type: 'text', text: 'future scheduled' } },
305+
'local-future',
306+
now + 60_000
307+
)
308+
store.messages.addMessage(
309+
session.id,
310+
{ role: 'user', content: { type: 'text', text: 'mature scheduled' } },
311+
'local-mature',
312+
now - 1
313+
)
314+
315+
expect(store.messages.countFutureScheduledLocalMessages(session.id, now)).toBe(1)
316+
})
317+
318+
it('batch query returns counts keyed by session id', () => {
319+
const store = makeStore()
320+
const sessionA = makeSession(store, 'sched-batch-a')
321+
const sessionB = makeSession(store, 'sched-batch-b')
322+
const now = Date.now()
323+
324+
store.messages.addMessage(
325+
sessionA.id,
326+
{ role: 'user', content: { type: 'text', text: 'a1' } },
327+
'a-1',
328+
now + 60_000
329+
)
330+
store.messages.addMessage(
331+
sessionA.id,
332+
{ role: 'user', content: { type: 'text', text: 'a2' } },
333+
'a-2',
334+
now + 120_000
335+
)
336+
store.messages.addMessage(
337+
sessionB.id,
338+
{ role: 'user', content: { type: 'text', text: 'immediate' } },
339+
'b-1'
340+
)
341+
342+
const counts = store.messages.countFutureScheduledBySessionIds([sessionA.id, sessionB.id], now)
343+
expect(counts.get(sessionA.id)).toBe(2)
344+
expect(counts.get(sessionB.id)).toBeUndefined()
345+
})
346+
})

hub/src/store/messages.ts

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -233,6 +233,53 @@ export function getImmediateQueuedLocalMessages(
233233
return rows.map(toStoredMessage)
234234
}
235235

236+
/** Count uninvoked local messages scheduled for a future time (session list indicator). */
237+
export function countFutureScheduledLocalMessages(
238+
db: Database,
239+
sessionId: string,
240+
now: number
241+
): number {
242+
const row = db.prepare(`
243+
SELECT COUNT(*) AS count
244+
FROM messages
245+
WHERE session_id = ?
246+
AND invoked_at IS NULL
247+
AND local_id IS NOT NULL
248+
AND scheduled_at IS NOT NULL
249+
AND scheduled_at > ?
250+
`).get(sessionId, now) as { count: number } | undefined
251+
return row?.count ?? 0
252+
}
253+
254+
/** Batch variant for GET /sessions — one query for all session IDs in a namespace. */
255+
export function countFutureScheduledBySessionIds(
256+
db: Database,
257+
sessionIds: string[],
258+
now: number
259+
): Map<string, number> {
260+
const counts = new Map<string, number>()
261+
if (sessionIds.length === 0) {
262+
return counts
263+
}
264+
265+
const placeholders = sessionIds.map(() => '?').join(',')
266+
const rows = db.prepare(`
267+
SELECT session_id, COUNT(*) AS count
268+
FROM messages
269+
WHERE session_id IN (${placeholders})
270+
AND invoked_at IS NULL
271+
AND local_id IS NOT NULL
272+
AND scheduled_at IS NOT NULL
273+
AND scheduled_at > ?
274+
GROUP BY session_id
275+
`).all(...sessionIds, now) as { session_id: string; count: number }[]
276+
277+
for (const row of rows) {
278+
counts.set(row.session_id, row.count)
279+
}
280+
return counts
281+
}
282+
236283
export function getMaxSeq(db: Database, sessionId: string): number {
237284
const row = db.prepare(
238285
'SELECT COALESCE(MAX(seq), 0) AS maxSeq FROM messages WHERE session_id = ?'

hub/src/sync/messageService.test.ts

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -851,6 +851,59 @@ describe('MessageService.releaseMatureScheduledMessages', () => {
851851
expect(cliEmitted).toHaveLength(1)
852852
})
853853

854+
it('emits scheduled-matured once per session for web session-list refresh', async () => {
855+
const store = makeStore()
856+
const session = makeSession(store, 'release-sse')
857+
const publisher = makePublisher()
858+
const { io } = makeTrackingIo()
859+
860+
const now = Date.now()
861+
const past = now - 1000
862+
store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: 'one' } }, 'local-a', past)
863+
store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: 'two' } }, 'local-b', past)
864+
865+
const service = new MessageService(store, io, publisher as any)
866+
service.releaseMatureScheduledMessages(now)
867+
868+
const matured = publisher.events.filter((event) => event.type === 'scheduled-matured')
869+
expect(matured).toEqual([{ type: 'scheduled-matured', sessionId: session.id }])
870+
})
871+
872+
it('does NOT re-emit scheduled-matured on later ticks while CLI ack is pending', async () => {
873+
const store = makeStore()
874+
const session = makeSession(store, 'release-sse-no-repeat')
875+
const publisher = makePublisher()
876+
const { io } = makeTrackingIo()
877+
878+
const now = Date.now()
879+
const past = now - 1000
880+
store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: 'hi' } }, 'local-repeat', past)
881+
882+
const service = new MessageService(store, io, publisher as any)
883+
service.releaseMatureScheduledMessages(now)
884+
service.releaseMatureScheduledMessages(now + 60_000)
885+
886+
const matured = publisher.events.filter((event) => event.type === 'scheduled-matured')
887+
expect(matured).toHaveLength(1)
888+
})
889+
890+
it('emits scheduled-matured when first scan is long after scheduled_at', async () => {
891+
const store = makeStore()
892+
const session = makeSession(store, 'release-sse-late-scan')
893+
const publisher = makePublisher()
894+
const { io } = makeTrackingIo()
895+
896+
const now = Date.now()
897+
const past = now - 60_000
898+
store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: 'hi' } }, 'local-late', past)
899+
900+
const service = new MessageService(store, io, publisher as any)
901+
service.releaseMatureScheduledMessages(now)
902+
903+
const matured = publisher.events.filter((event) => event.type === 'scheduled-matured')
904+
expect(matured).toEqual([{ type: 'scheduled-matured', sessionId: session.id }])
905+
})
906+
854907
it('does NOT call markMessagesInvoked (pitfall #2 guard): message is re-emitted on next tick', async () => {
855908
const store = makeStore()
856909
const session = makeSession(store, 'release-no-mark')

hub/src/sync/messageService.ts

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@ function toVisibleDecryptedMessages(messages: StoredMessageForDelivery[]): Decry
2828
}
2929

3030
export class MessageService {
31+
/** One scheduled-matured SSE per localId per hub process (cleared on cancel/consume paths here). */
32+
private readonly scheduledMatureNotifiedLocalIds = new Set<string>()
33+
3134
constructor(
3235
private readonly store: Store,
3336
private readonly io: Server,
@@ -36,6 +39,12 @@ export class MessageService {
3639
) {
3740
}
3841

42+
private forgetScheduledMatureNotified(localIds: Iterable<string>): void {
43+
for (const localId of localIds) {
44+
this.scheduledMatureNotifiedLocalIds.delete(localId)
45+
}
46+
}
47+
3948
getMessages(sessionId: string, limit: number = 200): DecryptedMessage[] {
4049
const stored = this.store.messages.getMessages(sessionId, limit)
4150
return toVisibleDecryptedMessages(stored)
@@ -196,6 +205,7 @@ export class MessageService {
196205
const now = Date.now()
197206
if (scheduledAt !== null && scheduledAt > now) {
198207
this.store.messages.deleteQueuedMessageById(sessionId, resolvedId)
208+
this.forgetScheduledMatureNotified([localId])
199209
this.publisher.emit({
200210
type: 'message-cancelled',
201211
sessionId,
@@ -224,6 +234,7 @@ export class MessageService {
224234
const recheck = this.store.messages.lookupQueuedMessage(sessionId, resolvedId)
225235
if (recheck.status === 'invoked') {
226236
// CLI beat us — treat identically to Race-B (ack returned not-found).
237+
this.forgetScheduledMatureNotified([localId])
227238
this.publisher.emit({
228239
type: 'messages-consumed',
229240
sessionId,
@@ -233,6 +244,7 @@ export class MessageService {
233244
return recheck
234245
}
235246
// Row is gone (absent) — clean cancel.
247+
this.forgetScheduledMatureNotified([localId])
236248
this.publisher.emit({
237249
type: 'message-cancelled',
238250
sessionId,
@@ -257,6 +269,7 @@ export class MessageService {
257269
// DB write failed — let the HTTP 500 surface to the caller.
258270
throw err
259271
}
272+
this.forgetScheduledMatureNotified([localId])
260273
// Notify all SSE subscribers (other open tabs) that this queued row is now
261274
// invoked so they remove it from the floating bar. Without this emit, only
262275
// the tab that sent the DELETE request learns about the status change via the
@@ -282,6 +295,7 @@ export class MessageService {
282295

283296
// Phase 3: CLI confirmed removal. Now DELETE the DB row and broadcast SSE.
284297
this.store.messages.deleteQueuedMessageById(sessionId, resolvedId)
298+
this.forgetScheduledMatureNotified([localId])
285299
this.publisher.emit({
286300
type: 'message-cancelled',
287301
sessionId,
@@ -455,6 +469,7 @@ export class MessageService {
455469
.filter((id): id is string => typeof id === 'string')
456470
if (localIds.length === 0) return null
457471
this.store.messages.markMessagesInvoked(sessionId, localIds, invokedAt)
472+
this.forgetScheduledMatureNotified(localIds)
458473
this.publisher.emit({ type: 'messages-consumed', sessionId, localIds, invokedAt })
459474
return { localIds, invokedAt }
460475
}
@@ -476,7 +491,13 @@ export class MessageService {
476491
* expected behaviour. */
477492
releaseMatureScheduledMessages(now: number): void {
478493
const mature = this.store.messages.getMatureScheduledMessages(now)
494+
const maturedSessionIds = new Set<string>()
479495
for (const msg of mature) {
496+
const localId = msg.localId
497+
if (typeof localId === 'string' && !this.scheduledMatureNotifiedLocalIds.has(localId)) {
498+
this.scheduledMatureNotifiedLocalIds.add(localId)
499+
maturedSessionIds.add(msg.sessionId)
500+
}
480501
const update = {
481502
id: msg.id,
482503
seq: msg.seq,
@@ -497,5 +518,8 @@ export class MessageService {
497518
// NOTE: do NOT call markMessagesInvoked here (pitfall #2).
498519
// CLI ack (messages-consumed) will handle invoked_at stamping.
499520
}
521+
for (const sessionId of maturedSessionIds) {
522+
this.publisher.emit({ type: 'scheduled-matured', sessionId })
523+
}
500524
}
501525
}

hub/src/sync/syncEngine.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -176,6 +176,10 @@ export class SyncEngine {
176176
return this.sessionCache.getSessionsByNamespace(namespace)
177177
}
178178

179+
getFutureScheduledMessageCounts(sessionIds: string[], now: number = Date.now()): Map<string, number> {
180+
return this.store.messages.countFutureScheduledBySessionIds(sessionIds, now)
181+
}
182+
179183
getSession(sessionId: string): Session | undefined {
180184
return this.sessionCache.getSession(sessionId) ?? this.sessionCache.refreshSession(sessionId) ?? undefined
181185
}

hub/src/web/routes/sessions.ts

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -64,7 +64,7 @@ export function createSessionsRoutes(getSyncEngine: () => SyncEngine | null): Ho
6464
const getPendingCount = (s: Session) => s.agentState?.requests ? Object.keys(s.agentState.requests).length : 0
6565

6666
const namespace = c.get('namespace')
67-
const sessions = engine.getSessionsByNamespace(namespace)
67+
const sessionRecords = engine.getSessionsByNamespace(namespace)
6868
.sort((a, b) => {
6969
// Active sessions first
7070
if (a.active !== b.active) {
@@ -79,7 +79,14 @@ export function createSessionsRoutes(getSyncEngine: () => SyncEngine | null): Ho
7979
// Then by updatedAt
8080
return b.updatedAt - a.updatedAt
8181
})
82-
.map(toSessionSummary)
82+
const scheduledCounts = engine.getFutureScheduledMessageCounts(sessionRecords.map((session) => session.id))
83+
const sessions = sessionRecords.map((session) => {
84+
const summary = toSessionSummary(session)
85+
return {
86+
...summary,
87+
futureScheduledMessageCount: scheduledCounts.get(session.id) ?? 0
88+
}
89+
})
8390

8491
return c.json({ sessions })
8592
})

package.json

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,10 +18,11 @@
1818
"typecheck:cli": "cd cli && bun run typecheck",
1919
"typecheck:hub": "cd hub && bun run typecheck",
2020
"typecheck:web": "cd web && bun run typecheck",
21-
"test": "bun run test:cli && bun run test:hub && bun run test:web",
21+
"test": "bun run test:cli && bun run test:hub && bun run test:web && bun run test:shared",
2222
"test:cli": "cd cli && bun run test",
2323
"test:hub": "cd hub && bun run test",
2424
"test:web": "cd web && bun run test",
25+
"test:shared": "cd shared && bun run test",
2526
"clean-session": "bun run hub/scripts/cleanup-sessions.ts",
2627
"release-all": "cd cli && bun run release-all"
2728
},

shared/package.json

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,9 @@
1818
"./voice": "./src/voice.ts"
1919
},
2020
"sideEffects": false,
21+
"scripts": {
22+
"test": "bun test"
23+
},
2124
"dependencies": {
2225
"zod": "^4.2.1"
2326
}

0 commit comments

Comments
 (0)