Repository navigation
fix(observability): address code review findings from PR #860 #875
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -179,8 +179,12 @@ export class MetricsAggregator { | |
| recordSpan(span: SpanData): void { | ||
| // Enforce maximum spans limit | ||
| if (this.spans.length >= this.config.maxSpansRetained) { | ||
| this.spans.shift(); // Remove oldest span | ||
| // Note: We keep aggregated metrics, only raw spans are trimmed | ||
| const evicted = this.spans.shift(); // Remove oldest span | ||
| // Only trim latencyValues when the evicted span had a duration recorded | ||
| if (evicted?.durationMs !== undefined) { | ||
| this.latencyValues.shift(); | ||
| } | ||
|
Comment on lines
181
to
+186
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Keep latency trimming tied to the evicted span, not only to max length. Current logic can leave stale latency samples when the removed span had 💡 Suggested fix- if (this.spans.length >= this.config.maxSpansRetained) {
- this.spans.shift(); // Remove oldest span
- // Trim latencyValues in sync to prevent unbounded memory growth
- if (this.latencyValues.length >= this.config.maxSpansRetained) {
- this.latencyValues.shift();
- }
+ if (this.spans.length >= this.config.maxSpansRetained) {
+ const removedSpan = this.spans.shift(); // Remove oldest span
+ // Trim latencyValues in sync with the removed span
+ if (removedSpan?.durationMs !== undefined && this.latencyValues.length > 0) {
+ this.latencyValues.shift();
+ }
// Note: We keep aggregated metrics, only raw spans and latency values are trimmed
}🤖 Prompt for AI Agents |
||
| // Note: We keep aggregated metrics, only raw spans and latency values are trimmed | ||
| } | ||
|
|
||
| this.spans.push(span); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -102,6 +102,8 @@ export class RedactionProcessor implements SpanProcessor { | |
| "credentials", | ||
| "private_key", | ||
| "privateKey", | ||
| "stack", | ||
| "error.stack", | ||
| ], | ||
| ); | ||
| this.redactedValue = config?.redactedValue ?? "[REDACTED]"; | ||
|
|
@@ -308,11 +310,24 @@ export class BatchProcessor implements SpanProcessor { | |
| }, this.flushIntervalMs); | ||
| } | ||
|
|
||
| // Note: flush() is intentionally synchronous. The onBatchReady callback is | ||
| // typed as `(spans: SpanData[]) => void` — callers must not pass async | ||
| // exporters. If async export is needed, the callback should handle its own | ||
| // error reporting (e.g. fire-and-forget with promise error handlers). | ||
| private flush(): void { | ||
| if (this.batch.length > 0 && this.onBatchReady) { | ||
| const spans = [...this.batch]; | ||
| this.batch = []; | ||
| this.onBatchReady(spans); | ||
| try { | ||
| this.onBatchReady(spans); | ||
| this.batch = []; | ||
| } catch (flushError: unknown) { | ||
| // Keep spans for next flush attempt, but cap backlog growth | ||
| void flushError; // acknowledged — error is expected during exporter outages | ||
| const maxBacklog = this.batchSize * 20; | ||
| if (this.batch.length > maxBacklog) { | ||
| this.batch = this.batch.slice(this.batch.length - maxBacklog); | ||
| } | ||
| } | ||
|
Comment on lines
+320
to
+330
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🧩 Analysis chain🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "1) Locate BatchProcessor constructions and onBatchReady assignments"
rg -nP --type=ts -C3 'new\s+BatchProcessor\s*\(|onBatchReady\s*:' src
echo
echo "2) Detect explicitly async onBatchReady callbacks"
rg -nP --type=ts -C2 'onBatchReady\s*:\s*async\b' src
echo
echo "Expected result:"
echo "- If any async onBatchReady exists, current flush() path can drop spans on rejection."
echo "- If none exist today, this is still a latent correctness bug for future callsites."Repository: juspay/neurolink Length of output: 127 🏁 Script executed: #!/bin/bash
set -euo pipefail
# Check if the file exists and get its size
echo "=== Checking file existence and size ==="
if [ -f "src/lib/observability/spanProcessor.ts" ]; then
wc -l src/lib/observability/spanProcessor.ts
echo "File exists"
else
echo "File not found at src/lib/observability/spanProcessor.ts"
# Try to find the file
echo "Searching for spanProcessor files..."
fd -e ts spanProcessor
fiRepository: juspay/neurolink Length of output: 155 🏁 Script executed: #!/bin/bash
set -euo pipefail
# Read lines around 316-325 from the spanProcessor file
echo "=== Code at lines 316-325 ==="
sed -n '310,335p' src/lib/observability/spanProcessor.ts
echo
echo "=== Looking for flush method definition ==="
rg -n 'flush\s*\(' src/lib/observability/spanProcessor.ts | head -5
echo
echo "=== Looking for onBatchReady definition ==="
rg -n 'onBatchReady' src/lib/observability/spanProcessor.ts | head -10Repository: juspay/neurolink Length of output: 1235 🏁 Script executed: #!/bin/bash
set -euo pipefail
# Search for BatchProcessor class definition and onBatchReady signature
echo "=== Search for class definition ==="
rg -n 'class\s+BatchProcessor' src/lib/observability/spanProcessor.ts
echo
echo "=== Get full context around flush method ==="
# Use ast-grep to find flush method with context
ast-grep --pattern $'private flush() {
$$$
}'
echo
echo "=== Search for any async callbacks ==="
rg -n 'async\s*\(' src/lib/observability/spanProcessor.ts | head -10Repository: juspay/neurolink Length of output: 200 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== Where is BatchProcessor instantiated? ==="
rg -n 'new\s+BatchProcessor' src --type=ts
echo
echo "=== All onBatchReady assignments/usages ==="
rg -n 'onBatchReady' src --type=ts -B2 -A2
echo
echo "=== Check if any async callbacks exist in the codebase ==="
# Look for patterns where async callbacks might be passed
rg -n 'async\s*\(\s*spans' src --type=ts
</parameter>Repository: juspay/neurolink Length of output: 105 🏁 Script executed: #!/bin/bash
set -euo pipefail
# Read the full config type and constructor
echo "=== Full BatchProcessor constructor and config type ==="
sed -n '272,295p' src/lib/observability/spanProcessor.ts
echo
echo "=== Check shutdown behavior ==="
sed -n '331,350p' src/lib/observability/spanProcessor.ts
</parameter>Repository: juspay/neurolink Length of output: 1515 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== Search for BatchProcessor usage across entire codebase ==="
rg 'BatchProcessor' src --type=ts
echo
echo "=== Search for any async function passed as callback ==="
rg 'onBatchReady\s*[:=]' src --type=ts -A1Repository: juspay/neurolink Length of output: 447 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== Full shutdown method ==="
sed -n '331,345p' src/lib/observability/spanProcessor.ts
echo
echo "=== Check SpanProcessor interface to understand expectations ==="
rg -n 'interface SpanProcessor' src --type=ts -A8Repository: juspay/neurolink Length of output: 1150 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== All flush() call sites ==="
sed -n '295,340p' src/lib/observability/spanProcessor.ts
echo
echo "=== Check if flush is called from timer/synchronous context ==="
rg -B5 'this.flush()' src/lib/observability/spanProcessor.tsRepository: juspay/neurolink Length of output: 2071 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== Verify TypeScript behavior: can async be assigned to void-returning type? ==="
echo "Checking if there are any TypeScript compiler directives or strict settings..."
fd -e json tsconfig | head -3 | xargs cat 2>/dev/null | head -30
echo
echo "=== Check if BatchProcessor is exported for public use ==="
rg -n 'export.*BatchProcessor' src --type=ts
</parameter>Repository: juspay/neurolink Length of output: 1307 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== Check the learning from previous PR about Promise handling ==="
echo "The learning suggests using .then()/.catch() for backward compatibility..."
echo ""
echo "=== Look for similar Promise handling patterns in the codebase ==="
rg 'onBatchReady\s*\(.*\)\.then|\.catch' src --type=ts
echo
echo "=== Check if there are examples of .then/.catch handling elsewhere ==="
rg '\.then\(|\.catch\(' src/lib/observability --type=ts -B2 -A2 | head -40Repository: juspay/neurolink Length of output: 5878 🏁 Script executed: #!/bin/bash
set -euo pipefail
echo "=== Verify: Can an async function be assigned to a void-returning type in TypeScript strict mode? ==="
cat > /tmp/test.ts << 'EOF'
// Test TypeScript assignability with strict mode
type SyncCallback = (data: string[]) => void;
type AsyncCallback = (data: string[]) => Promise<void>;
const asyncFn: AsyncCallback = async (data) => { console.log(data); };
const syncFn: SyncCallback = asyncFn; // Does this error in strict mode?
EOF
echo "TypeScript strict mode allows assigning async functions to void-returning types"
echo "This is because Promise<void> is compatible with void in JavaScript execution"
echo
echo "=== Confirm the actual issue: unhandled promise rejection scenario ==="
echo "Current code path:"
echo "1. Line 317: this.onBatchReady(spans) // If async, returns Promise but not awaited"
echo "2. Line 318: this.batch = [] // Clears batch immediately"
echo "3. If Promise rejects later: catch block doesn't execute, spans lost"Repository: juspay/neurolink Length of output: 629
Line 317 calls This is critical during Fix: Use Suggested fix (backward-compatible Promise handling) private flush(): void {
if (this.batch.length > 0 && this.onBatchReady) {
const spans = [...this.batch];
this.batch = [];
- try {
- this.onBatchReady(spans);
- } catch {
- // Keep spans for next flush attempt, but cap backlog growth
- const maxBacklog = this.batchSize * 20;
- if (this.batch.length > maxBacklog) {
- this.batch = this.batch.slice(this.batch.length - maxBacklog);
- }
+ const result = this.onBatchReady(spans);
+ if (result instanceof Promise) {
+ result.catch(() => {
+ // Keep spans for next flush attempt on rejection, but cap backlog growth
+ this.batch.unshift(...spans);
+ const maxBacklog = this.batchSize * 20;
+ if (this.batch.length > maxBacklog) {
+ this.batch = this.batch.slice(this.batch.length - maxBacklog);
+ }
+ });
+ }
}
}This keeps 🤖 Prompt for AI Agents |
||
| } | ||
| } | ||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,31 @@ | ||
| /** | ||
| * Safe metadata filtering for observability exporters. | ||
| * | ||
| * Only these attribute keys are forwarded to third-party backends as trace | ||
| * metadata. User prompts (input), LLM responses (output), error stacks, and | ||
| * any other potentially sensitive data are excluded to prevent PII leaks. | ||
| */ | ||
|
|
||
| import type { SpanAttributes } from "../types/spanTypes.js"; | ||
|
|
||
| // Only ai.* keys are forwarded as metadata. Stream metrics (chunk_count, | ||
| // content_length) should be accessed via span attributes directly, not via | ||
| // metadata sent to third-party backends. | ||
| export const SAFE_METADATA_KEYS = new Set([ | ||
| "ai.provider", | ||
| "ai.model", | ||
| "ai.temperature", | ||
| "ai.max_tokens", | ||
| ]); | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
|
|
||
| export function filterSafeMetadata( | ||
| attributes: SpanAttributes, | ||
| ): Record<string, unknown> { | ||
| const filtered: Record<string, unknown> = {}; | ||
| for (const key of SAFE_METADATA_KEYS) { | ||
| if (attributes[key] !== undefined) { | ||
| filtered[key] = attributes[key]; | ||
| } | ||
| } | ||
| return filtered; | ||
| } | ||
Uh oh!
There was an error while loading. Please reload this page.