Skip to content

Commit 0993d12

Browse files
committed
fix(artifact): repair completed folder publications
Persisted transport completion proves a generation was published once, but it cannot prove the target still retains every referenced object. Allow explicit full synchronization to revalidate the bounded authoritative closure and recreate missing objects without duplicating immutable journal authority. Treat Windows access-denied responses from directory handles as unsupported directory flush behavior, while continuing to propagate file-sync and real storage failures.
1 parent 5b8b289 commit 0993d12

8 files changed

Lines changed: 228 additions & 20 deletions

internal/artifact/sync.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -96,8 +96,9 @@ func SyncWithRepository(
9696
transport, err := OpenFolderTransport(
9797
opts.Target,
9898
FolderTransportOptions{
99-
ForbiddenRoots: forbidden,
100-
StateStore: databaseFolderTransportState{database: database},
99+
ForbiddenRoots: forbidden,
100+
StateStore: databaseFolderTransportState{database: database},
101+
RepairPublished: opts.Full,
101102
},
102103
)
103104
if err != nil {

internal/artifact/sync_test.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,57 @@ func TestArtifactSyncTwoNodeFolderRoundTripAndReplay(t *testing.T) {
127127
assert.Equal(t, "updated response", updatedMessages[1].Content)
128128
}
129129

130+
func TestArtifactSyncFullRepairsMissingPublishedObjectWithoutJournalGrowth(
131+
t *testing.T,
132+
) {
133+
t.Parallel()
134+
135+
target := t.TempDir()
136+
database := testDB(t)
137+
repository, err := OpenRepository(t.Context(), t.TempDir())
138+
require.NoError(t, err)
139+
t.Cleanup(func() { require.NoError(t, repository.Close()) })
140+
origin := "laptop-a1b2c3"
141+
seedSession(t, database, "one", "alpha")
142+
opts := SyncOptions{Target: target, Origin: origin}
143+
144+
initial, err := SyncWithRepository(t.Context(), database, repository, opts)
145+
require.NoError(t, err)
146+
assert.Positive(t, initial.PublishedArtifacts)
147+
head, found, err := database.GetArtifactCheckpointHead(t.Context(), origin)
148+
require.NoError(t, err)
149+
require.True(t, found)
150+
checkpointRef, err := NewRef(
151+
origin,
152+
KindCheckpoints,
153+
fmt.Sprintf("cp-%010d.json", head.Sequence),
154+
)
155+
require.NoError(t, err)
156+
checkpointWire, err := ToWireRef(checkpointRef)
157+
require.NoError(t, err)
158+
checkpointPath := filepath.Join(
159+
target,
160+
checkpointWire.Origin,
161+
string(checkpointWire.Kind),
162+
checkpointWire.Name,
163+
)
164+
require.NoError(t, os.Remove(checkpointPath))
165+
journalSequence := readTestFolderJournalSequence(t, target)
166+
167+
noOp, err := SyncWithRepository(t.Context(), database, repository, opts)
168+
require.NoError(t, err)
169+
assert.Zero(t, noOp.PublishedArtifacts)
170+
assert.NoFileExists(t, checkpointPath)
171+
172+
opts.Full = true
173+
repaired, err := SyncWithRepository(t.Context(), database, repository, opts)
174+
require.NoError(t, err)
175+
assert.Equal(t, 1, repaired.PublishedArtifacts)
176+
assert.FileExists(t, checkpointPath)
177+
assert.Equal(t, journalSequence, readTestFolderJournalSequence(t, target),
178+
"repairing an already-journaled object must not duplicate the journal")
179+
}
180+
130181
func readTestFolderJournalSequence(t *testing.T, target string) int64 {
131182
t.Helper()
132183
root, err := os.OpenRoot(filepath.Join(target, folderJournalDirectory))

internal/artifact/transport.go

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -25,13 +25,15 @@ type ExchangeResult struct {
2525
More bool
2626
}
2727

28-
// FolderTransportOptions defines roots that must remain disjoint from the
29-
// external artifact target.
28+
// FolderTransportOptions configures a bounded folder exchange. RepairPublished
29+
// verifies an already-completed authoritative generation without adding new
30+
// journal events, and is intended for explicit full synchronization.
3031
type FolderTransportOptions struct {
31-
ForbiddenRoots []string
32-
MaxObjects int
33-
MaxBytes int64
34-
StateStore FolderTransportStateStore
32+
ForbiddenRoots []string
33+
MaxObjects int
34+
MaxBytes int64
35+
StateStore FolderTransportStateStore
36+
RepairPublished bool
3537
}
3638

3739
// FolderTransportStateStore persists target-bound continuation state between

internal/artifact/transport_folder.go

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ type folderTransport struct {
4848
syncDirectory func(*os.Root) error
4949
maxObjects int
5050
maxBytes int64
51+
repairPublished bool
5152
pushCursor folderPushCursor
5253
stateStore FolderTransportStateStore
5354
stateLoaded bool
@@ -62,6 +63,7 @@ type folderPushCursor struct {
6263
Offset int `json:"offset,omitempty"`
6364
SegmentIndex int `json:"segment_index,omitempty"`
6465
PublicationSessionID string `json:"publication_session_id,omitempty"`
66+
Repair bool `json:"repair,omitempty"`
6567
}
6668

6769
// OpenFolderTransport opens or initializes a marked artifact exchange target.
@@ -93,12 +95,13 @@ func OpenFolderTransport(
9395
}
9496
}()
9597
transport := &folderTransport{
96-
target: canonical,
97-
root: root,
98-
rootIdentity: identity,
99-
maxObjects: opts.MaxObjects,
100-
maxBytes: opts.MaxBytes,
101-
stateStore: opts.StateStore,
98+
target: canonical,
99+
root: root,
100+
rootIdentity: identity,
101+
maxObjects: opts.MaxObjects,
102+
maxBytes: opts.MaxBytes,
103+
repairPublished: opts.RepairPublished,
104+
stateStore: opts.StateStore,
102105
}
103106
if transport.maxObjects <= 0 {
104107
transport.maxObjects = folderExchangeMaxObjects

internal/artifact/transport_folder_push.go

Lines changed: 18 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -119,10 +119,21 @@ func (t *folderTransport) pushAuthoritativeOriginLocked(
119119
origin string,
120120
) (published int, more bool, retErr error) {
121121
generation := pager.folderTransportGeneration()
122+
continuing := t.pushCursor.Origin == origin &&
123+
t.pushCursor.Generation == generation
122124
if generation == t.publishedGeneration {
123-
return 0, false, nil
124-
}
125-
if t.pushCursor.Origin != origin || t.pushCursor.Generation != generation {
125+
switch {
126+
case continuing && t.pushCursor.Repair:
127+
case !t.repairPublished:
128+
return 0, false, nil
129+
default:
130+
t.pushCursor = folderPushCursor{
131+
Generation: generation,
132+
Origin: origin,
133+
Repair: true,
134+
}
135+
}
136+
} else if !continuing || t.pushCursor.Repair {
126137
t.pushCursor = folderPushCursor{
127138
Generation: generation,
128139
Origin: origin,
@@ -174,8 +185,10 @@ func (t *folderTransport) pushAuthoritativeOriginLocked(
174185
if created {
175186
published++
176187
}
177-
if err := t.appendFolderJournalLocked(ctx, entry); err != nil {
178-
return published, false, err
188+
if !t.pushCursor.Repair {
189+
if err := t.appendFolderJournalLocked(ctx, entry); err != nil {
190+
return published, false, err
191+
}
179192
}
180193
}
181194
t.pushCursor = next

internal/artifact/transport_folder_sync_windows.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,8 @@ import (
99
)
1010

1111
func isFolderDirectorySyncUnsupported(err error) bool {
12-
return errors.Is(err, windows.ERROR_INVALID_FUNCTION) ||
12+
return errors.Is(err, windows.ERROR_ACCESS_DENIED) ||
13+
errors.Is(err, windows.ERROR_INVALID_FUNCTION) ||
1314
errors.Is(err, windows.ERROR_INVALID_HANDLE) ||
1415
errors.Is(err, windows.ERROR_NOT_SUPPORTED)
1516
}
Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
//go:build windows
2+
3+
package artifact
4+
5+
import (
6+
"os"
7+
"testing"
8+
9+
"github.com/stretchr/testify/assert"
10+
"golang.org/x/sys/windows"
11+
)
12+
13+
func TestWindowsDirectorySyncAccessDeniedIsUnsupported(t *testing.T) {
14+
t.Parallel()
15+
16+
assert.True(t, isFolderDirectorySyncUnsupported(&os.PathError{
17+
Op: "sync",
18+
Path: "artifact-folder",
19+
Err: windows.ERROR_ACCESS_DENIED,
20+
}))
21+
assert.False(t, isFolderDirectorySyncUnsupported(&os.PathError{
22+
Op: "sync",
23+
Path: "artifact-folder",
24+
Err: windows.ERROR_DISK_FULL,
25+
}))
26+
}

internal/artifact/transport_publication_store_test.go

Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@ package artifact
33
import (
44
"context"
55
"fmt"
6+
"os"
7+
"path/filepath"
68
"testing"
79

810
"github.com/stretchr/testify/assert"
@@ -237,3 +239,112 @@ func TestFolderTransportNoOpSkipsAuthoritativePublicationPages(t *testing.T) {
237239
assert.Equal(t, 2, authority.pageCalls,
238240
"an unchanged head must not inspect any publication page")
239241
}
242+
243+
func TestFolderTransportResumesBoundedPublishedRepairAfterReopen(t *testing.T) {
244+
t.Parallel()
245+
246+
origin := "local-a1b2c3"
247+
content := newTestArtifactStore(t)
248+
segmentBodies := [][]byte{
249+
[]byte("first segment\n"),
250+
[]byte("second segment\n"),
251+
}
252+
segmentHashes := make([]string, 0, len(segmentBodies))
253+
for _, body := range segmentBodies {
254+
ref := testContentRef(t, origin, KindSegments, body, ".ndjson")
255+
createTestStoreArtifact(t, content, ref, body)
256+
segmentHashes = append(segmentHashes, identityForBytes(t, body).SHA256)
257+
}
258+
manifestBody, err := canonicalJSON(manifest{
259+
Version: manifestFormatVersion,
260+
Origin: origin,
261+
Segments: segmentHashes,
262+
})
263+
require.NoError(t, err)
264+
manifestIdentity := identityForBytes(t, manifestBody)
265+
manifestRef, err := NewRef(
266+
origin,
267+
KindManifests,
268+
manifestIdentity.SHA256+".json",
269+
)
270+
require.NoError(t, err)
271+
createTestStoreArtifact(t, content, manifestRef, manifestBody)
272+
checkpointBody := []byte("checkpoint")
273+
checkpointIdentity := identityForBytes(t, checkpointBody)
274+
checkpointRef, err := NewRef(origin, KindCheckpoints, "cp-0000000001.json")
275+
require.NoError(t, err)
276+
createTestStoreArtifact(t, content, checkpointRef, checkpointBody)
277+
authority := &countingPublicationAuthority{
278+
head: db.ArtifactCheckpointHead{
279+
Origin: origin,
280+
Sequence: 1,
281+
PublicationRevision: 7,
282+
SessionMapSHA256: emptyArtifactPublicationMapSHA256,
283+
CheckpointSHA256: checkpointIdentity.SHA256,
284+
CheckpointSize: checkpointIdentity.Size,
285+
},
286+
publications: []db.ArtifactPublication{{
287+
Origin: origin, SessionID: "one",
288+
ManifestHash: manifestIdentity.SHA256,
289+
}},
290+
}
291+
publishedStore, err := newAuthoritativePublicationStore(
292+
t.Context(), authority, content, origin,
293+
)
294+
require.NoError(t, err)
295+
state := &testFolderTransportStateStore{}
296+
target := t.TempDir()
297+
initial, err := OpenFolderTransport(target, FolderTransportOptions{
298+
MaxObjects: 10,
299+
StateStore: state,
300+
})
301+
require.NoError(t, err)
302+
initialResult, err := initial.Exchange(t.Context(), publishedStore, origin)
303+
require.NoError(t, err)
304+
assert.Equal(t, 4, initialResult.Published)
305+
assert.False(t, initialResult.More)
306+
require.NoError(t, initial.Close())
307+
308+
checkpointWire, err := ToWireRef(checkpointRef)
309+
require.NoError(t, err)
310+
checkpointPath := filepath.Join(
311+
target,
312+
checkpointWire.Origin,
313+
string(checkpointWire.Kind),
314+
checkpointWire.Name,
315+
)
316+
require.NoError(t, os.Remove(checkpointPath))
317+
journalSequence := readTestFolderJournalSequence(t, target)
318+
repair, err := OpenFolderTransport(target, FolderTransportOptions{
319+
MaxObjects: 1,
320+
StateStore: state,
321+
RepairPublished: true,
322+
})
323+
require.NoError(t, err)
324+
firstRepair, err := repair.Exchange(t.Context(), publishedStore, origin)
325+
require.NoError(t, err)
326+
assert.Zero(t, firstRepair.Published)
327+
assert.True(t, firstRepair.More)
328+
require.NoError(t, repair.Close())
329+
330+
published := 0
331+
more := true
332+
for attempts := 0; more && attempts < 10; attempts++ {
333+
resumed, openErr := OpenFolderTransport(target, FolderTransportOptions{
334+
MaxObjects: 1,
335+
StateStore: state,
336+
})
337+
require.NoError(t, openErr)
338+
result, exchangeErr := resumed.Exchange(
339+
t.Context(), publishedStore, origin,
340+
)
341+
require.NoError(t, exchangeErr)
342+
require.NoError(t, resumed.Close())
343+
published += result.Published
344+
more = result.More
345+
}
346+
assert.False(t, more)
347+
assert.Equal(t, 1, published)
348+
assert.FileExists(t, checkpointPath)
349+
assert.Equal(t, journalSequence, readTestFolderJournalSequence(t, target))
350+
}

0 commit comments

Comments
 (0)