Repository navigation
fix(mothership): pass typed fork refusals through and stop overlapping retries #8564
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Closed
Closed
Changes from all commits
Commits
Show all changes
6 commits
Select commit
Hold shift + click to select a range
7f42068
fix(mothership): pass typed fork refusals through and stop overlappin…
waleedlatif1 3f142fa
fix(mothership): classify a fork refusal even when its body fails to …
waleedlatif1 0f0449d
fix(mothership): retry a fork only after a recognized socket failure
waleedlatif1 9d295c7
fix(core): treat EHOSTUNREACH as a retryable network failure
waleedlatif1 27b3177
Merge remote-tracking branch 'origin/staging' into fa-sync-8564
waleedlatif1 b0ebe2c
fix(mothership): show a fork refusal's own message
waleedlatif1 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,35 +1,73 @@ | ||
| import { isRetryableNetworkError } from '@/lib/core/errors/retryable-infrastructure' | ||
| import { OrchestrationError } from '@/lib/core/orchestration/types' | ||
| import { type ForkChatRequest, ForkChatResponse } from '@/lib/mothership/generated/protocol' | ||
| import { fetchGo } from '@/lib/mothership/request/go/fetch' | ||
| import { mothershipRequestHeaders } from '@/lib/mothership/request/headers' | ||
| import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url' | ||
|
|
||
| /** A lost copy acknowledgement retries the same destination and immutable request. */ | ||
| /** | ||
| * How long one attempt waits for the copy. The worker rolls a copy back once its caller's | ||
| * connection closes, so an attempt that times out leaves no conversation behind. | ||
| */ | ||
| const ATTEMPT_TIMEOUT_MS = 45_000 | ||
|
|
||
| /** Gateway failures: the request may never have reached a worker, so one more attempt is safe. */ | ||
| const RETRYABLE_STATUSES = new Set([502, 503, 504]) | ||
|
|
||
| /** The worker's typed refusals, each one the caller can act on. */ | ||
| function workerRefusal(status: number): Error { | ||
| if (status === 404) return new OrchestrationError('not_found', 'Chat not found') | ||
| if (status === 409) | ||
| return new OrchestrationError( | ||
| 'conflict', | ||
| 'The selected response has not finished. Retry the fork once it completes.' | ||
| ) | ||
| if (status === 413) | ||
| return new OrchestrationError( | ||
| 'payload_too_large', | ||
| 'This conversation is too long to fork. Fork from an earlier message.' | ||
| ) | ||
| return new Error('The conversation could not be copied. Retry the fork.') | ||
| } | ||
|
|
||
| type Attempt = { kind: 'receipt'; body: unknown } | { kind: 'status'; status: number } | ||
|
|
||
| /** | ||
| * A lost copy acknowledgement retries the same destination and immutable request; the | ||
| * worker answers a repeat with the first attempt's receipt. Only a refused, reset or dropped | ||
| * socket or a gateway status is retried: a timed-out attempt may still be copying. | ||
| */ | ||
| export async function copyWorkerConversation(request: ForkChatRequest): Promise<void> { | ||
| const baseURL = await getMothershipBaseURL({ userId: request.userId }) | ||
| const body = JSON.stringify(request) | ||
| for (let attempt = 0; attempt < 2; attempt++) { | ||
| for (let attempt = 0; ; attempt++) { | ||
| const canRetry = attempt === 0 | ||
| let outcome: Attempt | ||
| try { | ||
| const response = await fetchGo(`${baseURL}/api/chats/fork`, { | ||
| method: 'POST', | ||
| headers: mothershipRequestHeaders(), | ||
| body, | ||
| signal: AbortSignal.timeout(15_000), | ||
| signal: AbortSignal.timeout(ATTEMPT_TIMEOUT_MS), | ||
| spanName: 'sim → worker /api/chats/fork', | ||
| operation: 'fork_chat', | ||
| }) | ||
| if (!response.ok) { | ||
| if (response.status >= 500 && attempt === 0) { | ||
| await response.body?.cancel() | ||
| continue | ||
| } | ||
| throw new Error('The conversation could not be copied. Retry the fork.') | ||
| if (response.ok) outcome = { kind: 'receipt', body: await response.json() } | ||
| else { | ||
| outcome = { kind: 'status', status: response.status } | ||
| // Releasing the unread error body is best effort; the status is already the answer. | ||
| await response.body?.cancel().catch(() => undefined) | ||
| } | ||
| const receipt = ForkChatResponse.parse(await response.json()) | ||
| if (receipt.chatId !== request.newChatId) | ||
| throw new Error('The fork returned a different chat') | ||
| return | ||
| } catch (error) { | ||
| if (attempt === 1) throw error | ||
| if (canRetry && isRetryableNetworkError(error)) continue | ||
| throw error | ||
| } | ||
| if (outcome.kind === 'status') { | ||
| if (canRetry && RETRYABLE_STATUSES.has(outcome.status)) continue | ||
| throw workerRefusal(outcome.status) | ||
| } | ||
| const receipt = ForkChatResponse.parse(outcome.body) | ||
| if (receipt.chatId !== request.newChatId) throw new Error('The fork returned a different chat') | ||
| return | ||
| } | ||
| } | ||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.