Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,11 @@ make dev # DevSpace + Kind with hot reload
make dev-stop # stop the development environment
make dev-reset # reset local development data
make dev-clear-cache
scripts/dev-ui.sh up # this worktree's frontend (Vite, hot reload) against the shared Kind backend
```

`scripts/dev-ui.sh` is safe with several worktrees running at once: each gets its own Vite port and they share one backend port-forward. The Kind backend itself is shared, so a backend image loaded from one worktree serves all of them.

The local Kubernetes workflow requires DevSpace and Kind. Python packages are managed by `uv`; JavaScript packages are managed by `pnpm`.

DevSpace synchronizes source into the containers. If a Svelte store change appears stale after hot reload, force recompilation of the importing component or clear the Vite cache; a browser refresh alone may retain a stale module reference.
Expand All @@ -74,12 +77,9 @@ Common commands:
make fmt # format and check files changed from main
make lint # lint files changed from main
pnpm check # workspace frontend/type checks
cd frontend && pnpm exec playwright test --project=fast --project=mobile
```

The Playwright `fast` and `mobile` projects use seeded data. The `e2e` project uses a real agent integration and must remain serial where conversation state is shared. Verify message submission before waiting for an agent response, use explicit button clicks, and wait for inputs to become enabled between messages.

Some historical tests and Make targets invoke live agents, external services, containers, or Kubernetes. Do not run the `e2e` project, live-agent tests, subscription-consuming commands, deployments, destructive resets, or production commands unless the task explicitly requires them and their target is known. Default automated tests for new native-agent adapters must use sanitized fixtures or fakes; keep live proofs opt-in and bounded.
Some historical tests and Make targets invoke live agents, external services, containers, or Kubernetes. Do not run the browser `e2e` tests, live-agent tests, subscription-consuming commands, deployments, destructive resets, or production commands unless the task explicitly requires them and their target is known. Default automated tests for new native-agent adapters must use sanitized fixtures or fakes; keep live proofs opt-in and bounded.

For Kubernetes commands, always specify the intended context. Tests must not rely on a developer's current context or mutate production resources.

