[full-ci] chore: bump reva to v2.44.0 update opencloud version 6.2.0 (#2734)
This commit is contained in:
Generated
Vendored
+10
-4
@@ -182,9 +182,12 @@ func (s *service) GetUser(ctx context.Context, req *userpb.GetUserRequest) (*use
|
||||
user, err := s.usermgr.GetUser(ctx, req.UserId, req.SkipFetchingUserGroups)
|
||||
if err != nil {
|
||||
res := &userpb.GetUserResponse{}
|
||||
if _, ok := err.(errtypes.NotFound); ok {
|
||||
switch err.(type) {
|
||||
case errtypes.NotFound:
|
||||
res.Status = status.NewNotFound(ctx, "user not found")
|
||||
} else {
|
||||
case errtypes.Unavailable:
|
||||
res.Status = status.NewUnavailable(ctx, "user provider temporarily unavailable")
|
||||
default:
|
||||
res.Status = status.NewInternal(ctx, "error getting user")
|
||||
}
|
||||
return res, nil
|
||||
@@ -205,9 +208,12 @@ func (s *service) GetUserByClaim(ctx context.Context, req *userpb.GetUserByClaim
|
||||
user, err := s.usermgr.GetUserByClaim(ctx, req.Claim, req.Value, tenantID, req.SkipFetchingUserGroups)
|
||||
if err != nil {
|
||||
res := &userpb.GetUserByClaimResponse{}
|
||||
if _, ok := err.(errtypes.NotFound); ok {
|
||||
switch err.(type) {
|
||||
case errtypes.NotFound:
|
||||
res.Status = status.NewNotFound(ctx, fmt.Sprintf("user not found %s %s", req.Claim, req.Value))
|
||||
} else {
|
||||
case errtypes.Unavailable:
|
||||
res.Status = status.NewUnavailable(ctx, "user provider temporarily unavailable")
|
||||
default:
|
||||
res.Status = status.NewInternal(ctx, "error getting user by claim")
|
||||
}
|
||||
return res, nil
|
||||
|
||||
+21
@@ -203,6 +203,15 @@ func (e TooEarly) Error() string { return "error: too early: " + string(e) }
|
||||
// IsTooEarly implements the IsTooEarly interface.
|
||||
func (e TooEarly) IsTooEarly() {}
|
||||
|
||||
// Unavailable is the error to use when a backend service (e.g. LDAP, database) is
|
||||
// temporarily unreachable. Callers should treat this as a transient failure and retry.
|
||||
type Unavailable string
|
||||
|
||||
func (e Unavailable) Error() string { return "error: unavailable: " + string(e) }
|
||||
|
||||
// IsUnavailable implements the IsUnavailable interface.
|
||||
func (e Unavailable) IsUnavailable() {}
|
||||
|
||||
// IsNotFound is the interface to implement
|
||||
// to specify that a resource is not found.
|
||||
type IsNotFound interface {
|
||||
@@ -293,6 +302,12 @@ type IsTooEarly interface {
|
||||
IsTooEarly()
|
||||
}
|
||||
|
||||
// IsUnavailable is the interface to implement to specify that a backend service is
|
||||
// temporarily unavailable and the caller should retry.
|
||||
type IsUnavailable interface {
|
||||
IsUnavailable()
|
||||
}
|
||||
|
||||
// NewErrtypeFromStatus maps a rpc status to an errtype
|
||||
func NewErrtypeFromStatus(status *rpc.Status) error {
|
||||
switch status.Code {
|
||||
@@ -329,6 +344,8 @@ func NewErrtypeFromStatus(status *rpc.Status) error {
|
||||
return BadRequest(status.Message)
|
||||
case rpc.Code_CODE_TOO_EARLY:
|
||||
return TooEarly(status.Message)
|
||||
case rpc.Code_CODE_UNAVAILABLE:
|
||||
return Unavailable(status.Message)
|
||||
default:
|
||||
return InternalError(status.Message)
|
||||
}
|
||||
@@ -363,6 +380,8 @@ func NewErrtypeFromHTTPStatusCode(code int, message string) error {
|
||||
return PartialContent(message)
|
||||
case http.StatusTooEarly:
|
||||
return TooEarly(message)
|
||||
case http.StatusServiceUnavailable:
|
||||
return Unavailable(message)
|
||||
case StatusChecksumMismatch:
|
||||
return ChecksumMismatch(message)
|
||||
default:
|
||||
@@ -399,6 +418,8 @@ func NewHTTPStatusCodeFromErrtype(err error) int {
|
||||
return http.StatusPartialContent
|
||||
case TooEarly:
|
||||
return http.StatusTooEarly
|
||||
case Unavailable:
|
||||
return http.StatusServiceUnavailable
|
||||
case ChecksumMismatch:
|
||||
return StatusChecksumMismatch
|
||||
default:
|
||||
|
||||
+25
-21
@@ -71,9 +71,8 @@ type RawStream struct {
|
||||
c Config
|
||||
}
|
||||
|
||||
func FromConfig(ctx context.Context, name string, cfg Config) (Stream, error) {
|
||||
var s Stream
|
||||
b := backoff.NewExponentialBackOff()
|
||||
func JetStream(ctx context.Context, name string, cfg Config) (jetstream.JetStream, error) {
|
||||
var js jetstream.JetStream
|
||||
|
||||
connect := func() error {
|
||||
var tlsConf *tls.Config
|
||||
@@ -120,27 +119,32 @@ func FromConfig(ctx context.Context, name string, cfg Config) (Stream, error) {
|
||||
return err
|
||||
}
|
||||
|
||||
jsConn, err := jetstream.New(conn)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
js, err := jsConn.Stream(ctx, events.MainQueueName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
s = &RawStream{
|
||||
js: js,
|
||||
c: cfg,
|
||||
}
|
||||
return nil
|
||||
js, err = jetstream.New(conn)
|
||||
return err
|
||||
}
|
||||
err := backoff.Retry(connect, b)
|
||||
|
||||
err := backoff.Retry(connect, backoff.NewExponentialBackOff())
|
||||
if err != nil {
|
||||
return s, errors.Wrap(err, "could not connect to nats jetstream")
|
||||
return nil, errors.Wrap(err, "could not connect to nats jetstream")
|
||||
}
|
||||
return s, nil
|
||||
return js, nil
|
||||
}
|
||||
|
||||
func FromConfig(ctx context.Context, name string, cfg Config) (Stream, error) {
|
||||
jsConn, err := JetStream(ctx, name, cfg)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
js, err := jsConn.Stream(ctx, events.MainQueueName)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &RawStream{
|
||||
js: js,
|
||||
c: cfg,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *RawStream) Consume(group string, evs ...events.Unmarshaller) (<-chan Event, error) {
|
||||
|
||||
+35
-1
@@ -2,6 +2,7 @@ package stream
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"errors"
|
||||
@@ -11,7 +12,9 @@ import (
|
||||
|
||||
"github.com/cenkalti/backoff"
|
||||
"github.com/go-micro/plugins/v4/events/natsjs"
|
||||
"github.com/nats-io/nats.go/jetstream"
|
||||
"github.com/opencloud-eu/reva/v2/pkg/events"
|
||||
"github.com/opencloud-eu/reva/v2/pkg/events/raw"
|
||||
"github.com/opencloud-eu/reva/v2/pkg/logger"
|
||||
)
|
||||
|
||||
@@ -65,7 +68,38 @@ func NatsFromConfig(connName string, disableDurability bool, cfg NatsConfig) (ev
|
||||
opts = append(opts, natsjs.DisableDurableStreams())
|
||||
}
|
||||
|
||||
return Nats(opts...)
|
||||
s, err := Nats(opts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// apply a MaxAge to the main queue to prevent it from filling up
|
||||
ctx := context.Background()
|
||||
jsConn, err := raw.JetStream(ctx, connName, raw.Config{
|
||||
Endpoint: cfg.Endpoint,
|
||||
Cluster: cfg.Cluster,
|
||||
TLSInsecure: cfg.TLSInsecure,
|
||||
TLSRootCACertificate: cfg.TLSRootCACertificate,
|
||||
EnableTLS: cfg.EnableTLS,
|
||||
AuthUsername: cfg.AuthUsername,
|
||||
AuthPassword: cfg.AuthPassword,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
streamCfg := jetstream.StreamConfig{
|
||||
Name: "main-queue",
|
||||
MaxAge: 7 * 24 * time.Hour,
|
||||
}
|
||||
_, err = jsConn.CreateStream(ctx, streamCfg)
|
||||
if err != nil {
|
||||
// If the stream already exists, update its configuration
|
||||
if err == jetstream.ErrStreamNameAlreadyInUse {
|
||||
_, _ = jsConn.UpdateStream(ctx, streamCfg)
|
||||
}
|
||||
}
|
||||
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// nats returns a nats streaming client
|
||||
|
||||
+13
@@ -68,6 +68,15 @@ func NewInternal(ctx context.Context, msg string) *rpc.Status {
|
||||
}
|
||||
}
|
||||
|
||||
// NewUnavailable returns a Status with CODE_UNAVAILABLE.
|
||||
func NewUnavailable(ctx context.Context, msg string) *rpc.Status {
|
||||
return &rpc.Status{
|
||||
Code: rpc.Code_CODE_UNAVAILABLE,
|
||||
Message: msg,
|
||||
Trace: getTrace(ctx),
|
||||
}
|
||||
}
|
||||
|
||||
// NewUnauthenticated returns a Status with CODE_UNAUTHENTICATED.
|
||||
func NewUnauthenticated(ctx context.Context, err error, msg string) *rpc.Status {
|
||||
return &rpc.Status{
|
||||
@@ -191,6 +200,10 @@ func NewStatusFromErrType(ctx context.Context, msg string, err error) *rpc.Statu
|
||||
return NewUnimplemented(ctx, err, msg+":"+err.Error())
|
||||
case errtypes.BadRequest:
|
||||
return NewInvalid(ctx, msg+":"+err.Error())
|
||||
case errtypes.Unavailable:
|
||||
return NewUnavailable(ctx, msg+": "+err.Error())
|
||||
case errtypes.IsUnavailable:
|
||||
return NewUnavailable(ctx, msg+": "+err.Error())
|
||||
}
|
||||
|
||||
// map GRPC status codes coming from the auth middleware
|
||||
|
||||
+1
-1
@@ -72,7 +72,7 @@ func NewNatsKeyValueFromJetStream(c Config, js jetstream.JetStream) (jetstream.K
|
||||
if err != nil {
|
||||
kvConfig := jetstream.KeyValueConfig{
|
||||
Bucket: c.Database,
|
||||
TTL: 0, // we don't do TTLs for this store
|
||||
TTL: c.TTL,
|
||||
}
|
||||
if c.DisablePersistence {
|
||||
kvConfig.Storage = jetstream.MemoryStorage
|
||||
|
||||
+11
-1
@@ -24,6 +24,7 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/rs/zerolog"
|
||||
tusd "github.com/tus/tusd/v2/pkg/handler"
|
||||
@@ -58,7 +59,8 @@ func init() {
|
||||
type posixFS struct {
|
||||
storage.FS
|
||||
|
||||
um usermapper.Mapper
|
||||
tree *tree.Tree
|
||||
um usermapper.Mapper
|
||||
}
|
||||
|
||||
// New returns an implementation to of the storage.FS interface that talk to
|
||||
@@ -70,6 +72,7 @@ func NewDefault(m map[string]interface{}, stream events.Stream, log *zerolog.Log
|
||||
}
|
||||
|
||||
o.IDCache.Database += "_v2" // Use a versioned bucket name to avoid conflicts with previous implementations
|
||||
o.IDCache.TTL = 0 // Disable TTL for the ID cache, as the posix driver relies on it for caching file IDs and we don't want them to expire
|
||||
kv, err := cache.NewNatsKeyValue(o.IDCache)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "could not create nats key value store")
|
||||
@@ -80,6 +83,7 @@ func NewDefault(m map[string]interface{}, stream events.Stream, log *zerolog.Log
|
||||
}
|
||||
|
||||
o.IDCache.Database += "_history" // Use a versioned bucket name to avoid conflicts with previous implementations
|
||||
o.IDCache.TTL = 24 * 60 * time.Minute
|
||||
historyKv, err := cache.NewNatsKeyValue(o.IDCache)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "could not create nats key value store")
|
||||
@@ -215,11 +219,17 @@ func New(o *options.Options, stream events.Stream, cache, historyCache *idcache.
|
||||
|
||||
mw := middleware.NewFS(dfs, hooks...)
|
||||
fs.FS = mw
|
||||
fs.tree = tp
|
||||
fs.um = um
|
||||
|
||||
return fs, nil
|
||||
}
|
||||
|
||||
// WarmupIDCache allows triggering a posix fs scan and id cache warmup manually.
|
||||
func (fs *posixFS) WarmupIDCache(root string, assimilate, onlyDirty bool) error {
|
||||
return fs.tree.WarmupIDCache(root, assimilate, onlyDirty)
|
||||
}
|
||||
|
||||
// ListUploadSessions returns the upload sessions matching the given filter
|
||||
func (fs *posixFS) ListUploadSessions(ctx context.Context, filter storage.UploadSessionFilter) ([]storage.UploadSession, error) {
|
||||
return fs.FS.(storage.UploadSessionLister).ListUploadSessions(ctx, filter)
|
||||
|
||||
+3
@@ -870,6 +870,9 @@ func (t *Tree) WarmupIDCache(root string, assimilate, onlyDirty bool) error {
|
||||
isTrash(path) ||
|
||||
t.isUpload(path) ||
|
||||
t.isIndex(path) {
|
||||
if info.IsDir() {
|
||||
return filepath.SkipDir
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if t.isRootPath(path) {
|
||||
|
||||
+48
-41
@@ -266,15 +266,10 @@ func (i *Identity) GetLDAPUserByFilter(ctx context.Context, lc ldap.Client, filt
|
||||
res, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
log.Debug().Str("backend", "ldap").Err(err).Str("userfilter", filter).Msg("Error looking up user by filter")
|
||||
var errmsg string
|
||||
if lerr, ok := err.(*ldap.Error); ok {
|
||||
if lerr.ResultCode == ldap.LDAPResultSizeLimitExceeded {
|
||||
errmsg = fmt.Sprintf("too many results searching for user '%s'", filter)
|
||||
}
|
||||
}
|
||||
span.SetAttributes(attribute.String("ldap.error", errmsg))
|
||||
span.SetStatus(codes.Error, errmsg)
|
||||
return nil, errtypes.NotFound(errmsg)
|
||||
classified := classifySearchError(err, fmt.Sprintf("too many results searching for user '%s'", filter))
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
if len(res.Entries) == 0 {
|
||||
return nil, errtypes.NotFound(filter)
|
||||
@@ -306,9 +301,10 @@ func (i *Identity) GetLDAPUserByDN(ctx context.Context, lc ldap.Client, dn strin
|
||||
res, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
log.Debug().Str("backend", "ldap").Err(err).Str("dn", dn).Msg("Error looking up user by DN")
|
||||
span.SetAttributes(attribute.String("ldap.error", err.Error()))
|
||||
span.SetStatus(codes.Error, "")
|
||||
return nil, errtypes.NotFound(dn)
|
||||
classified := classifySearchError(err, "")
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
span.SetStatus(codes.Ok, "")
|
||||
if len(res.Entries) == 0 {
|
||||
@@ -337,9 +333,10 @@ func (i *Identity) GetLDAPUsers(ctx context.Context, lc ldap.Client, query, tena
|
||||
sr, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
log.Debug().Str("backend", "ldap").Err(err).Str("filter", filter).Msg("Error searching users")
|
||||
span.SetAttributes(attribute.String("ldap.error", err.Error()))
|
||||
span.SetStatus(codes.Error, "")
|
||||
return nil, errtypes.NotFound(query)
|
||||
classified := classifySearchError(err, "")
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
|
||||
span.SetAttributes(attribute.Int("ldap.result_count", len(sr.Entries)))
|
||||
@@ -376,7 +373,8 @@ func (i *Identity) IsLDAPUserInDisabledGroup(ctx context.Context, lc ldap.Client
|
||||
sr, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
log.Error().Str("backend", "ldap").Err(err).Str("filter", filter).Msg("Error looking up error group")
|
||||
// Err on the side of caution.
|
||||
// Err on the side of caution: treat search failures (including network
|
||||
// errors) as if the user is in the disabled group.
|
||||
span.SetAttributes(attribute.String("ldap.error", err.Error()))
|
||||
span.SetStatus(codes.Error, "")
|
||||
return true
|
||||
@@ -423,10 +421,10 @@ func (i *Identity) GetLDAPUserGroups(ctx context.Context, lc ldap.Client, userEn
|
||||
// not having any groups in LDAP
|
||||
return []string{}, nil
|
||||
}
|
||||
|
||||
span.SetAttributes(attribute.String("ldap.error", err.Error()))
|
||||
span.SetStatus(codes.Error, "")
|
||||
return []string{}, err
|
||||
classified := classifySearchError(err, "")
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
span.SetStatus(codes.Ok, "")
|
||||
span.SetAttributes(attribute.Int("ldap.result_count", len(sr.Entries)))
|
||||
@@ -504,15 +502,10 @@ func (i *Identity) GetLDAPGroupByFilter(ctx context.Context, lc ldap.Client, fil
|
||||
res, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
log.Debug().Str("backend", "ldap").Err(err).Str("filter", filter).Msg("Error looking up group by filter")
|
||||
var errmsg string
|
||||
if lerr, ok := err.(*ldap.Error); ok {
|
||||
if lerr.ResultCode == ldap.LDAPResultSizeLimitExceeded {
|
||||
errmsg = fmt.Sprintf("too many results searching for group '%s'", filter)
|
||||
}
|
||||
}
|
||||
span.SetAttributes(attribute.String("ldap.error", errmsg))
|
||||
span.SetStatus(codes.Error, "")
|
||||
return nil, errtypes.NotFound(errmsg)
|
||||
classified := classifySearchError(err, fmt.Sprintf("too many results searching for group '%s'", filter))
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
if len(res.Entries) == 0 {
|
||||
return nil, errtypes.NotFound(filter)
|
||||
@@ -543,10 +536,11 @@ func (i *Identity) GetLDAPGroups(ctx context.Context, lc ldap.Client, query stri
|
||||
setLDAPSearchSpanAttributes(span, searchRequest)
|
||||
sr, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
span.SetAttributes(attribute.String("ldap.error", err.Error()))
|
||||
span.SetStatus(codes.Error, "")
|
||||
log.Debug().Str("backend", "ldap").Err(err).Str("query", query).Msg("Error search for groups")
|
||||
return nil, errtypes.NotFound(query)
|
||||
classified := classifySearchError(err, "")
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
span.SetStatus(codes.Ok, "")
|
||||
return sr.Entries, nil
|
||||
@@ -919,15 +913,10 @@ func (i *Identity) GetLDAPTenantByFilter(ctx context.Context, lc ldap.Client, fi
|
||||
res, err := lc.Search(searchRequest)
|
||||
if err != nil {
|
||||
log.Debug().Str("backend", "ldap").Err(err).Str("tenantfilter", filter).Msg("Error looking up tenant by filter")
|
||||
var errmsg string
|
||||
if lerr, ok := err.(*ldap.Error); ok {
|
||||
if lerr.ResultCode == ldap.LDAPResultSizeLimitExceeded {
|
||||
errmsg = fmt.Sprintf("too many results searching for tenant '%s'", filter)
|
||||
}
|
||||
}
|
||||
span.SetAttributes(attribute.String("ldap.error", errmsg))
|
||||
span.SetStatus(codes.Error, errmsg)
|
||||
return nil, errtypes.NotFound(errmsg)
|
||||
classified := classifySearchError(err, fmt.Sprintf("too many results searching for tenant '%s'", filter))
|
||||
span.SetAttributes(attribute.String("ldap.error", classified.Error()))
|
||||
span.SetStatus(codes.Error, classified.Error())
|
||||
return nil, classified
|
||||
}
|
||||
if len(res.Entries) == 0 {
|
||||
return nil, errtypes.NotFound(filter)
|
||||
@@ -980,6 +969,24 @@ func (i *Identity) getTenantAttributeFilter(attribute, value string) (string, er
|
||||
), nil
|
||||
}
|
||||
|
||||
// classifySearchError maps a raw error from lc.Search to the appropriate
|
||||
// errtypes value:
|
||||
// - ldap.ErrorNetwork → errtypes.Unavailable (transient; caller should retry)
|
||||
// - ldap.LDAPResultSizeLimitExceeded → errtypes.NotFound(sizeExceededMsg)
|
||||
// - anything else → errtypes.NotFound("") (preserving prior behaviour)
|
||||
//
|
||||
// The sizeExceededMsg is only used for the SizeLimitExceeded case; pass an
|
||||
// empty string if the caller does not need a custom message for that case.
|
||||
func classifySearchError(err error, sizeExceededMsg string) error {
|
||||
if ldap.IsErrorWithCode(err, ldap.ErrorNetwork) {
|
||||
return errtypes.Unavailable("ldap server unreachable: " + err.Error())
|
||||
}
|
||||
if sizeExceededMsg != "" && ldap.IsErrorWithCode(err, ldap.LDAPResultSizeLimitExceeded) {
|
||||
return errtypes.NotFound(sizeExceededMsg)
|
||||
}
|
||||
return errtypes.NotFound("")
|
||||
}
|
||||
|
||||
func setLDAPSearchSpanAttributes(span trace.Span, request *ldap.SearchRequest) {
|
||||
span.SetAttributes(
|
||||
attribute.String("ldap.basedn", request.BaseDN),
|
||||
|
||||
Reference in New Issue
Block a user