fix(sessions): fence lifecycle transcript writers (#121284)

This commit is contained in:
Peter Steinberger
2026-08-09 17:55:56 -07:00
committed by GitHub
parent 1fb5d4f101
commit d7cdb2e60b
14 changed files with 653 additions and 184 deletions
+44 -44
View File
@@ -3,22 +3,22 @@
71522995185b956a0cc4927a472cc8d1153e5e998874bfd9a750513175174713 module/account-id
d768139934447ff3ecf15470dc1fe613d36509fc5b93fc47ef09d20829cefa57 module/account-resolution
4fbb1c87e99399f842a20d75d5e35a4b7064a1b7f02115c23f9a2a7cdcfb57ee module/agent-config-primitives
30b474660d851867df45515fd3c4e7c07120932e29779277ec881fae6add4c4d module/agent-harness
217415168d269b097d5d47765f66458f0f5cbaa0e414366806ccb2c919f80c88 module/agent-harness-runtime
4f26de54af184620dde1c5a3c9199bd9f155b95b6363fb5be7a8ea891e8f1f93 module/agent-harness
efc3643c6d96869c3795275047556913c6e3b0871c0f394f2a1510e097e31024 module/agent-harness-runtime
d6097cfa1b410f4b5267a56a7bd19c2a33fbaf6642683dae6aec68e998e48f6c module/agent-media-payload
34269cbdc975c509888ba94d53bc6ae20a8e987606cdb2b8fec6c2c34f144de6 module/agent-runtime
3eb909151540256fe88d02a34c2341f14c1836f6996ca527f1ae92399f43a12c module/agent-runtime
b57a3cb274a9977c48c5387772df7707e50dde9a6132e532e7367a27225dc142 module/agent-scope-runtime
8fecb210e22bce4532b6ab649b09465f0bd2c857a44abf40db7d683d6491e6da module/allow-from
bf66a0447f3d1e16a757843d068e56275f4781e9b427221fca26a854c3597467 module/allowlist-config-edit
15a4e709ba127fa17473499e1a874f3a5735ffae88f021052a04cd70c4fa3319 module/approval-auth-runtime
afb2c5728dab26c6c7c9ba555bb616ab45a36e3046d2916d2b4b8c2c51a444dc module/approval-client-runtime
d2238639ffc895409090021e8998a04dddada0e55d0c7be12b10824513f9fa0b module/approval-delivery-runtime
bb2fc2239e306b2f6323eff68961fb2d63f57bc8c1b50d1bd2383545ad15a628 module/approval-auth-runtime
de6b92a0328ee31ef89690639c7e5734baf6472266c662e8c34337a25aa78789 module/approval-client-runtime
2895ef48c4e7aa0cbc469df02c49a1479c4e20699d4cffe57fe8617a4ee0f994 module/approval-delivery-runtime
3be30812656a1cd2da700865f74f087473dac8b3958915d4eb9322e8f4fa042c module/approval-gateway-runtime
00583c41e60108630d09a5f6d22404fdc005156026d7f912141f126c5b2d9e37 module/approval-handler-adapter-runtime
4ea77f879914b3ab92272aa8a5f33d07da4350b061f750c422a1e11240ebadba module/approval-handler-runtime
236d858fa9493d1e8ec3f0014fb69d1a85ac04ec5623f6260612d32c84cb69b4 module/approval-native-runtime
06fa81619f39f3ef0f3e2aa124c536885b54ce46b01bb35ae912bffc01238c80 module/approval-native-runtime
6ecc8f2ddeb9ed348a68d845a876e9b03916613e6b58864c76aa6abc078aa3a5 module/approval-reply-runtime
a348239ac7eedcd42f1366062ba0a1d83b0a8c72bb4163867366c99c48c59b6c module/approval-runtime
a1820bb8b4d0f0bc3b74d121f89509f8dc302ef737feaa1a76da17ea3f3a6ab2 module/approval-runtime
01ca912836b8dec672f705e294f72d346e778557e4c591317d67558ea7669c0b module/archive
d7e53de63b0ac11a266e4abdc18ba6e9401b80309f5c8f5f6a72a00f65dfe3bd module/boolean-param
b11b9d8fb991e26acd6e7efcb96b77635455d2c9b94d0a7e840dbdac94ad5b85 module/channel-actions
@@ -26,46 +26,46 @@ afad33fdaada25984db53504dc6f665ff0478c12f68b79f36a39d28ceb13355e module/channel
c2cc71d5070b6071c51248b0648d1ad1a9468d3737df890adc77ec02025e8853 module/channel-config-primitives
a6cca5706f3aba6abb2178b175a0986d921ade98c2c05d2a54451e2fb7e16825 module/channel-config-schema
37925e2b8c74b4444a14ea85b831ab569df9d46efebe89714527ee638719c100 module/channel-contract
1b52c802c98bcf60996c40a260be85387485945d8ad44fe51b056e87b66b5c4f module/channel-core
301edf190b5ec4928ed70a59cf68343a605744c4aba9fedf56247398bd4ede67 module/channel-core
f4a9870d37f3b4e824bc7b0f634e4eb868ae7dd4c5a9693f0105a6677c2ff5f9 module/channel-dm-policy
6bf7c17cde6bf073ca20b163e138ac9ca4d90020cd89de54f014488d38a1d035 module/channel-entry-contract
e0f0dbf074b2ca3f9a8f6961e11c094f2bd9528213b7206613bd4f2a64cb93a8 module/channel-entry-contract
47cf8765e76c151ae7d2991d41beca62922a8837c1521e10a3fd23f9992c2d7c module/channel-feedback
008e6c083399fb5155ed108a3bb8dbe22bc4913a0a4e5249bfe374f61ce64f25 module/channel-inbound
339faec638b81f50cebec82e4e683ab2ad7f754953626b2f35a9b160424ae082 module/channel-inbound
3115366026efa38e07bfe0bfecd483e93454cc7bcad432385139c956777accff module/channel-inbound-debounce
bc59c696ee45fb500d105c619c0ec81b8186bebe99135f23dc818ac16004deb0 module/channel-ingress-runtime
bae1492066e55ea4f41b6400ab6d2ee454d56ed99b9027a484659fbd1257c46f module/channel-lifecycle
0e47457e38d1df0bd572e1408cde2ca6a788b65205f43c585316b5ad3a8f2f16 module/channel-logging
e990d15fb50cf91f684a74993782e00317683aa9f59abcb9dc7b410126b25282 module/channel-message
b28f15439530b0bcc6c696eea3a31864c3eb6702b44f3e92bb10a5038edf996e module/channel-outbound
b28bff5ff9856a25948306de293171655cf9cfd14576ae278a1ab486a858f9a1 module/channel-pairing
13007e28d9e2cb04e7a11098ba0f3ff59c2369bc09ddb160a92714f6f47f31b1 module/channel-plugin-common
47782dd6bb0ff1157ef21fe2ab2503e1d5ec152fb6f45ef645ff69d97a774f02 module/channel-message
34cc2cfd9f0665591db26392ddafd473bea2649f643a0e4a379e4819aa438472 module/channel-outbound
67344a2e73056dfdaed53dc0c67be1e6c3c4fd85c2cdab03107bb60cd242344e module/channel-pairing
33b3c1c9fdfaf31f37f28566e52bee87301d32149e31b40fc0e5349ebfa474ac module/channel-plugin-common
1786ca2cb6867de9fe386b4d2869c2adada32ddf7b047699aba56fa6dfc00ea9 module/channel-policy
b82d23ae060319e1e0cbdc58268024bd4e6a4a3beab208cc515917249c155beb module/channel-reply-pipeline
2f3023c9ff4269d01fc53baa8aeade334bb8149223aec8c7c5937ec09455faf9 module/channel-reply-pipeline
482370e60135db9bfaf07f24bab549e5fde09ab265a6061a1f587c5d93929e91 module/channel-runtime-context
241050b8979ed905727789e99011535631dd8c205aa6967401677d78281900de module/channel-secret-basic-runtime
2edd63d294a8ee128a18c1815f6ae53be915569c49ebbad0dd139eb60475b0e6 module/channel-secret-runtime
c822fb4ced0d7d5fb4bd8e79fed3aaedc457e376ef1af00f2e4446735caf8e6e module/channel-send-result
0116c425524b3476b7bdbfdb7d85b4162834c9ee3df642e7baba2e844257fd08 module/channel-send-result
b2fcc6f55310e17d13cacf15a8c9250f3c834a3795fa2fed998abf940c543ec0 module/channel-setup
2abcf0cba7e06ed6aae6bd4d35ff7f4c7de11332d5532bc2c61f5e666fcb8d28 module/channel-status
f6c25ae55d49d90431682ca53c9dd836304942a738425d6ff028146a0c6d2c97 module/channel-streaming
b2f920ff4a6b4190e6d6ea0a3effb001751e092f0e3ac0cf296721ff8c383d86 module/channel-streaming-config
fdeffe356c7c4edeec9f8fd03edcadc375eabc7a9412e582b10c3180e3ef40fc module/cli-argv
ad12670dbfe538f8d0ebf4fb2b68080e93a760278278e6b1ce9bb129d4b2d533 module/collection-runtime
1e7b3c2313e589380d6d988b14e0f94427382dcfa83fd06e7cb5621f61129b10 module/command-auth
4eda4571e1965b4ba8ce422cd5a2e0bc821430a506912986ff8685edbb0040dc module/command-auth-native
ad8df7a996a185e843cd06ac84bbd085823c1d00136cbbc69fe25832d7825990 module/command-auth
a7375507d07cba14960f20269c51e5c33fbf6fa95dc13adec847a5fabea252f2 module/command-auth-native
50c24235bca2c1d3c011f6bc266b57078b78a76d2d27ee12e7b3245f2947b493 module/command-detection
e24382c2cca7fd2cf69ac353e258153dab80b279daf9f0ef54a727981a7dff3c module/command-primitives-runtime
1d6ebcdc843b7072a7fc4475eadadb90a1111d49c1da1e4326775800873cf9f1 module/command-status
ab86235fcfff7c7cf0021fafeca6e92afe2c257ccf6a4a38441e0d41532998ff module/config-contracts
0d99f5cb8c4978ed760e5fb4e476543759fbd8fd5bf73cc50a50c1d550203826 module/config-mutation
1c79d1356d7f41c22a0e85ff43766fc839cf13734435a3dcf7a4739512d513bf module/config-runtime
a98c3e81edc15b3b9a3b956f70d64eb25f58903c07de98f3622c05f45ebe105a module/config-runtime
d9e5f2ae27e29a40a6d6084c4d59e4811b4669235faa6f0d6a4fbfdd62e51ca3 module/conversation-runtime
807f792edbfe1c7ee2a161a3cd1f34f1ccbd2ea5c3f22d2f1f279fa64968d230 module/core
2ccc2c2011e0703068581ee3547295014d2dfc32b4730c29d66933387df0fbaa module/dedupe-runtime
3e6419ba07df923f75489143efcbde5648f14d7b8e2caeade7541d92227d4886 module/core
3b00dd1e358eeff70b459cadc15f0fd31a4490caf065c2b4af7f3feb2ea83240 module/dedupe-runtime
ebef0e650ab45e44c9335e2b3e15588c968cea6dadd125364a076f9c50ad1e8c module/device-bootstrap
21d86413166ef815581d606f678b6a216a1cc73ffe470b841f5bc4a131bff6df module/diagnostic-runtime
2af0c3b8867148d3deaf0b6d112743a42b206014ed1b5a5cf4ef939f1976d420 module/directory-runtime
5a8a3cf9049b9fac36413b5a4fc33c11b6257c3653bbf8250516880b14bf2b81 module/discord
d0d94423925eb46f3c4c16fc7dd8f9e489c0dbc56f90b2a5cd2d6289b0fd54c0 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
17538c7e143bb6e556cbebc9086de7c2317cbecfa5ebaee64df61681bafb51cb module/inbound-reply-dispatch
ff775bd21ab88126de07c99d507c1bf1a31d8cb93de0516675dbd41d8ff411d2 module/infra-runtime
fbebe0bc2fc26a644db8652c7d2a1160402eb20ad99246959d410fe91f16272c module/inbound-reply-dispatch
164863d241cad378f35d3f01b8d00b19e3a9162b2bd186676eb3cb9867a51f77 module/infra-runtime
ce73721421f1b903dd04ead4df173582e59ea3e9990248102c448b419cc6d272 module/ingress-effect-once
93aa5d6a73719dd9ba6cc65a57e6d48b6c6a1191e3979b84a7d0ffd2c467191f module/interactive-runtime
408d257ab5cc4b88a22b7e7595039cb8fc524b261c44141b294fbd0100ba62ee module/json-store
@@ -85,32 +85,32 @@ cc2d0e1c1e7b9491eb254f314ed2c081a33256936781017eb5d8f9f11808f35b module/logging
f1ca4ced4305d0769c2d8cc1291137ac7002fe0e6eaec2c1a71edad2204c8311 module/matrix
4d18b3bcec2c4291085e6001ed444198fe74f12708956e1c40c964d84e53bec8 module/media-local-roots
f74d7295fe716aa140aa0bc9300d6259d71dab826de0808fca6bb02592bf5d6e module/media-mime
6ca070f839c2f46b29a3b4454f672799370a02660a7b388a3935f746edcd0dd0 module/media-runtime
b3be288c506f146714d3123a89a630f8c4d69f95bc7b543c6b66d1c567af6c58 module/media-runtime
6a52f93107335f88751704352cc01e62add06f854a5b7d765e2a5ee87c0313b6 module/media-store
b7e71516842300c041d2423d822da0080881db6ddc6b9b9ce38cabd8d546676b module/media-understanding
a206a1486f6a6bed3091795b324e95dd710b4b3b87a4a8f18788e6158a48a922 module/media-understanding-runtime
ddd07beaf8a56b140b6eae4ba14b0cd67ace511ed2c926ba7cd445c5850af89a module/meeting-runtime
7e4f82e8924f2420c3852518259a4d83c724556841630d1335bfe36a3bd5a86a module/meeting-runtime
3312468e2e8f3423b765fac6bb17944b800ea2c84acffeb64342c99040b2f482 module/memory-core-host-engine-foundation
8b530380f1ad01d38977fa4d43ee4274516c9a476e0338cd3e2f5df915d0c183 module/memory-host-core
734c9815d91f90283938130abe24d53aa31ef40bab2f9cab69ce118af0fe3842 module/memory-host-core
1efa0aadc4261d1c6073058cbf3dcc9fa681424819bdd14333e19b249bbc4b18 module/messaging-targets
6c43c704519f1178c8ea42d5a0596bec719a81f360b7678ecc00db3ccdb94be7 module/model-session-runtime
339379a69a25d12393fc36863c6c7aabcb2d5b2a9c2dab78281cc5d29d229851 module/models-provider-runtime
f2824762fa0533cbfdfcb77da61ee608bfd108d05a95e9eebdb75b912d39c094 module/model-session-runtime
af2a785fb07f3bd8c6333aa7c58a7070b5cc4f399d30ef31f59ab758a410a6d9 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
280b241f912fe8c405129dbc25a033c53d3f63bc44f09cfbe85a61b867d96b3b module/plugin-entry
4ab0070240ce2b6afff1b8583ccf35a022ede5905ba9a02271181ac112eaee37 module/plugin-runtime
dc834ed4730337a03ba7986daa91eadedfedb54a225d36e329cfded4bb62817f module/provider-auth
39d45f2d8eaff9d634d284ea0d586968cf4c26eec420461ef21874e3c112e762 module/provider-catalog-runtime
2e4a330610c115e16381ee54383c300c3fadddc00d091dcb54b443bc30076908 module/plugin-entry
c6005a9c2c8c6455a593d34d6c0a30cfb88a60f148e4a62052464a1a2bf398ff module/plugin-runtime
744c71dbcacc467da292a43857c68883b0842ca8b022837cd35b58f4acf55221 module/provider-auth
852146d002349be7741a93b1e61b716250425d5a639e99cb77774c7b7ee62367 module/provider-catalog-runtime
8131147d699394bd06503e2ea2f5f1a50b1594a87dded6d118b74a8d0328c8f6 module/proxy-capture
4a698efc36d896c4702de8df831e36b06e85c82aa30bf86c9b1066e6cad4b700 module/question-gateway-runtime
1171a76ea0485b36f77c9e12601c44a0e858043140e1d4669bb86d1f9e35930b module/reply-chunking
7b3e22d8670341edd5cb6b2d5edd092f8afd44455d367670defe12e23998a420 module/reply-dispatch-runtime
cca158d176653ff126b6b855206b42610dc34968d5342635da8eec030e5c1da2 module/reply-dispatch-runtime
73f861fa3179d5af1159853c5acab0eec7a6c8f9398dcb75ea770e784fca6727 module/reply-history
fedbd80588bfc407af75dec097f52d90794e64165e79209253beec5ff615a3a8 module/reply-payload
689c1a6c31f496f10818197ee9e62e6c086caa24dd658da53ec3477073f051b1 module/reply-runtime
7d322a1d48099c4b72dba1417df688d7e4f6d04f4302c54a3de1365385724c18 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
cd04fa399ec561bc4cd2e17ab6fbfc8499ca1e94fcc203cccf3c397856221c56 module/runtime-store
c9d2c8907f0fd6ba90b08119148fc4047043de176ef540262e2bb4b30a41b011 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
d03681d33846765af8a41a498812f0ce7d571b4d37d2f9eff950f8922a248da3 module/session-catalog
7bc12cad4cbf01f43632006a47c3971dc37be1d80354db642d229710746dbcea module/session-discussion
cadcebdf79a8cdbc44d4ecec489babce7fecc52d573e7aee0337a2977b737347 module/session-store-runtime
b1c2023569a3c01a2899a99e67e5f0d89aee5027b7a4617ffab57601be746318 module/session-catalog
99b2bb454b3f515275129fb206792a4b34991be171b927ae30910047deef6bdf module/session-discussion
8d188c2b78e85ce54fc5db9dbd59015ee89ff785fc249bbf069357806014a74c module/session-store-runtime
fb0af0a51ba93e070d8862701c016375cf8576c47be1d25e20537ec86b0bcf73 module/setup
3b37b8371bee9239682a4b2e9c066ef423ffba129d06ded0d7030359ec0dd68e module/setup-runtime
44d37e0d9131ad2859f41068f2604090c784e65f1bd6ebda8e051b6f2e5e1660 module/setup-tools
8ec6ca8a40d4117c669fbd0f20241d54e7e6954d7cff8c5b4afc67ec9626e216 module/skill-commands-runtime
b4b700aa02ce7a77a25d47fbca4feb905fda0caab51409fbf243653663669b11 module/speech-settings
bd0072bb2256c82eab006a71e99ec505201c289899d08b7b225527fff3d694ed module/skill-commands-runtime
c076d03f3832a9549fe7a0e6dc3586678752f7cef6f8d4bcd035d684aea59a32 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
20e8335d6698114f55c9322570e7be3dd7704aeb079c2d85046952783c43ccbf module/tool-plugin
d5f0e369a56b99758185378f8a80546f8a5345515a37b7712284b9a468ff3fd3 module/tool-plugin
dc1a073c59ab61e2789533b777b3f0cb9af689d64a97796b10e8aa82552510db module/tool-results
a5eef5a532c439b1489236711f4bc2721abe612a9db3911430e1344c239f9861 module/tool-send
cda105b721d498df23a554c6b68be150b8fe66b8b9172185c31a0b3b0646b1dc module/web-media
e0b7fae87f43cc90d528df383ce48fbde104af954f5275afe36fe499a796c169 module/webhook-ingress
90790b3039228a62889104a97f22e83e8ebab180d56310d8a349b465ddff8d10 module/webhook-ingress
e3a199a9ce0b85d203e9e8a29b500db29c6b7af307e3145d0a311e29d598925b module/webhook-request-guards
de59e86e126b75d13251cba7ebbe27b44d9b5588785d98df5ff4d6722374c81f module/widget-html
9161b36ec0ab062ea41b363c894fcd672a7727f21cb726739f99f9c184fce69d module/zod
@@ -1,6 +1,8 @@
import {
loadSessionEntry,
replaceSessionEntrySync,
replaceTranscriptEventsSync,
withTranscriptWriteTransaction,
} from "../../config/sessions/session-accessor.js";
import { projectCanonicalSessionEntryShape } from "../../config/sessions/store-entry-shape.js";
import { CURRENT_SESSION_VERSION } from "../../config/sessions/version.js";
@@ -102,7 +104,7 @@ export class SessionManagerBranching extends SessionManagerEntries {
return { entries, opaqueEntries, tailId, usedIds };
}
createBranchedSession(leafId: string): string | undefined {
async createBranchedSession(leafId: string): Promise<string | undefined> {
const previousSessionId = this.sessionId;
const branchPath = this.collectBranchedSessionPath(leafId);
if (branchPath.entries.length === 0) {
@@ -152,35 +154,63 @@ export class SessionManagerBranching extends SessionManagerEntries {
this.fileEntries = [header, ...branchPath.entries, ...labelEntries];
this.opaqueFileEntries = branchPath.opaqueEntries;
this.sessionId = newSessionId;
if (persistenceTarget) {
const updatedAt = Date.now();
const previousEntry = loadSessionEntry({
agentId: persistenceTarget.agentId,
sessionKey: persistenceTarget.sessionKey,
storePath: persistenceTarget.storePath,
});
const canonicalPreviousEntry = previousEntry
? projectCanonicalSessionEntryShape(previousEntry as unknown as Record<string, unknown>)
: { updatedAt };
this.persistenceTarget = { ...persistenceTarget, sessionId: newSessionId };
replaceSessionEntrySync(
{
agentId: persistenceTarget.agentId,
sessionKey: persistenceTarget.sessionKey,
storePath: persistenceTarget.storePath,
},
{
...canonicalPreviousEntry,
sessionId: newSessionId,
updatedAt,
},
);
this.buildIndex();
this.replacePersistedTranscript();
return newSessionId;
this.buildIndex();
if (!persistenceTarget) {
return undefined;
}
this.buildIndex();
return undefined;
const entryScope = {
agentId: persistenceTarget.agentId,
sessionKey: persistenceTarget.sessionKey,
storePath: persistenceTarget.storePath,
};
const previousEntry = loadSessionEntry(entryScope);
const updatedAt = Date.now();
const nextTarget = { ...persistenceTarget, sessionId: newSessionId };
const nextEntry = {
...(previousEntry
? projectCanonicalSessionEntryShape(previousEntry as unknown as Record<string, unknown>)
: { updatedAt }),
sessionId: newSessionId,
updatedAt,
};
try {
const persisted = await withTranscriptWriteTransaction(persistenceTarget, () => {
const currentEntry = loadSessionEntry(entryScope);
if (
currentEntry?.sessionId !== previousSessionId ||
currentEntry.lifecycleRevision !== previousEntry?.lifecycleRevision
) {
return false;
}
replaceSessionEntrySync(entryScope, nextEntry);
if (!replaceTranscriptEventsSync(nextTarget, this.getPersistedFileEntries())) {
throw new Error("Branched session transcript was not persisted");
}
return true;
});
if (!persisted) {
const actualEntry = loadSessionEntry(entryScope);
const cause = actualEntry
? {
actualSessionId: actualEntry.sessionId,
code: "session-rebound" as const,
expectedSessionId: previousSessionId,
sessionKey: persistenceTarget.sessionKey,
}
: {
code: "session-entry-missing" as const,
expectedSessionId: previousSessionId,
sessionKey: persistenceTarget.sessionKey,
};
throw new Error(`Branched session was not persisted: ${cause.code}`, { cause });
}
} catch (error) {
this.setSessionTarget(persistenceTarget);
throw error;
}
this.persistenceTarget = nextTarget;
this.persistenceHeaderPending = false;
return newSessionId;
}
}
+69 -1
View File
@@ -13,6 +13,7 @@ import {
loadTranscriptEvents,
readTranscriptRawDelta,
replaceTranscriptEventsSync,
updateSessionEntry,
upsertSessionEntry,
} from "../../config/sessions/session-accessor.js";
import {
@@ -891,7 +892,7 @@ describe("SessionManager.open", () => {
});
const sessionManager = openMarker(marker, sessionKey, dir);
const branchedMarker = sessionManager.createBranchedSession(assistant.messageId);
const branchedMarker = await sessionManager.createBranchedSession(assistant.messageId);
const branchedSessionId = sessionManager.getSessionId();
expect(branchedMarker).toBe(branchedSessionId);
@@ -923,6 +924,73 @@ describe("SessionManager.open", () => {
]);
});
it("rejects a queued branch when lifecycle ownership changes before persistence", async () => {
const dir = tempDirs.make("openclaw-session-manager-");
const storePath = path.join(dir, "sessions.json");
const sessionId = "sqlite-branch-race-source";
const sessionKey = "agent:main:dashboard:sqlite-branch-race-source";
const marker = formatSqliteSessionFileMarker({ agentId: "main", sessionId, storePath });
const scope = { agentId: "main", sessionId, sessionKey, storePath };
await upsertSessionEntry(scope, {
lifecycleRevision: "branch-original-revision",
sessionFile: marker,
sessionId,
updatedAt: 10,
});
const user = await appendTranscriptMessage(scope, {
cwd: dir,
eventId: "branch-race-user",
message: { role: "user", content: "question before raced branch" },
});
const assistant = await appendTranscriptMessage(scope, {
cwd: dir,
eventId: "branch-race-assistant",
message: buildAssistantMessage("answer before raced branch"),
parentId: user.messageId,
});
const sessionManager = openMarker(marker, sessionKey, dir);
let releaseOwnerChange = () => {};
const ownerChangeGate = new Promise<void>((resolve) => {
releaseOwnerChange = resolve;
});
let markOwnerChangeStarted = () => {};
const ownerChangeStarted = new Promise<void>((resolve) => {
markOwnerChangeStarted = resolve;
});
const ownerChange = updateSessionEntry(scope, async () => {
markOwnerChangeStarted();
await ownerChangeGate;
return { lifecycleRevision: "branch-replacement-revision" };
});
await ownerChangeStarted;
const branch = sessionManager.createBranchedSession(assistant.messageId);
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
releaseOwnerChange();
await ownerChange;
await expect(branch).rejects.toMatchObject({
cause: {
code: "session-rebound",
expectedSessionId: sessionId,
sessionKey,
},
});
expect(loadSessionEntry(scope)).toMatchObject({
lifecycleRevision: "branch-replacement-revision",
sessionId,
});
expect(sessionManager.getSessionId()).toBe(sessionId);
await expect(loadTranscriptEvents(scope)).resolves.toEqual([
expect.objectContaining({ id: sessionId, type: "session" }),
expect.objectContaining({ id: user.messageId, type: "message" }),
expect.objectContaining({ id: assistant.messageId, type: "message" }),
]);
});
it("persists user turns when a SQLite marker has no external recorder", async () => {
const dir = tempDirs.make("openclaw-session-manager-");
const storePath = path.join(dir, "sessions.json");
@@ -2144,6 +2144,7 @@ describe("sqlite session normalization", () => {
const result = await branchSqliteCompactionCheckpointSession({
agentId: "main",
env,
expectedState: sourceEntry,
storePath: paths.sqlitePath,
sourceKey: sourceEntryScope.sessionKey,
nextKey: branchKey,
@@ -2232,6 +2233,7 @@ describe("sqlite session normalization", () => {
const result = await branchSqliteCompactionCheckpointSession({
agentId: "main",
env,
expectedState: { sessionId: "source-session", lifecycleRevision: undefined },
storePath: paths.sqlitePath,
sourceKey: sourceEntryScope.sessionKey,
nextKey: "agent:main:checkpoint-post-fallback",
@@ -2312,6 +2314,7 @@ describe("sqlite session normalization", () => {
const result = await restoreSqliteCompactionCheckpointSession({
agentId: "main",
env,
expectedState: { sessionId: "current-session", lifecycleRevision: undefined },
storePath: paths.sqlitePath,
sessionKey: sourceEntryScope.sessionKey,
checkpointId: checkpoint.checkpointId,
@@ -30,17 +30,20 @@ export function resolveSessionTranscriptActiveLeafEntryId(
export async function rewindSessionToMessage(
params: SessionMessageCutMutationParams,
): Promise<SessionMessageCutMutationResult> {
return await rewindSqliteSessionToMessage(params);
const result = await rewindSqliteSessionToMessage(params);
return result.status === "conflict" ? { status: "failed" } : result;
}
export async function forkSessionAtMessage(
params: SessionMessageCutMutationParams & { targetKey: string },
): Promise<SessionMessageCutMutationResult> {
return await forkSqliteSessionAtMessage(params);
const result = await forkSqliteSessionAtMessage(params);
return result.status === "conflict" ? { status: "failed" } : result;
}
export async function switchSessionBranch(
params: SessionBranchSwitchMutationParams,
): Promise<SessionBranchSwitchMutationResult> {
return await switchSqliteSessionBranch(params);
const result = await switchSqliteSessionBranch(params);
return result.status === "conflict" ? { status: "failed" } : result;
}
@@ -49,6 +49,8 @@ type SqliteCompactionCheckpointLegacySource = {
totalTokens?: number;
};
type SessionEntryExpectedState = Pick<SessionEntry, "lifecycleRevision" | "sessionId">;
/** Result from SQLite compaction checkpoint branch or restore operations. */
type SqliteCompactionCheckpointSessionMutationResult =
| {
@@ -61,6 +63,7 @@ type SqliteCompactionCheckpointSessionMutationResult =
| { status: "missing-checkpoint" }
| { status: "missing-boundary" }
| { status: "model-selection-locked" }
| { status: "conflict" }
| { status: "failed" };
/** Parameters for branching a SQLite session from a compaction checkpoint. */
@@ -72,6 +75,7 @@ type SqliteBranchCheckpointSessionParams = {
sourceStoreKey?: string;
nextKey: string;
checkpointId: string;
expectedState: SessionEntryExpectedState;
legacySource?: SqliteCompactionCheckpointLegacySource;
};
@@ -83,6 +87,7 @@ type SqliteRestoreCheckpointSessionParams = {
sessionKey: string;
sessionStoreKey?: string;
checkpointId: string;
expectedState: SessionEntryExpectedState;
legacySource?: SqliteCompactionCheckpointLegacySource;
};
@@ -110,6 +115,7 @@ export async function branchSqliteCompactionCheckpointSession(
previousIdentity = readSqliteSessionIdentitySnapshot(database, identityKeys);
result = branchSqliteCompactionCheckpointSessionInTransaction(database, {
checkpointId: params.checkpointId,
expectedState: params.expectedState,
parentSessionKey: requestedSourceKey,
legacySource: params.legacySource,
resolved,
@@ -147,6 +153,7 @@ export async function restoreSqliteCompactionCheckpointSession(
previousIdentity = readSqliteSessionIdentitySnapshot(database, identityKeys);
result = restoreSqliteCompactionCheckpointSessionInTransaction(database, {
checkpointId: params.checkpointId,
expectedState: params.expectedState,
legacySource: params.legacySource,
resolved,
sourceKey: sessionKey,
@@ -165,6 +172,7 @@ function branchSqliteCompactionCheckpointSessionInTransaction(
database: OpenClawAgentDatabase,
params: {
checkpointId: string;
expectedState: SessionEntryExpectedState;
legacySource?: SqliteCompactionCheckpointLegacySource;
parentSessionKey: string;
resolved: ResolvedSqliteScope;
@@ -176,6 +184,12 @@ function branchSqliteCompactionCheckpointSessionInTransaction(
if (!currentEntry?.sessionId) {
return { status: "missing-session" };
}
if (
currentEntry.sessionId !== params.expectedState.sessionId ||
currentEntry.lifecycleRevision !== params.expectedState.lifecycleRevision
) {
return { status: "conflict" };
}
if (currentEntry.modelSelectionLocked === true) {
return { status: "model-selection-locked" };
}
@@ -215,6 +229,7 @@ function restoreSqliteCompactionCheckpointSessionInTransaction(
database: OpenClawAgentDatabase,
params: {
checkpointId: string;
expectedState: SessionEntryExpectedState;
legacySource?: SqliteCompactionCheckpointLegacySource;
resolved: ResolvedSqliteScope;
sourceKey: string;
@@ -225,6 +240,12 @@ function restoreSqliteCompactionCheckpointSessionInTransaction(
if (!currentEntry?.sessionId) {
return { status: "missing-session" };
}
if (
currentEntry.sessionId !== params.expectedState.sessionId ||
currentEntry.lifecycleRevision !== params.expectedState.lifecycleRevision
) {
return { status: "conflict" };
}
if (currentEntry.modelSelectionLocked === true) {
return { status: "model-selection-locked" };
}
@@ -21,6 +21,7 @@ import {
readSessionTranscriptMessageEvents,
rewindSessionToMessage,
switchSessionBranch,
updateSessionEntry,
upsertSessionEntry,
} from "./session-accessor.js";
import { listSqliteSessionBranches } from "./session-accessor.sqlite.js";
@@ -29,6 +30,10 @@ import type { InternalSessionEntry } from "./types.js";
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
const agentId = "main";
const sessionKey = "agent:main:message-cut";
const sourceExpectedState = {
lifecycleRevision: "source-lifecycle-revision",
sessionId: "message-cut-source",
};
afterEach(() => {
vi.restoreAllMocks();
@@ -231,7 +236,12 @@ describe("SQLite session message cuts", () => {
const result =
mode === "rewind"
? await rewindSessionToMessage({ agentId, env, entryId: "user-2", sessionKey })
? await rewindSessionToMessage({
agentId,
env,
entryId: "user-2",
sessionKey,
})
: mode === "switch"
? await switchSessionBranch({
agentId,
@@ -267,6 +277,69 @@ describe("SQLite session message cuts", () => {
},
);
it.each(["rewind", "switch", "fork"] as const)(
"rejects %s when the source lifecycle changes in the writer queue",
async (mode) => {
const { env, scope } = await createSession();
let releaseOwnerChange = () => {};
const ownerChangeGate = new Promise<void>((resolve) => {
releaseOwnerChange = resolve;
});
let markOwnerChangeStarted = () => {};
const ownerChangeStarted = new Promise<void>((resolve) => {
markOwnerChangeStarted = resolve;
});
const ownerChange = updateSessionEntry(scope, async () => {
markOwnerChangeStarted();
await ownerChangeGate;
return { lifecycleRevision: "replacement-lifecycle-revision" };
});
await ownerChangeStarted;
const targetKey = `${sessionKey}:raced-fork`;
const mutation =
mode === "rewind"
? rewindSessionToMessage({
agentId,
env,
entryId: "user-2",
sessionKey,
})
: mode === "switch"
? switchSessionBranch({
agentId,
env,
leafEntryId: "off-path-user",
sessionKey,
})
: forkSessionAtMessage({
agentId,
env,
entryId: "user-2",
sessionKey,
targetKey,
});
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
releaseOwnerChange();
await ownerChange;
await expect(mutation).resolves.toEqual({ status: "failed" });
expect(loadSessionEntry(scope)).toMatchObject({
lifecycleRevision: "replacement-lifecycle-revision",
sessionId: sourceExpectedState.sessionId,
});
expect(loadSessionEntry({ agentId, env, sessionKey: targetKey })).toBeUndefined();
await expect(listSessionBranches({ agentId, env, sessionKey })).resolves.toMatchObject({
status: "ok",
branches: expect.arrayContaining([
expect.objectContaining({ active: true, leafEntryId: "assistant-2" }),
]),
});
},
);
it("keeps branch summaries isolated between sessions in the same store", async () => {
const { env } = await createSession();
const sibling = await createSiblingSession({
@@ -374,7 +447,12 @@ describe("SQLite session message cuts", () => {
const { env } = await createSession();
await expect(
switchSessionBranch({ agentId, env, leafEntryId, sessionKey }),
switchSessionBranch({
agentId,
env,
leafEntryId,
sessionKey,
}),
).resolves.toMatchObject({ status });
});
@@ -528,7 +606,12 @@ describe("SQLite session message cuts", () => {
const { env } = await createSession();
await expect(
rewindSessionToMessage({ agentId, env, entryId, sessionKey }),
rewindSessionToMessage({
agentId,
env,
entryId,
sessionKey,
}),
).resolves.toMatchObject({ status });
});
});
@@ -58,9 +58,11 @@ type MessageCut = {
type SessionTranscriptMutationResult =
| SessionMessageCutMutationResult
| SessionBranchSwitchMutationResult;
| SessionBranchSwitchMutationResult
| { status: "conflict" };
type SessionTranscriptMutationMode = "fork" | "rewind" | "switch";
type SessionEntryExpectedState = Pick<SessionEntry, "lifecycleRevision" | "sessionId">;
const BRANCH_HEADLINE_MAX_CHARS = 120;
const SESSION_BRANCH_CACHE_MAX_ENTRIES = 32;
@@ -167,34 +169,44 @@ export function resolveSessionTranscriptActiveLeafEntryId(
export async function rewindSqliteSessionToMessage(
params: SessionMessageCutMutationParams,
): Promise<SessionMessageCutMutationResult> {
return await mutateSqliteSessionAtMessage(params, "rewind");
expectedState?: SessionEntryExpectedState,
): Promise<SessionMessageCutMutationResult | { status: "conflict" }> {
return await mutateSqliteSessionAtMessage(params, "rewind", expectedState);
}
export async function forkSqliteSessionAtMessage(
params: SessionMessageCutMutationParams & { targetKey: string },
): Promise<SessionMessageCutMutationResult> {
return await mutateSqliteSessionAtMessage(params, "fork");
expectedState?: SessionEntryExpectedState,
): Promise<SessionMessageCutMutationResult | { status: "conflict" }> {
return await mutateSqliteSessionAtMessage(params, "fork", expectedState);
}
export async function switchSqliteSessionBranch(
params: SessionBranchSwitchMutationParams,
): Promise<SessionBranchSwitchMutationResult> {
return await mutateSqliteSessionAtMessage({ ...params, entryId: params.leafEntryId }, "switch");
expectedState?: SessionEntryExpectedState,
): Promise<SessionBranchSwitchMutationResult | { status: "conflict" }> {
return await mutateSqliteSessionAtMessage(
{ ...params, entryId: params.leafEntryId },
"switch",
expectedState,
);
}
function mutateSqliteSessionAtMessage(
params: SessionMessageCutMutationParams,
mode: "fork" | "rewind",
): Promise<SessionMessageCutMutationResult>;
expectedState?: SessionEntryExpectedState,
): Promise<SessionMessageCutMutationResult | { status: "conflict" }>;
function mutateSqliteSessionAtMessage(
params: SessionMessageCutMutationParams,
mode: "switch",
): Promise<SessionBranchSwitchMutationResult>;
expectedState?: SessionEntryExpectedState,
): Promise<SessionBranchSwitchMutationResult | { status: "conflict" }>;
async function mutateSqliteSessionAtMessage(
params: SessionMessageCutMutationParams,
mode: SessionTranscriptMutationMode,
expectedState?: SessionEntryExpectedState,
): Promise<SessionTranscriptMutationResult> {
const canonicalSourceKey = normalizeSqliteSessionKey(params.sessionKey);
const sourceKey = normalizeSqliteSessionKey(params.sessionStoreKey ?? params.sessionKey);
@@ -206,6 +218,18 @@ async function mutateSqliteSessionAtMessage(
sessionKey: sourceKey,
...(params.storePath ? { storePath: params.storePath } : {}),
});
const preparedEntry = readSessionEntryRow(
openOpenClawAgentDatabase(toDatabaseOptions(resolved)),
sourceKey,
)?.entry;
const preparedExpectedState =
expectedState ??
(preparedEntry?.sessionId
? {
sessionId: preparedEntry.sessionId,
lifecycleRevision: preparedEntry.lifecycleRevision,
}
: undefined);
return await runExclusiveSqliteSessionWrite(resolved, async () => {
let previousIdentity = new Map<string, SessionEntry>();
let currentIdentity = new Map<string, SessionEntry>();
@@ -222,6 +246,7 @@ async function mutateSqliteSessionAtMessage(
canonicalSourceKey,
creation: params.creation,
mode,
expectedState: preparedExpectedState,
sourceKey,
targetKey,
});
@@ -248,6 +273,7 @@ function mutateSqliteSessionAtMessageInTransaction(
canonicalSourceKey: string;
creation?: SessionMessageCutMutationParams["creation"];
entryId: string;
expectedState: SessionEntryExpectedState | undefined;
mode: SessionTranscriptMutationMode;
sourceKey: string;
targetKey: string;
@@ -257,6 +283,13 @@ function mutateSqliteSessionAtMessageInTransaction(
if (!currentEntry?.sessionId) {
return { status: "missing-session" };
}
if (
!params.expectedState ||
currentEntry.sessionId !== params.expectedState.sessionId ||
currentEntry.lifecycleRevision !== params.expectedState.lifecycleRevision
) {
return { status: "conflict" };
}
const events = loadSqliteTranscriptEventsFromDatabase(database, currentEntry.sessionId);
const cut = params.mode === "switch" ? undefined : resolveMessageCut(events, params.entryId);
if (cut && "status" in cut) {
@@ -34,6 +34,22 @@ const compactionCheckpointStore = createFileBackedCompactionCheckpointStore();
const MODEL_SELECTION_LOCKED_CHECKPOINT_MESSAGE =
"Checkpoint branch and restore are unavailable while model selection is locked.";
function respondCheckpointConflict(
key: string,
action: "branch" | "restore",
respond: Parameters<GatewayRequestHandlers[string]>[0]["respond"],
): void {
respond(
false,
undefined,
errorShape(
ErrorCodes.INVALID_REQUEST,
`Session ${key} changed before checkpoint ${action}. Retry.`,
{ details: { reason: SESSION_LIFECYCLE_CHANGED_ERROR_REASON } },
),
);
}
export const sessionCheckpointHandlers: GatewayRequestHandlers = {
"sessions.compaction.branch": async ({ params, respond, context }) => {
if (
@@ -89,6 +105,10 @@ export const sessionCheckpointHandlers: GatewayRequestHandlers = {
const nextKey = buildDashboardSessionKey(target.agentId);
const branchedSession = await compactionCheckpointStore.branchCheckpointSession({
agentId: target.agentId,
expectedState: {
sessionId: entry.sessionId,
lifecycleRevision: entry.lifecycleRevision,
},
storePath,
sourceKey: canonicalKey,
sourceStoreKey: sessionStoreKey,
@@ -122,6 +142,10 @@ export const sessionCheckpointHandlers: GatewayRequestHandlers = {
);
return;
}
if (branchedSession.status === "conflict") {
respondCheckpointConflict(key, "branch", respond);
return;
}
if (branchedSession.status === "failed") {
respond(
false,
@@ -364,6 +388,10 @@ export const sessionCheckpointHandlers: GatewayRequestHandlers = {
const restoredSession = await compactionCheckpointStore.restoreCheckpointSession({
agentId: requestedAgent.agentId,
expectedState: {
sessionId: current.entry.sessionId,
lifecycleRevision: current.entry.lifecycleRevision,
},
storePath,
sessionKey: current.canonicalKey,
sessionStoreKey: current.sessionStoreKey,
@@ -396,6 +424,10 @@ export const sessionCheckpointHandlers: GatewayRequestHandlers = {
);
return;
}
if (restoredSession.status === "conflict") {
respondCheckpointConflict(key, "restore", respond);
return;
}
if (restoredSession.status === "failed") {
respond(
false,
+60 -42
View File
@@ -11,14 +11,16 @@ import { resolveDefaultAgentId } from "../../agents/agent-scope.js";
import { listRegisteredAgentHarnesses } from "../../agents/harness/registry.js";
import { clearSessionQueues } from "../../auto-reply/reply/queue/cleanup.js";
import {
forkSessionAtMessage,
listSessionBranches,
rewindSessionToMessage,
switchSessionBranch,
type SessionBranchListResult,
type SessionBranchSwitchMutationResult,
type SessionMessageCutMutationResult,
} from "../../config/sessions/session-accessor.js";
import {
forkSqliteSessionAtMessage as forkSessionAtMessage,
rewindSqliteSessionToMessage as rewindSessionToMessage,
switchSqliteSessionBranch as switchSessionBranch,
} from "../../config/sessions/session-accessor.sqlite.js";
import { MEDIA_MAX_BYTES, readMediaBuffer } from "../../media/store.js";
import {
isCompetingSessionWorkAdmissionActive,
@@ -44,6 +46,10 @@ import type { GatewayRequestHandlerOptions, GatewayRequestHandlers } from "./typ
import { assertValidParams } from "./validation.js";
type MessageCutAction = "fork" | "rewind" | "switch";
type MessageCutMutationResult =
| SessionMessageCutMutationResult
| SessionBranchSwitchMutationResult
| { status: "conflict" };
const EXTERNAL_CONVERSATION_ERROR =
"Session history changes are unavailable because this session is owned by an external agent harness.";
@@ -344,6 +350,10 @@ async function mutateSessionAtMessage(
}
const targetKey =
action === "fork" ? buildDashboardSessionKey(current.target.agentId) : current.canonicalKey;
const expectedState = {
sessionId: current.entry.sessionId,
lifecycleRevision: current.entry.lifecycleRevision,
};
const upstreamForkHarness = upstreamLink
? resolveUpstreamForkHarness(upstreamLink)
: undefined;
@@ -411,33 +421,42 @@ async function mutateSessionAtMessage(
});
return;
}
let result: SessionMessageCutMutationResult | SessionBranchSwitchMutationResult;
let result: MessageCutMutationResult;
try {
result = await (action === "fork"
? forkSessionAtMessage({
agentId: current.target.agentId,
entryId,
sessionKey: current.canonicalKey,
sessionStoreKey: current.sessionStoreKey,
storePath: current.storePath,
targetKey,
creation: resolveOperatorSessionCreation(client),
})
: action === "rewind"
? rewindSessionToMessage({
? forkSessionAtMessage(
{
agentId: current.target.agentId,
entryId,
sessionKey: current.canonicalKey,
sessionStoreKey: current.sessionStoreKey,
storePath: current.storePath,
})
: switchSessionBranch({
agentId: current.target.agentId,
leafEntryId: entryId,
sessionKey: current.canonicalKey,
sessionStoreKey: current.sessionStoreKey,
storePath: current.storePath,
}));
targetKey,
creation: resolveOperatorSessionCreation(client),
},
expectedState,
)
: action === "rewind"
? rewindSessionToMessage(
{
agentId: current.target.agentId,
entryId,
sessionKey: current.canonicalKey,
sessionStoreKey: current.sessionStoreKey,
storePath: current.storePath,
},
expectedState,
)
: switchSessionBranch(
{
agentId: current.target.agentId,
leafEntryId: entryId,
sessionKey: current.canonicalKey,
sessionStoreKey: current.sessionStoreKey,
storePath: current.storePath,
},
expectedState,
));
} catch {
respond(
false,
@@ -501,31 +520,30 @@ async function mutateSessionAtMessage(
}
function respondMessageCutError(
result: Exclude<
SessionMessageCutMutationResult | SessionBranchSwitchMutationResult,
{ status: "created" }
>,
result: Exclude<MessageCutMutationResult, { status: "created" }>,
action: MessageCutAction,
entryId: string,
respond: GatewayRequestHandlerOptions["respond"],
): void {
const actionLabel = action === "switch" ? "branch switch" : action;
const message =
result.status === "missing-session"
? "session not found"
: result.status === "missing-entry"
? `${action === "switch" ? "branch" : "message"} entry not found: ${entryId}`
: result.status === "not-branch-tip"
? `entry is not a branch tip: ${entryId}`
: result.status === "already-active"
? `branch is already active: ${entryId}`
: result.status === "not-user-message"
? `entry is not a user message: ${entryId}`
: result.status === "off-active-path"
? `message entry is not on the active path: ${entryId}`
: result.status === "unsupported-storage"
? `session transcript storage does not support ${actionLabel}`
: `failed to ${actionLabel} session`;
result.status === "conflict"
? `Session changed; retry ${action}.`
: result.status === "missing-session"
? "session not found"
: result.status === "missing-entry"
? `${action === "switch" ? "branch" : "message"} entry not found: ${entryId}`
: result.status === "not-branch-tip"
? `entry is not a branch tip: ${entryId}`
: result.status === "already-active"
? `branch is already active: ${entryId}`
: result.status === "not-user-message"
? `entry is not a user message: ${entryId}`
: result.status === "off-active-path"
? `message entry is not on the active path: ${entryId}`
: result.status === "unsupported-storage"
? `session transcript storage does not support ${actionLabel}`
: `failed to ${actionLabel} session`;
respond(
false,
undefined,
@@ -13,7 +13,9 @@ import { formatSqliteSessionFileMarker } from "../config/sessions/legacy-sqlite-
import {
appendTranscriptEvent,
appendTranscriptMessage,
loadSessionEntry,
loadTranscriptEvents,
updateSessionEntry,
upsertSessionEntry,
} from "../config/sessions/session-accessor.js";
import {
@@ -26,6 +28,10 @@ const tempDirs: string[] = [];
const MAIN_AGENT_ID = "main";
const MAIN_SESSION_KEY = "agent:main:main";
function checkpointExpectedState(sessionId: string) {
return { lifecycleRevision: undefined, sessionId };
}
function requireNonEmptyString(value: string | null | undefined, message: string): string {
if (!value) {
throw new Error(message);
@@ -141,6 +147,7 @@ describe("session-compaction-checkpoints", () => {
const store = createFileBackedCompactionCheckpointStore();
const branchKey = "agent:main:checkpoint-branch";
const branched = await store.branchCheckpointSession({
expectedState: checkpointExpectedState(sessionId),
storePath,
sourceKey: sessionKey,
sourceStoreKey: sessionStoreKey,
@@ -148,6 +155,7 @@ describe("session-compaction-checkpoints", () => {
checkpointId: checkpoint.checkpointId,
});
const restored = await store.restoreCheckpointSession({
expectedState: checkpointExpectedState(sessionId),
storePath,
sessionKey,
sessionStoreKey,
@@ -181,6 +189,105 @@ describe("session-compaction-checkpoints", () => {
).toBe(true);
});
test.each(["branch", "restore"] as const)(
"checkpoint %s rejects a lifecycle change queued before its transaction",
async (mode) => {
const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-checkpoint-sqlite-race-"));
tempDirs.push(dir);
const storePath = path.join(dir, "openclaw-agent.sqlite");
const sessionId = `sqlite-checkpoint-${mode}-race`;
const sessionKey = MAIN_SESSION_KEY;
const scope = { agentId: MAIN_AGENT_ID, sessionId, sessionKey, storePath };
const expectedState = {
lifecycleRevision: "checkpoint-original-revision",
sessionId,
};
await upsertSessionEntry(scope, {
...expectedState,
updatedAt: 10,
});
await appendTranscriptEvent(scope, {
type: "session",
version: CURRENT_SESSION_VERSION,
id: sessionId,
timestamp: "2026-06-26T12:00:00.000Z",
cwd: dir,
});
const sourceMessage = await appendTranscriptMessage(scope, {
message: { role: "user", content: "checkpoint race source", timestamp: 1 },
now: Date.parse("2026-06-26T12:00:01.000Z"),
});
const checkpoint: SessionCompactionCheckpoint = {
checkpointId: `sqlite-checkpoint-${mode}-conflict`,
sessionKey,
sessionId,
createdAt: Date.now(),
reason: "manual",
preCompaction: {
sessionId,
leafId: sourceMessage.messageId,
entryId: sourceMessage.messageId,
},
postCompaction: {
sessionId,
leafId: sourceMessage.messageId,
entryId: sourceMessage.messageId,
},
};
await upsertSessionEntry(scope, { compactionCheckpoints: [checkpoint] });
let releaseOwnerChange = () => {};
const ownerChangeGate = new Promise<void>((resolve) => {
releaseOwnerChange = resolve;
});
let markOwnerChangeStarted = () => {};
const ownerChangeStarted = new Promise<void>((resolve) => {
markOwnerChangeStarted = resolve;
});
const ownerChange = updateSessionEntry(scope, async () => {
markOwnerChangeStarted();
await ownerChangeGate;
return { lifecycleRevision: "checkpoint-replacement-revision" };
});
await ownerChangeStarted;
const branchKey = `${sessionKey}:${mode}-conflict`;
const store = createFileBackedCompactionCheckpointStore();
const mutation =
mode === "branch"
? store.branchCheckpointSession({
expectedState,
storePath,
sourceKey: sessionKey,
nextKey: branchKey,
checkpointId: checkpoint.checkpointId,
})
: store.restoreCheckpointSession({
expectedState,
storePath,
sessionKey,
checkpointId: checkpoint.checkpointId,
});
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
releaseOwnerChange();
await ownerChange;
await expect(mutation).resolves.toEqual({ status: "conflict" });
expect(loadSessionEntry(scope)).toMatchObject({
lifecycleRevision: "checkpoint-replacement-revision",
sessionId,
});
expect(loadSessionEntry({ agentId: MAIN_AGENT_ID, sessionKey: branchKey, storePath })).toBe(
undefined,
);
await expect(loadTranscriptEvents(scope)).resolves.toEqual(
expect.arrayContaining([expect.objectContaining({ id: sourceMessage.messageId })]),
);
},
);
test("checkpoint store branches row-backed checkpoints when entry sessionFile is stale", async () => {
const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-checkpoint-sqlite-stale-"));
tempDirs.push(dir);
@@ -274,6 +381,7 @@ describe("session-compaction-checkpoints", () => {
const branchKey = "agent:main:stale-checkpoint-branch";
const branched = await createFileBackedCompactionCheckpointStore().branchCheckpointSession({
expectedState: checkpointExpectedState(sessionId),
storePath,
sourceKey: sessionKey,
nextKey: branchKey,
@@ -297,6 +405,7 @@ describe("session-compaction-checkpoints", () => {
const markerBranched =
await createFileBackedCompactionCheckpointStore().branchCheckpointSession({
expectedState: checkpointExpectedState(sessionId),
storePath,
sourceKey: sessionKey,
nextKey: "agent:main:stale-marker-checkpoint-branch",
@@ -370,6 +479,7 @@ describe("session-compaction-checkpoints", () => {
);
const branched = await createFileBackedCompactionCheckpointStore().branchCheckpointSession({
expectedState: checkpointExpectedState(sessionId),
storePath,
sourceKey: sessionKey,
nextKey: "agent:main:legacy-checkpoint-branch",
@@ -66,10 +66,14 @@ export function resolveCompactionCheckpointTranscriptPosition(params: {
};
}
type CompactionCheckpointSessionMutationResult = SessionCompactionCheckpointMutationResult;
type CompactionCheckpointSessionMutationResult =
| SessionCompactionCheckpointMutationResult
| { status: "conflict" };
type SessionEntryExpectedState = Pick<SessionEntry, "lifecycleRevision" | "sessionId">;
type BranchCheckpointSessionParams = {
agentId?: string;
expectedState: SessionEntryExpectedState;
storePath: string;
sourceKey: string;
sourceStoreKey?: string;
@@ -79,6 +83,7 @@ type BranchCheckpointSessionParams = {
type RestoreCheckpointSessionParams = {
agentId?: string;
expectedState: SessionEntryExpectedState;
storePath: string;
sessionKey: string;
sessionStoreKey?: string;
@@ -540,6 +545,7 @@ async function branchCheckpointSessionFromStoredBoundary(
sourceKey: params.sourceKey,
nextKey: params.nextKey,
checkpointId: params.checkpointId,
expectedState: params.expectedState,
...(params.sourceStoreKey ? { sourceStoreKey: params.sourceStoreKey } : {}),
...(legacySource ? { legacySource } : {}),
});
@@ -561,6 +567,7 @@ async function restoreCheckpointSessionFromStoredBoundary(
storePath: params.storePath,
sessionKey: params.sessionKey,
checkpointId: params.checkpointId,
expectedState: params.expectedState,
...(params.sessionStoreKey ? { sessionStoreKey: params.sessionStoreKey } : {}),
...(legacySource ? { legacySource } : {}),
});
@@ -12,6 +12,7 @@ import {
loadSessionEntry,
loadTranscriptEvents,
resolveSessionTranscriptRuntimeTarget,
updateSessionEntry,
upsertSessionEntry,
} from "../../config/sessions/session-accessor.js";
import type { OpenClawConfig } from "../../config/types.openclaw.js";
@@ -177,6 +178,7 @@ describe("worker transcript commit application", () => {
await upsertSessionEntry(
{ agentId: "main", sessionKey: SESSION_KEY, storePath },
{
lifecycleRevision: "worker-original-revision",
sessionId: SESSION_ID,
updatedAt: 10,
},
@@ -362,6 +364,42 @@ describe("worker transcript commit application", () => {
expect(reopened.getLeafId()).toBe(first.result.newLeafId);
});
it("rejects a commit when lifecycle ownership changes in the writer queue", async () => {
let releaseOwnerChange = () => {};
const ownerChangeGate = new Promise<void>((resolve) => {
releaseOwnerChange = resolve;
});
let markOwnerChangeStarted = () => {};
const ownerChangeStarted = new Promise<void>((resolve) => {
markOwnerChangeStarted = resolve;
});
const ownerChange = updateSessionEntry(
{ agentId: "main", sessionKey: SESSION_KEY, storePath },
async () => {
markOwnerChangeStarted();
await ownerChangeGate;
return { lifecycleRevision: "worker-replacement-revision" };
},
);
await ownerChangeStarted;
const commit = committer.commit({ identity: IDENTITY, request: createRequest() });
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
releaseOwnerChange();
await ownerChange;
await expect(commit).resolves.toEqual({ ok: false, reason: "invalid-batch" });
expect(loadSessionEntry({ agentId: "main", sessionKey: SESSION_KEY, storePath })).toMatchObject(
{
lifecycleRevision: "worker-replacement-revision",
sessionId: SESSION_ID,
},
);
expect(SessionManager.open(sessionTarget).getEntries()).toEqual([]);
});
it("replays the same tuple without duplicates and rejects a changed payload", async () => {
const request = createRequest();
const first = await committer.commit({ identity: IDENTITY, request });
@@ -47,6 +47,8 @@ type PersistedCommitResolution =
| { kind: "ambiguous" | "missing" }
| { kind: "found"; messages: AppliedTranscriptMessage[] };
const WORKER_TRANSCRIPT_SESSION_CONFLICT = new Error("worker transcript session changed");
function cloneContentPart(
part: WorkerTranscriptMessage["content"][number],
): WorkerTranscriptMessage["content"][number] {
@@ -320,67 +322,88 @@ async function applyWorkerTranscriptCommit(params: {
const redactedMessages = params.messages.map(
(message) => redactTranscriptMessage(message, params.config) as CommittedAgentMessage,
);
const applied = await withTranscriptWriteTransaction(params.target, (transcriptTarget) => {
const currentEntry = loadSessionEntry(params.target);
if (!currentEntry || currentEntry.sessionId !== params.sessionId) {
return { ok: false as const, reason: "session-not-attached" as const };
}
const expectedState = {
sessionId: params.sessionId,
lifecycleRevision: params.target.sessionEntry.lifecycleRevision,
};
let applied: ApplyTranscriptCommitResult;
try {
applied = await withTranscriptWriteTransaction(params.target, (transcriptTarget) => {
const currentEntry = loadSessionEntry(params.target);
if (!currentEntry || currentEntry.sessionId !== expectedState.sessionId) {
return { ok: false as const, reason: "session-not-attached" as const };
}
if (currentEntry.lifecycleRevision !== expectedState.lifecycleRevision) {
return { ok: false as const, reason: "invalid-batch" as const };
}
const manager = SessionManager.open(transcriptTarget);
if (params.recoverPersistedBatch) {
// Only a pending ledger row may prove an off-branch batch: the agent DB
// can commit before the shared replay ledger records its terminal result.
const recovered = resolvePersistedCommitAcrossDag({
const manager = SessionManager.open(transcriptTarget);
if (params.recoverPersistedBatch) {
// Only a pending ledger row may prove an off-branch batch: the agent DB
// can commit before the shared replay ledger records its terminal result.
const recovered = resolvePersistedCommitAcrossDag({
baseLeafId: params.requestedBaseLeafId,
manager,
messages: redactedMessages,
});
if (recovered.kind === "found") {
return { ok: true as const, messages: recovered.messages };
}
if (recovered.kind === "ambiguous") {
return { ok: false as const, reason: "invalid-batch" as const };
}
}
const prefix = resolveActiveCommitPrefix({
baseLeafId: params.requestedBaseLeafId,
manager,
messages: redactedMessages,
});
if (recovered.kind === "found") {
return { ok: true as const, messages: recovered.messages };
if (!prefix.ok) {
return { ok: false as const, reason: "stale-base-leaf" as const };
}
if (recovered.kind === "ambiguous") {
return { ok: false as const, reason: "invalid-batch" as const };
const messages = [...prefix.recoveredMessages];
let nextMessageSeq = prefix.activeVisibleEntryCount;
for (const message of redactedMessages.slice(prefix.recoveredMessages.length)) {
const messageId = manager.appendMessage(message, {
config: params.config,
// Active-path recovery owns dedupe. A global key scan could reuse an
// id from an abandoned branch while SessionManager advances another id.
idempotencyLookup: "caller-checked",
});
nextMessageSeq += 1;
messages.push({
appended: true,
message,
messageId,
messageSeq: nextMessageSeq,
});
}
}
const prefix = resolveActiveCommitPrefix({
baseLeafId: params.requestedBaseLeafId,
manager,
messages: redactedMessages,
const freshEntry = loadSessionEntry(params.target);
if (
!freshEntry ||
freshEntry.sessionId !== expectedState.sessionId ||
freshEntry.lifecycleRevision !== expectedState.lifecycleRevision
) {
throw WORKER_TRANSCRIPT_SESSION_CONFLICT;
}
const appendedCount = messages.filter((message) => message.appended).length;
const nextEntry = {
...freshEntry,
...(appendedCount > 0
? { updatedAt: Math.max(freshEntry.updatedAt ?? 0, Date.now()) }
: {}),
};
replaceSessionEntrySync(params.target, nextEntry);
return { ok: true as const, messages };
});
if (!prefix.ok) {
return { ok: false as const, reason: "stale-base-leaf" as const };
} catch (error) {
if (error === WORKER_TRANSCRIPT_SESSION_CONFLICT) {
return { ok: false, reason: "invalid-batch" };
}
const messages = [...prefix.recoveredMessages];
let nextMessageSeq = prefix.activeVisibleEntryCount;
for (const message of redactedMessages.slice(prefix.recoveredMessages.length)) {
const messageId = manager.appendMessage(message, {
config: params.config,
// Active-path recovery owns dedupe. A global key scan could reuse an
// id from an abandoned branch while SessionManager advances another id.
idempotencyLookup: "caller-checked",
});
nextMessageSeq += 1;
messages.push({
appended: true,
message,
messageId,
messageSeq: nextMessageSeq,
});
}
const freshEntry = loadSessionEntry(params.target);
if (!freshEntry || freshEntry.sessionId !== params.sessionId) {
return { ok: false as const, reason: "session-not-attached" as const };
}
const appendedCount = messages.filter((message) => message.appended).length;
const nextEntry = {
...freshEntry,
...(appendedCount > 0 ? { updatedAt: Math.max(freshEntry.updatedAt ?? 0, Date.now()) } : {}),
};
replaceSessionEntrySync(params.target, nextEntry);
return { ok: true as const, messages };
});
throw error;
}
if (!applied.ok) {
return applied;
}