Expand Down
82 changes: 51 additions & 31 deletions backend/src/mainloop/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -342,7 +342,10 @@ async def get_main_thread_info(user_id: str = Header(alias="X-User-ID", default=
from mainloop.runtime import delegation, native_sessions

binding = await delegation.ensure_main_session(user_id)
await native_sessions.sync(binding["session_id"])
# A rotation holds the session lock for the cut; do not queue behind it, so the UI can
# show "rotating" while it happens (the reconcile loop mirrors journal evidence anyway).
if not native_sessions.is_rotating(binding["session_id"]):
await native_sessions.sync(binding["session_id"])
session = await db.get_session(binding["session_id"])
topics = await delegation._topic_lines(user_id)
return MainThreadInfo(
Expand Down Expand Up @@ -428,7 +431,7 @@ async def get_conversation(conversation_id: str):
WHERE b.role='main' AND s.conversation_id=$1""",
conversation_id,
)
if main_sid:
if main_sid and not native_sessions.is_rotating(main_sid):
await native_sessions.sync(
main_sid
) # mirror new native-journal evidence first
Expand Down Expand Up @@ -885,21 +888,18 @@ async def send_session_message(
return {"status": "ok", "message_id": message.id}


class SessionLogsResponse(BaseModel):
"""Response for session logs."""

logs: str
source: str
session_status: str
# Done: nothing more will run, so the session can be cleared from the list.
FINISHED_STATUSES = frozenset(
{SessionStatus.COMPLETED, SessionStatus.FAILED, SessionStatus.CANCELLED}
)


@app.get("/sessions/{session_id}/logs", response_model=SessionLogsResponse)
async def get_session_logs(
@app.post("/sessions/{session_id}/cancel")
async def cancel_session(
session_id: str,
tail: int = 100,
user_id: str = Header(alias="X-User-ID", default=None),
):
"""Get execution logs for a session."""
"""Cancel a running session."""
if not user_id:
user_id = get_user_id_from_cf_header()

Expand All @@ -910,40 +910,60 @@ async def get_session_logs(
if session.user_id != user_id:
raise HTTPException(status_code=403, detail="Not your session")

# TODO: Get logs from K8s pod
logs = ""
source = "none"
if session.status in FINISHED_STATUSES:
return {"status": session.status.value, "agent": "not_running"}

return SessionLogsResponse(
logs=logs,
source=source,
session_status=session.status.value,
)
from mainloop.runtime import native_sessions

if await native_sessions.get_binding(session_id):
try:
agent = await native_sessions.cancel(session_id)
except ValueError as exc:
raise HTTPException(status_code=409, detail=str(exc)) from exc
else:
# TODO: Cancel the legacy session worker workflow
await db.update_session(session_id, status=SessionStatus.CANCELLED)
agent = "not_running"

return {"status": "cancelled", "agent": agent}

@app.post("/sessions/{session_id}/cancel")
async def cancel_session(

@app.post("/sessions/archive-finished")
async def archive_finished_sessions(
user_id: str = Header(alias="X-User-ID", default=None),
):
"""Clear every finished session (done, failed, cancelled) from the list.

Rows are kept for audit; sessions still running or waiting are left alone.
"""
if not user_id:
user_id = get_user_id_from_cf_header()
archived = await db.archive_sessions(user_id)
return {"archived": archived}


@app.post("/sessions/{session_id}/archive")
async def archive_session(
session_id: str,
user_id: str = Header(alias="X-User-ID", default=None),
):
"""Cancel a running session."""
"""Clear one finished session from the list. A live one must be cancelled first."""
if not user_id:
user_id = get_user_id_from_cf_header()

session = await db.get_session(session_id)
if not session:
raise HTTPException(status_code=404, detail="Session not found")

if session.user_id != user_id:
raise HTTPException(status_code=403, detail="Not your session")
if session.status not in FINISHED_STATUSES:
raise HTTPException(
status_code=409,
detail="This session is still running or waiting; cancel it first.",
)

await db.update_session(
session_id, status=SessionStatus.FAILED, error="Cancelled by user"
)

# TODO: Cancel session worker workflow

return {"status": "cancelled"}
await db.archive_sessions(user_id, session_ids=[session_id])
return {"status": "archived"}


# ============= Notification Endpoints =============
Expand Down
51 changes: 50 additions & 1 deletion backend/src/mainloop/db/postgres.py
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,17 @@ def _parse_json_field(value: Any) -> list | dict | None:
ALTER TABLE native_bindings ADD COLUMN IF NOT EXISTS turns_in_lineage INTEGER NOT NULL DEFAULT 0;
ALTER TABLE native_bindings ADD COLUMN IF NOT EXISTS reported_at TIMESTAMPTZ;
ALTER TABLE native_bindings ADD COLUMN IF NOT EXISTS continuations INTEGER NOT NULL DEFAULT 0;
ALTER TABLE sessions ADD COLUMN IF NOT EXISTS archived_at TIMESTAMPTZ;
-- Cancelling used to record status failed + this error text (and agent sync could then revive
-- it). Cancelled is its own status now; correct the old rows. Idempotent.
UPDATE sessions SET status = 'cancelled', error = NULL
WHERE error = 'Cancelled by user' AND status <> 'cancelled';
-- A child that has reported is done, not waiting on the user. New reports set this directly;
-- this corrects children that reported before that, which no sync would revisit. Idempotent.
UPDATE sessions SET status = 'completed'
WHERE status = 'waiting_on_user'
AND EXISTS (SELECT 1 FROM native_bindings b
WHERE b.session_id = sessions.id AND b.role = 'child' AND b.reported_at IS NOT NULL);
ALTER TABLE native_deliveries ADD COLUMN IF NOT EXISTS source TEXT NOT NULL DEFAULT 'user';
CREATE UNIQUE INDEX IF NOT EXISTS idx_native_bindings_token ON native_bindings(token_hash) WHERE token_hash IS NOT NULL;
CREATE INDEX IF NOT EXISTS idx_native_bindings_parent ON native_bindings(parent_session_id);
Expand Down Expand Up @@ -1386,8 +1397,9 @@ async def list_sessions(
user_id: str,
status: SessionStatus | None = None,
limit: int = 50,
include_archived: bool = False,
) -> list[Session]:
"""List sessions for a user."""
"""List sessions for a user (cleared sessions are hidden unless asked for)."""
if not self._pool:
return []

Expand All @@ -1398,6 +1410,9 @@ async def list_sessions(
)
params: list[Any] = [user_id]

if not include_archived:
query += " AND archived_at IS NULL"

if status:
query += f" AND status = ${len(params) + 1}"
params.append(status.value)
Expand All @@ -1409,6 +1424,39 @@ async def list_sessions(
rows = await conn.fetch(query, *params)
return [self._row_to_session(row) for row in rows]

async def archive_sessions(
self,
user_id: str,
session_ids: list[str] | None = None,
parent_session_id: str | None = None,
) -> list[str]:
"""Clear finished sessions from the list; returns the ids that were archived.

Only finished sessions (completed, failed, cancelled) qualify, and never the main
thread. Rows are kept for audit. ``session_ids`` narrows to those sessions and
``parent_session_id`` to the direct children of that session; both omitted means every
finished session of the user.
"""
if not self._pool:
return []
async with self.connection() as conn:
rows = await conn.fetch(
"""UPDATE sessions SET archived_at = NOW()
WHERE user_id = $1 AND archived_at IS NULL
AND status IN ('completed', 'failed', 'cancelled')
AND NOT EXISTS (SELECT 1 FROM native_bindings m
WHERE m.session_id = sessions.id AND m.role = 'main')
AND ($2::text[] IS NULL OR id = ANY($2))
AND ($3::text IS NULL OR EXISTS (
SELECT 1 FROM native_bindings c
WHERE c.session_id = sessions.id AND c.parent_session_id = $3))
RETURNING id""",
user_id,
session_ids,
parent_session_id,
)
return [r["id"] for r in rows]

async def update_session(
self,
session_id: str,
Expand Down Expand Up @@ -1559,6 +1607,7 @@ def _row_to_session(self, row: asyncpg.Record) -> Session:
created_at=row["created_at"],
started_at=row.get("started_at"),
completed_at=row.get("completed_at"),
archived_at=row.get("archived_at"),
summary=row.get("summary"),
error=row.get("error"),
# Code work fields
Expand Down
84 changes: 84 additions & 0 deletions backend/src/mainloop/runtime/agent_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,10 @@ def hash_token(token: str) -> str:
return hashlib.sha256(token.encode()).hexdigest()


# A session in one of these is done: nothing more will run and it can be cleared from the list.
FINISHED_STATUSES = frozenset({"completed", "failed", "cancelled"})


class Store(Protocol):
async def binding_by_token_hash(self, token_hash: str) -> dict | None: ...
async def get_binding(self, session_id: str) -> dict | None: ...
Expand All @@ -57,6 +61,10 @@ async def children_state(self, parent_session_id: str) -> list[dict]: ...
async def messages(
self, session_id: str, offset: int, limit: int
) -> list[dict]: ...
async def cancel_session(self, session_id: str) -> str: ...
async def archive_children(
self, user_id: str, parent_session_id: str, session_ids: list[str] | None
) -> list[str]: ...
async def spawn_child(
self, parent: dict, topic: dict, kind: str, title: str, brief: str
) -> str: ...
Expand Down Expand Up @@ -190,6 +198,64 @@ async def report(self, ctx: Ctx, summary: str, *, fallback: bool = False) -> dic
"message_id": mid,
}

# -- cleanup: end a child, clear finished ones from the list -------------------------------
async def _child(self, ctx: Ctx, session_id: str) -> dict:
rows = await self.store.children_state(ctx.binding["session_id"])
match = [r for r in rows if r["session_id"].startswith(session_id)]
if not match:
raise HTTPException(status_code=404, detail="no such child in your tree")
if len(match) > 1:
raise HTTPException(
status_code=400, detail="ambiguous session id; give more characters"
)
return match[0]

@staticmethod
def _require_manager(ctx: Ctx) -> None:
try:
policy.may_manage_children(ctx.actor)
except PolicyError as exc:
raise HTTPException(
status_code=403, detail=f"[{exc.code}] {exc.message}"
) from exc

async def cancel(self, ctx: Ctx, session_id: str) -> dict:
self._require_manager(ctx)
child = await self._child(ctx, session_id)
short = child["session_id"][:8]
if child["status"] in FINISHED_STATUSES:
return {"text": f"{short} is already {child['status']}; nothing to cancel"}
outcome = await self.store.cancel_session(child["session_id"])
text = f"cancelled {short}"
if outcome == "unknown":
text += "; stopping its agent could not be confirmed, so it may still be running"
return {"text": text, "agent": outcome}

async def clear(self, ctx: Ctx, session_id: str | None) -> dict:
"""Clear finished children from the user's list (kept for audit, never deleted)."""
self._require_manager(ctx)
parent = ctx.binding["session_id"]
only = (
[(await self._child(ctx, session_id))["session_id"]] if session_id else None
)
archived = await self.store.archive_children(
ctx.binding["user_id"], parent, only
)
left = [
r
for r in await self.store.children_state(parent)
if (only is None or r["session_id"] in only)
and r["status"] not in FINISHED_STATUSES
]
text = f"cleared {len(archived)} finished child session(s)"
if left:
names = ", ".join(r["session_id"][:8] for r in left)
text += (
f"; not cleared because they are still running or waiting: {names} "
"(cancel one with `mainloop cancel <id>` first)"
)
return {"text": text, "cleared": archived}

# -- state, answered from Postgres only (no native turn) ----------------------------------
async def status(self, ctx: Ctx, session_id: str | None) -> dict:
rows = await self.store.children_state(ctx.binding["session_id"])
Expand Down Expand Up @@ -285,6 +351,14 @@ class DelegateIn(BaseModel):
brief: str


class CancelIn(BaseModel):
session: str = Field(..., min_length=1)


class ClearIn(BaseModel):
session: str | None = None


class ReportIn(BaseModel):
summary: str = Field(..., min_length=1)

Expand Down Expand Up @@ -326,6 +400,16 @@ async def delegate(
return await s.delegate(ctx, body.topic, body.kind, body.title, body.brief)


@router.post("/cancel")
async def cancel(body: CancelIn, ctx: CtxDep, s: SvcDep):
return await s.cancel(ctx, body.session)


@router.post("/clear")
async def clear(body: ClearIn, ctx: CtxDep, s: SvcDep):
return await s.clear(ctx, body.session)


@router.post("/report")
async def report(body: ReportIn, ctx: CtxDep, s: SvcDep):
return await s.report(ctx, body.summary)
Expand Down
Loading
Loading