feat: add Gin HTTP API with SSE streaming chat

This commit is contained in:
2026-05-05 11:37:19 +08:00
parent a25de092b1
commit 6bbe013170
8 changed files with 377 additions and 10 deletions
+93
View File
@@ -0,0 +1,93 @@
package api
import (
"encoding/json"
"fmt"
"io"
"net/http"
"github.com/cloudwego/eino/adk"
"github.com/cloudwego/eino/schema"
"github.com/gin-gonic/gin"
)
type ChatHandler struct {
supervisor adk.Agent
}
func NewChatHandler(supervisor adk.Agent) *ChatHandler {
return &ChatHandler{supervisor: supervisor}
}
func (h *ChatHandler) Chat(c *gin.Context) {
var req ChatRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
c.Header("Content-Type", "text/event-stream")
c.Header("Cache-Control", "no-cache")
c.Header("Connection", "keep-alive")
c.Header("X-Accel-Buffering", "no")
iter := h.supervisor.Run(c.Request.Context(), &adk.AgentInput{
Messages: []adk.Message{
schema.UserMessage(req.Message),
},
EnableStreaming: true,
})
c.Stream(func(w io.Writer) bool {
event, ok := iter.Next()
if !ok {
return false
}
if event.Err != nil {
writeSSE(w, "error", map[string]string{"error": event.Err.Error()})
return false
}
if event.Output == nil || event.Output.MessageOutput == nil {
return true
}
mv := event.Output.MessageOutput
if mv.IsStreaming && mv.MessageStream != nil {
stream := mv.MessageStream
for {
msg, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
break
}
if msg.Content != "" {
writeSSE(w, "message", map[string]string{
"agent": event.AgentName,
"content": msg.Content,
"role": string(mv.Role),
})
}
}
} else if mv.Message != nil {
if mv.Message.Content != "" {
writeSSE(w, "message", map[string]string{
"agent": event.AgentName,
"content": mv.Message.Content,
"role": string(mv.Role),
})
}
}
return true
})
}
func writeSSE(w io.Writer, event string, data any) {
b, _ := json.Marshal(data)
fmt.Fprintf(w, "event: %s\ndata: %s\n\n", event, string(b))
}
+27
View File
@@ -0,0 +1,27 @@
package api
import (
"net/http"
"eino-test/config"
"eino-test/internal/rag"
"github.com/gin-gonic/gin"
)
type IndexHandler struct {
pipeline *rag.RAGPipeline
cfg *config.Config
}
func NewIndexHandler(pipeline *rag.RAGPipeline, cfg *config.Config) *IndexHandler {
return &IndexHandler{pipeline: pipeline, cfg: cfg}
}
func (h *IndexHandler) RebuildIndex(c *gin.Context) {
if err := h.pipeline.IndexAllFromDir(c.Request.Context(), h.cfg.DataDir); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"message": "index rebuilt"})
}
+94
View File
@@ -0,0 +1,94 @@
package api
import (
"net/http"
"time"
"eino-test/internal/rag"
"eino-test/internal/store"
"github.com/gin-gonic/gin"
)
type NoteHandler struct {
pipeline *rag.RAGPipeline
}
func NewNoteHandler(pipeline *rag.RAGPipeline) *NoteHandler {
return &NoteHandler{pipeline: pipeline}
}
func (h *NoteHandler) ListNotes(c *gin.Context) {
notes, err := h.pipeline.NoteStore().List(c.Request.Context())
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
var resp []NoteResponse
for _, n := range notes {
resp = append(resp, NoteResponse{
ID: n.ID,
Title: n.Title,
Tags: n.Tags,
CreatedAt: n.CreatedAt.Format(time.RFC3339),
})
}
c.JSON(http.StatusOK, resp)
}
func (h *NoteHandler) GetNote(c *gin.Context) {
id := c.Param("id")
note, err := h.pipeline.NoteStore().GetByID(c.Request.Context(), id)
if err != nil {
c.JSON(http.StatusNotFound, gin.H{"error": "note not found"})
return
}
c.JSON(http.StatusOK, NoteResponse{
ID: note.ID,
Title: note.Title,
Content: note.Content,
Tags: note.Tags,
CreatedAt: note.CreatedAt.Format(time.RFC3339),
})
}
func (h *NoteHandler) CreateNote(c *gin.Context) {
var req NoteRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
note := &store.Note{
ID: req.Title,
Title: req.Title,
Content: req.Content,
Tags: req.Tags,
CreatedAt: time.Now(),
UpdatedAt: time.Now(),
}
if err := h.pipeline.NoteStore().Create(c.Request.Context(), note); err != nil {
c.JSON(http.StatusConflict, gin.H{"error": err.Error()})
return
}
if err := h.pipeline.IndexNote(c.Request.Context(), note); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusCreated, NoteResponse{
ID: note.ID,
Title: note.Title,
Content: note.Content,
Tags: note.Tags,
CreatedAt: note.CreatedAt.Format(time.RFC3339),
})
}
func (h *NoteHandler) DeleteNote(c *gin.Context) {
id := c.Param("id")
if err := h.pipeline.NoteStore().Delete(c.Request.Context(), id); err != nil {
c.JSON(http.StatusNotFound, gin.H{"error": "note not found"})
return
}
h.pipeline.RemoveNote(id)
c.JSON(http.StatusOK, gin.H{"message": "deleted"})
}
+42
View File
@@ -0,0 +1,42 @@
package api
import (
"eino-test/config"
"eino-test/internal/rag"
"github.com/cloudwego/eino/adk"
"github.com/gin-contrib/cors"
"github.com/gin-gonic/gin"
)
func NewRouter(cfg *config.Config, pipeline *rag.RAGPipeline, supervisor adk.Agent) *gin.Engine {
r := gin.Default()
r.Use(cors.New(cors.Config{
AllowOrigins: []string{"*"},
AllowMethods: []string{"GET", "POST", "PUT", "DELETE", "OPTIONS"},
AllowHeaders: []string{"Origin", "Content-Type", "Authorization"},
AllowCredentials: true,
}))
chatHandler := NewChatHandler(supervisor)
noteHandler := NewNoteHandler(pipeline)
indexHandler := NewIndexHandler(pipeline, cfg)
api := r.Group("/api")
{
api.POST("/chat", chatHandler.Chat)
api.GET("/notes", noteHandler.ListNotes)
api.GET("/notes/:id", noteHandler.GetNote)
api.POST("/notes", noteHandler.CreateNote)
api.DELETE("/notes/:id", noteHandler.DeleteNote)
api.POST("/index", indexHandler.RebuildIndex)
}
r.Static("/static", "./web/dist")
r.NoRoute(func(c *gin.Context) {
c.File("./web/dist/index.html")
})
return r
}
+19
View File
@@ -0,0 +1,19 @@
package api
type ChatRequest struct {
Message string `json:"message" binding:"required"`
}
type NoteRequest struct {
Title string `json:"title" binding:"required"`
Content string `json:"content" binding:"required"`
Tags []string `json:"tags"`
}
type NoteResponse struct {
ID string `json:"id"`
Title string `json:"title"`
Content string `json:"content"`
Tags []string `json:"tags"`
CreatedAt string `json:"created_at"`
}