Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 9 additions & 9 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ jobs:
- name: Setup Go
uses: actions/setup-go@v5
with:
go-version: "1.26.x"
go-version: "1.26.9"
check-latest: true
cache-dependency-path: server/go.sum

Expand Down Expand Up @@ -349,8 +349,8 @@ jobs:
uses: actions/setup-go@v5
with:
# Deliberately independent from the minimum patch in server/go.mod:
# CI follows the newest 1.26 patch available to setup-go.
go-version: "1.26.x"
# CI pins Go 1.26.9 for a consistent patched toolchain.
go-version: "1.26.9"
check-latest: true
cache: false

Expand Down Expand Up @@ -450,7 +450,7 @@ jobs:
id: go
uses: actions/setup-go@v5
with:
go-version: "1.26.x"
go-version: "1.26.9"
check-latest: true
cache: false

Expand Down Expand Up @@ -489,8 +489,8 @@ jobs:
uses: actions/setup-go@v5
with:
# Deliberately independent from the minimum patch in server/go.mod:
# CI follows the newest 1.26 patch available to setup-go.
go-version: "1.26.x"
# CI pins Go 1.26.9 for a consistent patched toolchain.
go-version: "1.26.9"
check-latest: true
cache-dependency-path: server/go.sum

Expand Down Expand Up @@ -528,8 +528,8 @@ jobs:
uses: actions/setup-go@v5
with:
# Deliberately independent from the minimum patch in server/go.mod:
# CI follows the newest 1.26 patch available to setup-go.
go-version: "1.26.x"
# CI pins Go 1.26.9 for a consistent patched toolchain.
go-version: "1.26.9"
check-latest: true
cache-dependency-path: server/go.sum

Expand Down Expand Up @@ -788,7 +788,7 @@ jobs:
- uses: actions/checkout@v6
- uses: actions/setup-go@v5
with:
go-version: "1.26.x"
go-version: "1.26.9"
cache-dependency-path: server/go.sum
- name: Verify Cursor background lifecycle and watchdog races
shell: bash
Expand Down
134 changes: 132 additions & 2 deletions server/cmd/server/listeners.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package main

