public_gateway.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411
  1. import logging
  2. import time
  3. import uuid
  4. from constants import DEVICE_INVALID_MESSAGE
  5. from tools.list_order_filter_options import ListOrderFilterOptionsTool
  6. from tools.list_customer_filter_options import ListCustomerFilterOptionsTool
  7. from tools.list_receivable_cost_filter_options import (
  8. ListReceivableCostFilterOptionsTool,
  9. )
  10. from tools.list_outbound_filter_options import ListOutboundFilterOptionsTool
  11. from tools.export_pending_outbound_orders import ExportPendingOutboundOrdersTool
  12. from tools.export_out_of_province_port_data import (
  13. ExportOutOfProvincePortDataTool,
  14. )
  15. from tools.export_receivable_cost_list import ExportReceivableCostListTool
  16. from tools.list_pending_outbound_export_filter_options import (
  17. ListPendingOutboundExportFilterOptionsTool,
  18. )
  19. from tools.query_order import QueryOrderTool
  20. from tools.query_customs_declaration_files import (
  21. QueryCustomsDeclarationFilesTool,
  22. )
  23. from tools.query_order_exact import QueryOrderExactTool
  24. from tools.query_order_detail import QueryOrderDetailTool
  25. from tools.query_export_task import QueryExportTaskTool
  26. from tools.query_outbound_detail import QueryOutboundDetailTool
  27. from tools.query_outbound_list import QueryOutboundListTool
  28. from tools.query_track import QueryTrackTool
  29. from tools.query_customer_list import QueryCustomerListTool
  30. from tools.query_customer_payment_followup import QueryCustomerPaymentFollowupTool
  31. from tools.query_customer_unverified_bill_details import QueryCustomerUnverifiedBillDetailsTool
  32. from tools.query_customer_payment_records import QueryCustomerPaymentRecordsTool
  33. from tools.query_order_receivable_cost_details import (
  34. QueryOrderReceivableCostDetailsTool,
  35. )
  36. from tools.query_receivable_cost_list import QueryReceivableCostListTool
  37. from tools.query_payable_cost_list import QueryPayableCostListTool
  38. from tools.list_payable_cost_filter_options import ListPayableCostFilterOptionsTool
  39. from tools.export_payable_cost_list import ExportPayableCostListTool
  40. from tools.export_pallet_data import ExportPalletDataTool
  41. from tools.query_destination_trailer_list import QueryDestinationTrailerListTool
  42. from tools.list_destination_trailer_filter_options import (
  43. ListDestinationTrailerFilterOptionsTool,
  44. )
  45. from tools.query_order_abnormal_list import QueryOrderAbnormalListTool
  46. from tools.list_order_abnormal_filter_options import (
  47. ListOrderAbnormalFilterOptionsTool,
  48. )
  49. from tools.query_receive_volume_list import QueryReceiveVolumeListTool
  50. from tools.list_receive_volume_filter_options import (
  51. ListReceiveVolumeFilterOptionsTool,
  52. )
  53. from tools.query_container_timeliness_list import QueryContainerTimelinessListTool
  54. from tools.export_container_timeliness_report import (
  55. ExportContainerTimelinessReportTool,
  56. )
  57. from tools.query_headhaul_document_list import QueryHeadhaulDocumentListTool
  58. from tools.list_headhaul_document_filter_options import (
  59. ListHeadhaulDocumentFilterOptionsTool,
  60. )
  61. from tools.prepare_headhaul_document_upload import (
  62. PrepareHeadhaulDocumentUploadTool,
  63. )
  64. from tools.save_headhaul_document import SaveHeadhaulDocumentTool
  65. from tools.delete_headhaul_document import DeleteHeadhaulDocumentTool
  66. from utils.security import hash_gateway_session_id
  67. logger = logging.getLogger(__name__)
  68. class PublicGatewayApp:
  69. def __init__(self, session_store, api_client, auth_client=None):
  70. self.session_store = session_store
  71. self.api_client = api_client
  72. self._tools = {
  73. 'query_order': QueryOrderTool(api_client=None),
  74. 'query_track': QueryTrackTool(api_client=None),
  75. 'query_order_exact': QueryOrderExactTool(api_client=None),
  76. 'query_order_detail': QueryOrderDetailTool(api_client=None),
  77. 'query_customs_declaration_files':
  78. QueryCustomsDeclarationFilesTool(api_client=None),
  79. 'query_outbound_list': QueryOutboundListTool(api_client=None),
  80. 'query_outbound_detail': QueryOutboundDetailTool(api_client=None),
  81. 'query_customer_list': QueryCustomerListTool(api_client=None),
  82. 'query_customer_payment_followup': QueryCustomerPaymentFollowupTool(
  83. api_client=None
  84. ),
  85. 'query_customer_unverified_bill_details':
  86. QueryCustomerUnverifiedBillDetailsTool(api_client=None),
  87. 'query_customer_payment_records': QueryCustomerPaymentRecordsTool(
  88. api_client=None
  89. ),
  90. 'query_order_receivable_cost_details':
  91. QueryOrderReceivableCostDetailsTool(api_client=None),
  92. 'query_receivable_cost_list':
  93. QueryReceivableCostListTool(api_client=None),
  94. 'query_payable_cost_list':
  95. QueryPayableCostListTool(api_client=None),
  96. 'list_outbound_filter_options': ListOutboundFilterOptionsTool(
  97. api_client=None
  98. ),
  99. 'list_order_filter_options': ListOrderFilterOptionsTool(
  100. api_client=None
  101. ),
  102. 'list_customer_filter_options': ListCustomerFilterOptionsTool(
  103. api_client=None
  104. ),
  105. 'list_receivable_cost_filter_options':
  106. ListReceivableCostFilterOptionsTool(api_client=None),
  107. 'list_payable_cost_filter_options':
  108. ListPayableCostFilterOptionsTool(api_client=None),
  109. 'export_pending_outbound_orders': ExportPendingOutboundOrdersTool(
  110. api_client=None
  111. ),
  112. 'export_out_of_province_port_data':
  113. ExportOutOfProvincePortDataTool(api_client=None),
  114. 'export_receivable_cost_list': ExportReceivableCostListTool(
  115. api_client=None
  116. ),
  117. 'export_payable_cost_list': ExportPayableCostListTool(
  118. api_client=None
  119. ),
  120. 'export_pallet_data': ExportPalletDataTool(
  121. api_client=None
  122. ),
  123. 'query_export_task': QueryExportTaskTool(api_client=None),
  124. 'list_pending_outbound_export_filter_options':
  125. ListPendingOutboundExportFilterOptionsTool(api_client=None),
  126. 'query_destination_trailer_list': QueryDestinationTrailerListTool(
  127. api_client=None
  128. ),
  129. 'list_destination_trailer_filter_options':
  130. ListDestinationTrailerFilterOptionsTool(api_client=None),
  131. 'query_order_abnormal_list': QueryOrderAbnormalListTool(
  132. api_client=None
  133. ),
  134. 'list_order_abnormal_filter_options':
  135. ListOrderAbnormalFilterOptionsTool(api_client=None),
  136. 'query_receive_volume_list': QueryReceiveVolumeListTool(
  137. api_client=None
  138. ),
  139. 'list_receive_volume_filter_options':
  140. ListReceiveVolumeFilterOptionsTool(api_client=None),
  141. 'query_container_timeliness_list': QueryContainerTimelinessListTool(
  142. api_client=None
  143. ),
  144. 'export_container_timeliness_report':
  145. ExportContainerTimelinessReportTool(api_client=None),
  146. 'query_headhaul_document_list': QueryHeadhaulDocumentListTool(
  147. api_client=None
  148. ),
  149. 'list_headhaul_document_filter_options':
  150. ListHeadhaulDocumentFilterOptionsTool(api_client=None),
  151. 'prepare_headhaul_document_upload':
  152. PrepareHeadhaulDocumentUploadTool(api_client=None),
  153. 'save_headhaul_document': SaveHeadhaulDocumentTool(
  154. api_client=None
  155. ),
  156. 'delete_headhaul_document': DeleteHeadhaulDocumentTool(
  157. api_client=None
  158. ),
  159. }
  160. from services.headhaul_document_app import HeadhaulSlotMemory
  161. self._headhaul_slots = HeadhaulSlotMemory()
  162. def registered_tool_names(self):
  163. return tuple(self._tools.keys())
  164. def _require_session(self, gateway_session_id, diagnostic_emitter=None):
  165. session = self.session_store.get(gateway_session_id)
  166. if not session or not session.get('mcp_token'):
  167. if diagnostic_emitter is not None:
  168. diagnostic_emitter.emit(
  169. stage='gateway_session',
  170. status='failed',
  171. event_code='GATEWAY_SESSION_NOT_FOUND',
  172. session_credential=gateway_session_id,
  173. context={'transport': 'http'},
  174. )
  175. raise RuntimeError(DEVICE_INVALID_MESSAGE)
  176. if diagnostic_emitter is not None:
  177. diagnostic_emitter.set_defaults(
  178. session_credential=gateway_session_id,
  179. admin_id=session.get('admin_id'),
  180. company_id=session.get('company_id'),
  181. context={'transport': 'http'},
  182. )
  183. diagnostic_emitter.emit(
  184. stage='gateway_session',
  185. status='succeeded',
  186. event_code='GATEWAY_SESSION_RESOLVED',
  187. )
  188. return session
  189. def _enabled_tool_names(self, response):
  190. if not isinstance(response, dict):
  191. raise RuntimeError('invalid enabled tool response')
  192. if response.get('code') != 'MCP_0000':
  193. raise RuntimeError(response.get('msg') or 'list enabled tools failed')
  194. data = response.get('data')
  195. codes = data.get('tool_codes') if isinstance(data, dict) else None
  196. if not isinstance(codes, list):
  197. raise RuntimeError('invalid enabled tool response')
  198. return {
  199. code.strip().lower()
  200. for code in codes
  201. if isinstance(code, str) and code.strip()
  202. }
  203. def _load_enabled_tool_names(self, token, request_id=''):
  204. response = self.api_client.list_enabled_tools(
  205. token,
  206. request_id=request_id,
  207. )
  208. return self._enabled_tool_names(response)
  209. def list_tools(self, gateway_session_id, request_id=''):
  210. session = self._require_session(gateway_session_id)
  211. if hasattr(self.session_store, 'touch_session'):
  212. self.session_store.touch_session(gateway_session_id)
  213. request_id = self.build_request_id(request_id)
  214. enabled = self._load_enabled_tool_names(
  215. session['mcp_token'],
  216. request_id,
  217. )
  218. return [
  219. tool.metadata()
  220. for name, tool in self._tools.items()
  221. if name in enabled
  222. ]
  223. def list_resources(self):
  224. return []
  225. def remember_headhaul_slot(self, gateway_session_id, token):
  226. self._headhaul_slots.remember(gateway_session_id, token)
  227. def session_for_headhaul_token(self, token):
  228. return self._headhaul_slots.session_for_token(token)
  229. def read_resource(self, uri, gateway_session_id=''):
  230. return None
  231. def upload_headhaul_document(
  232. self,
  233. gateway_session_id,
  234. upload_token,
  235. filename,
  236. file_bytes,
  237. content_type='',
  238. request_id='',
  239. client_ip='',
  240. ):
  241. session = self._require_session(gateway_session_id)
  242. request_id = self.build_request_id(request_id)
  243. enabled = self._load_enabled_tool_names(
  244. session['mcp_token'],
  245. request_id,
  246. )
  247. if 'prepare_headhaul_document_upload' not in enabled:
  248. raise RuntimeError('tool disabled: prepare_headhaul_document_upload')
  249. if not hasattr(self.api_client, 'upload_headhaul_document'):
  250. raise RuntimeError('upload client unavailable')
  251. return self.api_client.upload_headhaul_document(
  252. token=session['mcp_token'],
  253. upload_token=upload_token,
  254. filename=filename,
  255. file_bytes=file_bytes,
  256. content_type=content_type,
  257. request_id=request_id,
  258. client_ip=client_ip,
  259. )
  260. def build_request_id(self, request_id=''):
  261. request_id = str(request_id or '').strip()
  262. return request_id or 'rq_{0}'.format(uuid.uuid4().hex[:16])
  263. def call_tool(
  264. self,
  265. gateway_session_id,
  266. name,
  267. arguments=None,
  268. request_id='',
  269. client_ip='',
  270. diagnostic_emitter=None,
  271. ):
  272. session = self._require_session(gateway_session_id, diagnostic_emitter)
  273. request_id = self.build_request_id(request_id)
  274. if name not in self._tools:
  275. if diagnostic_emitter is not None:
  276. diagnostic_emitter.emit(
  277. stage='backend_call',
  278. status='failed',
  279. event_code='TOOL_NOT_REGISTERED',
  280. context={'transport': 'http'},
  281. )
  282. raise KeyError('tool not registered: {0}'.format(name))
  283. try:
  284. enabled_tools = self._load_enabled_tool_names(
  285. session['mcp_token'],
  286. request_id,
  287. )
  288. except Exception:
  289. if diagnostic_emitter is not None:
  290. diagnostic_emitter.emit(
  291. stage='backend_call',
  292. status='failed',
  293. event_code='ENABLED_TOOL_LOOKUP_FAILED',
  294. tool_code=name,
  295. context={'transport': 'http'},
  296. )
  297. raise
  298. if name not in enabled_tools:
  299. if diagnostic_emitter is not None:
  300. diagnostic_emitter.emit(
  301. stage='backend_call',
  302. status='failed',
  303. event_code='TOOL_DISABLED',
  304. tool_code=name,
  305. context={'transport': 'http'},
  306. )
  307. raise RuntimeError('tool disabled: {0}'.format(name))
  308. tool = self._tools[name]
  309. if diagnostic_emitter is not None:
  310. diagnostic_emitter.set_defaults(tool_code=name)
  311. session_hash = hash_gateway_session_id(gateway_session_id)[:12]
  312. admin_id = session.get('admin_id')
  313. company_id = session.get('company_id')
  314. logger.info(
  315. 'MCP public tool call',
  316. extra={
  317. 'request_id': request_id,
  318. 'tool_code': name,
  319. 'session_hash': session_hash,
  320. 'admin_id': admin_id,
  321. 'company_id': company_id,
  322. },
  323. )
  324. try:
  325. started_at = time.monotonic()
  326. if diagnostic_emitter is not None:
  327. diagnostic_emitter.emit(
  328. stage='backend_call',
  329. status='started',
  330. event_code='BACKEND_CALL_STARTED',
  331. )
  332. result = self.api_client.call_tool(
  333. token=session['mcp_token'],
  334. tool_code=tool.name,
  335. route_path=tool.route_path,
  336. payload=arguments or {},
  337. request_id=request_id,
  338. client_ip=client_ip,
  339. )
  340. if hasattr(self.session_store, 'touch_session'):
  341. self.session_store.touch_session(gateway_session_id)
  342. if name == 'prepare_headhaul_document_upload' and isinstance(result, dict):
  343. data = result.get('data')
  344. token = ''
  345. if isinstance(data, dict):
  346. token = str(data.get('upload_token') or '').strip()
  347. if token:
  348. self.remember_headhaul_slot(gateway_session_id, token)
  349. logger.info(
  350. 'MCP public tool success',
  351. extra={
  352. 'request_id': request_id,
  353. 'tool_code': name,
  354. 'session_hash': session_hash,
  355. 'response_code': result.get('code'),
  356. },
  357. )
  358. if diagnostic_emitter is not None:
  359. diagnostic_emitter.emit(
  360. stage='backend_call',
  361. status='succeeded',
  362. event_code='BACKEND_CALL_COMPLETED',
  363. response_code=(
  364. result.get('code') if isinstance(result, dict) else None
  365. ),
  366. cost_ms=max(0, int((time.monotonic() - started_at) * 1000)),
  367. )
  368. return result
  369. except Exception as e:
  370. logger.error(
  371. 'MCP public tool failed',
  372. extra={
  373. 'request_id': request_id,
  374. 'tool_code': name,
  375. 'session_hash': session_hash,
  376. 'response_code': 'MCP_9001',
  377. 'diagnostic_reason': 'UNEXPECTED_EXCEPTION',
  378. 'exception_class': e.__class__.__name__,
  379. },
  380. )
  381. if diagnostic_emitter is not None:
  382. diagnostic_emitter.emit(
  383. stage='backend_call',
  384. status='failed',
  385. event_code='UNEXPECTED_EXCEPTION',
  386. response_code='MCP_9001',
  387. )
  388. raise