first commit
This commit is contained in:
@@ -0,0 +1,21 @@
|
||||
const { apiKey } = require("./config");
|
||||
|
||||
function requireApiKey(req, res, next) {
|
||||
// allow health check without api key (optional)
|
||||
if (req.path === "/healthz") return next();
|
||||
|
||||
if (!apiKey) {
|
||||
return res.status(500).json({ error: "SMS_GATEWAY_API_KEY is not set on server" });
|
||||
}
|
||||
|
||||
const key = req.get("X-Api-Key");
|
||||
if (!key || key !== apiKey) {
|
||||
return res.status(401).json({ error: "unauthorized" });
|
||||
}
|
||||
next();
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
requireApiKey,
|
||||
};
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
const dotenv = require("dotenv");
|
||||
const path = require("path");
|
||||
|
||||
// 无论从哪个目录启动,都强制读取 backend 根目录的 .env
|
||||
const envPath = path.resolve(__dirname, "..", ".env");
|
||||
dotenv.config({ path: envPath });
|
||||
|
||||
function env(name, fallback) {
|
||||
const v = process.env[name];
|
||||
if (v === undefined || v === "") return fallback;
|
||||
return v;
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
port: Number(env("PORT", 7788)),
|
||||
dataDir: env("DATA_DIR", "./data"),
|
||||
dbFile: env("DB_FILE", "./data/sms-gateway.db"),
|
||||
apiKey: env("SMS_GATEWAY_API_KEY", ""),
|
||||
};
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
const fs = require("fs");
|
||||
const path = require("path");
|
||||
const Database = require("better-sqlite3");
|
||||
const { dbFile, dataDir } = require("./config");
|
||||
|
||||
function ensureDir(dir) {
|
||||
if (!fs.existsSync(dir)) fs.mkdirSync(dir, { recursive: true });
|
||||
}
|
||||
|
||||
ensureDir(path.resolve(dataDir));
|
||||
|
||||
const db = new Database(dbFile);
|
||||
db.pragma("journal_mode = WAL");
|
||||
|
||||
function initSchema() {
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS inbound_sms (
|
||||
id TEXT PRIMARY KEY,
|
||||
device_id TEXT NOT NULL,
|
||||
sender TEXT,
|
||||
content TEXT NOT NULL,
|
||||
received_at INTEGER NOT NULL,
|
||||
parsed_code TEXT,
|
||||
parse_status TEXT,
|
||||
raw_pdu_hash TEXT,
|
||||
created_at INTEGER NOT NULL,
|
||||
UNIQUE(device_id, raw_pdu_hash, received_at)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS outbound_tasks (
|
||||
task_id TEXT PRIMARY KEY,
|
||||
device_id TEXT NOT NULL,
|
||||
phone TEXT NOT NULL,
|
||||
content TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'pending', -- pending/sending/success/failed
|
||||
retry_count INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT,
|
||||
created_at INTEGER NOT NULL,
|
||||
updated_at INTEGER NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS device_heartbeats (
|
||||
device_id TEXT PRIMARY KEY,
|
||||
last_heartbeat_at INTEGER NOT NULL,
|
||||
last_ip TEXT
|
||||
);
|
||||
`);
|
||||
}
|
||||
|
||||
initSchema();
|
||||
|
||||
module.exports = {
|
||||
db,
|
||||
};
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
const express = require("express");
|
||||
const cors = require("cors");
|
||||
const { port } = require("./config");
|
||||
const { requireApiKey } = require("./auth");
|
||||
const routes = require("./routes");
|
||||
|
||||
const app = express();
|
||||
|
||||
app.use(cors());
|
||||
app.use(express.json({ limit: "1mb" }));
|
||||
|
||||
// 统一鉴权(除 healthz)
|
||||
app.use(requireApiKey);
|
||||
|
||||
app.use(routes);
|
||||
|
||||
app.use((req, res) => {
|
||||
res.status(404).json({ error: "not found" });
|
||||
});
|
||||
|
||||
app.listen(port, () => {
|
||||
// eslint-disable-next-line no-console
|
||||
console.log(`[sms-gateway-backend] listening on http://localhost:${port}`);
|
||||
});
|
||||
|
||||
@@ -0,0 +1,237 @@
|
||||
const express = require("express");
|
||||
const { db } = require("./db");
|
||||
const { nowMs, uuid, safeString } = require("./utils");
|
||||
|
||||
const router = express.Router();
|
||||
|
||||
function requireJsonKeys(req, res, keys) {
|
||||
for (const k of keys) {
|
||||
if (req.body?.[k] === undefined || req.body?.[k] === null || req.body?.[k] === "") {
|
||||
return res.status(400).json({ error: `missing/empty field: ${k}` });
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function getDeviceIdFromHeader(req) {
|
||||
// device_id 在 DB 中用来“归属”网关;
|
||||
// 这里用 apiKey 本身作为网关标识,实现“只要 apiKey 不要 deviceId”的效果。
|
||||
return safeString(req.get("X-Api-Key"));
|
||||
}
|
||||
|
||||
// 1) 安卓端上报入站短信
|
||||
router.post("/api/v1/sms/inbound", (req, res) => {
|
||||
const err = requireJsonKeys(req, res, ["sender", "content", "parseStatus"]);
|
||||
if (err) return;
|
||||
|
||||
const deviceId = getDeviceIdFromHeader(req);
|
||||
if (!deviceId) return res.status(401).json({ error: "missing api key for device identification" });
|
||||
const sender = safeString(req.body.sender);
|
||||
const content = safeString(req.body.content);
|
||||
const receivedAt = req.body.receivedAt !== undefined ? Number(req.body.receivedAt) : null;
|
||||
const parsedCode = req.body.parsedCode ? safeString(req.body.parsedCode) : null;
|
||||
const parseStatus = safeString(req.body.parseStatus);
|
||||
const rawPduHash = req.body.rawPduHash ? safeString(req.body.rawPduHash) : null;
|
||||
|
||||
const id = uuid();
|
||||
const createdAt = nowMs();
|
||||
|
||||
try {
|
||||
const stmt = db.prepare(`
|
||||
INSERT OR IGNORE INTO inbound_sms
|
||||
(id, device_id, sender, content, received_at, parsed_code, parse_status, raw_pdu_hash, created_at)
|
||||
VALUES
|
||||
(@id, @device_id, @sender, @content, @received_at, @parsed_code, @parse_status, @raw_pdu_hash, @created_at)
|
||||
`);
|
||||
|
||||
stmt.run({
|
||||
id,
|
||||
device_id: deviceId,
|
||||
sender,
|
||||
content,
|
||||
received_at: Number.isFinite(receivedAt) ? receivedAt : createdAt,
|
||||
parsed_code: parsedCode,
|
||||
parse_status: parseStatus,
|
||||
raw_pdu_hash: rawPduHash,
|
||||
created_at: createdAt,
|
||||
});
|
||||
|
||||
res.json({ ackId: id });
|
||||
} catch (e) {
|
||||
res.status(500).json({ error: "failed to store inbound sms", detail: String(e?.message || e) });
|
||||
}
|
||||
});
|
||||
|
||||
// 2) 安卓端拉取设备任务
|
||||
router.get("/api/v1/device/tasks", (req, res) => {
|
||||
const deviceId = getDeviceIdFromHeader(req);
|
||||
if (!deviceId) return res.status(401).json({ error: "missing api key for device identification" });
|
||||
|
||||
const limit = Math.min(Number(req.query.limit || 5), 20);
|
||||
const t = nowMs();
|
||||
|
||||
try {
|
||||
const rows = db
|
||||
.prepare(
|
||||
`
|
||||
SELECT task_id, phone, content, retry_count, last_error, created_at, updated_at
|
||||
FROM outbound_tasks
|
||||
WHERE device_id = ?
|
||||
AND status = 'pending'
|
||||
ORDER BY created_at ASC
|
||||
LIMIT ?
|
||||
`
|
||||
)
|
||||
.all(deviceId, limit);
|
||||
|
||||
if (rows.length === 0) return res.json({ tasks: [] });
|
||||
|
||||
// 标记为 sending(防止多端重复拉取)
|
||||
const mark = db.prepare(`
|
||||
UPDATE outbound_tasks
|
||||
SET status = 'sending',
|
||||
updated_at = ?,
|
||||
retry_count = retry_count + 1
|
||||
WHERE task_id IN (${rows.map(() => "?").join(",")})
|
||||
`);
|
||||
mark.run(t, ...rows.map((r) => r.task_id));
|
||||
|
||||
res.json({
|
||||
tasks: rows.map((r) => ({
|
||||
taskId: r.task_id,
|
||||
phone: r.phone,
|
||||
content: r.content,
|
||||
retryCount: r.retry_count,
|
||||
lastError: r.last_error,
|
||||
})),
|
||||
});
|
||||
} catch (e) {
|
||||
res.status(500).json({ error: "failed to fetch tasks", detail: String(e?.message || e) });
|
||||
}
|
||||
});
|
||||
|
||||
// 3) 安卓端回传结果
|
||||
router.post("/api/v1/sms/outbound/result", (req, res) => {
|
||||
const err = requireJsonKeys(req, res, ["taskId", "status"]);
|
||||
if (err) return;
|
||||
|
||||
const deviceId = getDeviceIdFromHeader(req);
|
||||
if (!deviceId) return res.status(401).json({ error: "missing api key for device identification" });
|
||||
const taskId = safeString(req.body.taskId);
|
||||
const status = safeString(req.body.status);
|
||||
const error = req.body.error ? safeString(req.body.error) : null;
|
||||
|
||||
if (!["success", "failed"].includes(status)) {
|
||||
return res.status(400).json({ error: "invalid status" });
|
||||
}
|
||||
|
||||
try {
|
||||
db.prepare(
|
||||
`
|
||||
UPDATE outbound_tasks
|
||||
SET status = @status,
|
||||
last_error = @last_error,
|
||||
updated_at = @t
|
||||
WHERE task_id = @task_id AND device_id = @device_id
|
||||
`
|
||||
).run({
|
||||
status,
|
||||
last_error: error,
|
||||
t: nowMs(),
|
||||
task_id: taskId,
|
||||
device_id: deviceId,
|
||||
});
|
||||
|
||||
res.json({ ok: true });
|
||||
} catch (e) {
|
||||
res.status(500).json({ error: "failed to store result", detail: String(e?.message || e) });
|
||||
}
|
||||
});
|
||||
|
||||
// 4) 业务侧入队发送任务(MVP 用)
|
||||
router.post("/api/v1/business/outbound-tasks", (req, res) => {
|
||||
const err = requireJsonKeys(req, res, ["phone", "content"]);
|
||||
if (err) return;
|
||||
|
||||
const deviceId = getDeviceIdFromHeader(req);
|
||||
if (!deviceId) return res.status(401).json({ error: "missing api key for device identification" });
|
||||
const phone = safeString(req.body.phone);
|
||||
const content = safeString(req.body.content);
|
||||
|
||||
const taskId = uuid();
|
||||
const t = nowMs();
|
||||
|
||||
try {
|
||||
db.prepare(
|
||||
`
|
||||
INSERT INTO outbound_tasks
|
||||
(task_id, device_id, phone, content, status, retry_count, last_error, created_at, updated_at)
|
||||
VALUES
|
||||
(@task_id, @device_id, @phone, @content, 'pending', 0, NULL, @t, @t)
|
||||
`
|
||||
).run({
|
||||
task_id: taskId,
|
||||
device_id: deviceId,
|
||||
phone,
|
||||
content,
|
||||
t,
|
||||
});
|
||||
|
||||
res.json({ taskId });
|
||||
} catch (e) {
|
||||
res.status(500).json({ error: "failed to enqueue task", detail: String(e?.message || e) });
|
||||
}
|
||||
});
|
||||
|
||||
// 5) 业务侧读取入站验证码(MVP 用)
|
||||
router.get("/api/v1/business/inbound-sms", (req, res) => {
|
||||
const deviceId = getDeviceIdFromHeader(req);
|
||||
if (!deviceId) return res.status(401).json({ error: "missing api key for device identification" });
|
||||
|
||||
const limit = Math.min(Number(req.query.limit || 20), 100);
|
||||
const since = Number(req.query.since || 0);
|
||||
const onlyMatched = String(req.query.onlyMatched || "").toLowerCase() === "true";
|
||||
|
||||
try {
|
||||
const where = ["device_id = ?", "received_at >= ?"];
|
||||
const params = [deviceId, since];
|
||||
|
||||
if (onlyMatched) {
|
||||
where.push("parse_status = 'matched'");
|
||||
where.push("parsed_code IS NOT NULL");
|
||||
}
|
||||
|
||||
const rows = db
|
||||
.prepare(
|
||||
`
|
||||
SELECT id, sender, content, received_at, parsed_code, parse_status, raw_pdu_hash
|
||||
FROM inbound_sms
|
||||
WHERE ${where.join(" AND ")}
|
||||
ORDER BY received_at DESC
|
||||
LIMIT ?
|
||||
`
|
||||
)
|
||||
.all(...params, limit);
|
||||
|
||||
res.json({
|
||||
inbounds: rows.map((r) => ({
|
||||
id: r.id,
|
||||
sender: r.sender,
|
||||
content: r.content,
|
||||
receivedAt: r.received_at,
|
||||
parsedCode: r.parsed_code,
|
||||
parseStatus: r.parse_status,
|
||||
rawPduHash: r.raw_pdu_hash,
|
||||
})),
|
||||
});
|
||||
} catch (e) {
|
||||
res.status(500).json({ error: "failed to query inbound sms", detail: String(e?.message || e) });
|
||||
}
|
||||
});
|
||||
|
||||
router.get("/healthz", (req, res) => {
|
||||
res.json({ ok: true, ts: nowMs() });
|
||||
});
|
||||
|
||||
module.exports = router;
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
const crypto = require("crypto");
|
||||
|
||||
function nowMs() {
|
||||
return Date.now();
|
||||
}
|
||||
|
||||
function uuid() {
|
||||
if (crypto.randomUUID) return crypto.randomUUID();
|
||||
return crypto.randomBytes(16).toString("hex");
|
||||
}
|
||||
|
||||
function safeString(v) {
|
||||
if (v === undefined || v === null) return "";
|
||||
return String(v);
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
nowMs,
|
||||
uuid,
|
||||
safeString,
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user