mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-26 12:35:35 +02:00
* 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>
545 lines
19 KiB
Go
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 ""
|
|
}
|