Skip to content
Open
Show file tree
Hide file tree
Changes from 3 commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
6396b1c
UploadTracker
Dreamsorcerer Aug 29, 2026
94e6aac
Update CHANGES.rst
Dreamsorcerer Aug 29, 2026
f388966
Apply batched suggestions from code review
Dreamsorcerer Aug 29, 2026
3ea4c12
Fix clipped sizes
Dreamsorcerer Aug 29, 2026
61176f4
Coverage
Dreamsorcerer Aug 29, 2026
3c027c3
Merge branch 'upload-tracker' of github.com:aio-libs/aiohttp into upl…
Dreamsorcerer Aug 29, 2026
d341158
Apply batched suggestions from code review
Dreamsorcerer Aug 29, 2026
205cc53
Rename 13427.deprecation.rst to 13579.deprecation.rst
Dreamsorcerer Aug 29, 2026
74236ef
Rename 13427.feature.rst to 13579.feature.rst
Dreamsorcerer Aug 29, 2026
b2fe939
Fix
Dreamsorcerer Aug 29, 2026
d9c5977
Fix
Dreamsorcerer Aug 29, 2026
adef9ca
Tests
Dreamsorcerer Aug 29, 2026
7b1751c
New exception type
Dreamsorcerer Aug 29, 2026
2de18b9
Fix
Dreamsorcerer Aug 29, 2026
c2ffdec
Fix
Dreamsorcerer Aug 29, 2026
90bdfdc
Fix docs
Dreamsorcerer Aug 29, 2026
1194c8b
Fix traces
Dreamsorcerer Aug 29, 2026
f848f01
Merge branch 'upload-tracker' of github.com:aio-libs/aiohttp into upl…
Dreamsorcerer Aug 29, 2026
aab57a2
Drop pointless extra method
Dreamsorcerer Aug 30, 2026
7974522
Apply batched suggestions from code review
Dreamsorcerer Aug 30, 2026
443be9e
Merge branch 'upload-tracker' of github.com:aio-libs/aiohttp into upl…
Dreamsorcerer Aug 30, 2026
eb2a617
Extra changelog
Dreamsorcerer Aug 30, 2026
68d85b8
Fix resend
Dreamsorcerer Aug 30, 2026
4848325
Fix timeout issue
Dreamsorcerer Aug 30, 2026
182974a
Apply batched suggestions from code review
Dreamsorcerer Aug 30, 2026
12a5932
Apply batched suggestions from code review
Dreamsorcerer Aug 30, 2026
da025c4
Coverage
Dreamsorcerer Aug 30, 2026
cca3e3d
Nitpick
Dreamsorcerer Aug 30, 2026
a4a1d9b
Apply suggestion from @Dreamsorcerer
Dreamsorcerer Aug 30, 2026
813c732
Coverage
Dreamsorcerer Aug 31, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions CHANGES.rst
Original file line number Diff line number Diff line change
Expand Up @@ -472,8 +472,7 @@ Features



- Added :attr:`~aiohttp.ClientResponse.output_size` and
:attr:`~aiohttp.ClientResponse.upload_complete` -- by :user:`Dreamsorcerer`.
- Added ``ClientResponse.output_size`` and ``ClientResponse.upload_complete`` -- by :user:`Dreamsorcerer`.


*Related issues and pull requests on GitHub:*
Expand Down
2 changes: 2 additions & 0 deletions CHANGES/13427.deprecation.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
Deprecated ``ClientResponse.output_size`` and ``ClientResponse.upload_complete``;
use ``aiohttp.UploadTracker`` instead -- by :user:`Dreamsorcerer`.
1 change: 1 addition & 0 deletions CHANGES/13427.feature.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Added :class:`aiohttp.UploadTracker` for observing a client request's upload progress -- by :user:`Dreamsorcerer`.
2 changes: 2 additions & 0 deletions aiohttp/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
TCPConnector,
TooManyRedirects,
UnixConnector,
UploadTracker,
WSMessageTypeError,
WSServerHandshakeError,
request,
Expand Down Expand Up @@ -156,6 +157,7 @@
"TCPConnector",
"TooManyRedirects",
"UnixConnector",
"UploadTracker",
"NamedPipeConnector",
"WSServerHandshakeError",
"request",
Expand Down
6 changes: 6 additions & 0 deletions aiohttp/abc.py
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,12 @@ class AbstractStreamWriter(ABC):
buffer_size: int = 0
output_size: int = 0
length: int | None = 0
# Called with each accepted body chunk's byte count (before any
# transport-level transformation such as compression or chunked
# framing). Assigned by the client request machinery for upload
# progress tracking; write()/write_eof() implementations should
# invoke it for every body chunk they accept.
on_body_write: Callable[[int], None] | None = None

@abstractmethod
async def write(
Expand Down
151 changes: 89 additions & 62 deletions aiohttp/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@
Fingerprint,
RequestInfo,
ResponseParams,
UploadTracker,
)
from .client_ws import (
DEFAULT_WS_CLIENT_TIMEOUT,
Expand Down Expand Up @@ -143,6 +144,7 @@
"ClientResponse",
"Fingerprint",
"RequestInfo",
"UploadTracker",
# connector
"BaseConnector",
"TCPConnector",
Expand Down Expand Up @@ -194,6 +196,7 @@
max_field_size: int | None
max_headers: int | None
middlewares: Sequence[ClientMiddlewareType] | None
upload_tracker: UploadTracker | None


class _WSConnectOptions(TypedDict, total=False):
Expand Down Expand Up @@ -494,6 +497,7 @@
max_field_size: int | None = None,
max_headers: int | None = None,
middlewares: Sequence[ClientMiddlewareType] | None = None,
upload_tracker: UploadTracker | None = None,
Comment thread
Dreamsorcerer marked this conversation as resolved.
) -> ClientResponse:
# NOTE: timeout clamps existing connect and read timeouts. We cannot
# set the default to None because we need to detect if the user wants
Expand Down Expand Up @@ -522,83 +526,100 @@
else:
data = payload.JsonPayload(json, dumps=self._json_serialize)

