Skip to content

AgentEvent

Agent events: one conversation’s commits as events shaped like a coding agent’s session events (message_start, message_update, tool_execution_*, turn_*, run_*, …), derived from its view (see View).

import { AgentEvent, type Ids } from "@effective-harness/core"
import * as Console from "effect/Console"
import * as Effect from "effect/Effect"
import * as Stream from "effect/Stream"
declare const conversationId: Ids.ConversationId
const program = Effect.gen(function* () {
const { snapshot, events } = yield* AgentEvent.watch(conversationId)
yield* Stream.runForEach(events, (event) => Console.log(JSON.stringify(event)))
})

A stream starts with the snapshot of the current state. A consumer that falls more than Session.MaxPendingWatchFrames commits behind gets one fresh snapshot event instead of the batches it missed.

Since v0.0.0


Signature

type AgentEvent =
| SnapshotEvent
| { readonly type: "run_start"; readonly inputs: ReadonlyArray<SubmissionId> }
| { readonly type: "run_end"; readonly inputs: ReadonlyArray<SubmissionId> }
| { readonly type: "turn_start" }
| { readonly type: "turn_end" }
| { readonly type: "message_start"; readonly message: Message }
/** `usage` is the partial's current usage. */
| { readonly type: "message_update"; readonly usage: Usage; readonly changes: ReadonlyArray<MessageChange> }
| { readonly type: "message_end"; readonly entry: EntryRecord }
| {
readonly type: "tool_execution_start"
readonly toolCallId: string
readonly toolName: string
readonly args: JsonObject
}
| {
readonly type: "tool_execution_update"
readonly toolCallId: string
readonly toolName: string
readonly output?: OutputChange
readonly details?: Json
readonly diagnostics?: ReadonlyArray<ToolDiagnostic>
}
/** `entry` is absent when the tool task faulted or was orphaned. */
| { readonly type: "tool_execution_end"; readonly toolCallId: string; readonly toolName: string; entry?: EntryRecord }
| { readonly type: "inbox_update"; readonly items: ReadonlyArray<QueuedItem> }
| { readonly type: "submission"; readonly record: SubmissionRecord }
| { readonly type: "auto_retry_start"; readonly attempt: number; readonly at: number; readonly errorMessage: string }
| { readonly type: "auto_retry_end"; readonly attempt: number }
/** A deferred response gained or moved its next status check. */
| { readonly type: "deferred_poll"; readonly pollAt: number }
| { readonly type: "entry_appended"; readonly entry: EntryRecord }
| { readonly type: "agent_changed"; readonly agent: AgentState }
| { readonly type: "usage_changed"; readonly usage: UsageState }
| { readonly type: "task_failed"; readonly taskId: TaskId; readonly kind: string; readonly message: string }
| {
readonly type: "compaction_start"
readonly taskId: TaskId
readonly reason: CompactionReason
readonly blocking: boolean
}
/** The task's receipt tells whether it produced a summary; the summary entry has its own events. */
| { readonly type: "compaction_end"; readonly taskId: TaskId; readonly reason: CompactionReason }

Source

Since v0.0.0

One conversation’s agent events.

Signature

export interface EventStream {
/** The `snapshot` event at attachment. */
readonly snapshot: SnapshotEvent
/**
* The events of every later commit, one batch per commit. An overflow replaces the undelivered batches with one
* batch holding a `snapshot` of the newest state. Ends when the Session closes. Consume either this or `events`.
*/
readonly batches: Stream.Stream<ReadonlyArray<AgentEvent>>
/** `batches`, flattened. */
readonly events: Stream.Stream<AgentEvent>
}

Source

Since v0.0.0

One change to the in-flight assistant message, relative to that message.

Signature

type MessageChange =
| {
readonly type: "text_start" | "thinking_start" | "toolcall_start"
readonly contentIndex: number
readonly block: Block
}
| { readonly type: "text_delta" | "thinking_delta"; readonly contentIndex: number; readonly delta: string }
| {
readonly type: "toolcall_delta"
readonly contentIndex: number
readonly path: ReadonlyArray<string | number>
readonly delta: string
}
| { readonly type: "block"; readonly contentIndex: number; readonly block: Block }
| { readonly type: "message"; readonly message: AssistantMessage }

Source

Since v0.0.0

Output change of a running tool: a front trim and then an append of the retained window, or its replacement.

Signature

type OutputChange = { readonly trimStart?: number; readonly append?: string } | { readonly set: string }

Source

Since v0.0.0

A queued submission in the inbox.

Signature

export interface QueuedItem {
readonly id: SubmissionId
readonly mode: InboxItem["mode"]
}

Source

Since v0.0.0

The current state of a conversation, at attachment or after an overflow.

Signature

export interface SnapshotEvent {
readonly type: "snapshot"
readonly entries: ReadonlyArray<EntryRecord>
readonly run?: { readonly inputs: ReadonlyArray<SubmissionId> }
/** Current generation attempt: its in-flight partial, retry backoff or deferred poll. */
readonly generation?: NonNullable<LiveState["generation"]>
readonly tools: ReadonlyArray<ToolSlot>
/** `pi.live.compactions`: live compactions with their attempt and retry backoff. */
readonly compactions: ReadonlyArray<CompactionStatus>
readonly inbox: ReadonlyArray<QueuedItem>
/** `pi.agent`; `{}` when absent. */
readonly agent: AgentState
readonly usage: UsageState
}

Source

Since v0.0.0

The snapshot event of a view.

Signature

declare const snapshotOf: (view: ConversationView) => SnapshotEvent

Source

Since v0.0.0

Attach to one conversation’s agent events. The snapshot and the registration for later commits are captured atomically on the commit line. The stream lives until the scope closes.

Signature

declare const watch: (
conversationId: ConversationId
) => Effect.Effect<EventStream, ConversationNotFound | ReadError, Session | Scope.Scope>

Source

Since v0.0.0