mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-24 19:41:14 +02:00
* feat(slack): add typing reaction on inbound message (MUL-3874) Mirror the Feishu typing indicator on Slack: react with 👀 on the user's message when it is ingested, then remove the reaction when the agent's run finishes (EventChatDone) or fails (EventTaskFailed). - New slack.TypingIndicatorManager: Add on ingest, Clear on terminal run events; state keyed by chat_session_id, bot token re-resolved from the DB on clear (never held in memory), all failures logged and swallowed (best-effort). - Wire via the channel-agnostic engine.TypingNotifier seam (slackTypingNotifier in the ResolverSet) — the Router already calls OnIngested off the ACK path. - Clear subscribes to the event bus directly so a failed run also drops the reaction (the outbound replier only handles EventChatDone). - Skip messages older than 2m so Socket Mode reconnect replays don't restamp. Requires the installed Slack app to hold the reactions:write scope. Co-authored-by: multica-agent <github@multica.ai> * fix(slack): clear typing reaction when no task runs; document reactions:write (MUL-3874) Addresses review feedback on the typing-indicator PR. 1. Stuck reaction on offline/archived agent. The debounced flush (flushChatRun) enqueues no task when the agent has no runtime or is archived (or on any enqueue/reload error), so no task lifecycle event is ever published and the bus-driven clear never fires — leaving the 👀 (and Feishu's Typing) reaction stuck on the user's message. Fix at the shared engine seam: add TypingNotifier.OnSettled(ctx, sessionID), which the Router calls from the flush on every no-task exit (before any offline/archived notice). Both the Slack and Feishu notifiers route it to manager.Clear, so the latent Feishu case is fixed too. Adds engine coverage (offline/archived clear, success does not) and a Slack OnSettled test. 2. Missing reactions:write scope in docs. reactions.add/remove silently fail without the scope, but the BYO app manifest/docs never listed it. Add reactions:write to the manifest + scope table and a reinstall note across all four locales (en/zh/ja/ko). Co-authored-by: multica-agent <github@multica.ai> --------- Co-authored-by: J <agent-j@multica.ai> Co-authored-by: multica-agent <github@multica.ai>
227 lines
9.3 KiB
Go
227 lines
9.3 KiB
Go
package engine
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
|
|
"github.com/multica-ai/multica/server/internal/integrations/channel"
|
|
"github.com/multica-ai/multica/server/internal/service"
|
|
db "github.com/multica-ai/multica/server/pkg/db/generated"
|
|
)
|
|
|
|
// This file defines the pluggable seams the Router runs the inbound pipeline
|
|
// through. Everything platform-specific lives behind these interfaces; a
|
|
// platform registers a ResolverSet and the channel-agnostic Router stays
|
|
// unchanged. The Feishu implementation is the first ResolverSet.
|
|
|
|
// Outcome categorizes what the Router decided to do with an inbound message.
|
|
// Values match the legacy lark outcomes 1:1 so behavior and dashboards carry
|
|
// over unchanged.
|
|
type Outcome string
|
|
|
|
const (
|
|
OutcomeDropped Outcome = "dropped"
|
|
OutcomeNeedsBinding Outcome = "needs_binding"
|
|
OutcomeIngested Outcome = "ingested"
|
|
OutcomeAgentOffline Outcome = "agent_offline"
|
|
OutcomeAgentArchived Outcome = "agent_archived"
|
|
)
|
|
|
|
// DropReason enumerates the drop-audit categories. Values match the legacy
|
|
// lark drop reasons 1:1.
|
|
type DropReason string
|
|
|
|
const (
|
|
DropReasonUnboundUser DropReason = "unbound_user"
|
|
DropReasonNonWorkspaceMember DropReason = "non_workspace_member"
|
|
DropReasonNotAddressedInGroup DropReason = "not_addressed_in_group"
|
|
DropReasonDuplicate DropReason = "duplicate"
|
|
DropReasonRevokedInstallation DropReason = "revoked_installation"
|
|
DropReasonInvalidEvent DropReason = "invalid_event"
|
|
)
|
|
|
|
// Result is the typed verdict the Router produces for one inbound message,
|
|
// consumed by the outbound side (OutboundReplier / typing). It mirrors the
|
|
// legacy lark.DispatchResult.
|
|
type Result struct {
|
|
Outcome Outcome
|
|
DropReason DropReason
|
|
InstallationID pgtype.UUID
|
|
ChatSessionID pgtype.UUID
|
|
// Sender is the platform-native sender id (e.g. Lark open_id), so the
|
|
// replier can target a binding prompt back to the sender.
|
|
Sender string
|
|
IssueID pgtype.UUID
|
|
IssueNumber int32
|
|
IssueIdentifier string
|
|
IssueTitle string
|
|
}
|
|
|
|
// ResolvedInstallation is the channel-agnostic installation context the Router
|
|
// needs after routing. Platform carries the adapter's own installation value
|
|
// opaquely so the set's other ports (binder, replier, typing) reuse it without
|
|
// a re-fetch; the Router never reads Platform.
|
|
type ResolvedInstallation struct {
|
|
ID pgtype.UUID
|
|
WorkspaceID pgtype.UUID
|
|
AgentID pgtype.UUID
|
|
InstallerUserID pgtype.UUID
|
|
Active bool
|
|
Platform any
|
|
}
|
|
|
|
// ResolvedIdentity is the sender mapped to a Multica user.
|
|
type ResolvedIdentity struct {
|
|
UserID pgtype.UUID
|
|
}
|
|
|
|
// EnsureSessionParams carries the inputs for SessionBinder.EnsureSession.
|
|
// Sender is the resolved session creator (the sole human for p2p, the
|
|
// installer for group chats — the Router decides which and passes it here).
|
|
type EnsureSessionParams struct {
|
|
Installation ResolvedInstallation
|
|
Sender pgtype.UUID
|
|
Message channel.InboundMessage
|
|
}
|
|
|
|
// AppendParams carries the inputs for SessionBinder.AppendMessage. ClaimToken
|
|
// is the dedup owner-fence token; the binder runs the dedup Mark INSIDE its
|
|
// chat_message+session tx so the durable write and the Mark commit atomically.
|
|
type AppendParams struct {
|
|
SessionID pgtype.UUID
|
|
Sender pgtype.UUID
|
|
InstallationID pgtype.UUID
|
|
Message channel.InboundMessage
|
|
ClaimToken pgtype.UUID
|
|
}
|
|
|
|
// AppendResult reports what AppendMessage decided.
|
|
type AppendResult struct {
|
|
// IssueCommand is non-nil when the message was an /issue command.
|
|
IssueCommand *IssueCommand
|
|
// DedupMarked is true when AppendMessage finalized the dedup claim in its
|
|
// own tx; the Router then skips the post-pipeline finalize.
|
|
DedupMarked bool
|
|
}
|
|
|
|
// IssueCommand is the parsed /issue command.
|
|
type IssueCommand struct {
|
|
Title string
|
|
Description string
|
|
}
|
|
|
|
// Sentinel errors the resolvers return so the Router can map them to the right
|
|
// product outcome instead of an infrastructure failure.
|
|
var (
|
|
// ErrInstallationNotFound: no installation matches the message's routing
|
|
// key → invalid_event drop.
|
|
ErrInstallationNotFound = errors.New("engine: installation not found")
|
|
// ErrSenderUnbound: the sender has no identity binding → needs_binding.
|
|
ErrSenderUnbound = errors.New("engine: sender unbound")
|
|
// ErrSenderNotMember: the sender is bound but not a workspace member →
|
|
// non_workspace_member drop.
|
|
ErrSenderNotMember = errors.New("engine: sender not a workspace member")
|
|
// ErrDuplicate: Claim found the message already processed / in flight →
|
|
// duplicate drop.
|
|
ErrDuplicate = errors.New("engine: duplicate message")
|
|
// ErrClaimLost: a concurrent reclaim rotated the dedup token mid-flight →
|
|
// treated as a duplicate.
|
|
ErrClaimLost = errors.New("engine: dedup claim lost")
|
|
)
|
|
|
|
// InstallationResolver routes an inbound message to its installation. The
|
|
// adapter reads whatever platform routing key it needs from the message
|
|
// (Source or Raw). Return ErrInstallationNotFound when none matches; return a
|
|
// ResolvedInstallation with Active=false when it exists but is revoked.
|
|
type InstallationResolver interface {
|
|
ResolveInstallation(ctx context.Context, msg channel.InboundMessage) (ResolvedInstallation, error)
|
|
}
|
|
|
|
// IdentityResolver maps the message sender to a Multica user within the
|
|
// installation, re-checking workspace membership. Return ErrSenderUnbound or
|
|
// ErrSenderNotMember for the product cases.
|
|
type IdentityResolver interface {
|
|
ResolveSender(ctx context.Context, inst ResolvedInstallation, msg channel.InboundMessage) (ResolvedIdentity, error)
|
|
}
|
|
|
|
// Deduper is the two-phase idempotency seam. Claim mints an owner-fence token
|
|
// (ErrDuplicate when already processed / in flight); Mark/Release are fenced on
|
|
// the token (a no-op on token mismatch is not an error).
|
|
type Deduper interface {
|
|
Claim(ctx context.Context, installationID pgtype.UUID, messageID string) (claimToken pgtype.UUID, err error)
|
|
Mark(ctx context.Context, installationID pgtype.UUID, messageID string, claimToken pgtype.UUID) error
|
|
Release(ctx context.Context, installationID pgtype.UUID, messageID string, claimToken pgtype.UUID) error
|
|
}
|
|
|
|
// SessionBinder ensures the chat_session and appends the message (with the
|
|
// in-tx dedup Mark). AppendMessage returns ErrClaimLost when the token was
|
|
// rotated mid-flight.
|
|
type SessionBinder interface {
|
|
EnsureSession(ctx context.Context, p EnsureSessionParams) (pgtype.UUID, error)
|
|
AppendMessage(ctx context.Context, p AppendParams) (AppendResult, error)
|
|
}
|
|
|
|
// Auditor records a dropped inbound event (no message body — drop-audit
|
|
// policy). instID may be the zero UUID for installation-less events.
|
|
type Auditor interface {
|
|
RecordDrop(ctx context.Context, instID pgtype.UUID, msg channel.InboundMessage, reason DropReason) error
|
|
}
|
|
|
|
// OutboundReplier delivers the verdict-driven reply (binding prompt, offline /
|
|
// archived notice, /issue confirmation). Optional; nil disables outbound
|
|
// replies. Driven off the ACK critical path by the Router.
|
|
type OutboundReplier interface {
|
|
Reply(ctx context.Context, inst ResolvedInstallation, msg channel.InboundMessage, res Result)
|
|
}
|
|
|
|
// TypingNotifier shows a "processing" indicator when a message is ingested and
|
|
// clears it once the message reaches a terminal outcome. Optional; nil disables
|
|
// it.
|
|
type TypingNotifier interface {
|
|
// OnIngested shows the indicator for a successfully ingested message.
|
|
OnIngested(ctx context.Context, inst ResolvedInstallation, msg channel.InboundMessage, sessionID pgtype.UUID)
|
|
// OnSettled clears the indicator for a session whose run trigger produced no
|
|
// task (agent offline / archived, or an enqueue failure). In that case no
|
|
// task lifecycle event is ever published, so the platform's own bus-driven
|
|
// clear (on chat-done / task-failed) would never fire and the indicator would
|
|
// stick. The Router calls this from the debounced flush. Idempotent: a
|
|
// session with no indicator is a no-op.
|
|
OnSettled(ctx context.Context, sessionID pgtype.UUID)
|
|
}
|
|
|
|
// ResolverSet is the per-platform bundle the Router runs the pipeline through.
|
|
// Installation/Identity/Dedup/Session/Audit are required; Replier/Typing are
|
|
// optional. OriginType is the issue.origin_type label written for /issue
|
|
// commands from this channel (Feishu: "lark_chat").
|
|
type ResolverSet struct {
|
|
Installation InstallationResolver
|
|
Identity IdentityResolver
|
|
Dedup Deduper
|
|
Session SessionBinder
|
|
Audit Auditor
|
|
Replier OutboundReplier
|
|
Typing TypingNotifier
|
|
OriginType string
|
|
}
|
|
|
|
// IssueCreator is the narrow subset of service.IssueService the Router needs
|
|
// for the /issue command. Shared across platforms.
|
|
type IssueCreator interface {
|
|
Create(ctx context.Context, p service.IssueCreateParams, opts service.IssueCreateOpts) (service.IssueCreateResult, error)
|
|
}
|
|
|
|
// TaskEnqueuer is the narrow subset of service.TaskService the Router needs to
|
|
// trigger a chat run. Shared across platforms.
|
|
type TaskEnqueuer interface {
|
|
EnqueueChatTask(ctx context.Context, session db.ChatSession, initiatorUserID pgtype.UUID, forceFreshSession bool) (db.AgentTaskQueue, error)
|
|
}
|
|
|
|
// SessionReader reads the rows the debounced flush + /issue identifier need.
|
|
// Shared across platforms; backed by *db.Queries (the channel-backed store).
|
|
type SessionReader interface {
|
|
GetChatSession(ctx context.Context, id pgtype.UUID) (db.ChatSession, error)
|
|
GetWorkspace(ctx context.Context, id pgtype.UUID) (db.Workspace, error)
|
|
}
|