Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
52862fd
feat(OCISDEV-900): extract upload coordinator out of decomposedfs
LarsJurgensen Jun 30, 2026
42070b3
feat(OCISDEV-900): define generic Session interface; remove OcisSessi…
LarsJurgensen Jun 30, 2026
b421e7c
feat(OCISDEV-900): coordinator owns full upload lifecycle
LarsJurgensen Jun 30, 2026
6dae362
feat(OCISDEV-900): wire coordinator from storageprovider/dataprovider…
LarsJurgensen Jul 1, 2026
61f9a11
feat(OCISDEV-900): replace UploadAborter with Delete; write AV scan r…
LarsJurgensen Jul 1, 2026
1a34b55
feat(OCISDEV-900): check quota in InitiateUpload before accepting bytes
LarsJurgensen Jul 1, 2026
34e4c01
filestore & session implementation
LarsJurgensen Jul 2, 2026
7c62e78
refactor(upload): decouple Coordinator from storage.FS
LarsJurgensen Jul 7, 2026
1ecd3ba
cleanup
LarsJurgensen Jul 8, 2026
3e4032b
tests
LarsJurgensen Jul 8, 2026
d3fe375
revert unused changes
LarsJurgensen Jul 8, 2026
724e094
bug fixes
LarsJurgensen Jul 9, 2026
9af4a28
TouchFile after upload
LarsJurgensen Jul 9, 2026
beefc36
Implement chunking v1 for coordinator
LarsJurgensen Jul 9, 2026
7d8022d
permission checks
LarsJurgensen Jul 9, 2026
cb4e72f
fixes
LarsJurgensen Jul 10, 2026
5691ae6
resolve permission issue
LarsJurgensen Jul 13, 2026
e643b8b
fix group permissions
LarsJurgensen Jul 13, 2026
7bef459
handle processing flag in propfind
LarsJurgensen Jul 13, 2026
27b11ba
test fixes
LarsJurgensen Jul 13, 2026
da2703e
handle missing nodeId
LarsJurgensen Jul 17, 2026
9667062
allow getMD when processing flag is set
LarsJurgensen Jul 21, 2026
b6d910c
feat(OCISDEV-900): new upload flow
LarsJurgensen Jul 27, 2026
ab47734
feat(OCISDEV-900): propagate disable version
LarsJurgensen Jul 28, 2026
02a8e36
feat(OCISDEV-900): set upload directory for nextcloud
LarsJurgensen Jul 28, 2026
dc405ca
feat(OCISDEV-900): fix nextcloud test
LarsJurgensen Jul 29, 2026
8ba9f04
feat(OCISDEV-900): rollbackUpload
LarsJurgensen Aug 3, 2026
255039a
feat(OCISDEV-900): tmp verbose test output
LarsJurgensen Aug 6, 2026
a5a194b
feat(OCISDEV-900): tmp verbose test output
LarsJurgensen Aug 6, 2026
b71398b
feat(OCISDEV-900): rebase followup
LarsJurgensen Aug 6, 2026
4cdb547
feat(OCISDEV-900): logs
LarsJurgensen Aug 6, 2026
8ba8e7b
feat(OCISDEV-900): logs
LarsJurgensen Aug 6, 2026
b28c524
feat(OCISDEV-900): logs
LarsJurgensen Aug 6, 2026
f01a473
feat(OCISDEV-900): logs
LarsJurgensen Aug 7, 2026
9a07e8d
feat(OCISDEV-900): debug logs & avoid race condition
LarsJurgensen Aug 7, 2026
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
42 changes: 21 additions & 21 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -13,22 +13,22 @@ jobs:
check-go-generate:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- run: make go-generate && git diff --exit-code

