Skip to content

Commit 8078555

Browse files
committed
s3-compatibility API: transition to dpq (datapath query)
* update top target `/s3` handler and its downstream methods * remove the following MPT query names from the dpq exceptions list: - s3.QparamMptUploadID - s3.QparamMptUploads - s3.QparamMptPartNo * (native multi-part API uses the same exact constants) * unescape s3.QparamMptUploadID (and therefore, apc.QparamMptUploadID) * micro-optimize: remove redundant dpq.alloc/parse (that would co-exist with query parsing in certain cases) ---- * unescape apc.QparamETLTransformArgs (fix) * add dpq parsing unit test Signed-off-by: Alex Aizman <alex.aizman@gmail.com>
1 parent aa836dc commit 8078555

14 files changed

Lines changed: 125 additions & 92 deletions

File tree

ais/dpq.go

Lines changed: 12 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -95,12 +95,6 @@ type (
9595
var _except = map[string]bool{
9696
apc.QparamDontHeadRemote: false,
9797

98-
// flows that utilize the following query parameters perform conventional r.URL.Query()
99-
// TODO -- FIXME: multipart must definitely transition to dpq
100-
s3.QparamMptUploadID: false,
101-
s3.QparamMptUploads: false,
102-
s3.QparamMptPartNo: false,
103-
10498
s3.QparamAccessKeyID: false,
10599
s3.QparamExpires: false,
106100
s3.QparamSignature: false,
@@ -142,6 +136,11 @@ func dpqFree(d *dpq) {
142136

143137
func (dpq *dpq) get(qparam string) string { return dpq.m[qparam] }
144138

139+
func (dpq *dpq) has(qparam string) bool {
140+
_, ok := dpq.m[qparam]
141+
return ok
142+
}
143+
145144
func (dpq *dpq) parse(rawQuery string) error {
146145
var (
147146
query = rawQuery // r.URL.RawQuery
@@ -217,8 +216,10 @@ func (dpq *dpq) parse(rawQuery string) error {
217216
dpq.sv.sig = value // (base64.RawURLEncoding)
218217

219218
// Parameters that don't need unescaping
220-
case apc.QparamMptUploadID, apc.QparamMptPartNo, apc.QparamFltPresence, apc.QparamBinfoWithOrWithoutRemote,
221-
apc.QparamAppendType, apc.QparamETLName, apc.QparamETLTransformArgs,
219+
case apc.QparamMptUploads, apc.QparamMptPartNo,
220+
apc.QparamFltPresence, apc.QparamBinfoWithOrWithoutRemote,
221+
apc.QparamETLName,
222+
apc.QparamAppendType,
222223
apc.QparamNewCustom,
223224
apc.QparamKeepRemote,
224225
apc.QparamTID:
@@ -234,7 +235,9 @@ func (dpq *dpq) parse(rawQuery string) error {
234235
continue
235236
}
236237

237-
// Default: unescape and store in map
238+
// Default: unescape and store in map - e.g.:
239+
// - apc.QparamMptUploadID == s3.QparamMptUploadID (case in point: s3 uploadId)
240+
// - apc.QparamETLTransformArgs (arbitrary JSON encoded via url.Values)
238241
dpq.m[key], err = _unescape(value)
239242
}
240243

ais/dpq_internal_test.go

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,32 @@
1+
// Package ais: internal unit tests
2+
/*
3+
* Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved.
4+
*/
5+
package ais
6+
7+
import (
8+
"net/url"
9+
"testing"
10+
11+
"github.com/NVIDIA/aistore/api/apc"
12+
"github.com/NVIDIA/aistore/tools/tassert"
13+
)
14+
15+
func TestDpqETLTransformArgs(t *testing.T) {
16+
for _, args := range []string{
17+
`{"value":"a+b%20&c"}`,
18+
`hello world`,
19+
`a/b?c=d&e=f`,
20+
} {
21+
q := url.Values{}
22+
q.Set(apc.QparamETLTransformArgs, args)
23+
24+
dpq := dpqAlloc()
25+
err := dpq.parse(q.Encode())
26+
tassert.CheckFatal(t, err)
27+
28+
actual := dpq.get(apc.QparamETLTransformArgs)
29+
tassert.Errorf(t, actual == args, "expected %q, got %q", args, actual)
30+
dpqFree(dpq)
31+
}
32+
}

ais/emd_internal_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
// Package ais: internal unit tests
22
/*
3-
* Copyright (c) 2021-2025, NVIDIA CORPORATION. All rights reserved.
3+
* Copyright (c) 2021-2026, NVIDIA CORPORATION. All rights reserved.
44
*/
55
package ais
66

ais/multiproxy_internal_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
// Package ais: internal unit tests
22
/*
3-
* Copyright (c) 2018-2025, NVIDIA CORPORATION. All rights reserved.
3+
* Copyright (c) 2018-2026, NVIDIA CORPORATION. All rights reserved.
44
*/
55
package ais
66

ais/s3/const.go

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,11 @@
11
// Package s3 provides Amazon S3 compatibility layer
22
/*
3-
* Copyright (c) 2018-2025, NVIDIA CORPORATION. All rights reserved.
3+
* Copyright (c) 2018-2026, NVIDIA CORPORATION. All rights reserved.
44
*/
55
package s3
66

77
import (
8+
"github.com/NVIDIA/aistore/api/apc"
89
"github.com/NVIDIA/aistore/cmn"
910
"github.com/NVIDIA/aistore/cmn/cos"
1011
)
@@ -23,10 +24,12 @@ const (
2324
QparamStartAfter = "start-after" // Start listing after this object key
2425
QparamDelimiter = "delimiter" // Delimiter for grouping object keys
2526

26-
// multipart
27-
QparamMptUploads = "uploads" // Start multipart upload or list active uploads
28-
QparamMptUploadID = "uploadId" // Complete, abort, or list parts of specific multipart upload
29-
QparamMptPartNo = "partNumber" // Part number for multipart upload
27+
// AIS native multipart APIs use
28+
// canonical S3 constants
29+
QparamMptUploads = apc.QparamMptUploads // Start multipart upload or list active uploads
30+
QparamMptUploadID = apc.QparamMptUploadID // Complete, abort, or list parts of specific multipart upload
31+
QparamMptPartNo = apc.QparamMptPartNo // Part number for multipart upload
32+
3033
QparamMptMaxUploads = "max-uploads"
3134
QparamMptUploadIDMarker = "upload-id-marker"
3235

ais/test/multiproxy_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
// Package integration_test.
22
/*
3-
* Copyright (c) 2018-2025, NVIDIA CORPORATION. All rights reserved.
3+
* Copyright (c) 2018-2026, NVIDIA CORPORATION. All rights reserved.
44
*/
55
package integration_test
66

ais/tgts3.go

Lines changed: 48 additions & 54 deletions
Original file line numberDiff line numberDiff line change
@@ -40,28 +40,41 @@ func (t *target) s3Handler(w http.ResponseWriter, r *http.Request) {
4040
switch r.Method {
4141
case http.MethodHead:
4242
t.headObjS3(w, r, apiItems)
43+
return
44+
case http.MethodGet, http.MethodPut, http.MethodDelete, http.MethodPost:
45+
// dpq parsed below
46+
default:
47+
cmn.WriteErr405(w, r, http.MethodDelete, http.MethodGet, http.MethodHead, http.MethodPut, http.MethodPost)
48+
return
49+
}
50+
51+
dpq := dpqAlloc()
52+
defer dpqFree(dpq)
53+
if err := dpq.parse(r.URL.RawQuery); err != nil {
54+
s3.WriteErr(w, r, s3.ErrInfo{Err: err})
55+
return
56+
}
57+
dpq.isS3 = true
58+
59+
switch r.Method {
4360
case http.MethodGet:
44-
t.getObjS3(w, r, apiItems)
61+
t.getObjS3(w, r, apiItems, dpq)
4562
case http.MethodPut:
46-
config := cmn.GCO.Get()
47-
t.putCopyMpt(w, r, config, apiItems)
63+
t.putCopyMpt(w, r, cmn.GCO.Get(), dpq, apiItems)
4864
case http.MethodDelete:
49-
q := r.URL.Query()
50-
if q.Has(s3.QparamMptUploadID) {
51-
t.abortMptS3(w, r, apiItems, q)
65+
if dpq.has(s3.QparamMptUploadID) {
66+
t.abortMptS3(w, r, dpq, apiItems)
5267
} else {
5368
t.delObjS3(w, r, apiItems)
5469
}
5570
case http.MethodPost:
56-
t.postObjS3(w, r, apiItems)
57-
default:
58-
cmn.WriteErr405(w, r, http.MethodDelete, http.MethodGet, http.MethodHead, http.MethodPut, http.MethodPost)
71+
t.postObjS3(w, r, apiItems, dpq)
5972
}
6073
}
6174

6275
// PUT /s3/<bucket-name>/<object-name>
6376
// [switch] mpt | put | copy
64-
func (t *target) putCopyMpt(w http.ResponseWriter, r *http.Request, config *cmn.Config, items []string) {
77+
func (t *target) putCopyMpt(w http.ResponseWriter, r *http.Request, config *cmn.Config, dpq *dpq, items []string) {
6578
cs := fs.Cap()
6679
if cs.IsOOS() {
6780
s3.WriteErr(w, r, s3.ErrInfo{Err: cs.Err(), Status: http.StatusInsufficientStorage})
@@ -77,9 +90,8 @@ func (t *target) putCopyMpt(w http.ResponseWriter, r *http.Request, config *cmn.
7790
s3.HandleAwsChunked(r)
7891
}
7992

80-
q := r.URL.Query()
8193
switch {
82-
case q.Has(s3.QparamMptPartNo) && q.Has(s3.QparamMptUploadID):
94+
case dpq.has(s3.QparamMptPartNo) && dpq.has(s3.QparamMptUploadID):
8395
if r.Header.Get(cos.S3HdrObjSrc) != "" {
8496
// TODO:
8597
// copy another object (or its range) => part of the specified multipart upload.
@@ -89,16 +101,16 @@ func (t *target) putCopyMpt(w http.ResponseWriter, r *http.Request, config *cmn.
89101
return
90102
}
91103
if cmn.Rom.V(5, cos.ModS3) {
92-
nlog.Infoln("putPartMpt", bck.String(), items, q)
104+
nlog.Infoln("putPartMpt", bck.String(), items, dpq.m)
93105
}
94-
t.putPartMptS3(w, r, items, q, bck)
106+
t.putPartMptS3(w, r, dpq, bck, items)
95107
case r.Header.Get(cos.S3HdrObjSrc) == "":
96108
objName, errN := s3.JoinValidateOname(w, r, items)
97109
if errN != nil {
98110
return
99111
}
100112
lom := core.AllocLOM(objName)
101-
t.putObjS3(w, r, bck, config, lom)
113+
t.putObjS3(w, r, bck, config, lom, dpq)
102114
core.FreeLOM(lom)
103115
default:
104116
t.copyObjS3(w, r, config, items)
@@ -201,7 +213,7 @@ func (t *target) copyObjS3(w http.ResponseWriter, r *http.Request, config *cmn.C
201213
sgl.Free()
202214
}
203215

204-
func (t *target) putObjS3(w http.ResponseWriter, r *http.Request, bck *meta.Bck, config *cmn.Config, lom *core.LOM) {
216+
func (t *target) putObjS3(w http.ResponseWriter, r *http.Request, bck *meta.Bck, config *cmn.Config, lom *core.LOM, dpq *dpq) {
205217
if err := lom.InitBck(bck); err != nil {
206218
if cmn.IsErrRemoteBckNotFound(err) {
207219
t.BMDVersionFixup(r)
@@ -215,14 +227,6 @@ func (t *target) putObjS3(w http.ResponseWriter, r *http.Request, bck *meta.Bck,
215227
started := time.Now()
216228
lom.SetAtimeUnix(started.UnixNano())
217229

218-
// TODO: dual checksumming, e.g. lom.SetCustom(apc.AWS, ...)
219-
220-
dpq := dpqAlloc()
221-
if err := dpq.parse(r.URL.RawQuery); err != nil {
222-
s3.WriteErr(w, r, s3.ErrInfo{Err: err})
223-
dpqFree(dpq)
224-
return
225-
}
226230
poi := allocPOI()
227231
{
228232
poi.atime = started.UnixNano()
@@ -240,26 +244,22 @@ func (t *target) putObjS3(w http.ResponseWriter, r *http.Request, bck *meta.Bck,
240244
} else {
241245
s3.SetS3Headers(w.Header(), lom)
242246
}
243-
dpqFree(dpq)
244247
}
245248

246249
// GET s3/<bucket-name[/<object-name>]
247-
func (t *target) getObjS3(w http.ResponseWriter, r *http.Request, items []string) {
250+
func (t *target) getObjS3(w http.ResponseWriter, r *http.Request, items []string, dpq *dpq) {
248251
bucket := items[0]
249252
bck, ecode, err := meta.InitByNameOnly(bucket, t.owner.bmd)
250253
if err != nil {
251254
s3.WriteErr(w, r, s3.ErrInfo{Err: err, Status: ecode})
252255
return
253256
}
254257

255-
// TODO -- FIXME: transition to dpq
256-
q := r.URL.Query()
257-
258-
if len(items) == 1 && q.Has(s3.QparamMptUploads) {
258+
if len(items) == 1 && dpq.has(s3.QparamMptUploads) {
259259
if cmn.Rom.V(5, cos.ModS3) {
260-
nlog.Infoln("listUploadsMpt", bck.String(), q)
260+
nlog.Infoln("listUploadsMpt", bck.String(), dpq.m)
261261
}
262-
t.listUploadsMptS3(w, bck, q)
262+
t.listUploadsMptS3(w, bck, dpq)
263263
return
264264
}
265265
if len(items) < 2 {
@@ -271,32 +271,25 @@ func (t *target) getObjS3(w http.ResponseWriter, r *http.Request, items []string
271271
if errN != nil {
272272
return
273273
}
274-
if q.Has(s3.QparamMptPartNo) {
274+
if dpq.has(s3.QparamMptPartNo) {
275275
if cmn.Rom.V(5, cos.ModS3) {
276-
nlog.Infoln("getMptPart", bck.String(), objName, q)
276+
nlog.Infoln("getMptPart", bck.String(), objName, dpq.m)
277277
}
278278
lom := core.AllocLOM(objName)
279-
t.getPartMptS3(w, r, bck, lom, q)
279+
t.getPartMptS3(w, r, bck, lom, dpq)
280280
core.FreeLOM(lom)
281281
return
282282
}
283-
uploadID := q.Get(s3.QparamMptUploadID)
283+
uploadID := dpq.get(s3.QparamMptUploadID)
284284
if uploadID != "" {
285285
if cmn.Rom.V(5, cos.ModS3) {
286-
nlog.Infoln("listPartsMpt", bck.String(), objName, q)
286+
nlog.Infoln("listPartsMpt", bck.String(), objName, dpq.m)
287287
}
288-
t.listPartsMptS3(w, r, bck, objName, q)
288+
t.listPartsMptS3(w, r, bck, objName, dpq)
289289
return
290290
}
291291

292-
dpq := dpqAlloc()
293-
if err := dpq.parse(r.URL.RawQuery); err != nil {
294-
dpqFree(dpq)
295-
s3.WriteErr(w, r, s3.ErrInfo{Err: err})
296-
return
297-
}
298292
lom := core.AllocLOM(objName)
299-
dpq.isS3 = true
300293
lom, err = t.getObject(w, r, dpq, bck, lom)
301294
core.FreeLOM(lom)
302295

@@ -307,7 +300,6 @@ func (t *target) getObjS3(w http.ResponseWriter, r *http.Request, items []string
307300
}
308301
s3.WriteErr(w, r, ei)
309302
}
310-
dpqFree(dpq)
311303
}
312304

313305
// HEAD /s3/<bucket-name>/<object-name> (TODO: s3.HdrMptCnt)
@@ -417,27 +409,29 @@ func (t *target) delObjS3(w http.ResponseWriter, r *http.Request, items []string
417409
}
418410

419411
// POST /s3/<bucket-name>/<object-name>
420-
func (t *target) postObjS3(w http.ResponseWriter, r *http.Request, items []string) {
412+
func (t *target) postObjS3(w http.ResponseWriter, r *http.Request, items []string, dpq *dpq) {
421413
bck, ecode, err := meta.InitByNameOnly(items[0], t.owner.bmd)
422414
if err != nil {
423415
s3.WriteErr(w, r, s3.ErrInfo{Err: err, Status: ecode})
424416
return
425417
}
426-
q := r.URL.Query()
427-
if q.Has(s3.QparamMptUploads) {
418+
419+
if dpq.has(s3.QparamMptUploads) {
428420
if cmn.Rom.V(5, cos.ModS3) {
429-
nlog.Infoln("startMpt", bck.String(), items, q)
421+
nlog.Infoln("startMpt", bck.String(), items, dpq.m)
430422
}
431-
t.startMptS3(w, r, items, bck)
423+
t.startMptS3(w, r, bck, items)
432424
return
433425
}
434-
if q.Has(s3.QparamMptUploadID) {
426+
427+
if dpq.has(s3.QparamMptUploadID) {
435428
if cmn.Rom.V(5, cos.ModS3) {
436-
nlog.Infoln("completeMpt", bck.String(), items, q)
429+
nlog.Infoln("completeMpt", bck.String(), items, dpq.m)
437430
}
438-
t.completeMptS3(w, r, items, q, bck)
431+
t.completeMptS3(w, r, dpq, bck, items)
439432
return
440433
}
434+
441435
err = fmt.Errorf("set query parameter %q to start multipart upload or %q to complete the upload",
442436
s3.QparamMptUploads, s3.QparamMptUploadID)
443437
s3.WriteErr(w, r, s3.ErrInfo{Err: err})

0 commit comments

Comments
 (0)