diff --git a/.env.example b/.env.example index a38c4c89a6e..bdfbfe6a8cb 100644 --- a/.env.example +++ b/.env.example @@ -635,6 +635,11 @@ MULTICA_LARK_CALLBACK_BASE_URL= # handshakes. Leave empty to use standard HTTP_PROXY / HTTPS_PROXY / NO_PROXY # environment handling. MULTICA_LARK_WS_PROXY_URL= +# Optional durable group event consumer. Empty disables forwarding. +# JSON: installation_id, app_id, chat_id, endpoint, secret, disposition. +# HTTPS endpoint; random >=32-byte secret; disposition is observe or consume. +# See docs/runbooks/lark-ingress-relay.md for the versioned receiver contract. +MULTICA_LARK_INGRESS_RELAY= # DingTalk bot integration (Settings → Integrations "Bind to DingTalk") # Off until MULTICA_DINGTALK_SECRET_KEY is set — a base64-encoded 32-byte key diff --git a/AGENTS.md b/AGENTS.md index c87d399ee1e..2e05ddf034a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -121,6 +121,8 @@ Workspace-scoped queries filter by `workspace_id`; membership gates access and ` ## Change and Delivery Rules +- Keep installation-specific automation and policy outside the product repository. Prefer existing APIs; any necessary product extension must have a minimal, reusable contract aligned with upstream. + - Keep changes scoped; reuse existing patterns. Code comments are English. - Do not add internal compatibility shims, dual writes, fallback paths, or legacy adapters unless requested. This does not relax API response compatibility above. - New global pre-workspace routes use a single word or `/{noun}/{verb}`, not hyphenated root names. Update `server/internal/handler/reserved_slugs.json`, run `pnpm generate:reserved-slugs`, and commit `packages/core/paths/reserved-slugs.ts` when changing reserved slugs. diff --git a/apps/docs/content/docs/environment-variables.fr.mdx b/apps/docs/content/docs/environment-variables.fr.mdx index 3fb7ecebb74..64b7a4d7444 100644 --- a/apps/docs/content/docs/environment-variables.fr.mdx +++ b/apps/docs/content/docs/environment-variables.fr.mdx @@ -180,6 +180,7 @@ Les limites de débit d'authentification nécessitent `REDIS_URL` ; sans cette v | GitHub | `GITHUB_APP_ID` | Nécessaire pour le statut CI et la possibilité de fusion sur les cartes de PR, ainsi que pour le sélecteur de dépôts « Choisir depuis GitHub » | | GitHub | `GITHUB_APP_PRIVATE_KEY` | Clé privée PEM complète associée à l'App ID ; mêmes usages que ci-dessus | | Lark | `MULTICA_LARK_SECRET_KEY` | Clé de chiffrement des identifiants de 32 octets, encodée en Base64 | +| Lark | `MULTICA_LARK_INGRESS_RELAY` | Configuration JSON facultative pour transmettre les événements de groupe à un récepteur durable ; voir `.env.example` et le guide du relais. | | Slack | `MULTICA_SLACK_SECRET_KEY` | Clé de chiffrement des jetons de 32 octets, encodée en Base64 | | Telegram | `MULTICA_TELEGRAM_SECRET_KEY` | Clé de chiffrement du jeton du bot de 32 octets, encodée en Base64 | | Composio | `COMPOSIO_API_KEY` | Active les connexions d'outils Composio | diff --git a/apps/docs/content/docs/environment-variables.mdx b/apps/docs/content/docs/environment-variables.mdx index 871e6660d93..2fd2cd00535 100644 --- a/apps/docs/content/docs/environment-variables.mdx +++ b/apps/docs/content/docs/environment-variables.mdx @@ -180,6 +180,7 @@ The auth rate limits require `REDIS_URL`; without it, the startup log notes that | GitHub | `GITHUB_APP_ID` | Needed for CI status and mergeability on PR cards and the "pick from GitHub" repository picker | | GitHub | `GITHUB_APP_PRIVATE_KEY` | Full PEM private key paired with the App ID; same uses as above | | Lark | `MULTICA_LARK_SECRET_KEY` | Base64-encoded 32-byte credential encryption key | +| Lark | `MULTICA_LARK_INGRESS_RELAY` | Optional JSON configuration for a durable group event receiver; see `.env.example` and the ingress relay runbook. | | Slack | `MULTICA_SLACK_SECRET_KEY` | Base64-encoded 32-byte token encryption key | | Telegram | `MULTICA_TELEGRAM_SECRET_KEY` | Base64-encoded 32-byte Bot token encryption key | | Composio | `COMPOSIO_API_KEY` | Enables Composio tool connections | diff --git a/apps/docs/content/docs/environment-variables.zh.mdx b/apps/docs/content/docs/environment-variables.zh.mdx index 59df60d2f5a..441d185353c 100644 --- a/apps/docs/content/docs/environment-variables.zh.mdx +++ b/apps/docs/content/docs/environment-variables.zh.mdx @@ -180,6 +180,7 @@ MULTICA_PUBLIC_URL=https://api.multica.example.com | GitHub | `GITHUB_APP_ID` | PR 卡片的 CI 状态、可合并性和"从 GitHub 选择"仓库都需要它 | | GitHub | `GITHUB_APP_PRIVATE_KEY` | 与 App ID 配套的完整 PEM 私钥,用途同上 | | 飞书 | `MULTICA_LARK_SECRET_KEY` | base64 编码的 32 字节凭据加密密钥 | +| Lark | `MULTICA_LARK_INGRESS_RELAY` | 可选 JSON 配置,将群消息转发到持久化接收端;参见 `.env.example` 和转发配置指南。 | | Slack | `MULTICA_SLACK_SECRET_KEY` | base64 编码的 32 字节 token 加密密钥 | | Telegram | `MULTICA_TELEGRAM_SECRET_KEY` | base64 编码的 32 字节 Bot token 加密密钥 | | Composio | `COMPOSIO_API_KEY` | 启用 Composio 工具连接 | diff --git a/docs/runbooks/lark-ingress-relay.md b/docs/runbooks/lark-ingress-relay.md new file mode 100644 index 00000000000..152f5d09ba6 --- /dev/null +++ b/docs/runbooks/lark-ingress-relay.md @@ -0,0 +1,36 @@ +# Lark group ingress relay + +The optional relay forwards literal messages from one configured group to an operator-owned durable receiver. It reuses the native connection. It introduces no worker, storage schema, model, business policy or document parser. + +Leave `MULTICA_LARK_INGRESS_RELAY` empty to retain current behavior. Configure a JSON object with exactly these fields: + +| Field | Meaning | +| --- | --- | +| `installation_id` | Active native installation UUID | +| `app_id` | The installation's app ID (`cli_...`) | +| `chat_id` | The opted-in group (`oc_...`) | +| `endpoint` | Trusted HTTPS receiver URL; plaintext is permitted only on loopback for local testing | +| `secret` | Random shared secret of at least 32 bytes; keep it in deployment secrets | +| `disposition` | `observe` forwards then continues native routing; `consume` forwards and stops native routing for the group | + +The receiver must already be available and able to commit before enabling the relay. This configuration does not grant provider permissions or authorize new access. DMs and nonmatching groups retain their native route. In consume mode, the receiver owns matching bot and unsupported events as well; it must acknowledge them deliberately without opening a private-agent route. + +## HTTP contract (version 1) + +The relay POSTs JSON with `version: 1`, `installation_id`, `workspace_id`, `agent_id`, `app_id`, `event_id`, `event_type`, `chat_id`, `chat_type`, `message_id`, `sender_id`, `sender_type`, `message_type`, `content`, `create_time`, `parent_id`, `root_id`, `thread_id`, `addressed_to_bot` and `command_body`. IDs and time values are strings; `addressed_to_bot` is boolean. `content` is the literal provider JSON string. The payload excludes enriched quotes/history, private sessions and installation credentials. All source text is untrusted evidence. + +`X-Multica-Timestamp` is Unix seconds. `X-Multica-Signature` is the lowercase hex HMAC-SHA256 of `timestamp + "." + exact_request_body`, keyed by the shared secret. Verify the signature with constant-time comparison, check a bounded timestamp window and revalidate every configured identity pin. HTTPS authenticates the receiver and protects both messages and receipts. + +Commit the event idempotently before responding. Return a 2xx status with: + +```json +{"version":1,"event_id":"the-request-event-id","durable":true,"disposition":"consume"} +``` + +The receipt disposition must match the configured value, including on duplicate events. An observe receipt cannot reopen native routing when consume is configured. Unsupported version, missing durable confirmation, mismatched ID/mode, non-2xx response or transport failure returns an error to the existing connector, which NACKs and reconnects. There is no fallthrough on failure and no group-root response from this relay. + +Requests are limited to 512 KiB and two seconds; receipts to 4 KiB. Redirects are rejected. Keep the receiver's synchronous work to authentication, validation and durable commit. OCR, model calls and external writes belong after that commit. Provider retries are finite: receivers need independent reconciliation and must tolerate commit-success/response-loss replays. + +## Verification and rollout + +Use a local fake receiver to verify signature, durable duplicate handling, observe/consume matching, timeout and rejection behavior. Verify an active installation and actual durable receipt before an authorized production cutover. Coordinate mode changes; mismatched modes intentionally interrupt delivery. Receiver outages never select observe automatically. Stopping forwarding or returning to observe is an explicit operator routing decision. diff --git a/server/cmd/server/main.go b/server/cmd/server/main.go index 6f6773cf898..b33997c03de 100644 --- a/server/cmd/server/main.go +++ b/server/cmd/server/main.go @@ -22,6 +22,7 @@ import ( "github.com/multica-ai/multica/server/internal/dbstartup" "github.com/multica-ai/multica/server/internal/events" "github.com/multica-ai/multica/server/internal/handler" + "github.com/multica-ai/multica/server/internal/integrations/lark" "github.com/multica-ai/multica/server/internal/integrations/wecom" "github.com/multica-ai/multica/server/internal/logger" "github.com/multica-ai/multica/server/internal/maintenance" @@ -661,6 +662,11 @@ func main() { // Validate the LLM retry budget before the router exists: an operator who // typed a value we cannot honor should see the boot stop, the same way a // malformed feature-flag file does above. + larkIngress, err := lark.NewIngressRelay(os.Getenv("MULTICA_LARK_INGRESS_RELAY")) + if err != nil { + slog.Error("invalid Lark ingress relay configuration") + os.Exit(1) + } llmMaxRetries, err := parseLLMMaxRetries(os.Getenv("MULTICA_LLM_MAX_RETRIES")) if err != nil { slog.Error("invalid MULTICA_LLM_MAX_RETRIES", "error", err) @@ -692,6 +698,7 @@ func main() { FeatureFlags: flags, HeartbeatScheduler: heartbeatScheduler, LLMMaxRetries: llmMaxRetries, + LarkIngress: larkIngress, LLMDisableThinking: llmDisableThinking, }) var replicaQueries *db.Queries diff --git a/server/cmd/server/router.go b/server/cmd/server/router.go index adca0985ddf..16e4d3965be 100644 --- a/server/cmd/server/router.go +++ b/server/cmd/server/router.go @@ -215,6 +215,8 @@ func NewRouter(pool *pgxpool.Pool, hub *realtime.Hub, bus *events.Bus, analytics } type RouterOptions struct { + LarkIngress *lark.IngressRelay + HTTPMetrics *obsmetrics.HTTPMetrics BusinessMetrics *obsmetrics.BusinessMetrics ChannelLeaseMetrics *obsmetrics.ChannelLeaseMetrics @@ -679,9 +681,13 @@ func NewRouterWithOptions(pool *pgxpool.Pool, hub *realtime.Hub, bus *events.Bus Logger: slog.Default(), }) mediaResolver := lark.NewFeishuMediaResolver(larkClient, installSvc, store, engine.NewDBMediaIntentLedger(queries), slog.Default()) - channelRouter.Register(channel.TypeFeishu, lark.NewFeishuResolverSet( + feishuResolvers := lark.NewFeishuResolverSet( cs, feishuSession, auditLogger, resolverReplier, typingIndicator, mediaResolver, - )) + ) + if opts.LarkIngress != nil { + feishuResolvers.Ingress = opts.LarkIngress + } + channelRouter.Register(channel.TypeFeishu, feishuResolvers) slog.Info("lark inbound pipeline wired", "connector", connectorLabel) // One-shot union_id backfill for installations created diff --git a/server/internal/integrations/channel/engine/resolvers.go b/server/internal/integrations/channel/engine/resolvers.go index cd8aeb9e807..209fd028885 100644 --- a/server/internal/integrations/channel/engine/resolvers.go +++ b/server/internal/integrations/channel/engine/resolvers.go @@ -396,11 +396,19 @@ type TypingNotifier interface { OnSettled(ctx context.Context, sessionID pgtype.UUID) } +// InboundInterceptor persists an opt-in source-only route after installation +// validation and before private-agent dedup, mention and identity handling. +// A handled event MUST NOT continue into a private agent session. +type InboundInterceptor interface { + Capture(context.Context, ResolvedInstallation, channel.InboundMessage) (handled bool, err error) +} + // ResolverSet is the per-platform bundle the Router runs the pipeline through. // Installation/Identity/Dedup/Session/Audit are required; Replier/Typing are // optional. OriginType is the issue.origin_type label written for /issue // commands from this channel (Feishu: "lark_chat"). type ResolverSet struct { + Ingress InboundInterceptor Installation InstallationResolver Identity IdentityResolver Dedup Deduper diff --git a/server/internal/integrations/channel/engine/router.go b/server/internal/integrations/channel/engine/router.go index d049cb5a50d..6c18b2a0410 100644 --- a/server/internal/integrations/channel/engine/router.go +++ b/server/internal/integrations/channel/engine/router.go @@ -300,6 +300,18 @@ func (r *Router) dispatch(ctx context.Context, set ResolverSet, msg channel.Inbo return r.drop(ctx, set, msg, inst.ID, DropReasonRevokedInstallation), inst, nil } + // Opt-in durable source capture is independent of private-agent authority. + // Returning success here means the source/job transaction committed. + if set.Ingress != nil { + handled, err := set.Ingress.Capture(ctx, inst, msg) + if err != nil { + return Result{}, inst, fmt.Errorf("capture source: %w", err) + } + if handled { + return Result{Outcome: OutcomeDropped}, inst, nil + } + } + // 2. Two-phase dedup claim with owner fencing — before group filter and // identity so a reconnect replay cannot re-trigger a binding prompt, // re-write a drop audit, or re-touch the session. Empty MessageID diff --git a/server/internal/integrations/channel/engine/router_test.go b/server/internal/integrations/channel/engine/router_test.go index 5de2571f24a..18d08e4f229 100644 --- a/server/internal/integrations/channel/engine/router_test.go +++ b/server/internal/integrations/channel/engine/router_test.go @@ -2245,3 +2245,58 @@ func TestRouter_MediaDeadlineStartsBeforeAppend(t *testing.T) { t.Fatal("resolver did not run") } } + +type captureInterceptor struct { + called int + handled bool + err error +} + +func (c *captureInterceptor) Capture(context.Context, ResolvedInstallation, channel.InboundMessage) (bool, error) { + c.called++ + return c.handled, c.err +} +func TestRouterSourceCaptureBeforePrivateAuthority(t *testing.T) { + for _, tc := range []struct { + name string + active, handled, fail bool + }{ + {"passive group", true, true, false}, {"database failure", true, true, true}, {"capture-only preserves private gate", true, false, false}, {"revoked installation", false, true, false}, + } { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + h.inst.inst.Active = tc.active + h.ident.err = ErrSenderUnbound + capture := &captureInterceptor{handled: tc.handled} + if tc.fail { + capture.err = errors.New("database down") + } + set := h.router.sets[channel.TypeFeishu] + set.Ingress = capture + h.router.Register(channel.TypeFeishu, set) + msg := p2pMessage(t) + msg.Source.ChatType = channel.ChatTypeGroup + msg.AddressedToBot = false + err := h.router.Handle(context.Background(), msg) + if (err != nil) != tc.fail { + t.Fatalf("error=%v", err) + } + wantCalls := 1 + if !tc.active { + wantCalls = 0 + } + if capture.called != wantCalls { + t.Fatalf("capture calls=%d", capture.called) + } + if (tc.handled || !tc.active) && h.dedup.claimCalls != 0 { + t.Fatal("source route reached private dedup") + } + if !tc.handled && h.dedup.claimCalls != 1 { + t.Fatal("capture-only bypassed existing private route") + } + if h.binder.ensureCalls != 0 { + t.Fatal("passive capture created private session") + } + }) + } +} diff --git a/server/internal/integrations/lark/feishu_types.go b/server/internal/integrations/lark/feishu_types.go index 9f9a06c64cb..ee601308601 100644 --- a/server/internal/integrations/lark/feishu_types.go +++ b/server/internal/integrations/lark/feishu_types.go @@ -21,6 +21,7 @@ type InboundMessage struct { ChatType ChatType MessageID string SenderOpenID OpenID + SenderType string Body string // Content is the raw msg_type-specific JSON string Lark sends in // event.message.content. Text/post decoding consumes it immediately; media diff --git a/server/internal/integrations/lark/ingress_relay.go b/server/internal/integrations/lark/ingress_relay.go new file mode 100644 index 00000000000..77feff4353b --- /dev/null +++ b/server/internal/integrations/lark/ingress_relay.go @@ -0,0 +1,139 @@ +package lark + +import ( + "bytes" + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "io" + "net/http" + "net/url" + "strconv" + "strings" + "time" + + "github.com/multica-ai/multica/server/internal/integrations/channel" + "github.com/multica-ai/multica/server/internal/integrations/channel/engine" + "github.com/multica-ai/multica/server/internal/util" +) + +// IngressRelay forwards one explicitly configured group to a durable consumer. +// It adds no queue: an unavailable receiver fails the connector ACK, allowing +// provider retries. Receivers must reconcile missed events independently. +type IngressRelay struct { + config ingressRelayConfig + client *http.Client +} +type ingressRelayConfig struct { + InstallationID string `json:"installation_id"` + AppID string `json:"app_id"` + ChatID string `json:"chat_id"` + Endpoint string `json:"endpoint"` + Secret string `json:"secret"` + Disposition string `json:"disposition"` +} + +func NewIngressRelay(raw string) (*IngressRelay, error) { + if strings.TrimSpace(raw) == "" { + return nil, nil + } + var c ingressRelayConfig + dec := json.NewDecoder(strings.NewReader(raw)) + dec.DisallowUnknownFields() + if dec.Decode(&c) != nil || dec.Decode(new(any)) != io.EOF { + return nil, errors.New("invalid ingress relay JSON") + } + u, err := url.Parse(c.Endpoint) + _, idErr := util.ParseUUID(c.InstallationID) + if err != nil || u.Host == "" || u.User != nil || u.Fragment != "" || u.RawQuery != "" || (u.Scheme != "https" && !(u.Scheme == "http" && (u.Hostname() == "localhost" || u.Hostname() == "127.0.0.1" || u.Hostname() == "::1"))) || idErr != nil || !strings.HasPrefix(c.AppID, "cli_") || !strings.HasPrefix(c.ChatID, "oc_") || len(c.Secret) < 32 || (c.Disposition != "consume" && c.Disposition != "observe") { + return nil, errors.New("invalid ingress relay target, secret or disposition") + } + return &IngressRelay{config: c, client: &http.Client{Timeout: 2 * time.Second, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}}, nil +} + +// The v1 envelope deliberately contains literal source content only. Enriched +// quotes/history, credentials and private session context never cross this seam. +type ingressEnvelope struct { + Version int `json:"version"` + InstallationID string `json:"installation_id"` + WorkspaceID string `json:"workspace_id"` + AgentID string `json:"agent_id"` + AppID string `json:"app_id"` + EventID string `json:"event_id"` + EventType string `json:"event_type"` + ChatID string `json:"chat_id"` + ChatType string `json:"chat_type"` + MessageID string `json:"message_id"` + SenderID string `json:"sender_id"` + SenderType string `json:"sender_type"` + MessageType string `json:"message_type"` + Content string `json:"content"` + CreateTime string `json:"create_time"` + ParentID string `json:"parent_id"` + RootID string `json:"root_id"` + ThreadID string `json:"thread_id"` + AddressedToBot bool `json:"addressed_to_bot"` + CommandBody string `json:"command_body"` +} + +func (r *IngressRelay) Capture(ctx context.Context, inst engine.ResolvedInstallation, msg channel.InboundMessage) (bool, error) { + if !inst.Active || uuidString(inst.ID) != r.config.InstallationID || msg.Source.ChatType != channel.ChatTypeGroup { + return false, nil + } + m, err := larkMsgFromRaw(msg) + if err != nil { + return false, errors.New("ingress source unavailable") + } + if string(m.ChatID) != r.config.ChatID || m.AppID != r.config.AppID { + return false, nil + } + // Use the server-resolved installation rather than accepting an app identity + // solely from the event payload. Never serialize Platform (it holds secrets). + native, ok := inst.Platform.(Installation) + if !ok || native.AppID != r.config.AppID { + return false, errors.New("ingress installation mismatch") + } + if m.EventID == "" { + return false, errors.New("ingress event ID missing") + } + body, err := json.Marshal(ingressEnvelope{Version: 1, InstallationID: uuidString(inst.ID), WorkspaceID: uuidString(inst.WorkspaceID), AgentID: uuidString(inst.AgentID), AppID: m.AppID, EventID: m.EventID, EventType: m.EventType, ChatID: string(m.ChatID), ChatType: string(m.ChatType), MessageID: m.MessageID, SenderID: string(m.SenderOpenID), SenderType: m.SenderType, MessageType: m.MessageType, Content: m.Content, CreateTime: m.CreateTime, ParentID: m.ParentID, RootID: m.RootID, ThreadID: m.ThreadID, AddressedToBot: m.AddressedToBot, CommandBody: m.CommandBody}) + if err != nil || len(body) > 512<<10 { + return false, errors.New("ingress source exceeds limit") + } + timestamp := strconv.FormatInt(time.Now().Unix(), 10) + mac := hmac.New(sha256.New, []byte(r.config.Secret)) + mac.Write([]byte(timestamp + ".")) + mac.Write(body) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.config.Endpoint, bytes.NewReader(body)) + if err != nil { + return false, errors.New("ingress request invalid") + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("X-Multica-Timestamp", timestamp) + req.Header.Set("X-Multica-Signature", hex.EncodeToString(mac.Sum(nil))) + resp, err := r.client.Do(req) + if err != nil { + return false, errors.New("ingress receiver unavailable") + } + defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return false, errors.New("ingress receiver rejected event") + } + b, err := io.ReadAll(io.LimitReader(resp.Body, 4097)) + if err != nil || len(b) > 4096 { + return false, errors.New("ingress receipt unavailable") + } + var receipt struct { + Version int `json:"version"` + EventID string `json:"event_id"` + Durable bool `json:"durable"` + Disposition string `json:"disposition"` + } + if json.Unmarshal(b, &receipt) != nil || receipt.Version != 1 || receipt.EventID != m.EventID || !receipt.Durable || receipt.Disposition != r.config.Disposition { + return false, errors.New("ingress receipt mismatch") + } + return r.config.Disposition == "consume", nil +} diff --git a/server/internal/integrations/lark/ingress_relay_test.go b/server/internal/integrations/lark/ingress_relay_test.go new file mode 100644 index 00000000000..8794be107c7 --- /dev/null +++ b/server/internal/integrations/lark/ingress_relay_test.go @@ -0,0 +1,91 @@ +package lark + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/multica-ai/multica/server/internal/integrations/channel" + "github.com/multica-ai/multica/server/internal/integrations/channel/engine" + "github.com/multica-ai/multica/server/internal/util" +) + +func TestIngressRelayReceiptAndLiteralBoundary(t *testing.T) { + const id = "11111111-1111-4111-8111-111111111111" + secret := strings.Repeat("s", 32) + for _, tc := range []struct { + name, receipt, disposition string + status int + wantErr, consume bool + }{ + {"consume", `{"version":1,"event_id":"evt","durable":true,"disposition":"consume"}`, "consume", 200, false, true}, + {"observe", `{"version":1,"event_id":"evt","durable":true,"disposition":"observe"}`, "observe", 200, false, false}, + {"mode drift", `{"version":1,"event_id":"evt","durable":true,"disposition":"observe"}`, "consume", 200, true, false}, + {"wrong event", `{"version":1,"event_id":"other","durable":true,"disposition":"consume"}`, "consume", 200, true, false}, + {"not committed", `{"version":1,"event_id":"evt","durable":false,"disposition":"consume"}`, "consume", 200, true, false}, + {"invalid", "{}", "consume", 200, true, false}, + {"unavailable", "{}", "consume", 503, true, false}, + {"redirect", "{}", "consume", 302, true, false}, + } { + t.Run(tc.name, func(t *testing.T) { + calls := 0 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + b, _ := io.ReadAll(r.Body) + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write([]byte(r.Header.Get("X-Multica-Timestamp") + ".")) + mac.Write(b) + if hex.EncodeToString(mac.Sum(nil)) != r.Header.Get("X-Multica-Signature") { + t.Error("signature mismatch") + } + if strings.Contains(string(b), "PRIVATE") { + t.Error("enriched context leaked") + } + w.WriteHeader(tc.status) + _, _ = w.Write([]byte(tc.receipt)) + })) + defer server.Close() + raw, _ := json.Marshal(ingressRelayConfig{InstallationID: id, AppID: "cli_test", ChatID: "oc_test", Endpoint: server.URL, Secret: secret, Disposition: tc.disposition}) + relay, err := NewIngressRelay(string(raw)) + if err != nil { + t.Fatal(err) + } + uid, _ := util.ParseUUID(id) + inst := engine.ResolvedInstallation{ID: uid, Active: true, Platform: Installation{AppID: "cli_test", AppSecretEncrypted: []byte("PRIVATE")}} + lm := InboundMessage{EventID: "evt", AppID: "cli_test", ChatID: "oc_test", ChatType: ChatTypeGroup, Body: "PRIVATE", Content: `{"text":"hello"}`} + b, _ := json.Marshal(lm) + m := channel.InboundMessage{Raw: b, Source: channel.Source{ChatType: channel.ChatTypeGroup}} + for i := 0; i < 2; i++ { + got, err := relay.Capture(context.Background(), inst, m) + if (err != nil) != tc.wantErr || got != tc.consume { + t.Fatalf("got %v %v", got, err) + } + } + if calls != 2 { + t.Fatal("duplicates must be durably handled by receiver") + } + m.Source.ChatType = channel.ChatTypeP2P + got, err := relay.Capture(context.Background(), inst, m) + if got || err != nil || calls != 2 { + t.Fatal("DM forwarded") + } + }) + } +} +func TestIngressRelayConfig(t *testing.T) { + if r, e := NewIngressRelay(""); r != nil || e != nil { + t.Fatal("disabled relay") + } + for _, s := range []string{`{}`, `{"unknown":true}`, `{} {}`} { + if _, e := NewIngressRelay(s); e == nil { + t.Fatal("invalid accepted") + } + } +} diff --git a/server/internal/integrations/lark/ws_frame_decoder.go b/server/internal/integrations/lark/ws_frame_decoder.go index 10148d3a23b..2fa36b34a49 100644 --- a/server/internal/integrations/lark/ws_frame_decoder.go +++ b/server/internal/integrations/lark/ws_frame_decoder.go @@ -92,6 +92,7 @@ func (d *LarkJSONFrameDecoder) Decode(payload []byte, inst Installation) (Inboun ChatType: normalizeChatType(evt.Message.ChatType), MessageID: evt.Message.MessageID, SenderOpenID: OpenID(evt.Sender.SenderID.OpenID), + SenderType: evt.Sender.SenderType, MessageType: evt.Message.MessageType, Content: evt.Message.Content, CreateTime: evt.Message.CreateTime,