Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@
"test:mcp": "npx tsx test/continuous-test-suite-mcp-http.ts",
"test:media": "npx tsx test/continuous-test-suite-media-gen.ts",
"test:memory": "npx tsx test/continuous-test-suite-memory.ts",
"test:middleware": "npx tsx test/continuous-test-suite-middleware.ts",
"test:observability": "npx tsx test/continuous-test-suite-observability.ts",
"test:ppt": "npx tsx test/continuous-test-suite-ppt.ts",
"test:providers": "npx tsx test/continuous-test-suite-providers.ts",
Expand Down
1 change: 1 addition & 0 deletions src/lib/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,7 @@ export {
createAnalyticsMiddleware,
getAnalyticsMetrics,
} from "./middleware/builtin/analytics.js";
export { createLifecycleMiddleware } from "./middleware/builtin/lifecycle.js";
export { MiddlewareFactory } from "./middleware/factory.js";
export { ExporterRegistry } from "./observability/exporterRegistry.js";
export { NoOpExporter } from "./observability/exporters/baseExporter.js";
Expand Down
201 changes: 201 additions & 0 deletions src/lib/middleware/builtin/lifecycle.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,201 @@
/**
* Lifecycle Middleware
*
* Provides onFinish, onError, and onChunk callbacks for observing
* generation and streaming lifecycle events.
*
* This middleware is automatically enabled when lifecycle callbacks
* (onFinish, onError, onChunk) are passed in GenerateOptions or StreamOptions.
*/

import type { LanguageModelV1Middleware } from "ai";
import type {
NeuroLinkMiddleware,
NeuroLinkMiddlewareMetadata,
LifecycleMiddlewareConfig,
} from "../../types/middlewareTypes.js";
import { logger } from "../../utils/logger.js";
import { isRecoverableError } from "../../utils/errorHandling.js";

