Skip to content

Commit c80c79a

Browse files
committed
check storage ID first
1 parent 565251d commit c80c79a

2 files changed

Lines changed: 27 additions & 27 deletions

File tree

cmd/curio/storacha-migration.go

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -126,14 +126,19 @@ var importPiecesCmd = &cli.Command{
126126
return xerrors.Errorf("resolving target path: %w", err)
127127
}
128128

129+
storageID, err := pieceStorageID(targetPath)
130+
if err != nil {
131+
return xerrors.Errorf("getting storage ID: %w", err)
132+
}
133+
129134
db, err := deps.MakeDB(cctx)
130135
if err != nil {
131136
return err
132137
}
133138

134139
si := paths.NewDBIndex(curioalerting.NewAlertingSystem(), db)
135140

136-
out, err := runImportPieces(ctx, db, si, sourcePath, targetPath, cctx.Int("batch-size"))
141+
out, err := runImportPieces(ctx, db, si, storageID, sourcePath, targetPath, cctx.Int("batch-size"))
137142
if err != nil {
138143
return err
139144
}
@@ -147,14 +152,9 @@ var importPiecesCmd = &cli.Command{
147152
},
148153
}
149154

150-
func runImportPieces(ctx context.Context, db *harmonydb.DB, si paths.SectorIndex, sourcePath, targetPath string, batchSize int) (importPiecesOutput, error) {
155+
func runImportPieces(ctx context.Context, db *harmonydb.DB, si paths.SectorIndex, storageID storiface.ID, sourcePath, targetPath string, batchSize int) (importPiecesOutput, error) {
151156
var out importPiecesOutput
152157

153-
storageID, err := pieceStorageID(targetPath)
154-
if err != nil {
155-
return out, err
156-
}
157-
158158
staging := filepath.Join(targetPath, "storacha-staging")
159159
existing, err := os.ReadDir(staging)
160160
if err != nil {

cmd/curio/storacha-migration_test.go

Lines changed: 20 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,7 @@ func TestStorachaMigrationFreshImport(t *testing.T) {
9494
fx := pieces.next()
9595
writeSourcePiece(t, env, fx)
9696

97-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
97+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
9898
require.NoError(t, err)
9999
require.Equal(t, 1, out.Count)
100100
require.Equal(t, []string{fx.cidV2}, out.Pieces)
@@ -121,7 +121,7 @@ func TestStorachaMigrationExistingCompleteMovesStagedWhenFinalMissing(t *testing
121121
pieceID := insertParkedPiece(t, env, fx, true, false, nil)
122122
writeStagedPiece(t, env, fx)
123123

124-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
124+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
125125
require.NoError(t, err)
126126
require.Equal(t, 1, out.Count)
127127

@@ -144,7 +144,7 @@ func TestStorachaMigrationExistingCompleteRemovesStagedDuplicateWhenFinalExists(
144144
writeStagedPiece(t, env, fx)
145145
writeFinalPiece(t, env, pieceID, finalBytes)
146146

147-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
147+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
148148
require.NoError(t, err)
149149
require.Equal(t, 1, out.Count)
150150

@@ -162,7 +162,7 @@ func TestStorachaMigrationClaimsIncompleteParkedPieceWithNoTask(t *testing.T) {
162162
pieceID := insertParkedPiece(t, env, fx, false, false, nil)
163163
writeStagedPiece(t, env, fx)
164164

165-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
165+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
166166
require.NoError(t, err)
167167
require.Equal(t, 1, out.Count)
168168

@@ -184,7 +184,7 @@ func TestStorachaMigrationClaimsIncompleteParkedPieceWithStaleTask(t *testing.T)
184184
pieceID := insertParkedPiece(t, env, fx, false, false, &staleTaskID)
185185
writeStagedPiece(t, env, fx)
186186

187-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
187+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
188188
require.NoError(t, err)
189189
require.Equal(t, 1, out.Count)
190190

@@ -204,7 +204,7 @@ func TestStorachaMigrationLeavesLiveTaskPieceAlone(t *testing.T) {
204204
pieceID := insertParkedPiece(t, env, fx, false, false, &taskID)
205205
writeStagedPiece(t, env, fx)
206206

207-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
207+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
208208
require.NoError(t, err)
209209
require.Equal(t, 0, out.Count)
210210

@@ -227,7 +227,7 @@ func TestStorachaMigrationReusesExistingRefs(t *testing.T) {
227227
pdpID := insertPDPRef(t, env, refID, serviceName, fx.info.PieceCIDV1.String())
228228
writeStagedPiece(t, env, fx)
229229

230-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
230+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
231231
require.NoError(t, err)
232232
require.Equal(t, 1, out.Count)
233233

@@ -246,7 +246,7 @@ func TestStorachaMigrationInsertsMissingPDPRefForExistingParkedRef(t *testing.T)
246246
refID := insertStorachaParkedRef(t, env, pieceID)
247247
writeStagedPiece(t, env, fx)
248248

249-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
249+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
250250
require.NoError(t, err)
251251
require.Equal(t, 1, out.Count)
252252
require.Equal(t, []int64{refID}, storachaRefIDs(t, env, pieceID))
@@ -262,7 +262,7 @@ func TestStorachaMigrationDuplicateStorachaRefsErrors(t *testing.T) {
262262
insertStorachaParkedRef(t, env, pieceID)
263263
writeStagedPiece(t, env, fx)
264264

265-
_, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
265+
_, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
266266
require.ErrorContains(t, err, "expected at most 1")
267267

268268
pp := parkedPieceByID(t, env, pieceID)
@@ -280,7 +280,7 @@ func TestStorachaMigrationMismatchedPDPRefErrors(t *testing.T) {
280280
insertPDPRef(t, env, refID, serviceName, "wrong-piece-cid")
281281
writeStagedPiece(t, env, fx)
282282

283-
_, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
283+
_, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
284284
require.ErrorContains(t, err, "does not match storacha migration piece")
285285

286286
pp := parkedPieceByID(t, env, pieceID)
@@ -298,7 +298,7 @@ func TestStorachaMigrationRecoversFinalFile(t *testing.T) {
298298
insertPDPRef(t, env, refID, serviceName, fx.info.PieceCIDV1.String())
299299
writeFinalPiece(t, env, pieceID, fx.data)
300300

301-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
301+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
302302
require.NoError(t, err)
303303
require.Equal(t, 1, out.Count)
304304

@@ -318,7 +318,7 @@ func TestStorachaMigrationFinalRecoveryMissingFileNoops(t *testing.T) {
318318
refID := insertStorachaParkedRef(t, env, pieceID)
319319
insertPDPRef(t, env, refID, serviceName, fx.info.PieceCIDV1.String())
320320

321-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
321+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
322322
require.NoError(t, err)
323323
require.Equal(t, 0, out.Count)
324324

@@ -336,7 +336,7 @@ func TestStorachaMigrationFinalRecoveryDirectoryErrors(t *testing.T) {
336336
insertPDPRef(t, env, refID, serviceName, fx.info.PieceCIDV1.String())
337337
require.NoError(t, os.MkdirAll(storachaFinalPiecePath(env.targetDir, pieceID), 0755))
338338

339-
_, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
339+
_, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
340340
require.ErrorContains(t, err, "final piece path is a directory")
341341

342342
pp := parkedPieceByID(t, env, pieceID)
@@ -354,7 +354,7 @@ func TestStorachaMigrationFinalRecoveryLeavesLiveTaskPieceAlone(t *testing.T) {
354354
insertPDPRef(t, env, refID, serviceName, fx.info.PieceCIDV1.String())
355355
writeFinalPiece(t, env, pieceID, fx.data)
356356

357-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
357+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
358358
require.NoError(t, err)
359359
require.Equal(t, 0, out.Count)
360360

@@ -373,7 +373,7 @@ func TestStorachaMigrationBatchProcessesStagingBeforeSource(t *testing.T) {
373373
writeStagedPiece(t, env, staged)
374374
writeSourcePiece(t, env, source)
375375

376-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 1)
376+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 1)
377377
require.NoError(t, err)
378378
require.Equal(t, 1, out.Count)
379379
require.Equal(t, []string{staged.cidV2}, out.Pieces)
@@ -395,7 +395,7 @@ func TestStorachaMigrationBatchProcessesFinalRecoveryBeforeSource(t *testing.T)
395395
writeFinalPiece(t, env, pieceID, recovered.data)
396396
writeSourcePiece(t, env, source)
397397

398-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 1)
398+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 1)
399399
require.NoError(t, err)
400400
require.Equal(t, 1, out.Count)
401401
require.Equal(t, []string{recovered.cidV2}, out.Pieces)
@@ -413,7 +413,7 @@ func TestStorachaMigrationInvalidStagedFileDoesNotConsumeBatch(t *testing.T) {
413413
writeFile(t, invalidStagedPath, []byte("invalid staged input"))
414414
writeSourcePiece(t, env, source)
415415

416-
out, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 1)
416+
out, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 1)
417417
require.NoError(t, err)
418418
require.Equal(t, 1, out.Count)
419419
require.Equal(t, []string{source.cidV2}, out.Pieces)
@@ -427,7 +427,7 @@ func TestStorachaMigrationMissingSourcePathErrors(t *testing.T) {
427427
env := setupStorachaMigrationTest(t)
428428
missingSource := filepath.Join(filepath.Dir(env.sourceDir), "missing-source")
429429

430-
_, err := runImportPieces(env.ctx, env.db, env.si, missingSource, env.targetDir, 20)
430+
_, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, missingSource, env.targetDir, 20)
431431
require.ErrorContains(t, err, "reading directory")
432432
}
433433

@@ -438,7 +438,7 @@ func TestStorachaMigrationSourceImportErrorsWhenStagingPathExists(t *testing.T)
438438
writeSourcePiece(t, env, source)
439439
require.NoError(t, os.MkdirAll(filepath.Join(env.stagingDir, source.fileName), 0755))
440440

441-
_, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
441+
_, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
442442
require.ErrorContains(t, err, "staging file already exists")
443443
assertFileBytes(t, filepath.Join(env.sourceDir, source.fileName), source.data)
444444
}
@@ -451,7 +451,7 @@ func TestStorachaMigrationStagedImportFinalDirectoryErrors(t *testing.T) {
451451
writeStagedPiece(t, env, fx)
452452
require.NoError(t, os.MkdirAll(storachaFinalPiecePath(env.targetDir, pieceID), 0755))
453453

454-
_, err := runImportPieces(env.ctx, env.db, env.si, env.sourceDir, env.targetDir, 20)
454+
_, err := runImportPieces(env.ctx, env.db, env.si, env.storageID, env.sourceDir, env.targetDir, 20)
455455
require.ErrorContains(t, err, "final piece path is a directory")
456456

457457
pp := parkedPieceByID(t, env, pieceID)

0 commit comments

Comments
 (0)