app.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225
  1. import cors from 'cors';
  2. import express, { Router, type ErrorRequestHandler, type RequestHandler } from 'express';
  3. import type { Pool } from 'pg';
  4. import { ZodError } from 'zod';
  5. import type { AppConfig } from './config/env.js';
  6. import { ApiError } from './http/api-error.js';
  7. import { createDomesticVocRouter, type DomesticSnapshotProvider } from './modules/domestic-voc/routes.js';
  8. import { SyncJobRepository, type SyncJobStore } from './modules/domestic-voc/repositories/sync-job.repository.js';
  9. import { SnapshotService } from './modules/domestic-voc/services/snapshot.service.js';
  10. import { SyncService } from './modules/domestic-voc/services/sync.service.js';
  11. import { createAuthenticationMiddleware, createAuthenticator, WorkspaceAccessService } from './modules/saas-platform/auth.js';
  12. import type { PlatformRepository } from './modules/saas-platform/domain.js';
  13. import { PostgresPlatformRepository } from './modules/saas-platform/postgres-platform.repository.js';
  14. import { createSaasPlatformRouter } from './modules/saas-platform/routes.js';
  15. import { FmodeAiClient } from './modules/ai-gateway/client.js';
  16. import { createAiGatewayRouter } from './modules/ai-gateway/routes.js';
  17. import type { AiPromptConfigStore } from './modules/ai-gateway/prompt-config.repository.js';
  18. import { createProductKnowledgeRouter } from './modules/product-knowledge/routes.js';
  19. import type { ProductKnowledgeStore } from './modules/product-knowledge/product-knowledge.store.js';
  20. import { InMemoryListingAiRepository } from './modules/listing-ai/repositories/in-memory-listing-ai.repository.js';
  21. import { FmodeJdVocAiScoringProvider, FmodeListingAiScoringProvider, ListingAiService } from './modules/listing-ai/listing-ai.service.js';
  22. import { FmodeGeminiImageReviewProvider } from './modules/listing-ai/image-review/gemini-image-review.provider.js';
  23. import { createListingAiRouter } from './modules/listing-ai/routes.js';
  24. import type { ListingAiRepository } from './modules/listing-ai/domain.js';
  25. import type { CompetitorListingMonitorService } from './modules/competitor-listing-monitor/competitor-listing-monitor.service.js';
  26. import { createCompetitorListingMonitorRouter } from './modules/competitor-listing-monitor/routes.js';
  27. import { createCloudFunctionRouter, dispatchThroughRouter } from './cloud-functions/router.js';
  28. import { createSpecialActionHandler } from './cloud-functions/special-actions.js';
  29. import { FmodeVocEcommerceClient } from './modules/domestic-voc/upstream/fmode-client.js';
  30. export function createApp(input: {
  31. config: AppConfig;
  32. pool?: Pool;
  33. parseApp?: RequestHandler;
  34. platformRepository?: PlatformRepository;
  35. jobs?: SyncJobStore;
  36. snapshot?: DomesticSnapshotProvider;
  37. healthCheck?: () => Promise<{ ready: boolean; missingObjects: string[] }>;
  38. aiPromptConfigs?: AiPromptConfigStore;
  39. productKnowledge?: ProductKnowledgeStore;
  40. listingAiRepository?: ListingAiRepository;
  41. listingAiService?: ListingAiService;
  42. competitorListingMonitor?: CompetitorListingMonitorService;
  43. }) {
  44. const app = express();
  45. app.disable('x-powered-by');
  46. app.use(cors({
  47. origin(origin, callback) {
  48. if (!origin || input.config.corsOrigins.includes(origin)) return callback(null, true);
  49. return callback(new Error('Origin is not allowed'));
  50. },
  51. credentials: true,
  52. }));
  53. app.use(express.json({ limit: '1mb' }));
  54. app.get('/health', async (_request, response) => {
  55. try {
  56. if (input.healthCheck) {
  57. const readiness = await input.healthCheck();
  58. if (!readiness.ready) {
  59. response.status(503).json({
  60. service: 'saas-voc-server',
  61. status: 'unavailable',
  62. database: readiness.missingObjects.length ? 'migration_required' : 'unavailable',
  63. missingTables: readiness.missingObjects,
  64. storage: input.config.storageDriver,
  65. timestamp: new Date().toISOString(),
  66. });
  67. return;
  68. }
  69. response.json({
  70. service: 'saas-voc-server',
  71. status: 'ok',
  72. database: 'ready',
  73. storage: input.config.storageDriver,
  74. timestamp: new Date().toISOString(),
  75. });
  76. return;
  77. }
  78. if (!input.pool) throw new Error('Database pool is not configured');
  79. const result = await input.pool.query<Record<string, string | null>>(`
  80. SELECT
  81. to_regclass('voc.workspace')::text AS workspace,
  82. to_regclass('voc.product')::text AS product,
  83. to_regclass('voc.review')::text AS review,
  84. to_regclass('voc.sync_job')::text AS sync_job,
  85. to_regclass('voc.workspace_member')::text AS workspace_member,
  86. to_regclass('voc.analysis_run')::text AS analysis_run,
  87. to_regclass('voc.action_item')::text AS action_item,
  88. to_regclass('voc.alert')::text AS alert,
  89. to_regclass('voc.audit_log')::text AS audit_log
  90. ,to_regclass('voc.listing_source_snapshot')::text AS listing_source_snapshot
  91. ,to_regclass('voc.listing_score_job')::text AS listing_score_job
  92. ,to_regclass('voc.listing_score_item')::text AS listing_score_item
  93. ,to_regclass('voc.listing_current_score')::text AS listing_current_score
  94. ,to_regclass('voc.listing_version')::text AS listing_version
  95. `);
  96. const readiness = result.rows[0] ?? {};
  97. const missingTables = Object.entries(readiness)
  98. .filter((entry) => !entry[1])
  99. .map((entry) => `voc.${entry[0]}`);
  100. if (missingTables.length) {
  101. response.status(503).json({
  102. service: 'saas-voc-server',
  103. status: 'unavailable',
  104. database: 'migration_required',
  105. missingTables,
  106. timestamp: new Date().toISOString(),
  107. });
  108. return;
  109. }
  110. response.json({
  111. service: 'saas-voc-server',
  112. status: 'ok',
  113. database: 'ready',
  114. timestamp: new Date().toISOString(),
  115. });
  116. } catch {
  117. response.status(503).json({
  118. service: 'saas-voc-server',
  119. status: 'unavailable',
  120. database: 'unavailable',
  121. timestamp: new Date().toISOString(),
  122. });
  123. }
  124. });
  125. if (input.parseApp) app.use('/parse', input.parseApp);
  126. if (!input.platformRepository && !input.pool) throw new Error('Platform repository is not configured');
  127. const platform = input.platformRepository ?? new PostgresPlatformRepository(input.pool!);
  128. const access = new WorkspaceAccessService(platform);
  129. const businessRouter = Router();
  130. const authenticationMiddleware = createAuthenticationMiddleware(createAuthenticator(input.config));
  131. app.use('/api', authenticationMiddleware);
  132. const aiClient = new FmodeAiClient(input.config.ai);
  133. const domesticGateway = new FmodeVocEcommerceClient(input.config.fmode);
  134. businessRouter.use('/ai', createAiGatewayRouter(aiClient, input.aiPromptConfigs));
  135. const jobs = input.jobs ?? new SyncJobRepository(input.pool!);
  136. const sync = new SyncService(jobs);
  137. const snapshot = input.snapshot ?? new SnapshotService(input.pool!);
  138. businessRouter.use('/domestic-voc', createDomesticVocRouter({
  139. jobs,
  140. sync,
  141. snapshot,
  142. catalog: platform,
  143. access,
  144. defaultWorkspaceId: input.config.auth.defaultWorkspaceId,
  145. }));
  146. businessRouter.use('/saas', createSaasPlatformRouter({ repository: platform, access }));
  147. const listingAi = input.listingAiService ?? new ListingAiService(
  148. input.listingAiRepository ?? new InMemoryListingAiRepository(),
  149. new FmodeListingAiScoringProvider(aiClient, input.config.listingAi.model),
  150. () => new Date(),
  151. input.config.listingAi.concurrency,
  152. input.config.listingAi.maxAiItemsPerJob,
  153. input.config.listingAi.jdVocAiEnabled ? new FmodeJdVocAiScoringProvider(aiClient, input.config.listingAi.jdVocAiModel) : undefined,
  154. input.config.listingAi.jdVocImageShadowEnabled ? new FmodeGeminiImageReviewProvider({ baseUrl: process.env.FMODE_LLM_BASE_URL ?? input.config.ai.baseUrl, token: process.env.FMODE_LLM_API_KEY ?? input.config.ai.token, timeoutMs: 45_000 }) : undefined,
  155. input.config.listingAi.jdVocDisplayDefault,
  156. input.config.listingAi.jdVocEnabled,
  157. );
  158. businessRouter.use('/listing-ai', createListingAiRouter({
  159. service: listingAi,
  160. access,
  161. audit: platform,
  162. defaultWorkspaceId: input.config.auth.defaultWorkspaceId,
  163. }));
  164. queueMicrotask(() => {
  165. void listingAi.resumePendingJobs(input.config.auth.defaultWorkspaceId).catch((error: unknown) => {
  166. console.error('[listing-ai] failed to resume persisted score jobs', error);
  167. });
  168. });
  169. if (input.productKnowledge) {
  170. businessRouter.use('/knowledge', createProductKnowledgeRouter({
  171. store: input.productKnowledge,
  172. repository: platform,
  173. access,
  174. defaultWorkspaceId: input.config.auth.defaultWorkspaceId,
  175. }));
  176. }
  177. if (input.competitorListingMonitor) {
  178. businessRouter.use('/competitor-listings', createCompetitorListingMonitorRouter({
  179. service: input.competitorListingMonitor,
  180. access,
  181. defaultWorkspaceId: input.config.auth.defaultWorkspaceId,
  182. }));
  183. }
  184. app.use('/api', businessRouter);
  185. const localCloudFunctionRouter = createCloudFunctionRouter({
  186. defaultWorkspaceId: input.config.auth.defaultWorkspaceId,
  187. specialAction: createSpecialActionHandler({ ai: aiClient, domestic: domesticGateway }),
  188. dispatch: (request, response, actionPath, method, body, requestId) =>
  189. dispatchThroughRouter(businessRouter, request, response, actionPath, method, body, requestId),
  190. });
  191. app.use('/api/functions', localCloudFunctionRouter);
  192. // Parse REST storage has no embedded Parse Server. This compatibility mount
  193. // lets the browser's Parse SDK keep calling /parse/functions/* while the
  194. // request is still dispatched through the same allow-listed Cloud actions.
  195. if (!input.parseApp) app.use('/parse/functions', authenticationMiddleware, localCloudFunctionRouter);
  196. app.use((_request, response) => {
  197. response.status(404).json({ error: 'not_found' });
  198. });
  199. const errorHandler: ErrorRequestHandler = (error, _request, response, _next) => {
  200. if (error instanceof ApiError) {
  201. response.status(error.status).json({ error: error.code });
  202. return;
  203. }
  204. if (error instanceof ZodError) {
  205. response.status(400).json({
  206. error: 'invalid_request',
  207. issues: error.issues.map((issue) => ({ path: issue.path.join('.'), message: issue.message })),
  208. });
  209. return;
  210. }
  211. console.error('[http] unhandled request error', error);
  212. response.status(500).json({ error: 'internal_error' });
  213. };
  214. app.use(errorHandler);
  215. return app;
  216. }