Skip to content
Open
Show file tree
Hide file tree
Changes from 4 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
7 changes: 7 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,13 @@ lifecycle: `run` executes the bundle without creating a conversation, while
Profile selection independently controls agent settings and credentials. A
profile on run-scoped work limits the saved secrets available to its command.

A running automation can submit work for an external subject to
`POST /v1/runs/{run_id}/subject-turns` with a source, stable subject key,
prompt, and idempotency key. The service creates or resumes that subject's
deterministic conversation. The scanner never receives runtime credentials or
attaches to the conversation itself, and one short run can fan out several
independent conversations within the configured concurrency limit.

Only conversation-scoped execution reads the server's authoritative
`conversation_runtime`. The service uses the same profile-backed conversation
API in local and Docker workspaces, so automation code remains independent of
Expand Down
54 changes: 54 additions & 0 deletions migrations/versions/028_add_subject_turn_requests.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
"""Add idempotent subject-turn requests.

Revision ID: 028
Revises: 027
"""

import sqlalchemy as sa
from alembic import op


revision = "028"
down_revision = "027"
branch_labels = None
depends_on = None


def upgrade() -> None:
op.create_table(
"automation_subject_turns",
sa.Column("id", sa.Uuid(), nullable=False),
sa.Column("automation_id", sa.Uuid(), nullable=False),
sa.Column("requester_run_id", sa.Uuid(), nullable=False),
sa.Column("subject_run_id", sa.Uuid(), nullable=False),
sa.Column("source", sa.String(100), nullable=False),
sa.Column("subject_key", sa.String(500), nullable=False),
sa.Column("idempotency_key", sa.String(500), nullable=False),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("CURRENT_TIMESTAMP"),
nullable=False,
),
sa.ForeignKeyConstraint(
["automation_id"], ["automations.id"], ondelete="CASCADE"
),
sa.ForeignKeyConstraint(
["requester_run_id"], ["automation_runs.id"], ondelete="CASCADE"
),
sa.ForeignKeyConstraint(
["subject_run_id"], ["automation_runs.id"], ondelete="CASCADE"
),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint(
"automation_id",
"source",
"subject_key",
"idempotency_key",
name="uq_automation_subject_turn_idempotency",
),
)


def downgrade() -> None:
op.drop_table("automation_subject_turns")
2 changes: 2 additions & 0 deletions openhands/automation/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
from openhands.automation.router import router
from openhands.automation.scheduler import scheduler_loop
from openhands.automation.streams import stream_supervisor_loop
from openhands.automation.subject_router import router as subject_router
from openhands.automation.telemetry_router import router as telemetry_router
from openhands.automation.uploads import router as uploads_router
from openhands.automation.utils.version import get_sdk_version, get_server_version_info
Expand Down Expand Up @@ -295,6 +296,7 @@ def _create_app() -> FastAPI:
app.include_router(webhook_router, prefix=_base_path)
app.include_router(telemetry_router, prefix=_base_path)
app.include_router(git_sync_router, prefix=_base_path)
app.include_router(subject_router, prefix=_base_path)

app.include_router(kv_router, prefix=_base_path)
app.include_router(router, prefix=_base_path)
Expand Down
181 changes: 179 additions & 2 deletions openhands/automation/conversations.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
import logging
import uuid
from dataclasses import dataclass
from typing import Any, Final
from typing import Any, Final, Literal

from sqlalchemy import select, text
from sqlalchemy.ext.asyncio import AsyncSession
Expand All @@ -20,14 +20,19 @@
FilterEvaluationError,
evaluate_expression,
)
from openhands.automation.models import AutomationRun, AutomationRunStatus
from openhands.automation.models import (
AutomationRun,
AutomationRunStatus,
AutomationSubjectTurn,
)
from openhands.automation.schemas import EventTrigger
from openhands.automation.subjects import conversation_id_for
from openhands.automation.utils import utcnow
from openhands.automation.utils.conversation_turn import (
compose_turn,
send_conversation_turn,
)
from openhands.automation.utils.run import create_conversation_turn_run


logger = logging.getLogger("automation.conversations")
Expand All @@ -41,6 +46,12 @@
AutomationRunStatus.SKIPPED,
)

_RETRYABLE = (
AutomationRunStatus.FAILED,
AutomationRunStatus.CANCELLED,
AutomationRunStatus.SKIPPED,
)

# Matches AutomationRun.subject_key. Truncating would merge two subjects.
MAX_SUBJECT_KEY_LENGTH: Final[int] = 500

Expand Down Expand Up @@ -70,6 +81,15 @@ def needs_run(self) -> bool:
return self.conversation_id is None


@dataclass(frozen=True, slots=True)
class SubjectTurnResult:
"""Result returned to a poller that submitted work for one subject."""

disposition: Literal["created", "queued", "delivered", "deduplicated"]
run_id: uuid.UUID
conversation_id: str


