refactor(storage): drop StorageBundle + legacy EventStore, add HttpEventStore
ci / JVM build + tests (push) Successful in 6m13s
ci / JVM build + tests (push) Successful in 6m13s
- Remove legacy pw.binom.agentik.messageStore.events.EventStore (EventRecord,
EventType) and all three impls (in-memory, sqlite, ksqlite) + tests + .sq
- Drop :storage-bundle module entirely; ChatAgent / ChatConversation /
ConversationLoop / DebugRoutes now take stores individually
(conversationStore, messageStore, workingMemoryStore, reflectionStore,
eventStore) instead of StorageBundle
- Delete server endpoints /events/replay and /conversations/{id}/events/replay;
Route.agentikAgent no longer takes eventStore param
- Add :client/HttpEventStore implementing :event-store/EventStore over HTTP:
events() -> GET /events/all, agentEvents() -> GET /events,
conversationEvents(convId) -> GET /conversations/{id}/events;
exposed via AgentClient.eventStore
- :event-store: add macosX64/macosArm64/linuxArm64 targets to match :client KMP
- :working-memory-api: drop api dep on :message-store-api (no longer needed)
- :storage-{inmemory,sqlite,ksqlite}: drop deps on :storage-bundle
This commit is contained in:
@@ -15,6 +15,7 @@ import io.ktor.http.contentType
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.agentik.eventStore.EventStore
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
@@ -38,6 +39,12 @@ internal class AgentClient(
|
||||
|
||||
private val agentUrl: String = baseUrl.trimEnd('/')
|
||||
|
||||
/**
|
||||
* Единый канал событий (lifecycle + per-conversation). Под капотом —
|
||||
* [HttpEventStore]: каждый метод бьёт свой URL (см. KDoc).
|
||||
*/
|
||||
val eventStore: EventStore = HttpEventStore(httpClient = httpClient, baseUrl = agentUrl)
|
||||
|
||||
override fun createConversation(temp: Boolean): Conversation =
|
||||
runBlocking {
|
||||
val snapshot: ConversationSnapshot = httpClient.post("$agentUrl/conversations") {
|
||||
@@ -97,56 +104,4 @@ internal class AgentClient(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Catchup для /events/replay — пагинированно читает events после [afterId].
|
||||
* Caller делает несколько вызовов пока `result.size < limit` (= конец).
|
||||
*
|
||||
* @param afterId exclusive cursor. `null` = с начала.
|
||||
* @param limit max per-request (default 100, max 1000 на сервере).
|
||||
*/
|
||||
internal suspend fun replayAllEvents(
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecordDto> {
|
||||
val response = httpClient.get("$agentUrl/events/replay") {
|
||||
afterId?.let { parameter("after_id", it) }
|
||||
parameter("limit", limit)
|
||||
}
|
||||
return response.body()
|
||||
}
|
||||
|
||||
/**
|
||||
* Catchup для /conversations/{id}/events/replay — пагинированно.
|
||||
*/
|
||||
internal suspend fun replayConversationEvents(
|
||||
conversationId: String,
|
||||
afterId: String? = null,
|
||||
limit: Int = 100,
|
||||
): List<EventRecordDto> {
|
||||
val response = httpClient.get("$agentUrl/conversations/$conversationId/events/replay") {
|
||||
afterId?.let { parameter("after_id", it) }
|
||||
parameter("limit", limit)
|
||||
}
|
||||
return response.body()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Запись event'а, которую возвращает /events/replay endpoint.
|
||||
*
|
||||
* Это **мини-зеркало** `pw.binom.agentik.messageStore.events.EventRecord` — клиент
|
||||
* не зависит от `:storage-core`, поэтому определяет свою модель (wire-only).
|
||||
*
|
||||
* Формат полностью совместим с сервером — `agentikJson.encodeToString(...)` там
|
||||
* и `agentikJson.decodeFromString(...)` здесь.
|
||||
*/
|
||||
@kotlinx.serialization.Serializable
|
||||
internal data class EventRecordDto(
|
||||
val id: String,
|
||||
val conversationId: String? = null,
|
||||
val createdAt: kotlin.time.Instant,
|
||||
val type: String,
|
||||
val payload: String,
|
||||
)
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.request.prepareGet
|
||||
import io.ktor.client.statement.bodyAsChannel
|
||||
import io.ktor.http.HttpStatusCode
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import pw.binom.agentik.eventStore.EventStore
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.CommonEvent
|
||||
import pw.binom.agentik.proto.Event
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* HTTP-реализация [EventStore], ходящая в `:server`-фасад.
|
||||
*
|
||||
* **Хитрый план**: вместо того, чтобы все методы шли в один общий endpoint и
|
||||
* фильтровали client-side ([EventStore.events]/[filterIsInstance]), эта
|
||||
* реализация бьёт запросы по URL'ам в зависимости от того, какой класс
|
||||
* событий нужен:
|
||||
* - [events] → `GET /events/all` (полный поток CommonEvent)
|
||||
* - [agentEvents] → `GET /events` (только lifecycle диалогов)
|
||||
* - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events`
|
||||
*
|
||||
* Так серверный фильтр (SQL `WHERE` или разные буферы) работает на своей стороне,
|
||||
* а клиент получает ровно тот срез, который ему нужен, без лишнего трафика.
|
||||
*
|
||||
* Для [conversationEvents] с `conversationId == null` (события всех диалогов)
|
||||
* fallback на default [EventStore.conversationEvents] — общий поток `/events/all`
|
||||
* + фильтр client-side. Это редкий кейс (admin-дашборды), и оптимизировать его
|
||||
* отдельно нерационально.
|
||||
*
|
||||
* [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()`
|
||||
* (см. KDoc [EventStore.earliestEventDate] — для пустого буфера это и есть
|
||||
* контрактное значение). Клиент, который полагался на gap detection через
|
||||
* message store, продолжит работать — просто fallback никогда не сработает.
|
||||
*/
|
||||
internal class HttpEventStore(
|
||||
private val httpClient: HttpClient,
|
||||
private val baseUrl: String,
|
||||
) : EventStore {
|
||||
|
||||
private val agentUrl: String = baseUrl.trimEnd('/')
|
||||
|
||||
override fun events(after: Instant?): Flow<CommonEvent> = flow {
|
||||
val url = buildString {
|
||||
append("$agentUrl/events/all")
|
||||
if (after != null) append("?after=$after")
|
||||
}
|
||||
httpClient.prepareGet(url) { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"events: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(CommonEvent.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Override: идём в `/events` напрямую — сервер фильтрует только lifecycle-события.
|
||||
* Default из [EventStore.agentEvents] читал бы `/events/all` + `filterIsInstance`.
|
||||
*/
|
||||
override fun agentEvents(after: Instant?): Flow<CommonEvent.Agent> = flow {
|
||||
val url = buildString {
|
||||
append("$agentUrl/events")
|
||||
if (after != null) append("?after=$after")
|
||||
}
|
||||
httpClient.prepareGet(url) { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"agentEvents: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
val event = agentikJson.decodeFromString(AgentEvent.serializer(), payload)
|
||||
emit(CommonEvent.Agent(date = event.date, event = event))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Override с `conversationId != null` — идём в `/conversations/{id}/events`.
|
||||
* С `null` (события всех диалогов) — fallback на default impl из [EventStore]:
|
||||
* общий `/events/all` + filter.
|
||||
*/
|
||||
override fun conversationEvents(
|
||||
after: Instant?,
|
||||
conversationId: String?,
|
||||
): Flow<CommonEvent.Conversation> {
|
||||
if (conversationId == null) {
|
||||
return super.conversationEvents(after, null)
|
||||
}
|
||||
return flow {
|
||||
val url = buildString {
|
||||
append("$agentUrl/conversations/$conversationId/events")
|
||||
if (after != null) append("?after=$after")
|
||||
}
|
||||
httpClient.prepareGet(url) { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"conversationEvents: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
val event = agentikJson.decodeFromString(Event.serializer(), payload)
|
||||
emit(CommonEvent.Conversation(date = event.date, conversationId = conversationId, event = event))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* У HTTP-варианта нет своего endpoint'а для earliest-event-date.
|
||||
* Контракт [EventStore.earliestEventDate] для пустого буфера говорит
|
||||
* "сейчас" — для HTTP-клиента буфер на нашей стороне всегда "пуст"
|
||||
* (мы не держим своё состояние), поэтому возвращаем `Clock.System.now()`.
|
||||
*/
|
||||
override suspend fun earliestEventDate(): Instant = Clock.System.now()
|
||||
|
||||
override fun close() {
|
||||
// HttpClient закрывает владелец (AgentClient / AgentikAgent).
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user