Реализация интеграции с Pi Agent: разбор сообщений, повторы и отмена
Реализация интеграции с Pi Agent: разбор сообщений, повторы и отмена
При интеграции с AI-агентом в форме CLI нельзя обойти три вещи: как преобразовать его приватный поток событий в стабильные сообщения, кто отвечает за повторы после сбоев, и как корректно остановить процесс, когда пользователь нажимает «отмена». По сути эти три вопроса сводятся к «чёткому разделению ответственности», но только при реальной реализации понимаешь, насколько это глубоко.
Предыстория
В последнее время я работаю над проектом AI-помощника для кодирования, и одним из агентов, с которым нужно интегрироваться, является pi. Это TUI/CLI coding-агент, который при запуске выводит JSON-события построчно в stdout. Звучит просто — запускаем процесс, читаем вывод, парсим — но на практике вы обнаружите, что «интеграция с agent CLI» и «интеграция с обычным CLI» — это совершенно разные вещи.
С обычным CLI вы просто читаете stdout, получаете код выхода, и всё. Но agent CLI обладает тремя особенностями, которые вызывают головную боль:
Во-первых, его поток событий использует частный протокол. turn_start, session, message_update, message_end, turn_end, agent_end — это события, определённые самим pi, не какие-то отраслевые стандарты. Каждый верхний слой, который хочет их потреблять, должен обрабатывать их по-своему, что приводит к утечке внутренних деталей pi повсюду. Это как смотреть на человека издалека: вам кажется, что вы его разглядели, но на самом деле вы видите только то, что он хочет вам показать.
Во-вторых, семантика сбоев очень неоднозначна. Агент может работать, потом сеть может дернуться, модель может быть ограничена по запросам, процесс может упасть — в таких случаях нужно ли делать повторы? Где? Не нарушат ли повторы уже наполовину записанное состояние сессии? Это архитектурное решение, не такое, которое можно решить простым циклом for.
В-третьих, он долгий и прерываемый. Один turn может работать десятки секунд или даже несколько минут, и пользователь может захотеть отменить его в любой момент. При отмене процесс не должен стать сиротой, вызовы инструментов не должны оставаться в недоделанном состоянии, а уже выведенный контент нельзя терять. Здесь воды гораздо больше, чем кажется.
Чтобы решить эти проблемы, мы потратили время на упорядочивание пути интеграции. Ниже мы обсудим это подробно, но сначала раскроем главный секрет: настоящая сложность не в «запуске процесса», а в «чётком разделении ответственности».
О HagiCode
Решение, описанное в этой статье, происходит из проекта HagiCode — AI-помощника для кодирования с поддержкой нескольких моделей и нескольких agent CLI-бэкендов. GitHub-репозиторий: HagiCode-org/site, не забудьте поставить звёздочку. Весь код и все проблемы, с которыми мы столкнулись, действительно работают в этом проекте. На самом деле это написано просто для того, чтобы оставить себе память о пройденном пути.
Общая слоистая архитектура
HagiCode разделяет интеграцию AI-возможностей на два слоя:
- Нижний слой —
Hagicode.Libs, предоставляет повторно используемый provider-примитивICliProvider<TOptions>, специально отвечающий за «запуск CLI-агента и нормализацию его вывода в общий поток сообщений». - Верхний слой —
hagicode-core, предоставляет project-level thin adapterIAIProvider, отвечающий за «перевод бизнес-запросов в параметры provider, потребление общего потока сообщений и вывод унифицированных потоковых чанков».
Интеграция pi идёт по этому пути. Нижний слой PiProvider запускает процесс pi, читает поток JSON-событий, нормализует в общие сообщения; верхний слой PiCliProvider переводит AIRequest в PiOptions, потребляет CliMessage, выводит AIStreamingChunk.
Эти три вещи — разбор сообщений, повторы, отмена — 分别 находятся в трёх разных местах: PiJsonEventMapper, архивном предложении, и CliProcessManager. Ниже обсудим каждое из них.
Разбор сообщений: как приватные события pi становятся общими сообщениями
pi выводит JSON-события построчно в режиме --mode json --print. Этот набор событий является приватным для pi, и его нельзя напрямую передавать на верхний слой, иначе каждый потребитель будет связан с внутренними деталями pi, и при обновлении pi структуры событий весь проект будет следовать за изменениями. Такая утечка похожа на то, как будто бы вы написали свои мысли на лице — другим смотреть утомительно, и вам самому не особо комфортно.
Мы используем PiJsonEventMapper для выполнения слоя перевода, нормализуя события pi в общие CliMessage. CliMessage определяется в HagiCode.Libs.Core/Transport/CliMessage.cs, структура очень простая — это record (Type, Content). Соотношение примерно следующее:
| pi событие | Общее сообщение | Назначение |
|---|---|---|
session | session.started / session.resumed | Жизненный цикл сессии |
message_update (тип text) | assistant | Инкремент поточного текста |
message_update (тип thinking) | assistant.thought | Цепочка рассуждений |
message_update (тип tool) | tool.call / tool.update | Инициирование вызова инструмента |
message_end / turn_end (toolResult) | tool.completed / tool.failed | Результат инструмента |
turn_end / agent_end | terminal.completed | Конец текущего раунда |
| Ненулевой выход / ошибка парсинга | terminal.failed | Конечное состояние сбоя |
Эта таблица только для быстрого ознакомления, в ней есть две ключевые техники, которые мы нашли после наступания на грабли и которые стоит обсудить подробнее.
Техника 1: преобразование cumulative snapshot в delta
Это самое место, где можно «сломать машину». Событие message_update от pi передаёт не инкремент, а полный накопительный текст — каждый токен, он заново отправляет «полный текст до текущего момента».
Если вы напрямую пересылаете полученный контент во фронтенд, пользователь увидит повторяющийся контент: первая строка «你», вторая «你好», третья «你好,», четвёртая «你好,世»… фронтенд будет считать это четырьмя независимыми выводами. На самом деле повторение — это интересно один раз, а десять раз — уже просто надоедает.
Решение — сравнение префиксов, вычисление реального инкремента:
// Ключевой момент: pi отправляет накопительный snapshot, не инкремент// Сравнение префиксов позволяет вытащить инкремент, иначе фронтенд увидит повторяющийся контентif (text.StartsWith(_lastAssistantTextSnapshot, StringComparison.Ordinal)){ var delta = text[_lastAssistantTextSnapshot.Length..]; _lastAssistantTextSnapshot = text; return delta.Length == 0 ? null : delta;}Здесь есть ещё скрытая ловушка: повторное воспроизведение префикса между turn. Когда pi заканчивает вызов инструмента и assistant снова начинает говорить, он снова выведет этот текст с самого начала. Если вы запомните только один глобальный snapshot, вы примете воспроизведённый контент за инкремент, что приведёт к повторению после вызова инструмента. В PiProviderTests есть специальный тестовый случай ExecuteAsync_deduplicates_replayed_assistant_prefix_after_tool_turns, покрывающий этот сценарий. Иными словами, snapshot до и после вызова инструмента нужно выровнять при обработке, нельзя действовать по отдельности.
Техника 2: thinking нужно буферизировать до конца turn, а потом отправлять
Цепочка рассуждений (thinking) не должна отправляться наружу при каждом полученном токене. pi в середине вызова инструмента может вставить кучу фрагментов thinking, и при прямой пересылке порядок потока превратится в кашу —一会儿 assistant текст,一会儿 thinking фрагменты,一会儿 tool.call. Это имеет смысл? На самом деле никакого смысла нет, просто увеличивает хаос.
Наш подход: при получении события thinking сначала помещаем его в BufferThinkingSnapshot для временного хранения, и только когда message_end или turn_end и stopReason != "toolUse", мы единообразно DrainBufferedThinkingMessages. Так фрагменты thinking в середине вызова инструмента не загрязнят основной поток, а в конце turn будет выдан полный процесс рассуждений за раз.
Отказоустойчивость: плохие строки не должны разрушать поток
agent CLI не идеальная система из учебника, иногда он может вывести строку не в JSON, или JSON без поля type. Если вы здесь выбросите исключение, весь поток погибнет, и пользователь ничего не увидит. Ведь реальный мир всегда несовершенен, кто может гарантировать, что каждая строка будет аккуратной?
Наша стратегия: при любой неудаче парсинга строки мы не прерываем поток, а собираем в _invalidOutputLines. После завершения процесса в Complete() мы собираем эти «плохие строки» в диагностический текст terminal.failed. Так пользователь при просмотре ошибки может напрямую увидеть, какую ерунду pi вывел, вместо сухого «parse error».
Повторы: если слой provider не делает, кто делает?
Это самая большая ловушка во всей интеграции. Интуитивно кажется, что «интеграция с CLI должна включать повторы», но HagiCode в архивном предложении активно удалил всю автоматическую логику повторов из слоя provider. Предложение называется remove-provider-auto-retry-support.
Почему не делать автоматические повторы
Фон предложения написан очень прямо. Логика повторов изначально была разбросана в двух местах: в Hagicode.Libs была одна (replay fresh-runtime в стиле OpenCode), в hagicode-core ещё одна (ProviderErrorAutoRetryCoordinator). Обе стороны делали по-своему, что привело к тому, что «делать повторы или нет» стало скрытым поведением внутри provider, которое тайно изменяет момент сбоя, способ продолжения сессии и поток состояния чата.
Подумайте об этом — голова болит: пользователь отправляет сообщение, provider внутри сам делает три повтора, первые два раза не удаётся, в третий раз удаётся. Верхний слой вообще не знает, что произошло посередине, состояние сессии, счётчик токенов, прогресс UI — всё не совпадает. Такое скрытое поведение — это хронический яд в архитектуре.
Поэтому граница была сужена до одного предложения:
provider сужается до семантики единственной попытки, вызывающая сторона должна рассматривать состояние без повтора как нормальный результат единственного выполнения.
Как это выглядит на PiProvider
На уровне кода это три вещи:
- В
PiOptionsнет ни одного поля, связанного с повторами — нетmaxAttempts, нетretryDelay, нетretryClassifier. ExecuteAsyncзавершается после однократного запуска процесса pi, при сбое напрямую выдаётterminal.failed.- Те классификаторы, которые раньше служили автоматическим повторам (
ClaudeCodeRetryableTerminalFailureClassifier,CodexRetryableTerminalFailureClassifierи т.д.), если они служат только автоматическим повторам, все удалены из активного пути.
Но обратите внимание, возможность повторов не исчезла, просто поднялась выше. В предложении прямо написано «оставить стабильную границу для последующего единого управления повторами на более высоком уровне». DTO, нормализация, сериация, round-trip страницы настроек фронтенда для конфигурации providerErrorAutoRetry — всё сохранено, только это больше не управляет выполнением provider. Ведь некоторые вещи не исчезают, просто они сохраняются в другой форме.
А как делать повторы
Если вы хотите добавить повторы поверх pi, правильный способ — делать это на стороне вызывающего PiCliProvider — например, в вашем слое оркестрации сессий (в HagiCode это SessionGrain из Orleans, на фронтенде может быть слой оркестрации чата). После получения terminal.failed вы сами решаете, можно ли повторять, сами определяете задержку и количество раз, и снова отправляете ExecuteAsync.
Минимально рабочий режим выглядит так:
// Логика повторов помещается на стороне вызывающего, не вставляйте обратно в PiProvider// Иначе это нарушит границу «единственной попытки», только что установленную в providerasync Task<AIResponse> ExecuteWithRetryAsync(AIRequest req, int maxAttempts, CancellationToken ct){ for (var attempt = 1; ; attempt++) { var response = await provider.ExecuteAsync(req, ct);
// При успехе или достижении лимита возвращаем if (response.FinishReason != FinishReason.Unknown || attempt >= maxAttempts) return response;
// Повторяем только для конечных сбоев, которые можно повторить (сеть, 5xx, падение процесса) // model rejected, auth failure这类 повторения бессмысленны, не повторяйте await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, attempt)), ct); }}Классификационная логика для «можно повторять» теперь не находится в provider, вызывающая сторона определяет её сама. Конфигурация providerErrorAutoRetry (maxAttempts, retryDelay, enabled) всё ещё может быть прочитана со страницы настроек фронтенда, но реальное управление повторами выполняет ваш слой оркестрации, а не PiProvider. Пожалуйста, повторите это трижды.
Отмена: проксирование токена + трёхступенчатая остановка
С отменой PiProvider почти ничего не реализует, полностью делегирует CliProcessManager, PiProvider отвечает только за две вещи: передать CancellationToken вниз и сделать уборку при исключении.
Полноканальное проксирование
Цепочка такая, передаётся до самого низа:
CancellationToken вызывающей стороны → PiCliProvider.StreamCoreAsync(cancellationToken) → PiProvider.ExecuteProcessAsync([EnumeratorCancellation] cancellationToken) → ReadLineAsync(cancellationToken) / WaitForExitAsync(cancellationToken) → при исключении _processManager.StopAsync(handle, CancellationToken.None)Обратите внимание на последнюю строку: при очистке используется CancellationToken.None, а не токен, переданный пользователем. Это деталь, но крайне важна.
Причина: токен пользователя уже отменён. Если вы используете этот уже отменённый токен для очистки, задача очистки будет немедленно отменена, процесс станет сиротой — pi всё ещё работает в фоновом режиме, никто его не собирает, CPU и память занимаются зря. Поэтому очистка должна использовать CancellationToken.None, чтобы убедиться, что действие очистки обязательно выполнится до конца. На самом деле это похоже на людей: некоторые вещи нужно правильно завершить после того, как они полностью остановятся, иначе останется куча мусора.
Трёхступенчатая прогрессивная остановка
CliProcessManager.StopProcessAsync — это трёхступенчатый прогрессивный процесс остановки, временные константы определены в начале файла:
// Терпение для плавной остановки: сначала дайте процессу время на завершениеprivate static readonly TimeSpan GracefulStopTimeout = TimeSpan.FromSeconds(2);// Терпение после принудительного kill, чтобы убедиться, что процесс действительно вышелprivate static readonly TimeSpan StopWaitTimeout = TimeSpan.FromSeconds(5);Три ступени таковы:
- Сигнал прерывания.
TryInterruptAsyncсначала пишет\u0003в stdin (это символ Ctrl+C), под Unix дополнительноkill -INT <pid>. Этот шаг предназначен для того, чтобы pi сам плавно завершился — он может почувствовать прерывание и завершить то, что пишет. - Плавное ожидание. Максимум ждём 2 секунды, смотрим, вышел ли процесс сам.
- Принудительный kill. Если ещё не вышел, напрямую
Process.Kill(entireProcessTree: true), убиваем всё дерево процессов вместе, затем максимум 5 секунд ждём, чтобы убедиться, что он действительно мёртв.
Зачем entireProcessTree: true? Потому что при запуске инструментов pi создаёт дочерние процессы — например, процесс локальной модели, на который маршрутизируется provider, или bash-подпроцесс. Если убить только родительский процесс, дочерние процессы станут сиротами и продолжат работать. Убивать всё дерево вместе — это чисто.
В Windows нет SIGINT, можно полагаться только на символ Ctrl+C, поэтому поведение между платформами будет отличаться, об этом нужно знать.
Уборка исключений в PiProvider
ExecuteProcessAsync в PiProvider при исключении из ReadLineAsync использует ExceptionDispatchInfo.Capture для временного хранения исключения, выходит из цикла, вызывает StopAsync для очистки процесса, а затем pendingException.Throw() для повторного выбрасывания исходного исключения на верхний слой.
Зачем временно хранить и потом выбрасывать? Потому что если выбросить напрямую, процесс не успевает быть собран, становится сиротой; если выбросить до StopAsync, логика очистки вообще не выполняется. Временно храня, сначала гарантируем, что процесс обязательно будет собран, а затем полностью сохраняем исходную семантику OperationCanceledException для вызывающей стороны — вызывающая сторона, получив это исключение, может определить «о, это пользователь отменил», а не «произошла ошибка».
Единый контракт для сбоев запуска
Есть ещё одна деталь, которую стоит упомянуть отдельно. Если процесс не запускается — например, исполняемый файл pi не существует, неправильные разрешения — PiProvider не выбрасывает исключение, а синтезирует сообщение terminal.failed, затем yield break.
Зачем так? Потому что если выбросить исключение, верхний потребитель должен обрабатывать две совершенно разные семантики: одна — «нормальное сообщение в процессе потокового потребления», другая — «исключение выброшено до начала потока». Это сделает await foreach потребителя особенно сложным для написания.
После унификации в «всегда сначала дать вам сообщение, затем закончить поток» логика потребителя становится одинаковой: при получении terminal.failed считаем сбой, при получении terminal.completed считаем успех, не нужно разветвляться с try/catch. Это маленькое, но важное архитектурное решение, которое стабилизирует контракт.
Практика: правильный способ потребления потока
Обратитесь к PiScenarioMessageReader в HagiCode (тестовый сценарий консоли в libs) и PiCliProvider.StreamCoreAsync в hagicode-core (thin adapter), потребитель примерно выглядит так:
await foreach (var message in provider.ExecuteAsync(options, prompt, cancellationToken)){ // 1. При сбое нужно сразу выйти коротким путём, не обрабатывать последующие сообщения if (NormalizedAcpCliAdapter.TryGetFailureMessage(message.Content, out var failure)) { yield return new AIStreamingChunk { Type = StreamingChunkType.Error, ErrorMessage = failure }; yield break; // после terminal.failed поток заканчивается }
// 2. assistant текст — это cumulative snapshot, нужно самому ещё раз вычислить инкремент if (message.Type == "assistant" && TryGetText(message.Content, out var text)) { var delta = ReconcileSnapshot(text); // сравнение префиксов if (!string.IsNullOrEmpty(delta)) yield return Chunk(delta); }
// 3. terminal.completed — единственный надёжный сигнал «окончания» if (message.Type == "terminal.completed") break;}Быстрая проверка распространённых ловушек
Соберём все ловушки, с которыми столкнулись, в таблицу, чтобы будущим людям было проще:
| Явление | Причина | Обработка |
|---|---|---|
| Фронтенд видит повторяющийся текст assistant | Не выполнено преобразование cumulative в delta | Используйте ReconcileAssistantTextSnapshot для сравнения префиксов |
| После отмены процесс всё ещё работает | При очистке использован уже отменённый токен | Используйте CancellationToken.None для очистки |
| Повторы не работают | Повторы написаны в PiProvider, но provider имеет семантику единственной попытки | Поднять на слой оркестрации вызывающей стороны |
| Сообщения об ошибках pi теряются | Не прочитано диагностическое поле terminal.failed | Полностью проксируйте text / invalid_output_lines / stderr |
| Фрагменты thinking получаются в середине вызова инструмента | Прямая пересылка событий thinking | Буферизовать до конца turn, затем DrainBufferedThinkingMessages |
Как проверить
В слое libs используется StubCliProcessManager для mock процессов, unit tests покрывают построение параметров, нормализацию событий, дедупликацию инкрементов, передачу сбоев и другую чистую логику. Реальный путь CLI opt-in с помощью переменной окружения HAGICODE_REAL_CLI_TESTS, используются реальные модели для запуска trip сценариев. В слое core PiCliProviderTests проверяет проекцию AIStreamingChunk thin adapter и session binding.
# Запустить unit tests, связанные с Pi, в репозитории Hagicode.Libsdotnet test --filter "FullyQualifiedName~PiProviderTests"
# Запустить реальные интеграционные тесты CLI (нужно установить pi локально)HAGICODE_REAL_CLI_TESTS=1 dotnet test --filter "FullyQualifiedName~PiProviderTests.RealCli"Заключение
Если соединить эти три вещи, ментальная модель интеграции с pi фактически сводится к одному предложению: позволить каждому слою делать только свою работу.
- Разбор сообщений доверить
PiJsonEventMapper: приватные события нормализуются в общиеCliMessage, cumulative snapshot превращается в delta, thinking буферизуется до конца turn. - Повторение доверить вызывающей стороне: provider делает единственную попытку, кто хочет повторять, делает сам на верхнем уровне, конфигурация сохраняется, но больше не управляет provider.
- Отмена доверить
CliProcessManager:CancellationTokenпроксируется по всему каналу, при очистке используетсяCancellationToken.None, трёхступенчатая прогрессивная остановка (сигнал прерывания → плавное ожидание → принудительный kill всего дерева процессов).
После того как эти границы чётко разграничены, интеграция нового agent CLI практически стала конвейерной работой — вам нужно только написать новый XxxProvider и XxxJsonEventMapper, а вся сквозная логика повторов, отмены, контрактов сообщений и обработки ошибок полностью переиспользуется. Это также фундаментальная причина, по которой HagiCode может одновременно поддерживать несколько agent CLI-бэкендов (claude code, codex, pi, gemini cli и т.д.), не превращаясь в хаос.
Ещё раз скажем самую важную границу: не добавляйте повторы в слой provider. Как только вы поймёте это, интеграция с agent CLI пройдёт большую часть пути…
Резюме
Возвращаясь к теме «Реализация интеграции с Pi Agent: разбор сообщений, повторы и отмена», то, что действительно стоит неоднократно подтверждать, — это не разрозненные техники, а ясно ли видны ограничения, границы реализации и инженерные компромиссы.
Пока вы превращаете основания для суждений в статье в стабильные контрольные пункты, при дальнейшем столкновении с подобными проблемами вы сможете быстрее принимать надёжные решения.
开始使用 HagiCode
一次安装,几分钟上手
HagiCode for Windows 在 Microsoft Store 免费提供。打开商店即可安装并保持更新;也可以先对比各版本与定价,再决定从哪个渠道开始。