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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 5 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module github.com/owncloud/reva/v2

go 1.26
go 1.25.11

require (
bou.ke/monkey v1.0.2
Expand Down Expand Up @@ -84,7 +84,7 @@ require (
github.com/tus/tusd/v2 v2.10.0
github.com/wk8/go-ordered-map v1.0.0
go-micro.dev/v4 v4.11.0
go.etcd.io/etcd/client/v3 v3.7.1
go.etcd.io/etcd/client/v3 v3.6.13
go.opencensus.io v0.24.0
go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.69.0
go.opentelemetry.io/otel v1.44.0
Expand Down Expand Up @@ -148,6 +148,7 @@ require (
github.com/go-openapi/strfmt v0.26.2 // indirect
github.com/go-task/slim-sprig/v3 v3.0.0 // indirect
github.com/go-viper/mapstructure/v2 v2.5.0 // indirect
github.com/gogo/protobuf v1.3.2 // indirect
github.com/golang/groupcache v0.0.0-20241129210726-2c02b8208cf8 // indirect
github.com/google/go-querystring v1.1.0 // indirect
github.com/google/go-tpm v0.9.8 // indirect
Expand Down Expand Up @@ -204,8 +205,8 @@ require (
github.com/xrash/smetrics v0.0.0-20240521201337-686a1a2994c1 // indirect
github.com/yusufpapurcu/wmi v1.2.4 // indirect
github.com/zeebo/xxh3 v1.1.0 // indirect
go.etcd.io/etcd/api/v3 v3.7.1 // indirect
go.etcd.io/etcd/client/pkg/v3 v3.7.1 // indirect
go.etcd.io/etcd/api/v3 v3.6.13 // indirect
go.etcd.io/etcd/client/pkg/v3 v3.6.13 // indirect
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.44.0 // indirect
go.opentelemetry.io/otel/metric v1.44.0 // indirect
Expand Down
13 changes: 7 additions & 6 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,7 @@ github.com/gofrs/uuid/v5 v5.4.0 h1:EfbpCTjqMuGyq5ZJwxqzn3Cbr2d0rUZU7v5ycAk/e/0=
github.com/gofrs/uuid/v5 v5.4.0/go.mod h1:CDOjlDMVAtN56jqyRUZh58JT31Tiw7/oQyEXZV+9bD8=
github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ=
github.com/gogo/protobuf v1.2.1/go.mod h1:hp+jE20tsWTFYpLwKvXlhS1hjn+gTNwPg2I6zVXpSg4=
github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q=
github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q=
github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY=
github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
Expand Down Expand Up @@ -700,12 +701,12 @@ github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ=
github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0=
github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs=
github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s=
go.etcd.io/etcd/api/v3 v3.7.1 h1:KJG0/DcWGfe3Y1otDf/fsBf0TSSgpxZ5RO/L8SFt73E=
go.etcd.io/etcd/api/v3 v3.7.1/go.mod h1:8bXIpCMeV7E3/XL0Ix123ATn3dB+0V7d9zklHbB0m78=
go.etcd.io/etcd/client/pkg/v3 v3.7.1 h1:rKYsj3pRkR0eK3yjT3XOgrhqfmIfj9pzNgxjh7mfFv4=
go.etcd.io/etcd/client/pkg/v3 v3.7.1/go.mod h1:cnzZGIUzSfjEwLC6UBVsSXlEK1eepS/JUD7wE6PLRT0=
go.etcd.io/etcd/client/v3 v3.7.1 h1:0PEMMC0KuZmVIN+RAbdqfkZ45pYTgKVtmBEbRCvZFUg=
go.etcd.io/etcd/client/v3 v3.7.1/go.mod h1:ffNqALa8tRCYhYo1F9oR489y23K39Gz+BSR3ApAGYq0=
go.etcd.io/etcd/api/v3 v3.6.13 h1:AvHPZv15LYEe7tZDyFglv7xnbiuF6GMZpZqKpIzXTt0=
go.etcd.io/etcd/api/v3 v3.6.13/go.mod h1:X9+3gaKwzjlOxzo6TZ2u3b7HcHBcAL+Ph7EBPjI/VWk=
go.etcd.io/etcd/client/pkg/v3 v3.6.13 h1:7QeMOisYByx8dBA7/CKcwCaPWfjb5C0xpmrIov/8WyY=
go.etcd.io/etcd/client/pkg/v3 v3.6.13/go.mod h1:Dn2zUBOCu/6xYcd6iAjB7LgoY16OTQjDZfWHLwvuQj4=
go.etcd.io/etcd/client/v3 v3.6.13 h1:0E+9ZYGpMsi9KlOJVoCdONh9PUDawKDTy5mSNY8wOEI=
go.etcd.io/etcd/client/v3 v3.6.13/go.mod h1:rtVI3vwobljb8xlTGcp1Yhz7hBIuBWULXwB848kqJGw=
go.opencensus.io v0.21.0/go.mod h1:mSImk1erAIZhrmZN+AvHh14ztQfjbGwt4TtuofqLduU=
go.opencensus.io v0.22.0/go.mod h1:+kGneAE2xo2IficOXnaByMWTGM9T73dGwxeWcUqIpI8=
go.opencensus.io v0.22.2/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw=
Expand Down
38 changes: 30 additions & 8 deletions internal/grpc/services/storageprovider/storageprovider.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import (
"net/url"
"os"
"path"
"path/filepath"
"sort"
"strconv"
"strings"
Expand All @@ -47,6 +48,7 @@ import (
"github.com/owncloud/reva/v2/pkg/storage"
"github.com/owncloud/reva/v2/pkg/storage/fs/registry"
"github.com/owncloud/reva/v2/pkg/storagespace"
"github.com/owncloud/reva/v2/pkg/upload"
"github.com/owncloud/reva/v2/pkg/utils"
"github.com/pkg/errors"
"github.com/rs/zerolog"
Expand All @@ -69,6 +71,7 @@ type config struct {
AvailableXS map[string]uint32 `mapstructure:"available_checksums" docs:"nil;List of available checksums."`
CustomMimeTypesJSON string `mapstructure:"custom_mimetypes_json" docs:"nil;An optional mapping file with the list of supported custom file extensions and corresponding mime types."`
MountID string `mapstructure:"mount_id"`
UploadDirectory string `mapstructure:"upload_directory" docs:";Local directory for staging upload sessions. Overrides the driver's root. Required for drivers that have no local filesystem root."`
UploadExpiration int64 `mapstructure:"upload_expiration" docs:"0;Duration for how long uploads will be valid."`
Events eventconfig `mapstructure:"events" docs:"0;Event stream configuration"`
}
Expand Down Expand Up @@ -106,6 +109,7 @@ func (c *config) init() {
type Service struct {
conf *config
Storage storage.FS
Coordinator upload.Coordinator
dataServerURL *url.URL
availableXS []*provider.ResourceChecksumPriority
}
Expand Down Expand Up @@ -175,11 +179,33 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc.

c.init()

fs, err := getFS(c, log)
evstream, err := estreamFromConfig(c.Events)
if err != nil {
return nil, err
}

fs, err := getFS(c, evstream, log)
if err != nil {
return nil, err
}

// Build the coordinator-owned session store. UploadDirectory (service level)
// takes precedence over the driver root, so rootless drivers can still get a
// coordinator. The store points at the same root as the driver's data path so
// it can read the sessions the coordinator writes (decomposedfs on-disk format).
store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log)
if store == nil {
return nil, fmt.Errorf("storageprovider: cannot determine upload directory, set upload_directory in config or driver root")
}
if err := store.Setup(); err != nil {
return nil, fmt.Errorf("storageprovider: upload directory setup failed: %w", err)
}
// Deliberately no StartPostprocessing here: this service initiates uploads but
// never receives the bytes, so it has nothing to hand to postprocessing or to
// commit afterwards. That belongs to the data provider, which owns the PUT and
// TUS paths.
coordinator := upload.NewCoordinator(fs, store, filepath.Join(store.Root(), "uploads"), evstream)

// parse data server url
u, err := url.Parse(c.DataServerURL)
if err != nil {
Expand All @@ -205,6 +231,7 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc.
service := &Service{
conf: c,
Storage: fs,
Coordinator: coordinator,
dataServerURL: u,
availableXS: xsTypes,
}
Expand Down Expand Up @@ -427,7 +454,7 @@ func (s *Service) InitiateFileUpload(ctx context.Context, req *provider.Initiate
metadata["expires"] = strconv.Itoa(int(expirationTimestamp.Seconds))
}

uploadIDs, err := s.Storage.InitiateUpload(ctx, req.Ref, uploadLength, metadata)
uploadIDs, err := s.Coordinator.InitiateUpload(ctx, req.Ref, uploadLength, metadata)
if err != nil {
var st *rpc.Status
switch err.(type) {
Expand Down Expand Up @@ -1266,12 +1293,7 @@ func (s *Service) addMissingStorageProviderID(resourceID *provider.ResourceId, s
}
}

func getFS(c *config, log *zerolog.Logger) (storage.FS, error) {
evstream, err := estreamFromConfig(c.Events)
if err != nil {
return nil, err
}

func getFS(c *config, evstream events.Stream, log *zerolog.Logger) (storage.FS, error) {
if f, ok := registry.NewFuncs[c.Driver]; ok {
driverConf := c.Drivers[c.Driver]
driverConf["mount_id"] = c.MountID // pass the mount id to the driver
Expand Down
29 changes: 26 additions & 3 deletions internal/http/services/dataprovider/dataprovider.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ package dataprovider
import (
"fmt"
"net/http"
"path/filepath"

"github.com/mitchellh/mapstructure"
"github.com/rs/zerolog"
Expand All @@ -33,6 +34,7 @@ import (
"github.com/owncloud/reva/v2/pkg/rhttp/router"
"github.com/owncloud/reva/v2/pkg/storage"
"github.com/owncloud/reva/v2/pkg/storage/fs/registry"
pkgupload "github.com/owncloud/reva/v2/pkg/upload"
)

func init() {
Expand All @@ -51,6 +53,7 @@ type config struct {
NatsEnableTLS bool `mapstructure:"nats_enable_tls"`
NatsUsername string `mapstructure:"nats_username"`
NatsPassword string `mapstructure:"nats_password"`
UploadDirectory string `mapstructure:"upload_directory" docs:";Local directory for staging upload sessions. Overrides the driver root. Required for drivers without a local root."`
}

func (c *config) init() {
Expand Down Expand Up @@ -104,7 +107,27 @@ func New(m map[string]interface{}, log *zerolog.Logger) (global.Service, error)
return nil, err
}

dataTXs, err := getDataTXs(conf, fs, evstream, log)
// The data provider finishes uploads that the storage provider initiated, so
// its store must resolve to the same upload directory.
store := pkgupload.NewFileStoreFromConfig(conf.UploadDirectory, conf.Drivers[conf.Driver], log)
if store == nil {
return nil, fmt.Errorf("dataprovider: cannot determine upload directory, set upload_directory in config or driver root")
}
if err := store.Setup(); err != nil {
return nil, fmt.Errorf("dataprovider: upload directory setup failed: %w", err)
}
coord := pkgupload.NewCoordinator(fs, store, filepath.Join(store.Root(), "uploads"), evstream)

// This is the service that receives the bytes, so it is the one that hands them
// to postprocessing and commits them once the verdict comes back. Without this
// the coordinator commits inline and uploads are never scanned.
if ac := pkgupload.AsyncConfFromDriverConf(conf.Drivers[conf.Driver]); ac.Enabled {
if err := coord.RunPostprocessingConsumer(evstream, ac); err != nil {
return nil, fmt.Errorf("dataprovider: could not start postprocessing: %w", err)
}
}

dataTXs, err := getDataTXs(conf, coord, fs, evstream, log)
if err != nil {
return nil, err
}
Expand All @@ -126,7 +149,7 @@ func getFS(c *config, stream events.Stream, log *zerolog.Logger) (storage.FS, er
return nil, fmt.Errorf("driver not found: %s", c.Driver)
}

func getDataTXs(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (map[string]http.Handler, error) {
func getDataTXs(c *config, coord pkgupload.Coordinator, driver storage.FS, publisher events.Publisher, log *zerolog.Logger) (map[string]http.Handler, error) {
if c.DataTXs == nil {
c.DataTXs = make(map[string]map[string]interface{})
}
Expand All @@ -146,7 +169,7 @@ func getDataTXs(c *config, fs storage.FS, publisher events.Publisher, log *zerol
for t := range c.DataTXs {
if f, ok := datatxregistry.NewFuncs[t]; ok {
if tx, err := f(c.DataTXs[t], publisher, log); err == nil {
if handler, err := tx.Handler(fs); err == nil {
if handler, err := tx.Handler(coord, driver); err == nil {
txs[t] = handler
}
}
Expand Down
6 changes: 5 additions & 1 deletion internal/http/services/owncloud/ocdav/tus.go
Original file line number Diff line number Diff line change
Expand Up @@ -319,7 +319,11 @@ func (s *svc) handleTusPost(ctx context.Context, w http.ResponseWriter, r *http.
sReq.Ref.Path = uReq.Ref.GetPath()
sReq.Ref.ResourceId = nil
} else {
if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil {
// A new file has no node id until the upload finishes, which is after the
// data server wrote this header, so it may be absent. sReq.Ref then keeps
// the path-based reference it was built with, which already names the
// file and stats just as well.
if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil && resid.GetOpaqueId() != "" {
sReq.Ref = &provider.Reference{
ResourceId: &resid,
}
Expand Down
6 changes: 5 additions & 1 deletion pkg/rhttp/datatx/datatx.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,12 +28,16 @@ import (
provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
"github.com/owncloud/reva/v2/pkg/events"
"github.com/owncloud/reva/v2/pkg/storage"
pkgupload "github.com/owncloud/reva/v2/pkg/upload"
"github.com/owncloud/reva/v2/pkg/utils"
)

// DataTX provides an abstraction around various data transfer protocols.
//
// Uploads go through the coordinator so the finish path is the same for every
// driver; the driver itself is still needed for downloads.
type DataTX interface {
Handler(fs storage.FS) (http.Handler, error)
Handler(coord pkgupload.Coordinator, driver storage.FS) (http.Handler, error)
}

// EmitFileUploadedEvent is a helper function which publishes a FileUploaded event
Expand Down
7 changes: 4 additions & 3 deletions pkg/rhttp/datatx/manager/simple/simple.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import (
"github.com/owncloud/reva/v2/pkg/storage"
"github.com/owncloud/reva/v2/pkg/storage/cache"
"github.com/owncloud/reva/v2/pkg/storagespace"
pkgupload "github.com/owncloud/reva/v2/pkg/upload"
"github.com/owncloud/reva/v2/pkg/utils"
)

Expand Down Expand Up @@ -78,7 +79,7 @@ func New(m map[string]interface{}, publisher events.Publisher, log *zerolog.Logg
}, nil
}

func (m *manager) Handler(fs storage.FS) (http.Handler, error) {
func (m *manager) Handler(coord pkgupload.Coordinator, driver storage.FS) (http.Handler, error) {
h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
sublog := m.log.With().Str("path", r.URL.Path).Logger()
r = r.WithContext(appctx.WithLogger(r.Context(), &sublog))
Expand All @@ -92,7 +93,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) {
metrics.DownloadsActive.Sub(1)
}()
}
download.GetOrHeadFile(w, r, fs, "")
download.GetOrHeadFile(w, r, driver, "")
case "PUT":
metrics.UploadsActive.Add(1)
defer func() {
Expand All @@ -114,7 +115,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) {
ctx = ctxpkg.ContextSetLockID(ctx, lockID)
}

info, err := fs.Upload(ctx, storage.UploadRequest{
info, err := coord.Upload(ctx, storage.UploadRequest{
Ref: ref,
Body: r.Body,
Length: r.ContentLength,
Expand Down
7 changes: 4 additions & 3 deletions pkg/rhttp/datatx/manager/spaces/spaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ import (
"github.com/owncloud/reva/v2/pkg/storage"
"github.com/owncloud/reva/v2/pkg/storage/cache"
"github.com/owncloud/reva/v2/pkg/storagespace"
pkgupload "github.com/owncloud/reva/v2/pkg/upload"
"github.com/owncloud/reva/v2/pkg/utils"
)

Expand Down Expand Up @@ -80,7 +81,7 @@ func New(m map[string]interface{}, publisher events.Publisher, log *zerolog.Logg
}, nil
}

func (m *manager) Handler(fs storage.FS) (http.Handler, error) {
func (m *manager) Handler(coord pkgupload.Coordinator, driver storage.FS) (http.Handler, error) {
h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var spaceID string
spaceID, r.URL.Path = router.ShiftPath(r.URL.Path)
Expand All @@ -97,7 +98,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) {
metrics.DownloadsActive.Sub(1)
}()
}
download.GetOrHeadFile(w, r, fs, spaceID)
download.GetOrHeadFile(w, r, driver, spaceID)
case "PUT":
metrics.UploadsActive.Add(1)
defer func() {
Expand All @@ -117,7 +118,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) {
Path: fn,
}
var info *provider.ResourceInfo
info, err = fs.Upload(ctx, storage.UploadRequest{
info, err = coord.Upload(ctx, storage.UploadRequest{
Ref: ref,
Body: r.Body,
Length: r.ContentLength,
Expand Down
Loading