export function createLifecycleMiddleware(
config: LifecycleMiddlewareConfig = {},
): NeuroLinkMiddleware {
const metadata: NeuroLinkMiddlewareMetadata = {
id: "lifecycle",
name: "Lifecycle Callbacks",
description:
"Provides onFinish, onError, and onChunk callbacks for generation and streaming lifecycle events",
priority: 110,
defaultEnabled: false,
};

const middleware: LanguageModelV1Middleware = {
wrapGenerate: async ({ doGenerate }) => {
const startTime = Date.now();

try {
const result = await doGenerate();

if (config.onFinish) {
try {
const callbackResult = config.onFinish({
text: result.text ?? "",
usage: result.usage
? {
promptTokens: result.usage.promptTokens ?? 0,
completionTokens: result.usage.completionTokens ?? 0,
}
: undefined,
duration: Date.now() - startTime,
finishReason: result.finishReason,
});
Promise.resolve(callbackResult).catch((e) => {
logger.warn("[LifecycleMiddleware] onFinish callback error:", e);
});
} catch (e) {
logger.warn("[LifecycleMiddleware] onFinish callback error:", e);
}
}

return result;
} catch (error) {
if (config.onError) {
const err = error instanceof Error ? error : new Error(String(error));
try {
const callbackResult = config.onError({
error: err,
duration: Date.now() - startTime,
recoverable: isRecoverableError(err),
});
Promise.resolve(callbackResult).catch((e) => {
logger.warn("[LifecycleMiddleware] onError callback error:", e);
});
} catch (e) {
logger.warn("[LifecycleMiddleware] onError callback error:", e);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

throw error;
}
},

wrapStream: async ({ doStream }) => {
const startTime = Date.now();

try {
const result = await doStream();

if (!config.onChunk && !config.onFinish && !config.onError) {
return result;
}

let sequenceNumber = 0;
let accumulatedText = "";

const transformStream = new TransformStream({
transform(chunk, controller) {
try {
if (chunk.type === "text-delta") {
accumulatedText += chunk.textDelta;
}

if (config.onChunk && chunk.type) {
try {
const callbackResult = config.onChunk({
type: chunk.type,
textDelta:
chunk.type === "text-delta" ? chunk.textDelta : undefined,
sequenceNumber: sequenceNumber++,
});
Promise.resolve(callbackResult).catch((e) => {
logger.warn(
"[LifecycleMiddleware] onChunk callback error:",
e,
);
});
} catch (e) {
logger.warn(
"[LifecycleMiddleware] onChunk callback error:",
e,
);
}
}

controller.enqueue(chunk);
} catch (error) {
if (config.onError) {
const err =
error instanceof Error ? error : new Error(String(error));
try {
const callbackResult = config.onError({
error: err,
duration: Date.now() - startTime,
recoverable: isRecoverableError(err),
});
Promise.resolve(callbackResult).catch((e) => {
logger.warn(
"[LifecycleMiddleware] onError callback error:",
e,
);
});
} catch (e) {
logger.warn(
"[LifecycleMiddleware] onError callback error:",
e,
);
}
}
throw error;
}
},
flush() {
if (config.onFinish) {
try {
const callbackResult = config.onFinish({
text: accumulatedText,
duration: Date.now() - startTime,
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Promise.resolve(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),
};
Comment on lines +174 to +177

Copilot AI Mar 20, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copilot uses AI. Check for mistakes.
} catch (error) {
if (config.onError) {
const err = error instanceof Error ? error : new Error(String(error));
try {
const callbackResult = config.onError({
error: err,
duration: Date.now() - startTime,
recoverable: isRecoverableError(err),
});
Promise.resolve(callbackResult).catch((e) => {
logger.warn("[LifecycleMiddleware] onError callback error:", e);
});
} catch (e) {
logger.warn("[LifecycleMiddleware] onError callback error:", e);
}
}

throw error;
Comment on lines +82 to +195

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

🧩 Analysis chain

🏁 Script executed:

cat -n src/lib/middleware/builtin/lifecycle.ts | head -200

Repository: juspay/neurolink

Length of output: 6120


🏁 Script executed:

rg "pipeThrough|stream.*error|\.catch\(\)" src/lib --type ts -A 2 -B 2 | head -60

Repository: 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 -80

Repository: 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 -100

Repository: juspay/neurolink

Length of output: 5907


🏁 Script executed:

rg "error.*stream|stream.*error" src/lib/middleware --type ts -B 3 -A 3

Repository: juspay/neurolink

Length of output: 610


🏁 Script executed:

rg "wrapStream.*async.*doStream" src/lib/middleware -A 50 | head -150

Repository: juspay/neurolink

Length of output: 10781


🏁 Script executed:

rg "stream.pipeThrough|pipeThrough.*TransformStream" src/lib -B 5 -A 10 | head -120

Repository: 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 3

Repository: juspay/neurolink

Length of output: 42


🏁 Script executed:

rg "consumeAndEvaluateStream" src/lib --type ts -A 30 | head -80

Repository: 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 -100

Repository: juspay/neurolink

Length of output: 848


🏁 Script executed:

rg "stream\.getReader|reader\.read|ReadableStream" src/lib/middleware --type ts -B 2 -A 5

Repository: 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:

  1. Calls stream.getReader()
  2. Catches errors from reader.read() and iteration failures
  3. Invokes onError with the caught error
  4. 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.

}
},
};

return { ...middleware, metadata };
}
3 changes: 3 additions & 0 deletions src/lib/middleware/factory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import { MiddlewareRegistry } from "./registry.js";
import { createAnalyticsMiddleware } from "./builtin/analytics.js";
import { createGuardrailsMiddleware } from "./builtin/guardrails.js";
import { createAutoEvaluationMiddleware } from "./builtin/autoEvaluation.js";
import { createLifecycleMiddleware } from "./builtin/lifecycle.js";
import { logger } from "../utils/logger.js";

/**
Expand Down Expand Up @@ -43,6 +44,7 @@ export class MiddlewareFactory {
analytics: createAnalyticsMiddleware,
guardrails: createGuardrailsMiddleware,
autoEvaluation: createAutoEvaluationMiddleware,
lifecycle: createLifecycleMiddleware,
};

// Register built-in presets
Expand Down Expand Up @@ -192,6 +194,7 @@ export class MiddlewareFactory {
analytics: createAnalyticsMiddleware,
guardrails: createGuardrailsMiddleware,
autoEvaluation: createAutoEvaluationMiddleware,
lifecycle: createLifecycleMiddleware,
};
logger.debug("Getting creator for middleware ID:", id);
return builtInMiddlewareCreators[id];
Expand Down
3 changes: 3 additions & 0 deletions src/lib/middleware/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,5 +29,8 @@ export type { LanguageModelV1Middleware } from "ai";
// Factory for creating and applying middleware chains
export { MiddlewareFactory };

// Built-in middleware creators
export { createLifecycleMiddleware } from "./builtin/lifecycle.js";

// Export the factory as the default export for clean, direct usage
export default MiddlewareFactory;
45 changes: 44 additions & 1 deletion src/lib/neurolink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2933,13 +2933,13 @@
* @see {@link stream} for streaming generation
* @since 1.0.0
*/
async generate(

Check warning on line 2936 in src/lib/neurolink.ts

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async method 'generate' has too many lines (544). Maximum allowed is 300

Check warning on line 2936 in src/lib/neurolink.ts

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'generate' has too many lines (544). Maximum allowed is 300
optionsOrPrompt: GenerateOptions | string,
): Promise<GenerateResult> {
return tracers.sdk.startActiveSpan(
"neurolink.generate",
{ kind: SpanKind.INTERNAL },
async (generateSpan) => {

Check warning on line 2942 in src/lib/neurolink.ts

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async arrow function has too many lines (536). Maximum allowed is 300

Check warning on line 2942 in src/lib/neurolink.ts

View workflow job for this annotation

GitHub Actions / test (20)

Async arrow function has too many lines (536). Maximum allowed is 300
// Set metrics trace context for parent-child span linking.
// The generation span will be the root (no parentSpanId).
// Tool spans will be children of the root span via rootSpanId.
Expand Down Expand Up @@ -3022,6 +3022,27 @@
});
}

// Auto-inject lifecycle middleware when callbacks are provided
// (must happen before workflow/PPT early returns so those paths get middleware too)
if (options.onFinish || options.onError) {
options.middleware = {
...options.middleware,
middlewareConfig: {
...options.middleware?.middlewareConfig,
lifecycle: {
...options.middleware?.middlewareConfig?.lifecycle,
enabled: true,
config: {
...options.middleware?.middlewareConfig?.lifecycle
?.config,
onFinish: options.onFinish,
onError: options.onError,
},
},
},
};
}

// Check if workflow is requested
if (options.workflow || options.workflowConfig) {
return await this.generateWithWorkflow(options);
Expand Down Expand Up @@ -3241,6 +3262,7 @@
fileRegistry: this.fileRegistry,
abortSignal: options.abortSignal,
skipToolPromptInjection: options.skipToolPromptInjection,
middleware: options.middleware,
};

// Auto-map top-level sessionId/userId to context for convenience
Expand Down Expand Up @@ -3290,7 +3312,6 @@
});
}

// Use redesigned generation logic
const textResult =
await this.generateTextInternal(textOptions);

Expand Down Expand Up @@ -3415,6 +3436,7 @@
);

generateSpan.setStatus({ code: SpanStatusCode.OK });

return generateResult;
},
);
Expand Down Expand Up @@ -5610,6 +5632,27 @@

this.emitStreamStartEvents(options, startTime);

// Auto-inject lifecycle middleware when callbacks are provided
// (must happen before workflow early return so that path gets middleware too)
if (options.onFinish || options.onError || options.onChunk) {
options.middleware = {
...options.middleware,
middlewareConfig: {
...options.middleware?.middlewareConfig,
lifecycle: {
...options.middleware?.middlewareConfig?.lifecycle,
enabled: true,
config: {
...options.middleware?.middlewareConfig?.lifecycle?.config,
onFinish: options.onFinish,
onError: options.onError,
onChunk: options.onChunk,
},
},
},
};
}

// Check if workflow is requested
if (options.workflow || options.workflowConfig) {
const result = await this.streamWithWorkflow(options, startTime);
Expand Down
15 changes: 14 additions & 1 deletion src/lib/types/generateTypes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,11 @@ import type { JsonValue } from "./common.js";
import type { Content, ImageWithAltText } from "./content.js";
import type { ChatMessage, ConversationMemoryConfig } from "./conversation.js";
import type { EvaluationData } from "./evaluation.js";
import type { MiddlewareFactoryOptions } from "./middlewareTypes.js";
import type {
MiddlewareFactoryOptions,
OnFinishCallback,
OnErrorCallback,
} from "./middlewareTypes.js";
import type {
DirectorModeOptions,
DirectorSegment,
Expand Down Expand Up @@ -444,6 +448,15 @@ export type GenerateOptions = {
* @internal Set by NeuroLink SDK — not typically used directly by consumers.
*/
fileRegistry?: unknown;

/** Per-call middleware configuration. */
middleware?: import("./middlewareTypes.js").MiddlewareFactoryOptions;

/** Callback invoked when generation completes successfully. */
onFinish?: OnFinishCallback;

/** Callback invoked when generation encounters an error. */
onError?: OnErrorCallback;
};

/**
Expand Down
Loading
Loading