feat: 平台治理与对象存储增强,审批中心与运行日志整合

- 新增 storage/policy.py 落盘策略:按大小/类型决定文件存 MinIO 或内联数据库
- 数据处理源文件与生成结果写入 MinIO 并登记 storage_objects,支持失败回滚
- 算力节点训练产物按版本归档到 MinIO,登记 model_artifacts
- 数据转换任务输入输出对象化,支持从 MinIO 读写
- 新增审批中心(申请/我的/策略)、组织与权限、运行日志整合页面
- schema 与 docker 配置、前端路由侧边栏、治理文档同步更新

Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
wuyongtao
2026-08-19 16:10:02 +08:00
parent 81c2f85c3a
commit 78e3baa9ba
30 changed files with 1994 additions and 249 deletions

View File

@@ -35,6 +35,8 @@ from fastapi.responses import StreamingResponse
from psycopg.rows import dict_row
from app.core.auth import filter_accessible_resource_ids, get_current_user, is_admin
from app.core.config import get_settings
from app.db.platform_store import get_platform_store
from app.modules.data_process.algorithms import (
ParsedText,
canonical_record_json,
@@ -85,6 +87,8 @@ from app.modules.data_process.store import (
new_id,
repeat_task_id,
)
from app.modules.storage.minio_store import get_object_storage
from app.modules.storage.policy import should_store_in_minio
from app.schemas.data_process import (
DataProcessRegenerateRequest,
DataProcessRepeatRequest,
@@ -224,9 +228,38 @@ def _commit_source_batch(
staged: list[StagedSourceObject],
) -> list[dict[str, Any]]:
storage.publish(staged)
storage_object_ids: list[str] = []
try:
# The source reference remains in the task schema for compatibility,
# while storage_objects provides the authoritative MinIO index.
if get_settings().minio_enabled:
for item in prepared:
reference = str(item.get("storage_object_id") or "")
if not reference.startswith("minio://"):
continue
object_key = storage.object_key(reference)
metadata = get_object_storage().stat(object_key)
object_row = get_platform_store().create_storage_object({
"resource_type": "data_process_source",
"resource_id": task_id,
"version_id": str(item.get("id") or new_id("dpsf")),
"bucket": get_object_storage().bucket,
"object_key": object_key,
"file_name": item.get("name"),
"content_type": (item.get("metadata") or {}).get("content_type", "application/octet-stream"),
"byte_size": metadata.get("byte_size") or item.get("raw_size") or 0,
"checksum_sha256": item.get("checksum_sha256"),
"status": "available",
"created_by": item.get("created_by"),
})
storage_object_ids.append(str(object_row["id"]))
return store.add_source_files(task_id, prepared)
except Exception:
for object_id in storage_object_ids:
try:
get_platform_store().update_storage_object(object_id, {"status": "deleted"})
except Exception:
pass
for item in staged:
try:
storage.delete(item.reference)
@@ -774,6 +807,25 @@ def _run_generation(
duplicate_count=duplicate_count,
error_count=error_count,
)
result_bytes = "".join(
structured_json_dumps(item) + "\n" for item in accepted
).encode("utf-8")
if should_store_in_minio(len(result_bytes)):
result_key = f"data-process/{task_id}/results/{generation_run_id}.jsonl"
uploaded = get_object_storage().put_bytes(result_key, result_bytes, "application/jsonl")
get_platform_store().create_storage_object({
"resource_type": "data_process_result",
"resource_id": task_id,
"version_id": generation_run_id,
"bucket": uploaded["bucket"],
"object_key": result_key,
"file_name": f"{generation_run_id}.jsonl",
"content_type": "application/jsonl",
"byte_size": len(result_bytes),
"checksum_sha256": hashlib.sha256(result_bytes).hexdigest(),
"status": "available",
"created_by": (store.get_task(task_id) or {}).get("created_by"),
})
logger.info(
"data process generation completed task_id=%s generation_run_id=%s "
"output_count=%s filtered_count=%s duplicate_count=%s error_count=%s "
@@ -933,7 +985,7 @@ def _repeat_file_copies(
source = store.get_source_file(source_task_id, old_file_id, include_content=True)
new_file_id = new_id("dpsf")
old_reference = str(source.get("storage_object_id") or "")
if old_reference.startswith("local://data-process/"):
if old_reference.startswith(("local://data-process/", "minio://data-process/")):
staged_object = storage.stage_copy(
batch_id=batch_id,
source_reference=old_reference,

View File

@@ -23,8 +23,9 @@ from app.core.audit import audit_log, AuditActions
from app.core.op_log import op_log, OpModule, OpAction
from app.db.platform_store import get_platform_store
from app.modules.compute_gateway.client import ComputeNodeClient
from app.modules.compute_gateway.sync import fetch_eval_result_content, poll_compute_jobs_once
from app.modules.compute_gateway.sync import _archive_node_directory, fetch_eval_result_content, poll_compute_jobs_once
from app.modules.storage.minio_store import ObjectStorageError, get_object_storage
from app.modules.storage.policy import should_store_in_minio
router = APIRouter()
_LOGIN_FAILURES: dict[str, list[float]] = {}
@@ -100,6 +101,106 @@ async def _wait_for_object_storage() -> None:
await asyncio.sleep(max(1, settings.storage_check_interval_seconds))
def _store_json_snapshot(
store: Any,
resource_type: str,
resource_id: str,
version_id: str,
object_key: str,
payload: dict[str, Any],
created_by: str | None = None,
) -> dict[str, Any]:
"""Persist non-secret task/model parameters as an auditable MinIO snapshot."""
sanitized = {
key: value
for key, value in payload.items()
if key not in {"api_key", "secret_key", "password", "token", "access_token"}
}
raw = json.dumps(sanitized, ensure_ascii=False, sort_keys=True, default=str).encode("utf-8")
# Task/evaluation payloads are already persisted in PostgreSQL. Avoid an
# extra MinIO round trip for small non-secret parameter snapshots.
if not should_store_in_minio(len(raw), content_type="application/json", file_format="json"):
return {}
uploaded = get_object_storage().put_bytes(object_key, raw, "application/json")
return store.create_storage_object({
"resource_type": resource_type,
"resource_id": resource_id,
"version_id": version_id,
"bucket": uploaded["bucket"],
"object_key": object_key,
"file_name": Path(object_key).name,
"content_type": "application/json",
"byte_size": len(raw),
"checksum_sha256": hashlib.sha256(raw).hexdigest(),
"status": "available",
"created_by": created_by,
})
def _dataset_file_bytes(store: Any, file_id: str) -> bytes:
"""Return the canonical dataset bytes, lazily indexing legacy DB content."""
row = store.dataset_file(file_id)
if not get_settings().minio_enabled:
return str(row.get("content") or "").encode("utf-8")
storage_object = None
object_id = str(row.get("storage_object_id") or "")
if object_id:
try:
storage_object = store.storage_object(object_id)
except KeyError:
storage_object = None
if storage_object and storage_object.get("status") == "available":
try:
return get_object_storage().get_bytes(storage_object["object_key"])
except Exception as exc: # noqa: BLE001 - expose storage outage to callers
raise RuntimeError(f"dataset object is unavailable in MinIO: {exc}") from exc
# Compatibility migration for files created before MinIO was enabled.
raw = str(row.get("content") or "").encode("utf-8")
if not raw:
raise RuntimeError(f"dataset file has no MinIO object or legacy content: {file_id}")
# Small files intentionally remain database-backed. They can still be
# copied to a compute node directly when a task needs them.
if not should_store_in_minio(len(raw), file_format=row.get("file_format")):
return raw
object_key = (
f"datasets/{row['dataset_id']}/versions/"
f"{row.get('active_version_id') or row['id']}/{Path(str(row.get('name') or row['id'])).name}"
)
uploaded = get_object_storage().put_bytes(object_key, raw, "application/octet-stream")
created = store.create_storage_object({
"resource_type": "dataset",
"resource_id": str(row["dataset_id"]),
"version_id": str(row.get("active_version_id") or row["id"]),
"bucket": uploaded["bucket"],
"object_key": object_key,
"file_name": row.get("name"),
"content_type": "application/octet-stream",
"byte_size": len(raw),
"checksum_sha256": hashlib.sha256(raw).hexdigest(),
"status": "available",
})
store.link_dataset_file_storage_object(str(row["id"]), created["id"])
return raw
def _dataset_version_bytes(store: Any, file_id: str, version_id: str) -> bytes:
row = store.dataset_file(file_id)
try:
version = next(item for item in store.file_versions(file_id)["versions"] if item["id"] == version_id)
except StopIteration as exc:
raise KeyError(version_id) from exc
object_id = str(version.get("storage_object_id") or "")
if get_settings().minio_enabled and object_id:
try:
obj = store.storage_object(object_id)
return get_object_storage().get_bytes(obj["object_key"])
except KeyError:
pass
return _dataset_file_bytes(store, file_id)
def _select_eval_node(store: Any, preferred_node_id: str | None = None) -> dict[str, Any] | None:
"""Select the compute node for an eval job.
@@ -130,18 +231,63 @@ async def _prepare_resource_on_node(store: Any, resource_type: str, resource_id:
if not get_settings().minio_enabled or not resource_id:
return None
objects = store.storage_objects_for_resource(resource_type, resource_id)
if not objects and resource_type in {"model", "trained_model"}:
resource = None
if resource_type == "model":
try:
resource = store.model(resource_id)
except KeyError:
try:
resource = store.model_by_name(resource_id)
except KeyError:
resource = next((item for item in store.models() if item.get("path") == resource_id), None)
else:
resource = next(
(item for item in store.trained_models() if item.get("id") == resource_id or item.get("name") == resource_id),
None,
)
source_path = str(
(resource or {}).get("path")
or (resource or {}).get("merged_path")
or (resource or {}).get("artifact_dir")
or ""
)
resolved_id = str((resource or {}).get("id") or resource_id)
if source_path and resolved_id:
client = ComputeNodeClient(node["api_base_url"], timeout=900)
await _archive_node_directory(
store,
client,
node,
source_path,
resource_type,
resolved_id,
"legacy-import",
f"models/{resolved_id}" if resource_type == "model" else f"trained_models/{resolved_id}",
)
resource_id = resolved_id
objects = store.storage_objects_for_resource(resource_type, resource_id)
if not objects:
return None
client = ComputeNodeClient(node["api_base_url"], timeout=900)
root_name = "trained_models" if resource_type in {"trained_model", "model_artifact"} else f"{resource_type}s"
for obj in objects:
object_key = str(obj["object_key"])
marker = f"{root_name}/{resource_id}/versions/"
relative_name = Path(str(obj.get("file_name") or object_key)).name
if marker in object_key:
suffix = object_key.split(marker, 1)[1]
if "/" in suffix:
suffix = suffix.split("/", 1)[1]
if suffix:
relative_name = suffix
await client.prepare_cache({
"resource_id": resource_id,
"version_id": obj["version_id"],
"download_url": get_object_storage().presigned_get(obj["object_key"]),
"download_url": get_object_storage().presigned_get(object_key),
"checksum_sha256": obj.get("checksum_sha256") or "",
"byte_size": obj.get("byte_size") or 0,
"relative_path": f"{root_name}/{resource_id}/{Path(str(obj.get('file_name') or obj['object_key'])).name}",
"relative_path": f"{root_name}/{resource_id}/{relative_name}",
})
return f"/data/yg-ft/{root_name}/{resource_id}"
@@ -397,8 +543,23 @@ async def _submit_fine_tune_task(store: Any, payload: dict[str, Any]) -> dict[st
if not preflight["valid"]:
errors = "; ".join(preflight.get("errors") or ["preflight failed"])
raise RuntimeError(f"preflight failed: {errors}")
payload = {**payload, "compute_node_id": preflight["node"]["id"]}
prepared_job_payload = preflight.get("job_payload") or {}
payload = {
**payload,
"compute_node_id": preflight["node"]["id"],
"prepared_base_model_path": prepared_job_payload.get("model_name_or_path") or payload.get("prepared_base_model_path"),
}
task = store.start_task(payload)
if get_settings().minio_enabled:
_store_json_snapshot(
store,
"fine_tune",
str(task["id"]),
str(task["id"]),
f"training/{task['id']}/versions/{task['id']}/training-config.json",
task,
task.get("created_by"),
)
if get_settings().compute_mode == "simulator":
return task
node, job_payload = store.build_compute_job_payload(task["id"])
@@ -454,6 +615,20 @@ async def _fine_tune_preflight_with_job_payload(
if get_settings().minio_enabled and get_settings().compute_mode != "simulator":
try:
await _wait_for_object_storage()
base_model_path = str(job_payload.get("model_name_or_path") or job_payload.get("base_model") or "")
base_model_id = str(job_payload.get("base_model_id") or job_payload.get("model_id") or "")
base_model = None
if base_model_id:
try:
base_model = store.model(base_model_id)
except KeyError:
base_model = None
if base_model is None:
base_model = next((item for item in store.models() if item.get("path") == base_model_path), None)
if base_model:
prepared_model = await _prepare_resource_on_node(store, "model", str(base_model["id"]), node)
if prepared_model:
job_payload = {**job_payload, "base_model": prepared_model, "model_name_or_path": prepared_model}
except Exception as exc: # noqa: BLE001 - preflight exposes node storage failure
sync_errors.append(f"shared storage health check failed: {exc}")
if get_settings().compute_mode == "simulator":
@@ -1091,8 +1266,9 @@ async def merge_model(payload: dict[str, Any] = Body(...), current_user: dict =
@router.get("/dataset-manage/preview/{file_id}")
async def dataset_preview(file_id: str) -> dict[str, Any]:
try:
row = get_platform_store().dataset_file(file_id)
return ok({"content": row["content"]})
store = get_platform_store()
content = _dataset_file_bytes(store, file_id).decode("utf-8", errors="replace")
return ok({"content": content})
except KeyError:
raise fail(404, "dataset file not found")
@@ -1121,7 +1297,8 @@ async def dataset_version_content(file_id: str, version_id: str) -> dict[str, An
version = next((item for item in versions if item["id"] == version_id), None)
if not version:
raise KeyError(version_id)
return ok({"version": version, "content": row["content"]})
content = _dataset_version_bytes(get_platform_store(), file_id, version_id)
return ok({"version": version, "content": content.decode("utf-8", errors="replace")})
except KeyError:
raise fail(404, "dataset version not found")
@@ -1210,7 +1387,7 @@ async def _sync_training_dataset_to_compute_node(
) -> list[dict[str, Any]]:
if get_settings().minio_enabled:
files = store.training_dataset_files(dataset_id)
object_by_resource_name: dict[tuple[str, str], dict[str, Any]] = {}
object_by_resource_name: dict[tuple[str, str, str], dict[str, Any]] = {}
resource_ids = {str(dataset_id)} | {
str(item.get("dataset_id"))
for item in files
@@ -1219,35 +1396,58 @@ async def _sync_training_dataset_to_compute_node(
for resource_id in resource_ids:
for obj in store.storage_objects_for_resource("dataset", resource_id):
file_name = Path(str(obj.get("file_name") or obj.get("object_key") or "")).name
object_by_resource_name[(resource_id, file_name)] = obj
object_by_resource_name[(resource_id, file_name, str(obj.get("version_id") or ""))] = obj
results: list[dict[str, Any]] = []
client = ComputeNodeClient(node["api_base_url"])
for item in files:
target_name = Path(str(item.get("name") or f"{item['id']}.jsonl")).name
item_dataset_id = str(item.get("dataset_id") or dataset_id)
obj = object_by_resource_name.get((item_dataset_id, target_name))
version_id = str(item.get("active_version_id") or item["id"])
obj = object_by_resource_name.get((item_dataset_id, target_name, version_id))
if not obj and item.get("content"):
# 兼容 MinIO 接入前已经发布的数据处理数据集
# 预检时用数据库正文补建对象,避免要求用户重新处理数据集
# 兼容 MinIO 接入前已经发布的数据处理数据集。大文件补建
# MinIO 对象,小文件直接从数据库正文同步到目标节点
raw = str(item.get("content") or "").encode("utf-8")
version_id = str(item.get("active_version_id") or item["id"])
object_key = f"datasets/{item_dataset_id}/versions/{version_id}/{target_name}"
uploaded = get_object_storage().put_bytes(object_key, raw, "application/jsonl")
obj = store.create_storage_object({
"resource_type": "dataset",
"resource_id": item_dataset_id,
"version_id": version_id,
"bucket": uploaded["bucket"],
"object_key": object_key,
"file_name": target_name,
"content_type": "application/jsonl",
"byte_size": len(raw),
"checksum_sha256": hashlib.sha256(raw).hexdigest(),
"status": "available",
})
store.link_dataset_file_storage_object(str(item["id"]), obj["id"])
if should_store_in_minio(len(raw)):
object_key = f"datasets/{item_dataset_id}/versions/{version_id}/{target_name}"
uploaded = get_object_storage().put_bytes(object_key, raw, "application/jsonl")
obj = store.create_storage_object({
"resource_type": "dataset",
"resource_id": item_dataset_id,
"version_id": version_id,
"bucket": uploaded["bucket"],
"object_key": object_key,
"file_name": target_name,
"content_type": "application/jsonl",
"byte_size": len(raw),
"checksum_sha256": hashlib.sha256(raw).hexdigest(),
"status": "available",
})
store.link_dataset_file_storage_object(str(item["id"]), obj["id"])
else:
result = await client.upload_file(
target_name,
raw,
f"datasets/{dataset_id}/{target_name}",
resource_type="dataset",
resource_id=dataset_id,
)
store.upsert_resource_replica(
node["id"], "dataset", dataset_id, str(result.get("local_path") or "")
)
results.append({
"node_id": node["id"],
"node_code": node.get("code"),
"file_id": item.get("id"),
"name": target_name,
"local_path": result.get("local_path"),
"byte_size": result.get("byte_size"),
"checksum_sha256": result.get("checksum_sha256"),
"storage_backend": "database",
})
continue
if not obj:
raise RuntimeError(f"dataset file is not available in MinIO: {target_name}")
raise RuntimeError(f"dataset file is not available: {target_name}")
url = get_object_storage().presigned_get(obj["object_key"])
result = await client.prepare_cache({
"resource_id": dataset_id,
@@ -1263,7 +1463,13 @@ async def _sync_training_dataset_to_compute_node(
dataset_id,
str(result.get("local_path") or ""),
)
results.append({**result, "file_id": item.get("id"), "name": target_name, "node_id": node["id"]})
results.append({
**result,
"file_id": item.get("id"),
"name": target_name,
"node_id": node["id"],
"storage_backend": "minio",
})
return results
if not dataset_id:
raise RuntimeError("train_dataset_id is required")
@@ -1328,7 +1534,11 @@ async def upload_dataset_files(
created_file = store.add_dataset_file(conn, dataset_id, file.filename or "upload.jsonl", content)
created.append(created_file)
pending_sync.append((created_file["id"], created_file["name"], raw))
if get_settings().minio_enabled:
if should_store_in_minio(
len(raw),
content_type=file.content_type,
file_format=Path(created_file["name"]).suffix,
):
object_key = f"datasets/{dataset_id}/versions/{created_file.get('active_version_id') or created_file['id']}/{Path(created_file['name']).name}"
uploaded = get_object_storage().put_bytes(object_key, raw, file.content_type or "application/octet-stream")
storage_object = get_platform_store().create_storage_object({
@@ -1372,7 +1582,10 @@ async def download_dataset(dataset_id: str, current_user: dict = Depends(get_cur
full_file = store.dataset_file(str(item["id"]))
except KeyError:
continue
files.append({**item, "content": full_file.get("content") or ""})
files.append({
**item,
"content": _dataset_file_bytes(store, str(item["id"])).decode("utf-8", errors="replace"),
})
if not files:
raise fail(404, "dataset has no downloadable files")
@@ -1413,8 +1626,12 @@ async def download_dataset(dataset_id: str, current_user: dict = Depends(get_cur
@router.get("/dataset-manage/download/{dataset_id}/{file_id}")
async def download_dataset_file(dataset_id: str, file_id: str, version_id: str | None = Query(default=None)) -> PlainTextResponse:
row = get_platform_store().dataset_file(file_id)
return PlainTextResponse(row["content"], media_type="text/plain")
store = get_platform_store()
row = store.dataset_file(file_id)
if str(row.get("dataset_id")) != str(dataset_id):
raise fail(404, "dataset file not found")
content = _dataset_version_bytes(store, file_id, version_id) if version_id else _dataset_file_bytes(store, file_id)
return PlainTextResponse(content.decode("utf-8", errors="replace"), media_type="text/plain")
@router.get("/dataset-manage")
@@ -1880,12 +2097,26 @@ async def model_eval_start(payload: dict[str, Any] = Body(...), current_user: di
# 1. Create eval task record
payload.setdefault("created_by", current_user.get("id"))
task = store.create_eval_task({**payload, "status": "pending"})
if get_settings().minio_enabled:
_store_json_snapshot(
store,
"eval",
str(task["id"]),
str(task["id"]),
f"evaluations/{task['id']}/versions/{task['id']}/evaluation-config.json",
payload,
current_user.get("id"),
)
# 2. Resolve model path (supports both regular models and trained models)
model_id = str(payload.get("model_id", ""))
model_path = ""
adapter_path = payload.get("adapter_path", "")
model_node_id = ""
ds_files: list[dict[str, Any]] = []
model_resource_type = "model"
model_resource_id = model_id
adapter_resource_id = ""
try:
db_model = store.model(model_id)
model_path = db_model.get("path", "")
@@ -1894,6 +2125,9 @@ async def model_eval_start(payload: dict[str, Any] = Body(...), current_user: di
# Try trained_models table (IDs prefixed with tm_)
trained = next((m for m in store.trained_models() if m["id"] == model_id), None)
if trained:
model_resource_type = "trained_model" if trained.get("merged") else "model"
model_resource_id = trained.get("id") or model_id
adapter_resource_id = trained.get("id") or ""
model_node_id = trained.get("compute_node_id") or ""
merged_path = trained.get("merged_path", "")
base_path = trained.get("base_model_path", "")
@@ -1907,6 +2141,9 @@ async def model_eval_start(payload: dict[str, Any] = Body(...), current_user: di
adapter_path = merged_path
else:
model_path = merged_path or base_path
if model_resource_type == "model":
base_model = next((item for item in store.models() if item.get("path") == base_path), None)
model_resource_id = str((base_model or {}).get("id") or base_path)
if not model_path:
store.update_eval_task(task["id"], {"status": "failed", "error": "model not found or no path"})
return ok({"task_id": task["id"], "status": "failed", "error": "model not found or no path"})
@@ -1979,6 +2216,22 @@ async def model_eval_start(payload: dict[str, Any] = Body(...), current_user: di
store.update_eval_task(task["id"], {"status": "failed", "error": message})
return ok({"task_id": task["id"], "status": "failed", "error": message})
if get_settings().minio_enabled and get_settings().compute_mode != "simulator":
try:
prepared_model = await _prepare_resource_on_node(store, model_resource_type, model_resource_id, node)
if prepared_model:
model_path = prepared_model
if adapter_resource_id and model_resource_type == "model":
prepared_adapter = await _prepare_resource_on_node(store, "trained_model", adapter_resource_id, node)
if prepared_adapter:
adapter_path = prepared_adapter
dataset_sync = await _sync_training_dataset_to_compute_node(store, node, dataset_id)
if dataset_sync:
dataset_path = str(dataset_sync[0].get("local_path") or dataset_path)
except Exception as exc:
store.update_eval_task(task["id"], {"status": "failed", "error": f"MinIO resource preparation failed: {exc}"})
return ok({"task_id": task["id"], "status": "failed", "error": str(exc)})
node_gpus = {
int(item.get("id", item.get("gpu_index", -1))): item
for item in store.gpus()