diff --git a/docs/.generated/plugin-sdk-api-baseline.sha256 b/docs/.generated/plugin-sdk-api-baseline.sha256 index 5a9e4af947c9..0ba6caca357a 100644 --- a/docs/.generated/plugin-sdk-api-baseline.sha256 +++ b/docs/.generated/plugin-sdk-api-baseline.sha256 @@ -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 diff --git a/src/agents/sessions/session-manager-branching.ts b/src/agents/sessions/session-manager-branching.ts index a63bddf131d5..00cb9c3c1d9d 100644 --- a/src/agents/sessions/session-manager-branching.ts +++ b/src/agents/sessions/session-manager-branching.ts @@ -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 { 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) - : { 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) + : { 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; } } diff --git a/src/agents/sessions/session-manager.test.ts b/src/agents/sessions/session-manager.test.ts index 65475f084c8c..3beedae91956 100644 --- a/src/agents/sessions/session-manager.test.ts +++ b/src/agents/sessions/session-manager.test.ts @@ -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((resolve) => { + releaseOwnerChange = resolve; + }); + let markOwnerChangeStarted = () => {}; + const ownerChangeStarted = new Promise((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((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"); diff --git a/src/config/sessions/session-accessor.conformance.test.ts b/src/config/sessions/session-accessor.conformance.test.ts index 82fa664b86b8..3d163992d154 100644 --- a/src/config/sessions/session-accessor.conformance.test.ts +++ b/src/config/sessions/session-accessor.conformance.test.ts @@ -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, diff --git a/src/config/sessions/session-accessor.message-cut.ts b/src/config/sessions/session-accessor.message-cut.ts index 4dfff7bbde0f..996c733f4da3 100644 --- a/src/config/sessions/session-accessor.message-cut.ts +++ b/src/config/sessions/session-accessor.message-cut.ts @@ -30,17 +30,20 @@ export function resolveSessionTranscriptActiveLeafEntryId( export async function rewindSessionToMessage( params: SessionMessageCutMutationParams, ): Promise { - 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 { - return await forkSqliteSessionAtMessage(params); + const result = await forkSqliteSessionAtMessage(params); + return result.status === "conflict" ? { status: "failed" } : result; } export async function switchSessionBranch( params: SessionBranchSwitchMutationParams, ): Promise { - return await switchSqliteSessionBranch(params); + const result = await switchSqliteSessionBranch(params); + return result.status === "conflict" ? { status: "failed" } : result; } diff --git a/src/config/sessions/session-accessor.sqlite-checkpoint.ts b/src/config/sessions/session-accessor.sqlite-checkpoint.ts index 239424c6fbdf..c3ca0f404632 100644 --- a/src/config/sessions/session-accessor.sqlite-checkpoint.ts +++ b/src/config/sessions/session-accessor.sqlite-checkpoint.ts @@ -49,6 +49,8 @@ type SqliteCompactionCheckpointLegacySource = { totalTokens?: number; }; +type SessionEntryExpectedState = Pick; + /** 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" }; } diff --git a/src/config/sessions/session-accessor.sqlite-message-cut.test.ts b/src/config/sessions/session-accessor.sqlite-message-cut.test.ts index 7bc043a9a059..234741543797 100644 --- a/src/config/sessions/session-accessor.sqlite-message-cut.test.ts +++ b/src/config/sessions/session-accessor.sqlite-message-cut.test.ts @@ -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((resolve) => { + releaseOwnerChange = resolve; + }); + let markOwnerChangeStarted = () => {}; + const ownerChangeStarted = new Promise((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((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 }); }); }); diff --git a/src/config/sessions/session-accessor.sqlite-message-cut.ts b/src/config/sessions/session-accessor.sqlite-message-cut.ts index ef81cef32bd7..ebdcf816eb0e 100644 --- a/src/config/sessions/session-accessor.sqlite-message-cut.ts +++ b/src/config/sessions/session-accessor.sqlite-message-cut.ts @@ -58,9 +58,11 @@ type MessageCut = { type SessionTranscriptMutationResult = | SessionMessageCutMutationResult - | SessionBranchSwitchMutationResult; + | SessionBranchSwitchMutationResult + | { status: "conflict" }; type SessionTranscriptMutationMode = "fork" | "rewind" | "switch"; +type SessionEntryExpectedState = Pick; 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 { - return await mutateSqliteSessionAtMessage(params, "rewind"); + expectedState?: SessionEntryExpectedState, +): Promise { + return await mutateSqliteSessionAtMessage(params, "rewind", expectedState); } export async function forkSqliteSessionAtMessage( params: SessionMessageCutMutationParams & { targetKey: string }, -): Promise { - return await mutateSqliteSessionAtMessage(params, "fork"); + expectedState?: SessionEntryExpectedState, +): Promise { + return await mutateSqliteSessionAtMessage(params, "fork", expectedState); } export async function switchSqliteSessionBranch( params: SessionBranchSwitchMutationParams, -): Promise { - return await mutateSqliteSessionAtMessage({ ...params, entryId: params.leafEntryId }, "switch"); + expectedState?: SessionEntryExpectedState, +): Promise { + return await mutateSqliteSessionAtMessage( + { ...params, entryId: params.leafEntryId }, + "switch", + expectedState, + ); } function mutateSqliteSessionAtMessage( params: SessionMessageCutMutationParams, mode: "fork" | "rewind", -): Promise; + expectedState?: SessionEntryExpectedState, +): Promise; function mutateSqliteSessionAtMessage( params: SessionMessageCutMutationParams, mode: "switch", -): Promise; + expectedState?: SessionEntryExpectedState, +): Promise; async function mutateSqliteSessionAtMessage( params: SessionMessageCutMutationParams, mode: SessionTranscriptMutationMode, + expectedState?: SessionEntryExpectedState, ): Promise { 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(); let currentIdentity = new Map(); @@ -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) { diff --git a/src/gateway/server-methods/sessions-compaction-checkpoints.ts b/src/gateway/server-methods/sessions-compaction-checkpoints.ts index 9591d488191f..739a8a4a6016 100644 --- a/src/gateway/server-methods/sessions-compaction-checkpoints.ts +++ b/src/gateway/server-methods/sessions-compaction-checkpoints.ts @@ -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[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, diff --git a/src/gateway/server-methods/sessions-rewind.ts b/src/gateway/server-methods/sessions-rewind.ts index 756f210f59d9..274061bf85e8 100644 --- a/src/gateway/server-methods/sessions-rewind.ts +++ b/src/gateway/server-methods/sessions-rewind.ts @@ -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, 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, diff --git a/src/gateway/session-compaction-checkpoints.test.ts b/src/gateway/session-compaction-checkpoints.test.ts index b8539de6a340..09fea0e5f48d 100644 --- a/src/gateway/session-compaction-checkpoints.test.ts +++ b/src/gateway/session-compaction-checkpoints.test.ts @@ -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((resolve) => { + releaseOwnerChange = resolve; + }); + let markOwnerChangeStarted = () => {}; + const ownerChangeStarted = new Promise((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((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", diff --git a/src/gateway/session-compaction-checkpoints.ts b/src/gateway/session-compaction-checkpoints.ts index b7c761ca6e31..3bae86a446ad 100644 --- a/src/gateway/session-compaction-checkpoints.ts +++ b/src/gateway/session-compaction-checkpoints.ts @@ -66,10 +66,14 @@ export function resolveCompactionCheckpointTranscriptPosition(params: { }; } -type CompactionCheckpointSessionMutationResult = SessionCompactionCheckpointMutationResult; +type CompactionCheckpointSessionMutationResult = + | SessionCompactionCheckpointMutationResult + | { status: "conflict" }; +type SessionEntryExpectedState = Pick; 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 } : {}), }); diff --git a/src/gateway/worker-environments/transcript-commit.test.ts b/src/gateway/worker-environments/transcript-commit.test.ts index 34c36aef82de..427aed7bd39c 100644 --- a/src/gateway/worker-environments/transcript-commit.test.ts +++ b/src/gateway/worker-environments/transcript-commit.test.ts @@ -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((resolve) => { + releaseOwnerChange = resolve; + }); + let markOwnerChangeStarted = () => {}; + const ownerChangeStarted = new Promise((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((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 }); diff --git a/src/gateway/worker-environments/transcript-commit.ts b/src/gateway/worker-environments/transcript-commit.ts index 33f096ec13bd..00a048c64c64 100644 --- a/src/gateway/worker-environments/transcript-commit.ts +++ b/src/gateway/worker-environments/transcript-commit.ts @@ -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; }