mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-10 17:26:38 +08:00
feat(message-gateway): emit message_gateway.inbound domain events
This commit is contained in:
@@ -0,0 +1,39 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package listener
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/pkg/message_gateway"
|
||||||
|
)
|
||||||
|
|
||||||
|
// EventMessageGatewayInbound is the domain event name for authorized inbound messages.
|
||||||
|
const EventMessageGatewayInbound = "message_gateway.inbound"
|
||||||
|
|
||||||
|
// MessageGatewayInbound is emitted when a bound user sends a private message.
|
||||||
|
type MessageGatewayInbound struct {
|
||||||
|
Msg message_gateway.InboundMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
// MessageGatewayInboundHandler handles inbound messaging events.
|
||||||
|
type MessageGatewayInboundHandler func(ctx context.Context, event MessageGatewayInbound)
|
||||||
|
|
||||||
|
var messageGatewayInboundHandlers []MessageGatewayInboundHandler
|
||||||
|
|
||||||
|
// OnMessageGatewayInbound registers a handler. Call from bootstrap only.
|
||||||
|
func OnMessageGatewayInbound(handler MessageGatewayInboundHandler) {
|
||||||
|
messageGatewayInboundHandlers = append(messageGatewayInboundHandlers, handler)
|
||||||
|
}
|
||||||
|
|
||||||
|
// EmitMessageGatewayInbound dispatches a bound inbound message.
|
||||||
|
func EmitMessageGatewayInbound(ctx context.Context, msg message_gateway.InboundMessage) {
|
||||||
|
if msg.BindingUserID == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
event := MessageGatewayInbound{Msg: msg}
|
||||||
|
for _, handler := range messageGatewayInboundHandlers {
|
||||||
|
handler(ctx, event)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,25 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package listener
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/pkg/message_gateway"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestEmitMessageGatewayInbound_SkipsUnbound(t *testing.T) {
|
||||||
|
called := 0
|
||||||
|
OnMessageGatewayInbound(func(ctx context.Context, ev MessageGatewayInbound) { called++ })
|
||||||
|
EmitMessageGatewayInbound(context.Background(), message_gateway.InboundMessage{Text: "x"})
|
||||||
|
if called != 0 {
|
||||||
|
t.Fatal("unbound must not emit")
|
||||||
|
}
|
||||||
|
uid := uint64(9)
|
||||||
|
EmitMessageGatewayInbound(context.Background(), message_gateway.InboundMessage{BindingUserID: &uid, Text: "x"})
|
||||||
|
if called != 1 {
|
||||||
|
t.Fatalf("called=%d", called)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -15,6 +15,7 @@ import (
|
|||||||
admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events"
|
"github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events"
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/listener"
|
||||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||||
taskhandlers "github.com/Rain-kl/Wavelet/internal/infra/task/handlers"
|
taskhandlers "github.com/Rain-kl/Wavelet/internal/infra/task/handlers"
|
||||||
"github.com/Rain-kl/Wavelet/internal/model"
|
"github.com/Rain-kl/Wavelet/internal/model"
|
||||||
@@ -38,10 +39,11 @@ type CacheRegistry struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
registerTasksOnce sync.Once
|
registerTasksOnce sync.Once
|
||||||
registerPushDomainEventsOnce sync.Once
|
registerPushDomainEventsOnce sync.Once
|
||||||
registerTaskListenersOnce sync.Once
|
registerTaskListenersOnce sync.Once
|
||||||
initRuntimeOnce sync.Once
|
registerMessageGatewayListenersOnce sync.Once
|
||||||
|
initRuntimeOnce sync.Once
|
||||||
|
|
||||||
cacheRegistries = make(map[string]CacheRegistry)
|
cacheRegistries = make(map[string]CacheRegistry)
|
||||||
cacheRegistriesMu sync.RWMutex
|
cacheRegistriesMu sync.RWMutex
|
||||||
@@ -99,6 +101,25 @@ func RegisterPushDomainEvents() {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RegisterMessageGatewayListeners registers the default log-only inbound handler.
|
||||||
|
func RegisterMessageGatewayListeners() {
|
||||||
|
registerMessageGatewayListenersOnce.Do(func() {
|
||||||
|
listener.OnMessageGatewayInbound(func(ctx context.Context, event listener.MessageGatewayInbound) {
|
||||||
|
userID := uint64(0)
|
||||||
|
if event.Msg.BindingUserID != nil {
|
||||||
|
userID = *event.Msg.BindingUserID
|
||||||
|
}
|
||||||
|
logger.InfoF(ctx, "[%s] channel=%d type=%s user=%d platform_user=%s",
|
||||||
|
listener.EventMessageGatewayInbound,
|
||||||
|
event.Msg.ChannelID,
|
||||||
|
event.Msg.ChannelType,
|
||||||
|
userID,
|
||||||
|
event.Msg.PlatformUserID,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
// RegisterTaskListeners wires operational listeners to task framework hooks.
|
// RegisterTaskListeners wires operational listeners to task framework hooks.
|
||||||
func RegisterTaskListeners() {
|
func RegisterTaskListeners() {
|
||||||
registerTaskListenersOnce.Do(func() {
|
registerTaskListenersOnce.Do(func() {
|
||||||
@@ -110,12 +131,14 @@ func RegisterTaskListeners() {
|
|||||||
func RegisterAPI() {
|
func RegisterAPI() {
|
||||||
RegisterTasks()
|
RegisterTasks()
|
||||||
RegisterPushDomainEvents()
|
RegisterPushDomainEvents()
|
||||||
|
RegisterMessageGatewayListeners()
|
||||||
}
|
}
|
||||||
|
|
||||||
// RegisterWorker wires integrations required by the task worker process.
|
// RegisterWorker wires integrations required by the task worker process.
|
||||||
func RegisterWorker() {
|
func RegisterWorker() {
|
||||||
RegisterTasks()
|
RegisterTasks()
|
||||||
RegisterTaskListeners()
|
RegisterTaskListeners()
|
||||||
|
RegisterMessageGatewayListeners()
|
||||||
}
|
}
|
||||||
|
|
||||||
// RegisterScheduler wires integrations required by the task scheduler process.
|
// RegisterScheduler wires integrations required by the task scheduler process.
|
||||||
@@ -128,6 +151,7 @@ func RegisterAll() {
|
|||||||
RegisterTasks()
|
RegisterTasks()
|
||||||
RegisterPushDomainEvents()
|
RegisterPushDomainEvents()
|
||||||
RegisterTaskListeners()
|
RegisterTaskListeners()
|
||||||
|
RegisterMessageGatewayListeners()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Init runs shared runtime bootstrap exactly once per process.
|
// Init runs shared runtime bootstrap exactly once per process.
|
||||||
|
|||||||
Reference in New Issue
Block a user