Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
892e2ba
feat(api): persist queued session input
mmabrouk Sep 4, 2026
539e0ba
feat(api): promote pending session input
mmabrouk Sep 4, 2026
a9a7e35
feat(api): steer pending session input
mmabrouk Sep 4, 2026
033c703
feat(web): use durable session input queue
mmabrouk Sep 4, 2026
876d9e0
test(sessions): cover pending input transactions
mmabrouk Sep 4, 2026
f4cca69
fix(sessions): fence queued input ownership
mmabrouk Sep 5, 2026
ab75d64
fix(sessions): exclude failed runner terminals
mmabrouk Sep 5, 2026
f978137
fix(sessions): deduplicate idle input retries
mmabrouk Sep 5, 2026
cdc831d
fix(sessions): serialize input admission with settlement
mmabrouk Sep 5, 2026
8dd9452
fix(chat): route queue sends through durable admission
mmabrouk Sep 5, 2026
4c47b6a
fix(sessions): bind steer to collapsed stop
mmabrouk Sep 5, 2026
72611c8
fix(sessions): expose recoverable promoted inputs
mmabrouk Sep 5, 2026
8238f7d
fix(sessions): guard missing input delivery receipts
mmabrouk Sep 5, 2026
6b9d4a3
refactor(sessions): defer replay snapshot projection
mmabrouk Sep 5, 2026
3e2c82b
fix(sessions): queue behind promoted continuations
mmabrouk Sep 5, 2026
b5b30e8
fix(sessions): align steer settlement locks
mmabrouk Sep 5, 2026
db08307
fix(sessions): preserve running queue successors
mmabrouk Sep 5, 2026
744007f
fix(sessions): advertise durable queue capabilities
mmabrouk Sep 5, 2026
afd0b1d
fix(chat): keep stop available for durable runs
mmabrouk Sep 5, 2026
ea12e61
fix(chat): require durable queue admission
mmabrouk Sep 5, 2026
1b39397
fix(sessions): release completed turn before queue delivery
mmabrouk Sep 5, 2026
8356fc6
fix(chat): hide remote-run banner for idle queued input
mmabrouk Sep 5, 2026
737b343
fix(chat): deliver steer during active streams
mmabrouk Sep 5, 2026
973ed10
test(chat): cover running-elsewhere admission
mmabrouk Sep 5, 2026
d25b83d
test(chat): match the milestone 2 running-elsewhere copy
mmabrouk Sep 5, 2026
b8eaeb3
fix(web): bind the remote-run gate and the steer payload to milestone 2
mmabrouk Sep 5, 2026
b92f46e
style(web): format the merged liveness module
mmabrouk Sep 5, 2026
e7a3232
fix(sessions): unify the session snapshot contract
mmabrouk Sep 5, 2026
dcf4753
fix(sessions): keep the queue snapshot reachable without a stream
mmabrouk Sep 5, 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
22 changes: 21 additions & 1 deletion api/entrypoints/routers.py
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,7 @@
from oss.src.core.folders.service import FoldersService
from oss.src.core.workflows.service import WorkflowsService
from oss.src.core.workflows.service import SimpleWorkflowsService
from oss.src.core.workflows.dtos import WorkflowServiceRequest
from oss.src.core.workflows.static_catalog import StaticWorkflowCatalog
from oss.src.core.evaluators.service import EvaluatorsService
from oss.src.core.evaluators.service import SimpleEvaluatorsService
Expand Down Expand Up @@ -186,6 +187,9 @@
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.dbs.postgres.sessions.inputs.dbes import SessionInputDBE # noqa: F401
from oss.src.dbs.postgres.sessions.inputs.dao import SessionInputsDAO
from oss.src.core.sessions.inputs.service import SessionInputsService
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 @@ -603,6 +607,7 @@ async def lifespan(*args, **kwargs):
session_turns_dao = SessionTurnsDAO(engine=_transactions_engine)
session_commands_dao = SessionCommandsDAO(engine=_transactions_engine)
session_executions_dao = SessionExecutionsDAO(engine=_transactions_engine)
session_inputs_dao = SessionInputsDAO(engine=_transactions_engine)

