mirror of
https://github.com/element-hq/synapse.git
synced 2026-08-28 13:44:24 +00:00
Merge remote-tracking branch 'origin/develop' into jaywink/msc4262
This commit is contained in:
@@ -30,6 +30,7 @@ jobs:
|
||||
- name: Start postgres with a faked clock
|
||||
background: true
|
||||
id: postgres
|
||||
# Use faketime here for schema deltas that are wall-clock sensitive under Postgres
|
||||
run: |
|
||||
# Build a docker image with faketime
|
||||
mkdir /tmp/postgres-faketime
|
||||
@@ -58,8 +59,8 @@ jobs:
|
||||
with:
|
||||
fetch-depth: 0
|
||||
|
||||
- name: Install PostgreSQL client
|
||||
run: sudo apt-get -qq install postgresql-client
|
||||
- name: Install PostgreSQL client and faketime
|
||||
run: sudo apt-get -qq install postgresql-client faketime
|
||||
|
||||
- uses: matrix-org/setup-python-poetry@5bbf6603c5c930615ec8a29f1b5d7d258d905aa4 # v2.0.0
|
||||
with:
|
||||
@@ -77,8 +78,10 @@ jobs:
|
||||
PGHOST: localhost
|
||||
PGUSER: postgres
|
||||
PGPASSWORD: postgres
|
||||
# Use faketime here for schema deltas that are wall-clock sensitive under SQLite
|
||||
run: |
|
||||
poetry run python .ci/scripts/schema_diff.py \
|
||||
faketime -f "2001-05-25 12:42:42" \
|
||||
poetry run python .ci/scripts/schema_diff.py \
|
||||
--base origin/develop \
|
||||
> "${{ runner.temp }}/schema_diff.md"
|
||||
|
||||
|
||||
+41
@@ -1,3 +1,44 @@
|
||||
# Synapse 1.159.0rc1 (2026-08-11)
|
||||
|
||||
Administrators using the Debian/Ubuntu packages from `packages.matrix.org`, please check
|
||||
[the relevant section in the upgrade notes](https://github.com/element-hq/synapse/blob/release-v1.159/docs/upgrade.md#upgrading-to-v11590)
|
||||
as we have recently updated the expiry date on the repository's GPG signing key. The old version of the key will expire on `2027-03-15`.
|
||||
|
||||
## Features
|
||||
|
||||
- Add optional support for [MSC4429: Profile Updates for Legacy Sync](https://github.com/matrix-org/matrix-spec-proposals/pull/4429).
|
||||
Currently defaults to not enabled, and is limited to local users only for the sync results. ([\#19556](https://github.com/element-hq/synapse/issues/19556))
|
||||
|
||||
## Bugfixes
|
||||
|
||||
- Fix thumbnail generation failing for MPO images. Animations that cannot be decoded now fall back to a static thumbnail. ([\#20025](https://github.com/element-hq/synapse/issues/20025))
|
||||
- Fix the `quarantined_media` replication stream never being sent when the configured `quarantined_media_changes` stream writer is a worker. Introduced in v1.152.0. ([\#20085](https://github.com/element-hq/synapse/issues/20085))
|
||||
|
||||
## Updates to the Docker image
|
||||
|
||||
- Run with `PYTHONUNBUFFERED=1` to ensure that we can always see log output when things go wrong. ([\#20075](https://github.com/element-hq/synapse/issues/20075))
|
||||
|
||||
## Improved Documentation
|
||||
|
||||
- Correct the documentation for the `on_media_upload_limit_exceeded` module callback with regards to where it is called from. ([\#20018](https://github.com/element-hq/synapse/issues/20018))
|
||||
- Add upgrade notes to point out updated Debian package signing key. ([\#20066](https://github.com/element-hq/synapse/issues/20066))
|
||||
- Update stream cheatsheet docs to re-link `synapse/config/workers.py` which has more references. ([\#20086](https://github.com/element-hq/synapse/issues/20086))
|
||||
|
||||
## Internal Changes
|
||||
|
||||
- Fix tests that use `homeserver_to_use=GenericWorkerServer` not being able to be run standalone. ([\#20017](https://github.com/element-hq/synapse/issues/20017))
|
||||
- Fix `RemoteJoinHelper` test helper to handle room version "12" rooms. Contributed by @famedly @jason-famedly. ([\#20021](https://github.com/element-hq/synapse/issues/20021))
|
||||
- Fix release script announcement to link to correct release branch of changelog. ([\#20023](https://github.com/element-hq/synapse/issues/20023))
|
||||
- Dust off `make_full_schema` and add CI using it to show schema diffs. ([\#20027](https://github.com/element-hq/synapse/issues/20027))
|
||||
- Remove broken `DROP` statements for SQLite in `make_full_schema` script. ([\#20028](https://github.com/element-hq/synapse/issues/20028))
|
||||
- Document how to capture a JSON snapshot of a Grafana dashboard to aid in debugging. ([\#20048](https://github.com/element-hq/synapse/issues/20048))
|
||||
- Routinely purge old cancelled tasks from the database. ([\#20068](https://github.com/element-hq/synapse/issues/20068))
|
||||
- Introduce an `RdataSafeValue` type and correct some minor type annotation mistakes. ([\#20071](https://github.com/element-hq/synapse/issues/20071))
|
||||
- Set `idle_in_transaction_session_timeout` (default 30 minutes) on new PostgreSQL connections, so that wedged connections don't hold locks or block vacuum indefinitely. ([\#20077](https://github.com/element-hq/synapse/issues/20077))
|
||||
|
||||
|
||||
|
||||
|
||||
# Synapse 1.158.0 (2026-08-04)
|
||||
|
||||
## Deprecations and Removals
|
||||
|
||||
@@ -1,2 +0,0 @@
|
||||
Add optional support for [MSC4429: Profile Updates for Legacy Sync](https://github.com/matrix-org/matrix-spec-proposals/pull/4429).
|
||||
Currently defaults to not enabled, and is limited to local users only for the sync results.
|
||||
@@ -0,0 +1 @@
|
||||
Document lighttpd reverse proxy configuration example from matrix.jaxlug.ngo, a contribution from the JaxLUG, the Jacksonville Linux Users Group Inc..
|
||||
@@ -1 +0,0 @@
|
||||
Fix tests that use `homeserver_to_use=GenericWorkerServer` not being able to be run standalone.
|
||||
@@ -1 +0,0 @@
|
||||
Correct the documentation for the `on_media_upload_limit_exceeded` module callback with regards to where it is called from.
|
||||
@@ -1 +0,0 @@
|
||||
Fix `RemoteJoinHelper` test helper to handle room version "12" rooms. Contributed by @famedly @jason-famedly.
|
||||
@@ -1 +0,0 @@
|
||||
Fix release script announcement to link to correct release branch of changelog.
|
||||
@@ -1 +0,0 @@
|
||||
Dust off `make_full_schema` and add CI using it to show schema diffs.
|
||||
@@ -1 +0,0 @@
|
||||
Remove broken `DROP` statements for SQLite in `make_full_schema` script.
|
||||
@@ -1 +0,0 @@
|
||||
Add package build targets for Ubuntu 26.04 'Resolute Raccoon'.
|
||||
@@ -1 +0,0 @@
|
||||
Remove package build targets for Ubuntu 25.10 'Questing Quokka' (end-of-life 2026-07-01).
|
||||
@@ -1 +0,0 @@
|
||||
Document how to capture a JSON snapshot of a Grafana dashboard to aid in debugging.
|
||||
@@ -1 +0,0 @@
|
||||
Add upgrade notes to point out updated Debian package signing key.
|
||||
@@ -0,0 +1 @@
|
||||
Allow specifying multiple `action_name` and `status` query parameters when listing scheduled tasks via the admin API.
|
||||
@@ -1 +0,0 @@
|
||||
Routinely purge old cancelled tasks from the database.
|
||||
@@ -1 +0,0 @@
|
||||
Introduce an `RdataSafeValue` type and correct some minor type annotation mistakes.
|
||||
@@ -1 +0,0 @@
|
||||
Run with `PYTHONUNBUFFERED=1` so to ensure can always see log output when things go wrong.
|
||||
@@ -1 +0,0 @@
|
||||
Set `idle_in_transaction_session_timeout` (default 30 minutes) on new PostgreSQL connections, so that wedged connections don't hold locks or block vacuum indefinitely.
|
||||
@@ -1 +0,0 @@
|
||||
Fix the `quarantined_media` replication stream never being sent when the configured `quarantined_media_changes` stream writer is a worker. Introduced in v1.152.0.
|
||||
@@ -0,0 +1 @@
|
||||
Fix the documentation on the `federation_domain_whitelist` config option.
|
||||
@@ -0,0 +1 @@
|
||||
Fix a bug where presence updates could stop being sent to clients (the presence stream position becoming stuck) if a `/sync` request was cancelled while a presence write was allocating a stream ID. Contributed by @FrenchGithubUser @Famedly.
|
||||
@@ -0,0 +1 @@
|
||||
Update release script to check more often for actions being completed so you don't have to wait around as much.
|
||||
@@ -0,0 +1 @@
|
||||
Fix the schema diff CI not using `faketime` for SQLite.
|
||||
Vendored
+6
@@ -1,3 +1,9 @@
|
||||
matrix-synapse-py3 (1.159.0~rc1) stable; urgency=medium
|
||||
|
||||
* New synapse release 1.159.0rc1.
|
||||
|
||||
-- Synapse Packaging team <packages@matrix.org> Tue, 11 Aug 2026 16:25:42 +0000
|
||||
|
||||
matrix-synapse-py3 (1.158.0) stable; urgency=medium
|
||||
|
||||
* New synapse release 1.158.0.
|
||||
|
||||
@@ -31,8 +31,10 @@ It returns a JSON body like the following:
|
||||
**Query parameters:**
|
||||
|
||||
* `action_name`: string - Is optional. Returns only the scheduled tasks with the given action name.
|
||||
May be given multiple times to return tasks matching any of the given action names.
|
||||
* `resource_id`: string - Is optional. Returns only the scheduled tasks with the given resource id.
|
||||
* `status`: string - Is optional. Returns only the scheduled tasks matching the given status, one of
|
||||
* `status`: string - Is optional. Returns only the scheduled tasks matching the given status.
|
||||
May be given multiple times to return tasks matching any of the given statuses. The status must be one of
|
||||
- "scheduled" - Task is scheduled but not active
|
||||
- "active" - Task is active and probably running, and if not will be run on next scheduler loop run
|
||||
- "complete" - Task has completed successfully
|
||||
|
||||
@@ -162,7 +162,7 @@ necessary registration and event handling.
|
||||
- Update `synapse/_scripts/synapse_port_db.py` so it knows about your new `SEQUENCE`: [add a new `_setup_sequence(...)`](https://github.com/element-hq/synapse/blob/35b55e962aa0bed3b2da5a3c12e3783ddf7604ca/synapse/_scripts/synapse_port_db.py#L883C24-L888)
|
||||
- [create a stream class and stream row class](https://github.com/element-hq/synapse/blob/4367fb2d078c52959aeca0fe6874539c53e8360d/synapse/replication/tcp/streams/_base.py#L728)
|
||||
- will need an [ID generator](https://github.com/element-hq/synapse/blob/4367fb2d078c52959aeca0fe6874539c53e8360d/synapse/storage/databases/main/thread_subscriptions.py#L75)
|
||||
- may need [writer configuration](https://github.com/element-hq/synapse/blob/4367fb2d078c52959aeca0fe6874539c53e8360d/synapse/config/workers.py#L177), if there isn't already an obvious source of configuration for which workers should be designated as writers to your new stream.
|
||||
- may need [writer configuration](https://github.com/element-hq/synapse/blob/62a4bc46203880dd5034483b0e84156d03a3a8c6/synapse/config/workers.py#L184-L187), if there isn't already an obvious source of configuration for which workers should be designated as writers to your new stream.
|
||||
- if adding new writer configuration, add Docker-worker configuration, which lets us configure the writer worker in Complement tests: [[1]](https://github.com/element-hq/synapse/blob/4367fb2d078c52959aeca0fe6874539c53e8360d/docker/configure_workers_and_start.py#L331), [[2]](https://github.com/element-hq/synapse/blob/4367fb2d078c52959aeca0fe6874539c53e8360d/docker/configure_workers_and_start.py#L440)
|
||||
- Ensure that it's been correctly added to `synapse/replication/tcp/handler.py` and it's `streams_to_replicate` attribute to ensure that changes are actually replicated.
|
||||
- most of the time, you will likely introduce a new datastore class for the concept represented by the new stream, unless there is already an obvious datastore that covers it.
|
||||
|
||||
+86
-2
@@ -4,8 +4,10 @@ It is recommended to put a reverse proxy such as
|
||||
[nginx](https://nginx.org/en/docs/http/ngx_http_proxy_module.html),
|
||||
[Apache](https://httpd.apache.org/docs/current/mod/mod_proxy_http.html),
|
||||
[Caddy](https://caddyserver.com/docs/quick-starts/reverse-proxy),
|
||||
[HAProxy](https://www.haproxy.org/) or
|
||||
[relayd](https://man.openbsd.org/relayd.8) in front of Synapse.
|
||||
[HAProxy](https://www.haproxy.org/),
|
||||
[relayd](https://man.openbsd.org/relayd.8) or
|
||||
[lighttpd](https://www.lighttpd.net/)
|
||||
in front of Synapse.
|
||||
This has the advantage of being able to expose the default HTTPS port (443) to Matrix
|
||||
clients without requiring Synapse to bind to a privileged port (port numbers less than
|
||||
1024), avoiding the need for `CAP_NET_BIND_SERVICE` or running as root.
|
||||
@@ -312,6 +314,88 @@ relay "matrix_federation" {
|
||||
}
|
||||
```
|
||||
|
||||
### lighttpd
|
||||
```conf
|
||||
server.modules = (
|
||||
"mod_rewrite",
|
||||
"mod_redirect",
|
||||
"mod_access",
|
||||
"mod_setenv",
|
||||
"mod_openssl",
|
||||
"mod_proxy",
|
||||
"mod_accesslog"
|
||||
)
|
||||
|
||||
server.username = "lighttpd"
|
||||
server.groupname = "lighttpd"
|
||||
|
||||
# We set this to "disable" and use IPv6 `[::]` explicitly below,
|
||||
# in order to listen on all incoming IPv6 addresses.
|
||||
#
|
||||
# If you only want to listen on specific IPv6 addresses, set this
|
||||
# to "enable" and specify said addresses below.
|
||||
server.use-ipv6 = "disable"
|
||||
|
||||
ssl.pemfile = "/etc/lighttpd/cert+privkey.pem"
|
||||
ssl.ca-file = "/etc/lighttpd/fullchain.pem"
|
||||
|
||||
# redirect HTTP traffic to HTTPS, same for IPv6 below
|
||||
$SERVER["socket"] == "0.0.0.0:80" {
|
||||
url.redirect = (
|
||||
"" => "https://${url.authority.noport}${url.path}${qsa}"
|
||||
)
|
||||
}
|
||||
$SERVER["socket"] == "0.0.0.0:443" { ssl.engine = "enable" }
|
||||
$SERVER["socket"] == "0.0.0.0:8448" { ssl.engine = "enable" }
|
||||
$SERVER["socket"] == "[::]:80" {
|
||||
url.redirect = (
|
||||
"" => "https://${url.authority.noport}${url.path}${qsa}"
|
||||
)
|
||||
}
|
||||
$SERVER["socket"] == "[::]:443" { ssl.engine = "enable" }
|
||||
$SERVER["socket"] == "[::]:8448" { ssl.engine = "enable" }
|
||||
|
||||
|
||||
|
||||
# both lighttpd and synapse need permissions for socket r/w
|
||||
$HTTP["url"] =~ "(/_matrix|_synapse/admin|/_synapse/client)" {
|
||||
proxy.balance = "hash"
|
||||
proxy.server = (
|
||||
"" => (
|
||||
"backend-socket" => (
|
||||
"host" => "/var/lib/synapse/main_public.sock",
|
||||
"port" => 0
|
||||
)
|
||||
)
|
||||
)
|
||||
proxy.forwarded = (
|
||||
"for" => 1,
|
||||
"proto" => 1,
|
||||
"host" => 1,
|
||||
)
|
||||
}
|
||||
# protect admin access IPv6 ULA only
|
||||
$HTTP["remoteip"] !="fd00::/8" {
|
||||
$HTTP["url"] =~ "^/_synapse/admin/" {
|
||||
url.access-deny = ( "" )
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
[Delegation](delegate.md) example:
|
||||
```conf
|
||||
url.rewrite-once = (
|
||||
"^/\.well-known/matrix/client$" => "/.well-known/matrix/client.json",
|
||||
"^/\.well-known/matrix/server$" => "/.well-known/matrix/server.json"
|
||||
)
|
||||
|
||||
# This condition intentionally matches the post-rewrite URLs.
|
||||
$HTTP["url"] =~ "^/\.well-known/matrix/(client|server)\.json$" {
|
||||
mimetype.assign = ( ".json" => "application/json" )
|
||||
setenv.set-response-header = ( "Access-Control-Allow-Origin" => "*" )
|
||||
}
|
||||
```
|
||||
|
||||
|
||||
## Health check endpoint
|
||||
|
||||
|
||||
@@ -1287,11 +1287,10 @@ Options related to federation.
|
||||
---
|
||||
### `federation_domain_whitelist`
|
||||
|
||||
*(array)* Restrict federation to the given whitelist of domains. N.B. we recommend also firewalling your federation listener to limit inbound federation traffic as early as possible, rather than relying purely on this application-layer restriction. If not specified, the default is to whitelist everything.
|
||||
|
||||
Note: this does not stop a server from joining rooms that servers not on the whitelist are in. As such, this option is really only useful to establish a "private federation", where a group of servers all whitelist each other and have the same whitelist.
|
||||
|
||||
Defaults to `[]`.
|
||||
*(array)* Restrict federation to the given whitelist of domains. N.B. we recommend also firewalling your federation listener to limit inbound federation traffic as early as possible, rather than relying purely on this application-layer restriction.
|
||||
If specified as an empty list (`[]`), federation will be denied with all servers. Specifying an empty list (`[]`) here is the recommended way of disabling federation.
|
||||
If not specified, the default is to allow federation with all servers.
|
||||
Note: this does not stop a server from joining rooms that servers not on the whitelist are in. As such, this option is really only useful to establish a "private federation", where a group of servers all whitelist each other and have the same whitelist. There is no default for this option.
|
||||
|
||||
Example configuration:
|
||||
```yaml
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "matrix-synapse"
|
||||
version = "1.158.0"
|
||||
version = "1.159.0rc1"
|
||||
description = "Homeserver for the Matrix decentralised comms protocol"
|
||||
readme = "README.rst"
|
||||
authors = [
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
$schema: https://element-hq.github.io/synapse/latest/schema/v1/meta.schema.json
|
||||
$id: https://element-hq.github.io/synapse/schema/synapse/v1.158/synapse-config.schema.json
|
||||
$id: https://element-hq.github.io/synapse/schema/synapse/v1.159/synapse-config.schema.json
|
||||
type: object
|
||||
properties:
|
||||
modules:
|
||||
@@ -280,10 +280,10 @@ properties:
|
||||
description: >-
|
||||
Use this option to include updates of other users' profiles in sync responses,
|
||||
for users who share rooms.
|
||||
|
||||
|
||||
Requires an [MSC4429](https://github.com/matrix-org/matrix-spec-proposals/pull/4429)
|
||||
compatible client, and is currently limited to legacy sync and local users only.
|
||||
|
||||
|
||||
This feature is under development and should be used with caution on busy servers or
|
||||
servers which depend on `limit_profile_requests_to_users_who_share_rooms` for ensuring
|
||||
profile information doesn't leak across rooms.
|
||||
@@ -1578,9 +1578,12 @@ properties:
|
||||
Restrict federation to the given whitelist of domains. N.B. we recommend
|
||||
also firewalling your federation listener to limit inbound federation
|
||||
traffic as early as possible, rather than relying purely on this
|
||||
application-layer restriction. If not specified, the default is to
|
||||
whitelist everything.
|
||||
application-layer restriction.
|
||||
|
||||
If specified as an empty list (`[]`), federation will be denied with all servers.
|
||||
Specifying an empty list (`[]`) here is the recommended way of disabling federation.
|
||||
|
||||
If not specified, the default is to allow federation with all servers.
|
||||
|
||||
Note: this does not stop a server from joining rooms that servers not on
|
||||
the whitelist are in. As such, this option is really only useful to
|
||||
@@ -1588,7 +1591,6 @@ properties:
|
||||
each other and have the same whitelist.
|
||||
items:
|
||||
type: string
|
||||
default: []
|
||||
examples:
|
||||
- - lon.example.com
|
||||
- nyc.example.com
|
||||
|
||||
@@ -600,9 +600,15 @@ def _wait_for_actions(gh_token: str | None) -> None:
|
||||
headers["authorization"] = f"token {gh_token}"
|
||||
req = urllib.request.Request(url, headers=headers)
|
||||
|
||||
# Initially, wait 10 minutes as we know the CI typically takes 15m+ anyway (no need
|
||||
# to check over and over when we know it won't be finished yet)
|
||||
time.sleep(10 * 60)
|
||||
while True:
|
||||
time.sleep(5 * 60)
|
||||
# Then check once every minute. Short enough to not have to wait around too long
|
||||
# while not spamming the GitHub API and running into the unauthenticated API
|
||||
# request rate limit (60 requests per hour so 1 request/minute perfectly aligns
|
||||
# to not run into any problems)
|
||||
time.sleep(1 * 60)
|
||||
response = urllib.request.urlopen(req)
|
||||
resp = json.loads(response.read())
|
||||
|
||||
|
||||
@@ -79,6 +79,10 @@ class Thumbnailer:
|
||||
# format in this list becomes part of our trusted computing base.
|
||||
PILLOW_FORMATS = ("jpeg", "png", "webp", "gif")
|
||||
|
||||
# Pillow reports MPO (a JPEG holding a stereo pair) as multi-frame, so
|
||||
# frame count alone doesn't tell us whether something is an animation.
|
||||
ANIMATED_FORMATS = frozenset({"GIF", "PNG", "WEBP"})
|
||||
|
||||
@staticmethod
|
||||
def set_limits(max_image_pixels: int) -> None:
|
||||
Image.MAX_IMAGE_PIXELS = max_image_pixels
|
||||
@@ -86,6 +90,12 @@ class Thumbnailer:
|
||||
def __init__(self, input_path: str):
|
||||
# Have we closed the image?
|
||||
self._closed = False
|
||||
# Whether attempting to thumbnail the image failed for some reason. The
|
||||
# thumbnailing code should fallback to treating the image as static in
|
||||
# this case.
|
||||
#
|
||||
# Cached so a broken animation isn't re-decoded for every thumbnail size.
|
||||
self._animation_broken = False
|
||||
|
||||
try:
|
||||
self.image = Image.open(input_path, formats=self.PILLOW_FORMATS)
|
||||
@@ -166,30 +176,52 @@ class Thumbnailer:
|
||||
|
||||
@property
|
||||
def is_animated(self) -> bool:
|
||||
if self._animation_broken:
|
||||
return False
|
||||
if self.image.format not in self.ANIMATED_FORMATS:
|
||||
return False
|
||||
return getattr(self.image, "is_animated", False)
|
||||
|
||||
def _encode_animated_from(
|
||||
self, transform: Callable[[Image.Image], Image.Image]
|
||||
) -> BytesIO:
|
||||
) -> BytesIO | None:
|
||||
"""Apply `transform` to every frame of the source image and encode the
|
||||
result as an animated thumbnail."""
|
||||
result as an animated thumbnail, or None if the source could not be
|
||||
decoded as an animation and the caller should fall back to a static one.
|
||||
"""
|
||||
frames = []
|
||||
durations = []
|
||||
loop = self.image.info.get("loop", 0)
|
||||
for frame in ImageSequence.Iterator(self.image):
|
||||
# Copy the frame to avoid referencing the original image memory.
|
||||
f = frame.copy()
|
||||
if f.mode != "RGBA":
|
||||
f = f.convert("RGBA")
|
||||
frames.append(transform(f))
|
||||
# A duration of 0 is valid (interpretation is implementation-defined,
|
||||
# see RFC 9649 section 2.7.1.1), so only fall back when unset. The
|
||||
# 100ms default matches libwebp's animation tools.
|
||||
duration = frame.info.get("duration")
|
||||
if duration is None:
|
||||
duration = self.image.info.get("duration", 100)
|
||||
durations.append(duration)
|
||||
return self._encode_animated(frames, durations, loop)
|
||||
try:
|
||||
for frame in ImageSequence.Iterator(self.image):
|
||||
# Copy the frame to avoid referencing the original image memory.
|
||||
f = frame.copy()
|
||||
if f.mode != "RGBA":
|
||||
f = f.convert("RGBA")
|
||||
frames.append(transform(f))
|
||||
# A duration of 0 is valid (interpretation is implementation-defined,
|
||||
# see RFC 9649 section 2.7.1.1), so only fall back when unset. The
|
||||
# 100ms default matches libwebp's animation tools.
|
||||
duration = frame.info.get("duration")
|
||||
if duration is None:
|
||||
duration = self.image.info.get("duration", 100)
|
||||
durations.append(duration)
|
||||
return self._encode_animated(frames, durations, loop)
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"Failed to generate an animated thumbnail, falling back to a "
|
||||
"static one: %s",
|
||||
e,
|
||||
)
|
||||
self._animation_broken = True
|
||||
|
||||
try:
|
||||
# Leave the source on its first frame for the static fallback.
|
||||
self.image.seek(0)
|
||||
except Exception as e:
|
||||
logger.warning("Failed to rewind image to its first frame: %s", e)
|
||||
|
||||
return None
|
||||
|
||||
@trace
|
||||
def scale(
|
||||
@@ -204,9 +236,11 @@ class Thumbnailer:
|
||||
The bytes of the encoded image ready to be written to disk
|
||||
"""
|
||||
if animated and self.is_animated:
|
||||
return self._encode_animated_from(
|
||||
output = self._encode_animated_from(
|
||||
lambda f: self._resize_image(f, width, height)
|
||||
)
|
||||
if output is not None:
|
||||
return output
|
||||
|
||||
with self._resize_image(self.image, width, height) as scaled:
|
||||
return self._encode_image(scaled, output_type)
|
||||
@@ -242,9 +276,11 @@ class Thumbnailer:
|
||||
crop = (crop_left, 0, crop_right, height)
|
||||
|
||||
if animated and self.is_animated:
|
||||
return self._encode_animated_from(
|
||||
output = self._encode_animated_from(
|
||||
lambda f: self._resize_image(f, scaled_width, scaled_height).crop(crop)
|
||||
)
|
||||
if output is not None:
|
||||
return output
|
||||
|
||||
with self._resize_image(
|
||||
self.image, scaled_width, scaled_height
|
||||
|
||||
@@ -15,7 +15,12 @@
|
||||
#
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from synapse.http.servlet import RestServlet, parse_integer, parse_string
|
||||
from synapse.http.servlet import (
|
||||
RestServlet,
|
||||
parse_integer,
|
||||
parse_string,
|
||||
parse_strings_from_args,
|
||||
)
|
||||
from synapse.http.site import SynapseRequest
|
||||
from synapse.rest.admin import admin_patterns, assert_requester_is_admin
|
||||
from synapse.types import JsonDict, TaskStatus
|
||||
@@ -38,19 +43,33 @@ class ScheduledTasksRestServlet(RestServlet):
|
||||
async def on_GET(self, request: SynapseRequest) -> tuple[int, JsonDict]:
|
||||
await assert_requester_is_admin(self._auth, request)
|
||||
|
||||
# twisted.web.server.Request.args is incorrectly defined as Any | None
|
||||
args: dict[bytes, list[bytes]] = request.args # type: ignore
|
||||
|
||||
# extract query params
|
||||
action_name = parse_string(request, "action_name")
|
||||
actions = parse_strings_from_args(args, "action_name")
|
||||
resource_id = parse_string(request, "resource_id")
|
||||
status = parse_string(request, "status")
|
||||
status_strings = parse_strings_from_args(
|
||||
args,
|
||||
"status",
|
||||
allowed_values=[status.value for status in TaskStatus],
|
||||
)
|
||||
# This parameter was historically called `job_status`, while the Admin API docs
|
||||
# defined it as `status`. We now support both, as `status` is generally
|
||||
# a nicer name. A v2 of this endpoint should keep only `status`.
|
||||
if status is None:
|
||||
status = parse_string(request, "job_status")
|
||||
if status_strings is None:
|
||||
status_strings = parse_strings_from_args(
|
||||
args,
|
||||
"job_status",
|
||||
allowed_values=[status.value for status in TaskStatus],
|
||||
)
|
||||
max_timestamp = parse_integer(request, "max_timestamp")
|
||||
|
||||
actions = [action_name] if action_name else None
|
||||
statuses = [TaskStatus(status)] if status else None
|
||||
statuses = (
|
||||
[TaskStatus(status) for status in status_strings]
|
||||
if status_strings
|
||||
else None
|
||||
)
|
||||
|
||||
tasks = await self._store.get_scheduled_tasks(
|
||||
actions=actions,
|
||||
|
||||
@@ -906,14 +906,44 @@ class _MultiWriterCtxManager:
|
||||
stream_ids: list[int] = attr.Factory(list)
|
||||
|
||||
async def __aenter__(self) -> int | list[int]:
|
||||
def _load(txn: LoggingTransaction) -> list[int]:
|
||||
ids = self.id_gen._load_next_mult_id_txn(txn, self.multiple_ids or 1)
|
||||
# Record the allocated IDs on the context manager as a side effect
|
||||
# (rather than only via the return value), so that if this coroutine
|
||||
# is cancelled after the transaction has committed we still know
|
||||
# which IDs to release below.
|
||||
self.stream_ids = ids
|
||||
return ids
|
||||
|
||||
# It's safe to run this in autocommit mode as fetching values from a
|
||||
# sequence ignores transaction semantics anyway.
|
||||
self.stream_ids = await self.id_gen._db.runInteraction(
|
||||
"_load_next_mult_id",
|
||||
self.id_gen._load_next_mult_id_txn,
|
||||
self.multiple_ids or 1,
|
||||
db_autocommit=True,
|
||||
)
|
||||
try:
|
||||
await self.id_gen._db.runInteraction(
|
||||
"_load_next_mult_id",
|
||||
_load,
|
||||
db_autocommit=True,
|
||||
)
|
||||
except BaseException:
|
||||
# We catch `BaseException` rather than `Exception`,
|
||||
# because request cancellation surfaces here as exceptions that are
|
||||
# not `Exception` subclasses: `asyncio.CancelledError`
|
||||
# and `GeneratorExit` (raised when a paused coroutine is garbage
|
||||
# collected).
|
||||
#
|
||||
# If we're interrupted (e.g. the enclosing request was cancelled)
|
||||
# after the transaction allocated the IDs but before we returned,
|
||||
# then `__aexit__` will never run, because Python only invokes it
|
||||
# once `__aenter__` has returned. The allocated IDs would then be
|
||||
# leaked into `_unfinished_ids` forever, permanently pinning the
|
||||
# persisted stream position and, e.g., wedging presence.
|
||||
#
|
||||
# So mark them as finished here to unblock the position. This mirrors
|
||||
# what `__aexit__` does on the failure path (marking the IDs finished
|
||||
# and notifying replication, but not persisting a new position).
|
||||
if self.stream_ids:
|
||||
self.id_gen._mark_ids_as_finished(self.stream_ids)
|
||||
self.notifier.notify_replication()
|
||||
raise
|
||||
|
||||
if self.multiple_ids is None:
|
||||
return self.stream_ids[0] * self.id_gen._return_factor
|
||||
|
||||
@@ -1428,6 +1428,22 @@ def _make_animated_gif() -> bytes:
|
||||
return out.getvalue()
|
||||
|
||||
|
||||
def _make_mpo() -> bytes:
|
||||
"""Build a two-image MPO: a JPEG holding a stereo pair, not an animation."""
|
||||
frames = [Image.new("RGB", (64, 64), color) for color in ((255, 0, 0), (0, 0, 255))]
|
||||
out = BytesIO()
|
||||
frames[0].save(out, format="MPO", save_all=True, append_images=frames[1:])
|
||||
return out.getvalue()
|
||||
|
||||
|
||||
def _make_stale_mpo() -> bytes:
|
||||
"""Build an MPO whose trailing image is stripped but still advertised."""
|
||||
data = _make_mpo()
|
||||
with Image.open(BytesIO(data)) as image:
|
||||
primary_size = image.mpinfo[0xB002][0]["Size"] # type: ignore[attr-defined]
|
||||
return data[:primary_size]
|
||||
|
||||
|
||||
class ThumbnailerAnimatedTestCase(unittest.TestCase):
|
||||
"""Tests that the thumbnailer only animates when explicitly asked to."""
|
||||
|
||||
@@ -1444,6 +1460,24 @@ class ThumbnailerAnimatedTestCase(unittest.TestCase):
|
||||
with open(self.png_path, "wb") as f:
|
||||
f.write(SMALL_PNG)
|
||||
|
||||
self.mpo_path = os.path.join(self.tempdir, "stereo.jpg")
|
||||
with open(self.mpo_path, "wb") as f:
|
||||
f.write(_make_mpo())
|
||||
|
||||
self.stale_mpo_path = os.path.join(self.tempdir, "stale.jpg")
|
||||
with open(self.stale_mpo_path, "wb") as f:
|
||||
f.write(_make_stale_mpo())
|
||||
|
||||
def assert_is_first_frame(self, output: BytesIO) -> None:
|
||||
"""Raises an `AssertionError` unless the given image is red."""
|
||||
pixel = Image.open(output).convert("RGB").getpixel((16, 16))
|
||||
assert isinstance(pixel, tuple)
|
||||
red, green, blue = pixel
|
||||
# The first frame of every source here is red.
|
||||
# WebP is lossy, so allow *some* green/blue to be present.
|
||||
self.assertGreater(red, 200)
|
||||
self.assertLess(max(green, blue), 50)
|
||||
|
||||
def test_scale_static_by_default(self) -> None:
|
||||
"""An animated source produces a static thumbnail unless animated=True."""
|
||||
with Thumbnailer(self.gif_path) as thumbnailer:
|
||||
@@ -1475,3 +1509,83 @@ class ThumbnailerAnimatedTestCase(unittest.TestCase):
|
||||
out = thumbnailer.scale(1, 1, ANIMATED_THUMBNAIL_TYPE, animated=True)
|
||||
result = Image.open(out)
|
||||
self.assertFalse(getattr(result, "is_animated", False))
|
||||
|
||||
@parameterized.expand([("GIF", "gif"), ("PNG", "apng"), ("WEBP", "webp")])
|
||||
def test_animated_formats(self, fmt: str, ext: str) -> None:
|
||||
"""Every animated format we accept produces an animated thumbnail."""
|
||||
frames = [
|
||||
Image.new("RGBA", (64, 64), color)
|
||||
for color in ((255, 0, 0, 255), (0, 0, 255, 255))
|
||||
]
|
||||
out = BytesIO()
|
||||
frames[0].save(
|
||||
out,
|
||||
format=fmt,
|
||||
save_all=True,
|
||||
append_images=frames[1:],
|
||||
duration=100,
|
||||
loop=0,
|
||||
)
|
||||
path = os.path.join(self.tempdir, f"animated.{ext}")
|
||||
with open(path, "wb") as f:
|
||||
f.write(out.getvalue())
|
||||
|
||||
with Thumbnailer(path) as thumbnailer:
|
||||
self.assertTrue(thumbnailer.is_animated)
|
||||
thumbnail = thumbnailer.scale(
|
||||
32, 32, ANIMATED_THUMBNAIL_TYPE, animated=True
|
||||
)
|
||||
result = Image.open(thumbnail)
|
||||
self.assertEqual(result.format, "WEBP")
|
||||
self.assertTrue(getattr(result, "is_animated", False))
|
||||
self.assertEqual(getattr(result, "n_frames", 1), 2)
|
||||
|
||||
@parameterized.expand(["scale", "crop"])
|
||||
def test_mpo_is_not_animated(self, method: str) -> None:
|
||||
"""An MPO packs several stills into one JPEG; not an animation."""
|
||||
with Thumbnailer(self.mpo_path) as thumbnailer:
|
||||
self.assertFalse(thumbnailer.is_animated)
|
||||
out = getattr(thumbnailer, method)(
|
||||
32, 32, ANIMATED_THUMBNAIL_TYPE, animated=True
|
||||
)
|
||||
self.assertFalse(getattr(Image.open(out), "is_animated", False))
|
||||
self.assert_is_first_frame(out)
|
||||
|
||||
@parameterized.expand(["scale", "crop"])
|
||||
def test_stale_mpo_index_does_not_raise(self, method: str) -> None:
|
||||
"""An MPO advertising frames that are not in the file still thumbnails.
|
||||
|
||||
Regression test for https://github.com/element-hq/synapse/issues/20024.
|
||||
"""
|
||||
with Thumbnailer(self.stale_mpo_path) as thumbnailer:
|
||||
out = getattr(thumbnailer, method)(
|
||||
32, 32, ANIMATED_THUMBNAIL_TYPE, animated=True
|
||||
)
|
||||
self.assertEqual(Image.open(out).format, "WEBP")
|
||||
self.assert_is_first_frame(out)
|
||||
|
||||
def test_fallback_thumbnails_the_first_frame(self) -> None:
|
||||
"""Failing after the frames are read leaves the source parked on the
|
||||
last one, so the fallback has to rewind."""
|
||||
with patch.object(Thumbnailer, "_encode_animated", side_effect=ValueError):
|
||||
with Thumbnailer(self.gif_path) as thumbnailer:
|
||||
out = thumbnailer.scale(32, 32, ANIMATED_THUMBNAIL_TYPE, animated=True)
|
||||
self.assert_is_first_frame(out)
|
||||
|
||||
@parameterized.expand(["scale", "crop"])
|
||||
def test_undecodable_animation_falls_back_to_static(self, method: str) -> None:
|
||||
"""If the frames can't be decoded we still serve a static thumbnail of
|
||||
the first frame rather than failing the request."""
|
||||
# Force the stale MPO down the animated path so decoding it fails.
|
||||
with patch.object(Thumbnailer, "ANIMATED_FORMATS", frozenset({"MPO"})):
|
||||
with Thumbnailer(self.stale_mpo_path) as thumbnailer:
|
||||
self.assertTrue(thumbnailer.is_animated)
|
||||
out = getattr(thumbnailer, method)(
|
||||
32, 32, ANIMATED_THUMBNAIL_TYPE, animated=True
|
||||
)
|
||||
# A broken source isn't retried for every other thumbnail size.
|
||||
self.assertFalse(thumbnailer.is_animated)
|
||||
|
||||
self.assertEqual(Image.open(out).format, "WEBP")
|
||||
self.assertFalse(getattr(Image.open(out), "is_animated", False))
|
||||
self.assert_is_first_frame(out)
|
||||
|
||||
@@ -190,3 +190,56 @@ class ScheduledTasksAdminApiTestCase(unittest.HomeserverTestCase):
|
||||
# only the task with the matching resource id should have been returned
|
||||
self.assertEqual(len(found_tasks), 1)
|
||||
self.assertEqual(found_tasks[0]["resource_id"], "failed_task")
|
||||
|
||||
def test_filtering_scheduled_tasks_multiple_values(self) -> None:
|
||||
"""
|
||||
Test that the `action_name` and `status` filters can be given multiple
|
||||
times, returning tasks matching any of the given values.
|
||||
"""
|
||||
# filter via multiple statuses
|
||||
channel = self.make_request(
|
||||
"GET",
|
||||
"/_synapse/admin/v1/scheduled_tasks?status=active&status=failed",
|
||||
content={},
|
||||
access_token=self.admin_user_tok,
|
||||
)
|
||||
self.assertEqual(200, channel.code, msg=channel.json_body)
|
||||
found_tasks = self.check_scheduled_tasks_response(
|
||||
channel.json_body["scheduled_tasks"]
|
||||
)
|
||||
|
||||
# the active and failed tasks should have been returned
|
||||
self.assertEqual(len(found_tasks), 2)
|
||||
self.assertEqual({task["status"] for task in found_tasks}, {"active", "failed"})
|
||||
|
||||
# filter via multiple action names
|
||||
channel = self.make_request(
|
||||
"GET",
|
||||
"/_synapse/admin/v1/scheduled_tasks?action_name=test_task&action_name=finished_test_task",
|
||||
content={},
|
||||
access_token=self.admin_user_tok,
|
||||
)
|
||||
self.assertEqual(200, channel.code, msg=channel.json_body)
|
||||
found_tasks = self.check_scheduled_tasks_response(
|
||||
channel.json_body["scheduled_tasks"]
|
||||
)
|
||||
|
||||
# only the tasks with the given action names should have been returned
|
||||
self.assertEqual(len(found_tasks), 2)
|
||||
self.assertEqual(
|
||||
{task["action"] for task in found_tasks},
|
||||
{"test_task", "finished_test_task"},
|
||||
)
|
||||
|
||||
def test_filtering_scheduled_tasks_invalid_status(self) -> None:
|
||||
"""
|
||||
Test that an invalid `status` value is rejected with a 400 error.
|
||||
"""
|
||||
channel = self.make_request(
|
||||
"GET",
|
||||
"/_synapse/admin/v1/scheduled_tasks?status=unknown_status",
|
||||
content={},
|
||||
access_token=self.admin_user_tok,
|
||||
)
|
||||
self.assertEqual(400, channel.code, msg=channel.json_body)
|
||||
self.assertEqual(Codes.INVALID_PARAM, channel.json_body["errcode"])
|
||||
|
||||
@@ -3441,6 +3441,23 @@ class AnimatedThumbnailTestCase(unittest.HomeserverTestCase):
|
||||
)
|
||||
return out.getvalue()
|
||||
|
||||
def _make_mpo(self, stale: bool = False) -> bytes:
|
||||
"""Build a two-image MPO: a JPEG holding a stereo pair, not an
|
||||
animation. If `stale` is set, the trailing image is stripped but still
|
||||
advertised."""
|
||||
frames = [
|
||||
Image.new("RGB", (64, 64), color) for color in ((255, 0, 0), (0, 0, 255))
|
||||
]
|
||||
out = io.BytesIO()
|
||||
frames[0].save(out, format="MPO", save_all=True, append_images=frames[1:])
|
||||
data = out.getvalue()
|
||||
if not stale:
|
||||
return data
|
||||
|
||||
with Image.open(io.BytesIO(data)) as image:
|
||||
primary_size = image.mpinfo[0xB002][0]["Size"] # type: ignore[attr-defined]
|
||||
return data[:primary_size]
|
||||
|
||||
def _upload(self, data: bytes, content_type: str) -> MXCUri:
|
||||
return self.get_success(
|
||||
self.repo.create_or_update_content(
|
||||
@@ -3515,6 +3532,37 @@ class AnimatedThumbnailTestCase(unittest.HomeserverTestCase):
|
||||
result = Image.open(io.BytesIO(channel.result["body"]))
|
||||
self.assertFalse(getattr(result, "is_animated", False))
|
||||
|
||||
@parameterized.expand([("intact", False), ("stale_index", True)])
|
||||
def test_mpo_is_thumbnailed_as_a_still(self, _name: str, stale: bool) -> None:
|
||||
"""An MPO holds several stills rather than an animation, so it never
|
||||
gets an animated thumbnail.
|
||||
|
||||
Regression test for https://github.com/element-hq/synapse/issues/20024.
|
||||
"""
|
||||
mpo_uri = self._upload(self._make_mpo(stale=stale), "image/jpeg")
|
||||
|
||||
thumbnails = self.get_success(
|
||||
self.store.get_local_media_thumbnails(mpo_uri.media_id)
|
||||
)
|
||||
self.assertTrue(thumbnails)
|
||||
self.assertNotIn("image/webp", [info.type for info in thumbnails])
|
||||
|
||||
channel = self._thumbnail(mpo_uri, animated="true")
|
||||
result = Image.open(io.BytesIO(channel.result["body"]))
|
||||
self.assertFalse(getattr(result, "is_animated", False))
|
||||
|
||||
@parameterized.expand([("intact", False), ("stale_index", True)])
|
||||
@override_config({"dynamic_thumbnails": True})
|
||||
def test_dynamic_thumbnails_of_mpo_are_static(
|
||||
self, _name: str, stale: bool
|
||||
) -> None:
|
||||
"""The same holds when the thumbnail is generated on demand."""
|
||||
mpo_uri = self._upload(self._make_mpo(stale=stale), "image/jpeg")
|
||||
|
||||
channel = self._thumbnail(mpo_uri, animated="true")
|
||||
result = Image.open(io.BytesIO(channel.result["body"]))
|
||||
self.assertFalse(getattr(result, "is_animated", False))
|
||||
|
||||
@override_config({"dynamic_thumbnails": True})
|
||||
def test_dynamic_thumbnails_generates_and_caches_animated(self) -> None:
|
||||
"""With dynamic thumbnails, animated thumbnails are generated on demand
|
||||
|
||||
@@ -19,8 +19,12 @@
|
||||
#
|
||||
#
|
||||
|
||||
from unittest import mock
|
||||
|
||||
from twisted.internet.defer import CancelledError, Deferred, ensureDeferred
|
||||
from twisted.internet.testing import MemoryReactor
|
||||
|
||||
from synapse.logging.context import LoggingContext, make_deferred_yieldable
|
||||
from synapse.server import HomeServer
|
||||
from synapse.storage.database import (
|
||||
DatabasePool,
|
||||
@@ -225,6 +229,90 @@ class MultiWriterIdGeneratorTestCase(MultiWriterIdGeneratorBase):
|
||||
self.assertEqual(id_gen.get_positions(), {"master": 8})
|
||||
self.assertEqual(id_gen.get_current_token_for_writer("master"), 8)
|
||||
|
||||
def test_cancelled_enter_does_not_wedge_position(self) -> None:
|
||||
"""Reproduces presence getting stuck.
|
||||
|
||||
If the `get_next()` async context manager is cancelled while
|
||||
`__aenter__` is allocating a stream ID, the DB interaction that runs the
|
||||
sequence has already added the ID to `_unfinished_ids`, but `__aexit__`
|
||||
is never called (Python only invokes `__aexit__` if `__aenter__`
|
||||
returned). The abandoned ID is therefore leaked into `_unfinished_ids`
|
||||
forever, which permanently pins the persisted stream position: new rows
|
||||
keep getting higher IDs, but `get_current_token()` can never advance past
|
||||
`leaked_id - 1` until the process restarts.
|
||||
|
||||
This mirrors a `/sync` request being cancelled part-way through
|
||||
persisting a presence update. `/sync` became `@cancellable` in #19499,
|
||||
and on a monolith the presence write in `PresenceStore.update_presence`
|
||||
is awaited inside that cancellable request scope.
|
||||
"""
|
||||
# Prefill table with 7 rows written by 'master'; position starts at 7.
|
||||
self._insert_rows("master", 7)
|
||||
|
||||
id_gen = self._create_id_generator()
|
||||
self.assertEqual(id_gen.get_current_token_for_writer("master"), 7)
|
||||
|
||||
# We model the cancellation at the seam it actually happens in
|
||||
# production: `__aenter__` awaits `runInteraction("_load_next_mult_id")`,
|
||||
# whose transaction runs in a thread pool and so *always* completes -
|
||||
# allocating stream ID 8 and adding it to `_unfinished_ids` - but the
|
||||
# awaiting coroutine is handed a `CancelledError` because the enclosing
|
||||
# `/sync` request was cancelled. We reproduce that by letting the real
|
||||
# interaction run (applying its side effects) and then failing the
|
||||
# awaited deferred with `CancelledError`.
|
||||
cancel_enter: "Deferred[None]" = Deferred()
|
||||
original_run_interaction = id_gen._db.runInteraction
|
||||
|
||||
async def blocking_run_interaction(desc, func, *args, **kwargs): # type: ignore[no-untyped-def]
|
||||
result = await original_run_interaction(desc, func, *args, **kwargs)
|
||||
if desc == "_load_next_mult_id":
|
||||
# Stream ID 8 is now allocated and recorded in `_unfinished_ids`.
|
||||
# Deliver the cancellation here, exactly as a cancelled `/sync`
|
||||
# would land it on this `await`.
|
||||
await make_deferred_yieldable(cancel_enter)
|
||||
return result
|
||||
|
||||
async def presence_like_write() -> None:
|
||||
# Mirrors `PresenceStore.update_presence`: allocate an ID and
|
||||
# "persist" under the context manager.
|
||||
with LoggingContext(name="sync", server_name=self.hs.hostname):
|
||||
async with id_gen.get_next():
|
||||
pass
|
||||
|
||||
with mock.patch.object(
|
||||
id_gen._db, "runInteraction", new=blocking_run_interaction
|
||||
):
|
||||
write = ensureDeferred(presence_like_write())
|
||||
|
||||
# The write is now blocked inside `__aenter__`, i.e. after stream ID
|
||||
# 8 has been allocated and added to `_unfinished_ids`.
|
||||
self.assertNoResult(write)
|
||||
|
||||
# The client goes away and the `/sync` request is cancelled.
|
||||
cancel_enter.errback(CancelledError())
|
||||
|
||||
# The cancellation must surface as a `CancelledError`.
|
||||
self.get_failure(write, CancelledError)
|
||||
|
||||
# The cancelled write never persisted a row for ID 8, so the generator
|
||||
# must not let that abandoned ID wedge the position. A subsequent
|
||||
# *successful* write should be able to advance the persisted token.
|
||||
async def _successful_write() -> None:
|
||||
async with id_gen.get_next():
|
||||
pass
|
||||
|
||||
self.get_success(_successful_write())
|
||||
|
||||
# On the buggy code the token is still stuck at 7 (ID 8 is leaked in
|
||||
# `_unfinished_ids`, blocking everything behind it). Once the leak is
|
||||
# fixed, the token advances to 9: ID 8 was allocated (and abandoned) by
|
||||
# the cancelled write, so the successful write above takes ID 9.
|
||||
self.assertEqual(
|
||||
id_gen.get_current_token_for_writer("master"),
|
||||
9,
|
||||
"presence stream position is wedged by the cancelled allocation",
|
||||
)
|
||||
|
||||
def test_out_of_order_finish(self) -> None:
|
||||
"""Test that IDs persisted out of order are correctly handled"""
|
||||
|
||||
|
||||
Reference in New Issue
Block a user