Skip to content

Commit 1675929

Browse files
jopemachineclaude
andcommitted
feat(BA-6159): add /upload/status endpoint and progress headers
Standard TUS HEAD only returns Upload-Offset, which is ambiguous when chunks arrive out of order: the client cannot tell which byte ranges are still missing. This sub-task surfaces the session metadata so resume is precise. - GET /upload/status?token=… returns JSON: session, total_size, committed_offset, chunks_received (list of offsets), missing_ranges (list of {offset, length} gaps), progress_percent, status. Authenticated via the same upload JWT. - PATCH responses gain three additive headers: X-Backend-Ai-Chunks-Received, X-Backend-Ai-Progress-Percent, X-Backend-Ai-Total-Expected. Headers are exposed via CORS Access-Control-Expose-Headers so browser clients can read them. Tests cover the status endpoint for empty / partial / out-of-order / missing-session cases plus the new PATCH response headers. Resolves BA-6159. Part of epic BA-6153 (implements BA-3974). Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
1 parent 583cb92 commit 1675929

3 files changed

Lines changed: 223 additions & 3 deletions

File tree

changes/11770.feature.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Add `GET /upload/status` endpoint and `X-Backend-Ai-*` progress response headers so clients can recover after a storage proxy crash by resending only the byte ranges still missing.

src/ai/backend/storage/api/client.py

Lines changed: 85 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,10 @@
5050
from ai.backend.storage.services.file_stream.zip import (
5151
ZipArchiveStreamReader,
5252
)
53-
from ai.backend.storage.services.upload.tus_session import TusUploadSession
53+
from ai.backend.storage.services.upload.tus_session import (
54+
SessionState,
55+
TusUploadSession,
56+
)
5457
from ai.backend.storage.types import SENTINEL, TusChunkUploadStreamReader
5558
from ai.backend.storage.utils import (
5659
CheckParamSource,
@@ -459,21 +462,90 @@ class Params(TypedDict):
459462
upload_offset=state.committed_offset,
460463
upload_length=total_size,
461464
)
465+
headers.update(_progress_response_headers(state))
462466
return web.Response(status=HTTPStatus.NO_CONTENT, headers=headers)
463467

464468

469+
async def tus_upload_status(request: web.Request) -> web.Response:
470+
"""
471+
GET /upload/status — JSON snapshot of an in-flight upload session.
472+
473+
Used by clients to recover after a Storage Proxy crash or network error:
474+
they discover which byte ranges are still missing and resend only those.
475+
"""
476+
ctx: RootContext = request.app["ctx"]
477+
secret = ctx.local_config.storage_proxy.secret
478+
479+
class Params(TypedDict):
480+
token: UploadTokenData
481+
dst_dir: str
482+
483+
async with cast(
484+
AbstractAsyncContextManager[Params],
485+
check_params(
486+
request,
487+
t.Dict(
488+
{
489+
t.Key("token"): tx.JsonWebToken(
490+
secret=secret,
491+
inner_iv=upload_token_data_iv,
492+
),
493+
t.Key("dst_dir", default=None): t.Null | t.String,
494+
},
495+
),
496+
read_from=CheckParamSource.QUERY,
497+
),
498+
) as params:
499+
token_data = params["token"]
500+
total_size = int(token_data["size"])
501+
async with ctx.get_volume(token_data["volume"]) as volume:
502+
session_dir = _resolve_session_dir(volume, token_data)
503+
if not session_dir.exists():
504+
raise web.HTTPNotFound(
505+
body=dump_json_str(
506+
{
507+
"title": "No such upload session",
508+
"type": ("https://api.backend.ai/probs/storage/no-such-upload-session"),
509+
},
510+
),
511+
content_type="application/problem+json",
512+
)
513+
session = TusUploadSession(
514+
session_dir,
515+
session_id=token_data["session"],
516+
total_size=total_size,
517+
)
518+
state = await session.read_state()
519+
body = {
520+
"session": token_data["session"],
521+
"total_size": total_size,
522+
"committed_offset": state.committed_offset,
523+
"chunks_received": [rec.offset for rec in state.received],
524+
"missing_ranges": [
525+
{"offset": offset, "length": length}
526+
for offset, length in state.missing_ranges()
527+
],
528+
"progress_percent": state.progress_percent(),
529+
"status": state.status,
530+
}
531+
return web.json_response(body, status=HTTPStatus.OK)
532+
533+
465534
def _resolve_session_dir(volume: AbstractVolume, token_data: UploadTokenData) -> Path:
466535
return volume.mangle_vfpath(token_data["vfid"]) / ".upload" / token_data["session"]
467536

