sync_account_token.py 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081
  1. import asyncio
  2. import json
  3. import time
  4. from datetime import datetime, timedelta
  5. from sqlalchemy.future import select
  6. from app.config.agent_base_url import RG_USER_LOGIN, DF_USER_LOGIN, RG_PING, DF_PING
  7. from app.config.config import settings
  8. from app.config.const import chat_server, workflow_server, http_200
  9. from app.models import UserTokenModel
  10. from app.models.app_token_model import AppToken
  11. from app.models.base_model import SessionLocal
  12. from app.models.postgresql_base_model import get_pdb, PostgresqlSessionLocal
  13. from app.service.ragflow import RagflowService
  14. from app.service.v2.app_driver.chat_data import ChatBaseApply
  15. from app.utils.password_handle import password_decrypted
  16. async def sync_token():
  17. """
  18. 1.获取到user_token表中的账号
  19. 2.判断是否超过12小时,超过则获取新的token
  20. 3.未超过12小时,则测试token有效性,chat_ping接口,返回401,则重新获取token
  21. 4.token获取:login
  22. df:/console/api/workspaces
  23. rg:/v1/system/version
  24. 5.跟新本地token和kong网关token
  25. :return:
  26. """
  27. app_data = [{"token_id": chat_server, "url": f"{settings.fwr_base_url}{RG_USER_LOGIN}",
  28. "data": {"email": "", "password": ""}, "is_crypt": True, "ping_url": f"{settings.fwr_base_url}{RG_PING}",
  29. "token": "{}"},
  30. {"token_id": workflow_server, "url": f"{settings.dify_base_url}{DF_USER_LOGIN}",
  31. "data": {"email": "", "password": "", "remember_me": True, "language": "zh-Hans"}, "is_crypt": False
  32. , "ping_url": f"{settings.dify_base_url}{DF_PING}", "token": "Bearer {}"}]
  33. async def sync_token_chat(token_id, url, data, is_crypt, ping_url, token):
  34. db = SessionLocal()
  35. # pdb = PostgresqlSessionLocal()
  36. current_time = datetime.now() - timedelta(hours=24)
  37. try:
  38. user_token = db.query(UserTokenModel).filter(UserTokenModel.id == token_id).first()
  39. chat = ChatBaseApply()
  40. if user_token and (user_token.updated_at < current_time or not user_token.access_token or await chat.chat_ping(ping_url, {}, await chat.get_chat_headers(token.format(user_token.access_token))) != http_200):
  41. data["email"] = user_token.account
  42. if is_crypt:
  43. data["password"] = await chat.password_encrypt(await password_decrypted(user_token.password))
  44. else:
  45. data["password"] = await password_decrypted(user_token.password)
  46. res = await chat.chat_login(url, data, {'Content-Type': 'application/json'})
  47. if res:
  48. access_token = res["data"]["access_token"]
  49. user_token.access_token = access_token
  50. user_token.updated_at = datetime.now()
  51. db.commit()
  52. except Exception as e:
  53. print(e)
  54. finally:
  55. db.close()
  56. # await pdb.close()
  57. tasks = []
  58. for app in app_data:
  59. tasks.append(asyncio.create_task(sync_token_chat(app["token_id"], app["url"], app["data"], app["is_crypt"], app["ping_url"], app["token"])))
  60. done, pending = await asyncio.wait(tasks, return_when=asyncio.ALL_COMPLETED)
  61. def start_sync_token_task():
  62. asyncio.run(sync_token())
  63. if __name__ == "__main__":
  64. start_sync_token_task()