Repository navigation
feat(lark): add optional durable group ingress relay #16
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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. |
| 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") { | ||||||||||||||||||||||||
| 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
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Missing error handling: If
Suggested change
|
||||||||||||||||||||||||
| 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
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 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
|
||||||||||||||||||||||||
| 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 | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
There was a problem hiding this comment.
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.