mirror of
https://github.com/dtzp555-max/ocp.git
synced 2026-07-22 05:25:08 +00:00
fix(server): assemble every assistant message so agentic turns return the final answer (#183)
* fix(server): assemble every assistant message so agentic turns return the final answer
`/v1/chat/completions` returned only the agent's opening preamble ("I'll find the
repo…") and silently dropped the post-tool-use final answer on every tool-using
(agentic) turn. The work still ran; OpenAI-compat clients (OpenClaw via ocp-connect,
OpenAI-SDK scripts) just received the preamble — the turn looked like it "did nothing."
Root cause: `parseStreamJsonEvent` extracted aggregate `assistant` text only when
`isFirstDelta` was true, and `isFirstDelta` flips false after the first text. Run
without `--include-partial-messages` the claude CLI emits NO content_block_delta
events — each assistant message arrives as its own aggregate `assistant` event — and
an agentic turn emits SEVERAL (preamble → one per tool round → final answer). So only
the first message's text survived; every later message, including the final answer,
was discarded.
Fix: guard on `sawTextDelta` (set only by a real content_block_delta) instead of
`isFirstDelta`. In aggregate mode (no deltas) accumulate the text of EVERY `assistant`
event, joined with a blank line; the delta+aggregate double-count case is still deduped
(a delta was seen ⇒ ignore the aggregate). Applied to both the buffered (`-p`) and
SSE-streaming paths, plus the mirrored parser + tests in test-features.mjs. 430 passed,
0 failed; added a multi-message agentic regression test.
Evidence — verified live, claude CLI 2.1.206 (`-p --output-format stream-json --verbose
--allowedTools Bash`): a preamble+tool+answer turn emits four `assistant` events, two
carrying text — #2 "Starting the task now." (preamble) and #4 "It printed
AGG_EVIDENCE_99." (final answer, after the Bash tool). Old code returned only #2; new
code returns both.
Endpoint class: B.1 (OpenAI-compatibility surface, `/v1/chat/completions`).
Specification: OpenAI chat/completions — the assistant `message.content` (and the
concatenation of streamed `choices[].delta.content`) carries the model's full response
text (https://platform.openai.com/docs/api-reference/chat/create).
Authorizing ADR: ADR 0006 — OpenAI shim scope. cli.js does NOT perform this operation
(it speaks Anthropic's protocol, not OpenAI's); scope is justified under ADR 0006 Class
B.1. No new endpoint, header, request field, or response field — this only fixes how
OCP assembles claude CLI stream-json output into the existing `content` field. No Class A
(cli.js-mirror) surface is touched.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Merge origin/main into #183 + align streaming-path separator to the buffered guard (reviewer LOW-cosmetic parity)
---------
Co-authored-by: vvlasy-openclaw <vvlasy@gmail.com>
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-authored-by: dtzp555 <dtzp555@gmail.com>
This commit is contained in:
+40
-17
@@ -251,8 +251,8 @@ function parseStreamJsonLines(buffered) {
|
||||
// Reference: OLP lib/providers/anthropic.mjs anthropicStreamJsonEventToIR (commit 97e7d16).
|
||||
//
|
||||
// @param {object} event — parsed NDJSON event
|
||||
// @param {boolean} isFirstDelta — true if no content has been yielded yet
|
||||
function parseStreamJsonEvent(event, isFirstDelta) {
|
||||
// @param {boolean} sawTextDelta — true if a streaming content_block_delta text was already seen
|
||||
function parseStreamJsonEvent(event, sawTextDelta) {
|
||||
const t = event?.type;
|
||||
|
||||
// system/* — first-event init + other system meta (api_retry etc.)
|
||||
@@ -264,19 +264,23 @@ function parseStreamJsonEvent(event, isFirstDelta) {
|
||||
if (t === "stream_event") {
|
||||
const inner = event.event ?? event;
|
||||
if (inner?.type === "content_block_delta" && inner.delta?.type === "text_delta") {
|
||||
return { text: inner.delta.text ?? "" };
|
||||
return { text: inner.delta.text ?? "", fromDelta: true };
|
||||
}
|
||||
// Other stream_event sub-types (content_block_start, message_delta, etc.) — consumed
|
||||
return null;
|
||||
}
|
||||
|
||||
// assistant — aggregate message (fallback when no prior content_block_delta seen)
|
||||
// Empirically (claude CLI without --include-partial-messages, verified v2.1.104 through v2.1.158): fast/short
|
||||
// responses may emit ONLY the aggregate assistant event, no content_block_delta events.
|
||||
// If isFirstDelta is true, extract text here; otherwise it's a duplicate, ignore.
|
||||
// assistant — aggregate message. claude CLI without --include-partial-messages emits NO
|
||||
// content_block_delta events; each assistant message arrives as its own aggregate `assistant`
|
||||
// event. An agentic/tool-using turn has SEVERAL (preamble + one per tool round + final answer),
|
||||
// so we must accumulate the text of EVERY such event. The prior `isFirstDelta` guard kept only
|
||||
// the FIRST message's text and dropped the rest — silently losing the post-tool-use final answer
|
||||
// on every tool-using turn (verified v2.1.104 through v2.1.211; see PR body capture).
|
||||
// The only real hazard is the delta+aggregate DOUBLE-COUNT: if streaming deltas were already
|
||||
// seen (sawTextDelta), the aggregate duplicates them — ignore it.
|
||||
// Reference: OLP commit 65f945c (assistant-aggregate fallback, fold-in).
|
||||
if (t === "assistant") {
|
||||
if (isFirstDelta) {
|
||||
if (!sawTextDelta) {
|
||||
const blocks = event.message?.content;
|
||||
if (Array.isArray(blocks)) {
|
||||
const text = blocks
|
||||
@@ -1501,7 +1505,7 @@ async function callClaude(model, messages, conversationId, keyName, res) {
|
||||
const { proc, cliModel, conversationId: convId, t0, cleanup, handleSessionFailure, markFirstByte } = ctx;
|
||||
let lineBuffer = "";
|
||||
let assembledText = "";
|
||||
let isFirstDelta = true;
|
||||
let sawTextDelta = false;
|
||||
let resultEventSeen = false;
|
||||
let stderr = "";
|
||||
|
||||
@@ -1511,11 +1515,18 @@ async function callClaude(model, messages, conversationId, keyName, res) {
|
||||
const { events, remainder } = parseStreamJsonLines(lineBuffer);
|
||||
lineBuffer = remainder;
|
||||
for (const event of events) {
|
||||
const parsed = parseStreamJsonEvent(event, isFirstDelta);
|
||||
const parsed = parseStreamJsonEvent(event, sawTextDelta);
|
||||
if (!parsed) continue;
|
||||
if (parsed.text !== undefined) {
|
||||
assembledText += parsed.text;
|
||||
isFirstDelta = false;
|
||||
if (parsed.fromDelta) {
|
||||
assembledText += parsed.text;
|
||||
sawTextDelta = true;
|
||||
} else {
|
||||
// aggregate assistant message — separate successive messages so the preamble and
|
||||
// the post-tool-use final answer don't run together.
|
||||
if (assembledText && !assembledText.endsWith("\n")) assembledText += "\n\n";
|
||||
assembledText += parsed.text;
|
||||
}
|
||||
} else if (parsed.stop) {
|
||||
resultEventSeen = true;
|
||||
} else if (parsed.error) {
|
||||
@@ -1937,9 +1948,10 @@ async function callClaudeStreaming(model, messages, conversationId, res, authInf
|
||||
let stderr = "";
|
||||
let headersSent = false;
|
||||
let totalChars = 0;
|
||||
let streamEndsWithNewline = false; // tracks whether emitted text ends in "\n" — see the separator guard below
|
||||
let cachedContent = ""; // accumulate for cache write-back
|
||||
let lineBuffer = "";
|
||||
let isFirstDelta = true;
|
||||
let sawTextDelta = false;
|
||||
let resultEventSeen = false;
|
||||
// Separate flag for is_error result — must NOT be conflated with resultEventSeen.
|
||||
// If errored===true the close handler must not cache the response or record success
|
||||
@@ -1976,15 +1988,26 @@ async function callClaudeStreaming(model, messages, conversationId, res, authInf
|
||||
lineBuffer = remainder;
|
||||
|
||||
for (const event of events) {
|
||||
const parsed = parseStreamJsonEvent(event, isFirstDelta);
|
||||
const parsed = parseStreamJsonEvent(event, sawTextDelta);
|
||||
if (!parsed) continue;
|
||||
|
||||
if (parsed.text !== undefined) {
|
||||
// content_block_delta text — forward as SSE delta
|
||||
const text = parsed.text;
|
||||
// Streamed delta, or an aggregate assistant-message text (agentic turns emit several).
|
||||
// For an aggregate message after earlier text, prepend a separator so the preamble and
|
||||
// the post-tool-use final answer don't run together in the forwarded stream.
|
||||
let text = parsed.text;
|
||||
if (parsed.fromDelta) {
|
||||
sawTextDelta = true;
|
||||
} else if (totalChars > 0 && !streamEndsWithNewline) {
|
||||
// Mirror the buffered path's guard (assembledText.endsWith("\n")): only inject the
|
||||
// blank-line separator when the already-emitted text doesn't already end in a newline,
|
||||
// so a message ending in "\n" doesn't produce a triple newline here while the buffered
|
||||
// path produces a single. Keeps the two assembly paths byte-identical. (PR #183 review.)
|
||||
text = "\n\n" + text;
|
||||
}
|
||||
streamEndsWithNewline = text.endsWith("\n");
|
||||
totalChars += text.length;
|
||||
if (CACHE_TTL > 0) cachedContent += text;
|
||||
isFirstDelta = false;
|
||||
|
||||
if (!ensureHeaders()) continue;
|
||||
sendSSE(res, {
|
||||
|
||||
Reference in New Issue
Block a user