diff --git a/package.json b/package.json index 6cbc319..3cc6d1b 100644 --- a/package.json +++ b/package.json @@ -39,6 +39,7 @@ "class-transformer": "^0.5.1", "class-validator": "^0.14.1", "compression": "^1.7.5", + "cookie-parser": "^1.4.7", "cors": "^2.8.5", "csv": "^6.3.11", "dayjs": "^1.11.13", @@ -59,6 +60,7 @@ "multer-s3": "^3.0.1", "node-cache": "^5.1.2", "nodemailer": "^6.9.16", + "openai": "^4.0.0", "passport": "^0.7.0", "passport-jwt": "^4.0.1", "reflect-metadata": "^0.2.2", @@ -66,6 +68,7 @@ "socket.io": "^4.8.1", "strip-ansi": "6.0.1", "swagger-ui-express": "^5.0.1", + "ulid": "^3.0.1", "winston": "^3.17.0" }, "devDependencies": { @@ -73,6 +76,7 @@ "@commitlint/config-conventional": "^19.5.0", "@types/bcrypt": "^5.0.2", "@types/compression": "^1.7.5", + "@types/cookie-parser": "^1.4.10", "@types/cors": "^2.8.17", "@types/express": "^4.17.21", "@types/jsonwebtoken": "^9.0.7", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 47dd938..6382377 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -38,6 +38,9 @@ importers: compression: specifier: ^1.7.5 version: 1.7.5 + cookie-parser: + specifier: ^1.4.7 + version: 1.4.7 cors: specifier: ^2.8.5 version: 2.8.5 @@ -98,6 +101,9 @@ importers: nodemailer: specifier: ^6.9.16 version: 6.9.16 + openai: + specifier: ^4.0.0 + version: 4.104.0(zod@3.25.76) passport: specifier: ^0.7.0 version: 0.7.0 @@ -119,6 +125,9 @@ importers: swagger-ui-express: specifier: ^5.0.1 version: 5.0.1(express@4.21.1) + ulid: + specifier: ^3.0.1 + version: 3.0.1 winston: specifier: ^3.17.0 version: 3.17.0 @@ -135,6 +144,9 @@ importers: '@types/compression': specifier: ^1.7.5 version: 1.7.5 + '@types/cookie-parser': + specifier: ^1.4.10 + version: 1.4.10(@types/express@4.17.21) '@types/cors': specifier: ^2.8.17 version: 2.8.17 @@ -1014,6 +1026,11 @@ packages: '@types/conventional-commits-parser@5.0.0': resolution: {integrity: sha512-loB369iXNmAZglwWATL+WRe+CRMmmBPtpolYzIebFaX4YA3x+BEfLqhUAV9WanycKI3TG1IMr5bMJDajDKLlUQ==} + '@types/cookie-parser@1.4.10': + resolution: {integrity: sha512-B4xqkqfZ8Wek+rCOeRxsjMS9OgvzebEzzLYw7NHYuvzb7IdxOkI0ZHGgeEBX4PUM7QGVvNSK60T3OvWj3YfBRg==} + peerDependencies: + '@types/express': '*' + '@types/cookie@0.4.1': resolution: {integrity: sha512-XW/Aa8APYr6jSVVA1y/DEIZX0/GMKLEVekNG727R8cs56ahETkRAy/3DR7+fJyh7oUgGwNQaRfXCun0+KbWY7Q==} @@ -1050,6 +1067,12 @@ packages: '@types/multer@1.4.12': resolution: {integrity: sha512-pQ2hoqvXiJt2FP9WQVLPRO+AmiIm/ZYkavPlIQnx282u4ZrVdztx0pkh3jjpQt0Kz+YI0YhSG264y08UJKoUQg==} + '@types/node-fetch@2.6.13': + resolution: {integrity: sha512-QGpRVpzSaUs30JBSGPjOg4Uveu384erbHBoT1zeONvyCfwQxIkUshLAOqN/k9EjGviPRmWTTe6aH2qySWKTVSw==} + + '@types/node@18.19.130': + resolution: {integrity: sha512-GRaXQx6jGfL8sKfaIDD6OupbIHBr9jv7Jnaml9tB7l4v068PAOXqfcujMMo5PhbIs6ggR1XODELqahT2R8v0fg==} + '@types/node@20.17.6': resolution: {integrity: sha512-VEI7OdvK2wP7XHnsuXbAJnEpEkF6NjSN45QJlL4VGqZSXsnicpesdTWsg9RISeSdYd3yeRj/y3k5KGjUXYnFwQ==} @@ -1162,6 +1185,10 @@ packages: abbrev@1.1.1: resolution: {integrity: sha512-nne9/IiQ/hzIhY6pdDnbBtz7DjPTKrY00P/zvPSm5pOFkl6xuGrGnXn/VtTNNfNtAfZ9/1RtehkszU9qcTii0Q==} + abort-controller@3.0.0: + resolution: {integrity: sha512-h8lQ8tacZYnR3vNQTgibj+tODHI5/+l06Au2Pcriv/Gmet0eaj4TwWH41sO9wnHDiQsEj19q0drzdWdeAHtweg==} + engines: {node: '>=6.5'} + accepts@1.3.8: resolution: {integrity: sha512-PYAthTa2m2VKxuvSD3DPC/Gy+U+sOA1LAuT8mkmRuvw+NACSaeXEQ+NHcVF7rONl6qcaxV3Uuemwawk+7+SJLw==} engines: {node: '>= 0.6'} @@ -1184,6 +1211,10 @@ packages: resolution: {integrity: sha512-RZNwNclF7+MS/8bDg70amg32dyeZGZxiDuQmZxKLAlQjr3jGyLx+4Kkk58UO7D2QdgFIQCovuSuZESne6RG6XQ==} engines: {node: '>= 6.0.0'} + agentkeepalive@4.6.0: + resolution: {integrity: sha512-kja8j7PjmncONqaTsB8fQ+wE2mSU2DJ9D4XKoJ5PFWIdRMa6SLSN1ff4mOr4jCbfRSsxR4keIiySJU0N9T5hIQ==} + engines: {node: '>= 8.0.0'} + ajv@6.12.6: resolution: {integrity: sha512-j3fVLgvTo527anyYyJOGTYJbG+vnnQYvE0m5mmkc1TK+nxAppkCLMIL0aZ4dblVCNoGShhm+kzE4ZUykBoMg4g==} @@ -1355,6 +1386,10 @@ packages: resolution: {integrity: sha512-/Nf7TyzTx6S3yRJObOAV7956r8cr2+Oj8AC5dt8wSP3BQAoeX58NoHyCU8P8zGkNXStjTSi6fzO6F0pBdcYbEg==} engines: {node: '>= 0.8'} + call-bind-apply-helpers@1.0.2: + resolution: {integrity: sha512-Sp1ablJ0ivDkSzjcaJdxEunN5/XvksFJ2sMBFfq6x0ryhQV/2b/KwFe21cMpmHtPOSij8K99/wSfoEuTObmuMQ==} + engines: {node: '>= 0.4'} + call-bind@1.0.7: resolution: {integrity: sha512-GHTSNSYICQ7scH7sZ+M2rFopRoLh8t2bLSW6BbgrtLsahOIB5iyAVJf9GjWK3cYTDaMj4XdBpM1cA6pIS0Kv2w==} engines: {node: '>= 0.4'} @@ -1507,6 +1542,10 @@ packages: engines: {node: '>=16'} hasBin: true + cookie-parser@1.4.7: + resolution: {integrity: sha512-nGUvgXnotP3BsjiLX2ypbQnWoGUPIIfHQNZkkC668ntrzGWEZVW70HDEB1qnNGMicPje6EttlIgzo51YSwNQGw==} + engines: {node: '>= 0.8.0'} + cookie-signature@1.0.6: resolution: {integrity: sha512-QADzlaHc8icV8I7vbaJXJwod9HWYp8uCqf1xa4OfNu1T7JVxQIrUgOWtHdNDtPiywmFbiS12VjotIXLrKM3orQ==} @@ -1667,6 +1706,10 @@ packages: resolution: {integrity: sha512-ZmdL2rui+eB2YwhsWzjInR8LldtZHGDoQ1ugH85ppHKwpUHL7j7rN0Ti9NCnGiQbhaZ11FpR+7ao1dNsmduNUg==} engines: {node: '>=12'} + dunder-proto@1.0.1: + resolution: {integrity: sha512-KIN/nDJBQRcXw0MLVhZE9iQHmG68qAVIBg9CqmUYjmQIhgij9U5MFvrqkUL5FbtyyzZuOeOt0zdeRe4UY7ct+A==} + engines: {node: '>= 0.4'} + eastasianwidth@0.2.0: resolution: {integrity: sha512-I88TYZWc9XiYHRQ4/3c5rjjfgkjhLyW2luGIheGERbNQ6OY7yTybanSpDXZa8y7VUP9YmDcYa+eyq4ca7iLqWA==} @@ -1723,6 +1766,10 @@ packages: resolution: {integrity: sha512-jxayLKShrEqqzJ0eumQbVhTYQM27CfT1T35+gCgDFoL82JLsXqTJ76zv6A0YLOgEnLUMvLzsDsGIrl8NFpT2gQ==} engines: {node: '>= 0.4'} + es-define-property@1.0.1: + resolution: {integrity: sha512-e3nRfgfUZ4rNGL232gUgX06QNyyez04KdjFrF+LTRoOXmrOgFKDg4BCdsjW8EnT69eqdYGmRpJwiPVYNrCaW3g==} + engines: {node: '>= 0.4'} + es-errors@1.3.0: resolution: {integrity: sha512-Zf5H2Kxt2xjTvbJvP2ZWLEICxA6j+hAmMzIlypy4xcBg1vKVnx89Wy0GbS+kf5cwCVFFzdCFh2XSCFNULS6csw==} engines: {node: '>= 0.4'} @@ -1731,10 +1778,18 @@ packages: resolution: {integrity: sha512-MZ4iQ6JwHOBQjahnjwaC1ZtIBH+2ohjamzAO3oaHcXYup7qxjF2fixyH+Q71voWHeOkI2q/TnJao/KfXYIZWbw==} engines: {node: '>= 0.4'} + es-object-atoms@1.1.1: + resolution: {integrity: sha512-FGgH2h8zKNim9ljj7dankFPcICIK9Cp5bm+c2gQSYePhpaG5+esrLODihIorn+Pe6FGJzWhXQotPv73jTaldXA==} + engines: {node: '>= 0.4'} + es-set-tostringtag@2.0.3: resolution: {integrity: sha512-3T8uNMC3OQTHkFUsFq8r/BwAXLHvU/9O9mE0fBc/MY5iq/8H7ncvO947LmYA6ldWw9Uh8Yhf25zu6n7nML5QWQ==} engines: {node: '>= 0.4'} + es-set-tostringtag@2.1.0: + resolution: {integrity: sha512-j6vWzfrGVfyXxge+O0x5sh6cvxAog0a/4Rdd2K36zCMV5eJ+/+tOAngRO8cODMNWbVRdVlmGZQL2YS3yR8bIUA==} + engines: {node: '>= 0.4'} + es-shim-unscopables@1.0.2: resolution: {integrity: sha512-J3yBRXCzDu4ULnQwxyToo/OjdMx6akgVC7K6few0a7F/0wLtmKKN7I73AH5T2836UuXRqN7Qg+IIUw/+YJksRw==} @@ -1854,6 +1909,10 @@ packages: resolution: {integrity: sha512-aIL5Fx7mawVa300al2BnEE4iNvo1qETxLrPI/o05L7z6go7fCw1J6EQmbK4FmJ2AS7kgVF/KEZWufBfdClMcPg==} engines: {node: '>= 0.6'} + event-target-shim@5.0.1: + resolution: {integrity: sha512-i/2XbnSz/uxRCU6+NdVJgKWDTM427+MqYbkQzD321DuCQJUqOuJKIA0IM2+W2xtYHdKOmZ4dR6fExsd4SXL+WQ==} + engines: {node: '>=6'} + eventemitter3@5.0.1: resolution: {integrity: sha512-GWkBvjiSZK87ELrYOSESUYeVIc9mvLLf/nXalMOS5dYrgZq9o5OVkbZAVM06CVxYsCwH9BDZFPlQTlPA1j4ahA==} @@ -1965,10 +2024,21 @@ packages: resolution: {integrity: sha512-Ld2g8rrAyMYFXBhEqMz8ZAHBi4J4uS1i/CxGMDnjyFWddMXLVcDp051DZfu+t7+ab7Wv6SMqpWmyFIj5UbfFvg==} engines: {node: '>=14'} + form-data-encoder@1.7.2: + resolution: {integrity: sha512-qfqtYan3rxrnCk1VYaA4H+Ms9xdpPqvLZa6xmMgFvhO32x7/3J/ExcTd6qpxM0vH2GdMI+poehyBZvqfMTto8A==} + form-data@4.0.1: resolution: {integrity: sha512-tzN8e4TX8+kkxGPK8D5u0FNmjPUjw3lwC9lSLxxoB/+GtsJG91CO8bSWy73APlgAZzZbXEYZJuxjkHH2w+Ezhw==} engines: {node: '>= 6'} + form-data@4.0.5: + resolution: {integrity: sha512-8RipRLol37bNs2bhoV67fiTEvdTrbMUYcFTiy3+wuuOnUog2QBHCZWXDRijWQfAkhBj2Uf5UnVaiWwA5vdd82w==} + engines: {node: '>= 6'} + + formdata-node@4.4.1: + resolution: {integrity: sha512-0iirZp3uVDjVGt9p49aTaqjk84TrglENEDuqfdlZQ1roC9CWlPk6Avf8EEnZNcAqPonwkG35x4n3ww/1THYAeQ==} + engines: {node: '>= 12.20'} + forwarded@0.2.0: resolution: {integrity: sha512-buRG0fpBtRHSTCOASe6hD258tEubFoRLb4ZNA6NxMVHNw2gOcwHo9wyablzMzOA5z9xA9L1KNjk/Nt6MT9aYow==} engines: {node: '>= 0.6'} @@ -2016,6 +2086,14 @@ packages: resolution: {integrity: sha512-5uYhsJH8VJBTv7oslg4BznJYhDoRI6waYCxMmCdnTrcCrHA/fCFKoTFz2JKKE0HdDFUF7/oQuhzumXJK7paBRQ==} engines: {node: '>= 0.4'} + get-intrinsic@1.3.0: + resolution: {integrity: sha512-9fSjSaos/fRIVIp+xSJlE6lfwhES7LNtKaCBIamHsjr2na1BiABJPo0mOjjz8GJDURarmCPGqaiVg5mfjb98CQ==} + engines: {node: '>= 0.4'} + + get-proto@1.0.1: + resolution: {integrity: sha512-sTSfBjoXBp89JvIKIefqw7U2CCebsc74kiY6awiGogKtoSGbgjYE/G/+l9sF3MWFPNc9IcoOC4ODfKHfxFmp0g==} + engines: {node: '>= 0.4'} + get-stream@8.0.1: resolution: {integrity: sha512-VaUJspBffn/LMCJVoMvSAdmscJyS1auj5Zulnn5UoYcY531UWmdwhRWkcGKnGU93m5HSXP9LP2usOryrBtQowA==} engines: {node: '>=16'} @@ -2064,6 +2142,10 @@ packages: gopd@1.0.1: resolution: {integrity: sha512-d65bNlIadxvpb/A2abVdlqKqV563juRnZ1Wtk6s1sIR8uNsXR70xqIzVqxVf1eTqDunwT2MkczEeaezCKTZhwA==} + gopd@1.2.0: + resolution: {integrity: sha512-ZUKRh6/kUFoAiTAtTYPZJ3hw9wNxx+BIBOijnlG9PnrJsCcSjs1wyyD6vJpaYtgnzDrKYRSqf3OO6Rfa93xsRg==} + engines: {node: '>= 0.4'} + graphemer@1.4.0: resolution: {integrity: sha512-EtKwoO6kxCL9WO5xipiHTZlSzBm7WLT627TqC/uVRd0HKmq8NXyebnNYxDoBi7wt8eTWrUrKXCOVaFq9x1kgag==} @@ -2094,6 +2176,10 @@ packages: resolution: {integrity: sha512-l3LCuF6MgDNwTDKkdYGEihYjt5pRPbEg46rtlmnSPlUbgmB8LOIrKJbYYFBSbnPaJexMKtiPO8hmeRjRz2Td+A==} engines: {node: '>= 0.4'} + has-symbols@1.1.0: + resolution: {integrity: sha512-1cDNdwJ2Jaohmb3sg4OmKaMBwuC48sYni5HUw2DvsC8LjGTLK9h+eb1X6RyuOHe4hT0ULCW68iomhjUoKUqlPQ==} + engines: {node: '>= 0.4'} + has-tostringtag@1.0.2: resolution: {integrity: sha512-NqADB8VjPFLM2V0VvHUewwwsw0ZWBaIdgo+ieHtK3hasLz4qeCRjYcqfB6AQrBggRKppKF8L52/VqdVsO47Dlw==} engines: {node: '>= 0.4'} @@ -2123,6 +2209,9 @@ packages: resolution: {integrity: sha512-AXcZb6vzzrFAUE61HnN4mpLqd/cSIwNQjtNWR0euPm6y0iqx3G4gOXaIDdtdDwZmhwe82LA6+zinmW4UBWVePQ==} engines: {node: '>=16.17.0'} + humanize-ms@1.2.1: + resolution: {integrity: sha512-Fl70vYtsAFb/C06PTS9dZBo7ihau+Tu/DNCk/OyHhea07S+aeMWpFFkUaXRa8fI+ScZbEI8dfSxwY7gxZ9SAVQ==} + husky@9.1.6: resolution: {integrity: sha512-sqbjZKK7kf44hfdE94EoX8MZNk0n7HeW37O4YrVGCF4wzgQjp+akPAkfUK5LZ6KuR/6sqeAVuXHji+RzQgOn5A==} engines: {node: '>=18'} @@ -2504,6 +2593,10 @@ packages: make-error@1.3.6: resolution: {integrity: sha512-s8UhlNe7vPKomQhC1qFelMokr/Sc3AgNbso3n74mVPA5LTZwkB9NlXf4XPamLxJE8h0gh73rM94xvwRT2CVInw==} + math-intrinsics@1.1.0: + resolution: {integrity: sha512-/IXtbwEk5HTPyEwyKX6hGkYXxM9nbj64B+ilVJnC/R6B0pH5G4V3b0pVbL7DBj4tkhBAppbQUlf6F6Xl9LHu1g==} + engines: {node: '>= 0.4'} + media-typer@0.3.0: resolution: {integrity: sha512-dq+qelQ9akHpcOl/gUVRTxVIOkAJ1wR3QAvb4RsVjS8oVoFjDGTc679wJYmUmknUF5HwMLOgb5O+a3KxfWapPQ==} engines: {node: '>= 0.6'} @@ -2720,6 +2813,11 @@ packages: resolution: {integrity: sha512-t1QzWwnk4sjLWaQAS8CHgOJ+RAfmHpxFWmc36IWTiWHQfs0w5JDMBS1b1ZxQteo0vVVuWJvIUKHDkkeK7vIGCg==} engines: {node: '>= 8.0.0'} + node-domexception@1.0.0: + resolution: {integrity: sha512-/jKZoMpw0F8GRwl4/eLROPA3cfcXtLApP0QzLmUT/HuPCZWyB7IY9ZrMeKw2O/nFIqPQB3PVM9aYm0F312AXDQ==} + engines: {node: '>=10.5.0'} + deprecated: Use your platform's native DOMException instead + node-fetch@2.7.0: resolution: {integrity: sha512-c4FRfUm/dbcWZ7U+1Wq0AwCyFL+3nt2bEw05wfxSz+DWpWsitgmSgYmy2dQdWyKC1694ELPqMs/YzUSNozLt8A==} engines: {node: 4.x || >=6.0.0} @@ -2817,6 +2915,18 @@ packages: resolution: {integrity: sha512-VXJjc87FScF88uafS3JllDgvAm+c/Slfz06lorj2uAY34rlUu0Nt+v8wreiImcrgAjjIHp1rXpTDlLOGw29WwQ==} engines: {node: '>=18'} + openai@4.104.0: + resolution: {integrity: sha512-p99EFNsA/yX6UhVO93f5kJsDRLAg+CTA2RBqdHK4RtK8u5IJw32Hyb2dTGKbnnFmnuoBv5r7Z2CURI9sGZpSuA==} + hasBin: true + peerDependencies: + ws: ^8.18.0 + zod: ^3.23.8 + peerDependenciesMeta: + ws: + optional: true + zod: + optional: true + optionator@0.9.4: resolution: {integrity: sha512-6IpQ7mKUxRcZNLIObR0hz7lxsapSSIYNZJwXPGeF0mTVqGKFIXj1DQcMoT22S3ROcLyY/rz0PWaWZ9ayWmad9g==} engines: {node: '>= 0.8.0'} @@ -3413,12 +3523,19 @@ packages: resolution: {integrity: sha512-u3xV3X7uzvi5b1MncmZo3i2Aw222Zk1keqLA1YkHldREkAhAqi65wuPfe7lHx8H/Wzy+8CE7S7uS3jekIM5s8g==} engines: {node: '>=8'} + ulid@3.0.1: + resolution: {integrity: sha512-dPJyqPzx8preQhqq24bBG1YNkvigm87K8kVEHCD+ruZg24t6IFEFv00xMWfxcC4djmFtiTLdFuADn4+DOz6R7Q==} + hasBin: true + unbox-primitive@1.0.2: resolution: {integrity: sha512-61pPlCD9h51VoreyJ0BReideM3MDKMKnh6+V9L08331ipq6Q8OFXZYiqP6n/tbHx4s5I9uRhcye6BrbkizkBDw==} undefsafe@2.0.5: resolution: {integrity: sha512-WxONCrssBM8TSPRqN5EmsjVrsv4A8X12J4ArBiiayv3DyyG3ZlIg6yysuuSYdZsVz3TKcTg2fd//Ujd4CHV1iA==} + undici-types@5.26.5: + resolution: {integrity: sha512-JlCMO+ehdEIKqlFxk6IfVoAUVmgz7cU7zD/h9XZ0qzeosSHmUJVOzSQvvYSYWXkFXC+IfLKSIffhv0sVZup6pA==} + undici-types@6.19.8: resolution: {integrity: sha512-ve2KP6f/JnbPBFyobGHuerC9g1FYGn/F8n1LWTwNxCEzd6IfqTwUQcNXgEtmmQ6DlRrC1hrSrBnCZPokRrDHjw==} @@ -3458,6 +3575,10 @@ packages: wcwidth@1.0.1: resolution: {integrity: sha512-XHPEwS0q6TaxcvG85+8EYkbiCux2XtWG2mkc47Ng2A77BQu9+DqIOJldST4HgPkuea7dvKSj5VgX3P1d4rW8Tg==} + web-streams-polyfill@4.0.0-beta.3: + resolution: {integrity: sha512-QW95TCTaHmsYfHDybGMwO5IJIM93I/6vTRk+daHTWFPhwh+C8Cg7j7XyKrwrj8Ib6vYXe0ocYNrmzY4xAAN6ug==} + engines: {node: '>= 14'} + webidl-conversions@3.0.1: resolution: {integrity: sha512-2JAn3z8AR6rjK8Sm8orRC0h/bcl/DqL7tRPdGZ4I1CjdF+EaMLmYxBHyXuKL849eucPFhvBoxMsflfOb8kxaeQ==} @@ -3573,6 +3694,9 @@ packages: resolution: {integrity: sha512-b4JR1PFR10y1mKjhHY9LaGo6tmrgjit7hxVIeAmyMw3jegXR4dhYqLaQF5zMXZxY7tLpMyJeLjr1C4rLmkVe8g==} engines: {node: '>=12.20'} + zod@3.25.76: + resolution: {integrity: sha512-gzUt/qt81nXsFGKIFcC3YnfEAx5NkunCfnDlvuBSSFS02bcXu4Lmea0AFIUwbLWxWPx3d9p8S5QoaujKcNQxcQ==} + snapshots: '@aws-crypto/crc32@5.2.0': @@ -4885,6 +5009,10 @@ snapshots: dependencies: '@types/node': 20.17.6 + '@types/cookie-parser@1.4.10(@types/express@4.17.21)': + dependencies: + '@types/express': 4.17.21 + '@types/cookie@0.4.1': {} '@types/cors@2.8.17': @@ -4940,6 +5068,15 @@ snapshots: dependencies: '@types/express': 4.17.21 + '@types/node-fetch@2.6.13': + dependencies: + '@types/node': 20.17.6 + form-data: 4.0.5 + + '@types/node@18.19.130': + dependencies: + undici-types: 5.26.5 + '@types/node@20.17.6': dependencies: undici-types: 6.19.8 @@ -5087,6 +5224,10 @@ snapshots: abbrev@1.1.1: {} + abort-controller@3.0.0: + dependencies: + event-target-shim: 5.0.1 + accepts@1.3.8: dependencies: mime-types: 2.1.35 @@ -5108,6 +5249,10 @@ snapshots: transitivePeerDependencies: - supports-color + agentkeepalive@4.6.0: + dependencies: + humanize-ms: 1.2.1 + ajv@6.12.6: dependencies: fast-deep-equal: 3.1.3 @@ -5324,6 +5469,11 @@ snapshots: bytes@3.1.2: {} + call-bind-apply-helpers@1.0.2: + dependencies: + es-errors: 1.3.0 + function-bind: 1.1.2 + call-bind@1.0.7: dependencies: es-define-property: 1.0.0 @@ -5486,6 +5636,11 @@ snapshots: meow: 12.1.1 split2: 4.2.0 + cookie-parser@1.4.7: + dependencies: + cookie: 0.7.2 + cookie-signature: 1.0.6 + cookie-signature@1.0.6: {} cookie@0.7.1: {} @@ -5622,6 +5777,12 @@ snapshots: dotenv@16.4.5: {} + dunder-proto@1.0.1: + dependencies: + call-bind-apply-helpers: 1.0.2 + es-errors: 1.3.0 + gopd: 1.2.0 + eastasianwidth@0.2.0: {} ecdsa-sig-formatter@1.0.11: @@ -5722,18 +5883,31 @@ snapshots: dependencies: get-intrinsic: 1.2.4 + es-define-property@1.0.1: {} + es-errors@1.3.0: {} es-object-atoms@1.0.0: dependencies: es-errors: 1.3.0 + es-object-atoms@1.1.1: + dependencies: + es-errors: 1.3.0 + es-set-tostringtag@2.0.3: dependencies: get-intrinsic: 1.2.4 has-tostringtag: 1.0.2 hasown: 2.0.2 + es-set-tostringtag@2.1.0: + dependencies: + es-errors: 1.3.0 + get-intrinsic: 1.3.0 + has-tostringtag: 1.0.2 + hasown: 2.0.2 + es-shim-unscopables@1.0.2: dependencies: hasown: 2.0.2 @@ -5909,6 +6083,8 @@ snapshots: etag@1.8.1: {} + event-target-shim@5.0.1: {} + eventemitter3@5.0.1: {} events@3.3.0: {} @@ -6061,12 +6237,27 @@ snapshots: cross-spawn: 7.0.5 signal-exit: 4.1.0 + form-data-encoder@1.7.2: {} + form-data@4.0.1: dependencies: asynckit: 0.4.0 combined-stream: 1.0.8 mime-types: 2.1.35 + form-data@4.0.5: + dependencies: + asynckit: 0.4.0 + combined-stream: 1.0.8 + es-set-tostringtag: 2.1.0 + hasown: 2.0.2 + mime-types: 2.1.35 + + formdata-node@4.4.1: + dependencies: + node-domexception: 1.0.0 + web-streams-polyfill: 4.0.0-beta.3 + forwarded@0.2.0: {} fresh@0.5.2: {} @@ -6115,6 +6306,24 @@ snapshots: has-symbols: 1.0.3 hasown: 2.0.2 + get-intrinsic@1.3.0: + dependencies: + call-bind-apply-helpers: 1.0.2 + es-define-property: 1.0.1 + es-errors: 1.3.0 + es-object-atoms: 1.1.1 + function-bind: 1.1.2 + get-proto: 1.0.1 + gopd: 1.2.0 + has-symbols: 1.1.0 + hasown: 2.0.2 + math-intrinsics: 1.1.0 + + get-proto@1.0.1: + dependencies: + dunder-proto: 1.0.1 + es-object-atoms: 1.1.1 + get-stream@8.0.1: {} get-symbol-description@1.0.2: @@ -6176,6 +6385,8 @@ snapshots: dependencies: get-intrinsic: 1.2.4 + gopd@1.2.0: {} + graphemer@1.4.0: {} handlebars@4.7.8: @@ -6201,6 +6412,8 @@ snapshots: has-symbols@1.0.3: {} + has-symbols@1.1.0: {} + has-tostringtag@1.0.2: dependencies: has-symbols: 1.0.3 @@ -6232,6 +6445,10 @@ snapshots: human-signals@5.0.0: {} + humanize-ms@1.2.1: + dependencies: + ms: 2.1.3 + husky@9.1.6: {} iconv-lite@0.4.24: @@ -6611,6 +6828,8 @@ snapshots: make-error@1.3.6: {} + math-intrinsics@1.1.0: {} + media-typer@0.3.0: {} memory-pager@1.5.0: {} @@ -6831,6 +7050,8 @@ snapshots: dependencies: clone: 2.1.2 + node-domexception@1.0.0: {} + node-fetch@2.7.0: dependencies: whatwg-url: 5.0.0 @@ -6934,6 +7155,20 @@ snapshots: dependencies: mimic-function: 5.0.1 + openai@4.104.0(zod@3.25.76): + dependencies: + '@types/node': 18.19.130 + '@types/node-fetch': 2.6.13 + abort-controller: 3.0.0 + agentkeepalive: 4.6.0 + form-data-encoder: 1.7.2 + formdata-node: 4.4.1 + node-fetch: 2.7.0 + optionalDependencies: + zod: 3.25.76 + transitivePeerDependencies: + - encoding + optionator@0.9.4: dependencies: deep-is: 0.1.4 @@ -7551,6 +7786,8 @@ snapshots: dependencies: '@lukeed/csprng': 1.1.0 + ulid@3.0.1: {} + unbox-primitive@1.0.2: dependencies: call-bind: 1.0.7 @@ -7560,6 +7797,8 @@ snapshots: undefsafe@2.0.5: {} + undici-types@5.26.5: {} + undici-types@6.19.8: {} unicorn-magic@0.1.0: {} @@ -7586,6 +7825,8 @@ snapshots: dependencies: defaults: 1.0.4 + web-streams-polyfill@4.0.0-beta.3: {} + webidl-conversions@3.0.1: {} webidl-conversions@7.0.0: {} @@ -7706,3 +7947,6 @@ snapshots: yocto-queue@0.1.0: {} yocto-queue@1.1.1: {} + + zod@3.25.76: + optional: true diff --git a/src/IOC/ioc.config.ts b/src/IOC/ioc.config.ts index 0fb9b01..804b469 100644 --- a/src/IOC/ioc.config.ts +++ b/src/IOC/ioc.config.ts @@ -54,6 +54,11 @@ import { ChatService } from "../modules/chat/chat.service"; import { ChatRepo, createChatRepo } from "../modules/chat/repository/chat.repository"; import { ChatMessageRepo, createChatMessageRepo } from "../modules/chat/repository/message.repository"; import { WsAuthService } from "../modules/chat/wsAuth.service"; +import { ChatbotGateway } from "../modules/chatbot/chatbot.gateway"; +import { ChatbotService } from "../modules/chatbot/providers/chatbot.service"; +import { DataContextService } from "../modules/chatbot/providers/data-context.service"; +import { LLMService } from "../modules/chatbot/providers/llm.service"; +import { WebSocketAuthService } from "../modules/chatbot/providers/websocket-auth.service"; import { ContactUsRepo, CreateContactUsRepo } from "../modules/contact-us/contactUs.repository"; import { ContactUsService } from "../modules/contact-us/contactUs.service"; import { CouponRepo, CouponUsageRepo, createCouponRepo, createCouponUsageRepo } from "../modules/Coupon/coupon.repository"; @@ -197,6 +202,11 @@ const containerModules = new AsyncContainerModule(async (bind) => { bind(IOCTYPES.LearningService).to(LearningService).inSingletonScope(); bind(IOCTYPES.RedisService).to(RedisService).inSingletonScope(); bind(IOCTYPES.CacheService).to(CacheService).inSingletonScope(); + bind(IOCTYPES.ChatbotService).to(ChatbotService).inSingletonScope(); + bind(IOCTYPES.ChatbotLLMService).to(LLMService).inSingletonScope(); + bind(IOCTYPES.ChatbotDataContextService).to(DataContextService).inSingletonScope(); + bind(IOCTYPES.ChatbotWebSocketAuthService).to(WebSocketAuthService).inSingletonScope(); + bind(IOCTYPES.ChatbotGateway).to(ChatbotGateway).inSingletonScope(); // #endregion // #region repository bind(IOCTYPES.CategoryRepository).toDynamicValue(createCategoryRepo).inSingletonScope(); @@ -279,7 +289,6 @@ const containerModules = new AsyncContainerModule(async (bind) => { bind(IOCTYPES.AboutUsRepo).toDynamicValue(CreateAboutUsRepo).inSingletonScope(); bind(IOCTYPES.SiteSettingRepo).toDynamicValue(CreateSiteSettingRepo).inSingletonScope(); bind(IOCTYPES.NewsletterRepo).toDynamicValue(CreateNewsletterRepo).inSingletonScope(); - // #endregion }); export { containerModules }; diff --git a/src/IOC/ioc.types.ts b/src/IOC/ioc.types.ts index df912a0..c343176 100644 --- a/src/IOC/ioc.types.ts +++ b/src/IOC/ioc.types.ts @@ -46,6 +46,11 @@ export const IOCTYPES = { NewsletterService: Symbol.for("NewsletterService"), PricingService: Symbol.for("PricingService"), OrderQueue: Symbol.for("OrderQueue"), + ChatbotService: Symbol.for("ChatbotService"), + ChatbotLLMService: Symbol.for("ChatbotLLMService"), + ChatbotDataContextService: Symbol.for("ChatbotDataContextService"), + ChatbotGateway: Symbol.for("ChatbotGateway"), + ChatbotWebSocketAuthService: Symbol.for("ChatbotWebSocketAuthService"), // #endregion // #region repository CategoryRepository: Symbol.for("CategoryRepository"), @@ -129,6 +134,8 @@ export const IOCTYPES = { ContactUsRepo: Symbol.for("ContactUsRepo"), AboutUsRepo: Symbol.for("AboutUsRepo"), NewsletterRepo: Symbol.for("NewsletterRepo"), + ChatbotChatSessionRepository: Symbol.for("ChatbotChatSessionRepository"), + ChatbotChatMessageRepository: Symbol.for("ChatbotChatMessageRepository"), // #endregion Logger: Symbol.for("Logger"), ZarinPalGateway: Symbol.for("ZarinPalGateway"), diff --git a/src/app.ts b/src/app.ts index f05d625..dd46728 100644 --- a/src/app.ts +++ b/src/app.ts @@ -1,4 +1,7 @@ +// import path from "path"; + import compression from "compression"; +import cookieParser from "cookie-parser"; import cors from "cors"; import express, { Application } from "express"; import { Container } from "inversify"; @@ -31,6 +34,7 @@ import "./modules/blog/blog.controller"; import "./modules/landing/landing.controller"; import "./modules/ticket/ticket.controller"; import "./modules/chat/chat.controller"; +import "./modules/chatbot/chatbot.controller"; import "./modules/job/job.controller"; import "./modules/faq/faq.controller"; import "./modules/notification/notification.controller"; @@ -120,6 +124,11 @@ class App { */ app.use(compression()); + /** + * configure cookie parser + */ + app.use(cookieParser()); + /** * configure body parser */ @@ -135,6 +144,11 @@ class App { */ app.use(cors(appConfig.cors)); + /** + * serve static files from project root (for test files) + */ + app.use(express.static(process.cwd())); + /** * config passport */ diff --git a/src/common/enums/message.enum.ts b/src/common/enums/message.enum.ts index 090b233..bb4ecc6 100644 --- a/src/common/enums/message.enum.ts +++ b/src/common/enums/message.enum.ts @@ -329,3 +329,17 @@ export const enum RoleMessage { PermissionsNotEmpty = "دسترسی‌ها نباید خالی باشد", RoleExist = "این نقش قبلا ثبت شده است", } + +export const enum WebSocketMessage { + AUTHENTICATED = "احراز هویت موفقیت‌آمیز", + AUTHENTICATION_REQUIRED = "احراز هویت ضروری است", + INVALID_TOKEN = "توکن احراز هویت نامعتبر است", + TOKEN_EXPIRED = "توکن منقضی شده است", + USER_NOT_AUTHENTICATED = "کاربر احراز هویت نشده است", + CONNECTION_FAILED = "اتصال برقرار نشد", + SESSION_CREATED = "جلسه چت با موفقیت ایجاد شد", + SESSION_CREATION_FAILED = "ایجاد جلسه چت با شکست مواجه شد", + SESSION_NOT_FOUND = "جلسه چت یافت نشد", + CHAT_JOINED = "با موفقیت به چت متصل شدید", + LLM_SERVICE_ERROR = "سرویس هوش مصنوعی موقتاً در دسترس نیست", +} diff --git a/src/modules/chatbot/DTO/send-message.dto.ts b/src/modules/chatbot/DTO/send-message.dto.ts new file mode 100644 index 0000000..033a0c8 --- /dev/null +++ b/src/modules/chatbot/DTO/send-message.dto.ts @@ -0,0 +1,29 @@ +import { Expose } from "class-transformer"; +import { IsNotEmpty, IsOptional, IsString, MaxLength } from "class-validator"; + +import { ApiProperty } from "../../../common/decorator/swggerDocs"; + +export class SendMessageDto { + @Expose() + @IsString() + @MaxLength(2000) + @ApiProperty({ type: "string", description: "Message content", example: "How can I cancel my subscription?" }) + content!: string; + + @Expose() + @IsNotEmpty() + @IsString() + @ApiProperty({ type: "string", description: "Chat session ID (ULID)" }) + sessionId!: string; + + @Expose() + @IsOptional() + @IsString() + @ApiProperty({ type: "string", description: "ID of message being replied to", required: false }) + responseToId?: string; + + @Expose() + @IsOptional() + @ApiProperty({ type: "object", description: "Additional metadata for the message", required: false }) + metadata?: Record; +} diff --git a/src/modules/chatbot/DTO/session-id.param.dto.ts b/src/modules/chatbot/DTO/session-id.param.dto.ts new file mode 100644 index 0000000..55f29b3 --- /dev/null +++ b/src/modules/chatbot/DTO/session-id.param.dto.ts @@ -0,0 +1,12 @@ +import { Expose } from "class-transformer"; +import { IsNotEmpty, IsString } from "class-validator"; + +import { ApiProperty } from "../../../common/decorator/swggerDocs"; + +export class SessionIdParamDto { + @Expose() + @IsNotEmpty({ message: "Session ID is required" }) + @IsString() + @ApiProperty({ type: "string", description: "Session id (ULID)", example: "01ARZ3NDEKTSV4RRFFQ69G5FAV" }) + sessionId!: string; +} diff --git a/src/modules/chatbot/DTO/websocket-events.dto.ts b/src/modules/chatbot/DTO/websocket-events.dto.ts new file mode 100644 index 0000000..a7223b0 --- /dev/null +++ b/src/modules/chatbot/DTO/websocket-events.dto.ts @@ -0,0 +1,62 @@ +import { IsNotEmpty, IsOptional, IsString, IsUUID, MaxLength } from "class-validator"; + +export class AuthenticateDto { + @IsNotEmpty() + @IsString() + token!: string; +} + +export class CreateSessionDto { + @IsNotEmpty() + @IsString() + @MaxLength(100) + title!: string; +} + +export class JoinChatDto { + @IsNotEmpty() + @IsUUID() + sessionId!: string; + + @IsNotEmpty() + @IsString() + userId!: string; +} + +export class LeaveChatDto { + @IsNotEmpty() + @IsUUID() + sessionId!: string; +} + +export class SendMessageWebSocketDto { + @IsNotEmpty() + @IsUUID() + sessionId!: string; + + @IsNotEmpty() + @IsString() + @MaxLength(4000) + content!: string; + + @IsNotEmpty() + @IsString() + userId!: string; + + @IsOptional() + @IsUUID() + responseToId?: string; + + @IsOptional() + metadata?: Record; +} + +export class TypingDto { + @IsNotEmpty() + @IsUUID() + sessionId!: string; + + @IsNotEmpty() + @IsString() + userId!: string; +} diff --git a/src/modules/chatbot/chatbot.controller.ts b/src/modules/chatbot/chatbot.controller.ts new file mode 100644 index 0000000..bb5d509 --- /dev/null +++ b/src/modules/chatbot/chatbot.controller.ts @@ -0,0 +1,135 @@ +import { Request, Response } from "express"; +import { inject } from "inversify"; +import { controller, httpGet, httpPost, httpPut, queryParam, request, requestBody, requestParam, response } from "inversify-express-utils"; + +import { SendMessageDto } from "./DTO/send-message.dto"; +import { SessionIdParamDto } from "./DTO/session-id.param.dto"; +import { ChatbotService } from "./providers/chatbot.service"; +import { getOrCreateChatbotUlid } from "./utils/ulid.util"; +import { BaseController } from "../../common/base/controller"; +import { ApiModel, ApiOperation, ApiParam, ApiResponse, ApiTags } from "../../common/decorator/swggerDocs"; +import { PaginationDTO } from "../../common/dto/pagination.dto"; +import { HttpStatus } from "../../common/enums/httpStatus.enum"; +import { ValidationMiddleware } from "../../core/middlewares/validator.middleware"; +import { IOCTYPES } from "../../IOC/ioc.types"; + +@controller("/chatbot") +@ApiTags("Chatbot") +class ChatbotController extends BaseController { + @inject(IOCTYPES.ChatbotService) private chatbotService: ChatbotService; + + @ApiOperation("Create a new chat session") + @ApiResponse("Chat session created successfully", HttpStatus.Created) + @httpPost("/sessions") + public async createSession(@request() req: Request, @response() res: Response) { + const ulid = getOrCreateChatbotUlid(req, res); + const data = await this.chatbotService.createChatSession(ulid); + return this.response({ data }, HttpStatus.Created); + } + + @ApiOperation("Get user's chat sessions") + @ApiResponse("Chat sessions retrieved successfully") + @httpGet("/sessions", ValidationMiddleware.validateQuery(PaginationDTO)) + public async getUserSessions(@request() req: Request, @response() res: Response, @queryParam() queryDto: PaginationDTO) { + const ulid = getOrCreateChatbotUlid(req, res); + const limit = queryDto.limit || 10; + const data = await this.chatbotService.getUserChatSessions(ulid, limit); + return this.response({ data }); + } + + @ApiOperation("Get a specific chat session with messages") + @ApiResponse("Chat session retrieved successfully") + @ApiParam("sessionId", "Session ID", true) + @httpGet("/sessions/:sessionId", ValidationMiddleware.validateParameter(SessionIdParamDto)) + public async getSession(@request() req: Request, @response() res: Response, @requestParam() param: SessionIdParamDto) { + const ulid = getOrCreateChatbotUlid(req, res); + const data = await this.chatbotService.getChatSession(param.sessionId, ulid); + return this.response({ data }); + } + + @ApiOperation("Send a message in a chat session") + @ApiResponse("Message sent successfully", HttpStatus.Created) + @ApiModel(SendMessageDto) + @httpPost("/messages", ValidationMiddleware.validateInput(SendMessageDto)) + public async sendMessage(@request() req: Request, @response() res: Response, @requestBody() sendDto: SendMessageDto) { + const ulid = getOrCreateChatbotUlid(req, res); + const data = await this.chatbotService.sendMessage(ulid, sendDto); + return this.response({ data }, HttpStatus.Created); + } + + @ApiOperation("Send a message and get streaming response") + @ApiResponse("Streaming response initiated") + @ApiModel(SendMessageDto) + @httpPost("/messages/stream", ValidationMiddleware.validateInput(SendMessageDto)) + public async sendMessageStream( + @request() req: Request, + @requestBody() sendDto: SendMessageDto, + @response() res: Response, + ): Promise { + const ulid = getOrCreateChatbotUlid(req, res); + try { + const { userMessage, streamGenerator } = await this.chatbotService.sendMessageStream(ulid, sendDto); + + // Set headers for Server-Sent Events + res.setHeader("Content-Type", "text/event-stream"); + res.setHeader("Cache-Control", "no-cache"); + res.setHeader("Connection", "keep-alive"); + res.setHeader("Access-Control-Allow-Origin", "*"); + + // Send initial user message + res.write(`data: ${JSON.stringify({ type: "user_message", data: userMessage })}\n\n`); + + // Start streaming bot response + res.write(`data: ${JSON.stringify({ type: "bot_response_start" })}\n\n`); + + let fullResponse = ""; + const stream = await streamGenerator(); + for await (const chunk of stream) { + if (chunk) { + fullResponse += chunk; + res.write(`data: ${JSON.stringify({ type: "bot_response_chunk", data: chunk })}\n\n`); + } + } + + // Save bot response to cache after streaming completes + if (fullResponse) { + await this.chatbotService.saveStreamedBotResponse(sendDto.sessionId, userMessage.id, fullResponse, ulid); + } + + // End streaming + res.write(`data: ${JSON.stringify({ type: "bot_response_end" })}\n\n`); + res.end(); + } catch (error) { + console.error(error); + res.write(`data: ${JSON.stringify({ type: "error", data: "خطا در تولید پاسخ هوشمند. لطفاً دوباره تلاش کنید ⚠️" })}\n\n`); + res.end(); + } + } + + @ApiOperation("Close a chat session") + @ApiResponse("Chat session closed successfully") + @ApiParam("sessionId", "Session ID", true) + @httpPut("/sessions/:sessionId/close", ValidationMiddleware.validateParameter(SessionIdParamDto)) + public async closeSession(@request() req: Request, @response() res: Response, @requestParam() param: SessionIdParamDto) { + const ulid = getOrCreateChatbotUlid(req, res); + const data = await this.chatbotService.closeChatSession(param.sessionId, ulid); + return this.response({ data }); + } + + @ApiOperation("Mark messages as read") + @ApiResponse("Messages marked as read successfully") + @ApiParam("sessionId", "Session ID", true) + @httpPut("/sessions/:sessionId/messages/read", ValidationMiddleware.validateParameter(SessionIdParamDto)) + public async markMessagesAsRead( + @request() req: Request, + @response() res: Response, + @requestParam() param: SessionIdParamDto, + @requestBody() body: { messageIds: string[] }, + ) { + const ulid = getOrCreateChatbotUlid(req, res); + const data = await this.chatbotService.markMessagesAsRead(param.sessionId, ulid, body.messageIds); + return this.response({ data }); + } +} + +export { ChatbotController }; diff --git a/src/modules/chatbot/chatbot.gateway.ts b/src/modules/chatbot/chatbot.gateway.ts new file mode 100644 index 0000000..fee9f52 --- /dev/null +++ b/src/modules/chatbot/chatbot.gateway.ts @@ -0,0 +1,339 @@ +import { Server as HttpServer } from "http"; + +import { inject, injectable } from "inversify"; +import { Server, Socket } from "socket.io"; + +import { ChatbotService } from "./providers/chatbot.service"; +import { WebSocketAuthService } from "./providers/websocket-auth.service"; +import { SendMessageDto } from "./DTO/send-message.dto"; +import { WEBSOCKET_EVENTS } from "./constants/chatbot.constants"; +import { IOCTYPES } from "../../IOC/ioc.types"; +import { Logger } from "../../core/logging/logger"; +import { AuthenticatedSocket, WebSocketResponse } from "./interfaces/websocket.interface"; +import * as ulidLib from "ulid"; + +@injectable() +export class ChatbotGateway { + private io: Server; + private readonly logger: Logger; + private readonly CHATBOT_ULID_COOKIE_NAME = "chatbot_session_id"; + + constructor( + @inject(IOCTYPES.ChatbotService) private chatbotService: ChatbotService, + @inject(IOCTYPES.ChatbotWebSocketAuthService) private wsAuthService: WebSocketAuthService, + ) { + this.logger = new Logger("ChatbotGateway"); + } + + public initialize(server: HttpServer): Server { + this.io = new Server(server, { + path: "/ws-chatbot", + cors: { + origin: true, + allowedHeaders: ["Authorization"], + credentials: true, + methods: ["GET", "POST"], + }, + }); + + this.io.on("connection", (socket: Socket) => this.handleConnection(socket)); + this.logger.info("ChatbotGateway initialized on /ws-chatbot"); + return this.io; + } + + private async handleConnection(socket: Socket): Promise { + try { + // Try to authenticate if token is provided + const token = socket.handshake.query.token as string; + let ulid: string | undefined; + + if (token) { + // Authenticated connection + const authResult = await this.wsAuthService.authenticateClient(socket as AuthenticatedSocket); + if (!authResult.success) { + this.wsAuthService.handleAuthenticationFailure(socket, authResult.error || "Authentication failed"); + return; + } + this.wsAuthService.emitAuthenticationSuccess(socket, authResult.user!); + ulid = authResult.user!.id; + } else { + // Anonymous connection - generate or get ULID from cookie + const cookieHeader = socket.handshake.headers.cookie; + ulid = this.extractUlidFromCookie(cookieHeader); + if (!ulid) { + ulid = ulidLib.ulid(); + } + socket.data = { user: { id: ulid, sub: ulid } }; + } + + this.logger.info(`Client connected: ${socket.id}, ULID: ${ulid}`); + + // Register event handlers + this.registerEventHandlers(socket, ulid); + } catch (error) { + this.logger.error("Connection error", error); + socket.emit(WEBSOCKET_EVENTS.ERROR, { + status: "error", + message: "Connection failed", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + socket.disconnect(); + } + } + + private extractUlidFromCookie(cookieHeader?: string | string[]): string | undefined { + if (!cookieHeader || typeof cookieHeader !== "string") { + return undefined; + } + + const cookies = cookieHeader.split(";").reduce( + (acc, cookie) => { + const [key, value] = cookie.trim().split("="); + if (key && value) { + acc[key] = decodeURIComponent(value); + } + return acc; + }, + {} as Record, + ); + + return cookies[this.CHATBOT_ULID_COOKIE_NAME]; + } + + private registerEventHandlers(socket: Socket, ulid: string): void { + // Session management + socket.on(WEBSOCKET_EVENTS.CREATE_SESSION, () => this.handleCreateSession(socket, ulid)); + socket.on(WEBSOCKET_EVENTS.JOIN_CHAT, (sessionId: string) => this.handleJoinChat(socket, ulid, sessionId)); + socket.on(WEBSOCKET_EVENTS.LEAVE_CHAT, (sessionId: string) => this.handleLeaveChat(socket, sessionId)); + + // Message handling + socket.on(WEBSOCKET_EVENTS.SEND_MESSAGE, (data: SendMessageDto) => this.handleSendMessage(socket, ulid, data)); + socket.on(WEBSOCKET_EVENTS.SEND_MESSAGE + "_stream", (data: SendMessageDto) => this.handleSendMessageStream(socket, ulid, data)); + + // Typing indicators + socket.on(WEBSOCKET_EVENTS.TYPING_START, (sessionId: string) => this.handleTypingStart(socket, sessionId)); + socket.on(WEBSOCKET_EVENTS.TYPING_STOP, (sessionId: string) => this.handleTypingStop(socket, sessionId)); + + // Disconnect + socket.on(WEBSOCKET_EVENTS.DISCONNECT, () => this.handleDisconnect(socket)); + } + + private async handleCreateSession(socket: Socket, ulid: string): Promise { + try { + const session = await this.chatbotService.createChatSession(ulid); + socket.emit(WEBSOCKET_EVENTS.SESSION_CREATED, { + status: "success", + data: session, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } catch (error) { + this.logger.error("Error creating session", error); + socket.emit(WEBSOCKET_EVENTS.SESSION_ERROR, { + status: "error", + message: "Failed to create session", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } + + private async handleJoinChat(socket: Socket, ulid: string, sessionId: string): Promise { + try { + const session = await this.chatbotService.getChatSession(sessionId, ulid); + socket.join(`session_${sessionId}`); + (socket as AuthenticatedSocket).data.sessionId = sessionId; + + socket.emit(WEBSOCKET_EVENTS.CHAT_JOINED, { + status: "success", + data: session, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Notify others in the session (if any) + socket.to(`session_${sessionId}`).emit(WEBSOCKET_EVENTS.USER_JOINED, { + status: "success", + data: { userId: ulid, sessionId }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } catch (error) { + this.logger.error("Error joining chat", error); + socket.emit(WEBSOCKET_EVENTS.SESSION_ERROR, { + status: "error", + message: "Failed to join chat session", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } + + private handleLeaveChat(socket: Socket, sessionId: string): void { + socket.leave(`session_${sessionId}`); + delete (socket as AuthenticatedSocket).data.sessionId; + + socket.emit(WEBSOCKET_EVENTS.CHAT_LEFT, { + status: "success", + message: "Left chat session", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + socket.to(`session_${sessionId}`).emit(WEBSOCKET_EVENTS.USER_LEFT, { + status: "success", + data: { sessionId }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + + private async handleSendMessage(socket: Socket, ulid: string, data: SendMessageDto): Promise { + try { + // Save user message + const userMessage = await this.chatbotService.sendMessage(ulid, data); + + // Emit user message to client + socket.emit(WEBSOCKET_EVENTS.MESSAGE_RECEIVED, { + status: "success", + data: userMessage, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Emit to others in the session (if any) + socket.to(`session_${data.sessionId}`).emit(WEBSOCKET_EVENTS.MESSAGE_RECEIVED, { + status: "success", + data: userMessage, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Generate bot response asynchronously + this.generateBotResponseAsync(socket, data.sessionId, userMessage.id, ulid); + } catch (error) { + this.logger.error("Error sending message", error); + socket.emit(WEBSOCKET_EVENTS.MESSAGE_ERROR, { + status: "error", + message: "Failed to send message", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } + + private async handleSendMessageStream(socket: Socket, ulid: string, data: SendMessageDto): Promise { + try { + const { userMessage, streamGenerator } = await this.chatbotService.sendMessageStream(ulid, data); + + // Emit user message + socket.emit(WEBSOCKET_EVENTS.MESSAGE_RECEIVED, { + status: "success", + data: userMessage, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Emit to others in the session (if any) + socket.to(`session_${data.sessionId}`).emit(WEBSOCKET_EVENTS.MESSAGE_RECEIVED, { + status: "success", + data: userMessage, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Start streaming bot response + socket.emit(WEBSOCKET_EVENTS.BOT_RESPONSE_START, { + status: "success", + data: { userMessageId: userMessage.id }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + let fullResponse = ""; + const stream = await streamGenerator(); + for await (const chunk of stream) { + if (chunk) { + fullResponse += chunk; + socket.emit(WEBSOCKET_EVENTS.BOT_RESPONSE_CHUNK, { + status: "success", + data: { chunk, userMessageId: userMessage.id }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } + + // Save bot response after streaming completes + if (fullResponse) { + await this.chatbotService.saveStreamedBotResponse(data.sessionId, userMessage.id, fullResponse, ulid); + } + + // Emit final bot response + socket.emit(WEBSOCKET_EVENTS.BOT_RESPONSE_END, { + status: "success", + data: { userMessageId: userMessage.id, fullResponse }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } catch (error) { + this.logger.error("Error sending stream message", error); + socket.emit(WEBSOCKET_EVENTS.MESSAGE_ERROR, { + status: "error", + message: "Failed to send message", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } + + private async generateBotResponseAsync(socket: Socket, sessionId: string, userMessageId: string, ulid: string): Promise { + try { + // Get the session to check if it exists + const session = await this.chatbotService.getChatSession(sessionId, ulid); + if (!session) { + return; + } + + // Emit that bot is typing + socket.emit(WEBSOCKET_EVENTS.BOT_RESPONSE_START, { + status: "success", + data: { userMessageId }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Generate bot response and wait for it + const botMessage = await this.chatbotService.generateBotResponse(sessionId, userMessageId, ulid); + + if (botMessage) { + // Emit bot response to the client + socket.emit(WEBSOCKET_EVENTS.BOT_RESPONSE, { + status: "success", + data: botMessage, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + + // Also emit to others in the session (if any) + socket.to(`session_${sessionId}`).emit(WEBSOCKET_EVENTS.BOT_RESPONSE, { + status: "success", + data: botMessage, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } catch (error) { + this.logger.error("Error generating bot response", error); + socket.emit(WEBSOCKET_EVENTS.BOT_RESPONSE, { + status: "error", + message: "Failed to generate bot response", + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + } + + private handleTypingStart(socket: Socket, sessionId: string): void { + socket.to(`session_${sessionId}`).emit(WEBSOCKET_EVENTS.TYPING_START, { + status: "success", + data: { sessionId, userId: socket.data?.user?.id }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + + private handleTypingStop(socket: Socket, sessionId: string): void { + socket.to(`session_${sessionId}`).emit(WEBSOCKET_EVENTS.TYPING_STOP, { + status: "success", + data: { sessionId, userId: socket.data?.user?.id }, + timestamp: new Date().toISOString(), + } as WebSocketResponse); + } + + private handleDisconnect(socket: Socket): void { + const user = socket.data?.user; + if (user) { + this.logger.info(`Client disconnected: ${socket.id}, ULID: ${user.id}`); + } + } +} diff --git a/src/modules/chatbot/constants/chatbot.constants.ts b/src/modules/chatbot/constants/chatbot.constants.ts new file mode 100644 index 0000000..cb5de5f --- /dev/null +++ b/src/modules/chatbot/constants/chatbot.constants.ts @@ -0,0 +1,340 @@ +export const CHATBOT_CONSTANTS = { + DEFAULT_MODEL: "gpt-4o", + DEFAULT_TEMPERATURE: 0.7, + DEFAULT_MAX_TOKENS: 2000, + MAX_CONVERSATION_HISTORY: 20, + + // General system prompt focused on company guidance + SYSTEM_PROMPT: `شما یک دستیار هوشمند AI برای پلتفرم DZone هستید که در راهنمایی کاربران درباره شرکت‌ها و صنایع تخصص دارید. شما باید تمامی پاسخ‌ها را به زبان فارسی ارائه دهید. + +## وظایف اصلی شما: +🏢 **راهنمایی شرکت‌ها**: کمک به کاربران برای یافتن شرکت‌های مناسب و اطلاعات تماس آن‌ها +🏭 **معرفی صنایع**: ارائه اطلاعات کامل درباره صنایع مختلف و شرکت‌های فعال در هر حوزه +💼 **محصولات و خدمات**: راهنمایی درباره محصولات و خدمات ارائه شده توسط شرکت‌ها +🔍 **جستجو و مقایسه**: کمک به انتخاب بهترین شرکت برای نیازهای خاص کاربران +📞 **اطلاعات تماس**: ارائه اطلاعات کامل تماس شامل تلفن، ایمیل و آدرس شرکت‌ها + +## اصول پاسخگویی: +✅ همیشه به زبان فارسی پاسخ دهید +✅ از اطلاعات پایگاه داده شرکت‌ها و صنایع برای پاسخ‌های دقیق استفاده کنید +✅ در صورت سوال درباره شرکت خاص، اطلاعات کامل آن را ارائه دهید +✅ همیشه اطلاعات تماس (تلفن، ایمیل، آدرس) شرکت‌ها را اضافه کنید +✅ اگر چیزی نمی‌دانید، صادقانه اعتراف کنید و راه حل جایگزین پیشنهاد دهید +✅ از ایموجی‌ها برای بهتر کردن تجربه کاربر استفاده کنید +✅ پاسخ‌های خود را ساختاربندی و خوانا ارائه دهید + +## قالب پاسخ برای شرکت‌ها: +📍 **نام شرکت**: [نام کامل] +📞 **تلفن**: [شماره تماس] +📧 **ایمیل**: [آدرس ایمیل] +🏠 **آدرس**: [آدرس کامل] +🏭 **صنعت**: [نوع صنعت] +💼 **محصولات/خدمات**: [فهرست محصولات و خدمات] + +اگر اطلاعات شرکت‌ها یا صنایع مربوطه از پایگاه داده در دسترس است، حتماً از آن استفاده کرده و منبع آن را مشخص کنید.`, + + // New company guidance system prompt + COMPANY_GUIDANCE_SYSTEM_PROMPT: `شما یک دستیار هوشمند برای راهنمایی کاربران درباره شرکت‌ها و صنایع موجود در پلتفرم DZone هستید. + +## ماموریت شما: +- راهنمایی کاربران برای یافتن شرکت‌های مناسب +- ارائه اطلاعات کامل درباره شرکت‌ها، محصولات و خدمات +- معرفی صنایع مختلف و شرکت‌های فعال در هر حوزه +- کمک به انتخاب بهترین شرکت برای نیازهای کاربران + +## اطلاعات شرکت‌ها و صنایع: +{context} + +## تاریخچه مکالمه: +{conversation_history} + +## سوال کاربر: +{question} + +## دستورالعمل پاسخ: +بر اساس اطلاعات شرکت‌ها و صنایع موجود، پاسخی جامع و کاربردی به زبان فارسی ارائه دهید: + +**نکات مهم:** +- اگر کاربر درباره شرکت خاصی سوال می‌پرسد، اطلاعات کامل آن شرکت (نام، آدرس، تلفن، محصولات، خدمات) را ارائه دهید +- اگر سوال درباره صنعت خاصی است، شرکت‌های فعال در آن صنعت را معرفی کنید +- در صورت جستجوی محصول یا خدمت، شرکت‌های مرتبط را پیشنهاد دهید +- همیشه اطلاعات تماس (تلفن، ایمیل، آدرس) شرکت‌ها را ارائه دهید +- از ایموجی‌های مناسب برای بهتر شدن خوانایی استفاده کنید +- پاسخ خود را ساختاربندی و منظم ارائه دهید + +پاسخ شما:`, + + // Company data templates + COMPANY_DATA_TEMPLATES: { + COMPANY_INFO: `نام شرکت: {name} +مدیر عامل: {ceo} +ایمیل: {email} +تلفن: {phone} +شماره ثبت: {registrationNumber} +تاریخ تاسیس: {establishmentDate} +آدرس: {address} +وب‌سایت: {website} +توضیحات: {description} +صنعت فعالیت: {industry} +وضعیت: {status} +کسب و کار والد: {business} + +این شرکت در پلتفرم DZone ثبت شده و آماده ارائه خدمات و محصولات خود به مشتریان می‌باشد.`, + + COMPANY_PRODUCT: `محصول شرکت {companyName}: +نام محصول: {productTitle} +شرکت تولیدکننده: {companyName} +صنعت: {industry} +توضیحات شرکت: {companyDescription} +آدرس شرکت: {companyAddress} +تلفن تماس: {companyPhone} +ایمیل: {companyEmail} + +این محصول توسط شرکت {companyName} در صنعت {industry} تولید و ارائه می‌شود.`, + + COMPANY_SERVICE: `خدمات شرکت {companyName}: +نام خدمت: {serviceTitle} +شرکت ارائه‌دهنده: {companyName} +صنعت: {industry} +توضیحات شرکت: {companyDescription} +آدرس شرکت: {companyAddress} +تلفن تماس: {companyPhone} +ایمیل: {companyEmail} + +این خدمت توسط شرکت {companyName} در صنعت {industry} ارائه می‌شود.`, + + INDUSTRY_INFO: `صنعت: {title} +وضعیت: {status} +تعداد شرکت‌های فعال: {companiesCount} +کسب و کار والد: {business} + +شرکت‌های فعال در این صنعت: +{companiesList} + +این صنعت شامل شرکت‌هایی است که در حوزه {title} فعالیت می‌کنند و خدمات و محصولات مختلفی ارائه می‌دهند.`, + }, + + // Company guidance documents + COMPANY_GUIDANCE_DOCUMENTS: { + COMPANY_GUIDE: `راهنمای شرکت‌ها در پلتفرم DZone: + +🏢 **اطلاعات شرکت‌ها**: +در پلتفرم DZone شرکت‌های مختلفی در صنایع متنوع فعالیت می‌کنند. هر شرکت دارای اطلاعات کاملی شامل: +- نام و مشخصات کامل شرکت +- اطلاعات تماس (تلفن، ایمیل، آدرس) +- صنعت فعالیت و حوزه کاری +- محصولات و خدمات ارائه شده +- وضعیت فعالیت و مجوزهای لازم + +💼 **محصولات و خدمات**: +هر شرکت می‌تواند محصولات و خدمات متنوعی ارائه دهد: +- محصولات: کالاها و تولیدات شرکت +- خدمات: سرویس‌ها و راهکارهای ارائه شده +- اطلاعات تماس برای استعلام و سفارش +- جزئیات کامل درباره هر محصول و خدمت + +🏭 **دسته‌بندی صنایع**: +شرکت‌ها بر اساس صنعت فعالیت دسته‌بندی شده‌اند: +- هر صنعت شامل شرکت‌های مرتبط +- امکان جستجو بر اساس نوع صنعت +- مشاهده تمام شرکت‌های یک صنعت خاص +- اطلاعات کامل درباره حوزه فعالیت`, + + COMPANY_FAQ: `سوالات متداول درباره شرکت‌ها: + +❓ **چگونه اطلاعات یک شرکت را پیدا کنم؟** +می‌توانید با ذکر نام شرکت یا صنعت مورد نظر، اطلاعات کامل شامل آدرس، تلفن و محصولات آن را دریافت کنید. + +❓ **چه شرکت‌هایی در یک صنعت خاص فعالیت می‌کنند؟** +با ذکر نام صنعت، لیست کامل شرکت‌های فعال در آن حوزه را مشاهده خواهید کرد. + +❓ **چگونه با یک شرکت تماس بگیرم؟** +اطلاعات تماس شامل تلفن، ایمیل و آدرس تمام شرکت‌ها در پلتفرم موجود است. + +❓ **چه محصولات و خدماتی ارائه می‌شود؟** +هر شرکت لیست کاملی از محصولات و خدمات خود ارائه داده که قابل مشاهده است. + +❓ **چگونه شرکت مناسب برای نیاز خودم پیدا کنم؟** +با توصیف نیاز خود یا نوع محصول/خدمت مورد نظر، شرکت‌های مرتبط را معرفی خواهم کرد.`, + + SEARCH_GUIDE: `راهنمای جستجو و یافتن شرکت‌ها: + +🔍 **نحوه جستجو**: +- جستجو بر اساس نام شرکت +- جستجو بر اساس نوع صنعت +- جستجو بر اساس محصول یا خدمت +- جستجو بر اساس منطقه جغرافیایی + +📊 **اطلاعات قابل دریافت**: +- اطلاعات کامل شرکت (نام، آدرس، تلفن) +- لیست محصولات و توضیحات +- لیست خدمات ارائه شده +- صنعت و حوزه فعالیت +- وضعیت فعالیت شرکت + +🎯 **توصیه‌های مفید**: +- برای یافتن شرکت مناسب، نیاز خود را واضح بیان کنید +- از کلمات کلیدی مرتبط با صنعت استفاده کنید +- در صورت نیاز، اطلاعات تماس برای پیگیری بیشتر درخواست کنید +- برای مقایسه، چندین شرکت در یک صنعت را بررسی کنید`, + + INDUSTRY_GUIDE: `راهنمای صنایع و حوزه‌های فعالیت: + +🏭 **انواع صنایع موجود**: +در پلتفرم DZone شرکت‌های مختلفی در صنایع متنوع حضور دارند: +- صنایع تولیدی و ساختمانی +- خدمات فناوری اطلاعات +- صنایع غذایی و کشاورزی +- خدمات مالی و بیمه +- حمل و نقل و لجستیک +- آموزش و تربیت +- سلامت و درمان +- و سایر صنایع + +💡 **مزایای دسته‌بندی صنعتی**: +- یافتن آسان شرکت‌های مرتبط +- مقایسه خدمات در یک صنعت +- انتخاب بهترین گزینه برای نیاز شما +- دسترسی به اطلاعات تخصصی هر حوزه + +📈 **نحوه استفاده**: +- ابتدا صنعت مورد نظر خود را مشخص کنید +- سپس شرکت‌های فعال در آن صنعت را بررسی کنید +- محصولات و خدمات مناسب را انتخاب کنید +- برای اطلاعات بیشتر با شرکت تماس بگیرید`, + }, + + // Logging messages for LangChain service + LANGCHAIN_MESSAGES: { + LOADING_DATA: "Loading company and industry data for model training...", + NO_DATA_FOUND: "No data found for model training", + DOCUMENTS_SPLIT: "{totalDocs} documents split into {chunks} chunks", + VECTOR_STORE_READY: "Vector store successfully initialized - ready to guide users about companies and industries", + VECTOR_STORE_ERROR: "Error initializing vector store", + DOCUMENTS_LOADED: "{count} documents loaded from company and industry database", + LOADING_ERROR: "Error loading company and industry documents", + SERVICE_NOT_INITIALIZED: "LangChain service is not initialized", + RESPONSE_ERROR: "Error generating LangChain response", + STREAM_RESPONSE_ERROR: "Error generating LangChain stream response", + REFRESHING_VECTOR_STORE: "Refreshing vector store with new company and industry data...", + }, + + // Source type labels in Farsi + SOURCE_TYPE_LABELS: { + company: "شرکت", + product: "محصول", + service: "خدمت", + industry: "صنعت", + company_guide: "راهنمای شرکت‌ها", + company_faq: "سوالات متداول", + search_guide: "راهنمای جستجو", + industry_guide: "راهنمای صنایع", + }, + + // Default values for missing data + DEFAULT_VALUES: { + NO_WEBSITE: "ندارد", + UNSPECIFIED: "مشخص نشده", + ACTIVE: "فعال", + INACTIVE: "غیرفعال", + }, + + // WEBSOCKET_EVENTS: { + // // Connection events + // CONNECT: "connect", + // DISCONNECT: "disconnect", + // CONNECTION_ERROR: "connection_failed", + + // // Authentication events + // AUTHENTICATE: "authenticate", + // AUTHENTICATED: "authenticated", + // UNAUTHORIZED: "unauthorized", + + // // Session events + // CREATE_SESSION: "create_session", + // SESSION_CREATED: "session_created", + // JOIN_CHAT: "join_chat", + // CHAT_JOINED: "chat_joined", + // LEAVE_CHAT: "leave_chat", + // CHAT_LEFT: "chat_left", + + // // Message events + // SEND_MESSAGE: "send_message", + // MESSAGE_RECEIVED: "message_received", + // BOT_RESPONSE_START: "bot_response_start", + // BOT_RESPONSE_CHUNK: "bot_response_chunk", + // BOT_RESPONSE_END: "bot_response_end", + // BOT_RESPONSE: "bot_response", + + // // Status events + // TYPING_START: "typing_start", + // TYPING_STOP: "typing_stop", + // USER_JOINED: "user_joined", + // USER_LEFT: "user_left", + + // // Error events + // ERROR: "error", + // WEBSOCKET_ERROR: "websocket_error", + // SESSION_ERROR: "session_error", + // MESSAGE_ERROR: "message_error", + // }, + + ERROR_MESSAGES: { + SESSION_NOT_FOUND: "جلسه چت یافت نشد", + UNAUTHORIZED_ACCESS: "دسترسی غیرمجاز به جلسه چت", + MESSAGE_TOO_LONG: "پیام از حداکثر طول مجاز تجاوز کرده است", + RATE_LIMIT_EXCEEDED: "تعداد پیام‌های ارسالی بیش از حد مجاز است", + LLM_SERVICE_ERROR: "سرویس هوش مصنوعی موقتاً در دسترس نیست", + INVALID_TOKEN: "توکن احراز هویت نامعتبر است", + AUTHENTICATION_REQUIRED: "احراز هویت ضروری است", + SESSION_CREATION_FAILED: "ایجاد جلسه چت با شکست مواجه شد", + CONNECTION_FAILED: "اتصال برقرار نشد", + WEBSOCKET_ERROR: "خطا در ارتباط WebSocket", + }, + + SUCCESS_MESSAGES: { + SESSION_CREATED: "جلسه چت با موفقیت ایجاد شد", + CHAT_JOINED: "با موفقیت به چت متصل شدید", + MESSAGE_SENT: "پیام ارسال شد", + AUTHENTICATED: "احراز هویت موفقیت‌آمیز", + }, +} as const; + +export const enum WEBSOCKET_EVENTS { + CONNECT = "connect", + DISCONNECT = "disconnect", + CONNECTION_ERROR = "connection_failed", + + // Authentication events + AUTHENTICATE = "authenticate", + AUTHENTICATED = "authenticated", + UNAUTHORIZED = "unauthorized", + + // Session events + CREATE_SESSION = "create_session", + SESSION_CREATED = "session_created", + JOIN_CHAT = "join_chat", + CHAT_JOINED = "chat_joined", + LEAVE_CHAT = "leave_chat", + CHAT_LEFT = "chat_left", + + // Message events + SEND_MESSAGE = "send_message", + MESSAGE_RECEIVED = "message_received", + BOT_RESPONSE_START = "bot_response_start", + BOT_RESPONSE_CHUNK = "bot_response_chunk", + BOT_RESPONSE_END = "bot_response_end", + BOT_RESPONSE = "bot_response", + + // Status events + TYPING_START = "typing_start", + TYPING_STOP = "typing_stop", + USER_JOINED = "user_joined", + USER_LEFT = "user_left", + + // Error events + ERROR = "error", + WEBSOCKET_ERROR = "websocket_error", + SESSION_ERROR = "session_error", + MESSAGE_ERROR = "message_error", +} diff --git a/src/modules/chatbot/interfaces/chatbot.interface.ts b/src/modules/chatbot/interfaces/chatbot.interface.ts new file mode 100644 index 0000000..af11748 --- /dev/null +++ b/src/modules/chatbot/interfaces/chatbot.interface.ts @@ -0,0 +1,36 @@ +export interface IChatbotResponse { + message: string; + confidence: number; + sources?: string[]; + tokensUsed: number; + context?: Record; +} + +export interface IChatContext { + userId: string; + sessionId: string; + conversationHistory: IChatMessage[]; + userPreferences?: Record; + businessContext?: Record; +} + +export interface IChatMessage { + content: string; + type: "user" | "bot" | "system"; + timestamp: Date; + metadata?: Record; +} + +export interface ILLMConfig { + model: string; + temperature: number; + maxTokens: number; + topP: number; + apiKey: string; +} + +export interface IDataSource { + query: (question: string, context?: Record) => Promise; + type: "database" | "vector" | "api"; + name: string; +} diff --git a/src/modules/chatbot/interfaces/websocket.interface.ts b/src/modules/chatbot/interfaces/websocket.interface.ts new file mode 100644 index 0000000..47b850c --- /dev/null +++ b/src/modules/chatbot/interfaces/websocket.interface.ts @@ -0,0 +1,50 @@ +import { Socket } from "socket.io"; + +/** + * Extended Socket interface with authenticated user data + */ +export interface AuthenticatedSocket extends Socket { + data: { + user: { id: string; sub: string }; + sessionId?: string; + }; +} + +/** + * User connection information stored in gateway + */ +export interface ConnectedUserInfo { + userId: string; + sessionId?: string; +} + +/** + * WebSocket error context for logging and debugging + */ +export interface WebSocketErrorContext { + clientId: string; + userId?: string; + sessionId?: string; + event?: string; + data?: any; + timestamp: Date; +} + +/** + * WebSocket event response structure + */ +export interface WebSocketResponse { + status: "success" | "error"; + message?: string; + data?: T; + timestamp?: string; +} + +/** + * Authentication result for WebSocket connections + */ +export interface AuthenticationResult { + success: boolean; + user?: { id: string; sub: string }; + error?: string; +} diff --git a/src/modules/chatbot/models/Abstraction/IChatMessage.ts b/src/modules/chatbot/models/Abstraction/IChatMessage.ts new file mode 100644 index 0000000..15956ec --- /dev/null +++ b/src/modules/chatbot/models/Abstraction/IChatMessage.ts @@ -0,0 +1,29 @@ +import { Types } from "mongoose"; + +export enum MessageType { + USER = "user", + BOT = "bot", + SYSTEM = "system", +} + +export enum MessageStatus { + SENT = "sent", + DELIVERED = "delivered", + READ = "read", + FAILED = "failed", +} + +export interface IChatMessage { + _id: Types.ObjectId; + content: string; + type: MessageType; + status: MessageStatus; + session: Types.ObjectId; + sender?: Types.ObjectId; + metadata?: Record; + responseToId?: string; + tokensUsed: number; + createdAt: Date; + updatedAt: Date; +} + diff --git a/src/modules/chatbot/models/Abstraction/IChatSession.ts b/src/modules/chatbot/models/Abstraction/IChatSession.ts new file mode 100644 index 0000000..6f94fa3 --- /dev/null +++ b/src/modules/chatbot/models/Abstraction/IChatSession.ts @@ -0,0 +1,19 @@ +import { Types } from "mongoose"; + +export enum ChatSessionStatus { + ACTIVE = "active", + CLOSED = "closed", + ARCHIVED = "archived", +} + +export interface IChatSession { + _id: Types.ObjectId; + title: string; + user: Types.ObjectId; + status: ChatSessionStatus; + context?: Record; + lastMessageAt?: Date; + createdAt: Date; + updatedAt: Date; +} + diff --git a/src/modules/chatbot/providers/chatbot.service.ts b/src/modules/chatbot/providers/chatbot.service.ts new file mode 100644 index 0000000..c7ad6d9 --- /dev/null +++ b/src/modules/chatbot/providers/chatbot.service.ts @@ -0,0 +1,455 @@ +import { inject, injectable } from "inversify"; +import * as ulidLib from "ulid"; + +import { LLMService } from "./llm.service"; +import { BadRequestError, NotFoundError } from "../../../core/app/app.errors"; +import { Logger } from "../../../core/logging/logger"; +import { IOCTYPES } from "../../../IOC/ioc.types"; +import { CacheService } from "../../../utils/cache.service"; +import { SendMessageDto } from "../DTO/send-message.dto"; +import { IChatContext } from "../interfaces/chatbot.interface"; +import { MessageStatus, MessageType } from "../models/Abstraction/IChatMessage"; +import { ChatSessionStatus } from "../models/Abstraction/IChatSession"; + +interface CachedSession { + id: string; + ulid: string; + status: ChatSessionStatus; + lastMessageAt?: string; + createdAt: string; +} + +interface CachedMessage { + id: string; + content: string; + type: MessageType; + status: MessageStatus; + sessionId: string; + metadata?: Record; + responseToId?: string; + tokensUsed: number; + createdAt: string; +} + +@injectable() +export class ChatbotService { + private readonly logger: Logger; + private readonly CACHE_TTL_SECONDS = 15 * 60; // 15 minutes + + constructor( + @inject(IOCTYPES.ChatbotLLMService) private llmService: LLMService, + @inject(IOCTYPES.CacheService) private cacheService: CacheService, + ) { + this.logger = new Logger("ChatbotService"); + this.logger.info("Using OpenAI as LLM provider"); + } + + private getSessionCacheKey(ulid: string, sessionId: string): string { + return `chatbot:session:${ulid}:${sessionId}`; + } + + private getSessionsListCacheKey(ulid: string): string { + return `chatbot:sessions:${ulid}`; + } + + private getMessagesCacheKey(sessionId: string): string { + return `chatbot:messages:${sessionId}`; + } + + async createChatSession(ulid: string) { + const sessionId = ulidLib.ulid(); + const now = new Date().toISOString(); + + const session: CachedSession = { + id: sessionId, + ulid, + status: ChatSessionStatus.ACTIVE, + lastMessageAt: now, + createdAt: now, + }; + + // Save session to cache + await this.cacheService.setAsync(this.getSessionCacheKey(ulid, sessionId), JSON.stringify(session), this.CACHE_TTL_SECONDS); + + // Add to sessions list + const sessionsListKey = this.getSessionsListCacheKey(ulid); + const existingSessionsJson = await this.cacheService.getAsync(sessionsListKey); + const sessionsList: string[] = existingSessionsJson ? JSON.parse(existingSessionsJson) : []; + sessionsList.push(sessionId); + await this.cacheService.setAsync(sessionsListKey, JSON.stringify(sessionsList), this.CACHE_TTL_SECONDS); + + // Initialize messages array for this session + await this.cacheService.setAsync(this.getMessagesCacheKey(sessionId), JSON.stringify([]), this.CACHE_TTL_SECONDS); + + return this.mapSessionToDto(session); + } + + async getUserChatSessions(ulid: string, limit = 10) { + const sessionsListKey = this.getSessionsListCacheKey(ulid); + const sessionsListJson = await this.cacheService.getAsync(sessionsListKey); + + if (!sessionsListJson) { + return []; + } + + const sessionIds: string[] = JSON.parse(sessionsListJson); + const sessions: CachedSession[] = []; + + // Get sessions in reverse order (newest first) + for (let i = sessionIds.length - 1; i >= 0 && sessions.length < limit; i--) { + const sessionId = sessionIds[i]; + const sessionKey = this.getSessionCacheKey(ulid, sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + + if (sessionJson) { + const session: CachedSession = JSON.parse(sessionJson); + // Get messages for this session + const messagesJson = await this.cacheService.getAsync(this.getMessagesCacheKey(sessionId)); + const messages: CachedMessage[] = messagesJson ? JSON.parse(messagesJson) : []; + sessions.push({ ...session, messages } as any); + } + } + + return sessions.map((session) => this.mapSessionToDto(session)); + } + + async getChatSession(sessionId: string, ulid: string) { + const sessionKey = this.getSessionCacheKey(ulid, sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + + if (!sessionJson) { + throw new NotFoundError("Chat session not found"); + } + + const session: CachedSession = JSON.parse(sessionJson); + + // Get messages for this session + const messagesJson = await this.cacheService.getAsync(this.getMessagesCacheKey(sessionId)); + const messages: CachedMessage[] = messagesJson ? JSON.parse(messagesJson) : []; + + return this.mapSessionToDto({ ...session, messages } as any); + } + + async sendMessage(ulid: string, sendDto: SendMessageDto) { + const sessionKey = this.getSessionCacheKey(ulid, sendDto.sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + + if (!sessionJson) { + throw new NotFoundError("Chat session not found"); + } + + const session: CachedSession = JSON.parse(sessionJson); + + if (session.status !== ChatSessionStatus.ACTIVE) { + throw new BadRequestError("Cannot send message to inactive session"); + } + + // Create user message + const messageId = ulidLib.ulid(); + const now = new Date().toISOString(); + const userMessage: CachedMessage = { + id: messageId, + content: sendDto.content, + type: MessageType.USER, + status: MessageStatus.SENT, + sessionId: sendDto.sessionId, + responseToId: sendDto.responseToId, + metadata: sendDto.metadata, + tokensUsed: 0, + createdAt: now, + }; + + // Save message to cache + const messagesKey = this.getMessagesCacheKey(sendDto.sessionId); + const messagesJson = await this.cacheService.getAsync(messagesKey); + const messages: CachedMessage[] = messagesJson ? JSON.parse(messagesJson) : []; + messages.push(userMessage); + await this.cacheService.setAsync(messagesKey, JSON.stringify(messages), this.CACHE_TTL_SECONDS); + + // Update session last message time + session.lastMessageAt = now; + await this.cacheService.setAsync(sessionKey, JSON.stringify(session), this.CACHE_TTL_SECONDS); + + // Note: Bot response will be generated by the gateway after emitting user message + + return this.mapMessageToDto(userMessage); + } + + async sendMessageStream(ulid: string, sendDto: SendMessageDto) { + const sessionKey = this.getSessionCacheKey(ulid, sendDto.sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + + if (!sessionJson) { + throw new NotFoundError("Chat session not found"); + } + + const session: CachedSession = JSON.parse(sessionJson); + + if (session.status !== ChatSessionStatus.ACTIVE) { + throw new BadRequestError("Cannot send message to inactive session"); + } + + // Create user message + const messageId = ulidLib.ulid(); + const now = new Date().toISOString(); + const userMessage: CachedMessage = { + id: messageId, + content: sendDto.content, + type: MessageType.USER, + status: MessageStatus.SENT, + sessionId: sendDto.sessionId, + responseToId: sendDto.responseToId, + metadata: sendDto.metadata, + tokensUsed: 0, + createdAt: now, + }; + + // Save message to cache + const messagesKey = this.getMessagesCacheKey(sendDto.sessionId); + const messagesJson = await this.cacheService.getAsync(messagesKey); + const messages: CachedMessage[] = messagesJson ? JSON.parse(messagesJson) : []; + messages.push(userMessage); + await this.cacheService.setAsync(messagesKey, JSON.stringify(messages), this.CACHE_TTL_SECONDS); + + // Update session last message time + session.lastMessageAt = now; + await this.cacheService.setAsync(sessionKey, JSON.stringify(session), this.CACHE_TTL_SECONDS); + + // Get conversation history for context + const history = messages.slice(-20).map((msg) => ({ + content: msg.content, + type: msg.type as "user" | "bot" | "system", + timestamp: new Date(msg.createdAt), + metadata: msg.metadata, + })); + + // Build context for LLM + const context: IChatContext = { + userId: ulid, + sessionId: sendDto.sessionId, + conversationHistory: history, + }; + + // Return a function that generates the stream using LLM service + const streamGenerator = () => this.llmService.generateStreamResponse(sendDto.content, context); + + return { + userMessage: this.mapMessageToDto(userMessage), + streamGenerator, + }; + } + + async generateBotResponse( + sessionId: string, + userMessageId: string, + ulid: string, + ): Promise | null> { + try { + // Get messages from cache + const messagesKey = this.getMessagesCacheKey(sessionId); + const messagesJson = await this.cacheService.getAsync(messagesKey); + if (!messagesJson) { + throw new NotFoundError("Session not found"); + } + + const messages: CachedMessage[] = JSON.parse(messagesJson); + const userMessage = messages.find((msg) => msg.id === userMessageId); + if (!userMessage) { + throw new NotFoundError("User message not found"); + } + + // Get session + const sessionsListKey = this.getSessionsListCacheKey(ulid); + const sessionsListJson = await this.cacheService.getAsync(sessionsListKey); + if (!sessionsListJson) { + throw new NotFoundError("Session not found"); + } + + const sessionKey = this.getSessionCacheKey(ulid, sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + if (!sessionJson) { + throw new NotFoundError("Session not found"); + } + + const session: CachedSession = JSON.parse(sessionJson); + + // Get conversation history + const history = messages.slice(-20).map((msg) => ({ + content: msg.content, + type: msg.type as "user" | "bot" | "system", + timestamp: new Date(msg.createdAt), + metadata: msg.metadata, + })); + + // Build context for LLM + const context: IChatContext = { + userId: ulid, + sessionId, + conversationHistory: history, + }; + + // Generate response using LLM + const llmResponse = await this.llmService.generateResponse(userMessage.content, context); + + // Create bot message + const botMessageId = ulidLib.ulid(); + const now = new Date().toISOString(); + const botMessage: CachedMessage = { + id: botMessageId, + content: llmResponse.message, + type: MessageType.BOT, + status: MessageStatus.SENT, + sessionId, + responseToId: userMessageId, + tokensUsed: llmResponse.tokensUsed, + metadata: { + confidence: llmResponse.confidence, + sources: llmResponse.sources, + llmContext: llmResponse.context, + }, + createdAt: now, + }; + + // Save bot message to cache + messages.push(botMessage); + await this.cacheService.setAsync(messagesKey, JSON.stringify(messages), this.CACHE_TTL_SECONDS); + + // Update session last message time + session.lastMessageAt = now; + await this.cacheService.setAsync(sessionKey, JSON.stringify(session), this.CACHE_TTL_SECONDS); + + this.logger.info(`Generated bot response for session ${sessionId}`); + return this.mapMessageToDto(botMessage); + } catch (error) { + this.logger.error(`Failed to generate bot response for session ${sessionId}`, error); + + // Create error message + const messagesKey = this.getMessagesCacheKey(sessionId); + const messagesJson = await this.cacheService.getAsync(messagesKey); + const messages: CachedMessage[] = messagesJson ? JSON.parse(messagesJson) : []; + + const errorMessageId = ulidLib.ulid(); + const now = new Date().toISOString(); + const errorMessage: CachedMessage = { + id: errorMessageId, + content: "متأسفم، در حال حاضر مشکلی در پاسخگویی دارم. لطفاً دوباره تلاش کنید. 🙏", + type: MessageType.BOT, + status: MessageStatus.FAILED, + sessionId, + responseToId: userMessageId, + tokensUsed: 0, + createdAt: now, + }; + + messages.push(errorMessage); + await this.cacheService.setAsync(messagesKey, JSON.stringify(messages), this.CACHE_TTL_SECONDS); + return this.mapMessageToDto(errorMessage); + } + } + + async closeChatSession(sessionId: string, ulid: string) { + const sessionKey = this.getSessionCacheKey(ulid, sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + + if (!sessionJson) { + throw new NotFoundError("Chat session not found"); + } + + const session: CachedSession = JSON.parse(sessionJson); + session.status = ChatSessionStatus.CLOSED; + + await this.cacheService.setAsync(sessionKey, JSON.stringify(session), this.CACHE_TTL_SECONDS); + + return { message: "Chat session closed successfully" }; + } + + async markMessagesAsRead(sessionId: string, ulid: string, messageIds: string[]) { + const sessionKey = this.getSessionCacheKey(ulid, sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + + if (!sessionJson) { + throw new NotFoundError("Chat session not found"); + } + + // Get messages and update their status + const messagesKey = this.getMessagesCacheKey(sessionId); + const messagesJson = await this.cacheService.getAsync(messagesKey); + + if (messagesJson) { + const messages: CachedMessage[] = JSON.parse(messagesJson); + messages.forEach((msg) => { + if (messageIds.includes(msg.id)) { + msg.status = MessageStatus.READ; + } + }); + await this.cacheService.setAsync(messagesKey, JSON.stringify(messages), this.CACHE_TTL_SECONDS); + } + + return { message: "Messages marked as read successfully" }; + } + + async saveStreamedBotResponse(sessionId: string, userMessageId: string, botResponse: string, ulid: string): Promise { + try { + const messagesKey = this.getMessagesCacheKey(sessionId); + const messagesJson = await this.cacheService.getAsync(messagesKey); + if (!messagesJson) { + this.logger.warn(`Cannot save bot response: session ${sessionId} not found in cache`); + return; + } + + const messages: CachedMessage[] = JSON.parse(messagesJson); + const botMessageId = ulidLib.ulid(); + const now = new Date().toISOString(); + + const botMessage: CachedMessage = { + id: botMessageId, + content: botResponse, + type: MessageType.BOT, + status: MessageStatus.SENT, + sessionId, + responseToId: userMessageId, + tokensUsed: 0, // Could estimate tokens if needed + createdAt: now, + }; + + messages.push(botMessage); + await this.cacheService.setAsync(messagesKey, JSON.stringify(messages), this.CACHE_TTL_SECONDS); + + // Update session last message time + const sessionKey = this.getSessionCacheKey(ulid, sessionId); + const sessionJson = await this.cacheService.getAsync(sessionKey); + if (sessionJson) { + const session: CachedSession = JSON.parse(sessionJson); + session.lastMessageAt = now; + await this.cacheService.setAsync(sessionKey, JSON.stringify(session), this.CACHE_TTL_SECONDS); + } + + this.logger.info(`Saved streamed bot response for session ${sessionId}`); + } catch (error) { + this.logger.error(`Failed to save streamed bot response for session ${sessionId}`, error); + } + } + + private mapSessionToDto(session: CachedSession | any) { + return { + id: session.id, + status: session.status, + createdAt: session.createdAt, + lastMessageAt: session.lastMessageAt, + messages: session.messages ? session.messages.map((msg: CachedMessage) => this.mapMessageToDto(msg)) : [], + }; + } + + private mapMessageToDto(message: CachedMessage) { + return { + id: message.id, + content: message.content, + type: message.type, + status: message.status, + createdAt: message.createdAt, + responseToId: message.responseToId, + metadata: message.metadata, + tokensUsed: message.tokensUsed, + }; + } +} diff --git a/src/modules/chatbot/providers/data-context.service.ts b/src/modules/chatbot/providers/data-context.service.ts new file mode 100644 index 0000000..2263a7f --- /dev/null +++ b/src/modules/chatbot/providers/data-context.service.ts @@ -0,0 +1,79 @@ +import { injectable } from "inversify"; + +import { IChatContext } from "../interfaces/chatbot.interface"; +import { Logger } from "../../../core/logging/logger"; + +/** + * Service for retrieving relevant context data from the database + * This is a simplified version - can be extended to fetch actual data from your database models + */ +@injectable() +export class DataContextService { + private readonly logger: Logger; + + constructor() { + this.logger = new Logger("DataContextService"); + } + + async getRelevantContext(_message: string, _context: IChatContext): Promise<{ data: string[]; sources: string[] }> { + const relevantData: string[] = []; + const sources: string[] = []; + + try { + // Analyze message to determine what data to fetch + // const keywords = this._extractKeywords(_message); // TODO: Use when implementing keyword-based filtering + + // TODO: Implement actual data fetching based on your project's models + // For now, this returns empty data - you can extend this to fetch: + // - Product data from ProductModel + // - Category data from CategoryModel + // - Order data from OrderModel + // - User data from UserModel + // etc. + + // Example: If you want to fetch product data + // if (this.hasProductKeywords(keywords)) { + // const products = await ProductModel.find({ ... }).limit(5); + // products.forEach((product) => { + // relevantData.push(`Product: ${product.title_fa} - Price: ${product.price}`); + // sources.push("Products Database"); + // }); + // } + + return { data: relevantData, sources }; + } catch (error) { + this.logger.error("Failed to get relevant context", error); + return { data: [], sources: [] }; + } + } + + // TODO: Use this method when implementing keyword-based filtering + // @ts-expect-error - Method reserved for future use + // eslint-disable-next-line @typescript-eslint/no-unused-vars + private _extractKeywords(_message: string): string[] { + const commonWords = ["the", "is", "at", "which", "on", "a", "an", "and", "or", "but", "in", "with", "to", "for", "of", "as", "by"]; + + return _message + .toLowerCase() + .replace(/[^\w\s]/g, "") + .split(/\s+/) + .filter((word: string) => word.length > 2 && !commonWords.includes(word)); + } + + // Helper methods for keyword detection - can be extended + // TODO: Implement when needed + // private hasProductKeywords(keywords: string[]): boolean { + // const productTerms = ["product", "item", "goods", "merchandise", "کالا", "محصول"]; + // return keywords.some(keyword => productTerms.includes(keyword)); + // } + + // private hasOrderKeywords(keywords: string[]): boolean { + // const orderTerms = ["order", "purchase", "buy", "cart", "سفارش", "خرید"]; + // return keywords.some(keyword => orderTerms.includes(keyword)); + // } + + // private hasCategoryKeywords(keywords: string[]): boolean { + // const categoryTerms = ["category", "type", "kind", "دسته", "دسته‌بندی"]; + // return keywords.some(keyword => categoryTerms.includes(keyword)); + // } +} diff --git a/src/modules/chatbot/providers/llm.service.ts b/src/modules/chatbot/providers/llm.service.ts new file mode 100644 index 0000000..07a95bf --- /dev/null +++ b/src/modules/chatbot/providers/llm.service.ts @@ -0,0 +1,218 @@ +import { inject, injectable } from "inversify"; +import OpenAI from "openai"; + +import { DataContextService } from "./data-context.service"; +import { Logger } from "../../../core/logging/logger"; +import { IOCTYPES } from "../../../IOC/ioc.types"; +import { CHATBOT_CONSTANTS } from "../constants/chatbot.constants"; +import { IChatContext, IChatbotResponse, ILLMConfig } from "../interfaces/chatbot.interface"; + +@injectable() +export class LLMService { + private readonly logger: Logger; + private openai: OpenAI; + private config: ILLMConfig; + + constructor(@inject(IOCTYPES.ChatbotDataContextService) private dataContextService: DataContextService) { + this.logger = new Logger("LLMService"); + this.config = { + model: process.env.OPENAI_MODEL || CHATBOT_CONSTANTS.DEFAULT_MODEL, + temperature: Number(process.env.OPENAI_TEMPERATURE || CHATBOT_CONSTANTS.DEFAULT_TEMPERATURE), + maxTokens: Number(process.env.OPENAI_MAX_TOKENS || CHATBOT_CONSTANTS.DEFAULT_MAX_TOKENS), + topP: Number(process.env.OPENAI_TOP_P || "0.95"), + apiKey: process.env.OPENAI_API_KEY || "", + }; + + if (!this.config.apiKey) { + throw new Error("OPENAI_API_KEY is required"); + } + + this.openai = new OpenAI({ + apiKey: this.config.apiKey, + }); + } + + async generateResponse(message: string, context: IChatContext): Promise { + try { + // Get relevant data from database + const relevantData = await this.dataContextService.getRelevantContext(message, context); + + // Build system instruction with context + const systemInstruction = this.buildSystemInstruction(relevantData); + + // Build conversation history + const conversationHistory = this.buildConversationHistory(context); + + // Build messages array for OpenAI + const messages: OpenAI.Chat.Completions.ChatCompletionMessageParam[] = [ + { + role: "system", + content: systemInstruction, + }, + ...conversationHistory, + { + role: "user", + content: message, + }, + ]; + + // Generate response using OpenAI + const completion = await this.openai.chat.completions.create({ + model: this.config.model, + messages, + temperature: this.config.temperature, + max_tokens: this.config.maxTokens, + top_p: this.config.topP, + }); + + const botMessage = + completion.choices[0]?.message?.content || "متأسفم، نمی‌توانم پاسخی تولید کنم. لطفاً سوال خود را دوباره مطرح کنید. 🤖"; + + // Calculate token usage from OpenAI response + const tokensUsed = (completion.usage?.total_tokens || 0) + this.estimateTokens(systemInstruction); + + return { + message: botMessage, + confidence: this.calculateConfidence(completion), + sources: relevantData.sources, + tokensUsed, + context: { + model: this.config.model, + relevantDataFound: relevantData.data.length > 0, + finishReason: completion.choices[0]?.finish_reason || "unknown", + usage: completion.usage, + }, + }; + } catch (error) { + this.logger.error("Failed to generate OpenAI response", error); + throw new Error("Failed to generate response from AI service"); + } + } + + async generateStreamResponse(message: string, context: IChatContext): Promise> { + try { + // Get relevant data from database + const relevantData = await this.dataContextService.getRelevantContext(message, context); + + // Build system instruction with context + const systemInstruction = this.buildSystemInstruction(relevantData); + + // Build conversation history + const conversationHistory = this.buildConversationHistory(context); + + // Build messages array for OpenAI + const messages: OpenAI.Chat.Completions.ChatCompletionMessageParam[] = [ + { + role: "system", + content: systemInstruction, + }, + ...conversationHistory, + { + role: "user", + content: message, + }, + ]; + + // Generate streaming response using OpenAI + const stream = await this.openai.chat.completions.create({ + model: this.config.model, + messages, + temperature: this.config.temperature, + max_tokens: this.config.maxTokens, + top_p: this.config.topP, + stream: true, + }); + + return this.createAsyncIterableFromStream(stream); + } catch (error) { + this.logger.error("Failed to generate streaming OpenAI response", error); + throw new Error("Failed to generate streaming response from AI service"); + } + } + + private buildSystemInstruction(relevantData: { data: string[]; sources: string[] }): string { + let instruction = CHATBOT_CONSTANTS.SYSTEM_PROMPT; + + // Add relevant data context + if (relevantData.data.length > 0) { + instruction += "\n\nRelevant information from our database:\n"; + instruction += relevantData.data.map((data, index) => `${index + 1}. ${data}`).join("\n"); + instruction += "\n\nUse this information to provide accurate and specific answers."; + } + + return instruction; + } + + private buildConversationHistory(context: IChatContext): OpenAI.Chat.Completions.ChatCompletionMessageParam[] { + if (!context.conversationHistory || context.conversationHistory.length === 0) { + return []; + } + + const history: OpenAI.Chat.Completions.ChatCompletionMessageParam[] = []; + const recentHistory = context.conversationHistory + .slice(-CHATBOT_CONSTANTS.MAX_CONVERSATION_HISTORY) + .filter((msg) => msg.type !== "system"); + + for (const msg of recentHistory) { + history.push({ + role: msg.type === "user" ? "user" : "assistant", + content: msg.content, + }); + } + + return history; + } + + private async *createAsyncIterableFromStream(stream: AsyncIterable): AsyncIterable { + try { + for await (const chunk of stream) { + const content = chunk.choices[0]?.delta?.content; + if (content) { + yield content; + } + } + } catch (error) { + this.logger.error("Error in streaming response", error); + throw new Error("Failed to process streaming response"); + } + } + + private calculateConfidence(completion: OpenAI.Chat.Completions.ChatCompletion): number { + const text = completion.choices[0]?.message?.content || ""; + + // Check finish reason - lower confidence if stopped early + const finishReason = completion.choices[0]?.finish_reason; + if (finishReason === "length" || finishReason === "content_filter") { + return 0.4; + } + + // Check for uncertainty indicators + if ( + text.includes("I don't know") || + text.includes("I'm not sure") || + text.includes("uncertain") || + text.includes("نمی‌دانم") || + text.includes("مطمئن نیستم") + ) { + return 0.3; + } + + // Check response length and quality + if (text.length < 50) { + return 0.6; + } + + // Check if response uses provided data + if (text.includes("based on") || text.includes("according to") || text.includes("بر اساس") || text.includes("طبق")) { + return 0.95; + } + + return 0.8; + } + + private estimateTokens(text: string): number { + // Rough estimation: 1 token ≈ 4 characters for English text + // OpenAI uses similar tokenization to other models + return Math.ceil(text.length / 4); + } +} diff --git a/src/modules/chatbot/providers/websocket-auth.service.ts b/src/modules/chatbot/providers/websocket-auth.service.ts new file mode 100644 index 0000000..68595b5 --- /dev/null +++ b/src/modules/chatbot/providers/websocket-auth.service.ts @@ -0,0 +1,203 @@ +import { inject, injectable } from "inversify"; +import { Socket } from "socket.io"; + +import { WebSocketMessage } from "../../../common/enums/message.enum"; +import { jwtExpiredErr } from "../../../core/app/app.errors"; +import { Logger } from "../../../core/logging/logger"; +import { IOCTYPES } from "../../../IOC/ioc.types"; +import { TokenService } from "../../token/token.service"; +import { WEBSOCKET_EVENTS } from "../constants/chatbot.constants"; +import { AuthenticationResult, WebSocketErrorContext } from "../interfaces/websocket.interface"; + +/** + * Service responsible for WebSocket authentication logic + * Follows single responsibility principle and provides reusable authentication methods + */ +@injectable() +export class WebSocketAuthService { + private readonly logger: Logger; + + constructor(@inject(IOCTYPES.TokenService) private readonly tokenService: TokenService) { + this.logger = new Logger("WebSocketAuthService"); + } + + /** + * Extracts token from WebSocket client + */ + private extractTokenFromClient(client: Socket): string | null { + // Try to get token from query parameters + const token = client.handshake.query.token as string; + if (token) { + return token; + } + + // Try to get token from auth header + const authHeader = client.handshake.headers.authorization; + if (authHeader && typeof authHeader === "string") { + const parts = authHeader.split(" "); + if (parts.length === 2 && parts[0] === "Bearer") { + return parts[1]; + } + } + + return null; + } + + /** + * Authenticates a WebSocket client and returns the result + * + * @param client - Socket.IO client to authenticate + * @returns Promise with authentication result + */ + async authenticateClient(client: Socket): Promise { + try { + const token = this.extractTokenFromClient(client); + + if (!token) { + this.logAuthenticationAttempt(client, false, "No token provided"); + return { + success: false, + error: WebSocketMessage.AUTHENTICATION_REQUIRED, + }; + } + + const payload = this.tokenService.verifyToken(token); + + if (!payload?.sub) { + this.logAuthenticationAttempt(client, false, "Invalid token payload"); + return { + success: false, + error: WebSocketMessage.INVALID_TOKEN, + }; + } + + // Store user data in client (convert sub to id for compatibility) + client.data.user = { id: payload.sub, sub: payload.sub }; + + this.logAuthenticationAttempt(client, true, undefined, payload.sub); + + return { + success: true, + user: { id: payload.sub, sub: payload.sub }, + }; + } catch (error) { + const errorMessage = this.getAuthErrorMessage(error); + this.logAuthenticationAttempt(client, false, errorMessage); + + return { + success: false, + error: errorMessage, + }; + } + } + + /** + * Emits authentication failure event to client and disconnects + * + * @param client - Socket.IO client + * @param error - Error message to send + */ + handleAuthenticationFailure(client: Socket, error: string): void { + client.emit(WEBSOCKET_EVENTS.UNAUTHORIZED, { + message: error, + timestamp: new Date().toISOString(), + }); + + // Graceful disconnect with delay to ensure message is received + setTimeout(() => { + client.disconnect(true); + }, 100); + } + + /** + * Emits successful authentication event to client + * + * @param client - Socket.IO client + * @param user - Authenticated user payload + */ + emitAuthenticationSuccess(client: Socket, user: { id: string; sub: string }): void { + client.emit(WEBSOCKET_EVENTS.AUTHENTICATED, { + message: WebSocketMessage.AUTHENTICATED, + userId: user.id, + timestamp: new Date().toISOString(), + }); + } + + /** + * Creates error context object for logging and monitoring + * + * @param client - Socket.IO client + * @param event - Event name (optional) + * @param data - Event data (optional) + * @returns WebSocket error context + */ + createErrorContext(client: Socket, event?: WEBSOCKET_EVENTS, data?: any): WebSocketErrorContext { + return { + clientId: client.id, + userId: client.data?.user?.id, + sessionId: client.data?.sessionId, + event, + data, + timestamp: new Date(), + }; + } + + /** + * Validates if client is authenticated + * + * @param client - Socket.IO client to validate + * @returns True if client has valid user data + */ + isClientAuthenticated(client: Socket): boolean { + return !!client.data?.user?.id; + } + + /** + * Gets user ID from authenticated client + * + * @param client - Socket.IO client + * @returns User ID or undefined if not authenticated + */ + getUserId(client: Socket): string | undefined { + return client.data?.user?.id; + } + + /** + * Logs authentication attempts with structured data + */ + private logAuthenticationAttempt(client: Socket, success: boolean, error?: string, userId?: string): void { + const logData = { + clientId: client.id, + clientIP: client.handshake.address, + success, + userId, + error, + timestamp: new Date().toISOString(), + }; + + if (success) { + this.logger.info(`WebSocket authentication successful`, logData); + } else { + this.logger.warn(`WebSocket authentication failed`, logData); + } + } + + /** + * Maps JWT errors to user-friendly messages + */ + private getAuthErrorMessage(error: any): string { + if (error instanceof jwtExpiredErr || error?.name === "TokenExpiredError") { + return WebSocketMessage.TOKEN_EXPIRED; + } + + if (error?.name === "JsonWebTokenError") { + return WebSocketMessage.INVALID_TOKEN; + } + + if (error?.name === "NotBeforeError") { + return WebSocketMessage.INVALID_TOKEN; + } + + return WebSocketMessage.INVALID_TOKEN; + } +} diff --git a/src/modules/chatbot/utils/ulid.util.ts b/src/modules/chatbot/utils/ulid.util.ts new file mode 100644 index 0000000..b46a8e3 --- /dev/null +++ b/src/modules/chatbot/utils/ulid.util.ts @@ -0,0 +1,28 @@ +import { Request, Response } from "express"; +import * as ulidLib from "ulid"; + +const CHATBOT_ULID_COOKIE_NAME = "chatbot_session_id"; + +/** + * Get or create ULID for chatbot session from cookies + * @param req Express request object + * @param res Express response object + * @returns ULID string + */ +export function getOrCreateChatbotUlid(req: Request, res: Response): string { + let ulidValue = req.cookies?.[CHATBOT_ULID_COOKIE_NAME]; + + if (!ulidValue) { + ulidValue = ulidLib.ulid(); + // Set cookie with 1 year expiration (or adjust as needed) + res.cookie(CHATBOT_ULID_COOKIE_NAME, ulidValue, { + maxAge: 365 * 24 * 60 * 60 * 1000, // 1 year + httpOnly: true, + secure: process.env.NODE_ENV === "production", + sameSite: "lax", + }); + } + + return ulidValue; +} + diff --git a/src/server.ts b/src/server.ts index e0b438d..490ff36 100644 --- a/src/server.ts +++ b/src/server.ts @@ -11,6 +11,7 @@ import { Logger } from "./core/logging/logger"; import { connectMongo } from "./db/connection"; import { IOCTYPES } from "./IOC/ioc.types"; import { ChatGateway } from "./modules/chat/chat.gateway"; +import { ChatbotGateway } from "./modules/chatbot/chatbot.gateway"; import { StartWorker } from "./queues"; import { gracefulShutdown } from "./utils/shutdown.utils"; @@ -28,6 +29,9 @@ async function bootStrap() { const chatGateway = container.get(IOCTYPES.ChatGateway); const wsServer = chatGateway.initialize(server); + const chatbotGateway = container.get(IOCTYPES.ChatbotGateway); + const chatbotWsServer = chatbotGateway.initialize(server); + const worker = await StartWorker(); // @@ -36,9 +40,9 @@ async function bootStrap() { logger.info(`swagger documentation is serving on http://localhost:${PORT}/api-docs`); }); - process.on("SIGINT", () => gracefulShutdown("SIGINT", server, wsServer, worker[0], logger)); + process.on("SIGINT", () => gracefulShutdown("SIGINT", server, wsServer, worker[0], logger, chatbotWsServer)); - process.on("SIGTERM", () => gracefulShutdown("SIGTERM", server, wsServer, worker[0], logger)); + process.on("SIGTERM", () => gracefulShutdown("SIGTERM", server, wsServer, worker[0], logger, chatbotWsServer)); } catch (error) { logger.error("Error starting the server.", error); process.exit(1); diff --git a/src/utils/shutdown.utils.ts b/src/utils/shutdown.utils.ts index 51d4ff0..a100e51 100644 --- a/src/utils/shutdown.utils.ts +++ b/src/utils/shutdown.utils.ts @@ -5,16 +5,29 @@ import { Server as WSServer } from "socket.io"; import { Logger } from "../core/logging/logger"; -export const gracefulShutdown = async (signal: string, server: Server, wsServer: WSServer, worker: Worker, logger: Logger) => { +export const gracefulShutdown = async ( + signal: string, + server: Server, + wsServer: WSServer, + worker: Worker, + logger: Logger, + chatbotWsServer?: WSServer, +) => { try { logger.warn(`${signal} signal received: closing HTTP server and worker...`); logger.warn("close all socket connection.."); wsServer.disconnectSockets(); + if (chatbotWsServer) { + chatbotWsServer.disconnectSockets(); + } logger.warn("all socket connection disconnected"); logger.warn("Closing ws server..."); wsServer.close(); + if (chatbotWsServer) { + chatbotWsServer.close(); + } logger.warn("ws server closed"); server.close(async () => { diff --git a/test-chatbot-socket.html b/test-chatbot-socket.html new file mode 100644 index 0000000..8ad4fd1 --- /dev/null +++ b/test-chatbot-socket.html @@ -0,0 +1,293 @@ + + + + + + Chatbot Socket.IO Test + + + + +
+

Chatbot Socket.IO Test Client

+ +
+
+ +
+ +
+
+ +
+ +
+ + + + +
+ +
+
+ +
+ +
+
+ + + +
+ +
+

Event Log:

+
+
+
+ + + + + diff --git a/test-chatbot-socket.js b/test-chatbot-socket.js new file mode 100644 index 0000000..e7e3f5c --- /dev/null +++ b/test-chatbot-socket.js @@ -0,0 +1,99 @@ +// Test script for ChatbotGateway Socket.IO connection +// Run with: node test-chatbot-socket.js + +const io = require("socket.io-client"); + +// Connect to the chatbot gateway +// The path option matches the path configured in ChatbotGateway +const socket = io("http://localhost:4000", { + path: "/ws-chatbot", + transports: ["websocket"], + // Optional: Add token for authenticated connection + // query: { + // token: "your-jwt-token-here" + // } +}); + +socket.on("connect", () => { + console.log("✅ Connected to ChatbotGateway!"); + console.log("Socket ID:", socket.id); + + // Test: Create a session + console.log("\n📤 Creating session..."); + socket.emit("create_session"); +}); + +socket.on("authenticated", (data) => { + console.log("✅ Authenticated:", data); +}); + +socket.on("unauthorized", (data) => { + console.log("❌ Unauthorized:", data); +}); + +socket.on("session_created", (data) => { + console.log("✅ Session created:", JSON.stringify(data, null, 2)); + + // Test: Join the session + if (data.data && data.data.id) { + console.log("\n📤 Joining session:", data.data.id); + socket.emit("join_chat", data.data.id); + } +}); + +socket.on("chat_joined", (data) => { + console.log("✅ Joined chat:", JSON.stringify(data, null, 2)); + + // Test: Send a message + if (data.data && data.data.id) { + console.log("\n📤 Sending test message..."); + socket.emit("send_message", { + sessionId: data.data.id, + content: "Hello, chatbot!", + }); + } +}); + +socket.on("message_received", (data) => { + console.log("📨 Message received:", JSON.stringify(data, null, 2)); +}); + +socket.on("bot_response_start", (data) => { + console.log("🤖 Bot response started:", JSON.stringify(data, null, 2)); +}); + +socket.on("bot_response_chunk", (data) => { + process.stdout.write(data.data?.chunk || ""); +}); + +socket.on("bot_response_end", (data) => { + console.log("\n✅ Bot response completed:", JSON.stringify(data, null, 2)); +}); + +socket.on("bot_response", (data) => { + console.log("🤖 Bot response:", JSON.stringify(data, null, 2)); +}); + +socket.on("error", (data) => { + console.log("❌ Error:", JSON.stringify(data, null, 2)); +}); + +socket.on("session_error", (data) => { + console.log("❌ Session error:", JSON.stringify(data, null, 2)); +}); + +socket.on("message_error", (data) => { + console.log("❌ Message error:", JSON.stringify(data, null, 2)); +}); + +socket.on("disconnect", (reason) => { + console.log("❌ Disconnected:", reason); +}); + +socket.on("connect_error", (error) => { + console.log("❌ Connection error:", error.message); +}); + +// Keep the script running +process.stdin.resume(); +