Compare commits

...
13 Commits
Author SHA1 Message Date
HampusandGitHub 4e730832c7 fix(installer): pull images before the first start (#3248) 2026-10-07 02:37:35 +02:00
HampusandGitHub d7c00d4556 fix(schema): keep template topics optional after trimming (#3247) 2026-10-07 01:42:37 +02:00
HampusandGitHub 154b65afe5 docs(self-hosting): fix LiveKit CSP and backup guidance (#3246) 2026-10-07 00:09:39 +02:00
HampusandGitHub fcc2a3f64b docs(github): keep vulnerability reports out of chats (#3245) 2026-10-06 23:25:53 +02:00
HampusandGitHub 80456861ac fix(api): accept long forum topics in imported templates (#3244) 2026-10-06 22:46:47 +02:00
0c4f016ba2 feat(config): read secrets from NAME_FILE variables (#1421)
Co-authored-by: Hampus <[email protected]>
2026-10-06 21:43:08 +02:00
HampusandGitHub cc5545c333 fix(api): sync stripe customer email on change (#3243) 2026-10-06 21:38:25 +02:00
HampusandGitHub 6e28092cdc fix(app): keep mention highlight when mentions are suppressed (#3242) 2026-10-06 20:58:28 +02:00
HampusandGitHub 8b6910d505 chore(github): send bug reports and ideas to feedback.fluxer.com (#3241) 2026-10-06 19:14:26 +02:00
HampusandGitHub d87e31efaf fix(app): drop the reply when its target message is deleted (#3237) 2026-10-06 02:14:00 +02:00
HampusandGitHub 6618a6baf4 fix(installer): replace a stale installer before upgrading (#3235) 2026-10-05 21:50:16 +02:00
HampusandGitHub 22b8f5454b fix(app): skip forwarded messages when editing with arrow up (#3234) 2026-10-05 21:34:05 +02:00
HampusandGitHub 801bd3f106 fix(app): cycle dms in sidebar order with the keyboard (#3233) 2026-10-05 21:17:40 +02:00
64 changed files with 1332 additions and 552 deletions
+9 -9
View File
@@ -1,24 +1,24 @@
# Contributing to Fluxer
This policy applies to all issues, discussions, commits and pull requests.
This policy applies to all commits and pull requests.
## Scope
To prevent spam, only approved contributors may submit pull requests.
To request approval, comment on an existing issue and ask to implement it. For work that extends beyond a defect fix, open a [discussion](https://github.com/orgs/fluxerapp/discussions) first.
To request approval, comment on the [feedback.fluxer.com](https://feedback.fluxer.com) post you want to implement and ask to work on it. For work that extends beyond a defect fix, post a feature request there first.
Every pull request must:
- Target the repository's default branch.
- Include a closing reference for each repository issue it resolves.
- Link each feedback.fluxer.com post it resolves.
- Receive approval from a maintainer before it is merged.
Place each closing reference on a separate line:
Place each link on a separate line:
```text
Closes #123
Closes #456
Resolves https://feedback.fluxer.com/p/123
Resolves https://feedback.fluxer.com/p/456
```
## Authorship
@@ -78,11 +78,11 @@ Complete every section of the pull request template. Clearly describe:
## Reports and other contributions
Use the [bug report form](https://github.com/fluxerapp/fluxer/issues/new?template=bug-report.yaml) to report reproducible defects.
Report bugs and request features at [feedback.fluxer.com](https://feedback.fluxer.com).
Report security vulnerabilities privately through the channels specified in the [security policy](https://github.com/fluxerapp/fluxer/blob/main/.github/SECURITY.md). Do not report vulnerabilities in public issues or discussions.
Report security vulnerabilities privately through [fluxer.app/security](https://fluxer.app/security). Never post them publicly.
Use [discussions](https://github.com/orgs/fluxerapp/discussions) for feature proposals and self-hosting questions.
Read the [operator documentation](https://fluxer.dev) for self-hosting questions.
Submit translations through [Weblate](https://weblate.fluxer.tools), not through pull requests.
-41
View File
@@ -1,41 +0,0 @@
# yaml-language-server: $schema=https://www.schemastore.org/github-discussion.json
body:
- type: markdown
attributes:
value: |
Search existing discussions before posting a feature proposal.
Report vulnerabilities through the [private form](https://github.com/fluxerapp/fluxer/security/advisories/new) or <[email protected]>.
- type: textarea
id: problem
attributes:
label: Current problem
description: State what you are trying to do and what prevents it.
validations:
required: true
- type: textarea
id: proposal
attributes:
label: Proposed change
description: State the expected behaviour.
validations:
required: true
- type: textarea
id: notes
attributes:
label: Additional information
description: Optional. Include constraints, trade-offs, related discussions, screenshots or mockups.
validations:
required: false
- type: checkboxes
id: checks
attributes:
label: Acknowledgements
options:
- label: I searched existing discussions.
required: true
-83
View File
@@ -1,83 +0,0 @@
# yaml-language-server: $schema=https://www.schemastore.org/github-issue-forms.json
name: Bug report
description: Report a reproducible defect in Fluxer.
type: Bug
body:
- type: markdown
attributes:
value: |
Search [open and closed issues](https://github.com/fluxerapp/fluxer/issues?q=is%3Aissue) before filing a report.
Report vulnerabilities through the [private form](https://github.com/fluxerapp/fluxer/security/advisories/new) or <[email protected]>. Send account and billing requests to <[email protected]>.
- type: textarea
id: summary
attributes:
label: Observed behaviour
description: State what happened and what you expected.
validations:
required: true
- type: textarea
id: steps
attributes:
label: Reproduction steps
description: Give numbered steps starting from a fresh app or session.
placeholder: |
1. Go to ...
2. Select ...
3. Observe ...
validations:
required: true
- type: input
id: build
attributes:
label: Build information
description: >-
Open User Settings, scroll to the bottom of the left sidebar, and select
the build information. Fluxer copies it to the clipboard.
validations:
required: true
- type: dropdown
id: surface
attributes:
label: Affected surface
multiple: true
options:
- Desktop app
- Web app
- Voice, video, or Go Live
- Self-hosted instance
- HTTP API or Gateway
- Documentation site
validations:
required: true
- type: input
id: instance
attributes:
label: Instance
description: For a self-hosted instance, include the release tag and database backend.
placeholder: fluxer.app
validations:
required: false
- type: textarea
id: evidence
attributes:
label: Evidence
description: Attach relevant logs, screenshots or recordings. Remove tokens, keys, private messages and other personal data. Configuration files may contain secrets.
validations:
required: false
- type: checkboxes
id: checks
attributes:
label: Acknowledgements
options:
- label: I searched open and closed issues.
required: true
- label: I removed secrets and unrelated personal data from the report.
required: true
-18
View File
@@ -1,18 +0,0 @@
# yaml-language-server: $schema=https://www.schemastore.org/github-issue-config.json
blank_issues_enabled: false
contact_links:
- name: Mobile client bugs
url: https://github.com/fluxerapp/flutter_client#bug-reporting
about: Read the reporting instructions for the Fluxer mobile client.
- name: Account and billing support
url: https://fluxer.app/help
about: Find account help and support contact details.
- name: Feature proposals
url: https://github.com/orgs/fluxerapp/discussions
about: Propose a feature in a discussion.
- name: Translations
url: https://weblate.fluxer.tools
about: Improve an existing locale or start a new one.
- name: Self-hosting support
url: https://fluxer.dev
about: Read the operator documentation, then open a discussion if the problem remains.
-44
View File
@@ -1,44 +0,0 @@
# yaml-language-server: $schema=https://www.schemastore.org/github-issue-forms.json
name: Documentation
description: Report incorrect, missing or unclear documentation.
type: Task
labels:
- docs
body:
- type: markdown
attributes:
value: |
This form covers <https://fluxer.dev> and operator documentation.
- type: textarea
id: issue
attributes:
label: Documentation defect
description: State what the page says and what is correct. For missing content, state what information you needed.
validations:
required: true
- type: input
id: location
attributes:
label: Location
description: Provide the page URL or file path and heading.
placeholder: https://fluxer.dev/gateway/overview/
validations:
required: false
- type: textarea
id: suggestion
attributes:
label: Proposed wording
description: Optional.
validations:
required: false
- type: checkboxes
id: checks
attributes:
label: Acknowledgements
options:
- label: I searched open and closed issues.
required: true
+2 -2
View File
@@ -1,7 +1,7 @@
# Security policy
Do not report a vulnerability in an issue, pull request, or discussion.
Do not report a vulnerability in a pull request, on feedback.fluxer.com, in a Fluxer community, or in a direct message to staff.
Submit a report through [GitHub private vulnerability reporting](https://github.com/fluxerapp/fluxer/security/advisories/new) or email <security@fluxer.com>. Include the affected component, impact, reproduction steps, and supporting evidence. Remove unrelated personal data and secrets.
Submit a report through <https://fluxer.app/security> or email <security@fluxer.com>. Include the affected component, impact, reproduction steps, and supporting evidence. Remove unrelated personal data and secrets.
The programme scope, testing rules, safe harbour, disclosure process, and reward terms are published at <https://fluxer.app/security>. That page is authoritative.
+2 -2
View File
@@ -1,6 +1,6 @@
Closes #
Resolves https://feedback.fluxer.com/p/
<!-- Repeat this line for each resolved issue, up to 20. Remove the placeholder only if no issue is resolved and the approval gate does not apply. -->
<!-- Repeat this line for each feedback.fluxer.com post this resolves, up to 20. Remove the placeholder only if no post is resolved and the approval gate does not apply. -->
## Summary
+3
View File
@@ -23,6 +23,9 @@
# Fluxer
> [!IMPORTANT]
> Bug reports and feature requests have moved to [feedback.fluxer.com](https://feedback.fluxer.com). Sign in with your Fluxer account to post, vote and follow updates. GitHub Issues and Discussions are closed. Report security vulnerabilities privately through [fluxer.app/security](https://fluxer.app/security).
Fluxer is a free and open source instant messaging and VoIP chat app built for friends, groups, and communities.
<p align="center">
+1 -1
View File
@@ -263,7 +263,7 @@ FLUXER_VAPID_PRIVATE_KEY=CHANGE_ME
# only when a browser must reach an origin the defaults do not cover. Separate
# several with spaces or commas. The three values below are illustrations.
#FLUXER_CSP_EXTRA_DEFAULT_SRC=
#FLUXER_CSP_EXTRA_CONNECT_SRC=wss://livekit.example.com:7881
#FLUXER_CSP_EXTRA_CONNECT_SRC=wss://livekit.example.com
#FLUXER_CSP_EXTRA_IMG_SRC=https://cdn.example.com
#FLUXER_CSP_EXTRA_MEDIA_SRC=
#FLUXER_CSP_EXTRA_FONT_SRC=
@@ -12,6 +12,7 @@ import type {EntityAssetService, PreparedAssetUpload} from '@app/api/infrastruct
import {Logger} from '@app/api/Logger';
import {getInstanceConfigRepository} from '@app/api/middleware/ServiceSingletons';
import type {User} from '@app/api/models/User';
import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer';
import {assertNoDiscriminatorChange, reserveUsername, type UsernameReservation} from '@app/api/user/UniqueUsernames';
import {USERNAME_MODE_DISCRIMINATOR} from '@app/api/user/UserTag';
import {TagAlreadyTakenError} from '@fluxer/errors/src/domains/user/TagAlreadyTakenError';
@@ -218,6 +219,7 @@ export class AdminUserProfileService {
users: userRepository,
cache: cacheService,
contactChangeLog: contactChangeLogService,
worker: workerService,
} = this.deps.apiContext.services;
const {auditService, updatePropagator} = this.deps;
const userId = createUserID(data.user_id);
@@ -240,6 +242,7 @@ export class AdminUserProfileService {
reason: 'admin_action',
actorUserId: adminUserId,
});
await enqueueStripeCustomerEmailSync(workerService, user, updatedUser);
await auditService.createAuditLog({
adminUserId,
targetType: 'user',
@@ -2,11 +2,12 @@
import {createTestAccount, setUserACLs} from '@app/api/auth/tests/AuthTestUtils';
import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHarness';
import {NoopWorkerService} from '@app/api/test/NoopWorkerService';
import {HTTP_STATUS, TEST_CREDENTIALS} from '@app/api/test/TestConstants';
import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder';
import {AdminACLs} from '@fluxer/constants/src/AdminACLs';
import {UserFlags} from '@fluxer/constants/src/UserConstants';
import {afterAll, beforeAll, beforeEach, describe, expect, test} from 'vitest';
import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, test, vi} from 'vitest';
interface ChangeLogResponse {
entries: Array<{
@@ -44,6 +45,9 @@ describe('Admin User Change Log and Flags', () => {
beforeEach(async () => {
await harness.reset();
});
afterEach(() => {
vi.restoreAllMocks();
});
afterAll(async () => {
await harness?.shutdown();
});
@@ -176,6 +180,37 @@ describe('Admin User Change Log and Flags', () => {
.execute();
});
});
describe('PATCH /admin/users/{user_id}/email', () => {
test('queues a Stripe customer email sync for users with a Stripe customer', async () => {
const admin = await createTestAccount(harness);
await setUserACLs(harness, admin, [AdminACLs.AUTHENTICATE, AdminACLs.WILDCARD]);
const target = await createTestAccount(harness);
await createBuilderWithoutAuth(harness)
.post(`/test/users/${target.userId}/premium`)
.body({stripe_customer_id: 'cus_admin_email_sync'})
.expect(HTTP_STATUS.OK)
.execute();
const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob');
await createBuilder(harness, `${admin.token}`)
.patch(`/admin/users/${target.userId}/email`)
.body({email: `admin-changed-${Date.now()}@example.com`})
.expect(HTTP_STATUS.OK)
.execute();
expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: target.userId});
});
test('does not queue a Stripe customer email sync for users without a Stripe customer', async () => {
const admin = await createTestAccount(harness);
await setUserACLs(harness, admin, [AdminACLs.AUTHENTICATE, AdminACLs.WILDCARD]);
const target = await createTestAccount(harness);
const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob');
await createBuilder(harness, `${admin.token}`)
.patch(`/admin/users/${target.userId}/email`)
.body({email: `admin-changed-${Date.now()}@example.com`})
.expect(HTTP_STATUS.OK)
.execute();
expect(addJob).not.toHaveBeenCalledWith('syncStripeCustomerEmail', expect.anything());
});
});
describe('PUT /admin/users/{user_id}/email-verification', () => {
test('verifying email clears email_bounced', async () => {
const admin = await createTestAccount(harness);
+3 -1
View File
@@ -6,6 +6,7 @@ import * as AuthSession from '@app/api/auth/AuthSession';
import * as AuthUtility from '@app/api/auth/AuthUtility';
import {createEmailRevertToken} from '@app/api/BrandedTypes';
import type {User} from '@app/api/models/User';
import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer';
import {mapUserToPrivateResponse} from '@app/api/user/UserMappers';
import {ValidationErrorCodes} from '@fluxer/constants/src/ValidationErrorCodes';
import {InputValidationError} from '@fluxer/errors/src/domains/core/InputValidationError';
@@ -44,7 +45,7 @@ export async function revertEmailChange(
user_id: string;
token: string;
}> {
const {users, gateway, contactChangeLog, config} = ctx.services;
const {users, gateway, contactChangeLog, config, worker} = ctx.services;
const {token, password, request} = params;
const tokenData = await users.getEmailRevertToken(token);
if (!tokenData) {
@@ -101,5 +102,6 @@ export async function revertEmailChange(
reason: 'user_requested',
actorUserId: user.id,
});
await enqueueStripeCustomerEmailSync(worker, user, updatedUser);
return {user_id: updatedUser.id.toString(), token: authToken};
}
@@ -9,8 +9,9 @@ import {
type TestAccount,
} from '@app/api/auth/tests/AuthTestUtils';
import type {ApiTestHarness} from '@app/api/test/ApiTestHarness';
import {NoopWorkerService} from '@app/api/test/NoopWorkerService';
import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder';
import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest';
import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi} from 'vitest';
interface BouncedEmailRequestNewResponse {
ticket: string;
@@ -51,12 +52,21 @@ describe('Bounced email recovery flow', () => {
await harness.reset();
await clearTestEmails(harness);
});
afterEach(() => {
vi.restoreAllMocks();
});
afterAll(async () => {
await harness?.shutdown();
});
it('allows bounced users to replace email without original-email verification', async () => {
const account = await createTestAccount(harness);
await markEmailAsBounced(harness, account);
await createBuilderWithoutAuth(harness)
.post(`/test/users/${account.userId}/premium`)
.body({stripe_customer_id: 'cus_bounced_email_sync'})
.expect(200)
.execute();
const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob');
const initialMe = await createBuilder<UserPrivateResponse>(harness, account.token)
.get('/users/@me')
.expect(200)
@@ -94,6 +104,7 @@ describe('Bounced email recovery flow', () => {
const finalMe = await createBuilder<UserPrivateResponse>(harness, account.token).get('/users/@me').execute();
expect(finalMe.email).toBe(replacementEmail);
expect(finalMe.email_bounced).toBe(false);
expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: account.userId});
});
it('rejects bounced-email recovery for accounts that are not marked as bounced', async () => {
const account = await createTestAccount(harness);
@@ -12,8 +12,9 @@ import {
type TestAccount,
} from '@app/api/auth/tests/AuthTestUtils';
import type {ApiTestHarness} from '@app/api/test/ApiTestHarness';
import {NoopWorkerService} from '@app/api/test/NoopWorkerService';
import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder';
import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest';
import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi} from 'vitest';
interface EmailChangeStartResponse {
ticket: string;
@@ -129,6 +130,9 @@ describe('Email change flow', () => {
await harness.reset();
await clearTestEmails(harness);
});
afterEach(() => {
vi.restoreAllMocks();
});
afterAll(async () => {
await harness?.shutdown();
});
@@ -341,13 +345,14 @@ describe('Email change flow', () => {
.execute();
expect(updated.email).toBe(newEmail);
});
it('applies email changes for users who have ever purchased', async () => {
it('applies email changes for users who have ever purchased and syncs their Stripe customer', async () => {
const account = await createTestAccount(harness);
await createBuilderWithoutAuth(harness)
.post(`/test/users/${account.userId}/premium`)
.body({has_ever_purchased: true})
.body({has_ever_purchased: true, stripe_customer_id: 'cus_email_change_sync'})
.expect(200)
.execute();
const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob');
const startResp = await startEmailChange(harness, account, account.password);
let originalProof: string;
if (startResp.require_original) {
@@ -387,9 +392,11 @@ describe('Email change flow', () => {
.execute();
expect(updated.email).toBe(newEmail);
expect(updated.has_ever_purchased).toBe(true);
expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: account.userId});
});
it('applies ordinary claimed email changes', async () => {
const account = await createTestAccount(harness);
const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob');
const startResp = await startEmailChange(harness, account, account.password);
const emails = await listTestEmails(harness, {recipient: account.email});
const originalEmail = findLastTestEmail(emails, 'email_change_original');
@@ -425,6 +432,7 @@ describe('Email change flow', () => {
.execute();
expect(updated.email).toBe(newEmail);
expect(updated.verified).toBe(true);
expect(addJob).not.toHaveBeenCalledWith('syncStripeCustomerEmail', expect.anything());
});
it('requires MFA (not password) for email_token apply when user has TOTP enabled', async () => {
const account = await createTestAccount(harness);
@@ -10,8 +10,9 @@ import {
type TestAccount,
} from '@app/api/auth/tests/AuthTestUtils';
import type {ApiTestHarness} from '@app/api/test/ApiTestHarness';
import {NoopWorkerService} from '@app/api/test/NoopWorkerService';
import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder';
import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest';
import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi} from 'vitest';
interface EmailChangeStartResponse {
ticket: string;
@@ -130,11 +131,20 @@ describe('Email revert flow', () => {
await harness.reset();
await clearTestEmails(harness);
});
afterEach(() => {
vi.restoreAllMocks();
});
afterAll(async () => {
await harness?.shutdown();
});
it('restores original email and clears mfa', async () => {
const account = await createTestAccount(harness);
await createBuilderWithoutAuth(harness)
.post(`/test/users/${account.userId}/premium`)
.body({stripe_customer_id: 'cus_email_revert_sync'})
.expect(200)
.execute();
const addJob = vi.spyOn(NoopWorkerService.prototype, 'addJob');
const startResp = await startEmailChange(harness, account, account.password);
let originalProof: string;
if (startResp.require_original) {
@@ -178,6 +188,7 @@ describe('Email revert flow', () => {
expect(revertEmail?.metadata?.token).toBeDefined();
const revertToken = revertEmail!.metadata!.token!;
const newPassword = uniquePassword();
addJob.mockClear();
const revertResp = await createBuilderWithoutAuth<EmailRevertResponse>(harness)
.post('/auth/email-revert')
.body({
@@ -186,6 +197,7 @@ describe('Email revert flow', () => {
})
.execute();
expect(revertResp.token.length).toBeGreaterThan(0);
expect(addJob).toHaveBeenCalledWith('syncStripeCustomerEmail', {userId: account.userId});
await createBuilder(harness, account.token).get('/users/@me').expect(401).execute();
const user = await createBuilder<UserPrivateResponse>(harness, revertResp.token).get('/users/@me').execute();
expect(user.email).toBe(account.email);
@@ -6,7 +6,7 @@ import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHa
import {createBuilder} from '@app/api/test/TestRequestBuilder';
import {ChannelTypes, Permissions} from '@fluxer/constants/src/ChannelConstants';
import {SystemChannelFlags} from '@fluxer/constants/src/GuildConstants';
import {VOICE_CHANNEL_USER_LIMIT_MAX} from '@fluxer/constants/src/LimitConstants';
import {CHANNEL_TOPIC_MAX_LENGTH, VOICE_CHANNEL_USER_LIMIT_MAX} from '@fluxer/constants/src/LimitConstants';
import type {GuildResponse} from '@fluxer/schema/src/domains/guild/GuildResponseSchemas';
import {afterAll, beforeAll, beforeEach, describe, expect, test} from 'vitest';
@@ -215,7 +215,7 @@ describe('Guild Template Import', () => {
['a slowmode above the channel maximum', {rate_limit_per_user: 1_000_000_000}],
['a fractional position', {position: 0.5}],
['a negative position', {position: -3}],
['a topic above the channel maximum', {topic: 'x'.repeat(1025)}],
['a topic above the template maximum', {topic: 'x'.repeat(4097)}],
['a name above the channel maximum', {name: 'x'.repeat(101)}],
['a negative user limit', {type: ChannelTypes.GUILD_VOICE, user_limit: -1}],
['a voice connection limit above the maximum', {type: ChannelTypes.GUILD_VOICE, voice_connection_limit: 100_000}],
@@ -273,6 +273,26 @@ describe('Guild Template Import', () => {
expect(channels.find((channel) => channel.name === 'town-hall')?.user_limit).toBe(VOICE_CHANNEL_USER_LIMIT_MAX);
expect(channels.find((channel) => channel.name === 'lounge')?.voice_connection_limit).toBe(100);
});
test('accepts forum-length topics and shortens them to the channel topic limit', async () => {
const account = await createTestAccount(harness);
const longTopic = `${'a'.repeat(CHANNEL_TOPIC_MAX_LENGTH - 1)}\u{1F600}${'b'.repeat(300)}`;
const guild = await createBuilder<GuildResponse>(harness, account.token)
.post('/guilds')
.body({
name: 'Forum Guild',
template: buildMinimalTemplate({
channels: [
{id: 6001, type: ChannelTypes.GUILD_TEXT, name: 'general', position: 0, topic: longTopic},
{id: 6002, type: 15, name: 'projects', position: 1, topic: 'c'.repeat(1356)},
],
}),
})
.execute();
const channels = await getGuildChannels(harness, account.token, guild.id);
const general = channels.find((channel) => channel.name === 'general');
expect(general?.topic).toBe('a'.repeat(CHANNEL_TOPIC_MAX_LENGTH - 1));
expect(channels.find((channel) => channel.name === 'projects')).toBeUndefined();
});
});
const DEFAULT_EVERYONE_PERMISSIONS = Permissions.VIEW_CHANNEL.toString();
+2 -2
View File
@@ -33075,8 +33075,8 @@
"anyOf": [{"type": "string", "maxLength": 100}, {"type": "null"}]
},
"topic": {
"description": "The channel topic",
"anyOf": [{"type": "string", "maxLength": 1024}, {"type": "null"}]
"description": "The channel topic, shortened to the Fluxer channel topic limit",
"anyOf": [{"type": "string", "maxLength": 4096}, {"type": "null"}]
},
"position": {"description": "The position of the channel", "$ref": "#/components/schemas/Int32Type"},
"parent_id": {
+119
View File
@@ -0,0 +1,119 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Logger} from '@app/api/Logger';
import {getBillingRepository} from '@app/api/middleware/ServiceRegistry';
import type {User} from '@app/api/models/User';
import type {IUserRepository} from '@app/api/user/IUserRepository';
import type {WorkerTaskName} from '@app/api/worker/WorkerLaneConfig';
import {StripeError} from '@fluxer/errors/src/domains/payment/StripeError';
import type {ICacheService} from '@pkgs/cache/src/ICacheService';
import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService';
import {seconds} from 'itty-time';
import type Stripe from 'stripe';
const CUSTOMER_LOCK_TTL_SECONDS = seconds('30 seconds');
interface EnsureStripeCustomerParams {
stripe: Stripe;
user: User;
userRepository: IUserRepository;
cacheService: ICacheService;
}
export function isStripeResourceMissingError(error: unknown): boolean {
return typeof error === 'object' && error !== null && 'code' in error && error.code === 'resource_missing';
}
export async function ensureStripeCustomer({
stripe,
user,
userRepository,
cacheService,
}: EnsureStripeCustomerParams): Promise<User> {
if (user.stripeCustomerId) {
try {
await syncStripeCustomerEmail(stripe, user);
} catch (error) {
Logger.warn(
{error, userId: user.id, customerId: user.stripeCustomerId},
'Failed to sync Stripe customer email, continuing with the stored customer',
);
}
return user;
}
const lockKey = `stripe_customer_create_lock:${user.id}`;
const lockToken = await cacheService.acquireLock(lockKey, CUSTOMER_LOCK_TTL_SECONDS);
if (!lockToken) {
const freshUser = await userRepository.findUnique(user.id);
if (freshUser?.stripeCustomerId) {
return freshUser;
}
throw new StripeError('Failed to acquire customer creation lock');
}
try {
const freshUser = await userRepository.findUnique(user.id);
if (freshUser?.stripeCustomerId) {
return freshUser;
}
const customer = await stripe.customers.create({
email: user.email ?? undefined,
metadata: {
userId: user.id.toString(),
},
});
await mirrorStripeCustomer(customer, user);
const updatedUser = await userRepository.patchUpsert(user.id, {stripe_customer_id: customer.id}, user.toRow());
Logger.debug({userId: user.id, customerId: customer.id}, 'Stripe customer created');
return updatedUser;
} finally {
try {
const released = await cacheService.releaseLock(lockKey, lockToken);
if (!released) {
Logger.warn({userId: user.id, lockKey}, 'Customer creation lock token no longer matched on release');
}
} catch (error) {
Logger.error({error, userId: user.id, lockKey}, 'Failed to release customer creation lock');
}
}
}
export async function syncStripeCustomerEmail(stripe: Stripe, user: User): Promise<void> {
const customerId = user.stripeCustomerId;
const email = user.email;
if (!customerId || !email) {
return;
}
const mirrored = await getBillingRepository().customers.findById(customerId);
if (mirrored?.deleted || mirrored?.email === email) {
return;
}
const customer = await stripe.customers.update(customerId, {email});
await mirrorStripeCustomer(customer, user);
Logger.info({userId: user.id, customerId}, 'Synced Stripe customer email with account email');
}
export async function enqueueStripeCustomerEmailSync(
workerService: IWorkerService<WorkerTaskName>,
oldUser: User,
newUser: User,
): Promise<void> {
if (!newUser.stripeCustomerId || !newUser.email || oldUser.email === newUser.email) {
return;
}
try {
await workerService.addJob('syncStripeCustomerEmail', {userId: newUser.id.toString()});
} catch (error) {
Logger.warn(
{error, userId: newUser.id, customerId: newUser.stripeCustomerId},
'Failed to enqueue Stripe customer email sync, checkout will sync it instead',
);
}
}
async function mirrorStripeCustomer(customer: Stripe.Customer, user: User): Promise<void> {
try {
await getBillingRepository().customers.upsertFromStripe(customer, {knownUserId: user.id});
} catch (mirrorErr) {
Logger.error({mirrorErr, customerId: customer.id}, 'Mirror upsert failed after Stripe write; reconciler will heal');
}
}
@@ -5,7 +5,7 @@ import {Config} from '@app/api/Config';
import type {IGatewayService} from '@app/api/infrastructure/IGatewayService';
import {Logger} from '@app/api/Logger';
import {getBillingRepository} from '@app/api/middleware/ServiceRegistry';
import type {User} from '@app/api/models/User';
import {ensureStripeCustomer} from '@app/api/stripe/StripeCustomer';
import type {IUserRepository} from '@app/api/user/IUserRepository';
import {mapUserToPrivateResponse} from '@app/api/user/UserMappers';
import {UserFlags} from '@fluxer/constants/src/UserConstants';
@@ -14,11 +14,8 @@ import {StripeError} from '@fluxer/errors/src/domains/payment/StripeError';
import {StripePaymentNotAvailableError} from '@fluxer/errors/src/domains/payment/StripePaymentNotAvailableError';
import {UnknownUserError} from '@fluxer/errors/src/domains/user/UnknownUserError';
import type {ICacheService} from '@pkgs/cache/src/ICacheService';
import {seconds} from 'itty-time';
import type Stripe from 'stripe';
const CUSTOMER_LOCK_TTL_SECONDS = seconds('30 seconds');
export class AgeVerificationService {
constructor(
private stripe: Stripe | null,
@@ -38,7 +35,12 @@ export class AgeVerificationService {
if (user.flags & UserFlags.AGE_VERIFIED_ADULT) {
throw new AgeVerificationAlreadyVerifiedError();
}
const customerUser = await this.ensureStripeCustomer(user);
const customerUser = await ensureStripeCustomer({
stripe: this.stripe,
user,
userRepository: this.userRepository,
cacheService: this.cacheService,
});
const customerId = customerUser.stripeCustomerId;
if (!customerId) {
throw new StripeError('Stripe customer id missing after customer setup');
@@ -134,58 +136,4 @@ export class AgeVerificationService {
});
Logger.info({userId}, 'Age verification completed successfully');
}
private async ensureStripeCustomer(user: User): Promise<User> {
if (user.stripeCustomerId) {
return user;
}
if (!this.stripe) {
throw new StripePaymentNotAvailableError();
}
const lockKey = `stripe_customer_create_lock:${user.id}`;
const lockToken = await this.cacheService.acquireLock(lockKey, CUSTOMER_LOCK_TTL_SECONDS);
if (!lockToken) {
const freshUser = await this.userRepository.findUnique(user.id);
if (freshUser?.stripeCustomerId) {
return freshUser;
}
throw new StripeError('Failed to acquire customer creation lock');
}
try {
const freshUser = await this.userRepository.findUnique(user.id);
if (freshUser?.stripeCustomerId) {
return freshUser;
}
const customer = await this.stripe.customers.create({
email: user.email ?? undefined,
metadata: {
userId: user.id.toString(),
},
});
try {
await getBillingRepository().customers.upsertFromStripe(customer, {knownUserId: user.id});
} catch (mirrorErr) {
Logger.error(
{mirrorErr, customerId: customer.id},
'Mirror upsert failed after Stripe write; reconciler will heal',
);
}
const updatedUser = await this.userRepository.patchUpsert(
user.id,
{stripe_customer_id: customer.id},
user.toRow(),
);
Logger.debug({userId: user.id, customerId: customer.id}, 'Stripe customer created for age verification');
return updatedUser;
} finally {
try {
const released = await this.cacheService.releaseLock(lockKey, lockToken);
if (!released) {
Logger.warn({userId: user.id, lockKey}, 'Customer creation lock token no longer matched on release');
}
} catch (error) {
Logger.error({error, userId: user.id, lockKey}, 'Failed to release customer creation lock');
}
}
}
}
@@ -11,6 +11,7 @@ import type {StoreEntitlementService} from '@app/api/store_billing/StoreEntitlem
import {getBillingBranding} from '@app/api/stripe/BillingBranding';
import {getEffectiveBillingConfig, isCurrentCatalogPriceId} from '@app/api/stripe/BillingConfigCache';
import type {ProductInfo, ProductRegistry} from '@app/api/stripe/ProductRegistry';
import {ensureStripeCustomer, isStripeResourceMissingError} from '@app/api/stripe/StripeCustomer';
import {getCachedStripePriceSummary, type StripePriceSummary} from '@app/api/stripe/StripePriceSummaryCache';
import {
canProvisionPremiumFromSubscriptionStatus,
@@ -42,15 +43,10 @@ import {UnclaimedAccountCannotMakePurchasesError} from '@fluxer/errors/src/domai
import {UnknownUserError} from '@fluxer/errors/src/domains/user/UnknownUserError';
import type {CheckoutPaymentMethod} from '@fluxer/schema/src/domains/premium/GiftCodeSchemas';
import type {ICacheService} from '@pkgs/cache/src/ICacheService';
import {seconds} from 'itty-time';
import type Stripe from 'stripe';
export const EU_WITHDRAWAL_WAIVER_TEXT_VERSION = '2026-04-23';
function isStripeResourceMissingError(error: unknown): boolean {
return typeof error === 'object' && error !== null && 'code' in error && error.code === 'resource_missing';
}
type CheckoutSessionCreateParams = Stripe.Checkout.SessionCreateParams;
type CheckoutSessionMode = CheckoutSessionCreateParams['mode'];
type CheckoutSessionPaymentMethodType = NonNullable<CheckoutSessionCreateParams['payment_method_types']>[number];
@@ -653,8 +649,6 @@ export class StripeCheckoutService {
}
}
private static readonly CUSTOMER_LOCK_TTL_SECONDS = seconds('30 seconds');
private resolveConfiguredPriceIds(countryCode?: string): ResolvedPriceIds {
const recurringCurrencyPreferences = getCurrencyPreferences(countryCode);
const giftCurrencyPreferences = getGiftCurrencyPreferences(countryCode);
@@ -888,59 +882,14 @@ export class StripeCheckoutService {
}
private async ensureStripeCustomer(existingUser: User): Promise<User> {
const user = await this.clearStaleStripeCustomer(existingUser);
if (user.stripeCustomerId) {
return user;
}
if (!this.stripe) {
throw new StripePaymentNotAvailableError();
}
const lockKey = `stripe_customer_create_lock:${user.id}`;
const lockToken = await this.cacheService.acquireLock(lockKey, StripeCheckoutService.CUSTOMER_LOCK_TTL_SECONDS);
if (!lockToken) {
const freshUser = await this.userRepository.findUnique(user.id);
if (freshUser?.stripeCustomerId) {
return freshUser;
}
throw new StripeError('Failed to acquire customer creation lock');
}
try {
const freshUser = await this.userRepository.findUnique(user.id);
if (freshUser?.stripeCustomerId) {
return freshUser;
}
const customer = await this.stripe.customers.create({
email: user.email ?? undefined,
metadata: {
userId: user.id.toString(),
},
});
try {
await getBillingRepository().customers.upsertFromStripe(customer, {knownUserId: user.id});
} catch (mirrorErr) {
Logger.error(
{mirrorErr, customerId: customer.id},
'Mirror upsert failed after Stripe write; reconciler will heal',
);
}
const updatedUser = await this.userRepository.patchUpsert(
user.id,
{
stripe_customer_id: customer.id,
},
user.toRow(),
);
Logger.debug({userId: user.id, customerId: customer.id}, 'Stripe customer created');
return updatedUser;
} finally {
try {
const released = await this.cacheService.releaseLock(lockKey, lockToken);
if (!released) {
Logger.warn({userId: user.id, lockKey}, 'Customer creation lock token no longer matched on release');
}
} catch (error) {
Logger.error({error, userId: user.id, lockKey}, 'Failed to release customer creation lock');
}
}
return ensureStripeCustomer({
stripe: this.stripe,
user: await this.clearStaleStripeCustomer(existingUser),
userRepository: this.userRepository,
cacheService: this.cacheService,
});
}
}
@@ -0,0 +1,233 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {createTestAccount, type TestAccount} from '@app/api/auth/tests/AuthTestUtils';
import {createUserID} from '@app/api/BrandedTypes';
import {Config} from '@app/api/Config';
import {getBillingRepository} from '@app/api/middleware/ServiceRegistry';
import {getUserRepository} from '@app/api/middleware/ServiceSingletons';
import {getStripeClient} from '@app/api/stripe/StripeClient';
import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHarness';
import {NoopLogger} from '@app/api/test/mocks/NoopLogger';
import {createStripeApiHandlers, type StripeApiHandlers} from '@app/api/test/msw/handlers/StripeApiHandlers';
import {server} from '@app/api/test/msw/server';
import {HTTP_STATUS} from '@app/api/test/TestConstants';
import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder';
import {PaymentRepository} from '@app/api/user/repositories/PaymentRepository';
import reconcileUserPayments from '@app/api/worker/tasks/ReconcileUserPayments';
import syncStripeCustomerEmail from '@app/api/worker/tasks/SyncStripeCustomerEmail';
import {clearWorkerDependencies, setWorkerDependenciesForTest} from '@app/api/worker/WorkerContext';
import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask';
import type Stripe from 'stripe';
import {afterAll, afterEach, beforeAll, beforeEach, describe, expect, test} from 'vitest';
const MOCK_PRICES = {
monthlyUsd: 'price_email_sync_monthly_usd',
yearlyUsd: 'price_email_sync_yearly_usd',
gift1MonthUsd: 'price_email_sync_gift_1_month_usd',
gift1YearUsd: 'price_email_sync_gift_1_year_usd',
};
const MOCK_PRICE_SEEDS = {
[MOCK_PRICES.monthlyUsd]: {unit_amount: 499, currency: 'usd', interval: 'month' as const},
[MOCK_PRICES.yearlyUsd]: {unit_amount: 4999, currency: 'usd', interval: 'year' as const},
};
const STALE_EMAIL = '[email protected]';
function createHelpers(): WorkerTaskHelpers {
return {
logger: new NoopLogger(),
jobId: 0n,
addJob: async () => 0n,
reportProgress: async () => {},
shouldCancel: async () => false,
setContextLink: async () => {},
};
}
async function currentEmail(account: TestAccount): Promise<string> {
const user = await getUserRepository().findUnique(createUserID(BigInt(account.userId)));
return user!.email!;
}
async function mirrorCustomer(customerId: string, account: TestAccount, email: string | null): Promise<void> {
await getBillingRepository().customers.upsertFromStripe(
{
id: customerId,
object: 'customer',
email,
created: 1_600_000_000,
livemode: false,
metadata: {userId: account.userId},
} as unknown as Stripe.Customer,
{knownUserId: BigInt(account.userId)},
);
}
describe('Stripe customer email sync', () => {
let harness: ApiTestHarness;
let stripeHandlers: StripeApiHandlers;
let originalPrices: typeof Config.stripe.prices | undefined;
async function createCustomerAccount(customerId: string): Promise<TestAccount> {
const account = await createTestAccount(harness);
await createBuilderWithoutAuth(harness)
.post(`/test/users/${account.userId}/security-flags`)
.body({email_verified: true})
.expect(HTTP_STATUS.OK)
.execute();
await createBuilderWithoutAuth(harness)
.post(`/test/users/${account.userId}/premium`)
.body({stripe_customer_id: customerId})
.expect(HTTP_STATUS.OK)
.execute();
return account;
}
async function runTask(account: TestAccount): Promise<void> {
await syncStripeCustomerEmail({userId: account.userId}, createHelpers());
}
function emailUpdatesFor(customerId: string): Array<unknown> {
return stripeHandlers.spies.updatedCustomers
.filter((update) => update.id === customerId && 'email' in update.params)
.map((update) => update.params.email);
}
beforeAll(async () => {
originalPrices = Config.stripe.prices;
Config.stripe.prices = MOCK_PRICES;
harness = await createApiTestHarness();
});
afterAll(async () => {
await harness.shutdown();
Config.stripe.prices = originalPrices;
});
beforeEach(async () => {
await harness.reset();
Config.stripe.prices = MOCK_PRICES;
stripeHandlers = createStripeApiHandlers({prices: MOCK_PRICE_SEEDS, subscriptionsListEmpty: true});
server.use(...stripeHandlers.handlers);
const stripe = getStripeClient();
expect(stripe).not.toBeNull();
setWorkerDependenciesForTest({userRepository: getUserRepository(), stripe});
});
afterEach(() => {
clearWorkerDependencies();
});
describe('worker task', () => {
test('pushes the account email to a customer that still has the old address', async () => {
const account = await createCustomerAccount('cus_sync_stale');
await mirrorCustomer('cus_sync_stale', account, STALE_EMAIL);
await runTask(account);
const email = await currentEmail(account);
expect(emailUpdatesFor('cus_sync_stale')).toEqual([email]);
expect((await getBillingRepository().customers.findById('cus_sync_stale'))?.email).toBe(email);
});
test('pushes the account email when the customer is not mirrored yet', async () => {
const account = await createCustomerAccount('cus_sync_unmirrored');
await runTask(account);
expect(emailUpdatesFor('cus_sync_unmirrored')).toEqual([await currentEmail(account)]);
});
test('does nothing when the customer already has the account email', async () => {
const account = await createCustomerAccount('cus_sync_current');
await mirrorCustomer('cus_sync_current', account, await currentEmail(account));
await runTask(account);
expect(stripeHandlers.spies.updatedCustomers).toHaveLength(0);
});
test('does nothing for a customer mirrored as deleted', async () => {
const account = await createCustomerAccount('cus_sync_deleted');
await getBillingRepository().customers.upsertFromStripe(
{id: 'cus_sync_deleted', object: 'customer', deleted: true} as Stripe.DeletedCustomer,
{knownUserId: BigInt(account.userId)},
);
await runTask(account);
expect(stripeHandlers.spies.updatedCustomers).toHaveLength(0);
});
test('completes without retrying when Stripe no longer has the customer', async () => {
server.use(...createStripeApiHandlers({customerShouldFail: true}).handlers);
const account = await createCustomerAccount('cus_sync_missing');
await mirrorCustomer('cus_sync_missing', account, STALE_EMAIL);
await expect(runTask(account)).resolves.toBeUndefined();
expect((await getBillingRepository().customers.findById('cus_sync_missing'))?.email).toBe(STALE_EMAIL);
});
test('does nothing for users without a Stripe customer', async () => {
const account = await createTestAccount(harness);
await runTask(account);
expect(stripeHandlers.spies.updatedCustomers).toHaveLength(0);
});
});
describe('payment reconciliation', () => {
test('heals a customer email that drifted before the account email changed', async () => {
setWorkerDependenciesForTest({
userRepository: getUserRepository(),
paymentRepository: new PaymentRepository(),
stripe: getStripeClient(),
});
const account = await createCustomerAccount('cus_sync_reconcile');
await mirrorCustomer('cus_sync_reconcile', account, STALE_EMAIL);
await reconcileUserPayments({userId: account.userId}, createHelpers());
expect(emailUpdatesFor('cus_sync_reconcile')).toEqual([await currentEmail(account)]);
});
test('leaves a customer alone when the email already matches', async () => {
setWorkerDependenciesForTest({
userRepository: getUserRepository(),
paymentRepository: new PaymentRepository(),
stripe: getStripeClient(),
});
const account = await createCustomerAccount('cus_sync_reconcile_current');
await mirrorCustomer('cus_sync_reconcile_current', account, await currentEmail(account));
await reconcileUserPayments({userId: account.userId}, createHelpers());
expect(emailUpdatesFor('cus_sync_reconcile_current')).toHaveLength(0);
});
});
describe('checkout', () => {
test('updates a stale customer email before opening checkout with that customer', async () => {
const account = await createCustomerAccount('cus_checkout_stale');
await mirrorCustomer('cus_checkout_stale', account, STALE_EMAIL);
await createBuilder(harness, account.token)
.post('/stripe/checkout/subscription')
.body({price_id: MOCK_PRICES.monthlyUsd})
.expect(HTTP_STATUS.OK)
.execute();
expect(emailUpdatesFor('cus_checkout_stale')).toEqual([await currentEmail(account)]);
expect(stripeHandlers.spies.createdCheckoutSessions).toHaveLength(1);
expect(stripeHandlers.spies.createdCheckoutSessions[0]?.customer).toBe('cus_checkout_stale');
});
test('leaves the customer alone when its email already matches', async () => {
const account = await createCustomerAccount('cus_checkout_current');
await mirrorCustomer('cus_checkout_current', account, await currentEmail(account));
await createBuilder(harness, account.token)
.post('/stripe/checkout/subscription')
.body({price_id: MOCK_PRICES.monthlyUsd})
.expect(HTTP_STATUS.OK)
.execute();
expect(emailUpdatesFor('cus_checkout_current')).toHaveLength(0);
expect(stripeHandlers.spies.createdCheckoutSessions[0]?.customer).toBe('cus_checkout_current');
});
test('still opens checkout when the email update fails', async () => {
server.use(
...createStripeApiHandlers({prices: MOCK_PRICE_SEEDS, subscriptionsListEmpty: true, customerShouldFail: true})
.handlers,
);
const account = await createCustomerAccount('cus_checkout_update_fails');
await mirrorCustomer('cus_checkout_update_fails', account, STALE_EMAIL);
await createBuilder(harness, account.token)
.post('/stripe/checkout/subscription')
.body({price_id: MOCK_PRICES.monthlyUsd})
.expect(HTTP_STATUS.OK)
.execute();
});
});
});
@@ -1565,6 +1565,19 @@ export function createStripeApiHandlers(config: StripeApiMockConfig = {}): Strip
const formData = await request.formData();
const updateParams = parseFormDataToObject(formData);
spies.updatedCustomers.push({id: customerId, params: updateParams});
if (config.customerShouldFail) {
return HttpResponse.json(
{
error: {
type: 'invalid_request_error',
message: 'No such customer',
code: 'resource_missing',
param: 'id',
},
},
{status: 404},
);
}
const customer = getCustomer(customerId);
const invoiceSettingsUpdate =
updateParams.invoice_settings && typeof updateParams.invoice_settings === 'object'
@@ -1572,6 +1585,7 @@ export function createStripeApiHandlers(config: StripeApiMockConfig = {}): Strip
: null;
const updatedCustomer: MockStripeCustomer = {
...customer,
email: typeof updateParams.email === 'string' ? updateParams.email : customer.email,
invoice_settings: invoiceSettingsUpdate
? {
...customer.invoice_settings,
@@ -5,6 +5,7 @@ import type {ApiContext} from '@app/api/ApiContext';
import * as AuthPassword from '@app/api/auth/AuthPassword';
import {assertEmailNotBlocklisted} from '@app/api/auth/EmailBlocklist';
import type {User} from '@app/api/models/User';
import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer';
import type {EmailChangeRepository} from '@app/api/user/repositories/auth/EmailChangeRepository';
import {
assertChangeCooldown,
@@ -300,7 +301,7 @@ export class EmailChangeService {
}
async verifyBouncedNew(user: User, ticket: string, code: string): Promise<User> {
const {users} = this.apiContext.services;
const {users, worker} = this.apiContext.services;
this.ensureBouncedEmailRecoveryAllowed(user);
const row = await getActiveChangeTicketForUser(this.repo, ticket, user.id);
if (row.require_original || !row.original_proof) {
@@ -314,6 +315,7 @@ export class EmailChangeService {
user.toRow(),
);
await this.deleteToken(emailToken);
await enqueueStripeCustomerEmailSync(worker, user, updatedUser);
return updatedUser;
}
@@ -14,6 +14,7 @@ import {Logger} from '@app/api/Logger';
import type {LimitConfigService} from '@app/api/limits/LimitConfigService';
import type {AuthSession} from '@app/api/models/AuthSession';
import type {User} from '@app/api/models/User';
import {enqueueStripeCustomerEmailSync} from '@app/api/stripe/StripeCustomer';
import type {IUserAccountRepository} from '@app/api/user/repositories/IUserAccountRepository';
import type {IUserChannelRepository} from '@app/api/user/repositories/IUserChannelRepository';
import type {IUserRelationshipRepository} from '@app/api/user/repositories/IUserRelationshipRepository';
@@ -182,6 +183,7 @@ export class UserAccountService {
reason: 'user_requested',
actorUserId: user.id,
}),
() => enqueueStripeCustomerEmailSync(this.apiContext.services.worker, user, updatedUser),
async () => {
try {
await this.profileService.commitAssetChanges(profileResult);
@@ -49,6 +49,7 @@ const LANE_CONFIG = {
'harvestUserData',
'batchGuildAuditLogMessageDeletes',
'reconcileUserPayments',
'syncStripeCustomerEmail',
'processAppStoreNotification',
'processGooglePlayNotification',
'refreshStorePurchase',
@@ -49,6 +49,7 @@ import syncCrosspostCopies from '@app/api/worker/tasks/SyncCrosspostCopies';
import syncCrosspostedMessage from '@app/api/worker/tasks/SyncCrosspostedMessage';
import syncDiscoveryIndex from '@app/api/worker/tasks/SyncDiscoveryIndex';
import syncFileShaBlocklists from '@app/api/worker/tasks/SyncFileShaBlocklists';
import syncStripeCustomerEmail from '@app/api/worker/tasks/SyncStripeCustomerEmail';
import syncUrlBlocklists from '@app/api/worker/tasks/SyncUrlBlocklists';
import userProcessPendingDeletion from '@app/api/worker/tasks/UserProcessPendingDeletion';
import userProcessPendingDeletions from '@app/api/worker/tasks/UserProcessPendingDeletions';
@@ -100,6 +101,7 @@ export const workerTasks: Record<WorkerTaskName, WorkerTaskHandler> = {
refreshSearchIndex,
removeChannelFollowers,
sendSystemDm,
syncStripeCustomerEmail,
syncFileShaBlocklists,
syncUrlBlocklists,
syncDiscoveryIndex,
@@ -6,6 +6,7 @@ import {mapGiftDurationMonthsToFields} from '@app/api/models/GiftCode';
import type {Payment} from '@app/api/models/Payment';
import type {User} from '@app/api/models/User';
import {getProductRegistry} from '@app/api/stripe/ProductRegistry';
import {syncStripeCustomerEmail} from '@app/api/stripe/StripeCustomer';
import {extractId} from '@app/api/stripe/StripeUtils';
import type {PaymentRepository} from '@app/api/user/repositories/PaymentRepository';
import {mapUserToPrivateResponse} from '@app/api/user/UserMappers';
@@ -332,6 +333,11 @@ const reconcileUserPayments: WorkerTaskHandler = async (payload, helpers) => {
if (!user.stripeCustomerId) {
return;
}
try {
await syncStripeCustomerEmail(stripe, user);
} catch (error) {
Logger.warn({userId: userIdStr, error}, 'Failed to sync Stripe customer email during payment reconciliation');
}
const payments = await paymentRepository.findPaymentsByUserId(userId);
let reconciledCount = 0;
for (const payment of payments) {
@@ -0,0 +1,37 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {createUserID} from '@app/api/BrandedTypes';
import {isStripeResourceMissingError, syncStripeCustomerEmail} from '@app/api/stripe/StripeCustomer';
import {getWorkerDependencies} from '@app/api/worker/WorkerContext';
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
const syncStripeCustomerEmailTask: WorkerTaskHandler = async (payload, helpers) => {
const {userRepository, stripe} = getWorkerDependencies();
if (!stripe) {
helpers.logger.debug('Stripe is disabled, skipping customer email sync');
return;
}
const userIdStr = payload.userId as string;
if (!userIdStr) {
helpers.logger.warn({payload}, 'Stripe customer email sync task missing userId');
return;
}
const user = await userRepository.findUnique(createUserID(BigInt(userIdStr)));
if (!user) {
return;
}
try {
await syncStripeCustomerEmail(stripe, user);
} catch (error) {
if (isStripeResourceMissingError(error)) {
helpers.logger.warn(
{userId: userIdStr, customerId: user.stripeCustomerId},
'Stripe customer no longer exists, skipping email sync',
);
return;
}
throw error;
}
};
export default syncStripeCustomerEmailTask;
@@ -69,6 +69,7 @@ import {shouldUseKeyboardShortcutsOverlayFallbackFromEvent} from '@app/features/
import {jsKeyToUiohookKeycode} from '@app/features/input/utils/UiohookKeycodes';
import type {Message} from '@app/features/messaging/models/MessagingMessage';
import MessageFocus from '@app/features/messaging/state/MessageFocus';
import {getSortedDmChannels} from '@app/features/messaging/utils/DmChannelUtils';
import * as NavigationCommands from '@app/features/navigation/commands/NavigationCommands';
import Navigation from '@app/features/navigation/state/Navigation';
import SelectedChannel from '@app/features/navigation/state/SelectedChannel';
@@ -374,7 +375,7 @@ class KeybindManager {
}
private cycleDirectMessageContext(direction: 1 | -1): void {
const dmChannels = Channels.dmChannels;
const dmChannels = getSortedDmChannels(Channels.dmChannels, Authentication.currentUserId);
const slotCount = dmChannels.length + 1;
const currentChannelId = this.currentChannelId;
const currentIndex = currentChannelId ? dmChannels.findIndex((channel) => channel.id === currentChannelId) + 1 : 0;
@@ -926,6 +926,7 @@ export async function crosspost(i18n: I18n, channelId: string, messageId: string
export function deleteLocal(channelId: string, messageId: string): void {
logger.debug(`Deleting message ${messageId} locally in channel ${channelId}`);
Messages.handleMessageDelete({id: messageId, channelId});
MessageReply.handleMessageDelete(channelId, messageId);
}
export function revealMessage(channelId: string, messageId: string | null): void {
@@ -3,6 +3,7 @@
import ChannelPins from '@app/features/channel/state/ChannelPins';
import type {GatewayHandlerContext} from '@app/features/gateway/events/EventRouter';
import MessageReferences from '@app/features/messaging/state/MessageReferences';
import MessageReply from '@app/features/messaging/state/MessageReply';
import Messages from '@app/features/messaging/state/MessagingMessages';
import SavedMessages from '@app/features/messaging/state/SavedMessages';
import MentionFeed from '@app/features/notification/state/MentionFeed';
@@ -20,6 +21,7 @@ export function handleMessageDelete(data: MessageDeletePayload, _context: Gatewa
ChannelPins.handleMessageDelete(data.channel_id, data.id);
Messages.handleMessageDelete({channelId: data.channel_id, id: data.id});
MessageReferences.handleMessageDelete(data.channel_id, data.id);
MessageReply.handleMessageDelete(data.channel_id, data.id);
ReadStates.handleMessageDelete({channelId: data.channel_id});
MentionFeed.handleMessageDelete(data.id);
Notification.handleMessageDelete({channelId: data.channel_id});
@@ -2,6 +2,7 @@
import type {GatewayHandlerContext} from '@app/features/gateway/events/EventRouter';
import MessageReferences from '@app/features/messaging/state/MessageReferences';
import MessageReply from '@app/features/messaging/state/MessageReply';
import Messages from '@app/features/messaging/state/MessagingMessages';
import ReadStates from '@app/features/read_state/state/ReadStates';
import Notification from '@app/features/ui/state/Notification';
@@ -14,6 +15,7 @@ interface MessageDeleteBulkPayload {
export function handleMessageDeleteBulk(data: MessageDeleteBulkPayload, _context: GatewayHandlerContext): void {
Messages.handleMessageDeleteBulk({channelId: data.channel_id, ids: data.ids});
MessageReferences.handleMessageDeleteBulk(data.channel_id, data.ids);
MessageReply.handleMessageDeleteBulk(data.channel_id, data.ids);
ReadStates.handleMessageDelete({channelId: data.channel_id});
Notification.handleMessageDelete({channelId: data.channel_id});
}
@@ -1,170 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import * as MessageCommands from '@app/features/messaging/commands/MessageCommands';
import type {Message} from '@app/features/messaging/models/MessagingMessage';
import MessageEdit from '@app/features/messaging/state/MessageEdit';
import MessageFocus from '@app/features/messaging/state/MessageFocus';
import {insertTextAtCursor} from '@app/features/messaging/utils/TextInputEditUtils';
import {ComponentBus} from '@app/features/platform/utils/ComponentBus';
import {canFocusTextarea, safeFocus} from '@app/features/platform/utils/InputFocusManager';
import {isTextInputKeyEvent} from '@app/features/platform/utils/IsTextInputKeyEvent';
import QuickSwitcher from '@app/features/search/state/QuickSwitcher';
import ContextMenu from '@app/features/ui/state/ContextMenu';
import KeyboardMode from '@app/features/ui/state/KeyboardMode';
import MobileLayout from '@app/features/ui/state/MobileLayout';
import {useCallback, useEffect} from 'react';
interface UseTextareaKeyboardOptions {
channelId: string;
isFocused: boolean;
textareaRef: React.RefObject<HTMLTextAreaElement | null>;
value: string;
setValue: React.Dispatch<React.SetStateAction<string>>;
handleTextChange: (newValue: string, previousValue: string) => void;
previousValueRef: React.RefObject<string>;
clearSegments: () => void;
replyingMessage: {
messageId: string;
mentioning: boolean;
} | null;
editingMessage: Message | null;
getLastEditableMessage: () => Message | null;
enabled: boolean;
}
type ArrowUpEditShortcutEvent = Pick<
React.KeyboardEvent | KeyboardEvent,
'key' | 'altKey' | 'ctrlKey' | 'metaKey' | 'shiftKey' | 'defaultPrevented'
>;
export const shouldStartLastMessageEditFromArrowUp = (event: ArrowUpEditShortcutEvent, value: string): boolean => {
if (event.defaultPrevented) return false;
if (event.key !== 'ArrowUp') return false;
if (event.altKey || event.ctrlKey || event.metaKey || event.shiftKey) return false;
return value.length === 0;
};
export const useTextareaKeyboard = ({
channelId,
isFocused,
textareaRef,
value,
setValue,
handleTextChange,
previousValueRef,
clearSegments,
replyingMessage,
editingMessage,
getLastEditableMessage,
enabled,
}: UseTextareaKeyboardOptions) => {
const mobileLayout = MobileLayout;
const editingMessageId = MessageEdit.getEditingMessageId(channelId);
useEffect(() => {
if (!enabled) {
return;
}
const handleKeyDown = (event: KeyboardEvent) => {
const textarea = textareaRef.current;
if (!canFocusTextarea(textarea || undefined)) {
return;
}
if (isFocused) {
return;
}
if (QuickSwitcher.getIsOpen()) {
return;
}
if (ContextMenu.contextMenu) {
return;
}
if (KeyboardMode.keyboardModeEnabled && MessageFocus.focusedMessageId) {
return;
}
if (!isTextInputKeyEvent(event)) {
return;
}
if (!textarea) {
return;
}
if (event.key === 'Dead') {
safeFocus(textarea, true);
return;
}
event.preventDefault();
safeFocus(textarea, true);
const inserted = insertTextAtCursor(textarea, event.key);
if (!inserted) {
setValue((prev) => {
const newValue = prev + event.key;
handleTextChange(newValue, previousValueRef.current ?? '');
return newValue;
});
}
};
window.addEventListener('keydown', handleKeyDown);
return () => {
window.removeEventListener('keydown', handleKeyDown);
};
}, [
editingMessageId,
isFocused,
mobileLayout.enabled,
handleTextChange,
previousValueRef,
textareaRef,
setValue,
enabled,
]);
useEffect(() => {
if (!enabled) {
return;
}
const handleKeyDown = (event: KeyboardEvent) => {
if (event.key === 'Escape') {
const isEditingInline = MessageEdit.getEditingMessageId(channelId) != null;
if (isEditingInline) {
event.preventDefault();
event.stopPropagation();
MessageCommands.stopEdit(channelId);
return;
}
if (editingMessage && mobileLayout.enabled) {
event.preventDefault();
MessageCommands.stopEditMobile(channelId);
setValue('');
clearSegments();
} else if (replyingMessage) {
event.preventDefault();
MessageCommands.stopReply(channelId);
} else {
event.preventDefault();
ComponentBus.dispatch('ESCAPE_PRESSED');
}
}
};
window.addEventListener('keydown', handleKeyDown);
return () => {
window.removeEventListener('keydown', handleKeyDown);
};
}, [channelId, replyingMessage, editingMessage, mobileLayout.enabled, clearSegments, setValue, enabled]);
const handleArrowUp = useCallback(
(event: React.KeyboardEvent) => {
if (!shouldStartLastMessageEditFromArrowUp(event, value)) {
return;
}
if (KeyboardMode.keyboardModeEnabled) {
event.preventDefault();
ComponentBus.dispatch('FOCUS_BOTTOMMOST_MESSAGE', {channelId});
return;
}
const message = getLastEditableMessage();
if (!message) {
return;
}
event.preventDefault();
MessageCommands.startEdit(channelId, message.id, message.content);
},
[channelId, value, getLastEditableMessage],
);
return {handleArrowUp};
};
@@ -13,7 +13,6 @@ import {emojiEquals} from '@app/features/messaging/utils/ReactionUtils';
import Relationships from '@app/features/relationship/state/Relationships';
import * as ThemeUtils from '@app/features/theme/utils/ThemeUtils';
import {User} from '@app/features/user/models/User';
import UserGuildSettings from '@app/features/user/state/UserGuildSettings';
import Users from '@app/features/user/state/Users';
import {LRUMap} from '@app/lib/list/ListLruMap';
import {MessageFlags, MessageStates, MessageTypes} from '@fluxer/constants/src/ChannelConstants';
@@ -560,15 +559,23 @@ export class Message {
}
}
export const messageMentionsCurrentUser = (message: WireMessage): boolean => {
interface MentionSuppression {
suppressEveryone?: boolean;
suppressRoles?: boolean;
}
export const messageMentionsCurrentUser = (
message: WireMessage,
{suppressEveryone = false, suppressRoles = false}: MentionSuppression = {},
): boolean => {
const channel = Channels.getChannel(message.channel_id);
if (!channel) return false;
if (message.mention_everyone && !UserGuildSettings.isEveryoneMentionSuppressed(channel.guildId ?? null)) return true;
if (message.mention_everyone && !suppressEveryone) return true;
if (message.mentions?.some((user) => user.id === Authentication.currentUserId)) {
return true;
}
if (!channel.guildId) return false;
if (UserGuildSettings.isRoleMentionSuppressed(channel.guildId)) return false;
if (suppressRoles) return false;
const guild = Guilds.getGuild(channel.guildId);
if (!guild) return false;
const guildMember = GuildMembers.getMember(guild.id, Authentication.currentUserId);
@@ -56,6 +56,19 @@ class MessageReply {
delete this.replyingMessageIds[channelId];
}
handleMessageDelete(channelId: string, messageId: string): void {
if (this.replyingMessageIds[channelId]?.messageId === messageId) {
delete this.replyingMessageIds[channelId];
}
}
handleMessageDeleteBulk(channelId: string, messageIds: Array<string>): void {
const current = this.replyingMessageIds[channelId];
if (current && messageIds.includes(current.messageId)) {
delete this.replyingMessageIds[channelId];
}
}
highlightMessage(messageId: string): void {
this.highlightMessageId = messageId;
}
@@ -232,7 +232,12 @@ class Messages {
getLastEditableMessage(channelId: string): Message | undefined {
return this.getMessages(channelId).searchFromNewest((message) => {
return message.isCurrentUserAuthor() && message.state === MessageStates.SENT && message.isUserMessage();
return (
message.isCurrentUserAuthor() &&
message.state === MessageStates.SENT &&
message.isUserMessage() &&
!message.messageSnapshots
);
});
}
@@ -11,6 +11,7 @@ import {
} from '@app/features/notification/utils/MentionFeedFilters';
import Relationships from '@app/features/relationship/state/Relationships';
import {makeSyncedField} from '@app/features/user/state/SyncedField';
import UserGuildSettings from '@app/features/user/state/UserGuildSettings';
import {ChannelTypes} from '@fluxer/constants/src/ChannelConstants';
import {MAX_MESSAGES_PER_CHANNEL} from '@fluxer/constants/src/LimitConstants';
import type {Channel as WireChannel} from '@fluxer/schema/src/domains/channel/ChannelSchemas';
@@ -179,14 +180,19 @@ class MentionFeed {
}
handleMessageCreate(message: WireMessage): void {
if (!messageMentionsCurrentUser(message)) {
const channel = Channels.getChannel(message.channel_id);
if (!channel) return;
const guildId = channel.guildId ?? null;
const mentioned = messageMentionsCurrentUser(message, {
suppressEveryone: UserGuildSettings.isEveryoneMentionSuppressed(guildId),
suppressRoles: UserGuildSettings.isRoleMentionSuppressed(guildId),
});
if (!mentioned) {
return;
}
if (Relationships.isBlocked(message.author.id)) {
return;
}
const channel = Channels.getChannel(message.channel_id);
if (!channel) return;
if (!this.isMessageIncludedByFilters(message, channel)) {
return;
}
+154 -2
View File
@@ -1,6 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use std::{env, path::Path};
use std::{env, fs, path::Path};
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum GeoipSourceConfig {
@@ -124,7 +124,59 @@ fn percent_decode(value: &str) -> String {
}
pub fn env_value(name: &str) -> Option<String> {
env::var(name).ok().filter(|value| !value.trim().is_empty())
resolve_env_value(name, |key| env::var(key).ok()).unwrap_or_else(|error| panic!("{error}"))
}
pub fn resolve_env_value<F>(name: &str, get: F) -> Result<Option<String>, String>
where
F: Fn(&str) -> Option<String>,
{
let non_blank = |value: Option<String>| value.filter(|value| !value.trim().is_empty());
let value = non_blank(get(name));
let Some(path) = non_blank(get(&format!("{name}_FILE"))) else {
return Ok(value);
};
if value.is_some() {
return Err(format!("{name} and {name}_FILE are both set, set only one"));
}
let contents = fs::read_to_string(&path)
.map_err(|error| format!("{name}_FILE could not read {path} ({error})"))?;
let contents = contents
.strip_suffix('\n')
.map_or(contents.as_str(), |rest| {
rest.strip_suffix('\r').unwrap_or(rest)
});
Ok(non_blank(Some(contents.to_owned())))
}
pub fn resolve_env_files<I>(vars: I) -> Result<Vec<(String, String)>, String>
where
I: IntoIterator<Item = (String, String)>,
{
let vars: Vec<(String, String)> = vars.into_iter().collect();
let get = |key: &str| {
vars.iter()
.find_map(|(name, value)| (name == key).then(|| value.clone()))
};
let mut resolved = Vec::new();
for (key, _) in &vars {
let Some(name) = key
.strip_suffix("_FILE")
.filter(|name| name.starts_with("FLUXER_"))
else {
continue;
};
if let Some(value) = resolve_env_value(name, get)? {
resolved.push((name.to_owned(), value));
}
}
let rest: Vec<(String, String)> = vars
.iter()
.filter(|(key, _)| !resolved.iter().any(|(name, _)| name == key))
.cloned()
.collect();
resolved.extend(rest);
Ok(resolved)
}
pub fn read_env(name: &str, fallback: &str) -> String {
@@ -763,6 +815,106 @@ mod tests {
assert_eq!(None, port);
}
fn secret_file(dir: &tempfile::TempDir, name: &str, contents: &str) -> String {
let path = dir.path().join(name);
fs::write(&path, contents).expect("write secret file");
path.to_string_lossy().into_owned()
}
fn pairs_reader(pairs: Vec<(String, String)>) -> impl Fn(&str) -> Option<String> {
move |key| {
pairs
.iter()
.find_map(|(name, value)| (name == key).then(|| value.clone()))
}
}
#[test]
fn env_value_reads_name_file_when_name_is_blank() {
let dir = tempfile::tempdir().expect("tempdir");
let path = secret_file(&dir, "secret", "from-file\r\n");
let unset = pairs_reader(vec![("X_FILE".to_owned(), path.clone())]);
let blank = pairs_reader(vec![
("X".to_owned(), " ".to_owned()),
("X_FILE".to_owned(), path),
]);
assert_eq!(
Ok(Some("from-file".to_owned())),
resolve_env_value("X", unset)
);
assert_eq!(
Ok(Some("from-file".to_owned())),
resolve_env_value("X", blank)
);
}
#[test]
fn env_value_trims_only_one_trailing_newline() {
let dir = tempfile::tempdir().expect("tempdir");
let pem = secret_file(&dir, "pem", "-----BEGIN-----\nabc\n-----END-----\n\n");
let empty = secret_file(&dir, "empty", "\n");
assert_eq!(
Ok(Some("-----BEGIN-----\nabc\n-----END-----\n".to_owned())),
resolve_env_value("X", pairs_reader(vec![("X_FILE".to_owned(), pem)]))
);
assert_eq!(
Ok(None),
resolve_env_value("X", pairs_reader(vec![("X_FILE".to_owned(), empty)]))
);
}
#[test]
fn env_value_rejects_name_and_name_file_together() {
let reader = pairs_reader(vec![
("X".to_owned(), "direct".to_owned()),
("X_FILE".to_owned(), "/run/secrets/x".to_owned()),
]);
assert_eq!(
Err("X and X_FILE are both set, set only one".to_owned()),
resolve_env_value("X", reader)
);
}
#[test]
fn env_value_names_the_file_when_it_is_missing() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("missing").to_string_lossy().into_owned();
let error = resolve_env_value("X", pairs_reader(vec![("X_FILE".to_owned(), path.clone())]))
.expect_err("missing file fails");
assert!(error.starts_with(&format!("X_FILE could not read {path} (")));
}
#[test]
fn resolve_env_files_replaces_fluxer_names_with_file_values() {
let dir = tempfile::tempdir().expect("tempdir");
let path = secret_file(&dir, "secret", "from-file\n");
let missing = dir.path().join("missing").to_string_lossy().into_owned();
let resolved = resolve_env_files(vec![
("FLUXER_X".to_owned(), String::new()),
("FLUXER_X_FILE".to_owned(), path),
("FLUXER_Y".to_owned(), "plain".to_owned()),
("FLUXER_Y_FILE".to_owned(), String::new()),
("SSL_CERT_FILE".to_owned(), missing),
])
.expect("resolves");
let values = |key: &str| {
resolved
.iter()
.filter(|(name, _)| name == key)
.map(|(_, value)| value.as_str())
.collect::<Vec<_>>()
};
assert_eq!(vec!["from-file"], values("FLUXER_X"));
assert_eq!(vec!["plain"], values("FLUXER_Y"));
assert!(
resolve_env_files(vec![
("FLUXER_X".to_owned(), "direct".to_owned()),
("FLUXER_X_FILE".to_owned(), "/run/secrets/x".to_owned()),
])
.is_err()
);
}
#[test]
fn a_blank_value_reads_as_unset() {
unsafe {
@@ -47,7 +47,7 @@
<image type="source" width="2560" height="2227">https://fluxer.app/static/img/screenshots-desktop-2560w.948e9d88f4c9807a.avif</image>
</screenshot>
</screenshots>
<url type="bugtracker">https://github.com/fluxerapp/fluxer/issues</url>
<url type="bugtracker">https://feedback.fluxer.com</url>
<url type="homepage">https://fluxer.app</url>
<url type="donation">https://fluxer.app/donate</url>
<url type="contact">https://fluxer.app/company-information#contact</url>
@@ -47,7 +47,7 @@
<image type="source" width="1920" height="1670">https://fluxer.app/static/img/screenshots-desktop-1920w.4976d2d1f2238c05.avif</image>
</screenshot>
</screenshots>
<url type="bugtracker">https://github.com/fluxerapp/fluxer/issues</url>
<url type="bugtracker">https://feedback.fluxer.com</url>
<url type="homepage">https://canary.fluxer.app</url>
<url type="donation">https://fluxer.app/donate</url>
<url type="vcs-browser">https://github.com/fluxerapp/fluxer.git</url>
+1 -1
View File
@@ -215,7 +215,7 @@ function buildTemplate(): Array<MenuItemConstructorOptions> {
{
label: t('desktop.appMenu.reportIssue'),
click: async () => {
await openExternalDeduped('https://github.com/fluxerapp/fluxer/issues');
await openExternalDeduped('https://feedback.fluxer.com');
},
},
{type: 'separator'},
@@ -2198,6 +2198,7 @@ async function verifyInstallerExecution(installerRoot: string): Promise<Array<st
env: {
...process.env,
PATH: `${stubBin}${path.delimiter}${process.env.PATH ?? ''}`,
FLUXER_INSTALLER_REFRESHED: '1',
...(cwd == null ? {} : {PWD: cwd}),
},
});
@@ -36,6 +36,34 @@ Precedence, highest first:
Use `true` or `false` for booleans, decimal integers for integer settings, and the specified object or array for JSON settings. Defaults and accepted values are listed below. Every service built from this version of the stack files or later reads an empty or blank value the same as an unset one, so an empty value restores the default. Older images do not, which [Match the images to the stack files](/operator/upgrading/#match-the-images-to-the-stack-files) covers.
### Reading a secret from a file
Any setting a Fluxer service reads for its configuration can come from a file instead, such as a Docker secret. Set the variable's name with `_FILE` appended to the file's path, and leave the variable itself empty. The service reads the file at startup and drops one trailing newline. Startup fails, naming both variables, when the variable and its `_FILE` form are both set. It also fails, naming the path, when the file is missing, unreadable or not valid UTF-8. The name to use is the one the service sees, so `POSTGRES_PASSWORD` in `.env` becomes `FLUXER_POSTGRES_PASSWORD_FILE`.
Add the secrets through a local Compose override. `FLUXER_SUDO_MODE_SECRET` is read by `api` and `worker` only:
```yaml
secrets:
sudo_mode_secret:
file: ./secrets/sudo_mode_secret
services:
api:
secrets: [sudo_mode_secret]
environment:
FLUXER_SUDO_MODE_SECRET: ''
FLUXER_SUDO_MODE_SECRET_FILE: /run/secrets/sudo_mode_secret
worker:
secrets: [sudo_mode_secret]
environment:
FLUXER_SUDO_MODE_SECRET: ''
FLUXER_SUDO_MODE_SECRET_FILE: /run/secrets/sudo_mode_secret
```
Compose still checks the required names in `.env` before it applies the override, so keep a placeholder value there. Every Fluxer service gets the same shared settings, so a secret several services read, such as `FLUXER_POSTGRES_PASSWORD`, needs the same two entries on each of them: `api`, `worker`, `gateway`, `media-proxy`, `push`, `admin`, and `snowflakes`, `users`, `gifs`, `messages` and `unfurl` with their `-shard` services. Any service left out keeps the placeholder. The secret file must be readable by the user the container runs as.
The bundled backing containers read their own settings. The `postgres` container takes `POSTGRES_PASSWORD: ''` and `POSTGRES_PASSWORD_FILE` in its own `environment`, and on a fresh volume it creates the database with that password. `seaweedfs-init`, `meilisearch` and `livekit` have no `_FILE` form, so a file-based S3, Meilisearch or LiveKit secret still needs its real value in `.env` for them. `FLUXER_ENV`, `LOG_LEVEL`, `FLUXER_DISABLE_RATE_LIMITS` and the `FLUXER_ERLANG_*` settings other than `FLUXER_ERLANG_COOKIE` are read directly, so set them as plain values.
## Core identity and public address
`FLUXER_DOMAIN` is required. Everything else here is optional.
@@ -2137,13 +2165,15 @@ CORS origins are exactly the app endpoint, the `FLUXER_APP_ORIGIN_ALIASES` origi
| --- | --- | --- |
| postgres-data | Every account, message, and configuration row | Yes |
| seaweedfs-data | Every uploaded file | Yes |
| valkey-data | Deletion queues and shared cache | Yes |
| valkey-data | Pending file deletions, other queues and shared cache | Yes |
| nats-data | Pending and failed background jobs | Yes |
| edge-data | Issued TLS certificates | Optional, a loss only costs a re-issue |
| edge-config | The edge's own state | No |
| meilisearch-data | The search index, rebuildable | No |
| meilisearch-data | The search index, rebuildable from Postgres | Optional |
Losing `valkey-data` can delay scheduled account and bulk-message deletions while their queues are rebuilt. Pending asset deletions and cache purges can be lost, so do not treat this volume as disposable cache.
Losing `valkey-data` can delay scheduled account and bulk-message deletions while their queues are rebuilt. Pending asset deletions and cache purges exist only here and are lost with it, so do not treat this volume as disposable cache.
Losing `meilisearch-data` leaves search empty until it is rebuilt. Nothing rebuilds it automatically, because Postgres still records every channel as indexed. Rebuild each index with [`POST /admin/search/indexes/{index_name}/refreshes`](/admin-api/search-indexes/). The `channel_messages` and `guild_members` indexes rebuild one community at a time and need a `guild_id` for each. A backup of the volume avoids those calls on a large instance.
`nats-data` retains pending jobs in `JOBS` for up to 7 days and failed jobs in `JOBS_DLQ` for up to 30 days. A full jobs stream rejects new work. A full dead-letter stream drops its oldest entries, so investigate failures promptly. If dead-letter storage is unavailable, failed jobs remain in `JOBS` only until they expire. Losing this volume loses queued work, which is not automatically recovered from the database.
@@ -240,7 +240,7 @@ Once email is on, the address goes through a DNS check before the account exists
Sign in to the admin dashboard at `https://chat.example.com/admin` with the account holding the wildcard ACL. The **Instance Config** page has everything the wizard asked, plus registration mode, approvals and integration keys. **Limit Config** holds the instance limits published to clients. **Voice Regions** and **Voice Servers** come seeded, so voice needs no setup there.
The desktop client opens the hosted web app for its release channel, so reach your own instance in a browser. [Issue #1088](https://github.com/fluxerapp/fluxer/issues/1088) tracks that, and is being worked on.
The desktop client opens the hosted web app for its release channel, so reach your own instance in a browser. [This feedback post](https://feedback.fluxer.com/p/641) tracks that, and is being worked on.
[Configuration](/operator/configuration/) lists every runtime setting and the environment variable each one overrides. [Deployment availability](/http-api/deployment-availability/) lists the routes that exist only on the hosted deployment.
@@ -522,10 +522,13 @@ LiveKit signalling goes through `/livekit/*` like everything else. WebRTC media
Both ports are published directly by the stack and must reach the host. A proxy or tunnel in front of `443` does nothing for them.
Hosting LiveKit on a hostname other than `FLUXER_DOMAIN` means widening the Content-Security-Policy the web app runs under. The line goes in `.env`, beside `FLUXER_DOMAIN`:
Signalling can use a hostname other than `FLUXER_DOMAIN`. Two lines go in `.env`, beside `FLUXER_DOMAIN`. The first points clients at the new hostname, and the second lets the web app's Content-Security-Policy connect to it:
```ini
FLUXER_CSP_EXTRA_CONNECT_SRC=wss://livekit.example.com:7881
FLUXER_LIVEKIT_URL=wss://livekit.example.com/livekit
FLUXER_CSP_EXTRA_CONNECT_SRC=wss://livekit.example.com
```
`app-proxy` builds the policy and reads its environment at container start, so apply the change with `docker compose up -d app-proxy`. `docker compose restart app-proxy` reuses the existing container with its old environment. [Content Security Policy](/operator/configuration/#content-security-policy) has the other variables.
The policy entry is the origin alone, with no port or path. The hostname needs a valid TLS certificate and must reach the `/livekit` route on the edge with WebSocket upgrades. The media ports above never belong in the policy, because WebRTC media is not governed by `connect-src`. A separate hostname is rarely needed, since media goes to the address LiveKit advertises, not to a hostname.
`api` writes `FLUXER_LIVEKIT_URL` into the default voice server at start, and `app-proxy` builds the policy at container start, so apply the change with `docker compose up -d api app-proxy`. `docker compose restart` reuses the existing containers with their old environment. [Content Security Policy](/operator/configuration/#content-security-policy) has the other variables.
@@ -128,6 +128,8 @@ Removing the line from `.env` also lets the upgrade run, and it is the wrong fix
Keep the installer beside the stack files. [Get started](/operator/get-started/#step-4-bring-up-the-instance) has the download and the checksum check for a fresh copy.
`--update` first compares the installer against the digest at `https://fluxer.dev/install.sh.sha256`. When the copy differs, it downloads the current one, checks it against that digest, and asks before it replaces the copy and runs it with the same options. Answering no, `--non-interactive`, or a run without a terminal stops before anything changes. On Windows the same check reads `https://fluxer.dev/install.ps1.sha256`.
See the plan first:
```bash
+80 -1
View File
@@ -59,11 +59,17 @@ param(
[string[]]$Rest = @()
)
$FluxerScriptArguments = @{}
foreach ($entry in $PSBoundParameters.GetEnumerator()) {
$FluxerScriptArguments[$entry.Key] = $entry.Value
}
Set-StrictMode -Version Latest
$ErrorActionPreference = 'Stop'
$ProgressPreference = 'SilentlyContinue'
$FluxerRawBase = 'https://raw.githubusercontent.com/fluxerapp/fluxer'
$FluxerInstallerUrl = 'https://fluxer.dev/install.ps1'
$FluxerStackPath = 'deploy/self-hosting'
$FluxerHealthPath = '/_health'
$FluxerInitService = 'seaweedfs-init'
@@ -1980,6 +1986,71 @@ function Assert-FluxerComposeFiles([string]$TargetDir, [string]$EnvPath) {
}
}
function Get-FluxerSha256([string]$Path) {
return (Get-FileHash -LiteralPath $Path -Algorithm SHA256).Hash.ToLowerInvariant()
}
function Update-FluxerInstaller {
if ($env:FLUXER_INSTALLER_REFRESHED) {
return
}
$self = $PSCommandPath
if (-not $self) {
return
}
$selfDir = Split-Path -Parent $self
if ($DryRun) {
$staging = New-FluxerStagingDirectory ([System.IO.Path]::GetTempPath())
} else {
$staging = New-FluxerStagingDirectory $selfDir
}
try {
$digestPath = Join-Path $staging 'install.ps1.sha256'
try {
Invoke-WebRequest -Uri "$FluxerInstallerUrl.sha256" -OutFile $digestPath -UseBasicParsing -MaximumRedirection 5 -TimeoutSec 60
} catch {
Write-FluxerLine "Could not reach $FluxerInstallerUrl.sha256, so this run goes on with $self."
return
}
$published = ([System.IO.File]::ReadAllText($digestPath).Trim() -split '\s+')[0].ToLowerInvariant()
if ($published -notmatch '^[0-9a-f]{64}$') {
Stop-Fluxer "$FluxerInstallerUrl.sha256 holds no sha256 digest. Nothing was changed." $FluxerExitDownload
}
if ((Get-FluxerSha256 $self) -eq $published) {
return
}
$fresh = Join-Path $staging 'install.ps1'
try {
Invoke-WebRequest -Uri $FluxerInstallerUrl -OutFile $fresh -UseBasicParsing -MaximumRedirection 5 -TimeoutSec 120
} catch {
Stop-Fluxer "Download failed for $FluxerInstallerUrl. Nothing was changed." $FluxerExitDownload
}
if ((Get-FluxerSha256 $fresh) -ne $published) {
Stop-Fluxer "$FluxerInstallerUrl does not match the digest in $FluxerInstallerUrl.sha256. Nothing was changed." $FluxerExitDownload
}
Write-FluxerLine "$self differs from the installer $FluxerInstallerUrl serves. The stack files an upgrade downloads can require .env keys that only the current installer writes."
if ($DryRun) {
Write-FluxerLine 'The run asks to replace it with the current installer before it changes anything. The plan below is the one this copy would follow.'
return
}
if ($NonInteractive -or [Console]::IsInputRedirected) {
Stop-Fluxer "Nothing was changed. Download the current installer and run it:`n Invoke-WebRequest -Uri $FluxerInstallerUrl -OutFile install.ps1 -UseBasicParsing" $FluxerExitRefused
}
$answer = Read-Host -Prompt "Replace $self with the current installer and run that? [y/N]"
if ($null -eq $answer -or $answer.Trim() -notmatch '^(y|yes)$') {
Stop-Fluxer "Kept $self. Nothing was changed. Read the current installer at $FluxerInstallerUrl and run it once it is in place." $FluxerExitRefused
}
Move-Item -LiteralPath $fresh -Destination $self -Force
} finally {
Remove-FluxerStagingDirectory $staging
}
Write-FluxerLine "Replaced $self. Running it."
$env:FLUXER_INSTALLER_REFRESHED = '1'
$global:LASTEXITCODE = 0
& $self @FluxerScriptArguments
exit $LASTEXITCODE
}
function Invoke-FluxerInstall {
if ($Help) {
Show-FluxerUsage
@@ -2022,6 +2093,10 @@ function Invoke-FluxerInstall {
Invoke-FluxerPreflight
if ($Update) {
Update-FluxerInstaller
}
$targetPath = $Dir
$adoptedCwd = $false
$fellBack = $false
@@ -2192,10 +2267,14 @@ function Invoke-FluxerInstall {
exit 0
}
Write-FluxerLine 'Pulling images. The first start pulls eighteen of them, which takes several minutes.'
if ((Invoke-FluxerDocker @('compose', 'pull')) -ne 0) {
Stop-Fluxer 'docker compose pull failed. Nothing was started.' $FluxerExitDownload
}
if ((Invoke-FluxerDocker @('compose', 'up', '-d')) -ne 0) {
Stop-Fluxer 'docker compose up -d failed.' $FluxerExitUnhealthy
}
Wait-FluxerStack 'Waiting for the stack to report healthy. The first start pulls images and takes several minutes.'
Wait-FluxerStack 'Waiting for the stack to report healthy.'
$readyOrigin = Get-FluxerPublicOrigin $envPath
if ($readyOrigin.Length -eq 0) {
$readyOrigin = "https://$domainValue"
+90 -1
View File
@@ -51,6 +51,7 @@ LC_ALL=C
export LC_ALL
FLUXER_RAW_BASE='https://raw.githubusercontent.com/fluxerapp/fluxer'
FLUXER_INSTALLER_URL='https://fluxer.dev/install.sh'
FLUXER_STACK_PATH='deploy/self-hosting'
FLUXER_MIN_ENGINE='24.0.0'
# Podman numbers its releases on its own scale, so the Docker Engine floor says
@@ -305,6 +306,18 @@ opt_no_volume_compression=0
opt_skip_backup=0
opt_allow_root=0
fluxer_self=''
if [ -f "$0" ]; then
case $0 in
/*) fluxer_self=$0 ;;
*) fluxer_self="$(pwd)/$0" ;;
esac
fi
fluxer_self_args=''
for fluxer_arg in "$@"; do
fluxer_self_args="$fluxer_self_args '$(printf '%s' "$fluxer_arg" | sed "s/'/'\\\\''/g")'"
done
while [ $# -gt 0 ]; do
case $1 in
--domain)
@@ -651,6 +664,74 @@ fluxer_prompt() {
return 1
}
fluxer_sha256() {
openssl dgst -sha256 < "$1" | awk '{print $NF}'
}
fluxer_drop_scratch() {
fluxer_cleanup
fluxer_scratch=''
}
fluxer_refresh_installer() {
[ -z "${FLUXER_INSTALLER_REFRESHED:-}" ] || return 0
[ -n "$fluxer_self" ] || return 0
fluxer_self_dir=$(dirname "$fluxer_self")
if [ "$opt_dry_run" -eq 1 ]; then
fluxer_open_scratch "${TMPDIR:-/tmp}"
elif [ -w "$fluxer_self_dir" ]; then
fluxer_open_scratch "$fluxer_self_dir"
else
fluxer_say "$fluxer_self_dir is not writable, so this run cannot check $fluxer_self against $FLUXER_INSTALLER_URL and goes on with it."
return 0
fi
if ! curl -fsSL --proto '=https' --tlsv1.2 -o "$fluxer_scratch/install.sh.sha256" "$FLUXER_INSTALLER_URL.sha256"; then
fluxer_say "Could not reach $FLUXER_INSTALLER_URL.sha256, so this run goes on with $fluxer_self."
fluxer_drop_scratch
return 0
fi
fluxer_published=$(awk 'NR == 1 {print $1}' "$fluxer_scratch/install.sh.sha256")
case $fluxer_published in
''|*[!0-9a-f]*) fluxer_fail 4 "$FLUXER_INSTALLER_URL.sha256 holds no sha256 digest. Nothing was changed." ;;
esac
if [ "$(fluxer_sha256 "$fluxer_self")" = "$fluxer_published" ]; then
fluxer_drop_scratch
return 0
fi
if ! curl -fsSL --proto '=https' --tlsv1.2 -o "$fluxer_scratch/install.sh" "$FLUXER_INSTALLER_URL"; then
fluxer_fail 4 "Download failed for $FLUXER_INSTALLER_URL. Nothing was changed."
fi
if [ "$(fluxer_sha256 "$fluxer_scratch/install.sh")" != "$fluxer_published" ]; then
fluxer_fail 4 "$FLUXER_INSTALLER_URL does not match the digest in $FLUXER_INSTALLER_URL.sha256. Nothing was changed."
fi
fluxer_say "$fluxer_self differs from the installer $FLUXER_INSTALLER_URL serves. The stack files an upgrade downloads can require .env keys that only the current installer writes."
if [ "$opt_dry_run" -eq 1 ]; then
fluxer_say 'The run asks to replace it with the current installer before it changes anything. The plan below is the one this copy would follow.'
fluxer_drop_scratch
return 0
fi
if [ "$opt_non_interactive" -eq 1 ] || [ ! -t 0 ]; then
fluxer_fail 3 "Nothing was changed. Download the current installer and run it:
curl -fsSLO $FLUXER_INSTALLER_URL"
fi
printf 'Replace %s with the current installer and run that? [y/N] ' "$fluxer_self" >&2
fluxer_answer=''
read -r fluxer_answer || true
case $fluxer_answer in
y|Y|yes|Yes|YES) ;;
*) fluxer_fail 3 "Kept $fluxer_self. Nothing was changed. Read the current installer at $FLUXER_INSTALLER_URL and run it once it is in place." ;;
esac
if [ -x "$fluxer_self" ]; then
chmod +x "$fluxer_scratch/install.sh"
fi
mv "$fluxer_scratch/install.sh" "$fluxer_self"
fluxer_drop_scratch
fluxer_say "Replaced $fluxer_self. Running it."
FLUXER_INSTALLER_REFRESHED=1
export FLUXER_INSTALLER_REFRESHED
eval "exec sh \"\$fluxer_self\" $fluxer_self_args"
}
fluxer_resolve_values() {
if [ "$opt_update" -eq 1 ] || [ "$opt_rollback" -eq 1 ]; then
return 0
@@ -2100,6 +2181,10 @@ fluxer_preflight
fluxer_validate_options
fluxer_resolve_values
if [ "$opt_update" -eq 1 ]; then
fluxer_refresh_installer
fi
if [ "$opt_update" -eq 1 ] || [ "$opt_rollback" -eq 1 ]; then
fluxer_require_instance
fluxer_resolve_ref
@@ -2160,11 +2245,15 @@ if [ "$opt_no_start" -eq 1 ]; then
exit 0
fi
fluxer_say 'Pulling images. The first start pulls eighteen of them, which takes several minutes.'
if ! $fluxer_engine compose pull; then
fluxer_fail 4 "$fluxer_engine compose pull failed in $opt_dir. Nothing was started."
fi
fluxer_say 'Starting the stack.'
if ! $fluxer_engine compose up -d; then
fluxer_fail 6 "$fluxer_engine compose up -d failed in $opt_dir. Read $fluxer_engine compose logs there."
fi
fluxer_say 'Waiting for every service to report ready. This takes several minutes on the first start, which pulls eighteen images.'
fluxer_say 'Waiting for every service to report ready.'
if ! fluxer_wait_ready; then
fluxer_fail 6 "The stack is not ready after $FLUXER_READY_TIMEOUT seconds. $(fluxer_not_ready_detail)
Read $fluxer_engine compose logs in $opt_dir."
@@ -70,6 +70,21 @@ else
rm -f "$node_name_file"
fi
case "${FLUXER_ERLANG_COOKIE_FILE:-}" in
*[![:space:]]*)
case "${FLUXER_ERLANG_COOKIE:-}" in
*[![:space:]]*)
echo 'FLUXER_ERLANG_COOKIE and FLUXER_ERLANG_COOKIE_FILE are both set, set only one.' >&2
exit 1
;;
esac
if ! FLUXER_ERLANG_COOKIE="$(cat -- "$FLUXER_ERLANG_COOKIE_FILE")"; then
echo "FLUXER_ERLANG_COOKIE_FILE could not read $FLUXER_ERLANG_COOKIE_FILE." >&2
exit 1
fi
;;
esac
if [ -z "${FLUXER_ERLANG_COOKIE:-}" ]; then
echo 'FLUXER_ERLANG_COOKIE is required.' >&2
exit 1
@@ -18,4 +18,8 @@ if [ -r "$node_name_file" ]; then
FLUXER_ERLANG_NODE_NAME="$(cat "$node_name_file")"
export FLUXER_ERLANG_NODE_NAME
fi
if [ -r "${FLUXER_ERLANG_COOKIE_FILE:-}" ] && [ -z "$(printf '%s' "${FLUXER_ERLANG_COOKIE:-}" | tr -d '[:space:]')" ]; then
FLUXER_ERLANG_COOKIE="$(cat "$FLUXER_ERLANG_COOKIE_FILE")"
export FLUXER_ERLANG_COOKIE
fi
exec "$script_dir/fluxer_gateway.real" "$@"
@@ -305,11 +305,44 @@ env_string(Name, Default) ->
-spec env_value(string()) -> string() | undefined.
env_value(Name) ->
Value = non_blank_env(Name),
FileName = Name ++ "_FILE",
case non_blank_env(FileName) of
undefined -> Value;
Path when Value =:= undefined -> read_env_file(FileName, Path);
_ -> erlang:error({ambiguous_env, Name, FileName})
end.
-spec non_blank_env(string()) -> string() | undefined.
non_blank_env(Name) ->
case os:getenv(Name) of
false -> undefined;
Value -> non_blank(Value)
end.
-spec read_env_file(string(), string()) -> string() | undefined.
read_env_file(FileName, Path) ->
case file:read_file(Path) of
{ok, Contents} -> env_file_value(FileName, Path, strip_newline(Contents));
{error, Reason} -> erlang:error({unreadable_env_file, FileName, Path, Reason})
end.
-spec env_file_value(string(), string(), binary()) -> string() | undefined.
env_file_value(FileName, Path, Contents) ->
case unicode:characters_to_list(Contents) of
Value when is_list(Value) -> non_blank(Value);
_ -> erlang:error({invalid_env_file, FileName, Path})
end.
-spec strip_newline(binary()) -> binary().
strip_newline(Contents) ->
Size = byte_size(Contents),
case Contents of
<<Rest:(Size - 2)/binary, "\r\n">> -> Rest;
<<Rest:(Size - 1)/binary, "\n">> -> Rest;
_ -> Contents
end.
-spec non_blank(string()) -> string() | undefined.
non_blank(Value) ->
case string:trim(Value) of
@@ -272,6 +272,64 @@ public_endpoints_defaults_test() ->
?assertEqual(undefined, maps:get(media_proxy_endpoint, Config)),
?assertEqual(<<"http://localhost:8088">>, maps:get(static_cdn_endpoint, Config)).
env_value_reads_name_file_test() ->
with_env_file(<<"from-file\r\n">>, fun(Path) ->
with_envs(
[{"FLUXER_GATEWAY_TEST_SECRET", ""}, {"FLUXER_GATEWAY_TEST_SECRET_FILE", Path}],
fun() ->
?assertEqual(
"from-file", fluxer_gateway_config:env_value("FLUXER_GATEWAY_TEST_SECRET")
)
end
)
end).
env_value_trims_only_one_newline_test() ->
with_env_file(<<"line1\nline2\n\n">>, fun(Path) ->
with_env("FLUXER_GATEWAY_TEST_SECRET_FILE", Path, fun() ->
?assertEqual(
"line1\nline2\n", fluxer_gateway_config:env_value("FLUXER_GATEWAY_TEST_SECRET")
)
end)
end).
env_value_rejects_name_and_name_file_test() ->
with_envs(
[
{"FLUXER_GATEWAY_TEST_SECRET", "direct"},
{"FLUXER_GATEWAY_TEST_SECRET_FILE", "/run/secrets/x"}
],
fun() ->
?assertError(
{ambiguous_env, "FLUXER_GATEWAY_TEST_SECRET",
"FLUXER_GATEWAY_TEST_SECRET_FILE"},
fluxer_gateway_config:env_value("FLUXER_GATEWAY_TEST_SECRET")
)
end
).
env_value_rejects_missing_name_file_test() ->
Path = "/nonexistent/fluxer-gateway-test-secret",
with_env("FLUXER_GATEWAY_TEST_SECRET_FILE", Path, fun() ->
?assertError(
{unreadable_env_file, "FLUXER_GATEWAY_TEST_SECRET_FILE", Path, enoent},
fluxer_gateway_config:env_value("FLUXER_GATEWAY_TEST_SECRET")
)
end).
with_env_file(Contents, Fun) ->
Path = filename:join(
filename:basedir(user_cache, "fluxer_gateway_tests"),
integer_to_list(erlang:unique_integer([positive]))
),
ok = filelib:ensure_dir(Path),
ok = file:write_file(Path, Contents),
try
Fun(Path)
after
file:delete(Path)
end.
with_envs([], Fun) ->
Fun();
with_envs([{Name, Value} | Rest], Fun) ->
+3 -1
View File
@@ -48,7 +48,9 @@ pub enum StorageBackendArg {
}
pub fn load_config(args: &Args) -> anyhow::Result<Config> {
load_config_from_iter(args, std::env::vars())
let vars =
fluxer_common::config::resolve_env_files(std::env::vars()).map_err(anyhow::Error::msg)?;
load_config_from_iter(args, vars)
}
pub fn load_config_from_iter<I, K, V>(args: &Args, vars: I) -> anyhow::Result<Config>
+1 -5
View File
@@ -13,7 +13,7 @@ use parse::{
parse_mode_env, parse_policy_mode, parse_storage_backend, parse_u16, parse_u64, parse_usize,
validate_read_endpoint,
};
use std::{env, path::PathBuf};
use std::path::PathBuf;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum StorageBackend {
@@ -144,10 +144,6 @@ pub struct Config {
}
impl Config {
pub fn load_from_env() -> anyhow::Result<Self> {
Self::load_from_iter(env::vars())
}
pub fn load_from_iter<I, K, V>(vars: I) -> anyhow::Result<Self>
where
I: IntoIterator<Item = (K, V)>,
+3 -1
View File
@@ -27,7 +27,9 @@ pub enum Command {
}
pub fn load_config(args: &Args) -> anyhow::Result<Config> {
load_config_from_iter(args, std::env::vars())
let vars =
fluxer_svc::config::resolve_env_files(std::env::vars()).map_err(anyhow::Error::msg)?;
load_config_from_iter(args, vars)
}
fn load_config_from_iter<I, K, V>(args: &Args, vars: I) -> anyhow::Result<Config>
+99 -2
View File
@@ -1,6 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use std::env;
use std::fs;
use std::net::{IpAddr, SocketAddr};
use std::time::Duration;
@@ -52,7 +53,7 @@ pub enum DatabaseBackend {
impl ServiceConfig {
pub fn from_env() -> anyhow::Result<Self> {
Self::from_env_reader(|name| env::var(name).ok())
Self::from_env_reader(env_var)
}
fn from_env_reader<F>(get: F) -> anyhow::Result<Self>
@@ -203,7 +204,63 @@ fn default_max_concurrent_requests(service_name: &str) -> usize {
}
pub fn optional_env(name: &str) -> Option<String> {
optional_from(&|key| env::var(key).ok(), name)
env_var(name)
}
fn env_var(name: &str) -> Option<String> {
resolve_env_value(name, |key| env::var(key).ok()).unwrap_or_else(|error| panic!("{error}"))
}
fn resolve_env_value<F>(name: &str, get: F) -> Result<Option<String>, String>
where
F: Fn(&str) -> Option<String>,
{
let non_blank = |value: Option<String>| value.filter(|value| !value.trim().is_empty());
let value = non_blank(get(name));
let Some(path) = non_blank(get(&format!("{name}_FILE"))) else {
return Ok(value);
};
if value.is_some() {
return Err(format!("{name} and {name}_FILE are both set, set only one"));
}
let contents = fs::read_to_string(&path)
.map_err(|error| format!("{name}_FILE could not read {path} ({error})"))?;
let contents = contents
.strip_suffix('\n')
.map_or(contents.as_str(), |rest| {
rest.strip_suffix('\r').unwrap_or(rest)
});
Ok(non_blank(Some(contents.to_owned())))
}
pub fn resolve_env_files<I>(vars: I) -> Result<Vec<(String, String)>, String>
where
I: IntoIterator<Item = (String, String)>,
{
let vars: Vec<(String, String)> = vars.into_iter().collect();
let get = |key: &str| {
vars.iter()
.find_map(|(name, value)| (name == key).then(|| value.clone()))
};
let mut resolved = Vec::new();
for (key, _) in &vars {
let Some(name) = key
.strip_suffix("_FILE")
.filter(|name| name.starts_with("FLUXER_"))
else {
continue;
};
if let Some(value) = resolve_env_value(name, get)? {
resolved.push((name.to_owned(), value));
}
}
let rest: Vec<(String, String)> = vars
.iter()
.filter(|(key, _)| !resolved.iter().any(|(name, _)| name == key))
.cloned()
.collect();
resolved.extend(rest);
Ok(resolved)
}
fn optional_from<F>(get: &F, name: &str) -> Option<String>
@@ -276,6 +333,46 @@ mod tests {
.unwrap()
}
#[test]
fn resolves_name_file_entries() {
let dir = env::temp_dir().join(format!("fluxer-svc-env-file-{}", std::process::id()));
fs::create_dir_all(&dir).unwrap();
let path = dir.join("token");
fs::write(&path, "from-file\r\n").unwrap();
let path = path.to_string_lossy().into_owned();
let missing = dir.join("missing").to_string_lossy().into_owned();
let resolved = resolve_env_files(vec![
("FLUXER_NATS_AUTH_TOKEN".to_owned(), String::new()),
("FLUXER_NATS_AUTH_TOKEN_FILE".to_owned(), path.clone()),
])
.unwrap();
let both = resolve_env_files(vec![
("FLUXER_X".to_owned(), "direct".to_owned()),
("FLUXER_X_FILE".to_owned(), path),
]);
let unreadable = resolve_env_files(vec![("FLUXER_X_FILE".to_owned(), missing.clone())]);
let unrelated = resolve_env_files(vec![("SSL_CERT_FILE".to_owned(), missing.clone())]);
fs::remove_dir_all(&dir).unwrap();
assert_eq!(
vec!["from-file"],
resolved
.iter()
.filter(|(name, _)| name == "FLUXER_NATS_AUTH_TOKEN")
.map(|(_, value)| value.as_str())
.collect::<Vec<_>>()
);
assert_eq!(
Err("FLUXER_X and FLUXER_X_FILE are both set, set only one".to_owned()),
both
);
assert!(
unreadable
.unwrap_err()
.starts_with(&format!("FLUXER_X_FILE could not read {missing} ("))
);
assert_eq!(Ok(vec![("SSL_CERT_FILE".to_owned(), missing)]), unrelated);
}
#[test]
fn extracts_statefulset_ordinal() {
assert_eq!(shard_id_from_pod_name("my-service-0"), Some(0));
+1 -6
View File
@@ -4,7 +4,6 @@ use super::{ResolveContext, Resolver, ResolverResult};
use crate::http_fetch;
use crate::media_proxy::{MediaMetadata, embed_media_flags};
use crate::types::{EmbedMedia, EmbedProvider, MessageEmbed};
use fluxer_svc::config::optional_env;
use std::future::Future;
use std::pin::Pin;
use std::time::Duration;
@@ -49,7 +48,7 @@ impl Resolver for KlipyResolver {
ctx: &'a ResolveContext<'_>,
) -> Pin<Box<dyn Future<Output = anyhow::Result<ResolverResult>> + Send + 'a>> {
Box::pin(async move {
let Some(api_key) = ctx.klipy_api_key.clone().or_else(klipy_api_key) else {
let Some(api_key) = ctx.klipy_api_key.clone() else {
return Ok(ResolverResult { embeds: vec![] });
};
let formats = match resolve_media_via_api(ctx, &api_key).await {
@@ -136,10 +135,6 @@ fn klipy_resource(kind: &str) -> &'static str {
}
}
fn klipy_api_key() -> Option<String> {
optional_env("FLUXER_KLIPY_API_KEY").or_else(|| optional_env("KLIPY_API_KEY"))
}
async fn resolve_media_via_api(
ctx: &ResolveContext<'_>,
api_key: &str,
+1 -6
View File
@@ -5,7 +5,6 @@ use crate::http_fetch;
use crate::media_proxy::embed_media_flags;
use crate::text_limits;
use crate::types::{EmbedAuthor, EmbedMedia, EmbedProvider, MessageEmbed};
use fluxer_svc::config::optional_env;
use serde::Deserialize;
use std::future::Future;
use std::pin::Pin;
@@ -61,7 +60,7 @@ impl Resolver for YouTubeResolver {
}
};
let Some(api_key) = ctx.youtube_api_key.clone().or_else(youtube_api_key) else {
let Some(api_key) = ctx.youtube_api_key.clone() else {
tracing::debug!("No YouTube API key configured");
return Ok(ResolverResult { embeds: vec![] });
};
@@ -251,10 +250,6 @@ struct YouTubeThumbnail {
height: Option<u32>,
}
fn youtube_api_key() -> Option<String> {
optional_env("FLUXER_YOUTUBE_API_KEY").or_else(|| optional_env("YOUTUBE_API_KEY"))
}
fn build_youtube_api_url(video_id: &str, api_key: &str) -> anyhow::Result<Url> {
let mut url = Url::parse(YOUTUBE_API_BASE)?;
url.query_pairs_mut()
+14 -2
View File
@@ -24,6 +24,8 @@ pub struct UnfurlShard {
resolvers: Vec<Box<dyn crate::resolvers::Resolver>>,
media_proxy: MediaProxyClient,
self_hosted: bool,
youtube_api_key: Option<String>,
klipy_api_key: Option<String>,
}
impl UnfurlShard {
@@ -78,6 +80,10 @@ impl UnfurlShard {
resolvers,
media_proxy,
self_hosted,
youtube_api_key: optional_env("FLUXER_YOUTUBE_API_KEY")
.or_else(|| optional_env("YOUTUBE_API_KEY")),
klipy_api_key: optional_env("FLUXER_KLIPY_API_KEY")
.or_else(|| optional_env("KLIPY_API_KEY")),
}
}
@@ -100,6 +106,8 @@ impl UnfurlShard {
internal_http_client(),
),
self_hosted: false,
youtube_api_key: None,
klipy_api_key: None,
}
}
@@ -125,8 +133,12 @@ impl UnfurlShard {
nsfw_mode,
media_proxy: &self.media_proxy,
self_hosted: self.self_hosted,
youtube_api_key: youtube_api_key.map(str::to_owned),
klipy_api_key: klipy_api_key.map(str::to_owned),
youtube_api_key: youtube_api_key
.map(str::to_owned)
.or_else(|| self.youtube_api_key.clone()),
klipy_api_key: klipy_api_key
.map(str::to_owned)
.or_else(|| self.klipy_api_key.clone()),
};
if let Some(idx) = matched_resolver_idx {
+3 -2
View File
@@ -1,5 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use fluxer_svc::config::optional_env;
use hmac::{KeyInit, Mac};
use std::sync::{LazyLock, OnceLock};
@@ -46,8 +47,8 @@ pub fn configure(secret: Option<String>, environment: Option<&str>) -> anyhow::R
pub fn configure_from_env() -> anyhow::Result<()> {
configure(
std::env::var(SECRET_ENV).ok(),
std::env::var("FLUXER_ENV").ok().as_deref(),
optional_env(SECRET_ENV),
optional_env("FLUXER_ENV").as_deref(),
)
}
@@ -1,6 +1,9 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {generateKeyPairSync} from 'node:crypto';
import {mkdtempSync, rmSync, writeFileSync} from 'node:fs';
import {tmpdir} from 'node:os';
import {join} from 'node:path';
import {getConfig, loadConfig, resetConfig} from '@fluxer/config/src/ConfigLoader';
import {afterEach, beforeEach, describe, expect, test, vi} from 'vitest';
@@ -80,6 +83,24 @@ describe('ConfigLoader', () => {
expect((await loadConfig()).domain.base_domain).toBe('localhost');
});
test('loadConfig reads secrets from NAME_FILE', async () => {
const dir = mkdtempSync(join(tmpdir(), 'fluxer-config-file-'));
try {
const path = join(dir, 'postgres_password');
writeFileSync(path, 'from-secret-file\n');
stubMinimalEnv({FLUXER_POSTGRES_PASSWORD: '', FLUXER_POSTGRES_PASSWORD_FILE: path});
const config = await loadConfig();
expect(config.database.postgres.password).toBe('from-secret-file');
} finally {
rmSync(dir, {recursive: true, force: true});
}
});
test('loadConfig rejects NAME and NAME_FILE together', async () => {
stubMinimalEnv({FLUXER_SUDO_MODE_SECRET_FILE: '/run/secrets/sudo'});
await expect(loadConfig()).rejects.toThrow('FLUXER_SUDO_MODE_SECRET and FLUXER_SUDO_MODE_SECRET_FILE are both set');
});
test('getConfig throws when config is not loaded', () => {
expect(() => getConfig()).toThrow('Config not loaded');
});
@@ -1,7 +1,14 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {buildNamedFluxerEnvOverrides, setNestedValue} from '@fluxer/config/src/config_loader/EnvironmentOverrides';
import {describe, expect, test} from 'vitest';
import {mkdtempSync, rmSync, writeFileSync} from 'node:fs';
import {tmpdir} from 'node:os';
import {join} from 'node:path';
import {
buildNamedFluxerEnvOverrides,
readEnvValue,
setNestedValue,
} from '@fluxer/config/src/config_loader/EnvironmentOverrides';
import {afterAll, describe, expect, test} from 'vitest';
describe('setNestedValue', () => {
test('sets a top-level key', () => {
@@ -316,3 +323,69 @@ describe('buildNamedFluxerEnvOverrides', () => {
);
});
});
describe('readEnvValue with NAME_FILE', () => {
const dir = mkdtempSync(join(tmpdir(), 'fluxer-env-file-'));
afterAll(() => rmSync(dir, {recursive: true, force: true}));
function secretFile(name: string, contents: string): string {
const path = join(dir, name);
writeFileSync(path, contents);
return path;
}
test('reads the value from NAME_FILE when NAME is unset', () => {
const path = secretFile('plain', 'from-file\n');
expect(readEnvValue({FLUXER_SUDO_MODE_SECRET_FILE: path}, 'FLUXER_SUDO_MODE_SECRET')).toBe('from-file');
});
test('reads the value from NAME_FILE when NAME is blank', () => {
const path = secretFile('blank', 'from-file');
expect(
readEnvValue({FLUXER_SUDO_MODE_SECRET: ' ', FLUXER_SUDO_MODE_SECRET_FILE: path}, 'FLUXER_SUDO_MODE_SECRET'),
).toBe('from-file');
});
test('trims only one trailing newline', () => {
const crlf = secretFile('crlf', 'value\r\n');
const multi = secretFile('multi', '-----BEGIN-----\nabc\n-----END-----\n\n');
expect(readEnvValue({X_FILE: crlf}, 'X')).toBe('value');
expect(readEnvValue({X_FILE: multi}, 'X')).toBe('-----BEGIN-----\nabc\n-----END-----\n');
});
test('treats an empty file as unset', () => {
const path = secretFile('empty', '\n');
expect(readEnvValue({X_FILE: path}, 'X')).toBeUndefined();
});
test('keeps NAME when NAME_FILE is blank', () => {
expect(readEnvValue({X: 'direct', X_FILE: ''}, 'X')).toBe('direct');
});
test('rejects NAME and NAME_FILE together', () => {
const path = secretFile('both', 'from-file');
expect(() => readEnvValue({X: 'direct', X_FILE: path}, 'X')).toThrow('X and X_FILE are both set, set only one');
});
test('names NAME_FILE and the path when the file is missing', () => {
const path = join(dir, 'missing');
expect(() => readEnvValue({X_FILE: path}, 'X')).toThrow(`X_FILE could not read ${path} (ENOENT)`);
});
test('rejects a file that is not valid UTF-8', () => {
const path = join(dir, 'binary');
writeFileSync(path, Buffer.from([0xff, 0x61]));
expect(() => readEnvValue({X_FILE: path}, 'X')).toThrow(`X_FILE could not read ${path} (`);
});
test('feeds named overrides and aliases', () => {
const secret = secretFile('stripe', 'sk_test_file\n');
const nats = secretFile('nats', 'nats://nats:4222\n');
expect(
buildNamedFluxerEnvOverrides({FLUXER_STRIPE_SECRET_KEY_FILE: secret, FLUXER_NATS_CORE_URL_FILE: nats}),
).toMatchObject({
integrations: {stripe: {secret_key: 'sk_test_file'}},
services: {nats: {core_url: 'nats://nats:4222'}},
});
});
});
@@ -1,5 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {readFileSync} from 'node:fs';
import {type ConfigObject, isConfigObject} from '@fluxer/config/src/config_loader/ConfigObject';
type ConfigPathKey = string | number;
@@ -420,11 +421,29 @@ const NAMED_FLUXER_ENV_ALIASES: Record<string, string | undefined> = {
export const NAMED_FLUXER_ENV_NAMES = Object.keys(NAMED_FLUXER_ENV_OVERRIDES);
export function readEnvValue(env: NodeJS.ProcessEnv, name: string): string | undefined {
const value = env[name];
function nonBlank(value: string | undefined): string | undefined {
return value === undefined || value.trim().length === 0 ? undefined : value;
}
export function readEnvValue(env: NodeJS.ProcessEnv, name: string): string | undefined {
const value = nonBlank(env[name]);
const filePath = nonBlank(env[`${name}_FILE`]);
if (filePath === undefined) {
return value;
}
if (value !== undefined) {
throw new Error(`${name} and ${name}_FILE are both set, set only one`);
}
let contents: string;
try {
contents = new TextDecoder('utf-8', {fatal: true}).decode(readFileSync(filePath));
} catch (error) {
const reason = error instanceof Error && 'code' in error ? String(error.code) : 'unreadable';
throw new Error(`${name}_FILE could not read ${filePath} (${reason})`);
}
return nonBlank(contents.replace(/\r?\n$/, ''));
}
export function buildNamedFluxerEnvOverrides(env: NodeJS.ProcessEnv): ConfigObject {
const overrides: ConfigObject = {};
for (const [envKey, mapping] of Object.entries(NAMED_FLUXER_ENV_OVERRIDES)) {
@@ -12,6 +12,14 @@ import {ColorType, createStringType, Int32Type} from '@fluxer/schema/src/primiti
import {z} from 'zod';
const TEMPLATE_NAME_MAX_LENGTH = 100;
const TEMPLATE_TOPIC_MAX_LENGTH = 4096;
function clipTopic(value: string): string {
if (value.length <= CHANNEL_TOPIC_MAX_LENGTH) return value;
const last = value.charCodeAt(CHANNEL_TOPIC_MAX_LENGTH - 1);
const end = last >= 0xd800 && last <= 0xdbff ? CHANNEL_TOPIC_MAX_LENGTH - 1 : CHANNEL_TOPIC_MAX_LENGTH;
return value.slice(0, end);
}
const TemplateEntityId = z
.union([
@@ -50,7 +58,12 @@ export const TemplateChannel = z.object({
.nullish()
.transform((value) => value ?? '')
.describe('The name of the channel'),
topic: z.string().max(CHANNEL_TOPIC_MAX_LENGTH).nullish().describe('The channel topic'),
topic: z
.string()
.max(TEMPLATE_TOPIC_MAX_LENGTH)
.transform(clipTopic)
.nullish()
.describe('The channel topic, shortened to the Fluxer channel topic limit'),
position: Int32Type.describe('The position of the channel'),
parent_id: TemplateEntityId.nullish().describe('The template-local ID of the parent category'),
bitrate: z.number().int().nonnegative().nullish().describe('The bitrate for voice channels'),