fix: storage and directories bug
This commit is contained in:
@@ -16,6 +16,7 @@ from common.db.models import (
|
|||||||
)
|
)
|
||||||
from common.ids import new_ulid
|
from common.ids import new_ulid
|
||||||
from common.storage.schemas import ServerObjectRequest
|
from common.storage.schemas import ServerObjectRequest
|
||||||
|
from common.storage import build_storage_uri
|
||||||
from fastapi import (
|
from fastapi import (
|
||||||
APIRouter,
|
APIRouter,
|
||||||
BackgroundTasks,
|
BackgroundTasks,
|
||||||
@@ -89,6 +90,15 @@ def safe_directory_name(value: str) -> str:
|
|||||||
return name
|
return name
|
||||||
|
|
||||||
|
|
||||||
|
# Version/run artifacts are not part of the user's workspace directory tree.
|
||||||
|
TREE_EXCLUDED_USAGE_TYPES = (
|
||||||
|
"version_artifact",
|
||||||
|
"snapshot",
|
||||||
|
"run_log",
|
||||||
|
"run_result",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def user_relative_path(context: RequestContext, child_path: str = "") -> str:
|
def user_relative_path(context: RequestContext, child_path: str = "") -> str:
|
||||||
base = f"users/{context.user.username}"
|
base = f"users/{context.user.username}"
|
||||||
normalized = normalize_user_path(child_path)
|
normalized = normalize_user_path(child_path)
|
||||||
@@ -464,7 +474,7 @@ async def create_script_record(
|
|||||||
bucket_name=bucket_name,
|
bucket_name=bucket_name,
|
||||||
object_key=object_key,
|
object_key=object_key,
|
||||||
object_key_hash=hashlib.sha256(object_key.encode("utf-8")).digest(),
|
object_key_hash=hashlib.sha256(object_key.encode("utf-8")).digest(),
|
||||||
storage_uri=f"s3://{bucket_name}/{object_key}",
|
storage_uri=build_storage_uri(bucket_name, object_key),
|
||||||
file_name=name,
|
file_name=name,
|
||||||
file_extension=PurePosixPath(jupyter_name).suffix.lower() or None,
|
file_extension=PurePosixPath(jupyter_name).suffix.lower() or None,
|
||||||
mime_type=mime_type,
|
mime_type=mime_type,
|
||||||
@@ -625,6 +635,7 @@ async def get_workspace_tree(
|
|||||||
StorageObjects.object_status == "available",
|
StorageObjects.object_status == "available",
|
||||||
StorageObjects.is_deleted == 0,
|
StorageObjects.is_deleted == 0,
|
||||||
StorageObjects.relative_path.like(like_prefix),
|
StorageObjects.relative_path.like(like_prefix),
|
||||||
|
StorageObjects.usage_type.notin_(TREE_EXCLUDED_USAGE_TYPES),
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
).all()
|
).all()
|
||||||
@@ -691,18 +702,20 @@ async def list_workspace_directories(
|
|||||||
|
|
||||||
rows = (
|
rows = (
|
||||||
await session.execute(
|
await session.execute(
|
||||||
select(StorageObjects.relative_path, StorageObjects.object_type).where(
|
select(StorageObjects.relative_path).where(
|
||||||
StorageObjects.workspace_id == context.workspace.workspace_id,
|
StorageObjects.workspace_id == context.workspace.workspace_id,
|
||||||
StorageObjects.object_status == "available",
|
StorageObjects.object_status == "available",
|
||||||
StorageObjects.is_deleted == 0,
|
StorageObjects.is_deleted == 0,
|
||||||
StorageObjects.relative_path.like(f"{descendant_prefix}%"),
|
StorageObjects.relative_path.like(f"{descendant_prefix}%"),
|
||||||
~StorageObjects.relative_path.like(f"{descendant_prefix}%/%"),
|
~StorageObjects.relative_path.like(f"{descendant_prefix}%/%"),
|
||||||
|
StorageObjects.object_type == "directory",
|
||||||
|
StorageObjects.usage_type.notin_(TREE_EXCLUDED_USAGE_TYPES),
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
).all()
|
).all()
|
||||||
|
|
||||||
directories: dict[str, dict[str, Any]] = {}
|
directories: dict[str, dict[str, Any]] = {}
|
||||||
for relative, _object_type in rows:
|
for (relative,) in rows:
|
||||||
if not relative or not relative.startswith(descendant_prefix):
|
if not relative or not relative.startswith(descendant_prefix):
|
||||||
continue
|
continue
|
||||||
suffix = relative[len(descendant_prefix) :]
|
suffix = relative[len(descendant_prefix) :]
|
||||||
@@ -720,7 +733,8 @@ async def list_workspace_directories(
|
|||||||
)
|
)
|
||||||
|
|
||||||
for directory in directories.values():
|
for directory in directories.values():
|
||||||
child_prefix = f"{target_prefix}/{directory['path']}/"
|
# directory['path'] is already workspace-relative and includes the parent segment.
|
||||||
|
child_prefix = f"{scoped_prefix}/{directory['path']}/"
|
||||||
has_children = await session.scalar(
|
has_children = await session.scalar(
|
||||||
select(StorageObjects.storage_object_id).where(
|
select(StorageObjects.storage_object_id).where(
|
||||||
StorageObjects.workspace_id == context.workspace.workspace_id,
|
StorageObjects.workspace_id == context.workspace.workspace_id,
|
||||||
@@ -728,6 +742,7 @@ async def list_workspace_directories(
|
|||||||
StorageObjects.is_deleted == 0,
|
StorageObjects.is_deleted == 0,
|
||||||
StorageObjects.relative_path.like(f"{child_prefix}%"),
|
StorageObjects.relative_path.like(f"{child_prefix}%"),
|
||||||
~StorageObjects.relative_path.like(f"{child_prefix}%/%"),
|
~StorageObjects.relative_path.like(f"{child_prefix}%/%"),
|
||||||
|
StorageObjects.usage_type.notin_(TREE_EXCLUDED_USAGE_TYPES),
|
||||||
).limit(1)
|
).limit(1)
|
||||||
)
|
)
|
||||||
directory["has_children"] = has_children is not None
|
directory["has_children"] = has_children is not None
|
||||||
|
|||||||
@@ -47,6 +47,7 @@ from common.storage.schemas import (
|
|||||||
DownloadUrlRequest,
|
DownloadUrlRequest,
|
||||||
ServerObjectRequest,
|
ServerObjectRequest,
|
||||||
)
|
)
|
||||||
|
from common.storage import build_storage_uri
|
||||||
|
|
||||||
|
|
||||||
# ── shared low-level helpers (module-private) ────────────────────────────
|
# ── shared low-level helpers (module-private) ────────────────────────────
|
||||||
@@ -90,7 +91,7 @@ def _build_storage_object(
|
|||||||
bucket_name=upload.bucket_name,
|
bucket_name=upload.bucket_name,
|
||||||
object_key=upload.object_key,
|
object_key=upload.object_key,
|
||||||
object_key_hash=upload.object_key_hash,
|
object_key_hash=upload.object_key_hash,
|
||||||
storage_uri=f"s3://{upload.bucket_name}/{upload.object_key}",
|
storage_uri=build_storage_uri(upload.bucket_name, upload.object_key),
|
||||||
file_name=safe_name,
|
file_name=safe_name,
|
||||||
file_extension=PurePosixPath(safe_name).suffix.lower() or None,
|
file_extension=PurePosixPath(safe_name).suffix.lower() or None,
|
||||||
mime_type=content_type,
|
mime_type=content_type,
|
||||||
|
|||||||
@@ -19,7 +19,7 @@ from common.db.models import (
|
|||||||
Workspaces,
|
Workspaces,
|
||||||
)
|
)
|
||||||
from common.ids import new_ulid
|
from common.ids import new_ulid
|
||||||
from common.storage import actual_bucket_name
|
from common.storage import actual_bucket_name, build_storage_uri
|
||||||
from common.storage.schemas import (
|
from common.storage.schemas import (
|
||||||
CreateUploadRequest,
|
CreateUploadRequest,
|
||||||
DownloadUrlRequest,
|
DownloadUrlRequest,
|
||||||
@@ -350,7 +350,7 @@ async def upload_bytes_to_session(
|
|||||||
bucket_name=upload.bucket_name,
|
bucket_name=upload.bucket_name,
|
||||||
object_key=upload.object_key,
|
object_key=upload.object_key,
|
||||||
object_key_hash=upload.object_key_hash,
|
object_key_hash=upload.object_key_hash,
|
||||||
storage_uri=f"s3://{upload.bucket_name}/{upload.object_key}",
|
storage_uri=build_storage_uri(upload.bucket_name, upload.object_key),
|
||||||
file_name=file_name,
|
file_name=file_name,
|
||||||
file_extension=PurePosixPath(file_name).suffix.lower() or None,
|
file_extension=PurePosixPath(file_name).suffix.lower() or None,
|
||||||
mime_type=upload.content_type,
|
mime_type=upload.content_type,
|
||||||
|
|||||||
@@ -23,6 +23,7 @@ from .factory import (
|
|||||||
PURPOSE_BUCKETS,
|
PURPOSE_BUCKETS,
|
||||||
RCLONE_REMOTE_NAME,
|
RCLONE_REMOTE_NAME,
|
||||||
actual_bucket_name,
|
actual_bucket_name,
|
||||||
|
build_storage_uri,
|
||||||
build_storage_config,
|
build_storage_config,
|
||||||
create_storage,
|
create_storage,
|
||||||
rclone_remote_spec,
|
rclone_remote_spec,
|
||||||
@@ -34,6 +35,7 @@ __all__ = [
|
|||||||
"create_storage",
|
"create_storage",
|
||||||
"build_storage_config",
|
"build_storage_config",
|
||||||
"actual_bucket_name",
|
"actual_bucket_name",
|
||||||
|
"build_storage_uri",
|
||||||
"workspaces_root",
|
"workspaces_root",
|
||||||
"rclone_remote_spec",
|
"rclone_remote_spec",
|
||||||
"RCLONE_REMOTE_NAME",
|
"RCLONE_REMOTE_NAME",
|
||||||
|
|||||||
@@ -78,6 +78,24 @@ def actual_bucket_name(purpose: str) -> str:
|
|||||||
return getattr(settings, f"s3_{purpose}_bucket")
|
return getattr(settings, f"s3_{purpose}_bucket")
|
||||||
|
|
||||||
|
|
||||||
|
def build_storage_uri(bucket_name: str, object_key: str) -> str:
|
||||||
|
"""根据 ``settings.storage_backend`` 构造对象的 storage_uri。
|
||||||
|
|
||||||
|
- s3 模式:``s3://{bucket_name}/{object_key}``
|
||||||
|
- local 模式:``file://{absolute_bucket_path}/{object_key}``
|
||||||
|
|
||||||
|
local 模式下 ``bucket_name`` 是文件系统路径(见 ``actual_bucket_name``),
|
||||||
|
不能直接用 ``s3://`` 前缀,否则会产生 ``s3:///data/version/...`` 这种
|
||||||
|
非法 URI。因此返回 ``file://`` URI,并把相对路径先转成绝对路径。
|
||||||
|
"""
|
||||||
|
from common.config import settings # 延迟 import 避免循环
|
||||||
|
|
||||||
|
if settings.storage_backend == "local":
|
||||||
|
absolute_path = Path(bucket_name).absolute().as_posix()
|
||||||
|
return f"file://{absolute_path}/{object_key}"
|
||||||
|
return f"s3://{bucket_name}/{object_key}"
|
||||||
|
|
||||||
|
|
||||||
def build_storage_config(bucket_name: str) -> Dict[str, Any]:
|
def build_storage_config(bucket_name: str) -> Dict[str, Any]:
|
||||||
"""根据 ``settings.storage_backend`` 构造 ``create_storage()`` 的入参。
|
"""根据 ``settings.storage_backend`` 构造 ``create_storage()`` 的入参。
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user