mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-20 01:21:41 -06:00
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.
This commit is contained in:
committed by
GitHub
parent
cfe6ebcd1e
commit
c71c29ecae
@@ -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
|
||||
|
||||
@@ -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,
|
||||
),
|
||||
};
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -119,12 +119,14 @@ interface FinishedSession {
|
||||
tail: string;
|
||||
truncated: boolean;
|
||||
totalOutputChars: number;
|
||||
unreadOutput?: ReturnType<typeof drainSession>;
|
||||
terminalPollObserved?: boolean;
|
||||
notifyOnExitRemoval?: NotifyOnExitRemoval;
|
||||
}
|
||||
|
||||
const runningSessions = new Map<string, ProcessSession>();
|
||||
const finishedSessions = new Map<string, FinishedSession>();
|
||||
let finishedSessionsByProcess = new WeakMap<ProcessSession, FinishedSession>();
|
||||
const activeBackgroundExecSessionIds = new Set<string>();
|
||||
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();
|
||||
|
||||
@@ -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<string, unknown>] {
|
||||
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",
|
||||
);
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<ManagedRun>) => unknown;
|
||||
};
|
||||
|
||||
export async function startDeferredNotifyRun(params: {
|
||||
spawn: SupervisorSpawnMock;
|
||||
sessionKey: string;
|
||||
notifyDeliveryContext?: DeliveryContext;
|
||||
}) {
|
||||
const exit = createDeferred<Awaited<ReturnType<ManagedRun["wait"]>>>();
|
||||
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;
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -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<typeof import("../infra/heartbeat-wake.js")>()),
|
||||
requestHeartbeat: requestHeartbeatMock,
|
||||
}));
|
||||
vi.mock("../infra/secure-random.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../infra/secure-random.js")>()),
|
||||
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<ReturnType<typeof startNotifyRun>> | 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}`]);
|
||||
});
|
||||
@@ -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<typeof appendOversizedPendingOutput> | 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");
|
||||
},
|
||||
);
|
||||
|
||||
|
||||
@@ -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<ReturnType<typeof getFinishedSession>>;
|
||||
|
||||
function finishedPollResult(
|
||||
sessionId: string,
|
||||
finished: FinishedSession,
|
||||
): AgentToolResult<unknown> {
|
||||
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 } : {}),
|
||||
},
|
||||
};
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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";
|
||||
|
||||
Reference in New Issue
Block a user