mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(orchestration): strip delivery plumbing from send and reply receipts
This commit is contained in:
@@ -13,6 +13,10 @@ const INTERNAL_MESSAGE_COLUMNS = [
|
||||
|
||||
export type MailboxMessageReceipt = Omit<MessageRow, (typeof INTERNAL_MESSAGE_COLUMNS)[number]>
|
||||
|
||||
export function exposeMessage(message: MessageRow): MailboxMessageReceipt {
|
||||
return exposeMessages([message])[0]!
|
||||
}
|
||||
|
||||
export function exposeMessages(messages: MessageRow[]): MailboxMessageReceipt[] {
|
||||
return messages.map((message) => {
|
||||
const exposed: Partial<MessageRow> = { ...message }
|
||||
|
||||
@@ -9,6 +9,7 @@ import {
|
||||
readMutationReplayNudge,
|
||||
stripMutationReplayNudge
|
||||
} from '../../../orchestration-mutation-executor'
|
||||
import { exposeMessage } from './mailbox-message-receipt'
|
||||
import { recordReceiptBeforeNudge, replayMutationNudge } from './mutation-replay-nudge'
|
||||
import {
|
||||
ReplyParams,
|
||||
@@ -82,7 +83,7 @@ export const ORCHESTRATION_MESSAGE_METHODS: RpcMethod[] = [
|
||||
})
|
||||
const federated = db.getFederatedDispatch(question.dispatch_id)
|
||||
const receipt = {
|
||||
message: answered.message,
|
||||
message: exposeMessage(answered.message),
|
||||
question: answered.question,
|
||||
duplicate: answered.duplicate
|
||||
}
|
||||
@@ -120,7 +121,7 @@ export const ORCHESTRATION_MESSAGE_METHODS: RpcMethod[] = [
|
||||
runId: original.run_id
|
||||
})
|
||||
|
||||
const receipt = { message: reply }
|
||||
const receipt = { message: exposeMessage(reply) }
|
||||
return recordReceiptBeforeNudge(recordMutationReceipt, receipt, () =>
|
||||
runtime.notifyMessageArrived(reply.to_handle, reply.type)
|
||||
)
|
||||
|
||||
@@ -4,6 +4,7 @@ import { OrchestrationError } from '../../../../orchestration/orchestration-erro
|
||||
import { resolveGroupAddress } from '../../../../orchestration/groups'
|
||||
import { resolveBareOrchestrationRecipient } from './recipient-routing'
|
||||
import { legacyWorkerDeliveryContract } from '../routing'
|
||||
import { exposeMessages } from './mailbox-message-receipt'
|
||||
import { recordReceiptBeforeNudge } from './mutation-replay-nudge'
|
||||
import type { SendRecipientWarning } from './recipient-routing'
|
||||
import type { SendParams } from '../schemas'
|
||||
@@ -120,7 +121,7 @@ export async function sendGroupMessage(args: {
|
||||
resolution.ok ? (resolution.warning ? [resolution.warning] : []) : [resolution.warning]
|
||||
)
|
||||
const receipt = {
|
||||
messages,
|
||||
messages: exposeMessages(messages),
|
||||
recipients: messages.length,
|
||||
...(groupWarnings.length > 0 ? { warnings: groupWarnings } : {})
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import { bindCoordinatorMutationPayload } from '../../../../orchestration/dispat
|
||||
import { isDispatchMutationMessageType, parseMessageTaskId } from '../schemas'
|
||||
import type { SendParams } from '../schemas'
|
||||
import { legacyWorkerDeliveryContract } from '../routing'
|
||||
import { exposeMessage } from './mailbox-message-receipt'
|
||||
import { recordReceiptForPostCommitNudge } from './mutation-replay-nudge'
|
||||
import { sweepSettledWorkerResumeFences } from '../../settled-worker-resume-fence-sweep'
|
||||
import type { SendRecipientWarning } from './recipient-routing'
|
||||
@@ -93,7 +94,7 @@ export function sendPointToPointMessage(args: {
|
||||
const rejection =
|
||||
db.convertLifecycleMessageToRejection(msg.id, authority.code, authority.reason) ?? msg
|
||||
const receipt = withSendWarnings({
|
||||
message: rejection,
|
||||
message: exposeMessage(rejection),
|
||||
lifecycle: {
|
||||
action: 'rejected',
|
||||
code: authority.code,
|
||||
@@ -112,25 +113,30 @@ export function sendPointToPointMessage(args: {
|
||||
if (reconciled.action === 'suppressed') {
|
||||
return recordReceiptForPostCommitNudge(
|
||||
recordMutationReceipt,
|
||||
withSendWarnings({ message: msg }),
|
||||
withSendWarnings({ message: exposeMessage(msg) }),
|
||||
() => undefined
|
||||
)
|
||||
}
|
||||
if (reconciled.action === 'rejected') {
|
||||
const rejection = db.getMessageById(msg.id) ?? msg
|
||||
const receipt = withSendWarnings({ message: rejection, lifecycle: reconciled })
|
||||
const receipt = withSendWarnings({
|
||||
message: exposeMessage(rejection),
|
||||
lifecycle: reconciled
|
||||
})
|
||||
return recordReceiptForPostCommitNudge(recordMutationReceipt, receipt, () =>
|
||||
runtime.notifyMessageArrived(rejection.to_handle, rejection.type)
|
||||
)
|
||||
}
|
||||
const receipt = withSendWarnings(
|
||||
msg.type === 'worker_done' ? { message: msg, lifecycle: reconciled } : { message: msg }
|
||||
msg.type === 'worker_done'
|
||||
? { message: exposeMessage(msg), lifecycle: reconciled }
|
||||
: { message: exposeMessage(msg) }
|
||||
)
|
||||
return recordReceiptForPostCommitNudge(recordMutationReceipt, receipt, () =>
|
||||
runtime.notifyMessageArrived(msg.to_handle, msg.type)
|
||||
)
|
||||
}
|
||||
const receipt = withSendWarnings({ message: msg })
|
||||
const receipt = withSendWarnings({ message: exposeMessage(msg) })
|
||||
return recordReceiptForPostCommitNudge(recordMutationReceipt, receipt, () =>
|
||||
runtime.notifyMessageArrived(msg.to_handle, msg.type)
|
||||
)
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { RpcContext } from '../../../core'
|
||||
import type { OrchestrationDb } from '../../../../orchestration/db'
|
||||
import type { OrcaRuntimeService } from '../../../../orca-runtime'
|
||||
import type { RuntimeTerminalSummary } from '../../../../../../shared/runtime-types'
|
||||
import { createOrchestrationRpcHarness } from '../rpc-test-harness'
|
||||
|
||||
// The same delivery plumbing `check` already strips; a send/reply receipt is the same mailbox row.
|
||||
const INTERNAL_COLUMNS = [
|
||||
'read',
|
||||
'sequence',
|
||||
'sender_pane_key',
|
||||
'pointer_enter_pending',
|
||||
'pointer_pty_id',
|
||||
'pointer_process_incarnation'
|
||||
]
|
||||
|
||||
function terminalSummary(handle: string): RuntimeTerminalSummary {
|
||||
return {
|
||||
handle,
|
||||
ptyId: `pty_${handle}`,
|
||||
worktreeId: 'wt_default',
|
||||
worktreePath: '/tmp/wt',
|
||||
branch: 'main',
|
||||
tabId: 'tab_1',
|
||||
leafId: handle,
|
||||
title: null,
|
||||
connected: true,
|
||||
writable: true,
|
||||
lastOutputAt: null,
|
||||
preview: ''
|
||||
}
|
||||
}
|
||||
|
||||
describe('orchestration send and reply receipts', () => {
|
||||
const h = createOrchestrationRpcHarness()
|
||||
let db: OrchestrationDb
|
||||
let runtime: OrcaRuntimeService
|
||||
let ctx: RpcContext
|
||||
let activeRunId: string | undefined
|
||||
|
||||
afterEach(() => h.cleanup())
|
||||
|
||||
function setup(): void {
|
||||
;({ db, runtime, ctx, activeRunId } = h.setup())
|
||||
}
|
||||
|
||||
it('keeps delivery plumbing out of a point-to-point send receipt', async () => {
|
||||
setup()
|
||||
|
||||
const result = (await h.call(
|
||||
'orchestration.send',
|
||||
{ from: 'term_coord', to: `run:${activeRunId}`, subject: 'plumbing' },
|
||||
ctx
|
||||
)) as { message: Record<string, unknown> }
|
||||
|
||||
expect(result.message).toMatchObject({ subject: 'plumbing' })
|
||||
for (const column of INTERNAL_COLUMNS) {
|
||||
expect(result.message).not.toHaveProperty(column)
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps delivery plumbing out of a group send receipt', async () => {
|
||||
setup()
|
||||
const terminals = [terminalSummary('term_a'), terminalSummary('term_b')]
|
||||
vi.spyOn(runtime, 'listTerminals').mockResolvedValue({
|
||||
terminals,
|
||||
totalCount: terminals.length,
|
||||
truncated: false
|
||||
})
|
||||
vi.mocked(runtime.getTerminalPaneKey).mockImplementation((handle) => {
|
||||
const terminal = terminals.find((candidate) => candidate.handle === handle)
|
||||
return terminal ? `${terminal.tabId}:${terminal.leafId}` : null
|
||||
})
|
||||
|
||||
const result = (await h.call(
|
||||
'orchestration.send',
|
||||
{ from: 'term_a', to: '@all', subject: 'group plumbing' },
|
||||
ctx
|
||||
)) as { messages: Record<string, unknown>[] }
|
||||
|
||||
expect(result.messages).toHaveLength(1)
|
||||
for (const message of result.messages) {
|
||||
for (const column of INTERNAL_COLUMNS) {
|
||||
expect(message).not.toHaveProperty(column)
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps delivery plumbing out of a reply receipt', async () => {
|
||||
setup()
|
||||
const original = db.insertMessage({
|
||||
from: 'term_worker',
|
||||
to: `run:${activeRunId}`,
|
||||
subject: 'Need an answer'
|
||||
})
|
||||
|
||||
const result = (await h.call(
|
||||
'orchestration.reply',
|
||||
{ id: original.id, body: 'One durable answer', from: 'term_coord' },
|
||||
ctx
|
||||
)) as { message: Record<string, unknown> }
|
||||
|
||||
expect(result.message).toMatchObject({ subject: 'Re: Need an answer' })
|
||||
for (const column of INTERNAL_COLUMNS) {
|
||||
expect(result.message).not.toHaveProperty(column)
|
||||
}
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user