opengraph stuff

This commit is contained in:
2026-02-24 21:47:36 +00:00
parent 05f0d0e8f7
commit e36908cb0b
61 changed files with 5870 additions and 178 deletions
+113
View File
@@ -0,0 +1,113 @@
import { readdir } from 'node:fs/promises';
import { join } from 'node:path';
import { simpleParser } from 'mailparser';
import type { EmailSummary, EmailMessage } from 'types';
import { createRouter } from '../../create-router';
import { getUserEmailDir } from '@@/data-path';
type CacheEntry = {
summaries: EmailSummary[];
fileCount: number;
};
const cache = new Map<string, CacheEntry>();
const parseHeadersOnly = async (filePath: string, id: string): Promise<EmailSummary | null> => {
try {
const file = Bun.file(filePath);
const buffer = Buffer.from(await file.arrayBuffer());
const parsed = await simpleParser(buffer, { skipHtmlToText: true, skipTextToHtml: true, skipImageLinks: true });
const text = parsed.text ?? '';
const snippet = text.slice(0, 120).replace(/\s+/g, ' ').trim();
return {
id,
from: parsed.from?.text ?? '',
to: parsed.to ? (Array.isArray(parsed.to) ? parsed.to.map((a) => a.text).join(', ') : parsed.to.text) : '',
subject: parsed.subject ?? '(no subject)',
date: (parsed.date ?? new Date()).toISOString(),
snippet,
};
} catch {
return null;
}
};
export const emailRouter = createRouter();
emailRouter.get('/messages', async (ctx) => {
const email = ctx.get('user').email;
const dir = getUserEmailDir(email);
let filenames: string[];
try {
const entries = await readdir(dir);
filenames = entries.filter((f) => f.endsWith('.eml'));
} catch {
return ctx.json({ messages: [], total: 0 });
}
const cached = cache.get(email);
if (cached && cached.fileCount === filenames.length) {
const page = Number(ctx.req.query('page') ?? '1');
const limit = Number(ctx.req.query('limit') ?? '50');
const start = (page - 1) * limit;
return ctx.json({ messages: cached.summaries.slice(start, start + limit), total: cached.summaries.length });
}
const summaries: EmailSummary[] = [];
for (const filename of filenames) {
const id = filename.replace(/\.eml$/, '');
const summary = await parseHeadersOnly(join(dir, filename), id);
if (summary) summaries.push(summary);
}
summaries.sort((a, b) => new Date(b.date).getTime() - new Date(a.date).getTime());
cache.set(email, { summaries, fileCount: filenames.length });
const page = Number(ctx.req.query('page') ?? '1');
const limit = Number(ctx.req.query('limit') ?? '50');
const start = (page - 1) * limit;
return ctx.json({ messages: summaries.slice(start, start + limit), total: summaries.length });
});
emailRouter.get('/messages/:id', async (ctx) => {
const email = ctx.get('user').email;
const id = ctx.req.param('id');
const filePath = join(getUserEmailDir(email), `${id}.eml`);
const file = Bun.file(filePath);
if (!(await file.exists())) {
return ctx.text('Not found', 404);
}
try {
const buffer = Buffer.from(await file.arrayBuffer());
const parsed = await simpleParser(buffer);
const attachments = (parsed.attachments ?? []).map((a) => ({
filename: a.filename ?? 'unknown',
size: a.size,
contentType: a.contentType,
}));
const message: EmailMessage = {
id,
from: parsed.from?.text ?? '',
to: parsed.to ? (Array.isArray(parsed.to) ? parsed.to.map((a) => a.text).join(', ') : parsed.to.text) : '',
cc: parsed.cc ? (Array.isArray(parsed.cc) ? parsed.cc.map((a) => a.text).join(', ') : parsed.cc.text) : undefined,
subject: parsed.subject ?? '(no subject)',
date: (parsed.date ?? new Date()).toISOString(),
snippet: (parsed.text ?? '').slice(0, 120).replace(/\s+/g, ' ').trim(),
html: parsed.html || undefined,
text: parsed.text || undefined,
attachments,
};
return ctx.json(message);
} catch {
return ctx.text('Failed to parse email', 500);
}
});
+18 -2
View File
@@ -1,10 +1,11 @@
import { join, relative } from "path";
import { homedir } from "node:os";
import { readdirSync, existsSync, mkdirSync, writeFileSync, readFileSync } from "node:fs";
import type { Subprocess } from "bun";
import type { PiEvent, MessageCost } from "./types";
import { readApiKeys } from "../server-settings/pi-mono";
import { readSearxngConfig } from "../server-settings/searxng";
import { PI_CONFIG_DIR, DATA_PATH, getGlobalSkillsDir, getUserSkillsDir, getGlobalExtensionsDir, getUserExtensionsDir, getGlobalToolsDir, getUserToolsDir, getNativeResourcesDir, getGlobalResourcesDir } from "../../data-path";
import { PI_CONFIG_DIR, DATA_PATH, getHomeDir, getGlobalSkillsDir, getUserSkillsDir, getGlobalExtensionsDir, getUserExtensionsDir, getGlobalToolsDir, getUserToolsDir, getNativeResourcesDir, getGlobalResourcesDir } from "../../data-path";
import { ensureDockerContainer } from "../terminal/websocket";
import { logger } from "./logger";
import { parseFrontmatter } from "../skills/skills";
@@ -153,6 +154,14 @@ function buildResourcesEnv(): string {
return JSON.stringify(result);
}
function getGoogleConfigPath(): string {
return join(homedir(), '.config', 'officer.dev', 'google-oauth.json');
}
function getGoogleTokenPath(email: string): string {
return join(DATA_PATH, email, 'integrations', 'google.json');
}
type SandboxOptions = {
userId: number;
username: string;
@@ -203,12 +212,19 @@ export async function spawnPi(
const resourcesEnv = buildResourcesEnv();
const googleConfigHost = getGoogleConfigPath();
const googleTokenHost = join(DATA_PATH, sandbox.email, 'integrations');
const envFlags = [
'-e', `PI_CODING_AGENT_DIR=${containerPiConfig}`,
'-e', `HOME=${containerHome}`,
'-e', `OFFICER_USER_HOME=${containerHome}`,
'-e', `OFFICER_USER_ROOT=/officer/user`,
'-e', `PI_TOOLS_DIRS=/officer/tools:/officer/user/tools`,
'-e', `PI_SEARXNG_URL=${searxng.url}`,
'-e', `OFFICER_RESOURCES=${resourcesEnv}`,
'-e', `OFFICER_GOOGLE_CONFIG_PATH=/officer/google-oauth.json`,
'-e', `OFFICER_GOOGLE_TOKEN_PATH=/officer/user/integrations/google.json`,
];
for (const [key, value] of Object.entries(storedKeys)) {
if (value?.trim()) envFlags.push('-e', `${key}=${value.trim()}`);
@@ -258,7 +274,7 @@ export async function spawnPi(
stdin: 'pipe',
stdout: 'pipe',
stderr: 'pipe',
env: { ...process.env, ...storedKeys, PI_CODING_AGENT_DIR: PI_CONFIG_DIR, PI_TOOLS_DIRS: toolsDirs, PI_SEARXNG_URL: searxng.url, OFFICER_RESOURCES: buildResourcesEnv() },
env: { ...process.env, ...storedKeys, HOME: getHomeDir(email), OFFICER_USER_HOME: getHomeDir(email), OFFICER_USER_ROOT: join(DATA_PATH, email), PI_CODING_AGENT_DIR: PI_CONFIG_DIR, PI_TOOLS_DIRS: toolsDirs, PI_SEARXNG_URL: searxng.url, OFFICER_RESOURCES: buildResourcesEnv(), OFFICER_GOOGLE_CONFIG_PATH: getGoogleConfigPath(), OFFICER_GOOGLE_TOKEN_PATH: getGoogleTokenPath(email) },
});
logger.info('Spawned Pi locally', {
+42
View File
@@ -0,0 +1,42 @@
import { createRouter } from '../../create-router';
import { enqueue, cancelJob, readJob, listAllJobs } from '../../queue';
import { NOT_FOUND } from '../../custom-errors';
export const queueRouter = createRouter();
queueRouter.get('/jobs', async (ctx) => {
const user = ctx.get('user');
const lane = ctx.req.query('lane');
const type = ctx.req.query('type');
const status = ctx.req.query('status');
let jobs = await listAllJobs();
jobs = jobs.filter((j) => j.userId === user.email);
if (lane) jobs = jobs.filter((j) => j.lane === lane);
if (type) jobs = jobs.filter((j) => j.type === type);
if (status) jobs = jobs.filter((j) => j.status === status);
return ctx.json(jobs);
});
queueRouter.get('/jobs/:id', async (ctx) => {
const job = await readJob(ctx.req.param('id'));
if (!job) throw NOT_FOUND('Job not found');
return ctx.json(job);
});
queueRouter.post('/jobs', async (ctx) => {
const user = ctx.get('user');
const body = ctx.get('body');
const { lane, type, meta } = body as { lane: string; type: string; meta?: Record<string, unknown> };
const job = await enqueue({ lane, type, userId: user.email, meta });
return ctx.json(job, 201);
});
queueRouter.delete('/jobs/:id', async (ctx) => {
const job = await cancelJob(ctx.req.param('id'));
if (!job) throw NOT_FOUND('Job not found');
return ctx.json(job);
});
+17 -7
View File
@@ -1,5 +1,6 @@
import type { ServerWebSocket } from 'bun';
import { mkdirSync, statSync } from 'node:fs';
import { existsSync, mkdirSync, statSync } from 'node:fs';
import { homedir } from 'node:os';
import { dirname, join } from 'node:path';
import { fileURLToPath } from 'node:url';
import { getHomeDir, getGlobalSkillsDir, getGlobalToolsDir, getGlobalExtensionsDir, getUserSkillsDir, getUserToolsDir, DATA_PATH } from '@@/data-path';
@@ -101,9 +102,9 @@ const ensureDockerImage = () => {
dockerImageReady = true;
};
// Check whether a container already has the officer resource mounts.
// We test for the global skills dir as a proxy for all mounts being present.
const containerHasResourceMounts = (dockerId: string): boolean => {
// Check whether a container has all expected volume mounts.
// Tests for multiple mount sources — if any is missing, the container should be recreated.
const containerHasExpectedMounts = (dockerId: string): boolean => {
const dockerPath = Bun.which('docker') ?? 'docker';
const result = Bun.spawnSync({
cmd: [dockerPath, 'inspect', '--format', '{{range .Mounts}}{{.Source}}\n{{end}}', dockerId],
@@ -111,7 +112,8 @@ const containerHasResourceMounts = (dockerId: string): boolean => {
stderr: 'ignore',
});
if (result.exitCode !== 0) return false;
return result.stdout.toString().includes(getGlobalSkillsDir());
const mounts = result.stdout.toString();
return mounts.includes(getGlobalSkillsDir()) && mounts.includes('google-oauth.json');
};
const startDockerSidecar = (port: number, homeDir: string, userId: number, username: string, email: string): { dockerId: string } => {
@@ -136,6 +138,11 @@ const startDockerSidecar = (port: number, homeDir: string, userId: number, usern
}
const containerHome = `/home/${username}`;
const googleConfigHost = join(homedir(), '.config', 'officer.dev', 'google-oauth.json');
const googleMounts: string[] = existsSync(googleConfigHost)
? ['-v', `${googleConfigHost}:/officer/google-oauth.json:ro`]
: [];
const run = Bun.spawnSync({
cmd: [
dockerPath,
@@ -164,6 +171,8 @@ const startDockerSidecar = (port: number, homeDir: string, userId: number, usern
'-v', `${getUserSkillsDir(email)}:/officer/user/skills:ro`,
'-v', `${getUserToolsDir(email)}:/officer/user/tools:ro`,
'-v', `${join(DATA_PATH, '.generated')}:/officer/generated:ro`,
...googleMounts,
'-v', `${join(DATA_PATH, email, 'integrations')}:/officer/user/integrations:ro`,
'-w', containerHome,
tag,
],
@@ -231,13 +240,14 @@ export const ensureDockerContainer = async (email: string, userId: number, homeD
// Ensure user-specific resource dirs exist before mounting (Docker creates them as root if missing)
mkdirSync(getUserSkillsDir(email), { recursive: true });
mkdirSync(getUserToolsDir(email), { recursive: true });
mkdirSync(join(DATA_PATH, email, 'integrations'), { recursive: true });
const map = await loadContainerMap();
const existing = map[email];
if (existing && dockerContainerRunning(existing.dockerId)) {
// Recreate if resource mounts are missing (e.g. first run after feature was added)
if (!containerHasResourceMounts(existing.dockerId)) {
if (!containerHasExpectedMounts(existing.dockerId)) {
console.log(`[terminal] recreating container for ${email} — resource mounts missing`);
stopDockerSidecar(existing.dockerId);
} else {
@@ -246,7 +256,7 @@ export const ensureDockerContainer = async (email: string, userId: number, homeD
}
if (existing && dockerContainerExists(existing.dockerId)) {
if (!containerHasResourceMounts(existing.dockerId)) {
if (!containerHasExpectedMounts(existing.dockerId)) {
stopDockerSidecar(existing.dockerId);
} else if (dockerStart(existing.dockerId)) {
return existing;
+5
View File
@@ -10,6 +10,7 @@ import { syncSeedTools } from './sync-tools';
import { syncSeedExtensions } from './sync-extensions';
import { syncSeedResources } from './sync-resources';
import { migrateSettingsToResources } from './migrate-resources';
import { initQueue } from './queue';
mkdirSync(DATA_PATH, { recursive: true });
mkdirSync(PI_CONFIG_DIR, { recursive: true });
@@ -87,4 +88,8 @@ function seedPiConfig(): void {
await syncAllUserPiConfigs().catch(err => {
console.error('[bootstrap] Failed to sync user Pi configs:', err);
});
await initQueue().catch(err => {
console.error('[bootstrap] Failed to initialize queue:', err);
});
})();
+2
View File
@@ -88,3 +88,5 @@ export const getAttachmentsDir = (email: string, provider: 'claude' | 'opencode'
join(DATA_PATH, email, 'chat_sessions', provider, sessionId, 'attachments');
export const getUserDockFile = (email: string) => join(DATA_PATH, email, 'dock', 'dock.json');
export const getUserEmailDir = (email: string) => join(DATA_PATH, email, 'Gmail', 'emails');
+4
View File
@@ -21,6 +21,8 @@ import { piRestRouter } from './api/pi/rest';
import { devServerRouter, devServerProxyRouter } from './api/dev-server/router';
import { dockRouter } from './api/dock/dock';
import { integrationsRouter, googleCallbackHandler } from './api/integrations/integrations';
import { queueRouter } from './api/queue/queue';
import { emailRouter } from './api/email/email';
import { CustomError } from './custom-errors';
import { userMiddleware, bodyParser } from './_middlewares';
@@ -64,6 +66,8 @@ protectedRouter.route('/file-browser', fileBrowserRouter);
protectedRouter.route('/dev-server', devServerRouter);
protectedRouter.route('/dock', dockRouter);
protectedRouter.route('/integrations', integrationsRouter);
protectedRouter.route('/queue', queueRouter);
protectedRouter.route('/email', emailRouter);
protectedRouter.route('/', piRestRouter);
honoServer.route('/api', protectedRouter);
+217
View File
@@ -0,0 +1,217 @@
import type { Job, JobProgress, EnqueueParams, StepContext } from './types';
import { readJob, writeJob, listAllJobs } from './storage';
import { getHandler } from './handler-registry';
import { sendMail } from 'emailer';
const activeLanes = new Map<string, boolean>();
const PROGRESS_THROTTLE_MS = 1000;
export async function enqueue(params: EnqueueParams): Promise<Job> {
const handler = getHandler(params.type);
if (!handler) throw new Error(`No handler registered for job type: ${params.type}`);
const job: Job = {
id: crypto.randomUUID(),
lane: params.lane,
type: params.type,
userId: params.userId,
status: 'queued',
steps: handler.steps.map((s) => ({ name: s.name, status: 'pending' as const })),
currentStep: 0,
createdAt: Date.now(),
meta: params.meta,
};
await writeJob(job);
console.log(`[queue] Enqueued job ${job.id} (${job.type}) in lane ${job.lane}`);
kickLane(job.lane);
return job;
}
export async function cancelJob(id: string): Promise<Job | null> {
const job = await readJob(id);
if (!job) return null;
if (job.status === 'completed' || job.status === 'failed' || job.status === 'cancelled') return job;
job.status = 'cancelled';
job.completedAt = Date.now();
for (const step of job.steps) {
if (step.status === 'pending' || step.status === 'running') {
step.status = 'failed';
step.error = 'Cancelled';
}
}
await writeJob(job);
console.log(`[queue] Cancelled job ${job.id}`);
return job;
}
export async function resumeInterruptedJobs() {
const jobs = await listAllJobs();
const lanesToKick = new Set<string>();
for (const job of jobs) {
if (job.status === 'running') {
job.status = 'queued';
job.startedAt = undefined;
for (const step of job.steps) {
if (step.status === 'running') {
step.status = 'pending';
step.startedAt = undefined;
}
}
await writeJob(job);
console.log(`[queue] Reset interrupted job ${job.id} back to queued`);
lanesToKick.add(job.lane);
} else if (job.status === 'queued') {
lanesToKick.add(job.lane);
}
}
for (const lane of lanesToKick) {
kickLane(lane);
}
}
function kickLane(lane: string) {
if (activeLanes.get(lane)) return;
activeLanes.set(lane, true);
processNextInLane(lane);
}
async function processNextInLane(lane: string) {
try {
const jobs = await listAllJobs();
const next = jobs
.filter((j) => j.lane === lane && j.status === 'queued')
.sort((a, b) => a.createdAt - b.createdAt)[0];
if (!next) {
activeLanes.set(lane, false);
return;
}
await runJob(next);
} catch (err) {
console.error(`[queue] Lane ${lane} processing error:`, err);
} finally {
const jobs = await listAllJobs();
const hasMore = jobs.some((j) => j.lane === lane && j.status === 'queued');
if (hasMore) {
processNextInLane(lane);
} else {
activeLanes.set(lane, false);
}
}
}
async function runJob(job: Job) {
const handler = getHandler(job.type);
if (!handler) {
job.status = 'failed';
job.error = `No handler for type: ${job.type}`;
job.completedAt = Date.now();
await writeJob(job);
return;
}
job.status = 'running';
job.startedAt = Date.now();
await writeJob(job);
console.log(`[queue] Running job ${job.id} (${job.type})`);
const sharedMeta: Record<string, unknown> = { ...(job.meta ?? {}) };
for (let i = 0; i < handler.steps.length; i++) {
const fresh = await readJob(job.id);
if (!fresh || fresh.status === 'cancelled') {
console.log(`[queue] Job ${job.id} was cancelled, stopping`);
return;
}
const handlerStep = handler.steps[i]!;
const step = fresh.steps[i]!;
fresh.currentStep = i;
step.status = 'running';
step.startedAt = Date.now();
await writeJob(fresh);
console.log(`[queue] Job ${fresh.id} step ${i + 1}/${handler.steps.length}: "${handlerStep.name}"`);
let lastProgressWrite = 0;
let pendingProgress: JobProgress | null = null;
const updateProgress = async (progress: JobProgress) => {
step.progress = progress;
const now = Date.now();
if (now - lastProgressWrite >= PROGRESS_THROTTLE_MS) {
lastProgressWrite = now;
pendingProgress = null;
await writeJob(fresh);
} else {
pendingProgress = progress;
}
};
const ctx: StepContext = { job: fresh, step, updateProgress, meta: sharedMeta };
try {
await handlerStep.run(ctx);
if (pendingProgress) {
step.progress = pendingProgress;
}
step.status = 'completed';
step.completedAt = Date.now();
await writeJob(fresh);
} catch (err) {
const errorMessage = err instanceof Error ? err.message : String(err);
step.status = 'failed';
step.error = errorMessage;
step.completedAt = Date.now();
fresh.status = 'failed';
fresh.error = `Step "${step.name}" failed: ${errorMessage}`;
fresh.completedAt = Date.now();
await writeJob(fresh);
console.error(`[queue] Job ${fresh.id} failed at step "${step.name}":`, errorMessage);
await notifyFailure(fresh);
return;
}
}
const final = await readJob(job.id);
if (final && final.status === 'running') {
final.status = 'completed';
final.completedAt = Date.now();
await writeJob(final);
console.log(`[queue] Job ${final.id} completed`);
await notifyCompletion(final);
}
}
async function notifyCompletion(job: Job) {
try {
await sendMail({
template: 'JobCompleted',
subject: `Job completed: ${job.type}`,
to: job.userId,
data: { job },
});
} catch {
// SMTP might not be configured — non-fatal
}
}
async function notifyFailure(job: Job) {
try {
await sendMail({
template: 'JobFailed',
subject: `Job failed: ${job.type}`,
to: job.userId,
data: { job },
});
} catch {
// SMTP might not be configured — non-fatal
}
}
+11
View File
@@ -0,0 +1,11 @@
import type { JobHandler } from './types';
const handlers = new Map<string, JobHandler>();
export function registerHandler(handler: JobHandler) {
handlers.set(handler.type, handler);
}
export function getHandler(type: string): JobHandler | undefined {
return handlers.get(type);
}
+271
View File
@@ -0,0 +1,271 @@
import { mkdirSync, readdirSync, writeFileSync } from 'node:fs';
import { join } from 'node:path';
import { homedir } from 'node:os';
import type { JobHandler } from '../types';
import { registerHandler } from '../handler-registry';
import { DATA_PATH } from '../../data-path';
type GoogleCredentials = {
accessToken: string;
refreshToken: string;
expiresAt: number;
clientId: string;
clientSecret: string;
};
const googleConfigPath = join(homedir(), '.config', 'officer.dev', 'google-oauth.json');
async function loadCredentials(userId: string): Promise<GoogleCredentials> {
const config = await Bun.file(googleConfigPath).json().catch(() => null);
if (!config?.clientId || !config?.clientSecret) {
throw new Error('Google OAuth not configured — ask your admin to set up credentials');
}
const tokenPath = join(DATA_PATH, userId, 'integrations', 'google.json');
const token = await Bun.file(tokenPath).json().catch(() => null);
if (!token?.accessToken) {
throw new Error('Google account not connected — connect in Settings → Integrations');
}
return {
accessToken: token.accessToken,
refreshToken: token.refreshToken ?? '',
expiresAt: token.expiresAt ?? 0,
clientId: config.clientId,
clientSecret: config.clientSecret,
};
}
let cachedAccessToken: string | null = null;
let cachedExpiresAt = 0;
async function getValidAccessToken(creds: GoogleCredentials): Promise<string> {
if (cachedAccessToken && cachedExpiresAt > Date.now() + 5 * 60 * 1000) {
return cachedAccessToken;
}
if (creds.expiresAt > Date.now() + 5 * 60 * 1000) {
cachedAccessToken = creds.accessToken;
cachedExpiresAt = creds.expiresAt;
return creds.accessToken;
}
if (!creds.refreshToken) throw new Error('Token expired and no refresh token available');
const res = await fetch('https://oauth2.googleapis.com/token', {
method: 'POST',
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
body: new URLSearchParams({
client_id: creds.clientId,
client_secret: creds.clientSecret,
refresh_token: creds.refreshToken,
grant_type: 'refresh_token',
}),
});
if (!res.ok) {
const error = await res.text().catch(() => '');
throw new Error(`Token refresh failed (${res.status}): ${error}`);
}
const data = (await res.json()) as { access_token: string; expires_in?: number };
cachedAccessToken = data.access_token;
cachedExpiresAt = Date.now() + (data.expires_in ?? 3600) * 1000;
return data.access_token;
}
const GMAIL_BASE = 'https://gmail.googleapis.com/gmail/v1/users/me';
async function gmailGet(token: string, path: string, params?: Record<string, string>): Promise<unknown> {
const url = new URL(`${GMAIL_BASE}${path}`);
if (params) {
for (const [k, v] of Object.entries(params)) {
if (v) url.searchParams.set(k, v);
}
}
const res = await fetch(url.toString(), {
headers: { Authorization: `Bearer ${token}` },
});
if (!res.ok) {
const error = await res.text().catch(() => '');
throw new Error(`Gmail API error (${res.status}): ${error}`);
}
return res.json();
}
function slugify(text: string, maxLen = 60): string {
return text
.toLowerCase()
.replace(/[^a-z0-9]+/g, '-')
.replace(/^-+|-+$/g, '')
.slice(0, maxLen)
.replace(/-+$/, '');
}
function buildEmlFilename(id: string, internalDate: string | undefined, rawEmail: string): string {
const subjectMatch = rawEmail.match(/^Subject:\s*(.+)$/mi);
const subject = subjectMatch?.[1]?.trim() || 'no-subject';
const ts = parseInt(internalDate || '0');
const d = new Date(ts);
const dateStr =
ts > 0
? `${d.getFullYear()}-${String(d.getMonth() + 1).padStart(2, '0')}-${String(d.getDate()).padStart(2, '0')}`
: 'unknown-date';
return `${dateStr}_${slugify(subject)}_${id}.eml`;
}
type SyncProgress = { saved: number; skipped: number; errors: number; page: number };
type OnProgress = (progress: SyncProgress) => void;
async function syncInbox(
token: string,
outputDir: string,
query?: string,
onProgress?: OnProgress,
): Promise<{ saved: number; skipped: number; errors: number }> {
mkdirSync(outputDir, { recursive: true });
const existingIds = new Set<string>();
try {
for (const file of readdirSync(outputDir)) {
const match = file.match(/_([a-f0-9]+)\.eml$/i);
if (match) existingIds.add(match[1]!);
}
} catch {
/* dir might not exist yet */
}
let saved = 0;
let skipped = 0;
let errors = 0;
let page = 0;
let pageToken: string | undefined;
do {
const params: Record<string, string> = { maxResults: '100' };
if (query) params.q = query;
if (pageToken) params.pageToken = pageToken;
const list = (await gmailGet(token, '/messages', params)) as {
messages?: Array<{ id: string }>;
nextPageToken?: string;
};
const messages = list.messages ?? [];
if (messages.length === 0) break;
for (let i = 0; i < messages.length; i += 5) {
const batch = messages.slice(i, i + 5);
await Promise.all(
batch.map(async ({ id }) => {
if (existingIds.has(id)) {
skipped++;
return;
}
try {
const msg = (await gmailGet(token, `/messages/${id}`, { format: 'raw' })) as {
id: string;
internalDate?: string;
raw: string;
};
const rawEmail = Buffer.from(msg.raw, 'base64url').toString('utf-8');
const filename = buildEmlFilename(msg.id, msg.internalDate, rawEmail);
writeFileSync(join(outputDir, filename), rawEmail);
existingIds.add(id);
saved++;
} catch {
errors++;
}
}),
);
}
page++;
console.log(`[gmail-sync] Page ${page}: saved ${saved}, skipped ${skipped}, errors ${errors}`);
onProgress?.({ saved, skipped, errors, page });
pageToken = list.nextPageToken;
} while (pageToken);
return { saved, skipped, errors };
}
const MONTH_NAMES = ['January', 'February', 'March', 'April', 'May', 'June', 'July', 'August', 'September', 'October', 'November', 'December'];
function buildMonthRanges(year: number): Array<{ label: string; after: string; before: string }> {
const now = new Date();
const currentMonth = now.getFullYear() === year ? now.getMonth() : 11;
const ranges: Array<{ label: string; after: string; before: string }> = [];
for (let m = 0; m <= currentMonth; m++) {
const after = `${year}/${m + 1}/1`;
const before = m < 11 ? `${year}/${m + 2}/1` : `${year + 1}/1/1`;
ranges.push({ label: `${MONTH_NAMES[m]!} ${year}`, after, before });
}
return ranges;
}
const gmailSyncHandler: JobHandler = {
type: 'gmail-sync',
steps: [
{
name: 'Verify connection',
run: async (ctx) => {
const creds = await loadCredentials(ctx.job.userId);
const token = await getValidAccessToken(creds);
// Store token in shared meta for the next step
ctx.meta.accessToken = token;
ctx.meta.outputDir = join(DATA_PATH, ctx.job.userId, 'Gmail', 'emails');
},
},
{
name: 'Sync emails',
run: async (ctx) => {
const token = ctx.meta.accessToken as string;
const outputDir = ctx.meta.outputDir as string;
const year = ctx.meta.year as number | undefined;
let totalSaved = 0;
let totalSkipped = 0;
let totalErrors = 0;
if (year) {
// Year-scoped sync: month by month with progress
const months = buildMonthRanges(year);
for (let i = 0; i < months.length; i++) {
const month = months[i]!;
await ctx.updateProgress({ current: i, total: months.length, label: month.label });
const query = `after:${month.after} before:${month.before}`;
const { saved, skipped, errors } = await syncInbox(token, outputDir, query, (p) => {
const label = `${month.label}${p.saved} saved`;
ctx.updateProgress({ current: i, total: months.length, label });
});
totalSaved += saved;
totalSkipped += skipped;
totalErrors += errors;
}
await ctx.updateProgress({ current: months.length, total: months.length, label: 'Done' });
} else {
// Full sync: all emails with per-page progress
await ctx.updateProgress({ current: 0, total: 0, label: 'Starting sync' });
const { saved, skipped, errors } = await syncInbox(token, outputDir, undefined, (p) => {
const label = `Saved ${p.saved}, skipped ${p.skipped} (page ${p.page})`;
ctx.updateProgress({ current: p.saved + p.skipped + p.errors, total: 0, label });
});
totalSaved = saved;
totalSkipped = skipped;
totalErrors = errors;
await ctx.updateProgress({ current: 1, total: 1, label: 'Done' });
}
console.log(`[gmail-sync] Saved ${totalSaved}, skipped ${totalSkipped}, errors ${totalErrors}`);
},
},
],
};
registerHandler(gmailSyncHandler);
+1
View File
@@ -0,0 +1 @@
import './gmail-sync';
+14
View File
@@ -0,0 +1,14 @@
import { ensureQueueDir } from './storage';
import { resumeInterruptedJobs } from './engine';
import './handlers';
export { enqueue, cancelJob } from './engine';
export { readJob, listAllJobs } from './storage';
export { registerHandler } from './handler-registry';
export type { Job, JobStep, JobStatus, JobProgress, StepContext, JobHandler, JobHandlerStep, EnqueueParams } from './types';
export async function initQueue() {
await ensureQueueDir();
await resumeInterruptedJobs();
console.log('[queue] Initialized');
}
+51
View File
@@ -0,0 +1,51 @@
import type { Job } from './types';
import { join } from 'node:path';
import { readdir, mkdir, unlink } from 'node:fs/promises';
import { DATA_PATH } from '../data-path';
const QUEUE_DIR = join(DATA_PATH, 'queue', 'jobs');
export async function ensureQueueDir() {
await mkdir(QUEUE_DIR, { recursive: true });
}
export async function readJob(id: string): Promise<Job | null> {
const file = Bun.file(join(QUEUE_DIR, `${id}.json`));
if (!(await file.exists())) return null;
try {
return await file.json();
} catch {
return null;
}
}
export async function writeJob(job: Job): Promise<void> {
await Bun.write(join(QUEUE_DIR, `${job.id}.json`), JSON.stringify(job, null, 2));
}
export async function listAllJobs(): Promise<Job[]> {
try {
const entries = await readdir(QUEUE_DIR);
const jobs: Job[] = [];
for (const entry of entries) {
if (!entry.endsWith('.json')) continue;
const file = Bun.file(join(QUEUE_DIR, entry));
try {
jobs.push(await file.json());
} catch {
// corrupted file — skip
}
}
return jobs.sort((a, b) => b.createdAt - a.createdAt);
} catch {
return [];
}
}
export async function deleteJobFile(id: string): Promise<void> {
try {
await unlink(join(QUEUE_DIR, `${id}.json`));
} catch {
// file doesn't exist — fine
}
}
+57
View File
@@ -0,0 +1,57 @@
export type JobStatus = 'queued' | 'running' | 'completed' | 'failed' | 'cancelled';
export type JobStepStatus = 'pending' | 'running' | 'completed' | 'failed';
export type JobProgress = {
current: number;
total: number;
label?: string;
};
export type JobStep = {
name: string;
status: JobStepStatus;
startedAt?: number;
completedAt?: number;
error?: string;
progress?: JobProgress;
};
export type Job = {
id: string;
lane: string;
type: string;
userId: string;
status: JobStatus;
steps: JobStep[];
currentStep: number;
error?: string;
createdAt: number;
startedAt?: number;
completedAt?: number;
meta?: Record<string, unknown>;
};
export type StepContext = {
job: Job;
step: JobStep;
updateProgress: (progress: JobProgress) => Promise<void>;
meta: Record<string, unknown>;
};
export type JobHandlerStep = {
name: string;
run: (ctx: StepContext) => Promise<void>;
};
export type JobHandler = {
type: string;
steps: JobHandlerStep[];
};
export type EnqueueParams = {
lane: string;
type: string;
userId: string;
meta?: Record<string, unknown>;
};