Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 20 additions & 18 deletions agent/conversation_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,25 @@ def _ra():
return run_agent


def _relay_available_reasoning(agent, assistant_content: str) -> None:
"""Relay completed inline reasoning without truncating the live fallback."""
think_text = re.sub(
r'</?(?:REASONING_SCRATCHPAD|think|reasoning)>', '', assistant_content.strip()
).strip()
first_line = think_text.split('\n')[0][:80] if think_text else ""

if first_line and getattr(agent, '_delegate_depth', 0) > 0:
try:
agent.tool_progress_callback("_thinking", first_line)
except Exception:
pass
elif think_text:
try:
agent.tool_progress_callback("reasoning.available", "_thinking", think_text, None)
except Exception:
pass


def _nous_entitlement_message(capability: str) -> str:
try:
from hermes_cli.nous_account import (
Expand Down Expand Up @@ -4422,24 +4441,7 @@ def _perform_api_call(next_api_kwargs):
# Notify progress callback of model's thinking (used by subagent
# delegation to relay the child's reasoning to the parent display).
if (assistant_message.content and agent.tool_progress_callback):
_think_text = assistant_message.content.strip()
# Strip reasoning XML tags that shouldn't leak to parent display
_think_text = re.sub(
r'</?(?:REASONING_SCRATCHPAD|think|reasoning)>', '', _think_text
).strip()
# For subagents: relay first line to parent display (existing behaviour).
# For all agents with a structured callback: emit reasoning.available event.
first_line = _think_text.split('\n')[0][:80] if _think_text else ""
if first_line and getattr(agent, '_delegate_depth', 0) > 0:
try:
agent.tool_progress_callback("_thinking", first_line)
except Exception:
pass
elif _think_text:
try:
agent.tool_progress_callback("reasoning.available", "_thinking", _think_text[:500], None)
except Exception:
pass
_relay_available_reasoning(agent, assistant_message.content)

# Check for incomplete <REASONING_SCRATCHPAD> (opened but never closed)
# This means the model ran out of output tokens mid-reasoning — retry up to 2 times
Expand Down
22 changes: 22 additions & 0 deletions apps/desktop/src/app/session/hooks/use-message-stream.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
import { describe, expect, it } from 'vitest'

import { reasoningPart, textPart } from '@/lib/chat-messages'

import { applyReasoningAvailable } from './use-message-stream'

describe('applyReasoningAvailable', () => {
it('inserts available reasoning before existing assistant text', () => {
const parts = applyReasoningAvailable([textPart('Final answer.')], 'Late reasoning.')

expect(parts.map(part => part.type)).toEqual(['reasoning', 'text'])
expect(parts[0]).toEqual(reasoningPart('Late reasoning.'))
expect(parts[1]).toEqual(textPart('Final answer.'))
})

it('keeps the current reasoning part when one already exists', () => {
const existing = [reasoningPart('Streaming reasoning.'), textPart('Final answer.')]
const parts = applyReasoningAvailable(existing, 'Late reasoning.')

expect(parts).toBe(existing)
})
})
26 changes: 15 additions & 11 deletions apps/desktop/src/app/session/hooks/use-message-stream/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,20 @@ interface QueuedStreamDeltas {
reasoning: string
}

export function applyReasoningAvailable(parts: ChatMessagePart[], text: string): ChatMessagePart[] {
if (parts.some(part => part.type === 'reasoning')) {
return parts
}

const textIndex = parts.findIndex(part => part.type === 'text')

if (textIndex < 0) {
return [...parts, reasoningPart(text)]
}

return [...parts.slice(0, textIndex), reasoningPart(text), ...parts.slice(textIndex)]
}

export function useMessageStream({
activeSessionIdRef,
hydrateFromStoredSession,
Expand Down Expand Up @@ -267,17 +281,7 @@ export function useMessageStream({

mutateStream(
sessionId,
(parts, message) => {
if (replace && chatMessageText(message).trim()) {
return parts
}

if (replace) {
return [...parts.filter(part => part.type !== 'reasoning'), reasoningPart(delta)]
}

return appendReasoningPart(parts, delta)
},
parts => (replace ? applyReasoningAvailable(parts, delta) : appendReasoningPart(parts, delta)),
() => [reasoningPart(delta)]
)
},
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
import { QueryClient } from '@tanstack/react-query'
import { act, cleanup, render, waitFor } from '@testing-library/react'
import { useEffect, useRef } from 'react'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'

import type { ClientSessionState } from '@/app/types'
import type { ChatMessagePart } from '@/lib/chat-messages'
import { createClientSessionState } from '@/lib/chat-runtime'
import type { RpcEvent } from '@/types/hermes'

import { useMessageStream } from './index'

const SID = 'session-1'
let handleEvent: ((event: RpcEvent) => void) | null = null
let sessionStateByRuntimeId: Map<string, ClientSessionState>
let reasoningSnapshots: ChatMessagePart[]

function Harness() {
const activeSessionIdRef = useRef<string | null>(SID)
const sessionStateByRuntimeIdRef = useRef(sessionStateByRuntimeId)
const queryClientRef = useRef(new QueryClient())

const stream = useMessageStream({
activeSessionIdRef,
hydrateFromStoredSession: vi.fn(async () => undefined),
queryClient: queryClientRef.current,
refreshHermesConfig: vi.fn(async () => undefined),
refreshSessions: vi.fn(async () => undefined),
sessionStateByRuntimeIdRef,
updateSessionState: (sessionId, updater) => {
const current = sessionStateByRuntimeIdRef.current.get(sessionId) ?? createClientSessionState()
const next = updater(current)
sessionStateByRuntimeIdRef.current.set(sessionId, next)

const message = next.messages.find(item => item.id === next.streamId)
const reasoning = message?.parts.find(part => part.type === 'reasoning')

if (reasoning) {
reasoningSnapshots.push(reasoning)
}

return next
}
})

useEffect(() => {
handleEvent = stream.handleGatewayEvent
}, [stream.handleGatewayEvent])

return null
}

async function mountStream() {
render(<Harness />)
await waitFor(() => expect(handleEvent).not.toBeNull())
}

function emit(type: RpcEvent['type'], payload: RpcEvent['payload'] = {}) {
act(() => handleEvent!({ payload, session_id: SID, type }))
}

function streamedParts(): ChatMessagePart[] {
const state = sessionStateByRuntimeId.get(SID)
const message = state?.messages.find(item => item.id === state.streamId)

return message?.parts ?? []
}

describe('useMessageStream reasoning.available fallback', () => {
beforeEach(() => {
handleEvent = null
sessionStateByRuntimeId = new Map()
reasoningSnapshots = []
})

afterEach(() => {
cleanup()
vi.restoreAllMocks()
})

it('preserves streamed reasoning when the late available fallback arrives', async () => {
await mountStream()
const streamedReasoning = 'streamed reasoning '.repeat(40)

emit('reasoning.delta', { text: streamedReasoning })
emit('reasoning.available', { text: streamedReasoning.slice(0, 500) })

const reasoningParts = streamedParts().filter(part => part.type === 'reasoning')

expect(reasoningParts).toHaveLength(1)
expect(reasoningParts[0]).toMatchObject({ text: streamedReasoning })
expect(reasoningSnapshots.at(-1)).toBe(reasoningSnapshots.at(-2))
})

it('inserts late available reasoning before assistant text when no delta streamed', async () => {
await mountStream()

emit('message.delta', { text: 'Final answer.' })
emit('reasoning.available', { text: 'Late reasoning.' })

expect(streamedParts()).toMatchObject([
{ text: 'Late reasoning.', type: 'reasoning' },
{ text: 'Final answer.', type: 'text' }
])
})
})
32 changes: 32 additions & 0 deletions tests/agent/test_reasoning_available.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
from types import SimpleNamespace

from agent.conversation_loop import _relay_available_reasoning


def test_available_reasoning_relay_preserves_full_content():
events = []
reasoning = "reasoning " * 80
agent = SimpleNamespace(
_delegate_depth=0,
tool_progress_callback=lambda *args: events.append(args),
)

_relay_available_reasoning(
agent,
f"<REASONING_SCRATCHPAD>{reasoning}</REASONING_SCRATCHPAD>",
)

assert events == [("reasoning.available", "_thinking", reasoning.strip(), None)]
assert len(events[0][2]) > 500


def test_available_reasoning_relay_keeps_subagent_preview_bounded():
events = []
agent = SimpleNamespace(
_delegate_depth=1,
tool_progress_callback=lambda *args: events.append(args),
)

_relay_available_reasoning(agent, "first line " * 20 + "\nsecond line")

assert events == [("_thinking", ("first line " * 20)[:80])]