diff --git a/apps/cf-ai-backend/src/env.d.ts b/apps/cf-ai-backend/src/env.d.ts index 4d6f675e..e760cba3 100644 --- a/apps/cf-ai-backend/src/env.d.ts +++ b/apps/cf-ai-backend/src/env.d.ts @@ -4,7 +4,7 @@ interface Env { SECURITY_KEY: string; OPENAI_API_KEY: string; GOOGLE_AI_API_KEY: string; - MY_QUEUE: Queue; + MY_QUEUE: Queue; KV: KVNamespace; } @@ -14,4 +14,5 @@ interface TweetData { authorName: string; handle: string; time: string; + saveToUser: string; } diff --git a/apps/cf-ai-backend/src/index.ts b/apps/cf-ai-backend/src/index.ts index ccaf06bd..317203c4 100644 --- a/apps/cf-ai-backend/src/index.ts +++ b/apps/cf-ai-backend/src/index.ts @@ -4,6 +4,7 @@ import { CloudflareVectorizeStore } from '@langchain/cloudflare'; import { OpenAIEmbeddings } from './OpenAIEmbedder'; import { GoogleGenerativeAI } from '@google/generative-ai'; import routeMap from './routes'; +import { queue } from './routes/queue'; function isAuthorized(request: Request, env: Env): boolean { return request.headers.get('X-Custom-Auth-Key') === env.SECURITY_KEY; @@ -45,4 +46,5 @@ export default { } return await handler(request, store, embeddings, model, env, ctx); }, + queue, }; diff --git a/apps/cf-ai-backend/src/routes.ts b/apps/cf-ai-backend/src/routes.ts index a8a3249a..0349c397 100644 --- a/apps/cf-ai-backend/src/routes.ts +++ b/apps/cf-ai-backend/src/routes.ts @@ -3,6 +3,7 @@ import * as apiAdd from './routes/add'; import * as apiQuery from './routes/query'; import * as apiAsk from './routes/ask'; import * as apiChat from './routes/chat'; +import * as apiBatchUploadTweets from './routes/batchUploadTweets'; import { OpenAIEmbeddings } from './OpenAIEmbedder'; import { GenerativeModel } from '@google/generative-ai'; import { Request } from '@cloudflare/workers-types'; @@ -26,6 +27,8 @@ routeMap.set('/ask', apiAsk); routeMap.set('/chat', apiChat); +routeMap.set('/batchUploadTweets', apiBatchUploadTweets); + // Add more route mappings as needed // routeMap.set('/api/otherRoute', { ... }); diff --git a/apps/cf-ai-backend/src/routes/add.ts b/apps/cf-ai-backend/src/routes/add.ts index b7fe073a..940287e3 100644 --- a/apps/cf-ai-backend/src/routes/add.ts +++ b/apps/cf-ai-backend/src/routes/add.ts @@ -2,6 +2,7 @@ import { Request } from '@cloudflare/workers-types'; import { type CloudflareVectorizeStore } from '@langchain/cloudflare'; import { OpenAIEmbeddings } from '../OpenAIEmbedder'; import { GenerativeModel } from '@google/generative-ai'; +import { seededRandom } from '../util'; export async function POST(request: Request, store: CloudflareVectorizeStore, _: OpenAIEmbeddings, m: GenerativeModel, env: Env) { const body = (await request.json()) as { @@ -20,15 +21,6 @@ export async function POST(request: Request, store: CloudflareVectorizeStore, _: const ourID = `${body.url}-${body.user}`; - // WHY? Because this helps us to prevent duplicate entries for the same URL and user - function seededRandom(seed: string) { - let x = [...seed].reduce((acc, cur) => acc + cur.charCodeAt(0), 0); - return () => { - x = (x * 9301 + 49297) % 233280; - return x / 233280; - }; - } - const random = seededRandom(ourID); const uuid = random().toString(36).substring(2, 15) + random().toString(36).substring(2, 15); diff --git a/apps/cf-ai-backend/src/routes/batchUploadTweets.ts b/apps/cf-ai-backend/src/routes/batchUploadTweets.ts new file mode 100644 index 00000000..0370436b --- /dev/null +++ b/apps/cf-ai-backend/src/routes/batchUploadTweets.ts @@ -0,0 +1,38 @@ +import { Request } from '@cloudflare/workers-types'; +import { type CloudflareVectorizeStore } from '@langchain/cloudflare'; +import { OpenAIEmbeddings } from '../OpenAIEmbedder'; +import { GenerativeModel } from '@google/generative-ai'; + +export async function POST(request: Request, store: CloudflareVectorizeStore, _: OpenAIEmbeddings, m: GenerativeModel, env: Env) { + const body = (await request.json()) as TweetData[] | undefined; + + if (!body) { + return new Response(JSON.stringify({ message: 'Body is missing' }), { status: 400 }); + } + + const bytes = new TextEncoder().encode(JSON.stringify(body)).length; + + if (bytes < 128000) { + await env.MY_QUEUE.send(body); + } else { + let bytesTillNow = 0; + let batches: TweetData[] = []; + + const getByteLength = (data: string) => new TextEncoder().encode(data).length; + + for (let i = 0; i < body.length; i++) { + const byteLength = getByteLength(JSON.stringify(body[i])); + + if (bytesTillNow + byteLength < 100000) { + bytesTillNow += byteLength; + batches.push(body[i]); + } else { + await env.MY_QUEUE.send(batches); + batches = [body[i]]; + bytesTillNow = byteLength; + } + } + } + + return new Response(JSON.stringify({ message: 'Document Added' }), { status: 200 }); +} diff --git a/apps/cf-ai-backend/src/routes/queue.ts b/apps/cf-ai-backend/src/routes/queue.ts new file mode 100644 index 00000000..c00afb5b --- /dev/null +++ b/apps/cf-ai-backend/src/routes/queue.ts @@ -0,0 +1,96 @@ +import { CloudflareVectorizeStore } from '@langchain/cloudflare'; +import { OpenAIEmbeddings } from '../OpenAIEmbedder'; +import { seededRandom } from '../util'; + +export const queue = async (batch: MessageBatch, env: Env): Promise => { + const messages = batch.messages[0].body as TweetData[]; + + const token = messages[0].saveToUser; + + if (!token) { + return; + } + + const limits = (await fetch('https://supermemory.dhr.wtf/api/getCount', { + headers: { + Authorization: `Bearer ${token}`, + }, + }).then((res) => res.json())) as { limit: number; tweetsCount: number; user: string }; + + if (messages.length > limits.limit - limits.tweetsCount) { + messages.splice(limits.limit - limits.tweetsCount); + } + + if (messages.length === 0) { + return; + } + + const embeddings = new OpenAIEmbeddings({ + apiKey: env.OPENAI_API_KEY, + modelName: 'text-embedding-3-small', + }); + + const store = new CloudflareVectorizeStore(embeddings, { + index: env.VECTORIZE_INDEX, + }); + + const collectedDocsUUIDs: { + document: { + pageContent: string; + metadata: { title: string; description: string; space: string; url: string; user: string }; + id: string; + }; + }[] = []; + + messages.forEach(async (message) => { + const ourID = `${message.postUrl}-${limits.user}`; + + const random = seededRandom(ourID); + const uuid = random().toString(36).substring(2, 15) + random().toString(36).substring(2, 15); + + await env.KV.put(uuid, ourID); + const pageContent = `This is a tweet from ${message.authorName}, it was posted on ${message.time}. The tweet reads: ${message.tweetText}`; + + collectedDocsUUIDs.push({ + document: { + pageContent, + metadata: { + title: 'Twitter Bookmark', + description: '', + space: 'Bookmarked Tweets', + url: message.postUrl, + user: limits.user, + }, + id: uuid, + }, + }); + }); + + console.log(collectedDocsUUIDs); + + await store.addDocuments( + collectedDocsUUIDs.map(({ document }) => document), + { + ids: collectedDocsUUIDs.map(({ document }) => document.id), + }, + ); + + console.log(token); + + const res = await fetch('https://supermemory.dhr.wtf/api/addTweetsToDb', { + method: 'POST', + headers: { + Authorization: `Bearer ${token}`, + }, + body: JSON.stringify(messages), + }); + + console.log(res.status, res.statusText); + + if (res.status !== 200) { + console.log(await res.json()); + console.error('Error adding tweets to db'); + } + + console.log(`consumed from our queue: ${messages}`); +}; diff --git a/apps/cf-ai-backend/src/util.ts b/apps/cf-ai-backend/src/util.ts new file mode 100644 index 00000000..81f7628c --- /dev/null +++ b/apps/cf-ai-backend/src/util.ts @@ -0,0 +1,7 @@ +export function seededRandom(seed: string) { + let x = [...seed].reduce((acc, cur) => acc + cur.charCodeAt(0), 0); + return () => { + x = (x * 9301 + 49297) % 233280; + return x / 233280; + }; +} diff --git a/apps/cf-ai-backend/wrangler.toml b/apps/cf-ai-backend/wrangler.toml index 6ef02af0..e1fc018a 100644 --- a/apps/cf-ai-backend/wrangler.toml +++ b/apps/cf-ai-backend/wrangler.toml @@ -13,6 +13,9 @@ binding = "AI" queue = "batch-vector-queue" binding = "MY_QUEUE" + [[queues.consumers]] + queue = "batch-vector-queue" + [[kv_namespaces]] binding = "KV" id = "37a90353da63401e84e20e71165531d0" diff --git a/apps/extension/src/App.tsx b/apps/extension/src/App.tsx index 89227432..c29d98a2 100644 --- a/apps/extension/src/App.tsx +++ b/apps/extension/src/App.tsx @@ -44,7 +44,7 @@ function App() { getUserData(); }, []); - // TODO: Implement getting bookmarks from API directly + // TODO: Implement getting bookmarks from Twitter API directly // const [status, setStatus] = useState(''); // const [bookmarks, setBookmarks] = useState([]); @@ -149,7 +149,7 @@ function App() { ); } -// TODO: Implement getting bookmarks from API directly +// TODO: Implement getting bookmarks from Twitter API directly // async function fetchAllBookmarks( // authorizationHeader: string, // csrfToken: string, diff --git a/apps/extension/src/SideBar.tsx b/apps/extension/src/SideBar.tsx index 38b43a30..b4f50207 100644 --- a/apps/extension/src/SideBar.tsx +++ b/apps/extension/src/SideBar.tsx @@ -32,12 +32,12 @@ function sendUrlToAPI(spaces: number[]) { } else { // const content = Entire page content, but cleaned up for the LLM. No ads, no scripts, no styles, just the text. if article, just the importnat info abou tit. const content = document.documentElement.innerText; - chrome.runtime.sendMessage({ type: "urlChange", content, url }); + chrome.runtime.sendMessage({ type: "urlChange", content, url, spaces }); } } -function SideBar() { - // TODO: Implement getting bookmarks from API directly +function SideBar({ jwt }: { jwt: string }) { + // TODO: Implement getting bookmarks from Twitter API directly // chrome.runtime.onMessage.addListener(function (request) { // if (request.action === 'showProgressIndicator') { // // TODO: SHOW PROGRESS INDICATOR @@ -56,12 +56,25 @@ function SideBar() { const [spaces, setSpaces] = useState(); const [selectedSpaces, setSelectedSpaces] = useState([]); + const [isImportingTweets, setIsImportingTweets] = useState(false); + + const [log, setLog] = useState([]); + interface TweetData { tweetText: string; postUrl: string; authorName: string; handle: string; time: string; + saveToUser: string; + } + + function sendBookmarkedTweetsToAPI(tweets: TweetData[], token: string) { + chrome.runtime.sendMessage({ + type: "sendBookmarkedTweets", + jwt: token, + tweets, + }); } const fetchSpaces = async () => { @@ -76,6 +89,9 @@ function SideBar() { const fetchBookmarks = () => { const tweets: TweetData[] = []; // Initialize an empty array to hold all tweet elements + setIsImportingTweets(true); + console.log("Importing tweets"); + const scrollInterval = 1000; const scrollStep = 5000; // Pixels to scroll on each step @@ -88,10 +104,12 @@ function SideBar() { if (currentTweetCount === previousTweetCount) { unchangedCount++; if (unchangedCount >= 2) { - // Stop if the count has not changed 5 times - console.log("Scraping complete"); - console.log("Total tweets scraped: ", tweets.length); - console.log("Downloading tweets as JSON..."); + setLog([ + ...log, + "Scraping complete", + `Total tweets scraped: ${tweets.length}`, + "Downloading tweets as JSON...", + ]); clearInterval(scrollToEndIntervalID); // Stop scrolling observer.disconnect(); // Stop observing DOM changes downloadTweetsAsJson(tweets); // Download the tweets list as a JSON file @@ -140,8 +158,10 @@ function SideBar() { tweetText, time: time ?? "", postUrl, + saveToUser: jwt, }); - console.log("Tweets capturados: ", tweets.length); + + setLog([...log, `Scraped tweet: ${tweets.length}`]); } }); } @@ -162,15 +182,43 @@ function SideBar() { observer.observe(document.body, { childList: true, subtree: true }); function downloadTweetsAsJson(tweetsArray: TweetData[]) { - const jsonData = JSON.stringify(tweetsArray); // Convert the array to JSON - - // TODO: SEND jsonData to server - console.log(jsonData); + setLog([...log, "Saving the tweets to our database..."]); + sendBookmarkedTweetsToAPI(tweetsArray, jwt); + setIsImportingTweets(false); } }; return ( <> + {isImportingTweets && ( +
+
+
+ + + +

Importing your tweets...

+
+ {log.map((message, index) => ( +

{message}

+ ))} +
+
+
+
+ )} +
{window.location.href.includes("twitter.com") || @@ -227,9 +275,6 @@ function SideBar() { >