Skip to content
This repository was archived by the owner on Oct 9, 2023. It is now read-only.

Commit 2e80fba

Browse files
author
Nick Müller
committed
Implement acquiring/releasing of reserverations as bulk operation
Implement deleting of artifacts as bulk operation Signed-off-by: Nick Müller <nmueller@blackshark.ai>
1 parent b6aaa3c commit 2e80fba

14 files changed

Lines changed: 1538 additions & 47 deletions

go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ module github.com/flyteorg/datacatalog
33
go 1.18
44

55
replace (
6-
github.com/flyteorg/flyteidl => github.com/blackshark-ai/flyteidl v0.24.22-0.20230104143947-9cc1f12a643f
6+
github.com/flyteorg/flyteidl => github.com/blackshark-ai/flyteidl v0.24.22-0.20230119104851-e9fb728f4733
77
github.com/flyteorg/flytestdlib => github.com/blackshark-ai/flytestdlib v1.0.1-0.20230104151410-d6ec6dba8697
88
)
99

go.sum

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -185,8 +185,8 @@ github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6r
185185
github.com/bgentry/speakeasy v0.1.0/go.mod h1:+zsyZBPWlz7T6j88CTgSN5bM796AkVf0kBD4zp0CCIs=
186186
github.com/bketelsen/crypt v0.0.3-0.20200106085610-5cbc8cc4026c/go.mod h1:MKsuJmJgSg28kpZDP6UIiPt0e0Oz0kqKNGyRaWEPv84=
187187
github.com/bketelsen/crypt v0.0.4/go.mod h1:aI6NrJ0pMGgvZKL1iVgXLnfIFJtfV+bKCoqOes/6LfM=
188-
github.com/blackshark-ai/flyteidl v0.24.22-0.20230104143947-9cc1f12a643f h1:BFpozetWgkWdrOz2wDQY/2pcvm8Wx60uI4mJup+0yrg=
189-
github.com/blackshark-ai/flyteidl v0.24.22-0.20230104143947-9cc1f12a643f/go.mod h1:sgOlQA2lnugarwSN8M+9gWoCZmzYNFI8gpShZrm+wmo=
188+
github.com/blackshark-ai/flyteidl v0.24.22-0.20230119104851-e9fb728f4733 h1:xvr9DovH3m3S10+U71+lJ6zmJDFoOHPLSZfXSRtZxEo=
189+
github.com/blackshark-ai/flyteidl v0.24.22-0.20230119104851-e9fb728f4733/go.mod h1:sgOlQA2lnugarwSN8M+9gWoCZmzYNFI8gpShZrm+wmo=
190190
github.com/blackshark-ai/flytestdlib v1.0.1-0.20230104151410-d6ec6dba8697 h1:N3ch1D89Hhe72LK9zVuSwZ9DrRl9G3M5fsntD6ydZEY=
191191
github.com/blackshark-ai/flytestdlib v1.0.1-0.20230104151410-d6ec6dba8697/go.mod h1:ojJnMQ9sDf1VW6zrHv8TauKrMV0+Pf+6n+uProvhzac=
192192
github.com/blang/semver v3.5.0+incompatible/go.mod h1:kRBLl5iJ+tD4TcOOxsy/0fnwebNt5EWlYSAyrTnjyyk=

pkg/errors/errors.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,7 @@ func NewCollectedErrors(code codes.Code, errors []error) error {
5151
errorCollection[idx] = err.Error()
5252
}
5353

54-
return NewDataCatalogError(code, strings.Join((errorCollection), ", "))
54+
return NewDataCatalogError(code, strings.Join(errorCollection, ", "))
5555
}
5656

