Skip to content

Commit 1fcdce3

Browse files
phernandezclaude
andcommitted
fix(cli): replay the opened file on upload retries and read redirect bodies
An upload retry re-opened the source path, so a file replaced during a Retry-After wait could be sent under the original Content-Length and mtime. The file is now opened once; each attempt rewinds that handle, and size and mtime come from fstat on it. A streamed 3xx response was left unread, so describing the failure raised ResponseNotRead instead of naming the status. Every non-2xx response is now read before it is handed back. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01APFUk2bjEwMptMQqpRhjea Signed-off-by: phernandez <paul@basicmachines.co>
1 parent 7fa5a9a commit 1fcdce3

2 files changed

Lines changed: 72 additions & 26 deletions

File tree

‎src/basic_memory/cli/commands/cloud/webdav.py‎

Lines changed: 35 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
"""
2121

2222
import asyncio
23+
import os
2324
import re
2425
import xml.etree.ElementTree as ElementTree
2526
from collections.abc import AsyncIterator, Callable
@@ -211,8 +212,9 @@ async def _rate_limited_request(
211212
that makes a retry something other than the same request again.
212213
213214
A successful response body is left unread, so a download can go to disk as it
214-
arrives. An error response is read in full: it is small, and the caller's
215-
error message quotes it. The response is closed when the block exits.
215+
arrives. Any other response (a redirect from a proxy as well as an error) is
216+
read in full: it is small, and the caller's error message quotes it. The
217+
response is closed when the block exits.
216218
217219
A 429 on the final attempt is yielded rather than raised, so the caller's
218220
own error handling reports it with the rate-limit detail attached.
@@ -232,20 +234,24 @@ async def _rate_limited_request(
232234
attempt += 1
233235

234236
try:
235-
if response.is_error:
237+
if not response.is_success:
236238
await response.aread()
237239
yield response
238240
finally:
239241
await response.aclose()
240242

241243

242-
def _file_chunks(source: Path) -> Callable[[], AsyncIterator[bytes]]:
243-
"""A body factory that reads ``source`` from the start on every call."""
244+
def _file_chunks(stream: BinaryIO) -> Callable[[], AsyncIterator[bytes]]:
245+
"""A body factory that rewinds one open file and reads it on every call.
246+
247+
The file is opened once by the caller, so a retry replays the file that was
248+
validated and measured, even if the path is replaced during a Retry-After wait.
249+
"""
244250

245251
async def chunks() -> AsyncIterator[bytes]:
246-
with source.open("rb") as stream:
247-
while chunk := stream.read(_STREAM_CHUNK_BYTES):
248-
yield chunk
252+
stream.seek(0)
253+
while chunk := stream.read(_STREAM_CHUNK_BYTES):
254+
yield chunk
249255

250256
return chunks
251257

@@ -361,25 +367,28 @@ async def upload_file(
361367
WebdavError: If the service refuses the upload for any other reason.
362368
"""
363369
request_path = webdav_path(project, rel_path)
364-
stat = source.stat()
365-
headers = {
366-
"X-OC-Mtime": str(int(stat.st_mtime)),
367-
"Content-Length": str(stat.st_size),
368-
}
369-
if create_only:
370-
headers["If-None-Match"] = "*"
370+
with source.open("rb") as stream:
371+
# Measured from the open handle, so the declared size and mtime describe
372+
# exactly the file every attempt sends.
373+
stat = os.fstat(stream.fileno())
374+
headers = {
375+
"X-OC-Mtime": str(int(stat.st_mtime)),
376+
"Content-Length": str(stat.st_size),
377+
}
378+
if create_only:
379+
headers["If-None-Match"] = "*"
371380

372-
try:
373-
async with _rate_limited_request(
374-
client, "PUT", request_path, content=_file_chunks(source), headers=headers
375-
) as response:
376-
# Checked before raise_for_status: a refused precondition is the answer
377-
# this call asked for, not a failure.
378-
if create_only and response.status_code == httpx.codes.PRECONDITION_FAILED:
379-
return False
380-
response.raise_for_status()
381-
except httpx.HTTPError as exc:
382-
raise WebdavError(f"Failed to upload {rel_path}: {_describe(exc)}") from exc
381+
try:
382+
async with _rate_limited_request(
383+
client, "PUT", request_path, content=_file_chunks(stream), headers=headers
384+
) as response:
385+
# Checked before raise_for_status: a refused precondition is the
386+
# answer this call asked for, not a failure.
387+
if create_only and response.status_code == httpx.codes.PRECONDITION_FAILED:
388+
return False
389+
response.raise_for_status()
390+
except httpx.HTTPError as exc:
391+
raise WebdavError(f"Failed to upload {rel_path}: {_describe(exc)}") from exc
383392

384393
return True
385394

‎tests/cli/cloud/test_webdav_client.py‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77

88
import io
99
import os
10+
import sys
1011
from contextlib import asynccontextmanager
1112
from datetime import datetime, timezone
1213
from pathlib import Path
@@ -535,6 +536,42 @@ async def handler(request: httpx.Request) -> httpx.Response:
535536
assert fetched == ["/webdav/research/a#draft.md"]
536537

537538

539+
@pytest.mark.asyncio
540+
async def test_a_redirect_is_reported_with_its_status():
541+
"""A proxy's 3xx is not a success; the message must name it, not a stream error."""
542+
543+
async def handler(request: httpx.Request) -> httpx.Response:
544+
return httpx.Response(302, text="moved", headers={"Location": "/login"})
545+
546+
async with _client(handler) as client:
547+
with pytest.raises(WebdavError, match="HTTP 302 - moved"):
548+
await download_file(client, "research", "a.md", io.BytesIO())
549+
550+
551+
# Windows refuses to replace a file another handle holds open, so the race this
552+
# guards against cannot happen there.
553+
@pytest.mark.skipif(sys.platform == "win32", reason="cannot replace an open file on Windows")
554+
@pytest.mark.asyncio
555+
async def test_upload_retry_sends_the_file_it_opened(recorded_waits, tmp_path):
556+
"""Replacing the path during a Retry-After wait does not change what is replayed."""
557+
source = _local_file(tmp_path, b"original")
558+
bodies: list[bytes] = []
559+
560+
async def handler(request: httpx.Request) -> httpx.Response:
561+
bodies.append(request.content)
562+
if len(bodies) == 1:
563+
replacement = tmp_path / "replacement.md"
564+
replacement.write_bytes(b"swapped!")
565+
os.replace(replacement, source)
566+
return _rate_limited()
567+
return httpx.Response(201)
568+
569+
async with _client(handler) as client:
570+
await upload_file(client, "research", "a.md", source=source)
571+
572+
assert bodies == [b"original", b"original"]
573+
574+
538575
@pytest.mark.asyncio
539576
async def test_upload_file_create_only_sends_the_conditional_header(tmp_path):
540577
seen: dict[str, object] = {}

0 commit comments

Comments
 (0)