package main import ( "bytes" "context" "encoding/json" "fmt" "net/http" "os" "os/signal" "syscall" "time" discordstub "github.com/yourname/cyrene-ai/platform-bridge/internal/adapter/discord" feishustub "github.com/yourname/cyrene-ai/platform-bridge/internal/adapter/feishu" qqadapter "github.com/yourname/cyrene-ai/platform-bridge/internal/adapter/qq" telegramadapter "github.com/yourname/cyrene-ai/platform-bridge/internal/adapter/telegram" wechatstub "github.com/yourname/cyrene-ai/platform-bridge/internal/adapter/wechat" webhookadapter "github.com/yourname/cyrene-ai/platform-bridge/internal/adapter/webhook" "github.com/yourname/cyrene-ai/platform-bridge/internal/bridge" "github.com/yourname/cyrene-ai/platform-bridge/internal/config" "github.com/yourname/cyrene-ai/platform-bridge/internal/handler" "github.com/yourname/cyrene-ai/platform-bridge/internal/logging" "github.com/yourname/cyrene-ai/platform-bridge/internal/permissions" ) func main() { cfg := config.Load() // Config store for platform adapter configs. configStore, err := config.NewStore("platform_configs.json") if err != nil { fmt.Printf("FATAL: config store: %v\n", err) os.Exit(1) } // Message logger. msgLogger, err := logging.NewLogger("logs") if err != nil { fmt.Printf("FATAL: logger: %v\n", err) os.Exit(1) } defer msgLogger.Close() // Core components. mapper := bridge.NewIdentityMapper() checker := permissions.NewChecker() router := bridge.NewPlatformRouter(mapper, checker) // Seed default identities from environment. seedIdentities(mapper) // Register platform adapters based on stored configs or defaults. adapters := createAdapters(cfg, configStore) for _, a := range adapters { router.RegisterAdapter(a) } // Set message handler with logging. router.SetMessageHandler(func(msg *bridge.UnifiedMessage) (*bridge.UnifiedResponse, error) { // Log incoming. msgLogger.Log(logging.LogEntry{ Timestamp: time.Now(), Direction: "incoming", Platform: msg.Platform, ChannelID: msg.ChannelID, SenderID: msg.SenderID, SenderName: msg.SenderName, Content: msg.Content, ContentType: msg.ContentType, MessageID: msg.MessageID, Success: true, }) response, err := forwardToAICore(cfg, msg) if err != nil { msgLogger.Log(logging.LogEntry{ Timestamp: time.Now(), Direction: "outgoing", Platform: msg.Platform, ChannelID: msg.ChannelID, SenderID: msg.SenderID, Success: false, Error: err.Error(), }) return nil, err } // Log outgoing. for _, rm := range response.Messages { msgLogger.Log(logging.LogEntry{ Timestamp: time.Now(), Direction: "outgoing", Platform: msg.Platform, ChannelID: msg.ChannelID, SenderID: msg.SenderID, SenderName: "Cyrene", Content: rm.Content, ContentType: "text", Success: true, }) } return response, nil }) // Connect all adapters. ctx := context.Background() for _, a := range adapters { if err := a.Connect(ctx); err != nil { fmt.Printf("WARN: connect %s failed: %v\n", a.PlatformName(), err) } else { fmt.Printf("Platform adapter connected: %s\n", a.PlatformName()) } } // Setup HTTP server. mux := http.NewServeMux() bh := handler.NewBridgeHandler(router) bh.RegisterRoutes(mux) // Config and log handlers. ch := handler.NewConfigHandler(configStore, router) ch.RegisterRoutes(mux) lh := handler.NewLogHandler(msgLogger) lh.RegisterRoutes(mux) // Start QQ message reader loop. qq, _ := router.GetAdapter("qq") if qqa, ok := qq.(*qqadapter.Adapter); ok { qqMsgCh := make(chan *qqadapter.OBv11Message, 100) go qqa.ReadMessages(ctx, qqMsgCh) go func() { for msg := range qqMsgCh { response, err := router.RouteMessage("qq", msg) if err != nil { fmt.Printf("[qq] route error: %v\n", err) continue } msgs, err := router.SendResponse(response) if err != nil { fmt.Printf("[qq] send error: %v\n", err) continue } _ = msgs } }() } addr := ":" + cfg.Port srv := &http.Server{Addr: addr, Handler: mux} go func() { fmt.Printf("Platform Bridge listening on port %s\n", cfg.Port) if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { fmt.Printf("FATAL: %v\n", err) os.Exit(1) } }() quit := make(chan os.Signal, 1) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) <-quit fmt.Println("Shutting down Platform Bridge...") shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() for _, a := range adapters { a.Disconnect(shutdownCtx) } srv.Shutdown(shutdownCtx) fmt.Println("Platform Bridge stopped") } // createAdapters builds platform adapters, preferring stored configs over defaults. func createAdapters(cfg *config.Config, store *config.Store) []bridge.PlatformAdapter { allNames := []string{"qq", "telegram", "webhook", "wechat", "feishu", "discord"} var adapters []bridge.PlatformAdapter for _, name := range allNames { stored, _ := store.Get(name) if stored != nil && !stored.Enabled { fmt.Printf("Platform %s is disabled in config, skipping\n", name) continue } var a bridge.PlatformAdapter fields := mergeFields(cfg, name, stored) switch name { case "qq": port := cfg.QQBotPort if p, ok := fields["bot_port"]; ok && p != "" { port = p } a = qqadapter.NewAdapter(port) case "telegram": token := cfg.TelegramToken if t, ok := fields["bot_token"]; ok && t != "" { token = t } webhookURL := cfg.TelegramWebhookURL if w, ok := fields["webhook_url"]; ok && w != "" { webhookURL = w } a = telegramadapter.NewAdapter(token, webhookURL) case "webhook": a = webhookadapter.NewAdapter("webhook") case "wechat": a = wechatstub.NewAdapter() case "feishu": a = feishustub.NewAdapter() case "discord": a = discordstub.NewAdapter() } if a != nil { adapters = append(adapters, a) } } return adapters } // mergeFields returns fields from stored config, falling back to env defaults. func mergeFields(cfg *config.Config, name string, stored *config.PlatformConfig) map[string]string { fields := make(map[string]string) if stored != nil { for k, v := range stored.Fields { fields[k] = v } } // Apply env var defaults if fields are missing. if fields["bot_token"] == "" && cfg.TelegramToken != "" && name == "telegram" { fields["bot_token"] = cfg.TelegramToken } if fields["webhook_url"] == "" && cfg.TelegramWebhookURL != "" && name == "telegram" { fields["webhook_url"] = cfg.TelegramWebhookURL } if fields["bot_port"] == "" && cfg.QQBotPort != "" && name == "qq" { fields["bot_port"] = cfg.QQBotPort } return fields } // forwardToAICore sends a unified message to AI-Core's chat endpoint and returns the response. func forwardToAICore(cfg *config.Config, msg *bridge.UnifiedMessage) (*bridge.UnifiedResponse, error) { reqBody, _ := json.Marshal(map[string]interface{}{ "user_id": msg.SenderID, "session_id": fmt.Sprintf("platform_%s_%s", msg.Platform, msg.ChannelID), "message": msg.Content, "mode": "text", "source": map[string]string{ "platform": msg.Platform, "channel_id": msg.ChannelID, "channel_type": msg.ChannelType, "sender_name": msg.SenderName, }, }) url := cfg.AICoreURL + "/api/v1/chat" req, _ := http.NewRequest("POST", url, bytes.NewReader(reqBody)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Accept", "text/event-stream") if cfg.InternalToken != "" { req.Header.Set("X-Internal-Token", cfg.InternalToken) } client := &http.Client{Timeout: 120 * time.Second} resp, err := client.Do(req) if err != nil { return nil, fmt.Errorf("forward to ai-core failed: %w", err) } defer resp.Body.Close() var result struct { Content string `json:"content"` Error string `json:"error"` } if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { buf := new(bytes.Buffer) buf.ReadFrom(resp.Body) return &bridge.UnifiedResponse{ Messages: []bridge.ResponseMessage{ {DisplayType: "chat", Content: buf.String(), FormatMode: "plain"}, }, Platform: msg.Platform, }, nil } if result.Error != "" { return &bridge.UnifiedResponse{ Messages: []bridge.ResponseMessage{ {DisplayType: "system_info", Content: result.Error, FormatMode: "plain"}, }, Platform: msg.Platform, }, nil } return &bridge.UnifiedResponse{ Messages: []bridge.ResponseMessage{ {DisplayType: "chat", Content: result.Content, FormatMode: "plain"}, }, Platform: msg.Platform, }, nil } // seedIdentities loads default identity mappings. func seedIdentities(m *bridge.IdentityMapper) { if qqAdmin := os.Getenv("QQ_ADMIN_UID"); qqAdmin != "" { m.Register(permissions.PlatformIdentity{ Platform: "qq", PlatformUID: qqAdmin, CyreneUser: "admin", Nickname: "开拓者", PermissionLevel: "admin", }) } if tgAdmin := os.Getenv("TELEGRAM_ADMIN_UID"); tgAdmin != "" { m.Register(permissions.PlatformIdentity{ Platform: "telegram", PlatformUID: tgAdmin, CyreneUser: "admin", Nickname: "开拓者", PermissionLevel: "admin", }) } }