Commit c2408631 authored by Suleimi Ahmed's avatar Suleimi Ahmed 🔴 Committed by Hayley Swimelar
Browse files

feat(notification): add blob download meta object to notifications

parent d05e93aa
Loading
Loading
Loading
Loading
+2 −1
Changes for blobs.go: 2 added lines, 1 removed line.
Original line number Diff line number Diff line
@@ -8,6 +8,7 @@ import (
	"net/http"
	"time"

	"github.com/docker/distribution/notifications/meta"
	"github.com/docker/distribution/reference"
	"github.com/opencontainers/go-digest"
	v1 "github.com/opencontainers/image-spec/specs-go/v1"
@@ -172,7 +173,7 @@ type BlobServer interface {
	// The implementation may serve the same blob from a different digest
	// domain. The appropriate headers will be set for the blob, unless they
	// have already been set by the caller.
	ServeBlob(ctx context.Context, w http.ResponseWriter, r *http.Request, dgst digest.Digest) error
	ServeBlob(ctx context.Context, w http.ResponseWriter, r *http.Request, dgst digest.Digest) (*meta.Blob, error)
}

// BlobIngester ingests blob data.
+21 −0
Changes for docs-gitlab/notifications/index.md: 21 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -37,6 +37,12 @@ Below is a sample payload of an event notification sent by the registry on diffe
      "source": {
        "addr": "127.0.0.1:5000",
        "instanceID": "45681f21-a006-42f2-ab7c-4cc37d8906b4"
      },
      "meta":{
        "blob":{
            "redirected": true,
            "storageBackend": "s3aws"
        }
      }
    }
  ]
