Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
46be7e6
feat(chat): render sender from shared session events
mmabrouk Sep 4, 2026
e22c629
feat(sessions): detach shared sender execution
mmabrouk Sep 4, 2026
389d887
feat(chat): reconnect sender to running turn
mmabrouk Sep 4, 2026
3011ea0
feat(mobile): follow shared sender execution
mmabrouk Sep 4, 2026
6445d5d
fix(chat): hide shared invoke acceptance rows
mmabrouk Sep 4, 2026
2006dfa
fix(sessions): apply shared sender review
mmabrouk Sep 4, 2026
64cc80b
fix(chat): clear the running strip when the turn ends after a reload
mmabrouk Sep 4, 2026
73d74a5
fix(chat): stop persisting a dead sender stream as a failed turn
mmabrouk Sep 4, 2026
4a23a87
fix(chat): surface accepted disconnect as connection state
mmabrouk Sep 4, 2026
e14b4dd
docs(sessions): remove stale flag consolidation note
mmabrouk Sep 4, 2026
3a59296
chore(ci): satisfy readiness checks
mmabrouk Sep 4, 2026
8716021
test(sessions): align shared sender contracts
mmabrouk Sep 4, 2026
539b722
fix(chat): keep accepted runs pending after disconnect
mmabrouk Sep 5, 2026
b7d60ec
fix(chat): latch one response source per turn
mmabrouk Sep 5, 2026
dde620a
fix(chat): preserve shared runner error provenance
mmabrouk Sep 5, 2026
15f3542
fix(chat): reset turn acceptance before request build
mmabrouk Sep 5, 2026
11451d7
test(chat): align shared sender fixtures after rebase
mmabrouk Sep 5, 2026
7ddea9e
fix(chat): show activity for connected shared readers
mmabrouk Sep 5, 2026
25e8c73
fix(chat): sync watched interactions live
mmabrouk Sep 5, 2026
d4bfc1b
fix(chat): bound shared sender acceptance
mmabrouk Sep 5, 2026
703054a
fix(runner): batch live session frames
mmabrouk Sep 5, 2026
1e9e748
fix(api): protect runner record ingest from plan throttling
mmabrouk Sep 5, 2026
754b2d5
fix(chat): restore flag-off running banner
mmabrouk Sep 5, 2026
fe2e5d1
chore(ci): refresh shared sender fixture fingerprint
mmabrouk Sep 5, 2026
28cb5da
docs(sessions): clarify live event retention and cursors
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
2 changes: 2 additions & 0 deletions .gitleaksignore
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,8 @@ ad74134f522cde71f860cb59b6363a8fdf0a64c6:ee/setup_agenta_web.sh:generic-api-key:
590578c803d94d8ccb1a6ca977471f3d44b43fc3:hosting/helm/oss/templates/config/app-configmap.yaml:generic-api-key:45
1d8f08b267675726441fcaaae24572bb635c5eac:api/oss/src/utils/env.py:generic-api-key:53
55f27e52327062382beb299b162f94895268d766:web/oss/public/__ENV.js:generic-api-key:1
012ae6318c880d19944ffe1ff51740da0613cd8a:services/runner/tests/unit/server.test.ts:generic-api-key:586
e22c629b3113e522f4fdd918a9b0971e62145056:services/runner/tests/unit/server.test.ts:generic-api-key:889
c98a5da1a33d2c0986e3c66329eaa5237fbccf3d:hosting/docker-compose/ee/aws/docker-compose.oss.prod.yml:generic-api-key:73
bf0cd42bffc2581b1df6f56fa6e4b20ff9b68c33:hosting/docker-compose/ee/aws/docker-compose.oss.aws.yml:generic-api-key:61
52cd40cefd3121eea2e21205e8208712b093529a:core/hosting/docker-compose/ee/docker-compose.dev.yml:generic-api-key:18
Expand Down
18 changes: 14 additions & 4 deletions api/ee/src/middlewares/throttling.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

from oss.src.utils.caching import get_cache, set_cache
from oss.src.utils.logging import get_module_logger
from oss.src.middlewares.auth import SECRET_RESOLVE_GRANT, request_has_grant
from oss.src.utils.throttling import Algorithm, check_throttles

