mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 00:02:29 +00:00
fix Claude structured session blockers
This commit is contained in:
@@ -42,6 +42,7 @@ export type ClaudeJournalTranslatorDeps = {
|
||||
bindPromptItemId?: (journalItemId: string, promptKey: string, questionId?: string) => void
|
||||
coalesceMs?: number
|
||||
schedule?: AgentSessionDeltaCoalescerDeps['schedule']
|
||||
fallbackIdPrefix?: string
|
||||
}
|
||||
|
||||
export type ClaudeJournalTranslator = {
|
||||
@@ -52,11 +53,13 @@ export type ClaudeJournalTranslator = {
|
||||
|
||||
export function createClaudeSessionJournalTranslator(
|
||||
sink: StructuredAgentSessionEventSink | undefined,
|
||||
prompts: ClaudePromptRegistry
|
||||
prompts: ClaudePromptRegistry,
|
||||
fallbackIdPrefix: string
|
||||
): ClaudeJournalTranslator | null {
|
||||
return sink
|
||||
? createClaudeJournalTranslator({
|
||||
sink,
|
||||
fallbackIdPrefix,
|
||||
bindPromptItemId: (itemId, promptKey, questionId) =>
|
||||
prompts.bindJournalItemId(itemId, promptKey, questionId)
|
||||
})
|
||||
@@ -96,7 +99,10 @@ export function createClaudeJournalTranslator(
|
||||
const latestStreamText = new Map<string, string>()
|
||||
const checkpointLengths = new Map<string, number>()
|
||||
let currentTurn: { sessionId: string; turnId: string } | null = null
|
||||
const providerFallback = createClaudeProviderFrameFallback(deps.sink)
|
||||
const providerFallback = createClaudeProviderFrameFallback(
|
||||
deps.sink,
|
||||
deps.fallbackIdPrefix ?? 'acquisition'
|
||||
)
|
||||
|
||||
const publishLifecycle = (sessionId: string, turnId: string, running: boolean): void => {
|
||||
const identity = lifecycleIdentity(sessionId, turnId)
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { CLAUDE_SPAWN_TOKEN_ENV, claudeProcessIdentity } from './claude-structured-owner-identity'
|
||||
|
||||
const IDENTITY = {
|
||||
sessionId: 'session-identity',
|
||||
workspaceId: 'workspace-1',
|
||||
hostId: 'local',
|
||||
agent: 'claude' as const,
|
||||
providerHandle: { kind: 'claude' as const, sessionId: 'session-1', leafUuid: 'leaf-1' }
|
||||
}
|
||||
|
||||
describe('claude structured owner identity', () => {
|
||||
it('exports the spawn token env and records the observed process identity', async () => {
|
||||
expect(CLAUDE_SPAWN_TOKEN_ENV).toBe('ORCA_AGENT_SESSION_SPAWN_TOKEN')
|
||||
await expect(
|
||||
claudeProcessIdentity(
|
||||
{ identity: IDENTITY, spawnToken: 'spawn-a', pid: 4242 },
|
||||
async () => 123
|
||||
)
|
||||
).resolves.toEqual({
|
||||
hostId: 'local',
|
||||
pid: 4242,
|
||||
processStartTimeMs: 123,
|
||||
spawnToken: 'spawn-a'
|
||||
})
|
||||
})
|
||||
|
||||
it('retries a failed start-time read before giving up', async () => {
|
||||
const readStartTime = vi
|
||||
.fn<(pid: number) => Promise<number | null>>()
|
||||
.mockResolvedValueOnce(null)
|
||||
.mockResolvedValueOnce(null)
|
||||
.mockResolvedValueOnce(456)
|
||||
await expect(
|
||||
claudeProcessIdentity({ identity: IDENTITY, spawnToken: 'spawn-a', pid: 4242 }, readStartTime)
|
||||
).resolves.toMatchObject({ processStartTimeMs: 456 })
|
||||
expect(readStartTime).toHaveBeenCalledTimes(3)
|
||||
})
|
||||
})
|
||||
@@ -1,4 +1,7 @@
|
||||
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
|
||||
import type { AgentSessionProviderHandleLink } from '../../shared/agent-session-provider-handle'
|
||||
import type { AgentSessionProcessIdentity } from '../../shared/agent-session-record'
|
||||
import { readProcessStartTimeMs } from '../runtime/agent-session-process-identity-probe'
|
||||
|
||||
export function claudeProviderHandleLink(input: {
|
||||
sessionId: string
|
||||
@@ -19,3 +22,41 @@ export function claudeProviderHandleLink(input: {
|
||||
observedAt: input.observedAt
|
||||
}
|
||||
}
|
||||
|
||||
/** The child echoes its spawn token here so the owner probe can tell a live
|
||||
* child of this reservation from a same-pid stranger. */
|
||||
export const CLAUDE_SPAWN_TOKEN_ENV = 'ORCA_AGENT_SESSION_SPAWN_TOKEN'
|
||||
|
||||
const START_TIME_READ_ATTEMPTS = 3
|
||||
|
||||
export async function claudeProcessIdentity(
|
||||
input: {
|
||||
identity: AgentSessionJournalIdentity
|
||||
spawnToken: string
|
||||
pid: number | undefined
|
||||
},
|
||||
readStartTime: (pid: number) => Promise<number | null> = readProcessStartTimeMs
|
||||
): Promise<AgentSessionProcessIdentity> {
|
||||
if (input.pid === undefined) {
|
||||
throw new Error('claude app-server started without a pid')
|
||||
}
|
||||
let processStartTimeMs: number | null = null
|
||||
for (
|
||||
let attempt = 0;
|
||||
attempt < START_TIME_READ_ATTEMPTS && processStartTimeMs === null;
|
||||
attempt += 1
|
||||
) {
|
||||
processStartTimeMs = await readStartTime(input.pid)
|
||||
}
|
||||
if (processStartTimeMs === null) {
|
||||
// Why: recording null makes every later owner probe indeterminate — a durable latch.
|
||||
// Failing here reaps the child and leaves a retryable refusal instead.
|
||||
throw new Error(`claude app-server start time for pid ${input.pid} could not be read`)
|
||||
}
|
||||
return {
|
||||
hostId: input.identity.hostId,
|
||||
pid: input.pid,
|
||||
processStartTimeMs,
|
||||
spawnToken: input.spawnToken
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,117 @@
|
||||
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 type {
|
||||
AgentJournalItemBody,
|
||||
AgentSessionJournalIdentity
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
import { openAgentSessionJournal } from '../native-chat/agent-session-journal/journal-store'
|
||||
import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
|
||||
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
|
||||
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
|
||||
|
||||
const IDENTITY: AgentSessionJournalIdentity = {
|
||||
sessionId: 'session-1',
|
||||
workspaceId: 'workspace-1',
|
||||
hostId: 'host-1',
|
||||
agent: 'claude',
|
||||
providerHandle: { kind: 'claude', sessionId: 'provider-1', leafUuid: 'leaf-1' }
|
||||
}
|
||||
|
||||
let root = ''
|
||||
|
||||
function message(
|
||||
role: 'assistant' | 'user',
|
||||
uuid: string,
|
||||
content: unknown[]
|
||||
): ClaudeStructuredSessionEvent {
|
||||
return {
|
||||
type: 'message',
|
||||
sessionId: 'orca-session',
|
||||
message: {
|
||||
type: role,
|
||||
uuid,
|
||||
session_id: 'provider-1',
|
||||
message: { role, content }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-claude-provider-fallback-'))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await rm(root, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('Claude provider fallback', () => {
|
||||
it('drops suppressed init frames instead of dereferencing a null translation', () => {
|
||||
const items: { identity: unknown; body: AgentJournalItemBody }[] = []
|
||||
const sink = {
|
||||
appendItem: (identity: unknown, body: AgentJournalItemBody) => {
|
||||
items.push({ identity, body })
|
||||
},
|
||||
appendTombstone: vi.fn(),
|
||||
publish: vi.fn()
|
||||
}
|
||||
const translator = createClaudeJournalTranslator({ sink })
|
||||
const initEvent: ClaudeStructuredSessionEvent = {
|
||||
type: 'message',
|
||||
sessionId: 'orca-session',
|
||||
message: {
|
||||
type: 'system',
|
||||
subtype: 'init',
|
||||
session_id: 'provider-1',
|
||||
uuid: 'init-1'
|
||||
}
|
||||
}
|
||||
|
||||
expect(() => translator.handle(initEvent)).not.toThrow()
|
||||
expect(items).toEqual([])
|
||||
})
|
||||
|
||||
it('keeps provider-fallback rows distinct across acquisitions', async () => {
|
||||
const journal = await openAgentSessionJournal({
|
||||
identity: IDENTITY,
|
||||
journalDir: root,
|
||||
now: () => 1_700_000_000_000,
|
||||
mintEpoch: () => 'epoch-1'
|
||||
})
|
||||
const deferred = createDeferredStructuredAgentSessionEventSink()
|
||||
deferred.bind({
|
||||
journal,
|
||||
fence: 1,
|
||||
publish: vi.fn()
|
||||
})
|
||||
|
||||
const first = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: '1' })
|
||||
const second = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: '2' })
|
||||
|
||||
first.handle(message('assistant', 'assistant-1', [{ type: 'future_event', message: 'first' }]))
|
||||
await deferred.drained()
|
||||
second.handle(
|
||||
message('assistant', 'assistant-2', [{ type: 'future_event', message: 'second' }])
|
||||
)
|
||||
await deferred.drained()
|
||||
|
||||
const fallbackRows = journal
|
||||
.snapshot()
|
||||
.items.filter(
|
||||
(item) =>
|
||||
item.body.kind === 'status' &&
|
||||
item.body.providerFrame?.kind === 'message:assistant:content:future_event'
|
||||
)
|
||||
|
||||
expect(fallbackRows).toHaveLength(2)
|
||||
expect(fallbackRows.map(statusText)).toEqual(['first', 'second'])
|
||||
})
|
||||
})
|
||||
|
||||
function statusText(row: { body: AgentJournalItemBody }): string {
|
||||
if (row.body.kind !== 'status') {
|
||||
throw new Error('expected status row')
|
||||
}
|
||||
return row.body.text
|
||||
}
|
||||
@@ -30,7 +30,10 @@ export function isModeledClaudeContent(value: unknown): boolean {
|
||||
return part.type === 'thinking' && claudeText(part.thinking) !== null
|
||||
}
|
||||
|
||||
export function createClaudeProviderFrameFallback(sink: StructuredAgentSessionEventSink): {
|
||||
export function createClaudeProviderFrameFallback(
|
||||
sink: StructuredAgentSessionEventSink,
|
||||
acquisitionId: string
|
||||
): {
|
||||
append: (kind: string, payload: unknown) => void
|
||||
} {
|
||||
let sequence = 0
|
||||
@@ -38,8 +41,14 @@ export function createClaudeProviderFrameFallback(sink: StructuredAgentSessionEv
|
||||
append: (kind, payload) => {
|
||||
sequence += 1
|
||||
const translated = unhandledProviderFrameJournalItem('claude', kind, payload)
|
||||
if (!translated) {
|
||||
return
|
||||
}
|
||||
sink.appendItem(
|
||||
{ provider: 'orca', clientMessageId: `provider-frame:claude:${sequence}` },
|
||||
{
|
||||
provider: 'orca',
|
||||
clientMessageId: `provider-frame:claude:${acquisitionId}:${sequence}`
|
||||
},
|
||||
translated.body,
|
||||
translated.blobs
|
||||
)
|
||||
|
||||
@@ -70,7 +70,11 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
async acquire(input: StructuredAgentSessionAcquireInput): Promise<AgentSessionAcquisition> {
|
||||
const sessionId = input.identity.sessionId
|
||||
const prompts = new ClaudePromptRegistry()
|
||||
const translator = createClaudeSessionJournalTranslator(input.events, prompts)
|
||||
const translator = createClaudeSessionJournalTranslator(
|
||||
input.events,
|
||||
prompts,
|
||||
String(input.fence)
|
||||
)
|
||||
const { previous, attempt } = this.acquisitions.start(sessionId, prompts)
|
||||
let liveSession: ClaudeSession | null = null
|
||||
let observedLeafUuid: string | null = null
|
||||
@@ -278,11 +282,11 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
readOptions = (input: { sessionId: string; fence: number }) =>
|
||||
readClaudeStructuredSessionOptions(this.session(input.sessionId), this.deps.requestTimeoutMs)
|
||||
|
||||
releaseAcquisition(input: { sessionId: string }): Promise<void> {
|
||||
releaseAcquisition(input: { sessionId: string }): Promise<boolean> {
|
||||
return this.closeSession(input.sessionId)
|
||||
}
|
||||
|
||||
async closeSession(sessionId: string): Promise<void> {
|
||||
async closeSession(sessionId: string): Promise<boolean> {
|
||||
const attempt = this.acquisitions.get(sessionId)
|
||||
if (attempt) {
|
||||
attempt.cancelled = true
|
||||
@@ -290,6 +294,7 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
await attempt.finished
|
||||
}
|
||||
await this.closePublishedSession(sessionId)
|
||||
return true
|
||||
}
|
||||
|
||||
private async closePublishedSession(sessionId: string): Promise<void> {
|
||||
|
||||
Reference in New Issue
Block a user