GitHub
09/25/2026, 6:32 PM<https://github.com/linkedin/venice/tree/main|main> by minhmo1620
<https://github.com/linkedin/venice/commit/6c40bf824f0c562a049fe21d32609bd988575075|6c40bf82> - [vpj] Encrypt push job writes for encryption-enabled stores (#3037)
linkedin/veniceGitHub
09/25/2026, 7:12 PMVenice-BlobTransfer-Snapshot-Cleanup-Scheduler
• Venice-BlobTransfer-Replica-Blob-Fetch-Executor
This made blob-transfer thread-priority behavior inconsistent across the blob-transfer stack.
Solution
Centralized blob-transfer thread priority in `BlobTransferUtils`:
• Added BlobTransferUtils.BLOB_TRANSFER_THREAD_PRIORITY = Thread.NORM_PRIORITY - 1
• Updated all blob-transfer thread factories to use the shared constant:
• Netty client event-loop threads
• client host-connect executor
• client timeout checker
• checksum validation executor
• server boss/worker Netty event loops
• snapshot cleanup scheduler
• replica blob fetch executor
• Removed duplicated local priority constants.
• Updated the existing blob-transfer client thread-priority test to assert against the shared constant.
This keeps existing behavior for already-tuned blob-transfer threads and tunes the two remaining blob-transfer executors down to the same priority.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
No synchronization, locking, queueing, or task execution logic was changed. This PR only changes thread priority configuration at executor/thread-factory creation time.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
Validated with:
./gradlew spotlessCheck clientsda-vinci-client:test --tests com.linkedin.davinci.blobtransfer.client.TestNettyFileTransferClient
linkedin/veniceGitHub
09/25/2026, 7:46 PM<https://github.com/linkedin/venice/tree/main|main> by eldernewborn
<https://github.com/linkedin/venice/commit/b1615e44b38ae3d6a3ee490cf159405942ae3f40|b1615e44> - Tune blob transfer thread priority (#3038)
linkedin/veniceGitHub
09/25/2026, 9:06 PMremove(t) used the wrong key type.
• The countdown map was keyed by topic-name strings, but storeToCountdownForDeletion.remove(t) passed a topic object. The removal did nothing.
• Countdown entries remained even after topics disappeared. As a side effect, a topic that finished its safety wait stayed eligible on every later scan.
Bug 2: the countdown was cleared before deletion, rather than after the topic was gone.
• The cleaner removed the countdown as soon as a topic became eligible and entered the deletion queue—not after it was deleted.
• A queue rebuild after the 60-second refresh interval could drop a topic still waiting its turn. A silent pre-delete check or failed delete call could also leave an eligible topic undeleted.
• With the countdown already removed, the next scan restarted the full safety wait instead of retrying deletion. Repeated missed opportunities could therefore delay store hard deletion for hours.
#2896 fixed Bug 1 by changing the call to remove(countdownKey), which made Bug 2 take effect. The regression is the now-effective removal on eligibility, not the per-fabric split itself. This PR retains eligibility until a scan observes the topic's absence, addressing Bug 2 without restoring Bug 1's stale-entry leak.
The earlier investigation also found the fabric-level DeletableCount rising from roughly 5–25 before September 10 to hundreds on September 10 and thousands by September 24–25. This supports increased cleanup pressure; aggregate metrics alone do not establish deployment timing or production impact.
Solution
Scans own countdown state; deletion calls do not.
1. Keep each per-broker countdown at zero once its delay has elapsed. A topic remains delay-eligible across queue refreshes, failures, and silent skips. All existing retention/resource safety filters still apply.
2. Use the existing successful topic-retention listing to discard countdowns for topics absent from that broker. No additional existence RPC is needed. A same-name topic appearing after an observed absence gets a fresh countdown.
3. Attempt each topic at most once per cleanup pass, excluding already-attempted topics when rebuilding the queue. This prevents failing RT/VT topics from being retried on every refresh while other topics wait. RT priority and RT deletion safety checks remain in place.
The default delay remains 20, with eligibility on the 21st qualifying scan (initialization plus 20 decrements). Zero-delay listing queries do not mutate countdown state.
This simplified version leaves deleteTopic, deletion metrics, and existing logs unchanged. It removes the earlier revision's extra existence check, immediate countdown clearing, per-topic eligibility logs, and countdown-state test accessor. Countdowns are grouped by broker so pruning one broker cannot reset another broker's state.
Code changes
• Added new code behind a config. No new config; existing delay and refresh settings are unchanged.
• Introduced new log lines. None.
• Confirmed if logs need to be rate limited to avoid excessive logging. N/A: no new logs.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues in the existing single-cleanup-worker usage. Zero-delay API queries bypass mutable countdown state.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed. No new synchronization; cleanup-thread ownership is retained.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation. No new locks or broker calls.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList). Existing HashMap ownership is retained, and the attempted-topic set is local to a pass; no new concurrent collection is required.
• Validated proper exception handling in multi-threaded code to avoid silent thread termination. Existing exception handling is unchanged; a failed listing cannot reach pruning.
How was this PR tested?
Testing Done
./gradlew --no-daemon spotlessApply spotlessCheck :services:venice-controller:test --tests 'com.linkedin.venice.controller.kafka.TestTopicCleanupService' -q
17 tests passed; 0 failures, errors, or skips. Assertions exercise cleanup behavior rather than inspecting the countdown map:
• Failure, a simulated silent deletion-precheck return, and eventual successful deletion with missing store metadata; a recreated topic waits again after observed absence.
• Pruning a countdown before its delay elapses; a reappearing topic receives the full delay.
• Multiple eligible VTs and a failing RT topic across repeated queue refreshes, with each topic attempted once per pass. The test forces the refresh interval rather than sleeping for 60 seconds.
• Independent countdowns and pruning for the same topic name in two brokers; a zero-delay query does not advance or reset the cleaner's countdown.
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests. Existing pre-PR tests are retained; four regression tests are added.
• Verified backward compatibility (if applicable). Existing tests pass; initial delay, RT priority, safety filters, and deletion-call behavior are preserved.
Risk, rollout, and follow-ups
• Pruning relies on a successful complete topic listing. Deletion and recreation entirely between scans cannot be distinguished by topic name alone.
• Rollout should monitor deletable-topic backlog and store hard-deletion completion. Rolling back reintroduces the reset-on-eligibility regression.
• Deferred: distinguishing a silent precheck return from a successful deletion in metrics, plus queue/scan/eligible-age instrumentation. A silent return still leaves the topic eligible here as long as it remains listed.
• Production impact has not been established by this investigation.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section. No API/configuration changes; this fixes cleanup retries.
• Yes. Clearly explain the behavior change and its impact.
linkedin/veniceGitHub
09/25/2026, 10:01 PMVenice Publication Pipeline (build-and-upload-archives-on-schedule.yml) has been failing at the make-tag step because there is no new commit on main since the last tag (0.5.7):
Not pushing new tag as there is no new commit since the last tag 0.5.7
##[error]Process completed with exit code 1.
Because build-and-publish depends (needs:) on make-tag, this also blocks the JFrog publish step for the existing tag, which in turn blocks downstream ELR resolution for that artifact.
Solution
Empty commit with no file changes, so the next Publication Pipeline run has a new commit to tag and can proceed through make-tag \u2192 build-and-publish successfully. No functional code changes.
How was this PR tested?
N/A \u2014 empty commit, no code changes.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
linkedin/veniceGitHub
09/25/2026, 10:08 PM<https://github.com/linkedin/venice/tree/main|main> by KaiSernLim
<https://github.com/linkedin/venice/commit/79aac66cc5b967957c4eef29466b720637944336|79aac66c> - [controller] Preserve Topic Cleanup Eligibility Across Cleanup Scans (#3039)
linkedin/veniceGitHub
09/25/2026, 11:51 PMPubSubProducerAdapterContext and propagate it through VeniceWriterFactory. KafkaStoreIngestionService now keeps the existing factory as the default with encryption disabled and creates a second encrypted factory. Ingestion tasks and view writers select the encrypted factory only when Store.isEncryptionEnabled() is true, while metadata and push-status writers continue using the default factory.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
No new external configuration or log lines are introduced.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
The change only selects immutable, shared writer factories during ingestion-task construction and adds no synchronization, blocking calls, collections, or exception paths.
How was this PR tested?
• Local code review completed
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
• ./gradlew :internal:venice-common:test --tests com.linkedin.venice.pubsub.api.PubSubProducerAdapterContextTest --tests com.linkedin.venice.writer.VeniceWriterFactoryTest
• ./gradlew :clients:da-vinci-client:test --tests com.linkedin.davinci.kafka.consumer.StoreIngestionTaskFactoryTest --tests com.linkedin.davinci.kafka.consumer.KafkaStoreIngestionServiceTest --tests com.linkedin.davinci.kafka.consumer.LeaderFollowerStoreIngestionTaskTest
• Repository pre-commit Spotless and generated-GHCI checks passed.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
🤖 Generated with GitHub Copilot CLI
linkedin/veniceGitHub
09/26/2026, 12:21 AM<https://github.com/linkedin/venice/tree/main|main> by minhmo1620
<https://github.com/linkedin/venice/commit/887ec5691159b028d3ad81d60fcbdb9ed8072c21|887ec569> - Select writer factory by store encryption (#3041)
linkedin/veniceGitHub
09/27/2026, 1:52 AMupdate_store that did not intend to change lifecycle hooks could silently wipe a store's configured hooks.
store_lifecycle_hooks_list has no "unset" representation on the wire: the UpdateStore Avro field is a non-nullable array defaulting to [], so a present-but-empty list is indistinguishable from "the caller never touched this field". StoreLifecycleHooksPolicy.validateLifecycleHooks only preserved the current hooks when the value was *absent*; a present-but-empty list was returned as-is and applied as a clear.
This bit the direct child-controller REST path (StoreConfigUpdater.applyOnChild), which has no updatedConfigsList gate — an update that only flipped storage_mode also carried the empty-default hooks list and cleared the hooks. The parent → Kafka → child path was protected only because empty hooks are never added to updatedConfigsList and get filtered out by AdminExecutionTask, a fragile, path-dependent safety net (and replicateAllConfigs=true bypasses it entirely).
Fix
Treat a present-but-empty list the same as an absent value (no change) in the shared validateLifecycleHooks, which is called by both applyOnChild and applyOnParent. Every apply path now preserves the current hooks, matching how every other field falls back to its current value when unspecified. Clearing hooks via an empty list is intentionally unsupported; a genuine clear would require a distinct sentinel.
Testing Done
• `StoreLifecycleHooksPolicyTest`: present-empty preserves existing hooks; absent preserves; present-empty with no existing hooks stays empty (plus existing blank-class-name/trim cases).
• `StoreConfigUpdaterTest.testApplyOnChild_EmptyLifecycleHooks_DoesNotWipeExistingHooks`: reproduces the failure (a storage_mode change carrying an empty hooks list on a store that already has a hook) and asserts the existing hook is preserved, not cleared.
./gradlew :services:venice-controller:test --tests "com.linkedin.venice.controller.storeconfig.StoreLifecycleHooksPolicyTest" --tests "com.linkedin.venice.controller.StoreConfigUpdaterTest" — BUILD SUCCESSFUL.
linkedin/veniceGitHub
09/28/2026, 8:13 PMproducerEncryptionEnabled = true whenever a pubSubEncryptionKeyUrnLookup function was supplied, regardless of whether the target store is actually encryption-enabled. Callers that always construct a non-null lookup function (e.g. `VenicePushJob`'s data writer) trip this path unconditionally, so producer-encryption capability was effectively always on, not scoped to encrypted stores.
• Separately, the controller (VeniceHelixAdmin) built this same lookup by checking whether any cluster it hosts has cluster.encryption.enabled=true, then attached it to its single shared VeniceWriterFactory — used for writes that aren't scoped to one store (the admin topic, which carries UpdateStore messages for every store in a cluster, has no store name in the topic name). A controller hosting even one encrypted cluster would mark its shared writer factory encryption-capable for every write it does, including for non-encrypted clusters/stores.
• The controller never writes actual versioned store data, so it never legitimately needs producer-side encryption at all.
Solution
• `VeniceWriterFactory`: derive producerEncryptionEnabled from pubSubEncryptionKeyUrnLookup != null instead of hardcoding true, so the flag reflects whether a lookup was actually supplied.
• VeniceHelixAdmin (controller): remove the encryption write path entirely — delete createPubSubEncryptionKeyUrnLookup() and its now-unused private helper, and stop wiring any lookup into TopicManagerContext / the shared VeniceWriterFactory.
• KafkaStoreIngestionService (server/DVC): stop passing the encryption key URN lookup to the plain (non-encrypted) `veniceWriterFactory`; only encryptedVeniceWriterFactory needs it. StoreIngestionTaskFactory.getVeniceWriterFactory(store) already dispatches per-store via store.isEncryptionEnabled() between the two factories, so the plain factory never legitimately needs lookup capability.
• Net effect: producer-encryption capability is now scoped to the actual per-store encrypted writer path (server/DVC ingestion, driven by metadataRepo::getPubSubEncryptionKeyUrn), not leaked to every writer that happens to receive a non-null lookup function or to the controller's non-store-scoped shared writer.
• Performance/trade-offs: none. This is a boolean-derivation change plus dead-code removal — no new allocations, locking, or I/O.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used where needed (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
N/A — this change is a boolean-derivation fix plus removal of dead code; no new threading, locking, or blocking I/O is introduced.
How was this PR tested?
• New unit tests added. — VeniceWriterFactoryTest#testProducerEncryptionFlagDerivedFromLookup.
• New integration tests added. — TestPushJobWithEncryptionKeyUrn, exercising a push job through the admin-topic UpdateStore flow with `cluster.encryption.enabled=true`: fails fast without a key URN, succeeds once one is set.
• Modified or extended existing tests. — Removed TestVeniceHelixAdminWithoutCluster#testEncryptionKeyLookupIsOnlyCreatedForEncryptionClusters (the only test of the deleted controller method).
• Verified backward compatibility (if applicable). — Confirmed no in-repo production code currently reads PubSubProducerAdapterContext#isProducerEncryptionEnabled() in this snapshot, and the removed controller method (createPubSubEncryptionKeyUrnLookup) was package-private with no external callers.
Additional local verification:
• ./gradlew :internal:venice-common:test --tests VeniceWriterFactoryTest — pass
• ./gradlew :internal:venice-test-common:integrationTest --tests TestPushJobWithEncryptionKeyUrn — pass
• ./gradlew :services:venice-controller:test --tests TestVeniceHelixAdminWithoutCluster — pass (29 tests)
• Clean compile of all touched modules (main + test source sets)
• Regression check on the real VenicePushJob write path (driver + MR + Spark data-writer jobs): confirmed pubSubEncryptionKeyUrn is only ever propagated for encryption-enabled stores (and explicitly cleared otherwise) before reaching VeniceWriterFactory. Ran VenicePushJobTest, TestDataWriterMRJob, PubSubEncryptionUtilsTest, AbstractDataWriterSparkJobTest, AbstractPartitionWriterTest — 59/59 pass, 0 failures, confirming no unintended change to push-job encryption behavior.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
No in-repo consumer currently reads isProducerEncryptionEnabled(), so there is no observable behavior change within this repository today. The fix scopes the flag correctly for any current or future PubSubProducerAdapterFactory implementation that does honor it, so it sees encryption capability tied to the store/writer that actually has a key URN, instead of true leaking to every writer that happens to receive a non-null lookup function (or to the controller's shared, non-store-scoped writer).
linkedin/veniceGitHub
09/29/2026, 12:03 AM<https://github.com/linkedin/venice/tree/main|main> by minhmo1620
<https://github.com/linkedin/venice/commit/948f8a36fd44765ccd9d0f08b1904a168bba5b5f|948f8a36> - [producer][controller][da-vinci] Scope producer encryption to actual encrypted stores/writers (#3043)
linkedin/veniceGitHub
09/29/2026, 12:59 AMKafkaConsumerService creates them with RandomAccessDaemonThreadFactory without an explicit priority.
These consumer threads are part of Da Vinci ingestion background processing. They should run below normal priority so they yield CPU to more latency-sensitive read/write paths under contention, matching the existing pattern used by other ingestion/background thread pools.
Solution
This PR sets shared consumer thread priority to Thread.NORM_PRIORITY - 1.
Changes:
• Added a priority-aware constructor to RandomAccessDaemonThreadFactory.
• Updated KafkaConsumerService to create shared consumer threads with explicit below-normal priority.
• Added test coverage to verify RandomAccessDaemonThreadFactory applies configured priority while preserving indexed thread lookup.
Please note that this is just priority annotation, and has no effect unless the container supports it directly, and the thread policy is explicitly passed.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
This change only sets thread priority at thread creation time. It does not change queueing, locking, polling, subscription, or exception-handling behavior.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
Validated with:
./gradlew spotlessCheck internalvenice-client-common:test --tests com.linkedin.venice.utils.DaemonThreadFactoryTest clientsda-vinci-client:test --tests com.linkedin.davinci.kafka.consumer.KafkaConsumerServiceTest
linkedin/veniceGitHub
09/29/2026, 2:45 AM<https://github.com/linkedin/venice/tree/main|main> by eldernewborn
<https://github.com/linkedin/venice/commit/127db552cbbbedde48a439ec28f267c18454af89|127db552> - [da-vinci] Tune shared consumer thread priority *(#3044)
linkedin/veniceGitHub
09/29/2026, 6:59 PMwriteQuotaEnabled as a boolean with default false to:
• StoreProperties in StoreMetaValue v50.
• StoreCreation and UpdateStore in AdminOperation v105.
The default lets new readers consume older records without enabling write quota enforcement. Keep code generation pinned to StoreMetaValue v49 and AdminOperation v104, following the existing protocol-staging convention. This PR adds no runtime behavior; implementation and protocol activation are separate follow-ups.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
Not applicable: this change adds schemas and build-time version pins only.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
Not applicable: no executable concurrency code changes.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
Validation on JDK 11:
• ./gradlew :internal:venice-common:compileAvro :services:venice-controller:compileAvro --console=plain passed. Generated classes remain on v49/v104.
• Direct Avro 1.10.2 validation parsed all 155 schemas and passed backward-compatibility checks against all 153 historical schemas, plus forward-compatibility checks against both immediate predecessors.
• Fifteen binary checks across the three affected records verified missing-field defaults, `true`/`false` roundtrips, and predecessor readers ignoring the new field while preserving existing fields.
• JSON comparison confirmed that the new schemas preserve every predecessor field and add only the requested field. git diff --check passed.
No test files were added. The full suite was not run. Gradle reported a deprecated-feature warning.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
linkedin/veniceGitHub
09/29/2026, 7:58 PM0 delta every interval, and gauges keep reporting a stale value.
• Placeholders instead of no data. For example 0, -1 or Long.MAX_VALUE before the first update or while a value is unknown, NaN mapped to 0, and Infinity when a store's quota is smaller than its partition count.
• Samples in the wrong population. Heartbeat and record delay recorded a 0 under the readiness state a follower was not in, which lowers percentiles of the real samples, and always counted leaders as ready to serve. Every parent controller reported system store health for every cluster, including clusters it didn't lead.
Solution
Metrics framework (venice-client-common)
• MetricScope groups a component's observable metric states and closes them together. A high-performance counter reports only while a scope owns it; after closing, it reports its final totals through the next collection (for one export interval when there may be several metric readers), so every reader exports the last interval.
• Async gauge factories require a MetricScope, a state resolver and a value resolver. A sample is emitted only when the state is non-null and the value is finite (GaugeObservation); callback exceptions are recorded as metric failures. The supplier-based create(...) factories and `registerObservableLongGauge`/`registerObservableDoubleGauge` are removed, registerObservableGauge requires a MetricScope that closes the gauge, and AsyncMetricApiContractTest checks that code creating an async gauge without a scope and resolvers doesn't compile.
• AbstractVeniceStats adds getMetricScope(), closeOtelMetrics() and `close()`; AbstractVeniceAggStats.removeStore() closes one store's stats.
• AbstractVeniceAggVersionedStats keeps one StoreOtelStats per store: created on first use, updated with the store's current and future versions, and closed on store deletion. Ingestion, DIV, storage engine, blob transfer and heartbeat stats use it.
Behavior changes
| Metrics | Before (OTel) | After (OTel) | Tehuti |
| ----------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------- |
| Per-store ingestion (ingestion.) and DIV (ingestion.div.) | Kept exporting (counters: 0 every interval) after store deletion or after the store's last ingestion task on the host stopped | Closed when the store's last local ingestion task stops, or on store deletion | Unchanged |
| Active / unique ingested key count | -1 / 0 for a replica type with no partition on the host | Omitted | Unchanged |
| Storage engine disk usage, key count estimate | Kept reporting stale or 0 values after a version's storage engine left the host | That version stops reporting; the store's stats close when no engine remains | Unchanged |
| Heartbeat and record delay | Follower delay under its readiness plus 0 under the other; leaders always READY_TO_SERVE; version role could go stale while the store had no replica on the host | Recorded once, under the replica's actual readiness (leaders included); version info kept current | Unchanged (the inactive follower sensor still gets 0) |
| Read quota usage ratio | 0 when there is no local version or quota share | Omitted; the store's quota stats close on store deletion | Store deletion now unregisters its read-quota sensors |
| Disk quota used | Infinity when the store quota is smaller than its partition count | Usage over this host's exact share of the quota (finite for a positive quota; a zero quota has no ratio and is omitted) | storage_quota_used is finite instead of Infinity for those stores with a positive quota; other stores change negligibly |
| RocksDB stats, RMD block cache …
linkedin/veniceGitHub
09/29/2026, 8:49 PM<https://github.com/linkedin/venice/tree/main|main> by misyel
<https://github.com/linkedin/venice/commit/47f7c405f60b5b80e632f27a9ae503911ebe65ef|47f7c405> - [protocol][controller] Stage write quota enabled schemas (#3045)
linkedin/veniceGitHub
09/29/2026, 9:28 PMtrue or false. The schema must land separately from the Java implementation.
Solution
Add one field to `CreateStoreGrpcRequest`: optional bool writeQuotaEnabled = 7;.
This PR contains only that protobuf schema line. It changes no Java code, tests, Avro schemas, or runtime behavior. Field 7 is unused in this message, and all existing field numbers stay unchanged. Proto3 optional presence distinguishes an unset field from explicit `false`; the generated getter returns false when unset. The new-user-store default of true belongs to the implementation in #3047, not this schema change.
Merge this prerequisite first. #3047's mixed-schema check is expected to remain blocked until its base includes this field and the implementation branch is updated against that base.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
No concurrency changes are included. The joint reviewer/author checks above remain open.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
Validation on this schema-only branch:
• ./gradlew :internal:venice-common:generateProto :internal:venice-common:compileJava :internal:venice-common:test --tests com.linkedin.venice.controllerapi.ControllerClientTest passed on JDK 17. The existing test class ran 4 tests with zero failures, errors, or skips.
• git diff --check passed. Normal pre-commit hooks (spotlessApply and generateGHCI) passed and left the change at one file, one added line.
• No test source changes are included; behavior tests remain in #3047. The full repository test suite was not rerun. CI is pending. Gradle reported its existing deprecated-feature warning.
• Compatibility inspection: this is an additive optional field with a new field number. Older readers ignore the field; older senders leave it unset. Mixed-version behavior still needs the implementation in #3047.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
Generated with GitHub Copilot CLI.
linkedin/veniceGitHub
09/29/2026, 10:55 PMwriteQuotaEnabled schemas. The setting still needs to be exposed in store metadata and carried through controller updates, migration, and recovery.
Solution
Wire the setting through the existing store configuration paths:
• Add accessors, cloning, read-only delegation, and StoreInfo serialization. Older JSON and Avro metadata retain the false default.
• Remove the schema-generation overrides so generation uses the latest schemas, and advance active protocol versions to StoreMetaValue v50 and AdminOperation v105.
• Expose write_quota_enabled on UPDATE_STORE and --write-quota-enabled true|false on the admin tool's --update-store command.
• Preserve the current value when an update omits the setting. Add it to updatedConfigsList only when explicitly requested, and carry it through parent/child admin-message consumption.
• Enable the setting for new user stores through the parent controller's existing post-create updateStore call, alongside storageNodeReadQuotaEnabled. Initial store metadata remains `false`; the following UPDATE_STORE sets it to true. System-store defaults remain unchanged.
• Preserve the source value during migration, metadata copying, and recovery through their existing follow-up update calls. A destination created through the parent may briefly have the setting enabled before the source value is restored.
Creation requests retain their original REST/gRPC fields and Java signatures. This PR has no .proto or .avsc changes and needs no additional schema PR. The historical StoreCreation.writeQuotaEnabled field from #3045 remains unchanged and unused by this implementation.
This PR adds configuration plumbing only; it does not implement write quota enforcement.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• `writeQuotaEnabled`: metadata default `false`; enabled for new user stores by the parent post-create update. Explicit updates accept either boolean value. No enforcement logic is introduced.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
The changes reuse the existing store creation locks and post-create metadata-update path. They introduce no new threads, collections, or synchronization mechanisms. Creation and its follow-up update are separate operations, as they already are for storageNodeReadQuotaEnabled. The joint reviewer/author checks remain open.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
• Independent read-only code review completed; no significant issues reported.
Validation for the post-create revision:
| JDK 17 focused unit-test scope | Tests passed |
| ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------ |
| TestZKStore, TestStoreJsonSerializer, UpdateStoreQueryParamsTest | 70 |
| TestVeniceParentHelixAdmin, StoreMigrationHelperTest, StoreConfigUpdaterTest, AdminExecutionTaskTest, AdminConsumptionTaskTest, TestVeniceHelixAdminStoreDeletion, AdminOperationSerializerTest, CreateStoreTest, StoreRequestHandlerTest | 327 |
| TestAdminTool, TestRecoverStoreMetadata | 25 |
| Total | 422 |
These were selected with Gradle `--tests`; unit-test XML reports show zero failures, errors, or skips. Coverage includes both boolean values, omitted fields, old JSON/Avro defaults, the unused v104/v105 creation field, selective and replicate-all updates, read-only snapshots, system-store defaults, and migration/recovery follow-up updates.
On JDK 11, 13 targeted integration cases passed:
• TestVeniceHelixAdminWithSharedEnvironment.testWriteQuotaEnabledDefaultsFalseUntilUpdate.
• VeniceParentHelixAdminTest.testParentCreationEnablesWriteQuotaThroughUpdate, including child-controller propagation and preservation across an unrelated update.
• Seven methods in `TestAdminOperationWithPreviousVersion`: store creation, pause, value schema, metadata schema, superset schema, migration, and update.
• All four methods in TestMultiDataCenterAdminOperations.
The first integration run had one multi-region second-push timeout, followed by two failures because that test had not reached its protocol reset. All three passed on automatic retry. A separate rerun of the four-test multi-region class then passed with zero failures, errors, or skips, without changing code, timeouts, or assertions. The rerun count is not additive.
All four affected SpotBugs tasks passed on JDK 17 with zero findings:
./gradlew --no-daemon \
internalvenice-common:spotbugsTest \
servicesvenice-controller:spotbugsTest \
clientsvenice-admin-tool:spotbugsTest \
internalvenice-test-common:spotbugsIntegrationTest \
-Pspotallbugs -Pspotbugs.reports.all=true
The legacy Avro test uses compatible record/default initialization instead of GenericRecordBuilder. git diff --check and the normal pre-commit hooks (spotlessApply and generateGHCI) passed; the committed tree matches the validated tree. The full repository suite was not run locally, and this PR remains a draft pending CI. Existing Gradle/JDK deprecation and Hadoop reflective-access warnings were observed.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
The store API now exposes writeQuotaEnabled, and updates accept write_quota_enabled. The admin tool accepts --write-quota-enabled true|false with --update-store. There is no new creation-request parameter.
New user stores created through the parent start with false metadata, then become true through the existing post-create UPDATE_STORE. Existing metadata without the field remains `false`; unrelated updates preserve its value. Local/child creation retains the false default until an update arrives. System-store defaults are unchanged.
A rejected post-create update can return a creation error after the store already exists with `writeQuotaEnabled=false`; there is no atomic rollback. Migration and recovery restore the source value through their subsequent update, including resetting it to false.
Active protocol versions advance to AdminOperation v105 and `Sto…
linkedin/veniceGitHub
09/29/2026, 11:26 PM<https://github.com/linkedin/venice/tree/main|main> by misyel
<https://github.com/linkedin/venice/commit/0a6f6fd422a411d1b662cb0bcab380d55d0b537c|0a6f6fd4> - [compat][controller] Wire write quota enabled store config (#3047)
linkedin/veniceGitHub
09/30/2026, 1:06 AMStoreMigrationHelper.validateEncryptionClusterMigration() used OR logic (srcEncryptionCluster || destEncryptionCluster), which blocked every store migration that touched an encryption cluster on either side — including migrations between two encryption clusters. That case was never intentionally blocked: the original design (introduced in the commit that added this check) allowed it, and no later change ever added an explicit encryption-to-encryption exception. This left an unverified gap with no unit test covering the (true, true) case, and it was only discovered when a real migration between two encryption clusters (encrypt-0 -> cert-encrypt-0) failed with an HTTP 400.
Solution
Switch the guard from OR to XOR (srcEncryptionCluster != destEncryptionCluster) so only migrations that cross the encryption boundary are rejected. Migrating between two encryption clusters (or two non-encryption clusters) is now allowed. Update the exception message to describe the new invariant ("migrating between an encryption cluster and a non-encryption cluster is not allowed").
Why this is safe
encrypt-0 and cert-encrypt-0 are both newly created clusters: CLUSTER_ENCRYPTION_ENABLED was set to true before either cluster held any stores, not toggled on afterward on a cluster that already had stores. A store's encryptionEnabled flag is set once, at store-creation time, from its cluster's encryption policy at that moment, and is never changed afterward. Because both clusters were encryption-enabled from creation, every store that has ever existed in either one is guaranteed to have encryptionEnabled=true with a pubSubEncryptionKeyUrn already configured — there is no "legacy" store that predates the cluster's encryption policy and could slip through this check without a key. This is why the fix only needs to compare the two clusters' encryption policies (the XOR check) and does not need to separately validate the source store's own encryption state.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues. (Stateless boolean guard; no shared mutable state.)
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed. (N/A — no shared state.)
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Validated thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList). (N/A — no collections involved.)
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
• Local code review completed
./gradlew :services:venice-controller:test --tests com.linkedin.venice.controller.StoreMigrationHelperTest — all 6 tests pass: the 4 encryption-policy tests (including the new testAllowsMigrationBetweenEncryptionClusters covering the previously-untested case) plus 2 pre-existing parameterized write-quota tests, unaffected by this change.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
Previously, --migrate-store requests between two clusters that are BOTH configured as encryption clusters were rejected with HTTP 400 ("migration from or to an encryption cluster is not allowed"), even though this case was intended to be supported per the original design. After this change, such migrations are allowed; only migrations that cross the encryption boundary (one side encrypted, the other not) remain blocked, with an updated error message reflecting this narrower restriction.
🤖 Generated with GitHub Copilot CLI
linkedin/veniceGitHub
09/30/2026, 1:29 AMpubSubEncryptionKeyUrn to the destination (StoreMigrationHelper#cloneDestinationStoreAndSyncConfigs → UpdateStoreQueryParams(StoreInfo, boolean)) is pre-existing and unmodified by that PR, and is unit-tested with a mocked ControllerClient (StoreMigrationHelperTest#testMigrationPropagatesPubSubEncryptionKeyUrn).
Gap
There is no integration-level test that exercises this propagation with real encrypted source and destination clusters, verifying:
• the destination's persisted pubSubEncryptionKeyUrn after the migration's updateStore call actually replays through the admin channel, and
• that version creation on the destination succeeds only once that key is persisted (mirroring the existing coverage in TestEncryptionClusterStoreConfig#testParentRequiresPubSubEncryptionKeyBeforeCreatingVersion, but through the migration-clone path specifically rather than a direct updateStore call).
Suggested work
Add a dedicated integration test (likely a new test class, since `TestStoreMigration`'s shared two-cluster @BeforeClass setup has no CLUSTER_ENCRYPTION_ENABLED and is shared by ~15 unrelated tests) that:
1. Stands up two encryption-enabled clusters (source + destination) under a parent controller.
2. Creates a store with a populated pubSubEncryptionKeyUrn in the source cluster.
3. Performs a real store migration (migrate-store / StoreMigrationManager) to the destination cluster.
4. Asserts the destination's persisted URN matches the source's, and that a version can be created on the destination afterward.
linkedin/veniceGitHub
09/30/2026, 3:32 AMVeniceSystemProducer starts can stay blocked on the controller long after Flink asked it to stop. The controller client treated the interrupt as an ordinary failure: ControllerTransport cleared the interrupt flag, and every retry loop above it started another attempt. The defaults allow 10 attempts of up to 17 minutes each, while Flink kills the whole TaskManager when a task misses its 180-second cancellation deadline, so one slow controller can turn a routine cancellation into mass TaskManager restarts.
The producer also needs the controllers just to find the cluster that hosts its store. Its three /discover_cluster lookups (the store, the KAFKA_MESSAGE_ENVELOPE system store, and transport re-initialization) are read-only and any router can answer them, but they go to the controllers, which also serve leader-bound calls such as /request_topic. The system store lookup had only 2 attempts.
Solution
Commit 1 stops retrying once the caller is interrupted, and leaves the interrupt flag set for the caller:
• ControllerTransport restores the flag when the response wait is interrupted, and keeps it across close(), where shutting down the async client's I/O reactor otherwise swallows it.
• ControllerClient#request, retryableRequest, and the URL-based leader and cluster discovery stop at the first interrupted attempt or backoff, and the error says the request was aborted because the calling thread was interrupted.
• D2ClientUtils restores the flag, and D2ControllerClient stops trying further D2 clients once interrupted.
• VeniceSystemProducer#controllerRequestWithRetry stops instead of starting another attempt.
• VenicePushJob reports and kills a failed push through the same client, so it clears the interrupt for that cleanup and restores it afterwards.
• Failures without an interrupt are retried exactly as before.
The transport, ControllerClient#request and URL-based discovery changes build on a patch by Yanan Hao.
Commit 2 moves the producer's cluster discovery to the routers:
• In D2 mode, the three lookups go to the routers of the producer's own region, through the child colo D2 client and the cluster discovery D2 service. The service defaults to venice-discovery, the thin client's default discovery service (ClientConfig.DEFAULT_CLUSTER_DISCOVERY_D2_SERVICE_NAME), which a standalone RouterServer announces when router.d2.announce.enabled is set.
• Each lookup gets 10 attempts.
• /request_topic and the other leader-bound calls stay on the controllers.
• URL mode (venice.controller.discovery.url) is unchanged.
Known limitation: an interrupt that arrives while the async HTTP client itself is closing is still swallowed by httpcore-nio. A later interrupt, such as Flink's periodic re-interrupt of a cancelling task, still stops the retries.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• venice.cluster.discovery.d2.service, default venice-discovery (ClientConfig.DEFAULT_CLUSTER_DISCOVERY_D2_SERVICE_NAME). Existing jobs need no config change.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging. Not needed: VeniceSystemFactory logs the discovery service when it creates a producer, and retryableRequest logs one warning when it stops on an interrupt.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues. The changes only read, clear and restore the calling thread's own interrupt flag.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed. No new shared state.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList). No new collections.
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
How was this PR tested?
• New unit tests added.
• ControllerInterruptHandlingTest covers the transport, request, retryableRequest, and the URL and D2 leader and cluster discovery, against a server that accepts connections and never answers, and against mocked D2 clients.
• VeniceSystemProducerInterruptTest#testStartExitsPromptlyWhenInterruptedWhileTheControllerStalls reproduces the problem above: start() blocked in /request_topic on a leader controller that never answers fails within seconds of one interrupt, without resending /request_topic, and with the interrupt flag still set.
• VeniceSystemProducerClusterDiscoveryTest covers the three router lookups, the configured D2 service and the 10 attempts.
• TestVenicePushJobCheckpoints#testFailedPushIsReportedAndKilledWhenTheThreadWasInterrupted, D2ClientUtilsTest, and config, builder and factory tests for the new setting.
• Each interrupt and router discovery test fails with the change it covers reverted, and the two tests that check ordinary failures still retry pass either way. That test also fails when only the ControllerTransport#close() fix is removed.
• Modified or extended existing tests.
• Integration tests that build a D2-mode producer set the discovery service to venice-discovery_test, the service that test routers announce. VeniceSystemFactoryTest, PartialUpdateNonAATest#testWriteComputeWithSamzaBatchJob (aggregate mode), TestVTConsistencyCheckerJob and ActiveActiveReplicationForHybridTest#testAAReplicationCanResolveConflicts pass, and VeniceSystemFactoryTest#testGetProducer fails without that setting.
• Verified backward compatibility (if applicable).
• Failures without an interrupt keep their retry counts, URL mode is unchanged, and toBuilder() keeps the new setting.
Does this PR introduce any user-facing or breaking changes?
• Yes. Clearly explain the behavior change and its impact.
• A thread that calls the controller client with its interrupt flag already set now fails at once, where it used to succeed on a retry 5 seconds later.
• D2-mode producers now discover their cluster through the routers of their own region, including producers in aggregate mode, which used to ask the parent controller. Those routers must announce the cluster discovery D2 service, which a standalone RouterServer does when D2 announcement is enabled.
linkedin/veniceGitHub
09/30/2026, 4:56 AM<https://github.com/linkedin/venice/tree/main|main> by minhmo1620
<https://github.com/linkedin/venice/commit/0d7c61b9ce6c2fbf96a86a9c8cb30ce10d79ccd1|0d7c61b9> - [controller] Allow store migration between two encryption clusters (#3049)
linkedin/veniceGitHub
09/30/2026, 10:23 PM<https://github.com/linkedin/venice/tree/main|main> by ymuppala
<https://github.com/linkedin/venice/commit/9e1b95f65e537658d919cb166c54d547290192b7|9e1b95f6> - [server] Skip batch count checks for compacted topics (#3023)
linkedin/veniceGitHub
09/30/2026, 10:27 PM<https://github.com/linkedin/venice/tree/main|main> by ymuppala
<https://github.com/linkedin/venice/commit/087115dc3ea500c04506789edf855c2586aa6288|087115dc> - [fast-client] Make dictionary fetch replica-aware and stop dropping 404s (#3021)
linkedin/veniceGitHub
10/01/2026, 5:05 PMStoreBackupVersionCleanupService.cleanupBackupVersion could delete the previous current version when a higher non-current version (e.g. one left at PUSHED by a kill that finished bootstrap before propagating KILLED to child metadata) sat between the prior current and the current. The keep-newest fallback sorts sub-current versions in descending order and removes index 0 from the deletion list, which preserves the highest version but lets the actually-served prior current through to deletion. DaVinci nodes still subscribed to the prior current then fail at version swap.
Solution
Use the existing Version.previousCurrentVersion field, auto-stamped by ZKStore.setCurrentVersion at swap time since StoreMetaValue v40, to preserve the prior current in the non-repush keep-newest branch. Falls back to keep-newest when the field is unset (pre-v40 stores or rollback paths that bypass the auto-stamp).
For sequential pushes the prior current is the newest sub-current version, so behavior matches the legacy keep-newest path; only the lingering-higher-version case changes.
Tracks: VENG-12676
Code changes
• Added new code behind a config.
• Introduced new log lines.
*Concurrency-Specific Checks
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms are used where needed.
• No blocking calls inside critical sections.
• Verified thread-safe collections are used.
• Validated proper exception handling in multi-threaded code.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
Unit tests cover four scenarios via a data provider: the bug reproduction, a longer-history variant, sequential pushes (legacy match), and back-compat when priorCurrent is unset. Integration test reproduces the full lix-member flow across two regions with a target-region deferred-swap push, kill, follow-up regular push, and asserts the cleanup picks v2 (lingering PUSHED) for deletion while preserving v1 (prior current).
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
linkedin/veniceGitHub
10/01/2026, 9:17 PMwrite.path.thread.priority and validates that configured values are within Java’s valid thread priority range.
Default:
• write.path.thread.priority = Thread.NORM_PRIORITY - 1
Valid range:
• Thread.MIN_PRIORITY through Thread.MAX_PRIORITY
The config is propagated to:
• shared Kafka consumer threads
• Kafka consumer batch-unsubscribe executor threads
• cross-TP parallel processing threads
• store ingestion task executor threads
• store buffer/drainer writer threads
• AA/WC workload processing threads
• AA/WC ingestion storage lookup threads
• blob-transfer client/server Netty event-loop threads
• blob-transfer host-connect, timeout, checksum-validation, snapshot-cleanup, and replica-fetch executors
This keeps write-path and blob-transfer work slightly below normal priority by default, while allowing operators to override the value when needed. The trade-off is that write-path work may yield more readily to normal-priority work under CPU contention unless explicitly configured otherwise.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• write.path.thread.priority
• Default: Thread.NORM_PRIORITY - 1
• Valid range: Thread.MIN_PRIORITY through Thread.MAX_PRIORITY
• No new log lines introduced.
• Log rate limiting not applicable.
Concurrency-Specific Checks
Both reviewer and PR author to verify:
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used where needed.
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Modified existing integration test mock setup for the new config.
• Verified backward compatibility.
Test coverage includes:
• default and overridden write.path.thread.priority config parsing
• invalid priority rejection
• store-writer thread factory priority propagation
• Kafka consumer adapter creation while temporarily applying configured priority
• Kafka consumer batch-unsubscribe executor priority propagation
• cross-TP processing pool priority propagation
• AA/WC workload and storage-lookup thread factory priority propagation
• blob-transfer Netty worker thread priority behavior
• integration-test mock setup for KafkaConsumptionTest
Validated with targeted tests:
./gradlew clientsda-vinci-client:test \
--tests com.linkedin.davinci.kafka.consumer.KafkaConsumerServiceTest \
--tests com.linkedin.davinci.kafka.consumer.KafkaConsumerServiceDelegatorTest \
--tests com.linkedin.davinci.kafka.consumer.AggKafkaConsumerServiceTest \
--tests com.linkedin.davinci.kafka.consumer.KafkaStoreIngestionServiceTest \
--tests com.linkedin.davinci.kafka.consumer.SeparatedStoreBufferServiceTest \
--tests com.linkedin.davinci.kafka.consumer.StoreBufferServiceTest \
--tests com.linkedin.davinci.DaVinciBackendTest
./gradlew internalvenice-test-common:integrationTests_87 \
--tests com.linkedin.venice.kafka.KafkaConsumptionTest
./gradlew spotlessJavaCheck
git diff --check
linkedin/veniceGitHub
10/02/2026, 12:01 AMPUSHED versions are retained indefinitely by version-retirement selection, and non-forced offline-push kills skip them because bootstrap completion includes PUSHED.
Roll-forward also selects a candidate before acquiring the store metadata write lock, allowing concurrent deletion, kill, or supersession to invalidate it. Parent roll-forward uses direct child RPCs while kills use admin messages, so a roll-forward RPC can overtake an already-queued kill.
Solution
• Allow existing retirement selection to remove non-latest, non-current PUSHED versions when the store is not migrating.
• Allow non-forced kills of non-current PUSHED versions to persist KILLED before sending kill messages. Preserve ONLINE/current protection, terminal duplicate suppression, and forced-kill behavior.
• Revalidate candidate existence, latest-version eligibility, status, and non-current status inside the metadata write lock before promotion.
• Route parent roll-forward through the existing ROLLFORWARD_CURRENT_VERSION admin message so roll-forward and kill share the per-store child queue.
• Wait for targeted children's store-scoped execution acknowledgements and expected current versions before reporting success. Parent metadata updates remain dependent on aggregate child state.
• Add a single-attempt metadata-read overload with a supplied timeout for completion polling, without changing existing client defaults.
Existing migration, repush, and retirement guards remain unchanged. No global changes are made to VersionStatus.isBootstrapCompleted or canDelete. No new configuration or admin schema changes are introduced.
Code changes
• Added new code behind a config. If so list the config names and their default values in the PR description.
• Introduced new log lines.
• Confirmed if logs need to be rate limited to avoid excessive logging.
No new configuration or logging sites. An existing kill-skip log message is updated to describe ONLINE/current protection.
*Concurrency-Specific Checks
Both reviewer and PR author to verify
• Code has no race conditions or thread safety issues.
• Proper synchronization mechanisms (e.g., synchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation.
• Verified thread-safe collections are used (e.g., ConcurrentHashMap, CopyOnWriteArrayList).
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
Regression coverage checks stale candidates and both admin-message orderings; it does not establish that every possible race is eliminated. No new shared mutable collections are introduced.
The existing per-store parent admin-message lock remains held during completion waiting. Poll reads use request timeouts, but retries and leader discovery can extend the total wait. Failed consumption retains the queue head for retry and can delay subsequent store operations. Metadata-lock validation remains necessary for mutations outside the admin queue.
How was this PR tested?
• New unit tests added.
• New integration tests added.
• Modified or extended existing tests.
• Verified backward compatibility (if applicable).
• Local code review completed.
Focused Gradle validation passed: 93 unit cases and 3 integration cases, with zero failures, errors, or skips.
| Test class | Executed cases |
| ----------------------------------------------------------- | -------------- |
| TestZKStore | 40 |
| ControllerClientTest | 6 |
| TestVeniceHelixAdmin — selected roll-forward and kill tests | 32 |
| TestVeniceParentHelixAdmin — roll-forward tests | 9 |
| AdminExecutionTaskTest — roll-forward queue tests | 6 |
| ParticipantStoreTest — selected integration tests | 3 |
Coverage includes PUSHED retention and migration, kill metadata/message effects, serving-version protection, forced kills, duplicate suppression, stale candidate rejection, region filtering, child acknowledgement and promotion, request timeouts, and queue ordering/retry.
The existing roll-forward payload was serialized and deserialized using admin schema v76. Mixed-version controller deployment was not exercised.
Scoped uncached Spotless checks, git diff --check, and commit hooks passed. The full repository test suite and performance testing were not run.
Does this PR introduce any user-facing or breaking changes?
• No. You can skip the rest of this section.
• Yes. Clearly explain the behavior change and its impact.
Superseded, non-current PUSHED versions now become eligible for existing retirement. Non-forced kills can terminate non-current PUSHED versions while preserving serving-version safeguards.
Parent roll-forward now uses durable, ordered admin messages instead of direct child RPCs. Success waits for targeted child consumption and promotion. The child selects its future version when the command is consumed; the command does not pin a version at submission.
A response timeout does not cancel the queued command. It may execute later, and readiness retries can delay subsequent commands for the store. No public API fields or admin schema definitions change.
🤖 Generated with GitHub Copilot CLI
linkedin/veniceGitHub
10/02/2026, 12:22 AM<https://github.com/linkedin/venice/tree/main|main> by ymuppala
<https://github.com/linkedin/venice/commit/885f7706993bf2b6f6e15c9525905d355b9b5914|885f7706> - [protocol][compat] Stage metadata v5 with generation pinned to v4 (#3031)
linkedin/veniceGitHub
10/02/2026, 6:24 PMsynchronized, RWLock) are used where needed.
• No blocking calls inside critical sections that could lead to deadlocks or performance degradation. Cleanup intentionally holds the shared lock through metadata, Helix, and ZooKeeper calls to prevent recreation races. Slow calls can delay same-store operations or cluster teardown; the exceptional cluster-name case uses a cluster write lock. The worker does not sleep while holding these locks.
• Verified thread-safe collections are used where needed. Discovery maps and queues are local to each sweep; the stop flag is atomic and the existing lock manager owns shared lock state.
• Validated proper exception handling in multi-threaded code to avoid silent thread termination.
Risk and rollback
This enables deletion of previously stranded metadata, so the main risk is misclassifying a live store as deleted. Owner checks, shared locks, leadership checks, and unchanged retention rules constrain that risk. Slow cleanup I/O and repeated error logging should be monitored during rollout.
Rollback is a controller code revert and redeployment. It stops the new orphan cleanup behavior but does not restore metadata already removed. No schema migration or special component upgrade order is required.
How was this PR tested?
• New unit tests added.
• New integration tests added. None added; regression coverage uses the real metadata repository adapter with synthetic stores and mocked external operations.
• Modified or extended existing tests.
• Verified backward compatibility for existing-store eligibility and retention behavior.
Testing Done
Validated on JDK 17 in an isolated checkout based on upstream main at 885f7706993bf2b6f6e15c9525905d355b9b5914. The committed tree is identical to the tree used for the focused and full controller test runs.
• ./gradlew :services:venice-controller:test --tests com.linkedin.venice.pushmonitor.LeakedPushStatusCleanUpServiceTest spotlessJavaCheck --no-daemon --console=plain -DmaxParallelForks=1 passed: 57 tests, zero failures, skips, or retries.
• ./gradlew :services:venice-controller:test --no-daemon --console=plain -DmaxParallelForks=1 passed: 1,539 tests, zero failures, skips, or retries.
• Controller jacocoTestReport and jacocoTestCoverageVerification passed. Cleanup service coverage, including its worker and comparator, is 97.16% lines and 92.68% branches.
• spotlessCheck and controller diffCoverage passed on the commit. Diff branch coverage is 92.31%, with zero violations.
• git diff --check passed. The normal commit hook ran spotlessApply and generateGHCI successfully without changing the approved patch.
• Formatter dependency setup initially failed because the configured npm registry lacked a public package. Using the public npm registry for that command resolved it; no checks were disabled or repository configuration changed.
• Existing Gradle deprecation warnings remain. Remote CI has not been verified. No integration suite, operational canary, rollout, or performance load test was run.
Does this PR introduce any user-facing or breaking changes?
• No. No public API, wire-format, or configuration changes; this repairs background orphan cleanup.
• Yes.
Checklist
• Tests added/updated.
• Documentation updated in the cleanup service's class documentation.
• CI passing; pending remote verification.
• No breaking changes.
This contribution is original work and is provided under the project's BSD 2-Clause license.
linkedin/venice