Skip to content

Stage 13 · 条件请求与 CAS

目标

把 ETag 变成缓存校验器与串行化的 compare-and-swap 前置条件。

动手任务

从stage-12开始,实现 ETag 匹配,并在持有 _lock 时处理对象操作的 if_match/if_none_match。 行为必须留在下列源码同构边界中;不要先复制补丁。

交付文件

  • src/minis3/__init__.py
  • src/minis3/conditional.py
  • src/minis3/errors.py
  • src/minis3/store.py
  • tests/test_conditional.py

自查

  1. 本阶段的可见性或状态迁移由谁负责?

    答案

    比较与发布必须处于同一临界区,否则两个写者都可能成功。

  2. 如果绕过新边界,哪个测试会最先失败?

    答案

    阅读 tests.txt,找出最窄的新节点,并说出它覆盖的公开调用。

通关命令

uv run pytest -q $(cat journey/stages/13-conditional-cas/tests.txt)

对应真实 S3 的一课

比较与发布必须处于同一临界区,否则两个写者都可能成功。

教材

第 7 章

在 GitHub 查看阶段差异

完成后可运行 git checkout stage-13 对照你的结果。

先做后看:stage.patch
diff --git a/src/minis3/__init__.py b/src/minis3/__init__.py
index 0c23aea..3f6e582 100644
--- a/src/minis3/__init__.py
+++ b/src/minis3/__init__.py
@@ -7,3 +7,4 @@ from .storage import InjectedCrash
 from .listing import ListedObject, ListedVersion, ListObjectsResult, ListObjectVersionsResult
 from .errors import EntityTooSmall, InvalidPart, InvalidPartOrder, NoSuchUpload
 from .multipart import MIN_PART_SIZE, MultipartPart, MultipartUpload
+from .errors import NotModified, PreconditionFailed
diff --git a/src/minis3/conditional.py b/src/minis3/conditional.py
new file mode 100644
index 0000000..d13ccf3
--- /dev/null
+++ b/src/minis3/conditional.py
@@ -0,0 +1,35 @@
+"""Pure HTTP-style ETag precondition evaluation.
+
+The service evaluates these functions while holding its mutation lock. That
+placement matters: checking an ETag and publishing a replacement must be one
+serialized compare-and-swap operation, not two individually safe calls.
+"""
+
+from __future__ import annotations
+
+from .errors import NotModified, PreconditionFailed
+
+
+def etag_matches(condition: str, current_etag: str | None) -> bool:
+    """Return whether a simplified ETag header matches the current object."""
+
+    candidates = tuple(item.strip() for item in condition.split(","))
+    if "*" in candidates:
+        return current_etag is not None
+    return current_etag is not None and current_etag in candidates
+
+
+def require_if_match(current_etag: str | None, condition: str | None) -> None:
+    """Raise S3's named 412 outcome when If-Match is not satisfied."""
+
+    if condition is not None and not etag_matches(condition, current_etag):
+        raise PreconditionFailed(condition)
+
+
+def require_if_none_match(
+    current_etag: str | None, condition: str | None
+) -> None:
+    """Raise the body-less 304 control outcome when a cached ETag matches."""
+
+    if condition is not None and etag_matches(condition, current_etag):
+        raise NotModified(condition)
diff --git a/src/minis3/errors.py b/src/minis3/errors.py
index 9db3b4c..5f255e0 100644
--- a/src/minis3/errors.py
+++ b/src/minis3/errors.py
@@ -43,3 +43,11 @@ class InvalidPartOrder(MiniS3Error):

 class EntityTooSmall(MiniS3Error):
     """A non-final multipart part is below the configured minimum size."""
+
+
+class PreconditionFailed(MiniS3Error):
+    """An If-Match condition failed (the S3-shaped HTTP 412 outcome)."""
+
+
+class NotModified(MiniS3Error):
+    """An If-None-Match condition matched (the HTTP 304 control outcome)."""
diff --git a/src/minis3/store.py b/src/minis3/store.py
index 9b50aa2..e47e1ac 100644
--- a/src/minis3/store.py
+++ b/src/minis3/store.py
@@ -8,8 +8,9 @@ from pathlib import Path
 from threading import RLock
 from time import time

+from .conditional import require_if_match, require_if_none_match
 from .bucket import Bucket, SequenceCounter, VersioningState
