diff --git a/.github/workflows/complement_tests.yml b/.github/workflows/complement_tests.yml index a891802ac8..629a48dcc7 100644 --- a/.github/workflows/complement_tests.yml +++ b/.github/workflows/complement_tests.yml @@ -39,7 +39,7 @@ jobs: steps: - name: Checkout synapse codebase - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 with: path: synapse @@ -65,7 +65,7 @@ jobs: - name: Prepare Complement's Prerequisites run: synapse/.ci/scripts/setup_complement_prerequisites.sh - - uses: actions/setup-go@4a3601121dd01d1626a1e23e37211e3254c1c06c # v6.4.0 + - uses: actions/setup-go@924ae3a1cded613372ab5595356fb5720e22ba16 # v6.5.0 with: cache-dependency-path: complement/go.sum go-version-file: complement/go.mod diff --git a/.github/workflows/docker.yml b/.github/workflows/docker.yml index 39ddb61918..8ecf8fb723 100644 --- a/.github/workflows/docker.yml +++ b/.github/workflows/docker.yml @@ -31,7 +31,7 @@ jobs: uses: docker/setup-buildx-action@d7f5e7f509e45cec5c76c4d5afdd7de93d0b3df5 # v4.1.0 - name: Checkout repository - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Extract version from pyproject.toml # Note: explicitly requesting bash will mean bash is invoked with `-eo pipefail`, see diff --git a/.github/workflows/docs-pr.yaml b/.github/workflows/docs-pr.yaml index e1a5b7be89..a1b355e74f 100644 --- a/.github/workflows/docs-pr.yaml +++ b/.github/workflows/docs-pr.yaml @@ -13,7 +13,7 @@ jobs: name: GitHub Pages runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 with: # Fetch all history so that the schema_versions script works. fetch-depth: 0 @@ -24,7 +24,7 @@ jobs: mdbook-version: '0.5.2' - name: Setup python - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" @@ -50,7 +50,7 @@ jobs: name: Check links in documentation runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Setup mdbook uses: peaceiris/actions-mdbook@ee69d230fe19748b7abf22df32acaa93833fad08 # v2.0.0 diff --git a/.github/workflows/docs.yaml b/.github/workflows/docs.yaml index 7236bf99d9..f8479a8c51 100644 --- a/.github/workflows/docs.yaml +++ b/.github/workflows/docs.yaml @@ -50,7 +50,7 @@ jobs: needs: - pre steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 with: # Fetch all history so that the schema_versions script works. fetch-depth: 0 @@ -64,7 +64,7 @@ jobs: run: echo 'window.SYNAPSE_VERSION = "${{ needs.pre.outputs.branch-version }}";' > ./docs/website_files/version.js - name: Setup python - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" diff --git a/.github/workflows/fix_lint.yaml b/.github/workflows/fix_lint.yaml index e0817698f4..88decb33d9 100644 --- a/.github/workflows/fix_lint.yaml +++ b/.github/workflows/fix_lint.yaml @@ -20,7 +20,7 @@ jobs: steps: - name: Checkout repository - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master diff --git a/.github/workflows/latest_deps.yml b/.github/workflows/latest_deps.yml index 815593ffcd..0b39c5c372 100644 --- a/.github/workflows/latest_deps.yml +++ b/.github/workflows/latest_deps.yml @@ -42,7 +42,7 @@ jobs: if: needs.check_repo.outputs.should_run_workflow == 'true' runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master with: @@ -77,7 +77,7 @@ jobs: postgres-version: "14" steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -93,7 +93,7 @@ jobs: -e POSTGRES_PASSWORD=postgres \ -e POSTGRES_INITDB_ARGS="--lc-collate C --lc-ctype C --encoding UTF8" \ postgres:${{ matrix.postgres-version }} - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - run: pip install .[all,test] @@ -151,7 +151,7 @@ jobs: BLACKLIST: ${{ matrix.workers && 'synapse-blacklist-with-workers' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -201,7 +201,7 @@ jobs: runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - uses: JasonEtco/create-an-issue@1b14a70e4d8dc185e5cc76d3bec9eab20257b2c5 # v2.9.2 env: GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/.github/workflows/poetry_lockfile.yaml b/.github/workflows/poetry_lockfile.yaml index 06545bd18a..856e51e7ab 100644 --- a/.github/workflows/poetry_lockfile.yaml +++ b/.github/workflows/poetry_lockfile.yaml @@ -16,8 +16,8 @@ jobs: name: "Check locked dependencies have sdists" runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: '3.x' - run: pip install tomli diff --git a/.github/workflows/push_complement_image.yml b/.github/workflows/push_complement_image.yml index 6f4c966cdc..34f86520b0 100644 --- a/.github/workflows/push_complement_image.yml +++ b/.github/workflows/push_complement_image.yml @@ -33,17 +33,17 @@ jobs: packages: write steps: - name: Checkout specific branch (debug build) - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 if: github.event_name == 'workflow_dispatch' with: ref: ${{ inputs.branch }} - name: Checkout clean copy of develop (scheduled build) - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 if: github.event_name == 'schedule' with: ref: develop - name: Checkout clean copy of master (on-push) - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 if: github.event_name == 'push' with: ref: master diff --git a/.github/workflows/release-artifacts.yml b/.github/workflows/release-artifacts.yml index c6b9f60baf..a54722e95b 100644 --- a/.github/workflows/release-artifacts.yml +++ b/.github/workflows/release-artifacts.yml @@ -27,8 +27,8 @@ jobs: name: "Calculate list of debian distros" runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - id: set-distros @@ -61,7 +61,7 @@ jobs: steps: - name: Checkout - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 with: path: src @@ -70,7 +70,7 @@ jobs: uses: docker/setup-buildx-action@d7f5e7f509e45cec5c76c4d5afdd7de93d0b3df5 # v4.1.0 - name: Set up docker layer caching - uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5 + uses: actions/cache@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6.1.0 with: path: /tmp/.buildx-cache key: ${{ runner.os }}-buildx-${{ github.sha }} @@ -78,7 +78,7 @@ jobs: ${{ runner.os }}-buildx- - name: Set up python - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" @@ -129,9 +129,9 @@ jobs: os: "ubuntu-24.04-arm" steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: # setup-python@v4 doesn't impose a default python version. Need to use 3.x # here, because `python` on osx points to Python 2.7. @@ -167,8 +167,8 @@ jobs: if: ${{ !startsWith(github.ref, 'refs/pull/') }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.10" diff --git a/.github/workflows/schema.yaml b/.github/workflows/schema.yaml index e36114d354..c4d46a1058 100644 --- a/.github/workflows/schema.yaml +++ b/.github/workflows/schema.yaml @@ -14,8 +14,8 @@ jobs: name: Ensure Synapse config schema is valid runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - name: Install check-jsonschema @@ -40,8 +40,8 @@ jobs: name: Ensure generated documentation is up-to-date runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - name: Install PyYAML diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 45fa2b8cad..58974ad7db 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -106,7 +106,7 @@ jobs: if: ${{ needs.changes.outputs.linting == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master with: @@ -126,8 +126,8 @@ jobs: if: ${{ needs.changes.outputs.linting == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - run: "pip install 'click==8.1.1' 'GitPython>=3.1.20' 'sqlglot>=28.0.0'" @@ -136,8 +136,8 @@ jobs: check-lockfile: runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - run: .ci/scripts/check_lockfile.py @@ -149,7 +149,7 @@ jobs: steps: - name: Checkout repository - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Setup Poetry uses: matrix-org/setup-python-poetry@5bbf6603c5c930615ec8a29f1b5d7d258d905aa4 # v2.0.0 @@ -171,7 +171,7 @@ jobs: steps: - name: Checkout repository - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -194,7 +194,7 @@ jobs: # Cribbed from # https://github.com/AustinScola/mypy-cache-github-action/blob/85ea4f2972abed39b33bd02c36e341b28ca59213/src/restore.ts#L10-L17 - name: Restore/persist mypy's cache - uses: actions/cache@27d5ce7f107fe9357f9df03efb73ab90386fccae # v5.0.5 + uses: actions/cache@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6.1.0 with: path: | .mypy_cache @@ -207,7 +207,7 @@ jobs: lint-crlf: runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Check line endings run: scripts-dev/check_line_terminators.sh @@ -216,11 +216,11 @@ jobs: if: ${{ github.event_name == 'pull_request' && (github.base_ref == 'develop' || contains(github.base_ref, 'release-')) && github.event.pull_request.user.login != 'dependabot[bot]' }} runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 with: ref: ${{ github.event.pull_request.head.sha }} fetch-depth: 0 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - run: "pip install 'towncrier>=18.6.0rc1'" @@ -234,7 +234,7 @@ jobs: if: ${{ needs.changes.outputs.rust == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -253,7 +253,7 @@ jobs: if: ${{ needs.changes.outputs.rust == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -271,7 +271,7 @@ jobs: steps: - name: Checkout repository - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -307,7 +307,7 @@ jobs: if: ${{ needs.changes.outputs.rust == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -326,9 +326,9 @@ jobs: if: ${{ needs.changes.outputs.golangci == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - - uses: actions/setup-go@4a3601121dd01d1626a1e23e37211e3254c1c06c # v6.4.0 + - uses: actions/setup-go@924ae3a1cded613372ab5595356fb5720e22ba16 # v6.5.0 with: cache-dependency-path: complement/go.sum go-version-file: complement/go.mod @@ -344,8 +344,8 @@ jobs: needs: changes if: ${{ needs.changes.outputs.linting_readme == 'true' }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - run: "pip install rstcheck" @@ -393,8 +393,8 @@ jobs: needs: linting-done runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.x" - id: get-matrix @@ -414,7 +414,7 @@ jobs: job: ${{ fromJson(needs.calculate-test-jobs.outputs.trial_test_matrix) }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - run: sudo apt-get -qq install xmlsec1 - name: Set up PostgreSQL ${{ matrix.job.postgres-version }} if: ${{ matrix.job.postgres-version }} @@ -470,7 +470,7 @@ jobs: - changes runs-on: ubuntu-22.04 steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -485,7 +485,7 @@ jobs: sudo apt-get -qq install build-essential libffi-dev python3-dev \ libxml2-dev libxslt-dev xmlsec1 zlib1g-dev libjpeg-dev libwebp-dev - - uses: actions/setup-python@a309ff8b426b58ec0e2a45f0f869d46889d02405 # v6.2.0 + - uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: python-version: "3.10" @@ -533,7 +533,7 @@ jobs: extras: ["all"] steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 # Install libs necessary for PyPy to build binary wheels for dependencies - run: sudo apt-get -qq install xmlsec1 libxml2-dev libxslt-dev - uses: matrix-org/setup-python-poetry@5bbf6603c5c930615ec8a29f1b5d7d258d905aa4 # v2.0.0 @@ -583,7 +583,7 @@ jobs: job: ${{ fromJson(needs.calculate-test-jobs.outputs.sytest_test_matrix) }} steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Prepare test blacklist run: cat sytest-blacklist .ci/worker-blacklist > synapse-blacklist-with-workers @@ -630,7 +630,7 @@ jobs: --health-retries 5 steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - run: sudo apt-get -qq install xmlsec1 postgresql-client - uses: matrix-org/setup-python-poetry@5bbf6603c5c930615ec8a29f1b5d7d258d905aa4 # v2.0.0 with: @@ -673,7 +673,7 @@ jobs: --health-retries 5 steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Add PostgreSQL apt repository # We need a version of pg_dump that can handle the version of # PostgreSQL being tested against. The Ubuntu package repository lags @@ -726,7 +726,7 @@ jobs: - changes steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -747,7 +747,7 @@ jobs: - changes steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master diff --git a/.github/workflows/triage_labelled.yml b/.github/workflows/triage_labelled.yml index 85d7be7b34..f6880fef1f 100644 --- a/.github/workflows/triage_labelled.yml +++ b/.github/workflows/triage_labelled.yml @@ -22,7 +22,7 @@ jobs: # This field is case-sensitive. TARGET_STATUS: Needs info steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 with: # Only clone the script file we care about, instead of the whole repo. sparse-checkout: .ci/scripts/triage_labelled_issue.sh diff --git a/.github/workflows/twisted_trunk.yml b/.github/workflows/twisted_trunk.yml index 1b906f7f44..80984ceeb3 100644 --- a/.github/workflows/twisted_trunk.yml +++ b/.github/workflows/twisted_trunk.yml @@ -42,7 +42,7 @@ jobs: runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -69,7 +69,7 @@ jobs: runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - run: sudo apt-get -qq install xmlsec1 - name: Install Rust @@ -115,7 +115,7 @@ jobs: - ${{ github.workspace }}:/src steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - name: Install Rust uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9 # master @@ -172,7 +172,7 @@ jobs: runs-on: ubuntu-latest steps: - - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3 + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0 - uses: JasonEtco/create-an-issue@1b14a70e4d8dc185e5cc76d3bec9eab20257b2c5 # v2.9.2 env: GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/Cargo.lock b/Cargo.lock index 87f50ea5df..7ca2265aca 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -78,9 +78,9 @@ checksum = "46c5e41b57b8bba42a04676d81cb89e9ee8e859a1a66f80a5a72e1cb76b34d43" [[package]] name = "bytes" -version = "1.11.1" +version = "1.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1e748733b7cbc798e1434b6ac524f0c1ff2ab456fe201501e6497c8417a4fc33" +checksum = "8ae3f5d315924270530207e2a68396c3cc547f6dca3fbdca317cfb1a51edb593" [[package]] name = "cc" @@ -718,9 +718,9 @@ checksum = "241eaef5fd12c88705a01fc1066c48c4b36e0dd4377dcdc7ec3942cea7a69956" [[package]] name = "log" -version = "0.4.32" +version = "0.4.33" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "953f07c43838f8e6f9758cab68bf5bed85465e7587ebe0b823f1bcd81978ad3a" +checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad" [[package]] name = "lru-slab" diff --git a/changelog.d/19539.bugfix b/changelog.d/19539.bugfix new file mode 100644 index 0000000000..41beb3a179 --- /dev/null +++ b/changelog.d/19539.bugfix @@ -0,0 +1 @@ +[MSC4140: Cancellable delayed events](https://github.com/matrix-org/matrix-spec-proposals/pull/4140): Update error responses to match their format in the current draft of the MSC. diff --git a/changelog.d/19539.feature b/changelog.d/19539.feature new file mode 100644 index 0000000000..93eb4cc1b9 --- /dev/null +++ b/changelog.d/19539.feature @@ -0,0 +1 @@ +[MSC4140: Cancellable delayed events](https://github.com/matrix-org/matrix-spec-proposals/pull/4140): Limit how many delayed events a user may have scheduled at once. diff --git a/changelog.d/19663.feature b/changelog.d/19663.feature new file mode 100644 index 0000000000..9f61182a90 --- /dev/null +++ b/changelog.d/19663.feature @@ -0,0 +1 @@ +Support [MSC4446](https://github.com/matrix-org/matrix-spec-proposals/pull/4446) for moving fully read markers backwards. Contributed by @SpiritCroc @ Beeper. diff --git a/changelog.d/19913.misc b/changelog.d/19913.misc deleted file mode 100644 index 31b5a4d37b..0000000000 --- a/changelog.d/19913.misc +++ /dev/null @@ -1 +0,0 @@ -Fix a flake in 3PID inhibit error unit tests, causing occasional failures in CI. \ No newline at end of file diff --git a/changelog.d/19916.misc b/changelog.d/19916.misc new file mode 100644 index 0000000000..aeec5ae412 --- /dev/null +++ b/changelog.d/19916.misc @@ -0,0 +1 @@ +Add note to 3PID email token request unit tests that the endpoint being tested can have an expected, artificial delay of up to 1s. \ No newline at end of file diff --git a/changelog.d/19928.bugfix b/changelog.d/19928.bugfix new file mode 100644 index 0000000000..d69e4febbe --- /dev/null +++ b/changelog.d/19928.bugfix @@ -0,0 +1 @@ +Fix a regression where application services that opted into ephemeral events using the legacy `de.sorunome.msc2409.push_ephemeral` registration flag stopped receiving ephemeral events (including to-device messages used for encryption). Introduced in v1.156.0. diff --git a/changelog.d/19935.feature b/changelog.d/19935.feature new file mode 100644 index 0000000000..cd5fa804b8 --- /dev/null +++ b/changelog.d/19935.feature @@ -0,0 +1 @@ +Add an `exclude_rooms_from_presence` configuration option to stop presence being routed between users solely because they share one of the listed rooms. diff --git a/changelog.d/19939.misc b/changelog.d/19939.misc new file mode 100644 index 0000000000..59d480c98d --- /dev/null +++ b/changelog.d/19939.misc @@ -0,0 +1 @@ +Minor presence performance improvements for large servers. diff --git a/changelog.d/19941.misc b/changelog.d/19941.misc new file mode 100644 index 0000000000..e928f68bc5 --- /dev/null +++ b/changelog.d/19941.misc @@ -0,0 +1 @@ +Reduce replication traffic caused by presence. diff --git a/changelog.d/19947.bugfix b/changelog.d/19947.bugfix new file mode 100644 index 0000000000..c16d6deea0 --- /dev/null +++ b/changelog.d/19947.bugfix @@ -0,0 +1 @@ +Fix a bug causing device list pruning to skip some rows when the transaction gets retried. \ No newline at end of file diff --git a/changelog.d/19948.bugfix b/changelog.d/19948.bugfix new file mode 100644 index 0000000000..5cd34d1abf --- /dev/null +++ b/changelog.d/19948.bugfix @@ -0,0 +1 @@ +Fix presence states being shown to clients forever after presence is disabled, by marking any previously only users as offline. diff --git a/changelog.d/19949.bugfix b/changelog.d/19949.bugfix new file mode 100644 index 0000000000..24ddd4dc4f --- /dev/null +++ b/changelog.d/19949.bugfix @@ -0,0 +1 @@ +Fix `SYNAPSE_ASYNC_IO_REACTOR=1` on Python 3.14. diff --git a/complement/go.mod b/complement/go.mod index aa2333a5a0..cea013573b 100644 --- a/complement/go.mod +++ b/complement/go.mod @@ -51,7 +51,7 @@ require ( go.opentelemetry.io/otel/sdk v1.43.0 // indirect go.opentelemetry.io/otel/sdk/metric v1.43.0 // indirect go.opentelemetry.io/otel/trace v1.43.0 // indirect - golang.org/x/crypto v0.51.0 // indirect + golang.org/x/crypto v0.52.0 // indirect golang.org/x/net v0.55.0 // indirect golang.org/x/sync v0.20.0 // indirect golang.org/x/sys v0.45.0 // indirect diff --git a/complement/go.sum b/complement/go.sum index 49c3724e83..4858bc78c8 100644 --- a/complement/go.sum +++ b/complement/go.sum @@ -122,8 +122,8 @@ go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXd golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= -golang.org/x/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI= -golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8= +golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988= +golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc= golang.org/x/exp v0.0.0-20240719175910-8a7402abbf56 h1:2dVuKD2vS7b0QIHQbpyTISPd0LeHDbnYEryqj5Q1ug8= golang.org/x/exp v0.0.0-20240719175910-8a7402abbf56/go.mod h1:M4RDyNAINzryxdtnbRXRL/OHtkFuWGRjvuhBJpk2IlY= golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= diff --git a/docs/usage/configuration/config_documentation.md b/docs/usage/configuration/config_documentation.md index d8085e0e8a..ab2241ecfc 100644 --- a/docs/usage/configuration/config_documentation.md +++ b/docs/usage/configuration/config_documentation.md @@ -4320,6 +4320,16 @@ exclude_rooms_from_sync: - '!foo:example.com' ``` --- +### `exclude_rooms_from_presence` + +*(array)* A list of rooms to exclude from presence updates. Presence will not be routed between two users solely because they share one of these rooms. Users who also share a non-excluded room continue to exchange presence as normal. Defaults to `[]`. + +Example configuration: +```yaml +exclude_rooms_from_presence: +- '!foo:example.com' +``` +--- ## Opentracing Configuration options related to Opentracing support. diff --git a/rust/src/config/mod.rs b/rust/src/config/mod.rs index b59d61e5f9..1c97373c27 100644 --- a/rust/src/config/mod.rs +++ b/rust/src/config/mod.rs @@ -47,7 +47,7 @@ pub struct AuthConfig { } #[derive(FromPyObject, Clone)] pub struct ServerConfig { - pub max_event_delay_ms: Option, + pub msc4140_enabled: bool, pub include_profile_updates_in_sync: bool, } @@ -73,4 +73,5 @@ pub struct ExperimentalConfig { pub msc4222_enabled: bool, pub msc4491_enabled: bool, pub msc4143_enabled: bool, + pub msc4446_enabled: bool, } diff --git a/rust/src/handlers/versions.rs b/rust/src/handlers/versions.rs index 5bf36fa4e9..5d35b052bd 100644 --- a/rust/src/handlers/versions.rs +++ b/rust/src/handlers/versions.rs @@ -269,6 +269,9 @@ pub struct UnstableFeatureMap { /// MSC4143: Matrix RTC transports (LiveKit backend) #[serde(rename = "org.matrix.msc4143")] msc4143_enabled: bool, + /// MSC4446: Allow moving the fully read marker backwards. + #[serde(rename = "com.beeper.msc4446")] + msc4446_enabled: bool, // Whether new rooms will be set to encrypted or not (based on presets). #[serde(rename = "io.element.e2ee_forced.public")] @@ -307,10 +310,7 @@ pub fn synapse_config_to_global_unstable_feature_map( msc4028: config.experimental.msc4028_push_encrypted_events, msc4108: config.experimental.msc4108_enabled || (config.experimental.msc4108_delegation_endpoint.is_some()), - msc4140: config - .server - .max_event_delay_ms - .is_some_and(|max_event_delay_ms| max_event_delay_ms > 0), + msc4140: config.server.msc4140_enabled, msc3575: config.experimental.msc3575_enabled, msc4133: config.experimental.msc4133_enabled, msc4133_stable: true, @@ -323,6 +323,7 @@ pub fn synapse_config_to_global_unstable_feature_map( msc4445_initial_sync_timeline_topological_ordering: true, msc4491_enabled: config.experimental.msc4491_enabled, msc4143_enabled: config.experimental.msc4143_enabled, + msc4446_enabled: config.experimental.msc4446_enabled, e2ee_forced_public: config .room .encryption_enabled_by_default_for_room_presets diff --git a/schema/synapse-config.schema.yaml b/schema/synapse-config.schema.yaml index 34754de734..b9236fa523 100644 --- a/schema/synapse-config.schema.yaml +++ b/schema/synapse-config.schema.yaml @@ -5365,6 +5365,18 @@ properties: default: [] examples: - - "!foo:example.com" + exclude_rooms_from_presence: + type: array + description: >- + A list of rooms to exclude from presence updates. Presence will not be + routed between two users solely because they share one of these rooms. + Users who also share a non-excluded room continue to exchange presence as + normal. + items: + type: string + default: [] + examples: + - - "!foo:example.com" opentracing: type: object description: >- diff --git a/synapse/__init__.py b/synapse/__init__.py index 3acfc1a0d7..a223066f04 100644 --- a/synapse/__init__.py +++ b/synapse/__init__.py @@ -49,7 +49,9 @@ if strtobool(os.environ.get("SYNAPSE_ASYNC_IO_REACTOR", "0")): from twisted.internet import asyncioreactor - asyncioreactor.install(asyncio.get_event_loop()) + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + asyncioreactor.install(loop) # Twisted and canonicaljson will fail to import when this file is executed to # get the __version__ during a fresh install. That's OK and subsequent calls to diff --git a/synapse/config/experimental.py b/synapse/config/experimental.py index f99f7b139e..6ad9f53517 100644 --- a/synapse/config/experimental.py +++ b/synapse/config/experimental.py @@ -287,6 +287,10 @@ class ExperimentalConfig(Config): # (and MSC4308: Thread Subscriptions extension to Sliding Sync) self.msc4306_enabled: bool = experimental.get("msc4306_enabled", False) + # MSC4446: Allow moving the fully read marker backwards. + # Tracked in: https://github.com/element-hq/synapse/issues/19940 + self.msc4446_enabled: bool = experimental.get("msc4446_enabled", False) + # MSC4354: Sticky Events # Tracked in: https://github.com/element-hq/synapse/issues/19409 # Note that sticky events persisted before this feature is enabled will not be diff --git a/synapse/config/server.py b/synapse/config/server.py index 22d3d0c1b1..e9767e05e6 100644 --- a/synapse/config/server.py +++ b/synapse/config/server.py @@ -37,6 +37,7 @@ from twisted.conch.ssh.keys import Key from synapse.api.room_versions import KNOWN_ROOM_VERSIONS from synapse.types import JsonDict, StrSequence +from synapse.util.duration import Duration from synapse.util.module_loader import load_module from synapse.util.stringutils import parse_and_validate_server_name @@ -941,6 +942,10 @@ class ServerConfig(Config): config.get("exclude_rooms_from_sync") or [] ) + self.rooms_to_exclude_from_presence: list[str] = ( + config.get("exclude_rooms_from_presence") or [] + ) + delete_stale_devices_after: str | None = ( config.get("delete_stale_devices_after") or None ) @@ -955,13 +960,33 @@ class ServerConfig(Config): # The maximum allowed delay duration for delayed events (MSC4140). max_event_delay_duration = config.get("max_event_delay_duration") if max_event_delay_duration is not None: - self.max_event_delay_ms: int | None = self.parse_duration( - max_event_delay_duration - ) - if self.max_event_delay_ms <= 0: - raise ConfigError("max_event_delay_duration must be a positive value") + max_event_delay_ms = self.parse_duration(max_event_delay_duration) + if max_event_delay_ms <= 0: + raise ConfigError( + "'max_event_delay_duration' must be a positive value if set", + ("max_event_delay_duration",), + ) + self.max_event_delay_duration = Duration(milliseconds=max_event_delay_ms) else: - self.max_event_delay_ms = None + self.max_event_delay_duration = Duration() + + # The maximum number of delayed events a user may have scheduled at a time. + # (Defined here despite being experimental to be near the other MSC4140 config) + self.max_delayed_events_per_user: int = config.get( + "experimental_features", {} + ).get("msc4140_max_delayed_events_per_user", 100) + if ( + not isinstance(self.max_delayed_events_per_user, int) + or self.max_delayed_events_per_user < 0 + ): + raise ConfigError( + "'msc4140_max_delayed_events_per_user' must be a non-negative integer", + ("experimental", "msc4140_max_delayed_events_per_user"), + ) + + self.msc4140_enabled = bool( + self.max_delayed_events_per_user and self.max_event_delay_duration + ) def has_tls_listener(self) -> bool: return any(listener.is_tls() for listener in self.listeners) diff --git a/synapse/handlers/appservice.py b/synapse/handlers/appservice.py index 36b2f63e41..68b8aa71f1 100644 --- a/synapse/handlers/appservice.py +++ b/synapse/handlers/appservice.py @@ -313,7 +313,10 @@ class ApplicationServicesHandler: StreamKeyType.PRESENCE, StreamKeyType.TO_DEVICE, ) - and service.supports_ephemeral + # Honour both the stable `receive_ephemeral` registration flag and the + # legacy `de.sorunome.msc2409.push_ephemeral` one, matching the + # transaction body built in `ApplicationServiceApi.push_bulk`. + and (service.supports_ephemeral or service.supports_unstable_ephemeral) ) or ( stream_key == StreamKeyType.DEVICE_LIST diff --git a/synapse/handlers/delayed_events.py b/synapse/handlers/delayed_events.py index f016d95e31..13d6a54de2 100644 --- a/synapse/handlers/delayed_events.py +++ b/synapse/handlers/delayed_events.py @@ -13,12 +13,13 @@ # import logging +from http import HTTPStatus from typing import TYPE_CHECKING, Optional from twisted.internet.interfaces import IDelayedCall from synapse.api.constants import EventTypes, StickyEvent, StickyEventField -from synapse.api.errors import ShadowBanError, SynapseError +from synapse.api.errors import Codes, ShadowBanError, SynapseError from synapse.api.ratelimiting import Ratelimiter from synapse.config.workers import MAIN_PROCESS_INSTANCE_NAME from synapse.http.site import SynapseRequest @@ -330,7 +331,7 @@ class DelayedEventsHandler: state_key: str | None, origin_server_ts: int | None, content: JsonDict, - delay: int, + delay: Duration, sticky_duration_ms: int | None, ) -> str: """ @@ -344,20 +345,37 @@ class DelayedEventsHandler: origin_server_ts: The custom timestamp to send the event with. If None, the timestamp will be the actual time when the event is sent. content: The content of the event to be sent. - delay: How long (in milliseconds) to wait before automatically sending the event. + delay: How long to wait before automatically sending the event. sticky_duration_ms: If an MSC4354 sticky event: the sticky duration (in milliseconds). The event will be attempted to be reliably delivered to clients and remote servers during its sticky period. Returns: The ID of the added delayed event. Raises: - SynapseError: if the delayed event fails validation checks. + SynapseError: if the delayed event fails validation checks, or + if the requested delay is longer than allowed, or + if sending delayed events has been disallowed entirely. """ # Use standard request limiter for scheduling new delayed events. # TODO: Instead apply ratelimiting based on the scheduled send time. # See https://github.com/element-hq/synapse/issues/18021 await self._request_ratelimiter.ratelimit(requester) + if not self._config.server.msc4140_enabled: + raise SynapseError( + HTTPStatus.FORBIDDEN, + "Sending delayed events has been disallowed", + Codes.FORBIDDEN, + ) + if delay > self._config.server.max_event_delay_duration: + requested_delay = delay.as_millis() + max_delay = self._config.server.max_event_delay_duration.as_millis() + raise SynapseError( + HTTPStatus.FORBIDDEN, + f"The requested delay ({requested_delay}ms) exceeds the allowed maximum ({max_delay}ms)", + Codes.FORBIDDEN, + ) + self._event_creation_handler.validator.validate_builder( self._event_creation_handler.event_builder_factory.for_room_version( await self._store.get_room_version(room_id), @@ -384,6 +402,7 @@ class DelayedEventsHandler: content=content, delay=delay, sticky_duration_ms=sticky_duration_ms, + limit=self._config.server.max_delayed_events_per_user, ) if self._repl_client is not None: diff --git a/synapse/handlers/presence.py b/synapse/handlers/presence.py index c1994fe488..55dd5ffb59 100644 --- a/synapse/handlers/presence.py +++ b/synapse/handlers/presence.py @@ -122,6 +122,7 @@ from synapse.types import ( ) from synapse.util.async_helpers import Linearizer from synapse.util.duration import Duration +from synapse.util.iterutils import batch_iter from synapse.util.metrics import Measure from synapse.util.wheel_timer import WheelTimer @@ -222,6 +223,12 @@ class BasePresenceHandler(abc.ABC): self._presence_enabled = hs.config.server.presence_enabled self._track_presence = hs.config.server.track_presence + # Rooms which, on their own, should not cause presence to be routed + # between their members. See `exclude_rooms_from_presence` in the config. + self._rooms_to_exclude_from_presence = frozenset( + hs.config.server.rooms_to_exclude_from_presence + ) + # The (configurable) presence state machine timers. self._last_active_granularity = ( hs.config.server.presence_last_active_granularity @@ -435,6 +442,7 @@ class BasePresenceHandler(abc.ABC): self.store, self.presence_router, states, + self._rooms_to_exclude_from_presence, ) for destinations, host_states in hosts_to_states: @@ -530,10 +538,34 @@ class WorkerPresenceHandler(BasePresenceHandler): # syncing but we haven't notified the presence writer of that yet self._user_devices_going_offline: dict[tuple[str, str | None], int] = {} + # How often to relay an unchanged sync-driven presence state to the + # presence writer. The relayed updates are what feed the writer's device + # last_sync_ts/last_active_ts timers, so this must sit comfortably below + # the timers it feeds — the (configurable) sync online timeout and + # last-active granularity — or users would flap offline / lose + # "currently active" between relays. We use 5/6 of the tighter of the + # two, i.e. the historic 25s at the default 30s sync online timeout. + self._sync_presence_relay_interval = ( + min(self._sync_online_timeout, self._last_active_granularity) * 5 // 6 + ) + + # (user_id, device_id) -> (state, last_sent_ms) of the most recent + # sync-driven presence update we proxied to the presence writer. Used + # to suppress the per-sync-request set_state/bump calls, which are + # no-ops on the writer at finer granularity than its timers: while + # the state is unchanged there is no point relaying more than one + # update per relay interval. Entries older than the window are swept by + # `_sweep_last_sent_presence`. + self._last_sent_presence: dict[tuple[str, str | None], tuple[str, int]] = {} + self._bump_active_client = ReplicationBumpPresenceActiveTime.make_client(hs) self._set_state_client = ReplicationPresenceSetState.make_client(hs) - self.clock.looping_call(self.send_stop_syncing, UPDATE_SYNCING_USERS) + if self._track_presence: + self.clock.looping_call(self.send_stop_syncing, UPDATE_SYNCING_USERS) + self.clock.looping_call( + self._sweep_last_sent_presence, Duration(minutes=30) + ) hs.register_async_shutdown_handler( phase="before", @@ -576,6 +608,8 @@ class WorkerPresenceHandler(BasePresenceHandler): sending a stopped syncing immediately followed by a started syncing notification to the presence writer """ + if not self._track_presence: + return self._user_devices_going_offline[(user_id, device_id)] = self.clock.time_msec() def send_stop_syncing(self) -> None: @@ -589,6 +623,22 @@ class WorkerPresenceHandler(BasePresenceHandler): if now - last_sync_ms > UPDATE_SYNCING_USERS.as_millis(): self._user_devices_going_offline.pop((user_id, device_id), None) self.send_user_sync(user_id, device_id, False, last_sync_ms) + # Once the writer knows the device stopped syncing it may time + # the user out, so if the device comes back we must relay its + # state again rather than suppress it as a repeat. + self._last_sent_presence.pop((user_id, device_id), None) + + def _sweep_last_sent_presence(self) -> None: + """Drop expired presence-throttling entries. + + Entries should be dropped in `send_stop_syncing`, but we add a safety + net here to ensure that the dict deesn't grow unbounded. + """ + now = self.clock.time_msec() + + for key, (_, last_sent_ms) in list(self._last_sent_presence.items()): + if now - last_sent_ms >= self._sync_presence_relay_interval: + self._last_sent_presence.pop(key, None) async def user_syncing( self, @@ -606,7 +656,9 @@ class WorkerPresenceHandler(BasePresenceHandler): return _NullContextManager() # Note that this causes last_active_ts to be incremented which is not - # what the spec wants. + # what the spec wants. (This call is throttled in `set_state`: while + # the state is unchanged, only one update per relay interval is relayed + # to the presence writer.) await self.set_state( UserID.from_string(user_id), device_id, @@ -644,7 +696,12 @@ class WorkerPresenceHandler(BasePresenceHandler): async def notify_from_replication( self, states: list[UserPresenceState], stream_id: int ) -> None: - parties = await get_interested_parties(self.store, self.presence_router, states) + parties = await get_interested_parties( + self.store, + self.presence_router, + states, + self._rooms_to_exclude_from_presence, + ) room_ids_to_states, users_to_states = parties self.notifier.on_new_event( @@ -743,6 +800,28 @@ class WorkerPresenceHandler(BasePresenceHandler): if not self._track_presence: return + now = self.clock.time_msec() + if is_sync and not force_notify: + # Sync-driven updates arrive on every /sync request, which is far + # finer-grained than any of the writer's presence timers need: + # while the state is unchanged, relaying one update per relay + # interval is enough to keep them fed. State changes always go + # through immediately. + last_sent = self._last_sent_presence.get((user_id, device_id)) + if last_sent is not None: + last_presence, last_sent_ms = last_sent + if ( + presence == last_presence + and now - last_sent_ms < self._sync_presence_relay_interval + ): + return + self._last_sent_presence[(user_id, device_id)] = (presence, now) + else: + # An explicit (non-sync) update doesn't refresh the writer's + # last_sync_ts, so it must not count as a recent relay: drop any + # entry so the next sync-driven update goes through. + self._last_sent_presence.pop((user_id, device_id), None) + # Proxy request to instance that writes presence await self._set_state_client( instance_name=self._presence_writer_instance, @@ -763,8 +842,25 @@ class WorkerPresenceHandler(BasePresenceHandler): if not self._track_presence: return - # Proxy request to instance that writes presence user_id = user.to_string() + + # A bump's only effects on the writer are updating last_active_ts and + # flipping an idle device back online. Going idle takes far longer + # than the relay window, so if we relayed an *online* update within + # the window the user cannot have gone idle since, and this bump is a + # no-op: skip it. Bumps after any other state (or an unknown one) go + # through immediately, as they may un-idle the device. + now = self.clock.time_msec() + last_sent = self._last_sent_presence.get((user_id, device_id)) + if ( + last_sent is not None + and last_sent[0] == PresenceState.ONLINE + and now - last_sent[1] < self._sync_presence_relay_interval + ): + return + self._last_sent_presence[(user_id, device_id)] = (PresenceState.ONLINE, now) + + # Proxy request to instance that writes presence await self._bump_active_client( instance_name=self._presence_writer_instance, user_id=user_id, @@ -889,6 +985,14 @@ class PresenceHandler(BasePresenceHandler): Duration(minutes=1), ) + if not self._presence_enabled and self.user_to_current_state: + # Presence is disabled but the database still contains non-offline + # presence states, i.e. presence used to be enabled. Nothing writes + # to the presence stream while presence is disabled, so without + # intervention clients would show the stale states forever. Send + # out one final round of updates marking everyone as offline. + self.clock.call_when_running(self._mark_stale_presence_as_offline) + presence_wheel_timer_size_gauge.register_hook( homeserver_instance_id=hs.get_instance_id(), hook=lambda: {(self.server_name,): len(self.wheel_timer)}, @@ -946,6 +1050,36 @@ class PresenceHandler(BasePresenceHandler): [self.user_to_current_state[user_id] for user_id in unpersisted] ) + @wrap_as_background_process("PresenceHandler._mark_stale_presence_as_offline") + async def _mark_stale_presence_as_offline(self) -> None: + """One-off job, run at startup when presence is disabled, that marks + any non-offline presence states left over from when presence was + enabled as offline, and streams the changes out to clients. + """ + states = [ + state.copy_and_replace( + state=PresenceState.OFFLINE, + status_msg=None, + currently_active=False, + ) + for state in self.user_to_current_state.values() + if state.state != PresenceState.OFFLINE + ] + if not states: + return + + logger.info( + "Presence is disabled: marking %d stale presence states as offline", + len(states), + ) + + self.user_to_current_state.update({state.user_id: state for state in states}) + + # There may be a lot of stale states (e.g. everyone that was online + # when presence was disabled), so persist them in batches. + for batch in batch_iter(states, 500): + await self._persist_and_notify(list(batch)) + async def _update_states( self, new_states: Iterable[UserPresenceState], @@ -1058,6 +1192,7 @@ class PresenceHandler(BasePresenceHandler): self.store, self.presence_router, list(to_federation_ping.values()), + self._rooms_to_exclude_from_presence, ) for destinations, states in hosts_to_states: @@ -1331,7 +1466,12 @@ class PresenceHandler(BasePresenceHandler): """ stream_id, max_token = await self.store.update_presence(states) - parties = await get_interested_parties(self.store, self.presence_router, states) + parties = await get_interested_parties( + self.store, + self.presence_router, + states, + self._rooms_to_exclude_from_presence, + ) room_ids_to_states, users_to_states = parties self.notifier.on_new_event( @@ -1478,7 +1618,10 @@ class PresenceHandler(BasePresenceHandler): observed_user.to_string() ) - if observer_room_ids & observed_room_ids: + shared_room_ids = ( + observer_room_ids & observed_room_ids + ) - self._rooms_to_exclude_from_presence + if shared_room_ids: return True return False @@ -1589,6 +1732,12 @@ class PresenceHandler(BasePresenceHandler): to be handled. """ + # Excluded rooms should not, on their own, share presence between their + # members. This method is entirely per-room presence fan-out, so skip + # excluded rooms wholesale. + if room_id in self._rooms_to_exclude_from_presence: + return + # Sets of newly joined users. Note that if the local server is # joining a remote room for the first time we'll see both the joining # user and all remote users as newly joined. @@ -1846,6 +1995,9 @@ class PresenceEventSource(EventSource[int, UserPresenceState]): self.server_name = hs.hostname self.clock = hs.get_clock() self.store = hs.get_datastores().main + self._rooms_to_exclude_from_presence = frozenset( + hs.config.server.rooms_to_exclude_from_presence + ) async def get_new_events( self, @@ -1960,9 +2112,31 @@ class PresenceEventSource(EventSource[int, UserPresenceState]): **{SERVER_NAME_LABEL: self.server_name}, ).inc() - sharing_users = await self.store.do_users_share_a_room( - user_id, updated_users - ) + # An updated user is interesting if they share a + # (non-excluded) room with the syncing user. We check by + # intersecting the cached per-user room sets rather than via + # `do_users_share_a_room`: its per-pair cache has a + # quadratic working set and is cleared wholesale on every + # membership change, so on busy servers every check missed + # into SQL. + # + # For every presence update we need to run this code for + # every user that is currently syncing. The + # `get_rooms_for_user` will therefore be computed only once + # for each updated user regardless of the number of syncing + # users. + # + # The syncing user's rooms will also be cached as its needed + # during sync processing anyway. + my_rooms = await self.store.get_rooms_for_user(user_id) + if self._rooms_to_exclude_from_presence: + my_rooms = my_rooms - self._rooms_to_exclude_from_presence + rooms_by_user = await self.store.get_rooms_for_users(updated_users) + sharing_users = { + updated_user + for updated_user, rooms in rooms_by_user.items() + if not my_rooms.isdisjoint(rooms) + } interested_and_updated_users = ( sharing_users.union(additional_users_interested_in) @@ -1977,7 +2151,9 @@ class PresenceEventSource(EventSource[int, UserPresenceState]): ).inc() users_interested_in = ( - await self.store.get_users_who_share_room_with_user(user_id) + await self.store.get_users_who_share_room_with_user( + user_id, self._rooms_to_exclude_from_presence + ) ) users_interested_in.update(additional_users_interested_in) @@ -1990,7 +2166,9 @@ class PresenceEventSource(EventSource[int, UserPresenceState]): # No from_key has been specified. Return the presence for all users # this user is interested in interested_and_updated_users = ( - await self.store.get_users_who_share_room_with_user(user_id) + await self.store.get_users_who_share_room_with_user( + user_id, self._rooms_to_exclude_from_presence + ) ) interested_and_updated_users.update(additional_users_interested_in) @@ -2390,7 +2568,10 @@ def _combine_device_states( async def get_interested_parties( - store: DataStore, presence_router: PresenceRouter, states: list[UserPresenceState] + store: DataStore, + presence_router: PresenceRouter, + states: list[UserPresenceState], + excluded_rooms: AbstractSet[str] = frozenset(), ) -> tuple[dict[str, list[UserPresenceState]], dict[str, list[UserPresenceState]]]: """Given a list of states return which entities (rooms, users) are interested in the given states. @@ -2399,6 +2580,8 @@ async def get_interested_parties( store: The homeserver's data store. presence_router: A module for augmenting the destinations for presence updates. states: A list of incoming user presence updates. + excluded_rooms: Rooms which should not, on their own, cause presence to + be routed between their members. Returns: A 2-tuple of `(room_ids_to_states, users_to_states)`, @@ -2409,6 +2592,8 @@ async def get_interested_parties( for state in states: room_ids = await store.get_rooms_for_user(state.user_id) for room_id in room_ids: + if room_id in excluded_rooms: + continue room_ids_to_states.setdefault(room_id, []).append(state) # Always notify self @@ -2429,6 +2614,7 @@ async def get_interested_remotes( store: DataStore, presence_router: PresenceRouter, states: list[UserPresenceState], + excluded_rooms: AbstractSet[str] = frozenset(), ) -> list[tuple[StrCollection, Collection[UserPresenceState]]]: """Given a list of presence states figure out which remote servers should be sent which. @@ -2439,6 +2625,8 @@ async def get_interested_remotes( store: The homeserver's data store. presence_router: A module for augmenting the destinations for presence updates. states: A list of incoming user presence updates. + excluded_rooms: Rooms which should not, on their own, cause presence to + be routed to their remote members. Returns: A map from destinations to presence states to send to that destination. @@ -2452,6 +2640,8 @@ async def get_interested_remotes( room_ids = await store.get_rooms_for_user(state.user_id) hosts: set[str] = set() for room_id in room_ids: + if room_id in excluded_rooms: + continue room_hosts = await store.get_current_hosts_in_room(room_id) hosts.update(room_hosts) hosts_and_states.append((hosts, [state])) diff --git a/synapse/handlers/read_marker.py b/synapse/handlers/read_marker.py index 85d2dd62bb..3f3b9e6d8b 100644 --- a/synapse/handlers/read_marker.py +++ b/synapse/handlers/read_marker.py @@ -41,7 +41,11 @@ class ReadMarkerHandler: ) async def received_client_read_marker( - self, room_id: str, user_id: str, event_id: str + self, + room_id: str, + user_id: str, + event_id: str, + allow_backward: bool = False, ) -> None: """Updates the read marker for a given user in a given room if the event ID given is ahead in the stream relative to the current read marker. @@ -59,7 +63,7 @@ class ReadMarkerHandler: # Get event ordering, this also ensures we know about the event event_ordering = await self.store.get_event_ordering(event_id, room_id) - if existing_read_marker: + if existing_read_marker and not allow_backward: try: old_event_ordering = await self.store.get_event_ordering( existing_read_marker["event_id"], room_id diff --git a/synapse/handlers/sync.py b/synapse/handlers/sync.py index 9a8b1d5192..1cbe7fe1a1 100644 --- a/synapse/handlers/sync.py +++ b/synapse/handlers/sync.py @@ -1828,9 +1828,19 @@ class SyncHandler: await self._generate_sync_entry_for_account_data(sync_result_builder) # Presence data is included if the server has it enabled and not filtered out. - include_presence_data = bool( - self.hs_config.server.presence_enabled - and not sync_config.filter_collection.blocks_all_presence() + presence_enabled = bool(self.hs_config.server.presence_enabled) + if not presence_enabled and since_token is not None: + # Even with presence disabled we send down any presence updates the + # client hasn't yet seen, so that the "mark everyone as offline" + # updates written when presence was disabled reach clients that + # would otherwise show the old presence states forever. The stream + # doesn't advance while presence is disabled, so once clients have + # caught up this check stops any further presence work. + presence_enabled = ( + since_token.presence_key < sync_result_builder.now_token.presence_key + ) + include_presence_data = ( + presence_enabled and not sync_config.filter_collection.blocks_all_presence() ) # Device list updates are sent if a since token is provided. include_device_list_updates = bool(since_token and since_token.device_list_key) diff --git a/synapse/rest/client/capabilities.py b/synapse/rest/client/capabilities.py index 2be5f5849d..4ddaaeda74 100644 --- a/synapse/rest/client/capabilities.py +++ b/synapse/rest/client/capabilities.py @@ -109,6 +109,11 @@ class CapabilitiesRestServlet(RestServlet): "capabilities" ]["m.profile_fields"] + response["capabilities"]["org.matrix.msc4140.delayed_events"] = { + "max_delay_ms": self.config.server.max_event_delay_duration.as_millis(), + "max_scheduled": self.config.server.max_delayed_events_per_user, + } + if self.config.experimental.msc4267_enabled: response["capabilities"]["org.matrix.msc4267.forget_forced_upon_leave"] = { "enabled": self.config.room.forget_on_leave, diff --git a/synapse/rest/client/read_marker.py b/synapse/rest/client/read_marker.py index 874e7487bf..8e0f2a2e7a 100644 --- a/synapse/rest/client/read_marker.py +++ b/synapse/rest/client/read_marker.py @@ -23,6 +23,7 @@ import logging from typing import TYPE_CHECKING from synapse.api.constants import ReceiptTypes +from synapse.api.errors import Codes, SynapseError from synapse.http.server import HttpServer from synapse.http.servlet import RestServlet, parse_json_object_from_request from synapse.http.site import SynapseRequest @@ -66,6 +67,21 @@ class ReadMarkerRestServlet(RestServlet): body = parse_json_object_from_request(request) unrecognized_types = set(body.keys()) - self._known_receipt_types + + if self.config.experimental.msc4446_enabled: + allow_backward = body.get("com.beeper.allow_backward", False) + if not isinstance(allow_backward, bool): + raise SynapseError( + 400, + "com.beeper.allow_backward must be a boolean.", + Codes.INVALID_PARAM, + ) + + # Prevent considering the `allow_backward` field as a receipt type. + unrecognized_types -= {"com.beeper.allow_backward"} + else: + allow_backward = False + if unrecognized_types: # It's fine if there are unrecognized receipt types, but let's log # it to help debug clients that have typoed the receipt type. @@ -86,6 +102,7 @@ class ReadMarkerRestServlet(RestServlet): room_id, user_id=requester.user.to_string(), event_id=event_id, + allow_backward=allow_backward, ) else: await self.receipts_handler.received_client_receipt( diff --git a/synapse/rest/client/receipts.py b/synapse/rest/client/receipts.py index d3a43537bb..949a1e64ad 100644 --- a/synapse/rest/client/receipts.py +++ b/synapse/rest/client/receipts.py @@ -20,6 +20,7 @@ # import logging +from http import HTTPStatus from typing import TYPE_CHECKING from synapse.api.constants import MAIN_TIMELINE, ReceiptTypes @@ -50,6 +51,7 @@ class ReceiptRestServlet(RestServlet): self.read_marker_handler = hs.get_read_marker_handler() self.presence_handler = hs.get_presence_handler() self._main_store = hs.get_datastores().main + self._msc4446_enabled = hs.config.experimental.msc4446_enabled self._known_receipt_types = { ReceiptTypes.READ, @@ -73,6 +75,25 @@ class ReceiptRestServlet(RestServlet): body = parse_json_object_from_request(request) + if self._msc4446_enabled: + allow_backward = body.get("com.beeper.allow_backward", False) + if not isinstance(allow_backward, bool): + raise SynapseError( + HTTPStatus.BAD_REQUEST, + "com.beeper.allow_backward must be a boolean.", + Codes.INVALID_PARAM, + ) + + if allow_backward and receipt_type != ReceiptTypes.FULLY_READ: + raise SynapseError( + HTTPStatus.BAD_REQUEST, + "com.beeper.allow_backward is only allowed to be true for " + f"{ReceiptTypes.FULLY_READ}.", + Codes.INVALID_PARAM, + ) + else: + allow_backward = False + # Pull the thread ID, if one exists. thread_id = None if "thread_id" in body: @@ -108,6 +129,7 @@ class ReceiptRestServlet(RestServlet): room_id, user_id=requester.user.to_string(), event_id=event_id, + allow_backward=allow_backward, ) else: await self.receipts_handler.received_client_receipt( diff --git a/synapse/rest/client/room.py b/synapse/rest/client/room.py index 36f638e236..951b4ff08b 100644 --- a/synapse/rest/client/room.py +++ b/synapse/rest/client/room.py @@ -84,6 +84,7 @@ from synapse.types import JsonDict, Requester, StreamToken, ThirdPartyInstanceID from synapse.types.state import StateFilter from synapse.util.cancellation import cancellable from synapse.util.clock import Clock +from synapse.util.duration import Duration from synapse.util.events import generate_fake_event_id from synapse.util.stringutils import parse_and_validate_server_name @@ -215,7 +216,6 @@ class RoomStateEventRestServlet(RestServlet): self.auth = hs.get_auth() self.clock = hs.get_clock() self._event_serializer = hs.get_event_client_serializer() - self._max_event_delay_ms = hs.config.server.max_event_delay_ms self._spam_checker_module_callbacks = hs.get_module_api_callbacks().spam_checker self._msc4354_enabled = hs.config.experimental.msc4354_enabled @@ -343,7 +343,7 @@ class RoomStateEventRestServlet(RestServlet): if self._msc4354_enabled: sticky_duration_ms = parse_integer(request, StickyEvent.QUERY_PARAM_NAME) - delay = _parse_request_delay(request, self._max_event_delay_ms) + delay = _parse_request_for_delayed_event_delay(request) if delay is not None: delay_id = await self.delayed_events_handler.add( requester, @@ -416,7 +416,6 @@ class RoomSendEventRestServlet(TransactionRestServlet): self.event_creation_handler = hs.get_event_creation_handler() self.delayed_events_handler = hs.get_delayed_events_handler() self.auth = hs.get_auth() - self._max_event_delay_ms = hs.config.server.max_event_delay_ms self._msc4354_enabled = hs.config.experimental.msc4354_enabled def register(self, http_server: HttpServer) -> None: @@ -442,7 +441,7 @@ class RoomSendEventRestServlet(TransactionRestServlet): if self._msc4354_enabled: sticky_duration_ms = parse_integer(request, StickyEvent.QUERY_PARAM_NAME) - delay = _parse_request_delay(request, self._max_event_delay_ms) + delay = _parse_request_for_delayed_event_delay(request) if delay is not None: delay_id = await self.delayed_events_handler.add( requester, @@ -515,47 +514,20 @@ class RoomSendEventRestServlet(TransactionRestServlet): ) -def _parse_request_delay( - request: SynapseRequest, - max_delay: int | None, -) -> int | None: +def _parse_request_for_delayed_event_delay(request: SynapseRequest) -> Duration | None: """Parses from the request string the delay parameter for delayed event requests, and checks it for correctness. Args: request: the twisted HTTP request. - max_delay: the maximum allowed value of the delay parameter, - or None if no delay parameter is allowed. Returns: The value of the requested delay, or None if it was absent. Raises: - SynapseError: if the delay parameter is present and forbidden, - or if it exceeds the maximum allowed value. + SynapseError: if the delay parameter is present and invalid. """ - delay = parse_integer(request, "org.matrix.msc4140.delay") - if delay is None: - return None - if max_delay is None: - raise SynapseError( - HTTPStatus.BAD_REQUEST, - "Delayed events are not supported on this server", - Codes.UNKNOWN, - { - "org.matrix.msc4140.errcode": "M_MAX_DELAY_UNSUPPORTED", - }, - ) - if delay > max_delay: - raise SynapseError( - HTTPStatus.BAD_REQUEST, - "The requested delay exceeds the allowed maximum.", - Codes.UNKNOWN, - { - "org.matrix.msc4140.errcode": "M_MAX_DELAY_EXCEEDED", - "org.matrix.msc4140.max_delay": max_delay, - }, - ) - return delay + delay_ms = parse_integer(request, "org.matrix.msc4140.delay") + return Duration(milliseconds=delay_ms) if delay_ms is not None else None # TODO: Needs unit testing for room ID + alias joins diff --git a/synapse/storage/databases/main/delayed_events.py b/synapse/storage/databases/main/delayed_events.py index 1727f589e2..bb512611e4 100644 --- a/synapse/storage/databases/main/delayed_events.py +++ b/synapse/storage/databases/main/delayed_events.py @@ -17,7 +17,7 @@ from typing import TYPE_CHECKING, NewType import attr -from synapse.api.errors import NotFoundError +from synapse.api.errors import LimitExceededError, NotFoundError from synapse.storage._base import SQLBaseStore, db_to_json from synapse.storage.database import ( DatabasePool, @@ -28,6 +28,7 @@ from synapse.storage.database import ( from synapse.storage.engines import PostgresEngine from synapse.types import JsonDict, RoomID from synapse.util import stringutils +from synapse.util.duration import Duration from synapse.util.json import json_encoder if TYPE_CHECKING: @@ -122,20 +123,84 @@ class DelayedEventsStore(SQLBaseStore): state_key: str | None, origin_server_ts: int | None, content: JsonDict, - delay: int, + delay: Duration, sticky_duration_ms: int | None, + limit: int, ) -> tuple[DelayID, Timestamp]: """ Inserts a new delayed event in the DB. + Args: + user_localpart: The localpart of the requester of the delayed event, who will be its owner. + device_id: The device ID of the requester. + creation_ts: The timestamp of when the request to add the delayed event was made. + room_id: The ID of the room where the event should be sent to. + event_type: The type of event to be sent. + state_key: The state key of the event to be sent, or None if it is not a state event. + origin_server_ts: The custom timestamp to send the event with. + If None, the timestamp will be the actual time when the event is sent. + content: The content of the event to be sent. + delay: How long to wait before automatically sending the event. + sticky_duration_ms: If an MSC4354 sticky event: the sticky duration (in milliseconds). + The event will be attempted to be reliably delivered to clients and remote servers + during its sticky period. + limit: The maximum number of delayed events the DB may store for the given requester. + Must be greater than 0. Returns: The generated ID assigned to the added delayed event, and the send time of the next delayed event to be sent, which is either the event just added or one added earlier. + + Raises: + LimitExceededError: if the DB has reached the limit of + how many delayed events it may store for the given requester. + AssertionError: if the limit is not greater than 0. """ + assert limit > 0, "limit must be greater than 0" + delay_id = _generate_delay_id() - send_ts = Timestamp(creation_ts + delay) + delay_ms = delay.as_millis() + send_ts = creation_ts + delay_ms def add_delayed_event_txn(txn: LoggingTransaction) -> Timestamp: + num_existing: int = self.db_pool.simple_select_one_onecol_txn( + txn, + table="delayed_events", + keyvalues={"user_localpart": user_localpart}, + retcol="COUNT(*)", + ) + if num_existing >= limit: + # Find the send_ts threshold that will bring the queue back under the limit. + # When the amount of existing delayed events has reached the limit, + # this will be the send time of the next delayed event to be sent. + # When the amount has exceeded the limit (e.g., due to config changes), + # this will be the send time of the delayed event that will be sent + # once all earlier events that exceed the limit have been sent. + # + # FIXME: Remove "AS subquery" after dropping support for PostgreSQL <16 + txn.execute( + """ + SELECT MAX(send_ts) FROM ( + SELECT * FROM delayed_events + WHERE user_localpart = ? + ORDER BY send_ts ASC + LIMIT ? + ) AS subquery + """, + ( + user_localpart, + num_existing - limit + 1, + ), + ) + row = txn.fetchone() + assert row + retry_after_ms = row[0] - self.clock.time_msec() + err = LimitExceededError( + limiter_name="add_delayed_event", + retry_after_ms=retry_after_ms if retry_after_ms > 0 else None, + ) + err.msg = "The maximum number of delayed events has been reached." + raise err + self.db_pool.simple_insert_txn( txn, table="delayed_events", @@ -143,7 +208,7 @@ class DelayedEventsStore(SQLBaseStore): "delay_id": delay_id, "user_localpart": user_localpart, "device_id": device_id, - "delay": delay, + "delay": delay_ms, "send_ts": send_ts, "room_id": room_id, "event_type": event_type, diff --git a/synapse/storage/databases/main/devices.py b/synapse/storage/databases/main/devices.py index b293b47fea..0e4c8ac491 100644 --- a/synapse/storage/databases/main/devices.py +++ b/synapse/storage/databases/main/devices.py @@ -2561,9 +2561,14 @@ class DeviceWorkerStore(RoomMemberWorkerStore, EndToEndKeyWorkerStore): # We default to 0 here as that is less than all possible stream IDs. min_stream_id = 0 - def prune_device_lists_changes_in_room_txn(txn: LoggingTransaction) -> int: - nonlocal min_stream_id - + def prune_device_lists_changes_in_room_txn( + txn: LoggingTransaction, min_stream_id: int + ) -> tuple[int, int]: + """ + Returns tuple of: + - number of rows deleted + - new `min_stream_id` for the next iteration + """ delete_sql = """ DELETE FROM device_lists_changes_in_room WHERE stream_id IN ( @@ -2596,13 +2601,14 @@ class DeviceWorkerStore(RoomMemberWorkerStore, EndToEndKeyWorkerStore): updatevalues={"stream_id": min_stream_id}, ) - return num_deleted + return num_deleted, min_stream_id progress_num_rows_deleted = 0 while True: - batch_deleted = await self.db_pool.runInteraction( + batch_deleted, min_stream_id = await self.db_pool.runInteraction( "prune_device_lists_changes_in_room", prune_device_lists_changes_in_room_txn, + min_stream_id, ) finished = batch_deleted < PRUNE_DEVICE_LISTS_BATCH_SIZE diff --git a/synapse/storage/databases/main/roommember.py b/synapse/storage/databases/main/roommember.py index 6106e55014..68798884dd 100644 --- a/synapse/storage/databases/main/roommember.py +++ b/synapse/storage/databases/main/roommember.py @@ -998,12 +998,22 @@ class RoomMemberWorkerStore(EventsWorkerStore, CacheInvalidationWorkerStore): return {u for u, share_room in user_dict.items() if share_room} - async def get_users_who_share_room_with_user(self, user_id: str) -> set[str]: - """Returns the set of users who share a room with `user_id`""" + async def get_users_who_share_room_with_user( + self, user_id: str, excluded_rooms: AbstractSet[str] = frozenset() + ) -> set[str]: + """Returns the set of users who share a room with `user_id`. + + Args: + user_id: The user to find the co-occupants of. + excluded_rooms: Rooms which should not, on their own, count as a + shared room. + """ room_ids = await self.get_rooms_for_user(user_id) user_who_share_room: set[str] = set() for room_id in room_ids: + if room_id in excluded_rooms: + continue user_ids = await self.get_users_in_room(room_id) user_who_share_room.update(user_ids) diff --git a/synapse/storage/schema/main/delta/88/01_add_delayed_events.sql b/synapse/storage/schema/main/delta/88/01_add_delayed_events.sql index 78ba5129af..4abe0ccaf4 100644 --- a/synapse/storage/schema/main/delta/88/01_add_delayed_events.sql +++ b/synapse/storage/schema/main/delta/88/01_add_delayed_events.sql @@ -22,6 +22,8 @@ CREATE TABLE delayed_events ( state_key TEXT, origin_server_ts BIGINT, content bytea NOT NULL, + -- is_processed = TRUE means that the work of sending the delayed event has begun. + -- Once the send is complete, the delayed event is removed from this table. is_processed BOOLEAN NOT NULL DEFAULT FALSE, PRIMARY KEY (user_localpart, delay_id) ); diff --git a/tests/config/test_server.py b/tests/config/test_server.py index d3c59ae14c..41ea8fb5b1 100644 --- a/tests/config/test_server.py +++ b/tests/config/test_server.py @@ -18,11 +18,15 @@ # # + +from typing import Any + import yaml from synapse.config._base import ConfigError, RootConfig from synapse.config.homeserver import HomeServerConfig from synapse.config.server import ServerConfig, generate_ip_set, is_threepid_reserved +from synapse.types import JsonDict from tests import unittest @@ -189,6 +193,54 @@ class ServerConfigTestCase(unittest.TestCase): self.assertEqual(conf["listeners"], expected_listeners) + def test_max_delayed_events_enforces_positive(self) -> None: + """ + Test that the configured maximum allowed delay must be a positive value if set, + as per documentation + """ + + def generate_config(value: int) -> JsonDict: + return {"max_event_delay_duration": value} + + _read_config(generate_config(1)) + + with self.assertRaises(ConfigError): + _read_config(generate_config(0)) + + with self.assertRaises(ConfigError): + _read_config(generate_config(-1)) + + def test_max_delayed_events_per_user_enforces_non_negative_int(self) -> None: + """ + Test that the configured maximum number of delayed events must be a non-negative value if set, + as a negative limit can never be satisfied + """ + + def generate_config(value: Any) -> JsonDict: + return { + "experimental_features": {"msc4140_max_delayed_events_per_user": value} + } + + for allowed_value in (0, 1): + _read_config(generate_config(allowed_value)) + + for disallowed_value in (-1, 0.5): + with self.assertRaises(ConfigError): + _read_config(generate_config(disallowed_value)) + + +def _read_config(config_values: JsonDict) -> None: + ServerConfig(RootConfig()).read_config( + yaml.safe_load( + HomeServerConfig().generate_config( + config_dir_path="CONFDIR", + data_dir_path="/data_dir_path", + server_name="che.org", + ) + ) + | config_values + ) + class GenerateIpSetTestCase(unittest.TestCase): def test_empty(self) -> None: diff --git a/tests/handlers/test_presence.py b/tests/handlers/test_presence.py index 624754d0cc..a562050842 100644 --- a/tests/handlers/test_presence.py +++ b/tests/handlers/test_presence.py @@ -19,7 +19,7 @@ # # import itertools -from typing import cast +from typing import Any, cast from unittest.mock import Mock, call from parameterized import parameterized @@ -50,6 +50,9 @@ from synapse.handlers.presence import ( FEDERATION_PING_INTERVAL, FEDERATION_TIMEOUT, PresenceHandler, + WorkerPresenceHandler, + get_interested_parties, + get_interested_remotes, handle_timeout, handle_update, ) @@ -999,6 +1002,125 @@ class PresenceHandlerInitTestCase(unittest.HomeserverTestCase): ) self.assertEqual(state.state, sync_state) + @unittest.override_config({"presence": {"enabled": False}}) + def test_restored_presence_flushed_offline_when_presence_disabled(self) -> None: + """If presence is disabled, any non-offline presence states left in the + database from when presence was enabled should be marked as offline at + startup, and the updates streamed out to clients. + """ + main_store = self.hs.get_datastores().main + before_token = main_store.get_current_presence_token() + + # Get the handler, which schedules the startup flush. + presence_handler = self.hs.get_presence_handler() + + # Fire pending `call_when_running` hooks and let the flush complete. + self.reactor.run() + self.reactor.advance(0) + + # The user should now be offline, both in memory and in the database. + state = self.get_success( + presence_handler.get_state(UserID.from_string(self.user_id)) + ) + self.assertEqual(state.state, PresenceState.OFFLINE) + + db_state = self.get_success(main_store.get_presence_for_users([self.user_id]))[ + self.user_id + ] + self.assertEqual(db_state.state, PresenceState.OFFLINE) + + # The flush must advance the presence stream so that syncing clients + # are sent the offline updates. + self.assertGreater(main_store.get_current_presence_token(), before_token) + + +class PresenceDisabledSyncTestCase(unittest.HomeserverTestCase): + """Tests that stale presence states left over from when presence was + enabled reach clients over /sync, and that the startup flush marks them as + offline and sends the offline updates down /sync too. + """ + + servlets = [ + admin.register_servlets, + login.register_servlets, + room.register_servlets, + sync.register_servlets, + ] + + @unittest.override_config({"presence": {"enabled": False}}) + def test_stale_presence_flushed_offline_and_sent_on_sync(self) -> None: + user1 = self.register_user("alice", "pass") + user1_tok = self.login(user1, "pass") + user2 = self.register_user("bob", "pass") + user2_tok = self.login(user2, "pass") + + room_id = self.helper.create_room_as(user1, tok=user1_tok) + self.helper.join(room_id, user2, tok=user2_tok) + + channel = self.make_request("GET", "/sync", access_token=user2_tok) + self.assertEqual(channel.code, 200, channel.json_body) + next_batch = channel.json_body["next_batch"] + + # Seed a stale online presence state for user1, left over from when + # presence was enabled: in the database, and in the presence handler's + # in-memory state (which at startup is preloaded from the database). + now = self.clock.time_msec() + stale_state = UserPresenceState( + user_id=user1, + state=PresenceState.ONLINE, + last_active_ts=now, + last_federation_update_ts=now, + last_user_sync_ts=now, + status_msg=None, + currently_active=True, + ) + main_store = self.hs.get_datastores().main + self.get_success(main_store.update_presence([stale_state])) + + presence_handler = self.hs.get_presence_handler() + assert isinstance(presence_handler, PresenceHandler) + presence_handler.user_to_current_state[user1] = stale_state + + # The stale state comes down user2's incremental sync, even though + # presence is disabled. + channel = self.make_request( + "GET", f"/sync?since={next_batch}", access_token=user2_tok + ) + self.assertEqual(channel.code, 200, channel.json_body) + presence_events = channel.json_body["presence"]["events"] + self.assertEqual( + [(e["sender"], e["content"]["presence"]) for e in presence_events], + [(user1, PresenceState.ONLINE)], + ) + next_batch = channel.json_body["next_batch"] + + # Run the startup flush, as scheduled when the presence writer starts + # up with presence disabled. + self.get_success(presence_handler._mark_stale_presence_as_offline()) + + # The stale state should have been marked offline in the database... + db_state = self.get_success(main_store.get_presence_for_users([user1]))[user1] + self.assertEqual(db_state.state, PresenceState.OFFLINE) + + # ... and the offline update also comes down user2's sync. + channel = self.make_request( + "GET", f"/sync?since={next_batch}", access_token=user2_tok + ) + self.assertEqual(channel.code, 200, channel.json_body) + presence_events = channel.json_body["presence"]["events"] + self.assertEqual( + [(e["sender"], e["content"]["presence"]) for e in presence_events], + [(user1, PresenceState.OFFLINE)], + ) + + # Once caught up, further syncs include no presence. + next_batch = channel.json_body["next_batch"] + channel = self.make_request( + "GET", f"/sync?since={next_batch}", access_token=user2_tok + ) + self.assertEqual(channel.code, 200, channel.json_body) + self.assertEqual(channel.json_body.get("presence", {}).get("events", []), []) + # Timer values used by `PresenceConfigurableTimersTestCase`, all larger than # the corresponding defaults. @@ -2320,3 +2442,391 @@ class PresenceJoinTestCase(unittest.HomeserverTestCase): ) return event + + +class PresenceExcludeRoomsTestCase(unittest.HomeserverTestCase): + """Tests that `exclude_rooms_from_presence` stops presence being routed + between users solely because they share an excluded room.""" + + servlets = [ + admin.register_servlets, + login.register_servlets, + room.register_servlets, + ] + + def prepare(self, reactor: MemoryReactor, clock: Clock, hs: HomeServer) -> None: + self.hs = hs + self.store = hs.get_datastores().main + self.presence_router = hs.get_presence_router() + self.presence_handler = hs.get_presence_handler() + + self.user1 = self.register_user("user1", "pass") + self.token1 = self.login("user1", "pass") + self.user2 = self.register_user("user2", "pass") + self.token2 = self.login("user2", "pass") + + def test_excluded_rooms_not_routed(self) -> None: + # Two rooms that user1 is joined to. + excluded_room = self.helper.create_room_as(self.user1, tok=self.token1) + shared_room = self.helper.create_room_as(self.user1, tok=self.token1) + + state = UserPresenceState.default(self.user1) + + # Without any exclusions both rooms are interested in user1's presence. + room_ids_to_states, users_to_states = self.get_success( + get_interested_parties(self.store, self.presence_router, [state]) + ) + self.assertIn(excluded_room, room_ids_to_states) + self.assertIn(shared_room, room_ids_to_states) + + # Excluding one room drops it as an interested party, but the other + # (non-excluded) room still routes presence... + room_ids_to_states, users_to_states = self.get_success( + get_interested_parties( + self.store, + self.presence_router, + [state], + frozenset({excluded_room}), + ) + ) + self.assertNotIn(excluded_room, room_ids_to_states) + self.assertIn(shared_room, room_ids_to_states) + + # ...and the user always receives their own presence, even when all of + # their rooms are excluded. + room_ids_to_states, users_to_states = self.get_success( + get_interested_parties( + self.store, + self.presence_router, + [state], + frozenset({excluded_room, shared_room}), + ) + ) + self.assertNotIn(excluded_room, room_ids_to_states) + self.assertNotIn(shared_room, room_ids_to_states) + self.assertIn(self.user1, users_to_states) + + @override_config({"exclude_rooms_from_presence": ["!excluded:test"]}) + def test_config_populates_handler(self) -> None: + """The config option should be plumbed through to the presence handler + and the presence event source as a frozenset.""" + self.assertEqual( + self.presence_handler._rooms_to_exclude_from_presence, + frozenset({"!excluded:test"}), + ) + + event_source = self.hs.get_event_sources().sources.presence + self.assertEqual( + event_source._rooms_to_exclude_from_presence, + frozenset({"!excluded:test"}), + ) + + def test_is_visible_respects_excluded_rooms(self) -> None: + """`is_visible` (which drives the read side of /sync) should not + consider two users to share presence solely via an excluded room.""" + user1 = UserID.from_string(self.user1) + user2 = UserID.from_string(self.user2) + + # A single shared room: the two users can see each other's presence. + excluded_room = self.helper.create_room_as(self.user1, tok=self.token1) + self.helper.join(excluded_room, self.user2, tok=self.token2) + + self.assertTrue( + self.get_success(self.presence_handler.is_visible(user2, user1)) + ) + + # Excluding the only shared room hides presence between them. + self.presence_handler._rooms_to_exclude_from_presence = frozenset( + {excluded_room} + ) + self.assertFalse( + self.get_success(self.presence_handler.is_visible(user2, user1)) + ) + + # But a second, non-excluded shared room restores visibility. + shared_room = self.helper.create_room_as(self.user1, tok=self.token1) + self.helper.join(shared_room, self.user2, tok=self.token2) + self.assertTrue( + self.get_success(self.presence_handler.is_visible(user2, user1)) + ) + + def test_get_interested_remotes_respects_excluded_rooms(self) -> None: + """The federation fan-out side (`get_interested_remotes`) must not route + presence to servers reached solely via an excluded room.""" + excluded_room = self.helper.create_room_as(self.user1, tok=self.token1) + state = UserPresenceState.default(self.user1) + + def hosts_for(excluded: frozenset) -> set: + result = self.get_success( + get_interested_remotes( + self.store, self.presence_router, [state], excluded + ) + ) + hosts: set[str] = set() + for room_hosts, _ in result: + hosts.update(room_hosts) + return hosts + + # The local server is a host in the room (all members are local here), + # so presence would be routed there... + self.assertIn("test", hosts_for(frozenset())) + # ...but excluding the only room removes it as a source of destinations. + self.assertNotIn("test", hosts_for(frozenset({excluded_room}))) + + +class PresenceGetNewEventsStreamTestCase(unittest.HomeserverTestCase): + """Tests the incremental (`from_key`) branch of + `PresenceEventSource.get_new_events`, which decides which updated users are + interesting to the syncing user by intersecting their cached room sets. + """ + + servlets = [ + admin.register_servlets, + login.register_servlets, + room.register_servlets, + ] + + def prepare(self, reactor: MemoryReactor, clock: Clock, hs: HomeServer) -> None: + self.presence_handler = hs.get_presence_handler() + self.event_source = hs.get_event_sources().sources.presence + + self.user1 = self.register_user("user1", "pass") + self.token1 = self.login("user1", "pass") + self.user2 = self.register_user("user2", "pass") + self.token2 = self.login("user2", "pass") + self.user3 = self.register_user("user3", "pass") + self.token3 = self.login("user3", "pass") + + def _set_presence(self, user_id: str, state: str = "online") -> None: + self.get_success( + self.presence_handler.set_state( + UserID.from_string(user_id), "dev", {"presence": state} + ) + ) + + def _updated_users_seen_by(self, user_id: str, from_key: int) -> set[str]: + states, _ = self.get_success( + self.event_source.get_new_events( + user=UserID.from_string(user_id), from_key=from_key + ) + ) + return {state.user_id for state in states} + + def test_incremental_interest(self) -> None: + """A syncing user sees updates from users they share a room with (and + themselves), but not from strangers.""" + shared_room = self.helper.create_room_as(self.user1, tok=self.token1) + self.helper.join(shared_room, self.user2, tok=self.token2) + # user3 is in an unrelated room. + self.helper.create_room_as(self.user3, tok=self.token3) + + from_key = self.event_source.get_current_key() + self._set_presence(self.user1) + self._set_presence(self.user2) + self._set_presence(self.user3) + + seen = self._updated_users_seen_by(self.user2, from_key) + self.assertIn(self.user1, seen) + self.assertIn(self.user2, seen) # always sees own updates + self.assertNotIn(self.user3, seen) + + def test_incremental_interest_excluded_room(self) -> None: + """Sharing only an excluded room does not make an updated user + interesting; sharing an additional normal room does.""" + excluded_room = self.helper.create_room_as(self.user1, tok=self.token1) + self.helper.join(excluded_room, self.user2, tok=self.token2) + + self.event_source._rooms_to_exclude_from_presence = frozenset({excluded_room}) + + from_key = self.event_source.get_current_key() + self._set_presence(self.user1, "online") + seen = self._updated_users_seen_by(self.user2, from_key) + self.assertNotIn(self.user1, seen) + + # A second, non-excluded shared room restores interest. (Use a + # different presence state, as repeating the same one would not + # generate a new update.) + shared_room = self.helper.create_room_as(self.user1, tok=self.token1) + self.helper.join(shared_room, self.user2, tok=self.token2) + + from_key = self.event_source.get_current_key() + self._set_presence(self.user1, "unavailable") + seen = self._updated_users_seen_by(self.user2, from_key) + self.assertIn(self.user1, seen) + + +class WorkerPresenceThrottleTestCase(BaseMultiWorkerStreamTestCase): + """Tests that sync workers suppress the per-sync-request presence updates + that the presence writer would discard anyway, while relaying genuine + state changes immediately.""" + + def prepare(self, reactor: MemoryReactor, clock: Clock, hs: HomeServer) -> None: + self.user_id = f"@throttled:{hs.config.server.server_name}" + self.user_id_obj = UserID.from_string(self.user_id) + self.device_id = "dev-1" + # In this test setup the main process is the presence writer. + self.writer_handler = hs.get_presence_handler() + + def _make_sync_worker(self) -> tuple[Any, list, list]: + """Create a sync worker whose proxied presence calls are recorded.""" + worker = self.make_worker_hs( + "synapse.app.generic_worker", {"worker_name": "synchrotron"} + ) + presence = worker.get_presence_handler() + assert isinstance(presence, WorkerPresenceHandler) + + set_state_calls: list = [] + bump_calls: list = [] + real_set_state = presence._set_state_client + real_bump = presence._bump_active_client + + async def recording_set_state(**kwargs: Any) -> Any: + set_state_calls.append(kwargs) + return await real_set_state(**kwargs) + + async def recording_bump(**kwargs: Any) -> Any: + bump_calls.append(kwargs) + return await real_bump(**kwargs) + + presence._set_state_client = recording_set_state + presence._bump_active_client = recording_bump + return presence, set_state_calls, bump_calls + + def _sync(self, presence: Any, state: str = PresenceState.ONLINE) -> Any: + # Note: `get_success` only advances the fake clock by tiny epsilon steps + # while pumping the replication traffic; the throttle window is + # time-sensitive and these tests advance time explicitly. + return self.get_success( + presence.user_syncing(self.user_id, self.device_id, True, state), + ) + + def test_repeated_syncs_are_throttled(self) -> None: + presence, set_state_calls, _ = self._make_sync_worker() + + # Several syncs in quick succession only relay one set_state. + for _ in range(3): + self._sync(presence) + self.assertEqual(len(set_state_calls), 1) + + # The user did come online on the writer. + state = self.get_success(self.writer_handler.get_state(self.user_id_obj)) + self.assertEqual(state.state, PresenceState.ONLINE) + + # Once the relay window has passed, the next sync relays again. + self.reactor.advance(presence._sync_presence_relay_interval / 1000 + 1) + self._sync(presence) + self.assertEqual(len(set_state_calls), 2) + + def test_state_changes_are_relayed_immediately(self) -> None: + presence, set_state_calls, _ = self._make_sync_worker() + + self._sync(presence, PresenceState.ONLINE) + self._sync(presence, PresenceState.UNAVAILABLE) + self._sync(presence, PresenceState.ONLINE) + self.assertEqual(len(set_state_calls), 3) + + # A repeat of the current state within the window is suppressed. + self._sync(presence, PresenceState.ONLINE) + self.assertEqual(len(set_state_calls), 3) + + def test_resends_after_device_stops_syncing(self) -> None: + """After a USER_SYNC stop is sent the writer may time the user out, so + a device that reconnects within the window must be relayed afresh.""" + presence, set_state_calls, _ = self._make_sync_worker() + + with self._sync(presence): + pass + self.assertEqual(len(set_state_calls), 1) + + # Wait for the going-offline grace period to elapse: USER_SYNC stop is + # sent and the throttle entry evicted. The writer then times the user + # out to offline. Advance in steps (rather than one jump) so the + # replicated stop command is delivered before the writer's timeout + # loop fires (which only starts 30s after startup). + for _ in range(4): + self.reactor.advance(12) + state = self.get_success(self.writer_handler.get_state(self.user_id_obj)) + self.assertEqual(state.state, PresenceState.OFFLINE) + + # Reconnecting relays the state immediately, even though the presence + # value is unchanged from the last relayed one. + self._sync(presence) + self.assertEqual(len(set_state_calls), 2) + state = self.get_success(self.writer_handler.get_state(self.user_id_obj)) + self.assertEqual(state.state, PresenceState.ONLINE) + + def test_bumps_are_throttled(self) -> None: + presence, set_state_calls, bump_calls = self._make_sync_worker() + + # While we recently relayed an online state, bumps are suppressed. + self._sync(presence, PresenceState.ONLINE) + for _ in range(3): + self.get_success( + presence.bump_presence_active_time(self.user_id_obj, self.device_id), + ) + self.assertEqual(len(bump_calls), 0) + + # After the window passes, a bump goes through (and then suppresses + # further bumps). + self.reactor.advance(presence._sync_presence_relay_interval / 1000 + 1) + for _ in range(2): + self.get_success( + presence.bump_presence_active_time(self.user_id_obj, self.device_id), + ) + self.assertEqual(len(bump_calls), 1) + + def test_explicit_set_state_always_relayed_and_resets(self) -> None: + """An explicit (non-sync) set_state is always relayed, and resets the + throttle so the next sync-driven update is relayed afresh.""" + presence, set_state_calls, _ = self._make_sync_worker() + + self._sync(presence, PresenceState.ONLINE) + self.assertEqual(len(set_state_calls), 1) + + # An explicit update of the same state within the window still goes + # through (it isn't sync-driven)... + self.get_success( + presence.set_state( + self.user_id_obj, + self.device_id, + {"presence": PresenceState.ONLINE}, + ), + ) + self.assertEqual(len(set_state_calls), 2) + + # ...and the following sync-driven update is relayed rather than + # suppressed, re-establishing the writer's sync timestamps. + self._sync(presence, PresenceState.ONLINE) + self.assertEqual(len(set_state_calls), 3) + + def test_bump_after_non_online_state_goes_through(self) -> None: + presence, set_state_calls, bump_calls = self._make_sync_worker() + + # The user is unavailable; a bump may un-idle them so it must not be + # suppressed. + self._sync(presence, PresenceState.UNAVAILABLE) + self.get_success( + presence.bump_presence_active_time(self.user_id_obj, self.device_id), + ) + self.assertEqual(len(bump_calls), 1) + + @override_config({"presence": {"sync_online_timeout": "12s"}}) + def test_relay_interval_scales_with_config(self) -> None: + """The throttle window is derived from the configurable presence timers, + so it stays comfortably below a lowered sync online timeout rather than + being a hardcoded 25s (which would make users flap).""" + presence, set_state_calls, _ = self._make_sync_worker() + + # 5/6 of min(12s sync online timeout, 60s default last-active + # granularity). + self.assertEqual(presence._sync_presence_relay_interval, 10 * 1000) + + # A repeat within the (now shorter) window is still suppressed... + self._sync(presence) + self._sync(presence) + self.assertEqual(len(set_state_calls), 1) + + # ...and once it passes, the next sync relays again. + self.reactor.advance(presence._sync_presence_relay_interval / 1000 + 1) + self._sync(presence) + self.assertEqual(len(set_state_calls), 2) diff --git a/tests/rest/client/test_account.py b/tests/rest/client/test_account.py index 86fab2fe25..d92bddc825 100644 --- a/tests/rest/client/test_account.py +++ b/tests/rest/client/test_account.py @@ -325,13 +325,8 @@ class PasswordResetTestCase(unittest.HomeserverTestCase): email = "test@example.com" client_secret = "foobar" - session_id = self._request_token( - email, - client_secret, - # The endpoint intentionally adds up to 1000ms of jitter to avoid - # leaking whether the email address is bound to an account. - timeout_ms=3000, - ) + + session_id = self._request_token(email, client_secret) self.assertIsNotNone(session_id) @@ -370,18 +365,21 @@ class PasswordResetTestCase(unittest.HomeserverTestCase): client_secret: str, ip: str = "127.0.0.1", next_link: str | None = None, - timeout_ms: int = 1000, ) -> str: body = {"client_secret": client_secret, "email": email, "send_attempt": 1} if next_link is not None: body["next_link"] = next_link + channel = self.make_request( "POST", b"account/password/email/requestToken", body, client_ip=ip, - timeout_ms=timeout_ms, + await_result=False, ) + # Note: The endpoint intentionally adds up to 1000ms of jitter to avoid + # leaking whether the email address is bound to an account. + channel.await_result(timeout_ms=1000) if channel.code != 200: raise HttpResponseException( diff --git a/tests/rest/client/test_capabilities.py b/tests/rest/client/test_capabilities.py index c28e0605b5..42926c4359 100644 --- a/tests/rest/client/test_capabilities.py +++ b/tests/rest/client/test_capabilities.py @@ -26,6 +26,7 @@ from synapse.api.room_versions import KNOWN_ROOM_VERSIONS from synapse.rest.client import capabilities, login from synapse.server import HomeServer from synapse.util.clock import Clock +from synapse.util.duration import Duration from tests import unittest from tests.unittest import override_config, skip_unless @@ -203,6 +204,43 @@ class CapabilitiesTestCase(unittest.HomeserverTestCase): ["avatar_url"], ) + def test_get_delayed_events_capabilities_default_config_msc4140(self) -> None: + access_token = self.login(self.localpart, self.password) + + channel = self.make_request("GET", self.url, access_token=access_token) + capabilities = channel.json_body["capabilities"] + + self.assertEqual(channel.code, HTTPStatus.OK) + self.assertEqual( + capabilities["org.matrix.msc4140.delayed_events"]["max_delay_ms"], 0 + ) + self.assertEqual( + capabilities["org.matrix.msc4140.delayed_events"]["max_scheduled"], 100 + ) + + @override_config( + { + "max_event_delay_duration": "24h", + "experimental_features": { + "msc4140_max_delayed_events_per_user": 50, + }, + } + ) + def test_get_delayed_events_capabilities_custom_config_msc4140(self) -> None: + access_token = self.login(self.localpart, self.password) + + channel = self.make_request("GET", self.url, access_token=access_token) + capabilities = channel.json_body["capabilities"] + + self.assertEqual(channel.code, HTTPStatus.OK) + self.assertEqual( + capabilities["org.matrix.msc4140.delayed_events"]["max_delay_ms"], + Duration(days=1).as_millis(), + ) + self.assertEqual( + capabilities["org.matrix.msc4140.delayed_events"]["max_scheduled"], 50 + ) + @override_config({"enable_3pid_changes": False}) def test_get_change_3pid_capabilities_3pid_disabled(self) -> None: """Test if change 3pid is disabled that the server responds it.""" diff --git a/tests/rest/client/test_read_marker.py b/tests/rest/client/test_read_marker.py index c8bb0da5e6..ad13d3607e 100644 --- a/tests/rest/client/test_read_marker.py +++ b/tests/rest/client/test_read_marker.py @@ -66,6 +66,19 @@ class ReadMarkerTestCase(unittest.HomeserverTestCase): self.store = self.hs.get_datastores().main self.clock = self.hs.get_clock() + def _get_fully_read_marker(self, room_id: str) -> str | None: + content = self.get_success( + self.store.get_account_data_for_room_and_type( + self.owner, + room_id, + "m.fully_read", + ) + ) + if content is None: + return None + + return content.get("event_id") + def test_send_read_marker(self) -> None: room_id = self.helper.create_room_as(self.owner, tok=self.owner_tok) @@ -98,6 +111,123 @@ class ReadMarkerTestCase(unittest.HomeserverTestCase): ) self.assertEqual(channel.code, 200, channel.result) + def test_send_read_marker_does_not_move_backwards_by_default(self) -> None: + room_id = self.helper.create_room_as(self.owner, tok=self.owner_tok) + + older_event_id = self.helper.send( + room_id=room_id, body="1", tok=self.owner_tok + )["event_id"] + newer_event_id = self.helper.send( + room_id=room_id, body="2", tok=self.owner_tok + )["event_id"] + + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": newer_event_id}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(room_id), newer_event_id) + + # Expected to be a no-op. + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": older_event_id}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(room_id), newer_event_id) + + @unittest.override_config({"experimental_features": {"msc4446_enabled": True}}) + def test_send_read_marker_can_move_backwards_with_opt_in(self) -> None: + room_id = self.helper.create_room_as(self.owner, tok=self.owner_tok) + + older_event_id = self.helper.send( + room_id=room_id, body="1", tok=self.owner_tok + )["event_id"] + newer_event_id = self.helper.send( + room_id=room_id, body="2", tok=self.owner_tok + )["event_id"] + + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": newer_event_id}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": older_event_id, "com.beeper.allow_backward": True}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(room_id), older_event_id) + + @unittest.override_config({"experimental_features": {"msc4446_enabled": True}}) + def test_send_read_marker_does_not_move_backwards_with_explicit_opt_out( + self, + ) -> None: + room_id = self.helper.create_room_as(self.owner, tok=self.owner_tok) + + older_event_id = self.helper.send( + room_id=room_id, body="1", tok=self.owner_tok + )["event_id"] + newer_event_id = self.helper.send( + room_id=room_id, body="2", tok=self.owner_tok + )["event_id"] + + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": newer_event_id}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + + # Expected to be a no-op. + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={ + "m.fully_read": older_event_id, + "com.beeper.allow_backward": False, + }, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(room_id), newer_event_id) + + def test_send_read_marker_ignores_opt_in_when_feature_disabled(self) -> None: + room_id = self.helper.create_room_as(self.owner, tok=self.owner_tok) + older_event_id = self.helper.send( + room_id=room_id, body="1", tok=self.owner_tok + )["event_id"] + newer_event_id = self.helper.send( + room_id=room_id, body="2", tok=self.owner_tok + )["event_id"] + + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": newer_event_id}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + + channel = self.make_request( + "POST", + f"/rooms/{room_id}/read_markers", + content={"m.fully_read": older_event_id, "com.beeper.allow_backward": True}, + access_token=self.owner_tok, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(room_id), newer_event_id) + def test_send_read_marker_missing_previous_event(self) -> None: """ Test moving a read marker from an event that previously existed but was diff --git a/tests/rest/client/test_receipts.py b/tests/rest/client/test_receipts.py index 3a6a869c54..0835eec6de 100644 --- a/tests/rest/client/test_receipts.py +++ b/tests/rest/client/test_receipts.py @@ -24,6 +24,7 @@ from twisted.internet.testing import MemoryReactor import synapse.rest.admin from synapse.api.constants import EduTypes, EventTypes, HistoryVisibility, ReceiptTypes +from synapse.api.errors import Codes from synapse.rest.client import login, receipts, room, sync from synapse.server import HomeServer from synapse.types import JsonDict @@ -44,6 +45,7 @@ class ReceiptsTestCase(unittest.HomeserverTestCase): def prepare(self, reactor: MemoryReactor, clock: Clock, hs: HomeServer) -> None: self.url = "/sync?since=%s" self.next_batch = "s0" + self.store = hs.get_datastores().main # Register the first user self.user_id = self.register_user("kermit", "monkey") @@ -59,6 +61,19 @@ class ReceiptsTestCase(unittest.HomeserverTestCase): # Join the second user self.helper.join(room=self.room_id, user=self.user2, tok=self.tok2) + def _get_fully_read_marker(self) -> str | None: + content = self.get_success( + self.store.get_account_data_for_room_and_type( + self.user2, + self.room_id, + ReceiptTypes.FULLY_READ, + ) + ) + if content is None: + return None + + return content.get("event_id") + def test_send_receipt(self) -> None: # Send a message. res = self.helper.send(self.room_id, body="hello", tok=self.tok) @@ -258,6 +273,126 @@ class ReceiptsTestCase(unittest.HomeserverTestCase): self.assertEqual(channel.code, HTTPStatus.BAD_REQUEST) self.assertEqual(channel.json_body["errcode"], "M_NOT_JSON", channel.json_body) + def test_fully_read_receipt_does_not_move_backwards_by_default(self) -> None: + older_event_id = self.helper.send(self.room_id, body="1", tok=self.tok)[ + "event_id" + ] + newer_event_id = self.helper.send(self.room_id, body="2", tok=self.tok)[ + "event_id" + ] + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{newer_event_id}", + {}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(), newer_event_id) + + # Expected to be a no-op. + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{older_event_id}", + {}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(), newer_event_id) + + @unittest.override_config({"experimental_features": {"msc4446_enabled": True}}) + def test_fully_read_receipt_can_move_backwards_with_opt_in(self) -> None: + older_event_id = self.helper.send(self.room_id, body="1", tok=self.tok)[ + "event_id" + ] + newer_event_id = self.helper.send(self.room_id, body="2", tok=self.tok)[ + "event_id" + ] + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{newer_event_id}", + {}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{older_event_id}", + {"com.beeper.allow_backward": True}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(), older_event_id) + + @unittest.override_config({"experimental_features": {"msc4446_enabled": True}}) + def test_fully_read_receipt_does_not_move_backwards_with_explicit_opt_out( + self, + ) -> None: + older_event_id = self.helper.send(self.room_id, body="1", tok=self.tok)[ + "event_id" + ] + newer_event_id = self.helper.send(self.room_id, body="2", tok=self.tok)[ + "event_id" + ] + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{newer_event_id}", + {}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + + # Expected to be a no-op. + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{older_event_id}", + {"com.beeper.allow_backward": False}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(), newer_event_id) + + @unittest.override_config({"experimental_features": {"msc4446_enabled": True}}) + def test_allow_backward_is_rejected_for_read_receipts(self) -> None: + event_id = self.helper.send(self.room_id, body="1", tok=self.tok)["event_id"] + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.READ}/{event_id}", + {"com.beeper.allow_backward": True}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, HTTPStatus.BAD_REQUEST, channel.result) + self.assertEqual(channel.json_body["errcode"], Codes.INVALID_PARAM) + + def test_allow_backward_is_ignored_when_feature_disabled(self) -> None: + older_event_id = self.helper.send(self.room_id, body="1", tok=self.tok)[ + "event_id" + ] + newer_event_id = self.helper.send(self.room_id, body="2", tok=self.tok)[ + "event_id" + ] + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{newer_event_id}", + {}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + + channel = self.make_request( + "POST", + f"/rooms/{self.room_id}/receipt/{ReceiptTypes.FULLY_READ}/{older_event_id}", + {"com.beeper.allow_backward": True}, + access_token=self.tok2, + ) + self.assertEqual(channel.code, 200, channel.result) + self.assertEqual(self._get_fully_read_marker(), newer_event_id) + def _get_read_receipt(self) -> JsonDict | None: """Syncs and returns the read receipt.""" diff --git a/tests/rest/client/test_register.py b/tests/rest/client/test_register.py index a9f3ac2462..1b04ea9d5c 100644 --- a/tests/rest/client/test_register.py +++ b/tests/rest/client/test_register.py @@ -753,10 +753,11 @@ class RegisterRestServletTestCase(unittest.HomeserverTestCase): "POST", b"register/email/requestToken", {"client_secret": "foobar", "email": email, "send_attempt": 1}, - # The endpoint intentionally adds up to 1000ms of jitter to avoid - # leaking whether the email address is already bound to an account. - timeout_ms=3000, + await_result=False, ) + # Note: The endpoint intentionally adds up to 1000ms of jitter to avoid + # leaking whether the email address is bound to an account. + channel.await_result(timeout_ms=1000) self.assertEqual(200, channel.code, channel.result) self.assertIsNotNone(channel.json_body.get("sid")) diff --git a/tests/rest/client/test_rooms.py b/tests/rest/client/test_rooms.py index 793a66f343..279e01e21a 100644 --- a/tests/rest/client/test_rooms.py +++ b/tests/rest/client/test_rooms.py @@ -61,6 +61,7 @@ from synapse.rest.client import ( from synapse.server import HomeServer from synapse.types import JsonDict, JsonMapping, RoomAlias, UserID, create_requester from synapse.util.clock import Clock +from synapse.util.duration import Duration from synapse.util.stringutils import random_string from tests import unittest @@ -2503,7 +2504,12 @@ class RoomDelayedEventTestCase(RoomBase): {}, ) self.assertEqual(HTTPStatus.BAD_REQUEST, channel.code, channel.result) - self.assertNotIn("org.matrix.msc4140.errcode", channel.json_body) + # Assert that the standard error response uses a valid errcode. + # The specific errcode is irrelevant for the purpose of this test. + self.assertIsInstance( + channel.json_body.get("errcode"), + str, + ) def test_delayed_event_unsupported_by_default(self) -> None: """Test that sending a delayed event is unsupported with the default config.""" @@ -2515,10 +2521,35 @@ class RoomDelayedEventTestCase(RoomBase): ).encode("ascii"), {"body": "test", "msgtype": "m.text"}, ) - self.assertEqual(HTTPStatus.BAD_REQUEST, channel.code, channel.result) + self.assertEqual(HTTPStatus.FORBIDDEN, channel.code, channel.result) self.assertEqual( - "M_MAX_DELAY_UNSUPPORTED", - channel.json_body.get("org.matrix.msc4140.errcode"), + Codes.FORBIDDEN, + channel.json_body.get("errcode"), + channel.json_body, + ) + + @unittest.override_config( + { + "max_event_delay_duration": "24h", + "experimental_features": { + "msc4140_max_delayed_events_per_user": 0, + }, + } + ) + def test_delayed_event_disabled_by_limit(self) -> None: + """Test that delayed events are disabled by configuring the per-user limit to 0.""" + channel = self.make_request( + "PUT", + ( + "rooms/%s/send/m.room.message/mid1?org.matrix.msc4140.delay=2000" + % self.room_id + ).encode("ascii"), + {"body": "test", "msgtype": "m.text"}, + ) + self.assertEqual(HTTPStatus.FORBIDDEN, channel.code, channel.result) + self.assertEqual( + Codes.FORBIDDEN, + channel.json_body.get("errcode"), channel.json_body, ) @@ -2533,13 +2564,178 @@ class RoomDelayedEventTestCase(RoomBase): ).encode("ascii"), {"body": "test", "msgtype": "m.text"}, ) - self.assertEqual(HTTPStatus.BAD_REQUEST, channel.code, channel.result) + self.assertEqual(HTTPStatus.FORBIDDEN, channel.code, channel.result) self.assertEqual( - "M_MAX_DELAY_EXCEEDED", - channel.json_body.get("org.matrix.msc4140.errcode"), + Codes.FORBIDDEN, + channel.json_body.get("errcode"), channel.json_body, ) + @unittest.override_config( + { + "max_event_delay_duration": "24h", + "experimental_features": { + "msc4140_max_delayed_events_per_user": 1, + }, + } + ) + def test_delayed_event_user_limit_reached(self) -> None: + """Test that users cannot have more delayed events scheduled at once than allowed.""" + # Disable rate-limits for this user. We want to specifically test the storage-based limit, not the request limits + self.get_success( + self.hs.get_datastores().main.set_ratelimit_for_user(self.user_id, 0, 0) + ) + + make_delayed_event_request = lambda: self.make_request( + "POST", + ( + "rooms/%s/send/m.room.message?org.matrix.msc4140.delay=15000" + % self.room_id + ).encode("ascii"), + {"body": "test", "msgtype": "m.text"}, + ) + # Send a delayed event to eat up the limit + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.OK, channel.code, channel.result) + + # Try to send another delayed event (we expect to hit the limit on the max number of delayed events that can be scheduled at once) + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.TOO_MANY_REQUESTS, channel.code, channel.result) + self.assertEqual( + Codes.LIMIT_EXCEEDED, + channel.json_body["errcode"], + channel.json_body, + ) + # Confirm that the response includes the time remaining until the next of the user's + # delayed events to be sent, at which point another delayed event may be scheduled + # without exceeding the limit + retry_after_headers = channel.headers.getRawHeaders("Retry-After") + assert retry_after_headers + retry_after_sec = int(retry_after_headers[0]) + self.assertGreater(retry_after_sec, 0) + # Confirm that there is only a single value to the Retry-After header, as per RFC9110 + self.assertEqual(1, len(retry_after_headers)) + + # Wait until we're able to retry again (the retry time from the error response) + self.reactor.advance(retry_after_sec) + + # We should be able to send another delayed event again + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.OK, channel.code, channel.result) + + @unittest.override_config( + { + "max_event_delay_duration": "24h", + "experimental_features": { + "msc4140_max_delayed_events_per_user": 1, + }, + } + ) + def test_delayed_event_processed_user_limit_reached(self) -> None: + """ + Test that delayed events in the midst of being sent still count towards the limit of + how many delayed events a user may have scheduled at once. + """ + send_after = Duration(seconds=1) + make_delayed_event_request = lambda: self.make_request( + "POST", + ( + f"rooms/%s/send/m.room.message?org.matrix.msc4140.delay={send_after.as_millis()}" + % self.room_id + ).encode("ascii"), + {"body": "test", "msgtype": "m.text"}, + ) + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.OK, channel.code, channel.result) + + # Simulate the server taking a long time to persist delayed events + simulated_send_lag = Duration(seconds=5) + event_creation_handler = self.hs.get_event_creation_handler() + orig_send_fn = event_creation_handler.create_and_send_nonmember_event + + async def slow_send_fn(*args: Any, **kwargs: Any) -> Any: + await self.clock.sleep(simulated_send_lag) + return await orig_send_fn(*args, **kwargs) + + with patch.object(event_creation_handler, orig_send_fn.__name__, slow_send_fn): + self.reactor.advance(send_after.as_secs()) + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.TOO_MANY_REQUESTS, channel.code, channel.result) + self.assertEqual( + Codes.LIMIT_EXCEEDED, + channel.json_body["errcode"], + channel.json_body, + ) + # Confirm that the response lacks a Retry-After header, because the reason for this limit + # is the server taking an indeterminitely long time to process a delayed event, and the + # server doesn't know how much longer the client should wait before sending more requests + retry_after_headers = channel.headers.getRawHeaders("Retry-After") + assert not retry_after_headers + + # Wait until the delayed event gets persisted + self.reactor.advance(simulated_send_lag.as_secs()) + + # We should be able to send another delayed event again + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.OK, channel.code, channel.result) + + @unittest.override_config( + { + "max_event_delay_duration": "24h", + "experimental_features": { + "msc4140_max_delayed_events_per_user": 5, + }, + } + ) + def test_delayed_event_user_limit_exceeded(self) -> None: + """ + Test that delayed event limits work properly when + the number of already scheduled events exceeds the configured limit. + + This can be invoked by the server admin lowering the configured limit & restarting the server + while a user has fewer scheduled delayed events than the old limit, but more than the new limit. + """ + send_after: Duration + make_delayed_event_request = lambda: self.make_request( + "POST", + ( + f"rooms/%s/send/m.room.message?org.matrix.msc4140.delay={send_after.as_millis()}" + % self.room_id + ).encode("ascii"), + {"body": f"test (send after {send_after.as_secs()}s)", "msgtype": "m.text"}, + ) + + for i in range(4): + send_after = Duration(seconds=i) + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.OK, channel.code, channel.result) + + # Simulate restarting the server after having reconfigured the limit + # to be lower than the number of delayed events we just scheduled. + # + # Set the limit > 1 to test not having to wait for _all_ delayed events + # to be sent before being able to schedule a new one. + self.hs.config.server.max_delayed_events_per_user = 2 + + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.TOO_MANY_REQUESTS, channel.code, channel.result) + self.assertEqual( + Codes.LIMIT_EXCEEDED, + channel.json_body["errcode"], + channel.json_body, + ) + retry_after_header = channel.headers.getRawHeaders("Retry-After") + assert retry_after_header + retry_after_sec = int(retry_after_header[0]) + assert retry_after_sec > 0 + + # Wait until we're able to retry again (the retry time from the error response) + self.reactor.advance(retry_after_sec) + + # We should be able to send another delayed event again + channel = make_delayed_event_request() + self.assertEqual(HTTPStatus.OK, channel.code, channel.result) + @unittest.override_config({"max_event_delay_duration": "24h"}) def test_delayed_event_with_negative_delay(self) -> None: """Test that sending a delayed event fails if its delay is negative.""" @@ -2595,7 +2791,7 @@ class RoomDelayedEventTestCase(RoomBase): """ # Test that new delayed events are correctly ratelimited. - args = ( + make_delayed_event_request = lambda: self.make_request( "POST", ( "rooms/%s/send/m.room.message?org.matrix.msc4140.delay=2000" @@ -2603,9 +2799,9 @@ class RoomDelayedEventTestCase(RoomBase): ).encode("ascii"), {"body": "test", "msgtype": "m.text"}, ) - channel = self.make_request(*args) + channel = make_delayed_event_request() self.assertEqual(HTTPStatus.OK, channel.code, channel.result) - channel = self.make_request(*args) + channel = make_delayed_event_request() self.assertEqual(HTTPStatus.TOO_MANY_REQUESTS, channel.code, channel.result) # Add the current user to the ratelimit overrides, allowing them no ratelimiting. @@ -2614,7 +2810,7 @@ class RoomDelayedEventTestCase(RoomBase): ) # Test that the new delayed events aren't ratelimited anymore. - channel = self.make_request(*args) + channel = make_delayed_event_request() self.assertEqual(HTTPStatus.OK, channel.code, channel.result) diff --git a/tests/rest/client/test_versions.py b/tests/rest/client/test_versions.py index d656098469..bbdbe38e07 100644 --- a/tests/rest/client/test_versions.py +++ b/tests/rest/client/test_versions.py @@ -142,6 +142,17 @@ class VersionsTestCase(unittest.HomeserverTestCase): channel.json_body, ) + def test_msc4446_false_by_default(self) -> None: + channel = self.make_request("GET", "/_matrix/client/versions") + self.assertEqual(channel.code, 200, channel.result) + self.assertFalse(channel.json_body["unstable_features"]["com.beeper.msc4446"]) + + @unittest.override_config({"experimental_features": {"msc4446_enabled": True}}) + def test_msc4446_true_if_enabled(self) -> None: + channel = self.make_request("GET", "/_matrix/client/versions") + self.assertEqual(channel.code, 200, channel.result) + self.assertTrue(channel.json_body["unstable_features"]["com.beeper.msc4446"]) + def _sanity_check_versions_response(self, versions_response: JsonDict) -> None: """ Make sure this looks like a `/_matrix/client/versions` response diff --git a/tests/server.py b/tests/server.py index 15a7661c34..052f76d755 100644 --- a/tests/server.py +++ b/tests/server.py @@ -459,7 +459,6 @@ def make_request( await_result: bool = True, custom_headers: Iterable[CustomHeaderType] | None = None, client_ip: str = "127.0.0.1", - timeout_ms: int = 1000, ) -> FakeChannel: """ Make a web request using the given method, path and content, and render it @@ -488,8 +487,6 @@ def make_request( custom_headers: (name, value) pairs to add as request headers client_ip: The IP to use as the requesting IP. Useful for testing ratelimiting. - timeout_ms: if `await_result` is `True`, the amount of time to wait on - the request before timing out. Ignored otherwise. Returns: channel @@ -574,7 +571,7 @@ def make_request( req.requestReceived(method, path, b"1.1") if await_result: - channel.await_result(timeout_ms=timeout_ms) + channel.await_result() return channel diff --git a/tests/unittest.py b/tests/unittest.py index 202f7120ef..3d130c1d77 100644 --- a/tests/unittest.py +++ b/tests/unittest.py @@ -571,7 +571,6 @@ class HomeserverTestCase(TestCase): await_result: bool = True, custom_headers: Iterable[CustomHeaderType] | None = None, client_ip: str = "127.0.0.1", - timeout_ms: int = 1000, ) -> FakeChannel: """ Create a SynapseRequest at the path using the method and containing the @@ -600,8 +599,6 @@ class HomeserverTestCase(TestCase): client_ip: The IP to use as the requesting IP. Useful for testing ratelimiting. - timeout_ms: if `await_result` is `True`, the amount of time to wait on - the request before timing out. Ignored otherwise. Returns: The FakeChannel object which stores the result of the request. @@ -621,7 +618,6 @@ class HomeserverTestCase(TestCase): await_result, custom_headers, client_ip, - timeout_ms, ) def setup_test_homeserver(