diff --git a/CHANGELOG.md b/CHANGELOG.md index 93e2f3f..46e231f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,9 +9,11 @@ - 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. +- 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. - Add MCQA / SQLRT interactive sessions: an `SQLRT` task instance stays `Running` until the client stops it or it idles out, and its statements run as sub queries over the instance information KV (`?info`) using the Java SDK's own object status codes, with `query`/`cancel` writes and `status`/`progress`/`result`/`result_` reads. - Declare `mcqa` in `/capabilities`, and name the remaining session gaps (named-session attach, MaxQA v2) separately from the surface that works. - Download an MCQA sub query's result over the instance tunnel (`?data&cached&taskname=..&queryid=..`): a record stream that carries its own schema, the `odps-tunnel-record-count` the SDK pages on, `rowrange` paging, `sizelimit` truncation and the `READ_TABLE_MAX_ROW` cap, with `InstanceTypeNotSupported` for statements that have no result set. This is the Java SDK's default interactive fetch and JDBC MaxQA's read path. 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 14e2bf2..222a920 100644 --- a/internal/server/resources.go +++ b/internal/server/resources.go @@ -238,6 +238,25 @@ 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. +// +// 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 { @@ -272,9 +291,16 @@ 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)) { + // 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")) return } @@ -282,16 +308,17 @@ func (s *Server) mergeResourceParts(w http.ResponseWriter, r *http.Request, p, s 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..e7ff95c 100644 --- a/internal/server/resources_test.go +++ b/internal/server/resources_test.go @@ -180,7 +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. + // 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) @@ -190,12 +193,23 @@ 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) 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()) @@ -209,6 +223,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..72d41b4 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)) @@ -417,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") @@ -432,11 +528,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)