redirects = 0
history: list[ClientResponse] = []
version = self._version
params = params or {}
if upload_tracker is not None:
upload_tracker._bind()

# Merge with default headers and transform to CIMultiDict
headers = self._prepare_headers(headers)
real_timeout = (
self._timeout if timeout is sentinel or timeout is None else timeout
)
# timeout is cumulative for all request operations
# (request, redirects, responses, data consuming)
tm = TimeoutHandle(
self._loop,
real_timeout.total,
ceil_threshold=real_timeout.ceil_threshold,
)
handle: asyncio.TimerHandle | None = None

try:
url = self._build_url(str_or_url)
except ValueError as e:
raise InvalidUrlClientError(str_or_url) from e

assert self._connector is not None
if url.scheme not in self._connector.allowed_protocol_schema_set:
raise NonHttpUrlClientError(url)

skip_headers: Iterable[istr] | None
if skip_auto_headers is not None:
skip_headers = {
istr(i) for i in skip_auto_headers
} | self._skip_auto_headers
elif self._skip_auto_headers:
skip_headers = self._skip_auto_headers
else:
skip_headers = None
redirects = 0
history: list[ClientResponse] = []
version = self._version
params = params or {}

if proxy is None:
proxy = self._default_proxy
# Merge with default headers and transform to CIMultiDict
headers = self._prepare_headers(headers)

resolved_proxy_headers: CIMultiDict[str] | None
if proxy is None:
resolved_proxy_headers = None
else:
resolved_proxy_headers = self._prepare_headers(proxy_headers)
try:
proxy = URL(proxy)
url = self._build_url(str_or_url)
except ValueError as e:
raise InvalidURL(proxy) from e
raise InvalidUrlClientError(str_or_url) from e

assert self._connector is not None
if url.scheme not in self._connector.allowed_protocol_schema_set:
raise NonHttpUrlClientError(url)

skip_headers: Iterable[istr] | None
if skip_auto_headers is not None:
skip_headers = {
istr(i) for i in skip_auto_headers
} | self._skip_auto_headers
elif self._skip_auto_headers:
skip_headers = self._skip_auto_headers
else:
skip_headers = None

if timeout is sentinel or timeout is None:
real_timeout: ClientTimeout = self._timeout
else:
real_timeout = timeout
# timeout is cumulative for all request operations
# (request, redirects, responses, data consuming)
tm = TimeoutHandle(
self._loop, real_timeout.total, ceil_threshold=real_timeout.ceil_threshold
)
handle = tm.start()
if proxy is None:
proxy = self._default_proxy

resolved_proxy_headers: CIMultiDict[str] | None
if proxy is None:
resolved_proxy_headers = None
else:
resolved_proxy_headers = self._prepare_headers(proxy_headers)
try:
proxy = URL(proxy)
except ValueError as e:
raise InvalidURL(proxy) from e

if read_bufsize is None:
read_bufsize = self._read_bufsize
handle = tm.start()

if auto_decompress is None:
auto_decompress = self._auto_decompress
if read_bufsize is None:
read_bufsize = self._read_bufsize

if max_line_size is None:
max_line_size = self._max_line_size
if auto_decompress is None:
auto_decompress = self._auto_decompress

if max_field_size is None:
max_field_size = self._max_field_size
if max_line_size is None:
max_line_size = self._max_line_size

if max_headers is None:
max_headers = self._max_headers
if max_field_size is None:
max_field_size = self._max_field_size

traces = [
Trace(
self,
trace_config,
trace_config.trace_config_ctx(trace_request_ctx=trace_request_ctx),
)
for trace_config in self._trace_configs
]
if max_headers is None:
max_headers = self._max_headers

for trace in traces:
await trace.send_request_start(method, url.update_query(params), headers)
traces = [
Trace(
self,
trace_config,
trace_config.trace_config_ctx(trace_request_ctx=trace_request_ctx),
)
for trace_config in self._trace_configs
]

for trace in traces:
await trace.send_request_start(
method, url.update_query(params), headers
)
except BaseException as e:
tm.close()
if handle is not None:
handle.cancel()
handle = None
Comment thread
github-code-quality[bot] marked this conversation as resolved.
Fixed
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
if upload_tracker is not None:
upload_tracker._finalize(e)
raise

timer = tm.timer()
req: ClientRequest | None = None
Expand Down Expand Up @@ -709,6 +730,7 @@
traces=traces,
trust_env=self.trust_env,
)
req._upload_tracker = upload_tracker

# Apply middleware (if any) - per-request middleware overrides session middleware
effective_middlewares = (
Expand Down Expand Up @@ -883,6 +905,8 @@
await trace.send_request_end(
method, url.update_query(params), headers, resp
)
if upload_tracker is not None:
upload_tracker._finalize(None)
return resp

except BaseException as e:
Expand All @@ -892,6 +916,9 @@
handle.cancel()
handle = None

if upload_tracker is not None:
upload_tracker._finalize(e)

if req is not None and req._body is not None:
await req._body.close()

Expand Down
Loading
Loading