Files
multica/server/internal/integrations/lark/outbound.go
Bohan Jiang ce28d0aa0e feat(integrations): add platform-agnostic channel foundation (MUL-3515) (#4412)
* feat(integrations): add platform-agnostic channel foundation

Introduce server/internal/integrations/channel — the contract every
inbound IM integration implements, so the core never learns a platform's
event JSON. Four pieces:

- Channel interface (Type/Connect/Disconnect/Send/Capabilities) + Factory
  + Config (channel_type + opaque JSON blob, maps to channel_installation).
- Normalized InboundMessage/OutboundMessage envelopes + Source/MediaRef/
  ReplyCtx/MsgType/ChatType. Envelope holds only cross-platform-true
  fields; platform specifics live in Raw, read only by the adapter.
- Capability bitmask: declaration only, no degrade logic in core.
- Registry: Type->Factory map, last-writer-wins, concurrency-safe.

Pure package (no DB/network/platform deps). Foundation for MUL-3515; the
lark cutover + lark_*->channel_* generalization land in follow-up PRs.

MUL-3515

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

* feat(channel): generalize lark_* tables into channel_* (DB layer)

Migration 123 creates channel_installation / channel_user_binding /
channel_chat_session_binding / channel_inbound_message_dedup /
channel_inbound_audit / channel_outbound_card_message /
channel_binding_token. Each carries a channel_type discriminator and a
JSONB config for platform-specific identifiers/credentials; cross-platform
columns stay flat. Existing Feishu rows are backfilled (channel_type=
'feishu', app_secret_encrypted via base64). NO foreign keys / cascades
(MUL-3515 §4) — integrity moves to the app layer in the cutover.

queries/channel.sql ports the lark query surface to channel_*, JSONB-aware,
plus DeleteChannelUserBindingsByWorkspaceMember /
DeleteChannelChatSessionBindingBySession for the app-layer cleanup that
replaces the removed cascades.

lark_* tables/queries are left in place here and removed once the Go
cutover lands, so this commit ships green on its own.

Verified: sqlc generate, go build ./..., full migrate chain (1..123) on
Postgres 17, and a real-data backfill spot-check (base64 round-trip,
NULL-strip, functional unique index on (channel_type, app_id)).

MUL-3515

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

* fix(channel): name app_id query param + multi-IM install key + null-safe binding merge

Addresses review on MUL-3515 (PR #4412):

- GetChannelInstallationByAppID: explicitly name params and cast app_id to
  ::text so sqlc emits AppID string. A bare $2 next to `config ->> 'app_id'`
  was mis-attributed to the JSONB config column, generating Config []byte.

- channel_installation uniqueness -> (workspace_id, agent_id, channel_type),
  with the UpsertChannelInstallation conflict key matched. Lets one agent
  hold one installation per IM (feishu + slack + ...) instead of a later
  install clobbering an earlier one. Behaviorally identical in the current
  feishu-only world; "one agent, at most one IM overall" stays an app-layer
  rule per MUL-3515 §4, not a DB constraint.

- CreateChannelUserBinding merges jsonb_strip_nulls(EXCLUDED.config) so a
  re-bind carrying {"union_id": null} no longer erases an already-captured
  union_id, restoring the old COALESCE(EXCLUDED.union_id, ...) semantics.

Regenerated with sqlc v1.31.1. Verified on PG17: re-install replaces in
place, feishu+slack coexist, null re-bind keeps union_id, real union_id wins.

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

* feat(lark): channel-backed Feishu store + fix base64 backfill wrapping

Cutover step 1 of switching the lark Go code from lark_* onto the channel_*
tables (MUL-3515). Introduces the JSONB config boundary the rest of the
cutover sits on, and fixes a latent backfill bug surfaced while building it.

- migration 123: strip newlines from the app_secret_encrypted base64 backfill.
  PostgreSQL encode(...,'base64') MIME-wraps at 76 chars, and a secretbox-
  sealed ~72-byte secret exceeds that. Go's encoding/json decodes a JSON
  string into []byte with base64.StdEncoding, which rejects embedded newlines,
  so without the strip every migrated installation would fail to decrypt its
  app secret once reads move to channel_installation.config.

- store.go: flat domain types (Installation / UserBinding / ChatSessionBinding)
  with field parity to the retired db.Lark* rows, plus the feishu config codec.
  Row->domain mappers decode the JSONB config; the secret decoder is
  whitespace-tolerant so legacy MIME-wrapped data still round-trips, while the
  encoder emits unwrapped base64. Binding config encodes an absent union_id as
  "{}" so the upsert's jsonb_strip_nulls merge never clobbers a stored union_id.

- store_test.go: 72-byte secret round-trip, MIME-wrapped tolerance, optional
  null-strip, and flat-column preservation. Verified on PG17.

Field parity keeps the upcoming ~190 db.LarkInstallation call sites a
mechanical rename. No call sites switched yet; behavior unchanged.

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

* feat(lark): route inbound integration onto channel_* + explicit membership checks

Cutover step 2 (MUL-3515): switch the Feishu Go code from the lark_* queries to
channel_* via a ChannelStore adapter, and replace the removed member foreign key
with explicit application-layer membership checks. No user-visible behavior change.

- channel_store.go: ChannelStore embeds *db.Queries and SHADOWS the ~24 lark
  query methods with channel_*-backed equivalents, keeping the db.Lark*
  signatures so the dispatcher/hub/services and their ~20k lines of tests stay
  untouched; the feishu JSONB config is (de)coded by store.go. Adds
  IsWorkspaceMember and a tx-aware WithTx. Only production wiring swaps
  *db.Queries for *ChannelStore.

- Membership re-check (§4 removed the lark_user_binding -> member FK, so a
  binding row no longer proves current membership):
  * the dispatcher inbound identity step verifies membership after the binding
    lookup; a former member's stale binding is dropped as non_workspace_member
    + audited and never reaches chat_session (§4.3 safety property).
  * RedeemAndBind and BindInstallerTx replace the now-dead FK (23503) branch
    with an explicit IsWorkspaceMember gate, preserving the existing
    ErrBindingNotWorkspaceMember outcome without burning the token.

- router wires the ChannelStore into the patcher, typing indicator, dispatcher,
  hub, and the union_id/region backfills; constructor-based services wrap
  *db.Queries internally so their signatures and nil-check tests are unchanged.

Verified: go build ./... ; go vet ; gofmt ; go test -race ./internal/integrations/...
(full lark suite green unchanged + new membership drop/error tests). Adapter
field mappings (secret base64, union_id RMW, chat-id/open-id remaps, dedup,
token, card) checked end-to-end against a PG17 channel_* schema.

lark_* tables and queries remain (unused at runtime) until the S3 cleanup-hooks
and S4 drop-tables/rename commits.

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

* fix(channel): renumber generalization migration 123 -> 124

main merged 123_issue_stage after this branch forked, so the branch's 123_channel_generalization now collides on the migration number. The runner keys schema_migrations by full version string and would still apply both, but a duplicate number is a merge hazard and convention violation, so move the channel migration to the next free slot (124).

issue_stage (ALTER issue ADD COLUMN stage) and the channel generalization touch disjoint tables; verified on PG17 that 123_issue_stage applies cleanly on a DB already carrying 124_channel_generalization, so the two are order-independent. sqlc regenerated (v1.31.1): only the migration-number comment changed.

MUL-3515

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

* feat(channel): prune channel bindings on member removal + chat session delete

MUL-3515 §4 dropped every channel_* foreign key, so the old ON DELETE CASCADE that cleared a user's channel_user_binding when they left a workspace, and a chat's channel_chat_session_binding when its chat_session was deleted, no longer fires. Re-establish that integrity in the application layer, inside the existing transactions: revokeAndRemoveMember -> DeleteChannelUserBindingsByWorkspaceMember, DeleteChatSession -> DeleteChannelChatSessionBindingBySession.

Adds real-DB tests for both paths, including a scoping check that a remaining member's binding survives the prune. Verified on PG17: both new tests plus the existing revocation tests and the full handler package pass.

MUL-3515

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

* fix(channel): scope Lark/Feishu store reads to channel_type='feishu'

The S2 cutover routed the Feishu integration onto channel_*, but the Lark-facing ChannelStore wrappers read installation / chat-session-binding / outbound-card rows across ALL channel_type values. Once a second IM exists, that would let the Lark hub supervise a non-Feishu installation, the Lark install list show it, /lark/installations/{id} revoke another channel's row, and the outbound patcher / typing indicator act on a non-Feishu chat binding or card.

Add a channel_type predicate to the six read/list channel queries and pass channelTypeFeishu from every wrapper: GetChannelInstallation, GetChannelInstallationInWorkspace, ListChannelInstallationsByWorkspace, ListActiveChannelInstallations, GetChannelChatSessionBindingBySession, GetChannelOutboundCardByTask.

The S3 cleanup deletes (DeleteChannelUserBindingsByWorkspaceMember / DeleteChannelChatSessionBindingBySession) stay all-channel on purpose: a member leaving or a chat_session being deleted should clear every IM's binding. Adds a real-DB test that seeds a Slack installation/binding/card next to the Feishu ones and asserts the Lark wrappers never return them.

MUL-3515

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

* refactor(channel): replace db.Lark* translation layer with lark domain types

S2 introduced ChannelStore as a translation layer that read/wrote channel_* but kept the retired db.Lark* struct/param shapes so the dispatcher/hub/services and their ~20k lines of tests did not have to change. This collapses that layer: the store now takes and returns the package's flat domain types (Installation, UserBinding, ChatSessionBinding, InboundMessageDedup, BindingTokenRow, OutboundCardMessage) and the *Params types in params.go, with channel-neutral field names (ChannelUserID / ChannelChatID / ...). All call sites, fakes, and tests move to the domain types.

No behavior change: only channel_* is read/written (as before); db.Lark* is now unused, and the lark_* tables + queries/lark.sql are removed in the next commit. Verified on PG17: go build / vet / gofmt clean, go test -race ./internal/integrations/... green (the ~20k-line fake suite), and the lark + handler suites pass.

MUL-3515

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

* refactor(channel): drop lark_* tables and queries (remove old path)

The Go cutover (previous commit) moved the lark package entirely onto channel_* and the domain types, leaving the lark_* tables, queries/lark.sql, and the generated db.Lark* models unused. Remove them per the design (§5: replace, do not keep both): migration 125 drops the seven lark_* tables (data already lives in channel_* since migration 124), and queries/lark.sql is deleted + sqlc regenerated, removing the db.Lark* models and lark query methods.

The 125 down recreates the authoritative pre-drop schema (bot_union_id, region, per-installation dedup PK, thread-reply columns). Verified on PG17: fresh migrate up ends with lark_* gone + channel_* present; isolated 125 down/up round-trips correctly; go build / vet / gofmt clean; go test -race ./internal/integrations/... and the handler suite pass.

MUL-3515

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

* fix(migrations): remove trailing blank line at EOF of 125 down migration

git diff --check flagged a blank line at EOF of 125_drop_lark_tables.down.sql (a pg_dump-generation artifact). Whitespace only; the recreate SQL is unchanged.

MUL-3515

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

* refactor(channel): defer lark_* table drop to a follow-up migration

Preflight deploy review: dropping lark_* in the same release that cuts over (old migration 125) is not rollback/rolling-safe — the v0.3.27 release still reads lark_*, so a rolling deploy or a post-deploy code rollback would hit "relation does not exist". Remove the drop and keep the old tables for one release (standard expand/contract): migration 124 already backfilled lark_* -> channel_*, the new code reads/writes only channel_*, and the physical drop moves to a separate cleanup migration once this ships and is observed.

The lark_* tables remain in the schema, so sqlc regenerates the (now unused) db.Lark* models; queries/lark.sql stays deleted (the new code uses channel_*). No code path reads lark_* — only the destructive drop is deferred, keeping the design's no-compat-layer / no-dual-write rule while being deploy-safe.

MUL-3515

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

* fix(channel): skip orphaned installations in hub-boot active scan

Preflight deploy review: channel_installation dropped the workspace/agent FK (MUL-3515 §4), so unlike lark_installation it does not cascade away when its workspace is deleted or its agent is hard-deleted (e.g. runtime teardown). The hub-boot query then keeps opening a WebSocket for a bot whose owner is gone.

JOIN ListActiveChannelInstallations to live workspace + agent so an orphaned installation is never connected, uniformly for every deletion path. The JOIN matches the old ON DELETE CASCADE semantics (row existence, not agent archival), so an archived-but-present agent's installation is still listed; the orphaned row's encrypted secret is thereby never decrypted/used.

Tests: a real-DB handler test asserts a deleted-workspace/agent installation and a non-Feishu one are both excluded; the lark scope test's active-list assertion moved there since the JOIN now needs real workspace/agent fixtures. (Physically deleting dormant orphaned channel rows on workspace/agent deletion is a separate app-layer-cleanup follow-up.)

MUL-3515

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

* docs(channel): document non-rolling cutover constraint for the lark->channel migration

Elon deploy review: keeping the lark_* tables (deferred drop) stops old v0.3.27 code from crashing, but is not full expand/contract. Migration 124 is a one-time backfill; afterwards new code runs on channel_* (lease + dedup on channel_*) while pre-cutover code runs on lark_* (lease + dedup on lark_*). If both run concurrently during a rolling deploy, each side claims the same Feishu bot's WS lease on its own table and double-processes inbound events.

This release therefore requires a NON-ROLLING cutover (stop the old hub before applying migration 124 + starting new code; rollback is not lossless once new code writes channel_*). Documented where deployers/reviewers see it: migration 124 header gains a ROLLOUT note; the channel_store.go header is corrected (lark_* tables are retained one release for rollback safety, not "gone"; the store still never touches them). Comment-only — no schema/codegen/behavior change.

MUL-3515

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

* feat(lark): add MULTICA_LARK_HUB_DISABLED switch for the channel cutover

The lark_*->channel_* cutover needs a way to make the Feishu bot briefly unavailable WITHOUT taking down the whole multica-api process — the Lark hub is a goroutine inside it, not a separate Deployment. MULTICA_LARK_HUB_DISABLED=true parks the hub at startup: the API serves HTTP normally but never claims a WS lease or opens a Feishu connection.

Rollout (see migration 124 ROLLOUT note): ship the new release with the flag SET so new pods run API-only while old pods (hub on lark_*) drain during the rolling deploy — the two hubs never overlap. After the old pods are gone and migration 124 has run, flip the flag off; the new hub comes up on channel_*. The old backend does NOT need this switch — its hub stops when k8s terminates the old pods, not via a flag. Nil-ing LarkHub reuses the existing not-configured path so both the startup start and the shutdown join skip it.

MUL-3515

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

* docs(channel): point migration 124 ROLLOUT note at the hub-disable switch

Refine the rollout note to use MULTICA_LARK_HUB_DISABLED for a bot-only cutover (new pods serve API with the hub parked while old pods drain; flip the switch off after the migration), instead of the earlier whole-API recreate. Comment-only.

MUL-3515

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

* docs(channel): fix migration 124 rollout order and document self-host cutover

The previous ROLLOUT note shipped the new (channel_*) build before
running migration 124, so the channel_*-backed HTTP paths (installation
list/install/revoke, chat-session delete, member revoke) would 500 in
the window between new-pod boot and the deferred migration. Restate the
runbook around two explicit invariants — channel_* must exist before the
new build serves those paths, and the old/new hubs must never overlap —
and order the steps so channel_* is created first (park old hub -> snapshot
-> deploy parked new build -> unpark). Document that default self-host
(entrypoint migrate + single-replica Recreate) satisfies both invariants
automatically and needs no manual steps; only prd / multi-replica rolling
self-host needs the switch procedure. Clarify in main.go that the
hub-park switch is generation-agnostic (parks whichever hub the build
carries), which is what enables the preparatory release.

Refs MUL-3515

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-24 12:46:20 +08:00

545 lines
19 KiB
Go

package lark
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype"
"github.com/multica-ai/multica/server/internal/events"
db "github.com/multica-ai/multica/server/pkg/db/generated"
"github.com/multica-ai/multica/server/pkg/protocol"
)
// CardStatus mirrors lark_outbound_card_message.status. Kept as a typed
// alias so callers can't pass arbitrary strings into the status column.
type CardStatus string
const (
CardStatusPending CardStatus = "pending"
CardStatusStreaming CardStatus = "streaming"
CardStatusFinal CardStatus = "final"
CardStatusError CardStatus = "error"
)
// CardKind enumerates the small set of card variants the patcher
// renders. The Renderer is plug-replaceable so the on-wire card
// template can evolve without touching the patcher's transport / DB
// logic.
type CardKind string
const (
CardKindThinking CardKind = "thinking"
CardKindRunning CardKind = "running"
CardKindFinal CardKind = "final"
CardKindError CardKind = "error"
)
// CardRender is the rendered card body the Renderer produces. The
// patcher serializes the JSON before handing it to APIClient.
type CardRender struct {
JSON string
}
// RenderInput is the (typed) snapshot the Renderer sees when building
// or patching a card. Fields are populated as they become available
// during a task lifecycle — IssueNumber is set for `/issue` flows,
// Content is set for completed chat tasks, ErrorMessage for failed.
type RenderInput struct {
Kind CardKind
AgentName string
IssueNumber int32
IssueID pgtype.UUID
TaskID pgtype.UUID
Content string
ErrorMessage string
}
// Renderer turns a typed RenderInput into the actual Lark card JSON.
// Centralizing this lets us swap card templates (or A/B them) without
// touching event subscription or persistence code.
type Renderer interface {
Render(in RenderInput) (CardRender, error)
}
// defaultRenderer produces minimal text-only cards that work against
// Lark's generic interactive-card schema. The exact JSON layout will
// be refined when the real product card design lands; this default
// keeps the wiring real (the JSON deserializes against Lark's schema)
// without committing the product to a particular template.
type defaultRenderer struct{}
// NewDefaultRenderer returns the production-default Renderer. Override
// via PatcherConfig.Renderer when a custom template is needed.
func NewDefaultRenderer() Renderer { return &defaultRenderer{} }
func (defaultRenderer) Render(in RenderInput) (CardRender, error) {
header := "Multica"
if in.AgentName != "" {
header = in.AgentName
}
var body string
switch in.Kind {
case CardKindThinking:
body = "Thinking…"
case CardKindRunning:
body = "Working on it…"
case CardKindFinal:
body = in.Content
if body == "" {
body = "Done."
}
case CardKindError:
body = "Run failed."
if in.ErrorMessage != "" {
body = "Run failed: " + in.ErrorMessage
}
default:
return CardRender{}, fmt.Errorf("unknown card kind %q", in.Kind)
}
// update_multi MUST be true on every render: Lark refuses to apply
// PatchInteractiveCard to a card whose config does not declare it
// a "shared, updatable" card. Since this renderer drives the
// thinking → streaming → final/error lifecycle (the card is sent
// once and patched multiple times), an absent update_multi causes
// every patch after the first send to silently no-op on the
// Lark side while the local outbound status row still flips to
// streaming/final. Keep this on every kind — including thinking
// and error — because that initial JSON IS the body Lark stores
// and consults for subsequent patches.
doc := map[string]any{
"config": map[string]any{
"wide_screen_mode": true,
"update_multi": true,
},
"header": map[string]any{
"template": "blue",
"title": map[string]any{"tag": "plain_text", "content": header},
},
"elements": []any{
map[string]any{
"tag": "div",
"text": map[string]any{
"tag": "plain_text",
"content": body,
},
},
},
}
raw, err := json.Marshal(doc)
if err != nil {
return CardRender{}, err
}
return CardRender{JSON: string(raw)}, nil
}
// PatcherQueries is the narrow subset of *db.Queries the Patcher
// needs. Declared as an interface so the patcher is unit-testable
// without a real Postgres connection.
type PatcherQueries interface {
GetAgentTask(ctx context.Context, id pgtype.UUID) (db.AgentTaskQueue, error)
GetChatSession(ctx context.Context, id pgtype.UUID) (db.ChatSession, error)
GetAgent(ctx context.Context, id pgtype.UUID) (db.Agent, error)
GetLarkInstallation(ctx context.Context, id pgtype.UUID) (Installation, error)
GetLarkChatSessionBindingBySession(ctx context.Context, chatSessionID pgtype.UUID) (ChatSessionBinding, error)
GetLarkOutboundCardByTask(ctx context.Context, taskID pgtype.UUID) (OutboundCardMessage, error)
CreateLarkOutboundCardMessage(ctx context.Context, arg CreateOutboundCardMessageParams) (OutboundCardMessage, error)
UpdateLarkOutboundCardStatus(ctx context.Context, arg UpdateOutboundCardStatusParams) error
}
// CredentialsResolver decrypts an installation's app_secret for the
// transport layer. *InstallationService satisfies it directly; tests
// substitute a fake.
type CredentialsResolver interface {
DecryptAppSecret(inst Installation) (string, error)
}
// PatcherConfig tunes the outbound Patcher. Defaults via withDefaults;
// tests typically override Renderer / Now / Logger.
type PatcherConfig struct {
// Renderer drives the error card template used on the EventTaskFailed
// path. The success path (EventChatDone) bypasses the renderer
// entirely — it sends the raw assistant reply as a plain text IM
// message — so this only matters for the failure branch.
Renderer Renderer
Now func() time.Time
Logger *slog.Logger
}
func (c PatcherConfig) withDefaults() PatcherConfig {
if c.Renderer == nil {
c.Renderer = NewDefaultRenderer()
}
if c.Now == nil {
c.Now = time.Now
}
if c.Logger == nil {
c.Logger = slog.Default()
}
return c
}
// Patcher reacts to task-lifecycle events on the event bus and forwards
// chat replies to Lark as plain text IM messages. It is the outbound
// side of §4.5 — but the original "thinking → streaming → final card"
// lifecycle was reduced to a single plain-text reply on EventChatDone
// after Bohan reported the card chrome made replies feel like system
// notifications. The error path is the one survivor of card rendering:
// failed runs surface as a short error card on EventTaskFailed because
// the visual distinction from a normal reply is genuinely useful.
//
// Scope:
//
// - Only tasks whose chat_session has a lark_chat_session_binding
// produce outbound. Tasks born from the web UI or autopilot pass
// through unchanged.
//
// - Each EventChatDone yields one Lark text message; there is no
// streaming, no throttling, no DB row to track card-state.
//
// - Multi-replica safety is inherited from the inbound WS lease: at
// most one replica holds the installation lease at a time, the
// event bus is per-process, so exactly one Patcher reacts per run.
type Patcher struct {
queries PatcherQueries
credentials CredentialsResolver
client APIClient
typingIndicator *TypingIndicatorManager
cfg PatcherConfig
}
// NewPatcher constructs a Patcher bound to its dependencies. The
// patcher does not subscribe to the bus until Register is called.
func NewPatcher(queries PatcherQueries, credentials CredentialsResolver, client APIClient, cfg PatcherConfig) *Patcher {
cfg = cfg.withDefaults()
return &Patcher{
queries: queries,
credentials: credentials,
client: client,
cfg: cfg,
}
}
// SetTypingIndicatorManager wires the typing-indicator manager into the
// patcher so that replies clear the "processing" reaction before they
// are sent. Call once at boot after both the patcher and manager are
// constructed. Nil disables the clear step.
func (p *Patcher) SetTypingIndicatorManager(m *TypingIndicatorManager) {
p.typingIndicator = m
}
// Register subscribes the patcher to the task-lifecycle events it
// cares about on the supplied bus. Idempotent only if you call it
// against a fresh bus; call sites should invoke it exactly once
// during server boot (after the bus + patcher are constructed and
// before HTTP traffic starts).
//
// Subscriptions are deliberately minimal:
//
// - EventChatDone — the agent finished replying. The Patcher sends
// the reply as a plain text IM message (Lark's `msg_type=text`),
// not as an interactive card. The earlier card-based design (with
// thinking → running → final patches) made every reply look like
// a system notification nested in card chrome; flipping to plain
// text makes free-form chat feel native.
//
// - EventTaskFailed — the run failed; surface a short error card
// so the failure is visually distinct from a successful reply.
//
// We deliberately do NOT subscribe to EventTaskQueued / EventTaskRunning
// (no thinking-card lifecycle anymore — adds noise without value) or to
// EventTaskCompleted (chat tasks always emit EventChatDone first, which
// is what we care about; non-chat tasks have no Lark binding anyway and
// would early-return). Leaving EventTaskCompleted unsubscribed also
// avoids the prior "Done." overwrite regression where the no-content
// EventTaskCompleted payload would wipe the real reply.
func (p *Patcher) Register(bus *events.Bus) {
bus.Subscribe(protocol.EventTaskFailed, p.handleEvent)
bus.Subscribe(protocol.EventChatDone, p.handleEvent)
}
func (p *Patcher) handleEvent(e events.Event) {
// Use a fresh background ctx with a tight timeout: bus delivery is
// synchronous so a stuck Lark HTTP call would otherwise wedge the
// whole publish call site.
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := p.processEvent(ctx, e); err != nil {
p.cfg.Logger.Warn("lark patcher: event handling failed",
"event_type", e.Type,
"task_id", e.TaskID,
"chat_session_id", e.ChatSessionID,
"error", err,
)
}
}
func (p *Patcher) processEvent(ctx context.Context, e events.Event) error {
taskID, chatSessionID, ok := taskAndSessionFromEvent(e)
if !ok {
return nil
}
if !chatSessionID.Valid {
// Issue / autopilot tasks have no chat_session.
return nil
}
binding, err := p.queries.GetLarkChatSessionBindingBySession(ctx, chatSessionID)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
// Web-only chat session — not a Lark target.
return nil
}
return fmt.Errorf("lookup chat session binding: %w", err)
}
inst, err := p.queries.GetLarkInstallation(ctx, binding.InstallationID)
if err != nil {
return fmt.Errorf("load installation: %w", err)
}
if InstallationStatus(inst.Status) != InstallationActive {
// Revoked between trigger and event; nothing to patch.
return nil
}
creds, err := p.installationCredentials(inst)
if err != nil {
return err
}
agent, agentErr := p.queries.GetAgent(ctx, inst.AgentID)
agentName := ""
if agentErr == nil {
agentName = agent.Name
}
// Clear the "processing" reaction before the reply is visible so the
// user sees a clean transition. Best-effort: a failure here is logged
// but does not block the actual reply.
if p.typingIndicator != nil {
p.typingIndicator.Clear(ctx, chatSessionID)
}
switch e.Type {
case protocol.EventChatDone:
return p.sendChatReply(ctx, creds, binding, e.Payload)
case protocol.EventTaskFailed:
return p.fail(ctx, creds, binding, taskID, agentName, e.Payload)
}
return nil
}
// sendChatReply turns ChatDonePayload.Content into a Lark message.
// The wire shape is chosen per-reply based on whether the body
// contains any markdown syntax:
//
// - Plain prose (no markdown) → `msg_type=text`. A one-line "Hi!"
// reply should feel like a normal IM message, not a notification
// card with chrome around it.
//
// - Anything with markdown (headings, lists, code blocks, tables,
// bold/italic, links) → schema-2.0 interactive card with a
// `tag: "markdown"` body element so Lark's client renders the
// formatting instead of leaving raw `**bold**` characters in
// the transcript. The card is visually subtler than the legacy
// binding-prompt template — just a single markdown block, no
// header / icon / CTA buttons.
//
// Empty content is silently dropped: we'd rather show nothing than
// "Done." (the prior card fallback that confused Bohan in the live
// dev env). In practice an empty Content means the daemon completed
// the task without producing visible output, which only happens for
// edge cases like a chat task that just acknowledged a system event;
// not emitting a message there is the right product call.
func (p *Patcher) sendChatReply(ctx context.Context, creds InstallationCredentials, binding ChatSessionBinding, payload any) error {
content := chatDoneContent(payload)
if content == "" {
return nil
}
target := threadReplyTarget(binding)
if containsMarkdown(content) {
return sendWithThreadFallback(p.cfg.Logger, "send markdown card", target, func(t ReplyTarget) error {
_, err := p.client.SendMarkdownCard(ctx, SendMarkdownCardParams{
InstallationID: creds,
ChatID: ChatID(binding.ChannelChatID),
Markdown: content,
ReplyTarget: t,
})
return err
})
}
return sendWithThreadFallback(p.cfg.Logger, "send text message", target, func(t ReplyTarget) error {
_, err := p.client.SendTextMessage(ctx, SendTextParams{
InstallationID: creds,
ChatID: ChatID(binding.ChannelChatID),
Text: content,
ReplyTarget: t,
})
return err
})
}
// threadReplyTarget derives the outbound reply target from the chat
// binding's most-recent inbound trigger. We thread the reply ONLY when
// that trigger was itself inside a Lark topic (last_lark_thread_id
// present): normal group / p2p chats keep the unchanged chat-level send
// path, and only an @-mention that happened inside a thread gets a
// threaded reply (replying to last_lark_message_id with reply_in_thread).
// The zero ReplyTarget means "send at the chat level".
func threadReplyTarget(binding ChatSessionBinding) ReplyTarget {
if binding.LastThreadID.Valid && binding.LastThreadID.String != "" &&
binding.LastMessageID.Valid && binding.LastMessageID.String != "" {
return ReplyTarget{MessageID: binding.LastMessageID.String, InThread: true}
}
return ReplyTarget{}
}
// sendWithThreadFallback runs send with the thread reply target and,
// ONLY when the threaded attempt fails with a Lark error that means the
// topic reply legitimately cannot land (trigger message recalled, topic
// gone, topics disabled, aggregated message — see
// threadReplyUnsupportedCodes), retries once at the chat level so the
// reply is not silently lost. Any other failure — transport error,
// 5xx, timeout, rate limit, or an ambiguous "the server may have
// received it" error — is logged and returned as a failure rather than
// retried: a blind chat-level retry could duplicate the reply or leak a
// thread-only reply into the main group chat. When target is already
// chat-level there is nothing to fall back to and the error is returned.
//
// It is a package-level function (rather than a Patcher method) so the
// event-driven Patcher and the immediate OutcomeReplier share one
// classified fallback path.
func sendWithThreadFallback(log *slog.Logger, op string, target ReplyTarget, send func(ReplyTarget) error) error {
err := send(target)
if err == nil {
return nil
}
if target.IsSet() && isThreadReplyUnsupported(err) {
log.Warn("lark: thread reply unsupported for target, retrying at chat level",
"op", op, "reply_message_id", target.MessageID, "error", err)
if fallbackErr := send(ReplyTarget{}); fallbackErr != nil {
return fmt.Errorf("%s (chat-level fallback after thread-unsupported reply: %v): %w", op, err, fallbackErr)
}
return nil
}
if target.IsSet() {
log.Warn("lark: thread reply failed; not falling back (non-classified error)",
"op", op, "reply_message_id", target.MessageID, "error", err)
}
return fmt.Errorf("%s: %w", op, err)
}
func (p *Patcher) installationCredentials(inst Installation) (InstallationCredentials, error) {
if p.credentials == nil {
return InstallationCredentials{}, errors.New("lark patcher: credentials resolver missing")
}
secret, err := p.credentials.DecryptAppSecret(inst)
if err != nil {
return InstallationCredentials{}, fmt.Errorf("decrypt app_secret: %w", err)
}
creds := InstallationCredentials{
AppID: inst.AppID,
AppSecret: secret,
Region: RegionOrDefault(inst.Region),
}
if inst.TenantKey.Valid {
creds.TenantKey = inst.TenantKey.String
}
return creds, nil
}
// fail surfaces a short error card on task failure. Unlike the
// success path (plain text via sendChatReply), failures stay as cards
// because the user benefits from the visual distinction — a red /
// header-styled card is much harder to miss than a regular bubble,
// and these are rare enough that the card chrome isn't noisy.
//
// One-shot send (no patching, no DB row): if the task fails a second
// time we'd just send a second card, which is fine — failure is
// usually a single terminal event.
func (p *Patcher) fail(ctx context.Context, creds InstallationCredentials, binding ChatSessionBinding, taskID pgtype.UUID, agentName string, payload any) error {
render, err := p.cfg.Renderer.Render(RenderInput{
Kind: CardKindError,
AgentName: agentName,
TaskID: taskID,
ErrorMessage: errorMessageFromPayload(payload),
})
if err != nil {
return fmt.Errorf("render error card: %w", err)
}
return sendWithThreadFallback(p.cfg.Logger, "send error card", threadReplyTarget(binding), func(t ReplyTarget) error {
_, err := p.client.SendInteractiveCard(ctx, SendCardParams{
InstallationID: creds,
ChatID: ChatID(binding.ChannelChatID),
CardJSON: render.JSON,
ReplyTarget: t,
})
return err
})
}
// taskAndSessionFromEvent parses the typed-ish payload broadcastTaskEvent
// publishes — a map[string]any with `task_id` (always) and
// `chat_session_id` (chat tasks only). EventChatDone carries a
// ChatDonePayload struct instead.
func taskAndSessionFromEvent(e events.Event) (taskID, chatSessionID pgtype.UUID, ok bool) {
if e.TaskID != "" {
if err := taskID.Scan(e.TaskID); err != nil {
taskID = pgtype.UUID{}
}
}
if e.ChatSessionID != "" {
if err := chatSessionID.Scan(e.ChatSessionID); err != nil {
chatSessionID = pgtype.UUID{}
}
}
switch p := e.Payload.(type) {
case map[string]any:
if !taskID.Valid {
if s, _ := p["task_id"].(string); s != "" {
_ = taskID.Scan(s)
}
}
if !chatSessionID.Valid {
if s, _ := p["chat_session_id"].(string); s != "" {
_ = chatSessionID.Scan(s)
}
}
case protocol.ChatDonePayload:
if !taskID.Valid {
_ = taskID.Scan(p.TaskID)
}
if !chatSessionID.Valid {
_ = chatSessionID.Scan(p.ChatSessionID)
}
}
return taskID, chatSessionID, taskID.Valid
}
func chatDoneContent(payload any) string {
switch p := payload.(type) {
case protocol.ChatDonePayload:
return p.Content
case map[string]any:
if s, ok := p["content"].(string); ok {
return s
}
}
return ""
}
func errorMessageFromPayload(payload any) string {
if m, ok := payload.(map[string]any); ok {
if s, ok := m["error"].(string); ok {
return s
}
if s, ok := m["error_message"].(string); ok {
return s
}
}
return ""
}