mirror of
https://github.com/tiennm99/store-scraper-bot.git
synced 2026-10-03 07:13:32 +00:00
feat: migrate to Vercel + Upstash with KEY_PREFIX namespacing
Phases 1-5 of consolidate-vercel-upstash plan. Replaces Cloudflare Workers + KV with Vercel serverless functions + Upstash Redis. Inlines app-store-scraper / google-play-scraper npm libs (drops the store-scraper.vercel.app HTTP roundtrip). KEY_PREFIX (default 'store-scraper-bot:') namespaces all Redis keys so the Upstash DB can be safely shared with other Vercel projects. - vercel.json + .vercelignore + Vercel-aware package.json scripts - api/webhook.js + api/cron.js Vercel functions (with shared src/app-builder.js); cron auth fails closed when CRON_SECRET unset - src/repository/upstash.js replaces kv.js; all 4 repos take a handle bundling client + prefix - scripts/migrate-atlas-to-upstash.js writes legacy Java Atlas state directly to Upstash with --dry-run + --include-cache flags - .env.example refreshed for the new env surface Phases 6 (Vercel deploy + webhook cutover) and 7 (Docker + wrangler cleanup) remain operator-driven post-deploy.
This commit is contained in:
1 parent
134bce0826
commit
c2dd35b75f
33 files changed
+4011
-488
No files matched your search
+12
-15
@@ -1,7 +1,8 @@
|
||||
import store from 'app-store-scraper';
|
||||
import { newAppleApp } from '../models/apple-app.js';
|
||||
|
||||
// Mirrors Java AppStoreScraper (api/apple/AppStoreScraper.java).
|
||||
const BASE_URL = 'https://store-scraper.vercel.app/apple';
|
||||
// Calls the `app-store-scraper` npm lib directly (no HTTP roundtrip).
|
||||
|
||||
export function buildAppleRequestByTrackId(id, country) {
|
||||
return { id, country, ratings: true };
|
||||
@@ -11,23 +12,19 @@ export function buildAppleRequestByBundleId(appId, country) {
|
||||
return { appId, country, ratings: true };
|
||||
}
|
||||
|
||||
export function createAppleScraper(config, store) {
|
||||
export function createAppleScraper(config, repository) {
|
||||
const { logger } = config;
|
||||
const repo = store.appleApp;
|
||||
|
||||
async function rawApp(req) {
|
||||
const res = await fetch(`${BASE_URL}/app`, {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify(req),
|
||||
});
|
||||
if (!res.ok) throw new Error(`apple HTTP status ${res.status}`);
|
||||
return await res.text();
|
||||
}
|
||||
const repo = repository.appleApp;
|
||||
|
||||
async function app(req) {
|
||||
const text = await rawApp(req);
|
||||
return JSON.parse(text);
|
||||
return store.app(req);
|
||||
}
|
||||
|
||||
// rawApp returns a JSON-text representation of the parsed object so the
|
||||
// /rawappleapp command and any other text consumers stay parity-compatible
|
||||
// with the previous HTTP-text response.
|
||||
async function rawApp(req) {
|
||||
return JSON.stringify(await app(req));
|
||||
}
|
||||
|
||||
async function cache(resp) {
|
||||
|
||||
+12
-15
@@ -1,29 +1,26 @@
|
||||
import gplay from 'google-play-scraper';
|
||||
import { newGoogleApp } from '../models/google-app.js';
|
||||
|
||||
// Mirrors Java GooglePlayScraper (api/google/GooglePlayScraper.java).
|
||||
const BASE_URL = 'https://store-scraper.vercel.app/google';
|
||||
// Calls the `google-play-scraper` npm lib directly (no HTTP roundtrip).
|
||||
|
||||
export function buildGoogleRequest(appId, country) {
|
||||
return { appId, country: country || 'vn' };
|
||||
}
|
||||
|
||||
export function createGoogleScraper(config, store) {
|
||||
export function createGoogleScraper(config, repository) {
|
||||
const { logger } = config;
|
||||
const repo = store.googleApp;
|
||||
|
||||
async function rawApp(req) {
|
||||
const res = await fetch(`${BASE_URL}/app`, {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify(req),
|
||||
});
|
||||
if (!res.ok) throw new Error(`google HTTP status ${res.status}`);
|
||||
return await res.text();
|
||||
}
|
||||
const repo = repository.googleApp;
|
||||
|
||||
async function app(req) {
|
||||
const text = await rawApp(req);
|
||||
return JSON.parse(text);
|
||||
return gplay.app(req);
|
||||
}
|
||||
|
||||
// rawApp returns a JSON-text representation of the parsed object so the
|
||||
// /rawgoogleapp command and any other text consumers stay parity-compatible
|
||||
// with the previous HTTP-text response.
|
||||
async function rawApp(req) {
|
||||
return JSON.stringify(await app(req));
|
||||
}
|
||||
|
||||
async function cache(resp, fallbackId) {
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
// Per-invocation context wiring. Both Vercel handlers (api/webhook.js,
|
||||
// api/cron.js) call this once per request. Cheap — the Upstash client is
|
||||
// HTTP-based, so no socket/connection setup happens until the first command.
|
||||
|
||||
import { loadConfig } from './config.js';
|
||||
import { createUpstashClient } from './repository/upstash.js';
|
||||
import { createStore } from './repository/store.js';
|
||||
import { createAppleScraper } from './api/apple-scraper.js';
|
||||
import { createGoogleScraper } from './api/google-scraper.js';
|
||||
import { createBot } from './bot/bot.js';
|
||||
|
||||
export function buildApp(env) {
|
||||
const config = loadConfig(env);
|
||||
const handle = createUpstashClient(env);
|
||||
const store = createStore(handle, config.appCacheSeconds);
|
||||
const appleScraper = createAppleScraper(config, store);
|
||||
const googleScraper = createGoogleScraper(config, store);
|
||||
const { sender, commands } = createBot(config, store, appleScraper, googleScraper);
|
||||
return { config, store, appleScraper, googleScraper, sender, commands };
|
||||
}
|
||||
+2
-2
@@ -9,8 +9,8 @@ function parseAdminIds(raw) {
|
||||
.filter((n) => Number.isFinite(n));
|
||||
}
|
||||
|
||||
// Builds config from a Workers `env` binding. Called once per fetch / scheduled
|
||||
// invocation; cheap.
|
||||
// Builds config from a plain env dictionary (Vercel's `process.env` or any
|
||||
// dict-like). Called once per webhook / cron invocation; cheap.
|
||||
export function loadConfig(env) {
|
||||
const required = [
|
||||
'TELEGRAM_BOT_TOKEN',
|
||||
|
||||
@@ -1,74 +0,0 @@
|
||||
import { loadConfig } from './config.js';
|
||||
import { createStore } from './repository/store.js';
|
||||
import { createAppleScraper } from './api/apple-scraper.js';
|
||||
import { createGoogleScraper } from './api/google-scraper.js';
|
||||
import { createBot } from './bot/bot.js';
|
||||
import { dispatch } from './bot/dispatch.js';
|
||||
import { runDailyCheck } from './scheduler/scheduler.js';
|
||||
|
||||
// Builds the per-invocation context. Cheap — KV binding is exposed by the
|
||||
// runtime; no connection setup needed.
|
||||
function build(env) {
|
||||
if (!env.STORE_KV) throw new Error('STORE_KV binding missing');
|
||||
const config = loadConfig(env);
|
||||
const store = createStore(env, config.appCacheSeconds);
|
||||
const appleScraper = createAppleScraper(config, store);
|
||||
const googleScraper = createGoogleScraper(config, store);
|
||||
const { sender, commands } = createBot(config, store, appleScraper, googleScraper);
|
||||
return { config, store, appleScraper, googleScraper, sender, commands };
|
||||
}
|
||||
|
||||
export default {
|
||||
// Telegram webhook entry. Validates the `secret_token` header, acks fast,
|
||||
// then dispatches in `ctx.waitUntil` so Telegram doesn't retry on slow downstream calls.
|
||||
async fetch(request, env, ctx) {
|
||||
if (request.method !== 'POST') {
|
||||
return new Response('Not found', { status: 404 });
|
||||
}
|
||||
|
||||
let app;
|
||||
try {
|
||||
app = build(env);
|
||||
} catch (err) {
|
||||
console.log(JSON.stringify({ level: 'error', msg: 'config error', err: err.message }));
|
||||
return new Response('Server misconfigured', { status: 500 });
|
||||
}
|
||||
|
||||
const secret = request.headers.get('X-Telegram-Bot-Api-Secret-Token');
|
||||
if (secret !== app.config.telegramWebhookSecret) {
|
||||
return new Response('Unauthorized', { status: 401 });
|
||||
}
|
||||
|
||||
let update;
|
||||
try {
|
||||
update = await request.json();
|
||||
} catch {
|
||||
return new Response('Bad request', { status: 400 });
|
||||
}
|
||||
if (!update?.message) return new Response('OK');
|
||||
|
||||
ctx.waitUntil(
|
||||
dispatch(update.message, {
|
||||
sender: app.sender,
|
||||
commands: app.commands,
|
||||
config: app.config,
|
||||
logger: app.config.logger,
|
||||
}),
|
||||
);
|
||||
return new Response('OK');
|
||||
},
|
||||
|
||||
// Daily cron handler. Schedule lives in wrangler.toml.
|
||||
async scheduled(event, env, ctx) {
|
||||
let app;
|
||||
try {
|
||||
app = build(env);
|
||||
} catch (err) {
|
||||
console.log(JSON.stringify({ level: 'error', msg: 'config error', err: err.message }));
|
||||
return;
|
||||
}
|
||||
ctx.waitUntil(
|
||||
runDailyCheck(app.config, app.store, app.sender, app.appleScraper, app.googleScraper),
|
||||
);
|
||||
},
|
||||
};
|
||||
@@ -1,4 +1,4 @@
|
||||
import { getJson, putJson } from './kv.js';
|
||||
import { getJson, putJson } from './upstash.js';
|
||||
import {
|
||||
ADMIN_ID,
|
||||
adminAddGroup,
|
||||
@@ -7,22 +7,23 @@ import {
|
||||
newAdmin,
|
||||
} from '../models/admin.js';
|
||||
|
||||
// KV-backed admin singleton — Java parity at the document level
|
||||
// (key 'admin' holds the same shape Mongo stored at _id="admin").
|
||||
export function createAdminRepository(env) {
|
||||
// Upstash-backed admin singleton — Java parity at the document level
|
||||
// (logical key 'admin' holds the same shape Mongo stored at _id="admin").
|
||||
// The physical Redis key carries the configured KEY_PREFIX (handled by adapter).
|
||||
export function createAdminRepository(handle) {
|
||||
async function init() {
|
||||
const existing = await getJson(env, ADMIN_ID);
|
||||
const existing = await getJson(handle, ADMIN_ID);
|
||||
if (existing) return;
|
||||
await save(newAdmin());
|
||||
}
|
||||
|
||||
async function getAdmin() {
|
||||
const doc = await getJson(env, ADMIN_ID);
|
||||
const doc = await getJson(handle, ADMIN_ID);
|
||||
return doc ?? newAdmin();
|
||||
}
|
||||
|
||||
async function save(admin) {
|
||||
await putJson(env, ADMIN_ID, admin);
|
||||
await putJson(handle, ADMIN_ID, admin);
|
||||
}
|
||||
|
||||
async function addGroup(groupId) {
|
||||
|
||||
@@ -1,19 +1,20 @@
|
||||
import { getJson, putJson } from './kv.js';
|
||||
import { getJson, putJson } from './upstash.js';
|
||||
|
||||
// KV-backed Apple app cache. Key shape: `apple:{appId}`.
|
||||
// KV's expirationTtl replaces Java/Mongo's manual `(now - millis) > cacheMillis`
|
||||
// check — expired keys are deleted, so a get() returning null is the cache miss.
|
||||
export function createAppleAppRepository(env, appCacheSeconds) {
|
||||
// Upstash-backed Apple app cache. Logical key shape: `apple:{appId}`.
|
||||
// Redis EX (via expirationTtl) replaces Java/Mongo's manual
|
||||
// `(now - millis) > cacheMillis` check — expired keys are deleted, so a
|
||||
// get() returning null is the cache miss.
|
||||
export function createAppleAppRepository(handle, appCacheSeconds) {
|
||||
function key(appId) {
|
||||
return `apple:${appId}`;
|
||||
}
|
||||
|
||||
async function get(appId) {
|
||||
return getJson(env, key(appId));
|
||||
return getJson(handle, key(appId));
|
||||
}
|
||||
|
||||
async function save(entry) {
|
||||
await putJson(env, key(entry._id), entry, { expirationTtl: appCacheSeconds });
|
||||
await putJson(handle, key(entry._id), entry, { expirationTtl: appCacheSeconds });
|
||||
}
|
||||
|
||||
async function getCached(appId) {
|
||||
|
||||
@@ -1,19 +1,20 @@
|
||||
import { getJson, putJson } from './kv.js';
|
||||
import { getJson, putJson } from './upstash.js';
|
||||
|
||||
// KV-backed Google app cache. Key shape: `google:{appId}`.
|
||||
// KV's expirationTtl replaces Java/Mongo's manual `(now - millis) > cacheMillis`
|
||||
// check — expired keys are deleted, so a get() returning null is the cache miss.
|
||||
export function createGoogleAppRepository(env, appCacheSeconds) {
|
||||
// Upstash-backed Google app cache. Logical key shape: `google:{appId}`.
|
||||
// Redis EX (via expirationTtl) replaces Java/Mongo's manual
|
||||
// `(now - millis) > cacheMillis` check — expired keys are deleted, so a
|
||||
// get() returning null is the cache miss.
|
||||
export function createGoogleAppRepository(handle, appCacheSeconds) {
|
||||
function key(appId) {
|
||||
return `google:${appId}`;
|
||||
}
|
||||
|
||||
async function get(appId) {
|
||||
return getJson(env, key(appId));
|
||||
return getJson(handle, key(appId));
|
||||
}
|
||||
|
||||
async function save(entry) {
|
||||
await putJson(env, key(entry._id), entry, { expirationTtl: appCacheSeconds });
|
||||
await putJson(handle, key(entry._id), entry, { expirationTtl: appCacheSeconds });
|
||||
}
|
||||
|
||||
async function getCached(appId) {
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { del, getJson, putJson } from './kv.js';
|
||||
import { del, getJson, putJson } from './upstash.js';
|
||||
import {
|
||||
groupAddAppleApp,
|
||||
groupAddGoogleApp,
|
||||
@@ -8,24 +8,25 @@ import {
|
||||
newGroup,
|
||||
} from '../models/group.js';
|
||||
|
||||
// KV-backed per-group state. Key shape: `group:{chatId}`.
|
||||
export function createGroupRepository(env) {
|
||||
// Upstash-backed per-group state. Logical key shape: `group:{chatId}`.
|
||||
// Physical key gets KEY_PREFIX prepended by the adapter.
|
||||
export function createGroupRepository(handle) {
|
||||
function key(groupId) {
|
||||
return `group:${groupIdToKey(groupId)}`;
|
||||
}
|
||||
|
||||
async function exists(groupId) {
|
||||
const doc = await getJson(env, key(groupId));
|
||||
const doc = await getJson(handle, key(groupId));
|
||||
return doc !== null;
|
||||
}
|
||||
|
||||
async function getGroup(groupId) {
|
||||
const doc = await getJson(env, key(groupId));
|
||||
const doc = await getJson(handle, key(groupId));
|
||||
return doc ?? newGroup(groupId);
|
||||
}
|
||||
|
||||
async function saveGroup(group) {
|
||||
await putJson(env, key(group._id), group);
|
||||
await putJson(handle, key(group._id), group);
|
||||
}
|
||||
|
||||
async function initGroup(groupId) {
|
||||
@@ -34,7 +35,7 @@ export function createGroupRepository(env) {
|
||||
}
|
||||
|
||||
async function deleteGroup(groupId) {
|
||||
await del(env, key(groupId));
|
||||
await del(handle, key(groupId));
|
||||
}
|
||||
|
||||
async function mutateAndSave(groupId, mutator) {
|
||||
|
||||
@@ -1,38 +0,0 @@
|
||||
// Thin wrapper around the Cloudflare KV binding `env.STORE_KV`.
|
||||
// All four logical collections live in one namespace, separated by key prefix:
|
||||
// admin singleton
|
||||
// group:{chatId} per-group state
|
||||
// apple:{appId} cached Apple response (with KV TTL)
|
||||
// google:{appId} cached Google response (with KV TTL)
|
||||
|
||||
// KV's minimum expirationTtl is 60s. Java/Mongo had no such floor; clamp here
|
||||
// so a low APP_CACHE_SECONDS override doesn't make put() reject.
|
||||
const KV_MIN_TTL_SECONDS = 60;
|
||||
|
||||
export class KvUnavailable extends Error {
|
||||
constructor() {
|
||||
super('STORE_KV binding is missing — check wrangler.toml [[kv_namespaces]]');
|
||||
this.name = 'KvUnavailable';
|
||||
}
|
||||
}
|
||||
|
||||
function binding(env) {
|
||||
if (!env || !env.STORE_KV) throw new KvUnavailable();
|
||||
return env.STORE_KV;
|
||||
}
|
||||
|
||||
export async function getJson(env, key) {
|
||||
return binding(env).get(key, 'json');
|
||||
}
|
||||
|
||||
export async function putJson(env, key, value, opts = {}) {
|
||||
const putOpts = { ...opts };
|
||||
if (putOpts.expirationTtl != null) {
|
||||
putOpts.expirationTtl = Math.max(KV_MIN_TTL_SECONDS, putOpts.expirationTtl);
|
||||
}
|
||||
await binding(env).put(key, JSON.stringify(value), putOpts);
|
||||
}
|
||||
|
||||
export async function del(env, key) {
|
||||
await binding(env).delete(key);
|
||||
}
|
||||
@@ -3,13 +3,14 @@ import { createGroupRepository } from './group-repository.js';
|
||||
import { createAppleAppRepository } from './apple-app-repository.js';
|
||||
import { createGoogleAppRepository } from './google-app-repository.js';
|
||||
|
||||
// Single binding point for all repositories. Threads `env` once so command
|
||||
// handlers don't need to know about the Worker `env` argument or the KV binding.
|
||||
export function createStore(env, appCacheSeconds) {
|
||||
// Single binding point for all repositories. Threads the Upstash handle
|
||||
// (client + key prefix) once so command handlers don't need to know about
|
||||
// process.env or the Redis client construction.
|
||||
export function createStore(handle, appCacheSeconds) {
|
||||
return {
|
||||
admin: createAdminRepository(env),
|
||||
group: createGroupRepository(env),
|
||||
appleApp: createAppleAppRepository(env, appCacheSeconds),
|
||||
googleApp: createGoogleAppRepository(env, appCacheSeconds),
|
||||
admin: createAdminRepository(handle),
|
||||
group: createGroupRepository(handle),
|
||||
appleApp: createAppleAppRepository(handle, appCacheSeconds),
|
||||
googleApp: createGoogleAppRepository(handle, appCacheSeconds),
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
// Upstash Redis adapter — replaces the prior Cloudflare KV wrapper.
|
||||
//
|
||||
// Logical key namespace (unchanged from KV layer):
|
||||
// admin singleton
|
||||
// group:{chatId} per-group state
|
||||
// apple:{appId} cached Apple response (with TTL)
|
||||
// google:{appId} cached Google response (with TTL)
|
||||
//
|
||||
// Multi-tenancy: every physical Redis key carries a configurable prefix
|
||||
// (env.KEY_PREFIX, default 'store-scraper-bot:') so this bot can safely share
|
||||
// an Upstash database with other Vercel projects without collision. Repository
|
||||
// callers pass logical keys; the adapter applies the prefix transparently.
|
||||
//
|
||||
// 60s minimum TTL clamp is preserved from the KV days for parity safety,
|
||||
// even though Redis would accept lower values.
|
||||
|
||||
import { Redis } from '@upstash/redis';
|
||||
|
||||
const MIN_TTL_SECONDS = 60;
|
||||
const DEFAULT_KEY_PREFIX = 'store-scraper-bot:';
|
||||
|
||||
export class UpstashUnavailable extends Error {
|
||||
constructor(missing) {
|
||||
super(`Upstash env var missing: ${missing}`);
|
||||
this.name = 'UpstashUnavailable';
|
||||
}
|
||||
}
|
||||
|
||||
// Build a handle bundling the Redis client and the key prefix together.
|
||||
// The handle is what callers pass into getJson/putJson/del/scan — it stays
|
||||
// opaque so repositories never need to know about prefixing themselves.
|
||||
export function createUpstashClient(env) {
|
||||
if (!env?.UPSTASH_REDIS_REST_URL) throw new UpstashUnavailable('UPSTASH_REDIS_REST_URL');
|
||||
if (!env?.UPSTASH_REDIS_REST_TOKEN) throw new UpstashUnavailable('UPSTASH_REDIS_REST_TOKEN');
|
||||
const client = new Redis({
|
||||
url: env.UPSTASH_REDIS_REST_URL,
|
||||
token: env.UPSTASH_REDIS_REST_TOKEN,
|
||||
});
|
||||
const prefix = env.KEY_PREFIX ?? DEFAULT_KEY_PREFIX;
|
||||
return { client, prefix };
|
||||
}
|
||||
|
||||
function physicalKey(handle, key) {
|
||||
return `${handle.prefix}${key}`;
|
||||
}
|
||||
|
||||
// Upstash auto-deserializes values that look like JSON. We always store via
|
||||
// JSON.stringify, so reads can return the parsed object directly. Returns null
|
||||
// on missing key, matching the prior KV semantics.
|
||||
export async function getJson(handle, key) {
|
||||
const value = await handle.client.get(physicalKey(handle, key));
|
||||
if (value == null) return null;
|
||||
// Some SDK versions return strings, others return parsed objects depending
|
||||
// on content. Normalize: if string, parse; if object, pass through.
|
||||
return typeof value === 'string' ? JSON.parse(value) : value;
|
||||
}
|
||||
|
||||
export async function putJson(handle, key, value, opts = {}) {
|
||||
const ex =
|
||||
opts.expirationTtl != null ? Math.max(MIN_TTL_SECONDS, opts.expirationTtl) : null;
|
||||
const setOpts = ex != null ? { ex } : undefined;
|
||||
await handle.client.set(physicalKey(handle, key), JSON.stringify(value), setOpts);
|
||||
}
|
||||
|
||||
export async function del(handle, key) {
|
||||
await handle.client.del(physicalKey(handle, key));
|
||||
}
|
||||
|
||||
// Suffix-based scan. Caller passes a logical match like 'group:*'; adapter
|
||||
// prepends the key prefix so only this bot's keys are returned.
|
||||
// Returns the list of *logical* keys (prefix stripped) so callers stay
|
||||
// prefix-unaware.
|
||||
export async function scan(handle, matchSuffix) {
|
||||
const match = `${handle.prefix}${matchSuffix}`;
|
||||
const out = [];
|
||||
let cursor = '0';
|
||||
do {
|
||||
const [next, batch] = await handle.client.scan(cursor, { match, count: 100 });
|
||||
cursor = next;
|
||||
for (const physical of batch) {
|
||||
out.push(physical.startsWith(handle.prefix) ? physical.slice(handle.prefix.length) : physical);
|
||||
}
|
||||
} while (cursor !== '0');
|
||||
return out;
|
||||
}
|
||||
@@ -1,8 +1,8 @@
|
||||
import { buildTable, formatNumber, truncateString } from '../util/table.js';
|
||||
import { daysBetween, formatDateInTz, formatDateTimeInTz, weekdayInTz } from '../util/time.js';
|
||||
|
||||
// One-shot daily check, invoked from the Worker `scheduled` handler. The cron
|
||||
// schedule lives in wrangler.toml ("0 0 * * *" UTC = 7am Asia/Ho_Chi_Minh).
|
||||
// One-shot daily check, invoked from api/cron.js. The cron schedule lives in
|
||||
// vercel.json ("0 0 * * *" UTC = 7am Asia/Ho_Chi_Minh).
|
||||
export async function runDailyCheck(config, store, sender, appleScraper, googleScraper) {
|
||||
const logger = config.logger;
|
||||
const now = new Date();
|
||||
|
||||
Reference in new issue
Block a user