Publish documentation / publish (push) Successful in 1m27s
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation. Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
83 lines
3.5 KiB
Python
83 lines
3.5 KiB
Python
"""Apply only removals confirmed by a successful Catalog physical synchronization."""
|
|
|
|
import json
|
|
|
|
from pydantic import BaseModel, ConfigDict, Field
|
|
from sqlalchemy import text
|
|
|
|
from .models import MemoryConflict
|
|
from .review import digest
|
|
|
|
|
|
class RemovedColumn(BaseModel):
|
|
model_config = ConfigDict(extra="forbid")
|
|
table: str = Field(min_length=1, max_length=200)
|
|
column: str = Field(min_length=1, max_length=200)
|
|
|
|
|
|
class CleanupRequest(BaseModel):
|
|
model_config = ConfigDict(extra="forbid")
|
|
sync_id: str = Field(min_length=1, max_length=200)
|
|
database: str = Field(min_length=1, max_length=200)
|
|
schema_name: str = Field(min_length=1, max_length=200)
|
|
removed_tables: list[str] = Field(default_factory=list, max_length=10000)
|
|
removed_columns: list[RemovedColumn] = Field(default_factory=list, max_length=100000)
|
|
|
|
|
|
def cleanup(service, request: CleanupRequest):
|
|
service._admin()
|
|
fingerprint = digest(request.model_dump(mode="json"))
|
|
with service.repository.operation() as repo:
|
|
with repo.transaction() as connection:
|
|
params = {"w": repo.workspace_id, "sync": request.sync_id}
|
|
receipt = (
|
|
connection.execute(
|
|
text(
|
|
"SELECT request_hash,card_ids "
|
|
"FROM thoth_memory.cleanup_receipts WHERE workspace_id=:w AND sync_id=:sync"
|
|
),
|
|
params,
|
|
)
|
|
.mappings()
|
|
.first()
|
|
)
|
|
if receipt:
|
|
if receipt["request_hash"] != fingerprint:
|
|
raise MemoryConflict(
|
|
"Physical cleanup identity was reused with different removals"
|
|
)
|
|
identities = receipt["card_ids"]
|
|
else:
|
|
identities = list(
|
|
connection.execute(
|
|
text(
|
|
"SELECT DISTINCT card_id "
|
|
"FROM thoth_memory.dependencies d WHERE workspace_id=:w "
|
|
"AND database_id=:db AND schema_name=:schema AND (table_name=ANY(:tables) "
|
|
"OR EXISTS (SELECT 1 FROM jsonb_to_recordset(CAST(:columns AS jsonb)) "
|
|
'AS removed("table" text, "column" text) WHERE '
|
|
'removed."table"=d.table_name AND removed."column"=d.column_name))'
|
|
),
|
|
{
|
|
**params,
|
|
"db": request.database,
|
|
"schema": request.schema_name,
|
|
"tables": request.removed_tables,
|
|
"columns": json.dumps(
|
|
[c.model_dump() for c in request.removed_columns]
|
|
),
|
|
},
|
|
).scalars()
|
|
)
|
|
for identity in identities:
|
|
repo.delete(identity)
|
|
connection.execute(
|
|
text(
|
|
"INSERT INTO thoth_memory.cleanup_receipts "
|
|
"VALUES (:w,:sync,:hash,CAST(:ids AS jsonb))"
|
|
),
|
|
{**params, "hash": fingerprint, "ids": json.dumps(identities)},
|
|
)
|
|
results = [service._propagate(repo, identity) for identity in identities]
|
|
return {"deleted": len(identities), "indexed": all(r["indexed"] for r in results)}
|