| 
									
										
										
										
											2021-06-17 19:20:08 +02:00
										 |  |  | package streaming | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | import ( | 
					
						
							|  |  |  | 	"sync" | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	"github.com/sirupsen/logrus" | 
					
						
							| 
									
										
										
										
											2021-06-18 13:06:02 +02:00
										 |  |  | 	apimodel "github.com/superseriousbusiness/gotosocial/internal/api/model" | 
					
						
							| 
									
										
										
										
											2021-06-17 19:20:08 +02:00
										 |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/config" | 
					
						
							|  |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/db" | 
					
						
							|  |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/gtserror" | 
					
						
							|  |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/gtsmodel" | 
					
						
							|  |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/oauth" | 
					
						
							|  |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/typeutils" | 
					
						
							|  |  |  | 	"github.com/superseriousbusiness/gotosocial/internal/visibility" | 
					
						
							|  |  |  | ) | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | // Processor wraps a bunch of functions for processing streaming. | 
					
						
							|  |  |  | type Processor interface { | 
					
						
							|  |  |  | 	// AuthorizeStreamingRequest returns an oauth2 token info in response to an access token query from the streaming API | 
					
						
							|  |  |  | 	AuthorizeStreamingRequest(accessToken string) (*gtsmodel.Account, error) | 
					
						
							| 
									
										
										
										
											2021-06-18 13:06:02 +02:00
										 |  |  | 	OpenStreamForAccount(account *gtsmodel.Account, streamType string) (*gtsmodel.Stream, gtserror.WithCode) | 
					
						
							|  |  |  | 	StreamStatusToAccount(s *apimodel.Status, account *gtsmodel.Account) error | 
					
						
							| 
									
										
										
										
											2021-06-19 10:47:05 +02:00
										 |  |  | 	StreamNotificationToAccount(n *apimodel.Notification, account *gtsmodel.Account) error | 
					
						
							| 
									
										
										
										
											2021-06-17 19:20:08 +02:00
										 |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | type processor struct { | 
					
						
							|  |  |  | 	tc          typeutils.TypeConverter | 
					
						
							|  |  |  | 	config      *config.Config | 
					
						
							|  |  |  | 	db          db.DB | 
					
						
							|  |  |  | 	filter      visibility.Filter | 
					
						
							|  |  |  | 	log         *logrus.Logger | 
					
						
							|  |  |  | 	oauthServer oauth.Server | 
					
						
							|  |  |  | 	streamMap   *sync.Map | 
					
						
							|  |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | // New returns a new status processor. | 
					
						
							|  |  |  | func New(db db.DB, tc typeutils.TypeConverter, oauthServer oauth.Server, config *config.Config, log *logrus.Logger) Processor { | 
					
						
							|  |  |  | 	return &processor{ | 
					
						
							|  |  |  | 		tc:          tc, | 
					
						
							|  |  |  | 		config:      config, | 
					
						
							|  |  |  | 		db:          db, | 
					
						
							|  |  |  | 		filter:      visibility.NewFilter(db, log), | 
					
						
							|  |  |  | 		log:         log, | 
					
						
							|  |  |  | 		oauthServer: oauthServer, | 
					
						
							|  |  |  | 		streamMap:   &sync.Map{}, | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | } |