diff --git a/apps/backend/cmd/worker/main.go b/apps/backend/cmd/worker/main.go index d1eb070..9f37f48 100644 --- a/apps/backend/cmd/worker/main.go +++ b/apps/backend/cmd/worker/main.go @@ -408,12 +408,17 @@ func runRadarSweep(ctx context.Context, jobs *jobUC.Service, radar *radarUC.Serv } res, err := radar.RunSweep(ctx, j.OwnerUID, watchID, j.ID) if err != nil { + // RunSweep persists a user-facing reason on fetch failures. Return that + // instead of leaking the crawler/provider transport response into Jobs UI. + if res != nil && res.FetchFailed && strings.TrimSpace(res.FailedReason) != "" { + return errors.New(res.FailedReason) + } return err } - summary := fmt.Sprintf("雷達巡檢完成 · 新建 %d · 判定 %d · 截斷 %d", res.Created, res.Judged, res.Truncated) if res.FetchFailed { return fmt.Errorf("%s", res.FailedReason) } + summary := fmt.Sprintf("雷達巡檢完成 · 新建 %d · 再次命中 %d · 判定 %d · 截斷 %d", res.Created, res.Rematched, res.Judged, res.Truncated) if _, err := jobs.MarkRunningProgress(ctx, j.ID, 90, summary); err != nil { return err } diff --git a/apps/backend/internal/logic/radar/list_opportunities_logic.go b/apps/backend/internal/logic/radar/list_opportunities_logic.go index 19bbc29..d48e973 100644 --- a/apps/backend/internal/logic/radar/list_opportunities_logic.go +++ b/apps/backend/internal/logic/radar/list_opportunities_logic.go @@ -31,7 +31,7 @@ func (l *ListOpportunitiesLogic) ListOpportunities(req *types.ListOpportunitiesR postedFrom, postedTo := req.From, req.To now := domain.NowNano() if req.TimeScope == "today" { - postedFrom, postedTo = domain.UTCDayBounds(now) + postedFrom, postedTo = domain.LocalDayBounds(now, domain.DisplayLocation()) } else if req.TimeScope == "7d" { postedFrom, postedTo = now-7*int64(24*time.Hour), 0 } diff --git a/apps/backend/internal/module/job/usecase/service.go b/apps/backend/internal/module/job/usecase/service.go index 5ca0167..3c74377 100644 --- a/apps/backend/internal/module/job/usecase/service.go +++ b/apps/backend/internal/module/job/usecase/service.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "fmt" + "strconv" "strings" "sync" "time" @@ -227,6 +228,61 @@ func (s *Service) findRadarSweepForRef(ctx context.Context, ownerUID int64, ref return nil, nil } +// ScheduleManualRadarSweep starts an extra patrol now. +// Daily jobs are keyed watch+UTC day and must stay idempotent; a finished +// daily job must not swallow 「立即巡邏」. If this watch already has a +// queued/running sweep, that in-flight job is returned instead of stacking. +func (s *Service) ScheduleManualRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (*domain.Job, error) { + watchID = strings.TrimSpace(watchID) + if ownerUID <= 0 || watchID == "" { + return nil, domain.ErrForbidden + } + if runAt <= 0 { + runAt = domain.NowNano() + } + list, err := s.Repo.ListByOwner(ctx, ownerUID) + if err != nil { + return nil, err + } + prefix := watchID + ":" + for _, j := range list { + if j == nil || j.TemplateType != domain.TemplateRadarSweep { + continue + } + if !strings.HasPrefix(j.RefID, prefix) { + continue + } + if !domain.IsTerminal(j.Status) { + return j, nil + } + } + day := time.Unix(0, runAt).UTC().Format("2006-01-02") + ref := watchID + ":manual:" + strconv.FormatInt(runAt, 10) + body, err := json.Marshal(RadarSweepPayload{WatchID: watchID, Day: day}) + if err != nil { + return nil, err + } + now := domain.NowNano() + j := &domain.Job{ + ID: uuid.NewString(), + OwnerUID: ownerUID, + TemplateType: domain.TemplateRadarSweep, + Status: domain.StatusQueued, + RefID: ref, + Payload: string(body), + RunAfter: runAt, + ProgressSummary: "立即巡邏已排程 · 等待 worker", + ProgressPercent: 0, + CreatedAt: now, + UpdatedAt: now, + } + if err := s.Repo.Insert(ctx, j); err != nil { + return nil, err + } + s.notify(ctx, j) + return j, nil +} + func (s *Service) List(ctx context.Context, ownerUID int64) ([]*domain.Job, error) { return s.Repo.ListByOwner(ctx, ownerUID) } diff --git a/apps/backend/internal/module/job/usecase/service_test.go b/apps/backend/internal/module/job/usecase/service_test.go index 2561806..f30d5fa 100644 --- a/apps/backend/internal/module/job/usecase/service_test.go +++ b/apps/backend/internal/module/job/usecase/service_test.go @@ -374,6 +374,32 @@ func TestJobLease_HeartbeatRenewsWhileRunning(t *testing.T) { }, 200*time.Millisecond, 5*time.Millisecond) } +func TestScheduleManualRadarSweep_ReusesInFlightAndCreatesAfterSuccess(t *testing.T) { + ctx := context.Background() + svc := usecase.New(jobRepo.NewMemory()) + first, err := svc.ScheduleManualRadarSweep(ctx, 42, "watch-1", 1) + require.NoError(t, err) + require.Contains(t, first.RefID, ":manual:") + again, err := svc.ScheduleManualRadarSweep(ctx, 42, "watch-1", 2) + require.NoError(t, err) + require.Equal(t, first.ID, again.ID) + + claimed, err := svc.ClaimNext(ctx, "worker-manual") + require.NoError(t, err) + require.Equal(t, first.ID, claimed.ID) + _, err = svc.SucceedJob(ctx, first.ID, "done") + require.NoError(t, err) + second, err := svc.ScheduleManualRadarSweep(ctx, 42, "watch-1", 3) + require.NoError(t, err) + require.NotEqual(t, first.ID, second.ID) + require.Contains(t, second.RefID, ":manual:") + + daily, err := svc.ScheduleRadarSweep(ctx, 42, "watch-1", 1) + require.NoError(t, err) + require.NotEqual(t, second.ID, daily.ID) + require.NotContains(t, daily.RefID, ":manual:") +} + func TestJobLease_ReclaimsLegacyRunningDocumentWithoutLease(t *testing.T) { repo := jobRepo.NewMemory() now := domain.NowNano() diff --git a/apps/backend/internal/module/radar/domain/opportunity.go b/apps/backend/internal/module/radar/domain/opportunity.go index f3a54f2..4b9ec72 100644 --- a/apps/backend/internal/module/radar/domain/opportunity.go +++ b/apps/backend/internal/module/radar/domain/opportunity.go @@ -431,3 +431,51 @@ func UTCDayBounds(at int64) (start, end int64) { day := time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, time.UTC) return day.UnixNano(), day.Add(24 * time.Hour).UnixNano() } + +// DisplayLocation is the product default timezone (Harbor Desk is Taipei-first). +func DisplayLocation() *time.Location { + return time.FixedZone("Asia/Taipei", 8*60*60) +} + +// LocalDayBounds returns [start, end) unix ns for the local calendar day that contains at. +func LocalDayBounds(at int64, loc *time.Location) (start, end int64) { + if loc == nil { + loc = DisplayLocation() + } + if at <= 0 { + at = NowNano() + } + t := time.Unix(0, at).In(loc) + day := time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, loc) + return day.UnixNano(), day.Add(24 * time.Hour).UnixNano() +} + +// InInboxTimeScope is true when a row belongs on today/7d. +// Posted-in-range is the primary meaning; a just-finished sweep must still +// surface older or undated posts via last_matched_at / created_at, otherwise +// the default pending+today inbox looks empty after patrol. +func InInboxTimeScope(postedAt, lastMatchedAt, createdAt, from, to int64) bool { + if from <= 0 && to <= 0 { + return true + } + in := func(ts int64) bool { + if ts <= 0 { + return false + } + if from > 0 && ts < from { + return false + } + if to > 0 && ts >= to { + return false + } + return true + } + if in(postedAt) { + return true + } + discovered := lastMatchedAt + if discovered <= 0 { + discovered = createdAt + } + return in(discovered) +} diff --git a/apps/backend/internal/module/radar/domain/opportunity_test.go b/apps/backend/internal/module/radar/domain/opportunity_test.go index a0859b8..8d5134f 100644 --- a/apps/backend/internal/module/radar/domain/opportunity_test.go +++ b/apps/backend/internal/module/radar/domain/opportunity_test.go @@ -3,6 +3,7 @@ package domain import ( "errors" "testing" + "time" ) func TestBandFromScore_LockedThresholds(t *testing.T) { @@ -34,6 +35,18 @@ func TestUnknownPublishedTimeIsNotStale(t *testing.T) { } } +func TestTwentyDayOldPostIsNotHardRejected(t *testing.T) { + now := NowNano() + posted := now - 20*24*int64(time.Hour) + if IsStaleHardReject(posted, now) { + t.Fatal("a 20-day-old demand post must still be judgeable, not hard-rejected") + } + old := now - 40*24*int64(time.Hour) + if !IsStaleHardReject(old, now) { + t.Fatal("a 40-day-old post should still be hard-rejected") + } +} + func TestValidateReasons_RequiresAllFive(t *testing.T) { full := fiveReasons() if err := ValidateReasons(full); err != nil { diff --git a/apps/backend/internal/module/radar/domain/opportunity_time_test.go b/apps/backend/internal/module/radar/domain/opportunity_time_test.go new file mode 100644 index 0000000..52b86ff --- /dev/null +++ b/apps/backend/internal/module/radar/domain/opportunity_time_test.go @@ -0,0 +1,43 @@ +package domain + +import ( + "testing" + "time" +) + +func TestLocalDayBoundsUsesTaipeiNotUTC(t *testing.T) { + // 2026-08-13 02:44 UTC is 10:44 in Taipei, same local calendar day. + // 2026-08-13 00:30 Taipei is still 2026-08-12 16:30 UTC. + morningTaipei := time.Date(2026, 8, 13, 0, 30, 0, 0, DisplayLocation()).UnixNano() + start, end := LocalDayBounds(morningTaipei, DisplayLocation()) + utcStart, utcEnd := UTCDayBounds(morningTaipei) + if start == utcStart { + t.Fatalf("local day must not collapse to UTC day: local=%d utc=%d", start, utcStart) + } + if morningTaipei < start || morningTaipei >= end { + t.Fatalf("Taipei 00:30 should sit inside local today [%d, %d)", start, end) + } + if utcEnd <= start { + t.Fatal("expected UTC Aug 12 window to end at Taipei midnight Aug 13") + } +} + +func TestInInboxTimeScopeKeepsResweptAndUndatedRows(t *testing.T) { + from, to := LocalDayBounds(NowNano(), DisplayLocation()) + oldPost := from - int64(48*time.Hour) + if InInboxTimeScope(oldPost, 0, oldPost, from, to) { + t.Fatal("stale post that was not rematched today should stay out of today") + } + if !InInboxTimeScope(oldPost, from+1, oldPost, from, to) { + t.Fatal("a post rematched today must remain visible after patrol") + } + if !InInboxTimeScope(0, 0, from+1, from, to) { + t.Fatal("undated post created today must remain visible after patrol") + } + if !InInboxTimeScope(oldPost, 0, from+1, from, to) { + t.Fatal("older post created today without last_matched must remain visible (no-backfill)") + } + if !InInboxTimeScope(from+1, 0, oldPost, from, to) { + t.Fatal("post published today should stay on today even if created earlier") + } +} diff --git a/apps/backend/internal/module/radar/domain/scoring.go b/apps/backend/internal/module/radar/domain/scoring.go index bed3054..bd1ab6f 100644 --- a/apps/backend/internal/module/radar/domain/scoring.go +++ b/apps/backend/internal/module/radar/domain/scoring.go @@ -11,8 +11,10 @@ const ( WeightFit = 10 ) -// Freshness hard-reject after 14 days. +// Freshness scoring decays after 14 days. Hard-reject is wider so Exa/Threads +// hits from the last month still reach the inbox as stale, not as "found nothing". const MaxFreshnessDays = 14 +const StaleHardRejectDays = 30 // FreshnessScore maps age in hours to the 0–15 freshness dimension score. func FreshnessScore(hours int) int { @@ -50,12 +52,14 @@ func FreshnessHoursSince(postedAt, now int64) int { return h } -// IsStaleHardReject is true when the post is older than 14 days. +// IsStaleHardReject is true when the post is older than a month. +// 14–30 day posts stay scoreable (low freshness) so a just-run patrol can +// still surface a real pain that Threads/Exa returned as an older result. func IsStaleHardReject(postedAt, now int64) bool { if postedAt <= 0 { return false } - return FreshnessHoursSince(postedAt, now) > MaxFreshnessDays*24 + return FreshnessHoursSince(postedAt, now) > StaleHardRejectDays*24 } // SumReasonScores totals dimension scores (capped components assumed already). diff --git a/apps/backend/internal/module/radar/domain/suggest.go b/apps/backend/internal/module/radar/domain/suggest.go index f386bab..c874aab 100644 --- a/apps/backend/internal/module/radar/domain/suggest.go +++ b/apps/backend/internal/module/radar/domain/suggest.go @@ -52,8 +52,8 @@ func NormalizeSuggestUsage(s string) string { CleanSuggestions 收掉空白與重複,丟掉沒有理由的項目,並套用數量上限。 沒有理由的項目直接丟:補一句「AI 建議」等於假裝有理由,比少一則更糟。 -include 關鍵字必須通過 Threads 短詞規則(IsThreadsSearchable);不合規整條丟掉、不截短, -避免產出半截怪詞。exclude 仍用較寬的長度界線(訂閱排除詞可能較長)。 +include 必須能在 Threads 搜到:已合規的原詞保留;過長或超過兩個 token 的先收成短詞變體, +變體也沒有才丟掉。exclude 仍用較寬的長度界線(訂閱排除詞可能較長)。 */ func CleanSuggestions(in []WatchTermSuggestion, limit int) []WatchTermSuggestion { if limit <= 0 || limit > MaxSuggestions { @@ -68,9 +68,17 @@ func CleanSuggestions(in []WatchTermSuggestion, limit int) []WatchTermSuggestion continue } usage := NormalizeSuggestUsage(s.Usage) + basisText := strings.TrimSpace(s.BasisText) if usage == SuggestUsageInclude { if !IsThreadsSearchable(term) { - continue + variants := SearchableTermVariants(term, true) + if len(variants) == 0 { + continue + } + if basisText == "" { + basisText = term + } + term = variants[0] } } else { // exclude: keep broader length; still reject empty after normalize @@ -86,7 +94,7 @@ func CleanSuggestions(in []WatchTermSuggestion, limit int) []WatchTermSuggestion seen[key] = true out = append(out, WatchTermSuggestion{ Term: term, Reason: reason, Usage: usage, - BasisKind: strings.TrimSpace(s.BasisKind), BasisText: strings.TrimSpace(s.BasisText), + BasisKind: strings.TrimSpace(s.BasisKind), BasisText: basisText, }) if len(out) >= limit { break diff --git a/apps/backend/internal/module/radar/domain/term.go b/apps/backend/internal/module/radar/domain/term.go index 10112bc..b568bdd 100644 --- a/apps/backend/internal/module/radar/domain/term.go +++ b/apps/backend/internal/module/radar/domain/term.go @@ -9,11 +9,14 @@ import ( // Threads 搜尋短詞硬約束(中文斷詞差、長字串常查無結果)。 // 一組查詢 = 最多 2 個 token(半形空格分隔);中文 token 2–4 字;整組去掉空格後 ≤12 字元。 const ( - MaxThreadsTokens = 2 - MinCJKTokenRunes = 2 - MaxCJKTokenRunes = 4 - MaxThreadsTermRunes = 12 // 去掉空白後 - MaxExploreTerms = 6 + MaxThreadsTokens = 2 + MinCJKTokenRunes = 2 + MaxCJKTokenRunes = 4 + MaxThreadsTermRunes = 12 // 去掉空白後 + MaxExploreTerms = 6 + // 一則長句最多收成幾組可搜短詞,避免滑窗把訂閱或需求地圖塞滿半截字。 + MaxSearchableVariantsPerTerm = 3 + MaxSearchableVariants = 8 ) // NormalizeSearchTerm trims, converts full-width spaces to half-width, and collapses whitespace. @@ -109,3 +112,180 @@ func isCJK(r rune) bool { unicode.Is(unicode.Katakana, r) || (r >= 0x3000 && r <= 0x303F) // CJK punctuation block — treated as CJK char class but punct rejected above } + +var searchFillers = []string{ + "怎麼辦", "求推薦", "有沒有人", "有人知道", "請問一下", "請問", + "想問", "想找", "有沒有", "可以嗎", "好不好", +} + +// SearchableTermVariants turns a user/product phrase into Threads-searchable queries. +// If the input already passes IsThreadsSearchable it is the only result. +// Longer phrases are shortened: adjacent token pairs, filler tails stripped, +// then first/last 2–4 CJK windows. include=false keeps the broader exclude length. +func SearchableTermVariants(raw string, include bool) []string { + raw = NormalizeSearchTerm(raw) + if raw == "" { + return nil + } + seen := map[string]bool{} + out := make([]string, 0, MaxSearchableVariants) + add := func(term string) { + term = NormalizeSearchTerm(term) + if term == "" || seen[strings.ToLower(term)] { + return + } + if include { + if !IsThreadsSearchable(term) { + return + } + } else if n := utf8.RuneCountInString(term); n < MinTermLen || n > MaxTermLen { + return + } + seen[strings.ToLower(term)] = true + out = append(out, term) + } + + if include && IsThreadsSearchable(raw) { + add(raw) + return out + } + if !include { + add(raw) + if len(out) > 0 { + return out + } + } + + tokens := strings.Fields(raw) + if len(tokens) > MaxThreadsTokens { + for i := 0; i+1 < len(tokens) && len(out) < MaxSearchableVariants; i++ { + add(tokens[i] + " " + tokens[i+1]) + } + for _, tok := range tokens { + if len(out) >= MaxSearchableVariants { + break + } + for _, v := range SearchableTermVariants(tok, include) { + add(v) + if len(out) >= MaxSearchableVariants { + break + } + } + } + return out + } + + parts := splitSearchParts(raw) + if len(parts) == 0 { + parts = []string{raw} + } + for _, part := range parts { + if len(out) >= MaxSearchableVariants { + break + } + if include && IsThreadsSearchable(part) { + add(part) + continue + } + stripped := stripSearchFillers(part) + if stripped != "" && stripped != part { + add(stripped) + } + runes := []rune(stripped) + if len(runes) == 0 { + runes = []rune(part) + } + if len(runes) == 5 { + add(string(runes[:3])) + add(string(runes[3:])) + } + for _, width := range []int{4, 3, 2} { + if len(runes) < width { + continue + } + add(string(runes[:width])) + if len(runes) > width { + add(string(runes[len(runes)-width:])) + } + } + if len(out) >= 3 { + continue + } + for width := 4; width >= 2; width-- { + if len(runes) < width { + continue + } + for start := 0; start+width <= len(runes) && len(out) < MaxSearchableVariants; start++ { + add(string(runes[start : start+width])) + } + } + } + return out +} + +// ExpandSearchTerms keeps already-searchable include terms and shortens the rest. +// Used when persisting a watch and when fanning out a sweep. +func ExpandSearchTerms(terms []string, max int) []string { + if max <= 0 { + max = MaxWatchTerms + } + out := make([]string, 0, max) + seen := map[string]bool{} + add := func(term string) bool { + term = strings.ToLower(NormalizeSearchTerm(term)) + if term == "" || !IsThreadsSearchable(term) || seen[term] { + return len(out) < max + } + seen[term] = true + out = append(out, term) + return len(out) < max + } + for _, raw := range terms { + raw = NormalizeSearchTerm(raw) + if raw == "" { + continue + } + if IsThreadsSearchable(raw) { + if !add(raw) { + return out + } + continue + } + for i, v := range SearchableTermVariants(raw, true) { + if i >= MaxSearchableVariantsPerTerm { + break + } + if !add(v) { + return out + } + } + } + return out +} + +func stripSearchFillers(s string) string { + for _, filler := range searchFillers { + s = strings.ReplaceAll(s, filler, "") + } + return strings.Trim(s, " ,,、。!?!?::的了嗎呢啊喔唷") +} + +func splitSearchParts(raw string) []string { + var b strings.Builder + parts := make([]string, 0, 4) + flush := func() { + if value := strings.TrimSpace(b.String()); value != "" { + parts = append(parts, value) + } + b.Reset() + } + for _, r := range raw { + if unicode.IsLetter(r) || unicode.IsDigit(r) || unicode.Is(unicode.Han, r) || unicode.Is(unicode.Hiragana, r) || unicode.Is(unicode.Katakana, r) { + b.WriteRune(r) + } else { + flush() + } + } + flush() + return parts +} diff --git a/apps/backend/internal/module/radar/domain/term_test.go b/apps/backend/internal/module/radar/domain/term_test.go index 48ee55e..33c8a8e 100644 --- a/apps/backend/internal/module/radar/domain/term_test.go +++ b/apps/backend/internal/module/radar/domain/term_test.go @@ -34,7 +34,7 @@ func TestIsThreadsSearchable(t *testing.T) { bad := []string{ "", "a", - "台北 婚攝 推薦", // 3 tokens + "台北 婚攝 推薦", // 3 tokens "這是一個超長關鍵字超過十二字", // too long "保母 AND 求推薦", "保母#推薦", @@ -50,3 +50,51 @@ func TestIsThreadsSearchable(t *testing.T) { } } } + +func TestSearchableTermVariantsKeepsShortAndShortensLong(t *testing.T) { + got := SearchableTermVariants("保母 求推薦", true) + if len(got) != 1 || got[0] != "保母 求推薦" { + t.Fatalf("short term variants=%v", got) + } + got = SearchableTermVariants("晚上睡覺容易口乾舌燥怎麼辦", true) + if len(got) == 0 { + t.Fatal("long sentence produced no variants") + } + for _, term := range got { + if !IsThreadsSearchable(term) { + t.Fatalf("variant %q is not searchable", term) + } + } + hasDry := false + hasSleep := false + for _, term := range got { + if term == "口乾舌燥" { + hasDry = true + } + if term == "晚上睡覺" { + hasSleep = true + } + } + if !hasDry || !hasSleep { + t.Fatalf("variants=%v want 晚上睡覺 and 口乾舌燥", got) + } + got = SearchableTermVariants("台北 婚攝 推薦 價格", true) + if len(got) == 0 || got[0] != "台北 婚攝" { + t.Fatalf("multi-token variants=%v want 台北 婚攝 first", got) + } +} + +func TestExpandSearchTermsShortensUnsearchable(t *testing.T) { + got := ExpandSearchTerms([]string{"求推薦", "晚上睡覺容易口乾舌燥怎麼辦"}, 20) + if len(got) < 2 { + t.Fatalf("expanded=%v", got) + } + if got[0] != "求推薦" { + t.Fatalf("kept original first, got %v", got) + } + for _, term := range got { + if !IsThreadsSearchable(term) { + t.Fatalf("expanded term %q is not searchable", term) + } + } +} diff --git a/apps/backend/internal/module/radar/domain/watch.go b/apps/backend/internal/module/radar/domain/watch.go index d385897..8c72062 100644 --- a/apps/backend/internal/module/radar/domain/watch.go +++ b/apps/backend/internal/module/radar/domain/watch.go @@ -231,6 +231,13 @@ func (w *RadarWatch) Normalize() error { } func normalizeTerms(in []string, field string, max int) ([]string, error) { + if field == "terms" { + out := ExpandSearchTerms(in, max) + if len(out) > max { + return nil, fmt.Errorf("%w: %s exceeds %d items", ErrValidation, field, max) + } + return out, nil + } out := make([]string, 0, len(in)) seen := map[string]bool{} for _, raw := range in { diff --git a/apps/backend/internal/module/radar/repository/opportunity_list_query_test.go b/apps/backend/internal/module/radar/repository/opportunity_list_query_test.go new file mode 100644 index 0000000..e9a04cf --- /dev/null +++ b/apps/backend/internal/module/radar/repository/opportunity_list_query_test.go @@ -0,0 +1,155 @@ +package repository + +import ( + "context" + "fmt" + "strings" + "testing" + "time" + + "apps/backend/internal/module/radar/domain" + + "go.mongodb.org/mongo-driver/bson" +) + +func TestReviewStateListClauseLegacyPending(t *testing.T) { + pending := fmt.Sprintf("%#v", reviewStateListClause(domain.ReviewPending)) + for _, part := range []string{domain.ReviewPending, "$nin", domain.OppAccepted, domain.OppDismissed} { + if !strings.Contains(pending, part) { + t.Fatalf("legacy pending clause missing %q: %s", part, pending) + } + } + completed := fmt.Sprintf("%#v", reviewStateListClause(domain.ReviewCompleted)) + if !strings.Contains(completed, domain.ReviewCompleted) || !strings.Contains(completed, domain.OppAccepted) { + t.Fatalf("completed clause: %s", completed) + } + removed := fmt.Sprintf("%#v", reviewStateListClause(domain.ReviewRemoved)) + if !strings.Contains(removed, domain.ReviewRemoved) || !strings.Contains(removed, domain.OppDismissed) { + t.Fatalf("removed clause: %s", removed) + } +} + +func TestInboxTimeScopeClauseOldPostedCreatedTodayNoLastMatched(t *testing.T) { + from, to := domain.LocalDayBounds(domain.NowNano(), domain.DisplayLocation()) + old := from - int64(48*time.Hour) + clause := inboxTimeScopeClause(from, to) + + justFound := bson.M{"posted_at": old, "created_at": from + 1} + if !bsonMMatches(clause, justFound) { + t.Fatalf("posted=old created=today last_matched missing must match today $or; clause=%#v", clause) + } + zeroLast := bson.M{"posted_at": old, "created_at": from + 1, "last_matched_at": int64(0)} + if !bsonMMatches(clause, zeroLast) { + t.Fatalf("posted=old created=today last_matched=0 must match today $or; clause=%#v", clause) + } + staleMatch := bson.M{"posted_at": old, "created_at": from + 1, "last_matched_at": from - int64(24*time.Hour)} + if bsonMMatches(clause, staleMatch) { + t.Fatal("old last_matched must not fall back to created_at") + } +} + +func TestListOpportunitiesCreatedTodayOldPostedWithoutUpsert(t *testing.T) { + ctx := context.Background() + m := NewMemory() + now := domain.NowNano() + from, to := domain.LocalDayBounds(now, domain.DisplayLocation()) + row := sampleOpp("seed-no-upsert", []string{"漏水"}, 80) + row.ID = "seed-no-upsert" + row.PostedAt = from - int64(48*time.Hour) + row.CreatedAt = from + 1 + row.LastMatchedAt = 0 + row.ReviewState = "" + m.opportunities[row.ID] = row + + pending, total, err := m.ListOpportunities(ctx, 42, domain.OpportunityListFilter{ + ReviewState: domain.ReviewPending, TimeScope: "today", PostedFrom: from, PostedTo: to, + }) + if err != nil || total != 1 || len(pending) != 1 || pending[0].ID != "seed-no-upsert" { + t.Fatalf("no-backfill just-found row missing from pending today: total=%d list=%d err=%v", total, len(pending), err) + } +} + +func bsonMMatches(query bson.M, doc bson.M) bool { + if or, ok := query["$or"]; ok && !bsonOrMatches(or, doc) { + return false + } + if and, ok := query["$and"]; ok && !bsonAndMatches(and, doc) { + return false + } + for key, cond := range query { + if key == "$or" || key == "$and" { + continue + } + if !bsonFieldMatches(key, cond, doc) { + return false + } + } + return true +} + +func bsonOrMatches(raw any, doc bson.M) bool { + items, ok := raw.(bson.A) + if !ok { + return false + } + for _, item := range items { + clause, ok := item.(bson.M) + if ok && bsonMMatches(clause, doc) { + return true + } + } + return false +} + +func bsonAndMatches(raw any, doc bson.M) bool { + items, ok := raw.(bson.A) + if !ok { + return false + } + for _, item := range items { + clause, ok := item.(bson.M) + if !ok || !bsonMMatches(clause, doc) { + return false + } + } + return true +} + +func bsonFieldMatches(key string, cond any, doc bson.M) bool { + got, exists := doc[key] + switch c := cond.(type) { + case bson.M: + if flag, ok := c["$exists"]; ok { + want, _ := flag.(bool) + if exists != want { + return false + } + } + if gte, ok := c["$gte"]; ok { + if !exists || toQueryInt64(got) < toQueryInt64(gte) { + return false + } + } + if lt, ok := c["$lt"]; ok { + if !exists || toQueryInt64(got) >= toQueryInt64(lt) { + return false + } + } + return true + default: + return exists && toQueryInt64(got) == toQueryInt64(c) + } +} + +func toQueryInt64(v any) int64 { + switch n := v.(type) { + case int64: + return n + case int: + return int64(n) + case float64: + return int64(n) + default: + return 0 + } +} diff --git a/apps/backend/internal/module/radar/repository/opportunity_memory.go b/apps/backend/internal/module/radar/repository/opportunity_memory.go index e31acae..c82ee5e 100644 --- a/apps/backend/internal/module/radar/repository/opportunity_memory.go +++ b/apps/backend/internal/module/radar/repository/opportunity_memory.go @@ -253,11 +253,19 @@ func (m *Memory) ListOpportunities(_ context.Context, ownerUID int64, f domain.O if f.CreatedTo > 0 && o.CreatedAt >= f.CreatedTo { continue } - if f.PostedFrom > 0 && o.PostedAt < f.PostedFrom { - continue - } - if f.PostedTo > 0 && o.PostedAt >= f.PostedTo { - continue + if f.PostedFrom > 0 || f.PostedTo > 0 { + if f.TimeScope == "today" || f.TimeScope == "7d" { + if !domain.InInboxTimeScope(o.PostedAt, o.LastMatchedAt, o.CreatedAt, f.PostedFrom, f.PostedTo) { + continue + } + } else { + if f.PostedFrom > 0 && o.PostedAt < f.PostedFrom { + continue + } + if f.PostedTo > 0 && o.PostedAt >= f.PostedTo { + continue + } + } } matched = append(matched, cloneOpportunity(o)) } diff --git a/apps/backend/internal/module/radar/repository/opportunity_memory_test.go b/apps/backend/internal/module/radar/repository/opportunity_memory_test.go index 2aaff40..3ffc9af 100644 --- a/apps/backend/internal/module/radar/repository/opportunity_memory_test.go +++ b/apps/backend/internal/module/radar/repository/opportunity_memory_test.go @@ -4,6 +4,7 @@ import ( "context" "errors" "testing" + "time" "apps/backend/internal/module/radar/domain" ) @@ -146,6 +147,52 @@ func TestUpdateStatusAndOverride(t *testing.T) { } } +func TestListOpportunitiesLegacyPendingAndResweptToday(t *testing.T) { + ctx := context.Background() + m := NewMemory() + now := domain.NowNano() + from, to := domain.LocalDayBounds(now, domain.DisplayLocation()) + old := from - int64(72*time.Hour) + + legacy := sampleOpp("legacy-pending", []string{"漏水"}, 80) + legacy.ReviewState = "" + legacy.PostedAt = old + legacy.CreatedAt = old + legacy.LastMatchedAt = 0 + if _, err := m.UpsertByExternalID(ctx, legacy); err != nil { + t.Fatal(err) + } + // Re-sweep the same source: merge must not hide the card from today's inbox. + if _, err := m.UpsertByExternalID(ctx, &domain.Opportunity{ + OwnerUID: 42, ExternalID: "legacy-pending", MatchedTerms: []string{"抓漏"}, + }); err != nil { + t.Fatal(err) + } + + pending, total, err := m.ListOpportunities(ctx, 42, domain.OpportunityListFilter{ + ReviewState: domain.ReviewPending, TimeScope: "today", PostedFrom: from, PostedTo: to, + }) + if err != nil || total != 1 || len(pending) != 1 { + t.Fatalf("legacy rematch missing from pending today: total=%d list=%d err=%v", total, len(pending), err) + } + if pending[0].ReviewState != domain.ReviewPending { + t.Fatalf("legacy review_state = %q", pending[0].ReviewState) + } + + undated := sampleOpp("undated-today", []string{"漏水"}, 70) + undated.PostedAt = 0 + undated.CreatedAt = now + if _, err := m.UpsertByExternalID(ctx, undated); err != nil { + t.Fatal(err) + } + pending, total, err = m.ListOpportunities(ctx, 42, domain.OpportunityListFilter{ + ReviewState: domain.ReviewPending, TimeScope: "today", PostedFrom: from, PostedTo: to, + }) + if err != nil || total != 2 { + t.Fatalf("undated created-today row should join pending today: total=%d err=%v", total, err) + } +} + func TestReviewStateGuardAndTombstone(t *testing.T) { ctx := context.Background() m := NewMemory() diff --git a/apps/backend/internal/module/radar/repository/opportunity_mongo.go b/apps/backend/internal/module/radar/repository/opportunity_mongo.go index 8e42a71..49772fe 100644 --- a/apps/backend/internal/module/radar/repository/opportunity_mongo.go +++ b/apps/backend/internal/module/radar/repository/opportunity_mongo.go @@ -37,21 +37,17 @@ func (s *MonStore) UpsertByExternalID(ctx context.Context, o *domain.Opportunity "external_id": externalID, }) if err == nil { - terms := domain.NormalizeMatchedTerms(o.MatchedTerms) - if len(terms) > 0 { - _, uerr := s.opportunities.UpdateOne(ctx, - bson.M{"_id": existing.ID}, - bson.M{ - "$addToSet": bson.M{"matched_terms": bson.M{"$each": terms}}, - "$set": bson.M{"updated_at": domain.NowNano()}, - }) - if uerr != nil { - return nil, uerr - } - // re-read after merge - if rerr := s.opportunities.FindOne(ctx, &existing, bson.M{"_id": existing.ID}); rerr != nil { - return nil, rerr - } + now := domain.NowNano() + set := bson.M{"updated_at": now, "last_matched_at": now} + update := bson.M{"$set": set} + if terms := domain.NormalizeMatchedTerms(o.MatchedTerms); len(terms) > 0 { + update["$addToSet"] = bson.M{"matched_terms": bson.M{"$each": terms}} + } + if _, uerr := s.opportunities.UpdateOne(ctx, bson.M{"_id": existing.ID}, update); uerr != nil { + return nil, uerr + } + if rerr := s.opportunities.FindOne(ctx, &existing, bson.M{"_id": existing.ID}); rerr != nil { + return nil, rerr } for _, match := range o.ProductMatches { if _, merr := s.MergeProductMatch(ctx, o.OwnerUID, existing.ID, match); merr != nil { @@ -85,6 +81,9 @@ func (s *MonStore) UpsertByExternalID(ctx context.Context, o *domain.Opportunity if o.ReviewState == "" { o.ReviewState = reviewStateFor(o) } + if o.LastMatchedAt == 0 { + o.LastMatchedAt = now + } o.ExternalID = externalID _, err = s.opportunities.InsertOne(ctx, o) @@ -242,10 +241,7 @@ func (s *MonStore) ListOpportunities(ctx context.Context, ownerUID int64, f doma q["priority_band"] = f.PriorityBand } if f.ReviewState != "" { - q["$or"] = bson.A{ - bson.M{"review_state": f.ReviewState}, - bson.M{"review_state": bson.M{"$exists": false}, "status": map[string]string{domain.ReviewCompleted: domain.OppAccepted, domain.ReviewRemoved: domain.OppDismissed}[f.ReviewState]}, - } + appendAnd(q, reviewStateListClause(f.ReviewState)) } switch f.MatchState { case "eligible": @@ -255,9 +251,9 @@ func (s *MonStore) ListOpportunities(ctx context.Context, ownerUID int64, f doma case "excluded": q["product_matches.excluded"] = true case "generic": - q["$or"] = bson.A{bson.M{"product_matches": bson.M{"$exists": false}}, bson.M{"product_matches": bson.M{"$size": 0}}} + appendAnd(q, bson.M{"$or": bson.A{bson.M{"product_matches": bson.M{"$exists": false}}, bson.M{"product_matches": bson.M{"$size": 0}}}}) case "stale": - q["posted_at"] = bson.M{"$lt": domain.NowNano() - int64(domain.MaxFreshnessDays)*24*int64(time.Hour)} + appendAnd(q, bson.M{"posted_at": bson.M{"$lt": domain.NowNano() - int64(domain.MaxFreshnessDays)*24*int64(time.Hour)}}) } if f.CreatedFrom > 0 || f.CreatedTo > 0 { rng := bson.M{} @@ -270,14 +266,18 @@ func (s *MonStore) ListOpportunities(ctx context.Context, ownerUID int64, f doma q["created_at"] = rng } if f.PostedFrom > 0 || f.PostedTo > 0 { - rng := bson.M{} - if f.PostedFrom > 0 { - rng["$gte"] = f.PostedFrom + if f.TimeScope == "today" || f.TimeScope == "7d" { + appendAnd(q, inboxTimeScopeClause(f.PostedFrom, f.PostedTo)) + } else { + rng := bson.M{} + if f.PostedFrom > 0 { + rng["$gte"] = f.PostedFrom + } + if f.PostedTo > 0 { + rng["$lt"] = f.PostedTo + } + q["posted_at"] = rng } - if f.PostedTo > 0 { - rng["$lt"] = f.PostedTo - } - q["posted_at"] = rng } total, err := s.opportunities.CountDocuments(ctx, q) @@ -446,3 +446,77 @@ func (s *MonStore) SetOpportunityOverride(ctx context.Context, id string, ov *do } return nil } + +func appendAnd(q bson.M, clause bson.M) { + if len(clause) == 0 { + return + } + if existing, ok := q["$and"].(bson.A); ok { + q["$and"] = append(existing, clause) + return + } + q["$and"] = bson.A{clause} +} + +// reviewStateListClause implements spec §3.2: missing review_state is derived +// from OpportunityStatus so a list query never requires a one-shot migration. +func reviewStateListClause(state string) bson.M { + missing := bson.A{ + bson.M{"review_state": bson.M{"$exists": false}}, + bson.M{"review_state": ""}, + } + switch state { + case domain.ReviewPending: + return bson.M{"$or": bson.A{ + bson.M{"review_state": domain.ReviewPending}, + bson.M{"$and": bson.A{ + bson.M{"$or": missing}, + bson.M{"status": bson.M{"$nin": []string{domain.OppAccepted, domain.OppDismissed}}}, + }}, + }} + case domain.ReviewCompleted: + return bson.M{"$or": bson.A{ + bson.M{"review_state": domain.ReviewCompleted}, + bson.M{"$and": bson.A{ + bson.M{"$or": missing}, + bson.M{"status": domain.OppAccepted}, + }}, + }} + case domain.ReviewRemoved: + return bson.M{"$or": bson.A{ + bson.M{"review_state": domain.ReviewRemoved}, + bson.M{"$and": bson.A{ + bson.M{"$or": missing}, + bson.M{"status": domain.OppDismissed}, + }}, + }} + default: + return bson.M{"review_state": state} + } +} + +func inboxTimeScopeClause(from, to int64) bson.M { + rng := bson.M{} + if from > 0 { + rng["$gte"] = from + } + if to > 0 { + rng["$lt"] = to + } + // Same predicate as domain.InInboxTimeScope: posted in range, or + // last_matched in range, or last_matched missing/0 and created in range. + // Must not gate created_at on posted_at — a just-found older post has a + // real posted_at from yesterday and no last_matched_at (no-backfill). + missingLastMatched := bson.A{ + bson.M{"last_matched_at": bson.M{"$exists": false}}, + bson.M{"last_matched_at": int64(0)}, + } + return bson.M{"$or": bson.A{ + bson.M{"posted_at": rng}, + bson.M{"last_matched_at": rng}, + bson.M{"$and": bson.A{ + bson.M{"$or": missingLastMatched}, + bson.M{"created_at": rng}, + }}, + }} +} diff --git a/apps/backend/internal/module/radar/usecase/demand_map.go b/apps/backend/internal/module/radar/usecase/demand_map.go index 710d05b..2bf9436 100644 --- a/apps/backend/internal/module/radar/usecase/demand_map.go +++ b/apps/backend/internal/module/radar/usecase/demand_map.go @@ -62,23 +62,35 @@ func (s *Service) UpdateDemandMap(ctx context.Context, ownerUID int64, value *do } func baselineDemandMap(product *ProductCatalogProduct) *domain.DemandMap { - mapPhrase := func(text, kind, basisKind string) domain.DemandMapPhrase { - return domain.DemandMapPhrase{Text: strings.TrimSpace(text), Kind: kind, BasisKind: basisKind, BasisText: product.Label, Origin: "product", Enabled: strings.TrimSpace(text) != ""} + mapPhrase := func(text, kind, basisKind, basisText string) domain.DemandMapPhrase { + return domain.DemandMapPhrase{Text: strings.TrimSpace(text), Kind: kind, BasisKind: basisKind, BasisText: basisText, Origin: "product", Enabled: strings.TrimSpace(text) != ""} } - phrases := func(items []string, kind, basisKind string) []domain.DemandMapPhrase { + phrases := func(items []string, kind, basisKind string, include bool) []domain.DemandMapPhrase { out := make([]domain.DemandMapPhrase, 0, len(items)) for _, item := range items { - if strings.TrimSpace(item) != "" { - out = append(out, mapPhrase(item, kind, basisKind)) + raw := strings.TrimSpace(item) + if raw == "" { + continue + } + if !include { + out = append(out, mapPhrase(raw, kind, basisKind, product.Label)) + continue + } + variants := domain.SearchableTermVariants(raw, true) + if len(variants) > domain.MaxSearchableVariantsPerTerm { + variants = variants[:domain.MaxSearchableVariantsPerTerm] + } + for _, term := range variants { + out = append(out, mapPhrase(term, kind, basisKind, raw)) } } return out } - pain := phrases(product.PainPoints, "pain", "pain_point") - scenario := phrases([]string{product.ProductContext}, "scenario", "product_context") - outcomes := phrases(product.MatchTags, "outcome", "match_tag") - solution := phrases(product.ProviderCapabilityTerms, "solution", "provider_capability") - exclusions := phrases(product.ProviderExcludeTerms, "exclusion", "provider_exclude") + pain := phrases(product.PainPoints, "pain", "pain_point", true) + scenario := phrases([]string{product.ProductContext}, "scenario", "product_context", true) + outcomes := phrases(product.MatchTags, "outcome", "match_tag", true) + solution := phrases(product.ProviderCapabilityTerms, "solution", "provider_capability", true) + exclusions := phrases(product.ProviderExcludeTerms, "exclusion", "provider_exclude", false) state := "incomplete" if len(pain) > 0 && len(scenario) > 0 && len(solution) > 0 { state = "ready" diff --git a/apps/backend/internal/module/radar/usecase/demand_map_enrich.go b/apps/backend/internal/module/radar/usecase/demand_map_enrich.go index c9a96bc..ae9b1b3 100644 --- a/apps/backend/internal/module/radar/usecase/demand_map_enrich.go +++ b/apps/backend/internal/module/radar/usecase/demand_map_enrich.go @@ -41,7 +41,7 @@ func (s *Service) EnrichDemandMap(ctx context.Context, ownerUID int64, productID if s.AI == nil && s.AIRegistry == nil && s.ResolveAI == nil { return nil, fmt.Errorf("%w: AI provider unavailable", domain.ErrNotReady) } - prompt := fmt.Sprintf("請只輸出 JSON 物件,根據需求地圖補充使用者會說的短詞,不要使用品牌或產品名稱。痛點=%q;情境=%q;結果=%q;能力=%q。欄位為 pain_phrases、scenario_phrases、desired_outcomes、solution_signals、exclusion_signals、custom_phrases,每則含 text、kind、basis_kind、basis_text、origin=ai、enabled=true。", phraseTexts(current.PainPhrases), phraseTexts(current.ScenarioPhrases), phraseTexts(current.DesiredOutcomes), phraseTexts(current.SolutionSignals)) + prompt := fmt.Sprintf("請只輸出 JSON 物件,根據需求地圖補充使用者會說的短詞,不要使用品牌或產品名稱。痛點=%q;情境=%q;結果=%q;能力=%q。欄位為 pain_phrases、scenario_phrases、desired_outcomes、solution_signals、exclusion_signals、custom_phrases,每則含 text、kind、basis_kind、basis_text、origin=ai、enabled=true。include 類 text 必須是 Threads 可搜短詞:最多 2 個 token、中文每 token 2–4 字、去空格後 ≤12 字、禁止標點。", phraseTexts(current.PainPhrases), phraseTexts(current.ScenarioPhrases), phraseTexts(current.DesiredOutcomes), phraseTexts(current.SolutionSignals)) raw, err := s.completeAI(ctx, ownerUID, prompt) if err != nil { return nil, err @@ -61,9 +61,30 @@ func (s *Service) EnrichDemandMap(ctx context.Context, ownerUID int64, productID p.Text = strings.TrimSpace(p.Text) p.Origin = "ai" p.Enabled = p.Enabled && p.Text != "" - if p.Enabled && !seen[strings.ToLower(p.Text)] { - out = append(out, p) - seen[strings.ToLower(p.Text)] = true + if !p.Enabled { + continue + } + texts := []string{p.Text} + if p.Kind != "exclusion" { + texts = domain.SearchableTermVariants(p.Text, true) + if len(texts) > domain.MaxSearchableVariantsPerTerm { + texts = texts[:domain.MaxSearchableVariantsPerTerm] + } + } + basis := strings.TrimSpace(p.BasisText) + if basis == "" { + basis = p.Text + } + for _, text := range texts { + key := strings.ToLower(text) + if text == "" || seen[key] { + continue + } + cp := p + cp.Text = text + cp.BasisText = basis + out = append(out, cp) + seen[key] = true } } return out diff --git a/apps/backend/internal/module/radar/usecase/demand_map_test.go b/apps/backend/internal/module/radar/usecase/demand_map_test.go index 484aeea..64dd0ef 100644 --- a/apps/backend/internal/module/radar/usecase/demand_map_test.go +++ b/apps/backend/internal/module/radar/usecase/demand_map_test.go @@ -37,6 +37,11 @@ func TestDemandMapBaselineAndOptimisticVersion(t *testing.T) { if first.State != "ready" || first.MapVersion != 1 || len(first.DemandInputVersion) != len("demand-")+16 { t.Fatalf("unexpected baseline: %+v", first) } + for _, p := range append(append([]domain.DemandMapPhrase{}, first.PainPhrases...), append(first.ScenarioPhrases, first.SolutionSignals...)...) { + if !domain.IsThreadsSearchable(p.Text) { + t.Fatalf("baseline phrase %q is not searchable", p.Text) + } + } first.PainPhrases[0].Text = "新的痛點" updated, err := svc.UpdateDemandMap(context.Background(), owner, first, 1) if err != nil { diff --git a/apps/backend/internal/module/radar/usecase/explore.go b/apps/backend/internal/module/radar/usecase/explore.go index dd7a2b0..29926d2 100644 --- a/apps/backend/internal/module/radar/usecase/explore.go +++ b/apps/backend/internal/module/radar/usecase/explore.go @@ -77,10 +77,7 @@ func (s *Service) exploreOpportunities(ctx context.Context, ownerUID int64, rawT if productContext != nil { if dm, derr := s.GetDemandMap(ctx, ownerUID, productID); derr == nil { if plan, perr := BuildQueryPlan(dm, productContext.ProductLabel); perr == nil && plan != nil { - w.Terms = make([]string, 0, len(plan.Groups)) - for _, group := range plan.Groups { - w.Terms = append(w.Terms, group.Query) - } + w = mergeFetchWatch(w, plan) } } } @@ -110,7 +107,7 @@ func (s *Service) exploreOpportunities(ctx context.Context, ownerUID int64, rawT PrefilterReviewCount: prefilter.Review, PrefilterRejectedCount: prefilter.Rejected, }) - created, judged, truncated, failed, judgeCredits, perr := s.ProcessCandidates( + created, _, judged, truncated, failed, judgeCredits, perr := s.ProcessCandidates( ctx, ownerUID, w, profile, sw.ID, cands, nil, ) if perr != nil { diff --git a/apps/backend/internal/module/radar/usecase/judge_persist.go b/apps/backend/internal/module/radar/usecase/judge_persist.go index 9b0a07e..1684c66 100644 --- a/apps/backend/internal/module/radar/usecase/judge_persist.go +++ b/apps/backend/internal/module/radar/usecase/judge_persist.go @@ -29,18 +29,18 @@ func (s *Service) ProcessCandidates( sweepID string, cands []*domain.CandidatePost, alreadyJudged map[string]bool, -) (created, judged, truncated, failed int, credits int, err error) { +) (created, rematched, judged, truncated, failed int, credits int, err error) { matchEvaluated, matchMerged, fitRejected, budgetDeferred := 0, 0, 0, 0 if alreadyJudged == nil { alreadyJudged = map[string]bool{} } maxDaily, err := s.MaxDailyOpportunities(ctx, ownerUID) if err != nil { - return 0, 0, 0, 0, 0, err + return 0, 0, 0, 0, 0, 0, err } todayCount, err := s.Repo.CountToday(ctx, ownerUID, domain.NowNano()) if err != nil { - return 0, 0, 0, 0, 0, err + return 0, 0, 0, 0, 0, 0, err } remaining := maxDaily - int(todayCount) if remaining < 0 { @@ -64,6 +64,7 @@ func (s *Service) ProcessCandidates( pcredits, perr := s.mergeExistingProductCandidate(ctx, ownerUID, watch, c, existing) credits += pcredits judged++ + rematched++ if perr == nil { matchMerged++ } else { @@ -77,7 +78,7 @@ func (s *Service) ProcessCandidates( if sweepID != "" { _, _ = s.Repo.UpdateSweep(ctx, sweepID, domain.SweepDelta{JudgedCount: judged, TruncatedCount: truncated, BudgetDeferredCount: budgetDeferred, CreditsUsed: credits, CreditJudge: credits, MatchEvaluatedCount: matchEvaluated, MatchMergedCount: matchMerged, FitRejectedCount: fitRejected}) } - return 0, judged, truncated, failed, credits, nil + return 0, rematched, judged, truncated, failed, credits, nil } var scored []scoredCandidate @@ -112,6 +113,7 @@ func (s *Service) ProcessCandidates( failed++ } else { credits += pcredits + rematched++ if watch != nil && watch.ContextMode == domain.WatchContextProduct && productMatchFor(existing, watch.ProductID) == nil { matchMerged++ } @@ -158,11 +160,16 @@ func (s *Service) ProcessCandidates( res := sc.result // rejected always stored if reasons ok if res.Status == domain.OppRejected { - if perr := s.persistOne(ctx, ownerUID, watch, sc.cand, res); perr != nil { + inserted, perr := s.persistOne(ctx, ownerUID, watch, sc.cand, res) + if perr != nil { failed++ continue } - created++ + if inserted { + created++ + } else { + rematched++ + } continue } if remaining <= 0 { @@ -170,11 +177,16 @@ func (s *Service) ProcessCandidates( budgetDeferred++ continue } - if perr := s.persistOne(ctx, ownerUID, watch, sc.cand, res); perr != nil { + inserted, perr := s.persistOne(ctx, ownerUID, watch, sc.cand, res) + if perr != nil { failed++ continue } - created++ + if inserted { + created++ + } else { + rematched++ + } remaining-- } @@ -192,15 +204,15 @@ func (s *Service) ProcessCandidates( BudgetDeferredCount: budgetDeferred, } if _, uerr := s.Repo.UpdateSweep(ctx, sweepID, delta); uerr != nil { - return created, judged, truncated, failed, credits, uerr + return created, rematched, judged, truncated, failed, credits, uerr } } - return created, judged, truncated, failed, credits, nil + return created, rematched, judged, truncated, failed, credits, nil } -func (s *Service) persistOne(ctx context.Context, ownerUID int64, watch *domain.RadarWatch, cand *domain.CandidatePost, res *JudgeResult) error { +func (s *Service) persistOne(ctx context.Context, ownerUID int64, watch *domain.RadarWatch, cand *domain.CandidatePost, res *JudgeResult) (inserted bool, err error) { if err := domain.ValidateReasons(res.Reasons); err != nil { - return err + return false, err } watchID := "" if watch != nil { @@ -235,6 +247,7 @@ func (s *Service) persistOne(ctx context.Context, ownerUID int64, watch *domain. EvidenceQualityScore: priority.EvidenceQuality, FreshnessScore: priority.Freshness, DemandEvidence: priority.Evidence, + LastMatchedAt: domain.NowNano(), } if watch != nil && watch.ContextMode == domain.WatchContextProduct { if dm, derr := s.GetDemandMap(ctx, ownerUID, watch.ProductID); derr == nil && dm != nil { @@ -245,11 +258,15 @@ func (s *Service) persistOne(ctx context.Context, ownerUID int64, watch *domain. o.ProductMatches = []*domain.ProductMatch{domain.CloneProductMatch(res.ProductMatch)} } if err := mergeProductMatchIntoOpportunity(o, res.ProductMatch); res.ProductMatch != nil && err != nil { - return err + return false, err } if o.IntentBand == "" { o.ApplyBandFromScore() } - _, err := s.Repo.UpsertByExternalID(ctx, o) - return err + existing, gerr := s.Repo.GetByExternalID(ctx, ownerUID, cand.ExternalID) + already := gerr == nil && existing != nil + if _, err := s.Repo.UpsertByExternalID(ctx, o); err != nil { + return false, err + } + return !already, nil } diff --git a/apps/backend/internal/module/radar/usecase/m2_integration_test.go b/apps/backend/internal/module/radar/usecase/m2_integration_test.go index 49d0f5a..4253b17 100644 --- a/apps/backend/internal/module/radar/usecase/m2_integration_test.go +++ b/apps/backend/internal/module/radar/usecase/m2_integration_test.go @@ -27,10 +27,10 @@ func seedProfileWatch(t *testing.T, svc *Service, owner int64) *domain.RadarWatc t.Helper() ctx := context.Background() p := &domain.ServiceProfile{ - OwnerUID: owner, - Services: []domain.ServiceItem{{Name: "婚禮攝影", Currency: "TWD"}}, + OwnerUID: owner, + Services: []domain.ServiceItem{{Name: "婚禮攝影", Currency: "TWD"}}, ServiceAreas: []string{"TPE"}, - RemoteOk: false, + RemoteOk: false, } if err := p.Normalize(); err != nil { // minimal normalize via save path @@ -98,7 +98,7 @@ func TestM2_OP02_ProviderOfferRejected(t *testing.T) { svc.HitFetch = &fakeHits{ path: domain.SweepPathAPI, hits: []ThreadHit{{ - URL: "https://www.threads.net/@biz/post/2", + URL: "https://www.threads.net/@biz/post/2", Snippet: "婚攝接案中 歡迎洽詢我 限時優惠 dm me", }}, } @@ -232,6 +232,57 @@ func TestM2_QT01_AtCapNoError(t *testing.T) { } } +func TestM2_ManualTriggerRunsAgainAfterDailySuccess(t *testing.T) { + ctx := context.Background() + jobs := jobUC.New(jobRepo.NewMemory()) + svc := New(repository.NewMemory()) + svc.SweepJobs = radarManualJobs{jobs: jobs} + now := domain.NowNano() + w := &domain.RadarWatch{ + ID: "w-live", OwnerUID: 17, Terms: []string{"痛點"}, Status: domain.WatchActive, + CreatedAt: now, UpdatedAt: now, + } + if err := svc.Repo.SaveWatch(ctx, w); err != nil { + t.Fatal(err) + } + daily, err := jobs.ScheduleRadarSweep(ctx, 17, w.ID, now) + if err != nil { + t.Fatal(err) + } + claimed, err := jobs.ClaimNext(ctx, "w") + if err != nil || claimed.ID != daily.ID { + t.Fatalf("claim daily: %v", err) + } + if _, err := jobs.SucceedJob(ctx, daily.ID, "daily done"); err != nil { + t.Fatal(err) + } + manualID, err := svc.TriggerSweep(ctx, 17, w.ID) + if err != nil { + t.Fatal(err) + } + if manualID == daily.ID { + t.Fatal("立即巡邏 must not reuse the finished daily job") + } +} + +type radarManualJobs struct{ jobs *jobUC.Service } + +func (a radarManualJobs) ScheduleRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (string, error) { + j, err := a.jobs.ScheduleRadarSweep(ctx, ownerUID, watchID, runAt) + if err != nil { + return "", err + } + return j.ID, nil +} + +func (a radarManualJobs) ScheduleManualRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (string, error) { + j, err := a.jobs.ScheduleManualRadarSweep(ctx, ownerUID, watchID, runAt) + if err != nil { + return "", err + } + return j.ID, nil +} + func TestM2_ManualTriggerPausedRejected(t *testing.T) { ctx := context.Background() jobs := jobUC.New(jobRepo.NewMemory()) diff --git a/apps/backend/internal/module/radar/usecase/opportunity_ops.go b/apps/backend/internal/module/radar/usecase/opportunity_ops.go index f511d20..ab28909 100644 --- a/apps/backend/internal/module/radar/usecase/opportunity_ops.go +++ b/apps/backend/internal/module/radar/usecase/opportunity_ops.go @@ -73,7 +73,10 @@ func (s *Service) TriggerSweep(ctx context.Context, ownerUID int64, watchID stri if w.Status != domain.WatchActive { return "", fmt.Errorf("%w: only active watches can be swept (status=%s)", domain.ErrValidation, w.Status) } - // Manual trigger is due now. + // Manual trigger is due now and must not reuse a finished daily slot. + if manual, ok := s.SweepJobs.(ManualSweepJobScheduler); ok { + return manual.ScheduleManualRadarSweep(ctx, ownerUID, watchID, domain.NowNano()) + } return s.SweepJobs.ScheduleRadarSweep(ctx, ownerUID, watchID, domain.NowNano()) } diff --git a/apps/backend/internal/module/radar/usecase/product_fit_integration_test.go b/apps/backend/internal/module/radar/usecase/product_fit_integration_test.go index 5ac5a5f..abefb1e 100644 --- a/apps/backend/internal/module/radar/usecase/product_fit_integration_test.go +++ b/apps/backend/internal/module/radar/usecase/product_fit_integration_test.go @@ -16,7 +16,10 @@ func TestProductFitStaleCandidateStillPersistsCompleteMatch(t *testing.T) { if err != nil { t.Fatal(err) } - if result.Status != domain.OppRejected || result.ProductMatch == nil || len(result.ProductMatch.Reasons) != 4 { + if result.Status == domain.OppRejected { + t.Fatalf("a 20-day-old demand post must stay judgeable, got rejected: %+v", result) + } + if result.ProductMatch == nil || len(result.ProductMatch.Reasons) != 4 { t.Fatalf("stale candidate lost product evidence: %+v", result) } if result.ProductMatch.ProductFitScore == 0 { diff --git a/apps/backend/internal/module/radar/usecase/query_plan.go b/apps/backend/internal/module/radar/usecase/query_plan.go index 20868a2..490e7fe 100644 --- a/apps/backend/internal/module/radar/usecase/query_plan.go +++ b/apps/backend/internal/module/radar/usecase/query_plan.go @@ -37,13 +37,39 @@ func BuildQueryPlan(m *domain.DemandMap, productLabel string) (*domain.QueryPlan } productLabel = strings.ToLower(domain.NormalizeSearchTerm(productLabel)) include := func(list []domain.DemandMapPhrase) []domain.DemandMapPhrase { + out := make([]domain.DemandMapPhrase, 0, len(list)) + for _, p := range list { + if !p.Enabled { + continue + } + variants := domain.SearchableTermVariants(p.Text, true) + if len(variants) > domain.MaxSearchableVariantsPerTerm { + variants = variants[:domain.MaxSearchableVariantsPerTerm] + } + basis := strings.TrimSpace(p.BasisText) + if basis == "" { + basis = strings.TrimSpace(p.Text) + } + for _, text := range variants { + if text == "" || strings.ToLower(text) == productLabel { + continue + } + cp := p + cp.Text = text + cp.BasisText = basis + out = append(out, cp) + } + } + return out + } + excludePhrases := func(list []domain.DemandMapPhrase) []domain.DemandMapPhrase { out := make([]domain.DemandMapPhrase, 0, len(list)) for _, p := range list { if !p.Enabled { continue } text := domain.NormalizeSearchTerm(p.Text) - if text == "" || strings.ToLower(text) == productLabel || !domain.IsThreadsSearchable(text) { + if text == "" { continue } p.Text = text @@ -53,11 +79,14 @@ func BuildQueryPlan(m *domain.DemandMap, productLabel string) (*domain.QueryPlan } pains, scenarios := include(m.PainPhrases), include(m.ScenarioPhrases) outcomes, solutions := include(m.DesiredOutcomes), include(m.SolutionSignals) - exclusions := include(m.ExclusionSignals) + exclusions := excludePhrases(m.ExclusionSignals) groups := make([]domain.QueryPlanGroup, 0, domain.MaxExploreTerms) seen := map[string]bool{} add := func(parts ...domain.DemandMapPhrase) { + if len(groups) >= domain.MaxExploreTerms { + return + } terms := make([]string, 0, len(parts)) basisKinds := make([]string, 0, len(parts)) basisTexts := make([]string, 0, len(parts)) @@ -82,12 +111,18 @@ func BuildQueryPlan(m *domain.DemandMap, productLabel string) (*domain.QueryPlan } for _, pain := range pains { add(pain) + if len(groups) >= domain.MaxExploreTerms { + break + } for _, scenario := range scenarios { add(pain, scenario) if len(groups) >= domain.MaxExploreTerms { break } } + if len(groups) >= domain.MaxExploreTerms { + break + } for _, outcome := range outcomes { add(pain, outcome) if len(groups) >= domain.MaxExploreTerms { diff --git a/apps/backend/internal/module/radar/usecase/query_plan_test.go b/apps/backend/internal/module/radar/usecase/query_plan_test.go index 4dbd9ec..027148d 100644 --- a/apps/backend/internal/module/radar/usecase/query_plan_test.go +++ b/apps/backend/internal/module/radar/usecase/query_plan_test.go @@ -38,3 +38,24 @@ func TestBuildQueryPlanRejectsIncomplete(t *testing.T) { t.Fatal("incomplete map must not produce a plan") } } + +func TestBuildQueryPlanShortensLongDemandPhrases(t *testing.T) { + m := &domain.DemandMap{ + ProductID: "p1", DemandInputVersion: "demand-a", MapVersion: 1, State: "ready", + PainPhrases: []domain.DemandMapPhrase{{Text: "晚上睡覺容易口乾舌燥怎麼辦", Kind: "pain", Enabled: true}}, + ScenarioPhrases: []domain.DemandMapPhrase{{Text: "換季日常修護", Kind: "scenario", Enabled: true}}, + SolutionSignals: []domain.DemandMapPhrase{{Text: "保濕", Kind: "solution", Enabled: true}}, + } + plan, err := BuildQueryPlan(m, "舒緩精華") + if err != nil { + t.Fatal(err) + } + if len(plan.Groups) == 0 { + t.Fatal("long demand phrases must still compile into a plan") + } + for _, group := range plan.Groups { + if !domain.IsThreadsSearchable(group.Query) { + t.Fatalf("unsearchable query: %+v", group) + } + } +} diff --git a/apps/backend/internal/module/radar/usecase/suggest.go b/apps/backend/internal/module/radar/usecase/suggest.go index 3c12476..3f39cab 100644 --- a/apps/backend/internal/module/radar/usecase/suggest.go +++ b/apps/backend/internal/module/radar/usecase/suggest.go @@ -6,7 +6,6 @@ import ( "errors" "fmt" "strings" - "unicode" "apps/backend/internal/module/radar/domain" usageDomain "apps/backend/internal/module/usage/domain" @@ -294,70 +293,7 @@ func productSuggestionFallback(p *ProductContextSnapshot, limit int) []domain.Wa // productSearchTermVariants keeps fallback terms short without inventing // product names. Exact short fields win; longer CJK fields yield small windows. func productSearchTermVariants(raw string, include bool) []string { - raw = domain.NormalizeSearchTerm(raw) - if raw == "" { - return nil - } - parts := splitProductSearchParts(raw) - if len(parts) == 0 { - parts = []string{raw} - } - seen := map[string]bool{} - out := make([]string, 0, 8) - add := func(term string) { - term = domain.NormalizeSearchTerm(term) - if term == "" || seen[term] { - return - } - if include { - if !domain.IsThreadsSearchable(term) { - return - } - } else if n := len([]rune(term)); n < domain.MinTermLen || n > domain.MaxTermLen { - return - } - seen[term] = true - out = append(out, term) - } - if domain.IsThreadsSearchable(raw) { - add(raw) - } - for _, part := range parts { - if domain.IsThreadsSearchable(part) { - add(part) - continue - } - runes := []rune(part) - for width := 4; width >= 2; width-- { - if len(runes) < width { - continue - } - for start := 0; start+width <= len(runes) && len(out) < 8; start++ { - add(string(runes[start : start+width])) - } - } - } - return out -} - -func splitProductSearchParts(raw string) []string { - var b strings.Builder - parts := make([]string, 0, 4) - flush := func() { - if value := strings.TrimSpace(b.String()); value != "" { - parts = append(parts, value) - } - b.Reset() - } - for _, r := range raw { - if unicode.IsLetter(r) || unicode.IsDigit(r) || unicode.Is(unicode.Han, r) || unicode.Is(unicode.Hiragana, r) || unicode.Is(unicode.Katakana, r) { - b.WriteRune(r) - } else { - flush() - } - } - flush() - return parts + return domain.SearchableTermVariants(raw, include) } /* diff --git a/apps/backend/internal/module/radar/usecase/suggest_test.go b/apps/backend/internal/module/radar/usecase/suggest_test.go index 3ad8e7b..182ce7e 100644 --- a/apps/backend/internal/module/radar/usecase/suggest_test.go +++ b/apps/backend/internal/module/radar/usecase/suggest_test.go @@ -160,7 +160,7 @@ func TestSuggestRespectsLimit(t *testing.T) { } func TestSuggestDropsUnusableItems(t *testing.T) { - // 沒理由、太短、重複、超過 Threads 短詞規則的 include 都要丟掉。 + // 沒理由、太短、重複丟掉;過長 include 收成可搜短詞,不再整條丟。 svc, _, ctx := suggestService(t, `[ {"term":"婚攝 求推薦","reason":"在找攝影師的人常這樣問","usage":"include"}, {"term":"沒有理由的詞","reason":" ","usage":"include"}, @@ -173,8 +173,13 @@ func TestSuggestDropsUnusableItems(t *testing.T) { if err != nil { t.Fatalf("suggest: %v", err) } - if len(list) != 1 || list[0].Term != "婚攝 求推薦" { - t.Fatalf("got %+v, want only the one usable suggestion", list) + if len(list) != 2 || list[0].Term != "婚攝 求推薦" || list[1].Term != "台北 婚攝" { + t.Fatalf("got %+v, want shortened multi-token plus the original short term", list) + } + for _, item := range list { + if item.Usage == domain.SuggestUsageInclude && !domain.IsThreadsSearchable(item.Term) { + t.Fatalf("include not searchable: %+v", item) + } } } diff --git a/apps/backend/internal/module/radar/usecase/sweep_fetch.go b/apps/backend/internal/module/radar/usecase/sweep_fetch.go index e07f551..9d47b38 100644 --- a/apps/backend/internal/module/radar/usecase/sweep_fetch.go +++ b/apps/backend/internal/module/radar/usecase/sweep_fetch.go @@ -59,7 +59,7 @@ func (s *Service) FetchCandidates(ctx context.Context, ownerUID int64, w *domain if w == nil { return nil, "", 0, fmt.Errorf("%w: watch required", domain.ErrValidation) } - terms := w.Terms + terms := domain.ExpandSearchTerms(w.Terms, domain.MaxWatchTerms) if len(terms) == 0 { return nil, "", 0, fmt.Errorf("%w: watch has no terms", domain.ErrValidation) } @@ -104,11 +104,15 @@ func (s *Service) FetchCandidates(ctx context.Context, ownerUID int64, w *domain out := make([]*domain.CandidatePost, 0, len(hits)) for _, h := range hits { + permalink := strings.TrimSpace(h.URL) text := strings.TrimSpace(h.Snippet) if text == "" { text = strings.TrimSpace(h.Title) } - if text == "" { + if text == "" && permalink != "" { + text = permalink + } + if text == "" || permalink == "" { continue } blob := strings.ToLower(text + " " + h.Title) @@ -122,10 +126,6 @@ func (s *Service) FetchCandidates(ctx context.Context, ownerUID int64, w *domain if skip { continue } - permalink := strings.TrimSpace(h.URL) - if permalink == "" { - continue - } term := matchingWatchTerm(blob, terms) class := classifyCandidate(blob) out = append(out, &domain.CandidatePost{ diff --git a/apps/backend/internal/module/radar/usecase/sweep_fetch_test.go b/apps/backend/internal/module/radar/usecase/sweep_fetch_test.go new file mode 100644 index 0000000..59abd87 --- /dev/null +++ b/apps/backend/internal/module/radar/usecase/sweep_fetch_test.go @@ -0,0 +1,65 @@ +package usecase + +import ( + "context" + "fmt" + "strings" + "testing" + + "apps/backend/internal/module/radar/domain" +) + +func TestMergeFetchWatchKeepsSubscriberTerms(t *testing.T) { + w := &domain.RadarWatch{Terms: []string{"求推薦", "泛紅"}, ExcludeTerms: []string{"抽獎"}} + plan := &domain.QueryPlan{Groups: []domain.QueryPlanGroup{ + {Query: "換季 不適", Exclude: []string{"業配"}}, + {Query: "求推薦"}, + }} + got := mergeFetchWatch(w, plan) + if len(got.Terms) != 3 || got.Terms[0] != "求推薦" || got.Terms[1] != "泛紅" || got.Terms[2] != "換季 不適" { + t.Fatalf("terms=%v want subscriber keywords first, then new plan queries", got.Terms) + } + if len(got.ExcludeTerms) != 2 { + t.Fatalf("exclude=%v", got.ExcludeTerms) + } +} + +func TestMergeFetchWatchExpandsLongSubscriberTerms(t *testing.T) { + w := &domain.RadarWatch{Terms: []string{"晚上睡覺容易口乾舌燥怎麼辦"}} + got := mergeFetchWatch(w, nil) + if len(got.Terms) == 0 { + t.Fatal("long subscriber term must expand before fetch") + } + for _, term := range got.Terms { + if !domain.IsThreadsSearchable(term) { + t.Fatalf("fetch term %q is not searchable", term) + } + } +} + +func TestHumanFetchErrorCrawlerSessionExpired(t *testing.T) { + got := humanFetchError(fmt.Errorf(`Chrome crawler status 422: {"error":"crawler session expired"}`)) + if !strings.Contains(got, "Chrome 登入已過期") || !strings.Contains(got, "設定") { + t.Fatalf("human message=%q", got) + } +} + +func TestFetchCandidatesKeepsTitleOnlyHits(t *testing.T) { + svc := New(nil) + svc.HitFetch = HitFetcherFunc(func(context.Context, int64, []string, int) ([]ThreadHit, string, error) { + return []ThreadHit{{ + URL: "https://www.threads.net/@a/post/xyz", + Title: "皮膚泛紅怎麼辦", + }}, "api", nil + }) + // bill path: HitFetch returns api so FetchCandidates will try to bill. + // Avoid billing by setting path after... FetchCandidates bills API path. + // Use empty Usage so bill is a no-op if implemented that way. + cands, _, _, err := svc.FetchCandidates(context.Background(), 1, &domain.RadarWatch{Terms: []string{"泛紅"}}, 10) + if err != nil { + t.Fatalf("fetch: %v", err) + } + if len(cands) != 1 || cands[0].Text == "" { + t.Fatalf("title-only hit dropped: %+v err=%v", cands, err) + } +} diff --git a/apps/backend/internal/module/radar/usecase/sweep_run.go b/apps/backend/internal/module/radar/usecase/sweep_run.go index c2ecce9..c5a8e4d 100644 --- a/apps/backend/internal/module/radar/usecase/sweep_run.go +++ b/apps/backend/internal/module/radar/usecase/sweep_run.go @@ -12,6 +12,7 @@ import ( type SweepRunResult struct { Sweep *domain.RadarSweep Created int + Rematched int Judged int Truncated int FailedJudges int @@ -85,17 +86,10 @@ func (s *Service) RunSweep(ctx context.Context, ownerUID int64, watchID, jobID s already[id] = true } - fetchWatch := w + fetchWatch := mergeFetchWatch(w, nil) if w.ContextMode == domain.WatchContextProduct { if plan, perr := s.BuildProductQueryPlan(ctx, ownerUID, w.ProductID); perr == nil && plan != nil && len(plan.Groups) > 0 { - planned := *w - planned.Terms = make([]string, 0, len(plan.Groups)) - planned.ExcludeTerms = append([]string(nil), w.ExcludeTerms...) - for _, group := range plan.Groups { - planned.Terms = append(planned.Terms, group.Query) - planned.ExcludeTerms = append(planned.ExcludeTerms, group.Exclude...) - } - fetchWatch = &planned + fetchWatch = mergeFetchWatch(w, plan) } } cands, path, fetchCredits, ferr := s.FetchCandidates(ctx, ownerUID, fetchWatch, 40) @@ -135,7 +129,7 @@ func (s *Service) RunSweep(ctx context.Context, ownerUID int64, watchID, jobID s } } - created, judged, truncated, failed, judgeCredits, perr := s.ProcessCandidates( + created, rematched, judged, truncated, failed, judgeCredits, perr := s.ProcessCandidates( ctx, ownerUID, w, profile, sw.ID, cands, already, ) if perr != nil { @@ -162,6 +156,7 @@ func (s *Service) RunSweep(ctx context.Context, ownerUID int64, watchID, jobID s return &SweepRunResult{ Sweep: sw, Created: created, + Rematched: rematched, Judged: judged, Truncated: truncated, FailedJudges: failed, @@ -199,15 +194,50 @@ func (s *Service) notifySweepFailed(ctx context.Context, ownerUID int64, sweepID return s.Notifier.NotifySweepFailed(ctx, ownerUID, sweepID, watchID, reason) } +// mergeFetchWatch keeps the subscriber's own keywords and adds demand-map +// queries. Replacing the watch terms with only plan groups made patrols +// miss the phrases the user actually typed. +func mergeFetchWatch(w *domain.RadarWatch, plan *domain.QueryPlan) *domain.RadarWatch { + if w == nil { + return nil + } + planned := *w + terms := append([]string{}, w.Terms...) + excludes := append([]string{}, w.ExcludeTerms...) + seen := map[string]bool{} + for _, t := range terms { + seen[strings.ToLower(strings.TrimSpace(t))] = true + } + if plan != nil { + for _, group := range plan.Groups { + q := strings.TrimSpace(group.Query) + if q != "" && !seen[strings.ToLower(q)] { + terms = append(terms, q) + seen[strings.ToLower(q)] = true + } + excludes = append(excludes, group.Exclude...) + } + } + planned.Terms = domain.ExpandSearchTerms(terms, domain.MaxWatchTerms) + planned.ExcludeTerms = excludes + return &planned +} + func humanFetchError(err error) string { if err == nil { return "抓取失敗" } msg := err.Error() + low := strings.ToLower(msg) // never include token-like blobs - if strings.Contains(strings.ToLower(msg), "bearer ") { + if strings.Contains(low, "bearer ") { return "抓取路徑不可用" } + if strings.Contains(low, "crawler session expired") || + strings.Contains(low, "crawler session is invalid") || + strings.Contains(low, "crawler session required") { + return "今天沒巡到:Chrome 登入已過期。請到設定重新同步已登入的 Threads 分頁,或先關掉開發模式改走 API 搜尋。" + } if len(msg) > 200 { msg = msg[:200] } diff --git a/apps/backend/internal/module/radar/usecase/sweep_schedule.go b/apps/backend/internal/module/radar/usecase/sweep_schedule.go index 2349617..0026152 100644 --- a/apps/backend/internal/module/radar/usecase/sweep_schedule.go +++ b/apps/backend/internal/module/radar/usecase/sweep_schedule.go @@ -18,6 +18,12 @@ type SweepJobScheduler interface { ScheduleRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (jobID string, err error) } +// ManualSweepJobScheduler is optional. TriggerSweep uses it so 「立即巡邏」 +// is not swallowed by a finished same-day daily job. +type ManualSweepJobScheduler interface { + ScheduleManualRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (jobID string, err error) +} + // SweepJobSchedulerFunc adapts a function to SweepJobScheduler. type SweepJobSchedulerFunc func(ctx context.Context, ownerUID int64, watchID string, runAt int64) (jobID string, err error) diff --git a/apps/backend/internal/module/radar/usecase/watch_test.go b/apps/backend/internal/module/radar/usecase/watch_test.go index 647407a..5f49452 100644 --- a/apps/backend/internal/module/radar/usecase/watch_test.go +++ b/apps/backend/internal/module/radar/usecase/watch_test.go @@ -69,6 +69,25 @@ func TestCreateWatchRejectsBadInput(t *testing.T) { } } +func TestCreateWatchExpandsUnsearchableTerms(t *testing.T) { + svc, ctx := serviceWithProfile(t, 5) + w, err := svc.CreateWatch(ctx, 42, WatchInput{ + Terms: []string{"晚上睡覺容易口乾舌燥怎麼辦"}, + Enabled: true, + }) + if err != nil { + t.Fatalf("create: %v", err) + } + if len(w.Terms) == 0 { + t.Fatal("long sentence must expand into searchable terms") + } + for _, term := range w.Terms { + if !domain.IsThreadsSearchable(term) { + t.Fatalf("persisted term %q is not searchable", term) + } + } +} + func TestWatchPauseResumeRoundTrip(t *testing.T) { svc, ctx := serviceWithProfile(t, 5) diff --git a/apps/backend/internal/module/scout/domain/domain.go b/apps/backend/internal/module/scout/domain/domain.go index 6c19947..0d2be72 100644 --- a/apps/backend/internal/module/scout/domain/domain.go +++ b/apps/backend/internal/module/scout/domain/domain.go @@ -7,13 +7,14 @@ import ( ) var ( - ErrNotFound = errors.New("scout not found") - ErrForbidden = errors.New("scout forbidden") - ErrValidation = errors.New("scout validation") - ErrNoCrawlerSession = errors.New("crawler session required when dev_mode enabled") - ErrTopicRemoved = errors.New("ScoutTopic CRUD removed") - ErrHasProducts = errors.New("brand has products; remove products first") - ErrIllegalRunStatus = errors.New("illegal scout run status transition") + ErrNotFound = errors.New("scout not found") + ErrForbidden = errors.New("scout forbidden") + ErrValidation = errors.New("scout validation") + ErrNoCrawlerSession = errors.New("crawler session required when dev_mode enabled") + ErrCrawlerSessionExpired = errors.New("crawler session expired") + ErrTopicRemoved = errors.New("ScoutTopic CRUD removed") + ErrHasProducts = errors.New("brand has products; remove products first") + ErrIllegalRunStatus = errors.New("illegal scout run status transition") ) func NowNano() int64 { return time.Now().UTC().UnixNano() } diff --git a/apps/backend/internal/module/scout/usecase/chrome_crawler_provider.go b/apps/backend/internal/module/scout/usecase/chrome_crawler_provider.go index c07c0c5..64ebafb 100644 --- a/apps/backend/internal/module/scout/usecase/chrome_crawler_provider.go +++ b/apps/backend/internal/module/scout/usecase/chrome_crawler_provider.go @@ -4,11 +4,14 @@ import ( "bytes" "context" "encoding/json" + "errors" "fmt" "io" "net/http" "strings" "time" + + "apps/backend/internal/module/scout/domain" ) // ChromeCrawlerProvider is the private worker-to-browser boundary. The @@ -43,7 +46,12 @@ func (p *HTTPCrawlerProvider) ResolveMediaID(ctx context.Context, storageState, defer res.Body.Close() raw, _ := io.ReadAll(io.LimitReader(res.Body, 64<<10)) if res.StatusCode < 200 || res.StatusCode >= 300 { - return "", fmt.Errorf("Chrome resolver status %d: %s", res.StatusCode, truncate(string(raw), 160)) + body := truncate(string(raw), 160) + err := fmt.Errorf("Chrome resolver status %d: %s", res.StatusCode, body) + if isCrawlerSessionDead(err) { + return "", fmt.Errorf("%w (%s)", domain.ErrCrawlerSessionExpired, body) + } + return "", err } var out struct { MediaID string `json:"media_id"` @@ -104,7 +112,12 @@ func (p *HTTPCrawlerProvider) SearchChrome(ctx context.Context, storageState str defer res.Body.Close() raw, _ := io.ReadAll(io.LimitReader(res.Body, 1<<20)) if res.StatusCode < 200 || res.StatusCode >= 300 { - return nil, fmt.Errorf("Chrome crawler status %d: %s", res.StatusCode, truncate(string(raw), 160)) + body := truncate(string(raw), 160) + err := fmt.Errorf("Chrome crawler status %d: %s", res.StatusCode, body) + if isCrawlerSessionDead(err) { + return nil, fmt.Errorf("%w (%s)", domain.ErrCrawlerSessionExpired, body) + } + return nil, err } var out struct { Posts []struct { @@ -165,3 +178,16 @@ func isSoftAged(publishedAtNano int64, softDays int) bool { cutoff := time.Now().UTC().AddDate(0, 0, -softDays).UnixNano() return publishedAtNano < cutoff } + +func isCrawlerSessionDead(err error) bool { + if err == nil { + return false + } + if errors.Is(err, domain.ErrNoCrawlerSession) || errors.Is(err, domain.ErrCrawlerSessionExpired) { + return true + } + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "crawler session expired") || + strings.Contains(msg, "crawler session is invalid") || + strings.Contains(msg, "crawler session required") +} diff --git a/apps/backend/internal/module/scout/usecase/search_hits_only_test.go b/apps/backend/internal/module/scout/usecase/search_hits_only_test.go index be09a51..92aefba 100644 --- a/apps/backend/internal/module/scout/usecase/search_hits_only_test.go +++ b/apps/backend/internal/module/scout/usecase/search_hits_only_test.go @@ -2,9 +2,11 @@ package usecase import ( "context" + "fmt" "testing" "apps/backend/internal/module/scout/domain" + "apps/backend/internal/module/scout/repository" ) type recordingProvider struct { @@ -39,8 +41,8 @@ func TestSearchHitsOnlyFansOutPerTermAndDedupes(t *testing.T) { if path != domain.PathAPI { t.Fatalf("path=%q want %q", path, domain.PathAPI) } - if len(prov.calls) != 2 { - t.Fatalf("provider calls=%d want 2 (fan-out), calls=%v", len(prov.calls), prov.calls) + if len(prov.calls) < 2 { + t.Fatalf("provider calls=%d want >=2 (fan-out, plus sparse top-up)", len(prov.calls)) } for _, c := range prov.calls { if len(c) != 1 { @@ -129,6 +131,87 @@ func TestMergeHitsDedupe(t *testing.T) { } } +type alwaysDevMode struct{} + +func (alwaysDevMode) DevModeEnabled(context.Context, int64) (bool, error) { return true, nil } + +type expiredCrawler struct{ calls int } + +func (c *expiredCrawler) SearchChrome(context.Context, string, []string, int) ([]ThreadSearchResult, error) { + c.calls++ + return nil, fmt.Errorf(`Chrome crawler status 422: {"error":"crawler session expired"}`) +} + +func (c *expiredCrawler) ResolveMediaID(context.Context, string, string) (string, error) { + return "", fmt.Errorf("crawler session expired") +} + +func TestSearchHitsOnlyFallsBackToAPIWhenCrawlerSessionExpired(t *testing.T) { + crawler := &expiredCrawler{} + prov := &recordingProvider{} + svc := &Service{ + Provider: prov, + Crawler: crawler, + Settings: alwaysDevMode{}, + Repo: repository.NewMemory(), + SessionSecret: "test-crawler-session-secret", + } + if err := svc.SetCrawlerSession(context.Background(), 1, `{"cookies":[{"domain":".threads.net","expires":4102444800}]}`); err != nil { + t.Fatalf("seed session: %v", err) + } + hits, path, err := svc.SearchHitsOnly(context.Background(), 1, []string{"保母", "求推薦"}, 20) + if err != nil { + t.Fatalf("expired crawler must fall back to api: %v", err) + } + if path != domain.PathAPI { + t.Fatalf("path=%q want api after crawler session expired", path) + } + if len(hits) == 0 { + t.Fatal("api fallback returned no hits") + } + if crawler.calls != 1 { + t.Fatalf("dead session should stop crawler fan-out, calls=%d", crawler.calls) + } + if len(prov.calls) == 0 { + t.Fatal("api provider was not used") + } +} + +func TestSearchHitsOnlyFallsBackToAPIWhenCrawlerSessionMissing(t *testing.T) { + prov := &recordingProvider{} + svc := &Service{Provider: prov, Settings: alwaysDevMode{}} + hits, path, err := svc.SearchHitsOnly(context.Background(), 1, []string{"保母"}, 10) + if err != nil { + t.Fatalf("missing session must fall back to api: %v", err) + } + if path != domain.PathAPI || len(hits) == 0 { + t.Fatalf("path=%q hits=%d", path, len(hits)) + } +} + +func TestSearchHitsOnlyErrorsWhenCrawlerExpiredAndAPIUnavailable(t *testing.T) { + crawler := &expiredCrawler{} + svc := &Service{ + Crawler: crawler, + Settings: alwaysDevMode{}, + Repo: repository.NewMemory(), + SessionSecret: "test-crawler-session-secret", + } + if err := svc.SetCrawlerSession(context.Background(), 1, `{"cookies":[{"domain":".threads.net","expires":4102444800}]}`); err != nil { + t.Fatalf("seed session: %v", err) + } + _, path, err := svc.SearchHitsOnly(context.Background(), 1, []string{"保母"}, 10) + if err == nil { + t.Fatal("both paths down must error") + } + if path != domain.PathCrawler { + t.Fatalf("path=%q want crawler", path) + } + if !isCrawlerSessionDead(err) { + t.Fatalf("err=%v want session-dead", err) + } +} + func TestCanonicalPostIdentityDedupesThreadsURLAliases(t *testing.T) { aliases := []string{ "https://www.threads.net/@alice/post/AbC123?xmt=AQG", diff --git a/apps/backend/internal/module/scout/usecase/search_pipeline.go b/apps/backend/internal/module/scout/usecase/search_pipeline.go index 2db5e9c..27ca66b 100644 --- a/apps/backend/internal/module/scout/usecase/search_pipeline.go +++ b/apps/backend/internal/module/scout/usecase/search_pipeline.go @@ -138,6 +138,10 @@ func (r *searchPipelineRunner) runStageForTerms(perQuery int, source func(contex if err != nil { r.lastErr = err r.diagnostics.SourceUnavailable = true + if isCrawlerSessionDead(err) { + r.source = nil + break + } } if err == nil { r.sourceOK = true diff --git a/apps/backend/internal/module/scout/usecase/service.go b/apps/backend/internal/module/scout/usecase/service.go index d76e30e..af43bf7 100644 --- a/apps/backend/internal/module/scout/usecase/service.go +++ b/apps/backend/internal/module/scout/usecase/service.go @@ -515,31 +515,49 @@ func (s *Service) SearchHitsOnly(ctx context.Context, ownerUID int64, terms []st } } if devMode { - storageState, serr := s.GetCrawlerSessionToken(ctx, ownerUID) - if serr != nil { - return nil, domain.PathCrawler, domain.ErrNoCrawlerSession + var storageState string + var serr error + if s.Repo != nil { + storageState, serr = s.GetCrawlerSessionToken(ctx, ownerUID) + } else { + serr = domain.ErrNoCrawlerSession } - path = domain.PathCrawler - if s.Crawler == nil { - return nil, path, fmt.Errorf("Chrome crawler is not configured") + if serr != nil || strings.TrimSpace(storageState) == "" { + if s.Provider == nil { + return nil, domain.PathCrawler, domain.ErrNoCrawlerSession + } + logx.Infof("scout SearchHitsOnly: no crawler session uid=%d; falling back to api", ownerUID) + } else if s.Crawler == nil { + if s.Provider == nil { + return nil, domain.PathCrawler, fmt.Errorf("Chrome crawler is not configured") + } + logx.Infof("scout SearchHitsOnly: crawler not configured uid=%d; falling back to api", ownerUID) + } else { + path = domain.PathCrawler + hits, err = fanOutSearch(ctx, terms, perQuery, func(ctx context.Context, q string, n int) ([]ThreadSearchResult, error) { + return s.Crawler.SearchChrome(ctx, storageState, []string{q}, n) + }) + if err == nil { + hits = fillSearchHitsToTarget(ctx, s, terms, hits, limit, path, storageState) + return capHits(hits, limit), path, nil + } + if s.Provider == nil { + return nil, path, err + } + logx.Errorf("scout SearchHitsOnly: crawler failed uid=%d: %v; falling back to api", ownerUID, err) } - hits, err = fanOutSearch(ctx, terms, perQuery, func(ctx context.Context, q string, n int) ([]ThreadSearchResult, error) { - return s.Crawler.SearchChrome(ctx, storageState, []string{q}, n) - }) - if err != nil { - return nil, path, err - } - return capHits(hits, limit), path, nil } if s.Provider == nil { return nil, path, fmt.Errorf("scout search provider is not configured") } + path = domain.PathAPI hits, err = fanOutSearch(ctx, terms, perQuery, func(ctx context.Context, q string, n int) ([]ThreadSearchResult, error) { return s.Provider.SearchThreads(ctx, []string{q}, n) }) if err != nil { return nil, path, err } + hits = fillSearchHitsToTarget(ctx, s, terms, hits, limit, path, "") return capHits(hits, limit), path, nil } @@ -761,6 +779,9 @@ func fanOutSearch(ctx context.Context, terms []string, perQuery int, search func if firstErr == nil { firstErr = err } + if isCrawlerSessionDead(err) { + break + } continue } for _, hit := range hits { diff --git a/apps/backend/internal/svc/service_context.go b/apps/backend/internal/svc/service_context.go index e674a86..5fc4293 100644 --- a/apps/backend/internal/svc/service_context.go +++ b/apps/backend/internal/svc/service_context.go @@ -247,14 +247,8 @@ func NewServiceContext(c config.Config) *ServiceContext { // 商機回覆一鍵送出:同一條 Outbox 佇列+同一套 crawler media id 解析,不重造第二套送出路徑。 radarSvc.ReplyQueue = &scoutReplyQueue{Studio: studioSvc} radarSvc.MediaResolver = scoutSvc - // 每日巡與手動觸發共用 job.ScheduleRadarSweep(同 template、同日去重)。 - radarSvc.SweepJobs = radarUC.SweepJobSchedulerFunc(func(ctx context.Context, ownerUID int64, watchID string, runAt int64) (string, error) { - j, err := jobs.ScheduleRadarSweep(ctx, ownerUID, watchID, runAt) - if err != nil { - return "", err - } - return j.ID, nil - }) + // 每日巡走同日去重;立即巡邏另開 manual job,避免當天已成功的日巡把按鈕吞掉。 + radarSvc.SweepJobs = radarSweepJobs{jobs: jobs} zipPath := findExtensionZip() @@ -552,6 +546,32 @@ func (b *metaMediaBridge) ListProfilePosts(ctx context.Context, accessToken, use return out, nil } +type radarSweepJobs struct { + jobs *jobUC.Service +} + +func (a radarSweepJobs) ScheduleRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (string, error) { + if a.jobs == nil { + return "", fmt.Errorf("jobs not configured") + } + j, err := a.jobs.ScheduleRadarSweep(ctx, ownerUID, watchID, runAt) + if err != nil { + return "", err + } + return j.ID, nil +} + +func (a radarSweepJobs) ScheduleManualRadarSweep(ctx context.Context, ownerUID int64, watchID string, runAt int64) (string, error) { + if a.jobs == nil { + return "", fmt.Errorf("jobs not configured") + } + j, err := a.jobs.ScheduleManualRadarSweep(ctx, ownerUID, watchID, runAt) + if err != nil { + return "", err + } + return j.ID, nil +} + // personaJobBridge adapts job.Service → studio.PersonaAnalyzeScheduler (jobID only). type personaJobBridge struct { Jobs *jobUC.Service diff --git a/apps/web/src/App.tsx b/apps/web/src/App.tsx index 94ad146..9f78675 100644 --- a/apps/web/src/App.tsx +++ b/apps/web/src/App.tsx @@ -87,9 +87,9 @@ export default function App() { } /> } /> } /> - } /> + } /> } /> - } /> + } /> } /> } /> } /> diff --git a/apps/web/src/components/radar/OpportunityDetailDrawer.tsx b/apps/web/src/components/radar/OpportunityDetailDrawer.tsx index 81c3f96..d2fb213 100644 --- a/apps/web/src/components/radar/OpportunityDetailDrawer.tsx +++ b/apps/web/src/components/radar/OpportunityDetailDrawer.tsx @@ -38,15 +38,15 @@ export function OpportunityDetailDrawer({ opportunity, onClose, onAccept, onComp {opportunity.reasons.length ?

原始判定

{opportunity.reasons.map((reason) =>
{reason.dimension}{reason.score}{reason.reason}
)}
: null}
- {pending && !accepted && onAccept ? ( - - ) : null} {pending && onComplete ? ( - + + ) : null} + {pending && !accepted && onAccept ? ( + ) : null} 開啟 Threads 原文
- {pending ?

加入名單會保留這位使用者供後續追蹤;只標示已處理不會建立名單。兩者都不扣點。

: null} + {pending ?

先看痛點與產品理由。留下或丟掉即可;加入名單只在你要追這個人時才需要。

: null} ); diff --git a/apps/web/src/components/radar/OpportunityInboxCard.tsx b/apps/web/src/components/radar/OpportunityInboxCard.tsx index 1cec3fe..bba2b75 100644 --- a/apps/web/src/components/radar/OpportunityInboxCard.tsx +++ b/apps/web/src/components/radar/OpportunityInboxCard.tsx @@ -67,14 +67,14 @@ export function OpportunityInboxCard({ opportunity, onOpen, onAccept, onComplete 查看 Threads 原文
+ {pending ? : null} + {pending ? : null} + {(pending || reviewState === "completed") && !accepted ? ( - ) : null} - {pending ? : null} - {pending ? : null} - {reviewState === "removed" ? : null} {reviewState === "completed" && accepted && opportunity.contact_id ? ( 前往名單 @@ -83,15 +83,15 @@ export function OpportunityInboxCard({ opportunity, onOpen, onAccept, onComplete
{removeOpen ? (
- 為什麼不適合? -

選原因後,這筆會移到「已移除」且不會再次出現在待處理。這個動作不扣點。

+ 為什麼丟掉? +

選原因後會移出「新找到」,之後巡邏不會再把同一篇推上來。這個動作不扣點。

{reason === "other" ?