package workbench import ( "context" "encoding/json" "os" "testing" "aigateway.local/core/internal/platform/config" "aigateway.local/core/internal/platform/database" ) // TestInboxMaterializeAndBroadcast 验证站内消息核心链路:事件物化 → 未读计数 → // 重放幂等 → 已读回执 → 管理员广播。Redis 传 nil,走 DB 权威未读路径。 func TestInboxMaterializeAndBroadcast(t *testing.T) { databaseURL := os.Getenv("WORKBENCH_TEST_DATABASE_URL") if databaseURL == "" { t.Skip("WORKBENCH_TEST_DATABASE_URL is not set") } ctx := context.Background() pool, err := database.Open(ctx, config.Database{URL: databaseURL, MaxConns: 8, MinConns: 0}) if err != nil { t.Fatal(err) } defer pool.Close() adminID := "44444444-4444-4444-4444-444444444444" portalID := "55555555-5555-5555-5555-555555555555" cleanup := func() { // 消息按收件人删:materialize 落给 admin 收件人,广播落给 portal 收件人(payload 为空, // 不能只按 payload 匹配,否则广播消息残留导致重跑未读数累加)。 // 注意:同一 $1 同时比较 uuid 列与 jsonb text 提取,须显式 ::uuid / ::text, // 否则 PG 无法推断参数类型报 "text = uuid"(SQLSTATE 42883)。 if _, cErr := pool.Exec(ctx, `DELETE FROM gateway.inbox_messages WHERE recipient_user_id=$1::uuid OR recipient_user_id=$2::uuid OR payload->>'actor_id'=$1::text OR payload->>'portal_user_id'=$1::text`, adminID, portalID); cErr != nil { t.Logf("cleanup inbox DELETE failed: %v", cErr) } _, _ = pool.Exec(ctx, `DELETE FROM gateway.admin_accounts WHERE id=$1::uuid`, adminID) _, _ = pool.Exec(ctx, `DELETE FROM gateway.portal_users WHERE id=$1::uuid OR lower(account)='m8-inbox-portal'`, portalID) } cleanup() defer cleanup() _, err = pool.Exec(ctx, `INSERT INTO gateway.admin_accounts(id,username,password_hash,role,active) VALUES($1,'m8-inbox-admin','test','superadmin',true) ON CONFLICT(id) DO NOTHING`, adminID) if err != nil { t.Fatal(err) } _, err = pool.Exec(ctx, `INSERT INTO gateway.portal_users(id,account,password_hash,active) VALUES($1,'m8-inbox-portal','test',true) ON CONFLICT(id) DO NOTHING`, portalID) if err != nil { t.Fatal(err) } svc := NewInboxService(NewService(pool), nil, "") // 1) 事件物化:knowledge_document.ready → admin 收件人 eventID := "99999999-9999-4999-8999-999999999999" payload, _ := json.Marshal(map[string]any{"chunk_count": "8", "actor_id": adminID}) if err := svc.Materialize(ctx, eventID, "knowledge_document.ready", payload); err != nil { t.Fatalf("materialize failed: %v", err) } if count, err := svc.UnreadCount(ctx, "admin", adminID); err != nil || count != 1 { t.Fatalf("unread after materialize = %d, err=%v; want 1", count, err) } // 2) 同事件重放幂等:不新增行、不报错 if err := svc.Materialize(ctx, eventID, "knowledge_document.ready", payload); err != nil { t.Fatalf("replay failed: %v", err) } if count, _ := svc.UnreadCount(ctx, "admin", adminID); count != 1 { t.Fatalf("unread after replay = %d, want 1 (idempotent)", count) } // 3) 收件箱列出 + 已读回执 items, err := svc.List(ctx, "admin", adminID, 10) if err != nil || len(items) != 1 { t.Fatalf("list = %d items, err=%v; want 1", len(items), err) } if items[0].SenderType != "system" || items[0].Category != "system" || items[0].Title != "知识文档已入库" { t.Fatalf("unexpected message shape: %+v", items[0]) } changed, err := svc.MarkRead(ctx, items[0].ID, "admin", adminID) if err != nil || !changed { t.Fatalf("mark read changed=%v err=%v; want true", changed, err) } if count, _ := svc.UnreadCount(ctx, "admin", adminID); count != 0 { t.Fatalf("unread after mark-read = %d, want 0", count) } // 4) 管理员广播到全部 portal 用户(至少命中测试门户用户) sent, err := svc.Broadcast(ctx, InboxInput{RecipientKind: "portal", Category: "system", Title: "m8 升级公告", Body: "新增站内消息功能", Link: "/portal/inbox"}, nil, adminID) if err != nil { t.Fatalf("broadcast failed: %v", err) } if sent < 1 { t.Fatalf("broadcast sent = %d, want >=1", sent) } if count, _ := svc.UnreadCount(ctx, "portal", portalID); count != 1 { t.Fatalf("portal unread after broadcast = %d, want 1", count) } // 5) AdminList scope=broadcasts 能看到这条管理员广播。broadcasts 是全局视角 // (所有 admin→portal 广播),不能假设列表恰好 1 条,改为在其中找到本测试广播。 broadcasts, err := svc.AdminList(ctx, adminID, "broadcasts", 100) if err != nil { t.Fatalf("broadcasts list failed: %v", err) } found := false for _, b := range broadcasts { if b.SenderType == "admin" && b.Body == "新增站内消息功能" { found = true break } } if !found { t.Fatalf("test broadcast not found in broadcasts list (%d items)", len(broadcasts)) } // 6) MarkAllRead 清空门户未读 marked, err := svc.MarkAllRead(ctx, "portal", portalID) if err != nil || marked != 1 { t.Fatalf("mark-all-read = %d, err=%v; want 1", marked, err) } if count, _ := svc.UnreadCount(ctx, "portal", portalID); count != 0 { t.Fatalf("portal unread after mark-all = %d, want 0", count) } }