Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
84 commits
Select commit Hold shift + click to select a range
4d924d0
feat(runner): relay live session frames
mmabrouk Sep 4, 2026
7424913
feat(api): ingest temporary session frames
mmabrouk Sep 4, 2026
2371bc6
feat(api): relay live session frames
mmabrouk Sep 4, 2026
e541dab
feat(web): preview live session frames
mmabrouk Sep 4, 2026
c3745b8
fix(api): protect durable records from frame trimming
mmabrouk Sep 4, 2026
9117708
fix(web): bound live session preview state
mmabrouk Sep 4, 2026
fbd06c3
fix(agent): wire live session relay
mmabrouk Sep 4, 2026
1dd3fd8
fix(api): isolate disposable live frames
mmabrouk Sep 4, 2026
054b3e9
fix(api): bound live frame ingress
mmabrouk Sep 4, 2026
1be4a25
fix(api): prevent live event caching
mmabrouk Sep 4, 2026
fa73cc1
docs(sessions): separate live frame retention
mmabrouk Sep 4, 2026
3e6b17f
fix(frontend): log invalid live frames
mmabrouk Sep 4, 2026
29bb4de
fix(frontend): reject gapped live previews
mmabrouk Sep 4, 2026
ff55b05
fix(api): restore durable record retention
mmabrouk Sep 4, 2026
27b4caf
fix(api): isolate live relay startup
mmabrouk Sep 4, 2026
ecb2214
docs(sessions): separate shipped and target events
mmabrouk Sep 4, 2026
6e6df19
feat(sessions): add durable sequence allocation
mmabrouk Sep 4, 2026
ed8eb9e
feat(sessions): relay committed durable events
mmabrouk Sep 4, 2026
1fa2449
feat(sessions): add reconnect snapshot and transcript paging
mmabrouk Sep 4, 2026
2f3bf88
feat(sessions): replay durable events before live follow
mmabrouk Sep 4, 2026
282d616
feat(sessions): reconnect from durable sequence cursors
mmabrouk Sep 4, 2026
455368b
fix(sessions): ignore untyped relay results
mmabrouk Sep 4, 2026
374bb2e
fix(sessions): sequence migrations by project
mmabrouk Sep 4, 2026
6bd3feb
fix(sessions): accept sparse durable event sequences
mmabrouk Sep 4, 2026
64c4824
fix(sessions): avoid watermark-only transcript refetches
mmabrouk Sep 4, 2026
bcd4f6b
test(sessions): declare the analytics engine fixture with pytest_asyncio
mmabrouk Sep 4, 2026
a7cfb8a
refactor(chat): stop retaining durable event payloads
mmabrouk Sep 4, 2026
c151273
fix(sessions): backpressure durable replay
mmabrouk Sep 4, 2026
636c916
fix(sessions): isolate durable relay envelopes
mmabrouk Sep 5, 2026
a19e671
fix(sessions): order cursor locks before acknowledgement
mmabrouk Sep 5, 2026
63b3343
fix(chat): hydrate through the reconnect cursor
mmabrouk Sep 5, 2026
e9f10a8
fix(chat): back off live reconnects
mmabrouk Sep 5, 2026
ef23286
test(migrations): verify tracing sequence parent
mmabrouk Sep 5, 2026
f5e7948
fix(chat): gate live follow on transcript adoption
mmabrouk Sep 5, 2026
4fdd89b
fix(chat): join interaction state into live replay
mmabrouk Sep 5, 2026
c1522b3
fix(sessions): exclude quarantined records from reconnect
mmabrouk Sep 5, 2026
04d3795
fix(chat): return desktop transcript adoption status
mmabrouk Sep 5, 2026
1969ec6
fix(chat): separate reconnect sequence and row watermarks
mmabrouk Sep 5, 2026
1281c60
fix(chat): guard transcript adoption failures
mmabrouk Sep 5, 2026
f600abb
style(chat): format live preview retry fixture
mmabrouk Sep 5, 2026
c672dc5
fix(sessions): add cursor lifecycle columns
mmabrouk Sep 5, 2026
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
0cb6549
Merge pull request #6522 from Agenta-AI/feat/session-live-relay
mmabrouk Sep 5, 2026
c2958e3
Merge pull request #6524 from Agenta-AI/feat/session-durable-reconnect
mmabrouk Sep 5, 2026
c8721c2
Merge pull request #6531 from Agenta-AI/feat/session-shared-sender
mmabrouk Sep 5, 2026
1a82e38
fix(api): treat a non-dict record payload as absent
mmabrouk Sep 5, 2026
ce34eac
perf(api): trim the live-frame stream approximately
mmabrouk Sep 5, 2026
ee0c74d
test(api): close a leftover analytics engine in the sequence fixture
mmabrouk Sep 5, 2026
efa5de9
docs(api): record why the transcript window pages by offset
mmabrouk Sep 5, 2026
47e4ece
fix(runner): emit session-accepted only after the turn is admitted
mmabrouk Sep 5, 2026
2b15409
fix(runner): bound the live-frame ingest POST
mmabrouk Sep 5, 2026
6835357
fix(web): keep an accepted shared turn's transcript adoptable
mmabrouk Sep 5, 2026
1c09d79
test(web): give the SessionRecord fixtures their sequence field
mmabrouk Sep 5, 2026
4d51858
test(sdk,services): carry detached through the test doubles
mmabrouk Sep 5, 2026
a354ce4
docs: describe the durable-event contract as shipped
mmabrouk Sep 5, 2026
6bc2410
docs(sessions): match the shipped live relay topology
mmabrouk Sep 5, 2026
4eb4446
Merge pull request #6574 from Agenta-AI/fix/session-live-events-coder…
mmabrouk Sep 5, 2026
8256a80
fix(migrations): build the records sequence index concurrently
mmabrouk Sep 5, 2026
c47d89e
Merge pull request #6576 from Agenta-AI/fix/session-live-events-index
mmabrouk Sep 5, 2026
64dc2b9
Merge release/v0.115.1 into feat/session-live-events
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()
55 changes: 42 additions & 13 deletions api/entrypoints/worker_streams.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@
from oss.src.core.secrets.services import VaultService
from oss.src.core.sessions.interactions.service import SessionInteractionsService
from oss.src.core.sessions.records.service import RecordsService
from oss.src.core.sessions.records.streaming import (
LIVE_FRAME_STREAM_NAME,
RECORD_STREAM_NAME,
)
from oss.src.core.tracing.service import TracingService
from oss.src.dbs.postgres.events.dao import EventsDAO
from oss.src.dbs.postgres.secrets.dao import SecretsDAO
Expand All @@ -37,6 +41,7 @@
from oss.src.dbs.redis.sessions.watch import SessionsWatchPublisher
from oss.src.tasks.asyncio.events.worker import EventsWorker
from oss.src.tasks.asyncio.sessions.records_worker import RecordsWorker
from oss.src.tasks.asyncio.sessions.live_relay_worker import LiveRelayWorker
from oss.src.tasks.asyncio.shared.consumer import StreamConsumer
from oss.src.tasks.asyncio.tracing.worker import TracingWorker
from oss.src.tasks.asyncio.webhooks.dispatcher import WebhooksDispatcher
Expand Down Expand Up @@ -85,7 +90,7 @@ async def _build_records_worker(redis_client: Redis) -> StreamConsumer:
return RecordsWorker(
service=RecordsService(records_dao=RecordsDAO()),
redis_client=redis_client,
stream_name="streams:records",
stream_name=RECORD_STREAM_NAME,
consumer_group="worker-records",
# M3 live relay: post-append change notifications on the durable plane,
# reusing this process's durable connection.
Expand All @@ -103,6 +108,14 @@ async def _build_records_worker(redis_client: Redis) -> StreamConsumer:
)


