Spaces:
Running
Running
File size: 102,539 Bytes
a81ae9b 746c6f2 a81ae9b 746c6f2 4fd4067 a81ae9b 54e7394 a81ae9b 746c6f2 a81ae9b 54e7394 a81ae9b 746c6f2 a81ae9b 746c6f2 a81ae9b 54e7394 bd7fab4 54e7394 bd7fab4 54e7394 a81ae9b 54e7394 746c6f2 a81ae9b 54e7394 a81ae9b af74e31 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 7ebbe37 7afbe42 a81ae9b af74e31 8da06be af74e31 8da06be a81ae9b af74e31 7afbe42 af74e31 a81ae9b af74e31 7afbe42 af74e31 a81ae9b 746c6f2 a81ae9b af74e31 a81ae9b 746c6f2 a81ae9b d9b7874 746c6f2 d9b7874 8a4684d a81ae9b d9b7874 a81ae9b 8a4684d a81ae9b 746c6f2 a81ae9b 746c6f2 a81ae9b 746c6f2 a81ae9b 54e7394 a81ae9b 24da3f2 a81ae9b | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 1279 1280 1281 1282 1283 1284 1285 1286 1287 1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499 1500 1501 1502 1503 1504 1505 1506 1507 1508 1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526 1527 1528 1529 1530 1531 1532 1533 1534 1535 1536 1537 1538 1539 1540 1541 1542 1543 1544 1545 1546 1547 1548 1549 1550 1551 1552 1553 1554 1555 1556 1557 1558 1559 1560 1561 1562 1563 1564 1565 1566 1567 1568 1569 1570 1571 1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596 1597 1598 1599 1600 1601 1602 1603 1604 1605 1606 1607 1608 1609 1610 1611 1612 1613 1614 1615 1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696 1697 1698 1699 1700 1701 1702 1703 1704 1705 1706 1707 1708 1709 1710 1711 1712 1713 1714 1715 1716 1717 1718 1719 1720 1721 1722 1723 1724 1725 1726 1727 1728 1729 1730 1731 1732 1733 1734 1735 1736 1737 1738 1739 1740 1741 1742 1743 1744 1745 1746 1747 1748 1749 1750 1751 1752 1753 1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765 1766 1767 1768 1769 1770 1771 1772 1773 1774 1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786 1787 1788 1789 1790 1791 1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802 1803 1804 1805 1806 1807 1808 1809 1810 1811 1812 1813 1814 1815 1816 1817 1818 1819 1820 1821 1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841 1842 1843 1844 1845 1846 1847 1848 1849 1850 1851 1852 1853 1854 1855 1856 1857 1858 1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909 1910 1911 1912 1913 1914 1915 1916 1917 1918 1919 1920 1921 1922 1923 1924 1925 1926 1927 1928 1929 1930 1931 1932 1933 1934 1935 1936 1937 1938 1939 1940 1941 1942 1943 1944 1945 1946 1947 1948 1949 1950 1951 1952 1953 1954 1955 1956 1957 1958 1959 1960 1961 1962 1963 1964 1965 1966 1967 1968 1969 1970 1971 1972 1973 1974 1975 1976 1977 1978 1979 1980 1981 1982 1983 1984 1985 1986 1987 1988 1989 1990 1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014 2015 2016 2017 2018 2019 2020 2021 2022 2023 2024 2025 2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040 2041 2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058 2059 2060 2061 2062 2063 2064 2065 2066 2067 2068 2069 2070 2071 2072 2073 2074 2075 2076 2077 2078 2079 2080 2081 2082 2083 2084 2085 2086 2087 2088 2089 2090 2091 2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113 2114 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 | import os, sys, json, secrets, logging, asyncio, re, time, threading, base64, tempfile, pathlib, hashlib
from html import escape
logging.basicConfig(level=logging.INFO, stream=sys.stdout)
logger = logging.getLogger("zalo-bot")
import requests
import gradio as gr
from fastapi import FastAPI, Request, Response
from starlette.responses import RedirectResponse
from huggingface_hub import HfApi, SpaceStage, hf_hub_download, upload_file
from datasets import Dataset, Features, Value
DEFAULT_BOT_TOKEN = os.getenv(
"DEFAULT_BOT_TOKEN",
"4179413508988279245:DmcFvOoFHHGiISQtmInHFchHwqfmAsaNWxhENixtvawrerrMGALunAbfhBvOzUcc",
)
WEBHOOK_SECRET = os.getenv("WEBHOOK_SECRET", "") or secrets.token_urlsafe(32)[:128]
SPACE_ID = os.getenv("SPACE_ID", "")
HF_TOKEN = os.getenv("HF_TOKEN", os.getenv("HF_API_TOKEN", ""))
if not HF_TOKEN:
try:
from huggingface_hub import get_token
HF_TOKEN = get_token() or ""
except Exception:
HF_TOKEN = ""
NAMESPACE = os.getenv("HF_NAMESPACE", "bep40")
MAIN_DATASET_ID = os.getenv("MAIN_DATASET_ID", f"{NAMESPACE}/zalo-products-all")
ZGR_SENDER_ID = "zgr-b7e1e71cf5701c2e4561"
OCR_MODEL_ID = os.getenv("OCR_MODEL_ID", "5CD-AI/Vintern-1B-v3_5")
if not HF_TOKEN:
logger.warning("[startup] HF_TOKEN not found — Space creation features will fail until HF_TOKEN secret is set")
logger.info("[startup] SPACE_ID=%s NAMESPACE=%s OCR_MODEL=%s", SPACE_ID or "(local)", NAMESPACE, OCR_MODEL_ID)
BOT_STATE = {
"bot_token": DEFAULT_BOT_TOKEN,
"webhook_url": "",
"webhook_secret": WEBHOOK_SECRET,
"connected": False,
"bot_info": {},
"logs": [],
"api_spaces": [],
"last_chat_id": "",
"last_sender_id": "",
}
def _load_proxy_spaces():
if not SPACE_ID or not HF_TOKEN:
return
try:
file_path = hf_hub_download(
repo_id=SPACE_ID, filename="proxy_spaces.json", repo_type="space", token=HF_TOKEN,
)
with open(file_path) as f:
data = json.load(f)
BOT_STATE["api_spaces"] = data.get("api_spaces", [])
logger.info("Loaded %d proxy spaces from disk", len(BOT_STATE["api_spaces"]))
except Exception as e:
logger.info("No persisted proxy data yet: %s", e)
def _save_proxy_spaces():
if not SPACE_ID or not HF_TOKEN:
return
try:
data = {"api_spaces": BOT_STATE.get("api_spaces", [])}
with tempfile.NamedTemporaryFile(mode="w", suffix=".json", delete=False) as tmp:
json.dump(data, tmp, indent=2, default=str)
tmp_path = tmp.name
upload_file(
path_or_fileobj=tmp_path,
path_in_repo="proxy_spaces.json",
repo_id=SPACE_ID,
repo_type="space",
token=HF_TOKEN,
commit_message="Update proxy spaces list",
)
pathlib.Path(tmp_path).unlink(missing_ok=True)
logger.info("Saved %d proxy spaces to disk", len(BOT_STATE.get("api_spaces", [])))
except Exception as e:
logger.error("Failed to save proxy spaces: %s", e)
_load_proxy_spaces()
def _save_chat_id(cid: str, sender_id: str):
if cid:
BOT_STATE["last_chat_id"] = cid
if sender_id:
BOT_STATE["last_sender_id"] = sender_id
class ZaloBotAPI:
BASE_URL = "https://bot-api.zaloplatforms.com"
def __init__(self, bt: str):
self.bt = bt
assert ":" in bt, "FAIL-FAST: sai dinh dang token"
self.api_base = f"{self.BASE_URL}/bot{bt}"
self.headers = {"Content-Type": "application/json"}
def get_me(self):
return requests.post(f"{self.api_base}/getMe", headers=self.headers, timeout=15).json()
def set_webhook(self, url: str, secret: str):
return requests.post(
f"{self.api_base}/setWebhook",
json={"url": url, "secret_token": secret},
headers=self.headers,
timeout=15,
).json()
def send_message(self, cid: str, text: str):
return requests.post(
f"{self.api_base}/sendMessage",
json={"chat_id": cid, "text": text, "parse_mode": "markdown"},
headers=self.headers,
timeout=15,
).json()
def get_webhook_url():
if SPACE_ID:
slug = SPACE_ID.replace("/", "-").replace("_", "-")
return f"https://{slug}.hf.space/webhooks"
port = os.getenv("PORT", "7860")
return f"http://localhost:{port}/webhooks"
def _extract_user_token(text: str):
m = re.search(r'HTTP\s*API\s*:\s*(\S+:\S+)', text, re.IGNORECASE)
if m:
return m.group(1).strip()
m = re.search(r'(\d+:[A-Za-z0-9_-]+)', text)
if m:
return m.group(1)
return None
def _safe_space_name(user_id: str) -> str:
return re.sub(r'[^a-zA-Z0-9-]', '', str(user_id))[:40]
def _get_user_dataset_id(user_id: str) -> str:
safe_id = _safe_space_name(user_id)
if not safe_id:
safe_id = "default"
return f"{NAMESPACE}/{safe_id}-zalo-data"
def _is_zgr_sender_local(sender_id):
sid = str(sender_id)
return ZGR_SENDER_ID in sid or sid == ZGR_SENDER_ID
def _ensure_main_dataset_schema():
"""Ensure the main dataset has proper schema (README + parquet) for viewer."""
if not HF_TOKEN or not MAIN_DATASET_ID:
return
try:
api = HfApi(token=HF_TOKEN)
# Check if dataset exists
try:
api.dataset_info(MAIN_DATASET_ID)
except Exception:
api.create_repo(repo_id=MAIN_DATASET_ID, repo_type="dataset", token=HF_TOKEN)
logger.info("Created main dataset: %s", MAIN_DATASET_ID)
# Upload README with schema
readme = """# Zalo Product Data - Main Dataset
Dữ liệu sản phẩm thu thập từ Zalo bot.
## Schema
| Column | Type | Description |
|--------|------|-------------|
| product_name | string | Tên sản phẩm |
| description | string | Nội dung mô tả |
| price | string | Giá sản phẩm |
| category | string | Chuyên mục |
| technical_specs | string | Thông số kỹ thuật |
| sender_id | string | ID người gửi |
| sender_name | string | Tên người gửi |
| chat_id | string | ID chat |
| timestamp | string | Thời gian ghi nhận |
| image | string | Đường link ảnh |
| text | string | Nội dung tin nhắn |
| message_type | string | Loại tin nhắn |
| is_zgr_group | string | Gửi từ nhóm ZGR |
## Usage
```python
from datasets import load_dataset
ds = load_dataset("bep40/zalo-products-all")
```
"""
with tempfile.NamedTemporaryFile(mode="w", suffix=".md", delete=False) as tmp:
tmp.write(readme)
tmp_path = tmp.name
try:
api.upload_file(
path_or_fileobj=tmp_path,
path_in_repo="README.md",
repo_id=MAIN_DATASET_ID,
repo_type="dataset",
token=HF_TOKEN,
commit_message="Update schema README",
)
logger.info("Updated main dataset README")
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
except Exception as e:
logger.error("Failed to ensure main dataset schema: %s", e)
def _append_to_main_dataset_parquet(new_records: list):
"""Merge new records into the main dataset as a parquet file (for viewer compatibility)."""
if not HF_TOKEN or not MAIN_DATASET_ID:
logger.warning("MAIN_DATASET_ID or HF_TOKEN not configured")
return
try:
api = HfApi(token=HF_TOKEN)
# Try to load existing parquet to merge
import pandas as pd
from datasets import Dataset
existing_df = None
try:
parquet_path = hf_hub_download(
repo_id=MAIN_DATASET_ID,
filename="data/train-00000-of-00001.parquet",
repo_type="dataset",
token=HF_TOKEN,
)
existing_ds = Dataset.from_parquet(parquet_path)
if len(existing_ds) > 1: # More than just placeholder
existing_df = existing_ds.to_pandas()
except Exception:
pass # No existing parquet, start fresh
# Prepare new dataframe
new_df = pd.DataFrame(new_records)
# Ensure all columns match
expected_cols = ["image", "product_name", "description", "price", "category",
"technical_specs", "sender_id", "sender_name", "chat_id",
"timestamp", "text", "message_type", "is_zgr_group"]
for col in expected_cols:
if col not in new_df.columns:
new_df[col] = ""
new_df = new_df[expected_cols]
# Merge with existing data (skip placeholder)
if existing_df is not None:
# Remove placeholder row
existing_df = existing_df[existing_df['sender_id'] != 'placeholder']
combined = pd.concat([existing_df, new_df], ignore_index=True)
else:
combined = new_df
# Convert to dataset with proper features
features = {col: Value("string") for col in expected_cols}
ds = Dataset.from_pandas(combined, preserve_index=False)
# Save as parquet
with tempfile.NamedTemporaryFile(suffix=".parquet", delete=False) as tmp:
ds.to_parquet(tmp.name)
tmp_path = tmp.name
try:
api.upload_file(
path_or_fileobj=tmp_path,
path_in_repo="data/train-00000-of-00001.parquet",
repo_id=MAIN_DATASET_ID,
repo_type="dataset",
token=HF_TOKEN,
commit_message=f"Added {len(new_records)} records",
)
logger.info("Merged %d records into parquet (total: %d)", len(new_records), len(combined))
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
except Exception as e:
logger.error("Failed to merge into parquet: %s", e)
def _ensure_user_dataset(user_id: str, token: str) -> tuple:
"""Ensure a user-specific dataset exists."""
if not HF_TOKEN:
raise RuntimeError("HF_TOKEN chưa được cấu hình.")
api = HfApi(token=HF_TOKEN)
dataset_id = _get_user_dataset_id(user_id)
created = False
try:
api.dataset_info(dataset_id)
logger.info("Dataset %s already exists", dataset_id)
except Exception:
with tempfile.TemporaryDirectory() as tmp:
readme_path = pathlib.Path(tmp, "README.md")
readme_path.write_text(
f"# Zalo Product Data — {user_id}\n\n"
f"Dữ liệu sản phẩm thu thập từ Zalo chat.\n\n"
f"## Cấu trúc (schema)\n"
f"| image | ảnh (binary/URL) | Hình ảnh sản phẩm |\n"
f"| product_name | text | Tên sản phẩm |\n"
f"| description | text | Nội dung mô tả |\n"
f"| price | number | Giá sản phẩm |\n"
f"| category | text | Chuyên mục |\n"
f"| sender_id | text | ID người gửi |\n"
f"| sender_name | text | Tên người gửi |\n"
f"| timestamp | text | Thời gian ghi nhận |\n"
)
try:
api.create_repo(repo_id=dataset_id, repo_type="dataset", exist_ok=True)
api.upload_file(
path_or_fileobj=str(readme_path),
path_in_repo="README.md",
repo_id=dataset_id,
repo_type="dataset",
token=HF_TOKEN,
commit_message="Initial dataset readme",
)
created = True
logger.info("Created dataset %s", dataset_id)
except Exception as e:
err_msg = str(e)
if "429" in err_msg or "rate limit" in err_msg:
raise RuntimeError("⚠️ Đã đạt giới hạn tạo repo (20/ngày). Vui lòng thử lại sau 24h.")
if "already exist" in err_msg or "conflict" in err_msg:
logger.info("Dataset %s already exists, skipping create", dataset_id)
else:
raise
return dataset_id, created
def _save_product_to_main_dataset(image_url, image_data_b64, description, price, category, sender_id, sender_name, product_name="", chat_id="", timestamp="", img_bytes=None, technical_specs=""):
if not HF_TOKEN or not MAIN_DATASET_ID:
logger.warning("MAIN_DATASET_ID or HF_TOKEN not configured")
return None
try:
api = HfApi(token=HF_TOKEN)
file_ts = timestamp or time.strftime("%Y%m%d_%H%M%S")
rec_ts = time.strftime("%Y-%m-%d %H:%M:%S")
safe_sender = _safe_space_name(sender_id) or "unknown"
img_filename = f"images/{file_ts}_{safe_sender}.jpg"
if img_bytes is None:
if image_data_b64:
try:
img_bytes = base64.b64decode(image_data_b64)
except Exception:
img_bytes = None
elif image_url:
try:
r = requests.get(image_url, timeout=15)
img_bytes = r.content
except Exception as e:
logger.error("Image download failed: %s", e)
uploaded_img = ""
if img_bytes:
with tempfile.NamedTemporaryFile(suffix=".jpg", delete=False) as tmp:
tmp.write(img_bytes)
tmp_path = tmp.name
try:
api.upload_file(
path_or_fileobj=tmp_path,
path_in_repo=img_filename,
repo_id=MAIN_DATASET_ID,
repo_type="dataset",
token=HF_TOKEN,
commit_message=f"Add product image from {sender_name}",
)
uploaded_img = img_filename
except Exception as e:
logger.error("Image upload to main dataset failed: %s", e)
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
record = {
"image": uploaded_img,
"product_name": str(product_name)[:200] if product_name else "",
"description": str(description)[:500] if description else "",
"price": str(price) if price else "",
"category": str(category) if category else "",
"technical_specs": str(technical_specs)[:500] if technical_specs else "",
"sender_id": str(sender_id),
"sender_name": str(sender_name),
"chat_id": str(chat_id),
"timestamp": rec_ts,
"text": "",
"message_type": "product",
"is_zgr_group": str(_is_zgr_sender_local(sender_id)),
}
_append_to_main_dataset_parquet([record])
logger.info("Saved product to main dataset: %s", MAIN_DATASET_ID)
return MAIN_DATASET_ID
except Exception as e:
logger.error("Failed to save to main dataset: %s", e)
return None
def extract_technical_specs_from_text(text):
"""Extract technical specifications from product text using keyword patterns."""
if not text:
return ""
tech_patterns = [
(r'Kích thước[^::]*[::]?\s*(.+?)(?:;|Chất liệu|Dòng sản phẩm|Bảo hành|Màu sắc|Khoang tủ|Chiều|$)', "Kích thước"),
(r'Chất liệu[^::]*[::]?\s*(.+?)(?:;|Kích thưỏi|Dòng sản phẩm|Bảo hành|Màu sắc|$)', "Chất liệu"),
(r'Dòng sản phẩm[^::]*[::]?\s*(.+?)(?:;|Kích thưỏi|Chất liệu|$)', "Dòng sản phẩm"),
(r'Chiều rộng tủ[^::]*[::]?\s*(.+?)(?:\n|$)', "Chiều rộng tủ"),
(r'Khoang tủ[^::]*[::]?\s*(.+?)(?:;|Kích thưỏi|Chất liệu|Chiều|$)', "Khoang tủ"),
(r'Bảo hành[^::]*[::]?\s*(.+?)(?:\n|$)', "Bảo hành"),
(r'Trọng lượng[^::]*[::]?\s*(.+?)(?:\n|$)', "Trọng lượng"),
(r'Màu sắc[^::]*[::]?\s*(.+?)(?:\n|$)', "Màu sắc"),
]
specs = []
seen = set()
for pattern, label in tech_patterns:
match = re.search(pattern, text, re.IGNORECASE | re.DOTALL)
if match:
value = match.group(1).strip()
if value and label.lower() not in seen:
value = value.rstrip(';').strip()
specs.append(f"{label}: {value}")
seen.add(label.lower())
return "\n".join(specs)[:500] if specs else ""
def _save_text_message_to_dataset(text, description, price, category, sender_id, sender_name, chat_id):
"""Save a text-only message to the main Zalo products dataset."""
try:
tech_specs = extract_technical_specs_from_text(description or text)
rec_ts = time.strftime("%Y-%m-%d %H:%M:%S")
record = {
"image": "",
"product_name": "",
"description": str(description or text)[:500] if description or text else "",
"price": str(price) if price else "",
"category": str(category) if category else "",
"technical_specs": str(tech_specs)[:500] if tech_specs else "",
"sender_id": str(sender_id),
"sender_name": str(sender_name),
"chat_id": str(chat_id),
"timestamp": rec_ts,
"text": str(text)[:1000] if text else "",
"message_type": "text",
"is_zgr_group": str(_is_zgr_sender_local(sender_id)),
}
_append_to_main_dataset_parquet([record])
logger.info("Text message saved to dataset: %s", MAIN_DATASET_ID)
return MAIN_DATASET_ID
except Exception as e:
logger.error("Failed to save text to dataset: %s", e)
return None
def _ocr_extract_text(image_bytes):
"""Use HF Inference API with Vietnamese OCR model (Vintern-1B) to extract text from image."""
try:
from huggingface_hub import InferenceClient
client = InferenceClient(model=OCR_MODEL_ID, token=HF_TOKEN)
# Convert bytes to PIL image
from PIL import Image
import io
img = Image.open(io.BytesIO(image_bytes)).convert("RGB")
prompt = "<image>\nTrích xuất toàn bộ văn bản trong hình ảnh và trả về dưới dạng markdown."
result = client.chat_completion(
messages=[{"role": "user", "content": [{"type": "text", "text": prompt}, {"type": "image_url", "image_url": img}]}],
max_tokens=2048,
)
ocr_text = result.choices[0].message.content.strip()
logger.info("OCR extracted %d chars of text", len(ocr_text))
return ocr_text
except Exception as e:
logger.error("OCR extraction failed: %s", e)
return ""
def _parse_ocr_product_info(ocr_text):
"""Parse OCR-extracted text to extract product fields."""
product_name, description, price, category = "", "", "", ""
# Try to extract price (Vietnamese đồng format: số, số, hoặc số vnđ)
price_match = re.search(r'(\d{1,3}(?:[.,]\d{3})*(?:[.,]\d{2,3})?(?:\s*(?:đ|vnđ|VND|dong))?)', ocr_text, re.IGNORECASE)
if price_match:
price = price_match.group(1)
# Try to extract category from common keywords
for kw in ["Chuyên mục", "Danh mục", "Loại", "Category"]:
m = re.search(kw + r'[:\s]*([^\n]+)', ocr_text, re.IGNORECASE)
if m:
category = m.group(1).strip()
break
# Try to extract product name from common keywords
for kw in ["Tên sp", "Tên sản phẩm", "Product", "Tên hàng"]:
m = re.search(kw + r'[:\s]*([^\n]+)', ocr_text, re.IGNORECASE)
if m:
product_name = m.group(1).strip()
break
# Use remaining text as description
desc_text = ocr_text
for kw in ["Chuyên mục", "Danh mục", "Loại", "Category", "Tên sp", "Tên sản phẩm", "Product", "Tên hàng", "Giá", "gia", "Price"]:
desc_text = re.sub(kw + r'[:\s]*[^\n]+', '', desc_text, flags=re.IGNORECASE)
description = desc_text.strip()[:500] if desc_text.strip() else ""
return product_name, description, price, category
def _create_api_proxy_space(token, user_id, sender_display):
safe_id = _safe_space_name(user_id)
token_suffix = token.split(":")[-1][:12] if ":" in token else re.sub(r'\W', '', token[:12])
unique_key = safe_id if safe_id else "u" + token_suffix
space_name = f"zalo-proxy-{unique_key}"
repo_id = f"{NAMESPACE}/{space_name}"
dataset_id = f"{NAMESPACE}/{unique_key}-zalo-data"
logger.info("Creating proxy space %s for user %s", repo_id, sender_display)
if not HF_TOKEN:
raise RuntimeError("HF_TOKEN secret chưa được cấu hình cho Space.")
api = HfApi(token=HF_TOKEN)
dockerfile = """FROM python:3.12-slim
RUN useradd -m -u 1000 user
USER user
ENV HOME=/home/user
ENV PATH=/home/user/.local/bin:$PATH
WORKDIR $HOME/app
COPY --chown=user requirements.txt .
RUN pip install --user --no-cache-dir -r requirements.txt
COPY --chown=user app.py .
EXPOSE 7860
CMD ["uvicorn", "app:app", "--host", "0.0.0.0", "--port", "7860"]
"""
requirements = "fastapi>=0.111.0\nuvicorn[standard]>=0.30.0\nrequests>=2.32.0\nhuggingface_hub>=0.30.0\n"
app_py = '''import os, json, requests, time, re, base64, tempfile, pathlib, sys
from html import escape as _escape
class _StderrLogger:
def __init__(self):
self._log = []
def write(self, s):
if s.strip():
self._log.append(s)
sys.__stderr__.write(s)
def flush(self): pass
sys.stderr = _StderrLogger()
from fastapi import FastAPI, Request, Response
app = FastAPI(title="Zalo Proxy Space")
BOT_TOKEN = "''' + token + '''"
TARGET_API = "https://bot-api.zaloplatforms.com"
PROXY_NAME = "''' + sender_display + '''"
HF_TOKEN = os.getenv("HF_TOKEN", "")
DATASET_ID = "''' + dataset_id + '''"
MAIN_DATASET_ID = "''' + MAIN_DATASET_ID + '''"
MAIN_SPACE_URL = "''' + SPACE_ID.replace("/", "-") + '''.hf.space"
ZGR_SENDER_ID = "zgr-b7e1e71cf5701c2e4561"
_logs = []
def _send(cid, text):
headers = {"Content-Type": "application/json"}
url = TARGET_API + "/bot" + BOT_TOKEN + "/sendMessage"
return requests.post(url, json={"chat_id": cid, "text": text}, headers=headers)
def _safe_name(name):
return re.sub(r'[^a-zA-Z0-9]', '_', str(name))[:30]
def _is_zgr_sender(sender_id):
sid = str(sender_id)
return ZGR_SENDER_ID in sid or sid == ZGR_SENDER_ID
def _log(event, sender_id, chat_id, text, sender_name="", chat_type=""):
entry = {"event": str(event), "sender_id": str(sender_id), "sender_name": str(sender_name), "chat_id": str(chat_id), "chat_type": str(chat_type), "text": str(text)[:500], "is_zgr": _is_zgr_sender(sender_id), "time": time.strftime("%Y-%m-%d %H:%M:%S")}
_logs.append(entry)
print("[WEBHOOK] event=" + str(event) + " sender=" + str(sender_id) + " chat=" + str(chat_id) + " is_zgr=" + str(entry["is_zgr"]) + " text=" + str(text)[:100], flush=True)
if len(_logs) > 200:
del _logs[:100]
_log("startup", "system", "SYSTEM", "Proxy space initialized. PROXY_NAME=" + PROXY_NAME)
def _save_to_main_dataset(image_url, image_data_b64, description, price, category, sender_id, sender_name, product_name="", chat_id="", text_message=None):
_log("main_dataset_save_start", sender_id, chat_id, "product_name=" + str(product_name) + " price=" + str(price))
if not HF_TOKEN or not MAIN_DATASET_ID:
_log("main_dataset_skip", sender_id, chat_id, "HF_TOKEN or MAIN_DATASET_ID missing")
return None
try:
from huggingface_hub import HfApi
api = HfApi(token=HF_TOKEN)
file_ts = time.strftime("%Y%m%d_%H%M%S")
rec_ts = time.strftime("%Y-%m-%d %H:%M:%S")
safe_sender = _safe_name(sender_id) or "unknown"
if text_message:
img_filename = ""
meta_filename = "data/" + file_ts + "_" + safe_sender + "_text.json"
else:
img_filename = "images/" + file_ts + "_" + safe_sender + ".jpg"
meta_filename = "data/" + file_ts + "_" + safe_sender + ".json"
img_bytes = None
if image_data_b64:
try:
img_bytes = base64.b64decode(image_data_b64)
except Exception:
img_bytes = None
elif image_url:
try:
r = requests.get(image_url, timeout=15)
img_bytes = r.content
_log("image_downloaded_from_url", sender_id, chat_id, image_url[:100])
except Exception as e:
_log("image_download_fail", sender_id, chat_id, str(e))
uploaded_img = None
if img_bytes:
with tempfile.NamedTemporaryFile(suffix=".jpg", delete=False) as tmp:
tmp.write(img_bytes)
tmp_path = tmp.name
try:
api.upload_file(path_or_fileobj=tmp_path, path_in_repo=img_filename, repo_id=MAIN_DATASET_ID, repo_type="dataset", token=HF_TOKEN, commit_message="Add product image from " + sender_name)
uploaded_img = img_filename
_log("image_uploaded", sender_id, chat_id, img_filename)
except Exception as e:
_log("image_upload_fail", sender_id, chat_id, str(e))
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
if text_message:
record = {"image": "", "product_name": str(product_name)[:200] if product_name else "", "category": str(category) if category else "", "description": str(description)[:500] if description else str(text_message)[:500], "price": str(price) if price else "", "sender_id": str(sender_id), "sender_name": str(sender_name), "chat_id": str(chat_id), "is_zgr_group": _is_zgr_sender(sender_id), "timestamp": rec_ts, "text": str(text_message)[:1000], "message_type": "text"}
else:
record = {"image": uploaded_img, "product_name": str(product_name)[:200] if product_name else "", "category": str(category) if category else "", "description": str(description)[:500] if description else "", "price": str(price) if price else "", "sender_id": str(sender_id), "sender_name": str(sender_name), "chat_id": str(chat_id), "is_zgr_group": _is_zgr_sender(sender_id), "timestamp": rec_ts}
with tempfile.NamedTemporaryFile(mode="w", suffix=".json", delete=False) as tmp:
json.dump(record, tmp, indent=2, ensure_ascii=False)
tmp_path = tmp.name
try:
api.upload_file(path_or_fileobj=tmp_path, path_in_repo=meta_filename, repo_id=MAIN_DATASET_ID, repo_type="dataset", token=HF_TOKEN, commit_message="Add product metadata from " + sender_name)
_log("dataset_save_to_main", sender_id, chat_id, "OK")
except Exception as e:
_log("dataset_save_fail", sender_id, chat_id, str(e))
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
return MAIN_DATASET_ID
except Exception as e:
_log("main_dataset_error", sender_id, chat_id, str(e))
return None
def _save_to_dataset(image_url, image_data_b64, description, price, category, sender_id, sender_name):
if not HF_TOKEN or not DATASET_ID:
_log("dataset_skip", sender_id, "N/A", "HF_TOKEN or DATASET_ID missing")
return None
try:
from huggingface_hub import HfApi
api = HfApi(token=HF_TOKEN)
file_ts = time.strftime("%Y%m%d_%H%M%S")
rec_ts = time.strftime("%Y-%m-%d %H:%M:%S")
safe_sender = _safe_name(sender_id) or "unknown"
img_filename = "images/" + file_ts + "_" + safe_sender + ".jpg"
meta_filename = "data/" + file_ts + "_" + safe_sender + ".json"
img_bytes = None
if image_data_b64:
try:
img_bytes = base64.b64decode(image_data_b64)
except Exception:
img_bytes = None
elif image_url:
try:
r = requests.get(image_url, timeout=15)
img_bytes = r.content
except Exception as e:
_log("image_download_fail", sender_id, "N/A", str(e))
uploaded_img = None
if img_bytes:
with tempfile.NamedTemporaryFile(suffix=".jpg", delete=False) as tmp:
tmp.write(img_bytes)
tmp_path = tmp.name
try:
api.upload_file(path_or_fileobj=tmp_path, path_in_repo=img_filename, repo_id=DATASET_ID, repo_type="dataset", token=HF_TOKEN, commit_message="Add product image from " + sender_name)
uploaded_img = img_filename
except Exception as e:
_log("image_upload_fail", sender_id, "N/A", str(e))
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
record = {"image": uploaded_img, "product_name": "", "category": str(category) if category else "", "description": str(description)[:500] if description else "", "price": str(price) if price else "", "sender_id": str(sender_id), "sender_name": str(sender_name), "timestamp": rec_ts}
with tempfile.NamedTemporaryFile(mode="w", suffix=".json", delete=False) as tmp:
json.dump(record, tmp, indent=2, ensure_ascii=False)
tmp_path = tmp.name
try:
api.upload_file(path_or_fileobj=tmp_path, path_in_repo=meta_filename, repo_id=DATASET_ID, repo_type="dataset", token=HF_TOKEN, commit_message="Add product metadata from " + sender_name)
finally:
pathlib.Path(tmp_path).unlink(missing_ok=True)
_log("dataset_saved", sender_id, "N/A", "Saved to " + DATASET_ID)
return DATASET_ID
except Exception as e:
_log("dataset_error", sender_id, "N/A", str(e))
return None
APP = FastAPI(title="Zalo Proxy Space")
@app.get("/")
async def root():
return {"status": "ok"}
@app.get("/health")
async def health():
return {"status": "ok", "dataset": DATASET_ID, "zgr_sender": ZGR_SENDER_ID, "is_zgr": _is_zgr_sender(ZGR_SENDER_ID)}
@app.get("/webhooks")
async def webhooks_get():
return Response(content=json.dumps({"message": "Success"}), media_type="application/json", status_code=200)
@app.post("/webhooks")
async def webhooks(request: Request):
body = await request.body()
body_str = body.decode("utf-8") if body else ""
_log("webhook_received", "N/A", "N/A", "Body length: " + str(len(body_str)))
try:
data = json.loads(body_str)
except Exception as e:
_log("parse_error", "N/A", "N/A", "Bad JSON: " + str(e) + " | body=" + body_str[:200])
return Response(content=json.dumps({"message": "Bad JSON"}), media_type="application/json", status_code=400)
result = data.get("result", data)
event = result.get("event_name", "unknown")
msg = result.get("message", {}); sender = msg.get("from", {}); chat = msg.get("chat", {})
text = msg.get("text", "")
sender_id = str(sender.get("id", "")); sender_name = sender.get("display_name") or sender.get("name") or sender_id
chat_id = str(chat.get("id", "")); chat_type = str(chat.get("chat_type", ""))
attachments = msg.get("attachment", {})
image_url = ""; image_data_b64 = ""
if attachments:
payload = attachments.get("payload", {})
if isinstance(payload, str):
try: payload = json.loads(payload)
except Exception: payload = {}
image_url = payload.get("url", "") or msg.get("photo", "") or msg.get("photo_url", "") or msg.get("image_url", "")
image_data_b64 = payload.get("data", "") or msg.get("image", "")
else:
image_url = msg.get("photo", "") or msg.get("photo_url", "") or msg.get("image_url", "")
image_data_b64 = msg.get("image", "")
is_zgr = _is_zgr_sender(sender_id)
_log(event, sender_id, chat_id, text, sender_name, chat_type)
_log("debug_info", sender_id, chat_id, "chat_type=" + str(chat_type) + " is_zgr=" + str(is_zgr) + " sender_id=" + str(sender_id) + " sender_name=" + str(sender_name) + " text_len=" + str(len(text)) + " has_attachment=" + str(bool(attachments)) + " image_url=" + str(image_url[:100]))
if event == "message.text.received" and chat_id:
description, price, category, product_name = "", "", "", ""
desc_match = re.search(r'(?:mo ta|description|desc|mota)[:\\\\s]*([^|\\n]+)', text, re.IGNORECASE)
price_match = re.search(r'(?:gia|price|don gia|donggia)[:\\\\s]*([\\d,.]+)', text, re.IGNORECASE)
cat_match = re.search(r'(?:chuyen muc|category|danh muc|loai)[:\\s]*([^|\\n]+?)(?:$|\\n)', text, re.IGNORECASE)
name_match = re.search(r'(?:ten sp|ten san pham|product name|name)[:\\s]*([^|\\n]+)', text, re.IGNORECASE)
if name_match:
product_name = name_match.group(1).strip()
if desc_match: description = desc_match.group(1).strip()
if price_match: price = price_match.group(1).strip()
if cat_match: category = cat_match.group(1).strip()
if image_url or image_data_b64 or product_name or description or price or category:
dataset_id = _save_to_dataset(image_url=image_url, image_data_b64=image_data_b64, description=description or text[:200], price=price, category=category, sender_id=sender_id, sender_name=sender_name)
_save_to_main_dataset(image_url=image_url, image_data_b64=image_data_b64, description=description or text[:200], price=price, category=category, sender_id=sender_id, sender_name=sender_name, product_name=product_name, chat_id=chat_id)
if dataset_id:
reply = "GOT IT! Product saved!"
else:
reply = "GOT IT! Saved to main dataset!"
elif is_zgr and text:
_save_to_main_dataset(image_url=image_url, image_data_b64=image_data_b64, description=text[:200], price=price, category=category, sender_id=sender_id, sender_name=sender_name, product_name=product_name, chat_id=chat_id, text_message=text)
reply = "👋 Xin chào " + str(sender_name) + " (Zalo ID: " + str(sender_id) + ")!\n\n" + HELP_INSTRUCTIONS
else:
reply = "Hi! Send image + product info to save."
try:
_send(chat_id, reply)
except Exception as e:
_log("send_reply_fail", sender_id, chat_id, str(e))
elif event == "message.image.received" and chat_id:
description, price, category, product_name = "", "", "", ""
desc_match = re.search(r'(?:mo ta|description|desc|mota)[:\\\\s]*([^|\\n]+)', text, re.IGNORECASE)
price_match = re.search(r'(?:gia|price|don gia|donggia)[:\\\\s]*([\\d,.]+)', text, re.IGNORECASE)
cat_match = re.search(r'(?:chuyen muc|category|danh muc|loai)[:\\s]*([^|\\n]+?)(?:$|\\n)', text, re.IGNORECASE)
name_match = re.search(r'(?:ten sp|ten san pham|product name|name)[:\\s]*([^|\\n]+)', text, re.IGNORECASE)
if name_match: product_name = name_match.group(1).strip()
if desc_match: description = desc_match.group(1).strip()
if price_match: price = price_match.group(1).strip()
if cat_match: category = cat_match.group(1).strip()
log_text = "photo_url=" + str(image_url[:100]) if image_url else "No photo_url in message"
if text: log_text += " | caption=" + str(text[:200])
if image_url:
dataset_id = _save_to_dataset(image_url=image_url, image_data_b64=image_data_b64, description=description or text[:200], price=price, category=category, sender_id=sender_id, sender_name=sender_name)
saved = _save_to_main_dataset(image_url=image_url, image_data_b64=image_data_b64, description=description or text[:200], price=price, category=category, sender_id=sender_id, sender_name=sender_name, product_name=product_name, chat_id=chat_id)
if dataset_id:
reply = "GOT IT! Product saved!"
else:
reply = "GOT IT! Saved to main dataset!"
_log("image_saved", sender_id, chat_id, log_text)
try: _send(chat_id, reply)
except Exception as e: _log("send_reply_fail", sender_id, chat_id, str(e))
else:
_log("image_no_url", sender_id, chat_id, log_text)
return Response(content=json.dumps({"message": "Success"}), media_type="application/json", status_code=200)
@app.get("/logs")
async def proxy_logs():
rows = ""
for log in reversed(_logs[-100:]):
is_zgr = _is_zgr_sender(log.get("sender_id", ""))
bg = "#e8f5e9" if is_zgr else "#ffffff"
zgr_badge = "[ZGR] " if is_zgr else ""
chat_type_val = _escape(str(log.get("chat_type", "")))
rows += "<div style='margin:6px 0;padding:8px;background:" + bg + ";border-radius:4px;border-left:3px solid #4CAF50'><b>" + zgr_badge + "[" + _escape(str(log["event"])) + "]</b> " + _escape(str(log.get("sender_name",""))) + " ID:<code>" + _escape(str(log.get("sender_id",""))) + "</code> chat:<code>" + _escape(str(log.get("chat_id",""))) + "</code> type:[" + chat_type_val + "]<br><span style='font-family:monospace;font-size:12px;color:#333'>" + _escape(str(log.get("text",""))[:300]) + "</span><br><small style='color:#999'>" + _escape(str(log.get("time",""))) + "</small></div>"
html_content = "<!DOCTYPE html><html><head><title>Proxy Logs</title><meta http-equiv='refresh' content='5'><style>body{font-family:Arial,sans-serif;max-width:1000px;margin:0 auto;padding:16px;}h1{color:#1a73e8;}.log-c{max-height:600px;overflow-y:auto;background:#fff;border-radius:8px;padding:8px;}</style></head><body><h1>Proxy Logs - " + _escape(PROXY_NAME) + "</h1><div class='log-c'>" + rows + "</div></body></html>"
return Response(content=html_content, media_type="text/html")
'''
readme = """---
title: Zalo Proxy
colorFrom: blue
colorTo: purple
sdk: docker
app_port: 7860
---
Zalo Webhook Proxy Space
"""
with tempfile.TemporaryDirectory() as tmp:
pathlib.Path(tmp, "Dockerfile").write_text(dockerfile)
pathlib.Path(tmp, "app.py").write_text(app_py)
pathlib.Path(tmp, "requirements.txt").write_text(requirements)
pathlib.Path(tmp, "README.md").write_text(readme)
try:
existing = api.space_info(repo_id=repo_id)
logger.info("Proxy space %s already exists (stage=%s)", repo_id, existing.stage)
except Exception:
try:
api.create_repo(repo_id=repo_id, repo_type="space", space_sdk="docker", exist_ok=True)
logger.info("Created new proxy space repo: %s", repo_id)
except Exception as e:
err_msg = str(e).lower()
if "429" in err_msg or "rate limit" in err_msg:
raise RuntimeError("Rate limit. Try again later.")
if "already exist" in err_msg or "conflict" in err_msg:
logger.info("Proxy space %s already exists, skipping create_repo", repo_id)
else:
raise
try:
api.upload_folder(folder_path=tmp, repo_id=repo_id, repo_type="space", commit_message="Initial proxy space")
except Exception as e:
err_msg = str(e)
if "404" in err_msg or "Repository Not Found" in err_msg:
raise RuntimeError("Proxy repo error")
raise
try:
api.wait_for_space(repo_id=repo_id, expected_stage=SpaceStage.RUNNING, timeout=180)
status = "RUNNING"
except Exception as e:
logger.warning("wait_for_space timeout: %s", e)
try:
rt = api.get_space_runtime(repo_id=repo_id)
status = str(rt.stage)
except Exception:
status = "UNKNOWN"
proxy_url = "https://" + repo_id.replace('/', '-') + ".hf.space/webhooks"
logger.info("Proxy space ready: %s -> %s (status=%s)", repo_id, proxy_url, status)
return repo_id, proxy_url, status
def _set_user_webhook(user_token, proxy_url, secret):
try:
zapi = ZaloBotAPI(user_token)
result = zapi.set_webhook(proxy_url, secret)
logger.info("setWebhook result: %s", result)
return result
except Exception as e:
logger.error("setWebhook failed: %s", e)
return {"ok": False, "message": str(e)}
def connect_bot(token: str):
if not token:
return "Nhap Bot Token"
BOT_STATE["bot_token"] = token
BOT_STATE["connected"] = False
BOT_STATE["bot_info"] = {}
try:
zapi = ZaloBotAPI(token)
except AssertionError as e:
return "Token sai: " + str(e)
try:
me = zapi.get_me()
if not me.get("ok"):
return "That bai: " + str(me.get('message', ''))
BOT_STATE["bot_info"] = me.get("result", {})
except Exception as e:
return "Loi: " + str(e)
wh = get_webhook_url()
sc = BOT_STATE["webhook_secret"]
try:
sw = zapi.set_webhook(wh, sc)
if sw.get("ok"):
BOT_STATE["webhook_url"] = wh
BOT_STATE["connected"] = True
return "Ket noi thanh cong! Webhook: " + wh
return "setWebhook that bai: " + str(sw.get('message',''))
except Exception as e:
return "Loi: " + str(e)
def send_msg(cid: str, text: str):
if not BOT_STATE["connected"]:
return "Chua ket noi bot."
if not cid or not text:
return "Nhap Chat ID va Noi dung"
try:
result = ZaloBotAPI(BOT_STATE["bot_token"]).send_message(cid, text)
return json.dumps(result, indent=2, ensure_ascii=False)
except Exception as e:
return "Loi: " + str(e)
def get_botinfo():
if BOT_STATE["bot_info"]:
info = BOT_STATE["bot_info"]
lines = ["Ten bot: " + str(info.get('name', '?')), "ID: " + str(info.get('id', ''))]
if BOT_STATE.get("webhook_url"):
lines.append("Webhook: " + BOT_STATE["webhook_url"])
lines.append("Ket noi: " + str(BOT_STATE.get("connected", False)))
return "\n".join(lines)
return "Chua ket noi"
def get_events():
if not BOT_STATE["logs"]:
return "Chua co su kien"
lines = []
for i, l in enumerate(BOT_STATE["logs"][-20:][::-1], 1):
is_zgr = _is_zgr_sender_local(l.get("sender_id", ""))
zgr_tag = " [ZGR]" if is_zgr else ""
lines.append("{}. [{}] {} Zalo:{} ID:{} chat:{} | {}".format(
i, l.get('event',''), zgr_tag, l.get('sender_name',''), l.get('sender_id',''), l.get('chat_id',''), str(l.get('text','')[:50])
))
return "\n".join(lines)
def get_proxy_spaces():
spaces = BOT_STATE.get("api_spaces", [])
if not spaces:
return "Chua co proxy space nao"
lines = []
for i, s in enumerate(spaces[-10:][::-1], 1):
lines.append("{}. {} ID:{} Space:{} Webhook:{} Status:{}".format(
i, s.get('sender_name',''), s.get('user_id',''), s.get('repo_id',''), s.get('proxy_url',''), s.get('status','')
))
return "\n".join(lines)
HELP_INSTRUCTIONS = (
"🎓 HƯỚNG DẪN CẤU HÌNH ZALO BOT CHI TIẾT\n\n"
"1️⃣ Cách đặt tên Zalobot (QUAN TRỌNG):\n"
"• Tên bot không được chứa 'Zalo' hoặc 'bot'\n"
"• Ví dụ đúng: Shop, ChămSóc, HỗTrợ247, CSKH-TựĐộng ✅\n"
"• Ví dụ sai: Zalo Support, ShopBot, ZaloBot ❌\n\n"
"2️⃣ Cách lấy HTTP API:\n"
"• Truy cập https://zalo.me/s/botcreator\n"
"• Chọn bot → Cài đặt → API/HTTP API\n"
"• Copy Bot token: 4179413508988279245:XXXXXXXXXXXXXXXXXXXXXX\n\n"
"3️⃣ Cách dùng:\n"
"• Gửi HTTP API: <bot_token> để tạo proxy tự động\n"
"• Gửi ảnh + mô tả sản phẩm để lưu vào dataset\n"
"• Mọi tin nhắn trong nhóm sẽ được lưu tự động"
)
async def handle_webhook(request: Request):
body = await request.body()
body_str = body.decode("utf-8") if body else ""
logger.info("webhook received: body=%s", body_str[:500])
try:
data = json.loads(body_str)
except Exception as e:
logger.error("JSON parse error: %s", e)
BOT_STATE["logs"].append({"event": "parse_error", "sender_id": "N/A", "chat_id": "N/A", "sender_name": "N/A", "chat_type": "N/A", "text": body_str[:200]})
return Response(content=json.dumps({"message": "Bad JSON", "error": str(e)}), media_type="application/json", status_code=400)
result = data.get("result", data)
event = result.get("event_name", "unknown")
msg = result.get("message", {}); sender = msg.get("from", {}); chat = msg.get("chat", {})
text = msg.get("text", "")
sender_id = str(sender.get("id", "")); sender_name = sender.get("display_name") or sender.get("name") or sender_id
chat_id = str(chat.get("id", "")); chat_type = str(chat.get("chat_type", ""))
BOT_STATE["logs"].append({
"event": str(event), "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type, "text": str(text)[:500],
"is_zgr": _is_zgr_sender_local(sender_id),
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
if len(BOT_STATE["logs"]) > 100:
BOT_STATE["logs"] = BOT_STATE["logs"][-100:]
logger.info("EVENT=%s SENDER_ID=%s CHAT_ID=%s SENDER_NAME=%s", event, sender_id, chat_id, sender_name)
if event == "message.text.received":
cid = chat.get("id") or sender.get("id") or ""
_save_chat_id(cid, sender_id)
user_token = _extract_user_token(text)
if user_token:
zapi = ZaloBotAPI(BOT_STATE["bot_token"])
try:
zapi.send_message(cid, "Processing your HTTP API...")
except Exception:
pass
def _create_and_setup():
try:
_r, proxy_url, status = _create_api_proxy_space(user_token, sender_id, sender_name)
dataset_id = None
try:
dataset_id, _ = _ensure_user_dataset(sender_id, user_token)
except Exception as de:
logger.error("Dataset setup failed: %s", de)
BOT_STATE["api_spaces"].append({
"repo_id": _r, "proxy_url": proxy_url, "user_token": user_token,
"user_id": sender_id, "sender_name": sender_name, "status": status, "dataset_id": dataset_id,
})
_save_proxy_spaces()
sw = _set_user_webhook(user_token, proxy_url, BOT_STATE["webhook_secret"])
user_bot_id = user_token.split(":")[0] if ":" in user_token else ""
user_bot_link = "https://zalo.me/s/" + user_bot_id if user_bot_id else "https://zalo.me/s/botcreator"
dataset_url = "https://huggingface.co/datasets/" + NAMESPACE + "/" + _safe_space_name(sender_id or "user") + "-zalo-data" if sender_id else ""
instructions = "BOT OK! Proxy: " + proxy_url + "\\nDataset: " + (dataset_url if dataset_url else "N/A") + "\\nManage bot at: " + user_bot_link
try:
zapi.send_message(cid, instructions)
except Exception:
pass
BOT_STATE["logs"].append({
"event": "proxy_created", "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type,
"text": "Proxy created: " + proxy_url,
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
except Exception as e:
logger.error("Failed: %s", e)
try:
ZaloBotAPI(BOT_STATE["bot_token"]).send_message(cid, "Loi tao proxy: " + str(e))
except Exception:
pass
threading.Thread(target=_create_and_setup, daemon=True).start()
return Response(content=json.dumps({"message": "Processing", "proxy_url": "pending"}), media_type="application/json")
# ─── Regular message ───
if cid:
zapi = ZaloBotAPI(BOT_STATE["bot_token"])
product_name, description, price, category = "", "", "", ""
image_url, image_data_b64 = "", ""
attachments = msg.get("attachment", {})
if attachments:
payload = attachments.get("payload", {})
if isinstance(payload, str):
try:
payload = json.loads(payload)
except Exception:
payload = {}
image_url = payload.get("url", "") or msg.get("photo", "") or msg.get("photo_url", "") or msg.get("image_url", "")
image_data_b64 = payload.get("data", "") or msg.get("image", "")
else:
image_url = msg.get("photo", "") or msg.get("photo_url", "") or msg.get("image_url", "")
image_data_b64 = msg.get("image", "")
name_match = re.search(r'(?:ten sp|ten san pham|product name|name)[:\s]*([^|\n]+)', text, re.IGNORECASE)
desc_match = re.search(r'(?:mo ta|description|desc|mota)[:\\s]*([^|\n]+)', text, re.IGNORECASE)
price_match = re.search(r'(?:gia|price|don gia|donggia)[:\\s]*([\\d,.]+)', text, re.IGNORECASE)
cat_match = re.search(r'(?:chuyen muc|category|danh muc|loai)[:\\s]*([^|\n]+?)(?:$|\n)', text, re.IGNORECASE)
if name_match: product_name = name_match.group(1).strip()
if desc_match: description = desc_match.group(1).strip()
if price_match: price = price_match.group(1).strip()
if cat_match: category = cat_match.group(1).strip()
if image_url or image_data_b64 or product_name or description or price or category:
_save_product_to_main_dataset(
image_url=image_url, image_data_b64=image_data_b64,
description=description, price=price, category=category,
sender_id=sender_id, sender_name=sender_name, product_name=product_name,
chat_id=chat_id,
)
BOT_STATE["logs"].append({
"event": "product_saved", "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type,
"text": "Saved! " + str(product_name)[:50],
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
elif _is_zgr_sender_local(sender_id) and text:
_save_text_message_to_dataset(
text=text, description=text[:200], price=price, category=category,
sender_id=sender_id, sender_name=sender_name, chat_id=chat_id,
)
BOT_STATE["logs"].append({
"event": "text_saved", "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type,
"text": "Text saved: " + str(text[:100]),
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
reply = "👋 Xin chào " + str(sender_name) + " (Zalo ID: " + str(sender_id) + ")!\n\n" + HELP_INSTRUCTIONS
asyncio.create_task(asyncio.to_thread(zapi.send_message, cid, reply))
elif event == "message.image.received" and chat_id:
cid = chat.get("id") or sender.get("id") or ""
_save_chat_id(cid, sender_id)
image_url = msg.get("photo", "") or msg.get("photo_url", "") or msg.get("image_url", "")
image_data_b64 = msg.get("image", "")
attachments = msg.get("attachment", {})
if attachments and not image_url:
payload = attachments.get("payload", {})
if isinstance(payload, str):
try:
payload = json.loads(payload)
except Exception:
payload = {}
image_url = payload.get("url", "")
image_data_b64 = payload.get("data", "")
product_name, description, price, category = "", "", "", ""
text = msg.get("caption", "") or text
name_match = re.search(r'(?:ten sp|ten san pham|product name|name)[:\s]*([^|\n]+)', text, re.IGNORECASE)
desc_match = re.search(r'(?:mo ta|description|desc|mota)[:\\s]*([^|\n]+)', text, re.IGNORECASE)
price_match = re.search(r'(?:gia|price|don gia|donggia)[:\\s]*([\\d,.]+)', text, re.IGNORECASE)
cat_match = re.search(r'(?:chuyen muc|category|danh muc|loai)[:\\s]*([^|\n]+?)(?:$|\n)', text, re.IGNORECASE)
if name_match: product_name = name_match.group(1).strip()
if desc_match: description = desc_match.group(1).strip()
if price_match: price = price_match.group(1).strip()
if cat_match: category = cat_match.group(1).strip()
log_text = "photo_url=" + str(image_url[:100]) if image_url else "No photo_url in message"
if text: log_text += " | caption=" + str(text[:200])
if image_url:
ts = time.strftime("%Y%m%d_%H%M%S")
log_text += " | image_url=" + str(image_url[:100])
logger.info("image.received from %s, url=%s", sender_id, image_url[:100])
img_bytes = None
if image_data_b64:
try:
img_bytes = base64.b64decode(image_data_b64)
except Exception:
img_bytes = None
if not img_bytes and image_url:
try:
r = requests.get(image_url, timeout=30)
img_bytes = r.content
logger.info("Downloaded image (%d bytes) from %s", len(img_bytes), image_url[:80])
except Exception as e:
logger.error("Image download failed: %s", e)
if img_bytes and not text:
logger.info("Running OCR extraction on image from %s", sender_id)
ocr_text = _ocr_extract_text(img_bytes)
if ocr_text:
ocr_name, ocr_desc, ocr_price, ocr_cat = _parse_ocr_product_info(ocr_text)
if not product_name: product_name = ocr_name
if not description: description = ocr_desc
if not price: price = ocr_price
if not category: category = ocr_cat
log_text += " | OCR: " + str(ocr_text[:200])
else:
log_text += " | OCR failed"
saved = _save_product_to_main_dataset(
image_url=image_url, image_data_b64=image_data_b64,
img_bytes=img_bytes,
description=description or text[:200], price=price, category=category,
sender_id=sender_id, sender_name=sender_name, product_name=product_name,
chat_id=chat_id, timestamp=ts,
)
if saved:
BOT_STATE["logs"].append({
"event": "product_saved", "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type,
"text": log_text + " | Image saved to dataset! " + str(product_name)[:50],
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
try:
zapi = ZaloBotAPI(BOT_STATE["bot_token"])
asyncio.create_task(asyncio.to_thread(zapi.send_message, cid, "GOT IT! Image saved to dataset!"))
except Exception as e:
logger.error("Reply failed: %s", e)
else:
BOT_STATE["logs"].append({
"event": "main_dataset_error", "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type,
"text": log_text + " | FAILED to save image to dataset",
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
else:
logger.warning("image.received but no photo_url found: %s", json.dumps(msg)[:300])
BOT_STATE["logs"].append({
"event": "image_no_url", "sender_id": sender_id, "chat_id": chat_id,
"sender_name": sender_name, "chat_type": chat_type,
"text": log_text,
"time": time.strftime("%Y-%m-%d %H:%M:%S"),
})
return Response(content=json.dumps({"message": "Success"}), media_type="application/json", status_code=200)
app = FastAPI(title="Zalo Bot Webhook")
@app.get("/")
async def root():
return RedirectResponse(url="/gradio/")
@app.get("/health")
async def health():
return {"status": "ok", "service": "zalo-bot-webhook", "main_dataset": MAIN_DATASET_ID, "zgr_sender_id": ZGR_SENDER_ID}
@app.post("/webhooks")
async def webhooks(request: Request):
return await handle_webhook(request)
@app.get("/logs")
async def logs_page():
log_lines = []
for log in reversed(BOT_STATE.get("logs", [])[-50:]):
sender_id = log.get("sender_id", "")
sender_name = log.get("sender_name", sender_id)
is_zgr = _is_zgr_sender_local(sender_id)
is_saved = log.get("event") in ("dataset_saved", "main_dataset_saved", "proxy_created", "product_saved", "image_saved")
is_error = log.get("event") in ("dataset_error", "main_dataset_error", "parse_error", "image_upload_fail", "image_no_url")
bg = "#e8f5e9" if is_zgr else "#ffffff"
header_color = "#4CAF50" if is_zgr else "#1a73e8"
zgr_badge = "[ZGR] " if is_zgr else ""
status_badge = "SUCCESS" if is_saved else ("ERROR" if is_error else "INFO")
rows += "<div style='margin:8px 0;padding:10px;background:" + bg + ";border-radius:6px;border-left:3px solid " + header_color + "'>"
+ "<div style='display:flex;gap:6px;align-items:center;flex-wrap:wrap'>"
+ "<b style='color:" + header_color + "'>" + status_badge + " [" + escape(str(log.get("event",""))) + "]</b>"
+ "<span style='color:#1a73e8'>👤 " + escape(str(sender_name)) + "</span>"
+ "<span style='color:#666'>🆔 " + escape(str(sender_id)) + "</span>"
+ "<span style='color:#666'>💬 " + escape(str(log.get("chat_id",""))) + "</span>"
+ "<span style='color:#888'>[" + escape(str(log.get("chat_type",""))) + "]</span>"
+ "<span style='color:#4CAF50'>" + zgr_badge + "</span>"
+ "</div>"
+ "<div style='margin-top:4px;color:#333;font-family:monospace;font-size:13px;word-break:break-word'>"
+ escape(str(log.get("text",""))[:300])
+ "</div>"
+ "<div style='margin-top:2px;color:#999;font-size:11px'>⏰ " + escape(str(log.get("time","")) or time.strftime('%Y-%m-%d %H:%M:%S')) + " | <a href='/logs/zgr-b7e1e71cf5701c2e4561'>zgr logs</a></div>"
+ "</div>"
total_logs = len(BOT_STATE.get("logs", []))
total_proxies = len(BOT_STATE.get("api_spaces", []))
connected_status = "✅" if BOT_STATE.get("connected") else "❌"
last_sender = escape(str(BOT_STATE.get("last_sender_id", "")[:8]) or "—")
zgr_count = sum(1 for l in BOT_STATE.get("logs", []) if _is_zgr_sender_local(l.get("sender_id", "")))
rows_html = "".join(log_lines) if log_lines else '<p style="color:#999">Chưa có sự kiện</p>'
html_content = (
'<!DOCTYPE html><html><head><title>Zalo Bot Logs</title>'
'<meta http-equiv="refresh" content="5">'
'<meta name="viewport" content="width=device-width, initial-scale=1">'
'<style>body { font-family: Arial, sans-serif; max-width: 1200px; margin: 0 auto; padding: 16px; background:#fafafa; }'
'h1 { color: #1a73e8; margin-bottom: 4px; }'
'.subtitle { color: #5f6368; font-size: 14px; margin-bottom: 16px; }'
'.stats { display: flex; gap: 16px; margin: 16px 0; flex-wrap: wrap; }'
'.stat-box { background: #e8f0fe; padding: 12px 24px; border-radius: 10px; min-width: 140px; }'
'.stat-value { font-size: 26px; font-weight: bold; color: #1a73e8; }'
'.stat-label { font-size: 12px; color: #5f6368; }'
'.log-container { max-height: 650px; overflow-y: auto; background:#fff; border-radius:8px; padding:8px; }'
'</style></head><body>'
'<h1>Zalo Bot Logs</h1>'
'<p class="subtitle">Event details</p>'
'<p>Links: <a href="/gradio/">Main UI</a> | <a href="/proxy-spaces">Proxy spaces</a></p>'
'<div class="stats">'
'<div class="stat-box"><div class="stat-value">' + str(total_logs) + '</div><div class="stat-label">Total Events</div></div>'
'<div class="stat-box"><div class="stat-value">' + str(total_proxies) + '</div><div class="stat-label">Proxies</div></div>'
'<div class="stat-box"><div class="stat-value">' + str(zgr_count) + '</div><div class="stat-label">ZGR Events</div></div>'
'<div class="stat-box"><div class="stat-value">' + connected_status + '</div><div class="stat-label">Bot Status</div></div>'
'<div class="stat-box"><div class="stat-value">' + last_sender + '</div><div class="stat-label">Last Sender</div></div>'
'</div>'
'<h2>Events (' + str(total_logs) + ')</h2>'
'<div class="log-container">' + rows_html + '</div>'
'<p><a href="/logs/zgr-b7e1e71cf5701c2e4561">Xem logs riêng cho nhóm ZGR</a></p>'
'</body></html>'
)
return Response(content=html_content, media_type="text/html")
@app.get("/logs/zgr-b7e1e71cf5701c2e4561")
async def zgr_logs_page():
zgr_logs = []
for log in BOT_STATE.get("logs", []):
sid = str(log.get("sender_id", ""))
if ZGR_SENDER_ID in sid or sid == ZGR_SENDER_ID:
zgr_logs.append(log)
if log.get("is_zgr", False) and not log.get("sender_id"):
zgr_logs.append(log)
log_lines = []
for idx, log in enumerate(reversed(zgr_logs[-50:])):
sender_id = log.get("sender_id", "")
sender_name = log.get("sender_name", sender_id)
is_saved = log.get("event") in ("dataset_saved", "main_dataset_saved", "product_saved")
is_error = log.get("event") in ("dataset_error", "main_dataset_error", "image_upload_fail")
bg = "#ffffff" if idx % 2 == 0 else "#fafafa"
status_color = "#4CAF50" if is_saved else ("#f44336" if is_error else "#1a73e8")
status_icon = "SUCCESS" if is_saved else ("ERROR" if is_error else "INFO")
rows = "<div style='margin:8px 0;padding:10px;background:" + bg + ";border-radius:6px;border-left:3px solid " + status_color + "'>"
+ "<div style='display:flex;gap:6px;align-items:center;flex-wrap:wrap'>"
+ "<b style='color:" + status_color + "'>" + status_icon + " [" + escape(str(log.get("event",""))) + "]</b>"
+ "<span style='color:#1a73e8;font-weight:bold'>👤 " + escape(str(sender_name)) + "</span>"
+ "<span style='color:#666'>🆔 " + escape(str(sender_id)[:20]) + "</span>"
+ "<span style='color:#666'>💬 " + escape(str(log.get("chat_id",""))[:20]) + "</span>"
+ "<span style='color:#666'>[" + escape(str(log.get("chat_type",""))) + "]</span>"
+ "</div>"
+ "<div style='margin-top:4px;color:#333;font-family:monospace;font-size:13px;word-break:break-word'>"
+ escape(str(log.get("text",""))[:300])
+ "</div>"
+ "<div style='margin-top:2px;color:#999;font-size:11px'>⏰ " + escape(str(log.get("time",""))) + " | <a href='https://huggingface.co/datasets/bep40/zalo-products-all' target='_blank'>Dataset</a></div>"
+ "</div>"
log_lines.append(rows)
saved_count = sum(1 for l in zgr_logs if l.get("event") in ("dataset_saved", "main_dataset_saved", "product_saved"))
error_count = sum(1 for l in zgr_logs if l.get("event") in ("dataset_error", "main_dataset_error", "image_upload_fail"))
total_zgr_logs = len(zgr_logs)
rows_html = "".join(log_lines) if log_lines else '<p style="color:#999">Chưa có sự kiện cho nhóm này. Gửi ảnh + thông tin sản phẩm để kiểm tra.</p>'
html_content = (
'<!DOCTYPE html><html><head><title>Zalo Logs - ZGR Group</title>'
'<meta http-equiv="refresh" content="5">'
'<meta name="viewport" content="width=device-width, initial-scale=1">'
'<style>'
'body { font-family: Arial, sans-serif; max-width: 1200px; margin: 0 auto; padding: 16px; background:#fafafa; }'
'h1 { color: #1a73e8; margin-bottom: 4px; }'
'.subtitle { color: #5f6368; font-size: 14px; margin-bottom: 16px; }'
'.stats { display: flex; gap: 16px; margin: 16px 0; flex-wrap: wrap; }'
'.stat-box { padding: 12px 24px; border-radius: 10px; min-width: 140px; }'
'.stat-value { font-size: 26px; font-weight: bold; }'
'.stat-saved { background: #e8f5e9; } .stat-saved .stat-value { color: #4CAF50; }'
'.stat-error { background: #ffebee; } .stat-error .stat-value { color: #f44336; }'
'.stat-total { background: #e8f0fe; } .stat-total .stat-value { color: #1a73e8; }'
'.log-container { max-height: 650px; overflow-y: auto; background:#fff; border-radius:8px; padding:8px; }'
'</style></head><body>'
'<h1>Zalo Logs - ZGR Group (zgr-b7e1e71cf5701c2e4561)</h1>'
'<p class="subtitle">All webhook events from this group</p>'
'<p>Links: <a href="/logs">All logs</a> | <a href="/gradio/">Main UI</a> | <a href="/proxy-spaces">Proxies</a></p>'
'<div class="stats">'
'<div class="stat-box stat-total"><div class="stat-value">' + str(total_zgr_logs) + '</div><div class="stat-label">Total Events</div></div>'
'<div class="stat-box stat-saved"><div class="stat-value">' + str(saved_count) + '</div><div class="stat-label">Saved to Dataset</div></div>'
'<div class="stat-box stat-error"><div class="stat-value">' + str(error_count) + '</div><div class="stat-label">Errors</div></div>'
'<div class="stat-box stat-total"><div class="stat-value"><a href="https://huggingface.co/datasets/bep40/zalo-products-all" target="_blank">Dataset</a></div><div class="stat-label">Main Dataset</div></div>'
'</div>'
'<h2>Events (' + str(total_zgr_logs) + ')</h2>'
'<div class="log-container">' + rows_html + '</div>'
'<p style="color:#5f6368;font-size:13px;margin-top:12px">Send image + "Tên sp: ..., Giá: ..., Chuyên mục: ..." to test.</p>'
'</body></html>'
)
return Response(content=html_content, media_type="text/html")
@app.get("/proxy-spaces")
async def proxy_spaces_page():
rows = []
for s in reversed(BOT_STATE.get("api_spaces", [])[-20:]):
repo_name = escape(str(s.get("repo_id", "").split("/")[-1]))
dataset_id_val = s.get("dataset_id", "")
dataset_link = "<a href='https://huggingface.co/datasets/" + escape(dataset_id_val) + "' target='_blank'>💾 dataset</a>" if dataset_id_val else ""
rows.append(
"<div style='margin:8px 0;padding:12px;background:#fff;border-radius:8px;border-left:4px solid #4CAF50;box-shadow:0 1px 3px rgba(0,0,0,0.1)'>"
"<div style='display:flex;gap:8px;align-items:center;flex-wrap:wrap;justify-content:space-between'>"
"<div>"
"<b style='color:#1a73e8'>👤 " + escape(str(s.get("sender_name", ""))) + "</b>"
"<span style='color:#666'>🆔 " + escape(str(s.get("user_id", ""))) + "</span>"
"<span style='color:#4CAF50;font-weight:bold'>[" + escape(str(s.get("status", ""))) + "]</span>"
"</div>"
"<div style='display:flex;gap:6px;flex-wrap:wrap'>"
"<a href='https://huggingface.co/spaces/bep40/" + repo_name + "' target='_blank'>Space</a>"
"<a href='" + escape(str(s.get("proxy_url", ""))) + "' target='_blank'>webhook</a>"
"<a href='" + escape(str(s.get("proxy_url", "")).replace("/webhooks", "/logs")) + "' target='_blank'>📊 logs</a>"
+ dataset_link +
"</div>"
"</div>"
"<div style='margin-top:6px'><span style='color:#5f6368'>Repo:</span> <code>" + escape(str(s.get("repo_id", ""))) + "</code></div>"
"<div style='margin-top:2px;color:#999;font-size:11px'>⏰ " + time.strftime('%Y-%m-%d %H:%M:%S') + "</div>"
"</div>"
)
total_proxies = len(BOT_STATE.get("api_spaces", []))
connected_status = "✅" if BOT_STATE.get("connected") else "❌"
rows_html = "".join(rows) if rows else '<p style="color:#999">Chưa có proxy space nào</p>'
html_content = (
'<!DOCTYPE html><html><head><title>Proxy Spaces</title>'
'<meta http-equiv="refresh" content="5">'
'<meta name="viewport" content="width=device-width, initial-scale=1">'
'<style>'
'body { font-family: Arial, sans-serif; max-width: 1200px; margin: 0 auto; padding: 16px; background:#fafafa; }'
'h1 { color: #1a73e8; } .subtitle { color: #5f6368; }'
'.stats { display: flex; gap: 16px; margin: 16px 0; flex-wrap: wrap; }'
'.stat-box { background: #e8f0fe; padding: 12px 24px; border-radius: 10px; min-width: 140px; }'
'.stat-value { font-size: 26px; font-weight: bold; color: #1a73e8; }'
'.stat-label { font-size: 12px; color: #5f6368; }'
'.container { background:#fff; border-radius:8px; padding:12px; }'
'a { color: #1a73e8; text-decoration: none; cursor: pointer; }'
'a:hover { text-decoration: underline; }'
'</style></head><body>'
'<h1>Quản lý Proxy Spaces</h1>'
'<p class="subtitle">Danh sách các space proxy đã tạo.</p>'
'<p>Links: <a href="/logs">Logs</a> | <a href="/gradio/">Main UI</a></p>'
'<div class="stats">'
'<div class="stat-box"><div class="stat-value">' + str(total_proxies) + '</div><div class="stat-label">Proxies</div></div>'
'<div class="stat-box"><div class="stat-value">' + connected_status + '</div><div class="stat-label">Bot Status</div></div>'
'</div>'
'<div class="container">'
+ rows_html +
'</div>'
'</body></html>'
)
return Response(content=html_content, media_type="text/html")
@app.get("/api/delete-proxy/{repo_name}")
async def delete_proxy(repo_name: str):
repo_id = NAMESPACE + "/" + repo_name
try:
api = HfApi(token=HF_TOKEN)
try:
api.delete_repo(repo_id=repo_id, repo_type="space", token=HF_TOKEN)
except Exception as e:
logger.warning("Could not delete HF Space %s: %s", repo_id, e)
for s in BOT_STATE.get("api_spaces", []):
if s.get("repo_id", "").split("/")[-1] == repo_name:
dataset_id = s.get("dataset_id", "")
if dataset_id:
try:
api.delete_repo(repo_id=dataset_id, repo_type="dataset", token=HF_TOKEN)
except Exception as e:
logger.warning("Could not delete dataset %s: %s", dataset_id, e)
BOT_STATE["api_spaces"] = [s for s in BOT_STATE.get("api_spaces", []) if s.get("repo_id", "").split("/")[-1] != repo_name]
_save_proxy_spaces()
return {"ok": True, "message": "Da xoa proxy " + repo_id}
except Exception as e:
logger.error("Delete proxy failed: %s", e)
return {"ok": False, "message": "Loi xoa: " + str(e)}
# ─── Data Import: file upload + URL scraping → extract structured data → save to main dataset ───
DATASET_COLUMNS = [
"product_name", "description", "price", "category", "technical_specs",
"sender_id", "sender_name", "chat_id", "timestamp",
"image", "text", "message_type", "is_zgr_group",
]
def _read_excel_file(filepath: str):
"""Read Excel (.xlsx/.xls) into a list of dicts."""
import pandas as pd
df = pd.read_excel(filepath, dtype=str)
return df.to_dict(orient="records")
def _read_csv_file(filepath: str):
"""Read CSV into a list of dicts."""
import pandas as pd
df = pd.read_csv(filepath, dtype=str)
return df.to_dict(orient="records")
def _read_docx_file(filepath: str):
"""Read Word (.docx) tables into a list of dicts."""
import docx
doc = docx.Document(filepath)
rows_data = []
for table in doc.tables:
rows = []
for row in table.rows:
cells = [cell.text.strip() for cell in row.cells]
if any(cells):
rows.append(cells)
if rows:
headers = rows[0]
for data_row in rows[1:]:
record = {}
for i, h in enumerate(headers):
record[h] = data_row[i] if i < len(data_row) else ""
rows_data.append(record)
if not rows_data:
paragraphs = [p.text.strip() for p in doc.paragraphs if p.text.strip()]
if paragraphs:
rows_data = [{"text": line, "description": line} for line in paragraphs]
return rows_data
def _read_txt_file(filepath: str):
"""Read plain text (.txt) into a list of dicts."""
import pandas as pd
from io import StringIO
with open(filepath, "r", encoding="utf-8", errors="replace") as f:
raw = f.read()
for sep in [",", "\t", ";", "|"]:
try:
df = pd.read_csv(StringIO(raw), sep=sep, dtype=str)
if len(df.columns) > 1:
return df.to_dict(orient="records")
except Exception:
continue
lines = [line.strip() for line in raw.splitlines() if line.strip()]
return [{"text": line, "description": line} for line in lines]
def _detect_file_type(filename: str) -> str:
"""Detect file type from filename extension."""
ext = filename.lower().rsplit(".", 1)[-1] if "." in filename else ""
mapping = {
"xlsx": "excel", "xls": "excel",
"csv": "csv", "txt": "txt",
"docx": "docx", "doc": "docx",
}
return mapping.get(ext, "unknown")
def _extract_file_data(filepath: str, file_type: str):
"""Dispatch to the right parser based on file_type."""
if file_type == "excel":
return _read_excel_file(filepath)
elif file_type == "csv":
return _read_csv_file(filepath)
elif file_type == "docx":
return _read_docx_file(filepath)
elif file_type == "txt":
return _read_txt_file(filepath)
else:
return [{"text": f"Unsupported file: {filepath}", "description": f"Unsupported file: {filepath}"}]
def _scrape_url_data(url: str):
"""Scrape tables and/or article content from a URL."""
import requests
from io import StringIO
rows_data = []
try:
headers = {"User-Agent": "Mozilla/5.0 (compatible; ZaloBotDataExtractor/1.0)"}
resp = requests.get(url, headers=headers, timeout=60)
resp.raise_for_status()
html = resp.text
except Exception as e:
logger.error("URL fetch failed for %s: %s", url, e)
return [{"text": f"Lỗi tải URL: {e}", "description": f"Lỗi tải URL: {e}"}], str(e)
try:
import pandas as pd
dfs = pd.read_html(StringIO(html))
for df in dfs:
df = df.astype(str)
rows_data.extend(df.to_dict(orient="records"))
except Exception as e:
logger.info("pd.read_html failed (maybe no tables): %s", e)
if not rows_data:
try:
from bs4 import BeautifulSoup
soup = BeautifulSoup(html, "html.parser")
tables = soup.find_all("table")
for table in tables:
rows = []
for tr in table.find_all("tr"):
cells = [td.get_text(strip=True) for td in tr.find_all(["td", "th"])]
if cells:
rows.append(cells)
if rows:
headers_bs = rows[0]
for data_row in rows[1:]:
record = {}
for i, h in enumerate(headers_bs):
record[h] = data_row[i] if i < len(data_row) else ""
rows_data.append(record)
except Exception as e:
logger.warning("BeautifulSoup table extraction failed: %s", e)
if not rows_data:
try:
from trafilatura import extract
text = extract(html, output_format="txt", include_tables=True)
if text:
lines = [line.strip() for line in text.splitlines() if line.strip()]
rows_data = [{"text": line, "description": line} for line in lines[:50]]
except Exception as e:
logger.warning("trafilatura extraction failed: %s", e)
if not rows_data:
rows_data = [{"text": url, "description": "Scraped from URL (no table detected)"}]
return rows_data, None
def _is_nan_value(value) -> bool:
"""Check if a value is NaN or equivalent string."""
if value is None:
return True
val_str = str(value).strip().lower()
if val_str in ("nan", "none", "null", "", "n/a", "na"):
return True
# Check for repeated "nan nan nan" patterns
if re.search(r'(nan\s+){2,}', val_str):
return True
return False
def _clean_text(value: str) -> str:
"""Clean text by removing NaN artifacts and repeated patterns."""
if value is None:
return ""
val_str = str(value).strip()
# Remove repeated "nan" patterns
val_str = re.sub(r'(?:\bnan\b\s*)+', '', val_str, flags=re.IGNORECASE)
# Remove multiple consecutive spaces/newlines
val_str = re.sub(r'\s+', ' ', val_str)
return val_str.strip()
def _classify_product_data(record: dict) -> dict:
"""Use AI-like pattern matching to classify product data into schema columns."""
rec = {col: "" for col in DATASET_COLUMNS}
# Collect all text content from the record
all_texts = []
for k, v in record.items():
if k and v and not _is_nan_value(v):
all_texts.append(f"{k}: {v}")
combined_text = "\n".join(all_texts)
# Try to extract product code/SKU
sku_match = re.search(r'(\b[A-Z0-9]{2,20}[-]\d{2,10}\b)', combined_text)
if sku_match:
rec["product_name"] = sku_match.group(1) if not rec.get("product_name") else rec["product_name"]
# Try to extract price (Vietnamese đồng format)
price_matches = re.findall(r'([\d.,]+)\s*(?:VNĐ|vnđ|đ|VND)', combined_text, re.IGNORECASE)
if price_matches:
rec["price"] = price_matches[0]
# Try to classify category
category_keywords = [
(r'xoong|nồi|inox|kitchen|cuisine', "Bếp nhà bếp"),
(r'tủ| Cabinet|tủ bếp|tủ âm', "Tủ nội thất"),
(r'ghế|chair|ghế sofa|ghế họp', "Đồ nội thất"),
(r'bàn|table|desk', "Bàn làm việc"),
(r'phòng|room|phòng học|phòng hội', "Nội thất phòng"),
]
for pattern, cat in category_keywords:
if re.search(pattern, combined_text, re.IGNORECASE):
rec["category"] = cat
break
# Build description from non-price, non-SKU text
desc_parts = []
for k, v in record.items():
if _is_nan_value(v):
continue
val_str = str(v).strip()
if not val_str:
continue
desc_parts.append(f"{k}: {val_str}" if k.lower() != 'text' else val_str)
rec["description"] = _clean_text("\n".join(desc_parts))[:500]
# Copy over any direct mappings
for key in ["product_name", "price", "category", "text"]:
val = record.get(key, "")
if not _is_nan_value(val) and val and not rec.get(key):
rec[key] = _clean_text(str(val))
return rec
def _normalize_timestamp(ts_val: str) -> str:
"""Normalize timestamp to ISO format."""
if not ts_val or _is_nan_value(ts_val):
return time.strftime("%Y-%m-%d %H:%M:%S")
# Try to parse and reformat
try:
# Handle "YYYYMMDD_HHMMSS" format
m = re.match(r'(\d{4})(\d{2})(\d{2})_(\d{2})(\d{2})(\d{2})', str(ts_val))
if m:
return f"{m.group(1)}-{m.group(2)}-{m.group(3)} {m.group(4)}:{m.group(5)}:{m.group(6)}"
# Try other common timestamp formats
for fmt in ["%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S", "%d/%m/%Y %H:%M:%S", "%Y/%m/%d %H:%M:%S", "%Y-%m-%d"]:
try:
parsed = time.strptime(str(ts_val), fmt)
return time.strftime("%Y-%m-%d %H:%M:%S", parsed)
except ValueError:
continue
except Exception:
pass
return str(ts_val)[:50] if ts_val else time.strftime("%Y-%m-%d %H:%M:%S")
def _is_product_row(record: dict) -> bool:
"""Aggressive filter - strict checks for product data. Filters out customer info, addresses, phones, headers, etc."""
all_text = " ".join(str(v) for v in record.values() if v and not _is_nan_value(v)).lower().strip()
if not all_text or len(all_text) < 10:
return False
# Các từ khóa nhận diện sản phẩm - BẮT BUỘC phải có ít nhất 1
product_indicators = [
"sp", "sản phẩm", "product", " hàng ", "hàng hóa", "mặt hàng",
"giá", "gia", "price", "đơn giá", "don gia", "cost", "costs",
"chuyên mục", "chuyen muc", "category", "danh mục", "danh muc", "loại", "loai",
"kích thước", "kich thuoc", "size", "chất liệu", "chat lieu", "material",
"bảo hành", "bao hanh", "warranty", "xuất xứ", "xuat xu", "origin", "nguyên liệu",
"dòng sản phẩm", "dong san pham", "variant",
"sku", "mã sp", "ma sp", "mã sản phẩm", "mã hàng",
"đặc tính", "thông số", "thong so", "thông số kỹ thuật",
"số lượng", "sl", "số lượng sp",
"giá nhập", "giá bán", "giá bán lẻ", "giá bán buôn",
"trùng tên", "tên sp", "tên hàng", "tên sản phẩm", "ten sp",
"in stock", "còn hàng", "hết hàng",
"voucher", "giảm giá", "khuyến mãi", "promotion",
"đơn hàng", "order", "xoong", "nồi", "đồ gỗ", "kim loại",
"model", "phiên bản", "version", "xuất xứ",
]
has_product_indicator = any(ind in all_text for ind in product_indicators)
# Các từ khóa loại trừ - NẾU CÓ THÌ LOẠI BỎ NGAY (kể cả có product indicators)
exclude_patterns = [
"người nhận hàng", "người giao hàng", "người lập phiếu",
"người nhận", "người giao", "lập phiếu",
"stt hình tên sản phẩm", "stt",
"thông tin khách hàng", "thông tin giao hàng",
"địa chỉ giao hàng", "địa chỉ nhận hàng", "địa chỉ",
"số điện thoại", "điện thoại liên hệ", "sdt", "mobile", "điện thoại",
"email", "@gmail", "@yahoo", "@zoho",
"tổng cộng", "tổng tiền", "tổng",
"thuế", "phí ship", "phí vận chuyển", "vận chuyển",
"hình thức", "httt", "chuyển khoản",
"cảm ơn", "xin cảm ơn", "kính thưa", "trân trọng",
"tên khách hàng", "ten khach hang",
"ngày", "ngày tạo", "ngày đặt hàng",
"trạng thái", "trang thái",
"ghi chú", "note",
"tên đơn hàng", "số đơn hàng", "mã đơn",
"thành tiền", "thanh toán", "cod", "cash on delivery",
"shop", "store", "website", "fanpage", "facebook",
"hotline", "lien he", "liên hệ",
"admin", "manager", "nhân viên", "staff",
"policy", "privacy", "terms", "chính sách",
"follow", "like", "share", "comment", "review",
]
for pattern in exclude_patterns:
if pattern in all_text:
# Nếu có product indicators mạnh, có thể vẫn giữ lại
# Trừ với header bảng
if pattern == "stt" and ("hình" in all_text or "tên sản phẩm" in all_text or "mã sp" in all_text):
return False
strong_indicators = ["giá", "price", "đơn giá", "kích thước", "chất liệu",
"bảo hành", "xuất xứ", "mã sp", "mã hàng", "giá nhập", "giá bán"]
has_strong = any(ind in all_text for ind in strong_indicators)
if not has_strong:
return False
# Skip pure phone numbers
phone_pattern = r"[\d+\-\s]{7,}"
if re.search(phone_pattern, all_text) and not has_product_indicator:
return False
if has_product_indicator:
return True
# Count populated product fields
has_name = bool(record.get("product_name", "").strip())
has_desc = bool(record.get("description", "").strip()) and len(str(record.get("description", ""))) > 10
has_price = bool(record.get("price", "").strip())
has_category = bool(record.get("category", "").strip())
product_fields = sum([has_name, has_desc, has_price, has_category])
return product_fields >= 2
def _normalize_columns(record: dict) -> dict:
"""Normalize record keys to match the dataset schema."""
normalized = {col: "" for col in DATASET_COLUMNS}
norm_map = {
"product_name": ["product_name", "productname", "tên sp", "ten sp", "tên sản phẩm", "ten san pham", "name", "tên hàng", "ten hang"],
"description": ["description", "mô tả", "mo ta", "desc", "mota", "nội dung", "noi dung"],
"price": ["price", "giá", "gia", "đơn giá", "don gia", "costs", "cost"],
"category": ["category", "chuyên mục", "chuyen muc", "danh mục", "danh muc", "loại", "loai", "type"],
"sender_id": ["sender_id", "sender id"],
"sender_name": ["sender_name", "sender name", "người gửi", "nguoi gui", "from"],
"chat_id": ["chat_id", "chat id"],
"image": ["image", "ảnh", "anh", "hình ảnh", "hinh anh", "photo", "photo_url", "image_url"],
"text": ["text", "nội dung", "noi dung", "message", "tin nhắn", "tin nhan"],
"timestamp": ["timestamp", "thời gian", "thoi gian", "time"],
}
for key, value in record.items():
if key is None:
continue
key_lower = str(key).strip().lower()
value_str = str(value).strip() if value is not None else ""
matched = False
for target, aliases in norm_map.items():
if key_lower in [a.lower() for a in aliases]:
normalized[target] = value_str
matched = True
break
if not matched:
normalized["text"] = normalized.get("text", "") + (f"\n{value_str}" if normalized.get("text") else value_str)
if normalized["text"] and not normalized["description"]:
normalized["description"] = _clean_text(normalized["text"])[:500]
# Clean all fields to remove "nan" artifacts
for key in DATASET_COLUMNS:
normalized[key] = _clean_text(normalized[key]) if normalized[key] else ""
# Classify product data using AI-like pattern matching
classified = _classify_product_data(normalized)
for key in DATASET_COLUMNS:
if classified.get(key) and not normalized.get(key):
normalized[key] = classified[key][:500] if key in ("description", "text") else classified[key][:200]
# Normalize timestamp to ISO format to prevent ArrowInvalid when loading dataset
normalized["timestamp"] = _normalize_timestamp(normalized.get("timestamp", ""))
return normalized
def import_data_process(input_type: str, file_objs=None, url: str = ""):
"""Main processing function for the Import Data tab."""
import pandas as pd
records = []
error_msg = ""
if input_type == "file" and file_objs:
for filepath in file_objs:
try:
filename = os.path.basename(filepath)
file_type = _detect_file_type(filename)
rows = _extract_file_data(filepath, file_type)
for row in rows:
# Skip rows that are all "nan" after cleaning
cleaned = {k: _clean_text(v) for k, v in row.items() if not _is_nan_value(v)}
if not cleaned or all(not v for v in cleaned.values()):
continue
normalized = _normalize_columns(row)
# Skip fully empty records
if not normalized.get("product_name") and not normalized.get("description") and \
not normalized.get("text") and not normalized.get("price"):
continue
# Filter out non-product rows (customer info, addresses, phone numbers, etc.)
if not _is_product_row(normalized):
logger.info("Filtered non-product row: %s", str(normalized.get("text",""))[:80])
continue
# Try to extract image from URL in the data
image_url = normalized.get("image", "")
if image_url and str(image_url).startswith("http"):
try:
ir = requests.get(image_url, timeout=30)
ir.raise_for_status()
img_data = ir.content
logger.info("Downloaded image (%d bytes) from URL in data", len(img_data))
except Exception as e:
logger.error("Image download from data URL failed: %s", e)
records.append(normalized)
except Exception as e:
logger.error("Failed to process file %s: %s", filepath, e)
records.append({col: "" for col in DATASET_COLUMNS})
records[-1]["description"] = f"Lỗi: {e}"
elif input_type == "url" and url:
rows, scrape_err = _scrape_url_data(url)
for row in rows:
cleaned = {k: _clean_text(v) for k, v in row.items() if not _is_nan_value(v)}
if not cleaned or all(not v for v in cleaned.values()):
continue
normalized = _normalize_columns(row)
if not normalized.get("product_name") and not normalized.get("description") and \
not normalized.get("text") and not normalized.get("price"):
continue
# Filter out non-product rows (customer info, addresses, phone numbers, etc.)
if not _is_product_row(normalized):
logger.info("Filtered non-product row from URL: %s", str(normalized.get("text",""))[:80])
continue
# Try to extract image from URL in scraped data
image_url = normalized.get("image", "")
if image_url and str(image_url).startswith("http"):
try:
ir = requests.get(image_url, timeout=30)
ir.raise_for_status()
img_data = ir.content
logger.info("Downloaded image (%d bytes) from URL in scraped data", len(img_data))
except Exception as e:
logger.error("Image download from scraped URL failed: %s", e)
records.append(normalized)
if scrape_err:
error_msg = scrape_err
if not records:
return gr.update(visible=True, value=pd.DataFrame(columns=DATASET_COLUMNS)), gr.update(visible=True, value="⚠️ Không có dữ liệu nào được tìm thấy.")
save_results = []
records_to_save = []
for rec in records:
try:
ts = time.strftime("%Y%m%d_%H%M%S")
rec_ts = time.strftime("%Y-%m-%d %H:%M:%S")
# Clean all fields to remove "nan" artifacts
rec = {col: _clean_text(str(rec.get(col, ""))) for col in DATASET_COLUMNS}
rec.setdefault("sender_id", "import_user")
rec.setdefault("sender_name", "Data Import")
rec.setdefault("chat_id", "")
rec.setdefault("timestamp", rec_ts)
rec.setdefault("message_type", "import")
rec.setdefault("is_zgr_group", "False")
if not rec.get("price") and not rec.get("product_name") and rec.get("text"):
rec["message_type"] = "text"
# ─── XỬ LÝ ẢNH TỪ URL HOỆC FILE ──────────────────────────────
image_url = rec.get("image", "")
img_bytes = None
uploaded_img = ""
if image_url:
# Nếu là URL http/https → tải ảnh về
if str(image_url).startswith("http"):
try:
r = requests.get(image_url, timeout=30)
r.raise_for_status()
img_bytes = r.content
logger.info("Downloaded image (%d bytes) from %s", len(img_bytes), image_url[:80])
except Exception as e:
logger.error("Image download failed from URL %s: %s", image_url[:80], e)
# Nếu là base64 → decode
elif str(image_url).startswith("data:"):
try:
b64_data = image_url.split("base64,")[-1]
img_bytes = base64.b64decode(b64_data)
except Exception as e:
logger.error("Base64 image decode failed: %s", e)
if img_bytes:
img_hash = hashlib.md5(img_bytes).hexdigest()[:8]
img_filename = f"images/{ts}_import_{img_hash}.jpg"
with tempfile.NamedTemporaryFile(suffix=".jpg", delete=False) as img_tmp:
img_tmp.write(img_bytes)
img_tmp_path = img_tmp.name
try:
api_upload = HfApi(token=HF_TOKEN) if HF_TOKEN else HfApi()
api_upload.upload_file(
path_or_fileobj=img_tmp_path,
path_in_repo=img_filename,
repo_id=MAIN_DATASET_ID,
repo_type="dataset",
token=HF_TOKEN,
commit_message=f"Add image for {rec.get('product_name','')[:30]}",
)
uploaded_img = img_filename
save_results.append(f"🖼️ Ảnh: {img_filename}")
except Exception as e:
logger.error("Image upload failed: %s", e)
save_results.append(f"❌ Lỗi ảnh: {e}")
finally:
pathlib.Path(img_tmp_path).unlink(missing_ok=True)
# Extract technical specs from description or text if not already populated
tech_specs = rec.get("technical_specs", "")
if not tech_specs:
tech_specs = extract_technical_specs_from_text(
rec.get("description", "") or rec.get("text", "")
)
record_to_save = {
"image": uploaded_img if uploaded_img else (rec.get("image", "")),
"product_name": rec.get("product_name", "")[:200],
"description": rec.get("description", "")[:500],
"price": rec.get("price", ""),
"category": rec.get("category", ""),
"technical_specs": tech_specs[:500] if tech_specs else "",
"sender_id": rec.get("sender_id", "import_user"),
"sender_name": rec.get("sender_name", "Data Import"),
"chat_id": rec.get("chat_id", ""),
"timestamp": rec_ts,
"text": rec.get("text", ""),
"message_type": rec.get("message_type", "import"),
"is_zgr_group": str(rec.get("is_zgr_group", False)),
}
records_to_save.append(record_to_save)
save_results.append(f"✅ Đã chuẩn bị lưu bản ghi")
except Exception as e:
logger.error("Record processing error: %s", e)
save_results.append(f"❌ Lỗi: {e}")
# Save all records to dataset as parquet merge
if records_to_save and HF_TOKEN:
_append_to_main_dataset_parquet(records_to_save)
df = pd.DataFrame(records)
status_parts = [f"📊 Đã xử lý {len(records)} bản ghi."]
if save_results:
saved_ok = sum(1 for s in save_results if s.startswith("✅"))
saved_fail = sum(1 for s in save_results if s.startswith("❌"))
status_parts.append(f"✅ Đã lưu {saved_ok} bản ghi vào {MAIN_DATASET_ID}")
if saved_fail:
status_parts.append(f"❌ {saved_fail} lỗi")
if error_msg:
status_parts.append(f"⚠️ Cảnh báo: {error_msg}")
status_msg = "\n".join(status_parts)
return gr.update(visible=True, value=df), gr.update(visible=True, value=status_msg)
def get_dataset_info():
"""Get info about the main dataset for display in the Import Data tab."""
info = f"**Dataset chính:** `{MAIN_DATASET_ID}`\n\n"
info += "**Cấu trúc (schema):**\n"
info += "| Trường | Kiểu | Mô tả |\n"
info += "|--------|------|-------|\n"
info += "| product_name | text | Tên sản phẩm |\n"
info += "| description | text | Nội dung mô tả |\n"
info += "| price | number | Giá sản phẩm |\n"
info += "| category | text | Chuyên mục |\n"
info += "| technical_specs | text | Thông số kỹ thuật |\n"
info += "| sender_id | text | ID người gửi |\n"
info += "| sender_name | text | Tên người gửi |\n"
info += "| chat_id | text | ID chat |\n"
info += "| timestamp | text | Thời gian ghi nhận |\n"
info += "| image | text | Đường link ảnh (nếu có) |\n"
info += "| text | text | Nội dung tin nhắn/văn bản |\n"
info += "| message_type | text | Loại tin (import/text/product) |\n"
info += "| is_zgr_group | bool | Gửi từ nhóm ZGR |\n\n"
info += "💡 **Cách dùng:**\n"
info += "1. Chọn chế độ: tải lên file hoặc nhập URL\n"
info += "2. Hỗ trợ: Excel (.xlsx), CSV, Word (.docx), TXT\n"
info += "3. Các cột trong file có thể dùng tiếng Việt hoặc tiếng Anh (ví dụ: 'Tên sp', 'Giá', 'Chuyên mục', 'Mô tả')\n"
info += "4. Kết quả sẽ được trích xuất và lưu vào dataset theo cấu trúc chuẩn"
return info
# ─── Initialize main dataset schema on startup ───
_ensure_main_dataset_schema()
with gr.Blocks(title="Zalo Bot Webhook") as demo:
gr.Markdown("# Zalo Bot Webhook Setup")
with gr.Tabs():
with gr.Tab("Kết nối Bot"):
tok = gr.Textbox(DEFAULT_BOT_TOKEN, label="Bot Token", type="password")
btn = gr.Button("Kết nối")
res = gr.Markdown("")
btn.click(fn=connect_bot, inputs=[tok], outputs=res)
gr.Textbox(value=get_botinfo, label="Thông tin bot", interactive=False, lines=8)
with gr.Tab("Gửi tin"):
with gr.Row():
cid = gr.Textbox(label="Chat ID")
txt = gr.Textbox("HTTP API: 4179413508988279245:abc123", label="Nội dung")
b = gr.Button("Gửi")
b.click(fn=send_msg, inputs=[cid, txt], outputs=gr.Textbox(label="Kết quả"))
gr.Textbox(value=get_events, label="Sự kiện nhận được", interactive=False, lines=20)
gr.Markdown(
"Links: [Proxy spaces](/proxy-spaces)\n\n"
"Each user sends HTTP API token to create their own proxy space."
)
with gr.Tab("Nhập dữ liệu"):
with gr.Row():
with gr.Column(scale=2):
input_choice = gr.Radio(
choices=["file", "url"],
value="file",
label="Chọn nguồn dữ liệu",
info="Chọn tải file lên hoặc nhập URL để cào dữ liệu",
)
file_input = gr.File(
file_types=[".xlsx", ".xls", ".csv", ".txt", ".docx", ".doc"],
file_count="multiple",
label="Tải lên file (Excel, CSV, TXT, Word)",
visible=True,
)
url_input = gr.Textbox(
label="Nhập URL (cào từ trang web bất kỳ)",
placeholder="https://example.com/products",
visible=False,
)
import_btn = gr.Button("Xuất khẩu dữ liệu", variant="primary")
with gr.Column(scale=1):
gr.Markdown(get_dataset_info)
df_output = gr.Dataframe(
headers=DATASET_COLUMNS,
interactive=False,
visible=False,
label="Dữ liệu được trích xuất",
)
status_output = gr.Markdown("", visible=False)
def _toggle_inputs(choice):
show_file = choice == "file"
show_url = choice == "url"
return [
gr.update(visible=show_file),
gr.update(visible=show_url),
]
input_choice.change(
fn=_toggle_inputs,
inputs=[input_choice],
outputs=[file_input, url_input],
)
import_btn.click(
fn=import_data_process,
inputs=[input_choice, file_input, url_input],
outputs=[df_output, status_output],
)
with gr.Tab("Hướng dẫn"):
gr.Markdown("1. Go to https://zalo.me/s/botcreator\n2. Copy HTTP API token\n3. Send to bot to auto-create proxy")
demo.queue()
app = gr.mount_gradio_app(app, demo, path="/gradio")
print("[startup] FastAPI app ready: /health, /webhooks, /logs, /logs/zgr-b7e1e71cf5701c2e4561, /proxy-spaces, /gradio/", flush=True)
if __name__ == "__main__":
port = int(os.getenv("PORT", "7860"))
server_name = os.getenv("GRADIO_SERVER_NAME", "0.0.0.0")
print("[launch] uvicorn on " + server_name + ":" + str(port), flush=True)
import uvicorn
uvicorn.run(app, host=server_name, port=port) |