from ee.src.core.access.entitlements.types import (
Expand Down Expand Up @@ -40,6 +41,14 @@ def _normalize_path(request: Request) -> str:
return path


def _is_runner_record_ingest(request: Request, method: str, path: str) -> bool:
return (
method == Method.POST.value
and path == "/sessions/records/ingest"
and request_has_grant(request, SECRET_RESOLVE_GRANT)
)


def _matches_endpoint(
method: str,
path: str,
Expand Down Expand Up @@ -168,6 +177,11 @@ async def throttling_middleware(request: Request, call_next):
if hasattr(request.state, "admin") and request.state.admin:
return await call_next(request)

method = request.method.lower()
path = _normalize_path(request)
if _is_runner_record_ingest(request, method, path):
return await call_next(request)

organization_id = (
request.state.organization_id
if hasattr(request.state, "organization_id")
Expand Down Expand Up @@ -221,10 +235,6 @@ async def throttling_middleware(request: Request, call_next):
if not throttles:
return await call_next(request)

method = request.method.lower()

path = _normalize_path(request)

# log.debug(
# "[throttling] START", org=organization_id, plan=plan, method=method, path=path
# )
Expand Down
101 changes: 101 additions & 0 deletions api/ee/tests/pytest/unit/test_throttling.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
from types import SimpleNamespace
from unittest.mock import AsyncMock, patch

import pytest
from fastapi import Request, Response

from ee.src.core.access.entitlements.types import (
Bucket,
Category,
Mode,
Throttle,
Tracker,
)
from ee.src.middlewares.throttling import throttling_middleware
from oss.src.middlewares.auth import SECRET_RESOLVE_GRANT


def _request(path: str, *, grants: tuple[str, ...] = ()) -> Request:
request = Request(
{
"type": "http",
"method": "POST",
"path": path,
"root_path": "/api" if path.startswith("/api/") else "",
"headers": [],
}
)
request.state.organization_id = "organization-1"
request.state.token_grants = grants
return request


async def test_runner_record_ingest_bypasses_plan_throttle():
request = _request(
"/api/sessions/records/ingest",
grants=(SECRET_RESOLVE_GRANT,),
)
call_next = AsyncMock(return_value=Response(status_code=204))

with (
patch(
"ee.src.middlewares.throttling._get_plan", new_callable=AsyncMock
) as get_plan,
patch(
"ee.src.middlewares.throttling.check_throttles",
new_callable=AsyncMock,
) as check_throttles,
):
response = await throttling_middleware(request, call_next)

assert response.status_code == 204
call_next.assert_awaited_once_with(request)
get_plan.assert_not_awaited()
check_throttles.assert_not_awaited()


@pytest.mark.parametrize(
("path", "grants"),
[
("/sessions/records/ingest", ()),
("/sessions/query", (SECRET_RESOLVE_GRANT,)),
],
)
async def test_throttle_still_counts_browser_ingest_and_other_runner_routes(
path: str,
grants: tuple[str, ...],
):
request = _request(path, grants=grants)
call_next = AsyncMock(return_value=Response(status_code=204))
standard = Throttle(
categories=[Category.STANDARD],
mode=Mode.INCLUDE,
bucket=Bucket(capacity=10, rate=10),
)
allowed = SimpleNamespace(
allow=True,
tokens_remaining=9,
retry_after_seconds=0,
)

with (
patch(
"ee.src.middlewares.throttling._get_plan",
new_callable=AsyncMock,
return_value="test-plan",
) as get_plan,
patch(
"ee.src.middlewares.throttling.get_plan_entitlements",
return_value={Tracker.THROTTLES: [standard]},
),
patch(
"ee.src.middlewares.throttling.check_throttles",
new_callable=AsyncMock,
return_value=[allowed],
) as check_throttles,
):
response = await throttling_middleware(request, call_next)

assert response.status_code == 204
get_plan.assert_awaited_once_with("organization-1")
check_throttles.assert_awaited_once()
12 changes: 11 additions & 1 deletion api/oss/src/apis/fastapi/sessions/models.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
from datetime import datetime
from typing import Annotated, Any, Dict, List, Literal, Optional
from typing import Annotated, Any, Dict, List, Literal, Optional, Union
from uuid import UUID

from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator
Expand Down Expand Up @@ -412,6 +412,16 @@ def validate_live_frame(self) -> "SessionRecordIngestRequest":
return self


SessionRecordIngestBatch = Annotated[
List[SessionRecordIngestRequest],
Field(min_length=1),
]
SessionRecordIngestBody = Union[
SessionRecordIngestRequest,
SessionRecordIngestBatch,
]


# ---------------------------------------------------------------------------
# Session control: durable commands (Stop)
# ---------------------------------------------------------------------------
Expand Down
66 changes: 44 additions & 22 deletions api/oss/src/apis/fastapi/sessions/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,7 @@
SessionStreamResponse,
SessionStreamsResponse,
# records
SessionRecordIngestRequest,
SessionRecordIngestBody,
SessionRecordQueryRequest,
SessionRecordResponse,
SessionRecordsQueryResponse,
Expand Down Expand Up @@ -912,7 +912,7 @@ async def get_record_event(
async def ingest_record_event(
self,
request: Request,
body: SessionRecordIngestRequest,
body: SessionRecordIngestBody,
) -> dict:
project_id = request.state.project_id
if not await check_action_access(
Expand All @@ -922,8 +922,28 @@ async def ingest_record_event(
):
raise FORBIDDEN_EXCEPTION

if body.kind == "frame":
_validate_session_id_http(body.session_id)
if isinstance(body, list):
if not body or any(item.kind != "frame" for item in body):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="Batched record ingest accepts live frames only.",
)
frames = body
else:
frames = [body] if body.kind == "frame" else []

if frames:
first = frames[0]
_validate_session_id_http(first.session_id)
if any(
frame.session_id != first.session_id
or frame.execution_id != first.execution_id
for frame in frames[1:]
):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="A live frame batch must share one session and execution.",
)
content_length = request.headers.get("content-length")
if content_length is not None:
try:
Expand All @@ -948,28 +968,30 @@ async def ingest_record_event(
current_execution_id = await get_running_owner(
get_lock_engine(),
project_id=str(project_id),
session_id=body.session_id,
session_id=first.session_id,
)
if current_execution_id != body.execution_id:
if current_execution_id != first.execution_id:
raise FORBIDDEN_EXCEPTION
await publish_live_frame(
organization_id=UUID(request.state.organization_id),
project_id=UUID(project_id),
frame=SessionLiveFrame(
version=body.version,
kind="frame",
session_id=body.session_id,
execution_id=body.execution_id,
frame_or_event_id=body.frame_or_event_id,
frame_index=body.frame_index,
entity_id=body.entity_id,
type=body.type,
payload=body.payload,
created_at=body.created_at,
),
)
for frame in frames:
await publish_live_frame(
organization_id=UUID(request.state.organization_id),
project_id=UUID(project_id),
frame=SessionLiveFrame(
version=frame.version,
kind="frame",
session_id=frame.session_id,
execution_id=frame.execution_id,
frame_or_event_id=frame.frame_or_event_id,
frame_index=frame.frame_index,
entity_id=frame.entity_id,
type=frame.type,
payload=frame.payload,
created_at=frame.created_at,
),
)
return {"ok": True}

assert not isinstance(body, list)
await publish_record(
organization_id=UUID(request.state.organization_id),
project_id=UUID(project_id),
Expand Down
17 changes: 17 additions & 0 deletions api/oss/src/core/sessions/records/dtos.py
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,11 @@ class ToolCompletedPayload(BaseModel):
status: str


class InteractionChangedPayload(BaseModel):
interaction_id: str
kind: Optional[str] = None


class SessionDurableEventBase(BaseModel):
"""Durable relay wire envelope.

Expand Down Expand Up @@ -166,6 +171,16 @@ class ToolCompletedEvent(SessionDurableEventBase):
payload: ToolCompletedPayload


class InteractionRequestedEvent(SessionDurableEventBase):
type: Literal["interaction.requested"]
payload: InteractionChangedPayload


class InteractionRespondedEvent(SessionDurableEventBase):
type: Literal["interaction.responded"]
payload: InteractionChangedPayload


SessionDurableEvent = Annotated[
Union[
ExecutionStartedEvent,
Expand All @@ -174,6 +189,8 @@ class ToolCompletedEvent(SessionDurableEventBase):
ExecutionLostEvent,
MessageCompletedEvent,
ToolCompletedEvent,
InteractionRequestedEvent,
InteractionRespondedEvent,
],
Field(discriminator="type"),
]
Expand Down
25 changes: 25 additions & 0 deletions api/oss/src/core/sessions/records/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@

from oss.src.core.sessions.records.dtos import (
SESSION_DURABLE_EVENT_TYPES,
InteractionRequestedEvent,
InteractionRespondedEvent,
MessageCompletedEvent,
SessionDurableEvent,
SessionRecord,
Expand Down Expand Up @@ -119,6 +121,29 @@ def durable_events_from_records(
if base is None:
continue

if record.record_type in {"interaction_request", "interaction_response"}:
payload = {
"interaction_id": entity_id,
"kind": attributes.get("kind"),
}
if record.record_type == "interaction_request":
events.append(
InteractionRequestedEvent(
**base,
type="interaction.requested",
payload=payload,
)
)
else:
events.append(
InteractionRespondedEvent(
**base,
type="interaction.responded",
payload=payload,
)
)
continue

if record.record_type == "message":
role = (
"assistant" if record.record_source == "agent" else record.record_source
Expand Down
3 changes: 1 addition & 2 deletions api/oss/src/utils/env.py
Original file line number Diff line number Diff line change
Expand Up @@ -1548,8 +1548,7 @@ class SessionsRedisConfig(BaseModel):
redis_contract.json) shared with the TypeScript runner. Do not change a default
without updating that fixture and the TS side in lockstep.

AGENTA_SESSIONS_SEQUENCE_WRITES folds into AGENTA_SESSIONS_HISTORY_WRITES after
PR #6517 and this increment merge, keeping one rollout switch per increment.
AGENTA_SESSIONS_SEQUENCE_WRITES independently gates atomic record sequencing.
"""

sequence_writes: bool = (
Expand Down
Loading
Loading