connections_dao = ConnectionsDAO(engine=_transactions_engine)
mounts_dao = MountsDAO(engine=_transactions_engine)
Expand Down Expand Up @@ -1164,9 +1169,23 @@ async def _dispatch_detached_run(*, project_id, user_id, request, run_id=None) -
],
control_command_id=command.id,
continuation_execution_id=command.target_turn_id,
)
),
continue_input=lambda command: workflows_service.invoke_workflow_detached(
project_id=command.project_id,
user_id=command.created_by_id,
request=WorkflowServiceRequest.model_validate(command.data["request"]),
run_id=command.target_turn_id,
control_command_id=command.id,
),
),
executions_dao=session_executions_dao,
inputs_dao=session_inputs_dao,
)
session_inputs_service = SessionInputsService(
inputs_dao=session_inputs_dao,
streams_service=session_streams_service,
executions_dao=session_executions_dao,
continuation_resumer=session_commands_service.resume_recoverable_continuation,
)
workflows_service.set_session_continuation_resumer(
session_commands_service.resume_recoverable_continuation
Expand All @@ -1183,6 +1202,7 @@ async def _dispatch_detached_run(*, project_id, user_id, request, run_id=None) -
turns_service=session_turns_service,
sessions_service=sessions_service,
commands_service=session_commands_service,
inputs_service=session_inputs_service,
respond_task=_interactions_worker.respond_interaction,
interactions_dispatcher=_interactions_dispatcher,
)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
"""add durable session pending inputs

Revision ID: oss000000028
Revises: oss000000027
Create Date: 2026-09-04 15:00:00.000000
"""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql


revision: str = "oss000000028"
down_revision: Union[str, None] = "oss000000027"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.drop_constraint("ck_session_commands_kind", "session_commands", type_="check")
op.create_check_constraint(
"ck_session_commands_kind",
"session_commands",
"kind IN ('cancel', 'continue_interaction', 'continue_input')",
)
op.create_table(
"session_inputs",
sa.Column("project_id", sa.UUID(as_uuid=True), nullable=False),
sa.Column("id", sa.UUID(as_uuid=True), nullable=False),
sa.Column("session_id", sa.String(), nullable=False),
sa.Column("content", postgresql.JSONB(astext_type=sa.Text()), nullable=False),
sa.Column("position", sa.BigInteger(), nullable=False),
sa.Column("state", sa.String(), server_default="pending", nullable=False),
sa.Column("policy", sa.String(), nullable=False),
sa.Column("idempotency_key", sa.String(), nullable=False),
sa.Column("request_fingerprint", sa.String(length=64), nullable=False),
sa.Column("promoted_execution_id", sa.String(), nullable=True),
sa.Column(
"created_at",
sa.TIMESTAMP(timezone=True),
server_default=sa.func.current_timestamp(),
nullable=True,
),
sa.Column("updated_at", sa.TIMESTAMP(timezone=True), nullable=True),
sa.Column("deleted_at", sa.TIMESTAMP(timezone=True), nullable=True),
sa.Column("created_by_id", sa.UUID(as_uuid=True), nullable=True),
sa.Column("updated_by_id", sa.UUID(as_uuid=True), nullable=True),
sa.Column("deleted_by_id", sa.UUID(as_uuid=True), nullable=True),
sa.CheckConstraint(
"state IN ('pending', 'promoted', 'removed')",
name="ck_session_inputs_state",
),
sa.CheckConstraint(
"policy IN ('queue', 'steer')", name="ck_session_inputs_policy"
),
sa.ForeignKeyConstraint(["project_id"], ["projects.id"], ondelete="CASCADE"),
sa.PrimaryKeyConstraint("project_id", "id"),
)
op.create_index("uq_session_inputs_id", "session_inputs", ["id"], unique=True)
op.create_index(
"uq_session_inputs_idempotency",
"session_inputs",
["project_id", "session_id", "idempotency_key"],
unique=True,
)
op.create_index(
"uq_session_inputs_position",
"session_inputs",
["project_id", "session_id", "position"],
unique=True,
)
op.create_index(
"ix_session_inputs_pending",
"session_inputs",
["project_id", "session_id", "position"],
postgresql_where=sa.text("state = 'pending'"),
)


