refactor(api): split files workflow routers

This commit is contained in:
2026-07-29 20:11:17 +02:00
parent 5b868272b9
commit 86a905a3a7
19 changed files with 4918 additions and 3565 deletions

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -0,0 +1 @@
"""Focused HTTP route modules for the Files API."""

View File

@@ -0,0 +1,134 @@
from __future__ import annotations
from io import BytesIO
from fastapi import APIRouter, Depends
from fastapi.responses import StreamingResponse
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.schemas import (
BulkDeleteRequest,
BulkDeleteResponse,
FileAssetResponse,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.files import (
get_asset_for_user,
read_asset_bytes,
soft_delete_assets,
)
from govoplan_files.backend.route_support import (
_asset_response,
_attachment_disposition,
_audit_connector_event,
_http_error,
_is_admin,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.get("/{file_id}", response_model=FileAssetResponse)
def get_file(
file_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
try:
asset = get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
is_admin=_is_admin(principal),
)
return _asset_response(session, asset, include_shares=True)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
@router.get("/{file_id}/download")
def download_file(
file_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:download")),
):
try:
asset = get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
is_admin=_is_admin(principal),
)
data, version, blob = read_asset_bytes(session, asset)
_audit_connector_event(
session,
principal,
action="files.connector.accessed",
asset=asset,
version=version,
blob=blob,
operation="download",
commit=True,
)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
headers = {"Content-Disposition": _attachment_disposition(asset.filename)}
return StreamingResponse(
BytesIO(data),
media_type=blob.content_type or "application/octet-stream",
headers=headers,
)
@router.delete("/{file_id}", response_model=BulkDeleteResponse)
def delete_file(
file_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:delete")),
):
try:
asset = get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
require_write=True,
is_admin=_is_admin(principal),
)
count = soft_delete_assets(session, [asset])
session.commit()
return BulkDeleteResponse(deleted_count=count)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc, not_found=True) from exc
@router.post("/bulk-delete", response_model=BulkDeleteResponse)
def bulk_delete_files(
payload: BulkDeleteRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:delete")),
):
try:
assets = [
get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
require_write=True,
is_admin=_is_admin(principal),
)
for file_id in payload.file_ids
]
count = soft_delete_assets(session, assets)
session.commit()
return BulkDeleteResponse(deleted_count=count)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc) from exc

View File

@@ -0,0 +1,252 @@
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_files.backend.schemas import (
FileConnectorBrowseItem,
FileConnectorBrowseResponse,
FileConnectorImportRequest,
FileConnectorSyncResponse,
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_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,
sync_file_asset_from_source,
)
from govoplan_files.backend.route_support import (
_asset_response,
_audit_connector_imports,
_audit_connector_sync,
_connector_browse_next_token,
_connector_policy_error,
_download_connector_payload,
_ensure_campaign_file_access,
_http_error,
_is_admin,
_visible_connector_profile,
)
router = APIRouter(prefix="/files", tags=["files"])
@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],
)

View File

