native macOS codings agent orchestrator prowl.onev.cat
Something went wrong. Try again.
Swift
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671import Foundationimport Networkimport Observation
@MainActor@Observablefinal class MirrorSession: Identifiable { enum Status: Equatable { case disconnected, connecting, choosingPane, subscribing, live, takenOver, hostStopped, paneClosed, incompatible
var label: String { switch self { case .disconnected: "Connection lost" case .connecting: "Connecting…" case .choosingPane: "Choose a pane" case .subscribing: "Opening mirror…" case .live: "Live" case .takenOver: "Taken over by another device" case .hostStopped: "Host stopped sharing" case .paneClosed: "Host pane closed" case .incompatible: "Update Host required" } } }
let id = UUID() private(set) var configuration: MirrorSavedConnection private(set) var panes: [MirrorPaneDescriptor] = [] private(set) var pane: MirrorPaneDescriptor? private(set) var status: Status = .disconnected private(set) var text = "" private(set) var revision: UInt64 = 0 private(set) var updatedAt: Date? private(set) var error: String? private(set) var supportsHistory = false private(set) var supportsLaunch = false private(set) var supportsRemoteScroll = false private(set) var isScrolling = false private(set) var scrollError: String? private(set) var completedScroll: ScrollCompletion? private(set) var scrollAtTop: Bool? private(set) var scrollAtBottom: Bool? @ObservationIgnored private var supportsScrollState = false @ObservationIgnored private var pendingScrollState: MirrorMessage.ScrollStatePayload? @ObservationIgnored private var pendingScroll: PendingScroll? @ObservationIgnored private var scrollTimeout: Task<Void, Never>? @ObservationIgnored private var supportsShellSend = false @ObservationIgnored private var supportsAgentInput = false @ObservationIgnored private var pendingCommand: PendingCommand? @ObservationIgnored private var commandTimeout: Task<Void, Never>? @ObservationIgnored private let clock: any Clock<Duration> @ObservationIgnored private let scrollConfirmationTimeout: Duration
private struct PendingCommand { let id: UUID let continuation: CheckedContinuation<MirrorJSON, any Error> } private struct PendingScroll { let id: UUID let direction: MirrorMessage.ScrollDirection let subscriptionID: UUID let baseline: UInt64 } struct ScrollCompletion: Equatable { let id: UUID let direction: MirrorMessage.ScrollDirection } var canScroll: Bool { status == .live && supportsRemoteScroll && subscriptionID != nil && !isScrolling && !showsHistory } func canScroll(_ direction: MirrorMessage.ScrollDirection) -> Bool { canScroll && (direction == .upward ? scrollAtTop != true : scrollAtBottom != true) } private enum CommandFailure: LocalizedError { case unavailable, disconnected, timedOut, shellUnavailable, agentInputUnavailable var errorDescription: String? { switch self { case .unavailable: "Host command service is unavailable or another command is pending." case .disconnected: "Connection lost before Host confirmed the command." case .timedOut: "Host did not confirm the command in time." case .agentInputUnavailable: "Update Host to support interactive Agent input." case .shellUnavailable: "This Host does not support shell submission. Choose an Agent pane." } } } private(set) var historyLines: [String] = [] private(set) var historyOffset = 0 private(set) var historyCapturedAt: Date? private(set) var historyTruncated = false private(set) var isLoadingHistory = false var showsHistory = false /// The reading view writes it on every scroll and reads it only when it appears, so a change /// must not update views. @ObservationIgnored var liveReadingAnchor: MirrorReadingAnchor? var historyReadingOffset: CGFloat = 0 @ObservationIgnored private var historyID: UUID? @ObservationIgnored private var historyBytes = 0 var draft = "" { didSet { if draft != oldValue { draftRevision = UUID() } } } var supportsSubmission: Bool { supportsLaunch } var isComposing = false private(set) var submission: Submission? @ObservationIgnored private var draftRevision = UUID() @ObservationIgnored private var hostRunID: UUID?
struct Submission { let id: UUID let text: String let paneID: UUID let runID: UUID let draftRevision: UUID let configuration: MirrorSavedConnection var outcome: MirrorDeliveryOutcome }
var canSubmit: Bool { status == .live && supportsSubmission && !isComposing && subscriptionID != nil && hostRunID != nil && submission?.outcome.status != .pending && submission?.outcome.status != .unknown && !draft.trimmingCharacters(in: .whitespacesAndNewlines).isEmpty && draft.utf8.count <= MirrorWire.maximumInput }
var submissionHint: String { if let submission, submission.outcome.status != .accepted { return submission.outcome.detail } if !supportsSubmission { return "This Host has not enabled message submission for this pane." } if status != .live { return "Reconnect to the pane before sending." } return "Host checks Agent readiness when you send." } var followsLatest = true var onVerifiedConnection: ((MirrorSavedConnection) -> Void)? @ObservationIgnored private var transport: (any MirrorTransport)? @ObservationIgnored private var subscriptionID: UUID? @ObservationIgnored private var generation = UUID() @ObservationIgnored private var intent: MirrorMessage.Intent = .ifFree @ObservationIgnored private var userDisconnected = false @ObservationIgnored private var supportsRefresh = false @ObservationIgnored private let makeTransport: (MirrorSavedConnection) throws -> any MirrorTransport
init( configuration: MirrorSavedConnection, clock: any Clock<Duration> = ContinuousClock(), scrollConfirmationTimeout: Duration = .seconds(5), makeTransport: @escaping (MirrorSavedConnection) throws -> any MirrorTransport = { configuration in MirrorRemoteConnection(configuration: configuration) } ) { self.clock = clock self.scrollConfirmationTimeout = scrollConfirmationTimeout self.configuration = configuration self.makeTransport = makeTransport }
func connect() { guard transport == nil else { return } userDisconnected = false error = nil status = .connecting supportsLaunch = false supportsShellSend = false supportsAgentInput = false supportsRemoteScroll = false supportsScrollState = false clearScrollState() clearScroll() scrollError = nil completedScroll = nil generation = UUID() let attempt = generation do { let channel = try makeTransport(configuration) transport = channel channel.onEnrolled = { [weak self] enrolled in guard let self, self.generation == attempt else { return } self.configuration = enrolled } channel.onReady = { [weak self] in guard let self, self.generation == attempt else { return } if let verified = self.transport?.verifiedConfiguration { self.configuration = verified } self.send(.list) } channel.onMessage = { [weak self] message in guard let self, self.generation == attempt else { return } self.receive(message) } channel.onClose = { [weak self] reason in guard let self, self.generation == attempt else { return } self.transport = nil self.finishCommand(.failure(CommandFailure.disconnected)) self.subscriptionID = nil self.clearScroll() self.clearScrollState() self.markDeliveryUncertain() self.isLoadingHistory = false if self.status == .live || self.status == .connecting || self.status == .subscribing || self.status == .choosingPane { self.status = .disconnected } self.error = self.error ?? reason } channel.start() } catch { status = .disconnected self.error = error.localizedDescription } }
func select(_ pane: MirrorPaneDescriptor) { guard status == .choosingPane else { return } self.pane = pane intent = pane.busy ? .takeover : .ifFree subscribe() }
func refreshPanes() { guard status == .choosingPane else { return } send(.list) }
func retry(takeover: Bool = false) { guard transport == nil else { return } intent = takeover ? .takeover : .ifFree connect() }
func foreground() { guard !userDisconnected else { return } if status == .disconnected { retry() } else if status == .live, supportsRefresh, let subscriptionID { send(.refresh(.init(subscriptionID: subscriptionID))) querySubmission() } }
func updateConnection(_ configuration: MirrorSavedConnection) { disconnect() if self.configuration.address != configuration.address || self.configuration.port != configuration.port { pane = nil panes = [] text = "" revision = 0 updatedAt = nil liveReadingAnchor = nil historyReadingOffset = 0 historyID = nil historyLines = [] historyBytes = 0 historyOffset = 0 historyCapturedAt = nil showsHistory = false } self.configuration = configuration intent = .ifFree connect() }
func disconnect() { userDisconnected = true finishCommand(.failure(CommandFailure.disconnected)) generation = UUID() let old = transport transport = nil old?.onClose = nil old?.close(nil) subscriptionID = nil clearScroll() clearScrollState() markDeliveryUncertain() isLoadingHistory = false status = .disconnected }
func command(_ command: MirrorCommandRequest.Command) async throws -> MirrorJSON { let liveListing: Bool = { if case .list = command { return status == .live } return false }() guard status == .choosingPane || liveListing, supportsLaunch, transport != nil, pendingCommand == nil else { throw CommandFailure.unavailable } let request = MirrorCommandRequest(requestID: UUID(), request: .init(command: command)) try Task.checkCancellation() return try await withTaskCancellationHandler { try await withCheckedThrowingContinuation { continuation in pendingCommand = PendingCommand(id: request.requestID, continuation: continuation) let clock = clock commandTimeout = Task { [weak self] in do { try await clock.sleep(for: .seconds(30)) } catch { return } guard self?.pendingCommand?.id == request.requestID else { return } self?.finishCommand(.failure(CommandFailure.timedOut)) } send(.command(.init(subscriptionID: subscriptionID, commandRequest: request))) } } onCancel: { Task { @MainActor [weak self] in guard self?.pendingCommand?.id == request.requestID else { return } self?.finishCommand(.failure(CancellationError())) } } }
private func finishCommand(_ result: Result<MirrorJSON, any Error>) { commandTimeout?.cancel() commandTimeout = nil let pending = pendingCommand pendingCommand = nil pending?.continuation.resume(with: result) }
func submitDraft() { guard canSubmit, let pane, let hostRunID, let subscriptionID else { return } let request = Submission( id: UUID(), text: draft, paneID: pane.id, runID: hostRunID, draftRevision: draftRevision, configuration: configuration, outcome: .init(status: .pending, detail: "Host is checking readiness and sending…")) submission = request Task { await routeSubmission(request, lease: subscriptionID) } let clock = clock Task { [weak self] in do { try await clock.sleep(for: .seconds(30)) } catch { return } guard self?.submission?.id == request.id else { return } self?.markDeliveryUncertain() } }
private func routeSubmission(_ request: Submission, lease: UUID) async { do { let listing = try await command(.list(.init())) guard status == .live, subscriptionID == lease, submission?.id == request.id else { return } let input = try MirrorInputRoute.command( listing: listing, paneID: request.paneID, text: request.text) if case .agentsInput = input, !supportsAgentInput { throw CommandFailure.agentInputUnavailable } if case .send = input, !supportsShellSend { throw CommandFailure.shellUnavailable } send( .command( .init( subscriptionID: lease, commandRequest: .init(requestID: request.id, request: .init(command: input))))) } catch { guard submission?.id == request.id else { return } submission?.outcome = .init( status: .rejected, detail: error.localizedDescription + " No input was sent.") } }
private func receiveDispatch(_ response: MirrorCommandResponse) { guard var current = submission, current.id == response.requestID, current.configuration == configuration, current.outcome.status == .pending || current.outcome.status == .unknown else { return } struct Result: Decodable { struct Payload: Decodable { struct Input: Decodable { let bytes: Int let trailing_enter_sent: Bool } let input: Input? } struct Failure: Decodable { let code: String let message: String } let ok: Bool let command: String? let data: Payload? let error: Failure? } guard let result = try? response.response.decode(Result.self) else { markDeliveryUncertain() return } if result.ok, ["send", "agents.input"].contains(result.command ?? ""), let input = result.data?.input, input.bytes == current.text.utf8.count, input.trailing_enter_sent { current.outcome = .init( status: .accepted, detail: "Sent. Agent completion is separate.") if current.draftRevision == draftRevision { draft = "" } } else if let error = result.error, !["REMOTE_COMMAND_UNCONFIRMED", "DISPATCH_FAILED", "SEND_FAILED"].contains(error.code) { current.outcome = .init(status: .rejected, detail: error.message) } else { current.outcome = .init( status: .unknown, detail: "Delivery is unconfirmed. Check Host before sending again.") } submission = current }
func querySubmission() { guard let submission, let subscriptionID, submission.configuration == configuration, submission.runID == hostRunID, submission.outcome.status == .pending || submission.outcome.status == .unknown else { return } send(.commandReceipt(.init(subscriptionID: subscriptionID, commandReceiptID: submission.id))) }
func acknowledgeUnknownSubmission() { guard submission?.outcome.status == .unknown else { return } submission = nil }
private func markDeliveryUncertain() { if submission?.outcome.status == .pending { submission?.outcome = .init( status: .unknown, detail: "Delivery is unconfirmed. Check the receipt or Host output before sending again.") } }
func loadHistory(refresh: Bool = false) { guard status == .live, supportsHistory, !isLoadingHistory, let subscriptionID else { return } if refresh { historyReadingOffset = 0 historyID = nil historyLines = [] historyBytes = 0 historyOffset = 0 historyCapturedAt = nil } guard historyID == nil || historyOffset > 0 else { return } isLoadingHistory = true showsHistory = true send( .history( .init( historyID: historyID, offset: historyID == nil ? nil : historyOffset, subscriptionID: subscriptionID))) }
func scroll(_ direction: MirrorMessage.ScrollDirection) { guard canScroll(direction), let subscriptionID else { return } let request = PendingScroll(id: UUID(), direction: direction, subscriptionID: subscriptionID, baseline: revision) pendingScroll = request isScrolling = true scrollError = nil followsLatest = false let clock = clock let timeout = scrollConfirmationTimeout scrollTimeout = Task { [weak self] in do { try await clock.sleep(for: timeout) } catch { return } guard self?.pendingScroll?.id == request.id else { return } self?.clearScroll() self?.scrollError = String( localized: "Host did not confirm the scroll in time. Check the screen before trying again.") } send(.scroll(.init(requestID: request.id, direction: direction, subscriptionID: subscriptionID))) }
private func clearScroll() { scrollTimeout?.cancel() scrollTimeout = nil pendingScroll = nil isScrolling = false }
private func clearScrollState() { pendingScrollState = nil scrollAtTop = nil scrollAtBottom = nil }
private func receiveScrollState(_ message: MirrorMessage) { guard case .scrollState(let state) = message, state.subscriptionID == subscriptionID, state.sequence > revision else { return } guard supportsScrollState else { invalidMessage() return } pendingScrollState = state }
private func commitScrollState(sequence: UInt64) -> Bool { guard supportsScrollState else { return true } guard let state = pendingScrollState, state.sequence == sequence, state.subscriptionID == subscriptionID else { return false } pendingScrollState = nil scrollAtTop = state.atTop scrollAtBottom = state.atBottom return true }
private func receiveScrollResult(_ message: MirrorMessage) { guard let pendingScroll, message.scrollRequestID == pendingScroll.id, message.subscriptionID == pendingScroll.subscriptionID, subscriptionID == pendingScroll.subscriptionID else { return } guard let sequence = message.sequence, sequence > pendingScroll.baseline, sequence <= revision else { invalidMessage() return } completedScroll = ScrollCompletion(id: pendingScroll.id, direction: pendingScroll.direction) clearScroll() }
private func consumeScrollFailure(_ message: MirrorMessage) -> Bool { guard message.error?.hasPrefix("SCROLL_UNAVAILABLE") == true || message.error?.hasPrefix("SCROLL_BUSY") == true else { return false } guard let pendingScroll, message.subscriptionID == pendingScroll.subscriptionID, message.scrollRequestID == pendingScroll.id else { return true } clearScroll() scrollError = message.error return true }
private func subscribe() { guard let pane else { return } status = .subscribing revision = 0 subscriptionID = nil clearScroll() clearScrollState() send( .subscribe( .init( paneID: pane.id, representation: .text, intent: intent, includeScrollState: supportsScrollState ? true : nil))) }
private func send(_ message: MirrorMessage) { transport?.send(message, closeAfterSending: false) }
private func receive(_ message: MirrorMessage) { switch message.kind { case .panes: guard message.capabilities?.contains("text-v1") == true else { status = .incompatible error = "Update Prowl on this Host to enable mobile text mirrors." transport?.close(nil) return } supportsLaunch = message.capabilities?.contains("launch-profile") == true supportsAgentInput = message.capabilities?.contains("agent-input") == true supportsShellSend = message.capabilities?.contains("shell-send") == true supportsRefresh = message.capabilities?.contains("refresh") == true supportsHistory = message.capabilities?.contains("history") == true supportsRemoteScroll = message.capabilities?.contains("remote-scroll") == true supportsScrollState = message.capabilities?.contains("scroll-state-v1") == true if supportsSubmission { querySubmission() } panes = message.panes ?? [] onVerifiedConnection?(configuration) if let pane { guard panes.contains(where: { $0.id == pane.id }) else { status = .paneClosed transport?.close(nil) return } subscribe() } else { status = .choosingPane } case .commandResult: guard let response = message.commandResponse else { invalidMessage() return } if response.requestID == submission?.id { receiveDispatch(response) return } guard response.requestID == pendingCommand?.id else { return } finishCommand(.success(response.response)) case .subscribed: guard status == .subscribing, message.paneID == pane?.id, message.hostRunID != nil, let lease = message.subscriptionID else { invalidMessage() return } subscriptionID = lease hostRunID = message.hostRunID historyID = nil historyLines = [] historyBytes = 0 historyOffset = 0 historyCapturedAt = nil isLoadingHistory = false showsHistory = false querySubmission() case .textFrame: guard let lease = subscriptionID, message.subscriptionID == lease, let sequence = message.sequence, sequence > revision, let replacement = message.text, replacement.utf8.count <= MirrorWire.maximumPayload / 8, commitScrollState(sequence: sequence) else { invalidMessage() return } text = replacement revision = sequence updatedAt = Date() status = .live send(.acknowledge(.init(sequence: sequence, subscriptionID: lease))) case .scrollResult: receiveScrollResult(message) case .scrollState: receiveScrollState(message) case .historyPage: guard isLoadingHistory, subscriptionID != nil, message.subscriptionID == subscriptionID, let id = message.historyID, let offset = message.offset, let lines = message.lines, let total = message.total, let capturedAt = message.capturedAt, total >= 0, total <= MirrorHistory.maximumBytes + 1, offset >= 0, offset <= total, capturedAt.isFinite, lines.count <= MirrorHistory.pageSize, offset + lines.count == (historyID == nil ? total : historyOffset), historyID == nil || historyID == id else { invalidMessage() return } let pageBytes = lines.reduce(0) { $0 + $1.utf8.count + 1 } guard pageBytes <= MirrorHistory.maximumBytes + 1 - historyBytes else { invalidMessage() return } historyBytes += pageBytes historyID = id historyOffset = offset historyLines.insert(contentsOf: lines, at: 0) historyCapturedAt = Date(timeIntervalSince1970: capturedAt) historyTruncated = message.truncated ?? false isLoadingHistory = false case .ended: guard let reason = message.reason else { invalidMessage() return } switch reason { case .takenOver: status = .takenOver case .hostStopped: status = .hostStopped case .paneClosed: status = .paneClosed } transport?.close(nil) case .failure: if consumeScrollFailure(message) { return } if message.error?.hasPrefix("HISTORY_UNAVAILABLE") == true, message.subscriptionID == subscriptionID { isLoadingHistory = false error = message.error return } if message.error?.hasPrefix("PANE_BUSY") == true { status = .takenOver } error = message.error ?? "Host rejected the request." transport?.close(nil) default: invalidMessage() } }
private func invalidMessage() { error = "Host sent an invalid mirror message." transport?.close(error) }}
nonisolated struct MirrorDeliveryOutcome: Equatable, Sendable { enum Status: String, Sendable { case pending, accepted, rejected, unknown } let status: Status let detail: String}