fix(onlyoffice): 修复云端台账保存与备份重开

This commit is contained in:
Cheng Zhou
2026-09-04 11:18:35 +08:00
parent 684b8b51bb
commit 90573f050f
8 changed files with 384 additions and 42 deletions
@@ -59,33 +59,55 @@ async def _active_session(
).order_by(CollaborationSession.created_at.desc())
)
if session:
if session.status == "ERROR":
recovered = await _recover_forgotten_content(session.document_key, item.file_type)
if recovered is not None:
revision, _ = await collaboration_service.append_revision(
db,
item,
recovered,
source="SERVER_RECOVERY",
created_by=user_id,
change_summary="自动恢复在线文档服务器备份",
)
session.base_revision_id = revision.id
# The failed key points to Document Server's recovery cache. A new
# generation must use a new key or every subsequent open falls back
# to the same unsaved backup again.
item.generation += 1
session.status = "RECOVERED" if recovered is not None else "CLOSED"
session.closed_at = datetime.now(timezone.utc)
await db.commit()
session = None
else:
if session.status != "ACTIVE":
session.status = "ACTIVE"
session.closed_at = None
await db.commit()
await db.refresh(session)
return session
should_recover = session.status == "ERROR"
if session.status == "ACTIVE":
created_at = session.created_at
if created_at.tzinfo is None:
created_at = created_at.replace(tzinfo=timezone.utc)
if (
session.last_callback_at is None
and datetime.now(timezone.utc) - created_at < timedelta(seconds=15)
):
# The editor config can be requested twice before the first
# browser has connected and emitted status 1. Keep a short
# connection grace period so the second request does not retire
# the freshly issued key as a false stale session.
return session
live_users = await _document_server_users(session.document_key)
if live_users:
return session
# A service restart or rejected final callback can leave the row
# ACTIVE after Document Server has already retired the editing
# process. Reusing that key opens its forgotten copy as an
# unbound server backup, so retire it exactly like a final callback.
should_recover = True
# A final callback ends the editing lifecycle for this key. Never
# reactivate it: Document Server can retain a cached or forgotten copy
# for the old key, especially across a container restart.
recovered = (
await _recover_forgotten_content(session.document_key, item.file_type)
if should_recover
else None
)
if recovered is not None:
revision, _ = await collaboration_service.append_revision(
db,
item,
recovered,
source="SERVER_RECOVERY",
created_by=user_id,
change_summary="自动恢复在线文档服务器备份",
)
session.base_revision_id = revision.id
# The retired key can still point to Document Server's cache. A new
# generation must use a new key or a later open can fall back to the
# same server-side copy again.
item.generation += 1
session.status = "RECOVERED" if recovered is not None else "CLOSED"
session.closed_at = datetime.now(timezone.utc)
await db.commit()
session = None
if not item.current_revision_id:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="协作文件尚无可编辑内容")
session = CollaborationSession(
@@ -413,6 +435,67 @@ async def _recover_forgotten_content(document_key: str, file_type: str) -> bytes
return content
async def _document_server_users(document_key: str) -> list[str]:
"""Return live editor ids for a key without trusting stale database state."""
command = {"c": "info", "key": document_key}
token = jwt.encode(command, settings.ONLYOFFICE_JWT_SECRET or "", algorithm="HS256")
command_url = f"{settings.ONLYOFFICE_INTERNAL_URL.rstrip('/')}/command"
try:
async with httpx.AsyncClient(timeout=10.0, follow_redirects=False) as client:
response = await client.post(
command_url,
params={"shardkey": document_key},
json={**command, "token": token},
)
if response.status_code != status.HTTP_200_OK:
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="在线文档会话检查服务暂不可用",
)
payload = response.json()
except (httpx.HTTPError, ValueError) as exc:
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="在线文档会话检查服务暂不可用",
) from exc
error = payload.get("error") if isinstance(payload, dict) else None
if error == 1:
# The database can retain an ACTIVE row after an interrupted callback,
# while Document Server no longer has a live editing process for it.
return []
users = payload.get("users") if isinstance(payload, dict) else None
if (
error != 0
or payload.get("key") != document_key
or not isinstance(users, list)
or any(not isinstance(user_id, str) or not user_id for user_id in users)
):
raise HTTPException(
status_code=status.HTTP_502_BAD_GATEWAY,
detail="在线文档服务器返回的会话信息无效",
)
return list(dict.fromkeys(users))
async def list_live_editing_sessions(db: AsyncSession) -> list[tuple[str, int]]:
"""List file titles and live editor counts for deployment safety checks."""
rows = (
await db.execute(
select(CollaborationSession.document_key, CollaborationFile.title)
.join(CollaborationFile, CollaborationFile.id == CollaborationSession.file_id)
.where(CollaborationSession.status == "ACTIVE")
.order_by(CollaborationFile.title)
)
).all()
active: list[tuple[str, int]] = []
for document_key, title in rows:
users = await _document_server_users(document_key)
if users:
active.append((title, len(users)))
return active
async def _callback_user(
db: AsyncSession, payload: CollaborationCallbackPayload, session: CollaborationSession
) -> User | None:
@@ -0,0 +1,44 @@
"""Abort a deployment when ONLYOFFICE still has live collaborative editors."""
import asyncio
import sys
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parent.parent
sys.path.insert(0, str(PROJECT_ROOT))
from fastapi import HTTPException # noqa: E402
from app.core.config import settings # noqa: E402
from app.db.session import SessionLocal # noqa: E402
from app.services.onlyoffice_collaboration_service import list_live_editing_sessions # noqa: E402
async def async_main() -> int:
if not settings.ONLYOFFICE_ENABLED:
print("ONLYOFFICE 未启用,跳过在线编辑会话检查")
return 0
try:
async with SessionLocal() as db:
active = await list_live_editing_sessions(db)
except HTTPException as exc:
print(f"无法确认 ONLYOFFICE 在线编辑状态:{exc.detail}", file=sys.stderr)
return 1
if not active:
print("未检测到 ONLYOFFICE 在线编辑者")
return 0
print("检测到仍在进行的 ONLYOFFICE 在线编辑,会中止本次部署:", file=sys.stderr)
for title, count in active:
print(f"- {title}:{count} 人在线", file=sys.stderr)
print("请通知用户退出编辑器,等待最终保存完成后重新执行部署。", file=sys.stderr)
return 2
def main() -> None:
raise SystemExit(asyncio.run(async_main()))
if __name__ == "__main__":
main()
+96 -1
View File
@@ -4,7 +4,7 @@ import uuid
import zipfile
from datetime import datetime, timedelta, timezone
from types import SimpleNamespace
from unittest.mock import ANY, AsyncMock
from unittest.mock import ANY, AsyncMock, call
import pytest
from fastapi import HTTPException
@@ -304,6 +304,18 @@ def test_result_download_url_is_restricted_and_public_proxy_urls_are_rewritten(m
assert onlyoffice_collaboration_service._validate_result_url(
"http://localhost:8888/onlyoffice/cache/result.docx?token=signed"
) == "http://onlyoffice/cache/result.docx?token=signed"
with pytest.raises(HTTPException) as mismatched_origin:
onlyoffice_collaboration_service._validate_result_url(
"https://ctms.example.com/onlyoffice/cache/result.xlsx?token=signed"
)
assert mismatched_origin.value.status_code == 422
monkeypatch.setattr(settings, "FRONTEND_PUBLIC_URL", "https://ctms.example.com")
assert onlyoffice_collaboration_service._validate_result_url(
"https://ctms.example.com/onlyoffice/cache/result.xlsx?token=signed"
) == "http://onlyoffice/cache/result.xlsx?token=signed"
monkeypatch.setattr(settings, "FRONTEND_PUBLIC_URL", "http://localhost:8888")
for value in (
"http://backend:8000/internal/file",
"http://onlyoffice.evil.example/cache/result.docx",
@@ -384,6 +396,89 @@ async def test_forgotten_document_command_rejects_a_damaged_backup(monkeypatch):
assert error.value.status_code == 502
@pytest.mark.asyncio
async def test_document_server_info_command_returns_unique_live_users(monkeypatch):
key = "ctms-collab-live-key"
request = {}
class FakeResponse:
status_code = 200
@staticmethod
def json():
return {"error": 0, "key": key, "users": ["user-1", "user-1", "user-2"]}
class FakeClient:
async def __aenter__(self):
return self
async def __aexit__(self, *_args):
return None
async def post(self, url, *, params, json):
request.update(url=url, params=params, body=json)
return FakeResponse()
monkeypatch.setattr(
onlyoffice_collaboration_service.httpx,
"AsyncClient",
lambda **_kwargs: FakeClient(),
)
users = await onlyoffice_collaboration_service._document_server_users(key)
assert users == ["user-1", "user-2"]
assert request["url"] == "http://onlyoffice/command"
assert request["params"] == {"shardkey": key}
assert jwt.decode(
request["body"]["token"], settings.ONLYOFFICE_JWT_SECRET, algorithms=["HS256"]
) == {"c": "info", "key": key}
@pytest.mark.asyncio
async def test_document_server_info_command_treats_unknown_key_as_no_live_users(monkeypatch):
class FakeResponse:
status_code = 200
@staticmethod
def json():
return {"error": 1}
class FakeClient:
async def __aenter__(self):
return self
async def __aexit__(self, *_args):
return None
async def post(self, *_args, **_kwargs):
return FakeResponse()
monkeypatch.setattr(
onlyoffice_collaboration_service.httpx,
"AsyncClient",
lambda **_kwargs: FakeClient(),
)
assert await onlyoffice_collaboration_service._document_server_users("retired-key") == []
@pytest.mark.asyncio
async def test_live_editing_session_check_filters_stale_active_rows(monkeypatch):
rows = SimpleNamespace(all=lambda: [
("live-key", "正在编辑.xlsx"),
("stale-key", "陈旧记录.xlsx"),
])
db = SimpleNamespace(execute=AsyncMock(return_value=rows))
lookup = AsyncMock(side_effect=[["user-1", "user-2"], []])
monkeypatch.setattr(onlyoffice_collaboration_service, "_document_server_users", lookup)
active = await onlyoffice_collaboration_service.list_live_editing_sessions(db)
assert active == [("正在编辑.xlsx", 2)]
assert lookup.await_args_list == [call("live-key"), call("stale-key")]
@pytest.mark.asyncio
async def test_editor_config_grants_edit_only_after_collaboration_permission(monkeypatch, tmp_path):
user_id = uuid.uuid4()
+72 -1
View File
@@ -1,7 +1,7 @@
import io
import uuid
import zipfile
from datetime import timezone
from datetime import datetime, timezone
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import AsyncMock
@@ -44,6 +44,7 @@ async def env(monkeypatch, tmp_path):
monkeypatch.setattr(collaboration, "COLLABORATION_ROOT", tmp_path)
monkeypatch.setattr(settings, "ONLYOFFICE_JWT_SECRET", "ledger-test-secret-long-enough-for-tests")
monkeypatch.setattr(onlyoffice_service, "ensure_onlyoffice_available", AsyncMock())
monkeypatch.setattr(office, "_document_server_users", AsyncMock(return_value=["live-user"]))
async with AsyncSession(engine, expire_on_commit=False) as db:
users = [User(id=uuid.uuid4(), email=f"ledger-{i}@example.com", password_hash="hash",
full_name=f"Ledger user {i}", clinical_department="test", is_admin=i == 0,
@@ -121,6 +122,76 @@ async def test_failed_ledger_session_recovers_server_backup_under_a_new_document
recovery.assert_awaited_once_with(failed.document_key, "cell")
@pytest.mark.asyncio
async def test_closed_ledger_session_starts_a_new_key_instead_of_reopening_server_cache(env, monkeypatch):
item = await env.db.get(CollaborationFile, env.initial[0].id)
first = await office.build_editor_config(env.db, item, env.admin)
closed = await env.db.scalar(select(CollaborationSession).where(
CollaborationSession.file_id == item.id,
CollaborationSession.generation == item.generation,
))
await office.process_callback(env.db, closed.id, CollaborationCallbackPayload(
key=closed.document_key,
status=4,
))
recovery = AsyncMock(return_value=None)
monkeypatch.setattr(office, "_recover_forgotten_content", recovery)
reopened = await office.build_editor_config(env.db, item, env.admin)
await env.db.refresh(item)
await env.db.refresh(closed)
sessions = (await env.db.scalars(select(CollaborationSession).where(
CollaborationSession.file_id == item.id,
).order_by(CollaborationSession.generation))).all()
assert first.config["document"]["key"] != reopened.config["document"]["key"]
assert item.generation == 2
assert closed.status == "CLOSED"
assert [session.generation for session in sessions] == [1, 2]
recovery.assert_not_awaited()
@pytest.mark.asyncio
async def test_stale_active_ledger_session_starts_a_new_key_when_document_server_has_no_users(
env, monkeypatch
):
item = await env.db.get(CollaborationFile, env.initial[0].id)
first = await office.build_editor_config(env.db, item, env.admin)
stale = await env.db.scalar(select(CollaborationSession).where(
CollaborationSession.file_id == item.id,
CollaborationSession.generation == item.generation,
))
stale.last_callback_at = datetime.now(timezone.utc)
await env.db.commit()
live_users = AsyncMock(return_value=[])
recovery = AsyncMock(return_value=None)
monkeypatch.setattr(office, "_document_server_users", live_users)
monkeypatch.setattr(office, "_recover_forgotten_content", recovery)
reopened = await office.build_editor_config(env.db, item, env.admin)
await env.db.refresh(item)
await env.db.refresh(stale)
assert first.config["document"]["key"] != reopened.config["document"]["key"]
assert item.generation == 2
assert stale.status == "CLOSED"
live_users.assert_awaited_once_with(stale.document_key)
recovery.assert_awaited_once_with(stale.document_key, "cell")
@pytest.mark.asyncio
async def test_new_active_session_is_reused_during_browser_connection_grace(env, monkeypatch):
item = await env.db.get(CollaborationFile, env.initial[0].id)
first = await office.build_editor_config(env.db, item, env.admin)
live_users = AsyncMock(return_value=[])
monkeypatch.setattr(office, "_document_server_users", live_users)
repeated = await office.build_editor_config(env.db, item, env.admin)
assert first.config["document"]["key"] == repeated.config["document"]["key"]
live_users.assert_not_awaited()
@pytest.mark.asyncio
async def test_account_grants_are_independent_and_all_file_routes_require_access(env):
await grant(env, env.editor, "EDITOR")