diff --git a/cmd/cocoon/main.go b/cmd/cocoon/main.go index 3adb24e..435c666 100644 --- a/cmd/cocoon/main.go +++ b/cmd/cocoon/main.go @@ -169,6 +169,14 @@ func main() { }, telemetry.CLIFlagDebug, telemetry.CLIFlagMetricsListenAddress, + &cli.StringFlag{ + Name: "subscribe-repos-service-url", + EnvVars: []string{"SUBSCRIBE_REPOS_SERVICE_URL"}, + }, + &cli.BoolFlag{ + Name: "push-based-events", + EnvVars: []string{"PUSH_BASED_EVENTS"}, + }, }, Commands: []*cli.Command{ runServe, @@ -249,10 +257,12 @@ var runServe = &cli.Command{ SecretKey: cmd.String("s3-secret-key"), CDNUrl: cmd.String("s3-cdn-url"), }, - SessionSecret: cmd.String("session-secret"), - SessionCookieKey: cmd.String("session-cookie-key"), - BlockstoreVariant: server.MustReturnBlockstoreVariant(cmd.String("blockstore-variant")), - FallbackProxy: cmd.String("fallback-proxy"), + SessionSecret: cmd.String("session-secret"), + SessionCookieKey: cmd.String("session-cookie-key"), + BlockstoreVariant: server.MustReturnBlockstoreVariant(cmd.String("blockstore-variant")), + FallbackProxy: cmd.String("fallback-proxy"), + PushBasedEvents: cmd.Bool("push-based-events"), + SubscribeReposServiceURL: cmd.String("subscribe-repos-service-url"), }) if err != nil { fmt.Printf("error creating cocoon: %v", err) diff --git a/server/event_emmiter.go b/server/event_emmiter.go new file mode 100644 index 0000000..fe3bfba --- /dev/null +++ b/server/event_emmiter.go @@ -0,0 +1,92 @@ +package server + +import ( + "bytes" + "context" + "net/http" + "time" + + "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/lex/util" +) + +func (s *Server) emmitEvents(ctx context.Context) error { + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + logger := s.logger.With("component", "event-emmiter") + ident := "self" + var since *int64 + // TODO: track since + + evts, evtManCancel, err := s.evtman.Subscribe(ctx, ident, func(evt *events.XRPCStreamEvent) bool { + return true + }, since) + if err != nil { + return err + } + defer evtManCancel() + + header := events.EventHeader{Op: events.EvtKindMessage} + for evt := range evts { + func() { + if ctx.Err() != nil { + logger.Error("context error", "err", err) + return + } + + var obj util.CBOR + switch { + case evt.Error != nil: + header.Op = events.EvtKindErrorFrame + obj = evt.Error + case evt.RepoCommit != nil: + header.MsgType = "#commit" + obj = evt.RepoCommit + case evt.RepoIdentity != nil: + header.MsgType = "#identity" + obj = evt.RepoIdentity + case evt.RepoAccount != nil: + header.MsgType = "#account" + obj = evt.RepoAccount + case evt.RepoInfo != nil: + header.MsgType = "#info" + obj = evt.RepoInfo + default: + logger.Warn("unrecognized event kind") + return + } + + buf := new(bytes.Buffer) + + if err := header.MarshalCBOR(buf); err != nil { + logger.Error("failed to marshal header to buffer", "err", err) + return + } + + if err := obj.MarshalCBOR(buf); err != nil { + logger.Error("failed to marshal event to buffer", "err", err) + return + } + + // TODO: use a HTTP client here not the default + _, err := http.Post(s.config.SubscribeReposServiceURL, "", buf) + if err != nil { + logger.Error("posting to web server", "error", err) + return + } + }() + } + + // we should tell the relay to request a new crawl at this point if we got disconnected + // use a new context since the old one might be cancelled at this point + go func() { + retryCtx, retryCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer retryCancel() + if err := s.requestCrawl(retryCtx); err != nil { + logger.Error("error requesting crawls", "err", err) + } + }() + + return nil +} diff --git a/server/server.go b/server/server.go index 53fd011..3b376e4 100644 --- a/server/server.go +++ b/server/server.go @@ -124,6 +124,9 @@ type Args struct { BlockstoreVariant BlockstoreVariant FallbackProxy string + + PushBasedEvents bool + SubscribeReposServiceURL string } type config struct { @@ -141,6 +144,9 @@ type config struct { SessionCookieKey string BlockstoreVariant BlockstoreVariant FallbackProxy string + + PushBasedEvents bool + SubscribeReposServiceURL string } type CustomValidator struct { @@ -435,20 +441,22 @@ func New(args *Args) (*Server, error) { plcClient: plcClient, privateKey: &pkey, config: &config{ - LogLevel: args.LogLevel, - Version: args.Version, - Did: args.Did, - Hostname: args.Hostname, - ContactEmail: args.ContactEmail, - EnforcePeering: false, - Relays: args.Relays, - AdminPassword: args.AdminPassword, - RequireInvite: args.RequireInvite, - SmtpName: args.SmtpName, - SmtpEmail: args.SmtpEmail, - SessionCookieKey: args.SessionCookieKey, - BlockstoreVariant: args.BlockstoreVariant, - FallbackProxy: args.FallbackProxy, + LogLevel: args.LogLevel, + Version: args.Version, + Did: args.Did, + Hostname: args.Hostname, + ContactEmail: args.ContactEmail, + EnforcePeering: false, + Relays: args.Relays, + AdminPassword: args.AdminPassword, + RequireInvite: args.RequireInvite, + SmtpName: args.SmtpName, + SmtpEmail: args.SmtpEmail, + SessionCookieKey: args.SessionCookieKey, + BlockstoreVariant: args.BlockstoreVariant, + FallbackProxy: args.FallbackProxy, + PushBasedEvents: args.PushBasedEvents, + SubscribeReposServiceURL: args.SubscribeReposServiceURL, }, evtman: events.NewEventManager(evtPersister), passport: identity.NewPassport(h, identity.NewMemCache(10_000)), @@ -635,6 +643,15 @@ func (s *Server) Serve(ctx context.Context) error { } }() + if s.config.PushBasedEvents { + slog.Info("pushed based events enabled") + go func() { + if err := s.emmitEvents(ctx); err != nil { + logger.Error("error emitting events", "err", err) + } + }() + } + <-ctx.Done() fmt.Println("shut down")