Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
58 commits
Select commit Hold shift + click to select a range
af14418
feat(sessions): settle executions whose runner cannot report an outcome
mmabrouk Sep 2, 2026
5eb30cf
fix(runner): probe the sandbox over HTTP, not through a local cache
mmabrouk Sep 2, 2026
208be17
docs(sessions): record the execution watchdog slice
mmabrouk Sep 2, 2026
59b5059
fix(sessions): key the watchdog off heartbeat age, not a lease
mmabrouk Sep 2, 2026
4e554e5
docs(sessions): name the risk in the threshold question
mmabrouk Sep 2, 2026
c540a64
docs(sessions): record the re-run at the 90-second threshold
mmabrouk Sep 2, 2026
4d6e363
docs(sessions): the watchdog stack is down, and how to bring it back
mmabrouk Sep 2, 2026
559a8d1
feat(sessions): give records a place to mark a late write
mmabrouk Sep 3, 2026
a630aff
feat(sessions): stamp the watchdog as the writer of the ending it writes
mmabrouk Sep 3, 2026
54923d4
feat(sessions): keep a quarantined record out of every transcript read
mmabrouk Sep 3, 2026
8f567b0
feat(sessions): quarantine records that arrive after the watchdog end…
mmabrouk Sep 3, 2026
1eabbb5
fix(sessions): write the ending for a stopped row whose runner died b…
mmabrouk Sep 3, 2026
2301f91
fix(sessions): release the dead turn's alive lock when the sweep writ…
mmabrouk Sep 3, 2026
df82ad2
feat(sessions): gate durable stop and late output
mmabrouk Sep 3, 2026
1ffffbd
fix(sessions): redeliver abandoned stop commands
mmabrouk Sep 3, 2026
c4e7cb9
fix(sessions): enforce one terminal execution outcome
mmabrouk Sep 3, 2026
16cf06e
fix(sessions): make stop settlement atomic
mmabrouk Sep 3, 2026
354f319
fix(sessions): make terminal settlement authoritative
mmabrouk Sep 3, 2026
378fc27
fix(sessions): bound terminal redis repair
mmabrouk Sep 3, 2026
76d3579
fix(sessions): default invalid late output policy
mmabrouk Sep 3, 2026
f43076f
fix(runner): release parked approvals on stop
mmabrouk Sep 3, 2026
a0d15a8
fix(runner): preserve fresh prompts after approvals
mmabrouk Sep 3, 2026
b77432c
fix(sessions): publish settled interaction cancellations
mmabrouk Sep 3, 2026
4ffabe4
chore(sessions): remove dead settlement field
mmabrouk Sep 3, 2026
0e398a3
fix(runner): settle parked approvals before repark
mmabrouk Sep 3, 2026
4342a65
fix(runner): watch re-gates while settling approvals
mmabrouk Sep 3, 2026
50418b3
fix(api): select watchdog endings by execution
mmabrouk Sep 4, 2026
cd7117a
fix(api): wire commands service into session watchdog
mmabrouk Sep 4, 2026
356c176
fix(api): clear dead session owner in watchdog
mmabrouk Sep 4, 2026
b76b701
fix(sessions): document cancel execution guard
mmabrouk Sep 4, 2026
7146130
fix(api): bound watchdog ending candidates
mmabrouk Sep 4, 2026
9160f68
chore(sessions): number the ending-marker migration 026
mmabrouk Sep 4, 2026
2fe0e92
fix(runner): repark a parked-approval Stop warm when no harness cance…
mmabrouk Sep 4, 2026
5d8058e
fix(api): keep the execution watchdog alive when a sweep pass raises
mmabrouk Sep 4, 2026
55db87f
fix(api): give MultiLogger an exception method
mmabrouk Sep 4, 2026
2727465
fix(api): let the watchdog settle commands past a row it cannot map
mmabrouk Sep 4, 2026
a629cf2
fix(api): let a runner claim commands past a row it cannot map
mmabrouk Sep 4, 2026
80531f0
fix(api): clear the running flag when the watchdog marks an execution…
mmabrouk Sep 4, 2026
1fc9aea
fix(api): tombstone a swept turn so a returning runner cannot re-set …
mmabrouk Sep 4, 2026
ca43e25
fix(sessions): reclaim session affinity from a replica that holds no …
mmabrouk Sep 4, 2026
7904daa
fix(sessions): persist the watchdog collapse through a Core UPDATE
mmabrouk Sep 4, 2026
fbc8259
fix(sessions): write the lost-turn is_running clear as a Core UPDATE
mmabrouk Sep 4, 2026
0f85872
test(sessions): pin watchdog settlement invariants
mmabrouk Sep 4, 2026
5251677
test(runner): pin watchdog teardown and probe routes
mmabrouk Sep 4, 2026
d7e58dc
fix(runner): end a turn when the provider says its sandbox is gone
mmabrouk Sep 4, 2026
8229202
test(runner): pin the sandbox-gone terminal end to end
mmabrouk Sep 4, 2026
4aecd1f
fix(sessions): reconcile watchdog rebase
mmabrouk Sep 4, 2026
9a15766
fix(sessions): guard watchdog stream updates
mmabrouk Sep 4, 2026
91fbab0
fix(sessions): retry watchdog lookup failures
mmabrouk Sep 4, 2026
c0a5ab9
fix(runner): expose quiet-run timeout overrides
mmabrouk Sep 4, 2026
0446be9
test(sessions): model guarded sweep rowcounts
mmabrouk Sep 4, 2026
7202cea
fix(api): fence watchdog settlement cleanup
mmabrouk Sep 4, 2026
3b3d8cc
fix(api): fence heartbeat row mirrors
mmabrouk Sep 4, 2026
b10767e
fix(api): generation-fence session affinity
mmabrouk Sep 4, 2026
e51ca17
fix(api): restore legacy watchdog cleanup
mmabrouk Sep 4, 2026
5c5355e
test(api): follow generated owner values
mmabrouk Sep 4, 2026
0e27a42
fix(api): fence heartbeat on stream state
mmabrouk Sep 4, 2026
9609237
fix(api): preserve owner generation on reclaim
mmabrouk Sep 4, 2026
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
16 changes: 14 additions & 2 deletions api/entrypoints/routers.py
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,7 @@
from oss.src.core.sessions.streams.service import SessionStreamsService
from oss.src.dbs.postgres.sessions.commands.dbes import SessionCommandDBE # noqa: F401
from oss.src.dbs.postgres.sessions.commands.dao import SessionCommandsDAO
from oss.src.dbs.postgres.sessions.executions.dao import SessionExecutionsDAO
from oss.src.core.sessions.commands.service import SessionCommandsService
from oss.src.dbs.http.sessions.control_delivery_direct import DirectControlDelivery
from oss.src.tasks.asyncio.sessions.orphan_sweep import orphan_sweep_loop
Expand Down Expand Up @@ -283,8 +284,16 @@ async def lifespan(*args, **kwargs):
except Exception as e: # noqa: BLE001
log.warning("Store bucket ensure failed at startup: %s", e)

