test_dataflow_create_saga_neo4j.py 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137
  1. from __future__ import annotations
  2. import os
  3. from concurrent.futures import ThreadPoolExecutor
  4. from threading import Barrier
  5. import pytest
  6. from neo4j import GraphDatabase
  7. from app import create_app
  8. from app.core.common.identifiers import new_governance_uid
  9. from app.core.data_flow.dataflows import DataFlowService
  10. def test_real_neo4j_concurrent_different_uids_same_name_is_closed():
  11. uri = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_URI")
  12. password = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_PASSWORD")
  13. if not uri or not password:
  14. pytest.skip("real Neo4j acceptance connection is not configured")
  15. user = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_USER", "neo4j")
  16. app = create_app()
  17. app.config.update(
  18. TESTING=True,
  19. NEO4J_URI=uri,
  20. NEO4J_USER=user,
  21. NEO4J_PASSWORD=password,
  22. NEO4J_ENCRYPTED=False,
  23. )
  24. name = f"同名并发生产线-{new_governance_uid()}"
  25. uids = [new_governance_uid(), new_governance_uid()]
  26. start = Barrier(2)
  27. driver = GraphDatabase.driver(
  28. uri, auth=(user, password), encrypted=False
  29. )
  30. def create(uid):
  31. start.wait()
  32. try:
  33. with app.app_context():
  34. DataFlowService._merge_governed_dataflow(
  35. {
  36. "uid": uid,
  37. "name_zh": name,
  38. "name_en": f"concurrent_{uid.replace('-', '')}",
  39. "script_type": "governed",
  40. "script_requirement": "{}",
  41. "script_path": "",
  42. }
  43. )
  44. return "created"
  45. except ValueError as exc:
  46. return str(exc)
  47. try:
  48. with driver.session() as session:
  49. session.run(
  50. "MATCH (n:DataFlow {name_zh: $name}) DETACH DELETE n",
  51. {"name": name},
  52. ).consume()
  53. with ThreadPoolExecutor(max_workers=2) as pool:
  54. outcomes = list(pool.map(create, uids))
  55. assert sorted(outcomes) == ["created", "dataflow_uid_conflict"]
  56. with driver.session() as session:
  57. count = session.run(
  58. "MATCH (n:DataFlow {name_zh: $name}) "
  59. "RETURN count(n) AS count",
  60. {"name": name},
  61. ).single()["count"]
  62. constraints = session.run(
  63. "SHOW CONSTRAINTS YIELD name "
  64. "WHERE name IN ['data_flow_uid', 'data_flow_name_zh'] "
  65. "RETURN collect(name) AS names"
  66. ).single()["names"]
  67. assert count == 1
  68. assert set(constraints) == {"data_flow_uid", "data_flow_name_zh"}
  69. finally:
  70. with driver.session() as session:
  71. session.run(
  72. "MATCH (n:DataFlow {name_zh: $name}) DETACH DELETE n",
  73. {"name": name},
  74. ).consume()
  75. driver.close()
  76. def test_real_neo4j_uid_constraint_merge_replay_and_conflict():
  77. uri = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_URI")
  78. password = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_PASSWORD")
  79. if not uri or not password:
  80. pytest.skip("real Neo4j acceptance connection is not configured")
  81. user = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_USER", "neo4j")
  82. app = create_app()
  83. app.config.update(
  84. TESTING=True,
  85. NEO4J_URI=uri,
  86. NEO4J_USER=user,
  87. NEO4J_PASSWORD=password,
  88. NEO4J_ENCRYPTED=False,
  89. )
  90. uid = new_governance_uid()
  91. node = {
  92. "uid": uid,
  93. "name_zh": f"Saga验收-{uid}",
  94. "name_en": f"saga-{uid}",
  95. "script_type": "governed",
  96. "script_requirement": '{"dataflow_spec":{"schema_version":"2.0"}}',
  97. "script_path": "",
  98. }
  99. driver = GraphDatabase.driver(uri, auth=(user, password), encrypted=False)
  100. try:
  101. with app.app_context():
  102. first_id, first = DataFlowService._merge_governed_dataflow(node)
  103. second_id, second = DataFlowService._merge_governed_dataflow(node)
  104. assert first_id == second_id
  105. assert first == second
  106. with pytest.raises(ValueError, match="dataflow_uid_conflict"):
  107. DataFlowService._merge_governed_dataflow(
  108. {**node, "name_zh": f"篡改-{uid}"}
  109. )
  110. with driver.session() as session:
  111. count = session.run(
  112. "MATCH (n:DataFlow {uid: $uid}) RETURN count(n) AS count",
  113. {"uid": uid},
  114. ).single()["count"]
  115. constraints = [
  116. record["name"]
  117. for record in session.run(
  118. "SHOW CONSTRAINTS YIELD name "
  119. "WHERE name IN ['data_flow_uid', 'data_flow_name_zh'] "
  120. "RETURN name"
  121. )
  122. ]
  123. assert count == 1
  124. assert set(constraints) == {"data_flow_uid", "data_flow_name_zh"}
  125. finally:
  126. with driver.session() as session:
  127. session.run("MATCH (n:DataFlow {uid: $uid}) DETACH DELETE n", {"uid": uid})
  128. driver.close()