
func NewAgent() *Agent {mcpClient, err := client.NewSSEMCPClient("http://127.0.0.1:8080/sse",)if err != nil {log.Fatalf("Failed to create MCP client: %v", err)}llm, err := openai.New(openai.WithBaseURL("https:///api/xxx"),openai.WithToken("xxxxx"),openai.WithModel("xxxx"),)if err != nil {log.Fatalf("Create client: %v", err)}return &Agent{mcpClient: mcpClient,llm: llm,}}
func NewAgent() *Agent {// Create an MCP client using stdiomcpClient, err := client.NewStdioMCPClient("./stdio/server/server", // Path to an MCP server executablenil, // Additional environment variables if needed)if err != nil {log.Fatalf("Failed to create MCP client: %v", err)}llm, err := openai.New(openai.WithBaseURL("https:///api/xxx"),openai.WithToken("xxxxx"),openai.WithModel("xxxx"),)if err != nil {log.Fatalf("Create client: %v", err)}return &Agent{mcpClient: mcpClient,llm: llm,}}
func (a *Agent) Init(ctx context.Context) {// Start the clientif err := a.mcpClient.Start(ctx); err != nil {fmt.Printf("Failed to start client: %v", err)}// Create the adapteradapter, err := langchaingo_mcp_adapter.New(a.mcpClient)if err != nil {log.Fatalf("Failed to create adapter: %v", err)}// Get all tools from MCP servermcpTools, err := adapter.Tools()if err != nil {log.Fatalf("Failed to get tools: %v", err)}for _, tool := range mcpTools {fmt.Println("Tools: ", tool.Name(), tool.Description())}// Create a agent with the toolsagent := agents.NewOneShotAgent(a.llm,mcpTools,agents.WithMaxIterations(3),)executor := agents.NewExecutor(agent)a.executor = executorfmt.Println("Agent executor initialized")}
func (a *Agent) Execute(ctx context.Context, question string) string {// Use the agentresult, err := chains.Run(ctx,a.executor,question,)if err != nil {log.Fatalf("Agent execution error: %v", err)}log.Printf("Agent result: %s", result)return result}
w.Header().Set("Content-Type", "text/event-stream")w.Header().Set("Cache-Control", "no-cache")w.Header().Set("connection", "keep-alive")flusher, ok := w.(http.Flusher)if !ok {http.Error(w, "Streaming unsupported!", http.StatusInternalServerError)return}fmt.Fprintf(w, "%s", string(res))flusher.Flush()
async function httpSend(loadingId,currentImplementation, message) {const response = await fetch('/api/chat', {method: 'POST',headers: {'Content-Type': 'text/event-stream','connection': 'keep-alive','cache-control': 'no-cache',},body: JSON.stringify({message: message,implementation: currentImplementation,}),});const data = await response.json();if (data.error) {showError(data.error);} else {addMessage('agent', "http 方式返回数据:"+data.message);}}
w.Header().Set("Content-Type", "text/event-stream")w.Header().Set("Cache-Control", "no-cache")w.Header().Set("connection", "keep-alive")flusher, ok := w.(http.Flusher)if !ok {fmt.Println("Streaming unsupported!")http.Error(w, "Streaming unsupported!", http.StatusInternalServerError)return}fmt.Println("send response", string(res))fmt.Fprintf(w, "data: %s\n\n", string(res))flusher.Flush()
eventSource = new EventSource("/api/sse?msg="+message,{ retry: 50000 });eventSource.onmessage = function(event) {console.log('SSE message:', event.data);addMessage('agent', "sse 方式返回数据:"+event.data);eventSource.close();};
func (s *Server) handleWebSocket(w http.ResponseWriter, r *http.Request) {// Upgrade HTTP connection to WebSocketconn, err := s.upgrader.Upgrade(w, r, nil)if err != nil {s.logger.Printf("WebSocket upgrade failed: %v", err)return}// Get agent ID from query parameteragentID := r.URL.Query().Get("agentId")client := NewWSClient(conn, s.wsHub, agentID)s.wsHub.register <- client// Start read and write pumpsgo client.WritePump()go client.ReadPump()s.logger.Printf("New WebSocket connection established (agentId: %s)", agentID)}
func (c *WSClient) WritePump() {for {select {case message, ok := <-c.send:c.conn.SetWriteDeadline(time.Now().Add(10 * time.Second))if !ok {// Hub closed the channelc.conn.WriteMessage(websocket.CloseMessage, []byte{})return}if err := c.conn.WriteJSON(message); err != nil {c.hub.logger.Printf("Error writing message: %v", err)return}
func (c *WSClient) ReadPump() {for {select {_, message, err := c.conn.ReadMessage()c.send <- &WSMessage{}
ws.send(JSON.stringify({type:'message',content: message,agentId: currentImplementation,}));
ws = new WebSocket(wsUrl);ws.onmessage = (event) => {try {const message = JSON.parse(event.data);handleWebSocketMessage(message);} catch (error) {console.error('Error parsing WebSocket message:', error);}};


文章转载自golang算法架构leetcode技术php,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