@@ -0,0 +1,261 @@
from __future__ import annotations
import json
from fastapi import APIRouter, Depends
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_any_scope
from govoplan_files.backend.change_tracking import (
FILES_CONNECTOR_PROFILES_COLLECTION,
)
from govoplan_files.backend.schemas import (
FileConnectorProfileResponse,
FileConnectorProfileUpdateRequest,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.connector_credential_deletion import (
delete_connector_profile_row,
)
from govoplan_files.backend.storage.connector_deployment import (
connector_effective_endpoint_url,
)
from govoplan_files.backend.storage.connector_profile_store import (
connector_profile_from_row,
get_connector_profile_row,
update_connector_profile_row,
)
from govoplan_files.backend.storage.connector_policy import (
ConnectorPolicyDenied,
)
from govoplan_files.backend.route_support import (
FILES_CONNECTOR_PROFILE_RESOURCE,
_can_read_disabled_connector_profiles,
_connector_policy_error,
_credential_row_for_profile,
_ensure_connector_configuration_allowed,
_ensure_connector_local_policy_allowed,
_http_error,
_record_connector_settings_change,
_require_connector_profile_write,
_visible_connector_profile,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.get(
"/connectors/profiles/{profile_id}", response_model=FileConnectorProfileResponse
)
def get_connector_profile(
profile_id: str,
campaign_id: str | None = None,
include_disabled: bool = False,
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,
include_disabled=include_disabled
and _can_read_disabled_connector_profiles(principal),
include_effective_policy=False,
)
except (OSError, ValueError, json.JSONDecodeError) as exc:
raise _http_error(exc) from exc
return FileConnectorProfileResponse(**profile.to_response())
@router.patch(
"/connectors/profiles/{profile_id}", response_model=FileConnectorProfileResponse
)
def update_connector_profile(
profile_id: str,
payload: FileConnectorProfileUpdateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
try:
row = get_connector_profile_row(
session,
tenant_id=principal.tenant_id,
profile_id=profile_id,
include_disabled=True,
)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
_require_connector_profile_write(principal, row.scope_type)
credentials = payload.credentials
try:
credential_profile_id = (
payload.credential_profile_id
if payload.credential_profile_id is not None
else row.credential_profile_id
)
provider = payload.provider if payload.provider is not None else row.provider
credential_row = _credential_row_for_profile(
session,
principal,
credential_profile_id=credential_profile_id,
provider=provider,
profile_id=row.id,
scope_type=row.scope_type,
scope_id=row.scope_id,
include_disabled=True,
)
_ensure_connector_configuration_allowed(
session,
principal,
connector_id=row.id,
credential_id=credential_profile_id,
provider=provider,
endpoint_url=connector_effective_endpoint_url(
provider=provider,
endpoint_url=payload.endpoint_url
if payload.endpoint_url is not None
else row.endpoint_url,
metadata=payload.metadata
if payload.metadata is not None
else row.metadata_,
),
base_path=payload.base_path
if payload.base_path is not None
else row.base_path,
scope_type=row.scope_type,
scope_id=row.scope_id,
operation="configure",
)
if payload.policy is not None:
_ensure_connector_local_policy_allowed(
session,
principal,
scope_type=row.scope_type,
scope_id=row.scope_id,
policy=payload.policy,
)
update_connector_profile_row(
session,
row,
user_id=principal.user.id,
label=payload.label,
provider=payload.provider,
endpoint_url=payload.endpoint_url,
base_path=payload.base_path,
enabled=payload.enabled,
credential_profile_id=payload.credential_profile_id,
credential_mode=payload.credential_mode,
username=credentials.username if credentials else None,
password=credentials.password if credentials else None,
token=credentials.token if credentials else None,
password_env=credentials.password_env if credentials else None,
token_env=credentials.token_env if credentials else None,
secret_ref=credentials.secret_ref if credentials else None,
clear_password=payload.clear_password,
clear_token=payload.clear_token,
capabilities=payload.capabilities,
policy=payload.policy,
metadata=payload.metadata,
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_PROFILES_COLLECTION,
resource_type=FILES_CONNECTOR_PROFILE_RESOURCE,
resource_id=row.id,
operation="updated",
principal=principal,
tenant_id=row.tenant_id,
payload={
"scope_type": row.scope_type,
"scope_id": row.scope_id,
"provider": row.provider,
},
)
session.commit()
session.refresh(row)
return FileConnectorProfileResponse(
**connector_profile_from_row(
row, credential_row=credential_row
).to_response()
)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.delete(
"/connectors/profiles/{profile_id}", response_model=FileConnectorProfileResponse
)
def deactivate_connector_profile(
profile_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
try:
row = get_connector_profile_row(
session,
tenant_id=principal.tenant_id,
profile_id=profile_id,
include_disabled=True,
)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
_require_connector_profile_write(principal, row.scope_type)
try:
changed = delete_connector_profile_row(
session,
row,
deletion_reason="api_delete",
user_id=principal.user.id,
api_key_id=principal.api_key.id if principal.api_key else None,
)
if changed:
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_PROFILES_COLLECTION,
resource_type=FILES_CONNECTOR_PROFILE_RESOURCE,
resource_id=row.id,
operation="deleted",
principal=principal,
tenant_id=row.tenant_id,
payload={
"scope_type": row.scope_type,
"scope_id": row.scope_id,
"provider": row.provider,
},
)
session.commit()
session.refresh(row)
return FileConnectorProfileResponse(
**connector_profile_from_row(row).to_response()
)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc) from exc
except Exception:
session.rollback()
raise

View File

@@ -0,0 +1,852 @@
from __future__ import annotations
import json
from typing import Literal
from fastapi import APIRouter, Depends, Query, status
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_any_scope
from govoplan_files.backend.change_tracking import (
FILES_CONNECTOR_CREDENTIALS_COLLECTION,
FILES_CONNECTOR_POLICIES_COLLECTION,
FILES_CONNECTOR_PROFILES_COLLECTION,
)
from govoplan_files.backend.schemas import (
FileConnectorCredentialCreateRequest,
FileConnectorCredentialResponse,
FileConnectorCredentialsResponse,
FileConnectorCredentialUpdateRequest,
FileConnectorDiscoveryRequest,
FileConnectorDiscoveryResponse,
FileConnectorSettingsDeltaResponse,
FileConnectorPolicyEvaluateRequest,
FileConnectorPolicyEvaluateResponse,
FileConnectorPolicyResponse,
FileConnectorPolicyUpdateRequest,
FileConnectorProfileCreateRequest,
FileConnectorProfileResponse,
FileConnectorProfilesResponse,
FileConnectorProviderResponse,
FileConnectorProvidersResponse,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.connector_credential_store import (
create_connector_credential_row,
get_connector_credential_row,
list_database_connector_credentials,
update_connector_credential_row,
)
from govoplan_files.backend.storage.connector_credential_deletion import (
delete_connector_credential_row,
)
from govoplan_files.backend.storage.connector_browse import (
ConnectorBrowseError,
ConnectorBrowseUnsupported,
browse_connector_profile,
)
from govoplan_files.backend.storage.connector_deployment import (
connector_effective_endpoint_url,
reject_api_controlled_deployment_references,
)
from govoplan_files.backend.storage.connector_profile_store import (
connector_profile_from_row,
create_connector_profile_row,
)
from govoplan_files.backend.storage.connector_providers import (
connector_provider_descriptors,
)
from govoplan_files.backend.storage.connector_policy import (
ConnectorAccessRequest,
ConnectorPolicyDenied,
connector_policy_decision,
connector_policy_sources_from_payload,
)
from govoplan_files.backend.storage.connector_policy_store import (
connector_policy_response,
set_connector_policy,
)
from govoplan_files.backend.route_support import (
FILES_CONNECTOR_CREDENTIAL_RESOURCE,
FILES_CONNECTOR_POLICY_RESOURCE,
FILES_CONNECTOR_PROFILE_RESOURCE,
_audit_connector_discovery_attempt,
_can_read_disabled_connector_profiles,
_connector_credential_response,
_connector_policy_error,
_credential_row_for_profile,
_discovery_profile_from_payload,
_ensure_campaign_file_access,
_ensure_connector_configuration_allowed,
_ensure_connector_credential_configuration_allowed,
_ensure_connector_local_policy_allowed,
_file_connector_policy_resource_id,
_file_connector_settings_entries,
_http_error,
_record_connector_settings_change,
_require_connector_credential_write,
_require_connector_policy_read,
_require_connector_profile_write,
_same_endpoint,
_visible_connector_profiles,
_webdav_discovery_candidates,
)
from govoplan_files.backend.services.connector_settings_delta import (
_full_file_connector_settings_delta_response,
_incremental_file_connector_settings_delta_response,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.post(
"/connector-policy/evaluate", response_model=FileConnectorPolicyEvaluateResponse
)
def evaluate_connector_policy(
payload: FileConnectorPolicyEvaluateRequest,
principal: ApiPrincipal = Depends(
require_any_scope("files:file:read", "files:file:upload", "files:file:download")
),
):
del principal
sources = connector_policy_sources_from_payload(
[item.model_dump(mode="json") for item in payload.policy_sources]
)
request = ConnectorAccessRequest.from_provenance(
payload.source_provenance.model_dump(mode="json", exclude_none=True),
operation=payload.operation,
)
return FileConnectorPolicyEvaluateResponse(
decision=connector_policy_decision(request, sources).to_dict()
)
@router.get(
"/connectors/settings/delta", response_model=FileConnectorSettingsDeltaResponse
)
def connector_settings_delta(
scope_type: str = Query(default="tenant"),
scope_id: str | None = Query(default=None),
provider: str | None = None,
campaign_id: str | None = None,
include_disabled: bool = False,
include_inactive: bool = False,
owner_type: Literal["user", "group"] | None = None,
owner_id: str | None = None,
since: str | None = None,
limit: int = Query(default=100, ge=1, le=500),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:read", "admin:settings:read"
)
),
):
scope_type = scope_type.strip().casefold()
_require_connector_policy_read(principal, scope_type)
try:
if since is None:
return _full_file_connector_settings_delta_response(
session,
principal,
scope_type=scope_type,
scope_id=scope_id,
provider=provider,
campaign_id=campaign_id,
include_disabled=include_disabled,
include_inactive=include_inactive,
owner_type=owner_type,
owner_id=owner_id,
)
entries, has_more = _file_connector_settings_entries(
session, tenant_id=principal.tenant_id, since=since, limit=limit
)
if entries is None:
return _full_file_connector_settings_delta_response(
session,
principal,
scope_type=scope_type,
scope_id=scope_id,
provider=provider,
campaign_id=campaign_id,
include_disabled=include_disabled,
include_inactive=include_inactive,
owner_type=owner_type,
owner_id=owner_id,
)
return _incremental_file_connector_settings_delta_response(
session,
principal,
entries=entries,
has_more=has_more,
scope_type=scope_type,
scope_id=scope_id,
provider=provider,
campaign_id=campaign_id,
include_disabled=include_disabled,
include_inactive=include_inactive,
owner_type=owner_type,
owner_id=owner_id,
)
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
raise _http_error(exc) from exc
@router.get("/connectors/providers", response_model=FileConnectorProvidersResponse)
def list_connector_providers(
principal: ApiPrincipal = Depends(
require_any_scope("files:file:read", "files:file:upload", "files:file:admin")
),
):
del principal
return FileConnectorProvidersResponse(
providers=[
FileConnectorProviderResponse(**item.to_response())
for item in connector_provider_descriptors()
]
)
@router.post("/connectors/discover", response_model=FileConnectorDiscoveryResponse)
def discover_connector_endpoint(
payload: FileConnectorDiscoveryRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
try:
reject_api_controlled_deployment_references(
password_env=payload.credentials.password_env,
token_env=payload.credentials.token_env,
secret_ref=payload.credentials.secret_ref,
metadata=payload.metadata,
)
except ValueError as exc:
raise _http_error(exc) from exc
if payload.provider not in {"webdav", "nextcloud"}:
return FileConnectorDiscoveryResponse(
provider=payload.provider,
endpoint_url=None,
base_path=payload.base_path,
status="unsupported",
message=f"Discovery is not implemented for {payload.provider} connectors yet",
)
candidates: list[dict[str, str]] = []
for endpoint_url in _webdav_discovery_candidates(payload):
profile = _discovery_profile_from_payload(
payload, endpoint_url, principal=principal
)
try:
_ensure_connector_configuration_allowed(
session,
principal,
connector_id=None,
credential_id=None,
provider=profile.provider,
endpoint_url=endpoint_url,
base_path=profile.base_path,
scope_type="tenant",
scope_id=principal.tenant_id,
operation="discover",
)
_audit_connector_discovery_attempt(
session,
principal,
provider=profile.provider,
endpoint_url=endpoint_url,
base_path=profile.base_path,
)
browse_connector_profile(profile, path=payload.base_path or "")
except ConnectorPolicyDenied as exc:
raise _connector_policy_error(exc) from exc
except (
ConnectorBrowseError,
ConnectorBrowseUnsupported,
OSError,
ValueError,
json.JSONDecodeError,
) as exc:
message = str(exc)
if "credentials were rejected" in message.casefold():
if payload.require_valid_credentials:
candidates.append(
{
"endpoint_url": endpoint_url,
"status": "credentials_rejected",
"message": "The endpoint exists, but the credentials were rejected.",
}
)
return FileConnectorDiscoveryResponse(
provider=payload.provider,
endpoint_url=endpoint_url,
base_path=payload.base_path,
status="credentials_rejected",
message="The endpoint was found, but login failed with these credentials.",
candidates=candidates,
metadata={"discovered_by": "webdav-auth-challenge"},
)
candidates.append(
{
"endpoint_url": endpoint_url,
"status": "found",
"message": "The endpoint exists, but credentials are required or were rejected.",
}
)
return FileConnectorDiscoveryResponse(
provider=payload.provider,
endpoint_url=endpoint_url,
base_path=payload.base_path,
status="found",
message="The endpoint was found. Add working credentials before saving or testing the connection.",
candidates=candidates,
metadata={"discovered_by": "webdav-auth-challenge"},
)
candidates.append(
{"endpoint_url": endpoint_url, "status": "failed", "message": message}
)
continue
status_value = (
"usable" if _same_endpoint(endpoint_url, payload.endpoint_url) else "found"
)
message = (
"The supplied URL is directly usable."
if status_value == "usable"
else "A usable connector endpoint was discovered."
)
candidates.append(
{"endpoint_url": endpoint_url, "status": status_value, "message": message}
)
return FileConnectorDiscoveryResponse(
provider=payload.provider,
endpoint_url=endpoint_url,
base_path=payload.base_path,
status=status_value,
message=message,
candidates=candidates,
metadata={"discovered_by": "webdav-propfind"},
)
return FileConnectorDiscoveryResponse(
provider=payload.provider,
endpoint_url=None,
base_path=payload.base_path,
status="not_found",
message="No usable WebDAV endpoint was found for this server URL.",
candidates=candidates,
)
@router.get("/connectors/credentials", response_model=FileConnectorCredentialsResponse)
def list_connector_credentials(
provider: str | None = None,
include_disabled: bool = False,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:read", "admin:settings:read"
)
),
):
provider_norm = provider.strip().casefold() if provider else None
credentials = list_database_connector_credentials(
session,
tenant_id=principal.tenant_id,
include_disabled=include_disabled
and _can_read_disabled_connector_profiles(principal),
)
return FileConnectorCredentialsResponse(
credentials=[
FileConnectorCredentialResponse(**credential.to_response())
for credential in credentials
if provider_norm is None or credential.provider in {None, provider_norm}
]
)
@router.get(
"/connectors/policies/{scope_type}", response_model=FileConnectorPolicyResponse
)
def read_connector_policy(
scope_type: str,
scope_id: str | None = Query(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:read", "admin:settings:read"
)
),
):
_require_connector_policy_read(principal, scope_type)
try:
return FileConnectorPolicyResponse(
**connector_policy_response(
session,
tenant_id=principal.tenant_id,
scope_type=scope_type,
scope_id=scope_id,
)
)
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
raise _http_error(exc, not_found=True) from exc
@router.put(
"/connectors/policies/{scope_type}", response_model=FileConnectorPolicyResponse
)
def write_connector_policy(
scope_type: str,
payload: FileConnectorPolicyUpdateRequest,
scope_id: str | None = Query(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
_require_connector_profile_write(principal, scope_type)
try:
set_connector_policy(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
scope_type=scope_type,
scope_id=scope_id,
policy=payload.policy,
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_POLICIES_COLLECTION,
resource_type=FILES_CONNECTOR_POLICY_RESOURCE,
resource_id=_file_connector_policy_resource_id(scope_type, scope_id),
operation="updated",
principal=principal,
tenant_id=None
if scope_type.strip().casefold() == "system"
else principal.tenant_id,
payload={"scope_type": scope_type.strip().casefold(), "scope_id": scope_id},
)
session.commit()
return FileConnectorPolicyResponse(
**connector_policy_response(
session,
tenant_id=principal.tenant_id,
scope_type=scope_type,
scope_id=scope_id,
)
)
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.post(
"/connectors/credentials",
response_model=FileConnectorCredentialResponse,
status_code=status.HTTP_201_CREATED,
)
def create_connector_credential(
payload: FileConnectorCredentialCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
_require_connector_credential_write(principal, payload.scope_type)
credentials = payload.credentials
try:
_ensure_connector_credential_configuration_allowed(
session,
principal,
credential_id=payload.id,
provider=payload.provider,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
operation="configure_credentials",
)
_ensure_connector_local_policy_allowed(
session,
principal,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
policy=payload.policy,
)
row = create_connector_credential_row(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
credential_id=payload.id,
label=payload.label,
provider=payload.provider,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
enabled=payload.enabled,
credential_mode=payload.credential_mode,
username=credentials.username,
password=credentials.password,
token=credentials.token,
password_env=credentials.password_env,
token_env=credentials.token_env,
secret_ref=credentials.secret_ref,
policy=payload.policy,
metadata=payload.metadata,
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_CREDENTIALS_COLLECTION,
resource_type=FILES_CONNECTOR_CREDENTIAL_RESOURCE,
resource_id=row.id,
operation="created",
principal=principal,
tenant_id=row.tenant_id,
payload={
"scope_type": row.scope_type,
"scope_id": row.scope_id,
"provider": row.provider,
},
)
session.commit()
session.refresh(row)
return _connector_credential_response(row)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.get(
"/connectors/credentials/{credential_id}",
response_model=FileConnectorCredentialResponse,
)
def get_connector_credential(
credential_id: str,
include_disabled: bool = False,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:read", "admin:settings:read"
)
),
):
try:
row = get_connector_credential_row(
session,
tenant_id=principal.tenant_id,
credential_id=credential_id,
include_disabled=include_disabled
and _can_read_disabled_connector_profiles(principal),
)
return _connector_credential_response(row)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
@router.patch(
"/connectors/credentials/{credential_id}",
response_model=FileConnectorCredentialResponse,
)
def update_connector_credential(
credential_id: str,
payload: FileConnectorCredentialUpdateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
try:
row = get_connector_credential_row(
session,
tenant_id=principal.tenant_id,
credential_id=credential_id,
include_disabled=True,
)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
_require_connector_credential_write(principal, row.scope_type)
credentials = payload.credentials
try:
provider = payload.provider if payload.provider is not None else row.provider
_ensure_connector_credential_configuration_allowed(
session,
principal,
credential_id=row.id,
provider=provider,
scope_type=row.scope_type,
scope_id=row.scope_id,
operation="configure_credentials",
)
if payload.policy is not None:
_ensure_connector_local_policy_allowed(
session,
principal,
scope_type=row.scope_type,
scope_id=row.scope_id,
policy=payload.policy,
)
update_connector_credential_row(
session,
row,
user_id=principal.user.id,
label=payload.label,
provider=payload.provider,
enabled=payload.enabled,
credential_mode=payload.credential_mode,
username=credentials.username if credentials else None,
password=credentials.password if credentials else None,
token=credentials.token if credentials else None,
password_env=credentials.password_env if credentials else None,
token_env=credentials.token_env if credentials else None,
secret_ref=credentials.secret_ref if credentials else None,
clear_password=payload.clear_password,
clear_token=payload.clear_token,
policy=payload.policy,
metadata=payload.metadata,
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_CREDENTIALS_COLLECTION,
resource_type=FILES_CONNECTOR_CREDENTIAL_RESOURCE,
resource_id=row.id,
operation="updated",
principal=principal,
tenant_id=row.tenant_id,
payload={
"scope_type": row.scope_type,
"scope_id": row.scope_id,
"provider": row.provider,
},
)
session.commit()
session.refresh(row)
return _connector_credential_response(row)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.delete(
"/connectors/credentials/{credential_id}",
response_model=FileConnectorCredentialResponse,
)
def deactivate_connector_credential(
credential_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
try:
row = get_connector_credential_row(
session,
tenant_id=principal.tenant_id,
credential_id=credential_id,
include_disabled=True,
)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
_require_connector_credential_write(principal, row.scope_type)
try:
deletion = delete_connector_credential_row(
session,
row,
deletion_reason="api_delete",
user_id=principal.user.id,
api_key_id=principal.api_key.id if principal.api_key else None,
)
if deletion.changed:
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_CREDENTIALS_COLLECTION,
resource_type=FILES_CONNECTOR_CREDENTIAL_RESOURCE,
resource_id=row.id,
operation="deleted",
principal=principal,
tenant_id=row.tenant_id,
payload={
"scope_type": row.scope_type,
"scope_id": row.scope_id,
"provider": row.provider,
},
)
for profile in deletion.affected_profiles:
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_PROFILES_COLLECTION,
resource_type=FILES_CONNECTOR_PROFILE_RESOURCE,
resource_id=profile.id,
operation="updated",
principal=principal,
tenant_id=profile.tenant_id,
payload={
"scope_type": profile.scope_type,
"scope_id": profile.scope_id,
"provider": profile.provider,
"reason": "credential_deleted",
},
)
session.commit()
session.refresh(row)
return _connector_credential_response(row)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc) from exc
except Exception:
session.rollback()
raise
@router.get("/connectors/profiles", response_model=FileConnectorProfilesResponse)
def list_connector_profiles(
provider: str | None = None,
campaign_id: str | None = None,
include_disabled: bool = False,
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:
if campaign_id:
_ensure_campaign_file_access(session, principal, campaign_id)
profiles = _visible_connector_profiles(
session,
principal,
provider=provider,
campaign_id=campaign_id,
include_disabled=include_disabled
and _can_read_disabled_connector_profiles(principal),
include_admin_scopes=_can_read_disabled_connector_profiles(principal),
include_effective_policy=False,
)
except (OSError, ValueError, json.JSONDecodeError) as exc:
raise _http_error(exc) from exc
return FileConnectorProfilesResponse(
profiles=[
FileConnectorProfileResponse(**profile.to_response())
for profile in profiles
]
)
@router.post(
"/connectors/profiles",
response_model=FileConnectorProfileResponse,
status_code=status.HTTP_201_CREATED,
)
def create_connector_profile(
payload: FileConnectorProfileCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope(
"files:file:admin", "system:settings:write", "admin:settings:write"
)
),
):
_require_connector_profile_write(principal, payload.scope_type)
credentials = payload.credentials
try:
credential_row = _credential_row_for_profile(
session,
principal,
credential_profile_id=payload.credential_profile_id,
provider=payload.provider,
profile_id=payload.id,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
)
_ensure_connector_configuration_allowed(
session,
principal,
connector_id=payload.id,
credential_id=payload.credential_profile_id,
provider=payload.provider,
endpoint_url=connector_effective_endpoint_url(
provider=payload.provider,
endpoint_url=payload.endpoint_url,
metadata=payload.metadata,
),
base_path=payload.base_path,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
operation="configure",
)
_ensure_connector_local_policy_allowed(
session,
principal,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
policy=payload.policy,
)
row = create_connector_profile_row(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
profile_id=payload.id,
label=payload.label,
provider=payload.provider,
scope_type=payload.scope_type,
scope_id=payload.scope_id,
endpoint_url=payload.endpoint_url,
base_path=payload.base_path,
enabled=payload.enabled,
credential_profile_id=payload.credential_profile_id,
credential_mode=payload.credential_mode,
username=credentials.username,
password=credentials.password,
token=credentials.token,
password_env=credentials.password_env,
token_env=credentials.token_env,
secret_ref=credentials.secret_ref,
capabilities=payload.capabilities,
policy=payload.policy,
metadata=payload.metadata,
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_PROFILES_COLLECTION,
resource_type=FILES_CONNECTOR_PROFILE_RESOURCE,
resource_id=row.id,
operation="created",
principal=principal,
tenant_id=row.tenant_id,
payload={
"scope_type": row.scope_type,
"scope_id": row.scope_id,
"provider": row.provider,
},
)
session.commit()
session.refresh(row)
return FileConnectorProfileResponse(
**connector_profile_from_row(
row, credential_row=credential_row
).to_response()
)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc

