Compare commits

..
Author SHA1 Message Date
Garry TanandClaude Fable 5 92656a221b fix(admin): regenerate admin-embedded manifest for rebuilt SPA bundle
The Sources-tab rebuild replaced admin/dist/assets/index-CoGEje3-.js with
index-BpDk4NI4.js but src/admin-embedded.ts (generated by
scripts/build-admin-embedded.ts) still imported the deleted file, so
'gbrain serve --http' crashed on startup (Cannot find module) and all 4
admin-embed E2E serial tests failed with 'never became ready'.

Regenerated via: bun run scripts/build-admin-embedded.ts

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 11:01:15 -07:00
6ec762cbcd feat(admin): Sources tab + federation management UI
Takeover of #1601 (stacked on the #1592 takeover), rebased onto current
master. Adds /admin/api/sources (buildSyncStatusReport over the new
queryAdminSources helper — deliberately no local_path filter so push-only
brains still list sources) and four federated-read routes behind
requireAdmin that reuse the same grantReadCore/revokeReadCore/
setFederatedReadCore helpers as the CLI. Admin SPA gains a Sources page
+ per-client manage-reads UI; admin/dist rebuilt from current admin/src.
New test/admin-sources.test.ts pins the sources SQL (archived filter,
null-local_path inclusion, JSONB config shape).

Co-authored-by: bitak1 <bitak1@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:33:21 -07:00
4a81c017a0 feat(auth): grant-read / revoke-read / set-federated-read / list-clients (atomic SQL race-safe)
Takeover of #1592, rebased onto current master. Adds federated-read
management CLI: atomic array_append/array_remove with NOT-ANY +
deleted_at guards (closes the read-modify-write race), source-id
validation at boundaries, terminal-control sanitization for
DCR-registered client names, and 75 PGLite tests.

