diff --git a/apps/.i18n/native-source.json b/apps/.i18n/native-source.json index ab05a2a2c038..e2dcbac43f9b 100644 --- a/apps/.i18n/native-source.json +++ b/apps/.i18n/native-source.json @@ -32995,7 +32995,7 @@ }, { "kind": "conditional-branch", - "line": 77, + "line": 118, "path": "apps/macos/Sources/OpenClaw/GatewayProcessManager.swift", "source": "Stopped", "surface": "apple", @@ -33003,7 +33003,7 @@ }, { "kind": "conditional-branch", - "line": 78, + "line": 119, "path": "apps/macos/Sources/OpenClaw/GatewayProcessManager.swift", "source": "Starting…", "surface": "apple", @@ -33011,7 +33011,7 @@ }, { "kind": "conditional-branch", - "line": 87, + "line": 128, "path": "apps/macos/Sources/OpenClaw/GatewayProcessManager.swift", "source": "Failed: \\(reason)", "surface": "apple", @@ -33019,7 +33019,7 @@ }, { "kind": "conditional-branch", - "line": 657, + "line": 669, "path": "apps/macos/Sources/OpenClaw/GatewayProcessManager.swift", "source": "not linked", "surface": "apple", @@ -33027,7 +33027,7 @@ }, { "kind": "conditional-branch", - "line": 670, + "line": 682, "path": "apps/macos/Sources/OpenClaw/GatewayProcessManager.swift", "source": "unknown error", "surface": "apple", diff --git a/apps/macos/Sources/OpenClaw/GatewayProcessManager.swift b/apps/macos/Sources/OpenClaw/GatewayProcessManager.swift index 024cf39d3d2b..3982683be93f 100644 --- a/apps/macos/Sources/OpenClaw/GatewayProcessManager.swift +++ b/apps/macos/Sources/OpenClaw/GatewayProcessManager.swift @@ -53,10 +53,22 @@ final class GatewayProcessManager { } } - private struct LaunchAgentStartupContext { + private struct GatewayReadinessContext { + let purpose: GatewayReadinessPurpose let port: Int + let generation: UInt64 let readinessPID: Int32? let readinessRevision: UInt64 + let readinessCandidate: LaunchAgentReadinessCandidate? + let readinessFailure: LaunchAgentReadinessFailure? + let endpointPIDBeforeProbe: Int32? + let launchAgentInstalled: Bool + } + + private enum GatewayReadinessPurpose { + case attach + case launchd + case audit } private enum GatewayProbeFailureDisposition: Equatable { @@ -65,6 +77,35 @@ final class GatewayProcessManager { case fail } + private enum GatewayReadinessDeadlinePolicy { + case migration(window: TimeInterval, tolerance: TimeInterval) + case fixed(timeout: TimeInterval) + } + + private enum GatewayReadinessFailure { + case attachProbe(String) + case responsiveProbe(String) + case timeoutWithRepairEvidence(LaunchAgentReadinessFailure) + case deadlineWithoutRepairEvidence + + var reason: String { + switch self { + case let .attachProbe(reason), let .responsiveProbe(reason): reason + case .timeoutWithRepairEvidence: "Gateway did not start in time" + case .deadlineWithoutRepairEvidence: "Gateway did not become ready in time" + } + } + } + + private enum GatewayReadinessTerminal { + case ready( + instance: PortGuardian.Descriptor?, + startingPID: Int32?, + snapshot: HealthSnapshot?) + case superseded + case failed(GatewayReadinessFailure) + } + enum Status: Equatable { case stopped case starting @@ -505,9 +546,15 @@ final class GatewayProcessManager { // MARK: - Internals - private func isCurrentGatewayStart(_ generation: UInt64?) -> Bool { - guard let generation else { return true } - return self.desiredActive && self.gatewayStartGeneration == generation + private func isCurrentGatewayStart(_ generation: UInt64) -> Bool { + self.desiredActive && self.gatewayStartGeneration == generation + } + + private func isCurrentGatewayReadiness(_ context: GatewayReadinessContext) -> Bool { + !Task.isCancelled && + self.isCurrentGatewayStart(context.generation) && + self.launchAgentReadinessRevision == context.readinessRevision && + self.launchAgentReadinessCandidate == context.readinessCandidate } private func attachExistingGatewayAfterPendingDisable( @@ -527,78 +574,35 @@ final class GatewayProcessManager { /// If successful, mark status as attached and skip launchd startup. private func attachExistingGatewayIfAvailable( port requestedPort: Int? = nil, - startGeneration: UInt64? = nil) async -> Bool + startGeneration: UInt64) async -> Bool { let port = requestedPort ?? GatewayEnvironment.gatewayPort() let instance = await PortGuardian.shared.describe(port: port) guard self.isCurrentGatewayStart(startGeneration) else { return true } - let instanceText = instance.map { self.describe(instance: $0) } let hasListener = instance != nil if hasListener, - await !(self.profileOwnsGateway(instance, port: port)) + await !(self.profileOwnsGateway( + instance, + port: port, + startGeneration: startGeneration)) { return true } - let attemptAttach = { - try await self.probeGatewayHealth(timeoutMs: 2000) + let context = self.gatewayReadinessContext( + purpose: .attach, + port: port, + generation: startGeneration, + readinessPID: instance?.pid) + let terminal = await self.observeGatewayReadiness( + context: context, + deadlinePolicy: .fixed(timeout: hasListener ? 6.5 : 2)) + if !hasListener, case .failed = terminal { + self.existingGatewayDetails = nil + return false } - - for attempt in 0..<(hasListener ? 3 : 1) { - guard self.isCurrentGatewayStart(startGeneration) else { return true } - do { - let data = try await attemptAttach() - guard self.isCurrentGatewayStart(startGeneration) else { return true } - let attachedInstance = await PortGuardian.shared.describe(port: port) - guard self.isCurrentGatewayStart(startGeneration) else { return true } - if await !(self.profileOwnsGateway(attachedInstance, port: port)) { - return true - } - let snap = decodeHealthSnapshot(from: data) - let attachedInstanceText = attachedInstance.map { self.describe(instance: $0) } - let details = self.describe(details: attachedInstanceText, port: port, snap: snap) - let endpointPIDChanged = Self.gatewayPIDChanged( - from: self.lastObservedGatewayPID, - to: attachedInstance?.pid) || Self.gatewayPIDChanged( - from: instance?.pid, - to: attachedInstance?.pid) - self.existingGatewayDetails = details - self.setLaunchAgentReadinessState(candidate: nil, failure: nil) - self.clearLastFailure() - self.status = .attachedExisting(details: details) - self.appendLog("[gateway] using existing instance: \(details)\n") - self.logger.info("gateway using existing instance details=\(details)") - self.refreshControlChannelIfNeeded( - reason: "attach existing", - force: endpointPIDChanged) - self.lastObservedGatewayPID = attachedInstance?.pid ?? self.lastObservedGatewayPID - self.refreshLog() - return true - } catch { - guard self.isCurrentGatewayStart(startGeneration) else { return true } - if attempt < 2, hasListener { - try? await Task.sleep(nanoseconds: 250_000_000) - continue - } - - if hasListener { - let reason = self.describeAttachFailure(error, port: port, instance: instance) - self.existingGatewayDetails = instanceText - self.status = .failed(reason) - self.lastFailureReason = reason - self.appendLog("[gateway] existing listener on port \(port) but attach failed: \(reason)\n") - self.logger.warning("gateway attach failed reason=\(reason)") - return true - } - - // No reachable Gateway (and no listener) — fall through to launchd startup. - self.existingGatewayDetails = nil - return false - } - } - - self.existingGatewayDetails = nil - return false + let published = await self.publishGatewayReadinessTerminal(terminal, context: context) + return hasListener || published || !self.isCurrentGatewayStart(startGeneration) } static func profileAllowsExistingGatewayAttachment( @@ -611,21 +615,29 @@ final class GatewayProcessManager { return listenerPID == managedServicePID } - private func profileOwnsGateway(_ instance: PortGuardian.Descriptor?, port: Int) async -> Bool { + private func profileOwnsGateway( + _ instance: PortGuardian.Descriptor?, + port: Int, + startGeneration: UInt64) async -> Bool + { guard AppProfile.current.isActive else { return true } let managedPID = await GatewayLaunchAgentManager.runningGatewayPID() + guard self.isCurrentGatewayStart(startGeneration) else { return false } guard Self.profileAllowsExistingGatewayAttachment( profile: .current, listenerPID: instance?.pid, managedServicePID: managedPID) else { - await self.failProfilePortOwnership(port: port) + await self.failProfilePortOwnership( + port: port, + startGeneration: startGeneration) return false } return true } - private func failProfilePortOwnership(port: Int) async { + private func failProfilePortOwnership(port: Int, startGeneration: UInt64) async { + guard self.isCurrentGatewayStart(startGeneration) else { return } let message = "Gateway port \(port) is already owned by another process or OpenClaw profile. " + "Set gateway.port to a free port for profile \(AppProfile.current.name ?? "named")." self.recordProfilePortConflict(message) @@ -700,7 +712,26 @@ final class GatewayProcessManager { } extension GatewayProcessManager { - private func prepareLaunchdGatewayStart(startGeneration: UInt64) async -> LaunchAgentStartupContext? { + private func gatewayReadinessContext( + purpose: GatewayReadinessPurpose, + port: Int, + generation: UInt64, + readinessPID: Int32? = nil, + launchAgentInstalled: Bool = false) -> GatewayReadinessContext + { + GatewayReadinessContext( + purpose: purpose, + port: port, + generation: generation, + readinessPID: readinessPID, + readinessRevision: self.launchAgentReadinessRevision, + readinessCandidate: self.launchAgentReadinessCandidate, + readinessFailure: self.launchAgentReadinessFailure, + endpointPIDBeforeProbe: self.lastObservedGatewayPID, + launchAgentInstalled: launchAgentInstalled) + } + + private func prepareLaunchdGatewayStart(startGeneration: UInt64) async -> GatewayReadinessContext? { guard self.isCurrentGatewayStart(startGeneration) else { return nil } self.existingGatewayDetails = nil if GatewayLaunchAgentManager.isLaunchAgentWriteDisabled() { @@ -729,77 +760,91 @@ extension GatewayProcessManager { let readinessPID = await GatewayLaunchAgentManager.reusableLoadedGatewayPID(port: port) guard self.isCurrentGatewayStart(startGeneration) else { return nil } - return LaunchAgentStartupContext( + return self.gatewayReadinessContext( + purpose: .launchd, port: port, + generation: startGeneration, readinessPID: readinessPID, - readinessRevision: self.launchAgentReadinessRevision) + launchAgentInstalled: enableResult.installed) } private func enableLaunchdGateway(startGeneration: UInt64) async { guard let context = await self.prepareLaunchdGatewayStart(startGeneration: startGeneration) else { return } - await self.observeLaunchdGatewayReadiness(context: context, startGeneration: startGeneration) + await self.observeLaunchdGatewayReadiness(context: context) } private func observeLaunchdGatewayReadiness( - context: LaunchAgentStartupContext, - startGeneration: UInt64, + context: GatewayReadinessContext, readinessWindow: TimeInterval = 6, // Fresh installs keep probing through the same first-run migration budget as the CLI. firstInstallReadinessBudget: TimeInterval = GatewayLaunchAgentManager.startupMigrationTolerance) async + { + let terminal = await self.observeGatewayReadiness( + context: context, + deadlinePolicy: .migration( + window: readinessWindow, + tolerance: firstInstallReadinessBudget)) + _ = await self.publishGatewayReadinessTerminal(terminal, context: context) + } + + private func observeGatewayReadiness( + context: GatewayReadinessContext, + deadlinePolicy: GatewayReadinessDeadlinePolicy) async -> GatewayReadinessTerminal { let startedAt = Date() - var deadline = startedAt.addingTimeInterval(readinessWindow) - let finalProbeDeadline = startedAt.addingTimeInterval( - max(readinessWindow, firstInstallReadinessBudget)) + let initialWindow: TimeInterval + let finalProbeDeadline: Date + switch deadlinePolicy { + case let .migration(window, tolerance): + initialWindow = window + finalProbeDeadline = startedAt.addingTimeInterval(max(window, tolerance)) + case let .fixed(timeout): + initialWindow = timeout + finalProbeDeadline = startedAt.addingTimeInterval(timeout) + } + var deadline = startedAt.addingTimeInterval(initialWindow) var latestRetryDisposition: GatewayProbeFailureDisposition? var readinessPID = context.readinessPID var freshInstallGraceAuthorized = false var responsiveStartupProgressObserved = false + var latestProbeError: Error? readinessLoop: while true { - guard !Task.isCancelled, self.isCurrentGatewayStart(startGeneration) else { return } + guard self.isCurrentGatewayReadiness(context) else { return .superseded } while Date() >= deadline { guard deadline < finalProbeDeadline else { break readinessLoop } + guard case .migration = deadlinePolicy else { break readinessLoop } let extensionAuthorization = await self.authorizeReadinessExtension( context: context, - startGeneration: startGeneration, responsiveStartupProgressObserved: responsiveStartupProgressObserved, freshInstallGraceAuthorized: freshInstallGraceAuthorized, readinessPID: readinessPID) + guard self.isCurrentGatewayReadiness(context) else { return .superseded } guard extensionAuthorization.allowed else { break readinessLoop } readinessPID = extensionAuthorization.readinessPID freshInstallGraceAuthorized = true deadline = min( - deadline.addingTimeInterval(readinessWindow), + deadline.addingTimeInterval(initialWindow), finalProbeDeadline) guard Date() < finalProbeDeadline else { break readinessLoop } } do { let remainingMs = max(1, deadline.timeIntervalSinceNow * 1000) - _ = try await self.probeGatewayHealth(timeoutMs: min(1500, remainingMs)) - guard !Task.isCancelled else { return } + let data = try await self.probeGatewayHealth(timeoutMs: min(1500, remainingMs)) + guard self.isCurrentGatewayReadiness(context) else { return .superseded } let instance = await PortGuardian.shared.describe(port: context.port) - guard await self.profileOwnsGateway(instance, port: context.port) else { return } - guard self.publishLaunchdGatewayReady( + guard self.isCurrentGatewayReadiness(context) else { return .superseded } + return .ready( instance: instance, - context: context, - startGeneration: startGeneration) - else { return } - return + startingPID: readinessPID, + snapshot: decodeHealthSnapshot(from: data)) } catch { - if Task.isCancelled || !self.isCurrentGatewayStart(startGeneration) { - return - } + guard self.isCurrentGatewayReadiness(context) else { return .superseded } + latestProbeError = error switch self.probeFailureDisposition(error) { case .fail: - await self.finishResponsiveGatewayProbeFailure( - error, - port: context.port, - startGeneration: startGeneration, - expectedCandidate: self.launchAgentReadinessCandidate, - expectedReadinessRevision: context.readinessRevision) - return + return await self.gatewayProbeFailureTerminal(error, context: context) case .retryWithRepair: latestRetryDisposition = .retryWithRepair case .retryWithoutRepair: @@ -809,30 +854,67 @@ extension GatewayProcessManager { responsiveStartupProgressObserved = true } } - let retryDelay = min(0.4, max(0, deadline.timeIntervalSinceNow)) + let retryDelay = min(0.3, max(0, deadline.timeIntervalSinceNow)) if retryDelay > 0 { try? await Task.sleep(nanoseconds: UInt64(retryDelay * 1_000_000_000)) } } } - guard !Task.isCancelled, self.isCurrentGatewayStart(startGeneration) else { return } - if latestRetryDisposition == .retryWithRepair { - await self.finishLaunchAgentReadinessFailure( - port: context.port, - startingPID: readinessPID, - startGeneration: startGeneration, - expectedReadinessRevision: context.readinessRevision) - } else { - self.finishGatewayReadinessDeadlineWithoutRepair( - startGeneration: startGeneration, - expectedReadinessRevision: context.readinessRevision) + return await self.gatewayReadinessTimeout( + context: context, + policy: deadlinePolicy, + latestDisposition: latestRetryDisposition, + latestError: latestProbeError, + readinessPID: readinessPID) + } + + private func gatewayReadinessTimeout( + context: GatewayReadinessContext, + policy: GatewayReadinessDeadlinePolicy, + latestDisposition: GatewayProbeFailureDisposition?, + latestError: Error?, + readinessPID: Int32?) async -> GatewayReadinessTerminal + { + guard self.isCurrentGatewayReadiness(context) else { return .superseded } + if case .attach = context.purpose, let latestError { + return await self.gatewayProbeFailureTerminal(latestError, context: context) } + let migration = if case .migration = policy { + true + } else { + false + } + guard latestDisposition == .retryWithRepair else { + return migration ? .failed(.deadlineWithoutRepairEvidence) : .superseded + } + let failure: LaunchAgentReadinessFailure? = if migration || context.readinessCandidate != nil { + await self.resolveLaunchAgentReadinessFailure( + port: context.port, + startingPID: readinessPID) + } else { + context.readinessFailure + } + guard self.isCurrentGatewayReadiness(context) else { return .superseded } + guard let failure else { + return migration ? .failed(.deadlineWithoutRepairEvidence) : .superseded + } + return .failed(.timeoutWithRepairEvidence(failure)) + } + + private func gatewayProbeFailureTerminal( + _ error: Error, + context: GatewayReadinessContext) async -> GatewayReadinessTerminal + { + let instance = await PortGuardian.shared.describe(port: context.port) + guard self.isCurrentGatewayReadiness(context) else { return .superseded } + let reason = self.describeAttachFailure(error, port: context.port, instance: instance) + if case .attach = context.purpose { return .failed(.attachProbe(reason)) } + return .failed(.responsiveProbe(reason)) } private func authorizeReadinessExtension( - context: LaunchAgentStartupContext, - startGeneration: UInt64, + context: GatewayReadinessContext, responsiveStartupProgressObserved: Bool, freshInstallGraceAuthorized: Bool, readinessPID: Int32?) async -> (allowed: Bool, readinessPID: Int32?) @@ -840,146 +922,17 @@ extension GatewayProcessManager { if responsiveStartupProgressObserved || freshInstallGraceAuthorized { // One live response or verified launchd owner authorizes the bounded migration window; // repeating launchd status at every boundary would expand the wall-clock budget. - let isCurrent = !Task.isCancelled && - self.isCurrentGatewayStart(startGeneration) && - self.launchAgentReadinessRevision == context.readinessRevision - return (isCurrent, readinessPID) + return (self.isCurrentGatewayReadiness(context), readinessPID) } - let reusablePID = await self.currentInstallReusableLaunchdPID( - context: context, - startGeneration: startGeneration) - return (reusablePID != nil, reusablePID) - } - - private func currentInstallReusableLaunchdPID( - context: LaunchAgentStartupContext, - startGeneration: UInt64) async -> Int32? - { - guard self.isCurrentFreshInstallReadiness( - context: context, - startGeneration: startGeneration) - else { return nil } + guard self.launchAgentFreshInstallGeneration == context.generation, + self.isCurrentGatewayReadiness(context) + else { return (false, nil) } guard let reusablePID = await self.reusableLaunchdPIDOwningPort(port: context.port) else { - return nil + return (false, nil) } - guard self.isCurrentFreshInstallReadiness( - context: context, - startGeneration: startGeneration) - else { return nil } - return reusablePID - } - - private func isCurrentFreshInstallReadiness( - context: LaunchAgentStartupContext, - startGeneration: UInt64) -> Bool - { - !Task.isCancelled && - self.launchAgentFreshInstallGeneration == startGeneration && - self.isCurrentGatewayStart(startGeneration) && - self.launchAgentReadinessRevision == context.readinessRevision - } - - private func publishLaunchdGatewayReady( - instance: PortGuardian.Descriptor?, - context: LaunchAgentStartupContext, - startGeneration: UInt64) -> Bool - { - guard !Task.isCancelled else { return false } - guard self.isCurrentGatewayStart(startGeneration) else { return false } - guard self.launchAgentReadinessRevision == context.readinessRevision else { return false } - let details = instance.map { "pid \($0.pid)" } - let endpointPIDChanged = if let readinessPID = context.readinessPID, - let observedPID = instance?.pid - { - readinessPID != observedPID - } else { - false - } - let previouslyObservedPIDChanged = Self.gatewayPIDChanged( - from: self.lastObservedGatewayPID, - to: instance?.pid) - let launchAgentReplaced = self.launchAgentInstallGeneration == startGeneration || - endpointPIDChanged || - previouslyObservedPIDChanged - self.setLaunchAgentReadinessState(candidate: nil, failure: nil) - self.clearLastFailure() - self.status = .running(details: details) - self.logger.info("gateway started details=\(details ?? "ok")") - self.refreshControlChannelIfNeeded( - reason: "gateway started", - force: launchAgentReplaced) - self.lastObservedGatewayPID = instance?.pid ?? self.lastObservedGatewayPID - if self.launchAgentInstallGeneration == startGeneration { - self.launchAgentInstallGeneration = nil - } - if self.launchAgentFreshInstallGeneration == startGeneration { - self.launchAgentFreshInstallGeneration = nil - } - self.refreshLog() - return true - } - - private func finishLaunchAgentReadinessFailure( - port: Int, - startingPID: Int32?, - startGeneration: UInt64, - expectedCandidate: LaunchAgentReadinessCandidate? = nil, - candidateMustMatch: Bool = false, - expectedReadinessRevision: UInt64? = nil) async - { - let failure = await self.resolveLaunchAgentReadinessFailure( - port: port, - startingPID: startingPID) - guard !Task.isCancelled else { return } - guard self.isCurrentGatewayStart(startGeneration) else { return } - if let expectedReadinessRevision, - self.launchAgentReadinessRevision != expectedReadinessRevision - { - return - } - if candidateMustMatch, self.launchAgentReadinessCandidate != expectedCandidate { - return - } - self.setLaunchAgentReadinessState(candidate: nil, failure: failure) - self.status = .failed("Gateway did not start in time") - self.lastFailureReason = "launchd start timeout" - self.logger.warning("gateway start timed out") - } - - private func finishResponsiveGatewayProbeFailure( - _ error: Error, - port: Int, - startGeneration: UInt64, - expectedCandidate: LaunchAgentReadinessCandidate?, - expectedReadinessRevision: UInt64) async - { - let instance = await PortGuardian.shared.describe(port: port) - guard !Task.isCancelled else { return } - guard self.isCurrentGatewayStart(startGeneration) else { return } - guard self.launchAgentReadinessRevision == expectedReadinessRevision else { return } - guard self.launchAgentReadinessCandidate == expectedCandidate else { return } - let reason = self.describeAttachFailure(error, port: port, instance: instance) - self.setLaunchAgentReadinessState(candidate: nil, failure: nil) - self.status = .failed(reason) - self.lastFailureReason = reason - self.appendLog("[gateway] responsive health probe failed: \(reason)\n") - self.logger.warning("gateway responsive health probe failed reason=\(reason)") - } - - private func finishGatewayReadinessDeadlineWithoutRepair( - startGeneration: UInt64, - expectedReadinessRevision: UInt64) - { - guard self.isCurrentGatewayStart(startGeneration) else { return } - guard self.launchAgentReadinessRevision == expectedReadinessRevision else { return } - // Transient RPC/cancellation responses do not prove the endpoint is unreachable. End the - // startup cleanly, but do not retain a PID that would authorize destructive repair. - self.setLaunchAgentReadinessState( - candidate: nil, - failure: self.launchAgentReadinessFailure) - self.status = .failed("Gateway did not become ready in time") - self.lastFailureReason = "gateway readiness deadline elapsed" - self.logger.warning("gateway readiness deadline elapsed without endpoint failure") + let allowed = self.launchAgentFreshInstallGeneration == context.generation && + self.isCurrentGatewayReadiness(context) + return (allowed, allowed ? reusablePID : nil) } private func probeFailureDisposition(_ error: Error) -> GatewayProbeFailureDisposition { @@ -1057,62 +1010,18 @@ extension GatewayProcessManager { let startGeneration = self.gatewayStartGeneration if await self.observeCurrentGatewayStart(generation: startGeneration) == true { return true } guard !Task.isCancelled, self.isCurrentGatewayStart(startGeneration) else { return false } - let readinessCandidate = self.launchAgentReadinessCandidate - let readinessFailure = self.launchAgentReadinessFailure - let readinessRevision = self.launchAgentReadinessRevision - let readinessPort = readinessCandidate?.failure.port + let readinessPort = self.launchAgentReadinessCandidate?.failure.port ?? GatewayEnvironment.gatewayPort() - let deadline = Date().addingTimeInterval(timeout) - let endpointPIDBeforeProbe = self.lastObservedGatewayPID - var latestRetryDisposition: GatewayProbeFailureDisposition? - while Date() < deadline { - guard !Task.isCancelled else { return false } - guard self.isCurrentGatewayStart(startGeneration) else { return false } - do { - let remainingMs = max(1, deadline.timeIntervalSinceNow * 1000) - _ = try await self.probeGatewayHealth(timeoutMs: min(1500, remainingMs)) - guard !Task.isCancelled else { return false } - let instance = await PortGuardian.shared.describe(port: readinessPort) - guard await self.profileOwnsGateway(instance, port: readinessPort) else { return false } - return self.publishGatewayReadinessSuccess( - instance: instance, - startGeneration: startGeneration, - readinessCandidate: readinessCandidate, - readinessRevision: readinessRevision, - launchAgentInstalled: launchAgentInstalled, - endpointPIDBeforeProbe: endpointPIDBeforeProbe) - } catch { - if Task.isCancelled || !self.isCurrentGatewayStart(startGeneration) { - return false - } - switch self.probeFailureDisposition(error) { - case .fail: - await self.finishResponsiveGatewayProbeFailure( - error, - port: readinessPort, - startGeneration: startGeneration, - expectedCandidate: readinessCandidate, - expectedReadinessRevision: readinessRevision) - return false - case .retryWithRepair: - latestRetryDisposition = .retryWithRepair - case .retryWithoutRepair: - // A responsive transient invalidates older connection-failure evidence. - latestRetryDisposition = .retryWithoutRepair - } - let retryDelay = min(0.3, max(0, deadline.timeIntervalSinceNow)) - if retryDelay > 0 { - try? await Task.sleep(nanoseconds: UInt64(retryDelay * 1_000_000_000)) - } - } - } - await self.finishGatewayReadinessTimeout( - startGeneration: startGeneration, - readinessCandidate: readinessCandidate, - readinessFailure: readinessFailure, - readinessRevision: readinessRevision, - latestRetryDisposition: latestRetryDisposition) - return false + let context = self.gatewayReadinessContext( + purpose: .audit, + port: readinessPort, + generation: startGeneration, + readinessPID: self.launchAgentReadinessCandidate?.failure.pid, + launchAgentInstalled: launchAgentInstalled) + let terminal = await self.observeGatewayReadiness( + context: context, + deadlinePolicy: .fixed(timeout: timeout)) + return await self.publishGatewayReadinessTerminal(terminal, context: context) } private func observeCurrentGatewayStart(generation: UInt64) async -> Bool? { @@ -1129,90 +1038,105 @@ extension GatewayProcessManager { } } - private func publishGatewayReadinessSuccess( - instance: PortGuardian.Descriptor?, - startGeneration: UInt64, - readinessCandidate: LaunchAgentReadinessCandidate?, - readinessRevision: UInt64, - launchAgentInstalled: Bool, - endpointPIDBeforeProbe: Int32?) -> Bool + private func publishGatewayReadinessTerminal( + _ terminal: GatewayReadinessTerminal, + context: GatewayReadinessContext) async -> Bool { - guard !Task.isCancelled else { return false } - guard self.desiredActive, self.gatewayStartGeneration == startGeneration else { return false } - guard self.launchAgentReadinessRevision == readinessRevision else { return false } - guard self.launchAgentReadinessCandidate == readinessCandidate else { return false } - let details = instance.map { "pid \($0.pid)" } - let launchAgentReplaced = launchAgentInstalled || - self.launchAgentInstallGeneration == startGeneration - self.setLaunchAgentReadinessState(candidate: nil, failure: nil) - self.clearLastFailure() - if case .attachedExisting = self.status { - self.status = launchAgentReplaced - ? .running(details: details) - : .attachedExisting(details: details) - } else { - self.status = .running(details: details) + // Fixed audits without fresh endpoint evidence return `.superseded`, so a terminal failure + // from this generation is replaced only by a later result carrying endpoint evidence. + switch terminal { + case let .ready(instance, startingPID, snapshot): + guard await self.canPublishGatewayReadiness(instance: instance, context: context) else { + return false + } + let replaced = context.launchAgentInstalled || + self.launchAgentInstallGeneration == context.generation || + Self.gatewayPIDChanged(from: context.endpointPIDBeforeProbe, to: instance?.pid) || + Self.gatewayPIDChanged(from: startingPID, to: instance?.pid) + let details: String? + let refreshReason: String + switch context.purpose { + case .attach: + details = self.describe( + details: instance.map { self.describe(instance: $0) }, + port: context.port, + snap: snapshot) + refreshReason = "attach existing" + case .launchd, .audit: + details = instance.map { "pid \($0.pid)" } + refreshReason = "gateway readiness recovered" + } + self.setLaunchAgentReadinessState(candidate: nil, failure: nil) + self.clearLastFailure() + if case .attach = context.purpose { + self.existingGatewayDetails = details + self.status = .attachedExisting(details: details) + self.appendLog("[gateway] using existing instance: \(details ?? "unknown")\n") + } else if case .attachedExisting = self.status, !replaced { + self.status = .attachedExisting(details: details) + } else { + self.status = .running(details: details) + } + // A replaced process can leave the old socket briefly marked connected. Routine audits + // retain the connected channel; only replacement evidence forces refresh. + self.refreshControlChannelIfNeeded(reason: refreshReason, force: replaced) + self.lastObservedGatewayPID = instance?.pid ?? self.lastObservedGatewayPID + if self.launchAgentInstallGeneration == context.generation { + self.launchAgentInstallGeneration = nil + } + if self.launchAgentFreshInstallGeneration == context.generation { + self.launchAgentFreshInstallGeneration = nil + } + self.refreshLog() + return true + + case let .failed(terminalFailure): + let instance = await PortGuardian.shared.describe(port: context.port) + guard await self.canPublishGatewayReadiness(instance: instance, context: context) else { + return false + } + let retainedFailure: LaunchAgentReadinessFailure? = switch terminalFailure { + case let .timeoutWithRepairEvidence(failure): failure + case .attachProbe, .responsiveProbe, .deadlineWithoutRepairEvidence: nil + } + self.setLaunchAgentReadinessState(candidate: nil, failure: retainedFailure) + self.status = .failed(terminalFailure.reason) + switch terminalFailure { + case .attachProbe: + self.lastFailureReason = terminalFailure.reason + self.appendLog("[gateway] existing listener attach failed: \(terminalFailure.reason)\n") + case .responsiveProbe: + self.lastFailureReason = terminalFailure.reason + self.appendLog("[gateway] responsive health probe failed: \(terminalFailure.reason)\n") + case .timeoutWithRepairEvidence: + self.lastFailureReason = if case .launchd = context.purpose { + "launchd start timeout" + } else { + "gateway readiness timeout" + } + case .deadlineWithoutRepairEvidence: + // Transient responsive/cancellation outcomes never retain a PID for repair. + self.lastFailureReason = "gateway readiness deadline elapsed" + } + self.logger.warning("gateway readiness failed reason=\(terminalFailure.reason)") + return false + + case .superseded: + return false } - let endpointPIDChanged = Self.gatewayPIDChanged( - from: endpointPIDBeforeProbe, - to: instance?.pid) || Self.gatewayPIDChanged( - from: readinessCandidate?.failure.pid, - to: instance?.pid) - // A replaced process can leave the old socket briefly marked connected. Routine audits - // retain the connected channel; only replacement evidence forces refresh. - self.refreshControlChannelIfNeeded( - reason: "gateway readiness recovered", - force: launchAgentReplaced || endpointPIDChanged) - self.lastObservedGatewayPID = instance?.pid ?? self.lastObservedGatewayPID - if self.launchAgentInstallGeneration == startGeneration { - self.launchAgentInstallGeneration = nil - } - if self.launchAgentFreshInstallGeneration == startGeneration { - self.launchAgentFreshInstallGeneration = nil - } - self.refreshLog() - return true } - private func finishGatewayReadinessTimeout( - startGeneration: UInt64, - readinessCandidate: LaunchAgentReadinessCandidate?, - readinessFailure: LaunchAgentReadinessFailure?, - readinessRevision: UInt64, - latestRetryDisposition: GatewayProbeFailureDisposition?) async + private func canPublishGatewayReadiness( + instance: PortGuardian.Descriptor?, + context: GatewayReadinessContext) async -> Bool { - guard !Task.isCancelled else { return } - guard self.isCurrentGatewayStart(startGeneration) else { return } - guard self.launchAgentReadinessRevision == readinessRevision else { return } - guard self.launchAgentReadinessCandidate == readinessCandidate else { return } - self.appendLog("[gateway] readiness wait timed out\n") - guard latestRetryDisposition == .retryWithRepair else { - self.logger.warning("gateway readiness wait ended without endpoint failure evidence") - return - } - if let readinessCandidate, - readinessCandidate.generation == startGeneration - { - await self.finishLaunchAgentReadinessFailure( - port: readinessCandidate.failure.port, - startingPID: readinessCandidate.failure.pid, - startGeneration: startGeneration, - expectedCandidate: readinessCandidate, - candidateMustMatch: true, - expectedReadinessRevision: readinessRevision) - return - } - - self.setLaunchAgentReadinessState(candidate: nil, failure: readinessFailure) - if case .failed = self.status { - // Startup or persistence already published a concrete launchd/configuration error. - // A follow-up reachability timeout must not replace that actionable diagnosis. - self.logger.warning("gateway readiness wait timed out; preserving existing failure") - } else { - self.status = .failed("Gateway did not start in time") - self.lastFailureReason = "gateway readiness timeout" - self.logger.warning("gateway readiness wait timed out") - } + guard self.isCurrentGatewayReadiness(context) else { return false } + guard await self.profileOwnsGateway( + instance, + port: context.port, + startGeneration: context.generation) + else { return false } + return self.isCurrentGatewayReadiness(context) } private func probeGatewayHealth(timeoutMs: Double) async throws -> Data { @@ -1317,7 +1241,10 @@ extension GatewayProcessManager { } func _testAttachExistingGatewayIfAvailable(port: Int) async -> Bool { - await self.attachExistingGatewayIfAvailable(port: port) + self.desiredActive = true + return await self.attachExistingGatewayIfAvailable( + port: port, + startGeneration: self.gatewayStartGeneration) } func _testAttachExistingGatewayAfterPendingDisable(port: Int) async -> Bool { @@ -1344,11 +1271,22 @@ extension GatewayProcessManager { } func _testFinishLaunchAgentReadinessFailure(port: Int, startingPID: Int32?) async { - let startGeneration = self.gatewayStartGeneration - await self.finishLaunchAgentReadinessFailure( + let context = self.gatewayReadinessContext( + purpose: .launchd, port: port, - startingPID: startingPID, - startGeneration: startGeneration) + generation: self.gatewayStartGeneration, + readinessPID: startingPID) + let failure = await self.resolveLaunchAgentReadinessFailure( + port: port, + startingPID: startingPID) + let terminalFailure: GatewayReadinessFailure = if let failure { + .timeoutWithRepairEvidence(failure) + } else { + .deadlineWithoutRepairEvidence + } + _ = await self.publishGatewayReadinessTerminal( + .failed(terminalFailure), + context: context) } func _testClearLaunchAgentReadinessFailure() { @@ -1398,15 +1336,6 @@ extension GatewayProcessManager { self.gatewayStartTaskGeneration = nil } - func _testFinishGatewayReadinessTimeout() async { - await self.finishGatewayReadinessTimeout( - startGeneration: self.gatewayStartGeneration, - readinessCandidate: self.launchAgentReadinessCandidate, - readinessFailure: self.launchAgentReadinessFailure, - readinessRevision: self.launchAgentReadinessRevision, - latestRetryDisposition: .retryWithRepair) - } - func _testStartLaunchdGatewayReadiness( port: Int, pid: Int32, @@ -1417,16 +1346,17 @@ extension GatewayProcessManager { self.status = .starting self.gatewayStartGeneration &+= 1 let generation = self.gatewayStartGeneration - let readinessRevision = self.launchAgentReadinessRevision self.launchAgentInstallGeneration = generation self.launchAgentFreshInstallGeneration = generation + let context = self.gatewayReadinessContext( + purpose: .launchd, + port: port, + generation: generation, + readinessPID: pid, + launchAgentInstalled: true) self.beginGatewayStartTask(generation: generation) { [weak self] in await self?.observeLaunchdGatewayReadiness( - context: LaunchAgentStartupContext( - port: port, - readinessPID: pid, - readinessRevision: readinessRevision), - startGeneration: generation, + context: context, readinessWindow: readinessWindow, firstInstallReadinessBudget: firstInstallReadinessBudget) } diff --git a/apps/macos/Tests/OpenClawIPCTests/AppProfileSourceInvariantTests.swift b/apps/macos/Tests/OpenClawIPCTests/AppProfileSourceInvariantTests.swift index 7edefc258af0..7da3fdca8ae0 100644 --- a/apps/macos/Tests/OpenClawIPCTests/AppProfileSourceInvariantTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/AppProfileSourceInvariantTests.swift @@ -59,7 +59,26 @@ struct AppProfileSourceInvariantTests { let gatewayManager = try String( contentsOf: sourceRoot.appendingPathComponent("GatewayProcessManager.swift"), encoding: .utf8) - #expect(gatewayManager.components(separatedBy: "profileOwnsGateway(").count - 1 >= 5) + let publisherStart = try #require(gatewayManager.range( + of: "private func publishGatewayReadinessTerminal")) + let ownershipGuardStart = try #require(gatewayManager.range( + of: "private func canPublishGatewayReadiness", + range: publisherStart.lowerBound..