Files

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",
]