mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 05:56:38 +08:00
chore(message-gateway): swagger and format
This commit is contained in:
@@ -24,7 +24,10 @@ import {
|
||||
SelectTrigger,
|
||||
SelectValue,
|
||||
} from '@/components/ui/select';
|
||||
import { TelegramForm, type TelegramFormValue } from '../channels/telegram/form';
|
||||
import {
|
||||
TelegramForm,
|
||||
type TelegramFormValue,
|
||||
} from '../channels/telegram/form';
|
||||
import { QQForm, type QQFormValue } from '../channels/qq/form';
|
||||
import type { CreateMessageChannelRequest } from '@/lib/services/message-gateway';
|
||||
|
||||
@@ -66,7 +69,9 @@ export function AddChannelDialog({
|
||||
const canSubmit =
|
||||
name.trim() !== '' &&
|
||||
((type === 'telegram' && telegram.bot_token.trim() !== '') ||
|
||||
(type === 'qq' && qq.app_id.trim() !== '' && qq.app_secret.trim() !== ''));
|
||||
(type === 'qq' &&
|
||||
qq.app_id.trim() !== '' &&
|
||||
qq.app_secret.trim() !== ''));
|
||||
|
||||
const handleSubmit = () => {
|
||||
if (!canSubmit || !type) return;
|
||||
|
||||
@@ -39,8 +39,7 @@ export function ChannelCard({
|
||||
}: ChannelCardProps) {
|
||||
const t = useTranslations('admin.messageGateway');
|
||||
const [confirmOpen, setConfirmOpen] = React.useState(false);
|
||||
const typeLabel =
|
||||
channel.type === 'qq' ? t('typeQQ') : t('typeTelegram');
|
||||
const typeLabel = channel.type === 'qq' ? t('typeQQ') : t('typeTelegram');
|
||||
|
||||
return (
|
||||
<div className='rounded-xl border bg-card p-4 space-y-4'>
|
||||
|
||||
@@ -119,7 +119,9 @@ export default function MessageGatewayAdminPage() {
|
||||
key={ch.id}
|
||||
channel={ch}
|
||||
toggling={toggleMutation.isPending}
|
||||
deleting={deleteMutation.isPending && deleteMutation.variables === ch.id}
|
||||
deleting={
|
||||
deleteMutation.isPending && deleteMutation.variables === ch.id
|
||||
}
|
||||
onToggle={(enabled) =>
|
||||
toggleMutation.mutate({ id: ch.id, enabled })
|
||||
}
|
||||
|
||||
@@ -51,7 +51,9 @@ export function BotBindingCard() {
|
||||
UserMessageGatewayService.bind({ channel_id: channelId, code }),
|
||||
onSuccess: () => {
|
||||
toast.success(t('bindSuccess'));
|
||||
queryClient.invalidateQueries({ queryKey: ['message-gateway', 'bindings'] });
|
||||
queryClient.invalidateQueries({
|
||||
queryKey: ['message-gateway', 'bindings'],
|
||||
});
|
||||
setOpen(false);
|
||||
setChannelId('');
|
||||
setCode('');
|
||||
@@ -65,7 +67,9 @@ export function BotBindingCard() {
|
||||
mutationFn: (id: string) => UserMessageGatewayService.unbind(id),
|
||||
onSuccess: () => {
|
||||
toast.success(t('unbindSuccess'));
|
||||
queryClient.invalidateQueries({ queryKey: ['message-gateway', 'bindings'] });
|
||||
queryClient.invalidateQueries({
|
||||
queryKey: ['message-gateway', 'bindings'],
|
||||
});
|
||||
},
|
||||
onError: (err: unknown) => {
|
||||
toast.error(t('unbindFailed') + ': ' + (err as Error).message);
|
||||
@@ -82,8 +86,12 @@ export function BotBindingCard() {
|
||||
<Bot className='size-4' />
|
||||
</div>
|
||||
<div>
|
||||
<h2 className='text-base font-semibold tracking-tight'>{t('title')}</h2>
|
||||
<p className='text-[11px] text-muted-foreground'>{t('description')}</p>
|
||||
<h2 className='text-base font-semibold tracking-tight'>
|
||||
{t('title')}
|
||||
</h2>
|
||||
<p className='text-[11px] text-muted-foreground'>
|
||||
{t('description')}
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
@@ -165,7 +173,9 @@ export function BotBindingCard() {
|
||||
</Select>
|
||||
{!channelsQuery.isPending &&
|
||||
(channelsQuery.data ?? []).length === 0 ? (
|
||||
<p className='text-xs text-muted-foreground'>{t('noChannels')}</p>
|
||||
<p className='text-xs text-muted-foreground'>
|
||||
{t('noChannels')}
|
||||
</p>
|
||||
) : null}
|
||||
</div>
|
||||
<div className='space-y-2'>
|
||||
|
||||
@@ -85,7 +85,11 @@ const adminItems: NavItem[] = [
|
||||
{ titleKey: 'storage', url: '/admin/files', icon: FolderOpen },
|
||||
{ titleKey: 'database', url: '/admin/database', icon: Database },
|
||||
{ titleKey: 'push', url: '/admin/push', icon: Bell },
|
||||
{ titleKey: 'messageGateway', url: '/admin/message-gateway', icon: MessagesSquare },
|
||||
{
|
||||
titleKey: 'messageGateway',
|
||||
url: '/admin/message-gateway',
|
||||
icon: MessagesSquare,
|
||||
},
|
||||
{ titleKey: 'logs', url: '/admin/logs', icon: Terminal },
|
||||
{ titleKey: 'system', url: '/admin/system', icon: ShieldCheck },
|
||||
{ titleKey: 'adminSettings', url: '/admin/settings', icon: Settings },
|
||||
|
||||
@@ -6,8 +6,8 @@ package message_gateway
|
||||
const (
|
||||
errNameRequired = "name is required"
|
||||
errTypeInvalid = "type must be telegram or qq"
|
||||
errTelegramTokenRequired = "bot_token is required"
|
||||
errQQCredentialsRequired = "app_id and app_secret are required"
|
||||
errTelegramTokenRequired = "telegram bot secret is required" //nolint:gosec // user-facing validation text
|
||||
errQQCredentialsRequired = "qq app id and secret are required" //nolint:gosec // user-facing validation text
|
||||
errChannelNotFound = "channel not found"
|
||||
errChannelProbeFailed = "channel probe failed"
|
||||
maskedSecret = "********"
|
||||
|
||||
@@ -318,8 +318,9 @@ func probeCredentials(ctx context.Context, typ string, creds, extra map[string]s
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
||||
defer func() { _ = resp.Body.Close() }()
|
||||
const probeBodyLimit = 4096
|
||||
body, _ := io.ReadAll(io.LimitReader(resp.Body, probeBodyLimit))
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("telegram getMe status %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package runner starts message-gateway adapters in the Worker process.
|
||||
package runner
|
||||
|
||||
import (
|
||||
|
||||
@@ -74,7 +74,6 @@ type gateway struct {
|
||||
node string
|
||||
mu sync.Mutex
|
||||
running map[uint64]*runningChannel
|
||||
seen map[uint64]time.Time
|
||||
}
|
||||
|
||||
func (r *gateway) sync(ctx context.Context) error {
|
||||
|
||||
@@ -15,9 +15,9 @@ import (
|
||||
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/risk_control"
|
||||
"github.com/Rain-kl/Wavelet/internal/listener"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
taskhandlers "github.com/Rain-kl/Wavelet/internal/infra/task/handlers"
|
||||
"github.com/Rain-kl/Wavelet/internal/listener"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/platform/lifecycle"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
@@ -39,11 +39,11 @@ type CacheRegistry struct {
|
||||
}
|
||||
|
||||
var (
|
||||
registerTasksOnce sync.Once
|
||||
registerPushDomainEventsOnce sync.Once
|
||||
registerTaskListenersOnce sync.Once
|
||||
registerTasksOnce sync.Once
|
||||
registerPushDomainEventsOnce sync.Once
|
||||
registerTaskListenersOnce sync.Once
|
||||
registerMessageGatewayListenersOnce sync.Once
|
||||
initRuntimeOnce sync.Once
|
||||
initRuntimeOnce sync.Once
|
||||
|
||||
cacheRegistries = make(map[string]CacheRegistry)
|
||||
cacheRegistriesMu sync.RWMutex
|
||||
|
||||
@@ -20,9 +20,10 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
configTypeSystem = "system"
|
||||
configValueTrue = "true"
|
||||
configValueFalse = "false"
|
||||
configTypeSystem = "system"
|
||||
configTypeBusiness = "business"
|
||||
configValueTrue = "true"
|
||||
configValueFalse = "false"
|
||||
)
|
||||
|
||||
// SetupTestEnvironment initializes an in-memory SQLite DB, seeds default configurations,
|
||||
@@ -139,7 +140,7 @@ func getSeedConfigsPart1() []model.SystemConfig {
|
||||
{
|
||||
Key: model.ConfigKeyMaxAPIKeysPerUser,
|
||||
Value: "5",
|
||||
Type: "business",
|
||||
Type: configTypeBusiness,
|
||||
Description: "限制每个普通用户可以创建的 API Key 最大数量",
|
||||
},
|
||||
{
|
||||
@@ -300,19 +301,19 @@ func getSeedConfigsPart2() []model.SystemConfig {
|
||||
{
|
||||
Key: model.ConfigKeyLogRetentionDaysPostgres,
|
||||
Value: "30",
|
||||
Type: "business",
|
||||
Type: configTypeBusiness,
|
||||
Description: "PostgreSQL 用户访问日志保留天数",
|
||||
},
|
||||
{
|
||||
Key: model.ConfigKeyLogRetentionDaysSQLite,
|
||||
Value: "30",
|
||||
Type: "business",
|
||||
Type: configTypeBusiness,
|
||||
Description: "SQLite 用户访问日志保留天数",
|
||||
},
|
||||
{
|
||||
Key: model.ConfigKeyLogRetentionDaysClickHouse,
|
||||
Value: "30",
|
||||
Type: "business",
|
||||
Type: configTypeBusiness,
|
||||
Description: "ClickHouse 用户访问日志保留天数",
|
||||
},
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package qq implements the official QQ Bot C2C adapter.
|
||||
package qq
|
||||
|
||||
import (
|
||||
@@ -69,10 +70,11 @@ func (a *Adapter) Connect(ctx context.Context) error {
|
||||
}
|
||||
|
||||
var api openapi.OpenAPI
|
||||
const apiTimeout = 5 * time.Second
|
||||
if strings.EqualFold(strings.TrimSpace(a.cfg.Extra["sandbox"]), "true") {
|
||||
api = botgo.NewSandboxOpenAPI(credentials.AppID, tokSrc).WithTimeout(5 * time.Second)
|
||||
api = botgo.NewSandboxOpenAPI(credentials.AppID, tokSrc).WithTimeout(apiTimeout)
|
||||
} else {
|
||||
api = botgo.NewOpenAPI(credentials.AppID, tokSrc).WithTimeout(5 * time.Second)
|
||||
api = botgo.NewOpenAPI(credentials.AppID, tokSrc).WithTimeout(apiTimeout)
|
||||
}
|
||||
|
||||
wsAP, err := api.WS(ctx, nil, "")
|
||||
@@ -92,7 +94,7 @@ func (a *Adapter) Connect(ctx context.Context) error {
|
||||
text = data.Content
|
||||
id = data.ID
|
||||
}
|
||||
a.handleEvent(qqEvent{Kind: "c2c", UserID: authorID, Text: text, MessageID: id})
|
||||
a.handleEvent(runCtx, qqEvent{Kind: "c2c", UserID: authorID, Text: text, MessageID: id})
|
||||
return nil
|
||||
}))
|
||||
|
||||
@@ -112,7 +114,7 @@ func (a *Adapter) Connect(ctx context.Context) error {
|
||||
}
|
||||
|
||||
// Disconnect stops token refresh and drops further inbound events.
|
||||
func (a *Adapter) Disconnect(ctx context.Context) error {
|
||||
func (a *Adapter) Disconnect(_ context.Context) error {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
a.disconnected = true
|
||||
@@ -138,7 +140,7 @@ func (a *Adapter) Send(ctx context.Context, to message_gateway.Recipient, msg me
|
||||
return err
|
||||
}
|
||||
|
||||
func (a *Adapter) handleEvent(ev qqEvent) {
|
||||
func (a *Adapter) handleEvent(ctx context.Context, ev qqEvent) {
|
||||
if ev.Kind != "c2c" {
|
||||
return
|
||||
}
|
||||
@@ -148,7 +150,7 @@ func (a *Adapter) handleEvent(ev qqEvent) {
|
||||
if disconnected || a.onInbound == nil {
|
||||
return
|
||||
}
|
||||
_ = a.onInbound(context.Background(), message_gateway.InboundMessage{
|
||||
_ = a.onInbound(ctx, message_gateway.InboundMessage{
|
||||
ChannelID: a.cfg.ID,
|
||||
ChannelType: message_gateway.ChannelTypeQQ,
|
||||
PlatformUserID: ev.UserID,
|
||||
|
||||
@@ -16,7 +16,7 @@ func TestHandleEvent_DropsNonC2C(t *testing.T) {
|
||||
got++
|
||||
return nil
|
||||
}}
|
||||
a.handleEvent(qqEvent{Kind: "group", UserID: "u1", Text: "hi"})
|
||||
a.handleEvent(context.Background(), qqEvent{Kind: "group", UserID: "u1", Text: "hi"})
|
||||
if got != 0 {
|
||||
t.Fatal("non-C2C must be ignored")
|
||||
}
|
||||
@@ -28,7 +28,7 @@ func TestHandleEvent_C2CText(t *testing.T) {
|
||||
got = msg
|
||||
return nil
|
||||
}}
|
||||
a.handleEvent(qqEvent{Kind: "c2c", UserID: "openid-1", Text: "hello", MessageID: "m1"})
|
||||
a.handleEvent(context.Background(), qqEvent{Kind: "c2c", UserID: "openid-1", Text: "hello", MessageID: "m1"})
|
||||
if got.Text != "hello" || got.PlatformUserID != "openid-1" || got.ChannelID != 3 {
|
||||
t.Fatalf("%+v", got)
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package telegram implements the Telegram private-chat adapter.
|
||||
package telegram
|
||||
|
||||
import (
|
||||
@@ -53,15 +54,15 @@ func (a *Adapter) Connect(ctx context.Context) error {
|
||||
}
|
||||
a.bot = bot
|
||||
bot.Handle(tele.OnText, func(c tele.Context) error {
|
||||
a.handleTeleMessage(c.Message())
|
||||
a.handleTeleMessage(ctx, c.Message())
|
||||
return nil
|
||||
})
|
||||
bot.Handle(tele.OnPhoto, func(c tele.Context) error {
|
||||
a.handleTeleMessage(c.Message())
|
||||
a.handleTeleMessage(ctx, c.Message())
|
||||
return nil
|
||||
})
|
||||
bot.Handle(tele.OnDocument, func(c tele.Context) error {
|
||||
a.handleTeleMessage(c.Message())
|
||||
a.handleTeleMessage(ctx, c.Message())
|
||||
return nil
|
||||
})
|
||||
go bot.Start()
|
||||
@@ -73,7 +74,7 @@ func (a *Adapter) Connect(ctx context.Context) error {
|
||||
}
|
||||
|
||||
// Disconnect stops the bot.
|
||||
func (a *Adapter) Disconnect(ctx context.Context) error {
|
||||
func (a *Adapter) Disconnect(_ context.Context) error {
|
||||
if a.bot != nil {
|
||||
a.bot.Stop()
|
||||
}
|
||||
@@ -81,7 +82,7 @@ func (a *Adapter) Disconnect(ctx context.Context) error {
|
||||
}
|
||||
|
||||
// Send replies to a private chat.
|
||||
func (a *Adapter) Send(ctx context.Context, to message_gateway.Recipient, msg message_gateway.OutboundMessage) error {
|
||||
func (a *Adapter) Send(_ context.Context, to message_gateway.Recipient, msg message_gateway.OutboundMessage) error {
|
||||
if a.bot == nil {
|
||||
return fmt.Errorf("telegram: not connected")
|
||||
}
|
||||
@@ -93,7 +94,7 @@ func (a *Adapter) Send(ctx context.Context, to message_gateway.Recipient, msg me
|
||||
return err
|
||||
}
|
||||
|
||||
func (a *Adapter) handleTeleMessage(m *tele.Message) {
|
||||
func (a *Adapter) handleTeleMessage(ctx context.Context, m *tele.Message) {
|
||||
if m == nil || m.Chat == nil || m.Chat.Type != tele.ChatPrivate {
|
||||
return
|
||||
}
|
||||
@@ -114,7 +115,7 @@ func (a *Adapter) handleTeleMessage(m *tele.Message) {
|
||||
if a.bot != nil {
|
||||
msg.Attachments = a.downloadMedia(m)
|
||||
}
|
||||
_ = a.onInbound(context.Background(), msg)
|
||||
_ = a.onInbound(ctx, msg)
|
||||
}
|
||||
|
||||
func (a *Adapter) downloadMedia(m *tele.Message) []message_gateway.Attachment {
|
||||
|
||||
@@ -17,7 +17,7 @@ func TestHandleUpdate_DropsGroups(t *testing.T) {
|
||||
got++
|
||||
return nil
|
||||
}}
|
||||
a.handleTeleMessage(&tele.Message{
|
||||
a.handleTeleMessage(context.Background(), &tele.Message{
|
||||
ID: 1,
|
||||
Text: "hi",
|
||||
Chat: &tele.Chat{ID: -100, Type: tele.ChatGroup},
|
||||
@@ -37,7 +37,7 @@ func TestHandleUpdate_PrivateText(t *testing.T) {
|
||||
return nil
|
||||
},
|
||||
}
|
||||
a.handleTeleMessage(&tele.Message{
|
||||
a.handleTeleMessage(context.Background(), &tele.Message{
|
||||
ID: 9,
|
||||
Text: "hi",
|
||||
Chat: &tele.Chat{ID: 42, Type: tele.ChatPrivate},
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package message_gateway defines channel adapters, pairing codes, and inbound types.
|
||||
package message_gateway
|
||||
|
||||
// ChannelTypeTelegram is the Telegram private-chat adapter type.
|
||||
|
||||
Reference in New Issue
Block a user