# The execution watchdog. It needs the records plane to write the terminal outcome a
# dead runner owed, and the watch publisher so an open browser sees the turn close.
_orphan_sweep_task = asyncio.create_task(
orphan_sweep_loop(_transactions_engine, _lock_engine)
orphan_sweep_loop(
_transactions_engine,
_lock_engine,
records_service=records_service,
watch_publisher=_sessions_watch_publisher,
commands_service=session_commands_service,
)
)

_attachment_sweep_task = asyncio.create_task(
Expand Down Expand Up @@ -591,6 +600,8 @@ async def lifespan(*args, **kwargs):
folders_dao = FoldersDAO(engine=_transactions_engine)
session_streams_dao = SessionStreamsDAO(engine=_transactions_engine)
session_turns_dao = SessionTurnsDAO(engine=_transactions_engine)
session_commands_dao = SessionCommandsDAO(engine=_transactions_engine)
session_executions_dao = SessionExecutionsDAO(engine=_transactions_engine)

connections_dao = ConnectionsDAO(engine=_transactions_engine)
mounts_dao = MountsDAO(engine=_transactions_engine)
Expand Down Expand Up @@ -625,6 +636,7 @@ async def lifespan(*args, **kwargs):

records_service = RecordsService(
records_dao=records_dao,
executions_dao=session_executions_dao,
)


Expand Down Expand Up @@ -1131,13 +1143,13 @@ async def _dispatch_detached_run(*, project_id, user_id, request) -> str:
"Only 'direct' is implemented; the long-poll adapter is a later change."
)

session_commands_dao = SessionCommandsDAO()
session_commands_service = SessionCommandsService(
commands_dao=session_commands_dao,
streams_service=session_streams_service,
interactions_service=interactions_service,
lock_engine=_lock_engine,
delivery=DirectControlDelivery(),
executions_dao=session_executions_dao,
)

