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
120 changes: 119 additions & 1 deletion src/lib/neurolink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,7 @@
MCPTool,
RoutingDecision,
MetricsTraceContext,
StreamGenerationEndContext,
} from "./types/index.js";
import { emergencyContentTruncation } from "./context/emergencyTruncation.js";
import {
Expand Down Expand Up @@ -476,6 +477,41 @@
*/
const metricsTraceContextStorage = new AsyncLocalStorage<MetricsTraceContext>();

/**
* Curator P2-4 dedup (concurrency-safe): native providers emit
* `generation:end` on the shared SDK emitter. We attach a fresh
* mutable `dedupContext` object directly to the per-call
* `StreamOptions` (under `_streamDedupContext`) so each stream gets
* its own instance — concurrent streams have different option objects
* and therefore different contexts, so they cannot interfere.
*
* Native provider emit sites read `options._streamDedupContext` and
* flip `.providerEmitted = true` before emitting; the orchestration's
* finally block reads the same closed-over reference and skips its
* own emit when the flag is set.
*
* This avoids the AsyncLocalStorage approach which doesn't reliably
* propagate through async-generator yield boundaries when iteration
* happens from outside the original `run()` scope (e.g. when the
* consumer drives `for await of result.stream` after `sdk.stream(...)`
* returns).
*/
export const STREAM_DEDUP_CONTEXT_KEY = "_streamDedupContext" as const;

/**
* Native providers call this from their `generation:end` emit sites,
* passing the same `options` object they received. Safe no-op when
* the field isn't set.
*/
export function markStreamProviderEmittedGenerationEnd(
options: { _streamDedupContext?: StreamGenerationEndContext } | undefined,
): void {
const ctx = options?._streamDedupContext;
if (ctx) {
ctx.providerEmitted = true;
}
}

export class NeuroLink {
private mcpInitialized = false;
private mcpSkipped = false;
Expand Down Expand Up @@ -3137,7 +3173,7 @@
* Initialize event listeners that feed span data to MetricsAggregator.
* Listens to generation:end, stream:complete, and tool:end events.
*/
private initializeMetricsListeners(): void {

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

View workflow job for this annotation

GitHub Actions / test (20)

Method 'initializeMetricsListeners' has too many lines (334). Maximum allowed is 300

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

View workflow job for this annotation

GitHub Actions / test (20)

Method 'initializeMetricsListeners' has too many lines (334). Maximum allowed is 300

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

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Method 'initializeMetricsListeners' has too many lines (334). Maximum allowed is 300
this.emitter.on("generation:end", ((...args: unknown[]) => {
const data = args[0] as Record<string, unknown>;
// A2 fix: When Pipeline A (AI SDK → @langfuse/otel) already creates a
Expand Down Expand Up @@ -5994,7 +6030,7 @@
/**
* Direct provider generation (no MCP, no recursion)
*/
private async directProviderGeneration(

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

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'directProviderGeneration' has too many lines (314). Maximum allowed is 300

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

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'directProviderGeneration' has too many lines (314). Maximum allowed is 300

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

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async method 'directProviderGeneration' has too many lines (314). Maximum allowed is 300
options: TextGenerationOptions,
): Promise<TextGenerationResult> {
const startTime = Date.now();
Expand Down Expand Up @@ -6808,7 +6844,7 @@
return result;
}

private async runStandardStreamRequest(params: {

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

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'runStandardStreamRequest' has too many lines (393). Maximum allowed is 300

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

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'runStandardStreamRequest' has too many lines (393). Maximum allowed is 300

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

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async method 'runStandardStreamRequest' has too many lines (393). Maximum allowed is 300
options: StreamOptions;
streamSpan: ReturnType<typeof tracers.sdk.startSpan>;
spanStartTime: number;
Expand Down Expand Up @@ -6888,8 +6924,27 @@
const streamStartTime = Date.now();
const sessionId = (enhancedOptions.context as Record<string, unknown>)
?.sessionId as string | undefined;
// Curator P2-4 dedup (concurrency-safe): native provider stream paths
// (Gemini 3 on Vertex / Google AI Studio) emit `generation:end`
// themselves. We attach a per-stream mutable flag directly to
// `enhancedOptions._streamDedupContext` — native providers receive
// these options and flip the flag before their emit; this finally
// block reads the same closed-over reference. Concurrent streams
// have different option objects so the contexts don't interfere.
const dedupContext: StreamGenerationEndContext = {
providerEmitted: false,
};
(
enhancedOptions as StreamOptions & {
_streamDedupContext?: StreamGenerationEndContext;
}
)._streamDedupContext = dedupContext;
const processedStream = (async function* () {
let streamError: unknown;
// Curator P2-4: hoist `resolvedUsage` so the finally block can emit a
// single `generation:end` event with cost data. Cost listeners
// subscribe here; previously the stream path never fired it.
let resolvedUsage: unknown;
try {
for await (const chunk of mcpStream) {
chunkCount++;
Expand Down Expand Up @@ -6932,7 +6987,7 @@
);
}

let resolvedUsage = streamUsage;
resolvedUsage = streamUsage;
if (!resolvedUsage && streamAnalytics) {
try {
const resolved = await Promise.resolve(streamAnalytics);
Expand Down Expand Up @@ -7010,6 +7065,69 @@
error: metadata.error,
},
);

// Curator P2-4: emit `generation:end` exactly once per stream so
// cost listeners receive the same contract as for `generate()`.
// The previous implementation only fired `stream:complete`, leaving
// any subscriber to `generation:end` with zero events.
//
// Dedup: native provider stream paths (Gemini 3 on Vertex / Google
// AI Studio) already emit `generation:end` themselves so Pipeline B
// (Langfuse) records a GENERATION observation. Skip our emit when
// they already fired — preserves their Pipeline B observation
// source and keeps the "exactly once" contract. Per-stream flag
// is concurrency-safe because it's scoped via AsyncLocalStorage.
if (!dedupContext.providerEmitted) {
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,
});
} catch (emitError) {
logger.debug(
"[NeuroLink.stream] generation:end listener threw — ignored",
{
error:
emitError instanceof Error
? emitError.message
: String(emitError),
},
);
}
}

self._disableToolCacheForCurrentRequest = false;
cleanupListeners();

Expand Down
19 changes: 18 additions & 1 deletion src/lib/providers/googleAiStudio.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,10 @@
import { BaseProvider } from "../core/baseProvider.js";
import { DEFAULT_MAX_STEPS } from "../core/constants.js";
import { streamAnalyticsCollector } from "../core/streamAnalytics.js";
import type { NeuroLink } from "../neurolink.js";
import {
markStreamProviderEmittedGenerationEnd,
type NeuroLink,
} from "../neurolink.js";
import { SpanStatusCode } from "@opentelemetry/api";
import { ATTR, tracers, withClientSpan } from "../telemetry/index.js";
import type {
Expand Down Expand Up @@ -788,7 +791,7 @@
* Execute stream using native @google/genai SDK for Gemini 3 models
* This bypasses @ai-sdk/google to properly handle thought_signature
*/
private async executeNativeGemini3Stream(

Check warning on line 794 in src/lib/providers/googleAiStudio.ts

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'executeNativeGemini3Stream' has too many lines (343). Maximum allowed is 300

Check warning on line 794 in src/lib/providers/googleAiStudio.ts

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'executeNativeGemini3Stream' has too many lines (343). Maximum allowed is 300

Check warning on line 794 in src/lib/providers/googleAiStudio.ts

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async method 'executeNativeGemini3Stream' has too many lines (343). Maximum allowed is 300
options: StreamOptions,
): Promise<StreamResult> {
const modelName = options.model || this.modelName;
Expand All @@ -804,7 +807,7 @@
[ATTR.NL_PROVIDER]: this.providerName,
},
},
async (span) => {

Check warning on line 810 in src/lib/providers/googleAiStudio.ts

View workflow job for this annotation

GitHub Actions / test (20)

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

Check warning on line 810 in src/lib/providers/googleAiStudio.ts

View workflow job for this annotation

GitHub Actions / test (20)

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

Check warning on line 810 in src/lib/providers/googleAiStudio.ts

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async arrow function has too many lines (325). Maximum allowed is 300
const startTime = Date.now();
const timeout = this.getTimeout(options);
const timeoutController = createTimeoutController(
Expand Down Expand Up @@ -1039,6 +1042,13 @@
// AI SDK so experimental_telemetry is never injected; we emit manually.
const nativeStreamEmitter = this.neurolink?.getEventEmitter();
if (nativeStreamEmitter) {
// Curator P2-4 dedup: flag the per-stream context attached
// to options so the orchestration skips its own emit.
markStreamProviderEmittedGenerationEnd(
options as unknown as Parameters<
typeof markStreamProviderEmittedGenerationEnd
>[0],
);
nativeStreamEmitter.emit("generation:end", {
provider: this.providerName,
responseTime,
Expand Down Expand Up @@ -1073,6 +1083,13 @@
// Emit failure generation:end so Pipeline B records the failed stream
const errorEmitter = this.neurolink?.getEventEmitter();
if (errorEmitter) {
// Curator P2-4 dedup: flag the per-stream context attached
// to options so the orchestration skips its own emit.
markStreamProviderEmittedGenerationEnd(
options as unknown as Parameters<
typeof markStreamProviderEmittedGenerationEnd
>[0],
);
errorEmitter.emit("generation:end", {
provider: this.providerName,
responseTime: Date.now() - startTime,
Expand Down
13 changes: 12 additions & 1 deletion src/lib/providers/googleVertex.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,10 @@
GLOBAL_LOCATION_MODELS,
} from "../core/constants.js";
import { ModelConfigurationManager } from "../core/modelConfiguration.js";
import type { NeuroLink } from "../neurolink.js";
import {
markStreamProviderEmittedGenerationEnd,
type NeuroLink,
} from "../neurolink.js";
import { createProxyFetch } from "../proxy/proxyFetch.js";
import { ATTR, tracers, withClientSpan } from "../telemetry/index.js";
import type {
Expand Down Expand Up @@ -1879,7 +1882,7 @@
const getOrCreateStep = (stepIndex: number | undefined): VertexToolStep => {
const key = makeKey(stepIndex);
if (stepMap.has(key)) {
return stepMap.get(key)!;

Check warning on line 1885 in src/lib/providers/googleVertex.ts

View workflow job for this annotation

GitHub Actions / test (20)

Forbidden non-null assertion

Check warning on line 1885 in src/lib/providers/googleVertex.ts

View workflow job for this annotation

GitHub Actions / test (20)

Forbidden non-null assertion

Check warning on line 1885 in src/lib/providers/googleVertex.ts

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Forbidden non-null assertion
}
const step: VertexToolStep = {
type: "tool_step",
Expand Down Expand Up @@ -2334,8 +2337,16 @@
// Emit generation:end so Pipeline B (Langfuse) creates a GENERATION
// observation. The native @google/genai stream path on Vertex bypasses the
// Vercel AI SDK so experimental_telemetry is never injected; we emit manually.
// Curator P2-4 dedup: flag the per-stream context attached to options
// so the orchestration in `runStandardStreamRequest` knows we already
// emitted and skips its own emit (preserving exactly-once).
const vertexStreamEmitter = this.neurolink?.getEventEmitter();
if (vertexStreamEmitter) {
markStreamProviderEmittedGenerationEnd(
params.options as unknown as Parameters<
typeof markStreamProviderEmittedGenerationEnd
>[0],
);
vertexStreamEmitter.emit("generation:end", {
provider: this.providerName,
responseTime,
Expand Down Expand Up @@ -2371,7 +2382,7 @@
* Execute generate using native @google/genai SDK for Gemini 3 models on Vertex AI
* This bypasses @ai-sdk/google-vertex to properly handle thought_signature
*/
private async executeNativeGemini3Generate(

Check warning on line 2385 in src/lib/providers/googleVertex.ts

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'executeNativeGemini3Generate' has too many lines (333). Maximum allowed is 300

Check warning on line 2385 in src/lib/providers/googleVertex.ts

View workflow job for this annotation

GitHub Actions / test (20)

Async method 'executeNativeGemini3Generate' has too many lines (333). Maximum allowed is 300

Check warning on line 2385 in src/lib/providers/googleVertex.ts

View workflow job for this annotation

GitHub Actions / 🛡️ Code Quality & Security Gate

Async method 'executeNativeGemini3Generate' has too many lines (333). Maximum allowed is 300
options: TextGenerationOptions,
): Promise<EnhancedGenerateResult> {
const modelName = this.resolveAlias(
Expand Down
3 changes: 3 additions & 0 deletions src/lib/types/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,3 +63,6 @@ export * from "./elicitation.js";

// Dynamic Arguments types
export * from "./dynamic.js";

// Curator P2-4 dedup: per-stream AsyncLocalStorage context
export * from "./streamDedup.js";
12 changes: 12 additions & 0 deletions src/lib/types/streamDedup.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
/**
* Curator P2-4 dedup (concurrency-safe): per-stream context that lets
* the orchestration's `runStandardStreamRequest` finally block know
* whether a *native provider* path within THIS stream's async chain
* already emitted `generation:end`. Native providers (Vertex / Google
* AI Studio for Gemini 3, etc.) emit on the shared SDK emitter; without
* scoping, a concurrent unrelated stream's emit on the same NeuroLink
* instance would suppress the wrong stream's orchestration emit.
*
* AsyncLocalStorage scopes each stream's flag to its own async chain.
*/
export type StreamGenerationEndContext = { providerEmitted: boolean };
Loading
Loading