Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
49 kB · 1511 lines
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512package cmd
import ( "bytes" "context" "crypto/rand" "crypto/tls" "encoding/json" "errors" "flag" "fmt" "io" "net/url" "os" "os/signal" "runtime" "runtime/pprof" "slices" "strconv" "stream.place/streamplace/pkg/acme" "stream.place/streamplace/pkg/branding" "strings" "syscall" "time"
"github.com/bluesky-social/indigo/carstore" "github.com/ethereum/go-ethereum/common/hexutil" "github.com/livepeer/go-livepeer/cmd/livepeer/starter" "github.com/streamplace/oatproxy/pkg/oatproxy" urfavecli "github.com/urfave/cli/v3" "stream.place/streamplace/pkg/aqhttp" "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/cdn/providers" "stream.place/streamplace/pkg/comatproto" "stream.place/streamplace/pkg/director" "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/pkg/ingestframe" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/localdb" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/muxl" "stream.place/streamplace/pkg/notifications" "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/replication/iroh_replicator" "stream.place/streamplace/pkg/replication/websocketrep" "stream.place/streamplace/pkg/rtmps" "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/storage" "stream.place/streamplace/pkg/upload" "stream.place/streamplace/pkg/viewlog" "stream.place/streamplace/pkg/vod"
"github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" awss3 "github.com/aws/aws-sdk-go-v2/service/s3"
_ "github.com/go-gst/go-glib/glib" _ "github.com/go-gst/go-gst/gst" "stream.place/streamplace/pkg/api" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/model")
// Additional jobs that can be injected by platformstype jobFunc func(ctx context.Context, cli *config.CLI) error
// parse the CLI and fire up an streamplace node!func start(build *config.BuildFlags, platformJobs []jobFunc) error { iroh_streamplace.InitLogging()
cli := config.CLI{Build: build} app := cli.NewCommand("streamplace") app.Usage = "decentralized live streaming platform" app.Version = build.Version app.Commands = []*urfavecli.Command{ makeSelfTestCommand(build), makeVODTestCommand(build), makeStreamCommand(build), makeIngestWorkerCommand(build), makeRTMPPushWorkerCommand(build), makeLiveCommand(build), makeWhepCommand(build), makeWhipCommand(build), makeCombineCommand(build), makeSplitCommand(build), makeLivepeerCommand(build), makeMigrateCommand(build), makeMigrateStateCommand(), makeSyncCommand(build), makeBrandingCommand(build), makeValidateConfigCommand(build), makeE2eCommand(build), } // Add the verbosity flag // app.Flags = append(app.Flags, &urfavecli.StringFlag{ // Name: "v", // Usage: "log verbosity level", // Value: "3", // }) app.Before = func(ctx context.Context, cmd *urfavecli.Command) (context.Context, error) { // validate-config must run wherever the binary runs, media stack or // not: diagnosing a crashloop with a broken gstreamer install is // exactly when you want pure config validation. if cmd.Name == "validate-config" { return ctx, nil } // Run self-test before starting selfTest := cmd.Name == "self-test" err := media.RunSelfTest(ctx) if err != nil { if selfTest { fmt.Println(err.Error()) os.Exit(1) } else { retryCount, _ := strconv.Atoi(os.Getenv("STREAMPLACE_SELFTEST_RETRY")) if retryCount >= 3 { log.Error(ctx, "gstreamer self-test failed 3 times, giving up", "error", err) return ctx, err } log.Log(ctx, "error in gstreamer self-test, attempting recovery", "error", err, "retry", retryCount+1) os.Setenv("STREAMPLACE_SELFTEST_RETRY", strconv.Itoa(retryCount+1)) err := syscall.Exec(os.Args[0], os.Args[1:], os.Environ()) if err != nil { log.Error(ctx, "error in gstreamer self-test, could not restart", "error", err) return ctx, err } panic("invalid code path: exec succeeded but we're still here???") } } return ctx, nil } app.Action = func(ctx context.Context, cmd *urfavecli.Command) error { return runMain(ctx, build, platformJobs, cmd, &cli) }
return app.Run(context.Background(), os.Args)}
func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFunc, cmd *urfavecli.Command, cli *config.CLI) error { _ = flag.Set("logtostderr", "true") vFlag := flag.Lookup("v")
err := cli.Validate(cmd) if err != nil { return err }
err = flag.CommandLine.Parse(nil) if err != nil { return err } verbosity := cmd.String("v") _ = vFlag.Value.Set(verbosity) log.SetColorLogger(cli.Color) ctx = log.WithDebugValue(ctx, cli.Debug)
log.Log(ctx, "streamplace", "version", build.Version, "buildTime", build.BuildTimeStr(), "uuid", build.UUID, "runtime.GOOS", runtime.GOOS, "runtime.GOARCH", runtime.GOARCH, "runtime.Version", runtime.Version())
muxl.Configure( uint64(cli.MuxlInitialMemoryMB)*1024*1024, uint64(cli.MuxlMaxMemoryMB)*1024*1024, )
signer, err := createSigner(ctx, cli) if err != nil { return err }
if len(os.Args) > 1 && os.Args[1] == "migrate" { return statedb.Migrate(cli) }
spmetrics.Version.WithLabelValues(build.Version).Inc() if cli.LivepeerHelp { lpFlags := flag.NewFlagSet("livepeer", flag.ContinueOnError) _ = starter.NewLivepeerConfig(lpFlags) lpFlags.VisitAll(func(f *flag.Flag) { adapted := config.ToSnakeCase(f.Name) fmt.Printf(" -%s\n", fmt.Sprintf("livepeer.%s", adapted)) usage := fmt.Sprintf(" %s", f.Usage) if f.DefValue != "" { usage = fmt.Sprintf("%s (default %s)", usage, f.DefValue) } fmt.Printf(" %s\n", usage) }) return nil }
aqhttp.UserAgent = fmt.Sprintf("streamplace/%s", build.Version)
err = os.MkdirAll(cli.DataDir, os.ModePerm) if err != nil { return fmt.Errorf("error creating streamplace dir at %s:%w", cli.DataDir, err) }
ldb, err := localdb.MakeDB(cli.LocalDBURL) if err != nil { return err }
mod, err := model.MakeDBConns(cli.DataFilePath([]string{"index"}), cli.IndexDBConnections) if err != nil { return err } var fbNotifier notifications.FirebaseNotifier if cli.FirebaseServiceAccount != "" { fbNotifier, err = notifications.MakeFirebaseNotifier(ctx, cli.FirebaseServiceAccount) if err != nil { return err } }
group, ctx := TimeoutGroupWithContext(ctx)
out := carstore.SQLiteStore{} err = out.Open(":memory:") if err != nil { return err } // The notifier is assembled after the DB exists, because the Web Push // notifier needs VAPID keys that are persisted in the Config table. The // queue processor nil-checks the notifier, so the brief window is safe. state, err := statedb.MakeDB(ctx, cli, nil, mod) if err != nil { return err }
// Build the Web Push notifier from VAPID keys generated/stored in the DB. vapidKeys, err := state.EnsureVAPIDKeys(ctx) if err != nil { return err } webNotifier := notifications.NewWebPushNotifier(vapidKeys, "") noter := notifications.NewMultiNotifier(fbNotifier, webNotifier) state.SetNotifier(noter) handle, err := atproto.MakeLexiconRepo(ctx, cli, mod, state) if err != nil { return err } defer handle.Close()
serverHandle, err := atproto.MakeServerRepo(ctx, cli, state) if err != nil { return err } defer serverHandle.Close()
jwk, err := state.EnsureJWK(ctx, "jwk") if err != nil { return err } cli.JWK = jwk
accessJWK, err := state.EnsureJWK(ctx, "access-jwk") if err != nil { return err } cli.AccessJWK = accessJWK
serviceAuthKey, err := state.EnsureServiceAuthKey(ctx) if err != nil { return err } cli.ServiceAuthKey = serviceAuthKey
b := bus.NewBus() atsync := &atproto.ATProtoSynchronizer{ CLI: cli, Model: mod, StatefulDB: state, Noter: noter, Bus: b, } mm, err := media.MakeMediaManager(ctx, cli, signer, mod, b, atsync, ldb) if err != nil { return err } // Every new playback session counts toward the streamer's running view // total, filed under their current livestream record. mm.SetViewRecorder(func(streamer string) { uri := "" if ls, err := mod.GetLatestLivestreamForRepo(streamer); err == nil && ls != nil { uri = ls.URI } if _, err := state.AddStreamView(ctx, streamer, uri); err != nil { log.Warn(ctx, "failed to count a view", "streamer", streamer, "err", err) } }) if cli.IsolatedIngest && !media.IngestIsolationSupported() { // The worker transport needs Unix fd-passing + Setsid (Linux today); fall // back to in-process ingest elsewhere rather than break. log.Log(ctx, "isolated ingest not supported on this platform; using in-process ingest", "goos", runtime.GOOS) cli.IsolatedIngest = false } if cli.IsolatedIngest { // Reconnect to any ingest workers still running from before this restart // and drain whatever they buffered while we were down (zero-downtime). mm.ResumeDetachedWorkers(ctx) }
ms, err := media.MakeMediaSigner(ctx, cli, cli.StreamerName, signer, mod) if err != nil { return err }
var clientMetadata *oatproxy.OAuthClientMetadata var host string if cli.PublicOAuth { u, err := url.Parse(cli.OwnPublicURL()) if err != nil { return err } host = u.Host clientMetadata = &oatproxy.OAuthClientMetadata{ Scope: atproto.OAuthString, ClientName: "Streamplace", RedirectURIs: []string{ fmt.Sprintf("%s/login", cli.OwnPublicURL()), fmt.Sprintf("%s/api/app-return", cli.OwnPublicURL()), }, } } else { host = cli.BroadcasterHost clientMetadata = &oatproxy.OAuthClientMetadata{ Scope: atproto.OAuthString, ClientName: "Streamplace", RedirectURIs: []string{ fmt.Sprintf("https://%s/login", cli.BroadcasterHost), fmt.Sprintf("https://%s/api/app-return", cli.BroadcasterHost), }, } }
op := oatproxy.New(&oatproxy.Config{ Host: host, CreateOAuthSession: state.CreateOAuthSession, UpdateOAuthSession: state.UpdateOAuthSession, GetOAuthSession: state.LoadOAuthSession, Lock: state.GetNamedLock, Scope: atproto.OAuthString, UpstreamJWK: cli.JWK, DownstreamJWK: cli.AccessJWK, ClientMetadata: clientMetadata, Public: cli.PublicOAuth, }) state.OATProxy = op // The synchronizer uses stored sessions to act as a streamer (e.g. ending // the source livestream when a teleport fires); without this it would // nil-deref the first time a teleport arrives for a logged-in streamer. atsync.OATProxy = op
// Sync every repo we know about, at boot and every --sweep-interval after: // a repair pass for repos left half-indexed by a previous run, and a head // check that finds the ones that drifted while we were not listening. // Nothing below depends on it, so it runs in the background off the serve // context -- shutdown cancels it -- and the node is up and serving in the // meantime. Started only now that atsync is fully wired (OATProxy above): // the sweep can index records and schedule callbacks that use every field. go atsync.SweepForever(ctx) // History is fetched by its own loop rather than by those passes: it takes // days where reconciliation takes minutes, nothing on the node waits for // it, and at full speed it is what buries a boot. It trickles in at // --deepen-rate windows a minute for as long as this node runs. go atsync.DeepenForever(ctx) // Trusted verifiers' records predate this node; pull them at boot and // hourly, the firehose keeps them current in between. go atsync.SeedVerificationsForever(ctx) go atsync.SeedLabelsForever(ctx)
var replicator replication.Replicator = nil if slices.Contains(cli.Replicators, config.ReplicatorIroh) { exists, err := cli.DataFileExists([]string{"iroh-kv-secret"}) if err != nil { return err } if !exists { secret := make([]byte, 32) _, err := rand.Read(secret) if err != nil { return fmt.Errorf("failed to generate random secret: %w", err) } err = cli.DataFileWrite([]string{"iroh-kv-secret"}, bytes.NewReader(secret), true) if err != nil { return err } } buf := bytes.Buffer{} err = cli.DataFileRead([]string{"iroh-kv-secret"}, &buf) if err != nil { return err } secret := buf.Bytes() var topic []byte if cli.IrohTopic != "" { topic, err = hexutil.Decode("0x" + cli.IrohTopic) if err != nil { return err } } replicator, err = iroh_replicator.NewSwarm(ctx, cli, secret, topic, mm, b, mod) if err != nil { return err } } if slices.Contains(cli.Replicators, config.ReplicatorWebsocket) { replicator = websocketrep.NewWebsocketReplicator(b, mod, mm, state) }
d := director.NewDirector(mm, mod, cli, b, op, state, replicator, ldb, atsync) um, err := upload.New(ctx, cli, state) if err != nil { return err } vodStore, err := makeVODStore(ctx, cli) if err != nil { return fmt.Errorf("make vod store: %w", err) } viewLog, err := makeViewLog(ctx, cli, vodStore, ldb) if err != nil { return fmt.Errorf("make view log: %w", err) } if viewLog != nil { group.Go(func() error { viewLog.Run(ctx) return nil }) defer func() { if err := viewLog.Close(); err != nil { log.Error(ctx, "view log close", "error", err) } }() } state.SetVODProcessor(func(ctx context.Context, t statedb.VODProcessTask) (string, error) { // Labeler enforcement: an account banned after starting an // upload (but before processing) doesn't get a video published. // The playback gates would hide it regardless, but skipping here // avoids the wasted transcode and a dead record. labels, err := mod.GetActiveLabels(t.RepoDID) if err != nil { return "", fmt.Errorf("vod-process: check account labels: %w", err) } if atproto.IsBanned(labels...) { return "", fmt.Errorf("vod-process: account %s is banned; skipping upload %s", t.RepoDID, t.UploadID) } return vod.ProcessVOD(ctx, cli, state, vodStore, vod.Input{ UploadID: t.UploadID, RepoDID: t.RepoDID, MimeType: t.MimeType, Filename: t.Filename, Size: t.Size, Backend: t.Backend, Location: t.Location, }) }) // Live-to-VOD finalize. Resolves the streamer's live signing key here // (where the model is in scope) so pkg/vod stays free of pkg/model, then // concatenates the recorded MUXL objects into a VOD. state.SetLivestreamVODFinalizer(func(ctx context.Context, t statedb.FinalizeLivestreamVODTask) (string, error) { signingKey, err := resolveLiveSigningKey(mod, t.RepoDID) if err != nil { return "", fmt.Errorf("finalize-livestream-vod: resolve signing key: %w", err) } return vod.FinalizeLivestreamVOD(ctx, cli, state, vodStore, vod.FinalizeInput{ UploadID: t.UploadID, RepoDID: t.RepoDID, LivestreamURI: t.LivestreamURI, LivestreamURIs: t.LivestreamURIs, SigningKey: signingKey, }) }) // Publishes the video record right after a finalize the operator asked // to publish (the finalize route), with the streamer's stored session. state.SetVideoPublisher(func(ctx context.Context, t statedb.FinalizeLivestreamVODTask) (string, string, error) { if t.Publish == nil { return "", "", nil } if vodStore == nil { return "", "", fmt.Errorf("finalize-livestream-vod: no VOD store configured") } return vod.PublishVideo(ctx, state, vodStore, t.RepoDID, t.UploadID, t.Publish.Record()) }) // View-count aggregator runs the log → record pipeline for one // window. Same function-pointer pattern as the VOD processor so // statedb stays free of viewlog's transitive deps. The scheduler // goroutine fires per --view-count-aggregate-interval; statedb's // unique TaskKey makes the cross-node race a no-op for losers. if vodStore != nil && cli.ViewCountAggregateInterval > 0 { // Resolver: for a blob CID, return strongRefs of every // place.stream.media.track record whose muxlTrack lives in // that blob, keyed by in-container trackId. Used by the // aggregator to attribute bytes/duration to the right track // record (the streamer's original track records or, later, // user-contributed transcript/transcode tracks). fetchTrackRefs := func(ctx context.Context, cid string) (map[string]comatproto.RepoStrongRef, error) { rows, err := mod.GetMediaTracksByBlob(ctx, cid) if err != nil { return nil, err } out := make(map[string]comatproto.RepoStrongRef, len(rows)) for _, row := range rows { rec, err := row.ToRecord() if err != nil { log.Warn(ctx, "viewlog refs: decode track record", "uri", row.URI, "error", err) continue } if rec.Track.MediaDefs_MuxlTrack == nil { continue } tid := rec.Track.MediaDefs_MuxlTrack.TrackId if tid == "" { continue } out[tid] = comatproto.RepoStrongRef{ LexiconTypeID: "com.atproto.repo.strongRef", Uri: row.URI, Cid: row.CID, } } return out, nil } state.SetViewCountAggregator(func(ctx context.Context, t statedb.ViewCountAggregateTask) error { return viewlog.RunAggregation(ctx, viewlog.RunAggregationInput{ Store: vodStore, CLI: cli, WindowStart: t.WindowStart, WindowEnd: t.WindowEnd, ReadMargin: 2 * cli.ViewLogFlushInterval, FetchTrackRefs: fetchTrackRefs, }) }) group.Go(func() error { return viewlog.ScheduleAggregations(ctx, state, viewlog.ScheduleConfig{ Interval: cli.ViewCountAggregateInterval, Lag: cli.ViewCountAggregateLag, }) }) } if err := wireCDNLogIngest(ctx, cli, state, vodStore, viewLog, ldb, group); err != nil { return err } a, err := api.MakeStreamplaceAPI(cli, mod, state, noter, mm, ms, b, atsync, d, op, ldb, um, vodStore, viewLog) if err != nil { return err }
ctx = log.WithLogValues(ctx, "version", build.Version)
group.Go(func() error { return handleSignals(ctx) })
group.Go(func() error { return state.ProcessQueue(ctx, cli.VODConcurrency) })
if cli.TracingEndpoint != "" { group.Go(func() error { return startTelemetry(ctx, cli.TracingEndpoint) }) }
if cli.ACME && !cli.Secure { log.Warn(ctx, "--acme has no effect without --secure; TLS is terminated elsewhere") } if cli.Secure { var rtmpsTLS *tls.Config if cli.ACME { mgr, err := acme.New(ctx, cli, state) if err != nil { return err } a.ACME = mgr rtmpsTLS = mgr.TLSConfig() group.Go(func() error { return mgr.Manage(ctx) }) } group.Go(func() error { return a.ServeHTTPS(ctx) }) group.Go(func() error { return a.ServeHTTPRedirect(ctx) }) if cli.RTMPServerAddon != "" { group.Go(func() error { return rtmps.ServeRTMPSAddon(ctx, cli, rtmpsTLS) }) } group.Go(func() error { return a.ServeRTMPS(ctx, cli) }) } else { group.Go(func() error { return a.ServeHTTP(ctx) }) group.Go(func() error { return a.ServeRTMP(ctx) }) }
group.Go(func() error { return a.ServeInternalHTTP(ctx) })
if !cli.NoFirehose { group.Go(func() error { return atsync.StartFirehose(ctx) }) }
// Make sure the beta-invite issuer's repo is registered so the // firehose path indexes its place.stream.beta.invite records. The // first call also backfills any invites that were published before // we came online; subsequent runs are cached and no-op. if cli.BetaInviteDID != "" { go func() { if _, err := atsync.SyncBlueskyRepoCached(ctx, cli.BetaInviteDID); err != nil { log.Error(ctx, "failed to sync beta-invite issuer repo; gating will rely on the firehose alone", "did", cli.BetaInviteDID, "err", err) } }() } for _, labeler := range cli.Labelers { group.Go(func() error { return atsync.StartLabelerFirehose(ctx, labeler) }) }
group.Go(func() error { return storage.StartSegmentCleaner(ctx, ldb, cli) })
if cli.LegacySegmentCleaner { group.Go(func() error { return ldb.StartSegmentCleaner(ctx) }) }
group.Go(func() error { return replicator.Start(ctx, cli) })
if cli.LivepeerGateway { // make a file to make sure the directory exists fd, err := cli.DataFileCreate([]string{"livepeer", "gateway", "empty"}, true) if err != nil { return err } fd.Close() if err != nil { return err } group.Go(func() error { err = GoLivepeer(ctx, config.LivepeerFlagSet) if err != nil { return err } // livepeer returns nil on error, so we need to check if we're responsible if ctx.Err() == nil { return fmt.Errorf("livepeer exited") } return nil }) }
group.Go(func() error { return d.Start(ctx) })
if cli.TestStream { atkey, err := atproto.ParsePubKey(signer.Public()) if err != nil { return err } did := atkey.DIDKey() testMediaSigner, err := media.MakeMediaSigner(ctx, cli, did, signer, mod) if err != nil { return err } err = mod.UpdateIdentity(&model.Identity{ ID: testMediaSigner.Pub().String(), Handle: "stream-self-tester", DID: "", }) if err != nil { return err } cli.AllowedStreams = append(cli.AllowedStreams, did) a.Aliases["self-test"] = did group.Go(func() error { return mm.TestSource(ctx, testMediaSigner) })
// Start a test stream that will run intermittently if err != nil { return err } atkey2, err := atproto.ParsePubKey(signer.Public()) if err != nil { return err } did2 := atkey2.DIDKey() intermittentMediaSigner, err := media.MakeMediaSigner(ctx, cli, did2, signer, mod) if err != nil { return err } err = mod.UpdateIdentity(&model.Identity{ ID: intermittentMediaSigner.Pub().String(), Handle: "stream-intermittent-tester", DID: "", }) if err != nil { return err } cli.AllowedStreams = append(cli.AllowedStreams, did2) a.Aliases["intermittent-self-test"] = did2
group.Go(func() error { for { // Start intermittent stream intermittentCtx, cancel := context.WithCancel(ctx) done := make(chan struct{}) go func() { _ = mm.TestSource(intermittentCtx, intermittentMediaSigner) close(done) }() // Stream ON for 15 seconds time.Sleep(15 * time.Second) // Stop stream cancel() <-done // Wait for TestSource to exit // Stream OFF for 15 seconds time.Sleep(15 * time.Second) } }) }
for _, job := range platformJobs { group.Go(func() error { return job(ctx, cli) }) }
if cli.WHIPTest != "" { group.Go(func() error { // Parse WHIPTest string using the whip command's flag parser whipCmd := makeWhipCommand(build) args := strings.Split(cli.WHIPTest, " ") err := whipCmd.Run(ctx, append([]string{"streamplace", "whip"}, args...)) log.Warn(ctx, "WHIP test complete, sleeping for 3 seconds and shutting down gstreamer") time.Sleep(time.Second * 3) // gst.Deinit() log.Warn(ctx, "gst deinit complete, exiting") return err }) }
return group.Wait()}
// makeVODStore picks the blob.Store backing VOD output for this// process. S3 if it's configured (production / multi-node); otherwise// the local DataDir. Either way the same Store is used to read the// user upload AND write the content-addressed VOD output — uploads// land under "uploads/" and content blobs land under "blobs/<cid>.mp4"// / "blobs/<cid>.json" (content-agnostic, since the blob doesn't know// what kind of video it's for).//// FileStore is rooted at DataDir so it can see the upload manager's// "uploads/<id>" tree alongside the content-addressed "blobs/" tree;// S3Store is rooted at the configured bucket for the same reason.//// Mirrors upload.New's multi-node-requires-S3 invariant: in single-node// file mode the upload and the produced VOD share local disk; in// multi-node S3 mode any station can pick up a queued VOD task and// land output in shared storage.func makeVODStore(ctx context.Context, cli *config.CLI) (blob.Store, error) { if cli.S3Configured() { s3client := awss3.New(awss3.Options{ Region: cli.S3Region, Credentials: credentials.NewStaticCredentialsProvider( cli.S3AccessKeyID, cli.S3SecretAccessKey, "", ), BaseEndpoint: aws.String(cli.S3Endpoint), UsePathStyle: true, }) log.Log(ctx, "VOD store: S3", "bucket", cli.S3Bucket) return blob.NewS3Store(s3client, cli.S3Bucket), nil } root := cli.DataFilePath(nil) log.Log(ctx, "VOD store: file", "root", root) return blob.NewFileStore(root)}
// wireCDNLogIngest installs the CDN access-log ingester when the// configured --vod-cdn-provider has a log source, and warns loudly// when a CDN is configured without one: segment fetches then bypass// the node entirely and every VOD's view count silently reads zero.func wireCDNLogIngest(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, vodStore blob.Store, viewLog *viewlog.Writer, ldb localdb.LocalDB, group *TimeoutGroup) error { provider, err := providers.FromConfig(cli) if err != nil { return err } if provider == nil { return nil } if provider.Logs == nil { log.Warn(ctx, "VOD CDN is configured without an access-log source: segment requests served by the CDN are invisible to view counting, so view counts for CDN-served VODs will be zero. Configure the provider's log ingestion (e.g. --bunny-log-storage-zone) or accept the gap.", "vod_cdn_url", cli.VODCDNURL, "provider", cli.VODCDNProvider) return nil } if vodStore == nil || cli.ViewCountAggregateInterval <= 0 { log.Warn(ctx, "VOD CDN log source configured but view-count aggregation is off; not ingesting CDN logs") return nil } if cli.CDNLogIngestInterval <= 0 { log.Log(ctx, "cdn log ingest: disabled (cdn-log-ingest-interval=0)") return nil } var salts *viewlog.SaltManager if viewLog != nil { salts = viewLog.Salts() } else { salts = viewlog.NewSaltManager(ldb) } sourceName := provider.Name state.SetCDNLogIngester(func(ctx context.Context) error { _, err := viewlog.RunIngest(ctx, viewlog.IngestInput{ Store: vodStore, Source: provider.Logs, SourceName: sourceName, Salts: salts, Cursor: state, Window: cli.ViewCountAggregateInterval, Reaggregate: func(ctx context.Context, partID string, start, end time.Time) error { _, err := state.EnqueueTask(ctx, statedb.TaskViewCountAggregate, statedb.ViewCountAggregateTask{WindowStart: start, WindowEnd: end}, statedb.WithTaskKey(viewlog.ReaggregateTaskKey(start, end, partID))) return err }, }) return err }) group.Go(func() error { return viewlog.ScheduleIngest(ctx, state, cli.CDNLogIngestInterval) }) log.Log(ctx, "cdn log ingest: enabled", "provider", sourceName, "interval", cli.CDNLogIngestInterval) return nil}
// makeViewLog returns the configured view-event log writer, or nil if// the operator has disabled it (--view-log-flush-interval=0) or there's// no place to write logs to (no VOD store). The writer reuses the VOD// blob.Store under a `view-logs/<server-did>/` prefix; if the operator// runs S3-backed VOD, view logs land in the same bucket alongside the// content blobs and pick up the bucket's lifecycle policy for free.func makeViewLog(ctx context.Context, cli *config.CLI, vodStore blob.Store, ldb localdb.LocalDB) (*viewlog.Writer, error) { if cli.ViewLogFlushInterval <= 0 { log.Log(ctx, "view log: disabled (view-log-flush-interval=0)") return nil, nil } if vodStore == nil { log.Log(ctx, "view log: no VOD store wired; skipping") return nil, nil } w, err := viewlog.NewWriter(viewlog.Config{ Store: vodStore, NodeDID: cli.ServerDID(), FlushAfter: cli.ViewLogFlushInterval, Salts: viewlog.NewSaltManager(ldb), }) if err != nil { return nil, err } log.Log(ctx, "view log: enabled", "flush_interval", cli.ViewLogFlushInterval, "node_did", cli.ServerDID()) return w, nil}
var ErrCaughtSignal = errors.New("caught signal")
func handleSignals(ctx context.Context) error { c := make(chan os.Signal, 1) signal.Notify(c, syscall.SIGQUIT, syscall.SIGTERM, syscall.SIGINT, syscall.SIGABRT) for { select { case s := <-c: if s == syscall.SIGABRT { if err := pprof.Lookup("goroutine").WriteTo(os.Stderr, 2); err != nil { log.Error(ctx, "failed to create pprof", "error", err) } } log.Log(ctx, "caught signal, attempting clean shutdown", "signal", s) return fmt.Errorf("%w signal=%v", ErrCaughtSignal, s) case <-ctx.Done(): return nil } }}
func makeSelfTestCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "self-test", Usage: "run gstreamer self-test", Action: func(ctx context.Context, cmd *urfavecli.Command) error { err := media.RunSelfTest(ctx) if err != nil { fmt.Println(err.Error()) os.Exit(1) } runtime.GC() if err := pprof.Lookup("goroutine").WriteTo(os.Stderr, 2); err != nil { log.Error(ctx, "error creating pprof", "error", err) } fmt.Println("self-test successful!") return nil }, }}
// makeValidateConfigCommand checks the current flags and environment exactly// the way a starting node would (cli.Validate) and exits without opening any// listeners, databases, or media pipelines. Point deploy tooling at it: it// fails before a misconfiguration can crashloop the node on restart.func makeValidateConfigCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} cmd := cli.NewCommand("validate-config") cmd.Usage = "validate configuration and exit without starting the node" // NewCommand's Before hook runs cli.Validate, so a bad configuration // fails the command before we ever get here. Do not check again in this // action: the checks must run against pristine flag values, and any // second pass would see the defaults PrepareConfig filled in as // operator-set conflicts. cmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { fmt.Println("configuration OK") return nil } return cmd}
// makeVODTestCommand runs the VOD gstreamer pipeline on a local file// and prints probe results. Useful for reproducing gstreamer-side// crashes against the static binary without needing the full server// scaffolding (DB, blob store, signer). The pipeline output is// discarded — this is a "did it crash or not" smoke test.func makeVODTestCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "vod-test", Usage: "run the VOD gstreamer pipeline on a local file", ArgsUsage: "[file]", Action: func(ctx context.Context, cmd *urfavecli.Command) error { args := cmd.Args() if args.Len() != 1 { return fmt.Errorf("usage: streamplace vod-test [file]") } return runVODTest(ctx, args.First()) }, }}
func runVODTest(ctx context.Context, path string) error { gstinit.InitGST()
f, err := os.Open(path) if err != nil { return fmt.Errorf("open %s: %w", path, err) } defer f.Close() info, err := f.Stat() if err != nil { return fmt.Errorf("stat %s: %w", path, err) } fmt.Printf("vod-test: processing %s (%d bytes)\n", path, info.Size())
start := time.Now() result, outBytes, err := vod.ProcessToDiscard(ctx, f, info.Size()) elapsed := time.Since(start)
if err != nil { fmt.Printf("vod-test: pipeline FAILED after %s: %v\n", elapsed, err) return err } fmt.Printf("vod-test: pipeline OK in %s, duration=%dms, output=%d bytes\n", elapsed, result.DurationMS, outBytes) if result.Video != nil { fmt.Printf(" video: codec=%s %dx%d fps=%d/%d\n", result.Video.Codec, result.Video.Width, result.Video.Height, result.Video.FPSNum, result.Video.FPSDen) } if result.Audio != nil { fmt.Printf(" audio: codec=%s rate=%d channels=%d\n", result.Audio.Codec, result.Audio.Rate, result.Audio.Channels) } return nil}
func makeStreamCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "stream", Usage: "stream command", ArgsUsage: "[user]", Action: func(ctx context.Context, cmd *urfavecli.Command) error { args := cmd.Args() if args.Len() != 1 { return fmt.Errorf("usage: streamplace stream [user]") } return Stream(args.First()) }, }}
// makeIngestWorkerCommand is the per-stream isolated ingest worker (Stage 1:// fMP4 / Mist pull). The node spawns it; it is not meant for direct use. It reads// the config handshake from fd 3, the fragmented-MP4 media from stdin, runs the mux + sign// pipeline, and writes signed canonical .m4s frames to fd 4 — dedicated fds so// stray stdout/stderr can't corrupt the frame stream. A clean run ends with an// End frame; a fatal error emits an Error frame before exiting non-zero.func makeIngestWorkerCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "ingest-worker", Usage: "internal: per-stream isolated ingest worker (spawned by the node)", ArgsUsage: "[streamer-did]", Hidden: true, Action: func(ctx context.Context, cmd *urfavecli.Command) error { // The streamer DID is passed on argv purely so the worker is // identifiable in a process listing (ps); the authoritative copy // still arrives in the fd-3 config. Thread it into the logger for // log correlation. if did := cmd.Args().First(); did != "" { ctx = log.WithLogValues(ctx, "streamer", did) log.Log(ctx, "ingest-worker starting") } cfgFile := os.NewFile(3, "ingest-config") if cfgFile == nil { return fmt.Errorf("ingest-worker: missing config fd 3") } cfgBytes, err := io.ReadAll(cfgFile) cfgFile.Close() if err != nil { return fmt.Errorf("ingest-worker: read config: %w", err) } var cfg media.IngestWorkerConfig if err := json.Unmarshal(cfgBytes, &cfg); err != nil { return fmt.Errorf("ingest-worker: parse config: %w", err) }
// WHIP transport: the worker owns the PeerConnection (built from the // offer in the config) and serves frames over the socket — no media fd. if cfg.Transport == media.IngestTransportWHIP { return media.ServeWHIPIngestWorkerSocket(ctx, cfg) }
// Detach/reattach transport: serve frames over a unix socket with // buffered reconnect (survives a main restart) instead of the fd-4 pipe. // Media comes from the fd-passed ingest connection (InputFD) when main // handed one off, else stdin. if cfg.SocketPath != "" { raw := io.Reader(os.Stdin) if cfg.InputFD > 0 { f := os.NewFile(uintptr(cfg.InputFD), "ingest-input") if f == nil { return fmt.Errorf("ingest-worker: bad input fd %d", cfg.InputFD) } defer f.Close() raw = f } return media.ServeMP4IngestWorkerSocket(ctx, cfg, media.WorkerInput(cfg, raw)) }
framesFile := os.NewFile(4, "ingest-frames") if framesFile == nil { return fmt.Errorf("ingest-worker: missing frames fd 4") } defer framesFile.Close() frames := ingestframe.NewWriter(framesFile)
// This fd-4 path has no back-channel for manifest updates, so the // manifest stays whatever main built at spawn. It's not used in prod // (api requires a hijackable connection); kept for the worker self-test. if err := media.RunMP4IngestWorker(ctx, cfg, os.Stdin, frames, func() []byte { return cfg.Manifest }); err != nil { _ = frames.Error(err.Error()) return err } return frames.End() }, }}
// makeRTMPPushWorkerCommand is the per-target isolated multistream egress// worker. The node spawns it; it is not meant for direct use. It reads the// config handshake (incl. the target URL + stream key) from fd 3, the assembled// fMP4 source stream from stdin, runs the native RTMP push pipeline, and writes// status Event frames to fd 4 — dedicated fds so stray stdout/stderr can't// corrupt the frame stream. A clean run ends with an End frame; a fatal error// emits an Error frame before exiting non-zero.func makeRTMPPushWorkerCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "rtmp-push-worker", Usage: "internal: per-target isolated RTMP push worker (spawned by the node)", ArgsUsage: "[streamer-did]", Hidden: true, Action: func(ctx context.Context, cmd *urfavecli.Command) error { // The streamer DID is passed on argv purely so the worker is // identifiable in a process listing (ps); the target URL stays on fd 3. if did := cmd.Args().First(); did != "" { ctx = log.WithLogValues(ctx, "streamer", did) log.Log(ctx, "rtmp-push-worker starting") } cfgFile := os.NewFile(3, "push-config") if cfgFile == nil { return fmt.Errorf("rtmp-push-worker: missing config fd 3") } cfgBytes, err := io.ReadAll(cfgFile) cfgFile.Close() if err != nil { return fmt.Errorf("rtmp-push-worker: read config: %w", err) } var cfg media.RTMPPushWorkerConfig if err := json.Unmarshal(cfgBytes, &cfg); err != nil { return fmt.Errorf("rtmp-push-worker: parse config: %w", err) }
eventsFile := os.NewFile(4, "push-events") if eventsFile == nil { return fmt.Errorf("rtmp-push-worker: missing events fd 4") } defer eventsFile.Close() events := ingestframe.NewWriter(eventsFile)
if err := media.RunRTMPPushWorker(ctx, cfg, os.Stdin, events); err != nil { _ = events.Error(err.Error()) return err } return events.End() }, }}
func makeLiveCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} liveCmd := cli.NewCommand("live") liveCmd.Usage = "start live stream (pipe fragmented MP4 to stdin)" liveCmd.ArgsUsage = "[stream-key]" liveCmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { args := cmd.Args() if args.Len() != 1 { return fmt.Errorf("usage: streamplace live [flags] [stream-key]") } return Live(args.First(), cli.HTTPInternalAddr) } return liveCmd}
func makeWhepCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "whep", Usage: "WHEP client", Flags: []urfavecli.Flag{ &urfavecli.IntFlag{ Name: "count", Usage: "number of concurrent streams (for load testing)", Value: 1, }, &urfavecli.DurationFlag{ Name: "duration", Usage: "stop after this long", }, &urfavecli.StringFlag{ Name: "endpoint", Usage: "endpoint to send the WHEP request to", }, }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { return WHEP( ctx, cmd.Int("count"), cmd.Duration("duration"), cmd.String("endpoint"), ) }, }}
func makeWhipCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "whip", Usage: "WHIP client", Flags: []urfavecli.Flag{ &urfavecli.StringFlag{ Name: "stream-key", Usage: "stream key", }, &urfavecli.IntFlag{ Name: "count", Usage: "number of concurrent streams (for load testing)", Value: 1, }, &urfavecli.IntFlag{ Name: "viewers", Usage: "number of viewers to simulate per stream", }, &urfavecli.DurationFlag{ Name: "duration", Usage: "total duration of the stream, including retries", }, &urfavecli.StringFlag{ Name: "file", Usage: "file to stream (needs to be an MP4 containing H264 video and Opus audio)", Required: true, }, &urfavecli.StringFlag{ Name: "endpoint", Usage: "endpoint to send the WHIP request to", Value: "http://127.0.0.1:38080", }, &urfavecli.DurationFlag{ Name: "freeze-after", Usage: "freeze the stream after the given duration", }, &urfavecli.DurationFlag{ Name: "retry", Usage: "restart after EOF or failure with this delay (e.g. 1s; 0 disables retries)", }, }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { ctx, stop := signal.NotifyContext(ctx, syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP) defer stop() return WHIP( ctx, cmd.String("stream-key"), cmd.Int("count"), cmd.Int("viewers"), cmd.Duration("duration"), cmd.String("file"), cmd.String("endpoint"), cmd.Duration("freeze-after"), cmd.Duration("retry"), ) }, }}
func makeCombineCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} combineCmd := cli.NewCommand("combine") combineCmd.Usage = "combine segments" combineCmd.ArgsUsage = "[output] [input1] [input2...]" combineCmd.Flags = []urfavecli.Flag{ &urfavecli.StringFlag{ Name: "debug-dir", Usage: "directory to write debug output", }, } combineCmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { args := cmd.Args() if args.Len() < 2 { return fmt.Errorf("usage: streamplace combine [--debug-dir dir] [output] [input1] [input2...]") } ctx = log.WithDebugValue(ctx, cli.Debug) return Combine( ctx, &cli, cmd.String("debug-dir"), args.Get(0), args.Slice()[1:], ) } return combineCmd}
func makeSplitCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} splitCmd := cli.NewCommand("split") splitCmd.Usage = "split video file" splitCmd.ArgsUsage = "[input file] [output directory]" splitCmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { args := cmd.Args() if args.Len() != 2 { return fmt.Errorf("usage: streamplace split [flags] [input file] [output directory]") } ctx = log.WithDebugValue(ctx, cli.Debug) gstinit.InitGST() return Split(ctx, args.Get(0), args.Get(1)) } return splitCmd}
func makeLivepeerCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "livepeer", Usage: "run livepeer gateway", Action: func(ctx context.Context, cmd *urfavecli.Command) error { return GoLivepeer(ctx, config.LivepeerFlagSet) }, }}
func makeMigrateCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} return &urfavecli.Command{ Name: "migrate", Usage: "run database migrations", Action: func(ctx context.Context, cmd *urfavecli.Command) error { return statedb.Migrate(&cli) }, }}
// makeMigrateStateCommand copies the state database from one engine to// another — the sqlite → Postgres move. Run it once while the old node is// still up to prove the target out, stop the node, run it again for the// delta (it is idempotent), then start the node with --db-url pointing at// the target.func makeMigrateStateCommand() *urfavecli.Command { var from, to string var batch int return &urfavecli.Command{ Name: "migrate-statedb", Usage: "copy the state database to another engine (sqlite → Postgres), then exit", Flags: []urfavecli.Flag{ &urfavecli.StringFlag{Name: "from", Usage: "source state database URL, e.g. sqlite:///data/state.sqlite", Required: true, Destination: &from}, &urfavecli.StringFlag{Name: "to", Usage: "target state database URL, e.g. postgres://user:pass@host/streamplace (created if missing)", Required: true, Destination: &to}, &urfavecli.IntFlag{Name: "batch-size", Usage: "rows per INSERT", Value: 500, Destination: &batch}, }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { reports, err := statedb.CopyState(ctx, from, to, batch) fmt.Fprintf(os.Stderr, "\n%-28s %10s %10s %10s\n", "table", "source", "inserted", "target") for _, r := range reports { fmt.Fprintf(os.Stderr, "%-28s %10d %10d %10d\n", r.Table, r.Source, r.Inserted, r.Target) } return err }, }}
// makeSyncCommand runs the backfill sweep to completion and exits, without// starting a node. It is for the case where a new index revision has to be warm// before traffic reaches it: run this, wait for it to finish, then start the// server -- rather than starting the server and serving from an index that is// still filling in behind it.func makeSyncCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} syncCmd := cli.NewCommand("sync") syncCmd.Usage = "index every repo this node knows about, then exit" syncCmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { return runSync(ctx, build, cmd, &cli) } return syncCmd}
// runSync builds the smallest stack a sweep needs -- the index, the state// database, an identity resolver -- and nothing else. No HTTP servers, no media// manager, no firehose: this process talks to other people's PDSes and to the// two databases, and then it is done.func runSync(ctx context.Context, build *config.BuildFlags, cmd *urfavecli.Command, cli *config.CLI) error { // The sync command's Before hook already ran cli.Validate; calling it // again would see PrepareConfig's defaults as operator-set conflicts. log.SetColorLogger(cli.Color) ctx = log.WithDebugValue(ctx, cli.Debug) log.Log(ctx, "streamplace sync", "version", build.Version, "dataDir", cli.DataDir)
if err := os.MkdirAll(cli.DataDir, os.ModePerm); err != nil { return fmt.Errorf("error creating streamplace dir at %s: %w", cli.DataDir, err) } mod, err := model.MakeDBConns(cli.DataFilePath([]string{"index"}), cli.IndexDBConnections) if err != nil { return err } state, err := statedb.MakeDB(ctx, cli, nil, mod) if err != nil { return err } atsync := &atproto.ATProtoSynchronizer{ CLI: cli, Model: mod, StatefulDB: state, Bus: bus.NewBus(), } // A sweep that could not sync a single repo is a broken node and exits // nonzero; anything less than that heals on the next run, so it is logged // and forgiven. return atsync.Sweep(ctx)}
// resolveLiveSigningKey returns the did:key whose private half signed a// streamer's live segments, for stamping on live-to-VOD place.stream.media.track// records so playback can verify them. It picks the most recently created// non-revoked place.stream.key for the repo; streamers normally have exactly// one. Errors if the repo has no active signing key.func resolveLiveSigningKey(mod model.Model, repoDID string) (string, error) { keys, err := mod.GetSigningKeysForRepo(repoDID) if err != nil { return "", err } var best *model.SigningKey for i := range keys { k := &keys[i] if k.RevokedAt != nil { continue } if best == nil || k.CreatedAt.After(best.CreatedAt) { best = k } } if best == nil { return "", fmt.Errorf("no active signing key for repo %s", repoDID) } return best.DID, nil}
// makeBrandingCommand exports and imports branding bundles straight against// the state database, without a running node or a signed-in admin: the way// to seed a fresh node before anyone can log in, or to copy a look between// hosts from a terminal.func makeBrandingCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} root := cli.NewCommand("branding") root.Usage = "export or import the node's branding as a bundle (zip with branding.yaml)" open := func(ctx context.Context, cmd *urfavecli.Command) (*statedb.StatefulDB, string, error) { // The branding command's Before hook already ran cli.Validate; see // runSync for why a second pass would be wrong. log.SetColorLogger(cli.Color) mod, err := model.MakeDBConns(cli.DataFilePath([]string{"index"}), cli.IndexDBConnections) if err != nil { return nil, "", err } state, err := statedb.MakeDB(ctx, &cli, nil, mod) if err != nil { return nil, "", err } return state, cli.BroadcasterDID(), nil } exportCmd := &urfavecli.Command{ Name: "export", Usage: "write the node's branding to a bundle", ArgsUsage: "<out.zip>", Action: func(ctx context.Context, cmd *urfavecli.Command) error { if cmd.Args().Len() != 1 { return fmt.Errorf("usage: streamplace branding export <out.zip>") } state, bid, err := open(ctx, cmd) if err != nil { return err } bs, err := branding.Export(ctx, state, bid) if err != nil { return err } if err := os.WriteFile(cmd.Args().First(), bs, 0o644); err != nil { return err } log.Log(ctx, "branding exported", "broadcaster", bid, "file", cmd.Args().First(), "bytes", len(bs)) return nil }, } var merge, dryRun bool importCmd := &urfavecli.Command{ Name: "import", Usage: "apply a bundle to the node (replaces branding unless --merge)", ArgsUsage: "<in.zip>", Flags: []urfavecli.Flag{ &urfavecli.BoolFlag{Name: "merge", Usage: "keep keys the bundle does not mention", Destination: &merge}, &urfavecli.BoolFlag{Name: "dry-run", Usage: "report what would change without writing", Destination: &dryRun}, }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { if cmd.Args().Len() != 1 { return fmt.Errorf("usage: streamplace branding import [--merge] [--dry-run] <in.zip>") } bs, err := os.ReadFile(cmd.Args().First()) if err != nil { return err } state, bid, err := open(ctx, cmd) if err != nil { return err } report, err := branding.Import(ctx, state, bid, bs, merge, dryRun) if err != nil { return err } for _, c := range report.Changes { if c.Detail != "" { fmt.Printf("%-10s %-24s %s\n", c.Action, c.Key, c.Detail) } else { fmt.Printf("%-10s %s\n", c.Action, c.Key) } } for _, w := range report.Warnings { fmt.Printf("warning: %s\n", w) } if !report.Applied { fmt.Println("dry run: nothing written") } return nil }, } root.Commands = append(root.Commands, exportCmd, importCmd) return root}