Initialer Versionstand des eBay Trackers
This commit is contained in:
8
.gitignore
vendored
Normal file
8
.gitignore
vendored
Normal file
@@ -0,0 +1,8 @@
|
||||
# Passwörter und API-Keys ignorieren
|
||||
.env
|
||||
|
||||
# Dependencies / Pakete ignorieren
|
||||
node_modules/
|
||||
|
||||
# Logs ignorieren
|
||||
*.log
|
||||
14
Dockerfile
Normal file
14
Dockerfile
Normal file
@@ -0,0 +1,14 @@
|
||||
FROM node:20-alpine
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
# Dependencies installieren
|
||||
COPY package*.json ./
|
||||
RUN npm ci --only=production
|
||||
|
||||
# Sourcecode kopieren
|
||||
COPY . .
|
||||
|
||||
EXPOSE 3001
|
||||
|
||||
CMD ["npm", "start"]
|
||||
11
docker-compose.yml
Normal file
11
docker-compose.yml
Normal file
@@ -0,0 +1,11 @@
|
||||
version: '3.8'
|
||||
|
||||
services:
|
||||
ebay-tracker:
|
||||
build: .
|
||||
container_name: ebay-tracker
|
||||
restart: unless-stopped
|
||||
ports:
|
||||
- "3001:3001"
|
||||
env_file:
|
||||
- .env
|
||||
1174
package-lock.json
generated
Normal file
1174
package-lock.json
generated
Normal file
File diff suppressed because it is too large
Load Diff
15
package.json
Normal file
15
package.json
Normal file
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"name": "ebay-tracker",
|
||||
"version": "1.1.0",
|
||||
"type": "module",
|
||||
"main": "src/server.js",
|
||||
"scripts": {
|
||||
"start": "node src/server.js"
|
||||
},
|
||||
"dependencies": {
|
||||
"axios": "^1.7.2",
|
||||
"dotenv": "^16.4.5",
|
||||
"express": "^4.19.2",
|
||||
"pg": "^8.12.0"
|
||||
}
|
||||
}
|
||||
84
src/db.js
Normal file
84
src/db.js
Normal file
@@ -0,0 +1,84 @@
|
||||
import pg from 'pg';
|
||||
import dotenv from 'dotenv';
|
||||
|
||||
dotenv.config();
|
||||
|
||||
const { Pool } = pg;
|
||||
|
||||
export const pool = new Pool({
|
||||
host: process.env.PGHOST,
|
||||
port: parseInt(process.env.PGPORT || '5432'),
|
||||
database: process.env.PGDATABASE,
|
||||
user: process.env.PGUSER,
|
||||
password: process.env.PGPASSWORD,
|
||||
connectionTimeoutMillis: 5000,
|
||||
});
|
||||
|
||||
// UTF-8 erzwingen
|
||||
pool.on('connect', (client) => {
|
||||
client.query("SET client_encoding TO 'UTF8'");
|
||||
});
|
||||
|
||||
export async function initDb() {
|
||||
const client = await pool.connect();
|
||||
try {
|
||||
console.log('🔄 Prüfe und erstelle Datenbank-Tabellen...');
|
||||
|
||||
// 1. Tabelle für Verkäufer-Whitelist
|
||||
await client.query(`
|
||||
CREATE TABLE IF NOT EXISTS sellers (
|
||||
id SERIAL PRIMARY KEY,
|
||||
username VARCHAR(255) UNIQUE NOT NULL,
|
||||
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
`);
|
||||
|
||||
// 2. Tabelle für Telegram-Suchaufträge
|
||||
await client.query(`
|
||||
CREATE TABLE IF NOT EXISTS tracked_searches (
|
||||
id SERIAL PRIMARY KEY,
|
||||
chat_id VARCHAR(255) NOT NULL,
|
||||
query TEXT NOT NULL,
|
||||
negative_keywords TEXT[] DEFAULT '{}',
|
||||
is_active BOOLEAN DEFAULT TRUE,
|
||||
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
`);
|
||||
|
||||
// 3. Tabelle für Artikel-Snapshots (Preisverlauf & Alerts)
|
||||
await client.query(`
|
||||
CREATE TABLE IF NOT EXISTS item_snapshots (
|
||||
id SERIAL PRIMARY KEY,
|
||||
item_id VARCHAR(255) UNIQUE NOT NULL,
|
||||
search_id INTEGER REFERENCES tracked_searches(id) ON DELETE CASCADE,
|
||||
seller_username VARCHAR(255),
|
||||
title TEXT NOT NULL,
|
||||
price NUMERIC(10, 2) NOT NULL,
|
||||
currency VARCHAR(10) DEFAULT 'EUR',
|
||||
item_url TEXT,
|
||||
image_url TEXT,
|
||||
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
|
||||
last_updated TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
`);
|
||||
|
||||
// In initDb() ergänzen:
|
||||
await client.query(`
|
||||
ALTER TABLE item_snapshots ADD COLUMN IF NOT EXISTS lowest_price NUMERIC(10, 2);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS price_history (
|
||||
id SERIAL PRIMARY KEY,
|
||||
item_id VARCHAR(255) NOT NULL,
|
||||
price NUMERIC(10, 2) NOT NULL,
|
||||
recorded_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
`);
|
||||
|
||||
console.log('✅ PostgreSQL Tabellen erfolgreich initialisiert (UTF-8).');
|
||||
} catch (error) {
|
||||
console.error('❌ Fehler bei der Tabellen-Initialisierung:', error.message);
|
||||
throw error;
|
||||
} finally {
|
||||
client.release();
|
||||
}
|
||||
}
|
||||
45
src/ebayAuth.js
Normal file
45
src/ebayAuth.js
Normal file
@@ -0,0 +1,45 @@
|
||||
import axios from 'axios';
|
||||
import dotenv from 'dotenv';
|
||||
|
||||
dotenv.config();
|
||||
|
||||
let tokenCache = {
|
||||
token: null,
|
||||
expiresAt: 0
|
||||
};
|
||||
|
||||
export async function getEbayToken() {
|
||||
const now = Date.now();
|
||||
|
||||
if (tokenCache.token && now < tokenCache.expiresAt - 5 * 60 * 1000) {
|
||||
return tokenCache.token;
|
||||
}
|
||||
|
||||
console.log('🔄 Hole neuen eBay OAuth Access Token...');
|
||||
|
||||
const credentials = Buffer.from(
|
||||
`${process.env.EBAY_APP_ID}:${process.env.EBAY_CERT_ID}`
|
||||
).toString('base64');
|
||||
|
||||
try {
|
||||
const response = await axios.post(
|
||||
'https://api.ebay.com/identity/v1/oauth2/token',
|
||||
'grant_type=client_credentials&scope=https://api.ebay.com/oauth/api_scope',
|
||||
{
|
||||
headers: {
|
||||
'Content-Type': 'application/x-www-form-urlencoded',
|
||||
'Authorization': `Basic ${credentials}`
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
tokenCache.token = response.data.access_token;
|
||||
tokenCache.expiresAt = now + (response.data.expires_in * 1000);
|
||||
|
||||
console.log('✅ eBay Token erfolgreich erneuert!');
|
||||
return tokenCache.token;
|
||||
} catch (error) {
|
||||
console.error('❌ eBay Auth Fehler:', error.response?.data || error.message);
|
||||
throw new Error('eBay Authentifizierung fehlgeschlagen.');
|
||||
}
|
||||
}
|
||||
86
src/ebayService.js
Normal file
86
src/ebayService.js
Normal file
@@ -0,0 +1,86 @@
|
||||
import axios from 'axios';
|
||||
import { getEbayToken } from './ebayAuth.js';
|
||||
|
||||
// Hilfsfunktion: Aufteilen eines Arrays in kleinere Häppchen (Chunks)
|
||||
function chunkArray(array, chunkSize) {
|
||||
const chunks = [];
|
||||
for (let i = 0; i < array.length; i += chunkSize) {
|
||||
chunks.push(array.slice(i, i + chunkSize));
|
||||
}
|
||||
return chunks;
|
||||
}
|
||||
|
||||
/**
|
||||
* Durchsucht eBay nach ALLEN Whitelist-Verkäufern in der Datenbank
|
||||
*/
|
||||
export async function searchEbayProducts(query, allowedSellers = [], negativeKeywords = []) {
|
||||
if (!allowedSellers || allowedSellers.length === 0) {
|
||||
console.log('⚠️ Keine Verkäufer in der Whitelist hinterlegt.');
|
||||
return [];
|
||||
}
|
||||
|
||||
try {
|
||||
const token = await getEbayToken();
|
||||
|
||||
// 1. Verkäufernamen bereinigen
|
||||
const validSellers = allowedSellers
|
||||
.map(s => String(s).trim())
|
||||
.filter(Boolean);
|
||||
|
||||
if (validSellers.length === 0) return [];
|
||||
|
||||
// 2. Verkäufer in 25er-Pakete aufteilen
|
||||
const sellerBatches = chunkArray(validSellers, 25);
|
||||
let allRawItems = [];
|
||||
|
||||
// 3. Jedes 25er-Paket nacheinander bei eBay abfragen
|
||||
for (const batch of sellerBatches) {
|
||||
const formattedSellers = batch.map(s => encodeURIComponent(s)).join('|');
|
||||
const sellerFilterString = `sellers:{${formattedSellers}}`;
|
||||
|
||||
const response = await axios.get('https://api.ebay.com/buy/browse/v1/item_summary/search', {
|
||||
headers: {
|
||||
'Authorization': `Bearer ${token}`,
|
||||
'X-EBAY-C-MARKETPLACE-ID': 'EBAY-DE'
|
||||
},
|
||||
params: {
|
||||
q: query,
|
||||
filter: sellerFilterString,
|
||||
limit: 50,
|
||||
sort: 'newlyListed'
|
||||
}
|
||||
});
|
||||
|
||||
const items = response.data.itemSummaries || [];
|
||||
allRawItems = allRawItems.concat(items);
|
||||
}
|
||||
|
||||
// 4. Duplikate über die itemId entfernen
|
||||
const uniqueItemsMap = new Map();
|
||||
for (const item of allRawItems) {
|
||||
uniqueItemsMap.set(item.itemId, item);
|
||||
}
|
||||
const uniqueItems = Array.from(uniqueItemsMap.values());
|
||||
|
||||
// 5. Negativ-Keywords lokal ausfiltern
|
||||
const filteredItems = uniqueItems.filter(item => {
|
||||
const titleLower = item.title.toLowerCase();
|
||||
return !negativeKeywords.some(kw => titleLower.includes(kw.toLowerCase().trim()));
|
||||
});
|
||||
|
||||
return filteredItems.map(item => ({
|
||||
itemId: item.itemId,
|
||||
title: item.title,
|
||||
price: item.price?.value,
|
||||
currency: item.price?.currency || 'EUR',
|
||||
sellerUsername: item.seller?.username || 'Unbekannt',
|
||||
itemUrl: item.itemWebUrl,
|
||||
imageUrl: item.image?.imageUrl || null
|
||||
}));
|
||||
|
||||
} catch (error) {
|
||||
const apiError = error.response?.data?.errors?.[0]?.message || error.message;
|
||||
console.error(`❌ Fehler bei der eBay-Suche für "${query}":`, apiError);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
47
src/ebayService_BK_1.js
Normal file
47
src/ebayService_BK_1.js
Normal file
@@ -0,0 +1,47 @@
|
||||
import axios from 'axios';
|
||||
import { getEbayToken } from './ebayAuth.js';
|
||||
|
||||
export async function searchEbayProducts(query, negativeKeywords = [], sellersWhitelist = []) {
|
||||
const token = await getEbayToken();
|
||||
|
||||
let filterParam = '';
|
||||
if (sellersWhitelist.length > 0) {
|
||||
filterParam = `sellers:{${sellersWhitelist.join('|')}}`;
|
||||
}
|
||||
|
||||
try {
|
||||
const response = await axios.get('https://api.ebay.com/buy/browse/v1/item_summary/search', {
|
||||
headers: {
|
||||
'Authorization': `Bearer ${token}`,
|
||||
'X-EBAY-C-MARKETPLACE-ID': 'EBAY-DE'
|
||||
},
|
||||
params: {
|
||||
q: query,
|
||||
filter: filterParam || undefined,
|
||||
limit: 50
|
||||
}
|
||||
});
|
||||
|
||||
const rawItems = response.data.itemSummaries || [];
|
||||
|
||||
const filteredItems = rawItems.filter(item => {
|
||||
const titleLower = item.title.toLowerCase();
|
||||
if (!negativeKeywords || negativeKeywords.length === 0) return true;
|
||||
return !negativeKeywords.some(neg => titleLower.includes(neg.toLowerCase().trim()));
|
||||
});
|
||||
|
||||
return filteredItems.map(item => ({
|
||||
itemId: item.itemId,
|
||||
title: item.title,
|
||||
price: parseFloat(item.price?.value || 0),
|
||||
currency: item.price?.currency || 'EUR',
|
||||
sellerUsername: item.seller?.username || 'Unbekannt',
|
||||
itemUrl: item.itemWebUrl,
|
||||
imageUrl: item.image?.imageUrl || null
|
||||
}));
|
||||
|
||||
} catch (error) {
|
||||
console.error(`❌ eBay API Fehler für Query "${query}":`, error.response?.data || error.message);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
169
src/server.js
Normal file
169
src/server.js
Normal file
@@ -0,0 +1,169 @@
|
||||
import express from 'express';
|
||||
import dotenv from 'dotenv';
|
||||
import { initDb, pool } from './db.js';
|
||||
import { runTrackingCycle } from './tracker.js';
|
||||
import { searchEbayProducts } from './ebayService.js';
|
||||
|
||||
dotenv.config();
|
||||
|
||||
const app = express();
|
||||
app.use(express.json());
|
||||
|
||||
const PORT = process.env.PORT || 3001;
|
||||
|
||||
// Hilfsfunktion: Repariert doppelt-kodierte Umlaute (Mojibake)
|
||||
function fixEncoding(str) {
|
||||
if (!str || typeof str !== 'string') return str;
|
||||
// Wenn der String typische Mojibake-Zeichen (wie ü für ü) enthält
|
||||
if (str.includes('Ã')) {
|
||||
try {
|
||||
return Buffer.from(str, 'latin1').toString('utf-8');
|
||||
} catch (e) {
|
||||
return str;
|
||||
}
|
||||
}
|
||||
return str;
|
||||
}
|
||||
|
||||
|
||||
|
||||
// Healthcheck
|
||||
app.get('/health', (req, res) => res.json({ status: 'ok', timestamp: new Date() }));
|
||||
|
||||
// Endpunkt 1: Suchauftrag aus Telegram / n8n in Postgres speichern
|
||||
app.post('/searches/add', async (req, res) => {
|
||||
const { chatId, query, negativeKeywords } = req.body;
|
||||
|
||||
if (!chatId || !query) {
|
||||
return res.status(400).json({ error: 'chatId und query sind Pflichtfelder.' });
|
||||
}
|
||||
|
||||
// 🛠️ FIX: Robustes Parsing für negativeKeywords
|
||||
let parsedKeywords = [];
|
||||
|
||||
if (Array.isArray(negativeKeywords)) {
|
||||
// Bereits ein echtes JS-Array
|
||||
parsedKeywords = negativeKeywords;
|
||||
} else if (typeof negativeKeywords === 'string') {
|
||||
const trimmed = negativeKeywords.trim();
|
||||
if (trimmed.startsWith('[')) {
|
||||
// Wurde als JSON-String gesendet: ["a", "b"]
|
||||
try {
|
||||
parsedKeywords = JSON.parse(trimmed);
|
||||
} catch (e) {
|
||||
parsedKeywords = [];
|
||||
}
|
||||
} else if (trimmed.length > 0) {
|
||||
// Wurde als Komma-getrennter String gesendet: "controller, headset"
|
||||
parsedKeywords = trimmed.split(',').map(item => item.trim());
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
const result = await pool.query(
|
||||
`INSERT INTO tracked_searches (chat_id, query, negative_keywords)
|
||||
VALUES ($1, $2, $3) RETURNING id`,
|
||||
[chatId, query, parsedKeywords] // 👈 Hier wird das saubere JS-Array übergeben
|
||||
);
|
||||
|
||||
res.status(201).json({
|
||||
success: true,
|
||||
message: 'Suchauftrag erfolgreich gespeichert.',
|
||||
searchId: result.rows[0].id
|
||||
});
|
||||
} catch (error) {
|
||||
console.error('Fehler beim Speichern der Suche:', error.message);
|
||||
res.status(500).json({ error: 'Datenbankfehler beim Speichern.' });
|
||||
}
|
||||
});
|
||||
|
||||
// Endpunkt: Statische Liste von Verkäufern hinzufügen (unterstützt Arrays, Komma- & Zeilenumbrüche)
|
||||
app.post('/sellers/add', async (req, res) => {
|
||||
const { sellers, username } = req.body;
|
||||
|
||||
// Akzeptiert sowohl "sellers" (Liste) als auch "username" (Einzelwert)
|
||||
const rawInput = sellers || username;
|
||||
|
||||
if (!rawInput) {
|
||||
return res.status(400).json({ error: 'Bitte gib mindestens einen Verkäufer an.' });
|
||||
}
|
||||
|
||||
let sellerList = [];
|
||||
|
||||
if (Array.isArray(rawInput)) {
|
||||
// Falls n8n ein JSON-Array schickt: ["shop1", "shop2"]
|
||||
sellerList = rawInput;
|
||||
} else if (typeof rawInput === 'string') {
|
||||
// Trennt bei Kommas ODER Zeilenumbrüchen (\n)
|
||||
sellerList = rawInput.split(/[\n,]+/).map(s => s.trim()).filter(Boolean);
|
||||
}
|
||||
|
||||
if (sellerList.length === 0) {
|
||||
return res.status(400).json({ error: 'Keine gültigen Verkäufernamen gefunden.' });
|
||||
}
|
||||
|
||||
try {
|
||||
const newlyAdded = [];
|
||||
for (const name of sellerList) {
|
||||
const cleanName = fixEncoding(name)
|
||||
const result = await pool.query(
|
||||
`INSERT INTO sellers (username) VALUES ($1) ON CONFLICT (username) DO NOTHING RETURNING username`,
|
||||
[cleanName]
|
||||
);
|
||||
if (result.rows.length > 0) {
|
||||
newlyAdded.push(name);
|
||||
}
|
||||
}
|
||||
|
||||
res.json({
|
||||
success: true,
|
||||
message: `${sellerList.length} Verkäufer verarbeitet (${newlyAdded.length} neu hinzugefügt).`,
|
||||
totalProcessed: sellerList.length,
|
||||
newlyAdded: newlyAdded
|
||||
});
|
||||
|
||||
} catch (error) {
|
||||
console.error('Fehler beim Speichern der Verkäufer:', error.message);
|
||||
res.status(500).json({ error: 'Datenbankfehler.' });
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
// Endpunkt 3: Der /scrape Endpunkt für n8n (Schedule Trigger)
|
||||
app.post('/scrape', async (req, res) => {
|
||||
const { webhookUrl } = req.body; // Optional: Kann in n8n dynamisch mitgegeben werden
|
||||
|
||||
// Sofort 202 Accepted melden, damit n8n nicht in HTTP-Timeouts rennt
|
||||
res.status(202).json({
|
||||
message: 'Scrape-Job akzeptiert und wird im Hintergrund ausgeführt.'
|
||||
});
|
||||
|
||||
// Scraping asynchron ausführen
|
||||
runTrackingCycle(webhookUrl).catch(err => {
|
||||
console.error('Hintergrund-Fehler beim Scrapen:', err.message);
|
||||
});
|
||||
});
|
||||
|
||||
// Endpunkt für n8n (Schedule): Führt den Scraper aus & gibt Preisänderungen zurück
|
||||
app.post('/tracker/run', async (req, res) => {
|
||||
try {
|
||||
const changes = await runTrackingCycle();
|
||||
res.json({
|
||||
success: true,
|
||||
count: changes.length,
|
||||
priceChanges: changes
|
||||
});
|
||||
} catch (error) {
|
||||
res.status(500).json({ error: error.message });
|
||||
}
|
||||
});
|
||||
|
||||
// Server-Start
|
||||
async function startServer() {
|
||||
await initDb();
|
||||
app.listen(PORT, '0.0.0.0', () => {
|
||||
console.log(`🚀 eBay-Tracker Server läuft auf Port ${PORT} (Steuerung via n8n)`);
|
||||
});
|
||||
}
|
||||
|
||||
startServer();
|
||||
115
src/tracker.js
Normal file
115
src/tracker.js
Normal file
@@ -0,0 +1,115 @@
|
||||
import { pool } from './db.js';
|
||||
import { searchEbayProducts } from './ebayService.js';
|
||||
|
||||
export async function runTrackingCycle() {
|
||||
console.log('🚀 Starte eBay-Tracking-Durchlauf mit Historien-Vergleich...');
|
||||
const notifications = [];
|
||||
|
||||
try {
|
||||
const sellersResult = await pool.query('SELECT username FROM sellers');
|
||||
const allowedSellers = sellersResult.rows.map(row => row.username);
|
||||
|
||||
if (allowedSellers.length === 0) return notifications;
|
||||
|
||||
const searchesResult = await pool.query('SELECT * FROM tracked_searches WHERE is_active = TRUE');
|
||||
const searches = searchesResult.rows;
|
||||
|
||||
for (const search of searches) {
|
||||
const foundItems = await searchEbayProducts(
|
||||
search.query,
|
||||
allowedSellers,
|
||||
search.negative_keywords || []
|
||||
);
|
||||
|
||||
for (const item of foundItems) {
|
||||
const newPrice = parseFloat(item.price);
|
||||
|
||||
// 1. Artikel in der DB suchen (Preis + bisherigen Tiefstpreis abfragen)
|
||||
const existing = await pool.query(
|
||||
'SELECT price, lowest_price FROM item_snapshots WHERE item_id = $1',
|
||||
[item.itemId]
|
||||
);
|
||||
|
||||
if (existing.rows.length === 0) {
|
||||
// 🆕 BRANDNEUER ARTIKEL
|
||||
notifications.push({
|
||||
type: 'NEW',
|
||||
title: item.title,
|
||||
oldPrice: null,
|
||||
newPrice: newPrice.toFixed(2),
|
||||
lowestPrice: newPrice.toFixed(2),
|
||||
isAllTimeLow: true,
|
||||
diff: 'NEU',
|
||||
seller: item.sellerUsername,
|
||||
itemUrl: item.itemUrl
|
||||
});
|
||||
|
||||
// In Preishistorie eintragen
|
||||
await pool.query(
|
||||
'INSERT INTO price_history (item_id, price) VALUES ($1, $2)',
|
||||
[item.itemId, newPrice]
|
||||
);
|
||||
|
||||
} else {
|
||||
// 🏷️ BEREITS BEKANNTER ARTIKEL
|
||||
const oldPrice = parseFloat(existing.rows[0].price);
|
||||
const currentLowest = existing.rows[0].lowest_price
|
||||
? parseFloat(existing.rows[0].lowest_price)
|
||||
: oldPrice;
|
||||
|
||||
if (oldPrice !== newPrice) {
|
||||
const diffVal = newPrice - oldPrice;
|
||||
const diffStr = diffVal > 0 ? `+${diffVal.toFixed(2)}` : `${diffVal.toFixed(2)}`;
|
||||
const isAllTimeLow = newPrice < currentLowest;
|
||||
const newLowest = Math.min(currentLowest, newPrice);
|
||||
|
||||
notifications.push({
|
||||
type: 'PRICE_CHANGE',
|
||||
title: item.title,
|
||||
oldPrice: oldPrice.toFixed(2),
|
||||
newPrice: newPrice.toFixed(2),
|
||||
lowestPrice: newLowest.toFixed(2),
|
||||
isAllTimeLow: isAllTimeLow,
|
||||
diff: diffStr,
|
||||
seller: item.sellerUsername,
|
||||
itemUrl: item.itemUrl
|
||||
});
|
||||
|
||||
// Preisänderung in Historie-Tabelle speichern
|
||||
await pool.query(
|
||||
'INSERT INTO price_history (item_id, price) VALUES ($1, $2)',
|
||||
[item.itemId, newPrice]
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Snapshot aktualisieren (LEAST stellt sicher, dass der lowest_price nie überschrieben wird)
|
||||
await pool.query(`
|
||||
INSERT INTO item_snapshots (
|
||||
item_id, search_id, seller_username, title, price, lowest_price, currency, item_url, image_url
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $5, $6, $7, $8)
|
||||
ON CONFLICT (item_id) DO UPDATE
|
||||
SET price = EXCLUDED.price,
|
||||
lowest_price = LEAST(item_snapshots.lowest_price, EXCLUDED.price),
|
||||
last_updated = CURRENT_TIMESTAMP
|
||||
`, [
|
||||
item.itemId,
|
||||
search.id,
|
||||
item.sellerUsername,
|
||||
item.title,
|
||||
item.price,
|
||||
item.currency,
|
||||
item.itemUrl,
|
||||
item.imageUrl
|
||||
]);
|
||||
}
|
||||
}
|
||||
|
||||
return notifications;
|
||||
|
||||
} catch (error) {
|
||||
console.error('❌ Fehler während des Tracking-Durchlaufs:', error.message);
|
||||
return notifications;
|
||||
}
|
||||
}
|
||||
110
src/tracker_BK_1.js
Normal file
110
src/tracker_BK_1.js
Normal file
@@ -0,0 +1,110 @@
|
||||
import axios from 'axios';
|
||||
import { pool } from './db.js';
|
||||
import { searchEbayProducts } from './ebayService.js';
|
||||
|
||||
export async function runTrackingCycle(targetWebhookUrl) {
|
||||
console.log(`\n🕒 [${new Date().toISOString()}] Starte Tracking-Durchlauf via n8n Trigger...`);
|
||||
|
||||
const webhookUrl = targetWebhookUrl || process.env.N8N_ALERT_WEBHOOK_URL;
|
||||
|
||||
try {
|
||||
// 1. Alle aktiven Suchaufträge aus Postgres holen
|
||||
const searchesRes = await pool.query('SELECT * FROM tracked_searches WHERE is_active = TRUE');
|
||||
const searches = searchesRes.rows;
|
||||
|
||||
// 2. Verkäufer-Whitelist aus Postgres holen
|
||||
const sellersRes = await pool.query('SELECT username FROM sellers');
|
||||
const sellerList = sellersRes.rows.map(r => r.username);
|
||||
|
||||
if (searches.length === 0) {
|
||||
console.log('ℹ️ Keine aktiven Suchen in der Datenbank.');
|
||||
return { status: 'success', processedSearches: 0, alertsFound: 0 };
|
||||
}
|
||||
|
||||
const alertsToSend = [];
|
||||
|
||||
for (const search of searches) {
|
||||
console.log(`🔎 Prüfe Suche ID ${search.id}: "${search.query}" (Chat: ${search.chat_id})...`);
|
||||
|
||||
const currentItems = await searchEbayProducts(
|
||||
search.query,
|
||||
search.negative_keywords,
|
||||
sellerList
|
||||
);
|
||||
|
||||
for (const item of currentItems) {
|
||||
// Prüfen, ob Artikel bereits in Postgres existiert
|
||||
const existingRes = await pool.query(
|
||||
'SELECT price FROM item_snapshots WHERE item_id = $1',
|
||||
[item.itemId]
|
||||
);
|
||||
|
||||
if (existingRes.rows.length === 0) {
|
||||
// NEUER ARTIKEL
|
||||
await pool.query(`
|
||||
INSERT INTO item_snapshots (item_id, search_id, seller_username, title, price, currency, item_url, image_url)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||
`, [item.itemId, search.id, item.sellerUsername, item.title, item.price, item.currency, item.itemUrl, item.imageUrl]);
|
||||
|
||||
alertsToSend.push({
|
||||
type: 'NEW_ITEM',
|
||||
chatId: search.chat_id,
|
||||
searchQuery: search.query,
|
||||
title: item.title,
|
||||
newPrice: item.price,
|
||||
oldPrice: null,
|
||||
currency: item.currency,
|
||||
seller: item.sellerUsername,
|
||||
url: item.itemUrl,
|
||||
imageUrl: item.imageUrl
|
||||
});
|
||||
|
||||
} else {
|
||||
const oldPrice = parseFloat(existingRes.rows[0].price);
|
||||
|
||||
if (item.price < oldPrice) {
|
||||
// PREISSTURZ
|
||||
await pool.query(`
|
||||
UPDATE item_snapshots
|
||||
SET price = $1, last_updated = CURRENT_TIMESTAMP
|
||||
WHERE item_id = $2
|
||||
`, [item.price, item.itemId]);
|
||||
|
||||
alertsToSend.push({
|
||||
type: 'PRICE_DROP',
|
||||
chatId: search.chat_id,
|
||||
searchQuery: search.query,
|
||||
title: item.title,
|
||||
newPrice: item.price,
|
||||
oldPrice: oldPrice,
|
||||
currency: item.currency,
|
||||
seller: item.sellerUsername,
|
||||
url: item.itemUrl,
|
||||
imageUrl: item.imageUrl
|
||||
});
|
||||
|
||||
} else {
|
||||
// Unverändert: Timestamp aktualisieren
|
||||
await pool.query('UPDATE item_snapshots SET last_updated = CURRENT_TIMESTAMP WHERE item_id = $1', [item.itemId]);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 3. Alerts per Webhook an n8n senden
|
||||
if (alertsToSend.length > 0 && webhookUrl) {
|
||||
console.log(`🚀 Sende ${alertsToSend.length} Alerts an n8n (${webhookUrl})...`);
|
||||
await axios.post(webhookUrl, { alerts: alertsToSend });
|
||||
}
|
||||
|
||||
return {
|
||||
status: 'success',
|
||||
processedSearches: searches.length,
|
||||
alertsFound: alertsToSend.length
|
||||
};
|
||||
|
||||
} catch (error) {
|
||||
console.error('❌ Fehler im Tracking-Durchlauf:', error.message);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user