From 150dc273f9fad8d9a8b538d6011a13e82d3e56f1 Mon Sep 17 00:00:00 2001 From: ember arlynx Date: Tue, 1 Sep 2026 15:12:56 -0400 Subject: [PATCH 1/2] fix(agent): retain events across active channels --- src/agent/event-queue.test.ts | 52 +++++++++++++++++++++++++++ src/agent/event-queue.ts | 11 ++++-- src/agent/loop.pending-events.test.ts | 52 +++++++++++++++++++++++++++ src/agent/loop.ts | 44 +++++++++++++++++------ 4 files changed, 145 insertions(+), 14 deletions(-) create mode 100644 src/agent/event-queue.test.ts create mode 100644 src/agent/loop.pending-events.test.ts diff --git a/src/agent/event-queue.test.ts b/src/agent/event-queue.test.ts new file mode 100644 index 0000000..cb2d260 --- /dev/null +++ b/src/agent/event-queue.test.ts @@ -0,0 +1,52 @@ +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]) + }) +}) diff --git a/src/agent/event-queue.ts b/src/agent/event-queue.ts index 7f431ba..24fee4f 100644 --- a/src/agent/event-queue.ts +++ b/src/agent/event-queue.ts @@ -17,7 +17,9 @@ export class EventQueue { /** * 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) { @@ -46,7 +48,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 } @@ -81,4 +87,3 @@ export class EventQueue { return ['message', 'reaction', 'edit', 'delete'].includes(event.type) } } - diff --git a/src/agent/loop.pending-events.test.ts b/src/agent/loop.pending-events.test.ts new file mode 100644 index 0000000..ae4398f --- /dev/null +++ b/src/agent/loop.pending-events.test.ts @@ -0,0 +1,52 @@ +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]) + + ;(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]) + }) +}) diff --git a/src/agent/loop.ts b/src/agent/loop.ts index 97fab4d..5528019 100644 --- a/src/agent/loop.ts +++ b/src/agent/loop.ts @@ -74,6 +74,10 @@ export class AgentLoop { private botMessageIds = new Set() // Track bot's own message IDs private mcpInitialized = false private activeChannels = new Set() // 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() private activationStore: ActivationStore private cacheDir: string private somaClient?: SomaClient // Optional Soma credit system client @@ -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) { @@ -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 @@ -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) @@ -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 } @@ -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) + for (const event of pending) { + this.queue.push(event) + } + logger.debug({ channelId, pendingEventCount: pending.length }, 'Requeued deferred channel events') + } private determineActivationReason(events: Event[]): { reason: 'mention' | 'reply' | 'random' | 'm_command' | 'reaction' | 'timer', @@ -4041,4 +4064,3 @@ export class AgentLoop { ) } } - From 4b6ab363f07d32961e4020a307a0f936952ec876 Mon Sep 17 00:00:00 2001 From: ember arlynx Date: Tue, 1 Sep 2026 15:18:26 -0400 Subject: [PATCH 2/2] fix(agent): preserve deferred event order --- src/agent/event-queue.test.ts | 13 +++++++++++++ src/agent/event-queue.ts | 8 ++++++++ src/agent/loop.pending-events.test.ts | 8 ++++++++ src/agent/loop.ts | 6 +++--- 4 files changed, 32 insertions(+), 3 deletions(-) diff --git a/src/agent/event-queue.test.ts b/src/agent/event-queue.test.ts index cb2d260..14e31dc 100644 --- a/src/agent/event-queue.test.ts +++ b/src/agent/event-queue.test.ts @@ -49,4 +49,17 @@ describe('EventQueue channel isolation', () => { 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]) + }) }) diff --git a/src/agent/event-queue.ts b/src/agent/event-queue.ts index 24fee4f..fdbfa8d 100644 --- a/src/agent/event-queue.ts +++ b/src/agent/event-queue.ts @@ -15,6 +15,14 @@ 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 consecutive events for one channel and the same type category. diff --git a/src/agent/loop.pending-events.test.ts b/src/agent/loop.pending-events.test.ts index ae4398f..072029f 100644 --- a/src/agent/loop.pending-events.test.ts +++ b/src/agent/loop.pending-events.test.ts @@ -43,10 +43,18 @@ describe('AgentLoop active-channel event buffering', () => { 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]) }) }) diff --git a/src/agent/loop.ts b/src/agent/loop.ts index 5528019..748465a 100644 --- a/src/agent/loop.ts +++ b/src/agent/loop.ts @@ -777,9 +777,9 @@ export class AgentLoop { if (!pending || pending.length === 0) return this.pendingChannelEvents.delete(channelId) - for (const event of pending) { - this.queue.push(event) - } + // 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') }