События¶
Журнал событий Control Plane — append-only история всего, что произошло в tenant'е: создание и смена статусов задач, claims, runs, approvals, артефакты, комментарии, цели. Он служит аудитом, источником синхронизации для харнессов и консоли и входом для памяти. Статья описывает модель события, надёжный курсор, чтение страницами и через WebSocket, каталог типов событий и хранение журнала. Фильтры подписки, версии данных событий и SDK потребителя — в статье Подписки на события.
Как событие появляется¶
Событие пишется той же транзакцией, что и изменение состояния:
sequenceDiagram
autonumber
participant Cmd as Команда (API / воркер)
participant DB as PostgreSQL
participant Hub as Realtime-хаб API
participant WS as WebSocket-клиент
Cmd->>DB: изменение таблиц состояния
Cmd->>DB: INSERT events (+ outbox)
Cmd->>DB: pg_notify('cp_events', …)
Cmd->>DB: COMMIT
DB-->>Hub: NOTIFY (только после commit)
Hub->>DB: чтение событий после позиции клиента
Hub-->>WS: события по порядку
- Commit публикует всё атомарно: состояние, событие и запись outbox; откат не оставляет ничего.
NOTIFYдоставляется только при commit, поэтому подписчики никогда не просыпаются по откаченным данным.- Журнал append-only: триггеры базы запрещают
UPDATE,DELETEиTRUNCATE. Единственное исключение — операторская архивация (см. ниже).
Модель события¶
{
"sequence": 48211,
"id": "…",
"tenantId": "…",
"type": "task.updated",
"schemaVersion": 1,
"entityType": "task",
"workspaceId": "<workspace-id>",
"entityId": "…",
"actorId": "<principal-id>",
"sessionId": null,
"correlationId": "…",
"causationId": null,
"requestId": "req_…",
"traceRunId": "…",
"iamActorId": "<iam-principal-id>",
"payload": {
"publicId": "TASK-000123",
"changes": {"status": "blocked", "customFields": true},
"fromStatus": "in_progress",
"status": "blocked",
"systemStatusCategory": "blocked",
"version": 9
},
"occurredAt": "2026-09-01T10:15:04.117Z",
"cursor": "ec1_…"
}
| Поле | Описание |
|---|---|
sequence |
Идентификатор события и порядок внутри транзакции. Не курсор воспроизведения |
id |
Идентификатор события (UUID) — ключ дедупликации у потребителя |
type |
Тип события, см. каталог ниже |
schemaVersion |
Версия схемы payload этого типа; версии только добавляют поля, см. Подписки на события |
entityType, entityId |
Сущность, в чей поток относится событие |
workspaceId |
Workspace сущности (или её задачи); null у событий уровня tenant |
actorId |
Principal, выполнивший действие (null у фоновых действий воркера) |
iamActorId |
IAM-identity актора; null у legacy-ключей |
sessionId |
Сессия, если действие выполнено в её рамках |
correlationId |
Из заголовка X-Correlation-ID или сгенерирован |
causationId |
Событие-причина (например, решение approval для событий его исхода) |
requestId |
Из X-Request-ID |
traceRunId |
Trace-корреляция из X-Run-Id (не доменный Run) |
payload |
Данные события — только ссылки и безопасные поля |
cursor |
Непрозрачный курсор позиции этого события |
Что в журнал не попадает¶
Журнал читают шире, чем сами сущности, и из него строится память. Поэтому в него сознательно не пишутся:
| Не пишется | Что пишется вместо |
|---|---|
| Содержимое custom fields | "customFields": true |
| Тексты комментариев | bodyLength |
content артефактов |
Ссылки: type, name, uri, ids |
| Данные checkpoints | checkpointId, seq, kind |
Тексты directive / reason управляющих сообщений |
ids, seq, операция, статус, causalPosition, safeBoundary |
| Acceptance и evidence задачи | Числа элементов |
Желаемое состояние цели, spec проверок |
desired_state: true, число критериев |
note и url в origin |
Сводка: kind, ref, ruleId, ids фактов |
| Run actions | Отдельная таблица аудита исполнения, не журнал |
| Title, summary и данные дочерней работы | ids, correlationId, исход, хэш результата |
Надёжный курсор¶
Порядок выдачи — пара (tx_id, sequence), где tx_id — 64-битный
идентификатор пишущей транзакции PostgreSQL. Выдаются только события ниже
стабильного горизонта — транзакции, которые уже гарантированно
завершились. Отсюда свойство: курсор, продвигающийся только по выданным
позициям, не может перешагнуть событие, которое закоммитится позже.
Почему не sequence
sequence назначается при INSERT, а tx_id — при первой записи
транзакции; у конкурентных команд эти порядки могут расходиться. Курсор по
sequence мог бы навсегда перескочить ещё невидимое событие с меньшим
номером. Цена надёжного курсора — задержка выдачи на время самой долгой
открытой пишущей транзакции: доставка откладывается, но не теряется.
Курсор — непрозрачная строка вида ec1_<base64url>. Клиенты не должны
разбирать, сравнивать или конструировать курсоры: храните последний
полученный и передавайте его обратно.
| Ситуация | Ответ |
|---|---|
| Малформированный курсор | 422 invalid_cursor |
| Курсор будущей версии формата | 422 unsupported_cursor_version |
| Курсор ниже границы удалённой истории | 422 cursor_below_journal_floor |
Для совместимости принимаются устаревшие формы: целое ?after=<sequence> и
старое кодирование nextCursor. При переходе с них возможна повторная выдача
уже виденных событий (at-least-once).
Чтение страницами¶
# С начала доступной истории
curl -s "$CP/events?limit=200" -H "Authorization: Bearer $TOKEN"
# Продолжение с сохранённого курсора
curl -s "$CP/events?cursor=ec1_…&limit=200" -H "Authorization: Bearer $TOKEN"
# Поток одной задачи
curl -s "$CP/events?entityType=task&entityId=<task-id>" -H "Authorization: Bearer $TOKEN"
# Последние 20 стабильных событий (диагностика)
curl -s "$CP/events?tail=20" -H "Authorization: Bearer $TOKEN"
| Параметр | Описание |
|---|---|
cursor |
Непрозрачный курсор; без него — чтение с начала доступной истории |
after |
Устаревший целочисленный курсор (sequence) |
limit |
1–200, по умолчанию 50 |
tail |
Последние N стабильных событий в порядке доставки (не больше limit) |
entityType, entityId |
Фильтр по потоку сущности |
types |
Префиксы типа (approval.), до 20 — см. Подписки на события |
workspaceId |
События поддерева workspace; право events.read проверяется на этом workspace |
Ответ всегда содержит nextCursor (на пустой странице — эхо входного
курсора) и hasMore:
Право: events.read на tenant, а с workspaceId — на этом workspace.
Цикл подписчика¶
cursor = load_cursor() # None при первом запуске
while True:
page = get("/api/v1/events", cursor=cursor, limit=200)
for event in page["items"]:
handle(event) # обработчик должен быть идемпотентным
cursor = event["cursor"]
save_cursor(cursor)
if not page["hasMore"]:
sleep(1) # или ждать WebSocket
cursor = page["nextCursor"]
Доставка — at-least-once: после сбоя между обработкой и сохранением курсора
событие придёт снова. Делайте обработку идемпотентной, например по
event["id"]. Готовый цикл с хранением курсора, дедупликацией и повторами —
EventConsumer из SDK, см. Подписки на события.
WebSocket¶
WebSocket — не источник истины, а сигнал «проснись и дочитай». Сервер
всегда читает события из таблицы в порядке (tx_id, sequence), отдаёт их
пачками по 200 и засыпает до NOTIFY своего tenant'а или до таймаута
CP_WS_POLL_INTERVAL_SECONDS (5 с по умолчанию) — потерянное уведомление не
теряет событий. Каждое сообщение — событие в той же форме, что и в GET
/events, с полем cursor.
Аутентификация — как у HTTP; право events.read. Ошибки передаются кодом
закрытия после установления соединения:
| Код закрытия | Причина |
|---|---|
4401 |
Нет или неверные credentials |
4403 |
Нет права events.read |
4404 |
Workspace фильтра не существует |
4400 |
Малформированный или неподдерживаемый курсор, неверный фильтр типов |
4503 |
Решение об авторизации не получено (PDP недоступен) — повтор имеет смысл |
1011 |
Внутренняя ошибка сервера |
При переподключении передайте ?after=<cursor последнего обработанного
события> — пропущенное будет дочитано.
Каталог событий¶
Задачи и работа¶
| Тип | Поток | Ключевые поля payload |
|---|---|---|
task.created |
task | publicId, title, status, systemStatusCategory, typeKey, typeVersion, priority, workspaceId, startDate, dueDate, customFields (флаг), goalId, origin (сводка), acceptanceChecks |
task.updated |
task | changes, version; при смене статуса — fromStatus, status, systemStatusCategory |
task.claimed |
task | claimId, sessionId, holderId, fencingToken, expiresAt, status, systemStatusCategory, version |
task.completed |
task | publicId, status, systemStatusCategory, version |
task.relation_added, task.relation_removed |
task | Связь |
task.comment_added, task.comment_edited |
task | commentId, authorPrincipalId, version, bodyLength, runId, artifactId |
task.external_reference_added, task.external_reference_updated |
task | Внешняя ссылка без metadata |
task_type.created, task_type.deprecated |
task_type | key, version и сводка lifecycle |
goal.created, goal.updated |
goal | См. Цели |
task.verification_started |
task | publicId, taskId, verificationId, attempt, trigger, checks |
task.verification_failed |
task | То же плюс results, failedCheck, reason, consecutiveFailures, blocked, fromStatus, status, systemStatusCategory |
task.verified |
task | publicId, taskId, verificationId, attempt, results, artifactId |
task.completion_work_executed, task.completion_work_failed |
task | Работа после завершения, объявленная типом задачи |
task.context_pack_recorded |
task | Пакет контекста, собранный при claim: ids и счётчики |
rule.created, .updated, .enabled, .disabled, .archived, .evaluated |
rule | Правила вывода работы |
work.derived, work.reconciled |
task | Работа, выведенная правилом: ruleId, ruleKey, evaluationId, taskId, dedupKey |
Исполнение¶
| Тип | Поток | Ключевые поля payload |
|---|---|---|
session.opened, session.closed, session.expired |
session | — |
claim.released |
claim | taskId, reason, taskStatus (и taskSystemStatusCategory при закрытии сессии) |
claim.expired |
claim | taskId, reason (expired, session_inactive) |
run.started |
run | taskId, claimId, attempt, fencingToken |
run.succeeded |
run | taskId, attempt, taskCompleted |
run.failed |
run | taskId, reason (в том числе superseded), attempt |
run.cancelled |
run | taskId, reason, attempt |
run.suspended |
run | taskId, reason, attempt, waitingForApprovalId |
run.checkpointed |
run | taskId, checkpointId, seq, kind |
run.handoff_prepared |
run | taskId, claimId, checkpointId, fencingToken, reason |
run.cancel_requested |
run | taskId, attempt |
run.control_message.accepted, .applied, .rejected, .superseded |
run | controlMessageId, seq, operation, status, causalPosition, safeBoundary |
run.manifest_compiled, run.manifest_ephemeral_recorded |
run | Версия манифеста исполнения |
run.child.launched, .started, .resolved, .revoked, .cancel_requested |
run | ids, correlationId, исход, хэш результата |
Approvals, артефакты, наблюдения¶
| Тип | Поток | Ключевые поля payload |
|---|---|---|
approval.requested (v2) |
approval | taskId, artifactId, requiredRoleId, assignedPrincipalId, gate; v2 — workspaceId, taskPublicId, taskTitle, requestedBy, comment |
approval.approved, approval.rejected (v2) |
approval | taskId, artifactId, outcomeStatus; v2 — decisionBy, comment, channel (канал решения, null для прямого вызова API) |
approval.cancelled (v2) |
approval | taskId; v2 — cancelledBy |
approval.outcome_executed, .outcome_failed, .outcome_deferred |
approval | См. Approvals |
artifact.created |
artifact | type, name, taskId, runId, uri, supersedesArtifactId |
observation.recorded |
observation | Явное «запомнить» (см. Контекст задачи и память) |
knowledge.snapshot_reconciled, knowledge.pack_registered, knowledge.packs_configured |
— | Счётчики без содержимого |
skill.invocation_requested, _claimed, _retry_scheduled, _succeeded, _failed, _cancelled |
skill_invocation | ids, skill и версия, попытка, код ошибки — без входов и выходов |
Организация и конфигурация¶
| Группа | Типы |
|---|---|
| Tenant и principals | tenant.bootstrapped, principal.created, api_key.created, api_key.revoked, iam_binding.created, iam_binding.revoked, delegation.created, delegation.revoked |
| Workspaces | workspace.created, .updated, .archived, .moved, .member_added, .member_removed; workspace_type.created, .updated, .archived |
| Роли и каталог | role.created, .updated, .assigned, .revoked; capability.created, .assigned, .revoked; skill.registered, .updated, .assigned, .revoked |
| Проекты | project_template.created, .deprecated; project.created, .updated, .archived, .status_changed, .config_revision_created, .config_revision_activated, .external_reference_added, .external_reference_updated |
| Операции | context_adapter.redriven, context_adapter.rebuilt, event_journal.archived, event_journal.pruned |
Статус: ключ или категория
События задач несут и пользовательский ключ status, и
systemStatusCategory. Подписчику, которому важен смысл («задача
завершена»), стоит реагировать на категорию: ключи у разных типов
задач разные.
Потребители журнала¶
| Потребитель | Как читает | Гарантия |
|---|---|---|
| Харнессы, runner'ы, консоль | GET /events, WebSocket |
At-least-once по курсору клиента |
Сервисы на SDK EventConsumer (например, сервис уведомлений) |
GET /events с фильтрами, WebSocket как будильник |
Курсор и отметки обработанных событий в базе потребителя: без потерь и без повторов для эффектов в этой базе |
| Context Adapter | Внутренний per-tenant курсор в event_consumer_cursors |
At-least-once; курсор двигается только после подтверждения памятью; сбойный tenant паркуется отдельно |
| Outbox воркера | Таблица outbox, FOR UPDATE SKIP LOCKED |
At-least-once; ограниченные повторы с backoff, после исчерпания запись остаётся с last_error. В базовой поставке точка доставки — структурированный лог |
Управление Context Adapter (GET /operations/context-adapter,
:redrive, :rebuild) описано в Контексте задачи и памяти.
Хранение журнала¶
Журнал растёт без ограничений, пока оператор не выполнит архивацию.
Операции требуют права operations.manage.
# Перенести подтверждённые и достаточно старые события в архив
curl -s -X POST "$CP/operations/journal:archive" \
-H "Authorization: Bearer $TOKEN" -H "Content-Type: application/json" \
-d '{"beforeSeconds": 2592000, "maxEvents": 50000}'
# Физически удалить заархивированное (необратимо)
curl -s -X POST "$CP/operations/journal:prune" \
-H "Authorization: Bearer $TOKEN" -H "Content-Type: application/json" \
-d '{"beforeSeconds": 7776000}'
flowchart LR
E["events<br/>(горячий журнал)"] -->|":archive"| A["event_archive"]
A -->|":prune"| X["удалено"]
E -. "journal floor" .- A
A -. "archive floor" .- X
:archiveпереносит события старшеbeforeSeconds(по умолчаниюCP_JOURNAL_RETENTION_MIN_AGE_SECONDS, 30 суток), но не дальше минимальной позиции курсоров потребителей и не дальше самого старого недоставленного outbox-события. То, что кому-то ещё нужно, не уезжает. Чтение черезGET /eventsпрозрачно охватывает архив — аудит не меняется.- Если ни один потребитель ещё не зарегистрировал курсор, архивация
отклоняется
409 retention_blocked_by_consumer. :prune— единственная операция, после которой данные теряются. Курсор ниже удалённой границы получает422 cursor_below_journal_floor, а не молчаливый пропуск; чтение без курсора начинается с первого сохранившегося события.- Обе операции пишут собственные события
event_journal.archived/event_journal.pruned.
Перед prune
Убедитесь, что резервные копии базы содержат нужную историю (см. Резервное копирование) и что память не придётся перестраивать из удаляемой части журнала.
См. также¶
- Подписки на события — фильтры, версии данных, каталог и SDK потребителя.
- Контекст задачи и память — Context Adapter и наблюдения.
- Харнесс-протокол — курсор в self-контексте харнесса.
- Мониторинг и здоровье
- Артефакты и комментарии