feat(middleware): add lifecycle middleware with onFinish, onError, onChunk callbacks - #888
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 introduces a lifecycle middleware system enabling consumers to attach Changes
Sequence DiagramsequenceDiagram
participant Consumer
participant NeuroLink
participant Middleware
participant Generator
participant Logger
Consumer->>NeuroLink: generate(text, {onFinish, onError})
NeuroLink->>NeuroLink: Merge lifecycle middleware config
NeuroLink->>Middleware: wrapGenerate()
Middleware->>Middleware: Record startTime
Middleware->>Generator: Call doGenerate()
alt Success
Generator-->>Middleware: Return {text, usage}
Middleware->>Middleware: Calculate duration
Middleware->>Middleware: Invoke onFinish({text, usage, duration}) async
Middleware->>Logger: Catch/log if callback fails
Middleware-->>NeuroLink: Return result
else Error
Generator-->>Middleware: Throw error
Middleware->>Middleware: Calculate duration
Middleware->>Middleware: Check isRecoverableError()
Middleware->>Middleware: Invoke onError({error, duration, recoverable}) async
Middleware->>Logger: Catch/log if callback fails
Middleware-->>NeuroLink: Rethrow error
end
NeuroLink-->>Consumer: Return result or throw
sequenceDiagram
participant Consumer
participant NeuroLink
participant Middleware
participant Generator
participant TransformStream
participant Logger
Consumer->>NeuroLink: stream(text, {onChunk, onFinish})
NeuroLink->>NeuroLink: Merge lifecycle middleware config
NeuroLink->>Middleware: wrapStream()
Middleware->>Middleware: Record startTime
Middleware->>Generator: Call doStream()
Generator-->>Middleware: Return result stream
Middleware->>TransformStream: Wrap stream with transform
loop Per Chunk
Consumer->>TransformStream: Read chunk
TransformStream->>Middleware: transform(chunk)
Middleware->>Middleware: Increment sequenceNumber
Middleware->>Middleware: Invoke onChunk({type, textDelta, sequenceNumber}) async
Middleware->>Logger: Catch/log if callback fails
TransformStream-->>Consumer: Forward chunk
end
Consumer->>TransformStream: Stream ends
TransformStream->>Middleware: flush()
Middleware->>Middleware: Calculate duration
Middleware->>Middleware: Invoke onFinish({text: "", duration}) async
Middleware->>Logger: Catch/log if callback fails
TransformStream-->>Consumer: Complete stream
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested labels
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 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 Tip You can disable poems in the walkthrough.Disable the |
✅ 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 |
There was a problem hiding this comment.
Pull request overview
Adds a built-in “lifecycle” middleware to NeuroLink so consumers can provide onFinish, onError, and onChunk callbacks directly on generate() / stream() calls, with automatic middleware injection and a supporting recoverable-error classifier plus a continuous test suite.
Changes:
- Introduces
createLifecycleMiddleware()(generate + stream wrappers) and lifecycle callback payload types. - Extends
GenerateOptions/StreamOptionswith lifecycle callback fields and auto-injects lifecycle middleware when provided. - Adds
isRecoverableError()and a new continuous middleware test runner + npm script.
Reviewed changes
Copilot reviewed 11 out of 11 changed files in this pull request and generated 7 comments.
Show a summary per file
| File | Description |
|---|---|
src/lib/middleware/builtin/lifecycle.ts |
Implements lifecycle middleware callbacks for generate/stream. |
src/lib/utils/errorHandling.ts |
Adds isRecoverableError() utility. |
src/lib/types/middlewareTypes.ts |
Adds lifecycle callback payload + config types. |
src/lib/types/generateTypes.ts |
Adds middleware and lifecycle callback options to GenerateOptions. |
src/lib/types/streamTypes.ts |
Adds lifecycle callback options to StreamOptions. |
src/lib/neurolink.ts |
Auto-injects lifecycle middleware based on provided callbacks. |
src/lib/middleware/factory.ts |
Registers lifecycle middleware creator. |
src/lib/middleware/index.ts |
Exports createLifecycleMiddleware. |
src/lib/index.ts |
Exports createLifecycleMiddleware from SDK entrypoint. |
test/continuous-test-suite-middleware.ts |
Adds continuous integration-style tests for lifecycle callbacks and error classification. |
package.json |
Adds test:middleware script. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| flush() { | ||
| if (config.onFinish) { | ||
| try { | ||
| const callbackResult = config.onFinish({ | ||
| text: "", | ||
| duration: Date.now() - startTime, | ||
| }); | ||
| if (callbackResult instanceof Promise) { |
There was a problem hiding this comment.
In wrapStream.flush(), onFinish is invoked with text: "", even though LifecycleFinishPayload.text is documented as “The generated text content”. This makes the callback payload misleading for streaming use cases. Consider accumulating textDelta chunks (or otherwise obtaining the final text) and passing the actual final text to onFinish for streams (and optionally include finishReason / usage when available).
| return { | ||
| ...result, | ||
| stream: result.stream.pipeThrough(transformStream), | ||
| }; |
There was a problem hiding this comment.
wrapStream only calls onError when doStream() throws. If the returned stream errors later during consumption, this middleware will never invoke onError, even though consumers will experience a stream failure. To reliably surface streaming errors, consider piping result.stream into a TransformStream via pipeTo(transformStream.writable) and attaching a .catch() handler that triggers onError (returning transformStream.readable to the caller), or otherwise intercept stream errors/cancellation explicitly.
| // Handle async callbacks non-blocking | ||
| if (callbackResult instanceof Promise) { | ||
| callbackResult.catch((e) => { | ||
| logger.warn( | ||
| "[LifecycleMiddleware] onChunk callback error:", | ||
| e, | ||
| ); | ||
| }); | ||
| } |
There was a problem hiding this comment.
The async-callback detection uses callbackResult instanceof Promise. That misses thenables and can fail across realms; it also skips error handling for async functions returning a non-native promise implementation. Prefer Promise.resolve(callbackResult).catch(...) (or checking typeof (callbackResult as any)?.then === "function") to make callback error isolation reliable.
| export function isRecoverableError(error: Error): boolean { | ||
| const message = error.message?.toLowerCase() || ""; | ||
|
|
||
| // Rate limit errors | ||
| if ( | ||
| message.includes("rate limit") || | ||
| message.includes("429") || | ||
| message.includes("too many requests") | ||
| ) { | ||
| return true; | ||
| } | ||
|
|
||
| // Timeout errors | ||
| if ( | ||
| message.includes("timeout") || | ||
| message.includes("etimedout") || | ||
| message.includes("timed out") | ||
| ) { | ||
| return true; | ||
| } | ||
|
|
||
| // Network errors | ||
| if ( | ||
| message.includes("econnreset") || | ||
| message.includes("econnrefused") || | ||
| message.includes("network") || | ||
| message.includes("socket") | ||
| ) { | ||
| return true; | ||
| } | ||
|
|
||
| // Server errors (5xx) | ||
| if ( | ||
| message.includes("500") || | ||
| message.includes("502") || | ||
| message.includes("503") || | ||
| message.includes("504") | ||
| ) { | ||
| return true; |
There was a problem hiding this comment.
isRecoverableError() duplicates most of isRetriableError() logic but (unlike isRetriableError) it ignores NeuroLinkError.retriable, which can lead to inconsistent classification for the same error type. Also, message.includes("500")/etc can match unrelated numbers. Consider: (1) returning error.retriable when error instanceof NeuroLinkError, and (2) using regex with word boundaries for HTTP status codes (or delegating to a shared pattern list) to avoid accidental matches.
| // ============================================ | ||
| // LIFECYCLE MIDDLEWARE TYPES | ||
| // ============================================ | ||
|
|
||
| /** | ||
| * Payload delivered to onFinish callbacks after generation or streaming completes. | ||
| */ | ||
| export type LifecycleFinishPayload = { | ||
| /** The generated text content */ | ||
| text: string; | ||
| /** Token usage from the provider */ | ||
| usage?: { promptTokens: number; completionTokens: number }; | ||
| /** Wall-clock duration in milliseconds */ | ||
| duration: number; | ||
| /** Why generation stopped */ | ||
| finishReason?: string; | ||
| }; | ||
|
|
||
| /** | ||
| * Payload delivered to onError callbacks when generation or streaming fails. | ||
| */ | ||
| export type LifecycleErrorPayload = { | ||
| /** The error that occurred */ | ||
| error: Error; | ||
| /** Wall-clock duration until failure in milliseconds */ | ||
| duration: number; | ||
| /** Whether the error is likely recoverable (rate limit, timeout, network) */ | ||
| recoverable: boolean; | ||
| }; | ||
|
|
||
| /** | ||
| * Payload delivered to onChunk callbacks for each streaming chunk. | ||
| */ | ||
| export type LifecycleChunkPayload = { | ||
| /** Chunk type from the AI SDK stream */ | ||
| type: string; | ||
| /** Text content for text-delta chunks */ | ||
| textDelta?: string; | ||
| /** Zero-based chunk sequence number */ | ||
| sequenceNumber: number; | ||
| }; | ||
|
|
||
| /** Callback invoked when generation or streaming finishes successfully. */ | ||
| export type OnFinishCallback = ( | ||
| payload: LifecycleFinishPayload, | ||
| ) => void | Promise<void>; | ||
|
|
||
| /** Callback invoked when generation or streaming encounters an error. */ | ||
| export type OnErrorCallback = ( | ||
| payload: LifecycleErrorPayload, | ||
| ) => void | Promise<void>; | ||
|
|
||
| /** Callback invoked for each chunk during streaming. */ | ||
| export type OnChunkCallback = ( | ||
| payload: LifecycleChunkPayload, | ||
| ) => void | Promise<void>; | ||
|
|
||
| /** | ||
| * Configuration for the lifecycle middleware. | ||
| * Pass callbacks to observe generation/streaming lifecycle events. | ||
| */ | ||
| export type LifecycleMiddlewareConfig = { | ||
| onFinish?: OnFinishCallback; | ||
| onError?: OnErrorCallback; | ||
| onChunk?: OnChunkCallback; | ||
| }; |
There was a problem hiding this comment.
The factory registers a built-in middleware with ID "lifecycle", and the SDK now auto-injects middlewareConfig.lifecycle, but BuiltInMiddlewareType (earlier in this file) does not include "lifecycle". This will cause typing drift/inconsistencies when users reference built-in middleware IDs. Add "lifecycle" to the BuiltInMiddlewareType union to keep the type system aligned with the registry.
| // Auto-inject lifecycle middleware when callbacks are provided | ||
| if (options.onFinish || options.onError) { | ||
| textOptions.middleware = { | ||
| ...textOptions.middleware, | ||
| middlewareConfig: { | ||
| ...textOptions.middleware?.middlewareConfig, | ||
| lifecycle: { | ||
| enabled: true, | ||
| config: { | ||
| onFinish: options.onFinish, | ||
| onError: options.onError, | ||
| }, | ||
| }, | ||
| }, |
There was a problem hiding this comment.
Auto-injection overwrites any existing middlewareConfig.lifecycle object with a new { enabled: true, config: { onFinish, onError } }. If a caller already provided middleware.middlewareConfig.lifecycle (e.g., with additional lifecycle options in the future), those settings will be lost. Consider merging with the existing lifecycle config/object (e.g., spread existing lifecycle and lifecycle.config) instead of replacing it wholesale.
| enhancedOptions.middleware = { | ||
| ...enhancedOptions.middleware, | ||
| middlewareConfig: { | ||
| ...enhancedOptions.middleware?.middlewareConfig, | ||
| lifecycle: { | ||
| enabled: true, | ||
| config: { | ||
| onFinish: options.onFinish, | ||
| onError: options.onError, | ||
| onChunk: options.onChunk, | ||
| }, |
There was a problem hiding this comment.
Same as generate(): this auto-injection replaces middlewareConfig.lifecycle rather than merging with any existing lifecycle configuration on enhancedOptions.middleware. If callers passed middleware options explicitly, their lifecycle config would be overwritten when onFinish/onError/onChunk are also provided. Consider merging existing lifecycle + lifecycle.config to avoid clobbering user configuration.
| enhancedOptions.middleware = { | |
| ...enhancedOptions.middleware, | |
| middlewareConfig: { | |
| ...enhancedOptions.middleware?.middlewareConfig, | |
| lifecycle: { | |
| enabled: true, | |
| config: { | |
| onFinish: options.onFinish, | |
| onError: options.onError, | |
| onChunk: options.onChunk, | |
| }, | |
| const existingMiddlewareConfig = | |
| enhancedOptions.middleware?.middlewareConfig; | |
| const existingLifecycle = existingMiddlewareConfig?.lifecycle; | |
| const existingLifecycleConfig = existingLifecycle?.config ?? {}; | |
| const mergedLifecycleConfig = { | |
| ...existingLifecycleConfig, | |
| ...(options.onFinish !== undefined | |
| ? { onFinish: options.onFinish } | |
| : {}), | |
| ...(options.onError !== undefined | |
| ? { onError: options.onError } | |
| : {}), | |
| ...(options.onChunk !== undefined | |
| ? { onChunk: options.onChunk } | |
| : {}), | |
| }; | |
| enhancedOptions.middleware = { | |
| ...enhancedOptions.middleware, | |
| middlewareConfig: { | |
| ...existingMiddlewareConfig, | |
| lifecycle: { | |
| ...existingLifecycle, | |
| // Ensure lifecycle is enabled when callbacks are provided, | |
| // but do not override an explicit existing `enabled` value. | |
| enabled: | |
| existingLifecycle?.enabled !== undefined | |
| ? existingLifecycle.enabled | |
| : true, | |
| config: mergedLifecycleConfig, |
There was a problem hiding this comment.
Actionable comments posted: 8
🧹 Nitpick comments (2)
src/lib/middleware/factory.ts (1)
40-48: Consider extracting duplicatedbuiltInMiddlewareCreatorsmap.The same
builtInMiddlewareCreatorsmap is defined in bothinitialize()(lines 40-48) andgetCreator()(lines 190-198). This duplication means any future middleware addition/removal requires updating two locations.Consider extracting this to a class-level constant or private property to maintain a single source of truth.
♻️ Optional: Extract to class property
export class MiddlewareFactory { public registry: MiddlewareRegistry; public presets = new Map<string, MiddlewarePreset>(); private options: MiddlewareFactoryOptions; + private static readonly builtInMiddlewareCreators: Record< + string, + (config?: Record<string, unknown>) => NeuroLinkMiddleware + > = { + analytics: createAnalyticsMiddleware, + guardrails: createGuardrailsMiddleware, + autoEvaluation: createAutoEvaluationMiddleware, + lifecycle: createLifecycleMiddleware, + };Then reference
MiddlewareFactory.builtInMiddlewareCreatorsin bothinitialize()andgetCreator().🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/middleware/factory.ts` around lines 40 - 48, Extract the duplicated builtInMiddlewareCreators map out of the two methods into a single shared location (e.g., a static class-level constant or a private property on MiddlewareFactory) and reference that single symbol from both initialize() and getCreator(); specifically, create MiddlewareFactory.builtInMiddlewareCreators (or this.builtInMiddlewareCreators) containing the mapping (analytics, guardrails, autoEvaluation, lifecycle -> their creator functions) and update initialize() and getCreator() to use that shared map so additions/removals only need one change.test/continuous-test-suite-middleware.ts (1)
133-145: Use the exported lifecycle payload types instead ofany.These tests are the main consumer contract for the new callbacks, but
anyerases the exact payload shape and lets breaking API changes compile silently. Typing them withLifecycleFinishPayload,LifecycleErrorPayload, andLifecycleChunkPayloadwill make the suite catch callback contract drift at compile time. As per coding guidelines, Maintain strict TypeScript type safety across all modules.Also applies to: 200-214, 324-335, 395-406, 461-475
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/continuous-test-suite-middleware.ts` around lines 133 - 145, Replace the untyped runtime payload variables with the exported lifecycle payload types to enforce compile-time contract checks: change finishPayload (used with the onFinish callback) to type LifecycleFinishPayload, any error callback payloads to LifecycleErrorPayload, and streaming/chunk payloads to LifecycleChunkPayload; update declarations where finishPayload, errorPayload, and chunkPayload are defined (and their associated callbacks like onFinish/onError/onUpdate) so the test suite imports and uses LifecycleFinishPayload, LifecycleErrorPayload, and LifecycleChunkPayload instead of any, and apply the same replacements at the other noted locations (around lines 200-214, 324-335, 395-406, 461-475) to ensure strict TypeScript typing across all callback handlers.
🤖 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/middleware/builtin/lifecycle.ts`:
- Around line 39-69: The wrapGenerate() lifecycle currently awaits consumer
callbacks (config.onFinish and config.onError) on the main request path which
can block request/streaming latency; change both usages in the lifecycle
middleware so they are invoked fire-and-forget instead of awaited: call
config.onFinish({...}) and config.onError({...}) without awaiting their Promise,
attach a .catch(...) to each invocation to log errors via logger.warn (as done
now inside the try/catch), and ensure in the onError branch you immediately
rethrow or return the original error path without waiting for the consumer
callback to complete; keep the same payload shape (text, usage, duration,
finishReason and error, duration, recoverable) and reference wrapGenerate(),
config.onFinish, and config.onError when making the change.
- Around line 88-119: The streaming path's flush() in the TransformStream calls
config.onFinish with text: "" instead of the accumulated final output; update
the lifecycle TransformStream to track and provide the final concatenated text
to config.onFinish (e.g., maintain a local accumulator variable updated in
transform(chunk, ...) when chunk.type === "text-delta" or chunk.textDelta
exists), then pass that accumulated text and the existing duration (Date.now() -
startTime) to config.onFinish in flush(); ensure errors from the callback are
handled the same way they are now.
- Around line 76-156: The stream returned from wrapStream must be wrapped so
mid-stream provider/network errors are caught and forwarded to config.onError;
update the code that returns result.stream.pipeThrough(transformStream) to
instead create a new ReadableStream (or TransformStream wrapper) that obtains a
reader via result.stream.getReader(), loops with reader.read(), enqueues chunks
to the downstream controller (and still runs the existing transform logic for
onChunk/onFinish), and in the catch path calls await config.onError({ error: err
instanceof Error ? err : new Error(String(err)), duration: Date.now() -
startTime, recoverable: isRecoverableError(err) }) (with the same try/catch
logging pattern used elsewhere) before controller.error(err) and rethrowing;
ensure the reader is released in finally and preserve sequenceNumber, startTime,
transformStream semantics and existing onChunk/onFinish handling in wrapStream,
transformStream, result.stream, config.onError, isRecoverableError.
In `@src/lib/neurolink.ts`:
- Around line 5675-5691: stream() only injects lifecycle middleware into
enhancedOptions on the happy path, so calls that return early (notably
streamWithWorkflow()) and the fallback in the catch (handleStreamError()) still
use the original options and never run onChunk/onFinish/onError; fix by moving
the lifecycle middleware injection to occur before any early returns or by
ensuring the fallback uses enhancedOptions: update streamWithWorkflow() (and any
early-return path) to apply the same enhancedOptions.middleware lifecycle wiring
(the block currently mutating enhancedOptions.middleware) before returning, or
change the catch path so handleStreamError() is invoked with enhancedOptions
instead of options so lifecycle handlers are present for workflow and
setup-fallback streams.
- Around line 3293-3308: generate() discards caller-provided middleware and can
skip lifecycle hooks because options.middleware is never copied into baseOptions
and lifecycle injection into textOptions happens after code paths that may
return early; fix by copying options.middleware into baseOptions when building
request options (ensure baseOptions.middleware = options.middleware or merged
with existing baseOptions.middleware) and move/perform the lifecycle injection
(merging into textOptions.middleware.middlewareConfig.lifecycle) before any
early returns in generate() and any branches that handle workflow or PPT
requests so onFinish/onError are always applied while preserving any existing
middleware entries.
In `@src/lib/types/middlewareTypes.ts`:
- Around line 268-333: BuiltInMiddlewareType is missing the "lifecycle" literal
so the public union is out of sync with the runtime registry; update the
exported BuiltInMiddlewareType union to include the "lifecycle" string literal
(e.g., add | "lifecycle") so callers can reference LifecycleMiddlewareConfig
without casting, and ensure any related exported types that enumerate built-in
keys (if present) are updated to include "lifecycle" as well; locate the
BuiltInMiddlewareType symbol in this file and add "lifecycle" to its union.
In `@test/continuous-test-suite-middleware.ts`:
- Around line 206-243: The test for sdk.generate currently treats a null
errorPayload as a PASS even if the failure occurred before middleware was
attached; change the logic in the test block that checks errorPayload
(variables: sdk.generate call, errorPayload, generationThrew) to mark the case
where generationThrew is true but errorPayload is null as SKIP (return null)
instead of PASS, update the logTest call to report "SKIP" and a SKIP message,
and apply the same change to the streaming variant test (the corresponding block
around lines 467-501) so skipped scenarios are represented by returning null
with SKIP status.
- Around line 153-175: The tests currently mark PASS when finishPayload (or
chunk payload) is missing, which hides missed lifecycle hooks; change the else
branch that currently logs PASS for a missing payload to instead log SKIP (or
similar skipped status) and return null so the suite doesn't false-pass when
onFinish/onChunk never fired; specifically update the block that inspects
finishPayload (variables finishPayload, content, and call site logTest("onFinish
fires after generation", ...)) to: on missing payload log SKIP and return null,
keep the existing type-check failure path to log FAIL and return false, and the
valid-payload path to log PASS; apply the exact same pattern to the streaming
onFinish and onChunk test blocks referenced.
---
Nitpick comments:
In `@src/lib/middleware/factory.ts`:
- Around line 40-48: Extract the duplicated builtInMiddlewareCreators map out of
the two methods into a single shared location (e.g., a static class-level
constant or a private property on MiddlewareFactory) and reference that single
symbol from both initialize() and getCreator(); specifically, create
MiddlewareFactory.builtInMiddlewareCreators (or this.builtInMiddlewareCreators)
containing the mapping (analytics, guardrails, autoEvaluation, lifecycle ->
their creator functions) and update initialize() and getCreator() to use that
shared map so additions/removals only need one change.
In `@test/continuous-test-suite-middleware.ts`:
- Around line 133-145: Replace the untyped runtime payload variables with the
exported lifecycle payload types to enforce compile-time contract checks: change
finishPayload (used with the onFinish callback) to type LifecycleFinishPayload,
any error callback payloads to LifecycleErrorPayload, and streaming/chunk
payloads to LifecycleChunkPayload; update declarations where finishPayload,
errorPayload, and chunkPayload are defined (and their associated callbacks like
onFinish/onError/onUpdate) so the test suite imports and uses
LifecycleFinishPayload, LifecycleErrorPayload, and LifecycleChunkPayload instead
of any, and apply the same replacements at the other noted locations (around
lines 200-214, 324-335, 395-406, 461-475) to ensure strict TypeScript typing
across all callback handlers.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 65a16b93-4f81-4119-a320-bf8a0dec4db8
📒 Files selected for processing (11)
package.jsonsrc/lib/index.tssrc/lib/middleware/builtin/lifecycle.tssrc/lib/middleware/factory.tssrc/lib/middleware/index.tssrc/lib/neurolink.tssrc/lib/types/generateTypes.tssrc/lib/types/middlewareTypes.tssrc/lib/types/streamTypes.tssrc/lib/utils/errorHandling.tstest/continuous-test-suite-middleware.ts
| wrapStream: async ({ doStream }) => { | ||
| const startTime = Date.now(); | ||
|
|
||
| try { | ||
| const result = await doStream(); | ||
|
|
||
| if (!config.onChunk && !config.onFinish) { | ||
| return result; | ||
| } | ||
|
|
||
| let sequenceNumber = 0; | ||
|
|
||
| const transformStream = new TransformStream({ | ||
| transform(chunk, controller) { | ||
| if (config.onChunk && chunk.type) { | ||
| try { | ||
| const callbackResult = config.onChunk({ | ||
| type: chunk.type, | ||
| textDelta: | ||
| chunk.type === "text-delta" ? chunk.textDelta : undefined, | ||
| sequenceNumber: sequenceNumber++, | ||
| }); | ||
| // Handle async callbacks non-blocking | ||
| if (callbackResult instanceof Promise) { | ||
| callbackResult.catch((e) => { | ||
| logger.warn( | ||
| "[LifecycleMiddleware] onChunk callback error:", | ||
| e, | ||
| ); | ||
| }); | ||
| } | ||
| } catch (e) { | ||
| logger.warn("[LifecycleMiddleware] onChunk callback error:", e); | ||
| } | ||
| } | ||
| controller.enqueue(chunk); | ||
| }, | ||
| flush() { | ||
| if (config.onFinish) { | ||
| try { | ||
| const callbackResult = config.onFinish({ | ||
| text: "", | ||
| duration: Date.now() - startTime, | ||
| }); | ||
| if (callbackResult instanceof Promise) { | ||
| callbackResult.catch((e) => { | ||
| logger.warn( | ||
| "[LifecycleMiddleware] onFinish callback error:", | ||
| e, | ||
| ); | ||
| }); | ||
| } | ||
| } catch (e) { | ||
| logger.warn( | ||
| "[LifecycleMiddleware] onFinish callback error:", | ||
| e, | ||
| ); | ||
| } | ||
| } | ||
| }, | ||
| }); | ||
|
|
||
| return { | ||
| ...result, | ||
| stream: result.stream.pipeThrough(transformStream), | ||
| }; | ||
| } catch (error) { | ||
| if (config.onError) { | ||
| const err = error instanceof Error ? error : new Error(String(error)); | ||
| try { | ||
| await config.onError({ | ||
| error: err, | ||
| duration: Date.now() - startTime, | ||
| recoverable: isRecoverableError(err), | ||
| }); | ||
| } catch (e) { | ||
| logger.warn("[LifecycleMiddleware] onError callback error:", e); | ||
| } | ||
| } | ||
|
|
||
| throw error; |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
cat -n src/lib/middleware/builtin/lifecycle.ts | head -200Repository: juspay/neurolink
Length of output: 6120
🏁 Script executed:
rg "pipeThrough|stream.*error|\.catch\(\)" src/lib --type ts -A 2 -B 2 | head -60Repository: juspay/neurolink
Length of output: 3372
🏁 Script executed:
rg "wrapStream|TransformStream" src/lib --type ts -B 2 -A 5 | grep -A 10 "TransformStream"Repository: juspay/neurolink
Length of output: 1941
🏁 Script executed:
rg "\.stream\s*=" src/lib --type ts -B 2 -A 3 | head -80Repository: juspay/neurolink
Length of output: 1593
🏁 Script executed:
cat -n src/lib/neurolink.ts | sed -n '3200,3350p'Repository: juspay/neurolink
Length of output: 8261
🏁 Script executed:
rg "for await.*stream|\.text\(\)|\.json\(\)" src/lib --type ts -B 2 -A 5 | head -100Repository: juspay/neurolink
Length of output: 5907
🏁 Script executed:
rg "error.*stream|stream.*error" src/lib/middleware --type ts -B 3 -A 3Repository: juspay/neurolink
Length of output: 610
🏁 Script executed:
rg "wrapStream.*async.*doStream" src/lib/middleware -A 50 | head -150Repository: juspay/neurolink
Length of output: 10781
🏁 Script executed:
rg "stream.pipeThrough|pipeThrough.*TransformStream" src/lib -B 5 -A 10 | head -120Repository: juspay/neurolink
Length of output: 3034
🏁 Script executed:
cat -n src/lib/middleware/builtin/autoEvaluation.ts | sed -n '80,130p'Repository: juspay/neurolink
Length of output: 1701
🏁 Script executed:
rg "readableStream.*error|addEventListener.*error" src/lib --type ts -B 2 -A 3Repository: juspay/neurolink
Length of output: 42
🏁 Script executed:
rg "consumeAndEvaluateStream" src/lib --type ts -A 30 | head -80Repository: juspay/neurolink
Length of output: 4380
🏁 Script executed:
cat -n src/lib/middleware/builtin/autoEvaluation.ts | sed -n '160,220p'Repository: juspay/neurolink
Length of output: 2309
🏁 Script executed:
cat -n src/lib/middleware/builtin/autoEvaluation.ts | sed -n '220,260p'Repository: juspay/neurolink
Length of output: 959
🏁 Script executed:
rg "generateStream|\.stream\)" src/lib/neurolink.ts -B 3 -A 5 | head -100Repository: juspay/neurolink
Length of output: 848
🏁 Script executed:
rg "stream\.getReader|reader\.read|ReadableStream" src/lib/middleware --type ts -B 2 -A 5Repository: juspay/neurolink
Length of output: 2603
Mid-stream failures never reach onError because error handling is not connected to stream iteration.
The try/catch block only wraps await doStream(). Once result.stream.pipeThrough(transformStream) is returned, any provider/network failure during chunk iteration—the primary failure mode for streaming—bypasses the catch block entirely. The TransformStream handlers only catch errors from callbacks, not source stream errors.
To fix this, wrap the returned stream with a reader loop that:
- Calls
stream.getReader() - Catches errors from
reader.read()and iteration failures - Invokes
onErrorwith the caught error - Propagates the error downstream via
controller.error()
Reference the pattern in autoEvaluation.ts (lines 180–224) which properly handles stream consumption errors with getReader() and try/finally.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@src/lib/middleware/builtin/lifecycle.ts` around lines 76 - 156, The stream
returned from wrapStream must be wrapped so mid-stream provider/network errors
are caught and forwarded to config.onError; update the code that returns
result.stream.pipeThrough(transformStream) to instead create a new
ReadableStream (or TransformStream wrapper) that obtains a reader via
result.stream.getReader(), loops with reader.read(), enqueues chunks to the
downstream controller (and still runs the existing transform logic for
onChunk/onFinish), and in the catch path calls await config.onError({ error: err
instanceof Error ? err : new Error(String(err)), duration: Date.now() -
startTime, recoverable: isRecoverableError(err) }) (with the same try/catch
logging pattern used elsewhere) before controller.error(err) and rethrowing;
ensure the reader is released in finally and preserve sequenceNumber, startTime,
transformStream semantics and existing onChunk/onFinish handling in wrapStream,
transformStream, result.stream, config.onError, isRecoverableError.
858c9e6 to
a74366a
Compare
Review Feedback Addressed (Cycle 1)All 15 inline review comments (10 deduplicated issues) from @copilot and @coderabbitai have been addressed. Changes Made
Validation
Requesting Re-review@copilot @coderabbitai — All feedback from this cycle has been addressed. Please re-review. |
|
Kicking off a full re-review now to verify all the changes. 🧠 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 |
a74366a to
16b8ec3
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 |
…Chunk callbacks Add a lifecycle middleware to the existing middleware system that provides consumer-facing callbacks on GenerateOptions and StreamOptions. When users pass onFinish/onError/onChunk, NeuroLink auto-injects the lifecycle middleware into the pipeline — no manual middleware setup needed. - New built-in middleware: src/lib/middleware/builtin/lifecycle.ts - Callback types in middlewareTypes.ts (LifecycleFinishPayload, etc.) - isRecoverableError utility in errorHandling.ts - Registered in MiddlewareFactory with priority 110, defaultEnabled: false - Continuous test suite: test/continuous-test-suite-middleware.ts
16b8ec3 to
cdbb563
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.30.0 🎉 The release is available on: Your semantic-release bot 📦🚀 |
Summary
src/lib/middleware/builtin/lifecycle.ts) to the existing middleware system that providesonFinish,onError, andonChunkcallbacks onGenerateOptionsandStreamOptionsisRecoverableError()utility toerrorHandling.tsfor classifying retryable errors (rate limit, timeout, network, 5xx)test/continuous-test-suite-middleware.ts) with 8 tests against real NeuroLink instancesUsage
Changes
src/lib/middleware/builtin/lifecycle.tssrc/lib/types/middlewareTypes.tssrc/lib/types/generateTypes.tsonFinish,onError,middlewarefields on GenerateOptionssrc/lib/types/streamTypes.tsonFinish,onError,onChunkfields on StreamOptionssrc/lib/utils/errorHandling.tsisRecoverableError()utilitysrc/lib/middleware/factory.tssrc/lib/middleware/index.tssrc/lib/index.tssrc/lib/neurolink.tspackage.jsontest/continuous-test-suite-middleware.tsTest plan
Summary by CodeRabbit
Release Notes
New Features
onFinish(triggered on successful completion),onError(triggered on operation failure), andonChunk(triggered during streaming chunks)Tests