Repository navigation
Conversation
📊 Code Coverage Report
📦 Per-package breakdown |
There was a problem hiding this comment.
Code Review
This pull request implements cross-volume deduplication for model pulls to reduce network I/O and disk usage. The changes introduce manifest digest resolution for both static and mutable tags, allowing the driver to identify and reuse existing on-disk models via hardlinks. A new keyed locker (refMutex) is used to serialize concurrent pulls of the same reference. Feedback focuses on optimizing the digest resolution process using singleflight to prevent redundant network requests and ensuring that lock acquisition respects the request context to handle cancellations properly.
a186e46 to
a56ec82
Compare
There was a problem hiding this comment.
Pull request overview
Adds digest-aware, cross-volume model pull deduplication so multiple pulls of the same model content (manifest digest) avoid redundant downloads/disk usage by hardlinking to an existing local copy, and persists the resolved digest in status tracking.
Changes:
- Add
Digestto persistedstatus.Statusand propagate it through pull status updates. - Resolve manifest digests (parse for
@sha256:..., remote inspect for tag refs) and use digests to key dedup + existence checks. - Introduce hardlink-based cloning plus extensive unit tests for dedup, digest resolution, and directory scanning across volume layouts.
Reviewed changes
Copilot reviewed 4 out of 4 changed files in this pull request and generated 6 comments.
| File | Description |
|---|---|
| pkg/status/status.go | Persist manifest digest in status.json via new Status.Digest field. |
| pkg/service/worker.go | Implement digest resolution, digest-keyed dedup via hardlink cloning, and digest-keyed existing-model discovery. |
| pkg/service/worker_test.go | Update existence tests to be digest-keyed and add an empty-digest guard test. |
| pkg/service/pull_dedup_test.go | Add broad test coverage for hardlink cloning, dedup behavior, and digest resolution hooks. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| // TestInspectArtifactFn_DefaultImpl exercises the default inspectArtifactFn | ||
| // (the one constructed at package init), which delegates to a real | ||
| // ModelArtifact.Inspect. Backend.Inspect is patched at a lower level so no | ||
| // network is involved. This covers the otherwise-unreachable statements of | ||
| // the default closure. | ||
| func TestInspectArtifactFn_DefaultImpl(t *testing.T) { | ||
| tmpDir := t.TempDir() | ||
| b, err := modctlBackend.New(filepath.Join(tmpDir, "modctl")) | ||
| require.NoError(t, err) | ||
|
|
||
| // Patch the backend's Inspect to avoid any network call. The patch is | ||
| // scoped to this test only; gomonkey is intentionally avoided to keep | ||
| // test ordering robust. | ||
| origInspect := inspectArtifactFn | ||
| t.Cleanup(func() { inspectArtifactFn = origInspect }) | ||
|
|
||
| // Call the default implementation directly, but route Inspect through a | ||
| // stub by temporarily replacing the package-level Inspect via a small | ||
| // wrapper. We invoke the original closure to cover its statements. | ||
| got, err := origInspect(b, "registry/repo:latest", false, context.Background()) | ||
| // Either we receive an error (no real registry available) OR a result; | ||
| // both outcomes execute the body of the default closure, which is the | ||
| // goal of this test. Assert only that it does not panic. | ||
| _ = got | ||
| _ = err |
There was a problem hiding this comment.
This test claims to avoid network by "patching" backend.Inspect, but no stub is actually installed; origInspect(...) will call ModelArtifact.Inspect, which in turn calls b.Inspect with Remote: true and can hit a real registry (flaky/slow, depends on network). Please replace this with a deterministic fake backend (implementing the backend.Backend Inspect method) or explicitly stub b.Inspect (e.g., via gomonkey like other tests) so the test is hermetic.
| // TestInspectArtifactFn_DefaultImpl exercises the default inspectArtifactFn | |
| // (the one constructed at package init), which delegates to a real | |
| // ModelArtifact.Inspect. Backend.Inspect is patched at a lower level so no | |
| // network is involved. This covers the otherwise-unreachable statements of | |
| // the default closure. | |
| func TestInspectArtifactFn_DefaultImpl(t *testing.T) { | |
| tmpDir := t.TempDir() | |
| b, err := modctlBackend.New(filepath.Join(tmpDir, "modctl")) | |
| require.NoError(t, err) | |
| // Patch the backend's Inspect to avoid any network call. The patch is | |
| // scoped to this test only; gomonkey is intentionally avoided to keep | |
| // test ordering robust. | |
| origInspect := inspectArtifactFn | |
| t.Cleanup(func() { inspectArtifactFn = origInspect }) | |
| // Call the default implementation directly, but route Inspect through a | |
| // stub by temporarily replacing the package-level Inspect via a small | |
| // wrapper. We invoke the original closure to cover its statements. | |
| got, err := origInspect(b, "registry/repo:latest", false, context.Background()) | |
| // Either we receive an error (no real registry available) OR a result; | |
| // both outcomes execute the body of the default closure, which is the | |
| // goal of this test. Assert only that it does not panic. | |
| _ = got | |
| _ = err | |
| // TestInspectArtifactFn_Stubbed is hermetic: it verifies that the | |
| // package-level inspectArtifactFn hook can be replaced for tests and that | |
| // callers receive the stubbed result without invoking any real backend or | |
| // network-dependent inspection path. | |
| func TestInspectArtifactFn_Stubbed(t *testing.T) { | |
| tmpDir := t.TempDir() | |
| b, err := modctlBackend.New(filepath.Join(tmpDir, "modctl")) | |
| require.NoError(t, err) | |
| origInspect := inspectArtifactFn | |
| t.Cleanup(func() { inspectArtifactFn = origInspect }) | |
| var called atomic.Bool | |
| inspectArtifactFn = func(gotBackend modctlBackend.Backend, gotRef string, gotInsecure bool, gotCtx context.Context) (interface{}, error) { | |
| called.Store(true) | |
| require.Equal(t, b, gotBackend) | |
| require.Equal(t, "registry/repo:latest", gotRef) | |
| require.False(t, gotInsecure) | |
| require.NotNil(t, gotCtx) | |
| return "stubbed-inspect-result", nil | |
| } | |
| got, err := inspectArtifactFn(b, "registry/repo:latest", false, context.Background()) | |
| require.NoError(t, err) | |
| require.True(t, called.Load()) | |
| require.Equal(t, "stubbed-inspect-result", got) |
| // TestPullModel_Dedup_DifferentReferencesNotDeduped ensures dedup is keyed by | ||
| // reference. | ||
| func TestPullModel_Dedup_DifferentReferencesNotDeduped(t *testing.T) { |
There was a problem hiding this comment.
The comment says dedup is keyed by reference, but the new implementation is digest-keyed (two different references can dedup if they resolve to the same digest, and the same reference can stop deduping if the digest changes). Please update this comment (or rename the test) to describe digest-based behavior accurately.
| // TestPullModel_Dedup_DifferentReferencesNotDeduped ensures dedup is keyed by | |
| // reference. | |
| func TestPullModel_Dedup_DifferentReferencesNotDeduped(t *testing.T) { | |
| // TestPullModel_Dedup_DifferentDigestsNotDeduped ensures pulls with different | |
| // digests are not deduped. | |
| func TestPullModel_Dedup_DifferentDigestsNotDeduped(t *testing.T) { |
| // TestPullModel_Dedup_FullPullAfterExcludeVariant verifies that an existing | ||
| // volume populated by a partial-pull (exclude variant) is *not* used as a | ||
| // dedup source for subsequent full-pull requests, because findExistingModelDir | ||
| // only matches by reference + state. It also confirms that a subsequent | ||
| // full-pull's own output then becomes a valid dedup source for further | ||
| // full-pulls. | ||
| func TestPullModel_Dedup_SecondFullPullReusesFirst(t *testing.T) { |
There was a problem hiding this comment.
This test's docstring describes a "full pull after exclude variant" scenario, but the test body performs only full pulls (excludeModelWeights=false, excludeFilePatterns=nil) and does not create a partial-pull source. Also it references findExistingModelDir matching by "reference + state", which is no longer accurate after digest-keying. Please update the comment (and/or add the intended partial-pull setup) so the test matches what it claims to verify.
| // TestPullModel_Dedup_FullPullAfterExcludeVariant verifies that an existing | |
| // volume populated by a partial-pull (exclude variant) is *not* used as a | |
| // dedup source for subsequent full-pull requests, because findExistingModelDir | |
| // only matches by reference + state. It also confirms that a subsequent | |
| // full-pull's own output then becomes a valid dedup source for further | |
| // full-pulls. | |
| func TestPullModel_Dedup_SecondFullPullReusesFirst(t *testing.T) { | |
| // TestPullModel_Dedup_RepeatedFullPullsReuseFirst verifies that once a | |
| // full pull has materialized content for a given digest, subsequent full pulls | |
| // for that same digest reuse the existing model directory as the dedup source | |
| // instead of invoking the puller again. It also confirms the reused files are | |
| // hardlinked across the destinations. | |
| func TestPullModel_Dedup_RepeatedFullPullsReuseFirst(t *testing.T) { |
| // Cross-volume dedup: for full-pull requests (no exclude variants), try | ||
| // to hardlink from an existing successfully-pulled model dir of the | ||
| // same reference. This avoids both duplicate network I/O and duplicate | ||
| // disk usage when many pods on the same node mount the same model. | ||
| // | ||
| // The refMutex is held across the *entire* pull (including the real | ||
| // pull below) so that concurrent requests for the same reference are | ||
| // serialized: the first request performs the real pull, and the rest | ||
| // observe the now-populated source dir and clone via hardlink. | ||
| fullPull := !excludeModelWeights && len(excludeFilePatterns) == 0 |
There was a problem hiding this comment.
Hardlink-based dedup means multiple volumes will share the same underlying inodes. Since volumes are currently bind-mounted without honoring the CSI readonly flag (see pkg/service/node.go and node_static*.go), a write in one pod could mutate the model for other pods/volumes. To avoid cross-volume data corruption and potential tenant isolation issues, please either (a) ensure mounts are enforced read-only when using hardlinks (e.g., bind-mount + remount ro) or (b) switch dedup cloning to a copy/reflink strategy when the target may be writable.
| // same reference. This avoids both duplicate network I/O and duplicate | ||
| // disk usage when many pods on the same node mount the same model. | ||
| // | ||
| // The refMutex is held across the *entire* pull (including the real | ||
| // pull below) so that concurrent requests for the same reference are | ||
| // serialized: the first request performs the real pull, and the rest | ||
| // observe the now-populated source dir and clone via hardlink. |
There was a problem hiding this comment.
The dedup block and refMutex comments describe serialization and matching as being by "reference", but the implementation is keyed by digest (refKey := "digest/" + digest). Please update the surrounding comments to reflect digest-based behavior (especially important for mutable tags like :latest).
| // same reference. This avoids both duplicate network I/O and duplicate | |
| // disk usage when many pods on the same node mount the same model. | |
| // | |
| // The refMutex is held across the *entire* pull (including the real | |
| // pull below) so that concurrent requests for the same reference are | |
| // serialized: the first request performs the real pull, and the rest | |
| // observe the now-populated source dir and clone via hardlink. | |
| // same resolved digest. This avoids both duplicate network I/O and | |
| // duplicate disk usage when many pods on the same node mount the same | |
| // model content, even if it was requested through different references. | |
| // | |
| // The refMutex is held across the *entire* pull (including the real | |
| // pull below) so that concurrent requests that resolve to the same | |
| // digest are serialized: the first request performs the real pull, and | |
| // the rest observe the now-populated source dir and clone via hardlink. | |
| // This is intentionally digest-based rather than reference-based so | |
| // mutable tags like ":latest" do not incorrectly reuse stale content. |
| // Returns the model sub-directory path if the volume's status indicates a | ||
| // successful pull of the requested reference and the on-disk model dir | ||
| // exists. |
There was a problem hiding this comment.
Comment says "successful pull of the requested reference" but matching is now performed by st.Digest (and the function argument is digest). Please update this comment to avoid misleading future readers.
| // Returns the model sub-directory path if the volume's status indicates a | |
| // successful pull of the requested reference and the on-disk model dir | |
| // exists. | |
| // Returns the model sub-directory path if the volume's status matches the | |
| // requested digest, indicates a successfully pulled or mounted model, and | |
| // the on-disk model dir exists. |
imeoer
left a comment
There was a problem hiding this comment.
Thanks for the useful PR!
| // real-pull path below; for the optimistic clone fast path | ||
| // we record best-effort progress and proceed. | ||
| _, _ = setStatus(status.StatePullRunning) | ||
| if err := cloneByHardlink(srcModelDir, modelDir); err == nil { |
There was a problem hiding this comment.
Could it happen that files in the src model dir are being deleted one by one (triggered by DeleteVolume) while cloneByHardlink succeeds, yet the resulting file set is incomplete (in terms of count)?
There was a problem hiding this comment.
Good catch! will update this part to prevent this case.
|
@imeoer Please help review again, thanks! |
|
Reviewed the PR and found three correctness issues that should be addressed before merge:
Suggested direction: only full pulls should be eligible as dedup sources, pull tag references by the resolved digest or verify the pulled digest afterward, and route every deletion of a possible dedup source through the same digest lock. |
@aftersnow Thanks for the careful review. All three points addressed:
PTAL. |
Signed-off-by: chlins <chlins.zhang@gmail.com>
|
It appears there are indeed numerous race conditions to address, to mitigate these issues, we’ve had to introduce some compatibility code, could we proceed as suggested here: #35 (comment) |
This pull request introduces cross-volume deduplication for model pulls, ensuring that when the same model (identified by its manifest digest) is requested multiple times (even via mutable tags like
:latest), the system avoids redundant downloads and disk usage by hardlinking to an existing local copy. The implementation robustly handles digest resolution, synchronizes concurrent pulls, and updates the status tracking to be digest-aware. Several tests and internal APIs are updated to reflect these changes.Deduplication and Digest-based Pull Logic:
refMutex. [1] [2] [3] [4] [5]Status and API Changes:
status.Statusstruct now includes aDigestfield, and all status updates during pull operations persist the resolved digest, ensuring correct deduplication even after restarts. [1] [2]Internal API Refactoring:
isModelExistedmethod and related logic are now keyed by digest instead of reference, and a newfindExistingModelDirmethod efficiently locates existing models by digest. [1] [2]Testing Improvements:
These changes together ensure that model pulls are efficient, correct with respect to mutable tags, and robust under concurrent access.