From fceb40b776eaf8463337f5413f85d1bc98cf265e Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 28 Oct 2025 15:00:58 -0700 Subject: [PATCH] replication: put irohreplicator behind an interface --- pkg/cmd/streamplace.go | 55 +++++++++++++++------------ pkg/config/config.go | 23 +++++++---- pkg/director/director.go | 10 ++--- pkg/director/stream_session.go | 9 +++-- pkg/media/media.go | 2 - pkg/replication/boring/boring.go | 47 ----------------------- pkg/replication/iroh_replicator/kv.go | 13 ++++--- pkg/replication/replicator.go | 15 +++++++- 8 files changed, 77 insertions(+), 97 deletions(-) delete mode 100644 pkg/replication/boring/boring.go diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 311cb54a..b8df6ade 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -14,6 +14,7 @@ import ( "path/filepath" "runtime" "runtime/pprof" + "slices" "strconv" "strings" "syscall" @@ -35,6 +36,7 @@ import ( "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/notifications" + "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/replication/iroh_replicator" "stream.place/streamplace/pkg/rtmps" v0 "stream.place/streamplace/pkg/schema/v0" @@ -404,38 +406,41 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { } } - exists, err := cli.DataFileExists([]string{"iroh-kv-secret"}) - if err != nil { - return err - } - if !exists { - secret := make([]byte, 32) - _, err := rand.Read(secret) + var replicator replication.Replicator = nil + if slices.Contains(cli.Replicators, config.ReplicatorIroh) { + exists, err := cli.DataFileExists([]string{"iroh-kv-secret"}) if err != nil { - return fmt.Errorf("failed to generate random secret: %w", err) + return err } - err = cli.DataFileWrite([]string{"iroh-kv-secret"}, bytes.NewReader(secret), true) + 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 } - } - 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) + 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 } } - swarm, err := iroh_replicator.NewSwarm(ctx, &cli, secret, topic, mm, b, mod) - if err != nil { - return err - } op := oatproxy.New(&oatproxy.Config{ Host: host, @@ -449,7 +454,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { ClientMetadata: clientMetadata, Public: cli.PublicOAuth, }) - d := director.NewDirector(mm, mod, &cli, b, op, state, swarm) + d := director.NewDirector(mm, mod, &cli, b, op, state, replicator) a, err := api.MakeStreamplaceAPI(&cli, mod, state, eip712signer, noter, mm, ms, b, atsync, d, op) if err != nil { return err @@ -517,7 +522,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { }) group.Go(func() error { - return swarm.Start(ctx, cli.Tickets) + return replicator.Start(ctx, &cli) }) if cli.LivepeerGateway { diff --git a/pkg/config/config.go b/pkg/config/config.go index 091e192e..01552c67 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -131,6 +131,7 @@ type CLI struct { DisableIrohRelay bool DevAccountCreds map[string]string StreamSessionTimeout time.Duration + Replicators []string } // ContentFilters represents the content filtering configuration @@ -144,6 +145,11 @@ type ContentFilters struct { } `json:"distribution_policy"` } +const ( + ReplicatorHTTP string = "http" + ReplicatorIroh string = "iroh" +) + func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { fs := flag.NewFlagSet("streamplace", flag.ExitOnError) fs.StringVar(&cli.DataDir, "data-dir", DefaultDataDir(), "directory for keeping all streamplace data") @@ -178,9 +184,9 @@ func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { fs.StringVar(&cli.LivepeerGatewayURL, "livepeer-gateway-url", "", "URL of the Livepeer Gateway to use for transcoding") fs.BoolVar(&cli.LivepeerGateway, "livepeer-gateway", false, "enable embedded Livepeer Gateway") fs.BoolVar(&cli.WideOpen, "wide-open", false, "allow ALL streams to be uploaded to this node (not recommended for production)") - cli.StringSliceFlag(fs, &cli.AllowedStreams, "allowed-streams", "", "if set, only allow these addresses or atproto DIDs to upload to this node") - cli.StringSliceFlag(fs, &cli.Peers, "peers", "", "other streamplace nodes to replicate to") - cli.StringSliceFlag(fs, &cli.Redirects, "redirects", "", "http 302s /path/one:/path/two,/path/three:/path/four") + cli.StringSliceFlag(fs, &cli.AllowedStreams, "allowed-streams", []string{}, "if set, only allow these addresses or atproto DIDs to upload to this node") + cli.StringSliceFlag(fs, &cli.Peers, "peers", []string{}, "other streamplace nodes to replicate to") + cli.StringSliceFlag(fs, &cli.Redirects, "redirects", []string{}, "http 302s /path/one:/path/two,/path/three:/path/four") cli.DebugFlag(fs, &cli.Debug, "debug", "", "modified log verbosity for specific functions or files in form func=ToHLS:3,file=gstreamer.go:4") fs.BoolVar(&cli.TestStream, "test-stream", false, "run a built-in test stream on boot") fs.BoolVar(&cli.NoFirehose, "no-firehose", false, "disable the bluesky firehose") @@ -205,7 +211,7 @@ func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { fs.BoolVar(&cli.NewWebRTCPlayback, "new-webrtc-playback", true, "enable new webrtc playback") fs.StringVar(&cli.AppleTeamID, "apple-team-id", "", "apple team id for deep linking") fs.StringVar(&cli.AndroidCertFingerprint, "android-cert-fingerprint", "", "android cert fingerprint for deep linking") - cli.StringSliceFlag(fs, &cli.Labelers, "labelers", "", "did of labelers that this instance should subscribe to") + cli.StringSliceFlag(fs, &cli.Labelers, "labelers", []string{}, "did of labelers that this instance should subscribe to") fs.StringVar(&cli.AtprotoDID, "atproto-did", "", "atproto did to respond to on /.well-known/atproto-did (default did:web:PUBLIC_HOST)") cli.JSONFlag(fs, &cli.ContentFilters, "content-filters", "{}", "JSON content filtering rules") fs.BoolVar(&cli.LivepeerHelp, "livepeer-help", false, "print help for livepeer flags and exit") @@ -213,11 +219,12 @@ func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { fs.BoolVar(&cli.SQLLogging, "sql-logging", false, "enable sql logging") fs.StringVar(&cli.SentryDSN, "sentry-dsn", "", "sentry dsn for error reporting") fs.BoolVar(&cli.LivepeerDebug, "livepeer-debug", false, "log livepeer segments to $SP_DATA_DIR/livepeer-debug") - cli.StringSliceFlag(fs, &cli.Tickets, "tickets", "[]", "tickets to join the swarm with") + cli.StringSliceFlag(fs, &cli.Tickets, "tickets", []string{}, "tickets to join the swarm with") fs.StringVar(&cli.IrohTopic, "iroh-topic", "", "topic to use for the iroh swarm (must be 32 bytes in hex)") fs.BoolVar(&cli.DisableIrohRelay, "disable-iroh-relay", false, "disable the iroh relay") cli.KVSliceFlag(fs, &cli.DevAccountCreds, "dev-account-creds", "", "(FOR DEVELOPMENT ONLY) did=password pairs for logging into test accounts without oauth") fs.DurationVar(&cli.StreamSessionTimeout, "stream-session-timeout", 60*time.Second, "how long to wait before considering a stream inactive on this node?") + cli.StringSliceFlag(fs, &cli.Replicators, "replicators", []string{ReplicatorIroh}, "list of replication protocols to use (http, iroh)") lpFlags := flag.NewFlagSet("livepeer", flag.ContinueOnError) _ = starter.NewLivepeerConfig(lpFlags) @@ -533,15 +540,15 @@ func (cli *CLI) AddressSliceFlag(fs *flag.FlagSet, dest *[]aqpub.Pub, name, defa }) } -func (cli *CLI) StringSliceFlag(fs *flag.FlagSet, dest *[]string, name, defaultValue, usage string) { - *dest = []string{} +func (cli *CLI) StringSliceFlag(fs *flag.FlagSet, dest *[]string, name string, defaultValue []string, usage string) { + *dest = defaultValue usage = fmt.Sprintf(`%s (default: "%s")`, usage, *dest) fs.Func(name, usage, func(s string) error { if s == "" { return nil } strs := strings.Split(s, ",") - *dest = append(*dest, strs...) + *dest = append([]string{}, strs...) return nil }) } diff --git a/pkg/director/director.go b/pkg/director/director.go index 1dcce845..75616d2e 100644 --- a/pkg/director/director.go +++ b/pkg/director/director.go @@ -12,7 +12,7 @@ import ( "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/model" - "stream.place/streamplace/pkg/replication/iroh_replicator" + "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/statedb" ) @@ -31,10 +31,10 @@ type Director struct { streamSessionsMu sync.Mutex op *oatproxy.OATProxy statefulDB *statedb.StatefulDB - swarm *iroh_replicator.IrohSwarm + replicator replication.Replicator } -func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus *bus.Bus, op *oatproxy.OATProxy, statefulDB *statedb.StatefulDB, swarm *iroh_replicator.IrohSwarm) *Director { +func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus *bus.Bus, op *oatproxy.OATProxy, statefulDB *statedb.StatefulDB, replicator replication.Replicator) *Director { return &Director{ mm: mm, mod: mod, @@ -44,7 +44,7 @@ func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus * streamSessionsMu: sync.Mutex{}, op: op, statefulDB: statefulDB, - swarm: swarm, + replicator: replicator, } } @@ -75,7 +75,7 @@ func (d *Director) Start(ctx context.Context) error { packets: make([]bus.PacketizedSegment, 0), started: make(chan struct{}), statefulDB: d.statefulDB, - swarm: d.swarm, + replicator: d.replicator, } d.streamSessions[not.Segment.RepoDID] = ss g.Go(func() error { diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 96279989..e2503234 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -23,7 +23,7 @@ import ( "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/renditions" - "stream.place/streamplace/pkg/replication/iroh_replicator" + "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/streamplace" @@ -50,7 +50,7 @@ type StreamSession struct { ctx context.Context packets []bus.PacketizedSegment statefulDB *statedb.StatefulDB - swarm *iroh_replicator.IrohSwarm + replicator replication.Replicator } func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotification) error { @@ -460,7 +460,10 @@ func (ss *StreamSession) UpdateBroadcastOrigin(ctx context.Context) error { Server: fmt.Sprintf("did:web:%s", ss.cli.ServerHost), Broadcaster: &broadcaster, UpdatedAt: time.Now().UTC().Format(util.ISO8601), - IrohTicket: &ss.swarm.NodeTicket, + } + err := ss.replicator.BuildOriginRecord(&origin) + if err != nil { + return fmt.Errorf("could not build origin record: %w", err) } client, err := ss.GetClientByDID(ss.repoDID) diff --git a/pkg/media/media.go b/pkg/media/media.go index 80a8714d..da0dabe4 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -26,7 +26,6 @@ import ( "stream.place/streamplace/pkg/streamplace" "stream.place/streamplace/pkg/log" - "stream.place/streamplace/pkg/replication" "github.com/piprate/json-gold/ld" @@ -42,7 +41,6 @@ var StreamplaceMetadata = "place.stream.metadata" type MediaManager struct { cli *config.CLI - replicator replication.Replicator hlsRunning map[string]*M3U8 hlsRunningMut sync.Mutex httpPipes map[string]io.Writer diff --git a/pkg/replication/boring/boring.go b/pkg/replication/boring/boring.go deleted file mode 100644 index ff0f75d4..00000000 --- a/pkg/replication/boring/boring.go +++ /dev/null @@ -1,47 +0,0 @@ -package boring - -import ( - "bytes" - "context" - "fmt" - "io" - "net/http" - - "stream.place/streamplace/pkg/aqhttp" - "stream.place/streamplace/pkg/log" -) - -// boring HTTP replication mechanism -type BoringReplicator struct { - Peers []string -} - -func (rep *BoringReplicator) NewSegment(ctx context.Context, bs []byte) { - for _, p := range rep.Peers { - go func(peer string) { - ctx := log.WithLogValues(ctx, "peer", peer) - err := sendSegment(ctx, peer, bs) - if err != nil { - log.Log(ctx, "error replicating segment", "error", err) - } - }(p) - } -} - -func sendSegment(ctx context.Context, peer string, bs []byte) error { - r := bytes.NewReader(bs) - peerURL := fmt.Sprintf("%s/api/segment", peer) - req, err := http.NewRequestWithContext(ctx, "POST", peerURL, r) - if err != nil { - return err - } - res, err := aqhttp.Client.Do(req) - if err != nil { - return err - } - if res.StatusCode != 204 { - body, _ := io.ReadAll(res.Body) - return fmt.Errorf("unexpected http code %d body=%s", res.StatusCode, body) - } - return nil -} diff --git a/pkg/replication/iroh_replicator/kv.go b/pkg/replication/iroh_replicator/kv.go index a43fed58..15ca1b9a 100644 --- a/pkg/replication/iroh_replicator/kv.go +++ b/pkg/replication/iroh_replicator/kv.go @@ -121,14 +121,18 @@ func NewSwarm(ctx context.Context, cli *config.CLI, secret []byte, topic []byte, return &swarm, nil } -func (swarm *IrohSwarm) Start(ctx context.Context, tickets []string) error { - if len(tickets) > 0 { - err := swarm.Node.JoinPeers(tickets) +func (swarm *IrohSwarm) BuildOriginRecord(origin *streamplace.BroadcastOrigin) error { + origin.IrohTicket = &swarm.NodeTicket + return nil +} + +func (swarm *IrohSwarm) Start(ctx context.Context, cli *config.CLI) error { + if len(cli.Tickets) > 0 { + err := swarm.Node.JoinPeers(cli.Tickets) if err != nil { return fmt.Errorf("failed to join peers: %w", err) } } - nodeId, err := swarm.Node.NodeId() if err != nil { return fmt.Errorf("failed to get node id: %w", err) @@ -189,7 +193,6 @@ func (swarm *IrohSwarm) startKV(ctx context.Context) error { log.Debug(ctx, "SubscribeItemOther", "other", item) } } - return nil } func (swarm *IrohSwarm) handleIrohMessage(ctx context.Context, item iroh_streamplace.SubscribeItemEntry) error { diff --git a/pkg/replication/replicator.go b/pkg/replication/replicator.go index 0535e8ec..5237b348 100644 --- a/pkg/replication/replicator.go +++ b/pkg/replication/replicator.go @@ -1,7 +1,18 @@ package replication -import "context" +import ( + "context" + + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/media" + "stream.place/streamplace/pkg/streamplace" +) type Replicator interface { - NewSegment(context.Context, []byte) + // start the replicator, ending on context cancellation. if your replicator doesn't need to start anything, you can just block on <-ctx.Done() + Start(context.Context, *config.CLI) error + // hey, we have a new segment! send it to whoever + SendSegment(context.Context, *media.NewSegmentNotification) error + // populate this origin record with whatever fields are pertinent to your replicator + BuildOriginRecord(*streamplace.BroadcastOrigin) error } -- 2.51.2