diff --git a/lexicons/place/stream/broadcast/defs.json b/lexicons/place/stream/broadcast/defs.json new file mode 100644 index 000000000..5c7ebaba4 --- /dev/null +++ b/lexicons/place/stream/broadcast/defs.json @@ -0,0 +1,19 @@ +{ + "lexicon": 1, + "id": "place.stream.broadcast.defs", + "defs": { + "broadcastOriginView": { + "type": "object", + "required": ["uri", "cid", "author", "record", "indexedAt"], + "properties": { + "uri": { "type": "string", "format": "at-uri" }, + "cid": { "type": "string", "format": "cid" }, + "author": { + "type": "ref", + "ref": "app.bsky.actor.defs#profileViewBasic" + }, + "record": { "type": "unknown" } + } + } + } +} diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 0fa5924b3..c5f35753e 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -390,7 +390,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD } case *streamplace.BroadcastOrigin: - _, err := atsync.SyncBlueskyRepoCached(ctx, userDID, atsync.Model) + repo, err := atsync.SyncBlueskyRepoCached(ctx, userDID, atsync.Model) if err != nil { return fmt.Errorf("failed to sync broadcast origin creator bluesky repo: %w", err) } @@ -402,6 +402,17 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD if err != nil { log.Error(ctx, "failed to update broadcast origin", "err", err) } + view := &streamplace.BroadcastDefs_BroadcastOriginView{ + Uri: aturi.String(), + Cid: cid, + Author: &bsky.ActorDefs_ProfileViewBasic{ + Did: userDID, + Handle: repo.Handle, + }, + Record: &lexutil.LexiconTypeDecoder{Val: rec}, + } + // publishes with an empty string because we're discovering the stream + go atsync.Bus.Publish("", view) default: log.Debug(ctx, "unhandled record type", "type", reflect.TypeOf(rec)) diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 2aa6839b5..9dfffe99c 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -428,7 +428,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return err } } - swarm, err := iroh_replicator.NewSwarm(ctx, cli.Tickets, secret, topic, mm) + swarm, err := iroh_replicator.NewSwarm(ctx, cli.Tickets, secret, topic, mm, b) if err != nil { return err } diff --git a/pkg/config/config.go b/pkg/config/config.go index de46549c2..d4c3f1dc3 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -125,6 +125,7 @@ type CLI struct { LivepeerDebug bool Tickets []string IrohTopic string + DID string } func (cli *CLI) NewFlagSet(name string) *flag.FlagSet { diff --git a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go index 5caeff686..48bbd66d6 100644 --- a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go +++ b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go @@ -361,6 +361,15 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_get_manifest_and_cert: UniFFI API checksum mismatch") } } + { + checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { + return C.uniffi_iroh_streamplace_checksum_func_node_id_from_ticket() + }) + if checksum != 36085 { + // If this happens try cleaning and rebuilding your project + panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_node_id_from_ticket: UniFFI API checksum mismatch") + } + } { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_func_sign() @@ -3183,6 +3192,111 @@ func (_ FfiDestroyerNodeAddrError) Destroy(value *NodeAddrError) { } } +// Error joining peers. +type ParseError struct { + err error +} + +// Convience method to turn *ParseError into error +// Avoiding treating nil pointer as non nil error interface +func (err *ParseError) AsError() error { + if err == nil { + return nil + } else { + return err + } +} + +func (err ParseError) Error() string { + return fmt.Sprintf("ParseError: %s", err.err.Error()) +} + +func (err ParseError) Unwrap() error { + return err.err +} + +// Err* are used for checking error type with `errors.Is` +var ErrParseErrorTicket = fmt.Errorf("ParseErrorTicket") + +// Variant structs +// Failed to parse a provided iroh node ticket. +type ParseErrorTicket struct { + Message string +} + +// Failed to parse a provided iroh node ticket. +func NewParseErrorTicket( + message string, +) *ParseError { + return &ParseError{err: &ParseErrorTicket{ + Message: message}} +} + +func (e ParseErrorTicket) destroy() { + FfiDestroyerString{}.Destroy(e.Message) +} + +func (err ParseErrorTicket) Error() string { + return fmt.Sprint("Ticket", + ": ", + + "Message=", + err.Message, + ) +} + +func (self ParseErrorTicket) Is(target error) bool { + return target == ErrParseErrorTicket +} + +type FfiConverterParseError struct{} + +var FfiConverterParseErrorINSTANCE = FfiConverterParseError{} + +func (c FfiConverterParseError) Lift(eb RustBufferI) *ParseError { + return LiftFromRustBuffer[*ParseError](c, eb) +} + +func (c FfiConverterParseError) Lower(value *ParseError) C.RustBuffer { + return LowerIntoRustBuffer[*ParseError](c, value) +} + +func (c FfiConverterParseError) Read(reader io.Reader) *ParseError { + errorID := readUint32(reader) + + switch errorID { + case 1: + return &ParseError{&ParseErrorTicket{ + Message: FfiConverterStringINSTANCE.Read(reader), + }} + default: + panic(fmt.Sprintf("Unknown error code %d in FfiConverterParseError.Read()", errorID)) + } +} + +func (c FfiConverterParseError) Write(writer io.Writer, value *ParseError) { + switch variantValue := value.err.(type) { + case *ParseErrorTicket: + writeInt32(writer, 1) + FfiConverterStringINSTANCE.Write(writer, variantValue.Message) + default: + _ = variantValue + panic(fmt.Sprintf("invalid error value `%v` in FfiConverterParseError.Write", value)) + } +} + +type FfiDestroyerParseError struct{} + +func (_ FfiDestroyerParseError) Destroy(value *ParseError) { + switch variantValue := value.err.(type) { + case ParseErrorTicket: + variantValue.destroy() + default: + _ = variantValue + panic(fmt.Sprintf("invalid error value `%v` in FfiDestroyerParseError.Destroy", value)) + } +} + type PublicKeyError struct { err error } @@ -3978,6 +4092,99 @@ func (_ FfiDestroyerSubscribeNextError) Destroy(value *SubscribeNextError) { } } +// Error when converting from ffi NodeAddr to iroh::NodeAddr +type TicketError struct { + err error +} + +// Convience method to turn *TicketError into error +// Avoiding treating nil pointer as non nil error interface +func (err *TicketError) AsError() error { + if err == nil { + return nil + } else { + return err + } +} + +func (err TicketError) Error() string { + return fmt.Sprintf("TicketError: %s", err.err.Error()) +} + +func (err TicketError) Unwrap() error { + return err.err +} + +// Err* are used for checking error type with `errors.Is` +var ErrTicketErrorParseError = fmt.Errorf("TicketErrorParseError") + +// Variant structs +type TicketErrorParseError struct { + message string +} + +func NewTicketErrorParseError() *TicketError { + return &TicketError{err: &TicketErrorParseError{}} +} + +func (e TicketErrorParseError) destroy() { +} + +func (err TicketErrorParseError) Error() string { + return fmt.Sprintf("ParseError: %s", err.message) +} + +func (self TicketErrorParseError) Is(target error) bool { + return target == ErrTicketErrorParseError +} + +type FfiConverterTicketError struct{} + +var FfiConverterTicketErrorINSTANCE = FfiConverterTicketError{} + +func (c FfiConverterTicketError) Lift(eb RustBufferI) *TicketError { + return LiftFromRustBuffer[*TicketError](c, eb) +} + +func (c FfiConverterTicketError) Lower(value *TicketError) C.RustBuffer { + return LowerIntoRustBuffer[*TicketError](c, value) +} + +func (c FfiConverterTicketError) Read(reader io.Reader) *TicketError { + errorID := readUint32(reader) + + message := FfiConverterStringINSTANCE.Read(reader) + switch errorID { + case 1: + return &TicketError{&TicketErrorParseError{message}} + default: + panic(fmt.Sprintf("Unknown error code %d in FfiConverterTicketError.Read()", errorID)) + } + +} + +func (c FfiConverterTicketError) Write(writer io.Writer, value *TicketError) { + switch variantValue := value.err.(type) { + case *TicketErrorParseError: + writeInt32(writer, 1) + default: + _ = variantValue + panic(fmt.Sprintf("invalid error value `%v` in FfiConverterTicketError.Write", value)) + } +} + +type FfiDestroyerTicketError struct{} + +func (_ FfiDestroyerTicketError) Destroy(value *TicketError) { + switch variantValue := value.err.(type) { + case TicketErrorParseError: + variantValue.destroy() + default: + _ = variantValue + panic(fmt.Sprintf("invalid error value `%v` in FfiDestroyerTicketError.Destroy", value)) + } +} + // A bound on time for filtering. type TimeBound interface { Destroy() @@ -4477,6 +4684,19 @@ func GetManifestAndCert(data []byte) (string, error) { } } +// Get this node's ticket. +func NodeIdFromTicket(ticketStr string) (*PublicKey, error) { + _uniffiRV, _uniffiErr := rustCallWithError[TicketError](FfiConverterTicketError{}, func(_uniffiStatus *C.RustCallStatus) unsafe.Pointer { + return C.uniffi_iroh_streamplace_fn_func_node_id_from_ticket(FfiConverterStringINSTANCE.Lower(ticketStr), _uniffiStatus) + }) + if _uniffiErr != nil { + var _uniffiDefaultValue *PublicKey + return _uniffiDefaultValue, _uniffiErr + } else { + return FfiConverterPublicKeyINSTANCE.Lift(_uniffiRV), nil + } +} + func Sign(manifest string, data []byte, certs []byte, gosigner GoSigner) ([]byte, error) { _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { return GoRustBuffer{ diff --git a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h index def03a225..8dc7969cf 100644 --- a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h +++ b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h @@ -738,6 +738,11 @@ uint64_t uniffi_iroh_streamplace_fn_method_writescope_put(void* ptr, RustBuffer RustBuffer uniffi_iroh_streamplace_fn_func_get_manifest_and_cert(RustBuffer data, RustCallStatus *out_status ); #endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_NODE_ID_FROM_TICKET +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_NODE_ID_FROM_TICKET +void* uniffi_iroh_streamplace_fn_func_node_id_from_ticket(RustBuffer ticket_str, RustCallStatus *out_status +); +#endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SIGN #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SIGN RustBuffer uniffi_iroh_streamplace_fn_func_sign(RustBuffer manifest, RustBuffer data, RustBuffer certs, void* gosigner, RustCallStatus *out_status @@ -1032,6 +1037,12 @@ void ffi_iroh_streamplace_rust_future_complete_void(uint64_t handle, RustCallSta #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_FUNC_GET_MANIFEST_AND_CERT uint16_t uniffi_iroh_streamplace_checksum_func_get_manifest_and_cert(void +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_FUNC_NODE_ID_FROM_TICKET +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_FUNC_NODE_ID_FROM_TICKET +uint16_t uniffi_iroh_streamplace_checksum_func_node_id_from_ticket(void + ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_FUNC_SIGN diff --git a/pkg/replication/iroh_replicator/kv.go b/pkg/replication/iroh_replicator/kv.go index d8fda1051..0d3c68549 100644 --- a/pkg/replication/iroh_replicator/kv.go +++ b/pkg/replication/iroh_replicator/kv.go @@ -6,13 +6,16 @@ import ( "crypto/rand" "encoding/json" "fmt" + "sync" "time" "github.com/bluesky-social/indigo/util" "golang.org/x/sync/errgroup" + "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" + "stream.place/streamplace/pkg/streamplace" ) type IrohSwarm struct { @@ -25,6 +28,8 @@ type IrohSwarm struct { NodeTicket string activeSubs map[string]*OriginInfo handleDataScoped func(topic string, data []byte) + bus *bus.Bus + originMutex sync.Mutex } // A message saying "hey I ingested node data at this time" @@ -33,7 +38,7 @@ type OriginInfo struct { Time string `json:"time"` } -func NewSwarm(ctx context.Context, tickets []string, secret []byte, topic []byte, mm *media.MediaManager) (*IrohSwarm, error) { +func NewSwarm(ctx context.Context, tickets []string, secret []byte, topic []byte, mm *media.MediaManager, bus *bus.Bus) (*IrohSwarm, error) { ctx = log.WithLogValues(ctx, "func", "StartKV") if topic == nil { @@ -55,6 +60,7 @@ func NewSwarm(ctx context.Context, tickets []string, secret []byte, topic []byte swarm := IrohSwarm{ mm: mm, activeSubs: make(map[string]*OriginInfo), + bus: bus, } // workaround to get context into the HandleData callback @@ -122,6 +128,9 @@ func (swarm *IrohSwarm) Start(ctx context.Context, tickets []string) error { <-ctx.Done() return swarm.Node.Shutdown() }) + g.Go(func() error { + return swarm.startBusSubscribe(ctx) + }) return g.Wait() } @@ -135,6 +144,7 @@ func (swarm *IrohSwarm) startKV(ctx context.Context) error { if err != nil { return fmt.Errorf("failed to get next subscription event: %w", err) } + if ev == nil { log.Debug(ctx, "Got empty event from sub.NextRaw(), pausing for a second") time.Sleep(1 * time.Second) @@ -156,53 +166,109 @@ func (swarm *IrohSwarm) startKV(ctx context.Context) error { log.Error(ctx, "could not unmarshal origin info", "error", err) continue } - oldSub, ok := swarm.activeSubs[keyStr] - if ok { - if oldSub.NodeID == info.NodeID { - log.Debug(ctx, "node hasn't changed", "streamer", keyStr) - // mmyep. same node still has the stream. great news. + err = swarm.checkOrigins(ctx, keyStr, info.NodeID) + if err != nil { + log.Error(ctx, "could not check origins", "error", err) + continue + } + case iroh_streamplace.SubscribeItemCurrentDone: + log.Debug(ctx, "SubscribeItemCurrentDone", "currentDone", item) + case iroh_streamplace.SubscribeItemExpired: + log.Debug(ctx, "SubscribeItemExpired", "expired", item) + case iroh_streamplace.SubscribeItemOther: + log.Debug(ctx, "SubscribeItemOther", "other", item) + } + } +} + +// subscribe to all streams +func (swarm *IrohSwarm) startBusSubscribe(ctx context.Context) error { + for { + select { + case <-ctx.Done(): + return ctx.Err() + case msg := <-swarm.bus.Subscribe(""): + if view, ok := msg.(*streamplace.BroadcastDefs_BroadcastOriginView); ok { + log.Debug(ctx, "got broadcast origin view", "view", view) + origin, ok := view.Record.Val.(*streamplace.BroadcastOrigin) + if !ok { + log.Error(ctx, "record is not a BroadcastOrigin", "record", view.Record) + continue + } + if view.Author.Did != origin.Streamer { + // currently, only streamers are allowed to advertise origins continue } - log.Log(ctx, "Stream origin changed, swapping to new node", "old_node", oldSub.NodeID, "new_node", info.NodeID, "streamer", keyStr) - pubKey, err := iroh_streamplace.PublicKeyFromString(oldSub.NodeID) + if origin.IrohTicket == nil { + log.Error(ctx, "origin has no iroh ticket", "origin", origin) + continue + } + pubKey, err := iroh_streamplace.NodeIdFromTicket(*origin.IrohTicket) if err != nil { - log.Error(ctx, "could not create public key", "error", err) + log.Error(ctx, "could not get node id from ticket", "error", err) continue } - // different node has the stream. we need to unsubscribe from the old node. - err = swarm.Node.Unsubscribe(keyStr, pubKey) + err = swarm.Node.AddTickets([]string{*origin.IrohTicket}) if err != nil { - log.Error(ctx, "could not unsubscribe from key", "error", err) + log.Error(ctx, "could not add tickets", "error", err) + continue + } + pubKeyStr := pubKey.String() + err = swarm.checkOrigins(ctx, origin.Streamer, pubKeyStr) + if err != nil { + log.Error(ctx, "could not check origin", "error", err) continue } - delete(swarm.activeSubs, keyStr) - } - if info.NodeID == swarm.NodeID { - log.Debug(ctx, "I already have this stream", "streamer", keyStr) - // oh, i have this stream. cool. do nothing. - continue - } - log.Log(ctx, "Subscribing to stream", "new_node", info.NodeID, "streamer", keyStr) - pubKey, err := iroh_streamplace.PublicKeyFromString(info.NodeID) - if err != nil { - log.Error(ctx, "could not create public key", "error", err) - continue - } - err = swarm.Node.Subscribe(keyStr, pubKey) - if err != nil { - log.Error(ctx, "could not subscribe to key", "error", err) - continue } - swarm.activeSubs[keyStr] = &info + } + } +} - case iroh_streamplace.SubscribeItemCurrentDone: - log.Debug(ctx, "SubscribeItemCurrentDone", "currentDone", item) - case iroh_streamplace.SubscribeItemExpired: - log.Debug(ctx, "SubscribeItemExpired", "expired", item) - case iroh_streamplace.SubscribeItemOther: - log.Debug(ctx, "SubscribeItemOther", "other", item) +func (swarm *IrohSwarm) checkOrigins(ctx context.Context, streamer string, nodeID string) error { + swarm.originMutex.Lock() + defer swarm.originMutex.Unlock() + oldSub, ok := swarm.activeSubs[streamer] + if ok { + if oldSub.NodeID == nodeID { + log.Debug(ctx, "node hasn't changed", "streamer", streamer) + // mmyep. same node still has the stream. great news. + return nil } + log.Log(ctx, "Stream origin changed, swapping to new node", "old_node", oldSub.NodeID, "new_node", nodeID, "streamer", streamer) + pubKey, err := iroh_streamplace.PublicKeyFromString(oldSub.NodeID) + if err != nil { + log.Error(ctx, "could not create public key", "error", err) + return err + } + // different node has the stream. we need to unsubscribe from the old node. + err = swarm.Node.Unsubscribe(streamer, pubKey) + if err != nil { + log.Error(ctx, "could not unsubscribe from key", "error", err) + return err + } + delete(swarm.activeSubs, streamer) + } + if nodeID == swarm.NodeID { + log.Debug(ctx, "I already have this stream", "streamer", streamer) + // oh, i have this stream. cool. do nothing. + return nil + } + log.Log(ctx, "Subscribing to stream", "new_node", nodeID, "streamer", streamer) + pubKey, err := iroh_streamplace.PublicKeyFromString(nodeID) + if err != nil { + log.Error(ctx, "could not create public key", "error", err) + return err } + err = swarm.Node.Subscribe(streamer, pubKey) + if err != nil { + log.Error(ctx, "could not subscribe to key", "error", err) + return err + } + swarm.activeSubs[streamer] = &OriginInfo{ + NodeID: nodeID, + Time: time.Now().Format(util.ISO8601), + } + return nil } func (swarm *IrohSwarm) startSegmentSender(ctx context.Context) error { diff --git a/pkg/streamplace/broadcastdefs.go b/pkg/streamplace/broadcastdefs.go new file mode 100644 index 000000000..5fc1abf07 --- /dev/null +++ b/pkg/streamplace/broadcastdefs.go @@ -0,0 +1,18 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package streamplace + +// schema: place.stream.broadcast.defs + +import ( + appbskytypes "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/lex/util" +) + +// BroadcastDefs_BroadcastOriginView is a "broadcastOriginView" in the place.stream.broadcast.defs schema. +type BroadcastDefs_BroadcastOriginView struct { + Author *appbskytypes.ActorDefs_ProfileViewBasic `json:"author" cborgen:"author"` + Cid string `json:"cid" cborgen:"cid"` + Record *util.LexiconTypeDecoder `json:"record" cborgen:"record"` + Uri string `json:"uri" cborgen:"uri"` +} diff --git a/rust/iroh-streamplace/src/db.rs b/rust/iroh-streamplace/src/db.rs index 945d4197f..bc47fc4fe 100644 --- a/rust/iroh-streamplace/src/db.rs +++ b/rust/iroh-streamplace/src/db.rs @@ -48,6 +48,14 @@ pub enum JoinPeersError { Irpc { message: String }, } +/// Error joining peers. +#[derive(Debug, Snafu, uniffi::Error)] +#[snafu(module)] +pub enum ParseError { + /// Failed to parse a provided iroh node ticket. + Ticket { message: String }, +} + /// Error putting a value into the database. #[derive(Debug, Snafu, uniffi::Error)] #[snafu(module)] diff --git a/rust/iroh-streamplace/src/node_addr.rs b/rust/iroh-streamplace/src/node_addr.rs index 966c8f141..9f6828ad1 100644 --- a/rust/iroh-streamplace/src/node_addr.rs +++ b/rust/iroh-streamplace/src/node_addr.rs @@ -1,3 +1,4 @@ +use iroh_base::ticket::NodeTicket; use std::{str::FromStr, sync::Arc}; use crate::public_key::PublicKey; @@ -87,3 +88,22 @@ impl From for NodeAddr { } } } + +/// Error when converting from ffi NodeAddr to iroh::NodeAddr +#[derive(Debug, snafu::Snafu, uniffi::Error)] +#[uniffi(flat_error)] +pub enum TicketError { + ParseError { message: String }, +} + +/// Get this node's ticket. +#[uniffi::export] +pub fn node_id_from_ticket( + ticket_str: String, +) -> Result, TicketError> { + let ticket = NodeTicket::from_str(&ticket_str).map_err(|e| TicketError::ParseError { + message: e.to_string(), + })?; + let node_addr = ticket.node_addr(); + Ok(Arc::new(node_addr.node_id.into())) +}