https://venicedb.org logo
Join Slack
Powered by
# github-notifications
  • g

    GitHub

    09/25/2026, 6:32 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/25/2026, 7:12 PM
    #3038 Tune and centeralize BlobTransfer thread priority Pull request opened by eldernewborn Problem Statement Blob transfer is not latency-sensitive, so its background threads should yield CPU to the read/write hot path under contention. Most blob-transfer thread pools already used below-normal priority, but two executors still used the default/inherited priority: •
    Venice-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/venice
    • 1
    • 1
  • g

    GitHub

    09/25/2026, 7:46 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/25/2026, 9:06 PM
    #3039 [controller] Preserve Topic Cleanup Eligibility Across Cleanup Scans Pull request opened by KaiSernLim Problem Statement Store hard deletion can stall after Venice metadata is gone because version topics repeatedly lose their place in the cleanup queue and restart the safety delay. In the reported incident, a store version was checked 147 times before deletion was requested. Why this started after #2896 There were two separate bugs. The first accidentally masked the second. Bug 1:
    remove(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/venice
    • 1
    • 1
  • g

    GitHub

    09/25/2026, 10:01 PM
    #3040 [build] Empty commit to trigger a build. Pull request opened by minhmo1620 Problem Statement The scheduled/manual
    Venice 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
    ):
    Copy code
    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/venice
  • g

    GitHub

    09/25/2026, 10:08 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/25/2026, 11:51 PM
    #3041 [da-vinci][server] Select writer factory by store encryption Pull request opened by minhmo1620 Problem Statement Venice servers currently construct all writers from one cluster-configured factory. In an encryption-enabled cluster, this also enables producer encryption for metadata and push-status system stores even though system stores are not encryption-enabled and do not have key-lineage URNs. Solution Add an explicit producer-encryption flag to
    PubSubProducerAdapterContext
    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/venice
    • 1
    • 1
  • g

    GitHub

    09/26/2026, 12:21 AM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/27/2026, 1:52 AM
    #2987 [controller] Preserve lifecycle hooks when update carries empty hooks list Pull request opened by misyel Problem An
    update_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/venice
    • 1
    • 1
  • g

    GitHub

    09/28/2026, 8:13 PM
    #3043 [producer][controller][da-vinci] Scope producer encryption to actual encrypted stores/writers Pull request opened by minhmo1620 Problem Statement • `VeniceWriterFactory`'s 5-arg constructor hardcoded
    producerEncryptionEnabled = 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/venice
    • 1
    • 1
  • g

    GitHub

    09/29/2026, 12:03 AM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/29/2026, 12:59 AM
    #3044 [da-vinci] Tune shared consumer thread priority Pull request opened by eldernewborn Problem Statement Shared Kafka consumer threads currently inherit the JVM/default thread priority because
    KafkaConsumerService
    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/venice
    • 1
    • 1
  • g

    GitHub

    09/29/2026, 2:45 AM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/29/2026, 6:59 PM
    #3045 [protocol][controller] Stage write quota enabled schemas Pull request opened by misyel Problem Statement The store metadata and admin operation schemas need a field to represent whether write quota enforcement is enabled for a store. The schema addition is staged separately from runtime implementation. Solution Add
    writeQuotaEnabled
    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/venice
    • 1
    • 2
  • g

    GitHub

    09/29/2026, 7:58 PM
    #3046  [dvc][server][controller][fc] Add OTel metric lifecycle infra to omit stale or placeholder values Pull request opened by m-nagarajan Problem Statement The OTel SDK calls async gauge and high-performance (observable) counter callbacks on every collection until the instrument is closed. Most were never closed when their store, version or client went away, and several callbacks reported placeholder values when there was nothing to measure. With delta export, these values go out every interval and skew aggregates across hosts: • Series outlive their owner. A store's metrics kept exporting after the store was deleted, after its last ingestion task on the host stopped, or after a version's storage engine left the host. Fast client metrics kept exporting after the client closed. Observable counters then export a
    0
    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/venice
  • g

    GitHub

    09/29/2026, 8:49 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/29/2026, 9:28 PM
    #3048 [protocol] Add write quota enabled creation field Pull request opened by misyel Problem Statement Prerequisite for #3047. Store creation needs an optional gRPC field to distinguish an omitted write quota setting from an explicit
    true
    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/venice
    • 1
    • 1
  • g

    GitHub

    09/29/2026, 10:55 PM
    #3047 [compat][controller] Wire write quota enabled store config Pull request opened by misyel Problem Statement #3045 staged the
    writeQuotaEnabled
    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/venice
    • 1
    • 2
  • g

    GitHub

    09/29/2026, 11:26 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/30/2026, 1:06 AM
    #3049 [controller] Allow store migration between two encryption clusters Pull request opened by minhmo1620 Problem Statement
    StoreMigrationHelper.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/venice
    • 1
    • 1
  • g

    GitHub

    09/30/2026, 1:29 AM
    #3050 [controller] Add encrypted-cluster-to-encrypted-cluster store migration integration test Issue created by minhmo1620 Context Follow-up from #3049 (Copilot review comment: #3049 (comment)). That PR allows store migration between two encryption clusters (previously blocked unconditionally). The clone/update logic that propagates the source store's
    pubSubEncryptionKeyUrn
    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/venice
    • 1
    • 1
  • g

    GitHub

    09/30/2026, 3:32 AM
    #3051 [samza][vpj][controller] Stop controller client retries on interrupt and discover producer clusters through routers Pull request opened by LeoLeo718 Problem Statement A Flink task that is cancelled while its
    VeniceSystemProducer
    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/venice
    • 1
    • 2
  • g

    GitHub

    09/30/2026, 4:56 AM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/30/2026, 10:23 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    09/30/2026, 10:27 PM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    10/01/2026, 5:05 PM
    #2803 [controller] Keep prior current version during backup cleanup Pull request opened by misyel Problem Statement
    StoreBackupVersionCleanupService.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/venice
    • 1
    • 1
  • g

    GitHub

    10/01/2026, 9:17 PM
    #3052 [da-vinci][server]: Configure ingestion and blob-transfer thread priority, added the config Pull request opened by eldernewborn Problem Statement Venice ingestion and blob-transfer write-path work runs across multiple executor pools, but thread priorities were either hardcoded in a few places or left to the JVM/default thread factory behavior. This makes it difficult to consistently tune ingestion/blob-transfer CPU scheduling relative to read-path work under host contention. This PR adds a single server/Da Vinci config to control write-path thread priority and wires it through the relevant ingestion and blob-transfer thread pools so deployments can tune this behavior consistently. Solution This PR introduces
    write.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/venice
  • g

    GitHub

    10/02/2026, 12:01 AM
    #3053 [controller] Retire superseded PUSHED versions and order roll-forward… Pull request opened by misyel Problem Statement Superseded, non-current
    PUSHED
    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/venice
    • 1
    • 1
  • g

    GitHub

    10/02/2026, 12:22 AM
    1 new commit pushed to
    <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/venice
  • g

    GitHub

    10/02/2026, 6:24 PM
    #3054 [controller] Reclaim orphaned push statuses for deleted stores Pull request opened by minhmo1620 Problem Statement Push statuses can outlive store deletion because server-side cleanup is asynchronous. The cleanup service currently requires store metadata to exist, so it skips these orphaned statuses and their Helix resources indefinitely. A cleanup failure can also stop the worker when its error handler tries to read an already-drained version queue. Solution • Reclaim push statuses for confirmed missing stores. For shared system stores, also require the owning user store to be absent; missing shared metadata alone is not evidence of deletion. • Use the controller's shared cluster/store lock manager from metadata lookup through cleanup. This protects store recreation and repository teardown, including the special lock mapping when an owner name equals the cluster name. • Preserve current, future, and metadata-retained versions of existing stores. Keep the newest eligible leaked version until its creation age strictly exceeds the configured linger time; retain it when the timestamp is unknown. Older eligible versions remain immediately eligible. • Keep leadership-checked Helix deletion first and remove residual ZooKeeper push statuses on a later sweep once Helix no longer contains the resource. • Isolate retryable failures, log stable store/topic identifiers, count completed cleanup actions, and stop cleanly on interruption. There are no new dependencies, configuration keys, or schema changes. The existing scan interval and retention settings are unchanged. Code changes • Added new code behind a config. No new config; this fixes the existing cleanup service. • Introduced new log lines. • Confirmed if logs need to be rate limited to avoid excessive logging. Logging occurs in the periodic cleanup sweep, not the request path. No additional limiter was added; persistent failures can log again each sweep. *Concurrency-Specific Checks Both reviewer and PR author to verify. The checked items reflect the author's review and tests; reviewer verification is still needed. • Code has no race conditions or thread safety issues in the reviewed recreation and teardown paths. • 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. 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