Skip to content

refactor: extract stream core to open-sse/utils/stream/streamCore.ts (Issue #3594) - #3927

Closed
oyi77 wants to merge 3 commits into
diegosouzapw:release/v3.8.28from
oyi77:pr/stream-core
Closed

oyi77 wants to merge 3 commits into
diegosouzapw:release/v3.8.28from
oyi77:pr/stream-core

Conversation

@oyi77

@oyi77 oyi77 commented Jun 15, 2026

Copy link
Copy Markdown
Contributor

Part of modularization effort (Issue #3594).

Extracts core streaming logic from the monolithic open-sse/utils/stream.ts into a dedicated module.

Changes:

  • New file: open-sse/utils/stream/streamCore.ts (2215 lines)
  • Exports: createSSEStream, createSSETransformStreamWithLogger, createPassthroughStreamWithLogger

Dependencies: Requires the following modules (submitted as prior PRs):

Testing: All stream tests pass (49/49 stream-utils, 14/14 stream-handler).

Note: This is the largest module (~2215 lines). It's the core streaming engine and is cohesive. The prior extracted modules are its dependencies.

@oyi77
oyi77 requested a review from diegosouzapw as a code owner June 15, 2026 20:29

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Code Review

This pull request introduces streamCore.ts to handle unified SSE transform streams with idle timeout protection, supporting both translation and passthrough modes for various LLM formats. The code review identifies several critical issues, including a large number of missing imports that will cause runtime failures, as well as duplicate and unused imports. Additionally, the reviewer notes that asynchronous onFailure callbacks are wrapped in synchronous try/catch blocks (failing to catch rejections) and that errors in onComplete callbacks are being silently swallowed, which violates the repository's style guide against swallowing errors in SSE streams.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread open-sse/utils/stream/streamCore.ts Outdated
import { normalizeStreamFailurePayload } from "./errors.ts";
import { buildErrorBody } from "../error.ts";
import { SSEStreamContext } from "./types.ts";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

critical

Several functions, constants, and utilities referenced in this file are not imported, which will cause compilation/runtime errors. Please import the following:

  • sanitizeStreamingChunk
  • extractThinkingFromContent
  • hasValidUsage
  • estimateUsage
  • filterUsageForFormat
  • addBufferToUsage
  • consumeToolFinishTime
  • formatSSE
  • hasValuableContent
  • fixInvalidId
  • OMIT_STREAMING_CHUNK_MARKER
  • markToolFinish
  • unwrapGeminiChunk
  • logUsage
  • buildStreamSummaryFromEvents
  • parseTextualToolCallCandidate

Comment thread open-sse/utils/stream/streamCore.ts Outdated
Comment on lines +1 to +3
import { convertOpenAIToResponsesToolCall } from "../handlers/responseTranslator.ts";
import { translateResponse, initState } from "../../translator/index.ts";
import { v4 as uuidv4 } from "uuid";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

The imports convertOpenAIToResponsesToolCall and uuidv4 are not used anywhere in this file and can be safely removed.

Suggested change
import { convertOpenAIToResponsesToolCall } from "../handlers/responseTranslator.ts";
import { translateResponse, initState } from "../../translator/index.ts";
import { v4 as uuidv4 } from "uuid";
import { translateResponse, initState } from "../../translator/index.ts";

Comment on lines +10 to +14
import {
STREAM_IDLE_TIMEOUT_MS,
FETCH_BODY_TIMEOUT_MS,
HTTP_STATUS,
} from "../../config/constants.ts";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

The import FETCH_BODY_TIMEOUT_MS is not used in this file and can be safely removed.

import {
  STREAM_IDLE_TIMEOUT_MS,
  HTTP_STATUS,
} from "../../config/constants.ts";

Comment thread open-sse/utils/stream/streamCore.ts Outdated
} from "../../config/constants.ts";
import { parseSSELine } from "../streamHelpers.ts";
import { recordToolLatency } from "../../services/toolLatencyTracker.ts";
import { processBufferedPassthroughLine } from "../passthroughTailProcessor.ts";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

The import processBufferedPassthroughLine is not used in this file and can be safely removed.

Comment thread open-sse/utils/stream/streamCore.ts Outdated
Comment on lines +66 to +67
import { normalizeStreamFailurePayload } from "./errors.ts";
import { buildErrorBody } from "../error.ts";

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

These imports are duplicates of the imports on lines 18-19 and should be removed.

Comment on lines +1413 to +1417
if (ctx.onFailure) {
try {
void ctx.onFailure(failurePayload);
} catch {}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

Since ctx.onFailure can return a Promise<void>, wrapping it in a synchronous try/catch block without await will not catch asynchronous rejections. This can lead to unhandled promise rejections. Wrapping it in Promise.resolve().catch() safely catches any asynchronous errors.

Suggested change
if (ctx.onFailure) {
try {
void ctx.onFailure(failurePayload);
} catch {}
}
if (ctx.onFailure) {
Promise.resolve(ctx.onFailure(failurePayload)).catch((err) => {
console.error("[STREAM] Error in onFailure callback:", err);
});
}

Comment on lines +1885 to +1887
});
} catch {}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

