import { client, xml } from "@xmpp/client"; import { OpenAI } from "openai"; import debug from "@xmpp/debug"; import { createSessionManager } from "./session.js"; import fs from "node:fs/promises"; import path from "node:path"; export const xmpp = client({ service: process.env.XMPP_SERVICE, domain: process.env.XMPP_DOMAIN, username: process.env.XMPP_USERNAME, password: process.env.XMPP_PASSWORD }); export const lmStudio = new OpenAI({ baseURL: process.env.API_BASE_URL, apiKey: process.env.API_KEY }); const sessions = createSessionManager(process.env.SESSION_MANAGER); debug(xmpp, false); xmpp.on("error", (err) => { console.error(err); }); xmpp.on("online", async (address) => { console.log("online as", address.toString()); await xmpp.send(xml("presence")); }); xmpp.on("stanza", async (stanza) => { if (!stanza.is("message") || stanza.attrs.type !== "chat") return; const body = stanza.getChildText("body"); if (!body) return; if (body.includes("@clr")) { sessions.flush(); console.log("@clr: All sessions cleared"); return; } const systemPromptReply = await readPromptFile(); const from = stanza.attrs.from; const sessionId = from.split("/")[0]; console.log(stanza.toString()); const history = sessions.getHistory(sessionId); await xmppSendTypingSignal(from); sessions.saveMessage(sessionId, { role: "user", content: body }); let reply = ""; let lastPlacedChunkIdx = 0; const responseStream = await lmStudio.chat.completions.create({ model: process.env.MODEL_ID, messages: [{ role: "system", content: systemPromptReply }, ...history.messages], stream: true, temperature: 0.85, top_p: 0.9, frequency_penalty: 0.2, presence_penalty: 0.1, reasoning_effort: "none", stream_options: { include_usage: true } }); for await (const chunk of responseStream) { const content = chunk.choices[0]?.delta?.content || ""; reply += content; const replyPart = reply.slice(lastPlacedChunkIdx); if (replyPart.length > 200) { lastPlacedChunkIdx = reply.length; await xmppSendMessage(from, replyPart); await xmppSendTypingSignal(from); } } await xmppSendMessage(from, reply.slice(lastPlacedChunkIdx)); if (reply) { sessions.saveMessage(sessionId, { role: "assistant", content: reply }); } await xmppSendActiveSignal(from); console.log(history); }); /** @param {string} to */ async function xmppSendTypingSignal(to) { await xmpp.send( xml( "message", { to, type: "chat" }, xml("composing", { xmlns: "http://jabber.org/protocol/chatstates" }) ) ); } /** * @param {string} to * @param {string} message */ async function xmppSendMessage(to, message) { await xmpp.send(xml("message", { type: "chat", to }, xml("body", {}, `${message}`))); } /** * @param {string} to */ async function xmppSendActiveSignal(to) { await xmpp.send( xml( "message", { type: "chat", to }, xml("active", { xmlns: "http://jabber.org/protocol/chatstates" }) ) ); } async function readPromptFile() { try { const sysPromptFilename = process.argv[2] || ""; const text = await fs.readFile(path.join(process.cwd(), sysPromptFilename), { encoding: "utf-8" }); return text; } catch (error) { console.error(`Error while reading system prompt file:`, error); return ""; } }