mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 08:02:21 +00:00
fix(native-chat): a chat you return to drops background tasks that finished while it was hidden (#24305)
* fix(native-chat): a chat you return to drops background tasks that finished while it was hidden A background command that finished while its chat pane was hidden stayed in the strip as running, its timer counting for hours, and Stop failed with "The background task wasn't stopped." The pane stops listening while hidden, so it missed the "no tasks" update; once the idle sweep stopped the agent, the host reopened the pane's subscription without any roster at all, and the client reads a missing roster as "unchanged". The roster now rides each subscriber's frames the way the slash-command list and the message queue already do. The subscriber registry reads it from the host's child records on a subscriber's first frame and on every snapshot and reset, and re-sends it to each subscriber whose last copy differs when the records change. The channel always answers: a conversation with no running child is "none" whether or not an agent holds it. Ordinary journal batches never read it. This also fixes the reverse case: a snapshot or reset sent to a live pane carried no roster, which cleared a still-running task from the strip until the roster next changed. The client treats a resumed subscription's first batch as stating the roster, so a pane reconnecting to an older host that still omits it drops its stale copy too. Fixes #24227 * refactor(native-chat): require the background-task roster read and drop the unused close observers The host's delivery layer now requires every hook it is built with, so production wiring cannot leave out the roster read (an opening frame without it reads as "no tasks" to current clients). The conversations map's close observers lost their only caller in this PR and are removed. The reducer test comment states the old-host omission rule precisely. * refactor(native-chat): keep structured-agent-session-host.ts within the line limit Main brought the host file to the 300-line lint limit, and this PR's roster wiring adds one line. Name the lease reconciler's type by its factory, as the neighbouring fields do. * test(native-chat): build the roster test's Claude provider handle with claudeProviderHandle Main made the provider handle opaque (#24991); build it with the existing helper, as main's own tests now do.
This commit is contained in:
@@ -5,7 +5,6 @@
|
||||
|
||||
import {
|
||||
AGENT_SESSION_HISTORY_MAX_LIMIT,
|
||||
type AgentSessionBackgroundTaskState,
|
||||
type AgentSessionSlashCommand,
|
||||
type AgentSessionSubscribeEvent,
|
||||
type AgentSessionTurnActivity
|
||||
@@ -40,16 +39,14 @@ export function deliverToSubscriber(
|
||||
journal: AgentSessionJournal
|
||||
hostNow: number
|
||||
emitCheckpoint: boolean
|
||||
backgroundTasks?: AgentSessionBackgroundTaskState | null | undefined
|
||||
activity?: AgentSessionTurnActivity | null | undefined
|
||||
}
|
||||
): void {
|
||||
const { subscriber, journal, hostNow, emitCheckpoint, backgroundTasks, activity } = input
|
||||
const { subscriber, journal, hostNow, emitCheckpoint, activity } = input
|
||||
const checkpointActivity = emitCheckpoint ? port.activity(subscriber.sessionId) : undefined
|
||||
const publishedActivity = activity !== undefined ? activity : checkpointActivity
|
||||
const shared = {
|
||||
hostNow,
|
||||
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
|
||||
...(publishedActivity !== undefined ? { activity: publishedActivity } : {})
|
||||
}
|
||||
// Caught up, so there are no rows to read: every publish behind a commit's own delivery.
|
||||
@@ -115,7 +112,6 @@ function emitCaughtUp(
|
||||
emitCheckpoint: boolean,
|
||||
shared: {
|
||||
hostNow: number
|
||||
backgroundTasks?: AgentSessionBackgroundTaskState | null
|
||||
activity?: AgentSessionTurnActivity | null
|
||||
}
|
||||
): void {
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
// Which per-emit fields ride one subscriber frame: the provider command catalog
|
||||
// and the queue publication (the draft list with the queue's pause). Both are
|
||||
// identity-deduplicated against the LAST VALUE SENT — never advanced on a frame
|
||||
// that withheld the field, or the final replacement would be suppressed — and
|
||||
// both attach whole to hydrating frames.
|
||||
// Which per-emit fields ride one subscriber frame: the provider command catalog,
|
||||
// the queue publication (the draft list with the queue's pause) and the strip's
|
||||
// background-task roster. Each is deduplicated against the LAST VALUE SENT to that
|
||||
// subscriber — never advanced on a frame that withheld the field, or the final
|
||||
// replacement would be suppressed — and each attaches whole to hydrating frames.
|
||||
|
||||
import type {
|
||||
AgentSessionBackgroundTaskState,
|
||||
AgentSessionSlashCommand,
|
||||
AgentSessionSubscribeEvent
|
||||
} from '../../../shared/agent-session-wire'
|
||||
@@ -15,11 +16,15 @@ export type SubscriberFieldState = {
|
||||
commands?: AgentSessionSlashCommand[] | null
|
||||
/** The last queue publication actually SENT. */
|
||||
queuePublication?: QueuePublication
|
||||
/** Fingerprint of the roster last SENT; absent until this subscriber's first frame. */
|
||||
backgroundTasks?: string
|
||||
}
|
||||
|
||||
export type SubscriberFieldHooks = {
|
||||
readCommands?: (sessionId: string) => AgentSessionSlashCommand[] | undefined
|
||||
readQueuePublication?: (sessionId: string) => QueuePublication | undefined
|
||||
/** Built from the host's child records, so it is read only when a frame owes it: never per token. */
|
||||
readBackgroundTasks?: (sessionId: string) => AgentSessionBackgroundTaskState | null
|
||||
}
|
||||
|
||||
export type SubscriberFrame = {
|
||||
@@ -27,6 +32,12 @@ export type SubscriberFrame = {
|
||||
commands: AgentSessionSlashCommand[] | null
|
||||
attachedQueued: boolean
|
||||
queued: QueuePublication | undefined
|
||||
/** The attached roster's fingerprint; undefined when the frame carries none. */
|
||||
backgroundTasks: string | undefined
|
||||
}
|
||||
|
||||
export function backgroundTaskFingerprint(state: AgentSessionBackgroundTaskState | null): string {
|
||||
return JSON.stringify(state)
|
||||
}
|
||||
|
||||
/** Builds the frame to emit; the caller stores the returned refs only after the
|
||||
@@ -49,20 +60,43 @@ export function buildSubscriberFrame(
|
||||
queued !== undefined &&
|
||||
event.type !== 'end' &&
|
||||
(event.type !== 'batch' || queued !== subscriber.queuePublication)
|
||||
const backgroundTasks = subscriberBackgroundTasks(hooks, subscriber, event)
|
||||
return {
|
||||
frame: {
|
||||
...event,
|
||||
...(includeCommands ? { commands: commands ?? null } : {}),
|
||||
...(attachedQueued && queued
|
||||
? { queuedMessages: queued.queuedMessages, queuePause: queued.queuePause }
|
||||
: {})
|
||||
: {}),
|
||||
...(backgroundTasks !== undefined ? { backgroundTasks } : {})
|
||||
},
|
||||
commands,
|
||||
attachedQueued,
|
||||
queued
|
||||
queued,
|
||||
backgroundTasks:
|
||||
backgroundTasks === undefined ? undefined : backgroundTaskFingerprint(backgroundTasks)
|
||||
}
|
||||
}
|
||||
|
||||
/** A subscriber's first frame and every hydrating frame state the roster, so a client resuming
|
||||
* from its cursor never keeps a roster that changed while it was away; a republish carries its
|
||||
* own. Other batches leave it out: the client keeps what it holds. */
|
||||
function subscriberBackgroundTasks(
|
||||
hooks: SubscriberFieldHooks,
|
||||
subscriber: SubscriberFieldState,
|
||||
event: AgentSessionSubscribeEvent
|
||||
): AgentSessionBackgroundTaskState | null | undefined {
|
||||
if (event.type === 'end') {
|
||||
return undefined
|
||||
}
|
||||
if (event.backgroundTasks !== undefined) {
|
||||
return event.backgroundTasks
|
||||
}
|
||||
return event.type !== 'batch' || subscriber.backgroundTasks === undefined
|
||||
? hooks.readBackgroundTasks?.(subscriber.sessionId)
|
||||
: undefined
|
||||
}
|
||||
|
||||
/** Whether a caught-up publish with no rows still owes this subscriber a frame:
|
||||
* draft inserts and pause changes write no journal row, so an unchanged cursor
|
||||
* must still deliver the changed publication. */
|
||||
|
||||
+33
-37
@@ -1,5 +1,5 @@
|
||||
// The strip channel keeps one fingerprint per conversation to skip an unchanged roster; a closed
|
||||
// conversation's goes with it.
|
||||
// The strip channel always answers a conversation's roster from the host's child records; each
|
||||
// subscriber's frames carry it (see structured-agent-session-subscribers.test.ts).
|
||||
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view'
|
||||
@@ -35,7 +35,7 @@ function channelOver(
|
||||
logger: recordingStructuredAgentSessionLogger().logger,
|
||||
now: () => 1
|
||||
})
|
||||
const sent = vi.fn()
|
||||
const republished = vi.fn()
|
||||
const channel = new StructuredAgentSessionBackgroundTaskChannel(
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the channel reads only the record store's fence and the adapter's stop capability.
|
||||
{
|
||||
@@ -43,8 +43,8 @@ function channelOver(
|
||||
adapter: { backgroundTaskStops: stops }
|
||||
} as unknown as StructuredAgentSessionHostDeps,
|
||||
sessions,
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: publish sends through `backgroundTasks` only.
|
||||
{ backgroundTasks: sent } as unknown as AgentSessionSubscribers,
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: publish sends through `republishBackgroundTasks` only.
|
||||
{ republishBackgroundTasks: republished } as unknown as AgentSessionSubscribers,
|
||||
async () => {
|
||||
throw new Error('not opened here')
|
||||
},
|
||||
@@ -55,43 +55,31 @@ function channelOver(
|
||||
journal: { observeCommits: () => {} },
|
||||
params: { provider: 'claude' }
|
||||
} as unknown as StructuredAgentSessionHostSession
|
||||
return { sessions, sent, channel, session }
|
||||
return { sessions, republished, channel, session }
|
||||
}
|
||||
|
||||
describe('the background-task channel', () => {
|
||||
it("forgets a closed conversation's roster, so reopening it sends the roster again", () => {
|
||||
const { sessions, sent, channel, session } = channelOver()
|
||||
it('lists running children with the stops the live provider offers', () => {
|
||||
const { sessions, channel, session } = channelOver()
|
||||
sessions.set('session-1', session)
|
||||
channel.publish('session-1')
|
||||
channel.publish('session-1')
|
||||
// An unchanged roster sends nothing.
|
||||
expect(sent).toHaveBeenCalledTimes(1)
|
||||
|
||||
sessions.delete('session-1')
|
||||
expect(channel['published'].size).toBe(0)
|
||||
sessions.set('session-1', session)
|
||||
channel.publish('session-1')
|
||||
expect(sent).toHaveBeenCalledTimes(2)
|
||||
expect(channel.read('session-1')).toMatchObject({
|
||||
state: 'monitoring',
|
||||
supportsTaskStop: true,
|
||||
children: [child]
|
||||
})
|
||||
})
|
||||
|
||||
it('sends "no children" once, not again on every change that leaves none', () => {
|
||||
let views: AgentChildWorkView[] = []
|
||||
const { sessions, sent, channel, session } = channelOver(() => views)
|
||||
// A task that finished while no pane listened, then an idle sweep that stopped the provider:
|
||||
// a pane resuming must read "none", not silence it would take as "unchanged" (#24227).
|
||||
it('answers "none" for a conversation no provider holds, not "unknown"', () => {
|
||||
const { sessions, channel, session } = channelOver(
|
||||
() => [],
|
||||
() => undefined
|
||||
)
|
||||
sessions.set('session-1', session)
|
||||
channel.publish('session-1')
|
||||
channel.publish('session-1')
|
||||
expect(sent.mock.calls.map(([, state]) => state)).toEqual([null])
|
||||
|
||||
views = [child]
|
||||
channel.publish('session-1')
|
||||
views = []
|
||||
channel.publish('session-1')
|
||||
channel.publish('session-1')
|
||||
expect(sent.mock.calls.map(([, state]) => state?.children?.length ?? null)).toEqual([
|
||||
null,
|
||||
1,
|
||||
null
|
||||
])
|
||||
expect(channel.read('session-1')).toBeNull()
|
||||
})
|
||||
|
||||
// Claude's release path settles the last child after the adapter let go of the session, so no
|
||||
@@ -99,17 +87,25 @@ describe('the background-task channel', () => {
|
||||
it('hides the strip when its last child settles after the provider let go of the session', () => {
|
||||
let views: AgentChildWorkView[] = [child]
|
||||
let held: Stops = { supportsTaskStop: true, supportsStopAll: true }
|
||||
const { sessions, sent, channel, session } = channelOver(
|
||||
const { sessions, channel, session } = channelOver(
|
||||
() => views,
|
||||
() => held
|
||||
)
|
||||
sessions.set('session-1', session)
|
||||
channel.publish('session-1')
|
||||
expect(sent.mock.calls.at(-1)?.[1]?.children).toHaveLength(1)
|
||||
expect(channel.read('session-1')?.children).toHaveLength(1)
|
||||
|
||||
held = undefined
|
||||
views = [{ ...child, state: 'done', membership: 'settled', outcome: 'unknown', settledAt: 2 }]
|
||||
expect(channel.read('session-1')).toBeNull()
|
||||
})
|
||||
|
||||
it('republishes only for an open conversation', () => {
|
||||
const { sessions, republished, channel, session } = channelOver()
|
||||
channel.publish('session-1')
|
||||
expect(sent.mock.calls.at(-1)?.[1]).toBeNull()
|
||||
expect(republished).not.toHaveBeenCalled()
|
||||
|
||||
sessions.set('session-1', session)
|
||||
channel.publish('session-1')
|
||||
expect(republished).toHaveBeenCalledWith('session-1', expect.any(Number))
|
||||
})
|
||||
})
|
||||
|
||||
+18
-41
@@ -22,10 +22,9 @@ import { structuredAgentSessionConversationFence } from './structured-agent-sess
|
||||
|
||||
/** The chat strip's roster: the host's running child records for one session, the same selection
|
||||
* the session list carries, with the legacy task rows an older client reads derived from them. No
|
||||
* running child is `null`, and the strip hides. */
|
||||
* running child is `null`, and the strip hides. Each subscriber's frames carry it; see
|
||||
* `agent-session-subscriber-frame-fields`. */
|
||||
export class StructuredAgentSessionBackgroundTaskChannel {
|
||||
private readonly published = new Map<string, string>()
|
||||
|
||||
constructor(
|
||||
private readonly deps: StructuredAgentSessionHostDeps,
|
||||
private readonly sessions: StructuredAgentSessionConversations,
|
||||
@@ -35,10 +34,7 @@ export class StructuredAgentSessionBackgroundTaskChannel {
|
||||
sessionId: string
|
||||
) => Promise<StructuredAgentSessionHostSession>,
|
||||
private readonly readChildWork: (sessionId: string) => AgentChildWorkView[] | undefined
|
||||
) {
|
||||
// A closed conversation's last roster is not kept for the host's lifetime.
|
||||
sessions.observeClose((sessionId) => this.published.delete(sessionId))
|
||||
}
|
||||
) {}
|
||||
|
||||
/** `scope` is for in-process readers; a wire request reads every agent's rows. */
|
||||
async history(
|
||||
@@ -52,7 +48,7 @@ export class StructuredAgentSessionBackgroundTaskChannel {
|
||||
request,
|
||||
scope
|
||||
})
|
||||
const backgroundTasks = this.state(request.sessionId)
|
||||
const backgroundTasks = this.read(request.sessionId)
|
||||
const queue = tryReadQueuePublication(journal)
|
||||
const hostNow = this.deps.now?.() ?? Date.now()
|
||||
return {
|
||||
@@ -65,7 +61,7 @@ export class StructuredAgentSessionBackgroundTaskChannel {
|
||||
...(queue !== undefined
|
||||
? { queuedMessages: queue.queuedMessages, queuePause: queue.queuePause }
|
||||
: {}),
|
||||
...(backgroundTasks !== undefined ? { backgroundTasks } : {})
|
||||
backgroundTasks
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -73,51 +69,32 @@ export class StructuredAgentSessionBackgroundTaskChannel {
|
||||
/** Resolves once the conversation is open and the subscriber holds its opening frame. */
|
||||
async subscribe(input: AgentSessionSubscribeInput): Promise<() => void> {
|
||||
const session = await this.conversation(input.sessionId)
|
||||
const backgroundTasks = this.state(input.sessionId)
|
||||
return this.subscribers.open({
|
||||
...input,
|
||||
journal: session.journal,
|
||||
fence: structuredAgentSessionConversationFence(this.deps.store, input.sessionId),
|
||||
...(backgroundTasks !== undefined ? { backgroundTasks } : {})
|
||||
fence: structuredAgentSessionConversationFence(this.deps.store, input.sessionId)
|
||||
})
|
||||
}
|
||||
|
||||
/** Re-read after the session's child records changed; an unchanged roster sends nothing. */
|
||||
/** The session's child records changed; subscribers whose roster differs get the new one. */
|
||||
publish(sessionId: string): void {
|
||||
const session = this.sessions.get(sessionId)
|
||||
// Explicit null once a roster was sent, not silence: a reader keeps its last roster on
|
||||
// `undefined`, and a closing provider stops answering before its records are gone.
|
||||
const read = this.state(sessionId)
|
||||
const state = read === undefined && this.published.has(sessionId) ? null : read
|
||||
if (!session || state === undefined) {
|
||||
return
|
||||
if (this.sessions.get(sessionId)) {
|
||||
this.subscribers.republishBackgroundTasks(
|
||||
sessionId,
|
||||
structuredAgentSessionConversationFence(this.deps.store, sessionId)
|
||||
)
|
||||
}
|
||||
const fingerprint = JSON.stringify(state)
|
||||
if (this.published.get(sessionId) === fingerprint) {
|
||||
return
|
||||
}
|
||||
// "None" is remembered too, so a session with no children sends it once, not on every change;
|
||||
// the entry goes when the conversation closes.
|
||||
this.published.set(sessionId, fingerprint)
|
||||
this.subscribers.backgroundTasks(
|
||||
sessionId,
|
||||
state,
|
||||
structuredAgentSessionConversationFence(this.deps.store, sessionId)
|
||||
)
|
||||
}
|
||||
|
||||
private state(sessionId: string): AgentSessionBackgroundTaskState | null | undefined {
|
||||
/** Always an answer, never "unknown": the child records are the roster whether or not a provider
|
||||
* holds the session, and a parent row the store lacks has no children in it. */
|
||||
read(sessionId: string): AgentSessionBackgroundTaskState | null {
|
||||
const session = this.sessions.get(sessionId)
|
||||
const stored = session ? this.readChildWork(sessionId) : undefined
|
||||
if (!session || stored === undefined) {
|
||||
return undefined
|
||||
const views = structuredRunningChildWork((session && this.readChildWork(sessionId)) ?? [])
|
||||
if (!session || views.length === 0) {
|
||||
return null
|
||||
}
|
||||
const views = structuredRunningChildWork(stored)
|
||||
const stops = this.deps.adapter.backgroundTaskStops?.(sessionId)
|
||||
if (views.length === 0) {
|
||||
// As before: a session no live provider holds says nothing, a live one says "none".
|
||||
return stops === undefined ? undefined : null
|
||||
}
|
||||
const { tasks, settledTasks } = structuredChildWorkLegacyTasks(views, session.params.provider)
|
||||
return {
|
||||
state: 'monitoring',
|
||||
|
||||
+2
-2
@@ -78,8 +78,8 @@ describe("the host's strip channel and summary list the same running children",
|
||||
})
|
||||
expect(summaries.at(-1)).toEqual(['run', 'owner', 'shell'])
|
||||
|
||||
// Nothing runs: no strip at all (a host that holds no provider says nothing).
|
||||
// Nothing runs: "none", even with no provider holding the session, so a reader drops its strip.
|
||||
records = [finished('done', 5), finished('run', 7)]
|
||||
expect(await strip()).toBeUndefined()
|
||||
expect(await strip()).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
+10
-6
@@ -2,6 +2,7 @@ import { AgentSessionRefusalError } from '../../../shared/agent-session-wire-ref
|
||||
import type { AgentChildWorkEvidence } from '../../../shared/agent-status-child-work-evidence'
|
||||
import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view'
|
||||
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
|
||||
import type { SubscriberFieldHooks } from './agent-session-subscriber-frame-fields'
|
||||
import { AgentSessionSubscribers } from './structured-agent-session-subscribers'
|
||||
import { tryReadQueuePublication } from './structured-agent-session-queued-publication'
|
||||
import type {
|
||||
@@ -32,17 +33,19 @@ export class StructuredAgentSessionClientDelivery {
|
||||
private readonly sessions: Map<string, StructuredAgentSessionHostSession>,
|
||||
now: () => number,
|
||||
private readonly deps: () => StructuredAgentSessionHostDeps,
|
||||
private readonly onJournalActivity?: (sessionId: string) => void,
|
||||
onAgentStarted?: (sessionId: string) => void,
|
||||
private readonly onJournalActivity: (sessionId: string) => void,
|
||||
onAgentStarted: (sessionId: string) => void,
|
||||
/** A session's child records changed; the chat strip republishes from them. */
|
||||
onChildWorkChanged?: (sessionId: string) => void
|
||||
onChildWorkChanged: (sessionId: string) => void,
|
||||
// Required: an opening frame without the roster reads as "no tasks" to current clients.
|
||||
readBackgroundTasks: NonNullable<SubscriberFieldHooks['readBackgroundTasks']>
|
||||
) {
|
||||
this.statusFeed = createStructuredAgentSessionHostStatusFeed({
|
||||
sessions,
|
||||
now,
|
||||
deps,
|
||||
...(onAgentStarted ? { onAgentStarted } : {}),
|
||||
...(onChildWorkChanged ? { onChildWorkChanged } : {})
|
||||
onAgentStarted,
|
||||
onChildWorkChanged
|
||||
})
|
||||
this.turnCompletionFeed = new StructuredAgentSessionTurnCompletionFeed({
|
||||
sessions,
|
||||
@@ -57,6 +60,7 @@ export class StructuredAgentSessionClientDelivery {
|
||||
readCommands: (sessionId) => this.readCommands(sessionId),
|
||||
readQueuePublication: (sessionId) =>
|
||||
tryReadQueuePublication(sessions.get(sessionId)?.journal),
|
||||
readBackgroundTasks,
|
||||
onJournalPublished: (sessionId, journal) => this.publishJournal(sessionId, journal)
|
||||
})
|
||||
}
|
||||
@@ -131,7 +135,7 @@ export class StructuredAgentSessionClientDelivery {
|
||||
// subscribed, which is the whole reason a backgrounded chat can complete at all. After the
|
||||
// status publish, so it reads the projection that publish cached.
|
||||
this.turnCompletionFeed.observe(sessionId, journal)
|
||||
this.onJournalActivity?.(sessionId)
|
||||
this.onJournalActivity(sessionId)
|
||||
}
|
||||
|
||||
private requireJournal(sessionId: string): AgentSessionJournal {
|
||||
|
||||
+12
-3
@@ -3,6 +3,7 @@ import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { expect, it, vi } from 'vitest'
|
||||
import type {
|
||||
AgentSessionBackgroundTaskState,
|
||||
AgentSessionSlashCommand,
|
||||
AgentSessionSubscribeEvent
|
||||
} from '../../../shared/agent-session-wire'
|
||||
@@ -65,7 +66,15 @@ it('delivers catalog changes through existing frames without resending them on o
|
||||
let commands: AgentSessionSlashCommand[] | undefined = [
|
||||
{ name: 'loaded', kind: 'command', kindUnspecified: true }
|
||||
]
|
||||
const subscribers = new AgentSessionSubscribers({ readCommands: () => commands })
|
||||
let roster: AgentSessionBackgroundTaskState | null = null
|
||||
const subscribers = new AgentSessionSubscribers({
|
||||
readCommands: () => commands,
|
||||
readBackgroundTasks: () => roster
|
||||
})
|
||||
const republishRoster = (): void => {
|
||||
roster = roster ? null : { state: 'monitoring' }
|
||||
subscribers.republishBackgroundTasks(sessionId, 7)
|
||||
}
|
||||
const close = subscribers.open({
|
||||
id: 'one',
|
||||
sessionId,
|
||||
@@ -75,13 +84,13 @@ it('delivers catalog changes through existing frames without resending them on o
|
||||
})
|
||||
expect(state.commands).toEqual(commands)
|
||||
for (let i = 0; i < 25; i++) {
|
||||
subscribers.backgroundTasks(sessionId, null, 7)
|
||||
republishRoster()
|
||||
}
|
||||
coalescer.flush()
|
||||
expect(events.filter((event) => 'commands' in event)).toHaveLength(1)
|
||||
commands = []
|
||||
subscribers.publish(sessionId, journal)
|
||||
subscribers.backgroundTasks(sessionId, null, 7)
|
||||
republishRoster()
|
||||
coalescer.flush()
|
||||
expect(state.commands).toEqual([])
|
||||
expect(events.filter((event) => 'commands' in event)).toHaveLength(2)
|
||||
|
||||
@@ -18,7 +18,6 @@ export class StructuredAgentSessionConversations extends Map<
|
||||
StructuredAgentSessionHostSession
|
||||
> {
|
||||
private readonly activity = new Map<string, number>()
|
||||
private readonly closeObservers = new Set<(sessionId: string) => void>()
|
||||
|
||||
constructor(
|
||||
private readonly delivery: {
|
||||
@@ -64,18 +63,7 @@ export class StructuredAgentSessionConversations extends Map<
|
||||
|
||||
override delete(sessionId: string): boolean {
|
||||
this.activity.delete(sessionId)
|
||||
const deleted = super.delete(sessionId)
|
||||
if (deleted) {
|
||||
for (const observer of this.closeObservers) {
|
||||
observer(sessionId)
|
||||
}
|
||||
}
|
||||
return deleted
|
||||
}
|
||||
|
||||
/** Told when a conversation leaves the map, so state kept per conversation dies with it. */
|
||||
observeClose(observer: (sessionId: string) => void): void {
|
||||
this.closeObservers.add(observer)
|
||||
return super.delete(sessionId)
|
||||
}
|
||||
|
||||
touch(sessionId: string): void {
|
||||
|
||||
@@ -81,14 +81,13 @@ export class StructuredAgentSessionHost {
|
||||
() => this.deps,
|
||||
(sessionId) => this.queued.onJournalActivity(sessionId),
|
||||
(sessionId) => this.restartResume.onAgentStarted(sessionId),
|
||||
(sessionId) => this.backgroundTasks.publish(sessionId)
|
||||
(sessionId) => this.backgroundTasks.publish(sessionId),
|
||||
(sessionId) => this.backgroundTasks.read(sessionId)
|
||||
)
|
||||
private readonly subscribers = this.clientDelivery.subscribers
|
||||
private readonly tasks = new StructuredAgentSessionTaskQueue()
|
||||
private readonly runtimeState: StructuredAgentSessionHostRuntimeState
|
||||
private readonly reconcileLeases: (
|
||||
sessionId: string
|
||||
) => Promise<SessionWire.AgentSessionWireRefusal | null>
|
||||
private readonly reconcileLeases: ReturnType<typeof createRestartReconciler>
|
||||
private readonly restore: ReturnType<typeof createStructuredAgentSessionHostRestore>
|
||||
private readonly lifetime: StructuredAgentSessionConversationLifetime
|
||||
private readonly conversationDelivery: ReturnType<
|
||||
|
||||
+197
@@ -0,0 +1,197 @@
|
||||
// The strip's roster rides each subscriber's frames: stated on its first frame and on every
|
||||
// hydrating frame, re-sent when it differs from what that subscriber last got, and never read for
|
||||
// an ordinary journal batch.
|
||||
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { AGENT_JOURNAL_THREAD_SCOPE } from '../../../shared/agent-session-journal-types'
|
||||
import { claudeProviderHandle } from '../../../shared/agent-session-provider-handle-encoding'
|
||||
import type {
|
||||
AgentSessionBackgroundTaskState,
|
||||
AgentSessionSubscribeEvent
|
||||
} from '../../../shared/agent-session-wire'
|
||||
import {
|
||||
EMPTY_STRUCTURED_AGENT_SESSION,
|
||||
reduceStructuredAgentSession,
|
||||
type StructuredAgentSessionState
|
||||
} from '../../../shared/structured-agent-session-reducer'
|
||||
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
|
||||
import { createTrackedJournalOpener } from '../agent-session-journal/journal-host-database-test-support'
|
||||
import { AgentSessionSubscribers } from './structured-agent-session-subscribers'
|
||||
|
||||
const SESSION = 'roster-session'
|
||||
const MONITORING: AgentSessionBackgroundTaskState = {
|
||||
state: 'monitoring',
|
||||
tasks: [{ id: 'bbpijar3m', kind: 'command', description: 'watch CI' }]
|
||||
}
|
||||
|
||||
let root: string
|
||||
const journals = createTrackedJournalOpener()
|
||||
|
||||
beforeEach(async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-subscriber-roster-'))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await journals.closeAll()
|
||||
await rm(root, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
function openJournal(name: string): Promise<AgentSessionJournal> {
|
||||
return journals.open({
|
||||
identity: {
|
||||
sessionId: SESSION,
|
||||
workspaceId: 'workspace-1',
|
||||
hostId: 'local',
|
||||
agent: 'claude',
|
||||
providerHandle: claudeProviderHandle('provider-1', null)
|
||||
},
|
||||
stateDirectory: join(root, name)
|
||||
})
|
||||
}
|
||||
|
||||
async function appendRow(journal: AgentSessionJournal, id: string): Promise<void> {
|
||||
await journal.appendItem(
|
||||
{ provider: 'orca', clientMessageId: id },
|
||||
{ kind: 'status', text: id },
|
||||
{ fence: 1, turnScope: AGENT_JOURNAL_THREAD_SCOPE }
|
||||
)
|
||||
}
|
||||
|
||||
function rosterOf(event: AgentSessionSubscribeEvent | undefined) {
|
||||
return event && event.type !== 'end' ? event.backgroundTasks : undefined
|
||||
}
|
||||
|
||||
describe('the background-task roster on subscriber frames', () => {
|
||||
it('a pane resuming from its cursor drops a task that ended while it was away (#24227)', async () => {
|
||||
const journal = await openJournal('resume')
|
||||
let roster: AgentSessionBackgroundTaskState | null = MONITORING
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks: () => roster })
|
||||
let pane: StructuredAgentSessionState = EMPTY_STRUCTURED_AGENT_SESSION
|
||||
const toPane = (event: AgentSessionSubscribeEvent): void => {
|
||||
pane = reduceStructuredAgentSession(pane, { type: 'event', event })
|
||||
}
|
||||
|
||||
const close = subscribers.open({
|
||||
id: 'pane-1',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
emit: toPane
|
||||
})
|
||||
expect(pane.backgroundTasks).toEqual(MONITORING)
|
||||
// The pane is hidden; the task ends with nobody listening, and the turn keeps journaling.
|
||||
close()
|
||||
roster = null
|
||||
subscribers.republishBackgroundTasks(SESSION, 1)
|
||||
await appendRow(journal, 'after')
|
||||
subscribers.publish(SESSION, journal)
|
||||
|
||||
// The idle sweep stopped the provider; the pane comes back from its cursor.
|
||||
subscribers.open({
|
||||
id: 'pane-2',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
cursor: pane.cursor!,
|
||||
emit: toPane
|
||||
})
|
||||
|
||||
expect(pane.backgroundTasks).toBeNull()
|
||||
})
|
||||
|
||||
it('states the roster on a caught-up resume, which carries no rows', async () => {
|
||||
const journal = await openJournal('caught-up')
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks: () => MONITORING })
|
||||
const events: AgentSessionSubscribeEvent[] = []
|
||||
|
||||
subscribers.open({
|
||||
id: 'pane',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
cursor: journal.cursor(),
|
||||
emit: (event) => events.push(event)
|
||||
})
|
||||
|
||||
expect(events).toHaveLength(1)
|
||||
expect(rosterOf(events[0])).toEqual(MONITORING)
|
||||
})
|
||||
|
||||
it('carries the current roster on snapshot frames, so a running task stays shown', async () => {
|
||||
const journal = await openJournal('replay')
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks: () => MONITORING })
|
||||
const events: AgentSessionSubscribeEvent[] = []
|
||||
subscribers.open({
|
||||
id: 'pane',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
emit: (e) => events.push(e)
|
||||
})
|
||||
|
||||
subscribers.snapshot(SESSION, journal, 2)
|
||||
|
||||
expect(events.map((event) => [event.type, rosterOf(event)])).toEqual([
|
||||
['snapshot', MONITORING],
|
||||
['snapshot', MONITORING]
|
||||
])
|
||||
})
|
||||
|
||||
it('neither reads nor resends the roster for an ordinary journal batch', async () => {
|
||||
const journal = await openJournal('ordinary')
|
||||
const readBackgroundTasks = vi.fn(() => MONITORING)
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks })
|
||||
const events: AgentSessionSubscribeEvent[] = []
|
||||
subscribers.open({
|
||||
id: 'pane',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
emit: (e) => events.push(e)
|
||||
})
|
||||
const reads = readBackgroundTasks.mock.calls.length
|
||||
|
||||
await appendRow(journal, 'token')
|
||||
subscribers.publish(SESSION, journal)
|
||||
|
||||
expect(readBackgroundTasks).toHaveBeenCalledTimes(reads)
|
||||
expect(events.at(-1)).toMatchObject({ type: 'batch' })
|
||||
expect(rosterOf(events.at(-1))).toBeUndefined()
|
||||
})
|
||||
|
||||
it('republishes to each subscriber whose last roster differs, and to no other', async () => {
|
||||
const journal = await openJournal('dedup')
|
||||
let roster: AgentSessionBackgroundTaskState | null = null
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks: () => roster })
|
||||
const first: AgentSessionSubscribeEvent[] = []
|
||||
const second: AgentSessionSubscribeEvent[] = []
|
||||
subscribers.open({
|
||||
id: 'first',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
emit: (e) => first.push(e)
|
||||
})
|
||||
|
||||
roster = MONITORING
|
||||
subscribers.republishBackgroundTasks(SESSION, 1)
|
||||
// Opens holding the current roster already.
|
||||
subscribers.open({
|
||||
id: 'second',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
emit: (e) => second.push(e)
|
||||
})
|
||||
subscribers.republishBackgroundTasks(SESSION, 1)
|
||||
expect([first.length, second.length]).toEqual([2, 1])
|
||||
|
||||
roster = null
|
||||
subscribers.republishBackgroundTasks(SESSION, 1)
|
||||
expect([rosterOf(first.at(-1)), rosterOf(second.at(-1))]).toEqual([null, null])
|
||||
expect([first.length, second.length]).toEqual([3, 2])
|
||||
})
|
||||
})
|
||||
+18
-8
@@ -7,6 +7,7 @@ import {
|
||||
AGENT_SESSION_JOURNAL_SCHEMA_VERSION
|
||||
} from '../../../shared/agent-session-journal-types'
|
||||
import type {
|
||||
AgentSessionBackgroundTaskState,
|
||||
AgentSessionStatusEvent,
|
||||
AgentSessionSubscribeEvent
|
||||
} from '../../../shared/agent-session-wire'
|
||||
@@ -115,8 +116,12 @@ describe('AgentSessionSubscribers', () => {
|
||||
stateDirectory: join(root, 'clock-journal')
|
||||
})
|
||||
let now = 1_000
|
||||
let roster: AgentSessionBackgroundTaskState | null = null
|
||||
const events: AgentSessionSubscribeEvent[] = []
|
||||
const subscribers = new AgentSessionSubscribers({ now: () => (now += 1) })
|
||||
const subscribers = new AgentSessionSubscribers({
|
||||
now: () => (now += 1),
|
||||
readBackgroundTasks: () => roster
|
||||
})
|
||||
const emit = (event: AgentSessionSubscribeEvent): void => {
|
||||
events.push(event)
|
||||
}
|
||||
@@ -128,7 +133,11 @@ describe('AgentSessionSubscribers', () => {
|
||||
{ fence: 1, turnScope: AGENT_JOURNAL_THREAD_SCOPE }
|
||||
)
|
||||
subscribers.publish(SESSION, journal)
|
||||
subscribers.backgroundTasks(SESSION, null, 1)
|
||||
roster = {
|
||||
state: 'monitoring',
|
||||
tasks: [{ id: 'task-1', kind: 'command', description: 'run the build' }]
|
||||
}
|
||||
subscribers.republishBackgroundTasks(SESSION, 1)
|
||||
subscribers.snapshot(SESSION, journal, 1)
|
||||
|
||||
expect(events.map((event) => ('hostNow' in event ? event.hostNow : null))).toEqual([
|
||||
@@ -296,23 +305,24 @@ describe('AgentSessionSubscribers', () => {
|
||||
},
|
||||
stateDirectory: join(root, 'background-journal')
|
||||
})
|
||||
const subscribers = new AgentSessionSubscribers()
|
||||
let roster: AgentSessionBackgroundTaskState | null = null
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks: () => roster })
|
||||
const events: AgentSessionSubscribeEvent[] = []
|
||||
subscribers.open({
|
||||
id: 'subscriber-1',
|
||||
sessionId: SESSION,
|
||||
journal,
|
||||
fence: 1,
|
||||
backgroundTasks: null,
|
||||
emit: (event) => events.push(event)
|
||||
})
|
||||
const cursor = journal.cursor()
|
||||
|
||||
const backgroundTasks = {
|
||||
state: 'monitoring' as const,
|
||||
tasks: [{ id: 'task-1', kind: 'command' as const, description: 'run the build' }]
|
||||
const backgroundTasks: AgentSessionBackgroundTaskState = {
|
||||
state: 'monitoring',
|
||||
tasks: [{ id: 'task-1', kind: 'command', description: 'run the build' }]
|
||||
}
|
||||
subscribers.backgroundTasks(SESSION, backgroundTasks, 2)
|
||||
roster = backgroundTasks
|
||||
subscribers.republishBackgroundTasks(SESSION, 2)
|
||||
|
||||
expect(journal.cursor()).toEqual(cursor)
|
||||
expect(events.at(-1)).toEqual({
|
||||
|
||||
@@ -6,12 +6,15 @@
|
||||
|
||||
import type { AgentJournalCursor } from '../../../shared/agent-session-journal-types'
|
||||
import type {
|
||||
AgentSessionBackgroundTaskState,
|
||||
AgentSessionSlashCommand,
|
||||
AgentSessionSubscribeEvent,
|
||||
AgentSessionTurnActivity
|
||||
} from '../../../shared/agent-session-wire'
|
||||
import { buildSubscriberFrame } from './agent-session-subscriber-frame-fields'
|
||||
import {
|
||||
backgroundTaskFingerprint,
|
||||
buildSubscriberFrame,
|
||||
type SubscriberFieldHooks
|
||||
} from './agent-session-subscriber-frame-fields'
|
||||
import type { QueuePublication } from './structured-agent-session-queued-publication'
|
||||
import { deliverToSubscriber } from './agent-session-subscriber-catch-up'
|
||||
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
|
||||
@@ -37,6 +40,8 @@ export type Subscriber = {
|
||||
/** The last draft list actually SENT — never advanced on a page that withheld
|
||||
* it, or the final replacement would be suppressed by the identity dedup. */
|
||||
queuePublication?: QueuePublication
|
||||
/** Fingerprint of the background-task roster last SENT. */
|
||||
backgroundTasks?: string
|
||||
}
|
||||
|
||||
export type AgentSessionSubscribersHooks = {
|
||||
@@ -44,6 +49,7 @@ export type AgentSessionSubscribersHooks = {
|
||||
/** Revision-stable per emit: an unchanged list keeps its reference, so token
|
||||
* streams never re-serialize it; any draft-table write changes it. */
|
||||
readQueuePublication?: (sessionId: string) => QueuePublication | undefined
|
||||
readBackgroundTasks?: SubscriberFieldHooks['readBackgroundTasks']
|
||||
/** Fires after publications that can change journal content. */
|
||||
onJournalPublished?: (sessionId: string, journal: AgentSessionJournal) => void
|
||||
now?: () => number
|
||||
@@ -70,7 +76,6 @@ export class AgentSessionSubscribers {
|
||||
fence: number
|
||||
emit: AgentSessionSubscriberEmit
|
||||
cursor?: AgentJournalCursor
|
||||
backgroundTasks?: AgentSessionBackgroundTaskState | null
|
||||
}): () => void {
|
||||
const liveCursor = input.journal.cursor()
|
||||
const subscriber: Subscriber = {
|
||||
@@ -86,7 +91,7 @@ export class AgentSessionSubscribers {
|
||||
|
||||
const hostNow = this.now()
|
||||
if (input.cursor) {
|
||||
this.deliver(subscriber, input.journal, hostNow, true, input.backgroundTasks)
|
||||
this.deliver(subscriber, input.journal, hostNow, true)
|
||||
} else {
|
||||
const page = readAgentSessionHydrationPage(input.journal, input.fence)
|
||||
this.emit(subscriber, {
|
||||
@@ -95,7 +100,6 @@ export class AgentSessionSubscribers {
|
||||
page,
|
||||
fence: input.fence,
|
||||
hostNow,
|
||||
...(input.backgroundTasks !== undefined ? { backgroundTasks: input.backgroundTasks } : {}),
|
||||
...this.activityField(input.sessionId)
|
||||
})
|
||||
subscriber.cursor = page.liveCursor ?? page.window.nextCursor
|
||||
@@ -132,7 +136,7 @@ export class AgentSessionSubscribers {
|
||||
}
|
||||
const hostNow = this.now()
|
||||
for (const subscriber of this.subscribers(sessionId)) {
|
||||
this.deliver(subscriber, journal, hostNow, false, undefined, activity)
|
||||
this.deliver(subscriber, journal, hostNow, false, activity)
|
||||
}
|
||||
if (activity === undefined) {
|
||||
this.hooks.onJournalPublished?.(sessionId, journal)
|
||||
@@ -140,12 +144,7 @@ export class AgentSessionSubscribers {
|
||||
}
|
||||
|
||||
/** Every subscriber back to a bounded tail page. */
|
||||
snapshot(
|
||||
sessionId: string,
|
||||
journal: AgentSessionJournal,
|
||||
fence: number,
|
||||
backgroundTasks?: AgentSessionBackgroundTaskState | null
|
||||
): void {
|
||||
snapshot(sessionId: string, journal: AgentSessionJournal, fence: number): void {
|
||||
const page = readAgentSessionHydrationPage(journal, fence)
|
||||
const hostNow = this.now()
|
||||
for (const subscriber of this.subscribers(sessionId)) {
|
||||
@@ -155,7 +154,6 @@ export class AgentSessionSubscribers {
|
||||
page,
|
||||
fence,
|
||||
hostNow,
|
||||
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
|
||||
...this.activityField(sessionId)
|
||||
})
|
||||
subscriber.cursor = page.liveCursor ?? page.window.nextCursor
|
||||
@@ -164,13 +162,19 @@ export class AgentSessionSubscribers {
|
||||
this.hooks.onJournalPublished?.(sessionId, journal)
|
||||
}
|
||||
|
||||
backgroundTasks(
|
||||
sessionId: string,
|
||||
state: AgentSessionBackgroundTaskState | null,
|
||||
fence: number
|
||||
): void {
|
||||
/** Re-sends the strip's roster to each subscriber whose last one differs; the session's child
|
||||
* records changed. */
|
||||
republishBackgroundTasks(sessionId: string, fence: number): void {
|
||||
const state = this.hooks.readBackgroundTasks?.(sessionId)
|
||||
if (state === undefined) {
|
||||
return
|
||||
}
|
||||
const fingerprint = backgroundTaskFingerprint(state)
|
||||
const hostNow = this.now()
|
||||
for (const subscriber of this.subscribers(sessionId)) {
|
||||
if (subscriber.backgroundTasks === fingerprint) {
|
||||
continue
|
||||
}
|
||||
this.emit(subscriber, {
|
||||
type: 'batch',
|
||||
sessionId,
|
||||
@@ -210,7 +214,6 @@ export class AgentSessionSubscribers {
|
||||
journal: AgentSessionJournal,
|
||||
hostNow: number,
|
||||
emitCheckpoint = false,
|
||||
backgroundTasks?: AgentSessionBackgroundTaskState | null,
|
||||
activity?: AgentSessionTurnActivity | null
|
||||
): void {
|
||||
deliverToSubscriber(
|
||||
@@ -220,7 +223,7 @@ export class AgentSessionSubscribers {
|
||||
isActive: (target) => this.isActive(target),
|
||||
activity: (sessionId) => this.activityField(sessionId).activity
|
||||
},
|
||||
{ subscriber, journal, hostNow, emitCheckpoint, backgroundTasks, activity }
|
||||
{ subscriber, journal, hostNow, emitCheckpoint, activity }
|
||||
)
|
||||
}
|
||||
|
||||
@@ -248,6 +251,9 @@ export class AgentSessionSubscribers {
|
||||
if (built.attachedQueued) {
|
||||
subscriber.queuePublication = built.queued
|
||||
}
|
||||
if (built.backgroundTasks !== undefined) {
|
||||
subscriber.backgroundTasks = built.backgroundTasks
|
||||
}
|
||||
} catch {
|
||||
this.drop(subscriber)
|
||||
}
|
||||
|
||||
+8
-3
@@ -6,7 +6,10 @@ import type {
|
||||
AgentJournalItemBody,
|
||||
AgentJournalItemIdentity
|
||||
} from '../../../shared/agent-session-journal-types'
|
||||
import type { AgentSessionSubscribeEvent } from '../../../shared/agent-session-wire'
|
||||
import type {
|
||||
AgentSessionBackgroundTaskState,
|
||||
AgentSessionSubscribeEvent
|
||||
} from '../../../shared/agent-session-wire'
|
||||
import { REMOTE_RUNTIME_MAX_OUTBOUND_JSON_BYTES } from '../../../shared/remote-runtime-memory-limits'
|
||||
import { mobileE2EETextPayloadAdmissionBytes } from '../../runtime/rpc/mobile-e2ee-outbound-admission'
|
||||
import {
|
||||
@@ -62,7 +65,8 @@ describe('structured agent-session outbound admission', () => {
|
||||
REMOTE_RUNTIME_MAX_OUTBOUND_JSON_BYTES
|
||||
)
|
||||
|
||||
const subscribers = new AgentSessionSubscribers()
|
||||
let roster: AgentSessionBackgroundTaskState | null = null
|
||||
const subscribers = new AgentSessionSubscribers({ readBackgroundTasks: () => roster })
|
||||
const initial: AgentSessionSubscribeEvent[] = []
|
||||
const dispose = subscribers.open({
|
||||
id: 'initial',
|
||||
@@ -75,7 +79,8 @@ describe('structured agent-session outbound admission', () => {
|
||||
expect(initial[0]).toMatchObject({ type: 'snapshot', page: { hasOlder: true } })
|
||||
expectAdmitted(initial[0])
|
||||
|
||||
subscribers.backgroundTasks(SESSION, null, 2)
|
||||
roster = { state: 'monitoring' }
|
||||
subscribers.republishBackgroundTasks(SESSION, 2)
|
||||
subscribers.snapshot(SESSION, journal, 2)
|
||||
expect(initial.slice(1)).toHaveLength(2)
|
||||
initial.slice(1).forEach(expectAdmitted)
|
||||
|
||||
@@ -222,7 +222,7 @@ function createReadOwner(
|
||||
apply({ type: 'loading' })
|
||||
}
|
||||
const transport = startStructuredAgentSessionReadTransport({
|
||||
applyEvent: (event) => apply({ type: 'event', event }),
|
||||
applyEvent: (event, options) => apply({ type: 'event', event, ...options }),
|
||||
applyError: (message, refusal) => apply({ type: 'error', message, refusal }),
|
||||
getCursor: () => snapshot.state.cursor,
|
||||
onHistoryReadInvalidated: invalidateOlderPages,
|
||||
|
||||
+47
-1
@@ -73,7 +73,10 @@ describe('structured agent-session read transport generations', () => {
|
||||
})
|
||||
})
|
||||
|
||||
function start(applyEvent: (event: AgentSessionSubscribeEvent) => void, applyError = vi.fn()) {
|
||||
function start(
|
||||
applyEvent: Parameters<typeof startStructuredAgentSessionReadTransport>[0]['applyEvent'],
|
||||
applyError = vi.fn()
|
||||
) {
|
||||
return startStructuredAgentSessionReadTransport({
|
||||
applyEvent,
|
||||
applyError,
|
||||
@@ -185,6 +188,49 @@ describe('structured agent-session read transport generations', () => {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it("applies each subscription's first batch alone and marked, and coalesces the rest", async () => {
|
||||
vi.useFakeTimers()
|
||||
try {
|
||||
const applied: [number | undefined, boolean][] = []
|
||||
const batch = (sequence: number): AgentSessionSubscribeEvent => ({
|
||||
type: 'batch',
|
||||
sessionId: 'session-a',
|
||||
batch: {
|
||||
cursor: { epoch: 'epoch-a', sequence },
|
||||
items: [],
|
||||
removedItemIds: [],
|
||||
submissions: []
|
||||
}
|
||||
})
|
||||
const transport = start((event, options) => {
|
||||
applied.push([
|
||||
event.type === 'batch' ? event.batch.cursor.sequence : undefined,
|
||||
options?.opensSubscription === true
|
||||
])
|
||||
})
|
||||
await flushPromises()
|
||||
attempts[0].onEvent(batch(1))
|
||||
attempts[0].onEvent(batch(2))
|
||||
attempts[0].onEvent(batch(3))
|
||||
expect(applied).toEqual([[1, true]])
|
||||
await vi.advanceTimersByTimeAsync(60)
|
||||
expect(applied).toEqual([
|
||||
[1, true],
|
||||
[3, false]
|
||||
])
|
||||
|
||||
attempts[0].closed.resolve({ unsubscribe: attempts[0].unsubscribe })
|
||||
await flushPromises()
|
||||
attempts[0].onClose()
|
||||
await vi.advanceTimersByTimeAsync(750)
|
||||
attempts[1].onEvent(batch(4))
|
||||
expect(applied.at(-1)).toEqual([4, true])
|
||||
transport.dispose()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
// The refusal a host raises for a session it holds no object for — after a chat close, or before
|
||||
|
||||
@@ -78,7 +78,7 @@ function createReconnectScheduler(args: { shouldStop: () => boolean; reconnect:
|
||||
}
|
||||
|
||||
export function startStructuredAgentSessionReadTransport(args: {
|
||||
applyEvent: (event: AgentSessionSubscribeEvent) => void
|
||||
applyEvent: (event: AgentSessionSubscribeEvent, options?: { opensSubscription: true }) => void
|
||||
/** `message` is the failure's own text, for logs; `refusal` is what a surface words. */
|
||||
applyError: (message: string, refusal?: AgentSessionRefusalReference) => void
|
||||
getCursor: () => AgentJournalCursor | null
|
||||
@@ -97,6 +97,8 @@ export function startStructuredAgentSessionReadTransport(args: {
|
||||
let unattachedSince: number | null = null
|
||||
let opening = false
|
||||
let openGeneration = 0
|
||||
/** The open generation whose first frame has not arrived yet. */
|
||||
let openingFrameGeneration = 0
|
||||
let stateGeneration = 0
|
||||
let unsubscribe = (): void => {}
|
||||
let shouldStopCoalescedEvent = (): boolean => true
|
||||
@@ -182,6 +184,19 @@ export function startStructuredAgentSessionReadTransport(args: {
|
||||
connected = false
|
||||
reconnectScheduler.schedule()
|
||||
}
|
||||
const opensSubscription = event.type !== 'end' && openingFrameGeneration === eventOpenGeneration
|
||||
if (opensSubscription) {
|
||||
openingFrameGeneration = 0
|
||||
}
|
||||
// A resumed subscription's first batch replaces state the client held while detached, so it is
|
||||
// applied alone, never merged into a later batch that would read as "unchanged".
|
||||
if (opensSubscription && event.type === 'batch') {
|
||||
coalescer.flush()
|
||||
if (isCurrentOpenGeneration(eventOpenGeneration)) {
|
||||
args.applyEvent(event, { opensSubscription: true })
|
||||
}
|
||||
return
|
||||
}
|
||||
shouldStopCoalescedEvent = captureHistoryReadGuard()
|
||||
coalescer.push(event)
|
||||
}
|
||||
@@ -200,6 +215,7 @@ export function startStructuredAgentSessionReadTransport(args: {
|
||||
return
|
||||
}
|
||||
const currentOpenGeneration = ++openGeneration
|
||||
openingFrameGeneration = currentOpenGeneration
|
||||
args.onHistoryReadInvalidated()
|
||||
unsubscribe()
|
||||
unsubscribe = (): void => {}
|
||||
|
||||
@@ -455,6 +455,38 @@ describe('structured agent session reducer', () => {
|
||||
expect(withoutCapability.backgroundTasks).toBeUndefined()
|
||||
})
|
||||
|
||||
// An older host omits the roster only when its records hold no running child; a pane resuming
|
||||
// from its cursor must not keep the one it held while away (#24227).
|
||||
it("drops the held roster on a resumed subscription's first batch that omits it", () => {
|
||||
const monitoring = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
|
||||
type: 'event',
|
||||
event: {
|
||||
type: 'snapshot',
|
||||
sessionId: 'session-a',
|
||||
fence: 1,
|
||||
page: hydrationPage([item('message', 1)]),
|
||||
backgroundTasks: { state: 'monitoring' }
|
||||
}
|
||||
})
|
||||
const batch = {
|
||||
type: 'batch' as const,
|
||||
sessionId: 'session-a',
|
||||
batch: { cursor: monitoring.cursor!, items: [], removedItemIds: [], submissions: [] },
|
||||
fence: 1
|
||||
}
|
||||
|
||||
expect(reduceStructuredAgentSession(monitoring, { type: 'event', event: batch })).toBe(
|
||||
monitoring
|
||||
)
|
||||
const resumed = reduceStructuredAgentSession(monitoring, {
|
||||
type: 'event',
|
||||
event: batch,
|
||||
opensSubscription: true
|
||||
})
|
||||
expect(resumed.backgroundTasks).toBeUndefined()
|
||||
expect(resumed.items).toBe(monitoring.items)
|
||||
})
|
||||
|
||||
it('projects ephemeral activity without changing transcript identity and clears it', () => {
|
||||
const initial = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
|
||||
type: 'event',
|
||||
|
||||
@@ -75,7 +75,9 @@ export type StructuredAgentSessionState = {
|
||||
export type StructuredAgentSessionAction =
|
||||
| { type: 'loading' }
|
||||
| { type: 'error'; message: string; refusal?: AgentSessionRefusalReference }
|
||||
| { type: 'event'; event: AgentSessionSubscribeEvent }
|
||||
/** `opensSubscription`: the first frame of a new subscription, which states the roster even
|
||||
* when it resumes from a cursor; a host that omits it there has none to report. */
|
||||
| { type: 'event'; event: AgentSessionSubscribeEvent; opensSubscription?: boolean }
|
||||
| { type: 'history-page'; page: AgentSessionHistoryPage }
|
||||
| { type: 'older-page'; requestedCursor: AgentJournalCursor; page: AgentSessionHistoryPage }
|
||||
|
||||
@@ -276,7 +278,9 @@ export function reduceStructuredAgentSession(
|
||||
const backgroundTasks =
|
||||
event.backgroundTasks !== undefined
|
||||
? admitAgentSessionBackgroundTaskState(event.backgroundTasks, state.backgroundTasks)
|
||||
: state.backgroundTasks
|
||||
: action.opensSubscription
|
||||
? undefined
|
||||
: state.backgroundTasks
|
||||
const activity = event.activity !== undefined ? event.activity : state.activity
|
||||
const liveItems = liveItemsWithinWindow(state, event.batch.items)
|
||||
// Every roster revision, the window's or not: a trimmed roster row keeps its sequence.
|
||||
@@ -333,7 +337,7 @@ export function reduceStructuredAgentSession(
|
||||
readRefusal: undefined,
|
||||
commands: event.commands !== undefined ? event.commands : state.commands,
|
||||
...queuePublicationField(event, state),
|
||||
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
|
||||
backgroundTasks,
|
||||
...(activity !== undefined ? { activity } : {}),
|
||||
...(lostTurnRow ? { unloadedTurnRevisions: (state.unloadedTurnRevisions ?? 0) + 1 } : {}),
|
||||
...hostClockField(event.hostNow, receivedAt, state.hostClock)
|
||||
|
||||
Reference in New Issue
Block a user