diff --git a/pkg/api/api.go b/pkg/api/api.go index 3ad0daa1..58c18ca2 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -87,7 +87,7 @@ type WebsocketTracker struct { mu sync.RWMutex } -func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.StatefulDB, signer *eip712.EIP712Signer, noter notifications.FirebaseNotifier, mm *media.MediaManager, ms media.MediaSigner, bus *bus.Bus, atsync *atproto.ATProtoSynchronizer, d *director.Director, op *oatproxy.OATProxy) (*StreamplaceAPI, error) { +func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.StatefulDB, noter notifications.FirebaseNotifier, mm *media.MediaManager, ms media.MediaSigner, bus *bus.Bus, atsync *atproto.ATProtoSynchronizer, d *director.Director, op *oatproxy.OATProxy) (*StreamplaceAPI, error) { updater, err := PrepareUpdater(cli) if err != nil { return nil, err @@ -96,7 +96,6 @@ func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, statefulDB *statedb.St Model: mod, StatefulDB: statefulDB, Updater: updater, - Signer: signer, FirebaseNotifier: noter, MediaManager: mm, MediaSigner: ms, diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go index 752470a8..ec9e8cfe 100644 --- a/pkg/api/api_internal.go +++ b/pkg/api/api_internal.go @@ -319,29 +319,6 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err w.WriteHeader(204) }) - router.GET("/settings", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { - w.Header().Set("Access-Control-Allow-Origin", "*") - w.Header().Set("Access-Control-Allow-Methods", "GET") - w.Header().Set("Access-Control-Allow-Headers", "Content-Type") - - id := a.Signer.Hex() - - ident, err := a.Model.GetIdentity(id) - if err != nil { - errors.WriteHTTPInternalServerError(w, "unable to get settings", err) - return - } - - bs, err := json.Marshal(ident) - if err != nil { - errors.WriteHTTPInternalServerError(w, "unable to marshal json", err) - return - } - if _, err := w.Write(bs); err != nil { - log.Error(ctx, "error writing response", "error", err) - } - }) - router.GET("/followers/:user", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { user := p.ByName("user") if user == "" { diff --git a/pkg/cmd/clip.go b/pkg/cmd/clip.go deleted file mode 100644 index c18bc930..00000000 --- a/pkg/cmd/clip.go +++ /dev/null @@ -1,29 +0,0 @@ -package cmd - -import ( - "context" - "fmt" - "os" - - "stream.place/streamplace/pkg/gstinit" - "stream.place/streamplace/pkg/log" - "stream.place/streamplace/pkg/media" -) - -func Clip(ctx context.Context, args []string, out string) error { - if out == "" { - return fmt.Errorf("out is required") - } - log.Log(ctx, "clip", "out", out) - gstinit.InitGST() - fd, err := os.Create(out) - if err != nil { - return err - } - defer fd.Close() - err = media.Clip(ctx, args, fd) - if err != nil { - return err - } - return nil -} diff --git a/pkg/cmd/combine.go b/pkg/cmd/combine.go new file mode 100644 index 00000000..aefb228c --- /dev/null +++ b/pkg/cmd/combine.go @@ -0,0 +1,38 @@ +package cmd + +import ( + "context" + "os" + + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/gstinit" + "stream.place/streamplace/pkg/media" +) + +func Combine(ctx context.Context, cli *config.CLI, args []string) error { + cryptoSigner, err := createSigner(ctx, cli) + if err != nil { + return err + } + ms, err := media.MakeMediaSigner(ctx, cli, "combine", cryptoSigner, nil) + if err != nil { + return err + } + outFile := args[0] + inputs := args[1:] + gstinit.InitGST() + outBs, err := media.CombineSegments(ctx, inputs, ms) + if err != nil { + return err + } + fd, err := os.Create(outFile) + if err != nil { + return err + } + defer fd.Close() + _, err = fd.Write(outBs) + if err != nil { + return err + } + return nil +} diff --git a/pkg/cmd/signer.go b/pkg/cmd/signer.go new file mode 100644 index 00000000..15cd6d9b --- /dev/null +++ b/pkg/cmd/signer.go @@ -0,0 +1,106 @@ +package cmd + +import ( + "context" + "crypto" + "fmt" + "os" + "strconv" + + "github.com/ThalesGroup/crypto11" + "golang.org/x/term" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/crypto/signers" + v0 "stream.place/streamplace/pkg/schema/v0" + + "stream.place/streamplace/pkg/crypto/signers/eip712" + "stream.place/streamplace/pkg/log" +) + +func createSigner(ctx context.Context, cli *config.CLI) (crypto.Signer, error) { + schema, err := v0.MakeV0Schema() + if err != nil { + return nil, err + } + eip712signer, err := eip712.MakeEIP712Signer(ctx, &eip712.EIP712SignerOptions{ + Schema: schema, + EthKeystorePath: cli.EthKeystorePath, + EthAccountAddr: cli.EthAccountAddr, + EthKeystorePassword: cli.EthPassword, + }) + if err != nil { + return nil, err + } + var signer crypto.Signer = eip712signer + if cli.PKCS11ModulePath != "" { + conf := &crypto11.Config{ + Path: cli.PKCS11ModulePath, + } + count := 0 + for _, val := range []string{cli.PKCS11TokenSlot, cli.PKCS11TokenLabel, cli.PKCS11TokenSerial} { + if val != "" { + count += 1 + } + } + if count != 1 { + return nil, fmt.Errorf("need exactly one of pkcs11-token-slot, pkcs11-token-label, or pkcs11-token-serial (got %d)", count) + } + if cli.PKCS11TokenSlot != "" { + num, err := strconv.ParseInt(cli.PKCS11TokenSlot, 10, 16) + if err != nil { + return nil, fmt.Errorf("error parsing pkcs11-slot: %w", err) + } + numint := int(num) + // why does crypto11 want this as a reference? odd. + conf.SlotNumber = &numint + } + if cli.PKCS11TokenLabel != "" { + conf.TokenLabel = cli.PKCS11TokenLabel + } + if cli.PKCS11TokenSerial != "" { + conf.TokenSerial = cli.PKCS11TokenSerial + } + pin := cli.PKCS11Pin + if pin == "" { + fmt.Printf("Please enter PKCS11 PIN: ") + password, err := term.ReadPassword(int(os.Stdin.Fd())) + fmt.Println("") + if err != nil { + return nil, fmt.Errorf("error reading PKCS11 password: %w", err) + } + pin = string(password) + } + conf.Pin = pin + + sc, err := crypto11.Configure(conf) + if err != nil { + return nil, fmt.Errorf("error initalizing PKCS11 HSM: %w", err) + } + var id []byte = nil + var label []byte = nil + if cli.PKCS11KeypairID != "" { + num, err := strconv.ParseInt(cli.PKCS11KeypairID, 10, 8) + if err != nil { + return nil, fmt.Errorf("error parsing pkcs11-keypair-id: %w", err) + } + id = []byte{byte(num)} + } + if cli.PKCS11KeypairLabel != "" { + label = []byte(cli.PKCS11KeypairLabel) + } + hwsigner, err := sc.FindKeyPair(id, label) + if err != nil { + return nil, fmt.Errorf("error finding keypair on PKCS11 token: %w", err) + } + if hwsigner == nil { + return nil, fmt.Errorf("keypair on token not found (tried id='%s' label='%s')", cli.PKCS11KeypairID, cli.PKCS11KeypairLabel) + } + addr, err := signers.HexAddrFromSigner(hwsigner) + if err != nil { + return nil, fmt.Errorf("error getting ethereum address for hardware keypair: %w", err) + } + log.Log(ctx, "successfully initialized hardware signer", "address", addr) + signer = hwsigner + } + return signer, nil +} diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index adb0c98f..5798c3ee 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -3,7 +3,6 @@ package cmd import ( "bytes" "context" - "crypto" "crypto/rand" "errors" "flag" @@ -11,7 +10,6 @@ import ( "net/url" "os" "os/signal" - "path/filepath" "runtime" "runtime/pprof" "slices" @@ -25,13 +23,11 @@ import ( "github.com/livepeer/go-livepeer/cmd/livepeer/starter" "github.com/peterbourgon/ff/v3" "github.com/streamplace/oatproxy/pkg/oatproxy" - "golang.org/x/term" "stream.place/streamplace/pkg/aqhttp" "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/bus" - "stream.place/streamplace/pkg/crypto/signers" - "stream.place/streamplace/pkg/crypto/signers/eip712" "stream.place/streamplace/pkg/director" + "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" @@ -40,12 +36,10 @@ import ( "stream.place/streamplace/pkg/replication/iroh_replicator" "stream.place/streamplace/pkg/replication/websocketrep" "stream.place/streamplace/pkg/rtmps" - v0 "stream.place/streamplace/pkg/schema/v0" "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/storage" - "github.com/ThalesGroup/crypto11" _ "github.com/go-gst/go-glib/glib" _ "github.com/go-gst/go-gst/gst" "stream.place/streamplace/pkg/api" @@ -127,23 +121,24 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return WHIP(os.Args[2:]) } - if len(os.Args) > 1 && os.Args[1] == "clip" { - cli := config.CLI{Build: build} - fs := cli.NewFlagSet("streamplace clip") - out := fs.String("out", "", "output file") + if len(os.Args) > 1 && os.Args[1] == "combine" { + gstinit.InitGST() + cli := &config.CLI{Build: build} + fs := cli.NewFlagSet("streamplace combine") err := cli.Parse(fs, os.Args[2:]) if err != nil { return err } + log.Debug(context.Background(), "combine command: starting", "args", fs.Args()) ctx := context.Background() ctx = log.WithDebugValue(ctx, cli.Debug) - return Clip(ctx, fs.Args(), *out) + return Combine(ctx, cli, fs.Args()) } if len(os.Args) > 1 && os.Args[1] == "segment" { cli := config.CLI{Build: build} - fs := cli.NewFlagSet("streamplace clip") + fs := cli.NewFlagSet("streamplace split") out := fs.String("out-dir", "", "output directory") err := cli.Parse(fs, os.Args[2:]) @@ -217,6 +212,10 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { if *version { return nil } + signer, err := createSigner(ctx, &cli) + if err != nil { + return err + } if len(os.Args) > 1 && os.Args[1] == "migrate" { return statedb.Migrate(&cli) @@ -244,90 +243,6 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { if err != nil { return fmt.Errorf("error creating streamplace dir at %s:%w", cli.DataDir, err) } - schema, err := v0.MakeV0Schema() - if err != nil { - return err - } - eip712signer, err := eip712.MakeEIP712Signer(ctx, &eip712.EIP712SignerOptions{ - Schema: schema, - EthKeystorePath: cli.EthKeystorePath, - EthAccountAddr: cli.EthAccountAddr, - EthKeystorePassword: cli.EthPassword, - }) - if err != nil { - return err - } - var signer crypto.Signer = eip712signer - if cli.PKCS11ModulePath != "" { - conf := &crypto11.Config{ - Path: cli.PKCS11ModulePath, - } - count := 0 - for _, val := range []string{cli.PKCS11TokenSlot, cli.PKCS11TokenLabel, cli.PKCS11TokenSerial} { - if val != "" { - count += 1 - } - } - if count != 1 { - return fmt.Errorf("need exactly one of pkcs11-token-slot, pkcs11-token-label, or pkcs11-token-serial (got %d)", count) - } - if cli.PKCS11TokenSlot != "" { - num, err := strconv.ParseInt(cli.PKCS11TokenSlot, 10, 16) - if err != nil { - return fmt.Errorf("error parsing pkcs11-slot: %w", err) - } - numint := int(num) - // why does crypto11 want this as a reference? odd. - conf.SlotNumber = &numint - } - if cli.PKCS11TokenLabel != "" { - conf.TokenLabel = cli.PKCS11TokenLabel - } - if cli.PKCS11TokenSerial != "" { - conf.TokenSerial = cli.PKCS11TokenSerial - } - pin := cli.PKCS11Pin - if pin == "" { - fmt.Printf("Please enter PKCS11 PIN: ") - password, err := term.ReadPassword(int(os.Stdin.Fd())) - fmt.Println("") - if err != nil { - return fmt.Errorf("error reading PKCS11 password: %w", err) - } - pin = string(password) - } - conf.Pin = pin - - sc, err := crypto11.Configure(conf) - if err != nil { - return fmt.Errorf("error initalizing PKCS11 HSM: %w", err) - } - var id []byte = nil - var label []byte = nil - if cli.PKCS11KeypairID != "" { - num, err := strconv.ParseInt(cli.PKCS11KeypairID, 10, 8) - if err != nil { - return fmt.Errorf("error parsing pkcs11-keypair-id: %w", err) - } - id = []byte{byte(num)} - } - if cli.PKCS11KeypairLabel != "" { - label = []byte(cli.PKCS11KeypairLabel) - } - hwsigner, err := sc.FindKeyPair(id, label) - if err != nil { - return fmt.Errorf("error finding keypair on PKCS11 token: %w", err) - } - if hwsigner == nil { - return fmt.Errorf("keypair on token not found (tried id='%s' label='%s')", cli.PKCS11KeypairID, cli.PKCS11KeypairLabel) - } - addr, err := signers.HexAddrFromSigner(hwsigner) - if err != nil { - return fmt.Errorf("error getting ethereum address for hardware keypair: %w", err) - } - log.Log(ctx, "successfully initialized hardware signer", "address", addr) - signer = hwsigner - } mod, err := model.MakeDB(cli.DataFilePath([]string{"index"})) if err != nil { @@ -473,7 +388,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { Public: cli.PublicOAuth, }) 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) + a, err := api.MakeStreamplaceAPI(&cli, mod, state, noter, mm, ms, b, atsync, d, op) if err != nil { return err } @@ -571,20 +486,12 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { }) if cli.TestStream { - // regular stream self-test - testSigner, err := eip712.MakeEIP712Signer(ctx, &eip712.EIP712SignerOptions{ - Schema: schema, - EthKeystorePath: filepath.Join(cli.DataDir, "test-signer"), - }) - if err != nil { - return err - } atkey, err := atproto.ParsePubKey(signer.Public()) if err != nil { return err } did := atkey.DIDKey() - testMediaSigner, err := media.MakeMediaSigner(ctx, &cli, did, testSigner, mod) + testMediaSigner, err := media.MakeMediaSigner(ctx, &cli, did, signer, mod) if err != nil { return err } @@ -603,19 +510,15 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { }) // Start a test stream that will run intermittently - intermittentSigner, err := eip712.MakeEIP712Signer(ctx, &eip712.EIP712SignerOptions{ - Schema: schema, - EthKeystorePath: filepath.Join(cli.DataDir, "intermittent-signer"), - }) if err != nil { return err } - atkey2, err := atproto.ParsePubKey(intermittentSigner.Public()) + atkey2, err := atproto.ParsePubKey(signer.Public()) if err != nil { return err } did2 := atkey2.DIDKey() - intermittentMediaSigner, err := media.MakeMediaSigner(ctx, &cli, did2, intermittentSigner, mod) + intermittentMediaSigner, err := media.MakeMediaSigner(ctx, &cli, did2, signer, mod) if err != nil { return err } diff --git a/pkg/media/clip.go b/pkg/media/clip.go index 86b91ea5..320d74d2 100644 --- a/pkg/media/clip.go +++ b/pkg/media/clip.go @@ -9,7 +9,6 @@ import ( "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" - "github.com/google/uuid" "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/log" ) @@ -31,20 +30,14 @@ func readFile(ctx context.Context, source string) (*bus.Seg, error) { return seg, nil } -// This function remains in scope for the duration of a single users' playback -func Clip(ctx context.Context, sources []string, w io.Writer) error { - uu, err := uuid.NewV7() - if err != nil { - return err - } - ctx = log.WithLogValues(ctx, "webrtcID", uu.String()) - ctx = log.WithLogValues(ctx, "mediafunc", "Clip") +func CombineSegmentsUnsigned(ctx context.Context, sources []string, w io.Writer) error { + ctx = log.WithLogValues(ctx, "mediafunc", "CombineSegmentsUnsigned") ctx, cancel := context.WithCancel(ctx) defer cancel() pipelineSlice := []string{ "mp4mux name=muxer ! appsink sync=false name=mp4sink", - "h264parse name=videoparse ! muxer.video_0", + "h264parse name=videoparse ! h264timestamper ! muxer.video_0", "opusparse name=audioparse ! muxer.audio_0", } diff --git a/pkg/media/clip_test.go b/pkg/media/clip_test.go index 9e5ad8d3..21c24325 100644 --- a/pkg/media/clip_test.go +++ b/pkg/media/clip_test.go @@ -26,7 +26,7 @@ func innerTestClip(t *testing.T) error { fName := getFixture("sample-segment.mp4") inputFiles := []string{fName, fName, fName} buf := bytes.NewBuffer(nil) - err := Clip(context.Background(), inputFiles, buf) + err := CombineSegmentsUnsigned(context.Background(), inputFiles, buf) require.NoError(t, err) require.Greater(t, buf.Len(), 2900000) require.Less(t, buf.Len(), 3100000) diff --git a/pkg/media/clip_user.go b/pkg/media/clip_user.go index 3ffb3894..4d8e798f 100644 --- a/pkg/media/clip_user.go +++ b/pkg/media/clip_user.go @@ -33,7 +33,7 @@ func ClipUser(ctx context.Context, mod model.Model, cli *config.CLI, user string } segmentFiles = append(segmentFiles, fpath) } - err = Clip(ctx, segmentFiles, writer) + err = CombineSegmentsUnsigned(ctx, segmentFiles, writer) if err != nil { return fmt.Errorf("unable to clip segments: %w", err) } diff --git a/pkg/media/deterministic_muxing_test.go b/pkg/media/deterministic_muxing_test.go index 70aa4ddb..0ea1a5bb 100644 --- a/pkg/media/deterministic_muxing_test.go +++ b/pkg/media/deterministic_muxing_test.go @@ -50,7 +50,7 @@ func splitAndCombineTest(t *testing.T, tempDir string, inputDir string) string { log.Log(context.Background(), "creating combined file", "file", outFilePath) require.NoError(t, err) defer outFile.Close() - err = Clip(context.Background(), firstReport.Segs, outFile) + err = CombineSegmentsUnsigned(context.Background(), firstReport.Segs, outFile) require.NoError(t, err) hash, err := hashFile(outFilePath) require.NoError(t, err) diff --git a/pkg/media/ingredient_test.go b/pkg/media/ingredient_test.go index c5b38692..6453ebe1 100644 --- a/pkg/media/ingredient_test.go +++ b/pkg/media/ingredient_test.go @@ -61,7 +61,7 @@ func TestIngredientConcat(t *testing.T) { require.NoError(t, err) ms := msInterface.(*MediaSignerLocal) buf := bytes.Buffer{} - err = Clip(context.Background(), testVids, &buf) + err = CombineSegmentsUnsigned(context.Background(), testVids, &buf) require.NoError(t, err) ingredients := [][]byte{} startTS, err := time.Parse(util.ISO8601, testTimestamp) diff --git a/pkg/media/segment_combine.go b/pkg/media/segment_combine.go new file mode 100644 index 00000000..a154a91a --- /dev/null +++ b/pkg/media/segment_combine.go @@ -0,0 +1,29 @@ +package media + +import ( + "bytes" + "context" + "os" +) + +// CombineSegments combines a list of segments into a single segment that maintains all of the manifests +func CombineSegments(ctx context.Context, sources []string, ms MediaSigner) ([]byte, error) { + buf := bytes.Buffer{} + err := CombineSegmentsUnsigned(ctx, sources, &buf) + if err != nil { + return nil, err + } + ingredients := [][]byte{} + for _, source := range sources { + bs, err := os.ReadFile(source) + if err != nil { + return nil, err + } + ingredients = append(ingredients, bs) + } + signedConcatBS, err := ms.SignConcatMP4(context.Background(), bytes.NewReader(buf.Bytes()), ingredients) + if err != nil { + return nil, err + } + return signedConcatBS, nil +}