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
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
# Dependency directories
node_modules
**/node_modules
.pnpm-store

# This monorepo uses pnpm — npm/yarn lockfiles should never be committed
package-lock.json
Expand Down
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ All notable changes to this project will be documented in this file.
## [Unreleased]

### Fixed
- Slack chat sessions now show Cyrus-branded in-flight status while work is underway, clear that status when the turn finishes or errors, and include guardrails for Slack streaming replies so a stream cannot mix text and task/progress modes. ([CYPACK-1367](https://linear.app/ceedar/issue/CYPACK-1367), [#1362](https://github.com/cyrusagents/cyrus/pull/1362))
- Forwarded and shared Slack messages are now included when you @mention Cyrus. Previously, forwarding a message (for example a Sentry alert) into a channel and @mentioning Cyrus passed along only your typed comment — the forwarded message's contents were dropped, so a forward with no comment gave Cyrus nothing to work with. The forwarded content is now part of the prompt. ([#1326](https://github.com/cyrusagents/cyrus/pull/1326))

### Changed
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
# Test Drive: CYPACK-1367 Slack In-Flight Status

**Date**: 2026-07-03
**Goal**: Verify the Slack chat session path still starts and completes after adding in-flight response status hooks.
**Test Repo**: `/tmp/f1-test-drive-cypack-1367`

## Verification Results

### F1 Setup
- [x] Test repository created with `apps/f1/f1 init-test-repo --path /tmp/f1-test-drive-cypack-1367`
- [x] F1 server started on `CYRUS_PORT=3600`
- [x] `apps/f1/f1 ping` and `apps/f1/f1 status` passed

### Slack Chat Path
- [x] Synthetic Slack chat event dispatched with `apps/f1/f1 start-chat-session`
- [x] Slack workspace directory created under `/tmp/cyrus-f1-*/slack-workspaces/`
- [x] Runner session started for the Slack event
- [x] Runner emitted a final `result`
- [x] No-token Slack status/reply paths skipped without blocking session completion

## Session Log

Initial run inherited a real `SLACK_BOT_TOKEN`, so the synthetic channel produced expected Slack API `channel_not_found` warnings for status, reaction, final reply, and status cleanup. The session still started, emitted `result`, and cleanup ran.

Clean no-token run:

```bash
env -u SLACK_BOT_TOKEN CYRUS_PORT=3600 CYRUS_REPO_PATH=/tmp/f1-test-drive-cypack-1367 bun run apps/f1/server.ts
CYRUS_PORT=3600 apps/f1/f1 start-chat-session --channel C_TEST_CHAN --user U_TEST_USER --text "Cyrus, give a one sentence reply"
```

Key observed events:

```text
Processing slack webhook: f1-1783118377.305
Cannot set Slack status: no slackBotToken available
[event:session_started]
[event:claude_session_id_assigned]
[event:message_emitted] {"messageType":"result"}
Session completed (subtype: success)
Cannot post Slack reply: no slackBotToken available
```

## Final Retrospective

F1 validates that the Slack ChatSessionHandler path remains functional and that missing Slack credentials do not prevent session startup or completion. F1 cannot visually verify Slack's rendered assistant status because it uses synthetic channels; the final acceptance criterion for branded Slack presentation still needs a live Slack thread with a real channel.
62 changes: 56 additions & 6 deletions packages/edge-worker/src/ChatSessionHandler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,12 @@ export interface ChatPlatformAdapter<TEvent> {
/** Fetch thread context as formatted string. Returns "" if not applicable */
fetchThreadContext(event: TEvent): Promise<string>;

/** Mark a platform thread as actively being handled, if supported. */
startInFlightResponse?(event: TEvent): Promise<void>;

/** Clear any active in-flight indicator for the platform thread, if supported. */
clearInFlightResponse?(event: TEvent): Promise<void>;

/** Post the agent's final response back to the platform */
postReply(event: TEvent, runner: IAgentRunner): Promise<void>;

Expand Down Expand Up @@ -198,6 +204,7 @@ export class ChatSessionHandler<TEvent> {
`Injecting follow-up prompt into running session ${existingSessionId} (thread ${threadKey})`,
);
this.enqueueReply(existingSessionId, event);
await this.startInFlightResponse(event);
existingRunner.addStreamMessage(taskInstructions);
} else {
// Runner can't accept mid-turn input (e.g. exec Codex). Queue the
Expand Down Expand Up @@ -339,6 +346,7 @@ export class ChatSessionHandler<TEvent> {
// stays open and the start() promise doesn't resolve until the
// whole session ends.
this.enqueueReply(sessionId, event);
await this.startInFlightResponse(event);
const startPromise =
runner.supportsStreamingInput && runner.startStreaming
? runner.startStreaming(userPrompt)
Expand All @@ -349,15 +357,16 @@ export class ChatSessionHandler<TEvent> {
`${this.adapter.platformName} session started: ${sessionInfo.sessionId}`,
);
})
.catch((error: unknown) => {
.catch(async (error: unknown) => {
this.logger.error(
`${this.adapter.platformName} session error for event ${eventId}`,
error instanceof Error ? error : new Error(String(error)),
);
// Runner died before emitting a final `result`. Drop any
// still-queued reply events for this session so a later
// resumeSession() doesn't pair them with a future turn.
this.clearPendingReplies(sessionId);
const clearedEvents = this.clearPendingReplies(sessionId);
await this.clearInFlightResponses(clearedEvents);
})
.finally(() => {
this.deps.onStateChange().catch((error: unknown) => {
Expand Down Expand Up @@ -452,6 +461,7 @@ export class ChatSessionHandler<TEvent> {
// warm sessions hold the streaming prompt open across turns so the
// start() promise only resolves when the whole session ends.
this.enqueueReply(sessionId, event);
await this.startInFlightResponse(event);
const startPromise =
runner.supportsStreamingInput && runner.startStreaming
? runner.startStreaming(taskInstructions)
Expand All @@ -462,12 +472,13 @@ export class ChatSessionHandler<TEvent> {
`${this.adapter.platformName} session resumed: ${sessionInfo.sessionId} (was ${resumeSessionId})`,
);
})
.catch((error: unknown) => {
.catch(async (error: unknown) => {
this.logger.error(
`${this.adapter.platformName} resume session error for ${sessionId}`,
error instanceof Error ? error : new Error(String(error)),
);
this.clearPendingReplies(sessionId);
const clearedEvents = this.clearPendingReplies(sessionId);
await this.clearInFlightResponses(clearedEvents);
});
}

Expand Down Expand Up @@ -499,6 +510,10 @@ export class ChatSessionHandler<TEvent> {
`Failed to post ${this.adapter.platformName} reply for session ${sessionId}`,
error instanceof Error ? error : new Error(String(error)),
);
} finally {
await this.clearInFlightResponses(
events.length > 0 ? events : [replyEvent],
);
}
// Fire-and-forget processed acknowledgement for every drained
// event (e.g., swap the receipt reaction) — runs even when
Expand Down Expand Up @@ -573,6 +588,40 @@ export class ChatSessionHandler<TEvent> {
this.lastReplyEvent.set(sessionId, event);
}

private async startInFlightResponse(event: TEvent): Promise<void> {
if (!this.adapter.startInFlightResponse) {
return;
}
try {
await this.adapter.startInFlightResponse(event);
} catch (error) {
this.logger.warn(
`Failed to start ${this.adapter.platformName} in-flight response status: ${error instanceof Error ? error.message : error}`,
);
}
}

private async clearInFlightResponses(events: TEvent[]): Promise<void> {
if (!this.adapter.clearInFlightResponse || events.length === 0) {
return;
}
const seen = new Set<string>();
for (const event of events) {
const key = this.adapter.getThreadKey(event);
if (seen.has(key)) {
continue;
}
seen.add(key);
try {
await this.adapter.clearInFlightResponse(event);
} catch (error) {
this.logger.warn(
`Failed to clear ${this.adapter.platformName} in-flight response status: ${error instanceof Error ? error.message : error}`,
);
}
}
}

private drainReplies(sessionId: string): TEvent[] {
const queue = this.pendingReplyEvents.get(sessionId);
if (!queue || queue.length === 0) return [];
Expand All @@ -586,14 +635,15 @@ export class ChatSessionHandler<TEvent> {
* resumeSession() on the same sessionId would pair the stale events with
* the first `result` of the new runner.
*/
private clearPendingReplies(sessionId: string): void {
private clearPendingReplies(sessionId: string): TEvent[] {
this.lastReplyEvent.delete(sessionId);
const queue = this.pendingReplyEvents.get(sessionId);
if (!queue || queue.length === 0) return;
if (!queue || queue.length === 0) return [];
this.logger.warn(
`Discarding ${queue.length} pending ${this.adapter.platformName} reply event(s) for session ${sessionId} after runner error`,
);
this.pendingReplyEvents.delete(sessionId);
return queue;
}

/**
Expand Down
43 changes: 43 additions & 0 deletions packages/edge-worker/src/SlackChatAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,17 @@ export const RECEIPT_REACTION = "eyes";
/** Reaction that replaces the receipt one once the agent finished its turn (✅) */
export const PROCESSED_REACTION = "white_check_mark";

/** Status Slack renders as "<App Name> <status>" while Cyrus is working. */
export const CYRUS_SLACK_ASSISTANT_STATUS = "is working through the request...";

/** Loading messages shown by Slack's assistant thread status indicator. */
export const CYRUS_SLACK_LOADING_MESSAGES = [
"Reading the thread...",
"Checking the workspace context...",
"Working through the request...",
"Verifying the result...",
] as const;

/**
* Slack implementation of ChatPlatformAdapter.
*
Expand Down Expand Up @@ -338,6 +349,38 @@ Supported mrkdwn syntax:
}
}

async startInFlightResponse(event: SlackWebhookEvent): Promise<void> {
const token = this.getSlackBotToken(event);
if (!token) {
this.logger.warn("Cannot set Slack status: no slackBotToken available");
return;
}

const threadTs = event.payload.thread_ts || event.payload.ts;
await new SlackMessageService().setAssistantThreadStatus({
token,
channel_id: event.payload.channel,
thread_ts: threadTs,
status: CYRUS_SLACK_ASSISTANT_STATUS,
loading_messages: [...CYRUS_SLACK_LOADING_MESSAGES],
});
}

async clearInFlightResponse(event: SlackWebhookEvent): Promise<void> {
const token = this.getSlackBotToken(event);
if (!token) {
return;
}

const threadTs = event.payload.thread_ts || event.payload.ts;
await new SlackMessageService().setAssistantThreadStatus({
token,
channel_id: event.payload.channel,
thread_ts: threadTs,
status: "",
});
}

async acknowledgeReceipt(event: SlackWebhookEvent): Promise<void> {
const token = this.getSlackBotToken(event);
if (!token) {
Expand Down
Loading
Loading