Skip to content

feat: plugins async fixes and added streaming to maxim plugin - #681

Merged
akshaydeo merged 1 commit into
mainfrom
10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin
Oct 25, 2025
Merged

akshaydeo merged 1 commit into
mainfrom
10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin

Conversation

@Pratham-Mishra04

Copy link
Copy Markdown
Collaborator

Summary

Improved method naming consistency and implemented deep copying for streaming responses to prevent data mutation between plugin accumulators.

Changes

  • Renamed conversion methods in provider schemas to be more specific and consistent (e.g., ToBifrostResponse → ToBifrostTextCompletionResponse)
  • Added deep copy functionality for ResponsesStreamResponse to prevent shared data mutation between different plugin accumulators
  • Fixed streaming response handling in the Maxim plugin to properly support streaming requests
  • Improved concurrency handling in plugins (Logging, Governance, OTEL) by moving response processing to goroutines
  • Updated the Anthropic provider to use the renamed conversion method

Type of change

  • Bug fix
  • Feature
  • Refactor
  • Documentation
  • Chore/CI

Affected areas

  • Core (Go)
  • Transports (HTTP)
  • Providers/Integrations
  • Plugins
  • UI (Next.js)
  • Docs

How to test

Test streaming responses with multiple plugins enabled:

# Core/Transports
go version
go test ./...

# Test with streaming requests
curl -X POST http://localhost:8080/v1/chat/completions \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer $API_KEY" \
  -d '{
    "model": "gpt-3.5-turbo",
    "messages": [{"role": "user", "content": "Hello"}],
    "stream": true
  }'

Breaking changes

  • Yes
  • No

Related issues

Fixes issues with data mutation between plugin accumulators during streaming responses.

Security considerations

No security implications.

Checklist

  • I added/updated tests where appropriate
  • I verified builds succeed (Go and UI)
  • I verified the CI pipeline passes locally if applicable

@coderabbitai

coderabbitai Bot commented Oct 25, 2025 •

Copy link
Copy Markdown
Contributor
📝 Walkthrough

Summary by CodeRabbit

  • New Features

    • Added support for Responses type in streaming operations.
  • Bug Fixes

    • Enhanced streaming response handling with improved copy mechanisms to prevent data conflicts.
    • Fixed streaming chunk ordering for consistent message construction.
  • Improvements

    • Optimized asynchronous processing in logging and observability pipelines.
    • Enhanced resource cleanup and error handling across plugins.

Walkthrough

Renamed many provider conversion methods to type-specific names (chat/text/embedding), updated all call sites, added deep-copy helpers and stable ordering for streaming responses, and refactored several plugins and transports to use asynchronous PostHook flows and a new Maxim Init/PostHook signature. (50 words)

Changes

