@@ -384,21 +384,50 @@ func (h *OpenAIGatewayHandler) pumpDevinStream(
384384 header .Set ("Connection" , "keep-alive" )
385385 header .Set ("X-Accel-Buffering" , "no" )
386386
387+ flusher , _ := c .Writer .(http.Flusher )
388+ return h .pumpDevinStreamCore (ctx , account , stream , reqModel , result , streamStarted ,
389+ func (chunk []byte ) error {
390+ if _ , err := fmt .Fprintf (c .Writer , "data: %s\n \n " , chunk ); err != nil {
391+ return err
392+ }
393+ if flusher != nil {
394+ flusher .Flush ()
395+ }
396+ return nil
397+ },
398+ func () error {
399+ if _ , err := fmt .Fprint (c .Writer , "data: [DONE]\n \n " ); err != nil {
400+ return err
401+ }
402+ if flusher != nil {
403+ flusher .Flush ()
404+ }
405+ return nil
406+ },
407+ )
408+ }
409+
410+ // pumpDevinStreamCore 把 devin 事件流转成连续的 chat.completion.chunk JSON,
411+ // 由 emitChunk 决定下游协议形态(OpenAI SSE / Responses SSE / Anthropic SSE),
412+ // finish 负责协议收尾([DONE] 或各协议的终止事件)。
413+ func (h * OpenAIGatewayHandler ) pumpDevinStreamCore (
414+ ctx context.Context ,
415+ account * service.Account ,
416+ stream * devinpkg.ChatStream ,
417+ reqModel string ,
418+ result * service.OpenAIForwardResult ,
419+ streamStarted * bool ,
420+ emitChunk func (chunkJSON []byte ) error ,
421+ finish func () error ,
422+ ) error {
387423 created := time .Now ().Unix ()
388424 responseID := "chatcmpl-" + strings .ReplaceAll (strings .ToLower (fmt .Sprintf ("%x" , created )), " " , "" )
389- flusher , _ := c .Writer .(http.Flusher )
390425 roleSent := false
391426 toolIndexes := make (map [string ]int )
392427 firstTokenAt := time .Now ()
393428
394429 writeChunk := func (delta map [string ]any , finishReason * string , usage map [string ]any ) error {
395- if _ , err := fmt .Fprintf (c .Writer , "data: %s\n \n " , devinChunkJSON (responseID , reqModel , created , delta , finishReason , usage )); err != nil {
396- return err
397- }
398- if flusher != nil {
399- flusher .Flush ()
400- }
401- return nil
430+ return emitChunk (devinChunkJSON (responseID , reqModel , created , delta , finishReason , usage ))
402431 }
403432 sendRole := func () error {
404433 if roleSent {
@@ -486,11 +515,7 @@ func (h *OpenAIGatewayHandler) pumpDevinStream(
486515 return err
487516 }
488517 }
489- _ , _ = fmt .Fprint (c .Writer , "data: [DONE]\n \n " )
490- if flusher != nil {
491- flusher .Flush ()
492- }
493- return nil
518+ return finish ()
494519 case devinpkg .EventError :
495520 if event .Message != nil {
496521 fillDevinUsage (result , event .Message .Usage )
@@ -518,14 +543,33 @@ func (h *OpenAIGatewayHandler) collectDevinResponse(
518543 result * service.OpenAIForwardResult ,
519544 streamStarted * bool ,
520545) error {
546+ response , err := h .collectDevinCCResponse (ctx , account , stream , reqModel , result )
547+ if err != nil {
548+ return err
549+ }
550+ * streamStarted = true
551+ c .Header ("Content-Type" , "application/json" )
552+ c .JSON (http .StatusOK , response )
553+ return nil
554+ }
555+
556+ // collectDevinCCResponse 排空 devin 流并聚合为 chat.completion JSON map,
557+ // 供 OpenAI 直写或协议桥接(Responses/Anthropic)二次转换。
558+ func (h * OpenAIGatewayHandler ) collectDevinCCResponse (
559+ ctx context.Context ,
560+ account * service.Account ,
561+ stream * devinpkg.ChatStream ,
562+ reqModel string ,
563+ result * service.OpenAIForwardResult ,
564+ ) (map [string ]any , error ) {
521565 var done * devinpkg.AssistantResult
522566 for {
523567 events , err := stream .Next ()
524568 if errors .Is (err , io .EOF ) {
525569 break
526570 }
527571 if err != nil {
528- return err
572+ return nil , err
529573 }
530574 for _ , event := range events {
531575 switch event .Kind {
@@ -539,14 +583,14 @@ func (h *OpenAIGatewayHandler) collectDevinResponse(
539583 h .devinGatewayService .MarkCredentialFailure (ctx , account , failure )
540584 }
541585 if event .Err != nil {
542- return event .Err
586+ return nil , event .Err
543587 }
544- return errors .New ("devin stream error" )
588+ return nil , errors .New ("devin stream error" )
545589 }
546590 }
547591 }
548592 if done == nil {
549- return errors .New ("devin stream ended without a terminal event" )
593+ return nil , errors .New ("devin stream ended without a terminal event" )
550594 }
551595 fillDevinUsage (result , done .Usage )
552596 firstMs := int (result .Duration .Milliseconds ())
@@ -595,8 +639,5 @@ func (h *OpenAIGatewayHandler) collectDevinResponse(
595639 if response ["id" ] == "chatcmpl-" {
596640 response ["id" ] = fmt .Sprintf ("chatcmpl-%d" , time .Now ().UnixNano ())
597641 }
598- * streamStarted = true
599- c .Header ("Content-Type" , "application/json" )
600- c .JSON (http .StatusOK , response )
601- return nil
642+ return response , nil
602643}
0 commit comments