Merge nucleic/plucky-opal-weasel into dev
This commit is contained in:
@@ -124,6 +124,9 @@ final class HostConnection {
|
|||||||
var castsReceived: ([WireCast]) -> Void = { _ in }
|
var castsReceived: ([WireCast]) -> Void = { _ in }
|
||||||
/// One catch-up batch for one (origin, channel) of the cast feed.
|
/// One catch-up batch for one (origin, channel) of the cast feed.
|
||||||
var castCatchUp: (CastCatchUp) -> Void = { _ in }
|
var castCatchUp: (CastCatchUp) -> Void = { _ in }
|
||||||
|
/// A mesh-dispatch ack from this host (mesh dispatch) — RemoteStore correlates it by
|
||||||
|
/// `requestID` to resume the waiting `dispatchChatToMesh`.
|
||||||
|
var chatStarted: (WireChatStarted) -> Void = { _ in }
|
||||||
}
|
}
|
||||||
private let callbacks: Callbacks
|
private let callbacks: Callbacks
|
||||||
|
|
||||||
@@ -677,10 +680,10 @@ final class HostConnection {
|
|||||||
let result = await PhoneIntelligenceExecutor.execute(request)
|
let result = await PhoneIntelligenceExecutor.execute(request)
|
||||||
self?.send(.intelligenceResult(result))
|
self?.send(.intelligenceResult(result))
|
||||||
}
|
}
|
||||||
case .chatStarted:
|
case .chatStarted(let outcome):
|
||||||
// Mesh-dispatch ack — the phone doesn't drive mesh dispatch yet (the Mac composer
|
// Mesh-dispatch ack — up to RemoteStore's requestID correlation map (the "Auto
|
||||||
// does); inert until the iOS "Auto (Mesh)" destination lands.
|
// (Mesh)" composer destination awaits it).
|
||||||
break
|
callbacks.chatStarted(outcome)
|
||||||
case .casts(let batch):
|
case .casts(let batch):
|
||||||
// Live mesh casts — up to the merged ledger (identity dedup across N connections).
|
// Live mesh casts — up to the merged ledger (identity dedup across N connections).
|
||||||
callbacks.castsReceived(batch.casts)
|
callbacks.castsReceived(batch.casts)
|
||||||
|
|||||||
@@ -80,6 +80,114 @@ final class RemoteStore: ObservableObject {
|
|||||||
UserDefaults.standard.set(feedLastSeenAt, forKey: "nucleic.feed.lastSeenAt")
|
UserDefaults.standard.set(feedLastSeenAt, forKey: "nucleic.feed.lastSeenAt")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// In-flight mesh dispatches awaiting their `chatStarted` ack, keyed by requestID.
|
||||||
|
private var chatStartedWaiters: [String: CheckedContinuation<WireChatStarted?, Never>] = [:]
|
||||||
|
|
||||||
|
/// The latest runner-presence card per host (mesh dispatch), decoded from the merged cast
|
||||||
|
/// ledger's `runner.presence` state channel — the dispatcher's candidate pool.
|
||||||
|
private var runnerPresenceByHost: [String: RunnerPresence] {
|
||||||
|
castLedger.stateValues(channel: MeshCastChannel.runnerPresence)
|
||||||
|
.compactMapValues { RunnerPresence.decode($0) }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Whether the composer should offer "Auto (Mesh)" for the project the phone has selected:
|
||||||
|
/// some connected host advertising `canAcknowledgeDispatch` (other than the project's own
|
||||||
|
/// owner) holds a matching repo. The descriptor comes from the owning host's presence card.
|
||||||
|
func meshDispatchAvailable(forProject projectID: ProjectID) -> Bool {
|
||||||
|
meshDispatchCandidates(forProject: projectID) != nil
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Resolve the (descriptor, candidates) for a mesh dispatch, or nil when it isn't offerable.
|
||||||
|
private func meshDispatchCandidates(
|
||||||
|
forProject projectID: ProjectID
|
||||||
|
) -> (descriptor: ProjectDescriptor, candidates: [MeshDispatchScorer.Candidate])? {
|
||||||
|
let presences = runnerPresenceByHost
|
||||||
|
// The chosen project belongs to one host; find its descriptor from that host's card.
|
||||||
|
guard let owner = connection(owningProject: projectID)?.hostID,
|
||||||
|
let descriptor = presences[owner]?.projects.first(where: { $0.projectID == projectID })
|
||||||
|
else { return nil }
|
||||||
|
var candidates: [MeshDispatchScorer.Candidate] = []
|
||||||
|
for (hostID, presence) in presences {
|
||||||
|
guard let conn = connections[hostID] else { continue }
|
||||||
|
candidates.append(MeshDispatchScorer.Candidate(
|
||||||
|
hostID: hostID, presence: presence,
|
||||||
|
live: conn.connectivity.isLive,
|
||||||
|
canAcknowledgeDispatch: conn.capabilities.canAcknowledgeDispatch))
|
||||||
|
}
|
||||||
|
// Offer only when at least one eligible host actually holds the repo.
|
||||||
|
guard !MeshDispatchScorer.rank(
|
||||||
|
candidates: candidates, project: descriptor, backend: nil).isEmpty
|
||||||
|
else { return nil }
|
||||||
|
return (descriptor, candidates)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Dispatch a new chat to the abstract **Mesh** 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`.
|
||||||
|
func dispatchChatToMesh(
|
||||||
|
projectID: ProjectID, message: String, model: String? = nil,
|
||||||
|
effort: String? = nil, auto: Bool? = nil
|
||||||
|
) {
|
||||||
|
let text = message.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||||
|
guard !text.isEmpty, let (descriptor, candidates) =
|
||||||
|
meshDispatchCandidates(forProject: projectID) else { return }
|
||||||
|
Task { @MainActor in
|
||||||
|
var excluded: Set<String> = []
|
||||||
|
var tried = 0
|
||||||
|
while true {
|
||||||
|
guard let pick = MeshDispatchScorer.rank(
|
||||||
|
candidates: candidates, project: descriptor,
|
||||||
|
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)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
tried += 1
|
||||||
|
// The phone hands over no clone URL, so every pick is a holder (non-nil id).
|
||||||
|
guard let targetProjectID = pick.projectID,
|
||||||
|
let conn = connections[pick.hostID], conn.connectivity.isLive else {
|
||||||
|
excluded.insert(pick.hostID); continue
|
||||||
|
}
|
||||||
|
let requestID = UUID().uuidString
|
||||||
|
let request = StartChatRequest(
|
||||||
|
projectID: targetProjectID, message: text, model: model, effort: effort,
|
||||||
|
auto: auto, requestID: requestID)
|
||||||
|
var outcome = await sendDispatch(request, on: conn, timeout: .seconds(10))
|
||||||
|
if outcome == nil {
|
||||||
|
outcome = await sendDispatch(request, on: conn, timeout: .seconds(5))
|
||||||
|
}
|
||||||
|
if let sessionID = outcome?.sessionID {
|
||||||
|
// The session lives on the target host; refresh its list and route to it.
|
||||||
|
conn.send(.listSessions)
|
||||||
|
conn.send(.listDashboard)
|
||||||
|
route(to: sessionID)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if let failure = outcome?.error {
|
||||||
|
showError("\(conn.hostName): \(failure)", sessionID: nil)
|
||||||
|
}
|
||||||
|
excluded.insert(pick.hostID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private func sendDispatch(
|
||||||
|
_ request: StartChatRequest, on conn: HostConnection, timeout: Duration
|
||||||
|
) async -> WireChatStarted? {
|
||||||
|
guard let requestID = request.requestID else { return nil }
|
||||||
|
conn.send(.startChat(request))
|
||||||
|
return await withCheckedContinuation { continuation in
|
||||||
|
chatStartedWaiters[requestID] = continuation
|
||||||
|
Task { @MainActor in
|
||||||
|
try? await Task.sleep(for: timeout)
|
||||||
|
self.chatStartedWaiters.removeValue(forKey: requestID)?.resume(returning: nil)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Apply freshly received casts (live or catch-up) and persist if anything changed.
|
/// Apply freshly received casts (live or catch-up) and persist if anything changed.
|
||||||
private func applyCasts(_ apply: (inout CastLedger) -> [WireCast]) {
|
private func applyCasts(_ apply: (inout CastLedger) -> [WireCast]) {
|
||||||
var ledger = castLedger
|
var ledger = castLedger
|
||||||
@@ -975,6 +1083,10 @@ final class RemoteStore: ObservableObject {
|
|||||||
cb.castCatchUp = { [weak self] batch in
|
cb.castCatchUp = { [weak self] batch in
|
||||||
self?.applyCasts { $0.applyCatchUp(batch) }
|
self?.applyCasts { $0.applyCatchUp(batch) }
|
||||||
}
|
}
|
||||||
|
cb.chatStarted = { [weak self] outcome in
|
||||||
|
// Mesh dispatch: resume the waiting `dispatchChatToMesh` by requestID.
|
||||||
|
self?.chatStartedWaiters.removeValue(forKey: outcome.requestID)?.resume(returning: outcome)
|
||||||
|
}
|
||||||
cb.meshRosterChanged = { [weak self] in
|
cb.meshRosterChanged = { [weak self] in
|
||||||
// Mesh "join": a Mac was learned or revoked via gossip. Reconnect to every paired Mac
|
// Mesh "join": a Mac was learned or revoked via gossip. Reconnect to every paired Mac
|
||||||
// (connecting the newcomer) and drop any that left — without switching the active host.
|
// (connecting the newcomer) and drop any that left — without switching the active host.
|
||||||
|
|||||||
@@ -18,6 +18,9 @@ struct StartChatComposer: View {
|
|||||||
/// project page's "+", so the chat always lands in that project.
|
/// project page's "+", so the chat always lands in that project.
|
||||||
var lockedProject: WireProject? = nil
|
var lockedProject: WireProject? = nil
|
||||||
@State private var projectID: ProjectID?
|
@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.
|
||||||
|
@State private var runOnMesh = false
|
||||||
@State private var autoOverride: Bool?
|
@State private var autoOverride: Bool?
|
||||||
@State private var model: String?
|
@State private var model: String?
|
||||||
@State private var effort = MobileEfforts.fallback
|
@State private var effort = MobileEfforts.fallback
|
||||||
@@ -39,6 +42,11 @@ struct StartChatComposer: View {
|
|||||||
lockedProject ?? projects.first { $0.id == projectID }
|
lockedProject ?? projects.first { $0.id == projectID }
|
||||||
}
|
}
|
||||||
private var controlled: Bool { selected?.isNucleicControlled ?? false }
|
private var controlled: Bool { selected?.isNucleicControlled ?? false }
|
||||||
|
/// Whether the "Run on: This Mac / Mesh" picker should appear for the current selection.
|
||||||
|
private var meshAvailable: Bool {
|
||||||
|
guard let id = selected?.id else { return false }
|
||||||
|
return store.meshDispatchAvailable(forProject: id)
|
||||||
|
}
|
||||||
/// Auto-approve for the chat about to start: the user's choice off-control, but forced on for a
|
/// Auto-approve for the chat about to start: the user's choice off-control, but forced on for a
|
||||||
/// Nucleic Control project. The host locks Auto on for those chats regardless of what the phone
|
/// Nucleic Control project. The host locks Auto on for those chats regardless of what the phone
|
||||||
/// sends (they run autonomously — nvrsion needs it: see `AppStore.createSession`), so the
|
/// sends (they run autonomously — nvrsion needs it: see `AppStore.createSession`), so the
|
||||||
@@ -79,6 +87,23 @@ struct StartChatComposer: View {
|
|||||||
? "Auto-approve is always on for Nucleic Control chats — they run autonomously."
|
? "Auto-approve is always on for Nucleic Control chats — they run autonomously."
|
||||||
: "Auto-approve safe actions; destructive ones still ask.")
|
: "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.
|
||||||
|
if meshAvailable {
|
||||||
|
Menu {
|
||||||
|
Picker("Run on", selection: $runOnMesh) {
|
||||||
|
Text("This Mac").tag(false)
|
||||||
|
Text("Mesh (automatic)").tag(true)
|
||||||
|
}
|
||||||
|
.pickerStyle(.inline)
|
||||||
|
} label: {
|
||||||
|
Label(
|
||||||
|
runOnMesh ? "Mesh" : "This Mac",
|
||||||
|
systemImage: runOnMesh ? "antenna.radiowaves.left.and.right" : "desktopcomputer")
|
||||||
|
.font(.caption)
|
||||||
|
.foregroundStyle(Palette.accent)
|
||||||
|
}
|
||||||
|
}
|
||||||
HStack {
|
HStack {
|
||||||
ModelMenu(model: $model, catalog: store.modelCatalog, backend: nil)
|
ModelMenu(model: $model, catalog: store.modelCatalog, backend: nil)
|
||||||
Spacer()
|
Spacer()
|
||||||
@@ -132,12 +157,20 @@ struct StartChatComposer: View {
|
|||||||
.focused(focus)
|
.focused(focus)
|
||||||
Button {
|
Button {
|
||||||
if let project = selected {
|
if let project = selected {
|
||||||
|
if runOnMesh && meshAvailable {
|
||||||
|
// Mesh dispatch: route to the best eligible host (attachments can't
|
||||||
|
// ride a dispatch, same as a remote start).
|
||||||
|
store.dispatchChatToMesh(
|
||||||
|
projectID: project.id, message: draft, model: model,
|
||||||
|
effort: effort, auto: effectiveAuto)
|
||||||
|
} else {
|
||||||
let branch = baseBranch.trimmingCharacters(in: .whitespaces)
|
let branch = baseBranch.trimmingCharacters(in: .whitespaces)
|
||||||
store.startChat(
|
store.startChat(
|
||||||
in: project.id, message: draft, model: model, effort: effort,
|
in: project.id, message: draft, model: model, effort: effort,
|
||||||
baseBranch: branch.isEmpty ? nil : branch,
|
baseBranch: branch.isEmpty ? nil : branch,
|
||||||
useWorktree: useWorktree, auto: effectiveAuto,
|
useWorktree: useWorktree, auto: effectiveAuto,
|
||||||
attachments: attachments.wireAttachments)
|
attachments: attachments.wireAttachments)
|
||||||
|
}
|
||||||
draft = ""
|
draft = ""
|
||||||
attachments = []
|
attachments = []
|
||||||
attachmentsOverflowed = false
|
attachmentsOverflowed = false
|
||||||
|
|||||||
Reference in New Issue
Block a user