From ffb1ec9af9c8046cc270cf934c688908c25e674c Mon Sep 17 00:00:00 2001 From: Tony Stack Date: Tue, 22 Sep 2026 08:09:25 +0800 Subject: [PATCH 1/2] fix: consume the parts a failed merge assembled from MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A chunked upload that reached the point of assembling its payload and was then rejected (digest mismatch, declared byte count mismatch) kept every part resource it had consumed. The service drops them together with the refusal; only requests turned down before any part was read — a create whose target already exists, a manifest that names nothing — keep the parts, because nothing was assembled yet. The divergence was visible to clients: the Java SDK reuses deterministic part names (..part.tmp.), so a refused duplicate create leaves a temp resource behind, and any later listing that sweeps for ".part.tmp." blames it on an unrelated upload. That is how the PyODPS contract probe started failing its merge-cleanup case whenever the Java suite had run first against the same instance: the leftover was the Java suite's, and the assertion was project-wide. - merge: delete the declared parts on the post-read verification failures, keep the existing consumption on success, and leave parts alone for the pre-read refusals (now pinned by tests on both sides of that line). - pyodps probe: scope the leftover check to the resource the case uploaded, assert a refused merge consumes its own part instead of the opposite, and let cleanup() sweep parts left by an aborted run of the probe itself. Verified: go test -race -count=1 ./internal/... (server 18.5s, engine 3.1s, wire 1.2s) all ok; PyODPS probe 30/30 against the built binary with a Java-style leftover seeded, and the pre-change probe still reports the old assertions verbatim, so the fix and the expectation moved together. --- CHANGELOG.md | 1 + internal/server/resources.go | 23 ++++++++++++++++++++--- internal/server/resources_test.go | 23 ++++++++++++++++++++++- tests/python/run.py | 28 ++++++++++++++++++++++------ 4 files changed, 65 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e9dc47f..5eab7a7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ - Extend test-mode fault injection to the metadata plane: `"plane":"rest"` rules match resources or functions by request verb and can inject `http_error` or `delay` without touching Tunnel sessions. - Log resources and functions requests as a separate `rest` event with action, object and project, keeping object names out of the log stream. - Require write grants to create or update resources and functions under strict authentication; the Tunnel read-session `create` exception no longer applies to metadata verbs. +- Consume the parts a chunked merge assembled from when that merge fails verification, while refusals decided before any part is read (a duplicate create, a malformed manifest) leave the uploaded parts in place. ## 1.1.0 — 2026-09-17 diff --git a/internal/server/resources.go b/internal/server/resources.go index 14e2bf2..728a247 100644 --- a/internal/server/resources.go +++ b/internal/server/resources.go @@ -238,6 +238,16 @@ func (s *Server) putResourcePart(w http.ResponseWriter, r *http.Request, p, sc, // mergeResourceParts assembles the declared parts into the final resource. The // body is "|[,...]"; the digest and the declared total // bytes are both verified before the payload becomes visible. +// +// A merge consumes the parts it declared, whether or not the target ends up +// published: a payload that fails verification is rejected together with the +// chunks it was built from, so a client that retries (the SDKs reuse +// deterministic part names) never finds a stale chunk in the way, and a +// session that gives up does not leave litter behind. Requests that were +// refused before any part was read keep the parts untouched: a malformed +// manifest names nothing to consume, a missing part is a refusal about that +// part, and creating over an existing resource is decided before the merge +// runs. func (s *Server) mergeResourceParts(w http.ResponseWriter, r *http.Request, p, sc, name string, mode engine.ResourceMode) { target, ok := s.resourceTarget(w, r, name) if !ok { @@ -272,26 +282,33 @@ func (s *Server) mergeResourceParts(w http.ResponseWriter, r *http.Request, p, s } merged = append(merged, content...) } + consumeParts := func() { + for _, part := range parts { + s.Engine.DeleteResource(r.Context(), p, sc, part) + } + } if size := r.Header.Get("x-odps-resource-merge-total-bytes"); size != "" { declared, e := strconv.ParseInt(size, 10, 64) if e != nil || declared != int64(len(merged)) { fail(w, r, 400, "InvalidParameter", fmt.Errorf("x-odps-resource-merge-total-bytes does not match the merged payload")) + consumeParts() return } } sum := md5.Sum(merged) if !strings.EqualFold(hex.EncodeToString(sum[:]), strings.TrimSpace(digest)) { fail(w, r, 400, "InvalidParameter", fmt.Errorf("merged payload does not match MD5 %s", strings.TrimSpace(digest))) + consumeParts() return } s.storeResource(w, r, p, sc, target, r.Header.Get("x-odps-resource-type"), r.Header.Get("x-odps-comment"), strings.EqualFold(r.Header.Get("x-odps-resource-istemp"), "true"), "", merged, mode) if w.Header().Get("Location") == "" { + // The target was refused (it already exists, or its kind is invalid), so + // nothing was assembled and the parts stay for the client to retry with. return } - for _, part := range parts { - s.Engine.DeleteResource(r.Context(), p, sc, part) - } + consumeParts() } func (s *Server) resourceMeta(w http.ResponseWriter, r *http.Request, p, sc, name string) { diff --git a/internal/server/resources_test.go b/internal/server/resources_test.go index 80dc3f1..e1b8857 100644 --- a/internal/server/resources_test.go +++ b/internal/server/resources_test.go @@ -180,7 +180,9 @@ func TestResourceRESTContract(t *testing.T) { if w.Code != 404 { t.Fatalf("parts must be cleaned after merge: %d %s", w.Code, w.Body.String()) } - // A mismatching digest or byte count must not publish a partial payload. + // A mismatching digest or byte count must not publish a partial payload, and + // the merge consumes the chunk it assembled from: an unverifiable part is + // dead weight, and both SDKs re-upload the same deterministic name on retry. for _, body := range []string{md5hex("nope") + "|" + partB, "0123456789abcdef0123456789abcdef|" + partB} { resourceRequest(t, s, "POST", "/projects/p/resources?rIsPart=true", partB, "file", temp, "x") w = resourceRequest(t, s, "POST", "/projects/p/resources?rOpMerge=true", "bad.py", "py", nil, body) @@ -190,6 +192,9 @@ func TestResourceRESTContract(t *testing.T) { if w = resourceRequest(t, s, "GET", "/projects/p/resources/bad.py", "", "", nil, ""); w.Code != 404 { t.Fatalf("failed merge must not publish: %d %s", w.Code, w.Body.String()) } + if w = resourceRequest(t, s, "GET", "/projects/p/resources/"+partB, "", "", nil, ""); w.Code != 404 { + t.Fatalf("refused merge must consume its part: %d %s", w.Code, w.Body.String()) + } } resourceRequest(t, s, "POST", "/projects/p/resources?rIsPart=true", partB, "file", temp, "x") w = resourceRequest(t, s, "POST", "/projects/p/resources?rOpMerge=true", "bad.py", "py", map[string]string{"x-odps-resource-merge-total-bytes": "99"}, md5hex("x")+"|"+partB) @@ -209,6 +214,22 @@ func TestResourceRESTContract(t *testing.T) { if w = resourceRequest(t, s, "GET", "/projects/p/resources/chunked.py", "", "", nil, ""); w.Body.String() != "print(2)" { t.Fatalf("chunked overwrite content: %q", w.Body.String()) } + // A create-merge whose target already exists is refused before the parts are + // read, so unlike a failed merge it leaves the uploaded chunk in place for + // the client to retry with (as an update, or after deleting the target). + if w = resourceRequest(t, s, "POST", "/projects/p/resources?rIsPart=true", partA, "file", temp, "print(3)"); w.Code != 200 { + t.Fatalf("part for the colliding merge: %d %s", w.Code, w.Body.String()) + } + w = resourceRequest(t, s, "POST", "/projects/p/resources?rOpMerge=true", "chunked.py", "py", nil, md5hex("print(3)")+"|"+partA) + if w.Code != 400 || !strings.Contains(w.Body.String(), "ResourceAlreadyExists") { + t.Fatalf("colliding merge: %d %s", w.Code, w.Body.String()) + } + if w = resourceRequest(t, s, "GET", "/projects/p/resources/"+partA, "", "", nil, ""); w.Code != 200 { + t.Fatalf("a refused create must keep the part it was given: %d %s", w.Code, w.Body.String()) + } + if w = resourceRequest(t, s, "GET", "/projects/p/resources/chunked.py", "", "", nil, ""); w.Body.String() != "print(2)" { + t.Fatalf("a colliding merge must not touch the target: %q", w.Body.String()) + } // A TABLE resource keeps metadata only and validates the referenced table. if _, err = e.Execute(context.Background(), "p", "default", "create table src(id bigint);"); err != nil { t.Fatal(err) diff --git a/tests/python/run.py b/tests/python/run.py index 474fbae..2e31e7a 100644 --- a/tests/python/run.py +++ b/tests/python/run.py @@ -132,6 +132,15 @@ def quiet(fn, *args, **kw): def cleanup(): + # A part left behind by an aborted run of *this* probe would otherwise be + # attributed to the next one. Sweep them by our own prefix, so the cleanup + # can never touch another client's in-flight upload. + try: + stale = [r.name for r in odps.list_resources(prefix=PREFIX) if ".part.tmp." in r.name] + except Exception: + stale = [] + for stale_part in stale: + quiet(odps.delete_resource, stale_part) for suffix in ("single.py", "bytes.bin", "sio.txt", "chunk.bin", "upd.bin", "schema.bin", "big.bin", "local.txt", "streamres.bin", "onlyschema.bin", "malformed.bin", "malformed2.bin", "badmd5.bin", "paged_%02d.bin"): @@ -203,8 +212,13 @@ def stream_write(): eq(fresh.size, len(BLOB), "merged size") eq(fresh.open("rb").read(), BLOB, "merged payload") eq(fresh.content_md5, hashlib.md5(BLOB).hexdigest(), "merged md5") - leftovers = [r.name for r in odps.list_resources() if ".part.tmp." in r.name] - eq(leftovers, [], "temp parts removed after merge") + # Scope the sweep to this case's own resource: a temp part left behind by a + # different client (a refused duplicate create keeps the chunks it uploaded, + # which is the service behaviour the probe pins further down) is not a + # regression in *our* merge, and asserting project-wide made this case depend + # on what else had run against the instance before it. + leftovers = [r.name for r in odps.list_resources(prefix=PREFIX + "_big.bin") if ".part.tmp." in r.name] + eq(leftovers, [], "this resource's temp parts removed after merge") return "size=%d md5=verified parts_left=%d" % (fresh.size, len(leftovers)) @@ -432,11 +446,13 @@ def merge_with_wrong_md5(): eq(code, 400, "status") eq("MD5" in body, True, "reason mentions MD5") eq(odps.get_resource(name).open("rb").read(), b"seed", "previous payload survives") - # a refused merge must leave the parts alone (they are not the client's to keep) - eq(any(".part.tmp." in r.name for r in odps.list_resources(prefix=name)), True, "part kept after refusal") - quiet(odps.delete_resource, name + ".part.tmp.000001.000000") + # The merge got as far as assembling the payload and then rejected it, so the + # chunk it consumed is gone with it; only refusals decided before any part was + # read (a duplicate create, a malformed manifest) keep the parts. + eq([r.name for r in odps.list_resources(prefix=name) if ".part.tmp." in r.name], [], + "part consumed by the refused merge") odps.delete_resource(name) - return "400 on MD5 mismatch, old payload intact" + return "400 on MD5 mismatch, old payload intact, part consumed" case("a merge whose MD5 does not match is refused", merge_with_wrong_md5) From 1eb8111ab702b72841c873765c1faa32844dfed6 Mon Sep 17 00:00:00 2001 From: Tony Stack Date: Sun, 4 Oct 2026 19:42:25 +0800 Subject: [PATCH 2/2] fix: make the part lifecycle after a refused merge follow the measured service Review comment on this pull request: the consumption behaviour was asserted, not observed, and the parts cannot be recovered once deleted. It is observed now, on a live MaxCompute project (three-tier, temp resources created and deleted by the probe itself, 2026-10-04 19:23-19:26), and the emulator moved to match it in two places. Measured, per failure shape, with the part resources read back afterwards: | merge request | service answer | its parts afterwards | | --- | --- | --- | | correct digest and size | accepted, target published | consumed | | digest does not match the assembled payload | `ODPS-0421213 Save resource error - Merge part temp files failed!` | **consumed** | | manifest names a part that was never uploaded | `ODPS-0421111 Resource not found` | kept (the one real part survives) | | target already exists | `ODPS-0421121 The resource has already existed` | kept | | declared `x-odps-resource-merge-total-bytes` wrong, digest right (4400 declared for 304 bytes) | **accepted**, target is the correct 304 bytes | consumed | Two conclusions, and they cut in different directions from what this branch claimed before: - consuming the parts on a digest-refused merge is right - the service really does it, so an SDK that retries with the same deterministic part names has to re-upload them, and the emulator should not be the place that quietly keeps stale chunks; - the declared-byte-count check is not a service check at all. The service compares the digest, and only uses the declared total against the project's maximum. A refusal on that local check therefore has no service precedent for what it does to the parts, so the emulator now refuses and **keeps** them: it should not destroy an upload on the strength of a rule the service does not have. That divergence is written into docs/protocol.md rather than left implicit; the HTTP status difference on the existing-target path (service `ODPS-0421121`, emulator 400) is pre-existing and deliberately untouched here, because that status was not read back. Pinned in tests, both directions: - `internal/server/resources_test.go`: the digest-refused merge consumes its part (as before), the declared-byte-count refusal now asserts the part is still readable (200), and the comment says which of the two is service behaviour and which is the emulator's own strictness. - `tests/python/run.py`: three new cases - digest refused/consumed, refusal decided before any part is read/kept, declared mismatch/kept. The existing-target case asserts the refusal and the retention, and explicitly does not assert a status code, because that one was not measured. Verified: `go test -race -count=1 ./internal/...` (engine 3.03s, server 35.81s, wire 1.28s) all ok, and the PyODPS probe against a freshly built binary from this tree: 33 passed, 0 failed (previously 30 cases). The probe's leftover sweep is scoped to the case's own resource names, so one client's aborted upload cannot be blamed on another's merge. No merge, no tag, no release. --- CHANGELOG.md | 1 + docs/protocol.md | 2 +- internal/server/resources.go | 30 +++++++---- internal/server/resources_test.go | 15 ++++-- tests/python/run.py | 82 +++++++++++++++++++++++++++++++ 5 files changed, 116 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7650ddb..46e231f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ - Add function registration (`/projects/p/registration/functions`) for Java, SQL and embedded functions with resource-reference validation; execution still returns `UnsupportedFeature`. - Declare `resources`, `functions` and `unsupported` surfaces in `/capabilities`. - Reject volume-backed resources and chunked-upload integrity mismatches with structured errors instead of publishing partial payloads. +- Follow the live service on what happens to a chunked upload's part resources after a merge is refused (measured 2026-10-04 on a real project, read-backs in the repo's probe): a merge rejected on its MD5 consumes the parts it read, a merge refused before any part was read keeps them, and the emulator's own stricter declared-byte-count check - which the service does not perform at all - keeps them too. - Extend test-mode fault injection to the metadata plane: `"plane":"rest"` rules match resources or functions by request verb and can inject `http_error` or `delay` without touching Tunnel sessions. - Log resources and functions requests as a separate `rest` event with action, object and project, keeping object names out of the log stream. - Require write grants to create or update resources and functions under strict authentication; the Tunnel read-session `create` exception no longer applies to metadata verbs. diff --git a/docs/protocol.md b/docs/protocol.md index c5f47b2..0fc5791 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -20,7 +20,7 @@ | POST /projects/p/authorization?sign_bearer_token | `` 占位 token,按请求生成;不参与鉴权 | | GET /connection/mcqa | 固定 404 UnsupportedOperation,消息写明只模拟 SQLRT(MCQA v2/MaxQA 未实现) | | GET/POST /projects/p/resources;GET/PUT/POST/DELETE /projects/p/resources/{name} | 资源列表(name 前缀、type 过滤、marker/maxitems)与创建;单资源读取按 `?meta` 返回头部元数据、否则返回 payload(支持 `rOffset`/`rSize` 与 `x-odps-resource-has-remaining`);PUT 覆盖,DELETE 删除 | -| POST /projects/p/resources?rIsPart=true;POST/PUT ...?rOpMerge=true | Java SDK 分片上传:分片按确定性临时名 upsert,合并请求体 `|[,..]` 校验 MD5 与 `x-odps-resource-merge-total-bytes` 后发布正式资源并删除分片 | +| POST /projects/p/resources?rIsPart=true;POST/PUT ...?rOpMerge=true | Java SDK 分片上传:分片按确定性临时名 upsert;合并请求体 `|[,..]` 以 MD5 为准,成功后发布正式资源并删除分片。分片生命周期按 2026-10-04 真实项目实测对齐:MD5 不符的合并连分片一起拒掉(分片被消费,重试必须重传);读取分片之前就定案的拒绝(目标已存在、manifest 里点名了不存在的分片)保留分片。`x-odps-resource-merge-total-bytes` 与实配合并长度的比对是本模拟器比服务端更严的本地检查(服务端只做上限校验,实测申报 4400 字节、实际 304 字节的合并被接受且内容正确),因此该路径不消费分片。目标已存在时服务端回 `ODPS-0421121 The resource has already existed`(HTTP 状态未回读),模拟器回 400 InvalidParameter——这个状态差异是既有行为,不在本次改动范围内 | | GET/POST /projects/p/registration/functions;GET/PUT/POST/DELETE .../functions/{name} | 函数列表与注册;别名在 XML 的 `Alias` 元素(兼容 `Name`),引用资源必须已存在 | HTTP 分区参数遵循 SDK 的 `ds=2026-09-14` 写法,也接受单引号值;SQL 中使用 `PARTITION(ds='2026-09-14')`。当前逗号是分区键分隔符,不支持分区值本身包含逗号。quotaName 解析到 default 或显式配置的命名 quota;不存在的命名 quota 返回 404 QuotaNotExist,不执行生产配额调度,asyncmode 接受后同步建立快照。Protobuf 的 raw_size 参数不裁剪行数;当前 CPP 只在 Arrow 路径发送它。 diff --git a/internal/server/resources.go b/internal/server/resources.go index 728a247..222a920 100644 --- a/internal/server/resources.go +++ b/internal/server/resources.go @@ -239,15 +239,24 @@ func (s *Server) putResourcePart(w http.ResponseWriter, r *http.Request, p, sc, // body is "|[,...]"; the digest and the declared total // bytes are both verified before the payload becomes visible. // -// A merge consumes the parts it declared, whether or not the target ends up -// published: a payload that fails verification is rejected together with the -// chunks it was built from, so a client that retries (the SDKs reuse -// deterministic part names) never finds a stale chunk in the way, and a -// session that gives up does not leave litter behind. Requests that were -// refused before any part was read keep the parts untouched: a malformed -// manifest names nothing to consume, a missing part is a refusal about that -// part, and creating over an existing resource is decided before the merge -// runs. +// What happens to the part resources is pinned against a live MaxCompute project +// (three_pangu2_odps2, 2026-10-04 19:23-19:26 CST, temp resources created and deleted by +// the probe; raw read-backs in the workspace work item evidence): +// +// - a merge that publishes consumes its parts; +// - a merge refused because the assembled payload does not match the declared MD5 also +// consumes them - the service answers ODPS-0421213 and both parts are gone when read +// back, so retrying with the same deterministic part names means re-uploading them; +// - a merge refused before any part was read keeps them: naming an absent part is +// ODPS-0421111 with the uploaded part still present, and creating over an existing +// resource is ODPS-0421121 with both parts still present. +// +// The declared byte count is the one deliberate difference below. The service does not +// compare it with the assembled payload at all (measured: declaring 4400 bytes for a +// 304-byte merge is accepted and merges correctly; it only checks the value against the +// project's maximum), so this emulator rejecting the mismatch is stricter than the +// service. Because no service behaviour covers a failure the service cannot produce, the +// parts are left alone on that path rather than consumed on an invented precedent. func (s *Server) mergeResourceParts(w http.ResponseWriter, r *http.Request, p, sc, name string, mode engine.ResourceMode) { target, ok := s.resourceTarget(w, r, name) if !ok { @@ -290,8 +299,9 @@ func (s *Server) mergeResourceParts(w http.ResponseWriter, r *http.Request, p, s if size := r.Header.Get("x-odps-resource-merge-total-bytes"); size != "" { declared, e := strconv.ParseInt(size, 10, 64) if e != nil || declared != int64(len(merged)) { + // Stricter than the service, and therefore no precedent for consuming parts: + // keep everything the client uploaded so it can correct the header and retry. fail(w, r, 400, "InvalidParameter", fmt.Errorf("x-odps-resource-merge-total-bytes does not match the merged payload")) - consumeParts() return } } diff --git a/internal/server/resources_test.go b/internal/server/resources_test.go index e1b8857..e7ff95c 100644 --- a/internal/server/resources_test.go +++ b/internal/server/resources_test.go @@ -180,9 +180,10 @@ func TestResourceRESTContract(t *testing.T) { if w.Code != 404 { t.Fatalf("parts must be cleaned after merge: %d %s", w.Code, w.Body.String()) } - // A mismatching digest or byte count must not publish a partial payload, and - // the merge consumes the chunk it assembled from: an unverifiable part is - // dead weight, and both SDKs re-upload the same deterministic name on retry. + // A mismatching digest must not publish a partial payload, and the merge consumes the + // chunk it assembled from. That matches the live service: a refused digest leaves both + // part resources gone (measured on three_pangu2_odps2, 2026-10-04 19:23), which is what + // makes the retry-cost note in resources.go worth reading before changing this. for _, body := range []string{md5hex("nope") + "|" + partB, "0123456789abcdef0123456789abcdef|" + partB} { resourceRequest(t, s, "POST", "/projects/p/resources?rIsPart=true", partB, "file", temp, "x") w = resourceRequest(t, s, "POST", "/projects/p/resources?rOpMerge=true", "bad.py", "py", nil, body) @@ -201,6 +202,14 @@ func TestResourceRESTContract(t *testing.T) { if w.Code != 400 || !strings.Contains(w.Body.String(), "merge-total-bytes") { t.Fatalf("declared size: %d %s", w.Code, w.Body.String()) } + // The emulator is stricter than the service here: the service never compares the + // declared byte count with what it assembled (measured: declaring 4400 for a 304-byte + // merge is accepted), so there is no service precedent for what it does with the parts + // on this path. It keeps them, and that is asserted rather than assumed - the whole + // point is that this emulator does not destroy an upload on a check the service lacks. + if w = resourceRequest(t, s, "GET", "/projects/p/resources/"+partB, "", "", nil, ""); w.Code != 200 { + t.Fatalf("a refusal on the emulator's own stricter check must keep the client's part: %d %s", w.Code, w.Body.String()) + } w = resourceRequest(t, s, "POST", "/projects/p/resources?rOpMerge=true", "bad.py", "py", nil, md5hex("x")+"|absent_part") if w.Code != 404 || !strings.Contains(w.Body.String(), "NoSuchObject") { t.Fatalf("missing part: %d %s", w.Code, w.Body.String()) diff --git a/tests/python/run.py b/tests/python/run.py index 2e31e7a..72d41b4 100644 --- a/tests/python/run.py +++ b/tests/python/run.py @@ -431,6 +431,88 @@ def merge_with_broken_manifest(): case("a merge body that is not a manifest is refused", merge_with_broken_manifest) +def parts_after_refusal(name, digest, declared, upload_part=True, precreate=False): + """Upload one part, refuse the merge, and report what happened to that part. + + The three refusals below are the shapes measured against a live MaxCompute project on + 2026-10-04 (raw read-backs in the workspace work item's evidence). The probe asserts the + emulator's answer in each case, so the difference between "the service consumed it" and + "the emulator kept it" is a checked fact rather than someone's recollection. + """ + if precreate: + try: + odps.delete_resource(name) + except Exception: + pass + odps.create_resource(name, "file", fileobj=b"already-here", temp=True) + part = name + ".part.tmp.probe.0" + headers = {"Content-Type": "application/octet-stream", "x-odps-resource-type": "file", + "x-odps-resource-name": part, "x-odps-resource-istemp": "true"} + raw("/projects/%s/resources?rIsPart&curr_project=%s" % (PROJECT, PROJECT), method="POST", + body=BLOB[:CHUNK], headers=headers) + manifest = (digest + "|" + part).encode() + merge_headers = {"Content-Type": "application/octet-stream", "x-odps-resource-type": "file", + "x-odps-resource-name": name, + "x-odps-resource-merge-total-bytes": str(declared)} + code, body, _ = raw_error("/projects/%s/resources?rOpMerge&curr_project=%s" % (PROJECT, PROJECT), + method="POST", body=manifest, headers=merge_headers) + kept = odps.exist_resource(part) + for leftover in (part, name): + try: + odps.delete_resource(leftover) + except Exception: + pass + return code, body, kept + + +def merge_wrong_digest_consumes_the_part(): + name = PREFIX + "_consumed.bin" + code, body, kept = parts_after_refusal(name, hashlib.md5(b"wrong").hexdigest(), CHUNK) + eq(code, 400, "status") + eq("MD5" in body, True, "reason mentions the digest") + # Measured on the service: a merge rejected on the digest also drops the parts it read. + eq(kept, False, "the refused merge consumed the part, as the service does") + return "400, part consumed" + + +case("a digest-refused merge consumes its part, like the service", merge_wrong_digest_consumes_the_part) + + +def merge_target_exists_keeps_the_part(): + name = PREFIX + "_dupmerge.bin" + code, body, kept = parts_after_refusal(name, hashlib.md5(BLOB[:CHUNK]).hexdigest(), CHUNK, + precreate=True) + # What the case asserts is the part lifecycle, not the status code: the service answers + # this request with `ODPS-0421121 The resource has already existed` (its HTTP status was + # not read back here), the emulator answers 400 InvalidParameter. That difference is + # pre-existing and out of this change's scope; both refuse before reading a part. + eq(200 <= code < 300, False, "the merge over an existing target is refused (%s %s)" % (code, body[:60])) + # Measured on the service: the refusal is decided before any part is read, so the part + # survives and the client can re-point the merge at another name. + eq(kept, True, "the part survived a refusal decided before it was read") + return "%s, part kept" % code + + +case("a merge refused over an existing target keeps its part", merge_target_exists_keeps_the_part) + + +def merge_declared_bytes_mismatch_keeps_the_part(): + name = PREFIX + "_declared.bin" + digest = hashlib.md5(BLOB[:CHUNK]).hexdigest() + code, body, kept = parts_after_refusal(name, digest, CHUNK + 4096) + eq(code, 400, "status") + eq("merge-total-bytes" in body, True, "reason names the header") + # The service has no such check at all - declaring the wrong total with a correct digest + # is accepted there and the merge lands. So this emulator is stricter, and it must not + # destroy the upload on a failure the service cannot produce: the part stays. + eq(kept, True, "the emulator's own stricter check does not eat the client's upload") + return "400 (stricter than the service), part kept" + + +case("a declared-byte-count refusal keeps the part (emulator-only strictness)", + merge_declared_bytes_mismatch_keeps_the_part) + + def merge_with_wrong_md5(): name = PREFIX + "_badmd5.bin" odps.create_resource(name, "file", fileobj=b"seed")