@@ -55,6 +61,7 @@ Below is a sample payload of an event notification sent by the registry on diffe
| `request`   | Object | Yes            | See [`request`](#request).                                                                                                                                  |
| `actor`     | Object | Yes            | See [`actor`](#actor).                                                                                                                                      |
| `source`    | Object | Yes            | See [`source`](#source).                                                                                                                                    |
| `meta`      | Object | No             | Meta contains additional (optional) information related to an event. See [`meta`](#meta).                                                                                                                                      |

#### `target`

@@ -103,3 +110,17 @@ Below is a sample payload of an event notification sent by the registry on diffe
|--------------|--------|----------------|-----------------------------------------------------------------------------------------|
| `addr`       | String | No             | Contains the IP or hostname and the port of the registry node that generated the event. |
| `instanceID` | String | No             | Identifies a running instance of the registry node.                                     |


#### `meta`

| Field        | Type   | Always present | Description                                                                             |
|--------------|--------|----------------|-----------------------------------------------------------------------------------------|
| `blob`       | Object | No             | Only present on blob download events. Contains additional metadata on downloaded blobs. See [`blob`](#blob)                               |

#### `blob`

| Field             | Type    | Always present | Description                                                                                                               |
|-------------------|---------|----------------|---------------------------------------------------------------------------------------------------------------------------|
| `redirected`      | Boolean | Yes            | Identifies if a blob download request was served via a redirect url to the requesting client.                                                                                                         |
| `storageBackend`  | String  | Yes            | Identifies the backend that was used to serve a blob download request. This is always the configured storage backend (and not the redirect url provider) until  https://gitlab.com/gitlab-org/container-registry/-/issues/1003 is addressed.  |
+11 −7
Changes for notifications/bridge.go: 11 added lines, 7 removed lines.
Original line number Diff line number Diff line
@@ -6,6 +6,7 @@ import (

	"github.com/docker/distribution"
	"github.com/docker/distribution/context"
	"github.com/docker/distribution/notifications/meta"
	"github.com/docker/distribution/reference"
	"github.com/docker/distribution/uuid"
	"github.com/opencontainers/go-digest"
@@ -69,15 +70,15 @@ func (b *bridge) ManifestDeleted(repo reference.Named, dgst digest.Digest) error
}

func (b *bridge) BlobPushed(repo reference.Named, desc distribution.Descriptor) error {
	return b.createBlobEventAndWrite(EventActionPush, repo, desc)
	return b.createBlobEventAndWrite(EventActionPush, repo, desc, nil)
}

func (b *bridge) BlobPulled(repo reference.Named, desc distribution.Descriptor) error {
	return b.createBlobEventAndWrite(EventActionPull, repo, desc)
func (b *bridge) BlobPulled(repo reference.Named, desc distribution.Descriptor, eventMeta *meta.Blob) error {
	return b.createBlobEventAndWrite(EventActionPull, repo, desc, eventMeta)
}

func (b *bridge) BlobMounted(repo reference.Named, desc distribution.Descriptor, fromRepo reference.Named) error {
	event, err := b.createBlobEvent(EventActionMount, repo, desc)
	event, err := b.createBlobEvent(EventActionMount, repo, desc, nil)
	if err != nil {
		return err
	}
@@ -161,8 +162,8 @@ func (b *bridge) createBlobDeleteEventAndWrite(action string, repo reference.Nam
	return b.sink.Write(event)
}

func (b *bridge) createBlobEventAndWrite(action string, repo reference.Named, desc distribution.Descriptor) error {
	event, err := b.createBlobEvent(action, repo, desc)
func (b *bridge) createBlobEventAndWrite(action string, repo reference.Named, desc distribution.Descriptor, eventMeta *meta.Blob) error {
	event, err := b.createBlobEvent(action, repo, desc, eventMeta)
	if err != nil {
		return err
	}
@@ -170,11 +171,14 @@ func (b *bridge) createBlobEventAndWrite(action string, repo reference.Named, de
	return b.sink.Write(event)
}

func (b *bridge) createBlobEvent(action string, repo reference.Named, desc distribution.Descriptor) (*Event, error) {
func (b *bridge) createBlobEvent(action string, repo reference.Named, desc distribution.Descriptor, eventMeta *meta.Blob) (*Event, error) {
	event := b.createEvent(action)
	event.Target.Descriptor = desc
	event.Target.Length = desc.Size
	event.Target.Repository = repo.Name()
	if eventMeta != nil {
		event.Meta = map[string]Meta{"blob": eventMeta}
	}

	ref, err := reference.WithDigest(repo, desc.Digest)
	if err != nil {
+9 −0
Changes for notifications/event.go: 9 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -63,8 +63,17 @@ type Event struct {
	// differently, while the actor "initiates" the event, the source
	// "generates" it.
	Source SourceRecord `json:"source,omitempty"`

	// Meta is a map of key-value pairs that are loosely associated to an event.
	// This field should be used only to propagate non-common/highly-specific details
	// of an event in cases where the detail can not be communicated in existing attributes of the `Event` struct.
	// Meta is not guaranteed to be in every event or to have simiar keys across event types.
	Meta map[string]Meta `json:"meta,omitempty"`
}

// Meta is the event meta type
type Meta interface{}

// Target uniquely describes the target of the event.
type Target struct {
	// TODO(stevvooe): Use http.DetectContentType for layers, maybe.
+11 −8
Changes for notifications/listener.go: 11 added lines, 8 removed lines.
Original line number Diff line number Diff line
@@ -7,6 +7,7 @@ import (
	"github.com/docker/distribution"

	dcontext "github.com/docker/distribution/context"
	"github.com/docker/distribution/notifications/meta"
	"github.com/docker/distribution/reference"
	"github.com/opencontainers/go-digest"
)
@@ -21,7 +22,7 @@ type ManifestListener interface {
// BlobListener describes a listener that can respond to layer related events.
type BlobListener interface {
	BlobPushed(repo reference.Named, desc distribution.Descriptor) error
	BlobPulled(repo reference.Named, desc distribution.Descriptor) error
	BlobPulled(repo reference.Named, desc distribution.Descriptor, eventMeta *meta.Blob) error
	BlobMounted(repo reference.Named, desc distribution.Descriptor, fromRepo reference.Named) error
	BlobDeleted(repo reference.Named, desc digest.Digest) error
}
@@ -146,7 +147,8 @@ func (bsl *blobServiceListener) Get(ctx context.Context, dgst digest.Digest) ([]
		if desc, err := bsl.Stat(ctx, dgst); err != nil {
			dcontext.GetLogger(ctx).Errorf("error resolving descriptor in ServeBlob listener: %v", err)
		} else {
			if err := bsl.parent.listener.BlobPulled(bsl.parent.Repository.Named(), desc); err != nil {
			// for now we ignore the blob download meta when the blob is not downloaded via an API request
			if err := bsl.parent.listener.BlobPulled(bsl.parent.Repository.Named(), desc, nil); err != nil {
				dcontext.GetLogger(ctx).Errorf("error dispatching layer pull to listener: %v", err)
			}
		}
@@ -164,7 +166,8 @@ func (bsl *blobServiceListener) Open(ctx context.Context, dgst digest.Digest) (d
		if desc, err := bsl.Stat(ctx, dgst); err != nil {
			dcontext.GetLogger(ctx).Errorf("error resolving descriptor in ServeBlob listener: %v", err)
		} else {
			if err := bsl.parent.listener.BlobPulled(bsl.parent.Repository.Named(), desc); err != nil {
			// for now we ignore the blob download meta when the blob is not downloaded via an API request
			if err := bsl.parent.listener.BlobPulled(bsl.parent.Repository.Named(), desc, nil); err != nil {
				dcontext.GetLogger(ctx).Errorf("error dispatching layer pull to listener: %v", err)
			}
		}
@@ -173,22 +176,22 @@ func (bsl *blobServiceListener) Open(ctx context.Context, dgst digest.Digest) (d
	return rc, err
}

func (bsl *blobServiceListener) ServeBlob(ctx context.Context, w http.ResponseWriter, r *http.Request, dgst digest.Digest) error {
	err := bsl.BlobStore.ServeBlob(ctx, w, r, dgst)
func (bsl *blobServiceListener) ServeBlob(ctx context.Context, w http.ResponseWriter, r *http.Request, dgst digest.Digest) (*meta.Blob, error) {
	eventMeta, err := bsl.BlobStore.ServeBlob(ctx, w, r, dgst)
	if err == nil {
		if bsl.fsMirroringDisabled {
			return nil
			return eventMeta, nil
		}
		if desc, err := bsl.Stat(ctx, dgst); err != nil {
			dcontext.GetLogger(ctx).Errorf("error resolving descriptor in ServeBlob listener: %v", err)
		} else {
			if err := bsl.parent.listener.BlobPulled(bsl.parent.Repository.Named(), desc); err != nil {
			if err := bsl.parent.listener.BlobPulled(bsl.parent.Repository.Named(), desc, eventMeta); err != nil {
				dcontext.GetLogger(ctx).Errorf("error dispatching layer pull to listener: %v", err)
			}
		}
	}

	return err
	return eventMeta, err
}

func (bsl *blobServiceListener) Put(ctx context.Context, mediaType string, p []byte) (distribution.Descriptor, error) {
Loading