mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 16:02:56 +00:00
fix(test): give federation tests a real read-after-write sync barrier
`syncOrchestrationFederation()` coalesces onto an already-in-flight relay-tick sync, which may have pulled from the peer before the caller's mutation existed. Tests used it as a barrier, so `keeps a timed-out remote question resumable` could reply against a home DB that had never imported the worker's question: the reply failed with `Message not found`, no `to_worker` relay was enqueued, and the resume ask surfaced it 5s later as a spurious timeout. Add `syncFederationBarrier()`, which chains each active dispatch past the current round via `syncOrchestrationFederatedDispatchAfterCurrent`, and use it at every barrier-purpose sync site. The two tests whose subject is the sync machinery itself keep the raw call. Also assert the reply response, so a failed reply fails at the reply instead of masquerading as a timeout. Production is unaffected: `syncOrchestrationFederation` has no production callers, real read-after-write paths already use the after-current sync, and relay ticks retry every second.
This commit is contained in:
+3
-2
@@ -8,6 +8,7 @@ import type { RpcRequest } from '../../../core'
|
||||
import { RpcDispatcher } from '../../../dispatcher'
|
||||
import { fingerprintAuthenticatedPairingCredential } from '../../../orchestration-mutation-executor'
|
||||
import { ORCHESTRATION_METHODS } from '../../orchestration'
|
||||
import { syncFederationBarrier } from './federation-sync-barrier.test-support'
|
||||
|
||||
describe('orchestration federation control mail', () => {
|
||||
const homeToken = 'run-home-device-token'
|
||||
@@ -160,7 +161,7 @@ describe('orchestration federation control mail', () => {
|
||||
})
|
||||
expect(homeDb.listPendingFederationRelay(dispatchId, 'to_worker')).toHaveLength(1)
|
||||
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
const checked = await workerDispatcher.dispatch(checkRequest('check-imported'))
|
||||
|
||||
expect(checked).toMatchObject({
|
||||
@@ -252,7 +253,7 @@ describe('orchestration federation control mail', () => {
|
||||
settleRemoteOutcome: 'succeeded'
|
||||
})
|
||||
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
|
||||
expect(homeDb.getWorkerDispatch(dispatchId)?.state).toBe('succeeded')
|
||||
expect(workerDb.getUnreadMessages(`dispatch:${dispatchId}`)).toHaveLength(0)
|
||||
|
||||
+17
@@ -0,0 +1,17 @@
|
||||
import type { OrcaRuntimeService } from '../../../../orca-runtime'
|
||||
import type { OrchestrationDb } from '../../../../orchestration/db'
|
||||
|
||||
// A run-wide sync coalesces onto whatever relay tick is already in flight, and that tick may have
|
||||
// read the peer before the test's latest mutation existed. Chain past the current round instead so
|
||||
// awaiting the barrier really means "everything enqueued before this call has been exchanged".
|
||||
export async function syncFederationBarrier(
|
||||
runtime: OrcaRuntimeService,
|
||||
db: OrchestrationDb
|
||||
): Promise<void> {
|
||||
const dispatches = db.listActiveFederatedDispatches()
|
||||
await Promise.allSettled(
|
||||
dispatches.map((dispatch) =>
|
||||
runtime.syncOrchestrationFederatedDispatchAfterCurrent(dispatch.dispatch_id)
|
||||
)
|
||||
)
|
||||
}
|
||||
@@ -11,6 +11,7 @@ import { RpcDispatcher } from '../../../dispatcher'
|
||||
import { ORCHESTRATION_METHODS } from '../../orchestration'
|
||||
import { createFederationWorkerStartRequest as startRequest } from './federation-request.test-support'
|
||||
import { configureFederationWorkerRuntime } from './federation-runtime.test-support'
|
||||
import { syncFederationBarrier } from './federation-sync-barrier.test-support'
|
||||
|
||||
describe('orchestration federation', () => {
|
||||
const databases: OrchestrationDb[] = []
|
||||
@@ -310,7 +311,7 @@ describe('orchestration federation', () => {
|
||||
expect(sent).toMatchObject({ ok: true, result: { lifecycle: { action: 'completed' } } })
|
||||
expect(homeDb.getTask(task.id)?.status).toBe('completed')
|
||||
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
|
||||
expect(homeDb.getTask(task.id)?.status).toBe('completed')
|
||||
expect(homeDb.getWorkerDispatch(dispatch.id)?.state).toBe('succeeded')
|
||||
@@ -362,7 +363,7 @@ describe('orchestration federation', () => {
|
||||
).toHaveLength(1)
|
||||
)
|
||||
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
const question = homeDb
|
||||
.getRunMailboxHistory(task.run_id, 10)
|
||||
.find((message) => message.type === 'question')
|
||||
@@ -383,7 +384,7 @@ describe('orchestration federation', () => {
|
||||
}
|
||||
})
|
||||
expect(reply).toMatchObject({ ok: true, result: { question: { status: 'answered' } } })
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
|
||||
await expect(ask).resolves.toMatchObject({
|
||||
ok: true,
|
||||
@@ -426,8 +427,8 @@ describe('orchestration federation', () => {
|
||||
})
|
||||
const questionId = (timedOut as { result: { messageId: string } }).result.messageId
|
||||
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await homeDispatcher.dispatch({
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
const lateReply = await homeDispatcher.dispatch({
|
||||
id: 'rpc_home_late_reply',
|
||||
authToken: 'coordinator-token',
|
||||
orchestrationContractVersion: ORCHESTRATION_CONTRACT_VERSION,
|
||||
@@ -435,6 +436,8 @@ describe('orchestration federation', () => {
|
||||
method: 'orchestration.reply',
|
||||
params: { id: questionId, body: 'yes', from: 'term_coord' }
|
||||
})
|
||||
// A rejected reply enqueues no relay, which would only surface as the resume timing out.
|
||||
expect(lateReply).toMatchObject({ ok: true, result: { question: { status: 'answered' } } })
|
||||
restartWorkerRuntime()
|
||||
const resumed = workerDispatcher.dispatch({
|
||||
id: 'rpc_remote_ask_resume',
|
||||
@@ -445,7 +448,7 @@ describe('orchestration federation', () => {
|
||||
method: 'orchestration.ask',
|
||||
params: { from: 'term_windows_worker', resume: questionId, timeoutMs: 5_000 }
|
||||
})
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
|
||||
await expect(resumed).resolves.toMatchObject({
|
||||
ok: true,
|
||||
@@ -476,8 +479,8 @@ describe('orchestration federation', () => {
|
||||
loseNextAckResponse = true
|
||||
const remoteCall = vi.spyOn(homeRuntime, 'callOrchestrationWorkerServer')
|
||||
|
||||
await expect(homeRuntime.syncOrchestrationFederation()).resolves.toBeUndefined()
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await expect(syncFederationBarrier(homeRuntime, homeDb)).resolves.toBeUndefined()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
|
||||
expect(
|
||||
homeDb
|
||||
@@ -612,7 +615,7 @@ describe('orchestration federation', () => {
|
||||
it('treats a worker runtime ID change as an epoch, not a new server', async () => {
|
||||
const task = createHomeTask()
|
||||
await homeDispatcher.dispatch(startRequest(task.id))
|
||||
await homeRuntime.syncOrchestrationFederation()
|
||||
await syncFederationBarrier(homeRuntime, homeDb)
|
||||
vi.spyOn(homeRuntime, 'ensureOrchestrationFederationRelay').mockImplementation(() => {})
|
||||
const dispatch = homeDb.getDispatchContext(task.id)!
|
||||
const oldEpoch = homeDb.getFederatedDispatch(dispatch.id)?.remote_runtime_epoch
|
||||
|
||||
Reference in New Issue
Block a user