diff --git a/pkg/api/api.go b/pkg/api/api.go index d25bae5a..1653ddc3 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -497,7 +497,7 @@ func (a *StreamplaceAPI) HandleNotification(ctx context.Context) http.HandlerFun func (a *StreamplaceAPI) HandleSegment(ctx context.Context) http.HandlerFunc { return func(w http.ResponseWriter, req *http.Request) { - err := a.MediaManager.ValidateMP4(ctx, req.Body) + err := a.MediaManager.ValidateMP4(ctx, req.Body, false) if err != nil { apierrors.WriteHTTPBadRequest(w, "could not ingest segment", err) return diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 6501bd33..2c26758c 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -32,8 +32,6 @@ 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/boring" "stream.place/streamplace/pkg/replication/iroh_replicator" "stream.place/streamplace/pkg/rtmps" v0 "stream.place/streamplace/pkg/schema/v0" @@ -307,7 +305,6 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { log.Log(ctx, "successfully initialized hardware signer", "address", addr) signer = hwsigner } - var rep replication.Replicator = &boring.BoringReplicator{Peers: cli.Peers} mod, err := model.MakeDB(cli.DataFilePath([]string{"index"})) if err != nil { @@ -363,7 +360,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return fmt.Errorf("failed to migrate: %w", err) } - mm, err := media.MakeMediaManager(ctx, &cli, signer, rep, mod, b, atsync) + mm, err := media.MakeMediaManager(ctx, &cli, signer, mod, b, atsync) if err != nil { return err } @@ -403,7 +400,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return err } secret := buf.Bytes() - swarm, err := iroh_replicator.StartKV(ctx, cli.Tickets, secret) + swarm, err := iroh_replicator.NewSwarm(ctx, cli.Tickets, secret, mm) if err != nil { return err } @@ -419,7 +416,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { DownstreamJWK: cli.AccessJWK, ClientMetadata: clientMetadata, }) - d := director.NewDirector(mm, mod, &cli, b, op, state, swarm) + d := director.NewDirector(mm, mod, &cli, b, op, state) a, err := api.MakeStreamplaceAPI(&cli, mod, state, eip712signer, noter, mm, ms, b, atsync, d, op) if err != nil { return err diff --git a/pkg/config/config.go b/pkg/config/config.go index b71f54d3..f8743cbf 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -167,7 +167,7 @@ func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { fs.StringVar(&cli.PublicHost, "public-host", "", "public host for this streamplace node (excluding https:// e.g. stream.place)") fs.BoolVar(&cli.Thumbnail, "thumbnail", true, "enable thumbnail generation") fs.BoolVar(&cli.SmearAudio, "smear-audio", false, "enable audio smearing to create 'perfect' segment timestamps") - fs.BoolVar(&cli.ExternalSigning, "external-signing", true, "enable external signing via exec (prevents potential memory leak)") + fs.BoolVar(&cli.ExternalSigning, "external-signing", false, "enable external signing via exec (prevents potential memory leak)") fs.StringVar(&cli.TracingEndpoint, "tracing-endpoint", "", "gRPC endpoint to send traces to") fs.IntVar(&cli.RateLimitPerSecond, "rate-limit-per-second", 0, "rate limit for requests per second per ip") fs.IntVar(&cli.RateLimitBurst, "rate-limit-burst", 0, "rate limit burst for requests per ip") diff --git a/pkg/director/director.go b/pkg/director/director.go index e9e75499..d708bcb4 100644 --- a/pkg/director/director.go +++ b/pkg/director/director.go @@ -2,11 +2,9 @@ package director import ( "context" - "encoding/json" "fmt" "sync" - "github.com/bluesky-social/indigo/util" "github.com/streamplace/oatproxy/pkg/oatproxy" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/bus" @@ -14,7 +12,6 @@ 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/statedb" ) @@ -33,10 +30,9 @@ type Director struct { streamSessionsMu sync.Mutex op *oatproxy.OATProxy statefulDB *statedb.StatefulDB - swarm *iroh_replicator.SwarmKV } -func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus *bus.Bus, op *oatproxy.OATProxy, statefulDB *statedb.StatefulDB, swarm *iroh_replicator.SwarmKV) *Director { +func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus *bus.Bus, op *oatproxy.OATProxy, statefulDB *statedb.StatefulDB) *Director { return &Director{ mm: mm, mod: mod, @@ -46,16 +42,10 @@ func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus * streamSessionsMu: sync.Mutex{}, op: op, statefulDB: statefulDB, - swarm: swarm, } } func (d *Director) Start(ctx context.Context) error { - nodeId, err := d.swarm.Node.NodeId() - if err != nil { - return fmt.Errorf("failed to get node id: %w", err) - } - newSeg := d.mm.NewSegment() ctx, cancel := context.WithCancel(ctx) defer cancel() @@ -96,27 +86,7 @@ func (d *Director) Start(ctx context.Context) error { }) } d.streamSessionsMu.Unlock() - go func() { - originInfo := iroh_replicator.OriginInfo{ - NodeID: nodeId.String(), - Time: not.Segment.StartTime.Format(util.ISO8601), - } - bs, err := json.Marshal(originInfo) - if err != nil { - log.Error(ctx, "could not marshal origin info", "error", err) - return - } - err = d.swarm.Put(ctx, not.Segment.RepoDID, bs) - if err != nil { - log.Error(ctx, "could not put segment to swarm", "error", err) - return - } - err = d.swarm.Node.SendSegment(not.Segment.RepoDID, not.Data) - if err != nil { - log.Error(ctx, "could not send segment to swarm", "error", err) - return - } - }() + err := ss.NewSegment(ctx, not) if err != nil { log.Error(ctx, "could not add segment to stream session", "error", err) diff --git a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go index e4125b08..d2b1feca 100644 --- a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go +++ b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go @@ -337,7 +337,6 @@ func readFloat64(reader io.Reader) float64 { func init() { FfiConverterDataHandlerINSTANCE.register() - FfiConverterDataHandlerOldINSTANCE.register() FfiConverterGoSignerINSTANCE.register() uniffiCheckChecksums() } @@ -389,15 +388,6 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_datahandler_handle_data: UniFFI API checksum mismatch") } } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_datahandlerold_handle_data() - }) - if checksum != 20343 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_datahandlerold_handle_data: UniFFI API checksum mismatch") - } - } { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_method_db_iter_with_opts() @@ -434,15 +424,6 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_db_write: UniFFI API checksum mismatch") } } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_endpoint_node_addr() - }) - if checksum != 17254 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_endpoint_node_addr: UniFFI API checksum mismatch") - } - } { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_method_filter_global() @@ -650,51 +631,6 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_publickey_fmt_short: UniFFI API checksum mismatch") } } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_receiver_node_addr() - }) - if checksum != 10730 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_receiver_node_addr: UniFFI API checksum mismatch") - } - } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_receiver_subscribe() - }) - if checksum != 24145 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_receiver_subscribe: UniFFI API checksum mismatch") - } - } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_receiver_unsubscribe() - }) - if checksum != 21760 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_receiver_unsubscribe: UniFFI API checksum mismatch") - } - } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_sender_node_addr() - }) - if checksum != 38541 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_sender_node_addr: UniFFI API checksum mismatch") - } - } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_method_sender_send() - }) - if checksum != 23930 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_sender_send: UniFFI API checksum mismatch") - } - } { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_method_subscriberesponse_next_raw() @@ -713,15 +649,6 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_writescope_put: UniFFI API checksum mismatch") } } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_constructor_endpoint_new() - }) - if checksum != 60672 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_constructor_endpoint_new: UniFFI API checksum mismatch") - } - } { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_constructor_filter_new() @@ -785,24 +712,6 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_constructor_publickey_from_string: UniFFI API checksum mismatch") } } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_constructor_receiver_new() - }) - if checksum != 18072 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_constructor_receiver_new: UniFFI API checksum mismatch") - } - } - { - checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { - return C.uniffi_iroh_streamplace_checksum_constructor_sender_new() - }) - if checksum != 56457 { - // If this happens try cleaning and rebuilding your project - panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_constructor_sender_new: UniFFI API checksum mismatch") - } - } } type FfiConverterUint64 struct{} @@ -1251,153 +1160,6 @@ func (c FfiConverterDataHandler) register() { C.uniffi_iroh_streamplace_fn_init_callback_vtable_datahandler(&UniffiVTableCallbackInterfaceDataHandlerINSTANCE) } -type DataHandlerOld interface { - HandleData(topic string, data []byte) -} -type DataHandlerOldImpl struct { - ffiObject FfiObject -} - -func (_self *DataHandlerOldImpl) HandleData(topic string, data []byte) { - _pointer := _self.ffiObject.incrementPointer("DataHandlerOld") - defer _self.ffiObject.decrementPointer() - uniffiRustCallAsync[error]( - nil, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) struct{} { - C.ffi_iroh_streamplace_rust_future_complete_void(handle, status) - return struct{}{} - }, - // liftFn - func(_ struct{}) struct{} { return struct{}{} }, - C.uniffi_iroh_streamplace_fn_method_datahandlerold_handle_data( - _pointer, FfiConverterStringINSTANCE.Lower(topic), FfiConverterBytesINSTANCE.Lower(data)), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_void(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_void(handle) - }, - ) - -} -func (object *DataHandlerOldImpl) Destroy() { - runtime.SetFinalizer(object, nil) - object.ffiObject.destroy() -} - -type FfiConverterDataHandlerOld struct { - handleMap *concurrentHandleMap[DataHandlerOld] -} - -var FfiConverterDataHandlerOldINSTANCE = FfiConverterDataHandlerOld{ - handleMap: newConcurrentHandleMap[DataHandlerOld](), -} - -func (c FfiConverterDataHandlerOld) Lift(pointer unsafe.Pointer) DataHandlerOld { - result := &DataHandlerOldImpl{ - newFfiObject( - pointer, - func(pointer unsafe.Pointer, status *C.RustCallStatus) unsafe.Pointer { - return C.uniffi_iroh_streamplace_fn_clone_datahandlerold(pointer, status) - }, - func(pointer unsafe.Pointer, status *C.RustCallStatus) { - C.uniffi_iroh_streamplace_fn_free_datahandlerold(pointer, status) - }, - ), - } - runtime.SetFinalizer(result, (*DataHandlerOldImpl).Destroy) - return result -} - -func (c FfiConverterDataHandlerOld) Read(reader io.Reader) DataHandlerOld { - return c.Lift(unsafe.Pointer(uintptr(readUint64(reader)))) -} - -func (c FfiConverterDataHandlerOld) Lower(value DataHandlerOld) unsafe.Pointer { - // TODO: this is bad - all synchronization from ObjectRuntime.go is discarded here, - // because the pointer will be decremented immediately after this function returns, - // and someone will be left holding onto a non-locked pointer. - pointer := unsafe.Pointer(uintptr(c.handleMap.insert(value))) - return pointer - -} - -func (c FfiConverterDataHandlerOld) Write(writer io.Writer, value DataHandlerOld) { - writeUint64(writer, uint64(uintptr(c.Lower(value)))) -} - -type FfiDestroyerDataHandlerOld struct{} - -func (_ FfiDestroyerDataHandlerOld) Destroy(value DataHandlerOld) { - if val, ok := value.(*DataHandlerOldImpl); ok { - val.Destroy() - } else { - panic("Expected *DataHandlerOldImpl") - } -} - -//export iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldMethod0 -func iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldMethod0(uniffiHandle C.uint64_t, topic C.RustBuffer, data C.RustBuffer, uniffiFutureCallback C.UniffiForeignFutureCompleteVoid, uniffiCallbackData C.uint64_t, uniffiOutReturn *C.UniffiForeignFuture) { - handle := uint64(uniffiHandle) - uniffiObj, ok := FfiConverterDataHandlerOldINSTANCE.handleMap.tryGet(handle) - if !ok { - panic(fmt.Errorf("no callback in handle map: %d", handle)) - } - - result := make(chan C.UniffiForeignFutureStructVoid, 1) - cancel := make(chan struct{}, 1) - guardHandle := cgo.NewHandle(cancel) - *uniffiOutReturn = C.UniffiForeignFuture{ - handle: C.uint64_t(guardHandle), - free: C.UniffiForeignFutureFree(C.iroh_streamplace_uniffiFreeGorutine), - } - - // Wait for compleation or cancel - go func() { - select { - case <-cancel: - case res := <-result: - C.call_UniffiForeignFutureCompleteVoid(uniffiFutureCallback, uniffiCallbackData, res) - } - }() - - // Eval callback asynchroniously - go func() { - asyncResult := &C.UniffiForeignFutureStructVoid{} - defer func() { - result <- *asyncResult - }() - - uniffiObj.HandleData( - FfiConverterStringINSTANCE.Lift(GoRustBuffer{ - inner: topic, - }), - FfiConverterBytesINSTANCE.Lift(GoRustBuffer{ - inner: data, - }), - ) - - }() -} - -var UniffiVTableCallbackInterfaceDataHandlerOldINSTANCE = C.UniffiVTableCallbackInterfaceDataHandlerOld{ - handleData: (C.UniffiCallbackInterfaceDataHandlerOldMethod0)(C.iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldMethod0), - - uniffiFree: (C.UniffiCallbackInterfaceFree)(C.iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldFree), -} - -//export iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldFree -func iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldFree(handle C.uint64_t) { - FfiConverterDataHandlerOldINSTANCE.handleMap.remove(uint64(handle)) -} - -func (c FfiConverterDataHandlerOld) register() { - C.uniffi_iroh_streamplace_fn_init_callback_vtable_datahandlerold(&UniffiVTableCallbackInterfaceDataHandlerOldINSTANCE) -} - // Iroh-streamplace specific metadata database. type DbInterface interface { IterWithOpts(filter *Filter) ([]Entry, error) @@ -1529,121 +1291,6 @@ func (_ FfiDestroyerDb) Destroy(value *Db) { value.Destroy() } -type EndpointInterface interface { - NodeAddr() *NodeAddr -} -type Endpoint struct { - ffiObject FfiObject -} - -// Create a new endpoint. -func NewEndpoint() (*Endpoint, error) { - res, err := uniffiRustCallAsync[Error]( - FfiConverterErrorINSTANCE, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) unsafe.Pointer { - res := C.ffi_iroh_streamplace_rust_future_complete_pointer(handle, status) - return res - }, - // liftFn - func(ffi unsafe.Pointer) *Endpoint { - return FfiConverterEndpointINSTANCE.Lift(ffi) - }, - C.uniffi_iroh_streamplace_fn_constructor_endpoint_new(), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_pointer(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_pointer(handle) - }, - ) - - if err == nil { - return res, nil - } - - return res, err -} - -func (_self *Endpoint) NodeAddr() *NodeAddr { - _pointer := _self.ffiObject.incrementPointer("*Endpoint") - defer _self.ffiObject.decrementPointer() - res, _ := uniffiRustCallAsync[error]( - nil, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) unsafe.Pointer { - res := C.ffi_iroh_streamplace_rust_future_complete_pointer(handle, status) - return res - }, - // liftFn - func(ffi unsafe.Pointer) *NodeAddr { - return FfiConverterNodeAddrINSTANCE.Lift(ffi) - }, - C.uniffi_iroh_streamplace_fn_method_endpoint_node_addr( - _pointer), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_pointer(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_pointer(handle) - }, - ) - - return res -} -func (object *Endpoint) Destroy() { - runtime.SetFinalizer(object, nil) - object.ffiObject.destroy() -} - -type FfiConverterEndpoint struct{} - -var FfiConverterEndpointINSTANCE = FfiConverterEndpoint{} - -func (c FfiConverterEndpoint) Lift(pointer unsafe.Pointer) *Endpoint { - result := &Endpoint{ - newFfiObject( - pointer, - func(pointer unsafe.Pointer, status *C.RustCallStatus) unsafe.Pointer { - return C.uniffi_iroh_streamplace_fn_clone_endpoint(pointer, status) - }, - func(pointer unsafe.Pointer, status *C.RustCallStatus) { - C.uniffi_iroh_streamplace_fn_free_endpoint(pointer, status) - }, - ), - } - runtime.SetFinalizer(result, (*Endpoint).Destroy) - return result -} - -func (c FfiConverterEndpoint) Read(reader io.Reader) *Endpoint { - return c.Lift(unsafe.Pointer(uintptr(readUint64(reader)))) -} - -func (c FfiConverterEndpoint) Lower(value *Endpoint) unsafe.Pointer { - // TODO: this is bad - all synchronization from ObjectRuntime.go is discarded here, - // because the pointer will be decremented immediately after this function returns, - // and someone will be left holding onto a non-locked pointer. - pointer := value.ffiObject.incrementPointer("*Endpoint") - defer value.ffiObject.decrementPointer() - return pointer - -} - -func (c FfiConverterEndpoint) Write(writer io.Writer, value *Endpoint) { - writeUint64(writer, uint64(uintptr(c.Lower(value)))) -} - -type FfiDestroyerEndpoint struct{} - -func (_ FfiDestroyerEndpoint) Destroy(value *Endpoint) { - value.Destroy() -} - // A filter for subscriptions and iteration. type FilterInterface interface { // Restrict to the global namespace, no per stream data. @@ -2562,334 +2209,6 @@ func (_ FfiDestroyerPublicKey) Destroy(value *PublicKey) { value.Destroy() } -type ReceiverInterface interface { - NodeAddr() *NodeAddr - // Subscribe to the given topic on the remote. - Subscribe(remoteId *PublicKey, topic string) error - // Unsubscribe from this topic on the remote. - Unsubscribe(remoteId *PublicKey, topic string) error -} -type Receiver struct { - ffiObject FfiObject -} - -// Create a new receiver. -func NewReceiver(endpoint *Endpoint, handler DataHandlerOld) (*Receiver, error) { - res, err := uniffiRustCallAsync[Error]( - FfiConverterErrorINSTANCE, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) unsafe.Pointer { - res := C.ffi_iroh_streamplace_rust_future_complete_pointer(handle, status) - return res - }, - // liftFn - func(ffi unsafe.Pointer) *Receiver { - return FfiConverterReceiverINSTANCE.Lift(ffi) - }, - C.uniffi_iroh_streamplace_fn_constructor_receiver_new(FfiConverterEndpointINSTANCE.Lower(endpoint), FfiConverterDataHandlerOldINSTANCE.Lower(handler)), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_pointer(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_pointer(handle) - }, - ) - - if err == nil { - return res, nil - } - - return res, err -} - -func (_self *Receiver) NodeAddr() *NodeAddr { - _pointer := _self.ffiObject.incrementPointer("*Receiver") - defer _self.ffiObject.decrementPointer() - res, _ := uniffiRustCallAsync[error]( - nil, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) unsafe.Pointer { - res := C.ffi_iroh_streamplace_rust_future_complete_pointer(handle, status) - return res - }, - // liftFn - func(ffi unsafe.Pointer) *NodeAddr { - return FfiConverterNodeAddrINSTANCE.Lift(ffi) - }, - C.uniffi_iroh_streamplace_fn_method_receiver_node_addr( - _pointer), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_pointer(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_pointer(handle) - }, - ) - - return res -} - -// Subscribe to the given topic on the remote. -func (_self *Receiver) Subscribe(remoteId *PublicKey, topic string) error { - _pointer := _self.ffiObject.incrementPointer("*Receiver") - defer _self.ffiObject.decrementPointer() - _, err := uniffiRustCallAsync[Error]( - FfiConverterErrorINSTANCE, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) struct{} { - C.ffi_iroh_streamplace_rust_future_complete_void(handle, status) - return struct{}{} - }, - // liftFn - func(_ struct{}) struct{} { return struct{}{} }, - C.uniffi_iroh_streamplace_fn_method_receiver_subscribe( - _pointer, FfiConverterPublicKeyINSTANCE.Lower(remoteId), FfiConverterStringINSTANCE.Lower(topic)), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_void(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_void(handle) - }, - ) - - if err == nil { - return nil - } - - return err -} - -// Unsubscribe from this topic on the remote. -func (_self *Receiver) Unsubscribe(remoteId *PublicKey, topic string) error { - _pointer := _self.ffiObject.incrementPointer("*Receiver") - defer _self.ffiObject.decrementPointer() - _, err := uniffiRustCallAsync[Error]( - FfiConverterErrorINSTANCE, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) struct{} { - C.ffi_iroh_streamplace_rust_future_complete_void(handle, status) - return struct{}{} - }, - // liftFn - func(_ struct{}) struct{} { return struct{}{} }, - C.uniffi_iroh_streamplace_fn_method_receiver_unsubscribe( - _pointer, FfiConverterPublicKeyINSTANCE.Lower(remoteId), FfiConverterStringINSTANCE.Lower(topic)), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_void(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_void(handle) - }, - ) - - if err == nil { - return nil - } - - return err -} -func (object *Receiver) Destroy() { - runtime.SetFinalizer(object, nil) - object.ffiObject.destroy() -} - -type FfiConverterReceiver struct{} - -var FfiConverterReceiverINSTANCE = FfiConverterReceiver{} - -func (c FfiConverterReceiver) Lift(pointer unsafe.Pointer) *Receiver { - result := &Receiver{ - newFfiObject( - pointer, - func(pointer unsafe.Pointer, status *C.RustCallStatus) unsafe.Pointer { - return C.uniffi_iroh_streamplace_fn_clone_receiver(pointer, status) - }, - func(pointer unsafe.Pointer, status *C.RustCallStatus) { - C.uniffi_iroh_streamplace_fn_free_receiver(pointer, status) - }, - ), - } - runtime.SetFinalizer(result, (*Receiver).Destroy) - return result -} - -func (c FfiConverterReceiver) Read(reader io.Reader) *Receiver { - return c.Lift(unsafe.Pointer(uintptr(readUint64(reader)))) -} - -func (c FfiConverterReceiver) Lower(value *Receiver) unsafe.Pointer { - // TODO: this is bad - all synchronization from ObjectRuntime.go is discarded here, - // because the pointer will be decremented immediately after this function returns, - // and someone will be left holding onto a non-locked pointer. - pointer := value.ffiObject.incrementPointer("*Receiver") - defer value.ffiObject.decrementPointer() - return pointer - -} - -func (c FfiConverterReceiver) Write(writer io.Writer, value *Receiver) { - writeUint64(writer, uint64(uintptr(c.Lower(value)))) -} - -type FfiDestroyerReceiver struct{} - -func (_ FfiDestroyerReceiver) Destroy(value *Receiver) { - value.Destroy() -} - -type SenderInterface interface { - NodeAddr() *NodeAddr - // Sends the given data to all subscribers that have subscribed to this `key`. - Send(key string, data []byte) error -} -type Sender struct { - ffiObject FfiObject -} - -// Create a new sender. -func NewSender(endpoint *Endpoint) *Sender { - res, _ := uniffiRustCallAsync[error]( - nil, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) unsafe.Pointer { - res := C.ffi_iroh_streamplace_rust_future_complete_pointer(handle, status) - return res - }, - // liftFn - func(ffi unsafe.Pointer) *Sender { - return FfiConverterSenderINSTANCE.Lift(ffi) - }, - C.uniffi_iroh_streamplace_fn_constructor_sender_new(FfiConverterEndpointINSTANCE.Lower(endpoint)), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_pointer(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_pointer(handle) - }, - ) - - return res -} - -func (_self *Sender) NodeAddr() *NodeAddr { - _pointer := _self.ffiObject.incrementPointer("*Sender") - defer _self.ffiObject.decrementPointer() - res, _ := uniffiRustCallAsync[error]( - nil, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) unsafe.Pointer { - res := C.ffi_iroh_streamplace_rust_future_complete_pointer(handle, status) - return res - }, - // liftFn - func(ffi unsafe.Pointer) *NodeAddr { - return FfiConverterNodeAddrINSTANCE.Lift(ffi) - }, - C.uniffi_iroh_streamplace_fn_method_sender_node_addr( - _pointer), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_pointer(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_pointer(handle) - }, - ) - - return res -} - -// Sends the given data to all subscribers that have subscribed to this `key`. -func (_self *Sender) Send(key string, data []byte) error { - _pointer := _self.ffiObject.incrementPointer("*Sender") - defer _self.ffiObject.decrementPointer() - _, err := uniffiRustCallAsync[Error]( - FfiConverterErrorINSTANCE, - // completeFn - func(handle C.uint64_t, status *C.RustCallStatus) struct{} { - C.ffi_iroh_streamplace_rust_future_complete_void(handle, status) - return struct{}{} - }, - // liftFn - func(_ struct{}) struct{} { return struct{}{} }, - C.uniffi_iroh_streamplace_fn_method_sender_send( - _pointer, FfiConverterStringINSTANCE.Lower(key), FfiConverterBytesINSTANCE.Lower(data)), - // pollFn - func(handle C.uint64_t, continuation C.UniffiRustFutureContinuationCallback, data C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_poll_void(handle, continuation, data) - }, - // freeFn - func(handle C.uint64_t) { - C.ffi_iroh_streamplace_rust_future_free_void(handle) - }, - ) - - if err == nil { - return nil - } - - return err -} -func (object *Sender) Destroy() { - runtime.SetFinalizer(object, nil) - object.ffiObject.destroy() -} - -type FfiConverterSender struct{} - -var FfiConverterSenderINSTANCE = FfiConverterSender{} - -func (c FfiConverterSender) Lift(pointer unsafe.Pointer) *Sender { - result := &Sender{ - newFfiObject( - pointer, - func(pointer unsafe.Pointer, status *C.RustCallStatus) unsafe.Pointer { - return C.uniffi_iroh_streamplace_fn_clone_sender(pointer, status) - }, - func(pointer unsafe.Pointer, status *C.RustCallStatus) { - C.uniffi_iroh_streamplace_fn_free_sender(pointer, status) - }, - ), - } - runtime.SetFinalizer(result, (*Sender).Destroy) - return result -} - -func (c FfiConverterSender) Read(reader io.Reader) *Sender { - return c.Lift(unsafe.Pointer(uintptr(readUint64(reader)))) -} - -func (c FfiConverterSender) Lower(value *Sender) unsafe.Pointer { - // TODO: this is bad - all synchronization from ObjectRuntime.go is discarded here, - // because the pointer will be decremented immediately after this function returns, - // and someone will be left holding onto a non-locked pointer. - pointer := value.ffiObject.incrementPointer("*Sender") - defer value.ffiObject.decrementPointer() - return pointer - -} - -func (c FfiConverterSender) Write(writer io.Writer, value *Sender) { - writeUint64(writer, uint64(uintptr(c.Lower(value)))) -} - -type FfiDestroyerSender struct{} - -func (_ FfiDestroyerSender) Destroy(value *Sender) { - value.Destroy() -} - // A response to a subscribe request. // // This can be used as a stream of [`SubscribeItem`]s. diff --git a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h index 49e04249..7f735d8a 100644 --- a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h +++ b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h @@ -391,20 +391,6 @@ static void call_UniffiCallbackInterfaceDataHandlerMethod0( } -#endif -#ifndef UNIFFI_FFIDEF_CALLBACK_INTERFACE_DATA_HANDLER_OLD_METHOD0 -#define UNIFFI_FFIDEF_CALLBACK_INTERFACE_DATA_HANDLER_OLD_METHOD0 -typedef void (*UniffiCallbackInterfaceDataHandlerOldMethod0)(uint64_t uniffi_handle, RustBuffer topic, RustBuffer data, UniffiForeignFutureCompleteVoid uniffi_future_callback, uint64_t uniffi_callback_data, UniffiForeignFuture* uniffi_out_return); - -// Making function static works arround: -// https://github.com/golang/go/issues/11263 -static void call_UniffiCallbackInterfaceDataHandlerOldMethod0( - UniffiCallbackInterfaceDataHandlerOldMethod0 cb, uint64_t uniffi_handle, RustBuffer topic, RustBuffer data, UniffiForeignFutureCompleteVoid uniffi_future_callback, uint64_t uniffi_callback_data, UniffiForeignFuture* uniffi_out_return) -{ - return cb(uniffi_handle, topic, data, uniffi_future_callback, uniffi_callback_data, uniffi_out_return); -} - - #endif #ifndef UNIFFI_FFIDEF_CALLBACK_INTERFACE_GO_SIGNER_METHOD0 #define UNIFFI_FFIDEF_CALLBACK_INTERFACE_GO_SIGNER_METHOD0 @@ -427,14 +413,6 @@ typedef struct UniffiVTableCallbackInterfaceDataHandler { UniffiCallbackInterfaceFree uniffiFree; } UniffiVTableCallbackInterfaceDataHandler; -#endif -#ifndef UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_DATA_HANDLER_OLD -#define UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_DATA_HANDLER_OLD -typedef struct UniffiVTableCallbackInterfaceDataHandlerOld { - UniffiCallbackInterfaceDataHandlerOldMethod0 handleData; - UniffiCallbackInterfaceFree uniffiFree; -} UniffiVTableCallbackInterfaceDataHandlerOld; - #endif #ifndef UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_GO_SIGNER #define UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_GO_SIGNER @@ -464,26 +442,6 @@ void uniffi_iroh_streamplace_fn_init_callback_vtable_datahandler(UniffiVTableCal uint64_t uniffi_iroh_streamplace_fn_method_datahandler_handle_data(void* ptr, RustBuffer topic, RustBuffer data ); #endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_DATAHANDLEROLD -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_DATAHANDLEROLD -void* uniffi_iroh_streamplace_fn_clone_datahandlerold(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_DATAHANDLEROLD -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_DATAHANDLEROLD -void uniffi_iroh_streamplace_fn_free_datahandlerold(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_INIT_CALLBACK_VTABLE_DATAHANDLEROLD -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_INIT_CALLBACK_VTABLE_DATAHANDLEROLD -void uniffi_iroh_streamplace_fn_init_callback_vtable_datahandlerold(UniffiVTableCallbackInterfaceDataHandlerOld* vtable -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_DATAHANDLEROLD_HANDLE_DATA -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_DATAHANDLEROLD_HANDLE_DATA -uint64_t uniffi_iroh_streamplace_fn_method_datahandlerold_handle_data(void* ptr, RustBuffer topic, RustBuffer data -); -#endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_DB #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_DB void* uniffi_iroh_streamplace_fn_clone_db(void* ptr, RustCallStatus *out_status @@ -514,27 +472,6 @@ void* uniffi_iroh_streamplace_fn_method_db_subscribe_with_opts(void* ptr, RustBu void* uniffi_iroh_streamplace_fn_method_db_write(void* ptr, RustBuffer secret, RustCallStatus *out_status ); #endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_ENDPOINT -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_ENDPOINT -void* uniffi_iroh_streamplace_fn_clone_endpoint(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_ENDPOINT -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_ENDPOINT -void uniffi_iroh_streamplace_fn_free_endpoint(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CONSTRUCTOR_ENDPOINT_NEW -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CONSTRUCTOR_ENDPOINT_NEW -uint64_t uniffi_iroh_streamplace_fn_constructor_endpoint_new(void - -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_ENDPOINT_NODE_ADDR -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_ENDPOINT_NODE_ADDR -uint64_t uniffi_iroh_streamplace_fn_method_endpoint_node_addr(void* ptr -); -#endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_FILTER #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_FILTER void* uniffi_iroh_streamplace_fn_clone_filter(void* ptr, RustCallStatus *out_status @@ -746,61 +683,6 @@ RustBuffer uniffi_iroh_streamplace_fn_method_publickey_fmt_short(void* ptr, Rust RustBuffer uniffi_iroh_streamplace_fn_method_publickey_uniffi_trait_display(void* ptr, RustCallStatus *out_status ); #endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_RECEIVER -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_RECEIVER -void* uniffi_iroh_streamplace_fn_clone_receiver(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_RECEIVER -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_RECEIVER -void uniffi_iroh_streamplace_fn_free_receiver(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CONSTRUCTOR_RECEIVER_NEW -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CONSTRUCTOR_RECEIVER_NEW -uint64_t uniffi_iroh_streamplace_fn_constructor_receiver_new(void* endpoint, void* handler -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_RECEIVER_NODE_ADDR -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_RECEIVER_NODE_ADDR -uint64_t uniffi_iroh_streamplace_fn_method_receiver_node_addr(void* ptr -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_RECEIVER_SUBSCRIBE -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_RECEIVER_SUBSCRIBE -uint64_t uniffi_iroh_streamplace_fn_method_receiver_subscribe(void* ptr, void* remote_id, RustBuffer topic -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_RECEIVER_UNSUBSCRIBE -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_RECEIVER_UNSUBSCRIBE -uint64_t uniffi_iroh_streamplace_fn_method_receiver_unsubscribe(void* ptr, void* remote_id, RustBuffer topic -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_SENDER -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_SENDER -void* uniffi_iroh_streamplace_fn_clone_sender(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_SENDER -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_SENDER -void uniffi_iroh_streamplace_fn_free_sender(void* ptr, RustCallStatus *out_status -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CONSTRUCTOR_SENDER_NEW -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CONSTRUCTOR_SENDER_NEW -uint64_t uniffi_iroh_streamplace_fn_constructor_sender_new(void* endpoint -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_SENDER_NODE_ADDR -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_SENDER_NODE_ADDR -uint64_t uniffi_iroh_streamplace_fn_method_sender_node_addr(void* ptr -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_SENDER_SEND -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_SENDER_SEND -uint64_t uniffi_iroh_streamplace_fn_method_sender_send(void* ptr, RustBuffer key, RustBuffer data -); -#endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_SUBSCRIBERESPONSE #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_SUBSCRIBERESPONSE void* uniffi_iroh_streamplace_fn_clone_subscriberesponse(void* ptr, RustCallStatus *out_status @@ -1153,12 +1035,6 @@ uint16_t uniffi_iroh_streamplace_checksum_func_subscribe_item_debug(void #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_DATAHANDLER_HANDLE_DATA uint16_t uniffi_iroh_streamplace_checksum_method_datahandler_handle_data(void -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_DATAHANDLEROLD_HANDLE_DATA -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_DATAHANDLEROLD_HANDLE_DATA -uint16_t uniffi_iroh_streamplace_checksum_method_datahandlerold_handle_data(void - ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_DB_ITER_WITH_OPTS @@ -1183,12 +1059,6 @@ uint16_t uniffi_iroh_streamplace_checksum_method_db_subscribe_with_opts(void #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_DB_WRITE uint16_t uniffi_iroh_streamplace_checksum_method_db_write(void -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_ENDPOINT_NODE_ADDR -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_ENDPOINT_NODE_ADDR -uint16_t uniffi_iroh_streamplace_checksum_method_endpoint_node_addr(void - ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_FILTER_GLOBAL @@ -1327,36 +1197,6 @@ uint16_t uniffi_iroh_streamplace_checksum_method_publickey_equal(void #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_PUBLICKEY_FMT_SHORT uint16_t uniffi_iroh_streamplace_checksum_method_publickey_fmt_short(void -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_RECEIVER_NODE_ADDR -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_RECEIVER_NODE_ADDR -uint16_t uniffi_iroh_streamplace_checksum_method_receiver_node_addr(void - -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_RECEIVER_SUBSCRIBE -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_RECEIVER_SUBSCRIBE -uint16_t uniffi_iroh_streamplace_checksum_method_receiver_subscribe(void - -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_RECEIVER_UNSUBSCRIBE -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_RECEIVER_UNSUBSCRIBE -uint16_t uniffi_iroh_streamplace_checksum_method_receiver_unsubscribe(void - -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_SENDER_NODE_ADDR -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_SENDER_NODE_ADDR -uint16_t uniffi_iroh_streamplace_checksum_method_sender_node_addr(void - -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_SENDER_SEND -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_SENDER_SEND -uint16_t uniffi_iroh_streamplace_checksum_method_sender_send(void - ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_SUBSCRIBERESPONSE_NEXT_RAW @@ -1369,12 +1209,6 @@ uint16_t uniffi_iroh_streamplace_checksum_method_subscriberesponse_next_raw(void #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_WRITESCOPE_PUT uint16_t uniffi_iroh_streamplace_checksum_method_writescope_put(void -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_ENDPOINT_NEW -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_ENDPOINT_NEW -uint16_t uniffi_iroh_streamplace_checksum_constructor_endpoint_new(void - ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_FILTER_NEW @@ -1417,18 +1251,6 @@ uint16_t uniffi_iroh_streamplace_checksum_constructor_publickey_from_bytes(void #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_PUBLICKEY_FROM_STRING uint16_t uniffi_iroh_streamplace_checksum_constructor_publickey_from_string(void -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_RECEIVER_NEW -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_RECEIVER_NEW -uint16_t uniffi_iroh_streamplace_checksum_constructor_receiver_new(void - -); -#endif -#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_SENDER_NEW -#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_CONSTRUCTOR_SENDER_NEW -uint16_t uniffi_iroh_streamplace_checksum_constructor_sender_new(void - ); #endif #ifndef UNIFFI_FFIDEF_FFI_IROH_STREAMPLACE_UNIFFI_CONTRACT_VERSION @@ -1440,8 +1262,6 @@ uint32_t ffi_iroh_streamplace_uniffi_contract_version(void void iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerMethod0(uint64_t uniffi_handle, RustBuffer topic, RustBuffer data, UniffiForeignFutureCompleteVoid uniffi_future_callback, uint64_t uniffi_callback_data, UniffiForeignFuture* uniffi_out_return); void iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerFree(uint64_t handle); - void iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldMethod0(uint64_t uniffi_handle, RustBuffer topic, RustBuffer data, UniffiForeignFutureCompleteVoid uniffi_future_callback, uint64_t uniffi_callback_data, UniffiForeignFuture* uniffi_out_return); - void iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerOldFree(uint64_t handle); void iroh_streamplace_cgo_dispatchCallbackInterfaceGoSignerMethod0(uint64_t uniffi_handle, RustBuffer data, RustBuffer* uniffi_out_return, RustCallStatus* callStatus ); void iroh_streamplace_cgo_dispatchCallbackInterfaceGoSignerFree(uint64_t handle); diff --git a/pkg/iroh/iroh_streamplace_test.go b/pkg/iroh/iroh_streamplace_test.go deleted file mode 100644 index 41736876..00000000 --- a/pkg/iroh/iroh_streamplace_test.go +++ /dev/null @@ -1,56 +0,0 @@ -package iroh_streamplace - -import ( - "testing" - - "github.com/stretchr/testify/assert" - - iroh "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" -) - -type Message struct { - topic string - data []byte -} - -type TestHandler struct { - messages chan Message -} - -func (handler TestHandler) HandleData(topic string, data []byte) { - handler.messages <- Message{topic, data} -} - -func TestBasicRoundtrip(t *testing.T) { - ep1, err := iroh.NewEndpoint() - assert.Nil(t, err) - sender := iroh.NewSender(ep1) - - messages := make(chan Message, 5) - handler := TestHandler{messages: messages} - - ep2, err := iroh.NewEndpoint() - assert.NoError(t, err) - receiver, err := iroh.NewReceiver(ep2, &handler) - assert.NoError(t, err) - - senderAddr := sender.NodeAddr() - senderId := senderAddr.NodeId() - - // subscribe - err = receiver.Subscribe(senderId, "foo") - assert.NoError(t, err) - - // send a few messages - for i := range 5 { - err = sender.Send("foo", []byte{byte(i), 0, 0, 0}) - assert.NoError(t, err) - } - - // make sure the receiver got them - for i := range 5 { - msg := <-messages - assert.Equal(t, msg.topic, "foo") - assert.Equal(t, msg.data, []byte{byte(i), 0, 0, 0}) - } -} diff --git a/pkg/media/media.go b/pkg/media/media.go index a4240927..d9f71f5e 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -58,6 +58,7 @@ type NewSegmentNotification struct { Segment *model.Segment Data []byte Metadata *SegmentMetadata + Local bool } func RunSelfTest(ctx context.Context) error { @@ -65,7 +66,7 @@ func RunSelfTest(ctx context.Context) error { return SelfTest(ctx) } -func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer, rep replication.Replicator, mod model.Model, bus *bus.Bus, atsync *atproto.ATProtoSynchronizer) (*MediaManager, error) { +func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer, mod model.Model, bus *bus.Bus, atsync *atproto.ATProtoSynchronizer) (*MediaManager, error) { gstinit.InitGST() err := SelfTest(ctx) if err != nil { @@ -121,7 +122,6 @@ func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer return &MediaManager{ cli: cli, - replicator: rep, hlsRunning: map[string]*M3U8{}, httpPipes: map[string]io.Writer{}, model: mod, @@ -135,7 +135,7 @@ func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer func (mm *MediaManager) HandleData(node *irohStreamplace.PublicKey, data []byte) { r := bytes.NewReader(data) ctx := context.Background() - err := mm.ValidateMP4(ctx, r) + err := mm.ValidateMP4(ctx, r, true) if err != nil { log.Log(ctx, "invalid incoming segment", "error", err) } diff --git a/pkg/media/media_test.go b/pkg/media/media_test.go index 9883e54f..13f875aa 100644 --- a/pkg/media/media_test.go +++ b/pkg/media/media_test.go @@ -13,7 +13,6 @@ import ( "stream.place/streamplace/pkg/config" ct "stream.place/streamplace/pkg/config/configtesting" "stream.place/streamplace/pkg/model" - "stream.place/streamplace/pkg/replication/boring" ) func getFixture(name string) string { @@ -40,7 +39,7 @@ func getStaticTestMediaManager(t *testing.T) (*MediaManager, MediaSigner) { StatefulDB: nil, // Test doesn't need StatefulDB for now Bus: bus.NewBus(), } - mm, err := MakeMediaManager(context.Background(), cli, nil, &boring.BoringReplicator{}, mod, bus.NewBus(), atsync) + mm, err := MakeMediaManager(context.Background(), cli, nil, mod, bus.NewBus(), atsync) require.NoError(t, err) // ms, err := MakeMediaSigner(context.Background(), cli, "test-person", signer) // require.NoError(t, err) @@ -126,6 +125,6 @@ func TestVerifyMP4(t *testing.T) { f, err := os.Open(getFixture("sample-segment.mp4")) require.NoError(t, err) mm, _ := getStaticTestMediaManager(t) - err = mm.ValidateMP4(context.Background(), f) + err = mm.ValidateMP4(context.Background(), f, true) require.NoError(t, err) } diff --git a/pkg/media/rtcrec_test.go b/pkg/media/rtcrec_test.go index a495a89d..de5348ef 100644 --- a/pkg/media/rtcrec_test.go +++ b/pkg/media/rtcrec_test.go @@ -11,7 +11,6 @@ import ( "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/crypto/spkey" "stream.place/streamplace/pkg/globalerror" - "stream.place/streamplace/pkg/replication/boring" "stream.place/streamplace/pkg/rtcrec" ) @@ -26,7 +25,7 @@ func TestRTCRecording(t *testing.T) { fs := cli.NewFlagSet("rtcrec-test") err = cli.Parse(fs, []string{"--data-dir", dir, "-wide-open=true"}) require.NoError(t, err) - mm, err := MakeMediaManager(context.Background(), cli, nil, &boring.BoringReplicator{}, nil, nil, nil) + mm, err := MakeMediaManager(context.Background(), cli, nil, nil, nil, nil) require.NoError(t, err) priv, pub, err := spkey.GenerateStreamKey() require.NoError(t, err) diff --git a/pkg/media/segmenter.go b/pkg/media/segmenter.go index 485e2a21..0aedd714 100644 --- a/pkg/media/segmenter.go +++ b/pkg/media/segmenter.go @@ -87,7 +87,7 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) return } - err = mm.ValidateMP4(ctx, bytes.NewReader(bs)) + err = mm.ValidateMP4(ctx, bytes.NewReader(bs), true) if err != nil { log.Error(ctx, "error validating segment", "error", err) globalerror.GlobalError(err) diff --git a/pkg/media/validate.go b/pkg/media/validate.go index df17c39a..01223f1c 100644 --- a/pkg/media/validate.go +++ b/pkg/media/validate.go @@ -23,7 +23,7 @@ type ManifestAndCert struct { Cert string `json:"cert"` } -func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader) error { +func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader, local bool) error { ctx, span := otel.Tracer("signer").Start(ctx, "ValidateMP4") defer span.End() buf, err := io.ReadAll(input) @@ -82,7 +82,7 @@ func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader) error return err } defer fd.Close() - go mm.replicator.NewSegment(ctx, buf) + r := bytes.NewReader(buf) if _, err := io.Copy(fd, r); err != nil { return err @@ -102,6 +102,7 @@ func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader) error Segment: seg, Data: buf, Metadata: meta, + Local: local, } for _, ch := range mm.newSegmentSubs { go func() { ch <- not }() diff --git a/pkg/replication/iroh_replicator/iroh.go b/pkg/replication/iroh_replicator/iroh.go index 138ab0ea..d389744e 100644 --- a/pkg/replication/iroh_replicator/iroh.go +++ b/pkg/replication/iroh_replicator/iroh.go @@ -2,36 +2,17 @@ package iroh_replicator import ( "context" - - "stream.place/streamplace/pkg/log" - - irohStreamplace "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" ) // IrohReplicator implements the replication mechanism using iroh type IrohReplicator struct { - topic string - sender *irohStreamplace.Sender } -func NewIrohReplicator(ctx context.Context, ep *irohStreamplace.Endpoint, topic string) (*IrohReplicator, error) { - sender := irohStreamplace.NewSender(ep) +func NewIrohReplicator(ctx context.Context) (*IrohReplicator, error) { - return &IrohReplicator{ - topic: topic, - sender: sender, - }, nil + return &IrohReplicator{}, nil } func (rep *IrohReplicator) NewSegment(ctx context.Context, bs []byte) { - go func(topic string) { - err := sendSegment(rep.sender, topic, bs) - if err != nil { - log.Log(ctx, "error replicating segment", "error", err) - } - }(rep.topic) -} -func sendSegment(endpoint *irohStreamplace.Sender, topic string, bs []byte) error { - return endpoint.Send(topic, bs) } diff --git a/pkg/replication/iroh_replicator/kv.go b/pkg/replication/iroh_replicator/kv.go index e11d63bb..a804006f 100644 --- a/pkg/replication/iroh_replicator/kv.go +++ b/pkg/replication/iroh_replicator/kv.go @@ -1,19 +1,26 @@ package iroh_replicator import ( + "bytes" "context" "encoding/json" "fmt" "time" + "github.com/bluesky-social/indigo/util" + "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/media" ) -type SwarmKV struct { - Node *iroh_streamplace.Node - DB *iroh_streamplace.Db - w *iroh_streamplace.WriteScope +type IrohSwarm struct { + Node *iroh_streamplace.Node + DB *iroh_streamplace.Db + w *iroh_streamplace.WriteScope + mm *media.MediaManager + segChan chan *media.NewSegmentNotification + nodeId string } // A message saying "hey I ingested node data at this time" @@ -22,14 +29,7 @@ type OriginInfo struct { Time string `json:"time"` } -type DataHandler struct{} - -func (handler *DataHandler) HandleData(topic string, data []byte) { - log.Log(context.Background(), "HandleData", "topic", topic, "data", len(data)) -} - -func StartKV(ctx context.Context, tickets []string, secret []byte) (*SwarmKV, error) { - handler := &DataHandler{} +func NewSwarm(ctx context.Context, tickets []string, secret []byte, mm *media.MediaManager) (*IrohSwarm, error) { ctx = log.WithLogValues(ctx, "func", "StartKV") log.Log(ctx, "Starting with tickets", "tickets", tickets) @@ -39,7 +39,12 @@ func StartKV(ctx context.Context, tickets []string, secret []byte) (*SwarmKV, er MaxSendDuration: 1000_000_000, // 1s } log.Log(ctx, "Config created", "config", config) - node, err := iroh_streamplace.NodeReceiver(config, handler) + + swarm := IrohSwarm{ + mm: mm, + } + + node, err := iroh_streamplace.NodeReceiver(config, &swarm) if err != nil { return nil, fmt.Errorf("failed to create NodeSender: %w", err) } @@ -47,11 +52,16 @@ func StartKV(ctx context.Context, tickets []string, secret []byte) (*SwarmKV, er db := node.Db() w := node.NodeScope() - node_id, err := node.NodeId() + swarm.DB = db + swarm.w = w + swarm.Node = node + + nodeId, err := node.NodeId() if err != nil { return nil, fmt.Errorf("failed to get NodeId: %w", err) } - log.Log(ctx, "Node ID:", "node_id", node_id) + log.Log(ctx, "Node ID:", "node_id", nodeId) + swarm.nodeId = nodeId.String() ticket, err := node.Ticket() if err != nil { @@ -59,17 +69,12 @@ func StartKV(ctx context.Context, tickets []string, secret []byte) (*SwarmKV, er } log.Log(ctx, "Ticket:", "ticket", ticket) - swarm := SwarmKV{ - Node: node, - DB: db, - w: w, - } return &swarm, nil } var activeSubs = make(map[string]bool) -func (swarm *SwarmKV) Start(ctx context.Context, tickets []string) error { +func (swarm *IrohSwarm) Start(ctx context.Context, tickets []string) error { if len(tickets) > 0 { err := swarm.Node.JoinPeers(tickets) if err != nil { @@ -84,6 +89,17 @@ func (swarm *SwarmKV) Start(ctx context.Context, tickets []string) error { nodeIdStr := nodeId.String() log.Log(ctx, "Node ID:", "node_id", nodeIdStr) + g, ctx := errgroup.WithContext(ctx) + g.Go(func() error { + return swarm.startKV(ctx) + }) + g.Go(func() error { + return swarm.startSegmentSender(ctx) + }) + return g.Wait() +} + +func (swarm *IrohSwarm) startKV(ctx context.Context) error { sub := swarm.DB.Subscribe(iroh_streamplace.NewFilter()) for { if ctx.Err() != nil { @@ -110,7 +126,7 @@ func (swarm *SwarmKV) Start(ctx context.Context, tickets []string) error { continue } if !activeSubs[keyStr] { - if info.NodeID == nodeIdStr { + if info.NodeID == swarm.nodeId { activeSubs[keyStr] = true continue } @@ -137,8 +153,52 @@ func (swarm *SwarmKV) Start(ctx context.Context, tickets []string) error { } } -func (swarm *SwarmKV) Put(ctx context.Context, key string, value []byte) error { - // streamerBs := []byte(streamer) - keyBs := []byte(key) - return swarm.w.Put(nil, keyBs, value) +func (swarm *IrohSwarm) startSegmentSender(ctx context.Context) error { + ch := swarm.mm.NewSegment() + for { + select { + case <-ctx.Done(): + return ctx.Err() + case not := <-ch: + err := swarm.SendSegment(ctx, not) + if err != nil { + log.Error(ctx, "could not send segment to swarm", "error", err) + } + continue + } + } +} + +func (swarm *IrohSwarm) HandleData(topic string, data []byte) { + err := swarm.mm.ValidateMP4(context.Background(), bytes.NewReader(data), false) + if err != nil { + log.Error(context.Background(), "could not validate segment", "error", err) + } +} + +func (swarm *IrohSwarm) SendSegment(ctx context.Context, not *media.NewSegmentNotification) error { + if !not.Local { + return nil + } + originInfo := OriginInfo{ + NodeID: swarm.nodeId, + Time: not.Segment.StartTime.Format(util.ISO8601), + } + bs, err := json.Marshal(originInfo) + if err != nil { + log.Error(ctx, "could not marshal origin info", "error", err) + return err + } + keyBs := []byte(not.Segment.RepoDID) + err = swarm.w.Put(nil, keyBs, bs) + if err != nil { + log.Error(ctx, "could not put segment to swarm", "error", err) + return err + } + err = swarm.Node.SendSegment(not.Segment.RepoDID, not.Data) + if err != nil { + log.Error(ctx, "could not send segment to swarm", "error", err) + return err + } + return nil } diff --git a/rust/iroh-streamplace/src/api.rs b/rust/iroh-streamplace/src/api.rs deleted file mode 100644 index 677e71ec..00000000 --- a/rust/iroh-streamplace/src/api.rs +++ /dev/null @@ -1,226 +0,0 @@ -//! Protocol API - -use std::collections::{BTreeMap, BTreeSet}; - -use bytes::Bytes; -use iroh::{Endpoint, NodeId, protocol::ProtocolHandler}; -use irpc::{Client, WithChannels, channel::oneshot, rpc::RemoteService, rpc_requests}; -use irpc_iroh::{IrohProtocol, IrohRemoteConnection}; -use n0_future::future::Boxed; -use serde::{Deserialize, Serialize}; -use tracing::{debug, warn}; - -/// Subscribe to the given `key` -#[derive(Debug, Serialize, Deserialize)] -struct Subscribe { - key: String, - // TODO: verify - remote_id: NodeId, -} - -/// Unsubscribe from the given `key` -#[derive(Debug, Serialize, Deserialize)] -struct Unsubscribe { - key: String, - // TODO: verify - remote_id: NodeId, -} - -#[derive(Debug, Serialize, Deserialize)] -struct SendSegment { - key: String, - data: Bytes, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -struct RecvSegment { - key: String, - data: Bytes, -} - -// Use the macro to generate both the Protocol and Message enums -// plus implement Channels for each type -#[rpc_requests(message = Message)] -#[derive(Serialize, Deserialize, Debug)] -enum Protocol { - #[rpc(tx=oneshot::Sender<()>)] - Subscribe(Subscribe), - #[rpc(tx=oneshot::Sender<()>)] - Unsubscribe(Unsubscribe), - #[rpc(tx=oneshot::Sender<()>)] - SendSegment(SendSegment), - #[rpc(tx=oneshot::Sender<()>)] - RecvSegment(RecvSegment), -} - -struct Actor { - endpoint: iroh::Endpoint, - recv: tokio::sync::mpsc::Receiver, - subscriptions: BTreeMap>, - connections: BTreeMap, - handler: Box) -> Boxed<()> + Send + Sync + 'static>, -} - -#[derive(Debug)] -struct Connection { - _id: NodeId, - rpc: Client, -} - -impl Actor { - fn spawn( - endpoint: &iroh::Endpoint, - handler: impl Fn(String, Vec) -> Boxed<()> + Send + Sync + 'static, - ) -> Api { - let (tx, rx) = tokio::sync::mpsc::channel(1); - let actor = Self { - endpoint: endpoint.clone(), - recv: rx, - subscriptions: BTreeMap::new(), - connections: BTreeMap::new(), - handler: Box::new(handler), - }; - n0_future::task::spawn(actor.run()); - Api { - inner: Client::local(tx), - } - } - - async fn run(mut self) { - while let Some(msg) = self.recv.recv().await { - self.handle(msg).await; - } - } - - async fn handle(&mut self, msg: Message) { - match msg { - Message::Subscribe(sub) => { - debug!("subscribe {:?}", sub); - let WithChannels { tx, inner, .. } = sub; - - self.subscriptions - .entry(inner.key) - .or_default() - .insert(inner.remote_id); - - tx.send(()).await.ok(); - } - Message::Unsubscribe(sub) => { - debug!("unsubscribe {:?}", sub); - let WithChannels { tx, inner, .. } = sub; - - if let Some(e) = self.subscriptions.get_mut(&inner.key) { - e.remove(&inner.remote_id); - } - - tx.send(()).await.ok(); - } - Message::SendSegment(segment) => { - debug!("send segment {:?}", segment); - let WithChannels { tx, inner, .. } = segment; - - let msg = RecvSegment { - key: inner.key.clone(), - data: inner.data.clone(), - }; - - for (key, remotes) in &self.subscriptions { - if key == &inner.key { - for remote in remotes { - debug!("sending to topic {}: {}", key, remote); - - // ensure connection - if !self.connections.contains_key(remote) { - let conn = IrohRemoteConnection::new( - self.endpoint.clone(), - (*remote).into(), - Api::ALPN.to_vec(), - ); - - let conn = Connection { - rpc: Client::boxed(conn), - _id: *remote, - }; - self.connections.insert(*remote, conn); - } - let conn = self.connections.get(remote).expect("just checked"); - - if let Err(err) = conn.rpc.rpc(msg.clone()).await { - warn!("failed to send to {}: {:?}", remote, err); - // remove conn - self.connections.remove(remote); - } - } - } - } - - tx.send(()).await.ok(); - } - Message::RecvSegment(segment) => { - debug!("recv segment {:?}", segment); - let WithChannels { tx, inner, .. } = segment; - (self.handler)(inner.key, inner.data.to_vec()).await; - tx.send(()).await.ok(); - } - } - } -} - -/// The actual API to interact with -pub(crate) struct Api { - inner: Client, -} - -impl Api { - pub(crate) const ALPN: &[u8] = b"/iroh/streamplace/1"; - - pub(crate) fn spawn(endpoint: &iroh::Endpoint) -> Self { - Actor::spawn(endpoint, |_, _| Box::pin(async move {})) - } - - pub(crate) fn spawn_with_handler( - endpoint: &iroh::Endpoint, - handler: impl Fn(String, Vec) -> Boxed<()> + Send + Sync + 'static, - ) -> Self { - Actor::spawn(endpoint, handler) - } - - pub(crate) fn connect(endpoint: Endpoint, addr: impl Into) -> Api { - let conn = IrohRemoteConnection::new(endpoint, addr.into(), Self::ALPN.to_vec()); - Api { - inner: Client::boxed(conn), - } - } - - pub(crate) fn expose(&self) -> impl ProtocolHandler { - let local = self - .inner - .as_local() - .expect("can not listen on remote service"); - IrohProtocol::new(Protocol::remote_handler(local)) - } - - pub(crate) async fn subscribe(&self, key: String, self_id: NodeId) -> irpc::Result<()> { - self.inner - .rpc(Subscribe { - key, - remote_id: self_id, - }) - .await - } - - pub(crate) async fn unsubscribe(&self, key: String, self_id: NodeId) -> irpc::Result<()> { - self.inner - .rpc(Unsubscribe { - key, - remote_id: self_id, - }) - .await - } - - /// Send this segment to all subscriptions. - pub(crate) async fn send_segment(&self, key: String, data: Bytes) -> irpc::Result<()> { - let msg = SendSegment { key, data }; - self.inner.rpc(msg).await - } -} diff --git a/rust/iroh-streamplace/src/endpoint.rs b/rust/iroh-streamplace/src/endpoint.rs deleted file mode 100644 index b3662d2e..00000000 --- a/rust/iroh-streamplace/src/endpoint.rs +++ /dev/null @@ -1,30 +0,0 @@ -use iroh::Watcher; - -use crate::{error::Error, node_addr::NodeAddr}; - -#[derive(uniffi::Object, Debug, Clone)] -pub struct Endpoint { - pub(crate) endpoint: iroh::Endpoint, -} - -#[uniffi::export] -impl Endpoint { - /// Create a new endpoint. - #[uniffi::constructor(async_runtime = "tokio")] - pub async fn new() -> Result { - let endpoint = iroh::Endpoint::builder() - .discovery_n0() - .discovery_local_network() - .bind() - .await?; - - Ok(Self { endpoint }) - } - - #[uniffi::method(async_runtime = "tokio")] - pub async fn node_addr(&self) -> NodeAddr { - let _ = self.endpoint.home_relay().initialized().await; - let addr = self.endpoint.node_addr().initialized().await; - addr.into() - } -} diff --git a/rust/iroh-streamplace/src/lib.rs b/rust/iroh-streamplace/src/lib.rs index 52558b38..6297bfb5 100644 --- a/rust/iroh-streamplace/src/lib.rs +++ b/rust/iroh-streamplace/src/lib.rs @@ -1,12 +1,7 @@ uniffi::setup_scaffolding!(); pub mod c2pa; -pub mod endpoint; pub mod error; -pub mod public_key; pub mod node; -pub mod receiver; -pub mod sender; pub mod node_addr; - -mod api; +pub mod public_key; diff --git a/rust/iroh-streamplace/src/receiver.rs b/rust/iroh-streamplace/src/receiver.rs deleted file mode 100644 index 23d2df7a..00000000 --- a/rust/iroh-streamplace/src/receiver.rs +++ /dev/null @@ -1,150 +0,0 @@ -use std::sync::Arc; - -use iroh::protocol::Router; - -use crate::{api::Api, endpoint::Endpoint, error::Error, public_key::PublicKey, node_addr::NodeAddr}; - -#[derive(uniffi::Object)] -pub struct Receiver { - endpoint: Endpoint, - _api: Api, - _router: iroh::protocol::Router, -} - -#[uniffi::export] -impl Receiver { - /// Create a new receiver. - #[uniffi::constructor(async_runtime = "tokio")] - pub async fn new( - endpoint: &Endpoint, - handler: Arc, - ) -> Result { - let api = Api::spawn_with_handler(&endpoint.endpoint, move |id, data| { - let handler = handler.clone(); - Box::pin(async move { - handler.handle_data(id, data).await; - }) - }); - let router = Router::builder(endpoint.endpoint.clone()) - .accept(Api::ALPN, api.expose()) - .spawn(); - - Ok(Receiver { - endpoint: endpoint.clone(), - _api: api, - _router: router, - }) - } - - /// Subscribe to the given topic on the remote. - #[uniffi::method(async_runtime = "tokio")] - pub async fn subscribe(&self, remote_id: Arc, topic: &str) -> Result<(), Error> { - let remote_id: iroh::NodeId = remote_id.as_ref().into(); - let api = Api::connect(self.endpoint.endpoint.clone(), remote_id); - api.subscribe(topic.to_string(), self.endpoint.endpoint.node_id()) - .await?; - Ok(()) - } - - /// Unsubscribe from this topic on the remote. - #[uniffi::method(async_runtime = "tokio")] - pub async fn unsubscribe( - &self, - remote_id: Arc, - topic: &str, - ) -> Result<(), Error> { - let remote_id: iroh::NodeId = remote_id.as_ref().into(); - let api = Api::connect(self.endpoint.endpoint.clone(), remote_id); - api.unsubscribe(topic.to_string(), self.endpoint.endpoint.node_id()) - .await?; - Ok(()) - } - - #[uniffi::method(async_runtime = "tokio")] - pub async fn node_addr(&self) -> NodeAddr { - self.endpoint.node_addr().await - } -} - -#[uniffi::export(with_foreign)] -#[async_trait::async_trait] -pub trait DataHandlerOld: Send + Sync { - async fn handle_data(&self, topic: String, data: Vec); -} - -#[cfg(test)] -mod tests { - - use super::*; - use crate::sender::Sender; - - #[tokio::test] - async fn test_roundtrip() { - tracing_subscriber::fmt() - .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) - .init(); - - let ep1 = Endpoint::new().await.unwrap(); - let sender = Sender::new(&ep1).await.unwrap(); - - let (s, mut r) = tokio::sync::mpsc::channel(5); - - #[derive(Debug, Clone)] - struct TestHandler { - messages: tokio::sync::mpsc::Sender<(String, Vec)>, - } - - #[async_trait::async_trait] - impl DataHandlerOld for TestHandler { - async fn handle_data(&self, topic: String, data: Vec) { - self.messages.send((topic, data)).await.unwrap(); - } - } - - let handler = TestHandler { messages: s }; - let ep2 = Endpoint::new().await.unwrap(); - let receiver = Receiver::new(&ep2, Arc::new(handler.clone())) - .await - .unwrap(); - - let sender_addr = sender.node_addr().await; - println!("sender addr: {sender_addr:?}"); - - let receiver_addr = receiver.node_addr().await; - println!("recv addr: {receiver_addr:?}"); - - // subscribe - receiver - .subscribe(Arc::new(sender_addr.node_id()), "foo") - .await - .unwrap(); - - // send a few messages - for i in 0u8..5 { - sender.send("foo", &[i, 0, 0, 0]).await.unwrap(); - } - - // make sure the receiver got them - for i in 0u8..5 { - let (topic, msg) = r.recv().await.unwrap(); - assert_eq!(topic, "foo"); - assert_eq!(msg, vec![i, 0, 0, 0]); - } - - // unsubscribe - receiver - .unsubscribe(Arc::new(sender_addr.node_id()), "foo") - .await - .unwrap(); - - // send a message, shouldn't error - sender.send("foo", &[1]).await.unwrap(); - - // no message received, times out - let res = tokio::time::timeout(std::time::Duration::from_millis(200), async { - r.recv().await.unwrap(); - }) - .await; - assert!(res.is_err()); - } -} diff --git a/rust/iroh-streamplace/src/sender.rs b/rust/iroh-streamplace/src/sender.rs deleted file mode 100644 index 9b665a7f..00000000 --- a/rust/iroh-streamplace/src/sender.rs +++ /dev/null @@ -1,43 +0,0 @@ -use bytes::Bytes; -use iroh::protocol::Router; - -use crate::{api::Api, c2pa::SPError, endpoint::Endpoint, error::Error, node_addr::NodeAddr}; - -#[derive(uniffi::Object)] -pub struct Sender { - endpoint: Endpoint, - api: Api, - _router: iroh::protocol::Router, -} - -#[uniffi::export] -impl Sender { - /// Create a new sender. - #[uniffi::constructor(async_runtime = "tokio")] - pub async fn new(endpoint: &Endpoint) -> Sender { - let api = Api::spawn(&endpoint.endpoint); - let router = Router::builder(endpoint.endpoint.clone()) - .accept(Api::ALPN, api.expose()) - .spawn(); - - Sender { - endpoint: endpoint.clone(), - api, - _router: router, - } - } - - /// Sends the given data to all subscribers that have subscribed to this `key`. - #[uniffi::method(async_runtime = "tokio")] - pub async fn send(&self, key: &str, data: &[u8]) -> Result<(), Error> { - self.api - .send_segment(key.to_string(), Bytes::copy_from_slice(data)) - .await?; - Ok(()) - } - - #[uniffi::method(async_runtime = "tokio")] - pub async fn node_addr(&self) -> NodeAddr { - self.endpoint.node_addr().await - } -}