chore: bump reva to 295fe6437

Signed-off-by: Jörn Friedrich Dreyer <jfd@butonic.de>
This commit is contained in:
Jörn Friedrich Dreyer committed 2026-10-02 12:24:00 +02:00
1 parent e848d469d5
commit fe507ff72a
23 files changed
+1069 -216

No files matched your search

+1 -1
View File
@@ -64,7 +64,7 @@ require (
github.com/open-policy-agent/opa v1.19.1
github.com/opencloud-eu/icap-client v0.0.0-20250930132611-28a2afe62d89
github.com/opencloud-eu/libre-graph-api-go v1.0.8-0.20260902170011-45af3945a067
github.com/opencloud-eu/reva/v2 v2.50.1-0.20261001091108-11d87fb6b985
github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002075447-9cf4d4365f4f
github.com/opensearch-project/opensearch-go/v4 v4.7.3
github.com/orcaman/concurrent-map v1.0.0
github.com/pkg/errors v0.9.1
+2 -2
View File
@@ -761,8 +761,8 @@ github.com/opencloud-eu/icap-client v0.0.0-20250930132611-28a2afe62d89 h1:W1ms+l
github.com/opencloud-eu/icap-client v0.0.0-20250930132611-28a2afe62d89/go.mod h1:vigJkNss1N2QEceCuNw/ullDehncuJNFB6mEnzfq9UI=
github.com/opencloud-eu/libre-graph-api-go v1.0.8-0.20260902170011-45af3945a067 h1:UkNMKauyJAzY6RE6mmthz9bQZLYkbvBuApm7ZDCparE=
github.com/opencloud-eu/libre-graph-api-go v1.0.8-0.20260902170011-45af3945a067/go.mod h1:lTM8JeGblNpoMySTW7Lui2+c5TTLI95mwxtdUIHHrhU=
github.com/opencloud-eu/reva/v2 v2.50.1-0.20261001091108-11d87fb6b985 h1:9NpeJKVzjzMKhiTaBPz6SqtRiv6k67E5wKPPc0zg/8M=
github.com/opencloud-eu/reva/v2 v2.50.1-0.20261001091108-11d87fb6b985/go.mod h1:ngLsDakKnt8zCFCeULhuSZFESGy5NIuSq5/JGeOhPdU=
github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002075447-9cf4d4365f4f h1:mlxI1BZr+ReoM2UPC7W4vlmzpzwSbCotcYNnfSK5YYw=
github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002075447-9cf4d4365f4f/go.mod h1:ngLsDakKnt8zCFCeULhuSZFESGy5NIuSq5/JGeOhPdU=
github.com/opencloud-eu/secure v0.0.0-20260312082735-b6f5cb2244e4 h1:l2oB/RctH+t8r7QBj5p8thfEHCM/jF35aAY3WQ3hADI=
github.com/opencloud-eu/secure v0.0.0-20260312082735-b6f5cb2244e4/go.mod h1:BmF5hyM6tXczk3MpQkFf1hpKSRqCyhqcbiQtiAF7+40=
github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U=
@@ -139,8 +139,12 @@ func (s *service) Authenticate(ctx context.Context, req *provider.AuthenticateRe
u, scope, err := s.authmgr.Authenticate(ctx, username, password)
if err != nil {
log.Debug().Str("client_id", username).Err(err).Msg("authsvc: error in Authenticate")
st := status.NewStatusFromErrType(ctx, "authsvc: error in Authenticate", err)
if entry := status.InnerErrorFromErr(err); entry != nil {
st.InnerError = entry
}
return &provider.AuthenticateResponse{
Status: status.NewStatusFromErrType(ctx, "authsvc: error in Authenticate", err),
Status: st,
}, nil
}
log.Info().Msgf("user %s authenticated", u.Id)
@@ -57,25 +57,9 @@ func (s *svc) Authenticate(ctx context.Context, req *gateway.AuthenticateRequest
ClientId: req.ClientId,
ClientSecret: req.ClientSecret,
}
res, err := c.Authenticate(ctx, authProviderReq)
switch {
case err != nil:
return &gateway.AuthenticateResponse{
Status: status.NewInternal(ctx, fmt.Sprintf("gateway: error calling Authenticate for type: %s", req.Type)),
}, nil
case res.Status.Code == rpc.Code_CODE_PERMISSION_DENIED:
fallthrough
case res.Status.Code == rpc.Code_CODE_UNAUTHENTICATED:
fallthrough
case res.Status.Code == rpc.Code_CODE_NOT_FOUND:
// normal failures, no need to log
return &gateway.AuthenticateResponse{
Status: res.Status,
}, nil
case res.Status.Code != rpc.Code_CODE_OK:
return &gateway.AuthenticateResponse{
Status: status.NewInternal(ctx, fmt.Sprintf("error authenticating credentials to auth provider for type: %s", req.Type)),
}, nil
res, callErr := c.Authenticate(ctx, authProviderReq)
if resp, done := translateProviderAuthenticateResult(ctx, req.Type, res, callErr); done {
return resp, nil
}
// validate valid userId
@@ -109,7 +93,7 @@ func (s *svc) Authenticate(ctx context.Context, req *gateway.AuthenticateRequest
return res, nil
}
if scope, ok := res.TokenScope["user"]; s.c.DisableHomeCreationOnLogin || !ok || scope.Role != authpb.Role_ROLE_OWNER || res.User.Id.Type == userpb.UserType_USER_TYPE_FEDERATED {
if scope, ok := res.TokenScope["user"]; s.c.DisableHomeCreationOnLogin || !ok || scope.Role != authpb.Role_ROLE_OWNER || res.User.Id.Type == userpb.UserType_USER_TYPE_FEDERATED || res.User.Id.Type == userpb.UserType_USER_TYPE_GUEST {
gwRes := &gateway.AuthenticateResponse{
Status: status.NewOK(ctx),
User: res.User,
@@ -151,6 +135,42 @@ func (s *svc) Authenticate(ctx context.Context, req *gateway.AuthenticateRequest
return gwRes, nil
}
// translateProviderAuthenticateResult inspects the result of calling the auth
// provider's Authenticate RPC and decides whether the gateway should return
// early with a translated response.
//
// Normal authentication failure statuses (CODE_UNAUTHENTICATED,
// CODE_PERMISSION_DENIED, CODE_NOT_FOUND, CODE_UNAVAILABLE) are passed
// through unchanged, preserving their Message, Trace and InnerError. Any
// other unexpected application status is turned into CODE_INTERNAL.
//
// returns done == true if the caller should return resp immediately.
func translateProviderAuthenticateResult(ctx context.Context, authType string, res *authpb.AuthenticateResponse, callErr error) (resp *gateway.AuthenticateResponse, done bool) {
switch {
case callErr != nil:
return &gateway.AuthenticateResponse{
Status: status.NewInternal(ctx, fmt.Sprintf("gateway: error calling Authenticate for type: %s", authType)),
}, true
case res.Status.Code == rpc.Code_CODE_PERMISSION_DENIED:
fallthrough
case res.Status.Code == rpc.Code_CODE_UNAUTHENTICATED:
fallthrough
case res.Status.Code == rpc.Code_CODE_NOT_FOUND:
fallthrough
case res.Status.Code == rpc.Code_CODE_UNAVAILABLE:
// normal failures, no need to log
return &gateway.AuthenticateResponse{
Status: res.Status,
}, true
case res.Status.Code != rpc.Code_CODE_OK:
return &gateway.AuthenticateResponse{
Status: status.NewInternal(ctx, fmt.Sprintf("error authenticating credentials to auth provider for type: %s", authType)),
}, true
}
return nil, false
}
func (s *svc) WhoAmI(ctx context.Context, req *gateway.WhoAmIRequest) (*gateway.WhoAmIResponse, error) {
u, _, err := s.tokenmgr.DismantleToken(ctx, req.Token)
if err != nil {
@@ -27,6 +27,7 @@ import (
ctxpkg "github.com/opencloud-eu/reva/v2/pkg/ctx"
sdk "github.com/opencloud-eu/reva/v2/pkg/sdk/common"
"github.com/opencloud-eu/reva/v2/pkg/storage/cache"
"github.com/opencloud-eu/reva/v2/pkg/storagespace"
"github.com/opencloud-eu/reva/v2/pkg/utils"
"github.com/pkg/errors"
"google.golang.org/grpc"
@@ -129,7 +130,17 @@ func (c *cachedSpacesAPIClient) UpdateStorageSpace(ctx context.Context, in *prov
return c.c.UpdateStorageSpace(ctx, in, opts...)
}
func (c *cachedSpacesAPIClient) DeleteStorageSpace(ctx context.Context, in *provider.DeleteStorageSpaceRequest, opts ...grpc.CallOption) (*provider.DeleteStorageSpaceResponse, error) {
return c.c.DeleteStorageSpace(ctx, in, opts...)
resp, err := c.c.DeleteStorageSpace(ctx, in, opts...)
switch {
case err != nil:
return nil, err
case resp.Status.Code != rpc.Code_CODE_OK:
return resp, nil
default:
_, spaceid, _, _ := storagespace.SplitID(in.GetId().GetOpaqueId())
_ = c.createPersonalSpaceCache.Delete(spaceid)
return resp, nil
}
}
/*
@@ -0,0 +1,61 @@
// Copyright 2026 OpenCloud GmbH <mail@opencloud.eu>
// SPDX-License-Identifier: Apache-2.0
package guestlinks
import (
"encoding/json"
types "github.com/cs3org/go-cs3apis/cs3/types/v1beta1"
)
const (
innerErrorType = "opencloud_guest_link_error"
reasonSessionExpired = "session_expired"
)
// errSessionExpired is returned when JWT expiry is the *only* validation
// failure. It implements errtypes.IsInvalidCredentials (-> CODE_UNAUTHENTICATED)
// and status.StatusInnerErrorProvider so that a safe, versioned JSON detail
// is attached to rpc.Status.InnerError.
type errSessionExpired struct {
shareID string
}
func (e *errSessionExpired) Error() string {
return "guestlinks: session expired"
}
// IsInvalidCredentials implements the errtypes.IsInvalidCredentials interface.
func (e *errSessionExpired) IsInvalidCredentials() {}
// innerErrorPayload is the JSON payload shape for the guest-link session
// expired InnerError.
type innerErrorPayload struct {
Type string `json:"type"`
Reason string `json:"reason"`
ShareID string `json:"share_id"`
}
// StatusInnerError implements status.StatusInnerErrorProvider.
func (e *errSessionExpired) StatusInnerError() *types.OpaqueEntry {
payload := innerErrorPayload{
Type: innerErrorType,
Reason: reasonSessionExpired,
ShareID: e.shareID,
}
// encoding a static, safe struct: an error here can only happen on
// programmer error (e.g. unmarshalable field), never in practice.
b, err := json.Marshal(payload)
if err != nil {
return nil
}
return &types.OpaqueEntry{
Decoder: "json",
Value: b,
}
}
func newSessionExpiredError(shareID string) error {
return &errSessionExpired{shareID: shareID}
}
@@ -0,0 +1,306 @@
// Copyright 2026 OpenCloud GmbH <mail@opencloud.eu>
// SPDX-License-Identifier: Apache-2.0
// Package guestlinks implements a Reva auth.Manager that authenticates
// OpenCloud guest-link sessions.
//
// It validates an OpenCloud guest-session JWT (issued by the OpenCloud
// proxy), re-validates the JWT's anchor collaborative share through the
// gateway GetShare API using a freshly minted service-account token, and
// returns the synthetic guest identity persisted as the share's grantee.
//
// See opencloud issue #3070 for the full design rationale.
package guestlinks
import (
"context"
"errors"
"net/mail"
"time"
authpb "github.com/cs3org/go-cs3apis/cs3/auth/provider/v1beta1"
userpb "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
rpc "github.com/cs3org/go-cs3apis/cs3/rpc/v1beta1"
collaboration "github.com/cs3org/go-cs3apis/cs3/sharing/collaboration/v1beta1"
storageprovider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
"github.com/go-viper/mapstructure/v2"
"github.com/golang-jwt/jwt/v5"
"github.com/rs/zerolog"
"github.com/opencloud-eu/reva/v2/pkg/auth"
"github.com/opencloud-eu/reva/v2/pkg/auth/manager/registry"
"github.com/opencloud-eu/reva/v2/pkg/auth/scope"
"github.com/opencloud-eu/reva/v2/pkg/errtypes"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/share"
"github.com/opencloud-eu/reva/v2/pkg/utils"
)
const (
jwtLeeway = 30 * time.Second
maxShareIDLength = 512
)
func init() {
registry.Register("guestlinks", New)
}
// guestClaims are the required custom JWT claims of a guest-session token.
type guestClaims struct {
// the attribute in the JWT is called permissionId (for consistency with the
// graph API), for us it really is the "shareId"
ShareID string `json:"permissionId"`
jwt.RegisteredClaims
}
// config holds the guestlinks auth manager configuration.
type config struct {
GatewayAddr string `mapstructure:"gateway_addr"`
JWTSecret string `mapstructure:"jwt_secret"`
ServiceAccountID string `mapstructure:"service_account_id"`
ServiceAccountSecret string `mapstructure:"service_account_secret"`
}
func (c *config) validate() error {
switch {
case c.GatewayAddr == "":
return errors.New("guestlinks: gateway_addr must not be empty")
case c.JWTSecret == "":
return errors.New("guestlinks: jwt_secret must not be empty")
case c.ServiceAccountID == "":
return errors.New("guestlinks: service_account_id must not be empty")
case c.ServiceAccountSecret == "":
return errors.New("guestlinks: service_account_secret must not be empty")
}
return nil
}
type manager struct {
c *config
log *zerolog.Logger
}
func parseConfig(m map[string]any) (*config, error) {
c := &config{}
if err := mapstructure.Decode(m, c); err != nil {
return nil, errors.New("guestlinks: error decoding conf: " + err.Error())
}
return c, nil
}
// New returns a new guestlinks auth.Manager.
func New(m map[string]any, log *zerolog.Logger) (auth.Manager, error) {
mgr := &manager{log: log}
if err := mgr.Configure(m); err != nil {
return nil, err
}
return mgr, nil
}
// Configure parses and validates the manager configuration.
func (m *manager) Configure(ml map[string]any) error {
c, err := parseConfig(ml)
if err != nil {
return err
}
if err := c.validate(); err != nil {
return err
}
m.c = c
return nil
}
// Authenticate implements auth.Manager.
//
// clientID must be empty as the guestlinks credential is carried entirely in
// clientSecret, which must is the raw guest-session JWT.
func (m *manager) Authenticate(ctx context.Context, clientID, clientSecret string) (*userpb.User, map[string]*authpb.Scope, error) {
if clientID != "" {
m.logOutcome("invalid", "non-empty client_id")
return nil, nil, errtypes.InvalidCredentials("non-empty client_id")
}
shareID, err := m.validateToken(clientSecret)
if err != nil {
if _, ok := errors.AsType[*errSessionExpired](err); ok {
m.logOutcome("expired", "token expired")
} else {
m.logOutcome("invalid", "token validation failed")
}
return nil, nil, err
}
foundShare, err := m.lookupShare(ctx, shareID)
if err != nil {
return nil, nil, err
}
u, err := guestUserFromShare(foundShare)
if err != nil {
m.logOutcome("share-invalid", "share grantee invalid")
return nil, nil, errtypes.InvalidCredentials("share grantee invalid")
}
sc, err := scope.AddOwnerScope(nil)
if err != nil {
m.logOutcome("internal", "error building scope")
return nil, nil, errtypes.InternalError("guestlinks: error building token scope: " + err.Error())
}
m.logOutcome("success", "")
return u, sc, nil
}
// validateToken validates the JWT contained in raw and returns the
// validated share_id.
//
// If the *only* validation failure is expiry, the returned error is a
// *errSessionExpired (still carrying the validated share_id), per the
// expiry-only classification rules. Any other validation failure is
// returned as a generic errtypes.InvalidCredentials with no detail.
func (m *manager) validateToken(raw string) (shareID string, err error) {
if raw == "" {
return "", errtypes.InvalidCredentials("empty token")
}
claims := &guestClaims{}
parser := jwt.NewParser(
jwt.WithValidMethods([]string{"HS256"}),
jwt.WithoutClaimsValidation(),
)
_, parseErr := parser.ParseWithClaims(raw, claims, func(t *jwt.Token) (any, error) {
if _, ok := t.Method.(*jwt.SigningMethodHMAC); !ok || t.Method.Alg() != "HS256" {
return nil, errors.New("unexpected signing method")
}
return []byte(m.c.JWTSecret), nil
})
if parseErr != nil {
return "", errtypes.InvalidCredentials("signature/algorithm/parse error")
}
if claims.ShareID == "" || len(claims.ShareID) > maxShareIDLength {
return "", errtypes.InvalidCredentials("missing/malformed share_id")
}
if claims.IssuedAt == nil || claims.ExpiresAt == nil {
return "", errtypes.InvalidCredentials("missing iat/exp")
}
now := time.Now()
// iat in the future (beyond leeway) is structurally invalid.
if claims.IssuedAt.After(now.Add(jwtLeeway)) {
return "", errtypes.InvalidCredentials("invalid iat")
}
// Everything but expiry has been validated at this point. Now check
// expiry last, so we can classify an expiry-only failure.
if now.After(claims.ExpiresAt.Add(jwtLeeway)) {
return claims.ShareID, newSessionExpiredError(claims.ShareID)
}
return claims.ShareID, nil
}
// Reads and re-validates the share associated with the authentication token
// using the service account.
func (m *manager) lookupShare(ctx context.Context, shareID string) (*collaboration.Share, error) {
gwc, err := pool.GetGatewayServiceClient(m.c.GatewayAddr)
if err != nil {
m.logOutcome("unavailable", "error getting gateway client")
return nil, errtypes.Unavailable("guestlinks: error getting gateway client: " + err.Error())
}
saCtx, err := utils.GetServiceUserContextWithContext(ctx, gwc, m.c.ServiceAccountID, m.c.ServiceAccountSecret)
if err != nil {
m.logOutcome("unavailable", "error minting service account token")
return nil, errtypes.Unavailable("guestlinks: error authenticating service account: " + err.Error())
}
getShareRes, err := gwc.GetShare(saCtx, &collaboration.GetShareRequest{
Ref: &collaboration.ShareReference{
Spec: &collaboration.ShareReference_Id{
Id: &collaboration.ShareId{OpaqueId: shareID},
},
},
})
if err != nil {
m.logOutcome("unavailable", "error calling GetShare")
return nil, errtypes.Unavailable("guestlinks: error calling GetShare: " + err.Error())
}
switch getShareRes.GetStatus().GetCode() {
case rpc.Code_CODE_OK:
// fall through
case rpc.Code_CODE_NOT_FOUND, rpc.Code_CODE_PERMISSION_DENIED, rpc.Code_CODE_UNAUTHENTICATED:
m.logOutcome("share-invalid", "share not found/inaccessible")
return nil, errtypes.InvalidCredentials("share not found/inaccessible")
case rpc.Code_CODE_UNAVAILABLE:
m.logOutcome("unavailable", "share provider unavailable")
return nil, errtypes.Unavailable("guestlinks: share provider unavailable")
default:
m.logOutcome("internal", "unexpected GetShare status")
return nil, errtypes.InternalError("guestlinks: unexpected GetShare status: " + getShareRes.GetStatus().GetCode().String())
}
foundShare := getShareRes.GetShare()
if foundShare == nil {
m.logOutcome("share-invalid", "nil share")
return nil, errtypes.InvalidCredentials("nil share")
}
if share.IsExpired(foundShare) {
m.logOutcome("share-invalid", "share expired")
return nil, errtypes.InvalidCredentials("share expired")
}
return foundShare, nil
}
// guestUserFromShare builds the synthetic guest user from the persisted
// share grantee. It never trusts JWT claims for identity data.
func guestUserFromShare(s *collaboration.Share) (*userpb.User, error) {
grantee := s.GetGrantee()
if grantee.GetType() != storageprovider.GranteeType_GRANTEE_TYPE_USER {
return nil, errors.New("guestlinks: grantee is not a user")
}
uid := grantee.GetUserId()
if uid == nil || uid.GetOpaqueId() == "" {
return nil, errors.New("guestlinks: incomplete grantee user id")
}
if uid.GetType() != userpb.UserType_USER_TYPE_GUEST {
return nil, errors.New("guestlinks: grantee is not a guest user")
}
addr, err := mail.ParseAddress(uid.GetOpaqueId())
if err != nil || addr.Name != "" || addr.Address != uid.GetOpaqueId() {
return nil, errors.New("guestlinks: grantee opaque id is not a bare email address")
}
email := uid.GetOpaqueId()
u := &userpb.User{
Id: &userpb.UserId{
Idp: uid.GetIdp(),
OpaqueId: uid.GetOpaqueId(),
Type: userpb.UserType_USER_TYPE_GUEST,
TenantId: uid.GetTenantId(),
ExternalIdentities: uid.GetExternalIdentities(),
},
Username: email,
DisplayName: email,
}
// TODO: OpenCloud usually, attaches a user role to the token, how can we do that here?
return u, nil
}
func (m *manager) logOutcome(outcome, detail string) {
if m.log == nil {
return
}
ev := m.log.Debug()
if outcome != "success" {
ev = m.log.Info()
}
ev.Str("outcome", outcome).Str("detail", detail).Msg("guestlinks: authenticate outcome")
}
@@ -22,6 +22,7 @@ import (
// Load core authentication managers.
_ "github.com/opencloud-eu/reva/v2/pkg/auth/manager/appauth"
_ "github.com/opencloud-eu/reva/v2/pkg/auth/manager/demo"
_ "github.com/opencloud-eu/reva/v2/pkg/auth/manager/guestlinks"
_ "github.com/opencloud-eu/reva/v2/pkg/auth/manager/impersonator"
_ "github.com/opencloud-eu/reva/v2/pkg/auth/manager/json"
_ "github.com/opencloud-eu/reva/v2/pkg/auth/manager/ldap"
+65 -4
View File
@@ -52,16 +52,25 @@ const (
RoleEditorWithVersions = "editor-with-versions"
// RoleEditorListGrants grants editor permission on a resource, including folders.
RoleEditorListGrants = "editor-list-grants"
// RoleEditorListGrantsWithVersions grants editor permission on a resource, including folders, and list versions.
RoleEditorListGrantsWithVersions = "editor-list-grants-with-versions"
// RoleSpaceEditor grants editor permission on a space.
RoleSpaceEditor = "spaceeditor"
// RoleSpaceEditorWithoutVersions grants editor permission without list/restore versions on a space.
RoleSpaceEditorWithoutVersions = "spaceeditor-without-versions"
// RoleSpaceEditorWithoutTrashbin grants editor permission without list/restore resources in trashbin on a space.
RoleSpaceEditorWithoutTrashbin = "spaceeditor-without-trashbin"
// RoleSpaceEditorWithoutVersionsWithoutTrashbin grants editor permission without list/restore versions
// and without list/restore resources in trashbin on a space.
RoleSpaceEditorWithoutVersionsWithoutTrashbin = "spaceeditor-without-versions-without-trashbin"
// RoleFileEditor grants editor permission on a single file.
RoleFileEditor = "file-editor"
// RoleFileEditorWithVersions grants editor permission on a single file, including list/restore versions.
RoleFileEditorWithVersions = "file-editor-with-versions"
// RoleFileEditorListGrants grants editor permission on a single file.
RoleFileEditorListGrants = "file-editor-list-grants"
// RoleFileEditorListGrantsWithVersions grants editor permission on a single file and list versions.
RoleFileEditorListGrantsWithVersions = "file-editor-list-grants-with-versions"
// RoleCoowner grants co-owner permissions on a resource.
RoleCoowner = "coowner"
// RoleEditorLite grants permission to upload and download to a resource.
@@ -186,14 +195,24 @@ func RoleFromName(name string) *Role {
return NewEditorWithVersionsRole()
case RoleEditorListGrants:
return NewEditorListGrantsRole()
case RoleEditorListGrantsWithVersions:
return NewEditorListGrantsWithVersionsRole()
case RoleSpaceEditorWithoutVersions:
return NewSpaceEditorWithoutVersionsRole()
case RoleSpaceEditor:
return NewSpaceEditorRole()
case RoleSpaceEditorWithoutTrashbin:
return NewSpaceEditorWithoutTrashbinRole()
case RoleSpaceEditorWithoutVersionsWithoutTrashbin:
return NewSpaceEditorWithoutVersionsWithoutTrashbinRole()
case RoleFileEditor:
return NewFileEditorRole()
case RoleFileEditorWithVersions:
return NewFileEditorWithVersionsRole()
case RoleFileEditorListGrants:
return NewFileEditorListGrantsRole()
case RoleFileEditorListGrantsWithVersions:
return NewFileEditorListGrantsWithVersionsRole()
case RoleUploader:
return NewUploaderRole()
case RoleManager:
@@ -311,6 +330,15 @@ func NewEditorListGrantsRole() *Role {
return role
}
// NewEditorListGrantsWithVersionsRole creates an editor role that can list the invited people
// and the file versions of a resource, including folders.
func NewEditorListGrantsWithVersionsRole() *Role {
role := NewEditorListGrantsRole()
role.Name = RoleEditorListGrantsWithVersions
role.cS3ResourcePermissions.ListFileVersions = true
return role
}
// NewEditorWithVersionsRole creates an editor role including list/restore versions. `sharing` indicates if sharing permission should be added
func NewEditorWithVersionsRole() *Role {
role := NewEditorRole()
@@ -346,8 +374,18 @@ func NewSpaceEditorRole() *Role {
// NewSpaceEditorWithoutVersionsRole creates an editor without list/restore versions role
func NewSpaceEditorWithoutVersionsRole() *Role {
role := NewSpaceEditorWithoutVersionsWithoutTrashbinRole()
role.Name = RoleSpaceEditorWithoutVersions
role.cS3ResourcePermissions.ListRecycle = true
role.cS3ResourcePermissions.RestoreRecycleItem = true
return role
}
// NewSpaceEditorWithoutVersionsWithoutTrashbinRole creates an editor role without list/restore
// versions and without list/restore resources in the trashbin on a space.
func NewSpaceEditorWithoutVersionsWithoutTrashbinRole() *Role {
return &Role{
Name: RoleSpaceEditorWithoutVersions,
Name: RoleSpaceEditorWithoutVersionsWithoutTrashbin,
cS3ResourcePermissions: &provider.ResourcePermissions{
CreateContainer: true,
Delete: true,
@@ -357,15 +395,23 @@ func NewSpaceEditorWithoutVersionsRole() *Role {
InitiateFileUpload: true,
ListContainer: true,
ListGrants: true,
ListRecycle: true,
Move: true,
RestoreRecycleItem: true,
Stat: true,
},
ocsPermissions: PermissionRead | PermissionCreate | PermissionWrite | PermissionDelete,
}
}
// NewSpaceEditorWithoutTrashbinRole creates an editor role without list/restore resources
// in the trashbin on a space.
func NewSpaceEditorWithoutTrashbinRole() *Role {
role := NewSpaceEditorWithoutVersionsWithoutTrashbinRole()
role.Name = RoleSpaceEditorWithoutTrashbin
role.cS3ResourcePermissions.ListFileVersions = true
role.cS3ResourcePermissions.RestoreFileVersion = true
return role
}
// NewFileEditorRole creates a file-editor role
func NewFileEditorRole() *Role {
p := PermissionRead | PermissionWrite
@@ -392,6 +438,15 @@ func NewFileEditorListGrantsRole() *Role {
return role
}
// NewFileEditorListGrantsWithVersionsRole creates a file-editor role that can list the invited
// people and the file versions of a single file.
func NewFileEditorListGrantsWithVersionsRole() *Role {
role := NewFileEditorListGrantsRole()
role.Name = RoleFileEditorListGrantsWithVersions
role.cS3ResourcePermissions.ListFileVersions = true
return role
}
// NewFileEditorWithVersionsRole creates a file-editor role including list/restore versions
func NewFileEditorWithVersionsRole() *Role {
role := NewFileEditorRole()
@@ -626,8 +681,11 @@ func RoleFromResourcePermissions(rp *provider.ResourcePermissions, islink bool)
rp.InitiateFileDownload {
r.ocsPermissions |= PermissionRead
}
// A role without trashbin access has no RestoreRecycleItem, so writing had to be
// inferred from Delete instead - otherwise the *WithoutTrashbin space editor roles
// come back without PermissionWrite and the web frontend renders them read-only.
if rp.InitiateFileUpload &&
rp.RestoreRecycleItem {
(rp.RestoreRecycleItem || (rp.Delete && !rp.ListRecycle)) {
r.ocsPermissions |= PermissionWrite
}
if rp.Stat &&
@@ -647,6 +705,9 @@ func RoleFromResourcePermissions(rp *provider.ResourcePermissions, islink bool)
r.Name = RoleEditor
if rp.ListGrants {
r.Name = RoleEditorListGrants
if rp.ListFileVersions {
r.Name = RoleEditorListGrantsWithVersions
}
}
if rp.RemoveGrant {
r.Name = RoleManager
+16
View File
@@ -216,96 +216,112 @@ func (e Unavailable) IsUnavailable() {}
// to specify that a resource is not found.
type IsNotFound interface {
IsNotFound()
error
}
// IsAlreadyExists is the interface to implement
// to specify that a resource already exists.
type IsAlreadyExists interface {
IsAlreadyExists()
error
}
// IsInternalError is the interface to implement
// to specify that there was some internal error
type IsInternalError interface {
IsInternalError()
error
}
// IsUserRequired is the interface to implement
// to specify that a user is required.
type IsUserRequired interface {
IsUserRequired()
error
}
// IsInvalidCredentials is the interface to implement
// to specify that credentials were wrong.
type IsInvalidCredentials interface {
IsInvalidCredentials()
error
}
// IsNotSupported is the interface to implement
// to specify that an action is not supported.
type IsNotSupported interface {
IsNotSupported()
error
}
// IsPermissionDenied is the interface to implement
// to specify that an action is denied.
type IsPermissionDenied interface {
IsPermissionDenied()
error
}
// IsLocked is the interface to implement
// to specify that a resource is locked.
type IsLocked interface {
IsLocked()
error
}
// IsAborted is the interface to implement
// to specify that a request was aborted.
type IsAborted interface {
IsAborted()
error
}
// IsPreconditionFailed is the interface to implement
// to specify that a precondition failed.
type IsPreconditionFailed interface {
IsPreconditionFailed()
error
}
// IsPartialContent is the interface to implement
// to specify that the client request has partial data.
type IsPartialContent interface {
IsPartialContent()
error
}
// IsBadRequest is the interface to implement
// to specify that the server cannot or will not process the request.
type IsBadRequest interface {
IsBadRequest()
error
}
// IsChecksumMismatch is the interface to implement
// to specify that a checksum does not match.
type IsChecksumMismatch interface {
IsChecksumMismatch()
error
}
// IsInsufficientStorage is the interface to implement
// to specify that there is insufficient storage.
type IsInsufficientStorage interface {
IsInsufficientStorage()
error
}
// IsTooEarly is the interface to implement
// to specify that there is some not finished job over resource is still in process.
type IsTooEarly interface {
IsTooEarly()
error
}
// IsUnavailable is the interface to implement to specify that a backend service is
// temporarily unavailable and the caller should retry.
type IsUnavailable interface {
IsUnavailable()
error
}
// NewErrtypeFromStatus maps a rpc status to an errtype
+17
View File
@@ -238,3 +238,20 @@ func (CleanUpload) Unmarshal(v []byte) (interface{}, error) {
err := json.Unmarshal(v, &e)
return e, err
}
// DeleteRevision can be emitted to delete a revision of a node. If Timestamp
// is set exactly that revision is deleted. If Timestamp is nil the node's
// current version is targeted and the node reverts to the pre-upload state,
// unstalling it - only while the node is actually stuck in processing, so a
// redelivered event is a no-op.
type DeleteRevision struct {
ResourceID *provider.ResourceId
Timestamp *types.Timestamp
}
// Unmarshal to fulfill umarshaller interface
func (DeleteRevision) Unmarshal(v []byte) (interface{}, error) {
e := DeleteRevision{}
err := json.Unmarshal(v, &e)
return e, err
}
+36
View File
@@ -0,0 +1,36 @@
// Copyright 2026 OpenCloud GmbH <mail@opencloud.eu>
// SPDX-License-Identifier: Apache-2.0
package status
import (
"errors"
types "github.com/cs3org/go-cs3apis/cs3/types/v1beta1"
)
// StatusInnerErrorProvider is an opt-in interface for error types that carry
// safe, already-encoded failure details meant to be transported in
// rpc.Status.InnerError.
type StatusInnerErrorProvider interface {
// StatusInnerError returns an already encoded OpaqueEntry (decoder and
// value) to be attached to rpc.Status.InnerError, or nil if no detail
// should be attached.
StatusInnerError() *types.OpaqueEntry
}
// InnerErrorFromErr looks for an error implementing StatusInnerErrorProvider
// in err's chain (using errors.As, so wrapped typed errors are supported)
// and returns the safe entry it provides.
func InnerErrorFromErr(err error) *types.OpaqueEntry {
if err == nil {
return nil
}
var provider StatusInnerErrorProvider
if !errors.As(err, &provider) {
return nil
}
return provider.StatusInnerError()
}
+16 -23
View File
@@ -175,35 +175,28 @@ func NewStatusFromErrType(ctx context.Context, msg string, err error) *rpc.Statu
switch e := err.(type) {
case nil:
return NewOK(ctx)
case errtypes.NotFound:
return NewNotFound(ctx, msg+": "+err.Error())
case errtypes.IsNotFound:
return NewNotFound(ctx, msg+": "+err.Error())
case errtypes.AlreadyExists:
return NewAlreadyExists(ctx, err, msg+": "+err.Error())
case errtypes.InvalidCredentials:
return NewPermissionDenied(ctx, e, msg+": "+err.Error())
return NewNotFound(ctx, msg+": "+e.Error())
case errtypes.IsAlreadyExists:
return NewAlreadyExists(ctx, e, msg+": "+e.Error())
case errtypes.IsInvalidCredentials:
// TODO this maps badly
return NewUnauthenticated(ctx, err, msg+": "+err.Error())
case errtypes.PermissionDenied:
return NewPermissionDenied(ctx, e, msg+": "+err.Error())
case errtypes.Locked:
return NewUnauthenticated(ctx, e, msg+": "+e.Error())
case errtypes.IsPermissionDenied:
return NewPermissionDenied(ctx, e, msg+": "+e.Error())
case errtypes.IsLocked:
// FIXME a locked error returns the current lockid
// FIXME use NewAborted as per the rpc code docs
return NewLocked(ctx, msg+": "+err.Error())
case errtypes.Aborted:
return NewAborted(ctx, e, msg+": "+err.Error())
case errtypes.PreconditionFailed:
return NewFailedPrecondition(ctx, e, msg+": "+err.Error())
return NewLocked(ctx, msg+": "+e.Error())
case errtypes.IsAborted:
return NewAborted(ctx, e, msg+": "+e.Error())
case errtypes.IsPreconditionFailed:
return NewFailedPrecondition(ctx, e, msg+": "+e.Error())
case errtypes.IsNotSupported:
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())
return NewUnimplemented(ctx, e, msg+": "+e.Error())
case errtypes.IsBadRequest:
return NewInvalid(ctx, msg+": "+e.Error())
case errtypes.IsUnavailable:
return NewUnavailable(ctx, msg+": "+err.Error())
return NewUnavailable(ctx, msg+": "+e.Error())
}
// map GRPC status codes coming from the auth middleware
@@ -83,6 +83,7 @@ var (
events.RestartPostprocessing{},
events.StartPostprocessingStep{},
events.CleanUpload{},
events.DeleteRevision{},
}
)
@@ -464,6 +465,31 @@ func (fs *Decomposedfs) handlePostprocessingEvent(ctx context.Context, event eve
return // NOTE: since we can't get the upload, we can't delete the blob
}
session.Cleanup(true, !ev.KeepUpload, !ev.KeepUpload, true)
case events.DeleteRevision:
sublog := log.With().Str("event", "DeleteRevision").Interface("nodeid", ev.ResourceID).Logger()
n, err := fs.lu.NodeFromID(ctx, ev.ResourceID)
if err != nil {
sublog.Error().Err(err).Msg("Failed to get node")
return
}
var deleteErr error
if ev.Timestamp == nil {
// the node's current revision is targeted - revert to the pre-upload
// state. Only stuck nodes are targeted: a reverted node has its
// processing flag removed, so a redelivered event is a no-op.
if !n.IsProcessing(ctx) {
sublog.Debug().Msg("node is not stuck, ignoring")
return
}
_, deleteErr = n.DeleteRevision(ctx, "")
} else {
versionID := time.Unix(int64(ev.Timestamp.Seconds), int64(ev.Timestamp.Nanos)).UTC().Format(time.RFC3339Nano)
deleteErr = fs.deleteRevisionFile(ctx, n, n.ID+node.RevisionIDDelimiter+versionID)
}
if deleteErr != nil {
sublog.Error().Err(deleteErr).Msg("Failed to delete revision")
}
case events.StartPostprocessingStep:
sublog := log.With().Str("event", "StartPostprocessingStep").Str("uploadid", ev.UploadID).Logger()
if ev.UploadID == "" {
@@ -1324,7 +1324,12 @@ func (n *Node) DeleteGrant(ctx context.Context, g *provider.Grant) (err error) {
// Purge removes a node from disk. It does not move it to the trash
func (n *Node) Purge(ctx context.Context) error {
return n.lu.PurgeNode(n)
if err := n.lu.PurgeNode(n); err != nil {
return err
}
// remove .mpk and .mlock files
return n.lu.MetadataBackend().Purge(ctx, n)
}
// ListGrants lists all grants of the current node.
@@ -1551,6 +1556,173 @@ func (n *Node) SetDTime(ctx context.Context, t *time.Time) (err error) {
return n.lu.TimeManager().SetDTime(ctx, n, t)
}
// RevertUpload reverts the upload that created the given revision: it
// restores the revision onto the node and removes the revision, including its
// metadata sidecars. The node's metadata lock is held for the duration of the
// operation.
func (n *Node) RevertUpload(ctx context.Context, versionID string) error {
if versionID == "" {
return errors.New("empty versionID")
}
unlock, err := n.lu.MetadataBackend().Lock(n)
if err != nil {
return err
}
defer func() {
_ = unlock()
}()
revisionNode := NewBaseNode(n.SpaceID, n.ID+RevisionIDDelimiter+versionID, n.lu)
revisionPath := revisionNode.InternalPath()
if _, err := os.Stat(revisionPath); err != nil {
appctx.GetLogger(ctx).Error().Str("versionpath", revisionPath).Err(err).Msg("revision does not exist")
return err
}
if err := n.lu.CopyMetadata(ctx, revisionNode, n, func(attributeName string, value []byte) (newValue []byte, copy bool) {
return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) ||
attributeName == prefixes.TypeAttr ||
attributeName == prefixes.BlobIDAttr ||
attributeName == prefixes.BlobsizeAttr ||
attributeName == prefixes.MTimeAttr
}); err != nil {
appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("restoring revision metadata failed")
return err
}
if err := os.RemoveAll(revisionPath); err != nil {
appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("error removing version")
return err
}
// remove the revision's metadata sidecars
if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionNode)); err != nil {
appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata, continuing")
}
if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionNode)); err != nil {
appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata lockfile, continuing")
}
if err := n.lu.MetadataBackend().Purge(ctx, revisionNode); err != nil {
appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not purge revision from cache, continuing")
}
return nil
}
// DeleteRevision deletes a revision of the node and returns the blob id of the
// deleted revision, so that the caller can delete the blob. If versionID is
// empty the node's current revision is deleted: the latest stored revision is
// restored onto the node (or the node is purged when there is none) and the
// processing flag is removed, unstalling the node. No blob is returned in that
// case.
func (n *Node) DeleteRevision(ctx context.Context, versionID string) (string, error) {
if versionID != "" {
return n.deleteRevision(ctx, n.ID+RevisionIDDelimiter+versionID)
}
revisionPath, err := n.getLatestRevision(ctx)
if err != nil {
return "", err
}
if revisionPath == "" {
// there is no revision - delete the node
unlock, err := n.lu.MetadataBackend().Lock(n)
if err != nil {
return "", err
}
defer func() {
_ = unlock()
}()
if err := n.Purge(ctx); err != nil {
appctx.GetLogger(ctx).Info().Str("nodepath", n.InternalPath()).Err(err).Msg("error purging node")
return "", err
}
return "", nil
}
latestID := strings.TrimPrefix(revisionPath, n.lu.VersionPath(n.SpaceID, n.ID, ""))
if err := n.RevertUpload(ctx, latestID); err != nil {
return "", err
}
// we just deleted the current revision - remove processing flag if set
if uploadid, err := n.ProcessingID(ctx); err == nil {
return "", n.UnmarkProcessing(ctx, uploadid)
}
return "", nil
}
// deleteRevision deletes the revision node identified by revisionID (the node's
// ID + RevisionIDDelimiter + timestamp) together with its metadata sidecars
// and returns its blob id. It is a no-op if the revision does not exist.
func (n *Node) deleteRevision(ctx context.Context, revisionID string) (string, error) {
log := appctx.GetLogger(ctx)
revisionNode := NewBaseNode(n.SpaceID, revisionID, n.lu)
revisionPath := revisionNode.InternalPath()
if _, err := os.Stat(revisionPath); err != nil {
log.Warn().Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("revision does not exist, nothing to delete")
return "", nil
}
unlock, err := n.lu.MetadataBackend().Lock(n)
if err != nil {
return "", err
}
defer func() {
_ = unlock()
}()
blobID, _, err := n.lu.ReadBlobIDAndSizeAttr(ctx, revisionNode, nil)
if err != nil {
return "", err
}
if err := os.RemoveAll(revisionPath); err != nil {
return "", err
}
if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionNode)); err != nil {
log.Warn().Err(err).Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("could not delete revision metadata, continuing")
}
if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionNode)); err != nil {
log.Warn().Err(err).Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("could not delete revision metadata lockfile, continuing")
}
if err := n.lu.MetadataBackend().Purge(ctx, revisionNode); err != nil {
log.Warn().Err(err).Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("could not purge revision from cache, continuing")
}
return blobID, nil
}
func (n *Node) getLatestRevision(ctx context.Context) (string, error) {
revPrefix := n.lu.VersionPath(n.SpaceID, n.ID, "")
revisions, err := filepath.Glob(n.lu.VersionPath(n.SpaceID, n.ID, "*"))
if err != nil {
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error reading revisions")
return "", err
}
revPath, latest := "", time.Time{}
for _, rev := range revisions {
if strings.HasSuffix(rev, ".mpk") || strings.HasSuffix(rev, ".mlock") {
continue
}
revDate, err := time.Parse(time.RFC3339Nano, strings.TrimPrefix(rev, revPrefix))
if err != nil {
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error parsing revision date")
continue
}
if revDate.After(latest) {
latest = revDate
revPath = rev
}
}
return revPath, nil
}
// ReadChildNodeFromLink reads the child node id from a link
func ReadChildNodeFromLink(ctx context.Context, path string) (string, error) {
_, span := tracer.Start(ctx, "readChildNodeFromLink")
@@ -162,11 +162,28 @@ func (fs *Decomposedfs) DeleteRevision(ctx context.Context, ref *provider.Refere
return err
}
if err := os.RemoveAll(fs.lu.InternalPath(n.SpaceID, revisionKey)); err != nil {
return err
return fs.deleteRevisionFile(ctx, n, revisionKey)
}
// deleteRevisionFile deletes the revision node identified by revisionKey
// (nodeID + RevisionIDDelimiter + timestamp) together with its metadata
// sidecars and its own blob. It is a no-op if the revision does not exist.
func (fs *Decomposedfs) deleteRevisionFile(ctx context.Context, n *node.Node, revisionKey string) error {
kp := strings.SplitN(revisionKey, node.RevisionIDDelimiter, 2)
if len(kp) != 2 {
return errtypes.NotFound(revisionKey)
}
return fs.tp.DeleteBlob(n)
blobID, err := n.DeleteRevision(ctx, kp[1])
if err != nil {
return err
}
if blobID == "" {
// no blob to delete (0-byte file or current-revision revert)
return nil
}
return fs.tp.DeleteBlob(&node.Node{BaseNode: node.BaseNode{SpaceID: n.SpaceID}, BlobID: blobID})
}
func (fs *Decomposedfs) getRevisionNode(ctx context.Context, ref *provider.Reference, revisionKey string, hasPermission func(*provider.ResourcePermissions) bool) (*node.Node, error) {
@@ -275,62 +275,6 @@ func (tp *Tree) DownloadRevision(ctx context.Context, ref *provider.Reference, r
return ri, reader, nil
}
// DeleteRevision deletes the specified revision of the resource
func (tp *Tree) DeleteRevision(ctx context.Context, ref *provider.Reference, revisionKey string) error {
_, span := tracer.Start(ctx, "DeleteRevision")
defer span.End()
n, err := tp.getRevisionNode(ctx, ref, revisionKey, func(rp *provider.ResourcePermissions) bool {
return rp.RestoreFileVersion
})
if err != nil {
return err
}
if err := os.RemoveAll(tp.lookup.InternalPath(n.SpaceID, revisionKey)); err != nil {
return err
}
return tp.DeleteBlob(n)
}
func (tp *Tree) getRevisionNode(ctx context.Context, ref *provider.Reference, revisionKey string, hasPermission func(*provider.ResourcePermissions) bool) (*node.Node, error) {
_, span := tracer.Start(ctx, "getRevisionNode")
defer span.End()
log := appctx.GetLogger(ctx)
// verify revision key format
kp := strings.SplitN(revisionKey, node.RevisionIDDelimiter, 2)
if len(kp) != 2 {
log.Error().Str("revisionKey", revisionKey).Msg("malformed revisionKey")
return nil, errtypes.NotFound(revisionKey)
}
log.Debug().Str("revisionKey", revisionKey).Msg("DownloadRevision")
spaceID := ref.ResourceId.SpaceId
// check if the node is available and has not been deleted
n, err := node.ReadNode(ctx, tp.lookup, spaceID, kp[0], "", false, nil, false)
if err != nil {
return nil, err
}
if !n.Exists {
err = errtypes.NotFound(filepath.Join(n.ParentID, n.Name))
return nil, err
}
p, err := tp.permissions.AssemblePermissions(ctx, n)
switch {
case err != nil:
return nil, err
case !hasPermission(p):
return nil, errtypes.PermissionDenied(filepath.Join(n.ParentID, n.Name))
}
// Set space owner in context
storagespace.ContextSendSpaceOwnerID(ctx, n.SpaceOwnerOrManager(ctx))
return n, nil
}
func (tp *Tree) RestoreRevision(ctx context.Context, sourceNode, targetNode metadata.MetadataNode, mtime time.Time) error {
err := tp.lookup.CopyMetadata(ctx, sourceNode, targetNode, func(attributeName string, value []byte) (newValue []byte, copy bool) {
return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) ||
@@ -353,22 +353,29 @@ func (session *DecomposedFsSession) Finalize(ctx context.Context) (err error) {
// another upload on this node is in progress or has finished since we started
if !isProcessing || processingID != session.ID() {
versionID := n.ID + node.RevisionIDDelimiter + session.MTime().UTC().Format(time.RFC3339Nano)
// There should be a revision node (created by the other upload that finished before us), read it and upload our blob there.
existingRevisionNode, revisionNodeUnlock, err := node.LockAndReadNode(ctx, session.store.lu, session.SpaceID(), versionID, "", false, spaceRoot, false)
if err != nil || !existingRevisionNode.Exists {
// The revision node has not been created. Likely because the file on disk was modified externally and re-assilimated (watchfs == true)
// Let's create the revision node now and upload the blob to it.
n, revisionNodeUnlock, err = session.createRevisionNodeForUpload(ctx, n, session.MTime().UTC().Format(time.RFC3339Nano))
if err != nil {
appctx.GetLogger(ctx).Debug().Err(err).Str("versionID", session.MTime().UTC().Format(time.RFC3339Nano)).Msg("failed to create revision node for upload finalization")
return err
}
// The node's current content is already this upload's content (e.g. a
// newer upload was reverted in the meantime) - no revision is needed,
// upload the blob to the node itself.
if attribs.String(prefixes.BlobIDAttr) == session.ID() {
appctx.GetLogger(ctx).Debug().Str("nodepath", n.InternalPath()).Msg("node already contains this upload's content, no revision needed")
} else {
n = existingRevisionNode
versionID := n.ID + node.RevisionIDDelimiter + session.MTime().UTC().Format(time.RFC3339Nano)
// There should be a revision node (created by the other upload that finished before us), read it and upload our blob there.
existingRevisionNode, revisionNodeUnlock, err := node.LockAndReadNode(ctx, session.store.lu, session.SpaceID(), versionID, "", false, spaceRoot, false)
if err != nil || !existingRevisionNode.Exists {
// The revision node has not been created. Likely because the file on disk was modified externally and re-assilimated (watchfs == true)
// Let's create the revision node now and upload the blob to it.
n, revisionNodeUnlock, err = session.createRevisionNodeForUpload(ctx, n, session.MTime().UTC().Format(time.RFC3339Nano))
if err != nil {
appctx.GetLogger(ctx).Debug().Err(err).Str("versionID", session.MTime().UTC().Format(time.RFC3339Nano)).Msg("failed to create revision node for upload finalization")
return err
}
} else {
n = existingRevisionNode
}
appctx.GetLogger(ctx).Debug().Str("new nodepath", n.InternalPath()).Msg("uploading to revision node, that was created for us by another upload")
defer func() { _ = revisionNodeUnlock() }()
}
appctx.GetLogger(ctx).Debug().Str("new nodepath", n.InternalPath()).Msg("uploading to revision node, that was created for us by another upload")
defer func() { _ = revisionNodeUnlock() }()
}
// upload the data to the blobstore
@@ -433,17 +440,6 @@ func checkHash(expected string, h hash.Hash) error {
return nil
}
func (session *DecomposedFsSession) removeNode(ctx context.Context) {
n, err := session.Node(ctx)
if err != nil {
appctx.GetLogger(ctx).Error().Str("session", session.ID()).Err(err).Msg("getting node from session failed")
return
}
if err := n.Purge(ctx); err != nil {
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("purging node failed")
}
}
// cleanup cleans up after the upload is finished
func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanInfo, unmarkPostprocessing bool) {
ctx := session.Context(context.Background())
@@ -455,35 +451,13 @@ func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanI
if err != nil {
sublog.Error().Err(err).Msg("reading node for session failed")
} else {
if session.NodeExists() && session.info.MetaData["versionID"] != "" {
versionID := session.info.MetaData["versionID"]
versionID := strings.TrimPrefix(session.info.MetaData["versionID"], n.ID+node.RevisionIDDelimiter)
if session.NodeExists() && versionID != "" {
sublog.Debug().Str("nodepath", n.InternalPath()).Str("versionID", versionID).Msg("restoring revision")
revisionNode, err := node.ReadNode(ctx, session.store.lu, session.SpaceID(), versionID, "", false, n.SpaceRoot, false)
if err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("reading revision node failed")
if err := n.RevertUpload(ctx, versionID); err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("reverting node metadata failed")
return
}
if !revisionNode.Exists {
sublog.Error().Str("versionID", versionID).Msg("revision node does not exist")
return
}
// restore the revision
mtime, err := revisionNode.GetMTime(ctx)
if err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("getting mtime of revision node failed")
mtime = time.Now()
}
if err := session.store.tp.RestoreRevision(ctx, revisionNode, n, mtime); err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("restoring revision node failed")
return
}
if err := os.RemoveAll(revisionNode.InternalPath()); err != nil {
sublog.Error().Err(err).Str("revisionpath", revisionNode.InternalPath()).Msg("removing restored revision file failed")
}
} else {
// if no other upload session is in progress (processing id != session id) or has finished (processing id == "")
latestSession, err := n.ProcessingID(ctx)
@@ -492,7 +466,9 @@ func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanI
}
if latestSession == session.ID() {
// actually delete the node
session.removeNode(ctx)
if err := n.Purge(ctx); err != nil {
sublog.Error().Err(err).Str("nodepath", n.InternalPath()).Msg("purging node failed")
}
}
// FIXME else if the upload has become a revision, delete the revision, or if it is the last one, delete the node
}
@@ -83,6 +83,7 @@ var (
events.PostprocessingStepFinished{},
events.RestartPostprocessing{},
events.CleanUpload{},
events.DeleteRevision{},
}
)
@@ -452,6 +453,31 @@ func (fs *Decomposedfs) handlePostprocessingEvent(ctx context.Context, event eve
return // NOTE: since we can't get the upload, we can't delete the blob
}
session.Cleanup(true, !ev.KeepUpload, !ev.KeepUpload, true)
case events.DeleteRevision:
sublog := log.With().Str("event", "DeleteRevision").Interface("nodeid", ev.ResourceID).Logger()
n, err := fs.lu.NodeFromID(ctx, ev.ResourceID)
if err != nil {
sublog.Error().Err(err).Msg("Failed to get node")
return
}
var deleteErr error
if ev.Timestamp == nil {
// the node's current revision is targeted - revert to the pre-upload
// state. Only stuck nodes are targeted: a reverted node has its
// processing flag removed, so a redelivered event is a no-op.
if !n.IsProcessing(ctx) {
sublog.Debug().Msg("node is not stuck, ignoring")
return
}
_, deleteErr = n.DeleteRevision(ctx, "")
} else {
versionID := time.Unix(int64(ev.Timestamp.Seconds), int64(ev.Timestamp.Nanos)).UTC().Format(time.RFC3339Nano)
deleteErr = fs.deleteRevisionFile(ctx, n, n.ID+node.RevisionIDDelimiter+versionID)
}
if deleteErr != nil {
sublog.Error().Err(deleteErr).Msg("Failed to delete revision")
}
case events.PostprocessingStepFinished:
sublog := log.With().Str("event", "PostprocessingStepFinished").Str("uploadid", ev.UploadID).Logger()
if ev.FinishedStep != events.PPStepAntivirus {
@@ -1165,7 +1165,12 @@ func (n *Node) Purge(ctx context.Context) error {
// remove child entry in parent
src := filepath.Join(n.ParentPath(), n.Name)
return os.Remove(src)
if err := os.Remove(src); err != nil {
return err
}
// remove .mpk and .mlock files
return n.lu.MetadataBackend().Purge(ctx, n.InternalPath())
}
// ListGrants lists all grants of the current node.
@@ -1391,3 +1396,164 @@ func (n *Node) GetDTime(ctx context.Context) (time.Time, error) {
func (n *Node) SetDTime(ctx context.Context, t *time.Time) (err error) {
return n.lu.TimeManager().SetDTime(ctx, n, t)
}
// RevertUpload reverts the upload that created the given revision: it
// restores the revision onto the node and removes the revision, including its
// metadata sidecars. The node's metadata lock is held for the duration of the
// operation.
func (n *Node) RevertUpload(ctx context.Context, versionID string) error {
if versionID == "" {
return errors.New("empty versionID")
}
lock, err := lockedfile.OpenFile(n.lu.MetadataBackend().LockfilePath(n.InternalPath()), os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return err
}
defer lock.Close()
revisionPath := n.InternalPath() + RevisionIDDelimiter + versionID
if _, err := os.Stat(revisionPath); err != nil {
appctx.GetLogger(ctx).Error().Str("versionpath", revisionPath).Err(err).Msg("revision does not exist")
return err
}
if err := n.lu.CopyMetadata(ctx, revisionPath, n.InternalPath(), func(attributeName string, value []byte) (newValue []byte, copy bool) {
return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) ||
attributeName == prefixes.TypeAttr ||
attributeName == prefixes.BlobIDAttr ||
attributeName == prefixes.BlobsizeAttr ||
attributeName == prefixes.MTimeAttr
}, false); err != nil {
appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("restoring revision metadata failed")
return err
}
if err := os.RemoveAll(revisionPath); err != nil {
appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("error removing version")
return err
}
// remove the revision's metadata sidecars
if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionPath)); err != nil {
appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata, continuing")
}
if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionPath)); err != nil {
appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata lockfile, continuing")
}
if err := n.lu.MetadataBackend().Purge(ctx, revisionPath); err != nil {
appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not purge revision from cache, continuing")
}
return nil
}
// DeleteRevision deletes a revision of the node and returns the blob id of the
// deleted revision, so that the caller can delete the blob. If versionID is
// empty the node's current revision is deleted: the latest stored revision is
// restored onto the node (or the node is purged when there is none) and the
// processing flag is removed, unstalling the node. No blob is returned in that
// case.
func (n *Node) DeleteRevision(ctx context.Context, versionID string) (string, error) {
if versionID != "" {
return n.deleteRevision(ctx, versionID)
}
revisionPath, err := n.getLatestRevision(ctx)
if err != nil {
return "", err
}
if revisionPath == "" {
// there is no revision - delete the node
lock, err := lockedfile.OpenFile(n.lu.MetadataBackend().LockfilePath(n.InternalPath()), os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return "", err
}
defer lock.Close()
if err := n.Purge(ctx); err != nil {
appctx.GetLogger(ctx).Info().Str("nodepath", n.InternalPath()).Err(err).Msg("error purging node")
return "", err
}
return "", nil
}
latestID := strings.TrimPrefix(revisionPath, n.InternalPath()+RevisionIDDelimiter)
if err := n.RevertUpload(ctx, latestID); err != nil {
return "", err
}
// we just deleted the current revision - remove processing flag if set
if uploadid, err := n.ProcessingID(ctx); err == nil {
return "", n.UnmarkProcessing(ctx, uploadid)
}
return "", nil
}
// deleteRevision deletes the revision identified by versionID (the timestamp
// following the node's InternalPath + RevisionIDDelimiter) together with its
// metadata sidecars and returns its blob id. It is a no-op if the revision
// does not exist.
func (n *Node) deleteRevision(ctx context.Context, versionID string) (string, error) {
log := appctx.GetLogger(ctx)
revisionPath := n.InternalPath() + RevisionIDDelimiter + versionID
if _, err := os.Stat(revisionPath); err != nil {
log.Warn().Str("nodeid", n.ID).Str("versionid", versionID).Msg("revision does not exist, nothing to delete")
return "", nil
}
lock, err := lockedfile.OpenFile(n.lu.MetadataBackend().LockfilePath(n.InternalPath()), os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return "", err
}
defer lock.Close()
blobID, _, err := n.lu.ReadBlobIDAndSizeAttr(ctx, revisionPath, nil)
if err != nil {
return "", err
}
if err := os.RemoveAll(revisionPath); err != nil {
return "", err
}
if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionPath)); err != nil {
log.Warn().Err(err).Str("nodeid", n.ID).Str("versionid", versionID).Msg("could not delete revision metadata, continuing")
}
if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionPath)); err != nil {
log.Warn().Err(err).Str("nodeid", n.ID).Str("versionid", versionID).Msg("could not delete revision metadata lockfile, continuing")
}
if err := n.lu.MetadataBackend().Purge(ctx, revisionPath); err != nil {
log.Warn().Err(err).Str("nodeid", n.ID).Str("versionid", versionID).Msg("could not purge revision from cache, continuing")
}
return blobID, nil
}
func (n *Node) getLatestRevision(ctx context.Context) (string, error) {
revPrefix := n.InternalPath() + RevisionIDDelimiter
revisions, err := filepath.Glob(revPrefix + "*")
if err != nil {
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error reading revisions")
return "", err
}
revPath, latest := "", time.Time{}
for _, rev := range revisions {
if strings.HasSuffix(rev, ".mpk") || strings.HasSuffix(rev, ".mlock") {
continue
}
revDate, err := time.Parse(time.RFC3339Nano, strings.TrimPrefix(rev, revPrefix))
if err != nil {
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error parsing revision date")
continue
}
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Str("revPath", revPath).Interface("time", revDate).Err(err).Msg("error parsing revision date")
if revDate.After(latest) {
latest = revDate
revPath = rev
}
}
return revPath, nil
}
@@ -345,11 +345,28 @@ func (fs *Decomposedfs) DeleteRevision(ctx context.Context, ref *provider.Refere
return err
}
if err := os.RemoveAll(fs.lu.InternalPath(n.SpaceID, revisionKey)); err != nil {
return err
return fs.deleteRevisionFile(ctx, n, revisionKey)
}
// deleteRevisionFile deletes the revision node identified by revisionKey
// (nodeID + RevisionIDDelimiter + timestamp) together with its metadata
// sidecars and its own blob. It is a no-op if the revision does not exist.
func (fs *Decomposedfs) deleteRevisionFile(ctx context.Context, n *node.Node, revisionKey string) error {
kp := strings.SplitN(revisionKey, node.RevisionIDDelimiter, 2)
if len(kp) != 2 {
return errtypes.NotFound(revisionKey)
}
return fs.tp.DeleteBlob(n)
blobID, err := n.DeleteRevision(ctx, kp[1])
if err != nil {
return err
}
if blobID == "" {
// no blob to delete (0-byte file or current-revision revert)
return nil
}
return fs.tp.DeleteBlob(&node.Node{SpaceID: n.SpaceID, BlobID: blobID})
}
func (fs *Decomposedfs) getRevisionNode(ctx context.Context, ref *provider.Reference, revisionKey string, hasPermission func(*provider.ResourcePermissions) bool) (*node.Node, error) {
@@ -316,57 +316,10 @@ func checkHash(expected string, h hash.Hash) error {
return nil
}
func (session *OcisSession) removeNode(ctx context.Context) {
n, err := session.Node(ctx)
if err != nil {
appctx.GetLogger(ctx).Error().Str("session", session.ID()).Err(err).Msg("getting node from session failed")
return
}
if err := n.Purge(ctx); err != nil {
appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("purging node failed")
}
}
// cleanup cleans up after the upload is finished
func (session *OcisSession) Cleanup(revertNodeMetadata, cleanBin, cleanInfo, unmarkPostprocessing bool) {
ctx := session.Context(context.Background())
if revertNodeMetadata {
n, err := session.Node(ctx)
if err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("sessionid", session.ID()).Msg("reading node for session failed")
} else {
if session.NodeExists() && session.info.MetaData["versionsPath"] != "" {
p := session.info.MetaData["versionsPath"]
if err := session.store.lu.CopyMetadata(ctx, p, n.InternalPath(), func(attributeName string, value []byte) (newValue []byte, copy bool) {
return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) ||
attributeName == prefixes.TypeAttr ||
attributeName == prefixes.BlobIDAttr ||
attributeName == prefixes.BlobsizeAttr ||
attributeName == prefixes.MTimeAttr
}, true); err != nil {
appctx.GetLogger(ctx).Info().Str("versionpath", p).Str("nodepath", n.InternalPath()).Err(err).Msg("renaming version node failed")
}
if err := os.RemoveAll(p); err != nil {
appctx.GetLogger(ctx).Info().Str("versionpath", p).Str("nodepath", n.InternalPath()).Err(err).Msg("error removing version")
}
} else {
// if no other upload session is in progress (processing id != session id) or has finished (processing id == "")
latestSession, err := n.ProcessingID(ctx)
if err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("spaceid", n.SpaceID).Str("nodeid", n.ID).Str("uploadid", session.ID()).Msg("reading processingid for session failed")
}
if latestSession == session.ID() {
// actually delete the node
session.removeNode(ctx)
}
// FIXME else if the upload has become a revision, delete the revision, or if it is the last one, delete the node
}
}
}
if cleanBin {
if err := os.Remove(session.binPath()); err != nil && !errors.Is(err, fs.ErrNotExist) {
appctx.GetLogger(ctx).Error().Str("path", session.binPath()).Err(err).Msg("removing upload failed")
@@ -376,8 +329,37 @@ func (session *OcisSession) Cleanup(revertNodeMetadata, cleanBin, cleanInfo, unm
if cleanInfo {
if err := os.Remove(session.infoPath()); err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("session", session.ID()).Msg("removing upload info failed")
}
}
if revertNodeMetadata {
n, err := session.Node(ctx)
if err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("sessionid", session.ID()).Msg("reading node for session failed")
return
}
versionID := strings.TrimPrefix(session.info.MetaData["versionsPath"], n.InternalPath()+node.RevisionIDDelimiter)
if session.NodeExists() && versionID != "" {
if err := n.RevertUpload(ctx, versionID); err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("nodepath", n.InternalPath()).Msg("reverting node metadata failed")
return
}
} else {
// if no other upload session is in progress (processing id != session id) or has finished (processing id == "")
latestSession, err := n.ProcessingID(ctx)
if err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("spaceid", n.SpaceID).Str("nodeid", n.ID).Str("uploadid", session.ID()).Msg("reading processingid for session failed")
}
if latestSession == session.ID() {
// actually delete the node
if err := n.Purge(ctx); err != nil {
appctx.GetLogger(ctx).Error().Err(err).Str("nodepath", n.InternalPath()).Msg("purging node failed")
return
}
}
// FIXME else if the upload has become a revision, delete the revision, or if it is the last one, delete the node
}
}
if unmarkPostprocessing {
+2 -1
View File
@@ -1362,7 +1362,7 @@ github.com/opencloud-eu/icap-client
# github.com/opencloud-eu/libre-graph-api-go v1.0.8-0.20260902170011-45af3945a067
## explicit; go 1.23
github.com/opencloud-eu/libre-graph-api-go
# github.com/opencloud-eu/reva/v2 v2.50.1-0.20261001091108-11d87fb6b985
# github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002075447-9cf4d4365f4f
## explicit; go 1.26.0
github.com/opencloud-eu/reva/v2/cmd/revad/internal/grace
github.com/opencloud-eu/reva/v2/cmd/revad/runtime
@@ -1472,6 +1472,7 @@ github.com/opencloud-eu/reva/v2/pkg/appctx
github.com/opencloud-eu/reva/v2/pkg/auth
github.com/opencloud-eu/reva/v2/pkg/auth/manager/appauth
github.com/opencloud-eu/reva/v2/pkg/auth/manager/demo
github.com/opencloud-eu/reva/v2/pkg/auth/manager/guestlinks
github.com/opencloud-eu/reva/v2/pkg/auth/manager/impersonator
github.com/opencloud-eu/reva/v2/pkg/auth/manager/json
github.com/opencloud-eu/reva/v2/pkg/auth/manager/ldap