from __future__ import annotations import hashlib import io import json import uuid import zipfile import base64 import xml.etree.ElementTree as ET from datetime import datetime, timezone from pathlib import Path import aiofiles from fastapi import HTTPException, UploadFile, status from sqlalchemy import and_, func, or_, select from sqlalchemy.ext.asyncio import AsyncSession from app.core.config import settings from app.core.deps import is_system_admin from app.core.project_permissions import role_has_api_permission from app.crud import member as member_crud from app.models.collaboration import ( CollaborationEditRequest, CollaborationFile, CollaborationFolder, CollaborationMember, CollaborationRevision, ) from app.models.study_member import StudyMember from app.models.user import User from app.schemas.collaboration import ( CollaborationCandidateRead, CollaborationEditRequestRead, CollaborationEditRequestResolve, CollaborationFileCollaboratorRead, CollaborationFileCreate, CollaborationFileRead, CollaborationFileUpdate, CollaborationFolderCreate, CollaborationFolderUpdate, CollaborationMemberRead, CollaborationMemberUpsert, CollaborationOwnershipTransferRequest, CollaborationRevisionCopyRequest, CollaborationRevisionRead, CollaborationRevisionUpdate, ) from app.services import notification_service COLLABORATION_ROOT = Path(__file__).resolve().parent.parent / "uploads" / "collaboration" FILE_TYPE_EXTENSION = {"word": "docx", "cell": "xlsx", "slide": "pptx"} EXTENSION_FILE_TYPE = {value: key for key, value in FILE_TYPE_EXTENSION.items()} MIME_TYPES = { "docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document", "xlsx": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", "pptx": "application/vnd.openxmlformats-officedocument.presentationml.presentation", } ROLE_ASSIGNMENT_PERMISSIONS = { "EDITOR": ("collaboration:read", "collaboration:edit"), "MANAGER": ( "collaboration:read", "collaboration:create", "collaboration:edit", "collaboration:manage", "collaboration:export", "collaboration:delete", ), } WORKBOOK_NAMESPACE = "http://schemas.openxmlformats.org/spreadsheetml/2006/main" RELATIONSHIP_NAMESPACE = "http://schemas.openxmlformats.org/officeDocument/2006/relationships" def _workbook_structure_password(file_id: uuid.UUID) -> str: # The password is never exposed to clients. It exists only to prevent an # editor from bypassing the file-level policy by clicking "Unprotect". seed = f"{settings.ONLYOFFICE_JWT_SECRET or ''}:sheet-structure:{file_id}" return hashlib.sha256(seed.encode("utf-8")).hexdigest()[:28] def apply_workbook_structure_policy( content: bytes, *, file_id: uuid.UUID, allow_sheet_structure_edit: bool, protection_backup: str | None, ) -> tuple[bytes, str | None]: """Protect or restore an XLSX workbook structure without rewriting cell data.""" if allow_sheet_structure_edit and protection_backup is None: return content, None if not zipfile.is_zipfile(io.BytesIO(content)): raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="Excel 文件格式无效") source = io.BytesIO(content) destination = io.BytesIO() ET.register_namespace("", WORKBOOK_NAMESPACE) ET.register_namespace("r", RELATIONSHIP_NAMESPACE) with zipfile.ZipFile(source, "r") as reader: try: workbook_xml = reader.read("xl/workbook.xml") except KeyError as exc: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="Excel 工作簿结构无效") from exc root = ET.fromstring(workbook_xml) protection_tag = f"{{{WORKBOOK_NAMESPACE}}}workbookProtection" existing = root.find(protection_tag) next_backup = protection_backup if allow_sheet_structure_edit: if existing is not None: root.remove(existing) if protection_backup: try: backup_payload = json.loads(protection_backup) backup_index = max(0, int(backup_payload["index"])) backup_xml = str(backup_payload["xml"]) except (KeyError, TypeError, ValueError, json.JSONDecodeError): # Compatibility with an early development backup that only # stored the element bytes. backup_index = 0 backup_xml = protection_backup root.insert( min(backup_index, len(root)), ET.fromstring(base64.b64decode(backup_xml.encode("ascii"))), ) next_backup = None else: if protection_backup is None: if existing is None: next_backup = "" else: next_backup = json.dumps({ "index": list(root).index(existing), "xml": base64.b64encode(ET.tostring(existing, encoding="utf-8")).decode("ascii"), }, separators=(",", ":")) if existing is not None: root.remove(existing) # ONLYOFFICE honours OOXML workbook structure protection for add, # delete, move, hide and rename operations. from openpyxl.utils.protection import hash_password protected = ET.Element(protection_tag, { "lockStructure": "1", "workbookPassword": hash_password(_workbook_structure_password(file_id)), }) insert_at = next( (index for index, child in enumerate(root) if child.tag in { f"{{{WORKBOOK_NAMESPACE}}}bookViews", f"{{{WORKBOOK_NAMESPACE}}}sheets", }), 0, ) root.insert(insert_at, protected) rewritten = ET.tostring(root, encoding="utf-8", xml_declaration=True) with zipfile.ZipFile(destination, "w") as writer: for entry in reader.infolist(): writer.writestr(entry, rewritten if entry.filename == "xl/workbook.xml" else reader.read(entry.filename)) return destination.getvalue(), next_backup def _safe_title(value: str, extension: str | None = None) -> str: title = (value or "").replace("\\", "/").rsplit("/", 1)[-1].strip() title = "".join(char for char in title if ord(char) >= 32 and ord(char) != 127) title = title[:240].strip(". ") or "未命名文件" if extension: suffix = f".{extension}" if not title.lower().endswith(suffix): title = f"{title}{suffix}" return title def _write_package(parts: dict[str, str]) -> bytes: buffer = io.BytesIO() with zipfile.ZipFile(buffer, "w", compression=zipfile.ZIP_DEFLATED) as archive: for name, content in parts.items(): archive.writestr(name, content.strip()) return buffer.getvalue() def _blank_docx() -> bytes: return _write_package({ "[Content_Types].xml": """ """, "_rels/.rels": """ """, "word/document.xml": """ """, }) def _blank_xlsx() -> bytes: return _write_package({ "[Content_Types].xml": """ """, "_rels/.rels": """ """, "xl/workbook.xml": """ """, "xl/_rels/workbook.xml.rels": """ """, "xl/worksheets/sheet1.xml": """ """, }) def _blank_pptx() -> bytes: return _write_package({ "[Content_Types].xml": """ """, "_rels/.rels": """ """, "ppt/presentation.xml": """ """, "ppt/_rels/presentation.xml.rels": """ """, "ppt/slides/slide1.xml": """ """, "ppt/slides/_rels/slide1.xml.rels": """ """, "ppt/slideLayouts/slideLayout1.xml": """ """, "ppt/slideLayouts/_rels/slideLayout1.xml.rels": """ """, "ppt/slideMasters/slideMaster1.xml": """ """, "ppt/slideMasters/_rels/slideMaster1.xml.rels": """ """, "ppt/theme/theme1.xml": """ """, }) def blank_file_bytes(file_type: str) -> bytes: if file_type == "word": return _blank_docx() if file_type == "cell": return _blank_xlsx() if file_type == "slide": return _blank_pptx() raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="不支持的协作文件类型") async def _folder_or_404(db: AsyncSession, study_id: uuid.UUID, folder_id: uuid.UUID) -> CollaborationFolder: folder = await db.get(CollaborationFolder, folder_id) if not folder or folder.study_id != study_id or folder.deleted_at: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作文件夹不存在") return folder async def get_file_or_404( db: AsyncSession, study_id: uuid.UUID, file_id: uuid.UUID, *, include_deleted: bool = False ) -> CollaborationFile: item = await db.get(CollaborationFile, file_id) if not item or item.study_id != study_id or (item.deleted_at and not include_deleted): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作文件不存在") return item async def _member_role(db: AsyncSession, file_id: uuid.UUID, user_id: uuid.UUID) -> str | None: return await db.scalar( select(CollaborationMember.role).where( CollaborationMember.file_id == file_id, CollaborationMember.user_id == user_id, ) ) async def _role_is_assignable( db: AsyncSession, study_id: uuid.UUID, user, membership, role: str, ) -> bool: if is_system_admin(user): return True if not membership or not membership.is_active: return False required_permissions = ROLE_ASSIGNMENT_PERMISSIONS.get(role, ()) return all([ await role_has_api_permission( db, study_id, membership.role_in_study, permission, ) for permission in required_permissions ]) async def _require_role_assignable( db: AsyncSession, study_id: uuid.UUID, user, membership, role: str, ) -> None: if await _role_is_assignable(db, study_id, user, membership, role): return role_label = "管理者或所有者" if role == "MANAGER" else "编辑者" raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail=f"该账号的项目角色不具备{role_label}资格", ) async def can_edit_file(db: AsyncSession, item: CollaborationFile, user) -> bool: if is_system_admin(user): return True membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: return False role = await _member_role(db, item.id, user.id) return item.owner_id == user.id or role in {"EDITOR", "MANAGER"} async def can_manage_file(db: AsyncSession, item: CollaborationFile, user) -> bool: if is_system_admin(user): return True membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: return False return item.owner_id == user.id or await _member_role(db, item.id, user.id) == "MANAGER" async def can_export_file(db: AsyncSession, item: CollaborationFile, user) -> bool: if is_system_admin(user): return True membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: return False role = await _member_role(db, item.id, user.id) is_manager = item.owner_id == user.id or role == "MANAGER" return is_manager or bool(item.allow_export and role == "EDITOR") async def can_create_file(db: AsyncSession, item: CollaborationFile, user) -> bool: if is_system_admin(user): return True membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: return False return await role_has_api_permission( db, item.study_id, membership.role_in_study, "collaboration:create" ) async def edit_request_status( db: AsyncSession, item: CollaborationFile, user_id: uuid.UUID ) -> str | None: return await db.scalar( select(CollaborationEditRequest.status) .where( CollaborationEditRequest.file_id == item.id, CollaborationEditRequest.requester_id == user_id, ) .order_by(CollaborationEditRequest.created_at.desc()) .limit(1) ) async def can_request_edit_file(db: AsyncSession, item: CollaborationFile, user) -> bool: if not item.allow_edit_request or await can_edit_file(db, item, user): return False membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: return False if not await role_has_api_permission(db, item.study_id, membership.role_in_study, "collaboration:read"): return False return await edit_request_status(db, item, user.id) != "PENDING" async def can_transfer_ownership(db: AsyncSession, item: CollaborationFile, user) -> bool: if is_system_admin(user): return True membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: return False return item.owner_id == user.id async def require_file_editor(db: AsyncSession, item: CollaborationFile, user) -> None: if not await can_edit_file(db, item, user): raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="您未被邀请编辑此协作文件") async def require_file_manager(db: AsyncSession, item: CollaborationFile, user) -> None: if not await can_manage_file(db, item, user): raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="您没有管理此协作文件的权限") async def require_file_exporter(db: AsyncSession, item: CollaborationFile, user) -> None: if not await can_export_file(db, item, user): raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="您没有下载或另存此协作文件的权限") async def require_ownership_transfer(db: AsyncSession, item: CollaborationFile, user) -> None: if not await can_transfer_ownership(db, item, user): raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="您没有转让此协作文件所有权的权限") async def create_folder(db: AsyncSession, study_id: uuid.UUID, payload: CollaborationFolderCreate, user) -> CollaborationFolder: if payload.parent_id: await _folder_or_404(db, study_id, payload.parent_id) duplicate = await db.scalar( select(CollaborationFolder.id).where( CollaborationFolder.study_id == study_id, CollaborationFolder.parent_id == payload.parent_id, CollaborationFolder.name == payload.name, CollaborationFolder.deleted_at.is_(None), ) ) if duplicate: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="同级文件夹名称已存在") folder = CollaborationFolder(study_id=study_id, created_by=user.id, **payload.model_dump()) db.add(folder) await db.commit() await db.refresh(folder) return folder async def list_folders(db: AsyncSession, study_id: uuid.UUID) -> list[CollaborationFolder]: rows = await db.scalars( select(CollaborationFolder) .where(CollaborationFolder.study_id == study_id, CollaborationFolder.deleted_at.is_(None)) .order_by(CollaborationFolder.sort_order, CollaborationFolder.name) ) return list(rows.all()) async def update_folder( db: AsyncSession, study_id: uuid.UUID, folder_id: uuid.UUID, payload: CollaborationFolderUpdate ) -> CollaborationFolder: folder = await _folder_or_404(db, study_id, folder_id) values = payload.model_dump(exclude_unset=True) parent_id = values.get("parent_id", folder.parent_id) if parent_id == folder.id: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="文件夹不能作为自身的上级") if parent_id: parent = await _folder_or_404(db, study_id, parent_id) seen = {folder.id} while parent.parent_id: if parent.id in seen: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="不能形成循环文件夹") seen.add(parent.id) parent = await _folder_or_404(db, study_id, parent.parent_id) if values.get("name") is not None: values["name"] = values["name"].strip() if not values["name"]: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="文件夹名称不能为空") target_name = values.get("name", folder.name) duplicate = await db.scalar( select(CollaborationFolder.id).where( CollaborationFolder.study_id == study_id, CollaborationFolder.parent_id == parent_id, CollaborationFolder.name == target_name, CollaborationFolder.id != folder.id, CollaborationFolder.deleted_at.is_(None), ) ) if duplicate: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="同级文件夹名称已存在") for key, value in values.items(): setattr(folder, key, value) await db.commit() await db.refresh(folder) return folder async def delete_folder(db: AsyncSession, study_id: uuid.UUID, folder_id: uuid.UUID) -> None: folder = await _folder_or_404(db, study_id, folder_id) has_children = await db.scalar( select(CollaborationFolder.id).where( CollaborationFolder.parent_id == folder.id, CollaborationFolder.deleted_at.is_(None), ).limit(1) ) has_files = await db.scalar( select(CollaborationFile.id).where( CollaborationFile.folder_id == folder.id, CollaborationFile.deleted_at.is_(None), ).limit(1) ) if has_children or has_files: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="文件夹非空,无法删除") folder.deleted_at = datetime.now(timezone.utc) await db.commit() async def _persist_revision_bytes( db: AsyncSession, item: CollaborationFile, content: bytes, *, source: str, created_by: uuid.UUID | None, change_summary: str | None = None, original_filename: str | None = None, ) -> tuple[CollaborationRevision, bool]: locked = await db.scalar(select(CollaborationFile).where(CollaborationFile.id == item.id).with_for_update()) if not locked: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作文件不存在") digest = hashlib.sha256(content).hexdigest() current = await db.get(CollaborationRevision, locked.current_revision_id) if locked.current_revision_id else None if current and current.file_hash == digest: return current, False last_no = await db.scalar( select(func.max(CollaborationRevision.revision_no)).where(CollaborationRevision.file_id == locked.id) ) revision_no = int(last_no or 0) + 1 directory = COLLABORATION_ROOT / str(locked.study_id) / str(locked.id) directory.mkdir(parents=True, exist_ok=True) destination = directory / f"{uuid.uuid4()}.{locked.extension}" async with aiofiles.open(destination, "wb") as stream: await stream.write(content) revision = CollaborationRevision( file_id=locked.id, revision_no=revision_no, parent_revision_id=locked.current_revision_id, file_uri=str(destination), original_filename=_safe_title(original_filename or locked.title, locked.extension), file_hash=digest, file_size=len(content), mime_type=MIME_TYPES[locked.extension], source=source, change_summary=change_summary, created_by=created_by, ) db.add(revision) await db.flush() locked.current_revision_id = revision.id locked.updated_at = datetime.now(timezone.utc) return revision, True async def append_revision( db: AsyncSession, item: CollaborationFile, content: bytes, *, source: str, created_by: uuid.UUID | None, change_summary: str | None = None, ) -> tuple[CollaborationRevision, bool]: if not content or len(content) > settings.COLLABORATION_MAX_FILE_BYTES: raise HTTPException(status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail="协作文件大小超出限制") return await _persist_revision_bytes( db, item, content, source=source, created_by=created_by, change_summary=change_summary ) async def _replace_revision_bytes(revision: CollaborationRevision, content: bytes) -> None: """Update the stored representation without creating a content-history entry.""" destination = Path(revision.file_uri) temporary = destination.with_name(f".{destination.name}.{uuid.uuid4().hex}.tmp") try: async with aiofiles.open(temporary, "wb") as stream: await stream.write(content) temporary.replace(destination) finally: temporary.unlink(missing_ok=True) revision.file_hash = hashlib.sha256(content).hexdigest() revision.file_size = len(content) async def _create_file_from_bytes( db: AsyncSession, study_id: uuid.UUID, *, title: str, file_type: str, folder_id: uuid.UUID | None, content: bytes, source: str, user, ) -> CollaborationFile: membership = await member_crud.get_member(db, study_id, user.id) await _require_role_assignable(db, study_id, user, membership, "MANAGER") if folder_id: await _folder_or_404(db, study_id, folder_id) extension = FILE_TYPE_EXTENSION[file_type] item = CollaborationFile( study_id=study_id, folder_id=folder_id, title=_safe_title(title, extension), file_type=file_type, extension=extension, owner_id=user.id, ) db.add(item) await db.flush() db.add(CollaborationMember(file_id=item.id, user_id=user.id, role="MANAGER", invited_by=user.id)) await _persist_revision_bytes(db, item, content, source=source, created_by=user.id, original_filename=item.title) await db.commit() await db.refresh(item) return item async def create_blank_file( db: AsyncSession, study_id: uuid.UUID, payload: CollaborationFileCreate, user ) -> CollaborationFile: return await _create_file_from_bytes( db, study_id, title=payload.title, file_type=payload.file_type, folder_id=payload.folder_id, content=blank_file_bytes(payload.file_type), source="CREATE", user=user, ) async def import_file( db: AsyncSession, study_id: uuid.UUID, folder_id: uuid.UUID | None, upload: UploadFile, user ) -> CollaborationFile: filename = _safe_title(upload.filename or "") extension = Path(filename).suffix.lower().lstrip(".") file_type = EXTENSION_FILE_TYPE.get(extension) if not file_type: raise HTTPException(status_code=status.HTTP_415_UNSUPPORTED_MEDIA_TYPE, detail="仅支持 DOCX、XLSX、PPTX") content = await upload.read(settings.COLLABORATION_MAX_FILE_BYTES + 1) if not content or len(content) > settings.COLLABORATION_MAX_FILE_BYTES: raise HTTPException(status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail="协作文件为空或大小超出限制") if not zipfile.is_zipfile(io.BytesIO(content)): raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="Office 文件格式无效或已损坏") return await _create_file_from_bytes( db, study_id, title=filename, file_type=file_type, folder_id=folder_id, content=content, source="IMPORT", user=user, ) async def _collaborators_by_file( db: AsyncSession, file_ids: list[uuid.UUID], ) -> dict[uuid.UUID, list[CollaborationFileCollaboratorRead]]: if not file_ids: return {} rows = (await db.execute( select(CollaborationMember, User) .join(User, User.id == CollaborationMember.user_id) .where(CollaborationMember.file_id.in_(file_ids)) .order_by(CollaborationMember.file_id, CollaborationMember.role.desc(), User.full_name) )).all() result: dict[uuid.UUID, list[CollaborationFileCollaboratorRead]] = { file_id: [] for file_id in file_ids } for member, collaborator in rows: result[member.file_id].append(CollaborationFileCollaboratorRead( user_id=member.user_id, full_name=collaborator.full_name, role=member.role, avatar_url=collaborator.avatar_url, )) return result async def _file_read( db: AsyncSession, item: CollaborationFile, user, *, collaborators: list[CollaborationFileCollaboratorRead] | None = None, ) -> CollaborationFileRead: owner = await db.get(User, item.owner_id) folder = await db.get(CollaborationFolder, item.folder_id) if item.folder_id else None revision = await db.get(CollaborationRevision, item.current_revision_id) if item.current_revision_id else None role = "MANAGER" if is_system_admin(user) else await _member_role(db, item.id, user.id) can_edit = await can_edit_file(db, item, user) request_status = None if can_edit else await edit_request_status(db, item, user.id) if collaborators is None: collaborators = (await _collaborators_by_file(db, [item.id])).get(item.id, []) return CollaborationFileRead.model_validate(item).model_copy(update={ "owner_name": owner.full_name if owner else None, "folder_name": folder.name if folder else None, "current_revision_no": revision.revision_no if revision else None, "current_revision_file_size": revision.file_size if revision else None, "current_revision_mime_type": revision.mime_type if revision else None, "current_revision_created_at": revision.created_at if revision else None, "collaboration_role": role, "collaborators": collaborators, "can_edit": can_edit, "can_manage": await can_manage_file(db, item, user), "can_export": await can_export_file(db, item, user), "can_request_edit": await can_request_edit_file(db, item, user), "edit_request_status": request_status, "can_transfer_ownership": await can_transfer_ownership(db, item, user), }) async def file_read(db: AsyncSession, item: CollaborationFile, user) -> CollaborationFileRead: return await _file_read(db, item, user) async def list_files( db: AsyncSession, study_id: uuid.UUID, user, *, folder_id: uuid.UUID | None = None, keyword: str | None = None, deleted: bool = False, ) -> list[CollaborationFileRead]: stmt = select(CollaborationFile).where(CollaborationFile.study_id == study_id) stmt = stmt.where(CollaborationFile.deleted_at.is_not(None) if deleted else CollaborationFile.deleted_at.is_(None)) if folder_id: stmt = stmt.where(CollaborationFile.folder_id == folder_id) if keyword and keyword.strip(): stmt = stmt.where(CollaborationFile.title.ilike(f"%{keyword.strip()}%")) rows = list((await db.scalars(stmt.order_by(CollaborationFile.updated_at.desc()))).all()) collaborators = await _collaborators_by_file(db, [item.id for item in rows]) return [ await _file_read(db, item, user, collaborators=collaborators.get(item.id, [])) for item in rows ] async def update_file( db: AsyncSession, item: CollaborationFile, payload: CollaborationFileUpdate, user ) -> CollaborationFile: await require_file_manager(db, item, user) values = payload.model_dump(exclude_unset=True) if "folder_id" in values and values["folder_id"]: await _folder_or_404(db, item.study_id, values["folder_id"]) if "title" in values: values["title"] = _safe_title(values["title"], item.extension) requested_sheet_policy = values.pop("allow_sheet_structure_edit", None) if requested_sheet_policy is not None and requested_sheet_policy != item.allow_sheet_structure_edit: if item.file_type != "cell": raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="仅 Excel 协作文件支持工作表增删权限", ) revision = await db.get(CollaborationRevision, item.current_revision_id) if item.current_revision_id else None if not revision or not Path(revision.file_uri).is_file(): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Excel 协作文件内容不存在") async with aiofiles.open(revision.file_uri, "rb") as stream: content = await stream.read(settings.COLLABORATION_MAX_FILE_BYTES + 1) protected_content, backup = apply_workbook_structure_policy( content, file_id=item.id, allow_sheet_structure_edit=requested_sheet_policy, protection_backup=item.sheet_structure_protection_backup, ) item.allow_sheet_structure_edit = requested_sheet_policy item.sheet_structure_protection_backup = backup if protected_content != content: await _replace_revision_bytes(revision, protected_content) item.generation += 1 for key, value in values.items(): setattr(item, key, value) await db.commit() await db.refresh(item) return item async def _next_copy_title(db: AsyncSession, item: CollaborationFile) -> str: suffix = f".{item.extension}" stem = item.title[:-len(suffix)] if item.title.lower().endswith(suffix) else item.title for index in range(1, 1001): marker = " - 副本" if index == 1 else f" - 副本 ({index})" candidate = _safe_title(f"{stem}{marker}", item.extension) duplicate = await db.scalar( select(CollaborationFile.id).where( CollaborationFile.study_id == item.study_id, CollaborationFile.folder_id == item.folder_id, CollaborationFile.title == candidate, CollaborationFile.deleted_at.is_(None), ).limit(1) ) if not duplicate: return candidate raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="无法生成可用的副本文件名") async def copy_file(db: AsyncSession, item: CollaborationFile, user) -> CollaborationFile: await require_file_exporter(db, item, user) revision = await db.get(CollaborationRevision, item.current_revision_id) if item.current_revision_id else None if not revision: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="协作文件尚无可复制版本") source = Path(revision.file_uri) if not source.is_file(): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作文件当前版本不存在") async with aiofiles.open(source, "rb") as stream: content = await stream.read(settings.COLLABORATION_MAX_FILE_BYTES + 1) if not content or len(content) > settings.COLLABORATION_MAX_FILE_BYTES: raise HTTPException(status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail="协作文件大小超出限制") return await _create_file_from_bytes( db, item.study_id, title=await _next_copy_title(db, item), file_type=item.file_type, folder_id=item.folder_id, content=content, source="COPY", user=user, ) async def prepare_download(db: AsyncSession, item: CollaborationFile, user) -> CollaborationRevision: await require_file_exporter(db, item, user) revision = await db.get(CollaborationRevision, item.current_revision_id) if item.current_revision_id else None if not revision: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="协作文件尚无可下载版本") if not Path(revision.file_uri).is_file(): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作文件当前版本不存在") return revision async def move_to_trash(db: AsyncSession, item: CollaborationFile, user) -> None: await require_file_manager(db, item, user) item.status = "DELETED" item.deleted_at = datetime.now(timezone.utc) await db.commit() async def restore_file(db: AsyncSession, item: CollaborationFile, user) -> CollaborationFile: await require_file_manager(db, item, user) item.status = "ACTIVE" item.deleted_at = None await db.commit() await db.refresh(item) return item async def record_export(db: AsyncSession, item: CollaborationFile, user, file_type: str) -> None: await require_file_exporter(db, item, user) async def record_download(db: AsyncSession, item: CollaborationFile, user, file_type: str) -> None: await require_file_exporter(db, item, user) async def list_revisions(db: AsyncSession, item: CollaborationFile) -> list[CollaborationRevisionRead]: rows = (await db.execute( select(CollaborationRevision, User) .outerjoin(User, User.id == CollaborationRevision.created_by) .where( CollaborationRevision.file_id == item.id, CollaborationRevision.deleted_at.is_(None), CollaborationRevision.source != "PERMISSION_CHANGE", ) .order_by(CollaborationRevision.revision_no.desc()) )).all() return [ CollaborationRevisionRead.model_validate(revision).model_copy(update={ "created_by_name": creator.full_name if creator else None, "created_by_avatar_url": creator.avatar_url if creator else None, }) for revision, creator in rows ] async def get_revision_or_404( db: AsyncSession, item: CollaborationFile, revision_id: uuid.UUID ) -> CollaborationRevision: revision = await db.get(CollaborationRevision, revision_id) if not revision or revision.file_id != item.id or getattr(revision, "deleted_at", None) is not None: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作修订不存在") return revision async def delete_revision( db: AsyncSession, item: CollaborationFile, revision_id: uuid.UUID, user, ) -> None: await require_file_manager(db, item, user) revision = await get_revision_or_404(db, item, revision_id) if item.current_revision_id == revision.id: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="当前版本不能删除") revision.deleted_at = datetime.now(timezone.utc) revision.deleted_by = user.id await db.commit() async def update_revision( db: AsyncSession, item: CollaborationFile, revision_id: uuid.UUID, payload: CollaborationRevisionUpdate, user, ) -> CollaborationRevision: await require_file_editor(db, item, user) revision = await get_revision_or_404(db, item, revision_id) revision.change_summary = payload.change_summary await db.commit() await db.refresh(revision) return revision async def prepare_revision_preview( db: AsyncSession, item: CollaborationFile, revision_id: uuid.UUID, user ) -> CollaborationRevision: revision = await get_revision_or_404(db, item, revision_id) if not Path(revision.file_uri).is_file(): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作修订内容不存在") return revision async def copy_revision( db: AsyncSession, item: CollaborationFile, revision_id: uuid.UUID, payload: CollaborationRevisionCopyRequest, user, ) -> CollaborationFile: await require_file_exporter(db, item, user) revision = await get_revision_or_404(db, item, revision_id) source = Path(revision.file_uri) if not source.is_file(): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作修订内容不存在") async with aiofiles.open(source, "rb") as stream: content = await stream.read(settings.COLLABORATION_MAX_FILE_BYTES + 1) if not content or len(content) > settings.COLLABORATION_MAX_FILE_BYTES: raise HTTPException(status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail="协作文件大小超出限制") if payload.folder_id: await _folder_or_404(db, item.study_id, payload.folder_id) return await _create_file_from_bytes( db, item.study_id, title=payload.title, file_type=item.file_type, folder_id=payload.folder_id, content=content, source="COPY", user=user, ) async def restore_revision( db: AsyncSession, item: CollaborationFile, revision_id: uuid.UUID, user, change_summary: str | None ) -> CollaborationRevision: await require_file_editor(db, item, user) revision = await get_revision_or_404(db, item, revision_id) async with aiofiles.open(revision.file_uri, "rb") as stream: content = await stream.read() restored, created = await append_revision( db, item, content, source="RESTORE", created_by=user.id, change_summary=change_summary or f"恢复到修订 {revision.revision_no}", ) if created: locked = await db.get(CollaborationFile, item.id) if locked: locked.generation += 1 await db.commit() await db.refresh(restored) return restored async def list_members(db: AsyncSession, item: CollaborationFile) -> list[CollaborationMemberRead]: rows = (await db.execute( select(CollaborationMember, User) .join(User, User.id == CollaborationMember.user_id) .where(CollaborationMember.file_id == item.id) .order_by(CollaborationMember.role.desc(), User.full_name) )).all() return [CollaborationMemberRead( id=member.id, file_id=member.file_id, user_id=member.user_id, role=member.role, invited_by=member.invited_by, full_name=user.full_name, email=user.email, created_at=member.created_at, ) for member, user in rows] async def upsert_member( db: AsyncSession, item: CollaborationFile, payload: CollaborationMemberUpsert, user ) -> CollaborationMemberRead: await require_file_manager(db, item, user) if payload.user_id == item.owner_id and payload.role != "MANAGER": raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="协作文件所有者必须保留管理者权限") study_member = await member_crud.get_member(db, item.study_id, payload.user_id) if not study_member or not study_member.is_active: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="只能邀请当前项目的有效成员") target = await db.get(User, payload.user_id) if not target: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="授权账号不存在") await _require_role_assignable(db, item.study_id, target, study_member, payload.role) member = await db.scalar(select(CollaborationMember).where( CollaborationMember.file_id == item.id, CollaborationMember.user_id == payload.user_id, )) if member: member.role = payload.role member.invited_by = user.id else: member = CollaborationMember(file_id=item.id, user_id=payload.user_id, role=payload.role, invited_by=user.id) db.add(member) await db.commit() await db.refresh(member) return CollaborationMemberRead( id=member.id, file_id=member.file_id, user_id=member.user_id, role=member.role, invited_by=member.invited_by, full_name=target.full_name if target else "", email=target.email if target else "", created_at=member.created_at, ) async def remove_member(db: AsyncSession, item: CollaborationFile, user_id: uuid.UUID, user) -> None: await require_file_manager(db, item, user) if user_id == item.owner_id: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="不能移除协作文件所有者") member = await db.scalar(select(CollaborationMember).where( CollaborationMember.file_id == item.id, CollaborationMember.user_id == user_id, )) if not member: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作成员不存在") await db.delete(member) await db.commit() async def list_candidates(db: AsyncSession, study_id: uuid.UUID) -> list[CollaborationCandidateRead]: rows = (await db.execute( select(StudyMember, User) .join(User, User.id == StudyMember.user_id) .where(StudyMember.study_id == study_id, StudyMember.is_active.is_(True)) .order_by(User.full_name) )).all() result: list[CollaborationCandidateRead] = [] for membership, user in rows: result.append(CollaborationCandidateRead( user_id=user.id, full_name=user.full_name, email=user.email, role_in_study=membership.role_in_study, can_be_editor=await _role_is_assignable(db, study_id, user, membership, "EDITOR"), can_be_manager=await _role_is_assignable(db, study_id, user, membership, "MANAGER"), )) return result def _edit_request_read(request: CollaborationEditRequest, requester: User) -> CollaborationEditRequestRead: return CollaborationEditRequestRead( id=request.id, file_id=request.file_id, requester_id=request.requester_id, requester_name=requester.full_name, requester_email=requester.email, status=request.status, resolved_by=request.resolved_by, resolved_at=request.resolved_at, created_at=request.created_at, ) async def create_edit_request( db: AsyncSession, item: CollaborationFile, user ) -> CollaborationEditRequestRead: if await can_edit_file(db, item, user): raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="您已拥有此文件的编辑权限") if not item.allow_edit_request: raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="此文件未开放编辑权限申请") membership = await member_crud.get_member(db, item.study_id, user.id) if not membership or not membership.is_active: raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="只有当前项目成员可以申请编辑权限") pending = await db.scalar( select(CollaborationEditRequest).where( CollaborationEditRequest.file_id == item.id, CollaborationEditRequest.requester_id == user.id, CollaborationEditRequest.status == "PENDING", ) ) if pending: return _edit_request_read(pending, user) request = CollaborationEditRequest(file_id=item.id, requester_id=user.id) db.add(request) await db.flush() manager_ids = set((await db.scalars( select(StudyMember.user_id) .outerjoin( CollaborationMember, and_( CollaborationMember.file_id == item.id, CollaborationMember.user_id == StudyMember.user_id, ), ) .where( StudyMember.study_id == item.study_id, StudyMember.is_active.is_(True), or_( StudyMember.user_id == item.owner_id, CollaborationMember.role == "MANAGER", ), ) )).all()) manager_ids.add(item.owner_id) requester_name = str(user.full_name or user.email or "项目成员") await notification_service.create_recipient_notifications( db, study_id=item.study_id, recipient_ids=manager_ids, category="COLLABORATION_EDIT_REQUEST", priority="NORMAL", title="新的编辑权限申请", message=f"{requester_name} 申请编辑“{item.title}”", action_path=f"/knowledge/collaboration?editRequestFile={item.id}", source_type="COLLABORATION_EDIT_REQUEST", source_id=str(request.id), dedupe_key=f"collaboration-edit-request:{request.id}", source_version="PENDING", ) await db.commit() await db.refresh(request) return _edit_request_read(request, user) async def list_edit_requests( db: AsyncSession, item: CollaborationFile, user, *, pending_only: bool = True ) -> list[CollaborationEditRequestRead]: await require_file_manager(db, item, user) stmt = ( select(CollaborationEditRequest, User) .join(User, User.id == CollaborationEditRequest.requester_id) .where(CollaborationEditRequest.file_id == item.id) ) if pending_only: stmt = stmt.where(CollaborationEditRequest.status == "PENDING") rows = (await db.execute(stmt.order_by(CollaborationEditRequest.created_at.desc()))).all() return [_edit_request_read(request, requester) for request, requester in rows] async def resolve_edit_request( db: AsyncSession, item: CollaborationFile, request_id: uuid.UUID, payload: CollaborationEditRequestResolve, user, ) -> CollaborationEditRequestRead: await require_file_manager(db, item, user) request = await db.scalar( select(CollaborationEditRequest) .where(CollaborationEditRequest.id == request_id, CollaborationEditRequest.file_id == item.id) .with_for_update() ) if not request: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="编辑权限申请不存在") if request.status != "PENDING": raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="编辑权限申请已处理") requester = await db.get(User, request.requester_id) membership = await member_crud.get_member(db, item.study_id, request.requester_id) if not requester or not membership or not membership.is_active: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="申请人已不是当前项目的有效成员") if payload.status == "APPROVED": await _require_role_assignable(db, item.study_id, requester, membership, "EDITOR") member = await db.scalar(select(CollaborationMember).where( CollaborationMember.file_id == item.id, CollaborationMember.user_id == request.requester_id, )) if not member: db.add(CollaborationMember( file_id=item.id, user_id=request.requester_id, role="EDITOR", invited_by=user.id, )) elif member.role != "MANAGER": member.role = "EDITOR" member.invited_by = user.id request.status = payload.status request.resolved_by = user.id request.resolved_at = datetime.now(timezone.utc) approved = payload.status == "APPROVED" await notification_service.create_recipient_notifications( db, study_id=item.study_id, recipient_ids=[request.requester_id], category="COLLABORATION_EDIT_REQUEST_RESULT", priority="NORMAL", title="编辑权限申请已通过" if approved else "编辑权限申请未通过", message=f"您对“{item.title}”的编辑权限申请已{'通过' if approved else '被拒绝'}", action_path=( f"/knowledge/collaboration/{item.id}" if approved else "/knowledge/collaboration" ), source_type="COLLABORATION_EDIT_REQUEST_RESULT", source_id=str(request.id), dedupe_key=f"collaboration-edit-request-result:{request.id}", requires_action=False, ) await notification_service.resolve_source_notifications( db, source_type="COLLABORATION_EDIT_REQUEST", source_id=str(request.id), ) await db.commit() await db.refresh(request) return _edit_request_read(request, requester) async def transfer_ownership( db: AsyncSession, item: CollaborationFile, payload: CollaborationOwnershipTransferRequest, user, ) -> CollaborationFile: await require_ownership_transfer(db, item, user) locked = await db.scalar(select(CollaborationFile).where(CollaborationFile.id == item.id).with_for_update()) if not locked: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="协作文件不存在") if payload.new_owner_id == locked.owner_id: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="所选联系人已经是文档所有者") membership = await member_crud.get_member(db, locked.study_id, payload.new_owner_id) target = await db.get(User, payload.new_owner_id) if not target or not membership or not membership.is_active: raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="只能转让给当前项目的有效成员") await _require_role_assignable(db, locked.study_id, target, membership, "MANAGER") previous_owner_id = locked.owner_id member = await db.scalar(select(CollaborationMember).where( CollaborationMember.file_id == locked.id, CollaborationMember.user_id == payload.new_owner_id, )) if member: member.role = "MANAGER" member.invited_by = user.id else: db.add(CollaborationMember( file_id=locked.id, user_id=payload.new_owner_id, role="MANAGER", invited_by=user.id, )) previous_owner_member = await db.scalar(select(CollaborationMember).where( CollaborationMember.file_id == locked.id, CollaborationMember.user_id == previous_owner_id, )) if previous_owner_member: previous_owner_member.role = "MANAGER" locked.owner_id = payload.new_owner_id locked.updated_at = datetime.now(timezone.utc) await db.commit() await db.refresh(locked) return locked