sessions = SessionsRouter(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
"""add authoritative session execution terminal outcomes

Revision ID: oss000000023
Revises: oss000000022
Create Date: 2026-09-03 22:00:00.000000
"""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa


revision: str = "oss000000023"
down_revision: Union[str, None] = "oss000000022"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.create_table(
"session_executions",
sa.Column("project_id", sa.UUID(as_uuid=True), nullable=False),
sa.Column("session_id", sa.String(), nullable=False),
sa.Column("execution_id", sa.String(), nullable=False),
sa.Column("terminal_outcome", sa.String(), nullable=False),
sa.Column("settled_by", sa.String(), nullable=False),
sa.Column("settled_at", sa.TIMESTAMP(timezone=True), nullable=False),
sa.ForeignKeyConstraint(["project_id"], ["projects.id"], ondelete="CASCADE"),
sa.PrimaryKeyConstraint("project_id", "session_id", "execution_id"),
)
op.create_index(
"ix_session_executions_project_session",
"session_executions",
["project_id", "session_id"],
)


def downgrade() -> None:
op.drop_index(
"ix_session_executions_project_session", table_name="session_executions"
)
op.drop_table("session_executions")
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
"""track execution Redis reconciliation

Revision ID: oss000000024
Revises: oss000000023
Create Date: 2026-09-03 22:30:00.000000
"""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa


revision: str = "oss000000024"
down_revision: Union[str, None] = "oss000000023"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.add_column(
"session_executions",
sa.Column("redis_reconciled_at", sa.TIMESTAMP(timezone=True), nullable=True),
)
op.create_index(
"ix_session_executions_redis_unreconciled",
"session_executions",
["settled_at"],
postgresql_where=sa.text(
"settled_by = 'runner' AND terminal_outcome = 'stopped' "
"AND redis_reconciled_at IS NULL"
),
)


def downgrade() -> None:
op.drop_index(
"ix_session_executions_redis_unreconciled",
table_name="session_executions",
)
op.drop_column("session_executions", "redis_reconciled_at")
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
"""add session execution ending marker

Revision ID: oss000000026
Revises: oss000000024
Create Date: 2026-09-04 12:00:00.000000
"""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa


revision: str = "oss000000026"
down_revision: Union[str, None] = "oss000000024"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.add_column(
"session_executions",
sa.Column("ending_written_at", sa.TIMESTAMP(timezone=True), nullable=True),
)
op.create_index(
"ix_session_executions_ending_unwritten",
"session_executions",
["settled_at"],
postgresql_where=sa.text("ending_written_at IS NULL"),
)


def downgrade() -> None:
op.drop_index(
"ix_session_executions_ending_unwritten",
table_name="session_executions",
)
op.drop_column("session_executions", "ending_written_at")
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
"""add_records_quarantined_at

Revision ID: oss000000005
Revises: oss000000004
Create Date: 2026-09-03 12:00:00.000000

"""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa

# revision identifiers, used by Alembic.
revision: str = "oss000000005"
down_revision: Union[str, None] = "oss000000004"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
# A record that reached ingest for a turn the execution watchdog had already ended.
# Nullable and forward-fill only, like every other column on this table: the tracing DB
# is never backfilled, and no existing row can be classified retroactively anyway.
#
# No index. Every read that filters on it is already scoped to one project and one
# session by an existing index, and the column is null on all but a handful of rows.
op.add_column(
"records",
sa.Column("quarantined_at", sa.TIMESTAMP(timezone=True), nullable=True),
)


def downgrade() -> None:
op.drop_column("records", "quarantined_at")
8 changes: 7 additions & 1 deletion api/oss/src/apis/fastapi/sessions/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -360,7 +360,13 @@ class SessionCancelRequest(BaseModel):
# refuses the request if another one is running. When absent, it cancels whichever
# execution is active when the request is applied. A person never types this: the browser
# fills it from the session's own state, and a first-party client always sends it.
expected_execution_id: Optional[str] = None
expected_execution_id: Optional[str] = Field(
default=None,
description=(
"Optional stale-request guard honored only in cancel mode; ignored for send, "
"steer, and attach."
),
)


class SessionCommandRef(BaseModel):
Expand Down
22 changes: 19 additions & 3 deletions api/oss/src/core/sessions/commands/interfaces.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@

from abc import ABC, abstractmethod
from datetime import datetime
from typing import List, NamedTuple, Optional
from typing import Any, AsyncContextManager, List, NamedTuple, Optional
from uuid import UUID

from pydantic import BaseModel
Expand Down Expand Up @@ -69,6 +69,10 @@ async def acknowledge(self, *, command_id: UUID, replica_id: str) -> None:


class SessionCommandsDAOInterface(ABC):
@abstractmethod
def transaction(self) -> AsyncContextManager[Any]:
"""Open a transaction that sibling session DAOs can share."""

@abstractmethod
async def create_command(
self,
Expand Down Expand Up @@ -147,11 +151,23 @@ async def claim_for_delivery(
a direct call. The long-poll adapter reaches the same transition through
`claim_commands`; both exist so the outcome route's guard reads the same either way."""

@abstractmethod
async def record_delivery_attempt(
self,
*,
project_id: UUID,
command_id: UUID,
now: datetime,
max_deliveries: int,
) -> Optional[SessionCommand]:
"""Reserve one bounded delivery attempt and return the updated command."""

@abstractmethod
async def settle_command(
self,
*,
settle: SessionCommandSettle,
transaction: Optional[Any] = None,
) -> Optional[SessionCommand]:
"""Terminal transition, guarded on `state='claimed' AND claimed_by=:replica_id`.
None means the claim had expired or somebody else settled it first."""
Expand All @@ -173,6 +189,6 @@ async def expire_claims(
*,
now: datetime,
max_deliveries: int,
pending_before: Optional[datetime] = None,
) -> List[SessionCommand]:
"""Commands whose claim lease has passed. The settlement sweep reads this. Not called
in this slice; the execution watchdog owns settlement (see the slice document)."""
"""Pending or claimed commands old enough for recovery."""
Loading
Loading