diff --git a/src/lib/atproto/client.ts b/src/lib/atproto/client.ts new file mode 100644 index 0000000..ffb963c --- /dev/null +++ b/src/lib/atproto/client.ts @@ -0,0 +1,51 @@ +/** + * AT Protocol authenticated client + * Makes authenticated XRPC calls through the Workers proxy + */ + +const WORKERS_API = import.meta.env.DEV + ? "http://localhost:8787" + : "https://connections.directory"; + +/** + * Make an authenticated XRPC call through the Workers proxy + */ +export async function xrpc( + method: string, + params?: Record, + data?: Record, +): Promise { + const response = await fetch(`${WORKERS_API}/api/xrpc`, { + method: "POST", + headers: { + "Content-Type": "application/json", + }, + credentials: "include", + body: JSON.stringify({ + method, + params, + data, + }), + }); + + if (!response.ok) { + const error = await response.json(); + throw new Error(error.error || error.message || "XRPC call failed"); + } + + return response.json(); +} + +// For backwards compatibility +export function initializeAgent(_pdsUrl: string): void { + // No-op now that we use the proxy +} + +export function clearAgent(): void { + // No-op now that we use the proxy +} + +export function getAgent(): any { + // Return a mock agent that throws if used + return null; +} diff --git a/src/lib/atproto/records.ts b/src/lib/atproto/records.ts index 9d2ab4d..b9a6f13 100644 --- a/src/lib/atproto/records.ts +++ b/src/lib/atproto/records.ts @@ -3,7 +3,9 @@ * Fetch and manage records from the user's PDS */ -export interface FollowRecord { +import { xrpc } from "./client"; + +export interface IdSifaGraphFollow { uri: string; cid: string; value: { @@ -36,7 +38,7 @@ async function getPdsUrl(did: string): Promise { */ export async function listFollowRecords( userDid: string, -): Promise { +): Promise { try { const pdsUrl = await getPdsUrl(userDid); @@ -86,36 +88,22 @@ export async function resolveHandle( } /** - * Create a follow record at the user's PDS via Workers backend - * The Workers has the OAuth session needed for authenticated requests + * Create a follow record at the user's PDS using the authenticated XRPC proxy */ export async function createFollowRecord( subjectDid: string, + userDid: string, ): Promise<{ uri: string; cid: string } | null> { try { - const WORKERS_API = import.meta.env.DEV - ? "http://localhost:8787" - : "https://connections.directory"; - - const response = await fetch(`${WORKERS_API}/api/sifa/follow`, { - method: "POST", - headers: { - "Content-Type": "application/json", - }, - credentials: "include", - body: JSON.stringify({ + const result = await xrpc("com.atproto.repo.createRecord", undefined, { + repo: userDid, + collection: "id.sifa.graph.follow", + record: { subject: subjectDid, - }), + createdAt: new Date().toISOString(), + }, }); - if (!response.ok) { - const error = await response.json(); - throw new Error( - error.message || error.error || "Failed to create follow record", - ); - } - - const result = await response.json(); return { uri: result.uri, cid: result.cid, @@ -127,31 +115,23 @@ export async function createFollowRecord( } /** - * Delete a follow record from the user's PDS via Workers backend + * Delete a follow record from the user's PDS using the authenticated XRPC proxy */ -export async function deleteFollowRecord(recordUri: string): Promise { +export async function deleteFollowRecord( + recordUri: string, + userDid: string, +): Promise { try { - const WORKERS_API = import.meta.env.DEV - ? "http://localhost:8787" - : "https://connections.directory"; - - const response = await fetch(`${WORKERS_API}/api/sifa/follow`, { - method: "DELETE", - headers: { - "Content-Type": "application/json", - }, - credentials: "include", - body: JSON.stringify({ - uri: recordUri, - }), + // Parse the record URI to get rkey + // Format: at://did:plc:abc.../id.sifa.graph.follow/rkey + const parts = recordUri.split("/"); + const rkey = parts[parts.length - 1]; + + await xrpc("com.atproto.repo.deleteRecord", undefined, { + repo: userDid, + collection: "id.sifa.graph.follow", + rkey: rkey, }); - - if (!response.ok) { - const error = await response.json(); - throw new Error( - error.message || error.error || "Failed to delete follow record", - ); - } } catch (error) { console.error("Failed to delete Sifa follow record:", error); throw error; diff --git a/src/pages/ScanPage.tsx b/src/pages/ScanPage.tsx index 6c8fa46..d1a3611 100644 --- a/src/pages/ScanPage.tsx +++ b/src/pages/ScanPage.tsx @@ -100,7 +100,7 @@ export default function ScanPage() { // Sync to Sifa graph try { const { createFollowRecord } = await import("../lib/atproto/records"); - const result = await createFollowRecord(scannedDid); + const result = await createFollowRecord(scannedDid, did!); if (result) { // Update connection with record URI/CID diff --git a/src/store/auth.ts b/src/store/auth.ts index 3e4e1cf..59c27c9 100644 --- a/src/store/auth.ts +++ b/src/store/auth.ts @@ -2,123 +2,138 @@ * Authentication state management with Zustand */ -import { create } from 'zustand' +import { create } from "zustand"; +import { initializeAgent, clearAgent } from "../lib/atproto/client"; interface AuthState { - authenticated: boolean - did: string | null - handle: string | null - loading: boolean - error: string | null + authenticated: boolean; + did: string | null; + handle: string | null; + pdsUrl: string | null; + loading: boolean; + error: string | null; // Actions - checkSession: () => Promise - login: (handle: string) => Promise - logout: () => Promise + checkSession: () => Promise; + login: (handle: string) => Promise; + logout: () => Promise; } const WORKERS_API = import.meta.env.DEV - ? 'http://localhost:8787' - : 'https://connections.directory' + ? "http://localhost:8787" + : "https://connections.directory"; export const useAuthStore = create((set) => ({ authenticated: false, did: null, handle: null, + pdsUrl: null, loading: false, error: null, checkSession: async () => { try { - set({ loading: true, error: null }) + set({ loading: true, error: null }); const response = await fetch(`${WORKERS_API}/api/auth/session`, { - credentials: 'include', - }) + credentials: "include", + }); - const data = await response.json() + const data = await response.json(); if (data.authenticated) { + // Initialize the authenticated agent + if (data.pdsUrl) { + await initializeAgent(data.pdsUrl); + } + set({ authenticated: true, did: data.did, handle: data.handle, + pdsUrl: data.pdsUrl, loading: false, - }) + }); } else { set({ authenticated: false, did: null, handle: null, + pdsUrl: null, loading: false, - }) + }); } } catch (error) { - console.error('Session check failed:', error) + console.error("Session check failed:", error); set({ authenticated: false, did: null, handle: null, + pdsUrl: null, loading: false, - error: 'Failed to check session', - }) + error: "Failed to check session", + }); } }, login: async (handle: string) => { try { - set({ loading: true, error: null }) + set({ loading: true, error: null }); const response = await fetch(`${WORKERS_API}/api/auth/login`, { - method: 'POST', + method: "POST", headers: { - 'Content-Type': 'application/json', + "Content-Type": "application/json", }, body: JSON.stringify({ handle }), - credentials: 'include', - }) + credentials: "include", + }); - const data = await response.json() + const data = await response.json(); if (data.authUrl) { // Redirect to OAuth provider - window.location.href = data.authUrl + window.location.href = data.authUrl; } else { set({ loading: false, - error: data.error || 'Login failed', - }) + error: data.error || "Login failed", + }); } } catch (error) { - console.error('Login failed:', error) + console.error("Login failed:", error); set({ loading: false, - error: 'Failed to initiate login', - }) + error: "Failed to initiate login", + }); } }, logout: async () => { try { - set({ loading: true, error: null }) + set({ loading: true, error: null }); await fetch(`${WORKERS_API}/api/auth/logout`, { - method: 'POST', - credentials: 'include', - }) + method: "POST", + credentials: "include", + }); + + // Clear the agent + clearAgent(); set({ authenticated: false, did: null, handle: null, + pdsUrl: null, loading: false, - }) + }); } catch (error) { - console.error('Logout failed:', error) + console.error("Logout failed:", error); set({ loading: false, - error: 'Failed to logout', - }) + error: "Failed to logout", + }); } }, -})) +})); diff --git a/workers/src/index.ts b/workers/src/index.ts index 5a82926..4c8babe 100644 --- a/workers/src/index.ts +++ b/workers/src/index.ts @@ -11,6 +11,7 @@ import { handleLogout, } from "./oauth/handlers"; import { handleCreateFollow, handleDeleteFollow } from "./sifa/handlers"; +import { handleXRPCProxy } from "./proxy/handlers"; export default { async fetch(request: Request, env: Env): Promise { @@ -81,7 +82,16 @@ export default { return response; } - // Sifa graph endpoints + // Authenticated XRPC proxy + if (url.pathname === "/api/xrpc" && request.method === "POST") { + const response = await handleXRPCProxy(request, env); + Object.entries(corsHeaders).forEach(([key, value]) => { + response.headers.set(key, value); + }); + return response; + } + + // Sifa graph endpoints (kept for backwards compatibility) if (url.pathname === "/api/sifa/follow" && request.method === "POST") { const response = await handleCreateFollow(request, env); Object.entries(corsHeaders).forEach(([key, value]) => { diff --git a/workers/src/oauth/handlers.ts b/workers/src/oauth/handlers.ts index 8437d99..49d3862 100644 --- a/workers/src/oauth/handlers.ts +++ b/workers/src/oauth/handlers.ts @@ -66,6 +66,15 @@ export async function handleCallback( // Handle OAuth callback const { session } = await client.callback(params); + // Get PDS URL from DID document + const didDoc = await fetch(`https://plc.directory/${session.did}`).then( + (r) => r.json(), + ); + const pdsService = didDoc.service?.find( + (s: any) => s.type === "AtprotoPersonalDataServer", + ); + const pdsUrl = pdsService?.serviceEndpoint || "https://bsky.social"; + // Create session ID const sessionId = crypto.randomUUID(); @@ -76,6 +85,7 @@ export async function handleCallback( did: session.did, handle: session.handle || session.did, sub: session.sub, // This is the key for OAuth client.restore() + pdsUrl: pdsUrl, }), { expirationTtl: 60 * 60 * 24 * 7 }, // 7 days ); @@ -150,6 +160,7 @@ export async function handleSession( authenticated: true, did: session.did, handle: session.handle, + pdsUrl: session.pdsUrl, }), { status: 200, diff --git a/workers/src/proxy/handlers.ts b/workers/src/proxy/handlers.ts new file mode 100644 index 0000000..4367ebe --- /dev/null +++ b/workers/src/proxy/handlers.ts @@ -0,0 +1,105 @@ +/** + * Authenticated proxy for XRPC calls + * Allows the frontend to make authenticated AT Protocol requests through the Workers + */ + +import { Env, createOAuthClient } from "../oauth/client"; + +export async function handleXRPCProxy( + request: Request, + env: Env, +): Promise { + try { + // Get session from cookie + const cookieHeader = request.headers.get("Cookie") || ""; + const cookies = Object.fromEntries( + cookieHeader.split(";").map((c) => c.trim().split("=")), + ); + const sessionId = cookies.session; + + if (!sessionId) { + return new Response(JSON.stringify({ error: "Not authenticated" }), { + status: 401, + headers: { "Content-Type": "application/json" }, + }); + } + + // Get session data + const sessionDataRaw = await env.OAUTH_SESSIONS.get( + `session:${sessionId}`, + "text", + ); + + if (!sessionDataRaw) { + return new Response(JSON.stringify({ error: "Session not found" }), { + status: 401, + headers: { "Content-Type": "application/json" }, + }); + } + + const sessionData = JSON.parse(sessionDataRaw); + + // Create OAuth client and restore session + const client = await createOAuthClient(env); + const oauthSession = await client.restore(sessionData.sub); + + if (!oauthSession) { + return new Response(JSON.stringify({ error: "Failed to restore OAuth session" }), { + status: 401, + headers: { "Content-Type": "application/json" }, + }); + } + + // Get the XRPC request body + const body = await request.json() as { + method: string; + params?: Record; + data?: Record; + }; + + // Build the XRPC URL + const xrpcUrl = `${sessionData.pdsUrl}/xrpc/${body.method}`; + const url = new URL(xrpcUrl); + + // Add query params if provided + if (body.params) { + Object.entries(body.params).forEach(([key, value]) => { + url.searchParams.set(key, String(value)); + }); + } + + // Make the authenticated request using dpopFetch + const fetchOptions: RequestInit = { + method: body.data ? "POST" : "GET", + headers: { + "Content-Type": "application/json", + }, + }; + + if (body.data) { + fetchOptions.body = JSON.stringify(body.data); + } + + const response = await oauthSession.dpopFetch(url.toString(), fetchOptions); + + // Return the response + const responseData = await response.json(); + + return new Response(JSON.stringify(responseData), { + status: response.status, + headers: { "Content-Type": "application/json" }, + }); + } catch (error) { + console.error("XRPC proxy error:", error); + return new Response( + JSON.stringify({ + error: "XRPC request failed", + details: error instanceof Error ? error.message : "Unknown error", + }), + { + status: 500, + headers: { "Content-Type": "application/json" }, + }, + ); + } +} diff --git a/workers/src/sifa/handlers.ts b/workers/src/sifa/handlers.ts index 004a0ec..5d34ce2 100644 --- a/workers/src/sifa/handlers.ts +++ b/workers/src/sifa/handlers.ts @@ -1,9 +1,9 @@ /** * Sifa graph handlers for creating follow records + * Uses OAuth session's dpopFetch for authenticated requests */ import { Env, createOAuthClient } from "../oauth/client"; -import { AtpAgent } from "@atproto/api"; export async function handleCreateFollow( request: Request, @@ -66,41 +66,45 @@ export async function handleCreateFollow( ); } - // Create an AtpAgent with the OAuth session - const agent = new AtpAgent({ - service: - oauthSession.serverUrl || - `https://${session.did.split(":")[2]}.host.bsky.network`, - }); - - // Set the session on the agent using the OAuth token - // The OAuth client should provide us with credentials - // For now, let's use the fetchHandler from the OAuth session - - // Create the follow record + // Create the follow record (without $type - it's inferred from collection) const record = { - $type: "id.sifa.graph.follow", subject: subject, createdAt: new Date().toISOString(), }; - // Use the agent to create the record - const result = await agent.api.com.atproto.repo.createRecord( - { - repo: session.did, - collection: "id.sifa.graph.follow", - record: record, - }, + // Get PDS URL from OAuth session + const pdsUrl = oauthSession.serverUrl; + + // Use OAuth session's dpopFetch for authenticated request + const response = await oauthSession.dpopFetch( + `${pdsUrl}/xrpc/com.atproto.repo.createRecord`, { - encoding: "application/json", - headers: await oauthSession.dpopFetch.createHeaders(), + method: "POST", + headers: { + "Content-Type": "application/json", + }, + body: JSON.stringify({ + repo: session.did, + collection: "id.sifa.graph.follow", + record: record, + }), }, ); + if (!response.ok) { + const errorText = await response.text(); + console.error("Create record error:", errorText); + throw new Error( + `Failed to create record: ${response.status} ${errorText}`, + ); + } + + const result = await response.json(); + return new Response( JSON.stringify({ - uri: result.data.uri, - cid: result.data.cid, + uri: result.uri, + cid: result.cid, }), { status: 200, @@ -181,26 +185,32 @@ export async function handleDeleteFollow( ); } - // Create an AtpAgent - const agent = new AtpAgent({ - service: - oauthSession.serverUrl || - `https://${session.did.split(":")[2]}.host.bsky.network`, - }); + // Get PDS URL + const pdsUrl = oauthSession.serverUrl; // Delete the record - await agent.api.com.atproto.repo.deleteRecord( - { - repo: session.did, - collection: "id.sifa.graph.follow", - rkey: uri.split("/").pop()!, - }, + const response = await oauthSession.dpopFetch( + `${pdsUrl}/xrpc/com.atproto.repo.deleteRecord`, { - encoding: "application/json", - headers: await oauthSession.dpopFetch.createHeaders(), + method: "POST", + headers: { + "Content-Type": "application/json", + }, + body: JSON.stringify({ + repo: session.did, + collection: "id.sifa.graph.follow", + rkey: uri.split("/").pop(), // Extract rkey from URI + }), }, ); + if (!response.ok) { + const errorText = await response.text(); + throw new Error( + `Failed to delete record: ${response.status} ${errorText}`, + ); + } + return new Response(JSON.stringify({ success: true }), { status: 200, headers: { "Content-Type": "application/json" },