diff --git a/src/modules/feeds/infrastructure/http/controllers/GetGlobalFeedController.ts b/src/modules/feeds/infrastructure/http/controllers/GetGlobalFeedController.ts new file mode 100644 index 00000000..8c404c87 --- /dev/null +++ b/src/modules/feeds/infrastructure/http/controllers/GetGlobalFeedController.ts @@ -0,0 +1,31 @@ +import { Request, Response } from 'express'; +import { GetGlobalFeedUseCase } from '../../../application/useCases/queries/GetGlobalFeedUseCase'; +import { BaseController } from '../../../../../shared/infrastructure/http/BaseController'; + +export class GetGlobalFeedController extends BaseController { + constructor(private getGlobalFeedUseCase: GetGlobalFeedUseCase) { + super(); + } + + async executeImpl(req: Request, res: Response): Promise { + try { + const page = parseInt(req.query.page as string) || 1; + const limit = parseInt(req.query.limit as string) || 20; + const beforeActivityId = req.query.beforeActivityId as string; + + const result = await this.getGlobalFeedUseCase.execute({ + page, + limit, + beforeActivityId, + }); + + if (result.isErr()) { + return this.fail(res, result.error.message); + } + + return this.ok(res, result.value); + } catch (error) { + return this.fail(res, 'An unexpected error occurred'); + } + } +} diff --git a/src/modules/feeds/infrastructure/http/routes/feedRoutes.ts b/src/modules/feeds/infrastructure/http/routes/feedRoutes.ts new file mode 100644 index 00000000..2b41ea83 --- /dev/null +++ b/src/modules/feeds/infrastructure/http/routes/feedRoutes.ts @@ -0,0 +1,20 @@ +import { Router } from 'express'; +import { GetGlobalFeedController } from '../controllers/GetGlobalFeedController'; +import { AuthMiddleware } from '../../../../../shared/infrastructure/http/middleware/AuthMiddleware'; + +export function createFeedRoutes( + authMiddleware: AuthMiddleware, + getGlobalFeedController: GetGlobalFeedController, +): Router { + const router = Router(); + + // Apply authentication middleware to all feed routes + router.use(authMiddleware.ensureAuthenticated()); + + // GET /api/feeds/global - Get global feed + router.get('/global', (req, res) => + getGlobalFeedController.execute(req, res), + ); + + return router; +} diff --git a/src/modules/feeds/tests/infrastructure/InMemoryFeedRepository.ts b/src/modules/feeds/tests/infrastructure/InMemoryFeedRepository.ts new file mode 100644 index 00000000..ceafe485 --- /dev/null +++ b/src/modules/feeds/tests/infrastructure/InMemoryFeedRepository.ts @@ -0,0 +1,83 @@ +import { Result, ok, err } from '../../../../shared/core/Result'; +import { + IFeedRepository, + FeedQueryOptions, + PaginatedFeedResult, +} from '../../domain/IFeedRepository'; +import { FeedActivity } from '../../domain/FeedActivity'; +import { ActivityId } from '../../domain/value-objects/ActivityId'; + +export class InMemoryFeedRepository implements IFeedRepository { + private activities: FeedActivity[] = []; + + async addActivity(activity: FeedActivity): Promise> { + try { + this.activities.push(activity); + // Sort by creation time descending + this.activities.sort((a, b) => + b.createdAt.getTime() - a.createdAt.getTime() + ); + return ok(undefined); + } catch (error) { + return err(error as Error); + } + } + + async getGlobalFeed( + options: FeedQueryOptions, + ): Promise> { + try { + const { page, limit, beforeActivityId } = options; + let filteredActivities = [...this.activities]; + + // Filter by cursor if provided + if (beforeActivityId) { + const beforeIndex = filteredActivities.findIndex( + (activity) => activity.activityId.equals(beforeActivityId) + ); + if (beforeIndex >= 0) { + filteredActivities = filteredActivities.slice(beforeIndex + 1); + } + } + + // Paginate + const offset = (page - 1) * limit; + const paginatedActivities = filteredActivities.slice(offset, offset + limit); + + const totalCount = this.activities.length; + const hasMore = offset + paginatedActivities.length < totalCount; + + let nextCursor: ActivityId | undefined; + if (hasMore && paginatedActivities.length > 0) { + nextCursor = paginatedActivities[paginatedActivities.length - 1]!.activityId; + } + + return ok({ + activities: paginatedActivities, + totalCount, + hasMore, + nextCursor, + }); + } catch (error) { + return err(error as Error); + } + } + + async findById(activityId: ActivityId): Promise> { + try { + const activity = this.activities.find((a) => a.activityId.equals(activityId)); + return ok(activity || null); + } catch (error) { + return err(error as Error); + } + } + + // Test helper methods + clear(): void { + this.activities = []; + } + + getAll(): FeedActivity[] { + return [...this.activities]; + } +} diff --git a/src/shared/infrastructure/http/app.ts b/src/shared/infrastructure/http/app.ts index 6cdbf07a..7740e064 100644 --- a/src/shared/infrastructure/http/app.ts +++ b/src/shared/infrastructure/http/app.ts @@ -4,6 +4,7 @@ import { Router } from 'express'; import { createUserRoutes } from '../../../modules/user/infrastructure/http/routes/userRoutes'; import { createAtprotoRoutes } from '../../../modules/atproto/infrastructure/atprotoRoutes'; import { createCardsModuleRoutes } from '../../../modules/cards/infrastructure/http/routes'; +import { createFeedRoutes } from '../../../modules/feeds/infrastructure/http/routes/feedRoutes'; import { EnvironmentConfigService } from '../config/EnvironmentConfigService'; import { RepositoryFactory } from './factories/RepositoryFactory'; import { ServiceFactory } from './factories/ServiceFactory'; @@ -71,10 +72,16 @@ export const createExpressApp = ( controllers.getMyCollectionsController, ); + const feedRouter = createFeedRoutes( + services.authMiddleware, + controllers.getGlobalFeedController, + ); + // Register routes app.use('/api/users', userRouter); app.use('/atproto', atprotoRouter); app.use('/api', cardsRouter); + app.use('/api/feeds', feedRouter); return app; }; diff --git a/src/shared/infrastructure/http/factories/ControllerFactory.ts b/src/shared/infrastructure/http/factories/ControllerFactory.ts index 3cf6050b..83f59a6c 100644 --- a/src/shared/infrastructure/http/factories/ControllerFactory.ts +++ b/src/shared/infrastructure/http/factories/ControllerFactory.ts @@ -16,6 +16,7 @@ import { UpdateCollectionController } from '../../../../modules/cards/infrastruc import { DeleteCollectionController } from '../../../../modules/cards/infrastructure/http/controllers/DeleteCollectionController'; import { GetCollectionPageController } from '../../../../modules/cards/infrastructure/http/controllers/GetCollectionPageController'; import { GetMyCollectionsController } from '../../../../modules/cards/infrastructure/http/controllers/GetMyCollectionsController'; +import { GetGlobalFeedController } from '../../../../modules/feeds/infrastructure/http/controllers/GetGlobalFeedController'; import { UseCases } from './UseCaseFactory'; import { GetMyProfileController } from 'src/modules/cards/infrastructure/http/controllers/GetMyProfileController'; import { LoginWithAppPasswordController } from 'src/modules/user/infrastructure/http/controllers/LoginWithAppPasswordController'; @@ -45,6 +46,8 @@ export interface Controllers { deleteCollectionController: DeleteCollectionController; getCollectionPageController: GetCollectionPageController; getMyCollectionsController: GetMyCollectionsController; + // Feed controllers + getGlobalFeedController: GetGlobalFeedController; } export class ControllerFactory { diff --git a/src/shared/infrastructure/http/factories/RepositoryFactory.ts b/src/shared/infrastructure/http/factories/RepositoryFactory.ts index 81b073ef..d37f18a6 100644 --- a/src/shared/infrastructure/http/factories/RepositoryFactory.ts +++ b/src/shared/infrastructure/http/factories/RepositoryFactory.ts @@ -29,6 +29,9 @@ import { NodeSavedStateStore, NodeSavedSessionStore, } from '@atproto/oauth-client-node'; +import { DrizzleFeedRepository } from '../../../../modules/feeds/infrastructure/repositories/DrizzleFeedRepository'; +import { InMemoryFeedRepository } from '../../../../modules/feeds/tests/infrastructure/InMemoryFeedRepository'; +import { IFeedRepository } from '../../../../modules/feeds/domain/IFeedRepository'; export interface Repositories { userRepository: IUserRepository; @@ -38,6 +41,7 @@ export interface Repositories { collectionRepository: ICollectionRepository; collectionQueryRepository: ICollectionQueryRepository; appPasswordSessionRepository: IAppPasswordSessionRepository; + feedRepository: IFeedRepository; oauthStateStore: NodeSavedStateStore; oauthSessionStore: NodeSavedSessionStore; } @@ -61,6 +65,7 @@ export class RepositoryFactory { ); const appPasswordSessionRepository = new InMemoryAppPasswordSessionRepository(); + const feedRepository = new InMemoryFeedRepository(); const oauthStateStore = new InMemoryStateStore(); const oauthSessionStore = new InMemorySessionStore(); @@ -72,6 +77,7 @@ export class RepositoryFactory { collectionRepository, collectionQueryRepository, appPasswordSessionRepository, + feedRepository, oauthStateStore, oauthSessionStore, }; @@ -92,6 +98,7 @@ export class RepositoryFactory { collectionRepository: new DrizzleCollectionRepository(db), collectionQueryRepository: new DrizzleCollectionQueryRepository(db), appPasswordSessionRepository: new DrizzleAppPasswordSessionRepository(db), + feedRepository: new DrizzleFeedRepository(db), oauthStateStore, oauthSessionStore, }; diff --git a/src/shared/infrastructure/http/factories/ServiceFactory.ts b/src/shared/infrastructure/http/factories/ServiceFactory.ts index b38feb40..b9d06617 100644 --- a/src/shared/infrastructure/http/factories/ServiceFactory.ts +++ b/src/shared/infrastructure/http/factories/ServiceFactory.ts @@ -44,6 +44,8 @@ import { IEventPublisher } from '../../../application/events/IEventPublisher'; import { QueueName } from '../../events/QueueConfig'; import { RedisFactory } from '../../redis/RedisFactory'; import { IEventSubscriber } from 'src/shared/application/events/IEventSubscriber'; +import { FeedService } from '../../../../modules/feeds/domain/services/FeedService'; +import { CardCollectionSaga } from '../../../../modules/feeds/application/sagas/CardCollectionSaga'; // Shared services needed by both web app and workers export interface SharedServices { @@ -52,6 +54,7 @@ export interface SharedServices { atProtoAgentService: IAgentService; metadataService: IMetadataService; profileService: IProfileService; + feedService: FeedService; } // Web app specific services (includes publishers, auth middleware) @@ -72,6 +75,7 @@ export interface WorkerServices extends SharedServices { redisConnection: Redis; eventPublisher: IEventPublisher; createEventSubscriber: (queueName: QueueName) => IEventSubscriber; + cardCollectionSaga: CardCollectionSaga; } // Legacy interface for backward compatibility @@ -185,11 +189,18 @@ export class ServiceFactory { return new BullMQEventSubscriber(redisConnection, { queueName }); }; + // Create saga for worker + const cardCollectionSaga = new CardCollectionSaga( + // We'll need to create this use case in the worker context + null as any // Will be set properly in worker + ); + return { ...sharedServices, redisConnection, eventPublisher, createEventSubscriber, + cardCollectionSaga, }; } @@ -239,12 +250,16 @@ export class ServiceFactory { ? new FakeBlueskyProfileService() : new BlueskyProfileService(atProtoAgentService); + // Feed Service + const feedService = new FeedService(repositories.feedRepository); + return { tokenService, userAuthService, atProtoAgentService, metadataService, profileService, + feedService, }; } } diff --git a/src/shared/infrastructure/http/factories/UseCaseFactory.ts b/src/shared/infrastructure/http/factories/UseCaseFactory.ts index b63e2aac..8eea2db4 100644 --- a/src/shared/infrastructure/http/factories/UseCaseFactory.ts +++ b/src/shared/infrastructure/http/factories/UseCaseFactory.ts @@ -21,6 +21,8 @@ import { Services } from './ServiceFactory'; import { GetMyProfileUseCase } from 'src/modules/cards/application/useCases/queries/GetMyProfileUseCase'; import { LoginWithAppPasswordUseCase } from 'src/modules/user/application/use-cases/LoginWithAppPasswordUseCase'; import { GenerateExtensionTokensUseCase } from 'src/modules/user/application/use-cases/GenerateExtensionTokensUseCase'; +import { GetGlobalFeedUseCase } from '../../../../modules/feeds/application/useCases/queries/GetGlobalFeedUseCase'; +import { AddActivityToFeedUseCase } from '../../../../modules/feeds/application/useCases/commands/AddActivityToFeedUseCase'; export interface UseCases { // User use cases @@ -46,6 +48,9 @@ export interface UseCases { deleteCollectionUseCase: DeleteCollectionUseCase; getCollectionPageUseCase: GetCollectionPageUseCase; getMyCollectionsUseCase: GetMyCollectionsUseCase; + // Feed use cases + getGlobalFeedUseCase: GetGlobalFeedUseCase; + addActivityToFeedUseCase: AddActivityToFeedUseCase; } export class UseCaseFactory { @@ -142,6 +147,17 @@ export class UseCaseFactory { repositories.collectionQueryRepository, services.profileService, ), + + // Feed use cases + getGlobalFeedUseCase: new GetGlobalFeedUseCase( + repositories.feedRepository, + services.profileService, + repositories.cardQueryRepository, + repositories.collectionRepository, + ), + addActivityToFeedUseCase: new AddActivityToFeedUseCase( + services.feedService, + ), }; } } diff --git a/src/workers/feed-worker.ts b/src/workers/feed-worker.ts index 5ead1d2e..25c8c29d 100644 --- a/src/workers/feed-worker.ts +++ b/src/workers/feed-worker.ts @@ -1,9 +1,11 @@ import { EnvironmentConfigService } from '../shared/infrastructure/config/EnvironmentConfigService'; import { RepositoryFactory } from '../shared/infrastructure/http/factories/RepositoryFactory'; import { ServiceFactory } from '../shared/infrastructure/http/factories/ServiceFactory'; +import { UseCaseFactory } from '../shared/infrastructure/http/factories/UseCaseFactory'; import { CardAddedToLibraryEventHandler } from '../modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler'; import { CardAddedToCollectionEventHandler } from '../modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler'; import { CollectionCreatedEventHandler } from '../modules/feeds/application/eventHandlers/CollectionCreatedEventHandler'; +import { CardCollectionSaga } from '../modules/feeds/application/sagas/CardCollectionSaga'; import { QueueNames } from '../shared/infrastructure/events/QueueConfig'; import { EventNames } from '../shared/infrastructure/events/EventConfig'; @@ -15,6 +17,7 @@ async function startFeedWorker() { // Create dependencies using factories const repositories = RepositoryFactory.create(configService); const services = ServiceFactory.createForWorker(configService, repositories); + const useCases = UseCaseFactory.create(repositories, services); // Test Redis connection try { @@ -28,19 +31,22 @@ async function startFeedWorker() { // Create subscriber for feeds queue const eventSubscriber = services.createEventSubscriber(QueueNames.FEEDS); - // Create event handlers with proper dependencies - const feedService = { - processCardAddedToLibrary: async (event: any) => { - console.log('Processing feed update for card added to library:', event); - // Your feed logic here - you can access all shared services here - // services.profileService, services.metadataService, etc. - return { isOk: () => true, isErr: () => false }; - }, - }; + // Create the saga with proper dependencies + const cardCollectionSaga = new CardCollectionSaga( + useCases.addActivityToFeedUseCase, + ); - const cardAddedToLibraryHandler = new CardAddedToLibraryEventHandler(feedService as any); - const cardAddedToCollectionHandler = new CardAddedToCollectionEventHandler(feedService as any); - const collectionCreatedHandler = new CollectionCreatedEventHandler(feedService as any); + // Create event handlers with the saga + const cardAddedToLibraryHandler = new CardAddedToLibraryEventHandler(cardCollectionSaga); + const cardAddedToCollectionHandler = new CardAddedToCollectionEventHandler(cardCollectionSaga); + + // For collection created, we'll create a simple handler for now + const collectionCreatedHandler = new CollectionCreatedEventHandler({ + handleCardEvent: async (event: any) => { + console.log('Processing collection created event:', event); + return { isOk: () => true, isErr: () => false }; + } + } as any); // Register handlers await eventSubscriber.subscribe(