Files
multica/server/internal/integrations/lark/channel_store.go
Bohan Jiang cb6616f530 feat(slack): Socket Mode channel.Channel adapter (MUL-3516) (#4523)
* feat(slack): Socket Mode channel.Channel adapter (MUL-3516)

First slice of the Slack adapter: implements channel.Channel (Type/Connect/Disconnect/Send/Capabilities) over Slack Socket Mode, normalizes inbound events to channel.InboundMessage (DM, channel @mention, thread reply; bot-loop + edit/delete guards), decodes the per-installation config/secret blob, and registers the Factory under TypeSlack. No engine, core, or channel_* schema change. Unit-tested (translation, capabilities, config decode, chunking, Send via httptest). Resolvers + engine wiring + Block Kit binding replier follow.

Co-authored-by: multica-agent <github@multica.ai>

* fix(slack): address adapter review (MUL-3516)

- Propagate InboundHandler errors through dispatchEventsAPI/handleSocketEvent to Connect so an infra failure tears down the connection for Supervisor reconnect/backoff instead of being silently swallowed (ACK still happens first).
- Capabilities: declare only CapText | CapThreadReply; drop CapRichCard/CapAttachment/CapMessageEdit until those Send paths are wired.
- slackChatType: map mpim (multi-party DM) to group, not p2p, so the 'must address bot' filter applies; only 1:1 im is p2p.
- Document the group-addressing decision: explicit @bot mention required in groups; mention-free thread continuation deferred to the session-aware layer.
- Tests: handler-error propagation, slackChatType table, mpim-requires-mention, capabilities negative assertions.

Co-authored-by: multica-agent <github@multica.ai>

* refactor(channel): shared channel-agnostic ChatSession service (MUL-3516)

Extract the session/append//issue machinery — currently locked inside the Feishu-pinned lark.chatSessionService — into a shared engine.ChatSession parameterized by channel_type + session titles, so every IM adapter reuses it instead of re-implementing it. Logic is verbatim (find-or-create session+binding with unique-violation race re-read; append+touch+reply-target+in-tx dedup Mark; /issue parse with bare-command previous-message fallback) but channel-neutral: command-parse source is supplied by the adapter (enrichment is platform-specific). Backed by a narrow SessionQueries interface so it is unit-tested with an in-memory fake (no DB). /issue parser moved to engine.ParseIssueCommand. Next: migrate Feishu onto it and wire Slack's ResolverSet, removing the lark duplicate.

Co-authored-by: multica-agent <github@multica.ai>

* fix(channel): decouple session binding key from outbound target (MUL-3516)

Addresses Elon's round-2 review. engine.ChatSession.EnsureSession previously keyed the binding on a raw chat id (EnsureSessionInput.ChatID), so a resolver wiring Slack straight through would collapse every @bot thread in one channel into a single chat_session and overwrite last_thread_id. Make the API un-misusable:

- EnsureSessionInput.ChatID -> BindingKey: the explicit session-isolation key (Feishu: chat id; Slack DM: channel id; Slack channel: channel id + thread root), documented so a raw threaded-platform chat id is never passed straight through.
- Add EnsureSessionInput.BindingConfig (opaque) persisted on the binding's config column, so the real outbound channel/thread is preserved when BindingKey is composite — outbound routing stays separate from the isolation key.
- channel.sql CreateChannelChatSessionBinding now writes config (additive, uses the existing NOT NULL column; lark caller passes '{}', no schema change, no Feishu regression).
- Tests: TestEnsureSession_ThreadRootIsolation (two thread roots in one channel -> two sessions; same root reuses) and TestEnsureSession_StoresBindingConfig.

No production wiring change yet (per review, the not-yet-wired shared service is an accepted preparatory state); this makes the API correct before Feishu/Slack are migrated onto it.

Co-authored-by: multica-agent <github@multica.ai>

* feat(slack): Slack ResolverSet with thread-root session isolation (MUL-3516)

Wires Slack into the channel-agnostic engine.Router via a ResolverSet built on the generic channel_* queries (installation route by team_id, identity + workspace-membership recheck, two-phase dedup, audit) plus the shared engine.ChatSession. No new query, no schema change.

slackSessionRouting is the per-message isolation rule (Elon round-2 / Niko round-3): a DM is one session per channel; a channel/group message is isolated by thread root (key = channel:threadRoot, root = inbound thread_ts or the message ts for a top-level @mention), so two @bot threads in one channel are two sessions. The real channel id rides in BindingConfig for outbound; the reply thread is returned separately. Tests cover DM/channel/thread routing, config, and that distinct thread roots isolate while a same-thread follow-up reuses its key.

Not yet wired into router.go (still a preparatory commit, per review); Feishu migration onto the shared service, router/config wiring, and the Slack outbound path follow.

Co-authored-by: multica-agent <github@multica.ai>

* feat(slack): Markdown->mrkdwn outbound formatting (MUL-3516)

Slack renders mrkdwn, not Markdown, so an unconverted agent reply shows literal ** , ## and [text](url). Add formatMrkdwn — a faithful Go port of Hermes Agent's slack format_message (MIT) — and apply it in slackChannel.Send before chunking/posting. Protects fenced+inline code, converted links, and existing Slack entities behind placeholders; converts headers/bold/italic/strike/links; escapes control chars. Unit tests cover each construct plus fenced-code protection and a link nested in bold.

Co-authored-by: multica-agent <github@multica.ai>

* docs(slack): preserve Hermes MIT notice for ported mrkdwn converter (MUL-3516)

Addresses Niko's review. formatMrkdwn is a substantial port of Hermes Agent's slack format_message; MIT requires preserving the copyright + permission notice. Add the full Hermes MIT copyright/permission notice + source URL as a header on mrkdwn.go (no repo-level third-party notice file exists, and the header cannot get separated from the ported code). Also add the suggested Send-layer regression test (TestSend_AppliesMrkdwn) that pins the wiring: slackChannel.Send converts Markdown to mrkdwn before posting.

Co-authored-by: multica-agent <github@multica.ai>

* refactor(lark): migrate Feishu onto shared engine.ChatSession, drop duplicate (MUL-3516)

Completes 'every IM reuses one shared session service' and removes the dual-path the reviewers flagged as temporary. Feishu's ResolverSet now drives the channel-agnostic engine.ChatSession (channel_type=feishu, Lark session titles preserved) instead of the Feishu-specific lark.chatSessionService, which is deleted. Behavior is unchanged: engine.ChatSession is the verbatim port of the old logic and is unit-tested; the new Feishu binder param-mapping (BindingKey=chat id, CommandText=un-enriched CommandBody from Raw) is covered by feishu_resolvers_test.go.

- Delete chat_service.go (chatSessionService + helpers) and issue_command.go/_test.go (parser now engine.ParseIssueCommand). Relocate the shared TxStarter interface to tx.go (still used by binding-token + registration services).
- chat.go keeps only the AuditLogger seam; remove the now-dead ChatSessionService / EnsureChatSessionParams / AppendUserMessageParams / AppendResult / IssueCommand types.
- router.go constructs engine.NewChatSession for Feishu; inbound_enricher_test + doc.go updated.

make-test parity: go build ./..., go vet, gofmt, and go test ./internal/integrations/{lark,channel/...,slack} all pass (full Feishu suite green).

Co-authored-by: multica-agent <github@multica.ai>

* feat(slack): wire Slack adapter + ResolverSet + outbound into router (MUL-3516)

Activates the full Slack pipeline, gated by MULTICA_SLACK_SECRET_KEY (the bot/app-token decryption key). When unset the block is skipped, so existing deployments are unaffected and Feishu is untouched.

- router.go registers slack.RegisterSlack (Socket Mode connect/send Factory) + channelRouter.Register(TypeSlack, NewSlackResolverSet) (inbound pipeline) + slack.NewOutbound(...).Register(bus) (outbound).
- New slack/outbound.go: an EventChatDone subscriber mirroring the Feishu Patcher. It finds the Slack chat binding for the finished session, recovers the real channel from the binding config (the channel_chat_id may be a composite thread-isolation key) + the reply thread from last_thread_id, and posts via slackChannel.Send (reusing formatMrkdwn / chunking / threading). Sessions with no Slack binding are ignored, so it coexists with the Feishu Patcher on the shared bus.
- Tests: posts to the bound channel/thread with the real channel id; ignores non-Slack sessions, empty completions, revoked installations, and non-chat events.

Slack now shares engine.ChatSession, channel_* tables, IssueService and TaskService with Feishu. Remaining: config-driven installation provisioning (an operator currently creates the channel_type='slack' row; the config block shape — which workspace/agent — is a product decision) and a live end-to-end smoke. go build ./..., go vet, gofmt, and go test ./internal/integrations/{slack,channel/...,lark} all pass.

Co-authored-by: multica-agent <github@multica.ai>

---------

Co-authored-by: J <j@multica.ai>
Co-authored-by: multica-agent <github@multica.ai>
2026-06-25 14:29:00 +08:00

402 lines
14 KiB
Go

package lark
// ChannelStore is the production data layer for the Feishu integration after
// MUL-3515 generalized lark_* into channel_*. It embeds *db.Queries (so every
// generic query — chat_session, chat_message, member, workspace, agent — is
// available unchanged) and adds the feishu-specific store methods, each backed
// by a channel_* query and translating at the JSONB-config boundary (store.go).
//
// The methods take and return the package's flat domain types (Installation,
// UserBinding, ChatSessionBinding, InboundMessageDedup, BindingTokenRow,
// OutboundCardMessage) and the *Params types in params.go. This store reads
// and writes only channel_*, never lark_*; queries/lark.sql is deleted. The
// physical lark_* tables are retained one release for rollout/rollback safety
// (see migration 124's ROLLOUT note) and dropped by a later cleanup migration.
import (
"context"
"errors"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype"
db "github.com/multica-ai/multica/server/pkg/db/generated"
)
// channelTypeFeishu is the channel_type discriminator for every row this
// Feishu-backed store reads or writes.
const channelTypeFeishu = "feishu"
type ChannelStore struct {
*db.Queries
}
// NewChannelStore wraps a *db.Queries so the lark package's DB seams resolve to
// channel_* rows.
func NewChannelStore(q *db.Queries) *ChannelStore {
return &ChannelStore{Queries: q}
}
// WithTx returns a ChannelStore bound to tx. It shadows db.Queries.WithTx (which
// returns *db.Queries) so transactional callers (chat ingest, token redemption)
// keep the channel-backed store methods inside their tx.
func (s *ChannelStore) WithTx(tx pgx.Tx) *ChannelStore {
return &ChannelStore{Queries: s.Queries.WithTx(tx)}
}
// IsWorkspaceMember reports whether userID is currently a member of
// workspaceID. With the lark_user_binding -> member foreign key removed
// (MUL-3515 §4), a binding row no longer proves membership, so the inbound
// identity step calls this to re-check it explicitly. ErrNoRows -> not a member.
func (s *ChannelStore) IsWorkspaceMember(ctx context.Context, workspaceID, userID pgtype.UUID) (bool, error) {
_, err := s.Queries.GetMemberByUserAndWorkspace(ctx, db.GetMemberByUserAndWorkspaceParams{
UserID: userID,
WorkspaceID: workspaceID,
})
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return false, nil
}
return false, err
}
return true, nil
}
// ---- installation ----
func (s *ChannelStore) GetLarkInstallationByAppID(ctx context.Context, appID string) (Installation, error) {
row, err := s.Queries.GetChannelInstallationByAppID(ctx, db.GetChannelInstallationByAppIDParams{
ChannelType: channelTypeFeishu,
AppID: appID,
})
if err != nil {
return Installation{}, err
}
return installationFromRow(row)
}
func (s *ChannelStore) GetLarkInstallation(ctx context.Context, id pgtype.UUID) (Installation, error) {
row, err := s.Queries.GetChannelInstallation(ctx, db.GetChannelInstallationParams{
ID: id,
ChannelType: channelTypeFeishu,
})
if err != nil {
return Installation{}, err
}
return installationFromRow(row)
}
func (s *ChannelStore) GetLarkInstallationInWorkspace(ctx context.Context, arg GetInstallationInWorkspaceParams) (Installation, error) {
row, err := s.Queries.GetChannelInstallationInWorkspace(ctx, db.GetChannelInstallationInWorkspaceParams{
ID: arg.ID,
WorkspaceID: arg.WorkspaceID,
ChannelType: channelTypeFeishu,
})
if err != nil {
return Installation{}, err
}
return installationFromRow(row)
}
func (s *ChannelStore) ListLarkInstallationsByWorkspace(ctx context.Context, workspaceID pgtype.UUID) ([]Installation, error) {
rows, err := s.Queries.ListChannelInstallationsByWorkspace(ctx, db.ListChannelInstallationsByWorkspaceParams{
WorkspaceID: workspaceID,
ChannelType: channelTypeFeishu,
})
if err != nil {
return nil, err
}
return installationsFromRows(rows)
}
func (s *ChannelStore) ListActiveLarkInstallations(ctx context.Context) ([]Installation, error) {
rows, err := s.Queries.ListActiveChannelInstallations(ctx, channelTypeFeishu)
if err != nil {
return nil, err
}
return installationsFromRows(rows)
}
func (s *ChannelStore) UpsertLarkInstallation(ctx context.Context, arg UpsertInstallationParams) (Installation, error) {
cfg, err := encodeInstallConfig(Installation{
AppID: arg.AppID,
AppSecretEncrypted: arg.AppSecretEncrypted,
TenantKey: arg.TenantKey,
BotOpenID: arg.BotOpenID,
BotUnionID: arg.BotUnionID,
Region: arg.Region,
})
if err != nil {
return Installation{}, err
}
row, err := s.Queries.UpsertChannelInstallation(ctx, db.UpsertChannelInstallationParams{
WorkspaceID: arg.WorkspaceID,
AgentID: arg.AgentID,
ChannelType: channelTypeFeishu,
Config: cfg,
InstallerUserID: arg.InstallerUserID,
})
if err != nil {
return Installation{}, err
}
return installationFromRow(row)
}
func (s *ChannelStore) SetLarkInstallationStatus(ctx context.Context, arg SetInstallationStatusParams) error {
return s.Queries.SetChannelInstallationStatus(ctx, db.SetChannelInstallationStatusParams{
ID: arg.ID,
Status: arg.Status,
})
}
// SetLarkInstallationBotUnionID folds bot_union_id into the JSONB config via a
// read-modify-write through SetChannelInstallationConfig (channel_installation
// has no dedicated union_id column). This is the operator union_id backfill,
// keyed by id and effectively single-writer, so the non-atomic RMW is safe —
// the same shape the channel.sql comment documents for this query.
func (s *ChannelStore) SetLarkInstallationBotUnionID(ctx context.Context, arg SetInstallationBotUnionIDParams) error {
row, err := s.Queries.GetChannelInstallation(ctx, db.GetChannelInstallationParams{
ID: arg.ID,
ChannelType: channelTypeFeishu,
})
if err != nil {
return err
}
inst, err := installationFromRow(row)
if err != nil {
return err
}
inst.BotUnionID = arg.BotUnionID
cfg, err := encodeInstallConfig(inst)
if err != nil {
return err
}
return s.Queries.SetChannelInstallationConfig(ctx, db.SetChannelInstallationConfigParams{
ID: arg.ID,
Config: cfg,
})
}
func (s *ChannelStore) BackfillLarkInstallationRegionToLark(ctx context.Context) (int64, error) {
return s.Queries.BackfillChannelInstallationRegionToFeishuLark(ctx)
}
// ---- WS lease ----
func (s *ChannelStore) AcquireLarkWSLease(ctx context.Context, arg AcquireWSLeaseParams) (Installation, error) {
row, err := s.Queries.AcquireChannelWSLease(ctx, db.AcquireChannelWSLeaseParams{
NewToken: arg.NewToken,
NewExpiresAt: arg.NewExpiresAt,
ID: arg.ID,
})
if err != nil {
return Installation{}, err
}
return installationFromRow(row)
}
func (s *ChannelStore) ReleaseLarkWSLease(ctx context.Context, arg ReleaseWSLeaseParams) error {
return s.Queries.ReleaseChannelWSLease(ctx, db.ReleaseChannelWSLeaseParams{
ID: arg.ID,
CurrentToken: arg.CurrentToken,
})
}
// ---- user binding ----
func (s *ChannelStore) GetLarkUserBindingByOpenID(ctx context.Context, arg GetUserBindingByOpenIDParams) (UserBinding, error) {
row, err := s.Queries.GetChannelUserBindingByUserID(ctx, db.GetChannelUserBindingByUserIDParams{
InstallationID: arg.InstallationID,
ChannelUserID: arg.ChannelUserID,
})
if err != nil {
return UserBinding{}, err
}
return userBindingFromRow(row)
}
func (s *ChannelStore) CreateLarkUserBinding(ctx context.Context, arg CreateUserBindingParams) (UserBinding, error) {
cfg, err := encodeBindingConfig(UserBinding{UnionID: arg.UnionID})
if err != nil {
return UserBinding{}, err
}
row, err := s.Queries.CreateChannelUserBinding(ctx, db.CreateChannelUserBindingParams{
WorkspaceID: arg.WorkspaceID,
MulticaUserID: arg.MulticaUserID,
InstallationID: arg.InstallationID,
ChannelType: channelTypeFeishu,
ChannelUserID: arg.ChannelUserID,
Config: cfg,
})
if err != nil {
return UserBinding{}, err
}
return userBindingFromRow(row)
}
// ---- chat session binding ----
func (s *ChannelStore) GetLarkChatSessionBinding(ctx context.Context, arg GetChatSessionBindingParams) (ChatSessionBinding, error) {
row, err := s.Queries.GetChannelChatSessionBinding(ctx, db.GetChannelChatSessionBindingParams{
InstallationID: arg.InstallationID,
ChannelChatID: arg.ChannelChatID,
})
if err != nil {
return ChatSessionBinding{}, err
}
return chatSessionBindingFromRow(row), nil
}
func (s *ChannelStore) GetLarkChatSessionBindingBySession(ctx context.Context, chatSessionID pgtype.UUID) (ChatSessionBinding, error) {
row, err := s.Queries.GetChannelChatSessionBindingBySession(ctx, db.GetChannelChatSessionBindingBySessionParams{
ChatSessionID: chatSessionID,
ChannelType: channelTypeFeishu,
})
if err != nil {
return ChatSessionBinding{}, err
}
return chatSessionBindingFromRow(row), nil
}
func (s *ChannelStore) CreateLarkChatSessionBinding(ctx context.Context, arg CreateChatSessionBindingParams) (ChatSessionBinding, error) {
row, err := s.Queries.CreateChannelChatSessionBinding(ctx, db.CreateChannelChatSessionBindingParams{
ChatSessionID: arg.ChatSessionID,
InstallationID: arg.InstallationID,
ChannelType: channelTypeFeishu,
ChannelChatID: arg.ChannelChatID,
ChatType: arg.ChatType,
// Feishu's channel_chat_id is the real chat id, so the key alone routes
// outbound; config stays the empty object (the column is NOT NULL).
Config: []byte("{}"),
})
if err != nil {
return ChatSessionBinding{}, err
}
return chatSessionBindingFromRow(row), nil
}
func (s *ChannelStore) UpdateLarkChatSessionBindingReplyTarget(ctx context.Context, arg UpdateChatSessionBindingReplyTargetParams) error {
return s.Queries.UpdateChannelChatSessionBindingReplyTarget(ctx, db.UpdateChannelChatSessionBindingReplyTargetParams{
ChatSessionID: arg.ChatSessionID,
LastMessageID: arg.LastMessageID,
LastThreadID: arg.LastThreadID,
})
}
// ---- inbound dedup ----
func (s *ChannelStore) ClaimLarkInboundDedup(ctx context.Context, arg ClaimInboundDedupParams) (InboundMessageDedup, error) {
row, err := s.Queries.ClaimChannelInboundDedup(ctx, db.ClaimChannelInboundDedupParams{
InstallationID: arg.InstallationID,
MessageID: arg.MessageID,
})
if err != nil {
return InboundMessageDedup{}, err
}
return dedupFromRow(row), nil
}
func (s *ChannelStore) MarkLarkInboundDedupProcessed(ctx context.Context, arg MarkInboundDedupProcessedParams) (int64, error) {
return s.Queries.MarkChannelInboundDedupProcessed(ctx, db.MarkChannelInboundDedupProcessedParams{
InstallationID: arg.InstallationID,
MessageID: arg.MessageID,
ClaimToken: arg.ClaimToken,
})
}
func (s *ChannelStore) ReleaseLarkInboundDedup(ctx context.Context, arg ReleaseInboundDedupParams) (int64, error) {
return s.Queries.ReleaseChannelInboundDedup(ctx, db.ReleaseChannelInboundDedupParams{
InstallationID: arg.InstallationID,
MessageID: arg.MessageID,
ClaimToken: arg.ClaimToken,
})
}
// ---- audit ----
func (s *ChannelStore) RecordLarkInboundDrop(ctx context.Context, arg RecordInboundDropParams) error {
return s.Queries.RecordChannelInboundDrop(ctx, db.RecordChannelInboundDropParams{
ChannelType: channelTypeFeishu,
EventType: arg.EventType,
DropReason: arg.DropReason,
InstallationID: arg.InstallationID,
ChannelChatID: arg.ChannelChatID,
ChannelEventID: arg.ChannelEventID,
ChannelMessageID: arg.ChannelMessageID,
})
}
// ---- binding token ----
func (s *ChannelStore) CreateLarkBindingToken(ctx context.Context, arg CreateBindingTokenParams) (BindingTokenRow, error) {
row, err := s.Queries.CreateChannelBindingToken(ctx, db.CreateChannelBindingTokenParams{
TokenHash: arg.TokenHash,
WorkspaceID: arg.WorkspaceID,
InstallationID: arg.InstallationID,
ChannelType: channelTypeFeishu,
ChannelUserID: arg.ChannelUserID,
ExpiresAt: arg.ExpiresAt,
})
if err != nil {
return BindingTokenRow{}, err
}
return bindingTokenFromRow(row), nil
}
func (s *ChannelStore) ConsumeLarkBindingToken(ctx context.Context, tokenHash string) (BindingTokenRow, error) {
row, err := s.Queries.ConsumeChannelBindingToken(ctx, tokenHash)
if err != nil {
return BindingTokenRow{}, err
}
return bindingTokenFromRow(row), nil
}
// ---- outbound card ----
func (s *ChannelStore) GetLarkOutboundCardByTask(ctx context.Context, taskID pgtype.UUID) (OutboundCardMessage, error) {
row, err := s.Queries.GetChannelOutboundCardByTask(ctx, db.GetChannelOutboundCardByTaskParams{
TaskID: taskID,
ChannelType: channelTypeFeishu,
})
if err != nil {
return OutboundCardMessage{}, err
}
return outboundCardFromRow(row), nil
}
func (s *ChannelStore) CreateLarkOutboundCardMessage(ctx context.Context, arg CreateOutboundCardMessageParams) (OutboundCardMessage, error) {
row, err := s.Queries.CreateChannelOutboundCardMessage(ctx, db.CreateChannelOutboundCardMessageParams{
ChatSessionID: arg.ChatSessionID,
ChannelType: channelTypeFeishu,
ChannelChatID: arg.ChannelChatID,
ChannelCardMessageID: arg.ChannelCardMessageID,
Status: arg.Status,
TaskID: arg.TaskID,
})
if err != nil {
return OutboundCardMessage{}, err
}
return outboundCardFromRow(row), nil
}
func (s *ChannelStore) UpdateLarkOutboundCardStatus(ctx context.Context, arg UpdateOutboundCardStatusParams) error {
return s.Queries.UpdateChannelOutboundCardStatus(ctx, db.UpdateChannelOutboundCardStatusParams{
ID: arg.ID,
Status: arg.Status,
})
}
// installationsFromRows maps a slice of channel_installation rows to domain
// Installations, surfacing the first config-decode error.
func installationsFromRows(rows []db.ChannelInstallation) ([]Installation, error) {
out := make([]Installation, len(rows))
for i, row := range rows {
inst, err := installationFromRow(row)
if err != nil {
return nil, err
}
out[i] = inst
}
return out, nil
}