mirror of
https://github.com/wu736139669/hapi.git
synced 2026-08-05 06:24:37 +00:00
fix(codex): stabilize goal status UI events (#652)
This commit is contained in:
@@ -0,0 +1,67 @@
|
||||
import { describe, expect, it } from 'bun:test'
|
||||
import { Store, type StoredSession } from '../../../store'
|
||||
import type { SyncEvent } from '../../../sync/syncEngine'
|
||||
import type { CliSocketWithData } from '../../socketTypes'
|
||||
import { registerSessionHandlers } from './sessionHandlers'
|
||||
|
||||
class FakeSocket {
|
||||
readonly roomEvents: Array<{ room: string; event: string; data: unknown }> = []
|
||||
private readonly handlers = new Map<string, (data: unknown) => void>()
|
||||
|
||||
on(event: string, handler: (data: unknown) => void): this {
|
||||
this.handlers.set(event, handler)
|
||||
return this
|
||||
}
|
||||
|
||||
to(room: string): { emit: (event: string, data: unknown) => void } {
|
||||
return {
|
||||
emit: (event: string, data: unknown) => {
|
||||
this.roomEvents.push({ room, event, data })
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
trigger(event: string, data: unknown): void {
|
||||
this.handlers.get(event)?.(data)
|
||||
}
|
||||
}
|
||||
|
||||
function redundantGoalStatusContent(message: string): unknown {
|
||||
return {
|
||||
role: 'agent',
|
||||
content: {
|
||||
id: `event-${message}`,
|
||||
type: 'event',
|
||||
data: { type: 'message', message }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
describe('cli session handlers', () => {
|
||||
it('drops redundant goal status events before persistence and broadcast', () => {
|
||||
const store = new Store(':memory:')
|
||||
const session = store.sessions.getOrCreateSession('goal-status-session', {}, null, 'default')
|
||||
const socket = new FakeSocket()
|
||||
const webEvents: SyncEvent[] = []
|
||||
|
||||
registerSessionHandlers(socket as unknown as CliSocketWithData, {
|
||||
store,
|
||||
resolveSessionAccess: () => ({ ok: true, value: session as StoredSession }),
|
||||
emitAccessError: () => {
|
||||
throw new Error('unexpected access error')
|
||||
},
|
||||
onWebappEvent: (event) => {
|
||||
webEvents.push(event)
|
||||
}
|
||||
})
|
||||
|
||||
socket.trigger('message', {
|
||||
sid: session.id,
|
||||
message: redundantGoalStatusContent('Goal active · 8016 tokens')
|
||||
})
|
||||
|
||||
expect(store.messages.getMessages(session.id)).toHaveLength(0)
|
||||
expect(socket.roomEvents).toHaveLength(0)
|
||||
expect(webEvents).toHaveLength(0)
|
||||
})
|
||||
})
|
||||
@@ -2,6 +2,7 @@ import type { ClientToServerEvents } from '@hapi/protocol'
|
||||
import { z } from 'zod'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import type { CodexCollaborationMode, PermissionMode } from '@hapi/protocol/types'
|
||||
import { isRedundantGoalStatusEventContent } from '@hapi/protocol/messages'
|
||||
import type { Store, StoredSession } from '../../../store'
|
||||
import type { SyncEvent } from '../../../sync/syncEngine'
|
||||
import { extractTodoWriteTodosFromMessageContent } from '../../../sync/todos'
|
||||
@@ -97,6 +98,10 @@ export function registerSessionHandlers(socket: CliSocketWithData, deps: Session
|
||||
}
|
||||
const session = sessionAccess.value
|
||||
|
||||
if (isRedundantGoalStatusEventContent(content)) {
|
||||
return
|
||||
}
|
||||
|
||||
const msg = store.messages.addMessage(sid, content, localId)
|
||||
if (shouldRecordSessionActivity(content)) {
|
||||
onSessionActivity?.(sid, msg.createdAt)
|
||||
|
||||
@@ -65,6 +65,68 @@ function makePublisher() {
|
||||
// Tests
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
describe('MessageService goal status filtering', () => {
|
||||
function redundantGoalStatusContent(message: string): unknown {
|
||||
return {
|
||||
role: 'agent',
|
||||
content: {
|
||||
id: `event-${message}`,
|
||||
type: 'event',
|
||||
data: { type: 'message', message }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
it('hides stored redundant goal status events but keeps actionable goal messages', () => {
|
||||
const store = makeStore()
|
||||
const session = makeSession(store, 'goal-status-filter')
|
||||
|
||||
store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: '/goal ship it' } })
|
||||
store.messages.addMessage(session.id, redundantGoalStatusContent('Goal active · 8016 tokens'))
|
||||
store.messages.addMessage(session.id, redundantGoalStatusContent('No goal to clear'))
|
||||
|
||||
const service = new MessageService(store, makeIo(() => {}), makePublisher() as any)
|
||||
const page = service.getMessagesPage(session.id, { limit: 10, beforeSeq: null })
|
||||
|
||||
expect(page.messages.map(message => message.content)).toEqual([
|
||||
{ role: 'user', content: { type: 'text', text: '/goal ship it' } },
|
||||
redundantGoalStatusContent('No goal to clear')
|
||||
])
|
||||
})
|
||||
|
||||
it('pages past hidden-only goal status rows', () => {
|
||||
const store = makeStore()
|
||||
const session = makeSession(store, 'goal-status-pagination')
|
||||
|
||||
const user = store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: '/goal ship it' } })
|
||||
store.messages.addMessage(session.id, redundantGoalStatusContent('Goal active'))
|
||||
|
||||
const service = new MessageService(store, makeIo(() => {}), makePublisher() as any)
|
||||
const latest = service.getMessagesPage(session.id, { limit: 1, beforeSeq: null })
|
||||
|
||||
expect(latest.messages).toHaveLength(1)
|
||||
expect(latest.messages[0]?.id).toBe(user.id)
|
||||
expect(latest.page.nextBeforeSeq).toBe(user.seq)
|
||||
expect(latest.page.hasMore).toBe(false)
|
||||
})
|
||||
|
||||
it('pages past hidden-only goal status rows in position pagination', () => {
|
||||
const store = makeStore()
|
||||
const session = makeSession(store, 'goal-status-position-pagination')
|
||||
|
||||
const user = store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: '/goal ship it' } })
|
||||
store.messages.addMessage(session.id, redundantGoalStatusContent('Goal active · 8016 tokens'))
|
||||
|
||||
const service = new MessageService(store, makeIo(() => {}), makePublisher() as any)
|
||||
const latest = service.getMessagesPageByPosition(session.id, { limit: 1, before: null })
|
||||
|
||||
expect(latest.messages).toHaveLength(1)
|
||||
expect(latest.messages[0]?.id).toBe(user.id)
|
||||
expect(latest.page.nextBeforeSeq).toBe(user.seq)
|
||||
expect(latest.page.hasMore).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
describe('MessageService.cancelQueuedMessage race scenarios', () => {
|
||||
describe('Race-A: CLI ack removed:true → DELETE + status=cancelled', () => {
|
||||
it('returns cancelled and emits message-cancelled SSE after CLI confirms removal', async () => {
|
||||
|
||||
@@ -1,9 +1,43 @@
|
||||
import type { AttachmentMetadata, DecryptedMessage } from '@hapi/protocol/types'
|
||||
import { isRedundantGoalStatusEventContent } from '@hapi/protocol/messages'
|
||||
import type { Server } from 'socket.io'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import type { Store, CancelQueuedMessageResult } from '../store'
|
||||
import { EventPublisher } from './eventPublisher'
|
||||
|
||||
type StoredMessageForDelivery = ReturnType<Store['messages']['getMessages']>[number]
|
||||
|
||||
function isWebVisibleStoredMessage(message: StoredMessageForDelivery): boolean {
|
||||
return !isRedundantGoalStatusEventContent(message.content)
|
||||
}
|
||||
|
||||
function toDecryptedMessage(message: StoredMessageForDelivery): DecryptedMessage {
|
||||
return {
|
||||
id: message.id,
|
||||
seq: message.seq,
|
||||
localId: message.localId,
|
||||
content: message.content,
|
||||
createdAt: message.createdAt,
|
||||
invokedAt: message.invokedAt,
|
||||
scheduledAt: message.scheduledAt
|
||||
}
|
||||
}
|
||||
|
||||
function toVisibleDecryptedMessages(messages: StoredMessageForDelivery[]): DecryptedMessage[] {
|
||||
return messages.filter(isWebVisibleStoredMessage).map(toDecryptedMessage)
|
||||
}
|
||||
|
||||
function getOldestSeq(messages: StoredMessageForDelivery[]): number | null {
|
||||
let oldestSeq: number | null = null
|
||||
for (const message of messages) {
|
||||
if (typeof message.seq !== 'number') continue
|
||||
if (oldestSeq === null || message.seq < oldestSeq) {
|
||||
oldestSeq = message.seq
|
||||
}
|
||||
}
|
||||
return oldestSeq
|
||||
}
|
||||
|
||||
export class MessageService {
|
||||
constructor(
|
||||
private readonly store: Store,
|
||||
@@ -15,13 +49,7 @@ export class MessageService {
|
||||
|
||||
getMessages(sessionId: string, limit: number = 200): DecryptedMessage[] {
|
||||
const stored = this.store.messages.getMessages(sessionId, limit)
|
||||
return stored.map((message) => ({
|
||||
id: message.id,
|
||||
seq: message.seq,
|
||||
localId: message.localId,
|
||||
content: message.content,
|
||||
createdAt: message.createdAt
|
||||
}))
|
||||
return toVisibleDecryptedMessages(stored)
|
||||
}
|
||||
|
||||
getMessagesPage(sessionId: string, options: { limit: number; beforeSeq: number | null }): {
|
||||
@@ -33,28 +61,21 @@ export class MessageService {
|
||||
hasMore: boolean
|
||||
}
|
||||
} {
|
||||
const stored = this.store.messages.getMessages(sessionId, options.limit, options.beforeSeq ?? undefined)
|
||||
const messages: DecryptedMessage[] = stored.map((message) => ({
|
||||
id: message.id,
|
||||
seq: message.seq,
|
||||
localId: message.localId,
|
||||
content: message.content,
|
||||
createdAt: message.createdAt,
|
||||
invokedAt: message.invokedAt,
|
||||
scheduledAt: message.scheduledAt
|
||||
}))
|
||||
let stored = this.store.messages.getMessages(sessionId, options.limit, options.beforeSeq ?? undefined)
|
||||
let messages = toVisibleDecryptedMessages(stored)
|
||||
let oldestSeq = getOldestSeq(stored)
|
||||
let hasMore = oldestSeq !== null
|
||||
&& this.store.messages.getMessages(sessionId, 1, oldestSeq).length > 0
|
||||
|
||||
let oldestSeq: number | null = null
|
||||
for (const message of messages) {
|
||||
if (typeof message.seq !== 'number') continue
|
||||
if (oldestSeq === null || message.seq < oldestSeq) {
|
||||
oldestSeq = message.seq
|
||||
}
|
||||
while (messages.length === 0 && hasMore && oldestSeq !== null) {
|
||||
stored = this.store.messages.getMessages(sessionId, options.limit, oldestSeq)
|
||||
messages = toVisibleDecryptedMessages(stored)
|
||||
oldestSeq = getOldestSeq(stored)
|
||||
hasMore = oldestSeq !== null
|
||||
&& this.store.messages.getMessages(sessionId, 1, oldestSeq).length > 0
|
||||
}
|
||||
|
||||
const nextBeforeSeq = oldestSeq
|
||||
const hasMore = nextBeforeSeq !== null
|
||||
&& this.store.messages.getMessages(sessionId, 1, nextBeforeSeq).length > 0
|
||||
|
||||
return {
|
||||
messages,
|
||||
@@ -79,53 +100,73 @@ export class MessageService {
|
||||
hasMore: boolean
|
||||
}
|
||||
} {
|
||||
const before = options.before ?? undefined
|
||||
const pageRows = this.store.messages.getMessagesByPosition(sessionId, options.limit, before)
|
||||
let before = options.before ?? undefined
|
||||
let pageRows = this.store.messages.getMessagesByPosition(sessionId, options.limit, before)
|
||||
|
||||
// Latest-page request (no cursor): also include uninvoked local user messages
|
||||
// out-of-band, so refresh / secondary clients can still see queued rows even
|
||||
// when their position key (createdAt) places them outside the latest page.
|
||||
// The cursor stays anchored to pageRows so out-of-band rows don't affect
|
||||
// pagination of older pages.
|
||||
const queuedRows = before === undefined
|
||||
let queuedRows = before === undefined
|
||||
? this.store.messages.getUninvokedLocalMessages(sessionId)
|
||||
: []
|
||||
|
||||
const byId = new Map<string, typeof pageRows[number]>()
|
||||
let byId = new Map<string, typeof pageRows[number]>()
|
||||
for (const row of pageRows) byId.set(row.id, row)
|
||||
for (const row of queuedRows) byId.set(row.id, row)
|
||||
|
||||
const stored = [...byId.values()].sort((a, b) => {
|
||||
let stored = [...byId.values()].sort((a, b) => {
|
||||
const at = (a.invokedAt ?? a.createdAt) - (b.invokedAt ?? b.createdAt)
|
||||
return at !== 0 ? at : a.seq - b.seq
|
||||
})
|
||||
|
||||
const messages: DecryptedMessage[] = stored.map((message) => ({
|
||||
id: message.id,
|
||||
seq: message.seq,
|
||||
localId: message.localId,
|
||||
content: message.content,
|
||||
createdAt: message.createdAt,
|
||||
invokedAt: message.invokedAt,
|
||||
scheduledAt: message.scheduledAt
|
||||
}))
|
||||
let messages = toVisibleDecryptedMessages(stored)
|
||||
|
||||
// The cursor is the oldest row in the actual position-ordered page (pageRows[0]).
|
||||
// Out-of-band queued rows are not part of the cursor — they are pinned to
|
||||
// every latest-page response.
|
||||
const oldest = pageRows[0] ?? null
|
||||
const oldestSeq: number | null = oldest?.seq ?? null
|
||||
const oldestPositionAt: number | null = oldest
|
||||
let oldest = pageRows[0] ?? null
|
||||
let oldestSeq: number | null = oldest?.seq ?? null
|
||||
let oldestPositionAt: number | null = oldest
|
||||
? oldest.invokedAt ?? oldest.createdAt
|
||||
: null
|
||||
|
||||
const hasMore = oldestSeq !== null && oldestPositionAt !== null
|
||||
let hasMore = oldestSeq !== null && oldestPositionAt !== null
|
||||
&& this.store.messages.getMessagesByPosition(
|
||||
sessionId,
|
||||
1,
|
||||
{ at: oldestPositionAt, seq: oldestSeq }
|
||||
).length > 0
|
||||
|
||||
while (messages.length === 0 && hasMore && oldestSeq !== null && oldestPositionAt !== null) {
|
||||
before = { at: oldestPositionAt, seq: oldestSeq }
|
||||
pageRows = this.store.messages.getMessagesByPosition(sessionId, options.limit, before)
|
||||
queuedRows = []
|
||||
|
||||
byId = new Map<string, typeof pageRows[number]>()
|
||||
for (const row of pageRows) byId.set(row.id, row)
|
||||
for (const row of queuedRows) byId.set(row.id, row)
|
||||
|
||||
stored = [...byId.values()].sort((a, b) => {
|
||||
const at = (a.invokedAt ?? a.createdAt) - (b.invokedAt ?? b.createdAt)
|
||||
return at !== 0 ? at : a.seq - b.seq
|
||||
})
|
||||
messages = toVisibleDecryptedMessages(stored)
|
||||
|
||||
oldest = pageRows[0] ?? null
|
||||
oldestSeq = oldest?.seq ?? null
|
||||
oldestPositionAt = oldest
|
||||
? oldest.invokedAt ?? oldest.createdAt
|
||||
: null
|
||||
hasMore = oldestSeq !== null && oldestPositionAt !== null
|
||||
&& this.store.messages.getMessagesByPosition(
|
||||
sessionId,
|
||||
1,
|
||||
{ at: oldestPositionAt, seq: oldestSeq }
|
||||
).length > 0
|
||||
}
|
||||
|
||||
return {
|
||||
messages,
|
||||
page: {
|
||||
|
||||
@@ -50,4 +50,24 @@ export function isClaudeChatVisibleMessage(message: { type: unknown; subtype?: u
|
||||
return isClaudeChatVisibleSystemSubtype(message.subtype)
|
||||
}
|
||||
|
||||
export function isRedundantGoalStatusMessageText(value: unknown): boolean {
|
||||
if (typeof value !== 'string') return false
|
||||
const message = value.trim()
|
||||
return message === 'Goal cleared'
|
||||
|| /^Goal (active|paused|complete|limited by budget)(?:$|\s+·\s+)/.test(message)
|
||||
}
|
||||
|
||||
export function isRedundantGoalStatusEventContent(value: unknown): boolean {
|
||||
const record = unwrapRoleWrappedRecordEnvelope(value)
|
||||
if (record?.role !== 'agent') return false
|
||||
|
||||
const eventContent = record.content
|
||||
if (!isObject(eventContent) || eventContent.type !== 'event') return false
|
||||
|
||||
const data = isObject(eventContent.data) ? eventContent.data : null
|
||||
if (!data || data.type !== 'message') return false
|
||||
|
||||
return isRedundantGoalStatusMessageText(data.message)
|
||||
}
|
||||
|
||||
export type { RoleWrappedRecord }
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { reduceChatBlocks } from './reducer'
|
||||
import { normalizeDecryptedMessage } from './normalize'
|
||||
import type { NormalizedMessage } from './types'
|
||||
import type { DecryptedMessage } from '@/types/api'
|
||||
import type { ThreadGoal, ThreadGoalStatus } from '@/types/api'
|
||||
|
||||
function userMessage(id: string, text: string, createdAt: number): NormalizedMessage {
|
||||
@@ -53,6 +55,30 @@ function goalClearedMessage(id: string, createdAt: number): NormalizedMessage {
|
||||
}
|
||||
}
|
||||
|
||||
function eventMessage(id: string, message: string, createdAt: number): NormalizedMessage {
|
||||
return {
|
||||
id,
|
||||
localId: null,
|
||||
createdAt,
|
||||
role: 'event',
|
||||
content: {
|
||||
type: 'message',
|
||||
message
|
||||
},
|
||||
isSidechain: false
|
||||
}
|
||||
}
|
||||
|
||||
function decryptedMessage(id: string, content: unknown, createdAt: number): DecryptedMessage {
|
||||
return {
|
||||
id,
|
||||
seq: 1,
|
||||
localId: null,
|
||||
content,
|
||||
createdAt
|
||||
}
|
||||
}
|
||||
|
||||
describe('reduceChatBlocks', () => {
|
||||
it('ignores child agent usage when calculating parent latest usage', () => {
|
||||
const messages: NormalizedMessage[] = [
|
||||
@@ -127,6 +153,21 @@ describe('reduceChatBlocks', () => {
|
||||
expect(reduced.latestGoal).toBeNull()
|
||||
})
|
||||
|
||||
it('can clear completed goal state using messages hidden from the rendered timeline', () => {
|
||||
const renderedMessages = [
|
||||
goalMessage('goal-complete', 'complete', 1)
|
||||
]
|
||||
const goalStateMessages = [
|
||||
...renderedMessages,
|
||||
userMessage('queued-user-later', 'start a new task', 2)
|
||||
]
|
||||
|
||||
const reduced = reduceChatBlocks(renderedMessages, null, { goalStateMessages })
|
||||
|
||||
expect(reduced.blocks).toHaveLength(0)
|
||||
expect(reduced.latestGoal).toBeNull()
|
||||
})
|
||||
|
||||
it('does not treat later goal slash commands as non-goal activity', () => {
|
||||
const reduced = reduceChatBlocks([
|
||||
goalMessage('goal-complete', 'complete', 1),
|
||||
@@ -155,4 +196,87 @@ describe('reduceChatBlocks', () => {
|
||||
|
||||
expect(reduced.latestGoal).toBeNull()
|
||||
})
|
||||
|
||||
it('uses goal events for latest goal state without rendering timeline prompts', () => {
|
||||
const reduced = reduceChatBlocks([
|
||||
goalMessage('goal-active', 'active', 1)
|
||||
], null)
|
||||
|
||||
expect(reduced.blocks).toHaveLength(0)
|
||||
expect(reduced.latestGoal).toMatchObject({
|
||||
threadId: 'thread-1',
|
||||
objective: 'ship goal support',
|
||||
status: 'active'
|
||||
})
|
||||
})
|
||||
|
||||
it('uses goal clear events to clear latest goal without rendering timeline prompts', () => {
|
||||
const reduced = reduceChatBlocks([
|
||||
goalMessage('goal-active', 'active', 1),
|
||||
goalClearedMessage('goal-cleared', 2)
|
||||
], null)
|
||||
|
||||
expect(reduced.blocks).toHaveLength(0)
|
||||
expect(reduced.latestGoal).toBeNull()
|
||||
})
|
||||
|
||||
it('hides redundant goal status messages but keeps actionable goal messages', () => {
|
||||
const reduced = reduceChatBlocks([
|
||||
eventMessage('goal-active-message', 'Goal active', 1),
|
||||
eventMessage('goal-active-usage-message', 'Goal active · 181737 tokens', 2),
|
||||
eventMessage('goal-complete-message', 'Goal complete', 3),
|
||||
eventMessage('goal-cleared-message', 'Goal cleared', 4),
|
||||
eventMessage('goal-actionable-message', 'No goal to clear', 5)
|
||||
], null)
|
||||
|
||||
expect(reduced.blocks).toHaveLength(1)
|
||||
expect(reduced.blocks[0]).toMatchObject({
|
||||
kind: 'agent-event',
|
||||
event: { type: 'message', message: 'No goal to clear' }
|
||||
})
|
||||
})
|
||||
|
||||
it('hides persisted goal status event envelopes alongside structured goal events', () => {
|
||||
const goal: ThreadGoal = {
|
||||
threadId: 'thread-1',
|
||||
objective: 'ship goal support',
|
||||
status: 'active',
|
||||
tokenBudget: null,
|
||||
tokensUsed: 8016,
|
||||
timeUsedSeconds: 10,
|
||||
createdAt: 1,
|
||||
updatedAt: 2
|
||||
}
|
||||
const normalized = [
|
||||
decryptedMessage('goal-status-envelope', {
|
||||
role: 'agent',
|
||||
content: {
|
||||
id: 'event-1',
|
||||
type: 'event',
|
||||
data: { type: 'message', message: 'Goal active · 8016 tokens' }
|
||||
}
|
||||
}, 1),
|
||||
decryptedMessage('goal-structured-envelope', {
|
||||
role: 'agent',
|
||||
content: {
|
||||
type: 'codex',
|
||||
data: {
|
||||
type: 'thread_goal_updated',
|
||||
thread_id: 'thread-1',
|
||||
goal
|
||||
}
|
||||
}
|
||||
}, 2)
|
||||
].map(message => normalizeDecryptedMessage(message))
|
||||
.filter((message): message is NormalizedMessage => message !== null)
|
||||
|
||||
const reduced = reduceChatBlocks(normalized, null)
|
||||
|
||||
expect(reduced.blocks).toHaveLength(0)
|
||||
expect(reduced.latestGoal).toMatchObject({
|
||||
threadId: 'thread-1',
|
||||
status: 'active',
|
||||
tokensUsed: 8016
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
+41
-3
@@ -5,6 +5,7 @@ import { traceMessages, type TracedMessage } from '@/chat/tracer'
|
||||
import { dedupeAgentEvents, foldApiErrorEvents } from '@/chat/reducerEvents'
|
||||
import { collectTitleChanges, collectToolIdsFromMessages, ensureToolBlock, getPermissions } from '@/chat/reducerTools'
|
||||
import { reduceTimeline } from '@/chat/reducerTimeline'
|
||||
import { isRedundantGoalStatusMessageText } from '@hapi/protocol/messages'
|
||||
|
||||
// Calculate context size from usage data
|
||||
function calculateContextSize(usage: UsageData): number {
|
||||
@@ -28,6 +29,10 @@ export type LatestUsage = {
|
||||
timestamp: number
|
||||
}
|
||||
|
||||
export type ReduceChatBlocksOptions = {
|
||||
goalStateMessages?: NormalizedMessage[]
|
||||
}
|
||||
|
||||
function getLatestThreadGoal(normalized: NormalizedMessage[]): ThreadGoal | null {
|
||||
let sawNewerNonGoalUserMessage = false
|
||||
for (let i = normalized.length - 1; i >= 0; i--) {
|
||||
@@ -52,9 +57,42 @@ function getLatestThreadGoal(normalized: NormalizedMessage[]): ThreadGoal | null
|
||||
return null
|
||||
}
|
||||
|
||||
function isRedundantGoalStatusMessage(event: AgentEvent): boolean {
|
||||
if (event.type !== 'message') return false
|
||||
return isRedundantGoalStatusMessageText(event.message)
|
||||
}
|
||||
|
||||
function isSilentGoalEventBlock(block: ChatBlock): boolean {
|
||||
return block.kind === 'agent-event'
|
||||
&& (
|
||||
block.event.type === 'thread-goal-updated'
|
||||
|| block.event.type === 'thread-goal-cleared'
|
||||
|| isRedundantGoalStatusMessage(block.event)
|
||||
)
|
||||
}
|
||||
|
||||
function filterSilentGoalBlocks(blocks: ChatBlock[]): ChatBlock[] {
|
||||
const filtered: ChatBlock[] = []
|
||||
|
||||
for (const block of blocks) {
|
||||
if (isSilentGoalEventBlock(block)) continue
|
||||
if (block.kind === 'tool-call' && block.children.length > 0) {
|
||||
filtered.push({
|
||||
...block,
|
||||
children: filterSilentGoalBlocks(block.children)
|
||||
})
|
||||
continue
|
||||
}
|
||||
filtered.push(block)
|
||||
}
|
||||
|
||||
return filtered
|
||||
}
|
||||
|
||||
export function reduceChatBlocks(
|
||||
normalized: NormalizedMessage[],
|
||||
agentState: AgentState | null | undefined
|
||||
agentState: AgentState | null | undefined,
|
||||
options: ReduceChatBlocksOptions = {}
|
||||
): { blocks: ChatBlock[]; hasReadyEvent: boolean; latestUsage: LatestUsage | null; latestGoal: ThreadGoal | null } {
|
||||
const permissionsById = getPermissions(agentState)
|
||||
const toolIdsInMessages = collectToolIdsFromMessages(normalized)
|
||||
@@ -142,9 +180,9 @@ export function reduceChatBlocks(
|
||||
}
|
||||
|
||||
return {
|
||||
blocks: dedupeAgentEvents(foldApiErrorEvents(rootResult.blocks)),
|
||||
blocks: filterSilentGoalBlocks(dedupeAgentEvents(foldApiErrorEvents(rootResult.blocks))),
|
||||
hasReadyEvent,
|
||||
latestUsage,
|
||||
latestGoal: getLatestThreadGoal(normalized)
|
||||
latestGoal: getLatestThreadGoal(options.goalStateMessages ?? normalized)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,31 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { shouldAutoClearPendingSchedule } from './SessionChat'
|
||||
import { buildGoalStateMessages, shouldAutoClearPendingSchedule } from './SessionChat'
|
||||
import type { PendingSchedule } from '@/components/AssistantChat/ScheduleTimePicker'
|
||||
import type { DecryptedMessage } from '@/types/api'
|
||||
|
||||
function userMessage(props: {
|
||||
id: string
|
||||
createdAt: number
|
||||
localId?: string | null
|
||||
invokedAt?: number | null
|
||||
scheduledAt?: number | null
|
||||
}): DecryptedMessage {
|
||||
return {
|
||||
id: props.id,
|
||||
seq: null,
|
||||
localId: props.localId ?? null,
|
||||
content: {
|
||||
role: 'user',
|
||||
content: {
|
||||
type: 'text',
|
||||
text: 'hello'
|
||||
}
|
||||
},
|
||||
createdAt: props.createdAt,
|
||||
invokedAt: props.invokedAt,
|
||||
scheduledAt: props.scheduledAt
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Unit tests for shouldAutoClearPendingSchedule.
|
||||
@@ -41,3 +66,58 @@ describe('shouldAutoClearPendingSchedule', () => {
|
||||
expect(shouldAutoClearPendingSchedule(expired)).toBe(true)
|
||||
})
|
||||
})
|
||||
|
||||
describe('buildGoalStateMessages', () => {
|
||||
it('keeps immediate queued user messages so completed goal status can clear before timeline render', () => {
|
||||
const now = 1_700_000_000_000
|
||||
const messages = [
|
||||
userMessage({
|
||||
id: 'local-immediate',
|
||||
localId: 'local-immediate',
|
||||
createdAt: now,
|
||||
invokedAt: null
|
||||
})
|
||||
]
|
||||
|
||||
expect(buildGoalStateMessages(messages).map((message) => message.id))
|
||||
.toEqual(['local-immediate'])
|
||||
})
|
||||
|
||||
it('includes pending messages that are outside the visible timeline window', () => {
|
||||
const now = 1_700_000_000_000
|
||||
const visible = [
|
||||
userMessage({ id: 'visible', createdAt: now - 10 })
|
||||
]
|
||||
const pending = [
|
||||
userMessage({ id: 'pending', createdAt: now })
|
||||
]
|
||||
|
||||
expect(buildGoalStateMessages(visible, pending).map((message) => message.id))
|
||||
.toEqual(['visible', 'pending'])
|
||||
})
|
||||
|
||||
it('ignores uninvoked scheduled messages, including mature prompts, until they are invoked', () => {
|
||||
const now = 1_700_000_000_000
|
||||
const futureQueued = userMessage({
|
||||
id: 'future',
|
||||
createdAt: now,
|
||||
invokedAt: null,
|
||||
scheduledAt: now + 60_000
|
||||
})
|
||||
const matureQueued = userMessage({
|
||||
id: 'mature',
|
||||
createdAt: now + 1,
|
||||
invokedAt: null,
|
||||
scheduledAt: now - 60_000
|
||||
})
|
||||
const invokedScheduled = userMessage({
|
||||
id: 'invoked',
|
||||
createdAt: now + 2,
|
||||
invokedAt: now + 30_000,
|
||||
scheduledAt: now - 60_000
|
||||
})
|
||||
|
||||
expect(buildGoalStateMessages([futureQueued, matureQueued, invokedScheduled]).map((message) => message.id))
|
||||
.toEqual(['invoked'])
|
||||
})
|
||||
})
|
||||
|
||||
@@ -17,7 +17,7 @@ import { reduceChatBlocks } from '@/chat/reducer'
|
||||
import { reconcileChatBlocks } from '@/chat/reconcile'
|
||||
import { buildConversationOutline } from '@/chat/outline'
|
||||
import { buildVisibleChatBlocks, isToolGroupBlock, type ToolGroupBlock } from '@/chat/toolGroups'
|
||||
import { isQueuedForInvocation } from '@/lib/messages'
|
||||
import { isQueuedForInvocation, mergeMessages } from '@/lib/messages'
|
||||
import { HappyComposer } from '@/components/AssistantChat/HappyComposer'
|
||||
import type { PendingSchedule } from '@/components/AssistantChat/ScheduleTimePicker'
|
||||
import { resolvePendingSchedule } from '@/components/AssistantChat/ScheduleTimePicker'
|
||||
@@ -49,6 +49,21 @@ export function shouldAutoClearPendingSchedule(pending: PendingSchedule | null):
|
||||
return pending !== null && pending.type === 'absolute'
|
||||
}
|
||||
|
||||
function isUninvokedScheduledMessage(message: DecryptedMessage): boolean {
|
||||
return message.invokedAt == null && message.scheduledAt != null
|
||||
}
|
||||
|
||||
export function buildGoalStateMessages(
|
||||
messages: DecryptedMessage[],
|
||||
pendingMessages: DecryptedMessage[] = []
|
||||
): DecryptedMessage[] {
|
||||
const eligibleMessages = messages.filter((message) => !isUninvokedScheduledMessage(message))
|
||||
const eligiblePendingMessages = pendingMessages.filter((message) => !isUninvokedScheduledMessage(message))
|
||||
return eligiblePendingMessages.length > 0
|
||||
? mergeMessages(eligibleMessages, eligiblePendingMessages)
|
||||
: eligibleMessages
|
||||
}
|
||||
|
||||
function getOutlineTitle(session: Session): string {
|
||||
if (session.metadata?.name) {
|
||||
return session.metadata.name
|
||||
@@ -83,6 +98,7 @@ export function SessionChat(props: {
|
||||
api: ApiClient
|
||||
session: Session
|
||||
messages: DecryptedMessage[]
|
||||
pendingMessages?: DecryptedMessage[]
|
||||
messagesWarning: string | null
|
||||
hasMoreMessages: boolean
|
||||
isLoadingMessages: boolean
|
||||
@@ -306,9 +322,25 @@ export function SessionChat(props: {
|
||||
return normalized
|
||||
}, [visibleMessages])
|
||||
|
||||
const goalStateSourceMessages = useMemo(
|
||||
() => buildGoalStateMessages(props.messages, props.pendingMessages ?? []),
|
||||
[props.messages, props.pendingMessages]
|
||||
)
|
||||
|
||||
const normalizedGoalStateMessages: NormalizedMessage[] = useMemo(() => {
|
||||
const normalized: NormalizedMessage[] = []
|
||||
for (const message of goalStateSourceMessages) {
|
||||
const next = normalizeDecryptedMessage(message)
|
||||
if (next) normalized.push(next)
|
||||
}
|
||||
return normalized
|
||||
}, [goalStateSourceMessages])
|
||||
|
||||
const reduced = useMemo(
|
||||
() => reduceChatBlocks(normalizedMessages, props.session.agentState),
|
||||
[normalizedMessages, props.session.agentState]
|
||||
() => reduceChatBlocks(normalizedMessages, props.session.agentState, {
|
||||
goalStateMessages: normalizedGoalStateMessages
|
||||
}),
|
||||
[normalizedMessages, normalizedGoalStateMessages, props.session.agentState]
|
||||
)
|
||||
const reconciled = useMemo(
|
||||
() => reconcileChatBlocks(reduced.blocks, blocksByIdRef.current),
|
||||
|
||||
@@ -28,6 +28,7 @@ export const EMPTY_STATE: MessageWindowState = {
|
||||
|
||||
export function useMessages(api: ApiClient | null, sessionId: string | null): {
|
||||
messages: DecryptedMessage[]
|
||||
pendingMessages: DecryptedMessage[]
|
||||
warning: string | null
|
||||
isLoading: boolean
|
||||
isLoadingMore: boolean
|
||||
@@ -88,6 +89,7 @@ export function useMessages(api: ApiClient | null, sessionId: string | null): {
|
||||
|
||||
return {
|
||||
messages: state.messages,
|
||||
pendingMessages: state.pending,
|
||||
warning: state.warning,
|
||||
isLoading: state.isLoading,
|
||||
isLoadingMore: state.isLoadingMore,
|
||||
|
||||
@@ -259,6 +259,7 @@ function SessionPage() {
|
||||
} = useSession(api, sessionId)
|
||||
const {
|
||||
messages,
|
||||
pendingMessages,
|
||||
warning: messagesWarning,
|
||||
isLoading: messagesLoading,
|
||||
isLoadingMore: messagesLoadingMore,
|
||||
@@ -371,6 +372,7 @@ function SessionPage() {
|
||||
api={api}
|
||||
session={session}
|
||||
messages={messages}
|
||||
pendingMessages={pendingMessages}
|
||||
messagesWarning={messagesWarning}
|
||||
hasMoreMessages={messagesHasMore}
|
||||
isLoadingMessages={messagesLoading}
|
||||
|
||||
Reference in New Issue
Block a user