View File

@@ -0,0 +1,151 @@
from __future__ import annotations
from typing import Literal
from fastapi import APIRouter, Depends, Query
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.schemas import (
FileFolderCreateRequest,
FileFolderDeleteRequest,
FileFolderDeleteResponse,
FileFolderResponse,
FileFoldersResponse,
)
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.folders import (
create_folder,
list_folders_for_user,
list_folders_for_user_window,
soft_delete_folder,
)
from govoplan_files.backend.route_support import (
_folder_response,
_http_error,
_is_admin,
)
from govoplan_files.backend.services.list_queries import (
FOLDERS_LIST_CURSOR_SCOPE,
_cursor_page_size,
_files_delta_watermark,
_folder_cursor_values,
_folders_list_fingerprint,
_next_folder_list_cursor,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.get("/folders", response_model=FileFoldersResponse)
def list_file_folders(
owner_type: Literal["user", "group"],
owner_id: str,
page_size: int | None = Query(default=None, ge=1, le=1000),
cursor: str | None = None,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
try:
watermark = _files_delta_watermark(session, principal.tenant_id)
effective_page_size = _cursor_page_size(
FOLDERS_LIST_CURSOR_SCOPE, cursor, page_size
)
if effective_page_size is None:
folders = list_folders_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
is_admin=_is_admin(principal),
)
return FileFoldersResponse(
folders=[_folder_response(folder) for folder in folders],
watermark=watermark,
)
fingerprint = _folders_list_fingerprint(
principal,
owner_type=owner_type,
owner_id=owner_id,
page_size=effective_page_size,
)
after_path, after_id = _folder_cursor_values(cursor, fingerprint=fingerprint)
folders, has_more = list_folders_for_user_window(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
is_admin=_is_admin(principal),
page_size=effective_page_size,
after_path=after_path,
after_id=after_id,
)
return FileFoldersResponse(
folders=[_folder_response(folder) for folder in folders],
cursor=cursor,
next_cursor=_next_folder_list_cursor(
principal,
folders,
owner_type=owner_type,
owner_id=owner_id,
page_size=effective_page_size,
has_more=has_more,
),
watermark=watermark,
)
except FileStorageError as exc:
raise _http_error(exc) from exc
@router.post("/folders", response_model=FileFolderResponse)
def create_file_folder(
payload: FileFolderCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:organize")),
):
try:
folder = create_folder(
session,
tenant_id=principal.tenant_id,
owner_type=payload.owner_type,
owner_id=payload.owner_id,
user_id=principal.user.id,
path=payload.path,
is_admin=_is_admin(principal),
)
session.commit()
return _folder_response(folder)
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.post("/folders/delete", response_model=FileFolderDeleteResponse)
def delete_file_folder(
payload: FileFolderDeleteRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:delete")),
):
try:
deleted_folders, deleted_files = soft_delete_folder(
session,
tenant_id=principal.tenant_id,
owner_type=payload.owner_type,
owner_id=payload.owner_id,
user_id=principal.user.id,
path=payload.path,
recursive=payload.recursive,
is_admin=_is_admin(principal),
)
session.commit()
return FileFolderDeleteResponse(
deleted_folders=deleted_folders, deleted_files=deleted_files
)
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc

View File

@@ -0,0 +1,142 @@
from __future__ import annotations
from typing import Literal
from fastapi import APIRouter, Depends, Query
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.schemas import (
FileDeltaResponse,
FileListResponse,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.files import (
list_assets_for_user,
list_assets_for_user_window,
)
from govoplan_files.backend.route_support import (
_asset_list_response,
_ensure_campaign_file_access,
_ensure_list_owner_access,
_is_admin,
)
from govoplan_files.backend.services.list_queries import (
FILES_LIST_CURSOR_SCOPE,
_cursor_page_size,
_file_cursor_values,
_files_delta_response,
_files_delta_watermark,
_files_list_fingerprint,
_full_file_delta_response,
_next_file_list_cursor,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.get("/delta", response_model=FileDeltaResponse)
def files_delta(
owner_type: Literal["user", "group"] | None = None,
owner_id: str | None = None,
campaign_id: str | None = None,
path_prefix: str | None = None,
since: str | None = None,
limit: int = Query(default=500, ge=1, le=1000),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
_ensure_list_owner_access(session, principal, owner_type, owner_id)
_ensure_campaign_file_access(session, principal, campaign_id)
if since is None:
return _full_file_delta_response(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
)
return _files_delta_response(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
since=since,
limit=limit,
)
@router.get("", response_model=FileListResponse)
def list_files(
owner_type: Literal["user", "group"] | None = None,
owner_id: str | None = None,
campaign_id: str | None = None,
path_prefix: str | None = None,
page_size: int | None = Query(default=None, ge=1, le=1000),
cursor: str | None = None,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
_ensure_list_owner_access(session, principal, owner_type, owner_id)
_ensure_campaign_file_access(session, principal, campaign_id)
watermark = _files_delta_watermark(session, principal.tenant_id)
effective_page_size = _cursor_page_size(FILES_LIST_CURSOR_SCOPE, cursor, page_size)
if effective_page_size is not None:
fingerprint = _files_list_fingerprint(
principal,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
page_size=effective_page_size,
)
after_display_path, after_updated_at, after_id = _file_cursor_values(
cursor, fingerprint=fingerprint
)
assets, has_more = list_assets_for_user_window(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
is_admin=_is_admin(principal),
page_size=effective_page_size,
after_display_path=after_display_path,
after_updated_at=after_updated_at,
after_id=after_id,
)
return FileListResponse(
files=_asset_list_response(session, assets, include_shares=True),
cursor=cursor,
next_cursor=_next_file_list_cursor(
principal,
assets,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
page_size=effective_page_size,
has_more=has_more,
),
watermark=watermark,
)
assets = list_assets_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
is_admin=_is_admin(principal),
)
return FileListResponse(
files=_asset_list_response(session, assets, include_shares=True),
watermark=watermark,
)

View File

@@ -0,0 +1,116 @@
from __future__ import annotations
from fastapi import APIRouter, Depends
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.schemas import (
BulkFileShareRequest,
BulkFileShareResponse,
FileShareRequest,
FileShareResponse,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.files import (
get_asset_for_user,
share_file,
share_files,
)
from govoplan_files.backend.route_support import (
_http_error,
_is_admin,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.post("/{file_id}/shares", response_model=FileShareResponse)
def create_share(
file_id: str,
payload: FileShareRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:share")),
):
try:
asset = get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
require_write=True,
is_admin=_is_admin(principal),
)
share = share_file(
session,
tenant_id=principal.tenant_id,
asset=asset,
target_type=payload.target_type,
target_id=payload.target_id,
permission=payload.permission,
user_id=principal.user.id,
)
session.commit()
return FileShareResponse(
id=share.id,
target_type=share.target_type,
target_id=share.target_id,
permission=share.permission,
created_at=share.created_at.isoformat(),
revoked_at=share.revoked_at.isoformat() if share.revoked_at else None,
)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc) from exc
@router.post("/bulk-shares", response_model=BulkFileShareResponse)
def create_bulk_shares(
payload: BulkFileShareRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:share")),
):
try:
file_ids = list(dict.fromkeys(payload.file_ids))
assets = [
get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
require_write=True,
is_admin=_is_admin(principal),
)
for file_id in file_ids
]
shares = share_files(
session,
tenant_id=principal.tenant_id,
assets=assets,
target_type=payload.target_type,
target_id=payload.target_id,
permission=payload.permission,
user_id=principal.user.id,
)
session.commit()
return BulkFileShareResponse(
shared_count=len(shares),
shares=[
FileShareResponse(
id=share.id,
target_type=share.target_type,
target_id=share.target_id,
permission=share.permission,
created_at=share.created_at.isoformat(),
revoked_at=share.revoked_at.isoformat()
if share.revoked_at
else None,
)
for share in shares
],
)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc) from exc

View File

@@ -0,0 +1,292 @@
from __future__ import annotations
import json
from typing import Literal
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.change_tracking import (
FILES_CONNECTOR_SPACES_COLLECTION,
)
from govoplan_files.backend.schemas import (
FileConnectorSpaceCreateRequest,
FileConnectorSpaceResponse,
FileConnectorSpacesResponse,
FileConnectorSpaceUpdateRequest,
FileSpaceResponse,
FileSpacesResponse,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.access import group_refs_for_ids, user_group_ids
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.connector_spaces import (
connector_space_owner_id,
create_connector_space,
get_connector_space_for_user,
list_connector_spaces_for_user,
soft_delete_connector_space,
update_connector_space,
)
from govoplan_files.backend.storage.connector_policy import (
ConnectorPolicyDenied,
)
from govoplan_files.backend.route_support import (
FILES_CONNECTOR_SPACE_RESOURCE,
_connector_policy_error,
_connector_space_file_space_response,
_connector_space_policy_decision,
_connector_space_response,
_http_error,
_is_admin,
_record_connector_settings_change,
_visible_connector_profile,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.get("/spaces", response_model=FileSpacesResponse)
def list_file_spaces(
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
spaces = [
FileSpaceResponse(
id=f"user:{principal.user.id}",
label="My files",
owner_type="user",
owner_id=principal.user.id,
description="Files owned by your user account.",
)
]
group_ids = user_group_ids(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
include_admin_groups=_is_admin(principal),
)
if group_ids:
groups = group_refs_for_ids(tenant_id=principal.tenant_id, group_ids=group_ids)
spaces.extend(
FileSpaceResponse(
id=f"group:{group.id}",
label=f"{group.name} files",
owner_type="group",
owner_id=group.id,
description="Files owned by this group.",
)
for group in groups
)
connector_spaces = list_connector_spaces_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
is_admin=_is_admin(principal),
)
spaces.extend(
_connector_space_file_space_response(space) for space in connector_spaces
)
return FileSpacesResponse(spaces=spaces)
@router.get("/connector-spaces", response_model=FileConnectorSpacesResponse)
def list_file_connector_spaces(
owner_type: Literal["user", "group"] | None = None,
owner_id: str | None = None,
include_inactive: bool = False,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
try:
spaces = list_connector_spaces_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
include_inactive=include_inactive and _is_admin(principal),
is_admin=_is_admin(principal),
)
return FileConnectorSpacesResponse(
spaces=[_connector_space_response(space) for space in spaces]
)
except FileStorageError as exc:
raise _http_error(exc) from exc
@router.post(
"/connector-spaces",
response_model=FileConnectorSpaceResponse,
status_code=status.HTTP_201_CREATED,
)
def create_file_connector_space(
payload: FileConnectorSpaceCreateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:organize")),
):
if payload.owner_type == "group" and not payload.owner_id:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="owner_id is required for group connector spaces",
)
target_owner = payload.owner_id or principal.user.id
try:
profile = _visible_connector_profile(
session, principal, payload.connector_profile_id
)
decision = _connector_space_policy_decision(
profile,
library_id=payload.library_id,
remote_path=payload.remote_path,
operation="link",
)
if not decision.allowed:
raise ConnectorPolicyDenied(decision)
space = create_connector_space(
session,
tenant_id=principal.tenant_id,
owner_type=payload.owner_type,
owner_id=target_owner,
user_id=principal.user.id,
label=payload.label,
profile=profile,
library_id=payload.library_id,
remote_path=payload.remote_path,
sync_mode=payload.sync_mode,
metadata=payload.metadata,
is_admin=_is_admin(principal),
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_SPACES_COLLECTION,
resource_type=FILES_CONNECTOR_SPACE_RESOURCE,
resource_id=space.id,
operation="created",
principal=principal,
tenant_id=space.tenant_id,
payload={
"owner_type": space.owner_type,
"owner_id": connector_space_owner_id(space),
},
)
session.commit()
return _connector_space_response(space)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.patch("/connector-spaces/{space_id}", response_model=FileConnectorSpaceResponse)
def update_file_connector_space(
space_id: str,
payload: FileConnectorSpaceUpdateRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:organize")),
):
try:
space = get_connector_space_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
space_id=space_id,
include_inactive=_is_admin(principal),
is_admin=_is_admin(principal),
)
except FileStorageError as exc:
raise _http_error(exc, not_found=True) from exc
try:
if payload.library_id is not None or payload.remote_path is not None:
profile = _visible_connector_profile(
session, principal, space.connector_profile_id
)
decision = _connector_space_policy_decision(
profile,
library_id=payload.library_id
if payload.library_id is not None
else space.library_id,
remote_path=payload.remote_path
if payload.remote_path is not None
else space.remote_path,
operation="link",
)
if not decision.allowed:
raise ConnectorPolicyDenied(decision)
update_connector_space(
session,
space,
user_id=principal.user.id,
label=payload.label,
library_id=payload.library_id,
remote_path=payload.remote_path,
sync_mode=payload.sync_mode,
is_active=payload.is_active,
metadata=payload.metadata,
is_admin=_is_admin(principal),
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_SPACES_COLLECTION,
resource_type=FILES_CONNECTOR_SPACE_RESOURCE,
resource_id=space.id,
operation="updated",
principal=principal,
tenant_id=space.tenant_id,
payload={
"owner_type": space.owner_type,
"owner_id": connector_space_owner_id(space),
},
)
session.commit()
return _connector_space_response(space)
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, ValueError, json.JSONDecodeError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.delete(
"/connector-spaces/{space_id}", response_model=FileConnectorSpaceResponse
)
def delete_file_connector_space(
space_id: str,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:organize")),
):
try:
space = get_connector_space_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
space_id=space_id,
include_inactive=True,
is_admin=_is_admin(principal),
)
soft_delete_connector_space(
session, space, user_id=principal.user.id, is_admin=_is_admin(principal)
)
_record_connector_settings_change(
session,
collection=FILES_CONNECTOR_SPACES_COLLECTION,
resource_type=FILES_CONNECTOR_SPACE_RESOURCE,
resource_id=space.id,
operation="deleted",
principal=principal,
tenant_id=space.tenant_id,
payload={
"owner_type": space.owner_type,
"owner_id": connector_space_owner_id(space),
},
)
session.commit()
return _connector_space_response(space)
except FileStorageError as exc:
session.rollback()
raise _http_error(exc, not_found=True) from exc

View File

@@ -0,0 +1,213 @@
from __future__ import annotations
import tempfile
from fastapi import APIRouter, Depends
from fastapi.responses import FileResponse
from starlette.background import BackgroundTask
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.schemas import (
ArchiveRequest,
PatternMatchResponse,
PatternResolveRequest,
PatternResolveResponse,
RenamePreviewItem,
RenameRequest,
RenameResponse,
TransferRequest,
TransferResponse,
_conflict_resolutions,
)
from govoplan_core.db.session import get_session
from govoplan_files.backend.storage.paths import (
UnsafeFilePathError,
filename_from_path,
normalize_logical_path,
)
from govoplan_files.backend.storage.archives import create_zip_file
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.files import (
get_asset_for_user,
list_assets_for_user,
)
from govoplan_files.backend.storage.search import resolve_patterns
from govoplan_files.backend.storage.transfers import (
rename_selection,
transfer_selection,
)
from govoplan_files.backend.route_support import (
_asset_response,
_attachment_disposition,
_audit_connector_access,
_cleanup_temp_file,
_ensure_campaign_file_access,
_ensure_list_owner_access,
_http_error,
_is_admin,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.post("/bulk-rename", response_model=RenameResponse)
def bulk_rename(
payload: RenameRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:organize")),
):
try:
plan = rename_selection(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
file_ids=payload.file_ids,
folder_paths=payload.folder_paths,
owner_type=payload.owner_type,
owner_id=payload.owner_id,
mode=payload.mode,
new_name=payload.new_name,
find=payload.find,
replacement=payload.replacement,
prefix=payload.prefix,
suffix=payload.suffix,
recursive=payload.recursive,
dry_run=payload.dry_run,
is_admin=_is_admin(principal),
)
if not payload.dry_run:
session.commit()
return RenameResponse(
dry_run=payload.dry_run,
items=[
RenamePreviewItem(
kind=item.kind,
id=item.id,
file_id=item.id if item.kind == "file" else None,
folder_path=item.old_path if item.kind == "folder" else None,
old_path=item.old_path,
new_path=item.new_path,
)
for item in plan
],
)
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.post("/transfer", response_model=TransferResponse)
def transfer_files(
payload: TransferRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:organize")),
):
try:
files, folders = transfer_selection(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
operation=payload.operation,
file_ids=payload.file_ids,
folder_paths=payload.folder_paths,
source_owner_type=payload.source_owner_type,
source_owner_id=payload.source_owner_id,
target_owner_type=payload.target_owner_type,
target_owner_id=payload.target_owner_id,
target_folder=payload.target_folder,
conflict_strategy=payload.conflict_strategy,
conflict_resolutions=_conflict_resolutions(payload.conflict_resolutions),
is_admin=_is_admin(principal),
)
session.commit()
return TransferResponse(
operation=payload.operation, files=files, folders=folders
)
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc
@router.post("/archive.zip")
def download_archive(
payload: ArchiveRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:download")),
):
try:
assets = [
get_asset_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
asset_id=file_id,
is_admin=_is_admin(principal),
)
for file_id in payload.file_ids
]
tmp = tempfile.NamedTemporaryFile(
prefix="govoplan-files-", suffix=".zip", delete=False
)
tmp_path = tmp.name
tmp.close()
try:
create_zip_file(session, assets, tmp_path)
except Exception:
_cleanup_temp_file(tmp_path)
raise
_audit_connector_access(session, principal, assets, operation="archive")
except FileStorageError as exc:
raise _http_error(exc) from exc
filename = filename_from_path(
normalize_logical_path(payload.filename, fallback_filename="files.zip")
)
headers = {"Content-Disposition": _attachment_disposition(filename)}
return FileResponse(
tmp_path,
media_type="application/zip",
headers=headers,
background=BackgroundTask(_cleanup_temp_file, tmp_path),
)
@router.post("/resolve-patterns", response_model=PatternResolveResponse)
def resolve_file_patterns(
payload: PatternResolveRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:read")),
):
_ensure_list_owner_access(session, principal, payload.owner_type, payload.owner_id)
_ensure_campaign_file_access(session, principal, payload.campaign_id)
try:
assets = list_assets_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=payload.owner_type,
owner_id=payload.owner_id,
campaign_id=payload.campaign_id,
path_prefix=payload.path_prefix,
is_admin=_is_admin(principal),
)
resolved, unmatched = resolve_patterns(
assets,
payload.patterns,
base_path=payload.path_prefix,
case_sensitive=payload.case_sensitive,
)
return PatternResolveResponse(
patterns=[
PatternMatchResponse(
pattern=item.pattern,
matches=[_asset_response(session, asset) for asset in item.matches],
)
for item in resolved
],
unmatched=[_asset_response(session, asset) for asset in unmatched]
if payload.include_unmatched
else [],
)
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
raise _http_error(exc) from exc

