Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 65 additions & 0 deletions src/agent/event-queue.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
import { describe, expect, it } from 'vitest'
import type { Event, EventType } from '../types.js'
import { EventQueue } from './event-queue.js'

function makeEvent(
type: EventType,
channelId: string,
id: string,
guildId = 'guild'
): Event {
return {
type,
channelId,
guildId,
data: { id },
timestamp: new Date('2026-09-01T12:00:00Z'),
}
}

describe('EventQueue channel isolation', () => {
it('batches consecutive Discord events from the same channel', () => {
const queue = new EventQueue()
const message = makeEvent('message', 'channel-a', 'message')
const reaction = makeEvent('reaction', 'channel-a', 'reaction')
queue.push(message)
queue.push(reaction)

expect(queue.pollBatch()).toEqual([message, reaction])
})

it('never combines Discord events from different channels', () => {
const queue = new EventQueue()
const first = makeEvent('message', 'channel-a', 'first')
const second = makeEvent('message', 'channel-b', 'second')
queue.push(first)
queue.push(second)

expect(queue.pollBatch()).toEqual([first])
expect(queue.pollBatch()).toEqual([second])
})

it('never combines internal events from different channels', () => {
const queue = new EventQueue()
const first = makeEvent('self_activation', 'channel-a', 'first')
const second = makeEvent('timer', 'channel-b', 'second')
queue.push(first)
queue.push(second)

expect(queue.pollBatch()).toEqual([first])
expect(queue.pollBatch()).toEqual([second])
})

it('prepends a deferred batch without reversing it', () => {
const queue = new EventQueue()
const deferredFirst = makeEvent('message', 'channel-a', 'deferred-first')
const deferredSecond = makeEvent('reaction', 'channel-a', 'deferred-second')
const newer = makeEvent('message', 'channel-b', 'newer')
queue.push(newer)

queue.prepend([deferredFirst, deferredSecond])

expect(queue.pollBatch()).toEqual([deferredFirst, deferredSecond])
expect(queue.pollBatch()).toEqual([newer])
})
})
19 changes: 16 additions & 3 deletions src/agent/event-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,19 @@ export class EventQueue {
this.queue.push(event)
}

/**
* Restore an older batch ahead of events that arrived after it was polled.
*/
prepend(events: Event[]): void {
if (events.length === 0) return
this.queue.unshift(...events)
}

