Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
48 kB · 1479 lines
Go
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480package 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), 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) { // 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 }, }}
// 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: "duration of the stream", }, &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", }, }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { 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"), ) }, }}
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 { if err := cli.Validate(cmd); err != nil { return err } 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) { if err := cli.Validate(cmd); err != nil { return nil, "", err } 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}