Skip to content

Commit 804f3a6

Browse files
LarsJurgensenf-firas
authored andcommitted
feat: add NewCoordinatorFromConfig
1 parent b529e2a commit 804f3a6

3 files changed

Lines changed: 25 additions & 33 deletions

File tree

‎internal/grpc/services/storageprovider/storageprovider.go‎

Lines changed: 3 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -212,9 +212,10 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc.
212212
return nil, err
213213
}
214214

215-
coord, err := getCoordinator(c, fs, evstream, log)
215+
// storageprovider only initiates uploads; the data path assembles chunks, so no chunking here.
216+
coord, err := upload.NewCoordinatorFromConfig(c.UploadDirectory, c.Drivers[c.Driver], fs, evstream, log, false)
216217
if err != nil {
217-
return nil, err
218+
return nil, fmt.Errorf("storageprovider: %w", err)
218219
}
219220

220221
service := &Service{
@@ -228,22 +229,6 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc.
228229
return service, nil
229230
}
230231

231-
// getCoordinator builds the coordinator that initiates uploads for the driver
232-
// this service mounts. It stages sessions in the same directory the dataprovider
233-
// appends bytes to, so an upload initiated here can be continued there.
234-
func getCoordinator(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (upload.Coordinator, error) {
235-
store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log)
236-
if store == nil {
237-
return nil, fmt.Errorf("storageprovider: cannot determine the upload directory, set upload_directory")
238-
}
239-
if err := store.Setup(); err != nil {
240-
return nil, fmt.Errorf("storageprovider: upload directory setup failed: %w", err)
241-
}
242-
243-
// No chunk folder: only the data path assembles chunks.
244-
return upload.NewCoordinator(fs, store, "", publisher), nil
245-
}
246-
247232
func (s *Service) SetArbitraryMetadata(ctx context.Context, req *provider.SetArbitraryMetadataRequest) (*provider.SetArbitraryMetadataResponse, error) {
248233
ctx = ctxpkg.ContextSetLockID(ctx, req.LockId)
249234

‎internal/http/services/dataprovider/dataprovider.go‎

Lines changed: 3 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -106,9 +106,10 @@ func New(m map[string]interface{}, log *zerolog.Logger) (global.Service, error)
106106
return nil, err
107107
}
108108

109-
coord, err := getCoordinator(conf, fs, evstream, log)
109+
// the data path assembles chunks, so enable chunking
110+
coord, err := upload.NewCoordinatorFromConfig(conf.UploadDirectory, conf.Drivers[conf.Driver], fs, evstream, log, true)
110111
if err != nil {
111-
return nil, err
112+
return nil, fmt.Errorf("dataprovider: %w", err)
112113
}
113114

114115
// only the data path consumes postprocessing results: one consumer group gets
@@ -141,19 +142,6 @@ func getFS(c *config, stream events.Stream, log *zerolog.Logger) (storage.FS, er
141142
return nil, fmt.Errorf("driver not found: %s", c.Driver)
142143
}
143144

144-
// getCoordinator builds the coordinator that owns the upload lifecycle for the
145-
// driver this service mounts.
146-
func getCoordinator(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (upload.Coordinator, error) {
147-
store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log)
148-
if store == nil {
149-
return nil, fmt.Errorf("dataprovider: cannot determine the upload directory, set upload_directory")
150-
}
151-
if err := store.Setup(); err != nil {
152-
return nil, fmt.Errorf("dataprovider: upload directory setup failed: %w", err)
153-
}
154-
return upload.NewCoordinator(fs, store, store.UploadDir(), publisher), nil
155-
}
156-
157145
func getDataTXs(c *config, coord upload.Coordinator, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (map[string]http.Handler, error) {
158146
if c.DataTXs == nil {
159147
c.DataTXs = make(map[string]map[string]interface{})

‎pkg/upload/coordinator.go‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import (
1111
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
1212
provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
1313
"github.com/google/uuid"
14+
"github.com/rs/zerolog"
1415
tusd "github.com/tus/tusd/v2/pkg/handler"
1516

1617
"github.com/owncloud/reva/v2/pkg/appctx"
@@ -62,6 +63,24 @@ func NewCoordinator(fs storage.FS, store SessionStore, chunkFolder string, pub e
6263
return c
6364
}
6465

66+
// NewCoordinatorFromConfig sets up a coordinator and its file store. Pass
67+
// withChunking=false when only initiating uploads: chunk assembly happens on the
68+
// data path, so the storageprovider doesn't need it.
69+
func NewCoordinatorFromConfig(uploadDir string, driverConf map[string]interface{}, fs storage.FS, pub events.Publisher, log *zerolog.Logger, withChunking bool) (Coordinator, error) {
70+
store := NewFileStoreFromConfig(uploadDir, driverConf, log)
71+
if store == nil {
72+
return nil, fmt.Errorf("cannot determine the upload directory, set upload_directory")
73+
}
74+
if err := store.Setup(); err != nil {
75+
return nil, fmt.Errorf("upload directory setup failed: %w", err)
76+
}
77+
chunkFolder := ""
78+
if withChunking {
79+
chunkFolder = store.UploadDir()
80+
}
81+
return NewCoordinator(fs, store, chunkFolder, pub), nil
82+
}
83+
6584
// InitiateUpload resolves the target, then creates and persists the session that
6685
// bytes are appended to.
6786
func (c *coordinator) InitiateUpload(ctx context.Context, ref *provider.Reference, uploadLength int64, metadata map[string]string) (map[string]string, error) {

0 commit comments

Comments
 (0)