Files
govoplan-files/src/govoplan_files/backend/routes/connector_io.py
T
zemion 378f4d6ac5
Module Package Release / publish-packages (push) Successful in 12s
feat(files): orchestrate connector folder sync
2026-08-21 22:49:47 +02:00

632 lines
22 KiB
Python

from __future__ import annotations
import json
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_any_scope, require_scope
from govoplan_core.audit.logging import audit_from_principal
from govoplan_files.backend.schemas import (
FileConnectorBrowseItem,
FileConnectorBrowseResponse,
FileConnectorFolderSyncItemResponse,
FileConnectorFolderSyncRequest,
FileConnectorFolderSyncResponse,
FileConnectorFolderSyncSummary,
FileConnectorImportRequest,
FileConnectorSyncResponse,
FileConnectorWriteRequest,
FileConnectorWriteResponse,
FileUploadResponse,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.paths import UnsafeFilePathError
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.connector_browse import (
ConnectorBrowseError,
ConnectorBrowseUnsupported,
browse_connector_profile,
normalize_connector_browse_path,
)
from govoplan_files.backend.storage.connector_imports import (
ConnectorImportError,
ConnectorImportUnsupported,
)
from govoplan_files.backend.storage.connector_folder_sync import (
connector_relative_path,
discover_connector_folder,
join_connector_path,
)
from govoplan_files.backend.storage.connector_deployment import (
connector_effective_endpoint_url,
)
from govoplan_files.backend.storage.connector_policy import (
ConnectorAccessRequest,
ConnectorPolicyDenied,
connector_policy_decision,
)
from govoplan_files.backend.storage.files import (
create_file_asset,
get_asset_for_user,
read_asset_bytes,
sync_file_asset_from_source,
)
from govoplan_files.backend.storage.connector_spaces import (
connector_space_owner_id,
get_connector_space_for_user,
)
from govoplan_files.backend.storage.connector_writes import write_connector_file
from govoplan_files.backend.route_support import (
_asset_response,
_audit_connector_imports,
_audit_connector_sync,
_connector_browse_next_token,
_connector_policy_error,
_connector_space_policy_decision,
_download_connector_payload,
_ensure_campaign_file_access,
_http_error,
_is_admin,
_visible_connector_profile,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.post(
"/connector-spaces/{space_id}/sync",
response_model=FileConnectorFolderSyncResponse,
)
def sync_connector_space_folder(
space_id: str,
payload: FileConnectorFolderSyncRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:upload")),
):
try:
space = get_connector_space_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
space_id=space_id,
is_admin=_is_admin(principal),
)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc, not_found=True) from exc
try:
if space.sync_mode != "manual":
raise FileStorageError("This connector space is not configured for manual sync")
profile = _visible_connector_profile(
session, principal, space.connector_profile_id
)
remote_path = join_connector_path(space.remote_path, payload.path)
target_folder = join_connector_path(payload.target_folder)
decision = _connector_space_policy_decision(
profile,
library_id=space.library_id,
remote_path=remote_path,
operation="sync",
)
if not decision.allowed:
raise ConnectorPolicyDenied(decision)
discovery = discover_connector_folder(
profile,
path=remote_path,
library_id=space.library_id,
recursive=payload.recursive,
max_files=payload.max_files,
max_depth=payload.max_depth,
)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except ConnectorBrowseUnsupported as exc:
session.rollback()
raise HTTPException(
status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)
) from exc
except (ConnectorBrowseError, FileStorageError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc
owner_id = connector_space_owner_id(space)
items: list[FileConnectorFolderSyncItemResponse] = [
FileConnectorFolderSyncItemResponse(
source_path=skipped.path,
action="skipped",
detail=skipped.reason,
)
for skipped in discovery.skipped
]
counts = {
"created": 0,
"updated": 0,
"unchanged": 0,
"skipped": len(discovery.skipped),
"conflicts": 0,
"policy_denied": 0,
"failed": 0,
}
for source in discovery.files:
source_path = normalize_connector_browse_path(source.path)
try:
relative_path = connector_relative_path(
space_root=space.remote_path,
item_path=source_path,
)
target_path = join_connector_path(target_folder, relative_path)
except (ConnectorBrowseError, ValueError) as exc:
counts["failed"] += 1
items.append(
FileConnectorFolderSyncItemResponse(
source_path=source_path,
action="failed",
detail=str(exc),
)
)
continue
try:
with session.begin_nested():
file_payload = FileConnectorImportRequest(
library_id=space.library_id or "",
path=source_path,
owner_type=space.owner_type,
owner_id=owner_id,
target_path=target_path,
conflict_strategy=payload.conflict_strategy,
source_revision=source.etag,
metadata={
**dict(source.metadata),
**payload.metadata,
"connector_space_id": space.id,
"browse_name": source.name,
"browse_path": source_path,
"browse_modified_at": source.modified_at,
"browse_etag": source.etag,
"folder_sync": True,
},
)
_source_path, downloaded, metadata = _download_connector_payload(
profile, file_payload, operation="sync"
)
stored, sync_action, previous_version_id = (
sync_file_asset_from_source(
session,
tenant_id=principal.tenant_id,
owner_type=space.owner_type,
owner_id=owner_id,
user_id=principal.user.id,
filename=downloaded.filename,
data=downloaded.data,
display_path=target_path,
content_type=downloaded.content_type,
metadata=metadata,
conflict_strategy=payload.conflict_strategy,
is_admin=_is_admin(principal),
)
)
_audit_connector_sync(
session,
principal,
stored.asset,
sync_action=sync_action,
previous_version_id=previous_version_id,
)
file_response = _asset_response(
session, stored.asset, include_shares=True
)
counts[sync_action] += 1
items.append(
FileConnectorFolderSyncItemResponse(
source_path=source_path,
target_path=target_path,
action=sync_action,
file=file_response,
previous_version_id=previous_version_id,
current_version_id=stored.version.id,
source_revision=file_response.source_revision,
)
)
except ConnectorPolicyDenied as exc:
counts["policy_denied"] += 1
items.append(
FileConnectorFolderSyncItemResponse(
source_path=source_path,
target_path=target_path,
action="policy_denied",
source_revision=source.etag,
detail=str(exc),
policy_decision=exc.decision.to_dict(),
)
)
except FileStorageError as exc:
detail = str(exc)
if detail.startswith("Skipped upload target:"):
action = "skipped"
counts["skipped"] += 1
elif detail.startswith("Target file already exists:"):
action = "conflict"
counts["conflicts"] += 1
else:
action = "failed"
counts["failed"] += 1
items.append(
FileConnectorFolderSyncItemResponse(
source_path=source_path,
target_path=target_path,
action=action,
source_revision=source.etag,
detail=detail,
)
)
except (
ConnectorImportError,
UnsafeFilePathError,
OSError,
ValueError,
json.JSONDecodeError,
) as exc:
counts["failed"] += 1
items.append(
FileConnectorFolderSyncItemResponse(
source_path=source_path,
target_path=target_path,
action="failed",
source_revision=source.etag,
detail=str(exc),
)
)
summary = FileConnectorFolderSyncSummary(
discovered=len(discovery.files),
**counts,
)
audit_from_principal(
session,
principal,
action="files.connector.folder_synced",
object_type="file_connector_space",
object_id=space.id,
details={
"connector_profile_id": profile.id,
"provider": profile.provider,
"library_id": space.library_id,
"remote_path": remote_path,
"target_folder": target_folder,
"recursive": payload.recursive,
"truncated": discovery.truncated,
"summary": summary.model_dump(),
"results": [
{
"source_path": item.source_path,
"target_path": item.target_path,
"action": item.action,
"file_id": item.file.id if item.file else None,
"current_version_id": item.current_version_id,
}
for item in items
],
},
)
session.commit()
return FileConnectorFolderSyncResponse(
connector_space_id=space.id,
connector_profile_id=profile.id,
provider=profile.provider,
remote_path=remote_path,
target_folder=target_folder,
recursive=payload.recursive,
truncated=discovery.truncated,
summary=summary,
items=items,
)
@router.post(
"/connector-spaces/{space_id}/write-back",
response_model=FileConnectorWriteResponse,
)
def write_back_connector_file(
space_id: str,
payload: FileConnectorWriteRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:connector:write")),
):
try:
space = get_connector_space_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
space_id=space_id,
is_admin=_is_admin(principal),
)
if space.read_only:
raise FileStorageError("This connector space is read-only")
profile = _visible_connector_profile(
session, principal, space.connector_profile_id
)
requested_path = normalize_connector_browse_path(payload.remote_path)
remote_path = normalize_connector_browse_path(
"/".join(part for part in (space.remote_path, requested_path) if part)
)
decision = connector_policy_decision(
ConnectorAccessRequest(
connector_id=profile.id,
credential_id=profile.credential_profile_id,
provider=profile.provider,
external_path=remote_path,
external_url=connector_effective_endpoint_url(
provider=profile.provider,
endpoint_url=profile.endpoint_url,
metadata=profile.metadata,
),
operation="write",
),
profile.policy_sources,
)
if not decision.allowed:
raise ConnectorPolicyDenied(decision)
asset = get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=payload.file_id,
require_write=True,
is_admin=_is_admin(principal),
)
data, version, blob = read_asset_bytes(session, asset)
connector_library_id = space.library_id
content_type = blob.content_type
# Close the read snapshot before the independent recovery transaction
# records authority for the external effect. This avoids upgrading an
# older SQLite read snapshot after the ledger commit.
session.commit()
result = write_connector_file(
profile,
tenant_id=principal.tenant_id,
library_id=connector_library_id,
remote_path=remote_path,
data=data,
content_type=content_type,
expected_revision=payload.expected_revision,
idempotency_key=payload.idempotency_key,
)
audit_from_principal(
session,
principal,
action="files.connector.written",
object_type="file",
object_id=asset.id,
details={
"file_version_id": version.id,
"file_blob_id": blob.id,
"checksum_sha256": blob.checksum_sha256,
"connector_space_id": space.id,
"connector_profile_id": profile.id,
"remote_path": remote_path,
"recovery_operation_id": result.recovery_operation_id,
"recovery_status": result.status,
"revision": result.revision,
},
)
session.commit()
return FileConnectorWriteResponse(
recovery_operation_id=result.recovery_operation_id,
status=result.status,
replayed=result.replayed,
provider=result.provider,
remote_path=result.remote_path,
revision=result.revision,
checksum_sha256=result.checksum_sha256,
size_bytes=result.size_bytes,
)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ConnectorBrowseError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.post(
"/connectors/profiles/{profile_id}/import", response_model=FileUploadResponse
)
def import_connector_file(
profile_id: str,
payload: FileConnectorImportRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:upload")),
):
try:
if payload.campaign_id:
_ensure_campaign_file_access(session, principal, payload.campaign_id)
profile = _visible_connector_profile(
session, principal, profile_id, campaign_id=payload.campaign_id
)
_source_path, downloaded, metadata = _download_connector_payload(
profile, payload, operation="import"
)
target_owner = payload.owner_id or principal.user.id
stored = create_file_asset(
session,
tenant_id=principal.tenant_id,
owner_type=payload.owner_type,
owner_id=target_owner,
user_id=principal.user.id,
filename=downloaded.filename,
data=downloaded.data,
folder=payload.target_folder,
display_path=payload.target_path,
content_type=downloaded.content_type,
metadata=metadata,
campaign_id=payload.campaign_id,
conflict_strategy=payload.conflict_strategy,
is_admin=_is_admin(principal),
)
_audit_connector_imports(session, principal, [stored.asset])
session.commit()
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except ConnectorImportUnsupported as exc:
session.rollback()
raise HTTPException(
status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)
) from exc
except (
ConnectorImportError,
FileStorageError,
UnsafeFilePathError,
ValueError,
json.JSONDecodeError,
) as exc:
session.rollback()
raise _http_error(exc) from exc
return FileUploadResponse(
files=[_asset_response(session, stored.asset, include_shares=True)]
)
@router.post(
"/connectors/profiles/{profile_id}/sync", response_model=FileConnectorSyncResponse
)
def sync_connector_file(
profile_id: str,
payload: FileConnectorImportRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:upload")),
):
try:
if payload.campaign_id:
_ensure_campaign_file_access(session, principal, payload.campaign_id)
profile = _visible_connector_profile(
session, principal, profile_id, campaign_id=payload.campaign_id
)
_source_path, downloaded, metadata = _download_connector_payload(
profile, payload, operation="sync"
)
target_owner = payload.owner_id or principal.user.id
stored, sync_action, previous_version_id = sync_file_asset_from_source(
session,
tenant_id=principal.tenant_id,
owner_type=payload.owner_type,
owner_id=target_owner,
user_id=principal.user.id,
filename=downloaded.filename,
data=downloaded.data,
folder=payload.target_folder,
display_path=payload.target_path,
content_type=downloaded.content_type,
metadata=metadata,
campaign_id=payload.campaign_id,
conflict_strategy=payload.conflict_strategy,
is_admin=_is_admin(principal),
)
_audit_connector_sync(
session,
principal,
stored.asset,
sync_action=sync_action,
previous_version_id=previous_version_id,
)
session.commit()
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except ConnectorImportUnsupported as exc:
session.rollback()
raise HTTPException(
status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)
) from exc
except (
ConnectorImportError,
FileStorageError,
UnsafeFilePathError,
ValueError,
json.JSONDecodeError,
) as exc:
session.rollback()
raise _http_error(exc) from exc
return FileConnectorSyncResponse(
file=_asset_response(session, stored.asset, include_shares=True),
action=sync_action,
previous_version_id=previous_version_id,
current_version_id=stored.version.id,
)
@router.get(
"/connectors/profiles/{profile_id}/browse",
response_model=FileConnectorBrowseResponse,
)
def browse_connector_profile_items(
profile_id: str,
path: str | None = None,
library_id: str | None = None,
continuation_token: str | None = None,
campaign_id: str | None = None,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:read",
"files:file:upload",
"files:file:download",
"files:file:admin",
"system:settings:read",
"admin:settings:read",
)
),
):
try:
profile = _visible_connector_profile(
session, principal, profile_id, campaign_id=campaign_id
)
browse_path = normalize_connector_browse_path(path)
decision = connector_policy_decision(
ConnectorAccessRequest(
connector_id=profile.id,
credential_id=profile.credential_profile_id,
provider=profile.provider,
external_path=browse_path,
external_url=connector_effective_endpoint_url(
provider=profile.provider,
endpoint_url=profile.endpoint_url,
metadata=profile.metadata,
),
operation="browse",
),
profile.policy_sources,
)
if not decision.allowed:
raise ConnectorPolicyDenied(decision)
items = browse_connector_profile(
profile,
path=browse_path,
library_id=library_id,
continuation_token=continuation_token,
)
except ConnectorPolicyDenied as exc:
raise _connector_policy_error(exc) from exc
except ConnectorBrowseUnsupported as exc:
raise HTTPException(
status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)
) from exc
except (ConnectorBrowseError, OSError, ValueError, json.JSONDecodeError) as exc:
raise _http_error(exc) from exc
return FileConnectorBrowseResponse(
profile_id=profile.id,
provider=profile.provider,
path=browse_path,
library_id=library_id,
next_continuation_token=_connector_browse_next_token(items),
has_more=any(bool(item.metadata.get("listing_truncated")) for item in items),
decision=decision.to_dict(),
items=[FileConnectorBrowseItem(**item.to_response()) for item in items],
)