mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(api): gate stream keys by channel type and cover previews (#2509)
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {Permissions} from '@fluxer/constants/src/ChannelConstants';
|
||||
import {ChannelTypes, Permissions} from '@fluxer/constants/src/ChannelConstants';
|
||||
import {InvalidChannelTypeError} from '@fluxer/errors/src/domains/channel/InvalidChannelTypeError';
|
||||
import {InvalidStreamKeyFormatError} from '@fluxer/errors/src/domains/channel/InvalidStreamKeyFormatError';
|
||||
import {InvalidStreamThumbnailPayloadError} from '@fluxer/errors/src/domains/channel/InvalidStreamThumbnailPayloadError';
|
||||
import {StreamKeyChannelMismatchError} from '@fluxer/errors/src/domains/channel/StreamKeyChannelMismatchError';
|
||||
@@ -90,6 +91,9 @@ export class StreamService {
|
||||
if (params.parsedKey.guildId !== channel.guildId.toString()) {
|
||||
throw new StreamKeyScopeMismatchError();
|
||||
}
|
||||
if (channel.type !== ChannelTypes.GUILD_VOICE) {
|
||||
throw new InvalidChannelTypeError();
|
||||
}
|
||||
const hasConnect = await this.gatewayService.checkPermission({
|
||||
guildId: channel.guildId,
|
||||
channelId: params.channelId,
|
||||
@@ -99,8 +103,13 @@ export class StreamService {
|
||||
if (!hasConnect) {
|
||||
throw new MissingPermissionsError();
|
||||
}
|
||||
} else if (params.parsedKey.scope !== 'dm') {
|
||||
throw new StreamKeyScopeMismatchError();
|
||||
} else {
|
||||
if (params.parsedKey.scope !== 'dm') {
|
||||
throw new StreamKeyScopeMismatchError();
|
||||
}
|
||||
if (channel.type !== ChannelTypes.DM && channel.type !== ChannelTypes.GROUP_DM) {
|
||||
throw new InvalidChannelTypeError();
|
||||
}
|
||||
}
|
||||
if (params.parsedKey.channelId !== params.channelId.toString()) {
|
||||
throw new StreamKeyChannelMismatchError();
|
||||
@@ -166,10 +175,7 @@ export class StreamService {
|
||||
channelId,
|
||||
parsedKey,
|
||||
});
|
||||
const preview = await this.streamPreviewService.getPreview(params.streamKey);
|
||||
if (preview) {
|
||||
}
|
||||
return preview;
|
||||
return this.streamPreviewService.getPreview(params.streamKey);
|
||||
}
|
||||
|
||||
async uploadPreview(params: {
|
||||
|
||||
@@ -0,0 +1,86 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {APIErrorCodes} from '@fluxer/constants/src/ApiErrorCodes';
|
||||
import {ChannelTypes} from '@fluxer/constants/src/ChannelConstants';
|
||||
import {afterAll, beforeAll, beforeEach, describe, it} from 'vitest';
|
||||
import {createTestAccount} from '../../auth/tests/AuthTestUtils';
|
||||
import {type ApiTestHarness, createApiTestHarness} from '../../test/ApiTestHarness';
|
||||
import {HTTP_STATUS} from '../../test/TestConstants';
|
||||
import {createBuilder} from '../../test/TestRequestBuilder';
|
||||
import {createChannel, createDmChannel, createFriendship, createGuild, getChannel} from './ChannelTestUtils';
|
||||
|
||||
const CONNECTION_ID = 'conn-channel-type';
|
||||
|
||||
describe('stream key channel type gate', () => {
|
||||
let harness: ApiTestHarness;
|
||||
|
||||
beforeAll(async () => {
|
||||
harness = await createApiTestHarness();
|
||||
});
|
||||
|
||||
beforeEach(async () => {
|
||||
await harness.reset();
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
await harness?.shutdown();
|
||||
});
|
||||
|
||||
async function createGuildChannels(): Promise<{
|
||||
token: string;
|
||||
textKey: string;
|
||||
voiceKey: string;
|
||||
}> {
|
||||
const owner = await createTestAccount(harness);
|
||||
const guild = await createGuild(harness, owner.token, 'Stream Type Guild');
|
||||
const textChannel = await getChannel(harness, owner.token, guild.system_channel_id!);
|
||||
const voiceChannel = await createChannel(harness, owner.token, guild.id, 'stream-voice', ChannelTypes.GUILD_VOICE);
|
||||
return {
|
||||
token: owner.token,
|
||||
textKey: `${guild.id}:${textChannel.id}:${CONNECTION_ID}`,
|
||||
voiceKey: `${guild.id}:${voiceChannel.id}:${CONNECTION_ID}`,
|
||||
};
|
||||
}
|
||||
|
||||
it('rejects a region update whose key names a guild text channel', async () => {
|
||||
const {token, textKey} = await createGuildChannels();
|
||||
await createBuilder(harness, token)
|
||||
.patch(`/streams/${textKey}/stream`)
|
||||
.body({region: 'us-east'})
|
||||
.expect(HTTP_STATUS.BAD_REQUEST, APIErrorCodes.INVALID_CHANNEL_TYPE)
|
||||
.execute();
|
||||
});
|
||||
|
||||
it('rejects a preview read whose key names a guild text channel', async () => {
|
||||
const {token, textKey} = await createGuildChannels();
|
||||
await createBuilder(harness, token)
|
||||
.get(`/streams/${textKey}/preview`)
|
||||
.expect(HTTP_STATUS.BAD_REQUEST, APIErrorCodes.INVALID_CHANNEL_TYPE)
|
||||
.execute();
|
||||
});
|
||||
|
||||
it('lets a guild voice key past the type gate and fail on the missing voice state', async () => {
|
||||
const {token, voiceKey} = await createGuildChannels();
|
||||
await createBuilder(harness, token)
|
||||
.patch(`/streams/${voiceKey}/stream`)
|
||||
.body({region: 'us-east'})
|
||||
.expect(HTTP_STATUS.FORBIDDEN, APIErrorCodes.ACCESS_DENIED)
|
||||
.execute();
|
||||
});
|
||||
|
||||
it('lets a guild voice key past the type gate and read an absent preview', async () => {
|
||||
const {token, voiceKey} = await createGuildChannels();
|
||||
await createBuilder(harness, token).get(`/streams/${voiceKey}/preview`).expect(HTTP_STATUS.NOT_FOUND).execute();
|
||||
});
|
||||
|
||||
it('lets a dm key past the type gate and read an absent preview', async () => {
|
||||
const owner = await createTestAccount(harness);
|
||||
const recipient = await createTestAccount(harness);
|
||||
await createFriendship(harness, owner, recipient);
|
||||
const dm = await createDmChannel(harness, owner.token, recipient.userId);
|
||||
await createBuilder(harness, owner.token)
|
||||
.get(`/streams/dm:${dm.id}:${CONNECTION_ID}/preview`)
|
||||
.expect(HTTP_STATUS.NOT_FOUND)
|
||||
.execute();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,76 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {STREAM_PREVIEW_CONTENT_TYPE_JPEG, STREAM_PREVIEW_MAX_BYTES} from '@fluxer/constants/src/StreamConstants';
|
||||
import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest';
|
||||
import {createTestAccount} from '../../auth/tests/AuthTestUtils';
|
||||
import {Config} from '../../Config';
|
||||
import {getCacheService} from '../../middleware/ServiceSingletons';
|
||||
import {type ApiTestHarness, createApiTestHarness} from '../../test/ApiTestHarness';
|
||||
import {createDmChannel, createFriendship} from './ChannelTestUtils';
|
||||
|
||||
const CONNECTION_ID = 'conn-oversized';
|
||||
|
||||
describe('stream preview object size ceiling', () => {
|
||||
let harness: ApiTestHarness;
|
||||
|
||||
beforeAll(async () => {
|
||||
harness = await createApiTestHarness();
|
||||
});
|
||||
|
||||
beforeEach(async () => {
|
||||
await harness.reset();
|
||||
harness.storageService.reset();
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
await harness?.shutdown();
|
||||
});
|
||||
|
||||
async function createStoredPreview(fileData: Uint8Array): Promise<{token: string; streamKey: string}> {
|
||||
const owner = await createTestAccount(harness);
|
||||
const viewer = await createTestAccount(harness);
|
||||
await createFriendship(harness, owner, viewer);
|
||||
const dm = await createDmChannel(harness, owner.token, viewer.userId);
|
||||
const streamKey = `dm:${dm.id}:${CONNECTION_ID}`;
|
||||
await getCacheService().set(
|
||||
`stream_preview:${streamKey}`,
|
||||
{
|
||||
bucket: Config.s3.buckets.uploads,
|
||||
key: `stream_previews/${dm.id}-${CONNECTION_ID}.jpg`,
|
||||
updatedAt: Date.now(),
|
||||
ownerId: owner.userId,
|
||||
channelId: dm.id,
|
||||
contentType: STREAM_PREVIEW_CONTENT_TYPE_JPEG,
|
||||
},
|
||||
60,
|
||||
);
|
||||
harness.storageService.configure({fileData});
|
||||
return {token: owner.token, streamKey};
|
||||
}
|
||||
|
||||
it('answers an empty 404 when the stored object is larger than the preview ceiling', async () => {
|
||||
const {token, streamKey} = await createStoredPreview(new Uint8Array(STREAM_PREVIEW_MAX_BYTES + 1));
|
||||
const response = await harness.requestJson({
|
||||
path: `/v1/streams/${streamKey}/preview`,
|
||||
headers: {authorization: token},
|
||||
});
|
||||
expect(response.status).toBe(404);
|
||||
expect(await response.text()).toBe('');
|
||||
expect(harness.storageService.readObjectSpy).toHaveBeenCalledWith(
|
||||
Config.s3.buckets.uploads,
|
||||
expect.any(String),
|
||||
STREAM_PREVIEW_MAX_BYTES,
|
||||
);
|
||||
});
|
||||
|
||||
it('serves a stored object that fits inside the preview ceiling', async () => {
|
||||
const {token, streamKey} = await createStoredPreview(new Uint8Array([0xff, 0xd8, 0x00, 0xff, 0xd9]));
|
||||
const response = await harness.requestJson({
|
||||
path: `/v1/streams/${streamKey}/preview`,
|
||||
headers: {authorization: token},
|
||||
});
|
||||
expect(response.status).toBe(200);
|
||||
expect(response.headers.get('content-type')).toBe(STREAM_PREVIEW_CONTENT_TYPE_JPEG);
|
||||
expect(new Uint8Array(await response.arrayBuffer())).toEqual(new Uint8Array([0xff, 0xd8, 0x00, 0xff, 0xd9]));
|
||||
});
|
||||
});
|
||||
@@ -71,9 +71,12 @@ class TestStorageService extends StorageService {
|
||||
this.readObjects.push(
|
||||
maxBytes === undefined ? {bucket: _bucket, key: _key} : {bucket: _bucket, key: _key, maxBytes},
|
||||
);
|
||||
return maxBytes !== undefined && this.sourceData.length > maxBytes
|
||||
? this.sourceData.slice(0, maxBytes)
|
||||
: this.sourceData;
|
||||
if (maxBytes !== undefined && this.sourceData.length > maxBytes) {
|
||||
throw new Error(
|
||||
`Stream exceeds maximum buffer size of ${maxBytes} bytes (got at least ${this.sourceData.length} bytes)`,
|
||||
);
|
||||
}
|
||||
return this.sourceData;
|
||||
}
|
||||
|
||||
override async copyObject(params: CopyObjectTestParams): Promise<void> {
|
||||
@@ -188,6 +191,40 @@ describe('StorageService.copyObjectWithMetadataStripping', () => {
|
||||
});
|
||||
});
|
||||
|
||||
interface S3ClientProbe {
|
||||
client: {send: (command: unknown) => Promise<{Body: Readable}>};
|
||||
}
|
||||
|
||||
function stubObjectBody(service: StorageService, chunks: Array<Uint8Array>): void {
|
||||
(service as unknown as S3ClientProbe).client = {
|
||||
send: async () => ({Body: Readable.from(chunks.map((chunk) => Buffer.from(chunk)))}),
|
||||
};
|
||||
}
|
||||
|
||||
describe('StorageService.readObject', () => {
|
||||
it('returns the whole object when it fits inside the ceiling', async () => {
|
||||
const service = new StorageService();
|
||||
stubObjectBody(service, [new Uint8Array([1, 2, 3]), new Uint8Array([4, 5])]);
|
||||
const data = await service.readObject('fluxer-uploads', 'stream_previews/preview.jpg', 5);
|
||||
expect(data).toEqual(new Uint8Array([1, 2, 3, 4, 5]));
|
||||
});
|
||||
|
||||
it('rejects an object larger than the ceiling instead of truncating it', async () => {
|
||||
const service = new StorageService();
|
||||
stubObjectBody(service, [new Uint8Array([1, 2, 3]), new Uint8Array([4, 5, 6])]);
|
||||
await expect(service.readObject('fluxer-uploads', 'stream_previews/preview.jpg', 5)).rejects.toThrow(
|
||||
/^Stream exceeds maximum buffer size of 5 bytes/,
|
||||
);
|
||||
});
|
||||
|
||||
it('reads the whole object when no ceiling is given', async () => {
|
||||
const service = new StorageService();
|
||||
stubObjectBody(service, [new Uint8Array([1, 2, 3]), new Uint8Array([4, 5, 6])]);
|
||||
const data = await service.readObject('fluxer-uploads', 'stream_previews/preview.jpg');
|
||||
expect(data).toEqual(new Uint8Array([1, 2, 3, 4, 5, 6]));
|
||||
});
|
||||
});
|
||||
|
||||
describe('provider selection', () => {
|
||||
interface ClientProbe {
|
||||
client: {config: {region: () => Promise<string>; endpoint?: () => Promise<{hostname: string}>}};
|
||||
|
||||
Reference in New Issue
Block a user