Skip to content
Closed
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
106 changes: 106 additions & 0 deletions migrations/versions/010_post_run_callbacks.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
"""Add post-run callback configuration and execution records.

Revision ID: 010
Revises: 009
Create Date: 2026-06-24
"""

from collections.abc import Sequence

import sqlalchemy as sa
from alembic import op


revision: str = "010"
down_revision: str = "009"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None


def _is_sqlite() -> bool:
return op.get_bind().dialect.name == "sqlite"


def upgrade() -> None:
op.add_column("automations", sa.Column("callbacks", sa.JSON(), nullable=True))

op.create_table(
"automation_run_callbacks",
sa.Column("id", sa.Uuid(), nullable=False),
sa.Column("run_id", sa.Uuid(), nullable=False),
sa.Column("name", sa.String(length=100), nullable=False),
sa.Column("trigger_status", sa.String(length=20), nullable=False),
sa.Column("entrypoint", sa.Text(), nullable=False),
sa.Column("timeout", sa.Integer(), nullable=True),
sa.Column("status", sa.String(length=20), nullable=False),
sa.Column("bash_command_id", sa.String(length=64), nullable=True),
sa.Column("error_detail", sa.Text(), nullable=True),
sa.Column("order", sa.Integer(), nullable=False),
sa.Column("started_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("timeout_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("completed_at", sa.DateTime(timezone=True), nullable=True),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("CURRENT_TIMESTAMP"),
nullable=False,
),
sa.ForeignKeyConstraint(["run_id"], ["automation_runs.id"], ondelete="CASCADE"),
sa.PrimaryKeyConstraint("id"),
)
op.create_index(
"ix_automation_run_callbacks_run_id",
"automation_run_callbacks",
["run_id"],
)
op.create_index(
"ix_automation_run_callbacks_status",
"automation_run_callbacks",
["status"],
)
op.create_index(
"ix_automation_run_callbacks_timeout_at",
"automation_run_callbacks",
["timeout_at"],
)
op.create_index(
"ix_automation_run_callbacks_run_order",
"automation_run_callbacks",
["run_id", "order"],
)
op.create_index(
"ix_automation_run_callbacks_status_timeout",
"automation_run_callbacks",
["status", "timeout_at"],
)

if not _is_sqlite():
op.execute(
"COMMENT ON COLUMN automations.callbacks IS "
"'Post-run callback configuration for automation runs.'"
)


def downgrade() -> None:
op.drop_index(
"ix_automation_run_callbacks_status_timeout",
table_name="automation_run_callbacks",
)
op.drop_index(
"ix_automation_run_callbacks_run_order",
table_name="automation_run_callbacks",
)
op.drop_index(
"ix_automation_run_callbacks_timeout_at",
table_name="automation_run_callbacks",
)
op.drop_index(
"ix_automation_run_callbacks_status",
table_name="automation_run_callbacks",
)
op.drop_index(
"ix_automation_run_callbacks_run_id",
table_name="automation_run_callbacks",
)
op.drop_table("automation_run_callbacks")
op.drop_column("automations", "callbacks")
11 changes: 11 additions & 0 deletions openhands/automation/backends/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,17 @@ async def get_execution_context(
TimeoutError: If sandbox doesn't become ready in time (Cloud mode)
"""

@abstractmethod
async def get_existing_execution_context(
self, client: httpx.AsyncClient
) -> ExecutionContext:
"""Reconnect to an existing run execution context.

Used by post-run callbacks after the main run has already dispatched.
For Cloud mode this discovers the still-running sandbox by run.sandbox_id.
For Local mode this returns the configured persistent agent server.
"""

@abstractmethod
async def release_context(
self, client: httpx.AsyncClient, ctx: ExecutionContext
Expand Down
22 changes: 22 additions & 0 deletions openhands/automation/backends/cloud.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
from openhands.automation.utils.sandbox import (
cleanup_sandbox,
delete_sandbox,
get_sandbox_agent_url,
verify_run_status,
)

Expand Down Expand Up @@ -174,6 +175,27 @@ async def _do_acquire() -> tuple[str, str, str]:
api_key=await self._ensure_api_key(),
)

async def get_existing_execution_context(
self, client: httpx.AsyncClient
) -> ExecutionContext:
"""Discover the existing sandbox's agent server context."""
sandbox_id = self._run.sandbox_id
if not sandbox_id:
raise RuntimeError("Run has no sandbox_id for callback execution")

api_key = await self._ensure_api_key()
result = await get_sandbox_agent_url(client, self.api_url, api_key, sandbox_id)
if result is None:
raise RuntimeError(f"Sandbox {sandbox_id} is not available")
agent_url, session_key = result
return ExecutionContext(
agent_url=agent_url,
session_key=session_key,
sandbox_id=sandbox_id,
api_url=self.api_url,
api_key=api_key,
)

async def release_context(
self, client: httpx.AsyncClient, ctx: ExecutionContext
) -> None:
Expand Down
7 changes: 7 additions & 0 deletions openhands/automation/backends/local.py
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,13 @@ async def get_execution_context(
sandbox_id=None, # No sandbox in local mode
)

async def get_existing_execution_context(
self,
client: httpx.AsyncClient, # noqa: ARG002
) -> ExecutionContext:
"""Return the persistent local agent server context."""
return await self.get_execution_context(client)

async def release_context(
self,
client: httpx.AsyncClient, # noqa: ARG002
Expand Down
Loading
Loading