Files
investor-flow/app/server/src/services/captureIngest.ts
T

153 lines
5.4 KiB
TypeScript
Raw Normal View History

// Investor Flow — X position-capture materializer (M21 follow-up).
//
// Reads a tracked fund manager's posts from x_cookie_posts (already fetched by
// the `x` timeline job), classifies them with captureParser, and upserts total
// position snapshots into fund_position_records with source='capture'.
//
// Recency rule (mirrorEngine): captures and 13F compete on MAX(as_of) per symbol.
// CLAIMS are explicitly skipped in v1: a claim ("added 865,000 shares") is a
// DELTA, not a total — letting it win recency would corrupt the Live Book.
// Claim folding (delta onto latest total) is future work.
//
// Local-only: no network, no rate limiting — this is a materializer, run from
// the x schedule branch after timeline jobs are enqueued.
import type { DatabaseSync } from 'node:sqlite';
import { extractCaptures } from '../mirror/captureParser.ts';
export interface CaptureIngestStats {
fundId: string;
postsScanned: number;
captures: number; // total-position posts materialized (inserted or refreshed)
inserted: number;
refreshed: number;
claims: number; // delta posts detected — recorded for observability, not materialized
skipped: number; // commentary / unparseable / missing symbol
disabled: boolean;
}
export interface FundWithHandle {
id: string;
x_handle: string | null;
enabled: number;
}
export function listTrackedFundHandles(db: DatabaseSync): FundWithHandle[] {
return db.prepare('SELECT id, x_handle, enabled FROM tracked_funds').all() as FundWithHandle[];
}
interface PostRow {
post_id: string;
body_text: string | null;
posted_at: string;
}
/**
* Materialize captures for one fund. Idempotent: converges on
* (fund_id, symbol, source='capture', evidence_url) — a re-run with the same
* post updates the row instead of duplicating it.
*/
export function ingestFundCaptures(db: DatabaseSync, fundId: string): CaptureIngestStats {
const fund = db.prepare('SELECT id, x_handle, enabled FROM tracked_funds WHERE id = ?').get(fundId) as
| FundWithHandle
| undefined;
const stats: CaptureIngestStats = {
fundId,
postsScanned: 0,
captures: 0,
inserted: 0,
refreshed: 0,
claims: 0,
skipped: 0,
disabled: false,
};
if (!fund || !fund.x_handle) return stats;
if (fund.enabled !== 1) { stats.disabled = true; return stats; }
const posts = db.prepare(
`SELECT post_id, body_text, posted_at FROM x_cookie_posts
WHERE lower(author_handle) = lower(?)
AND body_text IS NOT NULL AND body_text != ''
ORDER BY posted_at ASC`,
).all(fund.x_handle) as PostRow[];
stats.postsScanned = posts.length;
const selectExisting = db.prepare(
`SELECT id FROM fund_position_records WHERE fund_id=? AND symbol=? AND source='capture' AND evidence_url=?`,
);
const updateExisting = db.prepare(
`UPDATE fund_position_records SET shares=?, value_usd=?, cost_basis=?, as_of=?, notes=COALESCE(notes, ?)
WHERE id=?`,
);
const insertNew = db.prepare(
`INSERT INTO fund_position_records (id, fund_id, symbol, shares, value_usd, cost_basis, as_of, source, evidence_url, notes, created_at)
VALUES (?,?,?,?,?,?,?, 'capture', ?, ?, ?)`,
);
for (const p of posts) {
// One post can mention several names (Mike: SLNH fill + OPEN total in one tweet).
const parsedList = extractCaptures(p.body_text ?? '');
if (parsedList.length === 0) {
stats.skipped++;
continue;
}
let anyMaterialized = false;
for (const parsed of parsedList) {
if (parsed.class === 'claim') {
stats.claims++;
continue;
}
if (parsed.class !== 'capture') continue;
if (!parsed.symbol) continue;
const asOf = normalizePostedAtDate(p.posted_at);
const evidenceUrl = `https://x.com/${fund.x_handle}/status/${p.post_id}`;
// Preserve book_reset / pre_reset markers; only stamp instrument notes when empty.
const instrumentNote = parsed.instrument && parsed.instrument !== 'equity'
? parsed.instrument
: null;
const existing = selectExisting.get(fundId, parsed.symbol, evidenceUrl) as { id?: string } | undefined;
if (existing?.id) {
updateExisting.run(
parsed.shares ?? null, parsed.value_usd ?? null, parsed.cost_basis ?? null, asOf,
instrumentNote, existing.id,
);
stats.refreshed++;
} else {
insertNew.run(
crypto.randomUUID(), fundId, parsed.symbol,
parsed.shares ?? null, parsed.value_usd ?? null, parsed.cost_basis ?? null, asOf, evidenceUrl,
instrumentNote, new Date().toISOString(),
);
stats.inserted++;
}
stats.captures++;
anyMaterialized = true;
}
if (!anyMaterialized && parsedList.every((x) => x.class === 'claim')) {
// already counted as claims
} else if (!anyMaterialized) {
stats.skipped++;
}
}
return stats;
}
/** posted_at may be ISO or Twitter "Wed Jul 15 20:43:28 +0000 2026". */
function normalizePostedAtDate(postedAt: string): string {
if (/^\d{4}-\d{2}-\d{2}/.test(postedAt)) return postedAt.slice(0, 10);
const t = Date.parse(postedAt);
if (Number.isFinite(t)) return new Date(t).toISOString().slice(0, 10);
return postedAt.slice(0, 10);
}
/** Materialize captures for every enabled tracked fund with an x_handle. */
export function ingestAllFundCaptures(db: DatabaseSync): CaptureIngestStats[] {
return listTrackedFundHandles(db)
.filter((f) => f.x_handle && f.enabled === 1)
.map((f) => ingestFundCaptures(db, f.id));
}