From ba71600d6fe061f42aa0c6da85efae0cf3338dd5 Mon Sep 17 00:00:00 2001 From: steven-ji Date: Sat, 22 Aug 2026 12:08:13 +0800 Subject: [PATCH 1/3] feat(sessions): expose last message activity Persist and backfill last_message_at, support activity-based session sorting, and expose the field through both SDKs. Refs #965 --- CHANGELOG.md | 6 + docs/v3/openapi.json | 20 ++ ...faff339d519_add_session_last_message_at.py | 61 ++++++ sdks/python/CHANGELOG.md | 6 + sdks/python/src/honcho/aio.py | 12 ++ sdks/python/src/honcho/api_types.py | 1 + sdks/python/src/honcho/client.py | 8 + sdks/python/src/honcho/peer.py | 1 + sdks/python/src/honcho/scope.py | 1 + sdks/python/src/honcho/session.py | 15 ++ sdks/typescript/CHANGELOG.md | 6 + sdks/typescript/__tests__/helpers.ts | 7 + sdks/typescript/__tests__/peer.test.ts | 3 + sdks/typescript/__tests__/session.test.ts | 52 +++++ .../typescript/__tests__/session.unit.test.ts | 3 + sdks/typescript/src/client.ts | 18 +- sdks/typescript/src/peer.ts | 3 +- sdks/typescript/src/scope.ts | 3 +- sdks/typescript/src/session.ts | 17 +- sdks/typescript/src/types/api.ts | 1 + sdks/typescript/src/validation.ts | 3 +- src/crud/message.py | 38 +++- src/crud/session.py | 31 ++- src/models.py | 10 + src/routers/sessions.py | 5 + src/schemas/api.py | 1 + tests/alembic/revisions/__init__.py | 2 + ...faff339d519_add_session_last_message_at.py | 118 +++++++++++ tests/routes/test_sessions.py | 190 ++++++++++++++++++ tests/sdk/test_peer.py | 8 + tests/sdk/test_session.py | 107 ++++++++++ 31 files changed, 741 insertions(+), 16 deletions(-) create mode 100644 migrations/versions/cfaff339d519_add_session_last_message_at.py create mode 100644 tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py diff --git a/CHANGELOG.md b/CHANGELOG.md index e314f867..cbb1fe9d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](http://keepachangelog.com/) and this project adheres to [Semantic Versioning](http://semver.org/). +## [Unreleased] + +### Added + +- Session responses now expose nullable `last_message_at`, backfilled and maintained from the newest message timestamp. `POST /v3/workspaces/{workspace_id}/sessions/list` accepts `sort_by=created_at|last_message_at` alongside the existing `reverse` parameter, with stable ID tie-breaking and sessions without messages placed last in either direction (#965). + ## [3.0.12] - 2026-08-10 ### Added diff --git a/docs/v3/openapi.json b/docs/v3/openapi.json index 7fa1d194..0c49a55c 100644 --- a/docs/v3/openapi.json +++ b/docs/v3/openapi.json @@ -1095,6 +1095,19 @@ }, "description": "Whether to reverse the order of results" }, + { + "name": "sort_by", + "in": "query", + "required": false, + "schema": { + "enum": ["created_at", "last_message_at"], + "type": "string", + "description": "Session timestamp used to order results", + "default": "created_at", + "title": "Sort By" + }, + "description": "Session timestamp used to order results" + }, { "name": "page", "in": "query", @@ -3561,6 +3574,13 @@ "type": "string", "format": "date-time", "title": "Created At" + }, + "last_message_at": { + "anyOf": [ + { "type": "string", "format": "date-time" }, + { "type": "null" } + ], + "title": "Last Message At" } }, "type": "object", diff --git a/migrations/versions/cfaff339d519_add_session_last_message_at.py b/migrations/versions/cfaff339d519_add_session_last_message_at.py new file mode 100644 index 00000000..d4f011ef --- /dev/null +++ b/migrations/versions/cfaff339d519_add_session_last_message_at.py @@ -0,0 +1,61 @@ +"""add session last_message_at + +Revision ID: cfaff339d519 +Revises: e4eba9cfaa6f +Create Date: 2026-08-22 + +""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op + +from migrations.utils import get_schema + +# revision identifiers, used by Alembic. +revision: str = "cfaff339d519" +down_revision: str | None = "e4eba9cfaa6f" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + +schema = get_schema() +INDEX_NAME = "ix_sessions_workspace_last_message_at" + + +def upgrade() -> None: + """Add the nullable session activity timestamp.""" + op.add_column( + "sessions", + sa.Column("last_message_at", sa.DateTime(timezone=True), nullable=True), + schema=schema, + ) + op.execute( + sa.text( + f""" + UPDATE "{schema}"."sessions" AS session + SET last_message_at = activity.last_message_at + FROM ( + SELECT workspace_name, session_name, MAX(created_at) AS last_message_at + FROM "{schema}"."messages" + GROUP BY workspace_name, session_name + ) AS activity + WHERE session.workspace_name = activity.workspace_name + AND session.name = activity.session_name + """ + ) + ) + op.create_index( + INDEX_NAME, + "sessions", + ["workspace_name", "last_message_at", "id"], + unique=False, + schema=schema, + postgresql_where=sa.text("is_active"), + ) + + +def downgrade() -> None: + """Remove the session activity timestamp.""" + op.drop_index(INDEX_NAME, table_name="sessions", schema=schema) + op.drop_column("sessions", "last_message_at", schema=schema) diff --git a/sdks/python/CHANGELOG.md b/sdks/python/CHANGELOG.md index 751d9005..c6c3af53 100644 --- a/sdks/python/CHANGELOG.md +++ b/sdks/python/CHANGELOG.md @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](http://keepachangelog.com/) and this project adheres to [Semantic Versioning](http://semver.org/). +## [Unreleased] + +### Added + +- `Session.last_message_at` exposes the newest message timestamp, and sync/async `Honcho.sessions()` accept `sort_by="created_at" | "last_message_at"` while preserving `reverse` across pagination. Requires a Honcho server with the matching API support. + ## [2.3.0] - 2026-08-10 ### Added diff --git a/sdks/python/src/honcho/aio.py b/sdks/python/src/honcho/aio.py index 29c0f445..f5a1011d 100644 --- a/sdks/python/src/honcho/aio.py +++ b/sdks/python/src/honcho/aio.py @@ -321,6 +321,7 @@ class HonchoAio(AsyncMetadataConfigMixin): session_data.configuration.model_dump() ), created_at=session_data.created_at, + last_message_at=session_data.last_message_at, is_active=session_data.is_active, ) @@ -331,6 +332,7 @@ class HonchoAio(AsyncMetadataConfigMixin): page: int = 1, size: int = 50, reverse: bool = False, + sort_by: Literal["created_at", "last_message_at"] = "created_at", ) -> AsyncPage[SessionResponse, Session]: """ Get all sessions in the current workspace asynchronously. @@ -340,11 +342,14 @@ class HonchoAio(AsyncMetadataConfigMixin): page: Page number (1-indexed). Default: 1. size: Number of items per page. Default: 50. reverse: If True, reverses the default ordering. Default: False. + sort_by: Session timestamp used for ordering. """ await self._honcho._ensure_workspace_async() query: dict[str, Any] = {"page": page, "size": size} if reverse: query["reverse"] = "true" + if sort_by != "created_at": + query["sort_by"] = sort_by data = await self._honcho._async_http_client.post( routes.sessions_list(self._honcho.workspace_id), body={"filters": filters} if filters else None, @@ -359,6 +364,7 @@ class HonchoAio(AsyncMetadataConfigMixin): metadata=session.metadata, configuration=session.configuration, created_at=session.created_at, + last_message_at=session.last_message_at, is_active=session.is_active, ) @@ -367,6 +373,8 @@ class HonchoAio(AsyncMetadataConfigMixin): next_query: dict[str, Any] = {"page": next_page, "size": size} if reverse: next_query["reverse"] = "true" + if sort_by != "created_at": + next_query["sort_by"] = sort_by next_data = await self._honcho._async_http_client.post( routes.sessions_list(self._honcho.workspace_id), body={"filters": filters} if filters else None, @@ -849,6 +857,7 @@ class PeerAio(AsyncMetadataConfigMixin): session.configuration.model_dump() ), created_at=session.created_at, + last_message_at=session.last_message_at, is_active=session.is_active, ) @@ -1091,6 +1100,7 @@ class SessionAio(AsyncMetadataConfigMixin): session.configuration.model_dump() ) self._session._created_at = session.created_at + self._session._last_message_at = session.last_message_at self._session._is_active = session.is_active async def get_metadata(self) -> dict[str, object]: @@ -1311,6 +1321,7 @@ class SessionAio(AsyncMetadataConfigMixin): metadata=cloned.metadata, configuration=cloned.configuration, created_at=cloned.created_at, + last_message_at=cloned.last_message_at, is_active=cloned.is_active, ) @@ -1904,6 +1915,7 @@ class ScopeAio: metadata=response.metadata, configuration=response.configuration, created_at=response.created_at, + last_message_at=response.last_message_at, is_active=response.is_active, ) diff --git a/sdks/python/src/honcho/api_types.py b/sdks/python/src/honcho/api_types.py index 23692f2f..06d75f3d 100644 --- a/sdks/python/src/honcho/api_types.py +++ b/sdks/python/src/honcho/api_types.py @@ -265,6 +265,7 @@ class SessionResponse(BaseModel): default_factory=SessionConfigurationResponse ) created_at: datetime.datetime + last_message_at: datetime.datetime | None = None class SessionCreateParams(BaseModel): diff --git a/sdks/python/src/honcho/client.py b/sdks/python/src/honcho/client.py index 1527792e..f38b3e90 100644 --- a/sdks/python/src/honcho/client.py +++ b/sdks/python/src/honcho/client.py @@ -469,6 +469,7 @@ class Honcho(BaseModel, MetadataConfigMixin): # pyright: ignore[reportUnsafeMul session_data.configuration.model_dump() ), created_at=session_data.created_at, + last_message_at=session_data.last_message_at, is_active=session_data.is_active, ) @@ -479,6 +480,7 @@ class Honcho(BaseModel, MetadataConfigMixin): # pyright: ignore[reportUnsafeMul page: int = 1, size: int = 50, reverse: bool = False, + sort_by: Literal["created_at", "last_message_at"] = "created_at", ) -> SyncPage[SessionResponse, Session]: """ Get all sessions in the current workspace. @@ -488,6 +490,7 @@ class Honcho(BaseModel, MetadataConfigMixin): # pyright: ignore[reportUnsafeMul page: Page number (1-indexed). Default: 1. size: Number of items per page. Default: 50. reverse: If True, reverses the default ordering. Default: False. + sort_by: Session timestamp used for ordering. Returns: A SyncPage of Session objects representing all sessions in the workspace. @@ -496,6 +499,8 @@ class Honcho(BaseModel, MetadataConfigMixin): # pyright: ignore[reportUnsafeMul query: dict[str, Any] = {"page": page, "size": size} if reverse: query["reverse"] = "true" + if sort_by != "created_at": + query["sort_by"] = sort_by data = self._http.post( routes.sessions_list(self.workspace_id), body={"filters": filters} if filters else None, @@ -510,6 +515,7 @@ class Honcho(BaseModel, MetadataConfigMixin): # pyright: ignore[reportUnsafeMul metadata=session.metadata, configuration=session.configuration, created_at=session.created_at, + last_message_at=session.last_message_at, is_active=session.is_active, ) @@ -518,6 +524,8 @@ class Honcho(BaseModel, MetadataConfigMixin): # pyright: ignore[reportUnsafeMul next_query: dict[str, Any] = {"page": next_page, "size": size} if reverse: next_query["reverse"] = "true" + if sort_by != "created_at": + next_query["sort_by"] = sort_by next_data = self._http.post( routes.sessions_list(self.workspace_id), body={"filters": filters} if filters else None, diff --git a/sdks/python/src/honcho/peer.py b/sdks/python/src/honcho/peer.py index 38edf269..60dcf926 100644 --- a/sdks/python/src/honcho/peer.py +++ b/sdks/python/src/honcho/peer.py @@ -469,6 +469,7 @@ class Peer(PeerBase, MetadataConfigMixin): session.configuration.model_dump() ), created_at=session.created_at, + last_message_at=session.last_message_at, is_active=session.is_active, ) diff --git a/sdks/python/src/honcho/scope.py b/sdks/python/src/honcho/scope.py index 14a5752b..a93eb8ea 100644 --- a/sdks/python/src/honcho/scope.py +++ b/sdks/python/src/honcho/scope.py @@ -202,6 +202,7 @@ class Scope(ScopeBase): metadata=response.metadata, configuration=response.configuration, created_at=response.created_at, + last_message_at=response.last_message_at, is_active=response.is_active, ) diff --git a/sdks/python/src/honcho/session.py b/sdks/python/src/honcho/session.py index 2f1838c9..d94972a1 100644 --- a/sdks/python/src/honcho/session.py +++ b/sdks/python/src/honcho/session.py @@ -60,11 +60,13 @@ class Session(SessionBase, MetadataConfigMixin): fetched. Call get_metadata() for fresh data. configuration: Cached configuration for this session. May be stale if not recently fetched. Call get_configuration() for fresh data. + last_message_at: Cached timestamp of the newest message, or None when empty. """ _metadata: dict[str, object] | None = PrivateAttr(default=None) _configuration: SessionConfiguration | None = PrivateAttr(default=None) _created_at: datetime | None = PrivateAttr(default=None) + _last_message_at: datetime | None = PrivateAttr(default=None) _is_active: bool | None = PrivateAttr(default=None) _honcho: "Honcho" = PrivateAttr() @@ -83,6 +85,11 @@ class Session(SessionBase, MetadataConfigMixin): """Timestamp when this session was created. Only available if fetched from the API.""" return self._created_at + @property + def last_message_at(self) -> datetime | None: + """Timestamp of the newest message, or None when the session is empty.""" + return self._last_message_at + @property def is_active(self) -> bool | None: """Whether this session is active. Only available if fetched from the API.""" @@ -117,6 +124,7 @@ class Session(SessionBase, MetadataConfigMixin): session.configuration.model_dump() ) self._created_at = session.created_at + self._last_message_at = session.last_message_at self._is_active = session.is_active def get_metadata(self) -> dict[str, object]: @@ -223,6 +231,10 @@ class Session(SessionBase, MetadataConfigMixin): None, description="Timestamp when this session was created.", ), + last_message_at: datetime | None = Field( + None, + description="Timestamp of the newest message in this session.", + ), is_active: bool | None = Field( None, description="Whether this session is active.", @@ -241,6 +253,7 @@ class Session(SessionBase, MetadataConfigMixin): If set, will get/create session immediately with metadata. configuration: Optional configuration to set for this session. If set, will get/create session immediately with flags. + last_message_at: Timestamp of the newest message, if fetched. """ super().__init__( id=session_id, @@ -250,6 +263,7 @@ class Session(SessionBase, MetadataConfigMixin): self._metadata = metadata self._configuration = configuration # pyright: ignore[reportIncompatibleVariableOverride] self._created_at = created_at + self._last_message_at = last_message_at self._is_active = is_active def add_peers( @@ -553,6 +567,7 @@ class Session(SessionBase, MetadataConfigMixin): metadata=cloned.metadata, configuration=cloned.configuration, created_at=cloned.created_at, + last_message_at=cloned.last_message_at, is_active=cloned.is_active, ) diff --git a/sdks/typescript/CHANGELOG.md b/sdks/typescript/CHANGELOG.md index 8d0e5ec1..4e314f73 100644 --- a/sdks/typescript/CHANGELOG.md +++ b/sdks/typescript/CHANGELOG.md @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](http://keepachangelog.com/) and this project adheres to [Semantic Versioning](http://semver.org/). +## [Unreleased] + +### Added + +- `Session.lastMessageAt` exposes the newest message timestamp, and `Honcho.sessions()` accepts `sortBy: 'created_at' | 'last_message_at'` while preserving `reverse` across pagination. Requires a Honcho server with the matching API support. + ## [2.3.0] - 2026-08-10 ### Added diff --git a/sdks/typescript/__tests__/helpers.ts b/sdks/typescript/__tests__/helpers.ts index 83cefe24..d4a0b4cb 100644 --- a/sdks/typescript/__tests__/helpers.ts +++ b/sdks/typescript/__tests__/helpers.ts @@ -63,6 +63,13 @@ export function assertSessionShape(session: SessionResponse): void { expect(typeof session.configuration).toBe('object') expect(typeof session.created_at).toBe('string') expectValidDateString(session.created_at) + expect( + session.last_message_at === null || + typeof session.last_message_at === 'string' + ).toBe(true) + if (typeof session.last_message_at === 'string') { + expectValidDateString(session.last_message_at) + } } /** diff --git a/sdks/typescript/__tests__/peer.test.ts b/sdks/typescript/__tests__/peer.test.ts index cde6a0ad..5d0c89da 100644 --- a/sdks/typescript/__tests__/peer.test.ts +++ b/sdks/typescript/__tests__/peer.test.ts @@ -214,12 +214,15 @@ describe('Peer', () => { const session = await client.session('peer-sessions-test', { metadata: {} }) await session.addPeers([peer.id]) + await session.addMessages(peer.message('peer session activity')) const sessions = await peer.sessions() expect(sessions.items.length).toBeGreaterThanOrEqual(1) const sessionIds = sessions.items.map((s) => s.id) expect(sessionIds).toContain('peer-sessions-test') + const returned = sessions.items.find((item) => item.id === session.id) + expect(returned?.lastMessageAt).toBeDefined() }) test('sessions returns empty for peer in no sessions', async () => { diff --git a/sdks/typescript/__tests__/session.test.ts b/sdks/typescript/__tests__/session.test.ts index f7f2f838..b3cd0d34 100644 --- a/sdks/typescript/__tests__/session.test.ts +++ b/sdks/typescript/__tests__/session.test.ts @@ -55,6 +55,7 @@ describe('Session', () => { expect(session.id).toBe('simple-session') expect(session.workspaceId).toBe(client.workspaceId) expect(session.createdAt).toBeDefined() + expect(session.lastMessageAt).toBeNull() expect(session.isActive).toBe(true) }) @@ -171,6 +172,55 @@ describe('Session', () => { expect(page.items.length).toBe(1) expect(page.items[0].metadata.tag).toBe(tag) }) + + test('sessions sort by last message activity and keep empty sessions last', async () => { + const activityGroup = generateId('last-activity-group') + const recentId = generateId('last-activity-recent') + const olderId = generateId('last-activity-older') + const emptyId = generateId('last-activity-empty') + const peer = await client.peer(generateId('last-activity-peer')) + const recentTime = new Date('2026-01-10T00:00:00Z') + const olderTime = new Date('2026-01-05T00:00:00Z') + + const older = await client.session(olderId, { + metadata: { activityGroup }, + peers: [peer], + }) + const recent = await client.session(recentId, { + metadata: { activityGroup }, + peers: [peer], + }) + await client.session(emptyId, { + metadata: { activityGroup }, + peers: [peer], + }) + await recent.addMessages( + peer.message('recent activity', { createdAt: recentTime }) + ) + await older.addMessages( + peer.message('older activity', { createdAt: olderTime }) + ) + + const page = await client.sessions({ + filters: { metadata: { activityGroup } }, + sortBy: 'last_message_at', + reverse: true, + size: 1, + }) + const sessions = await page.toArray() + + expect(sessions.map((session) => session.id)).toEqual([ + recentId, + olderId, + emptyId, + ]) + const returnedActivity = sessions[0].lastMessageAt + expect(typeof returnedActivity).toBe('string') + expect(new Date(returnedActivity ?? '').getTime()).toBe( + recentTime.getTime() + ) + expect(sessions[2].lastMessageAt).toBeNull() + }) }) // =========================================================================== @@ -247,6 +297,8 @@ describe('Session', () => { const cloned = await original.clone() expect(cloned.id).not.toBe(original.id) + expect(cloned.lastMessageAt).toBeDefined() + expect(cloned.lastMessageAt).not.toBeNull() // Cloned session should have same messages const messages = await cloned.messages() expect(messages.items.length).toBe(1) diff --git a/sdks/typescript/__tests__/session.unit.test.ts b/sdks/typescript/__tests__/session.unit.test.ts index 71a440e6..86558758 100644 --- a/sdks/typescript/__tests__/session.unit.test.ts +++ b/sdks/typescript/__tests__/session.unit.test.ts @@ -14,6 +14,7 @@ function createSessionResponse( metadata: {}, configuration: {}, created_at: '2024-01-01T00:00:00Z', + last_message_at: null, ...overrides, } } @@ -35,6 +36,7 @@ describe('Session unit behavior', () => { createSessionResponse({ metadata: { topic: 'testing' }, created_at: '2024-02-01T12:00:00Z', + last_message_at: '2024-02-02T12:00:00Z', is_active: true, }), } as unknown as HonchoHTTPClient @@ -48,6 +50,7 @@ describe('Session unit behavior', () => { expect(metadata).toEqual({ topic: 'testing' }) expect(session.createdAt).toBe('2024-02-01T12:00:00Z') + expect(session.lastMessageAt).toBe('2024-02-02T12:00:00Z') expect(session.isActive).toBe(true) }) diff --git a/sdks/typescript/src/client.ts b/sdks/typescript/src/client.ts index 2b66b9c2..a69debfc 100644 --- a/sdks/typescript/src/client.ts +++ b/sdks/typescript/src/client.ts @@ -359,6 +359,7 @@ export class Honcho { page?: number size?: number reverse?: boolean + sortBy?: 'created_at' | 'last_message_at' } ): Promise> { return this._http.post>( @@ -369,6 +370,10 @@ export class Honcho { page: params?.page, size: params?.size, reverse: params?.reverse ? 'true' : undefined, + sort_by: + params?.sortBy && params.sortBy !== 'created_at' + ? params.sortBy + : undefined, }, } ) @@ -593,7 +598,8 @@ export class Honcho { sessionConfigFromApi(sessionData.configuration) ?? undefined, () => this._ensureWorkspace(), sessionData.created_at, - sessionData.is_active + sessionData.is_active, + sessionData.last_message_at ) } @@ -688,7 +694,7 @@ export class Honcho { * the current workspace. * * @param options - Either a legacy raw filter object or an options object with - * `filters`, `page`, `size`, and `reverse`. See + * `filters`, `page`, `size`, `reverse`, and `sortBy`. See * [search filters documentation](https://honcho.dev/docs/v3/documentation/core-concepts/features/using-filters). * @returns Promise resolving to a Page of Session objects representing all sessions * in the workspace. Returns an empty page if no sessions exist @@ -701,6 +707,7 @@ export class Honcho { page?: number size?: number reverse?: boolean + sortBy?: 'created_at' | 'last_message_at' } ): Promise> { await this._ensureWorkspace() @@ -709,16 +716,19 @@ export class Honcho { 'page', 'size', 'reverse', + 'sortBy', ]) const validatedFilter = normalizedOptions.filters ? FilterSchema.parse(normalizedOptions.filters) : undefined const reverse = normalizedOptions.reverse + const sortBy = normalizedOptions.sortBy const sessionsPage = await this._listSessions(this.workspaceId, { filters: validatedFilter, page: normalizedOptions.page, size: normalizedOptions.size, reverse, + sortBy, }) const fetchNextPage = async ( @@ -730,6 +740,7 @@ export class Honcho { page, size, reverse, + sortBy, }) } @@ -744,7 +755,8 @@ export class Honcho { sessionConfigFromApi(session.configuration) ?? undefined, () => this._ensureWorkspace(), session.created_at, - session.is_active + session.is_active, + session.last_message_at ), fetchNextPage ) diff --git a/sdks/typescript/src/peer.ts b/sdks/typescript/src/peer.ts index 9dab1658..2762a4c4 100644 --- a/sdks/typescript/src/peer.ts +++ b/sdks/typescript/src/peer.ts @@ -634,7 +634,8 @@ export class Peer { sessionConfigFromApi(session.configuration) ?? undefined, () => this._ensureWorkspace(), session.created_at, - session.is_active + session.is_active, + session.last_message_at ), fetchNextPage ) diff --git a/sdks/typescript/src/scope.ts b/sdks/typescript/src/scope.ts index 5a1114fe..8e5192ba 100644 --- a/sdks/typescript/src/scope.ts +++ b/sdks/typescript/src/scope.ts @@ -239,7 +239,8 @@ export class Scope { sessionConfigFromApi(session.configuration) ?? undefined, () => this._ensureWorkspace(), session.created_at, - session.is_active + session.is_active, + session.last_message_at ), fetchNextPage ) diff --git a/sdks/typescript/src/session.ts b/sdks/typescript/src/session.ts index ed6140a6..5891b74f 100644 --- a/sdks/typescript/src/session.ts +++ b/sdks/typescript/src/session.ts @@ -87,6 +87,7 @@ export class Session { private _metadata?: Record private _configuration?: SessionConfig private _createdAt?: string + private _lastMessageAt?: string | null private _isActive?: boolean private _ensureWorkspace: () => Promise @@ -119,6 +120,13 @@ export class Session { return this._createdAt } + /** + * Timestamp of the newest message. Null when the fetched session is empty. + */ + get lastMessageAt(): string | null | undefined { + return this._lastMessageAt + } + /** * Whether this session is active. Only available if fetched from the API. */ @@ -134,6 +142,7 @@ export class Session { * @param http - Reference to the HTTP client instance * @param metadata - Optional metadata to initialize the cached value * @param configuration - Optional configuration to initialize the cached value + * @param lastMessageAt - Optional newest-message timestamp from the API */ constructor( id: string, @@ -143,7 +152,8 @@ export class Session { configuration?: SessionConfig, ensureWorkspace: () => Promise = async () => undefined, createdAt?: string, - isActive?: boolean + isActive?: boolean, + lastMessageAt?: string | null ) { this.id = id this.workspaceId = workspaceId @@ -153,12 +163,14 @@ export class Session { this._ensureWorkspace = ensureWorkspace this._createdAt = createdAt this._isActive = isActive + this._lastMessageAt = lastMessageAt } private _applySessionResponse(session: SessionResponse): void { this._metadata = session.metadata || {} this._configuration = sessionConfigFromApi(session.configuration) || {} this._createdAt = session.created_at + this._lastMessageAt = session.last_message_at this._isActive = session.is_active } @@ -741,7 +753,8 @@ export class Session { sessionConfigFromApi(clonedSessionData.configuration) ?? undefined, () => this._ensureWorkspace(), clonedSessionData.created_at, - clonedSessionData.is_active + clonedSessionData.is_active, + clonedSessionData.last_message_at ) } diff --git a/sdks/typescript/src/types/api.ts b/sdks/typescript/src/types/api.ts index 10ed6037..c473ef69 100644 --- a/sdks/typescript/src/types/api.ts +++ b/sdks/typescript/src/types/api.ts @@ -126,6 +126,7 @@ export interface SessionResponse { metadata: Record configuration: SessionConfigApi created_at: string + last_message_at?: string | null } export interface SessionCreateParams { diff --git a/sdks/typescript/src/validation.ts b/sdks/typescript/src/validation.ts index 832690e2..09e9cb98 100644 --- a/sdks/typescript/src/validation.ts +++ b/sdks/typescript/src/validation.ts @@ -407,7 +407,8 @@ export const FilterSchema = z.record(z.string(), z.unknown()).optional() * shape are both accepted. * * Discriminates on the `filters` key: if the input has a `filters` property or - * any of the pagination-only keys (`page`, `size`, `reverse`) it is treated as + * any of the supplied option-only keys (for example `page`, `reverse`, or + * `sortBy`) it is treated as * the new options object. Otherwise it is treated as a legacy raw filter. */ export function normalizeListOptions( diff --git a/src/crud/message.py b/src/crud/message.py index fbfb1698..d5ef569e 100644 --- a/src/crud/message.py +++ b/src/crud/message.py @@ -4,10 +4,21 @@ from logging import getLogger from typing import Any from nanoid import generate as generate_nanoid -from sqlalchemy import ColumnElement, Select, and_, func, or_, select, text +from sqlalchemy import ( + ColumnElement, + Select, + and_, + case, + func, + or_, + select, + text, + update, +) from sqlalchemy.ext.asyncio import AsyncSession from src import models, schemas +from src.cache.client import safe_cache_delete from src.config import settings from src.dependencies import tracked_db from src.embedding_client import embedding_client @@ -18,7 +29,7 @@ from src.utils.types import embedding_call_purpose from src.vector_store import get_external_vector_store from .peer import reject_scope_peers -from .session import get_or_create_session +from .session import get_or_create_session, session_cache_key logger = getLogger(__name__) @@ -408,7 +419,30 @@ async def create_messages( if pending_rows: db.add_all(pending_rows) + await db.flush() + latest_message_at = max(message.created_at for message in message_objects) + await db.execute( + update(models.Session) + .where(models.Session.workspace_name == workspace_name) + .where(models.Session.name == session_name) + .values( + last_message_at=case( + ( + models.Session.last_message_at.is_(None), + latest_message_at, + ), + ( + models.Session.last_message_at < latest_message_at, + latest_message_at, + ), + else_=models.Session.last_message_at, + ) + ) + .execution_options(synchronize_session=False) + ) + await db.commit() + await safe_cache_delete(session_cache_key(workspace_name, session_name)) return message_objects diff --git a/src/crud/session.py b/src/crud/session.py index efa3f664..29859c65 100644 --- a/src/crud/session.py +++ b/src/crud/session.py @@ -2,7 +2,7 @@ from dataclasses import dataclass from logging import getLogger -from typing import Any +from typing import Any, Literal from typing import cast as typing_cast from cashews import NOT_NONE @@ -66,8 +66,8 @@ class SessionDeletionResult: conclusions_deleted: int -SESSION_CACHE_KEY_TEMPLATE = "v2:workspace:{workspace_name}:session:{session_name}" -SESSION_LOCK_PREFIX = f"{get_cache_namespace()}:lock:v2" +SESSION_CACHE_KEY_TEMPLATE = "v3:workspace:{workspace_name}:session:{session_name}" +SESSION_LOCK_PREFIX = f"{get_cache_namespace()}:lock:v3" def session_cache_key(workspace_name: str, session_name: str) -> str: @@ -116,6 +116,7 @@ async def _fetch_session( "internal_metadata": obj.internal_metadata, "configuration": obj.configuration, "created_at": obj.created_at, + "last_message_at": obj.last_message_at, } @@ -162,6 +163,7 @@ async def get_sessions( workspace_name: str, filters: dict[str, Any] | None = None, reverse: bool = False, + sort_by: Literal["created_at", "last_message_at"] = "created_at", ) -> Select[tuple[models.Session]]: """ Get all active sessions in a workspace. @@ -169,7 +171,8 @@ async def get_sessions( Args: workspace_name: Name of the workspace filters: Optional filters to apply to the query - reverse: If True, order by created_at descending; if False, ascending + reverse: If True, order descending; if False, ascending + sort_by: Session timestamp used for ordering Returns: Select statement for Session objects @@ -182,9 +185,22 @@ async def get_sessions( stmt = apply_filter(stmt, models.Session, filters) + sort_column = ( + models.Session.last_message_at + if sort_by == "last_message_at" + else models.Session.created_at + ) if reverse: - return stmt.order_by(models.Session.created_at.desc(), models.Session.id.desc()) - return stmt.order_by(models.Session.created_at.asc(), models.Session.id.asc()) + primary_order = sort_column.desc() + id_order = models.Session.id.desc() + else: + primary_order = sort_column.asc() + id_order = models.Session.id.asc() + + if sort_by == "last_message_at": + primary_order = primary_order.nulls_last() + + return stmt.order_by(primary_order, id_order) async def get_or_create_session( @@ -390,6 +406,7 @@ async def get_or_create_session( "internal_metadata": honcho_session.internal_metadata, "configuration": honcho_session.configuration, "created_at": honcho_session.created_at, + "last_message_at": honcho_session.last_message_at, }, expire=settings.CACHE.DEFAULT_TTL_SECONDS, ) @@ -835,6 +852,8 @@ async def clone_session( insert_stmt = insert(models.Message).returning(models.Message) result = await db.execute(insert_stmt, new_messages) + cloned_messages = result.scalars().all() + new_session.last_message_at = max(message.created_at for message in cloned_messages) # Clone peers from original session to new session (including their configurations) stmt = select(models.SessionPeer).where( diff --git a/src/models.py b/src/models.py index 6433225a..b3bf0654 100644 --- a/src/models.py +++ b/src/models.py @@ -178,6 +178,9 @@ class Session(Base): created_at: Mapped[datetime.datetime] = mapped_column( DateTime(timezone=True), server_default=func.now(), index=True ) + last_message_at: Mapped[datetime.datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) workspace_name: Mapped[str] = mapped_column( ForeignKey("workspaces.name"), nullable=False, index=True ) @@ -193,6 +196,13 @@ class Session(Base): __table_args__ = ( UniqueConstraint("name", "workspace_name"), + Index( + "ix_sessions_workspace_last_message_at", + "workspace_name", + "last_message_at", + "id", + postgresql_where=text("is_active"), + ), CheckConstraint("length(name) <= 512", name="name_length"), CheckConstraint("length(id) = 21", name="id_length"), CheckConstraint("id ~ '^[A-Za-z0-9_-]+$'", name="id_format"), diff --git a/src/routers/sessions.py b/src/routers/sessions.py index f0d5b056..0c29586e 100644 --- a/src/routers/sessions.py +++ b/src/routers/sessions.py @@ -3,6 +3,7 @@ import logging from contextlib import suppress from time import perf_counter +from typing import Literal from fastapi import APIRouter, Body, Depends, Path, Query, Response from fastapi_pagination import Page @@ -262,6 +263,9 @@ async def get_sessions( None, description="Filtering and pagination options for the sessions list" ), reverse: bool = Query(False, description="Whether to reverse the order of results"), + sort_by: Literal["created_at", "last_message_at"] = Query( + "created_at", description="Session timestamp used to order results" + ), db: AsyncSession = read_db, ): """Get all Sessions for a Workspace, paginated with optional filters.""" @@ -278,6 +282,7 @@ async def get_sessions( workspace_name=workspace_id, filters=filter_param, reverse=reverse, + sort_by=sort_by, ), ) diff --git a/src/schemas/api.py b/src/schemas/api.py index 55a6121a..9cb6dee1 100644 --- a/src/schemas/api.py +++ b/src/schemas/api.py @@ -464,6 +464,7 @@ class Session(SessionBase): ) configuration: dict[str, Any] = Field(default_factory=dict) created_at: datetime.datetime + last_message_at: datetime.datetime | None = None model_config = ConfigDict( # pyright: ignore from_attributes=True, populate_by_name=True diff --git a/tests/alembic/revisions/__init__.py b/tests/alembic/revisions/__init__.py index 77c15483..b5172ebc 100644 --- a/tests/alembic/revisions/__init__.py +++ b/tests/alembic/revisions/__init__.py @@ -21,6 +21,7 @@ from . import ( test_baa22cad81e2_standardize_constraint_names, test_bb6fb3a7a643_add_message_seq_in_session_column, test_c3828084f472_add_indexes_for_messages_and_, + test_cfaff339d519_add_session_last_message_at, test_d429de0e5338_adopt_peer_paradigm, test_e4eba9cfaa6f_make_document_session_name_nullable, test_e9b705f9adf9_add_server_defaults_to_timestamp_, @@ -49,6 +50,7 @@ __all__ = [ "test_baa22cad81e2_standardize_constraint_names", "test_bb6fb3a7a643_add_message_seq_in_session_column", "test_c3828084f472_add_indexes_for_messages_and_", + "test_cfaff339d519_add_session_last_message_at", "test_d429de0e5338_adopt_peer_paradigm", "test_e4eba9cfaa6f_make_document_session_name_nullable", "test_e9b705f9adf9_add_server_defaults_to_timestamp_", diff --git a/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py b/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py new file mode 100644 index 00000000..6dcab6ba --- /dev/null +++ b/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py @@ -0,0 +1,118 @@ +"""Hooks for revision cfaff339d519 (add session last_message_at).""" + +from __future__ import annotations + +import datetime + +import sqlalchemy as sa +from nanoid import generate as generate_nanoid +from sqlalchemy import text + +from tests.alembic.registry import register_after_upgrade, register_before_upgrade +from tests.alembic.verifier import MigrationVerifier + +WORKSPACE_NAME = generate_nanoid() +PEER_NAME = generate_nanoid() +ACTIVE_SESSION_NAME = generate_nanoid() +EMPTY_SESSION_NAME = generate_nanoid() +LATEST_MESSAGE_AT = datetime.datetime(2026, 1, 3, 12, 0, tzinfo=datetime.timezone.utc) +INDEX_NAME = "ix_sessions_workspace_last_message_at" + + +@register_before_upgrade("cfaff339d519") +def prepare_add_session_last_message_at(verifier: MigrationVerifier) -> None: + """Seed sessions and messages before the activity timestamp exists.""" + verifier.assert_column_exists("sessions", "last_message_at", exists=False) + verifier.assert_indexes_not_exist([("sessions", INDEX_NAME)]) + + schema = verifier.schema + connection = verifier.conn + connection.execute( + text( + f""" + INSERT INTO "{schema}"."workspaces" ("id", "name") + VALUES (:id, :name) + """ + ), + {"id": generate_nanoid(), "name": WORKSPACE_NAME}, + ) + connection.execute( + text( + f""" + INSERT INTO "{schema}"."peers" ("id", "name", "workspace_name") + VALUES (:id, :name, :workspace_name) + """ + ), + { + "id": generate_nanoid(), + "name": PEER_NAME, + "workspace_name": WORKSPACE_NAME, + }, + ) + for session_name in (ACTIVE_SESSION_NAME, EMPTY_SESSION_NAME): + connection.execute( + text( + f""" + INSERT INTO "{schema}"."sessions" + ("id", "name", "workspace_name", "is_active") + VALUES (:id, :name, :workspace_name, true) + """ + ), + { + "id": generate_nanoid(), + "name": session_name, + "workspace_name": WORKSPACE_NAME, + }, + ) + + for seq, created_at in enumerate( + ( + datetime.datetime(2026, 1, 1, 12, 0, tzinfo=datetime.timezone.utc), + LATEST_MESSAGE_AT, + ), + start=1, + ): + connection.execute( + text( + f""" + INSERT INTO "{schema}"."messages" + ("public_id", "session_name", "content", "token_count", + "seq_in_session", "created_at", "peer_name", "workspace_name") + VALUES + (:public_id, :session_name, :content, 1, + :seq_in_session, :created_at, :peer_name, :workspace_name) + """ + ), + { + "public_id": generate_nanoid(), + "session_name": ACTIVE_SESSION_NAME, + "content": f"message {seq}", + "seq_in_session": seq, + "created_at": created_at, + "peer_name": PEER_NAME, + "workspace_name": WORKSPACE_NAME, + }, + ) + + +@register_after_upgrade("cfaff339d519") +def verify_add_session_last_message_at(verifier: MigrationVerifier) -> None: + """Verify the nullable field, historical backfill, and sorting index.""" + verifier.assert_column_exists("sessions", "last_message_at", nullable=True) + verifier.assert_column_type("sessions", "last_message_at", sa.TIMESTAMP) + verifier.assert_indexes_exist([("sessions", INDEX_NAME)]) + + rows = verifier.conn.execute( + text( + f""" + SELECT name, last_message_at + FROM "{verifier.schema}"."sessions" + WHERE workspace_name = :workspace_name + """ + ), + {"workspace_name": WORKSPACE_NAME}, + ).all() + activity_by_session = {row.name: row.last_message_at for row in rows} + + assert activity_by_session[ACTIVE_SESSION_NAME] == LATEST_MESSAGE_AT + assert activity_by_session[EMPTY_SESSION_NAME] is None diff --git a/tests/routes/test_sessions.py b/tests/routes/test_sessions.py index 27a4151a..b298c56f 100644 --- a/tests/routes/test_sessions.py +++ b/tests/routes/test_sessions.py @@ -214,6 +214,119 @@ def test_get_sessions(client: TestClient, sample_data: tuple[Workspace, Peer]): assert data["items"][0]["workspace_id"] == test_workspace.name +def test_session_response_exposes_null_last_message_at_without_messages( + client: TestClient, sample_data: tuple[Workspace, Peer] +): + """A newly created session reports that no message activity exists yet.""" + test_workspace, test_peer = sample_data + session_id = f"last-activity-empty-{generate_nanoid()}" + + response = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions", + json={ + "id": session_id, + "peer_names": {test_peer.name: {}}, + }, + ) + + assert response.status_code in [200, 201] + assert response.json()["last_message_at"] is None + + +def test_session_response_tracks_latest_message_timestamp_after_cached_create( + client: TestClient, sample_data: tuple[Workspace, Peer] +): + """Appending messages refreshes the cached session activity timestamp.""" + test_workspace, test_peer = sample_data + session_id = f"last-activity-cached-{generate_nanoid()}" + older_timestamp = "2026-01-01T12:00:00Z" + latest_timestamp = "2026-01-03T12:00:00Z" + + created = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions", + json={ + "id": session_id, + "peer_names": {test_peer.name: {}}, + }, + ) + assert created.status_code in [200, 201] + assert created.json()["last_message_at"] is None + + messages = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions/{session_id}/messages", + json={ + "messages": [ + { + "content": "older activity", + "peer_id": test_peer.name, + "created_at": older_timestamp, + }, + { + "content": "latest activity", + "peer_id": test_peer.name, + "created_at": latest_timestamp, + }, + ] + }, + ) + assert messages.status_code == 201 + + refreshed = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions", + json={"id": session_id}, + ) + + assert refreshed.status_code == 200 + assert datetime.datetime.fromisoformat( + refreshed.json()["last_message_at"].replace("Z", "+00:00") + ) == datetime.datetime(2026, 1, 3, 12, 0, tzinfo=datetime.timezone.utc) + + +def test_session_last_message_at_does_not_move_backwards_for_backdated_message( + client: TestClient, sample_data: tuple[Workspace, Peer] +): + """Appending historical data cannot make a session appear less recent.""" + test_workspace, test_peer = sample_data + session_id = f"last-activity-backfill-{generate_nanoid()}" + + created = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions", + json={ + "id": session_id, + "peer_names": {test_peer.name: {}}, + }, + ) + assert created.status_code in [200, 201] + + for content, created_at in ( + ("current activity", "2026-01-03T12:00:00Z"), + ("historical import", "2026-01-01T12:00:00Z"), + ): + message_response = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions/{session_id}/messages", + json={ + "messages": [ + { + "content": content, + "peer_id": test_peer.name, + "created_at": created_at, + } + ] + }, + ) + assert message_response.status_code == 201 + + refreshed = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions", + json={"id": session_id}, + ) + + assert refreshed.status_code == 200 + assert datetime.datetime.fromisoformat( + refreshed.json()["last_message_at"].replace("Z", "+00:00") + ) == datetime.datetime(2026, 1, 3, 12, 0, tzinfo=datetime.timezone.utc) + + def test_get_sessions_with_empty_filter( client: TestClient, sample_data: tuple[Workspace, Peer] ): @@ -280,6 +393,75 @@ def test_get_sessions_with_reverse( ] +@pytest.mark.asyncio +async def test_get_sessions_sort_by_last_message_at_reverses_activity_and_keeps_nulls_last( + client: TestClient, + db_session: AsyncSession, + sample_data: tuple[Workspace, Peer], +): + """Activity sorting differs from creation order and leaves empty sessions last.""" + test_workspace, test_peer = sample_data + activity_group = f"last-activity-sort-{generate_nanoid()}" + most_recent_session = f"activity-recent-{generate_nanoid()}" + older_session = f"activity-older-{generate_nanoid()}" + empty_session = f"activity-empty-{generate_nanoid()}" + + db_session.add_all( + [ + models.Session( + name=most_recent_session, + workspace_name=test_workspace.name, + created_at=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc), + h_metadata={"activity_group": activity_group}, + ), + models.Session( + name=older_session, + workspace_name=test_workspace.name, + created_at=datetime.datetime(2026, 1, 2, tzinfo=datetime.timezone.utc), + h_metadata={"activity_group": activity_group}, + ), + models.Session( + name=empty_session, + workspace_name=test_workspace.name, + created_at=datetime.datetime(2026, 1, 3, tzinfo=datetime.timezone.utc), + h_metadata={"activity_group": activity_group}, + ), + ] + ) + await db_session.commit() + + for session_id, created_at in ( + (most_recent_session, "2026-01-10T00:00:00Z"), + (older_session, "2026-01-05T00:00:00Z"), + ): + message_response = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions/{session_id}/messages", + json={ + "messages": [ + { + "content": f"activity for {session_id}", + "peer_id": test_peer.name, + "created_at": created_at, + } + ] + }, + ) + assert message_response.status_code == 201 + + response = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions/list?sort_by=last_message_at&reverse=true", + json={"filters": {"metadata": {"activity_group": activity_group}}}, + ) + + assert response.status_code == 200 + assert [item["id"] for item in response.json()["items"]] == [ + most_recent_session, + older_session, + empty_session, + ] + assert response.json()["items"][-1]["last_message_at"] is None + + @pytest.mark.asyncio async def test_get_sessions_reverse_uses_id_tiebreaker( client: TestClient, @@ -580,6 +762,10 @@ def test_clone_session(client: TestClient, sample_data: tuple[Workspace, Peer]): assert response.status_code == 201 data = response.json() assert data["metadata"] == {"test": "key"} + assert data["last_message_at"] is not None + cloned_last_message_at = datetime.datetime.fromisoformat( + data["last_message_at"].replace("Z", "+00:00") + ) # Check messages were cloned response = client.post( @@ -597,6 +783,10 @@ def test_clone_session(client: TestClient, sample_data: tuple[Workspace, Peer]): assert data["items"][1]["content"] == "Test message 2" assert data["items"][1]["metadata"] == {"key": "value2"} + assert cloned_last_message_at == max( + datetime.datetime.fromisoformat(item["created_at"].replace("Z", "+00:00")) + for item in data["items"] + ) def test_clone_session_with_cutoff( diff --git a/tests/sdk/test_peer.py b/tests/sdk/test_peer.py index 0c019d9b..2a3ce9ce 100644 --- a/tests/sdk/test_peer.py +++ b/tests/sdk/test_peer.py @@ -56,6 +56,7 @@ async def test_peer_sessions(client_fixture: tuple[Honcho, str]): await session1.aio.add_peers(peer) await session2.aio.add_peers(peer) + await session1.aio.add_messages(peer.message("peer session activity")) sessions_page = await peer.aio.sessions() sessions = sessions_page.items @@ -63,6 +64,9 @@ async def test_peer_sessions(client_fixture: tuple[Honcho, str]): session_ids = {s.id for s in sessions} assert "s1" in session_ids assert "s2" in session_ids + sessions_by_id = {session.id: session for session in sessions} + assert sessions_by_id["s1"].last_message_at is not None + assert sessions_by_id["s2"].last_message_at is None else: peer = honcho_client.peer(id="test-peer-sessions") session1 = honcho_client.session(id="s1") @@ -70,6 +74,7 @@ async def test_peer_sessions(client_fixture: tuple[Honcho, str]): session1.add_peers(peer) session2.add_peers(peer) + session1.add_messages(peer.message("peer session activity")) sessions_page = peer.sessions() sessions = list(sessions_page) @@ -77,6 +82,9 @@ async def test_peer_sessions(client_fixture: tuple[Honcho, str]): session_ids = {s.id for s in sessions} assert "s1" in session_ids assert "s2" in session_ids + sessions_by_id = {session.id: session for session in sessions} + assert sessions_by_id["s1"].last_message_at is not None + assert sessions_by_id["s2"].last_message_at is None @pytest.mark.asyncio diff --git a/tests/sdk/test_session.py b/tests/sdk/test_session.py index a15693fa..4269e2dd 100644 --- a/tests/sdk/test_session.py +++ b/tests/sdk/test_session.py @@ -1,3 +1,6 @@ +import datetime +from typing import cast + import pytest from sdks.python.src.honcho.api_types import QueueStatusResponse @@ -92,6 +95,108 @@ async def test_session_fetch_methods_refresh_cached_status_fields( assert session.is_active is True +@pytest.mark.asyncio +async def test_session_refresh_populates_last_message_at( + client_fixture: tuple[Honcho, str], +): + """Session.refresh exposes the newest message timestamp cached by the SDK.""" + honcho_client, client_type = client_fixture + session_id = f"test-session-last-message-at-{client_type}" + peer_id = f"test-peer-last-message-at-{client_type}" + message_time = datetime.datetime(2026, 1, 3, 12, 0, tzinfo=datetime.timezone.utc) + + if client_type == "async": + peer = await honcho_client.aio.peer(id=peer_id) + created = await honcho_client.aio.session(id=session_id, peers=[peer]) + await created.aio.add_messages( + peer.message("session activity", created_at=message_time) + ) + session = Session(session_id, honcho_client) + await session.aio.refresh() + else: + peer = honcho_client.peer(id=peer_id) + created = honcho_client.session(id=session_id, peers=[peer]) + created.add_messages(peer.message("session activity", created_at=message_time)) + session = Session(session_id, honcho_client) + session.refresh() + + assert session.last_message_at == message_time + + +@pytest.mark.asyncio +async def test_client_sessions_sort_by_last_message_at( + client_fixture: tuple[Honcho, str], +): + """Workspace session listing exposes and preserves activity ordering.""" + honcho_client, client_type = client_fixture + group = f"sdk-last-activity-sort-{client_type}" + peer_id = f"sdk-last-activity-peer-{client_type}" + recent_id = f"sdk-last-activity-recent-{client_type}" + older_id = f"sdk-last-activity-older-{client_type}" + empty_id = f"sdk-last-activity-empty-{client_type}" + recent_time = datetime.datetime(2026, 1, 10, tzinfo=datetime.timezone.utc) + older_time = datetime.datetime(2026, 1, 5, tzinfo=datetime.timezone.utc) + + if client_type == "async": + peer = await honcho_client.aio.peer(id=peer_id) + older = await honcho_client.aio.session( + id=older_id, metadata={"activity_group": group}, peers=[peer] + ) + recent = await honcho_client.aio.session( + id=recent_id, metadata={"activity_group": group}, peers=[peer] + ) + await honcho_client.aio.session( + id=empty_id, metadata={"activity_group": group}, peers=[peer] + ) + await recent.aio.add_messages( + peer.message("recent activity", created_at=recent_time) + ) + await older.aio.add_messages( + peer.message("older activity", created_at=older_time) + ) + page = await honcho_client.aio.sessions( + {"metadata": {"activity_group": group}}, + sort_by="last_message_at", + reverse=True, + size=1, + ) + sessions = cast(list[Session], page.items) + while page.has_next_page(): + next_page = await page.get_next_page() + assert next_page is not None + sessions.extend(cast(list[Session], next_page.items)) + page = next_page + else: + peer = honcho_client.peer(id=peer_id) + older = honcho_client.session( + id=older_id, metadata={"activity_group": group}, peers=[peer] + ) + recent = honcho_client.session( + id=recent_id, metadata={"activity_group": group}, peers=[peer] + ) + honcho_client.session( + id=empty_id, metadata={"activity_group": group}, peers=[peer] + ) + recent.add_messages(peer.message("recent activity", created_at=recent_time)) + older.add_messages(peer.message("older activity", created_at=older_time)) + page = honcho_client.sessions( + {"metadata": {"activity_group": group}}, + sort_by="last_message_at", + reverse=True, + size=1, + ) + sessions = cast(list[Session], page.items) + while page.has_next_page(): + next_page = page.get_next_page() + assert next_page is not None + sessions.extend(cast(list[Session], next_page.items)) + page = next_page + + assert [session.id for session in sessions] == [recent_id, older_id, empty_id] + assert sessions[0].last_message_at == recent_time + assert sessions[-1].last_message_at is None + + @pytest.mark.asyncio async def test_session_peer_management( client_fixture: tuple[Honcho, str], @@ -624,6 +729,7 @@ async def test_session_clone(client_fixture: tuple[Honcho, str]): cloned = await session.aio.clone() assert isinstance(cloned, Session) assert cloned.id != session.id # Should have a different ID + assert cloned.last_message_at is not None # Verify cloned session has the same messages cloned_messages_page = await cloned.aio.messages() @@ -652,6 +758,7 @@ async def test_session_clone(client_fixture: tuple[Honcho, str]): cloned = session.clone() assert isinstance(cloned, Session) assert cloned.id != session.id # Should have a different ID + assert cloned.last_message_at is not None # Verify cloned session has the same messages cloned_messages_page = cloned.messages() From 90e05f67e1c38ebcafb175d0475db54e6d79f5c0 Mon Sep 17 00:00:00 2001 From: steven-ji Date: Sun, 23 Aug 2026 10:32:18 +0800 Subject: [PATCH 2/3] fix(sdks): keep session activity cache current --- sdks/python/src/honcho/aio.py | 8 ++- sdks/python/src/honcho/session.py | 19 ++++- sdks/typescript/__tests__/peer.test.ts | 2 +- sdks/typescript/src/types/api.ts | 2 +- tests/sdk/test_file_uploads.py | 1 + tests/sdk/test_session.py | 99 +++++++++++++++++++++++++- 6 files changed, 122 insertions(+), 9 deletions(-) diff --git a/sdks/python/src/honcho/aio.py b/sdks/python/src/honcho/aio.py index f5a1011d..e50bcad0 100644 --- a/sdks/python/src/honcho/aio.py +++ b/sdks/python/src/honcho/aio.py @@ -1256,10 +1256,12 @@ class SessionAio(AsyncMetadataConfigMixin): routes.messages(self._session.workspace_id, self._session.id), body={"messages": messages_data}, ) - return [ + created_messages = [ Message.from_api_response(MessageResponse.model_validate(msg)) for msg in data ] + self._session._update_last_message_at_from_messages(created_messages) + return created_messages async def messages( self, @@ -1559,10 +1561,12 @@ class SessionAio(AsyncMetadataConfigMixin): data=data_dict, ) - return [ + created_messages = [ Message.from_api_response(MessageResponse.model_validate(msg)) for msg in response ] + self._session._update_last_message_at_from_messages(created_messages) + return created_messages @validate_call(config=ConfigDict(arbitrary_types_allowed=True)) async def representation( diff --git a/sdks/python/src/honcho/session.py b/sdks/python/src/honcho/session.py index d94972a1..28415329 100644 --- a/sdks/python/src/honcho/session.py +++ b/sdks/python/src/honcho/session.py @@ -127,6 +127,17 @@ class Session(SessionBase, MetadataConfigMixin): self._last_message_at = session.last_message_at self._is_active = session.is_active + def _update_last_message_at_from_messages( + self, messages: Sequence[Message] + ) -> None: + """Advance the cached activity timestamp from locally written messages.""" + if not messages: + return + + newest_message_at = max(message.created_at for message in messages) + if self._last_message_at is None or newest_message_at > self._last_message_at: + self._last_message_at = newest_message_at + def get_metadata(self) -> dict[str, object]: """ Get metadata from the server and update the cache. @@ -451,10 +462,12 @@ class Session(SessionBase, MetadataConfigMixin): routes.messages(self.workspace_id, self.id), body={"messages": messages_data}, ) - return [ + created_messages = [ Message.from_api_response(MessageResponse.model_validate(msg)) for msg in data ] + self._update_last_message_at_from_messages(created_messages) + return created_messages @validate_call def messages( @@ -906,10 +919,12 @@ class Session(SessionBase, MetadataConfigMixin): data=data_dict, ) - return [ + created_messages = [ Message.from_api_response(MessageResponse.model_validate(msg)) for msg in response ] + self._update_last_message_at_from_messages(created_messages) + return created_messages @validate_call(config=ConfigDict(arbitrary_types_allowed=True)) def representation( diff --git a/sdks/typescript/__tests__/peer.test.ts b/sdks/typescript/__tests__/peer.test.ts index 5d0c89da..e754dd83 100644 --- a/sdks/typescript/__tests__/peer.test.ts +++ b/sdks/typescript/__tests__/peer.test.ts @@ -222,7 +222,7 @@ describe('Peer', () => { const sessionIds = sessions.items.map((s) => s.id) expect(sessionIds).toContain('peer-sessions-test') const returned = sessions.items.find((item) => item.id === session.id) - expect(returned?.lastMessageAt).toBeDefined() + expect(typeof returned?.lastMessageAt).toBe('string') }) test('sessions returns empty for peer in no sessions', async () => { diff --git a/sdks/typescript/src/types/api.ts b/sdks/typescript/src/types/api.ts index c473ef69..3f004f88 100644 --- a/sdks/typescript/src/types/api.ts +++ b/sdks/typescript/src/types/api.ts @@ -126,7 +126,7 @@ export interface SessionResponse { metadata: Record configuration: SessionConfigApi created_at: string - last_message_at?: string | null + last_message_at: string | null } export interface SessionCreateParams { diff --git a/tests/sdk/test_file_uploads.py b/tests/sdk/test_file_uploads.py index 777abaef..391ac86c 100644 --- a/tests/sdk/test_file_uploads.py +++ b/tests/sdk/test_file_uploads.py @@ -48,6 +48,7 @@ async def test_session_upload_file( assert text_content in messages[0].content assert messages[0].peer_id == user.id assert messages[0].session_id == session.id + assert session.last_message_at == max(message.created_at for message in messages) @pytest.mark.asyncio diff --git a/tests/sdk/test_session.py b/tests/sdk/test_session.py index 4269e2dd..cc9de362 100644 --- a/tests/sdk/test_session.py +++ b/tests/sdk/test_session.py @@ -1,15 +1,104 @@ import datetime +from types import SimpleNamespace from typing import cast +from unittest.mock import AsyncMock, MagicMock import pytest -from sdks.python.src.honcho.api_types import QueueStatusResponse +from sdks.python.src.honcho.api_types import MessageCreateParams, QueueStatusResponse from sdks.python.src.honcho.client import Honcho from sdks.python.src.honcho.message import Message from sdks.python.src.honcho.peer import Peer from sdks.python.src.honcho.session import Session, SessionPeerConfig +def _message_response( + message_id: str, created_at: datetime.datetime +) -> dict[str, object]: + """Build a complete message response fixture for SDK boundary tests.""" + return { + "id": message_id, + "content": message_id, + "peer_id": "peer-1", + "session_id": "session-1", + "workspace_id": "workspace-1", + "metadata": {}, + "created_at": created_at, + "token_count": 1, + } + + +def _session_write_honcho() -> ( + tuple[SimpleNamespace, datetime.datetime, datetime.datetime] +): + """Build sync and async HTTP stubs for local session write tests.""" + newest_time = datetime.datetime(2026, 1, 10, 12, 0, tzinfo=datetime.timezone.utc) + older_time = datetime.datetime(2026, 1, 5, 12, 0, tzinfo=datetime.timezone.utc) + older_response = [_message_response("backdated activity", older_time)] + newest_response = [_message_response("newest activity", newest_time)] + sync_http = SimpleNamespace( + post=MagicMock(side_effect=[older_response, older_response]), + upload=MagicMock(side_effect=[newest_response, older_response]), + ) + async_http = SimpleNamespace( + post=AsyncMock(side_effect=[older_response, older_response]), + upload=AsyncMock(side_effect=[newest_response, older_response]), + ) + honcho = SimpleNamespace( + workspace_id="workspace-1", + _http=sync_http, + _async_http_client=async_http, + _ensure_workspace=MagicMock(return_value=None), + _ensure_workspace_async=AsyncMock(return_value=None), + ) + return honcho, older_time, newest_time + + +def test_session_sync_writes_update_last_message_at_monotonically() -> None: + """Sync message and file writes keep the newest cached activity timestamp.""" + honcho, older_time, newest_time = _session_write_honcho() + session = Session("session-1", honcho) + + session.add_messages(MessageCreateParams(content="activity", peer_id="peer-1")) + assert session.last_message_at == older_time + + session.upload_file(("activity.txt", b"activity", "text/plain"), peer="peer-1") + assert session.last_message_at == newest_time + + session.add_messages(MessageCreateParams(content="activity", peer_id="peer-1")) + assert session.last_message_at == newest_time + + session.upload_file(("activity.txt", b"activity", "text/plain"), peer="peer-1") + assert session.last_message_at == newest_time + + +@pytest.mark.asyncio +async def test_session_async_writes_update_last_message_at_monotonically() -> None: + """Async message and file writes keep the newest cached activity timestamp.""" + honcho, older_time, newest_time = _session_write_honcho() + session = Session("session-1", honcho) + + await session.aio.add_messages( + MessageCreateParams(content="activity", peer_id="peer-1") + ) + assert session.last_message_at == older_time + + await session.aio.upload_file( + ("activity.txt", b"activity", "text/plain"), peer="peer-1" + ) + assert session.last_message_at == newest_time + + await session.aio.add_messages( + MessageCreateParams(content="activity", peer_id="peer-1") + ) + assert session.last_message_at == newest_time + + await session.aio.upload_file( + ("activity.txt", b"activity", "text/plain"), peer="peer-1" + ) + assert session.last_message_at == newest_time + + @pytest.mark.asyncio async def test_session_metadata(client_fixture: tuple[Honcho, str]): """ @@ -729,12 +818,14 @@ async def test_session_clone(client_fixture: tuple[Honcho, str]): cloned = await session.aio.clone() assert isinstance(cloned, Session) assert cloned.id != session.id # Should have a different ID - assert cloned.last_message_at is not None # Verify cloned session has the same messages cloned_messages_page = await cloned.aio.messages() cloned_messages = cloned_messages_page.items assert len(cloned_messages) == 2 + assert cloned.last_message_at == max( + message.created_at for message in cloned_messages + ) # Verify original session still has messages original_messages_page = await session.aio.messages() @@ -758,12 +849,14 @@ async def test_session_clone(client_fixture: tuple[Honcho, str]): cloned = session.clone() assert isinstance(cloned, Session) assert cloned.id != session.id # Should have a different ID - assert cloned.last_message_at is not None # Verify cloned session has the same messages cloned_messages_page = cloned.messages() cloned_messages = list(cloned_messages_page) assert len(cloned_messages) == 2 + assert cloned.last_message_at == max( + message.created_at for message in cloned_messages + ) # Verify original session still has messages original_messages_page = session.messages() From c8f52e1f9e81ae0c86170dd6632b0443ab65c5ec Mon Sep 17 00:00:00 2001 From: steven-ji Date: Fri, 4 Sep 2026 09:08:12 +0800 Subject: [PATCH 3/3] fix(sessions): harden activity ordering --- ...faff339d519_add_session_last_message_at.py | 6 +- src/models.py | 4 +- src/schemas/api.py | 10 +++ ...faff339d519_add_session_last_message_at.py | 18 +++++ tests/routes/test_sessions.py | 68 +++++++++++++++++++ 5 files changed, 103 insertions(+), 3 deletions(-) diff --git a/migrations/versions/cfaff339d519_add_session_last_message_at.py b/migrations/versions/cfaff339d519_add_session_last_message_at.py index d4f011ef..9006dc89 100644 --- a/migrations/versions/cfaff339d519_add_session_last_message_at.py +++ b/migrations/versions/cfaff339d519_add_session_last_message_at.py @@ -48,7 +48,11 @@ def upgrade() -> None: op.create_index( INDEX_NAME, "sessions", - ["workspace_name", "last_message_at", "id"], + [ + "workspace_name", + sa.text("last_message_at DESC NULLS LAST"), + sa.text("id DESC"), + ], unique=False, schema=schema, postgresql_where=sa.text("is_active"), diff --git a/src/models.py b/src/models.py index b3bf0654..77d218e4 100644 --- a/src/models.py +++ b/src/models.py @@ -199,8 +199,8 @@ class Session(Base): Index( "ix_sessions_workspace_last_message_at", "workspace_name", - "last_message_at", - "id", + text("last_message_at DESC NULLS LAST"), + text("id DESC"), postgresql_where=text("is_active"), ), CheckConstraint("length(name) <= 512", name="name_length"), diff --git a/src/schemas/api.py b/src/schemas/api.py index b0e65ef5..65e94593 100644 --- a/src/schemas/api.py +++ b/src/schemas/api.py @@ -339,6 +339,16 @@ class MessageCreate(MessageBase): def sanitize_content(cls, v: str) -> str: return strip_nul(v) + @field_validator("created_at", mode="after") + @classmethod + def normalize_created_at_timezone( + cls, value: datetime.datetime | None + ) -> datetime.datetime | None: + """Treat timezone-naive message timestamps as UTC.""" + if value is not None and value.tzinfo is None: + return value.replace(tzinfo=datetime.UTC) + return value + @property def encoded_message(self) -> list[int]: return self._encoded_message diff --git a/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py b/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py index bfac4af8..9f9f4f95 100644 --- a/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py +++ b/tests/alembic/revisions/test_cfaff339d519_add_session_last_message_at.py @@ -102,6 +102,24 @@ def verify_add_session_last_message_at(verifier: MigrationVerifier) -> None: verifier.assert_column_type("sessions", "last_message_at", sa.TIMESTAMP) verifier.assert_indexes_exist([("sessions", INDEX_NAME)]) + index_definition = verifier.conn.execute( + text( + """ + SELECT indexdef + FROM pg_indexes + WHERE schemaname = :schema + AND tablename = 'sessions' + AND indexname = :index_name + """ + ), + {"schema": verifier.schema, "index_name": INDEX_NAME}, + ).scalar_one() + normalized_index_definition = " ".join(index_definition.replace('"', "").split()) + assert "(workspace_name, last_message_at DESC NULLS LAST, id DESC)" in ( + normalized_index_definition + ) + assert "WHERE is_active" in normalized_index_definition + rows = verifier.conn.execute( text( f""" diff --git a/tests/routes/test_sessions.py b/tests/routes/test_sessions.py index 9efe1e78..2d48bd00 100644 --- a/tests/routes/test_sessions.py +++ b/tests/routes/test_sessions.py @@ -4,6 +4,7 @@ from typing import Any import pytest from fastapi.testclient import TestClient from nanoid import generate as generate_nanoid +from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession from src import models @@ -327,6 +328,48 @@ def test_session_last_message_at_does_not_move_backwards_for_backdated_message( ) == datetime.datetime(2026, 1, 3, 12, 0, tzinfo=datetime.UTC) +def test_message_batch_normalizes_naive_activity_timestamps( + client: TestClient, sample_data: tuple[Workspace, Peer] +): + """Mixed explicit and server-default timestamps remain comparable.""" + test_workspace, test_peer = sample_data + session_id = f"last-activity-mixed-timezone-{generate_nanoid()}" + + response = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions/{session_id}/messages", + json={ + "messages": [ + { + "content": "naive activity timestamp", + "peer_id": test_peer.name, + "created_at": "2024-01-01T10:00:00", + }, + { + "content": "server timestamp", + "peer_id": test_peer.name, + }, + ] + }, + ) + + assert response.status_code == 201 + message_timestamps = [ + datetime.datetime.fromisoformat(item["created_at"].replace("Z", "+00:00")) + for item in response.json() + ] + assert all(timestamp.utcoffset() is not None for timestamp in message_timestamps) + + session_response = client.post( + f"/v3/workspaces/{test_workspace.name}/sessions", + json={"id": session_id}, + ) + assert session_response.status_code == 200 + session_activity = datetime.datetime.fromisoformat( + session_response.json()["last_message_at"].replace("Z", "+00:00") + ) + assert session_activity == max(message_timestamps) + + def test_get_sessions_with_empty_filter( client: TestClient, sample_data: tuple[Workspace, Peer] ): @@ -462,6 +505,31 @@ async def test_get_sessions_sort_by_last_message_at_reverses_activity_and_keeps_ assert response.json()["items"][-1]["last_message_at"] is None +@pytest.mark.asyncio +async def test_session_activity_index_supports_reverse_ordering( + db_session: AsyncSession, +): + """The model index matches the primary DESC NULLS LAST query shape.""" + result = await db_session.execute( + text( + """ + SELECT indexdef + FROM pg_indexes + WHERE schemaname = current_schema() + AND tablename = 'sessions' + AND indexname = 'ix_sessions_workspace_last_message_at' + """ + ) + ) + index_definition = result.scalar_one() + normalized_index_definition = " ".join(index_definition.replace('"', "").split()) + + assert "(workspace_name, last_message_at DESC NULLS LAST, id DESC)" in ( + normalized_index_definition + ) + assert "WHERE is_active" in normalized_index_definition + + @pytest.mark.asyncio async def test_get_sessions_reverse_uses_id_tiebreaker( client: TestClient,