mirror of
https://github.com/stablyai/orca.git
synced 2026-09-22 00:02:31 +00:00
fix(orchestration): revalidate an attempted Enter instead of resending it (#19911)
When a PTY retires mid-delivery, every staged message was marked undelivered, which made all of them redeliverable. That is right for a pointer whose Enter never fired, but an Enter that was already written may have landed: redelivering it types the same mail into the pane a second time. The Enter timer is cleared at the top of retirement, so a RESERVED or WRITE_ATTEMPTED pointer provably never submitted and is released. An ENTER_ATTEMPTED pointer is ambiguous and now stays at its phase for the resume path to revalidate, matching the policy mailbox-pointer-submit.ts already documents for an unverifiable settlement. Co-authored-by: Merge Sim <sim@local>
This commit is contained in:
co-authored by
Merge Sim
parent
a6e6de93c4
commit
6e9de5fa58
@@ -9,6 +9,10 @@ import {
|
||||
OrchestrationMailboxPointerState,
|
||||
type OrchestrationMailboxDeliveryFlight
|
||||
} from './mailbox-pointer-state'
|
||||
import {
|
||||
MAILBOX_POINTER_RESERVED,
|
||||
MAILBOX_POINTER_WRITE_ATTEMPTED
|
||||
} from './db/messages/mailbox-pointer-enter-state'
|
||||
import { resumePendingOrchestrationMailboxPointer } from './mailbox-pointer-resume'
|
||||
import { stageOrchestrationMailboxPointer } from './mailbox-pointer-stage'
|
||||
|
||||
@@ -163,7 +167,19 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
|
||||
clearTimeout(flight.enterTimer)
|
||||
}
|
||||
if (flight?.stagedMessageIds.length) {
|
||||
this.deps.getDb()?.markAsUndelivered(flight.stagedMessageIds)
|
||||
const db = this.deps.getDb()
|
||||
if (db && flight.processIncarnation) {
|
||||
// Why: the Enter timer was just cleared, so a reserved or merely-written pointer provably
|
||||
// never submitted and is released. An attempted Enter may already have landed, so it stays
|
||||
// at its phase for the resume path to revalidate rather than being sent a second time.
|
||||
db.releaseMailboxPointerEnter(
|
||||
flight.stagedMessageIds,
|
||||
{ ptyId, processIncarnation: flight.processIncarnation },
|
||||
[MAILBOX_POINTER_RESERVED, MAILBOX_POINTER_WRITE_ATTEMPTED]
|
||||
)
|
||||
} else {
|
||||
db?.markAsUndelivered(flight.stagedMessageIds)
|
||||
}
|
||||
}
|
||||
for (const mailboxHandle of releasedMailboxes) {
|
||||
this.redrive(mailboxHandle, true)
|
||||
|
||||
@@ -200,3 +200,72 @@ describe('mailbox pointer staging watermark', () => {
|
||||
db.close()
|
||||
})
|
||||
})
|
||||
|
||||
describe('retiring a pty mid-delivery', () => {
|
||||
// Why: an Enter that was already written may have landed. Releasing it would send the same
|
||||
// mail a second time, so only phases that provably never submitted become redeliverable.
|
||||
it('leaves an attempted Enter at its phase instead of making it redeliverable', async () => {
|
||||
vi.useFakeTimers()
|
||||
const db = new OrchestrationDb(':memory:')
|
||||
const settlements: ((settlement: WriteSettlement) => void)[] = []
|
||||
const writePty = vi.fn(
|
||||
() =>
|
||||
new Promise<WriteSettlement>((resolve) => {
|
||||
settlements.push(resolve)
|
||||
}) as unknown as WriteSettlement
|
||||
)
|
||||
try {
|
||||
const message = db.insertMessage({ from: 'a', to: 'run:run-1', subject: 'mail' })
|
||||
const delivery = new OrchestrationMailboxPointerDelivery(pointerDeps(db, writePty) as never)
|
||||
delivery.deliver(LEAF, { mailboxHandle: 'run:run-1' })
|
||||
|
||||
// Settle the pointer write, so the pane reaches WRITE_ATTEMPTED and arms the Enter.
|
||||
settlements[0]?.(WRITE_ACCEPTED)
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(db.getMessageById(message.id)?.pointer_enter_pending).toBe(2)
|
||||
|
||||
// Fire the Enter but never settle it: this is the ambiguous state.
|
||||
await vi.advanceTimersByTimeAsync(600)
|
||||
expect(db.getMessageById(message.id)?.pointer_enter_pending).toBe(3)
|
||||
|
||||
delivery.retirePty('pty-1')
|
||||
expect(db.getMessageById(message.id)).toMatchObject({
|
||||
pointer_enter_pending: 3,
|
||||
read: 0
|
||||
})
|
||||
} finally {
|
||||
db.close()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it('releases a pointer whose Enter never fired', async () => {
|
||||
vi.useFakeTimers()
|
||||
const db = new OrchestrationDb(':memory:')
|
||||
const settlements: ((settlement: WriteSettlement) => void)[] = []
|
||||
const writePty = vi.fn(
|
||||
() =>
|
||||
new Promise<WriteSettlement>((resolve) => {
|
||||
settlements.push(resolve)
|
||||
}) as unknown as WriteSettlement
|
||||
)
|
||||
try {
|
||||
const message = db.insertMessage({ from: 'a', to: 'run:run-1', subject: 'mail' })
|
||||
const delivery = new OrchestrationMailboxPointerDelivery(pointerDeps(db, writePty) as never)
|
||||
delivery.deliver(LEAF, { mailboxHandle: 'run:run-1' })
|
||||
settlements[0]?.(WRITE_ACCEPTED)
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(db.getMessageById(message.id)?.pointer_enter_pending).toBe(2)
|
||||
|
||||
delivery.retirePty('pty-1')
|
||||
expect(db.getMessageById(message.id)).toMatchObject({
|
||||
pointer_enter_pending: 0,
|
||||
read: 0,
|
||||
delivered_at: null
|
||||
})
|
||||
} finally {
|
||||
db.close()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
@@ -58,6 +58,7 @@ export function stageOrchestrationMailboxPointer<TWaiter extends OrchestrationMe
|
||||
return
|
||||
}
|
||||
const flight = args.state.beginFlight(ptyId)
|
||||
flight.processIncarnation = expectedTarget.processIncarnation
|
||||
flight.stagedMessageIds = args.messages.map((message) => message.id)
|
||||
try {
|
||||
if (
|
||||
|
||||
@@ -6,6 +6,8 @@ export type OrchestrationMailboxDeliveryFlight = {
|
||||
submitEnter: (() => void) | null
|
||||
deferredUntilIdle: boolean
|
||||
idleObservedWhileDeferred: boolean
|
||||
/** The incarnation that staged this flight, so retirement can name the rows it owns. */
|
||||
processIncarnation?: string
|
||||
}
|
||||
|
||||
export type ParkedOrchestrationMailboxDelivery = {
|
||||
|
||||
@@ -386,6 +386,7 @@ function makeOrchestrationDbStub(toHandle: () => string) {
|
||||
runMailbox,
|
||||
markAsDelivered,
|
||||
markAsUndelivered,
|
||||
releaseMailboxPointerEnter,
|
||||
stageMailboxPointerEnter,
|
||||
insert(subject: string, type: StoredMessageRow['type'] = 'status'): void {
|
||||
rows.push({
|
||||
@@ -759,7 +760,7 @@ describe('push-on-idle orchestration delivery absence gate', () => {
|
||||
await vi.advanceTimersByTimeAsync(500)
|
||||
expect(write.mock.calls.filter(([, data]) => data === '\r')).toHaveLength(0)
|
||||
expect(stub.stageMailboxPointerEnter).toHaveBeenCalledOnce()
|
||||
expect(stub.markAsUndelivered).toHaveBeenCalledOnce()
|
||||
expect(stub.releaseMailboxPointerEnter).toHaveBeenCalledOnce()
|
||||
expect(stub.rows[0].delivered_at).toBeNull()
|
||||
|
||||
// The replacement's own delivery starts a fresh flight and completes —
|
||||
@@ -804,7 +805,7 @@ describe('push-on-idle orchestration delivery absence gate', () => {
|
||||
await vi.advanceTimersByTimeAsync(500)
|
||||
expect(write.mock.calls.filter(([, data]) => data === '\r')).toHaveLength(0)
|
||||
expect(stub.stageMailboxPointerEnter).toHaveBeenCalledOnce()
|
||||
expect(stub.markAsUndelivered).toHaveBeenCalledOnce()
|
||||
expect(stub.releaseMailboxPointerEnter).toHaveBeenCalledOnce()
|
||||
// No stray settle flushed the parked trigger into the dead pty.
|
||||
expect(write).toHaveBeenCalledTimes(1)
|
||||
expect(stub.rows.every((row) => row.delivered_at === null)).toBe(true)
|
||||
@@ -835,7 +836,7 @@ describe('push-on-idle orchestration delivery absence gate', () => {
|
||||
await vi.advanceTimersByTimeAsync(500)
|
||||
expect(write.mock.calls.filter(([, data]) => data === '\r')).toHaveLength(0)
|
||||
expect(stub.stageMailboxPointerEnter).toHaveBeenCalledOnce()
|
||||
expect(stub.markAsUndelivered).toHaveBeenCalledOnce()
|
||||
expect(stub.releaseMailboxPointerEnter).toHaveBeenCalledOnce()
|
||||
expect(stub.rows[0].delivered_at).toBeNull()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
|
||||
Reference in New Issue
Block a user