Co-authored-by: bitak1 <bitak1@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:28:25 -07:00
27 changed files with 1915 additions and 565 deletions
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
+1 -1
View File
@@ -7,7 +7,7 @@
<link rel="preconnect" href="https://fonts.googleapis.com" />
<link rel="preconnect" href="https://fonts.gstatic.com" crossorigin />
<link href="https://fonts.googleapis.com/css2?family=Inter:wght@400;500;600&family=JetBrains+Mono:wght@400;500&display=swap" rel="stylesheet" />
<script type="module" crossorigin src="/admin/assets/index-CoGEje3-.js"></script>
<script type="module" crossorigin src="/admin/assets/index-BpDk4NI4.js"></script>
<link rel="stylesheet" crossorigin href="/admin/assets/index-GxkWX7v3.css">
</head>
<body>
+6 -2
View File
@@ -5,13 +5,14 @@ import { AgentsPage } from './pages/Agents';
import { RequestLogPage } from './pages/RequestLog';
import { CalibrationPage } from './pages/Calibration';
import { JobsWatchPage } from './pages/JobsWatch';
import { SourcesPage } from './pages/Sources';
import { api } from './api';
type Page = 'login' | 'dashboard' | 'agents' | 'log' | 'calibration' | 'jobs';
type Page = 'login' | 'dashboard' | 'agents' | 'sources' | 'log' | 'calibration' | 'jobs';
function getPage(): Page {
const hash = window.location.hash.replace('#', '') || 'dashboard';
if (['login', 'dashboard', 'agents', 'log', 'calibration', 'jobs'].includes(hash)) return hash as Page;
if (['login', 'dashboard', 'agents', 'sources', 'log', 'calibration', 'jobs'].includes(hash)) return hash as Page;
return 'dashboard';
}
@@ -54,6 +55,8 @@ export function App() {
onClick={() => navigate('dashboard')}>Dashboard</a>
<a className={`nav-item ${page === 'agents' ? 'active' : ''}`}
onClick={() => navigate('agents')}>Agents</a>
<a className={`nav-item ${page === 'sources' ? 'active' : ''}`}
onClick={() => navigate('sources')}>Sources</a>
<a className={`nav-item ${page === 'log' ? 'active' : ''}`}
onClick={() => navigate('log')}>Request Log</a>
<a className={`nav-item ${page === 'calibration' ? 'active' : ''}`}
@@ -83,6 +86,7 @@ export function App() {
<main className="main">
{page === 'dashboard' && <DashboardPage />}
{page === 'agents' && <AgentsPage />}
{page === 'sources' && <SourcesPage />}
{page === 'log' && <RequestLogPage />}
{page === 'calibration' && <CalibrationPage />}
{page === 'jobs' && <JobsWatchPage />}
+18
View File
@@ -52,4 +52,22 @@ export const api = {
apiFetchText(`/admin/api/calibration/charts/${encodeURIComponent(type)}${holder ? `?holder=${encodeURIComponent(holder)}` : ''}`),
// v0.41 D2 — live minion-jobs dashboard snapshot.
jobsWatch: () => apiFetch('/admin/api/jobs/watch'),
// v0.41.29 Sources tab + federated-read management
sources: () => apiFetch('/admin/api/sources'),
agentsFederatedRead: () => apiFetch('/admin/api/agents/federated-read'),
grantRead: (clientId: string, sourceId: string) =>
apiFetch(`/admin/api/agents/${encodeURIComponent(clientId)}/grant-read`, {
method: 'POST',
body: JSON.stringify({ source_id: sourceId }),
}),
revokeRead: (clientId: string, sourceId: string) =>
apiFetch(`/admin/api/agents/${encodeURIComponent(clientId)}/revoke-read`, {
method: 'POST',
body: JSON.stringify({ source_id: sourceId }),
}),
setFederatedRead: (clientId: string, sourceIds: string[]) =>
apiFetch(`/admin/api/agents/${encodeURIComponent(clientId)}/set-federated-read`, {
method: 'POST',
body: JSON.stringify({ source_ids: sourceIds }),
}),
};
+196
View File
@@ -381,8 +381,16 @@ function CredentialsModal({ credentials, onClose }: {
);
}
interface FederationState {
source_id: string | null;
federated_read: string[];
}
function AgentDrawer({ agent, onClose, onRevoked }: { agent: Agent; onClose: () => void; onRevoked: () => void }) {
const [tab, setTab] = useState<'claude-code' | 'chatgpt' | 'claude-cowork' | 'perplexity' | 'cursor' | 'json'>('claude-code');
const [federation, setFederation] = useState<FederationState | null>(null);
const [allSources, setAllSources] = useState<string[]>([]);
const [showFederation, setShowFederation] = useState(false);
const copy = (text: string) => navigator.clipboard.writeText(text);
const serverUrl = window.location.origin;
@@ -390,6 +398,30 @@ function AgentDrawer({ agent, onClose, onRevoked }: { agent: Agent; onClose: ()
const isOAuth = agent.auth_type === 'oauth';
const agentName = agent.name || agent.client_name || 'unknown';
// Lazy-load federation state when the drawer opens for an OAuth client.
// The /admin/api/agents endpoint doesn't carry source_id / federated_read,
// so we fetch /admin/api/agents/federated-read separately and pair by id.
useEffect(() => {
if (!isOAuth || !cid) return;
let cancelled = false;
Promise.all([
api.agentsFederatedRead().catch(() => ({ clients: [] })),
api.sources().catch(() => ({ sources: [] })),
]).then(([feds, srcs]: any) => {
if (cancelled) return;
const me = (feds.clients || []).find((c: any) => c.client_id === cid);
setFederation(me ? { source_id: me.source_id, federated_read: me.federated_read || [] } : null);
setAllSources((srcs.sources || []).map((s: any) => s.source_id));
});
return () => { cancelled = true; };
}, [cid, isOAuth]);
const reloadFederation = async () => {
const feds: any = await api.agentsFederatedRead().catch(() => ({ clients: [] }));
const me = (feds.clients || []).find((c: any) => c.client_id === cid);
setFederation(me ? { source_id: me.source_id, federated_read: me.federated_read || [] } : null);
};
// For API keys, we can't show the actual token (it was shown once at creation).
// For OAuth, we show the client_id and tell them to use their secret.
@@ -553,6 +585,34 @@ function AgentDrawer({ agent, onClose, onRevoked }: { agent: Agent; onClose: ()
<span>{agent.token_ttl ? (agent.token_ttl >= 31536000 ? 'No expiry' : agent.token_ttl >= 86400 ? `${Math.floor(agent.token_ttl / 86400)}d` : agent.token_ttl >= 3600 ? `${Math.floor(agent.token_ttl / 3600)}h` : `${agent.token_ttl}s`) : '1h (default)'}</span>
</div>
{isOAuth && federation && (
<>
<div className="section-title" style={{ display: 'flex', alignItems: 'center', justifyContent: 'space-between' }}>
<span>Federation</span>
<button
className="btn btn-secondary"
style={{ padding: '4px 10px', fontSize: 12 }}
onClick={() => setShowFederation(true)}
>
Manage reads
</button>
</div>
<div style={{ display: 'grid', gridTemplateColumns: '120px 1fr', gap: '6px 12px', fontSize: 13 }}>
<span style={{ color: 'var(--text-secondary)' }}>Write source</span>
<span className="mono">{federation.source_id || '(none)'}</span>
<span style={{ color: 'var(--text-secondary)' }}>Federated reads</span>
<span style={{ fontSize: 12 }}>
{federation.federated_read.length === 0
? <span style={{ color: 'var(--text-muted)' }}>(empty no federated reads)</span>
: federation.federated_read.map((s) => (
<span key={s} className="badge badge-read" style={{ marginRight: 4, marginBottom: 2 }}>{s}</span>
))
}
</span>
</div>
</>
)}
{/*
Config Export visible for both auth_type=oauth AND auth_type=api_key.
Claude Code + Cursor + JSON tabs render real snippets regardless
@@ -628,6 +688,142 @@ function AgentDrawer({ agent, onClose, onRevoked }: { agent: Agent; onClose: ()
)}
</div>
</div>
{showFederation && federation && (
<FederationModal
clientId={cid}
clientName={agentName}
allSources={allSources}
currentReads={federation.federated_read}
writeSource={federation.source_id}
onClose={() => setShowFederation(false)}
onSaved={async () => {
await reloadFederation();
setShowFederation(false);
}}
/>
)}
</>
);
}
/**
* FederationModal admin counterpart of `gbrain auth set-federated-read`.
* Source checkbox list; "Save" submits the full new list via the
* race-safe atomic SQL path in setFederatedReadCore. Per the CLI's
* documented contract, this is wholesale-replace semantics concurrent
* grant/revoke from a CLI operator would be last-writer-wins against
* a Save here.
*/
function FederationModal({
clientId, clientName, allSources, currentReads, writeSource, onClose, onSaved,
}: {
clientId: string;
clientName: string;
allSources: string[];
currentReads: string[];
writeSource: string | null;
onClose: () => void;
onSaved: () => Promise<void> | void;
}) {
const [selected, setSelected] = useState<Set<string>>(new Set(currentReads));
const [saving, setSaving] = useState(false);
const [error, setError] = useState<string | null>(null);
const toggle = (id: string) => {
const next = new Set(selected);
if (next.has(id)) next.delete(id); else next.add(id);
setSelected(next);
};
const handleSave = async () => {
setSaving(true);
setError(null);
try {
await api.setFederatedRead(clientId, Array.from(selected));
await onSaved();
} catch (e: any) {
setError(e.message || 'save failed');
setSaving(false);
}
};
// Union: all known sources + any current reads not in the source list
// (e.g. orphan entries from before the source was deleted). The latter
// surface as "(missing source)" so operators can revoke them.
const allKnown = new Set([...allSources, ...currentReads]);
const ordered = Array.from(allKnown).sort();
return (
<div className="modal-overlay" onClick={onClose}>
<div className="modal" onClick={(e) => e.stopPropagation()} style={{ maxWidth: 520 }}>
<div className="modal-header">
<div style={{ fontSize: 16, fontWeight: 600 }}>Manage federated reads</div>
<div style={{ fontSize: 13, color: 'var(--text-secondary)', marginTop: 4 }}>
<strong>{clientName}</strong> pick which sources this client can read in addition to its
{writeSource ? <> write source <code className="mono">{writeSource}</code></> : <> write source</>}.
</div>
</div>
<div className="modal-body" style={{ maxHeight: '50vh', overflowY: 'auto' }}>
{ordered.length === 0 && (
<div style={{ color: 'var(--text-muted)', fontSize: 13 }}>
No sources registered. Use <code>gbrain sources add &lt;id&gt; --path &lt;dir&gt;</code> from the CLI first.
</div>
)}
{ordered.map((id) => {
const isOrphan = !allSources.includes(id);
const isWriteSource = id === writeSource;
return (
<label
key={id}
style={{
display: 'flex',
alignItems: 'center',
gap: 10,
padding: '8px 10px',
borderBottom: '1px solid var(--border)',
cursor: 'pointer',
fontSize: 13,
}}
>
<input
type="checkbox"
checked={selected.has(id)}
onChange={() => toggle(id)}
style={{ width: 16, height: 16, margin: 0, flexShrink: 0, cursor: 'pointer' }}
/>
<span className="mono" style={{ flex: 1, minWidth: 0, overflow: 'hidden', textOverflow: 'ellipsis', whiteSpace: 'nowrap' }}>{id}</span>
{isWriteSource && <span className="badge badge-write" style={{ fontSize: 10, flexShrink: 0 }}>write source</span>}
{isOrphan && <span className="badge badge-danger" style={{ fontSize: 10, flexShrink: 0 }}>missing source</span>}
</label>
);
})}
</div>
{error && (
<div style={{
background: 'rgba(239,68,68,0.08)',
border: '1px solid rgba(239,68,68,0.3)',
color: '#ef4444',
padding: '10px 12px',
borderRadius: 6,
margin: '12px 0',
fontSize: 12,
}}>
{error}
</div>
)}
<div className="modal-footer">
<button type="button" className="btn btn-secondary" onClick={onClose} disabled={saving}>Cancel</button>
<button
type="button"
className="btn btn-primary"
onClick={handleSave}
disabled={saving}
>
{saving ? 'Saving…' : `Save (${selected.size} source${selected.size === 1 ? '' : 's'})`}
</button>
</div>
</div>
</div>
);
}
+195
View File
@@ -0,0 +1,195 @@
import React, { useState, useEffect } from 'react';
import { api } from '../api';
interface SourceRow {
source_id: string;
name: string;
local_path: string | null;
sync_enabled: boolean;
last_sync_at: string | null;
staleness_hours: number | null;
staleness_class: 'fresh' | 'stale' | 'severe' | 'unknown';
last_commit: string | null;
pages: number;
chunks_total: number;
chunks_unembedded: number;
embedding_coverage_pct: number;
}
interface FederatedClient {
client_id: string;
client_name: string;
source_id: string | null;
federated_read: string[];
}
function timeAgo(iso: string | null): string {
if (!iso) return 'never';
const s = Math.floor((Date.now() - new Date(iso).getTime()) / 1000);
if (s < 0) return 'in the future?';
if (s < 60) return 'just now';
if (s < 3600) return `${Math.floor(s / 60)}m ago`;
if (s < 86400) return `${Math.floor(s / 3600)}h ago`;
return `${Math.floor(s / 86400)}d ago`;
}
function stalenessColor(cls: string): string {
switch (cls) {
case 'fresh': return '#4ade80';
case 'stale': return '#fbbf24';
case 'severe': return '#ef4444';
default: return 'var(--text-muted)';
}
}
function coverageColor(pct: number): string {
if (pct >= 99) return '#4ade80';
if (pct >= 90) return '#fbbf24';
return '#ef4444';
}
export function SourcesPage() {
const [sources, setSources] = useState<SourceRow[]>([]);
const [clients, setClients] = useState<FederatedClient[]>([]);
const [loading, setLoading] = useState(true);
const [error, setError] = useState<string | null>(null);
const load = async () => {
setLoading(true);
setError(null);
try {
const [srcReport, clientsResp] = await Promise.all([
api.sources(),
api.agentsFederatedRead(),
]);
setSources(srcReport.sources || []);
setClients(clientsResp.clients || []);
} catch (e: any) {
setError(e.message || 'load failed');
} finally {
setLoading(false);
}
};
useEffect(() => { load(); }, []);
// Reverse-lookup: for each source, which clients can read it?
const readersBySource = (sourceId: string): string[] =>
clients.filter((c) => c.federated_read.includes(sourceId)).map((c) => c.client_name);
// Reverse-lookup: which clients WRITE to this source (source_id == sourceId)?
const writersBySource = (sourceId: string): string[] =>
clients.filter((c) => c.source_id === sourceId).map((c) => c.client_name);
return (
<div style={{ padding: 24, maxWidth: 1200 }}>
<div style={{ display: 'flex', alignItems: 'center', justifyContent: 'space-between', marginBottom: 24 }}>
<h1 style={{ fontSize: 24, margin: 0 }}>Sources</h1>
<button
onClick={load}
style={{
background: 'transparent',
border: '1px solid var(--border)',
color: 'var(--text-secondary)',
padding: '6px 12px',
borderRadius: 6,
fontSize: 12,
cursor: 'pointer',
}}
>
Refresh
</button>
</div>
{loading && <div style={{ color: 'var(--text-muted)' }}>Loading</div>}
{error && (
<div style={{
background: 'rgba(239,68,68,0.08)',
border: '1px solid rgba(239,68,68,0.3)',
color: '#ef4444',
padding: 12,
borderRadius: 6,
marginBottom: 16,
fontSize: 13,
}}>
Failed to load sources: {error}
</div>
)}
{!loading && !error && sources.length === 0 && (
<div style={{ color: 'var(--text-muted)', padding: 16 }}>
No active sources with a local_path. Use{' '}
<code style={{ background: 'var(--bg-elevated)', padding: '2px 6px', borderRadius: 4 }}>
gbrain sources add &lt;id&gt; --path &lt;dir&gt;
</code>{' '}
to register one.
</div>
)}
{!loading && !error && sources.length > 0 && (
<div style={{ overflowX: 'auto' }}>
<table style={{ width: '100%', borderCollapse: 'collapse', fontSize: 13 }}>
<thead>
<tr style={{ borderBottom: '1px solid var(--border)', textAlign: 'left', color: 'var(--text-muted)' }}>
<th style={{ padding: '10px 12px' }}>ID</th>
<th style={{ padding: '10px 12px', textAlign: 'right' }}>Pages</th>
<th style={{ padding: '10px 12px', textAlign: 'right' }}>Chunks</th>
<th style={{ padding: '10px 12px', textAlign: 'right' }}>Embed%</th>
<th style={{ padding: '10px 12px' }}>Last Sync</th>
<th style={{ padding: '10px 12px' }}>Writers</th>
<th style={{ padding: '10px 12px' }}>Readers (federated)</th>
</tr>
</thead>
<tbody>
{sources.map((s) => {
const readers = readersBySource(s.source_id);
const writers = writersBySource(s.source_id);
return (
<tr key={s.source_id} style={{ borderBottom: '1px solid var(--border)' }}>
<td style={{ padding: '10px 12px', fontFamily: 'JetBrains Mono, monospace' }}>
<div>{s.source_id}</div>
{s.name !== s.source_id && (
<div style={{ fontSize: 11, color: 'var(--text-muted)', fontFamily: 'inherit' }}>{s.name}</div>
)}
</td>
<td style={{ padding: '10px 12px', textAlign: 'right', fontFamily: 'JetBrains Mono, monospace' }}>{s.pages.toLocaleString()}</td>
<td style={{ padding: '10px 12px', textAlign: 'right', fontFamily: 'JetBrains Mono, monospace' }}>{s.chunks_total.toLocaleString()}</td>
<td style={{ padding: '10px 12px', textAlign: 'right', color: coverageColor(s.embedding_coverage_pct), fontFamily: 'JetBrains Mono, monospace' }}>
{s.embedding_coverage_pct.toFixed(0)}%
</td>
<td style={{ padding: '10px 12px', color: s.local_path == null ? 'var(--text-muted)' : stalenessColor(s.staleness_class) }}>
{s.local_path == null ? 'push-only' : timeAgo(s.last_sync_at)}
</td>
<td style={{ padding: '10px 12px', fontSize: 12, color: 'var(--text-secondary)' }}>
{writers.length === 0 ? <span style={{ color: 'var(--text-muted)' }}>none</span> : writers.join(', ')}
</td>
<td style={{ padding: '10px 12px', fontSize: 12, color: 'var(--text-secondary)' }}>
{readers.length === 0 ? <span style={{ color: 'var(--text-muted)' }}>none</span> : readers.join(', ')}
</td>
</tr>
);
})}
</tbody>
</table>
</div>
)}
<div style={{
marginTop: 24,
padding: 12,
background: 'var(--bg-elevated)',
border: '1px solid var(--border)',
borderRadius: 6,
fontSize: 12,
color: 'var(--text-muted)',
lineHeight: 1.6,
}}>
<strong style={{ color: 'var(--text-secondary)' }}>Two scopes per OAuth client:</strong>{' '}
<em>Writers</em> = clients with this source as their <code>source_id</code> (write authority).{' '}
<em>Readers</em> = clients with this source in their <code>federated_read</code> list (read access via federation).
Manage federation per-client from the <a href="#agents" style={{ color: '#60a5fa' }}>Agents</a> tab using the
"Manage reads" action.
</div>
</div>
);
}
+3 -3
View File
@@ -1,13 +1,13 @@
// AUTO-GENERATED — do not edit by hand.
// Run `bun run scripts/build-admin-embedded.ts` to regenerate.
// Source: admin/dist/ at 2026-05-27.
// Source: admin/dist/ at 2026-07-22.
//
// Bun resolves the file: imports to a path that works at runtime even
// inside a compiled binary (`bun build --compile`). The manifest maps
// the request path the express handler sees to (resolved-path, mime).
// @ts-ignore — type: 'file' is Bun ESM, not in lib.d.ts
import A_0_assets_index_CoGEje3__js from '../admin/dist/assets/index-CoGEje3-.js' with { type: 'file' };
import A_0_assets_index_BpDk4NI4_js from '../admin/dist/assets/index-BpDk4NI4.js' with { type: 'file' };
// @ts-ignore — type: 'file' is Bun ESM, not in lib.d.ts
import A_1_assets_index_GxkWX7v3_css from '../admin/dist/assets/index-GxkWX7v3.css' with { type: 'file' };
// @ts-ignore — type: 'file' is Bun ESM, not in lib.d.ts
@@ -19,7 +19,7 @@ export interface AdminAsset {
}
export const ADMIN_ASSETS: Record<string, AdminAsset> = {
"/admin/assets/index-CoGEje3-.js": { path: A_0_assets_index_CoGEje3__js as unknown as string, mime: "application/javascript; charset=utf-8" },
"/admin/assets/index-BpDk4NI4.js": { path: A_0_assets_index_BpDk4NI4_js as unknown as string, mime: "application/javascript; charset=utf-8" },
"/admin/assets/index-GxkWX7v3.css": { path: A_1_assets_index_GxkWX7v3_css as unknown as string, mime: "text/css; charset=utf-8" },
"/admin/index.html": { path: A_2_index_html as unknown as string, mime: "text/html; charset=utf-8" },
};
+586 -1
View File
@@ -24,6 +24,8 @@ import { loadConfig, toEngineConfig } from '../core/config.ts';
import { createEngine } from '../core/engine-factory.ts';
import type { BrainEngine } from '../core/engine.ts';
import { sqlQueryForEngine, executeRawJsonb, type SqlQuery } from '../core/sql-query.ts';
import { pgArray } from '../core/oauth-provider.ts';
import { assertValidSourceId } from '../core/source-id.ts';
function hashToken(token: string): string {
return createHash('sha256').update(token).digest('hex');
@@ -165,6 +167,100 @@ async function list() {
});
}
/**
* `gbrain auth list-clients [--json]` read surface for OAuth 2.1 clients.
*
* The existing `gbrain auth list` shows LEGACY bearer tokens from
* `access_tokens`; this is the parallel for v0.26+ OAuth clients. Separate
* commands rather than merged output because the two models have different
* field sets (legacy: lifecycle dates; OAuth: scopes + source_id +
* federated_read).
*
* Human output is card-style (multi-line per client) instead of a fixed-
* width table federated_read can hold many ids per client and a wide
* single-line layout truncates / wraps badly on terminals < 200 cols.
* JSON output uses a `schema_version: 1` envelope; additive only.
*/
async function listClients(args: string[]) {
const json = args.includes('--json');
const includeDeleted = args.includes('--include-deleted');
await withConfiguredSql(async (sql) => {
// Codex finding #2 (medium): default-hide soft-deleted clients so admin
// soft-deletes are honored by the CLI surface. Opt-in via flag.
const rows = includeDeleted
? await sql`
SELECT client_id, client_name, scope, source_id, federated_read,
grant_types, created_at, deleted_at
FROM oauth_clients
ORDER BY client_name
`
: await sql`
SELECT client_id, client_name, scope, source_id, federated_read,
grant_types, created_at, deleted_at
FROM oauth_clients
WHERE deleted_at IS NULL
ORDER BY client_name
`;
if (json) {
const clients = rows.map((r) => ({
client_id: String(r.client_id),
client_name: String(r.client_name),
scope: r.scope == null ? null : String(r.scope),
source_id: r.source_id == null ? null : String(r.source_id),
federated_read: Array.isArray(r.federated_read)
? (r.federated_read as string[]).map(String)
: [],
grant_types: Array.isArray(r.grant_types)
? (r.grant_types as string[]).map(String)
: [],
created_at:
r.created_at instanceof Date
? r.created_at.toISOString()
: r.created_at == null
? null
: String(r.created_at),
deleted_at:
r.deleted_at instanceof Date
? r.deleted_at.toISOString()
: r.deleted_at == null
? null
: String(r.deleted_at),
}));
process.stdout.write(JSON.stringify({ schema_version: 1, clients }, null, 2) + '\n');
return;
}
if (rows.length === 0) {
console.log(
includeDeleted
? 'No OAuth clients found (including deleted). Register one: gbrain auth register-client <name>'
: 'No active OAuth clients found. Register one: gbrain auth register-client <name>'
+ '\n(Use --include-deleted to also show soft-deleted clients.)',
);
return;
}
for (let i = 0; i < rows.length; i++) {
const r = rows[i];
const fed = Array.isArray(r.federated_read)
? (r.federated_read as string[]).map(String)
: [];
const grants = Array.isArray(r.grant_types)
? (r.grant_types as string[]).map(String)
: [];
const deletedAt = r.deleted_at;
const status = deletedAt == null
? ''
: ` [SOFT-DELETED ${deletedAt instanceof Date ? deletedAt.toISOString() : String(deletedAt)}]`;
console.log(`${sanitizeForTerminal(String(r.client_name))}${status}`);
console.log(` client_id: ${sanitizeForTerminal(String(r.client_id))}`);
console.log(` scope: ${r.scope == null ? '(none)' : sanitizeForTerminal(String(r.scope))}`);
console.log(` grant types: ${grants.length ? sanitizeForTerminal(grants.join(', ')) : '(none)'}`);
console.log(` write source: ${r.source_id == null ? '(none)' : sanitizeForTerminal(String(r.source_id))}`);
console.log(` federated: ${fed.length ? sanitizeForTerminal(fed.join(', ')) : '(empty)'}`);
if (i < rows.length - 1) console.log('');
}
});
}
async function revoke(name: string) {
if (!name) { console.error('Usage: auth revoke <name>'); process.exit(1); }
await withConfiguredSql(async (sql) => {
@@ -301,6 +397,475 @@ async function test(url: string, token: string) {
console.log(`\n🧠 Your brain is live! (${elapsed}s)`);
}
/**
* Strip ANSI escapes + C0/C1 control characters from a string before
* printing it to the operator's terminal. Defense for the
* codex-flagged terminal-control-injection class: a client_name or
* source_id registered via DCR with `\x1b[2J` (clear-screen) or
* `\x1b]0;TITLE\x07` (OSC title-change) would poison
* `gbrain auth list-clients` output otherwise.
*
* Replaces unsafe bytes with their `\xNN` hex escape so the operator
* sees that something weird is in the field, instead of silent
* mutilation. Tab and newline are preserved as-is so legitimate
* multi-line values render.
*/
export function sanitizeForTerminal(s: string): string {
// ALL C0/C1 controls + DEL get escaped. Codex re-review caught that
// preserving `\n` lets a DCR-registered client_name spoof additional
// human-output lines in list-clients (a real attack — newline in the
// name visually adds a fake row to the operator's terminal). Tab is
// also escaped for the same reason — field-separator spoofing.
// C0: 0x00-0x1F. DEL: 0x7F. C1: 0x80-0x9F.
return s.replace(/[\x00-\x1f\x7f-\x9f]/g, (ch) =>
`\\x${ch.charCodeAt(0).toString(16).padStart(2, '0')}`,
);
}
export interface ResolvedClient {
client_id: string;
client_name: string;
source_id: string | null;
federated_read: string[];
deleted_at: Date | string | null;
}
export type FederatedReadOutcome =
| { kind: 'noop'; reason: 'already-granted' | 'not-present' | 'same-list'; client: ResolvedClient; current: string[] }
| { kind: 'updated'; client: ResolvedClient; before: string[]; after: string[] };
/**
* Resolve an OAuth client by client_id (exact) or client_name (unique).
* Errors on no-match and on ambiguous client_name (>1 row). client_id
* takes precedence if a long hash is passed and matches, returns
* immediately without ever querying by name.
*
* Legacy bearer tokens in `access_tokens` are NOT searched. Federated read
* scope is an OAuth-client concept (oauth_clients.federated_read column);
* legacy bearers have no source scope.
*/
/**
* Resolve an OAuth client. Codex finding #2 (medium): default-hide
* soft-deleted clients so admin-soft-deleted rows aren't mutated by the
* CLI. The `includeDeleted` opt is reserved for future read-side surfaces;
* grant/revoke/set ALWAYS filter active rows only.
*/
export async function resolveClient(
sql: SqlQuery,
nameOrId: string,
opts: { includeDeleted?: boolean } = {},
): Promise<ResolvedClient> {
const allowDeleted = opts.includeDeleted === true;
const byId = allowDeleted
? await sql`
SELECT client_id, client_name, source_id, federated_read, deleted_at
FROM oauth_clients WHERE client_id = ${nameOrId} LIMIT 1
`
: await sql`
SELECT client_id, client_name, source_id, federated_read, deleted_at
FROM oauth_clients WHERE client_id = ${nameOrId} AND deleted_at IS NULL LIMIT 1
`;
if (byId.length === 1) return normalizeClientRow(byId[0]);
const byName = allowDeleted
? await sql`
SELECT client_id, client_name, source_id, federated_read, deleted_at
FROM oauth_clients WHERE client_name = ${nameOrId}
`
: await sql`
SELECT client_id, client_name, source_id, federated_read, deleted_at
FROM oauth_clients WHERE client_name = ${nameOrId} AND deleted_at IS NULL
`;
if (byName.length === 0) {
throw new Error(
`No active OAuth client found with name or id "${nameOrId}". ` +
`Run \`gbrain auth register-client <name>\` to create one, ` +
`or \`gbrain auth list-clients\` to see what exists. ` +
`(Soft-deleted clients are hidden by default.)`,
);
}
if (byName.length > 1) {
const ids = byName.map((r) => ` ${String(r.client_id)}`).join('\n');
throw new Error(
`Multiple active OAuth clients named "${nameOrId}". Pass the full client_id instead:\n${ids}`,
);
}
return normalizeClientRow(byName[0]);
}
function normalizeClientRow(row: Record<string, unknown>): ResolvedClient {
const fed = row.federated_read;
return {
client_id: String(row.client_id),
client_name: String(row.client_name),
source_id: row.source_id == null ? null : String(row.source_id),
federated_read: Array.isArray(fed) ? (fed as string[]).map(String) : [],
deleted_at: row.deleted_at == null
? null
: (row.deleted_at as Date | string),
};
}
/**
* Validate the source_id shape AND DB existence. Codex finding #3 (medium):
* a manually-INSERTed source row with weird chars (e.g. comma, quote)
* would otherwise land in oauth_clients.federated_read as a never-deletable
* malformed entry. Fail at the boundary before the existence query so
* malformed input gets the validator's hint, not a "does not exist" hint
* pointing at a non-creatable id.
*/
export async function assertSourceExists(sql: SqlQuery, sourceId: string): Promise<void> {
assertValidSourceId(sourceId);
const rows = await sql`SELECT id FROM sources WHERE id = ${sourceId} LIMIT 1`;
if (rows.length === 0) {
throw new Error(
`Source "${sourceId}" does not exist. Run \`gbrain sources list\` to see registered sources, ` +
`or \`gbrain sources add ${sourceId}\` to create it.`,
);
}
}
/**
* Atomic append: array_append + NOT-ANY guard so the row-lock fully
* serializes concurrent grant/revoke against the same client. Codex
* finding #1 (HIGH): the previous read-modify-write shape allowed a
* concurrent revoke to be silently UNDONE by a racing grant.
*
* Returns the post-write federated_read array, or null when no rows
* matched (already-granted, soft-deleted, or missing client). Callers
* disambiguate via prior resolveClient + includes() check.
*
* `WHERE deleted_at IS NULL` is part of the atomic guard so a client
* soft-deleted between resolveClient and the UPDATE can't be mutated.
*/
async function appendFederatedReadAtomic(
sql: SqlQuery,
clientId: string,
sourceId: string,
): Promise<string[] | null> {
const rows = await sql`
UPDATE oauth_clients
SET federated_read = array_append(federated_read, ${sourceId})
WHERE client_id = ${clientId}
AND deleted_at IS NULL
AND NOT (${sourceId} = ANY(federated_read))
RETURNING federated_read
`;
if (rows.length === 0) return null;
const fed = rows[0].federated_read;
return Array.isArray(fed) ? (fed as string[]).map(String) : [];
}
/**
* Atomic remove: array_remove + ANY guard. Same race-correctness story
* as appendFederatedReadAtomic. Returns post-write array or null.
*/
async function removeFederatedReadAtomic(
sql: SqlQuery,
clientId: string,
sourceId: string,
): Promise<string[] | null> {
const rows = await sql`
UPDATE oauth_clients
SET federated_read = array_remove(federated_read, ${sourceId})
WHERE client_id = ${clientId}
AND deleted_at IS NULL
AND ${sourceId} = ANY(federated_read)
RETURNING federated_read
`;
if (rows.length === 0) return null;
const fed = rows[0].federated_read;
return Array.isArray(fed) ? (fed as string[]).map(String) : [];
}
/**
* Wholesale array overwrite for `set-federated-read`. Honors the
* deleted_at filter. Last-writer-wins semantics under concurrent
* `set` calls is acceptable the user is asserting "this exact list"
* intent; concurrent set+set just means whichever ran second wins.
* Concurrent set+grant or set+revoke is also last-writer-wins, which
* is the documented contract for `set`.
*/
async function replaceFederatedReadAtomic(
sql: SqlQuery,
clientId: string,
next: string[],
): Promise<string[] | null> {
// TEXT[] binding via pgArray() string-literal escaping (see helper
// for the security note). Our narrow SqlQuery surface
// (src/core/sql-query.ts) doesn't bind JS arrays directly.
const literal = pgArray(next);
const rows = await sql`
UPDATE oauth_clients
SET federated_read = ${literal}
WHERE client_id = ${clientId}
AND deleted_at IS NULL
RETURNING federated_read
`;
if (rows.length === 0) return null;
const fed = rows[0].federated_read;
return Array.isArray(fed) ? (fed as string[]).map(String) : [];
}
/**
* Pure helper: dedupe a comma-separated source-id list while preserving
* insertion order. Empty input empty array. Exported so the CLI parser
* and tests share one normalizer.
*/
export function parseSourceCsv(csv: string): string[] {
const requested = csv.split(',').map((s) => s.trim()).filter(Boolean);
const seen = new Set<string>();
const out: string[] = [];
for (const s of requested) {
if (!seen.has(s)) {
seen.add(s);
out.push(s);
}
}
return out;
}
export interface FederatedReadOpts {
/** When true, compute the outcome but skip the persisting UPDATE. */
dryRun?: boolean;
}
/**
* Core: append a source to the client's federated_read.
*
* Atomicity contract (Codex finding #1, HIGH):
* The actual write goes through `appendFederatedReadAtomic` which
* serializes at the row-lock so concurrent grant/revoke against the
* same client cannot lose updates. The race vector that previously
* silently restored revoked access is closed: under two operators
* racing `revoke-read sensitive` + `grant-read harmless`, postgres
* serializes the two UPDATEs and BOTH ops apply (sensitive removed,
* harmless added), instead of one clobbering the other.
*
* The reported `before` is the snapshot at resolveClient time, which
* may be stale relative to a concurrent racer. The `after` reflects
* the post-UPDATE state from RETURNING (always fresh).
*/
export async function grantReadCore(
sql: SqlQuery,
nameOrId: string,
sourceId: string,
opts: FederatedReadOpts = {},
): Promise<FederatedReadOutcome> {
const client = await resolveClient(sql, nameOrId);
await assertSourceExists(sql, sourceId);
if (client.federated_read.includes(sourceId)) {
return { kind: 'noop', reason: 'already-granted', client, current: client.federated_read };
}
if (opts.dryRun) {
// Compute the would-be result without touching the row. Last-known
// snapshot is best-effort under concurrent writes.
const projected = [...client.federated_read, sourceId];
return { kind: 'updated', client, before: client.federated_read, after: projected };
}
const after = await appendFederatedReadAtomic(sql, client.client_id, sourceId);
if (after === null) {
// Two equivalent failure modes: (a) racing grant-read already added
// the source and the NOT-ANY guard suppressed our UPDATE, or
// (b) the client was soft-deleted between resolveClient and UPDATE.
// (a) is the more common path. Re-resolve to confirm + report.
const reresolved = await resolveClient(sql, client.client_id, { includeDeleted: true });
if (reresolved.deleted_at != null) {
throw new Error(`Client "${client.client_name}" was soft-deleted before write could land.`);
}
return { kind: 'noop', reason: 'already-granted', client: reresolved, current: reresolved.federated_read };
}
return { kind: 'updated', client, before: client.federated_read, after };
}
/**
* Core: remove a source from the client's federated_read. Atomic via
* array_remove + ANY-guard. Same race-correctness rationale as
* grantReadCore concurrent ops serialize at the row lock.
*/
export async function revokeReadCore(
sql: SqlQuery,
nameOrId: string,
sourceId: string,
opts: FederatedReadOpts = {},
): Promise<FederatedReadOutcome> {
const client = await resolveClient(sql, nameOrId);
if (!client.federated_read.includes(sourceId)) {
return { kind: 'noop', reason: 'not-present', client, current: client.federated_read };
}
if (opts.dryRun) {
const projected = client.federated_read.filter((s) => s !== sourceId);
return { kind: 'updated', client, before: client.federated_read, after: projected };
}
const after = await removeFederatedReadAtomic(sql, client.client_id, sourceId);
if (after === null) {
// Same disambiguation as grant: either a concurrent revoke already
// removed the source (most common) or the client was soft-deleted.
const reresolved = await resolveClient(sql, client.client_id, { includeDeleted: true });
if (reresolved.deleted_at != null) {
throw new Error(`Client "${client.client_name}" was soft-deleted before write could land.`);
}
return { kind: 'noop', reason: 'not-present', client: reresolved, current: reresolved.federated_read };
}
return { kind: 'updated', client, before: client.federated_read, after };
}
/**
* Core: replace the whole federated_read list. Idempotent on same list.
*
* Race semantics: wholesale-overwrite + deleted_at guard. Concurrent
* set+set is last-writer-wins (documented contract for `set` the
* operator is asserting the exact list). Concurrent set+grant or
* set+revoke is also last-writer-wins. If a strict-merge semantics is
* needed, use grant-read / revoke-read individually.
*/
export async function setFederatedReadCore(
sql: SqlQuery,
nameOrId: string,
sourceCsv: string,
opts: FederatedReadOpts = {},
): Promise<FederatedReadOutcome> {
const next = parseSourceCsv(sourceCsv);
const client = await resolveClient(sql, nameOrId);
for (const s of next) {
await assertSourceExists(sql, s);
}
const prev = client.federated_read;
const same = prev.length === next.length && prev.every((v, i) => v === next[i]);
if (same) {
return { kind: 'noop', reason: 'same-list', client, current: prev };
}
if (opts.dryRun) {
return { kind: 'updated', client, before: prev, after: next };
}
const after = await replaceFederatedReadAtomic(sql, client.client_id, next);
if (after === null) {
throw new Error(`Client "${client.client_name}" was soft-deleted before write could land.`);
}
return { kind: 'updated', client, before: prev, after };
}
function printOutcome(
verb: 'grant' | 'revoke' | 'set',
sourceArg: string,
outcome: FederatedReadOutcome,
dryRun: boolean,
): void {
// Terminal-injection defense (Codex finding #5, low): a client_name
// registered via DCR with ANSI escapes or control chars would
// otherwise poison this output. Sanitize ALL strings that round-trip
// from the DB before printing.
const s = sanitizeForTerminal;
const prefix = dryRun ? '[dry-run] ' : '';
if (outcome.kind === 'noop') {
const name = s(outcome.client.client_name);
if (outcome.reason === 'already-granted') {
console.log(`${prefix}No change: "${name}" already reads "${s(sourceArg)}".`);
} else if (outcome.reason === 'not-present') {
console.log(`${prefix}No change: "${name}" did not read "${s(sourceArg)}".`);
} else {
console.log(`${prefix}No change: "${name}" federated_read already matches.`);
}
console.log(` federated_read: ${outcome.current.map(s).join(', ') || '(empty)'}`);
return;
}
const { client, before, after } = outcome;
const name = s(client.client_name);
const wouldOrDid = dryRun ? 'Would' : 'Did';
if (verb === 'grant') {
console.log(`${prefix}${wouldOrDid} grant: "${name}" can now read "${s(sourceArg)}".`);
console.log(` federated_read: ${after.map(s).join(', ')}`);
} else if (verb === 'revoke') {
console.log(`${prefix}${wouldOrDid} revoke: "${name}" no longer reads "${s(sourceArg)}".`);
console.log(` federated_read: ${after.map(s).join(', ') || '(empty — client has no federated reads)'}`);
} else {
console.log(`${prefix}${wouldOrDid} update "${name}" federated_read:`);
console.log(` before: ${before.map(s).join(', ') || '(empty)'}`);
console.log(` after: ${after.map(s).join(', ') || '(empty)'}`);
}
if (after.length === 0) {
console.log(
'Warning: client now reads no sources via federation. Queries through this ' +
'client will only see content scoped explicitly via its write source.',
);
}
}
/**
* Strip `--dry-run` from a positional-arg list. Returns the filtered list
* plus the flag value. Kept positional-tolerant the existing
* `auth grant-read alice source` shape MUST keep working, AND
* `auth grant-read alice source --dry-run` AND `auth grant-read --dry-run alice source`.
*/
export function extractDryRun(args: string[]): { dryRun: boolean; rest: string[] } {
let dryRun = false;
const rest: string[] = [];
for (const a of args) {
if (a === '--dry-run') {
dryRun = true;
continue;
}
rest.push(a);
}
return { dryRun, rest };
}
async function grantRead(args: string[]): Promise<void> {
const { dryRun, rest } = extractDryRun(args);
const [nameOrId, sourceId] = rest;
if (!nameOrId || !sourceId) {
console.error('Usage: gbrain auth grant-read <client-name-or-id> <source-id> [--dry-run]');
process.exit(1);
}
try {
await withConfiguredSql(async (sql) => {
const outcome = await grantReadCore(sql, nameOrId, sourceId, { dryRun });
printOutcome('grant', sourceId, outcome, dryRun);
});
} catch (e: any) {
console.error('Error:', e.message);
process.exit(1);
}
}
async function revokeRead(args: string[]): Promise<void> {
const { dryRun, rest } = extractDryRun(args);
const [nameOrId, sourceId] = rest;
if (!nameOrId || !sourceId) {
console.error('Usage: gbrain auth revoke-read <client-name-or-id> <source-id> [--dry-run]');
process.exit(1);
}
try {
await withConfiguredSql(async (sql) => {
const outcome = await revokeReadCore(sql, nameOrId, sourceId, { dryRun });
printOutcome('revoke', sourceId, outcome, dryRun);
});
} catch (e: any) {
console.error('Error:', e.message);
process.exit(1);
}
}
async function setFederatedRead(args: string[]): Promise<void> {
const { dryRun, rest } = extractDryRun(args);
const [nameOrId, sourceCsv] = rest;
if (!nameOrId || sourceCsv === undefined) {
console.error(
'Usage: gbrain auth set-federated-read <client-name-or-id> <source-id1,source-id2,...> [--dry-run]',
);
console.error('Pass an empty string ("") to clear all federated reads.');
process.exit(1);
}
try {
await withConfiguredSql(async (sql) => {
const outcome = await setFederatedReadCore(sql, nameOrId, sourceCsv, { dryRun });
printOutcome('set', sourceCsv, outcome, dryRun);
});
} catch (e: any) {
console.error('Error:', e.message);
process.exit(1);
}
}
async function revokeClient(clientId: string) {
if (!clientId) {
console.error('Usage: auth revoke-client <client_id>');
@@ -319,7 +884,7 @@ async function revokeClient(clientId: string) {
console.error(`No client found with id "${clientId}"`);
process.exit(1);
}
console.log(`OAuth client revoked: "${rows[0].client_name}" (${clientId})`);
console.log(`OAuth client revoked: "${sanitizeForTerminal(String(rows[0].client_name))}" (${clientId})`);
console.log('Tokens and authorization codes purged via cascade.');
});
} catch (e: any) {
@@ -440,6 +1005,15 @@ export function parseRegisterClientArgs(args: string[]): RegisterClientArgs {
if (!grantTypesSet && out.redirectUris.length > 0) {
out.grantTypes = ['authorization_code', 'refresh_token'];
}
// Codex re-review (medium): validate source_id shape at the CLI boundary
// so register-client can't seed malformed entries into source_id /
// federated_read that subsequent grant/revoke/set commands can't manage.
assertValidSourceId(out.sourceId);
if (out.federatedRead) {
for (const s of out.federatedRead) {
assertValidSourceId(s);
}
}
return out;
}
@@ -557,6 +1131,10 @@ export async function runAuth(args: string[]): Promise<void> {
}
case 'register-client': await registerClient(rest[0], rest.slice(1)); return;
case 'revoke-client': await revokeClient(rest[0]); return;
case 'list-clients': await listClients(rest); return;
case 'grant-read': await grantRead(rest); return;
case 'revoke-read': await revokeRead(rest); return;
case 'set-federated-read': await setFederatedRead(rest); return;
case 'test': {
const tokenIdx = rest.indexOf('--token');
const url = rest.find(a => !a.startsWith('--') && a !== rest[tokenIdx + 1]);
@@ -594,6 +1172,13 @@ Usage:
--bound-max-concurrent <n> Bound submit_agent concurrency (default: 1)
--budget-usd-per-day <usd> Bound submit_agent daily spend cap
gbrain auth revoke-client <client_id> Hard-delete an OAuth 2.1 client (cascades to tokens + codes)
gbrain auth list-clients [--json] List OAuth 2.1 clients with scope + write source + federated_read.
gbrain auth grant-read <name|client_id> <source-id> [--dry-run]
Add a source to the client's federated_read list (idempotent).
gbrain auth revoke-read <name|client_id> <source-id> [--dry-run]
Remove a source from the client's federated_read list (idempotent).
gbrain auth set-federated-read <name|client_id> "<id1,id2,...>" [--dry-run]
Replace the client's whole federated_read list. Pass "" to clear.
gbrain auth test <url> --token <token> Smoke-test a remote MCP server
`);
}
+5 -10
View File
@@ -170,14 +170,10 @@ export async function runImport(
// v0.22.13 (PR #490 Q2): shared parseWorkers helper rejects bad input
// (--workers 0, -3, "foo") with a loud error instead of silently falling
// through to 1. Mirrors sync.ts's flag handling.
const { parseWorkers, autoConcurrency } = await import('../core/sync-concurrency.ts');
// #1207: undefined (no --workers flag) defers to autoConcurrency below —
// the shared sync/import policy (PGLite → 1, >100 files → 4) — instead of
// hardcoding serial. Large Postgres imports stop paying one embedding
// round-trip per file in sequence.
let workerCount: number | undefined;
const { parseWorkers } = await import('../core/sync-concurrency.ts');
let workerCount: number;
try {
workerCount = parseWorkers(workersArg ?? undefined);
workerCount = parseWorkers(workersArg ?? undefined) ?? 1;
} catch (e) {
console.error(e instanceof Error ? e.message : String(e));
process.exit(1);
@@ -256,9 +252,8 @@ export async function runImport(
}
const files = resumeFilter(allFiles, dir, completed);
// Determine actual worker count. Explicit --workers wins; otherwise the
// shared autoConcurrency policy decides from engine kind + file count.
const actualWorkers = autoConcurrency(engine, files.length, workerCount);
// Determine actual worker count
const actualWorkers = workerCount > 1 ? workerCount : 1;
if (actualWorkers > 1) {
console.log(`Using ${actualWorkers} parallel workers`);
}
+141
View File
@@ -365,6 +365,42 @@ export interface AgentClientSpend {
inflight_count: number;
}
/**
* `/admin/api/sources` source list the input rows for buildSyncStatusReport.
*
* Queries the JSONB config column directly (listSources doesn't carry it,
* but buildSyncStatusReport needs syncEnabled / strategy fields).
*
* Deliberately does NOT filter on local_path: in a push-only deployment
* (content arrives via MCP put_page / capture / ingest, not `gbrain sync`
* of a server checkout) every source has a null local_path filtering on
* it would empty both the Sources tab AND the federation source-picker.
* buildSyncStatusReport does no disk I/O, so null-local_path sources
* report fine (pages/chunks from SQL, staleness 'unknown' / never-synced).
*/
export async function queryAdminSources(engine: BrainEngine): Promise<
Array<{ id: string; name: string; local_path: string | null; config: Record<string, unknown> }>
> {
const rows = await engine.executeRaw<{
id: string;
name: string;
local_path: string | null;
config: Record<string, unknown> | string | null;
}>(
`SELECT id, name, local_path, config FROM sources
WHERE archived IS NOT TRUE
ORDER BY id`,
);
return rows.map((r) => ({
id: r.id,
name: r.name,
local_path: r.local_path,
config: typeof r.config === 'string'
? (JSON.parse(r.config) as Record<string, unknown>)
: (r.config ?? {}),
}));
}
export async function queryAgentClientSpend(engine: BrainEngine): Promise<AgentClientSpend[]> {
const sql = sqlQueryForEngine(engine);
const rows = await sql`
@@ -1511,6 +1547,111 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
}
});
// ---------------------------------------------------------------------------
// Sources tab — read-only view of registered sources with sync + embed
// coverage stats. Drives the admin SPA's `Sources` page.
//
// Returns the same shape `gbrain sources status --json` prints, so the
// SPA stays in lockstep with the CLI surface.
// ---------------------------------------------------------------------------
app.get('/admin/api/sources', requireAdmin, async (_req: Request, res: Response) => {
try {
const { buildSyncStatusReport } = await import('./sync.ts');
const report = await buildSyncStatusReport(engine, await queryAdminSources(engine));
res.json(report);
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
res.status(503).json({ error: 'service_unavailable', detail: msg });
}
});
// ---------------------------------------------------------------------------
// Federated-read management (admin-side counterparts of the CLI commands
// `gbrain auth grant-read / revoke-read / set-federated-read`). All three
// route through the same *Core helpers as the CLI so race-safety,
// soft-delete filter, and source-id shape validation apply uniformly.
//
// The admin SPA's `Agents` page renders "Manage reads" actions per
// client backed by these endpoints.
// ---------------------------------------------------------------------------
app.get('/admin/api/agents/federated-read', requireAdmin, async (_req: Request, res: Response) => {
try {
const rows = await sql`
SELECT client_id, client_name, source_id, federated_read
FROM oauth_clients
WHERE deleted_at IS NULL
ORDER BY client_name
`;
const clients = rows.map((r) => ({
client_id: String(r.client_id),
client_name: String(r.client_name),
source_id: r.source_id == null ? null : String(r.source_id),
federated_read: Array.isArray(r.federated_read)
? (r.federated_read as string[]).map(String)
: [],
}));
res.json({ clients });
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
res.status(503).json({ error: 'service_unavailable', detail: msg });
}
});
app.post('/admin/api/agents/:clientId/grant-read', requireAdmin, express.json(), async (req: Request, res: Response) => {
const clientId = String(req.params.clientId ?? '');
const sourceId = String(req.body?.source_id ?? '').trim();
if (!clientId || !sourceId) {
res.status(400).json({ error: 'invalid_request', detail: 'clientId path param + source_id body required' });
return;
}
try {
const { grantReadCore } = await import('./auth.ts');
const outcome = await grantReadCore(sql, clientId, sourceId);
res.json({ outcome });
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
res.status(400).json({ error: 'mutation_failed', detail: msg });
}
});
app.post('/admin/api/agents/:clientId/revoke-read', requireAdmin, express.json(), async (req: Request, res: Response) => {
const clientId = String(req.params.clientId ?? '');
const sourceId = String(req.body?.source_id ?? '').trim();
if (!clientId || !sourceId) {
res.status(400).json({ error: 'invalid_request', detail: 'clientId path param + source_id body required' });
return;
}
try {
const { revokeReadCore } = await import('./auth.ts');
const outcome = await revokeReadCore(sql, clientId, sourceId);
res.json({ outcome });
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
res.status(400).json({ error: 'mutation_failed', detail: msg });
}
});
app.post('/admin/api/agents/:clientId/set-federated-read', requireAdmin, express.json(), async (req: Request, res: Response) => {
const clientId = String(req.params.clientId ?? '');
const rawIds = req.body?.source_ids;
if (!clientId || !Array.isArray(rawIds)) {
res.status(400).json({ error: 'invalid_request', detail: 'clientId path param + source_ids[] body required' });
return;
}
// Encode the array as CSV so the same setFederatedReadCore signature
// (string CSV input) the CLI uses applies here. Empty array → empty
// string → clears the list.
const csv = rawIds.map((s) => String(s).trim()).filter(Boolean).join(',');
try {
const { setFederatedReadCore } = await import('./auth.ts');
const outcome = await setFederatedReadCore(sql, clientId, csv);
res.json({ outcome });
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
res.status(400).json({ error: 'mutation_failed', detail: msg });
}
});
// ---------------------------------------------------------------------------
// SSE live activity feed
// ---------------------------------------------------------------------------
+6 -24
View File
@@ -1513,21 +1513,12 @@ export async function embed(texts: string[], opts?: EmbedOpts): Promise<Float32A
const embedding = recipe.touchpoints?.embedding;
const maxBatchTokens = embedding?.max_batch_tokens;
const maxBatchCount = embedding?.max_batch_count;
const charsPerToken = embedding?.chars_per_token ?? DEFAULT_CHARS_PER_TOKEN;
// Pre-split is gated on max_batch_tokens / max_batch_count. Recipes with
// neither (e.g. OpenAI) ride the fast path: one embedMany call, no
// recursion safety net.
const batches = (maxBatchTokens || maxBatchCount)
? splitByTokenBudget(
truncated,
maxBatchTokens
? Math.floor(maxBatchTokens * effectiveSafetyFactor(recipe))
: Number.MAX_SAFE_INTEGER,
charsPerToken,
maxBatchCount,
)
// Pre-split is gated on max_batch_tokens. Recipes without it (e.g. OpenAI)
// ride the fast path: one embedMany call, no recursion safety net.
const batches = maxBatchTokens
? splitByTokenBudget(truncated, Math.floor(maxBatchTokens * effectiveSafetyFactor(recipe)), charsPerToken)
: [truncated];
const allEmbeddings: Float32Array[] = [];
@@ -1577,9 +1568,6 @@ export async function embed(texts: string[], opts?: EmbedOpts): Promise<Float32A
* responsible for applying any safety-factor shrink before passing in.
* @param charsPerToken - Provider-specific character density. Defaults to
* `DEFAULT_CHARS_PER_TOKEN` (4) when omitted, matching OpenAI tiktoken.
* @param maxBatchCount - #1199: optional cap on INPUTS per sub-batch, for
* providers that reject batches by count (DashScope: 10). When omitted,
* only the token budget governs.
*
* @internal exported for tests; not part of the public gateway API.
*/
@@ -1587,17 +1575,15 @@ export function splitByTokenBudget(
texts: string[],
budgetTokens: number,
charsPerToken: number = DEFAULT_CHARS_PER_TOKEN,
maxBatchCount?: number,
): string[][] {
const ratio = charsPerToken > 0 ? charsPerToken : DEFAULT_CHARS_PER_TOKEN;
const maxCount = maxBatchCount !== undefined && maxBatchCount > 0 ? maxBatchCount : Infinity;
const batches: string[][] = [];
let current: string[] = [];
let currentTokens = 0;
for (const text of texts) {
const estTokens = Math.ceil(text.length / ratio);
if (current.length > 0 && (currentTokens + estTokens > budgetTokens || current.length >= maxCount)) {
if (current.length > 0 && currentTokens + estTokens > budgetTokens) {
batches.push(current);
current = [];
currentTokens = 0;
@@ -1623,11 +1609,7 @@ export function isTokenLimitError(err: unknown): boolean {
/token.*limit.*exceeded/i.test(msg) ||
// OpenAI embeddings: "Invalid 'input': maximum request size is 300000 tokens per request."
/maximum request size.*tokens/i.test(msg) ||
/max.*tokens.*per.*request/i.test(msg) ||
// DashScope: "batch size is invalid, it should not be larger than 10." (#1199)
// Count-cap error, but recursive halving shrinks count too, so the same
// safety net converges.
/batch size is invalid/i.test(msg)
/max.*tokens.*per.*request/i.test(msg)
);
}
-4
View File
@@ -31,10 +31,6 @@ export const dashscope: Recipe = {
// path. Conservative declaration so the gateway pre-splits before
// hitting whatever undocumented server-side limit exists.
max_batch_tokens: 8192,
// #1199: DashScope hard-caps embeddings at 10 inputs per request
// ("batch size is invalid, it should not be larger than 10"). The
// token budget alone admits far more than 10 short chunks per batch.
max_batch_count: 10,
// text-embedding-v3 mixes English + CJK heavily; the tokenizer is
// closer to Voyage density than OpenAI tiktoken for CJK-dominant
// content. Conservative chars_per_token=2 leaves headroom.
-9
View File
@@ -16,15 +16,6 @@ export const google: Recipe = {
dims_options: [768, 1536, 3072],
cost_per_1m_tokens_usd: 0.15,
price_last_verified: '2026-04-20',
// #970: Gemini's documented limits are per-INPUT (2048 tokens,
// silently truncated beyond) and per-REQUEST count (batchEmbedContents
// caps at 100 inputs). There is no separate per-request token cap, so
// the token budget is derived: 100 inputs × 2048 tokens. The count cap
// binds first for typical chunk sizes. Do NOT copy the 2048 per-input
// limit into max_batch_tokens — that would over-split 50×.
max_batch_tokens: 204_800,
chars_per_token: 4,
max_batch_count: 100,
},
expansion: {
models: ['gemini-2.0-flash', 'gemini-2.0-flash-lite'],
+1 -4
View File
@@ -58,8 +58,5 @@ export function getRecipe(id: string): Recipe | undefined {
}
export function listRecipes(): Recipe[] {
// Read the map (not ALL) so there is one source of truth — getRecipe,
// model-resolver, and listRecipes all see the same registry, and tests
// can inject a synthetic recipe via RECIPES to exercise registry walks.
return [...RECIPES.values()];
return [...ALL];
}
-10
View File
@@ -46,16 +46,6 @@ export interface EmbeddingTouchpoint {
* Only consulted when `max_batch_tokens` is also set.
*/
chars_per_token?: number;
/**
* #1199: maximum number of INPUTS per embedding request, for providers
* that hard-cap batch size by count rather than (or in addition to)
* tokens DashScope text-embedding-v3 rejects batches > 10 with
* `InvalidParameter`, Gemini batchEmbedContents caps at 100 requests.
* When set, the gateway's pre-split flushes a sub-batch at this count
* even if the token budget still has room. Independent of
* `max_batch_tokens`; either alone triggers the pre-split.
*/
max_batch_count?: number;
/**
* Budget-utilization ceiling in (0, 1]. The gateway pre-splits at
* `safety_factor × max_batch_tokens` to leave headroom for tokenizer
+6 -56
View File
@@ -79,34 +79,15 @@ export interface EmbedBatchOptions {
* and amplify rate-limit pressure.
*/
maxRetries?: number;
/**
* #1818: bounded parallelism across BATCH_SIZE sub-batches. Defaults to
* `GBRAIN_EMBED_BATCH_CONCURRENCY` env, else 4. Results are
* index-addressed so output order always matches input order. Set 1 to
* force the pre-v0.42 serial dispatch.
*/
concurrency?: number;
}
/**
* Embed a batch of texts via the gateway. Sub-batches of 100 so upstream
* progress callbacks fire incrementally on large imports. The gateway owns
* adaptive batch splitting and per-recipe token-budget logic; this paginator
* owns progress-callback granularity and (#1818) bounded parallel dispatch
* of the sub-batches the embed-stale.ts worker-pool pattern, scoped down.
* is purely about progress-callback granularity.
*/
const BATCH_SIZE = 100;
const DEFAULT_EMBED_BATCH_CONCURRENCY = 4;
function resolveEmbedBatchConcurrency(options: EmbedBatchOptions): number {
if (options.concurrency !== undefined) {
return Math.max(1, Math.floor(options.concurrency));
}
const env = Number(process.env.GBRAIN_EMBED_BATCH_CONCURRENCY);
if (Number.isFinite(env) && env >= 1) return Math.floor(env);
return DEFAULT_EMBED_BATCH_CONCURRENCY;
}
export async function embedBatch(
texts: string[],
options: EmbedBatchOptions = {},
@@ -122,44 +103,13 @@ export async function embedBatch(
if (texts.length <= BATCH_SIZE && !options.onBatchComplete) {
return gatewayEmbed(texts, gwOpts);
}
// #1818: dispatch sub-batches through a bounded worker pool instead of a
// serial loop. Results are written into a preallocated index-addressed
// array so output order matches input order regardless of completion
// order; onBatchComplete reports a monotonic completed-embedding count.
const slices: Array<{ start: number; texts: string[] }> = [];
const results: Float32Array[] = [];
for (let i = 0; i < texts.length; i += BATCH_SIZE) {
slices.push({ start: i, texts: texts.slice(i, i + BATCH_SIZE) });
const slice = texts.slice(i, i + BATCH_SIZE);
const out = await gatewayEmbed(slice, gwOpts);
results.push(...out);
options.onBatchComplete?.(results.length, texts.length);
}
const results = new Array<Float32Array>(texts.length);
let next = 0;
let done = 0;
const numWorkers = Math.min(resolveEmbedBatchConcurrency(options), slices.length);
// Once any sub-batch fails, `failed` stops the surviving workers from
// dispatching FURTHER slices — the whole call is rejecting anyway, so
// continuing would burn real provider spend in the background and fire
// onBatchComplete after the caller already saw the failure (worst with
// embedBatchWithBackoff, whose 429 backoff assumes nothing is in flight).
// In-flight sibling calls still run to completion (bounded by numWorkers-1).
let failed = false;
const worker = async (): Promise<void> => {
while (!failed && next < slices.length) {
// NOTE: no local aborted-check here — an aborted signal makes the next
// gatewayEmbed call throw (SDK-side), which rejects the pool. Returning
// silently instead would resolve with holes in `results`.
const slice = slices[next++];
let out: Float32Array[];
try {
out = await gatewayEmbed(slice.texts, gwOpts);
} catch (err) {
failed = true;
throw err;
}
for (let j = 0; j < out.length; j++) results[slice.start + j] = out[j];
done += out.length;
if (!failed) options.onBatchComplete?.(done, texts.length);
}
};
await Promise.all(Array.from({ length: numWorkers }, () => worker()));
return results;
}
+1 -1
View File
@@ -55,7 +55,7 @@ export interface AgentClientBindings {
* `redirect_uri` containing `,`) would be parsed by Postgres as MULTIPLE
* array elements, smuggling values past validation. See CSO finding #5.
*/
function pgArray(arr: string[]): string {
export function pgArray(arr: string[]): string {
if (!arr || arr.length === 0) return '{}';
const escaped = arr.map(s => `"${s.replace(/\\/g, '\\\\').replace(/"/g, '\\"')}"`);
return `{${escaped.join(',')}}`;
+82
View File
@@ -0,0 +1,82 @@
import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'bun:test';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { resetPgliteState } from './helpers/reset-pglite.ts';
import { queryAdminSources } from '../src/commands/serve-http.ts';
import { buildSyncStatusReport } from '../src/commands/sync.ts';
/**
* v0.41.29 Sources tab `/admin/api/sources` endpoint SQL.
*
* The endpoint is a thin Express handler over `queryAdminSources` +
* `buildSyncStatusReport`; the source-selection SQL is the load-bearing
* surface (same pattern as test/admin-agents-spend.test.ts).
*
* Pinned behaviors:
* - Excludes archived sources
* - INCLUDES sources with null local_path (push-only brains: filtering
* on local_path emptied the Sources tab + federation source-picker)
* - JSONB config surfaces as an object, defaulting to {}
* - Deterministic ORDER BY id
* - buildSyncStatusReport accepts the rows (no disk I/O on null paths)
*/
let engine: PGLiteEngine;
beforeAll(async () => {
engine = new PGLiteEngine();
await engine.connect({});
await engine.initSchema();
});
afterAll(async () => {
await engine.disconnect();
});
beforeEach(async () => {
await resetPgliteState(engine);
});
describe('queryAdminSources (/admin/api/sources SQL)', () => {
it('includes push-only sources with null local_path', async () => {
await engine.executeRaw(
`INSERT INTO sources (id, name, local_path, config)
VALUES ('push-only', 'push-only', NULL, '{}'::jsonb)`,
);
const sources = await queryAdminSources(engine);
const ids = sources.map((s) => s.id);
expect(ids).toContain('push-only');
expect(sources.find((s) => s.id === 'push-only')!.local_path).toBe(null);
});
it('excludes archived sources', async () => {
await engine.executeRaw(
`INSERT INTO sources (id, name, archived) VALUES ('gone', 'gone', true)`,
);
const sources = await queryAdminSources(engine);
expect(sources.map((s) => s.id)).not.toContain('gone');
});
it('surfaces JSONB config as an object and orders by id', async () => {
await engine.executeRaw(
`INSERT INTO sources (id, name, config)
VALUES ('bbb', 'bbb', '{"syncEnabled": true}'::jsonb),
('aaa', 'aaa', '{}'::jsonb)`,
);
const sources = await queryAdminSources(engine);
const ids = sources.map((s) => s.id);
expect(ids.indexOf('aaa')).toBeLessThan(ids.indexOf('bbb'));
expect(sources.find((s) => s.id === 'bbb')!.config).toEqual({ syncEnabled: true });
expect(sources.find((s) => s.id === 'aaa')!.config).toEqual({});
});
it('buildSyncStatusReport accepts the rows (null local_path does not throw)', async () => {
await engine.executeRaw(
`INSERT INTO sources (id, name, local_path, config)
VALUES ('push-only', 'push-only', NULL, '{}'::jsonb)`,
);
const report = await buildSyncStatusReport(engine, await queryAdminSources(engine));
expect(report.schema_version).toBe(1);
const row = report.sources.find((s) => s.source_id === 'push-only');
expect(row).toBeDefined();
});
});
+5 -109
View File
@@ -39,8 +39,6 @@ import {
__getShrinkStateForTests,
} from '../../src/core/ai/gateway.ts';
import { AIConfigError, AITransientError } from '../../src/core/ai/errors.ts';
import { RECIPES } from '../../src/core/ai/recipes/index.ts';
import type { Recipe } from '../../src/core/ai/types.ts';
// The last test in this file leaves the gateway configured with a remote
// provider + fake key and a REAL embed transport. Without a final reset,
@@ -95,14 +93,6 @@ function configureGoogle(): void {
});
}
function configureDashscope(): void {
configureGateway({
embedding_model: 'dashscope:text-embedding-v3',
embedding_dimensions: 1024,
env: { DASHSCOPE_API_KEY: 'sk-fake' },
});
}
// --------- 1. Pure helpers ---------
describe('splitByTokenBudget (pure helper)', () => {
@@ -159,27 +149,6 @@ describe('splitByTokenBudget (pure helper)', () => {
expect(splitByTokenBudget(texts, 96_000, 0)).toEqual(splitByTokenBudget(texts, 96_000, 4));
expect(splitByTokenBudget(texts, 96_000, -1)).toEqual(splitByTokenBudget(texts, 96_000, 4));
});
// #1199: count cap for providers that reject batches by input count.
test('max_batch_count flushes even when token budget has room', () => {
const texts = Array.from({ length: 25 }, (_, i) => `t${i}`);
const result = splitByTokenBudget(texts, 1_000_000, 4, 10);
expect(result.map(b => b.length)).toEqual([10, 10, 5]);
expect(result.flat()).toEqual(texts);
});
test('token budget still governs alongside max_batch_count', () => {
const texts = ['a'.repeat(50_000), 'b'.repeat(50_000), 'c'.repeat(50_000)];
const result = splitByTokenBudget(texts, 96_000, 1, 10);
expect(result).toHaveLength(3);
});
test('undefined / zero / negative max_batch_count is ignored', () => {
const texts = Array.from({ length: 25 }, () => 'x');
expect(splitByTokenBudget(texts, 1_000_000, 4, undefined)).toHaveLength(1);
expect(splitByTokenBudget(texts, 1_000_000, 4, 0)).toHaveLength(1);
expect(splitByTokenBudget(texts, 1_000_000, 4, -5)).toHaveLength(1);
});
});
describe('isTokenLimitError (pure helper)', () => {
@@ -210,12 +179,6 @@ describe('isTokenLimitError (pure helper)', () => {
expect(isTokenLimitError(new Error('Exceeded 300000 max tokens per request'))).toBe(true);
});
test('matches DashScope batch-count error (#1199)', () => {
expect(isTokenLimitError(new Error(
'InvalidParameter: batch size is invalid, it should not be larger than 10.',
))).toBe(true);
});
test('does not match unrelated errors', () => {
expect(isTokenLimitError(new Error('Connection refused'))).toBe(false);
expect(isTokenLimitError(new Error('Invalid API key'))).toBe(false);
@@ -424,92 +387,26 @@ describe('shrink-on-miss adaptive cache', () => {
});
});
// --------- 8. Pre-split count cap through public embed() (#1199 / #970) ---------
describe('embed() pre-split honors max_batch_count', () => {
beforeEach(() => resetGateway());
afterEach(() => __setEmbedTransportForTests(null));
test('dashscope never dispatches more than 10 inputs per call (#1199)', async () => {
configureDashscope();
const stub = mock(async ({ values }: { values: string[] }) => fakeEmbeddings(values, 1024));
__setEmbedTransportForTests(stub as any);
// 25 short texts fit trivially in the 8192-token budget; without the
// count cap they'd ship as ONE batch and DashScope would reject it.
const texts = Array.from({ length: 25 }, (_, i) => `short-${i}`);
const result = await embed(texts);
expect(result).toHaveLength(25);
const callLengths = stub.mock.calls.map(([arg]) => (arg as { values: string[] }).values.length);
expect(Math.max(...callLengths)).toBeLessThanOrEqual(10);
expect(callLengths.reduce((a, b) => a + b, 0)).toBe(25);
// Order preserved across sub-batches.
expect((stub.mock.calls[0][0] as { values: string[] }).values[0]).toBe('short-0');
});
test('google pre-splits at 100 inputs per batchEmbedContents call (#970)', async () => {
configureGoogle();
const stub = mock(async ({ values }: { values: string[] }) => fakeEmbeddings(values, 768));
__setEmbedTransportForTests(stub as any);
const texts = Array.from({ length: 250 }, (_, i) => `g${i}`);
const result = await embed(texts);
expect(result).toHaveLength(250);
const callLengths = stub.mock.calls.map(([arg]) => (arg as { values: string[] }).values.length);
expect(callLengths).toEqual([100, 100, 50]);
});
});
// --------- 7. Startup warning (D9-B) ---------
describe('startup warning for recipes missing max_batch_tokens', () => {
beforeEach(() => resetGateway());
// #970 closed google's missing cap, so no registered recipe is capless
// anymore. Inject a synthetic capless recipe to keep the warning path
// covered for the NEXT recipe that forgets the field.
const caplessRecipe: Recipe = {
id: 'capless-test',
name: 'Capless Test Provider',
tier: 'openai-compat',
implementation: 'openai-compatible',
base_url_default: 'https://example.invalid/v1',
auth_env: { required: [] },
touchpoints: {
embedding: { models: ['capless-embed-1'], default_dims: 768 },
},
};
function configureCapless(): void {
configureGateway({
embedding_model: 'capless-test:capless-embed-1',
embedding_dimensions: 768,
env: {},
});
}
test('configured missing-cap recipe warns once; unrelated recipes stay quiet', () => {
const warnings: string[] = [];
const original = console.warn;
console.warn = (msg: string) => warnings.push(String(msg));
RECIPES.set(caplessRecipe.id, caplessRecipe);
try {
configureOpenAI();
expect(warnings.length).toBe(0);
// #970 regression: google now declares max_batch_tokens → quiet.
configureGoogle();
expect(warnings.length).toBe(0);
configureCapless();
const firstCallCount = warnings.length;
// Reconfigure: the warning should NOT re-fire for the same recipes
// within one process (we already told the operator).
configureCapless();
configureGoogle();
expect(warnings.length).toBe(firstCallCount);
} finally {
console.warn = original;
RECIPES.delete(caplessRecipe.id);
}
// The warning text should match the documented contract.
@@ -518,12 +415,11 @@ describe('startup warning for recipes missing max_batch_tokens', () => {
);
expect(contractMatch.length).toBe(1);
// Voyage + google declare max_batch_tokens → suppressed. OpenAI is the
// canonical fast-path recipe → also suppressed by id. All must be
// absent from the warnings; only the synthetic capless recipe fires.
// Voyage declares max_batch_tokens → suppressed. OpenAI is the
// canonical fast-path recipe → also suppressed by id. Both must be
// absent from the warnings.
expect(warnings.find(w => w.includes('"voyage"'))).toBeUndefined();
expect(warnings.find(w => w.includes('"openai"'))).toBeUndefined();
expect(warnings.find(w => w.includes('"google"'))).toBeUndefined();
expect(warnings.find(w => w.includes('"capless-test"'))).toBeDefined();
expect(warnings.find(w => w.includes('"google"'))).toBeDefined();
});
});
+13 -13
View File
@@ -52,7 +52,16 @@ describe('v0.32 #779: no_batch_cap suppresses the missing-max_batch_tokens warni
}
});
test('configureGateway does NOT warn for google now that it declares batch caps (#970)', () => {
test('configureGateway warns for google only when google embedding is configured', () => {
warnSpy.mockClear();
resetGateway();
configureGateway({ env: {} });
let messages = warnSpy.mock.calls.map(c => String(c[0] ?? ''));
expect(
messages.some(m => m.includes('"google"') && m.includes('without max_batch_tokens')),
'google should not warn while OpenAI default is configured',
).toBe(false);
warnSpy.mockClear();
resetGateway();
configureGateway({
@@ -60,20 +69,11 @@ describe('v0.32 #779: no_batch_cap suppresses the missing-max_batch_tokens warni
embedding_dimensions: 768,
env: { GOOGLE_GENERATIVE_AI_API_KEY: 'fake' },
});
const messages = warnSpy.mock.calls.map(c => String(c[0] ?? ''));
messages = warnSpy.mock.calls.map(c => String(c[0] ?? ''));
expect(
messages.some(m => m.includes('"google"') && m.includes('without max_batch_tokens')),
'google declares max_batch_tokens/max_batch_count since #970 — no warning',
).toBe(false);
});
test('google recipe declares its derived batch caps (#970)', () => {
const e = getRecipe('google')!.touchpoints.embedding!;
// Count cap is the REAL Gemini limit (batchEmbedContents: 100 inputs);
// the token budget is derived (100 × 2048 per-input tokens), NOT the
// 2048 per-input limit — copying that verbatim would over-split 50×.
expect(e.max_batch_count).toBe(100);
expect(e.max_batch_tokens).toBe(204_800);
'google should warn when configured because it has fixed-cap models',
).toBe(true);
});
test('every recipe with empty models[] declares user_provided_models OR has openai-fast-path', () => {
-5
View File
@@ -55,11 +55,6 @@ describe('recipe: dashscope', () => {
expect(r.touchpoints.embedding!.chars_per_token).toBeGreaterThan(0);
});
test('declares max_batch_count: 10 — DashScope rejects larger batches (#1199)', () => {
const r = getRecipe('dashscope')!;
expect(r.touchpoints.embedding!.max_batch_count).toBe(10);
});
test('dimsProviderOptions threads dimensions for text-embedding-v3 (Matryoshka)', async () => {
// Codex finding #1: DashScope text-embedding-v3 is Matryoshka 64-1024.
// Without `dimensions` on the wire, user-selected non-default dims are
+568
View File
@@ -0,0 +1,568 @@
/**
* Tests for `gbrain auth grant-read|revoke-read|set-federated-read`.
*
* Pure helper: parseSourceCsv (no DB).
* DB-coupled: resolveClient, assertSourceExists, *Core fns exercised
* against a real PGLite via the canonical block.
*/
import { describe, expect, test, beforeAll, afterAll, beforeEach } from 'bun:test';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { resetPgliteState } from './helpers/reset-pglite.ts';
import { sqlQueryForEngine } from '../src/core/sql-query.ts';
import { pgArray } from '../src/core/oauth-provider.ts';
import {
parseSourceCsv,
resolveClient,
assertSourceExists,
grantReadCore,
revokeReadCore,
setFederatedReadCore,
extractDryRun,
sanitizeForTerminal,
} from '../src/commands/auth.ts';
let engine: PGLiteEngine;
beforeAll(async () => {
engine = new PGLiteEngine();
await engine.connect({});
await engine.initSchema();
});
afterAll(async () => {
await engine.disconnect();
});
beforeEach(async () => {
await resetPgliteState(engine);
});
// ---------------------------------------------------------------------------
// pure helpers
// ---------------------------------------------------------------------------
describe('parseSourceCsv', () => {
test('splits and trims', () => {
expect(parseSourceCsv('a,b,c')).toEqual(['a', 'b', 'c']);
expect(parseSourceCsv(' a , b ')).toEqual(['a', 'b']);
});
test('drops empty segments', () => {
expect(parseSourceCsv('a,,b,')).toEqual(['a', 'b']);
expect(parseSourceCsv(',,')).toEqual([]);
expect(parseSourceCsv('')).toEqual([]);
});
test('dedupes while preserving first-seen order', () => {
expect(parseSourceCsv('a,b,a,c,b')).toEqual(['a', 'b', 'c']);
});
});
describe('extractDryRun', () => {
test('absent flag → false', () => {
expect(extractDryRun(['alice', 'proj-x'])).toEqual({
dryRun: false,
rest: ['alice', 'proj-x'],
});
});
test('flag at end', () => {
expect(extractDryRun(['alice', 'proj-x', '--dry-run'])).toEqual({
dryRun: true,
rest: ['alice', 'proj-x'],
});
});
test('flag at start', () => {
expect(extractDryRun(['--dry-run', 'alice', 'proj-x'])).toEqual({
dryRun: true,
rest: ['alice', 'proj-x'],
});
});
test('flag in middle', () => {
expect(extractDryRun(['alice', '--dry-run', 'proj-x'])).toEqual({
dryRun: true,
rest: ['alice', 'proj-x'],
});
});
test('no args', () => {
expect(extractDryRun([])).toEqual({ dryRun: false, rest: [] });
});
});
// ---------------------------------------------------------------------------
// DB-coupled
// ---------------------------------------------------------------------------
async function seedSource(id: string): Promise<void> {
const sql = sqlQueryForEngine(engine);
await sql`INSERT INTO sources (id, name) VALUES (${id}, ${id}) ON CONFLICT (id) DO NOTHING`;
}
async function seedClient(name: string, federated: string[] = []): Promise<string> {
// Ensure write source FK is satisfied — every seeded client points at 'default'.
await seedSource('default');
const sql = sqlQueryForEngine(engine);
const clientId = `gbrain_cl_test_${name}_${Date.now()}_${Math.random().toString(16).slice(2, 8)}`;
const fedLit = pgArray(federated);
await sql`
INSERT INTO oauth_clients (client_id, client_name, client_secret_hash,
redirect_uris, grant_types, scope,
client_id_issued_at, source_id, federated_read)
VALUES (${clientId}, ${name}, ${'dummy-hash'},
${pgArray([])}, ${pgArray(['client_credentials'])}, ${'read'},
${Date.now()}, ${'default'}, ${fedLit})
`;
return clientId;
}
async function readFederated(clientId: string): Promise<string[]> {
const sql = sqlQueryForEngine(engine);
const rows = await sql`SELECT federated_read FROM oauth_clients WHERE client_id = ${clientId}`;
const fed = rows[0]?.federated_read;
return Array.isArray(fed) ? (fed as string[]).map(String) : [];
}
describe('resolveClient', () => {
test('matches by client_id', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const c = await resolveClient(sql, id);
expect(c.client_name).toBe('alice');
expect(c.federated_read).toEqual(['default']);
});
test('matches by client_name', async () => {
await seedSource('default');
await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const c = await resolveClient(sql, 'alice');
expect(c.client_name).toBe('alice');
});
test('errors loudly on no-match', async () => {
const sql = sqlQueryForEngine(engine);
await expect(resolveClient(sql, 'nobody')).rejects.toThrow(/No active OAuth client found/);
});
test('errors loudly on ambiguous client_name', async () => {
await seedSource('default');
await seedClient('bob', ['default']);
await seedClient('bob', ['default']);
const sql = sqlQueryForEngine(engine);
await expect(resolveClient(sql, 'bob')).rejects.toThrow(/Multiple active OAuth clients named/);
});
test('null source_id is preserved as null (legacy row tolerance)', async () => {
const sql = sqlQueryForEngine(engine);
const clientId = `gbrain_cl_test_null_${Date.now()}`;
await sql`
INSERT INTO oauth_clients (client_id, client_name, client_secret_hash,
redirect_uris, grant_types, scope,
client_id_issued_at, source_id, federated_read)
VALUES (${clientId}, ${'legacy'}, ${'dummy'},
${pgArray([])}, ${pgArray(['client_credentials'])}, ${'read'},
${Date.now()}, ${null}, ${pgArray([])})
`;
const c = await resolveClient(sql, clientId);
expect(c.source_id).toBeNull();
expect(c.federated_read).toEqual([]);
});
});
describe('assertSourceExists', () => {
test('passes when present', async () => {
await seedSource('proj-x');
const sql = sqlQueryForEngine(engine);
await expect(assertSourceExists(sql, 'proj-x')).resolves.toBeUndefined();
});
test('throws with paste-ready hint when missing', async () => {
const sql = sqlQueryForEngine(engine);
await expect(assertSourceExists(sql, 'ghost')).rejects.toThrow(
/Source "ghost" does not exist.*gbrain sources add ghost/s,
);
});
});
describe('grantReadCore', () => {
test('appends when not present and persists', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const outcome = await grantReadCore(sql, 'alice', 'proj-x');
expect(outcome.kind).toBe('updated');
if (outcome.kind === 'updated') {
expect(outcome.before).toEqual(['default']);
expect(outcome.after).toEqual(['default', 'proj-x']);
}
expect(await readFederated(id)).toEqual(['default', 'proj-x']);
});
test('is idempotent — second call is a noop, list unchanged', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
await grantReadCore(sql, 'alice', 'proj-x');
const outcome = await grantReadCore(sql, 'alice', 'proj-x');
expect(outcome.kind).toBe('noop');
if (outcome.kind === 'noop') {
expect(outcome.reason).toBe('already-granted');
}
expect(await readFederated(id)).toEqual(['default', 'proj-x']);
});
test('refuses unknown source (fails BEFORE mutating)', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
await expect(grantReadCore(sql, 'alice', 'ghost')).rejects.toThrow(/does not exist/);
expect(await readFederated(id)).toEqual(['default']);
});
test('refuses unknown client', async () => {
const sql = sqlQueryForEngine(engine);
await expect(grantReadCore(sql, 'nobody', 'whatever')).rejects.toThrow(/No active OAuth client found/);
});
test('accepts client_id resolution too', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
await grantReadCore(sql, id, 'proj-x');
expect(await readFederated(id)).toEqual(['default', 'proj-x']);
});
test('rejects malformed source_id BEFORE existence check (Codex finding #3)', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
// Even with a row in `sources` having a weird id, the validator at the
// boundary refuses. Closes the "manual SQL plants a row, CLI lets it
// become unmanageable in federated_read" vector.
await sql`INSERT INTO sources (id, name) VALUES (${'has,"weird"-bits'}, ${'weird'})`;
await expect(grantReadCore(sql, 'alice', 'has,"weird"-bits')).rejects.toThrow(/Invalid source_id/);
// DB unchanged.
expect(await readFederated(id)).toEqual(['default']);
});
});
describe('revokeReadCore', () => {
test('removes when present', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default', 'proj-x']);
const sql = sqlQueryForEngine(engine);
const outcome = await revokeReadCore(sql, 'alice', 'proj-x');
expect(outcome.kind).toBe('updated');
expect(await readFederated(id)).toEqual(['default']);
});
test('is idempotent — second call is a noop, list unchanged', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const outcome = await revokeReadCore(sql, 'alice', 'ghost-source');
expect(outcome.kind).toBe('noop');
if (outcome.kind === 'noop') {
expect(outcome.reason).toBe('not-present');
}
expect(await readFederated(id)).toEqual(['default']);
});
test('allows clearing the list down to empty (no implicit guard)', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
await revokeReadCore(sql, 'alice', 'default');
expect(await readFederated(id)).toEqual([]);
});
test('does NOT validate the source exists — operator may revoke stale references', async () => {
await seedSource('default');
// federated_read carries 'proj-x' but the source row was deleted.
const id = await seedClient('alice', ['default', 'proj-x']);
const sql = sqlQueryForEngine(engine);
const outcome = await revokeReadCore(sql, 'alice', 'proj-x');
expect(outcome.kind).toBe('updated');
expect(await readFederated(id)).toEqual(['default']);
});
});
describe('setFederatedReadCore', () => {
test('replaces list wholesale', async () => {
await seedSource('a');
await seedSource('b');
await seedSource('c');
const id = await seedClient('alice', ['a']);
const sql = sqlQueryForEngine(engine);
const outcome = await setFederatedReadCore(sql, 'alice', 'b,c');
expect(outcome.kind).toBe('updated');
expect(await readFederated(id)).toEqual(['b', 'c']);
});
test('dedupes CSV input', async () => {
await seedSource('a');
await seedSource('b');
const id = await seedClient('alice', []);
const sql = sqlQueryForEngine(engine);
await setFederatedReadCore(sql, 'alice', 'a,b,a,b,a');
expect(await readFederated(id)).toEqual(['a', 'b']);
});
test('empty string clears the list', async () => {
await seedSource('a');
const id = await seedClient('alice', ['a']);
const sql = sqlQueryForEngine(engine);
await setFederatedReadCore(sql, 'alice', '');
expect(await readFederated(id)).toEqual([]);
});
test('noop when result equals current list', async () => {
await seedSource('a');
await seedSource('b');
const id = await seedClient('alice', ['a', 'b']);
const sql = sqlQueryForEngine(engine);
const outcome = await setFederatedReadCore(sql, 'alice', 'a,b');
expect(outcome.kind).toBe('noop');
if (outcome.kind === 'noop') {
expect(outcome.reason).toBe('same-list');
}
expect(await readFederated(id)).toEqual(['a', 'b']);
});
test('refuses unknown source (fails BEFORE mutating)', async () => {
await seedSource('a');
const id = await seedClient('alice', ['a']);
const sql = sqlQueryForEngine(engine);
await expect(setFederatedReadCore(sql, 'alice', 'a,ghost')).rejects.toThrow(/does not exist/);
// Original list preserved.
expect(await readFederated(id)).toEqual(['a']);
});
test('order in CSV is the order persisted', async () => {
await seedSource('a');
await seedSource('b');
await seedSource('c');
const id = await seedClient('alice', ['a']);
const sql = sqlQueryForEngine(engine);
await setFederatedReadCore(sql, 'alice', 'c,a,b');
expect(await readFederated(id)).toEqual(['c', 'a', 'b']);
});
});
// ---------------------------------------------------------------------------
// --dry-run semantics
// ---------------------------------------------------------------------------
// ---------------------------------------------------------------------------
// Codex fixes: soft-delete filter, atomic-SQL race-safety, sanitizer
// ---------------------------------------------------------------------------
describe('sanitizeForTerminal', () => {
test('preserves printable ASCII unchanged', () => {
expect(sanitizeForTerminal('alice')).toBe('alice');
expect(sanitizeForTerminal('a b-c_d.e/f@g')).toBe('a b-c_d.e/f@g');
});
test('escapes ANSI escape sequences', () => {
expect(sanitizeForTerminal('\x1b[2J')).toBe('\\x1b[2J');
expect(sanitizeForTerminal('\x1b]0;TITLE\x07')).toBe('\\x1b]0;TITLE\\x07');
});
test('escapes ALL C0 controls including tab and newline', () => {
// Codex re-review: preserving \n lets a DCR-registered name spoof
// additional rows in list-clients output. Tab spoofs field separators.
// Both are now escaped.
expect(sanitizeForTerminal('\x00\x07\x08')).toBe('\\x00\\x07\\x08');
expect(sanitizeForTerminal('line1\nline2')).toBe('line1\\x0aline2');
expect(sanitizeForTerminal('col1\tcol2')).toBe('col1\\x09col2');
});
test('escapes DEL and C1 controls', () => {
expect(sanitizeForTerminal('\x7f')).toBe('\\x7f');
expect(sanitizeForTerminal('\x9b[31m')).toBe('\\x9b[31m');
});
test('passes through unicode', () => {
expect(sanitizeForTerminal('café')).toBe('café');
expect(sanitizeForTerminal('日本語')).toBe('日本語');
});
});
describe('soft-delete filter (Codex finding #2)', () => {
async function softDeleteClient(clientId: string): Promise<void> {
const sql = sqlQueryForEngine(engine);
await sql`UPDATE oauth_clients SET deleted_at = now() WHERE client_id = ${clientId}`;
}
test('resolveClient hides soft-deleted clients by default', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
await softDeleteClient(id);
const sql = sqlQueryForEngine(engine);
await expect(resolveClient(sql, 'alice')).rejects.toThrow(/No active OAuth client found/);
await expect(resolveClient(sql, id)).rejects.toThrow(/No active OAuth client found/);
});
test('resolveClient with includeDeleted finds soft-deleted clients', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
await softDeleteClient(id);
const sql = sqlQueryForEngine(engine);
const c = await resolveClient(sql, id, { includeDeleted: true });
expect(c.client_name).toBe('alice');
expect(c.deleted_at).not.toBeNull();
});
test('grantReadCore refuses to mutate soft-deleted clients', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default']);
await softDeleteClient(id);
const sql = sqlQueryForEngine(engine);
await expect(grantReadCore(sql, 'alice', 'proj-x')).rejects.toThrow(/No active OAuth client found/);
expect(await readFederated(id)).toEqual(['default']);
});
test('revokeReadCore refuses to mutate soft-deleted clients', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
await softDeleteClient(id);
const sql = sqlQueryForEngine(engine);
await expect(revokeReadCore(sql, 'alice', 'default')).rejects.toThrow(/No active OAuth client found/);
expect(await readFederated(id)).toEqual(['default']);
});
test('two clients with same name but only one active resolves to the active one', async () => {
await seedSource('default');
// Seed two clients with the same name; soft-delete the older one.
const sql = sqlQueryForEngine(engine);
const oldId = await seedClient('alice', ['default']);
await softDeleteClient(oldId);
const newId = await seedClient('alice', ['default']); // same name, new row
const c = await resolveClient(sql, 'alice');
expect(c.client_id).toBe(newId); // active row wins; ambiguity error suppressed
});
});
describe('atomic SQL race-safety (Codex finding #1, HIGH)', () => {
test('grant+revoke serialize at row-lock — sensitive stays revoked', async () => {
await seedSource('default');
await seedSource('sensitive');
await seedSource('harmless');
const id = await seedClient('alice', ['default', 'sensitive']);
const sql = sqlQueryForEngine(engine);
// Simulate concurrent revoke(sensitive) + grant(harmless). Real concurrency
// would race at the JS event loop boundary; here we await sequentially but
// each call goes through the ATOMIC SQL path. The contract: regardless of
// ordering, the final state has sensitive REMOVED and harmless ADDED.
await revokeReadCore(sql, 'alice', 'sensitive');
await grantReadCore(sql, 'alice', 'harmless');
const final1 = await readFederated(id);
expect(final1.sort()).toEqual(['default', 'harmless']);
// Reverse order, same final state. The pre-fix read-modify-write shape
// would have produced ['default', 'sensitive', 'harmless'] here (the
// resurrection bug Codex caught).
const id2 = await seedClient('bob', ['default', 'sensitive']);
await grantReadCore(sql, 'bob', 'harmless');
await revokeReadCore(sql, 'bob', 'sensitive');
const final2 = await readFederated(id2);
expect(final2.sort()).toEqual(['default', 'harmless']);
});
test('grant uses RETURNING to surface the post-write state', async () => {
await seedSource('default');
await seedSource('proj-x');
await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const outcome = await grantReadCore(sql, 'alice', 'proj-x');
expect(outcome.kind).toBe('updated');
if (outcome.kind === 'updated') {
// The `after` came from RETURNING, not from computing prev+sourceId
// in JS — proves the atomic path returned authoritative state.
expect(outcome.after).toEqual(['default', 'proj-x']);
}
});
test('grant noop path still survives without writing', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const outcome = await grantReadCore(sql, 'alice', 'default');
expect(outcome.kind).toBe('noop');
if (outcome.kind === 'noop') expect(outcome.reason).toBe('already-granted');
expect(await readFederated(id)).toEqual(['default']);
});
});
describe('dryRun mode', () => {
test('grantReadCore returns "updated" outcome but skips the write', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
const outcome = await grantReadCore(sql, 'alice', 'proj-x', { dryRun: true });
expect(outcome.kind).toBe('updated');
if (outcome.kind === 'updated') {
expect(outcome.before).toEqual(['default']);
expect(outcome.after).toEqual(['default', 'proj-x']);
}
// Crucially: the DB row is UNCHANGED.
expect(await readFederated(id)).toEqual(['default']);
});
test('revokeReadCore returns "updated" outcome but skips the write', async () => {
await seedSource('default');
await seedSource('proj-x');
const id = await seedClient('alice', ['default', 'proj-x']);
const sql = sqlQueryForEngine(engine);
const outcome = await revokeReadCore(sql, 'alice', 'proj-x', { dryRun: true });
expect(outcome.kind).toBe('updated');
expect(await readFederated(id)).toEqual(['default', 'proj-x']);
});
test('setFederatedReadCore returns "updated" outcome but skips the write', async () => {
await seedSource('a');
await seedSource('b');
await seedSource('c');
const id = await seedClient('alice', ['a']);
const sql = sqlQueryForEngine(engine);
const outcome = await setFederatedReadCore(sql, 'alice', 'b,c', { dryRun: true });
expect(outcome.kind).toBe('updated');
expect(await readFederated(id)).toEqual(['a']);
});
test('noop outcomes are surfaced identically with or without dryRun', async () => {
await seedSource('default');
await seedSource('proj-x');
await seedClient('alice', ['default', 'proj-x']);
const sql = sqlQueryForEngine(engine);
const live = await grantReadCore(sql, 'alice', 'proj-x', { dryRun: false });
const dry = await grantReadCore(sql, 'alice', 'proj-x', { dryRun: true });
expect(live.kind).toBe('noop');
expect(dry.kind).toBe('noop');
});
test('errors still fire in dryRun (operator sees the problem before commit)', async () => {
await seedSource('default');
const id = await seedClient('alice', ['default']);
const sql = sqlQueryForEngine(engine);
await expect(
grantReadCore(sql, 'alice', 'ghost', { dryRun: true }),
).rejects.toThrow(/does not exist/);
await expect(
grantReadCore(sql, 'nobody', 'default', { dryRun: true }),
).rejects.toThrow(/No active OAuth client found/);
// DB unchanged.
expect(await readFederated(id)).toEqual(['default']);
});
});
+23
View File
@@ -174,6 +174,29 @@ describe('parseRegisterClientArgs', () => {
});
describe('error cases', () => {
test('--source with malformed id throws (validates source_id shape — codex re-review)', () => {
// Defense for the "register-client seeds an unmanageable
// federated_read entry" vector. assertValidSourceId fires before the
// function returns so DB never sees a row with bad source scope.
expect(() => parseRegisterClientArgs(['--source', 'has,weird,bits'])).toThrow(/Invalid source_id/);
expect(() => parseRegisterClientArgs(['--source', 'UPPER'])).toThrow(/Invalid source_id/);
expect(() => parseRegisterClientArgs(['--source', ''])).toThrow(/Invalid source_id|requires a value/);
});
test('--federated-read with any malformed id throws', () => {
// Single-item bad.
expect(() => parseRegisterClientArgs(['--federated-read', 'bad,source!'])).toThrow(/Invalid source_id/);
// Mixed valid + invalid — fails on the first bad one.
expect(() => parseRegisterClientArgs(['--federated-read', 'good,bad source'])).toThrow(/Invalid source_id/);
});
test('--source default + --federated-read default,team passes (regression — common case)', () => {
// Sanity: the canonical real-world invocation still parses cleanly.
const out = parseRegisterClientArgs(['--source', 'default', '--federated-read', 'default,team']);
expect(out.sourceId).toBe('default');
expect(out.federatedRead).toEqual(['default', 'team']);
});
test('--redirect-uri without value → throws', () => {
expect(() => parseRegisterClientArgs(['--redirect-uri'])).toThrow(/requires a value/);
});
-161
View File
@@ -1,161 +0,0 @@
/**
* #1818: embedBatch dispatches its 100-input sub-batches through a bounded
* worker pool (the embed-stale.ts concurrency pattern) instead of a serial
* `for` loop. This file pins:
*
* - output order matches input order regardless of completion order
* (index-addressed results)
* - parallelism actually happens (max in-flight > 1) and stays bounded
* (max in-flight <= configured concurrency)
* - concurrency: 1 restores the serial pre-#1818 dispatch
* - GBRAIN_EMBED_BATCH_CONCURRENCY env is honored when the option is unset
* - onBatchComplete reports a monotonic completed count ending at total
*
* Transport is stubbed via the gateway's __setEmbedTransportForTests seam
* (same pattern as test/ai/adaptive-embed-batch.test.ts). OpenAI recipe =
* fast path (no pre-split), so each embedBatch sub-batch is exactly one
* transport call.
*/
import { afterAll, afterEach, beforeEach, describe, expect, test } from 'bun:test';
import {
configureGateway,
resetGateway,
__setEmbedTransportForTests,
} from '../src/core/ai/gateway.ts';
import { embedBatch } from '../src/core/embedding.ts';
import { withEnv } from './helpers/with-env.ts';
const DIMS = 1536;
function configureOpenAI(): void {
configureGateway({
embedding_model: 'openai:text-embedding-3-large',
embedding_dimensions: DIMS,
env: { OPENAI_API_KEY: 'sk-fake' },
});
}
/**
* Install a transport whose returned embedding encodes the GLOBAL input
* index in dim 0 (texts are `t<N>`), so order can be asserted end-to-end.
* Tracks the max number of concurrently in-flight transport calls.
*/
function installTrackingTransport(delayMs = 5): { maxInFlight: () => number } {
let inFlight = 0;
let maxInFlight = 0;
__setEmbedTransportForTests((async ({ values }: { values: string[] }) => {
inFlight++;
maxInFlight = Math.max(maxInFlight, inFlight);
await new Promise(r => setTimeout(r, delayMs));
inFlight--;
return {
embeddings: values.map(v => {
const idx = Number(v.slice(1));
return Array.from({ length: DIMS }, (_, j) => (j === 0 ? idx : 0.1));
}),
};
}) as any);
return { maxInFlight: () => maxInFlight };
}
const texts = Array.from({ length: 250 }, (_, i) => `t${i}`);
afterAll(() => resetGateway());
describe('embedBatch bounded parallelism (#1818)', () => {
beforeEach(() => {
resetGateway();
configureOpenAI();
});
afterEach(() => {
__setEmbedTransportForTests(null);
});
test('default pool dispatches sub-batches in parallel, order preserved', async () => {
const tracker = installTrackingTransport();
const result = await embedBatch(texts, { onBatchComplete: () => {} });
expect(result).toHaveLength(250);
for (let i = 0; i < 250; i++) {
expect(result[i][0]).toBe(i);
}
// 250 texts → 3 sub-batches; default concurrency 4 → all 3 in flight.
expect(tracker.maxInFlight()).toBeGreaterThan(1);
expect(tracker.maxInFlight()).toBeLessThanOrEqual(4);
});
test('concurrency: 1 keeps the serial dispatch', async () => {
const tracker = installTrackingTransport();
const result = await embedBatch(texts, { concurrency: 1, onBatchComplete: () => {} });
expect(result).toHaveLength(250);
expect(tracker.maxInFlight()).toBe(1);
});
test('GBRAIN_EMBED_BATCH_CONCURRENCY env bounds the pool when option unset', async () => {
const tracker = installTrackingTransport();
await withEnv({ GBRAIN_EMBED_BATCH_CONCURRENCY: '2' }, async () => {
await embedBatch(texts, { onBatchComplete: () => {} });
});
expect(tracker.maxInFlight()).toBeGreaterThan(1);
expect(tracker.maxInFlight()).toBeLessThanOrEqual(2);
});
test('onBatchComplete reports a monotonic count ending at total', async () => {
installTrackingTransport();
const seen: number[] = [];
await embedBatch(texts, {
onBatchComplete: (done, total) => {
expect(total).toBe(250);
seen.push(done);
},
});
expect(seen).toHaveLength(3); // 100 + 100 + 50 sub-batches
for (let i = 1; i < seen.length; i++) {
expect(seen[i]).toBeGreaterThan(seen[i - 1]);
}
expect(seen[seen.length - 1]).toBe(250);
});
test('a failing sub-batch rejects the whole call', async () => {
let call = 0;
__setEmbedTransportForTests((async ({ values }: { values: string[] }) => {
call++;
if (call === 2) throw new Error('boom');
await new Promise(r => setTimeout(r, 2));
return { embeddings: values.map(() => Array.from({ length: DIMS }, () => 0.1)) };
}) as any);
await expect(embedBatch(texts, { onBatchComplete: () => {} })).rejects.toThrow();
});
test('after a failure, surviving workers stop dispatching new slices', async () => {
// 1000 texts → 10 slices, concurrency 2. First call fails immediately;
// without the `failed` flag the second worker would keep draining all
// 10 slices in the background AFTER embedBatch already rejected —
// burning provider spend and firing onBatchComplete post-rejection.
let calls = 0;
const completions: number[] = [];
__setEmbedTransportForTests((async ({ values }: { values: string[] }) => {
calls++;
if (calls === 1) throw new Error('boom');
await new Promise(r => setTimeout(r, 5));
return { embeddings: values.map(() => Array.from({ length: DIMS }, () => 0.1)) };
}) as any);
const many = Array.from({ length: 1000 }, (_, i) => `t${i}`);
await expect(
embedBatch(many, { concurrency: 2, onBatchComplete: d => completions.push(d) }),
).rejects.toThrow('boom');
const callsAtRejection = calls;
await new Promise(r => setTimeout(r, 50)); // would-be background drain window
expect(calls).toBe(callsAtRejection); // no new dispatch after rejection
expect(calls).toBeLessThanOrEqual(2); // only the in-flight sibling ran
expect(completions).toHaveLength(0); // no progress reported after failure
});
test('single small batch without callback stays on the one-call fast path', async () => {
const tracker = installTrackingTransport(1);
const result = await embedBatch(['t0', 't1', 't2']);
expect(result).toHaveLength(3);
expect(result[1][0]).toBe(1);
expect(tracker.maxInFlight()).toBe(1);
});
});
+3 -27
View File
@@ -19,7 +19,7 @@
* overwrites this preload.
*/
import { configureGateway, getEmbeddingDimensions } from '../../src/core/ai/gateway.ts';
import { afterEach, beforeEach } from 'bun:test';
import { beforeEach } from 'bun:test';
const LEGACY_CONFIG = {
embedding_model: 'openai:text-embedding-3-large',
@@ -52,7 +52,7 @@ applyLegacy();
// 2. file-local beforeAll → may overwrite to ZE/1280
// Since beforeAll runs once per file BEFORE the first beforeEach,
// file-local beforeAll wins for that file's tests. ✓
function applyLegacyIfEmpty() {
beforeEach(() => {
try {
// Only re-apply if the gateway was reset (or never configured).
// Tests that explicitly configured a different model in their
@@ -62,28 +62,4 @@ function applyLegacyIfEmpty() {
} catch {
applyLegacy();
}
}
beforeEach(applyLegacyIfEmpty);
// PR #3130 shard-order fix: beforeEach alone leaves ONE window open — a file
// whose LAST afterEach calls resetGateway() poisons the NEXT file's
// beforeAll, which runs BEFORE any beforeEach fires. A beforeAll there that
// does engine.initSchema() then sizes the embedding column from the gateway
// DEFAULTS (zembed-1/1280d) instead of the pinned legacy 1536, and every
// 1536-d Float32Array fixture in that file dies with
// "expected 1280 dimensions, not 1536". Which file pair collides is a
// function of shard composition, so adding/removing ANY test file can
// surface it (that is exactly how it bit shard 9).
//
// Preload hooks are registered before any file-local hooks, and bun runs
// after-hooks inside-out (file-local afterEach first, then this one), so
// this repairs the empty slot immediately after the poisoning reset —
// before the next file's beforeAll can observe it.
//
// Known remaining window: a file whose afterAll() resets the gateway (no
// hook runs between its afterAll and the next file's beforeAll). Files
// that reset in afterAll and can precede a schema-creating file should
// re-apply their own config, or the victim file should configureGateway()
// explicitly in its beforeAll.
afterEach(applyLegacyIfEmpty);
});
-69
View File
@@ -1,69 +0,0 @@
/**
* #1207: `gbrain import` without `--workers` used to hardcode workerCount=1,
* so a large Postgres import paid one serial embedding round-trip per file.
* runImport now routes the default through the shared autoConcurrency policy
* (PGLite 1, >100 files on Postgres DEFAULT_PARALLEL_WORKERS), while an
* explicit `--workers N` still wins.
*
* The engine here is a minimal postgres-kind stub with no database_url in
* config runImport's parallel branch then falls back to serial processing
* (its PR #490 guard) but the WORKER-COUNT DECISION (the thing #1207 fixes)
* is still observable via the "Using N parallel workers" log line. Per-file
* imports fail against the stub engine and are swallowed by runImport's
* per-file catch; that's fine this test pins the policy, not the import.
*/
import { afterEach, beforeEach, describe, expect, test } from 'bun:test';
import { mkdtempSync, writeFileSync, mkdirSync, rmSync, realpathSync } from 'fs';
import { tmpdir } from 'os';
import { join } from 'path';
import { withEnv } from './helpers/with-env.ts';
import { runImport } from '../src/commands/import.ts';
const fakePostgresEngine = {
kind: 'postgres',
executeRaw: async () => [],
logIngest: async () => {},
setConfig: async () => {},
getConfig: async () => null,
} as any;
let workspace: string;
let brainDir: string;
let logs: string[];
const realLog = console.log;
beforeEach(() => {
workspace = mkdtempSync(join(tmpdir(), 'gbrain-import-workers-home-'));
mkdirSync(join(workspace, '.gbrain'), { recursive: true });
brainDir = realpathSync(mkdtempSync(join(tmpdir(), 'gbrain-import-workers-brain-')));
// 101 files: one past AUTO_CONCURRENCY_FILE_THRESHOLD (100).
for (let i = 0; i < 101; i++) {
writeFileSync(join(brainDir, `page-${i}.md`), `# Page ${i}\n\nbody ${i}\n`);
}
logs = [];
console.log = (msg?: unknown) => logs.push(String(msg));
});
afterEach(() => {
console.log = realLog;
rmSync(workspace, { recursive: true, force: true });
rmSync(brainDir, { recursive: true, force: true });
});
describe('import default worker count (#1207)', () => {
test('no --workers flag → autoConcurrency picks 4 for >100 files on Postgres', async () => {
await withEnv({ GBRAIN_HOME: join(workspace, '.gbrain'), GBRAIN_SOURCE: undefined }, async () => {
await runImport(fakePostgresEngine, [brainDir, '--no-embed'], { sourceId: 'default' });
});
expect(logs.some(l => l.includes('Using 4 parallel workers'))).toBe(true);
});
test('explicit --workers 2 still wins over the auto policy', async () => {
await withEnv({ GBRAIN_HOME: join(workspace, '.gbrain'), GBRAIN_SOURCE: undefined }, async () => {
await runImport(fakePostgresEngine, [brainDir, '--no-embed', '--workers', '2'], { sourceId: 'default' });
});
expect(logs.some(l => l.includes('Using 2 parallel workers'))).toBe(true);
expect(logs.some(l => l.includes('Using 4 parallel workers'))).toBe(false);
});
});