From fc269d2da862d8f577bef71c520dd98004924b1e Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 21:00:50 -0700 Subject: [PATCH 1/9] fix(t3-agent): normalize imported Hermes history --- .../components/chat/MessagesTimeline.test.tsx | 2 +- .../src/components/chat/MessagesTimeline.tsx | 2 +- apps/web/src/hermesLineage.ts | 17 +-- integrations/hermes/t3agent/adapter.py | 55 +++++++- .../hermes/t3agent/tests/test_adapter.py | 60 +++++++-- packages/client-runtime/package.json | 4 + .../client-runtime/src/state/entities.test.ts | 100 +++++++++++++++ .../src/state/hermesImportedHistory.test.ts | 119 ++++++++++++++++++ .../src/state/hermesImportedHistory.ts | 114 +++++++++++++++++ .../client-runtime/src/state/threadDetail.ts | 9 +- 10 files changed, 457 insertions(+), 25 deletions(-) create mode 100644 packages/client-runtime/src/state/hermesImportedHistory.test.ts create mode 100644 packages/client-runtime/src/state/hermesImportedHistory.ts diff --git a/apps/web/src/components/chat/MessagesTimeline.test.tsx b/apps/web/src/components/chat/MessagesTimeline.test.tsx index 52544d07191..55ee9937e43 100644 --- a/apps/web/src/components/chat/MessagesTimeline.test.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.test.tsx @@ -441,7 +441,7 @@ describe("MessagesTimeline", () => { expect(markup).not.toContain('data-maintain-scroll-at-end="enabled"'); expect(markup).toContain('data-maintain-visible-content-position="object"'); expect(markup).toContain('data-maintain-visible-content-position-data="true"'); - expect(markup).toContain('data-maintain-visible-content-position-size="false"'); + expect(markup).toContain('data-maintain-visible-content-position-size="true"'); expect(onAnchorReady).toHaveBeenCalledOnce(); expect(onAnchorReady).toHaveBeenCalledWith(secondEntry.message.id, 1); expect(onAnchorSizeChanged).toHaveBeenCalledWith(secondEntry.message.id, 240); diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 81ba744178c..075af209d2c 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -520,7 +520,7 @@ export const MessagesTimeline = memo(function MessagesTimeline({ } maintainVisibleContentPosition={{ data: true, - size: false, + size: true, }} onScroll={handleScroll} className={cn( diff --git a/apps/web/src/hermesLineage.ts b/apps/web/src/hermesLineage.ts index b1c2eb0052a..c3c8206f02d 100644 --- a/apps/web/src/hermesLineage.ts +++ b/apps/web/src/hermesLineage.ts @@ -1,15 +1,8 @@ -import { HermesLineageMetadata, type HermesLineageMetadata as Metadata } from "@t3tools/contracts"; -import * as Option from "effect/Option"; -import * as Schema from "effect/Schema"; - -export const HERMES_LINEAGE_PREFIX = "t3agent-lineage:"; - -const decodeLineage = Schema.decodeUnknownOption(Schema.fromJsonString(HermesLineageMetadata)); - -export function parseHermesLineageMessage(text: string): Metadata | null { - if (!text.startsWith(HERMES_LINEAGE_PREFIX)) return null; - return Option.getOrNull(decodeLineage(text.slice(HERMES_LINEAGE_PREFIX.length))); -} +import type { HermesLineageMetadata as Metadata } from "@t3tools/contracts"; +export { + HERMES_LINEAGE_PREFIX, + parseHermesLineageMessage, +} from "@t3tools/client-runtime/state/hermes-imported-history"; export function formatHermesSourceLabel(source: string): string { const normalized = source.trim().toLocaleLowerCase(); diff --git a/integrations/hermes/t3agent/adapter.py b/integrations/hermes/t3agent/adapter.py index 6f5b4d7fd90..562ded06a22 100644 --- a/integrations/hermes/t3agent/adapter.py +++ b/integrations/hermes/t3agent/adapter.py @@ -72,6 +72,23 @@ _OUTBOX_PATH_ENV = "T3_AGENT_OUTBOX_PATH" _INGRESS_LEDGER_PATH_ENV = "T3_AGENT_INGRESS_LEDGER_PATH" +_DISCORD_TRIGGERING_MESSAGE_PREFIX = re.compile( + r"^\[Triggering message id:[^\r\n]*\]\s*", + re.IGNORECASE, +) +_DISCORD_TEXT_DOCUMENT_PREFIX = re.compile( + r"^\[The user sent a text document: '([^'\r\n]+)'\. " + r"Its content has been included below\. The file is also saved at: " + r"[^\]\r\n]+\]\s*", + re.IGNORECASE, +) +_DISCORD_SENDER_PREFIX = re.compile(r"^\[([^\]\r\n:]{1,80})\]\s+") +_DISCORD_NON_SENDER_LABEL = re.compile(r"^(?:async\b|the user\b)", re.IGNORECASE) +_HERMES_ASYNC_DELEGATION_RESULT_PREFIX = re.compile( + r"^\[ASYNC DELEGATION BATCH COMPLETE\b", + re.IGNORECASE, +) + def _env_or_extra(config: PlatformConfig, env_name: str, key: str, default: Any = "") -> Any: value = os.getenv(env_name) @@ -225,6 +242,31 @@ def _gateway_session_resources(gateway_runner: Any) -> Tuple[Any, Any, Any]: return session_db, async_store, db +def _is_synthetic_history_user(content: Any) -> bool: + return isinstance(content, str) and bool( + _HERMES_ASYNC_DELEGATION_RESULT_PREFIX.match(content) + ) + + +def _normalize_history_user_content(source: str, content: str) -> str: + if source.strip().casefold() != "discord": + return content + + body = _DISCORD_TRIGGERING_MESSAGE_PREFIX.sub("", content, count=1) + document_match = _DISCORD_TEXT_DOCUMENT_PREFIX.match(body) + attachment_label = "" + if document_match is not None: + attachment_label = f"**Attached:** {document_match.group(1)}\n\n" + body = body[document_match.end() :] + + sender_match = _DISCORD_SENDER_PREFIX.match(body) + if sender_match is not None and not _DISCORD_NON_SENDER_LABEL.match( + sender_match.group(1) + ): + body = body[sender_match.end() :] + return f"{attachment_label}{body}" + + def _stored_model_config(session: Dict[str, Any]) -> Dict[str, Any]: raw_config = session.get("model_config") if isinstance(raw_config, dict): @@ -1256,7 +1298,9 @@ async def fork_once() -> Tuple[int, Dict[str, Any]]: user_turns = 0 for message in messages: role = str(message.get("role") or "") - if role == "user": + if role == "user" and not _is_synthetic_history_user( + message.get("content") + ): user_turns += 1 if raw_turn_count is not None and user_turns > raw_turn_count: break @@ -1342,6 +1386,7 @@ async def fork_once() -> Tuple[int, Dict[str, Any]]: set_reasoning(target_session_key, reasoning_config) history: List[Dict[str, Any]] = [] + source = str(source_session.get("source") or "unknown") for message in copied_messages: role = str(message.get("role") or "") content = message.get("content") @@ -1349,11 +1394,15 @@ async def fork_once() -> Tuple[int, Dict[str, Any]]: content, str ): continue + if role == "assistant" and not content.strip(): + continue + if role == "user" and _is_synthetic_history_user(content): + continue history.append( { "role": role, "content": ( - re.sub(r"^\[[^\]\n]{1,80}\]\s+", "", content) + _normalize_history_user_content(source, content) if role == "user" else content ), @@ -1366,7 +1415,7 @@ async def fork_once() -> Tuple[int, Dict[str, Any]]: "sourceSessionId": source_session_id, "childSessionId": child_session_id, "targetThreadId": target_thread_id, - "source": str(source_session.get("source") or "unknown"), + "source": source, "title": child_title, "messages": history, **( diff --git a/integrations/hermes/t3agent/tests/test_adapter.py b/integrations/hermes/t3agent/tests/test_adapter.py index 3f06d70bdb1..03dc9fb27d3 100644 --- a/integrations/hermes/t3agent/tests/test_adapter.py +++ b/integrations/hermes/t3agent/tests/test_adapter.py @@ -2,8 +2,9 @@ import asyncio import json +import os from types import SimpleNamespace -from typing import Any, Dict, List +from typing import Any, Dict, Iterator, List from aiohttp import ClientSession, web from aiohttp.test_utils import TestClient, TestServer @@ -15,6 +16,25 @@ from integrations.hermes.t3agent import adapter as adapter_module +@pytest.fixture(autouse=True) +def isolate_t3agent_environment() -> Iterator[None]: + """Keep lazy Hermes imports from leaking live bridge config between tests.""" + original = { + name: value + for name, value in os.environ.items() + if name.startswith("T3_AGENT_") + } + for name in original: + os.environ.pop(name, None) + try: + yield + finally: + for name in tuple(os.environ): + if name.startswith("T3_AGENT_"): + os.environ.pop(name, None) + os.environ.update(original) + + def make_config(**extra: Any) -> SimpleNamespace: defaults = { "instance_id": "hermes-test", @@ -348,10 +368,34 @@ def get_session(self, session_id: str) -> Any: def get_messages(self, session_id: str) -> List[Dict[str, Any]]: assert session_id == "discord-source" return [ - {"role": "user", "content": "[Ada] first", "timestamp": 10}, - {"role": "assistant", "content": "first answer", "timestamp": 11}, - {"role": "user", "content": "[Ada] second", "timestamp": 12}, - {"role": "assistant", "content": "second answer", "timestamp": 13}, + { + "role": "user", + "content": ( + "[Triggering message id: `123` — use as `message_id` for " + "reply/react/pin via the discord tools.]\n\n" + "[The user sent a text document: 'notes.txt'. Its content " + "has been included below. The file is also saved at: " + "/tmp/notes.txt]\n\n[Ada] first" + ), + "timestamp": 10, + }, + { + "role": "assistant", + "content": "", + "tool_calls": [{"name": "search"}], + "timestamp": 11, + }, + {"role": "assistant", "content": "first answer", "timestamp": 12}, + { + "role": "user", + "content": ( + "[ASYNC DELEGATION BATCH COMPLETE — deleg_123]\n" + "Internal tool result" + ), + "timestamp": 13, + }, + {"role": "user", "content": "[Ada] second", "timestamp": 14}, + {"role": "assistant", "content": "second answer", "timestamp": 15}, ] def get_next_title_in_lineage(self, title: str) -> str: @@ -424,7 +468,7 @@ def set_reasoning_override(session_key: str, config: Dict[str, Any]) -> None: assert payload["targetThreadId"] == target_thread_id assert payload["title"] == "Planning #2" assert [message["content"] for message in payload["messages"]] == [ - "first", + "**Attached:** notes.txt\n\nfirst", "first answer", ] assert created[0]["source"] == "t3agent" @@ -435,7 +479,9 @@ def set_reasoning_override(session_key: str, config: Dict[str, Any]) -> None: assert created[0]["model_config"]["_t3agent_imported_from"] == ( "discord-source" ) - assert replacements[0][1][-1]["content"] == "first answer" + assert replacements[0][1][-1]["content"].startswith( + "[ASYNC DELEGATION BATCH COMPLETE" + ) assert titles[0][1] == "Planning #2" assert route_sources[0].thread_id == target_thread_id assert switches[0][1] == payload["childSessionId"] diff --git a/packages/client-runtime/package.json b/packages/client-runtime/package.json index 4fa05f850e5..ec287d1079c 100644 --- a/packages/client-runtime/package.json +++ b/packages/client-runtime/package.json @@ -63,6 +63,10 @@ "types": "./src/state/git.ts", "default": "./src/state/git.ts" }, + "./state/hermes-imported-history": { + "types": "./src/state/hermesImportedHistory.ts", + "default": "./src/state/hermesImportedHistory.ts" + }, "./state/models": { "types": "./src/state/models.ts", "default": "./src/state/models.ts" diff --git a/packages/client-runtime/src/state/entities.test.ts b/packages/client-runtime/src/state/entities.test.ts index e08fd9e552f..e25c59b9ea1 100644 --- a/packages/client-runtime/src/state/entities.test.ts +++ b/packages/client-runtime/src/state/entities.test.ts @@ -1,5 +1,6 @@ import { EnvironmentId, + MessageId, ProjectId, ProviderInstanceId, ThreadId, @@ -367,4 +368,103 @@ describe("environment entity projections", () => { expect(harness.registry.get(messagesAtom)).toBe(messages); expect(harness.registry.get(activitiesAtom)).toBe(activities); }); + + it("normalizes inherited Discord history before exposing thread messages", () => { + const harness = makeHarness(); + const threadRef = { + environmentId: ENVIRONMENT_ID, + threadId: THREAD_ID, + }; + const messagesAtom = harness.threadDetails.messagesAtom(threadRef); + const createdAt = "2026-07-22T18:36:52.404Z"; + const lineageText = + 't3agent-lineage:{"kind":"import","label":"Imported from discord","sourceProvider":"discord","sourceSessionId":"discord-source"}'; + const detail = { + ...THREAD_SHELL, + deletedAt: null, + messages: [ + { + id: MessageId.make("discord-user"), + role: "user", + text: "[Triggering message id: `1529557863226150933` — use as `message_id` for reply/react/pin via the discord tools.]\n\n[Parker] Keep this strategy thread focused.", + turnId: null, + streaming: false, + createdAt, + updatedAt: createdAt, + }, + { + id: MessageId.make("discord-tool-placeholder"), + role: "assistant", + text: "", + turnId: null, + streaming: false, + createdAt, + updatedAt: createdAt, + }, + { + id: MessageId.make("discord-answer"), + role: "assistant", + text: "Understood.", + turnId: null, + streaming: false, + createdAt, + updatedAt: createdAt, + }, + { + id: MessageId.make("discord-lineage"), + role: "system", + text: lineageText, + turnId: null, + streaming: false, + createdAt, + updatedAt: createdAt, + }, + { + id: MessageId.make("native-follow-up"), + role: "user", + text: "[Parker] This native follow-up must remain untouched.", + turnId: null, + streaming: false, + createdAt, + updatedAt: createdAt, + }, + ], + proposedPlans: [], + activities: [], + checkpoints: [], + } satisfies OrchestrationThread; + + harness.registry.set( + harness.threadStateAtom(THREAD_ID), + AsyncResult.success({ + data: Option.some(detail), + status: "live", + error: Option.none(), + }), + ); + + expect( + harness.registry.get(messagesAtom).map((message) => ({ + id: message.id, + text: message.text, + })), + ).toEqual([ + { + id: MessageId.make("discord-user"), + text: "Keep this strategy thread focused.", + }, + { + id: MessageId.make("discord-answer"), + text: "Understood.", + }, + { + id: MessageId.make("discord-lineage"), + text: lineageText, + }, + { + id: MessageId.make("native-follow-up"), + text: "[Parker] This native follow-up must remain untouched.", + }, + ]); + }); }); diff --git a/packages/client-runtime/src/state/hermesImportedHistory.test.ts b/packages/client-runtime/src/state/hermesImportedHistory.test.ts new file mode 100644 index 00000000000..2945cf805c2 --- /dev/null +++ b/packages/client-runtime/src/state/hermesImportedHistory.test.ts @@ -0,0 +1,119 @@ +import { MessageId, type OrchestrationMessage } from "@t3tools/contracts"; +import { describe, expect, it } from "@effect/vitest"; + +import { HERMES_LINEAGE_PREFIX, normalizeImportedHermesHistory } from "./hermesImportedHistory.ts"; + +const CREATED_AT = "2026-07-06T17:52:27.000Z"; + +function message( + id: string, + role: OrchestrationMessage["role"], + text: string, + overrides: Partial = {}, +): OrchestrationMessage { + return { + id: MessageId.make(id), + role, + text, + turnId: null, + streaming: false, + createdAt: CREATED_AT, + updatedAt: CREATED_AT, + ...overrides, + }; +} + +function lineage(sourceProvider: string): OrchestrationMessage { + return message( + `${sourceProvider}-lineage`, + "system", + `${HERMES_LINEAGE_PREFIX}${JSON.stringify({ + kind: "import", + label: `Imported from ${sourceProvider}`, + sourceProvider, + sourceSessionId: `${sourceProvider}-source`, + })}`, + ); +} + +describe("normalizeImportedHermesHistory", () => { + it("drops Telegram tool placeholders without changing user content", () => { + const messages = [ + message("telegram-user", "user", "Research this event."), + message("telegram-tool", "assistant", ""), + message("telegram-answer", "assistant", "Here is the result."), + lineage("telegram"), + ]; + + const normalized = normalizeImportedHermesHistory(messages); + + expect(normalized.map(({ id, text }) => ({ id, text }))).toEqual([ + { id: MessageId.make("telegram-user"), text: "Research this event." }, + { id: MessageId.make("telegram-answer"), text: "Here is the result." }, + { id: MessageId.make("telegram-lineage"), text: messages[3]?.text }, + ]); + }); + + it("preserves empty assistant rows with attachments and rows after the import boundary", () => { + const inheritedAttachment = message("inherited-image", "assistant", "", { + attachments: [ + { + type: "image", + id: "image-1", + name: "chart.png", + mimeType: "image/png", + sizeBytes: 12, + }, + ], + }); + const nativeEmptyAssistant = message("native-empty", "assistant", ""); + const messages = [ + inheritedAttachment, + message("inherited-tool", "assistant", ""), + lineage("discord"), + nativeEmptyAssistant, + ]; + + const normalized = normalizeImportedHermesHistory(messages); + + expect(normalized).toEqual([inheritedAttachment, messages[2], nativeEmptyAssistant]); + }); + + it("simplifies Discord text-document wrappers and removes synthetic delegation results", () => { + const messages = [ + message( + "discord-document", + "user", + "[Triggering message id: `123` — use as `message_id` for reply/react/pin via the discord tools.]\n\n[The user sent a text document: 'copied-text.txt'. Its content has been included below. The file is also saved at: /home/hermes/.hermes/cache/documents/doc_123_copied-text.txt]\n\n[Parker] [Content of copied-text.txt]:\nUseful notes", + ), + message( + "delegation-result", + "user", + "[ASYNC DELEGATION BATCH COMPLETE — deleg_123]\nInternal tool result", + ), + lineage("discord"), + ]; + + expect(normalizeImportedHermesHistory(messages).map(({ id, text }) => ({ id, text }))).toEqual([ + { + id: MessageId.make("discord-document"), + text: "**Attached:** copied-text.txt\n\n[Content of copied-text.txt]:\nUseful notes", + }, + { + id: MessageId.make("discord-lineage"), + text: messages[2]?.text, + }, + ]); + }); + + it("returns ordinary and malformed-lineage message collections unchanged", () => { + const ordinary = [message("ordinary-user", "user", "[Parker] Keep this text.")]; + const malformed = [ + ...ordinary, + message("bad-lineage", "system", `${HERMES_LINEAGE_PREFIX}{not-json`), + ]; + + expect(normalizeImportedHermesHistory(ordinary)).toBe(ordinary); + expect(normalizeImportedHermesHistory(malformed)).toBe(malformed); + }); +}); diff --git a/packages/client-runtime/src/state/hermesImportedHistory.ts b/packages/client-runtime/src/state/hermesImportedHistory.ts new file mode 100644 index 00000000000..89628122150 --- /dev/null +++ b/packages/client-runtime/src/state/hermesImportedHistory.ts @@ -0,0 +1,114 @@ +import { + HermesLineageMetadata, + type HermesLineageMetadata as LineageMetadata, + type OrchestrationMessage, +} from "@t3tools/contracts"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; + +export const HERMES_LINEAGE_PREFIX = "t3agent-lineage:"; + +const decodeLineage = Schema.decodeUnknownOption(Schema.fromJsonString(HermesLineageMetadata)); +const discordTriggeringMessagePrefix = /^\[Triggering message id:[^\r\n]*\]\s*/iu; +const discordTextDocumentPrefix = + /^\[The user sent a text document: '([^'\r\n]+)'\. Its content has been included below\. The file is also saved at: [^\]\r\n]+\]\s*/iu; +const discordSenderPrefix = /^\[([^\]\r\n:]{1,80})\]\s+/u; +const discordNonSenderLabels = /^(?:async\b|the user\b)/iu; +const hermesAsyncDelegationResultPrefix = /^\[ASYNC DELEGATION BATCH COMPLETE\b/iu; + +export function parseHermesLineageMessage(text: string): LineageMetadata | null { + if (!text.startsWith(HERMES_LINEAGE_PREFIX)) { + return null; + } + return Option.getOrNull(decodeLineage(text.slice(HERMES_LINEAGE_PREFIX.length))); +} + +function normalizeImportedUserText(sourceProvider: string, text: string): string { + if (sourceProvider.trim().toLocaleLowerCase() !== "discord") { + return text; + } + + const withoutTrigger = text.replace(discordTriggeringMessagePrefix, ""); + const documentMatch = discordTextDocumentPrefix.exec(withoutTrigger); + const attachmentLabel = documentMatch?.[1] ? `**Attached:** ${documentMatch[1]}\n\n` : ""; + const body = documentMatch ? withoutTrigger.slice(documentMatch[0].length) : withoutTrigger; + const senderMatch = discordSenderPrefix.exec(body); + if (!senderMatch || discordNonSenderLabels.test(senderMatch[1] ?? "")) { + return `${attachmentLabel}${body}`; + } + return `${attachmentLabel}${body.slice(senderMatch[0].length)}`; +} + +function isEmptyAssistantPlaceholder(message: OrchestrationMessage): boolean { + return ( + message.role === "assistant" && + !message.streaming && + message.text.trim().length === 0 && + (message.attachments?.length ?? 0) === 0 + ); +} + +/** + * Imported Hermes history is copied verbatim into the child session for model + * continuity, but its display projection still contains gateway envelopes and + * contentless assistant rows representing tool calls. Normalize only the + * inherited prefix before the lineage marker so later T3 Agent turns retain + * their native content and empty-message semantics. + */ +export function normalizeImportedHermesHistory( + messages: ReadonlyArray, +): ReadonlyArray { + const lineageIndex = messages.findIndex((message) => { + if (message.role !== "system") { + return false; + } + return parseHermesLineageMessage(message.text)?.kind === "import"; + }); + if (lineageIndex < 0) { + return messages; + } + + const lineageMessage = messages[lineageIndex]; + if (!lineageMessage) { + return messages; + } + const lineage = parseHermesLineageMessage(lineageMessage.text); + if (lineage === null || lineage.kind !== "import") { + return messages; + } + + let changed = false; + const normalized: Array = []; + for (let index = 0; index < messages.length; index += 1) { + const message = messages[index]; + if (!message) { + continue; + } + if (index >= lineageIndex) { + normalized.push(message); + continue; + } + if (isEmptyAssistantPlaceholder(message)) { + changed = true; + continue; + } + if (message.role === "user" && hermesAsyncDelegationResultPrefix.test(message.text)) { + changed = true; + continue; + } + if (message.role !== "user") { + normalized.push(message); + continue; + } + + const text = normalizeImportedUserText(lineage.sourceProvider, message.text); + if (text === message.text) { + normalized.push(message); + continue; + } + changed = true; + normalized.push({ ...message, text }); + } + + return changed ? normalized : messages; +} diff --git a/packages/client-runtime/src/state/threadDetail.ts b/packages/client-runtime/src/state/threadDetail.ts index 770738ad9bb..4cd06af3b1b 100644 --- a/packages/client-runtime/src/state/threadDetail.ts +++ b/packages/client-runtime/src/state/threadDetail.ts @@ -13,6 +13,7 @@ import { AsyncResult, Atom } from "effect/unstable/reactivity"; import type { EnvironmentThread, EnvironmentThreadShell } from "./models.ts"; import { scopeThread } from "./models.ts"; +import { normalizeImportedHermesHistory } from "./hermesImportedHistory.ts"; import { EMPTY_ENVIRONMENT_THREAD_STATE, type EnvironmentThreadState } from "./threadState.ts"; import { parseThreadKey, threadKey } from "./entities.ts"; import { THREAD_STATE_IDLE_TTL_MS } from "./threadRetention.ts"; @@ -92,7 +93,13 @@ export function createEnvironmentThreadDetailAtoms( return previousValue; } previousSource = source; - previousValue = source === null ? null : scopeThread(ref.environmentId, source); + previousValue = + source === null + ? null + : scopeThread(ref.environmentId, { + ...source, + messages: normalizeImportedHermesHistory(source.messages), + }); return previousValue; }).pipe( Atom.setIdleTTL(THREAD_STATE_IDLE_TTL_MS), From be431693b64ebe1a7e4df821f3d9e7606b7b7b84 Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 21:37:31 -0700 Subject: [PATCH 2/9] Render imported Hermes tool history as native turns --- apps/server/src/orchestration/decider.ts | 2 +- .../hermes/HermesBridgeClient.test.ts | 25 ++ .../HermesConversationLifecycle.test.ts | 43 ++- .../hermes/HermesConversationLifecycle.ts | 55 ++- integrations/hermes/t3agent/adapter.py | 346 ++++++++++++++++-- .../hermes/t3agent/tests/test_adapter.py | 69 +++- packages/contracts/src/hermesBridge.ts | 4 + packages/contracts/src/orchestration.ts | 1 + 8 files changed, 504 insertions(+), 41 deletions(-) diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 767c713bd6d..356251850bb 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -921,7 +921,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" messageId: command.messageId, role: command.role, text: command.text, - turnId: null, + turnId: command.turnId ?? null, streaming: false, imported: true, createdAt: command.createdAt, diff --git a/apps/server/src/provider/hermes/HermesBridgeClient.test.ts b/apps/server/src/provider/hermes/HermesBridgeClient.test.ts index 6663741c86a..24088e0cc0f 100644 --- a/apps/server/src/provider/hermes/HermesBridgeClient.test.ts +++ b/apps/server/src/provider/hermes/HermesBridgeClient.test.ts @@ -132,6 +132,28 @@ describe("HermesBridgeClient", () => { content: "Continue this", createdAt: "2026-07-23T10:00:00.000Z", }, + { + role: "assistant", + content: "Done", + createdAt: "2026-07-23T10:00:02.000Z", + turnId: "hermes-turn-1", + }, + ], + activities: [ + { + id: "hermes-activity-1", + tone: "tool", + kind: "tool.completed", + summary: "Searched files", + payload: { + itemType: "mcp_tool_call", + status: "completed", + data: { toolCallId: "call-1" }, + }, + turnId: "hermes-turn-1", + sequence: 0, + createdAt: "2026-07-23T10:00:01.000Z", + }, ], }), ); @@ -148,6 +170,8 @@ describe("HermesBridgeClient", () => { }); assert.strictEqual(result.childSessionId, "t3-child"); + assert.strictEqual(result.messages[1]?.turnId, "hermes-turn-1"); + assert.strictEqual(result.activities?.[0]?.kind, "tool.completed"); const call = execute.mock.calls[0]; assert.ok(call); const [request] = call; @@ -223,6 +247,7 @@ describe("HermesBridgeClient", () => { protocolVersion: 1, requestId: "request-interrupt", type: "turn.interrupt", + chatId: "t3agent", threadId: "thread-1", }), ], diff --git a/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts b/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts index 411dbfb37f5..5cf49a62a8f 100644 --- a/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts +++ b/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts @@ -1,11 +1,13 @@ import { assert, describe, it, vi } from "@effect/vitest"; import { HERMES_BRIDGE_PROTOCOL_VERSION, + EventId, HermesBridgeRequestId, HermesBridgeSessionId, ProjectId, ProviderInstanceId, ThreadId, + TurnId, type HermesBridgeSessionListResponse, type OrchestrationCommand, } from "@t3tools/contracts"; @@ -27,6 +29,7 @@ const CHILD_SESSION_ID = HermesBridgeSessionId.make("t3-child"); const SOURCE_THREAD_ID = ThreadId.make("source-thread"); const LIVE_IMPORTED_THREAD_ID = ThreadId.make("live-import"); const DELETED_IMPORTED_THREAD_ID = ThreadId.make("deleted-import"); +const IMPORTED_TURN_ID = TurnId.make("hermes-turn-1"); const sessionList = { protocolVersion: HERMES_BRIDGE_PROTOCOL_VERSION, @@ -104,6 +107,23 @@ function makeHarness(input: { role: "assistant", content: "Continuing", createdAt: NOW, + turnId: IMPORTED_TURN_ID, + }, + ], + activities: [ + { + id: EventId.make("hermes-activity-1"), + tone: "tool", + kind: "tool.completed", + summary: "Searched files", + payload: { + itemType: "mcp_tool_call", + status: "completed", + data: { toolCallId: "call-1" }, + }, + turnId: IMPORTED_TURN_ID, + sequence: 0, + createdAt: NOW, }, ], }), @@ -232,7 +252,7 @@ describe("HermesConversationLifecycle", () => { }); }); - it.effect("reuses an existing live import unless another copy is requested", () => { + it.effect("refreshes an existing live import in place unless another copy is requested", () => { const { lifecycle, commands, forkSession } = makeHarness({}); return Effect.gen(function* () { const result = yield* lifecycle.forkConversation({ @@ -243,8 +263,17 @@ describe("HermesConversationLifecycle", () => { threadId: LIVE_IMPORTED_THREAD_ID, existing: true, }); - assert.lengthOf(commands, 0); - assert.strictEqual(forkSession.mock.calls.length, 0); + assert.deepStrictEqual( + commands.map((command) => command.type), + ["thread.message.import", "thread.message.import", "thread.activity.append"], + ); + const assistantImport = commands[1]; + assert.strictEqual(assistantImport?.type, "thread.message.import"); + if (assistantImport?.type === "thread.message.import") { + assert.strictEqual(assistantImport.turnId, IMPORTED_TURN_ID); + } + assert.strictEqual(forkSession.mock.calls.length, 1); + assert.strictEqual(forkSession.mock.calls[0]?.[0].targetThreadId, LIVE_IMPORTED_THREAD_ID); }); }); @@ -266,6 +295,7 @@ describe("HermesConversationLifecycle", () => { "thread.meta.update", "thread.message.import", "thread.message.import", + "thread.activity.append", "thread.message.import", ], ); @@ -284,6 +314,11 @@ describe("HermesConversationLifecycle", () => { options: [{ id: "reasoningEffort", value: "high" }], }); } + const assistantImport = commands[3]; + assert.strictEqual(assistantImport?.type, "thread.message.import"); + if (assistantImport?.type === "thread.message.import") { + assert.strictEqual(assistantImport.turnId, IMPORTED_TURN_ID); + } const lineage = commands.at(-1); assert.strictEqual(lineage?.type, "thread.message.import"); if (lineage?.type === "thread.message.import") { @@ -321,7 +356,7 @@ describe("HermesConversationLifecycle", () => { assert.strictEqual(first.threadId, second.threadId); assert.deepStrictEqual([first.existing, second.existing].sort(), [false, true]); - assert.strictEqual(forkSession.mock.calls.length, 1); + assert.strictEqual(forkSession.mock.calls.length, 2); assert.strictEqual(commands.filter((command) => command.type === "thread.create").length, 1); }); }); diff --git a/apps/server/src/provider/hermes/HermesConversationLifecycle.ts b/apps/server/src/provider/hermes/HermesConversationLifecycle.ts index 5b79efa2901..25515741c8d 100644 --- a/apps/server/src/provider/hermes/HermesConversationLifecycle.ts +++ b/apps/server/src/provider/hermes/HermesConversationLifecycle.ts @@ -10,6 +10,7 @@ import { ProviderInstanceId, ThreadId, type HermesBridgeSessionListResponse, + type HermesBridgeSessionForkResponse, type HermesConversationForkInput, type HermesConversationForkResult, type OrchestrationCommand, @@ -120,6 +121,33 @@ export function makeHermesConversationLifecycle( .pipe(Effect.ignoreCause({ log: true })); }); + const importHistory = Effect.fn("HermesConversationLifecycle.importHistory")(function* ( + threadId: ThreadId, + fork: HermesBridgeSessionForkResponse, + ) { + for (const [index, message] of fork.messages.entries()) { + yield* dependencies.dispatch({ + type: "thread.message.import", + commandId: yield* nextCommandId(`hermes-history-${index}`), + threadId, + messageId: MessageId.make(`hermes-history:${fork.childSessionId}:${index}`), + role: message.role, + text: message.content, + ...(message.turnId !== undefined ? { turnId: message.turnId } : {}), + createdAt: message.createdAt, + }); + } + for (const activity of fork.activities ?? []) { + yield* dependencies.dispatch({ + type: "thread.activity.append", + commandId: yield* nextCommandId(`hermes-activity-${activity.id}`), + threadId, + activity, + createdAt: activity.createdAt, + }); + } + }); + const listSessions = Effect.gen(function* () { const client = yield* dependencies.getClient(); const [response, snapshot] = yield* Effect.all([ @@ -186,6 +214,21 @@ export function makeHermesConversationLifecycle( source.source !== "t3agent" && existingThreadId !== undefined ) { + const childSessionId = HermesBridgeSessionId.make(`t3-${existingThreadId}`); + const fork = yield* client + .forkSession({ + protocolVersion: HERMES_BRIDGE_PROTOCOL_VERSION, + requestId: HermesBridgeRequestId.make( + `conversation-refresh:${source.sessionId}:${existingThreadId}`, + ), + type: "session.fork", + sourceSessionId: source.sessionId, + childSessionId, + targetThreadId: existingThreadId, + ...(input.userTurnCount !== undefined ? { userTurnCount: input.userTurnCount } : {}), + }) + .pipe(Effect.retry({ times: 1 })); + yield* importHistory(existingThreadId, fork); return { threadId: existingThreadId, existing: true, @@ -292,17 +335,7 @@ export function makeHermesConversationLifecycle( threadId, title: fork.title, }); - for (const [index, message] of fork.messages.entries()) { - yield* dependencies.dispatch({ - type: "thread.message.import", - commandId: yield* nextCommandId(`hermes-history-${index}`), - threadId, - messageId: MessageId.make(`hermes-history:${fork.childSessionId}:${index}`), - role: message.role, - text: message.content, - createdAt: message.createdAt, - }); - } + yield* importHistory(threadId, fork); const isT3Fork = source.source === "t3agent" && source.threadId !== undefined; const lineage: HermesLineageMetadata = { diff --git a/integrations/hermes/t3agent/adapter.py b/integrations/hermes/t3agent/adapter.py index 562ded06a22..532ab5078f1 100644 --- a/integrations/hermes/t3agent/adapter.py +++ b/integrations/hermes/t3agent/adapter.py @@ -267,6 +267,323 @@ def _normalize_history_user_content(source: str, content: str) -> str: return f"{attachment_label}{body}" +def _history_reasoning_text(message: Dict[str, Any]) -> str: + for key in ("reasoning_content", "reasoning"): + value = message.get(key) + if isinstance(value, str) and value.strip(): + return value.strip() + return "" + + +def _history_tool_call_name_args( + tool_call: Dict[str, Any], +) -> Tuple[str, Dict[str, Any]]: + function = ( + tool_call.get("function") + if isinstance(tool_call.get("function"), dict) + else {} + ) + name = str( + function.get("name") or tool_call.get("name") or "unknown_tool" + ).strip() + raw_args = ( + function.get("arguments") + or tool_call.get("arguments") + or tool_call.get("args") + or {} + ) + if isinstance(raw_args, str): + try: + raw_args = json.loads(raw_args) + except (TypeError, ValueError): + raw_args = {"raw": raw_args} + if not isinstance(raw_args, dict): + raw_args = {"value": raw_args} + return name or "unknown_tool", raw_args + + +def _history_tool_call_id(tool_call: Dict[str, Any]) -> str: + return str( + tool_call.get("id") + or tool_call.get("call_id") + or tool_call.get("tool_call_id") + or "" + ).strip() + + +def _history_tool_title(name: str) -> str: + titles = { + "terminal": "Terminal", + "execute_code": "Ran code", + "process": "Process", + "web_search": "Web search", + "web_extract": "Read web page", + "browser_navigate": "Navigated browser", + "browser_snapshot": "Inspected browser", + "browser_click": "Clicked browser", + "browser_type": "Typed in browser", + "computer_use": "Used computer", + "read_file": "Read file", + "search_files": "Searched files", + "write_file": "Wrote file", + "patch": "Edited files", + "skill_view": "Read skill", + "delegate_task": "Delegated task", + } + return titles.get(name, name.replace("_", " ").strip().capitalize() or "Tool") + + +def _history_tool_item_type(name: str) -> str: + if name in {"terminal", "execute_code", "process"}: + return "command_execution" + if name in {"write_file", "patch"}: + return "file_change" + if name in {"web_search", "web_extract"}: + return "web_search" + if name in {"computer_use", "image_view"}: + return "image_view" + if name in {"delegate_task"}: + return "collab_agent_tool_call" + return "mcp_tool_call" + + +def _truncate_history_tool_text(value: Any, limit: int = 12_000) -> Optional[str]: + if not isinstance(value, str): + return None + normalized = value.strip() + if not normalized: + return None + if len(normalized) <= limit: + return normalized + return f"{normalized[:limit].rstrip()}\n\n[Tool output truncated during import]" + + +def _history_tool_detail(value: Optional[str]) -> Optional[str]: + if not value: + return None + first_line = next( + (line.strip() for line in value.splitlines() if line.strip()), + "", + ) + if not first_line: + return None + return first_line if len(first_line) <= 240 else f"{first_line[:239].rstrip()}…" + + +def _build_imported_history( + source_session_id: str, + source: str, + copied_messages: List[Dict[str, Any]], +) -> Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]: + history: List[Dict[str, Any]] = [] + activities: List[Dict[str, Any]] = [] + active_tool_calls: Dict[str, Dict[str, Any]] = {} + current_turn_id: Optional[str] = None + sequence = 0 + + def append_tool_activity( + *, + call_id: str, + name: str, + args: Dict[str, Any], + turn_id: Optional[str], + created_at: str, + result: Any = None, + status: str = "completed", + ) -> None: + nonlocal sequence + title = _history_tool_title(name) + item_type = _history_tool_item_type(name) + result_text = _truncate_history_tool_text(result) + item: Dict[str, Any] = { + "toolCallId": call_id, + "name": name, + "input": args, + } + command = args.get("command") or args.get("cmd") + if item_type == "command_execution" and command is not None: + item["command"] = command + if result_text is not None: + item["result"] = {"output": result_text} + payload: Dict[str, Any] = { + "itemType": item_type, + "status": status, + "title": title, + "data": { + "toolCallId": call_id, + "item": item, + }, + } + detail = _history_tool_detail(result_text) + if detail is not None: + payload["detail"] = detail + activities.append( + { + "id": _canonical_id( + "hermes_activity", + { + "sessionId": source_session_id, + "toolCallId": call_id, + "kind": "tool.completed", + }, + ), + "tone": "tool", + "kind": "tool.completed", + "summary": title, + "payload": payload, + "turnId": turn_id, + "sequence": sequence, + "createdAt": created_at, + } + ) + sequence += 1 + + def flush_unfinished(created_at: str) -> None: + for call_id, call in list(active_tool_calls.items()): + append_tool_activity( + call_id=call_id, + name=call["name"], + args=call["args"], + turn_id=call["turnId"], + created_at=created_at, + status="stopped", + ) + active_tool_calls.pop(call_id, None) + + for index, message in enumerate(copied_messages): + role = str(message.get("role") or "") + content = message.get("content") + created_at = _iso_timestamp(message.get("timestamp")) + + if role == "user": + if not isinstance(content, str) or _is_synthetic_history_user(content): + continue + flush_unfinished(created_at) + current_turn_id = _canonical_id( + "hermes_turn", + { + "sessionId": source_session_id, + "messageId": message.get("id") or index, + "timestamp": message.get("timestamp"), + }, + ) + history.append( + { + "role": role, + "content": _normalize_history_user_content(source, content), + "createdAt": created_at, + } + ) + continue + + if role == "system": + if isinstance(content, str): + history.append( + { + "role": role, + "content": content, + "createdAt": created_at, + } + ) + continue + + if role == "assistant": + reasoning = _history_reasoning_text(message) + if reasoning: + activities.append( + { + "id": _canonical_id( + "hermes_activity", + { + "sessionId": source_session_id, + "messageId": message.get("id") or index, + "kind": "reasoning", + }, + ), + "tone": "info", + "kind": "task.progress", + "summary": "Reasoning update", + "payload": { + "taskId": f"hermes-reasoning:{source_session_id}:{index}", + "detail": _truncate_history_tool_text(reasoning, 8_000), + }, + "turnId": current_turn_id, + "sequence": sequence, + "createdAt": created_at, + } + ) + sequence += 1 + + tool_calls = message.get("tool_calls") + if isinstance(tool_calls, list): + for call_index, tool_call in enumerate(tool_calls): + if not isinstance(tool_call, dict): + continue + call_id = _history_tool_call_id(tool_call) or _canonical_id( + "hermes_tool", + { + "sessionId": source_session_id, + "messageId": message.get("id") or index, + "callIndex": call_index, + }, + ) + name, args = _history_tool_call_name_args(tool_call) + active_tool_calls[call_id] = { + "name": name, + "args": args, + "turnId": current_turn_id, + "createdAt": created_at, + } + + if isinstance(content, str) and content.strip(): + item = { + "role": role, + "content": content, + "createdAt": created_at, + } + if current_turn_id is not None: + item["turnId"] = current_turn_id + history.append(item) + continue + + if role == "tool": + call_id = str(message.get("tool_call_id") or "").strip() + call = active_tool_calls.pop(call_id, None) if call_id else None + name = str(message.get("tool_name") or "").strip() + if call is not None: + name = call["name"] + args = call["args"] + turn_id = call["turnId"] + else: + args = {} + turn_id = current_turn_id + if not call_id: + call_id = _canonical_id( + "hermes_tool", + { + "sessionId": source_session_id, + "messageId": message.get("id") or index, + "toolName": name, + }, + ) + append_tool_activity( + call_id=call_id, + name=name or "unknown_tool", + args=args, + turn_id=turn_id, + created_at=created_at, + result=content, + ) + + final_created_at = ( + _iso_timestamp(copied_messages[-1].get("timestamp")) + if copied_messages + else _iso_timestamp(None) + ) + flush_unfinished(final_created_at) + return history, activities + + def _stored_model_config(session: Dict[str, Any]) -> Dict[str, Any]: raw_config = session.get("model_config") if isinstance(raw_config, dict): @@ -1385,30 +1702,12 @@ async def fork_once() -> Tuple[int, Dict[str, Any]]: if callable(set_reasoning): set_reasoning(target_session_key, reasoning_config) - history: List[Dict[str, Any]] = [] source = str(source_session.get("source") or "unknown") - for message in copied_messages: - role = str(message.get("role") or "") - content = message.get("content") - if role not in {"user", "assistant", "system"} or not isinstance( - content, str - ): - continue - if role == "assistant" and not content.strip(): - continue - if role == "user" and _is_synthetic_history_user(content): - continue - history.append( - { - "role": role, - "content": ( - _normalize_history_user_content(source, content) - if role == "user" - else content - ), - "createdAt": _iso_timestamp(message.get("timestamp")), - } - ) + history, activities = _build_imported_history( + source_session_id, + source, + copied_messages, + ) return 201, { "protocolVersion": PROTOCOL_VERSION, "requestId": request_id, @@ -1418,6 +1717,7 @@ async def fork_once() -> Tuple[int, Dict[str, Any]]: "source": source, "title": child_title, "messages": history, + "activities": activities, **( {"modelSelection": model_selection} if model_selection is not None diff --git a/integrations/hermes/t3agent/tests/test_adapter.py b/integrations/hermes/t3agent/tests/test_adapter.py index 03dc9fb27d3..c59bbd70f3e 100644 --- a/integrations/hermes/t3agent/tests/test_adapter.py +++ b/integrations/hermes/t3agent/tests/test_adapter.py @@ -369,6 +369,7 @@ def get_messages(self, session_id: str) -> List[Dict[str, Any]]: assert session_id == "discord-source" return [ { + "id": 1, "role": "user", "content": ( "[Triggering message id: `123` — use as `message_id` for " @@ -380,12 +381,35 @@ def get_messages(self, session_id: str) -> List[Dict[str, Any]]: "timestamp": 10, }, { + "id": 2, "role": "assistant", "content": "", - "tool_calls": [{"name": "search"}], + "reasoning_content": "I should inspect the relevant files.", + "tool_calls": [ + { + "id": "call-search", + "function": { + "name": "search_files", + "arguments": '{"query":"notes"}', + }, + } + ], "timestamp": 11, }, - {"role": "assistant", "content": "first answer", "timestamp": 12}, + { + "id": 3, + "role": "tool", + "content": "notes.txt", + "tool_call_id": "call-search", + "tool_name": "search_files", + "timestamp": 11.5, + }, + { + "id": 4, + "role": "assistant", + "content": "first answer", + "timestamp": 12, + }, { "role": "user", "content": ( @@ -471,6 +495,47 @@ def set_reasoning_override(session_key: str, config: Dict[str, Any]) -> None: "**Attached:** notes.txt\n\nfirst", "first answer", ] + turn_id = payload["messages"][1]["turnId"] + assert turn_id.startswith("hermes_turn_") + assert payload["activities"] == [ + { + "id": payload["activities"][0]["id"], + "tone": "info", + "kind": "task.progress", + "summary": "Reasoning update", + "payload": { + "taskId": "hermes-reasoning:discord-source:1", + "detail": "I should inspect the relevant files.", + }, + "turnId": turn_id, + "sequence": 0, + "createdAt": "1970-01-01T00:00:11Z", + }, + { + "id": payload["activities"][1]["id"], + "tone": "tool", + "kind": "tool.completed", + "summary": "Searched files", + "payload": { + "itemType": "mcp_tool_call", + "status": "completed", + "title": "Searched files", + "data": { + "toolCallId": "call-search", + "item": { + "toolCallId": "call-search", + "name": "search_files", + "input": {"query": "notes"}, + "result": {"output": "notes.txt"}, + }, + }, + "detail": "notes.txt", + }, + "turnId": turn_id, + "sequence": 1, + "createdAt": "1970-01-01T00:00:11.500000Z", + }, + ] assert created[0]["source"] == "t3agent" assert created[0]["parent_session_id"] == "discord-source" assert created[0]["model_config"]["gateway_runtime"]["provider"] == ( diff --git a/packages/contracts/src/hermesBridge.ts b/packages/contracts/src/hermesBridge.ts index 8693ff570df..30d6a191b46 100644 --- a/packages/contracts/src/hermesBridge.ts +++ b/packages/contracts/src/hermesBridge.ts @@ -6,7 +6,9 @@ import { NonNegativeInt, ThreadId, TrimmedNonEmptyString, + TurnId, } from "./baseSchemas.ts"; +import { OrchestrationThreadActivity } from "./orchestration.ts"; export const HERMES_BRIDGE_MAX_IMAGES = 8; export const HERMES_BRIDGE_MAX_IMAGE_DATA_URL_CHARS = 14_000_000; @@ -440,6 +442,7 @@ export const HermesBridgeHistoryMessage = openStruct({ role: Schema.Literals(["user", "assistant", "system"]), content: Schema.String, createdAt: IsoDateTime, + turnId: Schema.optionalKey(TurnId), }); export type HermesBridgeHistoryMessage = typeof HermesBridgeHistoryMessage.Type; @@ -451,6 +454,7 @@ export const HermesBridgeSessionForkResponse = openStruct({ source: TrimmedNonEmptyString, title: TrimmedNonEmptyString, messages: Schema.Array(HermesBridgeHistoryMessage), + activities: Schema.optionalKey(Schema.Array(OrchestrationThreadActivity)), modelSelection: Schema.optionalKey(HermesBridgeModelSelection), }); export type HermesBridgeSessionForkResponse = typeof HermesBridgeSessionForkResponse.Type; diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 0ac3416af98..ff1ccd7a624 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -798,6 +798,7 @@ const ThreadMessageImportCommand = Schema.Struct({ messageId: MessageId, role: OrchestrationMessageRole, text: Schema.String, + turnId: Schema.optional(TurnId), createdAt: IsoDateTime, }); From 555493a30c13ad404b89599eab879e14e9559073 Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 21:44:16 -0700 Subject: [PATCH 3/9] Refresh existing Hermes imports before opening --- .../components/hermes/HermesSessionBrowser.logic.test.ts | 7 ++++--- .../src/components/hermes/HermesSessionBrowser.logic.ts | 8 ++------ 2 files changed, 6 insertions(+), 9 deletions(-) diff --git a/apps/web/src/components/hermes/HermesSessionBrowser.logic.test.ts b/apps/web/src/components/hermes/HermesSessionBrowser.logic.test.ts index 38b51c83761..d32bb9219b2 100644 --- a/apps/web/src/components/hermes/HermesSessionBrowser.logic.test.ts +++ b/apps/web/src/components/hermes/HermesSessionBrowser.logic.test.ts @@ -22,7 +22,7 @@ describe("resolveHermesConversationSelection", () => { }); }); - it("opens an existing imported copy in open mode", () => { + it("refreshes an existing imported copy before opening it", () => { expect( resolveHermesConversationSelection({ mode: "open", @@ -33,8 +33,9 @@ describe("resolveHermesConversationSelection", () => { }, }), ).toEqual({ - type: "open-thread", - threadId: "imported-thread", + type: "fork-session", + sessionId: "session-1", + forceNew: false, }); }); diff --git a/apps/web/src/components/hermes/HermesSessionBrowser.logic.ts b/apps/web/src/components/hermes/HermesSessionBrowser.logic.ts index 5cfa221d29b..4b4b019c7e4 100644 --- a/apps/web/src/components/hermes/HermesSessionBrowser.logic.ts +++ b/apps/web/src/components/hermes/HermesSessionBrowser.logic.ts @@ -34,14 +34,10 @@ export function resolveHermesConversationSelection(input: { }; } - const existingThreadId = - input.session.source === "t3agent" - ? input.session.threadId - : input.session.importedThreadIds?.[0]; - if (existingThreadId !== undefined) { + if (input.session.source === "t3agent" && input.session.threadId !== undefined) { return { type: "open-thread", - threadId: existingThreadId, + threadId: input.session.threadId, }; } From dde3415fb8f28120d45586d05dea77541926b05c Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 21:54:10 -0700 Subject: [PATCH 4/9] Preserve Hermes history slots during refresh --- integrations/hermes/t3agent/adapter.py | 70 ++++++++++++++----- .../hermes/t3agent/tests/test_adapter.py | 8 ++- 2 files changed, 60 insertions(+), 18 deletions(-) diff --git a/integrations/hermes/t3agent/adapter.py b/integrations/hermes/t3agent/adapter.py index 532ab5078f1..cf8d1b25203 100644 --- a/integrations/hermes/t3agent/adapter.py +++ b/integrations/hermes/t3agent/adapter.py @@ -450,6 +450,10 @@ def flush_unfinished(created_at: str) -> None: ) active_tool_calls.pop(call_id, None) + # Preserve one history slot per source row. The lifecycle derives stable + # message IDs from these positions, so compacting hidden tool/synthetic + # rows would overwrite the wrong legacy rows when an existing import is + # refreshed. for index, message in enumerate(copied_messages): role = str(message.get("role") or "") content = message.get("content") @@ -457,6 +461,18 @@ def flush_unfinished(created_at: str) -> None: if role == "user": if not isinstance(content, str) or _is_synthetic_history_user(content): + history.append( + { + "role": "assistant", + "content": "", + "createdAt": created_at, + **( + {"turnId": current_turn_id} + if current_turn_id is not None + else {} + ), + } + ) continue flush_unfinished(created_at) current_turn_id = _canonical_id( @@ -477,14 +493,13 @@ def flush_unfinished(created_at: str) -> None: continue if role == "system": - if isinstance(content, str): - history.append( - { - "role": role, - "content": content, - "createdAt": created_at, - } - ) + history.append( + { + "role": role, + "content": content if isinstance(content, str) else "", + "createdAt": created_at, + } + ) continue if role == "assistant": @@ -535,15 +550,14 @@ def flush_unfinished(created_at: str) -> None: "createdAt": created_at, } - if isinstance(content, str) and content.strip(): - item = { - "role": role, - "content": content, - "createdAt": created_at, - } - if current_turn_id is not None: - item["turnId"] = current_turn_id - history.append(item) + item = { + "role": role, + "content": content if isinstance(content, str) else "", + "createdAt": created_at, + } + if current_turn_id is not None: + item["turnId"] = current_turn_id + history.append(item) continue if role == "tool": @@ -574,6 +588,28 @@ def flush_unfinished(created_at: str) -> None: created_at=created_at, result=content, ) + history.append( + { + "role": "assistant", + "content": "", + "createdAt": created_at, + **({"turnId": turn_id} if turn_id is not None else {}), + } + ) + continue + + history.append( + { + "role": "assistant", + "content": "", + "createdAt": created_at, + **( + {"turnId": current_turn_id} + if current_turn_id is not None + else {} + ), + } + ) final_created_at = ( _iso_timestamp(copied_messages[-1].get("timestamp")) diff --git a/integrations/hermes/t3agent/tests/test_adapter.py b/integrations/hermes/t3agent/tests/test_adapter.py index c59bbd70f3e..b6d11bb1429 100644 --- a/integrations/hermes/t3agent/tests/test_adapter.py +++ b/integrations/hermes/t3agent/tests/test_adapter.py @@ -493,10 +493,16 @@ def set_reasoning_override(session_key: str, config: Dict[str, Any]) -> None: assert payload["title"] == "Planning #2" assert [message["content"] for message in payload["messages"]] == [ "**Attached:** notes.txt\n\nfirst", + "", + "", "first answer", + "", ] - turn_id = payload["messages"][1]["turnId"] + turn_id = payload["messages"][3]["turnId"] assert turn_id.startswith("hermes_turn_") + assert payload["messages"][1]["turnId"] == turn_id + assert payload["messages"][2]["turnId"] == turn_id + assert payload["messages"][4]["turnId"] == turn_id assert payload["activities"] == [ { "id": payload["activities"][0]["id"], From 5800195b84161c3dffe775eb76ec01ec0b3c96bf Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 21:58:06 -0700 Subject: [PATCH 5/9] Use fresh request IDs for Hermes refreshes --- .../provider/hermes/HermesConversationLifecycle.test.ts | 9 +++++++++ .../src/provider/hermes/HermesConversationLifecycle.ts | 7 ++++--- 2 files changed, 13 insertions(+), 3 deletions(-) diff --git a/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts b/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts index 5cf49a62a8f..11a24618143 100644 --- a/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts +++ b/apps/server/src/provider/hermes/HermesConversationLifecycle.test.ts @@ -274,6 +274,15 @@ describe("HermesConversationLifecycle", () => { } assert.strictEqual(forkSession.mock.calls.length, 1); assert.strictEqual(forkSession.mock.calls[0]?.[0].targetThreadId, LIVE_IMPORTED_THREAD_ID); + + yield* lifecycle.forkConversation({ + source: { type: "session", sessionId: SOURCE_SESSION_ID }, + }); + assert.strictEqual(forkSession.mock.calls.length, 2); + assert.notStrictEqual( + forkSession.mock.calls[0]?.[0].requestId, + forkSession.mock.calls[1]?.[0].requestId, + ); }); }); diff --git a/apps/server/src/provider/hermes/HermesConversationLifecycle.ts b/apps/server/src/provider/hermes/HermesConversationLifecycle.ts index 25515741c8d..71b8cbc267e 100644 --- a/apps/server/src/provider/hermes/HermesConversationLifecycle.ts +++ b/apps/server/src/provider/hermes/HermesConversationLifecycle.ts @@ -215,12 +215,13 @@ export function makeHermesConversationLifecycle( existingThreadId !== undefined ) { const childSessionId = HermesBridgeSessionId.make(`t3-${existingThreadId}`); + const refreshRequestId = HermesBridgeRequestId.make( + `conversation-refresh:${source.sessionId}:${existingThreadId}:${yield* dependencies.randomUuid}`, + ); const fork = yield* client .forkSession({ protocolVersion: HERMES_BRIDGE_PROTOCOL_VERSION, - requestId: HermesBridgeRequestId.make( - `conversation-refresh:${source.sessionId}:${existingThreadId}`, - ), + requestId: refreshRequestId, type: "session.fork", sourceSessionId: source.sessionId, childSessionId, From 4f34402afe745537a0bedcc13765259feb0bdd50 Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 22:00:43 -0700 Subject: [PATCH 6/9] Preserve Hermes source order in imported timelines --- integrations/hermes/t3agent/adapter.py | 17 +++++++++++++---- .../hermes/t3agent/tests/test_adapter.py | 3 ++- 2 files changed, 15 insertions(+), 5 deletions(-) diff --git a/integrations/hermes/t3agent/adapter.py b/integrations/hermes/t3agent/adapter.py index cf8d1b25203..db6b0f7c9cb 100644 --- a/integrations/hermes/t3agent/adapter.py +++ b/integrations/hermes/t3agent/adapter.py @@ -23,7 +23,7 @@ from pathlib import Path import re import secrets -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from typing import Any, Awaitable, Callable, Dict, List, Optional, Tuple from urllib.parse import quote, unquote_to_bytes, urlparse @@ -380,6 +380,15 @@ def _build_imported_history( active_tool_calls: Dict[str, Dict[str, Any]] = {} current_turn_id: Optional[str] = None sequence = 0 + last_created_at: Optional[datetime] = None + + def next_created_at(value: Any) -> str: + nonlocal last_created_at + candidate = datetime.fromisoformat(_iso_timestamp(value).replace("Z", "+00:00")) + if last_created_at is not None and candidate <= last_created_at: + candidate = last_created_at + timedelta(milliseconds=1) + last_created_at = candidate + return candidate.isoformat().replace("+00:00", "Z") def append_tool_activity( *, @@ -457,7 +466,7 @@ def flush_unfinished(created_at: str) -> None: for index, message in enumerate(copied_messages): role = str(message.get("role") or "") content = message.get("content") - created_at = _iso_timestamp(message.get("timestamp")) + created_at = next_created_at(message.get("timestamp")) if role == "user": if not isinstance(content, str) or _is_synthetic_history_user(content): @@ -612,8 +621,8 @@ def flush_unfinished(created_at: str) -> None: ) final_created_at = ( - _iso_timestamp(copied_messages[-1].get("timestamp")) - if copied_messages + last_created_at.isoformat().replace("+00:00", "Z") + if last_created_at is not None else _iso_timestamp(None) ) flush_unfinished(final_created_at) diff --git a/integrations/hermes/t3agent/tests/test_adapter.py b/integrations/hermes/t3agent/tests/test_adapter.py index b6d11bb1429..0d973098459 100644 --- a/integrations/hermes/t3agent/tests/test_adapter.py +++ b/integrations/hermes/t3agent/tests/test_adapter.py @@ -408,7 +408,7 @@ def get_messages(self, session_id: str) -> List[Dict[str, Any]]: "id": 4, "role": "assistant", "content": "first answer", - "timestamp": 12, + "timestamp": 11, }, { "role": "user", @@ -503,6 +503,7 @@ def set_reasoning_override(session_key: str, config: Dict[str, Any]) -> None: assert payload["messages"][1]["turnId"] == turn_id assert payload["messages"][2]["turnId"] == turn_id assert payload["messages"][4]["turnId"] == turn_id + assert payload["messages"][3]["createdAt"] == "1970-01-01T00:00:11.501000Z" assert payload["activities"] == [ { "id": payload["activities"][0]["id"], From 322767dca303784c511e0117bb555bdcdd03d750 Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 22:06:53 -0700 Subject: [PATCH 7/9] Replace imported history slots on refresh --- .../Layers/ProjectionPipeline.test.ts | 73 ++++++++++++++ .../Layers/ProjectionPipeline.ts | 5 +- .../src/orchestration/decider.import.test.ts | 96 +++++++++++++++++++ apps/server/src/orchestration/decider.ts | 1 + apps/server/src/orchestration/projector.ts | 2 + 5 files changed, 176 insertions(+), 1 deletion(-) create mode 100644 apps/server/src/orchestration/decider.import.test.ts diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 926182a3ef0..de3ee733a41 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -242,6 +242,79 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-base-")))( "OrchestrationProjectionPipeline", (it) => { + it.effect("fully replaces imported history slots on refresh", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-import-refresh"); + const messageId = MessageId.make("message-import-refresh"); + const originalAt = "2026-01-01T00:00:00.000Z"; + const refreshedAt = "2026-01-01T00:01:00.000Z"; + + yield* eventStore.append({ + type: "thread.message-sent", + eventId: EventId.make("evt-import-original"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: originalAt, + commandId: CommandId.make("cmd-import-original"), + causationEventId: null, + correlationId: CommandId.make("cmd-import-original"), + metadata: {}, + payload: { + threadId, + messageId, + role: "user", + text: "legacy provider wrapper", + turnId: null, + streaming: false, + createdAt: originalAt, + updatedAt: originalAt, + }, + }); + yield* eventStore.append({ + type: "thread.message-sent", + eventId: EventId.make("evt-import-refresh"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: refreshedAt, + commandId: CommandId.make("cmd-import-refresh"), + causationEventId: null, + correlationId: CommandId.make("cmd-import-refresh"), + metadata: {}, + payload: { + threadId, + messageId, + role: "assistant", + text: "", + replaceText: true, + turnId: null, + streaming: false, + imported: true, + createdAt: refreshedAt, + updatedAt: refreshedAt, + }, + }); + + yield* projectionPipeline.bootstrap; + + const rows = yield* sql<{ + readonly role: string; + readonly text: string; + readonly createdAt: string; + }>` + SELECT + role, + text, + created_at AS "createdAt" + FROM projection_thread_messages + WHERE message_id = ${messageId} + `; + assert.deepEqual(rows, [{ role: "assistant", text: "", createdAt: refreshedAt }]); + }), + ); + it.effect("stores message attachment references without mutating payloads", () => Effect.gen(function* () { const projectionPipeline = yield* OrchestrationProjectionPipeline; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 048db29a599..fc32155fb6e 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -879,7 +879,10 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti text: nextText, ...(nextAttachments !== undefined ? { attachments: [...nextAttachments] } : {}), isStreaming: event.payload.streaming, - createdAt: previousMessage?.createdAt ?? event.payload.createdAt, + createdAt: + event.payload.imported || previousMessage === undefined + ? event.payload.createdAt + : previousMessage.createdAt, updatedAt: event.payload.updatedAt, }); return; diff --git a/apps/server/src/orchestration/decider.import.test.ts b/apps/server/src/orchestration/decider.import.test.ts new file mode 100644 index 00000000000..a3bb8506680 --- /dev/null +++ b/apps/server/src/orchestration/decider.import.test.ts @@ -0,0 +1,96 @@ +import { + CommandId, + DEFAULT_PROVIDER_INTERACTION_MODE, + EventId, + MessageId, + ProjectId, + ProviderInstanceId, + ThreadId, +} from "@t3tools/contracts"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; + +import { decideOrchestrationCommand } from "./decider.ts"; +import { createEmptyReadModel, projectEvent } from "./projector.ts"; + +it.layer(NodeServices.layer)("decider imported messages", (it) => { + it.effect("replaces an existing history slot even when the imported text is empty", () => + Effect.gen(function* () { + const createdAt = "2026-01-01T00:00:00.000Z"; + const projectId = ProjectId.make("project-import"); + const threadId = ThreadId.make("thread-import"); + const initial = createEmptyReadModel(createdAt); + const withProject = yield* projectEvent(initial, { + sequence: 1, + eventId: EventId.make("event-project-import"), + aggregateKind: "project", + aggregateId: projectId, + type: "project.created", + occurredAt: createdAt, + commandId: CommandId.make("command-project-import"), + causationEventId: null, + correlationId: CommandId.make("command-project-import"), + metadata: {}, + payload: { + projectId, + title: "Imports", + workspaceRoot: "/tmp/imports", + defaultModelSelection: null, + scripts: [], + createdAt, + updatedAt: createdAt, + }, + }); + const readModel = yield* projectEvent(withProject, { + sequence: 2, + eventId: EventId.make("event-thread-import"), + aggregateKind: "thread", + aggregateId: threadId, + type: "thread.created", + occurredAt: createdAt, + commandId: CommandId.make("command-thread-import"), + causationEventId: null, + correlationId: CommandId.make("command-thread-import"), + metadata: {}, + payload: { + threadId, + projectId, + title: "Imported thread", + modelSelection: { + instanceId: ProviderInstanceId.make("hermes"), + model: "openai-codex::gpt-5.6-sol", + }, + runtimeMode: "full-access", + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + branch: null, + worktreePath: null, + createdAt, + updatedAt: createdAt, + }, + }); + + const event = yield* decideOrchestrationCommand({ + command: { + type: "thread.message.import", + commandId: CommandId.make("command-message-import"), + threadId, + messageId: MessageId.make("message-import"), + role: "assistant", + text: "", + createdAt: "2026-01-01T00:01:00.000Z", + }, + readModel, + }); + + const importedEvent = Array.isArray(event) ? event[0] : event; + expect(importedEvent?.type).toBe("thread.message-sent"); + if (importedEvent?.type !== "thread.message-sent") { + return; + } + expect(importedEvent.payload.imported).toBe(true); + expect(importedEvent.payload.replaceText).toBe(true); + expect(importedEvent.payload.text).toBe(""); + }), + ); +}); diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 356251850bb..7dc3ce3b594 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -921,6 +921,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" messageId: command.messageId, role: command.role, text: command.text, + replaceText: true, turnId: command.turnId ?? null, streaming: false, imported: true, diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index a9fe821dbc0..45a54ce72b9 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -443,6 +443,7 @@ export function projectEvent( entry.id === message.id ? { ...entry, + role: message.role, text: message.streaming ? `${entry.text}${message.text}` : payload.replaceText @@ -451,6 +452,7 @@ export function projectEvent( ? message.text : entry.text, streaming: message.streaming, + createdAt: payload.imported ? message.createdAt : entry.createdAt, updatedAt: message.updatedAt, turnId: message.turnId, ...(message.attachments !== undefined From 428638bfb03bca5312472d0ad40a1862fbb8dcb5 Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 22:12:06 -0700 Subject: [PATCH 8/9] Remove empty assistant rows from chat timelines --- apps/web/src/session-logic.test.ts | 29 +++++++++++++++++++++++++++++ apps/web/src/session-logic.ts | 20 ++++++++++++++------ 2 files changed, 43 insertions(+), 6 deletions(-) diff --git a/apps/web/src/session-logic.test.ts b/apps/web/src/session-logic.test.ts index 20dd3a198b9..d6da709eb92 100644 --- a/apps/web/src/session-logic.test.ts +++ b/apps/web/src/session-logic.test.ts @@ -1522,6 +1522,35 @@ describe("deriveWorkLogEntries", () => { }); describe("deriveTimelineEntries", () => { + it("omits empty settled assistant placeholders without hiding their work entries", () => { + const entries = deriveTimelineEntries( + [ + { + id: MessageId.make("empty-assistant"), + role: "assistant", + text: "", + createdAt: "2026-02-23T00:00:01.000Z", + turnId: TurnId.make("turn-1"), + updatedAt: "2026-02-23T00:00:01.000Z", + streaming: false, + }, + ], + [], + [ + { + id: "work-1", + createdAt: "2026-02-23T00:00:02.000Z", + turnId: TurnId.make("turn-1"), + label: "Read file", + tone: "tool", + }, + ], + ); + + expect(entries).toHaveLength(1); + expect(entries[0]).toMatchObject({ kind: "work", id: "work-1" }); + }); + it("includes proposed plans alongside messages and work entries in chronological order", () => { const entries = deriveTimelineEntries( [ diff --git a/apps/web/src/session-logic.ts b/apps/web/src/session-logic.ts index 16a85c16ba7..833ccbe5622 100644 --- a/apps/web/src/session-logic.ts +++ b/apps/web/src/session-logic.ts @@ -1340,12 +1340,20 @@ export function deriveTimelineEntries( proposedPlans: ReadonlyArray, workEntries: ReadonlyArray, ): TimelineEntry[] { - const messageRows: TimelineEntry[] = messages.map((message) => ({ - id: message.id, - kind: "message", - createdAt: message.createdAt, - message, - })); + const messageRows: TimelineEntry[] = messages + .filter( + (message) => + message.role !== "assistant" || + message.streaming || + message.text.trim().length > 0 || + (message.attachments?.length ?? 0) > 0, + ) + .map((message) => ({ + id: message.id, + kind: "message", + createdAt: message.createdAt, + message, + })); const proposedPlanRows: TimelineEntry[] = proposedPlans.map((proposedPlan) => ({ id: proposedPlan.id, kind: "proposed-plan", From 73bbd72d341251278e7981c0271252b882aa493e Mon Sep 17 00:00:00 2001 From: eimexdev Date: Thu, 23 Jul 2026 22:34:39 -0700 Subject: [PATCH 9/9] Render live Hermes tools as native activities --- .../Layers/ProviderRuntimeIngestion.test.ts | 30 ++++- .../Layers/ProviderRuntimeIngestion.ts | 3 + .../src/provider/Layers/HermesAdapter.test.ts | 71 ++++++++++++ .../src/provider/Layers/HermesAdapter.ts | 109 ++++++++++++++++++ integrations/hermes/t3agent/README.md | 7 +- integrations/hermes/t3agent/adapter.py | 50 ++++++++ .../hermes/t3agent/tests/test_adapter.py | 58 ++++++++++ packages/contracts/src/hermesBridge.test.ts | 32 +++++ packages/contracts/src/hermesBridge.ts | 24 ++++ 9 files changed, 376 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 0c26b439738..45953004bb8 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -2866,9 +2866,16 @@ describe("ProviderRuntimeIngestion", () => { turnId: asTurnId("turn-9"), payload: { itemType: "command_execution", - status: "in_progress", + status: "inProgress", title: "Read file", detail: "/tmp/file.ts", + data: { + toolCallId: "call-read", + item: { + name: "read_file", + input: { path: "/tmp/file.ts" }, + }, + }, }, }); @@ -2883,11 +2890,22 @@ describe("ProviderRuntimeIngestion", () => { ); expect(thread.session?.status).toBe("ready"); - expect( - thread.activities.some( - (activity: ProviderRuntimeTestActivity) => activity.kind === "tool.started", - ), - ).toBe(true); + const toolStarted = thread.activities.find( + (activity: ProviderRuntimeTestActivity) => activity.kind === "tool.started", + ); + expect(toolStarted?.payload).toEqual({ + itemType: "command_execution", + status: "inProgress", + title: "Read file", + detail: "/tmp/file.ts", + data: { + toolCallId: "call-read", + item: { + name: "read_file", + input: { path: "/tmp/file.ts" }, + }, + }, + }); }); it("consumes P1 runtime events into thread metadata, diff checkpoints, and activities", async () => { diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 9372eb0cc04..d02edb4a146 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -693,7 +693,10 @@ export function runtimeEventToActivities( summary: `${event.payload.title ?? "Tool"} started`, payload: { itemType: event.payload.itemType, + ...(event.payload.status ? { status: event.payload.status } : {}), + ...(event.payload.title ? { title: event.payload.title } : {}), ...(event.payload.detail ? { detail: truncateDetail(event.payload.detail) } : {}), + ...(event.payload.data !== undefined ? { data: event.payload.data } : {}), }, turnId: toTurnId(event.turnId) ?? null, ...maybeSequence, diff --git a/apps/server/src/provider/Layers/HermesAdapter.test.ts b/apps/server/src/provider/Layers/HermesAdapter.test.ts index 1f84d22c33d..8b64b6d98dc 100644 --- a/apps/server/src/provider/Layers/HermesAdapter.test.ts +++ b/apps/server/src/provider/Layers/HermesAdapter.test.ts @@ -236,6 +236,77 @@ it.layer(testLayer)("HermesAdapter", (it) => { }), ); + it.effect("emits correlated native tool lifecycle events with full tool data", () => + Effect.gen(function* () { + const { adapter } = yield* HermesAdapterTestHarness; + const threadId = ThreadId.make("hermes-tool-thread"); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("hermes"), + threadId, + runtimeMode: "full-access", + }); + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.take(3), + Stream.runCollect, + Effect.forkChild, + ); + yield* Effect.yieldNow; + const turn = yield* adapter.sendTurn({ threadId, input: "read the skill" }); + const sourceMessageId = `hermes-user:${turn.turnId}`; + yield* adapter.receiveCallback({ + protocolVersion: HERMES_BRIDGE_PROTOCOL_VERSION, + requestId: "callback-tool-start-request", + deliveryId: "callback-tool-start-delivery", + type: "tool.started", + chatId: "t3agent", + threadId, + sourceMessageId, + toolCallId: "call-skill", + name: "skill_view", + input: { name: "query" }, + }); + yield* adapter.receiveCallback({ + protocolVersion: HERMES_BRIDGE_PROTOCOL_VERSION, + requestId: "callback-tool-complete-request", + deliveryId: "callback-tool-complete-delivery", + type: "tool.completed", + chatId: "t3agent", + threadId, + sourceMessageId, + toolCallId: "call-skill", + name: "skill_view", + input: { name: "query" }, + result: "Skill loaded", + isError: false, + }); + + const events = Array.from(yield* Fiber.join(eventsFiber)); + NodeAssert.deepEqual( + events.map((event) => event.type), + ["turn.started", "item.started", "item.completed"], + ); + const started = events[1]; + const completed = events[2]; + NodeAssert.equal(started?.itemId, completed?.itemId); + NodeAssert.deepEqual(completed?.payload, { + itemType: "mcp_tool_call", + status: "completed", + title: "Read skill", + detail: "Skill loaded", + data: { + toolCallId: "call-skill", + item: { + toolCallId: "call-skill", + name: "skill_view", + input: { name: "query" }, + result: { output: "Skill loaded" }, + }, + }, + }); + }), + ); + it.effect("acknowledges a repeated delivery id as a duplicate", () => Effect.gen(function* () { const { adapter } = yield* HermesAdapterTestHarness; diff --git a/apps/server/src/provider/Layers/HermesAdapter.ts b/apps/server/src/provider/Layers/HermesAdapter.ts index ea13a70e975..70d0eeed2c1 100644 --- a/apps/server/src/provider/Layers/HermesAdapter.ts +++ b/apps/server/src/provider/Layers/HermesAdapter.ts @@ -230,6 +230,70 @@ function textDelta(previous: string, next: string): string { return ""; } +function hermesToolTitle(name: string): string { + const titles: Readonly> = { + terminal: "Terminal", + execute_code: "Ran code", + process: "Process", + web_search: "Web search", + web_extract: "Read web page", + browser_navigate: "Navigated browser", + browser_snapshot: "Inspected browser", + browser_click: "Clicked browser", + browser_type: "Typed in browser", + computer_use: "Used computer", + read_file: "Read file", + search_files: "Searched files", + write_file: "Wrote file", + patch: "Edited files", + skill_view: "Read skill", + delegate_task: "Delegated task", + }; + return titles[name] ?? name.replaceAll("_", " ").replace(/^\w/, (value) => value.toUpperCase()); +} + +function hermesToolItemType( + name: string, +): + | "command_execution" + | "file_change" + | "mcp_tool_call" + | "collab_agent_tool_call" + | "web_search" + | "image_view" { + if (name === "terminal" || name === "execute_code" || name === "process") { + return "command_execution"; + } + if (name === "write_file" || name === "patch") return "file_change"; + if (name === "web_search" || name === "web_extract") return "web_search"; + if (name === "computer_use" || name === "image_view") return "image_view"; + if (name === "delegate_task") return "collab_agent_tool_call"; + return "mcp_tool_call"; +} + +function toolResultData(result: unknown): unknown { + return typeof result === "string" ? { output: result } : result; +} + +function toolResultDetail(result: unknown): string | undefined { + const text = + typeof result === "string" + ? result + : (() => { + try { + return JSON.stringify(result); + } catch { + return ""; + } + })(); + const firstLine = text + .split(/\r?\n/) + .map((line) => line.trim()) + .find(Boolean); + if (!firstLine) return undefined; + return firstLine.length <= 240 ? firstLine : `${firstLine.slice(0, 239).trimEnd()}…`; +} + function firstAnswer(answers: ProviderUserInputAnswers): unknown { return Object.values(answers)[0] ?? ""; } @@ -454,6 +518,47 @@ export const makeHermesAdapter = Effect.fn("makeHermesAdapter")(function* ( }); }); + const handleToolCallback = Effect.fn("HermesAdapter.handleToolCallback")(function* ( + callback: Extract, + threadId: ThreadId, + ) { + const base = yield* eventBase( + callback, + threadId, + callback.type === "tool.started" ? "tool-started" : "tool-completed", + ); + const itemType = hermesToolItemType(callback.name); + const title = hermesToolTitle(callback.name); + const item = { + toolCallId: callback.toolCallId, + name: callback.name, + input: callback.input, + ...(callback.type === "tool.completed" ? { result: toolResultData(callback.result) } : {}), + }; + const detail = + callback.type === "tool.completed" ? toolResultDetail(callback.result) : undefined; + yield* publish({ + ...base, + type: callback.type === "tool.started" ? "item.started" : "item.completed", + itemId: RuntimeItemId.make(`hermes-tool:${callback.toolCallId}`), + payload: { + itemType, + status: + callback.type === "tool.started" + ? "inProgress" + : callback.isError + ? "failed" + : "completed", + title, + ...(detail ? { detail } : {}), + data: { + toolCallId: callback.toolCallId, + item, + }, + }, + }); + }); + const rememberApproval = Effect.fn("HermesAdapter.rememberApproval")(function* ( callback: HermesBridgeApprovalRequest, threadId: ThreadId, @@ -554,6 +659,10 @@ export const makeHermesAdapter = Effect.fn("makeHermesAdapter")(function* ( case "message.edit": yield* handleMessageCallback(callback, threadId, context); break; + case "tool.started": + case "tool.completed": + yield* handleToolCallback(callback, threadId); + break; case "message.delete": return yield* requestError( "callback.message.delete", diff --git a/integrations/hermes/t3agent/README.md b/integrations/hermes/t3agent/README.md index 7e5ab77f32d..de79424b50d 100644 --- a/integrations/hermes/t3agent/README.md +++ b/integrations/hermes/t3agent/README.md @@ -110,8 +110,11 @@ capabilities response includes the initial command catalog; T3 can use it for slash-command completion without reimplementing command behavior. Callback tags are `message.send`, `message.edit`, `message.delete`, -`typing.set`, `turn.complete`, `approval.request`, `clarification.request`, -`slash-confirmation.request`, and `thread.create`. Message send/edit content is +`tool.started`, `tool.completed`, `typing.set`, `turn.complete`, +`approval.request`, `clarification.request`, `slash-confirmation.request`, and +`thread.create`. Structured tool callbacks carry the tool-call ID, name, +arguments, and completion result, so T3 renders native expandable tool cards +instead of Hermes' compact emoji progress lines. Message send/edit content is always cumulative full content. `final` closes an individual message bubble; it does not imply that the whole agent turn is done. The plugin emits `turn.complete` from Hermes' post-delivery lifecycle hook only after the final diff --git a/integrations/hermes/t3agent/adapter.py b/integrations/hermes/t3agent/adapter.py index db6b0f7c9cb..5cd96c74f21 100644 --- a/integrations/hermes/t3agent/adapter.py +++ b/integrations/hermes/t3agent/adapter.py @@ -983,6 +983,7 @@ class T3AgentAdapter(BasePlatformAdapter): supports_async_delivery = True supports_code_blocks = True + supports_structured_tool_events = True def __init__(self, config: PlatformConfig): super().__init__(config, Platform("t3agent")) @@ -2194,6 +2195,55 @@ async def send( retryable=not ok, ) + def _tool_event_fields( + self, + chat_id: str, + tool_call_id: str, + name: str, + args: Any, + metadata: Optional[Dict[str, Any]], + ) -> Dict[str, Any]: + fields: Dict[str, Any] = { + "chatId": str(chat_id), + "toolCallId": str(tool_call_id), + "name": str(name), + "input": args if isinstance(args, dict) else {"value": args}, + } + thread_id = _metadata_value(metadata, "threadId", "thread_id") + if thread_id: + fields["threadId"] = thread_id + source_message_id = self._source_message_id(chat_id, metadata) + if source_message_id: + fields["sourceMessageId"] = source_message_id + return fields + + async def send_tool_started( + self, + chat_id: str, + tool_call_id: str, + name: str, + args: Any, + metadata: Optional[Dict[str, Any]] = None, + ) -> None: + fields = self._tool_event_fields(chat_id, tool_call_id, name, args, metadata) + await self._post_event("tool.started", fields) + + async def send_tool_completed( + self, + chat_id: str, + tool_call_id: str, + name: str, + args: Any, + result: Any, + metadata: Optional[Dict[str, Any]] = None, + *, + is_error: bool = False, + ) -> None: + fields = self._tool_event_fields(chat_id, tool_call_id, name, args, metadata) + fields["result"] = result + fields["isError"] = bool(is_error) + await self._post_event("tool.completed", fields) + async def send_image_file( self, chat_id: str, diff --git a/integrations/hermes/t3agent/tests/test_adapter.py b/integrations/hermes/t3agent/tests/test_adapter.py index 0d973098459..f862a4ba4eb 100644 --- a/integrations/hermes/t3agent/tests/test_adapter.py +++ b/integrations/hermes/t3agent/tests/test_adapter.py @@ -1019,6 +1019,64 @@ async def receive(request: web.Request) -> web.Response: await server.close() +@pytest.mark.asyncio +async def test_structured_tool_lifecycle_posts_full_tool_frames( + fake_platform: SimpleNamespace, +) -> None: + received: List[Dict[str, Any]] = [] + + async def receive(request: web.Request) -> web.Response: + frame = await request.json() + received.append(frame) + return web.json_response( + { + "protocolVersion": 1, + "requestId": frame["requestId"], + "deliveryId": frame["deliveryId"], + "status": "accepted", + } + ) + + app = web.Application() + app.router.add_post("/api/hermes/hermes-test/events", receive) + server = TestServer(app) + await server.start_server() + adapter = adapter_module.T3AgentAdapter(make_config(bridge_url=str(server.make_url("/")))) + adapter.bridge_url = adapter.bridge_url.rstrip("/") + adapter._client = ClientSession() + adapter._processing_sources[("chat-1", "thread-1")] = "hermes-user:turn-1" + try: + await adapter.send_tool_started( + "chat-1", + "call-1", + "skill_view", + {"name": "query"}, + metadata={"thread_id": "thread-1"}, + ) + await adapter.send_tool_completed( + "chat-1", + "call-1", + "skill_view", + {"name": "query"}, + "Skill loaded", + metadata={"thread_id": "thread-1"}, + ) + + assert [frame["type"] for frame in received] == [ + "tool.started", + "tool.completed", + ] + assert received[0]["sourceMessageId"] == "hermes-user:turn-1" + assert received[0]["toolCallId"] == "call-1" + assert received[0]["input"] == {"name": "query"} + assert received[1]["result"] == "Skill loaded" + assert received[1]["isError"] is False + finally: + await adapter._client.close() + adapter._client = None + await server.close() + + @pytest.mark.asyncio async def test_real_stream_consumer_does_not_complete_turn_at_segment_boundaries( fake_platform: SimpleNamespace, diff --git a/packages/contracts/src/hermesBridge.test.ts b/packages/contracts/src/hermesBridge.test.ts index 1049b99525a..f8e14b1c284 100644 --- a/packages/contracts/src/hermesBridge.test.ts +++ b/packages/contracts/src/hermesBridge.test.ts @@ -288,6 +288,38 @@ describe("Hermes bridge Hermes to T3 callbacks", () => { expect(edit.futurePresentation).toBe("markdown-v2"); }); + it("decodes structured tool lifecycle callbacks", () => { + const started = decodeHermesToT3({ + ...callbackFields, + type: "tool.started", + chatId: "t3agent", + threadId: "thread-1", + sourceMessageId: "hermes-user:turn-1", + toolCallId: "call-1", + name: "skill_view", + input: { name: "query" }, + }); + const completed = decodeHermesToT3({ + ...callbackFields, + requestId: "request-2", + deliveryId: "delivery-2", + type: "tool.completed", + chatId: "t3agent", + threadId: "thread-1", + sourceMessageId: "hermes-user:turn-1", + toolCallId: "call-1", + name: "skill_view", + input: { name: "query" }, + result: "Skill loaded", + isError: false, + }); + + expect(started.type).toBe("tool.started"); + expect(started.input).toEqual({ name: "query" }); + expect(completed.type).toBe("tool.completed"); + expect(completed.result).toBe("Skill loaded"); + }); + it("accepts delete, typing, interactions, confirmation, and thread creation", () => { const callbacks = [ { ...callbackFields, type: "message.delete", messageId: "message-1" }, diff --git a/packages/contracts/src/hermesBridge.ts b/packages/contracts/src/hermesBridge.ts index 30d6a191b46..483953505bb 100644 --- a/packages/contracts/src/hermesBridge.ts +++ b/packages/contracts/src/hermesBridge.ts @@ -232,6 +232,28 @@ export const HermesBridgeEditMessageRequest = openStruct({ }); export type HermesBridgeEditMessageRequest = typeof HermesBridgeEditMessageRequest.Type; +const ToolCallbackFields = { + ...CallbackFields, + ...DestinationFields, + toolCallId: TrimmedNonEmptyString, + name: TrimmedNonEmptyString, + input: Schema.Unknown, +} as const; + +export const HermesBridgeToolStartedRequest = openStruct({ + ...ToolCallbackFields, + type: Schema.Literal("tool.started"), +}); +export type HermesBridgeToolStartedRequest = typeof HermesBridgeToolStartedRequest.Type; + +export const HermesBridgeToolCompletedRequest = openStruct({ + ...ToolCallbackFields, + type: Schema.Literal("tool.completed"), + result: Schema.Unknown, + isError: Schema.Boolean, +}); +export type HermesBridgeToolCompletedRequest = typeof HermesBridgeToolCompletedRequest.Type; + export const HermesBridgeDeleteMessageRequest = openStruct({ ...CallbackFields, ...DestinationFields, @@ -307,6 +329,8 @@ export type HermesBridgeThreadCreateRequest = typeof HermesBridgeThreadCreateReq export const HermesBridgeHermesToT3Request = Schema.Union([ HermesBridgeSendMessageRequest, HermesBridgeEditMessageRequest, + HermesBridgeToolStartedRequest, + HermesBridgeToolCompletedRequest, HermesBridgeDeleteMessageRequest, HermesBridgeTypingRequest, HermesBridgeTurnCompleteRequest,