mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
Fix structured host session activity lifecycle
This commit is contained in:
@@ -38,4 +38,5 @@ export type StructuredAgentSessionAttachContext = {
|
||||
reconcileLeases: (sessionId: string) => Promise<AgentSessionWireRefusal | null>
|
||||
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
|
||||
now: () => number
|
||||
publishStatus?: (sessionId: string) => void
|
||||
}
|
||||
|
||||
@@ -22,6 +22,7 @@ export class StructuredAgentSessionEventRecovery {
|
||||
sessions: Map<string, StructuredAgentSessionHostSession>
|
||||
flushLifecycle: (sessionId: string) => Promise<StructuredAgentSessionSinkBarrier>
|
||||
publishFence: (sessionId: string, session: StructuredAgentSessionHostSession) => void
|
||||
publishStatus?: (sessionId: string) => void
|
||||
hasResumeCapableHolder: (sessionId: string) => boolean
|
||||
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
|
||||
now: () => number
|
||||
|
||||
@@ -25,6 +25,7 @@ type HostHandoffAccess = {
|
||||
flush: (sessionId: string) => Promise<void>
|
||||
serialize: (sessionId: string, task: () => Promise<void>) => Promise<void>
|
||||
subscribers: AgentSessionSubscribers
|
||||
publishStatus?: (sessionId: string) => void
|
||||
now: () => number
|
||||
}
|
||||
|
||||
@@ -82,6 +83,7 @@ export function createStructuredAgentSessionHostHandoff(
|
||||
return { state: 'live' }
|
||||
}
|
||||
host.session(sessionId).hasProviderChild = false
|
||||
host.publishStatus?.(sessionId)
|
||||
try {
|
||||
await host.flush(sessionId)
|
||||
host.eventSink(sessionId).unbind()
|
||||
@@ -243,6 +245,7 @@ export async function acquireNativeHandoffOwner(
|
||||
return rethrowAfterAgentSessionAcquisitionCleanup(deps.adapter, input.sessionId, error)
|
||||
}
|
||||
session.hasProviderChild = true
|
||||
host.publishStatus?.(input.sessionId)
|
||||
session.fence = proved.lease.runtimeFence
|
||||
session.acquisitionGeneration = acquired.acquisitionGeneration ?? null
|
||||
eventSink.bind({
|
||||
|
||||
@@ -117,6 +117,7 @@ export class StructuredAgentSessionHost {
|
||||
flush: (sessionId) => this.flushStreamedEvents(sessionId),
|
||||
serialize: (sessionId, task) => this.serialize(sessionId, task),
|
||||
subscribers: this.subscribers,
|
||||
publishStatus: (sessionId) => this.statusFeed.publish(sessionId),
|
||||
now: this.now
|
||||
})
|
||||
this.holds = createStructuredAgentSessionHolds(this.lifetimeContext(), {
|
||||
@@ -144,6 +145,7 @@ export class StructuredAgentSessionHost {
|
||||
flushLifecycle: (sessionId) => this.runtimeState.lifecycleBarrier(sessionId),
|
||||
publishFence: (sessionId, session) =>
|
||||
this.subscribers.snapshot(sessionId, session.journal, session.fence),
|
||||
publishStatus: (sessionId) => this.statusFeed.publish(sessionId),
|
||||
hasResumeCapableHolder: (sessionId) => this.holds.hasResumeCapableHolder(sessionId),
|
||||
serialize: (sessionId, task) => this.serialize(sessionId, task),
|
||||
now: () => this.now(),
|
||||
@@ -188,7 +190,8 @@ export class StructuredAgentSessionHost {
|
||||
subscribers: this.subscribers,
|
||||
tasks: this.tasks,
|
||||
reconcileLeases: (sessionId) => this.reconcileLeases(sessionId),
|
||||
serialize: (sessionId, task) => this.serialize(sessionId, task)
|
||||
serialize: (sessionId, task) => this.serialize(sessionId, task),
|
||||
publishStatus: (sessionId) => this.statusFeed.publish(sessionId)
|
||||
}
|
||||
}
|
||||
/** Releases a session's resources without ending the conversation: the record and journal stay
|
||||
@@ -197,6 +200,7 @@ export class StructuredAgentSessionHost {
|
||||
return this.serialize(sessionId, async () => {
|
||||
await this.handoffs.closeRetainedTuiOwner(sessionId)
|
||||
await evictHeldStructuredAgentSession(this.lifetimeContext(), sessionId)
|
||||
this.statusFeed.revokeLive(sessionId)
|
||||
// Whoever asked for the close, the surfaces that were holding this session are looking at a
|
||||
// session that no longer exists. A failed eviction throws above and keeps them.
|
||||
this.holds.forget(sessionId)
|
||||
|
||||
+46
-2
@@ -53,15 +53,24 @@ async function openJournal(sessionId = SESSION, now?: () => number) {
|
||||
})
|
||||
}
|
||||
|
||||
function indexed(session: { journal: Awaited<ReturnType<typeof openJournal>> }) {
|
||||
function indexed(session: {
|
||||
journal: Awaited<ReturnType<typeof openJournal>>
|
||||
hasProviderChild?: boolean
|
||||
}) {
|
||||
return {
|
||||
journal: session.journal,
|
||||
...(session.hasProviderChild !== undefined
|
||||
? { hasProviderChild: session.hasProviderChild }
|
||||
: {}),
|
||||
params: { location: { workspaceId: 'workspace-1' }, provider: 'codex' as const }
|
||||
}
|
||||
}
|
||||
|
||||
function feedFor(
|
||||
sessions: Map<string, { journal: Awaited<ReturnType<typeof openJournal>> }>,
|
||||
sessions: Map<
|
||||
string,
|
||||
{ journal: Awaited<ReturnType<typeof openJournal>>; hasProviderChild?: boolean }
|
||||
>,
|
||||
record: Partial<AgentSessionRecord> | null = null,
|
||||
onStatusChanged?: StructuredAgentSessionStatusFeedDeps['onStatusChanged']
|
||||
) {
|
||||
@@ -88,6 +97,41 @@ function feedFor(
|
||||
}
|
||||
|
||||
describe('StructuredAgentSessionStatusFeed', () => {
|
||||
it('publishes provider ownership transitions without changing journal time', async () => {
|
||||
const journal = await openJournal()
|
||||
const sessions = new Map([[SESSION, { journal, hasProviderChild: true }]])
|
||||
const { feed, events, dispose } = feedFor(sessions)
|
||||
events.length = 0
|
||||
await journal.appendItem(
|
||||
USER_IDENTITY,
|
||||
{ kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'hello' }] },
|
||||
{ fence: 1 }
|
||||
)
|
||||
feed.publish(SESSION, journal)
|
||||
expect(events.at(-1)).toEqual({
|
||||
type: 'status',
|
||||
session: expect.objectContaining({ hostExecutionOwned: true, updatedAt: expect.any(Number) })
|
||||
})
|
||||
const firstStatus = events.at(-1)
|
||||
expect(firstStatus?.type).toBe('status')
|
||||
if (firstStatus?.type !== 'status') {
|
||||
throw new Error('status publication missing')
|
||||
}
|
||||
const journalTime = firstStatus.session.updatedAt
|
||||
sessions.get(SESSION)!.hasProviderChild = false
|
||||
feed.publish(SESSION, journal)
|
||||
expect(events.at(-1)).toEqual({
|
||||
type: 'status',
|
||||
session: expect.objectContaining({ status: 'idle', updatedAt: journalTime })
|
||||
})
|
||||
const secondStatus = events.at(-1)
|
||||
expect(secondStatus?.type).toBe('status')
|
||||
if (secondStatus?.type === 'status') {
|
||||
expect(secondStatus.session).not.toHaveProperty('hostExecutionOwned')
|
||||
}
|
||||
dispose()
|
||||
})
|
||||
|
||||
it('opens with every readable session and reports no status before a persisted turn', async () => {
|
||||
const journal = await openJournal()
|
||||
const { events } = feedFor(new Map([[SESSION, { journal }]]))
|
||||
|
||||
@@ -30,6 +30,7 @@ export type StructuredAgentSessionStatusSubscriber = {
|
||||
type StatusFeedSession = {
|
||||
journal: AgentSessionJournal
|
||||
params: { location: { workspaceId: string }; provider: AgentSessionRecord['provider'] }
|
||||
hasProviderChild?: boolean
|
||||
}
|
||||
|
||||
export type StructuredAgentSessionStatusFeedDeps = {
|
||||
@@ -46,6 +47,7 @@ function summariesEqual(a: AgentSessionStatusSummary, b: AgentSessionStatusSumma
|
||||
a.workspaceId === b.workspaceId &&
|
||||
a.agent === b.agent &&
|
||||
a.status === b.status &&
|
||||
a.hostExecutionOwned === b.hostExecutionOwned &&
|
||||
a.rewindBlockedReason === b.rewindBlockedReason &&
|
||||
// Settled activity changes ranking; streaming active turns must stay quiet.
|
||||
(a.status !== 'idle' || a.updatedAt === b.updatedAt) &&
|
||||
@@ -110,6 +112,20 @@ export class StructuredAgentSessionStatusFeed {
|
||||
}
|
||||
}
|
||||
|
||||
/** Revoke live execution authority while retaining the last projection for reload history. */
|
||||
revokeLive(sessionId: string): void {
|
||||
const previous = this.published.get(sessionId)
|
||||
if (!previous) {
|
||||
return
|
||||
}
|
||||
const { hostExecutionOwned: _hostExecutionOwned, ...retained } = previous
|
||||
this.published.set(sessionId, retained)
|
||||
this.broadcast({
|
||||
type: 'status',
|
||||
session: retained
|
||||
})
|
||||
}
|
||||
|
||||
/** Re-projects one session after its journal changed; equal projections are not re-sent. */
|
||||
publish(sessionId: string, journal?: AgentSessionJournal, options?: { replay?: boolean }): void {
|
||||
const session = this.deps.sessions.get(sessionId)
|
||||
@@ -147,6 +163,7 @@ export class StructuredAgentSessionStatusFeed {
|
||||
sessionId,
|
||||
workspaceId: session.params.location.workspaceId,
|
||||
agent: session.params.provider,
|
||||
...(session.hasProviderChild ? { hostExecutionOwned: true as const } : {}),
|
||||
...projectStructuredAgentSessionStatusSummary(items),
|
||||
...(record?.rewind?.phase === 'prepared' || record?.rewind?.phase === 'provider-succeeded'
|
||||
? { rewindBlockedReason: 'outcome-unknown' as const }
|
||||
|
||||
@@ -32,6 +32,7 @@ export type StructuredAgentSessionUnexpectedExitContext = {
|
||||
sessions: Map<string, StructuredAgentSessionHostSession>
|
||||
flushLifecycle: (sessionId: string) => Promise<StructuredAgentSessionSinkBarrier>
|
||||
publishFence: (sessionId: string, session: StructuredAgentSessionHostSession) => void
|
||||
publishStatus?: (sessionId: string) => void
|
||||
hasResumeCapableHolder: (sessionId: string) => boolean
|
||||
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
|
||||
now: () => number
|
||||
@@ -59,6 +60,7 @@ export async function settleUnexpectedStructuredAgentSessionExit(
|
||||
if (!record || record.lease.handoffStage !== null) {
|
||||
// The handoff coordinator owns an already-started transition.
|
||||
session.hasProviderChild = false
|
||||
context.publishStatus?.(unexpectedEvent.sessionId)
|
||||
return null
|
||||
}
|
||||
|
||||
@@ -117,6 +119,7 @@ export async function settleUnexpectedStructuredAgentSessionExit(
|
||||
context.onBarrierError?.(unexpectedEvent.sessionId, error)
|
||||
} finally {
|
||||
session.hasProviderChild = false
|
||||
context.publishStatus?.(unexpectedEvent.sessionId)
|
||||
if (released) {
|
||||
session.fence = released.lease.runtimeFence
|
||||
context.publishFence(unexpectedEvent.sessionId, session)
|
||||
|
||||
@@ -21,6 +21,7 @@ export async function replaceClaudeRewindOwner(
|
||||
return rewindRefusal('outcome-unknown')
|
||||
}
|
||||
session.hasProviderChild = false
|
||||
context.publishStatus?.(sessionId)
|
||||
const head = agentSessionProviderHandleChainHead(
|
||||
context.deps.store.getRecord(sessionId)!.providerHandleChain
|
||||
)?.handle
|
||||
|
||||
@@ -23,12 +23,18 @@ function summary(over: Partial<AgentSessionStatusSummary> = {}): AgentSessionSta
|
||||
status: 'working',
|
||||
latestPrompt: 'ship the thing',
|
||||
updatedAt: 1_757_030_400_000,
|
||||
hostExecutionOwned: true,
|
||||
...over
|
||||
} as AgentSessionStatusSummary
|
||||
}
|
||||
|
||||
function attach(summaries: AgentSessionStatusSummary[]): RuntimeWorktreePsSummary {
|
||||
const row = { worktreeId: WORKTREE_ID, agents: [] } as unknown as RuntimeWorktreePsSummary
|
||||
const row = {
|
||||
worktreeId: WORKTREE_ID,
|
||||
status: 'inactive',
|
||||
hasHostSidebarActivity: false,
|
||||
agents: []
|
||||
} as unknown as RuntimeWorktreePsSummary
|
||||
const summariesById = new Map<string, RuntimeWorktreePsSummary>([[WORKTREE_ID, row]])
|
||||
attachRuntimeWorktreeAgentRows({
|
||||
summaries: summariesById,
|
||||
@@ -63,6 +69,12 @@ describe('worktree ps reports structured sessions', () => {
|
||||
expect(attach([summary({ status: 'idle' })]).agents[0]?.state).toBe('done')
|
||||
})
|
||||
|
||||
it('does not turn a completed host-held session into permission', () => {
|
||||
const row = attach([summary({ status: 'idle' })])
|
||||
expect(row.status).toBe('inactive')
|
||||
expect(row.hasHostSidebarActivity).toBe(false)
|
||||
})
|
||||
|
||||
it('reports the DERIVED pane key, never an orchestration credential', () => {
|
||||
const row = attach([summary()])
|
||||
const sessionId = 'a1b2c3d4-e5f6-4a7b-8c9d-0e1f2a3b4c5d'
|
||||
|
||||
@@ -81,7 +81,9 @@ export function attachRuntimeWorktreeAgentRows(args: {
|
||||
const monitoringSources: RuntimeWorktreeAgentSource[] = []
|
||||
for (const row of rows) {
|
||||
const source = rowSources.get(row.paneKey)
|
||||
if (source?.authority !== 'structured-host' && !isFreshNonDoneAgentStatus(row, now)) {
|
||||
const hostHeldStructuredSession =
|
||||
source?.authority === 'structured-host' && row.state !== 'done'
|
||||
if (!hostHeldStructuredSession && !isFreshNonDoneAgentStatus(row, now)) {
|
||||
continue
|
||||
}
|
||||
summary.hasHostSidebarActivity = true
|
||||
|
||||
@@ -130,6 +130,7 @@ describe('worktree ps and a closed structured chat', () => {
|
||||
const { feed } = await awaitingApproval()
|
||||
const aged = feed.liveSessionSummaries().map((summary) => ({
|
||||
...summary,
|
||||
hostExecutionOwned: true as const,
|
||||
updatedAt: Date.now() - 30 * 60 * 1000 - 1,
|
||||
status: 'working' as const
|
||||
}))
|
||||
@@ -144,6 +145,7 @@ describe('worktree ps and a closed structured chat', () => {
|
||||
const { feed } = await awaitingApproval()
|
||||
const aged = feed.liveSessionSummaries().map((summary) => ({
|
||||
...summary,
|
||||
hostExecutionOwned: true as const,
|
||||
updatedAt: Date.now() - 30 * 60 * 1000 - 1
|
||||
}))
|
||||
const row = worktreeFor(feed, aged)
|
||||
|
||||
@@ -41,7 +41,7 @@ export function structuredRuntimeWorktreeAgentSources(
|
||||
interrupted: false,
|
||||
stateStartedAt: summary.updatedAt,
|
||||
updatedAt: summary.updatedAt,
|
||||
authority: 'structured-host'
|
||||
...(summary.hostExecutionOwned ? { authority: 'structured-host' as const } : {})
|
||||
})
|
||||
}
|
||||
return sources
|
||||
|
||||
@@ -7,6 +7,7 @@ import type {
|
||||
AgentSessionStatusSummary
|
||||
} from '../../../../shared/agent-session-wire'
|
||||
import { resolveAttention } from '../sidebar/smart-attention'
|
||||
import { isExplicitAgentStatusFresh } from '@/lib/pane-agent-evidence'
|
||||
import type { AgentStatusEntry } from '../../../../shared/agent-status-types'
|
||||
import type { Tab } from '../../../../shared/tab-types'
|
||||
import type { AppState } from '@/store/types'
|
||||
@@ -88,6 +89,7 @@ function summary(overrides: Partial<AgentSessionStatusSummary> = {}): AgentSessi
|
||||
workspaceId: 'wt-1',
|
||||
agent: 'codex',
|
||||
status: 'working',
|
||||
hostExecutionOwned: true,
|
||||
latestPrompt: 'hello',
|
||||
providerSession,
|
||||
updatedAt: 1,
|
||||
@@ -180,6 +182,26 @@ describe('StructuredAgentSessionStatusBridge', () => {
|
||||
// Hiddenness is the host's side of this: see structured-agent-session-subscribers.test.ts,
|
||||
// which drives an unsubscribed journal through the feed. Here the transport is a mock, so
|
||||
// only the summary-to-store mapping is under test.
|
||||
it('keeps host-held working evidence active past the normal freshness window', async () => {
|
||||
render(<StructuredAgentSessionStatusBridge />)
|
||||
await waitFor(() => expect(mocks.subscribeStatus).toHaveBeenCalledOnce())
|
||||
const updatedAt = Date.now() - 30 * 60 * 1000 - 1
|
||||
act(() => feed().emit({ type: 'status', session: summary({ updatedAt }) }))
|
||||
const entry = statuses()[0]
|
||||
expect(entry).toEqual(expect.objectContaining({ state: 'working', structuredHostOwned: true }))
|
||||
expect(isExplicitAgentStatusFresh(entry, Date.now(), 30 * 60 * 1000)).toBe(true)
|
||||
})
|
||||
|
||||
it('clears host-held evidence when the status stream disconnects', async () => {
|
||||
render(<StructuredAgentSessionStatusBridge />)
|
||||
await waitFor(() => expect(mocks.subscribeStatus).toHaveBeenCalledOnce())
|
||||
act(() => feed().emit({ type: 'status', session: summary() }))
|
||||
expect(statuses()).toHaveLength(1)
|
||||
act(() => feed().emit({ type: 'end' }))
|
||||
expect(statuses()).toHaveLength(1)
|
||||
expect(statuses()[0]).not.toHaveProperty('structuredHostOwned')
|
||||
})
|
||||
|
||||
it('maps each host status onto the sidebar agent state', async () => {
|
||||
render(<StructuredAgentSessionStatusBridge />)
|
||||
await waitFor(() => expect(mocks.subscribeStatus).toHaveBeenCalledOnce())
|
||||
|
||||
@@ -98,6 +98,7 @@ function projectStatus(tab: StructuredTab, summary: AgentSessionStatusSummary |
|
||||
current.tabId === tab.id &&
|
||||
current.worktreeId === tab.worktreeId &&
|
||||
current.terminalResumeEligible === false &&
|
||||
current.structuredHostOwned === summary.hostExecutionOwned &&
|
||||
agentProviderSessionsEqual(
|
||||
tab.agentSessionAgent,
|
||||
current.providerSession,
|
||||
@@ -118,12 +119,13 @@ function projectStatus(tab: StructuredTab, summary: AgentSessionStatusSummary |
|
||||
desired.state !== 'done' && current?.state === desired.state
|
||||
? current.stateStartedAt
|
||||
: summary.updatedAt,
|
||||
evidenceObservedAt: Date.now()
|
||||
evidenceObservedAt: summary.updatedAt
|
||||
},
|
||||
{ tabId: tab.id, worktreeId: tab.worktreeId },
|
||||
{
|
||||
...(summary.providerSession ? { providerSession: summary.providerSession } : {}),
|
||||
terminalResumeEligible: false
|
||||
terminalResumeEligible: false,
|
||||
...(summary.hostExecutionOwned ? { structuredHostOwned: true as const } : {})
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
@@ -18,14 +18,20 @@ import {
|
||||
export function isExplicitAgentStatusFresh(
|
||||
entry: Pick<
|
||||
AgentStatusEntry,
|
||||
'updatedAt' | 'evidenceObservedAt' | 'mirroredEvidenceReceivedAt' | 'restoredUnconfirmed'
|
||||
| 'updatedAt'
|
||||
| 'evidenceObservedAt'
|
||||
| 'mirroredEvidenceReceivedAt'
|
||||
| 'restoredUnconfirmed'
|
||||
| 'structuredHostOwned'
|
||||
>,
|
||||
now: number,
|
||||
staleAfterMs: number
|
||||
): boolean {
|
||||
// Why: an unconfirmed hydrated row may describe a turn that ended while no receiver was up; never fresh.
|
||||
return (
|
||||
entry.restoredUnconfirmed !== true && now - agentStatusEvidenceObservedAt(entry) <= staleAfterMs
|
||||
entry.restoredUnconfirmed !== true &&
|
||||
(entry.structuredHostOwned === true ||
|
||||
now - agentStatusEvidenceObservedAt(entry) <= staleAfterMs)
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type {
|
||||
AgentSessionStatusEvent,
|
||||
AgentSessionStatusSummary
|
||||
} from '../../../shared/agent-session-wire'
|
||||
|
||||
const mocks = vi.hoisted(() => ({ subscribe: vi.fn() }))
|
||||
vi.mock('./structured-agent-session-client', () => ({
|
||||
subscribeStructuredAgentSessionStatus: mocks.subscribe
|
||||
}))
|
||||
vi.mock('./runtime-rpc-client', () => ({ runtimeEnvironmentSupportsCapability: vi.fn() }))
|
||||
|
||||
import {
|
||||
getStructuredAgentSessionStatusFeed,
|
||||
resetStructuredAgentSessionStatusFeedsForTests
|
||||
} from './structured-agent-session-status-feed'
|
||||
|
||||
type Subscription = {
|
||||
emit: (event: AgentSessionStatusEvent) => void
|
||||
unsubscribe: ReturnType<typeof vi.fn>
|
||||
}
|
||||
const subscriptions: Subscription[] = []
|
||||
const owned: AgentSessionStatusSummary = {
|
||||
sessionId: 'running',
|
||||
workspaceId: 'workspace',
|
||||
agent: 'codex',
|
||||
status: 'working',
|
||||
latestPrompt: 'work',
|
||||
updatedAt: 1,
|
||||
hostExecutionOwned: true
|
||||
}
|
||||
const done: AgentSessionStatusSummary = {
|
||||
...owned,
|
||||
sessionId: 'completed',
|
||||
status: 'idle',
|
||||
hostExecutionOwned: undefined
|
||||
}
|
||||
|
||||
function subscription(index = 0): Subscription {
|
||||
const value = subscriptions[index]
|
||||
if (!value) {
|
||||
throw new Error('missing subscription')
|
||||
}
|
||||
return value
|
||||
}
|
||||
|
||||
describe('structured status feed execution authority lifecycle', () => {
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
resetStructuredAgentSessionStatusFeedsForTests()
|
||||
subscriptions.length = 0
|
||||
mocks.subscribe.mockReset()
|
||||
mocks.subscribe.mockImplementation((_target, emit: Subscription['emit']) => {
|
||||
const unsubscribe = vi.fn(() => emit({ type: 'end' }))
|
||||
subscriptions.push({ emit, unsubscribe })
|
||||
return Promise.resolve({ unsubscribe })
|
||||
})
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
resetStructuredAgentSessionStatusFeedsForTests()
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('revokes on end before reentrant unsubscribe and ignores late frames', async () => {
|
||||
const feed = getStructuredAgentSessionStatusFeed({ kind: 'local' })
|
||||
feed.activate()
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
subscription().emit({ type: 'snapshot', sessions: [owned, done] })
|
||||
subscription().emit({ type: 'end' })
|
||||
expect(subscription().unsubscribe).toHaveBeenCalledOnce()
|
||||
expect(feed.getSnapshot().get('running')).toEqual({ ...owned, hostExecutionOwned: undefined })
|
||||
expect(feed.getSnapshot().get('completed')).toBe(done)
|
||||
subscription().emit({ type: 'status', session: owned })
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBeUndefined()
|
||||
expect(vi.getTimerCount()).toBe(1)
|
||||
await vi.advanceTimersByTimeAsync(250)
|
||||
subscription(1).emit({ type: 'snapshot', sessions: [done] })
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBeUndefined()
|
||||
subscription(1).emit({ type: 'status', session: owned })
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBe(true)
|
||||
})
|
||||
|
||||
it('retains history without ownership while stopped and until remount receives fresh evidence', async () => {
|
||||
const feed = getStructuredAgentSessionStatusFeed({ kind: 'local' })
|
||||
const deactivate = feed.activate()
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
subscription().emit({ type: 'snapshot', sessions: [owned, done] })
|
||||
deactivate()
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBeUndefined()
|
||||
expect(feed.getSnapshot().get('running')?.updatedAt).toBe(owned.updatedAt)
|
||||
expect(feed.getSnapshot().get('completed')).toBe(done)
|
||||
feed.activate()
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBeUndefined()
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
subscription().emit({ type: 'status', session: owned })
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBeUndefined()
|
||||
subscription(1).emit({ type: 'snapshot', sessions: [owned] })
|
||||
expect(feed.getSnapshot().get('running')?.hostExecutionOwned).toBe(true)
|
||||
})
|
||||
|
||||
it('does not notify or reallocate already unowned historical rows on teardown', async () => {
|
||||
const feed = getStructuredAgentSessionStatusFeed({ kind: 'local' })
|
||||
const deactivate = feed.activate()
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
subscription().emit({ type: 'snapshot', sessions: [done] })
|
||||
const previous = feed.getSnapshot()
|
||||
const listener = vi.fn()
|
||||
feed.subscribe(listener)
|
||||
deactivate()
|
||||
expect(feed.getSnapshot()).toBe(previous)
|
||||
expect(listener).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -82,6 +82,32 @@ function createOwner(target: RuntimeClientTarget): OwnedStatusFeed {
|
||||
handle?.unsubscribe()
|
||||
handle = null
|
||||
}
|
||||
const revokeSnapshotOwnership = (): void => {
|
||||
let next: Map<string, AgentSessionStatusSummary> | null = null
|
||||
for (const [sessionId, summary] of snapshot) {
|
||||
if (!summary.hostExecutionOwned) {
|
||||
continue
|
||||
}
|
||||
if (!next) {
|
||||
next = new Map(snapshot)
|
||||
}
|
||||
const { hostExecutionOwned: _owned, ...retained } = summary
|
||||
next.set(sessionId, retained)
|
||||
}
|
||||
if (next) {
|
||||
snapshot = next
|
||||
emit()
|
||||
}
|
||||
}
|
||||
const fenceCandidateAndReconnect = (candidate: number): void => {
|
||||
if (candidate !== generation) {
|
||||
return
|
||||
}
|
||||
generation += 1
|
||||
revokeSnapshotOwnership()
|
||||
dropHandle()
|
||||
scheduleReconnect(generation)
|
||||
}
|
||||
let open = (): void => {}
|
||||
const scheduleReconnect = (candidate: number): void => {
|
||||
if (!active(candidate) || reconnectTimer) {
|
||||
@@ -104,22 +130,19 @@ function createOwner(target: RuntimeClientTarget): OwnedStatusFeed {
|
||||
return
|
||||
}
|
||||
if (event.type === 'end') {
|
||||
dropHandle()
|
||||
scheduleReconnect(candidate)
|
||||
fenceCandidateAndReconnect(candidate)
|
||||
return
|
||||
}
|
||||
applyEvent(event)
|
||||
},
|
||||
() => {
|
||||
if (active(candidate)) {
|
||||
dropHandle()
|
||||
scheduleReconnect(candidate)
|
||||
fenceCandidateAndReconnect(candidate)
|
||||
}
|
||||
},
|
||||
() => {
|
||||
if (active(candidate)) {
|
||||
dropHandle()
|
||||
scheduleReconnect(candidate)
|
||||
fenceCandidateAndReconnect(candidate)
|
||||
}
|
||||
}
|
||||
)
|
||||
@@ -130,7 +153,13 @@ function createOwner(target: RuntimeClientTarget): OwnedStatusFeed {
|
||||
opened.unsubscribe()
|
||||
}
|
||||
})
|
||||
.catch(() => scheduleReconnect(candidate))
|
||||
.catch(() => {
|
||||
if (active(candidate)) {
|
||||
fenceCandidateAndReconnect(candidate)
|
||||
} else {
|
||||
scheduleReconnect(candidate)
|
||||
}
|
||||
})
|
||||
}
|
||||
open = (): void => {
|
||||
const candidate = ++generation
|
||||
@@ -163,6 +192,7 @@ function createOwner(target: RuntimeClientTarget): OwnedStatusFeed {
|
||||
generation += 1
|
||||
clearReconnect()
|
||||
dropHandle()
|
||||
revokeSnapshotOwnership()
|
||||
reconnectAttempt = 0
|
||||
}
|
||||
|
||||
|
||||
@@ -108,6 +108,8 @@ export type AgentStatusRouting = {
|
||||
}
|
||||
|
||||
export type AgentStatusMetadata = {
|
||||
/** Structured status rows remain fresh while the host owns the session; cleared on feed loss. */
|
||||
structuredHostOwned?: true
|
||||
providerSession?: AgentProviderSessionMetadata
|
||||
launchConfig?: SleepingAgentLaunchConfig
|
||||
launchToken?: string
|
||||
|
||||
@@ -227,6 +227,7 @@ export function buildAgentStatusLiveEntry(
|
||||
...(timing?.evidenceObservedAt !== undefined
|
||||
? { evidenceObservedAt: timing.evidenceObservedAt }
|
||||
: {}),
|
||||
...(metadata?.structuredHostOwned === true ? { structuredHostOwned: true as const } : {}),
|
||||
stateStartedAt,
|
||||
agentType: identity.agentType,
|
||||
model:
|
||||
|
||||
@@ -196,6 +196,8 @@ export type AgentSessionStatusSummary = {
|
||||
agent: AgentSessionRecord['provider']
|
||||
/** Null until the journal holds a persisted user or assistant message. */
|
||||
status: StructuredAgentSessionProjectedStatus | null
|
||||
/** Present only while this host has the provider child executing the session. */
|
||||
hostExecutionOwned?: true
|
||||
latestPrompt: string
|
||||
/** Provider model in force for the next turn; absent until the host has read the options. */
|
||||
model?: string
|
||||
|
||||
@@ -37,6 +37,7 @@ export function isFreshNonDoneAgentStatus(
|
||||
| 'evidenceObservedAt'
|
||||
| 'mirroredEvidenceReceivedAt'
|
||||
| 'restoredUnconfirmed'
|
||||
| 'structuredHostOwned'
|
||||
>
|
||||
| undefined,
|
||||
now = Date.now(),
|
||||
@@ -47,6 +48,7 @@ export function isFreshNonDoneAgentStatus(
|
||||
entry &&
|
||||
entry.state !== 'done' &&
|
||||
entry.restoredUnconfirmed !== true &&
|
||||
now - agentStatusEvidenceObservedAt(entry) <= staleAfterMs
|
||||
(entry.structuredHostOwned === true ||
|
||||
now - agentStatusEvidenceObservedAt(entry) <= staleAfterMs)
|
||||
)
|
||||
}
|
||||
|
||||
@@ -114,6 +114,8 @@ export type AgentStatusEntry = {
|
||||
* which is the delivery/ordering clock a relay reconnect must restamp to stay monotonic.
|
||||
* Absent for locally derived rows and old hosts; freshness falls back to `updatedAt`. */
|
||||
evidenceObservedAt?: number
|
||||
/** True only while a host-held structured session is represented by its live status feed. */
|
||||
structuredHostOwned?: true
|
||||
/** Timestamp (ms) when the current `state` was first reported.
|
||||
* Why: separate from updatedAt so tool/prompt pings (which reset updatedAt) don't move it. */
|
||||
stateStartedAt: number
|
||||
|
||||
Reference in New Issue
Block a user