import { type DatabaseClient, getDatabase } from "./database"; import { getObjectStore } from "./object-store"; import { getRuntimeEnv, isEnabled } from "./runtime-env"; import feishuSnapshot from "./feishu-source-snapshot.json"; type D1ResultRow = Record; export function getRawDb(): DatabaseClient { return getDatabase(); } export function getUploadBucket() { return getObjectStore(); } export async function ensureSchema(database?: DatabaseClient) { const db = database ?? getRawDb(); { await db.prepare("SELECT id FROM tasks LIMIT 1").all(); const tasksWithoutShare = await db .prepare("SELECT id FROM tasks WHERE share_token IS NULL OR share_token = ''") .all<{ id: string }>(); for (const task of tasksWithoutShare.results) { await db .prepare("UPDATE tasks SET share_token = ? WHERE id = ?") .bind(crypto.randomUUID().replaceAll("-", ""), task.id) .run(); } await db .prepare( `UPDATE distributions SET latest_likes = COALESCE(d7_likes, d5_likes, d2_likes), latest_comments = COALESCE(d7_comments, d5_comments, d2_comments), latest_collects = COALESCE(d7_collects, d5_collects, d2_collects), collection_status = 'success', collection_status_description = '历史采集数据已迁移', collection_updated_at = updated_at, last_collection_day = CASE WHEN d7_likes IS NOT NULL THEN 7 WHEN d5_likes IS NOT NULL THEN 5 WHEN d2_likes IS NOT NULL THEN 2 ELSE NULL END WHERE latest_likes IS NULL AND COALESCE(d7_likes, d5_likes, d2_likes) IS NOT NULL`, ) .run(); const tasksNeedingImages = await db .prepare( `SELECT DISTINCT t.id FROM tasks t JOIN contents c ON c.task_id = t.id WHERE t.source_sheet_id = ? AND (c.image_assets IS NULL OR c.image_assets = '' OR c.image_assets = '[]')`, ) .bind(feishuSnapshot.sheetId) .all<{ id: string }>(); for (const task of tasksNeedingImages.results) { const updates = feishuSnapshot.rows .filter((row) => row.images.length > 0) .map((row) => db .prepare( `UPDATE contents SET image_assets = ? WHERE task_id = ? AND source_row = ? AND (image_assets IS NULL OR image_assets = '' OR image_assets = '[]')`, ) .bind(JSON.stringify(row.images), task.id, row.sourceRow), ); if (updates.length > 0) await db.batch(updates); } } return; const statements = [ `CREATE TABLE IF NOT EXISTS partners ( id TEXT PRIMARY KEY, name TEXT NOT NULL, wecom_name TEXT NOT NULL, owner TEXT NOT NULL DEFAULT '运营组', claimed_total INTEGER NOT NULL DEFAULT 0, completed_total INTEGER NOT NULL DEFAULT 0, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS tasks ( id TEXT PRIMARY KEY, name TEXT NOT NULL, brand TEXT NOT NULL, quantity INTEGER NOT NULL, claimed_quantity INTEGER NOT NULL DEFAULT 0, due_at TEXT NOT NULL, status TEXT NOT NULL DEFAULT 'active', task_type TEXT NOT NULL DEFAULT 'content_publish', source_url TEXT NOT NULL DEFAULT '', source_sheet_id TEXT NOT NULL DEFAULT '', source_sheet_name TEXT NOT NULL DEFAULT '', source_synced_at TEXT, share_token TEXT, collection_start_date TEXT, collection_days TEXT NOT NULL DEFAULT '[]', collection_schedule_updated_at TEXT, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS contents ( id TEXT PRIMARY KEY, task_id TEXT NOT NULL, title TEXT NOT NULL, body TEXT NOT NULL DEFAULT '', image_assets TEXT NOT NULL DEFAULT '[]', status TEXT NOT NULL DEFAULT 'available', source TEXT NOT NULL DEFAULT '飞书内容表', source_row INTEGER, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS accounts ( id TEXT PRIMARY KEY, platform TEXT NOT NULL DEFAULT '小红书', platform_uid TEXT NOT NULL, public_account_id TEXT NOT NULL DEFAULT '', nickname TEXT NOT NULL, profile_url TEXT NOT NULL DEFAULT '', ip_location TEXT NOT NULL DEFAULT '待识别', followers INTEGER NOT NULL DEFAULT 0, post_count INTEGER NOT NULL DEFAULT 0, avg_views INTEGER NOT NULL DEFAULT 0, cooperation_source TEXT NOT NULL DEFAULT '', first_seen_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, last_seen_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE UNIQUE INDEX IF NOT EXISTS accounts_platform_uid_idx ON accounts(platform, platform_uid)`, `CREATE TABLE IF NOT EXISTS claims ( id TEXT PRIMARY KEY, task_id TEXT NOT NULL, partner_id TEXT NOT NULL, claimant_name TEXT NOT NULL, claim_token TEXT NOT NULL, quantity INTEGER NOT NULL, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS delegation_bundles ( id TEXT PRIMARY KEY, task_id TEXT NOT NULL, claim_id TEXT NOT NULL, partner_id TEXT NOT NULL, label TEXT NOT NULL, share_token TEXT NOT NULL, quantity INTEGER NOT NULL, status TEXT NOT NULL DEFAULT 'active', created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, revoked_at TEXT )`, `CREATE TABLE IF NOT EXISTS distributions ( id TEXT PRIMARY KEY, task_id TEXT NOT NULL, content_id TEXT NOT NULL, partner_id TEXT NOT NULL, claim_id TEXT, delegation_bundle_id TEXT, account_id TEXT, publish_url TEXT, publish_time TEXT, publish_screenshot_key TEXT, result_screenshot_key TEXT, result_submitted_at TEXT, status TEXT NOT NULL DEFAULT 'claimed', claimed_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, screenshot_key TEXT, ocr_status TEXT NOT NULL DEFAULT 'none', exposure INTEGER, views INTEGER, d2_likes INTEGER, d2_comments INTEGER, d2_collects INTEGER, d5_likes INTEGER, d5_comments INTEGER, d5_collects INTEGER, d7_likes INTEGER, d7_comments INTEGER, d7_collects INTEGER, latest_likes INTEGER, latest_comments INTEGER, latest_collects INTEGER, collection_status TEXT NOT NULL DEFAULT 'pending', collection_status_description TEXT, collection_updated_at TEXT, last_collection_day INTEGER, updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS collection_runs ( id TEXT PRIMARY KEY, task_id TEXT NOT NULL, distribution_id TEXT NOT NULL, scheduled_date TEXT NOT NULL, schedule_day INTEGER, scheduled_at TEXT NOT NULL, status TEXT NOT NULL DEFAULT 'pending', likes INTEGER, comments INTEGER, collects INTEGER, status_description TEXT, started_at TEXT, completed_at TEXT, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS users ( id TEXT PRIMARY KEY, username TEXT NOT NULL, password_hash TEXT NOT NULL, password_salt TEXT NOT NULL, password_iterations INTEGER NOT NULL, role TEXT NOT NULL, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS auth_sessions ( token_hash TEXT PRIMARY KEY, user_id TEXT NOT NULL, expires_at TEXT NOT NULL, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, `CREATE TABLE IF NOT EXISTS mcp_export_tokens ( token_hash TEXT PRIMARY KEY, kind TEXT NOT NULL, payload TEXT NOT NULL, expires_at TEXT NOT NULL, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP )`, ]; for (const statement of statements) { await db.prepare(statement).run(); } const ensureColumn = async ( table: string, column: string, definition: string, ) => { const info = await db .prepare(`PRAGMA table_info(${table})`) .all<{ name: string }>(); if (!info.results.some((item) => item.name === column)) { await db.prepare(`ALTER TABLE ${table} ADD COLUMN ${definition}`).run(); } }; await ensureColumn("tasks", "source_url", "source_url TEXT NOT NULL DEFAULT ''"); await ensureColumn( "tasks", "task_type", "task_type TEXT NOT NULL DEFAULT 'content_publish'", ); await ensureColumn( "tasks", "source_sheet_id", "source_sheet_id TEXT NOT NULL DEFAULT ''", ); await ensureColumn( "tasks", "source_sheet_name", "source_sheet_name TEXT NOT NULL DEFAULT ''", ); await ensureColumn("tasks", "source_synced_at", "source_synced_at TEXT"); await ensureColumn("tasks", "share_token", "share_token TEXT"); await ensureColumn( "tasks", "collection_start_date", "collection_start_date TEXT", ); await ensureColumn( "tasks", "collection_days", "collection_days TEXT NOT NULL DEFAULT '[]'", ); await ensureColumn( "tasks", "collection_schedule_updated_at", "collection_schedule_updated_at TEXT", ); await ensureColumn("contents", "source_row", "source_row INTEGER"); await ensureColumn( "contents", "image_assets", "image_assets TEXT NOT NULL DEFAULT '[]'", ); await ensureColumn( "accounts", "public_account_id", "public_account_id TEXT NOT NULL DEFAULT ''", ); await ensureColumn("distributions", "claim_id", "claim_id TEXT"); await ensureColumn( "distributions", "delegation_bundle_id", "delegation_bundle_id TEXT", ); await ensureColumn( "distributions", "publish_screenshot_key", "publish_screenshot_key TEXT", ); await ensureColumn( "distributions", "result_screenshot_key", "result_screenshot_key TEXT", ); await ensureColumn( "distributions", "result_submitted_at", "result_submitted_at TEXT", ); await ensureColumn( "distributions", "latest_likes", "latest_likes INTEGER", ); await ensureColumn( "distributions", "latest_comments", "latest_comments INTEGER", ); await ensureColumn( "distributions", "latest_collects", "latest_collects INTEGER", ); await ensureColumn( "distributions", "collection_status", "collection_status TEXT NOT NULL DEFAULT 'pending'", ); await ensureColumn( "distributions", "collection_status_description", "collection_status_description TEXT", ); await ensureColumn( "distributions", "collection_updated_at", "collection_updated_at TEXT", ); await ensureColumn( "distributions", "last_collection_day", "last_collection_day INTEGER", ); const tasksWithoutShare = await db .prepare( "SELECT id FROM tasks WHERE share_token IS NULL OR share_token = ''", ) .all<{ id: string }>(); for (const task of tasksWithoutShare.results) { await db .prepare("UPDATE tasks SET share_token = ? WHERE id = ?") .bind(crypto.randomUUID().replaceAll("-", ""), task.id) .run(); } await db .prepare( "CREATE UNIQUE INDEX IF NOT EXISTS tasks_share_token_idx ON tasks(share_token)", ) .run(); await db .prepare( "CREATE UNIQUE INDEX IF NOT EXISTS claims_claim_token_idx ON claims(claim_token)", ) .run(); await db .prepare( `CREATE INDEX IF NOT EXISTS claims_task_partner_created_idx ON claims(task_id, partner_id, created_at)`, ) .run(); await db .prepare( `CREATE UNIQUE INDEX IF NOT EXISTS delegation_bundles_share_token_idx ON delegation_bundles(share_token)`, ) .run(); await db .prepare( `CREATE INDEX IF NOT EXISTS delegation_bundles_claim_created_idx ON delegation_bundles(claim_id, created_at)`, ) .run(); await db .prepare( `CREATE UNIQUE INDEX IF NOT EXISTS collection_runs_distribution_date_idx ON collection_runs(distribution_id, scheduled_date)`, ) .run(); await db .prepare( `CREATE INDEX IF NOT EXISTS collection_runs_task_date_idx ON collection_runs(task_id, scheduled_date)`, ) .run(); await db .prepare("CREATE UNIQUE INDEX IF NOT EXISTS users_username_idx ON users(username)") .run(); await db .prepare( `CREATE UNIQUE INDEX IF NOT EXISTS users_single_super_admin_idx ON users(role) WHERE role = 'super_admin'`, ) .run(); await db .prepare( "CREATE INDEX IF NOT EXISTS auth_sessions_user_id_idx ON auth_sessions(user_id)", ) .run(); await db .prepare( "CREATE INDEX IF NOT EXISTS auth_sessions_expires_at_idx ON auth_sessions(expires_at)", ) .run(); await db .prepare( "CREATE INDEX IF NOT EXISTS mcp_export_tokens_expires_at_idx ON mcp_export_tokens(expires_at)", ) .run(); await db .prepare( `UPDATE distributions SET latest_likes = COALESCE(d7_likes, d5_likes, d2_likes), latest_comments = COALESCE(d7_comments, d5_comments, d2_comments), latest_collects = COALESCE(d7_collects, d5_collects, d2_collects), collection_status = 'success', collection_status_description = '历史采集数据已迁移', collection_updated_at = updated_at, last_collection_day = CASE WHEN d7_likes IS NOT NULL THEN 7 WHEN d5_likes IS NOT NULL THEN 5 WHEN d2_likes IS NOT NULL THEN 2 ELSE NULL END WHERE latest_likes IS NULL AND COALESCE(d7_likes, d5_likes, d2_likes) IS NOT NULL`, ) .run(); const tasksNeedingImages = await db .prepare( `SELECT DISTINCT t.id FROM tasks t JOIN contents c ON c.task_id = t.id WHERE t.source_sheet_id = ? AND (c.image_assets IS NULL OR c.image_assets = '' OR c.image_assets = '[]')`, ) .bind(feishuSnapshot.sheetId) .all<{ id: string }>(); for (const task of tasksNeedingImages.results) { const updates = feishuSnapshot.rows .filter((row) => row.images.length > 0) .map((row) => db .prepare( `UPDATE contents SET image_assets = ? WHERE task_id = ? AND source_row = ? AND (image_assets IS NULL OR image_assets = '' OR image_assets = '[]')`, ) .bind(JSON.stringify(row.images), task.id, row.sourceRow), ); if (updates.length > 0) await db.batch(updates); } } export async function seedIfEmpty() { if (!isEnabled(getRuntimeEnv().SEED_DEMO_DATA)) return; const db = getRawDb(); const row = await db.prepare("SELECT COUNT(*) AS count FROM tasks").first<{ count: number; }>(); if ((row?.count ?? 0) > 0) return; const batch = [ db .prepare( `INSERT INTO partners (id, name, wecom_name, owner, claimed_total, completed_total) VALUES (?, ?, ?, ?, ?, ?)`, ) .bind("partner-linlin", "林林KOC社群", "林林|母婴社群", "小吴", 24, 18), db .prepare( `INSERT INTO partners (id, name, wecom_name, owner, claimed_total, completed_total) VALUES (?, ?, ?, ?, ?, ?)`, ) .bind("partner-muzi", "木子", "木子同学", "小陈", 8, 7), db .prepare( `INSERT INTO partners (id, name, wecom_name, owner, claimed_total, completed_total) VALUES (?, ?, ?, ?, ?, ?)`, ) .bind("partner-xiaolu", "小鹿内容组", "鹿鹿日常", "小吴", 15, 12), db .prepare( `INSERT INTO tasks (id, name, brand, quantity, claimed_quantity, due_at, status) VALUES (?, ?, ?, ?, ?, ?, ?)`, ) .bind( "task-summer", "夏日轻盈计划", "青柠实验室", 30, 9, "2026-08-05", "active", ), db .prepare( `INSERT INTO tasks (id, name, brand, quantity, claimed_quantity, due_at, status) VALUES (?, ?, ?, ?, ?, ?, ?)`, ) .bind( "task-living", "理想生活家", "木屿家居", 18, 6, "2026-08-02", "active", ), ]; const summerTitles = [ "夏天轻松管理身材的小习惯", "一周清爽饮食记录", "通勤女生的轻负担早餐", "周末宅家也要好好吃饭", "最近让我状态变好的三件事", "忙碌上班族的饮食搭配", "清爽一夏的简单生活方式", "我的夏日冰箱常备清单", "下班后的低成本幸福感", "在家也能完成的状态管理", "夏日办公室好物分享", "高温天的清爽仪式感", "一人食也可以很认真", "近期值得回购的小东西", "我的轻盈生活观察", "夏天拒绝疲惫感", "简单好坚持的日常习惯", "最近的通勤包里有什么", "不费力的夏日松弛感", "周末恢复能量的小计划", "高效生活的三个小改变", "一个人的清爽晚餐", "办公室里的续航秘诀", "生活需要一点轻盈感", "最近在坚持的健康习惯", "从早餐开始认真生活", "我的低负担下午茶", "夏日宅家幸福清单", "忙碌生活中的小确幸", "值得记录的轻盈一天", ]; const livingTitles = Array.from( { length: 18 }, (_, index) => `理想生活家的空间灵感 ${String(index + 1).padStart(2, "0")}`, ); summerTitles.forEach((title, index) => { batch.push( db .prepare( `INSERT INTO contents (id, task_id, title, body, status) VALUES (?, ?, ?, ?, ?)`, ) .bind( `content-s-${index + 1}`, "task-summer", title, "来自飞书内容表的完整笔记正文与素材说明。", index < 9 ? "allocated" : "available", ), ); }); livingTitles.forEach((title, index) => { batch.push( db .prepare( `INSERT INTO contents (id, task_id, title, body, status) VALUES (?, ?, ?, ?, ?)`, ) .bind( `content-l-${index + 1}`, "task-living", title, "来自飞书内容表的完整笔记正文与素材说明。", index < 6 ? "allocated" : "available", ), ); }); const accountSeeds = [ ["account-01", "xhs-8af3", "小满的轻生活", "上海", 12800, 4, 4360], ["account-02", "xhs-2bd8", "橘子汽水日记", "杭州", 8600, 3, 2980], ["account-03", "xhs-7ca1", "木木在成长", "广东", 21400, 2, 7620], ["account-04", "xhs-91ee", "一颗软糖", "江苏", 5200, 2, 1850], ]; accountSeeds.forEach(([id, uid, nickname, ip, followers, posts, avg]) => { batch.push( db .prepare( `INSERT INTO accounts (id, platform, platform_uid, nickname, profile_url, ip_location, followers, post_count, avg_views) VALUES (?, '小红书', ?, ?, ?, ?, ?, ?, ?)`, ) .bind( id, uid, nickname, `https://www.xiaohongshu.com/user/profile/${uid}`, ip, followers, posts, avg, ), ); }); const distributionSeeds = [ [ "dist-01", "content-s-1", "partner-linlin", "account-01", "https://www.xiaohongshu.com/explore/demo01", "complete", 118, 24, 43, 168, 31, 61, 232, 38, 86, 18420, 9430, "recognized", ], [ "dist-02", "content-s-2", "partner-muzi", "account-02", "https://www.xiaohongshu.com/explore/demo02", "collecting", 76, 12, 28, 103, 17, 39, null, null, null, null, null, "none", ], [ "dist-03", "content-s-3", "partner-linlin", "account-03", "https://www.xiaohongshu.com/explore/demo03", "published", null, null, null, null, null, null, null, null, null, null, null, "none", ], [ "dist-04", "content-l-1", "partner-xiaolu", "account-04", "https://www.xiaohongshu.com/explore/demo04", "collecting", 42, 8, 16, 69, 12, 24, 91, 17, 35, null, null, "uploaded", ], ]; distributionSeeds.forEach((seed, index) => { const [ id, contentId, partnerId, accountId, publishUrl, status, d2Likes, d2Comments, d2Collects, d5Likes, d5Comments, d5Collects, d7Likes, d7Comments, d7Collects, exposure, views, ocrStatus, ] = seed; batch.push( db .prepare( `INSERT INTO distributions ( id, task_id, content_id, partner_id, account_id, publish_url, publish_time, status, claimed_at, ocr_status, exposure, views, d2_likes, d2_comments, d2_collects, d5_likes, d5_comments, d5_collects, d7_likes, d7_comments, d7_collects ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, ) .bind( id, index < 3 ? "task-summer" : "task-living", contentId, partnerId, accountId, publishUrl, `2026-07-${String(22 + index).padStart(2, "0")} 10:30:00`, status, `2026-07-${String(20 + index).padStart(2, "0")} 09:00:00`, ocrStatus, exposure, views, d2Likes, d2Comments, d2Collects, d5Likes, d5Comments, d5Collects, d7Likes, d7Comments, d7Collects, ), ); }); for (let index = 4; index < 9; index += 1) { batch.push( db .prepare( `INSERT INTO distributions (id, task_id, content_id, partner_id, status, claimed_at) VALUES (?, 'task-summer', ?, 'partner-linlin', 'claimed', ?)`, ) .bind( `dist-pending-${index}`, `content-s-${index + 1}`, `2026-07-${String(24 + (index % 2)).padStart(2, "0")} 11:00:00`, ), ); } await db.batch(batch); } export async function getDashboardData() { const db = getRawDb(); const [partnersResult, tasksResult, accountsResult, distributionsResult] = await Promise.all([ db.prepare("SELECT * FROM partners ORDER BY created_at DESC").all(), db.prepare("SELECT * FROM tasks ORDER BY created_at DESC").all(), db.prepare("SELECT * FROM accounts ORDER BY last_seen_at DESC").all(), db .prepare( `SELECT d.*, c.title AS content_title, p.name AS partner_name, a.nickname AS account_nickname, a.platform AS account_platform, t.name AS task_name, t.brand AS task_brand, t.task_type AS task_type, t.due_at AS due_at FROM distributions d JOIN contents c ON c.id = d.content_id JOIN partners p ON p.id = d.partner_id JOIN tasks t ON t.id = d.task_id LEFT JOIN accounts a ON a.id = d.account_id ORDER BY d.updated_at DESC, d.claimed_at DESC`, ) .all(), ]); return { partners: partnersResult.results as D1ResultRow[], tasks: tasksResult.results as D1ResultRow[], accounts: accountsResult.results as D1ResultRow[], distributions: distributionsResult.results as D1ResultRow[], portal_url: getRuntimeEnv().KOC_PORTAL_URL ?? "", }; } export function uid(prefix: string) { return `${prefix}-${crypto.randomUUID().slice(0, 8)}`; } export function hashText(value: string) { let hash = 2166136261; for (let index = 0; index < value.length; index += 1) { hash ^= value.charCodeAt(index); hash = Math.imul(hash, 16777619); } return Math.abs(hash >>> 0).toString(36); }