5757
func IsAlreadyExistsError(err error) bool {

pkg/manager/impl/artifact_manager.go

Lines changed: 71 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,9 @@ type artifactMetrics struct {
4949
deleteResponseTime labeled.StopWatch
5050
deleteSuccessCounter labeled.Counter
5151
deleteFailureCounter labeled.Counter
52+
bulkDeleteResponseTime labeled.StopWatch
53+
bulkDeleteSuccessCounter labeled.Counter
54+
bulkDeleteFailureCounter labeled.Counter
5255
}
5356

5457
type artifactManager struct {
@@ -57,7 +60,8 @@ type artifactManager struct {
5760
systemMetrics artifactMetrics
5861
}
5962

60-
// Create an Artifact along with the associated ArtifactData. The ArtifactData will be stored in an offloaded location.
63+
// CreateArtifact creates an Artifact along with the associated ArtifactData. The ArtifactData will be stored in an
64+
// offloaded location.
6165
func (m *artifactManager) CreateArtifact(ctx context.Context, request *datacatalog.CreateArtifactRequest) (*datacatalog.CreateArtifactResponse, error) {
6266
timer := m.systemMetrics.createResponseTime.Start(ctx)
6367
defer timer.Stop()
@@ -133,7 +137,7 @@ func (m *artifactManager) CreateArtifact(ctx context.Context, request *datacatal
133137
return &datacatalog.CreateArtifactResponse{}, nil
134138
}
135139

136-
// Get the Artifact and its associated ArtifactData. The request can query by ArtifactID or TagName.
140+
// GetArtifact retrieves the Artifact and its associated ArtifactData. The request can query by ArtifactID or TagName.
137141
func (m *artifactManager) GetArtifact(ctx context.Context, request *datacatalog.GetArtifactRequest) (*datacatalog.GetArtifactResponse, error) {
138142
timer := m.systemMetrics.getResponseTime.Start(ctx)
139143
defer timer.Stop()
@@ -244,6 +248,8 @@ func (m *artifactManager) getArtifactDataList(ctx context.Context, artifactDataM
244248
return artifactDataList, nil
245249
}
246250

251+
// ListArtifacts returns a paginated list of artifacts matching the provided filter expression, including their
252+
// associated artifact data.
247253
func (m *artifactManager) ListArtifacts(ctx context.Context, request *datacatalog.ListArtifactsRequest) (*datacatalog.ListArtifactsResponse, error) {
248254
err := validators.ValidateListArtifactRequest(request)
249255
if err != nil {
@@ -408,36 +414,22 @@ func (m *artifactManager) UpdateArtifact(ctx context.Context, request *datacatal
408414
}, nil
409415
}
410416

411-
// DeleteArtifact deletes the given artifact, removing all stored artifact data from the underlying blob storage.
412-
func (m *artifactManager) DeleteArtifact(ctx context.Context, request *datacatalog.DeleteArtifactRequest) (*datacatalog.DeleteArtifactResponse, error) {
413-
ctx = contextutils.WithProjectDomain(ctx, request.Dataset.Project, request.Dataset.Domain)
414-
415-
timer := m.systemMetrics.deleteResponseTime.Start(ctx)
416-
defer timer.Stop()
417-
418-
err := validators.ValidateDeleteArtifactRequest(request)
419-
if err != nil {
420-
logger.Warningf(ctx, "Invalid delete artifact request %v, err: %v", request, err)
421-
m.systemMetrics.validationErrorCounter.Inc(ctx)
422-
m.systemMetrics.deleteFailureCounter.Inc(ctx)
423-
return nil, err
424-
}
417+
func (m *artifactManager) deleteArtifact(ctx context.Context, datasetID *datacatalog.DatasetID, queryHandle artifactQueryHandle) error {
418+
ctx = contextutils.WithProjectDomain(ctx, datasetID.Project, datasetID.Domain)
425419

426420
// artifact must already exist, verify first
427-
artifactModel, err := m.findArtifact(ctx, request.GetDataset(), request)
421+
artifactModel, err := m.findArtifact(ctx, datasetID, queryHandle)
428422
if err != nil {
429-
logger.Errorf(ctx, "Failed to get artifact for delete artifact request %v, err: %v", request, err)
430-
m.systemMetrics.deleteFailureCounter.Inc(ctx)
431-
return nil, err
423+
logger.Errorf(ctx, "Failed to get artifact while trying to delete [%v], err: %v", queryHandle, err)
424+
return err
432425
}
433426

434427
// delete all artifact data from the blob storage
435428
for _, artifactData := range artifactModel.ArtifactData {
436429
if err := m.artifactStore.DeleteData(ctx, artifactData); err != nil {
437-
logger.Errorf(ctx, "Failed to delete artifact data [%v] during delete, err: %v", artifactData.Name, err)
430+
logger.Errorf(ctx, "Failed to delete artifact data [%v] while deleting artifact [%v], err: %v", artifactData.Name, artifactModel.ArtifactID, err)
438431
m.systemMetrics.deleteDataFailureCounter.Inc(ctx)
439-
m.systemMetrics.deleteFailureCounter.Inc(ctx)
440-
return nil, err
432+
return err
441433
}
442434

443435
m.systemMetrics.deleteDataSuccessCounter.Inc(ctx)
@@ -447,21 +439,71 @@ func (m *artifactManager) DeleteArtifact(ctx context.Context, request *datacatal
447439
err = m.repo.ArtifactRepo().Delete(ctx, artifactModel)
448440
if err != nil {
449441
if errors.IsDoesNotExistError(err) {
450-
logger.Warnf(ctx, "Artifact does not exist key: %+v, err %v", artifactModel.ArtifactID, err)
442+
logger.Warnf(ctx, "Artifact [%v] does not exist, err %v", artifactModel.ArtifactID, err)
451443
m.systemMetrics.doesNotExistCounter.Inc(ctx)
452444
} else {
453-
logger.Errorf(ctx, "Failed to delete artifact %v, err: %v", artifactModel, err)
445+
logger.Errorf(ctx, "Failed to delete artifact [%v], err: %v", artifactModel, err)
454446
}
447+
return err
448+
}
449+
450+
logger.Debugf(ctx, "Successfully deleted artifact [%v]", artifactModel.ArtifactID)
451+
return nil
452+
}
453+
454+
// DeleteArtifact deletes the given artifact, removing all stored artifact data from the underlying blob storage.
455+
func (m *artifactManager) DeleteArtifact(ctx context.Context, request *datacatalog.DeleteArtifactRequest) (*datacatalog.DeleteArtifactResponse, error) {
456+
ctx = contextutils.WithProjectDomain(ctx, request.Dataset.Project, request.Dataset.Domain)
457+
458+
timer := m.systemMetrics.deleteResponseTime.Start(ctx)
459+
defer timer.Stop()
460+
461+
err := validators.ValidateDeleteArtifactRequest(request)
462+
if err != nil {
463+
logger.Warningf(ctx, "Invalid delete artifacts request %v, err: %v", request, err)
464+
m.systemMetrics.validationErrorCounter.Inc(ctx)
455465
m.systemMetrics.deleteFailureCounter.Inc(ctx)
456466
return nil, err
457467
}
458468

459-
logger.Debugf(ctx, "Successfully deleted artifact id: %v", artifactModel.ArtifactID)
469+
if err := m.deleteArtifact(ctx, request.GetDataset(), request); err != nil {
470+
m.systemMetrics.deleteFailureCounter.Inc(ctx)
471+
return nil, err
472+
}
460473

461474
m.systemMetrics.deleteSuccessCounter.Inc(ctx)
462475
return &datacatalog.DeleteArtifactResponse{}, nil
463476
}
464477

478+
// DeleteArtifacts deletes the given artifacts, removing all stored artifact data from the underlying blob storage.
479+
func (m *artifactManager) DeleteArtifacts(ctx context.Context, request *datacatalog.DeleteArtifactsRequest) (*datacatalog.DeleteArtifactResponse, error) {
480+
timer := m.systemMetrics.bulkDeleteResponseTime.Start(ctx)
481+
defer timer.Stop()
482+
483+
err := validators.ValidateDeleteArtifactsRequest(request)
484+
if err != nil {
485+
logger.Warningf(ctx, "Invalid delete artifacts request %v, err: %v", request, err)
486+
m.systemMetrics.validationErrorCounter.Inc(ctx)
487+
m.systemMetrics.bulkDeleteFailureCounter.Inc(ctx)
488+
return nil, err
489+
}
490+
491+
for _, deleteArtifactReq := range request.Artifacts {
492+
if err := m.deleteArtifact(ctx, deleteArtifactReq.GetDataset(), deleteArtifactReq); err != nil {
493+
// bulk delete endpoint is idempotent, ignore errors regarding missing artifacts as they might've already
494+
// been deleted by a previous call.
495+
if errors.IsDoesNotExistError(err) {
496+
continue
497+
}
498+
m.systemMetrics.bulkDeleteFailureCounter.Inc(ctx)
499+
return nil, err
500+
}
501+
}
502+
503+
m.systemMetrics.bulkDeleteSuccessCounter.Inc(ctx)
504+
return &datacatalog.DeleteArtifactResponse{}, nil
505+
}
506+
465507
func NewArtifactManager(repo repositories.RepositoryInterface, store *storage.DataStore, storagePrefix storage.DataReference, artifactScope promutils.Scope) interfaces.ArtifactManager {
466508
artifactMetrics := artifactMetrics{
467509
scope: artifactScope,
@@ -489,6 +531,9 @@ func NewArtifactManager(repo repositories.RepositoryInterface, store *storage.Da
489531
deleteResponseTime: labeled.NewStopWatch("delete_duration", "The duration of the delete artifact calls.", time.Millisecond, artifactScope, labeled.EmitUnlabeledMetric),
490532
deleteSuccessCounter: labeled.NewCounter("delete_success_count", "The number of times delete artifact succeeded", artifactScope, labeled.EmitUnlabeledMetric),
491533
deleteFailureCounter: labeled.NewCounter("delete_failure_count", "The number of times delete artifact failed", artifactScope, labeled.EmitUnlabeledMetric),
534+
bulkDeleteResponseTime: labeled.NewStopWatch("bulk_delete_duration", "The duration of the bulk delete artifacts calls.", time.Millisecond, artifactScope, labeled.EmitUnlabeledMetric),
535+
bulkDeleteSuccessCounter: labeled.NewCounter("bulk_delete_success_count", "The number of times bulk delete artifacts succeeded", artifactScope, labeled.EmitUnlabeledMetric),
536+
bulkDeleteFailureCounter: labeled.NewCounter("bulk_delete_failure_count", "The number of times bulk delete artifacts failed", artifactScope, labeled.EmitUnlabeledMetric),
492537
}
493538

494539
return &artifactManager{

0 commit comments

Comments
 (0)