diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 74ca058..c9b5820 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -13,7 +13,7 @@ jobs: runs-on: ubuntu-latest timeout-minutes: 15 steps: - - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 - uses: pnpm/action-setup@0ebf47130e4866e96fce0953f49152a61190b271 # v6.0.9 with: version: 10.34.4 diff --git a/.prettierignore b/.prettierignore index aa883ed..cd32ef1 100644 --- a/.prettierignore +++ b/.prettierignore @@ -1,5 +1,6 @@ dist/ node_modules/ +pnpm-lock.yaml vendor/ TXODDS_STOPPAGE_MASTER_PLAN.md data/private/ diff --git a/README.md b/README.md index a32d35f..518ed88 100644 --- a/README.md +++ b/README.md @@ -397,7 +397,11 @@ ignored private capture directory, reconnects with bounded backoff, emits stream-health inputs into the same governor, and persists derived decision receipts and Certified Reopen sidecars separately. It also refreshes the fixture catalog every five minutes so new knockout fixtures become eligible without a -restart. +restart. Its private status file retains only sanitized failure classes, HTTP +status codes and timestamps for the latest stream or fixture-refresh failure; +successful recovery clears those diagnostics. Console snapshots are emitted +immediately for health-state changes and otherwise rate-limited, while the +status file remains fresh every second. For a container host, the compiled worker runs without development dependencies: diff --git a/package.json b/package.json index 9f3989e..826a227 100644 --- a/package.json +++ b/package.json @@ -56,12 +56,12 @@ }, "dependencies": { "@stoppage/sdk": "workspace:*", - "@fastify/static": "9.3.0", + "@fastify/static": "10.1.2", "@noble/hashes": "2.2.0", "fastify": "5.10.0", - "lucide-react": "1.24.0", - "react": "19.2.7", - "react-dom": "19.2.7", + "lucide-react": "1.27.0", + "react": "19.2.8", + "react-dom": "19.2.8", "tweetnacl": "1.0.3", "zod": "4.4.3" }, @@ -72,12 +72,12 @@ "@types/node": "26.1.1", "@types/react": "19.2.17", "@types/react-dom": "19.2.3", - "@vitejs/plugin-react": "6.0.3", - "concurrently": "10.0.3", - "prettier": "3.9.5", + "@vitejs/plugin-react": "6.0.4", + "concurrently": "10.0.4", + "prettier": "3.9.6", "tsx": "4.23.0", "typescript": "5.9.3", - "vite": "8.1.4", + "vite": "8.1.5", "vitest": "4.1.10" }, "pnpm": { diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index a96a0a9..cfc51e8 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -11,8 +11,8 @@ importers: .: dependencies: "@fastify/static": - specifier: 9.3.0 - version: 9.3.0 + specifier: 10.1.2 + version: 10.1.2 "@noble/hashes": specifier: 2.2.0 version: 2.2.0 @@ -23,14 +23,14 @@ importers: specifier: 5.10.0 version: 5.10.0 lucide-react: - specifier: 1.24.0 - version: 1.24.0(react@19.2.7) + specifier: 1.27.0 + version: 1.27.0(react@19.2.8) react: - specifier: 19.2.7 - version: 19.2.7 + specifier: 19.2.8 + version: 19.2.8 react-dom: - specifier: 19.2.7 - version: 19.2.7(react@19.2.7) + specifier: 19.2.8 + version: 19.2.8(react@19.2.8) tweetnacl: specifier: 1.0.3 version: 1.0.3 @@ -57,14 +57,14 @@ importers: specifier: 19.2.3 version: 19.2.3(@types/react@19.2.17) "@vitejs/plugin-react": - specifier: 6.0.3 - version: 6.0.3(vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)) + specifier: 6.0.4 + version: 6.0.4(vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)) concurrently: - specifier: 10.0.3 - version: 10.0.3 + specifier: 10.0.4 + version: 10.0.4 prettier: - specifier: 3.9.5 - version: 3.9.5 + specifier: 3.9.6 + version: 3.9.6 tsx: specifier: 4.23.0 version: 4.23.0 @@ -72,11 +72,11 @@ importers: specifier: 5.9.3 version: 5.9.3 vite: - specifier: 8.1.4 - version: 8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) + specifier: 8.1.5 + version: 8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) vitest: specifier: 4.1.10 - version: 4.1.10(@types/node@26.1.1)(vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)) + version: 4.1.10(@types/node@26.1.1)(vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)) packages/sdk: dependencies: @@ -398,10 +398,10 @@ packages: integrity: sha512-PxcYtKLbQ8Z+yApiqjK8FwxIwvEj38k2OiLc17u8dkJSlmfi2wHHPaSnaoqBPQqtvF8YVsDgDpP2snDCfFrpfw==, } - "@fastify/forwarded@3.0.1": + "@fastify/forwarded@3.0.2": resolution: { - integrity: sha512-JqDochHFqXs3C3Ml3gOY58zM7OqO9ENqPo0UqAjAjH8L01fRZqwX9iLeX34//kiJubF7r2ZQHtBRU36vONbLlw==, + integrity: sha512-NE8HgKLgYejV9lDpqkEFaDKMLYelJBVfHekhB0UKvX0ghagXRJqg68feg8er1NPXxG4N9i6vPxzt8E+3wHfcmA==, } "@fastify/merge-json-schemas@0.2.1": @@ -422,10 +422,10 @@ packages: integrity: sha512-TMYeQLCBSy2TOFmV95hQWkiTYgC/SEx7vMdV+wnZVX4tt8VBLKzmH8vV9OzJehV0+XBfg+WxPMt5wp+JBUKsVw==, } - "@fastify/static@9.3.0": + "@fastify/static@10.1.2": resolution: { - integrity: sha512-9YMYRpCOtMBrqKYWcqiw7ykOrn4D0jogHpJrFS0KGeSuOwzKMM5/mjj7B0CFLVoQ6htqKYw//Zs7APn9DBq05w==, + integrity: sha512-G/g18cG9tLutT/OVyN1AIsHIl9L1UwmJ+S3dkyhVpplIx0nEMicd7RGQ+uJLyhKKF4a3tTcQydccn3Mop1fX+Q==, } "@jridgewell/sourcemap-codec@1.5.5": @@ -764,10 +764,10 @@ packages: integrity: sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==, } - "@vitejs/plugin-react@6.0.3": + "@vitejs/plugin-react@6.0.4": resolution: { - integrity: sha512-vmFvco5/QuC2f9Oj+wTk0+9XeDFkHxSamwZKYc7MxYwKICfvUvlMhqKI0VuICPltGqh1neqBKDvO4kes1ya8vg==, + integrity: sha512-XcCQz0TBpBgljhj0gMuuDj49i6Ytqh5q1osT/Gp5uAVJUCTWxyskk/l1jwYYiu2xcNHHipdMz40EGfM1VdamVg==, } engines: { node: ^20.19.0 || >=22.12.0 } peerDependencies: @@ -888,10 +888,10 @@ packages: } engines: { node: ">=8.0.0" } - avvio@9.2.0: + avvio@9.3.0: resolution: { - integrity: sha512-2t/sy01ArdHHE0vRH5Hsay+RtCZt3dLPji7W7/MMOCEgze5b7SNDC4j5H6FnVgPkI1MTNFGzHdHrVXDDl7QSSQ==, + integrity: sha512-g2tQ7LE7oOSqDfwEm3M+ZCMTJc7KiZCdJ4UwyZJb5ckTKyYu50OYmvv0mCFXPuYXoM4zkSt8zM9XQ9KCvxA74A==, } balanced-match@4.0.4: @@ -925,12 +925,12 @@ packages: integrity: sha512-CLCsZGIBCFnPtkNnieW/a8wmreDmfUtjU2m9yHrzPXIlNbqVs0AQrSatSG6vdNYUqdc83tkQi2eHfF98ubzQLA==, } - brace-expansion@5.0.7: + brace-expansion@5.0.8: resolution: { - integrity: sha512-7oFy703dxfY3/NLxC1fh2SUCQ0H9rmAY+5EpDVfXjUTTs+HEwR2nYaqLv+GWcTsumwxPfiz6CzCNkwXwBUwqCA==, + integrity: sha512-JZyDyq3D4AUifKTPOB7DELf6XsB3WdPuNxCtob1vFXPsSXhdAiHBWJ/tJ8HAc9aH84BK+5JFZLNkJKx3G9kzQg==, } - engines: { node: 18 || 20 || >=22 } + engines: { node: 20 || >=22 } bs58@4.0.1: resolution: @@ -999,18 +999,18 @@ packages: integrity: sha512-GpVkmM8vF2vQUkj2LvZmD35JxeJOLCwJ9cUkugyk2nuhbv3+mJvpLYYt+0+USMxE+oj+ey/lJEnhZw75x/OMcQ==, } - concurrently@10.0.3: + concurrently@10.0.4: resolution: { - integrity: sha512-hc3LH4UaKWd/bbyDK/IGVa4RB6PtQ3CUYwtrkzqHn+wIG3Hr5fhpRlk0L/gCa8ZE1L/Ufj50Zho69cI5w8SQBA==, + integrity: sha512-trZql+7l/0+WRAsAnEdctr4+iiOS6ZrViI6H8QWcCF9MFS/LT0dKpe8vluB1to6it+OxSI4VospFTIFMW8DJRw==, } engines: { node: ">=22" } hasBin: true - content-disposition@1.1.0: + content-disposition@2.0.1: resolution: { - integrity: sha512-5jRCH9Z/+DRP7rkvY83B+yGIGX96OYdJmzngqnw2SBSxqCFPd0w2km3s5iawpGX8krnwSGmF0FW5Nhr0Hfai3g==, + integrity: sha512-e+H0ZXHSWYrENhQzw1LPuP4oF5MzVKmDU6d3hxlvaPEYLLg62MxtQNPRx4SYSuYJSBUgnQIG4HIN2tEtNv7Dog==, } engines: { node: ">=18" } @@ -1174,16 +1174,16 @@ packages: integrity: sha512-wpYMUmFu5f00Sm0cj2pfivpmawLZ0NKdviQ4w9zJeR8JVtOpOxHmLaJuj0vxvGqMJQWyP/COUkF75/57OKyRag==, } - fast-uri@3.1.3: + fast-uri@3.1.4: resolution: { - integrity: sha512-i70LwGWUduXqzicKXWshooq+sWL1K3WUU5rKZNG/0i3a1OSoX3HqhH5WbWwTmqWfor4urUakGPiRQcleRZTwOg==, + integrity: sha512-8JnbkQ4juDyvYs4mgFGQqg4yCYtFDtUtmp2QIQq11ZZe5CFQ5wcqm1rqDgAh/QdMySuBnPzMUiJUNZG5N/AiQw==, } - fast-uri@4.1.0: + fast-uri@4.1.1: resolution: { - integrity: sha512-ZodJ2cRiLVWGi9IgPb3mbgSqM4CD3LexCHkuv0FfBXHJI1ADfucTD06m6clO2Cy5RZYsw/SiCVl/dyrFI/SYWA==, + integrity: sha512-YPOs1zD5TG2+EZt+r88LwF6mclA7TPkpwMP7ZN3TO2HiHS8TXvq7QA/17iJsV9dubcLo/f8eEYqMBruyQV21hQ==, } fastify-plugin@6.0.0: @@ -1216,10 +1216,10 @@ packages: picomatch: optional: true - find-my-way@9.6.0: + find-my-way@9.7.0: resolution: { - integrity: sha512-Zf4Xve4RymLl7NgaavNebZ01joJ8MfVerOG43wy7SHLO+r+K0C6d/SE0BiR7AV5V1VOCFlOP7ecdo+I4qmiHrQ==, + integrity: sha512-f2JHn75x2JlwUwLenZypgczR7YWMb/uO9BvUXtus+JMgkbIkLADd38cI4EiV+OQqrGo1Zlq6V8wnqMJ8e62wUQ==, } engines: { node: ">=20" } @@ -1449,10 +1449,10 @@ packages: } engines: { node: 20 || >=22 } - lucide-react@1.24.0: + lucide-react@1.27.0: resolution: { - integrity: sha512-YT6mBD8lGKkg4nM39enlm94/sfJIiW0YKUT60fBy4YK8tai31ylg1VhGNWxkpSKHo9UagfnZqwIff3HTDQwXeA==, + integrity: sha512-rJicGl/3Fly/E0rOH1YmPZ6e49JCnKknh1ox1vpHnkfjujAkKA6sqUZvH3MTAaXXjgexyUwgNwTJzTtYuAFYJw==, } peerDependencies: react: ^16.5.1 || ^17.0.0 || ^18.0.0 || ^19.0.0 @@ -1471,10 +1471,10 @@ packages: engines: { node: ">=10.0.0" } hasBin: true - minimatch@10.2.5: + minimatch@10.2.6: resolution: { - integrity: sha512-MULkVLfKGYDFYejP07QOurDLLQpcjk7Fw+7jXS2R2czRQzR56yHRveU5NDJEOviH+hETZKSkIk5c+T23GjFUMg==, + integrity: sha512-vpLQEs+VLCr1nU0BXS07maYoFwlDAH0gngQuuttxIwutDFEMHq2blX+8vpgxDdK3J1PwjCJiep77OitTZ4Ll1A==, } engines: { node: 18 || 20 || >=22 } @@ -1491,10 +1491,10 @@ packages: integrity: sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==, } - nanoid@3.3.15: + nanoid@3.3.16: resolution: { - integrity: sha512-y7Wygv/7mEOvxTuEQDB8StXdMRBWf1kR/tlhAzBRUFkB2jfcLOAxO/SHmOO2zgz1pVgK29/kyupn059/bCHdjA==, + integrity: sha512-bzlKTyNJ7+LdGIIwy8ijFpIqEQIvafahV7eYykJ8Cvh42EdJeODoJ6gUJXpQJvej1BddH8OqTXZNE/KfbWAu8Q==, } engines: { node: ^10 || ^12 || ^13.7 || ^14 || >=15.0.1 } hasBin: true @@ -1599,17 +1599,17 @@ packages: engines: { node: ">=18" } hasBin: true - postcss@8.5.16: + postcss@8.5.24: resolution: { - integrity: sha512-vuwillviilfKZsg0VGj5R/YwwcHx4SLsIOI/7K6mQkWx+l5cUHTjj5g0AasTBcyXsbfTgrwsUNmVUb5xVwyPwg==, + integrity: sha512-8RyVklq0owXUTa4xlpzu4l9AaVKIdQvAcOHZWaMh98HgySsUtxRVf/chRe3dsSLqb6i40BzGRzEUddRaI+9TSw==, } engines: { node: ^10 || ^12 || >=14 } - prettier@3.9.5: + prettier@3.9.6: resolution: { - integrity: sha512-/FVl766LpUfB5vXgCYOYa0MeV/441Ia99AeICQIQFTY/Nw0roZwULcXpku5i1/m5kt/baz+s4Zogspd839HSMg==, + integrity: sha512-OpN0zzVdiaiAhxpuuj5efpIS4sY9j7bY6uR5mnj5yPzGkdkjNKSJeUThPb60Jw29QuAZgA4o+/iB49kFiaBX6g==, } engines: { node: ">=14" } hasBin: true @@ -1632,18 +1632,18 @@ packages: integrity: sha512-tYC1Q1hgyRuHgloV/YXs2w15unPVh8qfu/qCTfhTYamaw7fyhumKa2yGpdSo87vY32rIclj+4fWYQXUMs9EHvg==, } - react-dom@19.2.7: + react-dom@19.2.8: resolution: { - integrity: sha512-t0BRVXvbiE/o20Hfw669rLbMCDWtYZLvmJigy2f0MxsXF+71pxhR3xOkspmsO8h3ZlNzyibAmtCa3l4lYKk6gQ==, + integrity: sha512-rVprimfGBG3DR+Tq0IQG2DT5PxKth1WIGDmj5yPmlzr4YBe7uyE+Du4oVqTDXZSHGGGXRtTJEGSSePyQCMBglQ==, } peerDependencies: - react: ^19.2.7 + react: ^19.2.8 - react@19.2.7: + react@19.2.8: resolution: { - integrity: sha512-HNe9WslTbXmFK8o8cmwgAeJFSBvt1bPdHCVKtaaV+WlAN36mpT4hcRpwbf3fY56ar2oIXzsBpOAiIRHAdY0OlQ==, + integrity: sha512-PWaYA1L/q9u2u7xYQi+Y3L3Yfnie7XyLeaJICV1MGD6LprsBxcAqGjYyr0eY3p+QdsA+x/Irkt4Qif8D63+Sbw==, } engines: { node: ">=0.10.0" } @@ -1759,10 +1759,10 @@ packages: integrity: sha512-E5LDX7Wrp85Kil5bhZv46j8jOeboKq5JMmYM3gVGdGH8xFpPWXUMsNrlODCrkoxMEeNi/XZIwuRvY4XNwYMJpw==, } - shell-quote@1.8.4: + shell-quote@1.9.0: resolution: { - integrity: sha512-VsC6n6vz1ihYYyZZwX7YZSF5l5x36ca17OC+a69h94YqB7X6XLwf+5MOgynYir2SLFUbl8gIYvBo8K8RoNQ6bQ==, + integrity: sha512-Iov+JwFv/2HcTpcwNMKd8+IWNb8tboQJNQTkAY/LLVK7gGH9jy+LGkVqPxfekHl+yMmiqXszdGWXgkfml7hjqA==, } engines: { node: ">= 0.4" } @@ -1985,10 +1985,10 @@ packages: } hasBin: true - vite@8.1.4: + vite@8.1.5: resolution: { - integrity: sha512-bTT9PsdWO+MQMNG9ZXIP/qM9wGh37DFxTV/sPq9cFpHr3w4jkgef032PkAL9jAqhk3Nz8NQw3O8n6/xFkqO4QQ==, + integrity: sha512-7ULLwsCdYx/nRyrpiEwvqb5TFHrMVZyBt+rg/OAXT7rgj/z+DtTDyKFeLAdDkubDVDKD8jOsndmy7m55XcfUsw==, } engines: { node: ^20.19.0 || >=22.12.0 } hasBin: true @@ -2291,7 +2291,7 @@ snapshots: dependencies: ajv: 8.20.0 ajv-formats: 3.0.1(ajv@8.20.0) - fast-uri: 3.1.3 + fast-uri: 3.1.4 "@fastify/error@4.2.0": {} @@ -2299,7 +2299,7 @@ snapshots: dependencies: fast-json-stringify: 7.0.1 - "@fastify/forwarded@3.0.1": {} + "@fastify/forwarded@3.0.2": {} "@fastify/merge-json-schemas@0.2.1": dependencies: @@ -2307,7 +2307,7 @@ snapshots: "@fastify/proxy-addr@5.1.0": dependencies: - "@fastify/forwarded": 3.0.1 + "@fastify/forwarded": 3.0.2 ipaddr.js: 2.4.0 "@fastify/send@4.1.0": @@ -2318,11 +2318,12 @@ snapshots: http-errors: 2.0.1 mime: 3.0.0 - "@fastify/static@9.3.0": + "@fastify/static@10.1.2": dependencies: "@fastify/accept-negotiator": 2.0.1 + "@fastify/error": 4.2.0 "@fastify/send": 4.1.0 - content-disposition: 1.1.0 + content-disposition: 2.0.1 fastify-plugin: 6.0.0 fastq: 1.20.1 glob: 13.0.6 @@ -2497,10 +2498,10 @@ snapshots: dependencies: "@types/node": 26.1.1 - "@vitejs/plugin-react@6.0.3(vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0))": + "@vitejs/plugin-react@6.0.4(vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0))": dependencies: "@rolldown/pluginutils": 1.0.1 - vite: 8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) + vite: 8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) "@vitest/expect@4.1.10": dependencies: @@ -2511,13 +2512,13 @@ snapshots: chai: 6.2.2 tinyrainbow: 3.1.0 - "@vitest/mocker@4.1.10(vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0))": + "@vitest/mocker@4.1.10(vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0))": dependencies: "@vitest/spy": 4.1.10 estree-walker: 3.0.3 magic-string: 0.30.21 optionalDependencies: - vite: 8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) + vite: 8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) "@vitest/pretty-format@4.1.10": dependencies: @@ -2556,7 +2557,7 @@ snapshots: ajv@8.20.0: dependencies: fast-deep-equal: 3.1.3 - fast-uri: 3.1.3 + fast-uri: 3.1.4 json-schema-traverse: 1.0.0 require-from-string: 2.0.2 @@ -2568,7 +2569,7 @@ snapshots: atomic-sleep@1.0.0: {} - avvio@9.2.0: + avvio@9.3.0: dependencies: "@fastify/error": 4.2.0 fastq: 1.20.1 @@ -2589,7 +2590,7 @@ snapshots: bs58: 4.0.1 text-encoding-utf-8: 1.0.2 - brace-expansion@5.0.7: + brace-expansion@5.0.8: dependencies: balanced-match: 4.0.4 @@ -2625,16 +2626,16 @@ snapshots: commander@2.20.3: {} - concurrently@10.0.3: + concurrently@10.0.4: dependencies: chalk: 5.6.2 rxjs: 7.8.2 - shell-quote: 1.8.4 + shell-quote: 1.9.0 supports-color: 10.2.2 tree-kill: 1.2.2 yargs: 18.0.0 - content-disposition@1.1.0: {} + content-disposition@2.0.1: {} convert-source-map@2.0.0: {} @@ -2720,7 +2721,7 @@ snapshots: "@fastify/merge-json-schemas": 0.2.1 ajv: 8.20.0 ajv-formats: 3.0.1(ajv@8.20.0) - fast-uri: 4.1.0 + fast-uri: 4.1.1 json-schema-ref-resolver: 3.0.0 rfdc: 1.4.1 @@ -2730,9 +2731,9 @@ snapshots: fast-stable-stringify@1.0.0: {} - fast-uri@3.1.3: {} + fast-uri@3.1.4: {} - fast-uri@4.1.0: {} + fast-uri@4.1.1: {} fastify-plugin@6.0.0: {} @@ -2743,9 +2744,9 @@ snapshots: "@fastify/fast-json-stringify-compiler": 5.1.0 "@fastify/proxy-addr": 5.1.0 abstract-logging: 2.0.1 - avvio: 9.2.0 + avvio: 9.3.0 fast-json-stringify: 7.0.1 - find-my-way: 9.6.0 + find-my-way: 9.7.0 light-my-request: 6.6.0 pino: 10.3.1 process-warning: 5.0.0 @@ -2762,7 +2763,7 @@ snapshots: optionalDependencies: picomatch: 4.0.5 - find-my-way@9.6.0: + find-my-way@9.7.0: dependencies: fast-deep-equal: 3.1.3 fast-querystring: 1.1.2 @@ -2780,7 +2781,7 @@ snapshots: glob@13.0.6: dependencies: - minimatch: 10.2.5 + minimatch: 10.2.6 minipass: 7.1.3 path-scurry: 2.0.2 @@ -2889,9 +2890,9 @@ snapshots: lru-cache@11.5.2: {} - lucide-react@1.24.0(react@19.2.7): + lucide-react@1.27.0(react@19.2.8): dependencies: - react: 19.2.7 + react: 19.2.8 magic-string@0.30.21: dependencies: @@ -2899,15 +2900,15 @@ snapshots: mime@3.0.0: {} - minimatch@10.2.5: + minimatch@10.2.6: dependencies: - brace-expansion: 5.0.7 + brace-expansion: 5.0.8 minipass@7.1.3: {} ms@2.1.3: {} - nanoid@3.3.15: {} + nanoid@3.3.16: {} node-fetch@2.7.0: dependencies: @@ -2961,13 +2962,13 @@ snapshots: optionalDependencies: fsevents: 2.3.2 - postcss@8.5.16: + postcss@8.5.24: dependencies: - nanoid: 3.3.15 + nanoid: 3.3.16 picocolors: 1.1.1 source-map-js: 1.2.1 - prettier@3.9.5: {} + prettier@3.9.6: {} process-warning@4.0.1: {} @@ -2975,12 +2976,12 @@ snapshots: quick-format-unescaped@4.0.4: {} - react-dom@19.2.7(react@19.2.7): + react-dom@19.2.8(react@19.2.8): dependencies: - react: 19.2.7 + react: 19.2.8 scheduler: 0.27.0 - react@19.2.7: {} + react@19.2.8: {} real-require@0.2.0: {} @@ -3050,7 +3051,7 @@ snapshots: setprototypeof@1.2.0: {} - shell-quote@1.8.4: {} + shell-quote@1.9.0: {} siginfo@2.0.0: {} @@ -3140,11 +3141,11 @@ snapshots: uuid@14.0.1: {} - vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0): + vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0): dependencies: lightningcss: 1.32.0 picomatch: 4.0.5 - postcss: 8.5.16 + postcss: 8.5.24 rolldown: 1.1.5 tinyglobby: 0.2.17 optionalDependencies: @@ -3153,10 +3154,10 @@ snapshots: fsevents: 2.3.3 tsx: 4.23.0 - vitest@4.1.10(@types/node@26.1.1)(vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)): + vitest@4.1.10(@types/node@26.1.1)(vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)): dependencies: "@vitest/expect": 4.1.10 - "@vitest/mocker": 4.1.10(vite@8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)) + "@vitest/mocker": 4.1.10(vite@8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0)) "@vitest/pretty-format": 4.1.10 "@vitest/runner": 4.1.10 "@vitest/snapshot": 4.1.10 @@ -3173,7 +3174,7 @@ snapshots: tinyexec: 1.2.4 tinyglobby: 0.2.17 tinyrainbow: 3.1.0 - vite: 8.1.4(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) + vite: 8.1.5(@types/node@26.1.1)(esbuild@0.28.1)(tsx@4.23.0) why-is-node-running: 2.3.0 optionalDependencies: "@types/node": 26.1.1 diff --git a/src/live/txline-live-worker.test.ts b/src/live/txline-live-worker.test.ts index 01e7cf1..ed1bd8a 100644 --- a/src/live/txline-live-worker.test.ts +++ b/src/live/txline-live-worker.test.ts @@ -202,6 +202,77 @@ describe("TxLineLiveWorker", () => { } }); + it("records fixture refresh failures and clears them after recovery", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + let fixtureRequests = 0; + let now = 1_000; + const client = { + fetchFixtures: async () => { + fixtureRequests += 1; + if (fixtureRequests === 2) { + throw new Error( + "TxLINE fixture snapshot failed with HTTP 503: private", + ); + } + return []; + }, + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + fixtureRefreshMs: 1_000, + now: () => now, + callbacks: { onInput: () => undefined }, + capture: async () => undefined, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + + now = 2_000; + await vi.advanceTimersByTimeAsync(1_000); + expect(worker.status()).toMatchObject({ + fixtureRefreshes: 1, + fixtureRefreshFailures: 1, + lastFixtureRefreshFailure: { + reason: "fixture-refresh-error:Error:HTTP_503", + at: 2_000, + }, + }); + expect(JSON.stringify(worker.status())).not.toContain("private"); + + now = 3_000; + await vi.advanceTimersByTimeAsync(1_000); + expect(worker.status()).toMatchObject({ + fixtureRefreshes: 2, + fixtureRefreshFailures: 1, + lastFixtureRefreshAt: 3_000, + lastFixtureRefreshFailure: null, + }); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + } finally { + vi.useRealTimers(); + } + }); + it("publishes a stopped status when startup fails", async () => { const statuses: boolean[] = []; const client = { @@ -226,9 +297,468 @@ describe("TxLineLiveWorker", () => { "fixture snapshot unavailable", ); expect(worker.status().running).toBe(false); + expect(worker.status()).toMatchObject({ + fixtureRefreshFailures: 1, + lastFixtureRefreshFailure: { + reason: "fixture-refresh-error:Error", + }, + }); expect(statuses.at(-1)).toBe(false); }); + it("records stream failures and clears them when the stream reconnects", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + let now = 1_000; + let oddsConnections = 0; + const statuses: ReturnType[] = []; + const client = { + fetchFixtures: async () => [], + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + oddsConnections += 1; + if (oddsConnections === 1) { + throw new Error("TxLINE odds failed with HTTP 502: private"); + } + await callbacks.onOpen(); + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + reconnectBaseMs: 100, + now: () => now, + callbacks: { + onInput: () => undefined, + onStatus: (status) => { + statuses.push(status); + }, + }, + capture: async () => undefined, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + expect(worker.status()).toMatchObject({ + reconnects: { odds: 1, scores: 0 }, + lastStreamFailure: { + odds: { + reason: "stream-error:Error:HTTP_502", + at: 1_000, + }, + scores: null, + }, + }); + expect(JSON.stringify(worker.status())).not.toContain("private"); + expect( + statuses.some( + (status) => + status.lastStreamFailure?.odds?.reason === + "stream-error:Error:HTTP_502", + ), + ).toBe(true); + + now = 1_100; + await vi.advanceTimersByTimeAsync(100); + expect(oddsConnections).toBe(2); + expect(worker.status().lastStreamFailure?.odds).toBeNull(); + expect(worker.status().streamHealth.odds).toBe(true); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + } finally { + vi.useRealTimers(); + } + }); + + it("reconnects and recovers after a heartbeat timeout", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + let now = 1_000; + let oddsConnections = 0; + let scoreConnections = 0; + const client = { + fetchFixtures: async () => [], + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + oddsConnections += 1; + await callbacks.onOpen(); + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + scoreConnections += 1; + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + heartbeatTimeoutMs: 500, + reconnectBaseMs: 100, + now: () => now, + callbacks: { onInput: () => undefined }, + capture: async () => undefined, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + + now = 2_000; + await vi.advanceTimersByTimeAsync(1_000); + expect(worker.status().lastStreamFailure).toEqual({ + odds: { reason: "heartbeat-timeout", at: 2_000 }, + scores: { reason: "heartbeat-timeout", at: 2_000 }, + }); + expect(worker.status()).toMatchObject({ + reconnects: { odds: 1, scores: 1 }, + streamHealth: { odds: false, scores: false }, + }); + + now = 2_100; + await vi.advanceTimersByTimeAsync(100); + expect({ oddsConnections, scoreConnections }).toEqual({ + oddsConnections: 2, + scoreConnections: 2, + }); + expect(worker.status()).toMatchObject({ + reconnects: { odds: 1, scores: 1 }, + streamHealth: { odds: true, scores: true }, + lastStreamFailure: { odds: null, scores: null }, + }); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + } finally { + vi.useRealTimers(); + } + }); + + it("ignores data that resumes after its stream attempt is aborted", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + let now = 1_000; + let oddsConnections = 0; + let releaseCapture: (() => void) | undefined; + const captureBlocked = new Promise((resolve) => { + releaseCapture = resolve; + }); + const client = { + fetchFixtures: async () => [ + { + FixtureId: 77, + StartTime: 1_000, + Participant1: "Northbridge", + Participant2: "Eastport", + Participant1IsHome: true, + }, + ], + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + oddsConnections += 1; + await callbacks.onOpen(); + if (oddsConnections === 1) { + await callbacks.onRaw?.(oddsPayload(), "odds:delayed"); + await callbacks.onData(oddsPayload()); + return; + } + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + heartbeatTimeoutMs: 500, + reconnectBaseMs: 100, + now: () => now, + callbacks: { onInput: () => undefined }, + capture: async () => captureBlocked, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + + now = 2_000; + await vi.advanceTimersByTimeAsync(1_000); + expect(worker.status()).toMatchObject({ + streamHealth: { odds: false }, + lastStreamFailure: { + odds: { reason: "heartbeat-timeout", at: 2_000 }, + }, + oddsMessages: 0, + normalizedOdds: 0, + }); + + releaseCapture?.(); + await vi.advanceTimersByTimeAsync(0); + expect(worker.status()).toMatchObject({ + streamHealth: { odds: false }, + lastStreamFailure: { + odds: { reason: "heartbeat-timeout", at: 2_000 }, + }, + oddsMessages: 0, + normalizedOdds: 0, + }); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + } finally { + vi.useRealTimers(); + } + }); + + it("removes completed delay listeners before scheduling the next tick", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + const addListener = vi.spyOn(controller.signal, "addEventListener"); + const removeListener = vi.spyOn(controller.signal, "removeEventListener"); + const client = { + fetchFixtures: async () => [], + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + callbacks: { onInput: () => undefined }, + capture: async () => undefined, + }); + const abortListenerBalance = () => + addListener.mock.calls.filter(([type]) => type === "abort").length - + removeListener.mock.calls.filter(([type]) => type === "abort").length; + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + const initialBalance = abortListenerBalance(); + + await vi.advanceTimersByTimeAsync(3_000); + expect(abortListenerBalance()).toBe(initialBalance); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + expect(abortListenerBalance()).toBe(0); + } finally { + vi.useRealTimers(); + } + }); + + it("times out an unopened stream attempt and reconnects", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + let now = 1_000; + let oddsConnections = 0; + const statuses: ReturnType[] = []; + const client = { + fetchFixtures: async () => [], + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + oddsConnections += 1; + if (oddsConnections > 1) await callbacks.onOpen(); + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + streamConnectTimeoutMs: 100, + reconnectBaseMs: 50, + now: () => now, + callbacks: { + onInput: () => undefined, + onStatus: (status) => { + statuses.push(status); + }, + }, + capture: async () => undefined, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + expect(oddsConnections).toBe(1); + + now = 1_100; + await vi.advanceTimersByTimeAsync(100); + expect(worker.status()).toMatchObject({ + reconnects: { odds: 1, scores: 0 }, + streamHealth: { odds: false, scores: true }, + lastStreamFailure: { + odds: { reason: "stream-connect-timeout", at: 1_100 }, + scores: null, + }, + }); + expect( + statuses.some( + (status) => + status.lastStreamFailure?.odds?.reason === "stream-connect-timeout", + ), + ).toBe(true); + + await vi.advanceTimersByTimeAsync(49); + expect(oddsConnections).toBe(1); + now = 1_150; + await vi.advanceTimersByTimeAsync(1); + expect(oddsConnections).toBe(2); + expect(worker.status().streamHealth.odds).toBe(true); + expect(worker.status().lastStreamFailure?.odds).toBeNull(); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + expect(worker.status().reconnects.odds).toBe(1); + } finally { + vi.useRealTimers(); + } + }); + + it("does not classify parent shutdown as a connection timeout", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + const client = { + fetchFixtures: async () => [], + streamOdds: async ( + _callbacks: FakeCallbacks, + signal: AbortSignal, + ) => untilAborted(signal), + streamScores: async ( + _callbacks: FakeCallbacks, + signal: AbortSignal, + ) => untilAborted(signal), + }; + const worker = new TxLineLiveWorker({ + client, + streamConnectTimeoutMs: 100, + callbacks: { onInput: () => undefined }, + capture: async () => undefined, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + await vi.advanceTimersByTimeAsync(100); + + expect(worker.status()).toMatchObject({ + running: false, + reconnects: { odds: 0, scores: 0 }, + lastStreamFailure: { odds: null, scores: null }, + }); + } finally { + vi.useRealTimers(); + } + }); + + it("retries an AbortError that is not caused by parent shutdown", async () => { + vi.useFakeTimers(); + const controller = new AbortController(); + let scoreConnections = 0; + const client = { + fetchFixtures: async () => [], + streamOdds: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + await callbacks.onOpen(); + await untilAborted(signal); + }, + streamScores: async ( + callbacks: FakeCallbacks, + signal: AbortSignal, + ) => { + scoreConnections += 1; + if (scoreConnections === 1) { + const error = new Error("shared guest session request aborted"); + error.name = "AbortError"; + throw error; + } + await callbacks.onOpen(); + await untilAborted(signal); + }, + }; + const worker = new TxLineLiveWorker({ + client, + reconnectBaseMs: 100, + callbacks: { onInput: () => undefined }, + capture: async () => undefined, + }); + + try { + const run = worker.run(controller.signal); + await vi.advanceTimersByTimeAsync(0); + expect(worker.status()).toMatchObject({ + reconnects: { odds: 0, scores: 1 }, + streamHealth: { odds: true, scores: false }, + lastStreamFailure: { + odds: null, + scores: { reason: "stream-error:AbortError" }, + }, + }); + + await vi.advanceTimersByTimeAsync(100); + expect(scoreConnections).toBe(2); + expect(worker.status()).toMatchObject({ + streamHealth: { odds: true, scores: true }, + lastStreamFailure: { odds: null, scores: null }, + }); + + controller.abort(); + await vi.advanceTimersByTimeAsync(0); + await run; + } finally { + vi.useRealTimers(); + } + }); + it("resets reconnect backoff after a healthy heartbeat", async () => { vi.useFakeTimers(); const controller = new AbortController(); @@ -296,3 +826,17 @@ function untilAborted(signal: AbortSignal) { signal.addEventListener("abort", () => resolve(), { once: true }), ); } + +function oddsPayload(): OddsPayload { + return { + FixtureId: 77, + MessageId: "delayed", + Ts: 1_000, + Bookmaker: "StablePrice", + BookmakerId: 1, + SuperOddsType: "1X2", + InRunning: true, + PriceNames: ["Northbridge", "Draw", "Eastport"], + Pct: ["50.000", "25.000", "25.000"], + }; +} diff --git a/src/live/txline-live-worker.ts b/src/live/txline-live-worker.ts index de089ee..b784dfc 100644 --- a/src/live/txline-live-worker.ts +++ b/src/live/txline-live-worker.ts @@ -30,21 +30,35 @@ interface LiveWorkerOptions { client: LiveTxLineClient; callbacks: LiveWorkerCallbacks; heartbeatTimeoutMs?: number; + streamConnectTimeoutMs?: number; reconnectBaseMs?: number; fixtureRefreshMs?: number; now?: () => number; capture?: (name: string, value: unknown) => Promise; } +type StreamAbortReason = "heartbeat-timeout" | "stream-connect-timeout"; + +interface ActiveStreamAttempt { + abortReason: StreamAbortReason | null; + failureRecorded: boolean; + controller: AbortController; +} + export class TxLineLiveWorker { readonly #client: LiveTxLineClient; readonly #callbacks: LiveWorkerCallbacks; readonly #heartbeatTimeoutMs: number; + readonly #streamConnectTimeoutMs: number; readonly #reconnectBaseMs: number; readonly #fixtureRefreshMs: number; readonly #now: () => number; readonly #capture: (name: string, value: unknown) => Promise; readonly #participants = new Map(); + readonly #activeStreamAttempts: Record< + "odds" | "scores", + ActiveStreamAttempt | null + > = { odds: null, scores: null }; #status: LiveWorkerStatus = emptyStatus(); #lastHealthEmission = { odds: true, scores: true }; @@ -52,6 +66,7 @@ export class TxLineLiveWorker { this.#client = options.client; this.#callbacks = options.callbacks; this.#heartbeatTimeoutMs = options.heartbeatTimeoutMs ?? 45_000; + this.#streamConnectTimeoutMs = options.streamConnectTimeoutMs ?? 30_000; this.#reconnectBaseMs = options.reconnectBaseMs ?? 1_000; this.#fixtureRefreshMs = options.fixtureRefreshMs ?? 300_000; this.#now = options.now ?? Date.now; @@ -104,11 +119,23 @@ export class TxLineLiveWorker { } async #refreshFixtures(signal: AbortSignal) { - const fixtures = await this.#client.fetchFixtures({ signal }); - this.#loadParticipants(fixtures); - this.#status.fixturesLoaded = this.#participants.size; - this.#status.fixtureRefreshes += 1; - this.#status.lastFixtureRefreshAt = this.#now(); + try { + const fixtures = await this.#client.fetchFixtures({ signal }); + this.#loadParticipants(fixtures); + this.#status.fixturesLoaded = this.#participants.size; + this.#status.fixtureRefreshes += 1; + this.#status.lastFixtureRefreshAt = this.#now(); + this.#status.lastFixtureRefreshFailure = null; + } catch (error) { + if (!signal.aborted && (error as Error).name !== "AbortError") { + this.#status.fixtureRefreshFailures += 1; + this.#status.lastFixtureRefreshFailure = { + reason: `fixture-refresh-error:${safeErrorReason(error)}`, + at: this.#now(), + }; + } + throw error; + } } async #refreshFixtureLoop(signal: AbortSignal) { @@ -119,7 +146,6 @@ export class TxLineLiveWorker { await this.#refreshFixtures(signal); } catch (error) { if (signal.aborted || (error as Error).name === "AbortError") break; - this.#status.fixtureRefreshFailures += 1; } await this.#publishStatus(); } @@ -131,18 +157,53 @@ export class TxLineLiveWorker { const controller = new AbortController(); const abort = () => controller.abort(); signal.addEventListener("abort", abort, { once: true }); + let connected = false; + const activeAttempt: ActiveStreamAttempt = { + abortReason: null, + failureRecorded: false, + controller, + }; + this.#activeStreamAttempts[name] = activeAttempt; + const connectTimeout = setTimeout(() => { + if (connected || signal.aborted) return; + activeAttempt.abortReason = "stream-connect-timeout"; + controller.abort(); + }, this.#streamConnectTimeoutMs); + const markConnected = () => { + if ( + controller.signal.aborted || + this.#activeStreamAttempts[name] !== activeAttempt + ) { + return false; + } + connected = true; + clearTimeout(connectTimeout); + return true; + }; try { if (name === "odds") { await this.#client.streamOdds( { - onOpen: () => this.#touch(name), + onOpen: () => { + if (!markConnected()) return; + return this.#touch(name); + }, onHeartbeat: () => { + if (!markConnected()) return; attempt = 0; return this.#touch(name); }, - onRaw: (payload, eventId) => - this.#captureRaw(name, payload, eventId), + onRaw: (payload, eventId) => { + if ( + controller.signal.aborted || + this.#activeStreamAttempts[name] !== activeAttempt + ) { + return; + } + return this.#captureRaw(name, payload, eventId); + }, onData: (payload) => { + if (!markConnected()) return; attempt = 0; return this.#handleOdds(payload); }, @@ -152,14 +213,26 @@ export class TxLineLiveWorker { } else { await this.#client.streamScores( { - onOpen: () => this.#touch(name), + onOpen: () => { + if (!markConnected()) return; + return this.#touch(name); + }, onHeartbeat: () => { + if (!markConnected()) return; attempt = 0; return this.#touch(name); }, - onRaw: (payload, eventId) => - this.#captureRaw(name, payload, eventId), + onRaw: (payload, eventId) => { + if ( + controller.signal.aborted || + this.#activeStreamAttempts[name] !== activeAttempt + ) { + return; + } + return this.#captureRaw(name, payload, eventId); + }, onData: (payload) => { + if (!markConnected()) return; attempt = 0; return this.#handleScore(payload); }, @@ -169,19 +242,32 @@ export class TxLineLiveWorker { } if (!signal.aborted) { this.#status.reconnects[name] += 1; - await this.#setHealth(name, false, "stream-closed"); + const reason = activeAttempt.abortReason ?? "stream-closed"; + if (!activeAttempt.failureRecorded) { + this.#recordStreamFailure(name, reason); + activeAttempt.failureRecorded = true; + await this.#setHealth(name, false, reason); + } + await this.#publishStatus(); attempt += 1; } } catch (error) { - if (signal.aborted || (error as Error).name === "AbortError") return; + if (signal.aborted) return; this.#status.reconnects[name] += 1; - await this.#setHealth( - name, - false, - `stream-error:${(error as Error).name}`, - ); + const reason = + activeAttempt.abortReason ?? `stream-error:${safeErrorReason(error)}`; + if (!activeAttempt.failureRecorded) { + this.#recordStreamFailure(name, reason); + activeAttempt.failureRecorded = true; + await this.#setHealth(name, false, reason); + } + await this.#publishStatus(); attempt += 1; } finally { + clearTimeout(connectTimeout); + if (this.#activeStreamAttempts[name] === activeAttempt) { + this.#activeStreamAttempts[name] = null; + } signal.removeEventListener("abort", abort); } @@ -241,8 +327,13 @@ export class TxLineLiveWorker { } async #touch(name: "odds" | "scores", timestamp = this.#now()) { + const recovered = + !this.#status.streamHealth[name] || + Boolean(this.#status.lastStreamFailure?.[name]); this.#status.lastMessageAt[name] = timestamp; + this.#status.lastStreamFailure![name] = null; await this.#setHealth(name, true); + if (recovered) await this.#publishStatus(); } async #monitorHealth(signal: AbortSignal) { @@ -253,10 +344,23 @@ export class TxLineLiveWorker { for (const name of ["odds", "scores"] as const) { const lastMessage = this.#status.lastMessageAt[name]; if ( + this.#status.streamHealth[name] && lastMessage !== null && now - lastMessage > this.#heartbeatTimeoutMs ) { - await this.#setHealth(name, false, "heartbeat-timeout"); + const activeAttempt = this.#activeStreamAttempts[name]; + if (activeAttempt && !activeAttempt.controller.signal.aborted) { + activeAttempt.abortReason = "heartbeat-timeout"; + activeAttempt.failureRecorded = true; + this.#recordStreamFailure(name, "heartbeat-timeout"); + const healthUpdate = this.#setHealth( + name, + false, + "heartbeat-timeout", + ); + activeAttempt.controller.abort(); + await healthUpdate; + } } } await this.#callbacks.onInput({ kind: "tick", observedTs: now }); @@ -277,6 +381,13 @@ export class TxLineLiveWorker { }); } + #recordStreamFailure(name: "odds" | "scores", reason: string) { + this.#status.lastStreamFailure![name] = { + reason, + at: this.#now(), + }; + } + async #publishStatus() { await this.#callbacks.onStatus?.(this.status()); } @@ -295,12 +406,26 @@ function emptyStatus(): LiveWorkerStatus { fixtureRefreshes: 0, fixtureRefreshFailures: 0, lastFixtureRefreshAt: null, + lastFixtureRefreshFailure: null, streamHealth: { odds: false, scores: false }, lastMessageAt: { odds: null, scores: null }, + lastStreamFailure: { odds: null, scores: null }, startedAt: null, }; } +function safeErrorReason(error: unknown) { + const name = + error instanceof Error && /^[A-Za-z][A-Za-z0-9]{0,63}$/.test(error.name) + ? error.name + : "UnknownError"; + const status = + error instanceof Error + ? /\bHTTP\s+(\d{3})\b/.exec(error.message)?.[1] + : null; + return status ? `${name}:HTTP_${status}` : name; +} + function captureName() { return `live-${new Date().toISOString().slice(0, 10)}.jsonl`; } @@ -308,14 +433,15 @@ function captureName() { function abortableDelay(milliseconds: number, signal: AbortSignal) { if (signal.aborted) return Promise.resolve(); return new Promise((resolve) => { - const timeout = setTimeout(resolve, milliseconds); - signal.addEventListener( - "abort", - () => { - clearTimeout(timeout); - resolve(); - }, - { once: true }, - ); + let settled = false; + const settle = () => { + if (settled) return; + settled = true; + clearTimeout(timeout); + signal.removeEventListener("abort", settle); + resolve(); + }; + const timeout = setTimeout(settle, milliseconds); + signal.addEventListener("abort", settle, { once: true }); }); } diff --git a/src/live/types.ts b/src/live/types.ts index 8cdc79a..a24e6a2 100644 --- a/src/live/types.ts +++ b/src/live/types.ts @@ -12,11 +12,18 @@ export interface LiveWorkerStatus { fixtureRefreshes: number; fixtureRefreshFailures: number; lastFixtureRefreshAt: number | null; + lastFixtureRefreshFailure?: WorkerFailure | null; streamHealth: Record; lastMessageAt: Record; + lastStreamFailure?: Record; startedAt: string | null; } +export interface WorkerFailure { + reason: string; + at: number; +} + export interface LiveWorkerCallbacks { onInput: (input: GovernorInput) => void | Promise; onStatus?: (status: LiveWorkerStatus) => void | Promise; diff --git a/src/worker.ts b/src/worker.ts index a896201..7ae2069 100644 --- a/src/worker.ts +++ b/src/worker.ts @@ -25,6 +25,8 @@ if (!config.txlineApiToken) { throw new Error("TXLINE_API_TOKEN is required for the live worker"); } +const STATUS_COUNTER_LOG_INTERVAL_MS = 300_000; +const STATUS_HEARTBEAT_LOG_INTERVAL_MS = 900_000; const governor = new QuoteGovernor(); const controller = new AbortController(); const durationMs = readOptionalDuration(); @@ -33,6 +35,8 @@ const client = new TxLineClient({ apiToken: config.txlineApiToken, }); let lastStatusLogAt = 0; +let lastStatusStateSignature: string | null = null; +let lastStatusCounterSignature: string | null = null; let inputQueue = Promise.resolve(); const executionContexts = new LiveExecutionContextTracker(); const liveDecisionTape = config.liveDecisionTapeEnabled @@ -82,9 +86,39 @@ const worker = new TxLineLiveWorker({ updatedAt: new Date().toISOString(), }; await writeRuntimeState("worker-status.json", snapshot); - if (Date.now() - lastStatusLogAt >= 30_000 || !status.running) { + const now = Date.now(); + const stateSignature = JSON.stringify({ + running: status.running, + fixturesLoaded: status.fixturesLoaded, + fixtureRefreshFailure: status.lastFixtureRefreshFailure?.reason ?? null, + streamHealth: status.streamHealth, + streamFailures: { + odds: status.lastStreamFailure?.odds?.reason ?? null, + scores: status.lastStreamFailure?.scores?.reason ?? null, + }, + }); + const counterSignature = JSON.stringify({ + oddsMessages: status.oddsMessages, + scoreMessages: status.scoreMessages, + normalizedOdds: status.normalizedOdds, + normalizedEvents: status.normalizedEvents, + skippedOdds: status.skippedOdds, + reconnects: status.reconnects, + fixtureRefreshes: status.fixtureRefreshes, + fixtureRefreshFailures: status.fixtureRefreshFailures, + }); + const stateChanged = stateSignature !== lastStatusStateSignature; + const countersChanged = counterSignature !== lastStatusCounterSignature; + const counterUpdateDue = + countersChanged && + now - lastStatusLogAt >= STATUS_COUNTER_LOG_INTERVAL_MS; + const heartbeatDue = + now - lastStatusLogAt >= STATUS_HEARTBEAT_LOG_INTERVAL_MS; + if (stateChanged || counterUpdateDue || heartbeatDue || !status.running) { console.log(JSON.stringify(snapshot)); - lastStatusLogAt = Date.now(); + lastStatusLogAt = now; + lastStatusStateSignature = stateSignature; + lastStatusCounterSignature = counterSignature; } }, },