mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
fix(telegram): prioritize finals over CLI commentary (#134826)
Worked on by: - @VACInc Co-authored-by: roboclaw-bot <309084314+roboclaw-bot@users.noreply.github.com> Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com>
This commit is contained in:
parent
6d8ae28a75
commit
b7fa0f52c6
55 changed files with 1080 additions and 149 deletions
22
.github/workflows/qa-live-transports-convex.yml
vendored
22
.github/workflows/qa-live-transports-convex.yml
vendored
|
|
@ -251,21 +251,11 @@ jobs:
|
|||
if [[ "$selected_revision" == "$release_branch_sha" ]]; then
|
||||
trusted_reason="release-branch-head"
|
||||
fi
|
||||
else
|
||||
pr_head_count="$(
|
||||
gh api \
|
||||
-H "Accept: application/vnd.github+json" \
|
||||
"repos/${GITHUB_REPOSITORY}/commits/${selected_revision}/pulls" \
|
||||
--jq '[.[] | select(.state == "open" and .head.repo.full_name == "'"${GITHUB_REPOSITORY}"'" and .head.sha == "'"${selected_revision}"'")] | length'
|
||||
)"
|
||||
if [[ "$pr_head_count" != "0" ]]; then
|
||||
trusted_reason="open-pr-head"
|
||||
fi
|
||||
fi
|
||||
|
||||
if [[ -z "$trusted_reason" ]]; then
|
||||
echo "Ref '${INPUT_REF}' resolved to $selected_revision, which is not trusted for this secret-bearing QA run." >&2
|
||||
echo "Allowed refs must be on main, point to a release tag, match a release branch head, or match an open PR head in ${GITHUB_REPOSITORY}." >&2
|
||||
echo "Allowed refs must be on main, point to a release tag, match a release branch head, or be an exact revision selected by a trusted caller." >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
|
|
@ -280,7 +270,7 @@ jobs:
|
|||
run_mock_parity:
|
||||
name: Run QA Lab mock parity lane
|
||||
needs: [validate_selected_ref]
|
||||
if: inputs.expected_sha == '' || inputs.run_mock_parity
|
||||
if: (inputs.expected_sha == '' || inputs.run_mock_parity) && (github.event_name != 'workflow_dispatch' || inputs.scenario == '')
|
||||
runs-on: ${{ vars.OPENCLAW_CI_RUNNER_BACKEND == 'github' && 'ubuntu-24.04' || 'blacksmith-16vcpu-ubuntu-2404' }}
|
||||
timeout-minutes: 30
|
||||
env:
|
||||
|
|
@ -464,7 +454,7 @@ jobs:
|
|||
run_live_matrix:
|
||||
name: Run Matrix live QA lane
|
||||
needs: [authorize_actor, validate_selected_ref]
|
||||
if: inputs.expected_sha == '' || inputs.run_matrix
|
||||
if: (inputs.expected_sha == '' || inputs.run_matrix) && (github.event_name != 'workflow_dispatch' || inputs.scenario == '')
|
||||
runs-on: ${{ vars.OPENCLAW_CI_RUNNER_BACKEND == 'github' && 'ubuntu-24.04' || 'blacksmith-16vcpu-ubuntu-2404' }}
|
||||
timeout-minutes: 90
|
||||
concurrency:
|
||||
|
|
@ -768,7 +758,7 @@ jobs:
|
|||
run_live_discord:
|
||||
name: Run Discord live QA lane with Convex leases
|
||||
needs: [authorize_actor, validate_selected_ref]
|
||||
if: inputs.expected_sha == '' || inputs.run_discord
|
||||
if: (inputs.expected_sha == '' || inputs.run_discord) && (github.event_name != 'workflow_dispatch' || inputs.scenario == '')
|
||||
runs-on: ${{ vars.OPENCLAW_CI_RUNNER_BACKEND == 'github' && 'ubuntu-24.04' || 'blacksmith-16vcpu-ubuntu-2404' }}
|
||||
timeout-minutes: 60
|
||||
environment: qa-live-shared
|
||||
|
|
@ -854,7 +844,7 @@ jobs:
|
|||
run_live_whatsapp:
|
||||
name: Run WhatsApp live QA lane with Convex leases
|
||||
needs: [authorize_actor, validate_selected_ref]
|
||||
if: inputs.expected_sha == '' || inputs.run_whatsapp
|
||||
if: (inputs.expected_sha == '' || inputs.run_whatsapp) && (github.event_name != 'workflow_dispatch' || inputs.scenario == '')
|
||||
runs-on: ${{ vars.OPENCLAW_CI_RUNNER_BACKEND == 'github' && 'ubuntu-24.04' || 'blacksmith-16vcpu-ubuntu-2404' }}
|
||||
timeout-minutes: 60
|
||||
concurrency:
|
||||
|
|
@ -933,7 +923,7 @@ jobs:
|
|||
run_live_slack:
|
||||
name: Run Slack live QA lane with Convex leases
|
||||
needs: [authorize_actor, validate_selected_ref]
|
||||
if: inputs.expected_sha == '' || inputs.run_slack
|
||||
if: (inputs.expected_sha == '' || inputs.run_slack) && (github.event_name != 'workflow_dispatch' || inputs.scenario == '')
|
||||
runs-on: ${{ vars.OPENCLAW_CI_RUNNER_BACKEND == 'github' && 'ubuntu-24.04' || 'blacksmith-16vcpu-ubuntu-2404' }}
|
||||
timeout-minutes: 60
|
||||
environment: qa-live-shared
|
||||
|
|
|
|||
|
|
@ -334,10 +334,26 @@ may include final text, usage, an error, and a successor session id. Session ids
|
|||
reported by either event shape participate in resumed-session and fork
|
||||
persistence.
|
||||
|
||||
Lifecycle events are intentionally separate from this return union so existing
|
||||
plugins can continue to match it exhaustively. Use `parseJsonlLifecycleEvent`
|
||||
for backend-owned lifecycle records instead.
|
||||
|
||||
Tool events describe work the backend already performed. OpenClaw renders and
|
||||
summarizes them, but does not treat them as host tool execution, trusted
|
||||
diagnostics, loopback correlation, or message-delivery evidence.
|
||||
|
||||
### `parseJsonlLifecycleEvent`: provider-native lifecycle records
|
||||
|
||||
Set `parseJsonlLifecycleEvent` when a backend emits JSONL records for lifecycle
|
||||
state that is independent of assistant text, tools, sessions, and terminal
|
||||
results. The hook receives the same line and context as `parseJsonlEvent` and is
|
||||
tried first. Returning a lifecycle event consumes that line; returning `null`
|
||||
lets the source-compatible `parseJsonlEvent` hook or built-in parser handle it.
|
||||
|
||||
The current lifecycle contract supports native compaction start and end records.
|
||||
An end record includes `completed` so channels can distinguish successful and
|
||||
incomplete compaction without inferring an outcome from later messages.
|
||||
|
||||
### `ownsNativeCompaction`: opting out of OpenClaw compaction
|
||||
|
||||
If your backend runs an agent that compacts its **own** transcript, set
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import type {
|
|||
CliBackendPlugin,
|
||||
CliBackendPreparedExecution,
|
||||
} from "openclaw/plugin-sdk/cli-backend";
|
||||
import { parseClaudeCliJsonlEvent } from "./cli-output.js";
|
||||
import { parseClaudeCliJsonlEvent, parseClaudeCliJsonlLifecycleEvent } from "./cli-output.js";
|
||||
import {
|
||||
CLAUDE_CLI_BACKEND_ID,
|
||||
CLAUDE_CLI_DEFAULT_MODEL_REF,
|
||||
|
|
@ -276,6 +276,7 @@ export function buildAnthropicCliBackend(
|
|||
return supportProbe ? supportProbe.then(prepare) : prepare();
|
||||
},
|
||||
parseJsonlEvent: parseClaudeCliJsonlEvent,
|
||||
parseJsonlLifecycleEvent: parseClaudeCliJsonlLifecycleEvent,
|
||||
resolveExecutionArgs: (context) =>
|
||||
resolveClaudeCliExecutionArgs(context, {
|
||||
excludeDynamicSystemPromptSections: options.supportsDynamicSystemPromptSections?.(),
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import type { SDKStatusMessage } from "@anthropic-ai/claude-agent-sdk";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { buildAnthropicCliBackend } from "./cli-backend.js";
|
||||
|
||||
|
|
@ -25,6 +26,62 @@ function parseResult(result: string) {
|
|||
}
|
||||
|
||||
describe("Claude CLI output validation", () => {
|
||||
it("projects Claude SDK compaction status lifecycle without inferring from the boundary", () => {
|
||||
const backend = buildAnthropicCliBackend();
|
||||
const parseLifecycle = (event: unknown) =>
|
||||
backend.parseJsonlLifecycleEvent?.(JSON.stringify(event), {
|
||||
backendId: backend.id,
|
||||
backend: backend.config,
|
||||
});
|
||||
|
||||
const sdkCompactionEvents = [
|
||||
{
|
||||
type: "system",
|
||||
subtype: "status",
|
||||
status: "compacting",
|
||||
uuid: "00000000-0000-4000-8000-000000000001",
|
||||
session_id: "00000000-0000-4000-8000-000000000002",
|
||||
},
|
||||
{
|
||||
type: "system",
|
||||
subtype: "status",
|
||||
status: null,
|
||||
compact_result: "success",
|
||||
uuid: "00000000-0000-4000-8000-000000000003",
|
||||
session_id: "00000000-0000-4000-8000-000000000002",
|
||||
},
|
||||
{
|
||||
type: "system",
|
||||
subtype: "status",
|
||||
status: null,
|
||||
compact_result: "failed",
|
||||
uuid: "00000000-0000-4000-8000-000000000004",
|
||||
session_id: "00000000-0000-4000-8000-000000000002",
|
||||
},
|
||||
] satisfies SDKStatusMessage[];
|
||||
|
||||
expect(parseLifecycle(sdkCompactionEvents[0])).toEqual({
|
||||
kind: "compaction",
|
||||
phase: "start",
|
||||
});
|
||||
expect(parseLifecycle(sdkCompactionEvents[1])).toEqual({
|
||||
kind: "compaction",
|
||||
phase: "end",
|
||||
completed: true,
|
||||
});
|
||||
expect(parseLifecycle(sdkCompactionEvents[2])).toEqual({
|
||||
kind: "compaction",
|
||||
phase: "end",
|
||||
completed: false,
|
||||
});
|
||||
expect(parseLifecycle({ compact_result: "success" })).toEqual({
|
||||
kind: "compaction",
|
||||
phase: "end",
|
||||
completed: true,
|
||||
});
|
||||
expect(parseLifecycle({ type: "system", subtype: "compact_boundary" })).toBeNull();
|
||||
});
|
||||
|
||||
it("rejects mocked raw tool protocol returned as terminal assistant text", () => {
|
||||
expect(parseResult(MOCK_RAW_TOOL_OUTPUT)).toEqual({
|
||||
kind: "result",
|
||||
|
|
|
|||
|
|
@ -1,4 +1,7 @@
|
|||
import type { CliBackendParseJsonlEvent } from "openclaw/plugin-sdk/cli-backend";
|
||||
import type {
|
||||
CliBackendParseJsonlEvent,
|
||||
CliBackendParseJsonlLifecycleEvent,
|
||||
} from "openclaw/plugin-sdk/cli-backend";
|
||||
import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime";
|
||||
import { findCodeRegions, type CodeRegion } from "openclaw/plugin-sdk/text-chunking";
|
||||
|
||||
|
|
@ -166,6 +169,37 @@ function hasClaudeRawToolInvocation(text: string): boolean {
|
|||
return false;
|
||||
}
|
||||
|
||||
function parseClaudeJsonlRecord(line: string): Record<string, unknown> | null {
|
||||
try {
|
||||
const parsed: unknown = JSON.parse(line);
|
||||
return isRecord(parsed) ? parsed : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Project Claude-owned lifecycle records without widening the legacy parser event union. */
|
||||
export const parseClaudeCliJsonlLifecycleEvent: CliBackendParseJsonlLifecycleEvent = (line) => {
|
||||
if (!line.includes("compacting") && !line.includes("compact_result")) {
|
||||
return null;
|
||||
}
|
||||
const parsed = parseClaudeJsonlRecord(line);
|
||||
if (!parsed) {
|
||||
return null;
|
||||
}
|
||||
if (parsed.compact_result === "success" || parsed.compact_result === "failed") {
|
||||
return {
|
||||
kind: "compaction",
|
||||
phase: "end",
|
||||
completed: parsed.compact_result === "success",
|
||||
};
|
||||
}
|
||||
if (parsed.type === "system" && parsed.subtype === "status") {
|
||||
return parsed.status === "compacting" ? { kind: "compaction", phase: "start" } : null;
|
||||
}
|
||||
return null;
|
||||
};
|
||||
|
||||
/** Reject malformed terminal Claude results before the generic CLI runner accepts them as prose. */
|
||||
export const parseClaudeCliJsonlEvent: CliBackendParseJsonlEvent = (line) => {
|
||||
const mightContainRawToolProtocol =
|
||||
|
|
@ -178,15 +212,11 @@ export const parseClaudeCliJsonlEvent: CliBackendParseJsonlEvent = (line) => {
|
|||
if (!mightContainRawToolProtocol) {
|
||||
return null;
|
||||
}
|
||||
|
||||
let parsed: unknown;
|
||||
try {
|
||||
parsed = JSON.parse(line);
|
||||
} catch {
|
||||
const parsed = parseClaudeJsonlRecord(line);
|
||||
if (!parsed) {
|
||||
return null;
|
||||
}
|
||||
if (
|
||||
!isRecord(parsed) ||
|
||||
parsed.type !== "result" ||
|
||||
typeof parsed.result !== "string" ||
|
||||
!hasClaudeRawToolInvocation(parsed.result)
|
||||
|
|
|
|||
|
|
@ -117,6 +117,7 @@ describe("Telegram QA transport adapter", () => {
|
|||
});
|
||||
mocks.userbotStart.mockResolvedValue({
|
||||
assertHealthy: mocks.userbotAssertHealthy,
|
||||
chatId: -100123,
|
||||
close: mocks.userbotClose,
|
||||
send: mocks.userbotSend,
|
||||
});
|
||||
|
|
@ -124,6 +125,41 @@ describe("Telegram QA transport adapter", () => {
|
|||
mocks.shouldRetainQaGatewayCredentialLease.mockResolvedValue(false);
|
||||
});
|
||||
|
||||
it("targets the SUT DM for direct-message-only scenarios", async () => {
|
||||
mocks.userbotStart.mockResolvedValueOnce({
|
||||
assertHealthy: mocks.userbotAssertHealthy,
|
||||
chatId: 200,
|
||||
close: mocks.userbotClose,
|
||||
send: mocks.userbotSend,
|
||||
});
|
||||
const adapter = await createTelegramQaTransportAdapter({
|
||||
adapterOptions: { transportPolicy: { directMessageOnly: true } },
|
||||
messages: {},
|
||||
} as never);
|
||||
|
||||
expect(mocks.userbotStart).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ chatId: "@sut_bot" }),
|
||||
);
|
||||
expect(adapter.createGatewayConfig?.({ baseUrl: "http://127.0.0.1:1234" })).toMatchObject({
|
||||
channels: {
|
||||
telegram: {
|
||||
accounts: {
|
||||
sut: { allowFrom: ["100"], dmPolicy: "allowlist" },
|
||||
},
|
||||
},
|
||||
},
|
||||
});
|
||||
expect(adapter.buildAgentDelivery({ target: "dm:qa-operator" })).toEqual({
|
||||
channel: "telegram",
|
||||
to: "100",
|
||||
replyChannel: "telegram",
|
||||
replyTo: "100",
|
||||
});
|
||||
|
||||
await adapter.cleanup?.();
|
||||
await adapter.cleanupAfterGatewayStop?.();
|
||||
});
|
||||
|
||||
it("leases a Test Server userbot and isolates its shared group by default", async () => {
|
||||
const adapter = await createTelegramQaTransportAdapter({
|
||||
adapterOptions: {
|
||||
|
|
@ -162,6 +198,12 @@ describe("Telegram QA transport adapter", () => {
|
|||
},
|
||||
},
|
||||
});
|
||||
expect(adapter.buildAgentDelivery({ target: "group:qa-channel" })).toEqual({
|
||||
channel: "telegram",
|
||||
to: "-100123",
|
||||
replyChannel: "telegram",
|
||||
replyTo: "-100123",
|
||||
});
|
||||
|
||||
await adapter.cleanup?.();
|
||||
await adapter.cleanupAfterGatewayStop?.();
|
||||
|
|
@ -190,6 +232,7 @@ describe("Telegram QA transport adapter", () => {
|
|||
onUpdate = params.onUpdate;
|
||||
return {
|
||||
assertHealthy: mocks.userbotAssertHealthy,
|
||||
chatId: -100123,
|
||||
close: mocks.userbotClose,
|
||||
send: mocks.userbotSend,
|
||||
};
|
||||
|
|
@ -289,6 +332,7 @@ describe("Telegram QA transport adapter", () => {
|
|||
onUpdate = params.onUpdate;
|
||||
return {
|
||||
assertHealthy: mocks.userbotAssertHealthy,
|
||||
chatId: -100123,
|
||||
close: mocks.userbotClose,
|
||||
send: mocks.userbotSend,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -112,6 +112,11 @@ export async function createTelegramQaTransportAdapter(
|
|||
updateCount: 0,
|
||||
};
|
||||
const accountId = options.sutAccountId?.trim() || "sut";
|
||||
const directMessageOnly = options.transportPolicy?.directMessageOnly === true;
|
||||
const agentDeliveryTarget = directMessageOnly
|
||||
? credentialLease.payload.testerUserId
|
||||
: credentialLease.payload.groupId;
|
||||
let nativeChatId = Number(credentialLease.payload.groupId);
|
||||
let logicalConversationId = credentialLease.payload.groupId;
|
||||
let logicalConversationKind: "channel" | "direct" | "group" = "channel";
|
||||
const nativeMessageIds = new Map<string, number>();
|
||||
|
|
@ -148,7 +153,7 @@ export async function createTelegramQaTransportAdapter(
|
|||
observerState.updateCount += 1;
|
||||
observerState.relevantUpdateKinds.add(update.kind);
|
||||
if (
|
||||
update.chatId !== Number(credentialLease.payload.groupId) ||
|
||||
update.chatId !== nativeChatId ||
|
||||
update.senderId !== Number(credentialLease.payload.sutBotId)
|
||||
) {
|
||||
observerState.filteredCount += 1;
|
||||
|
|
@ -168,12 +173,13 @@ export async function createTelegramQaTransportAdapter(
|
|||
apiProxy = await skillRuntime.startApiProxy(leaseHealth);
|
||||
await apiProxy.drainUpdates(restored.sutToken);
|
||||
userbot = await TelegramUserbotDriver.start({
|
||||
chatId: restored.groupId,
|
||||
chatId: directMessageOnly ? `@${credentialLease.payload.sutUsername}` : restored.groupId,
|
||||
driverEnv: restored.driverEnv,
|
||||
leaseHealth,
|
||||
userDriverPath: skillRuntime.userDriverPath,
|
||||
onUpdate: observeUpdate,
|
||||
});
|
||||
nativeChatId = userbot.chatId;
|
||||
} catch (error) {
|
||||
const cleanupErrors: unknown[] = [];
|
||||
try {
|
||||
|
|
@ -275,6 +281,7 @@ export async function createTelegramQaTransportAdapter(
|
|||
createGatewayConfig: () =>
|
||||
buildTelegramQaConfig({} as OpenClawConfig, {
|
||||
apiRoot: activeApiProxy.apiRoot,
|
||||
directMessageOnly,
|
||||
groupId: credentialLease.payload.groupId,
|
||||
sutToken: credentialLease.payload.sutToken,
|
||||
testerUserId: credentialLease.payload.testerUserId,
|
||||
|
|
@ -287,9 +294,9 @@ export async function createTelegramQaTransportAdapter(
|
|||
}),
|
||||
buildAgentDelivery: () => ({
|
||||
channel: "telegram",
|
||||
to: credentialLease.payload.groupId,
|
||||
to: agentDeliveryTarget,
|
||||
replyChannel: "telegram",
|
||||
replyTo: credentialLease.payload.groupId,
|
||||
replyTo: agentDeliveryTarget,
|
||||
}),
|
||||
async handleAction() {
|
||||
throw new Error("Telegram live QA adapter does not implement transport actions");
|
||||
|
|
|
|||
|
|
@ -43,7 +43,6 @@ describe("Telegram QA profiles", () => {
|
|||
it("lets explicit scenarios override profile selection", () => {
|
||||
expect(
|
||||
resolveTelegramQaScenarioIds({
|
||||
profile: "release",
|
||||
providerMode: "live-frontier",
|
||||
scenarioIds: ["telegram-help-command"],
|
||||
}),
|
||||
|
|
@ -74,6 +73,15 @@ describe("Telegram QA profiles", () => {
|
|||
).toEqual(["telegram-queue-invalid-mode"]);
|
||||
});
|
||||
|
||||
it("selects the Claude CLI compaction final-priority regression explicitly", () => {
|
||||
expect(
|
||||
resolveTelegramQaScenarioIds({
|
||||
providerMode: "live-frontier",
|
||||
scenarioIds: ["telegram-claude-cli-compaction-final-priority"],
|
||||
}),
|
||||
).toEqual(["telegram-claude-cli-compaction-final-priority"]);
|
||||
});
|
||||
|
||||
it("rejects unknown profiles and channel-ineligible explicit scenarios", () => {
|
||||
expect(() =>
|
||||
resolveTelegramQaScenarioIds({ providerMode: "live-frontier", profile: "transport" }),
|
||||
|
|
|
|||
|
|
@ -36,6 +36,25 @@ describe("Telegram QA API boundary", () => {
|
|||
});
|
||||
});
|
||||
|
||||
it("allows only the leased tester in direct-message mode", () => {
|
||||
const config = buildTelegramQaConfig(
|
||||
{},
|
||||
{
|
||||
apiRoot: "http://127.0.0.1:8080",
|
||||
directMessageOnly: true,
|
||||
groupId: "-100123",
|
||||
sutToken: "placeholder",
|
||||
testerUserId: "1",
|
||||
sutAccountId: "sut",
|
||||
},
|
||||
);
|
||||
|
||||
expect(config.channels?.telegram?.accounts?.sut).toMatchObject({
|
||||
allowFrom: ["1"],
|
||||
dmPolicy: "allowlist",
|
||||
});
|
||||
});
|
||||
|
||||
it("waits for the selected Telegram account to connect", async () => {
|
||||
const call = vi
|
||||
.fn()
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ export function buildTelegramQaConfig(
|
|||
baseCfg: OpenClawConfig,
|
||||
params: {
|
||||
apiRoot: string;
|
||||
directMessageOnly?: boolean;
|
||||
groupId: string;
|
||||
sutAccountId: string;
|
||||
sutToken: string;
|
||||
|
|
@ -71,7 +72,9 @@ export function buildTelegramQaConfig(
|
|||
enabled: true,
|
||||
botToken: params.sutToken,
|
||||
apiRoot: params.apiRoot,
|
||||
dmPolicy: "disabled",
|
||||
...(params.directMessageOnly
|
||||
? { dmPolicy: "allowlist", allowFrom: [params.testerUserId] }
|
||||
: { dmPolicy: "disabled" }),
|
||||
groups: {
|
||||
[params.groupId]: {
|
||||
groupPolicy: "allowlist",
|
||||
|
|
|
|||
|
|
@ -63,6 +63,7 @@ describe("Telegram userbot driver runtime", () => {
|
|||
updates.push(update);
|
||||
},
|
||||
});
|
||||
expect(driver.chatId).toBe(-1001);
|
||||
try {
|
||||
await expect(driver.send({ text })).resolves.toMatchObject({
|
||||
messageId: 11,
|
||||
|
|
|
|||
|
|
@ -117,6 +117,7 @@ function waitForChildExit(child: ChildProcessWithoutNullStreams, timeoutMs: numb
|
|||
}
|
||||
|
||||
export class TelegramUserbotDriver {
|
||||
private activeChatId: number | undefined;
|
||||
private closing = false;
|
||||
private commandId = 0;
|
||||
private readonly pending = new Map<
|
||||
|
|
@ -211,6 +212,12 @@ export class TelegramUserbotDriver {
|
|||
return;
|
||||
}
|
||||
if (message.type === "ready") {
|
||||
const chatId = message.chatId;
|
||||
if (typeof chatId !== "number" || !Number.isInteger(chatId) || chatId === 0) {
|
||||
this.fail(new Error("Telegram userbot emitted an invalid ready chat id."));
|
||||
return;
|
||||
}
|
||||
this.activeChatId = chatId;
|
||||
this.readyResolve();
|
||||
return;
|
||||
}
|
||||
|
|
@ -271,6 +278,13 @@ export class TelegramUserbotDriver {
|
|||
}
|
||||
}
|
||||
|
||||
get chatId(): number {
|
||||
if (this.activeChatId === undefined) {
|
||||
throw new Error("Telegram userbot chat id is unavailable before readiness.");
|
||||
}
|
||||
return this.activeChatId;
|
||||
}
|
||||
|
||||
async send(params: { replyToMessageId?: number; text: string }): Promise<TelegramUserbotUpdate> {
|
||||
this.leaseHealth.assertHealthy();
|
||||
this.assertHealthy();
|
||||
|
|
|
|||
|
|
@ -68,6 +68,7 @@ const qaScenarioChannelSchema = z
|
|||
});
|
||||
|
||||
const qaScenarioTransportPolicySchema = z.object({
|
||||
directMessageOnly: z.literal(true).optional(),
|
||||
requireGroupMention: z.literal(true).optional(),
|
||||
senderAllowlist: z.array(z.string().trim().min(1)).min(1).optional(),
|
||||
topLevelReplies: z.literal(true).optional(),
|
||||
|
|
|
|||
|
|
@ -555,6 +555,7 @@ describe("qa suite planning helpers", () => {
|
|||
it("isolates and collects scenario-declared transport policy", () => {
|
||||
const scenario = makeQaSuiteTestScenario("sender-policy", {
|
||||
transportPolicy: {
|
||||
directMessageOnly: true,
|
||||
requireGroupMention: true,
|
||||
senderAllowlist: ["driver"],
|
||||
},
|
||||
|
|
@ -562,6 +563,7 @@ describe("qa suite planning helpers", () => {
|
|||
|
||||
expect(scenarioRequiresIsolatedQaSuiteWorker(scenario)).toBe(true);
|
||||
expect(collectQaSuiteTransportPolicy([scenario])).toEqual({
|
||||
directMessageOnly: true,
|
||||
requireGroupMention: true,
|
||||
senderAllowlist: ["driver"],
|
||||
});
|
||||
|
|
|
|||
|
|
@ -274,6 +274,7 @@ function collectQaSuiteGatewayRuntimeOptions(
|
|||
function collectQaSuiteTransportPolicy(
|
||||
scenarios: ReturnType<typeof readQaBootstrapScenarioCatalog>["scenarios"],
|
||||
) {
|
||||
let directMessageOnly = false;
|
||||
let requireGroupMention = false;
|
||||
let topLevelReplies = false;
|
||||
let senderAllowlist: readonly string[] | undefined;
|
||||
|
|
@ -282,6 +283,7 @@ function collectQaSuiteTransportPolicy(
|
|||
continue;
|
||||
}
|
||||
const policy = scenario.execution.transportPolicy;
|
||||
directMessageOnly ||= policy?.directMessageOnly === true;
|
||||
requireGroupMention ||= policy?.requireGroupMention === true;
|
||||
topLevelReplies ||= policy?.topLevelReplies === true;
|
||||
if (!policy?.senderAllowlist) {
|
||||
|
|
@ -295,8 +297,9 @@ function collectQaSuiteTransportPolicy(
|
|||
}
|
||||
senderAllowlist = policy.senderAllowlist;
|
||||
}
|
||||
return requireGroupMention || topLevelReplies || senderAllowlist
|
||||
return directMessageOnly || requireGroupMention || topLevelReplies || senderAllowlist
|
||||
? {
|
||||
...(directMessageOnly ? { directMessageOnly: true as const } : {}),
|
||||
...(requireGroupMention ? { requireGroupMention: true as const } : {}),
|
||||
...(senderAllowlist ? { senderAllowlist } : {}),
|
||||
...(topLevelReplies ? { topLevelReplies: true as const } : {}),
|
||||
|
|
|
|||
|
|
@ -439,6 +439,21 @@ async function materializeAnswerLaneBeforeRotation(turn: Turn): Promise<void> {
|
|||
await handlePreviewFinalizedResult(turn, result);
|
||||
}
|
||||
|
||||
async function cleanupProgressWithoutBlockingFinal(
|
||||
phase: "discard" | "teardown",
|
||||
cleanup: () => Promise<void>,
|
||||
): Promise<void> {
|
||||
try {
|
||||
await cleanup();
|
||||
} catch (err) {
|
||||
// Preview cleanup is best-effort; dropping the durable final is worse than
|
||||
// leaving stale progress visible for Telegram to expire or replace later.
|
||||
logVerbose(
|
||||
`telegram progress ${phase} failed before final delivery: ${formatErrorMessage(err)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async function deliverTelegramProgressModeFinalAnswer(
|
||||
turn: Turn,
|
||||
payload: ReplyPayload,
|
||||
|
|
@ -449,8 +464,15 @@ async function deliverTelegramProgressModeFinalAnswer(
|
|||
bindPendingFinalDelivery?: <T extends ReplyPayload>(payload: T) => T,
|
||||
): Promise<LaneDeliveryResult> {
|
||||
const afterAcceptedDraft = turn.answerLane.stream?.hasConsumedReplyTarget?.() === true;
|
||||
// Seal pending preview updates before the durable final send. This bounds
|
||||
// final latency to one in-flight edit and prevents stale progress overtaking it.
|
||||
await cleanupProgressWithoutBlockingFinal("discard", async () => {
|
||||
await turn.answerLane.stream?.discard?.();
|
||||
});
|
||||
if (payload.isError === true) {
|
||||
await teardownProgressWindow(turn);
|
||||
await cleanupProgressWithoutBlockingFinal("teardown", async () => {
|
||||
await teardownProgressWindow(turn);
|
||||
});
|
||||
const delivered = await sendPayload(turn, applyTextToPayload(payload, text), {
|
||||
afterAcceptedDraft,
|
||||
durable: true,
|
||||
|
|
@ -476,7 +498,9 @@ async function deliverTelegramProgressModeFinalAnswer(
|
|||
});
|
||||
// The final must dispatch before the activity window retires, so the answer
|
||||
// lane cannot accept follow-ups against a stale preview message.
|
||||
await teardownProgressWindow(turn);
|
||||
await cleanupProgressWithoutBlockingFinal("teardown", async () => {
|
||||
await teardownProgressWindow(turn);
|
||||
});
|
||||
if (!delivered) {
|
||||
return { kind: "skipped" };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -49,6 +49,26 @@ type TelegramProgressDraftState = {
|
|||
streamReasoningInProgressDraft: boolean;
|
||||
};
|
||||
|
||||
const TELEGRAM_COMPACTION_PROGRESS_ID = "context-compaction";
|
||||
|
||||
function buildTelegramCompactionProgressLine(
|
||||
phase: "start" | "complete" | "incomplete",
|
||||
): ChannelProgressDraftLine {
|
||||
const label = {
|
||||
start: "Compacting context...",
|
||||
complete: "Compaction complete",
|
||||
incomplete: "Compaction incomplete",
|
||||
}[phase];
|
||||
return {
|
||||
id: TELEGRAM_COMPACTION_PROGRESS_ID,
|
||||
kind: "item",
|
||||
icon: "🧹",
|
||||
label,
|
||||
text: `🧹 ${label}`,
|
||||
prefix: false,
|
||||
};
|
||||
}
|
||||
|
||||
export function createProgressState(
|
||||
config: TurnConfig,
|
||||
draftState: TelegramProgressDraftState,
|
||||
|
|
@ -113,6 +133,16 @@ export function canPushToolProgress(turn: Turn): boolean {
|
|||
);
|
||||
}
|
||||
|
||||
function canPushCompactionProgress(turn: Turn): boolean {
|
||||
return Boolean(
|
||||
turn.streamMode === "progress" &&
|
||||
turn.answerLane.stream &&
|
||||
!turn.answerLane.finalized &&
|
||||
!turn.finalAnswerDeliveryStarted &&
|
||||
!turn.finalAnswerDelivered,
|
||||
);
|
||||
}
|
||||
|
||||
async function pushProgressEvent(turn: Turn, event: () => Promise<boolean>): Promise<boolean> {
|
||||
return canPushToolProgress(turn) ? await event() : false;
|
||||
}
|
||||
|
|
@ -191,14 +221,39 @@ export async function handleToolStart(
|
|||
return await progressPromise;
|
||||
}
|
||||
|
||||
export async function handleCompactionStart(turn: Turn): Promise<boolean> {
|
||||
const progress = canPushCompactionProgress(turn)
|
||||
? turn.progressCompositor.pushToolProgress(buildTelegramCompactionProgressLine("start"), {
|
||||
startImmediately: true,
|
||||
flush: true,
|
||||
})
|
||||
: Promise.resolve(false);
|
||||
await turn.statusReactionController?.setCompacting();
|
||||
return await progress;
|
||||
}
|
||||
|
||||
export async function handleCompactionEnd(
|
||||
turn: Turn,
|
||||
payload?: CallbackPayload<"onCompactionEnd">,
|
||||
): Promise<boolean> {
|
||||
const progress = canPushCompactionProgress(turn)
|
||||
? turn.progressCompositor.pushToolProgress(
|
||||
buildTelegramCompactionProgressLine(
|
||||
payload?.completed === false ? "incomplete" : "complete",
|
||||
),
|
||||
{ startImmediately: true, flush: true },
|
||||
)
|
||||
: Promise.resolve(false);
|
||||
turn.statusReactionController?.cancelPending();
|
||||
await turn.statusReactionController?.setThinking();
|
||||
return await progress;
|
||||
}
|
||||
|
||||
export async function handleItemEvent(
|
||||
turn: Turn,
|
||||
payload: CallbackPayload<"onItemEvent">,
|
||||
): Promise<boolean> {
|
||||
if (payload.kind === "preamble") {
|
||||
if (turn.verboseProgressActive()) {
|
||||
return false;
|
||||
}
|
||||
let rendered = false;
|
||||
if (turn.streamMode === "progress") {
|
||||
rendered = await turn.progressCompositor.pushPreambleHeadline(payload.progressText, {
|
||||
|
|
|
|||
|
|
@ -371,16 +371,10 @@ export async function deliverReply(
|
|||
!turn.activeAnswerDraftIsToolProgressOnly &&
|
||||
!ownedByQueuedRotation &&
|
||||
segment.update.text.trimEnd() === turn.answerLane.lastPartialText.trimEnd();
|
||||
const isDurableProgressCommentary =
|
||||
turn.streamMode === "progress" &&
|
||||
info.kind === "block" &&
|
||||
effectivePayload.isCommentary === true;
|
||||
// CLI finals exclude separately classified commentary, so it must outlive the progress draft.
|
||||
const suppressProgressAnswerBlock =
|
||||
turn.streamMode === "progress" &&
|
||||
info.kind === "block" &&
|
||||
segment.lane === "answer" &&
|
||||
!isDurableProgressCommentary &&
|
||||
!reply.hasMedia &&
|
||||
!hasExecApprovalPayload(effectivePayload) &&
|
||||
telegramButtons === undefined;
|
||||
|
|
@ -429,7 +423,6 @@ export async function deliverReply(
|
|||
infoKind: info.kind,
|
||||
buttons: telegramButtons,
|
||||
...(isAskUserPayload ? { finalizePreview: true } : {}),
|
||||
allowStream: !isDurableProgressCommentary,
|
||||
onPlatformSendDispatch: info.onPlatformSendDispatch,
|
||||
assertPlatformSendAuthorized: info.assertPlatformSendAuthorized,
|
||||
bindPendingFinalDelivery: info.bindPendingFinalDelivery,
|
||||
|
|
|
|||
|
|
@ -26,6 +26,8 @@ import {
|
|||
canPushToolProgress,
|
||||
handleApprovalEvent,
|
||||
handleCommandOutput,
|
||||
handleCompactionEnd,
|
||||
handleCompactionStart,
|
||||
handleItemEvent,
|
||||
handlePatchSummary,
|
||||
handlePlanUpdate,
|
||||
|
|
@ -271,12 +273,11 @@ export async function runTelegramDispatchTurn(turn: Turn) {
|
|||
turn.streamMode === "progress" ? turn.commentaryProgressEnabled : undefined,
|
||||
progressPreambleEnabled: turn.progressPreambleEnabled,
|
||||
commentaryPayloadsEnabled: turn.progressPreambleEnabled,
|
||||
// Read the current getter after core freezes visibility so draft
|
||||
// and durable commentary cannot both own the same preamble.
|
||||
// The progress draft is the only commentary owner and retires before
|
||||
// the clean final. A durable copy would restore the queue burst this
|
||||
// owner boundary prevents; verbose still controls durable tool output.
|
||||
shouldDeliverCommentaryPayloads:
|
||||
turn.streamMode === "progress" && turn.commentaryProgressEnabled
|
||||
? () => turn.verboseProgressActive()
|
||||
: undefined,
|
||||
turn.progressPreambleEnabled === true ? () => false : undefined,
|
||||
reasoningPayloadsEnabled: turn.durableReasoningPayloadsEnabled,
|
||||
onToolStart: (payload) => handleToolStart(turn, payload),
|
||||
onItemEvent: (payload) => handleItemEvent(turn, payload),
|
||||
|
|
@ -303,19 +304,14 @@ export async function runTelegramDispatchTurn(turn: Turn) {
|
|||
},
|
||||
onCommandOutput: (payload) => handleCommandOutput(turn, payload),
|
||||
onPatchSummary: (payload) => handlePatchSummary(turn, payload),
|
||||
onCompactionStart: turn.statusReactionController
|
||||
? async () => {
|
||||
await turn.statusReactionController?.setCompacting();
|
||||
return false;
|
||||
}
|
||||
: undefined,
|
||||
onCompactionEnd: turn.statusReactionController
|
||||
? async () => {
|
||||
turn.statusReactionController?.cancelPending();
|
||||
await turn.statusReactionController?.setThinking();
|
||||
return false;
|
||||
}
|
||||
: undefined,
|
||||
// Ambient room events are intentionally invisible, including reactions.
|
||||
// User requests in group chats are not room_event turns and retain these callbacks.
|
||||
onCompactionStart: isRoomEvent
|
||||
? undefined
|
||||
: async () => await handleCompactionStart(turn),
|
||||
onCompactionEnd: isRoomEvent
|
||||
? undefined
|
||||
: async (payload) => await handleCompactionEnd(turn, payload),
|
||||
onModelSelected,
|
||||
},
|
||||
}),
|
||||
|
|
|
|||
|
|
@ -434,28 +434,64 @@ describeTelegramDispatch("dispatchTelegramMessage draft-failures-progress", () =
|
|||
expect(deliverReplies).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("keeps compaction replay on the same answer stream", async () => {
|
||||
it("shows compaction progress on the same answer stream", async () => {
|
||||
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
|
||||
const compactionFlushCounts: number[] = [];
|
||||
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(
|
||||
async ({ dispatcherOptions, replyOptions }) => {
|
||||
await replyOptions?.onPartialReply?.({ text: "Partial before compaction" });
|
||||
await replyOptions?.onToolStart?.({ name: "exec", phase: "start" });
|
||||
answerDraftStream.flush.mockClear();
|
||||
await replyOptions?.onCompactionStart?.();
|
||||
compactionFlushCounts.push(answerDraftStream.flush.mock.calls.length);
|
||||
await replyOptions?.onCompactionEnd?.({ completed: false });
|
||||
compactionFlushCounts.push(answerDraftStream.flush.mock.calls.length);
|
||||
await replyOptions?.onCompactionStart?.();
|
||||
compactionFlushCounts.push(answerDraftStream.flush.mock.calls.length);
|
||||
await replyOptions?.onCompactionEnd?.({ completed: true });
|
||||
compactionFlushCounts.push(answerDraftStream.flush.mock.calls.length);
|
||||
await replyOptions?.onPartialReply?.({ text: "Partial before compaction" });
|
||||
await dispatcherOptions.deliver({ text: "Final after compaction" }, { kind: "final" });
|
||||
return { queuedFinal: true };
|
||||
},
|
||||
);
|
||||
|
||||
await dispatchWithContext({ context: createContext() });
|
||||
await dispatchWithContext({ context: createContext(), streamMode: "progress" });
|
||||
|
||||
expect(answerDraftStream.forceNewMessage).not.toHaveBeenCalled();
|
||||
expect(answerDraftStream.update).toHaveBeenNthCalledWith(1, "Partial before compaction");
|
||||
expect(answerDraftStream.update).toHaveBeenNthCalledWith(
|
||||
2,
|
||||
"Final after compaction",
|
||||
expect.objectContaining({ onPlatformSendDispatch: expect.any(Function) }),
|
||||
expect(compactionFlushCounts).toEqual([1, 2, 3, 4]);
|
||||
expect(answerDraftStream.updatePreview).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ text: expect.stringContaining("Compacting context") }),
|
||||
);
|
||||
expect(deliverReplies).not.toHaveBeenCalled();
|
||||
expect(answerDraftStream.updatePreview).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ text: expect.stringContaining("Compaction incomplete") }),
|
||||
);
|
||||
expect(answerDraftStream.updatePreview).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ text: expect.stringContaining("Compaction complete") }),
|
||||
);
|
||||
expectDeliveredReply(0, { text: "Final after compaction" });
|
||||
expect(
|
||||
requireInvocationOrder(answerDraftStream.discard, 0, "compaction progress discard"),
|
||||
).toBeLessThan(requireInvocationOrder(deliverReplies, 0, "final reply delivery"));
|
||||
expectWindowRetiredAfterFinal(answerDraftStream, deliverReplies);
|
||||
});
|
||||
|
||||
it("keeps compaction reactions without rendering a draft outside progress mode", async () => {
|
||||
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
|
||||
const statusReactionController = createStatusReactionController();
|
||||
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(async ({ replyOptions }) => {
|
||||
await replyOptions?.onCompactionStart?.();
|
||||
await replyOptions?.onCompactionEnd?.({ completed: true });
|
||||
return { queuedFinal: true };
|
||||
});
|
||||
|
||||
await dispatchWithContext({
|
||||
context: createContext({ statusReactionController: statusReactionController as never }),
|
||||
streamMode: "partial",
|
||||
});
|
||||
|
||||
expect(answerDraftStream.updatePreview).not.toHaveBeenCalled();
|
||||
expect(statusReactionController.setCompacting).toHaveBeenCalledTimes(1);
|
||||
expect(statusReactionController.cancelPending).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("rotates a tool-progress-only answer draft before streaming the final answer", async () => {
|
||||
|
|
@ -690,8 +726,8 @@ describeTelegramDispatch("dispatchTelegramMessage draft-failures-progress", () =
|
|||
expectWindowRetiredAfterFinal(answerDraftStream, deliverReplies);
|
||||
});
|
||||
|
||||
it("sends the final answer before retiring the progress window", async () => {
|
||||
// Deliver first so removing the progress window cannot move the final off screen.
|
||||
it("seals pending progress before sending the final answer", async () => {
|
||||
// Seal the preview queue first so stale progress cannot overtake the final.
|
||||
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
|
||||
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(
|
||||
async ({ dispatcherOptions, replyOptions }) => {
|
||||
|
|
@ -708,9 +744,45 @@ describeTelegramDispatch("dispatchTelegramMessage draft-failures-progress", () =
|
|||
});
|
||||
|
||||
expectDeliveredReply(0, { text: "All done" });
|
||||
expect(
|
||||
requireInvocationOrder(answerDraftStream.discard, 0, "progress draft discard"),
|
||||
).toBeLessThan(requireInvocationOrder(deliverReplies, 0, "final reply delivery"));
|
||||
expectWindowRetiredAfterFinal(answerDraftStream, deliverReplies);
|
||||
});
|
||||
|
||||
it.each([false, true])(
|
||||
"delivers the final when progress cleanup fails (isError=%s)",
|
||||
async (isError) => {
|
||||
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
|
||||
answerDraftStream.discard.mockRejectedValueOnce(new Error("discard failed"));
|
||||
answerDraftStream.rotateToNewMessageDeferringDelete.mockRejectedValueOnce(
|
||||
new Error("teardown failed"),
|
||||
);
|
||||
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(
|
||||
async ({ dispatcherOptions, replyOptions }) => {
|
||||
await replyOptions?.onToolStart?.({ name: "exec", phase: "start" });
|
||||
await dispatcherOptions.deliver(
|
||||
{ text: "Final survives cleanup", ...(isError ? { isError: true } : {}) },
|
||||
{ kind: "final" },
|
||||
);
|
||||
return { queuedFinal: true };
|
||||
},
|
||||
);
|
||||
|
||||
await dispatchWithContext({
|
||||
context: createContext(),
|
||||
streamMode: "progress",
|
||||
telegramCfg: { streaming: { mode: "progress" } },
|
||||
});
|
||||
|
||||
expectDeliveredReply(0, {
|
||||
text: "Final survives cleanup",
|
||||
...(isError ? { isError: true } : {}),
|
||||
});
|
||||
expect(deliverReplies).toHaveBeenCalledTimes(1);
|
||||
},
|
||||
);
|
||||
|
||||
it("retires the progress window when the final answer send is skipped", async () => {
|
||||
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
|
||||
deliverReplies.mockResolvedValue({ delivered: false });
|
||||
|
|
|
|||
|
|
@ -132,21 +132,20 @@ describeTelegramDispatch("dispatchTelegramMessage progress-lifecycle", () => {
|
|||
expectWindowRetiredAfterFinal(answerDraftStream, deliverReplies);
|
||||
});
|
||||
|
||||
it("keeps CLI pre-tool commentary after the progress window retires", async () => {
|
||||
const markers = "Test markers: caribou-lampion-473, fromage-quantique, satellite-en-tricot";
|
||||
setupDraftStreams({ answerMessageId: 2001 });
|
||||
it("keeps verbose CLI commentary bounded in the progress window so the final wins", async () => {
|
||||
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
|
||||
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(
|
||||
async ({ dispatcherOptions, replyOptions }) => {
|
||||
expect(replyOptions?.commentaryPayloadsEnabled).toBe(true);
|
||||
expect(replyOptions?.shouldDeliverCommentaryPayloads).toBeUndefined();
|
||||
await replyOptions?.onItemEvent?.({
|
||||
kind: "preamble",
|
||||
itemId: "commentary-1",
|
||||
progressText: markers,
|
||||
suppressDurableProgress: true,
|
||||
});
|
||||
await replyOptions?.onBlockReplyQueued?.({ text: markers, isCommentary: true });
|
||||
await dispatcherOptions.deliver({ text: markers, isCommentary: true }, { kind: "block" });
|
||||
replyOptions?.onVerboseProgressVisibility?.(() => true);
|
||||
expect(replyOptions?.shouldDeliverCommentaryPayloads?.()).toBe(false);
|
||||
for (let index = 1; index <= 10; index += 1) {
|
||||
await replyOptions?.onItemEvent?.({
|
||||
kind: "preamble",
|
||||
itemId: `commentary-${index}`,
|
||||
progressText: `Commentary ${index}`,
|
||||
});
|
||||
}
|
||||
await replyOptions?.onToolStart?.({ name: "Bash", phase: "start" });
|
||||
await dispatcherOptions.deliver({ text: "TEST DONE" }, { kind: "final" });
|
||||
return { queuedFinal: true };
|
||||
|
|
@ -156,10 +155,15 @@ describeTelegramDispatch("dispatchTelegramMessage progress-lifecycle", () => {
|
|||
await dispatchWithContext({
|
||||
context: createContext(),
|
||||
streamMode: "progress",
|
||||
telegramCfg: { streaming: { mode: "progress" } },
|
||||
telegramCfg: { streaming: { mode: "progress", progress: { commentary: true } } },
|
||||
});
|
||||
|
||||
expect(allDeliveredReplyTexts()).toEqual([markers, "TEST DONE"]);
|
||||
const lastPreview = answerDraftStream.updatePreview.mock.calls.at(-1)?.[0].text ?? "";
|
||||
expect(lastPreview).not.toContain("Commentary 1\n");
|
||||
expect(lastPreview).not.toContain("Commentary 2\n");
|
||||
expect(lastPreview).toContain("Commentary 3");
|
||||
expect(lastPreview).toContain("Commentary 10");
|
||||
expect(allDeliveredReplyTexts()).toEqual(["TEST DONE"]);
|
||||
});
|
||||
|
||||
it("never streams an interim answer block into the progress window (Discord parity)", async () => {
|
||||
|
|
|
|||
|
|
@ -788,7 +788,7 @@ describeTelegramDispatch("dispatchTelegramMessage progress-updates", () => {
|
|||
["active", true],
|
||||
["inactive", false],
|
||||
])(
|
||||
"freezes the durable commentary owner to verbose visibility %s",
|
||||
"keeps the draft as commentary owner when verbose visibility becomes %s",
|
||||
async (_label, verboseActive) => {
|
||||
const draftStream = createSequencedDraftStream(2001);
|
||||
createTelegramDraftStream.mockReturnValue(draftStream);
|
||||
|
|
@ -797,7 +797,7 @@ describeTelegramDispatch("dispatchTelegramMessage progress-updates", () => {
|
|||
expect(replyOptions?.commentaryPayloadsEnabled).toBe(true);
|
||||
expect(replyOptions?.shouldDeliverCommentaryPayloads?.()).toBe(false);
|
||||
replyOptions?.onVerboseProgressVisibility?.(() => verboseActive);
|
||||
expect(replyOptions?.shouldDeliverCommentaryPayloads?.()).toBe(verboseActive);
|
||||
expect(replyOptions?.shouldDeliverCommentaryPayloads?.()).toBe(false);
|
||||
await replyOptions?.onItemEvent?.({
|
||||
kind: "preamble",
|
||||
itemId: "preamble-1",
|
||||
|
|
@ -820,13 +820,8 @@ describeTelegramDispatch("dispatchTelegramMessage progress-updates", () => {
|
|||
const updates = draftStream.updatePreview.mock.calls
|
||||
.map(([preview]) => preview.text)
|
||||
.join("\n");
|
||||
if (verboseActive) {
|
||||
// The durable lane owns commentary: the draft must not repeat it.
|
||||
expect(updates).not.toContain("Checking recent context");
|
||||
} else {
|
||||
// The draft owns commentary: exactly one visible copy per preamble.
|
||||
expect(updates.split("Checking recent context")).toHaveLength(2);
|
||||
}
|
||||
// The draft owns commentary: exactly one visible copy per preamble.
|
||||
expect(updates.split("Checking recent context")).toHaveLength(2);
|
||||
},
|
||||
);
|
||||
|
||||
|
|
|
|||
|
|
@ -149,7 +149,7 @@ type TelegramProgressCompositor = {
|
|||
cancel: () => void;
|
||||
pushToolProgress: (
|
||||
line?: string | ChannelProgressDraftLine,
|
||||
options?: { toolName?: string; startImmediately?: boolean },
|
||||
options?: { toolName?: string; startImmediately?: boolean; flush?: boolean },
|
||||
) => Promise<boolean>;
|
||||
pushReasoningProgress: (text?: string, options?: { snapshot?: boolean }) => Promise<boolean>;
|
||||
pushCommentaryProgress: (text?: string, options?: { itemId?: string }) => Promise<boolean>;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,205 @@
|
|||
title: Telegram Claude CLI compaction final priority
|
||||
|
||||
scenario:
|
||||
id: telegram-claude-cli-compaction-final-priority
|
||||
surface: channels
|
||||
coverage:
|
||||
primary:
|
||||
- channels.streaming-final-reply
|
||||
secondary:
|
||||
- anthropic.claude-cli-compatibility
|
||||
objective: Verify source-shaped Claude CLI compaction and commentary events stay in one Telegram progress draft and settle before one final reply.
|
||||
successCriteria:
|
||||
- The real Anthropic CLI parser projects a compaction start and completion from source-shaped stream-json records.
|
||||
- Telegram Test Server records commentary and compaction notices on one editable progress message.
|
||||
- The unique final marker is delivered exactly once, after all progress edits.
|
||||
docsRefs:
|
||||
- docs/concepts/compaction.md
|
||||
- docs/concepts/streaming.md
|
||||
codeRefs:
|
||||
- extensions/anthropic/cli-output.ts
|
||||
- src/agents/cli-output-events.ts
|
||||
- src/auto-reply/reply/agent-runner-cli-dispatch.ts
|
||||
- extensions/telegram/src/bot-message-dispatch-delivery.ts
|
||||
execution:
|
||||
kind: flow
|
||||
channel: telegram
|
||||
providerMode: live-frontier
|
||||
retryCount: 0
|
||||
suiteIsolation: isolated
|
||||
isolationReason: Replaces the ephemeral gateway model route with a source-shaped Claude CLI fixture and enables session verbose mode.
|
||||
transportPolicy:
|
||||
directMessageOnly: true
|
||||
summary: Run source-shaped Claude compaction and commentary records through the real parser and Telegram Test Server, then prove one bounded progress draft precedes one final.
|
||||
config:
|
||||
modelRef: anthropic/claude-haiku-4-5
|
||||
modelId: claude-haiku-4-5
|
||||
conversationId: telegram-claude-cli-compaction
|
||||
senderId: qa-claude-cli-operator
|
||||
commentaryOne: CLAUDE-COMMENTARY-ONE
|
||||
commentaryTwo: CLAUDE-COMMENTARY-TWO
|
||||
finalMarker: CLAUDE-COMPACTION-TELEGRAM-FINAL-OK
|
||||
fixtureSource: |-
|
||||
#!/usr/bin/env node
|
||||
if (process.argv.includes("auth") && process.argv.includes("status")) {
|
||||
process.stdout.write(`${JSON.stringify({ loggedIn: true, authMethod: "fixture" })}\n`);
|
||||
process.exit(0);
|
||||
}
|
||||
if (process.argv.includes("--version")) {
|
||||
process.stdout.write("2.1.252\n");
|
||||
process.exit(0);
|
||||
}
|
||||
const sessionIndex = process.argv.indexOf("--session-id");
|
||||
const sessionId = sessionIndex >= 0 ? process.argv[sessionIndex + 1] : "00000000-0000-4000-8000-000000000001";
|
||||
const frames = [
|
||||
{ type: "init", session_id: sessionId },
|
||||
{ type: "assistant", message: { id: "msg-commentary-1", content: [{ type: "text", text: "CLAUDE-COMMENTARY-ONE" }], stop_reason: null } },
|
||||
{ type: "assistant", message: { id: "msg-commentary-1", content: [{ type: "text", text: "CLAUDE-COMMENTARY-ONE" }, { type: "tool_use", id: "tool-1", name: "Read", input: { file_path: "fixture-one.txt" } }], stop_reason: null } },
|
||||
{ type: "user", message: { role: "user", content: [{ type: "tool_result", tool_use_id: "tool-1", content: "fixture one", is_error: false }] } },
|
||||
{ type: "system", subtype: "status", status: "compacting" },
|
||||
{ type: "system", subtype: "status", status: null, compact_result: "success" },
|
||||
{ type: "assistant", message: { id: "msg-commentary-2", content: [{ type: "text", text: "CLAUDE-COMMENTARY-TWO" }], stop_reason: null } },
|
||||
{ type: "assistant", message: { id: "msg-commentary-2", content: [{ type: "text", text: "CLAUDE-COMMENTARY-TWO" }, { type: "tool_use", id: "tool-2", name: "Read", input: { file_path: "fixture-two.txt" } }], stop_reason: null } },
|
||||
{ type: "user", message: { role: "user", content: [{ type: "tool_result", tool_use_id: "tool-2", content: "fixture two", is_error: false }] } },
|
||||
{ type: "assistant", message: { id: "msg-final", content: [{ type: "text", text: "CLAUDE-COMPACTION-TELEGRAM-FINAL-OK" }], stop_reason: "end_turn" } },
|
||||
{ type: "result", subtype: "success", session_id: sessionId, result: "CLAUDE-COMPACTION-TELEGRAM-FINAL-OK", is_error: false }
|
||||
];
|
||||
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
|
||||
(async () => {
|
||||
for (const frame of frames) {
|
||||
process.stdout.write(`${JSON.stringify(frame)}\n`);
|
||||
await sleep(frame.type === "system" && frame.status === "compacting" ? 1500 : 250);
|
||||
}
|
||||
})().catch((error) => {
|
||||
process.stderr.write(`${error instanceof Error ? error.stack : String(error)}\n`);
|
||||
process.exitCode = 1;
|
||||
});
|
||||
|
||||
flow:
|
||||
steps:
|
||||
- name: keeps Claude compaction progress ahead of one Telegram final
|
||||
actions:
|
||||
- assert:
|
||||
expr: env.providerMode === "live-frontier" && transport.id === "telegram"
|
||||
message: this proof requires the live Telegram Test Server adapter
|
||||
- set: fixtureDir
|
||||
value:
|
||||
expr: path.join(env.gateway.tempRoot, "claude-cli-fixture")
|
||||
- set: fixturePath
|
||||
value:
|
||||
expr: path.join(fixtureDir, "claude")
|
||||
- call: fs.mkdir
|
||||
args: [{ ref: fixtureDir }, { recursive: true }]
|
||||
- call: fs.writeFile
|
||||
args: [{ ref: fixturePath }, { ref: config.fixtureSource }, utf8]
|
||||
- call: fs.chmod
|
||||
args: [{ ref: fixturePath }, 493]
|
||||
- call: env.gateway.restartAfterStateMutation
|
||||
args:
|
||||
- lambda:
|
||||
async: true
|
||||
params: [ctx]
|
||||
expr: |-
|
||||
(async () => {
|
||||
const cfg = JSON.parse(await fs.readFile(ctx.configPath, "utf8"));
|
||||
ctx.runtimeEnv.PATH = `${fixtureDir}${path.delimiter}${ctx.runtimeEnv.PATH ?? ""}`;
|
||||
cfg.plugins = cfg.plugins ?? {};
|
||||
cfg.plugins.allow = [...new Set([...(cfg.plugins.allow ?? []), "anthropic"])];
|
||||
cfg.plugins.entries = cfg.plugins.entries ?? {};
|
||||
cfg.plugins.entries.anthropic = { enabled: true };
|
||||
cfg.models = cfg.models ?? { mode: "merge", providers: {} };
|
||||
cfg.models.mode = "merge";
|
||||
cfg.models.providers = cfg.models.providers ?? {};
|
||||
cfg.models.providers.anthropic = {
|
||||
models: [{ id: config.modelId, name: "Claude CLI fixture" }],
|
||||
};
|
||||
cfg.agents.defaults.model = { primary: config.modelRef };
|
||||
cfg.agents.defaults.models = {
|
||||
[config.modelRef]: { agentRuntime: { id: "claude-cli" } },
|
||||
};
|
||||
cfg.agents.defaults.modelPolicy = { allow: [config.modelRef] };
|
||||
cfg.agents.defaults.compaction = {
|
||||
...(cfg.agents.defaults.compaction ?? {}),
|
||||
notifyUser: true,
|
||||
};
|
||||
cfg.agents.entries.qa.model = { primary: config.modelRef };
|
||||
const telegram = cfg.channels.telegram.accounts?.[transport.accountId] ?? cfg.channels.telegram;
|
||||
telegram.streaming = {
|
||||
mode: "progress",
|
||||
progress: { commentary: true },
|
||||
};
|
||||
await fs.writeFile(ctx.configPath, `${JSON.stringify(cfg, null, 2)}\n`, "utf8");
|
||||
})()
|
||||
- call: waitForGatewayHealthy
|
||||
args: [{ ref: env }, 60000]
|
||||
- call: waitForTransportReady
|
||||
args: [{ ref: env }, 60000]
|
||||
- resetTransport: true
|
||||
- set: commandStartIndex
|
||||
value:
|
||||
expr: state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length
|
||||
- sendInbound:
|
||||
conversation:
|
||||
id: { ref: config.conversationId }
|
||||
kind: direct
|
||||
senderId: { ref: config.senderId }
|
||||
senderName: QA Claude CLI Operator
|
||||
text: /verbose on
|
||||
nativeCommand: { name: verbose }
|
||||
- waitForOutbound:
|
||||
sinceIndex: { ref: commandStartIndex }
|
||||
timeoutMs: 60000
|
||||
- set: startCursor
|
||||
value:
|
||||
expr: state.getSnapshot().cursor
|
||||
- sendInbound:
|
||||
conversation:
|
||||
id: { ref: config.conversationId }
|
||||
kind: direct
|
||||
senderId: { ref: config.senderId }
|
||||
senderName: QA Claude CLI Operator
|
||||
text: Exercise the Claude compaction fixture.
|
||||
- waitForOutboundSequence:
|
||||
sinceCursor: { ref: startCursor }
|
||||
finalTextIncludes: { ref: config.finalMarker }
|
||||
finalSettleMs: 1500
|
||||
minimumPreviewEvents: 3
|
||||
timeoutMs: 90000
|
||||
saveAs: sequence
|
||||
- set: turnEvents
|
||||
value:
|
||||
expr: "state.getSnapshot().events.filter((event) => event.cursor > startCursor && ['outbound-message', 'message-edited', 'message-deleted'].includes(event.kind) && event.message?.accountId === transport.accountId)"
|
||||
- set: markerEvents
|
||||
value:
|
||||
expr: "turnEvents.filter((event) => event.kind !== 'message-deleted' && !event.message.deleted && event.message.text.includes(config.finalMarker))"
|
||||
- set: progressEvents
|
||||
value:
|
||||
expr: "turnEvents.filter((event) => event.cursor < markerEvents[0].cursor && event.kind !== 'message-deleted' && !event.message.deleted)"
|
||||
- set: commentaryEvents
|
||||
value:
|
||||
expr: "progressEvents.filter((event) => event.message.text.includes(config.commentaryOne) || event.message.text.includes(config.commentaryTwo))"
|
||||
- assert:
|
||||
expr: markerEvents.length === 1
|
||||
message:
|
||||
expr: "`expected one final marker event, got ${markerEvents.length}: ${JSON.stringify(turnEvents)}`"
|
||||
- assert:
|
||||
expr: "turnEvents.every((event) => event.cursor <= markerEvents[0].cursor)"
|
||||
message:
|
||||
expr: "`progress changed after final delivery: ${JSON.stringify(turnEvents)}`"
|
||||
- assert:
|
||||
expr: "new Set(commentaryEvents.map((event) => event.message.id)).size === 1"
|
||||
message:
|
||||
expr: "`expected one bounded commentary message: ${JSON.stringify(commentaryEvents)}`"
|
||||
- assert:
|
||||
expr: "progressEvents.some((event) => event.message.text.includes(config.commentaryOne)) && progressEvents.some((event) => event.message.text.includes(config.commentaryTwo))"
|
||||
message:
|
||||
expr: "`Claude commentary did not stay in the progress draft: ${JSON.stringify(progressEvents)}`"
|
||||
- assert:
|
||||
expr: "progressEvents.some((event) => event.message.text.includes('Compacting context'))"
|
||||
message:
|
||||
expr: "`Claude compaction was not visible in Telegram progress: ${JSON.stringify(progressEvents)}`"
|
||||
- assert:
|
||||
expr: "progressEvents.some((event) => event.message.text.includes('Compaction complete'))"
|
||||
message:
|
||||
expr: "`Claude compaction completion was not visible in Telegram progress: ${JSON.stringify(progressEvents)}`"
|
||||
detailsExpr: "JSON.stringify({ verdict: 'PASS', scenario: scenario.id, providerPath: 'claude-cli-fixture', transport: transport.id, eventKinds: turnEvents.map((event) => event.kind), commentaryMessageId: commentaryEvents[0].message.id, finalMessageId: markerEvents[0].message.id, finalMarkerCount: markerEvents.length, noEventsAfterFinal: true })"
|
||||
|
|
@ -151,12 +151,16 @@ describe("resolveCliBackendConfig", () => {
|
|||
|
||||
it("preserves the plugin-owned JSONL parser through runtime resolution", () => {
|
||||
const parseJsonlEvent = vi.fn();
|
||||
const parseJsonlLifecycleEvent = vi.fn();
|
||||
cliBackendsTesting.setDepsForTest({
|
||||
resolveRuntimeCliBackends: () => [runtimeEntry({ parseJsonlEvent })],
|
||||
resolveRuntimeCliBackends: () => [
|
||||
runtimeEntry({ parseJsonlEvent, parseJsonlLifecycleEvent }),
|
||||
],
|
||||
resolvePluginSetupCliBackend: () => undefined,
|
||||
});
|
||||
|
||||
expect(requireBackend().parseJsonlEvent).toBe(parseJsonlEvent);
|
||||
expect(requireBackend().parseJsonlLifecycleEvent).toBe(parseJsonlLifecycleEvent);
|
||||
});
|
||||
|
||||
it("normalizes the registered adapter with agent and runtime config context", () => {
|
||||
|
|
|
|||
|
|
@ -60,6 +60,7 @@ export type ResolvedCliBackend = {
|
|||
resolveExecutionArgs?: CliBackendPlugin["resolveExecutionArgs"];
|
||||
resolveModelId?: CliBackendPlugin["resolveModelId"];
|
||||
parseJsonlEvent?: CliBackendPlugin["parseJsonlEvent"];
|
||||
parseJsonlLifecycleEvent?: CliBackendPlugin["parseJsonlLifecycleEvent"];
|
||||
toolAvailabilityEnforcement?: CliBackendToolAvailabilityEnforcement;
|
||||
nativeToolMode?: CliBackendNativeToolMode;
|
||||
sideQuestionToolMode?: CliBackendSideQuestionToolMode;
|
||||
|
|
@ -325,6 +326,7 @@ export function resolveCliBackendConfig(
|
|||
resolveExecutionArgs: backend.resolveExecutionArgs,
|
||||
resolveModelId: backend.resolveModelId,
|
||||
parseJsonlEvent: backend.parseJsonlEvent,
|
||||
parseJsonlLifecycleEvent: backend.parseJsonlLifecycleEvent,
|
||||
toolAvailabilityEnforcement: backend.toolAvailabilityEnforcement,
|
||||
nativeToolMode: backend.nativeToolMode,
|
||||
sideQuestionToolMode: backend.sideQuestionToolMode,
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ import type {
|
|||
CliBackendConfig,
|
||||
CliBackendJsonlUsage,
|
||||
CliBackendParseJsonlEvent,
|
||||
CliBackendParseJsonlLifecycleEvent,
|
||||
} from "../plugins/cli-backend.types.js";
|
||||
import type { AcceptedSessionSpawn } from "./accepted-session-spawn.js";
|
||||
import type {
|
||||
|
|
@ -95,6 +96,8 @@ export type CliThinkingProgress = {
|
|||
progressTokens: number;
|
||||
};
|
||||
|
||||
export type CliCompactionDelta = { phase: "start" } | { phase: "end"; completed: boolean };
|
||||
|
||||
/** Tool-call start event reconstructed from CLI stream output. */
|
||||
export type CliToolUseStartDelta = {
|
||||
toolCallId: string;
|
||||
|
|
@ -116,9 +119,11 @@ export type CliJsonlStreamingParserOptions = {
|
|||
backend: CliBackendConfig;
|
||||
providerId: string;
|
||||
parseJsonlEvent?: CliBackendParseJsonlEvent;
|
||||
parseJsonlLifecycleEvent?: CliBackendParseJsonlLifecycleEvent;
|
||||
onAssistantDelta: (delta: CliStreamingDelta) => void;
|
||||
onThinkingDelta?: (delta: CliThinkingDelta) => void;
|
||||
onThinkingProgress?: (progress: CliThinkingProgress) => void;
|
||||
onCompaction?: (delta: CliCompactionDelta) => void;
|
||||
onToolUseStart?: (delta: CliToolUseStartDelta) => void;
|
||||
onToolResult?: (delta: CliToolResultDelta) => void;
|
||||
onDisplayToolUseStart?: (delta: CliToolUseStartDelta) => void;
|
||||
|
|
|
|||
|
|
@ -1,5 +1,9 @@
|
|||
import { describe, expect, it, vi } from "vitest";
|
||||
import type { CliToolResultDelta, CliToolUseStartDelta } from "./cli-output-contracts.js";
|
||||
import type {
|
||||
CliCompactionDelta,
|
||||
CliToolResultDelta,
|
||||
CliToolUseStartDelta,
|
||||
} from "./cli-output-contracts.js";
|
||||
import { createCliJsonlStreamingParser } from "./cli-output-stream.js";
|
||||
|
||||
function joinJsonlFrames(...frames: unknown[]) {
|
||||
|
|
@ -284,6 +288,47 @@ describe("createCliJsonlStreamingParser events", () => {
|
|||
}
|
||||
});
|
||||
|
||||
it("projects lifecycle events without dropping fork-successor session metadata", () => {
|
||||
const compactionEvents: CliCompactionDelta[] = [];
|
||||
const sessionIds: string[] = [];
|
||||
const parseJsonlEvent = vi.fn(() => ({ kind: "result" as const, text: "done" }));
|
||||
const parser = createCliJsonlStreamingParser({
|
||||
backend: { command: "acme", output: "jsonl", sessionIdFields: ["session_id"] },
|
||||
providerId: "acme-cli",
|
||||
parseJsonlLifecycleEvent: (line) => {
|
||||
const event = JSON.parse(line) as { type?: string; completed?: boolean };
|
||||
return event.type === "compaction"
|
||||
? event.completed === undefined
|
||||
? { kind: "compaction", phase: "start" }
|
||||
: { kind: "compaction", phase: "end", completed: event.completed }
|
||||
: null;
|
||||
},
|
||||
parseJsonlEvent,
|
||||
onAssistantDelta: () => {},
|
||||
onCompaction: (event) => compactionEvents.push(event),
|
||||
onSessionId: (sessionId) => sessionIds.push(sessionId),
|
||||
});
|
||||
|
||||
parser.push(
|
||||
[
|
||||
JSON.stringify({ type: "compaction", session_id: "fork-successor" }),
|
||||
JSON.stringify({ type: "compaction", completed: true }),
|
||||
JSON.stringify({ type: "result" }),
|
||||
"",
|
||||
].join("\n"),
|
||||
);
|
||||
parser.finish();
|
||||
|
||||
expect(compactionEvents).toEqual([{ phase: "start" }, { phase: "end", completed: true }]);
|
||||
expect(sessionIds).toEqual(["fork-successor"]);
|
||||
expect(parseJsonlEvent).toHaveBeenCalledOnce();
|
||||
expect(parser.getOutput()).toEqual({
|
||||
text: "done",
|
||||
sessionId: "fork-successor",
|
||||
usage: undefined,
|
||||
});
|
||||
});
|
||||
|
||||
it("streams detailed Gemini error events over generic result errors", () => {
|
||||
const parser = createCliJsonlStreamingParser({
|
||||
backend: {
|
||||
|
|
|
|||
44
src/agents/cli-output-lifecycle.ts
Normal file
44
src/agents/cli-output-lifecycle.ts
Normal file
|
|
@ -0,0 +1,44 @@
|
|||
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
|
||||
import { formatErrorMessage } from "../infra/errors.js";
|
||||
import type {
|
||||
CliBackendConfig,
|
||||
CliBackendParseJsonlLifecycleEvent,
|
||||
CliBackendParsedJsonlLifecycleEvent,
|
||||
} from "../plugins/cli-backend.types.js";
|
||||
import type { CliCompactionDelta } from "./cli-output-contracts.js";
|
||||
|
||||
export function parseCliBackendLifecycleLine(params: {
|
||||
line: string;
|
||||
backendId: string;
|
||||
backend: CliBackendConfig;
|
||||
parse?: CliBackendParseJsonlLifecycleEvent;
|
||||
}): { events: readonly CliBackendParsedJsonlLifecycleEvent[] } | { errorText: string } | undefined {
|
||||
if (!params.parse) {
|
||||
return undefined;
|
||||
}
|
||||
try {
|
||||
const parsed = params.parse(params.line, {
|
||||
backendId: params.backendId,
|
||||
backend: params.backend,
|
||||
});
|
||||
if (parsed == null) {
|
||||
return undefined;
|
||||
}
|
||||
return { events: "kind" in parsed ? [parsed] : parsed };
|
||||
} catch (error) {
|
||||
return {
|
||||
errorText: truncateUtf16Safe(
|
||||
`CLI backend ${params.backendId} JSONL lifecycle parser failed: ${formatErrorMessage(error)}`,
|
||||
500,
|
||||
),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
export function projectCliBackendLifecycleEvent(
|
||||
event: CliBackendParsedJsonlLifecycleEvent,
|
||||
): CliCompactionDelta {
|
||||
return event.phase === "start"
|
||||
? { phase: "start" }
|
||||
: { phase: "end", completed: event.completed };
|
||||
}
|
||||
|
|
@ -539,6 +539,28 @@ describe("parseCliOutput", () => {
|
|||
|
||||
expect(result).toEqual(expected);
|
||||
});
|
||||
|
||||
it("keeps the missing-result failure after compaction-only metadata", () => {
|
||||
const result = parseCliOutput({
|
||||
raw: JSON.stringify({ type: "system", subtype: "status", status: "compacting" }),
|
||||
backend: {
|
||||
command: "claude",
|
||||
output: "jsonl",
|
||||
jsonlDialect: "claude-stream-json",
|
||||
sessionIdFields: ["session_id"],
|
||||
},
|
||||
providerId: "claude-cli",
|
||||
parseJsonlLifecycleEvent: () => ({ kind: "compaction", phase: "start" }),
|
||||
outputMode: "jsonl",
|
||||
});
|
||||
|
||||
expect(result).toEqual({
|
||||
text: "",
|
||||
sessionId: undefined,
|
||||
usage: undefined,
|
||||
errorText: "CLI stream-json output ended without a result event.",
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
describe("parseCliJsonl record usage", () => {
|
||||
|
|
|
|||
|
|
@ -59,6 +59,21 @@ export function isClaudeStreamJsonResult(params: {
|
|||
return supportsCliJsonlToolEvents(params) && params.parsed.type === "result";
|
||||
}
|
||||
|
||||
export function isClaudeSyntheticNoResponse(parsed: Record<string, unknown>): boolean {
|
||||
if (parsed.type !== "assistant" || !isRecord(parsed.message)) {
|
||||
return false;
|
||||
}
|
||||
const message = parsed.message;
|
||||
return (
|
||||
message.model === "<synthetic>" &&
|
||||
Array.isArray(message.content) &&
|
||||
message.content.length === 1 &&
|
||||
isRecord(message.content[0]) &&
|
||||
message.content[0].type === "text" &&
|
||||
message.content[0].text === "No response requested."
|
||||
);
|
||||
}
|
||||
|
||||
function extractJsonObjectCandidates(raw: string): string[] {
|
||||
return extractBalancedJsonFragments(raw, { openers: ["{"] }).map((fragment) => fragment.json);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,10 +26,12 @@ import {
|
|||
projectCliBackendEvent,
|
||||
projectCliTaggedReasoning,
|
||||
} from "./cli-output-events.js";
|
||||
import * as cliOutputLifecycle from "./cli-output-lifecycle.js";
|
||||
import {
|
||||
decodeCliRecords,
|
||||
isClaudeStreamJsonDialect,
|
||||
isClaudeStreamJsonResult,
|
||||
isClaudeSyntheticNoResponse,
|
||||
isClaudeSubagentRecord,
|
||||
isGeminiStreamJsonDialect,
|
||||
isStreamJsonDialect,
|
||||
|
|
@ -57,22 +59,6 @@ const CLI_STREAM_JSON_OUTPUT_LIMITS = Object.freeze({
|
|||
maxTurnLines: CLI_STREAM_JSON_DEFAULT_MAX_TURN_LINES,
|
||||
} satisfies CliStreamJsonOutputLimits);
|
||||
|
||||
function isClaudeSyntheticNoResponse(parsed: Record<string, unknown>): boolean {
|
||||
if (parsed.type !== "assistant" || !isRecord(parsed.message)) {
|
||||
return false;
|
||||
}
|
||||
const message = parsed.message;
|
||||
if (message.model !== "<synthetic>" || !Array.isArray(message.content)) {
|
||||
return false;
|
||||
}
|
||||
return (
|
||||
message.content.length === 1 &&
|
||||
isRecord(message.content[0]) &&
|
||||
message.content[0].type === "text" &&
|
||||
message.content[0].text === "No response requested."
|
||||
);
|
||||
}
|
||||
|
||||
/** Frames arbitrary stdout chunks while bounding each individual raw JSONL line. */
|
||||
function frameBoundedCliJsonlChunk(
|
||||
state: { pending: string },
|
||||
|
|
@ -274,10 +260,40 @@ export function createCliJsonlStreamingParser(params: CliJsonlStreamingParserOpt
|
|||
return false;
|
||||
};
|
||||
|
||||
const observeSessionId = (parsed: Record<string, unknown>) => {
|
||||
const parsedSessionId = pickCliSessionId(parsed, params.backend);
|
||||
if (parsedSessionId && parsedSessionId !== sessionId) {
|
||||
sessionId = parsedSessionId;
|
||||
params.onSessionId?.(parsedSessionId);
|
||||
}
|
||||
};
|
||||
|
||||
const handleCustomJsonlLine = (line: string, rawLine: string): boolean => {
|
||||
if (parseErrorText) {
|
||||
return true;
|
||||
}
|
||||
const lifecycle = cliOutputLifecycle.parseCliBackendLifecycleLine({
|
||||
line,
|
||||
backendId: params.providerId,
|
||||
backend: params.backend,
|
||||
parse: params.parseJsonlLifecycleEvent,
|
||||
});
|
||||
if (lifecycle) {
|
||||
if ("errorText" in lifecycle) {
|
||||
parseErrorText = lifecycle.errorText;
|
||||
} else {
|
||||
if (claudeStreamJson && !accountClaudeJsonlLine(rawLine.length)) {
|
||||
return true;
|
||||
}
|
||||
for (const parsed of decodeCliRecords(line)) {
|
||||
observeSessionId(parsed);
|
||||
}
|
||||
for (const event of lifecycle.events) {
|
||||
params.onCompaction?.(cliOutputLifecycle.projectCliBackendLifecycleEvent(event));
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
if (!params.parseJsonlEvent) {
|
||||
return false;
|
||||
}
|
||||
|
|
@ -313,14 +329,10 @@ export function createCliJsonlStreamingParser(params: CliJsonlStreamingParserOpt
|
|||
if (parseErrorText) {
|
||||
return;
|
||||
}
|
||||
const parsedSessionId = pickCliSessionId(parsed, params.backend);
|
||||
if (parsed.type === "result" && isStreamJsonDialect(params)) {
|
||||
sawTerminalResult = true;
|
||||
}
|
||||
if (parsedSessionId && parsedSessionId !== sessionId) {
|
||||
sessionId = parsedSessionId;
|
||||
params.onSessionId?.(parsedSessionId);
|
||||
}
|
||||
observeSessionId(parsed);
|
||||
const nextUsage = readCliUsage(parsed);
|
||||
const isClaudeTerminalResult =
|
||||
isClaudeStreamJsonDialect({
|
||||
|
|
|
|||
|
|
@ -4,7 +4,11 @@
|
|||
* reconstruction.
|
||||
*/
|
||||
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
|
||||
import type { CliBackendConfig, CliBackendParseJsonlEvent } from "../plugins/cli-backend.types.js";
|
||||
import type {
|
||||
CliBackendConfig,
|
||||
CliBackendParseJsonlEvent,
|
||||
CliBackendParseJsonlLifecycleEvent,
|
||||
} from "../plugins/cli-backend.types.js";
|
||||
import type { CliOutput } from "./cli-output-contracts.js";
|
||||
import {
|
||||
collectExplicitCliErrorText,
|
||||
|
|
@ -53,6 +57,7 @@ export function parseCliOutput(params: {
|
|||
backend: CliBackendConfig;
|
||||
providerId: string;
|
||||
parseJsonlEvent?: CliBackendParseJsonlEvent;
|
||||
parseJsonlLifecycleEvent?: CliBackendParseJsonlLifecycleEvent;
|
||||
outputMode?: "json" | "jsonl" | "text";
|
||||
fallbackSessionId?: string;
|
||||
}): CliOutput {
|
||||
|
|
@ -65,6 +70,7 @@ export function parseCliOutput(params: {
|
|||
backend: params.backend,
|
||||
providerId: params.providerId,
|
||||
parseJsonlEvent: params.parseJsonlEvent,
|
||||
parseJsonlLifecycleEvent: params.parseJsonlLifecycleEvent,
|
||||
onAssistantDelta: () => {},
|
||||
});
|
||||
parser.push(params.raw);
|
||||
|
|
|
|||
|
|
@ -97,6 +97,33 @@ describe("cli tool result events", () => {
|
|||
}
|
||||
});
|
||||
|
||||
it("emits canonical CLI compaction lifecycle events", () => {
|
||||
const runId = "run-compaction-events";
|
||||
const handlers = createCliEventHandlers({
|
||||
context: buildContext(runId),
|
||||
toolTracking: buildToolTracking(),
|
||||
getRunState: () => ({ failed: false, error: undefined }),
|
||||
});
|
||||
const events: AgentEventRuntimePayload[] = [];
|
||||
const dispose = onAgentEvent((event) => {
|
||||
if (event.runId === runId && event.stream === "compaction") {
|
||||
events.push(event);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
handlers.emitCliCompaction({ phase: "start" });
|
||||
handlers.emitCliCompaction({ phase: "end", completed: true });
|
||||
|
||||
expect(events.map((event) => event.data)).toEqual([
|
||||
{ phase: "start", backend: "claude-cli" },
|
||||
{ phase: "end", backend: "claude-cli", completed: true },
|
||||
]);
|
||||
} finally {
|
||||
dispose();
|
||||
}
|
||||
});
|
||||
|
||||
it("keeps correlated result args without adding them to display results", () => {
|
||||
const runId = "run-tool-result-args";
|
||||
const handlers = createCliEventHandlers({
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
import { emitAgentEvent } from "../../infra/agent-events.js";
|
||||
import { emitTrustedDiagnosticEvent } from "../../infra/diagnostic-events.js";
|
||||
import type {
|
||||
CliCompactionDelta,
|
||||
CliStreamingDelta,
|
||||
CliThinkingDelta,
|
||||
CliThinkingProgress,
|
||||
|
|
@ -288,6 +289,19 @@ export function createCliEventHandlers(params: {
|
|||
emitParsedToolTerminal(event);
|
||||
emitCliToolResult(event);
|
||||
};
|
||||
const emitCliCompaction = (event: CliCompactionDelta) => {
|
||||
observedCliActivity = true;
|
||||
if (emitLiveEvents) {
|
||||
emitAgentEvent({
|
||||
runId: runParams.runId,
|
||||
stream: "compaction",
|
||||
data: {
|
||||
...event,
|
||||
backend: context.backendResolved.id,
|
||||
},
|
||||
});
|
||||
}
|
||||
};
|
||||
const finalizeParsedTools = () => {
|
||||
for (const [toolCallId, activeTool] of Array.from(activeParsedTools)) {
|
||||
emitParsedToolTerminal({
|
||||
|
|
@ -380,6 +394,7 @@ export function createCliEventHandlers(params: {
|
|||
emitCliDisplayToolResult,
|
||||
emitParsedToolUseStart,
|
||||
emitParsedToolResult,
|
||||
emitCliCompaction,
|
||||
finalizeParsedTools,
|
||||
emitCliCommentaryText,
|
||||
emitCliAssistantDelta,
|
||||
|
|
|
|||
|
|
@ -114,9 +114,11 @@ export async function executeCliProcess(params: {
|
|||
backend: params.backend,
|
||||
providerId: context.backendResolved.id,
|
||||
parseJsonlEvent: context.backendResolved.parseJsonlEvent,
|
||||
parseJsonlLifecycleEvent: context.backendResolved.parseJsonlLifecycleEvent,
|
||||
onAssistantDelta: params.events.emitCliAssistantDelta,
|
||||
onThinkingDelta: params.events.emitCliThinkingDelta,
|
||||
onThinkingProgress: params.events.emitCliThinkingProgress,
|
||||
onCompaction: params.events.emitCliCompaction,
|
||||
onToolUseStart: params.events.emitParsedToolUseStart,
|
||||
onToolResult: params.events.emitParsedToolResult,
|
||||
onDisplayToolUseStart: params.events.emitCliDisplayToolUseStart,
|
||||
|
|
|
|||
|
|
@ -351,8 +351,10 @@ export type GetReplyOptions = {
|
|||
}) => Promise<ProgressCallbackResult> | ProgressCallbackResult;
|
||||
/** Called when context auto-compaction starts (allows UX feedback during the pause). */
|
||||
onCompactionStart?: () => Promise<ProgressCallbackResult> | ProgressCallbackResult;
|
||||
/** Called when context auto-compaction completes. */
|
||||
onCompactionEnd?: () => Promise<ProgressCallbackResult> | ProgressCallbackResult;
|
||||
/** Called when context auto-compaction ends; omitted outcome means completed for legacy callers. */
|
||||
onCompactionEnd?: (payload?: {
|
||||
completed: boolean;
|
||||
}) => Promise<ProgressCallbackResult> | ProgressCallbackResult;
|
||||
/** Called when the actual model is selected (including after fallback).
|
||||
* Use this to get model/provider/thinkLevel for responsePrefix template interpolation. */
|
||||
onModelSelected?: (ctx: ModelSelectedContext) => void;
|
||||
|
|
|
|||
|
|
@ -221,6 +221,8 @@ export async function runCliFallbackCandidate(
|
|||
onReasoningProgress: async (payload) => {
|
||||
await turn.opts?.onReasoningProgress?.(payload);
|
||||
},
|
||||
onCompactionStart: turn.opts?.onCompactionStart,
|
||||
onCompactionEnd: turn.opts?.onCompactionEnd,
|
||||
onToolEvent: async (payload) => {
|
||||
if (!params.preserveProgressCallbackStartOrder) {
|
||||
const commandBearing = await cliToolSummaryTracker.noteToolEvent(payload);
|
||||
|
|
|
|||
|
|
@ -52,6 +52,57 @@ afterEach(() => {
|
|||
});
|
||||
|
||||
describe("runCliAgentWithLifecycle", () => {
|
||||
it("bridges completed CLI compaction lifecycles to reply callbacks", async () => {
|
||||
cliDispatchState.runCliAgentMock.mockImplementationOnce(async (params: { runId: string }) => {
|
||||
emitAgentEvent({
|
||||
runId: params.runId,
|
||||
stream: "compaction",
|
||||
data: { phase: "start", backend: "claude-cli" },
|
||||
});
|
||||
emitAgentEvent({
|
||||
runId: params.runId,
|
||||
stream: "compaction",
|
||||
data: { phase: "end", backend: "claude-cli", completed: false },
|
||||
});
|
||||
emitAgentEvent({
|
||||
runId: params.runId,
|
||||
stream: "compaction",
|
||||
data: { phase: "start", backend: "claude-cli" },
|
||||
});
|
||||
emitAgentEvent({
|
||||
runId: params.runId,
|
||||
stream: "compaction",
|
||||
data: { phase: "end", backend: "claude-cli", completed: true },
|
||||
});
|
||||
return { payloads: [], meta: { durationMs: 1 } };
|
||||
});
|
||||
const callbacks: string[] = [];
|
||||
|
||||
await runCliAgentWithLifecycle({
|
||||
runId: "run-compaction-bridge",
|
||||
provider: "claude-cli",
|
||||
onCompactionStart: async () => {
|
||||
callbacks.push("start");
|
||||
},
|
||||
onCompactionEnd: async (payload) => {
|
||||
callbacks.push(payload?.completed === false ? "incomplete" : "end");
|
||||
},
|
||||
runParams: {
|
||||
sessionId: "session-1",
|
||||
sessionFile: "/tmp/session.jsonl",
|
||||
workspaceDir: "/tmp/workspace",
|
||||
prompt: "hello",
|
||||
provider: "claude-cli",
|
||||
model: "claude-opus-4-8",
|
||||
thinkLevel: "high",
|
||||
timeoutMs: 1_000,
|
||||
runId: "run-compaction-bridge",
|
||||
},
|
||||
});
|
||||
|
||||
expect(callbacks).toEqual(["start", "incomplete", "start", "end"]);
|
||||
});
|
||||
|
||||
it("bridges typed CLI plan events", async () => {
|
||||
cliDispatchState.runCliAgentMock.mockImplementationOnce(async (params: { runId: string }) => {
|
||||
emitAgentEvent({
|
||||
|
|
|
|||
|
|
@ -401,6 +401,8 @@ type RunCliAgentWithLifecycleParams = {
|
|||
onAssistantText?: (text: string) => Promise<boolean | void>;
|
||||
onReasoningText?: (payload: ReasoningTextPayload) => Promise<void>;
|
||||
onReasoningProgress?: (payload: ReasoningProgressPayload) => Promise<void>;
|
||||
onCompactionStart?: GetReplyOptions["onCompactionStart"];
|
||||
onCompactionEnd?: GetReplyOptions["onCompactionEnd"];
|
||||
onToolEvent?: (payload: CliToolEventPayload) => Promise<void>;
|
||||
onCommentaryText?: (payload: CommentaryTextPayload) => Promise<void>;
|
||||
onPlanUpdate?: GetReplyOptions["onPlanUpdate"];
|
||||
|
|
@ -539,6 +541,31 @@ async function runCliAgentWithLifecycleInternal(
|
|||
deliver: params.onReasoningProgress,
|
||||
startOrder: progressStartOrder,
|
||||
});
|
||||
const compactionBridge = createAgentEventBridge<
|
||||
{ phase: "start" } | { completed: boolean; phase: "end" }
|
||||
>({
|
||||
runId: params.runId,
|
||||
suppressed: params.suppressAssistantBridge,
|
||||
startOrder: progressStartOrder,
|
||||
deliver: async (event) => {
|
||||
if (event.phase === "start") {
|
||||
await params.onCompactionStart?.();
|
||||
} else {
|
||||
await params.onCompactionEnd?.({ completed: event.completed });
|
||||
}
|
||||
},
|
||||
read: (evt) => {
|
||||
if (evt.stream !== "compaction") {
|
||||
return undefined;
|
||||
}
|
||||
if (evt.data.phase === "start") {
|
||||
return { phase: "start" };
|
||||
}
|
||||
return evt.data.phase === "end"
|
||||
? { phase: "end", completed: evt.data.completed === true }
|
||||
: undefined;
|
||||
},
|
||||
});
|
||||
const toolBridge = createToolEventBridge({
|
||||
runId: params.runId,
|
||||
suppressed: params.suppressAssistantBridge,
|
||||
|
|
@ -567,6 +594,7 @@ async function runCliAgentWithLifecycleInternal(
|
|||
assistantBridge,
|
||||
reasoningBridge,
|
||||
reasoningProgressBridge,
|
||||
compactionBridge,
|
||||
toolBridge,
|
||||
commentaryBridge,
|
||||
planBridge,
|
||||
|
|
|
|||
|
|
@ -264,6 +264,7 @@ export function createAgentRunEventHandler(params: {
|
|||
return;
|
||||
}
|
||||
if (evt.data.completed !== true) {
|
||||
await params.turn.opts?.onCompactionEnd?.({ completed: false });
|
||||
await sendCompactionUserNotices("incomplete");
|
||||
return;
|
||||
}
|
||||
|
|
@ -288,7 +289,7 @@ export function createAgentRunEventHandler(params: {
|
|||
consoleMessage,
|
||||
});
|
||||
}
|
||||
await params.turn.opts?.onCompactionEnd?.();
|
||||
await params.turn.opts?.onCompactionEnd?.({ completed: true });
|
||||
await sendCompactionUserNotices("end");
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -75,7 +75,7 @@ describe("executeAgentTurn: compaction events", () => {
|
|||
|
||||
expect(result.kind).toBe("success");
|
||||
expect(onCompactionStart).toHaveBeenCalledTimes(1);
|
||||
expect(onCompactionEnd).toHaveBeenCalledTimes(1);
|
||||
expect(onCompactionEnd).toHaveBeenCalledWith({ completed: true });
|
||||
expect(onBlockReply).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
|
|
@ -432,6 +432,7 @@ describe("executeAgentTurn: compaction events", () => {
|
|||
|
||||
it("emits an incomplete compaction notice when compaction ends without completing", async () => {
|
||||
const onBlockReply = vi.fn();
|
||||
const onCompactionEnd = vi.fn();
|
||||
state.runEmbeddedAgentMock.mockImplementationOnce(async (params: EmbeddedAgentParams) => {
|
||||
await params.onAgentEvent?.({ stream: "compaction", data: { phase: "start" } });
|
||||
await params.onAgentEvent?.({
|
||||
|
|
@ -442,11 +443,12 @@ describe("executeAgentTurn: compaction events", () => {
|
|||
});
|
||||
|
||||
const result = await executeTestTurn(
|
||||
{ followupRun: createNotifyUserRun(), opts: { onBlockReply } },
|
||||
{ followupRun: createNotifyUserRun(), opts: { onBlockReply, onCompactionEnd } },
|
||||
{ commandBody: "hello" },
|
||||
);
|
||||
|
||||
expect(result.kind).toBe("success");
|
||||
expect(onCompactionEnd).toHaveBeenCalledWith({ completed: false });
|
||||
expectBlockReplyCall(onBlockReply, 0, {
|
||||
text: "🧹 Compacting context...",
|
||||
isCompactionNotice: true,
|
||||
|
|
|
|||
|
|
@ -277,7 +277,19 @@ export async function prepareDispatchExecution(state: ChooseDispatchRouteReadySt
|
|||
// Snapshot verbose progress visibility for this run: commentary
|
||||
// classification in the CLI runners is wired once at run start, so a
|
||||
// mid-run verbose toggle cannot move inter-tool commentary between lanes.
|
||||
const deliverStandaloneCommentaryProgress = shouldEmitVerboseProgress();
|
||||
const standaloneCommentaryProgressVisible = shouldEmitVerboseProgress();
|
||||
const resolveVerboseProgressVisibility = () =>
|
||||
standaloneCommentaryProgressVisible &&
|
||||
shouldSendVerboseProgressMessages() &&
|
||||
!shouldSuppressProgressDelivery();
|
||||
const { commentaryPayloadsEnabled, draftOwnsCommentaryProgress } =
|
||||
resolveTurnCommentaryProgressOwner({
|
||||
commentaryPayloadsEnabled: state.commentaryPayloadsEnabled,
|
||||
options: params.replyOptions,
|
||||
resolveVerboseProgressVisibility,
|
||||
});
|
||||
const deliverStandaloneCommentaryProgress =
|
||||
standaloneCommentaryProgressVisible && !draftOwnsCommentaryProgress;
|
||||
const itemEventForwardingOptions = {
|
||||
forwardWhenSourceDeliverySuppressed: true,
|
||||
requiresToolSummaryVisibility: true,
|
||||
|
|
@ -327,16 +339,6 @@ export async function prepareDispatchExecution(state: ChooseDispatchRouteReadySt
|
|||
return await forwardItemEvent?.(payload);
|
||||
}
|
||||
: undefined;
|
||||
const resolveVerboseProgressVisibility = () =>
|
||||
deliverStandaloneCommentaryProgress &&
|
||||
shouldSendVerboseProgressMessages() &&
|
||||
!shouldSuppressProgressDelivery();
|
||||
const { commentaryPayloadsEnabled } = resolveTurnCommentaryProgressOwner({
|
||||
commentaryPayloadsEnabled: state.commentaryPayloadsEnabled,
|
||||
options: params.replyOptions,
|
||||
resolveVerboseProgressVisibility,
|
||||
});
|
||||
|
||||
const replyResolver =
|
||||
params.replyResolver ??
|
||||
(
|
||||
|
|
|
|||
|
|
@ -98,6 +98,50 @@ describe("dispatchReplyFromConfig", () => {
|
|||
expect(activeDuringOffRun).toBe(false);
|
||||
});
|
||||
|
||||
it("keeps verbose commentary in the channel draft when that draft owns progress", async () => {
|
||||
setNoAbort();
|
||||
sessionStoreMocks.currentEntry = { verboseLevel: "on" };
|
||||
const dispatcher = createDispatcher();
|
||||
const onItemEvent = vi.fn();
|
||||
|
||||
await dispatchReplyFromConfig({
|
||||
ctx: buildTestCtx({
|
||||
Provider: "telegram",
|
||||
Surface: "telegram",
|
||||
ChatType: "group",
|
||||
From: "telegram:group:-100123",
|
||||
SessionKey: "agent:main:telegram:group:-100123",
|
||||
}),
|
||||
cfg: automaticGroupReplyConfig,
|
||||
dispatcher,
|
||||
replyResolver: async (_ctx, opts) => {
|
||||
expect(opts?.commentaryPayloadsEnabled).toBe(false);
|
||||
await opts?.onItemEvent?.({
|
||||
itemId: "commentary-draft-1",
|
||||
kind: "preamble",
|
||||
progressText: "Inspecting the dispatch path.",
|
||||
});
|
||||
return { text: "Done." } satisfies ReplyPayload;
|
||||
},
|
||||
replyOptions: {
|
||||
commentaryProgressEnabled: true,
|
||||
progressPreambleEnabled: true,
|
||||
commentaryPayloadsEnabled: true,
|
||||
shouldDeliverCommentaryPayloads: () => false,
|
||||
onItemEvent,
|
||||
},
|
||||
});
|
||||
|
||||
expect(onItemEvent).toHaveBeenCalledExactlyOnceWith({
|
||||
itemId: "commentary-draft-1",
|
||||
kind: "preamble",
|
||||
progressText: "Inspecting the dispatch path.",
|
||||
});
|
||||
expect(dispatcher.sendToolResult).not.toHaveBeenCalled();
|
||||
expect(dispatcher.sendBlockReply).not.toHaveBeenCalled();
|
||||
expect(dispatcher.sendFinalReply).toHaveBeenCalledExactlyOnceWith({ text: "Done." });
|
||||
});
|
||||
|
||||
it.each([
|
||||
{ surface: "slack", verboseLevel: "on", nextVerboseLevel: "off", durableCommentary: true },
|
||||
{ surface: "slack", verboseLevel: "off", nextVerboseLevel: "on", durableCommentary: false },
|
||||
|
|
|
|||
|
|
@ -240,10 +240,10 @@ export async function executeFollowupTurn(params: {
|
|||
)
|
||||
: undefined,
|
||||
onCompactionEnd: sourceOpts?.onCompactionEnd
|
||||
? () =>
|
||||
? (payload) =>
|
||||
enqueueProgressResult(async () =>
|
||||
progressAllowed()
|
||||
? (await settleProgressVisibilityCallbackResult(sourceOpts.onCompactionEnd!()))
|
||||
? (await settleProgressVisibilityCallbackResult(sourceOpts.onCompactionEnd!(payload)))
|
||||
.visible
|
||||
: false,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -365,7 +365,7 @@ export function createChannelProgressDraftCompositor(params: {
|
|||
|
||||
const noteProgress = async (
|
||||
line?: ChannelProgressDraftCompositorLine,
|
||||
options?: { toolName?: string; startImmediately?: boolean },
|
||||
options?: { toolName?: string; startImmediately?: boolean; flush?: boolean },
|
||||
) => {
|
||||
if (!params.active || finalReplyStarted || finalReplyDelivered) {
|
||||
return false;
|
||||
|
|
@ -425,7 +425,8 @@ export function createChannelProgressDraftCompositor(params: {
|
|||
return shouldStoreLine ? await publish() : false;
|
||||
}
|
||||
if (options?.startImmediately || params.shouldStartNow?.(line) || (summary && needsAttention)) {
|
||||
return await startAndRender(summary && needsAttention ? { flush: true } : undefined);
|
||||
const flush = options?.flush === true || (summary && needsAttention);
|
||||
return await startAndRender(flush ? { flush: true } : undefined);
|
||||
}
|
||||
const alreadyStarted = gate.hasStarted;
|
||||
const progressActive = await gate.noteWork();
|
||||
|
|
|
|||
26
src/plugin-sdk/cli-backend-compat.test.ts
Normal file
26
src/plugin-sdk/cli-backend-compat.test.ts
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
import type { CliBackendParsedJsonlEvent } from "openclaw/plugin-sdk/cli-backend";
|
||||
import { describe, expect, it } from "vitest";
|
||||
|
||||
function describeLegacyCliBackendEvent(event: CliBackendParsedJsonlEvent): string {
|
||||
switch (event.kind) {
|
||||
case "text":
|
||||
case "thinking":
|
||||
return event.text;
|
||||
case "toolStart":
|
||||
return event.name;
|
||||
case "toolResult":
|
||||
return event.toolCallId;
|
||||
case "result":
|
||||
return event.text ?? "result";
|
||||
case "sessionId":
|
||||
return event.sessionId;
|
||||
}
|
||||
const exhaustive: never = event;
|
||||
return exhaustive;
|
||||
}
|
||||
|
||||
describe("CLI backend Plugin SDK compatibility", () => {
|
||||
it("keeps existing parser events exhaustively matchable when lifecycle events are added", () => {
|
||||
expect(describeLegacyCliBackendEvent({ kind: "text", text: "ready" })).toBe("ready");
|
||||
});
|
||||
});
|
||||
|
|
@ -15,7 +15,9 @@ export type {
|
|||
CliBackendNativeToolMode,
|
||||
CliBackendParseJsonlEvent,
|
||||
CliBackendParseJsonlEventContext,
|
||||
CliBackendParseJsonlLifecycleEvent,
|
||||
CliBackendParsedJsonlEvent,
|
||||
CliBackendParsedJsonlLifecycleEvent,
|
||||
CliBackendPlugin,
|
||||
CliBackendPreparedExecution,
|
||||
CliBackendPromptContext,
|
||||
|
|
|
|||
|
|
@ -19,6 +19,7 @@ import type {
|
|||
} from "./qa-channel-protocol.js";
|
||||
|
||||
type QaRunnerTransportPolicy = {
|
||||
directMessageOnly?: true;
|
||||
requireGroupMention?: true;
|
||||
senderAllowlist?: readonly string[];
|
||||
topLevelReplies?: true;
|
||||
|
|
|
|||
|
|
@ -334,6 +334,19 @@ export type CliBackendParseJsonlEvent = (
|
|||
ctx: CliBackendParseJsonlEventContext,
|
||||
) => CliBackendParsedJsonlEvent | readonly CliBackendParsedJsonlEvent[] | null | undefined;
|
||||
|
||||
export type CliBackendParsedJsonlLifecycleEvent =
|
||||
| { kind: "compaction"; phase: "start" }
|
||||
| { kind: "compaction"; phase: "end"; completed: boolean };
|
||||
|
||||
export type CliBackendParseJsonlLifecycleEvent = (
|
||||
line: string,
|
||||
ctx: CliBackendParseJsonlEventContext,
|
||||
) =>
|
||||
| CliBackendParsedJsonlLifecycleEvent
|
||||
| readonly CliBackendParsedJsonlLifecycleEvent[]
|
||||
| null
|
||||
| undefined;
|
||||
|
||||
export type CliBackendAuthEpochMode = "combined" | "profile-only";
|
||||
|
||||
export type CliBackendNativeToolMode = "none" | "always-on" | "selectable";
|
||||
|
|
@ -519,6 +532,11 @@ type CliBackendPluginBase = {
|
|||
* renders them but does not treat them as host tool execution or delivery evidence.
|
||||
*/
|
||||
parseJsonlEvent?: CliBackendParseJsonlEvent;
|
||||
/**
|
||||
* Optional lifecycle parser kept separate from the legacy JSONL event union.
|
||||
* Existing plugins can continue exhaustively matching `parseJsonlEvent` results.
|
||||
*/
|
||||
parseJsonlLifecycleEvent?: CliBackendParseJsonlLifecycleEvent;
|
||||
/**
|
||||
* Whether this CLI backend can expose native tools outside OpenClaw's tool
|
||||
* catalog. Exact restricted runs require `selectable` plus a declared
|
||||
|
|
|
|||
|
|
@ -61,8 +61,8 @@ it.skipIf(process.platform === "win32").each([
|
|||
{ mode: "qa-branch", code: 0, reason: "release-branch-head" },
|
||||
{ mode: "qa-mismatch", code: 1, reason: "" },
|
||||
{ mode: "qa-untrusted", code: 1, reason: "" },
|
||||
{ mode: "qa-pr", code: 0, reason: "open-pr-head" },
|
||||
{ mode: "qa-api-error", code: 23, reason: "" },
|
||||
{ mode: "qa-pr", code: 1, reason: "" },
|
||||
{ mode: "qa-api-error", code: 1, reason: "" },
|
||||
{ mode: "qa-foreign-origin", code: 125, reason: "" },
|
||||
])(
|
||||
"QA selected-ref validation owns authenticated fetches without changing trust ($mode)",
|
||||
|
|
|
|||
|
|
@ -7669,15 +7669,20 @@ printf '%s\\n' "$DEEPSEEK_API_KEY" "$DEEPINFRA_API_KEY"`,
|
|||
for (const lane of ["mock_parity", "buzz", "telegram", "discord", "whatsapp", "slack"]) {
|
||||
expect(releaseJob.with?.[`run_${lane}`]).toBeUndefined();
|
||||
}
|
||||
const manualScenarioGuard =
|
||||
"(github.event_name != 'workflow_dispatch' || inputs.scenario == '')";
|
||||
expect(workflowJob(QA_LIVE_TRANSPORTS_WORKFLOW, "run_mock_parity").if).toBe(
|
||||
"inputs.expected_sha == '' || inputs.run_mock_parity",
|
||||
`(inputs.expected_sha == '' || inputs.run_mock_parity) && ${manualScenarioGuard}`,
|
||||
);
|
||||
expect(workflowJob(QA_LIVE_TRANSPORTS_WORKFLOW, "run_live_matrix").if).toBe(
|
||||
"inputs.expected_sha == '' || inputs.run_matrix",
|
||||
`(inputs.expected_sha == '' || inputs.run_matrix) && ${manualScenarioGuard}`,
|
||||
);
|
||||
for (const channel of ["telegram", "discord", "whatsapp", "slack"]) {
|
||||
expect(workflowJob(QA_LIVE_TRANSPORTS_WORKFLOW, "run_live_telegram").if).toBe(
|
||||
"inputs.expected_sha == '' || inputs.run_telegram",
|
||||
);
|
||||
for (const channel of ["discord", "whatsapp", "slack"]) {
|
||||
expect(workflowJob(QA_LIVE_TRANSPORTS_WORKFLOW, `run_live_${channel}`).if).toBe(
|
||||
`inputs.expected_sha == '' || inputs.run_${channel}`,
|
||||
`(inputs.expected_sha == '' || inputs.run_${channel}) && ${manualScenarioGuard}`,
|
||||
);
|
||||
}
|
||||
expect(releaseWorkflow).not.toContain("qa_live_matrix_release_checks");
|
||||
|
|
|
|||
|
|
@ -320,7 +320,11 @@ suite.define(() => {
|
|||
const messages: unknown[] = [];
|
||||
const appWindow = window as Window & {
|
||||
openclawNativeLinkMessages?: unknown[];
|
||||
webkit?: unknown;
|
||||
webkit?: {
|
||||
messageHandlers?: {
|
||||
openclawLink?: { postMessage: (message: unknown) => void };
|
||||
};
|
||||
};
|
||||
};
|
||||
appWindow.openclawNativeLinkMessages = messages;
|
||||
appWindow.webkit = {
|
||||
|
|
|
|||
|
|
@ -75,7 +75,11 @@ describeControlUiE2e("native link routing", () => {
|
|||
const messages: unknown[] = [];
|
||||
const host = window as Window & {
|
||||
openclawNativeLinkMessages?: unknown[];
|
||||
webkit?: unknown;
|
||||
webkit?: {
|
||||
messageHandlers?: {
|
||||
openclawLink?: { postMessage: (message: unknown) => void };
|
||||
};
|
||||
};
|
||||
};
|
||||
host.openclawNativeLinkMessages = messages;
|
||||
host.webkit = {
|
||||
|
|
|
|||
|
|
@ -122,7 +122,6 @@ describe("AppSidebar catalog row lifecycle", () => {
|
|||
|
||||
const updatedLabel = sidebar.querySelector<HTMLElement>(labelSelector);
|
||||
expect(updatedLabel?.textContent).toBe("Short");
|
||||
expect(updatedLabel).not.toBe(oldLabel);
|
||||
expect(updatedLabel?.classList.contains("hover-marquee--scrolling")).toBe(false);
|
||||
expect(updatedLabel?.style.getPropertyValue("--hover-marquee-shift")).toBe("");
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue