fix(observability): emit generation:end exactly once on stream finalize - #989
Conversation
|
The latest updates on your projects. Learn more about Vercel for GitHub.
|
|
Important Review skippedAuto incremental reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
WalkthroughThis PR implements per-stream concurrency-safe deduplication for Changes
Sequence DiagramsequenceDiagram
participant Client
participant Orchestration as Orchestration<br/>(neurolink)
participant ALS as AsyncLocalStorage<br/>Context
participant Provider as Native Provider<br/>(Stream)
Client->>Orchestration: call stream()
Orchestration->>ALS: set StreamGenerationEndContext<br/>(providerEmitted: false)
Orchestration->>Provider: initiate stream
Provider->>Provider: emit tokens/events
Provider->>Provider: stream completes
Provider->>ALS: call markStreamProviderEmittedGenerationEnd()
ALS->>ALS: set providerEmitted: true
Provider->>Provider: emit generation:end
Provider-->>Client: generation:end event
Orchestration->>ALS: check providerEmitted flag
alt providerEmitted is false
Orchestration->>Client: emit generation:end
else providerEmitted is true
Orchestration-->>Client: skip emission
end
Orchestration->>Client: return result
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
✅ Single Commit Policy - COMPLIANTStatus: Policy requirements met • 1 commit • Valid format • Ready for merge 📊 View validation details📝 Commit Details
✅ Validation Results
🤖 Automated validation by NeuroLink Single Commit Enforcement |
🤖 AI Review & Build Compliance ✅Status: AI analysis complete • Build rules validated • Ready for review 📊 View detailed analysis results🛡️ Analysis Complete
📋 Ready for Merge When
🤖 AI analysis complete - check individual code comments for specific feedback |
Documentation Validation Results🚀 Documentation validation passed!
📦 Build artifact uploaded successfully. Ready for deployment preview. Commit: |
There was a problem hiding this comment.
Pull request overview
This PR fixes an observability gap where sdk.stream() never emitted the generation:end event, preventing cost/audit/alerting listeners from receiving end-of-generation data for streaming calls.
Changes:
- Emit
generation:endfrom therunStandardStreamRequeststream finalizer (and hoistresolvedUsageso it’s available at finalize time). - Add a continuous test script to count
generation:endemissions forgenerate()vsstream()across real providers. - Add small test helpers and documentation capturing the Curator Issue #4 investigation and fix.
Reviewed changes
Copilot reviewed 4 out of 4 changed files in this pull request and generated 4 comments.
| File | Description |
|---|---|
src/lib/neurolink.ts |
Adds a generation:end emit in the stream generator finally block and hoists resolvedUsage for the final payload. |
test/helpers/envGuard.ts |
Adds env-var skip helper and provider-credential/network error classification helper for real-provider test scripts. |
test/continuous-test-suite-issue-04-generation-end-dedup.ts |
Adds a real-provider script that asserts exactly one generation:end emission for both generate() and stream(). |
docs/curator-feedback-fixes/issue-04-stream-generation-end.md |
Documents the reported issue, root cause, fix approach, and verification steps. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| // The previous implementation only fired `stream:complete`, leaving | ||
| // any subscriber to `generation:end` with zero events. | ||
| try { | ||
| const finalProvider = | ||
| metadata.fallbackProvider ?? providerName ?? "unknown"; | ||
| const finalModel = | ||
| metadata.fallbackModel ?? | ||
| streamModel ?? | ||
| enhancedOptions.model ?? | ||
| "unknown"; | ||
| const finalFinishReason = streamError | ||
| ? "error" | ||
| : (streamState.finishReason ?? "stop"); | ||
| self.emitter.emit("generation:end", { | ||
| provider: finalProvider, | ||
| model: finalModel, | ||
| responseTime: Date.now() - streamStartTime, | ||
| toolsUsed: streamState.toolCalls?.map((t) => t.toolName), | ||
| timestamp: Date.now(), | ||
| result: { | ||
| content: accumulatedContent, | ||
| usage: resolvedUsage, | ||
| model: finalModel, | ||
| provider: finalProvider, | ||
| finishReason: finalFinishReason, | ||
| }, | ||
| prompt: | ||
| enhancedOptions.input?.text || | ||
| (enhancedOptions as Record<string, unknown>).prompt, | ||
| temperature: enhancedOptions.temperature, | ||
| maxTokens: enhancedOptions.maxTokens, | ||
| success: !streamError, | ||
| error: streamError | ||
| ? streamError instanceof Error | ||
| ? streamError.message | ||
| : String(streamError) | ||
| : undefined, | ||
| pipelineAHandled: true, | ||
| }); |
| self.emitter.emit("generation:end", { | ||
| provider: finalProvider, | ||
| model: finalModel, | ||
| responseTime: Date.now() - streamStartTime, |
| prompt: | ||
| enhancedOptions.input?.text || | ||
| (enhancedOptions as Record<string, unknown>).prompt, |
| `generation:end` exactly once with the final stream state. A | ||
| `generationEndEmitted` boolean guards against any future double-emit. The | ||
| event payload mirrors the shape used by `generate()` so listeners receive | ||
| a consistent contract. | ||
|
|
||
| `src/lib/neurolink.ts` lines around 6500 and 6610: hoist `resolvedUsage` | ||
| and `generationEndEmitted` to the generator's scope, populate | ||
| `resolvedUsage` after the success path resolves it, then emit in | ||
| `finally`. |
d1b205a to
db8d093
Compare
|
Force-pushed
Final diff: Suite re-verified post-rebase: 6 PASS (vertex/google-ai-studio/litellm × generate+stream), 1 SKIP (OpenAI quota), 1 FAIL pre-existing OpenAI tool-injection env issue unrelated to this fix. @coderabbitai full review |
|
🧠 Learnings used✅ Actions performedFull review triggered. |
🤖 AI Review & Build Compliance ✅Status: AI analysis complete • Build rules validated • Ready for review 📊 View detailed analysis results🛡️ Analysis Complete
📋 Ready for Merge When
🤖 AI analysis complete - check individual code comments for specific feedback |
|
Force-pushed addressing reviewer Finding #1 (double-emit on native Google streams): The native Gemini 3 paths in Fix: register a counter listener at the top of
Suite re-verified: 6 PASS (vertex/google-ai-studio/litellm × generate+stream count=1 each). @coderabbitai full review |
db8d093 to
771e44e
Compare
🤖 AI Review & Build Compliance ✅Status: AI analysis complete • Build rules validated • Ready for review 📊 View detailed analysis results🛡️ Analysis Complete
📋 Ready for Merge When
🤖 AI analysis complete - check individual code comments for specific feedback |
🧩 Analysis chain
🏁 Script executed: Length of output: 2202 🏁 Script executed: Length of output: 8339 🏁 Script executed: Length of output: 15401 🏁 Script executed: Length of output: 2205
Overall assessmentThe counter-based dedup ( One new concern surfaces from the shared-emitter architecture — details below. Finding
|
771e44e to
7e095fb
Compare
|
Force-pushed addressing reviewer Finding #1 follow-up (concurrency-safe dedup): The previous fix used a counter-listener on the shared SDK emitter and would mis-attribute a concurrent unrelated stream's New approach: AsyncLocalStorage scoped per-stream:
Suite re-verified: 6 PASS (vertex/google-ai-studio/litellm × generate+stream count=1). @coderabbitai full review |
🤖 AI Review & Build Compliance ✅Status: AI analysis complete • Build rules validated • Ready for review 📊 View detailed analysis results🛡️ Analysis Complete
📋 Ready for Merge When
🤖 AI analysis complete - check individual code comments for specific feedback |
🧩 Analysis chain🏁 Script executed: Length of output: 1501 🏁 Script executed: Length of output: 2008 Good — the ALS scaffolding and call sites look structurally correct. Let me check two remaining surface areas before writing the full review. 🏁 Script executed: Length of output: 12969
Overall assessmentThe AsyncLocalStorage approach correctly solves Finding One new gap surfaces: Amazon Bedrock and Ollama stream-path native emits are not instrumented. Finding
|
| File | Line | Site | Marked? |
|---|---|---|---|
googleVertex.ts |
2345 | stream generation:end |
✅ |
googleAiStudio.ts |
1046 | stream success generation:end |
✅ |
googleAiStudio.ts |
1082 | stream failure generation:end |
✅ |
amazonBedrock.ts |
1670 | stream generation:end |
❌ |
ollama.ts |
1123 | stream generation:end |
❌ |
If amazonBedrock.ts:1670 and ollama.ts:1123 are called from within the async chain wrapped by streamGenerationEndContext.run(), the dedupContext.providerEmitted flag will remain false, the finally block will fire the orchestration emit, and the listener receives two generation:end events for those providers.
Fix (mirrors what was done for Google providers):
// amazonBedrock.ts — before line 1670
markStreamProviderEmittedGenerationEnd();
streamEmitter.emit("generation:end", { … });
// ollama.ts — before line 1123
markStreamProviderEmittedGenerationEnd();
ollamaEmitter.emit("generation:end", { … });If Bedrock/Ollama stream paths are routed outside runStandardStreamRequest and never enter the ALS scope, the mark call is a safe no-op (getStore() returns undefined) and causes no harm, so it's worth adding defensively regardless.
Correctness of the ALS implementation ✅
Propagation: streamGenerationEndContext.run(dedupContext, buildProcessedStream) invokes the factory synchronously inside the run() scope. The returned async generator inherits the ALS store through its async continuation chain — any await or yield within the generator, including those triggered by native provider code, will have access to the correct dedupContext. This is a well-established Node.js ALS behaviour.
Ordering: markStreamProviderEmittedGenerationEnd() is called before emit("generation:end", …) at every instrumented site — dedupContext.providerEmitted is therefore true before the finally block's check runs. ✅
Generate-path isolation: googleVertex.ts:2674 and googleAiStudio.ts:1335 (generate paths) do not call markStreamProviderEmittedGenerationEnd(). Since generate paths never enter the ALS scope, getStore() returns undefined and any accidental call would be a no-op anyway. ✅
Concurrent isolation: Two simultaneous sdk.stream() calls produce two independent dedupContext objects in two independent ALS scopes — they cannot observe or mutate each other's flag. Finding #2 is fully resolved. ✅
Non-blocking carry-overs
- Finding
#3(payload shape): The orchestration emit still includesmodelat the top level whilegenerate()'s emit nests it insideresult{}. Still a non-breaking asymmetry; worth a follow-up normalisation. - The
try/catcharound the orchestration emit (Line ~6724) and theemitter.offbefore the conditional check remain sound. ✅
Bottom line: The ALS design is the right fix and is correctly implemented for the Google native paths. Please add markStreamProviderEmittedGenerationEnd() at amazonBedrock.ts:1670 and ollama.ts:1123 before merging to close the remaining double-emit gap on those providers. Happy to confirm once those two sites are instrumented.
🧠 Learnings used
Learnt from: murdore
Repo: juspay/neurolink PR: 0
File: :0-0
Timestamp: 2026-04-17T04:30:25.521Z
Learning: In neurolink (src/lib/neurolink.ts), the `pipelineAHandled: true` flag on `generation:end` events is set ONLY by `finalizeGenerateRequestResult` (the SDK-level wrapper). Native providers (amazonBedrock.ts, ollama.ts, googleAiStudio.ts, googleVertex.ts) emit their own `generation:end` events WITHOUT this flag. The `initializeMetricsListeners()` early-return guard `if (data.pipelineAHandled) return` correctly prevents duplicates: it skips the SDK-level event while allowing native provider events through to Pipeline B.
✅ Actions performed
Full review triggered.
There was a problem hiding this comment.
Actionable comments posted: 4
🧹 Nitpick comments (2)
test/continuous-test-suite-issue-04-generation-end-dedup.ts (2)
159-166: Add a concurrent-stream scenario to actually exercise the ALS isolation.The whole point of switching to
AsyncLocalStorage(over the instance-level counter approach) was per-stream isolation across concurrent streams on the sameNeuroLinkinstance. The current loop is strictly sequential (generate, 1s, stream, 1s, next provider), so it would also pass with the prior, non-concurrency-safe counter implementation — i.e., it cannot catch a regression of the very issue the ALS rewrite was introduced to prevent.Consider adding a scenario that runs two
sdk.stream()calls in parallel on the samesdkinstance (for a provider whose native path emitsgeneration:end, e.g.vertexorgoogle-ai-studio) and asserts each call observes exactly onegeneration:endfor its own correlation id (e.g. by including a unique sentinel in the prompt or by partitioning events via the per-callresult/event metadata).🧪 Sketch of a concurrent-stream check
async function countEmissionsForConcurrentStreams( target: ProviderTarget, ): Promise<void> { const testName = `stream(x2) / ${target.provider} — concurrent generation:end isolation`; const skip = skipIfEnvMissing(...target.envVars); if (skip) { record(testName, "SKIP", skip); return; } const sdk = new NeuroLink(); const events: unknown[] = []; sdk.getEventEmitter().on("generation:end", (e) => events.push(e)); try { const model = target.modelEnv ? process.env[target.modelEnv] : undefined; const runOne = async () => { const r = await sdk.stream({ provider: target.provider as never, ...(model && { model }), input: { text: "Reply with the single word: hello" }, maxTokens: 32, disableTools: true, } as never); for await (const _ of r.stream) { /* drain */ } }; await Promise.all([runOne(), runOne()]); await new Promise((r) => setTimeout(r, 500)); if (events.length === 2) { record(testName, "PASS", `expected 2, got ${events.length}`); } else { record(testName, "FAIL", `expected 2, got ${events.length}`); } } catch (err) { const msg = err instanceof Error ? err.message : String(err); isExpectedProviderError(msg) ? record(testName, "SKIP", msg.slice(0, 120)) : record(testName, "FAIL", `unexpected error: ${msg.slice(0, 200)}`); } finally { await sdk.shutdown?.().catch(() => {}); } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/continuous-test-suite-issue-04-generation-end-dedup.ts` around lines 159 - 166, The current test loop (calling countEmissionsForGenerate and countEmissionsForStream sequentially over TARGETS) doesn't exercise AsyncLocalStorage concurrency; add a new test function (e.g., countEmissionsForConcurrentStreams) that creates a single NeuroLink instance and runs two sdk.stream() calls in parallel, listens for "generation:end" on sdk.getEventEmitter(), drains both streams, waits briefly, and asserts you saw exactly two generation:end events for that instance; then invoke this new function for appropriate targets (those with native generation:end like vertex or google-ai-studio) in main alongside the existing checks and ensure you handle skips, provider errors, and sdk.shutdown similar to the other helpers (use TARGETS, sdk.stream, countEmissionsForGenerate/countEmissionsForStream references to locate insertion points).
77-157: DRY the two near-identical scenario runners.
countEmissionsForGenerateandcountEmissionsForStreamdiffer only in thesdk.generatevssdk.streamcall (and the chunk drain). A small helper that takes arunCall(sdk, target) => Promise<{ provider; model }>and a label would remove ~60 lines of duplicated SKIP/PASS/FAIL/error/shutdown plumbing and make it trivial to add the concurrent-stream variant suggested above.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/continuous-test-suite-issue-04-generation-end-dedup.ts` around lines 77 - 157, The two functions countEmissionsForGenerate and countEmissionsForStream duplicate the same SKIP/PASS/FAIL/error/shutdown/event-listening plumbing; refactor by extracting a helper (e.g., runEmissionTest or runCall) that accepts (target, label, runner) where runner is an async function called with the sdk that returns { provider, model } (for stream the runner should also drain r.stream and wait the 250ms grace period before returning). Move the shared event subscription (sdk.getEventEmitter().on("generation:end", ...)), skip check, try/catch error handling (isExpectedProviderError), record(...) calls and finally sdk.shutdown?.() into the helper, and replace countEmissionsForGenerate/countEmissionsForStream with two tiny callers that pass a runner invoking sdk.generate(...) or sdk.stream(...) respectively.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/lib/neurolink.ts`:
- Around line 6622-6629: The stream path resolves token usage into resolvedUsage
but does not attach streamAnalytics to the synthetic "generation:end" result, so
listeners of generate() won't see stream cost data; update the code that emits
the synthetic generation:end result (the place that assigns resolvedUsage from
streamUsage/streamAnalytics) to also set result.analytics = result.analytics ??
{}; and copy the resolved streamAnalytics (or the resolved object containing
tokenUsage/cost) into result.analytics (e.g., result.analytics.cost or
result.analytics.tokenUsage as appropriate) so the emitted result mirrors the
non-stream path; apply the same change in the other stream branch around lines
where streamAnalytics is resolved (also referenced in the other block handling
stream results).
- Around line 6774-6778: The SDK span currently always sets
ATTR.GEN_AI_FINISH_REASON to "stop" for successful streams, which hides real
finish reasons; update the logic in the block that sets streamSpan.setAttribute
(around streamSpan, ATTR.GEN_AI_FINISH_REASON) to prefer
streamState.finishReason when available, otherwise fall back to setting "error"
when metadata.error or streamError is truthy, and finally "stop" as the last
fallback; ensure you still coerce to a string and handle undefined/null values.
- Around line 6828-6831: The AsyncLocalStorage context is lost because
streamGenerationEndContext.run wraps the generator creation
(buildProcessedStream) instead of the async-iterator itself; change to run the
context around the iterator so the store is active during iteration.
Specifically, call buildProcessedStream() to get the async iterator, then wrap
that iterator with a context-bound wrapper created by
streamGenerationEndContext.run(dedupContext, () => iterator) (or equivalent) so
each next()/throw()/return() invocation runs inside streamGenerationEndContext;
ensure processStreamResult still iterates the wrapped iterator and that
markStreamProviderEmittedGenerationEnd() can read/write
dedupContext.providerEmitted during iteration.
In `@test/continuous-test-suite-issue-04-generation-end-dedup.ts`:
- Around line 139-146: Replace the fragile fixed 250ms grace sleep with a
deadline-based poll: repeatedly check the events array (events.length) in a
short interval (e.g., 50ms) until events.length >= 1 or a deadline (e.g.,
2000ms) is reached, then once you observe events.length >= 1 wait an additional
stability window (~250ms) to ensure no late duplicates are emitted; after the
stability window, assert the final events.length is exactly 1 and call
record(testName, "PASS" or "FAIL", ...) using the same detail string (which
references chunks, r.provider, r.model) so the logic in this block (variables
events, chunks, r.provider, r.model, testName, record) is preserved but robust
to slow CI runners and late duplicate events.
---
Nitpick comments:
In `@test/continuous-test-suite-issue-04-generation-end-dedup.ts`:
- Around line 159-166: The current test loop (calling countEmissionsForGenerate
and countEmissionsForStream sequentially over TARGETS) doesn't exercise
AsyncLocalStorage concurrency; add a new test function (e.g.,
countEmissionsForConcurrentStreams) that creates a single NeuroLink instance and
runs two sdk.stream() calls in parallel, listens for "generation:end" on
sdk.getEventEmitter(), drains both streams, waits briefly, and asserts you saw
exactly two generation:end events for that instance; then invoke this new
function for appropriate targets (those with native generation:end like vertex
or google-ai-studio) in main alongside the existing checks and ensure you handle
skips, provider errors, and sdk.shutdown similar to the other helpers (use
TARGETS, sdk.stream, countEmissionsForGenerate/countEmissionsForStream
references to locate insertion points).
- Around line 77-157: The two functions countEmissionsForGenerate and
countEmissionsForStream duplicate the same
SKIP/PASS/FAIL/error/shutdown/event-listening plumbing; refactor by extracting a
helper (e.g., runEmissionTest or runCall) that accepts (target, label, runner)
where runner is an async function called with the sdk that returns { provider,
model } (for stream the runner should also drain r.stream and wait the 250ms
grace period before returning). Move the shared event subscription
(sdk.getEventEmitter().on("generation:end", ...)), skip check, try/catch error
handling (isExpectedProviderError), record(...) calls and finally
sdk.shutdown?.() into the helper, and replace
countEmissionsForGenerate/countEmissionsForStream with two tiny callers that
pass a runner invoking sdk.generate(...) or sdk.stream(...) respectively.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 59987b15-3a05-434f-8216-b11c3cc57684
📒 Files selected for processing (7)
src/lib/neurolink.tssrc/lib/providers/googleAiStudio.tssrc/lib/providers/googleVertex.tssrc/lib/types/index.tssrc/lib/types/streamDedup.tstest/continuous-test-suite-issue-04-generation-end-dedup.tstest/helpers/envGuard.ts
| resolvedUsage = streamUsage; | ||
| if (!resolvedUsage && streamAnalytics) { | ||
| try { | ||
| const resolved = await Promise.resolve(streamAnalytics); | ||
| if (resolved?.tokenUsage) { | ||
| resolvedUsage = resolved.tokenUsage; | ||
| } | ||
| } catch { |
There was a problem hiding this comment.
Carry analytics through the synthetic generation:end result.
Right now the stream path resolves usage but drops streamAnalytics from the emitted result. Any listener shared with generate() that reads data.result.analytics?.cost still won't see stream costs, so the contract is still narrower than the non-stream path.
Suggested fix
- let resolvedUsage: unknown;
+ let resolvedUsage: unknown;
+ let resolvedAnalytics: AnalyticsData | undefined;
try {
for await (const chunk of mcpStream) {
chunkCount++;
@@
- resolvedUsage = streamUsage;
- if (!resolvedUsage && streamAnalytics) {
+ resolvedUsage = streamUsage;
+ if (streamAnalytics) {
try {
- const resolved = await Promise.resolve(streamAnalytics);
- if (resolved?.tokenUsage) {
- resolvedUsage = resolved.tokenUsage;
+ resolvedAnalytics = await Promise.resolve(streamAnalytics);
+ if (!resolvedUsage && resolvedAnalytics?.tokenUsage) {
+ resolvedUsage = resolvedAnalytics.tokenUsage;
}
} catch {
// non-blocking
}
}
@@
result: {
content: accumulatedContent,
usage: resolvedUsage,
+ analytics: resolvedAnalytics,
model: finalModel,
provider: finalProvider,
finishReason: finalFinishReason,
},Also applies to: 6724-6749
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@src/lib/neurolink.ts` around lines 6622 - 6629, The stream path resolves
token usage into resolvedUsage but does not attach streamAnalytics to the
synthetic "generation:end" result, so listeners of generate() won't see stream
cost data; update the code that emits the synthetic generation:end result (the
place that assigns resolvedUsage from streamUsage/streamAnalytics) to also set
result.analytics = result.analytics ?? {}; and copy the resolved streamAnalytics
(or the resolved object containing tokenUsage/cost) into result.analytics (e.g.,
result.analytics.cost or result.analytics.tokenUsage as appropriate) so the
emitted result mirrors the non-stream path; apply the same change in the other
stream branch around lines where streamAnalytics is resolved (also referenced in
the other block handling stream results).
| const processedStream = streamGenerationEndContext.run( | ||
| dedupContext, | ||
| buildProcessedStream, | ||
| ); |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
node <<'NODE'
const { AsyncLocalStorage } = require('async_hooks');
const als = new AsyncLocalStorage();
const makeIterator = () => (async function* () {
console.log('first resume store =', als.getStore());
yield 1;
await Promise.resolve();
console.log('second resume store =', als.getStore());
})();
(async () => {
const it = als.run({ streamId: 'demo' }, makeIterator);
await it.next();
await it.next();
})();
NODERepository: juspay/neurolink
Length of output: 123
🏁 Script executed:
# First, let's check if the file exists and get its size
wc -l src/lib/neurolink.tsRepository: juspay/neurolink
Length of output: 87
🏁 Script executed:
# Read the specific lines mentioned in the review comment
sed -n '6825,6835p' src/lib/neurolink.tsRepository: juspay/neurolink
Length of output: 481
🏁 Script executed:
# Search for buildProcessedStream definition to understand what it does
rg -A 20 "const buildProcessedStream|function buildProcessedStream" src/lib/neurolink.tsRepository: juspay/neurolink
Length of output: 962
🏁 Script executed:
# Search for streamGenerationEndContext to understand how it's created and used
rg -B 5 -A 5 "streamGenerationEndContext" src/lib/neurolink.tsRepository: juspay/neurolink
Length of output: 1285
🏁 Script executed:
# Search for markStreamProviderEmittedGenerationEnd to understand its implementation
rg -B 5 -A 10 "markStreamProviderEmittedGenerationEnd" src/lib/neurolink.tsRepository: juspay/neurolink
Length of output: 1248
🏁 Script executed:
# Find all call sites of markStreamProviderEmittedGenerationEnd
rg -B 3 -A 3 "markStreamProviderEmittedGenerationEnd()" src/lib/Repository: juspay/neurolink
Length of output: 3993
🏁 Script executed:
# Also check if it's exported and used in external files
rg "markStreamProviderEmittedGenerationEnd" --type ts --type tsxRepository: juspay/neurolink
Length of output: 88
🏁 Script executed:
# Find processStreamResult to see how it iterates the stream
rg -B 5 -A 20 "processStreamResult" src/lib/neurolink.ts | head -80Repository: juspay/neurolink
Length of output: 1649
🏁 Script executed:
# Let me also check where dedupContext is used after the stream is processed
rg -A 10 "providerEmitted" src/lib/neurolink.tsRepository: juspay/neurolink
Length of output: 1270
AsyncLocalStorage context is lost during stream iteration; dedup will silently fail.
The reproducer proves that AsyncLocalStorage.run() does not preserve context across async generator resumptions—both logs returned undefined instead of the store value. Since buildProcessedStream() returns an async generator that executes during iteration (inside processStreamResult()), calls to markStreamProviderEmittedGenerationEnd() from provider code will execute outside the run() context. This means dedupContext.providerEmitted will never be set, and the dedup check will always emit a duplicate generation:end event even when the provider already emitted one.
Replace streamGenerationEndContext.run(dedupContext, buildProcessedStream) with an iterator-bound context that wraps the result of buildProcessedStream() rather than its creation.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@src/lib/neurolink.ts` around lines 6828 - 6831, The AsyncLocalStorage context
is lost because streamGenerationEndContext.run wraps the generator creation
(buildProcessedStream) instead of the async-iterator itself; change to run the
context around the iterator so the store is active during iteration.
Specifically, call buildProcessedStream() to get the async iterator, then wrap
that iterator with a context-bound wrapper created by
streamGenerationEndContext.run(dedupContext, () => iterator) (or equivalent) so
each next()/throw()/return() invocation runs inside streamGenerationEndContext;
ensure processStreamResult still iterates the wrapped iterator and that
markStreamProviderEmittedGenerationEnd() can read/write
dedupContext.providerEmitted during iteration.
| // small grace period for any post-stream async emit | ||
| await new Promise((r) => setTimeout(r, 250)); | ||
| const detail = `count=${events.length}; chunks=${chunks}; provider=${r.provider}; model=${r.model}`; | ||
| if (events.length === 1) { | ||
| record(testName, "PASS", `expected 1, ${detail}`); | ||
| } else { | ||
| record(testName, "FAIL", `expected 1, got ${events.length}: ${detail}`); | ||
| } |
There was a problem hiding this comment.
Fixed 250 ms grace period is fragile.
A hard-coded setTimeout(250) after stream drain may be too short on slow CI runners or for providers that flush generation:end from a microtask chain after the last chunk, producing a flaky count=0 FAIL. Conversely, if a buggy provider emits a duplicate >250 ms later, the test will pass and miss it.
Consider polling for the expected count with a deadline (e.g., wait until events.length >= 1 or a 2 s timeout, then verify it stays at 1 for an additional ~250 ms window) instead of a single fixed sleep.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@test/continuous-test-suite-issue-04-generation-end-dedup.ts` around lines 139
- 146, Replace the fragile fixed 250ms grace sleep with a deadline-based poll:
repeatedly check the events array (events.length) in a short interval (e.g.,
50ms) until events.length >= 1 or a deadline (e.g., 2000ms) is reached, then
once you observe events.length >= 1 wait an additional stability window (~250ms)
to ensure no late duplicates are emitted; after the stability window, assert the
final events.length is exactly 1 and call record(testName, "PASS" or "FAIL",
...) using the same detail string (which references chunks, r.provider, r.model)
so the logic in this block (variables events, chunks, r.provider, r.model,
testName, record) is preserved but robust to slow CI runners and late duplicate
events.
|
Force-pushed addressing the reviewer's concern that AsyncLocalStorage doesn't reliably propagate through async-generator yield boundaries. New approach: drop ALS entirely. Attach a fresh mutable Why it's concurrency-safe: each Wire-up:
Suite re-verified: 6 PASS (vertex/google-ai-studio/litellm × generate+stream count=1). @coderabbitai full review |
7e095fb to
9c67e68
Compare
🤖 AI Review & Build Compliance ✅Status: AI analysis complete • Build rules validated • Ready for review 📊 View detailed analysis results🛡️ Analysis Complete
📋 Ready for Merge When
🤖 AI analysis complete - check individual code comments for specific feedback |
|
Tip For best results, initiate chat on the files or code changes. 🧩 Analysis chain🏁 Script executed: Length of output: 9887 🏁 Script executed: Length of output: 9800
Overall assessmentThe closure-capture approach is correct and directly sidesteps the ALS yield-boundary propagation concern. Per-call isolation is guaranteed by the fact that each Item 1 — Stale "AsyncLocalStorage" comments in two placesThe previous ALS-based implementation left behind stale references in two spots:
This contradicts the current approach. Suggested replacement: /**
* Curator P2-4 dedup (concurrency-safe): per-stream context that lets
* the orchestration's `runStandardStreamRequest` finally block know
* whether a native provider path already emitted `generation:end` for
* THIS stream. Attached directly to `enhancedOptions._streamDedupContext`
* so each `sdk.stream()` call carries its own isolated instance —
* concurrent streams on the same NeuroLink instance cannot interfere.
*/
export type StreamGenerationEndContext = { providerEmitted: boolean };
Should be updated to:
Item 2 — Bedrock and Ollama: confirming no action neededFrom the script output,
Finding Correctness checklist
Carry-over non-blocking itemFinding Summary: The implementation is correct and concurrency-safe. Fix the two stale "AsyncLocalStorage" comments and this is ready to merge. Happy to confirm once those are updated. ✏️ Learnings added
🧠 Learnings used |
Curator P2-4: cost listeners that subscribe to `generation:end` previously received zero events from `sdk.stream()` calls. The doc described "fires twice"; the bug on shipped 9.56.x is the opposite direction — `stream()` emitted `stream:complete` but never `generation:end`, leaving any listener (cost-tracker, audit log, alerting) with nothing. Reproduction (all real providers): generate / vertex count=1 PASS stream / vertex count=0 FAIL ← bug generate / google-ai-studio count=1 PASS stream / google-ai-studio count=0 FAIL ← bug generate / litellm count=1 PASS stream / litellm count=0 FAIL ← bug Fix: in `runStandardStreamRequest`'s `processedStream` generator, emit `generation:end` exactly once in the `finally` block with the final stream state (provider, model, content, usage, finishReason, toolsUsed, prompt, temperature, maxTokens, success, error, pipelineAHandled). Hoist `resolvedUsage` to the generator scope so it's available in `finally`. The event payload mirrors `generate()`'s shape so listeners receive a consistent contract across both APIs. Backward compatibility: additive — listeners that previously received zero events on streams now receive one event with the same shape `generate()` produces.
9c67e68 to
9011edd
Compare
🤖 AI Review & Build Compliance ✅Status: AI analysis complete • Build rules validated • Ready for review 📊 View detailed analysis results🛡️ Analysis Complete
📋 Ready for Merge When
🤖 AI analysis complete - check individual code comments for specific feedback |
|
🎉 This PR is included in version 9.59.1 🎉 The release is available on: Your semantic-release bot 📦🚀 |
Summary
Curator P2-4: cost listeners that subscribe to
generation:endpreviously received zero events fromsdk.stream()calls. The doc described "fires twice"; the bug on shipped 9.56.x is the opposite direction —stream()emittedstream:completebut nevergeneration:end, leaving any listener (cost-tracker, audit log, alerting) with nothing.Reproduction (real providers, before fix)
After fix: all stream rows return
count=1.Fix
In
runStandardStreamRequest'sprocessedStreamgenerator, emitgeneration:endexactly once in thefinallyblock with the final stream state (provider, model, content, usage, finishReason, toolsUsed, prompt, temperature, maxTokens, success, error,pipelineAHandled: true). HoistresolvedUsageto the generator scope so it's available infinally. The event payload mirrorsgenerate()'s shape so listeners receive a consistent contract across both APIs.Backward compatibility
Additive only. Listeners that previously received zero events on streams now receive one event with the same shape
generate()produces. Any listener that already expected the event will start working; any listener that didn't will still ignore the new event.Verification
Expected: 6 passed (3 generate + 3 stream); the OpenAI rows may SKIP/FAIL on env-specific quota / tool-injection issues unrelated to this fix.
Test plan
releasefor vertex, google-ai-studio, litellmgenerate()paths (still count=1)pipelineAHandled: trueflag preserved for Langfuse exporter dedupSummary by CodeRabbit
Release Notes
Bug Fixes
Tests