const std = @import("std"); const atproto_identity = @import("../atproto/identity.zig"); const atproto_oauth = @import("../atproto/oauth.zig"); const atproto_preferences = @import("../atproto/preferences.zig"); const atproto_proxy = @import("../atproto/proxy.zig"); const atproto_repo = @import("../atproto/repo.zig"); const atproto_server = @import("../atproto/server.zig"); const atproto_space = @import("../atproto/space.zig"); const atproto_sync = @import("../atproto/sync.zig"); const build_options = @import("build_options"); const httpz = @import("httpz"); const log = @import("../core/log.zig"); const http_api = @import("api.zig"); const docs = @import("docs.zig"); const account_page = @import("../internal/account.zig"); const devtools = @import("../internal/devtools.zig"); const passkeys = @import("../internal/passkeys.zig"); const landing = @import("landing/mod.zig"); const router = @import("router.zig"); const stats = @import("stats.zig"); const telemetry = @import("../internal/telemetry.zig"); const http = std.http; const corsPreflight = http_api.corsPreflight; const json = http_api.json; const xrpcError = http_api.xrpcError; pub const Options = struct { host: []const u8 = "127.0.0.1", port: u16 = 2583, }; const App = struct { io: std.Io, pub const WebsocketHandler = atproto_sync.SubscribeReposClient; pub fn handle(app: *App, request: *httpz.Request, response: *httpz.Response) void { http_api.bindResponse(response); defer http_api.unbindResponse(); app.serveRequest(request) catch |err| { log.err("failed to serve {s}: {s}\n", .{ request.url.raw, @errorName(err) }); response.setStatus(.internal_server_error); response.body = "Internal Server Error"; }; } fn serveRequest(app: *App, request: *httpz.Request) !void { const route = router.route(request.method, request.url.raw); const start_us = telemetry.start(); var telemetry_route = route; var telemetry_class: telemetry.Class = if (route == .not_found) .unknown else .local; var telemetry_label = telemetryLabel(route, request); var handler_failed = true; const record_telemetry = route != .stats_page; defer if (record_telemetry) telemetry.record( telemetry_route, request.method, telemetry_class, telemetry_label, start_us, if (handler_failed) 500 else http_api.currentStatus(), handler_failed, ); log.debug("http {s} {s} route={s} start\n", .{ methodName(request), request.url.raw, @tagName(route) }); if (atproto_proxy.shouldProxy(request)) { telemetry_route = .proxy_xrpc; telemetry_class = .proxy; telemetry_label = stripQuery(request.url.raw); try atproto_proxy.xrpcProxy(request); handler_failed = false; return; } switch (route) { .cors_preflight => try corsPreflight(request), .root => try landing.serve(request), .api_docs => try docs.serve(request), .api_openapi => try docs.serveOpenApi(request), .stats_page => try stats.serve(request), .account_page => try account_page.page(request), .spaces_redirect => try account_page.redirectToSpaces(request), .favicon => try landing.serveFavicon(request), .og_image => try landing.serveOgImage(request), .health => try health(request), .did_json => try atproto_server.didJson(request), .oauth_protected_resource => try atproto_oauth.protectedResource(request), .oauth_authorization_server => try atproto_oauth.authorizationServer(request), .oauth_jwks => try atproto_oauth.jwks(request), .oauth_par => try atproto_oauth.par(request), .oauth_authorize => if (request.method == .GET) try atproto_oauth.authorizeGet(request) else try atproto_oauth.authorizePost(request), .oauth_token => try atproto_oauth.token(request), .oauth_introspect => try atproto_oauth.introspect(request), .oauth_revoke => try atproto_oauth.revoke(request), .oauth_passkey_options => try passkeys.loginStart(request), .oauth_passkey_finish => try passkeys.loginFinish(request), .passkeys_redirect => try passkeys.redirectToSecurity(request), .security_redirect => try account_page.redirectToSecurity(request), .admin_sessions_page => try passkeys.adminSessionsPage(request), .zds_list_account_sessions => try atproto_server.listSessions(request), .zds_list_admin_sessions => try atproto_server.adminListSessions(request), .atproto_did => try atproto_server.atprotoDid(request), .describe_server => try atproto_server.describeServer(request), .reserve_signing_key => try atproto_server.reserveSigningKey(request), .create_account => try atproto_server.createAccount(request), .create_invite_code => try atproto_server.createInviteCode(request), .create_invite_codes => try atproto_server.createInviteCodes(request), .get_account_invite_codes => try atproto_server.getAccountInviteCodes(request), .admin_update_subject_status => try atproto_server.updateSubjectStatus(request), .list_app_passwords => try atproto_server.listAppPasswords(request), .create_app_password => try atproto_server.createAppPassword(request), .revoke_app_password => try atproto_server.revokeAppPassword(request), .start_passkey_registration => try passkeys.xrpcStartRegistration(request), .finish_passkey_registration => try passkeys.xrpcFinishRegistration(request), .list_passkeys => try passkeys.xrpcList(request), .delete_passkey => try passkeys.xrpcDelete(request), .update_passkey => try passkeys.xrpcUpdate(request), .create_session => try atproto_server.createSession(request), .refresh_session => try atproto_server.refreshSession(request), .get_session => try atproto_server.getSession(request), .get_service_auth => try atproto_server.getServiceAuth(request), .activate_account => try atproto_server.activateAccount(request), .deactivate_account => try atproto_server.deactivateAccount(request), .request_email_confirmation => try atproto_server.requestEmailConfirmation(request), .confirm_email => try atproto_server.confirmEmail(request), .request_email_update => try atproto_server.requestEmailUpdate(request), .update_email => try atproto_server.updateEmail(request), .check_account_status => try atproto_server.checkAccountStatus(request), .app_preferences_get => try atproto_preferences.getPreferences(request), .app_preferences_put => try atproto_preferences.putPreferences(request), .repo_create_record => try atproto_repo.createRecord(request), .repo_put_record => try atproto_repo.putRecord(request), .repo_describe_repo => try atproto_repo.describeRepo(request), .repo_get_record => try atproto_repo.getRecord(request), .repo_list_records => try atproto_repo.listRecords(request), .repo_delete_record => try atproto_repo.deleteRecord(request), .repo_apply_writes => try atproto_repo.applyWrites(request), .repo_import_repo => try atproto_repo.importRepo(request), .repo_upload_blob => try atproto_repo.uploadBlob(app.io, request), .repo_list_missing_blobs => try atproto_repo.listMissingBlobs(request), .sync_get_blob => try atproto_sync.getBlob(request), .sync_get_repo => try atproto_sync.getRepo(request), .sync_get_latest_commit => try atproto_sync.getLatestCommit(request), .sync_list_repos => try atproto_sync.listRepos(request), .sync_list_blobs => try atproto_sync.listBlobs(request), .sync_subscribe_repos => try atproto_sync.subscribeRepos(request), .sync_get_repo_status => try atproto_sync.getRepoStatus(request), .sync_notify_of_update => try atproto_sync.notifyOfUpdate(request), .sync_request_crawl => try atproto_sync.requestCrawl(request), .identity_get_recommended_did_credentials => try atproto_identity.getRecommendedDidCredentials(request), .identity_request_plc_operation_signature => try atproto_identity.requestPlcOperationSignature(request), .identity_sign_plc_operation => try atproto_identity.signPlcOperation(request), .identity_submit_plc_operation => try atproto_identity.submitPlcOperation(request), .identity_resolve_handle => try atproto_identity.resolveHandle(request), .identity_update_handle => try atproto_identity.updateHandle(request), .permissioned_data => try atproto_space.dispatch(request), .devtools => try devtools.dispatch(request), .proxy_xrpc => unreachable, .not_found => try xrpcError(request, .not_found, "UnknownMethod", "Unknown XRPC method"), } handler_failed = false; } }; pub fn listen(io: std.Io, options: Options) !void { telemetry.init(io); var app: App = .{ .io = io }; var server = try httpz.Server(*App).init(io, std.heap.smp_allocator, .{ .address = .{ .ip = .{ .ip4 = try std.Io.net.Ip4Address.parse(options.host, options.port) } }, .request = .{ .max_body_size = 128 * 1024 * 1024, .lazy_read_size = 16 * 1024 * 1024, }, .workers = .{ .large_buffer_count = 4, .large_buffer_size = 64 * 1024, }, }, &app); defer { server.stop(); server.deinit(); } log.info("zds listening on http://{s}:{d}\n", .{ options.host, options.port }); try server.listen(); } fn health(request: *httpz.Request) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const body = try std.fmt.allocPrint( arena.allocator(), "{{\"version\":{f},\"status\":\"ok\"}}", .{std.json.fmt(build_options.version, .{})}, ); try json(request, .ok, body); } fn methodName(request: *const httpz.Request) []const u8 { return switch (request.method) { .GET => "GET", .HEAD => "HEAD", .POST => "POST", .PUT => "PUT", .PATCH => "PATCH", .DELETE => "DELETE", .OPTIONS => "OPTIONS", .CONNECT => "CONNECT", .OTHER => request.method_string, }; } fn telemetryLabel(route: router.Route, request: *const httpz.Request) []const u8 { if (route == .not_found) return stripQuery(request.url.raw); for (router.endpoints) |endpoint| { if (endpoint.route == route and methodMatches(endpoint.method, request.method)) return endpoint.path; } return @tagName(route); } fn methodMatches(expected: []const u8, actual: httpz.Method) bool { return switch (actual) { .GET => std.mem.eql(u8, expected, "GET"), .HEAD => std.mem.eql(u8, expected, "HEAD"), .POST => std.mem.eql(u8, expected, "POST"), else => false, }; } fn stripQuery(target: []const u8) []const u8 { const end = std.mem.indexOfScalar(u8, target, '?') orelse target.len; return target[0..end]; }