diff --git a/auth/basic_auth.go b/auth/basic_auth.go index 6253b10e..91e8e161 100644 --- a/auth/basic_auth.go +++ b/auth/basic_auth.go @@ -64,14 +64,14 @@ func (a *auth) BasicAuth() echo.MiddlewareFunc { printInMiddleware := true defer func() { if printInMiddleware { - a.logger.Log(ctx).Send() + a.logger.Log(ctx, nil) } }() if ctx.Request().RequestURI == "/v2/" { _, err := a.validateUser(username, password) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return false, ctx.NoContent(http.StatusUnauthorized) } @@ -87,12 +87,12 @@ func (a *auth) BasicAuth() echo.MiddlewareFunc { Message: "not authorised", Detail: nil, }) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + a.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return false, ctx.JSON(http.StatusForbidden, errMsg) } resp, err := a.validateUser(username, password) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return false, err } diff --git a/auth/jwt_middleware.go b/auth/jwt_middleware.go index 0f16653f..421f39f6 100644 --- a/auth/jwt_middleware.go +++ b/auth/jwt_middleware.go @@ -1,6 +1,7 @@ package auth import ( + "fmt" "net/http" "time" @@ -32,8 +33,7 @@ func (a *auth) JWT() echo.MiddlewareFunc { ErrorHandlerWithContext: func(err error, ctx echo.Context) error { // ErrorHandlerWithContext only logs the failing requtest ctx.Set(types.HandlerStartTime, time.Now()) - ctx.Set(types.HttpEndpointErrorKey, err.Error()) - a.logger.Log(ctx) + a.logger.Log(ctx, err) return ctx.NoContent(http.StatusUnauthorized) }, KeyFunc: middleware.DefaultJWTConfig.KeyFunc, @@ -51,7 +51,7 @@ func (a *auth) ACL() echo.MiddlewareFunc { return func(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) defer func() { - a.logger.Log(ctx) + a.logger.Log(ctx, nil) }() m := ctx.Request().Method @@ -61,13 +61,13 @@ func (a *auth) ACL() echo.MiddlewareFunc { token, ok := ctx.Get("user").(*jwt.Token) if !ok { - ctx.Set(types.HttpEndpointErrorKey, "ACL: unauthorized") + a.logger.Log(ctx, fmt.Errorf("ACL: unauthorized")) return ctx.NoContent(http.StatusUnauthorized) } claims, ok := token.Claims.(*Claims) if !ok { - ctx.Set(types.HttpEndpointErrorKey, "ACL: invalid claims") + a.logger.Log(ctx, fmt.Errorf("ACL: invalid claims")) return ctx.NoContent(http.StatusUnauthorized) } @@ -76,7 +76,7 @@ func (a *auth) ACL() echo.MiddlewareFunc { return hf(ctx) } - ctx.Set(types.HttpEndpointErrorKey, "ACL: username didn't match from token") + a.logger.Log(ctx, fmt.Errorf("ACL: username didn't match from token")) return ctx.NoContent(http.StatusUnauthorized) } } diff --git a/auth/signin.go b/auth/signin.go index 1d738619..a180f52c 100644 --- a/auth/signin.go +++ b/auth/signin.go @@ -2,6 +2,7 @@ package auth import ( "encoding/json" + "fmt" "net/http" "time" @@ -11,29 +12,22 @@ import ( func (a *auth) SignIn(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - a.logger.Log(ctx).Send() - }() - var user User + var user User if err := json.NewDecoder(ctx.Request().Body).Decode(&user); err != nil { return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), }) } if user.Email == "" && user.Username == "" { - errMsg := echo.Map{ - "error": "email and username cannot be empty, please provide at least one of them", - } - ctx.Set(types.HttpEndpointErrorKey, errMsg) + errMsg := fmt.Errorf("email and username cannot be empty, please provide at least one of them") + a.logger.Log(ctx, errMsg) return ctx.JSON(http.StatusBadRequest, errMsg) } if user.Password == "" { - errMsg := echo.Map{ - "error": "password cannot be empty", - } - ctx.Set(types.HttpEndpointErrorKey, errMsg) + errMsg := fmt.Errorf("password cannot be empty") + a.logger.Log(ctx, errMsg) return ctx.JSON(http.StatusBadRequest, errMsg) } @@ -41,7 +35,7 @@ func (a *auth) SignIn(ctx echo.Context) error { if user.Email != "" { if err := verifyEmail(user.Email); err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), }) @@ -54,15 +48,15 @@ func (a *auth) SignIn(ctx echo.Context) error { //bz, err := a.store.Get([]byte(key)) userFromDb, err := a.pgStore.GetUser(ctx.Request().Context(), key) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), }) } if !a.verifyPassword(userFromDb.Password, user.Password) { - errMsg := "invalid password" - ctx.Set(types.HttpEndpointErrorKey, errMsg) + errMsg := fmt.Errorf("invalid password") + a.logger.Log(ctx, errMsg) return ctx.JSON(http.StatusUnauthorized, errMsg) } @@ -74,12 +68,13 @@ func (a *auth) SignIn(ctx echo.Context) error { token, err := a.newToken(uu, tokenLife) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusInternalServerError, echo.Map{ "error": err.Error(), }) } + a.logger.Log(ctx, nil) return ctx.JSON(http.StatusOK, echo.Map{ "token": token, "expires_in": tokenLife, diff --git a/auth/signup.go b/auth/signup.go index 21eb770c..f8722ff0 100644 --- a/auth/signup.go +++ b/auth/signup.go @@ -140,14 +140,10 @@ func verifyPassword(password string) error { func (a *auth) SignUp(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - a.logger.Log(ctx).Send() - }() var u User - if err := json.NewDecoder(ctx.Request().Body).Decode(&u); err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), "message": "error decoding request body in sign-up", @@ -156,7 +152,7 @@ func (a *auth) SignUp(ctx echo.Context) error { _ = ctx.Request().Body.Close() if err := u.Validate(a.store); err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), }) @@ -164,7 +160,7 @@ func (a *auth) SignUp(ctx echo.Context) error { hpwd, err := a.hashPassword(u.Password) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusInternalServerError, echo.Map{ "error": err.Error(), }) @@ -179,12 +175,13 @@ func (a *auth) SignUp(ctx echo.Context) error { err = a.pgStore.AddUser(ctx.Request().Context(), newUser) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusInternalServerError, echo.Map{ "error": err.Error(), }) } + a.logger.Log(ctx, nil) return ctx.JSON(http.StatusCreated, echo.Map{ "message": "user successfully created", }) diff --git a/auth/token.go b/auth/token.go index beca3794..e372104c 100644 --- a/auth/token.go +++ b/auth/token.go @@ -18,21 +18,18 @@ func (a *auth) Token(ctx echo.Context) error { // TODO (jay-dee7) - check for all valid query params here like serive, client_id, offline_token, etc // more at this link - https://docs.docker.com/registry/spec/auth/token/ ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - a.logger.Log(ctx).Send() - }() authHeader := ctx.Request().Header.Get(AuthorizationHeaderKey) if authHeader != "" { username, password, err := a.getCredsFromHeader(ctx.Request()) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.NoContent(http.StatusUnauthorized) } creds, err := a.validateUser(username, password) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.JSON(http.StatusUnauthorized, echo.Map{ "error": err.Error(), }) @@ -47,7 +44,7 @@ func (a *auth) Token(ctx echo.Context) error { "error": err.Error(), "msg": "invalid scope provided", } - ctx.Set(types.HttpEndpointErrorKey, errMsg) + a.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSON(http.StatusBadRequest, errMsg) } @@ -56,7 +53,7 @@ func (a *auth) Token(ctx echo.Context) error { if len(scope.Actions) == 1 && scope.Actions["pull"] { token, err := a.newPublicPullToken() if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + a.logger.Log(ctx, err) return ctx.NoContent(http.StatusInternalServerError) } diff --git a/go.mod b/go.mod index 07505b3b..556d91cd 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,7 @@ require ( github.com/go-playground/validator/v10 v10.9.0 github.com/golang-jwt/jwt v3.2.2+incompatible github.com/google/uuid v1.3.0 + github.com/hashicorp/go-multierror v1.0.0 github.com/labstack/echo-contrib v0.11.0 github.com/labstack/echo/v4 v4.5.0 github.com/rs/zerolog v1.24.0 @@ -18,6 +19,8 @@ require ( golang.org/x/crypto v0.0.0-20210817164053-32db794688a5 ) +require github.com/hashicorp/errwrap v1.0.0 // indirect + require ( github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash v1.1.0 // indirect diff --git a/go.sum b/go.sum index c3be02ec..b6ef6851 100644 --- a/go.sum +++ b/go.sum @@ -266,10 +266,12 @@ github.com/hashicorp/consul/api v1.1.0/go.mod h1:VmuI/Lkw1nC05EYQWNKwWGbkg+FbDBt github.com/hashicorp/consul/api v1.3.0/go.mod h1:MmDNSzIMUjNpY/mQ398R4bk2FnqQLoPndWW5VkKPlCE= github.com/hashicorp/consul/sdk v0.1.1/go.mod h1:VKf9jXwCTEY1QZP2MOLRhb5i/I/ssyNV1vwHyQBF0x8= github.com/hashicorp/consul/sdk v0.3.0/go.mod h1:VKf9jXwCTEY1QZP2MOLRhb5i/I/ssyNV1vwHyQBF0x8= +github.com/hashicorp/errwrap v1.0.0 h1:hLrqtEDnRye3+sgx6z4qVLNuviH3MR5aQ0ykNJa/UYA= github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4= github.com/hashicorp/go-cleanhttp v0.5.1/go.mod h1:JpRdi6/HCYpAwUzNwuwqhbovhLtngrth3wmdIIUrZ80= github.com/hashicorp/go-immutable-radix v1.0.0/go.mod h1:0y9vanUI8NX6FsYoO3zeMjhV/C5i9g4Q3DwcSNZ4P60= github.com/hashicorp/go-msgpack v0.5.3/go.mod h1:ahLV/dePpqEmjfWmKiqvPkv/twdG7iPBM1vqhUKIvfM= +github.com/hashicorp/go-multierror v1.0.0 h1:iVjPR7a6H0tWELX5NxNe7bYopibicUzc7uPribsnS6o= github.com/hashicorp/go-multierror v1.0.0/go.mod h1:dHtQlpGsu+cZNNAkkCN/P3hoUDHhCYQXV3UM06sGGrk= github.com/hashicorp/go-rootcerts v1.0.0/go.mod h1:K6zTfqpRlCUIjkwsN4Z+hiSfzSTQa6eBIzfwKfwNnHU= github.com/hashicorp/go-sockaddr v1.0.0/go.mod h1:7Xibr9yA9JjQq1JpNB2Vw7kxv8xerXegt+ozgdvDeDU= diff --git a/main.go b/main.go index 89ed2fed..d2097c9c 100644 --- a/main.go +++ b/main.go @@ -44,7 +44,7 @@ func main() { os.Exit(1) } - logger := telemetry.ZLogger(telemetry.SetupLogger(), fluentBitCollector) + logger := telemetry.ZLogger(fluentBitCollector, cfg.Environment) authSvc := auth.New(localCache, cfg, pgStore, logger) skynetClient := skynet.NewClient(cfg) @@ -55,5 +55,5 @@ func main() { } router.Register(cfg, e, reg, authSvc, localCache, pgStore) - logger.Errorf("error initialising OpenRegistry Server: %s", e.Start(cfg.Registry.Address())) + color.Red("error initialising OpenRegistry Server: %s", e.Start(cfg.Registry.Address())) } diff --git a/registry/v2/blobs.go b/registry/v2/blobs.go index 47c2a020..fd4e5a24 100644 --- a/registry/v2/blobs.go +++ b/registry/v2/blobs.go @@ -33,9 +33,6 @@ func (b *blobs) errorResponse(code, msg string, detail map[string]interface{}) [ func (b *blobs) HEAD(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - b.registry.logger.Log(ctx).Send() - }() digest := ctx.Param("digest") @@ -45,8 +42,7 @@ func (b *blobs) HEAD(ctx echo.Context) error { "skynet": "layer not found", } errMsg := b.errorResponse(RegistryErrorCodeManifestBlobUnknown, err.Error(), details) - - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -57,12 +53,13 @@ func (b *blobs) HEAD(ctx echo.Context) error { "error": err.Error(), } errMsg := b.errorResponse(RegistryErrorCodeManifestBlobUnknown, "Manifest does not exist", details) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } ctx.Response().Header().Set("Content-Length", fmt.Sprintf("%d", metadata.ContentLength)) ctx.Response().Header().Set("Docker-Content-Digest", digest) + b.registry.logger.Log(ctx, nil) return ctx.String(http.StatusOK, "OK") } @@ -74,9 +71,6 @@ these will be part of the txn in StartUpload */ func (b *blobs) UploadBlob(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - b.registry.logger.Log(ctx).Send() - }() namespace := ctx.Param("username") + "/" + ctx.Param("imagename") contentRange := ctx.Request().Header.Get("Content-Range") @@ -89,8 +83,7 @@ func (b *blobs) UploadBlob(ctx echo.Context) error { "stream upload after first write are not allowed", nil, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } @@ -111,13 +104,14 @@ func (b *blobs) UploadBlob(ctx echo.Context) error { err.Error(), nil, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } locationHeader := fmt.Sprintf("/v2/%s/blobs/uploads/%s", namespace, uuid) ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Range", fmt.Sprintf("0-%d", len(buf.Bytes())-1)) + b.registry.logger.Log(ctx, nil) return ctx.NoContent(http.StatusAccepted) } @@ -129,13 +123,13 @@ func (b *blobs) UploadBlob(ctx echo.Context) error { "contentRange": contentRange, } errMsg := b.errorResponse(RegistryErrorCodeBlobUploadUnknown, err.Error(), details) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusRequestedRangeNotSatisfiable, errMsg) } if start != len(b.uploads[uuid]) { errMsg := b.errorResponse(RegistryErrorCodeBlobUploadUnknown, "content range mismatch", nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusRequestedRangeNotSatisfiable, errMsg) } @@ -147,7 +141,7 @@ func (b *blobs) UploadBlob(ctx echo.Context) error { "error while creating new buffer from existing blobs", nil, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusInternalServerError, errMsg) } // 10 ctx.Request().Body.Close() @@ -159,12 +153,13 @@ func (b *blobs) UploadBlob(ctx echo.Context) error { err.Error(), nil, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } locationHeader := fmt.Sprintf("/v2/%s/blobs/uploads/%s", namespace, uuid) ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Range", fmt.Sprintf("0-%d", buf.Len()-1)) + b.registry.logger.Log(ctx, nil) return ctx.NoContent(http.StatusAccepted) } diff --git a/registry/v2/helpers.go b/registry/v2/helpers.go index c2d486df..ad8b41b0 100644 --- a/registry/v2/helpers.go +++ b/registry/v2/helpers.go @@ -6,6 +6,8 @@ import ( "encoding/json" "fmt" "strings" + + "github.com/fatih/color" ) func (r *registry) errorResponse(code, msg string, detail map[string]interface{}) []byte { @@ -19,7 +21,7 @@ func (r *registry) errorResponse(code, msg string, detail map[string]interface{} bz, e := json.Marshal(err) if e != nil { - r.logger.Error(e.Error()) + color.Red("error marshalling error response: %w", err) } return bz diff --git a/registry/v2/registry.go b/registry/v2/registry.go index 70aa7fd4..46bd0de0 100644 --- a/registry/v2/registry.go +++ b/registry/v2/registry.go @@ -64,9 +64,6 @@ func (r *registry) LayerExists(ctx echo.Context) error { // OK func (r *registry) ManifestExists(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() namespace := ctx.Param("username") + "/" + ctx.Param("imagename") ref := ctx.Param("reference") // ref can be either tag or digest @@ -80,7 +77,7 @@ func (r *registry) ManifestExists(ctx echo.Context) error { } errMsg := r.errorResponse(RegistryErrorCodeManifestBlobUnknown, err.Error(), details) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -92,7 +89,7 @@ func (r *registry) ManifestExists(ctx echo.Context) error { } errMsg := r.errorResponse(RegistryErrorCodeManifestBlobUnknown, "Manifest does not exist", detail) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -102,9 +99,8 @@ func (r *registry) ManifestExists(ctx echo.Context) error { "foundDigest": manifest.Digest, "clientDigest": ref, } - r.logger.Error(details) - errMsg := r.errorResponse(RegistryErrorCodeManifestInvalid, "manifest digest does not match", nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + errMsg := r.errorResponse(RegistryErrorCodeManifestInvalid, "manifest digest does not match", details) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } @@ -112,7 +108,7 @@ func (r *registry) ManifestExists(ctx echo.Context) error { ctx.Response().Header().Set("Content-Type", "application/json") ctx.Response().Header().Set("Content-Length", fmt.Sprintf("%d", metadata.ContentLength)) ctx.Response().Header().Set("Docker-Content-Digest", manifest.Digest) - + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusOK) } @@ -121,9 +117,6 @@ func (r *registry) ManifestExists(ctx echo.Context) error { // OK func (r *registry) Catalog(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() queryParamPageSize := ctx.QueryParam("n") queryParamOffset := ctx.QueryParam("last") @@ -133,7 +126,7 @@ func (r *registry) Catalog(ctx echo.Context) error { if queryParamPageSize != "" { ps, err := strconv.ParseInt(ctx.QueryParam("n"), 10, 64) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + r.logger.Log(ctx, err) return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), }) @@ -144,7 +137,7 @@ func (r *registry) Catalog(ctx echo.Context) error { if queryParamOffset != "" { o, err := strconv.ParseInt(ctx.QueryParam("last"), 10, 64) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + r.logger.Log(ctx, err) return ctx.JSON(http.StatusBadRequest, echo.Map{ "error": err.Error(), }) @@ -154,18 +147,19 @@ func (r *registry) Catalog(ctx echo.Context) error { catalogList, err := r.store.GetCatalog(ctx.Request().Context(), namespace, pageSize, offset) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + r.logger.Log(ctx, err) return ctx.JSON(http.StatusInternalServerError, echo.Map{ "error": err.Error(), }) } total, err := r.store.GetCatalogCount(ctx.Request().Context()) if err != nil { - ctx.Set(types.HttpEndpointErrorKey, err.Error()) + r.logger.Log(ctx, err) return ctx.JSON(http.StatusInternalServerError, echo.Map{ "error": err.Error(), }) } + r.logger.Log(ctx, nil) return ctx.JSON(http.StatusOK, echo.Map{ "repositories": catalogList, "total": total, @@ -178,9 +172,6 @@ func (r *registry) Catalog(ctx echo.Context) error { // OK func (r *registry) ListTags(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() namespace := ctx.Param("username") + "/" + ctx.Param("imagename") limit := ctx.QueryParam("n") @@ -188,8 +179,7 @@ func (r *registry) ListTags(ctx echo.Context) error { tags, err := r.store.GetImageTags(ctx.Request().Context(), namespace) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeTagInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -197,8 +187,7 @@ func (r *registry) ListTags(ctx echo.Context) error { n, err := strconv.ParseInt(limit, 10, 32) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeTagInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } if n > 0 { @@ -209,6 +198,7 @@ func (r *registry) ListTags(ctx echo.Context) error { } } + r.logger.Log(ctx, nil) return ctx.JSON(http.StatusOK, echo.Map{ "name": namespace, "tags": tags, @@ -222,26 +212,28 @@ func (r *registry) List(ctx echo.Context) error { // GET /v2//manifests/ // OK func (r *registry) PullManifest(ctx echo.Context) error { + ctx.Set(types.HandlerStartTime, time.Now()) + namespace := ctx.Param("username") + "/" + ctx.Param("imagename") ref := ctx.Param("reference") manifest, err := r.store.GetManifestByReference(ctx.Request().Context(), namespace, ref) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeManifestUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } resp, err := r.skynet.Download(manifest.Skylink) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeManifestInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } bz, err := io.ReadAll(resp) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeManifestInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } _ = resp.Close() @@ -249,6 +241,7 @@ func (r *registry) PullManifest(ctx echo.Context) error { ctx.Response().Header().Set("X-Docker-Content-ID", manifest.Skylink) ctx.Response().Header().Set("Content-Type", manifest.MediaType) ctx.Response().Header().Set("Content-Length", fmt.Sprintf("%d", len(bz))) + r.logger.Log(ctx, nil) return ctx.JSONBlob(http.StatusOK, bz) } @@ -258,18 +251,14 @@ func (r *registry) PullManifest(ctx echo.Context) error { func (r *registry) PullLayer(ctx echo.Context) error { //namespace := ctx.Param("username") + "/" + ctx.Param("imagename") ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() clientDigest := ctx.Param("digest") layer, err := r.store.GetLayer(ctx.Request().Context(), clientDigest) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) - } if layer.SkynetLink == "" { @@ -278,7 +267,7 @@ func (r *registry) PullLayer(ctx echo.Context) error { } e := fmt.Errorf("skylink is empty").Error() errMsg := r.errorResponse(RegistryErrorCodeBlobUnknown, e, detail) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -289,13 +278,13 @@ func (r *registry) PullLayer(ctx echo.Context) error { "skylink": layer.SkynetLink, } errMsg := r.errorResponse(RegistryErrorCodeBlobUnknown, err.Error(), detail) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } buf := &bytes.Buffer{} if _, err := io.Copy(buf, resp); err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusInternalServerError, errMsg) } _ = resp.Close() @@ -311,59 +300,59 @@ func (r *registry) PullLayer(ctx echo.Context) error { "client digest is different than computed digest", details, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } ctx.Response().Header().Set("Content-Length", fmt.Sprintf("%d", len(buf.Bytes()))) ctx.Response().Header().Set("Docker-Content-Digest", dig) + r.logger.Log(ctx, nil) return ctx.Blob(http.StatusOK, "application/octet-stream", buf.Bytes()) } // MonolithicUpload // PUT /v2//blobs/uploads/?digest= func (r *registry) MonolithicUpload(ctx echo.Context) error { + ctx.Set(types.HandlerStartTime, time.Now()) + namespace := ctx.Param("username") + "/" + ctx.Param("imagename") uuid := ctx.Param("uuid") - digest := ctx.QueryParam("digest") - buf := &bytes.Buffer{} - if _, err := io.Copy(buf, ctx.Request().Body); err != nil { - errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + if _, ok := r.b.uploads[uuid]; ok { + errMsg := r.b.errorResponse( + RegistryErrorCodeBlobUploadInvalid, + "error in monolithic upload", + nil, + ) + r.b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } - _ = ctx.Request().Body.Close() - link, err := r.skynet.Upload(namespace, digest, buf.Bytes(), true) - if err != nil { - detail := echo.Map{ - "error": err.Error(), - "caller": "MonolithicUpload", - } - errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), detail) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - return ctx.JSONBlob(http.StatusInternalServerError, buf.Bytes()) + buf := &bytes.Buffer{} + if _, err := io.Copy(buf, ctx.Request().Body); err != nil { + return ctx.JSON(http.StatusBadRequest, echo.Map{ + "error": err.Error(), + "message": "error copying request body in monolithic upload blob", + }) } - metadata := types.Metadata{ - Namespace: namespace, - Manifest: types.ImageManifest{ - SchemaVersion: 2, - MediaType: "", - Layers: []*types.Layer{{MediaType: "", Size: len(buf.Bytes()), Digest: digest, SkynetLink: link, UUID: uuid}}, - }, - } + _ = ctx.Request().Body.Close() + r.b.uploads[uuid] = buf.Bytes() - err = r.localCache.Update([]byte(namespace), metadata.Bytes()) - if err != nil { - errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + if err := r.b.blobTransaction(ctx, buf.Bytes(), uuid); err != nil { + errMsg := r.b.errorResponse( + RegistryErrorCodeBlobUploadInvalid, + err.Error(), + nil, + ) + r.b.registry.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } - locationHeader := link + locationHeader := fmt.Sprintf("/v2/%s/blobs/uploads/%s", namespace, uuid) + ctx.Response().Header().Set("Location", locationHeader) + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusCreated) } @@ -381,9 +370,6 @@ registry.tnxMap[uuid] = {txn,blobs[],timeout} // POST /v2//blobs/uploads/ func (r *registry) StartUpload(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() namespace := ctx.Param("username") + "/" + ctx.Param("imagename") clientDigest := ctx.QueryParam("digest") @@ -400,8 +386,7 @@ func (r *registry) StartUpload(ctx echo.Context) error { "error while reading request body", details, ) - - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } _ = ctx.Request().Body.Close() // why defer? body is already read :) @@ -417,16 +402,14 @@ func (r *registry) StartUpload(ctx echo.Context) error { "client digest does not meet computed digest", details, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + r.logger.Log(ctx, fmt.Errorf("%s", fmt.Errorf("%s", errMsg))) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } skylink, err := r.skynet.Upload(namespace, dig, buf.Bytes(), true) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusRequestedRangeNotSatisfiable, errMsg) } @@ -442,24 +425,24 @@ func (r *registry) StartUpload(ctx echo.Context) error { txnOp, err := r.store.NewTxn(ctx.Request().Context()) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusInternalServerError, errMsg) } if err := r.store.SetLayer(ctx.Request().Context(), txnOp, layerV2); err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } if err := r.store.Commit(ctx.Request().Context(), txnOp); err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } link := r.getHttpUrlFromSkylink(skylink) ctx.Response().Header().Set("Location", link) + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusCreated) } @@ -472,7 +455,7 @@ func (r *registry) StartUpload(ctx echo.Context) error { err.Error(), nil, ) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusInternalServerError, errMsg) } r.txnMap[id.String()] = TxnStore{ @@ -484,15 +467,13 @@ func (r *registry) StartUpload(ctx echo.Context) error { ctx.Response().Header().Set("Content-Length", "0") ctx.Response().Header().Set("Docker-Upload-UUID", id.String()) ctx.Response().Header().Set("Range", fmt.Sprintf("0-%d", 0)) - + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusAccepted) } +//UploadProgress TODO func (r *registry) UploadProgress(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() namespace := ctx.Param("username") + "/" + ctx.Param("imagename") uuid := ctx.Param("uuid") @@ -503,7 +484,7 @@ func (r *registry) UploadProgress(ctx echo.Context) error { ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Range", "bytes=0-0") ctx.Response().Header().Set("Docker-Upload-UUID", uuid) - + r.logger.Log(ctx, err) return ctx.NoContent(http.StatusNoContent) } @@ -513,7 +494,7 @@ func (r *registry) UploadProgress(ctx echo.Context) error { ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Range", "bytes=0-0") ctx.Response().Header().Set("Docker-Upload-UUID", uuid) - + r.logger.Log(ctx, err) return ctx.NoContent(http.StatusNoContent) } @@ -521,7 +502,7 @@ func (r *registry) UploadProgress(ctx echo.Context) error { ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Range", fmt.Sprintf("bytes=0-%d", metadata.ContentLength)) ctx.Response().Header().Set("Docker-Upload-UUID", uuid) - + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusNoContent) } @@ -533,6 +514,8 @@ and inserted in the blob table thus committing the txn */ func (r *registry) CompleteUpload(ctx echo.Context) error { + ctx.Set(types.HandlerStartTime, time.Now()) + dig := ctx.QueryParam("digest") namespace := ctx.Param("username") + "/" + ctx.Param("imagename") id := ctx.Param("uuid") @@ -540,7 +523,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { buf := &bytes.Buffer{} if _, err := io.Copy(buf, ctx.Request().Body); err != nil { errMsg := r.errorResponse(RegistryErrorCodeDigestInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } _ = ctx.Request().Body.Close() @@ -555,7 +538,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { "headerDigest": dig, "serverSideDigest": ourHash, "bodyDigest": digest(buf.Bytes()), } errMsg := r.errorResponse(RegistryErrorCodeDigestInvalid, "digest mismatch", details) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } @@ -563,7 +546,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { skylink, err := r.skynet.Upload(blobNamespace, dig, ubuf.Bytes(), true) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUploadInvalid, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusRequestedRangeNotSatisfiable, errMsg) } @@ -578,7 +561,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { } if !ok { errMsg := r.errorResponse(RegistryErrorCodeUnknown, "transaction does not exist for uuid -"+id, nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } @@ -586,7 +569,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), echo.Map{ "error_detail": "set layer issues", }) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } @@ -594,7 +577,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), echo.Map{ "error_detail": "commitment issue", }) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } delete(r.txnMap, id) @@ -603,6 +586,7 @@ func (r *registry) CompleteUpload(ctx echo.Context) error { ctx.Response().Header().Set("Content-Length", "0") ctx.Response().Header().Set("Docker-Content-Digest", ourHash) ctx.Response().Header().Set("Location", locationHeader) + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusCreated) } @@ -617,13 +601,11 @@ func (r *registry) PushImage(ctx echo.Context) error { } func (r *registry) PushManifest(ctx echo.Context) error { + ctx.Set(types.HandlerStartTime, time.Now()) + namespace := ctx.Param("username") + "/" + ctx.Param("imagename") ref := ctx.Param("reference") contentType := ctx.Request().Header.Get("Content-Type") - ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() var manifest ImageManifest @@ -640,7 +622,7 @@ func (r *registry) PushManifest(ctx echo.Context) error { err = json.Unmarshal(buf.Bytes(), &manifest) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } dig := digest(buf.Bytes()) @@ -649,7 +631,7 @@ func (r *registry) PushManifest(ctx echo.Context) error { skylink, err := r.skynet.Upload(mfNamespace, dig, buf.Bytes(), true) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeManifestBlobUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -682,21 +664,21 @@ func (r *registry) PushManifest(ctx echo.Context) error { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), echo.Map{ "reason": "PG_ERR_CREATE_NEW_TXN", }) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) _ = r.store.Abort(ctx.Request().Context(), txnOp) return ctx.JSONBlob(http.StatusInternalServerError, errMsg) } if err := r.store.SetManifest(ctx.Request().Context(), txnOp, val); err != nil { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) _ = r.store.Abort(ctx.Request().Context(), txnOp) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } if err := r.store.SetConfig(ctx.Request().Context(), txnOp, mfc); err != nil { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) _ = r.store.Abort(ctx.Request().Context(), txnOp) return ctx.JSONBlob(http.StatusBadRequest, errMsg) } @@ -705,7 +687,7 @@ func (r *registry) PushManifest(ctx echo.Context) error { errMsg := r.errorResponse(RegistryErrorCodeUnknown, err.Error(), echo.Map{ "reason": "ERR_PG_COMMIT_TXN", }) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) _ = r.store.Abort(ctx.Request().Context(), txnOp) return ctx.JSONBlob(http.StatusInternalServerError, errMsg) } @@ -714,6 +696,7 @@ func (r *registry) PushManifest(ctx echo.Context) error { ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Docker-Content-Digest", dig) ctx.Response().Header().Set("X-Docker-Content-ID", skylink) + r.logger.Log(ctx, nil) return ctx.String(http.StatusCreated, "Created") } @@ -721,9 +704,6 @@ func (r *registry) PushManifest(ctx echo.Context) error { // POST /v2//blobs/uploads/ func (r *registry) PushLayer(ctx echo.Context) error { ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() elem := strings.Split(ctx.Request().URL.Path, "/") elem = elem[1:] @@ -733,8 +713,7 @@ func (r *registry) PushLayer(ctx echo.Context) error { // Must have a path of form /v2/{name}/blobs/{upload,sha256:} if len(elem) < 4 { errMsg := r.errorResponse(RegistryErrorCodeNameInvalid, "blobs must be attached to a repo", nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) - + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } @@ -744,7 +723,7 @@ func (r *registry) PushLayer(ctx echo.Context) error { ctx.Response().Header().Set("Location", locationHeader) ctx.Response().Header().Set("Docker-Upload-UUID", id.String()) ctx.Response().Header().Set("Range", "bytes=0-0") - + r.logger.Log(ctx, nil) return ctx.NoContent(http.StatusAccepted) } @@ -755,6 +734,8 @@ func (r *registry) CancelUpload(ctx echo.Context) error { // DeleteTagOrManifest // DELETE /v2//manifest/ or func (r *registry) DeleteTagOrManifest(ctx echo.Context) error { + ctx.Set(types.HandlerStartTime, time.Now()) + namespace := ctx.Param("username") + "/" + ctx.Param("imagename") ref := ctx.Param("reference") @@ -772,31 +753,28 @@ func (r *registry) DeleteTagOrManifest(ctx echo.Context) error { "digest": ref, } errMsg := r.errorResponse(RegistryErrorCodeManifestUnknown, err.Error(), details) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } - _ = r.store.Commit(ctx.Request().Context(), txnOp) + err := r.store.Commit(ctx.Request().Context(), txnOp) + r.logger.Log(ctx, err) return ctx.NoContent(http.StatusAccepted) } func (r *registry) DeleteLayer(ctx echo.Context) error { - //namespace := ctx.Param("username") + "/" + ctx.Param("imagename") - dig := ctx.Param("digest") - - //var m types.Metadata + ctx.Set(types.HandlerStartTime, time.Now()) + dig := ctx.Param("digest") layer, err := r.store.GetLayer(ctx.Request().Context(), dig) - //_, err := r.localCache.GetDigest(dig) if err != nil { errMsg := r.errorResponse(RegistryErrorCodeBlobUnknown, err.Error(), nil) - ctx.Set(types.HttpEndpointErrorKey, errMsg) + r.logger.Log(ctx, fmt.Errorf("%s", errMsg)) return ctx.JSONBlob(http.StatusNotFound, errMsg) } blobs := layer.BlobDigests - //err = r.localCache.DeleteLayer(namespace, dig) txnOp, _ := r.store.NewTxn(context.Background()) err = r.store.DeleteLayerV2(ctx.Request().Context(), txnOp, dig) if err != nil { @@ -807,47 +785,42 @@ func (r *registry) DeleteLayer(ctx echo.Context) error { bz, err := json.Marshal(logMsg) if err == nil { - ctx.Set(types.HttpEndpointErrorKey, logMsg) + r.logger.Log(ctx, err) } return ctx.JSONBlob(http.StatusInternalServerError, bz) } for i := range blobs { - //if err = r.localCache.DeleteDigest(dig); err != nil { if err = r.store.DeleteBlobV2(ctx.Request().Context(), txnOp, blobs[i]); err != nil { logMsg := echo.Map{ "error": err.Error(), "caller": "DeleteLayer", } - ctx.Set(types.HttpEndpointErrorKey, logMsg) + r.logger.Log(ctx, fmt.Errorf("%s", logMsg)) bz, err := json.Marshal(logMsg) if err != nil { r.log.Err(err).Send() } - return ctx.JSONBlob(http.StatusInternalServerError, bz) } } - _ = r.store.Commit(ctx.Request().Context(), txnOp) + err = r.store.Commit(ctx.Request().Context(), txnOp) + r.logger.Log(ctx, err) return ctx.NoContent(http.StatusAccepted) } // Should also look into 401 Code // https://docs.docker.com/registry/spec/api/ func (r *registry) ApiVersion(ctx echo.Context) error { - ctx.Set(types.HandlerStartTime, time.Now()) - defer func() { - r.logger.Log(ctx).Send() - }() ctx.Response().Header().Set(HeaderDockerDistributionApiVersion, "registry/2.0") - return ctx.String(http.StatusOK, "OK\n") } func (r *registry) GetImageNamespace(ctx echo.Context) error { + searchQuery := ctx.QueryParam("search_query") if searchQuery == "" { return ctx.JSON(http.StatusBadRequest, echo.Map{ diff --git a/telemetry/consoleWriter.go b/telemetry/consoleWriter.go new file mode 100644 index 00000000..0b09cb9c --- /dev/null +++ b/telemetry/consoleWriter.go @@ -0,0 +1,76 @@ +package telemetry + +import ( + "bytes" + "os" + "strings" + "time" + + "github.com/fatih/color" + "github.com/hashicorp/go-multierror" + "github.com/labstack/echo/v4" + "github.com/rs/zerolog" + "github.com/rs/zerolog/log" +) + +func (l logger) consoleWriter(ctx echo.Context, errMsg error) { + l.zlog = log.Output(zerolog.ConsoleWriter{Out: os.Stdout, TimeFormat: time.RFC822}) + l.zlog = l.zlog.With().Logger() + + buf := l.pool.Get().(*bytes.Buffer) + buf.Reset() + defer l.pool.Put(buf) + + req := ctx.Request() + res := ctx.Response() + + status := res.Status + level := zerolog.InfoLevel + switch { + case status >= 500: + level = zerolog.ErrorLevel + case status >= 400: + level = zerolog.WarnLevel + case status >= 300: + level = zerolog.ErrorLevel + } + + var e error + + _, err := buf.WriteString(req.Method + " ") + e = multierror.Append(e, err) + + _, err = buf.WriteString(color.GreenString("%d ", res.Status)) + e = multierror.Append(e, err) + + if level == zerolog.ErrorLevel { + e = multierror.Append(e, err) + } + if level == zerolog.WarnLevel { + e = multierror.Append(e, err) + } + + _, err = buf.WriteString(req.Host) + e = multierror.Append(e, err) + + _, err = buf.WriteString(req.RequestURI + " ") + e = multierror.Append(e, err) + + _, err = buf.WriteString(req.Proto + " ") + e = multierror.Append(e, err) + + _, err = buf.WriteString(req.UserAgent() + " ") + e = multierror.Append(e, err) + + if errMsg != nil { + _, err = buf.WriteString(color.YellowString(" %s", errMsg)) + e = multierror.Append(e, err) + } + + merr := e.(*multierror.Error) + if merr.ErrorOrNil() != nil { + buf.WriteString(strings.TrimSpace(merr.Error())) + } + + l.zlog.WithLevel(level).Msg(buf.String()) +} diff --git a/telemetry/log.go b/telemetry/log.go index 9b262111..6126938e 100644 --- a/telemetry/log.go +++ b/telemetry/log.go @@ -2,7 +2,6 @@ package telemetry import ( "bytes" - "encoding/json" "fmt" "io" "os" @@ -11,150 +10,34 @@ import ( "sync" "time" + "github.com/containerish/OpenRegistry/config" + fluentbit "github.com/containerish/OpenRegistry/telemetry/fluent-bit" - "github.com/containerish/OpenRegistry/types" "github.com/labstack/echo/v4" - "github.com/labstack/gommon/log" "github.com/rs/zerolog" "github.com/valyala/fasttemplate" ) type Logger interface { - echo.Logger - Log(ctx echo.Context) *zerolog.Event + Log(ctx echo.Context, err error) } -func SetupLogger() zerolog.Logger { +func SetupLogger(env string) zerolog.Logger { zerolog.TimeFieldFormat = zerolog.TimeFormatUnix - zerolog.SetGlobalLevel(zerolog.DebugLevel) - l := zerolog.New(os.Stdout) l = l.With().Caller().Logger() - return l -} - -//nolint:cyclop // insane amount of complexity because of templating -func ZerologMiddleware(baseLogger zerolog.Logger, fluentbitClient fluentbit.FluentBit) echo.MiddlewareFunc { - return func(hf echo.HandlerFunc) echo.HandlerFunc { - pool := &sync.Pool{ - New: func() interface{} { - return bytes.NewBuffer(make([]byte, 256)) - }, - } - - logFmt := `{"time":"${time_rfc3339}","x_request_id":"${request_id}","remote_ip":"${remote_ip}",` + - `"host":"${host}","method":"${method}","uri":"${uri}","user_agent":"${user_agent}",` + - `"status":${status},"error":"${error}","latency":${latency},"latency_human":"${latency_human}"` + - `,"bytes_in":${bytes_in},"bytes_out":${bytes_out}}` + "\n" - - template := fasttemplate.New(logFmt, "${", "}") - - return func(ctx echo.Context) (err error) { - buf := pool.Get().(*bytes.Buffer) - buf.Reset() - defer pool.Put(buf) - - req := ctx.Request() - res := ctx.Response() - start := time.Now() - if err = hf(ctx); err != nil { - ctx.Error(err) - } - stop := time.Now() - - var level zerolog.Level - if _, err = template.ExecuteFunc(buf, func(_ io.Writer, tag string) (int, error) { - switch tag { - case "time_rfc3339": - return buf.WriteString(time.Now().Format(time.RFC3339)) - case "request_id": - id := req.Header.Get(echo.HeaderXRequestID) - if id == "" { - id = res.Header().Get(echo.HeaderXRequestID) - } - return buf.WriteString(id) - case "remote_ip": - return buf.WriteString(ctx.RealIP()) - case "host": - return buf.WriteString(req.Host) - case "uri": - return buf.WriteString(req.RequestURI) - case "method": - return buf.WriteString(req.Method) - case "path": - p := req.URL.Path - if p == "" { - p = "/" - } - return buf.WriteString(p) - case "protocol": - return buf.WriteString(req.Proto) - case "referer": - return buf.WriteString(req.Referer()) - case "user_agent": - return buf.WriteString(req.UserAgent()) - case "status": - status := res.Status - level = zerolog.InfoLevel - switch { - case status >= 500: - level = zerolog.ErrorLevel - case status >= 400: - level = zerolog.WarnLevel - case status >= 300: - level = zerolog.ErrorLevel - } - - return buf.WriteString(strconv.FormatInt(int64(status), 10)) - case "error": - if err != nil { - // Error may contain invalid JSON e.g. `"` - b, _ := json.Marshal(err.Error()) - b = b[1 : len(b)-1] - return buf.Write(b) - } - - if ctxErr, ok := ctx.Get(types.HttpEndpointErrorKey).([]byte); ok { - return buf.Write(ctxErr) - } - case "latency": - l := stop.Sub(start) - return buf.WriteString(strconv.FormatInt(int64(l), 10)) - case "latency_human": - return buf.WriteString(stop.Sub(start).String()) - case "bytes_in": - cl := req.Header.Get(echo.HeaderContentLength) - if cl == "" { - cl = "0" - } - return buf.WriteString(cl) - case "bytes_out": - return buf.WriteString(strconv.FormatInt(res.Size, 10)) - default: - switch { - case strings.HasPrefix(tag, "header:"): - return buf.Write([]byte(ctx.Request().Header.Get(tag[7:]))) - case strings.HasPrefix(tag, "query:"): - return buf.Write([]byte(ctx.QueryParam(tag[6:]))) - case strings.HasPrefix(tag, "form:"): - return buf.Write([]byte(ctx.FormValue(tag[5:]))) - case strings.HasPrefix(tag, "cookie:"): - if cookie, cookieErr := ctx.Cookie(tag[7:]); cookieErr == nil { - return buf.Write([]byte(cookie.Value)) - } - } - } - return 0, nil - }); err != nil { - return - } - - bz := bytes.TrimSpace(buf.Bytes()) - baseLogger.WithLevel(level).RawJSON("msg", bz).Send() - fluentbitClient.Send(bz) - return + if env != config.Prod { + zerolog.SetGlobalLevel(zerolog.TraceLevel) + consoleWriter := zerolog.ConsoleWriter{ + Out: os.Stdout, + NoColor: false, + TimeFormat: time.RFC3339, } + l.Output(consoleWriter) + return l } + zerolog.SetGlobalLevel(zerolog.DebugLevel) + return l } type logger struct { @@ -163,9 +46,10 @@ type logger struct { pool *sync.Pool template *fasttemplate.Template zlog zerolog.Logger + env string } -func ZLogger(baseLogger zerolog.Logger, fluentbitClient fluentbit.FluentBit) Logger { +func ZLogger(fluentbitClient fluentbit.FluentBit, env string) Logger { pool := &sync.Pool{ New: func() interface{} { return bytes.NewBuffer(make([]byte, 256)) @@ -176,17 +60,25 @@ func ZLogger(baseLogger zerolog.Logger, fluentbitClient fluentbit.FluentBit) Log `"status":${status},"error":"${error}","latency":${latency},"latency_human":"${latency_human}"` + `,"bytes_in":${bytes_in},"bytes_out":${bytes_out}}` + "\n" + baseLogger := SetupLogger(env) + return &logger{ zlog: baseLogger, fluentBit: fluentbitClient, output: os.Stdout, pool: pool, template: fasttemplate.New(logFmt, "${", "}"), + env: env, } } //nolint:cyclop // insane amount of complexity because of templating -func (l logger) Log(ctx echo.Context) *zerolog.Event { +func (l logger) Log(ctx echo.Context, errMsg error) { + + if l.env != config.Prod { + l.consoleWriter(ctx, errMsg) + return + } start, ok := ctx.Get("start").(time.Time) if !ok { @@ -202,6 +94,10 @@ func (l logger) Log(ctx echo.Context) *zerolog.Event { req := ctx.Request() res := ctx.Response() + if errMsg != nil { + buf.WriteString(errMsg.Error()) + } + var level zerolog.Level if _, err := l.template.ExecuteFunc(buf, func(_ io.Writer, tag string) (int, error) { switch tag { @@ -246,10 +142,6 @@ func (l logger) Log(ctx echo.Context) *zerolog.Event { } return buf.WriteString(strconv.FormatInt(int64(status), 10)) - case "error": - if ctxErr, ok := ctx.Get(types.HttpEndpointErrorKey).([]byte); ok { - return buf.Write(ctxErr) - } case "latency": l := stop.Sub(start) return buf.WriteString(strconv.FormatInt(int64(l), 10)) @@ -284,118 +176,14 @@ func (l logger) Log(ctx echo.Context) *zerolog.Event { bz := bytes.TrimSpace(buf.Bytes()) l.fluentBit.Send(bz) - return l.zlog.WithLevel(level).RawJSON("msg", bz) -} - -func (l logger) Output() io.Writer { - return l.output -} - -func (l logger) SetOutput(w io.Writer) { - l.zlog.Output(w) + l.zlog.WithLevel(level).RawJSON("msg", bz).Send() } -// Prefix is not being used since zerologger is the only logger being used -func (l logger) Prefix() string { - return "" -} - -// SetPrefix is not being used since zerologger is the only logger being used -func (l logger) SetPrefix(p string) {} - -func (l logger) Level() log.Lvl { - level := l.zlog.GetLevel() - return log.Lvl(level + 1) -} - -// SetLevel - echo.loglvl starts from 1 while zerologger starts from 0 -func (l logger) SetLevel(v log.Lvl) { - l.zlog.Level(zerolog.Level(v - 1)) -} - -func (l logger) SetHeader(h string) { -} - -func (l logger) Print(i ...interface{}) { - l.zlog.WithLevel(l.zlog.GetLevel()).Msgf("%v", i...) -} - -func (l logger) Printf(format string, args ...interface{}) { - l.zlog.WithLevel(l.zlog.GetLevel()).Msgf(format, args...) -} - -func (l logger) Printj(j log.JSON) { - l.zlog.WithLevel(l.zlog.GetLevel()).Fields(j).Send() -} - -func (l logger) Debug(i ...interface{}) { - l.zlog.Debug().Msgf("%v", i...) -} - -func (l logger) Debugf(format string, args ...interface{}) { - l.zlog.Debug().Msgf(format, args...) -} - -func (l logger) Debugj(j log.JSON) { - l.zlog.Debug().Fields(j).Send() -} - -func (l logger) Info(i ...interface{}) { - l.zlog.Info().Msgf("%v", i...) -} - -func (l logger) Infof(format string, args ...interface{}) { - l.zlog.Info().Msgf(format, args...) -} - -func (l logger) Infoj(j log.JSON) { - l.zlog.Info().Fields(j).Send() -} - -func (l logger) Warn(i ...interface{}) { - l.zlog.Warn().Msgf("%v", i...) -} - -func (l logger) Warnf(format string, args ...interface{}) { - l.zlog.Warn().Msgf(format, args...) -} - -func (l logger) Warnj(j log.JSON) { - l.zlog.Warn().Fields(j).Send() -} - -func (l logger) Error(i ...interface{}) { - l.zlog.Error().Msgf("%v", i...) -} - -func (l logger) Errorf(format string, args ...interface{}) { - l.zlog.Error().Msgf(format, args...) -} - -func (l logger) Errorj(j log.JSON) { - l.zlog.Error().Fields(j).Send() -} - -func (l logger) Fatal(i ...interface{}) { - l.zlog.Fatal().Msgf("%v", i...) -} - -func (l logger) Fatalj(j log.JSON) { - l.zlog.Fatal().Fields(j).Send() -} - -func (l logger) Fatalf(format string, args ...interface{}) { - l.zlog.Fatal().Msgf(format, args...) -} - -func (l logger) Panic(i ...interface{}) { - l.zlog.Panic().Msgf("%v", i...) -} - -func (l logger) Panicj(j log.JSON) { - l.zlog.Panic().Fields(j).Send() -} - -func (l logger) Panicf(format string, args ...interface{}) { - l.zlog.Panic().Msgf(format, args...) -} +//func (l *logger) sanitizeErrors(errors ...error) { +// var e error +// for _, err := range errors { +// if err != nil { +// e = fmt.Errorf("%w", err) +// } +// } +//}