mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-15 17:32:37 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b0185306e1 |
Vendored
+56
File diff suppressed because one or more lines are too long
Vendored
-57
File diff suppressed because one or more lines are too long
Vendored
+1
-1
@@ -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-CzLjRij_.js"></script>
|
||||
<script type="module" crossorigin src="/admin/assets/index-CoGEje3-.js"></script>
|
||||
<link rel="stylesheet" crossorigin href="/admin/assets/index-GxkWX7v3.css">
|
||||
</head>
|
||||
<body>
|
||||
|
||||
@@ -39,7 +39,6 @@ export const api = {
|
||||
stats: () => apiFetch('/admin/api/stats'),
|
||||
health: () => apiFetch('/admin/api/health-indicators'),
|
||||
agents: () => apiFetch('/admin/api/agents'),
|
||||
sources: () => apiFetch('/admin/api/sources'),
|
||||
requests: (page = 1, qs = '') => apiFetch(`/admin/api/requests?page=${page}${qs}`),
|
||||
apiKeys: () => apiFetch('/admin/api/api-keys'),
|
||||
createApiKey: (name: string) => apiFetch('/admin/api/api-keys', { method: 'POST', body: JSON.stringify({ name }) }),
|
||||
|
||||
@@ -23,16 +23,9 @@ interface Agent {
|
||||
total_requests: number;
|
||||
requests_today: number;
|
||||
token_ttl: number | null;
|
||||
source_id: string | null;
|
||||
federated_read: string[] | null;
|
||||
status: 'active' | 'revoked';
|
||||
}
|
||||
|
||||
interface SourceOption {
|
||||
id: string;
|
||||
name: string | null;
|
||||
}
|
||||
|
||||
interface ApiKey {
|
||||
id: string;
|
||||
name: string;
|
||||
@@ -95,7 +88,6 @@ export function AgentsPage() {
|
||||
<th>Name</th>
|
||||
<th>Type</th>
|
||||
<th>Scopes</th>
|
||||
<th>Source</th>
|
||||
<th>Status</th>
|
||||
<th>Requests</th>
|
||||
<th>Last Used</th>
|
||||
@@ -116,9 +108,6 @@ export function AgentsPage() {
|
||||
<span key={s} className={`badge badge-${s}`} style={{ marginRight: 4 }}>{s}</span>
|
||||
))}
|
||||
</td>
|
||||
<td className="mono" style={{ fontSize: 12, color: 'var(--text-secondary)' }}>
|
||||
{a.auth_type === 'oauth' ? (a.source_id || 'default') : 'all'}
|
||||
</td>
|
||||
<td>
|
||||
<span className={`badge ${a.status === 'active' ? 'badge-success' : 'badge-danger'}`}>{a.status}</span>
|
||||
</td>
|
||||
@@ -268,29 +257,9 @@ function RegisterModal({ onClose, onRegistered }: {
|
||||
Object.fromEntries(ALLOWED_SCOPES_LIST.map(s => [s, s === 'read'])) as Record<Scope, boolean>,
|
||||
);
|
||||
const [ttl, setTtl] = useState('86400'); // 24h default
|
||||
const [redirectUris, setRedirectUris] = useState('');
|
||||
const [sources, setSources] = useState<SourceOption[]>([{ id: 'default', name: 'default' }]);
|
||||
const [sourceId, setSourceId] = useState('default');
|
||||
const [readSources, setReadSources] = useState<Record<string, boolean>>({ default: true });
|
||||
const [loading, setLoading] = useState(false);
|
||||
const [error, setError] = useState('');
|
||||
|
||||
useEffect(() => {
|
||||
api.sources()
|
||||
.then((rows: SourceOption[]) => {
|
||||
const nextSources = rows.length > 0 ? rows : [{ id: 'default', name: 'default' }];
|
||||
const nextSourceId = nextSources.some(s => s.id === sourceId) ? sourceId : nextSources[0].id;
|
||||
setSources(nextSources);
|
||||
setSourceId(nextSourceId);
|
||||
setReadSources(prev => {
|
||||
const next = Object.fromEntries(nextSources.map(s => [s.id, Boolean(prev[s.id])])) as Record<string, boolean>;
|
||||
next[nextSourceId] = true;
|
||||
return next;
|
||||
});
|
||||
})
|
||||
.catch(() => {});
|
||||
}, []);
|
||||
|
||||
const ttlOptions = [
|
||||
{ label: '1 hour', value: '3600' },
|
||||
{ label: '24 hours', value: '86400' },
|
||||
@@ -308,24 +277,11 @@ function RegisterModal({ onClose, onRegistered }: {
|
||||
try {
|
||||
// Use the CLI registration endpoint (POST to admin API)
|
||||
const selectedScopes = Object.entries(scopes).filter(([, v]) => v).map(([k]) => k).join(' ');
|
||||
const federatedRead = sources.map(s => s.id).filter(id => readSources[id]);
|
||||
if (federatedRead.length === 0) { setError('Select at least one read source'); setLoading(false); return; }
|
||||
// #1036: redirect URIs, one per line. Mirrors the CLI convention
|
||||
// (src/commands/auth.ts): presence of redirect URIs implies
|
||||
// authorization_code + refresh_token grants unless the user had none.
|
||||
const uris = redirectUris.split('\n').map(u => u.trim()).filter(Boolean);
|
||||
const res = await fetch('/admin/api/register-client', {
|
||||
method: 'POST',
|
||||
credentials: 'same-origin',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
name: name.trim(),
|
||||
scopes: selectedScopes,
|
||||
tokenTtl: ttl === '0' ? 315360000 : Number(ttl),
|
||||
source_id: sourceId,
|
||||
federated_read: federatedRead,
|
||||
...(uris.length > 0 ? { redirectUris: uris, grantTypes: ['authorization_code', 'refresh_token'] } : {}),
|
||||
}),
|
||||
body: JSON.stringify({ name: name.trim(), scopes: selectedScopes, tokenTtl: ttl === '0' ? 315360000 : Number(ttl) }),
|
||||
});
|
||||
if (!res.ok) throw new Error('Registration failed');
|
||||
const data = await res.json();
|
||||
@@ -356,42 +312,6 @@ function RegisterModal({ onClose, onRegistered }: {
|
||||
))}
|
||||
</div>
|
||||
</div>
|
||||
<div style={{ marginBottom: 16 }}>
|
||||
<label>Write Source</label>
|
||||
<select value={sourceId} onChange={e => {
|
||||
const id = e.target.value;
|
||||
setSourceId(id);
|
||||
setReadSources(p => ({ ...p, [id]: true }));
|
||||
}}
|
||||
style={{ width: '100%', background: 'var(--bg-secondary)', color: 'var(--text-primary)', border: '1px solid var(--border)', borderRadius: 6, padding: '6px 10px', fontSize: 14 }}>
|
||||
{sources.map(s => <option key={s.id} value={s.id}>{s.name && s.name !== s.id ? `${s.id} (${s.name})` : s.id}</option>)}
|
||||
</select>
|
||||
</div>
|
||||
<div style={{ marginBottom: 16 }}>
|
||||
<label>Read Sources</label>
|
||||
<div className="checkbox-group">
|
||||
{sources.map(s => (
|
||||
<label key={s.id} className="checkbox-label">
|
||||
<input
|
||||
type="checkbox"
|
||||
checked={Boolean(readSources[s.id])}
|
||||
onChange={e => setReadSources(p => ({ ...p, [s.id]: e.target.checked }))}
|
||||
/>
|
||||
{s.name && s.name !== s.id ? `${s.id} (${s.name})` : s.id}
|
||||
</label>
|
||||
))}
|
||||
</div>
|
||||
</div>
|
||||
<div style={{ marginBottom: 16 }}>
|
||||
<label>Redirect URIs <span style={{ color: 'var(--text-secondary)', fontWeight: 400 }}>(optional, one per line — enables authorization_code flow)</span></label>
|
||||
<textarea
|
||||
placeholder={'https://example.com/callback'}
|
||||
value={redirectUris}
|
||||
onChange={e => setRedirectUris(e.target.value)}
|
||||
rows={2}
|
||||
style={{ width: '100%', background: 'var(--bg-secondary)', color: 'var(--text-primary)', border: '1px solid var(--border)', borderRadius: 6, padding: '6px 10px', fontSize: 13, fontFamily: 'inherit', resize: 'vertical' }}
|
||||
/>
|
||||
</div>
|
||||
<div style={{ marginBottom: 20 }}>
|
||||
<label>Token Lifetime</label>
|
||||
<select value={ttl} onChange={e => setTtl(e.target.value)}
|
||||
@@ -608,8 +528,6 @@ function AgentDrawer({ agent, onClose, onRevoked }: { agent: Agent; onClose: ()
|
||||
client_name: agentName,
|
||||
auth_type: agent.auth_type,
|
||||
scope: agent.scope,
|
||||
source_id: agent.source_id,
|
||||
federated_read: agent.federated_read,
|
||||
}, null, 2),
|
||||
};
|
||||
|
||||
@@ -629,10 +547,6 @@ function AgentDrawer({ agent, onClose, onRevoked }: { agent: Agent; onClose: ()
|
||||
<span>{(agent.scope || '').split(' ').filter(Boolean).map(s => (
|
||||
<span key={s} className={`badge badge-${s}`} style={{ marginRight: 4 }}>{s}</span>
|
||||
))}</span>
|
||||
<span style={{ color: 'var(--text-secondary)' }}>Write Source</span>
|
||||
<span className="mono">{agent.source_id || (isOAuth ? 'default' : 'all')}</span>
|
||||
<span style={{ color: 'var(--text-secondary)' }}>Read Sources</span>
|
||||
<span className="mono">{isOAuth ? (agent.federated_read && agent.federated_read.length > 0 ? agent.federated_read.join(', ') : (agent.source_id || 'default')) : 'all'}</span>
|
||||
<span style={{ color: 'var(--text-secondary)' }}>Registered</span>
|
||||
<span>{new Date(agent.created_at).toLocaleDateString()}</span>
|
||||
<span style={{ color: 'var(--text-secondary)' }}>Token TTL</span>
|
||||
|
||||
@@ -21,14 +21,17 @@ GBrain is tuned for the Supabase **Transaction pooler** (port 6543): it
|
||||
auto-disables prepared statements there and routes `engine.transaction()`
|
||||
(migrations, DDL, sync imports) to a derived **direct** connection
|
||||
(`db.<ref>.supabase.co:5432`). That direct host is IPv6-only, so on an
|
||||
IPv4-only host, reads work but sync **silently skips most pages**. This is the
|
||||
number one cause of "sync ran but nothing happened."
|
||||
IPv4-only host it is unreachable. When that happens gbrain now falls back to
|
||||
the pooler automatically (one stderr warning, then single-pool mode for the
|
||||
rest of the process) — but the pooler's ~2-min statement timeout can truncate
|
||||
very long migrations or bulk imports.
|
||||
|
||||
Fix: make the direct connection reachable over IPv4. Either set
|
||||
`GBRAIN_DIRECT_DATABASE_URL` to the **Session pooler** string (port 5432 on the
|
||||
`pooler.supabase.com` host, IPv4), or enable Supabase's IPv4 add-on. Verify by
|
||||
running `gbrain sync` and checking that the page count in `gbrain stats` matches
|
||||
the syncable file count in the repo.
|
||||
`pooler.supabase.com` host, IPv4), or enable Supabase's IPv4 add-on.
|
||||
`GBRAIN_DISABLE_DIRECT_POOL=1` skips the direct pool (and the fallback warning)
|
||||
entirely. Verify by running `gbrain sync` and checking that the page count in
|
||||
`gbrain stats` matches the syncable file count in the repo.
|
||||
|
||||
### The Primitives
|
||||
|
||||
|
||||
+8
-5
@@ -2720,14 +2720,17 @@ GBrain is tuned for the Supabase **Transaction pooler** (port 6543): it
|
||||
auto-disables prepared statements there and routes `engine.transaction()`
|
||||
(migrations, DDL, sync imports) to a derived **direct** connection
|
||||
(`db.<ref>.supabase.co:5432`). That direct host is IPv6-only, so on an
|
||||
IPv4-only host, reads work but sync **silently skips most pages**. This is the
|
||||
number one cause of "sync ran but nothing happened."
|
||||
IPv4-only host it is unreachable. When that happens gbrain now falls back to
|
||||
the pooler automatically (one stderr warning, then single-pool mode for the
|
||||
rest of the process) — but the pooler's ~2-min statement timeout can truncate
|
||||
very long migrations or bulk imports.
|
||||
|
||||
Fix: make the direct connection reachable over IPv4. Either set
|
||||
`GBRAIN_DIRECT_DATABASE_URL` to the **Session pooler** string (port 5432 on the
|
||||
`pooler.supabase.com` host, IPv4), or enable Supabase's IPv4 add-on. Verify by
|
||||
running `gbrain sync` and checking that the page count in `gbrain stats` matches
|
||||
the syncable file count in the repo.
|
||||
`pooler.supabase.com` host, IPv4), or enable Supabase's IPv4 add-on.
|
||||
`GBRAIN_DISABLE_DIRECT_POOL=1` skips the direct pool (and the fallback warning)
|
||||
entirely. Verify by running `gbrain sync` and checking that the page count in
|
||||
`gbrain stats` matches the syncable file count in the repo.
|
||||
|
||||
### The Primitives
|
||||
|
||||
|
||||
@@ -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-07-21.
|
||||
// Source: admin/dist/ at 2026-05-27.
|
||||
//
|
||||
// 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_CzLjRij__js from '../admin/dist/assets/index-CzLjRij_.js' with { type: 'file' };
|
||||
import A_0_assets_index_CoGEje3__js from '../admin/dist/assets/index-CoGEje3-.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-CzLjRij_.js": { path: A_0_assets_index_CzLjRij__js as unknown as string, mime: "application/javascript; charset=utf-8" },
|
||||
"/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-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" },
|
||||
};
|
||||
|
||||
@@ -1078,6 +1078,9 @@ async function initPostgres(opts: {
|
||||
console.warn(' Direct connections are IPv6 only and fail in many environments.');
|
||||
console.warn(' Use the Transaction pooler connection string instead (port 6543):');
|
||||
console.warn(' Supabase Dashboard > Connect (top bar) > Connection String > Transaction pooler');
|
||||
console.warn(' (With a pooler URL, gbrain derives a direct connection for DDL and falls back');
|
||||
console.warn(' to the pooler automatically if that host is unreachable. Power users:');
|
||||
console.warn(' GBRAIN_DIRECT_DATABASE_URL overrides the derived URL; GBRAIN_DISABLE_DIRECT_POOL=1 disables it.)');
|
||||
console.warn('');
|
||||
}
|
||||
|
||||
@@ -1091,6 +1094,9 @@ async function initPostgres(opts: {
|
||||
if (databaseUrl.includes('supabase.co') && (msg.includes('ECONNREFUSED') || msg.includes('ETIMEDOUT'))) {
|
||||
console.error('Connection failed. Supabase direct connections (db.*.supabase.co:5432) are IPv6 only.');
|
||||
console.error('Use the Transaction pooler connection string instead (port 6543).');
|
||||
console.error('(gbrain derives its own direct connection from pooler URLs for DDL; if that host is');
|
||||
console.error('unreachable it falls back to the pooler. GBRAIN_DIRECT_DATABASE_URL overrides the');
|
||||
console.error('derived URL; GBRAIN_DISABLE_DIRECT_POOL=1 disables the direct pool entirely.)');
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
|
||||
@@ -39,7 +39,6 @@ import * as db from '../core/db.ts';
|
||||
import { sqlQueryForEngine, executeRawJsonb } from '../core/sql-query.ts';
|
||||
import { MinionQueue } from '../core/minions/queue.ts';
|
||||
import { isRetryableError } from '../core/retry-matcher.ts';
|
||||
import { normalizeAdminClientSourceScope } from '../core/admin-source-scope.ts';
|
||||
import {
|
||||
computeContentHash,
|
||||
validateIngestionEvent,
|
||||
@@ -1087,7 +1086,6 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
const oauthClients = await sql`
|
||||
SELECT c.client_id as id, c.client_name as name, 'oauth' as auth_type,
|
||||
c.grant_types, c.scope, c.created_at, c.token_ttl,
|
||||
c.source_id, c.federated_read,
|
||||
CASE WHEN c.deleted_at IS NOT NULL THEN 'revoked' ELSE 'active' END as status,
|
||||
(SELECT max(created_at) FROM mcp_request_log WHERE token_name = c.client_id) as last_used_at,
|
||||
(SELECT count(*)::int FROM mcp_request_log WHERE token_name = c.client_id) as total_requests,
|
||||
@@ -1097,7 +1095,6 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
const legacyKeys = await sql`
|
||||
SELECT a.id, a.name, 'api_key' as auth_type,
|
||||
'{"bearer"}' as grant_types, 'read write admin' as scope, a.created_at, null as token_ttl,
|
||||
null as source_id, null as federated_read,
|
||||
CASE WHEN a.revoked_at IS NOT NULL THEN 'revoked' ELSE 'active' END as status,
|
||||
a.last_used_at,
|
||||
(SELECT count(*)::int FROM mcp_request_log WHERE token_name = a.name) as total_requests,
|
||||
@@ -1110,17 +1107,6 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
}
|
||||
});
|
||||
|
||||
app.get('/admin/api/sources', requireAdmin, async (_req: Request, res: Response) => {
|
||||
try {
|
||||
const sources = await engine.listAllSources({ includeArchived: false });
|
||||
const options = sources.map(s => ({ id: s.id, name: s.name }));
|
||||
if (!options.some(s => s.id === 'default')) options.unshift({ id: 'default', name: 'default' });
|
||||
res.json(options);
|
||||
} catch (e) {
|
||||
res.status(503).json({ error: 'service_unavailable' });
|
||||
}
|
||||
});
|
||||
|
||||
// v0.38 Slice 4 — per-OAuth-client agent spend viewer. Pre-computes today's
|
||||
// spend (committed + pending reservations) per client so the Agents tab
|
||||
// can render a "$X / $Y today" cell. Read-side endpoint only — no mutation.
|
||||
@@ -1454,8 +1440,6 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
// missing, empty) and rejects the rest with a structured 400.
|
||||
const { name, tokenTtl, grantTypes, redirectUris, tokenEndpointAuthMethod } = req.body;
|
||||
const rawScopes = (req.body as Record<string, unknown>).scopes ?? (req.body as Record<string, unknown>).scope;
|
||||
const rawSourceId = (req.body as Record<string, unknown>).source_id ?? (req.body as Record<string, unknown>).sourceId;
|
||||
const rawFederatedRead = (req.body as Record<string, unknown>).federated_read ?? (req.body as Record<string, unknown>).federatedRead;
|
||||
if (!name) { res.status(400).json({ error: 'Name required' }); return; }
|
||||
let scopeString: string;
|
||||
try {
|
||||
@@ -1486,23 +1470,8 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
});
|
||||
return;
|
||||
}
|
||||
let sourceScope: ReturnType<typeof normalizeAdminClientSourceScope>;
|
||||
try {
|
||||
const sources = await engine.listAllSources({ includeArchived: false });
|
||||
sourceScope = normalizeAdminClientSourceScope(
|
||||
rawSourceId,
|
||||
rawFederatedRead,
|
||||
sources.map(s => s.id),
|
||||
);
|
||||
} catch (e) {
|
||||
res.status(400).json({
|
||||
error: 'invalid_source_scope',
|
||||
message: e instanceof Error ? e.message : String(e),
|
||||
});
|
||||
return;
|
||||
}
|
||||
const result = await oauthProvider.registerClientManual(
|
||||
name, grants, scopeString, uris, sourceScope.sourceId, sourceScope.federatedRead, validatedAuthMethod,
|
||||
name, grants, scopeString, uris, 'default', undefined, validatedAuthMethod,
|
||||
);
|
||||
// Set per-client TTL if specified
|
||||
if (tokenTtl && Number(tokenTtl) > 0) {
|
||||
|
||||
@@ -1,51 +0,0 @@
|
||||
import { isValidSourceId } from './source-id.ts';
|
||||
|
||||
export interface AdminClientSourceScope {
|
||||
sourceId: string;
|
||||
federatedRead: string[];
|
||||
}
|
||||
|
||||
function parseFederatedReadInput(raw: unknown, sourceId: string): string[] {
|
||||
if (raw === undefined || raw === null) return [sourceId];
|
||||
if (typeof raw === 'string') {
|
||||
const values = raw.split(',').map(s => s.trim()).filter(Boolean);
|
||||
return values.length > 0 ? values : [sourceId];
|
||||
}
|
||||
if (Array.isArray(raw)) {
|
||||
if (!raw.every(v => typeof v === 'string')) {
|
||||
throw new Error('federated_read must be an array of source ids');
|
||||
}
|
||||
const values = raw.map(v => v.trim()).filter(Boolean);
|
||||
return values.length > 0 ? values : [sourceId];
|
||||
}
|
||||
throw new Error('federated_read must be a string or array');
|
||||
}
|
||||
|
||||
export function normalizeAdminClientSourceScope(
|
||||
rawSourceId: unknown,
|
||||
rawFederatedRead: unknown,
|
||||
availableSourceIds: string[],
|
||||
): AdminClientSourceScope {
|
||||
const knownSources = new Set(availableSourceIds);
|
||||
const fallbackSourceId = knownSources.has('default') ? 'default' : (availableSourceIds[0] ?? 'default');
|
||||
const sourceId = rawSourceId === undefined || rawSourceId === null || rawSourceId === ''
|
||||
? fallbackSourceId
|
||||
: rawSourceId;
|
||||
if (!isValidSourceId(sourceId)) {
|
||||
throw new Error(`Invalid source_id: ${JSON.stringify(sourceId)}`);
|
||||
}
|
||||
if (knownSources.size > 0 && !knownSources.has(sourceId)) {
|
||||
throw new Error(`Unknown source_id: ${sourceId}`);
|
||||
}
|
||||
|
||||
const federatedRead = Array.from(new Set(parseFederatedReadInput(rawFederatedRead, sourceId)));
|
||||
for (const id of federatedRead) {
|
||||
if (!isValidSourceId(id)) {
|
||||
throw new Error(`Invalid federated_read source_id: ${JSON.stringify(id)}`);
|
||||
}
|
||||
if (knownSources.size > 0 && !knownSources.has(id)) {
|
||||
throw new Error(`Unknown federated_read source_id: ${id}`);
|
||||
}
|
||||
}
|
||||
return { sourceId, federatedRead };
|
||||
}
|
||||
@@ -167,6 +167,25 @@ export function deriveDirectUrl(url: string): string | null {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Error codes that mean "the direct host is unreachable from this network"
|
||||
* (#1641). The auto-derived db.<ref>.supabase.co host is IPv6-only without
|
||||
* the paid IPv4 add-on, so ENOTFOUND/ECONNREFUSED here is expected on
|
||||
* IPv4-only networks — we fall back to the pooler instead of failing init.
|
||||
*/
|
||||
const NETWORK_UNREACHABLE_CODES = [
|
||||
'ENOTFOUND', 'ECONNREFUSED', 'ENETUNREACH', 'EHOSTUNREACH',
|
||||
'ETIMEDOUT', 'CONNECT_TIMEOUT',
|
||||
];
|
||||
|
||||
/** True when err looks like a network-unreachable failure (not auth/SQL). */
|
||||
export function isNetworkUnreachableError(err: unknown): boolean {
|
||||
const code = (err as { code?: unknown } | null)?.code;
|
||||
if (typeof code === 'string' && NETWORK_UNREACHABLE_CODES.includes(code)) return true;
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
return NETWORK_UNREACHABLE_CODES.some(c => msg.includes(c));
|
||||
}
|
||||
|
||||
/**
|
||||
* Read kill-switch state from env. Subordinate to parent manager's state
|
||||
* when present (A2 inheritance).
|
||||
@@ -319,7 +338,30 @@ export class ConnectionManager {
|
||||
throw err;
|
||||
});
|
||||
}
|
||||
const pool = await this._directInit;
|
||||
let pool: Sql | null;
|
||||
try {
|
||||
pool = await this._directInit;
|
||||
} catch (err) {
|
||||
// #1641: the derived direct host (db.<ref>.supabase.co) is IPv6-only
|
||||
// without Supabase's IPv4 add-on. On IPv4-only networks the direct
|
||||
// pool can never connect — permanently fall back to the read pool
|
||||
// (self-activating kill-switch) instead of failing init/migrations.
|
||||
// Non-network errors (auth, SQL) still throw: they mean misconfig,
|
||||
// not unreachability.
|
||||
if (isNetworkUnreachableError(err)) {
|
||||
const alreadyWarned = this._killSwitch;
|
||||
this._killSwitch = true;
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
if (!alreadyWarned) console.error(
|
||||
`gbrain: direct connection to ${this._directUrl ? this.hostOnly(this._directUrl) : 'unknown host'} unreachable (${msg}); ` +
|
||||
'falling back to the pooler for DDL/bulk (long migrations may hit the pooler statement timeout). ' +
|
||||
'Set GBRAIN_DIRECT_DATABASE_URL to a reachable direct URL (e.g. the Session pooler, port 5432) or enable the Supabase IPv4 add-on; ' +
|
||||
'GBRAIN_DISABLE_DIRECT_POOL=1 silences this.',
|
||||
);
|
||||
return this.getReadPool();
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
if (!pool) {
|
||||
// Defensive — initDirectPool should have thrown.
|
||||
throw new Error('connection-manager: direct pool init returned null');
|
||||
@@ -350,8 +392,9 @@ export class ConnectionManager {
|
||||
},
|
||||
};
|
||||
const t0 = Date.now();
|
||||
let pool: Sql | null = null;
|
||||
try {
|
||||
const pool = postgres(this._directUrl, opts);
|
||||
pool = postgres(this._directUrl, opts);
|
||||
// Probe to validate connectivity early.
|
||||
await pool`SELECT 1`;
|
||||
logConnectionEvent({
|
||||
@@ -362,6 +405,9 @@ export class ConnectionManager {
|
||||
});
|
||||
return pool;
|
||||
} catch (err) {
|
||||
// Don't leak the failed pool's sockets/timers (#1641 fallback keeps
|
||||
// the process running afterward).
|
||||
if (pool) await endPoolBounded(pool);
|
||||
logConnectionEvent({
|
||||
pool: 'ddl',
|
||||
op: 'error',
|
||||
|
||||
@@ -1,44 +0,0 @@
|
||||
import { describe, expect, test } from 'bun:test';
|
||||
import { normalizeAdminClientSourceScope } from '../src/core/admin-source-scope.ts';
|
||||
|
||||
describe('admin register-client source scope normalization', () => {
|
||||
const sources = ['default', 'team-a', 'team-b'];
|
||||
|
||||
test('defaults write and read scope to default source', () => {
|
||||
expect(normalizeAdminClientSourceScope(undefined, undefined, sources)).toEqual({
|
||||
sourceId: 'default',
|
||||
federatedRead: ['default'],
|
||||
});
|
||||
});
|
||||
|
||||
test('accepts snake-case body fields with multiple read sources', () => {
|
||||
expect(normalizeAdminClientSourceScope('team-a', ['team-a', 'team-b'], sources)).toEqual({
|
||||
sourceId: 'team-a',
|
||||
federatedRead: ['team-a', 'team-b'],
|
||||
});
|
||||
});
|
||||
|
||||
test('accepts comma-separated federated_read for API callers', () => {
|
||||
expect(normalizeAdminClientSourceScope('team-a', 'team-b, default,team-b', sources)).toEqual({
|
||||
sourceId: 'team-a',
|
||||
federatedRead: ['team-b', 'default'],
|
||||
});
|
||||
});
|
||||
|
||||
test('falls back to selected write source when federated_read is empty', () => {
|
||||
expect(normalizeAdminClientSourceScope('team-b', [], sources)).toEqual({
|
||||
sourceId: 'team-b',
|
||||
federatedRead: ['team-b'],
|
||||
});
|
||||
});
|
||||
|
||||
test('rejects malformed source ids before registration', () => {
|
||||
expect(() => normalizeAdminClientSourceScope('../secret', undefined, sources)).toThrow(/Invalid source_id/);
|
||||
expect(() => normalizeAdminClientSourceScope('team-a', ['snake_case'], sources)).toThrow(/Invalid federated_read/);
|
||||
});
|
||||
|
||||
test('rejects valid-looking but unregistered source ids', () => {
|
||||
expect(() => normalizeAdminClientSourceScope('ghost', undefined, sources)).toThrow(/Unknown source_id/);
|
||||
expect(() => normalizeAdminClientSourceScope('team-a', ['ghost'], sources)).toThrow(/Unknown federated_read/);
|
||||
});
|
||||
});
|
||||
@@ -3,6 +3,7 @@ import {
|
||||
isSupabasePoolerUrl,
|
||||
deriveDirectUrl,
|
||||
readKillSwitchEnv,
|
||||
isNetworkUnreachableError,
|
||||
resolveDirectPoolSize,
|
||||
ConnectionManager,
|
||||
DEFAULT_DIRECT_POOL_SIZE,
|
||||
@@ -238,3 +239,65 @@ describe('ConnectionManager — parent inheritance (A2)', () => {
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe('isNetworkUnreachableError (#1641)', () => {
|
||||
test('classifies network codes as unreachable', () => {
|
||||
for (const code of ['ENOTFOUND', 'ECONNREFUSED', 'ENETUNREACH', 'EHOSTUNREACH', 'ETIMEDOUT', 'CONNECT_TIMEOUT']) {
|
||||
const err = Object.assign(new Error('connect failed'), { code });
|
||||
expect(isNetworkUnreachableError(err)).toBe(true);
|
||||
}
|
||||
});
|
||||
|
||||
test('classifies by message when code absent', () => {
|
||||
expect(isNetworkUnreachableError(new Error('getaddrinfo ENOTFOUND db.abc.supabase.co'))).toBe(true);
|
||||
});
|
||||
|
||||
test('auth/SQL errors are NOT unreachable', () => {
|
||||
expect(isNetworkUnreachableError(new Error('password authentication failed for user "postgres"'))).toBe(false);
|
||||
expect(isNetworkUnreachableError(new Error('syntax error at or near "SELEC"'))).toBe(false);
|
||||
expect(isNetworkUnreachableError(null)).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe('ConnectionManager — direct-pool fallback on unreachable host (#1641)', () => {
|
||||
let originalKillSwitch: string | undefined;
|
||||
let originalError: typeof console.error;
|
||||
let errLines: string[];
|
||||
beforeEach(() => {
|
||||
originalKillSwitch = process.env.GBRAIN_DISABLE_DIRECT_POOL;
|
||||
delete process.env.GBRAIN_DISABLE_DIRECT_POOL;
|
||||
originalError = console.error;
|
||||
errLines = [];
|
||||
console.error = (...args: unknown[]) => { errLines.push(args.join(' ')); };
|
||||
});
|
||||
afterEach(() => {
|
||||
console.error = originalError;
|
||||
if (originalKillSwitch === undefined) delete process.env.GBRAIN_DISABLE_DIRECT_POOL;
|
||||
else process.env.GBRAIN_DISABLE_DIRECT_POOL = originalKillSwitch;
|
||||
});
|
||||
|
||||
test('ddl() falls back to the read pool when the direct host is unreachable', async () => {
|
||||
const cm = new ConnectionManager({
|
||||
url: 'postgresql://postgres.abc:p@aws.pooler.supabase.com:6543/db',
|
||||
// 127.0.0.1:9 (discard) → instant ECONNREFUSED, the IPv4-only-network shape.
|
||||
directUrl: 'postgresql://postgres:p@127.0.0.1:9/db',
|
||||
});
|
||||
const fakeReadPool = {} as ReturnType<typeof ConnectionManager.prototype.read>;
|
||||
cm.setReadPool(fakeReadPool);
|
||||
expect(cm.isDualPoolActive()).toBe(true);
|
||||
|
||||
const pool = await cm.ddl(); // without the fix this throws ECONNREFUSED
|
||||
expect(pool).toBe(fakeReadPool);
|
||||
// Self-activating kill-switch: subsequent calls skip the direct pool.
|
||||
expect(cm.isKillSwitchActive()).toBe(true);
|
||||
expect(cm.isDualPoolActive()).toBe(false);
|
||||
expect(cm.describeMode().mode).toBe('single (kill-switch)');
|
||||
// One stderr line mentioning the power-user override.
|
||||
const warning = errLines.filter(l => l.includes('GBRAIN_DIRECT_DATABASE_URL'));
|
||||
expect(warning.length).toBe(1);
|
||||
|
||||
const again = await cm.ddl();
|
||||
expect(again).toBe(fakeReadPool);
|
||||
expect(errLines.filter(l => l.includes('GBRAIN_DIRECT_DATABASE_URL')).length).toBe(1);
|
||||
}, 20000);
|
||||
});
|
||||
|
||||
@@ -78,8 +78,6 @@ describe('v0.36.1.x #1077 — admin register-client supports PKCE public clients
|
||||
// (under either name) from req.body. Pin the fallback pattern so the
|
||||
// PKCE-fix regression contract stays load-bearing.
|
||||
expect(src).toMatch(/req\.body[^;]*scopes\s*\?\?\s*[^;]*scope\b/);
|
||||
expect(src).toMatch(/req\.body[^;]*source_id\s*\?\?\s*[^;]*sourceId\b/);
|
||||
expect(src).toMatch(/req\.body[^;]*federated_read\s*\?\?\s*[^;]*federatedRead\b/);
|
||||
// v0.41.3 (T4 atomicity fix, codex F4): admin endpoint now validates
|
||||
// tokenEndpointAuthMethod via the shared validator and passes it to
|
||||
// registerClientManual as a positional arg. Pre-v0.41.3 the route did
|
||||
@@ -89,7 +87,6 @@ describe('v0.36.1.x #1077 — admin register-client supports PKCE public clients
|
||||
// UPDATE block (the regex deliberately asserts the post-insert UPDATE
|
||||
// is GONE).
|
||||
expect(src).toMatch(/validateTokenEndpointAuthMethod\(tokenEndpointAuthMethod\)/);
|
||||
expect(src).toMatch(/normalizeAdminClientSourceScope\(/);
|
||||
expect(src).toMatch(/registerClientManual\([^)]*validatedAuthMethod[^)]*\)/);
|
||||
// Regression guard: post-insert UPDATE flipping client_secret_hash to
|
||||
// NULL based on a runtime check is exactly the non-atomic pattern T4
|
||||
@@ -270,26 +267,3 @@ describe('v0.42.43.0 #2095 — volunteer-events sink + cycle purge wiring (struc
|
||||
expect(src).toMatch(/purged_volunteer_events_count/);
|
||||
});
|
||||
});
|
||||
|
||||
describe('#1490 + #1036 — admin Register Agent form exposes source scope + redirect URIs', () => {
|
||||
test('Agents.tsx RegisterModal sends source_id + federated_read in the POST body (#1490)', () => {
|
||||
const src = readFileSync('admin/src/pages/Agents.tsx', 'utf8');
|
||||
expect(src).toMatch(/source_id:\s*sourceId/);
|
||||
expect(src).toMatch(/federated_read:\s*federatedRead/);
|
||||
// The source list comes from the requireAdmin-gated endpoint, not a hardcoded 'default'.
|
||||
expect(src).toMatch(/api\.sources\(\)/);
|
||||
});
|
||||
|
||||
test('serve-http.ts serves /admin/api/sources behind requireAdmin (#1490)', () => {
|
||||
const src = readFileSync('src/commands/serve-http.ts', 'utf8');
|
||||
expect(src).toMatch(/app\.get\('\/admin\/api\/sources',\s*requireAdmin/);
|
||||
});
|
||||
|
||||
test('Agents.tsx RegisterModal sends redirectUris + authorization_code grants when URIs are entered (#1036)', () => {
|
||||
const src = readFileSync('admin/src/pages/Agents.tsx', 'utf8');
|
||||
// One-per-line textarea split into an array...
|
||||
expect(src).toMatch(/redirectUris\.split\('\\n'\)/);
|
||||
// ...included in the body with the CLI's grant-type inference convention.
|
||||
expect(src).toMatch(/redirectUris:\s*uris,\s*grantTypes:\s*\['authorization_code',\s*'refresh_token'\]/);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -125,30 +125,6 @@ describe('client registration', () => {
|
||||
expect(client!.client_name).toBe('test-agent');
|
||||
});
|
||||
|
||||
test('registerClientManual persists source_id and federated_read into AuthInfo', async () => {
|
||||
await sql`
|
||||
INSERT INTO sources (id, name) VALUES (${'team-a'}, ${'Team A'}), (${'team-b'}, ${'Team B'})
|
||||
ON CONFLICT DO NOTHING
|
||||
`;
|
||||
const { clientId, clientSecret } = await provider.registerClientManual(
|
||||
'source-scoped-test',
|
||||
['client_credentials'],
|
||||
'read write',
|
||||
[],
|
||||
'team-a',
|
||||
['team-a', 'team-b'],
|
||||
);
|
||||
|
||||
const tokens = await provider.exchangeClientCredentials(clientId, clientSecret!, 'read');
|
||||
const authInfo = await provider.verifyAccessToken(tokens.access_token) as {
|
||||
sourceId?: string;
|
||||
allowedSources?: string[];
|
||||
};
|
||||
|
||||
expect(authInfo.sourceId).toBe('team-a');
|
||||
expect(authInfo.allowedSources).toEqual(['team-a', 'team-b']);
|
||||
});
|
||||
|
||||
test('getClient returns undefined for unknown client', async () => {
|
||||
const client = await provider.clientsStore.getClient('nonexistent');
|
||||
expect(client).toBeUndefined();
|
||||
|
||||
Reference in New Issue
Block a user