How to create a pipeline where i take each Instagram post and extract its insights at 24 hour and 15 days and then at day 30.
ls -la /home/daytona/workspace/ 2>/dev/null || echo "workspace empty or missing"
mkdir -p /home/daytona/workspace/instagram-insights-pipeline/src
~/instagram-insights-pipeline/schema.sql
-- ─────────────────────────────────────────────────────────────
-- Instagram Insights Pipeline — Job Queue Schema
-- Works with SQLite (default) or PostgreSQL (change AUTOINCREMENT
-- to SERIAL for Postgres, and INTEGER PRIMARY KEY → BIGSERIAL).
-- ─────────────────────────────────────────────────────────────
-- Tracks every Instagram post we need to monitor
CREATE TABLE IF NOT EXISTS posts (
id TEXT PRIMARY KEY, -- Instagram media_id (string)
account_id TEXT NOT NULL, -- IG Business Account ID
permalink TEXT,
media_type TEXT, -- IMAGE | VIDEO | CAROUSEL_ALBUM
published_at DATETIME NOT NULL, -- Exact publish timestamp (UTC)
created_at DATETIME DEFAULT (datetime('now'))
);
-- One row per (post × window). Three rows per post: 24h, 15d, 30d.
CREATE TABLE IF NOT EXISTS insight_jobs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
post_id TEXT NOT NULL REFERENCES posts(id),
window_label TEXT NOT NULL, -- '24h' | '15d' | '30d'
run_at DATETIME NOT NULL, -- When to execute the fetch
status TEXT NOT NULL DEFAULT 'pending', -- pending | running | done | failed
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT,
created_at DATETIME DEFAULT (datetime('now')),
updated_at DATETIME DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_jobs_run_at_status ON insight_jobs (run_at, status);
-- Stores the raw insights payload for each completed job
CREATE TABLE IF NOT EXISTS insights (
id INTEGER PRIMARY KEY AUTOINCREMENT,
job_id INTEGER NOT NULL REFERENCES insight_jobs(id),
post_id TEXT NOT NULL,
window_label TEXT NOT NULL,
fetched_at DATETIME DEFAULT (datetime('now')),
-- Core reach / engagement metrics
impressions INTEGER,
reach INTEGER,
likes INTEGER,
comments INTEGER,
shares INTEGER,
saves INTEGER,
video_views INTEGER,
-- Full raw response for future-proofing
raw_json TEXT
);
~/instagram-insights-pipeline/src/db.js
/**
* db.js — Lightweight SQLite adapter using the `better-sqlite3` package.
*
* If you already have a data store (Postgres, MySQL, etc.), replace this
* file with your own adapter that exposes the same exported functions.
* The rest of the codebase only imports from this file, so swapping is easy.
*/
'use strict';
const Database = require('better-sqlite3');
const path = require('path');
const fs = require('fs');
const DB_PATH = process.env.DB_PATH || path.join(__dirname, '..', 'data', 'pipeline.db');
// Ensure the data directory exists
fs.mkdirSync(path.dirname(DB_PATH), { recursive: true });
const db = new Database(DB_PATH);
// Enable WAL mode for better concurrent read performance
db.pragma('journal_mode = WAL');
db.pragma('foreign_keys = ON');
// ── Bootstrap schema on first run ──────────────────────────────────────────
const schemaPath = path.join(__dirname, '..', 'schema.sql');
const schema = fs.readFileSync(schemaPath, 'utf8');
db.exec(schema);
// ── Posts ───────────────────────────────────────────────────────────────────
/**
* Insert a new post. Silently ignores duplicates (idempotent).
* @param {{ id, account_id, permalink, media_type, published_at }} post
*/
function upsertPost(post) {
db.prepare(`
INSERT INTO posts (id, account_id, permalink, media_type, published_at)
VALUES (@id, @account_id, @permalink, @media_type, @published_at)
ON CONFLICT(id) DO NOTHING
`).run(post);
}
function getPost(postId) {
return db.prepare('SELECT * FROM posts WHERE id = ?').get(postId);
}
// ── Jobs ────────────────────────────────────────────────────────────────────
/**
* Insert the 3 scheduled jobs for a post (idempotent via UNIQUE constraint).
* @param {{ post_id, window_label, run_at }[]} jobs
*/
function insertJobs(jobs) {
const stmt = db.prepare(`
INSERT INTO insight_jobs (post_id, window_label, run_at)
VALUES (@post_id, @window_label, @run_at)
ON CONFLICT DO NOTHING
`);
const insertMany = db.transaction((rows) => rows.forEach((r) => stmt.run(r)));
insertMany(jobs);
}
/**
* Fetch all pending jobs whose run_at has passed.
* Marks them as 'running' atomically before returning.
*/
function claimDueJobs(limit = 20) {
const now = new Date().toISOString();
return db.transaction(() => {
const jobs = db.prepare(`
SELECT * FROM insight_jobs
WHERE status = 'pending' AND run_at <= ?
ORDER BY run_at ASC
LIMIT ?
`).all(now, limit);
if (jobs.length === 0) return [];
const ids = jobs.map((j) => j.id);
db.prepare(`
UPDATE insight_jobs
SET status = 'running', updated_at = datetime('now')
WHERE id IN (${ids.map(() => '?').join(',')})
`).run(...ids);
return jobs;
})();
}
function markJobDone(jobId) {
db.prepare(`
UPDATE insight_jobs
SET status = 'done', updated_at = datetime('now')
WHERE id = ?
`).run(jobId);
}
function markJobFailed(jobId, errorMessage) {
db.prepare(`
UPDATE insight_jobs
SET status = 'failed',
last_error = ?,
attempts = attempts + 1,
updated_at = datetime('now')
WHERE id = ?
`).run(errorMessage, jobId);
}
/**
* Reset a failed job to pending so it will be retried.
* Called automatically by runner.js after exponential backoff.
*/
function requeueJob(jobId, runAt) {
db.prepare(`
UPDATE insight_jobs
SET status = 'pending', run_at = ?, updated_at = datetime('now')
WHERE id = ?
`).run(runAt, jobId);
}
// ── Insights ─────────────────────────────────────────────────────────────────
/**
* Persist the fetched insights for a completed job.
* @param {object} params
*/
function saveInsights({ job_id, post_id, window_label, metrics, raw }) {
db.prepare(`
INSERT INTO insights
(job_id, post_id, window_label, impressions, reach, likes,
comments, shares, saves, video_views, raw_json)
VALUES
(@job_id, @post_id, @window_label, @impressions, @reach, @likes,
@comments, @shares, @saves, @video_views, @raw_json)
`).run({
job_id,
post_id,
window_label,
impressions: metrics.impressions ?? null,
reach: metrics.reach ?? null,
likes: metrics.likes ?? null,
comments: metrics.comments ?? null,
shares: metrics.shares ?? null,
saves: metrics.saves ?? null,
video_views: metrics.video_views ?? null,
raw_json: JSON.stringify(raw),
});
}
module.exports = {
upsertPost,
getPost,
insertJobs,
claimDueJobs,
markJobDone,
markJobFailed,
requeueJob,
saveInsights,
};
~/instagram-insights-pipeline/src/instagram.js
/**
* instagram.js — Instagram Graph API client
*
* Docs: https://developers.facebook.com/docs/instagram-api/guides/insights
*
* Required env vars:
* IG_ACCESS_TOKEN — Long-lived Page / System User access token
* IG_API_VERSION — e.g. "v19.0" (defaults to v19.0)
*/
'use strict';
const https = require('https');
const BASE = 'https://graph.facebook.com';
const API_VERSION = process.env.IG_API_VERSION || 'v19.0';
const ACCESS_TOKEN = process.env.IG_ACCESS_TOKEN;
if (!ACCESS_TOKEN) {
console.warn('[instagram] WARNING: IG_ACCESS_TOKEN is not set. API calls will fail.');
}
// ── Generic fetch helper (no external dependencies, uses built-in https) ────
function get(url) {
return new Promise((resolve, reject) => {
https.get(url, (res) => {
let body = '';
res.on('data', (chunk) => (body += chunk));
res.on('end', () => {
try {
const json = JSON.parse(body);
if (json.error) {
reject(new Error(`[IG API] ${json.error.message} (code ${json.error.code})`));
} else {
resolve(json);
}
} catch (e) {
reject(new Error(`[IG API] Invalid JSON response: ${body.slice(0, 200)}`));
}
});
}).on('error', reject);
});
}
// ── Metrics per media type ───────────────────────────────────────────────────
//
// Instagram returns different metric names depending on media type.
// VIDEO and REELS additionally support video_views.
// CAROUSEL_ALBUM supports the same as IMAGE.
//
// Reference: https://developers.facebook.com/docs/instagram-api/reference/ig-media/insights
const METRICS_BY_TYPE = {
IMAGE: 'impressions,reach,likes,comments,shares,saved',
CAROUSEL_ALBUM: 'impressions,reach,likes,comments,shares,saved',
VIDEO: 'impressions,reach,likes,comments,shares,saved,video_views',
REEL: 'reach,likes,comments,shares,saved,plays',
STORY: 'impressions,reach,replies,exits,taps_forward,taps_back',
};
function metricsForType(mediaType) {
return METRICS_BY_TYPE[mediaType?.toUpperCase()] || METRICS_BY_TYPE.IMAGE;
}
// ── Public API ───────────────────────────────────────────────────────────────
/**
* Fetch insights for a single media object.
*
* @param {string} mediaId — Instagram media ID
* @param {string} mediaType — IMAGE | VIDEO | CAROUSEL_ALBUM | REEL | STORY
* @returns {Promise<{ metrics: object, raw: object }>}
*/
async function fetchInsights(mediaId, mediaType = 'IMAGE') {
const metric = metricsForType(mediaType);
const url = `${BASE}/${API_VERSION}/${mediaId}/insights`
+ `?metric=${metric}&access_token=${ACCESS_TOKEN}`;
const raw = await get(url);
// Flatten the data array into a { metric_name: value } object
const metrics = {};
for (const item of (raw.data || [])) {
// Each item: { name, period, values: [{ value }], ... }
const val = item.values?.[0]?.value ?? item.value ?? null;
const key = item.name === 'saved' ? 'saves' : item.name;
metrics[key] = typeof val === 'object' ? JSON.stringify(val) : val;
}
return { metrics, raw };
}
/**
* Fetch basic fields for a media object (used to get media_type and timestamp
* when registering posts from a webhook that only delivers a media ID).
*
* @param {string} mediaId
* @returns {Promise<{ id, media_type, timestamp, permalink }>}
*/
async function fetchMediaFields(mediaId) {
const url = `${BASE}/${API_VERSION}/${mediaId}`
+ `?fields=id,media_type,timestamp,permalink&access_token=${ACCESS_TOKEN}`;
return get(url);
}
module.exports = { fetchInsights, fetchMediaFields };
~/instagram-insights-pipeline/src/scheduler.js
/**
* scheduler.js — Registers the three insight-fetch jobs for a new post.
*
* Call registerPost() as soon as a new post is detected (from webhook or poll).
* It is fully idempotent: safe to call multiple times for the same post_id.
*
* Windows:
* 24h → published_at + 24 hours
* 15d → published_at + 15 days
* 30d → published_at + 30 days
*/
'use strict';
const db = require('./db');
// ── Window definitions ────────────────────────────────────────────────────────
const WINDOWS = [
{ label: '24h', offsetMs: 24 * 60 * 60 * 1000 },
{ label: '15d', offsetMs: 15 * 24 * 60 * 60 * 1000 },
{ label: '30d', offsetMs: 30 * 24 * 60 * 60 * 1000 },
];
/**
* Register a post and schedule its three insight-fetch jobs.
*
* @param {object} params
* @param {string} params.id — Instagram media_id
* @param {string} params.account_id — IG Business Account ID
* @param {string} [params.permalink]
* @param {string} [params.media_type] — IMAGE | VIDEO | CAROUSEL_ALBUM | REEL | STORY
* @param {Date|string} params.published_at — Publish timestamp (UTC)
*/
function registerPost({ id, account_id, permalink, media_type, published_at }) {
const publishedDate = new Date(published_at);
if (isNaN(publishedDate.getTime())) {
throw new Error(`[scheduler] Invalid published_at value: ${published_at}`);
}
// 1. Persist the post (no-op if already exists)
db.upsertPost({
id,
account_id,
permalink: permalink || null,
media_type: media_type || 'IMAGE',
published_at: publishedDate.toISOString(),
});
// 2. Create the three scheduled jobs
const jobs = WINDOWS.map(({ label, offsetMs }) => ({
post_id: id,
window_label: label,
run_at: new Date(publishedDate.getTime() + offsetMs).toISOString(),
}));
db.insertJobs(jobs);
console.log(
`[scheduler] Registered post ${id} | jobs scheduled at: ` +
jobs.map((j) => `${j.window_label} → ${j.run_at}`).join(', ')
);
return jobs;
}
module.exports = { registerPost };
~/instagram-insights-pipeline/src/webhook.js
/**
* webhook.js — Express router for receiving Instagram webhook events.
*
* Mount this in index.js:
* app.use('/webhook', require('./webhook'));
*
* Required env vars:
* WEBHOOK_VERIFY_TOKEN — Your chosen verification token (set in Meta App Dashboard)
* IG_ACCOUNT_ID — Default IG Business Account ID (used when not in payload)
*
* Meta Webhook docs:
* https://developers.facebook.com/docs/graph-api/webhooks/getting-started
* https://developers.facebook.com/docs/instagram-api/guides/webhooks
*/
'use strict';
const { Router } = require('express');
const { fetchMediaFields } = require('./instagram');
const { registerPost } = require('./scheduler');
const router = Router();
const VERIFY_TOKEN = process.env.WEBHOOK_VERIFY_TOKEN;
const IG_ACCOUNT_ID = process.env.IG_ACCOUNT_ID;
// ── Step 1: Webhook verification (one-time, during Meta App setup) ────────────
router.get('/', (req, res) => {
const mode = req.query['hub.mode'];
const token = req.query['hub.verify_token'];
const challenge = req.query['hub.challenge'];
if (mode === 'subscribe' && token === VERIFY_TOKEN) {
console.log('[webhook] Verification successful');
return res.status(200).send(challenge);
}
console.warn('[webhook] Verification failed — token mismatch');
res.sendStatus(403);
});
// ── Step 2: Receive real-time events ─────────────────────────────────────────
router.post('/', async (req, res) => {
// Always respond quickly — Meta retries on non-2xx or slow responses
res.sendStatus(200);
const body = req.body;
if (body.object !== 'instagram') {
return; // Not an IG event
}
for (const entry of (body.entry || [])) {
for (const change of (entry.changes || [])) {
// We only care about new media publications
if (change.field !== 'media') continue;
const value = change.value || {};
const mediaId = value.media_id || value.id;
if (!mediaId) {
console.warn('[webhook] Received media event without media_id:', JSON.stringify(value));
continue;
}
// Webhooks often deliver minimal data — fetch full fields from the API
try {
const media = await fetchMediaFields(mediaId);
await registerPost({
id: media.id,
account_id: entry.id || IG_ACCOUNT_ID,
permalink: media.permalink,
media_type: media.media_type,
published_at: media.timestamp, // ISO 8601 string from IG API
});
} catch (err) {
console.error(`[webhook] Failed to register post ${mediaId}:`, err.message);
}
}
}
});
module.exports = router;
~/instagram-insights-pipeline/src/runner.js
/**
* runner.js — Cron job that processes due insight-fetch jobs.
*
* Run this file on a schedule (e.g. every 5 minutes via node-cron or an
* external cron daemon). It:
* 1. Claims all pending jobs whose run_at has passed
* 2. Fetches insights from the Instagram Graph API
* 3. Persists results and marks jobs done
* 4. On failure, reschedules with exponential backoff (max 3 attempts)
*
* Usage:
* node src/runner.js ← run once (suitable for external cron)
* RUN_INLINE=true node index.js ← run on embedded schedule (via node-cron)
*/
'use strict';
const db = require('./db');
const { fetchInsights } = require('./instagram');
// Max retry attempts before giving up
const MAX_ATTEMPTS = 3;
// Backoff delays (in ms) indexed by attempt number (0-based)
const BACKOFF_MS = [
5 * 60 * 1000, // attempt 1 failed → retry in 5 min
30 * 60 * 1000, // attempt 2 failed → retry in 30 min
2 * 60 * 60 * 1000, // attempt 3 failed → give up
];
/**
* Process all due jobs. Called by the cron schedule or directly.
*/
async function runDueJobs() {
const jobs = db.claimDueJobs(20);
if (jobs.length === 0) {
console.log('[runner] No jobs due.');
return;
}
console.log(`[runner] Processing ${jobs.length} due job(s)...`);
for (const job of jobs) {
await processJob(job);
}
}
async function processJob(job) {
const post = db.getPost(job.post_id);
if (!post) {
console.error(`[runner] Post ${job.post_id} not found in DB for job ${job.id}. Marking failed.`);
db.markJobFailed(job.id, 'Post record not found');
return;
}
console.log(`[runner] Fetching insights — post=${job.post_id} window=${job.window_label}`);
try {
const { metrics, raw } = await fetchInsights(post.id, post.media_type);
db.saveInsights({
job_id: job.id,
post_id: post.id,
window_label: job.window_label,
metrics,
raw,
});
db.markJobDone(job.id);
console.log(
`[runner] ✓ Saved insights for post=${post.id} window=${job.window_label}`,
`reach=${metrics.reach ?? 'n/a'} impressions=${metrics.impressions ?? 'n/a'}`
);
} catch (err) {
console.error(`[runner] ✗ Error for job ${job.id}:`, err.message);
const attempt = (job.attempts || 0) + 1;
if (attempt >= MAX_ATTEMPTS) {
db.markJobFailed(job.id, err.message);
console.error(`[runner] Job ${job.id} exceeded max attempts. Marked failed.`);
} else {
const delayMs = BACKOFF_MS[attempt - 1] || BACKOFF_MS[BACKOFF_MS.length - 1];
const retryAt = new Date(Date.now() + delayMs).toISOString();
db.requeueJob(job.id, retryAt);
console.warn(`[runner] Job ${job.id} will retry at ${retryAt} (attempt ${attempt}/${MAX_ATTEMPTS})`);
}
}
}
// ── Direct execution (node src/runner.js) ────────────────────────────────────
if (require.main === module) {
runDueJobs()
.then(() => {
console.log('[runner] Done.');
process.exit(0);
})
.catch((err) => {
console.error('[runner] Fatal error:', err);
process.exit(1);
});
}
module.exports = { runDueJobs };
~/instagram-insights-pipeline/index.js
/**
* index.js — Entry point
*
* Starts the Express webhook server and optionally runs the job runner
* on an embedded cron schedule.
*
* Environment variables (create a .env file or set them in your environment):
*
* IG_ACCESS_TOKEN — Long-lived Instagram Graph API access token
* IG_API_VERSION — e.g. "v19.0" (default: v19.0)
* IG_ACCOUNT_ID — Your IG Business Account ID
* WEBHOOK_VERIFY_TOKEN — Secret token you set in Meta App Dashboard
* PORT — HTTP server port (default: 3000)
* DB_PATH — SQLite DB path (default: ./data/pipeline.db)
* CRON_SCHEDULE — node-cron pattern (default: "*/5 * * * *" = every 5 min)
* RUN_RUNNER_INLINE — Set to "true" to run the job runner inside this process
* (alternative: run `node src/runner.js` via external cron)
*/
'use strict';
require('dotenv').config();
const express = require('express');
const cron = require('node-cron');
const webhookRouter = require('./src/webhook');
const { runDueJobs } = require('./src/runner');
const PORT = process.env.PORT || 3000;
const CRON_SCHEDULE = process.env.CRON_SCHEDULE || '*/5 * * * *';
const RUN_INLINE = process.env.RUN_RUNNER_INLINE === 'true';
// ── Express app ───────────────────────────────────────────────────────────────
const app = express();
app.use(express.json());
// Health check
app.get('/health', (_req, res) => res.json({ status: 'ok', ts: new Date().toISOString() }));
// Instagram webhook endpoint
app.use('/webhook/instagram', webhookRouter);
// ── Optional: manual trigger for testing ─────────────────────────────────────
// POST /admin/run-jobs (protect this in production!)
app.post('/admin/run-jobs', async (_req, res) => {
try {
await runDueJobs();
res.json({ ok: true });
} catch (err) {
res.status(500).json({ error: err.message });
}
});
// ── Start server ──────────────────────────────────────────────────────────────
app.listen(PORT, () => {
console.log(`[app] Server listening on port ${PORT}`);
console.log(`[app] Webhook endpoint: POST /webhook/instagram`);
// ── Inline cron runner ────────────────────────────────────────────────────
if (RUN_INLINE) {
if (!cron.validate(CRON_SCHEDULE)) {
console.error(`[app] Invalid CRON_SCHEDULE: "${CRON_SCHEDULE}"`);
process.exit(1);
}
console.log(`[app] Job runner cron active: "${CRON_SCHEDULE}"`);
cron.schedule(CRON_SCHEDULE, async () => {
try {
await runDueJobs();
} catch (err) {
console.error('[cron] Error running jobs:', err.message);
}
});
} else {
console.log('[app] Inline cron disabled. Run `node src/runner.js` via external cron.');
}
});
~/instagram-insights-pipeline/package.json
{
"name": "instagram-insights-pipeline",
"version": "1.0.0",
"description": "Automatically fetches Instagram post insights at 24h, 15d, and 30d after publication.",
"main": "index.js",
"scripts": {
"start": "node index.js",
"start:cron": "RUN_RUNNER_INLINE=true node index.js",
"runner": "node src/runner.js",
"dev": "RUN_RUNNER_INLINE=true node --watch index.js"
},
"dependencies": {
"better-sqlite3": "^9.4.3",
"dotenv": "^16.4.5",
"express": "^4.19.2",
"node-cron": "^3.0.3"
},
"engines": {
"node": ">=18.0.0"
}
}
~/instagram-insights-pipeline/README.md
# Instagram Insights Pipeline
Automatically fetches Instagram post insights at **24 hours**, **15 days**, and **30 days** after publication.
## How it works
```
New Post Published on Instagram
│
▼
[Webhook Receiver] (/webhook/instagram)
Receives media_id → fetches media fields → calls registerPost()
│
▼
[Scheduler]
Inserts 3 rows into insight_jobs table:
├─ run_at = published_at + 24h
├─ run_at = published_at + 15d
└─ run_at = published_at + 30d
│
▼ (every 5 minutes)
[Job Runner] (src/runner.js or cron)
SELECT jobs WHERE run_at <= NOW() AND status='pending'
For each:
1. Call Instagram Insights API
2. Save results to `insights` table
3. Mark job as 'done'
(On failure: retry with exponential backoff, max 3 attempts)
```
## Quick Start
### 1. Install dependencies
```bash
npm install
```
### 2. Configure environment
```bash
cp .env.example .env
# Fill in your values
```
### 3. Start the server (with embedded cron)
```bash
npm run start:cron
```
Or start the server and run the job processor via your own external cron:
```bash
npm start # webhook server
# Add to crontab:
# */5 * * * * /usr/bin/node /path/to/project/src/runner.js
```
## Environment Variables
| Variable | Required | Description |
|---|---|---|
| `IG_ACCESS_TOKEN` | ✅ | Long-lived Instagram Graph API access token |
| `IG_API_VERSION` | | API version, e.g. `v19.0` (default: `v19.0`) |
| `IG_ACCOUNT_ID` | ✅ | Your IG Business Account ID |
| `WEBHOOK_VERIFY_TOKEN` | ✅ | Secret token set in Meta App Dashboard |
| `PORT` | | HTTP port (default: `3000`) |
| `DB_PATH` | | SQLite file path (default: `./data/pipeline.db`) |
| `CRON_SCHEDULE` | | node-cron pattern (default: `*/5 * * * *`) |
| `RUN_RUNNER_INLINE` | | Set `true` to run cron inside the main process |
## Meta App Setup (Webhooks)
1. Go to [Meta Developer Dashboard](https://developers.facebook.com)
2. Your App → Webhooks → Instagram → Subscribe to `media` field
3. Callback URL: `https://your-domain.com/webhook/instagram`
4. Verify Token: same value as `WEBHOOK_VERIFY_TOKEN`
## Manually register a post (for backfilling)
```js
const { registerPost } = require('./src/scheduler');
registerPost({
id: '17854360229135492',
account_id: '1234567890',
media_type: 'IMAGE',
published_at: '2025-04-25T10:00:00Z',
permalink: 'https://www.instagram.com/p/abc123/',
});
```
## Database schema
- **posts** — one row per Instagram post
- **insight_jobs** — job queue (3 rows per post: 24h / 15d / 30d)
- **insights** — stored results with all metrics + raw JSON
## Swap the data store
Replace `src/db.js` with your own adapter. The public interface is:
```
upsertPost(post)
getPost(postId)
insertJobs(jobs[])
claimDueJobs(limit)
markJobDone(jobId)
markJobFailed(jobId, error)
requeueJob(jobId, retryAt)
saveInsights({ job_id, post_id, window_label, metrics, raw })
```
~/instagram-insights-pipeline/.env.example
IG_ACCESS_TOKEN=your_long_lived_access_token_here IG_API_VERSION=v19.0 IG_ACCOUNT_ID=your_ig_business_account_id WEBHOOK_VERIFY_TOKEN=choose_any_secret_string_here PORT=3000 DB_PATH=./data/pipeline.db # Cron: every 5 minutes (runs the job runner inside the main process) CRON_SCHEDULE=*/5 * * * * RUN_RUNNER_INLINE=true
file
File ·
instagram-insights-pipeline/
├── index.js ← Entry point (Express server + optional inline cron)
├── package.json
├── schema.sql ← DB schema (posts, insight_jobs, insights tables)
├── .env.example ← Copy to .env and fill in your tokens
├── README.md
└── src/
├── webhook.js ← Receives Instagram webhook events
├── scheduler.js ← Registers 3 jobs per new post
├── runner.js ← Cron job: fetches due insights from IG API
├── instagram.js ← Graph API client (no external HTTP libs)
└── db.js ← SQLite adapter (swap for your own store)
| Stage | File | What it does |
|---|---|---|
| Detection | webhook.js | Meta sends a media event → extract media_id |
| Registration | scheduler.js | registerPost() → insert post + 3 jobs (24h / 15d / 30d) |
| Execution | runner.js | Every 5 min → SELECT jobs WHERE run_at <= NOW() → fetch API → save |
| API | instagram.js | Calls /{media_id}/insights with correct metrics per media type |
| Storage | db.js | SQLite adapter — swap out for your existing DB |
# 1. Install
npm install
# 2. Configure
cp .env.example .env
# Fill in: IG_ACCESS_TOKEN, IG_ACCOUNT_ID, WEBHOOK_VERIFY_TOKEN
# 3. Start (server + embedded cron every 5 min)
npm run start:cron
media field → set callback to https://your-domain.com/webhook/instagramregisterPost() twice for the same post ID is safe (no duplicate jobs)published_at, not from when you discovered the postdb.js exposes a clean interface; replace with Postgres/MySQL/your ORM without touching any other fileinstagram.js handles this automatically