diff --git a/CHANGELOG.md b/CHANGELOG.md index 144c989..8171aa1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,16 @@ # Changelog +## 17.1.11 — 2026-09-28 + +Use published TCloud 0.6 for research chat and search. +Require Interface 2.13 to match the installed Sandbox dependency. +Keep the public `RouterError`, retry options, and injectable `RouterClient` contract. +Forward cancellation into HTTP and attribute cost from each response. +Pass run cancellation through default worker and verifier requests, including claim extraction. +Stop research after cancellation before fallback queries, source writes, or round events. +Set `usage().usd` to NaN if a successful response omits its cost receipt. +Reject search results without URLs and retain the caller-controlled reasoning timeout. + ## 17.1.10 — 2026-09-27 Allow consumers to install public Eval 0.199 alongside Knowledge after packed-package checks. Keep the previously supported Eval 0.182 through 0.198 releases; the 0.199.0 development pin replaces 0.198.0. diff --git a/package.json b/package.json index defc4d6..b8842cf 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@tangle-network/agent-knowledge", - "version": "17.1.10", + "version": "17.1.11", "description": "Build, search, evaluate, and improve source-backed knowledge bases.", "homepage": "https://github.com/tangle-network/agent-knowledge#readme", "repository": { @@ -81,18 +81,19 @@ "dependencies": { "@types/proper-lockfile": "4.1.4", "proper-lockfile": "4.1.2", - "zod": "4.5.4" + "zod": "4.5.4", + "@tangle-network/tcloud": ">=0.6.0 <0.7.0" }, "peerDependencies": { "@tangle-network/agent-eval": ">=0.182.0 <0.200.0", - "@tangle-network/agent-interface": "^2.0.0" + "@tangle-network/agent-interface": "^2.13.0" }, "devDependencies": { "@arethetypeswrong/cli": "^0.18.5", "@biomejs/biome": "^2.5.11", "@neo4j-labs/agent-memory": "0.4.1", "@tangle-network/agent-eval": "0.199.0", - "@tangle-network/agent-interface": "2.0.0", + "@tangle-network/agent-interface": "2.13.0", "@types/node": "^26.4.0", "mem0ai": "3.1.7", "oxc-parser": "0.147.0", @@ -119,7 +120,9 @@ "@tangle-network/agent-trace-contract", "esbuild", "vite", - "zod@4.5.4" + "zod@4.5.4", + "@tangle-network/tcloud", + "@tangle-network/sandbox" ], "overrides": { "@hono/node-server": "2.0.12", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 23ae409..e270259 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -15,6 +15,9 @@ importers: .: dependencies: + '@tangle-network/tcloud': + specifier: '>=0.6.0 <0.7.0' + version: 0.6.0(typescript@7.0.2)(zod@4.5.4) '@types/proper-lockfile': specifier: 4.1.4 version: 4.1.4 @@ -38,14 +41,14 @@ importers: specifier: 0.199.0 version: 0.199.0 '@tangle-network/agent-interface': - specifier: 2.0.0 - version: 2.0.0 + specifier: 2.13.0 + version: 2.13.0 '@types/node': specifier: ^26.4.0 version: 26.4.0 mem0ai: specifier: 3.1.7 - version: 3.1.7(@cloudflare/workers-types@4.20260702.1)(@types/jest@29.5.14)(@types/pg@8.11.0)(better-sqlite3@12.11.1)(compromise@14.16.0)(mongodb@7.5.0)(natural@8.1.1(@opentelemetry/api@1.9.1))(pg@8.11.3) + version: 3.1.7(@cloudflare/workers-types@4.20260702.1)(@types/jest@29.5.14)(@types/pg@8.11.0)(better-sqlite3@12.11.1)(compromise@14.16.0)(mongodb@7.5.0)(natural@8.1.1(@opentelemetry/api@1.9.1))(pg@8.11.3)(ws@8.21.1) oxc-parser: specifier: 0.147.0 version: 0.147.0 @@ -70,6 +73,9 @@ importers: packages: + '@adraffy/ens-normalize@1.11.1': + resolution: {integrity: sha512-nhCBV3quEgesuf7c7KYfperqSS14T8bYuvJ8PcLJp6znkZpFc0AuW4qBtr8eKVyPPe/8RSr7sglCWPU5eaxwKQ==} + '@andrewbranch/untar.js@1.0.3': resolution: {integrity: sha512-Jh15/qVmrLGhkKJBdXlK1+9tY4lZruYjsgkDFj08ZmDiWVBLJcqkok7Z0/R0In+i1rScBpJlSvrTS2Lm41Pbnw==} @@ -364,6 +370,18 @@ packages: resolution: {integrity: sha512-93e47GTUyQ/HxzwXM1h+iZtZtaUyUX1Bcd7yBpUm0SYqC/qE1MB5QO0/jP0pnPNtBwYq+nZac23Hsht5NpCMQA==} engines: {node: '>=20.0.0'} + '@noble/ciphers@1.3.0': + resolution: {integrity: sha512-2I0gnIVPtfnMw9ee9h1dJG7tp81+8Ob3OJb3Mv37rx5L40/b0i7djjCVvGOVqc9AEIQyvyu1i6ypKdFw8R8gQw==} + engines: {node: ^14.21.3 || >=16} + + '@noble/curves@1.9.1': + resolution: {integrity: sha512-k11yZxZg+t+gWvBbIswW0yoJlu8cHOC7dhunwOzoWH/mXGBiYyR4YY6hAEK/3EUs4UpB8la1RfdRpeGsFHkWsA==} + engines: {node: ^14.21.3 || >=16} + + '@noble/hashes@1.8.0': + resolution: {integrity: sha512-jCs9ldd7NwzpgXDIf6P3+NrHh9/sD6CQdxHyjQI+h/6rDNo88ypBxxz45UDuZHz9r3tNz7N/VInSVoVdtXEI4A==} + engines: {node: ^14.21.3 || >=16} + '@noble/hashes@2.4.0': resolution: {integrity: sha512-X5XaVWZIBCT7HHZGm5I7ZQXDwLG+bGXuSrMQAW+7Zvl87h1kmc1ZB1VSRJcpUfoUrGQp4Fkoxm5kZ+Ms+aW+eA==} engines: {node: '>= 20.19.0'} @@ -737,6 +755,15 @@ packages: '@rolldown/pluginutils@1.0.1': resolution: {integrity: sha512-2j9bGt5Jh8hj+vPtgzPtl72j0yRxHAyumoo6TNfAjsLB04UtpSvPbPcDcBMxz7n+9CYB0c1GxQFxYRg2jimqGw==} + '@scure/base@1.2.6': + resolution: {integrity: sha512-g/nm5FgUa//MCj1gV09zTJTaM6KBAHqLN907YVQqf7zC49+DcO4B1so4ZX07Ef10Twr6nuqYEH9GEggFXA4Fmg==} + + '@scure/bip32@1.7.0': + resolution: {integrity: sha512-E4FFX/N3f4B80AKWp5dP6ow+flD1LQZo/w8UnLGYZO674jS6YnYeepycOOksv+vLPSpgN35wgKgy+ybfTb2SMw==} + + '@scure/bip39@1.6.0': + resolution: {integrity: sha512-+lF0BbLiJNwVlev4eKelw1WWLaiKXw7sSl8T6FvBlWkdX+94aGJ4o8XjUdlyhTCjd8c+B3KT3JfS8P0bLRNU6A==} + '@sinclair/typebox@0.27.12': resolution: {integrity: sha512-hhyNJ+nbR6ZR7pToHvllEFun9TL0sbL+tk/ON75lo+Xas054uez98qRbsuNt7MBCyZKK4+8Yli/OAGZhmfBZ/g==} @@ -760,12 +787,42 @@ packages: engines: {node: '>=20.19.0'} hasBin: true - '@tangle-network/agent-interface@2.0.0': - resolution: {integrity: sha512-t4oXA4y/4irRd1NmsxZcYwfabUzE//5yP/Lqg0Fz9WlV2bBfBJdZ5sHybaTecgr+p+PcWu2RG01/OsGmfdzudQ==} + '@tangle-network/agent-interface@2.13.0': + resolution: {integrity: sha512-gBOlHyw7iDbMBLLqmZqIZ2hVg9bRBOU5Hmyo31w5sGOWxXQYqSAwqqUlGe5klHRefOW7souT71i5B/dEKdbbFA==} '@tangle-network/agent-trace-contract@1.3.0': resolution: {integrity: sha512-kvWKftAShL9TB8F945TII4Wj2vpH/MfCNtU2aoChqWZIB8Kt9AzqJKisZimUw7fRouLKSm1uOw9FsdX09vsXEQ==} + '@tangle-network/sandbox@0.54.2': + resolution: {integrity: sha512-TGMZlsnQXA/dDoE6M8t5RMLSypFcw3h3TNZCPmHh/MZyZKMugyX6WSmH4xaWSkBXj0hewMkXjm7PJpInx/17LQ==} + peerDependencies: + '@mastra/core': ^1.36.0 + '@modelcontextprotocol/sdk': ^1.30.0 + ai: ^6.0.175 || ^7.0.107 + openai: ^6.36.0 + viem: ^2.0.0 + peerDependenciesMeta: + '@mastra/core': + optional: true + '@modelcontextprotocol/sdk': + optional: true + ai: + optional: true + openai: + optional: true + viem: + optional: true + + '@tangle-network/tcloud-attestation@0.1.1': + resolution: {integrity: sha512-+TAF9s5t1jOWGyGHvKhIWe2FYmG7puVaxmmg0Et67ylAjGa7GqUAvISXGjG/6dzld7A170V0kQHK0WVdh2Wh0Q==} + engines: {node: '>=18'} + + '@tangle-network/tcloud@0.6.0': + resolution: {integrity: sha512-NdT4Rcwb1zRENqdglhWszRMCO66N6rws+FhSL8x1lj6lsoGajNMZS8OGOvtsBuhWPObhec0WZzJqVV5EOULVZA==, tarball: https://registry.npmjs.org/@tangle-network/tcloud/-/tcloud-0.6.0.tgz} + version: 0.6.0 + engines: {node: '>=20.19.0'} + hasBin: true + '@tybys/wasm-util@0.10.3': resolution: {integrity: sha512-F3fo1MYrRJYL3zER0OUOmkutjr1Vp23m7OsSgp7nq4SP6OqX6C/56XFIPAl5bt3zaBRjmW7SGz3u/6LwFpYcOg==} @@ -1097,6 +1154,17 @@ packages: '@yuku-toolchain/types@0.8.0': resolution: {integrity: sha512-hL/raFM5V9UT2lVE/lIWDTWvANqSB8TvMVp+PgICehBa96KhTl/UpGm370JXaacukR+QnaHM2Yv6O/UcrzAFGg==} + abitype@1.2.3: + resolution: {integrity: sha512-Ofer5QUnuUdTFsBRwARMoWKOH1ND5ehwYhJ3OJ/BQO+StkwQjHw0XyVh4vDttzHB7QOFhPHa/o413PJ82gU/Tg==} + peerDependencies: + typescript: '>=5.0.4' + zod: ^3.22.0 || ^4.0.0 + peerDependenciesMeta: + typescript: + optional: true + zod: + optional: true + abort-controller@3.0.0: resolution: {integrity: sha512-h8lQ8tacZYnR3vNQTgibj+tODHI5/+l06Au2Pcriv/Gmet0eaj4TwWH41sO9wnHDiQsEj19q0drzdWdeAHtweg==} engines: {node: '>=6.5'} @@ -1255,6 +1323,10 @@ packages: resolution: {integrity: sha512-y4Mg2tXshplEbSGzx7amzPwKKOCGuoSRP/CjEdwwk0FOGlUbq6lKuoyDZTNZkmxHdJtp54hdfY/JUrdL7Xfdug==} engines: {node: '>=14'} + commander@14.0.3: + resolution: {integrity: sha512-H+y0Jo/T1RZ9qPP4Eh1pkcQcLRglraJaSLoyOtHxu6AapkjWVCy2Sit1QQ4x3Dng8qDlSsZEet7g5Pq06MvTgw==} + engines: {node: '>=20'} + compromise@14.16.0: resolution: {integrity: sha512-4DFYl/Hl7sW4XWUDfx9S5vxqyYKpZDwwqrpXsQv5acdbVP+joKceIcIaLb0lhVWUpDBV0OnExk/o/dnYUwXnhQ==} engines: {node: '>=12.0.0'} @@ -1371,6 +1443,9 @@ packages: resolution: {integrity: sha512-i/2XbnSz/uxRCU6+NdVJgKWDTM427+MqYbkQzD321DuCQJUqOuJKIA0IM2+W2xtYHdKOmZ4dR6fExsd4SXL+WQ==} engines: {node: '>=6'} + eventemitter3@5.0.1: + resolution: {integrity: sha512-GWkBvjiSZK87ELrYOSESUYeVIc9mvLLf/nXalMOS5dYrgZq9o5OVkbZAVM06CVxYsCwH9BDZFPlQTlPA1j4ahA==} + expand-template@2.0.3: resolution: {integrity: sha512-XYfuKMvj4O35f/pOXLObndIRvyQ+/+6AhODh+OKWj9S9498pHHn/IMszH+gt0fBCRWMNfk1ZSp5x3AifmnI2vg==} engines: {node: '>=6'} @@ -1523,6 +1598,11 @@ packages: resolution: {integrity: sha512-41Cifkg6e8TylSpdtTpeLVMqvSBEVzTttHvERD741+pnZ8ANv0004MRL43QKPDlK9cGvNp6NZWZUBlbGXYxxng==} engines: {node: '>=0.12.0'} + isows@1.0.7: + resolution: {integrity: sha512-I1fSfDCZL5P0v33sVqeTDSpcstAg/N+wF5HS033mogOVIp4B+oHC7oOCsA3axAbBSGTJ8QubbNmnIRN/h8U7hg==} + peerDependencies: + ws: 8.21.1 + jest-diff@29.7.0: resolution: {integrity: sha512-LMIgiIrhigmPrs03JHpxUh2yISK3vLFPkAodPeo0+BuF7wA2FoQbkEg1u8gBYBThncu7e1oEDUfIXVuTqLRUjw==} engines: {node: ^14.15.0 || ^16.10.0 || >=18.0.0} @@ -1930,6 +2010,14 @@ packages: openapi3-ts@4.6.0: resolution: {integrity: sha512-a4sfn6L2sIShhtzJqmjGrARvxAW/3F2BJDdyRVvNF9VhAsZSh5hSyI3a9TNvmzBxXmq66nY5LNT5bQcBxYAZZg==} + ox@0.14.45: + resolution: {integrity: sha512-jpsQ+p0JZh9mKCxxzxz0d5qFHZZg2/O9xW18IpEC/dNXXPVqmDJ0c//9RMz0CILkbSGPA42+cbU0fC/ghwJ03w==} + peerDependencies: + typescript: '>=5.4.0' + peerDependenciesMeta: + typescript: + optional: true + oxc-parser@0.147.0: resolution: {integrity: sha512-5xaug6t7GfV3BO5Iv+xHW1rmQkDEQ3BEu3L8g3InsvWO5i8CYGc4tCZ2X985QcwWNycFJam+aOns6Nr2XAThTA==} engines: {node: ^20.19.0 || >=22.12.0} @@ -2392,6 +2480,14 @@ packages: resolution: {integrity: sha512-Njrh4U8UODGajoZ44QS2C/BsoEM9DTI/aCqY5swsizb+/ap0FamvnCMcZAxrR5+aoC0ZqkawEfpC/N2SBc+xeA==} engines: {node: '>=18.12.0'} + viem@2.56.9: + resolution: {integrity: sha512-UC1G0Qz+RbPKqwl6yRrp3dH6ltrOmUcFJMTcOyE1Mg0nO5lxrEBTONh/8OFWPnALGuN/O7LkNcz/IsSj83Svqw==} + peerDependencies: + typescript: '>=5.0.4' + peerDependenciesMeta: + typescript: + optional: true + vite@8.2.2: resolution: {integrity: sha512-cFKLV/PRgAUlIRm5WjMjJ86jrftzpqcgH+Us+DS8mI3CDNiH30Whrz8uHL3+MOLPAgqbMBAqWdAHAphOAM+z/Q==} engines: {node: ^20.19.0 || >=22.12.0} @@ -2510,6 +2606,18 @@ packages: wrappy@1.0.2: resolution: {integrity: sha512-l4Sp/DRseor9wL6EvV2+TuQn63dMkPjZ/sp9XkghTEbV9KlPS1xUsZ3u7/IQO4wxtcFB4bgpQPRcR3QCvezPcQ==} + ws@8.21.1: + resolution: {integrity: sha512-+0NTnW77fFN/DjQi6k/Sq/Yvk4Sgajw7urW8V+asjXnRgDs9gyGkdb7EzgfhA4goXsRIZKE28fzIXBHEzhuiWw==} + engines: {node: '>=10.0.0'} + peerDependencies: + bufferutil: ^4.0.1 + utf-8-validate: '>=5.0.2' + peerDependenciesMeta: + bufferutil: + optional: true + utf-8-validate: + optional: true + xtend@4.0.2: resolution: {integrity: sha512-LKYU1iAXJXUgAXn9URjiu+MWhyUXHsvfp7mcuYm9dSUKK0/CjtrUwFAxD82/mCWbtLsGjFIad0wIsod4zrTAEQ==} engines: {node: '>=0.4'} @@ -2548,6 +2656,8 @@ packages: snapshots: + '@adraffy/ens-normalize@1.11.1': {} + '@andrewbranch/untar.js@1.0.3': {} '@arethetypeswrong/cli@0.18.5': @@ -2760,6 +2870,14 @@ snapshots: '@neo4j-labs/agent-memory@0.4.1': {} + '@noble/ciphers@1.3.0': {} + + '@noble/curves@1.9.1': + dependencies: + '@noble/hashes': 1.8.0 + + '@noble/hashes@1.8.0': {} + '@noble/hashes@2.4.0': {} '@opentelemetry/api@1.9.1': @@ -2952,6 +3070,19 @@ snapshots: '@rolldown/pluginutils@1.0.1': {} + '@scure/base@1.2.6': {} + + '@scure/bip32@1.7.0': + dependencies: + '@noble/curves': 1.9.1 + '@noble/hashes': 1.8.0 + '@scure/base': 1.2.6 + + '@scure/bip39@1.6.0': + dependencies: + '@noble/hashes': 1.8.0 + '@scure/base': 1.2.6 + '@sinclair/typebox@0.27.12': {} '@sindresorhus/is@4.6.0': {} @@ -2960,7 +3091,7 @@ snapshots: '@tangle-network/agent-core@0.9.6': dependencies: - '@tangle-network/agent-interface': 2.0.0 + '@tangle-network/agent-interface': 2.13.0 zod: 4.5.4 '@tangle-network/agent-eval@0.199.0': @@ -2968,7 +3099,7 @@ snapshots: '@asteasolutions/zod-to-openapi': 9.1.0(zod@4.5.4) '@hono/node-server': 2.0.12(hono@4.12.32) '@tangle-network/agent-core': 0.9.6 - '@tangle-network/agent-interface': 2.0.0 + '@tangle-network/agent-interface': 2.13.0 '@tangle-network/agent-trace-contract': 1.3.0 hono: 4.12.32 linear-sum-assignment: 1.0.9 @@ -2978,7 +3109,7 @@ snapshots: transitivePeerDependencies: - '@modelcontextprotocol/sdk' - '@tangle-network/agent-interface@2.0.0': + '@tangle-network/agent-interface@2.13.0': dependencies: '@noble/hashes': 2.4.0 spdx-expression-parse: 5.0.0 @@ -2986,6 +3117,32 @@ snapshots: '@tangle-network/agent-trace-contract@1.3.0': {} + '@tangle-network/sandbox@0.54.2(viem@2.56.9(typescript@7.0.2)(zod@4.5.4))': + dependencies: + '@tangle-network/agent-core': 0.9.6 + '@tangle-network/agent-interface': 2.13.0 + zod: 4.5.4 + optionalDependencies: + viem: 2.56.9(typescript@7.0.2)(zod@4.5.4) + + '@tangle-network/tcloud-attestation@0.1.1': {} + + '@tangle-network/tcloud@0.6.0(typescript@7.0.2)(zod@4.5.4)': + dependencies: + '@tangle-network/sandbox': 0.54.2(viem@2.56.9(typescript@7.0.2)(zod@4.5.4)) + '@tangle-network/tcloud-attestation': 0.1.1 + commander: 14.0.3 + viem: 2.56.9(typescript@7.0.2)(zod@4.5.4) + transitivePeerDependencies: + - '@mastra/core' + - '@modelcontextprotocol/sdk' + - ai + - bufferutil + - openai + - typescript + - utf-8-validate + - zod + '@tybys/wasm-util@0.10.3': dependencies: tslib: 2.8.1 @@ -3223,6 +3380,11 @@ snapshots: '@yuku-toolchain/types@0.8.0': {} + abitype@1.2.3(typescript@7.0.2)(zod@4.5.4): + optionalDependencies: + typescript: 7.0.2 + zod: 4.5.4 + abort-controller@3.0.0: dependencies: event-target-shim: 5.0.1 @@ -3370,6 +3532,8 @@ snapshots: commander@10.0.1: {} + commander@14.0.3: {} + compromise@14.16.0: dependencies: efrt: 2.7.0 @@ -3477,6 +3641,8 @@ snapshots: event-target-shim@5.0.1: {} + eventemitter3@5.0.1: {} + expand-template@2.0.3: {} expect-type@1.4.0: {} @@ -3602,6 +3768,10 @@ snapshots: is-number@7.0.0: {} + isows@1.0.7(ws@8.21.1): + dependencies: + ws: 8.21.1 + jest-diff@29.7.0: dependencies: chalk: 4.1.2 @@ -3719,7 +3889,7 @@ snapshots: math-intrinsics@1.1.0: {} - mem0ai@3.1.7(@cloudflare/workers-types@4.20260702.1)(@types/jest@29.5.14)(@types/pg@8.11.0)(better-sqlite3@12.11.1)(compromise@14.16.0)(mongodb@7.5.0)(natural@8.1.1(@opentelemetry/api@1.9.1))(pg@8.11.3): + mem0ai@3.1.7(@cloudflare/workers-types@4.20260702.1)(@types/jest@29.5.14)(@types/pg@8.11.0)(better-sqlite3@12.11.1)(compromise@14.16.0)(mongodb@7.5.0)(natural@8.1.1(@opentelemetry/api@1.9.1))(pg@8.11.3)(ws@8.21.1): dependencies: '@cloudflare/workers-types': 4.20260702.1 '@types/jest': 29.5.14 @@ -3728,7 +3898,7 @@ snapshots: better-sqlite3: 12.11.1 compromise: 14.16.0 natural: 8.1.1(@opentelemetry/api@1.9.1) - openai: 4.104.0(zod@3.25.76) + openai: 4.104.0(ws@8.21.1)(zod@3.25.76) pg: 8.11.3 uuid: 11.1.1 zod: 3.25.76 @@ -3893,7 +4063,7 @@ snapshots: dependencies: wrappy: 1.0.2 - openai@4.104.0(zod@3.25.76): + openai@4.104.0(ws@8.21.1)(zod@3.25.76): dependencies: '@types/node': 18.19.130 '@types/node-fetch': 2.6.13 @@ -3903,6 +4073,7 @@ snapshots: formdata-node: 4.4.1 node-fetch: 2.7.0 optionalDependencies: + ws: 8.21.1 zod: 3.25.76 transitivePeerDependencies: - encoding @@ -3911,6 +4082,21 @@ snapshots: dependencies: yaml: 2.9.0 + ox@0.14.45(typescript@7.0.2)(zod@4.5.4): + dependencies: + '@adraffy/ens-normalize': 1.11.1 + '@noble/ciphers': 1.3.0 + '@noble/curves': 1.9.1 + '@noble/hashes': 1.8.0 + '@scure/bip32': 1.7.0 + '@scure/bip39': 1.6.0 + abitype: 1.2.3(typescript@7.0.2)(zod@4.5.4) + eventemitter3: 5.0.1 + optionalDependencies: + typescript: 7.0.2 + transitivePeerDependencies: + - zod + oxc-parser@0.147.0: dependencies: '@oxc-project/types': 0.147.0 @@ -4409,6 +4595,23 @@ snapshots: verkit@0.3.0: {} + viem@2.56.9(typescript@7.0.2)(zod@4.5.4): + dependencies: + '@noble/curves': 1.9.1 + '@noble/hashes': 1.8.0 + '@scure/bip32': 1.7.0 + '@scure/bip39': 1.6.0 + abitype: 1.2.3(typescript@7.0.2)(zod@4.5.4) + isows: 1.0.7(ws@8.21.1) + ox: 0.14.45(typescript@7.0.2)(zod@4.5.4) + ws: 8.21.1 + optionalDependencies: + typescript: 7.0.2 + transitivePeerDependencies: + - bufferutil + - utf-8-validate + - zod + vite@8.2.2(@types/node@26.4.0)(esbuild@0.28.1)(tsx@4.23.1)(yaml@2.9.0): dependencies: lightningcss: 1.33.0 @@ -4482,6 +4685,8 @@ snapshots: wrappy@1.0.2: {} + ws@8.21.1: {} + xtend@4.0.2: {} y18n@5.0.8: {} diff --git a/scripts/prove-research-live.mjs b/scripts/prove-research-live.mjs new file mode 100644 index 0000000..dfcb3ae --- /dev/null +++ b/scripts/prove-research-live.mjs @@ -0,0 +1,29 @@ +#!/usr/bin/env node +import assert from 'node:assert/strict' +import { randomUUID } from 'node:crypto' +import { createTangleRouterClient, RouterError } from '../dist/index.js' + +assert.ok(process.env.TANGLE_API_KEY, 'Set a funded TANGLE_API_KEY without printing it') +const marker = `knowledge-tcloud-ok-${randomUUID()}` +const model = process.env.TANGLE_PROOF_MODEL || 'gpt-4o-mini' +const client = createTangleRouterClient({ + apiKey: process.env.TANGLE_API_KEY, + baseUrl: process.env.TANGLE_ROUTER_URL || 'https://router.tangle.tools/v1', + model, + maxRetries: 0, + ...(process.env.TANGLE_SEARCH_PROVIDER ? { searchProvider: process.env.TANGLE_SEARCH_PROVIDER } : {}), + signal: AbortSignal.timeout(120000), +}) +const query = 'Tangle Sandbox SDK official documentation' +const hits = await client.search(query, { maxResults: 1 }) +assert.ok(hits.length > 0 && hits.every(hit => hit.url), 'No live search results') +const messages = [{ role: 'user', content: `Reply with exactly ${marker}` }] +const answer = await client.chat(messages, 1200) +assert.ok(answer.includes(marker), 'Chat omitted the unique request marker') +const usage = client.usage() +assert.equal(usage.chatCalls, 1) +assert.equal(usage.searchCalls, 1) +assert.ok(Number.isFinite(usage.usd) && usage.usd >= 0, 'Router omitted a cost receipt') +console.log(JSON.stringify({ proof: 'built-knowledge-router', model, marker, query, hits, messages, answer, usage }, null, 2)) +// Existing public compatibility facade remains constructible. +assert.equal(new RouterError(401, 'proof').status, 401) diff --git a/scripts/prove-research-transport.mjs b/scripts/prove-research-transport.mjs new file mode 100644 index 0000000..a6edf9b --- /dev/null +++ b/scripts/prove-research-transport.mjs @@ -0,0 +1,75 @@ +#!/usr/bin/env node +// Local HTTP proof of failures found in #222. No mocked fetch or unit runner. +import assert from 'node:assert/strict' +import { createServer } from 'node:http' +import { once } from 'node:events' +import { setTimeout as sleep } from 'node:timers/promises' +import { createTangleRouterClient, RouterError } from '../dist/index.js' + +const abort = new AbortController() +const timers = new Set() +let cancelledSocket = false +const observed = [] +const server = createServer(async (req, res) => { + let text = '' + for await (const chunk of req) text += chunk + const body = JSON.parse(text) + observed.push({ path: req.url, client: req.headers['x-tangle-client'], hasSignalField: 'signal' in body }) + if (body.query) { + res.setHeader('Content-Type', 'application/json') + res.end(JSON.stringify({ data: [{ title: 'source', url: 'https://example.com/proof' }, {}], usage: { billed_cost: 0.003 } })) + return + } + const prompt = body.messages[0].content + if (prompt === 'denied') { + res.writeHead(401, { 'Content-Type': 'application/json' }) + res.end('{"error":{"message":"proof denial"}}') + return + } + if (prompt === 'cancel') { + res.once('close', () => { cancelledSocket = true }) + res.writeHead(200, { 'Content-Type': 'application/json' }) + res.write('{"choices":[') + timers.add(setTimeout(() => abort.abort(new DOMException('proof cancelled', 'AbortError')), 30)) + timers.add(setTimeout(() => res.end(']}'), 2500)) + return + } + if (prompt === 'unbilled') { + res.writeHead(200, { + 'Content-Type': 'application/json', + 'X-Tangle-Price-Input': '0.000001', + 'X-Tangle-Price-Output': '0.000002', + }) + res.end(JSON.stringify({ choices: [{ message: { role: 'assistant', content: prompt } }], usage: { prompt_tokens: 2, completion_tokens: 1, total_tokens: 3 } })) + return + } + await sleep(prompt === 'first' ? 50 : 5) + res.writeHead(200, { 'Content-Type': 'application/json', 'X-Tangle-Cost-USD': prompt === 'first' ? '0.01' : '0.02' }) + res.end(JSON.stringify({ choices: [{ message: { role: 'assistant', content: prompt } }], usage: { prompt_tokens: 2, completion_tokens: 1, total_tokens: 3 } })) +}) +server.listen(0, '127.0.0.1') +await once(server, 'listening') +const options = { baseUrl: `http://127.0.0.1:${server.address().port}/v1`, apiKey: 'local-proof-only', maxRetries: 0, retryBaseMs: 1 } +const client = createTangleRouterClient(options) +try { + const answers = await Promise.all(['first', 'second'].map(content => client.chat([{ role: 'user', content }]))) + assert.deepEqual(answers, ['first', 'second']) + const hits = await client.search('proof', { maxResults: 2 }) + assert.equal(hits.length, 1) + assert.ok(Math.abs(client.usage().usd - 0.033) < 1e-10, 'concurrent requests double-counted cost') + assert.equal(client.usage().promptTokens, 4) + const unbilled = createTangleRouterClient(options) + assert.equal(await unbilled.chat([{ role: 'user', content: 'unbilled' }]), 'unbilled') + assert.equal(Number.isNaN(unbilled.usage().usd), true, 'rate estimate became a billing receipt') + await assert.rejects(client.chat([{ role: 'user', content: 'denied' }]), error => error instanceof RouterError && error.status === 401) + const cancellable = createTangleRouterClient({ ...options, signal: abort.signal }) + await assert.rejects(cancellable.chat([{ role: 'user', content: 'cancel' }]), { name: 'AbortError' }) + for (let i = 0; i < 50 && !cancelledSocket; i++) await sleep(10) + assert.equal(cancelledSocket, true, 'aborted caller left the HTTP request alive') + assert.ok(observed.every(request => request.client?.startsWith('tcloud-sdk/') && !request.hasSignalField)) + console.log(JSON.stringify({ proof: 'research-adapter-real-http', observed, usage: client.usage(), missingReceiptIsNaN: Number.isNaN(unbilled.usage().usd), cancelledSocket, errorFacade: 'RouterError(401)' }, null, 2)) +} finally { + for (const timer of timers) clearTimeout(timer) + server.closeAllConnections() + await new Promise(resolve => server.close(resolve)) +} diff --git a/src/research-driving-driver.ts b/src/research-driving-driver.ts index c9fd87b..9b7831a 100644 --- a/src/research-driving-driver.ts +++ b/src/research-driving-driver.ts @@ -382,8 +382,19 @@ function buildDriver( } } - function resolveRouter(): RouterClient { - return options.router ?? createTangleRouterClient(options.router_options) + function resolveRouter(signal?: AbortSignal): RouterClient { + if (options.router) return options.router + const configuredSignal = options.router_options?.signal + const requestSignal = + configuredSignal && signal + ? AbortSignal.any([configuredSignal, signal]) + : (configuredSignal ?? signal) + return createTangleRouterClient({ ...options.router_options, signal: requestSignal }) + } + + function throwIfRouterAborted(signal?: AbortSignal): void { + signal?.throwIfAborted() + if (!options.router) options.router_options?.signal?.throwIfAborted() } /** Record a claim from a source, growing its independent-source support. */ @@ -586,12 +597,14 @@ function buildDriver( source: ResearchSourceProposal, ctx: SourceVerificationContext, ): Promise { + throwIfRouterAborted(ctx.signal) const sourceSnapshot = snapshotSourceTextInput(source) const goalSnapshot = ctx.goal const roundSnapshot = ctx.round const sourceVersion = sourceVersionOfProposal(sourceSnapshot) bindGoal(goalSnapshot) - const extracted = await extractClaims(sourceSnapshot, goalSnapshot) + const extracted = await extractClaims(sourceSnapshot, goalSnapshot, ctx.signal) + throwIfRouterAborted(ctx.signal) if (extracted.length === 0) { return { accept: false, @@ -620,6 +633,7 @@ function buildDriver( // loop confirms source registration through `commitSources`, closing both // possible crash directions without a cross-store transaction. await persist() + throwIfRouterAborted(ctx.signal) return { accept: true } }, @@ -685,9 +699,11 @@ function buildDriver( async function extractClaims( source: ResearchSourceProposal, goal: string, + signal?: AbortSignal, ): Promise { const ledger = claimsForExtraction() - const fromLlm = await extractClaimsWithLlm(source, goal, ledger) + const fromLlm = await extractClaimsWithLlm(source, goal, ledger, signal) + throwIfRouterAborted(signal) if (fromLlm.length > 0) return fromLlm.slice(0, maxClaimsPerSource) if (deterministicFallback) return deterministicClaims(source).slice(0, maxClaimsPerSource) return [] @@ -714,11 +730,14 @@ function buildDriver( source: ResearchSourceProposal, goal: string, ledger: TrackedClaim[], + signal?: AbortSignal, ): Promise { let router: RouterClient try { - router = resolveRouter() - } catch { + router = resolveRouter(signal) + } catch (error) { + throwIfRouterAborted(signal) + if ((error as { name?: string } | null)?.name === 'AbortError') throw error return [] } const excerpt = source.text.slice(0, 1800) @@ -750,9 +769,12 @@ function buildDriver( ], 1200, ) - } catch { + } catch (error) { + throwIfRouterAborted(signal) + if ((error as { name?: string } | null)?.name === 'AbortError') throw error return [] } + throwIfRouterAborted(signal) return parseExtractedClaims(raw, ledger) } diff --git a/src/verified-research-loop.ts b/src/verified-research-loop.ts index 691ff99..73ad2ad 100644 --- a/src/verified-research-loop.ts +++ b/src/verified-research-loop.ts @@ -229,20 +229,25 @@ export interface VerifiedResearchLoopResult { export async function runVerifiedResearchLoop( options: VerifiedResearchLoopOptions, ): Promise { + const assertActive = () => options.signal?.throwIfAborted() + assertActive() const maxRounds = Math.max(1, options.maxRounds ?? 3) await initKnowledgeBase(options.root) + assertActive() const store = new FileSystemKbStore({ root: options.root }) const steps: VerifiedResearchRound[] = [] let index = await buildKnowledgeIndex(options.root) + assertActive() // Reconcile source writes that completed before a previous process died while // confirming them to the driver. The records carry both original URI and hash. await confirmRegisteredSources(options.driver, index.sources) + assertActive() let readiness = readinessFor(options, index) let ready = isReady(readiness?.report) && (options.driver.isComplete?.() ?? true) let steer: string | undefined for (let round = 1; round <= maxRounds && !ready; round++) { - if (options.signal?.aborted) throw new Error('Verified research loop aborted') + assertActive() const gaps = gapsFromReadiness(readiness) @@ -257,6 +262,7 @@ export async function runVerifiedResearchLoop( readiness: requireReadiness(readiness, options), signal: options.signal, }) + assertActive() // 2. DRIVER VERIFIES the worker's sources before they commit. const accepted: ResearchSourceProposal[] = [] @@ -285,6 +291,7 @@ export async function runVerifiedResearchLoop( acceptedThisRound: accepted, signal: options.signal, }) + assertActive() if (verdict.accept) accepted.push(source) else rejectedWorkerSources.push({ source, reason: verdict.reason }) } @@ -293,14 +300,18 @@ export async function runVerifiedResearchLoop( // pages — but only when at least one source survived verification, so a // page never cites a rejected source. const acceptedWorkerSources = await registerSources(options, accepted) + assertActive() await confirmRegisteredSources(options.driver, acceptedWorkerSources) + assertActive() const writtenPages: string[] = [] writtenPages.push( ...(await applyPages(options.root, workerContribution, acceptedWorkerSources)), ) + assertActive() // Re-index so the driver's gap-fill pass sees the worker's contribution. index = await buildKnowledgeIndex(options.root) + assertActive() readiness = readinessFor(options, index) // 3. DRIVER GAP-FILLS the gaps the worker left open (opt-in). @@ -317,11 +328,16 @@ export async function runVerifiedResearchLoop( readiness: requireReadiness(readiness, options), signal: options.signal, }) + assertActive() driverNotes = driverContribution.notes driverSources = await registerSources(options, driverContribution.sources ?? []) + assertActive() await confirmRegisteredSources(options.driver, driverSources) + assertActive() writtenPages.push(...(await applyPages(options.root, driverContribution, driverSources))) + assertActive() index = await buildKnowledgeIndex(options.root) + assertActive() readiness = readinessFor(options, index) } @@ -333,6 +349,7 @@ export async function runVerifiedResearchLoop( steer = undefined } else { await options.driver.prepareFold?.() + assertActive() steer = foldGaps(options.driver, remainingGaps) } @@ -365,12 +382,16 @@ export async function runVerifiedResearchLoop( // Commit the driver's state before publishing the round event. A persisted // event therefore never claims a round whose generated questions were lost. await options.driver.checkpoint?.() + assertActive() await store.putEvent(step.event) + assertActive() steps.push(step) await options.onRound?.(step) + assertActive() } + assertActive() return { root: options.root, goal: options.goal, @@ -436,8 +457,10 @@ async function registerSources( ): Promise { const records: SourceRecord[] = [] for (const candidate of sources) { + options.signal?.throwIfAborted() const source = snapshotSourceTextInput(candidate) records.push(await addSourceText(options.root, source, options.sourceOptions)) + options.signal?.throwIfAborted() } return records } diff --git a/src/web-research-worker.ts b/src/web-research-worker.ts index 130025a..74ecd2d 100644 --- a/src/web-research-worker.ts +++ b/src/web-research-worker.ts @@ -22,12 +22,12 @@ * worker ADDS; the driver GATES. Together they build a cleaner knowledge base * than a single agent at the same compute budget. * - * Dependency-free on purpose: it talks to the router over `fetch` directly with - * the published OpenAI-compatible chat shape and the `/v1/search` shape, so it - * works whether or not the `tcloud` CLI is installed. Point it at any router by - * passing `baseUrl`; supply the key via `apiKey` or `TANGLE_API_KEY`. + * Model and search requests use the published TCloud SDK. Knowledge owns the + * research policy, not HTTP, bearer headers, retries or provider rate cards. + * Point it at any router with `baseUrl`; supply `apiKey` or `TANGLE_API_KEY`. */ +import { type SearchProvider, TCloudClient, TCloudError } from '@tangle-network/tcloud' import { htmlToText } from './sources/html' import { politeFetch } from './sources/http' import type { SourceRecord } from './types' @@ -62,7 +62,7 @@ export interface WebSearchHit { /** * The two router capabilities the worker/driver need. Injectable so tests can - * stub the network; the default talks to the live Tangle router over `fetch`. + * supply a client; the default uses the published TCloud SDK. */ export interface RouterClient { /** Live web search — returns title/url/snippet hits. */ @@ -88,6 +88,7 @@ export interface RouterUsage { searchCalls: number promptTokens: number completionTokens: number + /** Reported cost only. NaN once a successful response omits its cost receipt. */ usd: number wallMs: number } @@ -101,53 +102,14 @@ export interface TangleRouterOptions { model?: string /** Optional preferred search provider (exa | you | perplexity | …). */ searchProvider?: string - /** - * Retries on a TRANSIENT upstream status (502/503/504/429) with exponential - * backoff. Default 4. A 4xx that isn't 429, and a 401, are NOT retried — those - * are not transient. After the budget is exhausted the call still fails loud - * with the original `RouterError`, so the fail-closed contract holds; this only - * stops a single upstream-capacity blip from voiding a whole multi-topic run. - */ + /** SDK retries for 429/502/503/504. Default 4. Kept for API compatibility. */ maxRetries?: number - /** Base backoff in ms (doubled each retry, ±25% jitter). Default 1500. */ + /** Initial SDK retry backoff in ms. Default 1500. */ retryBaseMs?: number signal?: AbortSignal } -/** Transient upstream statuses worth a retry (capacity / rate-limit / gateway). */ -const transientStatuses = new Set([429, 502, 503, 504]) - -/** - * POST with bounded exponential backoff on transient upstream statuses. Returns - * the first `res.ok` response, or the LAST response (so the caller throws the - * real status). A non-transient failure returns immediately — only 502/503/504/ - * 429 are retried. Aborts propagate at once. - */ -async function fetchWithRetry( - url: string, - init: RequestInit, - opts: { maxRetries: number; retryBaseMs: number; signal?: AbortSignal }, -): Promise { - let lastRes: Response | undefined - for (let attempt = 0; attempt <= opts.maxRetries; attempt += 1) { - if (opts.signal?.aborted) throw new RouterError(0, 'aborted') - const res = await fetch(url, init) - if (res.ok || !transientStatuses.has(res.status)) return res - lastRes = res - if (attempt === opts.maxRetries) break - // Drain the body so the socket frees before we wait. - await res.text().catch(() => '') - const backoff = opts.retryBaseMs * 2 ** attempt - const jitter = backoff * (0.75 + Math.random() * 0.5) - await new Promise((resolve) => setTimeout(resolve, jitter)) - } - // Exhausted: hand back the last transient response so the caller fails loud - // with its real status. - if (lastRes) return lastRes - throw new RouterError(0, 'fetchWithRetry produced no response') -} - -/** A small error so a failed router call fails loud rather than returning junk. */ +/** Compatibility error facade. HTTP execution and classification belong to TCloud. */ export class RouterError extends Error { constructor( public readonly status: number, @@ -158,27 +120,22 @@ export class RouterError extends Error { } } -/** - * Build a dependency-free Tangle router client over `fetch`. This is the same - * wire surface the `tcloud` SDK + `tcloud mcp` use (`/v1/search` for web search, - * `/v1/chat/completions` for chat) so it needs no CLI installed. - */ +/** Adapt Knowledge's public contract to the published SDK, without a second transport. */ export function createTangleRouterClient(options: TangleRouterOptions = {}): RouterClient { - const baseUrl = (options.baseUrl ?? DEFAULT_BASE_URL).replace(/\/$/, '') const apiKey = options.apiKey ?? process.env.TANGLE_API_KEY - if (!apiKey) { - throw new RouterError(401, 'no TANGLE_API_KEY (pass apiKey or set the env var)') - } - const model = options.model ?? DEFAULT_MODEL - const maxRetries = Math.max(0, options.maxRetries ?? 4) - const retryBaseMs = Math.max(1, options.retryBaseMs ?? 1500) - const headers = { - 'Content-Type': 'application/json', - Authorization: `Bearer ${apiKey}`, - } - - // glm-5.2 pricing (USD per token) + a cumulative accumulator. Read via usage(). - const price = { prompt: 0.95 / 1_000_000, completion: 3.0 / 1_000_000 } + if (!apiKey) throw new RouterError(401, 'no TANGLE_API_KEY (pass apiKey or set the env var)') + const client = new TCloudClient({ + baseURL: (options.baseUrl ?? DEFAULT_BASE_URL).replace(/\/+$/, ''), + apiKey, + model: options.model ?? DEFAULT_MODEL, + // Preserve the unbounded reasoning timeout. The caller can cancel the actual request. + timeout: 0, + retry: { + maxRetries: Math.max(0, options.maxRetries ?? 4), + initialBackoffMs: Math.max(1, options.retryBaseMs ?? 1500), + retryableStatuses: [429, 502, 503, 504], + }, + }) const acc: RouterUsage = { chatCalls: 0, searchCalls: 0, @@ -187,64 +144,60 @@ export function createTangleRouterClient(options: TangleRouterOptions = {}): Rou usd: 0, wallMs: 0, } + // NaN is deliberately not zero: an incomplete bill must not win a cost comparison. + const recordCost = (cost: number | undefined) => { + acc.usd += typeof cost === 'number' && Number.isFinite(cost) && cost >= 0 ? cost : Number.NaN + } + const translate = (error: unknown): never => { + // Preserve the caller's exact abort reason, including non-Error reasons. + options.signal?.throwIfAborted() + if (error instanceof TCloudError) throw new RouterError(error.status, error.message) + throw error + } return { async search(query, opts) { - const t0 = Date.now() - const res = await fetchWithRetry( - `${baseUrl}/search`, - { - method: 'POST', - headers, + options.signal?.throwIfAborted() + const started = Date.now() + try { + const response = await client.search({ + query, + ...(options.searchProvider ? { provider: options.searchProvider as SearchProvider } : {}), + ...(opts?.maxResults != null ? { maxResults: opts.maxResults } : {}), signal: options.signal, - body: JSON.stringify({ - query, - ...(options.searchProvider ? { provider: options.searchProvider } : {}), - ...(opts?.maxResults != null ? { maxResults: opts.maxResults } : {}), - }), - }, - { maxRetries, retryBaseMs, signal: options.signal }, - ) - acc.searchCalls += 1 - acc.wallMs += Date.now() - t0 - if (!res.ok) { - throw new RouterError(res.status, await res.text().catch(() => res.statusText)) + }) + recordCost(response.usage?.billed_cost) + return (response.data ?? []) + .filter((hit) => typeof hit?.url === 'string' && hit.url.length > 0) + .map((hit) => ({ title: hit.title ?? hit.url, url: hit.url, snippet: hit.snippet })) + } catch (error) { + return translate(error) + } finally { + acc.searchCalls += 1 + acc.wallMs += Date.now() - started } - const body = (await res.json()) as { data?: WebSearchHit[] } - return (body.data ?? []) - .filter((hit) => typeof hit?.url === 'string' && hit.url.length > 0) - .map((hit) => ({ title: hit.title ?? hit.url, url: hit.url, snippet: hit.snippet })) }, async chat(messages, maxTokens) { - // Reasoning-model floor: never let glm-5.2 spend the whole budget on - // hidden reasoning and return empty visible content. - const max_tokens = Math.max(MIN_MAX_TOKENS, maxTokens ?? MIN_MAX_TOKENS) - const t0 = Date.now() - const res = await fetchWithRetry( - `${baseUrl}/chat/completions`, - { - method: 'POST', - headers, + options.signal?.throwIfAborted() + const started = Date.now() + try { + const response = await client.chat({ + messages, + maxTokens: Math.max(MIN_MAX_TOKENS, maxTokens ?? MIN_MAX_TOKENS), + temperature: 0.2, signal: options.signal, - body: JSON.stringify({ model, messages, max_tokens, temperature: 0.2, stream: false }), - }, - { maxRetries, retryBaseMs, signal: options.signal }, - ) - if (!res.ok) { - throw new RouterError(res.status, await res.text().catch(() => res.statusText)) - } - const body = (await res.json()) as { - choices?: { message?: { content?: string } }[] - usage?: { prompt_tokens?: number; completion_tokens?: number } + }) + acc.chatCalls += 1 + acc.promptTokens += response.usage?.prompt_tokens ?? 0 + acc.completionTokens += response.usage?.completion_tokens ?? 0 + // Per-response cost, not a shared-client usage delta that races parallel calls. + recordCost(response.tangle?.costSource === 'receipt' ? response.tangle.costUsd : undefined) + return response.choices?.[0]?.message?.content ?? '' + } catch (error) { + return translate(error) + } finally { + acc.wallMs += Date.now() - started } - const promptTokens = body.usage?.prompt_tokens ?? 0 - const completionTokens = body.usage?.completion_tokens ?? 0 - acc.chatCalls += 1 - acc.promptTokens += promptTokens - acc.completionTokens += completionTokens - acc.usd += promptTokens * price.prompt + completionTokens * price.completion - acc.wallMs += Date.now() - t0 - return body.choices?.[0]?.message?.content ?? '' }, usage() { return { ...acc } @@ -270,12 +223,34 @@ export interface WebResearchWorkerOptions { maxTextChars?: number } -/** Resolve the router client lazily so a worker with an injected client never reads env. */ -function resolveRouter(opts: { - router?: RouterClient - router_options?: TangleRouterOptions -}): RouterClient { - return opts.router ?? createTangleRouterClient(opts.router_options) +/** Bind run cancellation to the default router without changing an injected client. */ +function resolveRouter( + opts: { + router?: RouterClient + router_options?: TangleRouterOptions + }, + signal?: AbortSignal, +): RouterClient { + if (opts.router) return opts.router + const requestSignal = combineRouterSignals(signal, opts.router_options?.signal) + return createTangleRouterClient({ ...opts.router_options, signal: requestSignal }) +} + +function combineRouterSignals( + runSignal: AbortSignal | undefined, + configuredSignal: AbortSignal | undefined, +): AbortSignal | undefined { + return runSignal && configuredSignal + ? AbortSignal.any([runSignal, configuredSignal]) + : (runSignal ?? configuredSignal) +} + +function throwIfRouterAborted( + runSignal: AbortSignal | undefined, + opts: { router?: RouterClient; router_options?: TangleRouterOptions }, +): void { + runSignal?.throwIfAborted() + if (!opts.router) opts.router_options?.signal?.throwIfAborted() } /** @@ -291,7 +266,12 @@ export function createWebResearchWorker(options: WebResearchWorkerOptions = {}): const maxTextChars = Math.max(minTextChars, options.maxTextChars ?? 4000) return async (ctx: WorkerResearchContext): Promise => { - const router = resolveRouter(options) + const assertActive = () => throwIfRouterAborted(ctx.signal, options) + assertActive() + const router = resolveRouter(options, ctx.signal) + const fetchSignal = options.router + ? ctx.signal + : combineRouterSignals(ctx.signal, options.router_options?.signal) // Target the BLOCKING gaps first; fall back to all gaps if none are blocking. const targetGaps = ctx.gaps.filter((gap) => gap.blocking) const gaps = targetGaps.length > 0 ? targetGaps : ctx.gaps @@ -299,28 +279,33 @@ export function createWebResearchWorker(options: WebResearchWorkerOptions = {}): return { sources: [], notes: 'no open gaps to research' } } - const queries = await formSearchQueries(router, ctx, gaps, queriesPerGap) + const queries = await formSearchQueries(router, ctx, gaps, queriesPerGap, assertActive) + assertActive() const proposals: ResearchSourceProposal[] = [] const seenUris = new Set() for (const query of queries) { if (proposals.length >= maxSourcesPerRound) break - if (ctx.signal?.aborted) break + assertActive() let hits: WebSearchHit[] try { hits = await router.search(query, { maxResults: resultsPerQuery }) } catch (error) { + assertActive() // A single failed query must not sink the round — record nothing, move on. - if ((error as { name?: string }).name === 'AbortError') break + if ((error as { name?: string } | null)?.name === 'AbortError') throw error continue } + assertActive() for (const hit of hits) { + assertActive() if (proposals.length >= maxSourcesPerRound) break if (seenUris.has(hit.url)) continue const fetched = await politeFetch(hit.url, { - signal: ctx.signal, + signal: fetchSignal, cacheDir: options.cacheDir, }) + assertActive() if (!fetched.verifiable) continue const text = htmlToText(fetched.body).slice(0, maxTextChars) if (text.length < minTextChars) continue @@ -346,6 +331,7 @@ export function createWebResearchWorker(options: WebResearchWorkerOptions = {}): } } + assertActive() return { sources: proposals, buildPages: buildCitingPages(proposals), @@ -364,6 +350,7 @@ async function formSearchQueries( ctx: WorkerResearchContext, gaps: KnowledgeGap[], queriesPerGap: number, + assertActive: () => void, ): Promise { const gapLines = gaps .map((gap, i) => `${i + 1}. ${gap.description} (readiness query: "${gap.query}")`) @@ -391,7 +378,10 @@ async function formSearchQueries( ], MIN_MAX_TOKENS, ) - } catch { + assertActive() + } catch (error) { + assertActive() + if ((error as { name?: string } | null)?.name === 'AbortError') throw error raw = '' } const parsed = parseQueryList(raw) @@ -521,7 +511,9 @@ export function createVerifyingResearchDriver( source: ResearchSourceProposal, ctx: SourceVerificationContext, ): Promise { - const router = resolveRouter(options) + const assertActive = () => throwIfRouterAborted(ctx.signal, options) + assertActive() + const router = resolveRouter(options, ctx.signal) const gapLines = ctx.gaps .map((gap) => `- ${gap.description} (query: "${gap.query}")`) .join('\n') @@ -554,8 +546,10 @@ export function createVerifyingResearchDriver( ], MIN_MAX_TOKENS, ) + assertActive() } catch (error) { - if ((error as { name?: string }).name === 'AbortError') throw error + assertActive() + if ((error as { name?: string } | null)?.name === 'AbortError') throw error // Router failure: fail-closed (reject) so an unverified source can't slip in. return acceptOnParseFailure ? { accept: true }