Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions apps/docs/content/docs/environment-variables.fr.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
1 change: 1 addition & 0 deletions apps/docs/content/docs/environment-variables.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
1 change: 1 addition & 0 deletions apps/docs/content/docs/environment-variables.zh.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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 工具连接 |
Expand Down
36 changes: 36 additions & 0 deletions docs/runbooks/lark-ingress-relay.md
Original file line number Diff line number Diff line change
@@ -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.
7 changes: 7 additions & 0 deletions server/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -692,6 +698,7 @@ func main() {
FeatureFlags: flags,
HeartbeatScheduler: heartbeatScheduler,
LLMMaxRetries: llmMaxRetries,
LarkIngress: larkIngress,
LLMDisableThinking: llmDisableThinking,
})
var replicaQueries *db.Queries
Expand Down
10 changes: 8 additions & 2 deletions server/cmd/server/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
8 changes: 8 additions & 0 deletions server/internal/integrations/channel/engine/resolvers.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions server/internal/integrations/channel/engine/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 55 additions & 0 deletions server/internal/integrations/channel/engine/router_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
})
}
}
1 change: 1 addition & 0 deletions server/internal/integrations/lark/feishu_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
139 changes: 139 additions & 0 deletions server/internal/integrations/lark/ingress_relay.go
Original file line number Diff line number Diff line change
@@ -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") {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🛑 Logic Error: Line 51 validation condition is excessively complex (285+ characters on one line) and difficult to verify for correctness. This creates high risk of logic bugs in security-critical configuration validation. Complex boolean expressions like this are prone to operator precedence errors and missing edge cases that could allow invalid configurations to pass validation.

Suggested change
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") {
if err != nil {
return nil, errors.New("invalid ingress relay target, secret or disposition")
}
if u.Host == "" || u.User != nil || u.Fragment != "" || u.RawQuery != "" {
return nil, errors.New("invalid ingress relay target, secret or disposition")
}
isLocalHTTP := u.Scheme == "http" && (u.Hostname() == "localhost" || u.Hostname() == "127.0.0.1" || u.Hostname() == "::1")
if u.Scheme != "https" && !isLocalHTTP {
return nil, errors.New("invalid ingress relay target, secret or disposition")
}
if idErr != nil || !strings.HasPrefix(c.AppID, "cli_") || !strings.HasPrefix(c.ChatID, "oc_") {
return nil, errors.New("invalid ingress relay target, secret or disposition")
}
if len(c.Secret) < 32 || (c.Disposition != "consume" && c.Disposition != "observe") {
return nil, errors.New("invalid ingress relay target, secret or disposition")
}

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")
}
Comment on lines +102 to +105

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Missing error handling: If json.Marshal returns an error on line 102, the function returns the error but leaves body uninitialized. The subsequent check len(body) > 512<<10 on line 103 will evaluate against a nil/empty slice instead of detecting the marshaling failure. This could mask serialization errors and allow processing to continue with invalid data.

Suggested change
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")
}
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 {
return false, errors.New("ingress source exceeds limit")
}
if 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)
Comment on lines +106 to +109

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🛑 Security Vulnerability: The timestamp in HMAC signature verification has no expiration check, allowing replay attacks. An attacker who intercepts a valid signed request can replay it indefinitely since time.Now().Unix() generates the timestamp but never validates its freshness.1

Add timestamp validation to reject requests outside an acceptable time window (e.g., ±5 minutes). The receiver should verify the timestamp freshness before processing the event.

Footnotes

  1. CWE-294: Authentication Bypass by Capture-replay - https://cwe.mitre.org/data/definitions/294.html ↩

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
}
Loading
Loading