Stage 13 · 条件请求与 CAS¶
目标¶
把 ETag 变成缓存校验器与串行化的 compare-and-swap 前置条件。
动手任务¶
从stage-12开始,实现 ETag 匹配,并在持有 _lock 时处理对象操作的 if_match/if_none_match。 行为必须留在下列源码同构边界中;不要先复制补丁。
交付文件¶
src/minis3/__init__.pysrc/minis3/conditional.pysrc/minis3/errors.pysrc/minis3/store.pytests/test_conditional.py
自查¶
-
本阶段的可见性或状态迁移由谁负责?
答案
比较与发布必须处于同一临界区,否则两个写者都可能成功。
-
如果绕过新边界,哪个测试会最先失败?
答案
阅读
tests.txt,找出最窄的新节点,并说出它覆盖的公开调用。
通关命令¶
uv run pytest -q $(cat journey/stages/13-conditional-cas/tests.txt)
对应真实 S3 的一课¶
比较与发布必须处于同一临界区,否则两个写者都可能成功。
教材¶
完成后可运行 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"}
+