Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion datacatalog/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ require (
go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.47.0
go.opentelemetry.io/otel v1.24.0
google.golang.org/grpc v1.62.1
google.golang.org/protobuf v1.34.1
gorm.io/driver/postgres v1.5.3
gorm.io/driver/sqlite v1.5.4
gorm.io/gorm v1.25.4
Expand Down Expand Up @@ -131,7 +132,6 @@ require (
google.golang.org/genproto v0.0.0-20240123012728-ef4313101c80 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20240123012728-ef4313101c80 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20240123012728-ef4313101c80 // indirect
google.golang.org/protobuf v1.34.1 // indirect
gopkg.in/inf.v0 v0.9.1 // indirect
gopkg.in/ini.v1 v1.66.4 // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect
Expand Down
36 changes: 36 additions & 0 deletions datacatalog/pkg/manager/impl/artifact_data_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import (
"context"
"fmt"

"google.golang.org/grpc/codes"

Expand All @@ -13,10 +14,14 @@
)

const artifactDataFile = "data.pb"
const futureDataFile = "future.pb"
const futureDataName = "future"

// ArtifactDataStore stores and retrieves ArtifactData values in a data.pb
type ArtifactDataStore interface {
PutData(ctx context.Context, artifact *datacatalog.Artifact, data *datacatalog.ArtifactData) (storage.DataReference, error)
PutFutureData(ctx context.Context, artifact *datacatalog.Artifact, djspec *core.DynamicJobSpec) (storage.DataReference, error)
GetFutureData(ctx context.Context, dataModel models.ArtifactData) (*core.DynamicJobSpec, error)
GetData(ctx context.Context, dataModel models.ArtifactData) (*core.Literal, error)
DeleteData(ctx context.Context, dataModel models.ArtifactData) error
}
Expand All @@ -31,6 +36,11 @@
return m.store.ConstructReference(ctx, m.storagePrefix, dataset.GetProject(), dataset.GetDomain(), dataset.GetName(), dataset.GetVersion(), artifact.GetId(), data.GetName(), artifactDataFile)
}

func (m *artifactDataStore) getFutureDataLocation(ctx context.Context, artifact *datacatalog.Artifact) (storage.DataReference, error) {
dataset := artifact.GetDataset()
return m.store.ConstructReference(ctx, m.storagePrefix, dataset.GetProject(), dataset.GetDomain(), dataset.GetName(), dataset.GetVersion(), artifact.GetId(), futureDataName, futureDataFile)
}

// Store marshalled data in data.pb under the storage prefix
func (m *artifactDataStore) PutData(ctx context.Context, artifact *datacatalog.Artifact, data *datacatalog.ArtifactData) (storage.DataReference, error) {
dataLocation, err := m.getDataLocation(ctx, artifact, data)
Expand All @@ -45,6 +55,20 @@
return dataLocation, nil
}

// Store marshalled future data in future.pb under the future storage prefix
func (m *artifactDataStore) PutFutureData(ctx context.Context, artifact *datacatalog.Artifact, djspec *core.DynamicJobSpec) (storage.DataReference, error) {
dataLocation, err := m.getFutureDataLocation(ctx, artifact)
if err != nil {
return "", errors.NewDataCatalogErrorf(codes.Internal, "Unable to generate data location %s, err %v", dataLocation.String(), err)
}

Check warning on line 63 in datacatalog/pkg/manager/impl/artifact_data_store.go

View check run for this annotation

Codecov / codecov/patch

datacatalog/pkg/manager/impl/artifact_data_store.go#L62-L63

Added lines #L62 - L63 were not covered by tests
err = m.store.WriteProtobuf(ctx, dataLocation, storage.Options{}, djspec)
if err != nil {
return "", errors.NewDataCatalogErrorf(codes.Internal, "Unable to store future artifact data in location %s, err %v", dataLocation.String(), err)
}

Check warning on line 67 in datacatalog/pkg/manager/impl/artifact_data_store.go

View check run for this annotation

Codecov / codecov/patch

datacatalog/pkg/manager/impl/artifact_data_store.go#L66-L67

Added lines #L66 - L67 were not covered by tests

return dataLocation, nil
}

// Retrieve the literal value of the ArtifactData from its specified location
func (m *artifactDataStore) GetData(ctx context.Context, dataModel models.ArtifactData) (*core.Literal, error) {
var value core.Literal
Expand All @@ -56,6 +80,18 @@
return &value, nil
}

// Retrieve the future data from the future storage prefix
func (m *artifactDataStore) GetFutureData(ctx context.Context, dataModel models.ArtifactData) (*core.DynamicJobSpec, error) {
var djSpec core.DynamicJobSpec
err := m.store.ReadProtobuf(ctx, storage.DataReference(dataModel.Location), &djSpec)
if err != nil {
return nil, errors.NewDataCatalogErrorf(codes.Internal, "Unable to read future artifact data from location %s, err %v", dataModel.Location, err)
}

Check warning on line 89 in datacatalog/pkg/manager/impl/artifact_data_store.go

View check run for this annotation

Codecov / codecov/patch

datacatalog/pkg/manager/impl/artifact_data_store.go#L88-L89

Added lines #L88 - L89 were not covered by tests
fmt.Printf("the dj spec: %+v\n", djSpec.GetNodes()[0])

return &djSpec, nil
}

// DeleteData removes the stored artifact data from the underlying blob storage
func (m *artifactDataStore) DeleteData(ctx context.Context, dataModel models.ArtifactData) error {
if err := m.store.Delete(ctx, storage.DataReference(dataModel.Location)); err != nil {
Expand Down
1 change: 1 addition & 0 deletions datacatalog/pkg/manager/impl/artifact_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,7 @@
m.systemMetrics.doesNotExistCounter.Inc(ctx)
} else {
logger.Errorf(ctx, "Unable to retrieve Artifact by tag %v, err: %v", key, err)
m.systemMetrics.getFailureCounter.Inc(ctx)

Check warning on line 211 in datacatalog/pkg/manager/impl/artifact_manager.go

View check run for this annotation

Codecov / codecov/patch

datacatalog/pkg/manager/impl/artifact_manager.go#L211

Added line #L211 was not covered by tests
}
return models.Artifact{}, err
}
Expand Down
Loading
Loading