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 1944 additions and 637 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
`);
}
+13 -43
View File
@@ -1,7 +1,6 @@
import type { BrainEngine } from '../core/engine.ts';
import { embedBatch, currentEmbeddingSignature } from '../core/embedding.ts';
import type { ChunkInput, ResolvedColumn } from '../core/types.ts';
import { resolveWriteColumnForEngine } from '../core/search/embedding-column.ts';
import type { ChunkInput } from '../core/types.ts';
import { chunkText } from '../core/chunkers/recursive.ts';
import { createProgress, type ProgressReporter } from '../core/progress.ts';
import { getCliOptions, cliOptsToProgressOptions } from '../core/cli-options.ts';
@@ -184,13 +183,8 @@ export class EmbeddingDimMismatchError extends Error {
* fresh-install bug class at the very first invocation instead of letting
* the worker pool hammer N pages with raw 22000 errors.
*/
async function preflightDimMismatch(engine: BrainEngine, dryRun: boolean, embeddingColumn?: ResolvedColumn): Promise<void> {
async function preflightDimMismatch(engine: BrainEngine, dryRun: boolean): Promise<void> {
if (dryRun) return; // dry-run never embeds, no risk
// #1262: an alt-column brain writes to `embeddingColumn`, not the legacy
// `embedding` column — the legacy column's dims are irrelevant, and the
// registry entry (validated at resolve time) pins the target's dims. Only
// the legacy default path needs the schema-vs-gateway dim comparison.
if (embeddingColumn && embeddingColumn.name !== 'embedding') return;
const { readContentChunksEmbeddingDim, embeddingMismatchMessage } = await import('../core/embedding-dim-check.ts');
const { getEmbeddingDimensions, getEmbeddingModel } = await import('../core/ai/gateway.ts');
let existing;
@@ -244,12 +238,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
// v0.37.11.0 (Lane D.2): pre-flight dim-mismatch check. Catches the headline
// fresh-install bug class before the worker pool spends 20 parallel calls
// hitting raw Postgres dimension errors.
// #1262: resolve the write-side embedding column ONCE at the boundary
// (merged config + gateway model) and thread the descriptor through every
// upsertChunks / stale-scan below. undefined => legacy `embedding` column.
const embeddingColumn = await resolveWriteColumnForEngine(engine);
await preflightDimMismatch(engine, !!opts.dryRun, embeddingColumn);
await preflightDimMismatch(engine, !!opts.dryRun);
const result: EmbedResult = {
embedded: 0,
@@ -264,7 +253,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
for (const s of opts.slugs) {
if (isAborted(opts.signal)) break; // #1737: stop the per-slug loop on abort
try {
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal, embeddingColumn);
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal);
} catch (e: unknown) {
serr(` Error embedding ${s}: ${e instanceof Error ? e.message : e}`);
}
@@ -358,7 +347,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
catchUp: opts.catchUp,
pacer,
paceMaxConcurrency,
}, opts.signal, embeddingColumn);
}, opts.signal);
} finally {
// E1: surface pacing telemetry (human + structured) when pacing was on.
const snap = pacer.snapshot();
@@ -387,7 +376,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
return result;
}
if (opts.slug) {
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal, embeddingColumn);
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal);
return result;
}
throw new Error('No embed target specified. Pass { slug }, { slugs }, { all }, or { stale }.');
@@ -532,13 +521,8 @@ async function embedPage(
result: EmbedResult,
sourceId?: string,
signal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
const opts = sourceId ? { sourceId } : undefined;
// #1262: write-side descriptor rides only on WRITE calls (upsertChunks).
const chunkOpts = (sourceId || embeddingColumn)
? { ...(sourceId && { sourceId }), ...(embeddingColumn && { embeddingColumn }) }
: undefined;
const page = await engine.getPage(slug, opts);
if (!page) {
throw new Error(`Page not found: ${slug}`);
@@ -570,7 +554,7 @@ async function embedPage(
}
if (inputs.length > 0) {
await engine.upsertChunks(slug, inputs, chunkOpts);
await engine.upsertChunks(slug, inputs, opts);
chunks = await engine.getChunks(slug, opts);
}
}
@@ -605,7 +589,7 @@ async function embedPage(
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await engine.upsertChunks(slug, updated, chunkOpts);
await engine.upsertChunks(slug, updated, opts);
// v0.41.31: stamp provenance so a later model/dims swap is detectable as
// stale. embedPage is the per-slug path used by `gbrain embed <slug>` AND
// by `gbrain sync`'s post-import embed step (runEmbedCore({slugs})).
@@ -638,7 +622,6 @@ async function embedAll(
paceMaxConcurrency?: number;
},
signal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
// v0.41.31: current embedding provenance signature. Stamped onto pages
// when their chunks are (re)embedded so a later model/dimension swap is
@@ -661,7 +644,7 @@ async function embedAll(
// D7: thread sourceId so `gbrain embed --stale --source X` actually scopes.
// v0.41.18.0 (A13): thread batchSize/priority/catchUp into the stale path.
// #1737: thread the external abort signal so the cycle embed phase bails.
return await embedAllStale(engine, sourceId, dryRun, result, onProgress, staleOpts, signature, signal, embeddingColumn);
return await embedAllStale(engine, sourceId, dryRun, result, onProgress, staleOpts, signature, signal);
}
// --all path: pacer (no-op when off). E-1: lower the worker count to the
@@ -742,10 +725,7 @@ async function embedAll(
embedding: embeddingMap.get(c.chunk_index) ?? undefined,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, {
...(pageSourceId && { sourceId: pageSourceId }),
...(embeddingColumn && { embeddingColumn }),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, pageOpts));
// v0.41.31: stamp embedding provenance so a later model swap is
// detectable as stale.
await observed(pacer, () =>
@@ -825,16 +805,10 @@ async function embedAllStale(
},
signature?: string,
externalSignal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
// D7: thread sourceId so source-scoped runs only count + visit
// that source's NULL embeddings.
// #1262: the stale predicate follows the write-side column — without it an
// alt-column brain would perpetually re-select (and re-pay for) chunks whose
// target column is already populated.
const sourceOpt = (sourceId || embeddingColumn)
? { ...(sourceId && { sourceId }), ...(embeddingColumn && { embeddingColumn }) }
: undefined;
const sourceOpt = sourceId ? { sourceId } : undefined;
// v0.41.31: re-embed pages whose embedding_signature drifted (model/dims
// swap). dry-run must NOT mutate, so it counts signature-stale via the
@@ -993,7 +967,6 @@ async function embedAllStale(
afterUpdatedAt,
}),
...(sourceId && { sourceId }),
...(embeddingColumn && { embeddingColumn }),
}),
);
if (batch.length === 0) {
@@ -1046,10 +1019,7 @@ async function embedAllStale(
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, {
sourceId: keySourceId,
...(embeddingColumn && { embeddingColumn }),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
// v0.41.31: stamp provenance after the page's chunks are embedded —
// but only when EVERY chunk was stale (fully re-embedded this pass).
// A partially-stale page keeps preserved chunks of unknown/old
@@ -1120,7 +1090,7 @@ async function embedAllStale(
// as a clean run — re-running won't help until the underlying failure is fixed.
if (staleOpts?.catchUp && !effectiveSignal.aborted && embedFailures > 0) {
const remaining = await engine.countStaleChunks(
signature ? { signature, ...sourceOpt } : sourceOpt,
signature ? { signature, ...(sourceId ? { sourceId } : {}) } : (sourceId ? { sourceId } : undefined),
);
if (remaining > 0) {
serr(`\n [embed] catch-up finished but ${remaining} chunk(s) remain stale after ${embedFailures} embed failure(s). These are not embeddable as-is; re-running won't clear them until the underlying error is resolved.`);
+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
// ---------------------------------------------------------------------------
+1 -8
View File
@@ -576,16 +576,9 @@ async function runInlineCostGate(
// Stale backlog: cheap single SQL; fail-open to 0 so a transient DB hiccup
// never blocks the sync. Signature-aware (model/dims swap surfaces here).
// #1262: follow the write-side embedding column — otherwise an alt-column
// brain's fully-embedded corpus counts as phantom backlog on every gate.
let staleChars = 0;
try {
const { resolveWriteColumnForEngine } = await import('../core/search/embedding-column.ts');
const embeddingColumn = await resolveWriteColumnForEngine(engine);
staleChars = await engine.sumStaleChunkChars({
signature: currentEmbeddingSignature(),
...(embeddingColumn && { embeddingColumn }),
});
staleChars = await engine.sumStaleChunkChars({ signature: currentEmbeddingSignature() });
} catch {
staleChars = 0;
}
-5
View File
@@ -61,7 +61,6 @@ import {
type SynopsisFailureKind,
} from './audit-synopsis.ts';
import type { BrainEngine } from './engine.ts';
import { resolveWriteColumnForEngine } from './search/embedding-column.ts';
import type { ChunkInput, CRMode, Page } from './types.ts';
import type { SourceRow } from './sources-ops.ts';
@@ -287,13 +286,9 @@ export async function reembedPageWithContextualRetrieval(
// ── PHASE 2: single DB transaction ───────────────────────────
try {
// #1262: contextual re-embeds write TEXT embeddings — thread the
// caller-resolved write column like every other embed path.
const embeddingColumn = await resolveWriteColumnForEngine(args.engine);
await args.engine.transaction(async (tx) => {
await tx.upsertChunks(args.pageSlug, phase1.embeddedChunks, {
sourceId: args.sourceId,
...(embeddingColumn && { embeddingColumn }),
});
await tx.updatePageContextualRetrievalState(
args.pageSlug,
+2 -13
View File
@@ -18,7 +18,7 @@
*/
import type { BrainEngine } from './engine.ts';
import type { ChunkInput, ResolvedColumn } from './types.ts';
import type { ChunkInput } from './types.ts';
import { embedBatchWithBackoff } from '../commands/embed.ts';
import { type DbPacer, createNoopPacer, observed } from './db-pacer.ts';
import { AbortError } from './abort-check.ts';
@@ -61,13 +61,6 @@ export interface EmbedStaleOpts {
* Omit to keep the legacy `embedding IS NULL`-only behavior.
*/
embeddingSignature?: string;
/**
* #1262: caller-resolved write-side embedding column. Threaded into BOTH
* listStaleChunks (staleness predicate) and upsertChunks (write target) so
* an alt-column brain converges instead of re-selecting embedded rows.
* Resolve at the boundary via `resolveWriteColumnForEngine()`.
*/
embeddingColumn?: ResolvedColumn;
/**
* DB-contention pacer (paced-backfill). When enabled it (a) supplies the
* worker count via the caller passing `concurrency = bundle.maxConcurrency`
@@ -163,7 +156,6 @@ export async function embedStaleForSource(
afterPageId,
afterChunkIndex,
sourceId,
...(opts.embeddingColumn && { embeddingColumn: opts.embeddingColumn }),
}),
);
if (batch.length === 0) {
@@ -231,10 +223,7 @@ export async function embedStaleForSource(
doc_comment: c.doc_comment ?? undefined,
symbol_name_qualified: c.symbol_name_qualified ?? undefined,
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, {
sourceId: keySourceId,
...(opts.embeddingColumn && { embeddingColumn: opts.embeddingColumn }),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
// v0.41.31: stamp provenance only when EVERY chunk was stale (fully
// re-embedded this pass) — a partially-stale page keeps preserved
// chunks of unknown provenance, so don't claim current. After the
+3 -22
View File
@@ -12,7 +12,6 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
CodeEdgeInput, CodeEdgeResult,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
@@ -988,13 +987,8 @@ export interface BrainEngine {
* Postgres rolls back automatically on conn drop, so commit-ambiguous
* failure replays to the same end state. Callers MUST NOT wrap externally;
* see {@link BatchOpts} retry-contract block.
*
* `opts.embeddingColumn` (optional) selects the content_chunks column that
* receives TEXT embeddings (#1262). The caller resolves the descriptor at
* the import/embed boundary via `resolveWriteColumn()`; engines never read
* config or choose columns themselves. Omitted => legacy `embedding`.
*/
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void>;
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void>;
/**
* Read every chunk for a page. `opts.sourceId` source-scopes the page
* lookup; without it, multi-source brains return chunks from every
@@ -1011,13 +1005,8 @@ export interface BrainEngine {
* counts across every source in the brain. Operators running
* `gbrain embed --stale --source media-corpus` expect only that
* source's NULLs touched; the caller threads `sourceId` here.
*
* `opts.embeddingColumn` switches the staleness predicate from the legacy
* `embedding` column to the resolved write-side column, so alt-column
* brains do not perpetually re-select rows whose target column is already
* populated (#1262). Must match the eventual upsertChunks target.
*/
countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number>;
countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number>;
/**
* Sum of LENGTH(chunk_text) over stale chunks the character-count
* backlog the embed phase / embed-backfill will process. Sibling of
@@ -1031,13 +1020,8 @@ export interface BrainEngine {
* model signature (a model/dims swap). NULL signature is GRANDFATHERED
* (never counted) so the post-migration corpus isn't flagged en masse.
* Omit `signature` for the legacy `embedding IS NULL`-only count.
*
* `opts.embeddingColumn` switches the staleness predicate to the resolved
* write-side column (#1262) same contract as countStaleChunks so the
* sync cost gate doesn't count an alt-column brain's fully-embedded corpus
* as phantom backlog.
*/
sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number>;
sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number>;
/**
* Stamp `pages.embedding_signature = signature` for one page. Called after
* a page's chunks are (re)embedded so a later model swap can detect it as
@@ -1085,9 +1069,6 @@ export interface BrainEngine {
// both round-trip TIMESTAMPTZ as Date | string; ISO string is the
// common denominator on the wire).
afterUpdatedAt?: string | null;
// #1262: staleness predicate targets this column when set (must match
// countStaleChunks and the eventual upsertChunks write target).
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]>;
/**
* Delete every chunk for a page. Internal page-id lookup is sourceId-scoped
+4 -25
View File
@@ -10,8 +10,7 @@ import { findChunkForOffset } from './chunkers/edge-extractor.ts';
import { extractCodeRefs, imageOfCandidates } from './link-extraction.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature } from './embedding.ts';
import { slugifyPath, slugifyCodePath, isCodeFilePath } from './sync.ts';
import type { ChunkInput, PageInput, PageType, ResolvedColumn } from './types.ts';
import { resolveWriteColumnForEngine } from './search/embedding-column.ts';
import type { ChunkInput, PageInput, PageType } from './types.ts';
import { computeEffectiveDate } from './effective-date.ts';
import { MARKDOWN_CHUNKER_VERSION } from './chunkers/recursive.ts';
import { logSlugFallback } from './audit-slug-fallback.ts';
@@ -741,14 +740,6 @@ export async function importFromContent(
// schema DEFAULT — required for multi-source brains; harmless ('default')
// for single-source callers.
const txOpts = sourceId ? { sourceId } : undefined;
// #1262: resolve the write-side embedding column once (merged config +
// gateway model) BEFORE the transaction; the descriptor rides only on
// upsertChunks so text embeddings land in the registered column.
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(sourceId || chunkWriteColumn)
? { ...(sourceId && { sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
await engine.transaction(async (tx) => {
if (existing) await tx.createVersion(slug, txOpts);
@@ -833,7 +824,7 @@ export async function importFromContent(
}
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks, chunkOpts);
await tx.upsertChunks(slug, chunks, txOpts);
// v0.41.31: stamp embedding provenance when this import actually
// embedded (not --no-embed), so a later model/dims swap is detectable
// as stale via embed --stale. The deferred/backfill + per-slug embed
@@ -1073,12 +1064,6 @@ export async function importCodeFile(
const title = `${relativePath} (${lang})`;
const sourceId = opts.sourceId;
const txOpts = sourceId ? { sourceId } : undefined;
// #1262: write-side embedding column descriptor (rides only on upsertChunks).
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(sourceId || chunkWriteColumn)
? { ...(sourceId && { sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
const byteLength = Buffer.byteLength(content, 'utf-8');
if (byteLength > MAX_FILE_SIZE) {
@@ -1198,7 +1183,7 @@ export async function importCodeFile(
await tx.addTag(slug, lang, txOpts);
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks, chunkOpts);
await tx.upsertChunks(slug, chunks, txOpts);
// v0.41.31: stamp embedding provenance ONLY when every chunk was
// freshly embedded with the current model this call (no reuse-by-hash
// carrying old-model vectors). Mixed pages stay unstamped rather than
@@ -1347,12 +1332,6 @@ export async function withImportTransaction(
): Promise<void> {
const sourceId = spec.sourceId ?? 'default';
const txOpts = spec.sourceId ? { sourceId: spec.sourceId } : undefined;
// #1262: write-side embedding column descriptor (rides only on upsertChunks).
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(spec.sourceId || chunkWriteColumn)
? { ...(spec.sourceId && { sourceId: spec.sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
await engine.transaction(async (tx) => {
if (spec.hadExisting) await tx.createVersion(spec.slug, txOpts);
await tx.putPage(spec.slug, spec.page, txOpts);
@@ -1368,7 +1347,7 @@ export async function withImportTransaction(
}
if (spec.chunks !== undefined) {
if (spec.chunks.length > 0) {
await tx.upsertChunks(spec.slug, spec.chunks, chunkOpts);
await tx.upsertChunks(spec.slug, spec.chunks, txOpts);
} else {
await tx.deleteChunks(spec.slug, txOpts);
}
@@ -35,7 +35,6 @@ import { tryAcquireDbLock } from '../../db-lock.ts';
import { BudgetTracker, BudgetExhausted } from '../../budget/budget-tracker.ts';
import { withBudgetTracker } from '../../ai/gateway.ts';
import { embedStaleForSource } from '../../embed-stale.ts';
import { resolveWriteColumnForEngine } from '../../search/embedding-column.ts';
import { currentEmbeddingSignature } from '../../embedding.ts';
import { type DbPacer, createDbPacer, createNoopPacer } from '../../db-pacer.ts';
import { resolvePaceMode, loadPaceModeConfig, readPaceEnv } from '../../pace-mode.ts';
@@ -165,16 +164,12 @@ export function makeEmbedBackfillHandler(engine: BrainEngine) {
// the supervisor, so pacing it is the headline win.
const { pacer, concurrency } = await resolveBackfillPacer(engine, job.data);
// #1262: resolve the write-side embedding column once at the job boundary.
const embeddingColumn = await resolveWriteColumnForEngine(engine);
try {
const result = await withBudgetTracker(tracker, async () =>
embedStaleForSource(engine, sourceId, {
batchSize,
signal: job.signal,
pacer,
...(embeddingColumn && { embeddingColumn }),
...(concurrency !== undefined && { concurrency }),
// v0.41.31: re-embed pages whose model signature drifted + stamp
// provenance as chunks land.
+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(',')}}`;
+22 -42
View File
@@ -40,7 +40,6 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
SalienceOpts, SalienceResult, AnomaliesOpts, AnomalyResult,
@@ -2231,20 +2230,12 @@ export class PGLiteEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void> {
return this.batchRetry(opts?.auditSite ?? 'upsertChunks', opts?.signal, () => this._upsertChunksOnce(slug, chunks, opts), chunks.length);
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn }): Promise<void> {
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// #1262: caller-resolved write target for TEXT embeddings. Descriptor
// names are identifier-validated + quoted by buildVectorCastFragment;
// omitted => legacy `embedding vector`. Mirrors postgres-engine.ts.
const targetFragment = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn)
: undefined;
const targetCol = targetFragment?.col ?? 'embedding';
const embeddingCast = targetFragment?.castSql.replace('$1::', '') ?? 'vector';
// Source-scope the page-id lookup so duplicate slugs in different sources
// do not return multiple rows or target the wrong page.
@@ -2279,7 +2270,7 @@ export class PGLiteEngine implements BrainEngine {
// list. Image chunks pass embedding=null + embedding_image=Float32Array
// (1024-dim Voyage). Text/code chunks pass embedding=Float32Array +
// embedding_image=null. Default modality='text' when omitted.
const cols = `(page_id, chunk_index, chunk_text, chunk_source, ${targetCol}, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)`;
const cols = '(page_id, chunk_index, chunk_text, chunk_source, embedding, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)';
const rowParts: string[] = [];
const params: unknown[] = [];
let paramIdx = 1;
@@ -2297,7 +2288,7 @@ export class PGLiteEngine implements BrainEngine {
const modality = chunk.modality ?? 'text';
// Inline ::vector NULL literals to avoid a per-branch placeholder.
const embeddingPh = embeddingStr ? `$${paramIdx++}::${embeddingCast}` : 'NULL';
const embeddingPh = embeddingStr ? `$${paramIdx++}::vector` : 'NULL';
const embeddedAtPh = embeddingStr ? 'now()' : 'NULL';
const embeddingImagePh = embeddingImageStr ? `$${paramIdx++}::vector` : 'NULL';
@@ -2336,19 +2327,19 @@ export class PGLiteEngine implements BrainEngine {
ON CONFLICT (page_id, chunk_index) DO UPDATE SET
chunk_text = EXCLUDED.chunk_text,
chunk_source = EXCLUDED.chunk_source,
${targetCol} = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.${targetCol}
WHEN content_chunks.${targetCol} IS NULL THEN EXCLUDED.${targetCol}
embedding = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.embedding
WHEN content_chunks.embedding IS NULL THEN EXCLUDED.embedding
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.${targetCol}
ELSE content_chunks.${targetCol}
THEN EXCLUDED.embedding
ELSE content_chunks.embedding
END,
model = COALESCE(EXCLUDED.model, content_chunks.model),
token_count = EXCLUDED.token_count,
embedded_at = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.${targetCol} IS NULL THEN NULL
WHEN content_chunks.${targetCol} IS NULL AND EXCLUDED.${targetCol} IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.embedding IS NULL THEN NULL
WHEN content_chunks.embedding IS NULL AND EXCLUDED.embedding IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedded_at
@@ -2386,19 +2377,14 @@ export class PGLiteEngine implements BrainEngine {
* drift (NULL grandfathered never stale). Shared by countStaleChunks +
* sumStaleChunkChars so they can't drift.
*/
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): { where: string; params: unknown[] } {
// #1262: staleness targets the caller-resolved write column when set
// (identifier-validated + quoted); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string }): { where: string; params: unknown[] } {
const params: unknown[] = [];
const conds: string[] = [];
if (opts?.signature !== undefined) {
params.push(opts.signature);
conds.push(`(cc.${staleCol} IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
conds.push(`(cc.embedding IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
} else {
conds.push(`cc.${staleCol} IS NULL`);
conds.push(`cc.embedding IS NULL`);
}
conds.push(`NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')`);
if (opts?.sourceId !== undefined) {
@@ -2408,7 +2394,7 @@ export class PGLiteEngine implements BrainEngine {
return { where: conds.join(' AND '), params };
}
async countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// D7: source-scoped count for `gbrain embed --stale --source X`. Always
// JOIN pages so embed-skip + signature predicates apply. PGLite is
// PostgreSQL 17.5 in WASM and supports the full JSONB operator set.
@@ -2424,7 +2410,7 @@ export class PGLiteEngine implements BrainEngine {
return Number(count);
}
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// Sibling of countStaleChunks: same stale predicate, summing chunk_text
// length for the sync cost preview. ::bigint guards int4 overflow.
const { where, params } = this.buildStaleChunkWhere(opts);
@@ -2477,17 +2463,11 @@ export class PGLiteEngine implements BrainEngine {
sourceId?: string;
orderBy?: 'page_id' | 'updated_desc';
afterUpdatedAt?: string | null;
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]> {
const limit = opts?.batchSize ?? 2000;
const afterPid = opts?.afterPageId ?? 0;
const afterIdx = opts?.afterChunkIndex ?? -1;
const orderBy = opts?.orderBy ?? 'page_id';
// #1262: staleness follows the caller-resolved write column (validated +
// quoted identifier); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
// v0.41.18.0 (A13, codex #9): --priority recent path. See postgres-engine
// sibling for full rationale. Same composite cursor + ORDER BY.
@@ -2501,7 +2481,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
LIMIT $1`,
@@ -2512,7 +2492,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
p.updated_at < $1::timestamptz
@@ -2531,7 +2511,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
@@ -2543,7 +2523,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
@@ -2568,7 +2548,7 @@ export class PGLiteEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > ($1, $2)
ORDER BY cc.page_id, cc.chunk_index
@@ -2582,7 +2562,7 @@ export class PGLiteEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > ($2, $3)
+22 -43
View File
@@ -50,7 +50,6 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
SalienceOpts, SalienceResult, AnomaliesOpts, AnomalyResult,
@@ -2381,21 +2380,13 @@ export class PostgresEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void> {
return this.batchRetry(opts?.auditSite ?? 'upsertChunks', opts?.signal, () => this._upsertChunksOnce(slug, chunks, opts), chunks.length);
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn }): Promise<void> {
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
// #1262: caller-resolved write target for TEXT embeddings. Descriptor
// names are identifier-validated + quoted by buildVectorCastFragment;
// omitted => legacy `embedding vector`.
const targetFragment = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn)
: undefined;
const targetCol = targetFragment?.col ?? 'embedding';
const embeddingCast = targetFragment?.castSql.replace('$1::', '') ?? 'vector';
// Source-scope the page-id lookup. Without this filter, multi-source
// brains where the slug exists in 2+ sources return >1 row and the
@@ -2422,7 +2413,7 @@ export class PostgresEngine implements BrainEngine {
// scope metadata through upserts.
// v0.27.1 (Phase 8): added `modality` + `embedding_image` to the column
// list. Image chunks pass embedding=null + embedding_image=Float32Array.
const cols = `(page_id, chunk_index, chunk_text, chunk_source, ${targetCol}, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)`;
const cols = '(page_id, chunk_index, chunk_text, chunk_source, embedding, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)';
const rows: string[] = [];
const params: unknown[] = [];
let paramIdx = 1;
@@ -2439,7 +2430,7 @@ export class PostgresEngine implements BrainEngine {
: null;
const modality = chunk.modality ?? 'text';
const embeddingPh = embeddingStr ? `$${paramIdx++}::${embeddingCast}` : 'NULL';
const embeddingPh = embeddingStr ? `$${paramIdx++}::vector` : 'NULL';
const embeddedAtPh = embeddingStr ? 'now()' : 'NULL';
const embeddingImagePh = embeddingImageStr ? `$${paramIdx++}::vector` : 'NULL';
@@ -2487,19 +2478,19 @@ export class PostgresEngine implements BrainEngine {
ON CONFLICT (page_id, chunk_index) DO UPDATE SET
chunk_text = EXCLUDED.chunk_text,
chunk_source = EXCLUDED.chunk_source,
${targetCol} = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.${targetCol}
WHEN content_chunks.${targetCol} IS NULL THEN EXCLUDED.${targetCol}
embedding = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.embedding
WHEN content_chunks.embedding IS NULL THEN EXCLUDED.embedding
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.${targetCol}
ELSE content_chunks.${targetCol}
THEN EXCLUDED.embedding
ELSE content_chunks.embedding
END,
model = COALESCE(EXCLUDED.model, content_chunks.model),
token_count = EXCLUDED.token_count,
embedded_at = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.${targetCol} IS NULL THEN NULL
WHEN content_chunks.${targetCol} IS NULL AND EXCLUDED.${targetCol} IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.embedding IS NULL THEN NULL
WHEN content_chunks.embedding IS NULL AND EXCLUDED.embedding IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedded_at
@@ -2539,19 +2530,14 @@ export class PostgresEngine implements BrainEngine {
* embedding_signature drift (NULL grandfathered). Shared by
* countStaleChunks + sumStaleChunkChars (parity with the PGLite sibling).
*/
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): { where: string; params: unknown[] } {
// #1262: staleness targets the caller-resolved write column when set
// (identifier-validated + quoted); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string }): { where: string; params: unknown[] } {
const params: unknown[] = [];
const conds: string[] = [];
if (opts?.signature !== undefined) {
params.push(opts.signature);
conds.push(`(cc.${staleCol} IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
conds.push(`(cc.embedding IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
} else {
conds.push(`cc.${staleCol} IS NULL`);
conds.push(`cc.embedding IS NULL`);
}
conds.push(`NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')`);
if (opts?.sourceId !== undefined) {
@@ -2561,7 +2547,7 @@ export class PostgresEngine implements BrainEngine {
return { where: conds.join(' AND '), params };
}
async countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// Always JOIN pages so the embed_skip + signature predicates apply.
// D7: source_id scoping. v0.41.31: optional signature widens staleness
// to embedding_signature drift (NULL grandfathered).
@@ -2579,7 +2565,7 @@ export class PostgresEngine implements BrainEngine {
});
}
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// Sibling of countStaleChunks: same stale predicate, summing chunk_text
// length for the sync cost preview. ::bigint guards int4 overflow.
const { where, params } = this.buildStaleChunkWhere(opts);
@@ -2632,18 +2618,11 @@ export class PostgresEngine implements BrainEngine {
sourceId?: string;
orderBy?: 'page_id' | 'updated_desc';
afterUpdatedAt?: string | null;
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]> {
const limit = opts?.batchSize ?? 2000;
const afterPid = opts?.afterPageId ?? 0;
const afterIdx = opts?.afterChunkIndex ?? -1;
const orderBy = opts?.orderBy ?? 'page_id';
// #1262: staleness follows the caller-resolved write column (validated +
// quoted identifier); legacy `embedding` otherwise. Interpolated below as
// an unsafe FRAGMENT (identifiers can't be bound parameters).
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
// RLS scope binding (opt-in via GBRAIN_RLS_SCOPE_BINDING).
return await this.withScopedReadTransaction(undefined, opts?.sourceId, async (tx) => {
@@ -2660,7 +2639,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
LIMIT ${limit}
@@ -2670,7 +2649,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
p.updated_at < ${afterUpdated}::timestamptz
@@ -2688,7 +2667,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
@@ -2699,7 +2678,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
@@ -2719,7 +2698,7 @@ export class PostgresEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > (${afterPid}, ${afterIdx})
ORDER BY cc.page_id, cc.chunk_index
@@ -2732,7 +2711,7 @@ export class PostgresEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > (${afterPid}, ${afterIdx})
-74
View File
@@ -443,80 +443,6 @@ export function resolveEmbeddingColumn(
};
}
/**
* Resolves the WRITE-side embedding column for the currently configured
* embedding model (#1262). The read-side resolver above answers "which
* column does this query search?"; this one answers "which column should
* newly produced text embeddings land in?".
*
* Unlike read-side search, writes take no per-call column override. The
* import/embed boundary resolves once from merged config + gateway state
* and passes the descriptor into `engine.upsertChunks`; engines stay
* config-free (same contract as the read-side descriptor).
*
* Behavior:
* - no user-declared `embedding_columns` => undefined (legacy brain,
* writes keep targeting the default `embedding` column)
* - a user-declared entry whose `provider` matches the current
* embedding model => that entry's descriptor
* - no provider match => undefined (fall back to legacy `embedding`)
*
* Only USER-declared entries are consulted never the cfg-derived
* builtins. The `embedding_image` builtin's provider is the multimodal
* model; matching it here would misroute text embeddings into the image
* column. The no-match fallback is intentional: switching models before
* registering a matching column must not silently write vectors into an
* arbitrary column.
*/
export function resolveWriteColumn(cfg: GBrainConfig): ResolvedColumn | undefined {
const userColumns = cfg.embedding_columns;
if (
!userColumns ||
typeof userColumns !== 'object' ||
Array.isArray(userColumns) ||
Object.keys(userColumns).length === 0
) {
return undefined;
}
// Same model-resolution chain as the registry builtin: cfg > gateway > default.
let gwModel: string | undefined;
try {
const gw = require('../ai/gateway.ts') as typeof import('../ai/gateway.ts');
gwModel = gw.getEmbeddingModel();
} catch {
// Gateway unconfigured — fall through to the canonical default.
}
const currentModel = cfg.embedding_model ?? gwModel ?? DEFAULT_EMBEDDING_MODEL;
for (const [name, entry] of Object.entries(userColumns)) {
if (!entry) continue;
validateColumnKey(name);
validateColumnConfig(name, entry);
if (entry.provider !== currentModel) continue;
return {
name,
type: entry.type,
dimensions: entry.dimensions,
embeddingModel: entry.provider,
};
}
return undefined;
}
/**
* Engine-boundary convenience: merged config (file/env + DB plane)
* resolveWriteColumn. Dynamic import keeps config.ts out of this module's
* static graph (mirrors the gateway require above).
*/
export async function resolveWriteColumnForEngine(
engine: { getConfig(key: string): Promise<string | null | undefined> },
): Promise<ResolvedColumn | undefined> {
const { loadConfigWithEngine } = await import('../config.ts');
const cfg = await loadConfigWithEngine(engine);
return cfg ? resolveWriteColumn(cfg) : undefined;
}
/**
* True when the resolved column is the default `embedding` name.
* Name-based check; does not compare embedding space.
+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();
});
});
+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/);
});
-133
View File
@@ -241,136 +241,3 @@ describe('buildVectorCastFragment — engine SQL composer (D3)', () => {
expect(castSql).toBe('$1::halfvec(2560)');
});
});
describe('PGLite engine: upsertChunks write-side ResolvedColumn descriptor (#1262)', () => {
test('halfvec descriptor writes the text embedding to the alternate column, not legacy embedding', async () => {
await engine.putPage('docs/write-alt-pglite', {
type: 'concept',
title: 'Write alt column PGLite',
compiled_truth: 'PGLite write-side alternate embedding column test.',
});
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
await engine.upsertChunks('docs/write-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'PGLite write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{
has_default: boolean;
has_ze: boolean;
has_embedded_at: boolean;
}>(
`SELECT embedding IS NOT NULL AS has_default,
embedding_ze IS NOT NULL AS has_ze,
embedded_at IS NOT NULL AS has_embedded_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-pglite'`,
);
expect(rows.length).toBe(1);
expect(rows[0].has_default).toBe(false);
expect(rows[0].has_ze).toBe(true);
expect(rows[0].has_embedded_at).toBe(true);
});
test('text-unchanged re-upsert without a vector preserves the alternate-column embedding', async () => {
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
// Same chunk_text, no embedding: the ON CONFLICT CASE must keep the
// existing alternate-column vector (D24 semantics follow the column).
await engine.upsertChunks('docs/write-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'PGLite write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{ has_ze: boolean }>(
`SELECT embedding_ze IS NOT NULL AS has_ze
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-pglite'`,
);
expect(rows).toEqual([{ has_ze: true }]);
});
});
describe('PGLite: embed --stale converges on an alt-column brain (#1262)', () => {
test('boundary resolves the write column; stale scan does not re-select embedded rows', async () => {
const { runEmbedCore } = await import('../../src/commands/embed.ts');
const local = new PGLiteEngine();
const previousHome = process.env.GBRAIN_HOME;
process.env.GBRAIN_HOME = `/tmp/gbrain-write-col-stale-${Date.now()}`;
try {
await local.connect({});
await local.initSchema();
await (local as any).db.exec(
`ALTER TABLE content_chunks ADD COLUMN IF NOT EXISTS embedding_ze halfvec(2560)`,
);
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
await local.setConfig('embedding_columns', JSON.stringify({
embedding_ze: { provider: 'zeroentropyai:zembed-1', dimensions: 2560, type: 'halfvec' },
}));
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
await local.putPage('docs/stale-alt-pglite', {
type: 'concept',
title: 'Dynamic stale column',
compiled_truth: 'A chunk that is embedded only in the dynamic column.',
});
await local.upsertChunks('docs/stale-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'A chunk that is embedded only in the dynamic column.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
// Engine-level contrast: legacy predicate still sees the row as stale;
// the alt-column predicate does not.
expect(await local.countStaleChunks()).toBe(1);
expect(await local.countStaleChunks({ embeddingColumn: descriptor })).toBe(0);
// sumStaleChunkChars feeds the sync cost gate — same predicate contract.
expect(await local.sumStaleChunkChars()).toBeGreaterThan(0);
expect(await local.sumStaleChunkChars({ embeddingColumn: descriptor })).toBe(0);
expect(await local.listStaleChunks({ embeddingColumn: descriptor, batchSize: 100 })).toHaveLength(0);
expect(await local.listStaleChunks({ batchSize: 100 })).toHaveLength(1);
// Boundary-level: `embed --stale --dry-run` resolves the write column
// from merged config + gateway and reports NOTHING to embed. Without
// the fix this reports 1 (perpetual re-embed loop).
const result = await runEmbedCore(local, { stale: true, dryRun: true });
expect(result.would_embed).toBe(0);
} finally {
await local.disconnect();
if (previousHome === undefined) delete process.env.GBRAIN_HOME;
else process.env.GBRAIN_HOME = previousHome;
resetGateway();
}
});
});
@@ -224,54 +224,4 @@ if (!dbUrl) {
await engine.executeRaw(`UPDATE content_chunks SET embedding_voyage = '${v}'::vector WHERE id = ${dogId}`);
});
});
describe('Postgres: upsertChunks write-side ResolvedColumn descriptor (#1262)', () => {
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
test('halfvec descriptor writes the text embedding to the alternate column, not legacy embedding', async () => {
await engine.putPage('docs/write-alt-postgres', {
type: 'concept',
title: 'Write alt column Postgres',
compiled_truth: 'Postgres write-side alternate embedding column test.',
});
await engine.upsertChunks('docs/write-alt-postgres', [
{
chunk_index: 0,
chunk_text: 'Postgres write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{
has_default: boolean;
has_ze: boolean;
}>(
`SELECT embedding IS NOT NULL AS has_default,
embedding_ze IS NOT NULL AS has_ze
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-postgres'`,
);
expect(rows.length).toBe(1);
expect(rows[0].has_default).toBe(false);
expect(rows[0].has_ze).toBe(true);
}, 30_000);
test('stale scan follows the write-side column (count + list parity with the write target)', async () => {
// Legacy predicate: cat/dog/write-alt rows all have embedding NULL.
expect(await engine.countStaleChunks()).toBeGreaterThan(0);
// Alt-column predicate: every chunk has embedding_ze populated.
expect(await engine.countStaleChunks({ embeddingColumn: descriptor })).toBe(0);
expect(await engine.listStaleChunks({ embeddingColumn: descriptor, batchSize: 100 })).toHaveLength(0);
expect((await engine.listStaleChunks({ batchSize: 100 })).length).toBeGreaterThan(0);
// updated_desc arm uses the same predicate.
expect(await engine.listStaleChunks({ embeddingColumn: descriptor, orderBy: 'updated_desc', batchSize: 100 })).toHaveLength(0);
}, 30_000);
});
}
+1 -110
View File
@@ -13,10 +13,9 @@
* throw on unknown string.
*/
import { describe, test, expect, afterAll, afterEach } from 'bun:test';
import { describe, test, expect } from 'bun:test';
import {
resolveEmbeddingColumn,
resolveWriteColumn,
getEmbeddingColumnRegistry,
buildVectorCastFragment,
quoteIdentifier,
@@ -35,28 +34,6 @@ import {
} from '../../src/core/search/embedding-column.ts';
import type { GBrainConfig } from '../../src/core/config.ts';
import type { ResolvedColumn } from '../../src/core/types.ts';
import { configureGateway, resetGateway } from '../../src/core/ai/gateway.ts';
/**
* Teardown: reset AND re-apply the legacy preload config
* (test/helpers/legacy-embedding-preload.ts). A bare resetGateway() would
* leave the slot empty for the NEXT file's beforeAll (the preload's
* per-test beforeEach only fires before tests, not before beforeAll), which
* would make sibling PGLite fixtures initSchema at the 1280 default instead
* of the legacy 1536 their seed vectors assume.
*/
function restorePreloadGateway() {
resetGateway();
configureGateway({
embedding_model: 'openai:text-embedding-3-large',
embedding_dimensions: 1536,
env: { ...process.env },
});
}
afterAll(() => {
restorePreloadGateway();
});
function cfg(overrides: Partial<GBrainConfig> = {}): GBrainConfig {
return { engine: 'pglite', ...overrides };
@@ -545,89 +522,3 @@ describe('codex /ship #4 — isCacheSafe (embedding-space-based skip)', () => {
expect(isCacheSafe(r, cfg())).toBe(true);
});
});
describe('resolveWriteColumn — write-side boundary resolution (#1262)', () => {
afterEach(() => {
restorePreloadGateway();
});
test('no registry / empty registry returns undefined (legacy single-column brain)', () => {
expect(resolveWriteColumn(cfg())).toBeUndefined();
expect(resolveWriteColumn(cfg({ embedding_columns: {} }))).toBeUndefined();
});
test('provider match via cfg.embedding_model returns the descriptor', () => {
const r = resolveWriteColumn(cfg({
embedding_model: 'voyage:voyage-3-large',
embedding_dimensions: 1024,
embedding_columns: {
embedding_voyage: { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
},
}));
expect(r).toEqual({
name: 'embedding_voyage',
type: 'vector',
dimensions: 1024,
embeddingModel: 'voyage:voyage-3-large',
});
});
test('provider match via gateway state (cfg.embedding_model unset) returns descriptor', () => {
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_ze: { provider: 'zeroentropyai:zembed-1', dimensions: 2560, type: 'halfvec' },
},
}));
expect(r).toEqual({
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
});
});
test('no provider match returns undefined instead of guessing a column', () => {
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_voyage: { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
},
}));
expect(r).toBeUndefined();
});
test('only USER-declared columns are consulted — multimodal builtin never captures text writes', () => {
// Current model equals the embedding_image BUILTIN's provider; a registry
// walk that consulted builtins would misroute text writes into the image
// column. resolveWriteColumn must return undefined here.
configureGateway({
embedding_model: 'voyage:voyage-multimodal-3',
embedding_dimensions: 1024,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_other: { provider: 'openai:text-embedding-3-large', dimensions: 1536, type: 'vector' },
},
}));
expect(r).toBeUndefined();
});
test('malformed registry entry throws loud (same validation as the read side)', () => {
expect(() => resolveWriteColumn(cfg({
embedding_model: 'voyage:voyage-3-large',
embedding_columns: {
'bad"col': { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
} as never,
}))).toThrow(EmbeddingColumnConfigError);
});
});