Silently swallowing errors thrown inside the onComplete callback makes debugging extremely difficult. Logging the error to the console provides visibility when issues occur.

                });
              } catch (err) {
                console.error("[STREAM] Error in onComplete callback:", err);
              }
            }

Comment on lines +1945 to +1954
if (ctx.onFailure) {
try {
void ctx.onFailure({
status: err.status,
message: err.message,
code: err.code,
type: err.type,
});
} catch {}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

Since ctx.onFailure can return a Promise<void>, wrapping it in a synchronous try/catch block without await will not catch asynchronous rejections. Wrapping it in Promise.resolve().catch() safely catches any asynchronous errors.

Suggested change
if (ctx.onFailure) {
try {
void ctx.onFailure({
status: err.status,
message: err.message,
code: err.code,
type: err.type,
});
} catch {}
}
if (ctx.onFailure) {
Promise.resolve(
ctx.onFailure({
status: err.status,
message: err.message,
code: err.code,
type: err.type,
})
).catch((err) => {
console.error("[STREAM] Error in onFailure callback:", err);
});
}

Comment on lines +2139 to +2141
});
} catch {}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

Silently swallowing errors thrown inside the onComplete callback makes debugging extremely difficult. Logging the error to the console provides visibility when issues occur.

              });
            } catch (err) {
              console.error("[STREAM] Error in onComplete callback:", err);
            }
          }

@diegosouzapw
diegosouzapw changed the base branch from release/v3.8.26 to release/v3.8.27 June 16, 2026 06:19
oyi77 added a commit to oyi77/OmniRoute that referenced this pull request Jun 17, 2026
…error handling

Address review feedback on diegosouzapw#3927:
- Add 15 missing imports (sanitizeStreamingChunk, estimateUsage, etc.)
- Remove unused imports (convertOpenAIToResponsesToolCall, uuidv4, FETCH_BODY_TIMEOUT_MS, processBufferedPassthroughLine, duplicates)
- Fix async onFailure with Promise.resolve().catch() pattern
- Add console.error logging in onComplete catch blocks
@diegosouzapw
diegosouzapw changed the base branch from release/v3.8.27 to release/v3.8.28 June 17, 2026 08:20
oyi77 added 3 commits June 17, 2026 15:39
Part of Issue diegosouzapw#3594 modularization. Extracted core streaming logic from the monolithic stream.ts.

Exports:
- createSSEStream, createSSETransformStreamWithLogger, createPassthroughStreamWithLogger
- buildSSEStreamContext (internal)

This is the main streaming engine handling passthrough, translate, and Responses API modes.
Part of Issue diegosouzapw#3594 modularization. The core streaming module is 2216 lines and is added as a frozen entry since it's a cohesive core module that cannot be easily split further.
…error handling

Address review feedback on diegosouzapw#3927:
- Add 15 missing imports (sanitizeStreamingChunk, estimateUsage, etc.)
- Remove unused imports (convertOpenAIToResponsesToolCall, uuidv4, FETCH_BODY_TIMEOUT_MS, processBufferedPassthroughLine, duplicates)
- Fix async onFailure with Promise.resolve().catch() pattern
- Add console.error logging in onComplete catch blocks
@diegosouzapw

Copy link
Copy Markdown
Owner

Thanks, @oyi77 🙏. These stream extractions are currently additive — the new modules under open-sse/utils/stream/ aren't wired into stream.ts yet, so they don't change runtime behavior on their own. Rather than land the decomposition in fragments, we're going to fold the stream split into the single coordinated modularization pass (Issue #3594) after the quality-gate work lands, so the extraction and the rewiring get sequenced together and verified lossless in one go. Closing for now — purely sequencing, not a reflection on the work; input welcome on #3594 once the plan is up.

@oyi77
oyi77 deleted the pr/stream-core branch August 7, 2026 21:09
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants