diff --git a/drizzle/0025_author_soft_delete.sql b/drizzle/0025_author_soft_delete.sql new file mode 100644 index 0000000..4714be8 --- /dev/null +++ b/drizzle/0025_author_soft_delete.sql @@ -0,0 +1,2 @@ +ALTER TABLE "topics" ADD COLUMN "is_author_deleted" boolean DEFAULT false NOT NULL; +ALTER TABLE "replies" ADD COLUMN "is_author_deleted" boolean DEFAULT false NOT NULL; diff --git a/drizzle/0026_composite_indexes.sql b/drizzle/0026_composite_indexes.sql new file mode 100644 index 0000000..a8d0066 --- /dev/null +++ b/drizzle/0026_composite_indexes.sql @@ -0,0 +1,3 @@ +CREATE INDEX IF NOT EXISTS "topics_community_category_activity_idx" ON "topics" ("community_did","category","last_activity_at"); +CREATE INDEX IF NOT EXISTS "replies_root_uri_created_at_idx" ON "replies" ("root_uri","created_at"); +CREATE INDEX IF NOT EXISTS "reactions_subject_uri_type_idx" ON "reactions" ("subject_uri","type"); diff --git a/drizzle/meta/_journal.json b/drizzle/meta/_journal.json index 109b91f..2723f4c 100644 --- a/drizzle/meta/_journal.json +++ b/drizzle/meta/_journal.json @@ -176,6 +176,20 @@ "when": 1771280352200, "tag": "0024_grey_skreet", "breakpoints": true + }, + { + "idx": 25, + "version": "7", + "when": 1771459200000, + "tag": "0025_author_soft_delete", + "breakpoints": true + }, + { + "idx": 26, + "version": "7", + "when": 1771459200001, + "tag": "0026_composite_indexes", + "breakpoints": true } ] } diff --git a/package.json b/package.json index a39938a..b355623 100644 --- a/package.json +++ b/package.json @@ -32,7 +32,7 @@ }, "dependencies": { "@atproto/api": "^0.18.21", - "@atproto/oauth-client-node": "^0.3.16", + "@atproto/oauth-client-node": "^0.3.17", "@atproto/tap": "^0.2.4", "@barazo-forum/lexicons": "^0.1.0", "@fastify/cookie": "^11.0.2", diff --git a/src/auth/oauth-client.ts b/src/auth/oauth-client.ts index ffd2b41..1c17d6b 100644 --- a/src/auth/oauth-client.ts +++ b/src/auth/oauth-client.ts @@ -102,17 +102,13 @@ export function createOAuthClient(env: Env, cache: Cache, logger: Logger): NodeO stateStore: new ValkeyStateStore(cache, logger), sessionStore: new ValkeySessionStore(cache, logger, env.OAUTH_SESSION_TTL), requestLock: createRequestLock(cache, logger), - }) - - // Log session lifecycle events for observability - client.addEventListener('updated', (event: CustomEvent) => { - const detail = event.detail as { sub: string } - logger.info({ sub: detail.sub }, 'OAuth session updated') - }) - - client.addEventListener('deleted', (event: CustomEvent) => { - const detail = event.detail as { sub: string; cause: unknown } - logger.info({ sub: detail.sub, cause: String(detail.cause) }, 'OAuth session deleted') + // Session lifecycle hooks for observability (replaces addEventListener in >=0.3.17) + onUpdate: (sub: string) => { + logger.info({ sub }, 'OAuth session updated') + }, + onDelete: (sub: string, cause: unknown) => { + logger.info({ sub, cause: String(cause) }, 'OAuth session deleted') + }, }) return client diff --git a/src/config/env.ts b/src/config/env.ts index e507892..918cf97 100644 --- a/src/config/env.ts +++ b/src/config/env.ts @@ -20,7 +20,7 @@ const positiveIntFromString = (defaultVal: string) => .transform((val) => Number(val)) .pipe(z.number().int().positive()) -export const envSchema = z.object({ +const baseEnvSchema = z.object({ // Required DATABASE_URL: z.url(), VALKEY_URL: z.url(), @@ -95,8 +95,28 @@ export const envSchema = z.object({ OZONE_LABELER_URL: z.string().default('https://mod.bsky.app'), }) +export const envSchema = baseEnvSchema.refine( + (data) => + data.COMMUNITY_MODE !== 'single' || (data.COMMUNITY_DID && data.COMMUNITY_DID.length > 0), + { + message: 'COMMUNITY_DID is required when COMMUNITY_MODE is "single"', + path: ['COMMUNITY_DID'], + } +) + export type Env = z.infer +/** + * Get the community DID, throwing if not set. + * Safe to call after env validation -- single mode requires COMMUNITY_DID at startup. + */ +export function getCommunityDid(env: Env): string { + if (!env.COMMUNITY_DID) { + throw new Error('COMMUNITY_DID is required in single mode but not set') + } + return env.COMMUNITY_DID +} + export function parseEnv(env: Record): Env { const result = envSchema.safeParse(env) if (!result.success) { diff --git a/src/db/schema/reactions.ts b/src/db/schema/reactions.ts index afcc287..6c5471f 100644 --- a/src/db/schema/reactions.ts +++ b/src/db/schema/reactions.ts @@ -21,5 +21,6 @@ export const reactions = pgTable( // communityDid intentionally excluded: AT URIs are globally unique, so a // reaction to a given subject is inherently community-scoped via the subject URI. unique('reactions_author_subject_type_uniq').on(table.authorDid, table.subjectUri, table.type), + index('reactions_subject_uri_type_idx').on(table.subjectUri, table.type), ] ) diff --git a/src/db/schema/replies.ts b/src/db/schema/replies.ts index ccb27c4..f1581ba 100644 --- a/src/db/schema/replies.ts +++ b/src/db/schema/replies.ts @@ -1,4 +1,4 @@ -import { pgTable, text, integer, timestamp, jsonb, index } from 'drizzle-orm/pg-core' +import { pgTable, text, integer, timestamp, jsonb, boolean, index } from 'drizzle-orm/pg-core' export const replies = pgTable( 'replies', @@ -18,6 +18,7 @@ export const replies = pgTable( reactionCount: integer('reaction_count').notNull().default(0), createdAt: timestamp('created_at', { withTimezone: true }).notNull(), indexedAt: timestamp('indexed_at', { withTimezone: true }).notNull().defaultNow(), + isAuthorDeleted: boolean('is_author_deleted').notNull().default(false), moderationStatus: text('moderation_status', { enum: ['approved', 'held', 'rejected'], }) @@ -42,5 +43,6 @@ export const replies = pgTable( index('replies_community_did_idx').on(table.communityDid), index('replies_moderation_status_idx').on(table.moderationStatus), index('replies_trust_status_idx').on(table.trustStatus), + index('replies_root_uri_created_at_idx').on(table.rootUri, table.createdAt), ] ) diff --git a/src/db/schema/topics.ts b/src/db/schema/topics.ts index 2a425fb..df0645f 100644 --- a/src/db/schema/topics.ts +++ b/src/db/schema/topics.ts @@ -22,6 +22,7 @@ export const topics = pgTable( isLocked: boolean('is_locked').notNull().default(false), isPinned: boolean('is_pinned').notNull().default(false), isModDeleted: boolean('is_mod_deleted').notNull().default(false), + isAuthorDeleted: boolean('is_author_deleted').notNull().default(false), moderationStatus: text('moderation_status', { enum: ['approved', 'held', 'rejected'], }) @@ -46,5 +47,10 @@ export const topics = pgTable( index('topics_community_did_idx').on(table.communityDid), index('topics_moderation_status_idx').on(table.moderationStatus), index('topics_trust_status_idx').on(table.trustStatus), + index('topics_community_category_activity_idx').on( + table.communityDid, + table.category, + table.lastActivityAt + ), ] ) diff --git a/src/firehose/clamp-timestamp.ts b/src/firehose/clamp-timestamp.ts new file mode 100644 index 0000000..adf82eb --- /dev/null +++ b/src/firehose/clamp-timestamp.ts @@ -0,0 +1,22 @@ +const MAX_FUTURE_MS = 5 * 60 * 1000 // 5 minutes +const MAX_PAST_MS = 60 * 60 * 1000 // 1 hour + +/** + * Clamp a client-declared createdAt timestamp to prevent feed manipulation. + * + * AT Protocol createdAt is client-declared — a malicious client can set it to + * any value. Clamping prevents future-dated records from permanently pinning to + * the top of feeds, and extremely backdated records from being buried. + * + * - Future timestamps (> 5 min ahead of now) → clamped to now + * - Very old timestamps (> 1 hour before now) → clamped to 1 hour ago + * - Otherwise → client timestamp used as-is + */ +export function clampCreatedAt(clientCreatedAt: Date, now: Date = new Date()): Date { + const maxFuture = new Date(now.getTime() + MAX_FUTURE_MS) + const maxPast = new Date(now.getTime() - MAX_PAST_MS) + + if (clientCreatedAt > maxFuture) return now + if (clientCreatedAt < maxPast) return maxPast + return clientCreatedAt +} diff --git a/src/firehose/handlers/record.ts b/src/firehose/handlers/record.ts index a9a93b0..0a1937d 100644 --- a/src/firehose/handlers/record.ts +++ b/src/firehose/handlers/record.ts @@ -1,10 +1,11 @@ import { eq } from 'drizzle-orm' import { users } from '../../db/schema/users.js' +import { replies } from '../../db/schema/replies.js' +import { reactions } from '../../db/schema/reactions.js' import type { Database } from '../../db/index.js' import type { Logger } from '../../lib/logger.js' import type { RecordEvent } from '../types.js' -import { COLLECTION_MAP, SUPPORTED_COLLECTIONS } from '../types.js' -import type { SupportedCollection } from '../types.js' +import { COLLECTION_MAP, isSupportedCollection } from '../types.js' import { validateRecord } from '../validation.js' import type { TopicIndexer } from '../indexers/topic.js' import type { ReplyIndexer } from '../indexers/reply.js' @@ -17,10 +18,6 @@ interface Indexers { reaction: ReactionIndexer } -function isSupportedCollection(collection: string): collection is SupportedCollection { - return (SUPPORTED_COLLECTIONS as readonly string[]).includes(collection) -} - export class RecordHandler { constructor( private indexers: Indexers, @@ -154,25 +151,53 @@ export class RecordHandler { did: params.did, }) break - case 'reply': - // For reply delete, we need the root URI to decrement the count. - // If the record is available (backfill), use it. Otherwise, the - // integration will handle the count via the stored rootUri. + case 'reply': { + // AT Protocol delete events don't include record data, so look up + // the rootUri from the DB before the indexer hard-deletes the row. + const replyRows = await this.db + .select({ rootUri: replies.rootUri }) + .from(replies) + .where(eq(replies.uri, params.uri)) + + const rootUri = replyRows[0]?.rootUri ?? '' + if (!rootUri) { + this.logger.debug( + { uri: params.uri }, + 'Reply not found in DB for delete, count decrement will be skipped' + ) + } + await this.indexers.reply.handleDelete({ uri: params.uri, rkey: params.rkey, did: params.did, - rootUri: '', + rootUri, }) break - case 'reaction': + } + case 'reaction': { + // Same pattern: look up subjectUri before the row is deleted. + const reactionRows = await this.db + .select({ subjectUri: reactions.subjectUri }) + .from(reactions) + .where(eq(reactions.uri, params.uri)) + + const subjectUri = reactionRows[0]?.subjectUri ?? '' + if (!subjectUri) { + this.logger.debug( + { uri: params.uri }, + 'Reaction not found in DB for delete, count decrement will be skipped' + ) + } + await this.indexers.reaction.handleDelete({ uri: params.uri, rkey: params.rkey, did: params.did, - subjectUri: '', + subjectUri, }) break + } } } diff --git a/src/firehose/indexers/reaction.ts b/src/firehose/indexers/reaction.ts index 093d3d7..be43d7b 100644 --- a/src/firehose/indexers/reaction.ts +++ b/src/firehose/indexers/reaction.ts @@ -4,10 +4,14 @@ import { topics } from '../../db/schema/topics.js' import { replies } from '../../db/schema/replies.js' import type { Database } from '../../db/index.js' import type { Logger } from '../../lib/logger.js' +import { clampCreatedAt } from '../clamp-timestamp.js' const TOPIC_COLLECTION = 'forum.barazo.topic.post' const REPLY_COLLECTION = 'forum.barazo.topic.reply' +/** Transaction type extracted from Database.transaction() callback parameter */ +type Transaction = Parameters[0]>[0] + interface CreateParams { uri: string rkey: string @@ -39,8 +43,10 @@ export class ReactionIndexer { ) {} async handleCreate(params: CreateParams): Promise { - const { uri, rkey, did, cid, record } = params + const { uri, rkey, did, cid, record, live } = params const subject = record['subject'] as { uri: string; cid: string } + const clientCreatedAt = new Date(record['createdAt'] as string) + const createdAt = live ? clampCreatedAt(clientCreatedAt) : clientCreatedAt await this.db.transaction(async (tx) => { await tx @@ -54,11 +60,11 @@ export class ReactionIndexer { type: record['type'] as string, communityDid: record['community'] as string, cid, - createdAt: new Date(record['createdAt'] as string), + createdAt, }) .onConflictDoNothing() - await this.incrementReactionCount(tx as never, subject.uri) + await this.incrementReactionCount(tx, subject.uri) }) this.logger.debug({ uri, did }, 'Indexed reaction') @@ -69,13 +75,13 @@ export class ReactionIndexer { await this.db.transaction(async (tx) => { await tx.delete(reactions).where(eq(reactions.uri, uri)) - await this.decrementReactionCount(tx as never, subjectUri) + await this.decrementReactionCount(tx, subjectUri) }) this.logger.debug({ uri }, 'Deleted reaction') } - private async incrementReactionCount(tx: Database, subjectUri: string): Promise { + private async incrementReactionCount(tx: Transaction, subjectUri: string): Promise { const collection = getCollectionFromUri(subjectUri) if (collection === TOPIC_COLLECTION) { @@ -91,7 +97,7 @@ export class ReactionIndexer { } } - private async decrementReactionCount(tx: Database, subjectUri: string): Promise { + private async decrementReactionCount(tx: Transaction, subjectUri: string): Promise { const collection = getCollectionFromUri(subjectUri) if (collection === TOPIC_COLLECTION) { diff --git a/src/firehose/indexers/reply.ts b/src/firehose/indexers/reply.ts index 9451059..ad9057b 100644 --- a/src/firehose/indexers/reply.ts +++ b/src/firehose/indexers/reply.ts @@ -4,6 +4,7 @@ import { topics } from '../../db/schema/topics.js' import type { Database } from '../../db/index.js' import type { Logger } from '../../lib/logger.js' import type { TrustStatus } from '../../services/account-age.js' +import { clampCreatedAt } from '../clamp-timestamp.js' interface CreateParams { uri: string @@ -39,10 +40,12 @@ export class ReplyIndexer { ) {} async handleCreate(params: CreateParams): Promise { - const { uri, rkey, did, cid, record, trustStatus } = params + const { uri, rkey, did, cid, record, live, trustStatus } = params const root = record['root'] as { uri: string; cid: string } const parent = record['parent'] as { uri: string; cid: string } + const clientCreatedAt = new Date(record['createdAt'] as string) + const createdAt = live ? clampCreatedAt(clientCreatedAt) : clientCreatedAt await this.db.transaction(async (tx) => { await tx @@ -60,7 +63,7 @@ export class ReplyIndexer { communityDid: record['community'] as string, cid, labels: (record['labels'] as { values: { val: string }[] } | undefined) ?? null, - createdAt: new Date(record['createdAt'] as string), + createdAt, trustStatus, }) .onConflictDoNothing() @@ -99,17 +102,19 @@ export class ReplyIndexer { const { uri, rootUri } = params await this.db.transaction(async (tx) => { - await tx.delete(replies).where(eq(replies.uri, uri)) + await tx.update(replies).set({ isAuthorDeleted: true }).where(eq(replies.uri, uri)) // Decrement reply count (floor at 0 via GREATEST) - await tx - .update(topics) - .set({ - replyCount: sql`GREATEST(${topics.replyCount} - 1, 0)`, - }) - .where(eq(topics.uri, rootUri)) + if (rootUri) { + await tx + .update(topics) + .set({ + replyCount: sql`GREATEST(${topics.replyCount} - 1, 0)`, + }) + .where(eq(topics.uri, rootUri)) + } }) - this.logger.debug({ uri }, 'Deleted reply') + this.logger.debug({ uri }, 'Soft-deleted reply (author delete)') } } diff --git a/src/firehose/indexers/topic.ts b/src/firehose/indexers/topic.ts index f1f1828..555f377 100644 --- a/src/firehose/indexers/topic.ts +++ b/src/firehose/indexers/topic.ts @@ -3,6 +3,7 @@ import { topics } from '../../db/schema/topics.js' import type { Database } from '../../db/index.js' import type { Logger } from '../../lib/logger.js' import type { TrustStatus } from '../../services/account-age.js' +import { clampCreatedAt } from '../clamp-timestamp.js' interface CreateParams { uri: string @@ -27,7 +28,9 @@ export class TopicIndexer { ) {} async handleCreate(params: CreateParams): Promise { - const { uri, rkey, did, cid, record, trustStatus } = params + const { uri, rkey, did, cid, record, live, trustStatus } = params + const clientCreatedAt = new Date(record['createdAt'] as string) + const createdAt = live ? clampCreatedAt(clientCreatedAt) : clientCreatedAt await this.db .insert(topics) @@ -43,8 +46,8 @@ export class TopicIndexer { communityDid: record['community'] as string, cid, labels: (record['labels'] as { values: { val: string }[] } | undefined) ?? null, - createdAt: new Date(record['createdAt'] as string), - lastActivityAt: new Date(record['createdAt'] as string), + createdAt, + lastActivityAt: createdAt, trustStatus, }) .onConflictDoUpdate({ @@ -87,8 +90,8 @@ export class TopicIndexer { async handleDelete(params: DeleteParams): Promise { const { uri } = params - await this.db.delete(topics).where(eq(topics.uri, uri)) + await this.db.update(topics).set({ isAuthorDeleted: true }).where(eq(topics.uri, uri)) - this.logger.debug({ uri }, 'Deleted topic') + this.logger.debug({ uri }, 'Soft-deleted topic (author delete)') } } diff --git a/src/firehose/service.ts b/src/firehose/service.ts index 1dfc63a..3eaa91c 100644 --- a/src/firehose/service.ts +++ b/src/firehose/service.ts @@ -99,12 +99,7 @@ export class FirehoseService { }) this.channel = this.tap.channel(indexer) - // Start in background (non-blocking) - void this.channel.start().catch((err: unknown) => { - this.logger.error({ err }, 'Firehose channel error') - this.connected = false - }) - + await this.channel.start() this.connected = true this.logger.info('Firehose subscription started') } catch (err) { diff --git a/src/firehose/types.ts b/src/firehose/types.ts index e5c7c88..a104635 100644 --- a/src/firehose/types.ts +++ b/src/firehose/types.ts @@ -53,6 +53,11 @@ export const SUPPORTED_COLLECTIONS = [ export type SupportedCollection = (typeof SUPPORTED_COLLECTIONS)[number] +/** Type guard: is the collection string one of the supported Barazo collections? */ +export function isSupportedCollection(collection: string): collection is SupportedCollection { + return (SUPPORTED_COLLECTIONS as readonly string[]).includes(collection) +} + /** Maps collection NSIDs to short indexer names. */ export const COLLECTION_MAP: Record = { 'forum.barazo.topic.post': 'topic', diff --git a/src/firehose/validation.ts b/src/firehose/validation.ts index f4993cd..da26446 100644 --- a/src/firehose/validation.ts +++ b/src/firehose/validation.ts @@ -1,6 +1,6 @@ import { topicPostSchema, topicReplySchema, reactionSchema } from '@barazo-forum/lexicons' import type { SupportedCollection } from './types.js' -import { SUPPORTED_COLLECTIONS } from './types.js' +import { isSupportedCollection } from './types.js' const MAX_RECORD_SIZE = 64 * 1024 // 64KB @@ -17,10 +17,6 @@ const schemaMap: Record< 'forum.barazo.interaction.reaction': reactionSchema, } -function isSupportedCollection(collection: string): collection is SupportedCollection { - return (SUPPORTED_COLLECTIONS as readonly string[]).includes(collection) -} - export function validateRecord(collection: string, record: unknown): ValidationResult { if (!isSupportedCollection(collection)) { return { success: false, error: `Unsupported collection: ${collection}` } diff --git a/src/lib/at-uri.ts b/src/lib/at-uri.ts new file mode 100644 index 0000000..d588b98 --- /dev/null +++ b/src/lib/at-uri.ts @@ -0,0 +1,14 @@ +/** + * Extract the rkey (record key) from an AT URI. + * Format: at://did:plc:xxx/collection/rkey + * + * @throws Error if the rkey is missing or empty + */ +export function extractRkey(uri: string): string { + const parts = uri.split('/') + const rkey = parts[parts.length - 1] + if (!rkey) { + throw new Error('Invalid AT URI: missing rkey') + } + return rkey +} diff --git a/src/lib/pds-client.ts b/src/lib/pds-client.ts index 46cd4c5..49b7e74 100644 --- a/src/lib/pds-client.ts +++ b/src/lib/pds-client.ts @@ -84,7 +84,7 @@ export function createPdsClient(oauthClient: NodeOAuthClient, logger: Logger): P const response = await agent.com.atproto.repo.createRecord({ repo: did, collection, - record, + record: { $type: collection, ...record }, }) return { uri: response.data.uri, cid: response.data.cid } @@ -108,7 +108,7 @@ export function createPdsClient(oauthClient: NodeOAuthClient, logger: Logger): P repo: did, collection, rkey, - record, + record: { $type: collection, ...record }, }) return { uri: response.data.uri, cid: response.data.cid } diff --git a/src/routes/admin-settings.ts b/src/routes/admin-settings.ts index f5d6e19..f9a8538 100644 --- a/src/routes/admin-settings.ts +++ b/src/routes/admin-settings.ts @@ -408,13 +408,13 @@ export function adminSettingsRoutes(): FastifyPluginCallback { async (_request, reply) => { const result = await db.execute(sql` SELECT - (SELECT COUNT(*) FROM topics WHERE is_mod_deleted = false) AS topic_count, - (SELECT COUNT(*) FROM replies) AS reply_count, + (SELECT COUNT(*) FROM topics WHERE is_mod_deleted = false AND is_author_deleted = false) AS topic_count, + (SELECT COUNT(*) FROM replies WHERE is_author_deleted = false) AS reply_count, (SELECT COUNT(*) FROM users) AS user_count, (SELECT COUNT(*) FROM categories) AS category_count, (SELECT COUNT(*) FROM reports WHERE status = 'pending') AS report_count, - (SELECT COUNT(*) FROM topics WHERE is_mod_deleted = false AND created_at > NOW() - INTERVAL '7 days') AS recent_topics, - (SELECT COUNT(*) FROM replies WHERE created_at > NOW() - INTERVAL '7 days') AS recent_replies, + (SELECT COUNT(*) FROM topics WHERE is_mod_deleted = false AND is_author_deleted = false AND created_at > NOW() - INTERVAL '7 days') AS recent_topics, + (SELECT COUNT(*) FROM replies WHERE is_author_deleted = false AND created_at > NOW() - INTERVAL '7 days') AS recent_replies, (SELECT COUNT(*) FROM users WHERE first_seen_at > NOW() - INTERVAL '7 days') AS recent_users `) diff --git a/src/routes/categories.ts b/src/routes/categories.ts index 9682dfe..4dccc68 100644 --- a/src/routes/categories.ts +++ b/src/routes/categories.ts @@ -1,6 +1,7 @@ import { randomUUID } from 'node:crypto' import { eq, and, count } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { notFound, badRequest, conflict } from '../lib/api-errors.js' import { isMaturityLowerThan } from '../lib/maturity.js' import { @@ -217,7 +218,7 @@ export function categoryRoutes(): FastifyPluginCallback { async (request, reply) => { const parsed = categoryQuerySchema.safeParse(request.query) const parentId = parsed.success ? parsed.data.parentId : undefined - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const conditions = [eq(categories.communityDid, communityDid)] if (parentId !== undefined) { @@ -261,7 +262,7 @@ export function categoryRoutes(): FastifyPluginCallback { }, async (request, reply) => { const { slug } = request.params as { slug: string } - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const rows = await db .select() @@ -328,7 +329,7 @@ export function categoryRoutes(): FastifyPluginCallback { } const { name, slug, description, parentId, sortOrder, maturityRating } = parsed.data - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Fetch community settings for maturity default const settingsRows = await db @@ -453,7 +454,7 @@ export function categoryRoutes(): FastifyPluginCallback { } const updates = parsed.data - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Fetch community settings for maturity validation const settingsRows = await db @@ -587,7 +588,7 @@ export function categoryRoutes(): FastifyPluginCallback { } // Check if category has topics within this community - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const topicCountResult = await db .select({ count: count() }) .from(topics) diff --git a/src/routes/moderation-queue.ts b/src/routes/moderation-queue.ts index 933624c..86cf2d0 100644 --- a/src/routes/moderation-queue.ts +++ b/src/routes/moderation-queue.ts @@ -1,5 +1,6 @@ import { eq, and, desc, sql } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { notFound, badRequest, conflict } from '../lib/api-errors.js' import { wordFilterSchema, queueActionSchema, queueQuerySchema } from '../validation/anti-spam.js' import { moderationQueue } from '../db/schema/moderation-queue.js' @@ -86,7 +87,7 @@ export function moderationQueueRoutes(): FastifyPluginCallback { const { db, env, authMiddleware } = app const requireModerator = createRequireModerator(db, authMiddleware, app.log) const requireAdmin = app.requireAdmin - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // ------------------------------------------------------------------- // GET /api/moderation/queue (moderator+) diff --git a/src/routes/moderation.ts b/src/routes/moderation.ts index 7605684..0140b71 100644 --- a/src/routes/moderation.ts +++ b/src/routes/moderation.ts @@ -1,5 +1,6 @@ import { eq, and, desc, sql } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { notFound, forbidden, badRequest, conflict } from '../lib/api-errors.js' import { lockTopicSchema, @@ -143,7 +144,7 @@ export function moderationRoutes(): FastifyPluginCallback { const { db, env, authMiddleware } = app const requireModerator = createRequireModerator(db, authMiddleware, app.log) const requireAdmin = app.requireAdmin - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const notificationService = createNotificationService(db, app.log) // ------------------------------------------------------------------- diff --git a/src/routes/onboarding.ts b/src/routes/onboarding.ts index eaf0f57..5af06be 100644 --- a/src/routes/onboarding.ts +++ b/src/routes/onboarding.ts @@ -1,5 +1,6 @@ import { eq, and, asc } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { notFound, badRequest, forbidden } from '../lib/api-errors.js' import { createOnboardingFieldSchema, @@ -115,7 +116,7 @@ export function onboardingRoutes(): FastifyPluginCallback { }, }, async (_request, reply) => { - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const fields = await db .select() @@ -165,7 +166,7 @@ export function onboardingRoutes(): FastifyPluginCallback { throw badRequest('Invalid onboarding field data') } - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const inserted = await db .insert(communityOnboardingFields) @@ -247,7 +248,7 @@ export function onboardingRoutes(): FastifyPluginCallback { throw badRequest('At least one field must be provided') } - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const dbUpdates: Record = { updatedAt: new Date() } if (updates.label !== undefined) dbUpdates.label = updates.label @@ -305,7 +306,7 @@ export function onboardingRoutes(): FastifyPluginCallback { }, }, async (request, reply) => { - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const deleted = await db .delete(communityOnboardingFields) @@ -375,7 +376,7 @@ export function onboardingRoutes(): FastifyPluginCallback { throw badRequest('Invalid reorder data') } - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Update each field's sort order for (const item of parsed.data) { @@ -429,7 +430,7 @@ export function onboardingRoutes(): FastifyPluginCallback { throw forbidden('Authentication required') } - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Get all fields for this community const fields = await db @@ -513,7 +514,7 @@ export function onboardingRoutes(): FastifyPluginCallback { throw badRequest('Invalid submission data') } - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Fetch all community fields to validate against const fields = await db diff --git a/src/routes/reactions.ts b/src/routes/reactions.ts index 0e40800..990f090 100644 --- a/src/routes/reactions.ts +++ b/src/routes/reactions.ts @@ -1,5 +1,6 @@ import { eq, and, sql, asc } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { createPdsClient } from '../lib/pds-client.js' import { notFound, forbidden, badRequest, conflict } from '../lib/api-errors.js' import { createReactionSchema, reactionQuerySchema } from '../validation/reactions.js' @@ -9,6 +10,7 @@ import { replies } from '../db/schema/replies.js' import { communitySettings } from '../db/schema/community-settings.js' import { checkOnboardingComplete } from '../lib/onboarding-gate.js' import { createNotificationService } from '../services/notification.js' +import { extractRkey } from '../lib/at-uri.js' // --------------------------------------------------------------------------- // Constants @@ -87,19 +89,6 @@ function decodeCursor(cursor: string): { createdAt: string; uri: string } | null } } -/** - * Extract the rkey from an AT URI. - * Format: at://did:plc:xxx/collection/rkey - */ -function extractRkey(uri: string): string { - const parts = uri.split('/') - const rkey = parts[parts.length - 1] - if (!rkey) { - throw badRequest('Invalid AT URI: missing rkey') - } - return rkey -} - /** * Get the collection NSID from an AT URI. * Format: at://did/collection/rkey -> returns "collection" @@ -180,7 +169,7 @@ export function reactionRoutes(): FastifyPluginCallback { } const { subjectUri, subjectCid, type: reactionType } = parsed.data - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Onboarding gate: block if user hasn't completed mandatory onboarding const onboarding = await checkOnboardingComplete(db, user.did, communityDid) @@ -372,7 +361,7 @@ export function reactionRoutes(): FastifyPluginCallback { const { uri } = request.params as { uri: string } const decodedUri = decodeURIComponent(uri) - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Fetch existing reaction (scoped to this community) const existing = await db @@ -472,7 +461,7 @@ export function reactionRoutes(): FastifyPluginCallback { } const { subjectUri, type: reactionType, cursor, limit } = parsed.data - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const conditions = [ eq(reactions.subjectUri, subjectUri), eq(reactions.communityDid, communityDid), diff --git a/src/routes/replies.ts b/src/routes/replies.ts index 424e93a..2cc8b8b 100644 --- a/src/routes/replies.ts +++ b/src/routes/replies.ts @@ -1,5 +1,6 @@ import { eq, and, sql, asc, notInArray } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { createPdsClient } from '../lib/pds-client.js' import { notFound, forbidden, badRequest } from '../lib/api-errors.js' import { resolveMaxMaturity, maturityAllows } from '../lib/content-filter.js' @@ -24,6 +25,7 @@ import { categories } from '../db/schema/categories.js' import { communitySettings } from '../db/schema/community-settings.js' import { checkOnboardingComplete } from '../lib/onboarding-gate.js' import { createNotificationService } from '../services/notification.js' +import { extractRkey } from '../lib/at-uri.js' // --------------------------------------------------------------------------- // Constants @@ -105,8 +107,8 @@ function serializeReply(row: typeof replies.$inferSelect) { uri: row.uri, rkey: row.rkey, authorDid: row.authorDid, - content: row.content, - contentFormat: row.contentFormat ?? null, + content: row.isAuthorDeleted ? '' : row.content, + contentFormat: row.isAuthorDeleted ? null : (row.contentFormat ?? null), rootUri: row.rootUri, rootCid: row.rootCid, parentUri: row.parentUri, @@ -116,6 +118,7 @@ function serializeReply(row: typeof replies.$inferSelect) { cid: row.cid, depth, reactionCount: row.reactionCount, + isAuthorDeleted: row.isAuthorDeleted, createdAt: row.createdAt.toISOString(), indexedAt: row.indexedAt.toISOString(), } @@ -146,19 +149,6 @@ function decodeCursor(cursor: string): { createdAt: string; uri: string } | null } } -/** - * Extract the rkey from an AT URI. - * Format: at://did:plc:xxx/collection/rkey - */ -function extractRkey(uri: string): string { - const parts = uri.split('/') - const rkey = parts[parts.length - 1] - if (!rkey) { - throw badRequest('Invalid AT URI: missing rkey') - } - return rkey -} - // --------------------------------------------------------------------------- // Reply routes plugin // --------------------------------------------------------------------------- @@ -536,7 +526,7 @@ export function replyRoutes(): FastifyPluginCallback { } // Maturity check: verify the topic's category is within the user's allowed level - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const catRows = await db .select({ maturityRating: categories.maturityRating }) .from(categories) @@ -619,8 +609,8 @@ export function replyRoutes(): FastifyPluginCallback { const ozoneMap = new Map() if (app.ozoneService) { const uniqueDids = [...new Set(serialized.map((r) => r.authorDid))] - for (const did of uniqueDids) { - const isSpam = await app.ozoneService.isSpamLabeled(did) + const spamMap = await app.ozoneService.batchIsSpamLabeled(uniqueDids) + for (const [did, isSpam] of spamMap) { ozoneMap.set(did, isSpam ? 'spam' : null) } } @@ -858,9 +848,12 @@ export function replyRoutes(): FastifyPluginCallback { await pdsClient.deleteRecord(user.did, COLLECTION, rkey) } - // Delete reply and update topic replyCount in a transaction + // Soft-delete reply and update topic replyCount in a transaction await db.transaction(async (tx) => { - await tx.delete(replies).where(eq(replies.uri, decodedUri)) + await tx + .update(replies) + .set({ isAuthorDeleted: true }) + .where(eq(replies.uri, decodedUri)) await tx .update(topics) .set({ diff --git a/src/routes/search.ts b/src/routes/search.ts index a9ad7fc..f4c40e2 100644 --- a/src/routes/search.ts +++ b/src/routes/search.ts @@ -521,6 +521,7 @@ async function searchTopicsFulltext( const conditions: ReturnType[] = [ sql`search_vector @@ websearch_to_tsquery('english', ${query})`, sql`is_mod_deleted = false`, + sql`is_author_deleted = false`, ] if (filters.category) { @@ -619,6 +620,7 @@ async function searchTopicsVector( const conditions: ReturnType[] = [ sql`embedding IS NOT NULL`, sql`is_mod_deleted = false`, + sql`is_author_deleted = false`, sql`embedding <=> ${embeddingStr}::vector < 0.5`, ] @@ -716,6 +718,7 @@ async function countSearchResults( const conditions: ReturnType[] = [ sql`search_vector @@ websearch_to_tsquery('english', ${query})`, sql`is_mod_deleted = false`, + sql`is_author_deleted = false`, ] if (filters.category) { diff --git a/src/routes/topics.ts b/src/routes/topics.ts index 96f53ed..5ba106b 100644 --- a/src/routes/topics.ts +++ b/src/routes/topics.ts @@ -1,5 +1,6 @@ import { eq, and, desc, sql, inArray, notInArray, isNotNull, ne, or } from 'drizzle-orm' import type { FastifyPluginCallback } from 'fastify' +import { getCommunityDid } from '../config/env.js' import { createPdsClient } from '../lib/pds-client.js' import { notFound, forbidden, badRequest } from '../lib/api-errors.js' import { resolveMaxMaturity, allowedRatings, maturityAllows } from '../lib/content-filter.js' @@ -20,12 +21,12 @@ import { import { tooManyRequests } from '../lib/api-errors.js' import { moderationQueue } from '../db/schema/moderation-queue.js' import { topics } from '../db/schema/topics.js' -import { replies } from '../db/schema/replies.js' import { users } from '../db/schema/users.js' import { categories } from '../db/schema/categories.js' import { communitySettings } from '../db/schema/community-settings.js' import { checkOnboardingComplete } from '../lib/onboarding-gate.js' import { createNotificationService } from '../services/notification.js' +import { extractRkey } from '../lib/at-uri.js' // --------------------------------------------------------------------------- // Constants @@ -104,9 +105,9 @@ function serializeTopic(row: typeof topics.$inferSelect, categoryMaturityRating: uri: row.uri, rkey: row.rkey, authorDid: row.authorDid, - title: row.title, - content: row.content, - contentFormat: row.contentFormat ?? null, + title: row.isAuthorDeleted ? '[Deleted by author]' : row.title, + content: row.isAuthorDeleted ? '' : row.content, + contentFormat: row.isAuthorDeleted ? null : (row.contentFormat ?? null), category: row.category, tags: row.tags ?? null, labels: row.labels ?? null, @@ -114,6 +115,7 @@ function serializeTopic(row: typeof topics.$inferSelect, categoryMaturityRating: cid: row.cid, replyCount: row.replyCount, reactionCount: row.reactionCount, + isAuthorDeleted: row.isAuthorDeleted, categoryMaturityRating, lastActivityAt: row.lastActivityAt.toISOString(), createdAt: row.createdAt.toISOString(), @@ -146,19 +148,6 @@ function decodeCursor(cursor: string): { lastActivityAt: string; uri: string } | } } -/** - * Extract the rkey from an AT URI. - * Format: at://did:plc:xxx/collection/rkey - */ -function extractRkey(uri: string): string { - const parts = uri.split('/') - const rkey = parts[parts.length - 1] - if (!rkey) { - throw badRequest('Invalid AT URI: missing rkey') - } - return rkey -} - // --------------------------------------------------------------------------- // Topic routes plugin // --------------------------------------------------------------------------- @@ -262,7 +251,7 @@ export function topicRoutes(): FastifyPluginCallback { const { title, content, category, tags, labels } = parsed.data const now = new Date().toISOString() - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Onboarding gate: block if user hasn't completed mandatory onboarding const onboarding = await checkOnboardingComplete(db, user.did, communityDid) @@ -624,7 +613,7 @@ export function topicRoutes(): FastifyPluginCallback { // Single mode: filter by the one configured community // --------------------------------------------------------------- - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) // Get category slugs matching allowed maturity levels const allowedCategories = await db @@ -652,8 +641,9 @@ export function topicRoutes(): FastifyPluginCallback { conditions.push(inArray(topics.category, allowedSlugs)) } - // Only show approved content in public listings + // Only show approved, non-deleted content in public listings conditions.push(eq(topics.moderationStatus, 'approved')) + conditions.push(eq(topics.isAuthorDeleted, false)) // Block/mute filtering: load the authenticated user's preferences const { blockedDids, mutedDids } = await loadBlockMuteLists(request.user?.did, db) @@ -704,8 +694,8 @@ export function topicRoutes(): FastifyPluginCallback { const ozoneMap = new Map() if (app.ozoneService) { const uniqueDids = [...new Set(serialized.map((t) => t.authorDid))] - for (const did of uniqueDids) { - const isSpam = await app.ozoneService.isSpamLabeled(did) + const spamMap = await app.ozoneService.batchIsSpamLabeled(uniqueDids) + for (const [did, isSpam] of spamMap) { ozoneMap.set(did, isSpam ? 'spam' : null) } } @@ -785,7 +775,7 @@ export function topicRoutes(): FastifyPluginCallback { } // Look up the category maturity rating - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const catRows = await db .select({ maturityRating: categories.maturityRating }) .from(categories) @@ -833,7 +823,7 @@ export function topicRoutes(): FastifyPluginCallback { } // Maturity check: verify the topic's category is within the user's allowed level - const communityDid = env.COMMUNITY_DID ?? 'did:plc:placeholder' + const communityDid = getCommunityDid(env) const catRows = await db .select({ maturityRating: categories.maturityRating }) .from(categories) @@ -1081,11 +1071,8 @@ export function topicRoutes(): FastifyPluginCallback { app.log.warn({ err, topicUri: decodedUri }, 'Failed to delete cross-posts') }) - // Cascade delete in a transaction for consistency - await db.transaction(async (tx) => { - await tx.delete(replies).where(eq(replies.rootUri, decodedUri)) - await tx.delete(topics).where(eq(topics.uri, decodedUri)) - }) + // Soft-delete: mark as author-deleted, preserve replies (they belong to other users) + await db.update(topics).set({ isAuthorDeleted: true }).where(eq(topics.uri, decodedUri)) return await reply.status(204).send() } catch (err: unknown) { diff --git a/src/services/cross-post.ts b/src/services/cross-post.ts index ea2d7c2..7929c77 100644 --- a/src/services/cross-post.ts +++ b/src/services/cross-post.ts @@ -6,6 +6,7 @@ import type { NotificationService } from './notification.js' import { generateOgImage } from './og-image.js' import { crossPosts } from '../db/schema/cross-posts.js' import { userPreferences } from '../db/schema/user-preferences.js' +import { extractRkey } from '../lib/at-uri.js' // --------------------------------------------------------------------------- // Constants @@ -52,15 +53,6 @@ export interface CrossPostConfig { // Helpers // --------------------------------------------------------------------------- -/** - * Extract the rkey from an AT URI. - * Format: at://did:plc:xxx/collection/rkey - */ -function extractRkey(uri: string): string { - const parts = uri.split('/') - return parts[parts.length - 1] ?? '' -} - /** * Truncate text to a maximum number of characters, appending ellipsis if needed. */ diff --git a/src/services/ozone.ts b/src/services/ozone.ts index 9165783..74d3cca 100644 --- a/src/services/ozone.ts +++ b/src/services/ozone.ts @@ -1,4 +1,4 @@ -import { eq, and } from 'drizzle-orm' +import { eq, and, inArray } from 'drizzle-orm' import type { Database } from '../db/index.js' import type { Cache } from '../cache/index.js' import type { Logger } from '../lib/logger.js' @@ -217,4 +217,71 @@ export class OzoneService { const labels = await this.getLabels(didOrUri) return labels.some((l) => SPAM_LABELS.has(l.val)) } + + /** + * Batch check which DIDs have spam-related labels. + * Uses a single DB query for all cache misses instead of N+1 individual queries. + */ + async batchIsSpamLabeled(dids: string[]): Promise> { + const result = new Map() + if (dids.length === 0) return result + + const uncached: string[] = [] + + // Check cache first + for (const did of dids) { + const cacheKey = `${CACHE_PREFIX}${did}` + try { + const cached = await this.cache.get(cacheKey) + if (cached) { + const labels = JSON.parse(cached) as CachedLabel[] + result.set( + did, + labels.some((l) => SPAM_LABELS.has(l.val)) + ) + continue + } + } catch { + // Fall through to DB + } + uncached.push(did) + } + + if (uncached.length === 0) return result + + // Batch DB query for all uncached DIDs + const rows = await this.db + .select({ + uri: ozoneLabels.uri, + val: ozoneLabels.val, + src: ozoneLabels.src, + neg: ozoneLabels.neg, + }) + .from(ozoneLabels) + .where(and(inArray(ozoneLabels.uri, uncached), eq(ozoneLabels.neg, false))) + + // Group by URI + const labelsByUri = new Map() + for (const row of rows) { + const labels = labelsByUri.get(row.uri) ?? [] + labels.push({ val: row.val, src: row.src, neg: row.neg }) + labelsByUri.set(row.uri, labels) + } + + // Cache and build results + for (const did of uncached) { + const labels = labelsByUri.get(did) ?? [] + result.set( + did, + labels.some((l) => SPAM_LABELS.has(l.val)) + ) + try { + await this.cache.set(`${CACHE_PREFIX}${did}`, JSON.stringify(labels), 'EX', CACHE_TTL) + } catch { + // Non-critical + } + } + + return result + } } diff --git a/tests/integration/firehose/record-processing.test.ts b/tests/integration/firehose/record-processing.test.ts index 049be83..a490a9c 100644 --- a/tests/integration/firehose/record-processing.test.ts +++ b/tests/integration/firehose/record-processing.test.ts @@ -154,7 +154,7 @@ describe('firehose record processing (integration)', () => { expect(topic.cid).toBe('bafytopic1v2') }) - it('deletes a topic', async () => { + it('soft-deletes a topic', async () => { await handler.handle(topicEvent) const deleteEvent: RecordEvent = { @@ -169,12 +169,14 @@ describe('firehose record processing (integration)', () => { await handler.handle(deleteEvent) - const result = await db - .select() - .from(topics) - .where(eq(topics.uri, 'at://did:plc:integ-user1/forum.barazo.topic.post/topic1')) + const topic = one( + await db + .select() + .from(topics) + .where(eq(topics.uri, 'at://did:plc:integ-user1/forum.barazo.topic.post/topic1')) + ) - expect(result).toHaveLength(0) + expect(topic.isAuthorDeleted).toBe(true) }) }) diff --git a/tests/integration/health.test.ts b/tests/integration/health.test.ts index 002e94e..fecf97f 100644 --- a/tests/integration/health.test.ts +++ b/tests/integration/health.test.ts @@ -32,6 +32,7 @@ describe('health routes (integration)', () => { LOG_LEVEL: 'silent', CORS_ORIGINS: 'http://localhost:3001', COMMUNITY_MODE: 'single' as const, + COMMUNITY_DID: 'did:plc:testcommunity', COMMUNITY_NAME: 'Test Community', RATE_LIMIT_AUTH: 10, RATE_LIMIT_WRITE: 10, diff --git a/tests/unit/auth/oauth-client.test.ts b/tests/unit/auth/oauth-client.test.ts index 0094784..1081d91 100644 --- a/tests/unit/auth/oauth-client.test.ts +++ b/tests/unit/auth/oauth-client.test.ts @@ -3,9 +3,8 @@ import type { Cache } from '../../../src/cache/index.js' import type { Logger } from '../../../src/lib/logger.js' import type { Env } from '../../../src/config/env.js' -// Track constructor calls and mock event listener +// Track constructor calls const constructorArgs: Record[] = [] -const mockAddEventListener = vi.fn() const mockJwks = { keys: [] } vi.mock('@atproto/oauth-client-node', () => { @@ -13,7 +12,6 @@ vi.mock('@atproto/oauth-client-node', () => { NodeOAuthClient: class MockNodeOAuthClient { clientMetadata: Record jwks: { keys: unknown[] } - addEventListener = mockAddEventListener constructor(options: { clientMetadata: Record }) { constructorArgs.push(options as Record) @@ -215,15 +213,45 @@ describe('createOAuthClient', () => { }) }) - describe('event listeners', () => { - it('registers updated and deleted event listeners', () => { + describe('session lifecycle hooks', () => { + it('passes onUpdate and onDelete hooks in constructor options', () => { const env = createMockEnv() createOAuthClient(env, cacheMocks.cache, logMocks.logger) - expect(mockAddEventListener).toHaveBeenCalledTimes(2) - expect(mockAddEventListener).toHaveBeenCalledWith('updated', expect.any(Function)) - expect(mockAddEventListener).toHaveBeenCalledWith('deleted', expect.any(Function)) + const options = getLastConstructorOptions() + expect(typeof options.onUpdate).toBe('function') + expect(typeof options.onDelete).toBe('function') + }) + + it('onUpdate hook logs session update', () => { + const env = createMockEnv() + + createOAuthClient(env, cacheMocks.cache, logMocks.logger) + + const options = getLastConstructorOptions() + const onUpdate = options.onUpdate as (sub: string) => void + onUpdate('did:plc:test123') + + expect(logMocks.infoFn).toHaveBeenCalledWith( + { sub: 'did:plc:test123' }, + 'OAuth session updated' + ) + }) + + it('onDelete hook logs session deletion with cause', () => { + const env = createMockEnv() + + createOAuthClient(env, cacheMocks.cache, logMocks.logger) + + const options = getLastConstructorOptions() + const onDelete = options.onDelete as (sub: string, cause: unknown) => void + onDelete('did:plc:test123', new Error('token_revoked')) + + expect(logMocks.infoFn).toHaveBeenCalledWith( + { sub: 'did:plc:test123', cause: 'Error: token_revoked' }, + 'OAuth session deleted' + ) }) }) diff --git a/tests/unit/config/env.test.ts b/tests/unit/config/env.test.ts index 2671035..78d3a64 100644 --- a/tests/unit/config/env.test.ts +++ b/tests/unit/config/env.test.ts @@ -1,5 +1,6 @@ import { describe, it, expect } from 'vitest' -import { envSchema, parseEnv } from '../../../src/config/env.js' +import { envSchema, parseEnv, getCommunityDid } from '../../../src/config/env.js' +import type { Env } from '../../../src/config/env.js' describe('envSchema', () => { const validEnv = { @@ -16,6 +17,7 @@ describe('envSchema', () => { LOG_LEVEL: 'info', CORS_ORIGINS: 'http://localhost:3001', COMMUNITY_MODE: 'single', + COMMUNITY_DID: 'did:plc:testcommunity123', } it('parses valid environment variables', () => { @@ -95,6 +97,7 @@ describe('envSchema', () => { OAUTH_CLIENT_ID: validEnv.OAUTH_CLIENT_ID, OAUTH_REDIRECT_URI: validEnv.OAUTH_REDIRECT_URI, SESSION_SECRET: validEnv.SESSION_SECRET, + COMMUNITY_DID: validEnv.COMMUNITY_DID, }) expect(result.success).toBe(true) if (result.success) { @@ -205,6 +208,61 @@ describe('envSchema', () => { }) }) +describe('COMMUNITY_DID validation', () => { + const baseEnv = { + DATABASE_URL: 'postgresql://barazo:barazo_dev@localhost:5432/barazo', + VALKEY_URL: 'redis://localhost:6379', + TAP_URL: 'http://localhost:2480', + TAP_ADMIN_PASSWORD: 'tap_dev_secret', + OAUTH_CLIENT_ID: + 'http://localhost?redirect_uri=http%3A%2F%2F127.0.0.1%3A3000%2Fapi%2Fauth%2Fcallback', + OAUTH_REDIRECT_URI: 'http://127.0.0.1:3000/api/auth/callback', + SESSION_SECRET: 'a-very-long-session-secret-that-is-at-least-32-characters', + } + + it('rejects single mode without COMMUNITY_DID', () => { + const result = envSchema.safeParse({ + ...baseEnv, + COMMUNITY_MODE: 'single', + }) + expect(result.success).toBe(false) + }) + + it('accepts single mode with COMMUNITY_DID', () => { + const result = envSchema.safeParse({ + ...baseEnv, + COMMUNITY_MODE: 'single', + COMMUNITY_DID: 'did:plc:testcommunity', + }) + expect(result.success).toBe(true) + }) + + it('accepts global mode without COMMUNITY_DID', () => { + const result = envSchema.safeParse({ + ...baseEnv, + COMMUNITY_MODE: 'global', + }) + expect(result.success).toBe(true) + }) + + it('rejects default mode (single) without COMMUNITY_DID', () => { + const result = envSchema.safeParse(baseEnv) + expect(result.success).toBe(false) + }) +}) + +describe('getCommunityDid', () => { + it('returns COMMUNITY_DID when set', () => { + const env = { COMMUNITY_DID: 'did:plc:test123' } as Env + expect(getCommunityDid(env)).toBe('did:plc:test123') + }) + + it('throws when COMMUNITY_DID is undefined', () => { + const env = { COMMUNITY_DID: undefined } as Env + expect(() => getCommunityDid(env)).toThrow('COMMUNITY_DID is required') + }) +}) + describe('parseEnv', () => { it('throws on invalid environment', () => { expect(() => parseEnv({})).toThrow() diff --git a/tests/unit/firehose/clamp-timestamp.test.ts b/tests/unit/firehose/clamp-timestamp.test.ts new file mode 100644 index 0000000..f481a0b --- /dev/null +++ b/tests/unit/firehose/clamp-timestamp.test.ts @@ -0,0 +1,52 @@ +import { describe, it, expect } from 'vitest' +import { clampCreatedAt } from '../../../src/firehose/clamp-timestamp.js' + +describe('clampCreatedAt', () => { + const now = new Date('2026-02-19T12:00:00.000Z') + + it('returns client timestamp when within acceptable range', () => { + const clientTime = new Date('2026-02-19T11:30:00.000Z') // 30 min ago + expect(clampCreatedAt(clientTime, now)).toEqual(clientTime) + }) + + it('returns client timestamp when slightly in the future (< 5 min)', () => { + const clientTime = new Date('2026-02-19T12:03:00.000Z') // 3 min ahead + expect(clampCreatedAt(clientTime, now)).toEqual(clientTime) + }) + + it('clamps future timestamps (> 5 min ahead) to now', () => { + const clientTime = new Date('2026-02-19T12:10:00.000Z') // 10 min ahead + expect(clampCreatedAt(clientTime, now)).toEqual(now) + }) + + it('clamps far-future timestamps to now', () => { + const clientTime = new Date('2027-01-01T00:00:00.000Z') // next year + expect(clampCreatedAt(clientTime, now)).toEqual(now) + }) + + it('clamps very old timestamps (> 1 hour ago) to max past', () => { + const clientTime = new Date('2026-02-19T10:00:00.000Z') // 2 hours ago + const maxPast = new Date('2026-02-19T11:00:00.000Z') // 1 hour ago + expect(clampCreatedAt(clientTime, now)).toEqual(maxPast) + }) + + it('clamps extremely old timestamps to max past', () => { + const clientTime = new Date('2020-01-01T00:00:00.000Z') // years ago + const maxPast = new Date('2026-02-19T11:00:00.000Z') + expect(clampCreatedAt(clientTime, now)).toEqual(maxPast) + }) + + it('returns client timestamp at exactly 5 min future boundary', () => { + const clientTime = new Date('2026-02-19T12:05:00.000Z') // exactly 5 min + expect(clampCreatedAt(clientTime, now)).toEqual(clientTime) + }) + + it('returns client timestamp at exactly 1 hour past boundary', () => { + const clientTime = new Date('2026-02-19T11:00:00.000Z') // exactly 1 hour ago + expect(clampCreatedAt(clientTime, now)).toEqual(clientTime) + }) + + it('returns client timestamp at exactly now', () => { + expect(clampCreatedAt(now, now)).toEqual(now) + }) +}) diff --git a/tests/unit/firehose/handlers/record.test.ts b/tests/unit/firehose/handlers/record.test.ts index 4d285f2..f074244 100644 --- a/tests/unit/firehose/handlers/record.test.ts +++ b/tests/unit/firehose/handlers/record.test.ts @@ -463,6 +463,112 @@ describe('RecordHandler', () => { }) }) + describe('delete with DB lookup', () => { + it('resolves rootUri from DB for reply deletes', async () => { + db.select.mockReturnValue({ + from: vi.fn().mockReturnValue({ + where: vi + .fn() + .mockResolvedValue([{ rootUri: 'at://did:plc:test/forum.barazo.topic.post/t1' }]), + }), + }) + + const event: RecordEvent = { + id: 20, + action: 'delete', + did: 'did:plc:test', + rev: 'rev3', + collection: 'forum.barazo.topic.reply', + rkey: 'reply1', + live: true, + } + + await handler.handle(event) + + expect(replyIndexer.handleDelete).toHaveBeenCalledWith({ + uri: 'at://did:plc:test/forum.barazo.topic.reply/reply1', + rkey: 'reply1', + did: 'did:plc:test', + rootUri: 'at://did:plc:test/forum.barazo.topic.post/t1', + }) + }) + + it('resolves subjectUri from DB for reaction deletes', async () => { + db.select.mockReturnValue({ + from: vi.fn().mockReturnValue({ + where: vi + .fn() + .mockResolvedValue([{ subjectUri: 'at://did:plc:test/forum.barazo.topic.post/t1' }]), + }), + }) + + const event: RecordEvent = { + id: 21, + action: 'delete', + did: 'did:plc:test', + rev: 'rev3', + collection: 'forum.barazo.interaction.reaction', + rkey: 'react1', + live: true, + } + + await handler.handle(event) + + expect(reactionIndexer.handleDelete).toHaveBeenCalledWith({ + uri: 'at://did:plc:test/forum.barazo.interaction.reaction/react1', + rkey: 'react1', + did: 'did:plc:test', + subjectUri: 'at://did:plc:test/forum.barazo.topic.post/t1', + }) + }) + + it('passes empty rootUri when reply not found in DB', async () => { + // Default mock returns [] — reply already deleted or never indexed + const event: RecordEvent = { + id: 22, + action: 'delete', + did: 'did:plc:test', + rev: 'rev3', + collection: 'forum.barazo.topic.reply', + rkey: 'missing1', + live: true, + } + + await handler.handle(event) + + expect(replyIndexer.handleDelete).toHaveBeenCalledWith({ + uri: 'at://did:plc:test/forum.barazo.topic.reply/missing1', + rkey: 'missing1', + did: 'did:plc:test', + rootUri: '', + }) + expect(logger.debug).toHaveBeenCalled() + }) + + it('passes empty subjectUri when reaction not found in DB', async () => { + // Default mock returns [] + const event: RecordEvent = { + id: 23, + action: 'delete', + did: 'did:plc:test', + rev: 'rev3', + collection: 'forum.barazo.interaction.reaction', + rkey: 'missing1', + live: true, + } + + await handler.handle(event) + + expect(reactionIndexer.handleDelete).toHaveBeenCalledWith({ + uri: 'at://did:plc:test/forum.barazo.interaction.reaction/missing1', + rkey: 'missing1', + did: 'did:plc:test', + subjectUri: '', + }) + expect(logger.debug).toHaveBeenCalled() + }) + }) + describe('live flag', () => { it('passes live flag through to indexer', async () => { const event: RecordEvent = { diff --git a/tests/unit/firehose/indexers/topic.test.ts b/tests/unit/firehose/indexers/topic.test.ts index 196b85b..7678eb3 100644 --- a/tests/unit/firehose/indexers/topic.test.ts +++ b/tests/unit/firehose/indexers/topic.test.ts @@ -103,7 +103,7 @@ describe('TopicIndexer', () => { }) describe('handleDelete', () => { - it('deletes a topic by URI', async () => { + it('soft-deletes a topic by URI', async () => { const db = createMockDb() const logger = createMockLogger() const indexer = new TopicIndexer(db as never, logger as never) @@ -114,7 +114,8 @@ describe('TopicIndexer', () => { did: baseParams.did, }) - expect(db.delete).toHaveBeenCalledTimes(1) + expect(db.update).toHaveBeenCalledTimes(1) + expect(db.delete).not.toHaveBeenCalled() }) }) }) diff --git a/tests/unit/firehose/service.test.ts b/tests/unit/firehose/service.test.ts index 358afd6..082a191 100644 --- a/tests/unit/firehose/service.test.ts +++ b/tests/unit/firehose/service.test.ts @@ -74,6 +74,7 @@ function createMinimalEnv(): Env { LOG_LEVEL: 'silent', CORS_ORIGINS: 'http://localhost:3001', COMMUNITY_MODE: 'single' as const, + COMMUNITY_DID: 'did:plc:testcommunity', COMMUNITY_NAME: 'Test Community', RATE_LIMIT_AUTH: 10, RATE_LIMIT_WRITE: 10, diff --git a/tests/unit/routes/health.test.ts b/tests/unit/routes/health.test.ts index fcfa1aa..ad252c7 100644 --- a/tests/unit/routes/health.test.ts +++ b/tests/unit/routes/health.test.ts @@ -63,6 +63,7 @@ describe('health routes', () => { LOG_LEVEL: 'silent', CORS_ORIGINS: 'http://localhost:3001', COMMUNITY_MODE: 'single' as const, + COMMUNITY_DID: 'did:plc:testcommunity', COMMUNITY_NAME: 'Test Community', RATE_LIMIT_AUTH: 10, RATE_LIMIT_WRITE: 10, diff --git a/tests/unit/routes/replies.test.ts b/tests/unit/routes/replies.test.ts index 7e3c423..df4ec16 100644 --- a/tests/unit/routes/replies.test.ts +++ b/tests/unit/routes/replies.test.ts @@ -894,7 +894,7 @@ describe('reply routes', () => { mutedDids: [], }, ]) - // eslint-disable-next-line @typescript-eslint/no-misused-promises -- thenable mock for Drizzle chain + selectChain.where.mockImplementationOnce(() => selectChain) // 6: replies .where const rows = [ @@ -1238,7 +1238,7 @@ describe('reply routes', () => { deleteRecordFn.mockResolvedValue(undefined) }) - it('deletes a reply when user is the author (deletes from PDS + DB)', async () => { + it('soft-deletes a reply when user is the author (deletes from PDS, soft-deletes in DB)', async () => { const existingRow = sampleReplyRow() selectChain.where.mockResolvedValueOnce([existingRow]) @@ -1255,14 +1255,12 @@ describe('reply routes', () => { expect(deleteRecordFn).toHaveBeenCalledOnce() expect(deleteRecordFn.mock.calls[0]?.[0]).toBe(TEST_DID) - // Should have deleted from DB - expect(mockDb.delete).toHaveBeenCalled() - - // Should have decremented topic replyCount + // Should have soft-deleted in DB (update, not delete) + decremented replyCount expect(mockDb.update).toHaveBeenCalled() + expect(mockDb.delete).not.toHaveBeenCalled() }) - it('deletes reply as moderator (index-only delete, not from PDS)', async () => { + it('soft-deletes reply as moderator (index-only soft-delete, not from PDS)', async () => { const modApp = await buildTestApp(testUser({ did: MOD_DID, handle: 'mod.bsky.social' })) const existingRow = sampleReplyRow({ authorDid: OTHER_DID }) @@ -1280,7 +1278,9 @@ describe('reply routes', () => { expect(response.statusCode).toBe(204) expect(deleteRecordFn).not.toHaveBeenCalled() - expect(mockDb.delete).toHaveBeenCalled() + // Soft-delete: update instead of delete + expect(mockDb.update).toHaveBeenCalled() + expect(mockDb.delete).not.toHaveBeenCalled() await modApp.close() }) diff --git a/tests/unit/routes/topics-replies-integration.test.ts b/tests/unit/routes/topics-replies-integration.test.ts index c1113a8..1de13b3 100644 --- a/tests/unit/routes/topics-replies-integration.test.ts +++ b/tests/unit/routes/topics-replies-integration.test.ts @@ -385,7 +385,7 @@ describe('topics + replies cross-endpoint integration', () => { deleteRecordFn.mockResolvedValue(undefined) }) - it('deleting a reply decrements replyCount using GREATEST', async () => { + it('soft-deleting a reply decrements replyCount using GREATEST', async () => { // Mock: reply lookup for delete const existingReply = sampleReplyRow() selectChain.where.mockResolvedValueOnce([existingReply]) @@ -399,17 +399,17 @@ describe('topics + replies cross-endpoint integration', () => { expect(response.statusCode).toBe(204) - // Verify: reply was deleted from DB - expect(mockDb.delete).toHaveBeenCalled() + // Verify: reply was soft-deleted (update, not delete) + expect(mockDb.update).toHaveBeenCalled() + expect(mockDb.delete).not.toHaveBeenCalled() // Verify: topic replyCount was decremented - expect(mockDb.update).toHaveBeenCalled() expect(updateChain.set).toHaveBeenCalled() - const setCall = updateChain.set.mock.calls[0]?.[0] as Record - expect(setCall).toBeDefined() - // replyCount should be a SQL expression with GREATEST - expect(setCall.replyCount).toBeDefined() + // One set call is the soft-delete (isAuthorDeleted), another is the replyCount decrement + const setCalls = updateChain.set.mock.calls as Array<[Record]> + const replyCountCall = setCalls.find((c) => c[0].replyCount !== undefined) + expect(replyCountCall).toBeDefined() }) it('delete reply also deletes from PDS when user is author', async () => { @@ -437,7 +437,7 @@ describe('topics + replies cross-endpoint integration', () => { // Delete topic cascades replies // ========================================================================= - describe('delete topic cascades replies', () => { + describe('delete topic preserves replies (soft-delete)', () => { let app: FastifyInstance beforeAll(async () => { @@ -454,9 +454,8 @@ describe('topics + replies cross-endpoint integration', () => { deleteRecordFn.mockResolvedValue(undefined) }) - it('deleting a topic deletes all its replies via transaction', async () => { + it('soft-deletes a topic without cascade-deleting replies', async () => { const existingTopic = sampleTopicRow() - // Topic lookup selectChain.where.mockResolvedValueOnce([existingTopic]) const encodedUri = encodeURIComponent(TEST_TOPIC_URI) @@ -471,16 +470,13 @@ describe('topics + replies cross-endpoint integration', () => { // Should have deleted from PDS (author delete) expect(deleteRecordFn).toHaveBeenCalledOnce() - // Transaction should have been used - expect(mockDb.transaction).toHaveBeenCalledOnce() - - // Inside the transaction, both replies and topic should be deleted. - // The transaction mock calls fn(mockDb), so mockDb.delete is called for: - // 1. replies (cascade), 2. topic itself, 3. cross-posts cleanup (fire-and-forget) - expect(mockDb.delete).toHaveBeenCalledTimes(3) + // Soft-delete: update isAuthorDeleted, no cascade delete of replies + expect(mockDb.update).toHaveBeenCalled() + // No transaction needed (no cascade) + expect(mockDb.transaction).not.toHaveBeenCalled() }) - it('topic cascade delete uses rootUri to find related replies', async () => { + it('soft-delete sets isAuthorDeleted on topic only', async () => { const existingTopic = sampleTopicRow() selectChain.where.mockResolvedValueOnce([existingTopic]) @@ -493,11 +489,10 @@ describe('topics + replies cross-endpoint integration', () => { expect(response.statusCode).toBe(204) - // Verify delete was called (replies, topic, and cross-posts cleanup) - expect(mockDb.delete).toHaveBeenCalledTimes(3) - - // All delete calls should have used .where() - expect(deleteChain.where).toHaveBeenCalledTimes(3) + // Verify the update was called with isAuthorDeleted + expect(updateChain.set).toHaveBeenCalled() + const setCall = updateChain.set.mock.calls[0]?.[0] as Record + expect(setCall.isAuthorDeleted).toBe(true) }) }) @@ -754,7 +749,7 @@ describe('topics + replies cross-endpoint integration', () => { resetAllDbMocks() deleteRecordFn.mockResolvedValue(undefined) - // Step 4: Delete topic (should cascade) + // Step 4: Delete topic (soft-delete, no cascade) selectChain.where.mockResolvedValueOnce([sampleTopicRow()]) const deleteResponse = await app.inject({ @@ -765,10 +760,8 @@ describe('topics + replies cross-endpoint integration', () => { expect(deleteResponse.statusCode).toBe(204) - // Should have cascade-deleted replies via transaction - // Delete count: 1. replies (cascade), 2. topic, 3. cross-posts cleanup - expect(mockDb.transaction).toHaveBeenCalledOnce() - expect(mockDb.delete).toHaveBeenCalledTimes(3) + // Soft-delete: update isAuthorDeleted, replies preserved + expect(mockDb.update).toHaveBeenCalled() }) }) }) diff --git a/tests/unit/routes/topics.test.ts b/tests/unit/routes/topics.test.ts index 052895b..daa6a72 100644 --- a/tests/unit/routes/topics.test.ts +++ b/tests/unit/routes/topics.test.ts @@ -1214,8 +1214,10 @@ describe('topic routes', () => { expect(deleteRecordFn).toHaveBeenCalledOnce() expect(deleteRecordFn.mock.calls[0]?.[0]).toBe(TEST_DID) - // Should have deleted from DB (replies + topics) - expect(mockDb.delete).toHaveBeenCalled() + // Should have soft-deleted topic in DB via update (not hard-delete) + expect(mockDb.update).toHaveBeenCalled() + // No transaction needed (no cascade delete of replies) + expect(mockDb.transaction).not.toHaveBeenCalled() }) it('deletes topic as moderator (index-only delete, not from PDS)', async () => { @@ -1239,8 +1241,8 @@ describe('topic routes', () => { // Moderator should NOT delete from PDS expect(deleteRecordFn).not.toHaveBeenCalled() - // But should delete from DB index - expect(mockDb.delete).toHaveBeenCalled() + // But should soft-delete in DB index + expect(mockDb.update).toHaveBeenCalled() await modApp.close() })