| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719 |
- """HTTP boundary for secret-free external data-source management."""
- import logging
- from flask import g, jsonify, request
- from app.api.data_source import bp
- from app.core.connectors.errors import (
- ConnectorAuthenticationError,
- ConnectorConfigurationError,
- ConnectorError,
- )
- from app.core.data_source.errors import DataSourceError
- from app.core.data_source.redaction import (
- redact_mapping,
- sanitize_exception,
- )
- from app.models.result import failed, success
- logger = logging.getLogger(__name__)
- def get_data_source_service():
- from app.core.data_source.runtime import get_data_source_manager
- from app.core.data_source.service import build_data_source_service
- return build_data_source_service(get_data_source_manager())
- def _actor_uid():
- identity = getattr(g, "current_user", {}) or {}
- return identity.get("id") or identity.get("sub")
- def _error_response(error):
- if isinstance(error, ConnectorError):
- logger.warning("连接器操作失败: category=%s", error.category)
- return jsonify(
- failed(
- "连接器操作失败",
- code=error.http_status,
- error={
- "code": "CONNECTOR_ERROR",
- "category": error.category,
- "retryable": error.retryable,
- },
- )
- ), error.http_status
- if isinstance(error, DataSourceError):
- logger.warning(
- "数据源操作失败: code=%s message=%s",
- error.code,
- sanitize_exception(error),
- )
- return (
- jsonify(
- failed(
- str(error),
- code=error.http_status,
- error={"code": error.code},
- )
- ),
- error.http_status,
- )
- logger.error(
- "数据源操作异常: %s",
- sanitize_exception(error),
- )
- return (
- jsonify(
- failed(
- "数据源操作失败",
- code=500,
- error={"code": "DATASOURCE_ERROR"},
- )
- ),
- 500,
- )
- @bp.route("/save", methods=["POST"])
- def data_source_save():
- payload = request.get_json(silent=True) or {}
- logger.debug("保存数据源请求: %s", redact_mapping(payload))
- try:
- service = get_data_source_service()
- definition, created = service.save(
- payload,
- actor_uid=_actor_uid(),
- )
- status = 201 if created else 200
- return jsonify(success(service.serialize(definition))), status
- except Exception as error:
- return _error_response(error)
- @bp.route("/list", methods=["POST"])
- def data_source_list():
- payload = request.get_json(silent=True) or {}
- try:
- service = get_data_source_service()
- definitions = service.list(payload)
- items = [service.serialize(item) for item in definitions]
- return jsonify(success({"data_source": items, "total": len(items)})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/delete", methods=["POST"])
- def data_source_delete():
- payload = request.get_json(silent=True) or {}
- logger.debug("删除数据源请求: %s", redact_mapping(payload))
- try:
- result = get_data_source_service().delete(
- payload.get("uid"),
- actor_uid=_actor_uid(),
- )
- return jsonify(success(result)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/conntest", methods=["POST"])
- def data_source_conn_test():
- payload = request.get_json(silent=True) or {}
- logger.debug("测试数据源连接请求: %s", redact_mapping(payload))
- try:
- result = get_data_source_service().test_connection(payload)
- return jsonify(success(result)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/valid", methods=["POST"])
- def data_source_connstr_valid():
- payload = request.get_json(silent=True) or {}
- logger.debug("验证数据源连接请求: %s", redact_mapping(payload))
- try:
- result = get_data_source_service().test_connection(payload)
- return jsonify(success({"exists": False, **result})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/parse", methods=["POST"])
- def data_source_connstr_parse():
- return (
- jsonify(
- failed(
- "连接字符串快捷解析已停用",
- code=410,
- error={"code": "DATASOURCE_PARSE_RETIRED"},
- )
- ),
- 410,
- )
- @bp.route("/pools", methods=["GET"])
- def data_source_pool_list():
- try:
- service = get_data_source_service()
- items = [
- service.serialize_pool_status(status) for status in service.pool_statuses()
- ]
- return jsonify(success({"pools": items, "total": len(items)})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/<data_source_uid>/pool", methods=["GET"])
- def data_source_pool_status(data_source_uid):
- try:
- service = get_data_source_service()
- status = service.pool_status(data_source_uid)
- serializer = getattr(
- service,
- "serialize_pool_status",
- None,
- )
- data = (
- serializer(status)
- if serializer is not None
- else {
- "data_source_uid": status.data_source_uid,
- "credential_version": status.credential_version,
- "pool_state": status.pool_state,
- "pool_size": status.pool_size,
- "checked_out": status.checked_out,
- "checked_in": status.checked_in,
- "overflow": status.overflow,
- "leases": status.leases,
- }
- )
- return jsonify(success(data)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/<data_source_uid>/pool/invalidate", methods=["POST"])
- def data_source_pool_invalidate(data_source_uid):
- payload = request.get_json(silent=True) or {}
- reason = payload.get("reason")
- if reason not in {
- "admin_reset",
- "configuration_changed",
- "credential_rotated",
- }:
- return (
- jsonify(
- failed(
- "连接池失效原因无效",
- code=400,
- error={"code": "DATASOURCE_CONFIGURATION_INVALID"},
- )
- ),
- 400,
- )
- try:
- result = get_data_source_service().invalidate_pool(
- data_source_uid,
- reason=reason,
- actor_uid=_actor_uid(),
- )
- return jsonify(success(result)), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/graph", methods=["POST"])
- def data_source_graph_relationship():
- from app import db
- from app.core.connectors.repository import ConnectorRepository
- payload = request.get_json(silent=True) or {}
- try:
- graph = ConnectorRepository(db.session).graph(
- source_uid=payload.get("source_uid"),
- business_domain_uid=payload.get("business_domain_uid"),
- process_key=payload.get("process_key"),
- run_uid=payload.get("run_uid"),
- limit=payload.get("limit", 500),
- )
- return jsonify(success(graph)), 200
- except Exception as error:
- return _error_response(error)
- def _connector_registry():
- from flask import current_app
- from app.core.connectors.builtin import register_builtin_connectors
- from app.core.connectors.registry import ConnectorRegistry
- from app.core.connectors.secrets import EnvironmentSecretResolver
- from app.core.data_source.runtime import get_data_source_manager
- registry = current_app.extensions.get("connector_registry")
- if registry is not None:
- return registry
- registry = ConnectorRegistry()
- register_builtin_connectors(
- registry,
- connection_provider=get_data_source_manager().connect,
- secret_resolver=current_app.config.get("CONNECTOR_SECRET_RESOLVER")
- or EnvironmentSecretResolver(),
- )
- current_app.extensions["connector_registry"] = registry
- return registry
- @bp.route("/connectors/manifests", methods=["GET"])
- def connector_manifests():
- try:
- return jsonify(success({"manifests": _connector_registry().manifests()})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/connectors/config/validate", methods=["POST"])
- def connector_config_validate():
- payload = request.get_json(silent=True) or {}
- try:
- config = _connector_registry().validate(
- payload.get("connector_id"), payload.get("version"), payload.get("config")
- )
- return jsonify(success({"valid": True, "config_keys": sorted(config)})), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/connectors/<connector_id>/<version>/health", methods=["POST"])
- def connector_health(connector_id, version):
- payload = request.get_json(silent=True) or {}
- try:
- connector = _connector_registry().resolve(connector_id, version)
- return jsonify(
- success(vars(connector.health(payload.get("config") or {})))
- ), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/connectors/<connector_id>/<version>/compatibility", methods=["GET"])
- def connector_compatibility(connector_id, version):
- try:
- connector = _connector_registry().resolve(connector_id, version)
- return jsonify(success(vars(connector.compatibility()))), 200
- except Exception as error:
- return _error_response(error)
- @bp.route("/connectors/runs", methods=["GET"])
- def connector_runs_list():
- from app import db
- from app.core.connectors.repository import ConnectorRepository
- try:
- repository = ConnectorRepository(db.session, _actor_uid())
- return jsonify(
- success({"runs": repository.list(request.args.get("limit", 100))})
- ), 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- def _require_human_dry_run(payload):
- if payload.get("dry_run") is not True:
- raise ConnectorConfigurationError("human connector runs must use dry_run=true")
- def _require_collection_operation(payload):
- if payload.get("operation") not in {
- "discover",
- "snapshot",
- "incremental",
- "lineage",
- "profile",
- "evidence",
- }:
- raise ConnectorConfigurationError("connector collection operation is invalid")
- @bp.route("/connectors/runs", methods=["POST"])
- def connector_run_create():
- from app import db
- from app.core.connectors.repository import ConnectorRepository
- try:
- repository = ConnectorRepository(
- db.session, _actor_uid(), run_access="human"
- )
- from dataclasses import asdict
- from app.core.connectors.runtime import ConnectorRuntime
- from app.core.connectors.sdk import OperationRequest
- payload = request.get_json(silent=True) or {}
- _require_human_dry_run(payload)
- _require_collection_operation(payload)
- registry = _connector_registry()
- for manifest in (
- registry.resolve(
- payload.get("connector_id"), payload.get("version")
- ).manifest,
- ):
- repository.register_manifest(manifest, _actor_uid())
- db.session.commit()
- operation_request = OperationRequest(
- source_uid=str(payload.get("source_uid") or ""),
- operation=str(payload.get("operation") or ""),
- config=payload.get("config") or {},
- scope=payload.get("scope") or {},
- cursor=payload.get("cursor") or {},
- checkpoint=payload.get("checkpoint") or {},
- idempotency_key=payload.get("idempotency_key"),
- dry_run=bool(payload.get("dry_run", False)),
- process_key=str(payload.get("process_key") or "dry-run"),
- )
- result = ConnectorRuntime(registry, store=repository).execute(
- payload.get("connector_id"), payload.get("version"), operation_request
- )
- return jsonify(success(asdict(result), code=201)), 201
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/machine/runs", methods=["POST"])
- def connector_machine_run():
- """Machine-only run boundary; a human bearer token is never accepted here."""
- from app import db
- from app.core.connectors.identity import ConnectorIdentityRepository
- from app.core.connectors.repository import ConnectorRepository
- from app.core.connectors.runtime import ConnectorRuntime
- from app.core.connectors.sdk import OperationRequest
- payload = request.get_json(silent=True) or {}
- token = request.headers.get("X-Connector-Credential", "").strip()
- if not token:
- return _error_response(
- ConnectorAuthenticationError("machine credential is required")
- )
- try:
- for field in (
- "connector_id",
- "version",
- "source_uid",
- "business_domain_uid",
- "environment",
- "process_key",
- "operation",
- ):
- if not str(payload.get(field) or "").strip():
- raise ConnectorConfigurationError(
- "machine connector binding is incomplete"
- )
- _require_collection_operation(payload)
- if "config" in payload:
- raise ConnectorConfigurationError(
- "machine connector config is supplied by the approved source binding"
- )
- identity = ConnectorIdentityRepository(db.session).authenticate(
- token,
- connector_id=str(payload.get("connector_id") or ""),
- connector_version=str(payload.get("version") or ""),
- source_uid=str(payload.get("source_uid") or ""),
- business_domain_uid=str(payload.get("business_domain_uid") or ""),
- environment=str(payload.get("environment") or ""),
- operation=str(payload.get("operation") or ""),
- scope=payload.get("scope") or {},
- )
- registry = _connector_registry()
- manifest = registry.resolve(
- payload.get("connector_id"), payload.get("version")
- ).manifest
- if not identity.get("source_binding_uid"):
- raise ConnectorConfigurationError(
- "machine connector principal has no approved source binding"
- )
- repository = ConnectorRepository(
- db.session,
- identity["principal_uid"],
- run_access="machine",
- principal_uid=identity["principal_uid"],
- )
- repository.register_manifest(manifest, identity["principal_uid"])
- db.session.commit()
- operation_request = OperationRequest(
- source_uid=str(payload.get("source_uid")),
- operation=str(payload.get("operation")),
- config=identity["approved_config"],
- scope=payload.get("scope") or {},
- cursor=payload.get("cursor") or {},
- checkpoint=payload.get("checkpoint") or {},
- idempotency_key=payload.get("idempotency_key"),
- dry_run=bool(payload.get("dry_run", False)),
- principal_uid=identity["principal_uid"],
- business_domain_uid=str(payload.get("business_domain_uid")),
- environment=str(payload.get("environment")),
- process_key=str(payload.get("process_key") or ""),
- source_binding_uid=identity["source_binding_uid"],
- source_binding_version=identity["source_binding_version"],
- )
- result = ConnectorRuntime(registry, store=repository).execute(
- payload.get("connector_id"), payload.get("version"), operation_request
- )
- from dataclasses import asdict
- return jsonify(success(asdict(result), code=201)), 201
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/runs/<idempotency_key>/cancel", methods=["POST"])
- def connector_run_cancel(idempotency_key):
- from app import db
- from app.core.connectors.repository import ConnectorRepository
- from app.core.connectors.runtime import ConnectorRuntime
- try:
- result = ConnectorRuntime(
- _connector_registry(),
- store=ConnectorRepository(
- db.session, _actor_uid(), run_access="human"
- ),
- ).cancel(idempotency_key)
- return jsonify(success(result)), 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/runs/<idempotency_key>/resume", methods=["POST"])
- def connector_run_resume(idempotency_key):
- from app import db
- from app.core.connectors.repository import ConnectorRepository
- from app.core.connectors.runtime import ConnectorRuntime
- try:
- result = ConnectorRuntime(
- _connector_registry(),
- store=ConnectorRepository(
- db.session, _actor_uid(), run_access="human"
- ),
- ).resume(idempotency_key)
- return jsonify(success(result)), 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- def _machine_run_action(idempotency_key, action):
- from app import db
- from app.core.connectors.identity import ConnectorIdentityRepository
- from app.core.connectors.repository import ConnectorRepository
- from app.core.connectors.runtime import ConnectorRuntime
- token = request.headers.get("X-Connector-Credential", "").strip()
- if not token:
- return _error_response(
- ConnectorAuthenticationError("machine credential is required")
- )
- try:
- record = ConnectorRepository(db.session).get(idempotency_key)
- if not record or not record.get("principal_uid") or record.get("dry_run"):
- raise ConnectorConfigurationError("machine connector run was not found")
- identity = ConnectorIdentityRepository(db.session).authenticate(
- token,
- connector_id=record["connector_id"],
- connector_version=record["connector_version"],
- source_uid=record["source_uid"],
- business_domain_uid=record["business_domain_uid"],
- environment=record["environment"],
- operation=action,
- scope=record.get("scope") or {},
- )
- if (
- identity["principal_uid"] != record["principal_uid"]
- or identity.get("source_binding_uid") != record.get("source_binding_uid")
- or identity.get("source_binding_version")
- != record.get("source_binding_version")
- ):
- raise ConnectorAuthenticationError(
- "machine credential run binding was rejected"
- )
- repository = ConnectorRepository(
- db.session,
- identity["principal_uid"],
- run_access="machine",
- principal_uid=identity["principal_uid"],
- )
- runtime = ConnectorRuntime(_connector_registry(), store=repository)
- result = getattr(runtime, action)(idempotency_key)
- if action == "resume":
- from dataclasses import asdict
- result = asdict(result)
- return jsonify(success(result)), 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route(
- "/connectors/machine/runs/<idempotency_key>/cancel", methods=["POST"]
- )
- def connector_machine_run_cancel(idempotency_key):
- """Cancel one bound machine run with a fresh one-time credential."""
- return _machine_run_action(idempotency_key, "cancel")
- @bp.route(
- "/connectors/machine/runs/<idempotency_key>/resume", methods=["POST"]
- )
- def connector_machine_run_resume(idempotency_key):
- """Resume one bound machine run with a fresh one-time credential."""
- return _machine_run_action(idempotency_key, "resume")
- @bp.route("/connectors/source-bindings", methods=["GET"])
- def connector_source_bindings_list():
- from app import db
- from app.core.connectors.bindings import ConnectorSourceBindingRepository
- try:
- items = ConnectorSourceBindingRepository(db.session).list_public(
- request.args.get("limit", 100)
- )
- return jsonify(success({"source_bindings": items})), 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/source-bindings", methods=["POST"])
- def connector_source_binding_approve():
- from app import db
- from app.core.connectors.bindings import ConnectorSourceBindingRepository
- from app.core.connectors.repository import ConnectorRepository
- payload = request.get_json(silent=True) or {}
- try:
- registry = _connector_registry()
- manifest = registry.resolve(
- payload.get("connector_id"), payload.get("version")
- ).manifest
- approved_config = registry.validate(
- payload.get("connector_id"),
- payload.get("version"),
- payload.get("approved_config") or {},
- )
- ConnectorRepository(db.session, _actor_uid()).register_manifest(
- manifest, _actor_uid()
- )
- result = ConnectorSourceBindingRepository(db.session).approve(
- connector_id=payload.get("connector_id"),
- connector_version=payload.get("version"),
- source_uid=payload.get("source_uid"),
- business_domain_uid=payload.get("business_domain_uid"),
- environment=payload.get("environment"),
- approved_config=approved_config,
- approved_by=_actor_uid(),
- binding_uid=payload.get("binding_uid"),
- )
- return jsonify(success(result, code=201)), 201
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/source-bindings/<binding_uid>/revoke", methods=["POST"])
- def connector_source_binding_revoke(binding_uid):
- from app import db
- from app.core.connectors.bindings import ConnectorSourceBindingRepository
- try:
- revoked = ConnectorSourceBindingRepository(db.session).revoke(
- binding_uid, _actor_uid()
- )
- return jsonify(success(revoked)), 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/principals", methods=["POST"])
- def connector_principal_create():
- from app import db
- from app.core.connectors.identity import ConnectorIdentityRepository
- payload = request.get_json(silent=True) or {}
- try:
- _connector_registry().resolve(
- payload.get("connector_id"), payload.get("version", "1.0.0")
- )
- uid = ConnectorIdentityRepository(db.session).create_principal(
- connector_id=payload.get("connector_id"),
- connector_version=payload.get("version", "1.0.0"),
- source_uid=payload.get("source_uid"),
- business_domain_uid=payload.get("business_domain_uid"),
- environment=payload.get("environment"),
- operations=payload.get("operations") or (),
- scopes=payload.get("scopes") or {},
- actor_uid=_actor_uid(),
- source_binding_uid=payload.get("source_binding_uid"),
- source_binding_version=payload.get("source_binding_version"),
- )
- return jsonify(success({"principal_uid": uid}, code=201)), 201
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/principals/<principal_uid>/credentials", methods=["POST"])
- def connector_credential_issue(principal_uid):
- from app import db
- from app.core.connectors.identity import ConnectorIdentityRepository
- payload = request.get_json(silent=True) or {}
- try:
- result = ConnectorIdentityRepository(db.session).issue(
- principal_uid,
- ttl_seconds=payload.get("ttl_seconds", 900),
- actor_uid=_actor_uid(),
- )
- response = jsonify(success(result, code=201))
- response.headers["Cache-Control"] = "no-store"
- return response, 201
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
- @bp.route("/connectors/credentials/<credential_uid>/<action>", methods=["POST"])
- def connector_credential_action(credential_uid, action):
- from app import db
- from app.core.connectors.identity import ConnectorIdentityRepository
- payload = request.get_json(silent=True) or {}
- try:
- repository = ConnectorIdentityRepository(db.session)
- if action == "revoke":
- result = {"revoked": repository.revoke(credential_uid, _actor_uid())}
- elif action == "rotate":
- result = repository.rotate(
- credential_uid,
- ttl_seconds=payload.get("ttl_seconds", 900),
- actor_uid=_actor_uid(),
- )
- else:
- raise ConnectorConfigurationError("credential action is invalid")
- response = jsonify(success(result))
- response.headers["Cache-Control"] = "no-store"
- return response, 200
- except Exception as error:
- db.session.rollback()
- return _error_response(error)
|