refactor(channels): replace direct constructors with factory registry in manager
This commit is contained in:
parent
57c1d37c22
commit
383687dc57
1 changed files with 48 additions and 130 deletions
|
|
@ -43,166 +43,84 @@ func NewManager(cfg *config.Config, messageBus *bus.MessageBus) (*Manager, error
|
||||||
return m, nil
|
return m, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// initChannel is a helper that looks up a factory by name and creates the channel.
|
||||||
|
func (m *Manager) initChannel(name, displayName string) {
|
||||||
|
f, ok := getFactory(name)
|
||||||
|
if !ok {
|
||||||
|
logger.WarnCF("channels", "Factory not registered", map[string]interface{}{
|
||||||
|
"channel": displayName,
|
||||||
|
})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
logger.DebugCF("channels", "Attempting to initialize channel", map[string]interface{}{
|
||||||
|
"channel": displayName,
|
||||||
|
})
|
||||||
|
ch, err := f(m.config, m.bus)
|
||||||
|
if err != nil {
|
||||||
|
logger.ErrorCF("channels", "Failed to initialize channel", map[string]interface{}{
|
||||||
|
"channel": displayName,
|
||||||
|
"error": err.Error(),
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
m.channels[name] = ch
|
||||||
|
logger.InfoCF("channels", "Channel enabled successfully", map[string]interface{}{
|
||||||
|
"channel": displayName,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (m *Manager) initChannels() error {
|
func (m *Manager) initChannels() error {
|
||||||
logger.InfoC("channels", "Initializing channel manager")
|
logger.InfoC("channels", "Initializing channel manager")
|
||||||
|
|
||||||
if m.config.Channels.Telegram.Enabled && m.config.Channels.Telegram.Token != "" {
|
if m.config.Channels.Telegram.Enabled && m.config.Channels.Telegram.Token != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize Telegram channel")
|
m.initChannel("telegram", "Telegram")
|
||||||
telegram, err := NewTelegramChannel(m.config, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize Telegram channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["telegram"] = telegram
|
|
||||||
logger.InfoC("channels", "Telegram channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.WhatsApp.Enabled && m.config.Channels.WhatsApp.BridgeURL != "" {
|
if m.config.Channels.WhatsApp.Enabled && m.config.Channels.WhatsApp.BridgeURL != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize WhatsApp channel")
|
m.initChannel("whatsapp", "WhatsApp")
|
||||||
whatsapp, err := NewWhatsAppChannel(m.config.Channels.WhatsApp, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize WhatsApp channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["whatsapp"] = whatsapp
|
|
||||||
logger.InfoC("channels", "WhatsApp channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.Feishu.Enabled {
|
if m.config.Channels.Feishu.Enabled {
|
||||||
logger.DebugC("channels", "Attempting to initialize Feishu channel")
|
m.initChannel("feishu", "Feishu")
|
||||||
feishu, err := NewFeishuChannel(m.config.Channels.Feishu, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize Feishu channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["feishu"] = feishu
|
|
||||||
logger.InfoC("channels", "Feishu channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.Discord.Enabled && m.config.Channels.Discord.Token != "" {
|
if m.config.Channels.Discord.Enabled && m.config.Channels.Discord.Token != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize Discord channel")
|
m.initChannel("discord", "Discord")
|
||||||
discord, err := NewDiscordChannel(m.config.Channels.Discord, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize Discord channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["discord"] = discord
|
|
||||||
logger.InfoC("channels", "Discord channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.MaixCam.Enabled {
|
if m.config.Channels.MaixCam.Enabled {
|
||||||
logger.DebugC("channels", "Attempting to initialize MaixCam channel")
|
m.initChannel("maixcam", "MaixCam")
|
||||||
maixcam, err := NewMaixCamChannel(m.config.Channels.MaixCam, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize MaixCam channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["maixcam"] = maixcam
|
|
||||||
logger.InfoC("channels", "MaixCam channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.QQ.Enabled {
|
if m.config.Channels.QQ.Enabled {
|
||||||
logger.DebugC("channels", "Attempting to initialize QQ channel")
|
m.initChannel("qq", "QQ")
|
||||||
qq, err := NewQQChannel(m.config.Channels.QQ, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize QQ channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["qq"] = qq
|
|
||||||
logger.InfoC("channels", "QQ channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.DingTalk.Enabled && m.config.Channels.DingTalk.ClientID != "" {
|
if m.config.Channels.DingTalk.Enabled && m.config.Channels.DingTalk.ClientID != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize DingTalk channel")
|
m.initChannel("dingtalk", "DingTalk")
|
||||||
dingtalk, err := NewDingTalkChannel(m.config.Channels.DingTalk, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize DingTalk channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["dingtalk"] = dingtalk
|
|
||||||
logger.InfoC("channels", "DingTalk channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.Slack.Enabled && m.config.Channels.Slack.BotToken != "" {
|
if m.config.Channels.Slack.Enabled && m.config.Channels.Slack.BotToken != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize Slack channel")
|
m.initChannel("slack", "Slack")
|
||||||
slackCh, err := NewSlackChannel(m.config.Channels.Slack, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize Slack channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["slack"] = slackCh
|
|
||||||
logger.InfoC("channels", "Slack channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.LINE.Enabled && m.config.Channels.LINE.ChannelAccessToken != "" {
|
if m.config.Channels.LINE.Enabled && m.config.Channels.LINE.ChannelAccessToken != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize LINE channel")
|
m.initChannel("line", "LINE")
|
||||||
line, err := NewLINEChannel(m.config.Channels.LINE, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize LINE channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["line"] = line
|
|
||||||
logger.InfoC("channels", "LINE channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.OneBot.Enabled && m.config.Channels.OneBot.WSUrl != "" {
|
if m.config.Channels.OneBot.Enabled && m.config.Channels.OneBot.WSUrl != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize OneBot channel")
|
m.initChannel("onebot", "OneBot")
|
||||||
onebot, err := NewOneBotChannel(m.config.Channels.OneBot, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize OneBot channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["onebot"] = onebot
|
|
||||||
logger.InfoC("channels", "OneBot channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.WeCom.Enabled && m.config.Channels.WeCom.Token != "" {
|
if m.config.Channels.WeCom.Enabled && m.config.Channels.WeCom.Token != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize WeCom channel")
|
m.initChannel("wecom", "WeCom")
|
||||||
wecom, err := NewWeComBotChannel(m.config.Channels.WeCom, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize WeCom channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["wecom"] = wecom
|
|
||||||
logger.InfoC("channels", "WeCom channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if m.config.Channels.WeComApp.Enabled && m.config.Channels.WeComApp.CorpID != "" {
|
if m.config.Channels.WeComApp.Enabled && m.config.Channels.WeComApp.CorpID != "" {
|
||||||
logger.DebugC("channels", "Attempting to initialize WeCom App channel")
|
m.initChannel("wecom_app", "WeCom App")
|
||||||
wecomApp, err := NewWeComAppChannel(m.config.Channels.WeComApp, m.bus)
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("channels", "Failed to initialize WeCom App channel", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
m.channels["wecom_app"] = wecomApp
|
|
||||||
logger.InfoC("channels", "WeCom App channel enabled successfully")
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
logger.InfoCF("channels", "Channel initialization completed", map[string]any{
|
logger.InfoCF("channels", "Channel initialization completed", map[string]interface{}{
|
||||||
"enabled_channels": len(m.channels),
|
"enabled_channels": len(m.channels),
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -226,11 +144,11 @@ func (m *Manager) StartAll(ctx context.Context) error {
|
||||||
go m.dispatchOutbound(dispatchCtx)
|
go m.dispatchOutbound(dispatchCtx)
|
||||||
|
|
||||||
for name, channel := range m.channels {
|
for name, channel := range m.channels {
|
||||||
logger.InfoCF("channels", "Starting channel", map[string]any{
|
logger.InfoCF("channels", "Starting channel", map[string]interface{}{
|
||||||
"channel": name,
|
"channel": name,
|
||||||
})
|
})
|
||||||
if err := channel.Start(ctx); err != nil {
|
if err := channel.Start(ctx); err != nil {
|
||||||
logger.ErrorCF("channels", "Failed to start channel", map[string]any{
|
logger.ErrorCF("channels", "Failed to start channel", map[string]interface{}{
|
||||||
"channel": name,
|
"channel": name,
|
||||||
"error": err.Error(),
|
"error": err.Error(),
|
||||||
})
|
})
|
||||||
|
|
@ -253,11 +171,11 @@ func (m *Manager) StopAll(ctx context.Context) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
for name, channel := range m.channels {
|
for name, channel := range m.channels {
|
||||||
logger.InfoCF("channels", "Stopping channel", map[string]any{
|
logger.InfoCF("channels", "Stopping channel", map[string]interface{}{
|
||||||
"channel": name,
|
"channel": name,
|
||||||
})
|
})
|
||||||
if err := channel.Stop(ctx); err != nil {
|
if err := channel.Stop(ctx); err != nil {
|
||||||
logger.ErrorCF("channels", "Error stopping channel", map[string]any{
|
logger.ErrorCF("channels", "Error stopping channel", map[string]interface{}{
|
||||||
"channel": name,
|
"channel": name,
|
||||||
"error": err.Error(),
|
"error": err.Error(),
|
||||||
})
|
})
|
||||||
|
|
@ -292,14 +210,14 @@ func (m *Manager) dispatchOutbound(ctx context.Context) {
|
||||||
m.mu.RUnlock()
|
m.mu.RUnlock()
|
||||||
|
|
||||||
if !exists {
|
if !exists {
|
||||||
logger.WarnCF("channels", "Unknown channel for outbound message", map[string]any{
|
logger.WarnCF("channels", "Unknown channel for outbound message", map[string]interface{}{
|
||||||
"channel": msg.Channel,
|
"channel": msg.Channel,
|
||||||
})
|
})
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := channel.Send(ctx, msg); err != nil {
|
if err := channel.Send(ctx, msg); err != nil {
|
||||||
logger.ErrorCF("channels", "Error sending message to channel", map[string]any{
|
logger.ErrorCF("channels", "Error sending message to channel", map[string]interface{}{
|
||||||
"channel": msg.Channel,
|
"channel": msg.Channel,
|
||||||
"error": err.Error(),
|
"error": err.Error(),
|
||||||
})
|
})
|
||||||
|
|
@ -315,13 +233,13 @@ func (m *Manager) GetChannel(name string) (Channel, bool) {
|
||||||
return channel, ok
|
return channel, ok
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *Manager) GetStatus() map[string]any {
|
func (m *Manager) GetStatus() map[string]interface{} {
|
||||||
m.mu.RLock()
|
m.mu.RLock()
|
||||||
defer m.mu.RUnlock()
|
defer m.mu.RUnlock()
|
||||||
|
|
||||||
status := make(map[string]any)
|
status := make(map[string]interface{})
|
||||||
for name, channel := range m.channels {
|
for name, channel := range m.channels {
|
||||||
status[name] = map[string]any{
|
status[name] = map[string]interface{}{
|
||||||
"enabled": true,
|
"enabled": true,
|
||||||
"running": channel.IsRunning(),
|
"running": channel.IsRunning(),
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue