mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-24 19:35:28 -06:00
fix(mac): bound node worker crash retries (#114974)
This commit is contained in:
committed by
GitHub
parent
30fda068aa
commit
4b59f06d27
@@ -5,6 +5,8 @@ import OSLog
|
||||
|
||||
extension Notification.Name {
|
||||
static let openclawNodeHostWorkerFailed = Notification.Name("openclaw.node-host-worker.failed")
|
||||
static let openclawNodeHostWorkerRetryExhausted = Notification.Name(
|
||||
"openclaw.node-host-worker.retry-exhausted")
|
||||
}
|
||||
|
||||
struct MacNodeHostManifest: Equatable, Sendable {
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
import Foundation
|
||||
|
||||
struct MacNodeHostWorkerRetryPolicy: Sendable {
|
||||
struct Input: Equatable, Sendable {
|
||||
let command: [String]
|
||||
let configurationGeneration: UInt64
|
||||
}
|
||||
|
||||
enum UnexpectedExitDisposition: Equatable, Sendable {
|
||||
case retry(attempt: Int, delayNanoseconds: UInt64)
|
||||
case giveUp(unexpectedExitCount: Int)
|
||||
}
|
||||
|
||||
struct RetryBudgetExhausted: LocalizedError, Equatable, Sendable {
|
||||
let unexpectedExitCount: Int
|
||||
|
||||
var errorDescription: String? {
|
||||
"node-host worker stopped after \(self.unexpectedExitCount) unexpected exits"
|
||||
}
|
||||
}
|
||||
|
||||
struct RetryBackoffPending: LocalizedError, Equatable, Sendable {
|
||||
var errorDescription: String? {
|
||||
"node-host worker retry backoff is still pending"
|
||||
}
|
||||
}
|
||||
|
||||
static let defaultMaximumRetryCount = 5
|
||||
static let defaultInitialDelayNanoseconds: UInt64 = 1_000_000_000
|
||||
static let defaultMaximumDelayNanoseconds: UInt64 = 10_000_000_000
|
||||
|
||||
private let maximumRetryCount: Int
|
||||
private let initialDelayNanoseconds: UInt64
|
||||
private let maximumDelayNanoseconds: UInt64
|
||||
private var input: Input?
|
||||
private var unexpectedExitCount = 0
|
||||
private var exhausted = false
|
||||
|
||||
init(
|
||||
maximumRetryCount: Int = Self.defaultMaximumRetryCount,
|
||||
initialDelayNanoseconds: UInt64 = Self.defaultInitialDelayNanoseconds,
|
||||
maximumDelayNanoseconds: UInt64 = Self.defaultMaximumDelayNanoseconds)
|
||||
{
|
||||
precondition(maximumRetryCount >= 0)
|
||||
precondition(initialDelayNanoseconds > 0)
|
||||
precondition(maximumDelayNanoseconds >= initialDelayNanoseconds)
|
||||
self.maximumRetryCount = maximumRetryCount
|
||||
self.initialDelayNanoseconds = initialDelayNanoseconds
|
||||
self.maximumDelayNanoseconds = maximumDelayNanoseconds
|
||||
}
|
||||
|
||||
mutating func prepareForStart(_ input: Input) throws {
|
||||
self.adopt(input)
|
||||
if self.exhausted {
|
||||
throw RetryBudgetExhausted(unexpectedExitCount: self.unexpectedExitCount)
|
||||
}
|
||||
}
|
||||
|
||||
mutating func recordUnexpectedExit(for input: Input) -> UnexpectedExitDisposition {
|
||||
self.adopt(input)
|
||||
guard !self.exhausted else {
|
||||
return .giveUp(unexpectedExitCount: self.unexpectedExitCount)
|
||||
}
|
||||
|
||||
self.unexpectedExitCount += 1
|
||||
guard self.unexpectedExitCount <= self.maximumRetryCount else {
|
||||
self.exhausted = true
|
||||
return .giveUp(unexpectedExitCount: self.unexpectedExitCount)
|
||||
}
|
||||
return .retry(
|
||||
attempt: self.unexpectedExitCount,
|
||||
delayNanoseconds: self.retryDelay(for: self.unexpectedExitCount))
|
||||
}
|
||||
|
||||
mutating func reset() {
|
||||
self.input = nil
|
||||
self.unexpectedExitCount = 0
|
||||
self.exhausted = false
|
||||
}
|
||||
|
||||
private mutating func adopt(_ input: Input) {
|
||||
guard self.input != input else { return }
|
||||
self.input = input
|
||||
self.unexpectedExitCount = 0
|
||||
self.exhausted = false
|
||||
}
|
||||
|
||||
private func retryDelay(for attempt: Int) -> UInt64 {
|
||||
var delay = self.initialDelayNanoseconds
|
||||
for _ in 1..<attempt {
|
||||
if delay >= self.maximumDelayNanoseconds / 2 {
|
||||
return self.maximumDelayNanoseconds
|
||||
}
|
||||
delay = min(delay * 2, self.maximumDelayNanoseconds)
|
||||
}
|
||||
return delay
|
||||
}
|
||||
}
|
||||
@@ -97,10 +97,14 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
private var endpointRefreshTask: Task<Void, Never>?
|
||||
private var reconnectProbeTask: Task<Void, Never>?
|
||||
private var routeInvalidationTask: Task<Void, Never>?
|
||||
private var nodeHostWorkerRetryTask: Task<Void, Never>?
|
||||
private var endpointAttemptGeneration: UInt64 = 0
|
||||
private var routeAuthorityGeneration: UInt64 = 0
|
||||
private var completedRouteAuthorityGeneration: UInt64 = 0
|
||||
private var nodeHostWorkerConfigurationGeneration: UInt64 = 0
|
||||
private var nodeHostWorkerRetryTaskGeneration: UInt64 = 0
|
||||
private var pendingEndpoint: GatewayConnection.EndpointSnapshot?
|
||||
private var activeNodeHostWorkerInput: MacNodeHostWorkerRetryPolicy.Input?
|
||||
private var lastObservedPaused: Bool
|
||||
private var lastObservedComputerControlEnabled: Bool
|
||||
private let runtime: MacNodeRuntime
|
||||
@@ -112,6 +116,7 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
private let refreshEvents: AsyncStream<Void>
|
||||
private let refreshContinuation: AsyncStream<Void>.Continuation
|
||||
private var tlsSessionCache = MacNodeGatewayTLSSessionCache()
|
||||
private var nodeHostWorkerRetryPolicy: MacNodeHostWorkerRetryPolicy
|
||||
|
||||
override private convenience init() {
|
||||
let session = GatewayNodeSession()
|
||||
@@ -143,7 +148,8 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
observeNotifications: Bool = false,
|
||||
initialPaused: Bool? = nil,
|
||||
initialComputerControlEnabled: Bool? = nil,
|
||||
routeInvalidationHook: (@Sendable () async -> Void)? = nil)
|
||||
routeInvalidationHook: (@Sendable () async -> Void)? = nil,
|
||||
nodeHostWorkerRetryPolicy: MacNodeHostWorkerRetryPolicy = MacNodeHostWorkerRetryPolicy())
|
||||
{
|
||||
let refreshEvents = AsyncStream.makeStream(of: Void.self, bufferingPolicy: .bufferingNewest(1))
|
||||
self.session = session
|
||||
@@ -152,6 +158,7 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
self.presenceReporter = presenceReporter
|
||||
self.notificationCenter = notificationCenter
|
||||
self.routeInvalidationHook = routeInvalidationHook
|
||||
self.nodeHostWorkerRetryPolicy = nodeHostWorkerRetryPolicy
|
||||
self.refreshEvents = refreshEvents.stream
|
||||
self.refreshContinuation = refreshEvents.continuation
|
||||
self.lastObservedPaused = initialPaused ?? UserDefaults.standard.bool(forKey: pauseDefaultsKey)
|
||||
@@ -241,6 +248,7 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
self.endpointRefreshTask = nil
|
||||
self.reconnectProbeTask?.cancel()
|
||||
self.reconnectProbeTask = nil
|
||||
self.resetNodeHostWorkerRetryState()
|
||||
}
|
||||
|
||||
func setPreferredGatewayStableID(
|
||||
@@ -431,6 +439,19 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
// actually change instead of rereading config and TCC state every second.
|
||||
guard await refreshIterator.next() != nil else { return }
|
||||
} catch {
|
||||
if error is MacNodeHostWorkerRetryPolicy.RetryBackoffPending {
|
||||
// The lifecycle-owned delayed wake is the only event allowed
|
||||
// to admit this same worker input after an unexpected exit.
|
||||
guard await refreshIterator.next() != nil else { return }
|
||||
continue
|
||||
}
|
||||
if error is MacNodeHostWorkerRetryPolicy.RetryBudgetExhausted {
|
||||
// Only a new worker command or startup-scoped configuration
|
||||
// generation can re-arm a terminally exhausted worker.
|
||||
guard await refreshIterator.next() != nil else { return }
|
||||
retryDelay = 1_000_000_000
|
||||
continue
|
||||
}
|
||||
if let tlsError = error as? GatewayTLSValidationError,
|
||||
let attemptedEndpoint,
|
||||
await GatewayTLSRepairCoordinator.shared.repair(
|
||||
@@ -727,6 +748,21 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
completedRouteAuthorityGeneration: self.completedRouteAuthorityGeneration,
|
||||
isPaused: isPaused)
|
||||
}
|
||||
|
||||
func prepareNodeHostWorkerRetryForTesting(command: [String]) throws {
|
||||
guard self.nodeHostWorkerRetryTask == nil else {
|
||||
throw MacNodeHostWorkerRetryPolicy.RetryBackoffPending()
|
||||
}
|
||||
let input = MacNodeHostWorkerRetryPolicy.Input(
|
||||
command: command,
|
||||
configurationGeneration: self.nodeHostWorkerConfigurationGeneration)
|
||||
try self.nodeHostWorkerRetryPolicy.prepareForStart(input)
|
||||
self.activeNodeHostWorkerInput = input
|
||||
}
|
||||
|
||||
func handleNodeHostWorkerFailureForTesting() {
|
||||
self.handleNodeHostWorkerFailure()
|
||||
}
|
||||
#endif
|
||||
|
||||
private func cancelReconnectProbe() {
|
||||
@@ -742,12 +778,14 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
|
||||
@objc private nonisolated func nodeHostWorkerFailed(_: Notification) {
|
||||
Task { @MainActor [weak self] in
|
||||
self?.enqueueRouteInvalidation(yieldRefresh: true)
|
||||
self?.handleNodeHostWorkerFailure()
|
||||
}
|
||||
}
|
||||
|
||||
@objc private nonisolated func nodeHostConfigurationChanged(_: Notification) {
|
||||
Task { @MainActor [weak self] in
|
||||
self?.nodeHostWorkerConfigurationGeneration &+= 1
|
||||
self?.resetNodeHostWorkerRetryState()
|
||||
// Worker code, plugin availability, and its manifest are startup-scoped.
|
||||
// Replace the process before reconnecting so updates cannot leave a stale route.
|
||||
self?.enqueueRouteInvalidation(yieldRefresh: true, restartNodeHostWorker: true)
|
||||
@@ -784,6 +822,9 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
|
||||
private func startNodeHostWorkerIfConfigured() async throws -> MacNodeHostManifest? {
|
||||
guard let nodeHostWorker else { return nil }
|
||||
guard self.nodeHostWorkerRetryTask == nil else {
|
||||
throw MacNodeHostWorkerRetryPolicy.RetryBackoffPending()
|
||||
}
|
||||
let executable: String
|
||||
if let projectExecutable = CommandResolver.projectOpenClawExecutable() {
|
||||
executable = projectExecutable
|
||||
@@ -794,7 +835,65 @@ final class MacNodeModeCoordinator: NSObject {
|
||||
throw MacNodeHostWorker.WorkerError.unavailable(status.message)
|
||||
}
|
||||
}
|
||||
return try await nodeHostWorker.start(command: [executable, "node", "worker"])
|
||||
let command = [executable, "node", "worker"]
|
||||
let input = MacNodeHostWorkerRetryPolicy.Input(
|
||||
command: command,
|
||||
configurationGeneration: self.nodeHostWorkerConfigurationGeneration)
|
||||
try self.nodeHostWorkerRetryPolicy.prepareForStart(input)
|
||||
self.activeNodeHostWorkerInput = input
|
||||
return try await nodeHostWorker.start(command: command)
|
||||
}
|
||||
|
||||
private func handleNodeHostWorkerFailure() {
|
||||
guard let input = self.activeNodeHostWorkerInput else {
|
||||
self.logger.error("node-host worker exited without an active startup input")
|
||||
self.enqueueRouteInvalidation(yieldRefresh: false)
|
||||
return
|
||||
}
|
||||
|
||||
self.cancelNodeHostWorkerRetryTask()
|
||||
let invalidation = self.enqueueRouteInvalidation(yieldRefresh: false)
|
||||
switch self.nodeHostWorkerRetryPolicy.recordUnexpectedExit(for: input) {
|
||||
case let .retry(attempt, delayNanoseconds):
|
||||
self.nodeHostWorkerRetryTaskGeneration &+= 1
|
||||
let taskGeneration = self.nodeHostWorkerRetryTaskGeneration
|
||||
let delaySeconds = Double(delayNanoseconds) / 1_000_000_000
|
||||
self.logger.error(
|
||||
"node-host worker retry \(attempt, privacy: .public) in \(delaySeconds, privacy: .public)s")
|
||||
self.nodeHostWorkerRetryTask = Task { @MainActor [weak self] in
|
||||
await invalidation.value
|
||||
do {
|
||||
try await Task.sleep(nanoseconds: delayNanoseconds)
|
||||
} catch {
|
||||
return
|
||||
}
|
||||
guard let self,
|
||||
self.nodeHostWorkerRetryTaskGeneration == taskGeneration,
|
||||
self.activeNodeHostWorkerInput == input
|
||||
else { return }
|
||||
self.nodeHostWorkerRetryTask = nil
|
||||
self.refreshContinuation.yield()
|
||||
}
|
||||
case let .giveUp(unexpectedExitCount):
|
||||
self.logger.critical(
|
||||
"node-host worker gave up after \(unexpectedExitCount, privacy: .public) unexpected exits")
|
||||
self.notificationCenter.post(
|
||||
name: .openclawNodeHostWorkerRetryExhausted,
|
||||
object: self,
|
||||
userInfo: ["unexpectedExitCount": unexpectedExitCount])
|
||||
}
|
||||
}
|
||||
|
||||
private func cancelNodeHostWorkerRetryTask() {
|
||||
self.nodeHostWorkerRetryTaskGeneration &+= 1
|
||||
self.nodeHostWorkerRetryTask?.cancel()
|
||||
self.nodeHostWorkerRetryTask = nil
|
||||
}
|
||||
|
||||
private func resetNodeHostWorkerRetryState() {
|
||||
self.cancelNodeHostWorkerRetryTask()
|
||||
self.activeNodeHostWorkerInput = nil
|
||||
self.nodeHostWorkerRetryPolicy.reset()
|
||||
}
|
||||
|
||||
private func buildSessionBox(url: URL, tls: GatewayTLSRoute?) -> WebSocketSessionBox? {
|
||||
|
||||
@@ -32,6 +32,49 @@ private actor StubMacNodeHostWorker: MacNodeHostWorking {
|
||||
|
||||
@Suite(.serialized)
|
||||
struct MacNodeHostWorkerTests {
|
||||
@Test func `worker crash retry budget is bounded and exponentially delayed`() throws {
|
||||
let input = MacNodeHostWorkerRetryPolicy.Input(
|
||||
command: ["/usr/local/bin/openclaw", "node", "worker"],
|
||||
configurationGeneration: 4)
|
||||
var policy = MacNodeHostWorkerRetryPolicy(maximumRetryCount: 5)
|
||||
|
||||
try policy.prepareForStart(input)
|
||||
let dispositions = (0..<20).map { _ in policy.recordUnexpectedExit(for: input) }
|
||||
|
||||
#expect(dispositions.prefix(5) == [
|
||||
.retry(attempt: 1, delayNanoseconds: 1_000_000_000),
|
||||
.retry(attempt: 2, delayNanoseconds: 2_000_000_000),
|
||||
.retry(attempt: 3, delayNanoseconds: 4_000_000_000),
|
||||
.retry(attempt: 4, delayNanoseconds: 8_000_000_000),
|
||||
.retry(attempt: 5, delayNanoseconds: 10_000_000_000),
|
||||
])
|
||||
#expect(dispositions.dropFirst(5).allSatisfy {
|
||||
$0 == .giveUp(unexpectedExitCount: 6)
|
||||
})
|
||||
#expect(throws: MacNodeHostWorkerRetryPolicy.RetryBudgetExhausted.self) {
|
||||
try policy.prepareForStart(input)
|
||||
}
|
||||
}
|
||||
|
||||
@Test func `new worker input resets an exhausted crash retry budget`() throws {
|
||||
let original = MacNodeHostWorkerRetryPolicy.Input(
|
||||
command: ["/usr/local/bin/openclaw", "node", "worker"],
|
||||
configurationGeneration: 4)
|
||||
let updated = MacNodeHostWorkerRetryPolicy.Input(
|
||||
command: original.command,
|
||||
configurationGeneration: 5)
|
||||
var policy = MacNodeHostWorkerRetryPolicy(maximumRetryCount: 1)
|
||||
|
||||
try policy.prepareForStart(original)
|
||||
#expect(policy.recordUnexpectedExit(for: original) ==
|
||||
.retry(attempt: 1, delayNanoseconds: 1_000_000_000))
|
||||
#expect(policy.recordUnexpectedExit(for: original) == .giveUp(unexpectedExitCount: 2))
|
||||
|
||||
try policy.prepareForStart(updated)
|
||||
#expect(policy.recordUnexpectedExit(for: updated) ==
|
||||
.retry(attempt: 1, delayNanoseconds: 1_000_000_000))
|
||||
}
|
||||
|
||||
@Test func `worker allows a generous cold-start window`() async throws {
|
||||
#expect(MacNodeHostWorker.defaultStartupTimeout == 300)
|
||||
|
||||
|
||||
@@ -176,6 +176,71 @@ struct MacNodeModeCoordinatorTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Test @MainActor func `terminal worker failure is reported instead of scheduling another restart`() async throws {
|
||||
let worker = CoordinatorNodeHostWorkerProbe()
|
||||
let session = GatewayNodeSession()
|
||||
let notificationCenter = NotificationCenter()
|
||||
let coordinator = MacNodeModeCoordinator(
|
||||
session: session,
|
||||
runtime: MacNodeRuntime(nodeHostWorker: worker),
|
||||
nodeHostWorker: worker,
|
||||
notificationCenter: notificationCenter,
|
||||
nodeHostWorkerRetryPolicy: MacNodeHostWorkerRetryPolicy(maximumRetryCount: 0))
|
||||
|
||||
try coordinator.prepareNodeHostWorkerRetryForTesting(
|
||||
command: ["/usr/local/bin/openclaw", "node", "worker"])
|
||||
await confirmation("terminal worker failure") { confirmed in
|
||||
let observer = notificationCenter.addObserver(
|
||||
forName: .openclawNodeHostWorkerRetryExhausted,
|
||||
object: coordinator,
|
||||
queue: nil)
|
||||
{ notification in
|
||||
#expect(notification.userInfo?["unexpectedExitCount"] as? Int == 1)
|
||||
confirmed()
|
||||
}
|
||||
coordinator.handleNodeHostWorkerFailureForTesting()
|
||||
notificationCenter.removeObserver(observer)
|
||||
}
|
||||
await coordinator.waitForRouteInvalidationForTesting()
|
||||
#expect(await worker.stops() == 0)
|
||||
}
|
||||
|
||||
@Test @MainActor func `worker cannot restart before its crash backoff expires`() async throws {
|
||||
let worker = CoordinatorNodeHostWorkerProbe()
|
||||
let session = GatewayNodeSession()
|
||||
let coordinator = MacNodeModeCoordinator(
|
||||
session: session,
|
||||
runtime: MacNodeRuntime(nodeHostWorker: worker),
|
||||
nodeHostWorker: worker,
|
||||
observeNotifications: false,
|
||||
nodeHostWorkerRetryPolicy: MacNodeHostWorkerRetryPolicy(
|
||||
maximumRetryCount: 1,
|
||||
initialDelayNanoseconds: 50_000_000,
|
||||
maximumDelayNanoseconds: 50_000_000))
|
||||
let command = ["/usr/local/bin/openclaw", "node", "worker"]
|
||||
|
||||
try coordinator.prepareNodeHostWorkerRetryForTesting(command: command)
|
||||
coordinator.handleNodeHostWorkerFailureForTesting()
|
||||
#expect(throws: MacNodeHostWorkerRetryPolicy.RetryBackoffPending.self) {
|
||||
try coordinator.prepareNodeHostWorkerRetryForTesting(command: command)
|
||||
}
|
||||
|
||||
try await self.waitUntil("node-host worker retry backoff") {
|
||||
do {
|
||||
try await MainActor.run {
|
||||
try coordinator.prepareNodeHostWorkerRetryForTesting(command: command)
|
||||
}
|
||||
return true
|
||||
} catch is MacNodeHostWorkerRetryPolicy.RetryBackoffPending {
|
||||
return false
|
||||
} catch {
|
||||
Issue.record("unexpected retry admission error: \(error)")
|
||||
return false
|
||||
}
|
||||
}
|
||||
await coordinator.waitForRouteInvalidationForTesting()
|
||||
}
|
||||
|
||||
@Test func `paused node state requires route disconnect`() {
|
||||
#expect(MacNodeModeCoordinator.pausedStateRequiresDisconnect(true))
|
||||
#expect(!MacNodeModeCoordinator.pausedStateRequiresDisconnect(false))
|
||||
|
||||
Reference in New Issue
Block a user