mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-26 04:25:46 +02:00
* feat(channel): add channel-agnostic engine Supervisor (MUL-3620) Stage-1 (MUL-3515) shipped the channel abstraction but nothing drove it. Add the generic engine that does: - channel.InboundHandler + Config.Handler: the single shared inbound entry the engine injects into every adapter (Hermes set_message_handler model). - channel.Channel.Connect now blocks for the connection lifetime (doc), so the supervisor can tie lease renewal to connection liveness. - new package channel/engine: Supervisor, generalized out of lark.Hub. It enumerates active installations across ALL channel types (no hard-coded feishu), fences each behind the WS lease CAS, builds the platform Channel via channel.Registry, drives Connect/Disconnect with backoff+jitter, and restarts on credential rotation. Knows nothing about any platform. channel.Channel is now driven by an engine; integrations/channel has an external consumer. Feishu adapter + boot cutover follow next. Tests: supervisor_test.go covers lease CAS, reclaim, reap-on-revoke, rotation restart + token fencing, backoff on build error, lease-loss teardown, bounded release, shutdown timeout. Race-clean. Co-authored-by: multica-agent <github@multica.ai> * feat(lark): drive Feishu through the channel engine; remove lark.Hub (MUL-3620) Refactor Feishu into the first channel.Channel and cut boot over to the channel-agnostic engine.Supervisor, removing the Feishu-only Hub. - feishuChannel implements channel.Channel: Connect runs the existing WS long-conn connector for one installation; Send posts a text reply via the Lark IM API; Capabilities declares Feishu's feature set. RegisterFeishu wires it to channel.TypeFeishu — adding a platform is now 'register a Factory', no engine edit. - FeishuRuntime extracts the former Hub.handleEvent / scheduleReply: runs the Dispatcher and drives the detached typing indicator + OutcomeReplier off the connector ACK path. main.go drains it on shutdown after the supervisor stops delivering events. - channelInstallationStore (engine.InstallationStore) enumerates active installations across ALL channel types via the new de-hardcoded query ListAllActiveChannelInstallations; the Supervisor routes each row to its registered Factory by channel_type. Generic per-row fingerprint replaces the feishu-specific one. - boot: engine.Supervisor replaces lark.Hub.Run; MULTICA_LARK_HUB_DISABLED keeps its name for runbook compatibility. - delete hub.go / hub_pgx.go / hub_test.go; relocate the connector contract (EventConnector/EventEmitter), uuidString, and the reply-path tests (-> feishu_runtime_test.go) so coverage is preserved. No channel_* schema change. Feishu behaviour unchanged; lark + channel + engine tests green under -race; go build/vet ./... clean. Remaining (follow-up): lift the Dispatcher pipeline into a channel- agnostic engine.Router over channel.InboundMessage + resolver interfaces, so the inbound core stops being Lark-shaped and adding a channel needs zero core edits (validated by Slack, MUL-3516). Co-authored-by: multica-agent <github@multica.ai> * feat(channel): add channel-agnostic engine.Router (inbound pipeline) (MUL-3620) Generalize lark.Dispatcher's inbound pipeline into engine.Router: the single shared channel.InboundHandler the Supervisor injects into every Channel. It routes by ChannelType to a registered ResolverSet and runs the same ordered pipeline for every platform (install route -> two-phase dedup -> group @bot filter -> identity+membership -> ensure session -> append+mark -> /issue -> debounced run), then drives the detached OutboundReplier + typing indicator. Platform specifics live behind resolver interfaces (InstallationResolver, IdentityResolver, Deduper, SessionBinder, Auditor, OutboundReplier, TypingNotifier) + shared services (IssueCreator/TaskEnqueuer/SessionReader). Adding a platform is 'register a ResolverSet', not 'edit the Router'. Outcome / DropReason values match the legacy lark ones 1:1. Additive: lark.Dispatcher untouched and still wired; the feishu ResolverSet, the cutover, and the old-path removal land next. channel.InboundMessage gains ForceFresh (the normalized /fresh affordance). Batcher moved into engine. router_test.go covers the pipeline invariants (routing, dedup finalize states, group filter, identity, membership, ensure/append, /issue, debounce, flush offline, force-fresh, drain) with generic fakes; race-clean. Co-authored-by: multica-agent <github@multica.ai> * feat(lark): cut Feishu over to engine.Router; remove lark.Dispatcher; core no longer Lark-shaped (MUL-3620) Wire the channel-agnostic engine.Router (added in the prior commit) as the shared inbound handler and refactor Feishu into a ResolverSet, completing the generic-engine cutover. The inbound core (engine.Router) now contains zero platform specifics. - Feishu ResolverSet (feishu_resolvers.go): InstallationResolver, IdentityResolver, Deduper, SessionBinder, Auditor, OutboundReplier, TypingNotifier — each backed by the existing ChannelStore / ChatSessionService / OutcomeReplier / typing indicator, translating at the channel.InboundMessage boundary (platform fields read from Raw). origin_type stays 'lark_chat'. - feishuChannel now produces a normalized channel.InboundMessage and hands it to the engine handler via channel.Config.Handler; the old Raw round-trip through lark.Dispatcher is gone. - Remove lark.Dispatcher, FeishuRuntime, and lark's pending_batcher (the engine owns the pipeline + batcher now); their behavioural coverage moved to engine.Router tests. Surviving native types (InboundMessage / Outcome / DispatchResult) relocated to feishu_types.go. elon review nits addressed: - The channel engine (Registry + Router + Supervisor) is now built UNCONDITIONALLY, outside the MULTICA_LARK_SECRET_KEY gate, so a non-Lark deployment runs it; Feishu registers its Factory + ResolverSet only when its key is present. - channel.Config.Raw is now genuinely the platform config JSONB (channel_installation.config): the feishu factory builds a credentials-only Installation from it, and the workspace/agent identity is resolved per message by the Router — no full-db-row marshaling. - feishuChannel gains direct unit tests: factory config decode, Send text + reply-target mapping, Capabilities, inbound normalization + Raw round-trip, msg-type + result mapping. No channel_* schema change. go build/vet ./... clean; channel + engine + lark green under -race. Feishu behaviour preserved (pipeline logic lifted verbatim, only generalized). Co-authored-by: multica-agent <github@multica.ai> * docs(channel): fix stale comments on the channel engine boot (MUL-3620) Address Elon's review nit: three comments still described the pre-cutover behavior. - handler.go: ChannelSupervisor is built UNCONDITIONALLY now, not nil when MULTICA_LARK_SECRET_KEY is unset. - main.go: same — the supervisor always exists; only MULTICA_LARK_HUB_DISABLED parks it. - router.go: with no platform registered the store still lists active rows; Registry.Build returns ErrUnknownType and the supervisor backs off (it does not 'find no installations'). Comment-only; no behavior change. Co-authored-by: multica-agent <github@multica.ai> --------- Co-authored-by: J <j@multica.ai> Co-authored-by: multica-agent <github@multica.ai>
100 lines
4.4 KiB
Go
100 lines
4.4 KiB
Go
package channel
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
)
|
|
|
|
// Type identifies an inbound channel platform — the discriminator the
|
|
// Registry keys on and the value persisted in the channel_type column of
|
|
// the generalized channel_* tables. Use the lower-case platform slug
|
|
// ("feishu", "slack", "wecom", …); keep it stable, it is durable data.
|
|
type Type string
|
|
|
|
const (
|
|
// TypeFeishu is the Feishu / Lark adapter — the only implementation
|
|
// in phase 1. It serves both the mainland Feishu cloud and the Lark
|
|
// international cloud; the cloud (region) is per-installation config,
|
|
// not a separate Type.
|
|
TypeFeishu Type = "feishu"
|
|
)
|
|
|
|
// Channel is the platform-agnostic contract every IM integration
|
|
// implements. An adapter keeps ALL platform specifics behind these five
|
|
// methods: the core supervisor calls Connect/Disconnect to manage the
|
|
// link, Send to deliver an outbound reply, and reads Capabilities to
|
|
// decide how to render; it never touches platform SDKs or wire formats.
|
|
//
|
|
// Inbound is intentionally NOT on this interface. A Channel pushes
|
|
// normalized InboundMessage values into the core router via the wiring
|
|
// established at construction (the adapter owns its receive loop); the
|
|
// core does not poll the Channel for messages.
|
|
type Channel interface {
|
|
// Type reports the platform discriminator. It MUST equal the Type
|
|
// the Channel was registered under, and is stable for the lifetime
|
|
// of the instance.
|
|
Type() Type
|
|
|
|
// Connect establishes the platform link (e.g. dials the outbound
|
|
// WebSocket long-conn, or starts the inbound HTTP listener) and then
|
|
// BLOCKS, running the receive loop, until the link ends. The
|
|
// connection mode is the implementation's choice and invisible to the
|
|
// core. It returns:
|
|
//
|
|
// - nil when ctx is cancelled (graceful shutdown / lease loss);
|
|
// - a non-nil error when the link drops and cannot be recovered
|
|
// locally — the supervisor treats this as "this attempt failed"
|
|
// and reconnects under exponential backoff.
|
|
//
|
|
// While Connect runs, the adapter delivers each inbound message by
|
|
// invoking the InboundHandler it captured at construction
|
|
// (Config.Handler). Send may be called concurrently from another
|
|
// goroutine for the lifetime of the connection. Implementations MUST
|
|
// tolerate repeated Connect calls on different contexts: the
|
|
// supervisor may Connect, return, and Connect again after backoff.
|
|
Connect(ctx context.Context) error
|
|
|
|
// Disconnect tears the platform link down and releases its
|
|
// resources. It is safe to call after a failed Connect and safe to
|
|
// call more than once; a Channel that is already disconnected
|
|
// returns nil.
|
|
Disconnect(ctx context.Context) error
|
|
|
|
// Send delivers a single outbound message and returns the platform's
|
|
// identifier for the delivered message. A non-nil error is reserved
|
|
// for real delivery failures (network, auth, rate limit) that the
|
|
// caller may retry.
|
|
Send(ctx context.Context, out OutboundMessage) (SendResult, error)
|
|
|
|
// Capabilities declares what this Channel supports. It is a pure
|
|
// declaration with no side effects and a stable result; callers read
|
|
// it to choose a rendering and degrade on their own (this package
|
|
// performs no degradation — see the Capability docs).
|
|
Capabilities() Capability
|
|
}
|
|
|
|
// Config is the normalized per-installation configuration a Factory
|
|
// consumes. Type is the platform discriminator; Raw is the platform's
|
|
// own credential/config blob (Feishu's app_id / encrypted app_secret /
|
|
// tenant_key / region, Slack's bot/app tokens, …), carried opaquely so
|
|
// the foundation never grows a per-platform field. It maps directly onto
|
|
// the channel_type column + JSONB config of a channel_installation row
|
|
// (MUL-3515 decision §3).
|
|
type Config struct {
|
|
Type Type
|
|
Raw json.RawMessage
|
|
|
|
// Handler is the shared inbound entry point the engine injects so the
|
|
// built Channel can deliver normalized InboundMessage values into the
|
|
// core (see InboundHandler). A Factory captures it and invokes it from
|
|
// the Channel's receive loop. It may be nil when a Channel is built
|
|
// purely for its outbound Send path (no inbound delivery needed).
|
|
Handler InboundHandler
|
|
}
|
|
|
|
// Factory builds a Channel from its per-installation Config. Each adapter
|
|
// registers exactly one Factory under its Type; the Registry calls it to
|
|
// instantiate a per-installation Channel. A Factory should validate Raw
|
|
// and return an error rather than a half-built Channel.
|
|
type Factory func(cfg Config) (Channel, error)
|