Compare commits

...
Author SHA1 Message Date
HampusandGitHub 14e772a751 fix(api): serialise elapsed temp bans as null for the admin panel (#1817) 2026-08-21 20:34:43 +02:00
HampusandGitHub 32afbf12d6 fix(api): stop treating accounts pending deletion as already deleted (#1816) 2026-08-21 20:31:32 +02:00
HampusandGitHub ffaf5119d8 fix(app): only warn about software encoding when no layer is accelerated (#1815) 2026-08-21 18:49:51 +02:00
HampusandGitHub 10bc8c1efa fix(desktop): stop the windows audio probe timing out against its own budget (#1814) 2026-08-21 17:29:12 +02:00
HampusandGitHub 85a03a9e39 fix(app): apply the screen share audio toggle to the surface being shared (#1813) 2026-08-21 17:04:22 +02:00
51ee6567b4 chore(i18n): update public marketing catalogs (#1812)
Co-authored-by: hampus-fluxer <[email protected]>
2026-08-21 16:31:01 +02:00
090220a29d chore(marketing): advance pointer f39eced → 02e3a2c (#1811)
Co-authored-by: hampus-fluxer <[email protected]>
2026-08-21 16:30:54 +02:00
HampusandGitHub e7973b8be0 fix(api): derive voice reconciliation candidate ttl from real sweep spacing (#1810) 2026-08-21 15:49:50 +02:00
HampusandGitHub b5324c9223 perf(gateway): stop materializing all members on the guild connect path (#1809) 2026-08-21 14:00:19 +02:00
HampusandGitHub edb8d80077 ci: source the s3 provider for downloads and static from repo variables (#1803) 2026-08-20 21:46:26 +02:00
HampusandGitHub ddee116339 feat(api): route downloads through the configured downloads provider (#1802) 2026-08-20 21:44:53 +02:00
HampusandGitHub bdacaea4a8 feat(config): add an optional separate s3 provider for downloads (#1801) 2026-08-20 21:34:19 +02:00
HampusandGitHub 9e28e02b5d ci(rust): pin the floating toolchains to the version images build with (#1800) 2026-08-20 21:27:12 +02:00
HampusandGitHub 3527dc95a2 fix(rust): silence result_large_err on axum response error paths (#1799) 2026-08-20 21:18:18 +02:00
HampusandGitHub 27c7b2722d chore(admin): regenerate openapi schemas for the admin acl cap (#1798) 2026-08-20 21:17:04 +02:00
HampusandGitHub ba96f52ed6 fix(schema): allow assigning every admin ACL to a user (#1797) 2026-08-20 21:09:56 +02:00
HampusandGitHub 2c8b3ff45c fix(build): build the messages and users images with scylla support (#1796) 2026-08-20 20:30:16 +02:00
HampusandGitHub 631bc2307a fix(slowmode): stop the local cooldown from outgrowing the channel setting (#1795) 2026-08-20 19:56:02 +02:00
HampusandGitHub 8f4f9a8601 refactor(media): share one external proxy url codec across every service (#1794) 2026-08-20 19:55:40 +02:00
100 changed files with 1466 additions and 628 deletions
+4 -4
View File
@@ -106,10 +106,10 @@ jobs:
- name: upload assets to S3 static bucket
env:
AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}
S3_ENDPOINT: https://ewr1.vultrobjects.com
STATIC_BUCKET: fluxer-static
AWS_ACCESS_KEY_ID: ${{ secrets.STATIC_AWS_ACCESS_KEY_ID || secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.STATIC_AWS_SECRET_ACCESS_KEY || secrets.AWS_SECRET_ACCESS_KEY }}
S3_ENDPOINT: ${{ vars.STATIC_S3_ENDPOINT }}
STATIC_BUCKET: ${{ vars.STATIC_S3_BUCKET }}
run: >-
cargo run --locked --quiet --manifest-path tools/ci/Cargo.toml -- build-app-proxy
--step upload_assets
+8 -8
View File
@@ -139,10 +139,10 @@ jobs:
SOURCE_SHA: ${{ needs.meta.outputs.source_sha }}
S3_DESKTOP_PREFIX: ${{ needs.meta.outputs.s3_prefix }}
DESKTOP_HANDOFF_PREFIX: _handoff/desktop/${{ needs.meta.outputs.build_channel }}/${{ needs.meta.outputs.version }}/${{ needs.meta.outputs.source_sha }}
S3_ENDPOINT: https://ewr1.vultrobjects.com
S3_BUCKET: fluxer-downloads
AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}
S3_ENDPOINT: ${{ vars.DOWNLOADS_S3_ENDPOINT }}
S3_BUCKET: ${{ vars.DOWNLOADS_S3_BUCKET }}
AWS_ACCESS_KEY_ID: ${{ secrets.DOWNLOADS_AWS_ACCESS_KEY_ID || secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.DOWNLOADS_AWS_SECRET_ACCESS_KEY || secrets.AWS_SECRET_ACCESS_KEY }}
DESKTOP_PLATFORM: ${{ matrix.platform }}
DESKTOP_ARCH: ${{ matrix.arch }}
DESKTOP_VARIANT: ${{ matrix.desktop_variant }}
@@ -524,11 +524,11 @@ jobs:
SOURCE_SHA: ${{ needs.meta.outputs.source_sha }}
S3_DESKTOP_PREFIX: ${{ needs.meta.outputs.s3_prefix }}
DESKTOP_HANDOFF_PREFIX: _handoff/desktop/${{ needs.meta.outputs.build_channel }}/${{ needs.meta.outputs.version }}/${{ needs.meta.outputs.source_sha }}
S3_ENDPOINT: https://ewr1.vultrobjects.com
S3_BUCKET: fluxer-downloads
S3_ENDPOINT: ${{ vars.DOWNLOADS_S3_ENDPOINT }}
S3_BUCKET: ${{ vars.DOWNLOADS_S3_BUCKET }}
PUBLIC_DL_BASE: https://api.fluxer.app/dl
AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}
AWS_ACCESS_KEY_ID: ${{ secrets.DOWNLOADS_AWS_ACCESS_KEY_ID || secrets.AWS_ACCESS_KEY_ID }}
AWS_SECRET_ACCESS_KEY: ${{ secrets.DOWNLOADS_AWS_SECRET_ACCESS_KEY || secrets.AWS_SECRET_ACCESS_KEY }}
steps:
- name: Checkout source
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0
+2 -2
View File
@@ -89,7 +89,7 @@ jobs:
- name: Install Rust toolchain
uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9
with:
toolchain: stable
toolchain: "1.93.0"
components: clippy, rustfmt
- name: Install pnpm
@@ -286,7 +286,7 @@ jobs:
- name: Install Rust toolchain
uses: dtolnay/rust-toolchain@e97e2d8cc328f1b50210efc529dca0028893a2d9
with:
toolchain: stable
toolchain: "1.93.0"
components: rustfmt
- name: Sync ci helper dependencies
Generated
+9
View File
@@ -1777,6 +1777,7 @@ dependencies = [
"axum",
"base64",
"clap",
"fluxer_common",
"hmac 0.13.0",
"hyper 1.10.1",
"hyper-util",
@@ -1800,6 +1801,7 @@ dependencies = [
"anyhow",
"base64",
"fluxer-svc",
"fluxer_common",
"hmac 0.13.0",
"moka",
"reqwest",
@@ -1836,6 +1838,7 @@ dependencies = [
"cc",
"clap",
"criterion",
"fluxer_common",
"hex",
"hmac 0.13.0",
"http 1.4.2",
@@ -1875,6 +1878,7 @@ dependencies = [
"chrono",
"criterion",
"fluxer-svc",
"fluxer_common",
"fluxer_markdown_parser",
"futures",
"hmac 0.13.0",
@@ -1943,6 +1947,7 @@ dependencies = [
"encoding_rs",
"entities",
"fluxer-svc",
"fluxer_common",
"hmac 0.13.0",
"infer",
"moka",
@@ -2038,10 +2043,14 @@ dependencies = [
"aws-credential-types",
"aws-sigv4",
"axum",
"base64",
"hmac 0.13.0",
"maxminddb",
"moka",
"reqwest",
"serde_json",
"sha2 0.11.0",
"thiserror",
"time",
"tracing",
"urlencoding",
+2 -2
View File
@@ -14913,7 +14913,7 @@
"pending_bulk_message_deletion_at": {"nullable": true, "type": "string"},
"deletion_reason_code": {"nullable": true, "allOf": [{"$ref": "#/components/schemas/Int32Type"}]},
"deletion_public_reason": {"nullable": true, "type": "string"},
"acls": {"type": "array", "items": {"type": "string"}, "maxItems": 100},
"acls": {"type": "array", "items": {"type": "string"}, "maxItems": 115},
"traits": {"type": "array", "items": {"type": "string"}, "maxItems": 100},
"has_totp": {"type": "boolean"},
"authenticator_types": {"type": "array", "items": {"$ref": "#/components/schemas/Int32Type"}, "maxItems": 10},
@@ -15481,7 +15481,7 @@
"acls": {
"type": "array",
"items": {"type": "string"},
"maxItems": 100,
"maxItems": 115,
"description": "List of access control permissions to assign"
}
},
+1
View File
@@ -112,6 +112,7 @@ fn is_urlencoded_form(request: &Request) -> bool {
})
}
#[allow(clippy::result_large_err)]
async fn extract_csrf_from_form_body(
request: Request,
) -> Result<(Request, Option<String>), Response> {
+10
View File
@@ -1,6 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {MasterConfig} from '@fluxer/config/src/MasterConfig';
import {resolveDownloadsProvider} from '@fluxer/config/src/S3DownloadsProvider';
import {parseIpAddress} from '@fluxer/ip_utils/src/IpAddress';
import {parseGeoipSourceConfig, resolveGeoipRuntimeSourceConfig} from '@pkgs/geoip/src/GeoipStartup';
import type {APIConfig, BlueskyOAuthConfig} from './config/APIConfig';
@@ -237,6 +238,7 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig {
cacheMinTtlSeconds: master.services.api.embeds.cache_min_ttl_seconds,
cacheRespectRemoteTtl: master.services.api.embeds.cache_respect_remote_ttl,
},
s3Downloads: resolveDownloadsProvider(master),
s3: {
endpoint: s3Config.endpoint,
presignedUrlBase: s3Config.presigned_url_base,
@@ -461,6 +463,14 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig {
taskName: apiWorkerConfig?.task as WorkerTaskName | undefined,
enableCronScheduler: apiWorkerConfig?.enable_cron_scheduler,
enableVoiceReconciliation: apiWorkerConfig?.enable_voice_reconciliation ?? true,
voiceReconciliation: {
intervalMs: apiWorkerConfig?.voice_reconciliation?.interval_ms,
staggerDelayMs: apiWorkerConfig?.voice_reconciliation?.stagger_delay_ms,
lockTtlSeconds: apiWorkerConfig?.voice_reconciliation?.lock_ttl_seconds,
cadenceTtlSeconds: apiWorkerConfig?.voice_reconciliation?.cadence_ttl_seconds,
gatewayOnlyGraceMs: apiWorkerConfig?.voice_reconciliation?.gateway_only_grace_ms,
liveKitOnlyGraceMs: apiWorkerConfig?.voice_reconciliation?.livekit_only_grace_ms,
},
laneConcurrencyOverrides: {
realtime: apiWorkerConfig?.lane_concurrency_overrides?.realtime,
unfurl: apiWorkerConfig?.lane_concurrency_overrides?.unfurl,
+2 -1
View File
@@ -84,7 +84,8 @@ export async function mapUserToAdminResponse(
premium_lifetime_sequence: user.premiumLifetimeSequence ?? null,
suspicious_activity_flags: user.suspiciousActivityFlags,
phone_verification_deferred: ((user.suspiciousActivityFlags ?? 0) & DEFERRED_PHONE_ON_COMMUNITY_JOIN) !== 0,
temp_banned_until: user.tempBannedUntil?.toISOString() ?? null,
temp_banned_until:
user.tempBannedUntil && user.tempBannedUntil.getTime() > Date.now() ? user.tempBannedUntil.toISOString() : null,
pending_deletion_at: user.pendingDeletionAt?.toISOString() ?? null,
pending_bulk_message_deletion_at: user.pendingBulkMessageDeletionAt?.toISOString() ?? null,
deletion_reason_code: user.deletionReasonCode,
@@ -929,9 +929,8 @@ export class MessageSendService {
algorithm: 'leaky_bucket',
});
if (!slowmodeResult.allowed) {
const retryAfter = Math.max(0, slowmodeResult.resetTime.getTime() - Date.now());
throw new SlowmodeRateLimitError({
retryAfter,
retryAfter: slowmodeResult.retryAfter,
retryAfterDecimal: slowmodeResult.retryAfterDecimal,
});
}
@@ -52,4 +52,27 @@ describe('Slowmode Enforcement', () => {
expect(messages).toHaveLength(1);
expect(messages[0]?.id).toBe(firstMessage.id);
});
it('reports the slowmode retry window in seconds on the Retry-After header', async () => {
const rateLimitPerUser = 5;
const {owner, members, guild} = await setupTestGuildWithMembers(harness, 1);
const member = members[0]!;
await ensureSessionStarted(harness, member.token);
const channel = await createChannel(harness, owner.token, guild.id, 'slowmode-channel');
await updateChannel(harness, owner.token, channel.id, {
rate_limit_per_user: rateLimitPerUser,
});
await sendChannelMessage(harness, member.token, channel.id, 'first message');
const {response, json} = await createBuilder<{code: string; retry_after: number}>(harness, member.token)
.post(`/channels/${channel.id}/messages`)
.body({content: 'second message'})
.expect(400, APIErrorCodes.SLOWMODE_RATE_LIMITED)
.executeWithResponse();
const headerRetryAfter = Number(response.headers.get('Retry-After'));
expect(Number.isInteger(headerRetryAfter)).toBe(true);
expect(headerRetryAfter).toBeGreaterThan(0);
expect(headerRetryAfter).toBeLessThanOrEqual(rateLimitPerUser);
expect(json.retry_after).toBeGreaterThan(0);
expect(json.retry_after).toBeLessThanOrEqual(rateLimitPerUser);
expect(headerRetryAfter - json.retry_after).toBeLessThan(1);
});
});
+10
View File
@@ -1,5 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {ResolvedDownloadsProvider} from '@fluxer/config/src/S3DownloadsProvider';
import type {WorkerTaskName} from '../worker/WorkerLaneConfig';
export type APIWorkerMode = 'all_lanes' | 'single_lane' | 'single_task';
@@ -150,6 +151,7 @@ export interface APIConfig {
static: string;
};
};
s3Downloads: ResolvedDownloadsProvider;
email: {
enabled: boolean;
provider: 'smtp' | 'none';
@@ -359,6 +361,14 @@ export interface APIConfig {
taskName?: WorkerTaskName;
enableCronScheduler?: boolean;
enableVoiceReconciliation: boolean;
voiceReconciliation: {
intervalMs: number | undefined;
staggerDelayMs: number | undefined;
lockTtlSeconds: number | undefined;
cadenceTtlSeconds: number | undefined;
gatewayOnlyGraceMs: number | undefined;
liveKitOnlyGraceMs: number | undefined;
};
laneConcurrencyOverrides: {
realtime?: number;
unfurl?: number;
@@ -187,3 +187,26 @@ describe('StorageService.copyObjectWithMetadataStripping', () => {
]);
});
});
describe('provider selection', () => {
interface ClientProbe {
client: {config: {region: () => Promise<string>; endpoint?: () => Promise<{hostname: string}>}};
}
it('defaults to the shared S3 configuration', async () => {
const service = new StorageService() as unknown as ClientProbe;
expect(await service.client.config.region()).toBe(Config.s3.region);
});
it('uses an explicitly supplied provider instead of the shared one', async () => {
const service = new StorageService({
endpoint: 'https://downloads.example.net',
forcePathStyle: false,
region: 'eu-central-9',
accessKeyId: 'DL_KEY',
secretAccessKey: 'DL_SECRET',
}) as unknown as ClientProbe;
expect(await service.client.config.region()).toBe('eu-central-9');
expect(await service.client.config.region()).not.toBe(Config.s3.region);
});
});
@@ -1,5 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {S3ProviderSettings} from '@fluxer/config/src/S3DownloadsProvider';
import assert from 'node:assert/strict';
import {createHash} from 'node:crypto';
import fs from 'node:fs';
@@ -105,27 +106,36 @@ function extractStreamFromGet(out: GetObjectCommandOutput): Readable {
export class StorageService implements IStorageService {
private readonly client: S3Client;
private readonly presignClient: S3Client;
private readonly provider: S3ProviderSettings;
constructor() {
this.client = buildPooledS3Client({
constructor(provider?: S3ProviderSettings) {
this.provider = provider ?? {
endpoint: Config.s3.endpoint,
presignedUrlBase: Config.s3.presignedUrlBase,
forcePathStyle: Config.s3.forcePathStyle,
region: Config.s3.region,
accessKeyId: Config.s3.accessKeyId,
secretAccessKey: Config.s3.secretAccessKey,
};
this.client = buildPooledS3Client({
endpoint: this.provider.endpoint,
region: this.provider.region,
accessKeyId: this.provider.accessKeyId,
secretAccessKey: this.provider.secretAccessKey,
forcePathStyle: true,
});
this.presignClient = buildPooledS3Client({
endpoint: this.resolvePresignEndpoint(),
region: Config.s3.region,
accessKeyId: Config.s3.accessKeyId,
secretAccessKey: Config.s3.secretAccessKey,
forcePathStyle: Config.s3.forcePathStyle,
region: this.provider.region,
accessKeyId: this.provider.accessKeyId,
secretAccessKey: this.provider.secretAccessKey,
forcePathStyle: this.provider.forcePathStyle,
});
}
private resolvePresignEndpoint(): string {
const fallbackEndpoint = Config.s3.endpoint;
const configuredEndpoint = Config.s3.presignedUrlBase;
const fallbackEndpoint = this.provider.endpoint;
const configuredEndpoint = this.provider.presignedUrlBase;
if (!configuredEndpoint) {
return fallbackEndpoint;
}
@@ -0,0 +1,18 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {describe, expect, it} from 'vitest';
import {Config} from '../Config';
import {createDownloadsStorageService} from './StorageServiceFactory';
describe('createDownloadsStorageService', () => {
it('returns null when no downloads override is configured', () => {
expect(Config.s3Downloads.isOverridden).toBe(false);
expect(createDownloadsStorageService()).toBeNull();
});
it('resolves the downloads provider to the shared provider by default', () => {
expect(Config.s3Downloads.settings.endpoint).toBe(Config.s3.endpoint);
expect(Config.s3Downloads.settings.region).toBe(Config.s3.region);
expect(Config.s3Downloads.settings.accessKeyId).toBe(Config.s3.accessKeyId);
});
});
@@ -1,8 +1,16 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Config} from '../Config';
import type {IStorageService} from './IStorageService';
import {StorageService} from './StorageService';
export function createStorageService(): IStorageService {
return new StorageService();
}
export function createDownloadsStorageService(): IStorageService | null {
if (!Config.s3Downloads.isOverridden) {
return null;
}
return new StorageService(Config.s3Downloads.settings);
}
@@ -60,7 +60,7 @@ import {KVActivityTracker} from '../infrastructure/KVActivityTracker';
import {KVBulkMessageDeletionQueueService} from '../infrastructure/KVBulkMessageDeletionQueueService';
import {NatsUnfurlerService} from '../infrastructure/NatsUnfurlerService';
import {PremiumStateReconciliationQueueService} from '../infrastructure/PremiumStateReconciliationQueueService';
import {createStorageService} from '../infrastructure/StorageServiceFactory';
import {createDownloadsStorageService, createStorageService} from '../infrastructure/StorageServiceFactory';
import {UserCacheService} from '../infrastructure/UserCacheService';
import {createUsersServiceClient} from '../infrastructure/UsersServiceClient';
import {VirusScanService} from '../infrastructure/VirusScanService';
@@ -193,6 +193,10 @@ export const getStorageService: () => IStorageService = (() => {
const fallback = singleton(() => createStorageService());
return () => _injectedStorageService ?? fallback();
})();
const getDownloadsStorageService: () => IStorageService = (() => {
const override = singleton(() => createDownloadsStorageService());
return () => override() ?? getStorageService();
})();
export const getErrorI18nService = singleton(() => new ErrorI18nService());
export const getLimitConfigService = singleton(
() => new LimitConfigService(getInstanceConfigRepository(), getCacheService(), getKVClient()),
@@ -262,7 +266,7 @@ export function getKVAccountDeletionQueue(): KVAccountDeletionQueueService {
return accountDeletionQueue;
}
export const getDownloadService = singleton(() => new DownloadService(getStorageService()));
export const getDownloadService = singleton(() => new DownloadService(getDownloadsStorageService()));
export const getThemeService = singleton(() => new ThemeService(getStorageService()));
const getNcmecReporter = singleton(() => new NcmecReporter({config: createNcmecApiConfig(), fetch}));
const getNcmecRepository = singleton(() => new NcmecRepository());
@@ -808,6 +808,9 @@ export class UserRelationshipService {
if (!user) {
return false;
}
if (user.pendingDeletionAt !== null) {
return false;
}
return (user.flags & UserFlags.DELETED) === UserFlags.DELETED;
}
}
@@ -19,6 +19,16 @@ import {
} from './RelationshipTestUtils';
import {fetchUserMe} from './UserTestUtils';
async function markUserScheduledForDeletion(harness: ApiTestHarness, userId: string): Promise<void> {
const pendingDeletionAt = new Date(Date.now() + 14 * 24 * 60 * 60 * 1000).toISOString();
await createBuilder(harness, '')
.post(`/test/users/${userId}/set-pending-deletion`)
.body({pending_deletion_at: pendingDeletionAt, set_self_deleted_flag: false})
.expect(HTTP_STATUS.OK)
.execute();
await markUserDeleted(harness, userId);
}
async function markUserDeleted(harness: ApiTestHarness, userId: string): Promise<void> {
await createBuilder(harness, '')
.patch(`/test/users/${userId}/flags`)
@@ -221,6 +231,16 @@ describe('UserRelationshipStateTransitions', () => {
.expect(HTTP_STATUS.BAD_REQUEST, 'FRIEND_REQUEST_BLOCKED')
.execute();
});
test('can accept a friend request from a user scheduled for deletion', async () => {
const alice = await createTestAccount(harness);
const bob = await createTestAccount(harness);
await sendFriendRequest(harness, bob.token, alice.userId);
await markUserScheduledForDeletion(harness, bob.userId);
const {json: friendship} = await acceptFriendRequest(harness, alice.token, bob.userId);
assertRelationshipType(friendship, RelationshipTypes.FRIEND);
const {json: aliceAfter} = await listRelationships(harness, alice.token);
assertRelationshipType(findRelationship(aliceAfter, bob.userId)!, RelationshipTypes.FRIEND);
});
test('cannot send friend request to user who blocked you', async () => {
const alice = await createTestAccount(harness);
const bob = await createTestAccount(harness);
@@ -0,0 +1,40 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {describe, expect, it} from 'vitest';
import {candidateTtlSecondsFor} from './VoiceReconciliationWorker';
const INTERVAL_MS = 15000;
const GATEWAY_ONLY_GRACE_MS = 10000;
function ttlFor(observedSweepSpacingMs: number): number {
return candidateTtlSecondsFor({
intervalMs: INTERVAL_MS,
observedSweepSpacingMs,
graceMs: GATEWAY_ONLY_GRACE_MS,
});
}
describe('candidateTtlSecondsFor', () => {
it('outlives the gap between two consecutive observations of the same key', () => {
for (const observedSweepSpacingMs of [0, 45_000, 136_000, 300_000, 596_000, 900_000]) {
expect(ttlFor(observedSweepSpacingMs) * 1000).toBeGreaterThan(observedSweepSpacingMs);
}
});
it('outlives a sweep gap far longer than the tick interval', () => {
expect(ttlFor(596_000) * 1000).toBeGreaterThan(596_000);
});
it('grows with the observed sweep spacing rather than the tick interval', () => {
expect(ttlFor(596_000)).toBeGreaterThan(ttlFor(136_000));
expect(ttlFor(136_000)).toBeGreaterThan(ttlFor(0));
});
it('keeps a floor that survives a single long sweep before any spacing is observed', () => {
expect(ttlFor(0)).toBeGreaterThanOrEqual(300);
});
it('stays bounded so a stale candidate cannot outlive its connection indefinitely', () => {
expect(ttlFor(Number.MAX_SAFE_INTEGER)).toBeLessThanOrEqual(3600);
});
});
@@ -96,9 +96,24 @@ const DEFAULT_STAGGER_DELAY_MS = 25;
const DEFAULT_LOCK_TTL_SECONDS = 180;
const DEFAULT_GATEWAY_ONLY_GRACE_MS = 10000;
const DEFAULT_LIVEKIT_ONLY_GRACE_MS = 60000;
const MIN_CANDIDATE_TTL_SECONDS = 300;
const MAX_CANDIDATE_TTL_SECONDS = 3600;
const CANDIDATE_TTL_SWEEP_MULTIPLIER = 3;
const LAST_SWEEP_KEY_TTL_SECONDS = 86400;
const ROOM_KEY_PREFIX = 'voice:room:server:';
export function candidateTtlSecondsFor(input: {
intervalMs: number;
observedSweepSpacingMs: number;
graceMs: number;
}): number {
const spacingMs = Math.max(input.intervalMs, input.observedSweepSpacingMs);
const ttlSeconds = Math.ceil((spacingMs * CANDIDATE_TTL_SWEEP_MULTIPLIER + input.graceMs * 2) / 1000);
return Math.min(MAX_CANDIDATE_TTL_SECONDS, Math.max(MIN_CANDIDATE_TTL_SECONDS, ttlSeconds));
}
const VOICE_RECONCILIATION_LOCK_KEY = 'voice:reconcile:lock';
const VOICE_RECONCILIATION_CADENCE_KEY = 'voice:reconcile:cadence';
const VOICE_RECONCILIATION_LAST_SWEEP_KEY = 'voice:reconcile:last-sweep-at';
const GATEWAY_ONLY_CANDIDATE_KEY_PREFIX = 'voice:reconcile:gateway-only:';
const LIVEKIT_ONLY_CANDIDATE_KEY_PREFIX = 'voice:reconcile:livekit-only:';
@@ -115,8 +130,7 @@ export class VoiceReconciliationWorker {
private readonly cadenceTtlSeconds: number;
private readonly gatewayOnlyGraceMs: number;
private readonly liveKitOnlyGraceMs: number;
private readonly gatewayOnlyCandidateTtlSeconds: number;
private readonly liveKitOnlyCandidateTtlSeconds: number;
private observedSweepSpacingMs = 0;
private intervalHandle: NodeJS.Timeout | null = null;
private reconciling = false;
private reconciliationLockLost = false;
@@ -136,14 +150,6 @@ export class VoiceReconciliationWorker {
this.cadenceTtlSeconds = options.cadenceTtlSeconds ?? Math.max(1, Math.ceil((this.intervalMs * 3) / 1000));
this.gatewayOnlyGraceMs = options.gatewayOnlyGraceMs ?? DEFAULT_GATEWAY_ONLY_GRACE_MS;
this.liveKitOnlyGraceMs = options.liveKitOnlyGraceMs ?? DEFAULT_LIVEKIT_ONLY_GRACE_MS;
this.gatewayOnlyCandidateTtlSeconds = Math.max(
60,
Math.ceil((this.intervalMs * 4 + this.gatewayOnlyGraceMs * 4) / 1000),
);
this.liveKitOnlyCandidateTtlSeconds = Math.max(
60,
Math.ceil((this.intervalMs * 4 + this.liveKitOnlyGraceMs * 4) / 1000),
);
}
start(): void {
@@ -156,6 +162,7 @@ export class VoiceReconciliationWorker {
intervalMs: this.intervalMs,
gatewayOnlyGraceMs: this.gatewayOnlyGraceMs,
liveKitOnlyGraceMs: this.liveKitOnlyGraceMs,
gatewayOnlyCandidateTtlSeconds: this.candidateTtlSeconds(this.gatewayOnlyGraceMs),
},
'Starting VoiceReconciliationWorker',
);
@@ -175,6 +182,7 @@ export class VoiceReconciliationWorker {
async reconcile(): Promise<void> {
const startTime = Date.now();
await this.recordSweepSpacing(startTime);
const discovery = await this.discoverActiveRooms();
this.logger.info(
{
@@ -233,6 +241,8 @@ export class VoiceReconciliationWorker {
const durationMs = Date.now() - startTime;
this.logger.info(
{
observedSweepSpacingMs: this.observedSweepSpacingMs,
gatewayOnlyCandidateTtlSeconds: this.candidateTtlSeconds(this.gatewayOnlyGraceMs),
roomsChecked,
totalConfirmed,
totalRepaired,
@@ -300,6 +310,37 @@ export class VoiceReconciliationWorker {
}
}
private candidateTtlSeconds(graceMs: number): number {
return candidateTtlSecondsFor({
intervalMs: this.intervalMs,
observedSweepSpacingMs: this.observedSweepSpacingMs,
graceMs,
});
}
private async recordSweepSpacing(startedAt: number): Promise<void> {
try {
const previous = await this.kvClient.get(VOICE_RECONCILIATION_LAST_SWEEP_KEY);
const previousAt = previous === null ? Number.NaN : Number(previous);
if (Number.isFinite(previousAt) && startedAt > previousAt) {
this.observedSweepSpacingMs = Math.max(this.observedSweepSpacingMs, startedAt - previousAt);
const gatewayOnlyCandidateTtlSeconds = this.candidateTtlSeconds(this.gatewayOnlyGraceMs);
if (this.observedSweepSpacingMs >= gatewayOnlyCandidateTtlSeconds * 1000) {
this.logger.warn(
{
observedSweepSpacingMs: this.observedSweepSpacingMs,
gatewayOnlyCandidateTtlSeconds,
},
'Reconciliation sweeps are further apart than the candidate TTL; divergent voice states will be deferred forever',
);
}
}
await this.kvClient.setex(VOICE_RECONCILIATION_LAST_SWEEP_KEY, LAST_SWEEP_KEY_TTL_SECONDS, String(startedAt));
} catch (error) {
this.logger.warn({error}, 'Failed to record reconciliation sweep spacing');
}
}
private async acquireCadenceLease(): Promise<boolean> {
try {
return await this.kvClient.setnx(VOICE_RECONCILIATION_CADENCE_KEY, '1', this.cadenceTtlSeconds);
@@ -993,7 +1034,7 @@ export class VoiceReconciliationWorker {
}
await this.kvClient.setex(
key,
this.gatewayOnlyCandidateTtlSeconds,
this.candidateTtlSeconds(this.gatewayOnlyGraceMs),
Number.isFinite(firstSeen) ? String(firstSeen) : String(now),
);
return false;
@@ -1024,7 +1065,7 @@ export class VoiceReconciliationWorker {
}
await this.kvClient.setex(
key,
this.liveKitOnlyCandidateTtlSeconds,
this.candidateTtlSeconds(this.liveKitOnlyGraceMs),
Number.isFinite(firstSeen) ? String(firstSeen) : String(now),
);
return false;
@@ -243,6 +243,12 @@ export async function initializeWorkerDependencies(snowflakeService: ISnowflakeS
voiceRoomStore,
kvClient,
logger: Logger,
intervalMs: Config.worker.voiceReconciliation.intervalMs,
staggerDelayMs: Config.worker.voiceReconciliation.staggerDelayMs,
lockTtlSeconds: Config.worker.voiceReconciliation.lockTtlSeconds,
cadenceTtlSeconds: Config.worker.voiceReconciliation.cadenceTtlSeconds,
gatewayOnlyGraceMs: Config.worker.voiceReconciliation.gatewayOnlyGraceMs,
liveKitOnlyGraceMs: Config.worker.voiceReconciliation.liveKitOnlyGraceMs,
})
: null;
if (Config.voice.enabled && voiceTopology !== null) {
@@ -834,7 +834,7 @@ export async function forward(
logger.warn(`Forward send failed in channel ${channelId}`);
return false;
}
SlowmodeCommands.recordMessageSend(channelId);
SlowmodeCommands.confirmMessageSend(channelId, forwardedMessage.timestamp);
if (optionalMessage) {
const commentNonce = SnowflakeUtils.fromTimestamp(Date.now() + 1);
const commentMessage = await send(channelId, {
@@ -845,7 +845,7 @@ export async function forward(
logger.warn(`Forward comment send failed in channel ${channelId}`);
return false;
}
SlowmodeCommands.recordMessageSend(channelId);
SlowmodeCommands.confirmMessageSend(channelId, commentMessage.timestamp);
}
}
logger.debug('Successfully forwarded message to all channels');
@@ -24,6 +24,7 @@ import {
normalizeMessageContent,
} from '@app/features/messaging/utils/MessageRequestUtils';
import * as MessageSubmitUtils from '@app/features/messaging/utils/MessageSubmitUtils';
import {resolveRetryAfterMs} from '@app/features/messaging/utils/RetryAfterUtils';
import {MatureContentRejectedModal} from '@app/features/moderation/components/alerts/MatureContentRejectedModal';
import {http} from '@app/features/platform/transport/RestTransport';
import {HttpError} from '@app/features/platform/types/EndpointError';
@@ -66,9 +67,7 @@ type ScheduledMessageRequest = MessageCreateRequest & {
interface ApiErrorBody {
code?: number | string;
retry_after?: number;
message?: string;
details?: unknown;
}
export interface ScheduleMessageParams {
@@ -372,37 +371,6 @@ const getApiErrorBody = (error: HttpError): ApiErrorBody | undefined => {
return typeof error.body === 'object' && error.body !== null ? (error.body as ApiErrorBody) : undefined;
};
function parseNestedRetryAfterSeconds(body: ApiErrorBody | undefined): number | undefined {
if (body === undefined || typeof body.details !== 'object' || body.details === null || Array.isArray(body.details)) {
return undefined;
}
const details = body.details as Record<string, unknown>;
const retry = details.retry;
if (typeof retry !== 'object' || retry === null || Array.isArray(retry)) return undefined;
const afterSeconds = (retry as Record<string, unknown>).after_seconds;
if (typeof afterSeconds !== 'number' || !Number.isFinite(afterSeconds) || afterSeconds <= 0) return undefined;
return afterSeconds;
}
function resolveRetryAfterSeconds(error: HttpError): number | undefined {
const body = getApiErrorBody(error);
const nestedRetryAfter = parseNestedRetryAfterSeconds(body);
if (nestedRetryAfter !== undefined) return nestedRetryAfter;
const bodyRetryAfter = body === undefined ? undefined : body.retry_after;
if (typeof bodyRetryAfter === 'number' && Number.isFinite(bodyRetryAfter) && bodyRetryAfter > 0) {
return bodyRetryAfter;
}
const responseHeaders: Record<string, string> | undefined = error.responseHeaders;
const header = responseHeaders === undefined ? undefined : responseHeaders['retry-after'];
if (header === undefined || header.trim() === '') return undefined;
const numeric = Number(header);
if (Number.isFinite(numeric) && numeric > 0) return numeric;
const deadline = Date.parse(header);
if (!Number.isFinite(deadline)) return undefined;
const remainingSeconds = (deadline - Date.now()) / 1000;
return remainingSeconds > 0 ? remainingSeconds : undefined;
}
function handleScheduleError(
i18n: I18n,
error: unknown,
@@ -430,8 +398,7 @@ function handleScheduleError(
return;
}
if (isSlowmodeError(error)) {
const retryAfterSeconds = resolveRetryAfterSeconds(error);
const retryAfterMs = SlowmodeCommands.retryAfterSecondsToMs(retryAfterSeconds);
const retryAfterMs = SlowmodeCommands.clampSlowmodeRetryAfterMs(resolveRetryAfterMs(error));
if (retryAfterMs <= 0) {
ModalCommands.push(
modal(() => (
@@ -490,8 +457,8 @@ function handleScheduleError(
}
function handleScheduleRateLimit(_i18n: I18n, error: HttpError): void {
const retryAfterSecondsValue = resolveRetryAfterSeconds(error);
const retryAfterSeconds = retryAfterSecondsValue === undefined ? undefined : Math.ceil(retryAfterSecondsValue);
const retryAfterMs = resolveRetryAfterMs(error);
const retryAfterSeconds = retryAfterMs === null ? undefined : Math.ceil(retryAfterMs / 1000);
ModalCommands.push(
modal(() => (
<MessageSendTooQuickModal
@@ -131,6 +131,7 @@ export const useMessageSubmission = ({channel, referencedMessage, replyingMessag
referenced_message: referencedMessage?.toJSON(),
});
SlowmodeCommands.prepareMessageSend(channel.id);
const pendingSend = SlowmodeCommands.recordPendingMessageSend(channel.id);
void MessageCommands.send(channel.id, {
content: message.content,
nonce,
@@ -141,11 +142,17 @@ export const useMessageSubmission = ({channel, referencedMessage, replyingMessag
stickers,
favoriteMemeId,
tts,
}).then((sentMessage) => {
if (sentMessage) {
SlowmodeCommands.recordMessageSend(channel.id);
}
});
})
.then((sentMessage) => {
if (sentMessage) {
SlowmodeCommands.confirmMessageSend(channel.id, sentMessage.timestamp, pendingSend);
return;
}
SlowmodeCommands.discardPendingMessageSend(channel.id, pendingSend);
})
.catch(() => {
SlowmodeCommands.discardPendingMessageSend(channel.id, pendingSend);
});
ComponentDispatch.dispatch('MESSAGE_SENT', {channelId: channel.id});
return true;
},
@@ -196,6 +203,7 @@ export const useMessageSubmission = ({channel, referencedMessage, replyingMessag
});
SlowmodeCommands.prepareMessageSend(channel.id);
const allowedMentions: AllowedMentions = {replied_user: replyingMessage?.mentioning ?? true};
const pendingSend = SlowmodeCommands.recordPendingMessageSend(channel.id);
void MessageCommands.send(channel.id, {
content: messageData.content,
nonce,
@@ -207,11 +215,17 @@ export const useMessageSubmission = ({channel, referencedMessage, replyingMessag
flags: 0,
stickers: messageData.stickers || [],
favoriteMemeId: sendOptions.favoriteMemeId,
}).then((sentMessage) => {
if (sentMessage) {
SlowmodeCommands.recordMessageSend(channel.id);
}
});
})
.then((sentMessage) => {
if (sentMessage) {
SlowmodeCommands.confirmMessageSend(channel.id, sentMessage.timestamp, pendingSend);
return;
}
SlowmodeCommands.discardPendingMessageSend(channel.id, pendingSend);
})
.catch(() => {
SlowmodeCommands.discardPendingMessageSend(channel.id, pendingSend);
});
ComponentDispatch.dispatch('MESSAGE_SENT', {channelId: channel.id});
},
[channel?.id, referencedMessage, replyingMessage],
@@ -38,6 +38,7 @@ import {
type MessageEditRequest,
normalizeMessageEditContent,
} from '@app/features/messaging/utils/MessageRequestUtils';
import {resolveRetryAfterMs} from '@app/features/messaging/utils/RetryAfterUtils';
import {MatureContentRejectedModal} from '@app/features/moderation/components/alerts/MatureContentRejectedModal';
import {http} from '@app/features/platform/transport/RestTransport';
import {HttpError} from '@app/features/platform/types/EndpointError';
@@ -141,9 +142,7 @@ export interface RetryError {
export interface ApiErrorBody {
code?: number | string;
retry_after?: number;
message?: string;
details?: unknown;
}
interface PresignedAttachmentUploadSinglepartResponse {
@@ -225,70 +224,25 @@ const getApiErrorBody = (error: HttpError): ApiErrorBody | undefined => {
};
interface MessageRateLimitRetry {
retryAfterSeconds: number | null;
retryAfterMs: number | null;
automaticRetryDelayMs: number | null;
}
function parsePositiveRetryAfterSeconds(value: unknown): number | null {
if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0) return null;
return value;
}
function parseRetryAfterHeaderSeconds(value: string | undefined): number | null {
if (value === undefined || value.trim() === '') return null;
const numeric = Number(value);
if (Number.isFinite(numeric) && numeric > 0) return numeric;
const deadline = Date.parse(value);
if (!Number.isFinite(deadline)) return null;
const remainingSeconds = (deadline - Date.now()) / 1000;
return remainingSeconds > 0 ? remainingSeconds : null;
}
function parseNestedRetryAfterSeconds(body: ApiErrorBody | undefined): number | null {
if (body === undefined || typeof body.details !== 'object' || body.details === null || Array.isArray(body.details)) {
return null;
}
const details = body.details as Record<string, unknown>;
const retry = details.retry;
if (typeof retry !== 'object' || retry === null || Array.isArray(retry)) return null;
const afterSeconds = (retry as Record<string, unknown>).after_seconds;
if (typeof afterSeconds !== 'number' || !Number.isFinite(afterSeconds) || afterSeconds <= 0) return null;
return afterSeconds;
}
function readRateLimitHeader(error: HttpError, name: string): string | undefined {
return error.responseHeaders[name.toLowerCase()];
}
function resolveMessageRateLimitRetry(error: HttpError): MessageRateLimitRetry {
const body = getApiErrorBody(error);
const candidates: Array<number> = [];
const nestedRetryAfter = parseNestedRetryAfterSeconds(body);
if (nestedRetryAfter !== null) candidates.push(nestedRetryAfter);
const bodyRetryAfter = body === undefined ? undefined : body.retry_after;
const parsedBodyRetryAfter = parsePositiveRetryAfterSeconds(bodyRetryAfter);
if (parsedBodyRetryAfter !== null) candidates.push(parsedBodyRetryAfter);
const parsedHeaderRetryAfter = parseRetryAfterHeaderSeconds(readRateLimitHeader(error, 'retry-after'));
if (parsedHeaderRetryAfter !== null) candidates.push(parsedHeaderRetryAfter);
const resetAfterHeader = readRateLimitHeader(error, 'x-ratelimit-reset-after');
if (resetAfterHeader !== undefined) {
const parsedResetAfter = parsePositiveRetryAfterSeconds(Number(resetAfterHeader));
if (parsedResetAfter !== null) candidates.push(parsedResetAfter);
}
if (candidates.length === 0) {
return {retryAfterSeconds: null, automaticRetryDelayMs: null};
}
const rawRetryAfter = Math.max(...candidates);
const retryAfterSeconds = Math.ceil(rawRetryAfter);
const retryAfterMs = Math.ceil(rawRetryAfter * 1000);
if (!Number.isSafeInteger(retryAfterSeconds) || !Number.isSafeInteger(retryAfterMs)) {
return {retryAfterSeconds: null, automaticRetryDelayMs: null};
const retryAfterMs = resolveRetryAfterMs(error);
if (retryAfterMs === null) {
return {retryAfterMs: null, automaticRetryDelayMs: null};
}
let automaticRetryDelayMs: number | null = null;
if (retryAfterMs <= MESSAGE_SEND_RATE_LIMIT_MAX_AUTOMATIC_DELAY_MS) {
automaticRetryDelayMs = retryAfterMs;
}
return {retryAfterSeconds, automaticRetryDelayMs};
return {retryAfterMs, automaticRetryDelayMs};
}
function retryAfterMsToWholeSeconds(retryAfterMs: number | null): number | null {
if (retryAfterMs === null) return null;
return Math.ceil(retryAfterMs / 1000);
}
const isAbortError = (error: unknown): boolean => {
return error instanceof DOMException && error.name === 'AbortError';
@@ -1344,7 +1298,7 @@ export class MessageQueue extends Queue<MessageQueuePayload, RestResponse<Messag
this.restoreFailedMessage(payload.channelId, payload.nonce);
}
completed(null, undefined, error);
this.handleRateLimitError(retry.retryAfterSeconds);
this.handleRateLimitError(retryAfterMsToWholeSeconds(retry.retryAfterMs));
}
private handleSendError(
@@ -1425,9 +1379,7 @@ export class MessageQueue extends Queue<MessageQueuePayload, RestResponse<Messag
);
} else if (error instanceof HttpError && isSlowmodeError(error)) {
const retry = resolveMessageRateLimitRetry(error);
const retryAfterMs = SlowmodeCommands.retryAfterSecondsToMs(
retry.retryAfterSeconds === null ? undefined : retry.retryAfterSeconds,
);
const retryAfterMs = SlowmodeCommands.clampSlowmodeRetryAfterMs(retry.retryAfterMs);
if (retryAfterMs <= 0) {
ModalCommands.push(
modal(() => (
@@ -1579,7 +1531,7 @@ export class MessageQueue extends Queue<MessageQueuePayload, RestResponse<Messag
return;
}
completed(null, undefined, error);
this.handleEditRateLimitError(retry.retryAfterSeconds);
this.handleEditRateLimitError(retryAfterMsToWholeSeconds(retry.retryAfterMs));
}
private showEditErrorModal(error: HttpError): void {
@@ -0,0 +1,65 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {resolveRetryAfterMs} from '@app/features/messaging/utils/RetryAfterUtils';
import {HttpError} from '@app/features/platform/types/EndpointError';
import {APIErrorCodes} from '@fluxer/constants/src/ApiErrorCodes';
import {describe, expect, it} from 'vitest';
function slowmodeRejection(body: unknown, responseHeaders: Record<string, string>): HttpError {
return new HttpError({
method: 'POST',
path: '/channels/1234567890123456789/messages',
status: 400,
body,
responseHeaders,
});
}
describe('resolveRetryAfterMs', () => {
it('uses the decimal seconds from the body instead of the Retry-After header', () => {
const error = slowmodeRejection(
{code: APIErrorCodes.SLOWMODE_RATE_LIMITED, retry_after: 4.44},
{'retry-after': '5'},
);
expect(resolveRetryAfterMs(error)).toBe(4440);
});
it('never multiplies a millisecond Retry-After header into a longer window than the body', () => {
const error = slowmodeRejection(
{code: APIErrorCodes.SLOWMODE_RATE_LIMITED, retry_after: 4.7},
{'retry-after': '4700'},
);
expect(resolveRetryAfterMs(error)).toBe(4700);
});
it('prefers a nested retry window over every other source', () => {
const error = slowmodeRejection(
{code: APIErrorCodes.SLOWMODE_RATE_LIMITED, retry_after: 30, details: {retry: {after_seconds: 2.5}}},
{'retry-after': '30'},
);
expect(resolveRetryAfterMs(error)).toBe(2500);
});
it('falls back to the Retry-After header when the body carries no window', () => {
const error = slowmodeRejection({code: APIErrorCodes.SLOWMODE_RATE_LIMITED}, {'retry-after': '3'});
expect(resolveRetryAfterMs(error)).toBe(3000);
});
it('falls back to the reset-after header when nothing else is present', () => {
const error = slowmodeRejection({code: APIErrorCodes.RATE_LIMITED}, {'x-ratelimit-reset-after': '1.25'});
expect(resolveRetryAfterMs(error)).toBe(1250);
});
it('reads an HTTP date Retry-After header as a remaining duration', () => {
const deadline = new Date(Date.now() + 4000).toUTCString();
const error = slowmodeRejection({code: APIErrorCodes.SLOWMODE_RATE_LIMITED}, {'retry-after': deadline});
const retryAfterMs = resolveRetryAfterMs(error);
expect(retryAfterMs).not.toBeNull();
expect(retryAfterMs!).toBeGreaterThan(0);
expect(retryAfterMs!).toBeLessThanOrEqual(4000);
});
it('returns null when no retry window is advertised', () => {
expect(resolveRetryAfterMs(slowmodeRejection({code: APIErrorCodes.SLOWMODE_RATE_LIMITED}, {}))).toBeNull();
});
});
@@ -0,0 +1,47 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {HttpError} from '@app/features/platform/types/EndpointError';
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value);
}
function parseSeconds(value: unknown): number | null {
if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0) return null;
return value;
}
function parseNestedSeconds(body: Record<string, unknown> | undefined): number | null {
if (body === undefined || !isRecord(body.details)) return null;
const retry = body.details.retry;
if (!isRecord(retry)) return null;
return parseSeconds(retry.after_seconds);
}
function parseHeaderSeconds(value: string | undefined): number | null {
if (value === undefined || value.trim() === '') return null;
const numeric = Number(value);
if (Number.isFinite(numeric)) return parseSeconds(numeric);
const deadline = Date.parse(value);
if (!Number.isFinite(deadline)) return null;
return parseSeconds((deadline - Date.now()) / 1000);
}
function resolveRetryAfterSeconds(error: HttpError): number | null {
const body = isRecord(error.body) ? error.body : undefined;
const nestedSeconds = parseNestedSeconds(body);
if (nestedSeconds !== null) return nestedSeconds;
const bodySeconds = body === undefined ? null : parseSeconds(body.retry_after);
if (bodySeconds !== null) return bodySeconds;
const headerSeconds = parseHeaderSeconds(error.responseHeaders['retry-after']);
if (headerSeconds !== null) return headerSeconds;
return parseHeaderSeconds(error.responseHeaders['x-ratelimit-reset-after']);
}
export function resolveRetryAfterMs(error: HttpError): number | null {
const retryAfterSeconds = resolveRetryAfterSeconds(error);
if (retryAfterSeconds === null) return null;
const retryAfterMs = Math.ceil(retryAfterSeconds * 1000);
if (!Number.isSafeInteger(retryAfterMs)) return null;
return retryAfterMs;
}
@@ -6,32 +6,47 @@ import {CHANNEL_RATE_LIMIT_PER_USER_MAX} from '@fluxer/constants/src/LimitConsta
const MAX_RETRY_AFTER_MS = CHANNEL_RATE_LIMIT_PER_USER_MAX * 1000;
function clearSendScopedSticker(channelId: string): void {
ChannelSticker.clearPendingStickerOnMessageSend(channelId);
export interface PendingMessageSend {
readonly previousSendTimestamp: number | null;
readonly pendingSendTimestamp: number;
}
function markSlowmodeSend(channelId: string): void {
Slowmode.recordMessageSend(channelId);
function clearSendScopedSticker(channelId: string): void {
ChannelSticker.clearPendingStickerOnMessageSend(channelId);
}
export function prepareMessageSend(channelId: string): void {
clearSendScopedSticker(channelId);
}
export function recordMessageSend(channelId: string): void {
markSlowmodeSend(channelId);
export function recordPendingMessageSend(channelId: string): PendingMessageSend {
const previousSendTimestamp = Slowmode.getLastSendTimestamp(channelId);
const pendingSendTimestamp = Slowmode.recordMessageSend(channelId);
return {previousSendTimestamp, pendingSendTimestamp};
}
export function confirmMessageSend(channelId: string, sentAt: string, pending?: PendingMessageSend): void {
const timestamp = Date.parse(sentAt);
if (!Number.isFinite(timestamp)) return;
const floor = pending?.pendingSendTimestamp;
const anchored = floor == null ? timestamp : Math.max(timestamp, floor);
Slowmode.updateSlowmodeTimestamp(channelId, anchored);
}
export function discardPendingMessageSend(channelId: string, pending: PendingMessageSend): void {
if (Slowmode.getLastSendTimestamp(channelId) !== pending.pendingSendTimestamp) return;
Slowmode.updateSlowmodeTimestamp(channelId, pending.previousSendTimestamp);
}
export function updateSlowmodeRemaining(channelId: string, retryAfterMs: number): void {
Slowmode.updateSlowmodeRemaining(channelId, retryAfterMs);
}
export function retryAfterSecondsToMs(retryAfterSeconds: number | undefined): number {
if (retryAfterSeconds == null || !Number.isFinite(retryAfterSeconds) || retryAfterSeconds <= 0) {
export function clampSlowmodeRetryAfterMs(retryAfterMs: number | null): number {
if (retryAfterMs == null || !Number.isSafeInteger(retryAfterMs) || retryAfterMs <= 0) {
return 0;
}
const retryAfterMs = Math.ceil(retryAfterSeconds * 1000);
if (!Number.isSafeInteger(retryAfterMs) || retryAfterMs > MAX_RETRY_AFTER_MS) {
if (retryAfterMs > MAX_RETRY_AFTER_MS) {
return 0;
}
return retryAfterMs;
@@ -52,7 +52,7 @@ class Slowmode {
this.pruneExpired(Date.now());
}
recordMessageSend(channelId: string): void {
recordMessageSend(channelId: string): number {
const now = Date.now();
this.pruneExpired(now);
const current = this.getEntry(channelId);
@@ -60,17 +60,22 @@ class Slowmode {
explicitExpiresAt: current.explicitExpiresAt,
lastSendTimestamp: now,
});
return now;
}
updateSlowmodeTimestamp(channelId: string, timestamp: number): void {
updateSlowmodeTimestamp(channelId: string, timestamp: number | null): void {
const now = Date.now();
if (!isValidTimestamp(timestamp, now)) return;
let nextTimestamp = timestamp;
if (nextTimestamp !== null) {
if (!isValidTimestamp(nextTimestamp, now)) return;
nextTimestamp = Math.min(nextTimestamp, now);
}
this.pruneExpired(now);
const current = this.getEntry(channelId);
if (current.lastSendTimestamp === timestamp) return;
if (current.lastSendTimestamp === nextTimestamp) return;
this.setEntry(channelId, {
explicitExpiresAt: current.explicitExpiresAt,
lastSendTimestamp: timestamp,
lastSendTimestamp: nextTimestamp,
});
}
@@ -0,0 +1,206 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {resolveRetryAfterMs} from '@app/features/messaging/utils/RetryAfterUtils';
import {HttpError} from '@app/features/platform/types/EndpointError';
import * as SlowmodeCommands from '@app/features/slowmode/commands/SlowmodeCommands';
import Slowmode from '@app/features/slowmode/state/Slowmode';
import {APIErrorCodes} from '@fluxer/constants/src/ApiErrorCodes';
import {afterEach, beforeEach, describe, expect, it, vi} from 'vitest';
const CHANNEL_ID = '1234567890123456789';
const RATE_LIMIT_PER_USER = 5;
const SLOWMODE_WINDOW_MS = RATE_LIMIT_PER_USER * 1000;
const T0 = Date.UTC(2026, 7, 20, 12, 0, 0);
function at(offsetMs: number): void {
vi.setSystemTime(T0 + offsetMs);
}
function shownRemainingMs(): number {
return Slowmode.getSlowmodeRemaining(CHANNEL_ID, RATE_LIMIT_PER_USER);
}
function shownCountdownSeconds(): number {
return Math.ceil(shownRemainingMs() / 1000);
}
function composerBlocksSend(): boolean {
return shownRemainingMs() > 0;
}
function serverTimestamp(offsetMs: number): string {
return new Date(T0 + offsetMs).toISOString();
}
function slowmodeRejection(retryAfterDecimalSeconds: number, retryAfterHeader: string): HttpError {
return new HttpError({
method: 'POST',
path: `/channels/${CHANNEL_ID}/messages`,
status: 400,
body: {code: APIErrorCodes.SLOWMODE_RATE_LIMITED, retry_after: retryAfterDecimalSeconds},
responseHeaders: {'retry-after': retryAfterHeader},
});
}
function applySlowmodeRejection(error: HttpError): number {
const retryAfterMs = SlowmodeCommands.clampSlowmodeRetryAfterMs(resolveRetryAfterMs(error));
SlowmodeCommands.updateSlowmodeRemaining(CHANNEL_ID, retryAfterMs);
return retryAfterMs;
}
describe('slowmode while an earlier send is still pending', () => {
beforeEach(() => {
vi.useFakeTimers();
at(0);
Slowmode.clearChannel(CHANNEL_ID);
});
afterEach(() => {
Slowmode.clearChannel(CHANNEL_ID);
vi.useRealTimers();
});
it('blocks the second message while the first one is still pending', () => {
at(0);
expect(composerBlocksSend()).toBe(false);
SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(120);
expect(composerBlocksSend()).toBe(true);
expect(shownRemainingMs()).toBe(SLOWMODE_WINDOW_MS - 120);
});
it('anchors the window to the timestamp the server assigned the message', () => {
at(0);
SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(450);
SlowmodeCommands.confirmMessageSend(CHANNEL_ID, serverTimestamp(200));
expect(Slowmode.getLastSendTimestamp(CHANNEL_ID)).toBe(T0 + 200);
expect(shownRemainingMs()).toBe(SLOWMODE_WINDOW_MS - 250);
});
it('never shows more than the channel setting across the reported ordering', () => {
const countdown: Array<{atMs: number; shownMs: number; shownSeconds: number}> = [];
const sample = (offsetMs: number): void => {
at(offsetMs);
countdown.push({atMs: offsetMs, shownMs: shownRemainingMs(), shownSeconds: shownCountdownSeconds()});
};
at(0);
SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
sample(0);
sample(120);
sample(449);
at(450);
SlowmodeCommands.confirmMessageSend(CHANNEL_ID, serverTimestamp(200));
sample(450);
sample(2000);
sample(5199);
sample(5200);
sample(10200);
for (const entry of countdown) {
expect(entry.shownMs).toBeLessThanOrEqual(SLOWMODE_WINDOW_MS);
expect(entry.shownSeconds).toBeLessThanOrEqual(RATE_LIMIT_PER_USER);
}
for (let index = 1; index < countdown.length; index++) {
expect(countdown[index]!.shownSeconds).toBeLessThanOrEqual(countdown[index - 1]!.shownSeconds);
}
expect(countdown.at(-3)!.shownMs).toBeGreaterThan(0);
expect(countdown.at(-2)!.shownMs).toBe(0);
expect(countdown.at(-1)!.shownMs).toBe(0);
});
it('keeps a server slowmode rejection inside the channel setting', () => {
at(0);
SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(450);
SlowmodeCommands.confirmMessageSend(CHANNEL_ID, serverTimestamp(200));
at(759);
const beforeRejection = shownRemainingMs();
at(760);
const storedMs = applySlowmodeRejection(slowmodeRejection(4.44, '5'));
expect(storedMs).toBe(4440);
expect(shownRemainingMs()).toBe(SLOWMODE_WINDOW_MS - 560);
expect(shownRemainingMs()).toBeLessThan(beforeRejection);
at(5200);
expect(shownRemainingMs()).toBe(0);
expect(composerBlocksSend()).toBe(false);
});
it('does not inflate the window when the rejection header carries milliseconds', () => {
at(0);
SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(450);
SlowmodeCommands.confirmMessageSend(CHANNEL_ID, serverTimestamp(200));
at(760);
const storedMs = applySlowmodeRejection(slowmodeRejection(4.7, '4700'));
expect(storedMs).toBe(4700);
expect(shownRemainingMs()).toBeLessThanOrEqual(SLOWMODE_WINDOW_MS);
at(5460);
expect(shownRemainingMs()).toBe(0);
});
it('releases the window when the pending send never reaches the server', () => {
at(0);
const pendingSend = SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(120);
expect(composerBlocksSend()).toBe(true);
at(300);
SlowmodeCommands.discardPendingMessageSend(CHANNEL_ID, pendingSend);
expect(Slowmode.getLastSendTimestamp(CHANNEL_ID)).toBeNull();
expect(composerBlocksSend()).toBe(false);
});
it('keeps a rejection window when the rejected send releases its own guess', () => {
at(0);
const pendingSend = SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(760);
applySlowmodeRejection(slowmodeRejection(4.44, '5'));
SlowmodeCommands.discardPendingMessageSend(CHANNEL_ID, pendingSend);
expect(Slowmode.getLastSendTimestamp(CHANNEL_ID)).toBeNull();
expect(shownRemainingMs()).toBe(4440);
});
it('ignores an anchor from a clock that runs behind the server', () => {
at(0);
SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
at(450);
SlowmodeCommands.confirmMessageSend(CHANNEL_ID, serverTimestamp(30_000));
expect(shownRemainingMs()).toBeLessThanOrEqual(SLOWMODE_WINDOW_MS);
at(5450);
expect(shownRemainingMs()).toBe(0);
});
});
describe('clock skew on the send anchor', () => {
const CHANNEL_ID = '900000000000000001';
const RATE_LIMIT_SECONDS = 5;
beforeEach(() => {
vi.useFakeTimers();
Slowmode.clearChannel(CHANNEL_ID);
});
afterEach(() => {
Slowmode.clearChannel(CHANNEL_ID);
vi.useRealTimers();
});
const remainingAfterAck = (clientAheadMs: number): number => {
vi.setSystemTime(new Date(1_000_000));
const pending = SlowmodeCommands.recordPendingMessageSend(CHANNEL_ID);
const serverAcceptedAt = new Date(1_000_000 - clientAheadMs).toISOString();
vi.setSystemTime(new Date(1_000_450));
SlowmodeCommands.confirmMessageSend(CHANNEL_ID, serverAcceptedAt, pending);
return Slowmode.getSlowmodeRemaining(CHANNEL_ID, RATE_LIMIT_SECONDS);
};
it('does not shorten the window when the client clock runs ahead of the server', () => {
for (const aheadMs of [0, 250, 1_000, 2_500, 4_800, 30_000, 3_600_000]) {
expect(remainingAfterAck(aheadMs)).toBe(4_550);
}
});
it('never lets the local guard reach zero while the window is live', () => {
for (const aheadMs of [4_800, 30_000, 3_600_000]) {
expect(remainingAfterAck(aheadMs)).toBeGreaterThan(0);
}
});
});
@@ -327,11 +327,17 @@ export async function pushActiveStreamSettings(
interface StreamSettingsMenuContentProps {
applyToLiveStream?: boolean;
shareContext?: StreamSettingsShareContext;
shareContextResolved?: boolean;
displayShareEnvironment: DisplayShareEnvironment;
}
export const StreamSettingsMenuContent = observer(
({applyToLiveStream = true, shareContext = 'display', displayShareEnvironment}: StreamSettingsMenuContentProps) => {
({
applyToLiveStream = true,
shareContext = 'display',
shareContextResolved = true,
displayShareEnvironment,
}: StreamSettingsMenuContentProps) => {
const {i18n} = useLingui();
useMediaEngineVersion();
const hasHigherVideoQuality = useHasHigherVideoQuality();
@@ -511,16 +517,22 @@ export const StreamSettingsMenuContent = observer(
);
const handleCaptureAudioToggle = useCallback(
(checked: boolean) => {
if (isAppShare) {
VoiceSettingsCommands.update({shareAppAudio: checked, muteStreamAudio: !checked});
} else if (isDeviceShare) {
if (isDeviceShare) {
VoiceSettingsCommands.update({shareDeviceAudio: checked, muteStreamAudio: !checked});
} else if (!shareContextResolved) {
VoiceSettingsCommands.update({
shareAppAudio: checked,
shareDesktopAudio: checked,
muteStreamAudio: !checked,
});
} else if (isAppShare) {
VoiceSettingsCommands.update({shareAppAudio: checked, muteStreamAudio: !checked});
} else {
VoiceSettingsCommands.update({shareDesktopAudio: checked, muteStreamAudio: !checked});
}
runApply({audioSettingsChanged: true});
},
[isAppShare, isDeviceShare, runApply],
[isAppShare, isDeviceShare, shareContextResolved, runApply],
);
const handleHidePreviewToggle = useCallback((checked: boolean) => {
VoiceSettingsCommands.update({hideStreamPreview: checked});
@@ -400,6 +400,7 @@ const VoiceControlBarInner = observer(function VoiceControlBarInner() {
applyToLiveStream={isScreenShareEnabled}
displayShareEnvironment={displayShareEnvironment}
shareContext={ActiveScreenShareSource.getSourceId()?.startsWith('window:') ? 'app' : 'display'}
shareContextResolved={ActiveScreenShareSource.getSourceId() != null}
data-flx="voice.voice-control-bar.render-screen-share-menu.stream-settings-menu-content"
/>
<MenuGroup data-flx="voice.voice-control-bar.render-screen-share-menu.menu-group--2">
@@ -414,10 +414,6 @@ export type ScreenShareEncoderVerificationFailure =
| (StalledVideoEncoderInfo & {reason: 'stalled'})
| MissingExpectedVideoEncoderInfo;
function isSoftwareEncoderStats(implementation: string | null, powerEfficientEncoder: boolean | null): boolean {
return classifyVideoEncoderAcceleration(implementation, powerEfficientEncoder) === 'software';
}
function getStatsKind(report: OutboundVideoStatsEntry, codecs: Map<string, CodecStatsEntry>): string | undefined {
if (report.kind || report.mediaType) return report.kind ?? report.mediaType;
if (!report.codecId) return undefined;
@@ -462,6 +458,7 @@ export function findSoftwareVideoEncoder(stats: RTCStatsReport, codec?: VideoCod
reports.push(report);
}
}
let softwareEncoder: SoftwareVideoEncoderInfo | null = null;
for (const report of reports) {
if (getStatsKind(report, codecs) !== 'video') continue;
if (report.codecId && !codecMatchesTarget(codecs.get(report.codecId)?.mimeType, codec)) continue;
@@ -471,13 +468,16 @@ export function findSoftwareVideoEncoder(stats: RTCStatsReport, codec?: VideoCod
: null;
const powerEfficientEncoder =
typeof report.powerEfficientEncoder === 'boolean' ? report.powerEfficientEncoder : null;
if (!isSoftwareEncoderStats(implementation, powerEfficientEncoder)) continue;
return {
implementation: implementation ?? UNKNOWN_ENCODER_IMPLEMENTATION,
powerEfficientEncoder,
};
const acceleration = classifyVideoEncoderAcceleration(implementation, powerEfficientEncoder);
if (acceleration === 'hardware') return null;
if (acceleration === 'software' && softwareEncoder === null) {
softwareEncoder = {
implementation: implementation ?? UNKNOWN_ENCODER_IMPLEMENTATION,
powerEfficientEncoder,
};
}
}
return null;
return softwareEncoder;
}
export function shouldTriggerSoftwareEncoderWarning(codec: VideoCodec): boolean {
@@ -832,7 +832,7 @@ export async function armNativeAudioForNextCapture(sourceId: string): Promise<bo
platform: electronApi.platform,
sourceId,
backend: availability?.backend ?? null,
reason: availability?.reason ?? 'os-version-too-old',
reason: availability?.reason ?? 'process-scope-unsupported',
detail: availability?.detail ?? null,
});
logger.warn('Cannot arm per-window audio capture: native audio addon unavailable', {
@@ -291,6 +291,14 @@ function getConfiguredScreenShareOptions(
sourceId,
preferredDisplaySurface,
);
if (!includeAudio && shareContext !== 'device' && supportsDesktopScreenShareAudioCapture()) {
logger.info('Screen share audio not requested for this surface', {
sourceId,
preferredDisplaySurface,
shareAppAudio: VoiceSettings.getShareAppAudio(),
shareDesktopAudio: VoiceSettings.getShareDesktopAudio(),
});
}
const {captureOptions, publishOptions} = buildScreenShareOptions({
resolution,
frameRate,
+1
View File
@@ -231,6 +231,7 @@ fn discovery_endpoint(discovery: &DiscoveryResponse, key: &str) -> Option<String
.map(ToOwned::to_owned)
}
#[allow(clippy::result_large_err)]
async fn load_spa_index_html(state: &AppState) -> Result<String, Response> {
if let Some(index_upstream_url) = &state.config.index_upstream_url {
let response = state
+4
View File
@@ -15,4 +15,8 @@ reqwest = { version = "0.13.4", default-features = false, features = ["blocking"
serde_json = "1.0.150"
time = { version = "0.3.47", features = ["formatting", "macros", "parsing"] }
tracing = "0.1.44"
base64 = "0.22"
hmac = "0.13.0"
sha2 = "0.11.0"
thiserror = "2"
urlencoding = "2.1.3"
+276
View File
@@ -0,0 +1,276 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use base64::prelude::*;
use thiserror::Error;
const V2_PREFIX: &str = "v2/";
#[derive(Debug, Error, Eq, PartialEq)]
pub enum ExternalPathError {
#[error("invalid external path")]
InvalidExternalPath,
#[error("invalid base64")]
InvalidBase64,
#[error("invalid utf-8")]
InvalidUtf8,
}
pub fn build_v2_external_media_proxy_path(input_url: &str) -> String {
format!(
"{V2_PREFIX}{}",
BASE64_URL_SAFE_NO_PAD.encode(input_url.as_bytes())
)
}
fn percent_encode_component(value: &str) -> String {
let mut out = String::with_capacity(value.len());
for byte in value.bytes() {
match byte {
b'A'..=b'Z'
| b'a'..=b'z'
| b'0'..=b'9'
| b'-'
| b'_'
| b'.'
| b'!'
| b'~'
| b'*'
| b'\''
| b'('
| b')' => out.push(byte as char),
_ => out.push_str(&format!("%{byte:02X}")),
}
}
out
}
pub fn build_external_media_proxy_path(input_url: &str) -> Result<String, ExternalPathError> {
let (scheme, remainder) = input_url
.split_once("://")
.ok_or(ExternalPathError::InvalidExternalPath)?;
if scheme.is_empty() || remainder.is_empty() {
return Err(ExternalPathError::InvalidExternalPath);
}
let (authority_and_path, query) = match remainder.split_once('?') {
Some((head, tail)) => (head, Some(tail)),
None => (remainder, None),
};
let (host, path) = match authority_and_path.split_once('/') {
Some((host, path)) => (host, path),
None => (authority_and_path, ""),
};
if host.is_empty() {
return Err(ExternalPathError::InvalidExternalPath);
}
let mut segments: Vec<String> = Vec::new();
if let Some(query) = query {
segments.push(percent_encode_component(&format!("?{query}")));
}
segments.push(scheme.to_owned());
segments.push(host.to_owned());
if !path.is_empty() {
segments.push(
path.split('/')
.map(percent_encode_component)
.collect::<Vec<_>>()
.join("/"),
);
}
Ok(segments.join("/"))
}
fn decode_v2(proxy_path: &str) -> Result<String, ExternalPathError> {
let encoded = proxy_path
.strip_prefix(V2_PREFIX)
.ok_or(ExternalPathError::InvalidExternalPath)?;
if encoded.is_empty() {
return Err(ExternalPathError::InvalidExternalPath);
}
let bytes = BASE64_URL_SAFE_NO_PAD
.decode(encoded)
.map_err(|_| ExternalPathError::InvalidBase64)?;
String::from_utf8(bytes).map_err(|_| ExternalPathError::InvalidUtf8)
}
fn hex_val(c: u8) -> Option<u8> {
match c {
b'0'..=b'9' => Some(c - b'0'),
b'a'..=b'f' => Some(c - b'a' + 10),
b'A'..=b'F' => Some(c - b'A' + 10),
_ => None,
}
}
pub fn percent_decode(input: &str, plus_as_space: bool) -> Vec<u8> {
let bytes = input.as_bytes();
let mut out = Vec::with_capacity(bytes.len());
let mut i = 0;
while i < bytes.len() {
let ch = bytes[i];
if ch == b'%'
&& i + 2 < bytes.len()
&& let (Some(hi), Some(lo)) = (hex_val(bytes[i + 1]), hex_val(bytes[i + 2]))
{
out.push((hi << 4) | lo);
i += 3;
continue;
}
out.push(if plus_as_space && ch == b'+' {
b' '
} else {
ch
});
i += 1;
}
out
}
pub fn percent_decode_string(input: &str, plus_as_space: bool) -> String {
String::from_utf8_lossy(&percent_decode(input, plus_as_space)).into_owned()
}
fn legacy_protocol_index(parts: &[&str]) -> Option<usize> {
for (index, part) in parts.iter().enumerate() {
if part.is_empty() {
continue;
}
if index == 0 && part.contains("%3D") {
continue;
}
let mut chars = part.bytes();
let first = chars.next()?;
if !first.is_ascii_alphabetic() {
continue;
}
if chars.all(|c| c.is_ascii_alphanumeric() || matches!(c, b'+' | b'.' | b'-')) {
return Some(index);
}
}
None
}
fn reconstruct_legacy(proxy_path: &str) -> Result<String, ExternalPathError> {
let parts: Vec<&str> = proxy_path.split('/').collect();
let protocol_index =
legacy_protocol_index(&parts).ok_or(ExternalPathError::InvalidExternalPath)?;
if protocol_index + 1 >= parts.len() {
return Err(ExternalPathError::InvalidExternalPath);
}
let protocol = parts[protocol_index];
let host_port = percent_decode_string(parts[protocol_index + 1], false);
let query_raw = parts[..protocol_index].join("/");
let path_raw = parts[protocol_index + 2..].join("/");
let query = percent_decode_string(&query_raw, false);
let path = percent_decode_string(&path_raw, false);
let normalized_query = query.strip_prefix('?').unwrap_or(&query);
Ok(format!(
"{protocol}://{host_port}/{path}{}{normalized_query}",
if normalized_query.is_empty() { "" } else { "?" }
))
}
pub fn reconstruct_original_url(proxy_path: &str) -> Result<String, ExternalPathError> {
if proxy_path.starts_with(V2_PREFIX) {
decode_v2(proxy_path)
} else {
reconstruct_legacy(proxy_path)
}
}
type HmacSha256 = hmac::Hmac<sha2::Sha256>;
pub fn create_signature(input: &str, secret: &[u8]) -> String {
use hmac::{KeyInit, Mac};
let mut mac = HmacSha256::new_from_slice(secret).expect("HMAC accepts any key length");
mac.update(input.as_bytes());
BASE64_URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes())
}
pub fn build_external_media_proxy_url(
public_endpoint: &str,
input_url: &str,
secret: &[u8],
) -> Option<String> {
let path = build_external_media_proxy_path(input_url).ok()?;
let signature = create_signature(&path, secret);
Some(format!("{public_endpoint}/external/{signature}/{path}"))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builds_the_plain_path_shape() {
assert_eq!(
"https/static.klipy.com/ii/c8/28/HkAKKCzZ.webp",
build_external_media_proxy_path("https://static.klipy.com/ii/c8/28/HkAKKCzZ.webp")
.unwrap()
);
}
#[test]
fn builds_an_encoded_query_ahead_of_the_scheme() {
assert_eq!(
"%3Fv%3Dquery_param%26goes%3Dhere/https/static.klipy.com/ii/HkAKKCzZ.webp",
build_external_media_proxy_path(
"https://static.klipy.com/ii/HkAKKCzZ.webp?v=query_param&goes=here"
)
.unwrap()
);
}
#[test]
fn keeps_a_non_default_port() {
assert_eq!(
"https/example.com:8443/a.png",
build_external_media_proxy_path("https://example.com:8443/a.png").unwrap()
);
}
#[test]
fn round_trips_every_shape() {
for url in [
"https://static.klipy.com/ii/c8/28/HkAKKCzZ.webp",
"https://static.klipy.com/ii/HkAKKCzZ.webp?v=query_param&goes=here",
"https://example.com:8443/a.png",
"https://avatars.githubusercontent.com/u/241303489?v=4",
"http://example.com/plain.gif",
"https://example.com/file.png?a=1&b=2&c=3",
] {
let path = build_external_media_proxy_path(url).unwrap();
assert_eq!(
url,
reconstruct_original_url(&path).unwrap(),
"round trip {url}"
);
}
}
#[test]
fn does_not_double_the_question_mark() {
let decoded = reconstruct_original_url("%3Fa%3D1/https/example.com/x.png").unwrap();
assert_eq!("https://example.com/x.png?a=1", decoded);
assert!(!decoded.contains("??"));
}
#[test]
fn still_accepts_a_query_without_a_leading_question_mark() {
assert_eq!(
"https://example.com/x.png?a=1",
reconstruct_original_url("a%3D1/https/example.com/x.png").unwrap()
);
}
#[test]
fn rejects_a_url_without_a_scheme() {
assert!(build_external_media_proxy_path("example.com/a.png").is_err());
}
#[test]
fn v2_path_roundtrip() {
let path = build_v2_external_media_proxy_path("https://example.com/a b.png?x=1");
let decoded = reconstruct_original_url(&path).unwrap();
assert_eq!("https://example.com/a b.png?x=1", decoded);
}
}
+1
View File
@@ -1,4 +1,5 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
pub mod config;
pub mod external_media_path;
pub mod geoip;
@@ -1071,6 +1071,7 @@ mod platform {
const BACKSLASH_UTF16: &[u16] = &[b'\\' as u16, 0];
const PROCESS_LOOPBACK_PROBE_TIMEOUT_MS: u32 = 1_500;
const PROCESS_LOOPBACK_PROBE_COUNT: u64 = 2;
const SESSION_MIXER_REFRESH_INTERVAL: Duration = Duration::from_millis(1_000);
const SESSION_MIXER_WAIT_TIMEOUT_MS: u32 = 250;
const MAX_SESSION_MIXER_CAPTURES: usize = 48;
@@ -2197,7 +2198,7 @@ mod platform {
});
match rx.recv_timeout(std::time::Duration::from_millis(
u64::from(PROCESS_LOOPBACK_PROBE_TIMEOUT_MS) + 500,
u64::from(PROCESS_LOOPBACK_PROBE_TIMEOUT_MS) * PROCESS_LOOPBACK_PROBE_COUNT + 500,
)) {
Ok((include_result, exclude_result)) => ProcessLoopbackRuntimeProbe {
include_supported: include_result.is_ok(),
@@ -17,6 +17,9 @@
-define(SESSION_CONNECT_MAX_WORKERS, 8).
-define(SESSION_CONNECT_DEFAULT_MAX_QUEUE, 1024).
-define(CONNECT_SNAPSHOT_HEAVY_MEMBER_KEYS, [
<<"members">>, members_normalized, <<"member_role_index">>, members_sorted_ids
]).
-spec ensure_session_connect_queue(term()) -> queue:queue().
ensure_session_connect_queue(Value) when is_list(Value) ->
@@ -363,14 +366,14 @@ queued_session_pid(_) ->
-spec start_worker(map(), map()) -> map().
start_worker(Item, State) ->
Self = self(),
Snapshot = build_connect_snapshot(State),
Snapshot = build_connect_snapshot(Item, State),
{_Pid, Ref} = spawn_monitor(fun() -> compute_and_send_done(Item, Self, Snapshot) end),
WorkerRefs = maps:get(session_connect_worker_refs, State, #{}),
State#{session_connect_worker_refs => WorkerRefs#{Ref => true}}.
-spec build_connect_snapshot(map()) -> map().
build_connect_snapshot(State) ->
maps:with(
-spec build_connect_snapshot(map(), map()) -> map().
build_connect_snapshot(Item, State) ->
Base = maps:with(
[
id,
data,
@@ -382,8 +385,76 @@ build_connect_snapshot(State) ->
virtual_channel_access
],
State
),
maybe_trim_connect_snapshot(Item, Base, State).
-spec maybe_trim_connect_snapshot(map(), map(), map()) -> map().
maybe_trim_connect_snapshot(Item, Base, State) ->
case has_members_ets(State) of
true -> trim_connect_snapshot(Item, Base);
false -> Base
end.
-spec has_members_ets(map()) -> boolean().
has_members_ets(#{data := #{members_ets := Tab}}) -> is_reference(Tab);
has_members_ets(_) -> false.
-spec trim_connect_snapshot(map(), map()) -> map().
trim_connect_snapshot(Item, #{data := Data} = Base) when is_map(Data) ->
Retained = retained_member_map(Item, Base, Data),
Trimmed = maps:without(?CONNECT_SNAPSHOT_HEAVY_MEMBER_KEYS, Data),
Base#{
data => Trimmed#{
<<"members">> => Retained,
members_normalized => Retained,
<<"member_role_index">> =>
guild_data_index_members:build_member_role_index(Retained)
}
};
trim_connect_snapshot(_Item, Base) ->
Base.
-spec retained_member_map(map(), map(), map()) -> #{integer() => map()}.
retained_member_map(Item, Base, Data) ->
UserIds = [connect_user_id(Item) | voice_state_user_ids(Base)],
lists:foldl(
fun(UserId, Acc) -> retain_member(UserId, Data, Acc) end,
#{},
UserIds
).
-spec retain_member(term(), map(), #{integer() => map()}) -> #{integer() => map()}.
retain_member(UserId, Data, Acc) when is_integer(UserId) ->
case guild_data_index_members:get_member_ets(UserId, Data) of
Member when is_map(Member) -> Acc#{UserId => Member};
_ -> Acc
end;
retain_member(_UserId, _Data, Acc) ->
Acc.
-spec connect_user_id(map()) -> integer() | undefined.
connect_user_id(Item) ->
Request = maps:get(request, Item, #{}),
case maps:get(user_id, Request, undefined) of
UserId when is_integer(UserId) -> UserId;
_ -> undefined
end.
-spec voice_state_user_ids(map()) -> [integer()].
voice_state_user_ids(Base) ->
maps:fold(
fun(_Key, VoiceState, Acc) -> add_voice_state_user_id(VoiceState, Acc) end,
[],
voice_state_utils:voice_states(Base)
).
-spec add_voice_state_user_id(term(), [integer()]) -> [integer()].
add_voice_state_user_id(VoiceState, Acc) ->
case voice_state_utils:voice_state_user_id(VoiceState) of
UserId when is_integer(UserId) -> [UserId | Acc];
_ -> Acc
end.
-spec compute_and_send_done(map(), pid(), map()) -> ok.
compute_and_send_done(Item, GuildPid, Snapshot) ->
Request = maps:get(request, Item, #{}),
+27 -3
View File
@@ -78,7 +78,6 @@ get_guild_state(UserId, State) ->
Data = guild_data_index:ensure_data_map(State),
GuildId = guild_id(State),
AllChannels = guild_data_channels:channels_from_data(Data),
AllMembers = guild_data_index:member_values(Data),
Member = guild_data_members:find_member_by_user_id(UserId, State),
{ViewableChannels, JoinedAt} = guild_data_channels:derive_member_view(
UserId, Member, State, AllChannels
@@ -89,9 +88,9 @@ get_guild_state(UserId, State) ->
AllVoiceStates = guild_voice:get_voice_states_list(StateWithVoice),
ViewableChannelIds = channel_id_set(ViewableChannels),
VoiceStates = filter_voice_states(AllVoiceStates, ViewableChannelIds),
VoiceMembers = guild_data_channels:voice_members_from_states(VoiceStates, AllMembers),
VoiceMembers = resolve_voice_members(VoiceStates, Data),
Members = guild_data_channels:merge_members(OwnMemberList, VoiceMembers),
MemberCount = maps:get(member_count, State, length(AllMembers)),
MemberCount = maps:get(member_count, State, guild_data_index:member_count(Data)),
build_guild_state_map(
GuildId,
Data,
@@ -103,6 +102,31 @@ get_guild_state(UserId, State) ->
JoinedAt
).
-spec resolve_voice_members([map()], map()) -> [map()].
resolve_voice_members(VoiceStates, Data) ->
lists:filtermap(fun(VS) -> resolve_voice_member_entry(VS, Data) end, VoiceStates).
-spec resolve_voice_member_entry(map(), map()) -> {true, map()} | false.
resolve_voice_member_entry(VoiceState, Data) ->
case maps:get(<<"member">>, VoiceState, undefined) of
Member when is_map(Member), map_size(Member) > 0 ->
{true, Member};
_ ->
resolve_voice_member_by_lookup(VoiceState, Data)
end.
-spec resolve_voice_member_by_lookup(map(), map()) -> {true, map()} | false.
resolve_voice_member_by_lookup(VoiceState, Data) ->
case voice_state_utils:voice_state_user_id(VoiceState) of
UserId when is_integer(UserId) ->
case guild_data_index_members:get_member_ets(UserId, Data) of
Member when is_map(Member) -> {true, Member};
_ -> false
end;
_ ->
false
end.
-spec fetch_latest_voice_states(guild_state()) -> guild_state().
fetch_latest_voice_states(State) ->
case maps:get(voice_server_pid, State, undefined) of
+1
View File
@@ -7,6 +7,7 @@ edition.workspace = true
license.workspace = true
[dependencies]
fluxer_common = { path = "../fluxer_common" }
anyhow = "1.0.102"
base64 = "0.22.1"
fluxer-svc = { path = "../fluxer_svc", default-features = false }
+6 -25
View File
@@ -1,12 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use base64::prelude::*;
use hmac::{Hmac, KeyInit, Mac};
use sha2::Sha256;
use url::Url;
const V2_PATH_PREFIX: &str = "v2/";
#[derive(Clone)]
pub struct MediaProxyUrlBuilder {
endpoint: String,
@@ -66,28 +61,14 @@ impl MediaProxyUrlBuilder {
return Some(input_url.to_owned());
}
let proxy_path = build_external_media_proxy_path(parsed.as_str());
let signature = create_signature(&proxy_path, &self.secret_key);
Some(format!(
"{}/external/{signature}/{proxy_path}",
self.endpoint
))
fluxer_common::external_media_path::build_external_media_proxy_url(
&self.endpoint,
parsed.as_str(),
self.secret_key.as_bytes(),
)
}
}
fn build_external_media_proxy_path(input_url: &str) -> String {
format!(
"{V2_PATH_PREFIX}{}",
BASE64_URL_SAFE_NO_PAD.encode(input_url)
)
}
fn create_signature(input: &str, secret: &str) -> String {
let mut mac = Hmac::<Sha256>::new_from_slice(secret.as_bytes()).expect("HMAC accepts any key");
mac.update(input.as_bytes());
BASE64_URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes())
}
#[cfg(test)]
mod tests {
use super::*;
@@ -153,7 +134,7 @@ mod tests {
.expect("proxy url");
assert!(url.starts_with("https://media.example.test/external/"));
assert!(url.contains("/v2/"));
assert!(url.contains("/https/"));
assert_eq!(
builder.external_proxy_url("https://media.example.test/external/existing"),
Some("https://media.example.test/external/existing".to_owned())
+1
View File
@@ -18,6 +18,7 @@ name = "core"
harness = false
[dependencies]
fluxer_common = { path = "../fluxer_common" }
anyhow = "1.0.102"
axum = {version = "0.8.9", default-features = false, features = ["http1", "json", "query", "tokio"]}
base64 = "0.22.1"
+4 -255
View File
@@ -1,257 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use base64::prelude::*;
use thiserror::Error;
const V2_PREFIX: &str = "v2/";
#[derive(Debug, Error, Eq, PartialEq)]
pub enum ExternalPathError {
#[error("invalid external path")]
InvalidExternalPath,
#[error("invalid base64")]
InvalidBase64,
#[error("invalid utf-8")]
InvalidUtf8,
}
pub fn build_v2_external_media_proxy_path(input_url: &str) -> String {
format!(
"{V2_PREFIX}{}",
BASE64_URL_SAFE_NO_PAD.encode(input_url.as_bytes())
)
}
fn percent_encode_component(value: &str) -> String {
let mut out = String::with_capacity(value.len());
for byte in value.bytes() {
match byte {
b'A'..=b'Z'
| b'a'..=b'z'
| b'0'..=b'9'
| b'-'
| b'_'
| b'.'
| b'!'
| b'~'
| b'*'
| b'\''
| b'('
| b')' => out.push(byte as char),
_ => out.push_str(&format!("%{byte:02X}")),
}
}
out
}
pub fn build_external_media_proxy_path(input_url: &str) -> Result<String, ExternalPathError> {
let (scheme, remainder) = input_url
.split_once("://")
.ok_or(ExternalPathError::InvalidExternalPath)?;
if scheme.is_empty() || remainder.is_empty() {
return Err(ExternalPathError::InvalidExternalPath);
}
let (authority_and_path, query) = match remainder.split_once('?') {
Some((head, tail)) => (head, Some(tail)),
None => (remainder, None),
};
let (host, path) = match authority_and_path.split_once('/') {
Some((host, path)) => (host, path),
None => (authority_and_path, ""),
};
if host.is_empty() {
return Err(ExternalPathError::InvalidExternalPath);
}
let mut segments: Vec<String> = Vec::new();
if let Some(query) = query {
segments.push(percent_encode_component(&format!("?{query}")));
}
segments.push(scheme.to_owned());
segments.push(host.to_owned());
if !path.is_empty() {
segments.push(
path.split('/')
.map(percent_encode_component)
.collect::<Vec<_>>()
.join("/"),
);
}
Ok(segments.join("/"))
}
fn decode_v2(proxy_path: &str) -> Result<String, ExternalPathError> {
let encoded = proxy_path
.strip_prefix(V2_PREFIX)
.ok_or(ExternalPathError::InvalidExternalPath)?;
if encoded.is_empty() {
return Err(ExternalPathError::InvalidExternalPath);
}
let bytes = BASE64_URL_SAFE_NO_PAD
.decode(encoded)
.map_err(|_| ExternalPathError::InvalidBase64)?;
String::from_utf8(bytes).map_err(|_| ExternalPathError::InvalidUtf8)
}
fn hex_val(c: u8) -> Option<u8> {
match c {
b'0'..=b'9' => Some(c - b'0'),
b'a'..=b'f' => Some(c - b'a' + 10),
b'A'..=b'F' => Some(c - b'A' + 10),
_ => None,
}
}
pub fn percent_decode(input: &str, plus_as_space: bool) -> Vec<u8> {
let bytes = input.as_bytes();
let mut out = Vec::with_capacity(bytes.len());
let mut i = 0;
while i < bytes.len() {
let ch = bytes[i];
if ch == b'%'
&& i + 2 < bytes.len()
&& let (Some(hi), Some(lo)) = (hex_val(bytes[i + 1]), hex_val(bytes[i + 2]))
{
out.push((hi << 4) | lo);
i += 3;
continue;
}
out.push(if plus_as_space && ch == b'+' {
b' '
} else {
ch
});
i += 1;
}
out
}
pub fn percent_decode_string(input: &str, plus_as_space: bool) -> String {
String::from_utf8_lossy(&percent_decode(input, plus_as_space)).into_owned()
}
fn legacy_protocol_index(parts: &[&str]) -> Option<usize> {
for (index, part) in parts.iter().enumerate() {
if part.is_empty() {
continue;
}
if index == 0 && part.contains("%3D") {
continue;
}
let mut chars = part.bytes();
let first = chars.next()?;
if !first.is_ascii_alphabetic() {
continue;
}
if chars.all(|c| c.is_ascii_alphanumeric() || matches!(c, b'+' | b'.' | b'-')) {
return Some(index);
}
}
None
}
fn reconstruct_legacy(proxy_path: &str) -> Result<String, ExternalPathError> {
let parts: Vec<&str> = proxy_path.split('/').collect();
let protocol_index =
legacy_protocol_index(&parts).ok_or(ExternalPathError::InvalidExternalPath)?;
if protocol_index + 1 >= parts.len() {
return Err(ExternalPathError::InvalidExternalPath);
}
let protocol = parts[protocol_index];
let host_port = percent_decode_string(parts[protocol_index + 1], false);
let query_raw = parts[..protocol_index].join("/");
let path_raw = parts[protocol_index + 2..].join("/");
let query = percent_decode_string(&query_raw, false);
let path = percent_decode_string(&path_raw, false);
let normalized_query = query.strip_prefix('?').unwrap_or(&query);
Ok(format!(
"{protocol}://{host_port}/{path}{}{normalized_query}",
if normalized_query.is_empty() { "" } else { "?" }
))
}
pub fn reconstruct_original_url(proxy_path: &str) -> Result<String, ExternalPathError> {
if proxy_path.starts_with(V2_PREFIX) {
decode_v2(proxy_path)
} else {
reconstruct_legacy(proxy_path)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builds_the_plain_path_shape() {
assert_eq!(
"https/static.klipy.com/ii/c8/28/HkAKKCzZ.webp",
build_external_media_proxy_path("https://static.klipy.com/ii/c8/28/HkAKKCzZ.webp")
.unwrap()
);
}
#[test]
fn builds_an_encoded_query_ahead_of_the_scheme() {
assert_eq!(
"%3Fv%3Dquery_param%26goes%3Dhere/https/static.klipy.com/ii/HkAKKCzZ.webp",
build_external_media_proxy_path(
"https://static.klipy.com/ii/HkAKKCzZ.webp?v=query_param&goes=here"
)
.unwrap()
);
}
#[test]
fn keeps_a_non_default_port() {
assert_eq!(
"https/example.com:8443/a.png",
build_external_media_proxy_path("https://example.com:8443/a.png").unwrap()
);
}
#[test]
fn round_trips_every_shape() {
for url in [
"https://static.klipy.com/ii/c8/28/HkAKKCzZ.webp",
"https://static.klipy.com/ii/HkAKKCzZ.webp?v=query_param&goes=here",
"https://example.com:8443/a.png",
"https://avatars.githubusercontent.com/u/241303489?v=4",
"http://example.com/plain.gif",
"https://example.com/file.png?a=1&b=2&c=3",
] {
let path = build_external_media_proxy_path(url).unwrap();
assert_eq!(
url,
reconstruct_original_url(&path).unwrap(),
"round trip {url}"
);
}
}
#[test]
fn does_not_double_the_question_mark() {
let decoded = reconstruct_original_url("%3Fa%3D1/https/example.com/x.png").unwrap();
assert_eq!("https://example.com/x.png?a=1", decoded);
assert!(!decoded.contains("??"));
}
#[test]
fn still_accepts_a_query_without_a_leading_question_mark() {
assert_eq!(
"https://example.com/x.png?a=1",
reconstruct_original_url("a%3D1/https/example.com/x.png").unwrap()
);
}
#[test]
fn rejects_a_url_without_a_scheme() {
assert!(build_external_media_proxy_path("example.com/a.png").is_err());
}
#[test]
fn v2_path_roundtrip() {
let path = build_v2_external_media_proxy_path("https://example.com/a b.png?x=1");
let decoded = reconstruct_original_url(&path).unwrap();
assert_eq!("https://example.com/a b.png?x=1", decoded);
}
}
pub use fluxer_common::external_media_path::{
ExternalPathError, build_external_media_proxy_path, build_v2_external_media_proxy_path,
percent_decode, percent_decode_string, reconstruct_original_url,
};
+1
View File
@@ -335,6 +335,7 @@ fn replace_image_extension(filename: &str, ext: AssetExtension) -> String {
}
}
#[allow(clippy::result_large_err)]
async fn rasterize_metadata_svg(
app: &Arc<AppState>,
input: InputData,
+1
View File
@@ -5,6 +5,7 @@ edition.workspace = true
license.workspace = true
[dependencies]
fluxer_common = { path = "../fluxer_common" }
anyhow = "1.0.102"
base64 = "0.22.1"
chrono = { version = "0.4.45", default-features = false, features = ["serde"] }
+1 -1
View File
@@ -6,7 +6,7 @@ WORKDIR /usr/src/app
COPY . .
RUN cargo build --release -p fluxer-messages
RUN cargo build --release -p fluxer-messages --features scylla
FROM debian:bookworm-slim
+4 -18
View File
@@ -15,13 +15,11 @@ use crate::types::{
MessageSnapshot, MessageStickerItem,
};
use crate::udt;
use base64::Engine;
use chrono::{DateTime, Utc};
use fluxer_svc::shard::ShardService;
use fluxer_svc::transport::NatsTransport;
use fluxer_svc::{postgres, postgres::BigIntBound, postgres::KeyPart};
use futures::stream::{self, StreamExt};
use hmac::{Hmac, KeyInit, Mac};
#[cfg(feature = "scylla")]
use scylla::DeserializeRow;
#[cfg(feature = "scylla")]
@@ -31,7 +29,6 @@ use scylla::response::query_result::QueryRowsResult;
#[cfg(feature = "scylla")]
use scylla::statement::prepared::PreparedStatement;
use serde::Deserialize;
use sha2::Sha256;
use std::collections::{HashMap, HashSet};
#[cfg(feature = "scylla")]
use std::sync::Arc;
@@ -2796,23 +2793,12 @@ fn external_media_proxy_url(input_url: &str, options: &ResponseBuildOptions) ->
Ok(url) => url,
Err(_) => return input_url.to_owned(),
};
let encoded =
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(parsed_url.as_str().as_bytes());
let path = format!("v2/{encoded}");
let signature = create_signature(&path, &options.media_proxy_secret_key);
format!(
"{}/external/{}/{}",
fluxer_common::external_media_path::build_external_media_proxy_url(
options.media_endpoint.trim_end_matches('/'),
signature,
path
parsed_url.as_str(),
options.media_proxy_secret_key.as_bytes(),
)
}
fn create_signature(input: &str, secret: &str) -> String {
let mut mac =
Hmac::<Sha256>::new_from_slice(secret.as_bytes()).expect("hmac accepts any key size");
mac.update(input.as_bytes());
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes())
.unwrap_or_else(|| input_url.to_owned())
}
fn dt_to_epoch_millis(dt: &DateTime<Utc>) -> i64 {
+1
View File
@@ -5,6 +5,7 @@ edition.workspace = true
license.workspace = true
[dependencies]
fluxer_common = { path = "../fluxer_common" }
ammonia = "4.1.2"
anyhow = "1.0.102"
base64 = "0.22.1"
+15 -29
View File
@@ -5,7 +5,6 @@ use std::time::Duration;
use url::Url;
const METADATA_TIMEOUT: Duration = Duration::from_secs(5);
const EXTERNAL_PROXY_VERSION_PREFIX: &str = "v2/";
pub struct MediaProxyClient {
http_client: reqwest::Client,
@@ -120,30 +119,14 @@ impl MediaProxyClient {
return Some(input_url.to_owned());
}
let parsed = Url::parse(input_url).ok()?;
let encoded = base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
fluxer_common::external_media_path::build_external_media_proxy_url(
&self.public_endpoint,
parsed.as_str(),
);
let path = format!("{EXTERNAL_PROXY_VERSION_PREFIX}{encoded}");
let signature = create_signature(&path, &self.secret_key);
Some(format!(
"{}/external/{}/{}",
self.public_endpoint, signature, path
))
self.secret_key.as_bytes(),
)
}
}
fn create_signature(input: &str, secret: &str) -> String {
use hmac::{Hmac, KeyInit, Mac};
let mut mac =
Hmac::<sha2::Sha256>::new_from_slice(secret.as_bytes()).expect("hmac accepts any key size");
mac.update(input.as_bytes());
base64::Engine::encode(
&base64::engine::general_purpose::URL_SAFE_NO_PAD,
mac.finalize().into_bytes(),
)
}
#[cfg(test)]
mod tests {
use super::*;
@@ -178,7 +161,7 @@ mod tests {
}
#[test]
fn external_proxy_url_uses_public_endpoint_and_v2_path() {
fn external_proxy_url_uses_public_endpoint_and_plain_path() {
let client = reqwest::Client::new();
let mp = MediaProxyClient::new_with_public_endpoint(
"http://media-proxy:8080/",
@@ -191,15 +174,18 @@ mod tests {
.expect("proxy url");
assert!(proxy.starts_with("https://media.example.test/external/"));
assert!(proxy.contains("/v2/"));
assert!(proxy.contains("/https/"));
let encoded = proxy.rsplit('/').next().expect("encoded path segment");
let decoded =
base64::Engine::decode(&base64::engine::general_purpose::URL_SAFE_NO_PAD, encoded)
.expect("valid base64");
let path = proxy
.split_once("/external/")
.expect("external segment")
.1
.split_once('/')
.expect("signature segment")
.1;
assert_eq!(
String::from_utf8(decoded).expect("utf8"),
"https://pbs.twimg.com/media/a.jpg?name=orig"
"https://pbs.twimg.com/media/a.jpg?name=orig",
fluxer_common::external_media_path::reconstruct_original_url(path).expect("decodes")
);
assert_eq!(
mp.external_proxy_url("https://media.example.test/external/already"),
+1 -1
View File
@@ -6,7 +6,7 @@ WORKDIR /usr/src/app
COPY . .
RUN cargo build --release -p fluxer-users
RUN cargo build --release -p fluxer-users --features scylla
FROM debian:bookworm-slim
+16
View File
@@ -79,6 +79,14 @@ export interface MasterConfig {
static: string;
};
};
s3_downloads?: {
endpoint: string;
presigned_url_base?: string;
force_path_style?: boolean;
region?: string;
access_key_id?: string;
secret_access_key?: string;
};
services: {
api: {
port: number;
@@ -104,6 +112,14 @@ export interface MasterConfig {
task?: string;
enable_cron_scheduler?: boolean;
enable_voice_reconciliation?: boolean;
voice_reconciliation?: {
interval_ms?: number;
stagger_delay_ms?: number;
lock_ttl_seconds?: number;
cadence_ttl_seconds?: number;
gateway_only_grace_ms?: number;
livekit_only_grace_ms?: number;
};
lane_concurrency_overrides?: {
realtime?: number;
unfurl?: number;
@@ -0,0 +1,47 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {MasterConfig} from '@fluxer/config/src/MasterConfig';
export interface S3ProviderSettings {
endpoint: string;
presignedUrlBase?: string;
forcePathStyle: boolean;
region: string;
accessKeyId: string;
secretAccessKey: string;
}
export interface ResolvedDownloadsProvider {
settings: S3ProviderSettings;
isOverridden: boolean;
}
export function resolveDownloadsProvider(master: Pick<MasterConfig, 's3' | 's3_downloads'>): ResolvedDownloadsProvider {
const base = master.s3;
if (!base) {
throw new Error('S3 configuration is required to resolve the downloads provider');
}
const baseSettings: S3ProviderSettings = {
endpoint: base.endpoint,
presignedUrlBase: base.presigned_url_base,
forcePathStyle: base.force_path_style,
region: base.region,
accessKeyId: base.access_key_id,
secretAccessKey: base.secret_access_key,
};
const override = master.s3_downloads;
if (!override?.endpoint) {
return {settings: baseSettings, isOverridden: false};
}
return {
settings: {
endpoint: override.endpoint,
presignedUrlBase: override.presigned_url_base ?? undefined,
forcePathStyle: override.force_path_style ?? baseSettings.forcePathStyle,
region: override.region ?? baseSettings.region,
accessKeyId: override.access_key_id ?? baseSettings.accessKeyId,
secretAccessKey: override.secret_access_key ?? baseSettings.secretAccessKey,
},
isOverridden: true,
};
}
@@ -0,0 +1,63 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {MasterConfig} from '@fluxer/config/src/MasterConfig';
import {resolveDownloadsProvider} from '@fluxer/config/src/S3DownloadsProvider';
import {describe, expect, test} from 'vitest';
const base: Pick<MasterConfig, 's3' | 's3_downloads'>['s3'] = {
endpoint: 'https://main.example.com',
presigned_url_base: 'https://public.example.com',
force_path_style: true,
region: 'us-east-1',
access_key_id: 'MAIN_KEY',
secret_access_key: 'MAIN_SECRET',
buckets: {cdn: 'cdn', uploads: 'uploads', downloads: 'downloads', reports: 'r', harvests: 'h', static: 's'},
};
describe('resolveDownloadsProvider', () => {
test('falls back to the main provider when no override is set', () => {
const resolved = resolveDownloadsProvider({s3: base});
expect(resolved.isOverridden).toBe(false);
expect(resolved.settings.endpoint).toBe('https://main.example.com');
expect(resolved.settings.accessKeyId).toBe('MAIN_KEY');
expect(resolved.settings.secretAccessKey).toBe('MAIN_SECRET');
expect(resolved.settings.region).toBe('us-east-1');
expect(resolved.settings.presignedUrlBase).toBe('https://public.example.com');
});
test('ignores a partial override that does not set an endpoint', () => {
const resolved = resolveDownloadsProvider({s3: base, s3_downloads: {endpoint: '', access_key_id: 'OTHER'}});
expect(resolved.isOverridden).toBe(false);
expect(resolved.settings.accessKeyId).toBe('MAIN_KEY');
});
test('uses the override provider when an endpoint is set', () => {
const resolved = resolveDownloadsProvider({
s3: base,
s3_downloads: {
endpoint: 'https://downloads.example.net',
region: 'eu-central-1',
access_key_id: 'DL_KEY',
secret_access_key: 'DL_SECRET',
},
});
expect(resolved.isOverridden).toBe(true);
expect(resolved.settings.endpoint).toBe('https://downloads.example.net');
expect(resolved.settings.region).toBe('eu-central-1');
expect(resolved.settings.accessKeyId).toBe('DL_KEY');
});
test('inherits unspecified fields from the main provider', () => {
const resolved = resolveDownloadsProvider({s3: base, s3_downloads: {endpoint: 'https://downloads.example.net'}});
expect(resolved.isOverridden).toBe(true);
expect(resolved.settings.endpoint).toBe('https://downloads.example.net');
expect(resolved.settings.region).toBe('us-east-1');
expect(resolved.settings.accessKeyId).toBe('MAIN_KEY');
expect(resolved.settings.forcePathStyle).toBe(true);
});
test('does not inherit the main presigned base for an override provider', () => {
const resolved = resolveDownloadsProvider({s3: base, s3_downloads: {endpoint: 'https://downloads.example.net'}});
expect(resolved.settings.presignedUrlBase).toBeUndefined();
});
});
@@ -65,6 +65,12 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record<string, NamedEnvOverride> = {
FLUXER_S3_BUCKET_REPORTS: {path: ['s3', 'buckets', 'reports']},
FLUXER_S3_BUCKET_HARVESTS: {path: ['s3', 'buckets', 'harvests']},
FLUXER_S3_BUCKET_STATIC: {path: ['s3', 'buckets', 'static']},
FLUXER_S3_DOWNLOADS_ENDPOINT: {path: ['s3_downloads', 'endpoint']},
FLUXER_S3_DOWNLOADS_PUBLIC_ENDPOINT: {path: ['s3_downloads', 'presigned_url_base']},
FLUXER_S3_DOWNLOADS_FORCE_PATH_STYLE: {path: ['s3_downloads', 'force_path_style'], parse: parseEnvValue},
FLUXER_S3_DOWNLOADS_REGION: {path: ['s3_downloads', 'region']},
FLUXER_S3_DOWNLOADS_ACCESS_KEY_ID: {path: ['s3_downloads', 'access_key_id']},
FLUXER_S3_DOWNLOADS_SECRET_ACCESS_KEY: {path: ['s3_downloads', 'secret_access_key']},
FLUXER_NATS_URL: {path: ['services', 'nats', 'core_url']},
FLUXER_NATS_CORE_URL: {path: ['services', 'nats', 'core_url']},
FLUXER_NATS_JETSTREAM_URL: {path: ['services', 'nats', 'jetstream_url']},
@@ -90,6 +96,30 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record<string, NamedEnvOverride> = {
path: ['services', 'api', 'worker', 'enable_voice_reconciliation'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_VOICE_RECONCILIATION_INTERVAL_MS: {
path: ['services', 'api', 'worker', 'voice_reconciliation', 'interval_ms'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_VOICE_RECONCILIATION_STAGGER_DELAY_MS: {
path: ['services', 'api', 'worker', 'voice_reconciliation', 'stagger_delay_ms'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_VOICE_RECONCILIATION_LOCK_TTL_SECONDS: {
path: ['services', 'api', 'worker', 'voice_reconciliation', 'lock_ttl_seconds'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_VOICE_RECONCILIATION_CADENCE_TTL_SECONDS: {
path: ['services', 'api', 'worker', 'voice_reconciliation', 'cadence_ttl_seconds'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_VOICE_RECONCILIATION_GATEWAY_ONLY_GRACE_MS: {
path: ['services', 'api', 'worker', 'voice_reconciliation', 'gateway_only_grace_ms'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_VOICE_RECONCILIATION_LIVEKIT_ONLY_GRACE_MS: {
path: ['services', 'api', 'worker', 'voice_reconciliation', 'livekit_only_grace_ms'],
parse: parseEnvValue,
},
FLUXER_API_WORKER_LANE_CONCURRENCY_OVERRIDES: {
path: ['services', 'api', 'worker', 'lane_concurrency_overrides'],
parse: parseEnvValue,
@@ -0,0 +1,46 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {SlowmodeRateLimitError} from '@fluxer/errors/src/domains/core/SlowmodeRateLimitError';
import {describe, expect, it} from 'vitest';
interface SlowmodeResponseBody {
code: string;
retry_after: number;
}
async function readSlowmodeResponse(error: SlowmodeRateLimitError): Promise<{
status: number;
header: string | null;
body: SlowmodeResponseBody;
}> {
const response = error.getResponse();
const body = (await response.json()) as SlowmodeResponseBody;
return {status: response.status, header: response.headers.get('Retry-After'), body};
}
describe('SlowmodeRateLimitError', () => {
it('reports Retry-After in whole seconds and the body in decimal seconds', async () => {
const {status, header, body} = await readSlowmodeResponse(
new SlowmodeRateLimitError({retryAfter: 5, retryAfterDecimal: 4.7}),
);
expect(status).toBe(400);
expect(body.code).toBe('SLOWMODE_RATE_LIMITED');
expect(body.retry_after).toBe(4.7);
expect(header).toBe('5');
});
it('keeps the header and the body within one second of each other', async () => {
const {header, body} = await readSlowmodeResponse(
new SlowmodeRateLimitError({retryAfter: 5, retryAfterDecimal: 4.997}),
);
const headerSeconds = Number(header);
expect(headerSeconds - body.retry_after).toBeLessThan(1);
expect(headerSeconds).toBeGreaterThanOrEqual(body.retry_after);
});
it('falls back to one second when the caller has no retry window', async () => {
const {header, body} = await readSlowmodeResponse(new SlowmodeRateLimitError({retryAfter: undefined}));
expect(header).toBe('1');
expect(body.retry_after).toBe(1);
});
});
@@ -1,25 +1,15 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {APIErrorCodes} from '@fluxer/constants/src/ApiErrorCodes';
import {
sanitizeRetryAfterDecimalSeconds,
sanitizeRetryAfterSeconds,
} from '@fluxer/errors/src/domains/core/RetryAfterSeconds';
import {ThrottledError} from '@fluxer/errors/src/domains/core/ThrottledError';
import type {FluxerErrorData} from '@fluxer/errors/src/FluxerError';
type RateLimitScope = 'global' | 'shared' | 'user';
function sanitizeRetryAfter(value: number | undefined | null): number {
if (value == null || !Number.isFinite(value) || value < 0) {
return 1;
}
return Math.max(1, Math.ceil(value));
}
function sanitizeRetryAfterDecimal(value: number | undefined | null, fallback: number): number {
if (value == null || !Number.isFinite(value) || value < 0) {
return fallback;
}
return Math.max(0.001, value);
}
function sanitizeResetTime(resetTime: Date): number {
const timestamp = resetTime.getTime();
if (!Number.isFinite(timestamp)) {
@@ -64,8 +54,8 @@ export class RateLimitError extends ThrottledError {
bucketHash?: string;
scope?: RateLimitScope;
}) {
const safeRetryAfter = sanitizeRetryAfter(retryAfter);
const safeRetryAfterDecimal = sanitizeRetryAfterDecimal(retryAfterDecimal, safeRetryAfter);
const safeRetryAfter = sanitizeRetryAfterSeconds(retryAfter);
const safeRetryAfterDecimal = sanitizeRetryAfterDecimalSeconds(retryAfterDecimal, safeRetryAfter);
const safeResetTimestamp = sanitizeResetTime(resetTime);
const safeLimit = Number.isFinite(limit) && limit > 0 ? limit : 1;
const safeResetAfterDecimal =
@@ -0,0 +1,15 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
export function sanitizeRetryAfterSeconds(value: number | undefined | null): number {
if (value == null || !Number.isFinite(value) || value < 0) {
return 1;
}
return Math.max(1, Math.ceil(value));
}
export function sanitizeRetryAfterDecimalSeconds(value: number | undefined | null, fallback: number): number {
if (value == null || !Number.isFinite(value) || value < 0) {
return fallback;
}
return Math.max(0.001, value);
}
@@ -2,22 +2,28 @@
import {APIErrorCodes} from '@fluxer/constants/src/ApiErrorCodes';
import {BadRequestError} from '@fluxer/errors/src/domains/core/BadRequestError';
import {
sanitizeRetryAfterDecimalSeconds,
sanitizeRetryAfterSeconds,
} from '@fluxer/errors/src/domains/core/RetryAfterSeconds';
export class SlowmodeRateLimitError extends BadRequestError {
constructor({
retryAfter,
retryAfterDecimal,
}: {
retryAfter: number;
retryAfter: number | undefined;
retryAfterDecimal?: number;
}) {
const safeRetryAfter = sanitizeRetryAfterSeconds(retryAfter);
const safeRetryAfterDecimal = sanitizeRetryAfterDecimalSeconds(retryAfterDecimal, safeRetryAfter);
super({
code: APIErrorCodes.SLOWMODE_RATE_LIMITED,
data: {
retry_after: retryAfterDecimal ?? retryAfter,
retry_after: safeRetryAfterDecimal,
},
headers: {
'Retry-After': retryAfter.toString(),
'Retry-After': safeRetryAfter.toString(),
},
});
}
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: ar\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: bg\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: cs\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: da\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: de\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: el\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: en-GB\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: en-US\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: es-419\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: es-ES\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: fi\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: fr\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: he\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: hi\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: hr\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: hu\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: id\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: it\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: ja\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: ko\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: lt\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: nl\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: no\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: pl\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: pt-BR\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: ro\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: ru\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: sv-SE\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: th\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: tr\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: uk\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: vi\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: zh-CN\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
+2 -2
View File
@@ -4,8 +4,8 @@
msgid ""
msgstr ""
"Project-Id-Version: fluxer-marketing\n"
"POT-Creation-Date: 2026-08-19 00:00+0000\n"
"PO-Revision-Date: 2026-08-19 00:00+0000\n"
"POT-Creation-Date: 2026-08-21 00:00+0000\n"
"PO-Revision-Date: 2026-08-21 00:00+0000\n"
"Language: zh-TW\n"
"MIME-Version: 1.0\n"
"Content-Type: text/plain; charset=UTF-8\n"
@@ -0,0 +1,19 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {AdminACLs} from '@fluxer/constants/src/AdminACLs';
import {SetUserAclsRequest} from '@fluxer/schema/src/domains/admin/AdminUserSchemas';
import {describe, expect, test} from 'vitest';
describe('SetUserAclsRequest', () => {
test('accepts every ACL the instance defines', () => {
const acls = Object.values(AdminACLs);
const result = SetUserAclsRequest.safeParse({user_id: '1', acls});
expect(result.success).toBe(true);
});
test('rejects more entries than there are ACLs', () => {
const acls = [...Object.values(AdminACLs), 'overflow:one'];
const result = SetUserAclsRequest.safeParse({user_id: '1', acls});
expect(result.success).toBe(false);
});
});
@@ -8,6 +8,7 @@ import {
UserFlags,
UserFlagsDescriptions,
} from '@fluxer/constants/src/UserConstants';
import {AdminACLs} from '@fluxer/constants/src/AdminACLs';
import {NSFWLevelSchema} from '@fluxer/schema/src/primitives/GuildValidators';
import {
createBitflagInt32Type,
@@ -21,6 +22,8 @@ import {
import {DiscriminatorType, EmailType, UsernameType} from '@fluxer/schema/src/primitives/UserValidators';
import {z} from 'zod';
const ADMIN_ACL_COUNT = Object.keys(AdminACLs).length;
export const UserAdminResponseSchema = z.object({
id: SnowflakeStringType,
username: z.string(),
@@ -65,7 +68,7 @@ export const UserAdminResponseSchema = z.object({
pending_bulk_message_deletion_at: z.string().nullable(),
deletion_reason_code: Int32Type.nullable(),
deletion_public_reason: z.string().nullable(),
acls: z.array(z.string()).max(100),
acls: z.array(z.string()).max(ADMIN_ACL_COUNT),
traits: z.array(z.string()).max(100),
has_totp: z.boolean(),
authenticator_types: z.array(Int32Type).max(10),
@@ -353,7 +356,7 @@ export type ScheduleAccountDeletionRequest = z.infer<typeof ScheduleAccountDelet
export const SetUserAclsRequest = z.object({
user_id: SnowflakeType.describe('ID of the user to set ACLs for'),
acls: z.array(createStringType(1, 64)).max(100).describe('List of access control permissions to assign'),
acls: z.array(createStringType(1, 64)).max(ADMIN_ACL_COUNT).describe('List of access control permissions to assign'),
});
export type SetUserAclsRequest = z.infer<typeof SetUserAclsRequest>;
+1
View File
@@ -5,6 +5,7 @@ edition.workspace = true
license.workspace = true
[dependencies]
fluxer_common = { path = "../../fluxer_common" }
anyhow = "1.0.102"
axum = "0.8.9"
base64 = "0.22.1"
+17 -18
View File
@@ -1,10 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use anyhow::{Context, Result};
use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD};
use clap::Args;
use hmac::{Hmac, KeyInit, Mac};
use sha2::Sha256;
#[derive(Debug, Clone, Args)]
pub struct SignExternalUrlArgs {
@@ -16,17 +13,12 @@ pub struct SignExternalUrlArgs {
}
pub fn sign_external_url(secret_key: &str, server_url: &str, upstream: &str) -> Result<String> {
let path = format!("v2/{}", URL_SAFE_NO_PAD.encode(upstream.as_bytes()));
let mut mac = Hmac::<Sha256>::new_from_slice(secret_key.as_bytes())
.context("failed to create HMAC signer")?;
mac.update(path.as_bytes());
let signature = URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes());
Ok(format!(
"{}/external/{}/{}",
fluxer_common::external_media_path::build_external_media_proxy_url(
server_url.trim_end_matches('/'),
signature,
path
))
upstream,
secret_key.as_bytes(),
)
.context("failed to build the external media proxy url")
}
#[cfg(test)]
@@ -41,13 +33,20 @@ mod tests {
"https://example.test/a b.jpg",
)
.unwrap();
let parts = signed.split('/').collect::<Vec<_>>();
assert!(signed.starts_with("http://127.0.0.1:19110/external/"));
assert_eq!(parts[5], "v2");
assert!(signed.ends_with("/https/example.test/a%20b.jpg"));
assert_eq!(
URL_SAFE_NO_PAD.decode(parts[6]).unwrap(),
b"https://example.test/a b.jpg"
"https://example.test/a b.jpg",
fluxer_common::external_media_path::reconstruct_original_url(
signed
.split_once("/external/")
.unwrap()
.1
.split_once('/')
.unwrap()
.1
)
.unwrap()
);
assert_eq!(URL_SAFE_NO_PAD.decode(parts[4]).unwrap().len(), 32);
}
}