Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
72 kB · 1831 lines
Go
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832package config
import ( "context" "crypto/rsa" "crypto/x509" "encoding/json" "encoding/pem" "errors" "flag" "fmt" "io" "net" "os" "path/filepath" "runtime" "slices" "strconv" "strings" "time"
"math/rand/v2"
"github.com/lestrrat-go/jwx/v2/jwk" "github.com/livepeer/go-livepeer/cmd/livepeer/starter" "github.com/lmittmann/tint" slogGorm "github.com/orandin/slog-gorm" urfavecli "github.com/urfave/cli/v3" "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/crypto/aqpub" "stream.place/streamplace/pkg/integrations/discord/discordtypes" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/moderation" placestream "stream.place/streamplace/pkg/placestream" "stream.place/streamplace/pkg/renditions" "stream.place/streamplace/pkg/s3")
const SPDataDir = "$SP_DATA_DIR"const SegmentsDir = "segments"const ThumbnailsDir = "thumbnails"
type BuildFlags struct { Version string BuildTime int64 UUID string}
func (b BuildFlags) BuildTimeStr() string { ts := time.Unix(b.BuildTime, 0) return ts.UTC().Format(time.RFC3339)}
func (b BuildFlags) BuildTimeStrExpo() string { ts := time.Unix(b.BuildTime, 0) return ts.UTC().Format("2006-01-02T15:04:05.000Z")}
type CLI struct { AdminAccount string Build *BuildFlags DataDir string DBURL string DBMaxOpenConns int LocalDBURL string EthAccountAddr string EthKeystorePath string EthPassword string FirebaseServiceAccount string FirebaseServiceAccountFile string GitLabURL string HTTPAddr string HTTPInternalAddr string HTTPSAddr string RTMPAddr string RTMPSAddr string RTMPSAddonAddr string Secure bool NoMist bool IsolatedIngest bool MistAdminPort int MistHTTPPort int MistRTMPPort int SigningKeyPath string TAURL string TLSCertPath string TLSKeyPath string // ACME enables automatic TLS certificates from an ACME CA (Let\'s Encrypt) // stored in statedb; see pkg/acme. ACME bool ACMEEmail string ACMECA string ACMEDomains []string PKCS11ModulePath string PKCS11Pin string PKCS11TokenSlot string PKCS11TokenLabel string PKCS11TokenSerial string PKCS11KeypairLabel string PKCS11KeypairID string StreamerName string RelayHost string Debug map[string]map[string]int AllowedStreams []string WideOpen bool Peers []string Redirects map[string]string TestStream bool FrontendProxy string Frontend string PublicOAuth bool AppBundleID string NoFirehose bool PrintChat bool Color string LivepeerGatewayURL string TranscodeRenditions string LivepeerGateway bool WHIPTest string Thumbnail bool ExternalSigning bool RTMPServerAddon string TracingEndpoint string BroadcasterHost string XXDeprecatedPublicHost string ServerHost string RateLimitPerSecond int RateLimitBurst int RateLimitWebsocket int JWK jwk.Key AccessJWK jwk.Key ServiceAuthKey jwk.Key dataDirFlags []*string DiscordWebhooks []*discordtypes.Webhook AppleTeamID string AndroidCertFingerprint string Labelers []string AtprotoDID string LivepeerHelp bool PLCURL string ContentFilters *ContentFilters ModerationDir string DefaultRecommendedStreamers []string SQLLogging bool SentryDSN string LivepeerDebug bool Tickets []string IrohTopic string DID string DisableIrohRelay bool DevAccountCreds map[string]string StreamSessionTimeout time.Duration LegacySegmentCleaner bool SegmentArchiveRetention time.Duration Replicators []string WebsocketURL string BehindHTTPSProxy bool SegmentDebugDir string AdminDIDs []string Syndicate []string PlayerTelemetry bool PlaybackWorkerURL string Ingests *placestream.IngestGetIngestUrls_Output S3Endpoint string S3Bucket string S3AccessKeyID string S3SecretAccessKey string S3Region string VODCDNURL string VODCDNProvider string BunnyTokenAuthKey string BunnyPullZone string BunnyLogStorageZone string BunnyLogStorageEndpoint string BunnyLogStorageKey string CDNLogIngestInterval time.Duration LiveCDNURL string LiveCDNProvider string LiveBunnyTokenAuthKey string LiveCDNTokenTTL time.Duration DisableSyndication bool MuxlInitialMemoryMB int MuxlMaxMemoryMB int GamesAPIURL string GamesAPIClientKey string GamesAPIClientSecret string BetaInviteDID string ViewLogFlushInterval time.Duration ViewCountAggregateInterval time.Duration ViewCountAggregateLag time.Duration VODConcurrency int MaximumLiveBitrate int SweepConcurrency int SweepInterval time.Duration SweepBootDelay time.Duration DeepenRate int FirehoseReplayWindow time.Duration IndexDBConnections int}
// DefaultSweepInterval is how often the atproto sweep re-runs when// --sweep-interval is unset.//// The sweep's first pass over a repo that is up to date is a single// getLatestCommit, so this is a per-repo request budget: six hours means an// indexed account is asked about four times a day, and drift -- a gap in the// firehose, a span missed while this node was down -- is found and repaired// within that. Any lower buys hours of detection latency for a proportional// increase in traffic against every PDS on the network.const DefaultSweepInterval = 6 * time.Hour
// DefaultSweepConcurrency is how many PDS hosts the atproto backfill sweep// works on at once when --sweep-concurrency is unset or zero.//// The sweep shards its work by host and gives each host one worker, so this// bounds remote servers rather than repos: 32 of them is a few hundred requests// per second spread across the whole network, and no more than one walk (5-7// requests per second) against any single PDS.const DefaultSweepConcurrency = 32
// DefaultDeepenRate is how many history windows a node walks per minute when// --deepen-rate is unset.//// History acquisition is the one part of the sync engine nothing waits for: a// repo's recent records are indexed by its shallow sync in seconds, and// everything older is a background trickle. Running it flat out is what a node// does exactly once -- at boot, where it replays years of every account's chat// as fast as the network allows and buries the reconciliation the node actually// serves from. So it is paced instead: 60 windows a minute is one window a// second across the whole node, which a fresh 20k-repo index's full history// (four to five windows a repo, 100k of them) trickles in over roughly a day.// Deliberately: nothing is waiting for it. 0 removes the cap entirely, which is// what an operator uses to rush an initial build in place; negative means this// default.const DefaultDeepenRate = 60
// DefaultFirehoseReplayWindow is how stale a stored relay cursor may be before// this node stops trying to replay from it and tails the live edge instead.//// The firehose is a latency optimization, not the sync engine: the sweep's head// check asks every repo's host one question and repairs the ones that have// drifted, so a gap of hours costs a few thousand cheap requests spread across// hundreds of hosts. Replaying that same gap costs the relay a full-rate flood// of every commit on the network -- including the overwhelming majority from// repos this node has never heard of -- which is how a two-hour-old cursor once// buried a node under millions of queued events. Fifteen minutes is long enough// to cover an ordinary restart or deploy, where replay genuinely is the cheaper// answer, and short enough that anything worse is handed to the mechanism built// for it. 0 disables the cap and always replays from the stored cursor.const DefaultFirehoseReplayWindow = 15 * time.Minute
// DefaultIndexDBConnections is how many sqlite connections the index database// pool holds. See --index-db-connections; 1 is the fallback to the historical// single-connection arrangement.const DefaultIndexDBConnections = 8
// DefaultSweepBootDelay is how long a warm-index boot holds its first sweep// (and the deepener's first scan). An ordinary upgrade-restart's gap is healed// by the firehose replaying from the stored cursor, so the boot sweep is// insurance, not repair -- it only needs to wait out the restart churn itself:// streams reconnecting, the replay catching up, caches warming. Two minutes// does that. Deliberately NOT longer: with deepening rate-capped and sweep// passes reduced to checks and repairs, the pass is either harmless -- in// which case it may as well run while the operator who just deployed is still// watching the graphs -- or it is a problem, and a longer delay only schedules// the problem for the moment they have stopped looking. A fresh index sweeps// immediately regardless, and a cursor too stale to replay kicks its own// sweep, so the cases that genuinely need boot-time sweeping keep it.const DefaultSweepBootDelay = 2 * time.Minute
// ContentFilters represents the content filtering configurationtype ContentFilters struct { ContentWarnings struct { Enabled bool `json:"enabled"` BlockedWarnings []string `json:"blocked_warnings"` } `json:"content_warnings"` DistributionPolicy struct { Enabled bool `json:"enabled"` } `json:"distribution_policy"`}
const ( ReplicatorWebsocket string = "websocket" ReplicatorIroh string = "iroh")
var LivepeerFlagSet *flag.FlagSetvar LivepeerConfig starter.LivepeerConfig
func (cli *CLI) NewCommand(name string) *urfavecli.Command { cmd := &urfavecli.Command{ Name: name, Usage: "streamplace server", Flags: []urfavecli.Flag{ &urfavecli.StringFlag{ Name: "data-dir", Usage: "directory for keeping all streamplace data", Value: DefaultDataDir(), Destination: &cli.DataDir, Sources: urfavecli.EnvVars("SP_DATA_DIR"), }, &urfavecli.StringFlag{ Name: "http-addr", Usage: "Public HTTP address", Value: ":38080", Destination: &cli.HTTPAddr, Sources: urfavecli.EnvVars("SP_HTTP_ADDR"), }, &urfavecli.StringFlag{ Name: "http-internal-addr", Usage: "Private, admin-only HTTP address", Value: "127.0.0.1:39090", Destination: &cli.HTTPInternalAddr, Sources: urfavecli.EnvVars("SP_HTTP_INTERNAL_ADDR"), }, &urfavecli.StringFlag{ Name: "https-addr", Usage: "Public HTTPS address", Value: ":38443", Destination: &cli.HTTPSAddr, Sources: urfavecli.EnvVars("SP_HTTPS_ADDR"), }, &urfavecli.BoolFlag{ Name: "secure", Usage: "Run with HTTPS. Required for WebRTC output", Value: false, Destination: &cli.Secure, Sources: urfavecli.EnvVars("SP_SECURE"), }, &urfavecli.BoolFlag{ Name: "isolated-ingest", Usage: "Run each MKV/RTMP-push ingest in an isolated worker subprocess (fault isolation)", Value: true, Destination: &cli.IsolatedIngest, Sources: urfavecli.EnvVars("SP_ISOLATED_INGEST"), }, &urfavecli.StringFlag{ Name: "tls-cert", Usage: fmt.Sprintf(`Path to TLS certificate (default: "%s")`, filepath.Join(SPDataDir, "tls", "tls.crt")), Destination: &cli.TLSCertPath, Value: filepath.Join(SPDataDir, "tls", "tls.crt"), Sources: urfavecli.EnvVars("SP_TLS_CERT"), }, &urfavecli.StringFlag{ Name: "tls-key", Usage: fmt.Sprintf(`Path to TLS key (default: "%s")`, filepath.Join(SPDataDir, "tls", "tls.key")), Destination: &cli.TLSKeyPath, Value: filepath.Join(SPDataDir, "tls", "tls.key"), Sources: urfavecli.EnvVars("SP_TLS_KEY"), }, &urfavecli.BoolFlag{ Name: "acme", Usage: "With --secure, obtain and renew the node's TLS certificates automatically from an ACME CA (Let's Encrypt by default) instead of reading --tls-cert/--tls-key. Certificates are kept in the state database so every node in a station shares them. The CA must be able to reach this node on port 80 (HTTP-01) or 443 (TLS-ALPN-01) at --broadcaster-host and --server-host. Using this means you agree to the CA's terms of service.", Value: false, Destination: &cli.ACME, Sources: urfavecli.EnvVars("SP_ACME"), }, &urfavecli.StringFlag{ Name: "acme-email", Usage: "Contact email registered with the ACME CA (expiry warnings, account recovery)", Destination: &cli.ACMEEmail, Sources: urfavecli.EnvVars("SP_ACME_EMAIL"), }, &urfavecli.StringFlag{ Name: "acme-ca", Usage: `ACME directory URL, or "letsencrypt" / "letsencrypt-staging" (default: letsencrypt)`, Destination: &cli.ACMECA, Sources: urfavecli.EnvVars("SP_ACME_CA"), }, &urfavecli.StringFlag{ Name: "acme-domains", Usage: `comma-separated extra hostnames to hold certificates for, on top of --broadcaster-host and --server-host (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.ACMEDomains = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_ACME_DOMAINS"), }, &urfavecli.StringFlag{ Name: "signing-key", Usage: "Path to signing key for pushing OTA updates to the app", Destination: &cli.SigningKeyPath, Sources: urfavecli.EnvVars("SP_SIGNING_KEY"), }, &urfavecli.StringFlag{ Name: "db-url", Usage: "URL of the database to use for storing private streamplace state", Value: "sqlite://$SP_DATA_DIR/state.sqlite", Destination: &cli.DBURL, Sources: urfavecli.EnvVars("SP_DB_URL"), }, &urfavecli.IntFlag{ Name: "db-max-open-conns", Usage: "most connections this node holds to a Postgres state database (one of them is the advisory-lock connection); requests wait for a free one rather than opening more. Size it so nodes × this stays under the server's max_connections. Ignored for sqlite", Value: 30, Destination: &cli.DBMaxOpenConns, Sources: urfavecli.EnvVars("SP_DB_MAX_OPEN_CONNS"), }, &urfavecli.StringFlag{ Name: "admin-account", Usage: "ethereum account that administrates this streamplace node", Destination: &cli.AdminAccount, Sources: urfavecli.EnvVars("SP_ADMIN_ACCOUNT"), }, &urfavecli.StringFlag{ Name: "firebase-service-account", Usage: "Base64-encoded JSON string of a firebase service account key", Destination: &cli.FirebaseServiceAccount, Sources: urfavecli.EnvVars("SP_FIREBASE_SERVICE_ACCOUNT"), }, &urfavecli.StringFlag{ Name: "firebase-service-account-file", Usage: "Path to a JSON file containing a firebase service account key", Destination: &cli.FirebaseServiceAccountFile, Sources: urfavecli.EnvVars("SP_FIREBASE_SERVICE_ACCOUNT_FILE"), }, &urfavecli.StringFlag{ Name: "gitlab-url", Usage: "gitlab url for generating download links", Value: "https://git.stream.place/api/v4/projects/1", Destination: &cli.GitLabURL, Sources: urfavecli.EnvVars("SP_GITLAB_URL"), }, &urfavecli.StringFlag{ Name: "eth-keystore-path", Usage: fmt.Sprintf(`path to ethereum keystore (default: "%s")`, filepath.Join(SPDataDir, "keystore")), Destination: &cli.EthKeystorePath, Value: filepath.Join(SPDataDir, "keystore"), Sources: urfavecli.EnvVars("SP_ETH_KEYSTORE_PATH"), }, &urfavecli.StringFlag{ Name: "eth-account-addr", Usage: "ethereum account address to use (if keystore contains more than one)", Destination: &cli.EthAccountAddr, Sources: urfavecli.EnvVars("SP_ETH_ACCOUNT_ADDR"), }, &urfavecli.StringFlag{ Name: "eth-password", Usage: "password for encrypting keystore", Destination: &cli.EthPassword, Sources: urfavecli.EnvVars("SP_ETH_PASSWORD"), }, &urfavecli.StringFlag{ Name: "ta-url", Usage: "timestamp authority server for signing", Value: "http://timestamp.digicert.com", Destination: &cli.TAURL, Sources: urfavecli.EnvVars("SP_TA_URL"), }, &urfavecli.StringFlag{ Name: "pkcs11-module-path", Usage: "path to a PKCS11 module for HSM signing, for example /usr/lib/x86_64-linux-gnu/opensc-pkcs11.so", Destination: &cli.PKCS11ModulePath, Sources: urfavecli.EnvVars("SP_PKCS11_MODULE_PATH"), }, &urfavecli.StringFlag{ Name: "pkcs11-pin", Usage: "PIN for logging into PKCS11 token. if not provided, will be prompted interactively", Destination: &cli.PKCS11Pin, Sources: urfavecli.EnvVars("SP_PKCS11_PIN"), }, &urfavecli.StringFlag{ Name: "pkcs11-token-slot", Usage: "slot number of PKCS11 token (only use one of slot, label, or serial)", Destination: &cli.PKCS11TokenSlot, Sources: urfavecli.EnvVars("SP_PKCS11_TOKEN_SLOT"), }, &urfavecli.StringFlag{ Name: "pkcs11-token-label", Usage: "label of PKCS11 token (only use one of slot, label, or serial)", Destination: &cli.PKCS11TokenLabel, Sources: urfavecli.EnvVars("SP_PKCS11_TOKEN_LABEL"), }, &urfavecli.StringFlag{ Name: "pkcs11-token-serial", Usage: "serial number of PKCS11 token (only use one of slot, label, or serial)", Destination: &cli.PKCS11TokenSerial, Sources: urfavecli.EnvVars("SP_PKCS11_TOKEN_SERIAL"), }, &urfavecli.StringFlag{ Name: "pkcs11-keypair-label", Usage: "label of signing keypair on PKCS11 token", Destination: &cli.PKCS11KeypairLabel, Sources: urfavecli.EnvVars("SP_PKCS11_KEYPAIR_LABEL"), }, &urfavecli.StringFlag{ Name: "pkcs11-keypair-id", Usage: "id of signing keypair on PKCS11 token", Destination: &cli.PKCS11KeypairID, Sources: urfavecli.EnvVars("SP_PKCS11_KEYPAIR_ID"), }, &urfavecli.StringFlag{ Name: "app-bundle-id", Usage: "bundle id of an app that we facilitate oauth login for", Destination: &cli.AppBundleID, Sources: urfavecli.EnvVars("SP_APP_BUNDLE_ID"), }, &urfavecli.StringFlag{ Name: "streamer-name", Usage: "name of the person streaming from this streamplace node", Destination: &cli.StreamerName, Sources: urfavecli.EnvVars("SP_STREAMER_NAME"), }, &urfavecli.StringFlag{ Name: "dev-frontend-proxy", Usage: "(FOR DEVELOPMENT ONLY) proxy frontend requests to this address instead of using the bundled frontend", Destination: &cli.FrontendProxy, Sources: urfavecli.EnvVars("SP_DEV_FRONTEND_PROXY"), Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "false" { cli.FrontendProxy = "" return nil } cli.FrontendProxy = s return nil }, }, &urfavecli.StringFlag{ Name: "frontend", Usage: "which bundled frontend to serve: 'app' (legacy Expo) or 'web' (Vite)", Value: "app", Destination: &cli.Frontend, Sources: urfavecli.EnvVars("SP_FRONTEND"), }, &urfavecli.BoolFlag{ Name: "dev-public-oauth", Usage: "(FOR DEVELOPMENT ONLY) enable public oauth login for http://127.0.0.1 development", Value: false, Destination: &cli.PublicOAuth, Sources: urfavecli.EnvVars("SP_DEV_PUBLIC_OAUTH"), }, &urfavecli.StringFlag{ Name: "livepeer-gateway-url", Usage: "URL of the Livepeer Gateway to use for transcoding", Destination: &cli.LivepeerGatewayURL, Sources: urfavecli.EnvVars("SP_LIVEPEER_GATEWAY_URL"), }, &urfavecli.StringFlag{ Name: "transcode-renditions", Usage: "the renditions to transcode, comma-separated names from the ladder (1080p,720p,360p,240p,160p); only those smaller than the source are made. Default: the whole ladder", Destination: &cli.TranscodeRenditions, Sources: urfavecli.EnvVars("SP_TRANSCODE_RENDITIONS"), }, &urfavecli.BoolFlag{ Name: "livepeer-gateway", Usage: "enable embedded Livepeer Gateway", Value: false, Destination: &cli.LivepeerGateway, Sources: urfavecli.EnvVars("SP_LIVEPEER_GATEWAY"), }, &urfavecli.BoolFlag{ Name: "wide-open", Usage: "allow ALL streams to be uploaded to this node (not recommended for production)", Value: false, Destination: &cli.WideOpen, Sources: urfavecli.EnvVars("SP_WIDE_OPEN"), }, &urfavecli.StringFlag{ Name: "allowed-streams", Usage: `if set, only allow these addresses or atproto DIDs to upload to this node (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.AllowedStreams = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_ALLOWED_STREAMS"), }, &urfavecli.StringFlag{ Name: "peers", Usage: `other streamplace nodes to replicate to (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.Peers = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_PEERS"), }, &urfavecli.StringFlag{ Name: "redirects", Usage: `http 302s /path/one:/path/two,/path/three:/path/four (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } return json.Unmarshal([]byte(s), &cli.Redirects) }, Sources: urfavecli.EnvVars("SP_REDIRECTS"), }, &urfavecli.StringFlag{ Name: "debug", Usage: "modified log verbosity for specific functions or files in form func=ToHLS:3,file=gstreamer.go:4", Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.Debug = map[string]map[string]int{} pairs := strings.SplitSeq(s, ",") for pair := range pairs { scoreSplit := strings.Split(pair, ":") if len(scoreSplit) != 2 { return fmt.Errorf("invalid debug flag: %s", pair) } score, err := strconv.Atoi(scoreSplit[1]) if err != nil { return fmt.Errorf("invalid debug flag: %s", pair) } selectorSplit := strings.Split(scoreSplit[0], "=") if len(selectorSplit) != 2 { return fmt.Errorf("invalid debug flag: %s", pair) } _, ok := cli.Debug[selectorSplit[0]] if !ok { cli.Debug[selectorSplit[0]] = map[string]int{} } cli.Debug[selectorSplit[0]][selectorSplit[1]] = score } return nil }, Sources: urfavecli.EnvVars("SP_DEBUG"), }, &urfavecli.BoolFlag{ Name: "test-stream", Usage: "run a built-in test stream on boot", Value: false, Destination: &cli.TestStream, Sources: urfavecli.EnvVars("SP_TEST_STREAM"), }, &urfavecli.BoolFlag{ Name: "no-firehose", Usage: "disable the bluesky firehose", Value: false, Destination: &cli.NoFirehose, Sources: urfavecli.EnvVars("SP_NO_FIREHOSE"), }, &urfavecli.BoolFlag{ Name: "print-chat", Usage: "print chat messages to stdout", Value: false, Destination: &cli.PrintChat, Sources: urfavecli.EnvVars("SP_PRINT_CHAT"), }, &urfavecli.StringFlag{ Name: "whip-test", Usage: "run a WHIP self-test with the given parameters", Destination: &cli.WHIPTest, Sources: urfavecli.EnvVars("SP_WHIP_TEST"), }, &urfavecli.StringFlag{ Name: "relay-host", Usage: "comma-separated url(s) for relay firehose(s); ws://, wss://, or moqt:// (MoQ-over-QUIC) relays may be mixed. Subscribing to several relays survives any one going down (duplicate events are deduped). Our own PDS firehose is always included so locally-published records are indexed immediately", Value: "wss://bsky.network", Destination: &cli.RelayHost, Sources: urfavecli.EnvVars("SP_RELAY_HOST"), }, &urfavecli.StringFlag{ Name: "color", Usage: "'true' to enable colorized logging, 'false' to disable", Destination: &cli.Color, Sources: urfavecli.EnvVars("SP_COLOR"), }, &urfavecli.StringFlag{ Name: "broadcaster-host", Usage: "public host for the broadcaster group that this node is a part of (excluding https:// e.g. stream.place)", Destination: &cli.BroadcasterHost, Sources: urfavecli.EnvVars("SP_BROADCASTER_HOST"), }, &urfavecli.StringFlag{ Name: "public-host", Usage: "deprecated, use broadcaster-host or server-host instead as appropriate", Destination: &cli.XXDeprecatedPublicHost, Sources: urfavecli.EnvVars("SP_PUBLIC_HOST"), }, &urfavecli.StringFlag{ Name: "server-host", Usage: "public host for this particular physical streamplace node. defaults to broadcaster-host and only must be set for multi-node broadcasters", Destination: &cli.ServerHost, Sources: urfavecli.EnvVars("SP_SERVER_HOST"), }, &urfavecli.BoolFlag{ Name: "thumbnail", Usage: "enable thumbnail generation", Value: true, Destination: &cli.Thumbnail, Sources: urfavecli.EnvVars("SP_THUMBNAIL"), }, &urfavecli.StringFlag{ Name: "tracing-endpoint", Usage: "gRPC endpoint to send traces to", Destination: &cli.TracingEndpoint, Sources: urfavecli.EnvVars("SP_TRACING_ENDPOINT"), }, &urfavecli.IntFlag{ Name: "rate-limit-per-second", Usage: "rate limit for requests per second per ip", Value: 0, Destination: &cli.RateLimitPerSecond, Sources: urfavecli.EnvVars("SP_RATE_LIMIT_PER_SECOND"), }, &urfavecli.IntFlag{ Name: "rate-limit-burst", Usage: "rate limit burst for requests per ip", Value: 0, Destination: &cli.RateLimitBurst, Sources: urfavecli.EnvVars("SP_RATE_LIMIT_BURST"), }, &urfavecli.IntFlag{ Name: "rate-limit-websocket", Usage: "number of concurrent websocket connections allowed per ip", Value: 10, Destination: &cli.RateLimitWebsocket, Sources: urfavecli.EnvVars("SP_RATE_LIMIT_WEBSOCKET"), }, &urfavecli.IntFlag{ Name: "muxl-initial-memory-mb", Usage: "initial wasm linear memory pre-allocation per muxl instance, in MiB. higher avoids realloc churn at the cost of holding more memory upfront", Value: 50, Destination: &cli.MuxlInitialMemoryMB, Sources: urfavecli.EnvVars("SP_MUXL_INITIAL_MEMORY_MB"), }, &urfavecli.IntFlag{ Name: "muxl-max-memory-mb", Usage: "hard ceiling on wasm linear memory per muxl instance, in MiB. signing fails if a segment requires more than this", Value: 1024, Destination: &cli.MuxlMaxMemoryMB, Sources: urfavecli.EnvVars("SP_MUXL_MAX_MEMORY_MB"), }, &urfavecli.StringFlag{ Name: "rtmp-server-addon", Usage: "address of external RTMP server to forward streams to", Destination: &cli.RTMPServerAddon, Sources: urfavecli.EnvVars("SP_RTMP_SERVER_ADDON"), }, &urfavecli.StringFlag{ Name: "rtmps-addon-addr", Usage: "address to listen for RTMPS on the addon server", Value: ":1936", Destination: &cli.RTMPSAddonAddr, Sources: urfavecli.EnvVars("SP_RTMPS_ADDON_ADDR"), }, &urfavecli.StringFlag{ Name: "rtmps-addr", Usage: "address to listen for RTMPS connections (when --secure=true)", Value: ":1935", Destination: &cli.RTMPSAddr, Sources: urfavecli.EnvVars("SP_RTMPS_ADDR"), }, &urfavecli.StringFlag{ Name: "rtmp-addr", Usage: "address to listen for RTMP connections (when --secure=false)", Value: ":1935", Destination: &cli.RTMPAddr, Sources: urfavecli.EnvVars("SP_RTMP_ADDR"), }, &urfavecli.StringFlag{ Name: "discord-webhooks", Usage: `JSON array of Discord webhooks to send notifications to (default: "[]")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } return json.Unmarshal([]byte(s), &cli.DiscordWebhooks) }, Sources: urfavecli.EnvVars("SP_DISCORD_WEBHOOKS"), }, &urfavecli.StringFlag{ Name: "apple-team-id", Usage: "apple team id for deep linking", Destination: &cli.AppleTeamID, Sources: urfavecli.EnvVars("SP_APPLE_TEAM_ID"), }, &urfavecli.StringFlag{ Name: "android-cert-fingerprint", Usage: "android cert fingerprint for deep linking", Destination: &cli.AndroidCertFingerprint, Sources: urfavecli.EnvVars("SP_ANDROID_CERT_FINGERPRINT"), }, &urfavecli.StringFlag{ Name: "labelers", Usage: `did of labelers that this instance should subscribe to (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.Labelers = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_LABELERS"), }, &urfavecli.StringFlag{ Name: "atproto-did", Usage: "atproto did to respond to on /.well-known/atproto-did (default did:web:PUBLIC_HOST)", Destination: &cli.AtprotoDID, Sources: urfavecli.EnvVars("SP_ATPROTO_DID"), }, &urfavecli.StringFlag{ Name: "content-filters", Usage: `JSON content filtering rules (default: "{}")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } return json.Unmarshal([]byte(s), &cli.ContentFilters) }, Sources: urfavecli.EnvVars("SP_CONTENT_FILTERS"), }, &urfavecli.StringFlag{ Name: "moderation-dir", Usage: "directory containing additional .txt profanity wordlists to load at startup", Destination: &cli.ModerationDir, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { moderation.ModerationDir = s return nil }, Sources: urfavecli.EnvVars("SP_MODERATION_DIR"), }, &urfavecli.StringFlag{ Name: "default-recommended-streamers", Usage: `comma-separated list of streamer DIDs to recommend by default when no other recommendations are available (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.DefaultRecommendedStreamers = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_DEFAULT_RECOMMENDED_STREAMERS"), }, &urfavecli.BoolFlag{ Name: "livepeer-help", Usage: "print help for livepeer flags and exit", Value: false, Destination: &cli.LivepeerHelp, Sources: urfavecli.EnvVars("SP_LIVEPEER_HELP"), }, &urfavecli.StringFlag{ Name: "plc-url", Usage: "url of the plc directory", Value: "https://plc.directory", Destination: &cli.PLCURL, Sources: urfavecli.EnvVars("SP_PLC_URL"), }, &urfavecli.BoolFlag{ Name: "sql-logging", Usage: "enable sql logging", Value: false, Destination: &cli.SQLLogging, Sources: urfavecli.EnvVars("SP_SQL_LOGGING"), }, &urfavecli.StringFlag{ Name: "sentry-dsn", Usage: "sentry dsn for error reporting", Destination: &cli.SentryDSN, Sources: urfavecli.EnvVars("SP_SENTRY_DSN"), }, &urfavecli.StringFlag{ Name: "playback-worker-url", Usage: "URL of the Cloudflare playback router worker", Destination: &cli.PlaybackWorkerURL, Sources: urfavecli.EnvVars("SP_PLAYBACK_WORKER_URL"), }, &urfavecli.StringFlag{ Name: "games-api-url", Usage: "URL of the games.gamesgamesgamesgames API (e.g. http://localhost:3001)", Destination: &cli.GamesAPIURL, Sources: urfavecli.EnvVars("SP_GAMES_API_URL"), Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { cli.GamesAPIURL = strings.TrimRight(s, "/") return nil }, }, &urfavecli.StringFlag{ Name: "games-api-client-key", Usage: "Client key for authenticating with the games.gamesgamesgamesgames API", Destination: &cli.GamesAPIClientKey, Sources: urfavecli.EnvVars("SP_GAMES_API_CLIENT_KEY"), }, &urfavecli.StringFlag{ Name: "games-api-client-secret", Usage: "Client secret for authenticating with the games.gamesgamesgamesgames API", Destination: &cli.GamesAPIClientSecret, Sources: urfavecli.EnvVars("SP_GAMES_API_CLIENT_SECRET"), }, &urfavecli.BoolFlag{ Name: "livepeer-debug", Usage: "log livepeer segments to $SP_DATA_DIR/livepeer-debug", Value: false, Destination: &cli.LivepeerDebug, Sources: urfavecli.EnvVars("SP_LIVEPEER_DEBUG"), }, &urfavecli.StringFlag{ Name: "segment-debug-dir", Usage: "directory to log segment validation to", Destination: &cli.SegmentDebugDir, Sources: urfavecli.EnvVars("SP_SEGMENT_DEBUG_DIR"), }, &urfavecli.StringFlag{ Name: "tickets", Usage: `tickets to join the swarm with (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.Tickets = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_TICKETS"), }, &urfavecli.StringFlag{ Name: "iroh-topic", Usage: "topic to use for the iroh swarm (must be 32 bytes in hex)", Destination: &cli.IrohTopic, Sources: urfavecli.EnvVars("SP_IROH_TOPIC"), }, &urfavecli.BoolFlag{ Name: "disable-iroh-relay", Usage: "disable the iroh relay", Value: false, Destination: &cli.DisableIrohRelay, Sources: urfavecli.EnvVars("SP_DISABLE_IROH_RELAY"), }, &urfavecli.StringFlag{ Name: "dev-account-creds", Usage: `(FOR DEVELOPMENT ONLY) did=password pairs for logging into test accounts without oauth (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.DevAccountCreds = map[string]string{} pairs := strings.Split(s, ",") for _, pair := range pairs { parts := strings.Split(pair, "=") if len(parts) != 2 { return fmt.Errorf("invalid kv flag: %s", pair) } cli.DevAccountCreds[parts[0]] = parts[1] } return nil }, Sources: urfavecli.EnvVars("SP_DEV_ACCOUNT_CREDS"), }, &urfavecli.DurationFlag{ Name: "stream-session-timeout", Usage: "how long to wait before considering a stream inactive on this node?", Value: 60 * time.Second, Destination: &cli.StreamSessionTimeout, Sources: urfavecli.EnvVars("SP_STREAM_SESSION_TIMEOUT"), }, &urfavecli.IntFlag{ Name: "vod-concurrency", Usage: "number of VOD processing tasks to run in parallel on this node", Value: 2, Destination: &cli.VODConcurrency, Sources: urfavecli.EnvVars("SP_VOD_CONCURRENCY"), }, &urfavecli.IntFlag{ Name: "sweep-concurrency", Usage: "how many PDS hosts the atproto backfill sweep talks to at once. Work is sharded by host and each host is walked by one worker, so this is a count of remote servers, not of repos; 0 for the default", Value: DefaultSweepConcurrency, Destination: &cli.SweepConcurrency, Sources: urfavecli.EnvVars("SP_SWEEP_CONCURRENCY"), }, &urfavecli.IntFlag{ Name: "deepen-rate", Usage: "how many history windows per minute this node walks in the background. History deepening is decoupled from the sweep -- the sweep checks and repairs, this fetches the past -- and nothing on the node waits for it, so it is paced rather than run flat out; 0 removes the cap (an initial build in a hurry), negative for the default", Value: DefaultDeepenRate, Destination: &cli.DeepenRate, Sources: urfavecli.EnvVars("SP_DEEPEN_RATE"), }, &urfavecli.IntFlag{ Name: "index-db-connections", Usage: "how many sqlite connections the index database pool holds. More than one lets reads run beside a reindex under WAL; 1 restores the old slower-but-safer single-connection arrangement; 0 for the default", Value: DefaultIndexDBConnections, Destination: &cli.IndexDBConnections, Sources: urfavecli.EnvVars("SP_INDEX_DB_CONNECTIONS"), }, &urfavecli.DurationFlag{ Name: "sweep-boot-delay", Usage: "how long a node with a warm index waits after boot before its first sweep, so the sweep's reindexing does not compound the busiest minutes of a restart. A fresh (empty) index always sweeps immediately, as does a --no-firehose node (no replay heals its gap), and 0 sweeps immediately in every case", Value: DefaultSweepBootDelay, Destination: &cli.SweepBootDelay, Sources: urfavecli.EnvVars("SP_SWEEP_BOOT_DELAY"), }, &urfavecli.DurationFlag{ Name: "sweep-interval", Usage: "how often to re-run the atproto sweep, which asks every indexed repo's host whether our copy is still current and repairs the ones that are not. 0 disables re-running; the sweep at startup always happens", Value: DefaultSweepInterval, Destination: &cli.SweepInterval, Sources: urfavecli.EnvVars("SP_SWEEP_INTERVAL"), }, &urfavecli.DurationFlag{ Name: "firehose-replay-window", Usage: "how old a stored relay cursor may be and still be replayed from on connect. A cursor whose newest event is older than this is discarded and we tail the relay's live edge instead, leaving the gap for the sweep's head check to repair -- which is far cheaper than making the relay re-send every commit on the network. 0 always replays from the stored cursor", Value: DefaultFirehoseReplayWindow, Destination: &cli.FirehoseReplayWindow, Sources: urfavecli.EnvVars("SP_FIREHOSE_REPLAY_WINDOW"), }, &urfavecli.StringFlag{ Name: "maximum-live-bitrate", Usage: "maximum allowed live ingest bitrate, measured per emitted segment. Accepts a bits-per-second number or a decimal SI suffix — e.g. 30M, 30000k, or 30000000 (all 30 Mbps). A stream whose bitrate exceeds this (plus a 10% margin) is disconnected and the streamer is shown a problem. 0 = unlimited", Value: "0", Sources: urfavecli.EnvVars("SP_MAXIMUM_LIVE_BITRATE"), Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { v, err := ParseSI(s) if err != nil { return fmt.Errorf("invalid --maximum-live-bitrate: %w", err) } cli.MaximumLiveBitrate = int(v) return nil }, }, &urfavecli.BoolFlag{ Name: "legacy-segment-cleaner", Usage: "re-enable the legacy segment cleaner. shouldn't be needed but can be useful in cases where localdb is too big.", Value: false, Destination: &cli.LegacySegmentCleaner, Sources: urfavecli.EnvVars("SP_LEGACY_SEGMENT_CLEANER"), }, &urfavecli.DurationFlag{ Name: "segment-archive-retention", Usage: "how long to keep on-disk segment files before cleaning them up (durable copies live in S3/VOD). 0 disables cleanup.", Value: 1 * time.Hour, Destination: &cli.SegmentArchiveRetention, Sources: urfavecli.EnvVars("SP_SEGMENT_ARCHIVE_RETENTION"), }, &urfavecli.StringFlag{ Name: "replicators", Usage: "comma-separated list of replication protocols to use (websocket, iroh)", Value: ReplicatorWebsocket, Sources: urfavecli.EnvVars("SP_REPLICATORS"), Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s != "" { cli.Replicators = strings.Split(s, ",") } return nil }, }, &urfavecli.StringFlag{ Name: "websocket-url", Usage: "override the websocket (ws:// or wss://) url to use for replication (normally not necessary, used for testing)", Destination: &cli.WebsocketURL, Sources: urfavecli.EnvVars("SP_WEBSOCKET_URL"), }, &urfavecli.BoolFlag{ Name: "behind-https-proxy", Usage: "set to true if this node is behind an https proxy and we should report https URLs even though the node isn't serving HTTPS", Value: false, Destination: &cli.BehindHTTPSProxy, Sources: urfavecli.EnvVars("SP_BEHIND_HTTPS_PROXY"), }, &urfavecli.StringFlag{ Name: "admin-dids", Usage: `comma-separated list of DIDs that are authorized to modify branding and other admin operations (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.AdminDIDs = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_ADMIN_DIDS"), }, &urfavecli.StringFlag{ Name: "syndicate", Usage: `list of DIDs that we should rebroadcast ('*' for everybody) (default: "")`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } cli.Syndicate = strings.Split(s, ",") return nil }, Sources: urfavecli.EnvVars("SP_SYNDICATE"), }, &urfavecli.BoolFlag{ Name: "disable-syndication", Usage: `entirely disable syndication in both directions. useful for local development.`, Value: false, Destination: &cli.DisableSyndication, Sources: urfavecli.EnvVars("SP_DISABLE_SYNDICATION"), }, &urfavecli.BoolFlag{ Name: "player-telemetry", Usage: "enable player telemetry", Value: true, Destination: &cli.PlayerTelemetry, Sources: urfavecli.EnvVars("SP_PLAYER_TELEMETRY"), }, &urfavecli.StringFlag{ Name: "local-db-url", Usage: "URL of the local database to use for storing local data", Value: "sqlite://$SP_DATA_DIR/localdb.sqlite", Destination: &cli.LocalDBURL, Sources: urfavecli.EnvVars("SP_LOCAL_DB_URL"), }, &urfavecli.StringFlag{ Name: "ingests", Usage: `JSON array of ingests to return from place.stream.ingest.getIngestUrls. Default is auto-generated ingests for RTMP and WHIP`, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } return json.Unmarshal([]byte(s), &cli.Ingests) }, Sources: urfavecli.EnvVars("SP_INGESTS"), }, &urfavecli.StringFlag{ Name: "s3-endpoint", Usage: "S3-compatible endpoint URL for segment archival uploads", Destination: &cli.S3Endpoint, Sources: urfavecli.EnvVars("SP_S3_ENDPOINT"), }, &urfavecli.StringFlag{ Name: "s3-bucket", Usage: "S3 bucket name for segment archival uploads", Destination: &cli.S3Bucket, Sources: urfavecli.EnvVars("SP_S3_BUCKET"), }, &urfavecli.StringFlag{ Name: "s3-access-key-id", Usage: "S3 access key ID for segment archival uploads", Destination: &cli.S3AccessKeyID, Sources: urfavecli.EnvVars("SP_S3_ACCESS_KEY_ID"), }, &urfavecli.StringFlag{ Name: "s3-secret-access-key", Usage: "S3 secret access key for segment archival uploads", Destination: &cli.S3SecretAccessKey, Sources: urfavecli.EnvVars("SP_S3_SECRET_ACCESS_KEY"), }, &urfavecli.StringFlag{ Name: "s3-region", Usage: "S3 region (default: us-east-1)", Value: "us-east-1", Destination: &cli.S3Region, Sources: urfavecli.EnvVars("SP_S3_REGION"), }, &urfavecli.StringFlag{ Name: "vod-cdn-url", Usage: "Static CDN URL fronting the VOD blob store. When set, HLS playlists emit segment + init-segment URLs of the form <vod-cdn-url>/<cid>.mp4?did=...&sid=... instead of the self-hosted getVideoBlob endpoint. Omit for self-contained deployments.", Destination: &cli.VODCDNURL, Sources: urfavecli.EnvVars("SP_VOD_CDN_URL"), }, &urfavecli.StringFlag{ Name: "vod-cdn-provider", Usage: "Which CDN product sits at --vod-cdn-url, selecting URL signing + access-log ingestion. Empty means a plain static CDN over a public bucket: unsigned URLs, and no way to count segment views. Supported: bunny (configure with the --bunny-* flags).", Destination: &cli.VODCDNProvider, Sources: urfavecli.EnvVars("SP_VOD_CDN_PROVIDER"), }, &urfavecli.StringFlag{ Name: "bunny-token-auth-key", Usage: "bunny.net: the pull zone's Token Authentication key. When set, every blob URL in an HLS playlist is signed (token + expires query params) so the CDN only serves blobs this node handed out; the token is bound to the blob path and expires well after the VOD's duration. Requires --vod-cdn-provider=bunny.", Destination: &cli.BunnyTokenAuthKey, Sources: urfavecli.EnvVars("SP_BUNNY_TOKEN_AUTH_KEY"), }, &urfavecli.StringFlag{ Name: "bunny-pull-zone", Usage: "bunny.net: the pull zone's name, as it appears in the Permanent Log Storage path (pullzone-logs/<name>/...). Required with --bunny-log-storage-zone.", Destination: &cli.BunnyPullZone, Sources: urfavecli.EnvVars("SP_BUNNY_PULL_ZONE"), }, &urfavecli.StringFlag{ Name: "bunny-log-storage-zone", Usage: "bunny.net: the Edge Storage zone the pull zone's Permanent Log Storage writes into. When set, the node periodically ingests archived access logs so segment requests served by the CDN count toward views. Requires --vod-cdn-provider=bunny, --bunny-pull-zone and --bunny-log-storage-key.", Destination: &cli.BunnyLogStorageZone, Sources: urfavecli.EnvVars("SP_BUNNY_LOG_STORAGE_ZONE"), }, &urfavecli.StringFlag{ Name: "bunny-log-storage-endpoint", Usage: "bunny.net: Edge Storage API endpoint for the log storage zone's region (e.g. https://ny.storage.bunnycdn.com).", Value: "https://storage.bunnycdn.com", Destination: &cli.BunnyLogStorageEndpoint, Sources: urfavecli.EnvVars("SP_BUNNY_LOG_STORAGE_ENDPOINT"), }, &urfavecli.StringFlag{ Name: "bunny-log-storage-key", Usage: "bunny.net: access key (password) for the log storage zone. A read-only key is sufficient.", Destination: &cli.BunnyLogStorageKey, Sources: urfavecli.EnvVars("SP_BUNNY_LOG_STORAGE_KEY"), }, &urfavecli.DurationFlag{ Name: "cdn-log-ingest-interval", Usage: "How often to pull newly archived CDN access logs into the view-log store and re-aggregate the view-count windows they touch. Only meaningful when the --vod-cdn-provider has a log source configured. Set to 0 to disable.", Value: 15 * time.Minute, Destination: &cli.CDNLogIngestInterval, Sources: urfavecli.EnvVars("SP_CDN_LOG_INGEST_INTERVAL"), }, &urfavecli.StringFlag{ Name: "live-cdn-url", Usage: "CDN URL fronting this node's live HLS segments (a pull zone whose origin is this node's public URL). When set, live media playlists emit segment URLs of the form <live-cdn-url>/live/<did>/<track>/<seq>.m4s instead of the self-hosted getLiveSegment endpoint; playlists and init segments stay on the node. Typically a different domain from --vod-cdn-url, since the origin is the node rather than the blob store. Omit for self-contained deployments.", Destination: &cli.LiveCDNURL, Sources: urfavecli.EnvVars("SP_LIVE_CDN_URL"), }, &urfavecli.StringFlag{ Name: "live-cdn-provider", Usage: "Which CDN product sits at --live-cdn-url, selecting URL signing. Empty means a plain CDN: unsigned URLs. Supported: bunny (configure with --live-bunny-token-auth-key). Live view counting rides on the playlist requests the node keeps serving, so no access-log ingestion is needed here.", Destination: &cli.LiveCDNProvider, Sources: urfavecli.EnvVars("SP_LIVE_CDN_PROVIDER"), }, &urfavecli.StringFlag{ Name: "live-bunny-token-auth-key", Usage: "bunny.net: the live pull zone's Token Authentication key (a separate zone from VOD means a separate key). When set, every live segment URL in a media playlist is signed so the CDN only serves segments this node handed out. Requires --live-cdn-provider=bunny.", Destination: &cli.LiveBunnyTokenAuthKey, Sources: urfavecli.EnvVars("SP_LIVE_BUNNY_TOKEN_AUTH_KEY"), }, &urfavecli.DurationFlag{ Name: "live-cdn-token-ttl", Usage: "How long a signed live segment URL stays valid. Expiry is rounded to this interval so every playlist rendered within it carries identical URLs and the CDN can cache them; a URL is therefore valid for between one and two intervals. A live player refetches its playlist every few seconds, so this only needs to outlive the segment window.", Value: 5 * time.Minute, Destination: &cli.LiveCDNTokenTTL, Sources: urfavecli.EnvVars("SP_LIVE_CDN_TOKEN_TTL"), }, &urfavecli.StringFlag{ Name: "beta-invite-did", Usage: "DID of the atproto account whose place.stream.beta.invite records this node trusts. When set, uploading VODs requires an invite from that account; when empty, falls back to the --allowed-streams allowlist used by livestreaming.", Destination: &cli.BetaInviteDID, Sources: urfavecli.EnvVars("SP_BETA_INVITE_DID"), }, &urfavecli.DurationFlag{ Name: "view-log-flush-interval", Usage: "How often the view-log writer rotates its buffer to the VOD blob store. Set to 0 to disable view-event logging entirely (no view counts will be available downstream). Files land at view-logs/<server-did>/<window>.jsonl.gz alongside the VOD content blobs.", Value: 5 * time.Minute, Destination: &cli.ViewLogFlushInterval, Sources: urfavecli.EnvVars("SP_VIEW_LOG_FLUSH_INTERVAL"), }, &urfavecli.DurationFlag{ Name: "view-count-aggregate-interval", Usage: "How often a node tries to enqueue a view-count aggregation task. Buckets align on UTC multiples of this interval; deduplication via statedb's unique task-key constraint ensures only one node per bucket actually runs the aggregation. Set to 0 to disable aggregation (capture continues but no place.stream.media.viewCount records are published).", Value: 5 * time.Minute, Destination: &cli.ViewCountAggregateInterval, Sources: urfavecli.EnvVars("SP_VIEW_COUNT_AGGREGATE_INTERVAL"), }, &urfavecli.DurationFlag{ Name: "view-count-aggregate-lag", Usage: "How long the aggregator waits after a bucket closes before processing it, so all writers have time to flush their buffers. Should be at least one --view-log-flush-interval; default 2× that.", Value: 10 * time.Minute, Destination: &cli.ViewCountAggregateLag, Sources: urfavecli.EnvVars("SP_VIEW_COUNT_AGGREGATE_LAG"), }, &urfavecli.BoolFlag{ Name: "external-signing", Usage: "DEPRECATED, does nothing.", Value: true, }, &urfavecli.BoolFlag{ Name: "insecure", Usage: "DEPRECATED, does nothing.", Value: false, }, }, Before: func(ctx context.Context, cmd *urfavecli.Command) (context.Context, error) { return ctx, cli.Validate(cmd) }, }
// Add data dir flags cli.dataDirFlags = append(cli.dataDirFlags, &cli.DBURL) cli.dataDirFlags = append(cli.dataDirFlags, &cli.LocalDBURL) cli.dataDirFlags = append(cli.dataDirFlags, &cli.TLSCertPath) cli.dataDirFlags = append(cli.dataDirFlags, &cli.TLSKeyPath) cli.dataDirFlags = append(cli.dataDirFlags, &cli.EthKeystorePath)
if runtime.GOOS == "linux" { cmd.Flags = append(cmd.Flags, &urfavecli.BoolFlag{ Name: "no-mist", Usage: "Disable MistServer", Value: true, Destination: &cli.NoMist, Sources: urfavecli.EnvVars("SP_NO_MIST"), }) cmd.Flags = append(cmd.Flags, &urfavecli.IntFlag{ Name: "mist-admin-port", Usage: "MistServer admin port (internal use only)", Value: 14242, Destination: &cli.MistAdminPort, Sources: urfavecli.EnvVars("SP_MIST_ADMIN_PORT"), }) cmd.Flags = append(cmd.Flags, &urfavecli.IntFlag{ Name: "mist-rtmp-port", Usage: "MistServer RTMP port (internal use only)", Value: 11935, Destination: &cli.MistRTMPPort, Sources: urfavecli.EnvVars("SP_MIST_RTMP_PORT"), }) cmd.Flags = append(cmd.Flags, &urfavecli.IntFlag{ Name: "mist-http-port", Usage: "MistServer HTTP port (internal use only) — ingest pulls Mist's live fMP4 output from this port, so it must match the running Mist config (docker/mistserver.json uses 28080, the default here)", Value: 28080, Destination: &cli.MistHTTPPort, Sources: urfavecli.EnvVars("SP_MIST_HTTP_PORT"), })
}
LivepeerFlagSet = flag.NewFlagSet("livepeer", flag.ContinueOnError) LivepeerConfig = starter.NewLivepeerConfig(LivepeerFlagSet) LivepeerFlagSet.VisitAll(func(f *flag.Flag) { adapted := LivepeerFlags.CamelToSnake[f.Name] cmd.Flags = append(cmd.Flags, &urfavecli.StringFlag{ Name: fmt.Sprintf("livepeer.%s", adapted), Usage: f.Usage, Sources: urfavecli.EnvVars(fmt.Sprintf("SP_LIVEPEER_%s", adapted)), Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { return LivepeerFlagSet.Set(f.Name, s) }, }) })
return cmd}
var StreamplaceSchemePrefix = "streamplace://"
// OwnPublicURL is the URL this process's own public listener answers on.//// With --secure we terminate TLS ourselves: the real handler is on HTTPSAddr// and the HTTPAddr listener only serves 307 redirects to it (ServeHTTPRedirect),// so http://<HTTPAddr> is not an address anything can actually be fetched from// — a websocket dial there gets the redirect instead of a 101 upgrade.// --behind-https-proxy is the opposite case: the proxy terminates TLS and we// really do serve the handler as plain HTTP on HTTPAddr, so only cli.Secure// flips this.func (cli *CLI) OwnPublicURL() string { // No errors because we know it's valid from AddrFlag addr, scheme := cli.HTTPAddr, "http" if cli.Secure { addr, scheme = cli.HTTPSAddr, "https" } host, port, _ := net.SplitHostPort(addr)
ip := net.ParseIP(host) if host == "" || ip.IsUnspecified() { host = "127.0.0.1" } return fmt.Sprintf("%s://%s", scheme, net.JoinHostPort(host, port))}
func (cli *CLI) OwnInternalURL() string { // No errors because we know it's valid from AddrFlag host, port, _ := net.SplitHostPort(cli.HTTPInternalAddr)
ip := net.ParseIP(host) if ip.IsUnspecified() { host = "127.0.0.1" } addr := net.JoinHostPort(host, port) return fmt.Sprintf("http://%s", addr)}
func (cli *CLI) ParseSigningKey() (*rsa.PrivateKey, error) { bs, err := os.ReadFile(cli.SigningKeyPath) if err != nil { return nil, err } block, _ := pem.Decode(bs) if block == nil { return nil, fmt.Errorf("no RSA key found in signing key") } key, err := x509.ParsePKCS1PrivateKey(block.Bytes) if err != nil { return nil, err } return key, nil}
func RandomTrailer(length int) string { const charset = "abcdefghijklmnopqrstuvwxyz0123456789"
res := make([]byte, length) for i := 0; i < length; i++ { res[i] = charset[rand.IntN(len(charset))] } return string(res)}
func DefaultDataDir() string { home, err := os.UserHomeDir() if err != nil { // not fatal unless the user doesn't set one later return "" } return filepath.Join(home, ".streamplace")}
var GormLogger = slogGorm.New( slogGorm.WithHandler(tint.NewHandler(os.Stderr, &tint.Options{ TimeFormat: time.RFC3339, })), slogGorm.WithTraceAll(),)
func DisableSQLLogging() { GormLogger = slogGorm.New( slogGorm.WithHandler(tint.NewHandler(os.Stderr, &tint.Options{ TimeFormat: time.RFC3339, })), )}
func EnableSQLLogging() { GormLogger = slogGorm.New( slogGorm.WithHandler(tint.NewHandler(os.Stderr, &tint.Options{ TimeFormat: time.RFC3339, })), slogGorm.WithTraceAll(), )}
func (cli *CLI) Validate(cmd *urfavecli.Command) error { if cli.DataDir == "" { return fmt.Errorf("could not determine default data dir (no $HOME) and none provided, please set --data-dir") } if cli.LivepeerGateway && cli.LivepeerGatewayURL != "" { return fmt.Errorf("defining both livepeer-gateway and livepeer-gateway-url doesn't make sense. do you want an embedded gateway or an external one?") } if _, err := renditions.Ladder(cli.TranscodeRenditions); err != nil { return fmt.Errorf("--transcode-renditions: %w", err) } if cli.LivepeerGateway { log.MonkeypatchStderr() // Livepeer gateway configuration will be handled in the caller cli.LivepeerGatewayURL = "http://127.0.0.1:8935" } for _, dest := range cli.dataDirFlags { *dest = strings.Replace(*dest, SPDataDir, cli.DataDir, 1) } if !cli.SQLLogging { DisableSQLLogging() } else { EnableSQLLogging() } if cli.XXDeprecatedPublicHost != "" && cli.BroadcasterHost == "" { log.Warn(context.Background(), "public-host is deprecated, use broadcaster-host or server-host instead as appropriate") cli.BroadcasterHost = cli.XXDeprecatedPublicHost } if cli.ServerHost == "" && cli.BroadcasterHost != "" { cli.ServerHost = cli.BroadcasterHost } if cli.PublicOAuth { log.Warn(context.Background(), "--dev-public-oauth is set, this is not recommended for production") } if cli.FirebaseServiceAccount != "" && cli.FirebaseServiceAccountFile != "" { return fmt.Errorf("defining both firebase-service-account and firebase-service-account-file doesn't make sense. do you want a base64-encoded string or a file?") } if cli.FirebaseServiceAccountFile != "" { bs, err := os.ReadFile(cli.FirebaseServiceAccountFile) if err != nil { return err } cli.FirebaseServiceAccount = string(bs) } // Set default replicator if none specified if len(cli.Replicators) == 0 { cli.Replicators = []string{ReplicatorWebsocket} } if err := cli.validateVODCDN(); err != nil { return err } if err := cli.validateLiveCDN(); err != nil { return err } return nil}
// validateLiveCDN is validateVODCDN for the live segment CDN: a key// without its provider, or a provider without a URL, is refused.func (cli *CLI) validateLiveCDN() error { switch cli.LiveCDNProvider { case "": if cli.LiveBunnyTokenAuthKey != "" { return fmt.Errorf("--live-bunny-token-auth-key is set but --live-cdn-provider is not; set --live-cdn-provider=bunny") } case "bunny": if cli.LiveCDNURL == "" { return fmt.Errorf("--live-cdn-provider=bunny requires --live-cdn-url (the pull zone hostname)") } default: return fmt.Errorf("unknown --live-cdn-provider %q (supported: bunny)", cli.LiveCDNProvider) } if cli.LiveCDNURL != "" && cli.LiveCDNTokenTTL <= 0 { return fmt.Errorf("--live-cdn-token-ttl must be positive") } return nil}
// validateVODCDN refuses half-configured CDN setups: provider flags// without their provider, a provider without a CDN URL, or a log// source missing one of its parts. Anything that passes here is a// combination cdn.FromConfig can assemble without surprises.func (cli *CLI) validateVODCDN() error { bunnyFlags := cli.BunnyTokenAuthKey != "" || cli.BunnyPullZone != "" || cli.BunnyLogStorageZone != "" || cli.BunnyLogStorageKey != "" switch cli.VODCDNProvider { case "": if bunnyFlags { return fmt.Errorf("--bunny-* flags are set but --vod-cdn-provider is not; set --vod-cdn-provider=bunny") } case "bunny": if cli.VODCDNURL == "" { return fmt.Errorf("--vod-cdn-provider=bunny requires --vod-cdn-url (the pull zone hostname)") } logFlags := cli.BunnyPullZone != "" || cli.BunnyLogStorageZone != "" || cli.BunnyLogStorageKey != "" if logFlags && (cli.BunnyPullZone == "" || cli.BunnyLogStorageZone == "" || cli.BunnyLogStorageKey == "") { return fmt.Errorf("bunny log ingestion needs all of --bunny-pull-zone, --bunny-log-storage-zone and --bunny-log-storage-key") } default: return fmt.Errorf("unknown --vod-cdn-provider %q (supported: bunny)", cli.VODCDNProvider) } return nil}
func (cli *CLI) DataFilePath(fpath []string) string { if cli.DataDir == "" { panic("no data dir configured") } // windows does not like colons safe := []string{} for _, p := range fpath { safe = append(safe, strings.ReplaceAll(p, ":", "-")) } fpath = append([]string{cli.DataDir}, safe...) fdpath := filepath.Join(fpath...) return fdpath}
// does a file exist in our data dir?func (cli *CLI) DataFileExists(fpath []string) (bool, error) { ddpath := cli.DataFilePath(fpath) _, err := os.Stat(ddpath) if err == nil { return true, nil } if errors.Is(err, os.ErrNotExist) { return false, nil } return false, err}
// write a file to our data dirfunc (cli *CLI) DataFileWrite(fpath []string, r io.Reader, overwrite bool) error { fd, err := cli.DataFileCreate(fpath, overwrite) if err != nil { return err } defer fd.Close() _, err = io.Copy(fd, r) if err != nil { return err }
return nil}
// create a file in our data dir. don't forget to close it!func (cli *CLI) DataFileCreate(fpath []string, overwrite bool) (*os.File, error) { ddpath := cli.DataFilePath(fpath) if !overwrite { exists, err := cli.DataFileExists(fpath) if err != nil { return nil, err } if exists { return nil, fmt.Errorf("refusing to overwrite file that exists: %s", ddpath) } } if len(fpath) > 1 { dirs, _ := filepath.Split(ddpath) err := os.MkdirAll(dirs, os.ModePerm) if err != nil { return nil, fmt.Errorf("error creating subdirectories for %s: %w", ddpath, err) } } return os.Create(ddpath)}
// get a path to a segment file in our databasefunc (cli *CLI) SegmentFilePath(user string, file string) (string, error) { ext := filepath.Ext(file) base := strings.TrimSuffix(file, ext) aqt, err := aqtime.FromString(base) if err != nil { return "", err } fname := fmt.Sprintf("%s%s", aqt.FileSafeString(), ext) yr, mon, day, hr, min, _, _ := aqt.Parts() return cli.DataFilePath([]string{SegmentsDir, user, yr, mon, day, hr, min, fname}), nil}
// get a path to a segment file in our databasefunc (cli *CLI) HLSDir(user string) (string, error) { return cli.DataFilePath([]string{SegmentsDir, "hls", user}), nil}
// create a segment file in our databasefunc (cli *CLI) SegmentFileCreate(user string, aqt aqtime.AQTime, ext string) (*os.File, error) { fname := fmt.Sprintf("%s.%s", aqt.FileSafeString(), ext) yr, mon, day, hr, min, _, _ := aqt.Parts() return cli.DataFileCreate([]string{SegmentsDir, user, yr, mon, day, hr, min, fname}, false)}
// ThumbnailFilePath returns the path to a user's current thumbnail. There is a// single, continually-overwritten thumbnail per user. The user is a DID// (e.g. did:plc:...); DataFilePath strips the colons so the filename is safe on// Windows.func (cli *CLI) ThumbnailFilePath(user string) string { return cli.DataFilePath([]string{ThumbnailsDir, fmt.Sprintf("%s.jpg", user)})}
// ThumbnailModTime returns the modification time of a user's thumbnail and// whether it exists. The mod time doubles as a "last seen live" signal.func (cli *CLI) ThumbnailModTime(user string) (time.Time, bool) { fi, err := os.Stat(cli.ThumbnailFilePath(user)) if err != nil { return time.Time{}, false } return fi.ModTime(), true}
// ThumbnailWrite atomically (re)writes a user's thumbnail. The image is written// to a temp file via the supplied function and renamed into place, so readers// (and PDS uploads) never observe a half-written thumbnail.func (cli *CLI) ThumbnailWrite(user string, write func(io.Writer) error) error { final := cli.ThumbnailFilePath(user) dir := filepath.Dir(final) if err := os.MkdirAll(dir, os.ModePerm); err != nil { return fmt.Errorf("error creating thumbnail dir %s: %w", dir, err) } tmp, err := os.CreateTemp(dir, "thumb-*.jpg") if err != nil { return err } defer os.Remove(tmp.Name()) // no-op once the rename below succeeds if err := write(tmp); err != nil { tmp.Close() return err } if err := tmp.Close(); err != nil { return err } return os.Rename(tmp.Name(), final)}
// read a file from our data dirfunc (cli *CLI) DataFileRead(fpath []string, w io.Writer) error { ddpath := cli.DataFilePath(fpath)
fd, err := os.Open(ddpath) if err != nil { return err } _, err = io.Copy(w, fd) if err != nil { return err }
return nil}
func (cli *CLI) HasMist() bool { return runtime.GOOS == "linux"}
// type for comma-separated ethereum addressesfunc (cli *CLI) AddressSliceFlag(name, defaultValue, usage string, dest *[]aqpub.Pub) urfavecli.Flag { *dest = []aqpub.Pub{} usage = fmt.Sprintf(`%s (default: "%s")`, usage, defaultValue)
return &urfavecli.StringFlag{ Name: name, Usage: usage, Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { if s == "" { return nil } strs := strings.Split(s, ",") for _, str := range strs { pub, err := aqpub.FromHexString(str) if err != nil { return err } *dest = append(*dest, pub) } return nil }, Sources: urfavecli.EnvVars(fmt.Sprintf("SP_%s", strings.ToUpper(strings.ReplaceAll(name, "-", "_")))), }}
func (cli *CLI) StreamIsAllowed(did string) error { if cli.WideOpen { return nil } // if the user set no test streams, anyone can stream openServer := len(cli.AllowedStreams) == 0 || (cli.TestStream && len(cli.AllowedStreams) == 1) // but only valid atproto accounts! did:key is only allowed for our local test stream isDIDKey := strings.HasPrefix(did, constants.DID_KEY_PREFIX) if openServer && !isDIDKey { return nil } if slices.Contains(cli.AllowedStreams, did) { return nil } return fmt.Errorf("user is not allowed to stream")}
func (cli *CLI) BroadcasterDID() string { return fmt.Sprintf("did:web:%s", cli.BroadcasterHost)}
func (cli *CLI) ServerDID() string { if cli.ServerHost == "" { return cli.BroadcasterDID() } return fmt.Sprintf("did:web:%s", cli.ServerHost)}
func (cli *CLI) HasHTTPS() bool { return cli.Secure || cli.BehindHTTPSProxy}
func (cli *CLI) DumpDebugSegment(ctx context.Context, name string, r io.Reader) { if cli.SegmentDebugDir == "" { return } go func() { err := os.MkdirAll(cli.SegmentDebugDir, 0755) if err != nil { log.Error(ctx, "failed to create debug directory", "error", err) return } now := aqtime.FromTime(time.Now()) outFile := filepath.Join(cli.SegmentDebugDir, fmt.Sprintf("%s-%s", now.FileSafeString(), strings.ReplaceAll(name, ":", "-"))) fd, err := os.Create(outFile) if err != nil { log.Error(ctx, "failed to create debug file", "error", err) return } defer fd.Close() _, err = io.Copy(fd, r) if err != nil { log.Error(ctx, "failed to copy debug file", "error", err) return } log.Log(ctx, "wrote debug file", "path", outFile) }()}
func (cli *CLI) S3Configured() bool { return cli.S3Endpoint != "" && cli.S3Bucket != "" && cli.S3AccessKeyID != "" && cli.S3SecretAccessKey != ""}
// S3Config assembles an s3.Config from the CLI's S3 flags.func (cli *CLI) S3Config() s3.Config { return s3.Config{ Endpoint: cli.S3Endpoint, Bucket: cli.S3Bucket, AccessKeyID: cli.S3AccessKeyID, SecretAccessKey: cli.S3SecretAccessKey, Region: cli.S3Region, }}
// SetS3Config applies an s3.Config to the CLI's S3 fields — the inverse of// S3Config, for processes (ingest workers) that receive the S3 destination over// a handshake instead of from flags.func (cli *CLI) SetS3Config(c s3.Config) { cli.S3Endpoint = c.Endpoint cli.S3Bucket = c.Bucket cli.S3AccessKeyID = c.AccessKeyID cli.S3SecretAccessKey = c.SecretAccessKey cli.S3Region = c.Region}
// DebugRecordingFile is the write target returned by DebugRecordingCreate: an// *os.File on local disk, or an S3 upload that commits on Close. Name() reports// the destination (path or object key) for logging.type DebugRecordingFile interface { io.WriteCloser Name() string}
// DebugRecordingCreate opens a write target for a debug recording (RTMP/MKV// dumps, WHIP rtcrec sessions). When S3 is configured the recording streams to// an S3 object at the key formed by joining fpath with "/" (so the bucket// mirrors the on-disk debug-recordings/<did>/<file> layout); otherwise it falls// back to a local file under DataDir — the dev default. The returned value must// be Closed to finalize (Close commits the S3 upload). overwrite only affects// the local-disk path (S3 puts always overwrite).func (cli *CLI) DebugRecordingCreate(ctx context.Context, fpath []string, contentType string, overwrite bool) (DebugRecordingFile, error) { if cli.S3Configured() { key := strings.Join(fpath, "/") // The recording outlives the ingest session's ctx: Close commits the upload // during teardown, after that ctx is typically cancelled — a cancelled ctx // here would abort the upload and lose the object. Callers bound the commit // with their own finalize waits instead. return s3.NewUploadWriter(context.WithoutCancel(ctx), s3.NewClient(cli.S3Config()), cli.S3Bucket, key, contentType) } return cli.DataFileCreate(fpath, overwrite)}
func (cli *CLI) ShouldSyndicate(did string) bool { if cli.DisableSyndication { return false } for _, d := range cli.Syndicate { if d == "*" { return true } if d == did { return true } } return false}
// RenditionLadder is the rendition ladder this node transcodes// (--transcode-renditions), validated at startup.func (cli *CLI) RenditionLadder() []renditions.Rendition { ladder, err := renditions.Ladder(cli.TranscodeRenditions) if err != nil { return renditions.DesiredRenditions } return ladder}