Compare commits

...
Author SHA1 Message Date
Zomatree 4586e4738a feat: publish events over amqp
Signed-off-by: Zomatree <me@zomatree.live>
2026-06-19 21:05:22 +01:00
Tom c70459b10c fix: openapi using old naming (#777)
* fix: openapi using old naming

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: remove january openapi security header

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix(docs): more Revolt usage

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

---------

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-06-07 16:56:00 -07:00
Asraye bebfe34922 fix: point docs favicon to correct location (#789)
chore(docs): update favicon

Signed-off-by: Asraye <asrayeofficial@gmail.com>
2026-06-03 11:06:38 -07:00
Tom 5b769b60de chore(docs): update header logo (#796)
* chore(docs): update header logo

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: don't step on toes

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

---------

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-06-02 17:42:04 -07:00
stoat-tofu[bot] 0896e68882 chore: modify renovate.json 2026-06-01 13:28:29 +00:00
stoat-tofu[bot] 65acc64034 chore: modify .github/workflows/renovate.yml 2026-06-01 13:28:27 +00:00
İspik bd987bf72a chore: update unicode emoji list (#781)
Signed-off-by: ispik <ispik@ispik.dev>
2026-05-28 14:55:49 -07:00
stoat-release[bot]andgithub-actions[bot] 7937179db7 chore(main): release 0.13.7 (#770)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-21 19:33:49 +01:00
İspik 2d308e03d5 fix: sanitize emoji input to handle variation selectors (#774)
Signed-off-by: ispik <ispik@ispik.dev>
2026-05-21 08:21:43 -07:00
Paul Makles b38499a05b ci: hard code packages, gh token limitation [skip ci] (#773) 2026-05-20 20:10:09 +01:00
Paul Makles 4815429952 ci: create Docker images for PR preview (#772) 2026-05-20 19:58:18 +01:00
Zomatree 0d9ae508d9 fix: update mention count badge for channel acks (#769)
Signed-off-by: Zomatree <me@zomatree.live>
2026-05-18 23:38:07 -07:00
stoat-release[bot]andgithub-actions[bot] 03b52655ff chore(main): release 0.13.6 (#762)
* chore(main): release 0.13.6

* chore: update Cargo.lock

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>

---------

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-18 15:53:42 -07:00
Angelo KontaxisandTom 5b1985381a chore: switch to lapin (#767)
* chore: begin switching to lapin fully

Signed-off-by: Zomatree <me@zomatree.live>

* chore: update rest of pushd to lapin

Signed-off-by: Zomatree <me@zomatree.live>

* chore: cleanup code

Signed-off-by: Zomatree <me@zomatree.live>

* chore: cleanup code

Signed-off-by: Zomatree <me@zomatree.live>

* fix: github webui sucks

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

---------

Signed-off-by: Zomatree <me@zomatree.live>
Signed-off-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
Co-authored-by: Tom <iamtomahawkx@gmail.com>
Release-As: 0.13.6
2026-05-18 15:46:17 -07:00
Tom 018afaf38f fix: set env var for publishing crates (#768)
* fix: set env var for publishing crates

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
Release-As: 0.13.6
2026-05-18 15:11:44 -07:00
İspik af0d8aad14 feat: user slowmode events (#760)
* feat: user slowmode events

Signed-off-by: ispik <ispik@ispik.dev>

* fix: remove debug print statement for slowmodes

Signed-off-by: ispik <ispik@ispik.dev>

* refactor: Send user slowmodes as websocket connects instead of trying to send it in ready payload

Signed-off-by: ispik <ispik@ispik.dev>

* refactor: optimize user slowmode handling with bulk operations

Signed-off-by: ispik <ispik@ispik.dev>

* chore: specify release version

Release-As: 0.13.6

---------

Signed-off-by: ispik <ispik@ispik.dev>
2026-05-18 14:51:43 -07:00
Tom acbc087982 feat: Update FCM payload for android notifications (#766)
* feat: modify fcm payload to jens will

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: add message id

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: rename field

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: whitespace

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

---------

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-18 10:57:23 -07:00
Tom 2871632382 fix: voice ingress crashing due to new Result in AMQP::new_auto() (#765)
Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-18 10:56:01 -07:00
Tom 494c8b7cab fix: Use proper headers to determine IP when not behind cloudflare (#764)
Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-18 10:51:19 -07:00
Paul Makles 26a8692677 ci: ignore test errors on main (#763) 2026-05-17 14:54:57 -05:00
Paul Makles 298742dbad fix: include minio region as tests need it (#761) 2026-05-17 14:41:44 -05:00
stoat-release[bot]andgithub-actions[bot] 6c920de03a chore(main): release 0.13.5 (#759)
* chore(main): release 0.13.5

* chore: update Cargo.lock

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>

---------

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-17 11:02:21 -07:00
Angelo Kontaxis c902077cf5 fix: dont panic on hash missing when deleting files (#755)
Signed-off-by: Zomatree <me@zomatree.live>
2026-05-17 10:59:56 -07:00
Tom 19ee535f45 Merge commit from fork
* fix: cache dns & block more ranges

* fix: idle time instead of ttl

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

---------

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-17 10:55:38 -07:00
stoat-release[bot]andgithub-actions[bot] ee4575470b chore(main): release 0.13.4 (#754)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-16 18:27:44 +01:00
Paul Makles 6cfee1f601 fix: add TLS feature to livekit-api crate (#753) 2026-05-16 18:23:45 +01:00
stoat-release[bot]andgithub-actions[bot] ab9b8ccfca chore(main): release 0.13.3 (#750)
* chore(main): release 0.13.3

* chore: update Cargo.lock

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>

---------

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-15 12:22:27 -07:00
Tom 7647cfc8d9 fix: don't automatically set up rabbitmq in delta (#749)
fix: don't declare queues which seem to cause the backend to crash in prod
for now, these exchanges/queues/bindings will need to be declared manually. Hopefully the lapin rewrite will fix this.

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-15 12:17:36 -07:00
stoat-release[bot]andgithub-actions[bot] 8157e1f6e9 chore(main): release 0.13.2 (#747)
* chore(main): release 0.13.2

* chore: update Cargo.lock

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>

---------

Signed-off-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-11 15:26:03 +01:00
Paul Makles fcb8091cd7 fix: update default exchange to revolt.default (#746)
Signed-off-by: Paul Makles <me@insrt.uk>
2026-05-11 15:21:34 +01:00
stoat-release[bot]andgithub-actions[bot] 260036488d chore(main): release 0.13.1 (#745)
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-10 15:34:41 +01:00
Tom 1100eaf46f fix: amqprs startup bug (#744) 2026-05-10 15:22:26 +01:00
stoat-release[bot]andgithub-actions[bot] d52e84c5d3 chore(main): release 0.13.0 (#722)
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-05-09 17:26:38 +01:00
İspik d76a71141f fix: don't strip ICC from exif (#735)
* fix: don't strip ICC from exif

Signed-off-by: ispik <ispik@ispik.dev>

* fix: refactor ICC profile handling in image processing

Signed-off-by: ispik <ispik@ispik.dev>

---------

Signed-off-by: ispik <ispik@ispik.dev>
2026-05-08 15:49:31 -07:00
TaureonandTaureon 6f3441cf4a feat: Add webhook endpoints for editing and deleting messages (#682)
* feat: ErrorType.CannotDeleteMessage, needed later

Signed-off-by: Taureon <taureon@noreply.codeberg.org>

* feat: webhook edit/delete message endpoints

Signed-off-by: Taureon <taureon@noreply.codeberg.org>

* lol, lmao even

Signed-off-by: Taureon <taureon@noreply.codeberg.org>

* fix contradictory comment

Signed-off-by: Taureon <taureon@noreply.codeberg.org>

---------

Signed-off-by: Taureon <taureon@noreply.codeberg.org>
Co-authored-by: Taureon <taureon@noreply.codeberg.org>
2026-05-08 15:37:00 -07:00
Gabriel 23ad135983 feat: add emoji rename endpoint (#714)
* feat: add rename endpoint with rename-only update

Signed-off-by: Gabriel <60961939+gabrielfordevelopment@users.noreply.github.com>

* fix: enforce detached emoji edit restrictions

Signed-off-by: Gabriel <60961939+gabrielfordevelopment@users.noreply.github.com>

* fix: always enforce emoji edit permissions

Signed-off-by: Gabriel <60961939+gabrielfordevelopment@users.noreply.github.com>

---------

Signed-off-by: Gabriel <60961939+gabrielfordevelopment@users.noreply.github.com>
2026-05-08 15:23:57 -07:00
İspikandIAmTomahawkx d46c7f7f3c feat: add embed support for YouTube Shorts (#734)
* feat: add embed support for YouTube Shorts

Signed-off-by: ispik <ispik@ispik.dev>

* feat: blacklist private ip ranges and add january domain blocklist

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: remove duplicates

Signed-off-by: ispik <ispik@ispik.dev>

---------

Signed-off-by: ispik <ispik@ispik.dev>
Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
Co-authored-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-08 15:23:30 -07:00
Tom 0719985ac5 fix: docker compose file had personal url in it (#742) 2026-05-08 15:22:10 -07:00
Tom ab5bd47a39 feat: Rewrite acks (#741)
* feat: rewrite ack system

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* feat: rewrite acks to crond + rabbit task

* fix: review changes

---------

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-08 14:39:16 -07:00
İspik 9fd7128f80 fix: encode filenames in redirects (#737)
* fix: encode filenames in redirects

Signed-off-by: ispik <ispik@ispik.dev>

* refactor: don't add another dependency when the one needed exists

Signed-off-by: ispik <ispik@ispik.dev>

---------

Signed-off-by: ispik <ispik@ispik.dev>
2026-05-08 14:38:57 -07:00
İspik df276ac40b chore: update emoji list (#740)
Signed-off-by: ispik <ispik@ispik.dev>
2026-05-08 14:37:35 -07:00
Tom 356491e934 fix: january ip redirects & domain resolver (#738)
* fix: properly block private ip ranges
I'm a dunce and forgot that domains do in fact resolve to ips.

* fix: reimplement max redirects

* fix: remove my debug error

* fix: actually check redirect urls
kind of the whole point of this thing.
2026-05-06 21:33:52 -07:00
İspik 21d82018cf feat: add legal links to root payload (#733)
Signed-off-by: ispik <ispik@ispik.dev>
2026-05-06 16:35:20 -07:00
Tom 6b41db984b feat: blacklist private ip ranges and add january domain blocklist (#731)
Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-05-06 16:25:31 -07:00
İspik 5378cd22b4 fix: use correct response for NoEffect errors (#732)
Signed-off-by: ispik <ispik@ispik.dev>
2026-05-05 11:57:42 -07:00
İspik 841985d3b9 feat: add role icon support (#724)
Signed-off-by: ispik <ispik@ispik.dev>
2026-05-02 17:04:28 -07:00
İspik 279f5d5fd7 fix: add new_user_hours to configuration limits (#729)
Signed-off-by: ispik <ispik@ispik.dev>
2026-04-26 22:47:52 -07:00
e93769786c feat: automatically sanitise usernames on create/update (#689)
* feat: automatically sanitise usernames on create/update

Signed-off-by: higgs01 <6546697+higgs01@users.noreply.github.com>

* test: add tests for validation and sanitasion

Signed-off-by: higgs01 <6546697+higgs01@users.noreply.github.com>

* fix: Return only the sanitised string in sanitise_username (#424)

* fix: Use as_str for role.id in insert_role
* Run rustfmt for code changes

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* feat: Make minimum username length configurable in Revolt.toml (#424)

* fix: Remove redundant call to is_match with blocked username regex

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* Update crates/core/database/src/models/users/model.rs

Co-authored-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: jarvarvarvis <53998846+jarvarvarvis@users.noreply.github.com>

* Update crates/core/database/src/models/users/model.rs

Co-authored-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: jarvarvarvis <53998846+jarvarvarvis@users.noreply.github.com>

* Update crates/core/database/src/models/users/model.rs

Co-authored-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: jarvarvarvis <53998846+jarvarvarvis@users.noreply.github.com>

* Update crates/core/database/src/models/users/model.rs

Co-authored-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: jarvarvarvis <53998846+jarvarvarvis@users.noreply.github.com>

* fix: Implement suggested changes and clean up last 4 commits

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* fix: Disallow stoat as username, update create_user test

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* fix: Use sanitised username to find updated discriminator

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* feat: Sanitise revolt.chat in username

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* fix: Implement discussed changes

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* test: Fix create_user test

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

* fix: don't overflow the stack
not entirely sure why this fixes it and I don't like it. But work it does.

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: revert odd file mode change

Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>

---------

Signed-off-by: higgs01 <6546697+higgs01@users.noreply.github.com>
Signed-off-by: jarvarvarvis <jarvistrigo@gmail.com>
Signed-off-by: jarvarvarvis <53998846+jarvarvarvis@users.noreply.github.com>
Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
Co-authored-by: higgs01 <6546697+higgs01@users.noreply.github.com>
Co-authored-by: Tom <iamtomahawkx@gmail.com>
2026-04-23 21:02:28 -07:00
İspik ed4fd5ebfe fix: update message length validation to remove upper limit (#723)
Signed-off-by: ispik <ispik@ispik.dev>
2026-04-23 20:39:13 -07:00
Angelo Kontaxis 89171e9bd0 fix: dont send notification in fcm (#721)
* fix: dont send notification in fcm

Signed-off-by: Zomatree <me@zomatree.live>

* fix: add notification type to fcm

Signed-off-by: Zomatree <me@zomatree.live>

* fix: switch to structured notification data

Signed-off-by: Zomatree <me@zomatree.live>

---------

Signed-off-by: Zomatree <me@zomatree.live>
2026-04-22 23:48:26 -07:00
Infiland 057f2bb8b3 fix: add reconnection policy to Redis subscriber to prevent ghost state (#708)
* fix: add reconnection policy to Redis subscriber to prevent ghost state

- Add ReconnectPolicy::new_exponential(0, 100, 30_000, 2) to the subscriber builder, unlimited retries with exponential backoff (100ms min, 30s max)
- Add on_reconnect handler that signals the listener loop to force a subscription reset, re-subscribing to all topics on the new connection
- Add warn-level logging to on_error for all Redis subscriber errors (previously only Canceled was handled, others were silently ignored)

Signed-off-by: Infiland <ljubica.citydesign@gmail.com>

* Update websocket.rs

Signed-off-by: Infiland <88491175+Infiland@users.noreply.github.com>

* Auto-manage subscriptions on reconnect

Call subscriber.manage_subscriptions() so the subscriber will automatically re-subscribe tracked channels after a Redis reconnect. Remove the manual reconnect channel and on_reconnect handler along with the select branch that forced SubscriptionStateChange::Reset.

Signed-off-by: Infiland <88491175+Infiland@users.noreply.github.com>

---------

Signed-off-by: Infiland <ljubica.citydesign@gmail.com>
Signed-off-by: Infiland <88491175+Infiland@users.noreply.github.com>
2026-04-17 20:15:54 -07:00
3675ff1a1f chore: migrate all local dependancies to workspace dependancies (#710)
* chore: start moving all deps to workspace deps

Signed-off-by: Zomatree <me@zomatree.live>

* chore: migrate all deps to workspace deps

Signed-off-by: Zomatree <me@zomatree.live>

* chore: add more dep groups

Signed-off-by: Zomatree <me@zomatree.live>

* fix: add migration to update existing files to be animated (#705)

* fix: add migration to update existing files to be animated

Signed-off-by: Zomatree <me@zomatree.live>

* Revert "fix: add migration to update existing files to be animated"

This reverts commit 4e1f1c116c.

Signed-off-by: Zomatree <me@zomatree.live>

* fix: calculate animated for existing files when fetched

Signed-off-by: Zomatree <me@zomatree.live>

---------

Signed-off-by: Zomatree <me@zomatree.live>

* fix: mise start + missing docker image (#564)

* fix: mise start + missing docker image

Signed-off-by: Damocles078 <hellodamocles078@gmail.com>

* fix: bump livekit version

Signed-off-by: Damocles <106018783+Damocles078@users.noreply.github.com>

---------

Signed-off-by: Damocles078 <hellodamocles078@gmail.com>
Signed-off-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: Damocles <106018783+Damocles078@users.noreply.github.com>
Co-authored-by: Tom <iamtomahawkx@gmail.com>

* docs: update donation link (#709)

Signed-off-by: Zomatree <me@zomatree.live>

* fix: remove usage of deprecated functions

Signed-off-by: Zomatree <me@zomatree.live>

---------

Signed-off-by: Zomatree <me@zomatree.live>
Signed-off-by: Damocles078 <hellodamocles078@gmail.com>
Signed-off-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: Damocles <106018783+Damocles078@users.noreply.github.com>
Co-authored-by: Damocles <106018783+Damocles078@users.noreply.github.com>
Co-authored-by: Tom <iamtomahawkx@gmail.com>
Co-authored-by: Paul Makles <me@insrt.uk>
2026-04-17 19:02:18 -07:00
infi 144e939c6b docs: add run in yaak button to endpoints page (#711) 2026-04-12 00:46:05 +01:00
stoat-release[bot]andgithub-actions[bot] 61fd13629f chore(main): release 0.12.1 (#700)
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2026-04-11 12:31:11 +01:00
Paul Makles 47336a7940 docs: pr template (#713) 2026-04-10 20:44:29 +01:00
Paul Makles 4ea8f90c0d docs: update donation link (#709) 2026-04-03 00:28:49 +01:00
DamoclesandTom fb8fe16557 fix: mise start + missing docker image (#564)
* fix: mise start + missing docker image

Signed-off-by: Damocles078 <hellodamocles078@gmail.com>

* fix: bump livekit version

Signed-off-by: Damocles <106018783+Damocles078@users.noreply.github.com>

---------

Signed-off-by: Damocles078 <hellodamocles078@gmail.com>
Signed-off-by: Tom <iamtomahawkx@gmail.com>
Signed-off-by: Damocles <106018783+Damocles078@users.noreply.github.com>
Co-authored-by: Tom <iamtomahawkx@gmail.com>
2026-04-01 19:51:14 -07:00
Angelo Kontaxis f2c056a151 fix: add migration to update existing files to be animated (#705)
* fix: add migration to update existing files to be animated

Signed-off-by: Zomatree <me@zomatree.live>

* Revert "fix: add migration to update existing files to be animated"

This reverts commit 4e1f1c116c.

Signed-off-by: Zomatree <me@zomatree.live>

* fix: calculate animated for existing files when fetched

Signed-off-by: Zomatree <me@zomatree.live>

---------

Signed-off-by: Zomatree <me@zomatree.live>
2026-04-01 19:48:22 -07:00
Tom f30b729ca9 fix: don't send self dm notifications (#706)
* fix: don't send self dm notifications

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: test failure due to wrong assertion (#707)

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

* fix: RA auto import moment

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>

---------

Signed-off-by: IAmTomahawkx <iamtomahawkx@gmail.com>
2026-03-31 15:25:36 -07:00
Tom f81e3291bd fix: test failure due to wrong assertion (#707) 2026-03-31 11:06:39 +01:00
Tom 1e80916b65 Fix/release please fix 4 (#703)
* fix: release-please-3

* fix: release please bs v4
2026-03-28 23:05:01 -07:00
Tom 8814bf4c23 fix: release-please-3 (#702) 2026-03-28 22:53:12 -07:00
Tom 5466ada9ae Fix: release please bs 2 electric boogaloo (#701)
* fix: update release (me from this) please

* fix: add dep to release please

* Revert "fix: update release (me from this) please"

This reverts commit 0618c492c7.
2026-03-28 22:42:13 -07:00
Tom f0e513ccae fix: update release (me from this) please (#699) 2026-03-28 22:35:51 -07:00
stoat-release[bot] 4d4b0dd864 chore(main): release 0.12.0 (#602)
Co-authored-by: stoat-release[bot] <245062572+stoat-release[bot]@users.noreply.github.com>
2026-03-28 22:22:37 -07:00
122 changed files with 10260 additions and 5702 deletions
+24
View File
@@ -0,0 +1,24 @@
<!-- Describe your changes -->
Fixes # (issue)
## How was this PR tested?
<!-- What did you do to test your changes? -->
- [ ] Test A
- [ ] Test B
## Checklist:
- [ ] I have carefully read [the contributing guidelines](https://developers.stoat.chat/developing/contrib/)
- [ ] I have performed a self-review of my own code
- [ ] I have made corresponding changes to the documentation if applicable
- [ ] I have no unrelated changes in the PR
- [ ] I have confirmed that any new dependencies are strictly necessary
- [ ] I have written tests for new code (if applicable)
- [ ] I have followed naming conventions/patterns in the surrounding code
## Please declare, if any, LLM usage involved in creating this PR
...
+46
View File
@@ -0,0 +1,46 @@
name: Docker PR Image Cleanup
on:
pull_request:
types:
- closed
permissions:
contents: read
packages: write
concurrency:
group: docker-cleanup-${{ github.event.pull_request.number }}
cancel-in-progress: false
jobs:
cleanup:
runs-on: ubuntu-latest
if: ${{ !github.event.pull_request.head.repo.fork }}
strategy:
fail-fast: false
matrix:
package:
- base
- api
- events
- file-server
- proxy
- gifbox
- crond
- pushd
- voice-ingress
steps:
- env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
ORG: stoatchat
PACKAGE: ${{ matrix.package }}
TAG: pr-${{ github.event.pull_request.number }}
run: |
set -euo pipefail
gh api --paginate \
"/orgs/${ORG}/packages/container/${PACKAGE}/versions" \
--jq ".[] | select(.metadata.container.tags | index(\"${TAG}\")) | .id" \
| while read -r id; do
gh api -X DELETE "/orgs/${ORG}/packages/container/${PACKAGE}/versions/${id}"
done
+23 -14
View File
@@ -5,8 +5,6 @@ on:
tags: tags:
- "*" - "*"
pull_request: pull_request:
paths:
- "Dockerfile"
workflow_dispatch: workflow_dispatch:
permissions: permissions:
@@ -19,9 +17,9 @@ concurrency:
jobs: jobs:
base: base:
name: Test base image build name: Test base image build (fork)
runs-on: arc-runner-set runs-on: arc-runner-set
if: github.event_name == 'pull_request' if: github.event_name == 'pull_request' && github.event.pull_request.head.repo.fork
steps: steps:
# Configure build environment # Configure build environment
- name: Checkout - name: Checkout
@@ -42,7 +40,7 @@ jobs:
publish: publish:
runs-on: arc-runner-set runs-on: arc-runner-set
if: github.event_name != 'pull_request' if: ${{ github.event_name != 'pull_request' || !github.event.pull_request.head.repo.fork }}
name: Publish Docker images name: Publish Docker images
steps: steps:
# Configure build environment # Configure build environment
@@ -59,6 +57,15 @@ jobs:
username: ${{ github.actor }} username: ${{ github.actor }}
password: ${{ secrets.GITHUB_TOKEN }} password: ${{ secrets.GITHUB_TOKEN }}
- name: Determine base image tag
id: base
run: |
if [ "${{ github.event_name }}" = "pull_request" ]; then
echo "tag=pr-${{ github.event.number }}" >> "$GITHUB_OUTPUT"
else
echo "tag=latest" >> "$GITHUB_OUTPUT"
fi
# Build the image # Build the image
- name: Build base image - name: Build base image
uses: docker/build-push-action@v4 uses: docker/build-push-action@v4
@@ -66,7 +73,9 @@ jobs:
context: . context: .
push: true push: true
platforms: linux/amd64,linux/arm64 platforms: linux/amd64,linux/arm64
tags: ghcr.io/${{ github.repository_owner }}/base:latest tags: ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
cache-from: type=gha,scope=buildx-base-multi-arch
cache-to: type=gha,scope=buildx-base-multi-arch,mode=max
# stoatchat/api # stoatchat/api
- name: Docker meta - name: Docker meta
@@ -84,7 +93,7 @@ jobs:
file: crates/delta/Dockerfile file: crates/delta/Dockerfile
tags: ${{ steps.meta-delta.outputs.tags }} tags: ${{ steps.meta-delta.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-delta.outputs.labels }} labels: ${{ steps.meta-delta.outputs.labels }}
# stoatchat/events # stoatchat/events
@@ -103,7 +112,7 @@ jobs:
file: crates/bonfire/Dockerfile file: crates/bonfire/Dockerfile
tags: ${{ steps.meta-bonfire.outputs.tags }} tags: ${{ steps.meta-bonfire.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-bonfire.outputs.labels }} labels: ${{ steps.meta-bonfire.outputs.labels }}
# stoatchat/file-server # stoatchat/file-server
@@ -122,7 +131,7 @@ jobs:
file: crates/services/autumn/Dockerfile file: crates/services/autumn/Dockerfile
tags: ${{ steps.meta-autumn.outputs.tags }} tags: ${{ steps.meta-autumn.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-autumn.outputs.labels }} labels: ${{ steps.meta-autumn.outputs.labels }}
# stoatchat/proxy # stoatchat/proxy
@@ -141,7 +150,7 @@ jobs:
file: crates/services/january/Dockerfile file: crates/services/january/Dockerfile
tags: ${{ steps.meta-january.outputs.tags }} tags: ${{ steps.meta-january.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-january.outputs.labels }} labels: ${{ steps.meta-january.outputs.labels }}
# stoatchat/gifbox # stoatchat/gifbox
@@ -160,7 +169,7 @@ jobs:
file: crates/services/gifbox/Dockerfile file: crates/services/gifbox/Dockerfile
tags: ${{ steps.meta-gifbox.outputs.tags }} tags: ${{ steps.meta-gifbox.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-gifbox.outputs.labels }} labels: ${{ steps.meta-gifbox.outputs.labels }}
# stoatchat/crond # stoatchat/crond
@@ -179,7 +188,7 @@ jobs:
file: crates/daemons/crond/Dockerfile file: crates/daemons/crond/Dockerfile
tags: ${{ steps.meta-crond.outputs.tags }} tags: ${{ steps.meta-crond.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-crond.outputs.labels }} labels: ${{ steps.meta-crond.outputs.labels }}
# stoatchat/pushd # stoatchat/pushd
@@ -198,7 +207,7 @@ jobs:
file: crates/daemons/pushd/Dockerfile file: crates/daemons/pushd/Dockerfile
tags: ${{ steps.meta-pushd.outputs.tags }} tags: ${{ steps.meta-pushd.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-pushd.outputs.labels }} labels: ${{ steps.meta-pushd.outputs.labels }}
# stoatchat/voice-ingress # stoatchat/voice-ingress
@@ -217,5 +226,5 @@ jobs:
file: crates/daemons/voice-ingress/Dockerfile file: crates/daemons/voice-ingress/Dockerfile
tags: ${{ steps.meta-voice-ingress.outputs.tags }} tags: ${{ steps.meta-voice-ingress.outputs.tags }}
build-args: | build-args: |
BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:latest BASE_IMAGE=ghcr.io/${{ github.repository_owner }}/base:${{ steps.base.outputs.tag }}
labels: ${{ steps.meta-voice-ingress.outputs.labels }} labels: ${{ steps.meta-voice-ingress.outputs.labels }}
+2
View File
@@ -21,4 +21,6 @@ jobs:
github-token: ${{ secrets.GITHUB_TOKEN }} github-token: ${{ secrets.GITHUB_TOKEN }}
- name: Publish - name: Publish
env:
CARGO_REGISTRY_TOKEN: ${{ secrets.CARGO_REGISTRY_TOKEN }}
run: mise publish --workspace run: mise publish --workspace
+30
View File
@@ -0,0 +1,30 @@
# DO NOT EDIT DIRECTLY IN REPOSITORY
# Managed in Terraform templates
name: Renovate
on:
workflow_dispatch:
schedule:
- cron: '0/15 * * * *'
jobs:
renovate:
runs-on: ubuntu-latest
steps:
- id: app-token
uses: actions/create-github-app-token@v2
with:
app-id: ${{ secrets.GH_STOAT_RELEASE_APP_ID }}
private-key: ${{ secrets.GH_STOAT_RELEASE_APP_PRIVATE_KEY }}
- name: Setup Mise
uses: immich-app/devtools/actions/use-mise@7b8610a904d57da241e4ddba17fa62b62b15aed4 # use-mise-action-v2.0.2
with:
github_token: ${{ steps.app-token.outputs.token }}
- name: Self-hosted Renovate
uses: renovatebot/github-action@v46.1.14
with:
token: '${{ steps.app-token.outputs.token }}'
env:
RENOVATE_PLATFORM_COMMIT: 'enabled'
RENOVATE_REPOSITORIES: '${{ github.repository }}'
+2
View File
@@ -37,6 +37,7 @@ jobs:
- name: Reference Test - name: Reference Test
env: env:
TEST_DB: REFERENCE TEST_DB: REFERENCE
continue-on-error: ${{ github.ref_name == 'main' }}
run: | run: |
mise test mise test
@@ -44,6 +45,7 @@ jobs:
env: env:
TEST_DB: MONGODB TEST_DB: MONGODB
MONGODB: mongodb://localhost MONGODB: mongodb://localhost
continue-on-error: ${{ github.ref_name == 'main' }}
run: | run: |
mise test mise test
+9
View File
@@ -16,4 +16,13 @@ idiomatic_version_file_enable_tools = ["rust"]
[tasks.start] [tasks.start]
description = "Run all services" description = "Run all services"
depends = ["docker:start", "build"] depends = ["docker:start", "build"]
wait_for = ["docker:start", "build"]
run = [{ task = "service:*" }] run = [{ task = "service:*" }]
[env]
BUILDER = "cargo"
DOCKER_NETWORK_NAME = "stoatchat_default"
DATABASE_PORT = "27017"
RABBIT_PORT = "5672"
REDIS_PORT = "6379"
_.file = { path = ".env", tools = true }
+1 -1
View File
@@ -2,4 +2,4 @@
#MISE description="Build project" #MISE description="Build project"
set -e set -e
cargo build "$@" ${BUILDER} build "$@"
+8
View File
@@ -3,3 +3,11 @@
set -e set -e
docker compose up -d docker compose up -d
docker run \
--network=${DOCKER_NETWORK_NAME} \
--name wait \
--rm dokku/wait -c \
rabbit:${RABBIT_PORT},\
database:${DATABASE_PORT},\
redis:${REDIS_PORT}
+1 -1
View File
@@ -1,3 +1,3 @@
{ {
".": "0.11.5" ".": "0.13.7"
} }
+139
View File
@@ -1,5 +1,144 @@
# Changelog # Changelog
## [0.13.7](https://github.com/stoatchat/stoatchat/compare/v0.13.6...v0.13.7) (2026-05-21)
### Bug Fixes
* sanitize emoji input to handle variation selectors ([#774](https://github.com/stoatchat/stoatchat/issues/774)) ([2d308e0](https://github.com/stoatchat/stoatchat/commit/2d308e03d58c19f27b5b4d65dc2a15ef20b56190))
* update mention count badge for channel acks ([#769](https://github.com/stoatchat/stoatchat/issues/769)) ([0d9ae50](https://github.com/stoatchat/stoatchat/commit/0d9ae508d9d2199f0e408b8ca634d20489be6f61))
## [0.13.6](https://github.com/stoatchat/stoatchat/compare/v0.13.5...v0.13.6) (2026-05-18)
### Features
* Update FCM payload for android notifications ([#766](https://github.com/stoatchat/stoatchat/issues/766)) ([acbc087](https://github.com/stoatchat/stoatchat/commit/acbc087982e9aeb05cabc5ab4c9b1291f67490ad))
* user slowmode events ([#760](https://github.com/stoatchat/stoatchat/issues/760)) ([af0d8aa](https://github.com/stoatchat/stoatchat/commit/af0d8aad14dc68d88159d0e1c714077d362e21e4))
### Bug Fixes
* include `minio` region as tests need it ([#761](https://github.com/stoatchat/stoatchat/issues/761)) ([298742d](https://github.com/stoatchat/stoatchat/commit/298742dbad4eafae356f976c56b9db23904b0c3a))
* set env var for publishing crates ([#768](https://github.com/stoatchat/stoatchat/issues/768)) ([018afaf](https://github.com/stoatchat/stoatchat/commit/018afaf38f6330d92dad2a68b640c0cb3f6b639a))
* Use proper headers to determine IP when not behind cloudflare ([#764](https://github.com/stoatchat/stoatchat/issues/764)) ([494c8b7](https://github.com/stoatchat/stoatchat/commit/494c8b7cabaae2a51039a7a5b559d5e2e5279554))
* voice ingress crashing due to new Result in AMQP::new_auto() ([#765](https://github.com/stoatchat/stoatchat/issues/765)) ([2871632](https://github.com/stoatchat/stoatchat/commit/2871632382395cb20cbe0047c542d3ac31ff3f03))
### Miscellaneous Chores
* switch to lapin ([#767](https://github.com/stoatchat/stoatchat/issues/767)) ([5b19853](https://github.com/stoatchat/stoatchat/commit/5b1985381ae829a92c80a19e91a414cd9dc4de93))
## [0.13.5](https://github.com/stoatchat/stoatchat/compare/v0.13.4...v0.13.5) (2026-05-17)
### Bug Fixes
* dont panic on hash missing when deleting files ([#755](https://github.com/stoatchat/stoatchat/issues/755)) ([c902077](https://github.com/stoatchat/stoatchat/commit/c902077cf51076fee11712eb732dc8a8f786fc4b))
## [0.13.4](https://github.com/stoatchat/stoatchat/compare/v0.13.3...v0.13.4) (2026-05-16)
### Bug Fixes
* add TLS feature to livekit-api crate ([#753](https://github.com/stoatchat/stoatchat/issues/753)) ([6cfee1f](https://github.com/stoatchat/stoatchat/commit/6cfee1f601c1e084df7c8f1e7a5e8a560d1dd514))
## [0.13.3](https://github.com/stoatchat/stoatchat/compare/v0.13.2...v0.13.3) (2026-05-15)
### Bug Fixes
* don't automatically set up rabbitmq in delta ([#749](https://github.com/stoatchat/stoatchat/issues/749)) ([7647cfc](https://github.com/stoatchat/stoatchat/commit/7647cfc8d93aba99f5faef13eb3d970097540d76))
* don't declare queues which seem to cause the backend to crash in prod ([7647cfc](https://github.com/stoatchat/stoatchat/commit/7647cfc8d93aba99f5faef13eb3d970097540d76))
## [0.13.2](https://github.com/stoatchat/stoatchat/compare/v0.13.1...v0.13.2) (2026-05-11)
### Bug Fixes
* update default exchange to `revolt.default` ([#746](https://github.com/stoatchat/stoatchat/issues/746)) ([fcb8091](https://github.com/stoatchat/stoatchat/commit/fcb8091cd7a00d7f26c798daa33aae4b923b2a8b))
## [0.13.1](https://github.com/stoatchat/stoatchat/compare/v0.13.0...v0.13.1) (2026-05-10)
### Bug Fixes
* amqprs startup bug ([#744](https://github.com/stoatchat/stoatchat/issues/744)) ([1100eaf](https://github.com/stoatchat/stoatchat/commit/1100eaf46f849f2509ae01ac497556ca33bde778))
## [0.13.0](https://github.com/stoatchat/stoatchat/compare/v0.12.1...v0.13.0) (2026-05-08)
### Features
* add embed support for YouTube Shorts ([#734](https://github.com/stoatchat/stoatchat/issues/734)) ([d46c7f7](https://github.com/stoatchat/stoatchat/commit/d46c7f7f3c04524c0639c3e0a122626f8e0b3bf7))
* add emoji rename endpoint ([#714](https://github.com/stoatchat/stoatchat/issues/714)) ([23ad135](https://github.com/stoatchat/stoatchat/commit/23ad1359834bb7d07a460b8678d6a6ebffc73eb0))
* add legal links to root payload ([#733](https://github.com/stoatchat/stoatchat/issues/733)) ([21d8201](https://github.com/stoatchat/stoatchat/commit/21d82018cf84ab0fdd10613d254b9562aea8eea3))
* add role icon support ([#724](https://github.com/stoatchat/stoatchat/issues/724)) ([841985d](https://github.com/stoatchat/stoatchat/commit/841985d3b994df1c6eefab2fc7ecbd77ab22c493))
* Add webhook endpoints for editing and deleting messages ([#682](https://github.com/stoatchat/stoatchat/issues/682)) ([6f3441c](https://github.com/stoatchat/stoatchat/commit/6f3441cf4acac2a8e6e1bf07a279a153b80f7956))
* automatically sanitise usernames on create/update ([#689](https://github.com/stoatchat/stoatchat/issues/689)) ([e937697](https://github.com/stoatchat/stoatchat/commit/e93769786c7669485a659ee471630740d3cea702))
* blacklist private ip ranges and add january domain blocklist ([#731](https://github.com/stoatchat/stoatchat/issues/731)) ([6b41db9](https://github.com/stoatchat/stoatchat/commit/6b41db984bb491b2e58324309cc70d8c14e0b814))
* Rewrite acks ([#741](https://github.com/stoatchat/stoatchat/issues/741)) ([ab5bd47](https://github.com/stoatchat/stoatchat/commit/ab5bd47a39ee889de0b5ae6e7b560620853daead))
### Bug Fixes
* add new_user_hours to configuration limits ([#729](https://github.com/stoatchat/stoatchat/issues/729)) ([279f5d5](https://github.com/stoatchat/stoatchat/commit/279f5d5fd7af2df55902c706859ec07f569cdb1e))
* add reconnection policy to Redis subscriber to prevent ghost state ([#708](https://github.com/stoatchat/stoatchat/issues/708)) ([057f2bb](https://github.com/stoatchat/stoatchat/commit/057f2bb8b359f8b942741a30ff54eeb8fbe3e0b1))
* docker compose file had personal url in it ([#742](https://github.com/stoatchat/stoatchat/issues/742)) ([0719985](https://github.com/stoatchat/stoatchat/commit/0719985ac5636590f91e6f9ec4b68f3eded70c13))
* don't strip ICC from exif ([#735](https://github.com/stoatchat/stoatchat/issues/735)) ([d76a711](https://github.com/stoatchat/stoatchat/commit/d76a71141f3e508f6308ba52fa28eaeb56fb3438))
* dont send notification in fcm ([#721](https://github.com/stoatchat/stoatchat/issues/721)) ([89171e9](https://github.com/stoatchat/stoatchat/commit/89171e9bd0f15711157e78c6eec0fe7b480de93a))
* encode filenames in redirects ([#737](https://github.com/stoatchat/stoatchat/issues/737)) ([9fd7128](https://github.com/stoatchat/stoatchat/commit/9fd7128f800badbd184baf943d4f799e601201e4))
* january ip redirects & domain resolver ([#738](https://github.com/stoatchat/stoatchat/issues/738)) ([356491e](https://github.com/stoatchat/stoatchat/commit/356491e934b274f9e895df883dd63ef0b3123510))
* update message length validation to remove upper limit ([#723](https://github.com/stoatchat/stoatchat/issues/723)) ([ed4fd5e](https://github.com/stoatchat/stoatchat/commit/ed4fd5ebfe6d0ea534a0898da4afdc1f4e2cd6c5))
* use correct response for NoEffect errors ([#732](https://github.com/stoatchat/stoatchat/issues/732)) ([5378cd2](https://github.com/stoatchat/stoatchat/commit/5378cd22b4c7d85f44c31a6af0dda00941b80d5c))
## [0.12.1](https://github.com/stoatchat/stoatchat/compare/v0.12.0...v0.12.1) (2026-04-10)
### Bug Fixes
* add migration to update existing files to be animated ([#705](https://github.com/stoatchat/stoatchat/issues/705)) ([f2c056a](https://github.com/stoatchat/stoatchat/commit/f2c056a1515be493b195f3f5db5886c2ddf36700))
* don't send self dm notifications ([#706](https://github.com/stoatchat/stoatchat/issues/706)) ([f30b729](https://github.com/stoatchat/stoatchat/commit/f30b729ca90d0be6853c57ab4935694e5e59ae56))
* mise start + missing docker image ([#564](https://github.com/stoatchat/stoatchat/issues/564)) ([fb8fe16](https://github.com/stoatchat/stoatchat/commit/fb8fe1655776791421284a6a093e86f0320c258a))
* test failure due to wrong assertion ([#707](https://github.com/stoatchat/stoatchat/issues/707)) ([f81e329](https://github.com/stoatchat/stoatchat/commit/f81e3291bdd57af9ceedb2987b111acc7051d69c))
## [0.12.0](https://github.com/stoatchat/stoatchat/compare/v0.11.5...v0.12.0) (2026-03-28)
### Features
* add bug report template for issue tracking ([#627](https://github.com/stoatchat/stoatchat/issues/627)) ([f777e28](https://github.com/stoatchat/stoatchat/commit/f777e2863c6ca50057c8b5d0a5be14915d287724))
* Add slowmode functionality to text channels ([#680](https://github.com/stoatchat/stoatchat/issues/680)) ([6107f24](https://github.com/stoatchat/stoatchat/commit/6107f242fd3ebaff71a15f9a16330ffbcb4f2d7b))
* Allow restricting server creation to specific users ([#685](https://github.com/stoatchat/stoatchat/issues/685)) ([edfa97d](https://github.com/stoatchat/stoatchat/commit/edfa97db108c9c81828547f98a1db5315cb5ba4a))
* compute thumbhash for images ([#596](https://github.com/stoatchat/stoatchat/issues/596)) ([c2d4369](https://github.com/stoatchat/stoatchat/commit/c2d4369e160f32d79bce0a0b0f14677f89de3669))
* Detect animation in image files for fetch_preview ([#574](https://github.com/stoatchat/stoatchat/issues/574)) ([3fa0abf](https://github.com/stoatchat/stoatchat/commit/3fa0abf47f5f42ddd8ee041fe4c44fbc5ba800c1))
* expose global and user limits in root API response ([#644](https://github.com/stoatchat/stoatchat/issues/644)) ([0b522eb](https://github.com/stoatchat/stoatchat/commit/0b522ebddc17f2e3f792ff5e2347793e9849fa23))
* implement time based message sweep on user ban ([#670](https://github.com/stoatchat/stoatchat/issues/670)) ([98c7b1b](https://github.com/stoatchat/stoatchat/commit/98c7b1b5a5b9fdac5c0ab83be10f0e23114dbfc9))
* load config from env vars ([#576](https://github.com/stoatchat/stoatchat/issues/576)) ([5191bd1](https://github.com/stoatchat/stoatchat/commit/5191bd16b2a905b8409838e34eb0baca96f08580))
* parse message push notification content and replace internal formatting ([#693](https://github.com/stoatchat/stoatchat/issues/693)) ([d1e72ce](https://github.com/stoatchat/stoatchat/commit/d1e72cee42c54e16f4e49af569897528b10a28ca))
* Transfer ownership ([#396](https://github.com/stoatchat/stoatchat/issues/396)) ([735d644](https://github.com/stoatchat/stoatchat/commit/735d644e043793cb86e74aab5b88bb4b8bc17ba2))
* update livekit ([#698](https://github.com/stoatchat/stoatchat/issues/698)) ([f181edc](https://github.com/stoatchat/stoatchat/commit/f181edc8f2ff3ce4b6d48938dfc73931ecfa2279))
### Bug Fixes
* add flag for disabling events instead of commenting them out ([#695](https://github.com/stoatchat/stoatchat/issues/695)) ([a5cd08a](https://github.com/stoatchat/stoatchat/commit/a5cd08a655dece4269f3ac84fa2387ae356709a5))
* add masquerade permission to default direct message settings ([#665](https://github.com/stoatchat/stoatchat/issues/665)) ([ab52569](https://github.com/stoatchat/stoatchat/commit/ab525699bd6663333f0e9fed6d2455e482e6a09f))
* Check for appropriate permission for removing other users avatar ([#657](https://github.com/stoatchat/stoatchat/issues/657)) ([d56135e](https://github.com/stoatchat/stoatchat/commit/d56135e0cbc713884c9378832952f7ad490fa315))
* default video resolution is a non-existent size ([#601](https://github.com/stoatchat/stoatchat/issues/601)) ([0698e11](https://github.com/stoatchat/stoatchat/commit/0698e115e8d003d615e468c4fb9654e6bbc9107f)), closes [#588](https://github.com/stoatchat/stoatchat/issues/588)
* **docs:** Update GitHub links ([#647](https://github.com/stoatchat/stoatchat/issues/647)) ([b830631](https://github.com/stoatchat/stoatchat/commit/b830631bd25a546844b7bdd30386084bb365e4de))
* don't use a bitop for OR ([#676](https://github.com/stoatchat/stoatchat/issues/676)) ([5701b5c](https://github.com/stoatchat/stoatchat/commit/5701b5c18c513f796af365169ceaea372a22638c))
* Fix typo for p256dh in vapid notification flow ([#622](https://github.com/stoatchat/stoatchat/issues/622)) ([a80ad1c](https://github.com/stoatchat/stoatchat/commit/a80ad1cbe58b8af5e45751e51d94d93c1cea1c9f))
* improve generated openapi.json ([#584](https://github.com/stoatchat/stoatchat/issues/584)) ([52ed510](https://github.com/stoatchat/stoatchat/commit/52ed5100c2446e0b261085639e123e7e124cab2c))
* no node state set on channel creation ([#653](https://github.com/stoatchat/stoatchat/issues/653)) ([24d0d2b](https://github.com/stoatchat/stoatchat/commit/24d0d2b7266f6f8a692d0a52704acfecf517674c))
* only show first line on commit messages ([#696](https://github.com/stoatchat/stoatchat/issues/696)) ([91783b9](https://github.com/stoatchat/stoatchat/commit/91783b906697fc85305dee683f7c15dda55f0c50))
* pass &str to Reference ([#697](https://github.com/stoatchat/stoatchat/issues/697)) ([ccda6f5](https://github.com/stoatchat/stoatchat/commit/ccda6f5c53ee043705f7ff6b5f6c393f020781de))
* redis_url vs redis_uri in config ([#666](https://github.com/stoatchat/stoatchat/issues/666)) ([b0b728f](https://github.com/stoatchat/stoatchat/commit/b0b728fb0dbc9ee28360301de1c3ea501bbbff1d))
* replace some links and Revolt mentions to current Stoat ([#515](https://github.com/stoatchat/stoatchat/issues/515)) ([d629e89](https://github.com/stoatchat/stoatchat/commit/d629e89304be2f0011e189293b278f07d346aa7d))
* send push notifications for DM and group messages ([#660](https://github.com/stoatchat/stoatchat/issues/660)) ([52c0d2f](https://github.com/stoatchat/stoatchat/commit/52c0d2f266b76d8975bba2d5e75c62bb30149c45))
* store server id in redis and in room metadata to be able to delete voice state in all scenarios ([#656](https://github.com/stoatchat/stoatchat/issues/656)) ([49c6289](https://github.com/stoatchat/stoatchat/commit/49c628958070e4f0a5edc764d3b48158589219d9))
* uname is missing from crond ([#675](https://github.com/stoatchat/stoatchat/issues/675)) ([dc4438b](https://github.com/stoatchat/stoatchat/commit/dc4438bc3c7b2cad8d442b3cd438afb9ed566a5e))
## [0.11.5](https://github.com/stoatchat/stoatchat/compare/v0.11.4...v0.11.5) (2026-02-17) ## [0.11.5](https://github.com/stoatchat/stoatchat/compare/v0.11.4...v0.11.5) (2026-02-17)
Generated
+2879 -1880
View File
File diff suppressed because it is too large Load Diff
+162 -9
View File
@@ -1,5 +1,5 @@
[workspace] [workspace]
resolver = "2" resolver = "3"
members = [ members = [
"crates/delta", "crates/delta",
@@ -20,24 +20,130 @@ lto = true
[workspace.dependencies] [workspace.dependencies]
# Async # Async
async-trait = "0.1.89" async-trait = "0.1.89"
tokio = { version = "1.49.0", features = ["macros", "rt"] } tokio = "1.49.0"
async-channel = "2.3.1"
futures = "0.3.32"
async-std = "1.8.0"
async-tungstenite = "0.17.0"
futures-locks = "0.7.1"
async-lock = "2.8.0"
async-recursion = "1.0.4"
# Error Handling # Error Handling
anyhow = "1.0.100" anyhow = "1.0.100"
thiserror = "2.0.18" thiserror = "2.0.18"
sentry = "0.31.5"
sentry-anyhow = "0.38.1"
# Other Utilities # Data Validation
uuid = { version = "1.19.0", features = ["v4"] } regex = "1.12.3"
validator = "0.16"
# Data Types
uuid = "1.19.0"
ulid = "1.2.1"
nanoid = "0.4.0"
typenum = "1.17.0"
num_enum = "0.6.1"
bitfield = "0.13.2"
# Time
chrono = "0.4.15"
iso8601-timestamp = "0.2.10"
# Data Collections
lru = "0.16.3"
indexmap = "2.13.1"
dashmap = "5.2.0"
moka = "0.12.8"
lru_time_cache = "0.11.11"
deadqueue = "0.2.4"
# Web scraping
scraper = "0.20.0"
encoding_rs = "0.8.34"
# Mail
lettre = "0.10.0-alpha.4"
# HTTP Requests
reqwest = "0.13.2"
isahc = "1.7"
# Notifications
fcm_v1 = "0.3.0"
web-push = "0.10.0"
revolt_a2 = "0.10"
# Parsing
logos = "0.15"
# SVG rendering
usvg = "0.44.0"
resvg = "0.44.0"
tiny-skia = "0.11.4"
# Logging
log = "0.4.29"
pretty_env_logger = "0.4.0"
# Redis
redis-kiss = "0.1.4"
fred = "8.0.1"
# Serialisation
bincode = "1.3.3"
serde_json = "1.0.79"
rmp-serde = "1.0.0"
serde = { version = "1", features = ["derive"] }
strum_macros = "0.26.4"
# MongoDB
bson = { version = "2.1.0" }
mongodb = { version = "3.1.0" }
# S3
aws-config = "1.5.5"
aws-sdk-s3 = "1.46.0"
# Axum (HTTP server) # Axum (HTTP server)
axum-macros = "0.4.1" axum-macros = "0.4.1"
axum_typed_multipart = "0.12.1" axum_typed_multipart = "0.12.1"
axum = { version = "0.7.5", features = ["multipart"] } axum = "0.7.5"
tower-http = { version = "0.5.2", features = ["cors", "trace"] } axum-extra = "0.9"
tower-http = "0.5.2"
# Rocket (HTTP server)
rocket = "0.5.1"
rocket_empty = "0.1.1"
revolt_rocket_okapi = "0.10.0"
rocket_cors = { git = "https://github.com/lawliet89/rocket_cors", rev = "072d90359b23e9b291df6b672c07c93de9c46011" }
rocket_authifier = "1.0.16"
rocket_prometheus = "0.10.0-rc.3"
# Spec Generation
utoipa = "4.2.3"
revolt_okapi = "0.9.1"
schemars = "0.8.8"
utoipa-scalar = "0.1.0"
# Image Processing # Image Processing
jxl-oxide = { version = "0.12.5", features = ["image"] } jxl-oxide = "0.12.5"
image = "0.25.9" sha2 = "0.10.8"
kamadak-exif = "0.5.4"
webp = "0.3.0"
image = "0.25.2" # avif encode requires dav1d system library: features = ["avif-native"]
thumbhash = "0.1.0"
lcms2 = "6.1.1" # for color profile processing
# File processing
revolt_clamav-client = "0.1.5"
simdutf8 = "0.1.4"
# Content type processing
infer = "0.16.0"
ffprobe = "0.4.0"
imagesize = "0.13.0"
# OpenTelemetry # OpenTelemetry
tracing = "0.1.44" tracing = "0.1.44"
@@ -47,4 +153,51 @@ tracing-subscriber = { version = "0.3.22", features = [
opentelemetry = { version = "0.31.0", features = ["logs"] } opentelemetry = { version = "0.31.0", features = ["logs"] }
opentelemetry_sdk = { version = "0.31.0", features = ["logs"] } opentelemetry_sdk = { version = "0.31.0", features = ["logs"] }
opentelemetry-otlp = { version = "0.31.0", features = ["logs"] } opentelemetry-otlp = { version = "0.31.0", features = ["logs"] }
opentelemetry-appender-tracing = { version = "0.31.1" } opentelemetry-appender-tracing = "0.31.1"
# Authifier
authifier = "1.0.16"
# RabbitMQ
lapin = "4.7.1"
# Voice
livekit-api = "0.4.4"
livekit-protocol = "0.7.4"
livekit-runtime = "0.4.0"
# Other Utilities
once_cell = "1.9.0"
config = "0.13.3"
cached = "0.44.0"
rand = "0.8.5"
base64 = "0.21.3"
decancer = "3.3.3"
linkify = "0.8.1"
url-escape = "0.1.1"
revolt_optional_struct = "0.2.0"
unicode-segmentation = "1.10.1"
querystring = "1.1.0"
tempfile = "3.12.0"
aes-gcm = "0.10.3"
auto_ops = "0.3.0"
url = "2.2.2"
impl_ops = "0.1.1"
lazy_static = "1.5.0"
mime = "0.3.17"
futures-lite = "2.6.1"
# Build Dependencies
vergen = "7.5.0"
# Local packages
revolt-coalesced = { version = "0.13.7", path = "crates/core/coalesced" }
revolt-config = { version = "0.13.7", path = "crates/core/config" }
revolt-database = { version = "0.13.7", path = "crates/core/database" }
revolt-files = { version = "0.13.7", path = "crates/core/files" }
revolt-models = { version = "0.13.7", path = "crates/core/models" }
revolt-parser = { version = "0.13.7", path = "crates/core/parser" }
revolt-permissions = { version = "0.13.7", path = "crates/core/permissions" }
revolt-presence = { version = "0.13.7", path = "crates/core/presence" }
revolt-ratelimits = { version = "0.13.7", path = "crates/core/ratelimits" }
revolt-result = { version = "0.13.7", path = "crates/core/result" }
+25 -21
View File
@@ -76,6 +76,14 @@ mise install
mise build mise build
``` ```
> [!TIP]
> You can override `BUILDER` in your `.env` file to run cargo with mold if you installed it:
>
> ```bash
> # .env
> BUILDER = "mold --run cargo"
> ```
A default configuration `Revolt.toml` is present in this project that is suited for development. A default configuration `Revolt.toml` is present in this project that is suited for development.
If you'd like to change anything, create a `Revolt.overrides.toml` file and specify relevant variables. If you'd like to change anything, create a `Revolt.overrides.toml` file and specify relevant variables.
@@ -112,7 +120,7 @@ If you'd like to change anything, create a `Revolt.overrides.toml` file and spec
> - "14672:15672" > - "14672:15672"
> ``` > ```
> >
> And corresponding Revolt configuration: > With the corresponding Revolt configuration:
> >
> ```toml > ```toml
> # Revolt.overrides.toml > # Revolt.overrides.toml
@@ -124,32 +132,25 @@ If you'd like to change anything, create a `Revolt.overrides.toml` file and spec
> [rabbit] > [rabbit]
> port = 14072 > port = 14072
> ``` > ```
>
> And mise configuration
>
> ```bash
> #.env
> DATABASE_PORT = "14017"
> RABBIT_PORT = "14072"
> REDIS_PORT = "14079"
> ```
Then continue: Then continue:
```bash ```bash
# start other necessary services cp livekit.example.yml livekit.yml
docker compose up -d
# run the API server mise start
cargo run --bin revolt-delta
# run the events server
cargo run --bin revolt-bonfire
# run the file server
cargo run --bin revolt-autumn
# run the proxy server
cargo run --bin revolt-january
# run the tenor proxy
cargo run --bin revolt-gifbox
# run the push daemon (not usually needed in regular development)
cargo run --bin revolt-pushd
# hint:
# mold -run <cargo build, cargo run, etc...>
# mold -run ./scripts/start.sh
``` ```
You can start a web client by doing the following: You can start a web client by doing the following in another terminal:
```bash ```bash
# if you do not have yarn yet and have a modern Node.js: # if you do not have yarn yet and have a modern Node.js:
@@ -163,6 +164,9 @@ cd stoat-web
When signing up, go to http://localhost:14080 to find confirmation/password reset emails. When signing up, go to http://localhost:14080 to find confirmation/password reset emails.
To stop all services, hit (CTRL + c) in the terminal you ran `mise start` and run `mise docker:stop`
## Deployment Guide ## Deployment Guide
### Cutting new crate releases ### Cutting new crate releases
@@ -198,7 +202,7 @@ If you have bumped the crate versions, proceed to [GitHub releases](https://gith
First, start the required services: First, start the required services:
```sh ```sh
docker compose -f docker-compose.db.yml up -d docker compose up -d
``` ```
Now run tests for whichever database: Now run tests for whichever database:
+13 -2
View File
@@ -8,10 +8,20 @@ services:
# MongoDB # MongoDB
database: database:
image: mongo image: mongo
command: mongod --replSet rs0
ports: ports:
- "27017:27017" - "27017:27017"
volumes: volumes:
- ./.data/db:/data/db - ./.data/db:/data/db
extra_hosts:
- "host.docker.internal:host-gateway"
healthcheck:
test: echo "try { rs.status() } catch (err) { rs.initiate({_id:'rs0',members:[{_id:0,host:'127.0.0.1:27017'}]}) }" | mongosh --port 27017 --quiet
interval: 5s
timeout: 30s
start_period: 0s
start_interval: 1s
retries: 30
ulimits: ulimits:
nofile: nofile:
soft: 65536 soft: 65536
@@ -19,11 +29,12 @@ services:
# MinIO # MinIO
minio: minio:
image: minio/minio image: firstfinger/minio:latest
command: server /data #command: server /data
environment: environment:
MINIO_ROOT_USER: minioautumn MINIO_ROOT_USER: minioautumn
MINIO_ROOT_PASSWORD: minioautumn MINIO_ROOT_PASSWORD: minioautumn
MINIO_REGION: minio
volumes: volumes:
- ./.data/minio:/data - ./.data/minio:/data
ports: ports:
+26 -30
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-bonfire" name = "revolt-bonfire"
version = "0.11.5" version = "0.13.7"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
edition = "2021" edition = "2021"
publish = false publish = false
@@ -9,42 +9,38 @@ publish = false
[dependencies] [dependencies]
# util # util
log = "*" log = { workspace = true }
sentry = "0.31.5" sentry = { workspace = true }
lru = "0.7.6" lru = { workspace = true }
ulid = "0.5.0" ulid = { workspace = true }
once_cell = "1.9.0" once_cell = { workspace = true }
redis-kiss = "0.1.4" redis-kiss = { workspace = true }
lru_time_cache = "0.11.11" lru_time_cache = { workspace = true }
async-channel = "2.3.1" async-channel = { workspace = true }
# parsing # parsing
querystring = "1.1.0" querystring = { workspace = true }
regex = "1.11.1" regex = { workspace = true }
# serde # serde
bincode = "1.3.3" bincode = { workspace = true }
serde_json = "1.0.79" serde_json = { workspace = true }
rmp-serde = "1.0.0" rmp-serde = { workspace = true }
serde = "1.0.136" serde = { workspace = true }
# async # async
futures = "0.3.21" futures = { workspace = true }
async-tungstenite = { version = "0.17.0", features = ["async-std-runtime"] } async-tungstenite = { workspace = true, features = ["async-std-runtime"] }
async-std = { version = "1.8.0", features = [ async-std = { workspace = true }
"tokio1",
"tokio02",
"attributes",
] }
# core # core
authifier = { version = "1.0.16" } authifier = { workspace = true }
revolt-result = { path = "../core/result" } revolt-result = { workspace = true }
revolt-models = { path = "../core/models" } revolt-models = { workspace = true }
revolt-config = { path = "../core/config" } revolt-config = { workspace = true }
revolt-database = { path = "../core/database", features = ["voice"] } revolt-database = { workspace = true, features = ["voice"] }
revolt-permissions = { path = "../core/permissions" } revolt-permissions = { workspace = true }
revolt-presence = { path = "../core/presence", features = ["redis-is-patched"] } revolt-presence = { workspace = true, features = ["redis-is-patched"] }
# redis # redis
fred = { version = "8.0.1", features = ["subscriber-client"] } fred = { workspace = true, features = ["subscriber-client"] }
+1
View File
@@ -1,6 +1,7 @@
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet};
use futures::future::join_all; use futures::future::join_all;
use redis_kiss::AsyncCommands;
use revolt_database::{ use revolt_database::{
events::client::{EventV1, ReadyPayloadFields}, events::client::{EventV1, ReadyPayloadFields},
util::permissions::DatabasePermissionQuery, util::permissions::DatabasePermissionQuery,
+2 -4
View File
@@ -1,7 +1,5 @@
use std::{ use std::{
collections::{HashMap, HashSet}, collections::{HashMap, HashSet}, num::NonZeroUsize, sync::Arc, time::Duration
sync::Arc,
time::Duration,
}; };
use async_std::sync::{Mutex, RwLock}; use async_std::sync::{Mutex, RwLock};
@@ -57,7 +55,7 @@ impl Default for Cache {
members: Default::default(), members: Default::default(),
servers: Default::default(), servers: Default::default(),
seen_events: LruCache::new(20), seen_events: LruCache::new(NonZeroUsize::new(20).unwrap()),
} }
} }
} }
+58 -5
View File
@@ -5,7 +5,7 @@ use authifier::AuthifierEvent;
use fred::{ use fred::{
error::RedisErrorKind, error::RedisErrorKind,
interfaces::{ClientLike, EventInterface, PubsubInterface}, interfaces::{ClientLike, EventInterface, PubsubInterface},
types::RedisConfig, types::{ReconnectPolicy, RedisConfig},
}; };
use futures::{ use futures::{
channel::oneshot, channel::oneshot,
@@ -13,7 +13,7 @@ use futures::{
stream::{SplitSink, SplitStream}, stream::{SplitSink, SplitStream},
FutureExt, SinkExt, StreamExt, TryStreamExt, FutureExt, SinkExt, StreamExt, TryStreamExt,
}; };
use redis_kiss::{PayloadType, REDIS_PAYLOAD_TYPE, REDIS_URI}; use redis_kiss::{get_connection, AsyncCommands, PayloadType, REDIS_PAYLOAD_TYPE, REDIS_URI};
use revolt_config::report_internal_error; use revolt_config::report_internal_error;
use revolt_database::{ use revolt_database::{
events::{client::EventV1, server::ClientMessage}, events::{client::EventV1, server::ClientMessage},
@@ -32,6 +32,7 @@ use sentry::Level;
use crate::config::{ProtocolConfiguration, WebsocketHandshakeCallback}; use crate::config::{ProtocolConfiguration, WebsocketHandshakeCallback};
use crate::events::state::{State, SubscriptionStateChange}; use crate::events::state::{State, SubscriptionStateChange};
use revolt_models::v0;
type WsReader = SplitStream<WebSocketStream<TcpStream>>; type WsReader = SplitStream<WebSocketStream<TcpStream>>;
type WsWriter = SplitSink<WebSocketStream<TcpStream>, async_tungstenite::tungstenite::Message>; type WsWriter = SplitSink<WebSocketStream<TcpStream>, async_tungstenite::tungstenite::Message>;
@@ -128,6 +129,14 @@ pub async fn client(db: &'static Database, stream: TcpStream, addr: SocketAddr)
return; return;
} }
let slowmodes = fetch_user_slowmodes(&user_id).await.unwrap_or_default();
if !slowmodes.is_empty() {
let event = EventV1::UserSlowmodes { slowmodes };
if report_internal_error!(write.send(config.encode(&event)).await).is_err() {
return;
}
}
// Create presence session. // Create presence session.
let (first_session, session_id) = create_session(&user_id, 0).await; let (first_session, session_id) = create_session(&user_id, 0).await;
@@ -225,9 +234,9 @@ async fn listener(
.unwrap_or(REDIS_URI.to_string()); .unwrap_or(REDIS_URI.to_string());
let redis_config = RedisConfig::from_url(&url).unwrap(); let redis_config = RedisConfig::from_url(&url).unwrap();
let subscriber = match report_internal_error!( let mut builder = fred::types::Builder::from_config(redis_config);
fred::types::Builder::from_config(redis_config).build_subscriber_client() builder.set_policy(ReconnectPolicy::new_exponential(8, 100, 30_000, 2));
) { let subscriber = match report_internal_error!(builder.build_subscriber_client()) {
Ok(subscriber) => subscriber, Ok(subscriber) => subscriber,
Err(_) => return, Err(_) => return,
}; };
@@ -236,16 +245,21 @@ async fn listener(
return; return;
} }
// Let Fred automatically re-subscribe to tracked channels on reconnect.
subscriber.manage_subscriptions();
// Handle Redis connection dropping // Handle Redis connection dropping
let (clean_up_s, clean_up_r) = async_channel::bounded(1); let (clean_up_s, clean_up_r) = async_channel::bounded(1);
let clean_up_s = Arc::new(Mutex::new(clean_up_s)); let clean_up_s = Arc::new(Mutex::new(clean_up_s));
subscriber.on_error(move |err| { subscriber.on_error(move |err| {
warn!("Redis subscriber error: {:?}", err);
if let RedisErrorKind::Canceled = err.kind() { if let RedisErrorKind::Canceled = err.kind() {
let clean_up_s = clean_up_s.clone(); let clean_up_s = clean_up_s.clone();
spawn(async move { spawn(async move {
clean_up_s.lock().await.send(()).await.ok(); clean_up_s.lock().await.send(()).await.ok();
}); });
} }
// Transient errors (IO, timeout) are handled by the reconnect policy.
Ok(()) Ok(())
}); });
@@ -523,3 +537,42 @@ async fn worker(
} }
} }
} }
async fn fetch_user_slowmodes(user_id: &str) -> Option<Vec<v0::ChannelSlowmode>> {
let mut conn = get_connection().await.ok()?.into_inner();
let idx_key = format!("slowmode_idx:{}", user_id);
let channel_ids: Vec<String> = conn.smembers(&idx_key).await.unwrap_or_default();
if channel_ids.is_empty() {
return Some(vec![]);
}
// Bulk fetch all TTLs in one round trip
let mut pipe = redis_kiss::redis::pipe();
for channel_id in &channel_ids {
pipe.ttl(format!("slowmode:{}:{}", user_id, channel_id));
}
let ttls: Vec<i64> = pipe.query_async(&mut conn).await.unwrap_or_default();
// Partition into alive/expired in one pass
let mut slowmodes = vec![];
let mut expired = vec![];
for (channel_id, ttl) in channel_ids.iter().zip(ttls.iter()) {
if *ttl > 0 {
slowmodes.push(v0::ChannelSlowmode {
channel_id: channel_id.clone(),
duration: *ttl as u64,
retry_after: *ttl as u64,
});
} else {
expired.push(channel_id.as_str());
}
}
// Bulk remove all expired members in one SREM call
if !expired.is_empty() {
conn.srem::<_, _, ()>(&idx_key, expired).await.ok();
}
Some(slowmodes)
}
+5 -5
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-coalesced" name = "revolt-coalesced"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "MIT" license = "MIT"
authors = ["Paul Makles <me@insrt.uk>", "Zomatree <me@zomatree.live>"] authors = ["Paul Makles <me@insrt.uk>", "Zomatree <me@zomatree.live>"]
@@ -15,12 +15,12 @@ cache = ["dep:lru"]
default = ["tokio"] default = ["tokio"]
[dependencies] [dependencies]
tokio = { version = "1.47.0", features = ["sync"], optional = true } tokio = { workspace = true, features = ["sync"], optional = true }
indexmap = { version = "2.13.0", optional = true } indexmap = { workspace = true, optional = true }
lru = { version = "0.16.3", optional = true } lru = { workspace = true, optional = true }
[dev-dependencies] [dev-dependencies]
tokio = { version = "1.47.0", features = [ tokio = { workspace = true, features = [
"rt", "rt",
"rt-multi-thread", "rt-multi-thread",
"macros", "macros",
+12 -12
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-config" name = "revolt-config"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "MIT" license = "MIT"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -18,24 +18,24 @@ default = ["test", "sentry"]
[dependencies] [dependencies]
# Utility # Utility
config = "0.13.3" config = { workspace = true }
cached = "0.44.0" cached = { workspace = true }
once_cell = "1.18.0" once_cell = { workspace = true }
# Serde # Serde
serde = { version = "1", features = ["derive"] } serde = { workspace = true }
# Async # Async
futures-locks = "0.7.1" futures-locks = { workspace = true }
async-std = { version = "1.8.0", features = ["attributes"], optional = true } async-std = { workspace = true, features = ["attributes"], optional = true }
# Logging # Logging
log = "0.4.14" log = { workspace = true }
pretty_env_logger = "0.4.0" pretty_env_logger = { workspace = true }
# Sentry # Sentry
sentry = { version = "0.31.5", optional = true } sentry = { workspace = true, optional = true }
sentry-anyhow = { version = "0.38.1", optional = true } sentry-anyhow = { workspace = true, optional = true }
# Core # Core
revolt-result = { version = "0.11.5", path = "../result", optional = true } revolt-result = { workspace = true, optional = true }
+15
View File
@@ -30,6 +30,11 @@ host = "rabbit"
port = 5672 port = 5672
username = "rabbituser" username = "rabbituser"
password = "rabbitpass" password = "rabbitpass"
default_exchange = "revolt.default"
[rabbit.queues]
acks = "internal.ack"
events = "internal.event"
[api] [api]
@@ -78,6 +83,8 @@ call_ring_duration = 30
[api.livekit.nodes] [api.livekit.nodes]
[api.users] [api.users]
# Minimum allowed length of usernames
min_username_length = 2
[pushd] [pushd]
# this changes the names of the queues to not overlap # this changes the names of the queues to not overlap
@@ -130,6 +137,8 @@ pkcs8 = ""
key_id = "" key_id = ""
team_id = "" team_id = ""
[january]
blocked_domains = []
[files] [files]
# Encryption key for stored files # Encryption key for stored files
@@ -313,6 +322,12 @@ emojis = 500_000
# default: 5 # default: 5
process_message_delay_limit = 5 process_message_delay_limit = 5
[features.legal_links]
# URLs for legal documents
terms_of_service = ""
privacy_policy = ""
guidelines = ""
[sentry] [sentry]
# Configuration for Sentry error reporting # Configuration for Sentry error reporting
api = "" api = ""
+26
View File
@@ -122,12 +122,20 @@ pub struct Database {
pub redis_pubsub: Option<String>, pub redis_pubsub: Option<String>,
} }
#[derive(Deserialize, Debug, Clone)]
pub struct RabbitQueues {
pub acks: String,
pub events: String,
}
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
pub struct Rabbit { pub struct Rabbit {
pub host: String, pub host: String,
pub port: u16, pub port: u16,
pub username: String, pub username: String,
pub password: String, pub password: String,
pub default_exchange: String,
pub queues: RabbitQueues,
} }
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
@@ -231,6 +239,7 @@ pub struct LiveKitNode {
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
pub struct ApiUsers { pub struct ApiUsers {
pub early_adopter_cutoff: Option<u64>, pub early_adopter_cutoff: Option<u64>,
pub min_username_length: usize,
} }
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
@@ -301,6 +310,11 @@ impl Pushd {
} }
} }
#[derive(Deserialize, Debug, Clone)]
pub struct January {
pub blocked_domains: Vec<String>,
}
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
pub struct FilesLimit { pub struct FilesLimit {
pub min_file_size: usize, pub min_file_size: usize,
@@ -376,6 +390,16 @@ pub struct FeaturesLimitsCollection {
pub roles: HashMap<String, FeaturesLimits>, pub roles: HashMap<String, FeaturesLimits>,
} }
#[derive(Deserialize, Debug, Clone)]
pub struct LegalLinks {
/// Terms of Service URL
pub terms_of_service: String,
/// Privacy Policy URL
pub privacy_policy: String,
/// Guidelines URL
pub guidelines: String,
}
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
pub struct FeaturesAdvanced { pub struct FeaturesAdvanced {
#[serde(default)] #[serde(default)]
@@ -393,6 +417,7 @@ impl Default for FeaturesAdvanced {
#[derive(Deserialize, Debug, Clone)] #[derive(Deserialize, Debug, Clone)]
pub struct Features { pub struct Features {
pub limits: FeaturesLimitsCollection, pub limits: FeaturesLimitsCollection,
pub legal_links: LegalLinks,
pub webhooks_enabled: bool, pub webhooks_enabled: bool,
pub mass_mentions_send_notifications: bool, pub mass_mentions_send_notifications: bool,
pub mass_mentions_enabled: bool, pub mass_mentions_enabled: bool,
@@ -420,6 +445,7 @@ pub struct Settings {
pub hosts: Hosts, pub hosts: Hosts,
pub api: Api, pub api: Api,
pub pushd: Pushd, pub pushd: Pushd,
pub january: January,
pub files: Files, pub files: Files,
pub features: Features, pub features: Features,
pub sentry: Sentry, pub sentry: Sentry,
+46 -54
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-database" name = "revolt-database"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -32,80 +32,72 @@ default = ["mongodb", "async-std-runtime", "tasks"]
[dependencies] [dependencies]
# Core # Core
revolt-config = { version = "0.11.5", path = "../config", features = [ revolt-config = { workspace = true, features = ["report-macros"] }
"report-macros", revolt-result = { workspace = true }
] } revolt-models = { workspace = true, features = ["validator"] }
revolt-result = { version = "0.11.5", path = "../result" } revolt-presence = { workspace = true }
revolt-models = { version = "0.11.5", path = "../models", features = [ revolt-permissions = { workspace = true, features = ["serde", "bson"] }
"validator", revolt-parser = { workspace = true }
] } revolt-coalesced = { workspace = true }
revolt-presence = { version = "0.11.5", path = "../presence" }
revolt-permissions = { version = "0.11.5", path = "../permissions", features = [
"serde",
"bson",
] }
revolt-parser = { version = "0.11.5", path = "../parser" }
# Utility # Utility
log = "0.4" log = { workspace = true }
lru = "0.11.0" lru = { workspace = true }
rand = "0.8.5" rand = { workspace = true }
ulid = "1.0.0" ulid = { workspace = true }
nanoid = "0.4.0" nanoid = { workspace = true }
base64 = "0.21.3" base64 = { workspace = true }
once_cell = "1.17" once_cell = { workspace = true }
indexmap = "1.9.1" indexmap = { workspace = true }
decancer = "1.6.2" decancer = { workspace = true }
deadqueue = "0.2.4" deadqueue = { workspace = true }
linkify = { optional = true, version = "0.8.1" } linkify = { workspace = true, optional = true }
url-escape = { optional = true, version = "0.1.1" } url-escape = { workspace = true, optional = true }
validator = { version = "0.16", features = ["derive"] } validator = { workspace = true, features = ["derive"] }
isahc = { optional = true, version = "1.7", features = ["json"] } isahc = { workspace = true, features = ["json"], optional = true }
# Serialisation # Serialisation
serde_json = "1" serde_json = { workspace = true }
revolt_optional_struct = "0.2.0" revolt_optional_struct = { workspace = true }
serde = { version = "1", features = ["derive"] } serde = { workspace = true }
iso8601-timestamp = { version = "0.2.10", features = ["serde", "bson"] } iso8601-timestamp = { workspace = true, features = ["serde", "bson"] }
# Events # Events
redis-kiss = { version = "0.1.4" } redis-kiss = { workspace = true }
# Database # Database
bson = { optional = true, version = "2.1.0" } bson = { workspace = true, optional = true }
mongodb = { optional = true, version = "3.1.0" } mongodb = { workspace = true, optional = true }
# Database Migration # Database Migration
unicode-segmentation = "1.10.1" unicode-segmentation = { workspace = true }
regex = "1" regex = { workspace = true }
# Async Language Features # Async Language Features
futures = "0.3.19" futures = { workspace = true }
async-lock = "2.8.0" async-lock = { workspace = true }
async-trait = "0.1.51" async-trait = { workspace = true }
async-recursion = "1.0.4" async-recursion = { workspace = true }
# Async # Async
async-std = { version = "1.8.0", features = ["attributes"], optional = true } async-std = { workspace = true, features = ["attributes"], optional = true }
# Axum Impl # Axum Impl
axum = { version = "0.7.5", optional = true } axum = { workspace = true, optional = true }
# Rocket Impl # Rocket Impl
schemars = { version = "0.8.8", optional = true } schemars = { workspace = true, optional = true }
rocket = { version = "0.5.1", default-features = false, features = [ rocket = { workspace = true, features = ["json"], optional = true }
"json", revolt_okapi = { workspace = true, optional = true }
], optional = true } revolt_rocket_okapi = { workspace = true, optional = true }
revolt_okapi = { version = "0.9.1", optional = true }
revolt_rocket_okapi = { version = "0.10.0", optional = true }
# Authifier # Authifier
authifier = { version = "1.0.16" } authifier = { workspace = true }
# RabbitMQ # RabbitMQ
amqprs = { version = "1.7.0" } lapin = { workspace = true, features = ["tokio"] }
# Voice # Voice
livekit-api = { version = "0.4.4", optional = true } livekit-api = { workspace = true, features = ["rustls-tls-native-roots"], optional = true }
livekit-protocol = { version = "0.4.0", optional = true } livekit-protocol = { workspace = true, optional = true }
livekit-runtime = { version = "0.3.1", features = ["tokio"], optional = true } livekit-runtime = { workspace = true, features = ["tokio"], optional = true }
+207 -110
View File
@@ -1,58 +1,92 @@
use std::collections::HashSet; use std::{
collections::HashSet,
sync::{Arc, OnceLock},
};
use crate::events::rabbit::*; use crate::events::{client::EventV1, rabbit::*};
use crate::User; use crate::User;
use amqprs::channel::{BasicPublishArguments, ExchangeDeclareArguments}; use lapin::{
use amqprs::connection::OpenConnectionArguments; options::BasicPublishOptions,
use amqprs::{channel::Channel, connection::Connection, error::Error as AMQPError}; protocol::basic::AMQPProperties,
use amqprs::{BasicProperties, FieldTable}; types::{AMQPValue, FieldTable},
BasicProperties, Channel, Connection, ConnectionProperties, Error as AMQPError,
};
use revolt_config::config;
use revolt_models::v0::PushNotification; use revolt_models::v0::PushNotification;
use revolt_presence::filter_online; use revolt_presence::filter_online;
use revolt_result::Result;
use serde_json::to_string; use serde_json::to_string;
static AMQP_INSTANCE: OnceLock<AMQP> = OnceLock::new();
pub fn get_amqp() -> &'static AMQP {
AMQP_INSTANCE.get().expect("No AMQP instance set.")
}
#[derive(Clone)] #[derive(Clone)]
pub struct AMQP { pub struct AMQP {
friend_request_accepted: Arc<Channel>,
friend_request_received: Arc<Channel>,
generic_message: Arc<Channel>,
message_sent: Arc<Channel>,
mass_mention_message_sent: Arc<Channel>,
ack_notification_message: Arc<Channel>,
dm_call_updated: Arc<Channel>,
process_ack: Arc<Channel>,
publish_event: Arc<Channel>,
#[allow(unused)] #[allow(unused)]
connection: Connection, connection: Arc<Connection>,
channel: Channel,
} }
impl AMQP { impl AMQP {
pub fn new(connection: Connection, channel: Channel) -> AMQP { pub async fn new(connection: Arc<Connection>) -> Self {
AMQP { let this = Self {
friend_request_accepted: Self::create_channel(&connection).await,
friend_request_received: Self::create_channel(&connection).await,
generic_message: Self::create_channel(&connection).await,
message_sent: Self::create_channel(&connection).await,
mass_mention_message_sent: Self::create_channel(&connection).await,
ack_notification_message: Self::create_channel(&connection).await,
dm_call_updated: Self::create_channel(&connection).await,
process_ack: Self::create_channel(&connection).await,
publish_event: Self::create_channel(&connection).await,
connection, connection,
channel, };
}
let _ = AMQP_INSTANCE.set(this.clone());
this
} }
pub async fn new_auto() -> AMQP { pub async fn new_auto() -> Self {
let config = revolt_config::config().await; let config = revolt_config::config().await;
let connection = Connection::open(&OpenConnectionArguments::new( let connection = Arc::new(
&config.rabbit.host, Connection::connect(
config.rabbit.port, &format!(
&config.rabbit.username, "amqp://{}:{}@{}:{}",
&config.rabbit.password, &config.rabbit.username,
)) &config.rabbit.password,
.await &config.rabbit.host,
.expect("Failed to connect to RabbitMQ"); &config.rabbit.port,
),
let channel = connection ConnectionProperties::default(),
.open_channel(None)
.await
.expect("Failed to open RabbitMQ channel");
channel
.exchange_declare(
ExchangeDeclareArguments::new(&config.pushd.exchange, "direct")
.durable(true)
.finish(),
) )
.await .await
.expect("Failed to declare exchange"); .expect("Failed to connect to RabbitMQ"),
);
AMQP::new(connection, channel) Self::new(connection).await
}
async fn create_channel(connection: &Connection) -> Arc<Channel> {
Arc::new(
connection
.create_channel()
.await
.expect("Failed to create channel"),
)
} }
pub async fn friend_request_accepted( pub async fn friend_request_accepted(
@@ -72,19 +106,20 @@ impl AMQP {
config.pushd.get_fr_accepted_routing_key(), config.pushd.get_fr_accepted_routing_key(),
payload payload
); );
self.channel
self.friend_request_accepted
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") config.pushd.get_fr_accepted_routing_key().into(),
.with_persistence(true) BasicPublishOptions::default(),
.finish(), payload.as_bytes(),
payload.into(), AMQPProperties::default()
BasicPublishArguments::new( .with_content_type("application/json".into())
&config.pushd.exchange, .with_delivery_mode(2),
&config.pushd.get_fr_accepted_routing_key(),
),
) )
.await .await?;
Ok(())
} }
pub async fn friend_request_received( pub async fn friend_request_received(
@@ -105,19 +140,19 @@ impl AMQP {
payload payload
); );
self.channel self.friend_request_received
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") config.pushd.get_fr_received_routing_key().into(),
.with_persistence(true) BasicPublishOptions::default(),
.finish(), payload.as_bytes(),
payload.into(), AMQPProperties::default()
BasicPublishArguments::new( .with_content_type("application/json".into())
&config.pushd.exchange, .with_delivery_mode(2),
&config.pushd.get_fr_received_routing_key(),
),
) )
.await .await?;
Ok(())
} }
pub async fn generic_message( pub async fn generic_message(
@@ -142,19 +177,19 @@ impl AMQP {
payload payload
); );
self.channel self.generic_message
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") config.pushd.get_generic_routing_key().into(),
.with_persistence(true) BasicPublishOptions::default(),
.finish(), payload.as_bytes(),
payload.into(), AMQPProperties::default()
BasicPublishArguments::new( .with_content_type("application/json".into())
&config.pushd.exchange, .with_delivery_mode(2),
&config.pushd.get_generic_routing_key(),
),
) )
.await .await?;
Ok(())
} }
pub async fn message_sent( pub async fn message_sent(
@@ -185,19 +220,19 @@ impl AMQP {
payload payload
); );
self.channel self.message_sent
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") config.pushd.get_message_routing_key().into(),
.with_persistence(true) BasicPublishOptions::default(),
.finish(), payload.as_bytes(),
payload.into(), AMQPProperties::default()
BasicPublishArguments::new( .with_content_type("application/json".into())
&config.pushd.exchange, .with_delivery_mode(2),
&config.pushd.get_message_routing_key(),
),
) )
.await .await?;
Ok(())
} }
pub async fn mass_mention_message_sent( pub async fn mass_mention_message_sent(
@@ -220,19 +255,24 @@ impl AMQP {
routing_key, payload routing_key, payload
); );
self.channel self.mass_mention_message_sent
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") routing_key.into(),
.with_persistence(true) BasicPublishOptions::default(),
.finish(), payload.as_bytes(),
payload.into(), AMQPProperties::default()
BasicPublishArguments::new(&config.pushd.exchange, routing_key.as_str()), .with_content_type("application/json".into())
.with_delivery_mode(2),
) )
.await .await?;
Ok(())
} }
pub async fn ack_message( /// # Sends an ack to pushd to update badges on iPhones.
/// Not to be confused with the process_ack function, which handles sending all acks to crond for processing.
pub async fn ack_notification_message(
&self, &self,
user_id: String, user_id: String,
channel_id: String, channel_id: String,
@@ -252,23 +292,25 @@ impl AMQP {
config.pushd.ack_queue, payload config.pushd.ack_queue, payload
); );
let mut headers = FieldTable::new(); let mut headers = FieldTable::default();
headers.insert( headers.insert(
"x-deduplication-header".try_into().unwrap(), "x-deduplication-header".into(),
format!("{}-{}", &user_id, &channel_id).into(), AMQPValue::LongString(format!("{}-{}", &user_id, &channel_id).into()),
); );
self.channel self.ack_notification_message
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") config.pushd.ack_queue.into(),
.with_persistence(true) BasicPublishOptions::default(),
//.with_headers(headers) payload.as_bytes(),
.finish(), AMQPProperties::default()
payload.into(), .with_content_type("application/json".into())
BasicPublishArguments::new(&config.pushd.exchange, &config.pushd.ack_queue), .with_delivery_mode(2),
) )
.await .await?;
Ok(())
} }
/// # DM Call Update /// # DM Call Update
@@ -302,18 +344,73 @@ impl AMQP {
payload payload
); );
self.channel self.dm_call_updated
.basic_publish( .basic_publish(
BasicProperties::default() config.pushd.exchange.clone().into(),
.with_content_type("application/json") config.pushd.get_dm_call_routing_key().into(),
.with_persistence(true) BasicPublishOptions::default(),
.finish(), payload.as_bytes(),
payload.into(), AMQPProperties::default()
BasicPublishArguments::new( .with_content_type("application/json".into())
&config.pushd.exchange, .with_delivery_mode(2),
&config.pushd.get_dm_call_routing_key(),
),
) )
.await .await?;
Ok(())
}
/// # Send an ack to crond for processing
pub async fn process_ack(
&self,
user_id: &str,
channel_id: Option<&str>,
server_id: Option<&str>,
) -> Result<(), AMQPError> {
let config = revolt_config::config().await;
let payload = AckEventPayload {
user_id: user_id.to_string(),
channel_id: channel_id.map(|value| value.to_string()),
server_id: server_id.map(|value| value.to_string()),
};
let payload = to_string(&payload).unwrap();
info!(
"Sending ack processor event on exchange {}, channel {}: {}",
config.rabbit.default_exchange, config.rabbit.queues.acks, payload
);
self.process_ack
.basic_publish(
config.rabbit.default_exchange.clone().into(),
config.rabbit.queues.acks.into(),
BasicPublishOptions::default(),
payload.as_bytes(),
AMQPProperties::default()
.with_content_type("application/json".into())
.with_delivery_mode(2),
)
.await?;
Ok(())
}
pub async fn publish_event(&self, channel: String, event: &EventV1) -> Result<(), AMQPError> {
let mut headers = FieldTable::default();
headers.insert("c".into(), AMQPValue::LongString(channel.into()));
let config = config().await;
self.publish_event
.basic_publish(
config.rabbit.default_exchange.clone().into(),
config.rabbit.queues.events.into(),
BasicPublishOptions::default(),
&serde_json::to_vec(event).unwrap(),
BasicProperties::default().with_headers(headers),
)
.await?;
Ok(())
} }
} }
+2
View File
@@ -1,2 +1,4 @@
#[allow(clippy::module_inception)] #[allow(clippy::module_inception)]
pub mod amqp; pub mod amqp;
pub use amqp::{AMQP, get_amqp};
+95 -25
View File
@@ -3,10 +3,15 @@ use revolt_result::Error;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use revolt_models::v0::{ use revolt_models::v0::{
AppendMessage, Channel, ChannelUnread, ChannelVoiceState, Emoji, FieldsChannel, FieldsMember, FieldsMessage, FieldsRole, FieldsServer, FieldsUser, FieldsWebhook, Member, MemberCompositeKey, Message, PartialChannel, PartialMember, PartialMessage, PartialRole, PartialServer, PartialUser, PartialUserVoiceState, PartialWebhook, PolicyChange, RemovalIntention, Report, Server, User, UserSettings, UserVoiceState, Webhook AppendMessage, Channel, ChannelSlowmode, ChannelUnread, ChannelVoiceState, Emoji,
FieldsChannel, FieldsMember, FieldsMessage, FieldsRole, FieldsServer, FieldsUser,
FieldsWebhook, Member, MemberCompositeKey, Message, PartialChannel, PartialEmoji,
PartialMember, PartialMessage, PartialRole, PartialServer, PartialUser, PartialUserVoiceState,
PartialWebhook, PolicyChange, RemovalIntention, Report, Server, User, UserSettings,
UserVoiceState, Webhook,
}; };
use crate::Database; use crate::{Database, amqp::get_amqp};
/// Ping Packet /// Ping Packet
#[derive(Serialize, Deserialize, Debug, Clone)] #[derive(Serialize, Deserialize, Debug, Clone)]
@@ -51,9 +56,13 @@ impl Default for ReadyPayloadFields {
#[serde(tag = "type")] #[serde(tag = "type")]
pub enum EventV1 { pub enum EventV1 {
/// Multiple events /// Multiple events
Bulk { v: Vec<EventV1> }, Bulk {
v: Vec<EventV1>,
},
/// Error event /// Error event
Error { data: Error }, Error {
data: Error,
},
/// Successfully authenticated /// Successfully authenticated
Authenticated, Authenticated,
@@ -84,7 +93,9 @@ pub enum EventV1 {
}, },
/// Ping response /// Ping response
Pong { data: Ping }, Pong {
data: Ping,
},
/// New message /// New message
Message(Message), Message(Message),
@@ -105,7 +116,10 @@ pub enum EventV1 {
}, },
/// Delete message /// Delete message
MessageDelete { id: String, channel: String }, MessageDelete {
id: String,
channel: String,
},
/// New reaction to a message /// New reaction to a message
MessageReact { MessageReact {
@@ -131,7 +145,10 @@ pub enum EventV1 {
}, },
/// Bulk delete messages /// Bulk delete messages
BulkMessageDelete { channel: String, ids: Vec<String> }, BulkMessageDelete {
channel: String,
ids: Vec<String>,
},
/// New server /// New server
ServerCreate { ServerCreate {
@@ -139,7 +156,7 @@ pub enum EventV1 {
server: Server, server: Server,
channels: Vec<Channel>, channels: Vec<Channel>,
emojis: Vec<Emoji>, emojis: Vec<Emoji>,
voice_states: Vec<ChannelVoiceState> voice_states: Vec<ChannelVoiceState>,
}, },
/// Update existing server /// Update existing server
@@ -151,7 +168,9 @@ pub enum EventV1 {
}, },
/// Delete server /// Delete server
ServerDelete { id: String }, ServerDelete {
id: String,
},
/// Update existing server member /// Update existing server member
ServerMemberUpdate { ServerMemberUpdate {
@@ -187,10 +206,16 @@ pub enum EventV1 {
}, },
/// Server role deleted /// Server role deleted
ServerRoleDelete { id: String, role_id: String }, ServerRoleDelete {
id: String,
role_id: String,
},
/// Server roles ranks updated /// Server roles ranks updated
ServerRoleRanksUpdate { id: String, ranks: Vec<String> }, ServerRoleRanksUpdate {
id: String,
ranks: Vec<String>,
},
/// Update existing user /// Update existing user
UserUpdate { UserUpdate {
@@ -202,9 +227,15 @@ pub enum EventV1 {
}, },
/// Relationship with another user changed /// Relationship with another user changed
UserRelationship { id: String, user: User }, UserRelationship {
id: String,
user: User,
},
/// Settings updated remotely /// Settings updated remotely
UserSettingsUpdate { id: String, update: UserSettings }, UserSettingsUpdate {
id: String,
update: UserSettings,
},
/// User has been platform banned or deleted their account /// User has been platform banned or deleted their account
/// ///
@@ -215,12 +246,23 @@ pub enum EventV1 {
/// - Server Memberships /// - Server Memberships
/// ///
/// User flags are specified to explain why a wipe is occurring though not all reasons will necessarily ever appear. /// User flags are specified to explain why a wipe is occurring though not all reasons will necessarily ever appear.
UserPlatformWipe { user_id: String, flags: i32 }, UserPlatformWipe {
user_id: String,
flags: i32,
},
/// New emoji /// New emoji
EmojiCreate(Emoji), EmojiCreate(Emoji),
/// Update existing emoji
EmojiUpdate {
id: String,
data: PartialEmoji,
},
/// Delete emoji /// Delete emoji
EmojiDelete { id: String }, EmojiDelete {
id: String,
},
/// New report /// New report
ReportCreate(Report), ReportCreate(Report),
@@ -236,19 +278,33 @@ pub enum EventV1 {
}, },
/// Delete channel /// Delete channel
ChannelDelete { id: String }, ChannelDelete {
id: String,
},
/// User joins a group /// User joins a group
ChannelGroupJoin { id: String, user: String }, ChannelGroupJoin {
id: String,
user: String,
},
/// User leaves a group /// User leaves a group
ChannelGroupLeave { id: String, user: String }, ChannelGroupLeave {
id: String,
user: String,
},
/// User started typing in a channel /// User started typing in a channel
ChannelStartTyping { id: String, user: String }, ChannelStartTyping {
id: String,
user: String,
},
/// User stopped typing in a channel /// User stopped typing in a channel
ChannelStopTyping { id: String, user: String }, ChannelStopTyping {
id: String,
user: String,
},
/// User acknowledged message in channel /// User acknowledged message in channel
ChannelAck { ChannelAck {
@@ -268,7 +324,9 @@ pub enum EventV1 {
}, },
/// Delete webhook /// Delete webhook
WebhookDelete { id: String }, WebhookDelete {
id: String,
},
/// Auth events /// Auth events
Auth(AuthifierEvent), Auth(AuthifierEvent),
@@ -286,7 +344,7 @@ pub enum EventV1 {
user: String, user: String,
from: String, from: String,
to: String, to: String,
state: UserVoiceState state: UserVoiceState,
}, },
UserVoiceStateUpdate { UserVoiceStateUpdate {
id: String, id: String,
@@ -298,7 +356,11 @@ pub enum EventV1 {
from: String, from: String,
to: String, to: String,
token: String, token: String,
} },
/// User's active slowmodes
UserSlowmodes {
slowmodes: Vec<ChannelSlowmode>,
},
} }
impl EventV1 { impl EventV1 {
@@ -310,8 +372,16 @@ impl EventV1 {
#[cfg(debug_assertions)] #[cfg(debug_assertions)]
info!("Publishing event to {channel}: {self:?}"); info!("Publishing event to {channel}: {self:?}");
#[cfg(debug_assertions)] // #[cfg(debug_assertions)]
redis_kiss::publish(channel, self).await.unwrap(); // redis_kiss::publish(channel, self).await.unwrap();
if let Err(e) = get_amqp().publish_event(channel, &self).await {
if cfg!(debug_assertions) {
panic!("{e:?}");
} else {
log::error!("{e:?}");
};
};
} }
/// Publish user event /// Publish user event
@@ -78,3 +78,11 @@ pub struct AckPayload {
pub channel_id: String, pub channel_id: String,
pub message_id: String, pub message_id: String,
} }
/// This is not the same as the AckPayload above, as the state for this event is stored in redis to allow for state updates while the event is queued.
#[derive(Serialize, Deserialize, Debug)]
pub struct AckEventPayload {
pub user_id: String,
pub channel_id: Option<String>,
pub server_id: Option<String>,
}
@@ -1,6 +1,7 @@
#![allow(deprecated)] #![allow(deprecated)]
use std::{borrow::Cow, collections::HashMap}; use std::{borrow::Cow, collections::HashMap};
use redis_kiss::get_connection;
use revolt_config::config; use revolt_config::config;
use revolt_models::v0::{self, MessageAuthor}; use revolt_models::v0::{self, MessageAuthor};
use revolt_permissions::OverrideField; use revolt_permissions::OverrideField;
@@ -212,7 +213,7 @@ impl Channel {
role_permissions: HashMap::new(), role_permissions: HashMap::new(),
nsfw: data.nsfw.unwrap_or(false), nsfw: data.nsfw.unwrap_or(false),
voice: data.voice.map(|voice| voice.into()), voice: data.voice.map(|voice| voice.into()),
slowmode: None slowmode: None,
}, },
v0::LegacyServerChannelType::Voice => Channel::TextChannel { v0::LegacyServerChannelType::Voice => Channel::TextChannel {
id: id.clone(), id: id.clone(),
@@ -225,7 +226,7 @@ impl Channel {
role_permissions: HashMap::new(), role_permissions: HashMap::new(),
nsfw: data.nsfw.unwrap_or(false), nsfw: data.nsfw.unwrap_or(false),
voice: Some(data.voice.unwrap_or_default().into()), voice: Some(data.voice.unwrap_or_default().into()),
slowmode: None slowmode: None,
}, },
}; };
@@ -643,7 +644,7 @@ impl Channel {
} }
/// Acknowledge a message /// Acknowledge a message
pub async fn ack(&self, user: &str, message: &str) -> Result<()> { pub async fn ack(&self, user: &str, message: &str, amqp: &AMQP) -> Result<()> {
EventV1::ChannelAck { EventV1::ChannelAck {
id: self.id().to_string(), id: self.id().to_string(),
user: user.to_string(), user: user.to_string(),
@@ -652,17 +653,7 @@ impl Channel {
.private(user.to_string()) .private(user.to_string())
.await; .await;
#[cfg(feature = "tasks")] crate::util::acker::ack_channel(user, self.id(), message, amqp).await
crate::tasks::ack::queue_ack(
self.id().to_string(),
user.to_string(),
crate::tasks::ack::AckEvent::AckMessage {
id: message.to_string(),
},
)
.await;
Ok(())
} }
/// Remove user from a group /// Remove user from a group
@@ -2,6 +2,7 @@ use std::collections::HashSet;
use std::str::FromStr; use std::str::FromStr;
use once_cell::sync::Lazy; use once_cell::sync::Lazy;
use revolt_models::v0;
use revolt_result::Result; use revolt_result::Result;
use ulid::Ulid; use ulid::Ulid;
@@ -11,7 +12,7 @@ use crate::Database;
static PERMISSIBLE_EMOJIS: Lazy<HashSet<String>> = Lazy::new(|| { static PERMISSIBLE_EMOJIS: Lazy<HashSet<String>> = Lazy::new(|| {
include_str!("unicode_emoji.txt") include_str!("unicode_emoji.txt")
.split('\n') .split('\n')
.map(|x| x.into()) .map(|x| x.replace('\u{FE0F}', ""))
.collect() .collect()
}); });
@@ -41,6 +42,12 @@ auto_derived!(
Server { id: String }, Server { id: String },
Detached, Detached,
} }
/// Partial representation of an emoji
pub struct PartialEmoji {
#[serde(skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
}
); );
#[allow(clippy::disallowed_methods)] #[allow(clippy::disallowed_methods)]
@@ -75,13 +82,34 @@ impl Emoji {
db.detach_emoji(&self).await db.detach_emoji(&self).await
} }
/// Update an emoji
pub async fn update(&mut self, db: &Database, partial: PartialEmoji) -> Result<()> {
if let Some(name) = partial.name.clone() {
self.name = name;
}
db.update_emoji(&self.id, &partial).await?;
EventV1::EmojiUpdate {
id: self.id.clone(),
data: v0::PartialEmoji {
name: partial.name.clone(),
},
}
.p(self.parent().to_string())
.await;
Ok(())
}
/// Check whether we can use a given emoji /// Check whether we can use a given emoji
pub async fn can_use(db: &Database, emoji: &str) -> Result<bool> { pub async fn can_use(db: &Database, emoji: &str) -> Result<bool> {
if Ulid::from_str(emoji).is_ok() { if Ulid::from_str(emoji).is_ok() {
db.fetch_emoji(emoji).await?; db.fetch_emoji(emoji).await?;
Ok(true) Ok(true)
} else { } else {
Ok(PERMISSIBLE_EMOJIS.contains(emoji)) let sanitized_emoji = emoji.replace('\u{FE0F}', "");
Ok(PERMISSIBLE_EMOJIS.contains(&sanitized_emoji))
} }
} }
} }
@@ -1,6 +1,6 @@
use revolt_result::Result; use revolt_result::Result;
use crate::Emoji; use crate::{Emoji, PartialEmoji};
#[cfg(feature = "mongodb")] #[cfg(feature = "mongodb")]
mod mongodb; mod mongodb;
@@ -20,6 +20,9 @@ pub trait AbstractEmojis: Sync + Send {
/// Fetch emoji by their parent ids /// Fetch emoji by their parent ids
async fn fetch_emoji_by_parent_ids(&self, parent_ids: &[String]) -> Result<Vec<Emoji>>; async fn fetch_emoji_by_parent_ids(&self, parent_ids: &[String]) -> Result<Vec<Emoji>>;
/// Update emoji with new information
async fn update_emoji(&self, emoji_id: &str, partial: &PartialEmoji) -> Result<()>;
/// Detach an emoji by its id /// Detach an emoji by its id
async fn detach_emoji(&self, emoji: &Emoji) -> Result<()>; async fn detach_emoji(&self, emoji: &Emoji) -> Result<()>;
} }
@@ -1,7 +1,7 @@
use bson::Document; use bson::Document;
use revolt_result::Result; use revolt_result::Result;
use crate::Emoji; use crate::{Emoji, PartialEmoji};
use crate::MongoDb; use crate::MongoDb;
use super::AbstractEmojis; use super::AbstractEmojis;
@@ -46,6 +46,11 @@ impl AbstractEmojis for MongoDb {
) )
} }
/// Update emoji with new information
async fn update_emoji(&self, emoji_id: &str, partial: &PartialEmoji) -> Result<()> {
query!(self, update_one_by_id, COL, emoji_id, partial, vec![], None).map(|_| ())
}
/// Detach an emoji by its id /// Detach an emoji by its id
async fn detach_emoji(&self, emoji: &Emoji) -> Result<()> { async fn detach_emoji(&self, emoji: &Emoji) -> Result<()> {
self.col::<Document>(COL) self.col::<Document>(COL)
@@ -1,6 +1,6 @@
use revolt_result::Result; use revolt_result::Result;
use crate::Emoji; use crate::{Emoji, PartialEmoji};
use crate::EmojiParent; use crate::EmojiParent;
use crate::ReferenceDb; use crate::ReferenceDb;
@@ -54,6 +54,19 @@ impl AbstractEmojis for ReferenceDb {
.collect()) .collect())
} }
/// Update emoji with new information
async fn update_emoji(&self, emoji_id: &str, partial: &PartialEmoji) -> Result<()> {
let mut emojis = self.emojis.lock().await;
if let Some(emoji) = emojis.get_mut(emoji_id) {
if let Some(name) = partial.name.clone() {
emoji.name = name;
}
Ok(())
} else {
Err(create_error!(NotFound))
}
}
/// Detach an emoji by its id /// Detach an emoji by its id
async fn detach_emoji(&self, emoji: &Emoji) -> Result<()> { async fn detach_emoji(&self, emoji: &Emoji) -> Result<()> {
let mut emojis = self.emojis.lock().await; let mut emojis = self.emojis.lock().await;
File diff suppressed because it is too large Load Diff
@@ -46,8 +46,7 @@ auto_derived!(
width: isize, width: isize,
height: isize, height: isize,
thumbhash: Option<Vec<u8>>, thumbhash: Option<Vec<u8>>,
#[serde(default)] animated: Option<bool>,
animated: bool,
}, },
/// File is a video with specific dimensions /// File is a video with specific dimensions
Video { width: isize, height: isize }, Video { width: isize, height: isize },
@@ -17,6 +17,12 @@ pub trait AbstractAttachmentHashes: Sync + Send {
/// Update an attachment hash nonce value. /// Update an attachment hash nonce value.
async fn set_attachment_hash_nonce(&self, hash: &str, nonce: &str) -> Result<()>; async fn set_attachment_hash_nonce(&self, hash: &str, nonce: &str) -> Result<()>;
/// Updates the attachments animated metadata value.
///
/// The primary use for this is to update the metadata for existing uploaded files, this
/// can only be used for images.
async fn set_attachment_hash_animated(&self, hash: &str, animated: bool) -> Result<()>;
/// Delete attachment hash by id. /// Delete attachment hash by id.
async fn delete_attachment_hash(&self, id: &str) -> Result<()>; async fn delete_attachment_hash(&self, id: &str) -> Result<()>;
} }
@@ -48,6 +48,29 @@ impl AbstractAttachmentHashes for MongoDb {
.map_err(|_| create_database_error!("update_one", COL)) .map_err(|_| create_database_error!("update_one", COL))
} }
/// Updates the attachments animated metadata value.
///
/// The primary use for this is to update the metadata for existing uploaded files, this
/// can only be used for images.
async fn set_attachment_hash_animated(&self, hash: &str, animated: bool) -> Result<()> {
self.col::<FileHash>(COL)
.update_one(
doc! {
"_id": hash,
"metadata.type": "Image",
"metadata.animated": { "$exists": false },
},
doc! {
"$set": {
"metadata.animated": animated
}
},
)
.await
.map(|_| ())
.map_err(|_| create_database_error!("update_one", COL))
}
/// Delete attachment hash by id. /// Delete attachment hash by id.
async fn delete_attachment_hash(&self, id: &str) -> Result<()> { async fn delete_attachment_hash(&self, id: &str) -> Result<()> {
query!(self, delete_one_by_id, COL, id).map(|_| ()) query!(self, delete_one_by_id, COL, id).map(|_| ())
@@ -1,7 +1,6 @@
use revolt_result::Result; use revolt_result::Result;
use crate::FileHash; use crate::{FileHash, Metadata, ReferenceDb};
use crate::ReferenceDb;
use super::AbstractAttachmentHashes; use super::AbstractAttachmentHashes;
@@ -39,6 +38,29 @@ impl AbstractAttachmentHashes for ReferenceDb {
} }
} }
/// Updates the attachments animated metadata value.
///
/// The primary use for this is to update the metadata for existing uploaded files, this
/// can only be used for images.
async fn set_attachment_hash_animated(&self, hash: &str, animated: bool) -> Result<()> {
let mut hashes = self.file_hashes.lock().await;
if let Some(FileHash {
metadata:
Metadata::Image {
animated: Some(animated_metadata),
..
},
..
}) = hashes.get_mut(hash)
{
*animated_metadata = animated;
Ok(())
} else {
Err(create_error!(NotFound))
}
}
/// Delete attachment hash by id. /// Delete attachment hash by id.
async fn delete_attachment_hash(&self, id: &str) -> Result<()> { async fn delete_attachment_hash(&self, id: &str) -> Result<()> {
let mut file_hashes = self.file_hashes.lock().await; let mut file_hashes = self.file_hashes.lock().await;
@@ -70,6 +70,7 @@ auto_derived!(
LegacyGroupIcon, LegacyGroupIcon,
ChannelIcon, ChannelIcon,
ServerIcon, ServerIcon,
RoleIcon,
} }
/// Information about what the file was used for /// Information about what the file was used for
@@ -239,4 +240,23 @@ impl File {
) )
.await .await
} }
/// Use a file for a role icon
pub async fn use_role_icon(
db: &Database,
id: &str,
parent: &str,
uploader_id: &str,
) -> Result<File> {
db.find_and_use_attachment(
id,
"icons",
FileUsedFor {
id: parent.to_owned(),
object_type: FileUsedForType::RoleIcon,
},
uploader_id.to_owned(),
)
.await
}
} }
@@ -705,7 +705,7 @@ impl Message {
Some( Some(
PushNotification::from( PushNotification::from(
self.clone().into_model(user, member), self.clone().into_model(user, member),
Some(author), Some(author.clone()),
channel.to_owned().into(), channel.to_owned().into(),
) )
.await, .await,
@@ -713,7 +713,11 @@ impl Message {
self.clone(), self.clone(),
match channel { match channel {
Channel::DirectMessage { recipients, .. } Channel::DirectMessage { recipients, .. }
| Channel::Group { recipients, .. } => recipients.clone(), | Channel::Group { recipients, .. } => recipients
.iter()
.filter(|uid| *uid != author.id())
.cloned()
.collect(),
Channel::TextChannel { .. } => { Channel::TextChannel { .. } => {
self.mentions.clone().unwrap_or_default() self.mentions.clone().unwrap_or_default()
} }
@@ -249,7 +249,7 @@ impl AbstractMessages for ReferenceDb {
let mut messages = self.messages.lock().await; let mut messages = self.messages.lock().await;
if let Some(message) = messages.get_mut(id) { if let Some(message) = messages.get_mut(id) {
if let Some(users) = message.reactions.get_mut(emoji) { if let Some(users) = message.reactions.get_mut(emoji) {
users.remove(&user.to_string()); users.swap_remove(&user.to_string());
} }
Ok(()) Ok(())
@@ -262,7 +262,7 @@ impl AbstractMessages for ReferenceDb {
async fn clear_reaction(&self, id: &str, emoji: &str) -> Result<()> { async fn clear_reaction(&self, id: &str, emoji: &str) -> Result<()> {
let mut messages = self.messages.lock().await; let mut messages = self.messages.lock().await;
if let Some(message) = messages.get_mut(id) { if let Some(message) = messages.get_mut(id) {
message.reactions.remove(emoji); message.reactions.swap_remove(emoji);
Ok(()) Ok(())
} else { } else {
Err(create_error!(NotFound)) Err(create_error!(NotFound))
@@ -86,6 +86,9 @@ auto_derived_partial!(
/// Ranking of this role /// Ranking of this role
#[serde(default)] #[serde(default)]
pub rank: i64, pub rank: i64,
/// Custom icon attachment
#[serde(skip_serializing_if = "Option::is_none")]
pub icon: Option<File>,
}, },
"PartialRole" "PartialRole"
); );
@@ -129,6 +132,7 @@ auto_derived!(
/// Optional fields on server object /// Optional fields on server object
pub enum FieldsRole { pub enum FieldsRole {
Colour, Colour,
Icon,
} }
); );
@@ -305,6 +309,7 @@ impl Role {
colour: self.colour, colour: self.colour,
hoist: Some(self.hoist), hoist: Some(self.hoist),
rank: Some(self.rank), rank: Some(self.rank),
icon: self.icon,
} }
} }
@@ -318,6 +323,7 @@ impl Role {
colour: None, colour: None,
hoist: false, hoist: false,
permissions: Default::default(), permissions: Default::default(),
icon: None,
}; };
db.insert_role(&server.id, &role).await?; db.insert_role(&server.id, &role).await?;
@@ -367,6 +373,7 @@ impl Role {
pub fn remove_field(&mut self, field: &FieldsRole) { pub fn remove_field(&mut self, field: &FieldsRole) {
match field { match field {
FieldsRole::Colour => self.colour = None, FieldsRole::Colour => self.colour = None,
FieldsRole::Icon => self.icon = None,
} }
} }
@@ -77,7 +77,7 @@ impl AbstractServers for MongoDb {
}, },
doc! { doc! {
"$set": { "$set": {
"roles.".to_owned() + &role.id: to_document(role) "roles.".to_owned() + role.id.as_str(): to_document(role)
.map_err(|_| create_database_error!("to_document", "role"))? .map_err(|_| create_database_error!("to_document", "role"))?
} }
}, },
@@ -172,6 +172,7 @@ impl IntoDocumentPath for FieldsRole {
fn as_path(&self) -> Option<&'static str> { fn as_path(&self) -> Option<&'static str> {
Some(match self { Some(match self {
FieldsRole::Colour => "colour", FieldsRole::Colour => "colour",
FieldsRole::Icon => "icon",
}) })
} }
} }
+153 -34
View File
@@ -7,6 +7,7 @@ use futures::future::join_all;
use iso8601_timestamp::Timestamp; use iso8601_timestamp::Timestamp;
use once_cell::sync::Lazy; use once_cell::sync::Lazy;
use rand::seq::SliceRandom; use rand::seq::SliceRandom;
use regex::{Regex, RegexBuilder};
use revolt_config::{config, FeaturesLimits}; use revolt_config::{config, FeaturesLimits};
use revolt_models::v0::{self, UserBadges, UserFlags}; use revolt_models::v0::{self, UserBadges, UserFlags};
use revolt_presence::filter_online; use revolt_presence::filter_online;
@@ -163,6 +164,13 @@ pub static DISCRIMINATOR_SEARCH_SPACE: Lazy<HashSet<String>> = Lazy::new(|| {
set.into_iter().collect() set.into_iter().collect()
}); });
static BLOCKED_USERNAME_PATTERNS: Lazy<Regex> = Lazy::new(|| {
RegexBuilder::new("`{3}|(discord|rvlt|guilded|stt)\\.gg|(revolt|stoat)\\.chat|https?:\\/\\/")
.case_insensitive(true)
.build()
.unwrap()
});
#[allow(clippy::derivable_impls)] #[allow(clippy::derivable_impls)]
impl Default for User { impl Default for User {
fn default() -> Self { fn default() -> Self {
@@ -198,11 +206,13 @@ impl User {
I: Into<Option<String>>, I: Into<Option<String>>,
D: Into<Option<PartialUser>>, D: Into<Option<PartialUser>>,
{ {
let username = User::validate_username(username)?; let new_username = User::sanitise_username(&username).await?;
User::validate_username(&new_username)?;
let mut user = User { let mut user = User {
id: account_id.into().unwrap_or_else(|| Ulid::new().to_string()), id: account_id.into().unwrap_or_else(|| Ulid::new().to_string()),
discriminator: User::find_discriminator(db, &username, None).await?, discriminator: User::find_discriminator(db, &new_username, None).await?,
username, username: new_username.clone(),
last_acknowledged_policy_change: Timestamp::now_utc(), last_acknowledged_policy_change: Timestamp::now_utc(),
..Default::default() ..Default::default()
}; };
@@ -278,39 +288,40 @@ impl User {
} }
} }
/// Sanitise and validate a username can be used /// Validate a username
pub fn validate_username(username: String) -> Result<String> { ///
// Copy the username for validation /// This will check if the username is a blocked name or contains a blocked pattern.
fn validate_username(username: &str) -> Result<()> {
let username_lowercase = username.to_lowercase(); let username_lowercase = username.to_lowercase();
// Block homoglyphs const BLOCKED_USERNAMES: &[&str] = &["admin", "revolt", "stoat"];
if decancer::cure(&username_lowercase).into_str() != username_lowercase {
if BLOCKED_USERNAMES.contains(&username_lowercase.as_str())
|| BLOCKED_USERNAME_PATTERNS.is_match(username)
{
return Err(create_error!(InvalidUsername)); return Err(create_error!(InvalidUsername));
} }
// Ensure the username itself isn't blocked Ok(())
const BLOCKED_USERNAMES: &[&str] = &["admin", "revolt"]; }
for username in BLOCKED_USERNAMES { /// Sanitise a username
if username_lowercase == *username { ///
return Err(create_error!(InvalidUsername)); /// This will clean up Unicode homoglyphs and pad to the min username length with underscores.
} async fn sanitise_username(username: &str) -> Result<String> {
} let options = decancer::Options::default().retain_capitalization();
let mut username = decancer::cure(username, options)
.map_err(|_| create_error!(InvalidUsername))?
.to_string();
// Ensure none of the following substrings show up in the username let config = revolt_config::config().await;
const BLOCKED_SUBSTRINGS: &[&str] = &[ let username_length_diff = config
"```", .api
"discord.gg", .users
"rvlt.gg", .min_username_length
"guilded.gg", .saturating_sub(username.len());
"https://", if username_length_diff > 0 {
"http://", username.push_str(&"_".repeat(username_length_diff))
];
for substr in BLOCKED_SUBSTRINGS {
if username_lowercase.contains(substr) {
return Err(create_error!(InvalidUsername));
}
} }
Ok(username) Ok(username)
@@ -416,12 +427,14 @@ impl User {
/// Update a user's username /// Update a user's username
pub async fn update_username(&mut self, db: &Database, username: String) -> Result<()> { pub async fn update_username(&mut self, db: &Database, username: String) -> Result<()> {
let username = User::validate_username(username)?; let new_username = User::sanitise_username(&username).await?;
if self.username.to_lowercase() == username.to_lowercase() { User::validate_username(&new_username)?;
if self.username.to_lowercase() == new_username.to_lowercase() {
self.update( self.update(
db, db,
PartialUser { PartialUser {
username: Some(username), username: Some(new_username),
..Default::default() ..Default::default()
}, },
vec![], vec![],
@@ -434,12 +447,12 @@ impl User {
discriminator: Some( discriminator: Some(
User::find_discriminator( User::find_discriminator(
db, db,
&username, &new_username,
Some((self.discriminator.to_string(), self.id.clone())), Some((self.discriminator.to_string(), self.id.clone())),
) )
.await?, .await?,
), ),
username: Some(username), username: Some(new_username),
..Default::default() ..Default::default()
}, },
vec![], vec![],
@@ -825,3 +838,109 @@ impl User {
badges badges
} }
} }
#[cfg(test)]
mod tests {
use crate::User;
#[test]
fn username_validation_blocked_names() {
let username_admin = "Admin";
let username_revolt = "Revolt";
let username_stoat = "Stoat";
let username_allowed = "Allowed";
assert!(User::validate_username(username_admin).is_err());
assert!(User::validate_username(username_revolt).is_err());
assert!(User::validate_username(username_stoat).is_err());
assert!(User::validate_username(username_allowed).is_ok());
}
#[test]
fn username_validation_blocked_patterns() {
let username_grave = "```_test";
let username_discord = "discord.gg_test";
let username_rvlt = "rvlt.gg_test";
let username_guilded = "guilded.gg_test";
let username_stt = "stt.gg_test";
let username_revolt = "revolt.chat_test";
let username_stoat = "stoat.chat_test";
let username_http = "http://_test";
let username_https = "https://_test";
assert!(User::validate_username(username_grave).is_err());
assert!(User::validate_username(username_discord).is_err());
assert!(User::validate_username(username_rvlt).is_err());
assert!(User::validate_username(username_guilded).is_err());
assert!(User::validate_username(username_stt).is_err());
assert!(User::validate_username(username_revolt).is_err());
assert!(User::validate_username(username_stoat).is_err());
assert!(User::validate_username(username_http).is_err());
assert!(User::validate_username(username_https).is_err());
}
#[async_std::test]
async fn username_sanitisation_clean() {
let username_clean = "Test";
let username_clean_sanitised = User::sanitise_username(username_clean).await;
assert!(username_clean_sanitised.is_ok());
assert_eq!(username_clean, username_clean_sanitised.unwrap());
}
#[async_std::test]
async fn username_sanitisation_homoglyphs() {
let username_homoglyphs = "𝔽𝕌Ňℕy";
let username_homoglyphs_sanitised =
User::sanitise_username(username_homoglyphs).await.unwrap();
assert_ne!(username_homoglyphs, username_homoglyphs_sanitised);
assert_eq!("funny", username_homoglyphs_sanitised);
}
#[async_std::test]
async fn username_sanitisation_padding() {
let username_padding = "a";
let username = User::sanitise_username(username_padding).await.unwrap();
assert_eq!("a_", username);
}
#[async_std::test]
async fn create_user() {
use revolt_result::Result;
database_test!(|db| async move {
let mut created_clean = User::create(&db, "Test".to_string(), None, None)
.await
.unwrap();
assert_eq!("Test", created_clean.username);
created_clean
.update_username(&db, "Test2".to_string())
.await
.unwrap();
assert_eq!("Test2", created_clean.username);
let created_invalid_result: Result<_> =
User::create(&db, "stoat.chat".to_string(), None, None).await;
assert!(created_invalid_result.is_err());
let mut updated_invalid = User::create(&db, "Test".to_string(), None, None)
.await
.unwrap();
let updated_invalid_update_result = updated_invalid
.update_username(&db, "http://test".to_string())
.await;
assert!(updated_invalid_update_result.is_err());
});
}
}
+6 -4
View File
@@ -105,7 +105,11 @@ pub async fn handle_ack_event(
if mentions_acked > 0 { if mentions_acked > 0 {
if let Err(err) = amqp if let Err(err) = amqp
.ack_message(user.to_string(), channel.to_string(), id.to_owned()) .ack_notification_message(
user.to_string(),
channel.to_string(),
id.to_owned(),
)
.await .await
{ {
revolt_config::capture_error(&err); revolt_config::capture_error(&err);
@@ -192,9 +196,7 @@ pub async fn handle_ack_event(
.expect("Failed to fetch channel from db"); .expect("Failed to fetch channel from db");
if let TextChannel { server, .. } = channel { if let TextChannel { server, .. } = channel {
if let Err(err) = if let Err(err) = amqp.mass_mention_message_sent(server, mass_mentions).await {
amqp.mass_mention_message_sent(server, mass_mentions).await
{
revolt_config::capture_error(&err); revolt_config::capture_error(&err);
} }
} else { } else {
+77
View File
@@ -0,0 +1,77 @@
use redis_kiss::{get_connection, AsyncCommands};
use revolt_permissions::{calculate_channel_permissions, ChannelPermission};
use revolt_result::{Result, ToRevoltError};
use crate::{events::client::EventV1, Channel, Database, Server, User, AMQP};
pub async fn ack_channel(user: &str, channel: &str, message: &str, amqp: &AMQP) -> Result<()> {
let mut redis = get_connection()
.await
.map_err(|_| create_error!(InternalError))?;
let old: Option<String> = redis
.getset(format!("acker:{user}+{channel}"), message)
.await
.to_internal_error()?;
if old.is_none() || old.unwrap() == message {
amqp.process_ack(user, Some(channel), None)
.await
.to_internal_error()?;
}
Ok(())
}
pub async fn ack_server(user: &User, server: &Server, db: &Database, amqp: &AMQP) -> Result<()> {
let mut redis = get_connection()
.await
.map_err(|_| create_error!(InternalError))?;
let channels = db.fetch_channels(&server.channels).await?;
let query = crate::util::permissions::DatabasePermissionQuery::new(db, user).server(server);
for channel in channels {
let channel_id = channel.id();
let mut q = query.clone().channel(&channel);
if calculate_channel_permissions(&mut q)
.await
.has_channel_permission(ChannelPermission::ViewChannel)
{
let channel_last_msg = match &channel {
Channel::TextChannel {
last_message_id, ..
} => last_message_id,
_ => unreachable!(),
}
.clone();
if let Some(channel_last_msg) = channel_last_msg {
let old: Option<String> = redis
.getset(
format!("acker:{}+{}", user.id, channel_id),
&channel_last_msg,
)
.await
.to_internal_error()?;
if old.is_none() || old.unwrap() == channel_last_msg {
amqp.process_ack(&user.id, Some(channel_id), Some(&server.id))
.await
.to_internal_error()?;
EventV1::ChannelAck {
id: channel_id.to_string(),
user: user.id.clone(),
message_id: channel_last_msg,
}
.private(user.id.clone())
.await;
}
}
}
}
Ok(())
}
+11 -5
View File
@@ -190,7 +190,7 @@ impl From<crate::Channel> for Channel {
role_permissions, role_permissions,
nsfw, nsfw,
voice, voice,
slowmode slowmode,
} => Channel::TextChannel { } => Channel::TextChannel {
id, id,
server, server,
@@ -202,7 +202,7 @@ impl From<crate::Channel> for Channel {
role_permissions, role_permissions,
nsfw, nsfw,
voice: voice.map(|voice| voice.into()), voice: voice.map(|voice| voice.into()),
slowmode slowmode,
}, },
} }
} }
@@ -256,7 +256,7 @@ impl From<Channel> for crate::Channel {
role_permissions, role_permissions,
nsfw, nsfw,
voice, voice,
slowmode slowmode,
} => crate::Channel::TextChannel { } => crate::Channel::TextChannel {
id, id,
server, server,
@@ -268,7 +268,7 @@ impl From<Channel> for crate::Channel {
role_permissions, role_permissions,
nsfw, nsfw,
voice: voice.map(|voice| voice.into()), voice: voice.map(|voice| voice.into()),
slowmode slowmode,
}, },
} }
} }
@@ -307,7 +307,7 @@ impl From<PartialChannel> for crate::PartialChannel {
default_permissions: value.default_permissions, default_permissions: value.default_permissions,
last_message_id: value.last_message_id, last_message_id: value.last_message_id,
voice: value.voice.map(|voice| voice.into()), voice: value.voice.map(|voice| voice.into()),
slowmode: value.slowmode slowmode: value.slowmode,
} }
} }
} }
@@ -926,6 +926,7 @@ impl From<crate::Role> for Role {
colour: value.colour, colour: value.colour,
hoist: value.hoist, hoist: value.hoist,
rank: value.rank, rank: value.rank,
icon: value.icon.map(|f| f.into()),
} }
} }
} }
@@ -939,6 +940,7 @@ impl From<Role> for crate::Role {
colour: value.colour, colour: value.colour,
hoist: value.hoist, hoist: value.hoist,
rank: value.rank, rank: value.rank,
icon: value.icon.map(|f| f.into()),
} }
} }
} }
@@ -952,6 +954,7 @@ impl From<crate::PartialRole> for PartialRole {
colour: value.colour, colour: value.colour,
hoist: value.hoist, hoist: value.hoist,
rank: value.rank, rank: value.rank,
icon: value.icon.map(|f| f.into()),
} }
} }
} }
@@ -965,6 +968,7 @@ impl From<PartialRole> for crate::PartialRole {
colour: value.colour, colour: value.colour,
hoist: value.hoist, hoist: value.hoist,
rank: value.rank, rank: value.rank,
icon: value.icon.map(|f| f.into()),
} }
} }
} }
@@ -973,6 +977,7 @@ impl From<crate::FieldsRole> for FieldsRole {
fn from(value: crate::FieldsRole) -> Self { fn from(value: crate::FieldsRole) -> Self {
match value { match value {
crate::FieldsRole::Colour => FieldsRole::Colour, crate::FieldsRole::Colour => FieldsRole::Colour,
crate::FieldsRole::Icon => FieldsRole::Icon,
} }
} }
} }
@@ -981,6 +986,7 @@ impl From<FieldsRole> for crate::FieldsRole {
fn from(value: FieldsRole) -> Self { fn from(value: FieldsRole) -> Self {
match value { match value {
FieldsRole::Colour => crate::FieldsRole::Colour, FieldsRole::Colour => crate::FieldsRole::Colour,
FieldsRole::Icon => crate::FieldsRole::Icon,
} }
} }
} }
+1
View File
@@ -1,3 +1,4 @@
pub mod acker;
pub mod bridge; pub mod bridge;
pub mod bulk_permissions; pub mod bulk_permissions;
mod funcs; mod funcs;
+17 -17
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-files" name = "revolt-files"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -16,33 +16,33 @@ tracing = { workspace = true }
tokio = { workspace = true } tokio = { workspace = true }
async-trait = { workspace = true } async-trait = { workspace = true }
ffprobe = "0.4.0" ffprobe = { workspace = true }
imagesize = "0.13.0" imagesize = { workspace = true }
tempfile = "3.12.0" tempfile = { workspace = true }
base64 = "0.22.1" base64 = { workspace = true }
aes-gcm = "0.10.3" aes-gcm = { workspace = true }
typenum = "1.17.0" typenum = { workspace = true }
aws-config = "1.5.5" aws-config = { workspace = true }
aws-sdk-s3 = { version = "1.46.0", features = ["behavior-version-latest"] } aws-sdk-s3 = { workspace = true, features = ["behavior-version-latest"] }
revolt-config = { version = "0.11.5", path = "../config", features = [ revolt-config = { workspace = true, features = [
"report-macros", "report-macros",
] } ] }
revolt-result = { version = "0.11.5", path = "../result" } revolt-result = { workspace = true }
# image processing # image processing
jxl-oxide = { workspace = true } jxl-oxide = { workspace = true, features = ["image"] }
image = { workspace = true } image = { workspace = true }
# svg rendering # svg rendering
usvg = "0.44.0" usvg = { workspace = true }
resvg = "0.44.0" resvg = { workspace = true }
tiny-skia = "0.11.4" tiny-skia = { workspace = true }
# encoding # encoding
webp = "0.3.0" webp = { workspace = true }
[dev-dependencies] [dev-dependencies]
uuid = { workspace = true } uuid = { workspace = true, features = ["v4"] }
+14 -14
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-models" name = "revolt-models"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "MIT" license = "MIT"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -21,26 +21,26 @@ default = ["serde", "partials", "rocket"]
[dependencies] [dependencies]
# Core # Core
revolt-config = { version = "0.11.5", path = "../config" } revolt-config = { workspace = true }
revolt-permissions = { version = "0.11.5", path = "../permissions" } revolt-permissions = { workspace = true }
# Utility # Utility
regex = "1.11" regex = { workspace = true }
indexmap = "1.9.3" indexmap = { workspace = true }
once_cell = "1.17.1" once_cell = { workspace = true }
num_enum = "0.6.1" num_enum = { workspace = true }
# Rocket # Rocket
rocket = { optional = true, version = "0.5.0-rc.2", default-features = false } rocket = { workspace = true, optional = true }
# Serialisation # Serialisation
revolt_optional_struct = { version = "0.2.0", optional = true } revolt_optional_struct = { workspace = true, optional = true }
serde = { version = "1", features = ["derive"], optional = true } serde = { workspace = true, optional = true }
iso8601-timestamp = { version = "0.2.11", features = ["schema", "bson"] } iso8601-timestamp = { workspace = true, features = ["schema", "bson"] }
# Spec Generation # Spec Generation
schemars = { version = "0.8.8", optional = true, features = ["indexmap1"] } schemars = { workspace = true, features = ["indexmap2"], optional = true }
utoipa = { version = "4.2.3", optional = true } utoipa = { workspace = true, optional = true }
# Validation # Validation
validator = { version = "0.16.0", optional = true, features = ["derive"] } validator = { workspace = true, features = ["derive"], optional = true }
+6
View File
@@ -314,6 +314,12 @@ auto_derived!(
/// Only used when the user is the first one connected. /// Only used when the user is the first one connected.
pub recipients: Option<Vec<String>>, pub recipients: Option<Vec<String>>,
} }
pub struct ChannelSlowmode {
pub channel_id: String,
pub duration: u64,
pub retry_after: u64,
}
); );
impl Channel { impl Channel {
+18
View File
@@ -54,4 +54,22 @@ auto_derived!(
#[serde(default)] #[serde(default)]
pub nsfw: bool, pub nsfw: bool,
} }
/// Partial emoji representation
#[derive(Default)]
pub struct PartialEmoji {
#[cfg_attr(feature = "serde", serde(skip_serializing_if = "Option::is_none"))]
pub name: Option<String>,
}
/// Edit emoji information
#[cfg_attr(feature = "validator", derive(Validate))]
pub struct DataEditEmoji {
/// Emoji name
#[cfg_attr(
feature = "validator",
validate(length(min = 1, max = 32), regex = "RE_EMOJI")
)]
pub name: Option<String>,
}
); );
+1 -1
View File
@@ -51,7 +51,7 @@ auto_derived!(
width: usize, width: usize,
height: usize, height: usize,
thumbhash: Option<Vec<u8>>, thumbhash: Option<Vec<u8>>,
animated: bool, animated: Option<bool>,
}, },
/// File is a video with specific dimensions /// File is a video with specific dimensions
Video { width: usize, height: usize }, Video { width: usize, height: usize },
+3 -2
View File
@@ -261,7 +261,7 @@ auto_derived!(
pub nonce: Option<String>, pub nonce: Option<String>,
/// Message content to send /// Message content to send
#[cfg_attr(feature = "validator", validate(length(min = 0, max = 2000)))] #[cfg_attr(feature = "validator", validate(length(min = 0)))]
pub content: Option<String>, pub content: Option<String>,
/// Attachments to include in message /// Attachments to include in message
pub attachments: Option<Vec<String>>, pub attachments: Option<Vec<String>>,
@@ -345,7 +345,7 @@ auto_derived!(
#[cfg_attr(feature = "validator", derive(Validate))] #[cfg_attr(feature = "validator", derive(Validate))]
pub struct DataEditMessage { pub struct DataEditMessage {
/// New message content /// New message content
#[cfg_attr(feature = "validator", validate(length(min = 1, max = 2000)))] #[cfg_attr(feature = "validator", validate(length(min = 1)))]
pub content: Option<String>, pub content: Option<String>,
/// Embeds to include in the message /// Embeds to include in the message
#[cfg_attr(feature = "validator", validate(length(min = 0, max = 10)))] #[cfg_attr(feature = "validator", validate(length(min = 0, max = 10)))]
@@ -391,6 +391,7 @@ auto_derived!(
); );
/// Message Author Abstraction /// Message Author Abstraction
#[derive(Clone)]
pub enum MessageAuthor<'a> { pub enum MessageAuthor<'a> {
User(&'a User), User(&'a User),
Webhook(&'a Webhook), Webhook(&'a Webhook),
+9
View File
@@ -106,6 +106,9 @@ auto_derived_partial!(
/// Ranking of this role /// Ranking of this role
#[cfg_attr(feature = "serde", serde(default))] #[cfg_attr(feature = "serde", serde(default))]
pub rank: i64, pub rank: i64,
/// Role icon
#[cfg_attr(feature = "serde", serde(skip_serializing_if = "Option::is_none"))]
pub icon: Option<File>,
}, },
"PartialRole" "PartialRole"
); );
@@ -123,6 +126,7 @@ auto_derived!(
/// Optional fields on server object /// Optional fields on server object
pub enum FieldsRole { pub enum FieldsRole {
Colour, Colour,
Icon,
} }
/// Channel category /// Channel category
@@ -278,6 +282,11 @@ auto_derived!(
/// ///
/// **Removed** - no effect, use the edit server role positions route /// **Removed** - no effect, use the edit server role positions route
pub rank: Option<i64>, pub rank: Option<i64>,
/// Role icon
///
/// Provide an Autumn attachment Id.
#[cfg_attr(feature = "validator", validate(length(min = 1, max = 128)))]
pub icon: Option<String>,
/// Fields to remove from role object /// Fields to remove from role object
#[cfg_attr(feature = "serde", serde(default))] #[cfg_attr(feature = "serde", serde(default))]
pub remove: Vec<FieldsRole>, pub remove: Vec<FieldsRole>,
+2 -2
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-parser" name = "revolt-parser"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "MIT" license = "MIT"
authors = ["Zomatree <me@zomatree.live>", "Paul Makles <me@insrt.uk>"] authors = ["Zomatree <me@zomatree.live>", "Paul Makles <me@insrt.uk>"]
@@ -8,4 +8,4 @@ description = "Revolt Backend: Message Parser"
repository = "https://github.com/stoatchat/stoatchat" repository = "https://github.com/stoatchat/stoatchat"
[dependencies] [dependencies]
logos = { version = "0.15" } logos = { workspace = true }
+1 -1
View File
@@ -262,7 +262,7 @@ mod tests {
) )
.collect::<Vec<_>>(); .collect::<Vec<_>>();
assert_eq!(output.len(), 6); assert_eq!(output.len(), 7);
assert_eq!(output[0], MessageToken::CodeblockMarker(1)); assert_eq!(output[0], MessageToken::CodeblockMarker(1));
assert_eq!( assert_eq!(
output[1], output[1],
+10 -10
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-permissions" name = "revolt-permissions"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "MIT" license = "MIT"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -18,23 +18,23 @@ try-from-primitive = ["dep:num_enum"]
[dev-dependencies] [dev-dependencies]
# Async # Async
async-std = { version = "1.8.0", features = ["attributes"] } async-std = { workspace = true, features = ["attributes"] }
[dependencies] [dependencies]
# Core # Core
revolt-result = { version = "0.11.5", path = "../result" } revolt-result = { workspace = true }
# Utility # Utility
auto_ops = "0.3.0" auto_ops = { workspace = true }
once_cell = "1.17" once_cell = { workspace = true }
num_enum = { version = "0.6.1", optional = true } num_enum = { workspace = true, optional = true }
# Async # Async
async-trait = "0.1.51" async-trait = { workspace = true }
# Serialisation # Serialisation
serde = { version = "1", features = ["derive"], optional = true } serde = { workspace = true, optional = true }
bson = { version = "2.1.0", optional = true } bson = { workspace = true, optional = true }
# Spec Generation # Spec Generation
schemars = { version = "0.8.8", optional = true } schemars = { workspace = true, optional = true }
+7 -7
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-presence" name = "revolt-presence"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -14,16 +14,16 @@ redis-is-patched = []
[dev-dependencies] [dev-dependencies]
# Async # Async
async-std = { version = "1.8.0", features = ["attributes"] } async-std = { workspace = true, features = ["attributes"] }
# Config for loading Redis URI # Config for loading Redis URI
revolt-config = { version = "0.11.5", path = "../config" } revolt-config = { workspace = true }
[dependencies] [dependencies]
# Utility # Utility
log = "0.4.17" log = { workspace = true }
rand = "0.8.5" rand = { workspace = true }
once_cell = "1.17.1" once_cell = { workspace = true }
# Redis # Redis
redis-kiss = "0.1.4" redis-kiss = { workspace = true }
+12 -12
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-ratelimits" name = "revolt-ratelimits"
version = "0.11.5" version = "0.13.7"
edition = "2024" edition = "2024"
license = "MIT" license = "MIT"
authors = ["Zomatree <me@zomatree.live>", "Paul Makles <me@insrt.uk>"] authors = ["Zomatree <me@zomatree.live>", "Paul Makles <me@insrt.uk>"]
@@ -18,17 +18,17 @@ axum = ["dep:axum", "revolt-database/axum-impl"]
default = ["rocket", "axum"] default = ["rocket", "axum"]
[dependencies] [dependencies]
revolt-database = { version = "0.11.5", path = "../database" } revolt-database = { workspace = true }
revolt-result = { version = "0.11.5", path = "../result" } revolt-result = { workspace = true }
revolt-config = { version = "0.11.5", path = "../config" } revolt-config = { workspace = true }
rocket = { version = "0.5.1", optional = true } rocket = { workspace = true, optional = true }
revolt_rocket_okapi = { version = "0.10.0", optional = true } revolt_rocket_okapi = { workspace = true, optional = true }
axum = { version = "0.7.5", optional = true, features = ["macros"] } axum = { workspace = true, optional = true, features = ["macros"] }
serde = { version = "1", features = ["derive"] } serde = { workspace = true }
authifier = { version = "1.0.16" } authifier = { workspace = true }
dashmap = "5.2.0" dashmap = { workspace = true }
async-trait = "0.1.81" async-trait = { workspace = true }
log = "0.4" log = { workspace = true }
+3 -3
View File
@@ -1,12 +1,12 @@
use async_trait::async_trait; use async_trait::async_trait;
use log::info; use log::info;
use revolt_config::config;
use rocket::fairing::{Fairing, Info, Kind}; use rocket::fairing::{Fairing, Info, Kind};
use rocket::http::uri::Origin; use rocket::http::uri::Origin;
use rocket::http::{Method, Status}; use rocket::http::{Method, Status};
use rocket::request::{FromRequest, Outcome}; use rocket::request::{FromRequest, Outcome};
use rocket::serde::json::Json; use rocket::serde::json::Json;
use rocket::{Data, Request, Response, State}; use rocket::{Data, Request, Response, State};
use revolt_config::config;
use revolt_rocket_okapi::r#gen::OpenApiGenerator; use revolt_rocket_okapi::r#gen::OpenApiGenerator;
use revolt_rocket_okapi::request::{OpenApiFromRequest, RequestHeaderInput}; use revolt_rocket_okapi::request::{OpenApiFromRequest, RequestHeaderInput};
@@ -28,8 +28,8 @@ pub type RatelimitStorage = crate::ratelimiter::RatelimitStorage<RocketRequestKi
/// Find the remote IP of the client /// Find the remote IP of the client
fn to_ip(request: &'_ rocket::Request<'_>) -> String { fn to_ip(request: &'_ rocket::Request<'_>) -> String {
request request
.remote() .client_ip()
.map(|x| x.ip().to_string()) .map(|r| r.to_string())
.unwrap_or_default() .unwrap_or_default()
} }
+11 -11
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-result" name = "revolt-result"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "MIT" license = "MIT"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
@@ -22,21 +22,21 @@ default = ["serde", "sentry"]
[dependencies] [dependencies]
# Serialisation # Serialisation
serde_json = { version = "1", optional = true } serde_json = { workspace = true, optional = true }
serde = { version = "1", features = ["derive"], optional = true } serde = { workspace = true, optional = true }
# Spec Generation # Spec Generation
schemars = { version = "0.8.8", optional = true } schemars = { workspace = true, optional = true }
utoipa = { version = "4.2.3", optional = true } utoipa = { workspace = true, optional = true }
# Rocket # Rocket
rocket = { optional = true, version = "0.5.0-rc.2", default-features = false } rocket = { workspace = true, optional = true }
revolt_rocket_okapi = { version = "0.10.0", optional = true } revolt_rocket_okapi = { workspace = true, optional = true }
revolt_okapi = { version = "0.9.1", optional = true } revolt_okapi = { workspace = true, optional = true }
# utilities # utilities
log = "0.4" log = { workspace = true }
# Axum # Axum
axum = { version = "0.7.5", optional = true } axum = { workspace = true, optional = true }
sentry = { version = "0.31.5", optional = true } sentry = { workspace = true, optional = true }
+3 -4
View File
@@ -24,6 +24,7 @@ impl IntoResponse for Error {
ErrorType::UnknownChannel => StatusCode::NOT_FOUND, ErrorType::UnknownChannel => StatusCode::NOT_FOUND,
ErrorType::UnknownMessage => StatusCode::NOT_FOUND, ErrorType::UnknownMessage => StatusCode::NOT_FOUND,
ErrorType::UnknownAttachment => StatusCode::BAD_REQUEST, ErrorType::UnknownAttachment => StatusCode::BAD_REQUEST,
ErrorType::CannotDeleteMessage => StatusCode::FORBIDDEN,
ErrorType::CannotEditMessage => StatusCode::FORBIDDEN, ErrorType::CannotEditMessage => StatusCode::FORBIDDEN,
ErrorType::CannotJoinCall => StatusCode::BAD_REQUEST, ErrorType::CannotJoinCall => StatusCode::BAD_REQUEST,
ErrorType::TooManyAttachments { .. } => StatusCode::BAD_REQUEST, ErrorType::TooManyAttachments { .. } => StatusCode::BAD_REQUEST,
@@ -36,9 +37,7 @@ impl IntoResponse for Error {
ErrorType::NotInGroup => StatusCode::NOT_FOUND, ErrorType::NotInGroup => StatusCode::NOT_FOUND,
ErrorType::AlreadyPinned => StatusCode::BAD_REQUEST, ErrorType::AlreadyPinned => StatusCode::BAD_REQUEST,
ErrorType::NotPinned => StatusCode::BAD_REQUEST, ErrorType::NotPinned => StatusCode::BAD_REQUEST,
ErrorType::InSlowmode { ErrorType::InSlowmode { retry_after: _ } => StatusCode::TOO_MANY_REQUESTS,
retry_after: _,
} => StatusCode::TOO_MANY_REQUESTS,
ErrorType::CantCreateServers => StatusCode::FORBIDDEN, ErrorType::CantCreateServers => StatusCode::FORBIDDEN,
ErrorType::UnknownServer => StatusCode::NOT_FOUND, ErrorType::UnknownServer => StatusCode::NOT_FOUND,
@@ -78,7 +77,7 @@ impl IntoResponse for Error {
ErrorType::DuplicateNonce => StatusCode::CONFLICT, ErrorType::DuplicateNonce => StatusCode::CONFLICT,
ErrorType::VosoUnavailable => StatusCode::BAD_REQUEST, ErrorType::VosoUnavailable => StatusCode::BAD_REQUEST,
ErrorType::NotFound => StatusCode::NOT_FOUND, ErrorType::NotFound => StatusCode::NOT_FOUND,
ErrorType::NoEffect => StatusCode::OK, ErrorType::NoEffect => StatusCode::BAD_REQUEST,
ErrorType::FailedValidation { .. } => StatusCode::BAD_REQUEST, ErrorType::FailedValidation { .. } => StatusCode::BAD_REQUEST,
ErrorType::LiveKitUnavailable => StatusCode::BAD_REQUEST, ErrorType::LiveKitUnavailable => StatusCode::BAD_REQUEST,
ErrorType::NotConnected => StatusCode::BAD_REQUEST, ErrorType::NotConnected => StatusCode::BAD_REQUEST,
+1
View File
@@ -78,6 +78,7 @@ pub enum ErrorType {
UnknownChannel, UnknownChannel,
UnknownAttachment, UnknownAttachment,
UnknownMessage, UnknownMessage,
CannotDeleteMessage,
CannotEditMessage, CannotEditMessage,
CannotJoinCall, CannotJoinCall,
TooManyAttachments { TooManyAttachments {
+3 -4
View File
@@ -30,6 +30,7 @@ impl<'r> Responder<'r, 'static> for Error {
ErrorType::UnknownChannel => Status::NotFound, ErrorType::UnknownChannel => Status::NotFound,
ErrorType::UnknownMessage => Status::NotFound, ErrorType::UnknownMessage => Status::NotFound,
ErrorType::UnknownAttachment => Status::BadRequest, ErrorType::UnknownAttachment => Status::BadRequest,
ErrorType::CannotDeleteMessage => Status::Forbidden,
ErrorType::CannotEditMessage => Status::Forbidden, ErrorType::CannotEditMessage => Status::Forbidden,
ErrorType::CannotJoinCall => Status::BadRequest, ErrorType::CannotJoinCall => Status::BadRequest,
ErrorType::TooManyAttachments { .. } => Status::BadRequest, ErrorType::TooManyAttachments { .. } => Status::BadRequest,
@@ -42,9 +43,7 @@ impl<'r> Responder<'r, 'static> for Error {
ErrorType::NotInGroup => Status::NotFound, ErrorType::NotInGroup => Status::NotFound,
ErrorType::AlreadyPinned => Status::BadRequest, ErrorType::AlreadyPinned => Status::BadRequest,
ErrorType::NotPinned => Status::BadRequest, ErrorType::NotPinned => Status::BadRequest,
ErrorType::InSlowmode { ErrorType::InSlowmode { retry_after: _ } => Status::TooManyRequests,
retry_after: _,
} => Status::TooManyRequests,
ErrorType::InvalidFlagValue => Status::BadRequest, ErrorType::InvalidFlagValue => Status::BadRequest,
ErrorType::CantCreateServers => Status::Forbidden, ErrorType::CantCreateServers => Status::Forbidden,
@@ -84,7 +83,7 @@ impl<'r> Responder<'r, 'static> for Error {
ErrorType::NotAuthenticated => Status::Unauthorized, ErrorType::NotAuthenticated => Status::Unauthorized,
ErrorType::DuplicateNonce => Status::Conflict, ErrorType::DuplicateNonce => Status::Conflict,
ErrorType::NotFound => Status::NotFound, ErrorType::NotFound => Status::NotFound,
ErrorType::NoEffect => Status::Ok, ErrorType::NoEffect => Status::BadRequest,
ErrorType::FailedValidation { .. } => Status::BadRequest, ErrorType::FailedValidation { .. } => Status::BadRequest,
ErrorType::LiveKitUnavailable => Status::BadRequest, ErrorType::LiveKitUnavailable => Status::BadRequest,
ErrorType::NotAVoiceChannel => Status::BadRequest, ErrorType::NotAVoiceChannel => Status::BadRequest,
+21 -7
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-crond" name = "revolt-crond"
version = "0.11.5" version = "0.13.7"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
authors = ["Paul Makles <me@insrt.uk>"] authors = ["Paul Makles <me@insrt.uk>"]
edition = "2021" edition = "2021"
@@ -11,13 +11,27 @@ publish = false
[dependencies] [dependencies]
# Utility # Utility
log = "0.4" log = { workspace = true }
# Async # Async
tokio = { version = "1" } tokio = { workspace = true }
# Redis
redis-kiss = { workspace = true }
# RabbitMQ
lapin = { workspace = true }
futures-lite = { workspace = true }
# Processing
serde_json = { workspace = true }
revolt_optional_struct = { workspace = true }
serde = { workspace = true }
iso8601-timestamp = { workspace = true, features = ["serde", "bson"] }
# Core # Core
revolt-database = { version = "0.11.5", path = "../../core/database" } revolt-database = { workspace = true }
revolt-result = { version = "0.11.5", path = "../../core/result" } revolt-result = { workspace = true }
revolt-config = { version = "0.11.5", path = "../../core/config" } revolt-config = { workspace = true }
revolt-files = { version = "0.11.5", path = "../../core/files" } revolt-files = { workspace = true }
revolt-permissions = { workspace = true }
+6 -3
View File
@@ -1,7 +1,7 @@
use revolt_config::configure; use revolt_config::configure;
use revolt_database::DatabaseInfo; use revolt_database::{DatabaseInfo, AMQP};
use revolt_result::Result; use revolt_result::Result;
use tasks::{file_deletion, prune_dangling_files, prune_members}; use tasks::{acks, file_deletion, prune_dangling_files, prune_members};
use tokio::try_join; use tokio::try_join;
pub mod tasks; pub mod tasks;
@@ -11,10 +11,13 @@ async fn main() -> Result<()> {
configure!(crond); configure!(crond);
let db = DatabaseInfo::Auto.connect().await.expect("database"); let db = DatabaseInfo::Auto.connect().await.expect("database");
let amqp = AMQP::new_auto().await;
try_join!( try_join!(
file_deletion::task(db.clone()), file_deletion::task(db.clone()),
prune_dangling_files::task(db.clone()), prune_dangling_files::task(db.clone()),
prune_members::task(db.clone()) prune_members::task(db.clone()),
acks::task(db.clone(), amqp.clone()),
) )
.map(|_| ()) .map(|_| ())
} }
+164
View File
@@ -0,0 +1,164 @@
use futures_lite::stream::StreamExt;
use lapin::{
options::*,
types::FieldTable,
uri::{AMQPAuthority, AMQPQueryString, AMQPUri, AMQPUserInfo},
ConnectionBuilder, ConnectionProperties, ExchangeKind,
};
use log::{debug, info};
use redis_kiss::{get_connection, AsyncCommands, Conn as RedisConnection};
use revolt_config::config;
use revolt_database::{events::rabbit::AckEventPayload, Database, AMQP};
use revolt_result::{Result, ToRevoltError};
use serde_json;
pub async fn task(db: Database, amqp: AMQP) -> Result<()> {
let config = config().await;
let mut redis = get_connection()
.await
.expect("Failed to get redis connection");
let uri = AMQPUri {
scheme: lapin::uri::AMQPScheme::AMQP,
authority: AMQPAuthority {
userinfo: AMQPUserInfo {
username: config.rabbit.username,
password: config.rabbit.password,
},
host: config.rabbit.host,
port: config.rabbit.port,
},
vhost: "/".to_string(),
query: AMQPQueryString::default(),
};
let connection = ConnectionBuilder::new()
.expect("Builder")
.with_uri(uri)
.with_properties(ConnectionProperties::default())
.connect()
.await
.expect("Failed to connect to rabbitmq");
let reader_channel = connection
.create_channel()
.await
.expect("Failed to create channel");
reader_channel
.exchange_declare(
config.rabbit.default_exchange.clone().into(),
ExchangeKind::Topic,
ExchangeDeclareOptions {
durable: true,
..Default::default()
},
FieldTable::default(),
)
.await
.expect("Failed to declare exchange");
reader_channel
.queue_declare(
config.rabbit.queues.acks.clone().into(),
QueueDeclareOptions {
durable: true,
..Default::default()
},
FieldTable::default(),
)
.await
.expect("Failed to bind queue");
reader_channel
.queue_bind(
config.rabbit.queues.acks.clone().into(),
config.rabbit.default_exchange.into(),
config.rabbit.queues.acks.clone().into(),
QueueBindOptions::default(),
FieldTable::default(),
)
.await
.expect("Failed to bind channel");
let mut consumer = reader_channel
.basic_consume(
config.rabbit.queues.acks.into(),
"crond-ack-consumer".into(),
BasicConsumeOptions::default(),
FieldTable::default(),
)
.await
.expect("Failed to create consumer");
while let Some(delivery) = consumer.next().await {
if let Ok(delivery) = delivery {
let payload = serde_json::from_slice::<AckEventPayload>(&delivery.data);
if let Ok(payload) = payload {
debug!("Received ack event: {payload:?}");
if let Err(e) = process_channel_ack(
&db,
&amqp,
payload.user_id,
payload.channel_id.unwrap(),
&mut redis,
)
.await
{
revolt_config::capture_error(&e);
_ = delivery.reject(BasicRejectOptions { requeue: false }).await;
} else {
_ = delivery.ack(BasicAckOptions { multiple: false }).await;
}
} else {
revolt_config::capture_message(
format!("Failed to decode ack data: {:?}", delivery.data).as_str(),
revolt_config::Level::Error,
);
}
}
}
Ok(())
}
#[allow(clippy::disallowed_methods)]
async fn process_channel_ack(
db: &Database,
amqp: &AMQP,
user: String,
channel: String,
redis: &mut RedisConnection,
) -> Result<()> {
let message_id: Option<String> = redis
.get_del(format!("acker:{user}+{channel}"))
.await
.to_internal_error()?;
if let Some(message_id) = message_id {
let unread = db.fetch_unread(&user, &channel).await?;
let updated = db.acknowledge_message(&channel, &user, &message_id).await?;
info!("Set new state for ack: {}:{}:{}", channel, user, message_id);
if let (Some(before), Some(after)) = (unread, updated) {
let before_mentions = before.mentions.unwrap_or_default().len();
let after_mentions = after.mentions.unwrap_or_default().len();
if after_mentions < before_mentions {
if let Err(err) = amqp
.ack_notification_message(user.to_string(), channel.to_string(), message_id)
.await
{
revolt_config::capture_error(&err);
}
};
}
Ok(())
} else {
Err(message_id.to_internal_error().expect_err("no err"))
}
}
+15 -13
View File
@@ -11,22 +11,24 @@ pub async fn task(db: Database) -> Result<()> {
let files = db.fetch_deleted_attachments().await?; let files = db.fetch_deleted_attachments().await?;
for file in files { for file in files {
let count = db if let Some(hash) = &file.hash {
.count_file_hash_references(file.hash.as_ref().expect("no `hash` present")) let count = db
.await?; .count_file_hash_references(hash)
// No other files reference this file on disk anymore
if count <= 1 {
let file_hash = db
.fetch_attachment_hash(file.hash.as_ref().expect("no `hash` present"))
.await?; .await?;
// Delete from S3 // No other files reference this file on disk anymore
delete_from_s3(&file_hash.bucket_id, &file_hash.path).await?; if count <= 1 {
let file_hash = db
.fetch_attachment_hash(hash)
.await?;
// Delete the hash // Delete from S3
db.delete_attachment_hash(&file_hash.id).await?; delete_from_s3(&file_hash.bucket_id, &file_hash.path).await?;
info!("Deleted file hash {}", file_hash.id);
// Delete the hash
db.delete_attachment_hash(&file_hash.id).await?;
info!("Deleted file hash {}", file_hash.id);
}
} }
// Delete the file // Delete the file
+1
View File
@@ -1,3 +1,4 @@
pub mod acks;
pub mod file_deletion; pub mod file_deletion;
pub mod prune_dangling_files; pub mod prune_dangling_files;
pub mod prune_members; pub mod prune_members;
+26 -33
View File
@@ -1,47 +1,40 @@
[package] [package]
name = "revolt-pushd" name = "revolt-pushd"
version = "0.11.5" version = "0.13.7"
edition = "2021" edition = "2021"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
publish = false publish = false
[dependencies] [dependencies]
revolt-result = { version = "0.11.5", path = "../../core/result" } revolt-result = { workspace = true }
revolt-config = { version = "0.11.5", path = "../../core/config", features = [ revolt-config = { workspace = true, features = ["report-macros", "anyhow"] }
"report-macros", revolt-database = { workspace = true }
"anyhow", revolt-models = { workspace = true, features = ["validator"] }
] } revolt-presence = { workspace = true, features = ["redis-is-patched"] }
revolt-database = { version = "0.11.5", path = "../../core/database" } revolt-parser = { workspace = true }
revolt-models = { version = "0.11.5", path = "../../core/models", features = [
"validator",
] }
revolt-presence = { version = "0.11.5", path = "../../core/presence", features = [
"redis-is-patched",
] }
revolt-parser = { version = "0.11.5", path = "../../core/parser" }
anyhow = { version = "1.0.98" } anyhow = { workspace = true }
amqprs = { version = "1.7.0" } lapin = { workspace = true }
fcm_v1 = "0.3.0" fcm_v1 = { workspace = true }
web-push = "0.10.0" web-push = { workspace = true }
isahc = { optional = true, version = "1.7", features = ["json"] } isahc = { workspace = true, features = ["json"], optional = true }
revolt_a2 = { version = "0.10", default-features = false, features = ["ring"] } revolt_a2 = { workspace = true, features = ["ring"] }
redis-kiss = "0.1.4" redis-kiss = { workspace = true }
tokio = "1.39.2" tokio = { workspace = true }
async-trait = "0.1.81" async-trait = { workspace = true }
ulid = "1.0.0" ulid = { workspace = true }
authifier = "1.0.16" authifier = { workspace = true }
log = "0.4.11" log = { workspace = true }
pretty_env_logger = "0.4.0" pretty_env_logger = { workspace = true }
regex = "1.12.3" regex = { workspace = true }
#serialization #serialization
serde_json = "1" serde_json = { workspace = true }
revolt_optional_struct = "0.2.0" revolt_optional_struct = { workspace = true }
serde = { version = "1", features = ["derive"] } serde = { workspace = true }
iso8601-timestamp = { version = "0.2.10", features = ["serde", "bson"] } iso8601-timestamp = { workspace = true, features = ["serde", "bson"] }
base64 = "0.22.1" base64 = { workspace = true }
@@ -1,96 +1,69 @@
use crate::consumers::inbound::internal::*; use std::sync::Arc;
use amqprs::{
channel::{BasicPublishArguments, Channel}, use crate::utils::Consumer;
connection::Connection, use anyhow::Result;
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
#[derive(Clone)]
#[allow(unused)]
pub struct AckConsumer { pub struct AckConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
} }
impl Channeled for AckConsumer { #[async_trait]
fn get_connection(&self) -> Option<&Connection> { impl Consumer for AckConsumer {
if self.conn.is_none() { async fn create(
None db: Database,
} else { authifier_db: authifier::Database,
Some(self.conn.as_ref().unwrap()) connection: Arc<Connection>,
} channel: Arc<Channel>,
} ) -> Self {
Self {
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
}
impl AckConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> AckConsumer {
AckConsumer {
db, db,
authifier_db, authifier_db,
conn: None, connection,
channel: None, channel,
} }
} }
}
#[allow(unused_variables)] fn channel(&self) -> &Arc<Channel> {
#[async_trait] &self.channel
impl AsyncConsumer for AckConsumer { }
/// This consumer processes all acks the platform receives, and sends relevant badge updates to apple platforms. /// This consumer processes all acks the platform receives, and sends relevant badge updates to apple platforms.
async fn consume( async fn consume(&self, delivery: Delivery) -> Result<()> {
&mut self, let payload: AckPayload = serde_json::from_slice(&delivery.data)?;
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
let content = String::from_utf8(content).unwrap();
let payload: AckPayload = serde_json::from_str(content.as_str()).unwrap();
// Step 1: fetch unreads and don't continue if there's no unreads // Step 1: fetch unreads and don't continue if there's no unreads
#[allow(clippy::disallowed_methods)] // #[allow(clippy::disallowed_methods)]
let unreads = self.db.fetch_unread_mentions(&payload.user_id).await;
debug!("Processing unreads for {:}", &payload.user_id); debug!("Processing unreads for {:}", &payload.user_id);
if let Ok(u) = &unreads { let unreads = if let Ok(u) = self.db.fetch_unread_mentions(&payload.user_id).await {
if u.is_empty() { if u.is_empty() {
debug!( debug!(
"Discarding unread task (no mentions found) for {:}", "Discarding unread task (no mentions found) for {:}",
&payload.user_id &payload.user_id
); );
return; return Ok(());
} };
u
} else { } else {
return; return Ok(());
} };
if let Ok(sessions) = self.authifier_db.find_sessions(&payload.user_id).await { if let Ok(sessions) = self.authifier_db.find_sessions(&payload.user_id).await {
let config = revolt_config::config().await; let config = revolt_config::config().await;
// Step 2: find any apple sessions, since we don't need to calculate this for anything else. // Step 2: find any apple sessions, since we don't need to calculate this for anything else.
// If there's no apple sessions, we can return early // If there's no apple sessions, we can return early
let apple_sessions: Vec<&authifier::models::Session> = sessions let mut apple_sessions = sessions
.iter() .into_iter()
.filter(|session| { .filter(|session| {
if let Some(sub) = &session.subscription { if let Some(sub) = &session.subscription {
sub.endpoint == "apn" sub.endpoint == "apn"
@@ -98,19 +71,19 @@ impl AsyncConsumer for AckConsumer {
false false
} }
}) })
.collect(); .peekable();
if apple_sessions.is_empty() { if apple_sessions.peek().is_none() {
debug!( debug!(
"Discarding unread task (no apn sessions found) for {:}", "Discarding unread task (no apn sessions found) for {:}",
&payload.user_id &payload.user_id
); );
return; return Ok(());
} }
// Step 3: calculate the actual mention count, since we have to send it out // Step 3: calculate the actual mention count, since we have to send it out
let mut mention_count = 0; let mut mention_count = 0;
for u in &unreads.unwrap() { for u in &unreads {
mention_count += u.mentions.as_ref().unwrap().len() mention_count += u.mentions.as_ref().unwrap().len()
} }
@@ -123,26 +96,22 @@ impl AsyncConsumer for AckConsumer {
token: session.subscription.as_ref().unwrap().auth.clone(), token: session.subscription.as_ref().unwrap().auth.clone(),
extras: Default::default(), extras: Default::default(),
}; };
let raw_service_payload = serde_json::to_string(&service_payload); let payload = serde_json::to_string(&service_payload)?;
if let Ok(p) = raw_service_payload { log::debug!(
let args = BasicPublishArguments::new( "Publishing ack to apn session {}",
config.pushd.exchange.as_str(), session.subscription.as_ref().unwrap().auth
config.pushd.apn.queue.as_str(), );
)
.finish();
log::debug!( self.publish_message(
"Publishing ack to apn session {}", payload.as_bytes(),
session.subscription.as_ref().unwrap().auth &config.pushd.exchange,
); &config.pushd.apn.queue,
)
publish_message(self, p.into(), args).await; .await?;
} else {
log::warn!("Failed to serialize ack badge update payload!");
revolt_config::capture_error(&raw_service_payload.unwrap_err());
}
} }
} }
Ok(())
} }
} }
@@ -1,70 +1,44 @@
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use crate::consumers::inbound::internal::*; use crate::utils::Consumer;
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::Connection,
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use anyhow::Result; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use log::debug; use log::debug;
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
#[derive(Clone)]
#[allow(unused)]
pub struct DmCallConsumer { pub struct DmCallConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
} }
impl Channeled for DmCallConsumer { #[async_trait]
fn get_connection(&self) -> Option<&Connection> { impl Consumer for DmCallConsumer {
if self.conn.is_none() { async fn create(
None db: Database,
} else { authifier_db: authifier::Database,
Some(self.conn.as_ref().unwrap()) connection: Arc<Connection>,
} channel: Arc<Channel>,
} ) -> Self {
Self {
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
}
impl DmCallConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> DmCallConsumer {
DmCallConsumer {
db, db,
authifier_db, authifier_db,
conn: None, connection,
channel: None, channel,
} }
} }
async fn consume_event( fn channel(&self) -> &Arc<Channel> {
&mut self, &self.channel
_channel: &Channel, }
_deliver: Deliver,
_basic_properties: BasicProperties, /// This consumer handles delegating messages into their respective platform queues.
content: Vec<u8>, async fn consume(&self, delivery: Delivery) -> Result<()> {
) -> Result<()> { let _p: InternalDmCallPayload = serde_json::from_slice(&delivery.data)?;
let content = String::from_utf8(content)?;
let _p: InternalDmCallPayload = serde_json::from_str(content.as_str())?;
let payload = _p.payload; let payload = _p.payload;
debug!("Received dm call start/stop event"); debug!("Received dm call start/stop event");
@@ -107,36 +81,27 @@ impl DmCallConsumer {
extras: HashMap::new(), extras: HashMap::new(),
}; };
let args: BasicPublishArguments; let routing_key = match sub.endpoint.as_str() {
"apn" => &config.pushd.apn.queue,
"fcm" => &config.pushd.fcm.queue,
endpoint => {
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), endpoint.to_string());
if sub.endpoint == "apn" { &config.pushd.vapid.queue
args = BasicPublishArguments::new( }
config.pushd.exchange.as_str(), };
config.pushd.apn.queue.as_str(),
)
.finish();
} else if sub.endpoint == "fcm" {
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.fcm.queue.as_str(),
)
.finish();
} else {
// web push (vapid)
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.vapid.queue.as_str(),
)
.finish();
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), sub.endpoint.clone());
}
let payload = serde_json::to_string(&sendable)?; let payload = serde_json::to_string(&sendable)?;
publish_message(self, payload.into(), args).await; self.publish_message(
payload.as_bytes(),
&config.pushd.exchange,
routing_key,
)
.await?;
} }
} }
} }
@@ -145,24 +110,3 @@ impl DmCallConsumer {
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for DmCallConsumer {
/// This consumer handles delegating messages into their respective platform queues.
async fn consume(
&mut self,
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
warn!("Failed to process dm call start/stop event: {err:?}");
}
}
}
@@ -1,70 +1,44 @@
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use crate::consumers::inbound::internal::*; use crate::utils::Consumer;
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::Connection,
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use anyhow::Result; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use log::debug; use log::debug;
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
#[derive(Clone)]
#[allow(unused)]
pub struct FRAcceptedConsumer { pub struct FRAcceptedConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
} }
impl Channeled for FRAcceptedConsumer { #[async_trait]
fn get_connection(&self) -> Option<&Connection> { impl Consumer for FRAcceptedConsumer {
if self.conn.is_none() { async fn create(
None db: Database,
} else { authifier_db: authifier::Database,
Some(self.conn.as_ref().unwrap()) connection: Arc<Connection>,
} channel: Arc<Channel>,
} ) -> Self {
Self {
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
}
impl FRAcceptedConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> FRAcceptedConsumer {
FRAcceptedConsumer {
db, db,
authifier_db, authifier_db,
conn: None, connection,
channel: None, channel,
} }
} }
async fn consume_event( fn channel(&self) -> &Arc<Channel> {
&mut self, &self.channel
_channel: &Channel, }
_deliver: Deliver,
_basic_properties: BasicProperties, /// This consumer handles delegating messages into their respective platform queues.
content: Vec<u8>, async fn consume(&self, delivery: Delivery) -> Result<()> {
) -> Result<()> { let payload: FRAcceptedPayload = serde_json::from_slice(&delivery.data)?;
let content = String::from_utf8(content)?;
let payload: FRAcceptedPayload = serde_json::from_str(content.as_str())?;
debug!("Received FR accept event"); debug!("Received FR accept event");
@@ -80,36 +54,23 @@ impl FRAcceptedConsumer {
extras: HashMap::new(), extras: HashMap::new(),
}; };
let args: BasicPublishArguments; let routing_key = match sub.endpoint.as_str() {
"apn" => &config.pushd.apn.queue,
"fcm" => &config.pushd.fcm.queue,
endpoint => {
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), endpoint.to_string());
if sub.endpoint == "apn" { &config.pushd.vapid.queue
args = BasicPublishArguments::new( }
config.pushd.exchange.as_str(), };
config.pushd.apn.queue.as_str(),
)
.finish();
} else if sub.endpoint == "fcm" {
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.fcm.queue.as_str(),
)
.finish();
} else {
// web push (vapid)
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.vapid.queue.as_str(),
)
.finish();
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), sub.endpoint.clone());
}
let payload = serde_json::to_string(&sendable)?; let payload = serde_json::to_string(&sendable)?;
publish_message(self, payload.into(), args).await; self.publish_message(payload.as_bytes(), &config.pushd.exchange, routing_key)
.await?;
} }
} }
} }
@@ -117,24 +78,3 @@ impl FRAcceptedConsumer {
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for FRAcceptedConsumer {
/// This consumer handles delegating messages into their respective platform queues.
async fn consume(
&mut self,
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process friend request accepted event: {err:?}");
}
}
}
@@ -1,70 +1,44 @@
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use crate::consumers::inbound::internal::*; use crate::utils::Consumer;
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::Connection,
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use anyhow::Result; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use log::debug; use log::debug;
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
#[derive(Clone)]
#[allow(unused)]
pub struct FRReceivedConsumer { pub struct FRReceivedConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
} }
impl Channeled for FRReceivedConsumer { #[async_trait]
fn get_connection(&self) -> Option<&Connection> { impl Consumer for FRReceivedConsumer {
if self.conn.is_none() { async fn create(
None db: Database,
} else { authifier_db: authifier::Database,
Some(self.conn.as_ref().unwrap()) connection: Arc<Connection>,
} channel: Arc<Channel>,
} ) -> Self {
Self {
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
}
impl FRReceivedConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> FRReceivedConsumer {
FRReceivedConsumer {
db, db,
authifier_db, authifier_db,
conn: None, connection,
channel: None, channel,
} }
} }
async fn consume_event( fn channel(&self) -> &Arc<Channel> {
&mut self, &self.channel
_channel: &Channel, }
_deliver: Deliver,
_basic_properties: BasicProperties, /// This consumer handles delegating messages into their respective platform queues.
content: Vec<u8>, async fn consume(&self, delivery: Delivery) -> Result<()> {
) -> Result<()> { let payload: FRReceivedPayload = serde_json::from_slice(&delivery.data)?;
let content = String::from_utf8(content)?;
let payload: FRReceivedPayload = serde_json::from_str(content.as_str())?;
debug!("Received FR received event"); debug!("Received FR received event");
@@ -80,36 +54,23 @@ impl FRReceivedConsumer {
extras: HashMap::new(), extras: HashMap::new(),
}; };
let args: BasicPublishArguments; let routing_key = match sub.endpoint.as_str() {
"apn" => &config.pushd.apn.queue,
"fcm" => &config.pushd.fcm.queue,
endpoint => {
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), endpoint.to_string());
if sub.endpoint == "apn" { &config.pushd.vapid.queue
args = BasicPublishArguments::new( }
config.pushd.exchange.as_str(), };
config.pushd.apn.queue.as_str(),
)
.finish();
} else if sub.endpoint == "fcm" {
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.fcm.queue.as_str(),
)
.finish();
} else {
// web push (vapid)
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.vapid.queue.as_str(),
)
.finish();
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), sub.endpoint.clone());
}
let payload = serde_json::to_string(&sendable)?; let payload = serde_json::to_string(&sendable)?;
publish_message(self, payload.into(), args).await; self.publish_message(payload.as_bytes(), &config.pushd.exchange, routing_key)
.await?;
} }
} }
} }
@@ -117,24 +78,3 @@ impl FRReceivedConsumer {
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for FRReceivedConsumer {
/// This consumer handles delegating messages into their respective platform queues.
async fn consume(
&mut self,
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process friend request received event: {err:?}");
}
}
}
@@ -1,70 +1,44 @@
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use crate::consumers::inbound::internal::*; use crate::utils::Consumer;
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::Connection,
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use anyhow::Result; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use log::debug; use log::debug;
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
#[derive(Clone)]
#[allow(unused)]
pub struct GenericConsumer { pub struct GenericConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
} }
impl Channeled for GenericConsumer { #[async_trait]
fn get_connection(&self) -> Option<&Connection> { impl Consumer for GenericConsumer {
if self.conn.is_none() { async fn create(
None db: Database,
} else { authifier_db: authifier::Database,
Some(self.conn.as_ref().unwrap()) connection: Arc<Connection>,
} channel: Arc<Channel>,
} ) -> Self {
Self {
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
}
impl GenericConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> GenericConsumer {
GenericConsumer {
db, db,
authifier_db, authifier_db,
conn: None, connection,
channel: None, channel,
} }
} }
async fn consume_event( fn channel(&self) -> &Arc<Channel> {
&mut self, &self.channel
_channel: &Channel, }
_deliver: Deliver,
_basic_properties: BasicProperties, /// This consumer handles delegating messages into their respective platform queues.
content: Vec<u8>, async fn consume(&self, delivery: Delivery) -> Result<()> {
) -> Result<()> { let payload: MessageSentPayload = serde_json::from_slice(&delivery.data)?;
let content = String::from_utf8(content)?;
let payload: MessageSentPayload = serde_json::from_str(content.as_str())?;
debug!("Received message event on origin"); debug!("Received message event on origin");
@@ -86,36 +60,23 @@ impl GenericConsumer {
extras: HashMap::new(), extras: HashMap::new(),
}; };
let args: BasicPublishArguments; let routing_key = match sub.endpoint.as_str() {
"apn" => &config.pushd.apn.queue,
"fcm" => &config.pushd.fcm.queue,
endpoint => {
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), endpoint.to_string());
if sub.endpoint == "apn" { &config.pushd.vapid.queue
args = BasicPublishArguments::new( }
config.pushd.exchange.as_str(), };
config.pushd.apn.queue.as_str(),
)
.finish();
} else if sub.endpoint == "fcm" {
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.fcm.queue.as_str(),
)
.finish();
} else {
// web push (vapid)
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.vapid.queue.as_str(),
)
.finish();
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), sub.endpoint.clone());
}
let payload = serde_json::to_string(&sendable)?; let payload = serde_json::to_string(&sendable)?;
publish_message(self, payload.into(), args).await; self.publish_message(payload.as_bytes(), &config.pushd.exchange, routing_key)
.await?;
} }
} }
} }
@@ -123,24 +84,3 @@ impl GenericConsumer {
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for GenericConsumer {
/// This consumer handles delegating messages into their respective platform queues.
async fn consume(
&mut self,
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process generic event: {err:?}");
}
}
}
@@ -1,53 +0,0 @@
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::{Connection, OpenConnectionArguments},
BasicProperties,
};
use log::{debug, warn};
pub(crate) trait Channeled {
#[allow(unused)]
fn get_connection(&self) -> Option<&Connection>;
fn get_channel(&self) -> Option<&Channel>;
fn set_connection(&mut self, conn: Connection);
fn set_channel(&mut self, channel: Channel);
}
pub(crate) async fn make_channel<T: Channeled>(consumer: &mut T) {
let config = revolt_config::config().await;
let args = OpenConnectionArguments::new(
&config.rabbit.host,
config.rabbit.port,
&config.rabbit.username,
&config.rabbit.password,
);
let conn = amqprs::connection::Connection::open(&args).await.unwrap();
let channel = conn.open_channel(None).await.unwrap();
consumer.set_connection(conn);
consumer.set_channel(channel);
}
pub(crate) async fn publish_message<T: Channeled>(
consumer: &mut T,
payload: Vec<u8>,
args: BasicPublishArguments,
) {
let routing_key = &args.routing_key.clone();
let mut channel = consumer.get_channel();
if channel.is_none() {
make_channel(consumer).await;
channel = consumer.get_channel();
}
if let Some(chnl) = channel {
chnl.basic_publish(BasicProperties::default(), payload.clone(), args.clone())
.await
.unwrap();
debug!("Sent message to queue for target {}", routing_key);
} else {
warn!("Failed to unwrap channel (including attempt to make a channel)!")
}
}
@@ -1,17 +1,13 @@
use std::{ use std::{
collections::{HashMap, HashSet}, collections::{HashMap, HashSet},
hash::RandomState, hash::RandomState,
sync::Arc,
}; };
use crate::{consumers::inbound::internal::*, utils}; use crate::utils::{render_notification_content, Consumer};
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::Connection,
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use anyhow::Result; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use revolt_database::{ use revolt_database::{
events::rabbit::*, util::bulk_permissions::BulkDatabasePermissionQuery, Database, Member, events::rabbit::*, util::bulk_permissions::BulkDatabasePermissionQuery, Database, Member,
MessageFlagsValue, MessageFlagsValue,
@@ -19,52 +15,18 @@ use revolt_database::{
use revolt_models::v0::{MessageFlags, PushNotification}; use revolt_models::v0::{MessageFlags, PushNotification};
use revolt_result::ToRevoltError; use revolt_result::ToRevoltError;
#[derive(Clone)]
#[allow(unused)]
pub struct MassMessageConsumer { pub struct MassMessageConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
}
impl Channeled for MassMessageConsumer {
fn get_connection(&self) -> Option<&Connection> {
if self.conn.is_none() {
None
} else {
Some(self.conn.as_ref().unwrap())
}
}
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
} }
impl MassMessageConsumer { impl MassMessageConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> MassMessageConsumer {
MassMessageConsumer {
db,
authifier_db,
conn: None,
channel: None,
}
}
async fn fire_notification_for_users( async fn fire_notification_for_users(
&mut self, &self,
push: &PushNotification, push: &PushNotification,
users: &[String], users: &[String],
) -> Result<()> { ) -> Result<()> {
@@ -84,56 +46,58 @@ impl MassMessageConsumer {
extras: HashMap::new(), extras: HashMap::new(),
}; };
let args: BasicPublishArguments; let routing_key = match sub.endpoint.as_str() {
"apn" => &config.pushd.apn.queue,
"fcm" => &config.pushd.fcm.queue,
endpoint => {
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), endpoint.to_string());
if sub.endpoint == "apn" { &config.pushd.vapid.queue
args = BasicPublishArguments::new( }
config.pushd.exchange.as_str(), };
config.pushd.apn.queue.as_str(),
)
.finish();
} else if sub.endpoint == "fcm" {
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.fcm.queue.as_str(),
)
.finish();
} else {
// web push (vapid)
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.vapid.queue.as_str(),
)
.finish();
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), sub.endpoint.clone());
}
let payload = serde_json::to_string(&sendable)?; let payload = serde_json::to_string(&sendable)?;
publish_message(self, payload.into(), args).await; self.publish_message(payload.as_bytes(), &config.pushd.exchange, routing_key)
.await?;
} }
} }
} }
Ok(()) Ok(())
} }
}
async fn consume_event( #[async_trait]
&mut self, impl Consumer for MassMessageConsumer {
_channel: &Channel, async fn create(
_deliver: Deliver, db: Database,
_basic_properties: BasicProperties, authifier_db: authifier::Database,
content: Vec<u8>, connection: Arc<Connection>,
) -> Result<()> { channel: Arc<Channel>,
) -> Self {
Self {
db,
authifier_db,
connection,
channel,
}
}
fn channel(&self) -> &Arc<Channel> {
&self.channel
}
/// This consumer handles adding mentions for all the users affected by a mass mention ping, and then sends out push notifications.
async fn consume(&self, delivery: Delivery) -> Result<()> {
let mut payload: MassMessageSentPayload = serde_json::from_slice(&delivery.data)?;
let config = revolt_config::config().await; let config = revolt_config::config().await;
let content = String::from_utf8(content)?;
let mut payload: MassMessageSentPayload = serde_json::from_str(content.as_str())?;
for push in payload.notifications.iter_mut() { for push in payload.notifications.iter_mut() {
if let Ok(body) = utils::render_notification_content(push, &self.db) if let Ok(body) = render_notification_content(push, &self.db)
.await .await
.to_internal_error() .to_internal_error()
{ {
@@ -280,24 +244,3 @@ impl MassMessageConsumer {
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for MassMessageConsumer {
/// This consumer handles adding mentions for all the users affected by a mass mention ping, and then sends out push notifications
async fn consume(
&mut self,
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process mass message event: {err:?}");
}
}
}
@@ -1,76 +1,46 @@
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use crate::{consumers::inbound::internal::*, utils}; use crate::utils::{render_notification_content, Consumer};
use amqprs::{
channel::{BasicPublishArguments, Channel},
connection::Connection,
consumer::AsyncConsumer,
BasicProperties, Deliver,
};
use anyhow::Result; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use lapin::{message::Delivery, Channel, Connection};
use log::debug; use log::debug;
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
use revolt_result::ToRevoltError;
#[derive(Clone)]
#[allow(unused)]
pub struct MessageConsumer { pub struct MessageConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database, authifier_db: authifier::Database,
conn: Option<Connection>, connection: Arc<Connection>,
channel: Option<Channel>, channel: Arc<Channel>,
} }
impl Channeled for MessageConsumer { #[async_trait]
fn get_connection(&self) -> Option<&Connection> { impl Consumer for MessageConsumer {
if self.conn.is_none() { async fn create(
None db: Database,
} else { authifier_db: authifier::Database,
Some(self.conn.as_ref().unwrap()) connection: Arc<Connection>,
} channel: Arc<Channel>,
} ) -> Self {
Self {
fn get_channel(&self) -> Option<&Channel> {
if self.channel.is_none() {
None
} else {
Some(self.channel.as_ref().unwrap())
}
}
fn set_connection(&mut self, conn: Connection) {
self.conn = Some(conn);
}
fn set_channel(&mut self, channel: Channel) {
self.channel = Some(channel)
}
}
impl MessageConsumer {
pub fn new(db: Database, authifier_db: authifier::Database) -> MessageConsumer {
MessageConsumer {
db, db,
authifier_db, authifier_db,
conn: None, connection,
channel: None, channel,
} }
} }
async fn consume_event( fn channel(&self) -> &Arc<Channel> {
&mut self, &self.channel
_channel: &Channel, }
_deliver: Deliver,
_basic_properties: BasicProperties,
content: Vec<u8>,
) -> Result<()> {
let content = String::from_utf8(content)?;
let mut payload: MessageSentPayload = serde_json::from_str(content.as_str())?;
if let Ok(body) = utils::render_notification_content(&payload.notification, &self.db) /// This consumer handles delegating messages into their respective platform queues.
.await async fn consume(&self, delivery: Delivery) -> Result<()> {
.to_internal_error() let mut payload: MessageSentPayload = serde_json::from_slice(&delivery.data)?;
{
if let Ok(body) = render_notification_content(&payload.notification, &self.db).await {
payload.notification.raw_body = Some(payload.notification.body); payload.notification.raw_body = Some(payload.notification.body);
payload.notification.body = body; payload.notification.body = body;
} }
@@ -95,36 +65,22 @@ impl MessageConsumer {
extras: HashMap::new(), extras: HashMap::new(),
}; };
let args: BasicPublishArguments; let routing_key = match sub.endpoint.as_str() {
"apn" => &config.pushd.apn.queue,
"fcm" => &config.pushd.fcm.queue,
endpoint => {
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), endpoint.to_string());
if sub.endpoint == "apn" { &config.pushd.vapid.queue
args = BasicPublishArguments::new( }
config.pushd.exchange.as_str(), };
config.pushd.apn.queue.as_str(),
)
.finish();
} else if sub.endpoint == "fcm" {
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.fcm.queue.as_str(),
)
.finish();
} else {
// web push (vapid)
args = BasicPublishArguments::new(
config.pushd.exchange.as_str(),
config.pushd.vapid.queue.as_str(),
)
.finish();
sendable.extras.insert("p256dh".to_string(), sub.p256dh);
sendable
.extras
.insert("endpoint".to_string(), sub.endpoint.clone());
}
let payload = serde_json::to_string(&sendable)?; let payload = serde_json::to_string(&sendable)?;
self.publish_message(payload.as_bytes(), &config.pushd.exchange, routing_key)
publish_message(self, payload.into(), args).await; .await?;
} }
} }
} }
@@ -132,24 +88,3 @@ impl MessageConsumer {
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for MessageConsumer {
/// This consumer handles delegating messages into their respective platform queues.
async fn consume(
&mut self,
channel: &Channel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process message event: {err:?}");
}
}
}
@@ -3,6 +3,5 @@ pub mod dm_call;
pub mod fr_accepted; pub mod fr_accepted;
pub mod fr_received; pub mod fr_received;
pub mod generic; pub mod generic;
mod internal;
pub mod mass_mention; pub mod mass_mention;
pub mod message; pub mod message;
@@ -1,12 +1,13 @@
use std::{borrow::Cow, collections::BTreeMap, io::Cursor}; use std::{borrow::Cow, collections::BTreeMap, io::Cursor, sync::Arc};
use amqprs::{channel::Channel as AmqpChannel, consumer::AsyncConsumer, BasicProperties, Deliver}; use crate::utils::Consumer;
use anyhow::{anyhow, Result}; use anyhow::Result;
use async_trait::async_trait; use async_trait::async_trait;
use base64::{ use base64::{
engine::{self}, engine::{self},
Engine as _, Engine as _,
}; };
use lapin::{message::Delivery, Channel as AMQPChannel, Connection};
use revolt_a2::{ use revolt_a2::{
request::{ request::{
notification::{DefaultAlert, NotificationOptions}, notification::{DefaultAlert, NotificationOptions},
@@ -42,7 +43,7 @@ impl<'a> PayloadLike for MessagePayload<'a> {
fn get_device_token(&self) -> &'a str { fn get_device_token(&self) -> &'a str {
self.device_token self.device_token
} }
fn get_options(&self) -> &NotificationOptions { fn get_options(&self) -> &NotificationOptions<'a> {
&self.options &self.options
} }
} }
@@ -68,16 +69,20 @@ impl<'a> PayloadLike for CallStartStopPayload<'a> {
fn get_device_token(&self) -> &'a str { fn get_device_token(&self) -> &'a str {
self.device_token self.device_token
} }
fn get_options(&self) -> &NotificationOptions { fn get_options(&self) -> &NotificationOptions<'a> {
&self.options &self.options
} }
} }
// region: consumer // region: consumer
#[derive(Clone)]
#[allow(unused)]
pub struct ApnsOutboundConsumer { pub struct ApnsOutboundConsumer {
#[allow(dead_code)]
db: Database, db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<AMQPChannel>,
client: Client, client: Client,
} }
@@ -117,15 +122,21 @@ impl ApnsOutboundConsumer {
} }
} }
impl ApnsOutboundConsumer { #[async_trait]
pub async fn new(db: Database) -> Result<ApnsOutboundConsumer, &'static str> { impl Consumer for ApnsOutboundConsumer {
async fn create(
db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<AMQPChannel>,
) -> Self {
let config = revolt_config::config().await; let config = revolt_config::config().await;
if config.pushd.apn.pkcs8.is_empty() if config.pushd.apn.pkcs8.is_empty()
|| config.pushd.apn.key_id.is_empty() || config.pushd.apn.key_id.is_empty()
|| config.pushd.apn.team_id.is_empty() || config.pushd.apn.team_id.is_empty()
{ {
return Err("Missing APN keys."); panic!("Missing APN keys.");
} }
let endpoint = if config.pushd.apn.sandbox { let endpoint = if config.pushd.apn.sandbox {
@@ -148,18 +159,21 @@ impl ApnsOutboundConsumer {
) )
.expect("could not create APN client"); .expect("could not create APN client");
Ok(ApnsOutboundConsumer { db, client }) Self {
db,
authifier_db,
connection,
channel,
client,
}
} }
async fn consume_event( fn channel(&self) -> &Arc<AMQPChannel> {
&mut self, &self.channel
_channel: &AmqpChannel, }
_deliver: Deliver,
_basic_properties: BasicProperties, async fn consume(&self, delivery: Delivery) -> Result<()> {
content: Vec<u8>, let payload: PayloadToService = serde_json::from_slice(&delivery.data)?;
) -> Result<()> {
let content = String::from_utf8(content)?;
let payload: PayloadToService = serde_json::from_str(content.as_str())?;
let payload_options = NotificationOptions { let payload_options = NotificationOptions {
apns_id: None, apns_id: None,
@@ -170,20 +184,15 @@ impl ApnsOutboundConsumer {
apns_collapse_id: None, apns_collapse_id: None,
}; };
let resp: Result<Response, Error>; let resp = match payload.notification {
match payload.notification {
PayloadKind::FRReceived(alert) => { PayloadKind::FRReceived(alert) => {
let loc_args = vec![Cow::from( let loc_args = vec![Cow::from(
alert alert.from_user.display_name.clone().unwrap_or_else(|| {
.from_user format!(
.display_name
.or(Some(format!(
"{}#{}", "{}#{}",
alert.from_user.username, alert.from_user.discriminator alert.from_user.username, alert.from_user.discriminator
))) )
.clone() }),
.ok_or_else(|| anyhow!("missing name"))?,
)]; )];
let apn_payload = Payload { let apn_payload = Payload {
@@ -216,20 +225,17 @@ impl ApnsOutboundConsumer {
"Sending friend request received for user: {:}", "Sending friend request received for user: {:}",
&payload.user_id &payload.user_id
); );
resp = self.client.send(apn_payload).await; self.client.send(apn_payload).await
} }
PayloadKind::FRAccepted(alert) => { PayloadKind::FRAccepted(alert) => {
let loc_args = vec![Cow::from( let loc_args = vec![Cow::from(
alert alert.accepted_user.display_name.clone().unwrap_or_else(|| {
.accepted_user format!(
.display_name
.or(Some(format!(
"{}#{}", "{}#{}",
alert.accepted_user.username, alert.accepted_user.discriminator alert.accepted_user.username, alert.accepted_user.discriminator
))) )
.clone() }),
.ok_or_else(|| anyhow!("missing name"))?,
)]; )];
let apn_payload = Payload { let apn_payload = Payload {
@@ -262,7 +268,7 @@ impl ApnsOutboundConsumer {
"Sending friend request accept for user: {:}", "Sending friend request accept for user: {:}",
&payload.user_id &payload.user_id
); );
resp = self.client.send(apn_payload).await; self.client.send(apn_payload).await
} }
PayloadKind::Generic(alert) => { PayloadKind::Generic(alert) => {
let apn_payload = Payload { let apn_payload = Payload {
@@ -295,7 +301,7 @@ impl ApnsOutboundConsumer {
"Sending generic notification for user: {:}", "Sending generic notification for user: {:}",
&payload.user_id &payload.user_id
); );
resp = self.client.send(apn_payload).await; self.client.send(apn_payload).await
} }
PayloadKind::MessageNotification(alert) => { PayloadKind::MessageNotification(alert) => {
@@ -334,7 +340,7 @@ impl ApnsOutboundConsumer {
"Sending message notification for user: {:}", "Sending message notification for user: {:}",
&payload.user_id &payload.user_id
); );
resp = self.client.send(apn_payload).await; self.client.send(apn_payload).await
} }
PayloadKind::BadgeUpdate(badge) => { PayloadKind::BadgeUpdate(badge) => {
@@ -349,7 +355,7 @@ impl ApnsOutboundConsumer {
}; };
debug!("Sending badge update for user: {:}", &payload.user_id); debug!("Sending badge update for user: {:}", &payload.user_id);
resp = self.client.send(apn_payload).await; self.client.send(apn_payload).await
} }
PayloadKind::DmCallStartEnd(alert) => { PayloadKind::DmCallStartEnd(alert) => {
@@ -378,58 +384,37 @@ impl ApnsOutboundConsumer {
"Sending call start/stop notification for user: {:}", "Sending call start/stop notification for user: {:}",
&payload.user_id &payload.user_id
); );
resp = self.client.send(apn_payload).await; self.client.send(apn_payload).await
} }
} };
if let Err(err) = resp { match resp {
match err { Err(Error::ResponseError(Response {
Error::ResponseError(Response { error:
error: Some(ErrorBody {
Some(ErrorBody { reason: ErrorReason::BadDeviceToken | ErrorReason::Unregistered,
reason: ErrorReason::BadDeviceToken | ErrorReason::Unregistered, ..
.. }),
}), ..
.. })) => {
}) => { info!(
info!( "Removing APNS subscription id {:} (user: {:}) due to invalid token",
"Removing APNS subscription id {:} (user: {:}) due to invalid token", &payload.session_id, &payload.user_id
&payload.session_id, &payload.user_id );
);
if let Err(err) = self if let Err(err) = self
.db .db
.remove_push_subscription_by_session_id(&payload.session_id) .remove_push_subscription_by_session_id(&payload.session_id)
.await .await
{ {
revolt_config::capture_error(&err);
}
}
err => {
revolt_config::capture_error(&err); revolt_config::capture_error(&err);
} }
} }
} resp => {
resp?;
}
};
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for ApnsOutboundConsumer {
async fn consume(
&mut self,
channel: &AmqpChannel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process APN event: {err:?}");
}
}
}
+183 -138
View File
@@ -1,51 +1,145 @@
use std::{collections::HashMap, time::Duration}; use std::{collections::HashMap, sync::Arc, time::Duration};
use amqprs::{channel::Channel as AmqpChannel, consumer::AsyncConsumer, BasicProperties, Deliver}; use crate::utils::Consumer;
use anyhow::{bail, Result};
use anyhow::{anyhow, bail, Result};
use async_trait::async_trait; use async_trait::async_trait;
use fcm_v1::{ use fcm_v1::{
android::{AndroidConfig, AndroidMessagePriority},
auth::{Authenticator, ServiceAccountKey}, auth::{Authenticator, ServiceAccountKey},
message::{Message, Notification}, message::Message,
Client, Error as FcmError, Client, Error as FcmError,
}; };
use lapin::{message::Delivery, Channel as AMQPChannel, Connection};
use revolt_config::config; use revolt_config::config;
use revolt_database::{events::rabbit::*, Database}; use revolt_database::{events::rabbit::*, Database};
use revolt_models::v0::{Channel, PushNotification};
use serde_json::Value; use serde_json::Value;
pub struct FcmOutboundConsumer { /// Custom notification data
db: Database, #[derive(Debug, Clone, PartialEq)]
client: Client, pub enum NotificationData {
FRReceived {
id: String,
username: String,
},
FRAccepted {
id: String,
username: String,
},
Generic {
title: String,
body: String,
image: Option<String>,
},
Message {
message: String,
body: String,
image: String,
channel: String,
author: String,
author_name: String,
},
DmCallStartEnd {
initiator_id: String,
channel_id: String,
started_at: String,
ended: bool,
duration: usize,
},
} }
impl FcmOutboundConsumer { impl NotificationData {
fn format_title(&self, notification: &PushNotification) -> String { pub fn get_type(&self) -> &str {
// ideally this changes depending on context match self {
// in a server, it would look like "Sendername, #channelname in servername" NotificationData::FRReceived { .. } => "push.fr.receive",
// in a group, it would look like "Sendername in groupname" NotificationData::FRAccepted { .. } => "push.fr.accept",
// in a dm it should just be "Sendername". NotificationData::Generic { .. } => "push.generic",
// not sure how feasible all those are given the PushNotification object as it currently stands. NotificationData::Message { .. } => "push.message",
NotificationData::DmCallStartEnd { .. } => "push.dm.call",
#[allow(deprecated)]
match &notification.channel {
Channel::DirectMessage { .. } => notification.author.clone(),
Channel::Group { name, .. } => format!("{}, #{}", notification.author, name),
Channel::TextChannel { name, .. } => {
format!("{} in #{}", notification.author, name)
}
_ => "Unknown".to_string(),
} }
} }
pub fn into_payload(self) -> HashMap<String, Value> {
let mut data = HashMap::new();
data.insert(
"type".to_string(),
Value::String(self.get_type().to_string()),
);
match self {
NotificationData::FRReceived { id, username } => {
data.insert("id".to_string(), Value::String(id));
data.insert("username".to_string(), Value::String(username));
}
NotificationData::FRAccepted { id, username } => {
data.insert("id".to_string(), Value::String(id));
data.insert("username".to_string(), Value::String(username));
}
NotificationData::Generic { title, body, image } => {
data.insert("title".to_string(), Value::String(title));
data.insert("body".to_string(), Value::String(body));
if let Some(image) = image {
data.insert("image".to_string(), Value::String(image));
}
}
NotificationData::Message {
message,
body,
image,
channel,
author,
author_name,
} => {
data.insert("message".to_string(), Value::String(message));
data.insert("body".to_string(), Value::String(body));
data.insert("image".to_string(), Value::String(image));
data.insert("channel".to_string(), Value::String(channel));
data.insert("author".to_string(), Value::String(author));
data.insert("author_name".to_string(), Value::String(author_name));
}
NotificationData::DmCallStartEnd {
initiator_id,
channel_id,
started_at,
ended,
duration,
} => {
data.insert("initiator_id".to_string(), Value::String(initiator_id));
data.insert("channel_id".to_string(), Value::String(channel_id));
data.insert("started_at".to_string(), Value::String(started_at));
data.insert("ended".to_string(), Value::Bool(ended));
data.insert("duration".to_string(), Value::Number(duration.into()));
}
}
data
}
} }
impl FcmOutboundConsumer { #[derive(Clone)]
pub async fn new(db: Database) -> Result<FcmOutboundConsumer, &'static str> { #[allow(unused)]
pub struct FcmOutboundConsumer {
db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<AMQPChannel>,
client: Client,
}
#[async_trait]
impl Consumer for FcmOutboundConsumer {
async fn create(
db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<AMQPChannel>,
) -> Self {
let config = revolt_config::config().await; let config = revolt_config::config().await;
Ok(FcmOutboundConsumer { Self {
db, db,
authifier_db,
connection,
channel,
client: Client::new( client: Client::new(
Authenticator::service_account::<&str>(ServiceAccountKey { Authenticator::service_account::<&str>(ServiceAccountKey {
key_type: Some(config.pushd.fcm.key_type), key_type: Some(config.pushd.fcm.key_type),
@@ -65,45 +159,36 @@ impl FcmOutboundConsumer {
false, false,
Duration::from_secs(5), Duration::from_secs(5),
), ),
}) }
} }
async fn consume_event( fn channel(&self) -> &Arc<AMQPChannel> {
&mut self, &self.channel
_channel: &AmqpChannel, }
_deliver: Deliver,
_basic_properties: BasicProperties, async fn consume(&self, delivery: Delivery) -> Result<()> {
content: Vec<u8>, let payload: PayloadToService = serde_json::from_slice(&delivery.data)?;
) -> Result<()> {
let content = String::from_utf8(content)?;
let payload: PayloadToService = serde_json::from_str(content.as_str())?;
#[allow(clippy::needless_late_init)] #[allow(clippy::needless_late_init)]
let resp: Result<Message, FcmError>; let resp: Result<Message, FcmError>;
match payload.notification { match payload.notification {
PayloadKind::FRReceived(alert) => { PayloadKind::FRReceived(alert) => {
let name = alert let name = alert.from_user.display_name.clone().unwrap_or_else(|| {
.from_user format!(
.display_name
.or(Some(format!(
"{}#{}", "{}#{}",
alert.from_user.username, alert.from_user.discriminator alert.from_user.username, alert.from_user.discriminator
))) )
.clone() });
.ok_or_else(|| anyhow!("missing name"))?;
let mut data = HashMap::new(); let data = NotificationData::FRReceived {
data.insert( id: alert.from_user.id,
"type".to_string(), username: name,
Value::String("push.fr.receive".to_string()), };
);
data.insert("id".to_string(), Value::String(alert.from_user.id));
data.insert("username".to_string(), Value::String(name));
let msg = Message { let msg = Message {
token: Some(payload.token), token: Some(payload.token),
data: Some(data), data: Some(data.into_payload()),
..Default::default() ..Default::default()
}; };
@@ -111,40 +196,36 @@ impl FcmOutboundConsumer {
} }
PayloadKind::FRAccepted(alert) => { PayloadKind::FRAccepted(alert) => {
let name = alert let name = alert.accepted_user.display_name.clone().unwrap_or_else(|| {
.accepted_user format!(
.display_name
.or(Some(format!(
"{}#{}", "{}#{}",
alert.accepted_user.username, alert.accepted_user.discriminator alert.accepted_user.username, alert.accepted_user.discriminator
))) )
.clone() });
.ok_or_else(|| anyhow!("missing name"))?;
let mut data: HashMap<String, Value> = HashMap::new(); let data = NotificationData::FRAccepted {
data.insert( id: alert.accepted_user.id,
"type".to_string(), username: name,
Value::String("push.fr.accept".to_string()), };
);
data.insert("id".to_string(), Value::String(alert.accepted_user.id));
data.insert("username".to_string(), Value::String(name));
let msg = Message { let msg = Message {
token: Some(payload.token), token: Some(payload.token),
data: Some(data), data: Some(data.into_payload()),
..Default::default() ..Default::default()
}; };
resp = self.client.send(&msg).await; resp = self.client.send(&msg).await;
} }
PayloadKind::Generic(alert) => { PayloadKind::Generic(alert) => {
let data = NotificationData::Generic {
title: alert.title,
body: alert.body,
image: alert.icon,
};
let msg = Message { let msg = Message {
token: Some(payload.token), token: Some(payload.token),
notification: Some(Notification { data: Some(data.into_payload()),
title: Some(alert.title),
body: Some(alert.body),
image: alert.icon,
}),
..Default::default() ..Default::default()
}; };
@@ -152,19 +233,18 @@ impl FcmOutboundConsumer {
} }
PayloadKind::MessageNotification(alert) => { PayloadKind::MessageNotification(alert) => {
let title = self.format_title(&alert); let data = NotificationData::Message {
message: alert.message.id,
body: alert.body,
image: alert.icon,
channel: alert.message.channel,
author: alert.message.author,
author_name: alert.author,
};
let msg = Message { let msg = Message {
token: Some(payload.token), token: Some(payload.token),
notification: Some(Notification { data: Some(data.into_payload()),
title: Some(title),
body: Some(alert.body),
image: Some(alert.icon),
}),
android: Some(AndroidConfig {
collapse_key: Some(alert.tag),
..Default::default()
}),
..Default::default() ..Default::default()
}; };
@@ -172,30 +252,17 @@ impl FcmOutboundConsumer {
} }
PayloadKind::DmCallStartEnd(alert) => { PayloadKind::DmCallStartEnd(alert) => {
let mut data: HashMap<String, Value> = HashMap::new(); let data = NotificationData::DmCallStartEnd {
data.insert( initiator_id: alert.initiator_id,
"initiator_id".to_string(), channel_id: alert.channel_id,
Value::String(alert.initiator_id), started_at: alert.started_at.unwrap_or_else(|| "".to_string()),
); ended: alert.ended,
data.insert("channel_id".to_string(), Value::String(alert.channel_id)); duration: config().await.api.livekit.call_ring_duration,
data.insert( };
"started_at".to_string(),
Value::String(alert.started_at.unwrap_or_else(|| "".to_string())),
);
data.insert("ended".to_string(), Value::Bool(alert.ended));
let msg = Message { let msg = Message {
token: Some(payload.token), token: Some(payload.token),
notification: None, data: Some(data.into_payload()),
data: Some(data),
android: Some(AndroidConfig {
priority: Some(AndroidMessagePriority::High),
ttl: Some(format!(
"{}s",
config().await.api.livekit.call_ring_duration
)),
..Default::default()
}),
..Default::default() ..Default::default()
}; };
@@ -207,43 +274,21 @@ impl FcmOutboundConsumer {
} }
} }
if let Err(err) = resp { match resp {
match err { Err(FcmError::Auth) => {
FcmError::Auth => { if let Err(err) = self
if let Err(err) = self .db
.db .remove_push_subscription_by_session_id(&payload.session_id)
.remove_push_subscription_by_session_id(&payload.session_id) .await
.await {
{
revolt_config::capture_error(&err);
}
}
err => {
revolt_config::capture_error(&err); revolt_config::capture_error(&err);
} }
} }
} res => {
res?;
}
};
Ok(()) Ok(())
} }
} }
#[allow(unused_variables)]
#[async_trait]
impl AsyncConsumer for FcmOutboundConsumer {
async fn consume(
&mut self,
channel: &AmqpChannel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process FCM event: {err:?}");
}
}
}
@@ -1,6 +1,6 @@
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use amqprs::{channel::Channel as AmqpChannel, consumer::AsyncConsumer, BasicProperties, Deliver}; use crate::utils::Consumer;
use anyhow::{anyhow, bail, Result}; use anyhow::{anyhow, bail, Result};
use async_trait::async_trait; use async_trait::async_trait;
@@ -8,46 +8,60 @@ use base64::{
engine::{self}, engine::{self},
Engine as _, Engine as _,
}; };
use lapin::{message::Delivery, Channel as AMQPChannel, Connection};
use revolt_database::{events::rabbit::*, util::format_display_name, Database}; use revolt_database::{events::rabbit::*, util::format_display_name, Database};
use web_push::{ use web_push::{
ContentEncoding, IsahcWebPushClient, SubscriptionInfo, SubscriptionKeys, VapidSignatureBuilder, ContentEncoding, IsahcWebPushClient, SubscriptionInfo, SubscriptionKeys, VapidSignatureBuilder,
WebPushClient, WebPushError, WebPushMessageBuilder, WebPushClient, WebPushError, WebPushMessageBuilder,
}; };
#[derive(Clone)]
#[allow(unused)]
pub struct VapidOutboundConsumer { pub struct VapidOutboundConsumer {
db: Database, db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<AMQPChannel>,
client: IsahcWebPushClient, client: IsahcWebPushClient,
pkey: Vec<u8>, pkey: Arc<Vec<u8>>,
} }
impl VapidOutboundConsumer { #[async_trait]
pub async fn new(db: Database) -> Result<VapidOutboundConsumer> { impl Consumer for VapidOutboundConsumer {
async fn create(
db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<AMQPChannel>,
) -> Self {
let config = revolt_config::config().await; let config = revolt_config::config().await;
if config.pushd.vapid.private_key.is_empty() | config.pushd.vapid.public_key.is_empty() { if config.pushd.vapid.private_key.is_empty() || config.pushd.vapid.public_key.is_empty() {
bail!("no Vapid keys present"); panic!("no Vapid keys present");
} }
let web_push_private_key = engine::general_purpose::URL_SAFE_NO_PAD let web_push_private_key = Arc::new(
.decode(config.pushd.vapid.private_key) engine::general_purpose::URL_SAFE_NO_PAD
.expect("valid `VAPID_PRIVATE_KEY`"); .decode(config.pushd.vapid.private_key)
.expect("valid `VAPID_PRIVATE_KEY`"),
);
Ok(VapidOutboundConsumer { Self {
db, db,
authifier_db,
connection,
channel,
client: IsahcWebPushClient::new().unwrap(), client: IsahcWebPushClient::new().unwrap(),
pkey: web_push_private_key, pkey: web_push_private_key,
}) }
} }
async fn consume_event( fn channel(&self) -> &Arc<AMQPChannel> {
&mut self, &self.channel
_channel: &AmqpChannel, }
_deliver: Deliver,
_basic_properties: BasicProperties, async fn consume(&self, delivery: Delivery) -> Result<()> {
content: Vec<u8>, let payload: PayloadToService = serde_json::from_slice(&delivery.data)?;
) -> Result<()> {
let content = String::from_utf8(content)?;
let payload: PayloadToService = serde_json::from_str(content.as_str())?;
let subscription = SubscriptionInfo { let subscription = SubscriptionInfo {
endpoint: payload endpoint: payload
@@ -65,10 +79,7 @@ impl VapidOutboundConsumer {
}, },
}; };
#[allow(clippy::needless_late_init)] let payload_body = match payload.notification {
let payload_body: String;
match payload.notification {
PayloadKind::FRReceived(alert) => { PayloadKind::FRReceived(alert) => {
let name = alert let name = alert
.from_user .from_user
@@ -83,7 +94,7 @@ impl VapidOutboundConsumer {
let mut body = HashMap::new(); let mut body = HashMap::new();
body.insert("body", format!("{} sent you a friend request", name)); body.insert("body", format!("{} sent you a friend request", name));
payload_body = serde_json::to_string(&body)?; serde_json::to_string(&body)?
} }
PayloadKind::FRAccepted(alert) => { PayloadKind::FRAccepted(alert) => {
let name = alert let name = alert
@@ -99,14 +110,10 @@ impl VapidOutboundConsumer {
let mut body = HashMap::new(); let mut body = HashMap::new();
body.insert("body", format!("{} accepted your friend request", name)); body.insert("body", format!("{} accepted your friend request", name));
payload_body = serde_json::to_string(&body)?; serde_json::to_string(&body)?
}
PayloadKind::Generic(alert) => {
payload_body = serde_json::to_string(&alert)?;
}
PayloadKind::MessageNotification(alert) => {
payload_body = serde_json::to_string(&alert)?;
} }
PayloadKind::Generic(alert) => serde_json::to_string(&alert)?,
PayloadKind::MessageNotification(alert) => serde_json::to_string(&alert)?,
PayloadKind::DmCallStartEnd(alert) => { PayloadKind::DmCallStartEnd(alert) => {
let initiator_name = if let Some(server_id) = let initiator_name = if let Some(server_id) =
self.db.fetch_channel(&alert.channel_id).await?.server() self.db.fetch_channel(&alert.channel_id).await?.server()
@@ -132,59 +139,41 @@ impl VapidOutboundConsumer {
_ => bail!("Invalid DmCallStart/End channel type"), _ => bail!("Invalid DmCallStart/End channel type"),
} }
payload_body = serde_json::to_string(&body)?; serde_json::to_string(&body)?
} }
PayloadKind::BadgeUpdate(_) => { PayloadKind::BadgeUpdate(_) => {
bail!("Vapid cannot handle badge updates and they should not be sent here."); bail!("Vapid cannot handle badge updates and they should not be sent here.");
} }
} };
match VapidSignatureBuilder::from_pem(std::io::Cursor::new(&self.pkey), &subscription) { let signature = VapidSignatureBuilder::from_pem(
Ok(sig_builder) => match sig_builder.build() { std::io::Cursor::new(self.pkey.as_ref()),
Ok(signature) => { &subscription,
let mut builder = WebPushMessageBuilder::new(&subscription); )?
builder.set_vapid_signature(signature); .build()?;
builder.set_payload(ContentEncoding::AesGcm, payload_body.as_bytes()); let mut builder = WebPushMessageBuilder::new(&subscription);
builder.set_vapid_signature(signature);
match builder.build() { builder.set_payload(ContentEncoding::AesGcm, payload_body.as_bytes());
Ok(msg) => {
if let Err(err) = self.client.send(msg).await {
if err == WebPushError::Unauthorized {
self.db
.remove_push_subscription_by_session_id(&payload.session_id)
.await?;
}
}
Ok(()) let msg = builder.build()?;
}
Err(err) => Err(err.into()), match self.client.send(msg).await {
} Err(WebPushError::Unauthorized) => {
if let Err(err) = self
.db
.remove_push_subscription_by_session_id(&payload.session_id)
.await
{
revolt_config::capture_error(&err);
} }
Err(err) => Err(err.into()), }
}, res => {
Err(err) => Err(err.into()), res?;
} }
} };
}
#[allow(unused_variables)] Ok(())
#[async_trait]
impl AsyncConsumer for VapidOutboundConsumer {
async fn consume(
&mut self,
channel: &AmqpChannel,
deliver: Deliver,
basic_properties: BasicProperties,
content: Vec<u8>,
) {
if let Err(err) = self
.consume_event(channel, deliver, basic_properties, content)
.await
{
revolt_config::capture_anyhow(&err);
eprintln!("Failed to process Vapid event: {err:?}");
}
} }
} }
+150 -101
View File
@@ -1,17 +1,16 @@
#[macro_use] #[macro_use]
extern crate log; extern crate log;
use amqprs::{ use std::sync::Arc;
channel::{
BasicConsumeArguments, Channel, ExchangeDeclareArguments, QueueBindArguments, use lapin::{
QueueDeclareArguments, options::{BasicConsumeOptions, ExchangeDeclareOptions, QueueBindOptions, QueueDeclareOptions},
}, types::{AMQPValue, FieldTable},
connection::{Connection, OpenConnectionArguments}, Channel, Connection, ConnectionProperties,
consumer::AsyncConsumer,
FieldTable,
}; };
use revolt_config::{config, Settings}; use revolt_config::{config, Settings};
use tokio::sync::Notify; use revolt_database::Database;
use tokio::signal::ctrl_c;
mod consumers; mod consumers;
mod utils; mod utils;
@@ -24,6 +23,8 @@ use consumers::{
outbound::{apn::ApnsOutboundConsumer, fcm::FcmOutboundConsumer, vapid::VapidOutboundConsumer}, outbound::{apn::ApnsOutboundConsumer, fcm::FcmOutboundConsumer, vapid::VapidOutboundConsumer},
}; };
use crate::utils::{Consumer, Delegate};
#[tokio::main(flavor = "multi_thread", worker_threads = 2)] #[tokio::main(flavor = "multi_thread", worker_threads = 2)]
async fn main() { async fn main() {
// Configure logging and environment // Configure logging and environment
@@ -43,7 +44,24 @@ async fn main() {
panic!("Mongo is not in use, can't connect via authifier!") panic!("Mongo is not in use, can't connect via authifier!")
} }
let mut connections: Vec<(Channel, Connection)> = Vec::new(); let config = config().await;
let connection = Arc::new(
Connection::connect(
&format!(
"amqp://{}:{}@{}:{}",
&config.rabbit.username,
&config.rabbit.password,
&config.rabbit.host,
&config.rabbit.port,
),
ConnectionProperties::default(),
)
.await
.expect("Failed to connect to RabbitMQ"),
);
let mut channels = Vec::new();
// An explainer of how this works: // An explainer of how this works:
// The inbound connections are on separate routing keys, such that they only receive the proper payload // The inbound connections are on separate routing keys, such that they only receive the proper payload
@@ -54,171 +72,178 @@ async fn main() {
// This'll require some interesting shimming if we need to add more events once this is in prod (different payloads between prod and test), // This'll require some interesting shimming if we need to add more events once this is in prod (different payloads between prod and test),
// but that sounds like a problem for future us. // but that sounds like a problem for future us.
let config = config().await; channels.push(
make_queue_and_consume::<GenericConsumer>(
// inbound: generic &db,
connections.push( &authifier,
make_queue_and_consume( &connection,
&config, &config,
&config.pushd.generic_queue, &config.pushd.generic_queue,
config.pushd.get_generic_routing_key().as_str(), &config.pushd.get_generic_routing_key(),
None, None,
GenericConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
// inbound: messages channels.push(
connections.push( make_queue_and_consume::<MessageConsumer>(
make_queue_and_consume( &db,
&authifier,
&connection,
&config, &config,
&config.pushd.message_queue, &config.pushd.message_queue,
config.pushd.get_message_routing_key().as_str(), &config.pushd.get_message_routing_key(),
None, None,
MessageConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
// inbound: FR received channels.push(
connections.push( make_queue_and_consume::<FRReceivedConsumer>(
make_queue_and_consume( &db,
&authifier,
&connection,
&config, &config,
&config.pushd.fr_received_queue, &config.pushd.fr_received_queue,
config.pushd.get_fr_received_routing_key().as_str(), &config.pushd.get_fr_received_routing_key(),
None, None,
FRReceivedConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
// inbound: FR accepted channels.push(
connections.push( make_queue_and_consume::<FRAcceptedConsumer>(
make_queue_and_consume( &db,
&authifier,
&connection,
&config, &config,
&config.pushd.fr_accepted_queue, &config.pushd.fr_accepted_queue,
config.pushd.get_fr_accepted_routing_key().as_str(), &config.pushd.get_fr_accepted_routing_key(),
None, None,
FRAcceptedConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
// inbound: Mass Mentions channels.push(
connections.push( make_queue_and_consume::<MassMessageConsumer>(
make_queue_and_consume( &db,
&authifier,
&connection,
&config, &config,
&config.pushd.mass_mention_queue, &config.pushd.mass_mention_queue,
config.pushd.get_mass_mention_routing_key().as_str(), &config.pushd.get_mass_mention_routing_key(),
None, None,
MassMessageConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
// inbound: Dm Calls channels.push(
connections.push( make_queue_and_consume::<DmCallConsumer>(
make_queue_and_consume( &db,
&authifier,
&connection,
&config, &config,
&config.pushd.dm_call_queue, &config.pushd.dm_call_queue,
config.pushd.get_dm_call_routing_key().as_str(), &config.pushd.get_dm_call_routing_key(),
None, None,
DmCallConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
if !config.pushd.apn.pkcs8.is_empty() { if !config.pushd.apn.pkcs8.is_empty() {
connections.push( channels.push(
make_queue_and_consume( make_queue_and_consume::<ApnsOutboundConsumer>(
&db,
&authifier,
&connection,
&config, &config,
&config.pushd.apn.queue, &config.pushd.apn.queue,
&config.pushd.apn.queue, &config.pushd.apn.queue,
None, None,
ApnsOutboundConsumer::new(db.clone()).await.unwrap(),
) )
.await, .await,
); );
let mut table = FieldTable::new(); let mut table = FieldTable::default();
table.insert("x-message-deduplication".try_into().unwrap(), "true".into()); table.insert("x-message-deduplication".into(), AMQPValue::Boolean(true));
connections.push( channels.push(
make_queue_and_consume( make_queue_and_consume::<AckConsumer>(
&db,
&authifier,
&connection,
&config, &config,
&config.pushd.ack_queue, &config.pushd.ack_queue,
&config.pushd.ack_queue, &config.pushd.ack_queue,
Some(table), Some(table),
AckConsumer::new(db.clone(), authifier.clone()),
) )
.await, .await,
); );
} }
if !config.pushd.fcm.auth_uri.is_empty() { if !config.pushd.fcm.auth_uri.is_empty() {
connections.push( channels.push(
make_queue_and_consume( make_queue_and_consume::<FcmOutboundConsumer>(
&db,
&authifier,
&connection,
&config, &config,
&config.pushd.fcm.queue, &config.pushd.fcm.queue,
&config.pushd.fcm.queue, &config.pushd.fcm.queue,
None, None,
FcmOutboundConsumer::new(db.clone()).await.unwrap(),
) )
.await, .await,
) );
} }
if !config.pushd.vapid.public_key.is_empty() { if !config.pushd.vapid.public_key.is_empty() {
connections.push( channels.push(
make_queue_and_consume( make_queue_and_consume::<VapidOutboundConsumer>(
&db,
&authifier,
&connection,
&config, &config,
&config.pushd.vapid.queue, &config.pushd.vapid.queue,
&config.pushd.vapid.queue, &config.pushd.vapid.queue,
None, None,
VapidOutboundConsumer::new(db.clone()).await.unwrap(),
) )
.await, .await,
) );
} }
let guard = Notify::new(); ctrl_c().await.unwrap();
guard.notified().await;
for (channel, conn) in connections { for channel in channels {
channel.close().await.expect("Unable to close channel"); let _ = channel.close(0, "close".into()).await;
conn.close().await.expect("Unable to close connection");
} }
} }
async fn make_queue_and_consume<F>( async fn make_queue_and_consume<F>(
db: &Database,
authifier_db: &authifier::Database,
connection: &Arc<Connection>,
config: &Settings, config: &Settings,
queue_name: &str, queue_name: &str,
routing_key: &str, routing_key: &str,
queue_args: Option<FieldTable>, queue_args: Option<FieldTable>,
consumer: F, ) -> Arc<Channel>
) -> (Channel, Connection)
where where
F: AsyncConsumer + Send + 'static, F: Consumer,
{ {
let connection = Connection::open(&OpenConnectionArguments::new( let channel = Arc::new(connection.create_channel().await.unwrap());
&config.rabbit.host,
config.rabbit.port,
&config.rabbit.username,
&config.rabbit.password,
))
.await
.unwrap();
let channel = connection.open_channel(None).await.unwrap();
channel channel
.exchange_declare( .exchange_declare(
ExchangeDeclareArguments::new(&config.pushd.exchange, "direct") config.pushd.exchange.clone().into(),
.durable(true) lapin::ExchangeKind::Direct,
.finish(), ExchangeDeclareOptions {
durable: true,
..Default::default()
},
FieldTable::default(),
) )
.await .await
.expect("Failed to declare pushd exchange"); .expect("Failed to declare exchange");
let mut queue_name = queue_name.to_string(); let mut queue_name = queue_name.to_string();
@@ -230,35 +255,59 @@ where
let queue_name = queue_name.as_str(); let queue_name = queue_name.as_str();
let mut args = QueueDeclareArguments::new(queue_name); let args = QueueDeclareOptions {
args.durable(true); durable: true,
..Default::default()
if let Some(arg) = queue_args { };
args.arguments(arg);
}
let args = args.finish();
_ = channel.queue_declare(args).await.unwrap().unwrap();
channel channel
.queue_bind(QueueBindArguments::new( .queue_declare(queue_name.into(), args, queue_args.unwrap_or_default())
queue_name, .await
&config.pushd.exchange, .unwrap();
routing_key,
)) channel
.queue_bind(
queue_name.into(),
config.pushd.exchange.clone().into(),
routing_key.into(),
QueueBindOptions::default(),
FieldTable::default(),
)
.await .await
.expect( .expect(
"This probably means the revolt.notifications exchange does not exist in rabbitmq!", "This probably means the revolt.notifications exchange does not exist in rabbitmq!",
); );
let args = BasicConsumeArguments::new(queue_name, "") let consumer = channel
.manual_ack(false) .basic_consume(
.finish(); queue_name.into(),
"".into(),
let routing_key = channel.basic_consume(consumer, args).await.unwrap(); BasicConsumeOptions {
no_ack: true,
..Default::default()
},
FieldTable::default(),
)
.await
.unwrap();
info!( info!(
"Consuming routing key {} as queue {}, tag {}", "Consuming routing key {} as queue {}, tag {}",
routing_key, queue_name, routing_key routing_key,
queue_name,
consumer.tag()
); );
(channel, connection)
let delegate = Delegate(
F::create(
db.clone(),
authifier_db.clone(),
connection.clone(),
channel.clone(),
)
.await,
);
consumer.set_delegate(delegate);
channel
} }
@@ -0,0 +1,91 @@
use std::{
future::{ready, Future},
pin::Pin,
sync::Arc,
};
use anyhow::Result;
use async_trait::async_trait;
use lapin::{
message::{Delivery, DeliveryResult},
options::BasicPublishOptions,
BasicProperties, Channel, Connection, ConsumerDelegate, Error as AMQPError,
};
use log::debug;
use revolt_database::Database;
#[async_trait]
pub trait Consumer: Clone + Send + Sync + 'static {
async fn create(
db: Database,
authifier_db: authifier::Database,
connection: Arc<Connection>,
channel: Arc<Channel>,
) -> Self;
fn channel(&self) -> &Arc<Channel>;
async fn consume(&self, delivery: Delivery) -> Result<()>;
async fn publish_message_with_options(
&self,
payload: &[u8],
exchange: &str,
routing_key: &str,
options: BasicPublishOptions,
properties: BasicProperties,
) -> Result<(), AMQPError> {
let channel = self.channel();
channel
.basic_publish(
exchange.into(),
routing_key.into(),
options,
payload,
properties,
)
.await?;
debug!("Sent message to queue for target {}", routing_key);
Ok(())
}
async fn publish_message(
&self,
payload: &[u8],
exchange: &str,
routing_key: &str,
) -> Result<(), AMQPError> {
self.publish_message_with_options(
payload,
exchange,
routing_key,
BasicPublishOptions::default(),
BasicProperties::default(),
)
.await
}
}
pub struct Delegate<C: Consumer>(pub C);
impl<C: Consumer> ConsumerDelegate for Delegate<C> {
fn on_new_delivery(
&self,
delivery: DeliveryResult,
) -> Pin<Box<dyn Future<Output = ()> + Send>> {
match delivery {
Ok(Some(delivery)) => {
let consumer = self.0.clone();
Box::pin(async move {
if let Err(e) = consumer.consume(delivery).await {
revolt_config::capture_anyhow(&e);
log::error!("{e:?}");
};
})
}
Ok(None) => Box::pin(ready(())),
Err(e) => Box::pin(async move { log::error!("Received bad delivery: {e:?}") }),
}
}
}
+3
View File
@@ -1,2 +1,5 @@
mod renderer; mod renderer;
mod consumer;
pub use renderer::render_notification_content; pub use renderer::render_notification_content;
pub use consumer::{Consumer, Delegate};
+22 -25
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-voice-ingress" name = "revolt-voice-ingress"
version = "0.11.5" version = "0.13.7"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
edition = "2021" edition = "2021"
publish = false publish = false
@@ -9,41 +9,38 @@ publish = false
[dependencies] [dependencies]
# util # util
log = "*" log = { workspace = true }
sentry = "0.31.5" sentry = { workspace = true }
lru = "0.7.6" lru = { workspace = true }
ulid = "0.5.0" ulid = { workspace = true }
redis-kiss = "0.1.4" redis-kiss = { workspace = true }
chrono = "0.4.15" chrono = { workspace = true }
# Serde # Serde
serde_json = "1.0.79" serde_json = { workspace = true }
rmp-serde = "1.0.0" rmp-serde = { workspace = true }
serde = "1.0.136" serde = { workspace = true }
# Http # Http
rocket = { version = "0.5.0-rc.2", features = ["json"] } rocket = { workspace = true, features = ["json"] }
rocket_empty = "0.1.1" rocket_empty = { workspace = true }
# Async # Async
futures = "0.3.21" futures = { workspace = true }
async-std = { version = "1.8.0", features = [ async-std = { workspace = true, features = [
"tokio1", "tokio1",
"tokio02", "tokio02",
"attributes", "attributes",
] } ] }
# Core # Core
revolt-result = { path = "../../core/result" } revolt-result = { workspace = true, features = ["rocket"] }
revolt-models = { path = "../../core/models" } revolt-models = { workspace = true }
revolt-config = { path = "../../core/config" } revolt-config = { workspace = true }
revolt-database = { path = "../../core/database", features = ["voice"] } revolt-database = { workspace = true, features = ["voice"] }
revolt-permissions = { path = "../../core/permissions" } revolt-permissions = { workspace = true }
# Voice # Voice
livekit-api = "0.4.4" livekit-api = { workspace = true }
livekit-protocol = "0.4.0" livekit-protocol = { workspace = true }
livekit-runtime = { version = "0.3.1", features = ["tokio"] } livekit-runtime = { workspace = true, features = ["tokio"] }
# RabbitMQ
amqprs = { version = "1.7.0" }
+44 -45
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "revolt-delta" name = "revolt-delta"
version = "0.11.5" version = "0.13.7"
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
authors = ["Paul Makles <paulmakles@gmail.com>"] authors = ["Paul Makles <paulmakles@gmail.com>"]
edition = "2018" edition = "2018"
@@ -10,84 +10,83 @@ publish = false
[dependencies] [dependencies]
# Test # Test
rand = "0.8.5" rand = { workspace = true }
redis-kiss = "0.1.4" redis-kiss = { workspace = true }
# Utility # Utility
lru = "0.7.0" lru = { workspace = true }
url = "2.2.2" url = { workspace = true }
log = "0.4.11" log = { workspace = true }
dashmap = "5.2.0" dashmap = { workspace = true }
linkify = "0.6.0" linkify = { workspace = true }
once_cell = "1.17.1" once_cell = { workspace = true }
env_logger = "0.7.1"
# Lang. Utilities # Lang. Utilities
regex = "1" regex = { workspace = true }
num_enum = "0.5.1" num_enum = { workspace = true }
impl_ops = "0.1.1" impl_ops = { workspace = true }
bitfield = "0.13.2" bitfield = { workspace = true }
# ID / key generation # ID / key generation
ulid = "0.4.1" ulid = { workspace = true }
nanoid = "0.4.0" nanoid = { workspace = true }
# serde # serde
serde_json = "1.0.57" serde_json = { workspace = true }
serde = { version = "1.0.115", features = ["derive"] } serde = { workspace = true }
validator = { version = "0.16", features = ["derive"] } validator = { workspace = true, features = ["derive"] }
iso8601-timestamp = { version = "0.2.11", features = [] } iso8601-timestamp = { workspace = true }
# async # async
futures = "0.3.8" futures = { workspace = true }
chrono = "0.4.15" chrono = { workspace = true }
async-channel = "1.6.1" async-channel = { workspace = true }
reqwest = { version = "0.11.4", features = ["json"] } reqwest = { workspace = true, features = ["json"] }
async-std = { version = "1.8.0", features = [ async-std = { workspace = true, features = [
"tokio1", "tokio1",
"tokio02", "tokio02",
"attributes", "attributes",
] } ] }
# internal util # internal util
lettre = "0.10.0-alpha.4" lettre = { workspace = true }
# web # web
rocket = { version = "0.5.1", default-features = false, features = ["json"] } rocket = { workspace = true, features = ["json"] }
rocket_cors = { git = "https://github.com/lawliet89/rocket_cors", rev = "072d90359b23e9b291df6b672c07c93de9c46011" } rocket_cors = { workspace = true }
rocket_empty = { version = "0.1.1", features = ["schema"] } rocket_empty = { workspace = true, features = ["schema"] }
rocket_authifier = { version = "1.0.16" } rocket_authifier = { workspace = true }
rocket_prometheus = "0.10.0-rc.3" rocket_prometheus = { workspace = true }
# spec generation # spec generation
schemars = "0.8.8" schemars = { workspace = true }
revolt_rocket_okapi = { version = "0.10.0", features = ["swagger"] } revolt_rocket_okapi = { workspace = true, features = ["swagger"] }
# rabbit # rabbit
amqprs = { version = "1.7.0" } lapin = { workspace = true, features = ["tokio"] }
# core # core
authifier = "1.0.16" authifier = { workspace = true }
revolt-config = { path = "../core/config" } revolt-config = { workspace = true }
revolt-database = { path = "../core/database", features = [ revolt-database = { workspace = true, features = [
"rocket-impl", "rocket-impl",
"redis-is-patched", "redis-is-patched",
"voice", "voice",
] } ] }
revolt-models = { path = "../core/models", features = [ revolt-models = { workspace = true, features = [
"schemas", "schemas",
"validator", "validator",
"rocket", "rocket",
] } ] }
revolt-presence = { path = "../core/presence" } revolt-presence = { workspace = true }
revolt-result = { path = "../core/result", features = ["rocket", "okapi"] } revolt-result = { workspace = true, features = ["rocket", "okapi"] }
revolt-permissions = { path = "../core/permissions", features = ["schemas"] } revolt-permissions = { workspace = true, features = ["schemas"] }
revolt-ratelimits = { path = "../core/ratelimits", features = ["rocket"] } revolt-ratelimits = { workspace = true, features = ["rocket"] }
# voice # voice
livekit-api = "0.4.4" livekit-api = { workspace = true }
livekit-protocol = "0.4.0" livekit-protocol = { workspace = true }
[build-dependencies] [build-dependencies]
vergen = "7.5.0" vergen = { workspace = true }
+4 -48
View File
@@ -9,8 +9,7 @@ pub mod routes;
pub mod util; pub mod util;
use revolt_config::config; use revolt_config::config;
use revolt_database::events::client::EventV1; use revolt_database::{AMQP, events::client::EventV1};
use revolt_database::AMQP;
use revolt_ratelimits::rocket as ratelimiter; use revolt_ratelimits::rocket as ratelimiter;
use rocket::{Build, Rocket}; use rocket::{Build, Rocket};
use rocket_cors::{AllowedOrigins, CorsOptions}; use rocket_cors::{AllowedOrigins, CorsOptions};
@@ -18,14 +17,10 @@ use rocket_prometheus::PrometheusMetrics;
use std::net::Ipv4Addr; use std::net::Ipv4Addr;
use std::str::FromStr; use std::str::FromStr;
use amqprs::{
channel::ExchangeDeclareArguments,
connection::{Connection, OpenConnectionArguments},
};
use async_std::channel::unbounded; use async_std::channel::unbounded;
use authifier::AuthifierEvent; use authifier::AuthifierEvent;
use rocket::data::ToByteUnit;
use revolt_database::voice::VoiceClient; use revolt_database::voice::VoiceClient;
use rocket::data::ToByteUnit;
pub async fn web() -> Rocket<Build> { pub async fn web() -> Rocket<Build> {
// Get settings // Get settings
@@ -36,7 +31,6 @@ pub async fn web() -> Rocket<Build> {
// Setup database // Setup database
let db = revolt_database::DatabaseInfo::Auto.connect().await.unwrap(); let db = revolt_database::DatabaseInfo::Auto.connect().await.unwrap();
log::info!("database_here {db:?}");
db.migrate_database().await.unwrap(); db.migrate_database().await.unwrap();
// Setup Authifier event channel // Setup Authifier event channel
@@ -93,49 +87,11 @@ pub async fn web() -> Rocket<Build> {
) )
.into(); .into();
let swagger_0_8 = revolt_rocket_okapi::swagger_ui::make_swagger_ui(
&revolt_rocket_okapi::swagger_ui::SwaggerUIConfig {
url: "/0.8/openapi.json".to_owned(),
..Default::default()
},
)
.into();
let swagger_0_8 = revolt_rocket_okapi::swagger_ui::make_swagger_ui(
&revolt_rocket_okapi::swagger_ui::SwaggerUIConfig {
url: "/0.8/openapi.json".to_owned(),
..Default::default()
},
)
.into();
// Voice handler // Voice handler
let voice_client = VoiceClient::new(config.api.livekit.nodes.clone()); let voice_client = VoiceClient::new(config.api.livekit.nodes.clone());
// Configure Rabbit // Configure Rabbit
let connection = Connection::open(&OpenConnectionArguments::new(
&config.rabbit.host,
config.rabbit.port,
&config.rabbit.username,
&config.rabbit.password,
))
.await
.expect("Failed to connect to RabbitMQ");
let channel = connection let amqp = AMQP::new_auto().await;
.open_channel(None)
.await
.expect("Failed to open RabbitMQ channel");
channel
.exchange_declare(
ExchangeDeclareArguments::new(&config.pushd.exchange, "direct")
.durable(true)
.finish(),
)
.await
.expect("Failed to declare exchange");
let amqp = AMQP::new(connection, channel);
// Launch background task workers // Launch background task workers
revolt_database::tasks::start_workers(db.clone(), amqp.clone()); revolt_database::tasks::start_workers(db.clone(), amqp.clone());
@@ -153,7 +109,6 @@ pub async fn web() -> Rocket<Build> {
.mount("/", rocket_cors::catch_all_options_routes()) .mount("/", rocket_cors::catch_all_options_routes())
.mount("/", ratelimiter::routes()) .mount("/", ratelimiter::routes())
.mount("/swagger/", swagger) .mount("/swagger/", swagger)
.mount("/0.8/swagger/", swagger_0_8)
.manage(authifier) .manage(authifier)
.manage(db) .manage(db)
.manage(amqp) .manage(amqp)
@@ -166,6 +121,7 @@ pub async fn web() -> Rocket<Build> {
limits: rocket::data::Limits::default().limit("string", 5.megabytes()), limits: rocket::data::Limits::default().limit("string", 5.megabytes()),
address: Ipv4Addr::new(0, 0, 0, 0).into(), address: Ipv4Addr::new(0, 0, 0, 0).into(),
port: 14702, port: 14702,
ip_header: Some("X-Forwarded-For".into()),
..Default::default() ..Default::default()
}) })
} }
@@ -1,6 +1,6 @@
use revolt_database::{ use revolt_database::{
util::{permissions::DatabasePermissionQuery, reference::Reference}, util::{permissions::DatabasePermissionQuery, reference::Reference},
Database, User, Database, User, AMQP,
}; };
use revolt_permissions::{calculate_channel_permissions, ChannelPermission}; use revolt_permissions::{calculate_channel_permissions, ChannelPermission};
use revolt_result::{create_error, Result}; use revolt_result::{create_error, Result};
@@ -14,6 +14,7 @@ use rocket_empty::EmptyResponse;
#[put("/<target>/ack/<message>")] #[put("/<target>/ack/<message>")]
pub async fn ack( pub async fn ack(
db: &State<Database>, db: &State<Database>,
amqp: &State<AMQP>,
user: User, user: User,
target: Reference<'_>, target: Reference<'_>,
message: Reference<'_>, message: Reference<'_>,
@@ -29,7 +30,7 @@ pub async fn ack(
.throw_if_lacking_channel_permission(ChannelPermission::ViewChannel)?; .throw_if_lacking_channel_permission(ChannelPermission::ViewChannel)?;
channel channel
.ack(&user.id, message.id) .ack(&user.id, message.id, amqp)
.await .await
.map(|_| EmptyResponse) .map(|_| EmptyResponse)
} }
@@ -1,4 +1,5 @@
use chrono::Utc; use std::time::Duration;
use revolt_database::{ use revolt_database::{
util::{permissions::DatabasePermissionQuery, reference::Reference}, util::{permissions::DatabasePermissionQuery, reference::Reference},
Database, Message, User, Database, Message, User,
@@ -36,10 +37,9 @@ pub async fn bulk_delete_messages(
if ulid::Ulid::from_string(id) if ulid::Ulid::from_string(id)
.map_err(|_| create_error!(InvalidOperation))? .map_err(|_| create_error!(InvalidOperation))?
.datetime() .datetime()
.signed_duration_since(Utc::now()) .elapsed()
.num_days() .expect("Time went backwards")
.abs() > Duration::from_hours(7 * 24) // 7 days
> 7
{ {
return Err(create_error!(InvalidOperation)); return Err(create_error!(InvalidOperation));
} }
@@ -1,11 +1,14 @@
use chrono::{Duration, Utc}; use std::time::Duration;
use redis_kiss::{get_connection, redis, AsyncCommands}; use redis_kiss::{get_connection, redis, AsyncCommands};
use revolt_database::events::client::EventV1;
use revolt_database::util::permissions::DatabasePermissionQuery; use revolt_database::util::permissions::DatabasePermissionQuery;
use revolt_database::{ use revolt_database::{
util::idempotency::IdempotencyKey, util::reference::Reference, Database, User, util::idempotency::IdempotencyKey, util::reference::Reference, Database, User,
}; };
use revolt_database::{Channel, Interactions, Message, AMQP}; use revolt_database::{Channel, Interactions, Message, AMQP};
use revolt_models::v0; use revolt_models::v0;
use revolt_models::v0::ChannelSlowmode;
use revolt_permissions::PermissionQuery; use revolt_permissions::PermissionQuery;
use revolt_permissions::{calculate_channel_permissions, ChannelPermission}; use revolt_permissions::{calculate_channel_permissions, ChannelPermission};
use revolt_result::{create_error, Result}; use revolt_result::{create_error, Result};
@@ -83,6 +86,16 @@ pub async fn message_send(
.await .await
.unwrap_or(None); .unwrap_or(None);
if set_result.is_some() {
let idx_key = format!("slowmode_idx:{}", user.id);
conn.sadd::<_, _, ()>(&idx_key, channel_id.as_str())
.await
.ok();
conn.expire::<_, ()>(&idx_key, *channel_slowmode as usize)
.await
.ok();
}
// If `set_result` is None, the `NX` condition failed because the key already exists. // If `set_result` is None, the `NX` condition failed because the key already exists.
// This means the user is currently in slowmode. // This means the user is currently in slowmode.
if set_result.is_none() { if set_result.is_none() {
@@ -91,10 +104,29 @@ pub async fn message_send(
// Redis returns positive integers for valid TTLs // Redis returns positive integers for valid TTLs
if ttl > 0 { if ttl > 0 {
EventV1::UserSlowmodes {
slowmodes: vec![ChannelSlowmode {
channel_id: channel_id.to_string(),
duration: *channel_slowmode,
retry_after: ttl as u64,
}],
}
.private(user.id.clone())
.await;
return Err(create_error!(InSlowmode { return Err(create_error!(InSlowmode {
retry_after: ttl as u64 retry_after: ttl as u64
})); }));
} }
} else {
EventV1::UserSlowmodes {
slowmodes: vec![ChannelSlowmode {
channel_id: channel_id.to_string(),
duration: *channel_slowmode,
retry_after: *channel_slowmode,
}],
}
.private(user.id.clone())
.await;
} }
} }
// If Redis connection fails, just skip the slowmode check // If Redis connection fails, just skip the slowmode check
@@ -111,8 +143,12 @@ pub async fn message_send(
// Disallow mentions for new users (TRUST-0: <12 hours age) in public servers // Disallow mentions for new users (TRUST-0: <12 hours age) in public servers
let allow_mentions = if let Some(server) = query.server_ref() { let allow_mentions = if let Some(server) = query.server_ref() {
if server.discoverable { if server.discoverable {
(Utc::now() - ulid::Ulid::from_string(&user.id).unwrap().datetime()) (ulid::Ulid::from_string(&user.id)
>= Duration::hours(12) .unwrap()
.datetime()
.elapsed()
.expect("Time went backwards"))
>= Duration::from_hours(12)
} else { } else {
true true
} }
@@ -0,0 +1,212 @@
use revolt_database::{
util::{permissions::DatabasePermissionQuery, reference::Reference},
Database, EmojiParent, PartialEmoji, User,
};
use revolt_models::v0;
use revolt_permissions::{calculate_server_permissions, ChannelPermission};
use revolt_result::{create_error, Result};
use rocket::{serde::json::Json, State};
use validator::Validate;
/// # Edit Emoji
///
/// Edit an emoji by its id.
#[openapi(tag = "Emojis")]
#[patch("/emoji/<emoji_id>", data = "<data>")]
pub async fn edit_emoji(
db: &State<Database>,
user: User,
emoji_id: Reference<'_>,
data: Json<v0::DataEditEmoji>,
) -> Result<Json<v0::Emoji>> {
let data = data.into_inner();
data.validate().map_err(|error| {
create_error!(FailedValidation {
error: error.to_string()
})
})?;
let mut emoji = emoji_id.as_emoji(db).await?;
match &emoji.parent {
EmojiParent::Server { id } => {
let server = db.fetch_server(id.as_str()).await?;
let mut query = DatabasePermissionQuery::new(db, &user).server(&server);
calculate_server_permissions(&mut query)
.await
.throw_if_lacking_channel_permission(ChannelPermission::ManageCustomisation)?;
}
EmojiParent::Detached => return Err(create_error!(NotAuthenticated)),
}
if data.name.is_none() {
return Ok(Json(emoji.into()));
}
let partial = PartialEmoji { name: data.name };
emoji.update(db, partial).await?;
Ok(Json(emoji.into()))
}
#[cfg(test)]
mod test {
use crate::util::test::TestHarness;
use revolt_database::{Emoji, EmojiParent, Member};
use revolt_models::v0;
use rocket::http::{ContentType, Header, Status};
use ulid::Ulid;
#[rocket::async_test]
async fn edit_emoji_name_as_creator() {
let harness = TestHarness::new().await;
let (_, session, user) = harness.new_user().await;
let (server, _) = harness.new_server(&user).await;
let emoji_id = Ulid::new().to_string();
let emoji = Emoji {
id: emoji_id.clone(),
parent: EmojiParent::Server {
id: server.id.clone(),
},
creator_id: user.id.clone(),
name: "initial_name".to_string(),
animated: false,
nsfw: false,
};
emoji.create(&harness.db).await.expect("`Emoji` created");
let response = harness
.client
.patch(format!("/custom/emoji/{emoji_id}"))
.header(Header::new("x-session-token", session.token.to_string()))
.header(ContentType::JSON)
.body(
json!(v0::DataEditEmoji {
name: Some("renamed_emoji".to_string()),
})
.to_string(),
)
.dispatch()
.await;
assert_eq!(response.status(), Status::Ok);
let edited: v0::Emoji = response.into_json().await.expect("`Emoji`");
assert_eq!(edited.name, "renamed_emoji");
}
#[rocket::async_test]
async fn reject_invalid_emoji_name() {
let harness = TestHarness::new().await;
let (_, session, user) = harness.new_user().await;
let (server, _) = harness.new_server(&user).await;
let emoji_id = Ulid::new().to_string();
let emoji = Emoji {
id: emoji_id.clone(),
parent: EmojiParent::Server {
id: server.id.clone(),
},
creator_id: user.id.clone(),
name: "valid_name".to_string(),
animated: false,
nsfw: false,
};
emoji.create(&harness.db).await.expect("`Emoji` created");
let response = harness
.client
.patch(format!("/custom/emoji/{emoji_id}"))
.header(Header::new("x-session-token", session.token.to_string()))
.header(ContentType::JSON)
.body(
json!(v0::DataEditEmoji {
name: Some("Invalid Name".to_string()),
})
.to_string(),
)
.dispatch()
.await;
assert_eq!(response.status(), Status::BadRequest);
}
#[rocket::async_test]
async fn reject_edit_for_detached_emoji() {
let harness = TestHarness::new().await;
let (_, session, user) = harness.new_user().await;
let emoji_id = Ulid::new().to_string();
let emoji = Emoji {
id: emoji_id.clone(),
parent: EmojiParent::Detached,
creator_id: user.id.clone(),
name: "detached_name".to_string(),
animated: false,
nsfw: false,
};
emoji.create(&harness.db).await.expect("`Emoji` created");
let response = harness
.client
.patch(format!("/custom/emoji/{emoji_id}"))
.header(Header::new("x-session-token", session.token.to_string()))
.header(ContentType::JSON)
.body(
json!(v0::DataEditEmoji {
name: Some("should_not_apply".to_string()),
})
.to_string(),
)
.dispatch()
.await;
assert_eq!(response.status(), Status::Unauthorized);
}
#[rocket::async_test]
async fn reject_edit_for_creator_without_manage_customisation() {
let harness = TestHarness::new().await;
let (_, _, owner) = harness.new_user().await;
let (_, creator_session, creator) = harness.new_user().await;
let (server, _) = harness.new_server(&owner).await;
Member::create(&harness.db, &server, &creator, None)
.await
.expect("`Member` created");
let emoji_id = Ulid::new().to_string();
let emoji = Emoji {
id: emoji_id.clone(),
parent: EmojiParent::Server {
id: server.id.clone(),
},
creator_id: creator.id.clone(),
name: "member_uploaded_name".to_string(),
animated: false,
nsfw: false,
};
emoji.create(&harness.db).await.expect("`Emoji` created");
let response = harness
.client
.patch(format!("/custom/emoji/{emoji_id}"))
.header(Header::new(
"x-session-token",
creator_session.token.to_string(),
))
.header(ContentType::JSON)
.body(
json!(v0::DataEditEmoji {
name: Some("renamed_without_permission".to_string()),
})
.to_string(),
)
.dispatch()
.await;
assert_eq!(response.status(), Status::Forbidden);
}
}
@@ -3,12 +3,14 @@ use rocket::Route;
mod emoji_create; mod emoji_create;
mod emoji_delete; mod emoji_delete;
mod emoji_edit;
mod emoji_fetch; mod emoji_fetch;
pub fn routes() -> (Vec<Route>, OpenApi) { pub fn routes() -> (Vec<Route>, OpenApi) {
openapi_get_routes_spec![ openapi_get_routes_spec![
emoji_create::create_emoji, emoji_create::create_emoji,
emoji_delete::delete_emoji, emoji_delete::delete_emoji,
emoji_edit::edit_emoji,
emoji_fetch::fetch_emoji emoji_fetch::fetch_emoji
] ]
} }
+25 -73
View File
@@ -64,47 +64,6 @@ pub fn mount(config: Settings, mut rocket: Rocket<Build>) -> Rocket<Build> {
}; };
} }
if config.features.webhooks_enabled {
mount_endpoints_and_merged_docs! {
rocket, "/0.8".to_owned(), settings,
"/" => (vec![], custom_openapi_spec()),
"" => openapi_get_routes_spec![root::root],
"/users" => users::routes(),
"/bots" => bots::routes(),
"/channels" => channels::routes(),
"/servers" => servers::routes(),
"/invites" => invites::routes(),
"/custom" => customisation::routes(),
"/safety" => safety::routes(),
"/auth/account" => rocket_authifier::routes::account::routes(),
"/auth/session" => rocket_authifier::routes::session::routes(),
"/auth/mfa" => rocket_authifier::routes::mfa::routes(),
"/onboard" => onboard::routes(),
"/push" => push::routes(),
"/sync" => sync::routes(),
"/webhooks" => webhooks::routes()
};
} else {
mount_endpoints_and_merged_docs! {
rocket, "/0.8".to_owned(), settings,
"/" => (vec![], custom_openapi_spec()),
"" => openapi_get_routes_spec![root::root],
"/users" => users::routes(),
"/bots" => bots::routes(),
"/channels" => channels::routes(),
"/servers" => servers::routes(),
"/invites" => invites::routes(),
"/custom" => customisation::routes(),
"/safety" => safety::routes(),
"/auth/account" => rocket_authifier::routes::account::routes(),
"/auth/session" => rocket_authifier::routes::session::routes(),
"/auth/mfa" => rocket_authifier::routes::mfa::routes(),
"/onboard" => onboard::routes(),
"/push" => push::routes(),
"/sync" => sync::routes()
};
}
rocket rocket
} }
@@ -115,8 +74,8 @@ fn custom_openapi_spec() -> OpenApi {
extensions.insert( extensions.insert(
"x-logo".to_owned(), "x-logo".to_owned(),
json!({ json!({
"url": "https://revolt.chat/header.png", "url": "https://stoat.chat/header.png",
"altText": "Revolt Header" "altText": "Stoat Header"
}), }),
); );
@@ -124,7 +83,7 @@ fn custom_openapi_spec() -> OpenApi {
"x-tagGroups".to_owned(), "x-tagGroups".to_owned(),
json!([ json!([
{ {
"name": "Revolt", "name": "Stoat",
"tags": [ "tags": [
"Core" "Core"
] ]
@@ -205,18 +164,21 @@ fn custom_openapi_spec() -> OpenApi {
OpenApi { OpenApi {
openapi: OpenApi::default_version(), openapi: OpenApi::default_version(),
info: Info { info: Info {
title: "Revolt API".to_owned(), title: "Stoat API".to_owned(),
description: Some("Open source user-first chat platform.".to_owned()), description: Some("Open source user-first chat platform.".to_owned()),
terms_of_service: Some("https://revolt.chat/terms".to_owned()), terms_of_service: Some("https://stoat.chat/terms".to_owned()),
contact: Some(Contact { contact: Some(Contact {
name: Some("Revolt Support".to_owned()), name: Some("Stoat".to_owned()),
url: Some("https://revolt.chat".to_owned()), url: Some("https://stoat.chat".to_owned()),
email: Some("contact@revolt.chat".to_owned()), email: Some("contact@stoat.chat".to_owned()),
..Default::default() ..Default::default()
}), }),
license: Some(License { license: Some(License {
name: "AGPLv3".to_owned(), name: "AGPLv3".to_owned(),
url: Some("https://github.com/stoatchat/stoatchat/blob/main/crates/delta/LICENSE".to_owned()), url: Some(
"https://github.com/stoatchat/stoatchat/blob/main/crates/delta/LICENSE"
.to_owned(),
),
..Default::default() ..Default::default()
}), }),
version: env!("CARGO_PKG_VERSION").to_string(), version: env!("CARGO_PKG_VERSION").to_string(),
@@ -224,29 +186,19 @@ fn custom_openapi_spec() -> OpenApi {
}, },
servers: vec![ servers: vec![
Server { Server {
url: "https://api.revolt.chat".to_owned(), url: "https://api.stoat.chat".to_owned(),
description: Some("Revolt Production".to_owned()), description: Some("Stoat Production".to_owned()),
..Default::default() ..Default::default()
}, },
Server { Server {
url: "https://revolt.chat/api".to_owned(), url: "https://beta.stoat.chat/api".to_owned(),
description: Some("Revolt Staging".to_owned()), description: Some("Stoat Beta".to_owned()),
..Default::default()
},
Server {
url: "http://local.revolt.chat:14702".to_owned(),
description: Some("Local Revolt Environment".to_owned()),
..Default::default()
},
Server {
url: "http://local.revolt.chat:14702/0.8".to_owned(),
description: Some("Local Revolt Environment (v0.8)".to_owned()),
..Default::default() ..Default::default()
}, },
], ],
external_docs: Some(ExternalDocs { external_docs: Some(ExternalDocs {
url: "https://developers.revolt.chat".to_owned(), url: "https://developers.stoat.chat".to_owned(),
description: Some("Revolt Developer Documentation".to_owned()), description: Some("Stoat Developer Documentation".to_owned()),
..Default::default() ..Default::default()
}), }),
extensions, extensions,
@@ -254,19 +206,19 @@ fn custom_openapi_spec() -> OpenApi {
Tag { Tag {
name: "Core".to_owned(), name: "Core".to_owned(),
description: Some( description: Some(
"Use in your applications to determine information about the Revolt node" "Use in your applications to determine information about the Stoat node"
.to_owned(), .to_owned(),
), ),
..Default::default() ..Default::default()
}, },
Tag { Tag {
name: "User Information".to_owned(), name: "User Information".to_owned(),
description: Some("Query and fetch users on Revolt".to_owned()), description: Some("Query and fetch users on Stoat".to_owned()),
..Default::default() ..Default::default()
}, },
Tag { Tag {
name: "Direct Messaging".to_owned(), name: "Direct Messaging".to_owned(),
description: Some("Direct message other users on Revolt".to_owned()), description: Some("Direct message other users on Stoat".to_owned()),
..Default::default() ..Default::default()
}, },
Tag { Tag {
@@ -283,7 +235,7 @@ fn custom_openapi_spec() -> OpenApi {
}, },
Tag { Tag {
name: "Channel Information".to_owned(), name: "Channel Information".to_owned(),
description: Some("Query and fetch channels on Revolt".to_owned()), description: Some("Query and fetch channels on Stoat".to_owned()),
..Default::default() ..Default::default()
}, },
Tag { Tag {
@@ -313,7 +265,7 @@ fn custom_openapi_spec() -> OpenApi {
}, },
Tag { Tag {
name: "Server Information".to_owned(), name: "Server Information".to_owned(),
description: Some("Query and fetch servers on Revolt".to_owned()), description: Some("Query and fetch servers on Stoat".to_owned()),
..Default::default() ..Default::default()
}, },
Tag { Tag {
@@ -349,7 +301,7 @@ fn custom_openapi_spec() -> OpenApi {
Tag { Tag {
name: "Onboarding".to_owned(), name: "Onboarding".to_owned(),
description: Some( description: Some(
"After signing up to Revolt, users must pick a unique username".to_owned(), "After signing up to Stoat, users must pick a unique username".to_owned(),
), ),
..Default::default() ..Default::default()
}, },
@@ -361,7 +313,7 @@ fn custom_openapi_spec() -> OpenApi {
Tag { Tag {
name: "Web Push".to_owned(), name: "Web Push".to_owned(),
description: Some( description: Some(
"Subscribe to and receive Revolt push notifications while offline".to_owned(), "Subscribe to and receive Stoat push notifications while offline".to_owned(),
), ),
..Default::default() ..Default::default()
}, },
+21
View File
@@ -57,6 +57,8 @@ pub struct RevoltFeatures {
pub livekit: VoiceFeature, pub livekit: VoiceFeature,
/// Limits /// Limits
pub limits: LimitsConfig, pub limits: LimitsConfig,
/// Legal links
pub legal_links: LegalLinks,
} }
/// # Limits For Users /// # Limits For Users
@@ -70,6 +72,17 @@ pub struct LimitsConfig {
pub default: UserLimits, pub default: UserLimits,
} }
/// # Legal links
#[derive(Serialize, JsonSchema, Debug)]
pub struct LegalLinks {
/// Terms of Service URL
pub terms_of_service: String,
/// Privacy Policy URL
pub privacy_policy: String,
/// Guidelines URL
pub guidelines: String,
}
/// # Global limits /// # Global limits
#[derive(Serialize, JsonSchema, Debug)] #[derive(Serialize, JsonSchema, Debug)]
pub struct GlobalLimits { pub struct GlobalLimits {
@@ -92,6 +105,8 @@ pub struct GlobalLimits {
/// restrict server creation to these users. /// restrict server creation to these users.
/// if blank, all users can create servers /// if blank, all users can create servers
pub restrict_server_creation: Vec<String>, pub restrict_server_creation: Vec<String>,
/// New user hours
new_user_hours: i64,
} }
/// # User Limits /// # User Limits
@@ -231,10 +246,16 @@ pub async fn root() -> Result<Json<RevoltConfig>> {
.limits .limits
.global .global
.restrict_server_creation, .restrict_server_creation,
new_user_hours: config.features.limits.global.new_user_hours as i64,
}, },
new_user: UserLimits::from_feature_limits(config.features.limits.new_user), new_user: UserLimits::from_feature_limits(config.features.limits.new_user),
default: UserLimits::from_feature_limits(config.features.limits.default), default: UserLimits::from_feature_limits(config.features.limits.default),
}, },
legal_links: LegalLinks {
terms_of_service: config.features.legal_links.terms_of_service,
privacy_policy: config.features.legal_links.privacy_policy,
guidelines: config.features.legal_links.guidelines,
},
}, },
ws: config.hosts.events, ws: config.hosts.events,
app: config.hosts.app, app: config.hosts.app,
+17 -3
View File
@@ -1,7 +1,7 @@
use revolt_database::{ use revolt_database::{
util::{permissions::DatabasePermissionQuery, reference::Reference}, util::{permissions::DatabasePermissionQuery, reference::Reference},
voice::{sync_voice_permissions, VoiceClient}, voice::{sync_voice_permissions, VoiceClient},
Database, PartialRole, User Database, File, PartialRole, User,
}; };
use revolt_models::v0; use revolt_models::v0;
use revolt_permissions::{calculate_server_permissions, ChannelPermission}; use revolt_permissions::{calculate_server_permissions, ChannelPermission};
@@ -47,14 +47,27 @@ pub async fn edit(
name, name,
colour, colour,
hoist, hoist,
icon,
remove, remove,
.. ..
} = data; } = data;
if remove.contains(&v0::FieldsRole::Icon) {
if let Some(existing_icon) = &role.icon {
db.mark_attachment_as_deleted(&existing_icon.id).await?;
}
}
let mut final_icon = None;
if let Some(icon_id) = icon {
final_icon = Some(File::use_role_icon(db, &icon_id, &role_id, &user.id).await?);
}
let partial = PartialRole { let partial = PartialRole {
name, name,
colour, colour,
hoist, hoist,
icon: final_icon,
..Default::default() ..Default::default()
}; };
@@ -69,8 +82,9 @@ pub async fn edit(
for channel_id in &server.channels { for channel_id in &server.channels {
let channel = Reference::from_unchecked(channel_id).as_channel(db).await?; let channel = Reference::from_unchecked(channel_id).as_channel(db).await?;
sync_voice_permissions(db, voice_client, &channel, Some(&server), Some(&role_id)).await?; sync_voice_permissions(db, voice_client, &channel, Some(&server), Some(&role_id))
}; .await?;
}
Ok(Json(role.into())) Ok(Json(role.into()))
} else { } else {
+10 -6
View File
@@ -1,6 +1,6 @@
use revolt_database::{ use revolt_database::{
util::{permissions::DatabasePermissionQuery, reference::Reference}, util::{acker, permissions::DatabasePermissionQuery, reference::Reference},
Database, User, Database, User, AMQP,
}; };
use revolt_permissions::PermissionQuery; use revolt_permissions::PermissionQuery;
use revolt_result::{create_error, Result}; use revolt_result::{create_error, Result};
@@ -12,7 +12,12 @@ use rocket_empty::EmptyResponse;
/// Mark all channels in a server as read. /// Mark all channels in a server as read.
#[openapi(tag = "Server Information")] #[openapi(tag = "Server Information")]
#[put("/<target>/ack")] #[put("/<target>/ack")]
pub async fn ack(db: &State<Database>, user: User, target: Reference<'_>) -> Result<EmptyResponse> { pub async fn ack(
db: &State<Database>,
amqp: &State<AMQP>,
user: User,
target: Reference<'_>,
) -> Result<EmptyResponse> {
if user.bot.is_some() { if user.bot.is_some() {
return Err(create_error!(IsBot)); return Err(create_error!(IsBot));
} }
@@ -23,7 +28,6 @@ pub async fn ack(db: &State<Database>, user: User, target: Reference<'_>) -> Res
return Err(create_error!(NotFound)); return Err(create_error!(NotFound));
} }
db.acknowledge_channels(&user.id, &server.channels) acker::ack_server(&user, &server, db, amqp).await?;
.await Ok(EmptyResponse)
.map(|_| EmptyResponse)
} }
+7 -3
View File
@@ -1,19 +1,23 @@
use rocket::Route;
use revolt_rocket_okapi::revolt_okapi::openapi3::OpenApi; use revolt_rocket_okapi::revolt_okapi::openapi3::OpenApi;
use rocket::Route;
mod webhook_delete; mod webhook_delete;
mod webhook_delete_message;
mod webhook_delete_token; mod webhook_delete_token;
mod webhook_edit; mod webhook_edit;
mod webhook_edit_message;
mod webhook_edit_token; mod webhook_edit_token;
mod webhook_execute; mod webhook_execute;
mod webhook_fetch_token;
mod webhook_fetch;
mod webhook_execute_github; mod webhook_execute_github;
mod webhook_fetch;
mod webhook_fetch_token;
pub fn routes() -> (Vec<Route>, OpenApi) { pub fn routes() -> (Vec<Route>, OpenApi) {
openapi_get_routes_spec![ openapi_get_routes_spec![
webhook_delete_message::webhook_delete_message,
webhook_delete_token::webhook_delete_token, webhook_delete_token::webhook_delete_token,
webhook_delete::webhook_delete, webhook_delete::webhook_delete,
webhook_edit_message::webhook_edit_message,
webhook_edit_token::webhook_edit_token, webhook_edit_token::webhook_edit_token,
webhook_edit::webhook_edit, webhook_edit::webhook_edit,
webhook_execute_github::webhook_execute_github, webhook_execute_github::webhook_execute_github,
@@ -0,0 +1,27 @@
use revolt_database::{util::reference::Reference, Database};
use revolt_result::{create_error, Result};
use rocket::State;
use rocket_empty::EmptyResponse;
/// # Deletes a webhook message
///
/// Deletes a message sent by a webhook
#[openapi(tag = "Webhooks")]
#[delete("/<webhook_id>/<token>/<message_id>")]
pub async fn webhook_delete_message(
db: &State<Database>,
webhook_id: Reference<'_>,
token: String,
message_id: Reference<'_>,
) -> Result<EmptyResponse> {
let webhook = webhook_id.as_webhook(db).await?;
webhook.assert_token(&token)?;
let message = message_id.as_message(db).await?;
if message.author != webhook.id {
return Err(create_error!(CannotDeleteMessage));
}
message.delete(db).await.map(|_| EmptyResponse)
}

Some files were not shown because too many files have changed in this diff Show More