1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
//! Bounded-concurrency external-source (JIRA / GitHub / …) resolution for the
//! classification pipeline.
//!
//! Why: prior to issue #2719 the pipeline resolved external tickets strictly
//! serially — one network round-trip per commit, each capable of taking up to
//! the 30s JIRA timeout — producing multi-minute "zero progress" stalls on
//! large corpora (~181k commits). This module fans the network work out with a
//! bounded `buffer_unordered` (mirroring the LLM fallback tier).
//!
//! Fetch-once-per-TICKET (not per-message) is achieved in two phases, per the
//! code-critic HIGH finding on PR #2720: deduping on raw commit *message* text
//! is not sufficient, because the resolver's cache is keyed by the extracted
//! ticket KEY — two differently-worded commits referencing the same ticket
//! produce two distinct messages, and resolving both concurrently races the
//! cache's non-atomic check→fetch→populate sequence. The fix:
//! 1. **Warm** ([`ExternalSourceResolver::warm_cache`]) — extract and dedupe
//! ticket keys/refs across the *entire* batch of unique messages, per
//! source, BEFORE any fetch is dispatched; fetch only the still-uncached
//! keys, bounded-concurrent. This guarantees each unique ticket is fetched
//! at most once, regardless of how many distinct messages reference it.
//! 2. **Apply** ([`resolve_unique`] over [`ExternalSourceResolver::resolve`]) —
//! now that every referenced ticket is cached, resolving each unique
//! message is cache-only (no further network I/O, no further race).
//!
//! What: [`resolve_external_signals`] collects the *unique* commit messages
//! that still need external resolution, warms the resolver's cache from that
//! full message set, then resolves each message (now cache-hit-only, so still
//! fanned out via [`resolve_unique`] for progress-reporting granularity), and
//! returns a `message -> ExternalSignal` map the caller applies to every
//! commit sharing that message.
//!
//! Test: `pipeline_external_tests` (dedupe + concurrency) and the wiremock
//! `pipeline_tests::pipeline_external_sources_*` integration tests — including
//! `pipeline_external_sources_dedupe_distinct_messages_same_ticket`, which
//! proves fetch-once-per-ticket across DISTINCT messages (not just identical
//! ones).
use ;
use Future;
use ;
use StreamExt;
use info;
use CommitRow;
use ;
use ClassificationResult;
/// Maximum number of external-source lookups in flight at once.
///
/// Why: external APIs (JIRA, GitHub) enforce rate limits, so unbounded fan-out
/// would trade a serial stall for a burst of throttled/failed requests. A modest
/// bound of 8 mirrors the LLM fallback tier's default concurrency and keeps the
/// HTTP budget proportional while collapsing the multi-minute serial stall.
/// What: the `buffer_unordered` bound used by [`resolve_external_signals`].
/// Test: exercised indirectly by the pipeline integration tests.
pub const EXTERNAL_SOURCE_CONCURRENCY: usize = 8;
/// Emit an `info!` progress line every this many resolved unique messages.
///
/// Why: in a non-TTY context (the Duetto cron ETL) the indicatif progress bar is
/// hidden, so an explicit periodic log line to stderr is what makes a slow source
/// observable rather than a silent hang.
/// What: the modulus cadence for the progress `info!` in [`resolve_external_signals`].
/// Test: covered indirectly; log output is side-effect-only.
const PROGRESS_LOG_EVERY: usize = 50;
/// Resolve external-source signals for every commit lacking a Tier-0 override.
///
/// Why: this is the pipeline's Tier-0.5 entry point; concentrating the
/// dedupe-then-concurrent-resolve logic here keeps `pipeline.rs` thin and under
/// its SLOC budget.
/// What: collects the unique messages of override-free commits, **warms the
/// resolver's cache from the full ticket-key set extracted across all of
/// them** (see the module doc — this is what makes fetch-once-per-ticket hold
/// even when the same ticket is referenced by differently-worded commits),
/// then resolves each unique message (now cache-hit-only) with bounded
/// concurrency via [`resolve_unique`], emitting progress to stderr, and
/// returns a `message -> ExternalSignal` map. Commits with a Tier-0 manual
/// override are excluded (overrides win). Returns an empty map when nothing
/// needs resolution.
/// Test: `pipeline_tests::pipeline_external_sources_dedupe_and_apply` and
/// `pipeline_external_sources_dedupe_distinct_messages_same_ticket`.
pub async
/// Collect the unique messages of commits that still need external resolution.
///
/// Why: dedupe must happen before fan-out so each referenced ticket is fetched
/// exactly once even under concurrency; this is the pure, HTTP-free core of that
/// step, kept separate so it can be unit-tested directly.
/// What: returns each distinct `message` (first-occurrence order) among commits
/// that do NOT carry a Tier-0 manual override. Override commits are excluded
/// because the override verdict wins regardless of any external signal.
/// Test: `pipeline_external_tests::unique_messages_dedupes_and_excludes_overrides`.
pub
/// Resolve a set of unique messages concurrently, deduped, HTTP-free core.
///
/// Why: isolating the dedupe + bounded-concurrency + collect logic behind a
/// `resolve` closure gives a test seam that counts invocations without any
/// network I/O, proving each unique message is resolved exactly once.
/// What: maps each unique message through `resolve` with at most `concurrency`
/// futures in flight (`buffer_unordered`), then collects the non-`None` results
/// into a `message -> ExternalSignal` map. Input messages must already be unique;
/// callers dedupe upstream.
/// Test: `pipeline_external_tests::resolve_unique_invokes_once_per_message` and
/// `resolve_unique_bounds_concurrency`.
pub async