Files
gbrain/src/core/parallel.ts
T
df86ea5f1d v0.40.5.0 Federated Sync v2 — parallel source sync + push triggers + per-source health (#1322)
* wip: federated sync v2 pre-merge snapshot

* v0.40.5.0 Federated Sync v2 — parallel source sync + push triggers + per-source health

Bump VERSION + package.json + CHANGELOG header + migration walkthrough filename
to v0.40.5.0 (claiming the next free slot in the v0.40.x patch series after
master's v0.40.1.0).

What ships (6 components, all behind sync.federated_v2 feature flag default-on):
1. Per-source sync lock — syncLockId(sourceId), phantom-redirect parity
2. Parallel sync --all — pMapAllSettled fan-out, --max-sources N cap
3. embed-backfill minion handler — D2 per-source lock + D6 $10/job budget + D15.1
   fire-and-forget submission + D19 source-level cooldown + 24h $25 rolling cap
4. sync trigger CLI + POST /webhooks/github — HMAC-verified (60 req/min/IP),
   X-GitHub-Event=push + ref filter against tracked_branch
5. sources status + federation_health doctor — batched GROUP BY pipeline
   (4 queries instead of 6×N per-source roundtrips)
6. sources federate/unfederate hook — auto-submit embed-backfill on flip

Correctness fixes (unconditional):
- D21: sync.ts:959 facts backstop now passes sourceId to engine.getPage
- D15.4: redactSourceConfig + CI guard prevent webhook_secret leak
- D15.5: safeHexEqual extracted to src/core/timing-safe.ts

Schema:
- Migration v89 (sources_github_repo_index): partial expression index on
  config->>'github_repo' for fast webhook source-lookup

Tests:
- 14 new test files, 112 cases. 4 IRON-RULE regressions pinned (SYNC_LOCK_ID
  back-compat, phantom per-source lock, embed-backfill kill+resume,
  webhook HMAC prefix-strip). All 9449 unit tests pass.

Caught at test-write time: the webhook handler had a Buffer.from('sha256=...',
'hex') truncation bug — without the prefix-strip, every signature would have
"matched" empty buffers. Pinned by a test/sources-webhook.test.ts IRON-RULE.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* fix(check-source-config-leak): tighten regex to source-row patterns only

The v0.40.5.0 wave added scripts/check-source-config-leak.sh with a
too-broad pattern (JSON\.stringify\(.*config) that flagged any variable
named 'config' — catching the GLOBAL gbrain config.json serializers in
src/commands/init.ts (status envelopes) and src/core/config.ts (the
config-file write site). On the CI runner without rg installed, the
grep -rE fallback fired correctly and produced 4 false positives that
broke the `verify` script.

Tightened the patterns to specifically match `(source|src|row|s).config`
property access — the actual risk shape (a sources-table row being
serialized whole). The global gbrain config has a different shape and
threat model (file-mode 0o600 at the write site), so it's safe to
exempt at the regex level rather than per-file whitelist.

Also fixed a latent bug: the rg branch used `--include='*.ts'` (grep's
flag, not rg's). rg silently rejected it and CANDIDATES came back empty,
so the local-dev runs (which have rg) would never have caught a real
leak. Now branches on tool availability: `-g '*.ts'` for rg, `--include`
for grep -rE. Both branches verified against a synthetic leak fixture.

Also added init.ts + config.ts to the whitelist as a belt-and-suspenders
since they handle gbrain-global config (not source rows) and could
otherwise reflect-back via regex iteration.

CI: `bun run verify` exit 0 locally with both the original false-positive
fixture (clean repo) and a synthetic leak fixture (correctly caught,
exit 1).

---------

Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-23 10:21:59 -07:00

60 lines
2.2 KiB
TypeScript

/**
* Bounded-concurrency Promise.allSettled.
*
* Runs `fn(item)` for each item with at most `concurrency` in flight at a time.
* Returns a `PromiseSettledResult` per input item, in input order, so callers
* can distinguish fulfilled-with-result from rejected-with-error per item.
*
* Used by `gbrain sync --all` (v0.40 Federated Sync v2) to fan out per-source
* syncs without overwhelming the embedding API or local disk.
*
* Why a semaphore + `Promise.allSettled` instead of a library:
* - Zero new deps. Bun's stdlib has every primitive we need.
* - `Promise.allSettled` semantics are what we want: one source failing
* must NOT short-circuit the others. Bare `Promise.all` rejects on
* first failure.
* - The concurrency cap is a hard ceiling, not a target. With N items and
* concurrency C, the function makes at most C concurrent `fn` calls at
* any moment.
*
* Order guarantee: results[i] always corresponds to items[i]. The execution
* order of `fn` calls is NOT guaranteed (a fast item submitted later can
* complete before a slow item submitted earlier).
*/
export async function pMapAllSettled<T, R>(
items: readonly T[],
concurrency: number,
fn: (item: T, index: number) => Promise<R>,
): Promise<Array<PromiseSettledResult<R>>> {
if (!Number.isInteger(concurrency) || concurrency < 1) {
throw new TypeError(
`pMapAllSettled: concurrency must be a positive integer, got ${concurrency}`,
);
}
if (items.length === 0) return [];
const results: Array<PromiseSettledResult<R>> = new Array(items.length);
const effectiveConcurrency = Math.min(concurrency, items.length);
let nextIndex = 0;
// Each worker pulls the next available index and runs fn until exhausted.
async function worker(): Promise<void> {
for (;;) {
const i = nextIndex++;
if (i >= items.length) return;
try {
const value = await fn(items[i], i);
results[i] = { status: 'fulfilled', value };
} catch (err) {
results[i] = { status: 'rejected', reason: err };
}
}
}
const workers = Array.from({ length: effectiveConcurrency }, () => worker());
await Promise.all(workers);
return results;
}