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
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
//! Static configuration for the native `EmbeddingService`: the entity message
//! name, the versioned outbox/work topics, the durable status tokens, the
//! top-k/quota bounds, the leader-pass batch knobs, the work-emitter cadence, and
//! the operator-tunable `Retrieve` score-threshold / fusion-weight knobs.
//! Extracted verbatim from the former god file; every value is byte-stable for
//! downstream audit/CDC consumers.
use OnceLock;
use Duration;
pub const EMBEDDING_SOURCE_MSG: &str = "udb.core.embedding.entity.v1.EmbeddingSource";
pub const EMBEDDING_MODEL_MSG: &str = "udb.core.embedding.entity.v1.EmbeddingModel";
pub const EMBEDDING_JOB_MSG: &str = "udb.core.embedding.entity.v1.EmbeddingJob";
pub const EMBEDDING_WORK_ITEM_MSG: &str = "udb.core.embedding.entity.v1.EmbeddingWorkItem";
pub const EMBEDDING_DOCUMENT_MSG: &str = "udb.core.embedding.entity.v1.EmbeddingDocument";
/// The change-driven embedding WORK topic the sidecar pool consumes. Its payload
/// carries ONLY the row pk + text + non-secret routing — never any credential.
pub const TOPIC_WORK: &str = "udb.embedding.work.v1";
pub const TOPIC_SOURCE_REGISTERED: &str = "udb.embedding.source.registered.v1";
pub const TOPIC_SOURCE_DELETED: &str = "udb.embedding.source.deleted.v1";
pub const TOPIC_BACKFILL_REQUESTED: &str = "udb.embedding.backfill.requested.v1";
pub const TOPIC_BACKFILL_COMPLETED: &str = "udb.embedding.backfill.completed.v1";
pub const TOPIC_SOURCE_CHANGE_COMPLETED: &str = "udb.embedding.source.change.completed.v1";
pub const TOPIC_MODEL_REGISTERED: &str = "udb.embedding.model.registered.v1";
pub const TOPIC_MODEL_STATUS_CHANGED: &str = "udb.embedding.model.status.changed.v1";
pub const TOPIC_MODEL_ALIAS_CUTOVER: &str = "udb.embedding.model.alias.cutover.v1";
pub const TOPIC_DOCUMENT_PARSE: &str = "udb.embedding.document.parse.v1";
pub const TOPIC_DOCUMENT_INGESTED: &str = "udb.embedding.document.ingested.v1";
pub const TOPIC_WORK_DEAD_LETTER: &str = "udb.embedding.work.dead.v1";
pub const TOPIC_METERED: &str = "udb.embedding.metered.v1";
pub const TOPIC_RETRIEVAL_SAMPLED: &str = "udb.embedding.retrieval.sampled.v1";
pub const TOPIC_RETRIEVAL_EVALUATED: &str = "udb.embedding.retrieval.evaluated.v1";
/// Completion marker for a deleted source's vector teardown, keyed by
/// `teardown_event_id` (the journal event id of the source-deleted event) so the
/// leader pass never re-runs a finished teardown (mirrors the backfill
/// requested/completed pairing).
pub const TOPIC_SOURCE_TEARDOWN_COMPLETED: &str =
"udb.embedding.source.teardown.completed.v1";
pub const STATUS_ACTIVE: &str = "ACTIVE";
pub const STATUS_DELETED: &str = "DELETED";
pub const STATUS_DEPRECATED: &str = "DEPRECATED";
pub const STATUS_RETIRED: &str = "RETIRED";
pub const WORK_PENDING: &str = "PENDING";
pub const WORK_ACKED: &str = "ACKED";
pub const WORK_DEAD: &str = "DEAD";
pub const JOB_PENDING: &str = "PENDING";
pub const EMBEDDING_WORK_EMITTER_BATCH: i64 = 200;
/// Page size for enumerating (and deleting) a deleted source's point ids during
/// vector teardown — bounds each journal scan and each vector-seam delete call.
pub const EMBEDDING_TEARDOWN_DELETE_BATCH: i64 = 200;
pub const EMBEDDING_BACKFILL_PAGE_LIMIT: i32 = 200;
const DEFAULT_EMBEDDING_WORK_EMITTER_INTERVAL_SECS: u64 = 30;
const EMBEDDING_WORK_EMITTER_INTERVAL_ENV: &str = "UDB_EMBEDDING_WORK_EMITTER_INTERVAL_SECS";
/// Minimum similarity score a vector hit must clear to be returned by `Retrieve`.
/// Operator-tunable (was a hardcoded `0.0`); resolved once via a `OnceLock`. `0.0`
/// keeps the historical "return everything the engine ranks" behavior by default.
/// A `RetrieveRequest.score_threshold` proto field (Part B.2.3) will later let a
/// caller raise this per query; until then this is the server-side floor.
const DEFAULT_EMBEDDING_RETRIEVE_SCORE_THRESHOLD: f32 = 0.0;
const EMBEDDING_RETRIEVE_SCORE_THRESHOLD_ENV: &str = "UDB_EMBEDDING_RETRIEVE_SCORE_THRESHOLD";
/// Comma-separated `lexical,vector` weights for hybrid-search fusion, e.g.
/// `"0.4,0.6"`. Empty (the default) preserves the delegated engine's built-in
/// fusion weighting (was a hardcoded empty vec).
const EMBEDDING_RETRIEVE_FUSION_WEIGHTS_ENV: &str = "UDB_EMBEDDING_RETRIEVE_FUSION_WEIGHTS";
/// Maximum characters of source text carried in a single `udb.embedding.work.v1`
/// event. Embedding models cap their input (roughly a few thousand tokens); an
/// unbounded row would make the sidecar's provider call fail, and — since a
/// failed embedding is never reported back — the row would silently stay
/// un-embedded. Bounding the text here keeps the request within a safe envelope.
/// Operator-tunable; a `<= 0`/malformed override falls back to the default.
const DEFAULT_EMBEDDING_MAX_TEXT_CHARS: usize = 8000;
const EMBEDDING_MAX_TEXT_CHARS_ENV: &str = "UDB_EMBEDDING_MAX_TEXT_CHARS";
/// Chunking knobs (Part B.2.2). A source row's text is split into overlapping
/// `chunk_size`-character windows (word-boundary-aware), `chunk_overlap`
/// characters shared between neighbors, capped at `max_chunks_per_row` points per
/// row. Defaults suit a general RAG corpus; all operator-tunable, resolved once.
/// A row whose text fits in one window keeps its bare `row_pk` point id (no
/// behavior change for short rows).
const DEFAULT_EMBEDDING_CHUNK_SIZE: usize = 1000;
const EMBEDDING_CHUNK_SIZE_ENV: &str = "UDB_EMBEDDING_CHUNK_SIZE";
const DEFAULT_EMBEDDING_CHUNK_OVERLAP: usize = 150;
const EMBEDDING_CHUNK_OVERLAP_ENV: &str = "UDB_EMBEDDING_CHUNK_OVERLAP";
const DEFAULT_EMBEDDING_MAX_CHUNKS_PER_ROW: usize = 256;
const EMBEDDING_MAX_CHUNKS_PER_ROW_ENV: &str = "UDB_EMBEDDING_MAX_CHUNKS_PER_ROW";
/// Fallback vector collection when a source row somehow carries no target (mirrors
/// `asset_service::DEFAULT_VECTOR_COLLECTION`; a source normally always specifies
/// its own `target_collection`, validated non-empty at register time).
pub const DEFAULT_VECTOR_COLLECTION: &str = "udb_asset_embeddings";
/// Default number of hits returned when the caller does not specify `top_k`.
const DEFAULT_TOP_K: i32 = 10;
/// Upper bound on `top_k` so one query cannot pull an unbounded result set.
const MAX_TOP_K: i32 = 200;
/// Per-tenant registered-source budget. Bounds the durable table so one tenant
/// cannot exhaust the shared store; a new source beyond this fails closed.
pub const MAX_SOURCES_PER_TENANT: usize = 128;
pub const MAX_MODELS_PER_TENANT: usize = 128;
pub const MAX_EMBEDDING_REPORT_BATCH: usize = 256;
pub const MAX_DOCUMENT_INGEST_BATCH: usize = 100;
pub const DEFAULT_WORK_MAX_ATTEMPTS: i32 = 5;
const DEFAULT_WORK_VISIBILITY_TIMEOUT_SECS: u64 = 120;
const WORK_VISIBILITY_TIMEOUT_ENV: &str = "UDB_EMBEDDING_WORK_VISIBILITY_TIMEOUT_SECS";
const DEFAULT_RETRY_SWEEP_LIMIT: i64 = 200;
const RETRY_SWEEP_LIMIT_ENV: &str = "UDB_EMBEDDING_RETRY_SWEEP_LIMIT";
const RETRIEVAL_EVAL_SAMPLE_RATE_ENV: &str = "UDB_EMBEDDING_EVAL_SAMPLE_RATE";
const FRESH_BUFFER_TTL_ENV: &str = "UDB_EMBEDDING_FRESH_BUFFER_TTL_SECONDS";
const FRESH_BUFFER_CAPACITY_ENV: &str = "UDB_EMBEDDING_FRESH_BUFFER_CAPACITY";
const DEFAULT_FRESH_BUFFER_TTL_SECS: u64 = 30;
const DEFAULT_FRESH_BUFFER_CAPACITY: usize = 10_000;
/// Clamp a requested `top_k` into `[1, MAX_TOP_K]`; non-positive → default.
pub
pub
/// Server-side minimum-score floor applied to every mediated `Retrieve`. Resolved
/// once (no per-request env read); a non-finite or negative override is ignored so
/// the floor is always a real, sane bound.
pub
/// Hybrid-search fusion weights, parsed once from a comma-separated list. Empty
/// (the default, and the fallback for any malformed/negative entry) hands fusion
/// weighting back to the delegated engine — identical to the previous hardcoded
/// `Vec::new()`.
pub
/// Maximum source-text characters per work event, resolved once. A `<= 0` or
/// malformed override is ignored so the bound is always a real, sane limit.
pub
/// Chunk window size in characters (resolved once). Bounded to
/// `max_embedding_text_chars` so a single chunk can never exceed the work-event
/// text envelope.
pub
/// Characters shared between neighboring chunks (resolved once). The chunker
/// clamps this below the window size to guarantee forward progress.
pub
/// Safety cap on chunks emitted per source row (resolved once), so a
/// pathologically large row cannot fan out an unbounded number of points.
pub
pub
pub
pub
pub
pub