Repository navigation
Fix: SSE flushing and responses done markers #117
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
8 commits
Select commit
Hold shift + click to select a range
654fe6c
Fix SSE flushing and responses done markers
SantiagoDePolonia 366f8d1
Harden responses done marker normalization
SantiagoDePolonia 363dda0
Harden streaming flush and done marker handling
SantiagoDePolonia 051d2fe
Update responses stream contract goldens
SantiagoDePolonia 4b33ed7
Handle EOF-terminated responses completion
SantiagoDePolonia 47f52d8
Merge remote-tracking branch 'origin/main' into fix/sse-flush-and-done
SantiagoDePolonia 12e8aad
Disable remote golangci-lint config verify
SantiagoDePolonia 1480f9b
Deflake streaming flush test
SantiagoDePolonia File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,201 @@ | ||
| package providers | ||
|
|
||
| import ( | ||
| "bytes" | ||
| "io" | ||
| ) | ||
|
|
||
| var responsesDoneMarker = []byte("data: [DONE]\n\n") | ||
|
|
||
| var responsesDoneLine = []byte("data: [DONE]") | ||
|
|
||
| var responsesDataPrefix = []byte("data: ") | ||
|
|
||
| var responsesCompletionPatterns = [][]byte{ | ||
| []byte(`"type":"response.completed"`), | ||
| []byte(`"type":"response.done"`), | ||
| } | ||
|
|
||
| // EnsureResponsesDone normalizes Responses API streams so clients always receive | ||
| // a terminal data: [DONE] marker when the upstream stream reaches a completed | ||
| // Responses event but closes at EOF before sending the final marker. | ||
| func EnsureResponsesDone(stream io.ReadCloser) io.ReadCloser { | ||
| if stream == nil { | ||
| return nil | ||
| } | ||
|
|
||
| return &responsesDoneWrapper{ | ||
| ReadCloser: stream, | ||
| atEventBoundary: true, | ||
| currentLineAtEventBoundary: true, | ||
| } | ||
| } | ||
|
|
||
| type responsesDoneWrapper struct { | ||
| io.ReadCloser | ||
| lineBuf []byte | ||
| pending []byte | ||
| sawDone bool | ||
| eventCompletedCandidate bool | ||
| completedEventReadyForDone bool | ||
| atEventBoundary bool | ||
| currentLineAtEventBoundary bool | ||
| emitted bool | ||
| } | ||
|
|
||
| func (w *responsesDoneWrapper) Read(p []byte) (int, error) { | ||
| if len(w.pending) > 0 { | ||
| n := copy(p, w.pending) | ||
| w.pending = w.pending[n:] | ||
| if len(w.pending) == 0 { | ||
| w.emitted = true | ||
| } | ||
| return n, nil | ||
| } | ||
|
|
||
| if w.emitted { | ||
| return 0, io.EOF | ||
| } | ||
|
|
||
| n, err := w.ReadCloser.Read(p) | ||
| if n > 0 { | ||
| w.trackStream(p[:n]) | ||
| } | ||
|
|
||
| if err == io.EOF { | ||
| if w.sawDone { | ||
| if n > 0 { | ||
| return n, nil | ||
| } | ||
| return 0, io.EOF | ||
| } | ||
|
|
||
| missingSuffix := w.synthesizeDoneSuffix() | ||
| if len(missingSuffix) == 0 { | ||
| if n > 0 { | ||
| return n, nil | ||
| } | ||
| return 0, io.EOF | ||
| } | ||
|
|
||
| if n > 0 { | ||
| w.pending = append(w.pending[:0], missingSuffix...) | ||
| return n, nil | ||
| } | ||
|
|
||
| n = copy(p, missingSuffix) | ||
| if n < len(missingSuffix) { | ||
| w.pending = append(w.pending[:0], missingSuffix[n:]...) | ||
| return n, nil | ||
| } | ||
|
|
||
| w.emitted = true | ||
| return n, nil | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| } | ||
|
|
||
| return n, err | ||
| } | ||
|
|
||
| func (w *responsesDoneWrapper) trackStream(data []byte) { | ||
| start := 0 | ||
| for i, b := range data { | ||
| if b != '\n' { | ||
| continue | ||
| } | ||
|
|
||
| w.lineBuf = append(w.lineBuf, data[start:i]...) | ||
| w.processLine(w.lineBuf) | ||
| w.lineBuf = w.lineBuf[:0] | ||
| start = i + 1 | ||
| w.currentLineAtEventBoundary = w.atEventBoundary | ||
| } | ||
|
|
||
| if start < len(data) { | ||
| w.lineBuf = append(w.lineBuf, data[start:]...) | ||
| } | ||
| } | ||
|
|
||
| func (w *responsesDoneWrapper) processLine(line []byte) { | ||
| line = bytes.TrimSuffix(line, []byte("\r")) | ||
| if len(line) == 0 { | ||
| if w.eventCompletedCandidate { | ||
| w.completedEventReadyForDone = true | ||
| } | ||
| w.eventCompletedCandidate = false | ||
| w.atEventBoundary = true | ||
| return | ||
| } | ||
|
|
||
| if w.completedEventReadyForDone && (!w.currentLineAtEventBoundary || !bytes.Equal(line, responsesDoneLine)) { | ||
| w.completedEventReadyForDone = false | ||
| } | ||
|
|
||
| if w.currentLineAtEventBoundary && bytes.Equal(line, responsesDoneLine) { | ||
| w.sawDone = true | ||
| } | ||
|
|
||
| if bytes.HasPrefix(line, responsesDataPrefix) { | ||
| if isCompletedDataLine(line) { | ||
| w.eventCompletedCandidate = true | ||
| } | ||
| } | ||
|
|
||
| w.atEventBoundary = false | ||
| } | ||
|
|
||
| func (w *responsesDoneWrapper) synthesizeDoneSuffix() []byte { | ||
| if w.sawDone { | ||
| return nil | ||
| } | ||
|
|
||
| if w.eventCompletedCandidate && len(w.lineBuf) == 0 { | ||
| return append([]byte{'\n'}, responsesDoneMarker...) | ||
| } | ||
|
|
||
| if isCompletedDataLine(w.lineBuf) { | ||
| return append([]byte("\n\n"), responsesDoneMarker...) | ||
| } | ||
|
|
||
| if !w.completedEventReadyForDone { | ||
| return nil | ||
| } | ||
|
|
||
| if len(w.lineBuf) == 0 { | ||
| if w.atEventBoundary { | ||
| return append([]byte(nil), responsesDoneMarker...) | ||
| } | ||
| return nil | ||
| } | ||
|
|
||
| if !w.currentLineAtEventBoundary || !isDoneLinePrefix(w.lineBuf) { | ||
| return nil | ||
| } | ||
|
|
||
| suffix := append([]byte(nil), responsesDoneLine[len(w.lineBuf):]...) | ||
| suffix = append(suffix, '\n', '\n') | ||
| return suffix | ||
| } | ||
|
|
||
| func isDoneLinePrefix(line []byte) bool { | ||
| if len(line) > len(responsesDoneLine) { | ||
| return false | ||
| } | ||
|
|
||
| return bytes.Equal(line, responsesDoneLine[:len(line)]) | ||
| } | ||
|
|
||
| func isCompletedDataLine(line []byte) bool { | ||
| line = bytes.TrimSuffix(line, []byte("\r")) | ||
| if !bytes.HasPrefix(line, responsesDataPrefix) { | ||
| return false | ||
| } | ||
|
|
||
| payload := line[len(responsesDataPrefix):] | ||
| for _, pattern := range responsesCompletionPatterns { | ||
| if bytes.Contains(payload, pattern) { | ||
| return true | ||
| } | ||
| } | ||
|
|
||
| return false | ||
| } | ||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.