468537

469538
_TUS_HEADER_LIST = "Tus-Resumable, Upload-Length, Upload-Metadata, Upload-Offset, Content-Type"
539+
_PROGRESS_HEADER_LIST = (
540+
"X-Backend-Ai-Chunks-Received, X-Backend-Ai-Progress-Percent, X-Backend-Ai-Total-Expected"
541+
)
470542

471543

472544
def _tus_response_headers(*, upload_offset: int, upload_length: int) -> dict[str, str]:
473545
return {
474546
"Access-Control-Allow-Origin": "*",
475547
"Access-Control-Allow-Headers": _TUS_HEADER_LIST,
476-
"Access-Control-Expose-Headers": _TUS_HEADER_LIST,
548+
"Access-Control-Expose-Headers": f"{_TUS_HEADER_LIST}, {_PROGRESS_HEADER_LIST}",
477549
"Access-Control-Allow-Methods": "*",
478550
"Cache-Control": "no-store",
479551
"Tus-Resumable": "1.0.0",
@@ -482,6 +554,14 @@ def _tus_response_headers(*, upload_offset: int, upload_length: int) -> dict[str
482554
}
483555

484556

557+
def _progress_response_headers(state: SessionState) -> dict[str, str]:
558+
return {
559+
"X-Backend-Ai-Chunks-Received": str(len(state.received)),
560+
"X-Backend-Ai-Progress-Percent": str(state.progress_percent()),
561+
"X-Backend-Ai-Total-Expected": str(state.total_size),
562+
}
563+
564+
485565
async def tus_options(request: web.Request) -> web.Response:
486566
"""
487567
Let clients discover the supported features of our tus.io server-side implementation.
@@ -580,4 +660,7 @@ async def init_client_app(ctx: RootContext) -> web.Application:
580660
r.add_route("HEAD", tus_check_session)
581661
r.add_route("PATCH", tus_upload_part)
582662

663+
status_resource = app.router.add_resource("/upload/status")
664+
status_resource.add_route("GET", tus_upload_status)
665+
583666
return app

tests/unit/storage/api/test_tus_upload.py

Lines changed: 137 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,14 +10,15 @@
1010
from __future__ import annotations
1111

1212
import dataclasses
13+
import json
1314
from pathlib import Path
1415
from typing import Any
1516
from unittest.mock import AsyncMock, MagicMock, patch
1617

1718
import pytest
1819
from aiohttp import web
1920

20-
from ai.backend.storage.api.client import tus_upload_part
21+
from ai.backend.storage.api.client import tus_upload_part, tus_upload_status
2122
from ai.backend.storage.errors import InvalidAPIParameters, UploadOffsetMismatchError
2223

2324

@@ -298,3 +299,138 @@ async def test_chunk_exceeding_declared_size_raises(self, patch_env: _PatchEnv)
298299
chunks_dir = patch_env.session_dir / "chunks"
299300
if chunks_dir.exists():
300301
assert list(chunks_dir.glob("*.tmp")) == []
302+
303+
304+
def _body_bytes(response: web.Response) -> bytes:
305+
body = response.body
306+
assert isinstance(body, (bytes, bytearray))
307+
return bytes(body)
308+
309+
310+
class TestProgressHeaders:
311+
async def test_patch_response_includes_progress_headers(self, patch_env: _PatchEnv) -> None:
312+
token_data = _token_data(session_id="test-session", total_size=2048, relpath="result.bin")
313+
cp = _patch_handler_params(token_data)
314+
try:
315+
request = _build_request(
316+
vfpath=patch_env.vfpath,
317+
session_id="test-session",
318+
total_size=2048,
319+
body=b"A" * 1024,
320+
offset_header="0",
321+
)
322+
response = await tus_upload_part(request)
323+
finally:
324+
cp.stop()
325+
326+
assert response.headers["X-Backend-Ai-Chunks-Received"] == "1"
327+
assert response.headers["X-Backend-Ai-Total-Expected"] == "2048"
328+
assert response.headers["X-Backend-Ai-Progress-Percent"] == "50.0"
329+
330+
331+
class TestUploadStatus:
332+
async def test_status_for_empty_session(self, patch_env: _PatchEnv) -> None:
333+
token_data = _token_data(session_id="test-session", total_size=2048, relpath="result.bin")
334+
request = _build_request(
335+
vfpath=patch_env.vfpath,
336+
session_id="test-session",
337+
total_size=2048,
338+
body=None,
339+
offset_header=None,
340+
)
341+
cp = _patch_handler_params(token_data)
342+
try:
343+
response = await tus_upload_status(request)
344+
finally:
345+
cp.stop()
346+
347+
payload = json.loads(_body_bytes(response))
348+
assert payload["total_size"] == 2048
349+
assert payload["committed_offset"] == 0
350+
assert payload["chunks_received"] == []
351+
assert payload["missing_ranges"] == [{"offset": 0, "length": 2048}]
352+
assert payload["progress_percent"] == 0.0
353+
assert payload["status"] == "pending"
354+
355+
async def test_status_after_partial_upload(self, patch_env: _PatchEnv) -> None:
356+
token_data = _token_data(session_id="test-session", total_size=4096, relpath="result.bin")
357+
cp = _patch_handler_params(token_data)
358+
try:
359+
patch_request = _build_request(
360+
vfpath=patch_env.vfpath,
361+
session_id="test-session",
362+
total_size=4096,
363+
body=b"A" * 1024,
364+
offset_header="0",
365+
)
366+
await tus_upload_part(patch_request)
367+
368+
status_request = _build_request(
369+
vfpath=patch_env.vfpath,
370+
session_id="test-session",
371+
total_size=4096,
372+
body=None,
373+
offset_header=None,
374+
)
375+
response = await tus_upload_status(status_request)
376+
finally:
377+
cp.stop()
378+
379+
payload = json.loads(_body_bytes(response))
380+
assert payload["committed_offset"] == 1024
381+
assert payload["chunks_received"] == [0]
382+
assert payload["missing_ranges"] == [{"offset": 1024, "length": 3072}]
383+
assert payload["progress_percent"] == 25.0
384+
assert payload["status"] == "pending"
385+
386+
async def test_status_with_out_of_order_chunks_reports_gaps(self, patch_env: _PatchEnv) -> None:
387+
token_data = _token_data(session_id="test-session", total_size=4096, relpath="result.bin")
388+
cp = _patch_handler_params(token_data)
389+
try:
390+
# Commit the *second* chunk first, leaving [0, 1024) missing.
391+
second = _build_request(
392+
vfpath=patch_env.vfpath,
393+
session_id="test-session",
394+
total_size=4096,
395+
body=b"B" * 1024,
396+
offset_header="1024",
397+
)
398+
await tus_upload_part(second)
399+
400+
status_request = _build_request(
401+
vfpath=patch_env.vfpath,
402+
session_id="test-session",
403+
total_size=4096,
404+
body=None,
405+
offset_header=None,
406+
)
407+
response = await tus_upload_status(status_request)
408+
finally:
409+
cp.stop()
410+
411+
payload = json.loads(_body_bytes(response))
412+
# Contiguous prefix has not advanced because [0,1024) is missing.
413+
assert payload["committed_offset"] == 0
414+
assert payload["chunks_received"] == [1024]
415+
assert payload["missing_ranges"] == [
416+
{"offset": 0, "length": 1024},
417+
{"offset": 2048, "length": 2048},
418+
]
419+
420+
async def test_status_for_missing_session_returns_404(self, tmp_path: Path) -> None:
421+
vfpath = tmp_path / "vfpath"
422+
vfpath.mkdir()
423+
token_data = _token_data(session_id="missing", total_size=1024, relpath="result.bin")
424+
request = _build_request(
425+
vfpath=vfpath,
426+
session_id="missing",
427+
total_size=1024,
428+
body=None,
429+
offset_header=None,
430+
)
431+
cp = _patch_handler_params(token_data)
432+
try:
433+
with pytest.raises(web.HTTPNotFound):
434+
await tus_upload_status(request)
435+
finally:
436+
cp.stop()

0 commit comments

Comments
 (0)