From 929f91c83286d789af5da68915f908110fdac780 Mon Sep 17 00:00:00 2001 From: Mat Manna Date: Wed, 29 Jul 2026 15:13:11 -0400 Subject: [PATCH] feat: linkify rss/api nodes & add feed db retries --- src/app.ts | 118 ++++++++++++++++++++++++++++++++++------------ src/web/app.tsx | 16 ++++++- src/web/style.css | 3 ++ 3 files changed, 106 insertions(+), 31 deletions(-) diff --git a/src/app.ts b/src/app.ts index 73529a0..7e31136 100644 --- a/src/app.ts +++ b/src/app.ts @@ -1789,42 +1789,48 @@ app.post("/slack/", slackSlashCommand); // === RSS Feed / JSON API === app.get("/feed/:channelId", async (c) => { - const env = c.env; - const rawId = c.req.param("channelId")!; - const wantsJson = rawId.endsWith(".json") || c.req.header("accept")?.includes("application/json"); - const channelId = wantsJson ? rawId.replace(/\.json$/, "") : rawId; - const store = await getStore(env.DATABASE_URL); - const ch = await store.getChannel(channelId); - if (!ch || !ch.enabled) { - return c.json({ error: "not found", id: channelId }, 404); - } - - const limit = Math.min(parseInt(c.req.query("limit") || "50") || 50, 200); - const offset = parseInt(c.req.query("offset") || "0") || 0; - const msgs = await store.getMessages(channelId, limit, offset); + try { + const env = c.env; + const rawId = c.req.param("channelId")!; + const wantsJson = rawId.endsWith(".json") || c.req.header("accept")?.includes("application/json"); + const channelId = wantsJson ? rawId.replace(/\.json$/, "") : rawId; + const store = await getStore(env.DATABASE_URL); + const ch = await store.getChannel(channelId); + if (!ch || !ch.enabled) { + return c.json({ error: "not found", id: channelId }, 404); + } - if (wantsJson) { - return c.json(msgs); - } + const limit = Math.min(parseInt(c.req.query("limit") || "50") || 50, 200); + const offset = parseInt(c.req.query("offset") || "0") || 0; + const msgs = await store.getMessages(channelId, limit, offset); - const baseUrl = env.BASE_URL || "http://localhost:8080"; - const feed = new Feed({ - title: `#${ch.name} — indigest`, - link: `${baseUrl}/feed/${channelId}`, - description: `Recent messages from #${ch.name}`, - }); + if (wantsJson) { + return c.json(msgs); + } - for (const m of msgs) { - feed.addItem({ - title: m.userName || "unknown", - description: m.text.substring(0, 500), - date: new Date(m.timestamp), - guid: `${channelId}:${m.slackTs}`, + const baseUrl = env.BASE_URL || "http://localhost:8080"; + const feed = new Feed({ + title: `#${ch.name} — indigest`, link: `${baseUrl}/feed/${channelId}`, + description: `Recent messages from #${ch.name}`, }); - } - return c.text(feed.rss2(), 200, { "Content-Type": "application/rss+xml; charset=utf-8" }); + for (const m of msgs) { + feed.addItem({ + title: m.userName || "unknown", + description: m.text.substring(0, 500), + date: new Date(m.timestamp), + guid: `${channelId}:${m.slackTs}`, + link: `${baseUrl}/feed/${channelId}`, + }); + } + + return c.text(feed.rss2(), 200, { "Content-Type": "application/rss+xml; charset=utf-8" }); + } catch (err: any) { + console.error("feed error:", err?.message || err); + c.header("Retry-After", "2"); + return c.text("Service Unavailable", 503); + } }); // === REST API === @@ -1921,6 +1927,57 @@ app.get("/api/graph", async (c) => { const storeCache = new Map(); const ddlRan = new Set(); +function sleep(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +function isTransientDbError(err: unknown): boolean { + const msg = (err instanceof Error ? err.message : String(err)).toLowerCase(); + return ( + msg.includes("econnreset") || + msg.includes("econnrefused") || + msg.includes("etimedout") || + msg.includes("timeout") || + msg.includes("network") || + msg.includes("fetch failed") || + msg.includes("dns") || + msg.includes("getaddrinfo") || + msg.includes("connection terminated") || + msg.includes("failed query") + ); +} + +async function withRetry(fn: () => Promise, attempts = 3): Promise { + let lastErr: unknown; + for (let i = 0; i < attempts; i++) { + try { + return await fn(); + } catch (err) { + lastErr = err; + if (!isTransientDbError(err) || i === attempts - 1) throw err; + await sleep(50 * Math.pow(2, i)); + } + } + throw lastErr; +} + +function retryingStore(store: Store): Store { + return { + getChannel: (id) => withRetry(() => store.getChannel(id)), + upsertChannel: (ch) => withRetry(() => store.upsertChannel(ch)), + listEnabledChannels: () => withRetry(() => store.listEnabledChannels()), + upsertMessage: (msg) => withRetry(() => store.upsertMessage(msg)), + deleteMessage: (channelId, slackTs) => withRetry(() => store.deleteMessage(channelId, slackTs)), + getMessages: (channelId, limit, offset) => withRetry(() => store.getMessages(channelId, limit, offset)), + addSubscription: (subscriberChannelId, sourceChannelId) => withRetry(() => store.addSubscription(subscriberChannelId, sourceChannelId)), + removeSubscription: (subscriberChannelId, sourceChannelId) => withRetry(() => store.removeSubscription(subscriberChannelId, sourceChannelId)), + getSubscribersBySource: (sourceChannelId) => withRetry(() => store.getSubscribersBySource(sourceChannelId)), + getSubscriptionsBySubscriber: (subscriberChannelId) => withRetry(() => store.getSubscriptionsBySubscriber(subscriberChannelId)), + getRecentMessages: (channelId, since) => withRetry(() => store.getRecentMessages(channelId, since)), + close: () => store.close(), + }; +} + async function runDdlOnce(databaseUrl: string) { if (ddlRan.has(databaseUrl)) return; ddlRan.add(databaseUrl); @@ -1962,6 +2019,7 @@ async function getStore(databaseUrl: string): Promise { const { PostgresStore } = await import("./store/pg"); store = new PostgresStore(databaseUrl); } + store = retryingStore(store); storeCache.set(databaseUrl, store); if (isNeon) runDdlOnce(databaseUrl).catch(() => {}); return store; diff --git a/src/web/app.tsx b/src/web/app.tsx index 43438ae..49e7619 100644 --- a/src/web/app.tsx +++ b/src/web/app.tsx @@ -83,10 +83,24 @@ function ChannelNode({ data }: { data: any }) { function FeedNode({ data }: { data: any }) { const labels: Record = { rss: "rss", api: "api", webhook: "webhook" }; + const linkFor = (): string | null => { + const channelId = data.channelId; + if (!channelId) return null; + if (data.feedType === "rss") return `/feed/${encodeURIComponent(channelId)}`; + if (data.feedType === "api") return `/api/messages?channel=${encodeURIComponent(channelId)}`; + if (data.feedType === "webhook") return data.url || null; + return null; + }; + const href = linkFor(); return (
- [{labels[data.feedType] || "feed"}] + [{labels[data.feedType] || "feed"}] + {href && ( + + open ↗ + + )}
); } diff --git a/src/web/style.css b/src/web/style.css index dbeed03..dbe55b2 100644 --- a/src/web/style.css +++ b/src/web/style.css @@ -254,6 +254,9 @@ body { border: 1px solid var(--border); border-radius: 3px; background: var(--surface); + display: inline-flex; + align-items: center; + gap: 0.4rem; } .node-subscription { padding: 0.25rem 0.5rem; -- 2.51.2