From 2ddeed6dcb6b5d1d70e768d3c43b348902bade22 Mon Sep 17 00:00:00 2001 From: briandevans <252620095+briandevans@users.noreply.github.com> Date: Tue, 4 Aug 2026 17:08:40 -0700 Subject: [PATCH 1/2] fix(gateway): admit exactly one /update per profile MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `/update` wrote `.update_pending.json` and spawned a detached `hermes update --gateway` with no admission control of any kind. Two concurrent invocations — a user double-tapping because the update runs for minutes with no acknowledgement, or two platforms on one multiplexed gateway both triggering it — each wrote the marker and each spawned an updater against the same checkout and virtualenv. The second write also replaced the first requester's routing metadata, so that user never learned their update finished. Both handlers also stage through the same `.update_pending.tmp` path, so the losing rename raises FileNotFoundError out of the handler. Reserve the profile-wide slot with os.open(O_CREAT | O_EXCL) before any routing metadata is written and before the updater is spawned. That collapses observe-and-claim into a single atomic syscall, so exactly one caller can ever win. The losing caller gets an "update already running" reply and still gets `_schedule_update_notification_watch()`, so it learns the outcome of the update that is actually running. The reservation carries a TTL: an updater killed before it can write its exit code (host reboot, OOM kill) would otherwise wedge /update for the lifetime of the profile. Any failure between the claim and the spawn releases the slot immediately so the user can retry. Supersedes #15539, which identified this collision. That guard patched `gateway/run.py`, where the handler no longer lives after 619bd7827, and used a check-then-write `if pending_path.exists(): return` that both callers can pass; its tests pre-seeded the marker, so they never exercised the check-to-create gap. --- gateway/slash_commands.py | 104 +++++++++++- locales/af.yaml | 1 + locales/ar.yaml | 1 + locales/de.yaml | 1 + locales/en.yaml | 1 + locales/es.yaml | 1 + locales/fr.yaml | 1 + locales/ga.yaml | 1 + locales/hu.yaml | 1 + locales/it.yaml | 1 + locales/ja.yaml | 1 + locales/ko.yaml | 1 + locales/pt.yaml | 1 + locales/ru.yaml | 1 + locales/tr.yaml | 1 + locales/uk.yaml | 1 + locales/zh-hant.yaml | 1 + locales/zh.yaml | 1 + tests/gateway/test_update_command.py | 245 +++++++++++++++++++++++++++ 19 files changed, 362 insertions(+), 4 deletions(-) diff --git a/gateway/slash_commands.py b/gateway/slash_commands.py index e49891749b5f..dab490ac9b95 100644 --- a/gateway/slash_commands.py +++ b/gateway/slash_commands.py @@ -55,6 +55,14 @@ # its worker thread. (#35994) _RESET_CLEANUP_TIMEOUT_S = 30.0 +# How long a /update reservation stays valid before another /update may take it +# over (see _claim_update_slot). An updater that is killed without writing its +# exit code — host reboot, OOM kill — leaves its marker behind forever, so the +# reservation needs a ceiling or /update wedges for the lifetime of the profile. +# Sized above _watch_update_progress's own 1800s watch timeout so a live update +# that is still being watched is never stolen. +_UPDATE_RESERVATION_TTL_S = 3600.0 + def _clean_str(value: Any) -> str: """Strip and return a non-empty string value, or empty string.""" @@ -123,6 +131,62 @@ def _home_thread_from_source(source) -> Optional[str]: return str(thread_id) +def _claim_update_slot(pending_path: Path, ttl_seconds: float) -> bool: + """Atomically reserve the profile-wide ``/update`` slot. + + ``/update`` is profile-global: it rewrites the checkout and the virtualenv + every session on the host shares, and the ``.update_*`` marker files are a + single-slot mailbox holding one requester's routing metadata. Two updaters + must therefore never run at once. + + A ``if pending_path.exists(): return`` guard cannot provide that. Two + handlers can both observe the marker missing and both go on to create it, + so the second one's metadata silently replaces the first one's and a second + updater is spawned against the same checkout. ``os.open`` with + ``O_CREAT | O_EXCL`` collapses the observe-and-create into a single atomic + syscall, so exactly one caller can ever win. + + Returns ``True`` when this caller now owns the slot (the marker exists and + is empty, ready for the caller to fill in with its routing metadata), and + ``False`` when another update already holds it. + """ + def _create_exclusive() -> bool: + try: + fd = os.open(pending_path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + except FileExistsError: + return False + os.close(fd) + return True + + if _create_exclusive(): + return True + + # The slot is held. Take it over only when it is provably abandoned: a host + # reboot or a SIGKILLed updater leaves a marker that nothing will ever + # clean up, and a permanently wedged /update would be a worse bug than the + # collision this guard exists to close. + try: + age = time.time() - pending_path.stat().st_mtime + except OSError: + # Vanished between the failed create and the stat — the previous update + # finished in that window. Retry once; if we lose again, the loser path + # is still correct. + return _create_exclusive() + + if age < ttl_seconds: + return False + + logger.warning( + "Reclaiming abandoned /update reservation %s (age %.0fs > %.0fs)", + pending_path, age, ttl_seconds, + ) + try: + pending_path.unlink() + except OSError: + return False + return _create_exclusive() + + class GatewaySlashCommandsMixin: """In-session slash-command handlers for GatewayRunner.""" @@ -5659,8 +5723,31 @@ async def _handle_update_command(self, event: MessageEvent) -> str: return t("gateway.update.hermes_cmd_not_found") pending_path = _hermes_home / ".update_pending.json" + claimed_path = _hermes_home / ".update_pending.claimed.json" output_path = _hermes_home / ".update_output.txt" exit_code_path = _hermes_home / ".update_exit_code" + + # Admission control. Reserve the profile-wide update slot BEFORE any + # routing metadata is written and BEFORE the updater is spawned, so two + # concurrent /update invocations (a double tap while the minutes-long + # update runs without an ack, or two platforms on one multiplexed + # gateway) cannot both start an updater against the same checkout and + # venv, with the second one's metadata clobbering the first requester's + # so they never learn their update finished. + # + # `.update_pending.claimed.json` is the notifier's transient rename of + # the marker (see _send_update_notification); while it exists an update + # is in flight, so reject. Checking it is not a check-then-create race: + # it can only ever produce an extra rejection, never an extra spawn, + # and the exclusive create below is what actually decides the winner. + if claimed_path.exists() or not _claim_update_slot( + pending_path, _UPDATE_RESERVATION_TTL_S + ): + # The loser still gets the outcome: the watcher is profile-wide and + # reports the result into the chat that started the update. + self._schedule_update_notification_watch() + return t("gateway.update.already_running") + session_key = self._session_key_for_source(event.source) pending = { "platform": event.source.platform.value, @@ -5674,10 +5761,6 @@ async def _handle_update_command(self, event: MessageEvent) -> str: pending["thread_id"] = event.source.thread_id if event.message_id: pending["message_id"] = event.message_id - _tmp_pending = pending_path.with_suffix(".tmp") - _tmp_pending.write_text(json.dumps(pending), encoding="utf-8") - _tmp_pending.replace(pending_path) - exit_code_path.unlink(missing_ok=True) # Spawn `hermes update --gateway` detached so it survives gateway restart. # --gateway enables file-based IPC for interactive prompts (stash @@ -5704,6 +5787,17 @@ async def _handle_update_command(self, event: MessageEvent) -> str: # so the simplest correct thing is: launch an inline Python helper # that runs the command and writes both outputs. try: + # Fill the reservation in with this requester's routing metadata. + # Write-to-temp + rename keeps the marker a complete JSON document + # for the readers in gateway/run.py, which poll it concurrently. + # Everything from the claim onwards must release the slot on + # failure, or a /update that never actually started would block the + # next one until the reservation TTL expires. + _tmp_pending = pending_path.with_suffix(".tmp") + _tmp_pending.write_text(json.dumps(pending), encoding="utf-8") + _tmp_pending.replace(pending_path) + exit_code_path.unlink(missing_ok=True) + if sys.platform == "win32": import textwrap from hermes_cli._subprocess_compat import windows_detach_popen_kwargs @@ -5764,6 +5858,8 @@ async def _handle_update_command(self, event: MessageEvent) -> str: start_new_session=True, ) except Exception as e: + # Release the reservation so the user can retry immediately. + pending_path.with_suffix(".tmp").unlink(missing_ok=True) pending_path.unlink(missing_ok=True) exit_code_path.unlink(missing_ok=True) return t("gateway.update.start_failed", error=e) diff --git a/locales/af.yaml b/locales/af.yaml index 6658c2da20d0..e6a3a7a82fe0 100644 --- a/locales/af.yaml +++ b/locales/af.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Kon nie die `hermes`-opdrag vind nie. Hermes loop, maar die opdateeropdrag kon nie die uitvoerbare lêer op PATH of via die huidige Python-vertolker vind nie. Probeer `hermes update` met die hand in jou terminale uitvoer." start_failed: "✗ Kon nie opdatering begin nie: {error}" starting: "⚕ Begin Hermes-opdatering… Ek sal vordering hier stroom." + already_running: "⚕ 'n Hermes-opdatering loop reeds vir hierdie profiel. Ek sal die resultaat rapporteer in die geselsie wat dit begin het." usage: rate_limits: "⏱️ **Tariefperke:** {state}" diff --git a/locales/ar.yaml b/locales/ar.yaml index 9d98db08eb9d..9e3b2c5151f7 100644 --- a/locales/ar.yaml +++ b/locales/ar.yaml @@ -370,6 +370,7 @@ gateway: hermes_cmd_not_found: "✗ تعذّر تحديد موقع أمر `hermes`. Hermes يعمل، لكن أمر التحديث لم يجد الملف التنفيذي على PATH أو عبر مُفسّر Python الحالي. جرّب تشغيل `hermes update` يدويًا في طرفيتك." start_failed: "✗ فشل بدء التحديث: {error}" starting: "⚕ جارٍ بدء تحديث Hermes… سأبثّ التقدّم هنا." + already_running: "⚕ يوجد بالفعل تحديث Hermes قيد التشغيل لهذا الملف الشخصي. سأبلّغ بالنتيجة في المحادثة التي بدأته." usage: rate_limits: "⏱️ **حدود المعدّل:** {state}" diff --git a/locales/de.yaml b/locales/de.yaml index 52d14ad237bd..f383f73fa586 100644 --- a/locales/de.yaml +++ b/locales/de.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Der Befehl `hermes` konnte nicht gefunden werden. Hermes läuft, aber der Update-Befehl konnte das ausführbare Programm weder im PATH noch über den aktuellen Python-Interpreter finden. Versuchen Sie, `hermes update` manuell im Terminal auszuführen." start_failed: "✗ Update konnte nicht gestartet werden: {error}" starting: "⚕ Hermes-Update wird gestartet… Ich streame den Fortschritt hier." + already_running: "⚕ Für dieses Profil läuft bereits ein Hermes-Update. Ich melde das Ergebnis in dem Chat, der es gestartet hat." usage: rate_limits: "⏱️ **Ratenlimits:** {state}" diff --git a/locales/en.yaml b/locales/en.yaml index d1069d7eac4d..2653e7d61e7c 100644 --- a/locales/en.yaml +++ b/locales/en.yaml @@ -381,6 +381,7 @@ gateway: hermes_cmd_not_found: "✗ Could not locate the `hermes` command. Hermes is running, but the update command could not find the executable on PATH or via the current Python interpreter. Try running `hermes update` manually in your terminal." start_failed: "✗ Failed to start update: {error}" starting: "⚕ Starting Hermes update… I'll stream progress here." + already_running: "⚕ A Hermes update is already running for this profile. I'll report the result in the chat that started it." usage: rate_limits: "⏱️ **Rate Limits:** {state}" diff --git a/locales/es.yaml b/locales/es.yaml index 16c030d616ce..5e1c54babfca 100644 --- a/locales/es.yaml +++ b/locales/es.yaml @@ -365,6 +365,7 @@ gateway: hermes_cmd_not_found: "✗ No se pudo localizar el comando `hermes`. Hermes está en ejecución, pero el comando de actualización no encontró el ejecutable en PATH ni a través del intérprete de Python actual. Intenta ejecutar `hermes update` manualmente en tu terminal." start_failed: "✗ No se pudo iniciar la actualización: {error}" starting: "⚕ Iniciando la actualización de Hermes… Transmitiré el progreso aquí." + already_running: "⚕ Ya hay una actualización de Hermes en curso para este perfil. Informaré del resultado en el chat que la inició." usage: rate_limits: "⏱️ **Límites de tasa:** {state}" diff --git a/locales/fr.yaml b/locales/fr.yaml index 77869000e24f..2425e67b021c 100644 --- a/locales/fr.yaml +++ b/locales/fr.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Impossible de localiser la commande `hermes`. Hermes est en cours d'exécution, mais la commande de mise à jour n'a pas pu trouver l'exécutable dans le PATH ni via l'interpréteur Python actuel. Essayez d'exécuter `hermes update` manuellement dans votre terminal." start_failed: "✗ Échec du démarrage de la mise à jour : {error}" starting: "⚕ Démarrage de la mise à jour Hermes… Je diffuserai la progression ici." + already_running: "⚕ Une mise à jour Hermes est déjà en cours pour ce profil. Je signalerai le résultat dans la conversation qui l'a lancée." usage: rate_limits: "⏱️ **Limites de débit :** {state}" diff --git a/locales/ga.yaml b/locales/ga.yaml index 1e9231a40525..22783a5f10a8 100644 --- a/locales/ga.yaml +++ b/locales/ga.yaml @@ -372,6 +372,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Níorbh fhéidir an t-ordú `hermes` a aimsiú. Tá Hermes ag rith, ach níorbh fhéidir leis an ordú nuashonraithe an inrite a aimsiú ar PATH ná tríd an léirmhínitheoir Python reatha. Bain triail as `hermes update` a rith de láimh i do theirminéal." start_failed: "✗ Theip ar nuashonrú a thosú: {error}" starting: "⚕ Ag tosú nuashonrú Hermes… Cuirfidh mé an dul chun cinn ar shruth anseo." + already_running: "⚕ Tá nuashonrú Hermes ar siúl cheana féin don phróifíl seo. Tuairisceoidh mé an toradh sa chomhrá a thosaigh é." usage: rate_limits: "⏱️ **Teorainneacha Ráta:** {state}" diff --git a/locales/hu.yaml b/locales/hu.yaml index fa4705deb28e..a56a1beeefd3 100644 --- a/locales/hu.yaml +++ b/locales/hu.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Nem sikerült megtalálni a `hermes` parancsot. A Hermes fut, de a frissítőparancs nem találta a futtatható fájlt a PATH-on vagy a jelenlegi Python interpreteren keresztül. Próbáld futtatni a `hermes update` parancsot manuálisan a terminálban." start_failed: "✗ Nem sikerült elindítani a frissítést: {error}" starting: "⚕ Hermes frissítés indítása… A folyamatot itt fogom közvetíteni." + already_running: "⚕ Ehhez a profilhoz már fut egy Hermes frissítés. Az eredményt abban a beszélgetésben jelentem, amelyik elindította." usage: rate_limits: "⏱️ **Sebességkorlátok:** {state}" diff --git a/locales/it.yaml b/locales/it.yaml index 654998047370..24f595a20149 100644 --- a/locales/it.yaml +++ b/locales/it.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Impossibile localizzare il comando `hermes`. Hermes è in esecuzione, ma il comando di aggiornamento non ha trovato l'eseguibile nel PATH o tramite l'interprete Python attuale. Prova a eseguire `hermes update` manualmente nel terminale." start_failed: "✗ Avvio dell'aggiornamento non riuscito: {error}" starting: "⚕ Avvio dell'aggiornamento di Hermes… mostrerò qui i progressi in streaming." + already_running: "⚕ Un aggiornamento di Hermes è già in corso per questo profilo. Riporterò il risultato nella chat che lo ha avviato." usage: rate_limits: "⏱️ **Limiti di frequenza:** {state}" diff --git a/locales/ja.yaml b/locales/ja.yaml index 17b66572b419..cb6908202da5 100644 --- a/locales/ja.yaml +++ b/locales/ja.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ `hermes` コマンドが見つかりません。Hermes は実行中ですが、更新コマンドは PATH 上にも現在の Python インタープリタ経由でも実行可能ファイルを見つけられませんでした。ターミナルで `hermes update` を手動で実行してみてください。" start_failed: "✗ 更新の開始に失敗しました: {error}" starting: "⚕ Hermes の更新を開始しています… 進捗をここにストリーミングします。" + already_running: "⚕ このプロファイルでは Hermes の更新がすでに実行中です。結果は更新を開始したチャットに報告します。" usage: rate_limits: "⏱️ **レート制限:** {state}" diff --git a/locales/ko.yaml b/locales/ko.yaml index cdc43eb6d0b3..37bc4008a5d1 100644 --- a/locales/ko.yaml +++ b/locales/ko.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ `hermes` 명령을 찾을 수 없습니다. Hermes는 실행 중이지만 PATH나 현재 Python 인터프리터를 통해 실행 파일을 찾을 수 없습니다. 터미널에서 `hermes update`를 직접 실행해 보세요." start_failed: "✗ 업데이트 시작 실패: {error}" starting: "⚕ Hermes 업데이트를 시작합니다… 진행 상황을 여기에 스트리밍하겠습니다." + already_running: "⚕ 이 프로필에서는 이미 Hermes 업데이트가 실행 중입니다. 결과는 업데이트를 시작한 채팅에 보고하겠습니다." usage: rate_limits: "⏱️ **요청 제한:** {state}" diff --git a/locales/pt.yaml b/locales/pt.yaml index 0cbf524b6853..e42d2cd6f461 100644 --- a/locales/pt.yaml +++ b/locales/pt.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Não foi possível localizar o comando `hermes`. O Hermes está em execução, mas o comando de atualização não conseguiu encontrar o executável no PATH nem através do interpretador Python atual. Tenta executar `hermes update` manualmente no teu terminal." start_failed: "✗ Falha ao iniciar a atualização: {error}" starting: "⚕ A iniciar a atualização do Hermes… Vou transmitir o progresso aqui." + already_running: "⚕ Já está em curso uma atualização do Hermes para este perfil. Vou comunicar o resultado na conversa que a iniciou." usage: rate_limits: "⏱️ **Limites de taxa:** {state}" diff --git a/locales/ru.yaml b/locales/ru.yaml index a0d8eb508c61..ca89b0161393 100644 --- a/locales/ru.yaml +++ b/locales/ru.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Не удалось найти команду `hermes`. Hermes запущен, но команда обновления не нашла исполняемый файл в PATH или через текущий интерпретатор Python. Попробуйте выполнить `hermes update` вручную в терминале." start_failed: "✗ Не удалось запустить обновление: {error}" starting: "⚕ Запуск обновления Hermes… Я буду транслировать прогресс сюда." + already_running: "⚕ Для этого профиля уже выполняется обновление Hermes. Я сообщу результат в чат, из которого оно было запущено." usage: rate_limits: "⏱️ **Ограничения скорости:** {state}" diff --git a/locales/tr.yaml b/locales/tr.yaml index f39605c4b41e..ad5f3fa424a7 100644 --- a/locales/tr.yaml +++ b/locales/tr.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ `hermes` komutu bulunamadı. Hermes çalışıyor, ancak güncelleme komutu yürütülebilir dosyayı PATH'te veya mevcut Python yorumlayıcısı aracılığıyla bulamadı. Terminalde `hermes update` komutunu manuel olarak çalıştırmayı deneyin." start_failed: "✗ Güncelleme başlatılamadı: {error}" starting: "⚕ Hermes güncellemesi başlatılıyor… İlerlemeyi buraya akıtacağım." + already_running: "⚕ Bu profil için zaten bir Hermes güncellemesi çalışıyor. Sonucu, güncellemeyi başlatan sohbette bildireceğim." usage: rate_limits: "⏱️ **Hız Sınırları:** {state}" diff --git a/locales/uk.yaml b/locales/uk.yaml index 08f40a849267..01aaa51921b6 100644 --- a/locales/uk.yaml +++ b/locales/uk.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ Не вдалося знайти команду `hermes`. Hermes запущено, але команда оновлення не знайшла виконуваний файл у PATH або через поточний інтерпретатор Python. Спробуйте виконати `hermes update` вручну у вашому терміналі." start_failed: "✗ Не вдалося запустити оновлення: {error}" starting: "⚕ Запуск оновлення Hermes… Я транслюватиму прогрес сюди." + already_running: "⚕ Для цього профілю вже виконується оновлення Hermes. Я повідомлю результат у чат, з якого його було запущено." usage: rate_limits: "⏱️ **Обмеження швидкості:** {state}" diff --git a/locales/zh-hant.yaml b/locales/zh-hant.yaml index c08ee22d68df..100cd9f8248f 100644 --- a/locales/zh-hant.yaml +++ b/locales/zh-hant.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ 找不到 `hermes` 指令。Hermes 正在執行,但更新指令無法在 PATH 上或透過目前的 Python 解譯器找到執行檔。請嘗試在終端機中手動執行 `hermes update`。" start_failed: "✗ 啟動更新失敗:{error}" starting: "⚕ 正在啟動 Hermes 更新…… 進度將在此處顯示。" + already_running: "⚕ 此設定檔已有一個 Hermes 更新正在執行。我會在發起該更新的聊天中回報結果。" usage: rate_limits: "⏱️ **速率限制:** {state}" diff --git a/locales/zh.yaml b/locales/zh.yaml index 2c8aafdae3ca..c46bf4f7bf0c 100644 --- a/locales/zh.yaml +++ b/locales/zh.yaml @@ -368,6 +368,7 @@ Future messages in this room will use that transcript until `/reset` or another hermes_cmd_not_found: "✗ 无法找到 `hermes` 命令。Hermes 正在运行,但更新命令无法在 PATH 上或通过当前 Python 解释器找到可执行文件。请尝试在终端中手动运行 `hermes update`。" start_failed: "✗ 启动更新失败:{error}" starting: "⚕ 正在启动 Hermes 更新…… 进度将在此处显示。" + already_running: "⚕ 该配置文件已有一个 Hermes 更新正在运行。我会在发起该更新的聊天中报告结果。" usage: rate_limits: "⏱️ **速率限制:** {state}" diff --git a/tests/gateway/test_update_command.py b/tests/gateway/test_update_command.py index 22cc9cd419b5..11db21f0fc83 100644 --- a/tests/gateway/test_update_command.py +++ b/tests/gateway/test_update_command.py @@ -4,7 +4,11 @@ the _send_update_notification startup hook (sends results after restart). """ +import asyncio import json +import os +import threading +import time from pathlib import Path from unittest.mock import patch, MagicMock, AsyncMock @@ -178,6 +182,247 @@ def which_no_setsid(x): assert "Starting Hermes update" in result +# --------------------------------------------------------------------------- +# Concurrent /update admission control +# --------------------------------------------------------------------------- + + +def _fake_checkout(tmp_path): + """Build a fake checkout + profile home for the /update pre-flight. + + ``_handle_update_command`` resolves ``project_root`` from + ``gateway/slash_commands.py``'s own ``__file__``, so the fake tree mirrors + that layout and carries a ``.git`` directory. Returns + ``(slash_commands_file, hermes_home)``. + """ + root = tmp_path / "project" + (root / "gateway").mkdir(parents=True) + (root / ".git").mkdir() + slash_commands_file = root / "gateway" / "slash_commands.py" + slash_commands_file.touch() + hermes_home = tmp_path / "hermes" + hermes_home.mkdir() + return str(slash_commands_file), hermes_home + + +class TestUpdateAdmissionControl: + """``/update`` must admit exactly one updater per profile. + + ``hermes update`` rewrites the checkout and the virtualenv every session on + the host shares, and the ``.update_*`` markers are a single-slot mailbox + holding one requester's routing metadata. Two updaters racing the same + checkout is a real failure mode: the user double-taps ``/update`` because + the update takes minutes with no acknowledgement, or two platforms on one + multiplexed gateway both trigger it. + + None of these tests pre-seed a marker file. A test that writes + ``.update_pending.json`` itself only proves the handler reads a file the + test created — it never exercises the window between observing the slot is + free and taking it, which is where the collision actually happens. + """ + + def test_claim_update_slot_admits_exactly_one_caller(self, tmp_path): + """The reservation primitive is exclusive, and creates its own marker.""" + from gateway.slash_commands import _claim_update_slot + + pending_path = tmp_path / ".update_pending.json" + assert not pending_path.exists() + + assert _claim_update_slot(pending_path, 3600.0) is True + assert pending_path.exists() + assert _claim_update_slot(pending_path, 3600.0) is False + assert _claim_update_slot(pending_path, 3600.0) is False + + def test_claim_update_slot_is_exclusive_under_parallel_callers(self, tmp_path): + """16 threads released together still yield exactly one winner. + + This is the check-to-create gap in isolation. A + ``if path.exists(): return`` guard fails here because every thread can + observe the marker missing before any of them creates it. + """ + from gateway.slash_commands import _claim_update_slot + + pending_path = tmp_path / ".update_pending.json" + assert not pending_path.exists() + + workers = 16 + barrier = threading.Barrier(workers) + results = [] + results_lock = threading.Lock() + + def _claim(): + barrier.wait(timeout=15) + won = _claim_update_slot(pending_path, 3600.0) + with results_lock: + results.append(won) + + threads = [threading.Thread(target=_claim) for _ in range(workers)] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=20) + + assert not any(thread.is_alive() for thread in threads) + assert len(results) == workers + assert results.count(True) == 1 + + def test_claim_update_slot_reclaims_an_abandoned_reservation(self, tmp_path): + """A reservation past its TTL is taken over, so /update never wedges. + + An updater killed before it writes its exit code (host reboot, OOM + kill) leaves a marker nothing will clean up. Without the TTL the guard + would block /update for the lifetime of the profile. + """ + from gateway.slash_commands import _claim_update_slot + + pending_path = tmp_path / ".update_pending.json" + assert _claim_update_slot(pending_path, 3600.0) is True + assert _claim_update_slot(pending_path, 3600.0) is False + + abandoned = time.time() - 7200 + os.utime(pending_path, (abandoned, abandoned)) + assert _claim_update_slot(pending_path, 3600.0) is True + + @pytest.mark.asyncio + async def test_second_update_is_rejected_without_preseeding_a_marker(self, tmp_path): + """The first /update creates the marker; the second is refused. + + Nothing is written to the profile home before the first call, so the + reject path is driven entirely by state the handler itself produced. + """ + runner = _make_runner() + slash_commands_file, hermes_home = _fake_checkout(tmp_path) + pending_path = hermes_home / ".update_pending.json" + assert not pending_path.exists() + + first = _make_event(platform=Platform.TELEGRAM, + chat_id="first-chat", user_id="first-user") + second = _make_event(platform=Platform.TELEGRAM, + chat_id="second-chat", user_id="second-user") + + mock_watch = MagicMock() + with patch("gateway.run._hermes_home", hermes_home), \ + patch("gateway.slash_commands.__file__", slash_commands_file), \ + patch.object(runner, "_schedule_update_notification_watch", mock_watch), \ + patch("shutil.which", side_effect=lambda x: f"/usr/bin/{x}"), \ + patch("subprocess.Popen") as mock_popen: + first_reply = await runner._handle_update_command(first) + second_reply = await runner._handle_update_command(second) + + # Exactly one updater was launched. + assert mock_popen.call_count == 1 + assert "Starting Hermes update" in first_reply + assert "already running" in second_reply.lower() + + # The winner's routing metadata survived intact — the loser did not + # overwrite it, so the completion notice still reaches the first chat. + data = json.loads(pending_path.read_text()) + assert data["chat_id"] == "first-chat" + assert data["user_id"] == "first-user" + + # Both callers get the notification watcher, so the loser still learns + # how the update turned out. + assert mock_watch.call_count == 2 + + def test_concurrent_update_commands_spawn_exactly_one_updater(self, tmp_path): + """Two callers racing into the handler produce one updater, not two. + + Both are released from a barrier inside ``_resolve_hermes_bin`` — the + last step before the slot is claimed — so they enter the claim window + together, in real threads, against the real filesystem, with no marker + pre-seeded. + """ + runner = _make_runner() + slash_commands_file, hermes_home = _fake_checkout(tmp_path) + pending_path = hermes_home / ".update_pending.json" + assert not pending_path.exists() + + callers = 2 + barrier = threading.Barrier(callers) + + def _resolve_hermes_bin_at_barrier(): + barrier.wait(timeout=15) + return ["/usr/bin/hermes"] + + spawns = [] + spawn_lock = threading.Lock() + + def _record_spawn(*args, **kwargs): + with spawn_lock: + spawns.append(args[0]) + return MagicMock() + + replies = [] + replies_lock = threading.Lock() + + def _invoke(chat_id): + event = _make_event(platform=Platform.TELEGRAM, + chat_id=chat_id, user_id=f"user-{chat_id}") + reply = asyncio.run(runner._handle_update_command(event)) + with replies_lock: + replies.append(reply) + + with patch("gateway.run._hermes_home", hermes_home), \ + patch("gateway.slash_commands.__file__", slash_commands_file), \ + patch("gateway.run._resolve_hermes_bin", _resolve_hermes_bin_at_barrier), \ + patch.object(runner, "_schedule_update_notification_watch", MagicMock()), \ + patch("shutil.which", side_effect=lambda x: f"/usr/bin/{x}"), \ + patch("subprocess.Popen", side_effect=_record_spawn): + threads = [ + threading.Thread(target=_invoke, args=(f"chat-{i}",)) + for i in range(callers) + ] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=20) + assert not any(thread.is_alive() for thread in threads) + + assert len(spawns) == 1, f"expected one updater, got {len(spawns)}" + assert len(replies) == callers + assert sum("already running" in reply.lower() for reply in replies) == 1 + + # The surviving marker is one caller's complete, parseable metadata — + # not a half-written or interleaved document. + data = json.loads(pending_path.read_text()) + assert data["chat_id"] in {"chat-0", "chat-1"} + assert data["user_id"] == f"user-{data['chat_id']}" + + @pytest.mark.asyncio + async def test_failed_spawn_releases_the_update_slot(self, tmp_path): + """A spawn that never started must not block the next /update. + + The slot is reserved before the spawn, so the failure path has to give + it back — otherwise one transient OSError would wedge /update until the + reservation TTL expired. + """ + runner = _make_runner() + slash_commands_file, hermes_home = _fake_checkout(tmp_path) + pending_path = hermes_home / ".update_pending.json" + event = _make_event(platform=Platform.TELEGRAM, chat_id="chat-1") + + with patch("gateway.run._hermes_home", hermes_home), \ + patch("gateway.slash_commands.__file__", slash_commands_file), \ + patch.object(runner, "_schedule_update_notification_watch", MagicMock()), \ + patch("shutil.which", side_effect=lambda x: f"/usr/bin/{x}"), \ + patch("subprocess.Popen", side_effect=OSError("cannot fork")): + failed_reply = await runner._handle_update_command(event) + + assert "Failed to start update" in failed_reply + assert not pending_path.exists() + assert not (hermes_home / ".update_pending.tmp").exists() + + with patch("gateway.run._hermes_home", hermes_home), \ + patch("gateway.slash_commands.__file__", slash_commands_file), \ + patch.object(runner, "_schedule_update_notification_watch", MagicMock()), \ + patch("shutil.which", side_effect=lambda x: f"/usr/bin/{x}"), \ + patch("subprocess.Popen") as mock_popen: + retry_reply = await runner._handle_update_command(event) + + assert mock_popen.call_count == 1 + assert "Starting Hermes update" in retry_reply + + # --------------------------------------------------------------------------- # Platform allowlist gate # --------------------------------------------------------------------------- From 0e2af35db73319f84a9364308f98ec46e4d67b2e Mon Sep 17 00:00:00 2001 From: briandevans <252620095+briandevans@users.noreply.github.com> Date: Tue, 4 Aug 2026 17:31:01 -0700 Subject: [PATCH 2/2] fix(gateway): re-check the claimed marker after reserving the update slot MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The `claimed_path.exists()` pre-check was not synchronized with the exclusive create. `_send_update_notification` renames `.update_pending.json` -> `.update_pending.claimed.json` while an update is in flight; landing that rename between the check and the create leaves `pending` momentarily absent, so `_claim_update_slot` succeeds and admits a second updater against the live one. The notifier's later `claimed.replace(pending)` then clobbers the new reservation. Re-check the claimed marker immediately after a successful claim and hand the reservation back if it is present. `_release_update_slot` only removes the marker while it is still the empty file the claim created, so a marker the notifier has already restored — someone else's routing metadata — survives instead of being deleted. --- gateway/slash_commands.py | 32 ++++++++++-- tests/gateway/test_update_command.py | 77 ++++++++++++++++++++++++++++ 2 files changed, 104 insertions(+), 5 deletions(-) diff --git a/gateway/slash_commands.py b/gateway/slash_commands.py index dab490ac9b95..2eecad2776e7 100644 --- a/gateway/slash_commands.py +++ b/gateway/slash_commands.py @@ -187,6 +187,22 @@ def _create_exclusive() -> bool: return _create_exclusive() +def _release_update_slot(pending_path: Path) -> None: + """Give back a reservation taken by :func:`_claim_update_slot`. + + Only removes the marker while it is still the empty file the claim created. + A non-empty marker means somebody else's routing metadata is now there — + the notifier restoring a live update's marker via + ``claimed_path.replace(pending_path)`` — and deleting that would lose their + completion notice. + """ + try: + if pending_path.stat().st_size == 0: + pending_path.unlink() + except OSError: + pass + + class GatewaySlashCommandsMixin: """In-session slash-command handlers for GatewayRunner.""" @@ -5737,12 +5753,18 @@ async def _handle_update_command(self, event: MessageEvent) -> str: # # `.update_pending.claimed.json` is the notifier's transient rename of # the marker (see _send_update_notification); while it exists an update - # is in flight, so reject. Checking it is not a check-then-create race: - # it can only ever produce an extra rejection, never an extra spawn, - # and the exclusive create below is what actually decides the winner. - if claimed_path.exists() or not _claim_update_slot( + # is in flight. It is checked both BEFORE and AFTER the exclusive + # create, because the notifier can rename pending -> claimed between + # the two: that leaves pending momentarily absent, so the create would + # otherwise succeed and admit a second updater against a live one. + admitted = not claimed_path.exists() and _claim_update_slot( pending_path, _UPDATE_RESERVATION_TTL_S - ): + ) + if admitted and claimed_path.exists(): + _release_update_slot(pending_path) + admitted = False + + if not admitted: # The loser still gets the outcome: the watcher is profile-wide and # reports the result into the chat that started the update. self._schedule_update_notification_watch() diff --git a/tests/gateway/test_update_command.py b/tests/gateway/test_update_command.py index 11db21f0fc83..742cb0cdefd8 100644 --- a/tests/gateway/test_update_command.py +++ b/tests/gateway/test_update_command.py @@ -388,6 +388,83 @@ def _invoke(chat_id): assert data["chat_id"] in {"chat-0", "chat-1"} assert data["user_id"] == f"user-{data['chat_id']}" + @pytest.mark.asyncio + async def test_rejects_when_notifier_claims_a_live_update_mid_admission(self, tmp_path): + """A marker renamed to `.claimed` mid-admission still blocks the caller. + + ``_send_update_notification`` renames pending -> claimed while an update + is in flight. Landing that between the ``claimed`` check and the + exclusive create leaves ``pending`` momentarily absent, so the create + would succeed and admit a second updater against the live one — and the + notifier's later ``claimed.replace(pending)`` would clobber the new + reservation. + + The claimed marker here is written from *inside* the admission window + to stand in for that concurrent notifier — it is not pre-seeded state + the guard is handed up front. + """ + from gateway import slash_commands as slash_commands_module + + runner = _make_runner() + slash_commands_file, hermes_home = _fake_checkout(tmp_path) + pending_path = hermes_home / ".update_pending.json" + claimed_path = hermes_home / ".update_pending.claimed.json" + event = _make_event(platform=Platform.TELEGRAM, chat_id="second-chat") + + real_claim = slash_commands_module._claim_update_slot + + def _claim_after_notifier_rename(path, ttl_seconds): + # The notifier claims a live update's marker right here, after this + # handler already saw claimed_path missing. + claimed_path.write_text( + json.dumps({ + "platform": "telegram", + "chat_id": "live-chat", + "user_id": "live-user", + }), + encoding="utf-8", + ) + return real_claim(path, ttl_seconds) + + mock_watch = MagicMock() + with patch("gateway.run._hermes_home", hermes_home), \ + patch("gateway.slash_commands.__file__", slash_commands_file), \ + patch.object(slash_commands_module, "_claim_update_slot", + _claim_after_notifier_rename), \ + patch.object(runner, "_schedule_update_notification_watch", mock_watch), \ + patch("shutil.which", side_effect=lambda x: f"/usr/bin/{x}"), \ + patch("subprocess.Popen") as mock_popen: + reply = await runner._handle_update_command(event) + + mock_popen.assert_not_called() + assert "already running" in reply.lower() + mock_watch.assert_called_once() + + # The reservation was handed back, and the live update's claimed marker + # is untouched so its owner still gets the completion notice. + assert not pending_path.exists() + assert json.loads(claimed_path.read_text())["chat_id"] == "live-chat" + + def test_release_update_slot_keeps_a_marker_that_holds_metadata(self, tmp_path): + """Releasing never deletes someone else's routing metadata. + + If the notifier restores a live update's marker with + ``claimed_path.replace(pending_path)`` before the release runs, the + marker is no longer this caller's empty reservation and must survive. + """ + from gateway.slash_commands import _claim_update_slot, _release_update_slot + + pending_path = tmp_path / ".update_pending.json" + + assert _claim_update_slot(pending_path, 3600.0) is True + _release_update_slot(pending_path) + assert not pending_path.exists() + + pending_path.write_text(json.dumps({"chat_id": "live-chat"}), encoding="utf-8") + _release_update_slot(pending_path) + assert pending_path.exists() + assert json.loads(pending_path.read_text())["chat_id"] == "live-chat" + @pytest.mark.asyncio async def test_failed_spawn_releases_the_update_slot(self, tmp_path): """A spawn that never started must not block the next /update.