From 07d585bc3cf4fb2abd594fd53fce494406e0cbbf Mon Sep 17 00:00:00 2001 From: ensan-hcl Date: Sun, 6 Sep 2026 09:06:12 +0900 Subject: [PATCH 1/3] fix(input): handle converter XPC key events synchronously --- .../Core/XPC/ConverterClientEventRouter.swift | 21 +- Core/Sources/Core/XPC/DeadlineReply.swift | 48 ++ .../ConverterClientEventRouterTests.swift | 41 +- .../XPCTests/DeadlineReplyTests.swift | 29 ++ .../CandidateWindow/CandidateView.swift | 2 +- .../ConverterServerClient.swift | 421 ++++++++---------- .../azooKeyMacInputController.swift | 190 +++----- .../ThinClientInputPipelineTests.swift | 260 +++++++---- 8 files changed, 550 insertions(+), 462 deletions(-) create mode 100644 Core/Sources/Core/XPC/DeadlineReply.swift create mode 100644 Core/Tests/CoreTests/XPCTests/DeadlineReplyTests.swift diff --git a/Core/Sources/Core/XPC/ConverterClientEventRouter.swift b/Core/Sources/Core/XPC/ConverterClientEventRouter.swift index f7a5495c..e1986de2 100644 --- a/Core/Sources/Core/XPC/ConverterClientEventRouter.swift +++ b/Core/Sources/Core/XPC/ConverterClientEventRouter.swift @@ -4,17 +4,16 @@ import Foundation /// /// 変換状態の本体は ConverterServer が所有する。Client は Server が最後に返した /// 読み取り専用の状態を使い、明らかな application shortcut を同期的に通す。 -/// Server の応答待ちがある間は状態が進んでいる可能性があるため、Command shortcut -/// 以外を保守的に consume し、生のキー入力が application へ漏れることを防ぐ。 +/// キーごとに Server の応答を同期的に反映するため、未応答のキーを推測する必要はない。 public enum ConverterClientEventDisposition: Sendable, Equatable { case sendToServer case fallthroughToApplication + case insertText(String) } public struct ConverterClientEventRoutingContext: Sendable, Equatable { public var acknowledgedInputState: ConverterInputState public var acknowledgedInputLanguage: InputLanguage - public var hasPendingKeyEvents: Bool public var liveConversionEnabled: Bool public var enableDebugWindow: Bool public var enableSuggestion: Bool @@ -23,7 +22,6 @@ public struct ConverterClientEventRoutingContext: Sendable, Equatable { public init( acknowledgedInputState: ConverterInputState = .none, acknowledgedInputLanguage: InputLanguage = .japanese, - hasPendingKeyEvents: Bool = false, liveConversionEnabled: Bool = true, enableDebugWindow: Bool = false, enableSuggestion: Bool = false, @@ -31,7 +29,6 @@ public struct ConverterClientEventRoutingContext: Sendable, Equatable { ) { self.acknowledgedInputState = acknowledgedInputState self.acknowledgedInputLanguage = acknowledgedInputLanguage - self.hasPendingKeyEvents = hasPendingKeyEvents self.liveConversionEnabled = liveConversionEnabled self.enableDebugWindow = enableDebugWindow self.enableSuggestion = enableSuggestion @@ -49,13 +46,6 @@ public enum ConverterClientEventRouter { return .fallthroughToApplication } - // 未応答イベントがある場合、acknowledgedInputState は古い可能性がある。 - // ここで fallthrough するとタイムアウトした文字が英字として漏れるため、 - // Server が順番に処理できるようイベントを consume する。 - if context.hasPendingKeyEvents { - return .sendToServer - } - let inputState = context.acknowledgedInputState.inputState let userAction = UserAction.getUserAction( eventCore: event, @@ -73,6 +63,13 @@ public enum ConverterClientEventRouter { if case .fallthrough = action { return .fallthroughToApplication } + if context.acknowledgedInputLanguage == .english, + context.acknowledgedInputState == .none, + case .insertWithoutMarkedText(let text) = action { + // 通常の直接入力は application に任せる。円記号・バックスラッシュ等、 + // azooKey の設定による置き換えが必要な場合だけ、その場で挿入する。 + return text == event.characters ? .fallthroughToApplication : .insertText(text) + } return .sendToServer } } diff --git a/Core/Sources/Core/XPC/DeadlineReply.swift b/Core/Sources/Core/XPC/DeadlineReply.swift new file mode 100644 index 00000000..be78ad16 --- /dev/null +++ b/Core/Sources/Core/XPC/DeadlineReply.swift @@ -0,0 +1,48 @@ +import Foundation + +/// 別キューから届く応答を期限付きで待つ。期限後の応答は採用しない。 +public final class DeadlineReply: @unchecked Sendable { + private enum State { + case waiting + case completed(Value) + case timedOut + } + + public let deadline: DispatchTime + private let lock = NSLock() + private let semaphore = DispatchSemaphore(value: 0) + private var state: State = .waiting + + public init(timeout: TimeInterval) { + self.deadline = .now() + max(0, timeout) + } + + @discardableResult + public func complete(_ value: Value) -> Bool { + lock.lock() + defer { lock.unlock() } + guard case .waiting = state else { + return false + } + guard DispatchTime.now() < deadline else { + state = .timedOut + semaphore.signal() + return false + } + state = .completed(value) + semaphore.signal() + return true + } + + /// メインキューを回さず待つ。complete は待機中のキューとは別のキューから呼ぶこと。 + public func wait() -> Value? { + _ = semaphore.wait(timeout: deadline) + lock.lock() + defer { lock.unlock() } + if case .completed(let value) = state { + return value + } + state = .timedOut + return nil + } +} diff --git a/Core/Tests/CoreTests/XPCTests/ConverterClientEventRouterTests.swift b/Core/Tests/CoreTests/XPCTests/ConverterClientEventRouterTests.swift index db2022d5..506a3517 100644 --- a/Core/Tests/CoreTests/XPCTests/ConverterClientEventRouterTests.swift +++ b/Core/Tests/CoreTests/XPCTests/ConverterClientEventRouterTests.swift @@ -4,15 +4,13 @@ import Testing private func disposition( event: KeyEventCore, state: ConverterInputState = .none, - language: InputLanguage = .japanese, - hasPendingKeyEvents: Bool = false + language: InputLanguage = .japanese ) -> ConverterClientEventDisposition { ConverterClientEventRouter.disposition( event: event, context: .init( acknowledgedInputState: state, - acknowledgedInputLanguage: language, - hasPendingKeyEvents: hasPendingKeyEvents + acknowledgedInputLanguage: language ) ) } @@ -43,7 +41,7 @@ private func disposition( ) } -@Test func backspaceIsConsumedWhileEarlierKeyEventIsPending() { +@Test func backspaceIsConsumedWhileComposing() { #expect( disposition( event: KeyEventCore( @@ -52,12 +50,12 @@ private func disposition( charactersIgnoringModifiers: nil, keyCode: 51 ), - hasPendingKeyEvents: true + state: .composing ) == .sendToServer ) } -@Test func commandShortcutAlwaysFallsThroughWhileServerIsDelayed() { +@Test func commandShortcutAlwaysFallsThroughWhileComposing() { #expect( disposition( event: KeyEventCore( @@ -66,12 +64,37 @@ private func disposition( charactersIgnoringModifiers: "c", keyCode: 8 ), - state: .composing, - hasPendingKeyEvents: true + state: .composing ) == .fallthroughToApplication ) } +@Test func directEnglishInputDoesNotNeedServer() { + for event in [ + KeyEventCore(modifierFlags: [], characters: "a", charactersIgnoringModifiers: "a", keyCode: 0), + KeyEventCore(modifierFlags: [], characters: "\r", charactersIgnoringModifiers: "\r", keyCode: 36), + KeyEventCore(modifierFlags: [], characters: "\u{7f}", charactersIgnoringModifiers: "\u{7f}", keyCode: 51), + KeyEventCore(modifierFlags: [], characters: "\t", charactersIgnoringModifiers: "\t", keyCode: 48), + KeyEventCore(modifierFlags: [], characters: " ", charactersIgnoringModifiers: " ", keyCode: 49) + ] { + #expect(disposition(event: event, language: .english) == .fallthroughToApplication) + #expect(disposition(event: event, state: .composing, language: .english) == .sendToServer) + } +} + +@Test func directEnglishInputPreservesBackslashSetting() { + let event = KeyEventCore(modifierFlags: [], characters: "¥", charactersIgnoringModifiers: "¥", keyCode: 93) + #expect(ConverterClientEventRouter.disposition( + event: event, + context: .init(acknowledgedInputLanguage: .english, typeBackSlash: true) + ) == .insertText("\\")) +} + +@Test func englishDeadKeyStillUsesServerState() { + let event = KeyEventCore(modifierFlags: [], characters: "a", charactersIgnoringModifiers: "a", keyCode: 0) + #expect(disposition(event: event, state: .attachDiacritic("´"), language: .english) == .sendToServer) +} + @Test func unknownControlShortcutIsConsumedOnlyDuringComposition() { let event = KeyEventCore( modifierFlags: [.control], diff --git a/Core/Tests/CoreTests/XPCTests/DeadlineReplyTests.swift b/Core/Tests/CoreTests/XPCTests/DeadlineReplyTests.swift new file mode 100644 index 00000000..6a659ef2 --- /dev/null +++ b/Core/Tests/CoreTests/XPCTests/DeadlineReplyTests.swift @@ -0,0 +1,29 @@ +import Core +import Foundation +import Testing + +@MainActor +@Test func deadlineReplyCanArriveWhileMainActorIsBlocked() { + let reply = DeadlineReply(timeout: 2) + DispatchQueue.global().async { reply.complete("response") } + #expect(reply.wait() == "response") +} + +@Test func deadlineReplyTimesOutAndRejectsLateReply() { + let reply = DeadlineReply(timeout: 0.01) + #expect(reply.wait() == nil) + #expect(!reply.complete("late response")) +} + +@Test func deadlineReplyDoesNotAcceptResponseAfterDeadlineBeforeWait() { + let reply = DeadlineReply(timeout: 0) + #expect(!reply.complete("too late")) + #expect(reply.wait() == nil) +} + +@Test func deadlineReplyOnlyAcceptsFirstCompletion() { + let reply = DeadlineReply(timeout: 2) + #expect(reply.complete("first")) + #expect(!reply.complete("second")) + #expect(reply.wait() == "first") +} diff --git a/azooKeyMac/InputController/CandidateWindow/CandidateView.swift b/azooKeyMac/InputController/CandidateWindow/CandidateView.swift index 2685f123..4029e16a 100644 --- a/azooKeyMac/InputController/CandidateWindow/CandidateView.swift +++ b/azooKeyMac/InputController/CandidateWindow/CandidateView.swift @@ -1,7 +1,7 @@ import Cocoa import Core -protocol CandidatesViewControllerDelegate: AnyObject { +@MainActor protocol CandidatesViewControllerDelegate: AnyObject { func candidateSubmitted() func candidateSelectionChanged(_ row: Int) } diff --git a/azooKeyMac/InputController/ConverterServerClient.swift b/azooKeyMac/InputController/ConverterServerClient.swift index 627f9a53..94697605 100644 --- a/azooKeyMac/InputController/ConverterServerClient.swift +++ b/azooKeyMac/InputController/ConverterServerClient.swift @@ -1,11 +1,7 @@ import Core import Foundation -private enum ConverterServerXPC { - static let machServiceName = "dev.ensan.inputmethod.azooKeyMac.ConverterServer" -} - -@objc private protocol ConverterServerXPCProtocol { +@objc protocol ConverterServerXPCProtocol { func openSession(with reply: @escaping @Sendable (String) -> Void) func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) @@ -14,118 +10,87 @@ private enum ConverterServerXPC { @MainActor final class ConverterServerClient { - private static let commandTimeout: TimeInterval = 1 - - private var connection: NSXPCConnection? - private var sessionID: String? - private var hasOpenedSession = false - private var shouldAttemptReconnect = false - private var nextReconnectAttemptDate = Date.distantPast - private let commandQueue = OrderedAsyncCommandQueue() - - nonisolated init() {} - - var onLog: ((String) -> Void)? - var hasOpenSession: Bool { - sessionID != nil + private enum Command { + case session((String) -> ConverterSessionCommand) + case global(ConverterServerCommand) } - var canSendOrReconnect: Bool { - sessionID != nil || !hasOpenedSession || (shouldAttemptReconnect && Date() >= nextReconnectAttemptDate) + + private enum Reply: Sendable { + case response(Data) + case failure(String) } - var pendingCommandCount: Int { - commandQueue.count + + private struct PendingCommand { + let id = UUID() + let command: Command + let timeout: TimeInterval + let completion: (ConverterServerResponse?) -> Void } - func closeSession() { - guard let sessionID else { - invalidateConnection() - return - } - remoteObjectProxy { [weak self] proxy in - proxy?.closeSession(sessionID) { _ in - Task { @MainActor in - self?.invalidateConnection() - } - } - } + private struct ActiveCommand { + let id: UUID + let reply: DeadlineReply + let openingSessionID: String? } - func ping(_ message: String, completion: @escaping (String?) -> Void) { - remoteObjectProxy { proxy in - proxy?.ping(message) { response in - completion(response) - } - if proxy == nil { - completion(nil) - } + // 通常の変換が遅いだけでキーを途中放棄しない長さを取る。 + private let keyEventTimeout: TimeInterval + private let commandTimeout: TimeInterval + private let connectionFactory: @Sendable () -> NSXPCConnection + private var connection: NSXPCConnection? + private var sessionID: String? + private var abandonedSessionIDs: [String] = [] + private var pendingCommands: [PendingCommand] = [] + private var activeCommand: ActiveCommand? + + var onLog: ((String) -> Void)? + var onSessionReset: (() -> Void)? + + nonisolated init( + keyEventTimeout: TimeInterval = 30, + commandTimeout: TimeInterval = 1, + connectionFactory: @escaping @Sendable () -> NSXPCConnection = { + NSXPCConnection(machServiceName: "dev.ensan.inputmethod.azooKeyMac.ConverterServer", options: []) } + ) { + self.keyEventTimeout = keyEventTimeout + self.commandTimeout = commandTimeout + self.connectionFactory = connectionFactory } func listSettings( capabilities: ConverterSettingClientCapabilities, completion: @escaping ([ConverterSettingDescriptor]?) -> Void ) { - send( - { _ in - .settings(.list(capabilities: capabilities)) - }, - completion: { response in - completion(response?.settings) - } - ) + send({ _ in .settings(.list(capabilities: capabilities)) }, completion: { completion($0?.settings) }) } - func updateSetting( - key: String, - value: ConverterSettingValue, - completion: @escaping (Bool) -> Void - ) { - send( - { _ in - .settings(.update(key: key, value: value)) - }, - completion: { response in - completion(response != nil) - } - ) + func updateSetting(key: String, value: ConverterSettingValue, completion: @escaping (Bool) -> Void) { + send({ _ in .settings(.update(key: key, value: value)) }, completion: { completion($0 != nil) }) } func restartServer(completion: @escaping (Bool) -> Void) { - enqueueGlobal(.shutdown) { [weak self] response in - self?.invalidateConnection() + enqueue(.global(.shutdown), timeout: commandTimeout) { [weak self] response in + self?.resetSession() completion(response != nil) } } - func synchronizeUserDictionary( - forceExport: Bool, - completion: @escaping (Bool) -> Void - ) { - enqueueGlobal(.maintenance(.synchronizeUserDictionary(forceExport: forceExport))) { response in - completion(response != nil) + func synchronizeUserDictionary(forceExport: Bool, completion: @escaping (Bool) -> Void) { + enqueue(.global(.maintenance(.synchronizeUserDictionary(forceExport: forceExport))), timeout: commandTimeout) { + completion($0 != nil) } } func resetLearningData(completion: @escaping (Bool) -> Void) { - enqueueGlobal(.maintenance(.resetLearningData)) { response in - completion(response != nil) - } + enqueue(.global(.maintenance(.resetLearningData)), timeout: commandTimeout) { completion($0 != nil) } } func send( _ commandBuilder: @escaping (String) -> ConverterSessionCommand, completion: @escaping (ConverterServerResponse?) -> Void ) { - enqueue(commandBuilder, retriesOnFailure: false, completion: completion) - } - - /// キーイベントはタイムアウトで捨てず、1件ずつ順番に Server へ送る。 - /// XPC が一時的に切断した場合も先頭イベントを保持して再接続後に再送する。 - func sendKeyEvent( - _ request: ConverterKeyEventRequest, - completion: @escaping (ConverterServerResponse?) -> Void - ) { - enqueue({ _ in .handleKeyEvent(request) }, retriesOnFailure: true, completion: completion) + enqueue(.session(commandBuilder), timeout: commandTimeout, completion: completion) } func sendIfSessionOpen( @@ -136,196 +101,164 @@ final class ConverterServerClient { completion(nil) return } - enqueue(commandBuilder, retriesOnFailure: false, completion: completion) + send(commandBuilder, completion: completion) } - private func enqueue( + /// 応答を受け取る XPC キューはブロックしない。呼び出し元だけが期限付きで待つ。 + func sendKeyEvent(_ request: ConverterKeyEventRequest) -> ConverterServerResponse? { + sendSynchronously { _ in .handleKeyEvent(request) } + } + + func sendSynchronously( _ commandBuilder: @escaping (String) -> ConverterSessionCommand, - retriesOnFailure: Bool, - completion: @escaping (ConverterServerResponse?) -> Void - ) { - var proposedSessionID: String? - commandQueue.enqueue( - timeout: Self.commandTimeout, - timeoutOutcome: retriesOnFailure ? .retry : .finish(nil), - onTimeout: { [weak self] in - self?.handleCommandTimeout() - }, - operation: { [weak self] finish in - guard let self else { - finish(.finish(nil)) - return - } - let reconnectDelay = self.nextReconnectAttemptDate.timeIntervalSinceNow - if self.shouldAttemptReconnect, reconnectDelay > 0 { - DispatchQueue.main.asyncAfter(deadline: .now() + reconnectDelay) { - finish(.retry) - } - return - } - let sessionID = self.sessionID ?? proposedSessionID ?? UUID().uuidString - proposedSessionID = sessionID - let sessionCommand = commandBuilder(sessionID) - let command: ConverterServerCommand = if self.sessionID == nil { - .openSession(sessionID: sessionID, command: sessionCommand) - } else { - .session(sessionID: sessionID, command: sessionCommand) - } - self.sendResolved(command) { [weak self] response in - guard let self else { - finish(.finish(nil)) - return - } - if response != nil, self.sessionID == nil { - self.acceptOpenedSession(sessionID) - } - if response == nil, retriesOnFailure { - self.recordReconnectFailure() - finish(.retry) - } else { - finish(.finish(response)) - } - } - }, - completion: { response in - completion(response) - } - ) + onlyIfSessionOpen: Bool = false + ) -> ConverterServerResponse? { + flushPendingCommands() + if onlyIfSessionOpen && sessionID == nil { + return nil + } + var response: ConverterServerResponse? + enqueue(.session(commandBuilder), timeout: keyEventTimeout) { response = $0 } + flushPendingCommands() + return response + } + + /// 先行するモード変更・候補選択等を反映してから次のキーを判定する。 + /// completion もここで実行し、古い応答が後から UI を巻き戻すことを防ぐ。 + func flushPendingCommands() { + while let activeCommand { + finishCommand(id: activeCommand.id) + } } - private func enqueueGlobal( - _ command: ConverterServerCommand, + private func enqueue( + _ command: Command, + timeout: TimeInterval, completion: @escaping (ConverterServerResponse?) -> Void ) { - commandQueue.enqueue( - timeout: Self.commandTimeout, - timeoutOutcome: .finish(nil), - onTimeout: { [weak self] in - self?.handleCommandTimeout() - }, - operation: { [weak self] finish in - guard let self else { - finish(.finish(nil)) - return - } - self.sendResolved(command) { response in - finish(.finish(response)) - } - }, - completion: { response in - completion(response) - } - ) + pendingCommands.append(.init(command: command, timeout: timeout, completion: completion)) + startNextCommand() } - private func remoteObjectProxy(completion: @escaping (ConverterServerXPCProtocol?) -> Void) { - let connection = ensureConnection() - guard let proxy = connection.remoteObjectProxyWithErrorHandler({ [weak self] error in - DispatchQueue.main.async { - self?.onLog?("ConverterServer XPC error: \(error.localizedDescription)") - self?.resetConnection(preservingSession: true) - completion(nil) - } - }) as? ConverterServerXPCProtocol else { - completion(nil) + private func startNextCommand() { + guard activeCommand == nil, let pending = pendingCommands.first else { return } - completion(proxy) - } + let command: ConverterServerCommand + var openingSessionID: String? + switch pending.command { + case .session(let builder): + if let sessionID { + command = .session(sessionID: sessionID, command: builder(sessionID)) + } else { + let newID = UUID().uuidString + openingSessionID = newID + command = .openSession(sessionID: newID, command: builder(newID)) + } + case .global(let global): + command = global + } - private func sendResolved( - _ command: ConverterServerCommand, - completion: @escaping (ConverterServerResponse?) -> Void - ) { + let reply = DeadlineReply(timeout: pending.timeout) + let id = pending.id + activeCommand = .init(id: id, reply: reply, openingSessionID: openingSessionID) + // reply の保存と signal は XPC の返信キューで行う。MainActor への移動はその後。 + let complete: @Sendable (Reply) -> Void = { [weak self] result in + reply.complete(result) + DispatchQueue.main.async { self?.finishCommand(id: id) } + } do { let data = try ConverterServerCodec.encode(command) - self.remoteObjectProxy { proxy in - guard let proxy else { - completion(nil) - return - } - proxy.handleCommand(data) { [weak self] responseData, errorMessage in - let errorDescription = errorMessage.map(String.init) - DispatchQueue.main.async { - if let errorDescription { - self?.onLog?("ConverterServer command failed: \(errorDescription)") - if errorDescription.hasPrefix("Unknown converter session:") { - self?.resetConnection(preservingSession: false) - } - completion(nil) - return - } - guard let responseData else { - completion(nil) - return - } - completion(try? ConverterServerCodec.decodeResponse(from: responseData)) + let connection = ensureConnection() + if let proxy = connection.remoteObjectProxyWithErrorHandler({ error in + complete(.failure(error.localizedDescription)) + }) as? ConverterServerXPCProtocol { + proxy.handleCommand(data) { data, error in + if let error { + complete(.failure(String(error))) + } else if let data { + complete(.response(data)) + } else { + complete(.failure("Empty ConverterServer response")) } } + } else { + complete(.failure("Failed to create ConverterServer proxy")) } } catch { - self.onLog?("ConverterServer encode failed: \(error.localizedDescription)") - completion(nil) + complete(.failure(error.localizedDescription)) + } + DispatchQueue.main.asyncAfter(deadline: reply.deadline) { [weak self] in + self?.finishCommand(id: id) } } - private func ensureConnection() -> NSXPCConnection { - if let connection { - return connection + private func finishCommand(id: UUID) { + guard let active = activeCommand, active.id == id else { + return } - let connection = NSXPCConnection(machServiceName: ConverterServerXPC.machServiceName, options: []) - connection.remoteObjectInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) - connection.interruptionHandler = { [weak self] in - DispatchQueue.main.async { - self?.onLog?("ConverterServer connection interrupted") - self?.resetConnection(preservingSession: true) + let result = active.reply.wait() + let pending = pendingCommands.removeFirst() + activeCommand = nil + var response: ConverterServerResponse? + switch result { + case .response(let data): + do { + response = try ConverterServerCodec.decodeResponse(from: data) + } catch { + onLog?("ConverterServer decode failed: \(error.localizedDescription)") } + case .failure(let message): + onLog?("ConverterServer command failed: \(message)") + case nil: + onLog?("ConverterServer command timed out") } - connection.invalidationHandler = { [weak self] in - DispatchQueue.main.async { - self?.onLog?("ConverterServer connection invalidated") - self?.resetConnection(preservingSession: true) + if response != nil { + if let openingSessionID = active.openingSessionID { + sessionID = openingSessionID + } + pending.completion(response) + startNextCommand() + } else { + // Server が処理済みかは不明。再送せず、新しい session へ切り替える。 + let abandoned = pendingCommands + pendingCommands.removeAll() + resetSession(openingSessionID: active.openingSessionID) + pending.completion(nil) + for command in abandoned { + command.completion(nil) } } - connection.resume() - self.connection = connection - return connection } - private func resetConnection(preservingSession: Bool) { - let connection = self.connection - self.connection = nil - connection?.interruptionHandler = nil - connection?.invalidationHandler = nil - connection?.invalidate() - if sessionID != nil || hasOpenedSession { - shouldAttemptReconnect = true + private func ensureConnection() -> NSXPCConnection { + if let connection { + return connection } - if !preservingSession { - self.sessionID = nil + let connection = connectionFactory() + connection.remoteObjectInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) + connection.resume() + self.connection = connection + // 切断前の計算は継続している場合がある。新しい接続で旧 session を回収する。 + if !abandonedSessionIDs.isEmpty, + let proxy = connection.remoteObjectProxyWithErrorHandler({ _ in }) as? ConverterServerXPCProtocol { + for id in abandonedSessionIDs { + proxy.closeSession(id) { _ in } + } + abandonedSessionIDs.removeAll() } + return connection } - private func invalidateConnection() { - resetConnection(preservingSession: false) - } - - private func recordReconnectFailure() { - shouldAttemptReconnect = true - nextReconnectAttemptDate = Date().addingTimeInterval(0.2) - } - - private func acceptOpenedSession(_ sessionID: String) { - self.sessionID = sessionID - hasOpenedSession = true - shouldAttemptReconnect = false - nextReconnectAttemptDate = .distantPast - onLog?("ConverterServer session opened: \(sessionID)") - } - - private func handleCommandTimeout() { - onLog?("ConverterServer command timed out") - recordReconnectFailure() - resetConnection(preservingSession: true) + private func resetSession(openingSessionID: String? = nil) { + // 切断は Server の計算をキャンセルしない。同じ session を使うと、 + // 時間切れになったキーが次回の候補へ混入するので再利用しない。 + if let id = sessionID ?? openingSessionID { + abandonedSessionIDs.append(id) + } + sessionID = nil + connection?.invalidate() + connection = nil + onSessionReset?() } } diff --git a/azooKeyMac/InputController/azooKeyMacInputController.swift b/azooKeyMac/InputController/azooKeyMacInputController.swift index f10ba353..09756b99 100644 --- a/azooKeyMac/InputController/azooKeyMacInputController.swift +++ b/azooKeyMac/InputController/azooKeyMacInputController.swift @@ -8,7 +8,6 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s private var currentConverterView: ConverterSessionSnapshot? private(set) var inputState: InputState = .none private var inputLanguage: InputLanguage = .japanese - private var pendingKeyEventCount = 0 private var nextKeyEventID: UInt64 = 0 private var activationGeneration: UInt64 = 0 private var pendingConverterServerActivation: ConverterSessionActivation? @@ -139,6 +138,9 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s self.converterServerClient.onLog = { [weak self] message in self?.appendDebugMessage(message) } + self.converterServerClient.onSessionReset = { [weak self] in + self?.recoverFromConverterServerFailure() + } } self.setupMenu() } @@ -188,17 +190,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s @MainActor override func commitComposition(_ sender: Any!) { - let activationGeneration = self.activationGeneration - self.converterServerClient.sendIfSessionOpen({ _ in .composition(.commit) }, completion: { [weak self] response in - Task { @MainActor in - guard let self, - self.activationGeneration == activationGeneration, - let response else { - return - } - self.apply(response) - } - }) + self.sendAndApply { _ in .composition(.commit) } } // MARK: - setValue: 状態同期のみ @@ -250,6 +242,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s guard event.type == .keyDown else { return false } + self.converterServerClient.flushPendingCommands() // カスタムプロンプトショートカットのチェック if let matchedPrompt = checkCustomPromptShortcut(event: event) { @@ -335,20 +328,26 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s enableSuggestion: Bool, optionDirectInputText: String? = nil ) -> Bool { + self.converterServerClient.flushPendingCommands() let disposition = ConverterClientEventRouter.disposition( event: event, context: .init( acknowledgedInputState: ConverterInputState(self.inputState), acknowledgedInputLanguage: self.inputLanguage, - hasPendingKeyEvents: self.pendingKeyEventCount > 0, liveConversionEnabled: Config.LiveConversion().value, enableDebugWindow: Config.DebugWindow().value, enableSuggestion: enableSuggestion, typeBackSlash: Config.TypeBackSlash().value ) ) - guard disposition == .sendToServer else { + switch disposition { + case .fallthroughToApplication: return false + case .insertText(let text): + self.client()?.insertText(text, replacementRange: NSRange(location: NSNotFound, length: 0)) + return true + case .sendToServer: + break } self.nextKeyEventID &+= 1 @@ -368,30 +367,31 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s activation: self.pendingConverterServerActivation ) self.pendingConverterServerActivation = nil - self.pendingKeyEventCount += 1 - let activationGeneration = self.activationGeneration - self.converterServerClient.sendKeyEvent(request) { [weak self] response in - Task { @MainActor in - guard let self else { - return - } - self.pendingKeyEventCount = max(0, self.pendingKeyEventCount - 1) - guard self.activationGeneration == activationGeneration else { - return - } - guard let response else { - self.appendDebugMessage("ConverterServer dropped key event \(request.eventID)") - return - } - if !response.handled { - // `handle` は既に同期的に consume 済み。未応答イベントがある場合は - // application へ後からイベントを戻せないため、ここでは漏らさない。 - self.appendDebugMessage("Consumed delayed fallthrough event \(request.eventID)") - } - self.apply(response) - } + guard let response = self.converterServerClient.sendKeyEvent(request) else { + return false } - return true + self.apply(response) + return response.handled + } + + @MainActor + private func recoverFromConverterServerFailure() { + // 最後に表示した文字列を保全し、Server にだけ残った未確認の操作は引き継がない。 + let text = self.currentMarkedText().elements.map(\.content).joined() + self.currentConverterView = nil + self.inputState = .none + self.activationGeneration &+= 1 + self.pendingConverterServerActivation = ConverterSessionActivation( + config: self.converterServerSessionConfig, + inputLanguage: self.inputLanguage + ) + if !text.isEmpty { + self.client()?.insertText(text, replacementRange: NSRange(location: NSNotFound, length: 0)) + self.refreshMarkedText() + } + self.refreshCandidateWindow() + self.hidePredictionWindow() + self.refreshReplaceSuggestionWindow() } @MainActor @@ -407,6 +407,13 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s ) } + @MainActor + private func sendAndApply(_ command: @escaping (String) -> ConverterSessionCommand) { + if let response = self.converterServerClient.sendSynchronously(command, onlyIfSessionOpen: true) { + self.apply(response) + } + } + @MainActor private func apply(_ response: ConverterServerResponse) { if let inputLanguage = response.inputLanguage { @@ -574,17 +581,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s let count = view.replaceSuggestionCandidates.count let current = view.replaceSuggestionSelectionIndex ?? (offset > 0 ? -1 : 0) let next = (current + offset + count) % count - self.converterServerClient.sendIfSessionOpen( - { _ in .replaceSuggestion(.selectReplaceSuggestionCandidate(index: next)) }, - completion: { [weak self] response in - Task { @MainActor in - guard let self, let response else { - return - } - self.apply(response) - } - } - ) + self.sendAndApply { _ in .replaceSuggestion(.selectReplaceSuggestionCandidate(index: next)) } } @MainActor private func showReplaceSuggestionError(message: String) { @@ -757,42 +754,20 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s } extension azooKeyMacInputController: CandidatesViewControllerDelegate { + @MainActor func candidateSubmitted() { - Task { @MainActor in - guard self.currentConverterView != nil else { - return - } - self.converterServerClient.sendIfSessionOpen( - { _ in .candidate(.submitSelectedCandidate(context: self.currentConverterTextContext())) }, - completion: { [weak self] response in - Task { @MainActor in - guard let self, let response else { - return - } - self.apply(response) - } - } - ) + guard self.currentConverterView != nil else { + return } + self.sendAndApply { _ in .candidate(.submitSelectedCandidate(context: self.currentConverterTextContext())) } } + @MainActor func candidateSelectionChanged(_ row: Int) { - Task { @MainActor in - guard self.currentConverterView != nil else { - return - } - self.converterServerClient.sendIfSessionOpen( - { _ in .candidate(.selectCandidate(index: row)) }, - completion: { [weak self] response in - Task { @MainActor in - guard let self, let response else { - return - } - self.apply(response) - } - } - ) + guard self.currentConverterView != nil else { + return } + self.sendAndApply { _ in .candidate(.selectCandidate(index: row)) } } } @@ -846,23 +821,11 @@ extension azooKeyMacInputController: ReplaceSuggestionsViewControllerDelegate { guard self.currentConverterView?.replaceSuggestionSelectionIndex != row else { return } - self.converterServerClient.sendIfSessionOpen( - { _ in .replaceSuggestion(.selectReplaceSuggestionCandidate(index: row)) }, - completion: { [weak self] response in - Task { @MainActor in - guard let self, let response else { - return - } - self.apply(response) - } - } - ) + self.sendAndApply { _ in .replaceSuggestion(.selectReplaceSuggestionCandidate(index: row)) } } func replaceSuggestionSubmitted() { - Task { @MainActor in - self.submitSelectedSuggestionCandidate() - } + self.submitSelectedSuggestionCandidate() } } @@ -882,26 +845,25 @@ extension azooKeyMacInputController { return } self.syncConverterServerSessionConfig() + let activationGeneration = self.activationGeneration self.converterServerClient.sendIfSessionOpen( { _ in .replaceSuggestion(.request(context: self.currentConverterTextContext())) }, completion: { [weak self] response in - Task { @MainActor in - guard let self else { - return - } - guard let response else { - self.showReplaceSuggestionError(message: "ConverterServerから候補を取得できませんでした") - return - } - guard self.currentConverterView?.convertTarget == response.snapshot.convertTarget else { - self.appendDebugMessage("候補ウィンドウ更新をスキップ: composition changed") - return - } - self.currentConverterView = response.snapshot - self.inputState = response.inputState.inputState - self.refreshMarkedText() - self.refreshReplaceSuggestionWindow() + guard let self, self.activationGeneration == activationGeneration else { + return + } + guard let response else { + self.showReplaceSuggestionError(message: "ConverterServerから候補を取得できませんでした") + return } + guard self.currentConverterView?.convertTarget == response.snapshot.convertTarget else { + self.appendDebugMessage("候補ウィンドウ更新をスキップ: composition changed") + return + } + self.currentConverterView = response.snapshot + self.inputState = response.inputState.inputState + self.refreshMarkedText() + self.refreshReplaceSuggestionWindow() } ) self.appendDebugMessage("requestReplaceSuggestion: 終了") @@ -914,17 +876,7 @@ extension azooKeyMacInputController { } @MainActor func submitSelectedSuggestionCandidate() { - self.converterServerClient.sendIfSessionOpen( - { _ in .replaceSuggestion(.submitSelectedReplaceSuggestion) }, - completion: { [weak self] response in - Task { @MainActor in - guard let self, let response else { - return - } - self.apply(response) - } - } - ) + self.sendAndApply { _ in .replaceSuggestion(.submitSelectedReplaceSuggestion) } } @MainActor private func finishReplaceSuggestionComposition() { diff --git a/azooKeyMacTests/ThinClientInputPipelineTests.swift b/azooKeyMacTests/ThinClientInputPipelineTests.swift index 3ed6ef38..f3f88e28 100644 --- a/azooKeyMacTests/ThinClientInputPipelineTests.swift +++ b/azooKeyMacTests/ThinClientInputPipelineTests.swift @@ -1,90 +1,196 @@ import Core +import Foundation import XCTest +@testable import azooKeyMac + +/// 実際の NSXPCConnection を使い、返信が MainActor を経由せずに同期待ちを解除することを検証する。 +private final class TestConverterService: NSObject, NSXPCListenerDelegate, ConverterServerXPCProtocol, @unchecked Sendable { + struct Response: Sendable { + var data: Data? + var delay: TimeInterval = 0 + var error: String? + } + + let listener = NSXPCListener.anonymous() + private let lock = NSLock() + private var connections: [NSXPCConnection] = [] + private var receivedCommands: [ConverterServerCommand] = [] + private let responses: [Response] + + init(responses: [Response]) { + self.responses = responses + super.init() + listener.delegate = self + listener.resume() + } + + func listener(_ listener: NSXPCListener, shouldAcceptNewConnection connection: NSXPCConnection) -> Bool { + connection.exportedInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) + connection.exportedObject = self + lock.lock() + connections.append(connection) + lock.unlock() + connection.resume() + return true + } + + var commands: [ConverterServerCommand] { + lock.lock() + defer { lock.unlock() } + return receivedCommands + } + + func stop() { + lock.lock() + let connections = self.connections + self.connections.removeAll() + lock.unlock() + connections.forEach { $0.invalidate() } + listener.invalidate() + } + + func openSession(with reply: @escaping @Sendable (String) -> Void) { reply(UUID().uuidString) } + func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) { reply(true) } + func ping(_ message: String, with reply: @escaping @Sendable (String) -> Void) { reply(message) } + + func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) { + guard let command = try? ConverterServerCodec.decodeCommand(from: data) else { + reply(nil, "invalid request") + return + } + lock.lock() + let index = receivedCommands.count + receivedCommands.append(command) + let response = responses.indices.contains(index) ? responses[index] : Response(error: "unexpected request") + lock.unlock() + DispatchQueue.global().asyncAfter(deadline: .now() + response.delay) { + reply(response.data, response.error.map { $0 as NSString }) + } + } +} + @MainActor final class ThinClientInputPipelineTests: XCTestCase { - func testDelayedServerReplyKeepsFollowingInputOwnedAndOrdered() async { - let printable = KeyEventCore( - modifierFlags: [], - characters: "a", - charactersIgnoringModifiers: "a", - keyCode: 0 - ) - XCTAssertEqual( - ConverterClientEventRouter.disposition( - event: printable, - context: .init(typeBackSlash: true) - ), - .sendToServer - ) + private func makeClient(_ service: TestConverterService, timeout: TimeInterval = 2) -> ConverterServerClient { + ConverterServerClient(keyEventTimeout: timeout, commandTimeout: timeout) { + NSXPCConnection(listenerEndpoint: service.listener.endpoint) + } + } - let backspace = KeyEventCore( - modifierFlags: [], - characters: "\u{7f}", - charactersIgnoringModifiers: "\u{7f}", - keyCode: 51 - ) - XCTAssertEqual( - ConverterClientEventRouter.disposition( - event: backspace, - context: .init(hasPendingKeyEvents: true, typeBackSlash: true) - ), - .sendToServer, - "未応答中の状態mirrorを信じてbackspaceをapplicationへ漏らしてはいけない" - ) + private func encodedResponse(handled: Bool = true, state: ConverterInputState = .none) throws -> Data { + try ConverterServerCodec.encode(ConverterServerResponse(handled: handled, inputState: state, snapshot: .empty)) + } - let queue = OrderedAsyncCommandQueue() - var firstFinish: OrderedAsyncCommandQueue.Finish? - var starts: [Int] = [] - var completions: [Int] = [] - let completed = expectation(description: "both key events completed") - completed.expectedFulfillmentCount = 2 - - queue.enqueue( - operation: { finish in - starts.append(1) - firstFinish = finish - }, - completion: { value in - completions.append(value) - completed.fulfill() - } - ) - queue.enqueue( - operation: { finish in - starts.append(2) - finish(.finish(2)) - }, - completion: { value in - completions.append(value) - completed.fulfill() - } + private func keyRequest(id: UInt64 = 1) -> ConverterKeyEventRequest { + .init( + eventID: id, + event: .init(modifierFlags: [], characters: "a", charactersIgnoringModifiers: "a", keyCode: 0), + inputStyle: .defaultRomanToKana, + liveConversionEnabled: true, + enableDebugWindow: false, + enableSuggestion: false ) + } + + func testSynchronousCallReceivesReplyWithoutRunningMainQueue() throws { + let service = TestConverterService(responses: [ + .init(data: try encodedResponse(state: .composing), delay: 0.02) + ]) + defer { service.stop() } + let client = makeClient(service) + + let response = client.sendKeyEvent(keyRequest()) - XCTAssertEqual(starts, [1], "2件目は1件目の遅延応答より先にServerへ送ってはいけない") - firstFinish?(.finish(1)) - await fulfillment(of: [completed], timeout: 1) - XCTAssertEqual(starts, [1, 2]) - XCTAssertEqual(completions, [1, 2]) + XCTAssertEqual(response?.inputState, .composing) + XCTAssertEqual(service.commands.count, 1) + guard case .openSession(_, .handleKeyEvent(let request)) = service.commands[0] else { + return XCTFail("The first key must open the session atomically") + } + XCTAssertEqual(request.eventID, 1) } - func testCommandShortcutFallsThroughEvenWhileServerReplyIsPending() { - let commandA = KeyEventCore( - modifierFlags: [.command], - characters: "a", - charactersIgnoringModifiers: "a", - keyCode: 0 - ) - XCTAssertEqual( - ConverterClientEventRouter.disposition( - event: commandA, - context: .init( - acknowledgedInputState: .composing, - hasPendingKeyEvents: true, - typeBackSlash: true - ) - ), - .fallthroughToApplication - ) + func testServerFallthroughIsReturnedSynchronously() throws { + let service = TestConverterService(responses: [.init(data: try encodedResponse(handled: false))]) + defer { service.stop() } + let client = makeClient(service) + + let response = client.sendKeyEvent(keyRequest()) + + XCTAssertEqual(response?.handled, false) + } + + func testEarlierAsyncCommandCompletionIsAppliedBeforeSynchronousRequest() async throws { + let service = TestConverterService(responses: [ + .init(data: try encodedResponse(), delay: 0.02), + .init(data: try encodedResponse(state: .composing)) + ]) + defer { service.stop() } + let client = makeClient(service) + var completions: [String] = [] + client.send({ _ in .lifecycle(.synchronizeInputLanguage(.japanese)) }) { _ in + completions.append("language") + } + + let response = client.sendSynchronously { _ in + XCTAssertEqual(completions, ["language"]) + return .composition(.snapshot) + } + completions.append("key") + + XCTAssertEqual(response?.inputState, .composing) + XCTAssertEqual(service.commands.count, 2) + guard case .openSession(let firstID, _) = service.commands[0], + case .session(let secondID, _) = service.commands[1] else { + return XCTFail("Expected one session shared by both commands") + } + XCTAssertEqual(firstID, secondID) + await Task.yield() + XCTAssertEqual(completions, ["language", "key"], "Queued main-thread callbacks must not apply a reply twice") + } + + func testTimeoutDropsOldReplyAndOpensFreshSession() async throws { + let service = TestConverterService(responses: [ + .init(data: try encodedResponse(state: .composing), delay: 0.3), + .init(data: try encodedResponse()) + ]) + defer { service.stop() } + let client = makeClient(service, timeout: 0.1) + var resets = 0 + client.onSessionReset = { resets += 1 } + + XCTAssertNil(client.sendSynchronously { _ in .composition(.snapshot) }) + XCTAssertEqual(resets, 1) + let response = client.sendSynchronously { _ in .composition(.snapshot) } + XCTAssertEqual(response?.inputState, ConverterInputState.none) + XCTAssertEqual(service.commands.count, 2) + guard case .openSession(let firstID, _) = service.commands[0], + case .openSession(let secondID, _) = service.commands[1] else { + return XCTFail("A timed-out session must not be reused") + } + XCTAssertNotEqual(firstID, secondID) + try await Task.sleep(nanoseconds: 350_000_000) + XCTAssertEqual(resets, 1, "A late reply/error must not invalidate the new session") + XCTAssertEqual(service.commands.count, 2, "A timed-out operation must not be retried") + } + + func testFailureAbandonsQueuedCommandsInsteadOfReplayingThem() throws { + let service = TestConverterService(responses: [.init(error: "test failure")]) + defer { service.stop() } + let client = makeClient(service) + var completions = 0 + client.send({ _ in .composition(.snapshot) }) { response in + XCTAssertNil(response) + completions += 1 + } + client.send({ _ in .composition(.commit) }) { response in + XCTAssertNil(response) + completions += 1 + } + + client.flushPendingCommands() + + XCTAssertEqual(completions, 2) + XCTAssertEqual(service.commands.count, 1) } } From 7c4998f1bdd7439eb6f2e6af8164a5209f61d9f8 Mon Sep 17 00:00:00 2001 From: ensan-hcl Date: Sun, 6 Sep 2026 13:41:25 +0900 Subject: [PATCH 2/3] fix(input): synchronize current modes and release converter sessions --- Core/Package.swift | 7 + .../ConverterServer+KeyEvent.swift | 3 +- .../ConverterServer/ConverterServer.swift | 330 +++++++++++++++++ .../ConverterServerConnection.swift | 81 +++++ Core/Sources/ConverterServer/main.swift | 331 +----------------- .../XPC/ConverterClientSessionState.swift | 32 ++ .../ConverterServerLifecycleTests.swift | 170 +++++++++ .../ConverterClientSessionStateTests.swift | 59 ++++ .../ConverterServerClient.swift | 44 ++- .../azooKeyMacInputController.swift | 62 ++-- .../ThinClientInputPipelineTests.swift | 52 ++- 11 files changed, 812 insertions(+), 359 deletions(-) create mode 100644 Core/Sources/ConverterServer/ConverterServer.swift create mode 100644 Core/Sources/ConverterServer/ConverterServerConnection.swift create mode 100644 Core/Sources/Core/XPC/ConverterClientSessionState.swift create mode 100644 Core/Tests/ConverterServerTests/ConverterServerLifecycleTests.swift create mode 100644 Core/Tests/CoreTests/XPCTests/ConverterClientSessionStateTests.swift diff --git a/Core/Package.swift b/Core/Package.swift index 87d7c3f0..68a8870f 100644 --- a/Core/Package.swift +++ b/Core/Package.swift @@ -61,6 +61,13 @@ targets.append( swiftSettings: [.interoperabilityMode(.Cxx)] ) ) +targets.append( + .testTarget( + name: "ConverterServerTests", + dependencies: ["ConverterServer"], + swiftSettings: [.interoperabilityMode(.Cxx)] + ) +) #endif let package = Package( diff --git a/Core/Sources/ConverterServer/ConverterServer+KeyEvent.swift b/Core/Sources/ConverterServer/ConverterServer+KeyEvent.swift index b90ba16c..7a9cea55 100644 --- a/Core/Sources/ConverterServer/ConverterServer+KeyEvent.swift +++ b/Core/Sources/ConverterServer/ConverterServer+KeyEvent.swift @@ -29,8 +29,7 @@ extension ConverterServer { session.manager.activate() } session.setContext(request.context) - Config.DebugPredictiveTyping().value = request.enablePredictiveTyping - Config.DebugTypoCorrection().value = request.enableTypoCorrection + applyRequestSettings(request) if request.enableOptionDirectFullWidthInput, let text = OptionDirectInputResolver.resolve( diff --git a/Core/Sources/ConverterServer/ConverterServer.swift b/Core/Sources/ConverterServer/ConverterServer.swift new file mode 100644 index 00000000..d1acdf32 --- /dev/null +++ b/Core/Sources/ConverterServer/ConverterServer.swift @@ -0,0 +1,330 @@ +import Core +import Darwin +import Foundation +import KanaKanjiConverterModuleWithDefaultDictionary + +enum ConverterServerXPC { + static let machServiceName = "dev.ensan.inputmethod.azooKeyMac.ConverterServer" +} + +@objc protocol ConverterServerXPCProtocol { + func openSession(with reply: @escaping @Sendable (String) -> Void) + func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) + func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) + func ping(_ message: String, with reply: @escaping @Sendable (String) -> Void) +} + +final class ConverterServer: @unchecked Sendable { + private static let learningDataCommitDelay: TimeInterval = 2 + + private var sessions: [String: ConverterSession] = [:] + private var sessionOwners: [String: UUID] = [:] + private let kanaKanjiConverter = KanaKanjiConverter.withDefaultDictionary() + private let learningDataCommitScheduler = DebouncedActionScheduler() + private let makeManager: @MainActor (KanaKanjiConverter) -> SegmentsManager + let applyRequestSettings: @Sendable (ConverterKeyEventRequest) -> Void + + init( + makeManager: @escaping @MainActor (KanaKanjiConverter) -> SegmentsManager = { + ConverterServer.makeSegmentsManager(kanaKanjiConverter: $0) + }, + applyRequestSettings: @escaping @Sendable (ConverterKeyEventRequest) -> Void = { + Config.DebugPredictiveTyping().value = $0.enablePredictiveTyping + Config.DebugTypoCorrection().value = $0.enableTypoCorrection + } + ) { + self.makeManager = makeManager + self.applyRequestSettings = applyRequestSettings + } + + @MainActor + func execute(_ command: ConverterServerCommand, owner: UUID) async throws -> ConverterServerResponse { + switch command { + case .openSession(let id, _): + sessionOwners[id] = owner + case .session(let id, _): + _ = try getSession(id) + sessionOwners[id] = owner + case .shutdown, .maintenance: + break + } + defer { learningDataCommitScheduler.postponeIfScheduled(after: Self.learningDataCommitDelay) } + return try await handle(command) + } + + @MainActor + func removeSessions(ownedBy owner: UUID) { + for id in sessionOwners.filter({ $0.value == owner }).map(\.key) { + removeSession(id) + } + } + + @MainActor + @discardableResult + func removeSession(_ id: String) -> Bool { + sessionOwners.removeValue(forKey: id) + guard let session = sessions.removeValue(forKey: id) else { + return false + } + kanaKanjiConverter.removeSession(session.conversionSessionID) + scheduleLearningDataCommit() + return true + } + + @MainActor + private func handle(_ command: ConverterServerCommand) async throws -> ConverterServerResponse { + switch command { + case .shutdown: + Self.scheduleShutdown() + return ConverterServerResponse(snapshot: .empty) + case .maintenance(let command): + return try handle(command) + case .openSession(let sessionID, let command): + createSessionIfNeeded(sessionID) + return try await handle(command, sessionID: sessionID) + case .session(let sessionID, let command): + return try await handle(command, sessionID: sessionID) + } + } + + @MainActor + private func createSessionIfNeeded(_ sessionID: String) { + guard sessions[sessionID] == nil else { + return + } + let conversionSessionID = kanaKanjiConverter.createSession() + sessions[sessionID] = ConverterSession( + manager: makeManager(kanaKanjiConverter), + conversionSessionID: conversionSessionID + ) + } + + @MainActor + private func handle(_ command: ConverterMaintenanceCommand) throws -> ConverterServerResponse { + switch command { + case .synchronizeUserDictionary(let forceExport): + let memoryDirectoryURL = AppGroup.memoryDirectoryURL() + if forceExport || !CompiledUserDictionaryStore.hasExportedDictionary(memoryDirectoryURL: memoryDirectoryURL) { + try CompiledUserDictionaryStore.exportCurrentDictionaries(memoryDirectoryURL: memoryDirectoryURL) + } + kanaKanjiConverter.updateUserDictionaryURL( + CompiledUserDictionaryStore.directoryURL(memoryDirectoryURL: memoryDirectoryURL), + forceReload: true + ) + case .resetLearningData: + kanaKanjiConverter.resetMemory() + } + return ConverterServerResponse(snapshot: .empty) + } + + @MainActor + private func handle(_ command: ConverterSessionCommand, sessionID: String) async throws -> ConverterServerResponse { + let session = try getSession(sessionID) + switch command { + case .lifecycle(let command): + return try withConverterSession(session) { + handle(command, session: session) + } + case .settings(let command): + return try withConverterSession(session) { + try handle(command, session: session) + } + case .updateConfig(let config): + return try withConverterSession(session) { + session.config = config + return makeResponse(for: session, inputState: .none) + } + case .handleKeyEvent(let request): + return try withConverterSession(session) { + try handleKeyEvent(sessionID: sessionID, request: request) + } + case .composition(let command): + return try withConverterSession(session) { + handle(command, session: session) + } + case .candidate(let command): + return try withConverterSession(session) { + handle(command, session: session) + } + case .replaceSuggestion(let command): + return try await handle(command, session: session) + } + } + + @MainActor + private func withConverterSession( + _ session: ConverterSession, + operation: () throws -> Result + ) throws -> Result { + try kanaKanjiConverter.withSession(session.conversionSessionID, operation: operation) + } + + @MainActor + private func handle( + _ command: ConverterSessionLifecycleCommand, + session: ConverterSession + ) -> ConverterServerResponse { + switch command { + case .activate: + session.manager.activate() + return makeResponse(for: session, inputState: session.inputState) + case .deactivate: + // アプリ切替直後のキー入力を、学習データの同期I/Oで塞がない。 + // 共有Converterはプロセス内に残るため、永続化だけ入力のアイドル時まで遅延できる。 + session.manager.deactivate(flushLearningData: false) + scheduleLearningDataCommit() + session.inputState = .none + session.clearReplaceSuggestions() + return makeResponse(for: session, inputState: session.inputState) + case .synchronizeInputLanguage(let language): + session.inputLanguage = language + if language == .english { + session.manager.stopJapaneseInput() + } + return makeResponse( + for: session, + inputState: session.inputState + ) + } + } + + @MainActor + private func handle( + _ command: ConverterSettingsCommand, + session: ConverterSession + ) throws -> ConverterServerResponse { + switch command { + case .list(let capabilities): + return makeResponse( + for: session, + inputState: .none, + settings: Self.makeSettingDescriptors(capabilities: capabilities) + ) + case .update(let key, let value): + try Self.updateSetting(key: key, value: value) + return makeResponse(for: session, inputState: .none) + } + } + + @MainActor + private func handle( + _ command: ConverterCompositionCommand, + session: ConverterSession + ) -> ConverterServerResponse { + switch command { + case .snapshot: + return makeResponse(for: session, inputState: session.inputState) + case .stopComposition: + session.manager.stopComposition() + session.inputState = .none + return makeResponse(for: session, inputState: session.inputState) + case .forgetMemory: + session.manager.forgetMemory() + return makeResponse(for: session, inputState: session.inputState) + case .commit: + let text = session.manager.commitMarkedText(inputState: session.inputState) + let effects: [ConverterClientEffect] = text.isEmpty ? [] : [.insertText(text)] + session.inputState = .none + return makeResponse( + for: session, + inputState: session.inputState, + effects: effects, + responseInputState: ConverterInputState.none + ) + } + } + + @MainActor + private func handle( + _ command: ConverterCandidateCommand, + session: ConverterSession + ) -> ConverterServerResponse { + switch command { + case .selectCandidate(let index): + session.manager.requestSelectingRow(index) + session.inputState = .selecting + return makeResponse(for: session, inputState: session.inputState) + case .submitSelectedCandidate(let context): + session.setContext(context) + var effects: [ConverterClientEffect] = [] + submitSelectedCandidate( + manager: session.manager, + leftSideContext: session.conversionLeftSideContext(), + effects: &effects + ) + let nextInputState: InputState = session.manager.isEmpty ? .none : .previewing + session.inputState = nextInputState + return makeResponse( + for: session, + inputState: nextInputState, + effects: effects, + responseInputState: ConverterInputState(nextInputState) + ) + } + } + + @MainActor + private func handle( + _ command: ConverterReplaceSuggestionCommand, + session: ConverterSession + ) async throws -> ConverterServerResponse { + switch command { + case .request(let context): + session.setContext(context) + try await requestReplaceSuggestion(session: session) + return try withConverterSession(session) { + session.inputState = .replaceSuggestion + return makeResponse( + for: session, + inputState: session.inputState, + responseInputState: .replaceSuggestion + ) + } + case .selectReplaceSuggestionCandidate(let index): + return try withConverterSession(session) { + session.selectReplaceSuggestion(at: index) + session.inputState = .replaceSuggestion + return makeResponse( + for: session, + inputState: session.inputState, + responseInputState: .replaceSuggestion + ) + } + case .submitSelectedReplaceSuggestion: + return try withConverterSession(session) { + var effects: [ConverterClientEffect] = [] + let didSubmit = submitSelectedReplaceSuggestion(session: session, effects: &effects) + let nextInputState: InputState = didSubmit ? .none : .replaceSuggestion + session.inputState = nextInputState + return makeResponse( + for: session, + inputState: nextInputState, + effects: effects, + responseInputState: ConverterInputState(nextInputState) + ) + } + } + } + + private static func scheduleShutdown() { + DispatchQueue.main.asyncAfter(deadline: .now() + 0.2) { + exit(EXIT_SUCCESS) + } + } + + @MainActor + private func scheduleLearningDataCommit() { + learningDataCommitScheduler.schedule(after: Self.learningDataCommitDelay) { [weak self] in + self?.kanaKanjiConverter.commitUpdateLearningData() + } + } + + @MainActor + func getSession(_ sessionID: String) throws -> ConverterSession { + guard let session = sessions[sessionID] else { + throw ConverterServerError.unknownSession(sessionID) + } + return session + } + +} diff --git a/Core/Sources/ConverterServer/ConverterServerConnection.swift b/Core/Sources/ConverterServer/ConverterServerConnection.swift new file mode 100644 index 00000000..df55403a --- /dev/null +++ b/Core/Sources/ConverterServer/ConverterServerConnection.swift @@ -0,0 +1,81 @@ +import Core +import Foundation + +/// 接続の終了後に届く要求を実行せず、その接続が残した session を回収する。 +final class ConverterServerConnection: NSObject, ConverterServerXPCProtocol, @unchecked Sendable { + private let server: ConverterServer + private let owner = UUID() + private let cleanupDelay: TimeInterval + private let lock = NSLock() + private var closed = false + + private var isClosed: Bool { + lock.lock() + defer { lock.unlock() } + return closed + } + + private func markClosed() -> Bool { + lock.lock() + defer { lock.unlock() } + guard !closed else { + return false + } + closed = true + return true + } + + init(server: ConverterServer, cleanupDelay: TimeInterval = 5) { + self.server = server + self.cleanupDelay = cleanupDelay + } + + func invalidate() { + guard markClosed() else { + return + } + Task { @MainActor in + // 旧Clientの再接続・同一sessionへの再送を短時間だけ許容する。 + // 新接続が所有権を取得したsessionは、古い接続のcleanupでは削除しない。 + DispatchQueue.main.asyncAfter(deadline: .now() + self.cleanupDelay) { + self.server.removeSessions(ownedBy: self.owner) + } + } + } + + func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) { + Task(priority: .userInitiated) { @MainActor in + guard !self.isClosed else { + reply(nil, "ConverterServer connection is closed") + return + } + do { + let command = try ConverterServerCodec.decodeCommand(from: data) + let response = try await self.server.execute(command, owner: self.owner) + reply(try ConverterServerCodec.encode(response), nil) + } catch { + reply(nil, error.localizedDescription as NSString) + } + } + } + + func openSession(with reply: @escaping @Sendable (String) -> Void) { + Task { @MainActor in + guard !self.isClosed else { + reply("") + return + } + let id = UUID().uuidString + _ = try? await self.server.execute(.openSession(sessionID: id, command: .composition(.snapshot)), owner: self.owner) + reply(id) + } + } + + func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) { + Task { @MainActor in reply(self.server.removeSession(sessionID)) } + } + + func ping(_ message: String, with reply: @escaping @Sendable (String) -> Void) { + reply("ConverterServer: \(message)") + } +} diff --git a/Core/Sources/ConverterServer/main.swift b/Core/Sources/ConverterServer/main.swift index 5e9b4bc5..2caf0494 100644 --- a/Core/Sources/ConverterServer/main.swift +++ b/Core/Sources/ConverterServer/main.swift @@ -1,337 +1,14 @@ -import Core -import Darwin import Foundation -import KanaKanjiConverterModuleWithDefaultDictionary - -private enum ConverterServerXPC { - static let machServiceName = "dev.ensan.inputmethod.azooKeyMac.ConverterServer" -} - -@objc private protocol ConverterServerXPCProtocol { - func openSession(with reply: @escaping @Sendable (String) -> Void) - func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) - func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) - func ping(_ message: String, with reply: @escaping @Sendable (String) -> Void) -} - -final class ConverterServer: NSObject, ConverterServerXPCProtocol, @unchecked Sendable { - private static let learningDataCommitDelay: TimeInterval = 2 - - private var sessions: [String: ConverterSession] = [:] - private let kanaKanjiConverter = KanaKanjiConverter.withDefaultDictionary() - private let learningDataCommitScheduler = DebouncedActionScheduler() - - func openSession(with reply: @escaping @Sendable (String) -> Void) { - DispatchQueue.main.async { - MainActor.assumeIsolated { - let sessionID = UUID().uuidString - self.createSessionIfNeeded(sessionID) - reply(sessionID) - } - } - } - - func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) { - DispatchQueue.main.async { - MainActor.assumeIsolated { - let session = self.sessions.removeValue(forKey: sessionID) - if let session { - self.kanaKanjiConverter.removeSession(session.conversionSessionID) - } - let removed = session != nil - reply(removed) - } - } - } - - func ping(_ message: String, with reply: @escaping @Sendable (String) -> Void) { - reply("ConverterServer: \(message)") - } - - func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) { - // キー入力の応答はユーザー操作のクリティカルパスなので、システム負荷が高い時も - // utility/background work より先に実行される優先度で Server actor へ渡す。 - Task(priority: .userInitiated) { @MainActor in - do { - let command = try ConverterServerCodec.decodeCommand(from: data) - let response = try await self.handle(command) - self.learningDataCommitScheduler.postponeIfScheduled( - after: Self.learningDataCommitDelay - ) - reply(try ConverterServerCodec.encode(response), nil) - } catch { - self.learningDataCommitScheduler.postponeIfScheduled( - after: Self.learningDataCommitDelay - ) - reply(nil, error.localizedDescription as NSString) - } - } - } - - @MainActor - private func handle(_ command: ConverterServerCommand) async throws -> ConverterServerResponse { - switch command { - case .shutdown: - Self.scheduleShutdown() - return ConverterServerResponse(snapshot: .empty) - case .maintenance(let command): - return try handle(command) - case .openSession(let sessionID, let command): - createSessionIfNeeded(sessionID) - return try await handle(command, sessionID: sessionID) - case .session(let sessionID, let command): - return try await handle(command, sessionID: sessionID) - } - } - - @MainActor - private func createSessionIfNeeded(_ sessionID: String) { - guard sessions[sessionID] == nil else { - return - } - let conversionSessionID = kanaKanjiConverter.createSession() - sessions[sessionID] = ConverterSession( - manager: Self.makeSegmentsManager(kanaKanjiConverter: kanaKanjiConverter), - conversionSessionID: conversionSessionID - ) - } - - @MainActor - private func handle(_ command: ConverterMaintenanceCommand) throws -> ConverterServerResponse { - switch command { - case .synchronizeUserDictionary(let forceExport): - let memoryDirectoryURL = AppGroup.memoryDirectoryURL() - if forceExport || !CompiledUserDictionaryStore.hasExportedDictionary(memoryDirectoryURL: memoryDirectoryURL) { - try CompiledUserDictionaryStore.exportCurrentDictionaries(memoryDirectoryURL: memoryDirectoryURL) - } - kanaKanjiConverter.updateUserDictionaryURL( - CompiledUserDictionaryStore.directoryURL(memoryDirectoryURL: memoryDirectoryURL), - forceReload: true - ) - case .resetLearningData: - kanaKanjiConverter.resetMemory() - } - return ConverterServerResponse(snapshot: .empty) - } - - @MainActor - private func handle(_ command: ConverterSessionCommand, sessionID: String) async throws -> ConverterServerResponse { - let session = try getSession(sessionID) - switch command { - case .lifecycle(let command): - return try withConverterSession(session) { - handle(command, session: session) - } - case .settings(let command): - return try withConverterSession(session) { - try handle(command, session: session) - } - case .updateConfig(let config): - return try withConverterSession(session) { - session.config = config - return makeResponse(for: session, inputState: .none) - } - case .handleKeyEvent(let request): - return try withConverterSession(session) { - try handleKeyEvent(sessionID: sessionID, request: request) - } - case .composition(let command): - return try withConverterSession(session) { - handle(command, session: session) - } - case .candidate(let command): - return try withConverterSession(session) { - handle(command, session: session) - } - case .replaceSuggestion(let command): - return try await handle(command, session: session) - } - } - - @MainActor - private func withConverterSession( - _ session: ConverterSession, - operation: () throws -> Result - ) throws -> Result { - try kanaKanjiConverter.withSession(session.conversionSessionID, operation: operation) - } - - @MainActor - private func handle( - _ command: ConverterSessionLifecycleCommand, - session: ConverterSession - ) -> ConverterServerResponse { - switch command { - case .activate: - session.manager.activate() - return makeResponse(for: session, inputState: session.inputState) - case .deactivate: - // アプリ切替直後のキー入力を、学習データの同期I/Oで塞がない。 - // 共有Converterはプロセス内に残るため、永続化だけ入力のアイドル時まで遅延できる。 - session.manager.deactivate(flushLearningData: false) - scheduleLearningDataCommit() - session.inputState = .none - session.clearReplaceSuggestions() - return makeResponse(for: session, inputState: session.inputState) - case .synchronizeInputLanguage(let language): - session.inputLanguage = language - if language == .english { - session.manager.stopJapaneseInput() - } - return makeResponse( - for: session, - inputState: session.inputState - ) - } - } - - @MainActor - private func handle( - _ command: ConverterSettingsCommand, - session: ConverterSession - ) throws -> ConverterServerResponse { - switch command { - case .list(let capabilities): - return makeResponse( - for: session, - inputState: .none, - settings: Self.makeSettingDescriptors(capabilities: capabilities) - ) - case .update(let key, let value): - try Self.updateSetting(key: key, value: value) - return makeResponse(for: session, inputState: .none) - } - } - - @MainActor - private func handle( - _ command: ConverterCompositionCommand, - session: ConverterSession - ) -> ConverterServerResponse { - switch command { - case .snapshot: - return makeResponse(for: session, inputState: session.inputState) - case .stopComposition: - session.manager.stopComposition() - session.inputState = .none - return makeResponse(for: session, inputState: session.inputState) - case .forgetMemory: - session.manager.forgetMemory() - return makeResponse(for: session, inputState: session.inputState) - case .commit: - let text = session.manager.commitMarkedText(inputState: session.inputState) - let effects: [ConverterClientEffect] = text.isEmpty ? [] : [.insertText(text)] - session.inputState = .none - return makeResponse( - for: session, - inputState: session.inputState, - effects: effects, - responseInputState: ConverterInputState.none - ) - } - } - - @MainActor - private func handle( - _ command: ConverterCandidateCommand, - session: ConverterSession - ) -> ConverterServerResponse { - switch command { - case .selectCandidate(let index): - session.manager.requestSelectingRow(index) - session.inputState = .selecting - return makeResponse(for: session, inputState: session.inputState) - case .submitSelectedCandidate(let context): - session.setContext(context) - var effects: [ConverterClientEffect] = [] - submitSelectedCandidate( - manager: session.manager, - leftSideContext: session.conversionLeftSideContext(), - effects: &effects - ) - let nextInputState: InputState = session.manager.isEmpty ? .none : .previewing - session.inputState = nextInputState - return makeResponse( - for: session, - inputState: nextInputState, - effects: effects, - responseInputState: ConverterInputState(nextInputState) - ) - } - } - - @MainActor - private func handle( - _ command: ConverterReplaceSuggestionCommand, - session: ConverterSession - ) async throws -> ConverterServerResponse { - switch command { - case .request(let context): - session.setContext(context) - try await requestReplaceSuggestion(session: session) - return try withConverterSession(session) { - session.inputState = .replaceSuggestion - return makeResponse( - for: session, - inputState: session.inputState, - responseInputState: .replaceSuggestion - ) - } - case .selectReplaceSuggestionCandidate(let index): - return try withConverterSession(session) { - session.selectReplaceSuggestion(at: index) - session.inputState = .replaceSuggestion - return makeResponse( - for: session, - inputState: session.inputState, - responseInputState: .replaceSuggestion - ) - } - case .submitSelectedReplaceSuggestion: - return try withConverterSession(session) { - var effects: [ConverterClientEffect] = [] - let didSubmit = submitSelectedReplaceSuggestion(session: session, effects: &effects) - let nextInputState: InputState = didSubmit ? .none : .replaceSuggestion - session.inputState = nextInputState - return makeResponse( - for: session, - inputState: nextInputState, - effects: effects, - responseInputState: ConverterInputState(nextInputState) - ) - } - } - } - - private static func scheduleShutdown() { - DispatchQueue.main.asyncAfter(deadline: .now() + 0.2) { - exit(EXIT_SUCCESS) - } - } - - @MainActor - private func scheduleLearningDataCommit() { - learningDataCommitScheduler.schedule(after: Self.learningDataCommitDelay) { [weak self] in - self?.kanaKanjiConverter.commitUpdateLearningData() - } - } - - @MainActor - func getSession(_ sessionID: String) throws -> ConverterSession { - guard let session = sessions[sessionID] else { - throw ConverterServerError.unknownSession(sessionID) - } - return session - } - -} private final class ServiceDelegate: NSObject, NSXPCListenerDelegate { private let server = ConverterServer() func listener(_ listener: NSXPCListener, shouldAcceptNewConnection connection: NSXPCConnection) -> Bool { + let handler = ConverterServerConnection(server: server) connection.exportedInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) - connection.exportedObject = server + connection.exportedObject = handler + connection.invalidationHandler = { handler.invalidate() } + connection.interruptionHandler = { handler.invalidate() } connection.resume() return true } diff --git a/Core/Sources/Core/XPC/ConverterClientSessionState.swift b/Core/Sources/Core/XPC/ConverterClientSessionState.swift new file mode 100644 index 00000000..875d653f --- /dev/null +++ b/Core/Sources/Core/XPC/ConverterClientSessionState.swift @@ -0,0 +1,32 @@ +/// Client が決める入力モードと、Server の再初期化が必要かだけを保持する。 +/// activation の値を保存しない。送信直前のモード・設定から作ることで、 +/// activate → モード変更 → 最初のキー、の順でも古いモードへ戻らない。 +public struct ConverterClientSessionState: Sendable { + public var inputLanguage: InputLanguage = .japanese + public private(set) var isActive = false + public private(set) var needsActivation = true + + public init() {} + + public mutating func activate() { + isActive = true + needsActivation = true + } + + public mutating func deactivate() { + isActive = false + needsActivation = true + } + + public mutating func connectionDidReset() { + needsActivation = true + } + + public mutating func takeActivation(config: @autoclosure () -> ConverterSessionConfig) -> ConverterSessionActivation? { + guard needsActivation else { + return nil + } + needsActivation = false + return .init(config: config(), inputLanguage: inputLanguage) + } +} diff --git a/Core/Tests/ConverterServerTests/ConverterServerLifecycleTests.swift b/Core/Tests/ConverterServerTests/ConverterServerLifecycleTests.swift new file mode 100644 index 00000000..c8e2b2a5 --- /dev/null +++ b/Core/Tests/ConverterServerTests/ConverterServerLifecycleTests.swift @@ -0,0 +1,170 @@ +@testable import ConverterServer +import Core +import Foundation +import Testing + +@MainActor +private func makeServer() -> ConverterServer { + ConverterServer(makeManager: { + SegmentsManager( + kanaKanjiConverter: $0, + applicationDirectoryURL: FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString), + containerURL: nil, + context: .init(useZenzai: false) + ) + }, applyRequestSettings: { _ in }) +} + +private let config = ConverterSessionConfig( + aiBackendPreference: .off, + openAIModelName: "test", + openAIEndpoint: "https://example.com", + openAIAPIKey: .init(""), + includeContextInAITransform: false +) + +private func key(_ activation: ConverterSessionActivation?, id: UInt64 = 1) -> ConverterSessionCommand { + .handleKeyEvent(.init( + eventID: id, + event: .init(modifierFlags: [], characters: "b", charactersIgnoringModifiers: "b", keyCode: 11), + inputStyle: .defaultRomanToKana, + liveConversionEnabled: false, + enableDebugWindow: false, + enableSuggestion: false, + activation: activation + )) +} + +@MainActor +@Test func firstJapaneseKeyAfterActivationInEnglishIsNotInsertedAsRomanText() async throws { + let server = makeServer() + let owner = UUID() + var state = ConverterClientSessionState() + state.inputLanguage = .english + state.activate() + state.inputLanguage = .japanese // setValue が activate より後に来る順序 + _ = try await server.execute(.openSession(sessionID: "session", command: .lifecycle(.synchronizeInputLanguage(.japanese))), owner: owner) + let response = try await server.execute( + .session(sessionID: "session", command: key(state.takeActivation(config: config))), owner: owner + ) + #expect(response.inputLanguage == .japanese) + #expect(response.inputState == .composing) + #expect(response.handled) + #expect(!response.snapshot.isEmpty) + #expect(!response.effects.contains(.insertText("b"))) + server.removeSessions(ownedBy: owner) +} + +@MainActor +@Test func oldCapturedActivationReproducesRomanInputBug() async throws { + let server = makeServer() + let owner = UUID() + let stale = ConverterSessionActivation(config: config, inputLanguage: .english) + _ = try await server.execute(.openSession(sessionID: "session", command: .lifecycle(.synchronizeInputLanguage(.japanese))), owner: owner) + let response = try await server.execute(.session(sessionID: "session", command: key(stale)), owner: owner) + #expect(response.inputLanguage == .english) + #expect(response.effects.contains(.insertText("b"))) + server.removeSessions(ownedBy: owner) +} + +@MainActor +@Test func disconnectedOwnersAreCollectedWithoutDeletingReconnectedSession() async throws { + let server = makeServer() + let oldOwner = UUID() + let newOwner = UUID() + _ = try await server.execute(.openSession(sessionID: "abandoned", command: .composition(.snapshot)), owner: oldOwner) + _ = try await server.execute(.openSession(sessionID: "resumed", command: .composition(.snapshot)), owner: oldOwner) + _ = try await server.execute(.session(sessionID: "resumed", command: .composition(.snapshot)), owner: newOwner) + server.removeSessions(ownedBy: oldOwner) + #expect(throws: (any Error).self) { try server.getSession("abandoned") } + #expect(throws: Never.self) { try server.getSession("resumed") } + server.removeSessions(ownedBy: newOwner) + #expect(throws: (any Error).self) { try server.getSession("resumed") } +} + +@MainActor +@Test func deactivationClearsCompositionBeforeNextActivation() async throws { + let server = makeServer() + let owner = UUID() + let activation = ConverterSessionActivation(config: config, inputLanguage: .japanese) + _ = try await server.execute(.openSession(sessionID: "session", command: key(activation)), owner: owner) + _ = try await server.execute(.session(sessionID: "session", command: .lifecycle(.deactivate)), owner: owner) + let response = try await server.execute(.session(sessionID: "session", command: key(activation, id: 2)), owner: owner) + #expect(response.snapshot.convertTarget == "b", "前回の b が復活して bb にならない") + server.removeSessions(ownedBy: owner) +} + +private final class TestListener: NSObject, NSXPCListenerDelegate, @unchecked Sendable { + let server: ConverterServer + let listener = NSXPCListener.anonymous() + + init(server: ConverterServer) { + self.server = server + super.init() + listener.delegate = self + listener.resume() + } + + func listener(_ listener: NSXPCListener, shouldAcceptNewConnection connection: NSXPCConnection) -> Bool { + let handler = ConverterServerConnection(server: server, cleanupDelay: 0) + connection.exportedInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) + connection.exportedObject = handler + connection.invalidationHandler = { handler.invalidate() } + connection.resume() + return true + } +} + +@MainActor +@Test func actualXPCDisconnectCollectsServerSession() async throws { + let server = makeServer() + let service = TestListener(server: server) + defer { service.listener.invalidate() } + let connection = NSXPCConnection(listenerEndpoint: service.listener.endpoint) + connection.remoteObjectInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) + connection.resume() + defer { connection.invalidate() } + let data = try ConverterServerCodec.encode(ConverterServerCommand.openSession( + sessionID: "xpc-session", command: .composition(.snapshot) + )) + let result: Data = try await withCheckedThrowingContinuation { continuation in + let proxy = connection.remoteObjectProxyWithErrorHandler { continuation.resume(throwing: $0) } + guard let proxy = proxy as? ConverterServerXPCProtocol else { + continuation.resume(throwing: NSError(domain: "Invalid proxy", code: 1)) + return + } + proxy.handleCommand(data) { response, error in + if let response { + continuation.resume(returning: response) + } else { + continuation.resume(throwing: NSError(domain: error.map(String.init) ?? "Missing response", code: 1)) + } + } + } + #expect(try ConverterServerCodec.decodeResponse(from: result).snapshot.isEmpty) + #expect(throws: Never.self) { try server.getSession("xpc-session") } + connection.invalidate() + let deadline = Date().addingTimeInterval(2) + while (try? server.getSession("xpc-session")) != nil, Date() < deadline { + try await Task.sleep(nanoseconds: 10_000_000) + } + #expect(throws: (any Error).self) { try server.getSession("xpc-session") } +} + +@MainActor +@Test func invalidatedConnectionCannotCreateAnotherSession() async throws { + let server = makeServer() + let handler = ConverterServerConnection(server: server, cleanupDelay: 0) + handler.invalidate() + let data = try ConverterServerCodec.encode(ConverterServerCommand.openSession( + sessionID: "late", command: .composition(.snapshot) + )) + let result: (Data?, String?) = await withCheckedContinuation { continuation in + handler.handleCommand(data) { response, error in + continuation.resume(returning: (response, error.map(String.init))) + } + } + #expect(result.0 == nil) + #expect(result.1 != nil) + #expect(throws: (any Error).self) { try server.getSession("late") } +} diff --git a/Core/Tests/CoreTests/XPCTests/ConverterClientSessionStateTests.swift b/Core/Tests/CoreTests/XPCTests/ConverterClientSessionStateTests.swift new file mode 100644 index 00000000..e104e970 --- /dev/null +++ b/Core/Tests/CoreTests/XPCTests/ConverterClientSessionStateTests.swift @@ -0,0 +1,59 @@ +import Core +import Testing + +private let sessionConfig = ConverterSessionConfig( + aiBackendPreference: .off, + openAIModelName: "model", + openAIEndpoint: "https://example.com", + openAIAPIKey: .init(""), + includeContextInAITransform: false +) + +@Test(arguments: [InputLanguage.japanese, .english]) +func activationUsesModeChangedAfterActivate(language: InputLanguage) { + var state = ConverterClientSessionState() + state.inputLanguage = language == .japanese ? .english : .japanese + state.activate() + state.inputLanguage = language + #expect(state.takeActivation(config: sessionConfig)?.inputLanguage == language) + #expect(state.takeActivation(config: sessionConfig) == nil) +} + +@Test(arguments: [InputLanguage.japanese, .english]) +func activationUsesModeChangedBeforeActivate(language: InputLanguage) { + var state = ConverterClientSessionState() + state.inputLanguage = language + state.activate() + #expect(state.takeActivation(config: sessionConfig)?.inputLanguage == language) +} + +@Test func resetAndReactivationUseLatestModeAndConfiguration() { + var state = ConverterClientSessionState() + state.activate() + _ = state.takeActivation(config: sessionConfig) + state.connectionDidReset() + state.inputLanguage = .english + var updated = sessionConfig + updated.openAIModelName = "updated" + #expect(state.takeActivation(config: updated) == .init(config: updated, inputLanguage: .english)) + state.deactivate() + #expect(!state.isActive) + state.connectionDidReset() + #expect(!state.isActive) + state.inputLanguage = .japanese + state.activate() + #expect(state.isActive) + #expect(state.takeActivation(config: updated)?.inputLanguage == .japanese) +} + +@Test func ordinaryKeysDoNotReadActivationConfigurationAgain() { + var state = ConverterClientSessionState() + var reads = 0 + func readConfig() -> ConverterSessionConfig { + reads += 1 + return sessionConfig + } + _ = state.takeActivation(config: readConfig()) + _ = state.takeActivation(config: readConfig()) + #expect(reads == 1) +} diff --git a/azooKeyMac/InputController/ConverterServerClient.swift b/azooKeyMac/InputController/ConverterServerClient.swift index 94697605..75ba8d34 100644 --- a/azooKeyMac/InputController/ConverterServerClient.swift +++ b/azooKeyMac/InputController/ConverterServerClient.swift @@ -1,6 +1,27 @@ import Core import Foundation +/// MainActor の Client が破棄された場合も接続を終了させる。 +/// deinit から actor へ Task を投げて Client 自身を延命しない。 +private final class ConverterConnectionLifetime: @unchecked Sendable { + let connection: NSXPCConnection + var sessionID: String? + + init(_ connection: NSXPCConnection) { self.connection = connection } + + deinit { + let connection = connection + guard let sessionID else { + connection.invalidate() + return + } + // 接続単位の回収に未対応の旧Serverにも、可能なら明示的に解放を通知する。 + let proxy = connection.remoteObjectProxyWithErrorHandler { _ in connection.invalidate() } + (proxy as? ConverterServerXPCProtocol)?.closeSession(sessionID) { _ in connection.invalidate() } + DispatchQueue.global().asyncAfter(deadline: .now() + 1) { connection.invalidate() } + } +} + @objc protocol ConverterServerXPCProtocol { func openSession(with reply: @escaping @Sendable (String) -> Void) func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) @@ -37,7 +58,8 @@ final class ConverterServerClient { private let keyEventTimeout: TimeInterval private let commandTimeout: TimeInterval private let connectionFactory: @Sendable () -> NSXPCConnection - private var connection: NSXPCConnection? + private var connectionLifetime: ConverterConnectionLifetime? + private var connection: NSXPCConnection? { connectionLifetime?.connection } private var sessionID: String? private var abandonedSessionIDs: [String] = [] private var pendingCommands: [PendingCommand] = [] @@ -97,7 +119,9 @@ final class ConverterServerClient { _ commandBuilder: @escaping (String) -> ConverterSessionCommand, completion: @escaping (ConverterServerResponse?) -> Void ) { - guard sessionID != nil else { + // session 作成中の deactivate 等を落とすと、Client だけ composition が + // 消え、次の activate で Server の古い入力が復活する。 + guard sessionID != nil || activeCommand?.openingSessionID != nil else { completion(nil) return } @@ -106,7 +130,14 @@ final class ConverterServerClient { /// 応答を受け取る XPC キューはブロックしない。呼び出し元だけが期限付きで待つ。 func sendKeyEvent(_ request: ConverterKeyEventRequest) -> ConverterServerResponse? { - sendSynchronously { _ in .handleKeyEvent(request) } + let started = Date() + let response = sendSynchronously { _ in .handleKeyEvent(request) } + let duration = Date().timeIntervalSince(started) + if response == nil || duration > 0.2 { + // 入力文字・前後文脈は記録しない。 + onLog?("ConverterServer key \(request.eventID): success=\(response != nil), elapsed=\(duration)s") + } + return response } func sendSynchronously( @@ -170,6 +201,7 @@ final class ConverterServerClient { do { let data = try ConverterServerCodec.encode(command) let connection = ensureConnection() + connectionLifetime?.sessionID = sessionID ?? openingSessionID if let proxy = connection.remoteObjectProxyWithErrorHandler({ error in complete(.failure(error.localizedDescription)) }) as? ConverterServerXPCProtocol { @@ -216,6 +248,7 @@ final class ConverterServerClient { if response != nil { if let openingSessionID = active.openingSessionID { sessionID = openingSessionID + connectionLifetime?.sessionID = openingSessionID } pending.completion(response) startNextCommand() @@ -238,7 +271,7 @@ final class ConverterServerClient { let connection = connectionFactory() connection.remoteObjectInterface = NSXPCInterface(with: ConverterServerXPCProtocol.self) connection.resume() - self.connection = connection + self.connectionLifetime = ConverterConnectionLifetime(connection) // 切断前の計算は継続している場合がある。新しい接続で旧 session を回収する。 if !abandonedSessionIDs.isEmpty, let proxy = connection.remoteObjectProxyWithErrorHandler({ _ in }) as? ConverterServerXPCProtocol { @@ -258,7 +291,8 @@ final class ConverterServerClient { } sessionID = nil connection?.invalidate() - connection = nil + connectionLifetime?.sessionID = nil + connectionLifetime = nil onSessionReset?() } } diff --git a/azooKeyMac/InputController/azooKeyMacInputController.swift b/azooKeyMac/InputController/azooKeyMacInputController.swift index 09756b99..d58bee01 100644 --- a/azooKeyMac/InputController/azooKeyMacInputController.swift +++ b/azooKeyMac/InputController/azooKeyMacInputController.swift @@ -1,3 +1,4 @@ +import Carbon import Cocoa import Core import InputMethodKit @@ -7,10 +8,13 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s let converterServerClient = ConverterServerClient() private var currentConverterView: ConverterSessionSnapshot? private(set) var inputState: InputState = .none - private var inputLanguage: InputLanguage = .japanese + private var converterSessionState = ConverterClientSessionState() + private var inputLanguage: InputLanguage { + get { self.converterSessionState.inputLanguage } + set { self.converterSessionState.inputLanguage = newValue } + } private var nextKeyEventID: UInt64 = 0 private var activationGeneration: UInt64 = 0 - private var pendingConverterServerActivation: ConverterSessionActivation? var liveConversionEnabled: Bool { Config.LiveConversion().value } @@ -153,10 +157,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s self.updateTransformSelectedTextMenuItemEnabledState() // ピン留めプロンプトのキャッシュを更新 self.reloadPinnedPromptsCache() - self.pendingConverterServerActivation = ConverterSessionActivation( - config: self.converterServerSessionConfig, - inputLanguage: self.inputLanguage - ) + self.converterSessionState.activate() if let client = sender as? IMKTextInput { client.overrideKeyboard(withKeyboardNamed: Config.KeyboardLayout().value.layoutIdentifier) @@ -177,7 +178,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s @MainActor override func deactivateServer(_ sender: Any!) { self.activationGeneration &+= 1 - self.pendingConverterServerActivation = nil + self.converterSessionState.deactivate() self.converterServerClient.sendIfSessionOpen({ _ in .lifecycle(.deactivate) }, completion: { _ in }) self.currentConverterView = nil self.inputState = .none @@ -200,7 +201,8 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s super.setValue(value, forTag: tag, client: sender) } - if let value = value as? NSString { + if tag == kTextServiceInputModePropertyTag, let value = value as? NSString, + value == "com.apple.inputmethod.Roman" || value == "com.apple.inputmethod.Japanese" { self.client()?.overrideKeyboard(withKeyboardNamed: Config.KeyboardLayout().value.layoutIdentifier) let englishMode = value == "com.apple.inputmethod.Roman" @@ -210,10 +212,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s // composing中でも英数キーMarkedTextを保ったまま英語入力へ移る。 if self.inputLanguage == .japanese { self.inputLanguage = .english - self.converterServerClient.send( - { _ in .lifecycle(.synchronizeInputLanguage(.english)) }, - completion: { _ in } - ) + self.synchronizeConverterInputLanguage() self.refreshCandidateWindow() self.refreshPredictionWindow() } @@ -221,10 +220,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s // 日本語モードへの切り替え if self.inputLanguage == .english { self.inputLanguage = .japanese - self.converterServerClient.send( - { _ in .lifecycle(.synchronizeInputLanguage(.japanese)) }, - completion: { _ in } - ) + self.synchronizeConverterInputLanguage() } } } @@ -234,6 +230,17 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s self.appMenu } + @MainActor + private func synchronizeConverterInputLanguage() { + // 初回キー前は activation に最新モードを同梱する。設定通知だけで + // session を開かず、短い非同期タイムアウトにも依存させない。 + guard !self.converterSessionState.needsActivation else { + return + } + let language = self.inputLanguage + self.sendAndApply { _ in .lifecycle(.synchronizeInputLanguage(language)) } + } + // swiftlint:disable:next cyclomatic_complexity @MainActor override func handle(_ event: NSEvent!, client sender: Any!) -> Bool { guard let event, let client = sender as? IMKTextInput else { @@ -364,9 +371,8 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s typeBackSlash: Config.TypeBackSlash().value, optionDirectInputText: optionDirectInputText, context: self.currentConverterTextContext(), - activation: self.pendingConverterServerActivation + activation: self.converterSessionState.takeActivation(config: self.converterServerSessionConfig) ) - self.pendingConverterServerActivation = nil guard let response = self.converterServerClient.sendKeyEvent(request) else { return false } @@ -381,10 +387,11 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s self.currentConverterView = nil self.inputState = .none self.activationGeneration &+= 1 - self.pendingConverterServerActivation = ConverterSessionActivation( - config: self.converterServerSessionConfig, - inputLanguage: self.inputLanguage - ) + self.converterSessionState.connectionDidReset() + // 非アクティブな client に遅れて届いた切断で、別アプリへ文字を挿入しない。 + guard self.converterSessionState.isActive else { + return + } if !text.isEmpty { self.client()?.insertText(text, replacementRange: NSRange(location: NSNotFound, length: 0)) self.refreshMarkedText() @@ -409,13 +416,16 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s @MainActor private func sendAndApply(_ command: @escaping (String) -> ConverterSessionCommand) { - if let response = self.converterServerClient.sendSynchronously(command, onlyIfSessionOpen: true) { + let generation = self.activationGeneration + if let response = self.converterServerClient.sendSynchronously(command, onlyIfSessionOpen: true), + self.activationGeneration == generation { self.apply(response) } } @MainActor private func apply(_ response: ConverterServerResponse) { + let hadMarkedText = !self.currentMarkedText().elements.isEmpty if let inputLanguage = response.inputLanguage { self.inputLanguage = inputLanguage } @@ -426,7 +436,11 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s self.apply(effect, client: client) } } - self.refreshMarkedText() + // composition がないモード変更等で空の marked text を送ると、 + // ホスト側の選択範囲まで置換してしまうことがある。 + if hadMarkedText || !self.currentMarkedText().elements.isEmpty { + self.refreshMarkedText() + } self.refreshCandidateWindow() self.refreshPredictionWindow() self.refreshReplaceSuggestionWindow() diff --git a/azooKeyMacTests/ThinClientInputPipelineTests.swift b/azooKeyMacTests/ThinClientInputPipelineTests.swift index f3f88e28..7047575a 100644 --- a/azooKeyMacTests/ThinClientInputPipelineTests.swift +++ b/azooKeyMacTests/ThinClientInputPipelineTests.swift @@ -16,6 +16,7 @@ private final class TestConverterService: NSObject, NSXPCListenerDelegate, Conve private let lock = NSLock() private var connections: [NSXPCConnection] = [] private var receivedCommands: [ConverterServerCommand] = [] + private var closedSessionIDs: [String] = [] private let responses: [Response] init(responses: [Response]) { @@ -41,6 +42,12 @@ private final class TestConverterService: NSObject, NSXPCListenerDelegate, Conve return receivedCommands } + var closedSessions: [String] { + lock.lock() + defer { lock.unlock() } + return closedSessionIDs + } + func stop() { lock.lock() let connections = self.connections @@ -51,7 +58,12 @@ private final class TestConverterService: NSObject, NSXPCListenerDelegate, Conve } func openSession(with reply: @escaping @Sendable (String) -> Void) { reply(UUID().uuidString) } - func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) { reply(true) } + func closeSession(_ sessionID: String, with reply: @escaping @Sendable (Bool) -> Void) { + lock.lock() + closedSessionIDs.append(sessionID) + lock.unlock() + reply(true) + } func ping(_ message: String, with reply: @escaping @Sendable (String) -> Void) { reply(message) } func handleCommand(_ data: Data, with reply: @escaping @Sendable (Data?, NSString?) -> Void) { @@ -193,4 +205,42 @@ final class ThinClientInputPipelineTests: XCTestCase { XCTAssertEqual(completions, 2) XCTAssertEqual(service.commands.count, 1) } + + func testDeactivationIsQueuedEvenWhileSessionIsOpening() throws { + let service = TestConverterService(responses: [ + .init(data: try encodedResponse(state: .composing), delay: 0.03), + .init(data: try encodedResponse()) + ]) + defer { service.stop() } + let client = makeClient(service) + client.send({ _ in .composition(.snapshot) }, completion: { _ in }) + var deactivated = false + client.sendIfSessionOpen({ _ in .lifecycle(.deactivate) }, completion: { response in + deactivated = response != nil + }) + client.flushPendingCommands() + XCTAssertTrue(deactivated) + XCTAssertEqual(service.commands.count, 2) + guard case .session(_, .lifecycle(.deactivate)) = service.commands.last else { + return XCTFail("deactivate must not be silently discarded while opening") + } + } + + func testClientDestructionClosesSessionWithoutAnotherKey() async throws { + let service = TestConverterService(responses: [.init(data: try encodedResponse())]) + defer { service.stop() } + var client: ConverterServerClient? = makeClient(service) + XCTAssertNotNil(client?.sendSynchronously { _ in .composition(.snapshot) }) + guard case .openSession(let id, _) = service.commands.first else { + return XCTFail("Expected a session") + } + weak var weakClient = client + client = nil + XCTAssertNil(weakClient) + let deadline = Date().addingTimeInterval(2) + while service.closedSessions.isEmpty && Date() < deadline { + try await Task.sleep(nanoseconds: 10_000_000) + } + XCTAssertEqual(service.closedSessions, [id]) + } } From 505345ac538c787e704bbd18850031e14f05700f Mon Sep 17 00:00:00 2001 From: ensan-hcl Date: Sun, 6 Sep 2026 14:38:22 +0900 Subject: [PATCH 3/3] refactor(input): simplify converter command routing and mode synchronization --- .../ConverterServer/ConverterServer.swift | 39 +++++--------- .../ConverterServerConnection.swift | 10 ++-- .../ConverterServerClient.swift | 19 +++---- .../azooKeyMacInputController.swift | 52 ++++++++----------- .../ThinClientInputPipelineTests.swift | 18 +++---- 5 files changed, 59 insertions(+), 79 deletions(-) diff --git a/Core/Sources/ConverterServer/ConverterServer.swift b/Core/Sources/ConverterServer/ConverterServer.swift index d1acdf32..e28aeef7 100644 --- a/Core/Sources/ConverterServer/ConverterServer.swift +++ b/Core/Sources/ConverterServer/ConverterServer.swift @@ -39,17 +39,22 @@ final class ConverterServer: @unchecked Sendable { @MainActor func execute(_ command: ConverterServerCommand, owner: UUID) async throws -> ConverterServerResponse { + defer { learningDataCommitScheduler.postponeIfScheduled(after: Self.learningDataCommitDelay) } switch command { - case .openSession(let id, _): - sessionOwners[id] = owner - case .session(let id, _): - _ = try getSession(id) - sessionOwners[id] = owner - case .shutdown, .maintenance: - break + case .shutdown: + Self.scheduleShutdown() + return ConverterServerResponse(snapshot: .empty) + case .maintenance(let command): + return try handle(command) + case .openSession(let sessionID, let command): + createSessionIfNeeded(sessionID) + sessionOwners[sessionID] = owner + return try await handle(command, sessionID: sessionID) + case .session(let sessionID, let command): + _ = try getSession(sessionID) + sessionOwners[sessionID] = owner + return try await handle(command, sessionID: sessionID) } - defer { learningDataCommitScheduler.postponeIfScheduled(after: Self.learningDataCommitDelay) } - return try await handle(command) } @MainActor @@ -71,22 +76,6 @@ final class ConverterServer: @unchecked Sendable { return true } - @MainActor - private func handle(_ command: ConverterServerCommand) async throws -> ConverterServerResponse { - switch command { - case .shutdown: - Self.scheduleShutdown() - return ConverterServerResponse(snapshot: .empty) - case .maintenance(let command): - return try handle(command) - case .openSession(let sessionID, let command): - createSessionIfNeeded(sessionID) - return try await handle(command, sessionID: sessionID) - case .session(let sessionID, let command): - return try await handle(command, sessionID: sessionID) - } - } - @MainActor private func createSessionIfNeeded(_ sessionID: String) { guard sessions[sessionID] == nil else { diff --git a/Core/Sources/ConverterServer/ConverterServerConnection.swift b/Core/Sources/ConverterServer/ConverterServerConnection.swift index df55403a..b1d94396 100644 --- a/Core/Sources/ConverterServer/ConverterServerConnection.swift +++ b/Core/Sources/ConverterServer/ConverterServerConnection.swift @@ -34,12 +34,10 @@ final class ConverterServerConnection: NSObject, ConverterServerXPCProtocol, @un guard markClosed() else { return } - Task { @MainActor in - // 旧Clientの再接続・同一sessionへの再送を短時間だけ許容する。 - // 新接続が所有権を取得したsessionは、古い接続のcleanupでは削除しない。 - DispatchQueue.main.asyncAfter(deadline: .now() + self.cleanupDelay) { - self.server.removeSessions(ownedBy: self.owner) - } + // 旧Clientの再接続・同一sessionへの再送を短時間だけ許容する。 + // 新接続が所有権を取得したsessionは、古い接続のcleanupでは削除しない。 + DispatchQueue.main.asyncAfter(deadline: .now() + self.cleanupDelay) { + self.server.removeSessions(ownedBy: self.owner) } } diff --git a/azooKeyMac/InputController/ConverterServerClient.swift b/azooKeyMac/InputController/ConverterServerClient.swift index 75ba8d34..1b1e3800 100644 --- a/azooKeyMac/InputController/ConverterServerClient.swift +++ b/azooKeyMac/InputController/ConverterServerClient.swift @@ -32,7 +32,8 @@ private final class ConverterConnectionLifetime: @unchecked Sendable { @MainActor final class ConverterServerClient { private enum Command { - case session((String) -> ConverterSessionCommand) + // 文脈の取得などは、先行コマンドの応答を反映してから行う。 + case session(() -> ConverterSessionCommand) case global(ConverterServerCommand) } @@ -84,11 +85,11 @@ final class ConverterServerClient { capabilities: ConverterSettingClientCapabilities, completion: @escaping ([ConverterSettingDescriptor]?) -> Void ) { - send({ _ in .settings(.list(capabilities: capabilities)) }, completion: { completion($0?.settings) }) + send({ .settings(.list(capabilities: capabilities)) }, completion: { completion($0?.settings) }) } func updateSetting(key: String, value: ConverterSettingValue, completion: @escaping (Bool) -> Void) { - send({ _ in .settings(.update(key: key, value: value)) }, completion: { completion($0 != nil) }) + send({ .settings(.update(key: key, value: value)) }, completion: { completion($0 != nil) }) } func restartServer(completion: @escaping (Bool) -> Void) { @@ -109,14 +110,14 @@ final class ConverterServerClient { } func send( - _ commandBuilder: @escaping (String) -> ConverterSessionCommand, + _ commandBuilder: @escaping () -> ConverterSessionCommand, completion: @escaping (ConverterServerResponse?) -> Void ) { enqueue(.session(commandBuilder), timeout: commandTimeout, completion: completion) } func sendIfSessionOpen( - _ commandBuilder: @escaping (String) -> ConverterSessionCommand, + _ commandBuilder: @escaping () -> ConverterSessionCommand, completion: @escaping (ConverterServerResponse?) -> Void ) { // session 作成中の deactivate 等を落とすと、Client だけ composition が @@ -131,7 +132,7 @@ final class ConverterServerClient { /// 応答を受け取る XPC キューはブロックしない。呼び出し元だけが期限付きで待つ。 func sendKeyEvent(_ request: ConverterKeyEventRequest) -> ConverterServerResponse? { let started = Date() - let response = sendSynchronously { _ in .handleKeyEvent(request) } + let response = sendSynchronously { .handleKeyEvent(request) } let duration = Date().timeIntervalSince(started) if response == nil || duration > 0.2 { // 入力文字・前後文脈は記録しない。 @@ -141,7 +142,7 @@ final class ConverterServerClient { } func sendSynchronously( - _ commandBuilder: @escaping (String) -> ConverterSessionCommand, + _ commandBuilder: @escaping () -> ConverterSessionCommand, onlyIfSessionOpen: Bool = false ) -> ConverterServerResponse? { flushPendingCommands() @@ -180,11 +181,11 @@ final class ConverterServerClient { switch pending.command { case .session(let builder): if let sessionID { - command = .session(sessionID: sessionID, command: builder(sessionID)) + command = .session(sessionID: sessionID, command: builder()) } else { let newID = UUID().uuidString openingSessionID = newID - command = .openSession(sessionID: newID, command: builder(newID)) + command = .openSession(sessionID: newID, command: builder()) } case .global(let global): command = global diff --git a/azooKeyMac/InputController/azooKeyMacInputController.swift b/azooKeyMac/InputController/azooKeyMacInputController.swift index d58bee01..9627c7a9 100644 --- a/azooKeyMac/InputController/azooKeyMacInputController.swift +++ b/azooKeyMac/InputController/azooKeyMacInputController.swift @@ -179,7 +179,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s override func deactivateServer(_ sender: Any!) { self.activationGeneration &+= 1 self.converterSessionState.deactivate() - self.converterServerClient.sendIfSessionOpen({ _ in .lifecycle(.deactivate) }, completion: { _ in }) + self.converterServerClient.sendIfSessionOpen({ .lifecycle(.deactivate) }, completion: { _ in }) self.currentConverterView = nil self.inputState = .none self.candidatesWindow.orderOut(nil) @@ -191,7 +191,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s @MainActor override func commitComposition(_ sender: Any!) { - self.sendAndApply { _ in .composition(.commit) } + self.sendAndApply { .composition(.commit) } } // MARK: - setValue: 状態同期のみ @@ -204,24 +204,16 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s if tag == kTextServiceInputModePropertyTag, let value = value as? NSString, value == "com.apple.inputmethod.Roman" || value == "com.apple.inputmethod.Japanese" { self.client()?.overrideKeyboard(withKeyboardNamed: Config.KeyboardLayout().value.layoutIdentifier) - let englishMode = value == "com.apple.inputmethod.Roman" - - if englishMode { - // 英語モードへの切り替え通知(実際の処理はhandleで行う) - // メニューバーやshortcut経由の切り替えに対応する。 - // composing中でも英数キーMarkedTextを保ったまま英語入力へ移る。 - if self.inputLanguage == .japanese { - self.inputLanguage = .english - self.synchronizeConverterInputLanguage() - self.refreshCandidateWindow() - self.refreshPredictionWindow() - } - } else { - // 日本語モードへの切り替え - if self.inputLanguage == .english { - self.inputLanguage = .japanese - self.synchronizeConverterInputLanguage() - } + let language: InputLanguage = value == "com.apple.inputmethod.Roman" ? .english : .japanese + guard self.inputLanguage != language else { + return + } + // メニューバーや shortcut 経由の切り替えも、最新モードを同期する。 + self.inputLanguage = language + self.synchronizeConverterInputLanguage() + if language == .english { + self.refreshCandidateWindow() + self.refreshPredictionWindow() } } } @@ -238,7 +230,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s return } let language = self.inputLanguage - self.sendAndApply { _ in .lifecycle(.synchronizeInputLanguage(language)) } + self.sendAndApply { .lifecycle(.synchronizeInputLanguage(language)) } } // swiftlint:disable:next cyclomatic_complexity @@ -415,7 +407,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s } @MainActor - private func sendAndApply(_ command: @escaping (String) -> ConverterSessionCommand) { + private func sendAndApply(_ command: @escaping () -> ConverterSessionCommand) { let generation = self.activationGeneration if let response = self.converterServerClient.sendSynchronously(command, onlyIfSessionOpen: true), self.activationGeneration == generation { @@ -505,7 +497,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s private func syncConverterServerSessionConfig() { let config = self.converterServerSessionConfig self.converterServerClient.send( - { _ in .updateConfig(config) }, + { .updateConfig(config) }, completion: { _ in } ) } @@ -525,7 +517,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s private func discardConverterServerComposition() { self.currentConverterView = nil self.converterServerClient.sendIfSessionOpen( - { _ in .composition(.stopComposition) }, + { .composition(.stopComposition) }, completion: { _ in } ) } @@ -595,7 +587,7 @@ class azooKeyMacInputController: IMKInputController, NSMenuItemValidation { // s let count = view.replaceSuggestionCandidates.count let current = view.replaceSuggestionSelectionIndex ?? (offset > 0 ? -1 : 0) let next = (current + offset + count) % count - self.sendAndApply { _ in .replaceSuggestion(.selectReplaceSuggestionCandidate(index: next)) } + self.sendAndApply { .replaceSuggestion(.selectReplaceSuggestionCandidate(index: next)) } } @MainActor private func showReplaceSuggestionError(message: String) { @@ -773,7 +765,7 @@ extension azooKeyMacInputController: CandidatesViewControllerDelegate { guard self.currentConverterView != nil else { return } - self.sendAndApply { _ in .candidate(.submitSelectedCandidate(context: self.currentConverterTextContext())) } + self.sendAndApply { .candidate(.submitSelectedCandidate(context: self.currentConverterTextContext())) } } @MainActor @@ -781,7 +773,7 @@ extension azooKeyMacInputController: CandidatesViewControllerDelegate { guard self.currentConverterView != nil else { return } - self.sendAndApply { _ in .candidate(.selectCandidate(index: row)) } + self.sendAndApply { .candidate(.selectCandidate(index: row)) } } } @@ -835,7 +827,7 @@ extension azooKeyMacInputController: ReplaceSuggestionsViewControllerDelegate { guard self.currentConverterView?.replaceSuggestionSelectionIndex != row else { return } - self.sendAndApply { _ in .replaceSuggestion(.selectReplaceSuggestionCandidate(index: row)) } + self.sendAndApply { .replaceSuggestion(.selectReplaceSuggestionCandidate(index: row)) } } func replaceSuggestionSubmitted() { @@ -861,7 +853,7 @@ extension azooKeyMacInputController { self.syncConverterServerSessionConfig() let activationGeneration = self.activationGeneration self.converterServerClient.sendIfSessionOpen( - { _ in .replaceSuggestion(.request(context: self.currentConverterTextContext())) }, + { .replaceSuggestion(.request(context: self.currentConverterTextContext())) }, completion: { [weak self] response in guard let self, self.activationGeneration == activationGeneration else { return @@ -890,7 +882,7 @@ extension azooKeyMacInputController { } @MainActor func submitSelectedSuggestionCandidate() { - self.sendAndApply { _ in .replaceSuggestion(.submitSelectedReplaceSuggestion) } + self.sendAndApply { .replaceSuggestion(.submitSelectedReplaceSuggestion) } } @MainActor private func finishReplaceSuggestionComposition() { diff --git a/azooKeyMacTests/ThinClientInputPipelineTests.swift b/azooKeyMacTests/ThinClientInputPipelineTests.swift index 7047575a..8220ec24 100644 --- a/azooKeyMacTests/ThinClientInputPipelineTests.swift +++ b/azooKeyMacTests/ThinClientInputPipelineTests.swift @@ -140,11 +140,11 @@ final class ThinClientInputPipelineTests: XCTestCase { defer { service.stop() } let client = makeClient(service) var completions: [String] = [] - client.send({ _ in .lifecycle(.synchronizeInputLanguage(.japanese)) }) { _ in + client.send({ .lifecycle(.synchronizeInputLanguage(.japanese)) }) { _ in completions.append("language") } - let response = client.sendSynchronously { _ in + let response = client.sendSynchronously { XCTAssertEqual(completions, ["language"]) return .composition(.snapshot) } @@ -171,9 +171,9 @@ final class ThinClientInputPipelineTests: XCTestCase { var resets = 0 client.onSessionReset = { resets += 1 } - XCTAssertNil(client.sendSynchronously { _ in .composition(.snapshot) }) + XCTAssertNil(client.sendSynchronously { .composition(.snapshot) }) XCTAssertEqual(resets, 1) - let response = client.sendSynchronously { _ in .composition(.snapshot) } + let response = client.sendSynchronously { .composition(.snapshot) } XCTAssertEqual(response?.inputState, ConverterInputState.none) XCTAssertEqual(service.commands.count, 2) guard case .openSession(let firstID, _) = service.commands[0], @@ -191,11 +191,11 @@ final class ThinClientInputPipelineTests: XCTestCase { defer { service.stop() } let client = makeClient(service) var completions = 0 - client.send({ _ in .composition(.snapshot) }) { response in + client.send({ .composition(.snapshot) }) { response in XCTAssertNil(response) completions += 1 } - client.send({ _ in .composition(.commit) }) { response in + client.send({ .composition(.commit) }) { response in XCTAssertNil(response) completions += 1 } @@ -213,9 +213,9 @@ final class ThinClientInputPipelineTests: XCTestCase { ]) defer { service.stop() } let client = makeClient(service) - client.send({ _ in .composition(.snapshot) }, completion: { _ in }) + client.send({ .composition(.snapshot) }, completion: { _ in }) var deactivated = false - client.sendIfSessionOpen({ _ in .lifecycle(.deactivate) }, completion: { response in + client.sendIfSessionOpen({ .lifecycle(.deactivate) }, completion: { response in deactivated = response != nil }) client.flushPendingCommands() @@ -230,7 +230,7 @@ final class ThinClientInputPipelineTests: XCTestCase { let service = TestConverterService(responses: [.init(data: try encodedResponse())]) defer { service.stop() } var client: ConverterServerClient? = makeClient(service) - XCTAssertNotNil(client?.sendSynchronously { _ in .composition(.snapshot) }) + XCTAssertNotNil(client?.sendSynchronously { .composition(.snapshot) }) guard case .openSession(let id, _) = service.commands.first else { return XCTFail("Expected a session") }