def downgrade() -> None:
op.drop_index("ix_session_inputs_pending", table_name="session_inputs")
op.drop_index("uq_session_inputs_position", table_name="session_inputs")
op.drop_index("uq_session_inputs_idempotency", table_name="session_inputs")
op.drop_index("uq_session_inputs_id", table_name="session_inputs")
op.drop_table("session_inputs")
op.drop_constraint("ck_session_commands_kind", "session_commands", type_="check")
op.create_check_constraint(
"ck_session_commands_kind",
"session_commands",
"kind IN ('cancel', 'continue_interaction')",
)
58 changes: 51 additions & 7 deletions api/oss/src/apis/fastapi/sessions/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
from oss.src.core.sessions.mounts.dtos import SessionMount, SessionMountQuery
from oss.src.core.sessions.turns.dtos import HarnessKind, SessionTurn, SessionTurnQuery
from oss.src.core.sessions.types import SessionReference
from oss.src.core.sessions.inputs.dtos import PendingInput
from oss.src.core.shared.dtos import OTelSpanId, Windowing
from oss.src.dbs.postgres.sessions.streams.dao import MAX_SESSION_QUERY_LIMIT

Expand Down Expand Up @@ -122,6 +123,35 @@ class SessionResponse(BaseModel):
session: Optional[SessionStream] = None


class SessionCapabilities(BaseModel):
durable_approvals: bool = False
queue: bool = False
steer: bool = False


class SessionExecutionSnapshot(BaseModel):
id: Optional[str] = None
state: Literal["idle", "running", "stopping"] = "idle"


class PendingInputResponse(BaseModel):
input: PendingInput


class PendingInputAdmissionRequest(BaseModel):
model_config = ConfigDict(extra="forbid")

session_id: SessionId
content: Dict[str, Any]
on_busy: Literal["reject", "queue", "steer"] = "reject"


class PendingInputAdmissionResponse(BaseModel):
action: Literal["execute", "pending"]
input: Optional[PendingInput] = None
execution_id: Optional[str] = None


# ---------------------------------------------------------------------------
# Streams request/response models
# ---------------------------------------------------------------------------
Expand All @@ -138,10 +168,6 @@ class SessionStreamQueryRequest(BaseModel):
is_running: Optional[bool] = None


class SessionCapabilities(BaseModel):
durable_approvals: bool = False


class SessionStreamResponse(BaseModel):
stream: Optional[SessionStream] = None
capabilities: SessionCapabilities = Field(default_factory=SessionCapabilities)
Expand Down Expand Up @@ -175,15 +201,33 @@ class SessionRecordsQueryResponse(BaseModel):


class SessionSnapshotPending(BaseModel):
inputs: List[Any] = Field(default_factory=list)
inputs: List[PendingInput] = Field(default_factory=list)
interactions: List[SessionInteraction] = Field(default_factory=list)


class SessionSnapshotResponse(BaseModel):
session: SessionStream
"""One snapshot for every reader of an open session.

`session`, `execution` and `read` are the nullable reconnect half: the stream row, the
latest turn (whose `end_time` says whether that turn is still live), and the durable
sequence watermark a reader replays from. They are absent when the shared reader is off or
before a fresh session has a stream row.

`execution_state` and `pending.inputs` are the queue half. `execution_state` is the
session's CURRENT lifecycle derived from the stream row, which is a different question from
`execution`: that names the last turn, this says whether anything is running right now.
`capabilities` reports the same flags the streams endpoint reports, from the same helper, so
a client never sees the two disagree.
"""

session: Optional[SessionStream] = None
execution: Optional[SessionTurn] = None
execution_state: SessionExecutionSnapshot = Field(
default_factory=SessionExecutionSnapshot
)
pending: SessionSnapshotPending
read: SessionRecordsReadState
read: Optional[SessionRecordsReadState] = None
capabilities: SessionCapabilities = Field(default_factory=SessionCapabilities)


class SessionRecordResponse(BaseModel):
Expand Down
Loading
Loading