View File

@@ -0,0 +1,207 @@
from __future__ import annotations
import json
from typing import Literal
from fastapi import APIRouter, Depends, File as FastAPIFile, Form, UploadFile
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, require_scope
from govoplan_files.backend.schemas import (
ConflictResolutionRequest,
FileUploadResponse,
_conflict_resolutions,
)
from govoplan_files.backend.db.models import FileAsset
from govoplan_core.db.session import get_session
from govoplan_files.backend.runtime import settings
from govoplan_files.backend.storage.paths import UnsafeFilePathError
from govoplan_files.backend.storage.archives import extract_zip_upload
from govoplan_files.backend.storage.common import FileStorageError
from govoplan_files.backend.storage.connector_policy import (
ConnectorPolicyDenied,
)
from govoplan_files.backend.storage.files import (
create_file_asset,
)
from govoplan_files.backend.route_support import (
_asset_response,
_audit_connector_imports,
_cleanup_temp_file,
_connector_policy_error,
_enforce_connector_policy,
_http_error,
_is_admin,
_read_limited_upload,
_source_metadata_from_form,
_spool_limited_upload_to_temp,
)
router = APIRouter(prefix="/files", tags=["files"])
@router.post("/upload", response_model=FileUploadResponse)
def upload_files(
files: list[UploadFile] = FastAPIFile(...),
owner_type: Literal["user", "group"] = Form(default="user"),
owner_id: str | None = Form(default=None),
path: str = Form(default=""),
campaign_id: str | None = Form(default=None),
unpack_zip: bool = Form(default=False),
conflict_strategy: Literal["reject", "overwrite", "rename"] = Form(
default="reject"
),
conflict_resolutions_json: str | None = Form(default=None),
source_provenance_json: str | None = Form(default=None),
source_revision: str | None = Form(default=None),
connector_policy_json: str | None = Form(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:upload")),
):
target_owner = owner_id or principal.user.id
uploaded_assets: list[FileAsset] = []
try:
raw_resolutions = (
json.loads(conflict_resolutions_json) if conflict_resolutions_json else []
)
upload_resolutions = _conflict_resolutions(
[ConflictResolutionRequest(**item) for item in raw_resolutions]
)
_enforce_connector_policy(
source_provenance_json, connector_policy_json, operation="import"
)
metadata = _source_metadata_from_form(source_provenance_json, source_revision)
for upload in files:
filename = upload.filename or "file"
content_type = upload.content_type or None
upload_limit = (
settings.file_upload_zip_max_bytes
if unpack_zip and filename.lower().endswith(".zip")
else settings.file_upload_max_bytes
)
if unpack_zip and filename.lower().endswith(".zip"):
zip_path = _spool_limited_upload_to_temp(
upload, max_bytes=upload_limit, suffix=".zip"
)
try:
extracted = extract_zip_upload(
session,
tenant_id=principal.tenant_id,
owner_type=owner_type,
owner_id=target_owner,
user_id=principal.user.id,
zip_data=zip_path,
folder=path,
campaign_id=campaign_id,
conflict_strategy=conflict_strategy,
conflict_resolutions=upload_resolutions,
metadata=metadata,
is_admin=_is_admin(principal),
max_file_bytes=settings.file_upload_max_bytes,
max_total_bytes=settings.file_upload_zip_max_bytes,
)
finally:
_cleanup_temp_file(zip_path)
uploaded_assets.extend(item.asset for item in extracted)
continue
data = _read_limited_upload(upload, max_bytes=upload_limit)
stored = create_file_asset(
session,
tenant_id=principal.tenant_id,
owner_type=owner_type,
owner_id=target_owner,
user_id=principal.user.id,
filename=filename,
data=data,
folder=path,
content_type=content_type,
campaign_id=campaign_id,
conflict_strategy=conflict_strategy,
conflict_resolutions=upload_resolutions,
metadata=metadata,
is_admin=_is_admin(principal),
)
uploaded_assets.append(stored.asset)
_audit_connector_imports(session, principal, uploaded_assets)
session.commit()
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc
return FileUploadResponse(
files=[
_asset_response(session, asset, include_shares=True)
for asset in uploaded_assets
]
)
@router.post("/upload-zip", response_model=FileUploadResponse)
def upload_zip(
file: UploadFile = FastAPIFile(...),
owner_type: Literal["user", "group"] = Form(default="user"),
owner_id: str | None = Form(default=None),
path: str = Form(default=""),
campaign_id: str | None = Form(default=None),
conflict_strategy: Literal["reject", "overwrite", "rename"] = Form(
default="reject"
),
conflict_resolutions_json: str | None = Form(default=None),
source_provenance_json: str | None = Form(default=None),
source_revision: str | None = Form(default=None),
connector_policy_json: str | None = Form(default=None),
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(require_scope("files:file:upload")),
):
target_owner = owner_id or principal.user.id
zip_path: str | None = None
try:
raw_resolutions = (
json.loads(conflict_resolutions_json) if conflict_resolutions_json else []
)
upload_resolutions = _conflict_resolutions(
[ConflictResolutionRequest(**item) for item in raw_resolutions]
)
_enforce_connector_policy(
source_provenance_json, connector_policy_json, operation="import"
)
metadata = _source_metadata_from_form(source_provenance_json, source_revision)
zip_path = _spool_limited_upload_to_temp(
file, max_bytes=settings.file_upload_zip_max_bytes, suffix=".zip"
)
extracted = extract_zip_upload(
session,
tenant_id=principal.tenant_id,
owner_type=owner_type,
owner_id=target_owner,
user_id=principal.user.id,
zip_data=zip_path,
folder=path,
campaign_id=campaign_id,
conflict_strategy=conflict_strategy,
conflict_resolutions=upload_resolutions,
metadata=metadata,
is_admin=_is_admin(principal),
max_file_bytes=settings.file_upload_max_bytes,
max_total_bytes=settings.file_upload_zip_max_bytes,
)
_audit_connector_imports(session, principal, [item.asset for item in extracted])
session.commit()
except ConnectorPolicyDenied as exc:
session.rollback()
raise _connector_policy_error(exc) from exc
except (FileStorageError, UnsafeFilePathError, ValueError) as exc:
session.rollback()
raise _http_error(exc) from exc
finally:
if zip_path:
_cleanup_temp_file(zip_path)
return FileUploadResponse(
files=[
_asset_response(session, item.asset, include_shares=True)
for item in extracted
]
)

View File

@@ -0,0 +1 @@
"""Query and response assembly services used by Files routes."""

View File

@@ -0,0 +1,292 @@
from __future__ import annotations
from typing import Literal
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal
from govoplan_core.api.v1.schemas import DeltaDeletedItem
from govoplan_core.core.change_sequence import (
ChangeSequenceEntry,
)
from govoplan_files.backend.change_tracking import (
FILES_CONNECTOR_CREDENTIALS_COLLECTION,
FILES_CONNECTOR_POLICIES_COLLECTION,
FILES_CONNECTOR_PROFILES_COLLECTION,
FILES_CONNECTOR_SPACES_COLLECTION,
)
from govoplan_files.backend.schemas import (
FileConnectorCredentialResponse,
FileConnectorSettingsDeltaResponse,
FileConnectorPolicyResponse,
FileConnectorProfileResponse,
)
from govoplan_files.backend.db.models import FileConnectorSpace
from govoplan_files.backend.storage.connector_credential_store import (
ConnectorCredential,
)
from govoplan_files.backend.storage.connector_profiles import ConnectorProfile
from govoplan_files.backend.storage.connector_policy_store import (
connector_policy_response,
)
from govoplan_files.backend.route_support import (
FILES_CONNECTOR_CREDENTIAL_RESOURCE,
FILES_CONNECTOR_PROFILE_RESOURCE,
FILES_CONNECTOR_SPACE_RESOURCE,
_can_read_disabled_connector_profiles,
_connector_deleted_item,
_connector_space_response,
_ensure_campaign_file_access,
_file_connector_settings_response_watermark,
_file_connector_settings_watermark,
_visible_connector_credentials,
_visible_connector_profiles,
_visible_connector_spaces,
)
def _full_file_connector_settings_delta_response(
session: Session,
principal: ApiPrincipal,
*,
scope_type: str,
scope_id: str | None,
provider: str | None,
campaign_id: str | None,
include_disabled: bool,
include_inactive: bool,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
) -> FileConnectorSettingsDeltaResponse:
if campaign_id:
_ensure_campaign_file_access(session, principal, campaign_id)
profiles = _visible_connector_profiles(
session,
principal,
provider=provider,
campaign_id=campaign_id,
include_disabled=include_disabled
and _can_read_disabled_connector_profiles(principal),
include_admin_scopes=_can_read_disabled_connector_profiles(principal),
include_effective_policy=False,
)
credentials = _visible_connector_credentials(
session,
principal,
provider=provider,
include_disabled=include_disabled,
)
spaces = _visible_connector_spaces(
session,
principal,
owner_type=owner_type,
owner_id=owner_id,
include_inactive=include_inactive,
)
return FileConnectorSettingsDeltaResponse(
profiles=[
FileConnectorProfileResponse(**profile.to_response())
for profile in profiles
],
credentials=[
FileConnectorCredentialResponse(**credential.to_response())
for credential in credentials
],
spaces=[_connector_space_response(space) for space in spaces],
policy=FileConnectorPolicyResponse(
**connector_policy_response(
session,
tenant_id=principal.tenant_id,
scope_type=scope_type,
scope_id=scope_id,
)
),
changed_sections=["profiles", "credentials", "spaces", "policy"],
deleted=[],
watermark=_file_connector_settings_watermark(
session, tenant_id=principal.tenant_id
),
has_more=False,
full=True,
)
def _changed_file_connector_setting_ids(
entries: list[ChangeSequenceEntry],
) -> tuple[set[str], set[str], set[str], bool]:
changed_profile_ids = {
entry.resource_id
for entry in entries
if entry.collection == FILES_CONNECTOR_PROFILES_COLLECTION
and entry.resource_type == FILES_CONNECTOR_PROFILE_RESOURCE
}
changed_credential_ids = {
entry.resource_id
for entry in entries
if entry.collection == FILES_CONNECTOR_CREDENTIALS_COLLECTION
and entry.resource_type == FILES_CONNECTOR_CREDENTIAL_RESOURCE
}
changed_space_ids = {
entry.resource_id
for entry in entries
if entry.collection == FILES_CONNECTOR_SPACES_COLLECTION
and entry.resource_type == FILES_CONNECTOR_SPACE_RESOURCE
}
policy_changed = any(
entry.collection == FILES_CONNECTOR_POLICIES_COLLECTION for entry in entries
)
return (
changed_profile_ids,
changed_credential_ids,
changed_space_ids,
policy_changed,
)
def _file_connector_settings_changed_sections(
*,
changed_profile_ids: set[str],
changed_credential_ids: set[str],
changed_space_ids: set[str],
policy_changed: bool,
profiles: list[ConnectorProfile],
) -> list[str]:
changed_sections = []
if changed_profile_ids or any(
profile.credential_profile_id
and profile.credential_profile_id in changed_credential_ids
for profile in profiles
):
changed_sections.append("profiles")
if changed_credential_ids:
changed_sections.append("credentials")
if changed_space_ids:
changed_sections.append("spaces")
if policy_changed:
changed_sections.append("policy")
return changed_sections
def _file_connector_settings_deleted_items(
entries: list[ChangeSequenceEntry],
*,
visible_profiles: dict[str, ConnectorProfile],
visible_credentials: dict[str, ConnectorCredential],
visible_spaces: dict[str, FileConnectorSpace],
) -> list[DeltaDeletedItem]:
return [
_connector_deleted_item(entry)
for entry in entries
if (
entry.resource_type == FILES_CONNECTOR_PROFILE_RESOURCE
and entry.resource_id not in visible_profiles
)
or (
entry.resource_type == FILES_CONNECTOR_CREDENTIAL_RESOURCE
and entry.resource_id not in visible_credentials
)
or (
entry.resource_type == FILES_CONNECTOR_SPACE_RESOURCE
and entry.resource_id not in visible_spaces
)
]
def _incremental_file_connector_settings_delta_response(
session: Session,
principal: ApiPrincipal,
*,
entries: list[ChangeSequenceEntry],
has_more: bool,
scope_type: str,
scope_id: str | None,
provider: str | None,
campaign_id: str | None,
include_disabled: bool,
include_inactive: bool,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
) -> FileConnectorSettingsDeltaResponse:
changed_profile_ids, changed_credential_ids, changed_space_ids, policy_changed = (
_changed_file_connector_setting_ids(entries)
)
profiles = _visible_connector_profiles(
session,
principal,
provider=provider,
campaign_id=campaign_id,
include_disabled=include_disabled
and _can_read_disabled_connector_profiles(principal),
include_admin_scopes=_can_read_disabled_connector_profiles(principal),
include_effective_policy=False,
)
visible_profiles = {
profile.id: profile
for profile in profiles
if profile.id in changed_profile_ids
or (
profile.credential_profile_id
and profile.credential_profile_id in changed_credential_ids
)
}
credentials = _visible_connector_credentials(
session,
principal,
provider=provider,
include_disabled=include_disabled,
)
visible_credentials = {
credential.id: credential
for credential in credentials
if credential.id in changed_credential_ids
}
spaces = _visible_connector_spaces(
session,
principal,
owner_type=owner_type,
owner_id=owner_id,
include_inactive=include_inactive,
)
visible_spaces = {
space.id: space for space in spaces if space.id in changed_space_ids
}
return FileConnectorSettingsDeltaResponse(
profiles=[
FileConnectorProfileResponse(**profile.to_response())
for profile in visible_profiles.values()
],
credentials=[
FileConnectorCredentialResponse(**credential.to_response())
for credential in visible_credentials.values()
],
spaces=[_connector_space_response(space) for space in visible_spaces.values()],
policy=FileConnectorPolicyResponse(
**connector_policy_response(
session,
tenant_id=principal.tenant_id,
scope_type=scope_type,
scope_id=scope_id,
)
)
if policy_changed
else None,
changed_sections=_file_connector_settings_changed_sections(
changed_profile_ids=changed_profile_ids,
changed_credential_ids=changed_credential_ids,
changed_space_ids=changed_space_ids,
policy_changed=policy_changed,
profiles=profiles,
),
deleted=_file_connector_settings_deleted_items(
entries,
visible_profiles=visible_profiles,
visible_credentials=visible_credentials,
visible_spaces=visible_spaces,
),
watermark=_file_connector_settings_response_watermark(
session, tenant_id=principal.tenant_id, entries=entries, has_more=has_more
),
has_more=has_more,
full=False,
)

View File

@@ -0,0 +1,611 @@
from __future__ import annotations
from datetime import datetime
from typing import Literal
from fastapi import HTTPException, status
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal
from govoplan_core.api.v1.schemas import DeltaDeletedItem
from govoplan_core.core.change_sequence import (
ChangeSequenceEntry,
decode_sequence_watermark,
encode_sequence_watermark,
max_sequence_id,
sequence_entries_since,
sequence_watermark_is_expired,
)
from govoplan_core.core.pagination import (
KeysetCursorError,
decode_keyset_cursor,
encode_keyset_cursor,
keyset_query_fingerprint,
)
from govoplan_files.backend.change_tracking import (
FILES_ASSETS_COLLECTION,
FILES_FOLDERS_COLLECTION,
FILES_MODULE_ID,
)
from govoplan_files.backend.schemas import (
FileDeltaResponse,
)
from govoplan_files.backend.db.models import FileAsset, FileFolder, FileShare
from govoplan_files.backend.storage.access import user_group_ids
from govoplan_files.backend.storage.files import (
list_assets_for_user,
)
from govoplan_files.backend.storage.folders import list_folders_for_user
from govoplan_files.backend.route_support import (
_asset_list_response,
_folder_response,
_is_admin,
)
_FILES_DELTA_COLLECTIONS = (FILES_ASSETS_COLLECTION, FILES_FOLDERS_COLLECTION)
FILES_LIST_CURSOR_SCOPE = "files.list.v1"
FOLDERS_LIST_CURSOR_SCOPE = "files.folders.list.v1"
DEFAULT_FILE_LIST_PAGE_SIZE = 500
def _cursor_http_error(exc: Exception) -> HTTPException:
return HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc))
def _cursor_page_size(
scope: str, cursor: str | None, explicit_page_size: int | None
) -> int | None:
if explicit_page_size is not None:
return explicit_page_size
if not cursor:
return DEFAULT_FILE_LIST_PAGE_SIZE
try:
values = decode_keyset_cursor(scope, cursor)
except KeysetCursorError as exc:
raise _cursor_http_error(exc) from exc
raw_page_size = values.get("page_size")
if not isinstance(raw_page_size, int) or raw_page_size < 1 or raw_page_size > 1000:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor"
)
return raw_page_size
def _files_list_fingerprint(
principal: ApiPrincipal,
*,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
page_size: int,
) -> str:
return keyset_query_fingerprint(
FILES_LIST_CURSOR_SCOPE,
{
"tenant_id": principal.tenant_id,
"actor": "admin" if _is_admin(principal) else principal.user.id,
"owner_type": owner_type,
"owner_id": owner_id,
"campaign_id": campaign_id,
"path_prefix": path_prefix or "",
"sort": "display_path.asc,updated_at.desc,id.asc",
"page_size": page_size,
},
)
def _folders_list_fingerprint(
principal: ApiPrincipal,
*,
owner_type: Literal["user", "group"],
owner_id: str,
page_size: int,
) -> str:
return keyset_query_fingerprint(
FOLDERS_LIST_CURSOR_SCOPE,
{
"tenant_id": principal.tenant_id,
"actor": "admin" if _is_admin(principal) else principal.user.id,
"owner_type": owner_type,
"owner_id": owner_id,
"sort": "path.asc,id.asc",
"page_size": page_size,
},
)
def _file_cursor_values(
cursor: str | None, *, fingerprint: str
) -> tuple[str | None, datetime | None, str | None]:
try:
values = decode_keyset_cursor(
FILES_LIST_CURSOR_SCOPE, cursor, fingerprint=fingerprint
)
except KeysetCursorError as exc:
raise _cursor_http_error(exc) from exc
if values is None:
return None, None, None
display_path = values.get("display_path")
updated_at = values.get("updated_at")
asset_id = values.get("id")
if (
not isinstance(display_path, str)
or not isinstance(updated_at, str)
or not isinstance(asset_id, str)
):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor"
)
try:
parsed_updated_at = datetime.fromisoformat(updated_at)
except ValueError as exc:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor"
) from exc
return display_path, parsed_updated_at, asset_id
def _folder_cursor_values(
cursor: str | None, *, fingerprint: str
) -> tuple[str | None, str | None]:
try:
values = decode_keyset_cursor(
FOLDERS_LIST_CURSOR_SCOPE, cursor, fingerprint=fingerprint
)
except KeysetCursorError as exc:
raise _cursor_http_error(exc) from exc
if values is None:
return None, None
folder_path = values.get("path")
folder_id = values.get("id")
if not isinstance(folder_path, str) or not isinstance(folder_id, str):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor"
)
return folder_path, folder_id
def _next_file_list_cursor(
principal: ApiPrincipal,
assets: list[FileAsset],
*,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
page_size: int,
has_more: bool,
) -> str | None:
if not has_more or not assets:
return None
last = assets[-1]
return encode_keyset_cursor(
FILES_LIST_CURSOR_SCOPE,
fingerprint=_files_list_fingerprint(
principal,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
page_size=page_size,
),
values={
"display_path": last.display_path,
"updated_at": last.updated_at,
"id": last.id,
"page_size": page_size,
},
)
def _next_folder_list_cursor(
principal: ApiPrincipal,
folders: list[FileFolder],
*,
owner_type: Literal["user", "group"],
owner_id: str,
page_size: int,
has_more: bool,
) -> str | None:
if not has_more or not folders:
return None
last = folders[-1]
return encode_keyset_cursor(
FOLDERS_LIST_CURSOR_SCOPE,
fingerprint=_folders_list_fingerprint(
principal, owner_type=owner_type, owner_id=owner_id, page_size=page_size
),
values={"path": last.path, "id": last.id, "page_size": page_size},
)
def _files_delta_watermark(session: Session, tenant_id: str) -> str:
return encode_sequence_watermark(
max_sequence_id(
session,
tenant_id=tenant_id,
module_id=FILES_MODULE_ID,
collections=_FILES_DELTA_COLLECTIONS,
)
)
def _full_file_delta_response(
session: Session,
*,
principal: ApiPrincipal,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
) -> FileDeltaResponse:
assets = list_assets_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
is_admin=_is_admin(principal),
)
folders = _visible_folders_for_delta(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
path_prefix=path_prefix,
)
return FileDeltaResponse(
files=_asset_list_response(session, assets, include_shares=True),
folders=[_folder_response(folder) for folder in folders],
deleted=[],
watermark=_files_delta_watermark(session, principal.tenant_id),
has_more=False,
full=True,
)
def _visible_folders_for_delta(
session: Session,
*,
principal: ApiPrincipal,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
path_prefix: str | None,
) -> list[FileFolder]:
if not owner_type or not owner_id:
return []
folders = list_folders_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
is_admin=_is_admin(principal),
)
if not path_prefix:
return folders
normalized = path_prefix.strip().strip("/")
if not normalized:
return folders
return [
folder
for folder in folders
if folder.path == normalized or folder.path.startswith(f"{normalized}/")
]
def _entry_path_matches(
entry_payload: dict[str, object], path_prefix: str | None
) -> bool:
if not path_prefix:
return True
normalized = path_prefix.strip().strip("/")
if not normalized:
return True
for key in ("path", "previous_path"):
value = entry_payload.get(key)
if isinstance(value, str) and (
value == normalized or value.startswith(f"{normalized}/")
):
return True
return False
def _entry_owner_matches(
entry_payload: dict[str, object],
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
) -> bool:
if not owner_type or not owner_id:
return True
return (
entry_payload.get("owner_type") == owner_type
and entry_payload.get("owner_id") == owner_id
)
def _entry_campaign_matches(session: Session, entry, campaign_id: str | None) -> bool:
if not campaign_id:
return True
payload = entry.payload or {}
if (
payload.get("share_target_type") == "campaign"
and payload.get("share_target_id") == campaign_id
):
return True
if entry.resource_type != "file":
return False
return (
session.query(FileShare)
.filter(
FileShare.tenant_id == entry.tenant_id,
FileShare.file_asset_id == entry.resource_id,
FileShare.target_type == "campaign",
FileShare.target_id == campaign_id,
FileShare.revoked_at.is_(None),
)
.first()
is not None
)
def _principal_group_ids_for_delta(
session: Session, principal: ApiPrincipal
) -> set[str]:
cache = session.info.setdefault("files_delta_group_ids", {})
key = (principal.tenant_id, principal.user.id)
if key not in cache:
cache[key] = set(
user_group_ids(
session, tenant_id=principal.tenant_id, user_id=principal.user.id
)
)
return cache[key]
def _entry_subject_matches_principal(
session: Session, principal: ApiPrincipal, entry_payload: dict[str, object]
) -> bool:
if _is_admin(principal):
return True
owner_type = entry_payload.get("owner_type")
owner_id = entry_payload.get("owner_id")
if owner_type == "user" and owner_id == principal.user.id:
return True
if owner_type == "group" and isinstance(owner_id, str):
if owner_id in _principal_group_ids_for_delta(session, principal):
return True
share_target_type = entry_payload.get("share_target_type")
share_target_id = entry_payload.get("share_target_id")
if share_target_type == "user" and share_target_id == principal.user.id:
return True
if share_target_type == "tenant" and share_target_id == principal.tenant_id:
return True
if share_target_type == "group" and isinstance(share_target_id, str):
if share_target_id in _principal_group_ids_for_delta(session, principal):
return True
return False
def _entry_matches_delta_scope(
session: Session,
principal: ApiPrincipal,
entry,
*,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
) -> bool:
payload = entry.payload or {}
if (
not owner_type
and not campaign_id
and not _entry_subject_matches_principal(session, principal, payload)
):
return False
return (
_entry_owner_matches(payload, owner_type, owner_id)
and _entry_path_matches(payload, path_prefix)
and _entry_campaign_matches(session, entry, campaign_id)
)
def _decode_files_delta_watermark(since: str) -> int:
try:
return decode_sequence_watermark(since)
except ValueError as exc:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)
) from exc
def _changed_delta_resource_ids(
entries: list[ChangeSequenceEntry],
) -> tuple[list[str], list[str]]:
file_ids = list(
dict.fromkeys(
entry.resource_id for entry in entries if entry.resource_type == "file"
)
)
folder_ids = list(
dict.fromkeys(
entry.resource_id for entry in entries if entry.resource_type == "folder"
)
)
return file_ids, folder_ids
def _visible_assets_for_delta(
session: Session,
*,
principal: ApiPrincipal,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
changed_file_ids: list[str],
) -> dict[str, FileAsset]:
return {
asset.id: asset
for asset in list_assets_for_user(
session,
tenant_id=principal.tenant_id,
user_id=principal.user.id,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
is_admin=_is_admin(principal),
)
if asset.id in changed_file_ids
}
def _changed_visible_folders_for_delta(
session: Session,
*,
principal: ApiPrincipal,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
path_prefix: str | None,
changed_folder_ids: list[str],
) -> dict[str, FileFolder]:
return {
folder.id: folder
for folder in _visible_folders_for_delta(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
path_prefix=path_prefix,
)
if folder.id in changed_folder_ids
}
def _deleted_delta_items(
session: Session,
*,
principal: ApiPrincipal,
entries: list[ChangeSequenceEntry],
visible_assets: dict[str, FileAsset],
visible_folders: dict[str, FileFolder],
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
) -> list[DeltaDeletedItem]:
deleted: dict[tuple[str, str], DeltaDeletedItem] = {}
for entry in entries:
if entry.resource_type == "file" and entry.resource_id in visible_assets:
continue
if entry.resource_type == "folder" and entry.resource_id in visible_folders:
continue
if not _entry_matches_delta_scope(
session,
principal,
entry,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
):
continue
deleted[(entry.resource_type, entry.resource_id)] = DeltaDeletedItem(
id=entry.resource_id,
resource_type=entry.resource_type,
revision=encode_sequence_watermark(entry.id),
deleted_at=entry.created_at if entry.operation == "deleted" else None,
)
return list(deleted.values())
def _files_delta_response(
session: Session,
*,
principal: ApiPrincipal,
owner_type: Literal["user", "group"] | None,
owner_id: str | None,
campaign_id: str | None,
path_prefix: str | None,
since: str,
limit: int,
) -> FileDeltaResponse:
since_sequence = _decode_files_delta_watermark(since)
if sequence_watermark_is_expired(
session,
since=since_sequence,
tenant_id=principal.tenant_id,
module_id=FILES_MODULE_ID,
collections=_FILES_DELTA_COLLECTIONS,
):
return _full_file_delta_response(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
)
entries_plus_one = sequence_entries_since(
session,
since=since_sequence,
tenant_id=principal.tenant_id,
module_id=FILES_MODULE_ID,
collections=_FILES_DELTA_COLLECTIONS,
limit=limit + 1,
)
has_more = len(entries_plus_one) > limit
entries = entries_plus_one[:limit]
changed_file_ids, changed_folder_ids = _changed_delta_resource_ids(entries)
visible_assets = _visible_assets_for_delta(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
changed_file_ids=changed_file_ids,
)
visible_folders = _changed_visible_folders_for_delta(
session,
principal=principal,
owner_type=owner_type,
owner_id=owner_id,
path_prefix=path_prefix,
changed_folder_ids=changed_folder_ids,
)
deleted = _deleted_delta_items(
session,
principal=principal,
entries=entries,
visible_assets=visible_assets,
visible_folders=visible_folders,
owner_type=owner_type,
owner_id=owner_id,
campaign_id=campaign_id,
path_prefix=path_prefix,
)
watermark = (
encode_sequence_watermark(entries[-1].id)
if has_more and entries
else _files_delta_watermark(session, principal.tenant_id)
)
return FileDeltaResponse(
files=_asset_list_response(
session, list(visible_assets.values()), include_shares=True
),
folders=[_folder_response(folder) for folder in visible_folders.values()],
deleted=deleted,
watermark=watermark,
has_more=has_more,
full=False,
)