unit-tests:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- run: sudo apt-get update && sudo apt-get install -y inotify-tools
- run: make test
- uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
- uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4
if: always()
with:
name: coverage
Expand All @@ -38,8 +38,8 @@ jobs:
needs: [unit-tests]
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- run: make dist
Expand All @@ -48,8 +48,8 @@ jobs:
needs: [unit-tests]
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- run: docker run -d --name redis --network host -e REDIS_DATABASES=1 redis:6-alpine
Expand All @@ -65,8 +65,8 @@ jobs:
matrix:
endpoint: [old-webdav, new-webdav, spaces-dav]
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- run: python3 tests/acceptance/run-litmus.py --endpoint ${{ matrix.endpoint }}
Expand All @@ -79,8 +79,8 @@ jobs:
matrix:
storage: [ocis] # TODO: s3ng blocked on Ceph networking — follow-up PR
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- run: python3 tests/acceptance/run-cs3api.py --storage ${{ matrix.storage }}
Expand All @@ -93,15 +93,15 @@ jobs:
matrix:
part: [1, 2] # TODO: parts 3,4 blocked on chunked Transfer-Encoding revad bug — follow-up PR
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- uses: shivammathur/setup-php@f3e473d116dcccaddc5834248c87452386958240 # v2
- uses: shivammathur/setup-php@accd6127cb78bee3e8082180cb391013d204ef9f # v2
with:
php-version: "8.4"
- run: python3 tests/acceptance/run-acceptance.py --storage ocis --total-parts 4 --run-part ${{ matrix.part }}
- uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
- uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4
if: failure()
with:
name: acceptance-ocis-part-${{ matrix.part }}
Expand All @@ -117,16 +117,16 @@ jobs:
matrix:
part: [1, 2] # TODO: parts 3,4 blocked on chunked Transfer-Encoding revad bug — follow-up PR
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
- uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4
- uses: actions/setup-go@40f1582b2485089dde7abd97c1529aa768e1baff # v5
with:
go-version-file: go.mod
- uses: shivammathur/setup-php@f3e473d116dcccaddc5834248c87452386958240 # v2
- uses: shivammathur/setup-php@accd6127cb78bee3e8082180cb391013d204ef9f # v2
with:
php-version: "8.4"
- run: sudo apt-get update && sudo apt-get install -y inotify-tools
- run: python3 tests/acceptance/run-acceptance.py --storage posixfs --total-parts 4 --run-part ${{ matrix.part }}
- uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
- uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4
if: failure()
with:
name: acceptance-posixfs-part-${{ matrix.part }}
Expand Down
54 changes: 52 additions & 2 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"
pkgupload "github.com/owncloud/reva/v2/pkg/upload"
"github.com/owncloud/reva/v2/pkg/utils"
"github.com/pkg/errors"
"github.com/rs/zerolog"
Expand All @@ -70,6 +72,7 @@ type config struct {
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"`
UploadExpiration int64 `mapstructure:"upload_expiration" docs:"0;Duration for how long uploads will be valid."`
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."`
Events eventconfig `mapstructure:"events" docs:"0;Event stream configuration"`
}

Expand All @@ -81,6 +84,9 @@ type eventconfig struct {
EnableTLS bool `mapstructure:"nats_enable_tls" docs:"events tls switch"`
AuthUsername string `mapstructure:"nats_username" docs:"event stream username"`
AuthPassword string `mapstructure:"nats_password" docs:"event stream password"`
ConsumerGroup string `mapstructure:"consumer_group" docs:"dcfs;Consumer group for the upload coordinator"`
NumConsumers int `mapstructure:"numconsumers" docs:"1;Number of concurrent postprocessing event consumers"`
AsyncUploads bool `mapstructure:"async_uploads" docs:"false;Require asynchronous upload processing; startup fails if the event stream is not configured"`
}

func (c *config) init() {
Expand All @@ -101,11 +107,20 @@ func (c *config) init() {
if len(c.AvailableXS) == 0 {
c.AvailableXS = map[string]uint32{"md5": 100, "unset": 1000}
}

if c.Events.ConsumerGroup == "" {
c.Events.ConsumerGroup = "dcfs"
}
if c.Events.NumConsumers <= 0 {
c.Events.NumConsumers = 1
}
}


type Service struct {
conf *config
Storage storage.FS
Coordinator pkgupload.Coordinator
dataServerURL *url.URL
availableXS []*provider.ResourceChecksumPriority
}
Expand Down Expand Up @@ -202,9 +217,36 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc.
return nil, err
}

store := pkgupload.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)
}
evstream, err := estreamFromConfig(c.Events)
if err != nil {
return nil, err
}
if c.Events.AsyncUploads && evstream == nil {
return nil, errors.New("storageprovider: async_uploads is enabled but no event stream is configured")
}
async := c.Events.AsyncUploads
coord, err := pkgupload.NewCoordinator(fs, store, evstream, async,
c.MountID, c.Events.ConsumerGroup, c.Events.NumConsumers, log,
filepath.Join(store.Root(), "uploads"))
if err != nil {
return nil, err
}
if async {
if err := coord.Start(evstream); err != nil {
return nil, err
}
}
service := &Service{
conf: c,
Storage: fs,
Coordinator: coord,
dataServerURL: u,
availableXS: xsTypes,
}
Expand Down Expand Up @@ -427,7 +469,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 @@ -1301,7 +1343,15 @@ func estreamFromConfig(c eventconfig) (events.Stream, error) {
return nil, nil
}

return stream.NatsFromConfig("storageprovider", false, stream.NatsConfig(c))
return stream.NatsFromConfig("storageprovider", false, stream.NatsConfig{
Endpoint: c.Endpoint,
Cluster: c.Cluster,
TLSInsecure: c.TLSInsecure,
TLSRootCACertificate: c.TLSRootCACertificate,
EnableTLS: c.EnableTLS,
AuthUsername: c.AuthUsername,
AuthPassword: c.AuthPassword,
})
}

