8 Commits

17 changed files with 949 additions and 20 deletions
+4
View File
@@ -118,6 +118,10 @@ dist
.yarn/build-state.yml .yarn/build-state.yml
.yarn/install-state.gz .yarn/install-state.gz
.pnp.* .pnp.*
# Local Python environments
translation-service/.venv/
translation-service/__pycache__/
# Logs # Logs
logs logs
*.log *.log
+29
View File
@@ -31,6 +31,27 @@ A step by step series of examples that tell you how to get a development env run
npm start npm start
``` ```
### Local MarianMT translation
The translation service is independent from the Node backend and provides `GET /health` and `POST /translate`.
1. Start it locally:
```
cd translation-service
python3 -m venv .venv
.venv/bin/pip install -r requirements.txt
.venv/bin/python server.py
```
2. Configure the Node backend:
```
TRANSLATION_PROVIDER=marian
MARIAN_TRANSLATION_URL=http://127.0.0.1:8000
```
With Docker Compose, the service runs as the internal `translation` service. Set `TRANSLATION_PROVIDER=marian` before running `docker compose up`. MarianMT models download only when first needed and are stored in the `translation-models` Docker volume.
Supported languages are English (`en`), Spanish (`es`), French (`fr`), Danish (`da`), and Arabic (`ar`). Non-English pairs translate through English. Keep the service on the internal Docker network; it has no public port mapping.
### API Documentation ### API Documentation
Once the server is running, you can access the interactive API documentation powered by Swagger UI at: Once the server is running, you can access the interactive API documentation powered by Swagger UI at:
@@ -130,6 +151,14 @@ The API is divided into several sections based on functionality. Most routes und
- `GET /chapters/:chapterId`: Get the content of a chapter. - `GET /chapters/:chapterId`: Get the content of a chapter.
- `GET /chapters/:chapterId/verses`: Get the verses of a chapter. - `GET /chapters/:chapterId/verses`: Get the verses of a chapter.
- `GET /search`: Search the Bible. - `GET /search`: Search the Bible.
- `GET /mine`: Get Bible highlights and notes for the active profile.
- `GET /verses/:verseKey/counters`: Get highlight and note counters for a verse key like `GEN.1.1`.
- `GET /verses/:verseKey/activity`: Get profiles that highlighted/noted a verse and recent notes, limited to 10 by default.
- `POST /verses/:verseKey/highlight`: Highlight a verse for the active profile.
- `DELETE /verses/:verseKey/highlight`: Remove a verse highlight for the active profile.
- `POST /verses/:verseKey/notes`: Add a note for a verse.
- `PUT /verses/:verseKey/notes/:noteId`: Update one of the active profile's notes.
- `DELETE /verses/:verseKey/notes/:noteId`: Delete one of the active profile's notes.
### Subsplash (`/subsplash`) ### Subsplash (`/subsplash`)
+287
View File
@@ -0,0 +1,287 @@
const DBName = "EMI_SOCIAL";
const toProfileIds = (DB, profileIds = []) => {
return profileIds
.filter((profileId) => DB.ObjectID.isValid(profileId))
.map((profileId) => DB.ObjectID(profileId));
};
const publicProfile = (profile) => {
if (!profile) return null;
return {
_id: profile._id,
profile: profile.profile,
username: profile.username,
isGroup: profile.isGroup,
isCourse: profile.isCourse,
};
};
const bibleDB = (DB) => {
DB.bibleVerseCols = DB.db.db(DBName).collection("bible_verses");
DB.bibleVerseCols.createIndex({ highlightedBy: 1 }).catch(console.error);
DB.bibleVerseCols.createIndex({ notedBy: 1 }).catch(console.error);
const getVerseDoc = async (verseKey) => {
return DB.bibleVerseCols.findOne({ _id: verseKey }).catch((err) => {
console.log(err);
return false;
});
};
DB.getBibleProfileData = async (profileid) => {
if (!DB.ObjectID.isValid(profileid)) return false;
const profile = await DB.profileCols.findOne(
{ _id: DB.ObjectID(profileid) },
{ projection: { bibleHighlights: 1, bibleNotes: 1 } }
).catch((err) => {
console.log(err);
return false;
});
return {
highlights: profile?.bibleHighlights || [],
notes: profile?.bibleNotes || [],
};
};
DB.getBibleVerseCounters = async (verseKey) => {
const verse = await getVerseDoc(verseKey);
const highlightedBy = Array.isArray(verse?.highlightedBy) ? verse.highlightedBy : [];
const notedBy = Array.isArray(verse?.notedBy) ? verse.notedBy : [];
return {
verseKey,
highlightCount: highlightedBy.length,
noteProfileCount: notedBy.length,
notesCount: verse?.notesCount || 0,
};
};
DB.getBibleVerseActivity = async (verseKey, limit = 10) => {
const verse = await getVerseDoc(verseKey);
const highlightedBy = Array.isArray(verse?.highlightedBy) ? verse.highlightedBy : [];
const notedBy = Array.isArray(verse?.notedBy) ? verse.notedBy : [];
const highlightProfiles = toProfileIds(DB, highlightedBy);
const noteProfiles = toProfileIds(DB, notedBy);
const [highlights, noteUsers, notes] = await Promise.all([
highlightProfiles.length
? DB.profileCols.find({ _id: { $in: highlightProfiles } }).project({
profile: 1,
username: 1,
isGroup: 1,
isCourse: 1,
}).toArray()
: [],
noteProfiles.length
? DB.profileCols.find({ _id: { $in: noteProfiles } }).project({
profile: 1,
username: 1,
isGroup: 1,
isCourse: 1,
}).toArray()
: [],
DB.profileCols.aggregate([
{ $match: { "bibleNotes.verseKey": verseKey } },
{ $unwind: "$bibleNotes" },
{ $match: { "bibleNotes.verseKey": verseKey } },
{ $sort: { "bibleNotes.updatedAt": -1, "bibleNotes.createdAt": -1 } },
{ $limit: limit },
{
$project: {
_id: 0,
note: "$bibleNotes",
profile: {
_id: "$_id",
profile: "$profile",
username: "$username",
isGroup: "$isGroup",
isCourse: "$isCourse",
}
}
}
]).toArray(),
]);
return {
counters: await DB.getBibleVerseCounters(verseKey),
highlightedBy: highlights.map(publicProfile),
notedBy: noteUsers.map(publicProfile),
notes,
};
};
DB.addBibleHighlight = async (profileid, verseKey) => {
if (!DB.ObjectID.isValid(profileid)) return false;
const now = new Date();
const _id = DB.ObjectID(profileid);
const profileUpdate = await DB.profileCols.updateOne(
{ _id, "bibleHighlights.verseKey": { $ne: verseKey } },
{
$push: {
bibleHighlights: {
verseKey,
createdAt: now,
}
},
$set: { lastUpdate: now }
}
).catch((err) => {
console.log(err);
return false;
});
await DB.bibleVerseCols.updateOne(
{ _id: verseKey },
{
$setOnInsert: { createdAt: now },
$set: { updatedAt: now },
$addToSet: { highlightedBy: profileid + "" }
},
{ upsert: true }
).catch((err) => {
console.log(err);
return false;
});
if (DB.clearProfileCache) DB.clearProfileCache(profileid);
return profileUpdate;
};
DB.removeBibleHighlight = async (profileid, verseKey) => {
if (!DB.ObjectID.isValid(profileid)) return false;
const now = new Date();
const _id = DB.ObjectID(profileid);
const profileUpdate = await DB.profileCols.updateOne(
{ _id },
{
$pull: { bibleHighlights: { verseKey } },
$set: { lastUpdate: now }
}
).catch((err) => {
console.log(err);
return false;
});
await DB.bibleVerseCols.updateOne(
{ _id: verseKey },
{
$pull: { highlightedBy: profileid + "" },
$set: { updatedAt: now }
}
).catch((err) => {
console.log(err);
return false;
});
if (DB.clearProfileCache) DB.clearProfileCache(profileid);
return profileUpdate;
};
DB.addBibleNote = async (profileid, verseKey, text) => {
if (!DB.ObjectID.isValid(profileid)) return false;
const now = new Date();
const note = {
_id: DB.ObjectID().toString(),
verseKey,
text,
createdAt: now,
updatedAt: now,
};
const _id = DB.ObjectID(profileid);
const profileUpdate = await DB.profileCols.updateOne(
{ _id },
{
$push: { bibleNotes: note },
$set: { lastUpdate: now }
}
).catch((err) => {
console.log(err);
return false;
});
await DB.bibleVerseCols.updateOne(
{ _id: verseKey },
{
$setOnInsert: { createdAt: now },
$set: { updatedAt: now },
$addToSet: { notedBy: profileid + "" },
$inc: { notesCount: 1 }
},
{ upsert: true }
).catch((err) => {
console.log(err);
return false;
});
if (DB.clearProfileCache) DB.clearProfileCache(profileid);
return profileUpdate ? note : false;
};
DB.updateBibleNote = async (profileid, verseKey, noteId, text) => {
if (!DB.ObjectID.isValid(profileid)) return false;
const now = new Date();
const _id = DB.ObjectID(profileid);
const result = await DB.profileCols.findOneAndUpdate(
{ _id, bibleNotes: { $elemMatch: { _id: noteId, verseKey } } },
{
$set: {
"bibleNotes.$.text": text,
"bibleNotes.$.updatedAt": now,
lastUpdate: now,
}
},
{
returnOriginal: false,
projection: { bibleNotes: 1 }
}
).catch((err) => {
console.log(err);
return false;
});
if (DB.clearProfileCache) DB.clearProfileCache(profileid);
return result?.value?.bibleNotes?.find((note) => note._id === noteId) || null;
};
DB.removeBibleNote = async (profileid, verseKey, noteId) => {
if (!DB.ObjectID.isValid(profileid)) return false;
const now = new Date();
const _id = DB.ObjectID(profileid);
const result = await DB.profileCols.updateOne(
{ _id, bibleNotes: { $elemMatch: { _id: noteId, verseKey } } },
{
$pull: { bibleNotes: { _id: noteId, verseKey } },
$set: { lastUpdate: now }
}
).catch((err) => {
console.log(err);
return false;
});
if (result?.modifiedCount) {
const stillHasNotes = await DB.profileCols.findOne(
{ _id, "bibleNotes.verseKey": verseKey },
{ projection: { _id: 1 } }
);
const verseUpdate = {
$inc: { notesCount: -1 },
$set: { updatedAt: now }
};
if (!stillHasNotes) {
verseUpdate.$pull = { notedBy: profileid + "" };
}
await DB.bibleVerseCols.updateOne({ _id: verseKey }, verseUpdate).catch((err) => {
console.log(err);
return false;
});
}
if (DB.clearProfileCache) DB.clearProfileCache(profileid);
return result;
};
}
module.exports = bibleDB;
+4
View File
@@ -5,6 +5,10 @@ let userProfileCache = {};
userDB = (DB) => { userDB = (DB) => {
DB.profileCols = DB.db.db(DBName).collection("profiles"); DB.profileCols = DB.db.db(DBName).collection("profiles");
DB.clearProfileCache = (profileid) => {
if (userProfileCache[profileid]) delete userProfileCache[profileid];
};
DB.newProfile = (profileObj) => { DB.newProfile = (profileObj) => {
return DB.profileCols.insertOne(profileObj.toObj()).catch((err) => { return DB.profileCols.insertOne(profileObj.toObj()).catch((err) => {
console.log(err); console.log(err);
+2
View File
@@ -17,6 +17,8 @@ class User {
this.lastUpdate = info.lastUpdate || new Date(); this.lastUpdate = info.lastUpdate || new Date();
this.newsFeedCache = info.newsFeedCache || []; this.newsFeedCache = info.newsFeedCache || [];
this.notifications = info.notifications || []; this.notifications = info.notifications || [];
this.bibleHighlights = info.bibleHighlights || [];
this.bibleNotes = info.bibleNotes || [];
//groupRelated //groupRelated
this.isGroup = info.isGroup || false; this.isGroup = info.isGroup || false;
+20 -1
View File
@@ -18,15 +18,32 @@ services:
- WEB_PUSH_EMAIL=${WEB_PUSH_EMAIL} - WEB_PUSH_EMAIL=${WEB_PUSH_EMAIL}
- EMAILPASS=${EMAILPASS} - EMAILPASS=${EMAILPASS}
- PORT=3001 - PORT=3001
- TRANSLATION_PROVIDER=${TRANSLATION_PROVIDER:-openai}
- MARIAN_TRANSLATION_URL=http://translation:8000
volumes: volumes:
- .:/app - .:/app
- '/app/node_modules' - '/app/node_modules'
#depends_on: #depends_on:
# - mongo # - mongo
command: node index.js command: node index.js
depends_on:
- translation
networks:
- emi-network
# networks: # networks:
# - emi-network # - emi-network
translation:
build:
context: ./translation-service
restart: unless-stopped
environment:
- MARIAN_HOST=0.0.0.0
volumes:
- translation-models:/models
networks:
- emi-network
#mongo: #mongo:
# image: mongo:latest # image: mongo:latest
# ports: # ports:
@@ -38,6 +55,8 @@ services:
# - ./dump:/dump # - ./dump:/dump
#entrypoint: mongodump ${MONGO_URL} && mongorestore --db EMI_SOCIAL dump/EMI_SOCIAL/ && mongod #entrypoint: mongodump ${MONGO_URL} && mongorestore --db EMI_SOCIAL dump/EMI_SOCIAL/ && mongod
#volumes: volumes:
translation-models:
driver: local
#mongodbdata: #mongodbdata:
# driver: local # This ensures the volume is created # driver: local # This ensures the volume is created
+1
View File
@@ -29,6 +29,7 @@ const limiter = rateLimit({
limit: 500, // Limit each IP to 100 requests per `window` (here, per 15 minutes). limit: 500, // Limit each IP to 100 requests per `window` (here, per 15 minutes).
standardHeaders: 'draft-8', // draft-6: `RateLimit-*` headers; draft-7 & draft-8: combined `RateLimit` header standardHeaders: 'draft-8', // draft-6: `RateLimit-*` headers; draft-7 & draft-8: combined `RateLimit` header
legacyHeaders: false, // Disable the `X-RateLimit-*` headers. legacyHeaders: false, // Disable the `X-RateLimit-*` headers.
skip: (req) => req.path.startsWith("/live-captions"),
keyGenerator: (req) => { keyGenerator: (req) => {
const forwarded = req.headers["x-forwarded-for"]?.split(",")[0]; // Take the first IP in the list const forwarded = req.headers["x-forwarded-for"]?.split(",")[0]; // Take the first IP in the list
const ip = forwarded || req.ip; // Fallback to req.ip const ip = forwarded || req.ip; // Fallback to req.ip
+2
View File
@@ -10,6 +10,7 @@ const profileDB = require("./dbTools/profile.js");
const paymentDB = require("./dbTools/payments.js"); const paymentDB = require("./dbTools/payments.js");
const songsDB = require("./dbTools/songs.js"); const songsDB = require("./dbTools/songs.js");
const chatDB = require("./dbTools/chat.js"); const chatDB = require("./dbTools/chat.js");
const bibleDB = require("./dbTools/bible.js");
console.log("Connecting to MongoDB..."); console.log("Connecting to MongoDB...");
const nodeMajorVersion = parseInt((process.versions.node || "0").split(".")[0], 10); const nodeMajorVersion = parseInt((process.versions.node || "0").split(".")[0], 10);
@@ -179,6 +180,7 @@ const getDB = new Promise((resolve, reject) => {
paymentDB(DB); paymentDB(DB);
songsDB(DB); songsDB(DB);
chatDB(DB); chatDB(DB);
bibleDB(DB);
resolve(DB); resolve(DB);
}); });
+1
View File
@@ -6,6 +6,7 @@
"scripts": { "scripts": {
"test": "npx mocha test/auth.test.js", "test": "npx mocha test/auth.test.js",
"start": "node index.js", "start": "node index.js",
"dev": "node --watch index.js",
"live-captions:test-sender": "node scripts/liveCaptionsTestSender.js", "live-captions:test-sender": "node scripts/liveCaptionsTestSender.js",
"docker": "docker compose up -d", "docker": "docker compose up -d",
"docker_restore": "docker-compose exec mongo mongorestore --db EMI_SOCIAL /dump/EMI_SOCIAL/", "docker_restore": "docker-compose exec mongo mongorestore --db EMI_SOCIAL /dump/EMI_SOCIAL/",
+311
View File
@@ -2,6 +2,7 @@ const axios = require('axios');
var express = require('express') var express = require('express')
var router = express.Router() var router = express.Router()
const DB = require("./../mongoDB.js"); const DB = require("./../mongoDB.js");
const { getProfileId } = require("./../utils/sessionUtils.js");
const fetchAPI = async (path) => { const fetchAPI = async (path) => {
baseUrl = "https://api.scripture.api.bible/v1/" baseUrl = "https://api.scripture.api.bible/v1/"
@@ -11,6 +12,20 @@ const fetchAPI = async (path) => {
const defaultBibleId = "592420522e16049f-01"; const defaultBibleId = "592420522e16049f-01";
const normalizeVerseKey = (verseKey) => {
if (!verseKey || typeof verseKey !== "string") return "";
return verseKey.trim().toUpperCase();
};
const isValidVerseKey = (verseKey) => {
return /^[A-Z0-9]+(\.[A-Z0-9]+)+$/.test(verseKey);
};
const getNoteText = (req) => {
const text = req.body?.text || req.body?.note || "";
return typeof text === "string" ? text.trim() : "";
};
//getMedia('y42zyf3').then(console.log) //getMedia('y42zyf3').then(console.log)
DB.getDB.then((DB) => { DB.getDB.then((DB) => {
/** /**
@@ -37,6 +52,302 @@ DB.getDB.then((DB) => {
return res.json(bibles); return res.json(bibles);
}); });
/**
* @swagger
* /bible/mine:
* get:
* summary: Get Bible highlights and notes for the active profile
* tags: [Bible]
* security:
* - cookieAuth: []
* responses:
* 200:
* description: OK
*/
router.get("/mine", async (req, res) => {
const profileid = getProfileId(req);
const bibleData = await DB.getBibleProfileData(profileid);
if (!bibleData) return res.status(400).json({ status: "Invalid profile" });
return res.json({
status: "ok",
...bibleData
});
});
/**
* @swagger
* /bible/verses/{verseKey}/counters:
* get:
* summary: Get highlight and note counters for a Bible verse
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* responses:
* 200:
* description: OK
*/
router.get("/verses/:verseKey/counters", async (req, res) => {
const verseKey = normalizeVerseKey(req.params.verseKey);
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
const counters = await DB.getBibleVerseCounters(verseKey);
return res.json({
status: "ok",
...counters
});
});
/**
* @swagger
* /bible/verses/{verseKey}/activity:
* get:
* summary: Get profiles and recent notes for a Bible verse
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* - in: query
* name: limit
* schema:
* type: integer
* default: 10
* responses:
* 200:
* description: OK
*/
router.get("/verses/:verseKey/activity", async (req, res) => {
const verseKey = normalizeVerseKey(req.params.verseKey);
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
const limit = Math.min(Math.max(parseInt(req.query.limit, 10) || 10, 1), 10);
const activity = await DB.getBibleVerseActivity(verseKey, limit);
return res.json({
status: "ok",
verseKey,
...activity
});
});
/**
* @swagger
* /bible/verses/{verseKey}/highlight:
* post:
* summary: Highlight a Bible verse for the active profile
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* responses:
* 200:
* description: OK
*/
router.post("/verses/:verseKey/highlight", async (req, res) => {
const profileid = getProfileId(req);
const verseKey = normalizeVerseKey(req.params.verseKey);
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
const result = await DB.addBibleHighlight(profileid, verseKey);
if (!result) return res.status(400).json({ status: "Could not add highlight" });
const counters = await DB.getBibleVerseCounters(verseKey);
return res.json({
status: "ok",
highlighted: true,
...counters
});
});
/**
* @swagger
* /bible/verses/{verseKey}/highlight:
* delete:
* summary: Remove a Bible verse highlight for the active profile
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* responses:
* 200:
* description: OK
*/
router.delete("/verses/:verseKey/highlight", async (req, res) => {
const profileid = getProfileId(req);
const verseKey = normalizeVerseKey(req.params.verseKey);
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
const result = await DB.removeBibleHighlight(profileid, verseKey);
if (!result) return res.status(400).json({ status: "Could not remove highlight" });
const counters = await DB.getBibleVerseCounters(verseKey);
return res.json({
status: "ok",
highlighted: false,
...counters
});
});
/**
* @swagger
* /bible/verses/{verseKey}/notes:
* post:
* summary: Add a Bible verse note for the active profile
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* responses:
* 200:
* description: OK
*/
router.post("/verses/:verseKey/notes", async (req, res) => {
const profileid = getProfileId(req);
const verseKey = normalizeVerseKey(req.params.verseKey);
const text = getNoteText(req);
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
if (!text) return res.status(400).json({ status: "Note text is required" });
const note = await DB.addBibleNote(profileid, verseKey, text);
if (!note) return res.status(400).json({ status: "Could not add note" });
const counters = await DB.getBibleVerseCounters(verseKey);
return res.json({
status: "ok",
verseKey,
note,
counters
});
});
/**
* @swagger
* /bible/verses/{verseKey}/notes/{noteId}:
* put:
* summary: Update a Bible verse note for the active profile
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* - in: path
* name: noteId
* required: true
* schema:
* type: string
* responses:
* 200:
* description: OK
*/
router.put("/verses/:verseKey/notes/:noteId", async (req, res) => {
const profileid = getProfileId(req);
const verseKey = normalizeVerseKey(req.params.verseKey);
const noteId = req.params.noteId;
const text = getNoteText(req);
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
if (!text) return res.status(400).json({ status: "Note text is required" });
const note = await DB.updateBibleNote(profileid, verseKey, noteId, text);
if (!note) return res.status(404).json({ status: "Note not found" });
return res.json({
status: "ok",
verseKey,
note
});
});
/**
* @swagger
* /bible/verses/{verseKey}/notes/{noteId}:
* delete:
* summary: Delete a Bible verse note for the active profile
* tags: [Bible]
* security:
* - cookieAuth: []
* parameters:
* - in: path
* name: verseKey
* required: true
* schema:
* type: string
* example: GEN.1.1
* - in: path
* name: noteId
* required: true
* schema:
* type: string
* responses:
* 200:
* description: OK
*/
router.delete("/verses/:verseKey/notes/:noteId", async (req, res) => {
const profileid = getProfileId(req);
const verseKey = normalizeVerseKey(req.params.verseKey);
const noteId = req.params.noteId;
if (!isValidVerseKey(verseKey)) {
return res.status(400).json({ status: "Invalid verse key" });
}
const result = await DB.removeBibleNote(profileid, verseKey, noteId);
if (!result) return res.status(400).json({ status: "Could not delete note" });
if (!result.modifiedCount) return res.status(404).json({ status: "Note not found" });
const counters = await DB.getBibleVerseCounters(verseKey);
return res.json({
status: "ok",
verseKey,
deleted: true,
counters
});
});
/** /**
* @swagger * @swagger
* /bible/books: * /bible/books:
+99 -10
View File
@@ -1,19 +1,36 @@
var express = require('express'); var express = require('express');
var router = express.Router(); var router = express.Router();
const { rateLimit } = require("express-rate-limit");
const sessionChecker = require("../middleware/sessionChecker.js"); const sessionChecker = require("../middleware/sessionChecker.js");
const MAX_BUFFER_SIZE = 300; const MAX_BUFFER_SIZE = 300;
const DEFAULT_INITIAL_LIMIT = 40; const DEFAULT_INITIAL_LIMIT = 40;
const MAX_INITIAL_LIMIT = 120; const MAX_INITIAL_LIMIT = 120;
const CAPTION_META_KEYS = new Set(["sequence", "createdAt", "original"]); const INACTIVITY_RESET_MS = 10 * 60 * 1000;
const CAPTION_META_KEYS = new Set(["sequence", "createdAt", "original", "draft", "sourceLang", "lang", "isDraft", "status", "translations"]);
const liveCaptionState = { const liveCaptionState = {
startedAt: Date.now(), startedAt: Date.now(),
lastIngestAt: 0,
latestSequence: 0, latestSequence: 0,
captions: [], captions: [],
}; };
const liveCaptionsLimiter = rateLimit({
windowMs: 10 * 60 * 1000,
limit: 6000,
standardHeaders: "draft-8",
legacyHeaders: false,
keyGenerator: (req) => {
const forwarded = req.headers["x-forwarded-for"]?.split(",")[0];
const ip = forwarded || req.ip || "";
return ip.includes(":") ? ip.split(":")[0] : ip;
},
});
router.use(liveCaptionsLimiter);
const normalizeLang = (lang = "") => { const normalizeLang = (lang = "") => {
const value = String(lang || "").trim().toLowerCase(); const value = String(lang || "").trim().toLowerCase();
if (!value) return ""; if (!value) return "";
@@ -33,8 +50,23 @@ const normalizeTranslations = (translations) => {
return normalized; return normalized;
}; };
const readText = (value) => {
if (typeof value === "string") return value.trim();
return "";
};
const extractDraftText = (body = {}) => {
const directDraft = readText(body?.draft);
if (directDraft) return directDraft;
const nestedDraft = readText(body?.draft?.text);
if (nestedDraft) return nestedDraft;
const fallbackText = readText(body?.text);
if (fallbackText) return fallbackText;
return "";
};
const buildTranslationsFromFlatPayload = (payload) => { const buildTranslationsFromFlatPayload = (payload) => {
const ignoredKeys = new Set(["original", "sourceLang", "translations"]); const ignoredKeys = new Set(["original", "draft", "sourceLang", "lang", "isDraft", "status", "translations"]);
const normalized = {}; const normalized = {};
for (const [key, value] of Object.entries(payload || {})) { for (const [key, value] of Object.entries(payload || {})) {
if (ignoredKeys.has(key)) continue; if (ignoredKeys.has(key)) continue;
@@ -67,8 +99,22 @@ const getAvailableLanguages = () => {
return Array.from(langs).filter(Boolean).sort(); return Array.from(langs).filter(Boolean).sort();
}; };
router.get("/stream", sessionChecker, async (req, res) => { const resetLiveCaptionState = () => {
liveCaptionState.startedAt = Date.now();
liveCaptionState.lastIngestAt = 0;
liveCaptionState.latestSequence = 0;
liveCaptionState.captions = [];
};
const maybeResetForInactivity = () => {
if (!liveCaptionState.lastIngestAt) return;
if ((Date.now() - liveCaptionState.lastIngestAt) < INACTIVITY_RESET_MS) return;
resetLiveCaptionState();
};
router.get("/stream", async (req, res) => {
try { try {
maybeResetForInactivity();
const sinceSequence = Number.parseInt(req.query?.sinceSequence, 10); const sinceSequence = Number.parseInt(req.query?.sinceSequence, 10);
const requestedLimit = Number.parseInt(req.query?.limit, 10); const requestedLimit = Number.parseInt(req.query?.limit, 10);
const initialLimit = Number.isFinite(requestedLimit) const initialLimit = Number.isFinite(requestedLimit)
@@ -103,16 +149,22 @@ router.get("/stream", sessionChecker, async (req, res) => {
router.post("/ingest", async (req, res) => { router.post("/ingest", async (req, res) => {
try { try {
// TODO: Add basic auth/API key validation before production roll-out. // TODO: Add basic auth/API key validation before production roll-out.
const original = typeof req.body?.original === "string" ? req.body.original.trim() : ""; const draft = extractDraftText(req.body || {});
const originalFromPayload = readText(req.body?.original);
const original = originalFromPayload || draft;
const requestedLang = normalizeLang(req.body?.lang);
const sourceLangFromRequest = normalizeLang(req.body?.sourceLang || (requestedLang && requestedLang !== "draft" ? requestedLang : ""));
const isDraft = !!draft || requestedLang === "draft" || sourceLangFromRequest === "draft" || req.body?.isDraft === true || req.body?.status === "draft";
const mapFromNested = normalizeTranslations(req.body?.translations); const mapFromNested = normalizeTranslations(req.body?.translations);
const mapFromFlat = buildTranslationsFromFlatPayload(req.body); const mapFromFlat = buildTranslationsFromFlatPayload(req.body);
const translations = { ...mapFromNested, ...mapFromFlat }; const translations = isDraft ? {} : { ...mapFromNested, ...mapFromFlat };
const sourceLang = normalizeLang(req.body?.sourceLang || inferSourceLangFromTranslations(original, translations)); const inferredSource = inferSourceLangFromTranslations(original, translations);
const sourceLang = isDraft ? "" : (sourceLangFromRequest || inferredSource);
if (!original) { if (!original) {
return res.status(400).json({ status: "Original text is required" }); return res.status(400).json({ status: "Original text is required" });
} }
if (sourceLang && sourceLang !== "original" && !translations[sourceLang]) { if (sourceLang && sourceLang !== "original" && sourceLang !== "draft" && !translations[sourceLang]) {
translations[sourceLang] = original; translations[sourceLang] = original;
} }
@@ -121,10 +173,15 @@ router.post("/ingest", async (req, res) => {
sequence, sequence,
createdAt: new Date().toISOString(), createdAt: new Date().toISOString(),
original, original,
sourceLang: sourceLang || undefined,
lang: isDraft ? "draft" : (sourceLang || undefined),
isDraft,
status: isDraft ? "draft" : "final",
...translations, ...translations,
}; };
liveCaptionState.latestSequence = sequence; liveCaptionState.latestSequence = sequence;
liveCaptionState.lastIngestAt = Date.now();
liveCaptionState.captions.push(caption); liveCaptionState.captions.push(caption);
if (liveCaptionState.captions.length > MAX_BUFFER_SIZE) { if (liveCaptionState.captions.length > MAX_BUFFER_SIZE) {
liveCaptionState.captions.splice(0, liveCaptionState.captions.length - MAX_BUFFER_SIZE); liveCaptionState.captions.splice(0, liveCaptionState.captions.length - MAX_BUFFER_SIZE);
@@ -145,9 +202,7 @@ router.post("/ingest", async (req, res) => {
router.post("/reset", async (_, res) => { router.post("/reset", async (_, res) => {
try { try {
// TODO: Add admin authorization before exposing this endpoint. // TODO: Add admin authorization before exposing this endpoint.
liveCaptionState.startedAt = Date.now(); resetLiveCaptionState();
liveCaptionState.latestSequence = 0;
liveCaptionState.captions = [];
return res.json({ status: "ok" }); return res.json({ status: "ok" });
} catch (error) { } catch (error) {
console.error("Error resetting live captions state", error); console.error("Error resetting live captions state", error);
@@ -155,4 +210,38 @@ router.post("/reset", async (_, res) => {
} }
}); });
router.get("/public/stream", async (req, res) => {
try {
maybeResetForInactivity();
const sinceSequence = Number.parseInt(req.query?.sinceSequence, 10);
const requestedLimit = Number.parseInt(req.query?.limit, 10);
const initialLimit = Number.isFinite(requestedLimit)
? Math.max(1, Math.min(requestedLimit, MAX_INITIAL_LIMIT))
: DEFAULT_INITIAL_LIMIT;
let captions = [];
if (Number.isFinite(sinceSequence) && sinceSequence >= 0) {
captions = liveCaptionState.captions.filter((item) => item.sequence > sinceSequence);
} else {
captions = liveCaptionState.captions.slice(-initialLimit);
}
return res.json({
status: "ok",
latestSequence: liveCaptionState.latestSequence,
startedAt: new Date(liveCaptionState.startedAt).toISOString(),
availableLanguages: getAvailableLanguages(),
captions,
});
} catch (error) {
console.error("Error getting public live captions stream", error);
return res.status(500).json({
status: "Internal server error",
latestSequence: liveCaptionState.latestSequence,
captions: [],
availableLanguages: [],
});
}
});
module.exports = router; module.exports = router;
+20 -5
View File
@@ -5,7 +5,7 @@ const axios = require("axios");
const baseUrl = (process.env.CAPTION_TEST_BASE_URL || process.env.BASE_URL || "http://localhost:3000").replace(/\/+$/, ""); const baseUrl = (process.env.CAPTION_TEST_BASE_URL || process.env.BASE_URL || "http://localhost:3000").replace(/\/+$/, "");
const ingestUrl = `${baseUrl}/live-captions/ingest`; const ingestUrl = `${baseUrl}/live-captions/ingest`;
const intervalMs = 5000; const intervalMs = 6000;
const samples = [ const samples = [
{ {
@@ -37,9 +37,8 @@ const samples = [
let sampleIndex = 0; let sampleIndex = 0;
let timer = null; let timer = null;
const sendNextSample = async () => { const postPayload = async (payload) => {
const payload = samples[sampleIndex]; const kind = payload?.draft ? "draft" : "final";
sampleIndex = (sampleIndex + 1) % samples.length;
try { try {
const response = await axios.post(ingestUrl, payload, { const response = await axios.post(ingestUrl, payload, {
@@ -47,7 +46,8 @@ const sendNextSample = async () => {
timeout: 10000, timeout: 10000,
}); });
const seq = response?.data?.caption?.sequence || response?.data?.latestSequence || "?"; const seq = response?.data?.caption?.sequence || response?.data?.latestSequence || "?";
console.log(`[live-captions:test-sender] sent sequence=${seq} original="${payload.original}"`); const text = payload?.draft || payload?.original || "";
console.log(`[live-captions:test-sender] sent ${kind} sequence=${seq} text="${text}"`);
} catch (error) { } catch (error) {
const status = error?.response?.status; const status = error?.response?.status;
const body = error?.response?.data; const body = error?.response?.data;
@@ -56,6 +56,21 @@ const sendNextSample = async () => {
} }
}; };
const sendNextSample = async () => {
const payload = samples[sampleIndex];
const draftWords = String(payload?.original || "").split(" ").filter(Boolean);
if (draftWords.length > 2) {
await postPayload({ draft: draftWords.slice(0, 2).join(" ") });
await new Promise((resolve) => setTimeout(resolve, 550));
await postPayload({ draft: draftWords.slice(0, 4).join(" ") });
await new Promise((resolve) => setTimeout(resolve, 550));
}
await postPayload(payload);
sampleIndex = (sampleIndex + 1) % samples.length;
};
const start = async () => { const start = async () => {
console.log(`[live-captions:test-sender] posting to ${ingestUrl} every ${intervalMs / 1000}s`); console.log(`[live-captions:test-sender] posting to ${ingestUrl} every ${intervalMs / 1000}s`);
await sendNextSample(); await sendNextSample();
+3
View File
@@ -0,0 +1,3 @@
.venv/
__pycache__/
*.pyc
+13
View File
@@ -0,0 +1,13 @@
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt ./
RUN pip install --no-cache-dir -r requirements.txt
COPY server.py ./
ENV MARIAN_HOST=0.0.0.0
ENV TRANSFORMERS_CACHE=/models
VOLUME ["/models"]
EXPOSE 8000
CMD ["python", "server.py"]
+4
View File
@@ -0,0 +1,4 @@
torch>=2.2,<3
transformers>=4.40,<5
sentencepiece>=0.2,<1
langdetect>=1.0.9,<2
+118
View File
@@ -0,0 +1,118 @@
import json
import os
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from langdetect import DetectorFactory, LangDetectException, detect
from transformers import MarianMTModel, MarianTokenizer
DetectorFactory.seed = 0
HOST = os.getenv("MARIAN_HOST", "127.0.0.1")
PORT = int(os.getenv("MARIAN_PORT", "8000"))
MAX_INPUT_LENGTH = int(os.getenv("MARIAN_MAX_INPUT_LENGTH", "1000"))
DEFAULT_SOURCE_LANGUAGE = os.getenv("MARIAN_DEFAULT_SOURCE_LANGUAGE", "en")
SUPPORTED_LANGUAGES = {"en", "es", "fr", "da", "ar"}
MODEL_CACHE = {}
def normalize_language(value):
language = str(value or "").strip().lower().split(",")[0].split("-")[0]
return language
def detect_source_language(text):
try:
detected = normalize_language(detect(text))
if detected in SUPPORTED_LANGUAGES:
return detected
except LangDetectException:
pass
return DEFAULT_SOURCE_LANGUAGE
def get_model(source, target):
model_name = f"Helsinki-NLP/opus-mt-{source}-{target}"
if model_name not in MODEL_CACHE:
MODEL_CACHE[model_name] = (
MarianTokenizer.from_pretrained(model_name),
MarianMTModel.from_pretrained(model_name),
)
return model_name, MODEL_CACHE[model_name]
def translate_once(text, source, target):
model_name, (tokenizer, model) = get_model(source, target)
encoded = tokenizer([text], return_tensors="pt", truncation=True)
generated = model.generate(**encoded)
return tokenizer.batch_decode(generated, skip_special_tokens=True)[0], model_name
def translate(text, source, target):
if source == "auto":
source = detect_source_language(text)
if source not in SUPPORTED_LANGUAGES or target not in SUPPORTED_LANGUAGES:
raise ValueError("Only en, es, fr, da, and ar are supported")
if source == target:
return text, source, "none"
if source == "en" or target == "en":
translated, model_name = translate_once(text, source, target)
return translated, source, model_name
english, first_model = translate_once(text, source, "en")
translated, second_model = translate_once(english, "en", target)
return translated, source, f"{first_model},{second_model}"
class TranslationHandler(BaseHTTPRequestHandler):
def send_json(self, status, body):
payload = json.dumps(body).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(payload)))
self.end_headers()
self.wfile.write(payload)
def do_GET(self):
if self.path != "/health":
self.send_json(404, {"status": "not found"})
return
self.send_json(200, {"status": "ok", "provider": "marianmt", "loadedModels": list(MODEL_CACHE)})
def do_POST(self):
if self.path != "/translate":
self.send_json(404, {"status": "not found"})
return
try:
content_length = int(self.headers.get("Content-Length", "0"))
body = json.loads(self.rfile.read(content_length).decode("utf-8"))
text = str(body.get("text") or "").strip()
source = normalize_language(body.get("sourceLang")) or "auto"
target = normalize_language(body.get("targetLang"))
if not text or not target:
self.send_json(400, {"status": "text and targetLang are required"})
return
if len(text) > MAX_INPUT_LENGTH:
self.send_json(400, {"status": f"text exceeds {MAX_INPUT_LENGTH} characters"})
return
translated, detected_source, model_name = translate(text, source, target)
self.send_json(200, {
"status": "ok",
"translatedText": translated,
"sourceLang": detected_source,
"targetLang": target,
"provider": "marianmt",
"model": model_name,
})
except (ValueError, json.JSONDecodeError) as error:
self.send_json(400, {"status": str(error)})
except Exception as error:
print(f"Translation failed: {error}", flush=True)
self.send_json(502, {"status": "Translation failed"})
def log_message(self, format_string, *args):
print(f"[marianmt] {self.address_string()} {format_string % args}", flush=True)
if __name__ == "__main__":
print(f"MarianMT translation service listening on {HOST}:{PORT}", flush=True)
ThreadingHTTPServer((HOST, PORT), TranslationHandler).serve_forever()
+27
View File
@@ -1,6 +1,8 @@
const axios = require("axios"); const axios = require("axios");
const DEFAULT_MODEL = process.env.OPENAI_TRANSLATION_MODEL || process.env.OPENAI_MODEL || "gpt-4o-mini"; const DEFAULT_MODEL = process.env.OPENAI_TRANSLATION_MODEL || process.env.OPENAI_MODEL || "gpt-4o-mini";
const TRANSLATION_PROVIDER = (process.env.TRANSLATION_PROVIDER || "openai").trim().toLowerCase();
const MARIAN_TRANSLATION_URL = (process.env.MARIAN_TRANSLATION_URL || "http://127.0.0.1:8000").replace(/\/$/, "");
const normalizeLanguageCode = (rawLanguage) => { const normalizeLanguageCode = (rawLanguage) => {
if (!rawLanguage || typeof rawLanguage !== "string") return "en"; if (!rawLanguage || typeof rawLanguage !== "string") return "en";
@@ -40,6 +42,31 @@ const translateText = async ({ text, sourceLang, targetLang }) => {
}; };
} }
if (TRANSLATION_PROVIDER === "marian") {
try {
const response = await axios.post(
`${MARIAN_TRANSLATION_URL}/translate`,
{ text, sourceLang: sourceLang || "auto", targetLang: normalizedTarget },
{ timeout: 30000, headers: { "Content-Type": "application/json" } }
);
const translatedText = response?.data?.translatedText?.trim();
if (!translatedText) return null;
return {
translatedText,
provider: response.data.provider || "marianmt",
model: response.data.model || "unknown",
};
} catch (error) {
console.error("Error translating with MarianMT", error?.response?.data || error?.message || error);
return null;
}
}
if (TRANSLATION_PROVIDER !== "openai") {
console.error(`Unsupported translation provider: ${TRANSLATION_PROVIDER}`);
return null;
}
const apiKey = process.env.OPENAI_API_KEY; const apiKey = process.env.OPENAI_API_KEY;
if (!apiKey) return null; if (!apiKey) return null;