Skip to content
Draft
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
2 changes: 1 addition & 1 deletion .github/workflows/preview_sweep.yml
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ jobs:
QUEUE_PR_NUMBERS=$(curl -fsS -H "Authorization: Bearer ${CLOUDFLARE_API_TOKEN}" \
"https://api.cloudflare.com/client/v4/accounts/${CLOUDFLARE_ACCOUNT_ID}/queues" \
| jq -r '.result[].queue_name // empty' \
| sed -n 's/^tipbot-pending-tip-pr\([0-9][0-9]*\)\(-dlq\)\?$/\1/p' || true)
| sed -n 's/^tipbot-\(pending-tip\|slack-reaction\)-pr\([0-9][0-9]*\)\(-dlq\)\?$/\2/p' || true)
COMMENT_PR_NUMBERS=$(gh api "repos/${{ github.repository }}/issues/comments" --paginate \
--jq '.[] | select(.body | contains("tipbot-preview")) | .issue_url | split("/")[-1]' \
2>/dev/null || true)
Expand Down
22 changes: 22 additions & 0 deletions .github/workflows/production.yml
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,28 @@ jobs:
- name: Setup pnpm
uses: ./.github/actions/setup-pnpm

- name: Setup production Queues
run: |
QUEUES=$(node -e '
const text = require("node:fs").readFileSync("wrangler.jsonc", "utf8")
.replace(/(?<=^[^"]*(?:"[^"]*"[^"]*)*)\/\/.*$/gm, "")
.replace(/,(\s*[\]}])/g, "$1")
const queues = JSON.parse(text).env.production.queues
const names = new Set([
...queues.producers.map((producer) => producer.queue),
...queues.consumers.flatMap((consumer) =>
consumer.dead_letter_queue ? [consumer.queue, consumer.dead_letter_queue] : [consumer.queue],
),
])
for (const name of names) console.log(name)
')
for QUEUE in $QUEUES; do
pnpm exec wrangler queues create "$QUEUE" 2>/dev/null || true
done
env:
CLOUDFLARE_ACCOUNT_ID: ${{ secrets.CLOUDFLARE_ACCOUNT_ID }}
CLOUDFLARE_API_TOKEN: ${{ secrets.CLOUDFLARE_API_TOKEN }}

- name: Run migrations
run: pnpm exec wrangler d1 migrations apply tipbot --env production --remote
env:
Expand Down
64 changes: 40 additions & 24 deletions src/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1071,6 +1071,45 @@ export const api = new Hono<{
)
return new Response('Invalid signature', { status: 401 })

const reaction = (() => {
let payload: unknown
try {
payload = JSON.parse(body)
} catch {
return null
}
const parsed = z
.object({
authorizations: Slack.reactionEventSchema.shape.authorizations,
event: Slack.reactionEventSchema.omit({
authorizations: true,
event_id: true,
team_id: true,
}),
event_id: z.string().min(1),
team_id: z.string().min(1),
type: z.literal('event_callback'),
})
.safeParse(payload)
if (!parsed.success) return null
return Slack.reactionEventSchema.parse({
...parsed.data.event,
authorizations: parsed.data.authorizations,
event_id: parsed.data.event_id,
team_id: parsed.data.team_id,
})
})()
if (reaction) {
if (reaction.type === 'reaction_removed') return new Response('', { status: 200 })
try {
await c.env.SLACK_REACTION_QUEUE.send(reaction)
return new Response('', { status: 200 })
} catch (error) {
console.error('Failed to enqueue signed Slack reaction event:', error)
return new Response('Queue unavailable', { status: 503 })
}
}

const params = request.headers
.get('content-type')
?.includes('application/x-www-form-urlencoded')
Expand Down Expand Up @@ -1124,30 +1163,7 @@ export const api = new Hono<{
actions ?? '',
].join(':')
})()
if (interaction) return interaction
if (params) return null

// Dedupe reaction Events API retries by Slack event_id. Chat SDK already
// dedupes message events, so keep this scoped to reactions only.
let payload: unknown
try {
payload = JSON.parse(body)
} catch {
return null
}
const parsed = z
.object({
event: z.looseObject({ type: z.string().min(1) }).optional(),
event_id: z.string().min(1).optional(),
type: z.string().min(1).optional(),
})
.safeParse(payload)
if (!parsed.success) return null
if (parsed.data.type !== 'event_callback') return null
if (!parsed.data.event_id) return null
if (!['reaction_added', 'reaction_removed'].includes(parsed.data.event?.type ?? ''))
return null
return `slack:webhook:${parsed.data.event_id}`
return interaction
})()
if (duplicateKey) {
await Chat.getChat().initialize()
Expand Down
99 changes: 98 additions & 1 deletion src/api.workers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,9 @@ beforeEach(async () => {
function isExpectedApiWorkerLog(args: unknown[]) {
const message = typeof args[0] === 'string' ? args[0] : ''
return (
message.startsWith('Twitter webhook ') || message.startsWith('Twitter OAuth callback failed:')
message.startsWith('Failed to enqueue signed Slack reaction event:') ||
message.startsWith('Twitter webhook ') ||
message.startsWith('Twitter OAuth callback failed:')
)
}

Expand Down Expand Up @@ -1324,6 +1326,81 @@ describe('/api/chat/slack', () => {
expect(response.status).toBe(401)
})

test('durably enqueues signed Slack reaction additions before acknowledging', async () => {
const sendSpy = vi.spyOn(env.SLACK_REACTION_QUEUE, 'send').mockResolvedValue({
metadata: { metrics: { backlogBytes: 0, backlogCount: 0 } },
})
const body = slackReactionBody('reaction_added')

const response = await client.api.chat.slack.$post(
{},
{
headers: {
...(await createSlackHeaders(body, env.SLACK_SIGNING_SECRET)),
'content-type': 'application/json',
},
init: { body },
},
)

expect(response.status).toBe(200)
expect(sendSpy).toHaveBeenCalledWith({
event_id: 'EvReactionQueue',
event_ts: '1700000000.000002',
item: {
channel: Constants.slack.channelId,
ts: '1700000000.000001',
type: 'message',
},
item_user: Constants.slack.memberUserId,
reaction: 'money_with_wings',
team_id: Constants.slack.teamId,
type: 'reaction_added',
user: Constants.slack.adminUserId,
})
expect(executionCtx.waitUntil).not.toHaveBeenCalled()
})

test('acknowledges Slack reaction removals without enqueueing', async () => {
const sendSpy = vi.spyOn(env.SLACK_REACTION_QUEUE, 'send')
const body = slackReactionBody('reaction_removed')

const response = await client.api.chat.slack.$post(
{},
{
headers: {
...(await createSlackHeaders(body, env.SLACK_SIGNING_SECRET)),
'content-type': 'application/json',
},
init: { body },
},
)

expect(response.status).toBe(200)
expect(sendSpy).not.toHaveBeenCalled()
expect(executionCtx.waitUntil).not.toHaveBeenCalled()
})

test('asks Slack to retry when reaction enqueueing fails', async () => {
vi.spyOn(env.SLACK_REACTION_QUEUE, 'send').mockRejectedValue(new Error('Queue unavailable'))
const body = slackReactionBody('reaction_added')

const response = await client.api.chat.slack.$post(
{},
{
headers: {
...(await createSlackHeaders(body, env.SLACK_SIGNING_SECRET)),
'content-type': 'application/json',
},
init: { body },
},
)

expect(response.status).toBe(503)
expect(await response.text()).toBe('Queue unavailable')
expect(executionCtx.waitUntil).not.toHaveBeenCalled()
})

test('publishes workspace missing Home tab on app_home_opened', async () => {
const providerId = `T${Nanoid.generate()}`
await Chat.getChat().initialize()
Expand Down Expand Up @@ -3077,6 +3154,26 @@ function slackFetchBodyParams(body: BodyInit | null | undefined) {
return new URLSearchParams()
}

function slackReactionBody(type: 'reaction_added' | 'reaction_removed') {
return JSON.stringify({
event: {
event_ts: '1700000000.000002',
item: {
channel: Constants.slack.channelId,
ts: '1700000000.000001',
type: 'message',
},
item_user: Constants.slack.memberUserId,
reaction: 'money_with_wings',
type,
user: Constants.slack.adminUserId,
},
event_id: 'EvReactionQueue',
team_id: Constants.slack.teamId,
type: 'event_callback',
})
}

async function slackFetchBodyJson(
call: Parameters<typeof fetch>,
): Promise<Record<string, unknown>> {
Expand Down
Loading
Loading