Cohort / File(s) Summary
Schema method renames — chat
core/schemas/providers/anthropic/chat.go, core/schemas/providers/bedrock/chat.go, core/schemas/providers/cohere/chat.go
Renamed chat conversion methods from ToBifrostResponse() → ToBifrostChatResponse() and updated comments; behavior unchanged.
Schema method renames — text completion
core/schemas/providers/anthropic/text.go, core/schemas/providers/bedrock/text.go
Renamed text methods: ToBifrostRequest() → ToBifrostTextCompletionRequest() and ToBifrostResponse() → ToBifrostTextCompletionResponse(); call sites updated.
Schema method renames — embedding
core/schemas/providers/cohere/embedding.go, core/schemas/providers/vertex/embedding.go
Renamed ToBifrostResponse() → ToBifrostEmbeddingResponse(); comments updated, logic unchanged.
OpenAI schema renames & callsite updates
core/schemas/providers/openai/*.go, transports/bifrost-http/integrations/openai.go
Renamed multiple ToBifrostRequest → type-specific conversions (Chat/Responses/Embedding/Speech/Text/Transcription) and updated all call sites.
Anthropic provider callsites
core/providers/anthropic.go, transports/bifrost-http/integrations/anthropic.go
Call sites switched to ToBifrostTextCompletionRequest() / ToBifrostTextCompletionResponse().
Streaming framework enhancements
framework/streaming/responses.go, framework/streaming/types.go
Added deep-copy helpers, store deep copies for non-OpenAI providers, sort accumulated chunks by sequence number, and added StreamTypeResponses handling in ToBifrostResponse.
Maxim plugin API & wiring
plugins/maxim/main.go, plugins/maxim/plugin_test.go, transports/bifrost-http/handlers/server.go
Init(config) → Init(config, logger); Plugin made public with logger, accumulator, mutex fields; PreHook/PostHook/Cleanup updated for streaming; tests and server Init call updated.
Logging plugin async refactor
plugins/logging/main.go
PostHook work moved to background goroutine; streaming-aware accumulator handling added; log updates deferred, enriched, and retried; cost calculation and pooled data propagation added.
Governance plugin concurrency
plugins/governance/main.go
Added sync.WaitGroup field; PostHook starts worker goroutine and delegates audit extraction via headers; Cleanup waits for workers.
OTEL plugin async emission
plugins/otel/main.go, plugins/otel/docker-compose.yml
PostHook emission moved to goroutine tracked by waitgroup; streaming path uses accumulator; non-streaming constructs final ResourceSpan before emit; minor docker-compose formatting tweaks.
Transport layer adjustments
transports/bifrost-http/handlers/server.go
Updated Maxim Init call to pass new logger argument.

Sequence Diagram(s)

sequenceDiagram
    participant Client
    participant Core
    participant Provider
    participant Schema
    participant PostHook as Async PostHook

    Client->>Core: Send request
    Core->>Provider: Call provider API
    Provider-->>Schema: Provider response object
    Schema->>Core: Convert via ToBifrostTextCompletionResponse/ToBifrostChatResponse/ToBifrostEmbeddingResponse
    Core->>Client: Return Bifrost response

    Client->>Core: triggers PostHook(result)
    Core->>PostHook: spawn goroutine (async)
    PostHook->>PostHook: compute usage/cost, prepare logs/spans
    alt streaming
      PostHook->>Accumulator: accumulate/process stream chunks
      Accumulator-->>PostHook: final aggregated result
    end
    PostHook->>External: emit/persist/log (async)
Loading
sequenceDiagram
    participant StreamSrc
    participant Framework
    participant Accumulator
    participant Builder

    StreamSrc->>Framework: send stream chunk (may be unordered)
    Framework->>Framework: deep-copy chunk (non-OpenAI providers)
    Framework->>Accumulator: append chunk
    Framework->>Framework: sort accumulated chunks by sequence
    Framework->>Builder: build complete message from ordered chunks
    Builder-->>Accumulator: final constructed message
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~45 minutes

  • Files/areas needing extra attention:
    • plugins/logging/main.go — concurrency, async error handling, accumulator lifecycle, retry/update logic, and pooled resource management.
    • plugins/maxim/main.go & tests — new Init signature, logger & accumulator wiring, test updates.
    • plugins/otel/main.go — emit waitgroup correctness and span lifecycle.
    • framework/streaming/responses.go & types.go — deep-copy correctness, memory/copy semantics, and chunk ordering.
    • Global search for missed call sites after numerous schema method renames.

Poem

🐰 I hopped through structs and renamed a few,
Chat, text, embed — now clearer in view,
Streams copied safe, then ordered with care,
Goroutines hum and quiet the air,
A rabbit nods — the repo breathes anew.

Pre-merge checks and finishing touches

✅ Passed checks (3 passed)
Check name Status Explanation
Title Check ✅ Passed The PR title "feat: plugins async fixes and added streaming to maxim plugin" refers to real aspects of the changeset, specifically the async concurrency improvements in plugins and the new streaming support in the Maxim plugin. However, it does not capture the most substantial portion of the changes: the systematic renaming of conversion methods across provider schemas (Anthropic, Bedrock, Cohere, Vertex, OpenAI). While the title is specific enough for the async and streaming work, it omits the primary refactoring effort that affects numerous files and provides consistency improvements to the conversion method naming convention.
Description Check ✅ Passed The PR description follows the provided template structure and includes all required sections: a clear summary explaining the method naming consistency and deep copying improvements, a detailed changes section with bullet points, type of change selections (Bug fix and Refactor), affected areas checkboxes, testing instructions with concrete examples, a security considerations statement, and a completed checklist. The description effectively communicates the purpose and scope of the changes, covering the main work items identified in the raw summary including method renaming, deep copy functionality, streaming in Maxim, and concurrency improvements across plugins.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
✨ Finishing touches
  • 📝 Generate docstrings
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Post copyable unit tests in a comment
  • Commit unit tests in branch 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin

Comment @coderabbitai help to get the list of available commands and usage tips.

Pratham-Mishra04 commented Oct 25, 2025 •

Copy link
Copy Markdown
Collaborator Author

@coderabbitai coderabbitai 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.

Actionable comments posted: 6

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (3)
core/schemas/providers/anthropic/chat.go (1)

991-1007: Nil-deref risk in error conversion.

ToAnthropicChatCompletionError dereferences bifrostErr.Error without a nil check. This can panic when Error is nil.

Apply a defensive fix:

 func ToAnthropicChatCompletionError(bifrostErr *schemas.BifrostError) *AnthropicMessageError {
   if bifrostErr == nil {
     return nil
   }

   // Provide blank strings for nil pointer fields
   errorType := ""
   if bifrostErr.Type != nil {
     errorType = *bifrostErr.Type
   }

-  // Handle nested error fields with nil checks
-  errorStruct := AnthropicMessageErrorStruct{
-    Type:    errorType,
-    Message: bifrostErr.Error.Message,
-  }
+  // Handle nested error fields with nil checks
+  msg := ""
+  if bifrostErr.Error != nil && bifrostErr.Error.Message != nil {
+    msg = *bifrostErr.Error.Message
+  }
+  errorStruct := AnthropicMessageErrorStruct{
+    Type:    errorType,
+    Message: msg,
+  }

   return &AnthropicMessageError{
     Type:  "error", // always "error" for Anthropic
     Error: errorStruct,
   }
 }
transports/bifrost-http/integrations/anthropic.go (1)

59-61: Add ToAnthropicResponsesError converter for non-streaming responses/messages errors.

The responses/messages route (line 42-69) uses response-specific converters for success responses (ToAnthropicResponsesResponse) and streaming errors (ToAnthropicResponsesStreamError), but non-streaming errors incorrectly use the chat-specific ToAnthropicChatCompletionError. This envelope mismatch will produce incorrect error formats for /v1/messages requests.

Create ToAnthropicResponsesError in core/schemas/providers/anthropic/responses.go (mirroring the non-SSE pattern of ToAnthropicResponsesStreamError) and update line 60 to use it:

-            ErrorConverter: func(err *schemas.BifrostError) interface{} {
-                return anthropic.ToAnthropicChatCompletionError(err)
-            },
+            ErrorConverter: func(err *schemas.BifrostError) interface{} {
+                return anthropic.ToAnthropicResponsesError(err)
+            },
framework/streaming/types.go (1)

167-189: Nil-safety: guard p.Data.OutputMessage before deref in Chat path.

This can panic when OutputMessage is nil.

Apply:

- if p.Data.OutputMessage.Content.ContentStr != nil {
+ if p.Data.OutputMessage != nil && p.Data.OutputMessage.Content != nil && p.Data.OutputMessage.Content.ContentStr != nil {
   ...
 }
- if p.Data.OutputMessage.ChatAssistantMessage != nil {
+ if p.Data.OutputMessage != nil && p.Data.OutputMessage.ChatAssistantMessage != nil {
   ...
 }
🧹 Nitpick comments (8)
transports/bifrost-http/integrations/anthropic.go (1)

39-41: Confirm text-completion error envelope.

For /v1/complete, ErrorConverter uses ToAnthropicChatCompletionError. If there’s a text-specific error mapper (e.g., ToAnthropicTextCompletionError), prefer that for consistency; otherwise confirm chat error mapping is intentional.

plugins/logging/main.go (3)

312-340: Also return pooled UpdateLogData (and cache debug) on error path.

Use the pool for UpdateLogData here and set SemanticCacheDebug from the error, mirroring success path.

-    logMsg.Operation = LogOperationUpdate
-    logMsg.UpdateData = &UpdateLogData{
-        Status:       "error",
-        ErrorDetails: bifrostErr,
-    }
+    logMsg.Operation = LogOperationUpdate
+    updateData := p.getUpdateLogData()
+    updateData.Status = "error"
+    updateData.ErrorDetails = bifrostErr
+    logMsg.UpdateData = updateData
+    if bifrostErr != nil {
+        logMsg.SemanticCacheDebug = bifrostErr.ExtraFields.CacheDebug
+    }
+    defer func() {
+        if logMsg.UpdateData != nil { p.putUpdateLogData(logMsg.UpdateData) }
+    }()

451-457: Remove late defer after adding early pool return.

Once the early defer is added, this defer (and the similar one in the streaming-final block) should be removed to avoid double put.

-    defer p.putLogMessage(logMsg) // Return to pool when done
     // Return pooled data structures to their respective pools
     defer func() {
         if logMsg.UpdateData != nil {
             p.putUpdateLogData(logMsg.UpdateData)
         }
     }()

330-337: Call callbacks without holding the mutex.

Lock only to read/copy the callback, then invoke it after unlocking to avoid blocking other log events on store I/O and user callbacks.

-    p.mu.Lock()
-    if p.logCallback != nil {
-        if updatedEntry, getErr := p.getLogEntry(p.ctx, logMsg.RequestID); getErr == nil {
-            p.logCallback(updatedEntry)
-        }
-    }
-    p.mu.Unlock()
+    p.mu.Lock()
+    cb := p.logCallback
+    p.mu.Unlock()
+    if cb != nil {
+        if updatedEntry, getErr := p.getLogEntry(p.ctx, logMsg.RequestID); getErr == nil {
+            cb(updatedEntry)
+        }
+    }

Also applies to: 360-367, 472-481

plugins/governance/main.go (1)

468-496: Success should be derived from bifrostErr, not just result presence.

Errors can arrive with partial results; using result != nil can mark failures as success. Compute success in PostHook and pass it into the worker.

Apply this minimal change:

- go p.postHookWorker(result, provider, model, requestType, virtualKey, requestID, headers, isCacheRead, isBatch, bifrost.IsFinalChunk(ctx))
+ success := (err == nil)
+ go p.postHookWorker(result, provider, model, requestType, virtualKey, requestID, headers, isCacheRead, isBatch, bifrost.IsFinalChunk(ctx), success)

And update worker signature/body:

-func (p *GovernancePlugin) postHookWorker(result *schemas.BifrostResponse, provider schemas.ModelProvider, model string, requestType schemas.RequestType, virtualKey, requestID string, headers map[string]string, isCacheRead, isBatch bool, isFinalChunk bool) {
-  // Determine if request was successful
-  success := (result != nil)
+func (p *GovernancePlugin) postHookWorker(result *schemas.BifrostResponse, provider schemas.ModelProvider, model string, requestType schemas.RequestType, virtualKey, requestID string, headers map[string]string, isCacheRead, isBatch bool, isFinalChunk bool, success bool) {
plugins/maxim/main.go (2)

391-416: Don’t fail the request if logging trace creation fails.

Logging should be best-effort. Returning an error from PreHook can block requests.

Apply:

- logger, err := plugin.getOrCreateLogger(effectiveLogRepoID)
- if err != nil {
-   return req, nil, fmt.Errorf("failed to create trace: %w", err)
- }
+ logger, err := plugin.getOrCreateLogger(effectiveLogRepoID)
+ if err != nil {
+   plugin.logger.Warn("%s failed to get/create logger for repo %s: %v", PluginLoggerPrefix, effectiveLogRepoID, err)
+   return req, nil, nil
+ }

488-491: Log logger creation errors in PostHook.

Silent return hides diagnostics.

Apply:

- if err != nil {
-   return
- }
+ if err != nil {
+   plugin.logger.Error("%s failed to get/create logger: %v", PluginLoggerPrefix, err)
+   return
+ }
framework/streaming/types.go (1)

198-214: Set Responses CreatedAt and optional ID for completeness.

BifrostResponsesResponse.CreatedAt is non-omitempty; setting it improves consistency. Including ID aids correlation.

Apply:

- responsesResp := &schemas.BifrostResponsesResponse{}
+ responsesResp := &schemas.BifrostResponsesResponse{
+   CreatedAt: int(p.Data.EndTimestamp.Unix()),
+ }
+ // Optionally set ID to requestID for traceability
+ if p.RequestID != "" {
+   responsesResp.ID = &p.RequestID
+ }
📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 7f62877 and 23f8a19.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (2 hunks)
🧰 Additional context used
🧬 Code graph analysis (10)
framework/streaming/types.go (2)
core/schemas/responses.go (1)
  • BifrostResponsesResponse (40-72)
core/schemas/bifrost.go (3)
  • BifrostResponseExtraFields (251-259)
  • RequestType (79-79)
  • ResponsesRequest (86-86)
transports/bifrost-http/handlers/server.go (1)
plugins/maxim/main.go (1)
  • Init (62-92)
transports/bifrost-http/integrations/anthropic.go (2)
transports/bifrost-http/integrations/utils.go (1)
  • RouteConfigTypeAnthropic (197-197)
core/schemas/bifrost.go (1)
  • TextCompletionRequest (82-82)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (3)
  • BifrostLLMUsage (547-553)
  • ChatMessage (336-345)
  • ChatMessageContent (348-351)
plugins/governance/main.go (2)
core/utils.go (1)
  • IsFinalChunk (130-145)
core/schemas/bifrost.go (3)
  • BifrostResponse (213-223)
  • ModelProvider (32-32)
  • RequestType (79-79)
core/schemas/providers/anthropic/text.go (2)
core/schemas/providers/anthropic/types.go (2)
  • AnthropicTextRequest (18-27)
  • AnthropicTextResponse (205-214)
core/schemas/textcompletions.go (2)
  • BifrostTextCompletionRequest (10-16)
  • BifrostTextCompletionResponse (64-72)
plugins/maxim/main.go (6)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/plugin.go (1)
  • Plugin (45-71)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (82-82)
  • TextCompletionStreamRequest (83-83)
  • ChatCompletionRequest (84-84)
  • ChatCompletionStreamRequest (85-85)
  • ResponsesRequest (86-86)
  • ResponsesStreamRequest (87-87)
  • BifrostContextKeyRequestID (101-101)
  • RequestType (79-79)
  • BifrostResponse (213-223)
  • BifrostError (319-328)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (2)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
framework/streaming/responses.go (2)
core/schemas/responses.go (5)
  • BifrostResponsesStreamResponse (1334-1367)
  • BifrostResponsesResponse (40-72)
  • ResponsesOutputMessageContentTextLogProb (425-430)
  • ResponsesOutputMessageContentText (406-409)
  • ResponsesOutputMessageContentRefusal (431-433)
core/schemas/bifrost.go (3)
  • OpenAI (35-35)
  • OpenRouter (48-48)
  • Azure (36-36)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (14)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (27)
plugins/otel/docker-compose.yml (1)

10-15: Formatting-only changes; YAML remains valid.

All changes in this file are cosmetic formatting adjustments (inlined comments, normalized spacing, command array reformatting) that preserve YAML validity and service configuration. These appear incidental to the main PR objectives (async/streaming plugin fixes).

Consider separating minor infrastructure formatting cleanups into a separate PR for clarity, but this is not a blocker.

Also applies to: 24-24, 29-29, 31-31, 229-229

plugins/otel/main.go (1)

78-79: LGTM on adding emitWg

Adding a WaitGroup to coordinate emissions makes sense. Usage details are addressed below.

transports/bifrost-http/integrations/anthropic.go (1)

30-31: Rename to ToBifrostTextCompletionRequest looks correct.

Call site aligns with specialized request type.

core/schemas/providers/anthropic/chat.go (1)

250-252: Rename to ToBifrostChatResponse is good.

Signature and behavior remain intact; naming aligns with API surface.

Please ensure all call sites were updated (grep for “ToBifrostResponse(”).

plugins/maxim/plugin_test.go (2)

32-37: Init now requires a logger — test helper change looks correct.

Creating a default logger and passing it into Init aligns with the updated API.


198-236: Init signature update wired through tests.

Logger is created once and reused across subtests; straightforward and correct.

transports/bifrost-http/handlers/server.go (1)

214-215: Server wiring updated for new Init signature.

Passing logger into maxim.Init is correct and consistent with plugin API.

plugins/logging/main.go (2)

120-125: Go version confirmed: compatible with for range syntax.

The repository's go.mod declares go 1.24, which exceeds the Go 1.22 requirement for for range over integers. The code is valid as-is; no refactoring needed.


341-369: Confirm whether filtering intermediate streaming updates is intentional.

The code definitively processes only final streaming responses (streamResponse.Type == streaming.StreamResponseTypeFinal). Intermediate delta updates are filtered out across logging, otel, and maxim plugins. This pattern is consistent but lacks design documentation.

If this is intentional for performance (reduced DB writes for partial updates), please document it in code comments. If live UI streaming updates are expected, apply the refactor to process all streaming responses:

-    } else if streamResponse != nil && streamResponse.Type == streaming.StreamResponseTypeFinal {
+    } else if streamResponse != nil {
         // Prepare final log data
         logMsg.Operation = LogOperationStreamUpdate
         logMsg.StreamResponse = streamResponse
         defer p.putLogMessage(logMsg) // Return to pool when done
         processingErr := retryOnNotFound(p.ctx, func() error {
             return p.updateStreamingLogEntry(p.ctx, logMsg.RequestID, logMsg.SemanticCacheDebug, logMsg.StreamResponse, streamResponse.Type == streaming.StreamResponseTypeFinal)
         })
plugins/governance/main.go (2)

435-436: Good async offload in PostHook.

Moving usage/cost work to a goroutine keeps the response path non-blocking. LGTM.


4-17: Go version requirement satisfied—no action needed.

The plugins/governance/go.mod specifies go 1.24.1, which exceeds the Go 1.22+ requirement for the math/rand/v2 import. The code change is compatible with the existing toolchain configuration.

plugins/maxim/main.go (2)

62-92: Init wiring looks good; accumulator and logger injected.

The new Init(config, logger) signature and accumulator setup are sensible for stream-aware logging. LGTM.


540-543: Accumulator cleanup on plugin shutdown looks good.

Proactive cleanup prevents background goroutines from leaking.

core/schemas/providers/cohere/chat.go (1)

183-189: Rename to ToBifrostChatResponse is consistent and non-breaking.

Conversion logic and ExtraFields are intact. LGTM.

core/schemas/providers/vertex/embedding.go (1)

69-71: Rename to ToBifrostEmbeddingResponse aligns with naming standard.

Implementation unchanged; output is correct. LGTM.

core/schemas/providers/bedrock/chat.go (1)

43-44: LGTM! Method rename improves API clarity.

The rename from ToBifrostResponse to ToBifrostChatResponse makes the conversion target explicit and aligns with the broader API naming standardization across providers.

core/providers/anthropic.go (1)

220-220: LGTM! Correctly uses the renamed conversion method.

The call site properly updated to use ToBifrostTextCompletionResponse(), matching the method rename in the Anthropic text schema.

core/schemas/providers/bedrock/text.go (2)

60-61: LGTM! Method rename aligns with naming conventions.

The rename to ToBifrostTextCompletionResponse explicitly indicates this converts to a text completion response format.


84-85: LGTM! Consistent method rename for Mistral text responses.

Matches the naming convention applied to the Anthropic text response converter above.

framework/streaming/responses.go (5)

13-136: Excellent addition of deep copy functionality to prevent shared mutation.

The deep copy implementation comprehensively handles all fields in the streaming response structure, addressing the data mutation issue mentioned in the PR objectives.


227-233: LGTM! Sorting ensures deterministic chunk processing.

Adding sort by sequence number before processing chunks improves consistency and prevents order-dependent bugs.


489-491: LGTM! Provider check uses constants and covers multiple providers.

Using provider constants (schemas.OpenAI, schemas.OpenRouter, schemas.Azure) instead of string literals is type-safe and properly handles all OpenAI-compatible providers that send complete responses in the final chunk.


561-562: LGTM! Deep copy prevents mutation between plugin accumulators.

Replacing direct assignment with deepCopyResponsesStreamResponse() ensures each plugin accumulator gets its own independent copy of the stream response.


78-83: The shallow copy concern is valid but lacks evidence of actual harm in current usage.

The struct ResponsesOutputMessageContentTextLogProb does contain slice fields (Bytes and TopLogProbs), and the current implementation at lines 78-83 performs a shallow copy rather than deep copy. However, codebase analysis found no mutations to these nested slices—all LogProbs usage is read-only (metadata extraction, telemetry, response mapping). The comment label "Deep copy LogProbs slice" is also misleading since the implementation is actually a shallow copy of elements.

If the "plugin accumulators" design guarantees that these nested slices remain immutable after copying, the current code is safe. However, if future code could modify these nested structures, deep copying would be necessary.

core/schemas/providers/anthropic/text.go (2)

48-49: LGTM! Method rename clarifies conversion purpose.

Renaming to ToBifrostTextCompletionRequest makes it explicit that this converts to a text completion request format.


80-81: LGTM! Consistent method rename for response conversion.

The rename to ToBifrostTextCompletionResponse mirrors the request conversion rename and maintains naming consistency.

core/schemas/providers/cohere/embedding.go (1)

64-65: LGTM! Method rename specifies embedding conversion.

The rename to ToBifrostEmbeddingResponse explicitly indicates this converts to an embedding response format, consistent with the PR's naming standardization effort.

Comment thread plugins/governance/main.go
Comment thread plugins/logging/main.go
Comment thread plugins/logging/main.go
Comment thread plugins/maxim/main.go
Comment thread plugins/maxim/main.go
Comment thread plugins/otel/main.go
@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin branch from 23f8a19 to 0e55852 Compare October 25, 2025 14:02

@coderabbitai coderabbitai 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.

Actionable comments posted: 0

♻️ Duplicate comments (1)
plugins/otel/main.go (1)

230-234: Handle Emit errors on streaming-final path

Emit errors are ignored here; the non-stream path logs them. Add consistent error handling to avoid silent drops.

Apply this diff:

-                if streamResponse != nil && streamResponse.Type == streaming.StreamResponseTypeFinal {
-                    defer p.ongoingSpans.Delete(traceID)
-                    p.client.Emit(p.ctx, []*ResourceSpan{completeResourceSpan(span, time.Now(), streamResponse.ToBifrostResponse(), bifrostErr, p.pricingManager)})
-                }
+                if streamResponse != nil && streamResponse.Type == streaming.StreamResponseTypeFinal {
+                    defer p.ongoingSpans.Delete(traceID)
+                    if err := p.client.Emit(p.ctx, []*ResourceSpan{
+                        completeResourceSpan(span, time.Now(), streamResponse.ToBifrostResponse(), bifrostErr, p.pricingManager),
+                    }); err != nil {
+                        logger.Error("failed to emit response span for request %s: %v", traceID, err)
+                    }
+                }
#!/bin/bash
# Verify all Otel client emissions are error-checked.
rg -n -C2 -g '!**/vendor/**' -P '\.Emit\('
🧹 Nitpick comments (1)
plugins/otel/main.go (1)

78-79: Minor: redundant WaitGroup initialization

The zero value of sync.WaitGroup is usable; explicit initialization in the struct literal isn’t needed. Safe to drop for brevity.

Apply this diff:

@@
 type OtelPlugin struct {
@@
-    emitWg sync.WaitGroup // Track in-flight emissions
+    emitWg sync.WaitGroup // Track in-flight emissions
 }
@@
-        emitWg:         sync.WaitGroup{},

Also applies to: 108-109

📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 23f8a19 and 0e55852.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (2 hunks)
🚧 Files skipped from review as they are similar to previous changes (8)
  • framework/streaming/types.go
  • core/schemas/providers/anthropic/chat.go
  • core/schemas/providers/cohere/chat.go
  • framework/streaming/responses.go
  • core/schemas/providers/anthropic/text.go
  • core/schemas/providers/bedrock/chat.go
  • core/schemas/providers/vertex/embedding.go
  • transports/bifrost-http/integrations/anthropic.go
🧰 Additional context used
🧬 Code graph analysis (6)
plugins/logging/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
plugins/governance/main.go (2)
core/utils.go (1)
  • IsFinalChunk (130-145)
core/schemas/bifrost.go (3)
  • BifrostResponse (213-223)
  • ModelProvider (32-32)
  • RequestType (79-79)
plugins/maxim/main.go (5)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (82-82)
  • TextCompletionStreamRequest (83-83)
  • ChatCompletionRequest (84-84)
  • ChatCompletionStreamRequest (85-85)
  • ResponsesRequest (86-86)
  • ResponsesStreamRequest (87-87)
  • BifrostContextKeyRequestID (101-101)
  • RequestType (79-79)
  • BifrostResponse (213-223)
  • BifrostError (319-328)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (1)
  • ProcessedStreamResponse (107-114)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
transports/bifrost-http/handlers/server.go (1)
plugins/maxim/main.go (1)
  • Init (62-92)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (19)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (13)
plugins/otel/docker-compose.yml (1)

10-16: Cosmetic formatting improvements approved.

All changes are consistent comment alignment and formatting adjustments. Port mappings, service configurations, and functional behavior remain unchanged. These improvements enhance code readability and consistency.

Also applies to: 24-24, 29-29, 31-31, 229-229

plugins/otel/main.go (1)

214-217: WaitGroup placement and lifecycle look correct

Add happens before launching the goroutine and Done is deferred inside. This fixes the earlier race/missed-tracking risk and ensures Cleanup waits for all emissions.

core/providers/anthropic.go (1)

220-220: LGTM! Method rename improves clarity.

The rename from ToBifrostResponse() to ToBifrostTextCompletionResponse() makes the conversion target type explicit and aligns with the broader naming consistency improvements across provider schemas.

core/schemas/providers/bedrock/text.go (1)

60-61: LGTM! Consistent naming improvements.

Both method renames from ToBifrostResponse() to ToBifrostTextCompletionResponse() improve clarity by making the target response type explicit for both Bedrock Anthropic and Mistral text responses.

Also applies to: 84-85

core/schemas/providers/cohere/embedding.go (1)

64-65: LGTM! Method rename clarifies conversion target.

The rename from ToBifrostResponse() to ToBifrostEmbeddingResponse() makes the embedding-specific conversion explicit and consistent with the broader naming improvements.

plugins/logging/main.go (1)

299-482: LGTM! Well-structured async refactoring.

The PostHook logic has been successfully moved to an asynchronous goroutine, which prevents blocking the request thread. The implementation correctly:

  • Manages pooled resources with defers (line 304 for LogMessage, lines 452-456 for UpdateData)
  • Handles three distinct paths: errors, streaming responses, and non-streaming responses
  • Cleans up streaming accumulators on errors (line 316)
  • Invokes callbacks with updated log entries after successful updates

The latency handling (lines 306-310) correctly defaults to 0 when result is nil, which is appropriate given that BifrostErrorExtraFields does not contain a latency field.

Based on learnings.

plugins/governance/main.go (1)

456-463: Verify case-insensitive header lookup for audit fields.

The current implementation uses case-sensitive map access for x-bf-team and x-bf-customer headers (lines 458, 461). HTTP headers are case-insensitive per RFC 7230, so these lookups may miss audit fields if headers arrive with different casing (e.g., X-Bf-Team, X-BF-TEAM).

Consider implementing case-insensitive header lookup:

// Helper function for case-insensitive header lookup
getHeader := func(m map[string]string, key string) string {
	if m == nil {
		return ""
	}
	lk := strings.ToLower(key)
	for k, v := range m {
		if strings.ToLower(k) == lk {
			return v
		}
	}
	return ""
}

// Extract team/customer info for audit trail
var teamID, customerID *string
if v := getHeader(headers, "x-bf-team"); v != "" {
	teamID = &v
}
if v := getHeader(headers, "x-bf-customer"); v != "" {
	customerID = &v
}

Note: This was flagged in a previous review and remains unaddressed.

plugins/maxim/main.go (4)

62-79: LGTM! Init signature updated to support streaming.

The Init function now accepts a logger parameter and properly initializes the streaming accumulator (line 78) and stores the logger (line 79). This enables the plugin to support streaming request/response logging with proper instrumentation.


423-433: LGTM! Streaming accumulator setup is correct.

PreHook now properly extracts the requestID and creates a stream accumulator for streaming requests. The defensive error logging on line 427 is appropriate, though per maintainer feedback, requestID is always set at the gateway level.

Based on learnings.


470-536: LGTM! PostHook async refactoring supports streaming correctly.

The PostHook logic has been successfully moved to an asynchronous goroutine, enabling:

  • Non-blocking request processing
  • Streaming response accumulation via accumulator.ProcessStreamingResponse() (line 476)
  • Early return for delta chunks to avoid redundant logging (lines 483-485)
  • Proper use of accumulated streaming data when finalizing generations (lines 507-525)
  • Cleanup of streaming accumulators on errors (line 503)

The implementation correctly uses streamResponse.ToBifrostResponse() to retrieve accumulated results for final logging.


269-270: LGTM! Streaming request types properly handled.

The PreHook now handles streaming request types (TextCompletionStreamRequest, ChatCompletionStreamRequest, ResponsesStreamRequest) in addition to their non-streaming counterparts. The PostHook switch statement (lines 507-525) correctly covers all these types when logging results.

Also applies to: 291-292, 325-326, 507-525

plugins/maxim/plugin_test.go (1)

32-36: LGTM! Tests updated for new Init signature.

All test functions have been correctly updated to create a logger and pass it to Init(), matching the new signature. The use of LogLevelDebug for tests ensures good visibility during test execution.

Also applies to: 198-234

transports/bifrost-http/handlers/server.go (1)

214-214: LGTM! Init call updated for new signature.

The maxim plugin initialization correctly passes the logger parameter, matching the updated Init signature.

@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin branch 2 times, most recently from 3f6e759 to 1a1a736 Compare October 25, 2025 14:22

@coderabbitai coderabbitai 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.

Actionable comments posted: 1

♻️ Duplicate comments (2)
plugins/governance/main.go (2)

456-463: Case-insensitive header lookup for x-bf-team/x-bf-customer.

Map access is case-sensitive; normalize or search case-insensitively to avoid missing audit fields.

Apply:

@@
-	// Extract team/customer info for audit trail
-	var teamID, customerID *string
-	if teamIDValue := headers["x-bf-team"]; teamIDValue != "" {
-		teamID = &teamIDValue
-	}
-	if customerIDValue := headers["x-bf-customer"]; customerIDValue != "" {
-		customerID = &customerIDValue
-	}
+	// Extract team/customer info for audit trail (case-insensitive)
+	var teamID, customerID *string
+	getHeader := func(m map[string]string, key string) string {
+		lk := strings.ToLower(key)
+		for k, v := range m {
+			if strings.ToLower(k) == lk {
+				return v
+			}
+		}
+		return ""
+	}
+	if v := getHeader(headers, "x-bf-team"); v != "" {
+		teamID = &v
+	}
+	if v := getHeader(headers, "x-bf-customer"); v != "" {
+		customerID = &v
+	}

435-435: Guard async worker lifecycle; copy headers before goroutine.

  • Without tracking, Cleanup can return before usage updates finish. Add a WaitGroup and wait in Cleanup.
  • Pass a copied headers map to avoid racy reads if upstream mutates it.

Apply:

@@
-	go p.postHookWorker(result, provider, model, requestType, virtualKey, requestID, headers, isCacheRead, isBatch, bifrost.IsFinalChunk(ctx))
+	// Copy headers to avoid concurrent map access after PostHook returns
+	headersCopy := make(map[string]string, len(headers))
+	for k, v := range headers {
+		headersCopy[k] = v
+	}
+	p.updateWg.Add(1)
+	go func() {
+		defer p.updateWg.Done()
+		p.postHookWorker(result, provider, model, requestType, virtualKey, requestID, headersCopy, isCacheRead, isBatch, bifrost.IsFinalChunk(ctx))
+	}()

And add:

@@
 type GovernancePlugin struct {
@@
 	inMemoryStore InMemoryStore
@@
 	isVkMandatory *bool
+	updateWg     sync.WaitGroup
 }
@@
 import (
 	"context"
 	"fmt"
 	"math/rand/v2"
 	"slices"
 	"sort"
 	"strings"
+	"sync"
@@
 func (p *GovernancePlugin) Cleanup() error {
 	if p.cancelFunc != nil {
 		p.cancelFunc()
 	}
+	// Wait for in-flight PostHook workers
+	p.updateWg.Wait()
 	if err := p.tracker.Cleanup(); err != nil {
 		return err
 	}
🧹 Nitpick comments (2)
plugins/logging/main.go (1)

299-311: Always return pooled UpdateLogData; use pool on error path.

Ensure UpdateLogData is pooled on all branches, not only the non-stream success path.

Apply:

@@
 	go func() {
 		requestType, _, _ := bifrost.GetResponseFields(result, bifrostErr)
 		// Queue the log update message (non-blocking) - use same pattern for both streaming and regular
 		logMsg := p.getLogMessage()
 		logMsg.RequestID = requestID
+		// Always return pooled UpdateData if set
+		defer func() {
+			if logMsg.UpdateData != nil {
+				p.putUpdateLogData(logMsg.UpdateData)
+			}
+		}()
 		defer p.putLogMessage(logMsg) // Return to pool when done
@@
-		// If response is nil, and there is an error, we update log with error
+		// If response is nil, and there is an error, we update log with error
 		if result == nil && bifrostErr != nil {
 			// If request type is streaming, then we trigger cleanup as well
 			if bifrost.IsStreamRequestType(requestType) {
 				p.accumulator.CleanupStreamAccumulator(requestID)
 			}
 			logMsg.Operation = LogOperationUpdate
-			logMsg.UpdateData = &UpdateLogData{
-				Status:       "error",
-				ErrorDetails: bifrostErr,
-			}
+			updateData := p.getUpdateLogData()
+			updateData.Status = "error"
+			updateData.ErrorDetails = bifrostErr
+			logMsg.UpdateData = updateData

And remove the later duplicate defer in the non-stream path:

@@
-			// Return pooled data structures to their respective pools
-			defer func() {
-				if logMsg.UpdateData != nil {
-					p.putUpdateLogData(logMsg.UpdateData)
-				}
-			}()
plugins/maxim/plugin_test.go (1)

32-36: Skip integration when env is missing (avoid CI flakes).

getPlugin errors on missing env; prefer skipping the integration test if API keys aren’t set.

Patch TestMaximLoggerPlugin:

@@
 func TestMaximLoggerPlugin(t *testing.T) {
 	ctx := context.Background()
+	if os.Getenv("MAXIM_API_KEY") == "" || os.Getenv("OPENAI_API_KEY") == "" {
+		t.Skip("MAXIM_API_KEY/OPENAI_API_KEY not set; skipping integration test")
+	}
 	// Initialize the Maxim plugin
 	plugin, err := getPlugin()
📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 0e55852 and 1a1a736.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (2 hunks)
✅ Files skipped from review due to trivial changes (1)
  • core/schemas/providers/anthropic/chat.go
🚧 Files skipped from review as they are similar to previous changes (8)
  • transports/bifrost-http/integrations/anthropic.go
  • core/schemas/providers/anthropic/text.go
  • core/schemas/providers/vertex/embedding.go
  • framework/streaming/responses.go
  • plugins/otel/docker-compose.yml
  • transports/bifrost-http/handlers/server.go
  • framework/streaming/types.go
  • core/providers/anthropic.go
🧰 Additional context used
🧬 Code graph analysis (5)
plugins/governance/main.go (2)
core/utils.go (1)
  • IsFinalChunk (130-145)
core/schemas/bifrost.go (3)
  • BifrostResponse (213-223)
  • ModelProvider (32-32)
  • RequestType (79-79)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (3)
  • BifrostLLMUsage (547-553)
  • ChatMessage (336-345)
  • ChatMessageContent (348-351)
plugins/maxim/main.go (6)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/plugin.go (1)
  • Plugin (45-71)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (82-82)
  • TextCompletionStreamRequest (83-83)
  • ChatCompletionRequest (84-84)
  • ChatCompletionStreamRequest (85-85)
  • ResponsesRequest (86-86)
  • ResponsesStreamRequest (87-87)
  • BifrostContextKeyRequestID (101-101)
  • RequestType (79-79)
  • BifrostResponse (213-223)
  • BifrostError (319-328)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (2)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (19)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (7)
core/schemas/providers/bedrock/chat.go (1)

43-44: LGTM! Improved method naming for clarity.

The rename from ToBifrostResponse to ToBifrostChatResponse makes the API more explicit about the response type being converted, improving code readability and maintainability.

core/schemas/providers/cohere/chat.go (1)

183-184: LGTM! Consistent naming improvement.

The rename aligns with the same pattern applied across other providers, making the conversion method more explicit about handling chat responses.

core/schemas/providers/bedrock/text.go (1)

60-61: LGTM! Clear distinction between response types.

Both renames from ToBifrostResponse to ToBifrostTextCompletionResponse make the API explicit about handling text completion responses (as opposed to chat responses), improving code clarity across Bedrock's Anthropic and Mistral implementations.

Also applies to: 84-85

core/schemas/providers/cohere/embedding.go (1)

64-65: LGTM!

All call sites have been properly updated. The ToBifrostEmbeddingResponse() method is now correctly called in both core/providers/cohere.go:762 and core/bedrock.go:876. The earlier search results showing ToBifrostResponse() calls are on the unrelated ProcessedStreamResponse type, which correctly retains its original method name and is not affected by this Cohere-specific rename.

plugins/otel/main.go (1)

213-244: LGTM: async emission, WG tracking, and error handling are correct.

Add-before-goroutine, defer Done inside, retain span until after emit, and log emit errors on both paths. Looks solid.

plugins/maxim/plugin_test.go (2)

198-198: LGTM: pass logger into Init to match new signature.

Tests reflect the updated Init(config, logger) API.


234-234: LGTM: negative-case Init validation covered.

Error-path test for missing API key aligns with new Init signature.

Comment thread plugins/maxim/main.go

@coderabbitai coderabbitai 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.

Actionable comments posted: 1

📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 0e55852 and 1a1a736.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (2 hunks)
✅ Files skipped from review due to trivial changes (2)
  • core/schemas/providers/cohere/chat.go
  • transports/bifrost-http/handlers/server.go
🚧 Files skipped from review as they are similar to previous changes (6)
  • core/schemas/providers/bedrock/chat.go
  • core/schemas/providers/anthropic/chat.go
  • core/providers/anthropic.go
  • plugins/otel/docker-compose.yml
  • plugins/governance/main.go
  • transports/bifrost-http/integrations/anthropic.go
🧰 Additional context used
🧬 Code graph analysis (7)
framework/streaming/responses.go (2)
core/schemas/responses.go (5)
  • BifrostResponsesStreamResponse (1334-1367)
  • BifrostResponsesResponse (40-72)
  • ResponsesOutputMessageContentTextLogProb (425-430)
  • ResponsesOutputMessageContentText (406-409)
  • ResponsesOutputMessageContentRefusal (431-433)
core/schemas/bifrost.go (3)
  • OpenAI (35-35)
  • OpenRouter (48-48)
  • Azure (36-36)
framework/streaming/types.go (2)
core/schemas/responses.go (1)
  • BifrostResponsesResponse (40-72)
core/schemas/bifrost.go (3)
  • BifrostResponseExtraFields (251-259)
  • RequestType (79-79)
  • ResponsesRequest (86-86)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/providers/anthropic/text.go (2)
core/schemas/providers/anthropic/types.go (2)
  • AnthropicTextRequest (18-27)
  • AnthropicTextResponse (205-214)
core/schemas/textcompletions.go (2)
  • BifrostTextCompletionRequest (10-16)
  • BifrostTextCompletionResponse (64-72)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
plugins/maxim/main.go (6)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/plugin.go (1)
  • Plugin (45-71)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (82-82)
  • TextCompletionStreamRequest (83-83)
  • ChatCompletionRequest (84-84)
  • ChatCompletionStreamRequest (85-85)
  • ResponsesRequest (86-86)
  • ResponsesStreamRequest (87-87)
  • BifrostContextKeyRequestID (101-101)
  • RequestType (79-79)
  • BifrostResponse (213-223)
  • BifrostError (319-328)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (2)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (4)
  • BifrostLLMUsage (547-553)
  • ChatMessage (336-345)
  • ChatMessageRoleAssistant (328-328)
  • ChatMessageContent (348-351)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (24)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (19)
framework/streaming/types.go (1)

198-213: LGTM!

The new StreamTypeResponses case correctly constructs a BifrostResponsesResponse from accumulated data, converts token usage appropriately, and populates extra fields consistently with other stream types.

core/schemas/providers/vertex/embedding.go (1)

69-70: LGTM!

Method rename improves API clarity and aligns with the broader naming convention update across providers.

plugins/maxim/plugin_test.go (2)

32-36: LGTM!

Test properly instantiates logger and passes it to Init, aligning with the updated signature.


198-234: LGTM!

All test cases correctly create and pass the logger parameter to Init.

core/schemas/providers/cohere/embedding.go (1)

64-65: LGTM!

Method rename enhances API consistency across providers.

core/schemas/providers/bedrock/text.go (1)

60-61: LGTM!

Method renames improve type specificity and maintain consistency with the broader provider API refactoring.

Also applies to: 84-85

plugins/logging/main.go (3)

299-340: LGTM!

The async refactor correctly handles:

  • Pool management with proper defer placement
  • Latency extraction from result when available
  • Error path with streaming accumulator cleanup
  • Retry logic with callback invocation

The goroutine safely captures result and bifrostErr as they're not modified post-launch.


342-368: LGTM!

Streaming path correctly:

  • Processes responses via accumulator
  • Updates log only on final chunk
  • Uses retry logic for database updates
  • Invokes callback after successful update

369-481: LGTM!

Non-streaming path thoroughly extracts:

  • Token usage from all response types
  • Chat, responses, embedding, speech, and transcription outputs
  • Raw response from extra fields
  • Cost calculation when pricing manager available

Pool management is correct with updateData obtained and deferred for return.

framework/streaming/responses.go (3)

13-136: LGTM!

The deep copy implementation correctly handles:

  • All pointer fields with proper nil checks
  • Nested structures (Response, Item, Part)
  • Slice fields (Output, LogProbs)
  • ExtraFields (shallow copy is safe as noted in comment)

This prevents shared data mutation between plugin accumulators during concurrent processing.


138-221: LGTM!

Helper functions properly deep copy nested message structures:

  • deepCopyResponsesMessage handles all ResponsesMessage fields
  • deepCopyResponsesMessageContentBlock copies content block variants
  • Both use value copies for strings/primitives and deep copies for nested structures

227-233: LGTM!

Key improvements:

  • Sorting chunks by sequence number (lines 227-233) ensures stable, correct message construction
  • Deep copying the stream response (line 562) prevents data races when multiple plugins access the same chunk

The OpenAI-compatible provider handling (line 491) correctly distinguishes providers that send complete responses in final chunks.

Also applies to: 489-562

plugins/otel/main.go (1)

213-244: LGTM!

The async emission refactor correctly addresses all previous review concerns:

  • Add(1) called before goroutine launch (line 214) ✓
  • defer Done() at goroutine start (line 216) ✓
  • No duplicate WaitGroup operations ✓
  • Error handling for both streaming (line 232) and non-streaming (line 240) emissions ✓
  • Span cleanup deferred after emission for both paths ✓

The implementation ensures all emissions are tracked and properly cleaned up during shutdown.

core/schemas/providers/anthropic/text.go (1)

48-81: LGTM! Method names now clearly indicate conversion types.

The renames from generic ToBifrostRequest/ToBifrostResponse to specific ToBifrostTextCompletionRequest/ToBifrostTextCompletionResponse improve code clarity and make the conversion types explicit. The implementation logic remains unchanged.

plugins/maxim/main.go (5)

11-11: LGTM! Clean import organization and constant additions.

The new imports support the streaming functionality, and the PluginLoggerPrefix constant enables consistent, structured logging throughout the plugin.

Also applies to: 14-14, 16-16, 23-26


78-78: Verify if nil pricingManager is intentional.

The accumulator is initialized with nil as the pricingManager. If pricing tracking is not required for the Maxim plugin's logging functionality, this is fine. However, if usage costs should be tracked, a valid pricingManager instance should be passed.

Can you confirm whether pricing tracking is needed for this plugin?


423-433: Streaming accumulator creation looks correct.

The code properly extracts the requestID and creates a streaming accumulator for streaming request types. The early return when requestID is missing serves as a defensive guard.

Based on learnings: The maintainer has confirmed that ctx is never nil and requestID is always set at the gateway level, making these checks safe guards that won't block legitimate requests.


475-486: LGTM! Streaming response handling is correctly implemented.

The code properly:

  • Processes streaming responses through the accumulator
  • Returns early for delta responses (intermediate chunks)
  • Uses accumulated data for final responses
  • Falls back to direct result data for non-streaming requests
  • Handles all three request types (TextCompletion, ChatCompletion, Responses) consistently

Also applies to: 508-525


540-542: LGTM! Proper cleanup of streaming resources.

The accumulator cleanup is correctly added while preserving the existing logger flush behavior.

Comment thread plugins/maxim/main.go
@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-24-feat_added_native_responses_support_to_azure_and_openrouter_streaming_and_openai_error_handling_optimisation branch from 7f62877 to 218d879 Compare October 25, 2025 16:26
@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin branch 2 times, most recently from efafe39 to 448221b Compare October 25, 2025 16:28
@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-24-feat_added_native_responses_support_to_azure_and_openrouter_streaming_and_openai_error_handling_optimisation branch from 218d879 to ee0dced Compare October 25, 2025 16:36
@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin branch from 448221b to 46d87f9 Compare October 25, 2025 16:36

@coderabbitai coderabbitai 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.

Actionable comments posted: 2

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (4)
core/schemas/providers/cohere/embedding.go (1)

90-101: Fix pointer-to-range-var in base64 loop (all pointers end up identical).

Taking &embedding in a for _, embedding := range ... loop captures the reused loop var; every EmbeddingStr will point to the last value. Use the slice element address instead.

Apply:

-       for i, embedding := range response.Embeddings.Base64 {
+       for i := range response.Embeddings.Base64 {
+               embedding := response.Embeddings.Base64[i]
                bifrostEmbedding := schemas.EmbeddingData{
                        Object: "embedding",
                        Index:  i,
                        Embedding: schemas.EmbeddingStruct{
-                               EmbeddingStr: &embedding,
+                               EmbeddingStr: &embedding, // local copy gives a stable address
                        },
                }
                bifrostEmbeddings = append(bifrostEmbeddings, bifrostEmbedding)
        }

Optionally, avoid the extra local and take the element address directly:

-       EmbeddingStr: &embedding,
+       EmbeddingStr: &response.Embeddings.Base64[i],
core/schemas/providers/bedrock/text.go (2)

11-14: Nil-deref risk: guard bifrostReq.Input before accessing fields.

bifrostReq.Input.PromptStr is accessed even when Input could be nil, causing panic.

Apply:

- if bifrostReq == nil || (bifrostReq.Input.PromptStr == nil && len(bifrostReq.Input.PromptArray) == 0) {
+ if bifrostReq == nil || bifrostReq.Input == nil || (bifrostReq.Input.PromptStr == nil && len(bifrostReq.Input.PromptArray) == 0) {
        return nil
 }

91-99: Fix pointer-to-range-var in Mistral outputs (choices share same address).

&output.Text and &output.StopReason take addresses of the loop variable; all choices end up pointing to the last element.

Apply:

- for i, output := range response.Outputs {
+ for i := range response.Outputs {
+     out := &response.Outputs[i]
      choices = append(choices, schemas.BifrostResponseChoice{
-         Index: i,
+         Index: i,
          TextCompletionResponseChoice: &schemas.TextCompletionResponseChoice{
-             Text: &output.Text,
+             Text: &out.Text,
          },
-         FinishReason: &output.StopReason,
+         FinishReason: &out.StopReason,
      })
 }
framework/streaming/responses.go (1)

283-294: Avoid duplicating tool-call messages on arguments deltas.

Appending resp.Item in every FunctionCallArgumentsDelta can produce multiple identical messages for the same call.

Apply:

-    if resp.Item != nil {
-        messages = append(messages, *resp.Item)
-    }
-    // Append arguments to the most recent message
+    // Append arguments to the most recent message (item was already added on OutputItemAdded)

If you need to seed the first delta when no item exists, gate it:

if len(messages) == 0 && resp.Item != nil {
    messages = append(messages, *resp.Item)
}
🧹 Nitpick comments (5)
core/schemas/providers/cohere/embedding.go (1)

79-89: Optional: avoid aliasing float slices.

EmbeddingArray: embedding copies the slice header, not the data. If upstream slices could be reused/mutated, defensively clone.

Apply:

-                       EmbeddingArray: embedding,
+                       EmbeddingArray: append([]float64(nil), embedding...),
framework/streaming/responses.go (1)

20-24: Minor: avoid naming locals copy.

Using copy shadows the built-in copy function and can confuse readers.

Rename to cloned or dup.

plugins/logging/main.go (1)

451-456: Optional: zero UpdateLogData before pooling to reduce retention.

If putUpdateLogData doesn’t reset fields, large slices (e.g., EmbeddingOutput, ResponsesOutput) may stick around.

Please confirm putUpdateLogData clears fields; if not, add:

func (p *LoggerPlugin) putUpdateLogData(d *UpdateLogData) {
    *d = UpdateLogData{} // zero
    p.updateDataPool.Put(d)
}
plugins/maxim/main.go (2)

476-485: On streaming processing error, ensure cleanup (and consider finalization).

If ProcessStreamingResponse returns an error, the goroutine returns without accumulator cleanup or generation/trace finalization.

Apply:

-        if bifrost.IsStreamRequestType(requestType) {
-            streamResponse, err = plugin.accumulator.ProcessStreamingResponse(ctx, result, bifrostErr)
-            if err != nil {
-                plugin.logger.Error("%s failed to process streaming response: %v", PluginLoggerPrefix, err)
-                return
-            }
+        if bifrost.IsStreamRequestType(requestType) {
+            streamResponse, err = plugin.accumulator.ProcessStreamingResponse(ctx, result, bifrostErr)
+            if err != nil {
+                plugin.logger.Error("%s failed to process streaming response: %v", PluginLoggerPrefix, err)
+                plugin.accumulator.CleanupStreamAccumulator(requestID) // prevent retention
+                // continue to finalize below (logger/end) if IDs exist
+            }

If you prefer to bail out on error, at least keep the cleanup.


423-434: Guardrail note: early return on missing requestID.

You return early if requestID is empty; if the gateway guarantee ever changes, streaming accumulators won’t be created and traces may remain open.

Log and proceed for non-streaming; only skip accumulator creation when empty. Based on learnings.

📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 1a1a736 and 46d87f9.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (1 hunks)
🚧 Files skipped from review as they are similar to previous changes (10)
  • core/schemas/providers/cohere/chat.go
  • core/schemas/providers/vertex/embedding.go
  • core/schemas/providers/anthropic/text.go
  • core/schemas/providers/bedrock/chat.go
  • plugins/otel/docker-compose.yml
  • core/schemas/providers/anthropic/chat.go
  • framework/streaming/types.go
  • transports/bifrost-http/integrations/anthropic.go
  • core/providers/anthropic.go
  • plugins/governance/main.go
🧰 Additional context used
🧬 Code graph analysis (6)
transports/bifrost-http/handlers/server.go (1)
plugins/maxim/main.go (1)
  • Init (62-92)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
plugins/maxim/main.go (6)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/plugin.go (1)
  • Plugin (45-71)
core/schemas/bifrost.go (7)
  • TextCompletionStreamRequest (84-84)
  • ChatCompletionStreamRequest (86-86)
  • ResponsesStreamRequest (88-88)
  • BifrostContextKeyRequestID (102-102)
  • RequestType (79-79)
  • BifrostResponse (216-226)
  • BifrostError (322-331)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (3)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
  • StreamResponseTypeFinal (24-24)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (6)
  • BifrostLLMUsage (547-553)
  • TextCompletionResponseChoice (507-509)
  • ChatMessage (336-345)
  • ChatMessageRoleAssistant (328-328)
  • ChatMessageContent (348-351)
  • ChatNonStreamResponseChoice (512-515)
framework/streaming/responses.go (2)
core/schemas/responses.go (5)
  • BifrostResponsesStreamResponse (1334-1367)
  • BifrostResponsesResponse (40-72)
  • ResponsesOutputMessageContentTextLogProb (425-430)
  • ResponsesOutputMessageContentText (406-409)
  • ResponsesOutputMessageContentRefusal (431-433)
core/schemas/bifrost.go (3)
  • OpenAI (35-35)
  • OpenRouter (48-48)
  • Azure (36-36)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (22)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (5)
plugins/logging/main.go (1)

299-311: Good: pooled LogMessage always returned.

defer p.putLogMessage(logMsg) ensures no leaks across all paths. Nice fix.

transports/bifrost-http/handlers/server.go (1)

214-214: LGTM! Logger parameter correctly added.

The change aligns with the updated maxim.Init signature that now requires a logger parameter for streaming accumulator support.

plugins/maxim/plugin_test.go (2)

32-36: LGTM! Test helper updated correctly.

The logger is properly initialized and passed to Init, matching the new signature.


198-234: LGTM! Test cases updated correctly.

The logger initialization and updated Init calls are consistent with the new plugin signature.

plugins/otel/main.go (1)

213-245: LGTM! Async emission refactoring correctly implemented.

The WaitGroup usage, error handling, and goroutine lifecycle are all correct:

  • Add(1) called before launching goroutine
  • defer Done() ensures cleanup in all paths
  • Both streaming and non-streaming emissions properly tracked and error-logged
  • Span cleanup appropriately deferred after emission

Comment thread framework/streaming/responses.go
Comment thread framework/streaming/responses.go

@coderabbitai coderabbitai 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.

Actionable comments posted: 1

♻️ Duplicate comments (1)
plugins/governance/main.go (1)

456-464: Case-sensitive header lookup for x-bf-team/x-bf-customer.

The header extraction is case-sensitive and may miss headers with alternate casing (e.g., X-Bf-Team, X-BF-CUSTOMER). HTTP headers are case-insensitive per RFC 7230.

This issue was previously identified and a solution was provided in past review comments. Please implement case-insensitive lookup using a helper function as suggested.

🧹 Nitpick comments (4)
framework/streaming/responses.go (2)

78-82: Consider deep copying nested slices in LogProbs.

The current implementation copies LogProbs structs by value, but ResponsesOutputMessageContentTextLogProb contains nested slices (Bytes []int and TopLogProbs []LogProb) that will remain shared between the original and copy. If these fields are mutated, the deep copy won't provide full isolation.

Consider this enhanced deep copy:

 	// Deep copy LogProbs slice if present
 	if original.LogProbs != nil {
 		copy.LogProbs = make([]schemas.ResponsesOutputMessageContentTextLogProb, len(original.LogProbs))
 		for i, logProb := range original.LogProbs {
-			copy.LogProbs[i] = logProb
+			// Deep copy the struct and its nested slices
+			copyLogProb := logProb
+			if logProb.Bytes != nil {
+				copyLogProb.Bytes = make([]int, len(logProb.Bytes))
+				copy_builtin(copyLogProb.Bytes, logProb.Bytes)
+			}
+			if logProb.TopLogProbs != nil {
+				copyLogProb.TopLogProbs = make([]schemas.LogProb, len(logProb.TopLogProbs))
+				copy_builtin(copyLogProb.TopLogProbs, logProb.TopLogProbs)
+			}
+			copy.LogProbs[i] = copyLogProb
 		}
 	}

173-191: Review completeness of ResponsesToolMessage copy.

The comment "Shallow copy for now" (line 175) suggests this implementation may be temporary. While the code does deep copy specific pointer fields (CallID, Name, Arguments), verify that all fields requiring isolation are covered or document any intentionally shared fields.

plugins/logging/main.go (2)

313-340: Use pool for UpdateData in error path for consistency.

The error path creates UpdateData inline (line 319) instead of obtaining it from the pool like the regular path does (line 373). This bypasses the pooling mechanism and causes unnecessary allocations during error conditions.

Apply this diff to use the pool consistently:

 		// If response is nil, and there is an error, we update log with error
 		if result == nil && bifrostErr != nil {
 			// If request type is streaming, then we trigger cleanup as well
 			if bifrost.IsStreamRequestType(requestType) {
 				p.accumulator.CleanupStreamAccumulator(requestID)
 			}
 			logMsg.Operation = LogOperationUpdate
-			logMsg.UpdateData = &UpdateLogData{
+			updateData := p.getUpdateLogData()
+			updateData.Status = "error"
+			updateData.ErrorDetails = bifrostErr
+			logMsg.UpdateData = updateData
+			defer func() {
+				if logMsg.UpdateData != nil {
+					p.putUpdateLogData(logMsg.UpdateData)
+				}
+			}()
-				Status:       "error",
-				ErrorDetails: bifrostErr,
-			}
 			processingErr := retryOnNotFound(p.ctx, func() error {

Alternatively, move the pooled-data return defer to the goroutine level (after line 304) so it covers all paths:

 	go func() {
 		requestType, _, _ := bifrost.GetResponseFields(result, bifrostErr)
 		// Queue the log update message (non-blocking) - use same pattern for both streaming and regular
 		logMsg := p.getLogMessage()
 		logMsg.RequestID = requestID
 		defer p.putLogMessage(logMsg) // Return to pool when done
+		defer func() {
+			if logMsg.UpdateData != nil {
+				p.putUpdateLogData(logMsg.UpdateData)
+			}
+		}()

Then obtain updateData from pool in the error path and remove the duplicate defer from the regular path (lines 452-456).


382-447: Consider extracting response field population into helper methods.

The logic for extracting token usage (lines 382-405) and outputs (lines 412-447) from various response types is correct but dense. Extracting these into separate helper methods would improve readability and testability.

For example:

func (p *LoggerPlugin) extractTokenUsage(result *schemas.BifrostResponse) *schemas.BifrostLLMUsage {
    // Lines 382-405 logic here
}

func (p *LoggerPlugin) populateOutputFields(updateData *UpdateLogData, result *schemas.BifrostResponse) {
    // Lines 412-447 logic here
}

Then simplify the PostHook goroutine:

updateData.TokenUsage = p.extractTokenUsage(result)
p.populateOutputFields(updateData, result)
📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 1a1a736 and 46d87f9.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (1 hunks)
✅ Files skipped from review due to trivial changes (1)
  • plugins/otel/docker-compose.yml
🚧 Files skipped from review as they are similar to previous changes (6)
  • transports/bifrost-http/integrations/anthropic.go
  • core/schemas/providers/vertex/embedding.go
  • transports/bifrost-http/handlers/server.go
  • core/schemas/providers/cohere/chat.go
  • core/providers/anthropic.go
  • core/schemas/providers/anthropic/chat.go
🧰 Additional context used
🧬 Code graph analysis (8)
core/schemas/providers/anthropic/text.go (2)
core/schemas/providers/anthropic/types.go (2)
  • AnthropicTextRequest (19-28)
  • AnthropicTextResponse (206-215)
core/schemas/textcompletions.go (2)
  • BifrostTextCompletionRequest (10-16)
  • BifrostTextCompletionResponse (64-72)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (4)
  • BifrostLLMUsage (547-553)
  • ChatMessage (336-345)
  • ChatMessageRoleAssistant (328-328)
  • ChatMessageContent (348-351)
plugins/governance/main.go (2)
core/utils.go (1)
  • IsFinalChunk (130-145)
core/schemas/bifrost.go (3)
  • BifrostResponse (216-226)
  • ModelProvider (32-32)
  • RequestType (79-79)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
plugins/maxim/main.go (6)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/plugin.go (1)
  • Plugin (45-71)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (83-83)
  • TextCompletionStreamRequest (84-84)
  • ChatCompletionRequest (85-85)
  • ChatCompletionStreamRequest (86-86)
  • ResponsesRequest (87-87)
  • ResponsesStreamRequest (88-88)
  • BifrostContextKeyRequestID (102-102)
  • RequestType (79-79)
  • BifrostResponse (216-226)
  • BifrostError (322-331)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (3)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
  • StreamResponseTypeFinal (24-24)
framework/streaming/responses.go (2)
core/schemas/responses.go (5)
  • BifrostResponsesStreamResponse (1334-1367)
  • BifrostResponsesResponse (40-72)
  • ResponsesOutputMessageContentTextLogProb (425-430)
  • ResponsesOutputMessageContentText (406-409)
  • ResponsesOutputMessageContentRefusal (431-433)
core/schemas/bifrost.go (3)
  • OpenAI (35-35)
  • OpenRouter (48-48)
  • Azure (36-36)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
framework/streaming/types.go (2)
core/schemas/responses.go (1)
  • BifrostResponsesResponse (40-72)
core/schemas/bifrost.go (3)
  • BifrostResponseExtraFields (254-262)
  • RequestType (79-79)
  • ResponsesRequest (87-87)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (31)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (22)
core/schemas/providers/bedrock/chat.go (1)

43-44: LGTM! Method rename improves clarity.

The rename from ToBifrostResponse to ToBifrostChatResponse makes the method name more explicit and aligns well with its return type (*schemas.BifrostChatResponse). This improves code readability and consistency across provider implementations.

core/schemas/providers/cohere/embedding.go (1)

64-65: Method rename verified and properly propagated across all call sites.

The rename from ToBifrostResponse to ToBifrostEmbeddingResponse has been successfully completed. No references to the old method name remain in the codebase, and the new method is actively being called in core/providers/cohere.go:824. The renaming is consistent across all embedding providers (Cohere, Vertex, Bedrock).

core/schemas/providers/anthropic/text.go (2)

48-49: LGTM! Method rename improves clarity.

The rename from ToBifrostRequest to ToBifrostTextCompletionRequest makes the method's purpose more explicit and aligns with the specific request type it handles.


80-81: LGTM! Consistent naming improvement.

The rename from ToBifrostResponse to ToBifrostTextCompletionResponse maintains consistency with the request method rename and clarifies the conversion target type.

framework/streaming/types.go (1)

198-213: LGTM! Proper streaming responses support added.

The implementation correctly handles the new StreamTypeResponses case with appropriate nil checks and field population. The conversion logic properly maps OutputMessages and TokenUsage to the BifrostResponsesResponse structure.

framework/streaming/responses.go (3)

196-221: LGTM with incremental approach noted.

The deep copy implementation covers the currently used content type fields. The comment on line 218 indicates an incremental approach—adding fields as they're needed—which is reasonable for maintaining the code.


489-491: LGTM! OpenAI-compatible provider handling clarified.

The updated comment and provider condition correctly identify OpenAI, OpenRouter, and Azure as OpenAI-compatible providers that return complete responses in the final chunk.


561-562: LGTM! Deep copy prevents shared mutation.

Using deepCopyResponsesStreamResponse instead of a direct reference successfully prevents shared data mutation between plugin accumulators, which aligns with the PR's stated objective.

plugins/otel/main.go (1)

78-78: LGTM! Async emission with proper WaitGroup tracking.

The asynchronous emission pattern is correctly implemented:

  • WaitGroup incremented before launching the goroutine
  • Deferred decrement at goroutine entry ensures proper cleanup
  • Both streaming final and non-streaming emit paths include error logging
  • Cleanup waits for all in-flight emissions to complete

This resolves the race conditions and resource tracking issues from earlier reviews.

Also applies to: 213-245

plugins/maxim/plugin_test.go (2)

32-36: LGTM!

The logger initialization and parameter passing align correctly with the new Init(config *Config, logger schemas.Logger) signature.


198-234: Logger setup looks good.

The logger is created once and reused across test cases. Note that the logger is only passed to Init when testing error conditions (line 234), while valid configuration tests skip actual initialization (lines 239-243). This is acceptable given the test's design to avoid requiring real API keys.

plugins/maxim/main.go (9)

11-11: LGTM!

The new imports (time and streaming) are properly utilized throughout the file. The const block organization with PluginLoggerPrefix improves consistency for logging.

Also applies to: 16-16, 23-26


36-52: LGTM!

The addition of accumulator and logger fields to the Plugin struct is appropriate for the new streaming functionality. Both fields are properly initialized in Init and used consistently throughout the lifecycle hooks.


62-62: Verify nil pricingManager is intentional.

The Init signature correctly accepts and wires the logger. However, line 78 passes nil for the pricingManager parameter when creating the accumulator:

accumulator: streaming.NewAccumulator(nil, logger),

From the relevant code snippets, NewAccumulator accepts pricingManager *pricing.PricingManager as its first parameter. If pricing/cost tracking is not required for the Maxim plugin's logging use case, this is acceptable. Otherwise, consider whether a pricing manager should be provided.

Can you confirm whether the Maxim plugin intentionally skips pricing tracking, or should it receive a pricing manager instance?

Also applies to: 78-79


269-269: LGTM!

Correctly added streaming request type cases (TextCompletionStreamRequest, ChatCompletionStreamRequest, ResponsesStreamRequest) to the switch statements for proper request handling.

Also applies to: 291-291, 325-325


423-433: Stream accumulator creation looks good.

The logic correctly extracts the requestID from context and creates a stream accumulator for streaming requests. Based on learnings, the requestID is guaranteed to be set at the gateway level, making the early return a safe guard rather than a functional issue.

Based on learnings


458-468: LGTM!

The early returns for missing effectiveLogRepoID and requestID are appropriate guards. Based on learnings, requestID is always set at the gateway level, so the check acts as a defensive safeguard.

Based on learnings


470-539: Verify goroutine usage for PostHook.

The entire PostHook logic now runs in a goroutine (line 470), with the function returning immediately (line 539). This makes the logging operations non-blocking, which can improve request latency. However, ensure this is intentional and that:

  1. Cleanup guarantees: The goroutine may still be running when Cleanup() is called during shutdown. Verify that logger.Flush() (line 537) and accumulator.CleanupStreamAccumulator() (lines 503, 527) complete before the application exits, or implement graceful shutdown coordination.

  2. Error visibility: Errors within the goroutine (lines 478, 490) are logged but not propagated to the caller. This is acceptable if logging failures should not block the main request flow.

  3. Context validity: The goroutine captures ctx, which may be cancelled after PostHook returns. Ensure ProcessStreamingResponse and logger operations handle cancelled contexts gracefully.

Additionally, the streaming cleanup logic looks correct:

  • Line 503: Cleanup on error
  • Line 527: Cleanup on final success

Can you confirm that:

  1. The async pattern is intentional and the plugin doesn't need to block until logging completes?
  2. Graceful shutdown ensures goroutines complete before process exit?

473-486: Streaming response handling looks correct.

The streaming logic properly:

  1. Processes responses via the accumulator (line 476)
  2. Returns early for delta/partial responses (lines 483-485)
  3. Adds accumulated results to generations for final responses (lines 508-524)
  4. Cleans up the accumulator on both error (line 503) and final success (line 527)

The switch statement correctly handles all streaming request types (TextCompletionStreamRequest, ChatCompletionStreamRequest, ResponsesStreamRequest) and uses the accumulated streaming response when available.

Also applies to: 506-528


543-545: LGTM!

The accumulator cleanup is properly guarded with a nil check and correctly integrated into the plugin's Cleanup() method. This ensures streaming resources are properly released during shutdown.

core/schemas/providers/bedrock/text.go (2)

84-109: LGTM! Method rename verified as complete and consistent.

The rename of ToBifrostResponse() to ToBifrostTextCompletionResponse() on BedrockMistralTextResponse has been successfully applied throughout the codebase. All call sites now use the new method name, including core/providers/bedrock.go:524, and the naming is consistent with BedrockAnthropicTextResponse.


60-82: ✓ Call site verification complete—refactoring is correct.

All references to BedrockAnthropicTextResponse and BedrockMistralTextResponse now use the new ToBifrostTextCompletionResponse() method name:

  • core/providers/bedrock.go:517 (Anthropic)
  • core/providers/bedrock.go:524 (Mistral)

No remaining calls to the old method name exist for these types. The rename improves clarity and is properly implemented across all call sites.

Comment thread plugins/governance/main.go Outdated

@coderabbitai coderabbitai 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.

Actionable comments posted: 3

♻️ Duplicate comments (1)
plugins/governance/main.go (1)

456-464: Case-insensitive lookup for x-bf-team/x-bf-customer.

Map keys from upstream headers can vary in casing. Do a case-insensitive lookup to avoid missing audit fields.

Apply locally within postHookWorker:

- // Extract team/customer info for audit trail
- var teamID, customerID *string
- if teamIDValue := headers["x-bf-team"]; teamIDValue != "" {
-   teamID = &teamIDValue
- }
- if customerIDValue := headers["x-bf-customer"]; customerIDValue != "" {
-   customerID = &customerIDValue
- }
+ // Extract team/customer info for audit trail (case-insensitive)
+ var teamID, customerID *string
+ getHeaderCI := func(m map[string]string, key string) string {
+   if m == nil { return "" }
+   lk := strings.ToLower(key)
+   for k, v := range m {
+     if strings.ToLower(k) == lk {
+       return v
+     }
+   }
+   return ""
+ }
+ if v := getHeaderCI(headers, "x-bf-team"); v != "" { teamID = &v }
+ if v := getHeaderCI(headers, "x-bf-customer"); v != "" { customerID = &v }
🧹 Nitpick comments (6)
framework/streaming/responses.go (1)

196-221: Consider comprehensive deep copy strategy for content blocks.

The function only deep copies ResponsesOutputMessageContentText and ResponsesOutputMessageContentRefusal fields, relying on the comment at line 218 to "Add other content type fields as they're used." This approach is error-prone—if new content types are introduced in ResponsesMessageContentBlock and this function isn't updated, shared data mutation bugs could silently reappear.

Consider one of these approaches:

  1. Use reflection to generically deep copy all fields, or
  2. Add a verification script to your CI that ensures all fields in ResponsesMessageContentBlock are explicitly handled here
plugins/governance/main.go (1)

227-246: Guard against zero/negative weights to avoid skewed selection.

If all weights are 0 (or non-positive), the loop won’t pick any and you fall back to the first provider, biasing selection.

- // Weighted random selection from allowed providers for the main model
- totalWeight := 0.0
- for _, config := range allowedProviderConfigs {
-   totalWeight += config.Weight
- }
+ // Filter out non-positive weights; if none remain, fall back to original list
+ filtered := make([]configstoreTables.TableVirtualKeyProviderConfig, 0, len(allowedProviderConfigs))
+ for _, c := range allowedProviderConfigs {
+   if c.Weight > 0 {
+     filtered = append(filtered, c)
+   }
+ }
+ if len(filtered) > 0 {
+   allowedProviderConfigs = filtered
+ }
+ // Weighted random selection from allowed providers for the main model
+ totalWeight := 0.0
+ for _, config := range allowedProviderConfigs {
+   totalWeight += config.Weight
+ }
plugins/otel/docker-compose.yml (1)

1-80: Pin image versions and consider resource limits.

Using :latest can introduce breaking changes; pin versions and optionally add basic cpu/mem limits for stability in CI/dev.

Also applies to: 178-229

plugins/logging/main.go (1)

319-326: Use pooled UpdateLogData in error path and return it via defer.

Aligns with success-path pooling and reduces GC churn.

- logMsg.Operation = LogOperationUpdate
- logMsg.UpdateData = &UpdateLogData{
-   Status:       "error",
-   ErrorDetails: bifrostErr,
- }
+ logMsg.Operation = LogOperationUpdate
+ ud := p.getUpdateLogData()
+ ud.Status = "error"
+ ud.ErrorDetails = bifrostErr
+ logMsg.UpdateData = ud
+ defer p.putUpdateLogData(ud)
plugins/maxim/main.go (2)

423-429: Unify ctx handling for consistency (non-blocking).

Maintainer note says ctx is never nil. Consider removing earlier ctx-nil checks in PreHook to reflect this guarantee, or add a guard here for consistency. Either path is fine—just be consistent.

Based on learnings


488-495: Avoid shadowing: rename local logger to improve readability.

The local logger (Maxim logger) shadows the plugin’s logger field; rename to mxLogger for clarity.

- logger, err := plugin.getOrCreateLogger(effectiveLogRepoID)
+ mxLogger, err := plugin.getOrCreateLogger(effectiveLogRepoID)
  if err != nil { return }
- logger.SetGenerationError(generationID, &genErr)
+ mxLogger.SetGenerationError(generationID, &genErr)
  ...
- logger.AddResultToGeneration(generationID, ...)
+ mxLogger.AddResultToGeneration(generationID, ...)
  ...
- logger.EndGeneration(generationID)
+ mxLogger.EndGeneration(generationID)
  ...
- logger.EndTrace(traceID)
+ mxLogger.EndTrace(traceID)
  ...
- logger.Flush()
+ mxLogger.Flush()
📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 1a1a736 and 46d87f9.

📒 Files selected for processing (18)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (2 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (1 hunks)
✅ Files skipped from review due to trivial changes (1)
  • core/schemas/providers/cohere/chat.go
🚧 Files skipped from review as they are similar to previous changes (6)
  • transports/bifrost-http/integrations/anthropic.go
  • plugins/maxim/plugin_test.go
  • core/schemas/providers/bedrock/chat.go
  • core/schemas/providers/cohere/embedding.go
  • plugins/otel/main.go
  • core/schemas/providers/vertex/embedding.go
🧰 Additional context used
🧬 Code graph analysis (7)
framework/streaming/responses.go (2)
core/schemas/responses.go (5)
  • BifrostResponsesStreamResponse (1334-1367)
  • BifrostResponsesResponse (40-72)
  • ResponsesOutputMessageContentTextLogProb (425-430)
  • ResponsesOutputMessageContentText (406-409)
  • ResponsesOutputMessageContentRefusal (431-433)
core/schemas/bifrost.go (3)
  • OpenAI (35-35)
  • OpenRouter (48-48)
  • Azure (36-36)
framework/streaming/types.go (2)
core/schemas/responses.go (1)
  • BifrostResponsesResponse (40-72)
core/schemas/bifrost.go (3)
  • BifrostResponseExtraFields (254-262)
  • RequestType (79-79)
  • ResponsesRequest (87-87)
core/schemas/providers/anthropic/text.go (2)
core/schemas/providers/anthropic/types.go (2)
  • AnthropicTextRequest (19-28)
  • AnthropicTextResponse (206-215)
core/schemas/textcompletions.go (2)
  • BifrostTextCompletionRequest (10-16)
  • BifrostTextCompletionResponse (64-72)
transports/bifrost-http/handlers/server.go (1)
plugins/maxim/main.go (1)
  • Init (62-92)
plugins/governance/main.go (2)
core/utils.go (1)
  • IsFinalChunk (130-145)
core/schemas/bifrost.go (3)
  • BifrostResponse (216-226)
  • ModelProvider (32-32)
  • RequestType (79-79)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (5)
  • BifrostLLMUsage (547-553)
  • ChatMessage (336-345)
  • ChatMessageRoleAssistant (328-328)
  • ChatMessageContent (348-351)
  • ChatNonStreamResponseChoice (512-515)
plugins/maxim/main.go (5)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (83-83)
  • TextCompletionStreamRequest (84-84)
  • ChatCompletionRequest (85-85)
  • ChatCompletionStreamRequest (86-86)
  • ResponsesRequest (87-87)
  • ResponsesStreamRequest (88-88)
  • BifrostContextKeyRequestID (102-102)
  • RequestType (79-79)
  • BifrostResponse (216-226)
  • BifrostError (322-331)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (3)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
  • StreamResponseTypeFinal (24-24)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (31)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (7)
framework/streaming/types.go (1)

198-213: LGTM! StreamTypeResponses handling implemented correctly.

The new case properly constructs a BifrostResponsesResponse from accumulated stream data, handles nil checks for OutputMessages and TokenUsage, and sets appropriate ExtraFields metadata. The implementation follows the established pattern used by other stream types.

framework/streaming/responses.go (2)

489-491: LGTM! Provider handling correctly extended.

Including OpenRouter and Azure alongside OpenAI in the special-case handling is appropriate, as these providers use OpenAI-compatible streaming response formats where the final chunk contains the complete accumulated response.


561-562: Good fix for preventing shared data mutation.

Using deepCopyResponsesStreamResponse instead of a direct reference correctly prevents shared data mutation between plugin accumulators during streaming response processing. Note that the effectiveness of this fix depends on the completeness of the deep copy implementation (see earlier comments about LogProbs and content block fields).

transports/bifrost-http/handlers/server.go (1)

214-216: The review comment incorrectly identifies a non-issue.

The two maxim.Init calls are from different packages:

  1. Line 214 (transports/bifrost-http/handlers/server.go): Calls the local Bifrost plugin's Init function with (maximConfig, logger) ✓
  2. Line 71 (plugins/maxim/main.go): Calls the external maxim SDK's Init function from package github.com/maximhq/maxim-go with only (&maxim.MaximSDKConfig{...})

These are distinct functions with independent signatures. The external SDK's Init is unrelated to the signature change of the local plugin's Init. Line 71 requires no update.

Likely an incorrect or invalid review comment.

core/schemas/providers/anthropic/chat.go (1)

250-251: LGTM: Method rename improves API clarity.

The rename from ToBifrostResponse to ToBifrostChatResponse makes the method's purpose more explicit, clearly indicating it converts to a chat-specific response format.

core/providers/anthropic.go (1)

287-287: LGTM: Call site updated correctly.

The call has been properly updated to use the renamed ToBifrostTextCompletionResponse() method, maintaining consistency with the API refactor.

core/schemas/providers/bedrock/text.go (1)

60-61: LGTM: Method renames enhance API specificity.

Both method renames from ToBifrostResponse to ToBifrostTextCompletionResponse improve clarity by explicitly indicating these convert to text completion response format. The changes are consistent across both Bedrock Anthropic and Bedrock Mistral implementations.

Also applies to: 84-85

Comment thread core/schemas/providers/anthropic/text.go
Comment thread framework/streaming/responses.go
Comment thread framework/streaming/responses.go
@Pratham-Mishra04
Pratham-Mishra04 force-pushed the 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin branch from 46d87f9 to e1cd99c Compare October 25, 2025 18:21

@coderabbitai coderabbitai 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.

Actionable comments posted: 0

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
plugins/logging/main.go (1)

299-482: Add panic recovery to background worker.

The goroutine that processes log updates lacks panic recovery. If any of the processing logic panics, the service could crash. This is especially important for background workers that handle errors, streaming responses, and database operations.

Apply this diff to add panic recovery at the start of the goroutine:

 go func() {
+    defer func() {
+        if r := recover(); r != nil {
+            p.logger.Error("panic in PostHook worker for request %s: %v", requestID, r)
+        }
+    }()
     requestType, _, _ := bifrost.GetResponseFields(result, bifrostErr)
     // Queue the log update message (non-blocking)
     logMsg := p.getLogMessage()
♻️ Duplicate comments (2)
plugins/governance/main.go (2)

437-441: Add panic recovery to background worker.

The goroutine lacks panic recovery, which could crash the service if postHookWorker encounters an unexpected error. This was suggested in a previous review but only the WaitGroup portion was implemented.

Apply this diff to add panic recovery:

 p.wg.Add(1)
 go func() {
     defer p.wg.Done()
+    defer func() {
+        if r := recover(); r != nil {
+            p.logger.Error("panic in postHookWorker: %v", r)
+        }
+    }()
     p.postHookWorker(result, provider, model, requestType, virtualKey, requestID, headers, isCacheRead, isBatch, bifrost.IsFinalChunk(ctx))
 }()

463-470: Implement case-insensitive header lookup.

HTTP headers are case-insensitive per RFC 7230. The current case-sensitive map access may miss x-bf-team and x-bf-customer headers if upstream services use different casing, resulting in incomplete audit trails.

Apply this diff to add case-insensitive lookup:

 // Extract team/customer info for audit trail
 var teamID, customerID *string
-if teamIDValue := headers["x-bf-team"]; teamIDValue != "" {
-    teamID = &teamIDValue
-}
-if customerIDValue := headers["x-bf-customer"]; customerIDValue != "" {
-    customerID = &customerIDValue
-}
+getHeader := func(m map[string]string, key string) string {
+    if m == nil { return "" }
+    lk := strings.ToLower(key)
+    for k, v := range m {
+        if strings.ToLower(k) == lk { return v }
+    }
+    return ""
+}
+if v := getHeader(headers, "x-bf-team"); v != "" { teamID = &v }
+if v := getHeader(headers, "x-bf-customer"); v != "" { customerID = &v }
🧹 Nitpick comments (5)
plugins/maxim/main.go (2)

423-433: Consider consistent nil handling for ctx.

While the maintainer confirmed that ctx is never nil, the PreHook function has multiple if ctx != nil checks earlier (lines 217-249), but this code dereferences ctx directly at line 424. For consistency and clarity, either:

  1. Remove all if ctx != nil checks in PreHook (since ctx is guaranteed non-nil), or
  2. Add a nil check here before dereferencing

The streaming accumulator creation logic itself looks good.

Based on learnings

-	requestID, ok := (*ctx).Value(schemas.BifrostContextKeyRequestID).(string)
+	var requestID string
+	if ctx != nil {
+		requestID, ok = (*ctx).Value(schemas.BifrostContextKeyRequestID).(string)
+	}
 	if !ok || requestID == "" {

Or remove the earlier nil checks:

-	if ctx != nil {
-		if existingGenerationID, ok := (*ctx).Value(GenerationIDKey).(string); ok && existingGenerationID != "" {
+	if existingGenerationID, ok := (*ctx).Value(GenerationIDKey).(string); ok && existingGenerationID != "" {

458-468: Consider consistent nil handling for ctx.

Similar to PreHook, line 465 dereferences ctx without a nil check. While ctx is guaranteed to be non-nil according to the maintainer, the getEffectiveLogRepoID method checks for nil (line 136). For consistency, consider adding a nil guard here or removing checks elsewhere.

Based on learnings

-	requestID, ok := (*ctx).Value(schemas.BifrostContextKeyRequestID).(string)
+	var requestID string
+	if ctx != nil {
+		requestID, ok = (*ctx).Value(schemas.BifrostContextKeyRequestID).(string)
+	}
 	if !ok || requestID == "" {
plugins/otel/main.go (3)

78-79: Good addition: explicit in-flight tracking with WaitGroup.

Nice improvement. One optional hardening: guard against new Add after Cleanup begins to avoid races during shutdown.

Apply this pattern:

 import (
   "context"
   "fmt"
   "os"
   "strings"
   "sync"
+  "sync/atomic"
   "time"
 )
@@
 type OtelPlugin struct {
   ...
   emitWg sync.WaitGroup // Track in-flight emissions
+  closing atomic.Bool   // Prevent new work during shutdown
 }
@@
 func (p *OtelPlugin) PostHook(ctx *context.Context, resp *schemas.BifrostResponse, bifrostErr *schemas.BifrostError) (*schemas.BifrostResponse, *schemas.BifrostError, error) {
+  if p.closing.Load() {
+    logger.Warn("otel plugin shutting down; skipping PostHook emit")
+    return resp, bifrostErr, nil
+  }
   ...
   p.emitWg.Add(1)
   go func() {
     defer p.emitWg.Done()
     ...
   }()
   return resp, bifrostErr, nil
 }
@@
 func (p *OtelPlugin) Cleanup() error {
+  p.closing.Store(true)
   p.emitWg.Wait()
   ...
 }

222-237: Defensive checks and small clarity improvements in streaming path.

  • Guard against both resp and bifrostErr being nil before calling GetResponseFields to avoid potential nil-deref. If the contract guarantees one is non-nil, add an early assert/log; otherwise bail out.
  • Avoid shadowing span when casting; use a new name to improve readability.
  • On ProcessStreamingResponse error, consider whether this can be terminal; if yes, delete the span and emit an error span to avoid relying on TTL cleanup.

Minimal diffs:

-    requestType, _, _ := bifrost.GetResponseFields(resp, bifrostErr)
-    if span, ok := span.(*ResourceSpan); ok {
+    if resp == nil && bifrostErr == nil {
+        logger.Warn("PostHook called with nil resp and error for request %s; skipping emit", traceID)
+        return
+    }
+    requestType, provider, model := bifrost.GetResponseFields(resp, bifrostErr)
+    if rspan, ok := span.(*ResourceSpan); ok {
       // We handle streaming responses differently, we will use the accumulator to process the response and then emit the final response
       if bifrost.IsStreamRequestType(requestType) {
         streamResponse, err := p.accumulator.ProcessStreamingResponse(ctx, resp, bifrostErr)
         if err != nil {
           logger.Error("failed to process streaming response: %v", err)
+          // Optional: if err is terminal for this stream, consider cleaning up:
+          // defer p.ongoingSpans.Delete(traceID)
+          // _ = p.client.Emit(p.ctx, []*ResourceSpan{completeResourceSpan(rspan, time.Now(), resp, bifrostErr, p.pricingManager)})
         }
         if streamResponse != nil && streamResponse.Type == streaming.StreamResponseTypeFinal {
           defer p.ongoingSpans.Delete(traceID)
-          if err := p.client.Emit(p.ctx, []*ResourceSpan{completeResourceSpan(span, time.Now(), streamResponse.ToBifrostResponse(), bifrostErr, p.pricingManager)}); err != nil {
-            logger.Error("failed to emit response span for request %s: %v", traceID, err)
+          if err := p.client.Emit(p.ctx, []*ResourceSpan{completeResourceSpan(rspan, time.Now(), streamResponse.ToBifrostResponse(), bifrostErr, p.pricingManager)}); err != nil {
+            logger.Error("failed to emit response span for request %s (provider=%s, model=%s): %v", traceID, provider, model, err)
           }
         }
         return
       }

If the framework guarantees at least one of (resp, bifrostErr) is non-nil, please confirm so we can simplify the guard. Otherwise, the early return above prevents a panic at the cost of leaving cleanup to the TTL map. Confirm which behavior you prefer.


239-241: Enrich non-stream emit error logs with provider/model.

Helps triage failures quickly.

-  rs := completeResourceSpan(span, time.Now(), resp, bifrostErr, p.pricingManager)
-  if err := p.client.Emit(p.ctx, []*ResourceSpan{rs}); err != nil {
-    logger.Error("failed to emit response span for request %s: %v", traceID, err)
+  // capture provider/model from earlier GetResponseFields call (see prior diff)
+  rs := completeResourceSpan(rspan, time.Now(), resp, bifrostErr, p.pricingManager)
+  if err := p.client.Emit(p.ctx, []*ResourceSpan{rs}); err != nil {
+    logger.Error("failed to emit response span for request %s (provider=%s, model=%s): %v", traceID, provider, model, err)
   }
📜 Review details

Configuration used: CodeRabbit UI

Review profile: CHILL

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 46d87f9 and e1cd99c.

📒 Files selected for processing (25)
  • core/providers/anthropic.go (1 hunks)
  • core/schemas/providers/anthropic/chat.go (1 hunks)
  • core/schemas/providers/anthropic/text.go (2 hunks)
  • core/schemas/providers/bedrock/chat.go (1 hunks)
  • core/schemas/providers/bedrock/text.go (2 hunks)
  • core/schemas/providers/cohere/chat.go (1 hunks)
  • core/schemas/providers/cohere/embedding.go (1 hunks)
  • core/schemas/providers/openai/chat.go (1 hunks)
  • core/schemas/providers/openai/embedding.go (1 hunks)
  • core/schemas/providers/openai/responses.go (2 hunks)
  • core/schemas/providers/openai/speech.go (1 hunks)
  • core/schemas/providers/openai/text.go (1 hunks)
  • core/schemas/providers/openai/transcription.go (1 hunks)
  • core/schemas/providers/vertex/embedding.go (1 hunks)
  • framework/streaming/responses.go (3 hunks)
  • framework/streaming/types.go (1 hunks)
  • plugins/governance/main.go (4 hunks)
  • plugins/logging/main.go (3 hunks)
  • plugins/maxim/main.go (9 hunks)
  • plugins/maxim/plugin_test.go (3 hunks)
  • plugins/otel/docker-compose.yml (3 hunks)
  • plugins/otel/main.go (2 hunks)
  • transports/bifrost-http/handlers/server.go (1 hunks)
  • transports/bifrost-http/integrations/anthropic.go (1 hunks)
  • transports/bifrost-http/integrations/openai.go (6 hunks)
🚧 Files skipped from review as they are similar to previous changes (9)
  • transports/bifrost-http/integrations/anthropic.go
  • core/schemas/providers/anthropic/chat.go
  • framework/streaming/responses.go
  • core/schemas/providers/bedrock/text.go
  • core/schemas/providers/cohere/chat.go
  • core/schemas/providers/vertex/embedding.go
  • framework/streaming/types.go
  • core/providers/anthropic.go
  • core/schemas/providers/anthropic/text.go
🧰 Additional context used
🧬 Code graph analysis (13)
core/schemas/providers/openai/speech.go (2)
core/schemas/providers/openai/types.go (1)
  • OpenAISpeechRequest (100-106)
core/schemas/speech.go (1)
  • BifrostSpeechRequest (9-15)
core/schemas/providers/openai/embedding.go (2)
core/schemas/providers/openai/types.go (1)
  • OpenAIEmbeddingRequest (27-32)
core/schemas/embedding.go (1)
  • BifrostEmbeddingRequest (9-15)
core/schemas/providers/openai/transcription.go (2)
core/schemas/providers/openai/types.go (1)
  • OpenAITranscriptionRequest (110-116)
core/schemas/transcriptions.go (1)
  • BifrostTranscriptionRequest (3-9)
plugins/maxim/plugin_test.go (3)
core/logger.go (1)
  • NewDefaultLogger (40-49)
core/schemas/logger.go (1)
  • LogLevelDebug (10-10)
plugins/maxim/main.go (2)
  • Init (62-92)
  • Config (31-34)
core/schemas/providers/openai/text.go (2)
core/schemas/providers/openai/types.go (1)
  • OpenAITextCompletionRequest (13-19)
core/schemas/textcompletions.go (1)
  • BifrostTextCompletionRequest (10-16)
transports/bifrost-http/handlers/server.go (1)
plugins/maxim/main.go (1)
  • Init (62-92)
plugins/otel/main.go (2)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/providers/openai/responses.go (2)
core/schemas/providers/openai/types.go (1)
  • OpenAIResponsesRequest (86-92)
core/schemas/responses.go (1)
  • BifrostResponsesRequest (32-38)
plugins/maxim/main.go (6)
core/schemas/logger.go (1)
  • Logger (28-55)
framework/streaming/accumulator.go (2)
  • Accumulator (15-31)
  • NewAccumulator (393-433)
core/schemas/plugin.go (1)
  • Plugin (45-71)
core/schemas/bifrost.go (10)
  • TextCompletionRequest (83-83)
  • TextCompletionStreamRequest (84-84)
  • ChatCompletionRequest (85-85)
  • ChatCompletionStreamRequest (86-86)
  • ResponsesRequest (87-87)
  • ResponsesStreamRequest (88-88)
  • BifrostContextKeyRequestID (102-102)
  • RequestType (79-79)
  • BifrostResponse (216-226)
  • BifrostError (322-331)
core/utils.go (2)
  • IsStreamRequestType (126-128)
  • GetResponseFields (148-155)
framework/streaming/types.go (3)
  • ProcessedStreamResponse (107-114)
  • StreamResponseTypeDelta (23-23)
  • StreamResponseTypeFinal (24-24)
plugins/governance/main.go (2)
core/utils.go (1)
  • IsFinalChunk (130-145)
core/schemas/bifrost.go (3)
  • BifrostResponse (216-226)
  • ModelProvider (32-32)
  • RequestType (79-79)
core/schemas/providers/openai/chat.go (2)
core/schemas/providers/openai/types.go (1)
  • OpenAIChatRequest (35-45)
core/schemas/chatcompletions.go (1)
  • BifrostChatRequest (11-17)
plugins/logging/main.go (3)
core/utils.go (2)
  • GetResponseFields (148-155)
  • IsStreamRequestType (126-128)
framework/streaming/types.go (1)
  • StreamResponseTypeFinal (24-24)
core/schemas/chatcompletions.go (3)
  • BifrostLLMUsage (547-553)
  • ChatMessage (336-345)
  • ChatMessageContent (348-351)
transports/bifrost-http/integrations/openai.go (2)
core/schemas/bifrost.go (5)
  • TextCompletionRequest (83-83)
  • ResponsesRequest (87-87)
  • EmbeddingRequest (89-89)
  • SpeechRequest (90-90)
  • TranscriptionRequest (92-92)
transports/bifrost-http/handlers/inference.go (5)
  • ChatRequest (165-169)
  • ResponsesRequest (195-199)
  • EmbeddingRequest (202-206)
  • SpeechRequest (208-212)
  • TranscriptionRequest (214-218)
⏰ Context from checks skipped due to timeout of 900000ms. You can increase the timeout in your CodeRabbit configuration to a maximum of 15 minutes (900000ms). (2)
  • GitHub Check: Graphite / mergeability_check
  • GitHub Check: Graphite / mergeability_check
🔇 Additional comments (27)
plugins/otel/docker-compose.yml (1)

10-31: Verify: Does this Docker Compose config change belong in this PR?

The changes to this file are purely formatting/cosmetic (comment alignment, command array normalization, volume entry styling) and do not affect functionality. However, this file appears orthogonal to the stated PR objectives, which focus on Go plugin async fixes, provider schema renaming, streaming response handling, and concurrency improvements in the core framework.

If this is an incidental cleanup, consider extracting it to a separate PR to keep the commit history focused on the core objectives.

Also applies to: 229-229

plugins/governance/main.go (2)

11-11: WaitGroup addition correctly manages goroutine lifecycle.

The addition of sync.WaitGroup to track active background workers and the corresponding wg.Wait() in Cleanup ensure that all goroutines complete before shutdown, preventing data loss.

Also applies to: 42-42, 448-448


459-459: Cleaner separation: header extraction delegated to worker.

Moving header extraction into the background worker (postHookWorker) reduces work on the hot path and keeps the PostHook function lightweight.

plugins/logging/main.go (2)

304-304: Correct placement of deferred pool return.

The defer p.putLogMessage(logMsg) ensures the pooled LogMessage is always returned, even on early returns or errors. This prevents pool depletion under load.


306-310: Latency extraction correctly handles nil result.

The logic correctly extracts latency from the result when available and defaults to 0 otherwise. Based on learnings, BifrostErrorExtraFields does not contain a latency field, so this implementation is appropriate.

Based on learnings

core/schemas/providers/openai/chat.go (1)

5-17: LGTM! Method rename improves API clarity.

The rename from ToBifrostRequest to ToBifrostChatRequest makes the method's purpose more explicit and aligns with the broader refactoring pattern across providers. The comment and implementation are correct.

core/schemas/providers/openai/speech.go (1)

5-17: LGTM! Consistent with the broader refactoring pattern.

The rename from ToBifrostRequest to ToBifrostSpeechRequest follows the same pattern as other OpenAI conversions, improving API specificity.

core/schemas/providers/openai/responses.go (1)

5-29: LGTM! Good addition of documentation.

The rename to ToBifrostResponsesRequest is consistent with the pattern, and the added comment for ToOpenAIResponsesRequest improves code documentation.

core/schemas/providers/bedrock/chat.go (1)

43-142: LGTM! Response conversion naming is now consistent.

The rename from ToBifrostResponse to ToBifrostChatResponse extends the specificity pattern to response conversions, making the Bedrock provider consistent with the broader refactoring.

core/schemas/providers/openai/text.go (1)

27-41: LGTM! Clear distinction between text completion and chat.

The rename to ToBifrostTextCompletionRequest provides clear differentiation from chat-based conversions, improving API discoverability.

core/schemas/providers/openai/transcription.go (1)

5-19: LGTM! Transcription conversion is now explicitly named.

The rename to ToBifrostTranscriptionRequest maintains consistency with the refactoring pattern across all OpenAI request types.

core/schemas/providers/openai/embedding.go (1)

7-19: LGTM! Embedding conversion follows the pattern.

The rename to ToBifrostEmbeddingRequest is consistent with the broader refactoring and makes the API more self-documenting.

core/schemas/providers/cohere/embedding.go (1)

64-124: Refactor verified—no issues found.

The ToBifrostResponse() calls that remain are on the generic streaming.ProcessedStreamResponse type (framework/streaming/types.go:117), which legitimately uses a generic method name since it handles multiple response types (text, chat, embeddings) in the streaming accumulator. The provider-specific response types have been correctly updated to use specific method names (e.g., ToBifrostEmbeddingResponse()). The refactor is complete.

transports/bifrost-http/integrations/openai.go (1)

115-115: Verified: All method renames complete with no remaining generic calls.

Confirmed that all six OpenAI request type conversions have been successfully renamed to their type-specific methods (ToBifrostTextCompletionRequest, ToBifrostChatRequest, ToBifrostResponsesRequest, ToBifrostEmbeddingRequest, ToBifrostSpeechRequest, ToBifrostTranscriptionRequest). No remaining generic ToBifrostRequest() calls detected, and the method implementations are properly defined in the provider schema files. The refactoring pattern is consistently applied across multiple provider integrations. This is a solid improvement to code clarity and maintainability.

transports/bifrost-http/handlers/server.go (1)

214-214: LGTM!

The logger parameter is correctly passed to maxim.Init, aligning with the updated plugin initialization signature.

plugins/maxim/plugin_test.go (2)

32-36: LGTM!

Test properly creates a logger and passes it to Init, aligning with the new plugin initialization signature.


198-234: LGTM!

Test updates correctly handle the new Init signature by creating and passing a logger instance for all test cases.

plugins/maxim/main.go (9)

11-16: LGTM!

The new imports (time, bifrost, streaming) are properly utilized for the streaming functionality added in this PR.


23-26: LGTM!

Adding PluginLoggerPrefix enhances logging consistency throughout the plugin.


36-52: LGTM!

The updated Plugin struct properly integrates streaming support with the new accumulator and logger fields.


269-269: LGTM!

Correctly includes TextCompletionStreamRequest in the case statement to handle streaming variants.


291-291: LGTM!

Correctly includes ChatCompletionStreamRequest in the case statement to handle streaming variants.


325-325: LGTM!

Correctly includes ResponsesStreamRequest in the case statement to handle streaming variants.


470-539: LGTM! Async PostHook pattern and cleanup logic are solid.

The asynchronous processing in a goroutine is appropriate for avoiding blocking, and the streaming accumulator cleanup is correctly handled in both error (line 503) and success (line 527) paths. The early return for delta responses (line 484) prevents unnecessary processing.


543-545: LGTM!

Proper cleanup of the accumulator with appropriate nil guard.


78-78: The nil pricingManager in maxim is valid but inconsistent with other plugins.

The accumulator's pricingManager field is actively used across the codebase (chat, audio, responses, transcription streams all call CalculateCostWithCacheDebug when the manager is not nil). All usages include nil checks, making nil a valid operational state.

However, maxim is the only plugin that doesn't receive pricingManager in its Init function. Other plugins—otel, logging, governance, and telemetry—all receive it, and telemetry explicitly notes that "all cost calculations will be skipped" without it. This architectural inconsistency warrants clarification: does maxim intentionally omit pricing features, or should it follow the pattern of receiving pricingManager like other plugins?

plugins/otel/main.go (1)

212-217: Async emission flow + concurrency-safety verified.

Both OtelClient.Emit implementations are concurrency-safe. The HTTP client uses Go's standard http.Client (explicitly thread-safe), and the gRPC implementation uses protobuf-generated stubs (designed for concurrent requests). Neither implementation mutates shared state; headers are read-only after initialization. The async emission pattern in main.go with WG accounting is correct.

akshaydeo commented Oct 25, 2025 •

Copy link
Copy Markdown
Contributor

Merge activity

  • Oct 25, 9:14 PM UTC: A user started a stack merge that includes this pull request via Graphite.
  • Oct 25, 9:16 PM UTC: @akshaydeo merged this pull request with Graphite.

@akshaydeo
akshaydeo changed the base branch from 10-24-feat_added_native_responses_support_to_azure_and_openrouter_streaming_and_openai_error_handling_optimisation to graphite-base/681 October 25, 2025 21:15
@akshaydeo
akshaydeo changed the base branch from graphite-base/681 to main October 25, 2025 21:15
@akshaydeo
akshaydeo merged commit bf83adb into main Oct 25, 2025
4 checks passed
@akshaydeo
akshaydeo deleted the 10-25-fix-plugin-async-fixes-and-streaming-added-to-maxim-plugin branch October 25, 2025 21:16
akshaydeo added a commit that referenced this pull request Nov 17, 2025
## Summary

Improved method naming consistency and implemented deep copying for streaming responses to prevent data mutation between plugin accumulators.

## Changes

- Renamed conversion methods in provider schemas to be more specific and consistent (e.g., `ToBifrostResponse` → `ToBifrostTextCompletionResponse`)
- Added deep copy functionality for `ResponsesStreamResponse` to prevent shared data mutation between different plugin accumulators
- Fixed streaming response handling in the Maxim plugin to properly support streaming requests
- Improved concurrency handling in plugins (Logging, Governance, OTEL) by moving response processing to goroutines
- Updated the Anthropic provider to use the renamed conversion method

## Type of change

- [x] Bug fix
- [ ] Feature
- [x] Refactor
- [ ] Documentation
- [ ] Chore/CI

## Affected areas

- [x] Core (Go)
- [ ] Transports (HTTP)
- [x] Providers/Integrations
- [x] Plugins
- [ ] UI (Next.js)
- [ ] Docs

## How to test

Test streaming responses with multiple plugins enabled:

```sh
# Core/Transports
go version
go test ./...

# Test with streaming requests
curl -X POST http://localhost:8080/v1/chat/completions \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer $API_KEY" \
  -d '{
    "model": "gpt-3.5-turbo",
    "messages": [{"role": "user", "content": "Hello"}],
    "stream": true
  }'
```

## Breaking changes

- [ ] Yes
- [x] No

## Related issues

Fixes issues with data mutation between plugin accumulators during streaming responses.

## Security considerations

No security implications.

## Checklist

- [x] I added/updated tests where appropriate
- [x] I verified builds succeed (Go and UI)
- [x] I verified the CI pipeline passes locally if applicable
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