async def _build_live_relay_worker(redis_client: Redis) -> StreamConsumer:
return LiveRelayWorker(
redis_client=redis_client,
stream_name=LIVE_FRAME_STREAM_NAME,
consumer_group="worker-session-live-relay",
)


async def _build_events_worker(redis_client: Redis) -> StreamConsumer:
events_service = EventsService(events_dao=EventsDAO())

Expand Down Expand Up @@ -136,6 +149,22 @@ async def _build_events_worker(redis_client: Redis) -> StreamConsumer:
)


async def _initialize_consumer(consumer: StreamConsumer) -> None:
await consumer.create_consumer_group()
removed = await prune_idle_consumers(
url=env.redis.uri_durable,
queue_name=consumer.stream_name,
consumer_group_name=consumer.consumer_group,
keep=consumer.consumer_name,
)
if removed:
log.info(
"[STREAMS] Pruned idle consumers",
stream=consumer.stream_name,
removed=removed,
)


async def main_async() -> int:
try:
streams = _selected_streams()
Expand Down Expand Up @@ -165,19 +194,19 @@ async def main_async() -> int:
]

for consumer in consumers:
await consumer.create_consumer_group()
removed = await prune_idle_consumers(
url=env.redis.uri_durable,
queue_name=consumer.stream_name,
consumer_group_name=consumer.consumer_group,
keep=consumer.consumer_name,
)
if removed:
log.info(
"[STREAMS] Pruned idle consumers",
stream=consumer.stream_name,
removed=removed,
await _initialize_consumer(consumer)

if env.sessions.shared_reader and "records" in streams:
try:
live_relay = await _build_live_relay_worker(redis_client)
await _initialize_consumer(live_relay)
except Exception:
log.error(
"[STREAMS] Live relay disabled after initialization failure",
exc_info=True,
)
else:
consumers.append(live_relay)

log.info("[STREAMS] Starting worker-streams", selected=streams)

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
"""add session sequence cursors

Revision ID: oss000000006
Revises: oss000000005
Create Date: 2026-09-04 00:00:00.000000

"""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa


revision: str = "oss000000006"
down_revision: Union[str, None] = "oss000000005"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.create_table(
"session_sequence_cursors",
sa.Column("project_id", sa.UUID(), nullable=False),
sa.Column("session_id", sa.String(), nullable=False),
sa.Column("latest_sequence", sa.BigInteger(), nullable=False),
sa.Column(
"created_at",
sa.TIMESTAMP(timezone=True),
server_default=sa.text("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(), nullable=True),
sa.Column("updated_by_id", sa.UUID(), nullable=True),
sa.Column("deleted_by_id", sa.UUID(), nullable=True),
sa.PrimaryKeyConstraint("project_id", "session_id"),
)
op.add_column("records", sa.Column("sequence", sa.BigInteger(), nullable=True))
with op.get_context().autocommit_block():
op.create_index(
"ux_records_session_id_sequence",
"records",
["project_id", "session_id", "sequence"],
unique=True,
postgresql_concurrently=True,
)


def downgrade() -> None:
with op.get_context().autocommit_block():
op.drop_index(
"ux_records_session_id_sequence",
table_name="records",
postgresql_concurrently=True,
)
op.drop_column("records", "sequence")
op.drop_table("session_sequence_cursors")
Loading
Loading