mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-26 20:45:37 +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>
677 lines
25 KiB
Go
677 lines
25 KiB
Go
package lark
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"encoding/base64"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"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"
|
|
)
|
|
|
|
// RegistrationSessionStatus is the discriminated state a `begin`
|
|
// session lives in. The HTTP status endpoint serializes the underlying
|
|
// string verbatim so the frontend can pattern-match without parsing
|
|
// prose.
|
|
type RegistrationSessionStatus string
|
|
|
|
const (
|
|
// RegistrationStatusPending means the QR has been minted and the
|
|
// background goroutine is still polling Lark. The frontend keeps
|
|
// polling our status endpoint at the cadence we return in
|
|
// poll_interval_seconds.
|
|
RegistrationStatusPending RegistrationSessionStatus = "pending"
|
|
|
|
// RegistrationStatusSuccess means the device-flow returned
|
|
// credentials AND the lark_installation + installer-binding pair
|
|
// committed. `installation_id` is populated. The frontend closes
|
|
// the dialog and invalidates the installations cache.
|
|
RegistrationStatusSuccess RegistrationSessionStatus = "success"
|
|
|
|
// RegistrationStatusError means the session reached a terminal
|
|
// failure (expired, user-denied, Lark protocol error, follow-up
|
|
// bot-info / DB error). `error_reason` is set to a stable code so
|
|
// the frontend can render the right copy without parsing
|
|
// `error_message`.
|
|
RegistrationStatusError RegistrationSessionStatus = "error"
|
|
)
|
|
|
|
// Reason codes the service stores on a failed session. Stable strings
|
|
// so the frontend can switch on them without parsing prose.
|
|
const (
|
|
RegistrationReasonExpired = "expired"
|
|
RegistrationReasonAccessDenied = "access_denied"
|
|
RegistrationReasonProtocol = "lark_protocol_error"
|
|
RegistrationReasonBotInfoFailed = "bot_info_failed"
|
|
RegistrationReasonInstallationConflict = "installation_conflict"
|
|
RegistrationReasonInstallerBindFailed = "installer_bind_failed"
|
|
RegistrationReasonInternalError = "internal_error"
|
|
)
|
|
|
|
// RegistrationServiceConfig configures the service.
|
|
type RegistrationServiceConfig struct {
|
|
// SessionTTL caps how long a successful or errored session stays in
|
|
// the in-process cache before GC. Default 30 minutes — long enough
|
|
// for the frontend to fetch the final status after the dialog
|
|
// closes, short enough that abandoned sessions do not pin memory
|
|
// forever. Independent of the device-flow expiry (Lark's
|
|
// expire_in, ~10 min).
|
|
SessionTTL time.Duration
|
|
|
|
// Now is overridable for deterministic expiry-bound tests.
|
|
Now func() time.Time
|
|
|
|
// Logger is used for protocol-level warnings (Lark error codes,
|
|
// post-success bot info failures). Nil uses slog.Default().
|
|
Logger *slog.Logger
|
|
}
|
|
|
|
func (c RegistrationServiceConfig) withDefaults() RegistrationServiceConfig {
|
|
if c.SessionTTL == 0 {
|
|
c.SessionTTL = 30 * time.Minute
|
|
}
|
|
if c.Now == nil {
|
|
c.Now = time.Now
|
|
}
|
|
if c.Logger == nil {
|
|
c.Logger = slog.Default()
|
|
}
|
|
return c
|
|
}
|
|
|
|
// RegistrationService owns the device-flow install lifecycle. It is the
|
|
// one place that:
|
|
//
|
|
// 1. opens a new device-flow session against Lark (Begin),
|
|
// 2. tracks the session's polling state in-process,
|
|
// 3. runs the background polling goroutine,
|
|
// 4. on success, calls APIClient.GetBotInfo with the freshly minted
|
|
// credentials, then writes lark_installation + the installer's
|
|
// lark_user_binding in a single transaction.
|
|
//
|
|
// In-process session storage is intentional: device-flow sessions are
|
|
// short-lived (<10 min), the QR has no value outside the same browser
|
|
// session that initiated it, and persisting half-completed installs
|
|
// into Postgres would add a migration + GC sweep without delivering any
|
|
// product capability the user can re-use across server restarts.
|
|
type RegistrationService struct {
|
|
cfg RegistrationServiceConfig
|
|
client *RegistrationClient
|
|
api APIClient
|
|
queries *ChannelStore
|
|
tx TxStarter
|
|
installs *InstallationService
|
|
binder InstallerBinder
|
|
authQueries authQueriesAdapter
|
|
|
|
// bus is optional. When wired (SetEventBus), a successful install
|
|
// publishes lark_installation:created the moment the row commits, so
|
|
// every workspace client refreshes its connection badge without
|
|
// waiting for a browser to poll the status endpoint to success. Nil
|
|
// is valid — install still works, it just won't push the WS frame.
|
|
bus *events.Bus
|
|
|
|
mu sync.Mutex
|
|
sessions map[string]*registrationSession
|
|
}
|
|
|
|
// authQueriesAdapter is the minimal lookup surface the service needs
|
|
// before kicking off a session: agent ↔ workspace ownership validation.
|
|
// Kept as an interface so tests can drop in a stub instead of a real
|
|
// *db.Queries + Postgres fixture.
|
|
type authQueriesAdapter interface {
|
|
GetAgentInWorkspace(ctx context.Context, params db.GetAgentInWorkspaceParams) (db.Agent, error)
|
|
}
|
|
|
|
// NewRegistrationService wires the device-flow client, the APIClient
|
|
// (for the post-success GetBotInfo lookup), and the DB write path. Any
|
|
// required dependency missing surfaces as a constructor error so a
|
|
// silent half-init at startup cannot leave the install button
|
|
// returning 500s at runtime.
|
|
func NewRegistrationService(
|
|
cfg RegistrationServiceConfig,
|
|
client *RegistrationClient,
|
|
api APIClient,
|
|
queries *db.Queries,
|
|
tx TxStarter,
|
|
installs *InstallationService,
|
|
binder InstallerBinder,
|
|
) (*RegistrationService, error) {
|
|
if client == nil {
|
|
return nil, errors.New("lark registration: RegistrationClient is required")
|
|
}
|
|
if api == nil {
|
|
return nil, errors.New("lark registration: APIClient is required")
|
|
}
|
|
if queries == nil {
|
|
return nil, errors.New("lark registration: queries is required")
|
|
}
|
|
if tx == nil {
|
|
return nil, errors.New("lark registration: TxStarter is required")
|
|
}
|
|
if installs == nil {
|
|
return nil, errors.New("lark registration: InstallationService is required")
|
|
}
|
|
if binder == nil {
|
|
return nil, errors.New("lark registration: InstallerBinder is required")
|
|
}
|
|
return &RegistrationService{
|
|
cfg: cfg.withDefaults(),
|
|
client: client,
|
|
api: api,
|
|
queries: NewChannelStore(queries),
|
|
tx: tx,
|
|
installs: installs,
|
|
binder: binder,
|
|
authQueries: queries,
|
|
sessions: make(map[string]*registrationSession),
|
|
}, nil
|
|
}
|
|
|
|
// SetEventBus wires the optional event bus AFTER construction so the
|
|
// six positional constructor-validation cases stay untouched and the
|
|
// bus remains nil-safe. With it set, finishSuccess publishes
|
|
// lark_installation:created at the row-commit point — the authoritative
|
|
// moment of truth — instead of relying on the HTTP status-poll handler
|
|
// to emit it only when a browser happens to poll to success.
|
|
func (s *RegistrationService) SetEventBus(bus *events.Bus) {
|
|
s.bus = bus
|
|
}
|
|
|
|
// publishInstalled emits lark_installation:created on the optional bus.
|
|
// Mirrors the revoke path (RevokeLarkInstallation publishes
|
|
// lark_installation:revoked from its handler): both events broadcast to
|
|
// the whole workspace via the SubscribeAll fanout, and the frontend
|
|
// invalidates larkKeys.installations on the lark_installation prefix, so
|
|
// every mounted surface (agent Integrations tab, inspector, Settings)
|
|
// refreshes its connection badge with no page reload. Covers fresh
|
|
// installs and revoked→active re-installs alike — both ride the same
|
|
// UpsertLarkInstallation write. Nil-safe.
|
|
func (s *RegistrationService) publishInstalled(workspaceID, installationID pgtype.UUID) {
|
|
if s.bus == nil {
|
|
return
|
|
}
|
|
s.bus.Publish(events.Event{
|
|
Type: protocol.EventLarkInstallationCreated,
|
|
WorkspaceID: uuidString(workspaceID),
|
|
ActorType: "system",
|
|
Payload: map[string]any{"installation_id": uuidString(installationID)},
|
|
})
|
|
}
|
|
|
|
// registrationSession is the in-memory state for one in-flight install.
|
|
type registrationSession struct {
|
|
id string
|
|
workspaceID pgtype.UUID
|
|
agentID pgtype.UUID
|
|
initiatorID pgtype.UUID
|
|
|
|
deviceCode string
|
|
domain string
|
|
qrCodeURL string
|
|
interval time.Duration
|
|
expiresAt time.Time
|
|
// region is the cloud the install was started against. The polling
|
|
// loop reads it as the initial value of its `region` local; if the
|
|
// poll stream surfaces a tenant_brand mid-flow, the local flips to
|
|
// RegionLark, but the session field stays at what the user picked
|
|
// (it is informational — the authoritative cloud flows back through
|
|
// finishSuccess via the loop's local).
|
|
region Region
|
|
|
|
mu sync.Mutex
|
|
status RegistrationSessionStatus
|
|
installationID pgtype.UUID
|
|
errorReason string
|
|
errorMessage string
|
|
gcAfter time.Time
|
|
}
|
|
|
|
func (s *registrationSession) snapshot() RegistrationSessionState {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
return RegistrationSessionState{
|
|
ID: s.id,
|
|
Status: s.status,
|
|
InstallationID: s.installationID,
|
|
ErrorReason: s.errorReason,
|
|
ErrorMessage: s.errorMessage,
|
|
}
|
|
}
|
|
|
|
func (s *registrationSession) markSuccess(installationID pgtype.UUID, gcAfter time.Time) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
s.status = RegistrationStatusSuccess
|
|
s.installationID = installationID
|
|
s.gcAfter = gcAfter
|
|
}
|
|
|
|
func (s *registrationSession) markError(reason, msg string, gcAfter time.Time) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
// Idempotent: if a parallel goroutine already terminated the
|
|
// session (e.g. expiry fired between status reads), don't clobber
|
|
// the first reason — the user already saw it.
|
|
if s.status != RegistrationStatusPending {
|
|
return
|
|
}
|
|
s.status = RegistrationStatusError
|
|
s.errorReason = reason
|
|
s.errorMessage = msg
|
|
s.gcAfter = gcAfter
|
|
}
|
|
|
|
// RegistrationSessionState is the read-only snapshot the handler
|
|
// serializes to the frontend. Internal mutex is hidden by construction.
|
|
type RegistrationSessionState struct {
|
|
ID string
|
|
Status RegistrationSessionStatus
|
|
InstallationID pgtype.UUID
|
|
ErrorReason string
|
|
ErrorMessage string
|
|
}
|
|
|
|
// BeginInstallParams is the trusted input from the handler — the
|
|
// workspace, agent, and initiating user have already been authenticated
|
|
// and authorized at the router (admin role on the workspace; agent
|
|
// belongs to the workspace).
|
|
type BeginInstallParams struct {
|
|
WorkspaceID pgtype.UUID
|
|
AgentID pgtype.UUID
|
|
InitiatorID pgtype.UUID
|
|
// Region picks which cloud's accounts host the device-flow begins
|
|
// against — Feishu (mainland, accounts.feishu.cn) or Lark
|
|
// (international, accounts.larksuite.com). The user picks this
|
|
// explicitly in the UI ("Bind to Feishu" vs "Bind to Lark") so the
|
|
// QR rendered up front already targets the right cloud and Lark
|
|
// users do not have to hit a Feishu URL first and rely on the
|
|
// tenant-brand auto-switch. Empty / unknown values fall back to
|
|
// Feishu, matching RegionOrDefault, so existing callers without
|
|
// the new field keep working.
|
|
Region Region
|
|
}
|
|
|
|
// BeginInstallResult is the public payload the handler echoes to the
|
|
// frontend. The session_id is the opaque handle the frontend uses to
|
|
// poll status; we deliberately do NOT echo the device_code or the
|
|
// polling interval (which is internal scheduling state).
|
|
type BeginInstallResult struct {
|
|
SessionID string
|
|
QRCodeURL string
|
|
ExpiresInSeconds int
|
|
PollIntervalSeconds int
|
|
}
|
|
|
|
// BeginInstall opens a fresh device-flow session and kicks off the
|
|
// background polling goroutine. The returned payload feeds the QR-code
|
|
// dialog on the frontend; the polling goroutine runs until success,
|
|
// terminal failure, or device_code expiry.
|
|
//
|
|
// The session_id is the only opaque token returned to the browser —
|
|
// the device_code is server-side only (Lark would honor a poll from
|
|
// anywhere if the device_code leaked, so we never echo it).
|
|
func (s *RegistrationService) BeginInstall(ctx context.Context, p BeginInstallParams) (BeginInstallResult, error) {
|
|
if !p.WorkspaceID.Valid || !p.AgentID.Valid || !p.InitiatorID.Valid {
|
|
return BeginInstallResult{}, errors.New("lark registration: workspace, agent, and initiator are required")
|
|
}
|
|
// Agent ownership pre-check — without this, a workspace admin
|
|
// could open an install session against another workspace's agent
|
|
// by guessing the UUID, and the device_code minted against Lark
|
|
// would still produce credentials. The handler does the same
|
|
// check; doing it here too keeps the service self-defending.
|
|
//
|
|
// We keep the agent: its name pre-fills the bot name on Lark's
|
|
// PersonalAgent creation form (see botNamePreset) so the installed
|
|
// bot reads "<agent> - Multica" instead of "{用户姓名}的智能助手".
|
|
agent, err := s.authQueries.GetAgentInWorkspace(ctx, db.GetAgentInWorkspaceParams{
|
|
ID: p.AgentID,
|
|
WorkspaceID: p.WorkspaceID,
|
|
})
|
|
if err != nil {
|
|
return BeginInstallResult{}, fmt.Errorf("lark registration: agent not in workspace: %w", err)
|
|
}
|
|
|
|
// Normalize the requested region: empty / unknown → Feishu, the same
|
|
// back-compat invariant the storage layer uses (RegionOrDefault).
|
|
// This both protects the device-flow client from a bogus value
|
|
// from the handler AND means a pre-region caller (omitting the
|
|
// field) keeps getting the historical mainland-first behaviour.
|
|
region := RegionOrDefault(string(p.Region))
|
|
|
|
begin, err := s.client.Begin(ctx, botNamePreset(agent.Name), region)
|
|
if err != nil {
|
|
return BeginInstallResult{}, fmt.Errorf("lark registration: begin: %w", err)
|
|
}
|
|
|
|
now := s.cfg.Now()
|
|
sessionID, err := randomSessionID()
|
|
if err != nil {
|
|
return BeginInstallResult{}, fmt.Errorf("lark registration: mint session id: %w", err)
|
|
}
|
|
sess := ®istrationSession{
|
|
id: sessionID,
|
|
workspaceID: p.WorkspaceID,
|
|
agentID: p.AgentID,
|
|
initiatorID: p.InitiatorID,
|
|
deviceCode: begin.DeviceCode,
|
|
domain: begin.Domain,
|
|
qrCodeURL: begin.QRCodeURL,
|
|
interval: begin.Interval,
|
|
expiresAt: now.Add(begin.ExpiresIn),
|
|
region: region,
|
|
status: RegistrationStatusPending,
|
|
}
|
|
s.mu.Lock()
|
|
s.sessions[sessionID] = sess
|
|
s.mu.Unlock()
|
|
|
|
// The polling goroutine outlives the request context, so we cannot
|
|
// reuse ctx here. We size its own context to the device_code
|
|
// expiry — the worst case is a session that quietly times out
|
|
// after Lark's window closes, which we surface to the user as
|
|
// RegistrationReasonExpired on the next status read.
|
|
go s.runPolling(sess)
|
|
|
|
return BeginInstallResult{
|
|
SessionID: sessionID,
|
|
QRCodeURL: begin.QRCodeURL,
|
|
ExpiresInSeconds: int(begin.ExpiresIn / time.Second),
|
|
PollIntervalSeconds: int(begin.Interval / time.Second),
|
|
}, nil
|
|
}
|
|
|
|
// GetSession returns the current state of an in-flight or recently-
|
|
// finished session. The workspace UUID is required so a session
|
|
// initiated by one workspace cannot be polled by another (the session
|
|
// id is unguessable but defense-in-depth costs nothing here).
|
|
//
|
|
// ErrRegistrationSessionNotFound is returned for unknown / expired /
|
|
// GC'd sessions; the frontend treats it the same as an error reason
|
|
// of "session_lost" — prompt the user to restart the install.
|
|
func (s *RegistrationService) GetSession(workspaceID pgtype.UUID, sessionID string) (RegistrationSessionState, error) {
|
|
if strings.TrimSpace(sessionID) == "" {
|
|
return RegistrationSessionState{}, ErrRegistrationSessionNotFound
|
|
}
|
|
s.gcExpiredLocked()
|
|
s.mu.Lock()
|
|
sess, ok := s.sessions[sessionID]
|
|
s.mu.Unlock()
|
|
if !ok {
|
|
return RegistrationSessionState{}, ErrRegistrationSessionNotFound
|
|
}
|
|
if !uuidEqual(sess.workspaceID, workspaceID) {
|
|
// Treat as not found — leaking "exists but wrong workspace"
|
|
// would let an attacker enumerate session ids across workspaces.
|
|
return RegistrationSessionState{}, ErrRegistrationSessionNotFound
|
|
}
|
|
return sess.snapshot(), nil
|
|
}
|
|
|
|
// runPolling is the background loop. The pattern matches the upstream
|
|
// SDK: wait → poll → branch on result. On any terminal outcome we
|
|
// either record the installation+binding (success) or mark the session
|
|
// errored.
|
|
func (s *RegistrationService) runPolling(sess *registrationSession) {
|
|
// Bound the entire polling life by Lark's expiry — once that
|
|
// window closes, no further poll can succeed.
|
|
ctx, cancel := context.WithDeadline(context.Background(), sess.expiresAt)
|
|
defer cancel()
|
|
|
|
interval := sess.interval
|
|
if interval <= 0 {
|
|
interval = time.Duration(registrationDefaultPollSeconds) * time.Second
|
|
}
|
|
domain := sess.domain
|
|
deviceCode := sess.deviceCode
|
|
// region tracks which cloud this install belongs to. It starts at
|
|
// whatever the user picked at begin-time (Feishu by default; the
|
|
// frontend now exposes an explicit Lark CTA that begins on
|
|
// accounts.larksuite.com directly). The SwitchedDomain branch
|
|
// below is still honored as a safety net — if a user clicks the
|
|
// Feishu CTA but actually authorizes with a Lark-international
|
|
// account, the poll stream surfaces tenant_brand="lark" and we
|
|
// flip the local accordingly. So at finishSuccess time `region`
|
|
// is the authoritative per-install cloud, derived first from the
|
|
// user's UI choice and then from the protocol's role-based switch
|
|
// — never by string-matching accounts hostnames (so staging/mock
|
|
// domains classify correctly too).
|
|
region := sess.region
|
|
if region == "" {
|
|
region = RegionFeishu
|
|
}
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
s.cfg.Logger.Info("lark registration: session expired",
|
|
"session_id", sess.id,
|
|
"workspace_id", uuidString(sess.workspaceID))
|
|
sess.markError(RegistrationReasonExpired, "QR expired before authorization", s.gcDeadline())
|
|
return
|
|
case <-time.After(interval):
|
|
}
|
|
|
|
res, err := s.client.Poll(ctx, domain, deviceCode)
|
|
if err != nil {
|
|
var re *RegistrationError
|
|
if errors.As(err, &re) {
|
|
s.cfg.Logger.Warn("lark registration: protocol error",
|
|
"session_id", sess.id, "code", re.Code, "desc", re.Description)
|
|
sess.markError(RegistrationReasonProtocol, re.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
// Transient transport error (DNS, network) — log and try
|
|
// again on the next tick rather than killing the session,
|
|
// which lets a 30-second cross-region blip self-heal.
|
|
s.cfg.Logger.Warn("lark registration: transport error, will retry",
|
|
"session_id", sess.id, "err", err)
|
|
continue
|
|
}
|
|
|
|
switch {
|
|
case res.SwitchedDomain != "":
|
|
// Tenant-brand switch — re-aim immediately without
|
|
// honoring the interval, matching the upstream SDK's
|
|
// behavior. Lark emits the brand hint exactly once on the
|
|
// transition poll and the credential-bearing response
|
|
// lands on the next call to the new domain.
|
|
//
|
|
// Both directions are honored (feishu→lark and lark→feishu)
|
|
// so the split-CTA UI's "wrong entry" path recovers
|
|
// regardless of which CTA the user picked. The new region
|
|
// rides on the same PollResult so we never have to
|
|
// re-derive it from the host string here — staging / mock
|
|
// accounts hosts then classify correctly without
|
|
// hostname-prefix matching.
|
|
domain = res.SwitchedDomain
|
|
region = res.SwitchedRegion
|
|
s.cfg.Logger.Info("lark registration: switched cloud after tenant-brand mismatch",
|
|
"session_id", sess.id, "domain", domain, "region", string(region))
|
|
continue
|
|
case res.ClientID != "" && res.ClientSecret != "":
|
|
s.finishSuccess(ctx, sess, res, region)
|
|
return
|
|
case res.Err != nil:
|
|
reason := RegistrationReasonProtocol
|
|
if res.Err.Code == "access_denied" {
|
|
reason = RegistrationReasonAccessDenied
|
|
} else if res.Err.Code == "expired_token" {
|
|
reason = RegistrationReasonExpired
|
|
}
|
|
s.cfg.Logger.Info("lark registration: terminal error",
|
|
"session_id", sess.id, "code", res.Err.Code, "desc", res.Err.Description)
|
|
sess.markError(reason, res.Err.Error(), s.gcDeadline())
|
|
return
|
|
case res.Status == "slow_down":
|
|
// Honor Lark's back-off — bump by 5s, per RFC 8628 §3.5.
|
|
interval += 5 * time.Second
|
|
default:
|
|
// authorization_pending — keep the interval, loop.
|
|
}
|
|
}
|
|
}
|
|
|
|
// finishSuccess runs the post-poll finalization: bot info lookup +
|
|
// installation insert + installer binding, all in a single DB
|
|
// transaction.
|
|
func (s *RegistrationService) finishSuccess(ctx context.Context, sess *registrationSession, res *PollResult, region Region) {
|
|
// Carry the detected region onto the credentials so the GetBotInfo
|
|
// call below hits the right open-platform host: a Lark-international
|
|
// install must reach open.larksuite.com, not the Feishu default.
|
|
creds := InstallationCredentials{AppID: res.ClientID, AppSecret: res.ClientSecret, Region: region}
|
|
info, err := s.api.GetBotInfo(ctx, creds)
|
|
if err != nil {
|
|
s.cfg.Logger.Warn("lark registration: bot info failed",
|
|
"session_id", sess.id, "err", err)
|
|
sess.markError(RegistrationReasonBotInfoFailed, err.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
if info.OpenID == "" {
|
|
sess.markError(RegistrationReasonBotInfoFailed, "bot info missing open_id", s.gcDeadline())
|
|
return
|
|
}
|
|
|
|
// Encrypt the app_secret before the transaction so the seal cost
|
|
// doesn't sit inside the DB lock. The InstallationService's Upsert
|
|
// would do this for us, but we need the encrypted blob inside the
|
|
// transaction-scoped queries handle so the installer-bind commits
|
|
// alongside the installation insert — replicate the Seal here.
|
|
sealed, err := s.installs.box.Seal([]byte(res.ClientSecret))
|
|
if err != nil {
|
|
s.cfg.Logger.Error("lark registration: seal app_secret",
|
|
"session_id", sess.id, "err", err)
|
|
sess.markError(RegistrationReasonInternalError, err.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
|
|
tx, err := s.tx.Begin(ctx)
|
|
if err != nil {
|
|
s.cfg.Logger.Error("lark registration: begin tx",
|
|
"session_id", sess.id, "err", err)
|
|
sess.markError(RegistrationReasonInternalError, err.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
defer tx.Rollback(ctx)
|
|
qtx := s.queries.WithTx(tx)
|
|
|
|
inst, err := qtx.UpsertLarkInstallation(ctx, UpsertInstallationParams{
|
|
WorkspaceID: sess.workspaceID,
|
|
AgentID: sess.agentID,
|
|
AppID: res.ClientID,
|
|
AppSecretEncrypted: sealed,
|
|
BotOpenID: string(info.OpenID),
|
|
BotUnionID: textOrNull(info.UnionID),
|
|
InstallerUserID: sess.initiatorID,
|
|
Region: string(region),
|
|
})
|
|
if err != nil {
|
|
s.cfg.Logger.Warn("lark registration: upsert installation",
|
|
"session_id", sess.id, "err", err)
|
|
sess.markError(RegistrationReasonInstallationConflict, err.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
|
|
if err := s.binder.BindInstallerTx(ctx, qtx, InstallerBindParams{
|
|
WorkspaceID: sess.workspaceID,
|
|
InstallationID: inst.ID,
|
|
MulticaUserID: sess.initiatorID,
|
|
LarkOpenID: res.OpenID,
|
|
}); err != nil {
|
|
s.cfg.Logger.Warn("lark registration: bind installer",
|
|
"session_id", sess.id, "err", err)
|
|
sess.markError(RegistrationReasonInstallerBindFailed, err.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
|
|
if err := tx.Commit(ctx); err != nil {
|
|
s.cfg.Logger.Error("lark registration: commit",
|
|
"session_id", sess.id, "err", err)
|
|
sess.markError(RegistrationReasonInternalError, err.Error(), s.gcDeadline())
|
|
return
|
|
}
|
|
sess.markSuccess(inst.ID, s.gcDeadline())
|
|
// Publish at the commit point so the connection badge updates on every
|
|
// workspace client without a page refresh — not only on the tab that
|
|
// happens to poll the status endpoint to success.
|
|
s.publishInstalled(sess.workspaceID, inst.ID)
|
|
s.cfg.Logger.Info("lark registration: install complete",
|
|
"session_id", sess.id,
|
|
"workspace_id", uuidString(sess.workspaceID),
|
|
"agent_id", uuidString(sess.agentID),
|
|
"installation_id", uuidString(inst.ID))
|
|
}
|
|
|
|
func (s *RegistrationService) gcDeadline() time.Time {
|
|
return s.cfg.Now().Add(s.cfg.SessionTTL)
|
|
}
|
|
|
|
// gcExpiredLocked drops any session whose `gcAfter` is in the past.
|
|
// Pending sessions are NOT GC'd here — runPolling will set their
|
|
// gcAfter when it terminates, and an expired-by-deadline session
|
|
// closes itself.
|
|
func (s *RegistrationService) gcExpiredLocked() {
|
|
now := s.cfg.Now()
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
for id, sess := range s.sessions {
|
|
sess.mu.Lock()
|
|
drop := !sess.gcAfter.IsZero() && sess.gcAfter.Before(now)
|
|
sess.mu.Unlock()
|
|
if drop {
|
|
delete(s.sessions, id)
|
|
}
|
|
}
|
|
}
|
|
|
|
// ErrRegistrationSessionNotFound is what the service returns for
|
|
// unknown / GC'd sessions. The handler maps it to 404.
|
|
var ErrRegistrationSessionNotFound = errors.New("lark registration: session not found")
|
|
|
|
func randomSessionID() (string, error) {
|
|
buf := make([]byte, 24)
|
|
if _, err := rand.Read(buf); err != nil {
|
|
return "", err
|
|
}
|
|
return base64.RawURLEncoding.EncodeToString(buf), nil
|
|
}
|
|
|
|
func uuidEqual(a, b pgtype.UUID) bool {
|
|
if !a.Valid || !b.Valid {
|
|
return false
|
|
}
|
|
return a.Bytes == b.Bytes
|
|
}
|
|
|
|
// botNamePreset builds the display name we pre-fill on Lark's
|
|
// PersonalAgent creation form so the installed bot reads
|
|
// "<agent> - Multica" instead of Lark's auto-generated
|
|
// "{用户姓名}的智能助手". Lark treats this as a default the installer can
|
|
// still edit; we never get to lock the final name. A blank agent name
|
|
// (defensive — Agent.Name is NOT NULL in schema) degrades to plain
|
|
// "Multica" rather than a dangling " - Multica".
|
|
func botNamePreset(agentName string) string {
|
|
name := strings.TrimSpace(agentName)
|
|
if name == "" {
|
|
return "Multica"
|
|
}
|
|
return name + " - Multica"
|
|
}
|
|
|
|
// uuidString is the package-local UUID-to-string helper defined in
|
|
// hub.go; redeclared `func uuidString(u pgtype.UUID) string` removed
|
|
// to avoid the symbol collision.
|
|
//
|
|
// InstallationService.box is unexported but reachable from this file
|
|
// because both live in package `lark`; we read it directly in
|
|
// finishSuccess so the Seal happens outside the DB transaction (which
|
|
// would otherwise hold a row lock across the crypto call).
|