func canLockPublicShare(ctx context.Context) bool {
Expand Down
42 changes: 39 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 @@ -44,13 +46,18 @@ type config struct {
Driver string `mapstructure:"driver" docs:"localhome;The storage driver to be used."`
Drivers map[string]map[string]interface{} `mapstructure:"drivers" docs:"url:pkg/storage/fs/localhome/localhome.go;The configuration for the storage driver"`
DataTXs map[string]map[string]interface{} `mapstructure:"data_txs" docs:"url:pkg/rhttp/datatx/manager/simple/simple.go;The configuration for the data tx protocols"`
MountID string `mapstructure:"mount_id"`
NatsAddress string `mapstructure:"nats_address"`
NatsClusterID string `mapstructure:"nats_clusterID"`
NatsTLSInsecure bool `mapstructure:"nats_tls_insecure"`
NatsRootCACertPath string `mapstructure:"nats_root_ca_cert_path"`
NatsEnableTLS bool `mapstructure:"nats_enable_tls"`
NatsUsername string `mapstructure:"nats_username"`
NatsPassword string `mapstructure:"nats_password"`
ConsumerGroup string `mapstructure:"consumer_group"`
NumConsumers int `mapstructure:"numconsumers"`
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."`
AsyncUploads bool `mapstructure:"async_uploads"`
}

func (c *config) init() {
Expand All @@ -60,6 +67,12 @@ func (c *config) init() {
if c.Driver == "" {
c.Driver = "localhome"
}
if c.ConsumerGroup == "" {
c.ConsumerGroup = "dcfs"
}
if c.NumConsumers <= 0 {
c.NumConsumers = 1
}
}

type svc struct {
Expand Down Expand Up @@ -104,7 +117,30 @@ func New(m map[string]interface{}, log *zerolog.Logger) (global.Service, error)
return nil, err
}

dataTXs, err := getDataTXs(conf, fs, evstream, log)
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)
}
if conf.AsyncUploads && evstream == nil {
return nil, fmt.Errorf("dataprovider: async_uploads is enabled but no event stream is configured")
}
async := conf.AsyncUploads
coord, err := pkgupload.NewCoordinator(fs, store, evstream, async,
conf.MountID, conf.ConsumerGroup, conf.NumConsumers, log,
filepath.Join(store.Root(), "uploads"))
if err != nil {
return nil, err
}
if async {
if err := coord.Start(evstream); err != nil {
return nil, err
}
}

dataTXs, err := getDataTXs(conf, coord, fs, evstream, log)
if err != nil {
return nil, err
}
Expand All @@ -126,7 +162,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 +182,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
7 changes: 7 additions & 0 deletions internal/http/services/owncloud/ocdav/copy.go
Original file line number Diff line number Diff line change
Expand Up @@ -666,6 +666,13 @@ func (s *svc) prepareCopy(ctx context.Context, w http.ResponseWriter, r *http.Re
return nil
}

if !srcStatRes.GetInfo().GetPermissionSet().GetInitiateFileDownload() {
w.WriteHeader(http.StatusForbidden)
b, err := errors.Marshal(http.StatusForbidden, "no permission to copy from this resource", "", "")
errors.HandleWebdavError(log, w, b, err)
return nil
}

dstStatReq := &provider.StatRequest{Ref: dstRef}
dstStatRes, err := client.Stat(ctx, dstStatReq)
switch {
Expand Down
4 changes: 4 additions & 0 deletions internal/http/services/owncloud/ocdav/propfind/propfind.go
Original file line number Diff line number Diff line change
Expand Up @@ -587,6 +587,10 @@ func (p *Handler) getResourceInfos(ctx context.Context, w http.ResponseWriter, r
var status *rpc.Status
info, status, err = p.statSpace(ctx, spaceRef, metadataKeys, fieldMaskPaths)
if err != nil || status.GetCode() != rpc.Code_CODE_OK {
if status.GetCode() == rpc.Code_CODE_TOO_EARLY {
w.WriteHeader(http.StatusTooEarly)
return nil, false, false
}
continue
}
}
Expand Down
2 changes: 1 addition & 1 deletion internal/http/services/owncloud/ocdav/tus.go
Original file line number Diff line number Diff line change
Expand Up @@ -319,7 +319,7 @@ 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 {
if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil && resid.OpaqueId != "" {
sReq.Ref = &provider.Reference{
ResourceId: &resid,
}
Expand Down
5 changes: 5 additions & 0 deletions pkg/logger/logger.go
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,11 @@ func fromConfig(conf *LogConf) (*zerolog.Logger, error) {
conf.Level = zerolog.DebugLevel.String()
}

// Set the global level to the minimum so only the per-logger level matters.
// oCIS does this in ocis-pkg/log when it embeds reva; the standalone revad
// binary otherwise inherits a global level that suppresses all output.
zerolog.SetGlobalLevel(zerolog.TraceLevel)

var opts []Option
opts = append(opts, WithLevel(conf.Level))

Expand Down
3 changes: 3 additions & 0 deletions pkg/ocm/storage/received/ocm.go
Original file line number Diff line number Diff line change
Expand Up @@ -539,6 +539,9 @@ func (d *driver) GetLock(ctx context.Context, ref *provider.Reference) (*provide
return nil, err
}

if token == "" {
return nil, errtypes.NotFound("no lock found")
}
return &provider.Lock{LockId: token, Type: provider.LockType_LOCK_TYPE_EXCL}, nil
}

Expand Down
3 changes: 2 additions & 1 deletion pkg/rhttp/datatx/datatx.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,12 +28,13 @@ 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.
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
Loading
Loading