Merge pull request #2349 from pikasTech/feat/2345-workbench-pinia-colada

[架构] Workbench 迁移 Pinia Colada server-state/query 层
This commit is contained in:
Lyon
2026-07-02 19:31:14 +08:00
committed by GitHub
12 changed files with 467 additions and 544 deletions
+33 -2
View File
@@ -5,11 +5,12 @@
"": {
"name": "hwlab-cloud-web",
"dependencies": {
"@pinia/colada": "0.15.3",
"@vueuse/core": "^10.7.0",
"axios": "^1.16.0",
"dompurify": "^3.3.1",
"marked": "^17.0.1",
"pinia": "^2.1.7",
"pinia": "^2.2.6",
"vue": "^3.4.0",
"vue-i18n": "^9.14.5",
"vue-router": "^4.2.5",
@@ -125,6 +126,8 @@
"@one-ini/wasm": ["@one-ini/wasm@0.1.1", "", {}, "sha512-XuySG1E38YScSJoMlqovLru4KTUNSjgVTIjyh7qMX6aNN5HY5Ct5LhRJdxO79JtTzKfzV/bnWpz+zquYrISsvw=="],
"@pinia/colada": ["@pinia/colada@0.15.3", "", { "dependencies": { "@vue/devtools-api": "^7.7.2" }, "peerDependencies": { "pinia": "^2.2.6 || ^3.0.0" } }, "sha512-nQrCW8zJ8zyIMPdqvCI1ye/z8C6YDhzlsoX7KfmMdXiHf23QfIwJ1KGNOBGDntG+nlTQ/j5hzE/Mq3zqSuMb4g=="],
"@pkgjs/parseargs": ["@pkgjs/parseargs@0.11.0", "", {}, "sha512-+1VkjdD0QBLPodGrJUeqarH8VAIvQODIbwh9XpP5Syisf7YoQgsJKPNFoqqLQlu+VQ/tVSshMR6loPMn8U+dPg=="],
"@playwright/test": ["@playwright/test@1.59.1", "", { "dependencies": { "playwright": "1.59.1" }, "bin": { "playwright": "cli.js" } }, "sha512-PG6q63nQg5c9rIi4/Z5lR5IVF7yU5MqmKaPOe0HSc0O2cX1fPi96sUQu5j7eo4gKCkB2AnNGoWt7y4/Xx3Kcqg=="],
@@ -223,7 +226,11 @@
"@vue/compiler-vue2": ["@vue/compiler-vue2@2.7.16", "", { "dependencies": { "de-indent": "^1.0.2", "he": "^1.2.0" } }, "sha512-qYC3Psj9S/mfu9uVi5WvNZIzq+xnXMhOwbTFKKDD7b1lhpnn71jXSFdTQ+WsIEk0ONCd7VV2IMm7ONl6tbQ86A=="],
"@vue/devtools-api": ["@vue/devtools-api@6.6.4", "", {}, "sha512-sGhTPMuXqZ1rVOk32RylztWkfXTRhuS7vgAKv0zjqk8gbsHkJ7xfFf+jbySxt7tWObEJwyKaHMikV/WGDiQm8g=="],
"@vue/devtools-api": ["@vue/devtools-api@7.7.10", "", { "dependencies": { "@vue/devtools-kit": "^7.7.10" } }, "sha512-KxtEpUOOpFz/qOGRrAwA36QF7DqIA+FXgCYit9mk9wjbaZt0sXOFz81ElOZtKA4HbWHUdwNjZHBFsFFyp5BZiA=="],
"@vue/devtools-kit": ["@vue/devtools-kit@7.7.10", "", { "dependencies": { "@vue/devtools-shared": "^7.7.10", "birpc": "^2.3.0", "hookable": "^5.5.3", "mitt": "^3.0.1", "perfect-debounce": "^1.0.0", "speakingurl": "^14.0.1", "superjson": "^2.2.2" } }, "sha512-3WNi2Kq4tbpVbmhml7RiphmAt0279oh3fKNeWMQIrltfX8Q91b4i5PL8DtyNKdwmcsGrV4fg+erwWOmD05CLIw=="],
"@vue/devtools-shared": ["@vue/devtools-shared@7.7.10", "", { "dependencies": { "rfdc": "^1.4.1" } }, "sha512-wOPslzB8vTvpxwdaOcR2qAbwmuSP0L+rhpoC6Cf56V3Jip+HWb7PQQXOUPgBNQARpXsbQX/+mvi8kKucmBGRwQ=="],
"@vue/language-core": ["@vue/language-core@2.2.12", "", { "dependencies": { "@volar/language-core": "2.4.15", "@vue/compiler-dom": "^3.5.0", "@vue/compiler-vue2": "^2.7.16", "@vue/shared": "^3.5.0", "alien-signals": "^1.0.3", "minimatch": "^9.0.3", "muggle-string": "^0.4.1", "path-browserify": "^1.0.1" }, "peerDependencies": { "typescript": "*" }, "optionalPeers": ["typescript"] }, "sha512-IsGljWbKGU1MZpBPN+BvPAdr55YPkj2nB/TBNGNC32Vy2qLG25DYu/NBN2vNtZqdRbTRjaoYrahLrToim2NanA=="],
@@ -275,6 +282,8 @@
"binary-extensions": ["binary-extensions@2.3.0", "", {}, "sha512-Ceh+7ox5qe7LJuLHoY0feh3pHuUDHAcRUeyL2VYghZwfpkNIy/+8Ocg0a3UuSoYzavmylwuLWQOf3hl0jjMMIw=="],
"birpc": ["birpc@2.9.0", "", {}, "sha512-KrayHS5pBi69Xi9JmvoqrIgYGDkD6mcSe/i6YKi3w5kekCLzrX4+nawcXqrj2tIp50Kw/mT/s3p+GVK0A0sKxw=="],
"brace-expansion": ["brace-expansion@2.1.1", "", { "dependencies": { "balanced-match": "^1.0.0" } }, "sha512-WR1cURNjuvBLMZBMbqM0UoE+WAfdUcEV1ccD8PVBVOI+Z3ND4+SZbN8RsfT2bMuG1qwz5RFvPukSZm5fF2D5eA=="],
"braces": ["braces@3.0.3", "", { "dependencies": { "fill-range": "^7.1.1" } }, "sha512-yQbXgO/OSZVD2IsiLlro+7Hf6Q18EJrKSEsdoMzKePKXct3gvD8oLcOQdIzGupr5Fj+EDe8gO/lxc1BzfMpxvA=="],
@@ -307,6 +316,8 @@
"config-chain": ["config-chain@1.1.13", "", { "dependencies": { "ini": "^1.3.4", "proto-list": "~1.2.1" } }, "sha512-qj+f8APARXHrM0hraqXYb2/bOVSV4PvJQlNZ/DVj0QrmNM2q2euizkeuVckQ57J+W0mRH6Hvi+k50M4Jul2VRQ=="],
"copy-anything": ["copy-anything@4.0.5", "", { "dependencies": { "is-what": "^5.2.0" } }, "sha512-7Vv6asjS4gMOuILabD3l739tsaxFQmC+a7pLZm02zyvs8p977bL3zEgq3yDk5rn9B0PbYgIv++jmHcuUab4RhA=="],
"cross-spawn": ["cross-spawn@7.0.6", "", { "dependencies": { "path-key": "^3.1.0", "shebang-command": "^2.0.0", "which": "^2.0.1" } }, "sha512-uV2QOWP2nWzsy2aMp8aRibhi9dlzF5Hgh5SHaB9OiTGEyDTiJJyx0uy51QXdyWbtAHNua4XJzUKca3OzKUd3vA=="],
"cssesc": ["cssesc@3.0.0", "", { "bin": { "cssesc": "bin/cssesc" } }, "sha512-/Tb/JcjK111nNScGob5MNtsntNM1aCNUDipB/TkwZFhyDrrE47SOx/18wF2bbjgc3ZzCSKW1T5nt5EbFoAz/Vg=="],
@@ -401,6 +412,8 @@
"he": ["he@1.2.0", "", { "bin": { "he": "bin/he" } }, "sha512-F/1DnUGPopORZi0ni+CvrCgHQ5FyEAHRLSApuYWMmrbSwoN2Mn/7k+Gl38gJnR7yyDZk6WLXwiGod1JOWNDKGw=="],
"hookable": ["hookable@5.5.3", "", {}, "sha512-Yc+BQe8SvoXH1643Qez1zqLRmbA5rCL+sSmk6TVos0LWVfNIB7PGncdlId77WzLGSIB5KaWgTaNTs2lNVEI6VQ=="],
"html-encoding-sniffer": ["html-encoding-sniffer@4.0.0", "", { "dependencies": { "whatwg-encoding": "^3.1.1" } }, "sha512-Y22oTqIU4uuPgEemfz7NDJz6OeKf12Lsu+QC+s3BVpda64lTiMYCyGwg5ki4vFxkMwQdeZDl2adZoqUgdFuTgQ=="],
"http-proxy-agent": ["http-proxy-agent@7.0.2", "", { "dependencies": { "agent-base": "^7.1.0", "debug": "^4.3.4" } }, "sha512-T1gkAiYYDWYx3V5Bmyu7HcfcvL7mUrTWiM6yOfa3PIphViJ/gFPbvidQ+veqSOHci/PxBcDabeUNCzpOODJZig=="],
@@ -425,6 +438,8 @@
"is-potential-custom-element-name": ["is-potential-custom-element-name@1.0.1", "", {}, "sha512-bCYeRA2rVibKZd+s2625gGnGF/t7DSqDs4dP7CrLA1m7jKWz6pps0LpYLJN8Q64HtmPKJ1hrN3nzPNKFEKOUiQ=="],
"is-what": ["is-what@5.5.0", "", {}, "sha512-oG7cgbmg5kLYae2N5IVd3jm2s+vldjxJzK1pcu9LfpGuQ93MQSzo0okvRna+7y5ifrD+20FE8FvjusyGaz14fw=="],
"isexe": ["isexe@2.0.0", "", {}, "sha512-RHxMLp9lnKHGHRng9QFhRCMbYAcVpn69smSGcq3f36xjgVVWThj4qqLbTLlq7Ssj8B+fIQ1EuCEGI2lKsyQeIw=="],
"jackspeak": ["jackspeak@3.4.3", "", { "dependencies": { "@isaacs/cliui": "^8.0.2" }, "optionalDependencies": { "@pkgjs/parseargs": "^0.11.0" } }, "sha512-OGlZQpz2yfahA/Rd1Y8Cd9SIEsqvXkLVoSw/cgwhnhFMDbsQFeZYoJJ7bIZBS9BcamUW96asq/npPWugM+RQBw=="],
@@ -463,6 +478,8 @@
"minipass": ["minipass@7.1.3", "", {}, "sha512-tEBHqDnIoM/1rXME1zgka9g6Q2lcoCkxHLuc7ODJ5BxbP5d4c2Z5cGgtXAku59200Cx7diuHTOYfSBD8n6mm8A=="],
"mitt": ["mitt@3.0.1", "", {}, "sha512-vKivATfr97l2/QBCYAkXYDbrIWPM2IIKEl7YPhjCvKlG3kE2gm+uBo6nEXK3M5/Ffh/FLpKExzOQ3JJoJGFKBw=="],
"ms": ["ms@2.1.3", "", {}, "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA=="],
"muggle-string": ["muggle-string@0.4.1", "", {}, "sha512-VNTrAak/KhO2i8dqqnqnAHOa3cYBwXEZe9h+D5h/1ZqFSTEFHdM65lR7RoIqq3tBBYavsOXV84NoHXZ0AkPyqQ=="],
@@ -499,6 +516,8 @@
"pathval": ["pathval@2.0.1", "", {}, "sha512-//nshmD55c46FuFw26xV/xFAaB5HF9Xdap7HJBBnrKdAd6/GxDBaNA1870O79+9ueg61cZLSVc+OaFlfmObYVQ=="],
"perfect-debounce": ["perfect-debounce@1.0.0", "", {}, "sha512-xCy9V055GLEqoFaHoC1SoLIaLmWctgCUaBaWxDZ7/Zx4CTyX7cJQLJOok/orfjZAh9kEYpjJa4d0KcJmCbctZA=="],
"picocolors": ["picocolors@1.1.1", "", {}, "sha512-xceH2snhtb5M9liqDsmEw56le376mTZkEX/jEb/RxNFyegNul7eNslCXP9FDj/Lcu0X8KEyMceP2ntpaHrDEVA=="],
"picomatch": ["picomatch@2.3.2", "", {}, "sha512-V7+vQEJ06Z+c5tSye8S+nHUfI51xoXIXjHQ99cQtKUkQqqO1kO/KCJUfZXuB47h/YBlDhah2H3hdUGXn8ie0oA=="],
@@ -549,6 +568,8 @@
"reusify": ["reusify@1.1.0", "", {}, "sha512-g6QUff04oZpHs0eG5p83rFLhHeV00ug/Yf9nZM6fLeUrPguBTkTQOdpAWWspMh55TZfVQDPaN3NQJfbVRAxdIw=="],
"rfdc": ["rfdc@1.4.1", "", {}, "sha512-q1b3N5QkRUWUl7iyylaaj3kOpIT0N2i9MqIEQXP73GVsN9cw3fdx8X63cEmWhJGi2PPCF23Ijp7ktmd39rawIA=="],
"rollup": ["rollup@4.61.0", "", { "dependencies": { "@types/estree": "1.0.9" }, "optionalDependencies": { "@rollup/rollup-android-arm-eabi": "4.61.0", "@rollup/rollup-android-arm64": "4.61.0", "@rollup/rollup-darwin-arm64": "4.61.0", "@rollup/rollup-darwin-x64": "4.61.0", "@rollup/rollup-freebsd-arm64": "4.61.0", "@rollup/rollup-freebsd-x64": "4.61.0", "@rollup/rollup-linux-arm-gnueabihf": "4.61.0", "@rollup/rollup-linux-arm-musleabihf": "4.61.0", "@rollup/rollup-linux-arm64-gnu": "4.61.0", "@rollup/rollup-linux-arm64-musl": "4.61.0", "@rollup/rollup-linux-loong64-gnu": "4.61.0", "@rollup/rollup-linux-loong64-musl": "4.61.0", "@rollup/rollup-linux-ppc64-gnu": "4.61.0", "@rollup/rollup-linux-ppc64-musl": "4.61.0", "@rollup/rollup-linux-riscv64-gnu": "4.61.0", "@rollup/rollup-linux-riscv64-musl": "4.61.0", "@rollup/rollup-linux-s390x-gnu": "4.61.0", "@rollup/rollup-linux-x64-gnu": "4.61.0", "@rollup/rollup-linux-x64-musl": "4.61.0", "@rollup/rollup-openbsd-x64": "4.61.0", "@rollup/rollup-openharmony-arm64": "4.61.0", "@rollup/rollup-win32-arm64-msvc": "4.61.0", "@rollup/rollup-win32-ia32-msvc": "4.61.0", "@rollup/rollup-win32-x64-gnu": "4.61.0", "@rollup/rollup-win32-x64-msvc": "4.61.0", "fsevents": "~2.3.2" }, "bin": { "rollup": "dist/bin/rollup" } }, "sha512-T9mWdbWfQtp0B5lv/HX+wrhYsmXRlcWnXXmJbXqKJhlRaoS6KMhq0gpyzW4UJfclcxrEdLnTgjT2NjruLONu0g=="],
"rrweb-cssom": ["rrweb-cssom@0.7.1", "", {}, "sha512-TrEMa7JGdVm0UThDJSx7ddw5nVm3UJS9o9CCIZ72B1vSyEZoziDqBYP3XIoi/12lKrJR8rE3jeFHMok2F/Mnsg=="],
@@ -571,6 +592,8 @@
"source-map-js": ["source-map-js@1.2.1", "", {}, "sha512-UXWMKhLOwVKb728IUtQPXxfYU+usdybtUrK/8uGE8CQMvrhOpwvzDBwj0QhSL7MQc7vIsISBG8VQ8+IDQxpfQA=="],
"speakingurl": ["speakingurl@14.0.1", "", {}, "sha512-1POYv7uv2gXoyGFpBCmpDVSNV74IfsWlDW216UPjbWufNf+bSU6GdbDsxdcxtfwb4xlI3yxzOTKClUosxARYrQ=="],
"stackback": ["stackback@0.0.2", "", {}, "sha512-1XMJE5fQo1jGH6Y/7ebnwPOBEkIEnT4QF32d5R1+VXdXveM0IBMJt8zfaxX1P3QhVwrYe+576+jkANtSS2mBbw=="],
"std-env": ["std-env@3.10.0", "", {}, "sha512-5GS12FdOZNliM5mAOxFRg7Ir0pWz8MdpYm6AY6VPkGpbA7ZzmbzNcBJQ0GPvvyWgcY7QAhCgf9Uy89I03faLkg=="],
@@ -585,6 +608,8 @@
"sucrase": ["sucrase@3.35.1", "", { "dependencies": { "@jridgewell/gen-mapping": "^0.3.2", "commander": "^4.0.0", "lines-and-columns": "^1.1.6", "mz": "^2.7.0", "pirates": "^4.0.1", "tinyglobby": "^0.2.11", "ts-interface-checker": "^0.1.9" }, "bin": { "sucrase": "bin/sucrase", "sucrase-node": "bin/sucrase-node" } }, "sha512-DhuTmvZWux4H1UOnWMB3sk0sbaCVOoQZjv8u1rDoTV0HTdGem9hkAZtl4JZy8P2z4Bg0nT+YMeOFyVr4zcG5Tw=="],
"superjson": ["superjson@2.2.6", "", { "dependencies": { "copy-anything": "^4" } }, "sha512-H+ue8Zo4vJmV2nRjpx86P35lzwDT3nItnIsocgumgr0hHMQ+ZGq5vrERg9kJBo5AWGmxZDhzDo+WVIJqkB0cGA=="],
"supports-preserve-symlinks-flag": ["supports-preserve-symlinks-flag@1.0.0", "", {}, "sha512-ot0WnXS9fgdkgIcePe6RHNk1WA8+muPa6cSjeR3V8K27q9BB1rTE3R1p7Hv0z1ZyAc8s6Vvv8DIyWf681MAt0w=="],
"symbol-tree": ["symbol-tree@3.2.4", "", {}, "sha512-9QNk5KwDF+Bvz+PyObkmSYjI5ksVUYtjW7AU22r2NKcfLJcXp96hkDWU3+XndOsUb+AQ9QhfzfCT2O+CNWT5Tw=="],
@@ -689,6 +714,8 @@
"fast-glob/glob-parent": ["glob-parent@5.1.2", "", { "dependencies": { "is-glob": "^4.0.1" } }, "sha512-AOIgSQCepiJYwP3ARnGx+5VnTu2HBYdzbGP45eLw1vr3zB3vZLeyed1sC9hnbcOc9/SrMyM5RPQrkGz4aS9Zow=="],
"pinia/@vue/devtools-api": ["@vue/devtools-api@6.6.4", "", {}, "sha512-sGhTPMuXqZ1rVOk32RylztWkfXTRhuS7vgAKv0zjqk8gbsHkJ7xfFf+jbySxt7tWObEJwyKaHMikV/WGDiQm8g=="],
"playwright/fsevents": ["fsevents@2.3.2", "", { "os": "darwin" }, "sha512-xiqMQR4xAeHTuB9uWm+fFRcIOgKBMiOBP+eXiyT7jsgVCq1bkVygt00oASowB7EdtpOHaaPgKt812P9ab+DDKA=="],
"string-width-cjs/emoji-regex": ["emoji-regex@8.0.0", "", {}, "sha512-MSjYzcWNOA0ewAHpz0MxpYFvwg6yjy1NG3xteoqz644VCo/RPgnr1/GGt+ic3iJTzQ8Eu3TdM14SawnVUmGE6A=="],
@@ -699,6 +726,10 @@
"tinyglobby/picomatch": ["picomatch@4.0.4", "", {}, "sha512-QP88BAKvMam/3NxH6vj2o21R6MjxZUAd6nlwAS/pnGvN9IVLocLHxGYIzFhg6fUQ+5th6P4dv4eW9jX3DSIj7A=="],
"vue-i18n/@vue/devtools-api": ["@vue/devtools-api@6.6.4", "", {}, "sha512-sGhTPMuXqZ1rVOk32RylztWkfXTRhuS7vgAKv0zjqk8gbsHkJ7xfFf+jbySxt7tWObEJwyKaHMikV/WGDiQm8g=="],
"vue-router/@vue/devtools-api": ["@vue/devtools-api@6.6.4", "", {}, "sha512-sGhTPMuXqZ1rVOk32RylztWkfXTRhuS7vgAKv0zjqk8gbsHkJ7xfFf+jbySxt7tWObEJwyKaHMikV/WGDiQm8g=="],
"wrap-ansi-cjs/ansi-styles": ["ansi-styles@4.3.0", "", { "dependencies": { "color-convert": "^2.0.1" } }, "sha512-zbB9rCJAT1rbjiVDb2hqKFHNYLxgtk8NURxZ3IZwD3F6NtxbXZQCnnSi1Lkx+IDohdPlFp222wVALIheZJQSEg=="],
"wrap-ansi-cjs/string-width": ["string-width@4.2.3", "", { "dependencies": { "emoji-regex": "^8.0.0", "is-fullwidth-code-point": "^3.0.0", "strip-ansi": "^6.0.1" } }, "sha512-wKyQRQpjJ0sIp62ErSZdGsjMJWsap5oRNihHhu6G7JVO/9jIB6UyevL+tXuOqrng8j/cxKTWyWUwvSTriiZz/g=="],
+2 -1
View File
@@ -16,11 +16,12 @@
"e2e:workbench:report": "bun run deps --quiet && bunx playwright show-report .state/workbench-e2e/html-report"
},
"dependencies": {
"@pinia/colada": "0.15.3",
"@vueuse/core": "^10.7.0",
"axios": "^1.16.0",
"dompurify": "^3.3.1",
"marked": "^17.0.1",
"pinia": "^2.1.7",
"pinia": "^2.2.6",
"vue": "^3.4.0",
"vue-i18n": "^9.14.5",
"vue-router": "^4.2.5"
+22 -12
View File
@@ -32,13 +32,16 @@ const requiredFiles = Object.freeze([
"src/utils/workbench-health.ts",
"src/utils/workbench-error-runtime.ts",
"src/utils/workbench-realtime-runtime.ts",
"src/utils/workbench-refresh-runtime.ts",
"src/utils/workbench-storage-runtime.ts",
"src/utils/workbench-stream-transport.ts",
"src/stores/index.ts",
"src/stores/app.ts",
"src/stores/auth.ts",
"src/stores/workbench.ts",
"src/stores/workbench-colada-keys.ts",
"src/stores/workbench-colada-queries.ts",
"src/stores/workbench-colada-mutations.ts",
"src/stores/workbench-colada-reducer.ts",
"src/stores/workbench-event-reducer.ts",
"src/stores/workbench-message-projection-runtime.ts",
"src/stores/workbench-realtime-plan.ts",
@@ -81,9 +84,9 @@ const html = readWeb("index.html");
const pkg = JSON.parse(readWeb("package.json")) as { dependencies?: Record<string, string>; devDependencies?: Record<string, string> };
const appSource = readCloudWebAppSource(rootDir);
const workbenchStoreSource = readWeb("src/stores/workbench.ts");
const workbenchColadaSource = `${readWeb("src/stores/workbench-colada-keys.ts")}\n${readWeb("src/stores/workbench-colada-queries.ts")}\n${readWeb("src/stores/workbench-colada-mutations.ts")}\n${readWeb("src/stores/workbench-colada-reducer.ts")}`;
const workbenchRuntimePolicySource = readWeb("src/config/workbench-runtime-policy.ts");
const workbenchRealtimeRuntimeSource = `${readWeb("src/utils/workbench-realtime-runtime.ts")}\n${readWeb("src/utils/workbench-stream-transport.ts")}`;
const workbenchRefreshRuntimeSource = readWeb("src/utils/workbench-refresh-runtime.ts");
const workbenchPerformanceSource = readWeb("src/utils/workbench-performance.ts");
const workbenchEventReducerSource = readWeb("src/stores/workbench-event-reducer.ts");
const workbenchMessageProjectionRuntimeSource = readWeb("src/stores/workbench-message-projection-runtime.ts");
@@ -105,7 +108,7 @@ assert.doesNotMatch(html, /\/v1\/workbench\/workspace/u, "HTML parse must not re
assert.doesNotMatch(html, /\/auth\/workspace-bootstrap/u, "HTML parse bootstrap must not use stale auth workspace bootstrap route");
assert.doesNotMatch(html, /src="\/src\/main\.tsx"|id="root"/u, "React entry must be gone");
for (const dep of ["vue", "pinia", "vue-router", "axios", "@vueuse/core", "vue-i18n"]) {
for (const dep of ["vue", "pinia", "@pinia/colada", "vue-router", "axios", "@vueuse/core", "vue-i18n"]) {
assert.ok(pkg.dependencies?.[dep], `missing Vue/Sub2API dependency ${dep}`);
}
for (const dep of ["@vitejs/plugin-vue", "tailwindcss", "postcss", "autoprefixer", "vue-tsc", "vitest"]) {
@@ -121,6 +124,7 @@ for (const directory of ["api", "stores", "composables", "components/common", "c
assertIncludes(appSource, "createApp", "Vue app must use createApp");
assertIncludes(appSource, "createPinia", "Pinia must be installed");
assertIncludes(appSource, "PiniaColada", "Pinia Colada must be installed after Pinia");
assertIncludes(appSource, "createRouter", "Vue Router must be installed");
assertIncludes(appSource, "defineStore", "Pinia stores must exist");
assertIncludes(appSource, "installChunkRecovery", "router must recover from stale dynamic route chunks");
@@ -138,9 +142,15 @@ assert.doesNotMatch(workbenchStoreSource, /const\s+(?:DEFAULT_CODE_AGENT_TIMEOUT
assertIncludes(workbenchRealtimeRuntimeSource, "connectWorkbenchEvents", "Workbench realtime runtime must own the unified SSE EventSource entry");
assertIncludes(workbenchRealtimeRuntimeSource, "WorkbenchStreamTransportRecovery", "SSE transport must own recovery actions, not just wrap EventSource");
assertIncludes(workbenchRealtimeRuntimeSource, "cursorByKey", "SSE transport runtime must own cursor replay state");
assertIncludes(workbenchRefreshRuntimeSource, "WorkbenchTraceHydrationQueueRuntime", "Trace hydration queue must live in refresh runtime");
assertIncludes(workbenchStoreSource, "createWorkbenchTraceHydrationQueueRuntime", "Workbench store must use the migrated trace hydration runtime");
assertIncludes(workbenchStoreSource, "createWorkbenchScheduledTaskRuntime", "Workbench store must use the migrated scheduled task runtime");
assertIncludes(workbenchColadaSource, "workbenchColadaKeys", "Workbench Colada modules must own query/mutation key hierarchy");
assertIncludes(workbenchColadaSource, "useQueryCache", "Workbench Colada modules must own query cache access");
assertIncludes(workbenchColadaSource, "useQuery", "Workbench Colada reducer must keep server-state in Colada cache");
assertIncludes(workbenchColadaSource, "useMutation", "Workbench mutations must execute through Pinia Colada");
assertIncludes(workbenchStoreSource, "useWorkbenchColadaQueries", "Workbench store must read server-state through Colada query facade");
assertIncludes(workbenchStoreSource, "useWorkbenchColadaMutations", "Workbench store must run admission/cancel through Colada mutations");
assertIncludes(workbenchStoreSource, "useWorkbenchColadaReducer", "Workbench store must use Colada cache as projection reducer authority");
assert.doesNotMatch(workbenchStoreSource, /createWorkbench(?:ReadHydration|ScheduledTask|TraceHydrationQueue)Runtime|runWorkbenchReadHydration/u, "Workbench store must not use legacy Workbench request runtimes");
assert.doesNotMatch(workbenchStoreSource, /api\.workbench\.(?:sessions|sessionMessages|turn|traceEvents)/u, "Workbench store must not call Workbench read APIs outside Colada queries");
assertIncludes(workbenchEventReducerSource, "reduceWorkbenchRealtimeEvent", "Realtime event reducer must own SSE event classification");
assertIncludes(workbenchEventReducerSource, "workbench-event-reducer", "Realtime event reducer must emit module diagnostics for monitor root cause");
assertIncludes(workbenchStoreSource, "reduceWorkbenchRealtimeEvent", "Workbench store must consume migrated realtime reducer actions");
@@ -178,8 +188,8 @@ assert.doesNotMatch(workbenchStoreSource, /case "session\.forget"[\s\S]{0,400}se
assert.doesNotMatch(workbenchStoreSource, /traceHydration(?:InFlight|Queued|Queue)|traceHydrationPumpActive|terminalRealtimeRefreshInFlight|realtimeSessionMessagesInFlight|realtimeErrorGapFillLastAtByKey|activeTraceRestGapFillTimers|realtimeOutboxSeqByKey|sessionListRefreshInFlight|sessionListRefreshTimers/u, "Workbench store must not reintroduce migrated in-flight/timer/cursor state");
assert.doesNotMatch(workbenchStoreSource, /scheduleRealtimeGapHydration|hydrateRealtimeGap/u, "Workbench realtime consumer must not repair projection through REST gap fill");
assert.doesNotMatch(workbenchStoreSource, /subscribeToTrace|TRACE_POLL_INTERVAL_MS/u, "Workbench store must not reintroduce active trace polling");
assertIncludes(workbenchRefreshRuntimeSource, "if (existing)", "Scheduled refresh runtime must coalesce instead of resetting timers under SSE error storms");
assertIncludes(workbenchRefreshRuntimeSource, "replaceTimer", "Scheduled refresh runtime must make replacement explicit instead of ad hoc timer resets");
assertIncludes(workbenchColadaSource, "invalidateQueries", "Realtime recovery must resolve to Colada invalidation/refetch instead of a bespoke REST scheduler");
assertIncludes(workbenchColadaSource, "staleTime", "Workbench query min-interval budget must be expressed as Colada staleTime");
assertIncludes(workbenchPerformanceSource, "recordWorkbenchRuntimeDiagnostic", "Workbench performance probe must record runtime diagnostics for monitor root cause visibility");
assertIncludes(workbenchPerformanceSource, "clearResourceTimings", "Workbench performance probe must bound browser ResourceTiming retention after API enrichment");
assertIncludes(workbenchStoreSource, "recordWorkbenchRuntimeDiagnostic", "Workbench store must surface SSE recovery diagnostics to the performance probe");
@@ -187,10 +197,10 @@ assertIncludes(workbenchRealtimePlanSource, "new Set(recovery.actions)", "Realti
assertIncludes(workbenchRealtimePlanSource, "actions.has(\"schedule-session-list\")", "Realtime stream errors must schedule bounded session list refreshes only when transport requests that action");
assertIncludes(workbenchRealtimePlanSource, "authority: \"automatic-recovery\"", "Realtime recovery planner must classify transport recovery as automatic recovery authority");
assert.doesNotMatch(workbenchRealtimePlanSource, /force:\s*true/u, "Realtime recovery planner must not turn transport recovery into force-refresh work");
assert.doesNotMatch(workbenchRefreshRuntimeSource, /!options\.force\s*&&\s*minIntervalMs/u, "Workbench refresh runtime must not let force bypass minInterval storm budget");
assertIncludes(workbenchRefreshRuntimeSource, "forceBudgetBypass = false", "Workbench refresh runtime diagnostics must make force budget suppression visible");
assertIncludes(workbenchColadaSource, "const state = await queryCache.refresh(entry);", "Workbench reads must preserve Colada staleTime/min-interval governance");
assert.doesNotMatch(workbenchColadaSource, /queryCache\.fetch\(entry/u, "Workbench reads must not call queryCache.fetch(entry), which bypasses Colada freshness governance");
assertIncludes(workbenchStoreSource, "runtimePolicy.sessionListRealtimeRefreshDelayMs", "Realtime recovery delay must come from runtime policy instead of store constants");
assertIncludes(workbenchStoreSource, "workbenchReadCooldownKey(\"session-detail\"", "Realtime session detail recovery must enter Workbench read hydration runtime");
assertIncludes(workbenchStoreSource, "workbenchColadaQueries.fetchSession", "Realtime session detail recovery must enter Colada query facade");
assertIncludes(workbenchStoreSource, "runtimePolicy.workbenchSessionDetailMinRefreshMs", "Realtime session detail recovery budget must come from runtime policy");
assert.doesNotMatch(workbenchStoreSource, /refreshRealtimeSessionMessages[\s\S]{0,900}refreshSessionMessageProjectionPage\(id, \{ force: true \}\)/u, "Realtime session message recovery must not force-bypass the message projection refresh budget");
assert.doesNotMatch(workbenchStoreSource, /handleRealtimeStreamError[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Realtime stream errors must not force-refresh the full session list");
@@ -204,7 +214,7 @@ assertIncludes(serverWorkbenchHttpSource, "workbench_facts_session_missing", "Re
assertIncludes(appSource, "/v1/workbench/events", "Workbench realtime client must use the RESTful same-origin events endpoint");
assertIncludes(appSource, "/v1/workbench/traces/", "trace hydration must use Workbench read-model trace API");
assert.doesNotMatch(appSource, /\/v1\/agent\/(?:turns|traces|chat\/result)\//u, "Cloud Web must not call legacy Code Agent read-through turn/trace/result APIs");
assert.doesNotMatch(workbenchStoreSource, /api\.agent\.(?:getAgentTurn|getAgentTrace)/u, "Workbench store turn/trace gap fill must read through api.workbench");
assert.doesNotMatch(workbenchStoreSource, /api\.(?:agent|workbench)\.(?:getAgentTurn|getAgentTrace|turn|traceEvents|sessions|sessionMessages)/u, "Workbench store turn/trace/session reads must go through Colada queries");
assert.doesNotMatch(appSource, /\/v1\/agent\/chat\/trace\//u, "Cloud Web must not use the legacy action-style chat trace API");
assert.doesNotMatch(appSource, /HWLAB_CLOUD_WEB_EARLY_WORKSPACE_BOOTSTRAP/u, "early workspace bootstrap helper must not remain in Cloud Web app source");
assert.doesNotMatch(appSource, /\/auth\/workspace-bootstrap/u, "stale auth workspace bootstrap route must not remain in Cloud Web app source");
@@ -15,7 +15,6 @@ import { composeWorkbenchScopedKey, splitWorkbenchScopedKey, workbenchPathKey, w
import { createSafeStorageRuntime, isStorageQuotaError, migrateLegacyStorage, normalizePersistedValue, readJsonStorage, removePersistedTarget, removeStorageKey, writeJsonStorage, type StorageLike } from "../src/utils/safe-storage.ts";
import { checkWorkbenchHealth, createWorkbenchHealthProbeCache } from "../src/utils/workbench-health.ts";
import { messageDiagnosticView } from "../src/utils/workbench-error-runtime.ts";
import { createWorkbenchReadHydrationRuntime } from "../src/utils/workbench-refresh-runtime.ts";
import { WORKBENCH_TIMELINE_OPENCODE_PARITY, buildWorkbenchTimelineRows, normalizeWorkbenchTimelineMessages, workbenchTimelineSignature } from "../src/stores/workbench-timeline-model.ts";
import { reduceWorkbenchRealtimeEvent } from "../src/stores/workbench-event-reducer.ts";
import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery } from "../src/stores/workbench-realtime-plan.ts";
@@ -365,25 +364,6 @@ test("realtime recovery planner gates refresh steps by transport actions and aut
assert.deepEqual(inactive.steps.map((step) => step.type), ["schedule-session-list", "refresh-turn-status", "hydrate-trace-events"]);
});
test("read hydration runtime keeps force inside the same storm budget", async () => {
let clock = 1_000;
let calls = 0;
const runtime = createWorkbenchReadHydrationRuntime({ concurrency: 1, failureCooldownMs: 1_000, now: () => clock });
const task = async () => ({ ok: true, status: 200, data: ++calls });
assert.deepEqual(await runtime.run(task, "turn:trc_1", { minIntervalMs: 500, force: true }), { ok: true, status: 200, data: 1 });
const throttled = await runtime.run(task, "turn:trc_1", { minIntervalMs: 500, force: true }) as { ok: boolean; status: number; diagnostic?: Record<string, unknown> };
assert.equal(throttled.ok, false);
assert.equal(throttled.status, 0);
assert.equal(throttled.diagnostic?.code, "workbench_read_hydration_throttled");
assert.equal(throttled.diagnostic?.forceRequested, true);
assert.equal(throttled.diagnostic?.forceBudgetBypass, false);
assert.equal(calls, 1);
clock = 1_501;
assert.deepEqual(await runtime.run(task, "turn:trc_1", { minIntervalMs: 500, force: true }), { ok: true, status: 200, data: 2 });
});
test("health probe cache records ok and unavailable states", async () => {
const cache = createWorkbenchHealthProbeCache({ cacheMs: 100 });
const ok = await cache.probe({ key: "workbench", fetcher: async () => ({ ready: true }), classify: (value) => value.ready ? "ok" : "degraded" });
+9
View File
@@ -1,9 +1,11 @@
import { createApp } from "vue";
import { createPinia } from "pinia";
import { PiniaColada } from "@pinia/colada";
import App from "./App.vue";
import router from "./router";
import i18n, { initI18n } from "./i18n";
import { useAppStore } from "@/stores/app";
import { workbenchRuntimePolicy } from "@/config/workbench-runtime-policy";
import { installWebRum } from "@/utils/rum";
import "./style.css";
import "./styles/workbench.css";
@@ -18,7 +20,14 @@ async function bootstrap(): Promise<void> {
initThemeClass();
const app = createApp(App);
const pinia = createPinia();
const runtimePolicy = workbenchRuntimePolicy();
app.use(pinia);
app.use(PiniaColada, {
queryOptions: {
staleTime: runtimePolicy.workbenchReadFailureCooldownMs,
gcTime: Math.max(runtimePolicy.workbenchReadFailureCooldownMs, runtimePolicy.sessionListMinRefreshIntervalMs)
}
});
useAppStore().initFromInjectedConfig();
await initI18n();
app.use(router);
@@ -0,0 +1,50 @@
// SPEC: pikasTech/HWLAB#2345 Workbench Pinia Colada migration.
// Responsibility: Stable Workbench query and mutation keys for Pinia Colada cache authority.
import type { EntryKey } from "@pinia/colada";
export interface WorkbenchSessionsKeyInput {
includeSessionId?: string | null;
cursor?: string | null;
limit?: number | null;
}
export interface WorkbenchSessionMessagesKeyInput {
cursor?: string | null;
limit?: number | null;
}
export interface WorkbenchTraceEventsKeyInput {
afterProjectedSeq?: number | null;
limit?: number | null;
}
export const workbenchColadaKeys = {
root: (): EntryKey => ["workbench"],
serverState: (): EntryKey => ["workbench", "server-state"],
sessionsRoot: (): EntryKey => ["workbench", "sessions"],
sessions: (input: WorkbenchSessionsKeyInput = {}): EntryKey => ["workbench", "sessions", normalizeKeyObject({ includeSessionId: input.includeSessionId ?? null, cursor: input.cursor ?? null, limit: finiteNumber(input.limit) })],
sessionRoot: (): EntryKey => ["workbench", "session"],
session: (sessionId: string): EntryKey => ["workbench", "session", sessionId],
sessionMessagesRoot: (): EntryKey => ["workbench", "session-messages"],
sessionMessages: (sessionId: string, input: WorkbenchSessionMessagesKeyInput = {}): EntryKey => ["workbench", "session-messages", sessionId, normalizeKeyObject({ cursor: input.cursor ?? null, limit: finiteNumber(input.limit) })],
turnRoot: (): EntryKey => ["workbench", "turn"],
turn: (traceId: string): EntryKey => ["workbench", "turn", traceId],
traceEventsRoot: (): EntryKey => ["workbench", "trace-events"],
traceEvents: (traceId: string, input: WorkbenchTraceEventsKeyInput = {}): EntryKey => ["workbench", "trace-events", traceId, normalizeKeyObject({ afterProjectedSeq: finiteNumber(input.afterProjectedSeq), limit: finiteNumber(input.limit) })],
providerProfiles: (): EntryKey => ["workbench", "provider-profiles"],
mutationSubmit: (): EntryKey => ["workbench", "submit"],
mutationSteer: (): EntryKey => ["workbench", "steer"],
mutationCancel: (traceId: string | null | undefined): EntryKey => ["workbench", "cancel", traceId ?? "unknown"],
mutationCreateSession: (): EntryKey => ["workbench", "create-session"],
mutationDeleteSession: (sessionId: string | null | undefined): EntryKey => ["workbench", "delete-session", sessionId ?? "unknown"]
};
function normalizeKeyObject<T extends Record<string, string | number | boolean | null | undefined>>(input: T): T {
return Object.fromEntries(Object.entries(input).sort(([left], [right]) => left.localeCompare(right))) as T;
}
function finiteNumber(value: number | null | undefined): number | null {
const number = Number(value);
return Number.isFinite(number) ? Math.trunc(number) : null;
}
@@ -0,0 +1,78 @@
// SPEC: pikasTech/HWLAB#2345 Workbench Pinia Colada migration.
// Responsibility: Colada-governed Workbench mutation execution and query invalidation.
import { useMutation, useQueryCache } from "@pinia/colada";
import { api } from "@/api";
import type { ApiRequestOptions } from "@/api/client";
import type { AgentChatResponse, ApiResult } from "@/types";
import { firstNonEmptyString } from "@/utils";
import { workbenchColadaKeys } from "./workbench-colada-keys";
export interface WorkbenchAgentMutationInput {
payload: Record<string, unknown>;
timeoutMs: number;
activityRef?: ApiRequestOptions["activityRef"];
sessionId?: string | null;
traceId?: string | null;
}
export interface WorkbenchCancelMutationInput {
traceId: string;
sessionId?: string | null;
threadId?: string | null;
}
export interface WorkbenchColadaMutations {
createAgentSession: (payload?: Record<string, unknown>) => Promise<ApiResult<{ session?: Record<string, unknown> }>>;
submitAgentMessage: (input: WorkbenchAgentMutationInput) => Promise<ApiResult<AgentChatResponse>>;
steerAgentMessage: (input: WorkbenchAgentMutationInput) => Promise<ApiResult<AgentChatResponse>>;
cancelAgentMessage: (input: WorkbenchCancelMutationInput) => Promise<ApiResult<AgentChatResponse & { canceled?: boolean; alreadyTerminal?: boolean; cancelStatus?: string; error?: { code?: string; message?: string; userMessage?: string } }>>;
deleteAgentSession: (input: { sessionId: string }) => Promise<ApiResult<{ ok?: boolean; status?: string; session?: Record<string, unknown>; archivedCount?: number }>>;
}
export function useWorkbenchColadaMutations(): WorkbenchColadaMutations {
const queryCache = useQueryCache();
const createSessionMutation = useMutation<ApiResult<{ session?: Record<string, unknown> }>, Record<string, unknown> | undefined>({
key: workbenchColadaKeys.mutationCreateSession(),
mutation: (payload) => api.agent.createAgentSession(payload ?? {}),
onSettled: () => { void queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionsRoot() }); }
});
const submitMutation = useMutation<ApiResult<AgentChatResponse>, WorkbenchAgentMutationInput>({
key: workbenchColadaKeys.mutationSubmit(),
mutation: (input) => api.agent.sendAgentMessage(input.payload, input.timeoutMs, input.activityRef),
onSettled: (result, _error, input) => { void invalidateMutationScope(queryCache, input.sessionId, firstNonEmptyString(result?.data?.traceId, input.traceId)); }
});
const steerMutation = useMutation<ApiResult<AgentChatResponse>, WorkbenchAgentMutationInput>({
key: workbenchColadaKeys.mutationSteer(),
mutation: (input) => api.agent.steerAgentMessage(input.payload, input.timeoutMs, input.activityRef),
onSettled: (result, _error, input) => { void invalidateMutationScope(queryCache, input.sessionId, firstNonEmptyString(result?.data?.traceId, input.traceId)); }
});
const cancelMutation = useMutation<ApiResult<AgentChatResponse & { canceled?: boolean; alreadyTerminal?: boolean; cancelStatus?: string; error?: { code?: string; message?: string; userMessage?: string } }>, WorkbenchCancelMutationInput>({
key: (input) => workbenchColadaKeys.mutationCancel(input.traceId),
mutation: (input) => api.agent.cancelAgentMessage({ traceId: input.traceId, sessionId: input.sessionId, threadId: input.threadId }),
onSettled: (_result, _error, input) => { void invalidateMutationScope(queryCache, input.sessionId, input.traceId); }
});
const deleteMutation = useMutation<ApiResult<{ ok?: boolean; status?: string; session?: Record<string, unknown>; archivedCount?: number }>, { sessionId: string }>({
key: (input) => workbenchColadaKeys.mutationDeleteSession(input.sessionId),
mutation: (input) => api.agent.deleteAgentSession(input.sessionId),
onSettled: (_result, _error, input) => { void invalidateMutationScope(queryCache, input.sessionId, null); }
});
return {
createAgentSession: (payload) => createSessionMutation.mutateAsync(payload),
submitAgentMessage: (input) => submitMutation.mutateAsync(input),
steerAgentMessage: (input) => steerMutation.mutateAsync(input),
cancelAgentMessage: (input) => cancelMutation.mutateAsync(input),
deleteAgentSession: (input) => deleteMutation.mutateAsync(input)
};
}
async function invalidateMutationScope(queryCache: ReturnType<typeof useQueryCache>, sessionId: string | null | undefined, traceId: string | null | undefined): Promise<void> {
await Promise.all([
queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionsRoot() }),
sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.session(sessionId) }) : Promise.resolve(),
sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionMessagesRoot().concat(sessionId) }) : Promise.resolve(),
traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.turn(traceId) }) : Promise.resolve(),
traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.traceEventsRoot().concat(traceId) }) : Promise.resolve()
]);
}
@@ -0,0 +1,83 @@
// SPEC: pikasTech/HWLAB#2345 Workbench Pinia Colada migration.
// Responsibility: Colada-governed Workbench read-model query execution, cache keys, invalidation and refetch.
import { useQueryCache, type EntryKey, type QueryCache } from "@pinia/colada";
import { api } from "@/api";
import type { ApiRequestOptions } from "@/api/client";
import type { AgentChatResultResponse, ApiResult } from "@/types";
import type { SessionListOptions, SessionMessageOptions, WorkbenchMessagePageResponse, WorkbenchSessionDetailResponse, WorkbenchSessionListResponse, WorkbenchTraceRequestOptions } from "@/api/workbench";
import { workbenchColadaKeys } from "./workbench-colada-keys";
interface QueryRunOptions {
force?: boolean;
minIntervalMs?: number | null;
}
export interface WorkbenchTurnQueryOptions extends QueryRunOptions {
timeoutMs?: number;
activityRef?: ApiRequestOptions["activityRef"];
}
export interface WorkbenchTraceEventsQueryOptions extends QueryRunOptions, WorkbenchTraceRequestOptions {
timeoutMs?: number;
activityRef?: ApiRequestOptions["activityRef"];
}
export interface WorkbenchColadaQueries {
queryCache: QueryCache;
fetchSessions: (options?: SessionListOptions & QueryRunOptions) => Promise<ApiResult<WorkbenchSessionListResponse>>;
fetchSession: (sessionId: string, options?: QueryRunOptions & { timeoutMs?: number | null }) => Promise<ApiResult<WorkbenchSessionDetailResponse>>;
fetchSessionMessages: (sessionId: string, options?: SessionMessageOptions & QueryRunOptions) => Promise<ApiResult<WorkbenchMessagePageResponse>>;
fetchTurn: (traceId: string, options?: WorkbenchTurnQueryOptions) => Promise<ApiResult<AgentChatResultResponse>>;
fetchTraceEvents: (traceId: string, options?: WorkbenchTraceEventsQueryOptions) => Promise<ApiResult<AgentChatResultResponse>>;
fetchProviderProfiles: (options?: QueryRunOptions) => Promise<ApiResult<unknown>>;
invalidateSessionList: () => Promise<unknown>;
invalidateSession: (sessionId: string | null | undefined) => Promise<unknown>;
invalidateSessionMessages: (sessionId: string | null | undefined) => Promise<unknown>;
invalidateTurn: (traceId: string | null | undefined) => Promise<unknown>;
invalidateTraceEvents: (traceId: string | null | undefined) => Promise<unknown>;
invalidateTraceScope: (input: { sessionId?: string | null; traceId?: string | null }) => Promise<unknown[]>;
}
export function useWorkbenchColadaQueries(): WorkbenchColadaQueries {
const queryCache = useQueryCache();
return {
queryCache,
fetchSessions: (options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.sessions(options), () => api.workbench.sessions(options), options),
fetchSession: (sessionId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.session(sessionId), () => api.workbench.session(sessionId, options.timeoutMs ?? null), options),
fetchSessionMessages: (sessionId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.sessionMessages(sessionId, options), () => api.workbench.sessionMessages(sessionId, options), options),
fetchTurn: (traceId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.turn(traceId), () => api.workbench.turn(traceId, options.timeoutMs ?? 8000, options.activityRef), options),
fetchTraceEvents: (traceId, options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.traceEvents(traceId, options), () => api.workbench.traceEvents(traceId, options.timeoutMs ?? 8000, options.activityRef, { afterProjectedSeq: options.afterProjectedSeq, limit: options.limit }), options),
fetchProviderProfiles: (options = {}) => runWorkbenchQuery(queryCache, workbenchColadaKeys.providerProfiles(), () => api.providerProfiles.catalog(), options),
invalidateSessionList: () => queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionsRoot() }),
invalidateSession: (sessionId) => sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.session(sessionId) }) : Promise.resolve(),
invalidateSessionMessages: (sessionId) => sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionMessagesRoot().concat(sessionId) }) : Promise.resolve(),
invalidateTurn: (traceId) => traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.turn(traceId) }) : Promise.resolve(),
invalidateTraceEvents: (traceId) => traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.traceEventsRoot().concat(traceId) }) : Promise.resolve(),
invalidateTraceScope: async (input) => Promise.all([
queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionsRoot() }),
input.sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.session(input.sessionId) }) : Promise.resolve(),
input.sessionId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.sessionMessagesRoot().concat(input.sessionId) }) : Promise.resolve(),
input.traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.turn(input.traceId) }) : Promise.resolve(),
input.traceId ? queryCache.invalidateQueries({ key: workbenchColadaKeys.traceEventsRoot().concat(input.traceId) }) : Promise.resolve()
])
};
}
async function runWorkbenchQuery<TResult>(queryCache: QueryCache, key: EntryKey, query: () => Promise<TResult>, options: QueryRunOptions): Promise<TResult> {
const entry = queryCache.ensure<TResult, Error, undefined>({
key,
query,
enabled: false,
staleTime: normalizeStaleTime(options.minIntervalMs)
});
// refresh() preserves Colada staleTime/min-interval governance; fetch() would bypass it.
const state = await queryCache.refresh(entry);
return state.data as TResult;
}
function normalizeStaleTime(value: number | null | undefined): number {
const number = Number(value);
return Number.isFinite(number) && number > 0 ? Math.trunc(number) : 0;
}
@@ -0,0 +1,41 @@
// SPEC: pikasTech/HWLAB#2345 Workbench Pinia Colada migration.
// Responsibility: Store Workbench reducer projection in Pinia Colada cache instead of a local Pinia ref.
import { computed, type ComputedRef } from "vue";
import { useQuery, useQueryCache, type QueryCache } from "@pinia/colada";
import { workbenchColadaKeys } from "./workbench-colada-keys";
import { createWorkbenchServerState, reduceWorkbenchServerState, type WorkbenchServerAction, type WorkbenchServerState } from "./workbench-server-state";
export interface WorkbenchColadaReducer {
queryCache: QueryCache;
serverState: ComputedRef<WorkbenchServerState>;
readServerState: () => WorkbenchServerState;
replaceServerState: (project: (state: WorkbenchServerState) => WorkbenchServerState) => void;
reduceServerState: (action: WorkbenchServerAction) => void;
}
export function useWorkbenchColadaReducer(): WorkbenchColadaReducer {
const queryCache = useQueryCache();
const key = workbenchColadaKeys.serverState();
const serverStateQuery = useQuery<WorkbenchServerState, Error, WorkbenchServerState>({
key,
query: async () => queryCache.getQueryData<WorkbenchServerState>(key) ?? createWorkbenchServerState(),
initialData: createWorkbenchServerState,
enabled: false
});
const serverState = computed<WorkbenchServerState>(() => serverStateQuery.data.value ?? createWorkbenchServerState());
function readServerState(): WorkbenchServerState {
return queryCache.getQueryData<WorkbenchServerState>(key) ?? serverState.value;
}
function replaceServerState(project: (state: WorkbenchServerState) => WorkbenchServerState): void {
queryCache.setQueryData<WorkbenchServerState>(key, (source) => project(source ?? createWorkbenchServerState()));
}
function reduceServerState(action: WorkbenchServerAction): void {
replaceServerState((source) => reduceWorkbenchServerState(source, action));
}
return { queryCache, serverState, readServerState, replaceServerState, reduceServerState };
}
@@ -0,0 +1,52 @@
import assert from "node:assert/strict";
import fs from "node:fs";
import path from "node:path";
import { fileURLToPath } from "node:url";
import { test } from "bun:test";
import { createApp } from "vue";
import { createPinia, setActivePinia } from "pinia";
import { PiniaColada, useQueryCache } from "@pinia/colada";
import { createWorkbenchServerState, type WorkbenchServerState } from "./workbench-server-state";
import { workbenchColadaKeys } from "./workbench-colada-keys";
import { useWorkbenchColadaReducer } from "./workbench-colada-reducer";
const storeDir = path.dirname(fileURLToPath(import.meta.url));
function installColadaContext(): ReturnType<typeof createApp> {
const app = createApp({});
const pinia = createPinia();
app.use(pinia);
app.use(PiniaColada, {});
setActivePinia(pinia);
return app;
}
test("workbench colada keys keep serializable hierarchical cache scopes", () => {
assert.deepEqual(workbenchColadaKeys.sessions({ limit: 50, includeSessionId: "ses_colada", cursor: "cur_2" }), ["workbench", "sessions", { cursor: "cur_2", includeSessionId: "ses_colada", limit: 50 }]);
assert.deepEqual(workbenchColadaKeys.sessionMessages("ses_colada", { limit: 100 }), ["workbench", "session-messages", "ses_colada", { cursor: null, limit: 100 }]);
assert.deepEqual(workbenchColadaKeys.traceEvents("trc_colada", { afterProjectedSeq: 4, limit: 25 }), ["workbench", "trace-events", "trc_colada", { afterProjectedSeq: 4, limit: 25 }]);
});
test("workbench colada reducer stores projection state in the query cache", () => {
const app = installColadaContext();
app.runWithContext(() => {
const reducer = useWorkbenchColadaReducer();
const queryCache = useQueryCache();
const key = workbenchColadaKeys.serverState();
assert.deepEqual(reducer.readServerState(), createWorkbenchServerState());
reducer.reduceServerState({ type: "session.detail", session: { sessionId: "ses_colada", status: "running", messages: [] } });
const cached = queryCache.getQueryData<WorkbenchServerState>(key);
assert.equal(cached?.sessionsById.ses_colada?.sessionId, "ses_colada");
assert.equal(reducer.serverState.value.sessionsById.ses_colada?.sessionId, "ses_colada");
});
});
test("workbench colada read path preserves freshness governance", () => {
const source = fs.readFileSync(path.join(storeDir, "workbench-colada-queries.ts"), "utf8");
assert.match(source, /const state = await queryCache\.refresh\(entry\);/u);
assert.doesNotMatch(source, /queryCache\.fetch\(entry/u);
});
+97 -181
View File
@@ -7,7 +7,6 @@ import { api } from "@/api";
import { workbenchRuntimePolicy } from "@/config/workbench-runtime-policy";
import { createWorkbenchHealthProbeCache } from "@/utils/workbench-health";
import { agentErrorFromProjection, normalizeApiErrorRecord, normalizeErrorDiagnostic, normalizeProjectionDiagnostic, projectionDiagnosticFromApiFailure, projectionDiagnosticFromFailure } from "@/utils/workbench-error-runtime";
import { createWorkbenchReadHydrationRuntime, createWorkbenchScheduledTaskRuntime, createWorkbenchTraceHydrationQueueRuntime, shouldCooldownWorkbenchReadFailure as shouldCooldownWorkbenchReadRuntimeFailure } from "@/utils/workbench-refresh-runtime";
import { readWorkbenchJson, readWorkbenchNumber, readWorkbenchString, removeWorkbenchStorageKey, writeWorkbenchJson, writeWorkbenchString } from "@/utils/workbench-storage-runtime";
import { createWorkbenchStreamTransportRuntime, type WorkbenchRealtimeEvent, type WorkbenchStreamTransportRecovery } from "@/utils/workbench-realtime-runtime";
import { mergeRunnerTrace, snapshotToRunnerTrace, type TraceSnapshot } from "@/composables/useTraceSubscription";
@@ -17,7 +16,7 @@ import { composeWorkbenchScopedKey } from "@/utils/workbench-key";
import { failWorkbenchSessionSwitch, failWorkbenchSubmitJourney, finishWorkbenchSessionSwitchFullLoad, markWorkbenchSubmitApiAccepted, markWorkbenchTraceEventsReceived, markWorkbenchTraceProjected, recordWorkbenchLoadingState, recordWorkbenchRuntimeDiagnostic, startWorkbenchSessionSwitch, startWorkbenchSubmitJourney } from "@/utils/workbench-performance";
import { RECENT_DRAFTS_STORAGE_KEY, appendSessionPage, defaultProviderProfileOptions, isArchivedSession, mergeSessionIntoList, normalizeChatMessageStatus, normalizeRecentDrafts, normalizeWorkbenchMessageTitle, providerProfileOptionsFromPayload, recordRecentDraft, resolveCancelableAgentMessage, resolveComposerState, selectActiveTurnStatusRefreshTraceIds, shouldShowSessionListLoading, sortSessionTabs, stableSessionList, type DraftEntry, type ProviderProfileOption, type TurnStatusAuthority } from "./workbench-session";
import { initialWorkbenchSessionIdFromLocation } from "./workbench-projection";
import { cleanupWorkbenchServerStateSessions, createWorkbenchServerState, reduceWorkbenchServerState, selectActiveMessages, selectActiveSession, selectSessionList, selectSessionStatusAuthority, selectTraceAuthorityById, selectTurnStatusAuthority, type WorkbenchServerAction } from "./workbench-server-state";
import { cleanupWorkbenchServerStateSessions, selectActiveMessages, selectActiveSession, selectSessionList, selectSessionStatusAuthority, selectTraceAuthorityById, selectTurnStatusAuthority, type WorkbenchServerAction } from "./workbench-server-state";
import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache";
import { reduceWorkbenchRealtimeEvent, type WorkbenchRealtimeAction } from "./workbench-event-reducer";
import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority";
@@ -64,6 +63,9 @@ import {
traceSnapshotError
} from "./workbench-message-projection-runtime";
import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery, type WorkbenchRealtimeApplyStep, type WorkbenchRealtimeRecoveryStep } from "./workbench-realtime-plan";
import { useWorkbenchColadaMutations } from "./workbench-colada-mutations";
import { useWorkbenchColadaQueries } from "./workbench-colada-queries";
import { useWorkbenchColadaReducer } from "./workbench-colada-reducer";
const WORKBENCH_SESSION_PROJECTION_SIGNAL_CHANNEL = "hwlab.workbench.sessionProjection.v1";
const WORKBENCH_SESSION_PROJECTION_SIGNAL_KEY = "hwlab.workbench.sessionProjectionSignal.v1";
@@ -81,6 +83,9 @@ interface SelectSessionOptions {
export const useWorkbenchStore = defineStore("workbench", () => {
const runtimePolicy = workbenchRuntimePolicy();
const workbenchColadaReducer = useWorkbenchColadaReducer();
const workbenchColadaQueries = useWorkbenchColadaQueries();
const workbenchColadaMutations = useWorkbenchColadaMutations();
const providerProfile = ref<ProviderProfile>(readString("hwlab.workbench.providerProfile.v1", "codex"));
const providerOptions = ref<ProviderProfileOption[]>(defaultProviderProfileOptions(providerProfile.value));
const recentDrafts = ref<DraftEntry[]>(readRecentDrafts());
@@ -102,22 +107,14 @@ export const useWorkbenchStore = defineStore("workbench", () => {
const explicitSessionId = ref<string | null>(initialWorkbenchSessionIdFromLocation());
const activeSelectionSource = ref<SessionSelectionSource>(explicitSessionId.value ? "route" : "system");
const selectionEpoch = ref(0);
const serverState = ref(createWorkbenchServerState());
const serverState = workbenchColadaReducer.serverState;
const sessions = computed(() => selectSessionList(serverState.value));
const sessionStatusAuthority = computed(() => selectSessionStatusAuthority(serverState.value));
const turnStatusAuthority = computed(() => selectTurnStatusAuthority(serverState.value));
const traceAuthorityById = computed(() => selectTraceAuthorityById(serverState.value));
const realtimeTransport = createWorkbenchStreamTransportRuntime();
const traceHydrationRuntime = createWorkbenchTraceHydrationQueueRuntime<ChatMessage>({ concurrency: runtimePolicy.traceHydrationBackgroundConcurrency, delayMs: runtimePolicy.traceHydrationBackgroundDelayMs });
const forcedTraceHydrationRetryRuntime = createWorkbenchScheduledTaskRuntime();
const terminalRealtimeRefreshRuntime = createWorkbenchScheduledTaskRuntime();
const realtimeSessionMessagesRuntime = createWorkbenchScheduledTaskRuntime();
const realtimeSessionRefreshRuntime = createWorkbenchScheduledTaskRuntime();
const activeTraceRestGapFillRuntime = createWorkbenchScheduledTaskRuntime();
const sessionListRefreshRuntime = createWorkbenchScheduledTaskRuntime();
const workbenchProjectionSignalSourceId = nextProtocolId("wbtab");
let workbenchProjectionSignalChannel: BroadcastChannel | null = null;
const workbenchReadHydrationRuntime = createWorkbenchReadHydrationRuntime({ concurrency: runtimePolicy.workbenchReadHydrationConcurrency, failureCooldownMs: runtimePolicy.workbenchReadFailureCooldownMs });
const workbenchHealthProbeCache = createWorkbenchHealthProbeCache({ cacheMs: runtimePolicy.workbenchReadFailureCooldownMs });
const projectedActiveSession = computed(() => selectActiveSession(serverState.value, explicitSessionId.value));
@@ -156,7 +153,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
recordWorkbenchLoadingState({ scope: "session_detail", active: true, reason: "hydrate", sessionId: routeSessionId });
setActiveSessionSelection(routeSessionId, "route");
}
const sessionsResult = await api.workbench.sessions({ includeSessionId, limit: runtimePolicy.sessionListPageLimit });
const sessionsResult = await workbenchColadaQueries.fetchSessions({ includeSessionId, limit: runtimePolicy.sessionListPageLimit, minIntervalMs: runtimePolicy.sessionListMinRefreshIntervalMs });
const listedSessions = sessionsResult.ok ? workbenchSessionsFromPayload(sessionsResult.data) : [];
if (sessionsResult.ok) {
applySessionPagination(sessionsResult.data);
@@ -231,7 +228,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
recordWorkbenchLoadingState({ scope: "workbench", active: true, reason: "create_session", sessionId: activeSessionId.value });
error.value = null;
void refreshProviderOptions();
const response = await api.agent.createAgentSession({ providerProfile: providerProfile.value });
const response = await workbenchColadaMutations.createAgentSession({ providerProfile: providerProfile.value });
loading.value = false;
recordWorkbenchLoadingState({ scope: "workbench", active: false, reason: "create_session", sessionId: activeSessionId.value });
const created = response.ok ? sessionFromWorkbenchSession(response.data?.session) : null;
@@ -307,7 +304,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
const sessionId = activeSessionId.value;
if (!sessionId) return;
error.value = null;
const response = await api.agent.deleteAgentSession(sessionId);
const response = await workbenchColadaMutations.deleteAgentSession({ sessionId });
if (!response.ok) {
error.value = response.error ?? "session delete failed";
return;
@@ -327,12 +324,11 @@ export const useWorkbenchStore = defineStore("workbench", () => {
async function refreshSessions(includeSessionId: string | null = activeSessionId.value, options: { force?: boolean } = {}): Promise<void> {
const requestIncludeSessionId = firstNonEmptyString(includeSessionId);
const requestLimit = currentSessionListLimit();
const requestKey = sessionListRefreshKey(requestIncludeSessionId, requestLimit);
await sessionListRefreshRuntime.run(requestKey, () => refreshSessionsNow(requestIncludeSessionId, requestLimit), { force: options.force, minIntervalMs: runtimePolicy.sessionListMinRefreshIntervalMs, reason: "session-list" });
await refreshSessionsNow(requestIncludeSessionId, requestLimit, { force: options.force });
}
async function refreshSessionsNow(includeSessionId: string | null, limit: number): Promise<void> {
const response = await api.workbench.sessions({ includeSessionId, limit });
async function refreshSessionsNow(includeSessionId: string | null, limit: number, options: { force?: boolean } = {}): Promise<void> {
const response = await workbenchColadaQueries.fetchSessions({ includeSessionId, limit, force: options.force, minIntervalMs: runtimePolicy.sessionListMinRefreshIntervalMs });
if (response.ok) {
const listed = workbenchSessionsFromPayload(response.data);
const activeSeed = activeSessionListSeed();
@@ -349,18 +345,12 @@ export const useWorkbenchStore = defineStore("workbench", () => {
}
function scheduleSessionListRefresh(includeSessionId: string | null | undefined = activeSessionId.value, delayMs = runtimePolicy.sessionListRealtimeRefreshDelayMs): void {
const requestIncludeSessionId = firstNonEmptyString(includeSessionId);
const requestLimit = currentSessionListLimit();
const requestKey = sessionListRefreshKey(requestIncludeSessionId, requestLimit);
void includeSessionId;
if (typeof window === "undefined") {
void refreshSessions(requestIncludeSessionId);
void workbenchColadaQueries.invalidateSessionList();
return;
}
sessionListRefreshRuntime.schedule(requestKey, () => refreshSessions(requestIncludeSessionId), { delayMs, minIntervalMs: runtimePolicy.sessionListMinRefreshIntervalMs, reason: "session-list" });
}
function sessionListRefreshKey(includeSessionId: string | null | undefined, limit: number): string {
return composeWorkbenchScopedKey("workbench.session-list", firstNonEmptyString(includeSessionId), limit);
window.setTimeout(() => { void workbenchColadaQueries.invalidateSessionList(); }, delayMs);
}
async function loadMoreSessions(): Promise<void> {
@@ -372,7 +362,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
}
sessionListLoadingMore.value = true;
sessionListLoadMoreError.value = null;
const response = await api.workbench.sessions({ includeSessionId: activeSessionId.value, limit: runtimePolicy.sessionListPageLimit, cursor });
const response = await workbenchColadaQueries.fetchSessions({ includeSessionId: activeSessionId.value, limit: runtimePolicy.sessionListPageLimit, cursor, force: true });
sessionListLoadingMore.value = false;
if (!response.ok) {
sessionListLoadMoreError.value = response.error ?? "session list load more failed";
@@ -397,7 +387,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
}
function reduceServerState(action: WorkbenchServerAction): void {
serverState.value = reduceWorkbenchServerState(serverState.value, action);
workbenchColadaReducer.reduceServerState(action);
}
function rememberSessionDetail(session: WorkbenchSessionRecord | null | undefined): void {
@@ -410,12 +400,12 @@ export const useWorkbenchStore = defineStore("workbench", () => {
const trimmed = trimWorkbenchSessionCache(beforeTrim, { retainSessionIds: [activeSessionId.value, explicitSessionId.value], maxSessions: currentSessionListLimit() });
const cleaned = cleanupDroppedWorkbenchSessionCaches(trimmed.state, trimmed.state.sessionOrder);
const authorityCleaned = cleanupWorkbenchServerStateSessions(beforeTrim, [...trimmed.evictedSessionIds, ...cleaned.droppedSessionIds]);
serverState.value = {
workbenchColadaReducer.replaceServerState(() => ({
...authorityCleaned,
sessionOrder: cleaned.state.sessionOrder,
sessionsById: cleaned.state.sessionsById,
messagesBySessionId: cleaned.state.messagesBySessionId
};
}));
}
function forgetSession(sessionId: string): void {
@@ -477,43 +467,14 @@ export const useWorkbenchStore = defineStore("workbench", () => {
void source;
}
function runWorkbenchReadHydration<T>(task: () => Promise<T>, cooldownKey?: string, options: { minIntervalMs?: number; force?: boolean } = {}): Promise<T> {
return workbenchReadHydrationRuntime.run(task, cooldownKey, {
minIntervalMs: options.minIntervalMs,
force: options.force,
shouldCooldownFailure: (value) => isApiResultLike(value) && shouldCooldownWorkbenchReadFailure(value),
reason: cooldownKey ?? null,
});
}
function workbenchReadCooldownKey(kind: string, id: string): string {
return composeWorkbenchScopedKey("workbench.read", kind, id);
}
function isApiResultLike(value: unknown): value is ApiResult<unknown> {
return Boolean(value && typeof value === "object" && "ok" in value && "status" in value);
}
function shouldCooldownWorkbenchReadFailure(result: ApiResult<unknown>): boolean {
return shouldCooldownWorkbenchReadRuntimeFailure(result);
}
function fetchWorkbenchTurnStatus(traceId: string, useActivityTimeout = shouldUseActivityTimeoutForTrace(traceId), options: { force?: boolean } = {}): Promise<ApiResult<AgentChatResultResponse>> {
const activitySource = useActivityTimeout ? () => activityRef.value : null;
return runWorkbenchReadHydration(
() => api.workbench.turn(traceId, 8000, activitySource),
workbenchReadCooldownKey("turn", traceId),
{ minIntervalMs: runtimePolicy.workbenchTurnStatusMinRefreshMs, force: options.force },
);
return workbenchColadaQueries.fetchTurn(traceId, { timeoutMs: 8000, activityRef: activitySource, minIntervalMs: runtimePolicy.workbenchTurnStatusMinRefreshMs, force: options.force });
}
function fetchWorkbenchTraceEvents(traceId: string, afterProjectedSeq: number, useActivityTimeout = shouldUseActivityTimeoutForTrace(traceId), options: { force?: boolean } = {}): Promise<ApiResult<AgentChatResultResponse>> {
const activitySource = useActivityTimeout ? () => activityRef.value : null;
return runWorkbenchReadHydration(
() => api.workbench.traceEvents(traceId, runtimePolicy.workbenchTraceEventsTimeoutMs, activitySource, { afterProjectedSeq, limit: runtimePolicy.traceHydrationPageLimit }),
workbenchReadCooldownKey("trace-events", traceId),
{ minIntervalMs: runtimePolicy.workbenchTraceEventsMinRefreshMs, force: options.force },
);
return workbenchColadaQueries.fetchTraceEvents(traceId, { timeoutMs: runtimePolicy.workbenchTraceEventsTimeoutMs, activityRef: activitySource, afterProjectedSeq, limit: runtimePolicy.traceHydrationPageLimit, minIntervalMs: runtimePolicy.workbenchTraceEventsMinRefreshMs, force: options.force });
}
function shouldUseActivityTimeoutForTrace(traceId: string | null | undefined): boolean {
@@ -530,11 +491,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
async function refreshSessionMessageProjectionPage(sessionId: string | null | undefined, options: { force?: boolean } = {}): Promise<void> {
const id = normalizeWorkbenchSessionId(sessionId);
if (!id) return;
const response = await runWorkbenchReadHydration(
() => api.workbench.sessionMessages(id, { limit: 100 }),
workbenchReadCooldownKey("session-messages", id),
{ minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force },
);
const response = await workbenchColadaQueries.fetchSessionMessages(id, { limit: 100, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force });
if (!response.ok || !response.data) return;
const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : [];
const merged = mergeMessageProjectionPage(id, pageMessages);
@@ -546,20 +503,14 @@ export const useWorkbenchStore = defineStore("workbench", () => {
async function refreshRealtimeSessionMessages(sessionId: string | null | undefined, reason: string, options: { force?: boolean } = {}): Promise<void> {
const id = normalizeWorkbenchSessionId(sessionId);
if (!id || id !== activeSessionId.value) return;
await realtimeSessionMessagesRuntime.run(composeWorkbenchScopedKey("workbench.realtime.session-messages", id), async () => {
recordActivity(reason);
await refreshSessionMessageProjectionPage(id, { force: options.force });
}, { force: options.force, minIntervalMs: runtimePolicy.workbenchRealtimeSessionMessagesMinRefreshMs, reason });
recordActivity(reason);
await refreshSessionMessageProjectionPage(id, { force: options.force });
}
async function refreshMessageProjectionForTrace(sessionId: string | null | undefined, traceId: string, options: { force?: boolean } = {}): Promise<void> {
const id = normalizeWorkbenchSessionId(sessionId);
if (!id) return;
const response = await runWorkbenchReadHydration(
() => api.workbench.sessionMessages(id, { limit: 100 }),
workbenchReadCooldownKey("session-messages", id),
{ minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force },
);
const response = await workbenchColadaQueries.fetchSessionMessages(id, { limit: 100, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force });
if (!response.ok || !response.data) {
if (shouldSuppressTransientWorkbenchReadFailure(response)) return;
if (traceHasTerminalResponse(traceId, messages.value)) return;
@@ -697,8 +648,9 @@ export const useWorkbenchStore = defineStore("workbench", () => {
rememberRecentDraft(value);
recordActivity("submit");
const payload = { message: value, prompt: value, sessionId, threadId: providerThreadId, traceId, providerProfile: providerProfile.value, gatewayShellTimeoutMs: gatewayShellTimeoutMs.value, targetTraceId: composer.value.targetTraceId, steerTraceId };
const route = steerMode ? api.agent.steerAgentMessage : api.agent.sendAgentMessage;
const response = await route(payload, codeAgentTimeoutMs.value, () => activityRef.value);
const response = steerMode
? await workbenchColadaMutations.steerAgentMessage({ payload, timeoutMs: codeAgentTimeoutMs.value, activityRef: () => activityRef.value, sessionId, traceId })
: await workbenchColadaMutations.submitAgentMessage({ payload, timeoutMs: codeAgentTimeoutMs.value, activityRef: () => activityRef.value, sessionId, traceId });
if (!response.ok || !response.data) {
const projection = projectionDiagnosticFromApiFailure(response, { code: "code_agent_admission_failed", message: response.error ?? "Code Agent 请求失败", health: response.status === 0 ? "unavailable" : "degraded" });
const failureError = agentErrorFromApiFailure(response, projection, "Code Agent 请求失败");
@@ -760,7 +712,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
const traceId = message.traceId ?? message.runnerTrace?.traceId;
if (!traceId) return;
const sessionId = message.sessionId ?? selectedSessionId.value;
const response = await api.agent.cancelAgentMessage({ traceId, sessionId: message.sessionId ?? selectedSessionId.value, threadId: message.threadId ?? selectedThreadId.value });
const response = await workbenchColadaMutations.cancelAgentMessage({ traceId, sessionId: message.sessionId ?? selectedSessionId.value, threadId: message.threadId ?? selectedThreadId.value });
if (!response.ok) {
applyProjectionDiagnostic(traceId, projectionDiagnosticFromApiFailure(response, { code: "code_agent_cancel_failed", message: response.error ?? "Code Agent 取消请求失败,当前 turn 状态保持 canonical projection。", health: response.status === 0 ? "unavailable" : "degraded" }));
return;
@@ -794,9 +746,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
async function hydrateTraceEventsForMessage(message: ChatMessage, options: { force?: boolean } = {}): Promise<void> {
const traceId = message.traceId ?? message.runnerTrace?.traceId;
if (!traceId) return;
await traceHydrationRuntime.run(traceId, () => hydrateTraceEventsForMessageNow(message, options), {
onAlreadyRunning: options.force ? () => scheduleForcedTraceHydrationRetry(message) : undefined
});
await hydrateTraceEventsForMessageNow(message, options);
}
async function hydrateTraceEventsForMessageNow(message: ChatMessage, options: { force?: boolean } = {}): Promise<void> {
@@ -819,18 +769,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
}
async function fetchTraceHydrationPage(traceId: string, afterProjectedSeq: number, options: { force?: boolean } = {}): Promise<ApiResult<AgentChatResultResponse>> {
let lastResult: ApiResult<AgentChatResultResponse> | null = null;
for (let attempt = 0; attempt < runtimePolicy.traceHydrationMaxAttempts; attempt += 1) {
const result = await fetchWorkbenchTraceEvents(traceId, afterProjectedSeq, shouldUseActivityTimeoutForTrace(traceId), { force: options.force });
if (result.ok && result.data) return result;
lastResult = result;
if (attempt < runtimePolicy.traceHydrationMaxAttempts - 1) await delayTraceHydrationRetry(runtimePolicy.traceHydrationRetryDelayMs * (attempt + 1));
}
return lastResult ?? { ok: false, status: 0, data: null, error: "trace_hydration_failed" };
}
function delayTraceHydrationRetry(ms: number): Promise<void> {
return new Promise((resolve) => window.setTimeout(resolve, ms));
return fetchWorkbenchTraceEvents(traceId, afterProjectedSeq, shouldUseActivityTimeoutForTrace(traceId), { force: options.force });
}
function hydrateTerminalTraceGaps(source: ChatMessage[], reason: string): void {
@@ -840,31 +779,14 @@ export const useWorkbenchStore = defineStore("workbench", () => {
void reason;
}
function scheduleForcedTraceHydrationRetry(message: ChatMessage): void {
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
if (!traceId || typeof window === "undefined") return;
forcedTraceHydrationRetryRuntime.schedule(composeWorkbenchScopedKey("workbench.trace-hydration.retry", traceId), async () => {
const ownerSessionId = traceOwnerSessionId(traceId, messageSessionAuthority(message));
const ownerMessages = ownerSessionId ? serverState.value.messagesBySessionId[ownerSessionId] ?? [] : messages.value;
const latest = latestMessageForTrace(traceId, ownerMessages) ?? message;
await hydrateTraceEventsForMessage(latest, { force: true });
}, { delayMs: runtimePolicy.traceHydrationRetryDelayMs, reason: "trace-hydration-retry" });
}
async function hydrateTraceEvents(source: ChatMessage[] = messages.value): Promise<void> {
for (const message of traceHydrationCandidates(source)) queueTraceHydration(message);
for (const message of traceHydrationCandidates(source)) void hydrateTraceEventsForMessage(message);
}
function traceHydrationCandidates(source: ChatMessage[]): ChatMessage[] {
return source.filter(messageNeedsTraceHydration).slice(-runtimePolicy.traceHydrationAutoQueueLimit).reverse();
}
function queueTraceHydration(message: ChatMessage): void {
const traceId = message.traceId ?? message.runnerTrace?.traceId;
if (!traceId) return;
traceHydrationRuntime.enqueue(traceId, () => hydrateTraceEventsForMessageNow(message));
}
function messagesWithTraceAuthority(source: ChatMessage[]): ChatMessage[] {
return source.map((message) => {
if (message.role !== "agent") return message;
@@ -960,7 +882,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
}
async function refreshProviderOptions(): Promise<boolean> {
const response = await api.providerProfiles.catalog();
const response = await workbenchColadaQueries.fetchProviderProfiles({ minIntervalMs: runtimePolicy.workbenchReadFailureCooldownMs });
providerOptions.value = response.ok ? providerProfileOptionsFromPayload(response.data, providerProfile.value) : defaultProviderProfileOptions(providerProfile.value);
return response.ok === true;
}
@@ -1057,14 +979,16 @@ export const useWorkbenchStore = defineStore("workbench", () => {
function scheduleActiveTraceRestGapFill(traceId: string | null | undefined, reason: string, delayMs = runtimePolicy.workbenchActiveTraceRestGapFillInitialMs): void {
const id = firstNonEmptyString(traceId);
if (!id || typeof window === "undefined") return;
activeTraceRestGapFillRuntime.schedule(composeWorkbenchScopedKey("workbench.active-trace-gap", id), () => refreshActiveTraceFromRest(id, reason), { delayMs, replaceTimer: true, reason });
if (!id) return;
if (typeof window === "undefined") {
void refreshActiveTraceFromRest(id, reason);
return;
}
window.setTimeout(() => { void refreshActiveTraceFromRest(id, reason); }, delayMs);
}
function clearActiveTraceRestGapFill(traceId: string | null | undefined): void {
const id = firstNonEmptyString(traceId);
if (!id || typeof window === "undefined") return;
activeTraceRestGapFillRuntime.cancelTimer(composeWorkbenchScopedKey("workbench.active-trace-gap", id));
void traceId;
}
async function refreshActiveTraceFromRest(traceId: string, reason: string): Promise<void> {
@@ -1241,21 +1165,19 @@ export const useWorkbenchStore = defineStore("workbench", () => {
async function refreshTerminalTraceFromRest(traceId: string, reason: string): Promise<void> {
const id = firstNonEmptyString(traceId);
if (!id) return;
await terminalRealtimeRefreshRuntime.run(composeWorkbenchScopedKey("workbench.terminal-trace-refresh", id), async () => {
const ownerBefore = traceOwnerSessionId(id, null);
if (ownerBefore === activeSessionId.value) recordActivity(reason);
await refreshTurnStatusByTraceId(id);
const ownerSessionId = traceOwnerSessionId(id, turnStatusAuthority.value[id]?.sessionId ?? null) ?? ownerBefore;
if (ownerSessionId) {
scheduleSessionListRefresh(ownerSessionId, runtimePolicy.sessionListTerminalRefreshDelayMs);
await refreshMessageProjectionForTrace(ownerSessionId, id);
}
const ownerMessages = ownerSessionId ? serverState.value.messagesBySessionId[ownerSessionId] ?? [] : messages.value;
const message = latestMessageForTrace(id, ownerMessages);
if (!message && ownerSessionId) {
await refreshRealtimeSessionMessages(ownerSessionId, `${reason}:message-gap`);
}
}, { reason });
const ownerBefore = traceOwnerSessionId(id, null);
if (ownerBefore === activeSessionId.value) recordActivity(reason);
await refreshTurnStatusByTraceId(id);
const ownerSessionId = traceOwnerSessionId(id, turnStatusAuthority.value[id]?.sessionId ?? null) ?? ownerBefore;
if (ownerSessionId) {
scheduleSessionListRefresh(ownerSessionId, runtimePolicy.sessionListTerminalRefreshDelayMs);
await refreshMessageProjectionForTrace(ownerSessionId, id);
}
const ownerMessages = ownerSessionId ? serverState.value.messagesBySessionId[ownerSessionId] ?? [] : messages.value;
const message = latestMessageForTrace(id, ownerMessages);
if (!message && ownerSessionId) {
await refreshRealtimeSessionMessages(ownerSessionId, `${reason}:message-gap`);
}
}
function installRealtimeVisibilityHandler(): void {
@@ -1335,16 +1257,10 @@ export const useWorkbenchStore = defineStore("workbench", () => {
async function refreshRealtimeSessionFromRest(sessionId: string, reason: string): Promise<void> {
const id = normalizeWorkbenchSessionId(sessionId);
if (!id || id !== activeSessionId.value) return;
await realtimeSessionRefreshRuntime.run(composeWorkbenchScopedKey("workbench.realtime.session-refresh", id), async () => {
recordActivity(reason);
const existing = sessions.value.find((item) => item.sessionId === id) ?? null;
const selected = await runWorkbenchReadHydration(
() => loadWorkbenchSession(id, existing),
workbenchReadCooldownKey("session-detail", id),
{ minIntervalMs: runtimePolicy.workbenchSessionDetailMinRefreshMs }
);
if (selected && activeSessionId.value === id) applySelectedSessionDetail(selected, "system");
}, { reason });
recordActivity(reason);
const existing = sessions.value.find((item) => item.sessionId === id) ?? null;
const selected = await loadWorkbenchSession(id, existing);
if (selected && activeSessionId.value === id) applySelectedSessionDetail(selected, "system");
}
function completeTrace(traceId: string, result: AgentChatResultResponse, options: { forceRead?: boolean } = {}): void {
@@ -1689,6 +1605,44 @@ export const useWorkbenchStore = defineStore("workbench", () => {
restartRealtime("session-read-unavailable");
}
async function loadWorkbenchSession(sessionId: string, seed: WorkbenchSessionRecord | null = null): Promise<WorkbenchSessionRecord | null> {
const requestId = normalizeWorkbenchSessionRouteId(sessionId);
if (!requestId) return null;
const normalizedRequestId = normalizeWorkbenchSessionId(requestId);
const eagerMessages = normalizedRequestId ? workbenchColadaQueries.fetchSessionMessages(normalizedRequestId, { limit: 100, force: true, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs }) : null;
const detail = await workbenchColadaQueries.fetchSession(requestId, { force: true, minIntervalMs: runtimePolicy.workbenchSessionDetailMinRefreshMs });
if (!detail.ok) return null;
const detailSession = sessionFromWorkbenchSession(detail.data?.session);
const id = detailSession?.sessionId ?? normalizedRequestId ?? seed?.sessionId;
if (!id) return null;
const messages = eagerMessages && id === normalizedRequestId ? await eagerMessages : await workbenchColadaQueries.fetchSessionMessages(id, { limit: 100, force: true, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs });
const base = detailSession ?? seed;
if (!base) return null;
const page = messages.ok ? messages.data : null;
const pageMessages = Array.isArray(page?.messages) ? await sealRestoredActiveTurnMessages(page.messages.map((message) => normalizeChatMessage(message as ChatMessage))) : base.messages;
return { ...base, sessionId: id, messages: pageMessages, messageCount: page?.total ?? pageMessages?.length ?? base.messageCount };
}
async function sealRestoredActiveTurnMessages(source: ChatMessage[]): Promise<ChatMessage[]> {
const targets = source.filter(messageNeedsRestoredTurnSeal).slice(-3);
if (targets.length === 0) return source;
const patches = new Map<string, Partial<ChatMessage>>();
await Promise.all(targets.map(async (message) => {
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
if (!traceId) return;
const response = await fetchWorkbenchTurnStatus(traceId, false, { force: true });
if (!response.ok || !response.data) return;
const patch = terminalMessagePatchFromTurnResult(message, response.data);
if (patch) patches.set(traceId, patch);
}));
if (patches.size === 0) return source;
return source.map((message) => {
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
const patch = traceId ? patches.get(traceId) : null;
return patch ? { ...message, ...patch } : message;
});
}
function reattachRestoredActiveTrace(): void {
void hydrateTraceEvents(messages.value);
const traceId = activeTraceIdFromMessages(messages.value, turnStatusAuthority.value);
@@ -1754,44 +1708,6 @@ function routeSelectedSessionRecord(sessionId: string | null | undefined, loadin
} as WorkbenchSessionRecord;
}
async function loadWorkbenchSession(sessionId: string, seed: WorkbenchSessionRecord | null = null): Promise<WorkbenchSessionRecord | null> {
const requestId = normalizeWorkbenchSessionRouteId(sessionId);
if (!requestId) return null;
const normalizedRequestId = normalizeWorkbenchSessionId(requestId);
const eagerMessages = normalizedRequestId ? api.workbench.sessionMessages(normalizedRequestId, { limit: 100 }) : null;
const detail = await api.workbench.session(requestId);
if (!detail.ok) return null;
const detailSession = sessionFromWorkbenchSession(detail.data?.session);
const id = detailSession?.sessionId ?? normalizedRequestId ?? seed?.sessionId;
if (!id) return null;
const messages = eagerMessages && id === normalizedRequestId ? await eagerMessages : await api.workbench.sessionMessages(id, { limit: 100 });
const base = detailSession ?? seed;
if (!base) return null;
const page = messages.ok ? messages.data : null;
const pageMessages = Array.isArray(page?.messages) ? await sealRestoredActiveTurnMessages(page.messages.map((message) => normalizeChatMessage(message as ChatMessage))) : base.messages;
return { ...base, sessionId: id, messages: pageMessages, messageCount: page?.total ?? pageMessages?.length ?? base.messageCount };
}
async function sealRestoredActiveTurnMessages(source: ChatMessage[]): Promise<ChatMessage[]> {
const targets = source.filter(messageNeedsRestoredTurnSeal).slice(-3);
if (targets.length === 0) return source;
const patches = new Map<string, Partial<ChatMessage>>();
await Promise.all(targets.map(async (message) => {
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
if (!traceId) return;
const response = await api.workbench.turn(traceId, 8000);
if (!response.ok || !response.data) return;
const patch = terminalMessagePatchFromTurnResult(message, response.data);
if (patch) patches.set(traceId, patch);
}));
if (patches.size === 0) return source;
return source.map((message) => {
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
const patch = traceId ? patches.get(traceId) : null;
return patch ? { ...message, ...patch } : message;
});
}
function messageNeedsRestoredTurnSeal(message: ChatMessage): boolean {
if (message.role !== "agent") return false;
if (messageHasTerminalResponse(message)) return false;
@@ -1,328 +0,0 @@
// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first.
// Responsibility: Workbench REST refresh queue, keyed single-flight, cooldown, scheduled gap-fill and request-storm suppression.
// Mechanical source references:
// - OpenCode util/queue.ts:1-32 AsyncQueue and bounded work.
// - OpenCode runtime.queue.ts:54-113 queue state/control.
// - OpenCode runtime.queue.ts:113-230 drain/send/abort flow.
// - OpenCode stream.transport.ts:83-130 per-turn wait/tick state consumed by recovery scheduling.
import { AsyncQueue } from "@/utils/scheduler/async-queue";
import { createKeyedSingleflight } from "@/utils/scheduler/keyed-singleflight";
import { normalizeErrorDiagnostic } from "@/utils/workbench-error-runtime";
interface ApiResultLike {
ok: boolean;
status: number;
data?: unknown;
error?: string | null;
apiError?: Record<string, unknown> | null;
diagnostic?: Record<string, unknown> | null;
}
export interface WorkbenchReadHydrationRuntimeOptions {
concurrency: number;
failureCooldownMs: number;
now?: () => number;
}
export interface WorkbenchReadHydrationRunOptions {
minIntervalMs?: number;
/** Priority marker only. Runtime budgets still apply to force requests. */
force?: boolean;
shouldCooldownFailure?: (value: unknown) => boolean;
reason?: string | null;
}
type QueuedTask = () => void;
type TimerHandle = ReturnType<typeof setTimeout>;
export interface WorkbenchKeyedTaskRunOptions {
minIntervalMs?: number;
/** Priority marker only. Runtime budgets still apply to force requests. */
force?: boolean;
replace?: boolean;
reason?: string | null;
}
export interface WorkbenchScheduledTaskOptions extends WorkbenchKeyedTaskRunOptions {
delayMs?: number;
replaceTimer?: boolean;
}
export interface WorkbenchScheduledTaskRuntimeOptions {
now?: () => number;
setTimeout?: (callback: () => void, delayMs: number) => TimerHandle;
clearTimeout?: (handle: TimerHandle) => void;
}
export interface WorkbenchTraceHydrationQueueRuntimeOptions {
concurrency: number;
delayMs: number;
setTimeout?: (callback: () => void, delayMs: number) => TimerHandle;
clearTimeout?: (handle: TimerHandle) => void;
}
export class WorkbenchReadHydrationRuntime {
private readonly queue = new AsyncQueue<QueuedTask>();
private readonly singleflight = createKeyedSingleflight<unknown>();
private readonly cooldownUntilByKey = new Map<string, number>();
private readonly lastStartedAtByKey = new Map<string, number>();
private readonly now: () => number;
private readonly failureCooldownMs: number;
constructor(options: WorkbenchReadHydrationRuntimeOptions) {
this.now = options.now ?? Date.now;
this.failureCooldownMs = Math.max(0, Math.trunc(options.failureCooldownMs));
const concurrency = Math.max(1, Math.trunc(options.concurrency));
for (let index = 0; index < concurrency; index += 1) void this.workerLoop();
}
run<T>(task: () => Promise<T>, cooldownKey?: string, options: WorkbenchReadHydrationRunOptions = {}): Promise<T> {
const key = normalizeOptionalKey(cooldownKey);
const execute = () => this.enqueue(() => this.executeTask(task, key, options));
if (!key) return execute();
return this.singleflight.run(key, () => execute(), { reason: options.reason ?? null }) as Promise<T>;
}
private enqueue<T>(task: () => Promise<T>): Promise<T> {
return new Promise<T>((resolve, reject) => {
this.queue.push(() => {
void task().then(resolve, reject);
});
});
}
private async executeTask<T>(task: () => Promise<T>, key: string | null, options: WorkbenchReadHydrationRunOptions): Promise<T> {
if (key) {
const now = this.now();
const cooldownUntilMs = this.cooldownUntilByKey.get(key) ?? 0;
if (cooldownUntilMs > now) return createWorkbenchReadCooldownResult(key, cooldownUntilMs, now) as T;
if (cooldownUntilMs > 0) this.cooldownUntilByKey.delete(key);
const minIntervalMs = Math.max(0, Math.trunc(options.minIntervalMs ?? 0));
const lastStartedAtMs = this.lastStartedAtByKey.get(key) ?? 0;
if (minIntervalMs > 0 && lastStartedAtMs > 0 && now - lastStartedAtMs < minIntervalMs) return createWorkbenchReadThrottledResult(key, lastStartedAtMs + minIntervalMs, now, options.force === true) as T;
if (minIntervalMs > 0) this.lastStartedAtByKey.set(key, now);
}
const value = await task();
if (key && options.shouldCooldownFailure?.(value)) this.cooldownUntilByKey.set(key, this.now() + this.failureCooldownMs);
return value;
}
private async workerLoop(): Promise<void> {
for (;;) {
const run = await this.queue.next();
run();
}
}
}
export class WorkbenchKeyedTaskRuntime {
private readonly singleflight = createKeyedSingleflight<unknown>();
private readonly lastStartedAtByKey = new Map<string, number>();
private readonly now: () => number;
constructor(options: { now?: () => number } = {}) {
this.now = options.now ?? Date.now;
}
run<T>(key: string, task: () => Promise<T>, options: WorkbenchKeyedTaskRunOptions = {}): Promise<T | undefined> {
const normalized = normalizeRequiredKey(key);
const now = this.now();
const minIntervalMs = Math.max(0, Math.trunc(options.minIntervalMs ?? 0));
const lastStartedAtMs = this.lastStartedAtByKey.get(normalized) ?? 0;
if (minIntervalMs > 0 && lastStartedAtMs > 0 && now - lastStartedAtMs < minIntervalMs) return Promise.resolve(undefined);
this.lastStartedAtByKey.set(normalized, now);
return this.singleflight.run(normalized, task, { reason: options.reason ?? null, replace: options.replace }) as Promise<T>;
}
cancel(key: string, reason = "workbench-keyed-task-canceled"): void {
this.singleflight.cancel(normalizeRequiredKey(key), reason);
}
clear(reason = "workbench-keyed-task-cleared"): void {
this.singleflight.clear(reason);
this.lastStartedAtByKey.clear();
}
}
export class WorkbenchScheduledTaskRuntime extends WorkbenchKeyedTaskRuntime {
private readonly timers = new Map<string, TimerHandle>();
private readonly setTimer: (callback: () => void, delayMs: number) => TimerHandle;
private readonly clearTimer: (handle: TimerHandle) => void;
constructor(options: WorkbenchScheduledTaskRuntimeOptions = {}) {
super({ now: options.now });
this.setTimer = options.setTimeout ?? ((callback, delayMs) => setTimeout(callback, delayMs));
this.clearTimer = options.clearTimeout ?? ((handle) => clearTimeout(handle));
}
schedule<T>(key: string, task: () => Promise<T>, options: WorkbenchScheduledTaskOptions = {}): boolean {
const normalized = normalizeRequiredKey(key);
const existing = this.timers.get(normalized);
if (existing) {
if (options.replaceTimer !== true) return false;
this.clearTimer(existing);
this.timers.delete(normalized);
}
const delayMs = Math.max(0, Math.trunc(options.delayMs ?? 0));
const timer = this.setTimer(() => {
this.timers.delete(normalized);
void this.run(normalized, task, options);
}, delayMs);
this.timers.set(normalized, timer);
return true;
}
cancelTimer(key: string): void {
const normalized = normalizeRequiredKey(key);
const timer = this.timers.get(normalized);
if (timer) this.clearTimer(timer);
this.timers.delete(normalized);
}
clearAll(reason = "workbench-scheduled-task-cleared"): void {
for (const timer of this.timers.values()) this.clearTimer(timer);
this.timers.clear();
this.clear(reason);
}
}
export class WorkbenchTraceHydrationQueueRuntime<T> {
private readonly queue = new AsyncQueue<QueuedTask>();
private readonly inFlight = new Set<string>();
private readonly queued = new Set<string>();
private readonly delayMs: number;
private readonly setTimer: (callback: () => void, delayMs: number) => TimerHandle;
constructor(options: WorkbenchTraceHydrationQueueRuntimeOptions) {
const concurrency = Math.max(1, Math.trunc(options.concurrency));
this.delayMs = Math.max(0, Math.trunc(options.delayMs));
this.setTimer = options.setTimeout ?? ((callback, delayMs) => setTimeout(callback, delayMs));
for (let index = 0; index < concurrency; index += 1) void this.workerLoop();
}
run(key: string, task: () => Promise<void>, options: { onAlreadyRunning?: () => void } = {}): Promise<boolean> {
const normalized = normalizeRequiredKey(key);
if (this.inFlight.has(normalized)) {
options.onAlreadyRunning?.();
return Promise.resolve(false);
}
this.inFlight.add(normalized);
return task().then(() => true, (error) => {
throw error;
}).finally(() => {
this.inFlight.delete(normalized);
});
}
enqueue(key: string, task: () => Promise<void>): boolean {
const normalized = normalizeRequiredKey(key);
if (this.inFlight.has(normalized) || this.queued.has(normalized)) return false;
this.queued.add(normalized);
this.queue.push(() => {
void this.run(normalized, task).finally(() => {
this.queued.delete(normalized);
});
});
return true;
}
private async workerLoop(): Promise<void> {
for (;;) {
const run = await this.queue.next();
run();
if (this.delayMs > 0) await new Promise<void>((resolve) => this.setTimer(resolve, this.delayMs));
}
}
}
export function createWorkbenchReadHydrationRuntime(options: WorkbenchReadHydrationRuntimeOptions): WorkbenchReadHydrationRuntime {
return new WorkbenchReadHydrationRuntime(options);
}
export function createWorkbenchKeyedTaskRuntime(options: { now?: () => number } = {}): WorkbenchKeyedTaskRuntime {
return new WorkbenchKeyedTaskRuntime(options);
}
export function createWorkbenchScheduledTaskRuntime(options: WorkbenchScheduledTaskRuntimeOptions = {}): WorkbenchScheduledTaskRuntime {
return new WorkbenchScheduledTaskRuntime(options);
}
export function createWorkbenchTraceHydrationQueueRuntime<T>(options: WorkbenchTraceHydrationQueueRuntimeOptions): WorkbenchTraceHydrationQueueRuntime<T> {
return new WorkbenchTraceHydrationQueueRuntime<T>(options);
}
export function shouldCooldownWorkbenchReadFailure(value: unknown): boolean {
const result = value as ApiResultLike | null;
if (!result || typeof result !== "object" || typeof result.status !== "number") return false;
const diagnostic = normalizeErrorDiagnostic(result.diagnostic, result.apiError?.diagnostic);
const code = firstStringOrNumber(result.apiError?.code, diagnostic?.code);
const category = firstNonEmptyString(result.apiError?.category, diagnostic?.category);
const source = firstNonEmptyString(result.apiError?.source, diagnostic?.source);
if (result.status === 0 && source === "browser" && code === "browser_network_error" && category === "network") return true;
if (result.status === 503) return code === "projection_store_unavailable" || code === "workbench_read_model_store_unavailable";
return false;
}
function createWorkbenchReadCooldownResult(key: string, cooldownUntilMs: number, now: number): ApiResultLike {
const retryAfterMs = Math.max(0, cooldownUntilMs - now);
const message = `Workbench read hydration is cooling down after a transient failure (${key}); retry after ${retryAfterMs}ms.`;
const diagnostic = diagnosticEnvelope("workbench_read_hydration_cooldown", "network", message, key, retryAfterMs);
return { ok: false, status: 503, data: null, error: message, apiError: { ...diagnostic, message, diagnostic }, diagnostic };
}
function createWorkbenchReadThrottledResult(key: string, nextAllowedAtMs: number, now: number, forceRequested = false): ApiResultLike {
const retryAfterMs = Math.max(0, nextAllowedAtMs - now);
const message = `Workbench read hydration is throttled (${key}); retry after ${retryAfterMs}ms.`;
const diagnostic = diagnosticEnvelope("workbench_read_hydration_throttled", "throttle", message, key, retryAfterMs);
diagnostic.forceRequested = forceRequested;
diagnostic.forceBudgetBypass = false;
return { ok: false, status: 0, data: null, error: message, apiError: { ...diagnostic, message, diagnostic }, diagnostic };
}
function diagnosticEnvelope(code: string, category: string, message: string, key: string, retryAfterMs: number): Record<string, unknown> {
return {
contractVersion: "hwlab-error-diagnostic-v1",
code,
category,
source: "workbench-web",
layer: "workbench-refresh-runtime",
message,
retryable: true,
recoveryAction: "suppress-duplicate-refresh",
rootCause: code,
scopedKey: key,
retryAfterMs,
valuesPrinted: false,
valuesRedacted: true
};
}
function normalizeOptionalKey(value: string | null | undefined): string | null {
const text = String(value ?? "").trim();
return text || null;
}
function normalizeRequiredKey(value: string): string {
const text = String(value ?? "").trim();
if (!text) throw new Error("workbench runtime key is required");
return text;
}
function firstNonEmptyString(...values: unknown[]): string | null {
for (const value of values) {
if (typeof value !== "string") continue;
const text = value.trim();
if (text) return text;
}
return null;
}
function firstStringOrNumber(...values: unknown[]): string | number | null {
for (const value of values) {
if (typeof value === "number" && Number.isFinite(value)) return value;
if (typeof value === "string" && value.trim()) return value.trim();
}
return null;
}