feat: implement async workspace deletion with active session checks (#378)

* feat: implement async workspace deletion with active session checks

- Updated the DELETE /workspaces/:id endpoint to return 202 Accepted, indicating that the deletion request is processed in the background.
- Added a check for active sessions before allowing workspace deletion, raising a ConflictException if any exist.
- Updated related tests to ensure proper handling of active sessions during workspace deletion.

* fix: Address review issues

---------

Co-authored-by: Vineeth Voruganti <13438633+VVoruganti@users.noreply.github.com>
This commit is contained in:
doria 2026-02-12 18:48:00 -05:00 committed by GitHub
parent 5f7dad97b9
commit 3aaced2cd1
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
9 changed files with 296 additions and 65 deletions

View File

@ -150,7 +150,7 @@ describe('Honcho Client', () => {
// ===========================================================================
describe('DELETE /workspaces/:id', () => {
test('deleteWorkspace removes workspace', async () => {
test('deleteWorkspace accepts deletion', async () => {
// Create a workspace to delete
const tempWorkspaceId = generateWorkspaceId('delete')
const tempClient = new Honcho({
@ -162,12 +162,8 @@ describe('Honcho Client', () => {
// Ensure it exists
await tempClient.getMetadata()
// Delete it (returns void)
// Delete it (returns 202 Accepted — deletion is processed in the background)
await client.deleteWorkspace(tempWorkspaceId)
// Verify it's gone from list
const page = await client.workspaces()
expect(page.items).not.toContain(tempWorkspaceId)
})
})

View File

@ -59,6 +59,7 @@ from .webhook import (
)
from .workspace import (
WorkspaceDeletionResult,
check_no_active_sessions,
delete_workspace,
get_all_workspaces,
get_or_create_workspace,
@ -128,6 +129,7 @@ __all__ = [
"list_webhook_endpoints",
# Workspace
"WorkspaceDeletionResult",
"check_no_active_sessions",
"delete_workspace",
"get_or_create_workspace",
"get_workspace",

View File

@ -3,7 +3,7 @@ from logging import getLogger
from typing import Any
from cashews import NOT_NONE
from sqlalchemy import Select, delete, func, select
from sqlalchemy import Select, delete, exists, func, select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncSession
@ -227,6 +227,33 @@ async def update_workspace(
return honcho_workspace
async def check_no_active_sessions(db: AsyncSession, workspace_name: str) -> None:
"""
Verify that a workspace has no active sessions.
Args:
db: Database session
workspace_name: Name of the workspace
Raises:
ConflictException: If active sessions exist in the workspace
"""
has_active_sessions: bool = bool(
await db.scalar(
select(
exists().where(
models.Session.workspace_name == workspace_name,
models.Session.is_active == True, # noqa: E712
)
)
)
)
if has_active_sessions:
raise ConflictException(
f"Cannot delete workspace '{workspace_name}': active session(s) remain. Delete all sessions first."
)
async def delete_workspace(
db: AsyncSession, workspace_name: str
) -> WorkspaceDeletionResult:
@ -250,6 +277,11 @@ async def delete_workspace(
logger.warning("Workspace %s not found", workspace_name)
raise ResourceNotFoundException()
# NOTE: No active session check here — that gate lives in the router.
# This crud method is called by the background worker, where a session
# could have been created after the user's request was accepted (202).
# The deletion should proceed and cascade-delete any new sessions.
# Create a snapshot of the workspace data before deletion
workspace_snapshot = schemas.Workspace(
name=honcho_workspace.name,

View File

@ -207,6 +207,8 @@ async def process_deletion(
resource_id = payload.resource_id
success = True
error_message: str | None = None
peers_deleted = 0
sessions_deleted = 0
messages_deleted = 0
conclusions_deleted = 0
@ -260,6 +262,30 @@ async def process_deletion(
str(e),
)
elif deletion_type == "workspace":
try:
result = await crud.delete_workspace(db, workspace_name=workspace_name)
peers_deleted = result.peers_deleted
sessions_deleted = result.sessions_deleted
messages_deleted = result.messages_deleted
conclusions_deleted = result.conclusions_deleted
logger.info(
"Successfully deleted workspace %s "
+ "(peers=%d, sessions=%d, messages=%d, conclusions=%d)",
workspace_name,
peers_deleted,
sessions_deleted,
messages_deleted,
conclusions_deleted,
)
except ResourceNotFoundException as e:
# Workspace not found - may have already been deleted, treat as success
logger.warning(
"Workspace %s not found during deletion (may already be deleted): %s",
workspace_name,
str(e),
)
else:
success = False
error_message = f"Unsupported deletion type: {deletion_type}"
@ -272,6 +298,8 @@ async def process_deletion(
deletion_type=deletion_type,
resource_id=resource_id,
success=success,
peers_deleted=peers_deleted,
sessions_deleted=sessions_deleted,
messages_deleted=messages_deleted,
conclusions_deleted=conclusions_deleted,
error_message=error_message,

View File

@ -558,7 +558,7 @@ async def enqueue_dream(
def create_deletion_record(
workspace_name: str,
deletion_type: Literal["session", "observation"],
deletion_type: Literal["session", "observation", "workspace"],
resource_id: str,
) -> dict[str, Any]:
"""
@ -589,7 +589,7 @@ def create_deletion_record(
async def enqueue_deletion(
workspace_name: str,
deletion_type: Literal["session", "observation"],
deletion_type: Literal["session", "observation", "workspace"],
resource_id: str,
db_session: AsyncSession | None = None,
) -> None:

View File

@ -9,10 +9,9 @@ from sqlalchemy.ext.asyncio import AsyncSession
from src import crud, models, schemas
from src.config import settings
from src.dependencies import db
from src.deriver.enqueue import enqueue_dream
from src.deriver.enqueue import enqueue_deletion, enqueue_dream
from src.exceptions import AuthenticationException
from src.security import JWTParams, require_auth
from src.telemetry.events import DeletionCompletedEvent, emit
from src.utils.search import search
logger = logging.getLogger(__name__)
@ -101,8 +100,7 @@ async def update_workspace(
@router.delete(
"/{workspace_id}",
status_code=204,
response_model=None,
status_code=202,
dependencies=[Depends(require_auth(workspace_name="workspace_id"))],
)
async def delete_workspace(
@ -110,26 +108,26 @@ async def delete_workspace(
db: AsyncSession = db,
):
"""
Delete a Workspace. This will permanently delete all sessions, peers, messages, and conclusions
associated with the workspace.
Delete a Workspace. This accepts the deletion request and processes it in the background,
permanently deleting all peers, messages, conclusions, and other resources associated
with the workspace.
Returns 409 Conflict if the workspace contains active sessions.
Delete all sessions first, then delete the workspace.
This action cannot be undone.
"""
result = await crud.delete_workspace(db, workspace_name=workspace_id)
# Verify workspace exists
await crud.get_workspace(db, workspace_name=workspace_id)
# Emit telemetry event with cascade counts
emit(
DeletionCompletedEvent(
workspace_name=workspace_id,
deletion_type="workspace",
resource_id=workspace_id,
success=True,
peers_deleted=result.peers_deleted,
sessions_deleted=result.sessions_deleted,
messages_deleted=result.messages_deleted,
conclusions_deleted=result.conclusions_deleted,
)
)
# Check for active sessions before accepting
await crud.check_no_active_sessions(db, workspace_name=workspace_id)
# Enqueue for background deletion
await enqueue_deletion(workspace_id, "workspace", workspace_id, db_session=db)
await db.commit()
return {"message": "Workspace deletion accepted"}
@router.post(

View File

@ -63,7 +63,7 @@ class DeletionPayload(BasePayload):
"""Payload for deletion tasks."""
task_type: Literal["deletion"] = "deletion"
deletion_type: Literal["session", "observation"]
deletion_type: Literal["session", "observation", "workspace"]
resource_id: str
@ -101,7 +101,7 @@ def create_dream_payload(
def create_deletion_payload(
deletion_type: Literal["session", "observation"],
deletion_type: Literal["session", "observation", "workspace"],
resource_id: str,
) -> dict[str, Any]:
"""Create a deletion payload."""

View File

@ -4,7 +4,7 @@ from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from src import crud, models, schemas
from src.exceptions import ResourceNotFoundException
from src.exceptions import ConflictException, ResourceNotFoundException
class TestWorkspaceCRUD:
@ -49,12 +49,12 @@ class TestWorkspaceCRUD:
assert len(peers) == 0
@pytest.mark.asyncio
async def test_delete_workspace_cascade_sessions(
async def test_check_no_active_sessions_raises(
self,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""Test that deleting a workspace cascades to delete sessions"""
"""Test that check_no_active_sessions raises ConflictException when active sessions exist"""
test_workspace, _test_peer = sample_data
# Create sessions
@ -67,21 +67,33 @@ class TestWorkspaceCRUD:
db_session.add_all([session1, session2])
await db_session.flush()
# Verify sessions exist
stmt = select(models.Session).where(
models.Session.workspace_name == test_workspace.name
# check_no_active_sessions should raise ConflictException
with pytest.raises(ConflictException, match="active session"):
await crud.check_no_active_sessions(db_session, test_workspace.name)
@pytest.mark.asyncio
async def test_delete_workspace_succeeds_with_active_sessions(
self,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""Test that delete_workspace cascade-deletes active sessions (no guard in crud)"""
test_workspace, _test_peer = sample_data
# Create active sessions
session1 = models.Session(
name=str(generate_nanoid()), workspace_name=test_workspace.name
)
result = await db_session.execute(stmt)
sessions = result.scalars().all()
assert len(sessions) == 2
session2 = models.Session(
name=str(generate_nanoid()), workspace_name=test_workspace.name
)
db_session.add_all([session1, session2])
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
# Verify sessions are deleted
result = await db_session.execute(stmt)
sessions = result.scalars().all()
assert len(sessions) == 0
# crud.delete_workspace should succeed — the guard lives in the router
result = await crud.delete_workspace(db_session, test_workspace.name)
assert result.workspace.name == test_workspace.name
assert result.sessions_deleted == 2
@pytest.mark.asyncio
async def test_delete_workspace_cascade_messages(
@ -125,6 +137,10 @@ class TestWorkspaceCRUD:
messages = result.scalars().all()
assert len(messages) == 2
# Mark session inactive so workspace deletion is allowed
session.is_active = False
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
@ -211,6 +227,10 @@ class TestWorkspaceCRUD:
documents = result.scalars().all()
assert len(documents) == 1
# Mark session inactive so workspace deletion is allowed
session.is_active = False
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
@ -254,6 +274,10 @@ class TestWorkspaceCRUD:
session_peers = result.all()
assert len(session_peers) == 1
# Mark session inactive so workspace deletion is allowed
session.is_active = False
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
@ -328,6 +352,10 @@ class TestWorkspaceCRUD:
queue_items = result.scalars().all()
assert len(queue_items) == 1
# Mark session inactive so workspace deletion is allowed
session.is_active = False
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
@ -366,6 +394,10 @@ class TestWorkspaceCRUD:
active_queues = result.scalars().all()
assert len(active_queues) == 1
# Mark session inactive so workspace deletion is allowed
session.is_active = False
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
@ -496,6 +528,11 @@ class TestWorkspaceCRUD:
assert len((await db_session.execute(document_stmt)).scalars().all()) == 1
assert len((await db_session.execute(webhook_stmt)).scalars().all()) == 1
# Mark sessions inactive so workspace deletion is allowed
session1.is_active = False
session2.is_active = False
await db_session.flush()
# Delete workspace
await crud.delete_workspace(db_session, test_workspace.name)
@ -506,3 +543,72 @@ class TestWorkspaceCRUD:
assert len((await db_session.execute(collection_stmt)).scalars().all()) == 0
assert len((await db_session.execute(document_stmt)).scalars().all()) == 0
assert len((await db_session.execute(webhook_stmt)).scalars().all()) == 0
@pytest.mark.asyncio
async def test_delete_workspace_allows_inactive_sessions(
self,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""Test that workspace deletion succeeds when all sessions are inactive"""
test_workspace, _test_peer = sample_data
# Create sessions and mark them inactive
session1 = models.Session(
name=str(generate_nanoid()),
workspace_name=test_workspace.name,
is_active=False,
)
session2 = models.Session(
name=str(generate_nanoid()),
workspace_name=test_workspace.name,
is_active=False,
)
db_session.add_all([session1, session2])
await db_session.flush()
# Delete workspace should succeed
result = await crud.delete_workspace(db_session, test_workspace.name)
assert result.workspace.name == test_workspace.name
# Verify sessions are cleaned up
stmt = select(models.Session).where(
models.Session.workspace_name == test_workspace.name
)
remaining = await db_session.execute(stmt)
assert len(remaining.scalars().all()) == 0
@pytest.mark.asyncio
async def test_check_no_active_sessions_passes_after_deletion(
self,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""Test that check_no_active_sessions passes after sessions are deleted"""
test_workspace, _test_peer = sample_data
# Create active sessions
session1 = models.Session(
name=str(generate_nanoid()), workspace_name=test_workspace.name
)
session2 = models.Session(
name=str(generate_nanoid()), workspace_name=test_workspace.name
)
db_session.add_all([session1, session2])
await db_session.flush()
# check_no_active_sessions should fail with active sessions
with pytest.raises(ConflictException):
await crud.check_no_active_sessions(db_session, test_workspace.name)
# Delete the sessions
await db_session.delete(session1)
await db_session.delete(session2)
await db_session.flush()
# Now check_no_active_sessions should pass (no exception)
await crud.check_no_active_sessions(db_session, test_workspace.name)
# And workspace deletion should succeed
result = await crud.delete_workspace(db_session, test_workspace.name)
assert result.workspace.name == test_workspace.name

View File

@ -278,7 +278,7 @@ def test_delete_workspace(client: TestClient):
# Delete the workspace
response = client.delete(f"/v3/workspaces/{name}")
assert response.status_code == 204
assert response.status_code == 202
# Verify the workspace no longer exists by trying to update it
response = client.put(
@ -321,24 +321,17 @@ def test_delete_workspace_with_peers(client: TestClient):
# Delete workspace
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 204
assert response.status_code == 202
def test_delete_workspace_with_sessions(client: TestClient):
"""Test deleting a workspace that has sessions"""
"""Test that deleting a workspace with active sessions returns 409"""
workspace_name = str(generate_nanoid())
# Create workspace
response = client.post("/v3/workspaces", json={"name": workspace_name})
assert response.status_code in [200, 201]
# Create peer
peer_name = str(generate_nanoid())
response = client.post(
f"/v3/workspaces/{workspace_name}/peers", json={"name": peer_name}
)
assert response.status_code in [200, 201]
# Create sessions
session1_name = str(generate_nanoid())
session2_name = str(generate_nanoid())
@ -351,9 +344,11 @@ def test_delete_workspace_with_sessions(client: TestClient):
)
assert response.status_code in [200, 201]
# Delete workspace
# Delete workspace should fail with 409
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 204
assert response.status_code == 409
data = response.json()
assert "active session" in data["detail"].lower()
def test_delete_workspace_with_messages(client: TestClient):
@ -397,9 +392,13 @@ def test_delete_workspace_with_messages(client: TestClient):
)
assert response.status_code == 201
# Delete session first (marks inactive, required before workspace deletion)
response = client.delete(f"/v3/workspaces/{workspace_name}/sessions/{session_name}")
assert response.status_code == 202
# Delete workspace
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 204
assert response.status_code == 202
def test_delete_workspace_with_webhooks(client: TestClient):
@ -421,7 +420,7 @@ def test_delete_workspace_with_webhooks(client: TestClient):
# Delete workspace
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 204
assert response.status_code == 202
# Verify webhook is deleted by checking workspace doesn't exist
response = client.get(f"/v3/workspaces/{workspace_name}/webhooks")
@ -479,13 +478,20 @@ def test_delete_workspace_cascade(client: TestClient):
)
assert response.status_code == 201
# Delete sessions first (marks inactive, required before workspace deletion)
for session_name in session_names:
response = client.delete(
f"/v3/workspaces/{workspace_name}/sessions/{session_name}"
)
assert response.status_code == 202
# Delete the workspace
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 204
assert response.status_code == 202
def test_delete_workspace_returns_no_content(client: TestClient):
"""Test that delete workspace returns 204 No Content"""
def test_delete_workspace_returns_accepted(client: TestClient):
"""Test that delete workspace returns 202 Accepted"""
name = str(generate_nanoid())
metadata = {"key": "value", "number": 42}
configuration = {"feature": True}
@ -499,4 +505,67 @@ def test_delete_workspace_returns_no_content(client: TestClient):
# Delete workspace
response = client.delete(f"/v3/workspaces/{name}")
assert response.status_code == 204
assert response.status_code == 202
def test_delete_workspace_blocked_by_sessions_returns_409(client: TestClient):
"""Test that deleting a workspace with active sessions returns 409 with descriptive message"""
workspace_name = str(generate_nanoid())
# Create workspace
response = client.post("/v3/workspaces", json={"name": workspace_name})
assert response.status_code in [200, 201]
# Create a session
session_name = str(generate_nanoid())
response = client.post(
f"/v3/workspaces/{workspace_name}/sessions", json={"name": session_name}
)
assert response.status_code in [200, 201]
# Attempt to delete workspace
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 409
data = response.json()
assert "detail" in data
assert "active session" in data["detail"]
assert "delete all sessions first" in data["detail"].lower()
def test_delete_workspace_after_session_deletion(client: TestClient):
"""Test that workspace deletion succeeds after all sessions are deleted"""
workspace_name = str(generate_nanoid())
# Create workspace
response = client.post("/v3/workspaces", json={"name": workspace_name})
assert response.status_code in [200, 201]
# Create sessions
session1_name = str(generate_nanoid())
session2_name = str(generate_nanoid())
response = client.post(
f"/v3/workspaces/{workspace_name}/sessions", json={"name": session1_name}
)
assert response.status_code in [200, 201]
response = client.post(
f"/v3/workspaces/{workspace_name}/sessions", json={"name": session2_name}
)
assert response.status_code in [200, 201]
# Workspace deletion should fail
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 409
# Delete all sessions (marks inactive)
response = client.delete(
f"/v3/workspaces/{workspace_name}/sessions/{session1_name}"
)
assert response.status_code == 202
response = client.delete(
f"/v3/workspaces/{workspace_name}/sessions/{session2_name}"
)
assert response.status_code == 202
# Now workspace deletion should succeed
response = client.delete(f"/v3/workspaces/{workspace_name}")
assert response.status_code == 202