feat: Unraid deploy, dealer-flow heatmap, confluence zones, 13F capture
CI / Test (push) Canceled after 0s
CI / Build and push (push) Canceled after 0s

Ship Node production images, Unraid compose, and Gitea CI/CD (test then
push registry images; cron script if no runner). Rebuild dealer flow as a
heatmap-first map with integrity gates and chart helpers. Add confluence
zone rules, session clock, capture evidence, and tighter 13F/queue/options
paths, plus the matching UI and tests.
This commit is contained in:
Investor Flow Build
2026-08-18 14:10:02 -04:00
parent f7efea2178
commit 90e1829d39
113 changed files with 7606 additions and 1947 deletions
+200 -19
View File
@@ -39,7 +39,16 @@ import { ConfluenceRepository, rackFromSlots } from '../db/confluenceRepository.
import { runSlotBacktest, signalHistoryToStats, type ConfluenceFireEvent } from '../confluence/confluenceBacktest.ts';
import { detectPictureChange } from '../confluence/confluenceRack.ts';
import { runCorridorBacktest, aggregateBacktest } from '../analysis/corridorBacktest.ts';
import { runConfluenceEvaluationCycle } from '../confluence/confluenceEngine.ts';
import {
runConfluenceEvaluationCycle,
runConfluenceReplay,
replayCoverage,
loadActivePredicates,
buildRecentZones,
scoreSymbolUnderPrior,
learningLedgerStatus,
} from '../confluence/confluenceEngine.ts';
import { REPLAY_LOOKBACK_DAYS } from '../confluence/confluenceSeed.ts';
import { CorridorRepository } from '../db/corridorRepository.ts';
// ---------------------------------------------------------------------------
@@ -376,12 +385,14 @@ const marketRouter = router({
const byKey = new Map(entries.map((e) => [e.key, e]));
const val = <T>(key: string): T | null => (byKey.get(key)?.value ?? null) as T | null;
const stale = (key: string): boolean => byKey.get(key)?.isStale ?? true;
const fetchedAt = (key: string): string | null => byKey.get(key)?.fetchedAt ?? null;
return {
symbol,
quote: val<Quote>(k.quote),
candles: val<PriceCandle[]>(k.candles),
sector: val<SymbolMeta>(k.sector),
stale: { quote: stale(k.quote), candles: stale(k.candles), sector: stale(k.sector) },
observed: { quote: fetchedAt(k.quote), candles: fetchedAt(k.candles), sector: fetchedAt(k.sector) },
};
}),
@@ -411,6 +422,7 @@ const marketRouter = router({
candles: PriceCandle[] | null;
sector: SymbolMeta | null;
stale: { quote: boolean; candles: boolean; sector: boolean };
observed: { quote: string | null; candles: string | null; sector: string | null };
}> = [];
for (const symbol of symbols) {
@@ -428,6 +440,11 @@ const marketRouter = router({
candles: byKey.get(kCandles)?.isStale ?? true,
sector: byKey.get(kSector)?.isStale ?? true,
},
observed: {
quote: byKey.get(kQuote)?.fetchedAt ?? null,
candles: byKey.get(kCandles)?.fetchedAt ?? null,
sector: byKey.get(kSector)?.fetchedAt ?? null,
},
});
}
@@ -617,17 +634,23 @@ const marketRouter = router({
}),
candles: publicProcedure
.input(z.object({ symbol: z.string().min(1), timeframe: z.enum(['1d', '1wk', '1mo']).default('1d') }))
.input(z.object({ symbol: z.string().min(1), timeframe: z.enum(['1m', '5m', '1d', '1wk', '1mo']).default('1d') }))
.query(async ({ ctx, input }) => {
const symbol = input.symbol.toUpperCase();
const key = `yfinance:candles:${symbol}:${input.timeframe}`;
const entry = await ctx.cache.get<PriceCandle[]>(key);
return { symbol, timeframe: input.timeframe, candles: (entry.value ?? []), isStale: entry.isStale };
return {
symbol,
timeframe: input.timeframe,
candles: (entry.value ?? []),
isStale: entry.isStale,
fetchedAt: entry.provenance?.fetchedAt ?? null,
};
}),
indicators: publicProcedure
.input(z.object({
symbol: z.string().min(1),
timeframe: z.enum(['1d', '1wk', '1mo']).default('1d'),
timeframe: z.enum(['1m', '5m', '1d', '1wk', '1mo']).default('1d'),
periods: z.object({
ema: z.array(z.number().int()).default([9, 21, 50, 200]),
rsi: z.number().int().default(14),
@@ -1859,8 +1882,41 @@ const adminRouter = router({
}),
xAccountsList: adminProcedure.query(({ ctx }) => {
const rows = ctx.db.prepare("SELECT id, symbol, handle, COALESCE(label, '') AS label, created_at FROM x_accounts ORDER BY symbol, handle").all();
return (rows ?? []) as Array<{id: string; symbol: string; handle: string; label: string; created_at: string}>;
const rows = ctx.db.prepare("SELECT id, symbol, handle, COALESCE(label, '') AS label, created_at FROM x_accounts ORDER BY symbol, handle").all() as Array<{id: string; symbol: string; handle: string; label: string; created_at: string}>;
const handles = [...new Set(rows.map((r) => r.handle))];
const pullByHandle = new Map<string, { status: string; lastAttempt: string | null; error: string | null }>();
const lastPostByHandle = new Map<string, string>();
if (handles.length > 0) {
const keys = handles.flatMap((h) => [`x:timeline:${h}`, `x:timeline:${h.toLowerCase()}`]);
const jobs = ctx.db.prepare(
`SELECT key, status, last_attempt, error FROM adapter_queue WHERE key IN (${keys.map(() => '?').join(',')})`,
).all(...keys) as Array<{ key: string; status: string; last_attempt: string | null; error: string | null }>;
for (const j of jobs) {
pullByHandle.set(j.key.slice('x:timeline:'.length), {
status: j.status,
lastAttempt: j.last_attempt,
error: j.error,
});
}
const posts = ctx.db.prepare(
`SELECT lower(author_handle) AS h, MAX(posted_at) AS last_post
FROM x_cookie_posts WHERE lower(author_handle) IN (${handles.map(() => '?').join(',')})
GROUP BY 1`,
).all(...handles.map((h) => h.toLowerCase())) as Array<{ h: string; last_post: string | null }>;
for (const p of posts) {
if (p.last_post) lastPostByHandle.set(p.h, p.last_post);
}
}
return rows.map((r) => {
const pull = pullByHandle.get(r.handle) ?? pullByHandle.get(r.handle.toLowerCase()) ?? null;
return {
...r,
pullStatus: pull?.status ?? null,
pullAt: pull?.lastAttempt ?? null,
pullError: pull?.error ?? null,
lastPostAt: lastPostByHandle.get(r.handle.toLowerCase()) ?? null,
};
});
}),
xAccountAdd: adminProcedure
@@ -4604,6 +4660,8 @@ const xRouter = router({
nextCursor: null,
configured: true as const,
health: null,
lastPullAt: null,
queuedHandles: 0,
};
}
@@ -4653,6 +4711,17 @@ const xRouter = router({
const healthRow = ctx.db.prepare('SELECT healthy, last_error, updated_at FROM x_credentials WHERE id=?').get('singleton') as { healthy?: number; last_error?: string | null } | undefined;
const pullKeys = handles.flatMap((h) => [`x:timeline:${h}`, `x:timeline:${h.toLowerCase()}`]);
const pullJobs = pullKeys.length
? ctx.db.prepare(
`SELECT key, status, last_attempt, error FROM adapter_queue WHERE key IN (${pullKeys.map(() => '?').join(',')})`,
).all(...pullKeys) as Array<{ key: string; status: string; last_attempt: string | null; error: string | null }>
: [];
const lastPullAt = pullJobs.map((j) => j.last_attempt).filter((t): t is string => Boolean(t)).sort().at(-1) ?? null;
const queuedHandles = new Set(
pullJobs.filter((j) => j.status === 'pending' && !j.last_attempt).map((j) => j.key.slice('x:timeline:'.length).toLowerCase()),
).size;
return {
cashtagPosts: page,
accountPosts: [],
@@ -4662,6 +4731,8 @@ const xRouter = router({
healthy: healthRow.healthy === 1 ? ('healthy' as const) : ('degraded' as const),
lastError: healthRow.last_error ?? null,
} : null,
lastPullAt,
queuedHandles,
};
}),
@@ -5238,17 +5309,22 @@ const symbolsRouter = router({
// Tier 1 — tracked funds holding the symbol (most recent record per fund), with weight in their disclosed book.
const trackedRows = db.prepare(
`SELECT tf.id AS fund_id, tf.fund_name, fpr.symbol, fpr.shares, fpr.value_usd, fpr.as_of, fpr.source
`SELECT tf.id AS fund_id, tf.fund_name, fpr.symbol, fpr.shares, fpr.value_usd,
fpr.as_of, fpr.source, fpr.notes
FROM fund_position_records fpr
JOIN tracked_funds tf ON tf.id = fpr.fund_id AND tf.enabled = 1
JOIN (
SELECT fund_id, MAX(as_of) AS max_as_of
SELECT fund_id,
CASE WHEN notes IN ('call','put') THEN notes ELSE '' END AS inst,
MAX(as_of) AS max_as_of
FROM fund_position_records
WHERE symbol = ?
GROUP BY fund_id
GROUP BY fund_id, CASE WHEN notes IN ('call','put') THEN notes ELSE '' END
) latest ON latest.fund_id = fpr.fund_id AND latest.max_as_of = fpr.as_of
AND CASE WHEN fpr.notes IN ('call','put') THEN fpr.notes ELSE '' END = latest.inst
WHERE fpr.symbol = ?
ORDER BY tf.fund_name ASC`,
AND (fpr.shares IS NULL OR fpr.shares > 0)
ORDER BY tf.fund_name ASC, fpr.notes ASC`,
).all(symbol, symbol) as Array<Record<string, any>>;
// Weight = position value / total disclosed book value at the fund's live book (per fund).
@@ -5268,6 +5344,7 @@ const symbolsRouter = router({
valueUsd: r.value_usd ?? null,
asOf: r.as_of,
source: r.source,
notes: (r.notes as string | null) ?? null,
weightPct: bookValueByFund.get(r.fund_id)
? ((r.value_usd ?? 0) / bookValueByFund.get(r.fund_id)!) * 100
: null,
@@ -5290,16 +5367,20 @@ const fundsRouter = router({
/** List operator-curated tracked funds (v1: Alpine Fox). */
list: publicProcedure
.query(async ({ ctx }) => {
const { listTrackedFunds } = await import('../db/fundRepository.ts');
return listTrackedFunds(ctx.db, { includeDisabled: true });
const { listTrackedFunds, fundsFreshness } = await import('../db/fundRepository.ts');
const funds = listTrackedFunds(ctx.db, { includeDisabled: true });
const fresh = fundsFreshness(ctx.db, funds);
return funds.map((f) => ({ ...f, freshness: fresh.get(f.id) ?? null }));
}),
/** Get a single tracked fund by id. */
get: publicProcedure
.input(z.object({ id: z.string().min(1) }))
.query(async ({ ctx, input }) => {
const { getTrackedFund } = await import('../db/fundRepository.ts');
return getTrackedFund(ctx.db, input.id);
const { getTrackedFund, fundFreshness } = await import('../db/fundRepository.ts');
const fund = getTrackedFund(ctx.db, input.id);
if (!fund) return null;
return { ...fund, freshness: fundFreshness(ctx.db, fund) };
}),
/** Live Book for a fund — most recent record per symbol, source-labeled. */
@@ -5307,7 +5388,20 @@ const fundsRouter = router({
.input(z.object({ fundId: z.string().min(1) }))
.query(async ({ ctx, input }) => {
const { liveBook } = await import('../db/fundRepository.ts');
return liveBook(ctx.db, input.fundId);
const book = liveBook(ctx.db, input.fundId);
if (ctx.xAdapter) {
const { annotateCaptureEvidence } = await import('../services/captureEvidence.ts');
// Cache only on the request path. Live bird reads here were serial and
// routinely blew the 8s client timeout, so the fund page rendered empty.
await annotateCaptureEvidence(
ctx.db,
book,
(id) => ctx.xAdapter!.readTweetStatus(id),
Date.now(),
{ live: false },
);
}
return book;
}),
/** Full append-only position timeline for a fund. */
@@ -5330,7 +5424,11 @@ const fundsRouter = router({
}))
.mutation(async ({ ctx, input }) => {
const { upsertTrackedFund } = await import('../db/fundRepository.ts');
return upsertTrackedFund(ctx.db, input);
const fund = upsertTrackedFund(ctx.db, input);
if (fund.enabled && fund.x_handle) {
try { await ctx.queue.queue(`x:timeline:${fund.x_handle}`); } catch { /* hourly schedule is the fallback */ }
}
return fund;
}),
adminSetEnabled: adminProcedure
@@ -5530,11 +5628,90 @@ const confluenceRouter = router({
body: (meta?.body ?? 'bull') as 'bull' | 'bear' | 'exit',
firedAt: f.firedAt,
qualityAtFire: f.qualityAtFire,
verdict: f.verdict,
priceConfirmed: f.priceConfirmed,
};
});
return { fires };
}),
/** Last few entry/exit windows for the open symbol, plus the active derived rules. */
recentZones: publicProcedure
.input(z.object({
symbol: z.string().min(1).max(12),
rackId: z.string().min(1).optional(),
entries: z.number().int().min(1).max(8).optional(),
exits: z.number().int().min(1).max(8).optional(),
}))
.query(async ({ ctx, input }) => {
const repo = new ConfluenceRepository(ctx.db);
const symbol = input.symbol.toUpperCase();
const rackId = input.rackId ?? (repo.listSystemRacks()[0]?.id ?? null);
if (!rackId) {
return {
symbol,
rackId: null,
entries: [],
exits: [],
entryRule: { label: 'n/a', source: 'baseline', derivedAt: null, sampleCaveat: 'No rack seeded.', train: null, validate: null },
exitRule: { label: 'n/a', source: 'baseline', derivedAt: null, sampleCaveat: 'No rack seeded.', train: null, validate: null },
symbolScore: {
entry: { zones: 0, resolved: 0, confirmed: 0, hitRate: null },
exit: { zones: 0, resolved: 0, confirmed: 0, hitRate: null },
},
learningLedger: { ready: false, filled: 0, total: 0 },
coverage: { evaluatedDays: 0, lookbackDays: REPLAY_LOOKBACK_DAYS, replayComplete: false },
};
}
const entry = await ctx.cache.get<PriceCandle[]>(`yfinance:candles:${symbol}:1d`);
const candles = (entry?.value ?? []) as PriceCandle[];
const zones = buildRecentZones(ctx.db, symbol, rackId, candles, input.entries ?? 3, input.exits ?? 3);
const preds = loadActivePredicates(ctx.db, rackId);
const coverage = replayCoverage(ctx.db, symbol, rackId, REPLAY_LOOKBACK_DAYS, candles.length);
if (!coverage.replayComplete) {
void runConfluenceReplay(ctx.db, ctx.cache, {
symbols: [symbol],
symbolsPerTick: 1,
budgetDaysPerSymbol: 80,
}).catch((e) => console.error('[confluence] on-read replay failed:', e));
}
const entryScore = scoreSymbolUnderPrior(ctx.db, symbol, rackId, candles, 'entry', preds.entry, preds.exit);
const exitScore = scoreSymbolUnderPrior(ctx.db, symbol, rackId, candles, 'exit', preds.exit, preds.entry);
const ledger = await learningLedgerStatus(ctx.db, ctx.cache);
const ruleCard = (
stored: typeof preds.entryRule,
pred: typeof preds.entry,
side: 'entry' | 'exit',
) => stored ? {
label: pred.label,
source: stored.source,
derivedAt: stored.derivedAt,
sampleCaveat: stored.sampleCaveat,
train: JSON.parse(stored.trainStatsJson),
validate: JSON.parse(stored.validateStatsJson),
} : {
label: pred.label,
source: 'baseline',
derivedAt: null,
sampleCaveat: ledger.ready
? 'Learning ledger is full. A sector-checked prior has not been stored yet.'
: 'Default window definition. Sector-checked prior is not ready (learning ledger still filling).',
train: null,
validate: null,
};
return {
symbol,
rackId,
entries: zones.entries,
exits: zones.exits,
entryRule: ruleCard(preds.entryRule, preds.entry, 'entry'),
exitRule: ruleCard(preds.exitRule, preds.exit, 'exit'),
symbolScore: { entry: entryScore, exit: exitScore },
learningLedger: { ready: ledger.ready, filled: ledger.symbols.filter((s) => s.ready).length, total: ledger.symbols.length },
coverage,
};
}),
/** Create or update a user-owned rack. */
saveRack: protectedProcedure
.input(z.object({
@@ -5662,10 +5839,14 @@ const confluenceRouter = router({
runEvaluationNow: protectedProcedure
.input(z.object({ symbol: z.string().min(1).max(12).optional() }))
.mutation(async ({ ctx, input }) => {
const summary = await runConfluenceEvaluationCycle(ctx.db, ctx.cache, {
symbols: input?.symbol ? [input.symbol] : undefined,
const symbols = input?.symbol ? [input.symbol] : undefined;
const replay = await runConfluenceReplay(ctx.db, ctx.cache, {
symbols: symbols ?? undefined,
symbolsPerTick: symbols ? 1 : 3,
budgetDaysPerSymbol: symbols ? 200 : 40,
});
return summary;
const summary = await runConfluenceEvaluationCycle(ctx.db, ctx.cache, { symbols });
return { ...summary, replay };
}),
});