def resolve_subject_key(
trigger: EventTrigger,
payload: dict[str, Any],
Expand Down Expand Up @@ -149,6 +169,14 @@ def _clean_key(value: str, origin: str) -> str | None:
return key


def clean_subject_key(value: str) -> str:
"""Validate a caller-supplied subject key without changing its identity."""
key = _clean_key(value, "subject turn")
if key is None:
raise ValueError("subject_key must be 1 to 500 non-whitespace characters")
return key


async def _take_subject_lock(
session: AsyncSession,
automation_id: uuid.UUID,
Expand Down Expand Up @@ -262,6 +290,32 @@ async def continue_conversation(
mid-run continues that conversation instead of racing a second run.
"""
await _take_subject_lock(session, automation_id, source, subject_key)
return await _continue_conversation_locked(
session,
org_id=org_id,
source=source,
subject_key=subject_key,
automation_id=automation_id,
event_key=event_key,
event_payload=event_payload,
turn_text=turn_text,
wake_agent=wake_agent,
)


async def _continue_conversation_locked(
session: AsyncSession,
*,
org_id: uuid.UUID,
source: str,
subject_key: str,
automation_id: uuid.UUID,
event_key: str,
event_payload: dict[str, Any] | None,
turn_text: str | None,
wake_agent: bool,
) -> ContinueResult:
"""Continue one subject after the caller has acquired its transaction lock."""

run = await _lock_subject_run(session, automation_id, source, subject_key)
if run is None:
Expand Down Expand Up @@ -294,3 +348,126 @@ async def continue_conversation(
return ContinueResult()

return ContinueResult(conversation_id=conversation_id)


async def submit_subject_turn(
session: AsyncSession,
*,
requester: AutomationRun,
source: str,
subject_key: str,
turn: str,
idempotency_key: str,
wake_agent: bool,
) -> SubjectTurnResult:
"""Create or continue conversation work selected by a running automation.

The requester chooses an external identity and prompt. The service owns
conversation identity, profile selection, runtime attachment, and
serialization. The transaction-scoped subject lock also orders duplicate
idempotency checks, so one retry cannot enqueue two conversations.
"""
subject_key = clean_subject_key(subject_key)
await _take_subject_lock(session, requester.automation_id, source, subject_key)

duplicate = (
(
await session.execute(
select(AutomationSubjectTurn).where(
AutomationSubjectTurn.automation_id == requester.automation_id,
AutomationSubjectTurn.source == source,
AutomationSubjectTurn.subject_key == subject_key,
AutomationSubjectTurn.idempotency_key == idempotency_key,
)
)
)
.scalars()
.first()
)
retry_record: AutomationSubjectTurn | None = None
if duplicate is not None:
duplicate_run = await session.get(AutomationRun, duplicate.subject_run_id)
if duplicate_run is None:
raise RuntimeError("Subject-turn idempotency record has no run")
conversation_id = duplicate_run.conversation_id or conversation_id_for(
requester.automation.org_id,
requester.automation.id,
source,
subject_key,
)
can_retry = duplicate_run.status in _RETRYABLE and (
duplicate_run.started_at is None
or duplicate_run.subject_released_at is not None
)
if not can_retry:
return SubjectTurnResult(
disposition="deduplicated",
run_id=duplicate_run.id,
conversation_id=conversation_id,
)
# A run canceled before dispatch never owned a runtime, but excluding it
# from the subject lookup still requires the normal released marker.
if duplicate_run.started_at is None:
duplicate_run.subject_released_at = utcnow()
retry_record = duplicate

outcome = await _continue_conversation_locked(
session,
org_id=requester.automation.org_id,
source=source,
subject_key=subject_key,
automation_id=requester.automation_id,
event_key="subject.turn",
event_payload=None,
turn_text=turn,
wake_agent=wake_agent,
)

subject_run = await _lock_subject_run(
session, requester.automation_id, source, subject_key
)
if outcome.needs_run:
if subject_run is not None and subject_run.status not in _FINISHED:
raise RuntimeError("The subject conversation is not reachable yet")
subject_run = create_conversation_turn_run(
requester,
source=source,
subject_key=subject_key,
turn=turn,
wake_agent=wake_agent,
)
conversation_id = conversation_id_for(
requester.automation.org_id,
requester.automation.id,
source,
subject_key,
)
session.add(subject_run)
disposition = "created"
else:
assert subject_run is not None
conversation_id = outcome.conversation_id
assert conversation_id is not None
disposition = "queued" if outcome.coalesced else "delivered"

if retry_record is None:
retry_record = AutomationSubjectTurn(
automation_id=requester.automation_id,
requester_run_id=requester.id,
subject_run_id=subject_run.id,
source=source,
subject_key=subject_key,
idempotency_key=idempotency_key,
)
session.add(retry_record)
else:
# Keep the unique idempotency record and point it at the new attempt.
# The superseded run remains in run history for diagnosis.
retry_record.requester_run_id = requester.id
retry_record.subject_run_id = subject_run.id
await session.flush()
return SubjectTurnResult(
disposition=disposition,
run_id=subject_run.id,
conversation_id=conversation_id,
)
Loading
Loading