From c71c29ecaecb8fcff7ef97b67e2a671a04504d80 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 9 Aug 2026 17:18:32 -0700 Subject: [PATCH] fix: preserve exec completion identity across poll and heartbeat (#120575) * fix(agents): bind terminal polls to exact process UUID-owned completion receipts and ProcessSession-bound finished snapshots prevent same-slug successor consumption. * chore(plugin-sdk): refresh API baseline Refresh declaration-closure hashes for the internal system-event receipt boundary. --- .../.generated/plugin-sdk-api-baseline.sha256 | 64 ++--- scripts/plugin-sdk-surface-report.mts | 3 +- src/agents/bash-process-registry.test.ts | 21 ++ src/agents/bash-process-registry.ts | 26 +- src/agents/bash-tools.exec-runtime.test.ts | 25 +- src/agents/bash-tools.exec-runtime.ts | 25 +- ...h-tools.notify-on-exit-ack.test-support.ts | 57 +++++ .../bash-tools.notify-on-exit-ack.test.ts | 161 +++++++++++++ .../bash-tools.process.poll-timeout.test.ts | 228 +++++++++++++++++- src/agents/bash-tools.process.ts | 154 ++++++------ src/gateway/server-cron.test.ts | 20 +- src/gateway/server-cron.ts | 14 +- src/infra/system-events.test.ts | 37 +++ src/infra/system-events.ts | 18 ++ src/plugin-sdk/infra-runtime.ts | 16 +- 15 files changed, 683 insertions(+), 186 deletions(-) create mode 100644 src/agents/bash-tools.notify-on-exit-ack.test-support.ts create mode 100644 src/agents/bash-tools.notify-on-exit-ack.test.ts diff --git a/docs/.generated/plugin-sdk-api-baseline.sha256 b/docs/.generated/plugin-sdk-api-baseline.sha256 index ff81b5782e4e..5a9e4af947c9 100644 --- a/docs/.generated/plugin-sdk-api-baseline.sha256 +++ b/docs/.generated/plugin-sdk-api-baseline.sha256 @@ -3,10 +3,10 @@ 71522995185b956a0cc4927a472cc8d1153e5e998874bfd9a750513175174713 module/account-id d768139934447ff3ecf15470dc1fe613d36509fc5b93fc47ef09d20829cefa57 module/account-resolution 4fbb1c87e99399f842a20d75d5e35a4b7064a1b7f02115c23f9a2a7cdcfb57ee module/agent-config-primitives -0eeb85850d6056f815e231a7f0acbca6fb4115699407b823a9cf4437bf627e69 module/agent-harness -e1754867b9fbdd210a21c3d00c3aa5f161517bd2ed10c1601f52e3f39714699a module/agent-harness-runtime +30b474660d851867df45515fd3c4e7c07120932e29779277ec881fae6add4c4d module/agent-harness +217415168d269b097d5d47765f66458f0f5cbaa0e414366806ccb2c919f80c88 module/agent-harness-runtime d6097cfa1b410f4b5267a56a7bd19c2a33fbaf6642683dae6aec68e998e48f6c module/agent-media-payload -3035368684499711a63540fbef74b4aae7d131073367a77a3179e3b9aef16ad8 module/agent-runtime +34269cbdc975c509888ba94d53bc6ae20a8e987606cdb2b8fec6c2c34f144de6 module/agent-runtime b57a3cb274a9977c48c5387772df7707e50dde9a6132e532e7367a27225dc142 module/agent-scope-runtime 8fecb210e22bce4532b6ab649b09465f0bd2c857a44abf40db7d683d6491e6da module/allow-from bf66a0447f3d1e16a757843d068e56275f4781e9b427221fca26a854c3597467 module/allowlist-config-edit @@ -26,19 +26,19 @@ afad33fdaada25984db53504dc6f665ff0478c12f68b79f36a39d28ceb13355e module/channel c2cc71d5070b6071c51248b0648d1ad1a9468d3737df890adc77ec02025e8853 module/channel-config-primitives a6cca5706f3aba6abb2178b175a0986d921ade98c2c05d2a54451e2fb7e16825 module/channel-config-schema 37925e2b8c74b4444a14ea85b831ab569df9d46efebe89714527ee638719c100 module/channel-contract -c6427c4fea5bd0a7574af77850445624a80e9988f7129af224a489297af041a3 module/channel-core +1b52c802c98bcf60996c40a260be85387485945d8ad44fe51b056e87b66b5c4f module/channel-core f4a9870d37f3b4e824bc7b0f634e4eb868ae7dd4c5a9693f0105a6677c2ff5f9 module/channel-dm-policy -3f78f6022fc2942f60bc07a554f57ef9f8763db773f10f6425fa940acdbe917a module/channel-entry-contract +6bf7c17cde6bf073ca20b163e138ac9ca4d90020cd89de54f014488d38a1d035 module/channel-entry-contract 47cf8765e76c151ae7d2991d41beca62922a8837c1521e10a3fd23f9992c2d7c module/channel-feedback -fbc3da5f43928af283c3d14bdb5a68137478d6ee1c1ed99e73b1517d6170bde9 module/channel-inbound +008e6c083399fb5155ed108a3bb8dbe22bc4913a0a4e5249bfe374f61ce64f25 module/channel-inbound 3115366026efa38e07bfe0bfecd483e93454cc7bcad432385139c956777accff module/channel-inbound-debounce bc59c696ee45fb500d105c619c0ec81b8186bebe99135f23dc818ac16004deb0 module/channel-ingress-runtime bae1492066e55ea4f41b6400ab6d2ee454d56ed99b9027a484659fbd1257c46f module/channel-lifecycle 0e47457e38d1df0bd572e1408cde2ca6a788b65205f43c585316b5ad3a8f2f16 module/channel-logging -58f1f78cce81d60b6860105d38a1150698855a882562ced6bdf4d96d3fb52271 module/channel-message -6d2d0e01bff4bbe4488b41ffa5aa4ea1651f48ed4523868b4cfa4781e29c0664 module/channel-outbound -fa5894555b41c297fa7d9e0f75cdf4b006c23f6ff1b61181f0473c4ad443dd21 module/channel-pairing -6aff4ff9a57e00f053789873a941f4a6e318ecdb963c646716dc9d9f3e79e748 module/channel-plugin-common +e990d15fb50cf91f684a74993782e00317683aa9f59abcb9dc7b410126b25282 module/channel-message +b28f15439530b0bcc6c696eea3a31864c3eb6702b44f3e92bb10a5038edf996e module/channel-outbound +b28bff5ff9856a25948306de293171655cf9cfd14576ae278a1ab486a858f9a1 module/channel-pairing +13007e28d9e2cb04e7a11098ba0f3ff59c2369bc09ddb160a92714f6f47f31b1 module/channel-plugin-common 1786ca2cb6867de9fe386b4d2869c2adada32ddf7b047699aba56fa6dfc00ea9 module/channel-policy b82d23ae060319e1e0cbdc58268024bd4e6a4a3beab208cc515917249c155beb module/channel-reply-pipeline 482370e60135db9bfaf07f24bab549e5fde09ab265a6061a1f587c5d93929e91 module/channel-runtime-context @@ -51,8 +51,8 @@ f6c25ae55d49d90431682ca53c9dd836304942a738425d6ff028146a0c6d2c97 module/channel b2f920ff4a6b4190e6d6ea0a3effb001751e092f0e3ac0cf296721ff8c383d86 module/channel-streaming-config fdeffe356c7c4edeec9f8fd03edcadc375eabc7a9412e582b10c3180e3ef40fc module/cli-argv ad12670dbfe538f8d0ebf4fb2b68080e93a760278278e6b1ce9bb129d4b2d533 module/collection-runtime -74c62a5fefe8513234218650c2669ac4ca8fc34deb2f99cbe31d66d49cdbe694 module/command-auth -55b4622736b6edbe50e740542e39903bdd6218938f5109e14e87088806746f92 module/command-auth-native +1e7b3c2313e589380d6d988b14e0f94427382dcfa83fd06e7cb5621f61129b10 module/command-auth +4eda4571e1965b4ba8ce422cd5a2e0bc821430a506912986ff8685edbb0040dc module/command-auth-native 50c24235bca2c1d3c011f6bc266b57078b78a76d2d27ee12e7b3245f2947b493 module/command-detection e24382c2cca7fd2cf69ac353e258153dab80b279daf9f0ef54a727981a7dff3c module/command-primitives-runtime 1d6ebcdc843b7072a7fc4475eadadb90a1111d49c1da1e4326775800873cf9f1 module/command-status @@ -60,12 +60,12 @@ ab86235fcfff7c7cf0021fafeca6e92afe2c257ccf6a4a38441e0d41532998ff module/config- 0d99f5cb8c4978ed760e5fb4e476543759fbd8fd5bf73cc50a50c1d550203826 module/config-mutation 1c79d1356d7f41c22a0e85ff43766fc839cf13734435a3dcf7a4739512d513bf module/config-runtime d9e5f2ae27e29a40a6d6084c4d59e4811b4669235faa6f0d6a4fbfdd62e51ca3 module/conversation-runtime -5bebdf2011850e732ef1f71f5bc7669b7ebc3c4ece653f48c5ee0e3aa4ebc369 module/core -2872c791a0b43fad1020a8f595e506d369bed6c57819733919858dddbd96fbaf module/dedupe-runtime +807f792edbfe1c7ee2a161a3cd1f34f1ccbd2ea5c3f22d2f1f279fa64968d230 module/core +2ccc2c2011e0703068581ee3547295014d2dfc32b4730c29d66933387df0fbaa module/dedupe-runtime ebef0e650ab45e44c9335e2b3e15588c968cea6dadd125364a076f9c50ad1e8c module/device-bootstrap 21d86413166ef815581d606f678b6a216a1cc73ffe470b841f5bc4a131bff6df module/diagnostic-runtime 2af0c3b8867148d3deaf0b6d112743a42b206014ed1b5a5cf4ef939f1976d420 module/directory-runtime -a589c6a22e936d596cdca5ce8d97ea091fd8ddd8d9df9407b57bbf1d683c8e82 module/discord +5a8a3cf9049b9fac36413b5a4fc33c11b6257c3653bbf8250516880b14bf2b81 module/discord 39fe343ed2119de714757c365eef2ccec89c2c82a0876c60a4bef8ce469f8c8a module/error-runtime ce4f1602bf5b5de968ca97cede4f498b6ae709a0a59393b5ce9255cf0e6a9d5e module/extension-shared dd9f6e0fd33cc88b22543c1ee30cc09cf4de4d8f30dff7b7f9cebef885c21543 module/gateway-method-runtime @@ -75,8 +75,8 @@ c76bd967ef989feff27494df1e59aab2abef6e01aabf96c222b1ea8b153af1fc module/gateway a35bf7d756fc781b860d131bac3710e4d3b2dd4403f7f98313899ddf80b564d3 module/hook-runtime a953bd0c23c562e29fc512bc2354ed45c497728169d8201cb2cdb75d9e056073 module/inbound-envelope 4928af5d2509f696b896f53ac790303a0742202dbcdae3e44fe6d1b434a9c1ba module/inbound-event-delivery -b249cb9753ef633d61a25dbd65e4a14062bb2aa147c7562179127dbbeac8c838 module/inbound-reply-dispatch -54329d247af8f8dce6aa9910153b2e33f6d9e9e4ad28689f44de7cfee7cf668c module/infra-runtime +17538c7e143bb6e556cbebc9086de7c2317cbecfa5ebaee64df61681bafb51cb module/inbound-reply-dispatch +ff775bd21ab88126de07c99d507c1bf1a31d8cb93de0516675dbd41d8ff411d2 module/infra-runtime ce73721421f1b903dd04ead4df173582e59ea3e9990248102c448b419cc6d272 module/ingress-effect-once 93aa5d6a73719dd9ba6cc65a57e6d48b6c6a1191e3979b84a7d0ffd2c467191f module/interactive-runtime 408d257ab5cc4b88a22b7e7595039cb8fc524b261c44141b294fbd0100ba62ee module/json-store @@ -89,28 +89,28 @@ f74d7295fe716aa140aa0bc9300d6259d71dab826de0808fca6bb02592bf5d6e module/media-m 6a52f93107335f88751704352cc01e62add06f854a5b7d765e2a5ee87c0313b6 module/media-store b7e71516842300c041d2423d822da0080881db6ddc6b9b9ce38cabd8d546676b module/media-understanding a206a1486f6a6bed3091795b324e95dd710b4b3b87a4a8f18788e6158a48a922 module/media-understanding-runtime -ffc96fca1e036592ed6b640c4917e4773a168b3d88aab1a11de4a4b4df84ce30 module/meeting-runtime +ddd07beaf8a56b140b6eae4ba14b0cd67ace511ed2c926ba7cd445c5850af89a module/meeting-runtime 3312468e2e8f3423b765fac6bb17944b800ea2c84acffeb64342c99040b2f482 module/memory-core-host-engine-foundation -e3d2db75fab4b4a4d8a77f7db8fc2df2047678d0570385843c68a2478d2be4e5 module/memory-host-core +8b530380f1ad01d38977fa4d43ee4274516c9a476e0338cd3e2f5df915d0c183 module/memory-host-core 1efa0aadc4261d1c6073058cbf3dcc9fa681424819bdd14333e19b249bbc4b18 module/messaging-targets 6c43c704519f1178c8ea42d5a0596bec719a81f360b7678ecc00db3ccdb94be7 module/model-session-runtime -1d8c9e43af38096bec0580967c99f666acb810ecc5271c514103233077f64559 module/models-provider-runtime +339379a69a25d12393fc36863c6c7aabcb2d5b2a9c2dab78281cc5d29d229851 module/models-provider-runtime 843158613b23664aea805911d372a6c88d160d3a93f44b76291f47423aee1665 module/native-command-config-runtime f6ffada942145ca2fcec4df90740b4578afa09e81d46cc09122a65d6e345fe02 module/native-command-registry a6b5532576fa4cfd609d0966927eec12806424f2e1efce9fe47d903dfeb8e4dd module/param-readers ca7a56bb1a6169b4cf9befbf5aa21da280a8086fdc49fca4eec520a7a7c98549 module/persistent-dedupe b31f5d86904097993a55377fed7973cd298b0e1f49cade7e37a1e28f6e724108 module/plugin-config-runtime -29e893378475bf70f9265e4bb160860f4fb9783dbbd7d5849cde4d10d3450cb6 module/plugin-entry -61364b898140f866d23aa282251eb2cbbcaf233dd2c966f0d47ea82e1b31d8dc module/plugin-runtime -515ca993ecfdcb462f2e006816fdfc4dfd7752b30fda55b095bc1d443744a231 module/provider-auth -a54b42a90c511494e32245f5f1477a65d7b5ca5aaa9993ddafe6132ae42be7c2 module/provider-catalog-runtime +280b241f912fe8c405129dbc25a033c53d3f63bc44f09cfbe85a61b867d96b3b module/plugin-entry +4ab0070240ce2b6afff1b8583ccf35a022ede5905ba9a02271181ac112eaee37 module/plugin-runtime +dc834ed4730337a03ba7986daa91eadedfedb54a225d36e329cfded4bb62817f module/provider-auth +39d45f2d8eaff9d634d284ea0d586968cf4c26eec420461ef21874e3c112e762 module/provider-catalog-runtime 8131147d699394bd06503e2ea2f5f1a50b1594a87dded6d118b74a8d0328c8f6 module/proxy-capture 4a698efc36d896c4702de8df831e36b06e85c82aa30bf86c9b1066e6cad4b700 module/question-gateway-runtime 1171a76ea0485b36f77c9e12601c44a0e858043140e1d4669bb86d1f9e35930b module/reply-chunking -cefa0341af067f4ac97e6bc53836ea43a86187f942b3b272d294046498ef2c3f module/reply-dispatch-runtime +7b3e22d8670341edd5cb6b2d5edd092f8afd44455d367670defe12e23998a420 module/reply-dispatch-runtime 73f861fa3179d5af1159853c5acab0eec7a6c8f9398dcb75ea770e784fca6727 module/reply-history fedbd80588bfc407af75dec097f52d90794e64165e79209253beec5ff615a3a8 module/reply-payload -05768d0e5e067dba8fb4f421f7bc38e4bc0658dc4d7056f2a4f1e43da93d0872 module/reply-runtime +689c1a6c31f496f10818197ee9e62e6c086caa24dd658da53ec3477073f051b1 module/reply-runtime aa07d85d99fdd2b1e0cbe9975fb6dcae66b8bdce2607c6bd5402ae68bb15118c module/root-walk ad6c5c5b16e22f8b06994f69ecfafbc9a9622b07c14670f3c4e8ee81cf3b8c4d module/routing 7877a7e58fa32a64107154e5b714c6d165e96989d4aa5f43e0afac085a187af0 module/run-command @@ -118,20 +118,20 @@ ad6c5c5b16e22f8b06994f69ecfafbc9a9622b07c14670f3c4e8ee81cf3b8c4d module/routing 0570a20fac6020880900a0cf9a889d8f28f47b89d0e1d17862265b557236d440 module/runtime-config-snapshot 52174d47f6a7edad215456ada8e7497c29a492a89caa4861f0799d453fd7547d module/runtime-env 49e9b6a8195c89704eaa80656f176444af7cacbf639b759f41f2c78ae6bfcfd9 module/runtime-group-policy -ad27ae199a4f66826b886f80e1955a2b41c9bbae2a42d86c9aaa400fce451a67 module/runtime-store +cd04fa399ec561bc4cd2e17ab6fbfc8499ca1e94fcc203cccf3c397856221c56 module/runtime-store d17862c40825af1ddf0257b44f1e1cbb9c375e8e5ed668fae75d530d1a465cf9 module/secret-file 8e2ac4d3973d8d8ce4478e3440d66ee5c0d9213b0fe9e927c421d14fd31e5e86 module/secret-input b1b0229280d7cc4a880e4db75130714eb6adc9d6aea31240dc9a8ff6fda55219 module/secret-input-runtime 0cbc3908bd9e1d4527585023a1692eba6c0bfe82f107f49a926a8d1c79aee2a0 module/secret-ref-runtime 66f58d7e35cd7db412f2091e5733729a7a6ff8dfde55270624391d7cc194c9b0 module/security-runtime -ea2f3efe8130bcf644028bf2590669894696f05706770256a0e88a5243b8770b module/session-catalog -a9ba9783ae7f17f4ed3e939c682652aaf12dcfa43a5573f8fd161411744036a6 module/session-discussion +d03681d33846765af8a41a498812f0ce7d571b4d37d2f9eff950f8922a248da3 module/session-catalog +7bc12cad4cbf01f43632006a47c3971dc37be1d80354db642d229710746dbcea module/session-discussion cadcebdf79a8cdbc44d4ecec489babce7fecc52d573e7aee0337a2977b737347 module/session-store-runtime fb0af0a51ba93e070d8862701c016375cf8576c47be1d25e20537ec86b0bcf73 module/setup 3b37b8371bee9239682a4b2e9c066ef423ffba129d06ded0d7030359ec0dd68e module/setup-runtime 44d37e0d9131ad2859f41068f2604090c784e65f1bd6ebda8e051b6f2e5e1660 module/setup-tools 8ec6ca8a40d4117c669fbd0f20241d54e7e6954d7cff8c5b4afc67ec9626e216 module/skill-commands-runtime -0b09bb506ca970ee39b58953d9c1f0407942b2762f54d89a4cd82c35b9bc6c8b module/speech-settings +b4b700aa02ce7a77a25d47fbca4feb905fda0caab51409fbf243653663669b11 module/speech-settings 0b3de0b219e431d26298c83688ffe53c262d58a336ca910e29619b5d0447266e module/ssrf-policy 6639ba57aacbb620a2f6ff802ecd94fc90b1969292e0306598546b5b4b9261cc module/ssrf-runtime 3855f0a23281d21063762f5b3be2b7485c8bdb3b4cd1692de4f2a17d4eeae496 module/state-paths @@ -141,11 +141,11 @@ f097d0096b21c8a052f0f649b7512ecf2aba4744ae6956f001950e053828b309 module/string- aef35bee2502cd6ed8765409b758e452aff8ac9469fd773e6a2a44c9a1bc3f66 module/temp-path 87fa81b9e58d8fc04a4b4202d2d37fca339615f5225687d9db905151439e0f4d module/text-chunking d05a2db4a97844950ff3bc30f90f07e7f45600b1219e08a19264174b399817dc module/text-runtime -6f22141e149df716f200321942bb9c6bc722ceaf49c3ea2f249fa5f8efddd4df module/tool-plugin +20e8335d6698114f55c9322570e7be3dd7704aeb079c2d85046952783c43ccbf module/tool-plugin dc1a073c59ab61e2789533b777b3f0cb9af689d64a97796b10e8aa82552510db module/tool-results a5eef5a532c439b1489236711f4bc2721abe612a9db3911430e1344c239f9861 module/tool-send cda105b721d498df23a554c6b68be150b8fe66b8b9172185c31a0b3b0646b1dc module/web-media -60465458d5674539012dabb6bf7861783798fd80b8111948812f1f451d848fd0 module/webhook-ingress +e0b7fae87f43cc90d528df383ce48fbde104af954f5275afe36fe499a796c169 module/webhook-ingress e3a199a9ce0b85d203e9e8a29b500db29c6b7af307e3145d0a311e29d598925b module/webhook-request-guards de59e86e126b75d13251cba7ebbe27b44d9b5588785d98df5ff4d6722374c81f module/widget-html 9161b36ec0ab062ea41b363c894fcd672a7727f21cb726739f99f9c184fce69d module/zod diff --git a/scripts/plugin-sdk-surface-report.mts b/scripts/plugin-sdk-surface-report.mts index 5f578ebfe77b..4e4ccee460ad 100644 --- a/scripts/plugin-sdk-surface-report.mts +++ b/scripts/plugin-sdk-surface-report.mts @@ -329,7 +329,8 @@ export function readPluginSdkSurfaceBudgets(env: NodeJS.ProcessEnv = process.env "OPENCLAW_PLUGIN_SDK_MAX_PUBLIC_WILDCARD_REEXPORTS", // -1: text-runtime now names its global-singleton exports explicitly. // -1: infra-runtime now names its error exports explicitly. - 80, + // -1: infra-runtime excludes the internal system-event receipt API. + 79, env, ), }; diff --git a/src/agents/bash-process-registry.test.ts b/src/agents/bash-process-registry.test.ts index 8d2670b066a3..0a0dd2f62734 100644 --- a/src/agents/bash-process-registry.test.ts +++ b/src/agents/bash-process-registry.test.ts @@ -12,9 +12,11 @@ import { appendOutput, createSessionSlug, deleteSession, + drainFinishedSession, drainSession, getActiveBackgroundExecSessionCount, getFinishedSession, + getFinishedSessionForProcess, listFinishedSessions, listRunningSessions, markBackgrounded, @@ -209,10 +211,29 @@ describe("bash process registry", () => { tail: "", truncated: false, totalOutputChars: 0, + unreadOutput: { stdout: "", stderr: "", outputDropped: false }, }, ]); }); + it("moves unread output into the exact finished snapshot and consumes it once", () => { + const session = createRegistrySession({ + id: "exact-finished-output", + maxOutputChars: 100, + pendingMaxOutputChars: 100, + backgrounded: true, + }); + addSession(session); + appendOutput(session, "stdout", "terminal output\n"); + markExited(session, 0, null, "completed"); + + const finished = getFinishedSessionForProcess(session); + expect(finished).toBe(getFinishedSession(session.id)); + expect(finished && drainFinishedSession(finished).stdout).toBe("terminal output\n"); + expect(finished && drainFinishedSession(finished).stdout).toBe(""); + expect(drainSession(session).stdout).toBe(""); + }); + it("evicts the oldest finished sessions when their count exceeds the retention limit", () => { for (let index = 0; index < 53; index += 1) { const session = createRegistrySession({ diff --git a/src/agents/bash-process-registry.ts b/src/agents/bash-process-registry.ts index 6aaedecb579c..e0a15591b115 100644 --- a/src/agents/bash-process-registry.ts +++ b/src/agents/bash-process-registry.ts @@ -119,12 +119,14 @@ interface FinishedSession { tail: string; truncated: boolean; totalOutputChars: number; + unreadOutput?: ReturnType; terminalPollObserved?: boolean; notifyOnExitRemoval?: NotifyOnExitRemoval; } const runningSessions = new Map(); const finishedSessions = new Map(); +let finishedSessionsByProcess = new WeakMap(); const activeBackgroundExecSessionIds = new Set(); let finishedSessionOutputChars = 0; @@ -157,6 +159,11 @@ export function getFinishedSession(id: string) { return finishedSessions.get(id); } +/** Returns the terminal snapshot owned by this exact process incarnation. */ +export function getFinishedSessionForProcess(session: ProcessSession) { + return finishedSessionsByProcess.get(session); +} + function deleteFinishedSession(id: string): boolean { const session = finishedSessions.get(id); if (!session) { @@ -237,6 +244,13 @@ export function drainSession(session: ProcessSession) { return { stdout, stderr, outputDropped }; } +/** Consumes the output transferred to one exact terminal snapshot. */ +export function drainFinishedSession(session: FinishedSession) { + const output = session.unreadOutput; + session.unreadOutput = undefined; + return output ?? { stdout: "", stderr: "", outputDropped: false }; +} + /** Moves a session to finished state and records exit metadata. */ export function markExited( session: ProcessSession, @@ -270,7 +284,7 @@ export function markBackgrounded(session: ProcessSession) { /** Records that a terminal process poll consumed the process result. */ export function markTerminalPollObserved(session: ProcessSession): void { session.terminalPollObserved = true; - const finished = finishedSessions.get(session.id); + const finished = finishedSessionsByProcess.get(session); if (finished) { finished.terminalPollObserved = true; } @@ -286,7 +300,7 @@ export function recordNotifyOnExitRemoval( return; } session.notifyOnExitRemoval = remove; - const finished = finishedSessions.get(session.id); + const finished = finishedSessionsByProcess.get(session); if (finished) { finished.notifyOnExitRemoval = remove; } @@ -354,7 +368,7 @@ function moveToFinished(session: ProcessSession, status: ProcessStatus) { // Keep full completed logs; evict older records rather than silently // truncating the process poll/log contract or dropping the newest result. deleteFinishedSession(session.id); - finishedSessions.set(session.id, { + const finished: FinishedSession = { id: session.id, command: session.command, scopeKey: session.scopeKey, @@ -372,9 +386,12 @@ function moveToFinished(session: ProcessSession, status: ProcessStatus) { tail: session.tail, truncated: session.truncated, totalOutputChars: session.totalOutputChars, + unreadOutput: drainSession(session), ...(session.terminalPollObserved ? { terminalPollObserved: true } : {}), ...(session.notifyOnExitRemoval ? { notifyOnExitRemoval: session.notifyOnExitRemoval } : {}), - }); + }; + finishedSessionsByProcess.set(session, finished); + finishedSessions.set(session.id, finished); finishedSessionOutputChars += session.aggregated.length; while ( finishedSessions.size > MAX_FINISHED_SESSION_COUNT || @@ -459,6 +476,7 @@ export function listFinishedSessions() { function resetProcessRegistryForTests() { runningSessions.clear(); finishedSessions.clear(); + finishedSessionsByProcess = new WeakMap(); finishedSessionOutputChars = 0; activeBackgroundExecSessionIds.clear(); stopSweeper(); diff --git a/src/agents/bash-tools.exec-runtime.test.ts b/src/agents/bash-tools.exec-runtime.test.ts index 6a161beb18d0..6eecccdc613c 100644 --- a/src/agents/bash-tools.exec-runtime.test.ts +++ b/src/agents/bash-tools.exec-runtime.test.ts @@ -21,8 +21,7 @@ import { MAX_SAFE_TIMEOUT_DELAY_MS } from "../utils/timer-delay.js"; import type { BashSandboxConfig } from "./bash-tools.shared.js"; const requestHeartbeatMock = vi.hoisted(() => vi.fn()); -const enqueueSystemEventMock = vi.hoisted(() => vi.fn()); -const consumeSelectedSystemEventEntriesMock = vi.hoisted(() => vi.fn(() => [])); +const enqueueSystemEventWithReceiptMock = vi.hoisted(() => vi.fn()); const supervisorMock = vi.hoisted(() => ({ spawn: vi.fn(), })); @@ -32,17 +31,7 @@ vi.mock("../infra/heartbeat-wake.js", () => ({ })); vi.mock("../infra/system-events.js", () => ({ - enqueueSystemEvent: enqueueSystemEventMock, - enqueueSystemEventEntry: (text: string, options: { deliveryContext?: unknown }) => { - enqueueSystemEventMock(text, options); - return { - text, - ts: Date.now(), - contextKey: null, - deliveryContext: options.deliveryContext, - }; - }, - consumeSelectedSystemEventEntries: consumeSelectedSystemEventEntriesMock, + enqueueSystemEventWithReceipt: enqueueSystemEventWithReceiptMock, })); vi.mock("../process/supervisor/index.js", () => ({ @@ -77,8 +66,8 @@ beforeEach(() => { resetGatewaySuspendCoordinatorForLifecycleRestart(); resetProcessRegistryForTests(); requestHeartbeatMock.mockClear(); - enqueueSystemEventMock.mockClear(); - consumeSelectedSystemEventEntriesMock.mockClear(); + enqueueSystemEventWithReceiptMock.mockReset(); + enqueueSystemEventWithReceiptMock.mockReturnValue(vi.fn(() => true)); supervisorMock.spawn.mockReset(); }); @@ -204,7 +193,7 @@ function expectExecTarget( } function requireSystemEventCall(): [string, Record] { - const call = enqueueSystemEventMock.mock.calls[0]; + const call = enqueueSystemEventWithReceiptMock.mock.calls[0]; if (!call) { throw new Error("expected system event call"); } @@ -551,7 +540,7 @@ describe("exec notifyOnExit suppression", () => { const outcome = await runBackgroundedExit({ reason: "manual-cancel" }); expect(outcome.status).toBe("failed"); - expect(enqueueSystemEventMock).not.toHaveBeenCalled(); + expect(enqueueSystemEventWithReceiptMock).not.toHaveBeenCalled(); expect(requestHeartbeatMock).not.toHaveBeenCalled(); }); @@ -728,7 +717,7 @@ describe("sandbox exec finalization suspension", () => { expect(finalizeExec).toHaveBeenCalledOnce(); expect(getActiveBackgroundExecSessionCount()).toBe(0); expect(run.session.finalizing).toBe(false); - expect(enqueueSystemEventMock).toHaveBeenCalledTimes(1); + expect(enqueueSystemEventWithReceiptMock).toHaveBeenCalledTimes(1); expect(requireSystemEventCall()[0]).toContain( expectedStatus === "failed" ? "Exec failed" : "Exec completed", ); diff --git a/src/agents/bash-tools.exec-runtime.ts b/src/agents/bash-tools.exec-runtime.ts index 1ee718e87289..dde74646bf12 100644 --- a/src/agents/bash-tools.exec-runtime.ts +++ b/src/agents/bash-tools.exec-runtime.ts @@ -17,10 +17,7 @@ import { } from "../infra/exec-approvals.js"; import { requestHeartbeat } from "../infra/heartbeat-wake.js"; import { findPathKey, mergePathPrepend, removePathPrepend } from "../infra/path-prepend.js"; -import { - consumeSelectedSystemEventEntries, - enqueueSystemEventEntry, -} from "../infra/system-events.js"; +import { enqueueSystemEventWithReceipt } from "../infra/system-events.js"; import { isSubagentSessionKey } from "../sessions/session-key-utils.js"; /** * Bash exec runtime. @@ -358,15 +355,17 @@ function maybeNotifyOnExit(session: ProcessSession, status: "completed" | "faile sessionScope: session.sessionScope, }; const eventSessionKey = resolveEventSessionKeyForPolicy(sessionKey, eventRouting); - const event = enqueueSystemEventEntry(eventText, { - sessionKey: eventSessionKey, - deliveryContext: session.notifyDeliveryContext, - }); - if (event) { - recordNotifyOnExitRemoval( - session, - () => consumeSelectedSystemEventEntries(eventSessionKey, [event]).length > 0, - ); + const remove = enqueueSystemEventWithReceipt( + eventText, + { + sessionKey: eventSessionKey, + contextKey: `exec:${session.id}`, + deliveryContext: session.notifyDeliveryContext, + }, + { allowDuplicate: true }, + ); + if (remove) { + recordNotifyOnExitRemoval(session, remove); } // Subagent sessions receive exec results via process poll and announce flow; // the heartbeat would fall back to the main session and cause spurious wakes. diff --git a/src/agents/bash-tools.notify-on-exit-ack.test-support.ts b/src/agents/bash-tools.notify-on-exit-ack.test-support.ts new file mode 100644 index 000000000000..1d39edd4e940 --- /dev/null +++ b/src/agents/bash-tools.notify-on-exit-ack.test-support.ts @@ -0,0 +1,57 @@ +import type { ManagedRun } from "../process/supervisor/index.js"; +import type { SpawnInput } from "../process/supervisor/types.js"; +import { createDeferred } from "../shared/deferred.js"; +import type { DeliveryContext } from "../utils/delivery-context.types.js"; +import { markBackgrounded } from "./bash-process-registry.js"; +import { runExecProcess } from "./bash-tools.exec-runtime.js"; + +type SupervisorSpawnMock = { + mockImplementationOnce: (implementation: (input: SpawnInput) => Promise) => unknown; +}; + +export async function startDeferredNotifyRun(params: { + spawn: SupervisorSpawnMock; + sessionKey: string; + notifyDeliveryContext?: DeliveryContext; +}) { + const exit = createDeferred>>(); + params.spawn.mockImplementationOnce(async (input) => { + input.onStdout?.("producer output\n"); + return { + runId: input.runId ?? "notify-on-exit", + startedAtMs: Date.now(), + wait: async () => await exit.promise, + cancel: () => undefined, + }; + }); + const run = await runExecProcess({ + command: "notify-command", + workdir: process.cwd(), + env: {}, + usePty: false, + warnings: [], + maxOutput: 1000, + pendingMaxOutput: 1000, + notifyOnExit: true, + sessionKey: params.sessionKey, + notifyDeliveryContext: params.notifyDeliveryContext, + timeoutSec: null, + }); + markBackgrounded(run.session); + return { + run, + finish: async () => { + exit.resolve({ + reason: "exit", + exitCode: 0, + exitSignal: null, + durationMs: 1, + stdout: "", + stderr: "", + timedOut: false, + noOutputTimedOut: false, + }); + await run.promise; + }, + }; +} diff --git a/src/agents/bash-tools.notify-on-exit-ack.test.ts b/src/agents/bash-tools.notify-on-exit-ack.test.ts new file mode 100644 index 000000000000..77dfc324eaf0 --- /dev/null +++ b/src/agents/bash-tools.notify-on-exit-ack.test.ts @@ -0,0 +1,161 @@ +import { afterEach, beforeEach, expect, it, vi } from "vitest"; +import type { OpenClawConfig } from "../config/config.js"; +import { runHeartbeatOnce } from "../infra/heartbeat-runner.js"; +import { + seedMainSessionStore, + setupTelegramHeartbeatPluginRuntimeForTests, + withTempTelegramHeartbeatSandbox, +} from "../infra/heartbeat-runner.test-utils.js"; +import { + consumeSelectedSystemEventEntries, + enqueueSystemEventEntry, + peekSystemEventEntries, + resetSystemEventsForTest, +} from "../infra/system-events.js"; +import { resetProcessRegistryForTests } from "./bash-process-registry.test-support.js"; +import { startDeferredNotifyRun } from "./bash-tools.notify-on-exit-ack.test-support.js"; +import { createProcessTool } from "./bash-tools.process.js"; + +const requestHeartbeatMock = vi.hoisted(() => vi.fn()); +const supervisorSpawnMock = vi.hoisted(() => vi.fn()); +const randomMock = vi.hoisted(() => vi.fn(() => 0)); + +vi.mock("../infra/heartbeat-wake.js", async (importOriginal) => ({ + ...(await importOriginal()), + requestHeartbeat: requestHeartbeatMock, +})); +vi.mock("../infra/secure-random.js", async (importOriginal) => ({ + ...(await importOriginal()), + generateSecureInt: randomMock, +})); +vi.mock("../process/supervisor/index.js", () => ({ + getProcessSupervisor: () => ({ spawn: supervisorSpawnMock, getRecord: vi.fn() }), +})); + +const QUEUE_KEY = "agent:main:notify-ack"; +const startNotifyRun = () => + startDeferredNotifyRun({ + spawn: supervisorSpawnMock, + sessionKey: QUEUE_KEY, + notifyDeliveryContext: { channel: "telegram", to: "-100123", threadId: 42 }, + }); +const processTool = createProcessTool(); +const execute = (action: "poll" | "clear", sessionId: string) => + processTool.execute(`${action}-${sessionId}`, { action, sessionId }); +const poll = (sessionId: string) => execute("poll", sessionId); +const contexts = () => peekSystemEventEntries(QUEUE_KEY).map((event) => event.contextKey); + +beforeEach(() => { + setupTelegramHeartbeatPluginRuntimeForTests(); + vi.spyOn(Date, "now").mockReturnValue(1_800_000_000_000); +}); +afterEach(() => { + resetProcessRegistryForTests(); + resetSystemEventsForTest(); + vi.clearAllMocks(); + vi.restoreAllMocks(); +}); + +it("isolates identical completions across exact full-slug reuse", async () => { + const first = await startNotifyRun(); + await first.finish(); + await execute("clear", first.run.session.id); + enqueueSystemEventEntry("unrelated", { sessionKey: QUEUE_KEY, contextKey: "marker" }); + const second = await startNotifyRun(); + await second.finish(); + + expect([first.run.session.id, second.run.session.id]).toEqual(["amber-atlas", "amber-atlas"]); + expect(contexts()).toEqual(["exec:amber-atlas", "marker", "exec:amber-atlas"]); + const queued = peekSystemEventEntries(QUEUE_KEY); + expect(queued[0]?.id).not.toBe(queued[2]?.id); + expect(queued[0]).toEqual({ ...queued[2], id: queued[0]?.id }); + + await poll(second.run.session.id); + expect(contexts()).toEqual(["exec:amber-atlas", "marker"]); + await poll(second.run.session.id); + expect(contexts()).toEqual(["exec:amber-atlas", "marker"]); +}); + +it("invalidates a heartbeat snapshot when terminal poll consumes its occurrence", async () => { + const process = await startNotifyRun(); + await process.finish(); + const snapshot = peekSystemEventEntries(QUEUE_KEY); + await poll(process.run.session.id); + + expect(peekSystemEventEntries(QUEUE_KEY)).toEqual([]); + expect(consumeSelectedSystemEventEntries(QUEUE_KEY, snapshot)).toEqual([]); +}); + +it("keeps an identical successor queued when heartbeat consumes a stale snapshot", async () => { + await withTempTelegramHeartbeatSandbox(async ({ tmpDir, storePath, replySpy }) => { + const cfg: OpenClawConfig = { + agents: { + defaults: { + workspace: tmpDir, + heartbeat: { every: "5m", target: "telegram" }, + }, + }, + channels: { telegram: { allowFrom: ["*"] } }, + session: { mainKey: "notify-ack", store: storePath }, + }; + const sessionKey = await seedMainSessionStore(storePath, cfg, { + lastChannel: "telegram", + lastProvider: "telegram", + lastTo: "-100123", + lastThreadId: 42, + }); + expect(sessionKey).toBe(QUEUE_KEY); + + const first = await startNotifyRun(); + await first.finish(); + const firstQueued = peekSystemEventEntries(QUEUE_KEY); + expect(firstQueued).toHaveLength(1); + + let successor: Awaited> | undefined; + replySpy.mockImplementation(async () => { + await poll(first.run.session.id); + await execute("clear", first.run.session.id); + const replacement = await startNotifyRun(); + successor = replacement; + await replacement.finish(); + const replacementQueued = peekSystemEventEntries(QUEUE_KEY); + expect(replacementQueued).toHaveLength(1); + expect(replacementQueued[0]?.id).not.toBe(firstQueued[0]?.id); + expect(replacementQueued[0]).toEqual({ ...firstQueued[0], id: replacementQueued[0]?.id }); + return { text: "Handled the exec completion" }; + }); + const sendTelegram = vi.fn().mockResolvedValue({ messageId: "m1", chatId: "100123" }); + + const result = await runHeartbeatOnce({ + cfg, + agentId: "main", + source: "exec-event", + intent: "event", + reason: "exec-event", + deps: { + getQueueSize: () => 0, + getReplyFromConfig: replySpy, + telegram: sendTelegram, + }, + }); + + expect(result.status).toBe("ran"); + expect(sendTelegram).toHaveBeenCalledOnce(); + if (!successor) { + throw new Error("heartbeat reply did not enqueue the successor completion"); + } + expect(successor.run.session.id).toBe(first.run.session.id); + expect(peekSystemEventEntries(QUEUE_KEY)).toHaveLength(1); + + await poll(successor.run.session.id); + expect(peekSystemEventEntries(QUEUE_KEY)).toEqual([]); + }); +}); + +it("keeps an unpolled completion deliverable after finished-session cleanup", async () => { + const process = await startNotifyRun(); + await process.finish(); + await execute("clear", process.run.session.id); + + expect(contexts()).toEqual([`exec:${process.run.session.id}`]); +}); diff --git a/src/agents/bash-tools.process.poll-timeout.test.ts b/src/agents/bash-tools.process.poll-timeout.test.ts index 4764fb861123..701995d746d5 100644 --- a/src/agents/bash-tools.process.poll-timeout.test.ts +++ b/src/agents/bash-tools.process.poll-timeout.test.ts @@ -7,8 +7,10 @@ import { resetDiagnosticSessionStateForTest } from "../logging/diagnostic-sessio import { addSession, appendOutput, + deleteSession, getFinishedSession, markExited, + recordNotifyOnExitRemoval, } from "./bash-process-registry.js"; import { createProcessSessionFixture } from "./bash-process-registry.test-helpers.js"; import { resetProcessRegistryForTests } from "./bash-process-registry.test-support.js"; @@ -110,6 +112,218 @@ test("process poll waits for completion when timeout is provided", async () => { }); }); +test("waiting poll returns only output appended since the previous poll", async () => { + vi.useFakeTimers(); + try { + const sessionId = "sess-incremental-terminal-output"; + const { processTool, session } = createProcessSessionHarness(sessionId); + + appendOutput(session, "stdout", "already-observed\n"); + const firstPoll = await pollSession(processTool, "toolcall-first", sessionId); + expect(firstPoll.content[0]).toMatchObject({ + type: "text", + text: expect.stringContaining("already-observed"), + }); + + const pollPromise = pollSession(processTool, "toolcall-terminal", sessionId, 2_000); + setTimeout(() => { + appendOutput(session, "stdout", "new-terminal-output\n"); + markExited(session, 0, null, "completed"); + }, 10); + + await vi.advanceTimersByTimeAsync(250); + const terminalPoll = await pollPromise; + const terminalText = + terminalPoll.content[0]?.type === "text" ? terminalPoll.content[0].text : ""; + const details = terminalPoll.details as { status?: string; aggregated?: string }; + + expect(details.status).toBe("completed"); + expect(details.aggregated).toContain("already-observed"); + expect(details.aggregated).toContain("new-terminal-output"); + expect(terminalText).toContain("new-terminal-output"); + expect(terminalText).not.toContain("already-observed"); + } finally { + vi.useRealTimers(); + } +}); + +test("waiting poll retains terminal state and its receipt after indexed cleanup", async () => { + vi.useFakeTimers(); + try { + const sessionId = "sess-cleared-while-waiting"; + const { processTool, session } = createProcessSessionHarness(sessionId); + const remove = vi.fn(() => true); + + setTimeout(() => { + appendOutput(session, "stdout", "done after cleanup\n"); + markExited(session, 0, null, "completed"); + recordNotifyOnExitRemoval(session, remove); + deleteSession(sessionId); + }, 10); + + const pollPromise = pollSession(processTool, "toolcall-cleanup", sessionId, 2_000); + await vi.advanceTimersByTimeAsync(250); + const poll = await pollPromise; + + expect(poll.details).toMatchObject({ + status: "completed", + aggregated: expect.stringContaining("done after cleanup"), + }); + expect(poll.content[0]).toMatchObject({ + type: "text", + text: expect.stringContaining("done after cleanup"), + }); + expect(remove).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } +}); + +test("waiting poll does not adopt a same-id successor after removal", async () => { + vi.useFakeTimers(); + try { + const sessionId = "sess-reused-while-waiting"; + const { processTool, session } = createProcessSessionHarness(sessionId); + const successorRemove = vi.fn(() => true); + + setTimeout(() => { + session.backgrounded = false; + deleteSession(sessionId); + markExited(session, 0, null, "completed"); + + const successor = createProcessSessionFixture({ + id: sessionId, + command: "successor", + backgrounded: true, + }); + addSession(successor); + appendOutput(successor, "stdout", "successor output\n"); + markExited(successor, 0, null, "completed"); + recordNotifyOnExitRemoval(successor, successorRemove); + }, 10); + + const originalPoll = pollSession(processTool, "toolcall-original", sessionId, 2_000); + await vi.advanceTimersByTimeAsync(250); + const removed = await originalPoll; + + expect(removed.details).toMatchObject({ status: "failed" }); + expect(removed.content[0]).toMatchObject({ + type: "text", + text: `No session found for ${sessionId}`, + }); + expect(successorRemove).not.toHaveBeenCalled(); + + const successorPoll = await pollSession(processTool, "toolcall-successor", sessionId); + expect(successorPoll.details).toMatchObject({ + status: "completed", + aggregated: expect.stringContaining("successor output"), + }); + expect(successorRemove).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } +}); + +test("waiting poll never recommends successor logs for omitted original output", async () => { + vi.useFakeTimers(); + try { + const sessionId = "sess-reused-after-omitted-output"; + const { processTool, session } = createProcessSessionHarness(sessionId); + const originalRemove = vi.fn(() => true); + const successorRemove = vi.fn(() => true); + let expected: ReturnType | undefined; + + setTimeout(() => { + expected = appendOversizedPendingOutput(session); + markExited(session, 0, null, "completed"); + recordNotifyOnExitRemoval(session, originalRemove); + deleteSession(sessionId); + + const successor = createProcessSessionFixture({ + id: sessionId, + command: "successor", + backgrounded: true, + }); + addSession(successor); + appendOutput(successor, "stdout", "successor output\n"); + markExited(successor, 7, null, "completed"); + recordNotifyOnExitRemoval(successor, successorRemove); + }, 10); + + const originalPoll = pollSession(processTool, "toolcall-original-omitted", sessionId, 2_000); + await vi.advanceTimersByTimeAsync(250); + const original = await originalPoll; + if (!expected) { + throw new Error("expected pending output to be appended"); + } + const originalText = original.content[0]?.type === "text" ? original.content[0].text : ""; + + expect(original.details).toMatchObject({ + status: "completed", + exitCode: 0, + aggregated: expected.aggregated, + }); + expect(originalText).not.toContain(expected.earlierMarker); + expect(originalText).toContain(expected.latestMarker); + expect(originalText).not.toContain("successor output"); + expect(originalText).not.toContain("use action=log"); + expect(originalText).toContain("omitted output is no longer available through action=log"); + expect(originalRemove).toHaveBeenCalledOnce(); + expect(successorRemove).not.toHaveBeenCalled(); + + const successorLog = await processTool.execute("toolcall-successor-log", { + action: "log", + sessionId, + }); + expect(successorLog.details).toMatchObject({ status: "completed", exitCode: 7 }); + expect(successorLog.content[0]).toMatchObject({ + type: "text", + text: expect.stringContaining("successor output"), + }); + expect(successorRemove).not.toHaveBeenCalled(); + } finally { + vi.useRealTimers(); + } +}); + +test.each([ + { name: "waiting", exitBeforePoll: false }, + { name: "already-finished", exitBeforePoll: true }, +])("$name terminal polls do not replay drained output", async ({ exitBeforePoll }) => { + vi.useFakeTimers(); + try { + const sessionId = `sess-no-replay-${exitBeforePoll ? "finished" : "waiting"}`; + const { processTool, session } = createProcessSessionHarness(sessionId); + const finish = () => { + appendOutput(session, "stdout", "only once\n"); + markExited(session, 0, null, "completed"); + }; + if (exitBeforePoll) { + finish(); + } else { + setTimeout(finish, 10); + } + + const firstPromise = pollSession( + processTool, + "toolcall-first-terminal", + sessionId, + exitBeforePoll ? undefined : 2_000, + ); + await vi.advanceTimersByTimeAsync(250); + const first = await firstPromise; + const second = await pollSession(processTool, "toolcall-second-terminal", sessionId); + const firstText = first.content[0]?.type === "text" ? first.content[0].text : ""; + const secondText = second.content[0]?.type === "text" ? second.content[0].text : ""; + + expect(firstText).toContain("only once"); + expect(secondText).not.toContain("only once"); + expect(secondText).toContain("no new output"); + } finally { + vi.useRealTimers(); + } +}); + test("process poll accepts string timeout values", async () => { await expectCompletedPollWithTimeout({ sessionId: "sess-2", @@ -318,9 +532,9 @@ test("process poll resets retryInMs when output appears and clears on completion test.each([ { name: "below the retained tail", outputLength: 1_999, expectsOmissionNote: false }, { name: "at the retained tail", outputLength: 2_000, expectsOmissionNote: false }, - { name: "above the retained tail", outputLength: 2_001, expectsOmissionNote: true }, + { name: "above the retained tail", outputLength: 2_001, expectsOmissionNote: false }, ])( - "process poll discloses omitted finished output $name", + "process poll returns unread finished output $name", async ({ outputLength, expectsOmissionNote }) => { const sessionId = `sess-finished-tail-${outputLength}`; const { processTool, session } = createProcessSessionHarness(sessionId); @@ -389,17 +603,13 @@ test.each([ expect(details.aggregated).toHaveLength(aggregateCap); expect(text).not.toContain(earlierMarker); - expect(text).toContain(latestMarker); + expect(text).not.toContain(latestMarker); + expect(text).toContain("no new output"); expect(text).toContain("discarded at the retention cap and cannot be recovered"); expect(runningLogText).toContain("discarded at the retention cap and cannot be recovered"); expect(runningPollText).toContain("discarded at the retention cap and cannot be recovered"); expect(finishedLogText).toContain("discarded at the retention cap and cannot be recovered"); - if (aggregateCap > 2_000) { - expect(text).toContain("earlier retained output is omitted"); - expect(text).toContain("action=log with offset and limit"); - } else { - expect(text).not.toContain("action=log with offset and limit"); - } + expect(text).not.toContain("action=log with offset and limit"); }, ); diff --git a/src/agents/bash-tools.process.ts b/src/agents/bash-tools.process.ts index c140962055fd..25094706c1fb 100644 --- a/src/agents/bash-tools.process.ts +++ b/src/agents/bash-tools.process.ts @@ -12,8 +12,10 @@ import { acknowledgeNotifyOnExit, type ProcessSession, deleteSession, + drainFinishedSession, drainSession, getFinishedSession, + getFinishedSessionForProcess, getSession, listFinishedSessions, listRunningSessions, @@ -155,6 +157,58 @@ function resetPollRetrySuggestion(sessionId: string): void { } } +type FinishedSession = NonNullable>; + +function finishedPollResult( + sessionId: string, + finished: FinishedSession, +): AgentToolResult { + resetPollRetrySuggestion(sessionId); + acknowledgeNotifyOnExit(finished); + const { stdout, stderr, outputDropped } = drainFinishedSession(finished); + const output = [stdout.trimEnd(), stderr.trimEnd()].filter(Boolean).join("\n").trim(); + // Omitted retained output is pageable only while this public id still owns + // the exact snapshot; a reused slug must never point the model at successor logs. + const retainedOutputNote = outputDropped + ? getFinishedSession(sessionId) === finished + ? "\n\n[earlier output is omitted from this poll; use action=log with offset and limit to inspect retained output]" + : "\n\n[earlier output is omitted from this poll; omitted output is no longer available through action=log]" + : ""; + return { + content: [ + { + type: "text", + text: appendExecTimeoutRetryGuidance( + (output || "(no new output)") + + retentionCapNote(finished) + + retainedOutputNote + + `\n\nProcess exited with ${renderExecExitLabel(finished)}.`, + finished.exitReason, + ), + }, + ], + details: { + status: finished.status === "completed" ? "completed" : "failed", + sessionId, + exitCode: finished.exitCode ?? undefined, + ...(finished.exitSignal != null ? { exitSignal: finished.exitSignal } : {}), + ...(finished.exitReason + ? { + exitReason: finished.exitReason, + timedOut: + finished.exitReason === "overall-timeout" || + finished.exitReason === "no-output-timeout", + } + : {}), + ...(finished.noOutputTimedOut !== undefined + ? { noOutputTimedOut: finished.noOutputTimedOut } + : {}), + aggregated: finished.aggregated, + name: deriveSessionName(finished.command), + }, + }; +} + function createAbortError(reason: unknown): Error { if (reason instanceof Error) { return reason; @@ -398,52 +452,7 @@ export function createProcessTool( case "poll": { if (!scopedSession) { if (scopedFinished) { - resetPollRetrySuggestion(params.sessionId); - acknowledgeNotifyOnExit(scopedFinished); - // Aggregate-cap loss is permanent; tail omission remains pageable. - const aggregateOutputNote = retentionCapNote(scopedFinished); - const retainedOutputNote = - scopedFinished.tail.length < scopedFinished.aggregated.length - ? "\n\n[earlier retained output is omitted; use action=log with offset and limit to page]" - : ""; - return { - content: [ - { - type: "text", - text: appendExecTimeoutRetryGuidance( - (scopedFinished.tail || - `(no output recorded${ - scopedFinished.truncated ? " — truncated to cap" : "" - })`) + - aggregateOutputNote + - retainedOutputNote + - `\n\nProcess exited with ${renderExecExitLabel(scopedFinished)}.`, - scopedFinished.exitReason, - ), - }, - ], - details: { - status: scopedFinished.status === "completed" ? "completed" : "failed", - sessionId: params.sessionId, - exitCode: scopedFinished.exitCode ?? undefined, - ...(scopedFinished.exitSignal != null - ? { exitSignal: scopedFinished.exitSignal } - : {}), - ...(scopedFinished.exitReason - ? { - exitReason: scopedFinished.exitReason, - timedOut: - scopedFinished.exitReason === "overall-timeout" || - scopedFinished.exitReason === "no-output-timeout", - } - : {}), - ...(scopedFinished.noOutputTimedOut !== undefined - ? { noOutputTimedOut: scopedFinished.noOutputTimedOut } - : {}), - aggregated: scopedFinished.aggregated, - name: deriveSessionName(scopedFinished.command), - }, - }; + return finishedPollResult(params.sessionId, scopedFinished); } resetPollRetrySuggestion(params.sessionId); return failText(`No session found for ${params.sessionId}`); @@ -458,30 +467,26 @@ export function createProcessTool( await sleepPollInterval(Math.max(0, Math.min(250, deadline - Date.now())), signal); } } - const { stdout, stderr, outputDropped } = drainSession(scopedSession); - const exited = scopedSession.exited; - if (exited) { + if (scopedSession.exited) { markTerminalPollObserved(scopedSession); - acknowledgeNotifyOnExit(scopedSession); + // Exit finalization owns the terminal transition. Re-read by process + // object because the public id may already index a successor. + const finishedAfterWait = getFinishedSessionForProcess(scopedSession); + if (finishedAfterWait && isInScope(finishedAfterWait)) { + return finishedPollResult(params.sessionId, finishedAfterWait); + } + resetPollRetrySuggestion(params.sessionId); + return failText(`No session found for ${params.sessionId}`); } - const status = exited - ? scopedSession.terminalStatus === "completed" - ? "completed" - : "failed" - : "running"; + const { stdout, stderr, outputDropped } = drainSession(scopedSession); const output = [stdout.trimEnd(), stderr.trimEnd()].filter(Boolean).join("\n").trim(); const aggregateOutputNote = retentionCapNote(scopedSession); const retainedOutputNote = outputDropped ? "\n\n[earlier output is omitted from this poll; use action=log with offset and limit to inspect retained output]" : ""; const hasNewOutput = output.length > 0; - const retryInMs = exited - ? undefined - : recordPollRetrySuggestion(params.sessionId, hasNewOutput); - if (exited) { - resetPollRetrySuggestion(params.sessionId); - } - const runtime = exited ? undefined : describeRunningSession(scopedSession); + const retryInMs = recordPollRetrySuggestion(params.sessionId, hasNewOutput); + const runtime = describeRunningSession(scopedSession); return { content: [ { @@ -490,34 +495,17 @@ export function createProcessTool( (output || "(no new output)") + aggregateOutputNote + retainedOutputNote + - (exited - ? `\n\nProcess exited with ${renderExecExitLabel(scopedSession)}.` - : buildInputWaitHint(runtime) || "\n\nProcess still running."), - exited ? scopedSession.exitReason : undefined, + (buildInputWaitHint(runtime) || "\n\nProcess still running."), + undefined, ), }, ], details: { - status, + status: "running", sessionId: params.sessionId, - exitCode: exited ? (scopedSession.exitCode ?? undefined) : undefined, - ...(exited && scopedSession.exitSignal != null - ? { exitSignal: scopedSession.exitSignal } - : {}), - ...(exited && scopedSession.exitReason - ? { - exitReason: scopedSession.exitReason, - timedOut: - scopedSession.exitReason === "overall-timeout" || - scopedSession.exitReason === "no-output-timeout", - } - : {}), - ...(exited && scopedSession.noOutputTimedOut !== undefined - ? { noOutputTimedOut: scopedSession.noOutputTimedOut } - : {}), aggregated: scopedSession.aggregated, name: deriveSessionName(scopedSession.command), - ...(runtime ? runningSessionInputDetails(runtime) : {}), + ...runningSessionInputDetails(runtime), ...(typeof retryInMs === "number" ? { retryInMs } : {}), }, }; diff --git a/src/gateway/server-cron.test.ts b/src/gateway/server-cron.test.ts index 9d930b542f6e..2fd1f51a076d 100644 --- a/src/gateway/server-cron.test.ts +++ b/src/gateway/server-cron.test.ts @@ -21,7 +21,7 @@ type RunCronIsolatedAgentTurnMock = (params: { const { enqueueSystemEventMock, - consumeSelectedSystemEventEntriesMock, + systemEventReceiptRemoveMock, requestHeartbeatMock, runHeartbeatOnceMock, loadConfigMock, @@ -40,7 +40,7 @@ const { isAgentDeletionBlockedMock, } = vi.hoisted(() => ({ enqueueSystemEventMock: vi.fn(), - consumeSelectedSystemEventEntriesMock: vi.fn((_sessionKey, entries) => entries ?? []), + systemEventReceiptRemoveMock: vi.fn(() => true), requestHeartbeatMock: vi.fn(), runHeartbeatOnceMock: vi.fn< (...args: unknown[]) => Promise<{ status: "ran"; durationMs: number }> @@ -103,19 +103,12 @@ function enqueueSystemEvent(text: string, opts?: unknown) { return enqueueSystemEventMock(text, opts); } -function enqueueSystemEventEntry(text: string, opts?: unknown) { +function enqueueSystemEventWithReceipt(text: string, opts?: unknown) { const result = enqueueSystemEventMock(text, opts); if (result === false || result === null) { return null; } - return { - text, - ts: Date.now(), - }; -} - -function consumeSelectedSystemEventEntries(sessionKey: string, entries: readonly unknown[]) { - return consumeSelectedSystemEventEntriesMock(sessionKey, entries); + return systemEventReceiptRemoveMock; } function requestHeartbeat(...args: unknown[]) { @@ -128,8 +121,7 @@ function runHeartbeatOnce(...args: unknown[]) { vi.mock("../infra/system-events.js", () => ({ enqueueSystemEvent, - enqueueSystemEventEntry, - consumeSelectedSystemEventEntries, + enqueueSystemEventWithReceipt, })); vi.mock("../infra/heartbeat-wake.js", async () => { @@ -315,7 +307,7 @@ describe("buildGatewayCronService", () => { beforeEach(() => { resetActiveCronTaskRunsForTests(); enqueueSystemEventMock.mockClear(); - consumeSelectedSystemEventEntriesMock.mockClear(); + systemEventReceiptRemoveMock.mockClear(); requestHeartbeatMock.mockClear(); runHeartbeatOnceMock.mockClear(); loadConfigMock.mockClear(); diff --git a/src/gateway/server-cron.ts b/src/gateway/server-cron.ts index c408cf54adfe..7aae24b7f784 100644 --- a/src/gateway/server-cron.ts +++ b/src/gateway/server-cron.ts @@ -55,10 +55,7 @@ import { runHeartbeatOnce } from "../infra/heartbeat-runner.js"; import { requestHeartbeat } from "../infra/heartbeat-wake.js"; import { mergeSsrFPolicies } from "../infra/net/ssrf.js"; import { listConfiguredMessageChannels } from "../infra/outbound/channel-selection.js"; -import { - consumeSelectedSystemEventEntries, - enqueueSystemEventEntry, -} from "../infra/system-events.js"; +import { enqueueSystemEventWithReceipt } from "../infra/system-events.js"; import { getChildLogger } from "../logging.js"; import { getGlobalHookRunner } from "../plugins/hook-runner-global.js"; import type { @@ -657,17 +654,12 @@ export function buildGatewayCronService(params: { if (!sessionKey) { throw new Error("Cron system event target did not resolve a session key."); } - const event = enqueueSystemEventEntry(text, { + const remove = enqueueSystemEventWithReceipt(text, { sessionKey, contextKey: opts?.contextKey, deliveryContext: opts?.deliveryContext, }); - return event - ? { - accepted: true, - remove: () => consumeSelectedSystemEventEntries(sessionKey, [event]).length > 0, - } - : { accepted: false }; + return remove ? { accepted: true, remove } : { accepted: false }; }, resolveOriginDeliveryContext: (opts) => { // Resolve the wake target the same way the enqueue/heartbeat deps do, diff --git a/src/infra/system-events.test.ts b/src/infra/system-events.test.ts index f48b314495f9..087a49e8ed1e 100644 --- a/src/infra/system-events.test.ts +++ b/src/infra/system-events.test.ts @@ -17,6 +17,7 @@ import { drainSystemEventEntries, enqueueSystemEvent, enqueueSystemEventEntry, + enqueueSystemEventWithReceipt, hasSystemEvents, isSystemEventContextChanged, peekSystemEventEntries, @@ -220,6 +221,42 @@ describe("system events (session routing)", () => { expect(peekSystemEvents(key)).toEqual(["second"]); }); + it("removes an exact receipt once while preserving its sibling", () => { + const key = "agent:main:test-receipt"; + const receipt = enqueueSystemEventWithReceipt("first", { + sessionKey: ` ${key} `, + contextKey: "exec:first", + }); + expect(receipt).not.toBeNull(); + enqueueSystemEvent("sibling", { sessionKey: key, contextKey: "exec:sibling" }); + + expect(receipt?.()).toBe(true); + expect(peekSystemEvents(key)).toEqual(["sibling"]); + expect(receipt?.()).toBe(false); + }); + + it("keeps structurally identical receipt-owned siblings distinct", () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-08-09T00:00:00Z")); + const key = "agent:main:test-identical-receipts"; + const options = { sessionKey: key, contextKey: "exec:reused-slug" }; + const first = enqueueSystemEventWithReceipt("completed", options, { + allowDuplicate: true, + }); + const second = enqueueSystemEventWithReceipt("completed", options, { + allowDuplicate: true, + }); + const queued = peekSystemEventEntries(key); + + expect(queued[0]).toEqual({ ...queued[1], id: queued[0]?.id }); + expect(queued[0]?.id).not.toBe(queued[1]?.id); + expect(second?.()).toBe(true); + expect(peekSystemEventEntries(key).map((event) => event.id)).toEqual([queued[0]?.id]); + expect(second?.()).toBe(false); + expect(first?.()).toBe(true); + expect(peekSystemEventEntries(key)).toStrictEqual([]); + }); + it.each([ { name: "prefix consume with object spread", diff --git a/src/infra/system-events.ts b/src/infra/system-events.ts index 284bc2c799d8..18a8be81037e 100644 --- a/src/infra/system-events.ts +++ b/src/infra/system-events.ts @@ -54,6 +54,8 @@ type SystemEventOptions = { replace?: boolean; }; +type ReceiptOptions = { allowDuplicate?: boolean }; + function requireSessionKey(key?: string | null): string { const trimmed = normalizeOptionalString(key) ?? ""; if (!trimmed) { @@ -127,6 +129,7 @@ export function enqueueSystemEventEntry( function enqueueOwnedSystemEventEntry( text: string, options: SystemEventOptions, + receiptOptions?: ReceiptOptions, ): SystemEvent | null { if (options.replace) { return replaceSystemEventEntry(text, options); @@ -141,6 +144,7 @@ function enqueueOwnedSystemEventEntry( const normalizedDeliveryContext = normalizeDeliveryContext(options.deliveryContext); const normalizedOwnerAgentId = resolveSystemEventOptionsOwnerAgentId(options); if ( + receiptOptions?.allowDuplicate !== true && findDuplicateInQueue( entry.queue, cleaned, @@ -173,6 +177,20 @@ export function enqueueSystemEvent(text: string, options: SystemEventOptions) { return enqueueSystemEventEntry(text, options) !== null; } +/** Enqueues one occurrence and returns one-use removal ownership for its UUID. */ +export function enqueueSystemEventWithReceipt( + text: string, + options: SystemEventOptions, + receiptOptions?: ReceiptOptions, +): (() => boolean) | null { + const event = enqueueOwnedSystemEventEntry(text, options, receiptOptions); + if (!event) { + return null; + } + const sessionKey = requireSessionKey(options.sessionKey); + return () => consumeSelectedSystemEventEntries(sessionKey, [event]).length > 0; +} + export function drainSystemEventEntries(sessionKey: string): SystemEvent[] { const key = requireSessionKey(sessionKey); const entry = getSessionQueue(key); diff --git a/src/plugin-sdk/infra-runtime.ts b/src/plugin-sdk/infra-runtime.ts index cb0c1dbed479..e9b28997c816 100644 --- a/src/plugin-sdk/infra-runtime.ts +++ b/src/plugin-sdk/infra-runtime.ts @@ -267,7 +267,21 @@ export { type SecretFileReadResult, } from "../infra/secret-file.js"; export * from "../infra/secure-random.js"; -export * from "../infra/system-events.js"; +export { + consumeSelectedSystemEventEntries, + consumeSystemEventEntries, + drainSystemEventEntries, + drainSystemEvents, + enqueueSystemEvent, + enqueueSystemEventEntry, + hasSystemEvents, + isSystemEventContextChanged, + peekSystemEventEntries, + peekSystemEvents, + resetSystemEventsForTest, + resolveSystemEventDeliveryContext, + type SystemEvent, +} from "../infra/system-events.js"; export * from "../infra/system-message.ts"; export * from "../infra/tmp-openclaw-dir.js"; export * from "../infra/transport-ready.js";