routes.py 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719
  1. """HTTP boundary for secret-free external data-source management."""
  2. import logging
  3. from flask import g, jsonify, request
  4. from app.api.data_source import bp
  5. from app.core.connectors.errors import (
  6. ConnectorAuthenticationError,
  7. ConnectorConfigurationError,
  8. ConnectorError,
  9. )
  10. from app.core.data_source.errors import DataSourceError
  11. from app.core.data_source.redaction import (
  12. redact_mapping,
  13. sanitize_exception,
  14. )
  15. from app.models.result import failed, success
  16. logger = logging.getLogger(__name__)
  17. def get_data_source_service():
  18. from app.core.data_source.runtime import get_data_source_manager
  19. from app.core.data_source.service import build_data_source_service
  20. return build_data_source_service(get_data_source_manager())
  21. def _actor_uid():
  22. identity = getattr(g, "current_user", {}) or {}
  23. return identity.get("id") or identity.get("sub")
  24. def _error_response(error):
  25. if isinstance(error, ConnectorError):
  26. logger.warning("连接器操作失败: category=%s", error.category)
  27. return jsonify(
  28. failed(
  29. "连接器操作失败",
  30. code=error.http_status,
  31. error={
  32. "code": "CONNECTOR_ERROR",
  33. "category": error.category,
  34. "retryable": error.retryable,
  35. },
  36. )
  37. ), error.http_status
  38. if isinstance(error, DataSourceError):
  39. logger.warning(
  40. "数据源操作失败: code=%s message=%s",
  41. error.code,
  42. sanitize_exception(error),
  43. )
  44. return (
  45. jsonify(
  46. failed(
  47. str(error),
  48. code=error.http_status,
  49. error={"code": error.code},
  50. )
  51. ),
  52. error.http_status,
  53. )
  54. logger.error(
  55. "数据源操作异常: %s",
  56. sanitize_exception(error),
  57. )
  58. return (
  59. jsonify(
  60. failed(
  61. "数据源操作失败",
  62. code=500,
  63. error={"code": "DATASOURCE_ERROR"},
  64. )
  65. ),
  66. 500,
  67. )
  68. @bp.route("/save", methods=["POST"])
  69. def data_source_save():
  70. payload = request.get_json(silent=True) or {}
  71. logger.debug("保存数据源请求: %s", redact_mapping(payload))
  72. try:
  73. service = get_data_source_service()
  74. definition, created = service.save(
  75. payload,
  76. actor_uid=_actor_uid(),
  77. )
  78. status = 201 if created else 200
  79. return jsonify(success(service.serialize(definition))), status
  80. except Exception as error:
  81. return _error_response(error)
  82. @bp.route("/list", methods=["POST"])
  83. def data_source_list():
  84. payload = request.get_json(silent=True) or {}
  85. try:
  86. service = get_data_source_service()
  87. definitions = service.list(payload)
  88. items = [service.serialize(item) for item in definitions]
  89. return jsonify(success({"data_source": items, "total": len(items)})), 200
  90. except Exception as error:
  91. return _error_response(error)
  92. @bp.route("/delete", methods=["POST"])
  93. def data_source_delete():
  94. payload = request.get_json(silent=True) or {}
  95. logger.debug("删除数据源请求: %s", redact_mapping(payload))
  96. try:
  97. result = get_data_source_service().delete(
  98. payload.get("uid"),
  99. actor_uid=_actor_uid(),
  100. )
  101. return jsonify(success(result)), 200
  102. except Exception as error:
  103. return _error_response(error)
  104. @bp.route("/conntest", methods=["POST"])
  105. def data_source_conn_test():
  106. payload = request.get_json(silent=True) or {}
  107. logger.debug("测试数据源连接请求: %s", redact_mapping(payload))
  108. try:
  109. result = get_data_source_service().test_connection(payload)
  110. return jsonify(success(result)), 200
  111. except Exception as error:
  112. return _error_response(error)
  113. @bp.route("/valid", methods=["POST"])
  114. def data_source_connstr_valid():
  115. payload = request.get_json(silent=True) or {}
  116. logger.debug("验证数据源连接请求: %s", redact_mapping(payload))
  117. try:
  118. result = get_data_source_service().test_connection(payload)
  119. return jsonify(success({"exists": False, **result})), 200
  120. except Exception as error:
  121. return _error_response(error)
  122. @bp.route("/parse", methods=["POST"])
  123. def data_source_connstr_parse():
  124. return (
  125. jsonify(
  126. failed(
  127. "连接字符串快捷解析已停用",
  128. code=410,
  129. error={"code": "DATASOURCE_PARSE_RETIRED"},
  130. )
  131. ),
  132. 410,
  133. )
  134. @bp.route("/pools", methods=["GET"])
  135. def data_source_pool_list():
  136. try:
  137. service = get_data_source_service()
  138. items = [
  139. service.serialize_pool_status(status) for status in service.pool_statuses()
  140. ]
  141. return jsonify(success({"pools": items, "total": len(items)})), 200
  142. except Exception as error:
  143. return _error_response(error)
  144. @bp.route("/<data_source_uid>/pool", methods=["GET"])
  145. def data_source_pool_status(data_source_uid):
  146. try:
  147. service = get_data_source_service()
  148. status = service.pool_status(data_source_uid)
  149. serializer = getattr(
  150. service,
  151. "serialize_pool_status",
  152. None,
  153. )
  154. data = (
  155. serializer(status)
  156. if serializer is not None
  157. else {
  158. "data_source_uid": status.data_source_uid,
  159. "credential_version": status.credential_version,
  160. "pool_state": status.pool_state,
  161. "pool_size": status.pool_size,
  162. "checked_out": status.checked_out,
  163. "checked_in": status.checked_in,
  164. "overflow": status.overflow,
  165. "leases": status.leases,
  166. }
  167. )
  168. return jsonify(success(data)), 200
  169. except Exception as error:
  170. return _error_response(error)
  171. @bp.route("/<data_source_uid>/pool/invalidate", methods=["POST"])
  172. def data_source_pool_invalidate(data_source_uid):
  173. payload = request.get_json(silent=True) or {}
  174. reason = payload.get("reason")
  175. if reason not in {
  176. "admin_reset",
  177. "configuration_changed",
  178. "credential_rotated",
  179. }:
  180. return (
  181. jsonify(
  182. failed(
  183. "连接池失效原因无效",
  184. code=400,
  185. error={"code": "DATASOURCE_CONFIGURATION_INVALID"},
  186. )
  187. ),
  188. 400,
  189. )
  190. try:
  191. result = get_data_source_service().invalidate_pool(
  192. data_source_uid,
  193. reason=reason,
  194. actor_uid=_actor_uid(),
  195. )
  196. return jsonify(success(result)), 200
  197. except Exception as error:
  198. return _error_response(error)
  199. @bp.route("/graph", methods=["POST"])
  200. def data_source_graph_relationship():
  201. from app import db
  202. from app.core.connectors.repository import ConnectorRepository
  203. payload = request.get_json(silent=True) or {}
  204. try:
  205. graph = ConnectorRepository(db.session).graph(
  206. source_uid=payload.get("source_uid"),
  207. business_domain_uid=payload.get("business_domain_uid"),
  208. process_key=payload.get("process_key"),
  209. run_uid=payload.get("run_uid"),
  210. limit=payload.get("limit", 500),
  211. )
  212. return jsonify(success(graph)), 200
  213. except Exception as error:
  214. return _error_response(error)
  215. def _connector_registry():
  216. from flask import current_app
  217. from app.core.connectors.builtin import register_builtin_connectors
  218. from app.core.connectors.registry import ConnectorRegistry
  219. from app.core.connectors.secrets import EnvironmentSecretResolver
  220. from app.core.data_source.runtime import get_data_source_manager
  221. registry = current_app.extensions.get("connector_registry")
  222. if registry is not None:
  223. return registry
  224. registry = ConnectorRegistry()
  225. register_builtin_connectors(
  226. registry,
  227. connection_provider=get_data_source_manager().connect,
  228. secret_resolver=current_app.config.get("CONNECTOR_SECRET_RESOLVER")
  229. or EnvironmentSecretResolver(),
  230. )
  231. current_app.extensions["connector_registry"] = registry
  232. return registry
  233. @bp.route("/connectors/manifests", methods=["GET"])
  234. def connector_manifests():
  235. try:
  236. return jsonify(success({"manifests": _connector_registry().manifests()})), 200
  237. except Exception as error:
  238. return _error_response(error)
  239. @bp.route("/connectors/config/validate", methods=["POST"])
  240. def connector_config_validate():
  241. payload = request.get_json(silent=True) or {}
  242. try:
  243. config = _connector_registry().validate(
  244. payload.get("connector_id"), payload.get("version"), payload.get("config")
  245. )
  246. return jsonify(success({"valid": True, "config_keys": sorted(config)})), 200
  247. except Exception as error:
  248. return _error_response(error)
  249. @bp.route("/connectors/<connector_id>/<version>/health", methods=["POST"])
  250. def connector_health(connector_id, version):
  251. payload = request.get_json(silent=True) or {}
  252. try:
  253. connector = _connector_registry().resolve(connector_id, version)
  254. return jsonify(
  255. success(vars(connector.health(payload.get("config") or {})))
  256. ), 200
  257. except Exception as error:
  258. return _error_response(error)
  259. @bp.route("/connectors/<connector_id>/<version>/compatibility", methods=["GET"])
  260. def connector_compatibility(connector_id, version):
  261. try:
  262. connector = _connector_registry().resolve(connector_id, version)
  263. return jsonify(success(vars(connector.compatibility()))), 200
  264. except Exception as error:
  265. return _error_response(error)
  266. @bp.route("/connectors/runs", methods=["GET"])
  267. def connector_runs_list():
  268. from app import db
  269. from app.core.connectors.repository import ConnectorRepository
  270. try:
  271. repository = ConnectorRepository(db.session, _actor_uid())
  272. return jsonify(
  273. success({"runs": repository.list(request.args.get("limit", 100))})
  274. ), 200
  275. except Exception as error:
  276. db.session.rollback()
  277. return _error_response(error)
  278. def _require_human_dry_run(payload):
  279. if payload.get("dry_run") is not True:
  280. raise ConnectorConfigurationError("human connector runs must use dry_run=true")
  281. def _require_collection_operation(payload):
  282. if payload.get("operation") not in {
  283. "discover",
  284. "snapshot",
  285. "incremental",
  286. "lineage",
  287. "profile",
  288. "evidence",
  289. }:
  290. raise ConnectorConfigurationError("connector collection operation is invalid")
  291. @bp.route("/connectors/runs", methods=["POST"])
  292. def connector_run_create():
  293. from app import db
  294. from app.core.connectors.repository import ConnectorRepository
  295. try:
  296. repository = ConnectorRepository(
  297. db.session, _actor_uid(), run_access="human"
  298. )
  299. from dataclasses import asdict
  300. from app.core.connectors.runtime import ConnectorRuntime
  301. from app.core.connectors.sdk import OperationRequest
  302. payload = request.get_json(silent=True) or {}
  303. _require_human_dry_run(payload)
  304. _require_collection_operation(payload)
  305. registry = _connector_registry()
  306. for manifest in (
  307. registry.resolve(
  308. payload.get("connector_id"), payload.get("version")
  309. ).manifest,
  310. ):
  311. repository.register_manifest(manifest, _actor_uid())
  312. db.session.commit()
  313. operation_request = OperationRequest(
  314. source_uid=str(payload.get("source_uid") or ""),
  315. operation=str(payload.get("operation") or ""),
  316. config=payload.get("config") or {},
  317. scope=payload.get("scope") or {},
  318. cursor=payload.get("cursor") or {},
  319. checkpoint=payload.get("checkpoint") or {},
  320. idempotency_key=payload.get("idempotency_key"),
  321. dry_run=bool(payload.get("dry_run", False)),
  322. process_key=str(payload.get("process_key") or "dry-run"),
  323. )
  324. result = ConnectorRuntime(registry, store=repository).execute(
  325. payload.get("connector_id"), payload.get("version"), operation_request
  326. )
  327. return jsonify(success(asdict(result), code=201)), 201
  328. except Exception as error:
  329. db.session.rollback()
  330. return _error_response(error)
  331. @bp.route("/connectors/machine/runs", methods=["POST"])
  332. def connector_machine_run():
  333. """Machine-only run boundary; a human bearer token is never accepted here."""
  334. from app import db
  335. from app.core.connectors.identity import ConnectorIdentityRepository
  336. from app.core.connectors.repository import ConnectorRepository
  337. from app.core.connectors.runtime import ConnectorRuntime
  338. from app.core.connectors.sdk import OperationRequest
  339. payload = request.get_json(silent=True) or {}
  340. token = request.headers.get("X-Connector-Credential", "").strip()
  341. if not token:
  342. return _error_response(
  343. ConnectorAuthenticationError("machine credential is required")
  344. )
  345. try:
  346. for field in (
  347. "connector_id",
  348. "version",
  349. "source_uid",
  350. "business_domain_uid",
  351. "environment",
  352. "process_key",
  353. "operation",
  354. ):
  355. if not str(payload.get(field) or "").strip():
  356. raise ConnectorConfigurationError(
  357. "machine connector binding is incomplete"
  358. )
  359. _require_collection_operation(payload)
  360. if "config" in payload:
  361. raise ConnectorConfigurationError(
  362. "machine connector config is supplied by the approved source binding"
  363. )
  364. identity = ConnectorIdentityRepository(db.session).authenticate(
  365. token,
  366. connector_id=str(payload.get("connector_id") or ""),
  367. connector_version=str(payload.get("version") or ""),
  368. source_uid=str(payload.get("source_uid") or ""),
  369. business_domain_uid=str(payload.get("business_domain_uid") or ""),
  370. environment=str(payload.get("environment") or ""),
  371. operation=str(payload.get("operation") or ""),
  372. scope=payload.get("scope") or {},
  373. )
  374. registry = _connector_registry()
  375. manifest = registry.resolve(
  376. payload.get("connector_id"), payload.get("version")
  377. ).manifest
  378. if not identity.get("source_binding_uid"):
  379. raise ConnectorConfigurationError(
  380. "machine connector principal has no approved source binding"
  381. )
  382. repository = ConnectorRepository(
  383. db.session,
  384. identity["principal_uid"],
  385. run_access="machine",
  386. principal_uid=identity["principal_uid"],
  387. )
  388. repository.register_manifest(manifest, identity["principal_uid"])
  389. db.session.commit()
  390. operation_request = OperationRequest(
  391. source_uid=str(payload.get("source_uid")),
  392. operation=str(payload.get("operation")),
  393. config=identity["approved_config"],
  394. scope=payload.get("scope") or {},
  395. cursor=payload.get("cursor") or {},
  396. checkpoint=payload.get("checkpoint") or {},
  397. idempotency_key=payload.get("idempotency_key"),
  398. dry_run=bool(payload.get("dry_run", False)),
  399. principal_uid=identity["principal_uid"],
  400. business_domain_uid=str(payload.get("business_domain_uid")),
  401. environment=str(payload.get("environment")),
  402. process_key=str(payload.get("process_key") or ""),
  403. source_binding_uid=identity["source_binding_uid"],
  404. source_binding_version=identity["source_binding_version"],
  405. )
  406. result = ConnectorRuntime(registry, store=repository).execute(
  407. payload.get("connector_id"), payload.get("version"), operation_request
  408. )
  409. from dataclasses import asdict
  410. return jsonify(success(asdict(result), code=201)), 201
  411. except Exception as error:
  412. db.session.rollback()
  413. return _error_response(error)
  414. @bp.route("/connectors/runs/<idempotency_key>/cancel", methods=["POST"])
  415. def connector_run_cancel(idempotency_key):
  416. from app import db
  417. from app.core.connectors.repository import ConnectorRepository
  418. from app.core.connectors.runtime import ConnectorRuntime
  419. try:
  420. result = ConnectorRuntime(
  421. _connector_registry(),
  422. store=ConnectorRepository(
  423. db.session, _actor_uid(), run_access="human"
  424. ),
  425. ).cancel(idempotency_key)
  426. return jsonify(success(result)), 200
  427. except Exception as error:
  428. db.session.rollback()
  429. return _error_response(error)
  430. @bp.route("/connectors/runs/<idempotency_key>/resume", methods=["POST"])
  431. def connector_run_resume(idempotency_key):
  432. from app import db
  433. from app.core.connectors.repository import ConnectorRepository
  434. from app.core.connectors.runtime import ConnectorRuntime
  435. try:
  436. result = ConnectorRuntime(
  437. _connector_registry(),
  438. store=ConnectorRepository(
  439. db.session, _actor_uid(), run_access="human"
  440. ),
  441. ).resume(idempotency_key)
  442. return jsonify(success(result)), 200
  443. except Exception as error:
  444. db.session.rollback()
  445. return _error_response(error)
  446. def _machine_run_action(idempotency_key, action):
  447. from app import db
  448. from app.core.connectors.identity import ConnectorIdentityRepository
  449. from app.core.connectors.repository import ConnectorRepository
  450. from app.core.connectors.runtime import ConnectorRuntime
  451. token = request.headers.get("X-Connector-Credential", "").strip()
  452. if not token:
  453. return _error_response(
  454. ConnectorAuthenticationError("machine credential is required")
  455. )
  456. try:
  457. record = ConnectorRepository(db.session).get(idempotency_key)
  458. if not record or not record.get("principal_uid") or record.get("dry_run"):
  459. raise ConnectorConfigurationError("machine connector run was not found")
  460. identity = ConnectorIdentityRepository(db.session).authenticate(
  461. token,
  462. connector_id=record["connector_id"],
  463. connector_version=record["connector_version"],
  464. source_uid=record["source_uid"],
  465. business_domain_uid=record["business_domain_uid"],
  466. environment=record["environment"],
  467. operation=action,
  468. scope=record.get("scope") or {},
  469. )
  470. if (
  471. identity["principal_uid"] != record["principal_uid"]
  472. or identity.get("source_binding_uid") != record.get("source_binding_uid")
  473. or identity.get("source_binding_version")
  474. != record.get("source_binding_version")
  475. ):
  476. raise ConnectorAuthenticationError(
  477. "machine credential run binding was rejected"
  478. )
  479. repository = ConnectorRepository(
  480. db.session,
  481. identity["principal_uid"],
  482. run_access="machine",
  483. principal_uid=identity["principal_uid"],
  484. )
  485. runtime = ConnectorRuntime(_connector_registry(), store=repository)
  486. result = getattr(runtime, action)(idempotency_key)
  487. if action == "resume":
  488. from dataclasses import asdict
  489. result = asdict(result)
  490. return jsonify(success(result)), 200
  491. except Exception as error:
  492. db.session.rollback()
  493. return _error_response(error)
  494. @bp.route(
  495. "/connectors/machine/runs/<idempotency_key>/cancel", methods=["POST"]
  496. )
  497. def connector_machine_run_cancel(idempotency_key):
  498. """Cancel one bound machine run with a fresh one-time credential."""
  499. return _machine_run_action(idempotency_key, "cancel")
  500. @bp.route(
  501. "/connectors/machine/runs/<idempotency_key>/resume", methods=["POST"]
  502. )
  503. def connector_machine_run_resume(idempotency_key):
  504. """Resume one bound machine run with a fresh one-time credential."""
  505. return _machine_run_action(idempotency_key, "resume")
  506. @bp.route("/connectors/source-bindings", methods=["GET"])
  507. def connector_source_bindings_list():
  508. from app import db
  509. from app.core.connectors.bindings import ConnectorSourceBindingRepository
  510. try:
  511. items = ConnectorSourceBindingRepository(db.session).list_public(
  512. request.args.get("limit", 100)
  513. )
  514. return jsonify(success({"source_bindings": items})), 200
  515. except Exception as error:
  516. db.session.rollback()
  517. return _error_response(error)
  518. @bp.route("/connectors/source-bindings", methods=["POST"])
  519. def connector_source_binding_approve():
  520. from app import db
  521. from app.core.connectors.bindings import ConnectorSourceBindingRepository
  522. from app.core.connectors.repository import ConnectorRepository
  523. payload = request.get_json(silent=True) or {}
  524. try:
  525. registry = _connector_registry()
  526. manifest = registry.resolve(
  527. payload.get("connector_id"), payload.get("version")
  528. ).manifest
  529. approved_config = registry.validate(
  530. payload.get("connector_id"),
  531. payload.get("version"),
  532. payload.get("approved_config") or {},
  533. )
  534. ConnectorRepository(db.session, _actor_uid()).register_manifest(
  535. manifest, _actor_uid()
  536. )
  537. result = ConnectorSourceBindingRepository(db.session).approve(
  538. connector_id=payload.get("connector_id"),
  539. connector_version=payload.get("version"),
  540. source_uid=payload.get("source_uid"),
  541. business_domain_uid=payload.get("business_domain_uid"),
  542. environment=payload.get("environment"),
  543. approved_config=approved_config,
  544. approved_by=_actor_uid(),
  545. binding_uid=payload.get("binding_uid"),
  546. )
  547. return jsonify(success(result, code=201)), 201
  548. except Exception as error:
  549. db.session.rollback()
  550. return _error_response(error)
  551. @bp.route("/connectors/source-bindings/<binding_uid>/revoke", methods=["POST"])
  552. def connector_source_binding_revoke(binding_uid):
  553. from app import db
  554. from app.core.connectors.bindings import ConnectorSourceBindingRepository
  555. try:
  556. revoked = ConnectorSourceBindingRepository(db.session).revoke(
  557. binding_uid, _actor_uid()
  558. )
  559. return jsonify(success(revoked)), 200
  560. except Exception as error:
  561. db.session.rollback()
  562. return _error_response(error)
  563. @bp.route("/connectors/principals", methods=["POST"])
  564. def connector_principal_create():
  565. from app import db
  566. from app.core.connectors.identity import ConnectorIdentityRepository
  567. payload = request.get_json(silent=True) or {}
  568. try:
  569. _connector_registry().resolve(
  570. payload.get("connector_id"), payload.get("version", "1.0.0")
  571. )
  572. uid = ConnectorIdentityRepository(db.session).create_principal(
  573. connector_id=payload.get("connector_id"),
  574. connector_version=payload.get("version", "1.0.0"),
  575. source_uid=payload.get("source_uid"),
  576. business_domain_uid=payload.get("business_domain_uid"),
  577. environment=payload.get("environment"),
  578. operations=payload.get("operations") or (),
  579. scopes=payload.get("scopes") or {},
  580. actor_uid=_actor_uid(),
  581. source_binding_uid=payload.get("source_binding_uid"),
  582. source_binding_version=payload.get("source_binding_version"),
  583. )
  584. return jsonify(success({"principal_uid": uid}, code=201)), 201
  585. except Exception as error:
  586. db.session.rollback()
  587. return _error_response(error)
  588. @bp.route("/connectors/principals/<principal_uid>/credentials", methods=["POST"])
  589. def connector_credential_issue(principal_uid):
  590. from app import db
  591. from app.core.connectors.identity import ConnectorIdentityRepository
  592. payload = request.get_json(silent=True) or {}
  593. try:
  594. result = ConnectorIdentityRepository(db.session).issue(
  595. principal_uid,
  596. ttl_seconds=payload.get("ttl_seconds", 900),
  597. actor_uid=_actor_uid(),
  598. )
  599. response = jsonify(success(result, code=201))
  600. response.headers["Cache-Control"] = "no-store"
  601. return response, 201
  602. except Exception as error:
  603. db.session.rollback()
  604. return _error_response(error)
  605. @bp.route("/connectors/credentials/<credential_uid>/<action>", methods=["POST"])
  606. def connector_credential_action(credential_uid, action):
  607. from app import db
  608. from app.core.connectors.identity import ConnectorIdentityRepository
  609. payload = request.get_json(silent=True) or {}
  610. try:
  611. repository = ConnectorIdentityRepository(db.session)
  612. if action == "revoke":
  613. result = {"revoked": repository.revoke(credential_uid, _actor_uid())}
  614. elif action == "rotate":
  615. result = repository.rotate(
  616. credential_uid,
  617. ttl_seconds=payload.get("ttl_seconds", 900),
  618. actor_uid=_actor_uid(),
  619. )
  620. else:
  621. raise ConnectorConfigurationError("credential action is invalid")
  622. response = jsonify(success(result))
  623. response.headers["Cache-Control"] = "no-store"
  624. return response, 200
  625. except Exception as error:
  626. db.session.rollback()
  627. return _error_response(error)