import (
"context"
"encoding/json"
"fmt"
"log/slog"
Expand All @@ -12,6 +13,104 @@ import (
"github.com/multica-ai/multica/server/pkg/protocol"
)

type eventAudienceResolver func(context.Context, events.Event) ([]string, bool)

func eventObjectReference(e events.Event) (objectType, objectID string, scoped bool) {
if e.TaskID != "" {
return "task", e.TaskID, true
}
var payload map[string]any
encoded, err := json.Marshal(e.Payload)
if err == nil {
_ = json.Unmarshal(encoded, &payload)
}
if payload != nil {
for _, key := range []string{"issue_id", "issueId"} {
if id := findPayloadString(payload, key); id != "" {
return "issue", id, true
}
}
for _, key := range []string{"project_id", "projectId"} {
if id := findPayloadString(payload, key); id != "" {
return "project", id, true
}
}
for _, key := range []string{"squad_id", "squadId"} {
if id := findPayloadString(payload, key); id != "" {
return "squad", id, true
}
}
for _, ref := range []struct{ key, typ string }{{"issue", "issue"}, {"project", "project"}, {"squad", "squad"}} {
if id := findNestedObjectID(payload, ref.key); id != "" {
return ref.typ, id, true
}
}
if id, ok := payload["task_id"].(string); ok && id != "" {
return "task", id, true
}
}
t := strings.ToLower(e.Type)
for _, prefix := range []string{"issue", "task", "project", "squad", "comment", "attachment", "reaction", "issue_metadata", "issue_labels", "issue_properties", "issue_status", "wakeup"} {
if strings.HasPrefix(t, prefix) {
return "", "", true
}
}
return "", "", false
}

func findPayloadString(value any, key string) string {
switch current := value.(type) {
case map[string]any:
if id, ok := current[key].(string); ok && id != "" {
return id
}
for _, child := range current {
if id := findPayloadString(child, key); id != "" {
return id
}
}
case []any:
for _, child := range current {
if id := findPayloadString(child, key); id != "" {
return id
}
}
}
return ""
}

func findNestedObjectID(value any, key string) string {
switch current := value.(type) {
case map[string]any:
if object, ok := current[key].(map[string]any); ok {
if id, ok := object["id"].(string); ok && id != "" {
return id
}
}
for _, child := range current {
if id := findNestedObjectID(child, key); id != "" {
return id
}
}
case []any:
for _, child := range current {
if id := findNestedObjectID(child, key); id != "" {
return id
}
}
}
return ""
}

func audienceContains(audience []string, userID string) bool {
for _, candidate := range audience {
if candidate == userID {
return true
}
}
return false
}

// internalOnlyPayloadKeys lists payload keys that exist purely for in-process
// listeners and must never be serialized to a WebSocket client.
//
Expand Down Expand Up @@ -76,7 +175,23 @@ func projectOutbound(eventType string, payload any) any {
// for a Redis-backed relay or a feature-flagged dual-write implementation
// without touching any of the event listeners below. This is Phase 0 of the
// horizontal-scaling plan tracked in MUL-1138.
func registerListeners(bus *events.Bus, b realtime.Broadcaster) {
func registerListeners(bus *events.Bus, b realtime.Broadcaster, audienceResolvers ...eventAudienceResolver) {
var resolveAudience eventAudienceResolver
if len(audienceResolvers) > 0 {
resolveAudience = audienceResolvers[0]
}
eventVisibleTo := func(e events.Event, recipientID string) bool {
if e.RecipientIDs != nil {
return audienceContains(e.RecipientIDs, recipientID)
}
if resolveAudience != nil {
audience, scoped := resolveAudience(context.Background(), e)
if scoped {
return audienceContains(audience, recipientID)
}
}
return true
}
// Personal events should NOT be broadcast to the whole workspace.
personalEvents := map[string]bool{
protocol.EventInboxNew: true,
Expand All @@ -93,7 +208,7 @@ func registerListeners(bus *events.Bus, b realtime.Broadcaster) {

// Helper: marshal event and send to a specific user.
sendToRecipient := func(b realtime.Broadcaster, e events.Event, recipientID string) {
if recipientID == "" {
if recipientID == "" || !eventVisibleTo(e, recipientID) {
return
}
data, err := json.Marshal(map[string]any{"type": e.Type, "payload": projectOutbound(e.Type, e.Payload), "actor_id": e.ActorID, "actor_type": e.ActorType})
Expand Down Expand Up @@ -245,6 +360,21 @@ func registerListeners(bus *events.Bus, b realtime.Broadcaster) {
if personalEvents[e.Type] {
return
}
if e.RecipientIDs != nil {
for _, recipientID := range e.RecipientIDs {
sendToRecipient(b, e, recipientID)
}
return
}
if resolveAudience != nil {
if recipients, scoped := resolveAudience(context.Background(), e); scoped {
e.RecipientIDs = recipients
for _, recipientID := range recipients {
sendToRecipient(b, e, recipientID)
}
return
}
}

msg := map[string]any{
"type": e.Type,
Expand Down
33 changes: 31 additions & 2 deletions server/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import (
"github.com/multica-ai/multica/server/internal/dbstartup"
"github.com/multica-ai/multica/server/internal/events"
"github.com/multica-ai/multica/server/internal/handler"
"github.com/multica-ai/multica/server/internal/util"
"github.com/multica-ai/multica/server/internal/integrations/wecom"
"github.com/multica-ai/multica/server/internal/logger"
"github.com/multica-ai/multica/server/internal/maintenance"
Expand Down Expand Up @@ -602,13 +603,41 @@ func main() {
channelLeaseRedis = newNamedRedisClient(opts, "channel-lease")
}
}
registerListeners(bus, broadcaster)

analyticsClient := analytics.NewFromEnv()
defer analyticsClient.Close()

queries := db.New(pool)
hub.SetAuthorizer(newScopeAuthorizer(queries))
visibilityGate := &handler.Handler{Queries: queries, DB: pool}
resolveEventAudience := func(ctx context.Context, e events.Event) ([]string, bool) {
objectType, objectID, scoped := eventObjectReference(e)
if !scoped {
return nil, false
}
if objectType == "task" {
taskID, err := util.ParseUUID(objectID)
if err != nil {
return nil, true
}
task, err := queries.GetAgentTask(ctx, taskID)
if err != nil {
return nil, true
}
if !task.IssueID.Valid {
// Standalone tasks do not yet have a business-object audience
// resolver. Keep them scoped and fail closed instead of falling
// back to workspace-wide delivery.
return nil, true
}
objectType, objectID = "issue", util.UUIDToString(task.IssueID)
}
audience, err := visibilityGate.BusinessObjectRecipients(ctx, e.WorkspaceID, objectType, objectID)
if err != nil {
return nil, true
}
return audience, true
}
registerListeners(bus, broadcaster, resolveEventAudience)
// Order matters: subscriber listeners must register BEFORE notification listeners.
// The notification listener queries the subscriber table to determine recipients,
// so subscribers must be written first within the same synchronous event dispatch.
Expand Down
Loading
Loading