diff --git a/aiui/packages/app/src/__tests__/providerBridge.test.ts b/aiui/packages/app/src/__tests__/providerBridge.test.ts
new file mode 100644
index 00000000..b8e874a6
--- /dev/null
+++ b/aiui/packages/app/src/__tests__/providerBridge.test.ts
@@ -0,0 +1,28 @@
+import { afterEach, describe, expect, it, vi } from 'vitest'
+import { archyBridge } from '@/services/archyBridge'
+const originalParent = window.parent
+const origin = 'https://node.example'
+afterEach(() => { archyBridge.destroy(); Object.defineProperty(window, 'parent', { value: originalParent, configurable: true }); vi.restoreAllMocks() })
+describe('trusted provider setup bridge', () => {
+ it('accepts configuration only from the embedding parent, rejects siblings and other origins', () => {
+ const parent = { postMessage: vi.fn() }
+ Object.defineProperty(window, 'parent', { value: parent, configurable: true })
+ archyBridge.init(origin)
+ const listener = vi.fn(); const unsubscribe = archyBridge.onProviderConfigured(listener)
+ const send = (source: unknown, from: string, provider = 'openai') => window.dispatchEvent(new MessageEvent('message', { source: source as Window, origin: from, data: { type: 'ai:provider-configured', provider, model: 'test-model' } }))
+ send({}, origin); send(parent, 'https://evil.example'); send(parent, origin, 'arbitrary')
+ expect(listener).not.toHaveBeenCalled()
+ send(parent, origin)
+ expect(listener).toHaveBeenCalledExactlyOnceWith({ provider: 'openai', model: 'test-model' })
+ archyBridge.requestAISetup()
+ expect(parent.postMessage).toHaveBeenLastCalledWith({ type: 'ai:setup-request' }, origin)
+ unsubscribe()
+ })
+ it('replays the selection when the composer mounts after the handshake', () => {
+ const parent = { postMessage: vi.fn() }; Object.defineProperty(window, 'parent', { value: parent, configurable: true })
+ archyBridge.init(origin)
+ window.dispatchEvent(new MessageEvent('message', { source: parent as unknown as Window, origin, data: { type: 'ai:provider-configured', provider: 'local' } }))
+ const listener = vi.fn(); const unsubscribe = archyBridge.onProviderConfigured(listener)
+ expect(listener).toHaveBeenCalledExactlyOnceWith({ provider: 'local', model: '' }); unsubscribe()
+ })
+})
diff --git a/aiui/packages/app/src/__tests__/useAI.test.ts b/aiui/packages/app/src/__tests__/useAI.test.ts
index 330add7a..dcad6757 100644
--- a/aiui/packages/app/src/__tests__/useAI.test.ts
+++ b/aiui/packages/app/src/__tests__/useAI.test.ts
@@ -249,6 +249,24 @@ describe('useAI', () => {
expect(chatStore.isStreaming).toBe(false)
})
+ it.each([502, 503, 429, 401])('distinguishes HTTP %s from a missing credential', async (status) => {
+ globalThis.fetch = vi.fn().mockResolvedValue({ ok: false, status, text: async () => 'Provider request failed' })
+ const store = useChatStore(); store.webSearchEnabled = false
+ const ai = useAI(); ai.needsApiKey.value = false
+ await ai.sendMessage('test')
+ expect(ai.needsApiKey.value).toBe(status === 401)
+ expect(store.isStreaming).toBe(false)
+ })
+
+ it('offers funding for a Routstr payment-required response without mislabeling it a key error', async () => {
+ globalThis.fetch = vi.fn().mockResolvedValue({ ok: false, status: 402, text: async () => JSON.stringify({ error: { message: 'Your spending allowance is exhausted' } }) })
+ const store = useChatStore(); store.webSearchEnabled = false
+ const ai = useAI(); ai.setProvider('routstr'); ai.needsApiKey.value = false; ai.needsFunding.value = false
+ await ai.sendMessage('test')
+ expect(ai.needsFunding.value).toBe(true); expect(ai.needsApiKey.value).toBe(false)
+ expect(store.messages.find(m => m.role === 'assistant')?.content).toContain('spending allowance')
+ })
+
it('handles connection errors gracefully', async () => {
globalThis.fetch = vi.fn().mockRejectedValue(new Error('Network failure'))
diff --git a/aiui/packages/app/src/components/chat/ChatHeader.vue b/aiui/packages/app/src/components/chat/ChatHeader.vue
index d51b3436..7f7c4ef8 100644
--- a/aiui/packages/app/src/components/chat/ChatHeader.vue
+++ b/aiui/packages/app/src/components/chat/ChatHeader.vue
@@ -123,6 +123,7 @@
:style="modelPickerDropdownStyle"
@click.stop
>
+
{{ provider.name }}
@@ -247,6 +248,7 @@ import { useAI } from '@/composables/useAI'
import { useContentPanel } from '@/composables/useContentPanel'
import { downloadConversation, type ExportFormat } from '@/utils/conversation-export'
import { parseImportFile } from '@/utils/conversation-import'
+import { archyBridge } from '@/services/archyBridge'
import { useComparisonMode } from '@/composables/useComparisonMode'
defineProps<{
@@ -332,8 +334,7 @@ const modelDisplayName = computed(() => {
})
function selectModel(providerId: string, modelId: string) {
- setProvider(providerId as 'routstr' | 'claude' | 'openrouter' | 'mock')
- setModel(modelId)
+ if (setProvider(providerId as Parameters[0])) setModel(modelId)
showModelPicker.value = false
}
diff --git a/aiui/packages/app/src/components/chat/ChatWindow.vue b/aiui/packages/app/src/components/chat/ChatWindow.vue
index 0e921682..7b8d2008 100644
--- a/aiui/packages/app/src/components/chat/ChatWindow.vue
+++ b/aiui/packages/app/src/components/chat/ChatWindow.vue
@@ -186,8 +186,9 @@ defineEmits<{
}>()
const chatStore = useChatStore()
-const { sendMessage, stopGeneration, editAndResend, regenerateLastResponse, activeModel, needsApiKey } = useAI()
+const { sendMessage, stopGeneration, editAndResend, regenerateLastResponse, activeModel, needsApiKey, needsFunding } = useAI()
const { updatePanelFromText, panelOpen, panelFilms, panelTitle, activeTab, availableTabs, setActiveTab, enterDesignSystemMode } = useContentPanel()
+import { archyBridge } from '@/services/archyBridge'
import { useCodeContext } from '@/composables/useCodeContext'
import { useVisualViewport } from '@/composables/useVisualViewport'
const codeContext = useCodeContext()
@@ -206,11 +207,19 @@ const showSettings = ref(false)
// without fixing anything).
watch(needsApiKey, (needs) => {
if (needs) {
- showSettings.value = true
+ if (archyBridge.isInArchy()) archyBridge.requestAISetup()
+ else showSettings.value = true
needsApiKey.value = false
}
})
+watch(needsFunding, needed => {
+ if (!needed) return
+ if (archyBridge.isInArchy()) archyBridge.requestAISetup('funding')
+ else showSettings.value = true
+ needsFunding.value = false
+})
+
// Scroll position memory per conversation
const scrollPositions = new Map()
diff --git a/aiui/packages/app/src/components/settings/ApiKeyManager.vue b/aiui/packages/app/src/components/settings/ApiKeyManager.vue
index 17783960..f091163b 100644
--- a/aiui/packages/app/src/components/settings/ApiKeyManager.vue
+++ b/aiui/packages/app/src/components/settings/ApiKeyManager.vue
@@ -1,5 +1,9 @@
-
+
+
AI connections and private keys are managed by this node.
+
+
+
API Keys
@@ -93,10 +97,12 @@
diff --git a/aiui/packages/app/src/composables/useAI.ts b/aiui/packages/app/src/composables/useAI.ts
index a65b922f..f1520a3a 100644
--- a/aiui/packages/app/src/composables/useAI.ts
+++ b/aiui/packages/app/src/composables/useAI.ts
@@ -13,7 +13,7 @@ import { useCodeContext } from '@/composables/useCodeContext'
import { apiFetch } from '@/utils/api-fetch'
import { useSettingsStore } from '@/stores/settings'
-type Provider = 'routstr' | 'claude' | 'openrouter' | 'mock'
+type Provider = 'routstr' | 'claude' | 'openrouter' | 'mock' | 'openai' | 'auto' | 'local'
// API paths are relative to the base URL so they work both in dev (/) and Archy (/aiui/)
const BASE = import.meta.env.BASE_URL || '/'
@@ -120,34 +120,15 @@ Prioritize Podcasting 2.0–friendly platforms: Fountain.fm, Podcast Index, Cast
Always include these tags so the UI can render rich cards. Write a brief reason why each is worth checking out.
${librarySection}`
-const activeProvider = ref
('claude')
+const activeProvider = ref(archyBridge.isInArchy() ? 'auto' : 'claude')
const activeModel = ref('claude-haiku-4.5')
-// One-shot signal a send/regenerate/edit failure looked like a missing or
-// invalid API key (or an unreachable proxy) rather than a transient/server
-// error — consumed by ChatWindow.vue to auto-open Settings so the user isn't
-// left in a dead end with no obvious next step. Deliberately narrow (401/403,
-// explicit "api key"/"unauthorized" text, or a connection-level failure to
-// reach the proxy at all) so a rate-limited or momentarily-flaky provider
-// response does NOT send the user to Settings for a problem Settings can't
-// fix. Reset to false by the consumer immediately after acting on it, so it
-// behaves as a pulse rather than sticky state (each new failure can re-fire).
+// Credentials require setup; network failures and provider outages require retry.
const needsApiKey = ref(false)
-
+const needsFunding = ref(false)
function looksLikeMissingApiKey(err: string): boolean {
- const lower = err.toLowerCase()
- return (
- /\b(401|403)\b/.test(err) ||
- lower.includes('api key') ||
- lower.includes('x-api-key') ||
- lower.includes('unauthorized') ||
- lower.includes('authentication_error') ||
- lower.includes('failed to fetch') ||
- lower.includes('econnrefused') ||
- lower.includes(' 502') ||
- lower.includes(' 503')
- )
+ return /\b(401|403)\b|api[ _-]?key|unauthorized|authentication_error|credential/i.test(err)
}
// ─── Routstr model catalog (fetched from the node's session-gated proxy) ───
@@ -161,7 +142,7 @@ async function refreshRoutstrModels() {
routstrModelsFetched = true
try {
const res = await apiFetch(ROUTSTR_MODELS_PATH)
- if (!res.ok) return
+ if (!res.ok) { routstrModelsFetched = false; return }
const data = await res.json()
if (Array.isArray(data?.data)) {
routstrModels.value = data.data
@@ -170,13 +151,21 @@ async function refreshRoutstrModels() {
id: m.id as string,
name: (m.name as string) || (m.id as string),
}))
- }
+ if (activeProvider.value === 'routstr' && activeModel.value === 'routstr-unavailable' && routstrModels.value[0]) activeModel.value = routstrModels.value[0].id
+ } else { routstrModelsFetched = false }
} catch {
routstrModelsFetched = false // allow a retry on the next send/open
}
}
const availableProviders = computed(() => {
+ if (archyBridge.isInArchy()) return [
+ { id: 'local' as Provider, name: 'Local AI', models: [{ id: 'node', name: 'Node configuration' }] },
+ { id: 'auto' as Provider, name: 'Node AI', models: [{ id: 'node', name: 'Node configuration' }] },
+ { id: 'claude' as Provider, name: 'Claude API', models: [{ id: 'node', name: 'Node configuration' }] },
+ { id: 'openai' as Provider, name: 'OpenAI API', models: [{ id: activeProvider.value === 'openai' ? activeModel.value : 'node', name: activeProvider.value === 'openai' ? activeModel.value : 'Configure model' }] },
+ { id: 'routstr' as Provider, name: 'Routstr (sats)', models: routstrModels.value.length ? routstrModels.value : [{ id: 'routstr-unavailable', name: 'Models unavailable — retry' }] },
+ ]
const providers: { id: Provider; name: string; models: { id: string; name: string }[] }[] = [
{
id: 'routstr',
@@ -187,7 +176,7 @@ const availableProviders = computed(() => {
},
{
id: 'claude',
- name: 'Claude (Max)',
+ name: 'Claude API',
models: [
{ id: 'claude-haiku-4.5', name: 'Claude 4.5 Haiku' },
{ id: 'claude-sonnet-4', name: 'Claude Sonnet 4' },
@@ -207,24 +196,36 @@ const availableProviders = computed(() => {
})
providers.push({
id: 'mock',
- name: 'Local (no API)',
+ name: 'Demo echo',
models: [{ id: 'echo', name: 'Echo (mirror input)' }],
})
return providers
})
function setProvider(provider: Provider) {
+ if (archyBridge.isInArchy()) {
+ if (provider === 'routstr' && activeProvider.value === 'routstr') return true
+ archyBridge.requestAISetup(); return false
+ }
activeProvider.value = provider
const p = availableProviders.value.find((pp) => pp.id === provider)
if (p && p.models.length > 0) {
activeModel.value = p.models[0].id
}
+ return true
}
function setModel(model: string) {
+ if (archyBridge.isInArchy() && activeProvider.value !== 'routstr') { archyBridge.requestAISetup(); return }
activeModel.value = model
}
+archyBridge.onProviderConfigured(({ provider, model }) => {
+ activeProvider.value = provider
+ activeModel.value = model || (provider === 'routstr' ? routstrModels.value[0]?.id || 'routstr-unavailable' : 'node')
+ if (provider === 'routstr') void refreshRoutstrModels()
+})
+
interface ChatMessage {
role: 'user' | 'assistant'
content: string
@@ -448,6 +449,7 @@ async function streamRoutstr(
})
const bodyText = await res.text().catch(() => '')
+ if (res.status === 402) needsFunding.value = true
if (!res.ok) {
// The node's refusals carry a plain-language error.message (budget not
// set, budget spent, wallet can't fund) — surface it verbatim.
@@ -974,5 +976,6 @@ export function useAI() {
setProvider,
setModel,
needsApiKey,
+ needsFunding,
}
}
diff --git a/aiui/packages/app/src/services/archyBridge.ts b/aiui/packages/app/src/services/archyBridge.ts
index c4fb3f5f..3cd2e71e 100644
--- a/aiui/packages/app/src/services/archyBridge.ts
+++ b/aiui/packages/app/src/services/archyBridge.ts
@@ -55,6 +55,10 @@ interface ThemeInfo {
type PermissionsCallback = (categories: AIContextCategory[]) => void
type ThemeCallback = (theme: ThemeInfo) => void
+export interface AIProviderSelection { provider: 'auto' | 'local' | 'claude' | 'openai' | 'routstr'; model: string }
+const providerCallbacks = new Set<(selection: AIProviderSelection) => void>()
+let currentProvider: AIProviderSelection | null = null
+
let requestId = 0
const pendingRequests = new Map void
@@ -80,12 +84,19 @@ function postToParent(msg: unknown) {
function handleMessage(event: MessageEvent) {
// Always validate origin — reject if not configured or mismatched
- if (!allowedOrigin || event.origin !== allowedOrigin) return
+ if (!allowedOrigin || event.origin !== allowedOrigin || event.source !== window.parent) return
const msg = event.data
if (!msg || typeof msg.type !== 'string') return
switch (msg.type) {
+ case 'ai:provider-configured': {
+ if (!['auto', 'local', 'claude', 'openai', 'routstr'].includes(msg.provider)) break
+ const selection = { provider: msg.provider as AIProviderSelection['provider'], model: typeof msg.model === 'string' ? msg.model : '' }
+ currentProvider = selection
+ for (const callback of providerCallbacks) callback(selection)
+ break
+ }
case 'context:response': {
const pending = pendingRequests.get(msg.id)
if (pending) {
@@ -223,11 +234,21 @@ export const archyBridge = {
}
},
+ requestAISetup(reason?: 'funding') { postToParent({ type: 'ai:setup-request', ...(reason ? { reason } : {}) }) },
+
+ onProviderConfigured(callback: (selection: AIProviderSelection) => void) {
+ providerCallbacks.add(callback)
+ if (currentProvider) callback(currentProvider)
+ return () => { providerCallbacks.delete(callback) }
+ },
+
/** Clean up listeners */
destroy() {
window.removeEventListener('message', handleMessage)
pendingRequests.clear()
initialized = false
+ currentProvider = null
+ allowedOrigin = null
},
/** Check if running inside Archy iframe */
diff --git a/core/archipelago/src/api/rpc/analytics.rs b/core/archipelago/src/api/rpc/analytics.rs
index 3532a812..8e6a7e97 100644
--- a/core/archipelago/src/api/rpc/analytics.rs
+++ b/core/archipelago/src/api/rpc/analytics.rs
@@ -352,20 +352,19 @@ impl RpcHandler {
(Some(u), Some(t)) if t > 0 => {
serde_json::json!((u as f64 / t as f64 * 100.0).round())
}
- _ => serde_json::json!(0),
+ _ => serde_json::Value::Null,
}
};
let apps = state.map(|s| s.apps.as_slice()).unwrap_or(&[]);
let reported_at = state
.map(|s| s.timestamp.clone())
- .or_else(|| n.last_seen.clone())
- .unwrap_or_else(|| n.added_at.clone());
+ .or_else(|| n.last_seen.clone());
let mut report = serde_json::json!({
"node_id": n.did,
"node_name": state.and_then(|s| s.node_name.clone()).or_else(|| n.name.clone()),
- "uptime_secs": state.and_then(|s| s.uptime_secs).unwrap_or(0),
- "cpu_pct": state.and_then(|s| s.cpu_usage_percent).map(|v| v.round()).unwrap_or(0.0),
+ "uptime_secs": state.and_then(|s| s.uptime_secs),
+ "cpu_pct": state.and_then(|s| s.cpu_usage_percent).filter(|v| v.is_finite() && (0.0..=100.0).contains(v)).map(|v| v.round()),
"mem_pct": pct(state.and_then(|s| s.mem_used_bytes), state.and_then(|s| s.mem_total_bytes)),
"disk_pct": pct(state.and_then(|s| s.disk_used_bytes), state.and_then(|s| s.disk_total_bytes)),
"container_count": apps.len(),
@@ -561,7 +560,7 @@ fn annotate_fleet_report(report: &mut serde_json::Value) {
let is_online = reported
.map(|dt| {
let age = chrono::Utc::now().signed_duration_since(dt);
- age.num_minutes() < 30
+ age.num_seconds() >= -60 && age.num_seconds() < 1800
})
.unwrap_or(false);
@@ -569,7 +568,9 @@ fn annotate_fleet_report(report: &mut serde_json::Value) {
.map(|dt| {
let age = chrono::Utc::now().signed_duration_since(dt);
let mins = age.num_minutes();
- if mins < 1 {
+ if age.num_seconds() < -60 {
+ "unknown (clock ahead)".to_string()
+ } else if mins < 1 {
"just now".to_string()
} else if mins < 60 {
format!("{}m ago", mins)
diff --git a/core/archipelago/src/api/rpc/federation/handlers.rs b/core/archipelago/src/api/rpc/federation/handlers.rs
index c901086d..98465e90 100644
--- a/core/archipelago/src/api/rpc/federation/handlers.rs
+++ b/core/archipelago/src/api/rpc/federation/handlers.rs
@@ -602,14 +602,39 @@ impl RpcHandler {
None
};
+ // Reuse the minute collector instead of running expensive probes for
+ // every peer. An absent/stalled collector is unknown, never zero load.
+ let now = chrono::Utc::now().timestamp();
+ let latest = self
+ .metrics_store
+ .latest()
+ .await
+ .filter(|sample| (0..=180).contains(&now.saturating_sub(sample.timestamp)));
+ let metrics = latest.as_ref().map(|sample| &sample.system);
+ let uptime = tokio::fs::read_to_string("/proc/uptime")
+ .await
+ .ok()
+ .and_then(|s| s.split_whitespace().next()?.parse::().ok())
+ .filter(|v| v.is_finite() && *v >= 0.0)
+ .map(|v| v as u64);
let state = federation::build_local_state(
apps,
- 0.0,
- 0,
- 0,
- 0,
- 0,
- 0,
+ metrics
+ .map(|m| m.cpu_percent)
+ .filter(|v| v.is_finite() && (0.0..=100.0).contains(v)),
+ metrics
+ .filter(|m| m.mem_total_bytes > 0)
+ .map(|m| m.mem_used_bytes),
+ metrics
+ .filter(|m| m.mem_total_bytes > 0)
+ .map(|m| m.mem_total_bytes),
+ metrics
+ .filter(|m| m.disk_total_bytes > 0)
+ .map(|m| m.disk_used_bytes),
+ metrics
+ .filter(|m| m.disk_total_bytes > 0)
+ .map(|m| m.disk_total_bytes),
+ uptime,
tor_active,
server_name,
nostr_npub,
@@ -1254,76 +1279,120 @@ impl RpcHandler {
);
}
+ let reply = self.prepare_peer_approval_reply(&req).await?;
+ // Persist the operator decision before transport. A relay outage must
+ // not require another approval or lose the already-authorized reply.
+ pending::decide(&self.config.data_dir, id, pending::PendingState::Approved).await?;
+ let delivered = self.deliver_peer_approval_reply(&reply).await?;
+ Ok(serde_json::json!({ "approved": true, "id": id, "delivery_pending": !delivered }))
+ }
+
+ async fn prepare_peer_approval_reply(
+ &self,
+ req: &pending::PendingPeerRequest,
+ ) -> Result {
+ use federation::handshake_delivery::{self, ApprovalReply};
+ if let Some(reply) = handshake_delivery::find(&self.config.data_dir, &req.id).await? {
+ anyhow::ensure!(
+ reply.recipient == req.from_nostr_pubkey && reply.expected_did == req.from_did,
+ "Approval recipient changed"
+ );
+ return Ok(reply);
+ }
let (data, _) = self.state_manager.get_snapshot().await;
let local_did = identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?;
let local_onion = data
.server_info
.tor_address
- .clone()
+ .as_deref()
.ok_or_else(|| anyhow::anyhow!("Tor address not available"))?;
- let local_pubkey = data.server_info.pubkey.clone();
-
- // Generate a one-shot federation invite. The code embeds OUR onion
- // and OUR pubkey, but it leaves this box only inside the NIP-44
- // ciphertext below.
- let identity_dir = self.config.data_dir.join("identity");
- let local_fips_npub = identity::fips_npub(&identity_dir).await.unwrap_or(None);
- // Discovery/connection-request approvals admit the requester as
- // Observer — the invite itself now carries that level, so both
- // sides converge on Observer without post-hoc demotion.
+ let local_fips_npub = identity::fips_npub(&self.config.data_dir.join("identity"))
+ .await
+ .unwrap_or(None);
let invite_code = federation::create_invite(
&self.config.data_dir,
&local_did,
- &local_onion,
- &local_pubkey,
+ local_onion,
+ &data.server_info.pubkey,
local_fips_npub.as_deref(),
TrustLevel::Observer,
)
.await?;
+ handshake_delivery::stage(
+ &self.config.data_dir,
+ ApprovalReply {
+ request_id: req.id.clone(),
+ recipient: req.from_nostr_pubkey.clone(),
+ expected_did: req.from_did.clone(),
+ invite_code,
+ attempts: 0,
+ next_attempt: 0,
+ },
+ )
+ .await
+ }
- // Pre-add the requester to OUR federation list as Observer so that
- // when their `federation.peer-joined` callback arrives over Tor we
- // already trust their pubkey enough to accept the join. Their DID
- // and pubkey come from the request — we'll cross-check the pubkey
- // against the eventual peer-joined signature in the existing
- // verification path (handlers.rs line ~365).
- if !req.from_did.is_empty() {
- // We don't know the requester's onion or ed25519 pubkey yet —
- // they'll send those in the federation.peer-joined callback
- // after they apply our invite. Until then we can't add a real
- // FederatedNode entry. We just store the pending row as
- // Approved so the UI shows progress, and trust the existing
- // peer-joined handler to admit them as Observer when they call.
- //
- // Caveat: peer-joined currently hardcodes TrustLevel::Trusted.
- // We override that below by demoting on success.
- debug!(
- requester_did = %req.from_did,
- "Approval pending — waiting for federation.peer-joined callback over Tor"
- );
- }
-
- // Encrypt + send the invite over NIP-44 to the requester.
- let identity_dir = self.config.data_dir.join("identity");
- nostr_handshake::send_peer_invite(
- &identity_dir,
- &req.from_nostr_pubkey,
- &invite_code,
- &self.config.nostr_relays,
+ async fn deliver_peer_approval_reply(
+ &self,
+ reply: &federation::handshake_delivery::ApprovalReply,
+ ) -> Result {
+ let Some(claimed) = federation::handshake_delivery::claim(
+ &self.config.data_dir,
+ &reply.request_id,
+ chrono::Utc::now().timestamp(),
+ )
+ .await?
+ else {
+ return Ok(false);
+ };
+ let result = nostr_handshake::send_peer_invite(
+ &self.config.data_dir.join("identity"),
+ &claimed.recipient,
+ &claimed.invite_code,
+ &self.handshake_relays().await,
self.config.nostr_tor_proxy.as_deref(),
)
- .await?;
+ .await;
+ if result.is_err() {
+ warn!(request_id = %reply.request_id, "Peer approval delivery deferred; durable retry scheduled");
+ }
+ Ok(result.is_ok())
+ }
- pending::set_state(&self.config.data_dir, id, pending::PendingState::Approved).await?;
- info!(
- id = %id,
- from = %req.from_nostr_pubkey,
- "Approved peer request and shipped invite over NIP-44"
- );
- Ok(serde_json::json!({
- "approved": true,
- "id": id,
- }))
+ /// Recover relay loss and legacy approvals without changing trust or
+ /// resurrecting a node that the operator explicitly removed.
+ pub(in crate::api::rpc) async fn retry_peer_approval_replies(&self) -> Result<()> {
+ let requests = pending::load_pending(&self.config.data_dir).await?;
+ let nodes = federation::load_nodes(&self.config.data_dir).await?;
+ let removed = federation::load_removed_dids(&self.config.data_dir).await?;
+ let cutoff = chrono::Utc::now() - chrono::Duration::days(30);
+ let mut sent = 0;
+ for req in requests {
+ if req.outbound || req.state != pending::PendingState::Approved {
+ federation::handshake_delivery::remove(&self.config.data_dir, &req.id).await?;
+ continue;
+ }
+ let expired = chrono::DateTime::parse_from_rfc3339(&req.received_at)
+ .map(|time| time < cutoff)
+ .unwrap_or(true);
+ if expired
+ || removed.contains(&req.from_did)
+ || nodes.iter().any(|node| node.did == req.from_did)
+ {
+ federation::handshake_delivery::remove(&self.config.data_dir, &req.id).await?;
+ continue;
+ }
+ if req.from_did.is_empty() || sent >= 4 {
+ continue;
+ }
+ let reply = self.prepare_peer_approval_reply(&req).await?;
+ if reply.next_attempt > chrono::Utc::now().timestamp() {
+ continue;
+ }
+ self.deliver_peer_approval_reply(&reply).await?;
+ sent += 1;
+ }
+ Ok(())
}
/// federation.reject-request — drop a pending request and, if requested,
@@ -1353,19 +1422,19 @@ impl RpcHandler {
);
}
+ pending::decide(&self.config.data_dir, id, pending::PendingState::Rejected).await?;
if notify {
let identity_dir = self.config.data_dir.join("identity");
let _ = nostr_handshake::send_peer_reject(
&identity_dir,
&req.from_nostr_pubkey,
reason,
- &self.config.nostr_relays,
+ &self.handshake_relays().await,
self.config.nostr_tor_proxy.as_deref(),
)
.await;
}
- pending::set_state(&self.config.data_dir, id, pending::PendingState::Rejected).await?;
info!(id = %id, from = %req.from_nostr_pubkey, "Rejected peer request");
Ok(serde_json::json!({ "rejected": true, "id": id }))
}
@@ -1410,7 +1479,7 @@ impl RpcHandler {
&identity_dir,
&req.from_nostr_pubkey,
reason,
- &self.config.nostr_relays,
+ &self.handshake_relays().await,
self.config.nostr_tor_proxy.as_deref(),
)
.await
diff --git a/core/archipelago/src/api/rpc/federation/handshake_tests.rs b/core/archipelago/src/api/rpc/federation/handshake_tests.rs
new file mode 100644
index 00000000..c4407d2b
--- /dev/null
+++ b/core/archipelago/src/api/rpc/federation/handshake_tests.rs
@@ -0,0 +1,276 @@
+//! Exercise encrypted replies through a relay configured in the UI only.
+use crate::federation::pending::{self, PendingState};
+use futures_util::{SinkExt, StreamExt};
+use nostr_sdk::prelude::{nip44, Event, Keys};
+use std::sync::Arc;
+use std::time::Duration;
+
+#[tokio::test]
+async fn managed_relay_receives_approval_rejection_and_cancellation() {
+ for (operation, accepted) in [
+ ("approve", true),
+ ("reject", true),
+ ("cancel", true),
+ ("approve", false),
+ ("retry", true),
+ ] {
+ let tmp = tempfile::tempdir().unwrap();
+ let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let relay_url = format!("ws://{}", listener.local_addr().unwrap());
+ let relay = tokio::spawn(async move {
+ let (socket, _) = listener.accept().await.unwrap();
+ let mut ws = tokio_tungstenite::accept_async(socket).await.unwrap();
+ while let Some(Ok(message)) = ws.next().await {
+ if !message.is_text() {
+ continue;
+ }
+ let value: serde_json::Value =
+ serde_json::from_str(message.to_text().unwrap()).unwrap();
+ if value[0] != "EVENT" {
+ continue;
+ }
+ let event: Event = serde_json::from_value(value[1].clone()).unwrap();
+ event.verify().unwrap();
+ ws.send(tokio_tungstenite::tungstenite::Message::Text(
+ serde_json::json!([
+ "OK",
+ event.id.to_hex(),
+ accepted,
+ "blocked: fixture rejection"
+ ])
+ .to_string(),
+ ))
+ .await
+ .unwrap();
+ return event;
+ }
+ panic!("relay closed without a signed event");
+ });
+ let mut config = crate::config::Config::default();
+ config.data_dir = tmp.path().to_path_buf();
+ config.nostr_relays.clear();
+ config.nostr_tor_proxy = None;
+ crate::nostr_relays::save_relays(
+ tmp.path(),
+ &crate::nostr_relays::RelayStore {
+ relays: vec![crate::nostr_relays::RelayConfig {
+ url: relay_url,
+ enabled: true,
+ added_at: chrono::Utc::now().to_rfc3339(),
+ }],
+ },
+ )
+ .await
+ .unwrap();
+ let sender = Keys::parse(&"11".repeat(32)).unwrap();
+ let recipient = Keys::parse(&"22".repeat(32)).unwrap();
+ let identity_dir = tmp.path().join("identity");
+ tokio::fs::create_dir_all(&identity_dir).await.unwrap();
+ tokio::fs::write(identity_dir.join("nostr_secret"), "11".repeat(32))
+ .await
+ .unwrap();
+ tokio::fs::write(identity_dir.join("node_key"), [0x33; 32])
+ .await
+ .unwrap();
+ let state = Arc::new(crate::state::StateManager::new());
+ state
+ .mutate_data(|data| {
+ data.server_info.pubkey = "33".repeat(32);
+ data.server_info.tor_address = Some(format!("{}.onion", "a".repeat(56)));
+ })
+ .await;
+ let handler = crate::api::rpc::RpcHandler::new(
+ config,
+ state,
+ Arc::new(crate::monitoring::MetricsStore::new()),
+ crate::session::SessionStore::new_for_tests(tmp.path().join("sessions.json")),
+ None,
+ None,
+ )
+ .await
+ .unwrap();
+ let row = if operation == "cancel" {
+ pending::insert_outbound(
+ tmp.path(),
+ recipient.public_key().to_hex(),
+ String::new(),
+ crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(),
+ None,
+ None,
+ )
+ .await
+ .unwrap()
+ } else {
+ pending::insert_inbound(
+ tmp.path(),
+ recipient.public_key().to_hex(),
+ String::new(),
+ crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(),
+ None,
+ None,
+ )
+ .await
+ .unwrap()
+ .unwrap()
+ };
+ if operation == "retry" {
+ pending::decide(tmp.path(), &row.id, PendingState::Approved)
+ .await
+ .unwrap();
+ }
+ let params = Some(serde_json::json!({"id": row.id, "notify": true}));
+ let action = async {
+ match operation {
+ "approve" => handler.handle_federation_approve_request(params).await,
+ "reject" => handler.handle_federation_reject_request(params).await,
+ "retry" => handler
+ .retry_peer_approval_replies()
+ .await
+ .map(|_| serde_json::json!({"ok": true})),
+ _ => handler.handle_federation_cancel_request(params).await,
+ }
+ };
+ let outcome = tokio::time::timeout(Duration::from_secs(20), action)
+ .await
+ .unwrap();
+ assert_eq!(outcome.is_ok(), accepted || operation == "approve");
+ let event = tokio::time::timeout(Duration::from_secs(5), relay)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(event.pubkey, sender.public_key());
+ let plaintext =
+ nip44::decrypt(recipient.secret_key(), &event.pubkey, &event.content).unwrap();
+ let message: serde_json::Value = serde_json::from_str(&plaintext).unwrap();
+ let expected = match operation {
+ "approve" | "retry" => "peer-invite",
+ "reject" => "peer-reject",
+ _ => "peer-cancel",
+ };
+ assert_eq!(message["type"], expected);
+ let saved = pending::find_by_id(tmp.path(), &row.id).await.unwrap();
+ if !accepted {
+ assert_eq!(
+ saved.unwrap().state,
+ if operation == "approve" {
+ PendingState::Approved
+ } else {
+ PendingState::Pending
+ }
+ );
+ if operation == "approve" {
+ assert_eq!(outcome.unwrap()["delivery_pending"], true);
+ let durable = crate::federation::handshake_delivery::find(tmp.path(), &row.id)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(durable.recipient, recipient.public_key().to_hex());
+ assert_eq!(durable.attempts, 1);
+ }
+ assert!(crate::federation::load_nodes(tmp.path())
+ .await
+ .unwrap()
+ .is_empty());
+ continue;
+ }
+ match operation {
+ "approve" | "retry" => {
+ assert_eq!(saved.unwrap().state, PendingState::Approved);
+ let first = crate::federation::handshake_delivery::find(tmp.path(), &row.id)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(first.attempts, 1);
+ // Concurrent/background polling honors the persisted backoff.
+ handler.retry_peer_approval_replies().await.unwrap();
+ let second = crate::federation::handshake_delivery::find(tmp.path(), &row.id)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(second.attempts, 1);
+ assert_eq!(first.invite_code, second.invite_code);
+ let invite =
+ crate::federation::parse_invite(message["invite_code"].as_str().unwrap())
+ .unwrap();
+ assert_eq!(invite.trust_level, crate::federation::TrustLevel::Observer);
+ assert!(!event.content.contains(".onion"));
+ }
+ "reject" => assert_eq!(saved.unwrap().state, PendingState::Rejected),
+ _ => assert!(saved.is_none()),
+ }
+ }
+}
+
+#[tokio::test]
+async fn federation_metrics_are_collected_values_or_unknown_never_placeholders() {
+ for age in [None, Some(0), Some(181), Some(-120)] {
+ let tmp = tempfile::tempdir().unwrap();
+ let mut config = crate::config::Config::default();
+ config.data_dir = tmp.path().to_path_buf();
+ let metrics = Arc::new(crate::monitoring::MetricsStore::new());
+ if let Some(age) = age {
+ metrics
+ .push(
+ serde_json::from_value(serde_json::json!({
+ "timestamp": chrono::Utc::now().timestamp() - age,
+ "system": {"cpu_percent": 37.5, "mem_used_bytes": 200,
+ "mem_total_bytes": 800, "disk_used_bytes": 600,
+ "disk_total_bytes": 1000, "net_rx_bytes": 0, "net_tx_bytes": 0,
+ "load_avg_1": 0.0, "load_avg_5": 0.0, "load_avg_15": 0.0},
+ "containers": [], "rpc_latency_ms": 0.0, "ws_connections": 0
+ }))
+ .unwrap(),
+ )
+ .await;
+ }
+ let handler = crate::api::rpc::RpcHandler::new(
+ config,
+ Arc::new(crate::state::StateManager::new()),
+ metrics,
+ crate::session::SessionStore::new_for_tests(tmp.path().join("sessions.json")),
+ None,
+ None,
+ )
+ .await
+ .unwrap();
+ let snapshot = handler.handle_federation_get_state().await.unwrap();
+ if age == Some(0) {
+ assert_eq!(snapshot["cpu_usage_percent"], 37.5);
+ assert_eq!(snapshot["mem_used_bytes"], 200);
+ assert_eq!(snapshot["disk_total_bytes"], 1000);
+ } else {
+ for field in [
+ "cpu_usage_percent",
+ "mem_used_bytes",
+ "mem_total_bytes",
+ "disk_used_bytes",
+ "disk_total_bytes",
+ ] {
+ assert!(
+ snapshot.get(field).is_none_or(|v| v.is_null()),
+ "{age:?}: {field}"
+ );
+ }
+ }
+ let peer = serde_json::from_value(serde_json::json!({
+ "did": "did:key:test", "pubkey": "11".repeat(32), "onion": "test.onion",
+ "trust_level": "trusted", "added_at": chrono::Utc::now().to_rfc3339(),
+ "last_state": snapshot
+ }))
+ .unwrap();
+ crate::federation::save_nodes(tmp.path(), &[peer])
+ .await
+ .unwrap();
+ let fleet = handler.handle_telemetry_fleet_status().await.unwrap();
+ let report = &fleet["nodes"][0];
+ if age == Some(0) {
+ assert_eq!(report["cpu_pct"], 38.0);
+ assert_eq!(report["mem_pct"], 25.0);
+ assert_eq!(report["disk_pct"], 60.0);
+ } else {
+ for field in ["cpu_pct", "mem_pct", "disk_pct"] {
+ assert!(report[field].is_null(), "{age:?}: {field}");
+ }
+ }
+ }
+}
diff --git a/core/archipelago/src/api/rpc/federation/mod.rs b/core/archipelago/src/api/rpc/federation/mod.rs
index fbefc121..2dd40a30 100644
--- a/core/archipelago/src/api/rpc/federation/mod.rs
+++ b/core/archipelago/src/api/rpc/federation/mod.rs
@@ -1,4 +1,6 @@
mod handlers;
+#[cfg(test)]
+mod handshake_tests;
use anyhow::Result;
diff --git a/core/archipelago/src/api/rpc/handshake.rs b/core/archipelago/src/api/rpc/handshake.rs
index 6399f7a1..2373256b 100644
--- a/core/archipelago/src/api/rpc/handshake.rs
+++ b/core/archipelago/src/api/rpc/handshake.rs
@@ -276,6 +276,11 @@ impl RpcHandler {
}
Err(e) => tracing::debug!("background handshake poll failed: {e:#}"),
}
+ if load_discovery_state(&self.config.data_dir).await.enabled {
+ if let Err(error) = self.retry_peer_approval_replies().await {
+ tracing::warn!("Peer approval retry could not complete: {error:#}");
+ }
+ }
}
pub(super) async fn handle_handshake_poll(&self) -> Result {
@@ -356,6 +361,16 @@ impl RpcHandler {
);
continue;
};
+ let scoped_invite = match crate::federation::restrict_discovery_invite(
+ invite_code,
+ &row.from_did,
+ ) {
+ Ok(code) => code,
+ Err(_) => {
+ tracing::warn!("Rejected peer invite with mismatched identity");
+ continue;
+ }
+ };
let row_id = row.id.clone();
let (data, _) = self.state_manager.get_snapshot().await;
let local_did =
@@ -373,7 +388,7 @@ impl RpcHandler {
let local_name = data.server_info.name.clone();
match crate::federation::accept_invite(
&self.config.data_dir,
- invite_code,
+ &scoped_invite,
&local_did,
&local_onion,
&local_pubkey,
diff --git a/core/archipelago/src/api/rpc/system/handlers.rs b/core/archipelago/src/api/rpc/system/handlers.rs
index 9f89721b..c1ec24e6 100644
--- a/core/archipelago/src/api/rpc/system/handlers.rs
+++ b/core/archipelago/src/api/rpc/system/handlers.rs
@@ -1025,10 +1025,46 @@ impl RpcHandler {
.ok_or_else(|| anyhow::anyhow!("Missing key"))?;
match key {
- "claude_api_key_set" => {
- let key_file = self.config.data_dir.join("secrets/claude-api-key");
- let has_key = tokio::fs::metadata(&key_file).await.is_ok();
- Ok(serde_json::json!({ "value": has_key }))
+ "claude_api_key_set" | "openai_api_key_set" => {
+ let provider = if key == "claude_api_key_set" {
+ "claude"
+ } else {
+ "openai"
+ };
+ Ok(
+ serde_json::json!({ "value": crate::settings::model_provider::has_key(&self.config.data_dir, provider).await }),
+ )
+ }
+ "ai_provider" => {
+ let settings =
+ crate::settings::model_provider::ModelProvider::load(&self.config.data_dir)
+ .await?;
+ Ok(serde_json::json!({ "value": settings }))
+ }
+ "ai_provider_status" => {
+ let settings =
+ crate::settings::model_provider::ModelProvider::load(&self.config.data_dir)
+ .await?;
+ let local = tokio::time::timeout(std::time::Duration::from_secs(4), async {
+ let (detected, _) = crate::api::rpc::mesh::assistant::detect_ollama().await;
+ detected
+ && crate::assistant::backends::ollama::model_supports_tools(
+ crate::assistant::backends::ollama::OLLAMA_BASE_URL,
+ crate::assistant::backends::ollama::OLLAMA_DEFAULT_MODEL,
+ )
+ .await
+ });
+ let (claude, openai, local) = tokio::join!(
+ crate::settings::model_provider::has_key(&self.config.data_dir, "claude"),
+ crate::settings::model_provider::has_key(&self.config.data_dir, "openai"),
+ local,
+ );
+ let budget = crate::assistant::AssistantBudget::load(&self.config.data_dir).await;
+ Ok(serde_json::json!({ "value": {
+ "schema": 1, "settings": settings, "claude_configured": claude,
+ "openai_configured": openai, "local_ready": local.ok(),
+ "routstr_remaining_sats": budget.remaining_sats(),
+ }}))
}
_ => Ok(serde_json::json!({ "value": null })),
}
@@ -1210,38 +1246,21 @@ impl RpcHandler {
let value = params.get("value").and_then(|v| v.as_str()).unwrap_or("");
match key {
- "claude_api_key" => {
- let secrets_dir = self.config.data_dir.join("secrets");
- tokio::fs::create_dir_all(&secrets_dir)
- .await
- .context("Failed to create secrets dir")?;
- let key_file = secrets_dir.join("claude-api-key");
-
- if value.is_empty() {
- // Remove key
- tokio::fs::remove_file(&key_file).await.ok();
- info!("Claude API key removed");
+ "claude_api_key" | "openai_api_key" => {
+ let provider = if key == "claude_api_key" {
+ "claude"
} else {
- // Save key
- tokio::fs::write(&key_file, value)
- .await
- .context("Failed to write API key")?;
- #[cfg(unix)]
- {
- use std::os::unix::fs::PermissionsExt;
- std::fs::set_permissions(&key_file, std::fs::Permissions::from_mode(0o600))
- .ok();
- }
- info!("Claude API key saved");
- }
-
- // `secrets/claude-api-key` (above) is deliberately the ONLY
- // Claude key ledger on this node (13-02-PLAN.md). A second
- // copy used to be written alongside it for a standalone,
- // unauthenticated sidecar process on port 3142 — that
- // sidecar and its key copy are retired; the session-gated
- // Rust daemon reads this one file directly.
-
+ "openai"
+ };
+ crate::settings::model_provider::save_key(&self.config.data_dir, provider, value)
+ .await?;
+ info!(provider, "AI provider credential updated");
+ Ok(serde_json::json!({ "saved": true }))
+ }
+ "ai_provider" => {
+ let settings: crate::settings::model_provider::ModelProvider =
+ serde_json::from_str(value).context("Invalid AI provider settings")?;
+ settings.save(&self.config.data_dir).await?;
Ok(serde_json::json!({ "saved": true }))
}
_ => anyhow::bail!("Unknown setting: {}", key),
diff --git a/core/archipelago/src/assistant/backends/mod.rs b/core/archipelago/src/assistant/backends/mod.rs
index d8f24057..72f1245c 100644
--- a/core/archipelago/src/assistant/backends/mod.rs
+++ b/core/archipelago/src/assistant/backends/mod.rs
@@ -11,6 +11,7 @@ use crate::api::rpc::RpcHandler;
pub mod claude;
pub mod ollama;
+pub mod openai;
pub mod routstr;
#[cfg(test)]
pub mod scripted;
@@ -39,6 +40,8 @@ pub trait Backend: Send + Sync {
pub enum BackendId {
Ollama,
Claude,
+ Openai,
+ Unavailable,
/// 13-13: the third D-04 leg. Not currently returned as the "primary"
/// id by `select_backend` (mirroring the existing convention that the
/// returned id names the primary attempt, not necessarily which leg of
@@ -52,11 +55,23 @@ impl std::fmt::Display for BackendId {
match self {
BackendId::Ollama => write!(f, "ollama"),
BackendId::Claude => write!(f, "claude"),
+ BackendId::Openai => write!(f, "openai"),
+ BackendId::Unavailable => write!(f, "unavailable"),
BackendId::Routstr => write!(f, "routstr"),
}
}
}
+struct InvalidProviderSettings;
+#[async_trait]
+impl Backend for InvalidProviderSettings {
+ async fn send(&self, _: &str, _: &[ToolDef], _: &[ChatMessage]) -> Result {
+ anyhow::bail!(
+ "AI connection settings could not be loaded. Review them before sending a message."
+ )
+ }
+}
+
/// D-04's per-call fallback: try `primary`'s `send()`, and on a transport
/// error fall through to `secondary` for that SAME call rather than
/// failing the whole turn — a local model that answers earlier turns and
@@ -115,6 +130,54 @@ fn ollama_is_selectable(detected: bool, tool_capable: bool) -> bool {
/// tools-free degrade.
pub async fn select_backend(handler: &RpcHandler) -> (Box, BackendId) {
let data_dir = handler.data_dir();
+ // An explicit provider is a privacy and billing choice. Never silently
+ // fall through to another provider if its credentials or network fail.
+ match crate::settings::model_provider::ModelProvider::load(data_dir).await {
+ Ok(settings) => match settings.provider {
+ crate::settings::model_provider::Provider::Openai => {
+ return (
+ Box::new(openai::OpenaiBackend::new(
+ data_dir.to_path_buf(),
+ settings.openai_model,
+ )),
+ BackendId::Openai,
+ )
+ }
+ crate::settings::model_provider::Provider::Claude => {
+ return (
+ Box::new(claude::ClaudeBackend::new(data_dir.to_path_buf())),
+ BackendId::Claude,
+ )
+ }
+ crate::settings::model_provider::Provider::Local => {
+ return (
+ Box::new(ollama::OllamaBackend::new(
+ ollama::OLLAMA_BASE_URL.to_string(),
+ ollama::OLLAMA_DEFAULT_MODEL.to_string(),
+ )),
+ BackendId::Ollama,
+ )
+ }
+ crate::settings::model_provider::Provider::Routstr => {
+ let budget = crate::assistant::AssistantBudget::load(data_dir).await;
+ let mints = crate::wallet::ecash::load_accepted_mints(data_dir)
+ .await
+ .map(|m| m.mints)
+ .unwrap_or_default();
+ return (
+ Box::new(routstr::RoutstrBackend::new(
+ data_dir.to_path_buf(),
+ budget.payment_policy(),
+ mints,
+ handler.nostr_tor_proxy(),
+ )),
+ BackendId::Routstr,
+ );
+ }
+ crate::settings::model_provider::Provider::Auto => {}
+ },
+ Err(_) => return (Box::new(InvalidProviderSettings), BackendId::Unavailable),
+ }
let (detected, _models) = crate::api::rpc::mesh::assistant::detect_ollama().await;
let model = ollama::OLLAMA_DEFAULT_MODEL;
let tool_capable = if detected {
diff --git a/core/archipelago/src/assistant/backends/openai.rs b/core/archipelago/src/assistant/backends/openai.rs
new file mode 100644
index 00000000..04ce733a
--- /dev/null
+++ b/core/archipelago/src/assistant/backends/openai.rs
@@ -0,0 +1,279 @@
+//! Explicit OpenAI API selection using the shared tool loop and egress policy.
+//! Keys stay node-side; no redirects, automatic retries, or provider fallback.
+use super::{Backend, BackendTurn};
+use crate::assistant::{
+ egress::{self, EgressVerdict},
+ tools::{ChatMessage, ToolCall, ToolDef},
+};
+use anyhow::{Context, Result};
+use async_trait::async_trait;
+use serde_json::{json, Value};
+use std::{path::PathBuf, time::Duration};
+
+const URL: &str = "https://api.openai.com/v1/chat/completions";
+const RESPONSE_LIMIT: usize = 2 * 1024 * 1024;
+pub struct OpenaiBackend {
+ data_dir: PathBuf,
+ model: String,
+}
+impl OpenaiBackend {
+ pub fn new(data_dir: PathBuf, model: String) -> Self {
+ Self { data_dir, model }
+ }
+ async fn send_at(
+ &self,
+ url: &str,
+ system: &str,
+ tools: &[ToolDef],
+ history: &[ChatMessage],
+ ) -> Result {
+ let key = tokio::fs::read_to_string(self.data_dir.join("secrets/openai-api-key"))
+ .await
+ .map_err(|_| {
+ anyhow::anyhow!("OpenAI API key is not configured. Open AI connection settings.")
+ })?;
+ anyhow::ensure!(!key.trim().is_empty(), "OpenAI API key is not configured");
+ anyhow::ensure!(
+ !self.model.is_empty(),
+ "Choose an OpenAI model in AI connection settings"
+ );
+ let mut messages = vec![json!({"role": "system", "content": system})];
+ messages.extend(history.iter().flat_map(super::routstr::message_to_wire));
+ let mut body = json!({"model": self.model, "messages": messages, "stream": false,
+ "store": false, "max_completion_tokens": 2048, "n": 1});
+ if !tools.is_empty() {
+ body["tools"] = json!(tools.iter().map(|tool| json!({"type": "function", "function": {
+ "name": tool.name, "description": tool.description, "parameters": tool.parameters,
+ }})).collect::>());
+ body["parallel_tool_calls"] = json!(false);
+ }
+ let context = egress::EgressContext::from_turn(
+ history,
+ &tools.iter().map(|tool| tool.name).collect::>(),
+ &self.data_dir.join("secrets"),
+ )
+ .await;
+ match egress::screen_outbound(&body.to_string(), &context) {
+ EgressVerdict::Allow => {}
+ EgressVerdict::Truncate(value) => {
+ body = serde_json::from_str(&value)
+ .context("Could not apply outbound privacy filter")?;
+ }
+ EgressVerdict::BlockFallBackLocal => {
+ crate::assistant::global_counters().note_blocked_egress();
+ anyhow::bail!("This message contains private key or recovery material and was not sent to OpenAI");
+ }
+ }
+ let client = reqwest::Client::builder()
+ .timeout(Duration::from_secs(180))
+ .connect_timeout(Duration::from_secs(15))
+ .redirect(reqwest::redirect::Policy::none())
+ .build()?;
+ let mut response = client.post(url).bearer_auth(key.trim()).json(&body).send().await
+ .map_err(|_| anyhow::anyhow!("OpenAI is temporarily unreachable. Your request was not retried automatically."))?;
+ if !response.status().is_success() {
+ anyhow::bail!("{}", error_message(response.status().as_u16()));
+ }
+ let mut bytes = Vec::new();
+ while let Some(chunk) = response
+ .chunk()
+ .await
+ .context("OpenAI response interrupted")?
+ {
+ anyhow::ensure!(
+ bytes.len().saturating_add(chunk.len()) <= RESPONSE_LIMIT,
+ "OpenAI response exceeded the size limit"
+ );
+ bytes.extend_from_slice(&chunk);
+ }
+ parse_response(
+ &serde_json::from_slice(&bytes).context("OpenAI returned an invalid response")?,
+ )
+ }
+}
+#[async_trait]
+impl Backend for OpenaiBackend {
+ async fn send(
+ &self,
+ system: &str,
+ tools: &[ToolDef],
+ history: &[ChatMessage],
+ ) -> Result {
+ self.send_at(URL, system, tools, history).await
+ }
+}
+fn error_message(status: u16) -> &'static str {
+ match status {
+ 401 | 403 => "OpenAI rejected the API key or project access. Check AI connection settings.",
+ 404 => "This OpenAI model is unavailable for your account. Choose another model in AI connection settings.",
+ 429 => "OpenAI usage or rate limit reached. Check your API billing and retry later.",
+ 500..=599 => "OpenAI is temporarily unavailable. Retry later.",
+ _ => "OpenAI rejected the request. Check the selected model and retry.",
+ }
+}
+fn parse_response(value: &Value) -> Result {
+ let choice = value["choices"]
+ .as_array()
+ .and_then(|items| items.first())
+ .context("OpenAI returned no answer")?;
+ anyhow::ensure!(
+ choice["finish_reason"] != "length",
+ "OpenAI reached the response limit. Try a shorter request."
+ );
+ let message = &choice["message"];
+ if let Some(calls) = message["tool_calls"]
+ .as_array()
+ .filter(|calls| !calls.is_empty())
+ {
+ let mut parsed = Vec::new();
+ for call in calls {
+ anyhow::ensure!(
+ call["type"] == "function",
+ "Unsupported OpenAI tool response"
+ );
+ let id = call["id"]
+ .as_str()
+ .filter(|id| !id.is_empty())
+ .context("Missing OpenAI tool call ID")?;
+ let name = call["function"]["name"]
+ .as_str()
+ .filter(|name| !name.is_empty())
+ .context("Missing OpenAI tool name")?;
+ let arguments: Value = serde_json::from_str(
+ call["function"]["arguments"]
+ .as_str()
+ .context("Invalid OpenAI tool arguments")?,
+ )
+ .context("Invalid OpenAI tool arguments")?;
+ anyhow::ensure!(
+ arguments.is_object()
+ && !parsed.iter().any(|previous: &ToolCall| previous.id == id),
+ "Invalid OpenAI tool call"
+ );
+ parsed.push(ToolCall {
+ id: id.into(),
+ name: name.into(),
+ arguments,
+ });
+ }
+ return Ok(BackendTurn::ToolCalls(parsed));
+ }
+ let text = message["content"]
+ .as_str()
+ .or_else(|| message["refusal"].as_str())
+ .filter(|text| !text.trim().is_empty())
+ .context("OpenAI returned no text; check model compatibility")?;
+ Ok(BackendTurn::Text(text.into()))
+}
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::assistant::tools::{Role, ToolResult};
+ #[test]
+ fn parses_text_and_rejects_incomplete_or_malformed_tool_calls() {
+ assert!(
+ matches!(parse_response(&json!({"choices":[{"message":{"content":"hello"}}]})).unwrap(), BackendTurn::Text(text) if text == "hello")
+ );
+ let valid = json!({"choices":[{"message":{"tool_calls":[{"id":"call_1","type":"function","function":{"name":"status","arguments":"{\"count\":1}"}}]}}]});
+ assert!(
+ matches!(parse_response(&valid).unwrap(), BackendTurn::ToolCalls(calls) if calls[0].arguments["count"] == 1)
+ );
+ for args in ["{", "null", "[]"] {
+ let mut invalid = valid.clone();
+ invalid["choices"][0]["message"]["tool_calls"][0]["function"]["arguments"] =
+ json!(args);
+ assert!(parse_response(&invalid).is_err());
+ }
+ assert!(parse_response(
+ &json!({"choices":[{"finish_reason":"length","message":{"content":"partial"}}]})
+ )
+ .is_err());
+ assert!(parse_response(&json!({"choices":[]})).is_err());
+ }
+ #[test]
+ fn errors_distinguish_credentials_limits_and_outages_without_raw_provider_data() {
+ assert!(error_message(401).contains("API key"));
+ assert!(error_message(429).contains("limit"));
+ assert!(!error_message(503).contains("key"));
+ let wire = super::super::routstr::message_to_wire(&ChatMessage {
+ role: Role::Tool,
+ text: None,
+ tool_calls: vec![],
+ tool_results: vec![ToolResult {
+ call_id: "call_1".into(),
+ content: "result".into(),
+ is_error: false,
+ }],
+ });
+ assert_eq!(wire[0]["tool_call_id"], "call_1");
+ }
+ #[tokio::test]
+ async fn real_http_adapter_sends_private_key_only_in_header_and_never_follows_redirect() {
+ use hyper::{
+ service::{make_service_fn, service_fn},
+ Body, Response, Server,
+ };
+ use std::sync::{Arc, Mutex};
+ let captured = Arc::new(Mutex::new(Vec::new()));
+ let capture = captured.clone();
+ let server = Server::bind(&([127, 0, 0, 1], 0).into()).serve(make_service_fn(move |_| {
+ let capture = capture.clone();
+ async move {
+ Ok::<_, hyper::Error>(service_fn(move |request: hyper::Request| {
+ let capture = capture.clone();
+ async move {
+ let (parts, body) = request.into_parts();
+ let body = hyper::body::to_bytes(body).await?;
+ capture.lock().unwrap().push((
+ parts.headers,
+ serde_json::from_slice::(&body).unwrap(),
+ ));
+ Ok::<_, hyper::Error>(
+ Response::builder()
+ .status(302)
+ .header("Location", "/leak")
+ .body(Body::empty())
+ .unwrap(),
+ )
+ }
+ }))
+ }
+ }));
+ let url = format!("http://{}/v1/chat/completions", server.local_addr());
+ let task = tokio::spawn(server);
+ let dir = tempfile::tempdir().unwrap();
+ crate::settings::model_provider::save_key(dir.path(), "openai", "fixture-private-key")
+ .await
+ .unwrap();
+ let backend = OpenaiBackend::new(dir.path().into(), "test-model".into());
+ let history = [ChatMessage {
+ role: Role::User,
+ text: Some("Hello".into()),
+ tool_calls: vec![],
+ tool_results: vec![],
+ }];
+ assert!(backend
+ .send_at(&url, "Be helpful", &[], &history)
+ .await
+ .is_err());
+ let requests = captured.lock().unwrap();
+ assert_eq!(requests.len(), 1);
+ assert_eq!(requests[0].0["authorization"], "Bearer fixture-private-key");
+ assert_eq!(requests[0].1["store"], false);
+ assert_eq!(requests[0].1["max_completion_tokens"], 2048);
+ assert!(!requests[0].1.to_string().contains("fixture-private-key"));
+ drop(requests);
+ let private = [ChatMessage {
+ role: Role::User,
+ text: Some("fixture-private-key".into()),
+ tool_calls: vec![],
+ tool_results: vec![],
+ }];
+ assert!(backend
+ .send_at(&url, "Be helpful", &[], &private)
+ .await
+ .is_err());
+ assert_eq!(captured.lock().unwrap().len(), 1);
+ task.abort();
+ }
+}
diff --git a/core/archipelago/src/assistant/backends/routstr.rs b/core/archipelago/src/assistant/backends/routstr.rs
index 10331833..f751d6f0 100644
--- a/core/archipelago/src/assistant/backends/routstr.rs
+++ b/core/archipelago/src/assistant/backends/routstr.rs
@@ -295,7 +295,7 @@ fn parse_openai_tool_calls(raw_calls: &[Value]) -> Vec {
/// (the wire-format inverse of `parse_openai_tool_calls`), and tool-result
/// turns carry `tool_call_id` so each call's id is echoed back exactly —
/// the OpenAI-shape contract this adapter's edge is responsible for.
-fn message_to_wire(msg: &ChatMessage) -> Vec {
+pub(super) fn message_to_wire(msg: &ChatMessage) -> Vec {
match msg.role {
Role::System => vec![],
Role::User => vec![json!({
diff --git a/core/archipelago/src/federation/handshake_delivery.rs b/core/archipelago/src/federation/handshake_delivery.rs
new file mode 100644
index 00000000..9be1d309
--- /dev/null
+++ b/core/archipelago/src/federation/handshake_delivery.rs
@@ -0,0 +1,199 @@
+//! Durable, node-encrypted approval replies. Relay acknowledgement is not peer
+//! acceptance: keep retrying the same invite until reciprocal membership exists.
+use anyhow::{Context, Result};
+use serde::{Deserialize, Serialize};
+use std::path::Path;
+use tokio::{fs, io::AsyncWriteExt};
+
+const FILE: &str = "federation/handshake-delivery.enc";
+const DOMAIN: &[u8] = b"archipelago-handshake-delivery-v1";
+static LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
+
+#[derive(Clone, Serialize, Deserialize)]
+pub(crate) struct ApprovalReply {
+ pub request_id: String,
+ pub recipient: String,
+ pub expected_did: String,
+ pub invite_code: String,
+ pub attempts: u32,
+ pub next_attempt: i64,
+}
+
+async fn load(data_dir: &Path) -> Result> {
+ let bytes = match fs::read(data_dir.join(FILE)).await {
+ Ok(bytes) => bytes,
+ Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
+ Err(e) => return Err(e.into()),
+ };
+ let key = crate::storage_crypto::derive_key(data_dir, DOMAIN).await?;
+ let plaintext = crate::storage_crypto::open(&bytes, &key)?;
+ serde_json::from_slice(&plaintext)
+ .context("Invalid handshake delivery store; preserved for recovery")
+}
+
+async fn save(data_dir: &Path, entries: &[ApprovalReply]) -> Result<()> {
+ let key = crate::storage_crypto::derive_key(data_dir, DOMAIN).await?;
+ let bytes = crate::storage_crypto::seal(&serde_json::to_vec(entries)?, &key)?;
+ let path = data_dir.join(FILE);
+ let parent = path.parent().context("Delivery parent missing")?;
+ fs::create_dir_all(parent).await?;
+ let temporary = parent.join(format!(".delivery-{}.tmp", uuid::Uuid::new_v4()));
+ let result = async {
+ let mut file = fs::OpenOptions::new()
+ .create_new(true)
+ .write(true)
+ .mode(0o600)
+ .open(&temporary)
+ .await?;
+ file.write_all(&bytes).await?;
+ file.sync_all().await?;
+ drop(file);
+ fs::rename(&temporary, &path).await?;
+ fs::File::open(parent).await?.sync_all().await?;
+ Ok::<_, anyhow::Error>(())
+ }
+ .await;
+ if result.is_err() {
+ let _ = fs::remove_file(temporary).await;
+ }
+ result
+}
+
+pub(crate) async fn find(data_dir: &Path, request_id: &str) -> Result
- No nodes reporting. Ensure telemetry is enabled on beta nodes.
+ No fleet nodes yet. Connect a trusted node to see its status.
@@ -32,7 +32,7 @@
{{ fleetNodeDisplayName(node) }}
@@ -49,10 +49,10 @@
-
{{ node.cpu_pct.toFixed(0) }}%
+
{{ formatMetric(node.cpu_pct) }}
- {{ node.mem_pct.toFixed(0) }}%
+ {{ formatMetric(node.mem_pct) }}