diff --git a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift index 92821b0..1dbac35 100644 --- a/NucleicRemote/NucleicRemote/Models/RemoteStore.swift +++ b/NucleicRemote/NucleicRemote/Models/RemoteStore.swift @@ -121,10 +121,13 @@ final class RemoteStore: ObservableObject { return (descriptor, candidates) } - /// Dispatch a new chat to the abstract **Mesh** destination from the phone: rank eligible + /// Dispatch a new chat to the abstract **Carbon** destination from the phone: rank eligible /// hosts by the shared scorer and try them in order, one at a time with a correlated, /// idempotent ack (a timeout retries the same candidate once before moving on). On success - /// routes to the landed session. Mirrors `AppStore.dispatchChatToMesh`. + /// routes to the landed session. Mirrors `AppStore.dispatchChatToMesh`. The phone stamps + /// itself as the Carbon origin, so the landed session is Carbon-managed (its owner keeps it + /// on the optimal host between turns); the phone can't own sessions, so the origin exclusion + /// is naturally satisfied by every candidate. func dispatchChatToMesh( projectID: ProjectID, message: String, model: String? = nil, effort: String? = nil, auto: Bool? = nil @@ -141,8 +144,8 @@ final class RemoteStore: ObservableObject { backend: nil, excluding: excluded).first else { showError(tried == 0 - ? "No mesh runner can take this chat right now." - : "No mesh runner could take this chat (\(tried) tried).", sessionID: nil) + ? "No Carbon host can take this chat right now." + : "No Carbon host could take this chat (\(tried) tried).", sessionID: nil) return } tried += 1 @@ -154,7 +157,8 @@ final class RemoteStore: ObservableObject { let requestID = UUID().uuidString let request = StartChatRequest( projectID: targetProjectID, message: text, model: model, effort: effort, - auto: auto, requestID: requestID) + auto: auto, requestID: requestID, + carbonOriginDeviceID: IdentityStore.deviceID()) var outcome = await sendDispatch(request, on: conn, timeout: .seconds(10)) if outcome == nil { outcome = await sendDispatch(request, on: conn, timeout: .seconds(5)) @@ -195,6 +199,7 @@ final class RemoteStore: ObservableObject { guard !fresh.isEmpty else { return } castLedger = ledger activityFeed = ActivityFeedItem.merged(activityFeed, adding: fresh) + updateComposerTyping(fresh) castPersistTask?.cancel() castPersistTask = Task { [ledger] in try? await Task.sleep(for: .milliseconds(500)) @@ -202,6 +207,123 @@ final class RemoteStore: ObservableObject { await CastCache.save(ledger) } } + + // MARK: - Live composer streaming (the "composer.typing" channel) + + /// The live remote typer per session (mesh composer streaming) — folded from + /// `composer.typing` casts, *other* devices only (this phone's own typing, republished by + /// the session's host, is skipped so the typer never locks itself). The session detail + /// locks its composer and renders the streamed draft while an entry is present; + /// expiry-guarded locally so a typer that vanished without its tombstone can't leave the + /// composer locked. + @Published private(set) var composerTypingBySession: [SessionID: ComposerTypingState] = [:] + /// Watcher-side expiry timers for `composerTypingBySession` (the crash guard). + private var composerExpiryTasks: [SessionID: Task] = [:] + /// Sender-side throttle + idle bookkeeping for streaming THIS device's drafts: the + /// trailing-edge send, the last send time, the newest not-yet-sent draft, the idle-stop + /// timer, and which sessions have an active report out (so `ended` sends exactly one stop). + private var composerSendTasks: [SessionID: Task] = [:] + private var composerLastSentAt: [SessionID: Date] = [:] + private var composerPendingText: [SessionID: String] = [:] + private var composerIdleTasks: [SessionID: Task] = [:] + private var composerStreamingSessions: Set = [] + + /// The composer's draft for `sessionID` changed on this phone — report it (throttled) to + /// the session's host, which republishes it on the `composer.typing` channel for the rest + /// of the mesh, and (re)arm the idle stop. An emptied field ends the typing immediately. + /// A no-op toward a host that never advertised `canStreamComposer` (an older host would + /// throw on the unknown tag). + func composerDraftChanged(_ sessionID: SessionID, text: String) { + guard !demoMode else { return } + guard !text.isEmpty else { composerDraftEnded(sessionID); return } + guard let conn = connection(owningSession: sessionID), conn.connectivity.isLive, + conn.capabilities.canStreamComposer else { return } + composerPendingText[sessionID] = text + scheduleComposerIdleStop(sessionID) + guard composerSendTasks[sessionID] == nil else { return } // trailing edge armed + let elapsed = Date().timeIntervalSince(composerLastSentAt[sessionID] ?? .distantPast) + let delay = max(0, ComposerTypingState.minPublishInterval - elapsed) + composerSendTasks[sessionID] = Task { [weak self] in + if delay > 0 { try? await Task.sleep(for: .seconds(delay)) } + guard !Task.isCancelled else { return } + self?.flushComposerDraft(sessionID) + } + } + + /// Typing for `sessionID` finished on this phone — sent, cleared, or the detail closed. + /// Sends the stop that tombstones the session's entry (the host also idle-stops on its own + /// after `idleTimeout`, so a phone that vanishes still unlocks everyone). + func composerDraftEnded(_ sessionID: SessionID) { + composerSendTasks[sessionID]?.cancel() + composerSendTasks[sessionID] = nil + composerPendingText[sessionID] = nil + composerIdleTasks[sessionID]?.cancel() + composerIdleTasks[sessionID] = nil + composerLastSentAt[sessionID] = nil + guard composerStreamingSessions.remove(sessionID) != nil else { return } + connection(owningSession: sessionID)?.send( + .composerTyping(WireComposerTyping(sessionID: sessionID, text: "", active: false))) + } + + /// One throttle-window flush: report the newest pending draft to the session's host. + private func flushComposerDraft(_ sessionID: SessionID) { + composerSendTasks[sessionID] = nil + guard let text = composerPendingText.removeValue(forKey: sessionID), + let conn = connection(owningSession: sessionID), conn.connectivity.isLive + else { return } + composerLastSentAt[sessionID] = Date() + composerStreamingSessions.insert(sessionID) + conn.send(.composerTyping(WireComposerTyping(sessionID: sessionID, text: text, active: true))) + } + + /// Re-arm the typer-side idle stop: `idleTimeout` with no keystrokes ends the typing. + private func scheduleComposerIdleStop(_ sessionID: SessionID) { + composerIdleTasks[sessionID]?.cancel() + composerIdleTasks[sessionID] = Task { [weak self] in + try? await Task.sleep(for: .seconds(ComposerTypingState.idleTimeout)) + guard !Task.isCancelled else { return } + self?.composerDraftEnded(sessionID) + } + } + + /// Fold composer-typing casts into the watcher projection: fresh entries from other + /// devices lock (and render into) that session's composer; tombstones — and local expiry, + /// for a typer that vanished without one — unlock it. + private func updateComposerTyping(_ casts: [WireCast]) { + for cast in casts where cast.channel == MeshCastChannel.composerTyping { + guard let state = ComposerTypingState.decode(cast) else { continue } + if cast.deleted { + guard composerTypingBySession[state.sessionID]?.deviceID == state.deviceID + else { continue } + composerTypingBySession[state.sessionID] = nil + composerExpiryTasks[state.sessionID]?.cancel() + composerExpiryTasks[state.sessionID] = nil + continue + } + // Own typing (republished by the session's host under this phone's deviceID) + // never locks this composer; an entry already past its window (a catch-up + // replay) never locks at all. + guard state.deviceID != IdentityStore.deviceID(), + state.isFresh(at: Date()) else { continue } + composerTypingBySession[state.sessionID] = state + scheduleComposerExpiry(state.sessionID, publishedAt: state.at) + } + } + + /// Arm (or re-arm) the watcher-side expiry for a session's typing entry. + private func scheduleComposerExpiry(_ sessionID: SessionID, publishedAt: Date) { + composerExpiryTasks[sessionID]?.cancel() + let deadline = publishedAt.addingTimeInterval(ComposerTypingState.staleTimeout) + let delay = max(0, deadline.timeIntervalSinceNow) + composerExpiryTasks[sessionID] = Task { [weak self] in + try? await Task.sleep(for: .seconds(delay)) + guard !Task.isCancelled, let self else { return } + guard let entry = self.composerTypingBySession[sessionID], + !entry.isFresh(at: Date()) else { return } + self.composerTypingBySession[sessionID] = nil + self.composerExpiryTasks[sessionID] = nil + } + } private var transcriptPersistTask: Task? @Published private(set) var capabilities = WireCapabilities(canModifyToolInput: false, allowAlwaysScopes: []) @Published private(set) var grantedScope: DeviceScope = .approve @@ -2002,7 +2124,11 @@ final class RemoteStore: ObservableObject { .intelligenceResult, .credentialManifest, .credentialProvision, .createProject, // Cast subscriptions are per-connection (each `HostConnection` subscribes on its // own ready, with the merged ledger's cursors) — nothing routes them here. - .castSubscribe: + .castSubscribe, + // Composer typing goes straight to the owning connection from + // `composerDraftChanged`/`flushComposerDraft` — ephemeral by design, it must + // never be queued for optimistic replay, so it skips this router entirely. + .composerTyping: break } } @@ -2195,7 +2321,10 @@ final class RemoteStore: ObservableObject { // Remote project creation (CLOUD_RUNTIME §4.3) — demo has no host to clone on. .createProject, // Mesh casting (demo has no host to cast). - .castSubscribe: + .castSubscribe, + // Live composer streaming — demo has no other devices to stream to, and + // `composerDraftChanged` already no-ops in demo mode before routing. + .composerTyping: break // passive / already handled by the seeded fixtures (demo has no mesh peers) } } diff --git a/NucleicRemote/NucleicRemote/Views/Composer.swift b/NucleicRemote/NucleicRemote/Views/Composer.swift index ab417d7..bb7e1b3 100644 --- a/NucleicRemote/NucleicRemote/Views/Composer.swift +++ b/NucleicRemote/NucleicRemote/Views/Composer.swift @@ -18,8 +18,10 @@ struct StartChatComposer: View { /// project page's "+", so the chat always lands in that project. var lockedProject: WireProject? = nil @State private var projectID: ProjectID? - /// Mesh dispatch: run on the project's own Mac (false) or auto-route to the best eligible - /// mesh host (true). Shown only when a connected peer could actually take it. + /// Carbon (mesh work queue): run on the project's own Mac (false) or on Carbon (true) — + /// auto-routed to the best eligible mesh host and kept on the optimal host between turns. + /// The phone itself can't run sessions, so every Carbon candidate is already "not here". + /// Shown only when a connected peer could actually take it. @State private var runOnMesh = false @State private var autoOverride: Bool? @State private var model: String? @@ -87,18 +89,19 @@ struct StartChatComposer: View { ? "Auto-approve is always on for Nucleic Control chats — they run autonomously." : "Auto-approve safe actions; destructive ones still ask.") } - // Mesh dispatch: when an eligible peer holds this repo, offer "Auto (Mesh)" so the - // chat routes to whichever host fits rather than always its own owner. + // Carbon (mesh work queue): when an eligible peer holds this repo, offer "Carbon" so + // the chat routes to whichever host fits — and moves between hosts as availability + // changes — rather than always running on its own owner. if meshAvailable { Menu { Picker("Run on", selection: $runOnMesh) { Text("This Mac").tag(false) - Text("Mesh (automatic)").tag(true) + Text("Carbon").tag(true) } .pickerStyle(.inline) } label: { Label( - runOnMesh ? "Mesh" : "This Mac", + runOnMesh ? "Carbon" : "This Mac", systemImage: runOnMesh ? "antenna.radiowaves.left.and.right" : "desktopcomputer") .font(.caption) .foregroundStyle(Palette.accent) diff --git a/NucleicRemote/NucleicRemote/Views/HomeView.swift b/NucleicRemote/NucleicRemote/Views/HomeView.swift index cab269c..afa4d6c 100644 --- a/NucleicRemote/NucleicRemote/Views/HomeView.swift +++ b/NucleicRemote/NucleicRemote/Views/HomeView.swift @@ -95,12 +95,12 @@ struct HomeView: View { .frame(maxWidth: .infinity, alignment: .leading) .card() - if !running.isEmpty { InProgressSessions(sessions: running) } - // Fleet activity (mesh casting): the merged cross-host feed's newest three // rows, with "See all" pushing the full list. Renders offline from the cache. RecentActivitySection() + if !running.isEmpty { InProgressSessions(sessions: running) } + QuickTodos() } .padding() diff --git a/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift b/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift index 81e6278..031f69e 100644 --- a/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift +++ b/NucleicRemote/NucleicRemote/Views/SessionDetailView.swift @@ -40,6 +40,18 @@ struct SessionDetailView: View { store.sessions.first { $0.sessionID == sessionID } } + /// Another device's live typing in this chat's composer (mesh composer streaming). While + /// present, this composer is locked and renders the streamed draft; it unlocks when the + /// typer sends or goes idle (their tombstone — or the local expiry guard — clears it). + private var remoteTyping: ComposerTypingState? { + store.composerTypingBySession[sessionID] + } + + /// The friendly label for a remote typer, never blank. + private func typerName(_ typing: ComposerTypingState) -> String { + typing.deviceName.isEmpty ? "Another device" : typing.deviceName + } + var body: some View { // The detail's height comes from the enclosing GeometryReader — a value fixed by the parent // (the navigation content area), never by anything inside it. The attention cards (approval / @@ -101,9 +113,43 @@ struct SessionDetailView: View { .onChange(of: proxy.size.height, initial: true) { _, height in availableHeight = height } + // Stream this composer's draft to the session's host (throttled in the store) + // so every other device viewing this chat sees it live and locks its own + // composer; an emptied field — including the clear on send — ends the typing + // and unlocks them. Leaving the detail ends it too. + .onChange(of: draft) { _, text in + store.composerDraftChanged(sessionID, text: text) + } + .onDisappear { store.composerDraftEnded(sessionID) } } } + /// The live view of another device's in-progress draft for this chat, shown above the + /// (locked) field. Head-truncated: the tail is where the typing is happening, so it's the + /// part that must stay visible. + private func remoteTypingRow(_ typing: ComposerTypingState) -> some View { + VStack(alignment: .leading, spacing: 4) { + HStack(spacing: 6) { + Image(systemName: "ellipsis.bubble") + Text("\(typerName(typing)) is typing…") + Spacer(minLength: 0) + } + .font(.footnote.weight(.semibold)) + .foregroundStyle(Palette.accent) + if !typing.text.isEmpty { + Text(typing.text) + .font(.footnote) + .foregroundStyle(.secondary) + .lineLimit(3) + .truncationMode(.head) + .frame(maxWidth: .infinity, alignment: .leading) + } + } + .frame(maxWidth: .infinity, alignment: .leading) + .padding(.horizontal, 8).padding(.vertical, 6) + .background(.quaternary.opacity(0.4), in: .rect(cornerRadius: 10)) + } + /// ⌘. interrupts a running session (the Mac's "stop" convention) — the action is otherwise /// only in the ⋯ menu. A hidden button carries the shortcut; present only when it applies. @ViewBuilder @@ -409,6 +455,12 @@ struct SessionDetailView: View { } if store.canControl, let summary { controlRow(summary) } if canCompose { + // Someone is typing in this chat on another device: their draft + // streams in live here while the field below is locked (one typer + // per session at a time, mirroring the Mac composer). + if let typing = remoteTyping { + remoteTypingRow(typing) + } // Staged attachments ride above the field, matching the queued-message // chips and the Mac composer. if !attachments.isEmpty { @@ -430,11 +482,16 @@ struct SessionDetailView: View { // No keyboard-accessory Done button here (it floats awkwardly // over the glass bar on iOS 26) — a drag on the transcript // dismisses the keyboard instead (`scrollDismissesKeyboard`). - TextField(running ? "Queue a follow-up…" : "Send a follow-up…", - text: $draft, axis: .vertical) + TextField( + remoteTyping.map { "\(typerName($0)) is typing…" } + ?? (running ? "Queue a follow-up…" : "Send a follow-up…"), + text: $draft, axis: .vertical) .textFieldStyle(.plain) .lineLimit(1...4) .padding(.vertical, 3) + // Locked while another device is typing here — the + // mesh-wide "one typer per session" contract. + .disabled(remoteTyping != nil) // Stop the in-flight turn (the Mac's ⌘. / "Interrupt"). Shown only // while running and at control scope; send stays at the far right so // its position never shifts. Mirrors `interruptShortcut`. @@ -458,8 +515,10 @@ struct SessionDetailView: View { Image(systemName: "arrow.up.circle.fill") .font(.title2) } - // An attachment-only follow-up (files, no typed text) is sendable. - .disabled((draft.trimmingCharacters(in: .whitespaces).isEmpty && attachments.isEmpty) || !store.connectivity.isLive) + // An attachment-only follow-up (files, no typed text) is sendable — + // but never while another device is mid-draft here (locked). + .disabled((draft.trimmingCharacters(in: .whitespaces).isEmpty && attachments.isEmpty) + || !store.connectivity.isLive || remoteTyping != nil) // Hardware-keyboard send (Magic Keyboard on iPad), mirroring the // Mac — plain Return stays newline in the multiline field. .keyboardShortcut(.return, modifiers: .command)