1+ // Copied from: https://github.com/zeppelin-social/atproto/blob/main/services/bsky/api.js
2+
3+ // @ts -check
4+ /* eslint-env node */
5+ /* eslint-disable import/order */
6+
7+ 'use strict'
8+
9+ const dd = require ( 'dd-trace' )
10+
11+ dd . tracer
12+ . init ( )
13+ . use ( 'http2' , {
14+ client : true , // calls into dataplane
15+ server : false ,
16+ } )
17+ . use ( 'express' , {
18+ hooks : {
19+ request : ( span , req ) => {
20+ maintainXrpcResource ( span , req )
21+ } ,
22+ } ,
23+ } )
24+
25+ // modify tracer in order to track calls to dataplane as a service with proper resource names
26+ const DATAPLANE_PREFIX = '/bsky.Service/'
27+ const origStartSpan = dd . tracer . _tracer . startSpan
28+ dd . tracer . _tracer . startSpan = function ( name , options ) {
29+ if (
30+ name !== 'http.request' ||
31+ options ?. tags ?. component !== 'http2' ||
32+ ! options ?. tags ?. [ 'http.url' ]
33+ ) {
34+ return origStartSpan . call ( this , name , options )
35+ }
36+ const uri = new URL ( options . tags [ 'http.url' ] )
37+ if ( ! uri . pathname . startsWith ( DATAPLANE_PREFIX ) ) {
38+ return origStartSpan . call ( this , name , options )
39+ }
40+ options . tags [ 'service.name' ] = 'dataplane-bsky'
41+ options . tags [ 'resource.name' ] = uri . pathname . slice ( DATAPLANE_PREFIX . length )
42+ return origStartSpan . call ( this , name , options )
43+ }
44+
45+ // Tracer code above must come before anything else
46+ const path = require ( 'node:path' )
47+ const assert = require ( 'node:assert' )
48+ const cluster = require ( 'node:cluster' )
49+ const { Secp256k1Keypair } = require ( '@atproto/crypto' )
50+ const bsky = require ( '@atproto/bsky' ) // import all bsky features
51+
52+ const appview = async ( ) => {
53+ const env = getEnv ( )
54+ const config = bsky . ServerConfig . readEnv ( )
55+ assert ( env . serviceSigningKey , 'must set BSKY_SERVICE_SIGNING_KEY' )
56+ assert ( env . dbPostgresUrl , 'must set BSKY_DB_POSTGRES_URL' )
57+ const signingKey = await Secp256k1Keypair . import ( env . serviceSigningKey )
58+
59+ const db = new bsky . Database ( {
60+ url : env . dbPostgresUrl ,
61+ schema : env . dbPostgresSchema ,
62+ poolSize : env . dbPoolSize ,
63+ } )
64+
65+ // ends: involve logics in packages/dev-env/src/bsky.ts <<<<<<<<<<<<<
66+
67+ assert ( env . bsyncPort , 'must set BSKY_BSYNC_PORT' )
68+ assert ( env . dataplanePort , 'must set BSKY_DATAPLANE_PORT' )
69+
70+
71+ const bsync = await bsky . MockBsync . create ( db , env . bsyncPort )
72+
73+ const dataplane = await bsky . DataPlaneServer . create (
74+ db ,
75+ env . dataplanePort ,
76+ config . didPlcUrl ,
77+ )
78+
79+ const migrationDb = new bsky . Database ( {
80+ url : env . dbPostgresUrl ,
81+ schema : env . dbPostgresSchema ,
82+ } )
83+ await migrationDb . migrateToLatestOrThrow ( )
84+ await migrationDb . close ( )
85+
86+ const server = bsky . BskyAppView . create ( { config, signingKey } )
87+
88+ assert ( env . repoProvider , 'must set BSKY_REPO_PROVIDER' )
89+
90+ const sub = new bsky . RepoSubscription ( {
91+ service : env . repoProvider ,
92+ db,
93+ idResolver : dataplane . idResolver ,
94+ } )
95+
96+
97+ await server . start ( )
98+
99+ sub . start ( )
100+ // Graceful shutdown (see also https://aws.amazon.com/blogs/containers/graceful-shutdowns-with-ecs/)
101+ const shutdown = async ( ) => {
102+ await sub . destroy ( )
103+ await server . destroy ( )
104+ await dataplane . destroy ( )
105+ await bsync . destroy ( )
106+ await db . close ( )
107+ }
108+ process . on ( 'SIGTERM' , shutdown )
109+ process . on ( 'disconnect' , shutdown ) // when clustering
110+ }
111+
112+ const getEnv = ( ) => ( {
113+ serviceSigningKey : process . env . BSKY_SERVICE_SIGNING_KEY || undefined ,
114+ dbPostgresUrl : process . env . BSKY_DB_POSTGRES_URL || undefined ,
115+ dbPostgresSchema : process . env . BSKY_DB_POSTGRES_SCHEMA || undefined ,
116+ dbPoolSize : maybeParseInt ( process . env . BSKY_DB_POOL_SIZE ) || undefined ,
117+ dataplanePort : maybeParseInt ( process . env . BSKY_DATAPLANE_PORT ) || undefined ,
118+ bsyncPort : maybeParseInt ( process . env . BSKY_BSYNC_PORT ) || undefined ,
119+ migration : process . env . ENABLE_MIGRATIONS === 'true' || undefined ,
120+ repoProvider : process . env . BSKY_REPO_PROVIDER || undefined ,
121+ } )
122+
123+ const maybeParseInt = ( str ) => {
124+ if ( ! str ) return
125+ const int = parseInt ( str , 10 )
126+ if ( isNaN ( int ) ) return
127+ return int
128+ }
129+
130+ const maintainXrpcResource = ( span , req ) => {
131+ // Show actual xrpc method as resource rather than the route pattern
132+ if ( span && req . originalUrl ?. startsWith ( '/xrpc/' ) ) {
133+ span . setTag (
134+ 'resource.name' ,
135+ [
136+ req . method ,
137+ path . posix . join ( req . baseUrl || '' , req . path || '' , '/' ) . slice ( 0 , - 1 ) , // Ensures no trailing slash
138+ ]
139+ . filter ( Boolean )
140+ . join ( ' ' ) ,
141+ )
142+ }
143+ }
144+
145+ appview ( ) . catch ( console . error )
0 commit comments