2879 lines
131 KiB
Python
2879 lines
131 KiB
Python
"""P8 inventory, backfill, shadow validation, and per-Workspace atomic cutover.
|
|
|
|
Inputs are strict operator-produced snapshots and remediation evidence. This module
|
|
never discovers or matches legacy product rows: opaque KnowledgeFS Space IDs are
|
|
registered directly into the independent control plane. Unknown access, unresolved
|
|
subjects, and unresolved cutover quarantine remain fail-closed and auditable.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
from collections.abc import Callable, Iterable
|
|
from datetime import UTC, datetime, timedelta
|
|
from typing import Literal, NamedTuple, cast
|
|
from uuid import NAMESPACE_URL, UUID, uuid5
|
|
|
|
import sqlalchemy as sa
|
|
from pydantic import BaseModel, ConfigDict, Field, JsonValue, ValidationError, field_validator, model_validator
|
|
from sqlalchemy.exc import IntegrityError
|
|
from sqlalchemy.orm import Session, sessionmaker
|
|
|
|
from libs.datetime_utils import naive_utc_now
|
|
from models.knowledge_fs import (
|
|
KnowledgeFSApiCredential,
|
|
KnowledgeFSApiCredentialStatus,
|
|
KnowledgeFSAuthorizationRevision,
|
|
KnowledgeFSControlSpace,
|
|
KnowledgeFSControlSpacePermission,
|
|
KnowledgeFSControlSpacePermissionRole,
|
|
KnowledgeFSControlSpaceState,
|
|
KnowledgeFSControlSpaceVisibility,
|
|
KnowledgeFSExternalAccessPolicy,
|
|
)
|
|
from models.knowledge_fs_cutover import (
|
|
KnowledgeFSCutoverRevisionWatermark,
|
|
KnowledgeFSCutoverSmokeResults,
|
|
KnowledgeFSMigrationIssue,
|
|
KnowledgeFSMigrationIssueKind,
|
|
KnowledgeFSMigrationIssueStatus,
|
|
KnowledgeFSMigrationQuarantine,
|
|
KnowledgeFSMigrationQuarantineDisposition,
|
|
KnowledgeFSMigrationQuarantineKind,
|
|
KnowledgeFSShadowAuthorizationDecision,
|
|
KnowledgeFSShadowAuthorizationDiff,
|
|
KnowledgeFSShadowAuthorizationObservation,
|
|
KnowledgeFSWorkspaceCutoverLedger,
|
|
KnowledgeFSWorkspaceCutoverPhase,
|
|
knowledge_fs_cutover_smoke_results_passed,
|
|
)
|
|
from repositories.knowledge_fs_cutover_repository import (
|
|
KnowledgeFSCutoverCASUpdate,
|
|
KnowledgeFSQuarantineCASUpdate,
|
|
KnowledgeFSShadowDiffCASUpdate,
|
|
)
|
|
from repositories.sqlalchemy_knowledge_fs_cutover_repository import SQLAlchemyKnowledgeFSCutoverRepository
|
|
from services.knowledge_fs.lifecycle_port import (
|
|
KnowledgeFSDifyIntegrationActivationAck,
|
|
KnowledgeFSDifyIntegrationActivationRequest,
|
|
KnowledgeFSDifyIntegrationFreezeAck,
|
|
KnowledgeFSDifyIntegrationFreezeRequest,
|
|
KnowledgeFSLifecycleRemoteError,
|
|
KnowledgeFSLifecycleRemotePort,
|
|
)
|
|
|
|
|
|
class StrictInput(BaseModel):
|
|
model_config = ConfigDict(extra="forbid", frozen=True)
|
|
|
|
|
|
class CutoverRevisionWatermarkInput(StrictInput):
|
|
membership_epoch: int = Field(ge=0)
|
|
space_acl_epoch: int = Field(ge=0)
|
|
external_access_epoch: int = Field(ge=0)
|
|
content_policy_revision: int = Field(ge=0)
|
|
|
|
def to_record(self) -> KnowledgeFSCutoverRevisionWatermark:
|
|
return cast(KnowledgeFSCutoverRevisionWatermark, self.model_dump())
|
|
|
|
|
|
class LegacyPermissionInventoryInput(StrictInput):
|
|
subject_id: str = Field(min_length=1, max_length=255)
|
|
account_id: UUID | None = None
|
|
role: Literal["owner", "editor", "viewer"]
|
|
|
|
|
|
class LegacyTaskInventoryInput(StrictInput):
|
|
task_id: str = Field(min_length=1, max_length=255)
|
|
subject_id: str | None = Field(default=None, min_length=1, max_length=255)
|
|
account_id: UUID | None = None
|
|
state: Literal["queued", "running", "paused", "completed", "failed", "canceled"]
|
|
expires_at: datetime | None = None
|
|
|
|
|
|
class LegacyApiKeyInventoryInput(StrictInput):
|
|
key_id: str = Field(min_length=1, max_length=255)
|
|
prefix: str = Field(min_length=1, max_length=32)
|
|
last4: str = Field(min_length=4, max_length=4)
|
|
|
|
|
|
class LegacySpaceInventoryInput(StrictInput):
|
|
knowledge_space_id: UUID
|
|
knowledge_space_revision: int = Field(ge=0)
|
|
provisioning_key: str = Field(min_length=1, max_length=255)
|
|
owner_subject_id: str = Field(min_length=1, max_length=255)
|
|
owner_account_id: UUID | None = None
|
|
visibility: Literal["only_me", "all_members", "partial_members", "unknown"]
|
|
external_access_enabled: bool | None = None
|
|
permissions: list[LegacyPermissionInventoryInput] = Field(default_factory=list)
|
|
legacy_api_keys: list[LegacyApiKeyInventoryInput] = Field(default_factory=list)
|
|
tasks: list[LegacyTaskInventoryInput] = Field(default_factory=list)
|
|
orphan_resource_ids: list[str] = Field(default_factory=list)
|
|
|
|
|
|
class WorkspaceInventoryInput(StrictInput):
|
|
tenant_id: UUID
|
|
source_revision_watermark: CutoverRevisionWatermarkInput
|
|
task_watermark: int = Field(ge=0)
|
|
spaces: list[LegacySpaceInventoryInput]
|
|
|
|
|
|
class ShadowAuthorizationObservationInput(StrictInput):
|
|
schema_version: Literal["knowledge-fs-p8-shadow-observation/v1"]
|
|
tenant_id: UUID
|
|
diff_key: str = Field(min_length=1, max_length=255)
|
|
control_space_id: UUID | None = None
|
|
principal: str = Field(min_length=1, max_length=255)
|
|
legacy_allowed: bool | None
|
|
dify_allowed: bool
|
|
reason: str = Field(min_length=1, max_length=4096)
|
|
observed_revision: CutoverRevisionWatermarkInput
|
|
observed_at: datetime
|
|
producer: str = Field(min_length=1, max_length=255)
|
|
|
|
@field_validator("observed_at")
|
|
@classmethod
|
|
def require_timezone(cls, value: datetime) -> datetime:
|
|
if value.tzinfo is None or value.utcoffset() is None:
|
|
raise ValueError("shadow observation timestamps require an explicit timezone")
|
|
return value
|
|
|
|
|
|
class ShadowCompletionInput(StrictInput):
|
|
schema_version: Literal["knowledge-fs-p8-shadow-completion/v1"]
|
|
tenant_id: UUID
|
|
expected_cas_version: int = Field(ge=0)
|
|
producer: str = Field(min_length=1, max_length=255)
|
|
completed_by_operator: str = Field(min_length=1, max_length=255)
|
|
completed_by_account_id: UUID
|
|
completed_at: datetime
|
|
traffic_zero: bool
|
|
window_started_at: datetime | None = None
|
|
window_ended_at: datetime | None = None
|
|
traffic_zero_evidence: dict[str, JsonValue] | None = Field(default=None, min_length=1, max_length=32)
|
|
|
|
@field_validator("completed_at", "window_started_at", "window_ended_at")
|
|
@classmethod
|
|
def require_timezone(cls, value: datetime | None) -> datetime | None:
|
|
if value is not None and (value.tzinfo is None or value.utcoffset() is None):
|
|
raise ValueError("shadow completion timestamps require an explicit timezone")
|
|
return value
|
|
|
|
@model_validator(mode="after")
|
|
def require_window_or_traffic_zero_evidence(self) -> ShadowCompletionInput:
|
|
if self.traffic_zero:
|
|
if self.traffic_zero_evidence is None:
|
|
raise ValueError("traffic-zero completion requires explicit evidence")
|
|
if self.window_started_at is not None or self.window_ended_at is not None:
|
|
raise ValueError("traffic-zero completion cannot also declare an observation window")
|
|
elif self.window_started_at is None or self.window_ended_at is None or self.traffic_zero_evidence is not None:
|
|
raise ValueError("nonzero shadow completion requires a window and no traffic-zero evidence")
|
|
return self
|
|
|
|
|
|
class FinalDeltaInput(StrictInput):
|
|
tenant_id: UUID
|
|
expected_cas_version: int = Field(ge=0)
|
|
final_revision_watermark: CutoverRevisionWatermarkInput
|
|
applied_revision_watermark: CutoverRevisionWatermarkInput
|
|
final_task_watermark: int = Field(ge=0)
|
|
applied_task_watermark: int = Field(ge=0)
|
|
|
|
|
|
class GeneralQuarantineResolutionEvidence(StrictInput):
|
|
schema_version: Literal["knowledge-fs-p8-general-resolution/v1"]
|
|
reference: str = Field(min_length=1, max_length=2048)
|
|
|
|
|
|
class LegacyCredentialRotationEvidence(StrictInput):
|
|
schema_version: Literal["knowledge-fs-p8-credential-rotation/v1"]
|
|
legacy_key_id: str = Field(min_length=1, max_length=255)
|
|
knowledge_space_id: UUID
|
|
control_space_id: UUID
|
|
dify_credential_id: UUID
|
|
dify_credential_revision: int = Field(ge=0)
|
|
legacy_revoked_at: datetime
|
|
verification_reference: str = Field(min_length=1, max_length=2048)
|
|
plaintext_migrated: Literal[False]
|
|
|
|
@field_validator("legacy_revoked_at")
|
|
@classmethod
|
|
def require_timezone(cls, value: datetime) -> datetime:
|
|
if value.tzinfo is None or value.utcoffset() is None:
|
|
raise ValueError("credential rotation timestamp requires an explicit timezone")
|
|
return value
|
|
|
|
|
|
class LegacyTaskResolutionEvidence(StrictInput):
|
|
schema_version: Literal["knowledge-fs-p8-task-resolution/v1"]
|
|
task_id: str = Field(min_length=1, max_length=255)
|
|
action: Literal["migrate", "wait", "cancel", "isolate"]
|
|
resulting_state: Literal["migrated", "completed", "failed", "canceled", "isolated"]
|
|
final_task_watermark: int = Field(ge=0)
|
|
verification_reference: str = Field(min_length=1, max_length=2048)
|
|
|
|
|
|
class QuarantineResolutionInput(StrictInput):
|
|
"""Operator-attested remediation evidence for one tenant-scoped quarantine row."""
|
|
|
|
schema_version: Literal["knowledge-fs-p8-quarantine-resolution/v1"]
|
|
tenant_id: UUID
|
|
source_kind: KnowledgeFSMigrationQuarantineKind
|
|
source_id: str = Field(min_length=1, max_length=255)
|
|
expected_row_version: int = Field(ge=0)
|
|
resolved_by_operator: str = Field(min_length=1, max_length=255)
|
|
resolved_by_account_id: UUID
|
|
evidence: GeneralQuarantineResolutionEvidence | LegacyCredentialRotationEvidence | LegacyTaskResolutionEvidence
|
|
resolved_at: datetime
|
|
|
|
@field_validator("resolved_at")
|
|
@classmethod
|
|
def require_timezone(cls, value: datetime) -> datetime:
|
|
if value.tzinfo is None or value.utcoffset() is None:
|
|
raise ValueError("quarantine resolution timestamps require an explicit timezone")
|
|
return value
|
|
|
|
@model_validator(mode="after")
|
|
def require_source_specific_evidence(self) -> QuarantineResolutionInput:
|
|
expected_type: type[StrictInput]
|
|
if self.source_kind is KnowledgeFSMigrationQuarantineKind.LEGACY_API_KEY:
|
|
expected_type = LegacyCredentialRotationEvidence
|
|
elif self.source_kind is KnowledgeFSMigrationQuarantineKind.TASK:
|
|
expected_type = LegacyTaskResolutionEvidence
|
|
else:
|
|
expected_type = GeneralQuarantineResolutionEvidence
|
|
if not isinstance(self.evidence, expected_type):
|
|
raise ValueError(f"{self.source_kind.value} requires its source-specific evidence schema")
|
|
return self
|
|
|
|
|
|
class CutoverSmokeChecksInput(StrictInput):
|
|
authorization: bool
|
|
list_spaces: bool
|
|
create_space: bool
|
|
query: bool
|
|
upload: bool
|
|
stream: bool
|
|
deletion: bool
|
|
|
|
|
|
class CutoverSmokeEvidenceReferencesInput(StrictInput):
|
|
authorization: str = Field(min_length=1, max_length=2048)
|
|
list_spaces: str = Field(min_length=1, max_length=2048)
|
|
create_space: str = Field(min_length=1, max_length=2048)
|
|
query: str = Field(min_length=1, max_length=2048)
|
|
upload: str = Field(min_length=1, max_length=2048)
|
|
stream: str = Field(min_length=1, max_length=2048)
|
|
deletion: str = Field(min_length=1, max_length=2048)
|
|
|
|
|
|
class CutoverSmokeResultsInput(StrictInput):
|
|
schema_version: Literal["knowledge-fs-p8-cutover-smoke/v1"]
|
|
tenant_id: UUID
|
|
environment: Literal["production"]
|
|
operator: str = Field(min_length=1, max_length=255)
|
|
operator_account_id: UUID
|
|
observed_at: datetime
|
|
checks: CutoverSmokeChecksInput
|
|
evidence_references: CutoverSmokeEvidenceReferencesInput
|
|
|
|
@field_validator("observed_at")
|
|
@classmethod
|
|
def require_timezone(cls, value: datetime) -> datetime:
|
|
if value.tzinfo is None or value.utcoffset() is None:
|
|
raise ValueError("smoke evidence timestamp requires an explicit timezone")
|
|
return value
|
|
|
|
def to_record(self) -> KnowledgeFSCutoverSmokeResults:
|
|
return cast(KnowledgeFSCutoverSmokeResults, self.model_dump(mode="json"))
|
|
|
|
@property
|
|
def passed(self) -> bool:
|
|
return all(self.checks.model_dump().values())
|
|
|
|
|
|
class LegacyDependencyInput(StrictInput):
|
|
dependency_key: str = Field(min_length=1, max_length=255)
|
|
kind: Literal["permission_snapshot", "foreign_key"]
|
|
table_name: str = Field(min_length=1, max_length=255)
|
|
column_name: str = Field(min_length=1, max_length=255)
|
|
constraint_name: str | None = Field(default=None, min_length=1, max_length=255)
|
|
active_rows: int = Field(ge=0)
|
|
migrated_rows: int = Field(ge=0)
|
|
|
|
|
|
class CutoverInventoryReport(NamedTuple):
|
|
tenant_id: str
|
|
spaces: int
|
|
permissions: int
|
|
unresolved_subjects: int
|
|
unknown_external_access: int
|
|
legacy_keys_requiring_rotation: int
|
|
tasks_by_disposition: dict[str, int]
|
|
orphan_resources: int
|
|
applied: bool
|
|
replayed: bool
|
|
|
|
|
|
class CutoverBackfillReport(NamedTuple):
|
|
tenant_id: str
|
|
control_spaces_registered: int
|
|
permissions_granted: int
|
|
quarantined: int
|
|
open_issues: int
|
|
phase: str
|
|
applied: bool
|
|
|
|
|
|
class ShadowAuthorizationReport(NamedTuple):
|
|
tenant_id: str
|
|
matches: int
|
|
tightened: int
|
|
expanded: int
|
|
unknown: int
|
|
open_diffs: int
|
|
applied: bool
|
|
recorded: int
|
|
replayed: int
|
|
|
|
|
|
class ShadowCompletionReport(NamedTuple):
|
|
tenant_id: str
|
|
observation_count: int
|
|
traffic_zero: bool
|
|
evidence_digest: str
|
|
latest_observed_revision: KnowledgeFSCutoverRevisionWatermark | None
|
|
applied: bool
|
|
replayed: bool
|
|
|
|
|
|
class LegacyDependencyDashboard(NamedTuple):
|
|
tenant_id: str
|
|
ready: bool
|
|
permission_snapshot_dependencies: int
|
|
foreign_key_dependencies: int
|
|
active_rows: int
|
|
dependencies: tuple[dict[str, object], ...]
|
|
applied: bool
|
|
|
|
|
|
class QuarantineResolutionReport(NamedTuple):
|
|
tenant_id: str
|
|
source_kind: str
|
|
source_id: str
|
|
disposition: str
|
|
row_version: int
|
|
applied: bool
|
|
replayed: bool
|
|
|
|
|
|
class KnowledgeFSCutoverError(RuntimeError):
|
|
"""Base error for an operator cutover request rejected before mutation."""
|
|
|
|
|
|
class KnowledgeFSCutoverNotFoundError(KnowledgeFSCutoverError):
|
|
pass
|
|
|
|
|
|
class KnowledgeFSCutoverConflictError(KnowledgeFSCutoverError):
|
|
pass
|
|
|
|
|
|
class KnowledgeFSCutoverGateBlockedError(KnowledgeFSCutoverError):
|
|
pass
|
|
|
|
|
|
_VISIBILITY_MAP = {
|
|
"only_me": KnowledgeFSControlSpaceVisibility.ONLY_ME,
|
|
"all_members": KnowledgeFSControlSpaceVisibility.ALL_TEAM_MEMBERS,
|
|
"partial_members": KnowledgeFSControlSpaceVisibility.PARTIAL_MEMBERS,
|
|
}
|
|
|
|
_CUTOVER_QUARANTINE_KINDS = (
|
|
KnowledgeFSMigrationQuarantineKind.CONTROL_SPACE,
|
|
KnowledgeFSMigrationQuarantineKind.TASK,
|
|
KnowledgeFSMigrationQuarantineKind.LEGACY_API_KEY,
|
|
KnowledgeFSMigrationQuarantineKind.ORPHAN_RESOURCE,
|
|
)
|
|
_DEFAULT_TRUSTED_SHADOW_PRODUCERS = frozenset({"dify-shadow-authorizer"})
|
|
_DEFAULT_TRUSTED_SHADOW_OPERATORS = frozenset({"knowledge-fs-cutover"})
|
|
_GREENFIELD_INITIALIZATION_MAX_ATTEMPTS = 32
|
|
_GREENFIELD_ROLLBACK_WINDOW = timedelta(days=30)
|
|
_ZERO_REVISION_WATERMARK: KnowledgeFSCutoverRevisionWatermark = {
|
|
"membership_epoch": 0,
|
|
"space_acl_epoch": 0,
|
|
"external_access_epoch": 0,
|
|
"content_policy_revision": 0,
|
|
}
|
|
|
|
|
|
class _ActivationAnchor(NamedTuple):
|
|
control_space_id: str
|
|
greenfield: bool
|
|
|
|
|
|
class KnowledgeFSWorkspaceCutoverService:
|
|
"""Coordinate one Workspace through the strict P8 phase sequence."""
|
|
|
|
_session_maker: sessionmaker[Session]
|
|
_clock: Callable[[], datetime]
|
|
_trusted_shadow_producers: frozenset[str]
|
|
_trusted_shadow_operators: frozenset[str]
|
|
_remote_factory: Callable[[], KnowledgeFSLifecycleRemotePort] | None
|
|
|
|
def __init__(
|
|
self,
|
|
session_maker: sessionmaker[Session],
|
|
clock: Callable[[], datetime] = naive_utc_now,
|
|
*,
|
|
trusted_shadow_producers: frozenset[str] = _DEFAULT_TRUSTED_SHADOW_PRODUCERS,
|
|
trusted_shadow_operators: frozenset[str] = _DEFAULT_TRUSTED_SHADOW_OPERATORS,
|
|
remote_factory: Callable[[], KnowledgeFSLifecycleRemotePort] | None = None,
|
|
):
|
|
self._session_maker = session_maker
|
|
self._clock = clock
|
|
self._trusted_shadow_producers = trusted_shadow_producers
|
|
self._trusted_shadow_operators = trusted_shadow_operators
|
|
self._remote_factory = remote_factory
|
|
|
|
def initialize_greenfield(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
owner_account_id: str,
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
"""Idempotently cut over a Workspace that has never owned KnowledgeFS state."""
|
|
|
|
try:
|
|
tenant_uuid = UUID(tenant_id)
|
|
owner_uuid = UUID(owner_account_id)
|
|
except ValueError as exc:
|
|
raise KnowledgeFSCutoverConflictError("Greenfield Workspace and owner identifiers must be UUIDs") from exc
|
|
if not self._trusted_shadow_producers or not self._trusted_shadow_operators:
|
|
raise KnowledgeFSCutoverGateBlockedError("Greenfield shadow attestation is not configured")
|
|
|
|
inventory = WorkspaceInventoryInput.model_validate(
|
|
{
|
|
"tenant_id": str(tenant_uuid),
|
|
"source_revision_watermark": _ZERO_REVISION_WATERMARK,
|
|
"task_watermark": 0,
|
|
"spaces": [],
|
|
}
|
|
)
|
|
for _ in range(_GREENFIELD_INITIALIZATION_MAX_ATTEMPTS):
|
|
ledger = self._greenfield_initialization_ledger(
|
|
tenant_id=tenant_id,
|
|
owner_account_id=owner_account_id,
|
|
)
|
|
if ledger is not None and _has_complete_product_cutover(ledger):
|
|
return ledger
|
|
try:
|
|
if ledger is None:
|
|
self.inventory(inventory, apply=True)
|
|
elif ledger.phase is KnowledgeFSWorkspaceCutoverPhase.INVENTORY:
|
|
self.backfill(inventory, apply=True)
|
|
elif ledger.phase is KnowledgeFSWorkspaceCutoverPhase.BACKFILL:
|
|
self.begin_shadow(
|
|
tenant_id=tenant_id,
|
|
expected_cas_version=ledger.cas_version,
|
|
started_at=self._clock(),
|
|
)
|
|
elif ledger.phase is KnowledgeFSWorkspaceCutoverPhase.SHADOW:
|
|
if ledger.shadow_completed_at is None:
|
|
completed_at = _naive_utc(self._clock()).replace(tzinfo=UTC)
|
|
self.complete_shadow(
|
|
ShadowCompletionInput.model_validate(
|
|
{
|
|
"schema_version": "knowledge-fs-p8-shadow-completion/v1",
|
|
"tenant_id": tenant_id,
|
|
"expected_cas_version": ledger.cas_version,
|
|
"producer": min(self._trusted_shadow_producers),
|
|
"completed_by_operator": min(self._trusted_shadow_operators),
|
|
"completed_by_account_id": str(owner_uuid),
|
|
"completed_at": completed_at,
|
|
"traffic_zero": True,
|
|
"traffic_zero_evidence": {
|
|
"schema_version": "knowledge-fs-greenfield-traffic-zero/v1",
|
|
"basis": "no-local-control-spaces",
|
|
},
|
|
}
|
|
),
|
|
apply=True,
|
|
)
|
|
elif not ledger.legacy_dependency_ready:
|
|
self.legacy_dependency_dashboard(
|
|
tenant_id=tenant_id,
|
|
dependencies=(),
|
|
expected_cas_version=ledger.cas_version,
|
|
checked_at=self._clock(),
|
|
apply=True,
|
|
)
|
|
else:
|
|
self.freeze(
|
|
tenant_id=tenant_id,
|
|
expected_cas_version=ledger.cas_version,
|
|
freeze_at=self._clock(),
|
|
)
|
|
elif ledger.phase is KnowledgeFSWorkspaceCutoverPhase.FROZEN:
|
|
if not _has_complete_greenfield_final_delta(ledger):
|
|
self.apply_final_delta(
|
|
FinalDeltaInput.model_validate(
|
|
{
|
|
"tenant_id": tenant_id,
|
|
"expected_cas_version": ledger.cas_version,
|
|
"final_revision_watermark": _ZERO_REVISION_WATERMARK,
|
|
"applied_revision_watermark": _ZERO_REVISION_WATERMARK,
|
|
"final_task_watermark": 0,
|
|
"applied_task_watermark": 0,
|
|
}
|
|
)
|
|
)
|
|
else:
|
|
cutover_at = _naive_utc(self._clock())
|
|
self.cutover(
|
|
tenant_id=tenant_id,
|
|
expected_cas_version=ledger.cas_version,
|
|
cutover_at=cutover_at,
|
|
rollback_cutoff_at=cutover_at + _GREENFIELD_ROLLBACK_WINDOW,
|
|
)
|
|
else:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Workspace is not eligible for automatic greenfield initialization"
|
|
)
|
|
except (IntegrityError, KnowledgeFSCutoverConflictError):
|
|
continue
|
|
raise KnowledgeFSCutoverConflictError("Greenfield Workspace initialization did not converge")
|
|
|
|
def inventory(self, payload: WorkspaceInventoryInput, *, apply: bool) -> CutoverInventoryReport:
|
|
"""Validate a read-only inventory and optionally create the initial CAS ledger."""
|
|
|
|
tenant_id = str(payload.tenant_id)
|
|
self._assert_unique_inventory(payload.spaces)
|
|
dispositions = self._task_disposition_counts(payload.spaces)
|
|
report = CutoverInventoryReport(
|
|
tenant_id=tenant_id,
|
|
spaces=len(payload.spaces),
|
|
permissions=sum(len(space.permissions) for space in payload.spaces),
|
|
unresolved_subjects=sum(space.owner_account_id is None for space in payload.spaces)
|
|
+ sum(permission.account_id is None for space in payload.spaces for permission in space.permissions),
|
|
unknown_external_access=sum(space.external_access_enabled is None for space in payload.spaces),
|
|
legacy_keys_requiring_rotation=sum(len(space.legacy_api_keys) for space in payload.spaces),
|
|
tasks_by_disposition=dispositions,
|
|
orphan_resources=sum(len(space.orphan_resource_ids) for space in payload.spaces),
|
|
applied=apply,
|
|
replayed=False,
|
|
)
|
|
if not apply:
|
|
return report
|
|
source_watermark = payload.source_revision_watermark.to_record()
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
current = repository.get_ledger(tenant_id=tenant_id)
|
|
if current is not None:
|
|
if (
|
|
current.source_revision_watermark != source_watermark
|
|
or current.source_task_watermark != payload.task_watermark
|
|
):
|
|
raise KnowledgeFSCutoverConflictError("Workspace inventory watermarks conflict with the ledger")
|
|
return report._replace(replayed=True)
|
|
repository.add_ledger(
|
|
KnowledgeFSWorkspaceCutoverLedger(
|
|
tenant_id=tenant_id,
|
|
source_revision_watermark=source_watermark,
|
|
applied_revision_watermark=source_watermark,
|
|
source_task_watermark=payload.task_watermark,
|
|
applied_task_watermark=payload.task_watermark,
|
|
)
|
|
)
|
|
return report
|
|
|
|
def backfill(self, payload: WorkspaceInventoryInput, *, apply: bool) -> CutoverBackfillReport:
|
|
"""Backfill only independent control-plane rows; unresolved inputs stay quarantined."""
|
|
|
|
tenant_id = str(payload.tenant_id)
|
|
self._assert_unique_inventory(payload.spaces)
|
|
if not apply:
|
|
return self._project_backfill(payload)
|
|
registered = 0
|
|
granted = 0
|
|
quarantined = 0
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
if ledger.phase not in {
|
|
KnowledgeFSWorkspaceCutoverPhase.INVENTORY,
|
|
KnowledgeFSWorkspaceCutoverPhase.BACKFILL,
|
|
}:
|
|
raise KnowledgeFSCutoverConflictError("Backfill is closed after shadow validation starts")
|
|
if (
|
|
ledger.source_revision_watermark != payload.source_revision_watermark.to_record()
|
|
or ledger.source_task_watermark != payload.task_watermark
|
|
):
|
|
raise KnowledgeFSCutoverConflictError("Backfill input does not match the inventoried watermarks")
|
|
for space in payload.spaces:
|
|
result = self._backfill_space(
|
|
session,
|
|
repository,
|
|
ledger,
|
|
space,
|
|
payload.source_revision_watermark.to_record(),
|
|
)
|
|
registered += result[0]
|
|
granted += result[1]
|
|
quarantined += result[2]
|
|
session.flush()
|
|
open_issues = repository.count_open_issues(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
if ledger.phase is KnowledgeFSWorkspaceCutoverPhase.INVENTORY and open_issues == 0:
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.BACKFILL,
|
|
),
|
|
)
|
|
phase = KnowledgeFSWorkspaceCutoverPhase.BACKFILL.value
|
|
else:
|
|
phase = ledger.phase.value
|
|
return CutoverBackfillReport(tenant_id, registered, granted, quarantined, open_issues, phase, True)
|
|
|
|
def _greenfield_initialization_ledger(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
owner_account_id: str,
|
|
) -> KnowledgeFSWorkspaceCutoverLedger | None:
|
|
with self._session_maker() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = repository.get_ledger(tenant_id=tenant_id)
|
|
control_spaces = tuple(
|
|
session.scalars(
|
|
sa.select(KnowledgeFSControlSpace)
|
|
.where(KnowledgeFSControlSpace.tenant_id == tenant_id)
|
|
.order_by(KnowledgeFSControlSpace.created_at, KnowledgeFSControlSpace.id)
|
|
).all()
|
|
)
|
|
if ledger is None:
|
|
if control_spaces:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Workspace with existing KnowledgeFS control state is not greenfield"
|
|
)
|
|
return None
|
|
if not _is_zero_space_cutover_ledger(ledger) or ledger.rolled_back_at is not None:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Workspace migration state is not eligible for automatic greenfield initialization"
|
|
)
|
|
if ledger.legacy_dependency_report not in (None, []):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Workspace with legacy dependencies is not eligible for automatic greenfield initialization"
|
|
)
|
|
if control_spaces:
|
|
anchor_id = _greenfield_activation_anchor_id(tenant_id)
|
|
provisioning_key = _greenfield_activation_provisioning_key(tenant_id)
|
|
if len(control_spaces) != 1:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Workspace with existing KnowledgeFS control state is not greenfield"
|
|
)
|
|
anchor = control_spaces[0]
|
|
if (
|
|
anchor.id != anchor_id
|
|
or anchor.state is not KnowledgeFSControlSpaceState.DELETED
|
|
or anchor.knowledge_space_id is not None
|
|
or anchor.provisioning_key != provisioning_key
|
|
or anchor.owner_account_id != owner_account_id
|
|
or anchor.lifecycle_operation_id != anchor_id
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Workspace greenfield activation audit anchor is inconsistent"
|
|
)
|
|
session.expunge(ledger)
|
|
return ledger
|
|
|
|
def begin_shadow(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
expected_cas_version: int,
|
|
started_at: datetime | None = None,
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
shadow_started_at = _naive_utc(started_at or self._clock())
|
|
if shadow_started_at > _naive_utc(self._clock()):
|
|
raise KnowledgeFSCutoverConflictError("Shadow start timestamp cannot be in the future")
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.BACKFILL, expected_cas_version)
|
|
self._assert_no_open_gates(repository, ledger)
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.SHADOW,
|
|
shadow_started_at=shadow_started_at,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def record_shadow_report(
|
|
self,
|
|
observations: Iterable[ShadowAuthorizationObservationInput],
|
|
*,
|
|
apply: bool,
|
|
) -> ShadowAuthorizationReport:
|
|
items = tuple(observations)
|
|
if not items:
|
|
raise KnowledgeFSCutoverConflictError("Shadow report requires at least one observation")
|
|
tenant_ids = {str(item.tenant_id) for item in items}
|
|
if len(tenant_ids) != 1:
|
|
raise KnowledgeFSCutoverConflictError("Shadow report cannot mix Workspaces")
|
|
tenant_id = tenant_ids.pop()
|
|
if len({item.producer for item in items}) != 1:
|
|
raise KnowledgeFSCutoverConflictError("Shadow report cannot mix producers")
|
|
decisions = [self._shadow_decision(item) for item in items]
|
|
counts = {decision: decisions.count(decision) for decision in KnowledgeFSShadowAuthorizationDecision}
|
|
replayed = 0
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
if ledger.phase is not KnowledgeFSWorkspaceCutoverPhase.SHADOW:
|
|
raise KnowledgeFSCutoverConflictError("Shadow observations require the shadow phase")
|
|
for item, decision in zip(items, decisions, strict=True):
|
|
self._validate_shadow_observation(ledger, item)
|
|
replayed += self._record_shadow_observation(
|
|
repository,
|
|
ledger,
|
|
item,
|
|
decision,
|
|
apply=apply,
|
|
)
|
|
if apply:
|
|
session.flush()
|
|
open_diffs = repository.count_unapproved_shadow_diffs(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
else:
|
|
open_diffs = sum(decision is not KnowledgeFSShadowAuthorizationDecision.MATCH for decision in decisions)
|
|
return ShadowAuthorizationReport(
|
|
tenant_id,
|
|
counts[KnowledgeFSShadowAuthorizationDecision.MATCH],
|
|
counts[KnowledgeFSShadowAuthorizationDecision.TIGHTENED],
|
|
counts[KnowledgeFSShadowAuthorizationDecision.EXPANDED],
|
|
counts[KnowledgeFSShadowAuthorizationDecision.UNKNOWN],
|
|
open_diffs,
|
|
apply,
|
|
len(items) - replayed,
|
|
replayed,
|
|
)
|
|
|
|
def complete_shadow(self, payload: ShadowCompletionInput, *, apply: bool) -> ShadowCompletionReport:
|
|
"""Close a trusted shadow window and persist an immutable digest of its evidence."""
|
|
|
|
tenant_id = str(payload.tenant_id)
|
|
completed_at = _naive_utc(payload.completed_at)
|
|
now = _naive_utc(self._clock())
|
|
if payload.producer not in self._trusted_shadow_producers:
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow completion producer is not trusted")
|
|
if payload.completed_by_operator not in self._trusted_shadow_operators:
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow completion operator is not trusted")
|
|
if completed_at > now:
|
|
raise KnowledgeFSCutoverConflictError("Shadow completion timestamp cannot be in the future")
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
if ledger.phase is not KnowledgeFSWorkspaceCutoverPhase.SHADOW:
|
|
raise KnowledgeFSCutoverConflictError("Shadow completion requires the shadow phase")
|
|
if ledger.shadow_started_at is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow start evidence is missing")
|
|
if completed_at < ledger.shadow_started_at:
|
|
raise KnowledgeFSCutoverConflictError("Shadow completion predates shadow start")
|
|
observations = repository.list_shadow_observations(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
latest_revision: KnowledgeFSCutoverRevisionWatermark | None
|
|
if payload.traffic_zero:
|
|
if observations:
|
|
raise KnowledgeFSCutoverConflictError(
|
|
"Traffic-zero completion conflicts with persisted observations"
|
|
)
|
|
latest_revision = None
|
|
evidence_digest = _shadow_traffic_zero_digest(payload)
|
|
else:
|
|
if not observations:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Shadow completion requires at least one observation or traffic-zero evidence"
|
|
)
|
|
if any(item.producer != payload.producer for item in observations):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Shadow observations were not produced by the completing trusted producer"
|
|
)
|
|
window_started_at = _naive_utc(cast(datetime, payload.window_started_at))
|
|
window_ended_at = _naive_utc(cast(datetime, payload.window_ended_at))
|
|
if window_started_at < ledger.shadow_started_at or window_ended_at <= window_started_at:
|
|
raise KnowledgeFSCutoverConflictError("Shadow observation window is invalid")
|
|
if window_ended_at > completed_at:
|
|
raise KnowledgeFSCutoverConflictError("Shadow observation window ends after completion")
|
|
if any(
|
|
item.observed_at < window_started_at or item.observed_at > window_ended_at for item in observations
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Persisted shadow observations fall outside the completed window"
|
|
)
|
|
latest_revision = _latest_shadow_revision(observations)
|
|
required_revision = ledger.final_revision_watermark or ledger.source_revision_watermark
|
|
if not _watermark_at_least(latest_revision, required_revision):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Shadow latest observed revision is below the applicable cutover watermark"
|
|
)
|
|
evidence_digest = _shadow_observation_set_digest(observations)
|
|
report = ShadowCompletionReport(
|
|
tenant_id,
|
|
len(observations),
|
|
payload.traffic_zero,
|
|
evidence_digest,
|
|
latest_revision,
|
|
apply,
|
|
ledger.shadow_completed_at is not None,
|
|
)
|
|
if ledger.shadow_completed_at is not None:
|
|
self._assert_shadow_completion_replay(ledger, payload, report)
|
|
return report
|
|
if ledger.cas_version != payload.expected_cas_version:
|
|
raise KnowledgeFSCutoverConflictError("Cutover ledger version changed")
|
|
if not apply:
|
|
return report
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=ledger.phase,
|
|
shadow_completed_at=completed_at,
|
|
shadow_evidence_digest=evidence_digest,
|
|
shadow_observation_count=len(observations),
|
|
shadow_window_started_at=(
|
|
_naive_utc(payload.window_started_at) if payload.window_started_at is not None else None
|
|
),
|
|
shadow_window_ended_at=(
|
|
_naive_utc(payload.window_ended_at) if payload.window_ended_at is not None else None
|
|
),
|
|
shadow_traffic_zero=payload.traffic_zero,
|
|
shadow_traffic_zero_evidence=(
|
|
cast(dict[str, object], payload.traffic_zero_evidence)
|
|
if payload.traffic_zero_evidence is not None
|
|
else None
|
|
),
|
|
shadow_latest_observed_revision=latest_revision,
|
|
shadow_producer=payload.producer,
|
|
shadow_completed_by_operator=payload.completed_by_operator,
|
|
shadow_completed_by_account_id=str(payload.completed_by_account_id),
|
|
),
|
|
)
|
|
return report
|
|
|
|
def approve_shadow_diff(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
diff_key: str,
|
|
account_id: str,
|
|
approved_at: datetime,
|
|
) -> None:
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
diff = repository.get_shadow_diff(tenant_id=tenant_id, ledger_id=ledger.id, diff_key=diff_key)
|
|
if diff is None:
|
|
raise KnowledgeFSCutoverNotFoundError("Shadow authorization diff was not found")
|
|
if diff.dify_allowed or diff.decision is KnowledgeFSShadowAuthorizationDecision.EXPANDED:
|
|
raise KnowledgeFSCutoverGateBlockedError("Access-expanding diffs cannot be approved as fail-closed")
|
|
changed = repository.set_shadow_diff_status(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=diff_key,
|
|
expected_status=KnowledgeFSMigrationIssueStatus.OPEN,
|
|
new_status=KnowledgeFSMigrationIssueStatus.APPROVED_FAIL_CLOSED,
|
|
account_id=account_id,
|
|
changed_at=_naive_utc(approved_at),
|
|
)
|
|
if not changed and diff.status is not KnowledgeFSMigrationIssueStatus.APPROVED_FAIL_CLOSED:
|
|
raise KnowledgeFSCutoverConflictError("Shadow authorization diff changed during approval")
|
|
|
|
def approve_issue_fail_closed(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
issue_key: str,
|
|
account_id: str,
|
|
approved_at: datetime,
|
|
) -> None:
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
issue = repository.get_issue(tenant_id=tenant_id, ledger_id=ledger.id, issue_key=issue_key)
|
|
if issue is None:
|
|
raise KnowledgeFSCutoverNotFoundError("Migration issue was not found")
|
|
approvable = issue.kind in {
|
|
KnowledgeFSMigrationIssueKind.UNKNOWN_EXTERNAL_ACCESS,
|
|
KnowledgeFSMigrationIssueKind.UNRESOLVED_SUBJECT,
|
|
} and not issue.issue_key.startswith("owner:")
|
|
if not approvable:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Only unknown fail-closed authorization can be approved; this issue must be resolved"
|
|
)
|
|
changed = repository.set_issue_status(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
issue_key=issue_key,
|
|
expected_status=KnowledgeFSMigrationIssueStatus.OPEN,
|
|
new_status=KnowledgeFSMigrationIssueStatus.APPROVED_FAIL_CLOSED,
|
|
account_id=account_id,
|
|
changed_at=_naive_utc(approved_at),
|
|
)
|
|
if not changed and issue.status is not KnowledgeFSMigrationIssueStatus.APPROVED_FAIL_CLOSED:
|
|
raise KnowledgeFSCutoverConflictError("Migration issue changed during approval")
|
|
|
|
def resolve_issue(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
issue_key: str,
|
|
account_id: str,
|
|
resolved_at: datetime,
|
|
) -> None:
|
|
"""Resolve remediated non-unknown evidence; unknown authorization requires fail-closed approval."""
|
|
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
issue = repository.get_issue(tenant_id=tenant_id, ledger_id=ledger.id, issue_key=issue_key)
|
|
if issue is None:
|
|
raise KnowledgeFSCutoverNotFoundError("Migration issue was not found")
|
|
if issue.kind in {
|
|
KnowledgeFSMigrationIssueKind.UNKNOWN_EXTERNAL_ACCESS,
|
|
KnowledgeFSMigrationIssueKind.UNRESOLVED_SUBJECT,
|
|
} and not issue.issue_key.startswith("owner:"):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Unknown authorization evidence can only be approved as fail-closed"
|
|
)
|
|
if issue.status is KnowledgeFSMigrationIssueStatus.RESOLVED:
|
|
return
|
|
changed = repository.set_issue_status(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
issue_key=issue_key,
|
|
expected_status=issue.status,
|
|
new_status=KnowledgeFSMigrationIssueStatus.RESOLVED,
|
|
account_id=account_id,
|
|
changed_at=_naive_utc(resolved_at),
|
|
)
|
|
if not changed:
|
|
raise KnowledgeFSCutoverConflictError("Migration issue changed during resolution")
|
|
|
|
def resolve_shadow_diff(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
diff_key: str,
|
|
account_id: str,
|
|
resolved_at: datetime,
|
|
) -> None:
|
|
"""Resolve a remediated known diff; unknown evidence remains approval-only."""
|
|
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
diff = repository.get_shadow_diff(tenant_id=tenant_id, ledger_id=ledger.id, diff_key=diff_key)
|
|
if diff is None:
|
|
raise KnowledgeFSCutoverNotFoundError("Shadow authorization diff was not found")
|
|
if diff.decision is KnowledgeFSShadowAuthorizationDecision.UNKNOWN:
|
|
raise KnowledgeFSCutoverGateBlockedError("Unknown shadow evidence can only be approved as fail-closed")
|
|
if diff.decision is KnowledgeFSShadowAuthorizationDecision.EXPANDED:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Expanded shadow evidence requires a safe shadow re-evaluation"
|
|
)
|
|
if diff.status is KnowledgeFSMigrationIssueStatus.RESOLVED:
|
|
return
|
|
changed = repository.set_shadow_diff_status(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=diff_key,
|
|
expected_status=diff.status,
|
|
new_status=KnowledgeFSMigrationIssueStatus.RESOLVED,
|
|
account_id=account_id,
|
|
changed_at=_naive_utc(resolved_at),
|
|
)
|
|
if not changed:
|
|
raise KnowledgeFSCutoverConflictError("Shadow authorization diff changed during resolution")
|
|
|
|
def resolve_quarantine(
|
|
self,
|
|
payload: QuarantineResolutionInput,
|
|
*,
|
|
apply: bool,
|
|
) -> QuarantineResolutionReport:
|
|
"""Resolve one quarantined source with immutable evidence and row-version CAS."""
|
|
|
|
tenant_id = str(payload.tenant_id)
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
item = repository.get_quarantine(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
source_kind=payload.source_kind,
|
|
source_id=payload.source_id,
|
|
)
|
|
if item is None:
|
|
raise KnowledgeFSCutoverNotFoundError("Quarantine item was not found")
|
|
resolved_at = _naive_utc(payload.resolved_at)
|
|
if resolved_at > _naive_utc(self._clock()):
|
|
raise KnowledgeFSCutoverConflictError("Quarantine resolution timestamp cannot be in the future")
|
|
evidence = cast(dict[str, object], payload.evidence.model_dump(mode="json"))
|
|
if item.disposition is KnowledgeFSMigrationQuarantineDisposition.RESOLVED:
|
|
if (
|
|
item.row_version != payload.expected_row_version + 1
|
|
or item.resolved_by_operator != payload.resolved_by_operator
|
|
or item.resolved_by_account_id != str(payload.resolved_by_account_id)
|
|
or item.evidence != evidence
|
|
or item.resolved_at != resolved_at
|
|
):
|
|
raise KnowledgeFSCutoverConflictError(
|
|
"Quarantine item was already resolved with different evidence"
|
|
)
|
|
return QuarantineResolutionReport(
|
|
tenant_id,
|
|
item.source_kind.value,
|
|
item.source_id,
|
|
item.disposition.value,
|
|
item.row_version,
|
|
apply,
|
|
True,
|
|
)
|
|
if ledger.phase in {
|
|
KnowledgeFSWorkspaceCutoverPhase.CUTOVER,
|
|
KnowledgeFSWorkspaceCutoverPhase.OBSERVING,
|
|
KnowledgeFSWorkspaceCutoverPhase.READY_FOR_CLEANUP,
|
|
}:
|
|
raise KnowledgeFSCutoverConflictError("Quarantine resolution is closed after cutover")
|
|
if item.row_version != payload.expected_row_version:
|
|
raise KnowledgeFSCutoverConflictError("Quarantine item version changed")
|
|
self._validate_quarantine_resolution(session, ledger, item, payload, resolved_at)
|
|
if not apply:
|
|
return QuarantineResolutionReport(
|
|
tenant_id,
|
|
item.source_kind.value,
|
|
item.source_id,
|
|
item.disposition.value,
|
|
item.row_version,
|
|
False,
|
|
False,
|
|
)
|
|
changed = repository.resolve_quarantine(
|
|
KnowledgeFSQuarantineCASUpdate(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
source_kind=item.source_kind,
|
|
source_id=item.source_id,
|
|
expected_disposition=item.disposition,
|
|
expected_row_version=item.row_version,
|
|
resolved_by_operator=payload.resolved_by_operator,
|
|
resolved_by_account_id=str(payload.resolved_by_account_id),
|
|
evidence=evidence,
|
|
resolved_at=resolved_at,
|
|
)
|
|
)
|
|
if not changed:
|
|
raise KnowledgeFSCutoverConflictError("Quarantine item changed during resolution")
|
|
return QuarantineResolutionReport(
|
|
tenant_id,
|
|
item.source_kind.value,
|
|
item.source_id,
|
|
KnowledgeFSMigrationQuarantineDisposition.RESOLVED.value,
|
|
item.row_version + 1,
|
|
True,
|
|
False,
|
|
)
|
|
|
|
def _validate_quarantine_resolution(
|
|
self,
|
|
session: Session,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
item: KnowledgeFSMigrationQuarantine,
|
|
payload: QuarantineResolutionInput,
|
|
resolved_at: datetime,
|
|
) -> None:
|
|
if isinstance(payload.evidence, LegacyCredentialRotationEvidence):
|
|
credential_evidence = payload.evidence
|
|
if credential_evidence.legacy_key_id != item.source_id:
|
|
raise KnowledgeFSCutoverConflictError("Credential rotation evidence names a different legacy key")
|
|
knowledge_space_id = str(credential_evidence.knowledge_space_id)
|
|
if item.details.get("knowledge_space_id") != knowledge_space_id:
|
|
raise KnowledgeFSCutoverConflictError(
|
|
"Credential rotation evidence names a different KnowledgeFS Space"
|
|
)
|
|
if _naive_utc(credential_evidence.legacy_revoked_at) > resolved_at:
|
|
raise KnowledgeFSCutoverConflictError("Legacy credential revocation occurred after resolution")
|
|
control_space = session.scalar(
|
|
sa.select(KnowledgeFSControlSpace).where(
|
|
KnowledgeFSControlSpace.id == str(credential_evidence.control_space_id),
|
|
KnowledgeFSControlSpace.tenant_id == ledger.tenant_id,
|
|
KnowledgeFSControlSpace.knowledge_space_id == knowledge_space_id,
|
|
KnowledgeFSControlSpace.state == KnowledgeFSControlSpaceState.ACTIVE,
|
|
)
|
|
)
|
|
if control_space is None:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Credential rotation requires the tenant-owned active control-space"
|
|
)
|
|
credential = session.scalar(
|
|
sa.select(KnowledgeFSApiCredential).where(
|
|
KnowledgeFSApiCredential.id == str(credential_evidence.dify_credential_id),
|
|
KnowledgeFSApiCredential.tenant_id == ledger.tenant_id,
|
|
KnowledgeFSApiCredential.control_space_id == control_space.id,
|
|
)
|
|
)
|
|
if (
|
|
credential is None
|
|
or credential.status is not KnowledgeFSApiCredentialStatus.ACTIVE
|
|
or credential.revision != credential_evidence.dify_credential_revision
|
|
or (credential.expires_at is not None and credential.expires_at <= resolved_at)
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Credential rotation requires a matching active Dify-managed credential"
|
|
)
|
|
if credential.credential_prefix == item.details.get(
|
|
"prefix"
|
|
) and credential.credential_last4 == item.details.get("last4"):
|
|
raise KnowledgeFSCutoverGateBlockedError("Credential rotation did not replace the legacy credential")
|
|
return
|
|
if isinstance(payload.evidence, LegacyTaskResolutionEvidence):
|
|
task_evidence = payload.evidence
|
|
expected_actions = {
|
|
KnowledgeFSMigrationQuarantineDisposition.MIGRATABLE: "migrate",
|
|
KnowledgeFSMigrationQuarantineDisposition.WAIT_FOR_COMPLETION: "wait",
|
|
KnowledgeFSMigrationQuarantineDisposition.CANCEL: "cancel",
|
|
KnowledgeFSMigrationQuarantineDisposition.ISOLATE: "isolate",
|
|
}
|
|
expected_results = {
|
|
"migrate": {"migrated"},
|
|
"wait": {"completed", "failed", "canceled"},
|
|
"cancel": {"canceled"},
|
|
"isolate": {"isolated"},
|
|
}
|
|
if task_evidence.task_id != item.source_id or task_evidence.action != expected_actions.get(
|
|
item.disposition
|
|
):
|
|
raise KnowledgeFSCutoverConflictError("Task evidence does not match its classified disposition")
|
|
if task_evidence.resulting_state not in expected_results[task_evidence.action]:
|
|
raise KnowledgeFSCutoverConflictError("Task evidence has an invalid resulting state")
|
|
if task_evidence.final_task_watermark < ledger.source_task_watermark:
|
|
raise KnowledgeFSCutoverGateBlockedError("Task evidence is below the inventoried task watermark")
|
|
|
|
def legacy_dependency_dashboard(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
dependencies: Iterable[LegacyDependencyInput],
|
|
expected_cas_version: int | None,
|
|
checked_at: datetime,
|
|
apply: bool,
|
|
) -> LegacyDependencyDashboard:
|
|
items = tuple(dependencies)
|
|
active_rows = sum(item.active_rows for item in items)
|
|
ready = active_rows == 0
|
|
report = LegacyDependencyDashboard(
|
|
tenant_id,
|
|
ready,
|
|
sum(item.kind == "permission_snapshot" for item in items),
|
|
sum(item.kind == "foreign_key" for item in items),
|
|
active_rows,
|
|
tuple(item.model_dump(mode="json") for item in items),
|
|
apply,
|
|
)
|
|
if not apply:
|
|
return report
|
|
if expected_cas_version is None:
|
|
raise KnowledgeFSCutoverConflictError("Applying a dependency dashboard requires expected_cas_version")
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
if ledger.cas_version != expected_cas_version:
|
|
raise KnowledgeFSCutoverConflictError("Cutover ledger version changed")
|
|
for item in items:
|
|
if item.active_rows > 0:
|
|
kind = (
|
|
KnowledgeFSMigrationIssueKind.LEGACY_SNAPSHOT_DEPENDENCY
|
|
if item.kind == "permission_snapshot"
|
|
else KnowledgeFSMigrationIssueKind.LEGACY_FOREIGN_KEY_DEPENDENCY
|
|
)
|
|
self._ensure_issue(
|
|
repository,
|
|
ledger,
|
|
issue_key=f"legacy:{item.dependency_key}",
|
|
kind=kind,
|
|
resource_type="legacy_dependency",
|
|
resource_id=item.dependency_key,
|
|
details=item.model_dump(mode="json"),
|
|
)
|
|
dependency_keys = {f"legacy:{item.dependency_key}" for item in items if item.active_rows > 0}
|
|
for issue in repository.list_issues(tenant_id=tenant_id, ledger_id=ledger.id):
|
|
if (
|
|
issue.kind
|
|
in {
|
|
KnowledgeFSMigrationIssueKind.LEGACY_SNAPSHOT_DEPENDENCY,
|
|
KnowledgeFSMigrationIssueKind.LEGACY_FOREIGN_KEY_DEPENDENCY,
|
|
}
|
|
and issue.issue_key not in dependency_keys
|
|
and issue.status is not KnowledgeFSMigrationIssueStatus.RESOLVED
|
|
):
|
|
repository.set_issue_status(
|
|
tenant_id=tenant_id,
|
|
ledger_id=ledger.id,
|
|
issue_key=issue.issue_key,
|
|
expected_status=issue.status,
|
|
new_status=KnowledgeFSMigrationIssueStatus.RESOLVED,
|
|
account_id=tenant_id,
|
|
changed_at=_naive_utc(checked_at),
|
|
)
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=ledger.phase,
|
|
legacy_dependency_report=[dict(dependency) for dependency in report.dependencies],
|
|
legacy_dependency_checked_at=_naive_utc(checked_at),
|
|
legacy_dependency_ready=ready,
|
|
),
|
|
)
|
|
return report
|
|
|
|
def freeze(
|
|
self, *, tenant_id: str, expected_cas_version: int, freeze_at: datetime
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
freeze_time = _naive_utc(freeze_at)
|
|
if freeze_time > _naive_utc(self._clock()):
|
|
raise KnowledgeFSCutoverConflictError("Freeze timestamp cannot be in the future")
|
|
|
|
# Snapshot all gates and the control identity, then release the transaction before the
|
|
# remote call. Local FROZEN state is not visible until KnowledgeFS durably acknowledges.
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.SHADOW, expected_cas_version)
|
|
self._assert_shadow_evidence_complete(ledger)
|
|
self._assert_no_open_gates(repository, ledger)
|
|
self._assert_no_unresolved_quarantine(repository, ledger)
|
|
if not ledger.legacy_dependency_ready or ledger.legacy_dependency_checked_at is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Legacy snapshot/FK dependency dashboard is not ready")
|
|
anchor = self._require_activation_anchor(session, ledger=ledger)
|
|
freeze_request = _freeze_request(ledger, control_space_id=anchor.control_space_id)
|
|
|
|
if self._remote_factory is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote freeze is not configured")
|
|
try:
|
|
remote = self._remote_factory()
|
|
except RuntimeError as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote freeze is not configured") from exc
|
|
if remote is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote freeze is not configured")
|
|
if anchor.greenfield:
|
|
self._assert_remote_greenfield_namespace_empty(
|
|
remote,
|
|
tenant_id=tenant_id,
|
|
control_space_id=anchor.control_space_id,
|
|
)
|
|
try:
|
|
raw_ack = remote.freeze_dify_workspace_integration(freeze_request)
|
|
except KnowledgeFSLifecycleRemoteError as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote freeze was not acknowledged") from exc
|
|
if anchor.greenfield:
|
|
self._assert_remote_greenfield_namespace_empty(
|
|
remote,
|
|
tenant_id=tenant_id,
|
|
control_space_id=anchor.control_space_id,
|
|
)
|
|
try:
|
|
ack = (
|
|
raw_ack
|
|
if isinstance(raw_ack, KnowledgeFSDifyIntegrationFreezeAck)
|
|
else KnowledgeFSDifyIntegrationFreezeAck.model_validate(raw_ack, from_attributes=True)
|
|
)
|
|
except (TypeError, ValueError, ValidationError) as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote freeze acknowledgement was malformed") from exc
|
|
_assert_exact_freeze_ack(freeze_request, ack)
|
|
acknowledged_at = _naive_utc(self._clock())
|
|
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.SHADOW, expected_cas_version)
|
|
self._assert_shadow_evidence_complete(ledger)
|
|
self._assert_no_open_gates(repository, ledger)
|
|
self._assert_no_unresolved_quarantine(repository, ledger)
|
|
if not ledger.legacy_dependency_ready or ledger.legacy_dependency_checked_at is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Legacy snapshot/FK dependency dashboard is not ready")
|
|
anchor = self._require_activation_anchor(
|
|
session,
|
|
ledger=ledger,
|
|
control_space_id=freeze_request.control_space_id,
|
|
)
|
|
if _freeze_request(ledger, control_space_id=anchor.control_space_id) != freeze_request:
|
|
raise KnowledgeFSCutoverConflictError("Freeze evidence changed after remote acknowledgement")
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.FROZEN,
|
|
freeze_at=freeze_time,
|
|
remote_freeze_id=ack.freeze_id,
|
|
remote_freeze_revision=ack.freeze_revision,
|
|
remote_freeze_digest=ack.source_revision_digest,
|
|
remote_freeze_task_watermark=ack.source_task_watermark,
|
|
remote_freeze_control_space_id=anchor.control_space_id,
|
|
remote_freeze_frozen_at=_naive_utc(ack.frozen_at),
|
|
remote_freeze_updated_at=_naive_utc(ack.updated_at),
|
|
remote_freeze_acknowledged_at=acknowledged_at,
|
|
remote_freeze_applied=ack.applied,
|
|
remote_freeze_replayed=ack.replayed,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def apply_final_delta(self, payload: FinalDeltaInput) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
tenant_id = str(payload.tenant_id)
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.FROZEN, payload.expected_cas_version)
|
|
final_watermark = payload.final_revision_watermark.to_record()
|
|
applied_watermark = payload.applied_revision_watermark.to_record()
|
|
if not _watermark_at_least(final_watermark, ledger.source_revision_watermark):
|
|
raise KnowledgeFSCutoverConflictError("Final authorization watermark moved backwards")
|
|
if payload.final_task_watermark < ledger.source_task_watermark:
|
|
raise KnowledgeFSCutoverConflictError("Final task watermark moved backwards")
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=ledger.phase,
|
|
final_revision_watermark=final_watermark,
|
|
applied_revision_watermark=applied_watermark,
|
|
final_task_watermark=payload.final_task_watermark,
|
|
applied_task_watermark=payload.applied_task_watermark,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def cutover(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
expected_cas_version: int,
|
|
cutover_at: datetime,
|
|
rollback_cutoff_at: datetime,
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
cutover_time = _naive_utc(cutover_at)
|
|
cutoff = _naive_utc(rollback_cutoff_at)
|
|
if cutoff <= cutover_time:
|
|
raise KnowledgeFSCutoverConflictError("Rollback cutoff must be after cutover")
|
|
|
|
# The first short transaction snapshots all gate and activation evidence. It
|
|
# deliberately commits before capability issuance or any remote HTTP call.
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_cutover_ready(repository, ledger, expected_cas_version)
|
|
anchor = self._require_activation_anchor(session, ledger=ledger)
|
|
activation_request = _activation_request(ledger, control_space_id=anchor.control_space_id)
|
|
|
|
if self._remote_factory is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote activation is not configured")
|
|
try:
|
|
remote = self._remote_factory()
|
|
except RuntimeError as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote activation is not configured") from exc
|
|
if remote is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote activation is not configured")
|
|
try:
|
|
raw_ack = remote.activate_dify_workspace_integration(activation_request)
|
|
except KnowledgeFSLifecycleRemoteError as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS remote activation was not acknowledged") from exc
|
|
try:
|
|
ack = (
|
|
raw_ack
|
|
if isinstance(raw_ack, KnowledgeFSDifyIntegrationActivationAck)
|
|
else KnowledgeFSDifyIntegrationActivationAck.model_validate(raw_ack, from_attributes=True)
|
|
)
|
|
except (TypeError, ValueError, ValidationError) as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"KnowledgeFS remote activation acknowledgement was malformed"
|
|
) from exc
|
|
_assert_exact_activation_ack(activation_request, ack)
|
|
acknowledged_at = _naive_utc(self._clock())
|
|
|
|
# Re-open a transaction only after the exact durable ACK. Re-evaluate every
|
|
# gate and identity, then publish all switches and ACK evidence in one CAS.
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_cutover_ready(repository, ledger, expected_cas_version)
|
|
anchor = self._require_activation_anchor(
|
|
session,
|
|
ledger=ledger,
|
|
control_space_id=activation_request.control_space_id,
|
|
)
|
|
if _activation_request(ledger, control_space_id=anchor.control_space_id) != activation_request:
|
|
raise KnowledgeFSCutoverConflictError(
|
|
"Cutover activation evidence changed after remote acknowledgement"
|
|
)
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.CUTOVER,
|
|
cutover_at=cutover_time,
|
|
rollback_cutoff_at=cutoff,
|
|
product_routes_enabled=True,
|
|
capability_v2_enabled=True,
|
|
integrated_mode_enabled=True,
|
|
legacy_acl_read_only=True,
|
|
clear_smoke_results=True,
|
|
remote_activation_id=ack.activation_id,
|
|
remote_activation_revision=ack.activation_revision,
|
|
remote_activation_digest=ack.source_revision_digest,
|
|
remote_activation_control_space_id=anchor.control_space_id,
|
|
remote_activation_activated_at=_naive_utc(ack.activated_at),
|
|
remote_activation_updated_at=_naive_utc(ack.updated_at),
|
|
remote_activation_acknowledged_at=acknowledged_at,
|
|
remote_activation_applied=ack.applied,
|
|
remote_activation_replayed=ack.replayed,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def record_smoke_results(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
expected_cas_version: int,
|
|
results: CutoverSmokeResultsInput,
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
if str(results.tenant_id) != tenant_id:
|
|
raise KnowledgeFSCutoverConflictError("Smoke evidence belongs to a different Workspace")
|
|
observed_at = _naive_utc(results.observed_at)
|
|
if observed_at > _naive_utc(self._clock()):
|
|
raise KnowledgeFSCutoverConflictError("Smoke evidence timestamp cannot be in the future")
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.CUTOVER, expected_cas_version)
|
|
if ledger.cutover_at is None or observed_at < ledger.cutover_at:
|
|
raise KnowledgeFSCutoverConflictError("Smoke evidence predates Workspace cutover")
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=ledger.phase,
|
|
smoke_results=results.to_record(),
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def begin_observation(
|
|
self,
|
|
*,
|
|
tenant_id: str,
|
|
expected_cas_version: int,
|
|
started_at: datetime,
|
|
window_ends_at: datetime,
|
|
maximum_task_expires_at: datetime,
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
started = _naive_utc(started_at)
|
|
window_end = _naive_utc(window_ends_at)
|
|
task_expiry = _naive_utc(maximum_task_expires_at)
|
|
if window_end <= started:
|
|
raise KnowledgeFSCutoverConflictError("Observation window must end after it starts")
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.CUTOVER, expected_cas_version)
|
|
if not knowledge_fs_cutover_smoke_results_passed(ledger.smoke_results):
|
|
raise KnowledgeFSCutoverGateBlockedError("All cutover smoke checks must pass before observation")
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.OBSERVING,
|
|
observation_started_at=started,
|
|
observation_window_ends_at=window_end,
|
|
maximum_task_expires_at=task_expiry,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def complete_observation(
|
|
self, *, tenant_id: str, expected_cas_version: int, observed_at: datetime
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
observed = _naive_utc(observed_at)
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.OBSERVING, expected_cas_version)
|
|
if (
|
|
ledger.observation_window_ends_at is None
|
|
or observed < ledger.observation_window_ends_at
|
|
or ledger.maximum_task_expires_at is None
|
|
or observed < ledger.maximum_task_expires_at
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Observation window and maximum task TTL must both elapse before cleanup"
|
|
)
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.READY_FOR_CLEANUP,
|
|
observation_completed_at=observed,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def rollback(
|
|
self, *, tenant_id: str, expected_cas_version: int, rolled_back_at: datetime
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
rollback_time = _naive_utc(rolled_back_at)
|
|
now = _naive_utc(self._clock())
|
|
with self._session_maker.begin() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
if ledger.phase not in {
|
|
KnowledgeFSWorkspaceCutoverPhase.CUTOVER,
|
|
KnowledgeFSWorkspaceCutoverPhase.OBSERVING,
|
|
KnowledgeFSWorkspaceCutoverPhase.READY_FOR_CLEANUP,
|
|
}:
|
|
raise KnowledgeFSCutoverConflictError("Rollback is only available after cutover and before cleanup")
|
|
if ledger.cas_version != expected_cas_version:
|
|
raise KnowledgeFSCutoverConflictError("Cutover ledger version changed")
|
|
if (
|
|
ledger.rollback_cutoff_at is None
|
|
or rollback_time >= ledger.rollback_cutoff_at
|
|
or now >= ledger.rollback_cutoff_at
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError("Rollback cutoff has elapsed")
|
|
if rollback_time > now:
|
|
raise KnowledgeFSCutoverConflictError("Rollback timestamp cannot be in the future")
|
|
if ledger.irreversible_cleanup_at is not None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Irreversible cleanup has already started")
|
|
# Legacy ACL remains read-only: post-cutover rollback is a safe traffic stop,
|
|
# never a return to the stale KnowledgeFS authorization source. Remote
|
|
# activation and its acknowledgement are intentionally immutable here.
|
|
self._cas(
|
|
repository,
|
|
KnowledgeFSCutoverCASUpdate(
|
|
tenant_id=tenant_id,
|
|
expected_phase=ledger.phase,
|
|
expected_cas_version=ledger.cas_version,
|
|
new_phase=KnowledgeFSWorkspaceCutoverPhase.FROZEN,
|
|
rolled_back_at=rollback_time,
|
|
product_routes_enabled=False,
|
|
),
|
|
)
|
|
return self._require_ledger(repository, tenant_id)
|
|
|
|
def status(self, *, tenant_id: str) -> dict[str, object]:
|
|
with self._session_maker() as session:
|
|
repository = SQLAlchemyKnowledgeFSCutoverRepository(session)
|
|
ledger = self._require_ledger(repository, tenant_id)
|
|
issues = repository.list_issues(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
quarantine = repository.list_quarantine(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
shadow_diffs = repository.list_shadow_diffs(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
shadow_observations = repository.list_shadow_observations(tenant_id=tenant_id, ledger_id=ledger.id)
|
|
return {
|
|
"tenant_id": ledger.tenant_id,
|
|
"phase": ledger.phase.value,
|
|
"cas_version": ledger.cas_version,
|
|
"source_revision_watermark": ledger.source_revision_watermark,
|
|
"final_revision_watermark": ledger.final_revision_watermark,
|
|
"applied_revision_watermark": ledger.applied_revision_watermark,
|
|
"source_task_watermark": ledger.source_task_watermark,
|
|
"final_task_watermark": ledger.final_task_watermark,
|
|
"applied_task_watermark": ledger.applied_task_watermark,
|
|
"feature_state": {
|
|
"product_routes_enabled": ledger.product_routes_enabled,
|
|
"capability_v2_enabled": ledger.capability_v2_enabled,
|
|
"integrated_mode_enabled": ledger.integrated_mode_enabled,
|
|
"legacy_acl_read_only": ledger.legacy_acl_read_only,
|
|
},
|
|
"smoke_results": ledger.smoke_results,
|
|
"legacy_dependency_report": ledger.legacy_dependency_report,
|
|
"legacy_dependency_ready": ledger.legacy_dependency_ready,
|
|
"shadow_started_at": _iso(ledger.shadow_started_at),
|
|
"shadow_completed_at": _iso(ledger.shadow_completed_at),
|
|
"shadow_evidence_digest": ledger.shadow_evidence_digest,
|
|
"shadow_observation_count": ledger.shadow_observation_count,
|
|
"shadow_window_started_at": _iso(ledger.shadow_window_started_at),
|
|
"shadow_window_ended_at": _iso(ledger.shadow_window_ended_at),
|
|
"shadow_traffic_zero": ledger.shadow_traffic_zero,
|
|
"shadow_traffic_zero_evidence": ledger.shadow_traffic_zero_evidence,
|
|
"shadow_latest_observed_revision": ledger.shadow_latest_observed_revision,
|
|
"shadow_producer": ledger.shadow_producer,
|
|
"shadow_completed_by_operator": ledger.shadow_completed_by_operator,
|
|
"shadow_completed_by_account_id": ledger.shadow_completed_by_account_id,
|
|
"remote_freeze_id": ledger.remote_freeze_id,
|
|
"remote_freeze_revision": ledger.remote_freeze_revision,
|
|
"remote_freeze_digest": ledger.remote_freeze_digest,
|
|
"remote_freeze_task_watermark": ledger.remote_freeze_task_watermark,
|
|
"remote_freeze_control_space_id": ledger.remote_freeze_control_space_id,
|
|
"remote_freeze_frozen_at": _iso(ledger.remote_freeze_frozen_at),
|
|
"remote_freeze_updated_at": _iso(ledger.remote_freeze_updated_at),
|
|
"remote_freeze_acknowledged_at": _iso(ledger.remote_freeze_acknowledged_at),
|
|
"remote_freeze_applied": ledger.remote_freeze_applied,
|
|
"remote_freeze_replayed": ledger.remote_freeze_replayed,
|
|
"remote_activation_id": ledger.remote_activation_id,
|
|
"remote_activation_revision": ledger.remote_activation_revision,
|
|
"remote_activation_digest": ledger.remote_activation_digest,
|
|
"remote_activation_control_space_id": ledger.remote_activation_control_space_id,
|
|
"remote_activation_activated_at": _iso(ledger.remote_activation_activated_at),
|
|
"remote_activation_updated_at": _iso(ledger.remote_activation_updated_at),
|
|
"remote_activation_acknowledged_at": _iso(ledger.remote_activation_acknowledged_at),
|
|
"remote_activation_applied": ledger.remote_activation_applied,
|
|
"remote_activation_replayed": ledger.remote_activation_replayed,
|
|
"open_issues": sum(issue.status is KnowledgeFSMigrationIssueStatus.OPEN for issue in issues),
|
|
"open_shadow_diffs": sum(diff.status is KnowledgeFSMigrationIssueStatus.OPEN for diff in shadow_diffs),
|
|
"quarantine_count": len(quarantine),
|
|
"unresolved_cutover_quarantine": sum(
|
|
item.source_kind in _CUTOVER_QUARANTINE_KINDS
|
|
and item.disposition is not KnowledgeFSMigrationQuarantineDisposition.RESOLVED
|
|
for item in quarantine
|
|
),
|
|
"issues": [
|
|
{
|
|
"issue_key": issue.issue_key,
|
|
"kind": issue.kind.value,
|
|
"status": issue.status.value,
|
|
"resource_type": issue.resource_type,
|
|
"resource_id": issue.resource_id,
|
|
"details": issue.details,
|
|
"approved_by_account_id": issue.approved_by_account_id,
|
|
"approved_at": _iso(issue.approved_at),
|
|
"resolved_by_account_id": issue.resolved_by_account_id,
|
|
"resolved_at": _iso(issue.resolved_at),
|
|
}
|
|
for issue in issues
|
|
],
|
|
"quarantine": [
|
|
{
|
|
"source_kind": item.source_kind.value,
|
|
"source_id": item.source_id,
|
|
"reason_code": item.reason_code,
|
|
"disposition": item.disposition.value,
|
|
"details": item.details,
|
|
"resolved_by_operator": item.resolved_by_operator,
|
|
"resolved_by_account_id": item.resolved_by_account_id,
|
|
"evidence": item.evidence,
|
|
"resolved_at": _iso(item.resolved_at),
|
|
"row_version": item.row_version,
|
|
}
|
|
for item in quarantine
|
|
],
|
|
"shadow_diffs": [
|
|
{
|
|
"diff_key": diff.diff_key,
|
|
"control_space_id": diff.control_space_id,
|
|
"principal": diff.principal,
|
|
"legacy_allowed": diff.legacy_allowed,
|
|
"dify_allowed": diff.dify_allowed,
|
|
"decision": diff.decision.value,
|
|
"reason": diff.reason,
|
|
"observed_revision": diff.observed_revision,
|
|
"status": diff.status.value,
|
|
"approved_by_account_id": diff.approved_by_account_id,
|
|
"approved_at": _iso(diff.approved_at),
|
|
"resolved_by_account_id": diff.resolved_by_account_id,
|
|
"resolved_at": _iso(diff.resolved_at),
|
|
"current_evidence_digest": diff.current_evidence_digest,
|
|
"last_observed_at": _iso(diff.last_observed_at),
|
|
"row_version": diff.row_version,
|
|
}
|
|
for diff in shadow_diffs
|
|
],
|
|
"shadow_observations": [
|
|
{
|
|
"diff_key": item.diff_key,
|
|
"producer": item.producer,
|
|
"control_space_id": item.control_space_id,
|
|
"principal": item.principal,
|
|
"legacy_allowed": item.legacy_allowed,
|
|
"dify_allowed": item.dify_allowed,
|
|
"decision": item.decision.value,
|
|
"reason": item.reason,
|
|
"observed_revision": item.observed_revision,
|
|
"observed_at": _iso(item.observed_at),
|
|
"evidence_digest": item.evidence_digest,
|
|
}
|
|
for item in shadow_observations
|
|
],
|
|
"freeze_at": _iso(ledger.freeze_at),
|
|
"cutover_at": _iso(ledger.cutover_at),
|
|
"rolled_back_at": _iso(ledger.rolled_back_at),
|
|
"rollback_cutoff_at": _iso(ledger.rollback_cutoff_at),
|
|
"observation_started_at": _iso(ledger.observation_started_at),
|
|
"observation_window_ends_at": _iso(ledger.observation_window_ends_at),
|
|
"maximum_task_expires_at": _iso(ledger.maximum_task_expires_at),
|
|
"observation_completed_at": _iso(ledger.observation_completed_at),
|
|
"legacy_dependency_checked_at": _iso(ledger.legacy_dependency_checked_at),
|
|
"irreversible_cleanup_at": _iso(ledger.irreversible_cleanup_at),
|
|
}
|
|
|
|
@staticmethod
|
|
def _task_disposition_counts(spaces: Iterable[LegacySpaceInventoryInput]) -> dict[str, int]:
|
|
counts = {disposition.value: 0 for disposition in KnowledgeFSMigrationQuarantineDisposition}
|
|
for space in spaces:
|
|
for task in space.tasks:
|
|
counts[_classify_task(task).value] += 1
|
|
return {key: value for key, value in counts.items() if value > 0}
|
|
|
|
@staticmethod
|
|
def _assert_unique_inventory(spaces: Iterable[LegacySpaceInventoryInput]) -> None:
|
|
knowledge_space_ids: set[UUID] = set()
|
|
provisioning_keys: set[str] = set()
|
|
for space in spaces:
|
|
if space.knowledge_space_id in knowledge_space_ids:
|
|
raise KnowledgeFSCutoverConflictError("Inventory contains a duplicate KnowledgeFS Space ID")
|
|
if space.provisioning_key in provisioning_keys:
|
|
raise KnowledgeFSCutoverConflictError("Inventory contains a duplicate provisioning key")
|
|
knowledge_space_ids.add(space.knowledge_space_id)
|
|
provisioning_keys.add(space.provisioning_key)
|
|
resolved_permissions: set[UUID] = set()
|
|
for permission in space.permissions:
|
|
if permission.account_id is None:
|
|
continue
|
|
if permission.account_id in resolved_permissions:
|
|
raise KnowledgeFSCutoverConflictError("Inventory contains duplicate resolved permission subjects")
|
|
if permission.account_id == space.owner_account_id and permission.role != "owner":
|
|
raise KnowledgeFSCutoverConflictError("Inventory cannot downgrade the registered owner")
|
|
resolved_permissions.add(permission.account_id)
|
|
|
|
@staticmethod
|
|
def _project_backfill(payload: WorkspaceInventoryInput) -> CutoverBackfillReport:
|
|
registered = sum(space.owner_account_id is not None for space in payload.spaces)
|
|
granted = sum(
|
|
(1 if space.owner_account_id is not None else 0)
|
|
+ sum(permission.account_id is not None for permission in space.permissions)
|
|
for space in payload.spaces
|
|
)
|
|
quarantined = sum(
|
|
(space.owner_account_id is None)
|
|
+ sum(permission.account_id is None for permission in space.permissions)
|
|
+ len(space.legacy_api_keys)
|
|
+ sum(
|
|
_classify_task(task) is not KnowledgeFSMigrationQuarantineDisposition.MIGRATABLE for task in space.tasks
|
|
)
|
|
+ len(space.orphan_resource_ids)
|
|
for space in payload.spaces
|
|
)
|
|
open_issues = sum(
|
|
(space.owner_account_id is None)
|
|
+ sum(permission.account_id is None for permission in space.permissions)
|
|
+ (space.external_access_enabled is None or space.visibility == "unknown")
|
|
for space in payload.spaces
|
|
)
|
|
return CutoverBackfillReport(
|
|
str(payload.tenant_id),
|
|
registered,
|
|
granted,
|
|
quarantined,
|
|
open_issues,
|
|
KnowledgeFSWorkspaceCutoverPhase.INVENTORY.value,
|
|
False,
|
|
)
|
|
|
|
def _backfill_space(
|
|
self,
|
|
session: Session,
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
space: LegacySpaceInventoryInput,
|
|
revision_watermark: KnowledgeFSCutoverRevisionWatermark,
|
|
) -> tuple[int, int, int]:
|
|
knowledge_space_id = str(space.knowledge_space_id)
|
|
quarantined = self._backfill_auxiliary_inventory(repository, ledger, space, knowledge_space_id)
|
|
if space.owner_account_id is None:
|
|
self._ensure_issue(
|
|
repository,
|
|
ledger,
|
|
issue_key=f"owner:{knowledge_space_id}",
|
|
kind=KnowledgeFSMigrationIssueKind.UNRESOLVED_SUBJECT,
|
|
resource_type="knowledge_space",
|
|
resource_id=knowledge_space_id,
|
|
details={"subject_id": space.owner_subject_id},
|
|
)
|
|
self._ensure_quarantine(
|
|
repository,
|
|
ledger,
|
|
source_kind=KnowledgeFSMigrationQuarantineKind.CONTROL_SPACE,
|
|
source_id=knowledge_space_id,
|
|
reason_code="UNRESOLVED_OWNER",
|
|
disposition=KnowledgeFSMigrationQuarantineDisposition.PENDING,
|
|
details={"subject_id": space.owner_subject_id},
|
|
)
|
|
return 0, 0, quarantined + 1
|
|
tenant_id = ledger.tenant_id
|
|
owner_account_id = str(space.owner_account_id)
|
|
existing = session.scalar(
|
|
sa.select(KnowledgeFSControlSpace).where(
|
|
KnowledgeFSControlSpace.tenant_id == tenant_id,
|
|
(KnowledgeFSControlSpace.knowledge_space_id == knowledge_space_id)
|
|
| (KnowledgeFSControlSpace.provisioning_key == space.provisioning_key),
|
|
)
|
|
)
|
|
registered = 0
|
|
if existing is not None and (
|
|
existing.knowledge_space_id != knowledge_space_id
|
|
or existing.provisioning_key != space.provisioning_key
|
|
or existing.owner_account_id != owner_account_id
|
|
):
|
|
self._ensure_issue(
|
|
repository,
|
|
ledger,
|
|
issue_key=f"registration:{knowledge_space_id}",
|
|
kind=KnowledgeFSMigrationIssueKind.REGISTRATION_CONFLICT,
|
|
resource_type="knowledge_space",
|
|
resource_id=knowledge_space_id,
|
|
details={"provisioning_key": space.provisioning_key},
|
|
)
|
|
return 0, 0, quarantined
|
|
visibility = _VISIBILITY_MAP.get(space.visibility, KnowledgeFSControlSpaceVisibility.ONLY_ME)
|
|
if existing is None:
|
|
existing = KnowledgeFSControlSpace(
|
|
tenant_id=tenant_id,
|
|
owner_account_id=owner_account_id,
|
|
provisioning_key=space.provisioning_key,
|
|
knowledge_space_id=knowledge_space_id,
|
|
knowledge_space_revision=space.knowledge_space_revision,
|
|
visibility=visibility,
|
|
state=KnowledgeFSControlSpaceState.ACTIVE,
|
|
lifecycle_operation_id=str(
|
|
uuid5(
|
|
NAMESPACE_URL,
|
|
f"dify-kfs-cutover-backfill:{tenant_id}:{knowledge_space_id}",
|
|
)
|
|
),
|
|
)
|
|
session.add(existing)
|
|
session.flush()
|
|
registered = 1
|
|
elif existing.knowledge_space_revision > space.knowledge_space_revision:
|
|
self._ensure_issue(
|
|
repository,
|
|
ledger,
|
|
issue_key=f"revision-drift:{knowledge_space_id}",
|
|
kind=KnowledgeFSMigrationIssueKind.REVISION_DRIFT,
|
|
resource_type="knowledge_space",
|
|
resource_id=knowledge_space_id,
|
|
details={
|
|
"existing_revision": existing.knowledge_space_revision,
|
|
"inventory_revision": space.knowledge_space_revision,
|
|
},
|
|
)
|
|
return registered, 0, quarantined
|
|
else:
|
|
if (
|
|
existing.knowledge_space_revision != space.knowledge_space_revision
|
|
or existing.visibility is not visibility
|
|
):
|
|
existing.resource_version += 1
|
|
existing.knowledge_space_revision = space.knowledge_space_revision
|
|
existing.visibility = visibility
|
|
revision = session.scalar(
|
|
sa.select(KnowledgeFSAuthorizationRevision).where(
|
|
KnowledgeFSAuthorizationRevision.tenant_id == tenant_id,
|
|
KnowledgeFSAuthorizationRevision.control_space_id == existing.id,
|
|
)
|
|
)
|
|
if revision is None:
|
|
revision = KnowledgeFSAuthorizationRevision(
|
|
tenant_id=tenant_id,
|
|
control_space_id=existing.id,
|
|
)
|
|
session.add(revision)
|
|
revision.membership_epoch = revision_watermark["membership_epoch"]
|
|
revision.space_acl_epoch = revision_watermark["space_acl_epoch"]
|
|
revision.external_access_epoch = revision_watermark["external_access_epoch"]
|
|
revision.content_policy_revision = revision_watermark["content_policy_revision"]
|
|
granted = self._ensure_permission(
|
|
session,
|
|
tenant_id=tenant_id,
|
|
control_space=existing,
|
|
account_id=owner_account_id,
|
|
role=KnowledgeFSControlSpacePermissionRole.OWNER,
|
|
)
|
|
for permission in space.permissions:
|
|
if permission.account_id is None:
|
|
continue
|
|
granted += self._ensure_permission(
|
|
session,
|
|
tenant_id=tenant_id,
|
|
control_space=existing,
|
|
account_id=str(permission.account_id),
|
|
role=KnowledgeFSControlSpacePermissionRole(permission.role),
|
|
)
|
|
external_enabled = space.external_access_enabled is True
|
|
external = session.scalar(
|
|
sa.select(KnowledgeFSExternalAccessPolicy).where(
|
|
KnowledgeFSExternalAccessPolicy.tenant_id == tenant_id,
|
|
KnowledgeFSExternalAccessPolicy.control_space_id == existing.id,
|
|
)
|
|
)
|
|
if external is None:
|
|
external = KnowledgeFSExternalAccessPolicy(
|
|
tenant_id=tenant_id,
|
|
control_space_id=existing.id,
|
|
)
|
|
session.add(external)
|
|
elif any(
|
|
(
|
|
external.service_api_enabled != external_enabled,
|
|
external.agent_enabled != external_enabled,
|
|
external.workflow_enabled != external_enabled,
|
|
external.mcp_enabled,
|
|
)
|
|
):
|
|
external.revision += 1
|
|
external.service_api_enabled = external_enabled
|
|
external.agent_enabled = external_enabled
|
|
external.workflow_enabled = external_enabled
|
|
external.mcp_enabled = False
|
|
return registered, granted, quarantined
|
|
|
|
def _backfill_auxiliary_inventory(
|
|
self,
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
space: LegacySpaceInventoryInput,
|
|
knowledge_space_id: str,
|
|
) -> int:
|
|
quarantined = 0
|
|
visibility = _VISIBILITY_MAP.get(space.visibility, KnowledgeFSControlSpaceVisibility.ONLY_ME)
|
|
for permission in space.permissions:
|
|
if permission.account_id is not None:
|
|
continue
|
|
self._ensure_issue(
|
|
repository,
|
|
ledger,
|
|
issue_key=f"subject:{knowledge_space_id}:{permission.subject_id}",
|
|
kind=KnowledgeFSMigrationIssueKind.UNRESOLVED_SUBJECT,
|
|
resource_type="subject",
|
|
resource_id=permission.subject_id,
|
|
details={"knowledge_space_id": knowledge_space_id},
|
|
)
|
|
self._ensure_quarantine(
|
|
repository,
|
|
ledger,
|
|
source_kind=KnowledgeFSMigrationQuarantineKind.SUBJECT,
|
|
source_id=f"{knowledge_space_id}:{permission.subject_id}",
|
|
reason_code="UNRESOLVED_SUBJECT",
|
|
disposition=KnowledgeFSMigrationQuarantineDisposition.ISOLATE,
|
|
details={"role": permission.role},
|
|
)
|
|
quarantined += 1
|
|
if space.external_access_enabled is None or space.visibility == "unknown":
|
|
self._ensure_issue(
|
|
repository,
|
|
ledger,
|
|
issue_key=f"unknown-access:{knowledge_space_id}",
|
|
kind=KnowledgeFSMigrationIssueKind.UNKNOWN_EXTERNAL_ACCESS,
|
|
resource_type="knowledge_space",
|
|
resource_id=knowledge_space_id,
|
|
details={"effective_external_access": False, "effective_visibility": visibility.value},
|
|
)
|
|
for legacy_key in space.legacy_api_keys:
|
|
self._ensure_quarantine(
|
|
repository,
|
|
ledger,
|
|
source_kind=KnowledgeFSMigrationQuarantineKind.LEGACY_API_KEY,
|
|
source_id=legacy_key.key_id,
|
|
reason_code="ROTATION_REQUIRED",
|
|
disposition=KnowledgeFSMigrationQuarantineDisposition.ROTATE_CREDENTIAL,
|
|
details={
|
|
"knowledge_space_id": knowledge_space_id,
|
|
"prefix": legacy_key.prefix,
|
|
"last4": legacy_key.last4,
|
|
},
|
|
)
|
|
quarantined += 1
|
|
for task in space.tasks:
|
|
disposition = _classify_task(task)
|
|
self._ensure_quarantine(
|
|
repository,
|
|
ledger,
|
|
source_kind=KnowledgeFSMigrationQuarantineKind.TASK,
|
|
source_id=task.task_id,
|
|
reason_code="LEGACY_TASK_CLASSIFICATION",
|
|
disposition=disposition,
|
|
details={
|
|
"knowledge_space_id": knowledge_space_id,
|
|
"state": task.state,
|
|
"subject_id": task.subject_id,
|
|
"account_id": str(task.account_id) if task.account_id else None,
|
|
"expires_at": task.expires_at.isoformat() if task.expires_at else None,
|
|
},
|
|
)
|
|
quarantined += disposition is not KnowledgeFSMigrationQuarantineDisposition.MIGRATABLE
|
|
for orphan_id in space.orphan_resource_ids:
|
|
self._ensure_quarantine(
|
|
repository,
|
|
ledger,
|
|
source_kind=KnowledgeFSMigrationQuarantineKind.ORPHAN_RESOURCE,
|
|
source_id=orphan_id,
|
|
reason_code="ORPHAN_RESOURCE",
|
|
disposition=KnowledgeFSMigrationQuarantineDisposition.PENDING,
|
|
details={"knowledge_space_id": knowledge_space_id},
|
|
)
|
|
quarantined += 1
|
|
return quarantined
|
|
|
|
@staticmethod
|
|
def _ensure_permission(
|
|
session: Session,
|
|
*,
|
|
tenant_id: str,
|
|
control_space: KnowledgeFSControlSpace,
|
|
account_id: str,
|
|
role: KnowledgeFSControlSpacePermissionRole,
|
|
) -> int:
|
|
existing = session.scalar(
|
|
sa.select(KnowledgeFSControlSpacePermission).where(
|
|
KnowledgeFSControlSpacePermission.tenant_id == tenant_id,
|
|
KnowledgeFSControlSpacePermission.control_space_id == control_space.id,
|
|
KnowledgeFSControlSpacePermission.account_id == account_id,
|
|
)
|
|
)
|
|
if existing is not None:
|
|
if existing.role is not role:
|
|
existing.role = role
|
|
existing.revision += 1
|
|
return 0
|
|
session.add(
|
|
KnowledgeFSControlSpacePermission(
|
|
tenant_id=tenant_id,
|
|
control_space_id=control_space.id,
|
|
account_id=account_id,
|
|
role=role,
|
|
granted_by_account_id=control_space.owner_account_id,
|
|
)
|
|
)
|
|
return 1
|
|
|
|
@staticmethod
|
|
def _ensure_issue(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
*,
|
|
issue_key: str,
|
|
kind: KnowledgeFSMigrationIssueKind,
|
|
resource_type: str,
|
|
resource_id: str,
|
|
details: dict[str, object],
|
|
) -> None:
|
|
if repository.get_issue(tenant_id=ledger.tenant_id, ledger_id=ledger.id, issue_key=issue_key) is None:
|
|
repository.add_issue(
|
|
KnowledgeFSMigrationIssue(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
issue_key=issue_key,
|
|
kind=kind,
|
|
resource_type=resource_type,
|
|
resource_id=resource_id,
|
|
details=details,
|
|
)
|
|
)
|
|
|
|
@staticmethod
|
|
def _ensure_quarantine(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
*,
|
|
source_kind: KnowledgeFSMigrationQuarantineKind,
|
|
source_id: str,
|
|
reason_code: str,
|
|
disposition: KnowledgeFSMigrationQuarantineDisposition,
|
|
details: dict[str, object],
|
|
) -> None:
|
|
if (
|
|
repository.get_quarantine(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
source_kind=source_kind,
|
|
source_id=source_id,
|
|
)
|
|
is None
|
|
):
|
|
repository.add_quarantine(
|
|
KnowledgeFSMigrationQuarantine(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
source_kind=source_kind,
|
|
source_id=source_id,
|
|
reason_code=reason_code,
|
|
disposition=disposition,
|
|
details=details,
|
|
)
|
|
)
|
|
|
|
def _validate_shadow_observation(
|
|
self,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
item: ShadowAuthorizationObservationInput,
|
|
) -> None:
|
|
if item.producer not in self._trusted_shadow_producers:
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow observation producer is not trusted")
|
|
if ledger.shadow_started_at is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow start evidence is missing")
|
|
observed_at = _naive_utc(item.observed_at)
|
|
if observed_at > _naive_utc(self._clock()):
|
|
raise KnowledgeFSCutoverConflictError("Shadow observation timestamp cannot be in the future")
|
|
if observed_at < ledger.shadow_started_at:
|
|
raise KnowledgeFSCutoverConflictError("Shadow observation predates shadow start")
|
|
required_revision = ledger.final_revision_watermark or ledger.source_revision_watermark
|
|
if not _watermark_at_least(item.observed_revision.to_record(), required_revision):
|
|
raise KnowledgeFSCutoverConflictError("Shadow observation revision is stale")
|
|
|
|
@staticmethod
|
|
def _record_shadow_observation(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
item: ShadowAuthorizationObservationInput,
|
|
decision: KnowledgeFSShadowAuthorizationDecision,
|
|
*,
|
|
apply: bool,
|
|
) -> bool:
|
|
evidence_digest = _shadow_observation_digest(item)
|
|
replay = repository.get_shadow_observation(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=item.diff_key,
|
|
evidence_digest=evidence_digest,
|
|
)
|
|
if replay is not None:
|
|
return True
|
|
if ledger.shadow_completed_at is not None:
|
|
raise KnowledgeFSCutoverConflictError("Shadow evidence is already complete")
|
|
existing = repository.get_shadow_diff(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=item.diff_key,
|
|
)
|
|
dify_allowed = False if item.legacy_allowed is None else item.dify_allowed
|
|
if existing is not None:
|
|
safe_reevaluation = (
|
|
existing.decision
|
|
in {
|
|
KnowledgeFSShadowAuthorizationDecision.EXPANDED,
|
|
KnowledgeFSShadowAuthorizationDecision.UNKNOWN,
|
|
}
|
|
and decision
|
|
in {
|
|
KnowledgeFSShadowAuthorizationDecision.MATCH,
|
|
KnowledgeFSShadowAuthorizationDecision.TIGHTENED,
|
|
}
|
|
and _watermark_at_least(item.observed_revision.to_record(), existing.observed_revision)
|
|
)
|
|
if not safe_reevaluation:
|
|
raise KnowledgeFSCutoverConflictError("Shadow diff key was reused with different evidence")
|
|
if not apply:
|
|
return False
|
|
repository.add_shadow_observation(
|
|
_shadow_observation_model(ledger=ledger, item=item, decision=decision, digest=evidence_digest)
|
|
)
|
|
changed = repository.reevaluate_shadow_diff(
|
|
KnowledgeFSShadowDiffCASUpdate(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=item.diff_key,
|
|
expected_row_version=existing.row_version,
|
|
control_space_id=str(item.control_space_id) if item.control_space_id else None,
|
|
principal=item.principal,
|
|
legacy_allowed=item.legacy_allowed,
|
|
dify_allowed=dify_allowed,
|
|
decision=decision,
|
|
reason=item.reason,
|
|
observed_revision=item.observed_revision.to_record(),
|
|
status=KnowledgeFSMigrationIssueStatus.RESOLVED,
|
|
current_evidence_digest=evidence_digest,
|
|
last_observed_at=_naive_utc(item.observed_at),
|
|
)
|
|
)
|
|
if not changed:
|
|
raise KnowledgeFSCutoverConflictError("Shadow diff changed during re-evaluation")
|
|
return False
|
|
if not apply:
|
|
return False
|
|
repository.add_shadow_observation(
|
|
_shadow_observation_model(ledger=ledger, item=item, decision=decision, digest=evidence_digest)
|
|
)
|
|
repository.add_shadow_diff(
|
|
KnowledgeFSShadowAuthorizationDiff(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=item.diff_key,
|
|
control_space_id=str(item.control_space_id) if item.control_space_id else None,
|
|
principal=item.principal,
|
|
legacy_allowed=item.legacy_allowed,
|
|
dify_allowed=dify_allowed,
|
|
decision=decision,
|
|
reason=item.reason,
|
|
observed_revision=item.observed_revision.to_record(),
|
|
current_evidence_digest=evidence_digest,
|
|
last_observed_at=_naive_utc(item.observed_at),
|
|
status=(
|
|
KnowledgeFSMigrationIssueStatus.RESOLVED
|
|
if decision is KnowledgeFSShadowAuthorizationDecision.MATCH
|
|
else KnowledgeFSMigrationIssueStatus.OPEN
|
|
),
|
|
)
|
|
)
|
|
return False
|
|
|
|
def _assert_shadow_completion_replay(
|
|
self,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
payload: ShadowCompletionInput,
|
|
report: ShadowCompletionReport,
|
|
) -> None:
|
|
evidence = (
|
|
cast(dict[str, object], payload.traffic_zero_evidence)
|
|
if payload.traffic_zero_evidence is not None
|
|
else None
|
|
)
|
|
if (
|
|
ledger.shadow_completed_at != _naive_utc(payload.completed_at)
|
|
or ledger.shadow_evidence_digest != report.evidence_digest
|
|
or ledger.shadow_observation_count != report.observation_count
|
|
or ledger.shadow_traffic_zero is not payload.traffic_zero
|
|
or ledger.shadow_traffic_zero_evidence != evidence
|
|
or ledger.shadow_latest_observed_revision != report.latest_observed_revision
|
|
or ledger.shadow_producer != payload.producer
|
|
or ledger.shadow_completed_by_operator != payload.completed_by_operator
|
|
or ledger.shadow_completed_by_account_id != str(payload.completed_by_account_id)
|
|
or ledger.shadow_window_started_at
|
|
!= (_naive_utc(payload.window_started_at) if payload.window_started_at is not None else None)
|
|
or ledger.shadow_window_ended_at
|
|
!= (_naive_utc(payload.window_ended_at) if payload.window_ended_at is not None else None)
|
|
):
|
|
raise KnowledgeFSCutoverConflictError("Shadow was already completed with different evidence")
|
|
|
|
@staticmethod
|
|
def _shadow_decision(item: ShadowAuthorizationObservationInput) -> KnowledgeFSShadowAuthorizationDecision:
|
|
if item.legacy_allowed is None:
|
|
return KnowledgeFSShadowAuthorizationDecision.UNKNOWN
|
|
if item.legacy_allowed == item.dify_allowed:
|
|
return KnowledgeFSShadowAuthorizationDecision.MATCH
|
|
if item.legacy_allowed and not item.dify_allowed:
|
|
return KnowledgeFSShadowAuthorizationDecision.TIGHTENED
|
|
return KnowledgeFSShadowAuthorizationDecision.EXPANDED
|
|
|
|
@staticmethod
|
|
def _require_ledger(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository, tenant_id: str
|
|
) -> KnowledgeFSWorkspaceCutoverLedger:
|
|
ledger = repository.get_ledger(tenant_id=tenant_id)
|
|
if ledger is None:
|
|
raise KnowledgeFSCutoverNotFoundError("Workspace cutover ledger was not found")
|
|
return ledger
|
|
|
|
@staticmethod
|
|
def _assert_phase(
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
expected_phase: KnowledgeFSWorkspaceCutoverPhase,
|
|
expected_cas_version: int,
|
|
) -> None:
|
|
if ledger.phase is not expected_phase:
|
|
raise KnowledgeFSCutoverConflictError(
|
|
f"Expected cutover phase {expected_phase.value}, found {ledger.phase.value}"
|
|
)
|
|
if ledger.cas_version != expected_cas_version:
|
|
raise KnowledgeFSCutoverConflictError("Cutover ledger version changed")
|
|
|
|
@staticmethod
|
|
def _assert_no_open_gates(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
) -> None:
|
|
if repository.count_open_issues(tenant_id=ledger.tenant_id, ledger_id=ledger.id) > 0:
|
|
raise KnowledgeFSCutoverGateBlockedError("Open migration issues block phase advancement")
|
|
if repository.count_unapproved_shadow_diffs(tenant_id=ledger.tenant_id, ledger_id=ledger.id) > 0:
|
|
raise KnowledgeFSCutoverGateBlockedError("Unapproved shadow differences block phase advancement")
|
|
|
|
def _assert_shadow_evidence_complete(
|
|
self,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
*,
|
|
required_revision: KnowledgeFSCutoverRevisionWatermark | None = None,
|
|
) -> None:
|
|
if (
|
|
ledger.shadow_started_at is None
|
|
or ledger.shadow_completed_at is None
|
|
or ledger.shadow_evidence_digest is None
|
|
or ledger.shadow_producer not in self._trusted_shadow_producers
|
|
or ledger.shadow_completed_by_operator not in self._trusted_shadow_operators
|
|
or ledger.shadow_completed_by_account_id is None
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow evidence is not complete and trusted")
|
|
if ledger.shadow_traffic_zero:
|
|
if ledger.shadow_observation_count != 0 or ledger.shadow_traffic_zero_evidence is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Traffic-zero shadow evidence is incomplete")
|
|
return
|
|
if (
|
|
ledger.shadow_observation_count <= 0
|
|
or ledger.shadow_window_started_at is None
|
|
or ledger.shadow_window_ended_at is None
|
|
or ledger.shadow_latest_observed_revision is None
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError("Shadow observation evidence is incomplete")
|
|
applicable_revision = required_revision or ledger.final_revision_watermark or ledger.source_revision_watermark
|
|
if not _watermark_at_least(ledger.shadow_latest_observed_revision, applicable_revision):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"Shadow latest observed revision is below the applicable cutover watermark"
|
|
)
|
|
|
|
def _assert_cutover_ready(
|
|
self,
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
expected_cas_version: int,
|
|
) -> None:
|
|
self._assert_phase(ledger, KnowledgeFSWorkspaceCutoverPhase.FROZEN, expected_cas_version)
|
|
# Final delta is applied only after freeze. Revalidation therefore uses the
|
|
# source watermark that was applicable to the frozen shadow gate.
|
|
self._assert_shadow_evidence_complete(ledger, required_revision=ledger.source_revision_watermark)
|
|
self._assert_no_open_gates(repository, ledger)
|
|
self._assert_no_unresolved_quarantine(repository, ledger)
|
|
if ledger.freeze_at is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Workspace freeze evidence is missing")
|
|
_assert_remote_freeze_evidence(ledger)
|
|
if not ledger.legacy_dependency_ready or ledger.legacy_dependency_checked_at is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Legacy snapshot/FK dependency dashboard is not ready")
|
|
if (
|
|
ledger.final_revision_watermark is None
|
|
or ledger.final_revision_watermark != ledger.applied_revision_watermark
|
|
or ledger.final_task_watermark is None
|
|
or ledger.final_task_watermark != ledger.applied_task_watermark
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError("Final delta watermarks are not fully applied")
|
|
|
|
@staticmethod
|
|
def _require_activation_anchor(
|
|
session: Session,
|
|
*,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
control_space_id: str | None = None,
|
|
) -> _ActivationAnchor:
|
|
tenant_id = ledger.tenant_id
|
|
statement = sa.select(KnowledgeFSControlSpace).where(
|
|
KnowledgeFSControlSpace.tenant_id == tenant_id,
|
|
KnowledgeFSControlSpace.state == KnowledgeFSControlSpaceState.ACTIVE,
|
|
KnowledgeFSControlSpace.knowledge_space_id.is_not(None),
|
|
)
|
|
if control_space_id is not None:
|
|
statement = statement.where(KnowledgeFSControlSpace.id == control_space_id)
|
|
control_space = session.scalar(
|
|
statement.order_by(KnowledgeFSControlSpace.created_at, KnowledgeFSControlSpace.id).limit(1)
|
|
)
|
|
if control_space is not None:
|
|
return _ActivationAnchor(control_space.id, False)
|
|
|
|
greenfield_anchor = _greenfield_activation_anchor_id(tenant_id)
|
|
greenfield_provisioning_key = _greenfield_activation_provisioning_key(tenant_id)
|
|
local_control_space_count = session.scalar(
|
|
sa.select(sa.func.count())
|
|
.select_from(KnowledgeFSControlSpace)
|
|
.where(KnowledgeFSControlSpace.tenant_id == tenant_id)
|
|
)
|
|
if not _is_zero_space_cutover_ledger(ledger) or control_space_id not in {None, greenfield_anchor}:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"A tenant-owned active control-space is required for activation audit"
|
|
)
|
|
|
|
persisted_anchor = session.scalar(
|
|
sa.select(KnowledgeFSControlSpace).where(
|
|
KnowledgeFSControlSpace.tenant_id == tenant_id,
|
|
KnowledgeFSControlSpace.id == greenfield_anchor,
|
|
)
|
|
)
|
|
if persisted_anchor is not None:
|
|
if (
|
|
local_control_space_count == 1
|
|
and persisted_anchor.state is KnowledgeFSControlSpaceState.DELETED
|
|
and persisted_anchor.knowledge_space_id is None
|
|
and persisted_anchor.provisioning_key == greenfield_provisioning_key
|
|
and persisted_anchor.owner_account_id == ledger.shadow_completed_by_account_id
|
|
and persisted_anchor.lifecycle_operation_id == greenfield_anchor
|
|
):
|
|
return _ActivationAnchor(greenfield_anchor, True)
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"KnowledgeFS greenfield activation audit anchor conflicts with local control-space state"
|
|
)
|
|
|
|
if local_control_space_count == 0 and ledger.shadow_completed_by_account_id is not None:
|
|
persisted_anchor = KnowledgeFSControlSpace(
|
|
tenant_id=tenant_id,
|
|
owner_account_id=ledger.shadow_completed_by_account_id,
|
|
provisioning_key=greenfield_provisioning_key,
|
|
state=KnowledgeFSControlSpaceState.DELETED,
|
|
lifecycle_operation_id=greenfield_anchor,
|
|
)
|
|
persisted_anchor.id = greenfield_anchor
|
|
session.add(persisted_anchor)
|
|
session.flush()
|
|
return _ActivationAnchor(greenfield_anchor, True)
|
|
raise KnowledgeFSCutoverGateBlockedError("A tenant-owned active control-space is required for activation audit")
|
|
|
|
@staticmethod
|
|
def _assert_remote_greenfield_namespace_empty(
|
|
remote: KnowledgeFSLifecycleRemotePort,
|
|
*,
|
|
tenant_id: str,
|
|
control_space_id: str,
|
|
) -> None:
|
|
try:
|
|
spaces = remote.list_spaces(
|
|
namespace_id=tenant_id,
|
|
control_space_id=control_space_id,
|
|
)
|
|
except KnowledgeFSLifecycleRemoteError as exc:
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"KnowledgeFS greenfield remote namespace could not be verified"
|
|
) from exc
|
|
if spaces:
|
|
raise KnowledgeFSCutoverGateBlockedError("KnowledgeFS greenfield remote namespace is not empty")
|
|
|
|
@staticmethod
|
|
def _assert_no_unresolved_quarantine(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
) -> None:
|
|
if (
|
|
repository.count_unresolved_quarantine(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
source_kinds=_CUTOVER_QUARANTINE_KINDS,
|
|
)
|
|
> 0
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError("Unresolved migration quarantine blocks freeze and cutover")
|
|
|
|
@staticmethod
|
|
def _cas(
|
|
repository: SQLAlchemyKnowledgeFSCutoverRepository,
|
|
update: KnowledgeFSCutoverCASUpdate,
|
|
) -> None:
|
|
if not repository.compare_and_set(update):
|
|
raise KnowledgeFSCutoverConflictError("Cutover ledger changed during transition")
|
|
|
|
|
|
def _classify_task(task: LegacyTaskInventoryInput) -> KnowledgeFSMigrationQuarantineDisposition:
|
|
if task.subject_id is None or task.subject_id.startswith("dify-workspace:") or task.account_id is None:
|
|
return KnowledgeFSMigrationQuarantineDisposition.ISOLATE
|
|
if task.state in {"running", "paused"}:
|
|
return KnowledgeFSMigrationQuarantineDisposition.WAIT_FOR_COMPLETION
|
|
if task.state == "queued":
|
|
return KnowledgeFSMigrationQuarantineDisposition.MIGRATABLE
|
|
if task.state in {"completed", "failed", "canceled"}:
|
|
return KnowledgeFSMigrationQuarantineDisposition.CANCEL
|
|
return KnowledgeFSMigrationQuarantineDisposition.CANCEL
|
|
|
|
|
|
def _shadow_observation_digest(item: ShadowAuthorizationObservationInput) -> str:
|
|
canonical = json.dumps(
|
|
item.model_dump(mode="json"),
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode()
|
|
return f"sha256:{hashlib.sha256(canonical).hexdigest()}"
|
|
|
|
|
|
def _shadow_observation_model(
|
|
*,
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
item: ShadowAuthorizationObservationInput,
|
|
decision: KnowledgeFSShadowAuthorizationDecision,
|
|
digest: str,
|
|
) -> KnowledgeFSShadowAuthorizationObservation:
|
|
return KnowledgeFSShadowAuthorizationObservation(
|
|
tenant_id=ledger.tenant_id,
|
|
ledger_id=ledger.id,
|
|
diff_key=item.diff_key,
|
|
producer=item.producer,
|
|
control_space_id=str(item.control_space_id) if item.control_space_id else None,
|
|
principal=item.principal,
|
|
legacy_allowed=item.legacy_allowed,
|
|
dify_allowed=False if item.legacy_allowed is None else item.dify_allowed,
|
|
decision=decision,
|
|
reason=item.reason,
|
|
observed_revision=item.observed_revision.to_record(),
|
|
observed_at=_naive_utc(item.observed_at),
|
|
evidence_digest=digest,
|
|
)
|
|
|
|
|
|
def _shadow_observation_set_digest(
|
|
observations: Iterable[KnowledgeFSShadowAuthorizationObservation],
|
|
) -> str:
|
|
canonical = "\n".join(sorted(item.evidence_digest for item in observations)).encode()
|
|
return f"sha256:{hashlib.sha256(canonical).hexdigest()}"
|
|
|
|
|
|
def _shadow_traffic_zero_digest(payload: ShadowCompletionInput) -> str:
|
|
canonical = json.dumps(
|
|
{
|
|
"completed_at": payload.completed_at.isoformat(),
|
|
"completed_by_account_id": str(payload.completed_by_account_id),
|
|
"completed_by_operator": payload.completed_by_operator,
|
|
"producer": payload.producer,
|
|
"traffic_zero_evidence": payload.traffic_zero_evidence,
|
|
},
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode()
|
|
return f"sha256:{hashlib.sha256(canonical).hexdigest()}"
|
|
|
|
|
|
def _greenfield_activation_anchor_id(tenant_id: str) -> str:
|
|
return str(uuid5(NAMESPACE_URL, f"dify-kfs-greenfield-activation:{tenant_id}"))
|
|
|
|
|
|
def _greenfield_activation_provisioning_key(tenant_id: str) -> str:
|
|
return f"dify-kfs-greenfield-activation:{tenant_id}"
|
|
|
|
|
|
def _is_zero_space_cutover_ledger(ledger: KnowledgeFSWorkspaceCutoverLedger) -> bool:
|
|
return (
|
|
ledger.source_revision_watermark == _ZERO_REVISION_WATERMARK
|
|
and ledger.applied_revision_watermark == _ZERO_REVISION_WATERMARK
|
|
and ledger.source_task_watermark == 0
|
|
and ledger.applied_task_watermark == 0
|
|
and (ledger.final_revision_watermark is None or ledger.final_revision_watermark == _ZERO_REVISION_WATERMARK)
|
|
and ledger.final_task_watermark in {None, 0}
|
|
)
|
|
|
|
|
|
def _has_complete_greenfield_final_delta(ledger: KnowledgeFSWorkspaceCutoverLedger) -> bool:
|
|
return (
|
|
ledger.final_revision_watermark == _ZERO_REVISION_WATERMARK
|
|
and ledger.applied_revision_watermark == _ZERO_REVISION_WATERMARK
|
|
and ledger.final_task_watermark == 0
|
|
and ledger.applied_task_watermark == 0
|
|
)
|
|
|
|
|
|
def _has_complete_product_cutover(ledger: KnowledgeFSWorkspaceCutoverLedger) -> bool:
|
|
return bool(
|
|
ledger.phase
|
|
in {
|
|
KnowledgeFSWorkspaceCutoverPhase.CUTOVER,
|
|
KnowledgeFSWorkspaceCutoverPhase.OBSERVING,
|
|
KnowledgeFSWorkspaceCutoverPhase.READY_FOR_CLEANUP,
|
|
}
|
|
and ledger.cutover_at is not None
|
|
and ledger.rolled_back_at is None
|
|
and ledger.product_routes_enabled
|
|
and ledger.capability_v2_enabled
|
|
and ledger.integrated_mode_enabled
|
|
and ledger.legacy_acl_read_only
|
|
)
|
|
|
|
|
|
def _revision_digest(revision: KnowledgeFSCutoverRevisionWatermark) -> str:
|
|
canonical = json.dumps(revision, sort_keys=True, separators=(",", ":")).encode()
|
|
return f"sha256:{hashlib.sha256(canonical).hexdigest()}"
|
|
|
|
|
|
def _canonical_freeze_id(
|
|
*,
|
|
freeze_revision: int,
|
|
namespace_id: str,
|
|
source_revision_digest: str,
|
|
source_task_watermark: int,
|
|
) -> str:
|
|
canonical = json.dumps(
|
|
{
|
|
"freezeRevision": freeze_revision,
|
|
"namespaceId": namespace_id,
|
|
"sourceRevisionDigest": source_revision_digest,
|
|
"sourceTaskWatermark": source_task_watermark,
|
|
},
|
|
ensure_ascii=False,
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode()
|
|
return f"sha256:{hashlib.sha256(canonical).hexdigest()}"
|
|
|
|
|
|
def _freeze_request(
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
*,
|
|
control_space_id: str,
|
|
) -> KnowledgeFSDifyIntegrationFreezeRequest:
|
|
freeze_revision = max(ledger.cas_version + 1, (ledger.remote_freeze_revision or 0) + 1)
|
|
source_revision_digest = _revision_digest(ledger.source_revision_watermark)
|
|
freeze_id = _canonical_freeze_id(
|
|
freeze_revision=freeze_revision,
|
|
namespace_id=ledger.tenant_id,
|
|
source_revision_digest=source_revision_digest,
|
|
source_task_watermark=ledger.source_task_watermark,
|
|
)
|
|
return KnowledgeFSDifyIntegrationFreezeRequest(
|
|
namespace_id=ledger.tenant_id,
|
|
control_space_id=control_space_id,
|
|
freeze_id=freeze_id,
|
|
freeze_revision=freeze_revision,
|
|
source_revision_digest=source_revision_digest,
|
|
source_task_watermark=ledger.source_task_watermark,
|
|
)
|
|
|
|
|
|
def _assert_exact_freeze_ack(
|
|
request: KnowledgeFSDifyIntegrationFreezeRequest,
|
|
ack: KnowledgeFSDifyIntegrationFreezeAck,
|
|
) -> None:
|
|
if (
|
|
ack.namespace_id != request.namespace_id
|
|
or ack.freeze_id != request.freeze_id
|
|
or ack.freeze_revision != request.freeze_revision
|
|
or ack.source_revision_digest != request.source_revision_digest
|
|
or ack.source_task_watermark != request.source_task_watermark
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"KnowledgeFS remote freeze acknowledgement did not exactly match the request"
|
|
)
|
|
|
|
|
|
def _assert_remote_freeze_evidence(ledger: KnowledgeFSWorkspaceCutoverLedger) -> None:
|
|
fields = (
|
|
ledger.remote_freeze_id,
|
|
ledger.remote_freeze_revision,
|
|
ledger.remote_freeze_digest,
|
|
ledger.remote_freeze_task_watermark,
|
|
ledger.remote_freeze_control_space_id,
|
|
ledger.remote_freeze_frozen_at,
|
|
ledger.remote_freeze_updated_at,
|
|
ledger.remote_freeze_acknowledged_at,
|
|
ledger.remote_freeze_applied,
|
|
ledger.remote_freeze_replayed,
|
|
)
|
|
if any(value is None for value in fields):
|
|
raise KnowledgeFSCutoverGateBlockedError("Remote Workspace freeze evidence is incomplete")
|
|
freeze_revision = cast(int, ledger.remote_freeze_revision)
|
|
digest = cast(str, ledger.remote_freeze_digest)
|
|
task_watermark = cast(int, ledger.remote_freeze_task_watermark)
|
|
if (
|
|
digest != _revision_digest(ledger.source_revision_watermark)
|
|
or task_watermark != ledger.source_task_watermark
|
|
or ledger.remote_freeze_id
|
|
!= _canonical_freeze_id(
|
|
freeze_revision=freeze_revision,
|
|
namespace_id=ledger.tenant_id,
|
|
source_revision_digest=digest,
|
|
source_task_watermark=task_watermark,
|
|
)
|
|
or ledger.remote_freeze_applied is ledger.remote_freeze_replayed
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError("Remote Workspace freeze evidence is inconsistent")
|
|
|
|
|
|
def _activation_request(
|
|
ledger: KnowledgeFSWorkspaceCutoverLedger,
|
|
*,
|
|
control_space_id: str,
|
|
) -> KnowledgeFSDifyIntegrationActivationRequest:
|
|
final_revision = ledger.final_revision_watermark
|
|
final_task_watermark = ledger.final_task_watermark
|
|
if final_revision is None or final_task_watermark is None:
|
|
raise KnowledgeFSCutoverGateBlockedError("Final delta evidence is missing")
|
|
digest = _activation_source_revision_digest(final_revision, final_task_watermark)
|
|
activation_revision = max(
|
|
ledger.cas_version + 1,
|
|
(ledger.remote_activation_revision or 0) + 1,
|
|
)
|
|
activation_id = _canonical_activation_id(
|
|
activation_revision=activation_revision,
|
|
namespace_id=ledger.tenant_id,
|
|
source_revision_digest=digest,
|
|
)
|
|
return KnowledgeFSDifyIntegrationActivationRequest(
|
|
namespace_id=ledger.tenant_id,
|
|
control_space_id=control_space_id,
|
|
activation_id=activation_id,
|
|
activation_revision=activation_revision,
|
|
source_revision_digest=digest,
|
|
)
|
|
|
|
|
|
def _assert_exact_activation_ack(
|
|
request: KnowledgeFSDifyIntegrationActivationRequest,
|
|
ack: KnowledgeFSDifyIntegrationActivationAck,
|
|
) -> None:
|
|
if (
|
|
ack.namespace_id != request.namespace_id
|
|
or ack.activation_id != request.activation_id
|
|
or ack.activation_revision != request.activation_revision
|
|
or ack.source_revision_digest != request.source_revision_digest
|
|
):
|
|
raise KnowledgeFSCutoverGateBlockedError(
|
|
"KnowledgeFS remote activation acknowledgement did not exactly match the request"
|
|
)
|
|
|
|
|
|
def _activation_source_revision_digest(
|
|
final_revision: KnowledgeFSCutoverRevisionWatermark,
|
|
final_task_watermark: int,
|
|
) -> str:
|
|
canonical = json.dumps(
|
|
{
|
|
"final_revision_watermark": final_revision,
|
|
"final_task_watermark": final_task_watermark,
|
|
},
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode()
|
|
return f"sha256:{hashlib.sha256(canonical).hexdigest()}"
|
|
|
|
|
|
def _canonical_activation_id(
|
|
*,
|
|
activation_revision: int,
|
|
namespace_id: str,
|
|
source_revision_digest: str,
|
|
) -> str:
|
|
activation_envelope = json.dumps(
|
|
{
|
|
"activationRevision": activation_revision,
|
|
"namespaceId": namespace_id,
|
|
"sourceRevisionDigest": source_revision_digest,
|
|
},
|
|
ensure_ascii=False,
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode()
|
|
return f"sha256:{hashlib.sha256(activation_envelope).hexdigest()}"
|
|
|
|
|
|
def knowledge_fs_remote_freeze_evidence_consistent(ledger: KnowledgeFSWorkspaceCutoverLedger) -> bool:
|
|
"""Return whether persisted KFS freeze evidence exactly binds the source watermarks."""
|
|
|
|
try:
|
|
_assert_remote_freeze_evidence(ledger)
|
|
except KnowledgeFSCutoverGateBlockedError:
|
|
return False
|
|
return True
|
|
|
|
|
|
def knowledge_fs_remote_activation_evidence_consistent(ledger: KnowledgeFSWorkspaceCutoverLedger) -> bool:
|
|
"""Return whether persisted KFS activation evidence exactly binds the final delta."""
|
|
|
|
fields = (
|
|
ledger.remote_activation_id,
|
|
ledger.remote_activation_revision,
|
|
ledger.remote_activation_digest,
|
|
ledger.remote_activation_control_space_id,
|
|
ledger.remote_activation_activated_at,
|
|
ledger.remote_activation_updated_at,
|
|
ledger.remote_activation_acknowledged_at,
|
|
ledger.remote_activation_applied,
|
|
ledger.remote_activation_replayed,
|
|
)
|
|
if any(value is None for value in fields):
|
|
return False
|
|
final_revision = ledger.final_revision_watermark
|
|
final_task_watermark = ledger.final_task_watermark
|
|
if final_revision is None or final_task_watermark is None:
|
|
return False
|
|
activation_revision = cast(int, ledger.remote_activation_revision)
|
|
digest = _activation_source_revision_digest(final_revision, final_task_watermark)
|
|
return (
|
|
ledger.remote_activation_digest == digest
|
|
and ledger.remote_activation_id
|
|
== _canonical_activation_id(
|
|
activation_revision=activation_revision,
|
|
namespace_id=ledger.tenant_id,
|
|
source_revision_digest=digest,
|
|
)
|
|
and ledger.remote_activation_control_space_id == ledger.remote_freeze_control_space_id
|
|
and ledger.remote_activation_applied is not ledger.remote_activation_replayed
|
|
)
|
|
|
|
|
|
def _latest_shadow_revision(
|
|
observations: Iterable[KnowledgeFSShadowAuthorizationObservation],
|
|
) -> KnowledgeFSCutoverRevisionWatermark:
|
|
values = tuple(observations)
|
|
return {
|
|
"membership_epoch": max(item.observed_revision["membership_epoch"] for item in values),
|
|
"space_acl_epoch": max(item.observed_revision["space_acl_epoch"] for item in values),
|
|
"external_access_epoch": max(item.observed_revision["external_access_epoch"] for item in values),
|
|
"content_policy_revision": max(item.observed_revision["content_policy_revision"] for item in values),
|
|
}
|
|
|
|
|
|
def _watermark_at_least(
|
|
candidate: KnowledgeFSCutoverRevisionWatermark,
|
|
baseline: KnowledgeFSCutoverRevisionWatermark,
|
|
) -> bool:
|
|
return (
|
|
candidate["membership_epoch"] >= baseline["membership_epoch"]
|
|
and candidate["space_acl_epoch"] >= baseline["space_acl_epoch"]
|
|
and candidate["external_access_epoch"] >= baseline["external_access_epoch"]
|
|
and candidate["content_policy_revision"] >= baseline["content_policy_revision"]
|
|
)
|
|
|
|
|
|
def _naive_utc(value: datetime) -> datetime:
|
|
if value.tzinfo is None:
|
|
return value
|
|
return value.astimezone(UTC).replace(tzinfo=None)
|
|
|
|
|
|
def _iso(value: datetime | None) -> str | None:
|
|
return value.replace(tzinfo=UTC).isoformat() if value is not None else None
|
|
|
|
|
|
__all__ = [
|
|
"CutoverBackfillReport",
|
|
"CutoverInventoryReport",
|
|
"CutoverRevisionWatermarkInput",
|
|
"CutoverSmokeResultsInput",
|
|
"FinalDeltaInput",
|
|
"KnowledgeFSCutoverConflictError",
|
|
"KnowledgeFSCutoverError",
|
|
"KnowledgeFSCutoverGateBlockedError",
|
|
"KnowledgeFSCutoverNotFoundError",
|
|
"KnowledgeFSWorkspaceCutoverService",
|
|
"LegacyDependencyDashboard",
|
|
"LegacyDependencyInput",
|
|
"QuarantineResolutionInput",
|
|
"QuarantineResolutionReport",
|
|
"ShadowAuthorizationObservationInput",
|
|
"ShadowAuthorizationReport",
|
|
"ShadowCompletionInput",
|
|
"ShadowCompletionReport",
|
|
"WorkspaceInventoryInput",
|
|
"knowledge_fs_remote_activation_evidence_consistent",
|
|
"knowledge_fs_remote_freeze_evidence_consistent",
|
|
]
|