/**
* Poll a batch of events
* Returns all consecutive Discord events (same type category)
* Returns consecutive events for one channel and the same type category.
* AgentLoop processes a batch using its first event's channel and guild, so
* mixing channels would route every later event through the wrong context.
*/
pollBatch(): Event[] {
if (this.queue.length === 0) {
Expand Down Expand Up @@ -46,7 +56,11 @@ export class EventQueue {
}

// Stop if we hit a different event category
if (this.isDiscordEvent(nextEvent) !== isDiscordEvent) {
if (
nextEvent.channelId !== firstEvent.channelId ||
nextEvent.guildId !== firstEvent.guildId ||
this.isDiscordEvent(nextEvent) !== isDiscordEvent
) {
break
}

Expand Down Expand Up @@ -81,4 +95,3 @@ export class EventQueue {
return ['message', 'reaction', 'edit', 'delete'].includes(event.type)
}
}

60 changes: 60 additions & 0 deletions src/agent/loop.pending-events.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
import { mkdtempSync, rmSync } from 'fs'
import { tmpdir } from 'os'
import { join } from 'path'
import { afterEach, describe, expect, it } from 'vitest'
import type { Event } from '../types.js'
import { EventQueue } from './event-queue.js'
import { AgentLoop } from './loop.js'

describe('AgentLoop active-channel event buffering', () => {
const temporaryDirectories: string[] = []

afterEach(() => {
for (const directory of temporaryDirectories.splice(0)) {
rmSync(directory, { recursive: true, force: true })
}
})

it('requeues events received while the channel is active', async () => {
const queue = new EventQueue()
const cacheDirectory = mkdtempSync(join(tmpdir(), 'chapterx-agent-loop-'))
temporaryDirectories.push(cacheDirectory)
const loop = new AgentLoop(
'bot',
queue,
{} as any,
{} as any,
{} as any,
{} as any,
{} as any,
cacheDirectory
)
const event: Event = {
type: 'message',
channelId: 'channel',
guildId: 'guild',
data: { id: 'message' },
timestamp: new Date('2026-09-01T12:00:00Z'),
}

;(loop as any).activeChannels.add('channel')
await (loop as any).processBatch([event])

expect(queue.isEmpty()).toBe(true)
expect((loop as any).pendingChannelEvents.get('channel')).toEqual([event])

const newerEvent: Event = {
...event,
channelId: 'other-channel',
data: { id: 'newer-message' },
}
queue.push(newerEvent)

;(loop as any).completeChannelActivation('channel')

expect((loop as any).activeChannels.has('channel')).toBe(false)
expect((loop as any).pendingChannelEvents.has('channel')).toBe(false)
expect(queue.pollBatch()).toEqual([event])
expect(queue.pollBatch()).toEqual([newerEvent])
})
})
44 changes: 33 additions & 11 deletions src/agent/loop.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,10 @@ export class AgentLoop {
private botMessageIds = new Set<string>() // Track bot's own message IDs
private mcpInitialized = false
private activeChannels = new Set<string>() // Track channels currently being processed
// Events that arrive during an activation must be replayed after it finishes.
// Coalescing per channel retains every trigger and command without creating
// a separate in-memory reply job for every arrival.
private pendingChannelEvents = new Map<string, Event[]>()
private activationStore: ActivationStore
private cacheDir: string
private somaClient?: SomaClient // Optional Soma credit system client
Expand Down Expand Up @@ -494,6 +498,19 @@ export class AgentLoop {
const firstEvent = events[0]
if (!firstEvent) return

const { channelId, guildId } = firstEvent

// The run loop keeps polling while activations execute asynchronously. Save
// active-channel events for a follow-up pass instead of discarding mentions,
// commands, reactions, edits, or deletes received in the meantime.
if (this.activeChannels.has(channelId)) {
const pending = this.pendingChannelEvents.get(channelId) || []
pending.push(...events)
this.pendingChannelEvents.set(channelId, pending)
logger.debug({ channelId, pendingEventCount: pending.length }, 'Channel active, deferred events')
return
}

// Profile: time from Discord event receipt to batch processing
const eventReceivedAt = (firstEvent as any).receivedAt
if (eventReceivedAt) {
Expand Down Expand Up @@ -523,8 +540,6 @@ export class AgentLoop {
}
logger.debug({ durationMs: Date.now() - shouldActivateStart }, 'shouldActivate completed')

const { channelId, guildId } = firstEvent

// Get triggering message ID for tool tracking (prefer non-system messages)
const triggeringEvent = this.findTriggeringMessageEvent(events)
const triggeringMessageId = triggeringEvent?.data?.id
Expand Down Expand Up @@ -629,12 +644,6 @@ export class AgentLoop {
logger.debug({ channelId }, 'Cancelled pending deferred activation due to new activity')
}

// Check if this channel is already being processed
if (this.activeChannels.has(channelId)) {
logger.debug({ channelId }, 'Channel already being processed, skipping')
return
}

// Mark channel as active and process asynchronously (don't await)
this.activeChannels.add(channelId)

Expand Down Expand Up @@ -662,7 +671,7 @@ export class AgentLoop {

if (somaCheckResult.status === 'blocked') {
// User doesn't have enough ichor - message already sent
this.activeChannels.delete(channelId)
this.completeChannelActivation(channelId)
return
}

Expand Down Expand Up @@ -756,9 +765,23 @@ export class AgentLoop {
}, 'Failed to handle activation')
})
.finally(() => {
this.activeChannels.delete(channelId)
this.completeChannelActivation(channelId)
})
}

/** Release a channel and put arrivals from its active turn back on the queue. */
private completeChannelActivation(channelId: string): void {
this.activeChannels.delete(channelId)

const pending = this.pendingChannelEvents.get(channelId)
if (!pending || pending.length === 0) return

this.pendingChannelEvents.delete(channelId)
// These events were polled before anything still in the shared queue, so
// restore them at the front to preserve arrival order.
this.queue.prepend(pending)
logger.debug({ channelId, pendingEventCount: pending.length }, 'Requeued deferred channel events')
}

private determineActivationReason(events: Event[]): {
reason: 'mention' | 'reply' | 'random' | 'm_command' | 'reaction' | 'timer',
Expand Down Expand Up @@ -4041,4 +4064,3 @@ export class AgentLoop {
)
}
}