-from .errors import BucketAlreadyExists, BucketNotEmpty, NoSuchBucket
+from .errors import BucketAlreadyExists, BucketNotEmpty, NoSuchBucket, NoSuchKey, NoSuchVersion
 from .listing import ListObjectsResult, ListObjectVersionsResult, list_object_versions, list_objects
 from .model import ObjectVersion, Version
 from .multipart import (
@@ -78,20 +79,39 @@ class MiniS3:
             self._buckets[name] = candidate


-    def put_object(self, bucket: str, key: str, body: bytes) -> Version:
+    def put_object(
+        self,
+        bucket: str,
+        key: str,
+        body: bytes,
+        *,
+        if_match: str | None = None,
+    ) -> Version:
         with self._lock:
             candidate = deepcopy(self._bucket(bucket))
-            result = candidate.put(key, body, self._counter)
+            require_if_match(self._current_etag(candidate, key), if_match)
+            result = candidate.put(
+                key, body, self._counter, now=self._clock()
+            )
             self._storage.persist_bucket(candidate)
             self._buckets[bucket] = candidate
             return result


     def get_object(
-        self, bucket: str, key: str, *, version_id: str | None = None
+        self,
+        bucket: str,
+        key: str,
+        *,
+        version_id: str | None = None,
+        if_match: str | None = None,
+        if_none_match: str | None = None,
     ) -> Version:
         with self._lock:
-            return self._bucket(bucket).get(key, version_id)
+            result = self._bucket(bucket).get(key, version_id)
+            require_if_match(result.etag, if_match)
+            require_if_none_match(result.etag, if_none_match)
+            return result


     def head_object(
@@ -103,11 +123,21 @@ class MiniS3:


     def delete_object(
-        self, bucket: str, key: str, *, version_id: str | None = None
+        self,
+        bucket: str,
+        key: str,
+        *,
+        version_id: str | None = None,
+        if_match: str | None = None,
     ) -> ObjectVersion | None:
         with self._lock:
             candidate = deepcopy(self._bucket(bucket))
-            result = candidate.delete(key, self._counter, version_id)
+            require_if_match(
+                self._addressed_etag(candidate, key, version_id), if_match
+            )
+            result = candidate.delete(
+                key, self._counter, version_id, now=self._clock()
+            )
             self._storage.persist_bucket(candidate)
             self._buckets[bucket] = candidate
             return result
@@ -220,6 +250,24 @@ class MiniS3:
             self._storage.remove_multipart_upload(bucket, key, upload_id)


+    @staticmethod
+    def _current_etag(bucket: Bucket, key: str) -> str | None:
+        try:
+            return bucket.get(key).etag
+        except NoSuchKey:
+            return None
+
+
+    @staticmethod
+    def _addressed_etag(
+        bucket: Bucket, key: str, version_id: str | None
+    ) -> str | None:
+        try:
+            return bucket.get(key, version_id).etag
+        except (NoSuchKey, NoSuchVersion):
+            return None
+
+
     def _bucket(self, name: str) -> Bucket:
         try:
             return self._buckets[name]
diff --git a/tests/test_conditional.py b/tests/test_conditional.py
new file mode 100644
index 0000000..137e7a1
--- /dev/null
+++ b/tests/test_conditional.py
@@ -0,0 +1,81 @@
+"""Conditional requests turn current ETags into an object-level CAS token."""
+
+from concurrent.futures import ThreadPoolExecutor
+from pathlib import Path
+from threading import Barrier
+
+import pytest
+
+from minis3 import MiniS3, NoSuchKey, NotModified, PreconditionFailed
+
+
+def test_get_if_none_match_has_304_semantics_and_if_match_has_412(
+    tmp_path: Path,
+) -> None:
+    store = MiniS3(tmp_path)
+    store.create_bucket("b")
+    current = store.put_object("b", "k", b"value")
+
+    with pytest.raises(NotModified):
+        store.get_object("b", "k", if_none_match=current.etag)
+    with pytest.raises(NotModified):
+        store.get_object("b", "k", if_none_match="*")
+    with pytest.raises(PreconditionFailed):
+        store.get_object(
+            "b", "k", if_match='"00000000000000000000000000000000"'
+        )
+    assert store.get_object("b", "k", if_match=current.etag) == current
+
+
+def test_put_and_delete_if_match_compare_against_current_visible_etag(
+    tmp_path: Path,
+) -> None:
+    store = MiniS3(tmp_path)
+    store.create_bucket("b")
+    initial = store.put_object("b", "k", b"old")
+    winner = store.put_object("b", "k", b"new", if_match=initial.etag)
+
+    with pytest.raises(PreconditionFailed):
+        store.put_object("b", "k", b"stale", if_match=initial.etag)
+    with pytest.raises(PreconditionFailed):
+        store.delete_object("b", "k", if_match=initial.etag)
+
+    removed = store.delete_object("b", "k", if_match=winner.etag)
+    assert removed is None
+    with pytest.raises(NoSuchKey):
+        store.get_object("b", "k")
+
+
+def test_if_match_wildcard_requires_a_current_visible_object(tmp_path: Path) -> None:
+    store = MiniS3(tmp_path)
+    store.create_bucket("b")
+
+    with pytest.raises(PreconditionFailed):
+        store.put_object("b", "missing", b"x", if_match="*")
+    with pytest.raises(PreconditionFailed):
+        store.delete_object("b", "missing", if_match="*")
+
+    store.put_object("b", "present", b"x")
+    assert store.put_object("b", "present", b"y", if_match="*").body == b"y"
+
+
+def test_two_conditional_writers_have_exactly_one_winner(tmp_path: Path) -> None:
+    store = MiniS3(tmp_path)
+    store.create_bucket("b")
+    observed = store.put_object("b", "counter", b"0").etag
+    barrier = Barrier(2)
+
+    def writer(value: bytes) -> str:
+        barrier.wait()
+        try:
+            store.put_object("b", "counter", value, if_match=observed)
+        except PreconditionFailed:
+            return "412"
+        return "stored"
+
+    with ThreadPoolExecutor(max_workers=2) as pool:
+        outcomes = list(pool.map(writer, (b"writer-a", b"writer-b")))
+
+    assert sorted(outcomes) == ["412", "stored"]
+    assert store.get_object("b", "counter").body in {b"writer-a", b"writer-b"}
+