Fix golint issues
This commit is contained in:
+14
-6
@@ -43,7 +43,7 @@ func prepareQueryString(params url.Values) string {
|
|||||||
return strings.Join(pieces, "&")
|
return strings.Join(pieces, "&")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Authenticate pusher
|
// Authentication Authenticate pusher
|
||||||
// see: https://gist.github.com/mloughran/376898
|
// see: https://gist.github.com/mloughran/376898
|
||||||
//
|
//
|
||||||
// The signature is a HMAC SHA256 hex digest.
|
// The signature is a HMAC SHA256 hex digest.
|
||||||
@@ -90,7 +90,7 @@ func Authentication(storage storage.Storage) func(http.Handler) http.Handler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if the application is disabled
|
// CheckAppDisabled Check if the application is disabled
|
||||||
func CheckAppDisabled(storage storage.Storage) func(http.Handler) http.Handler {
|
func CheckAppDisabled(storage storage.Storage) func(http.Handler) http.Handler {
|
||||||
return func(next http.Handler) http.Handler {
|
return func(next http.Handler) http.Handler {
|
||||||
fn := func(w http.ResponseWriter, r *http.Request) {
|
fn := func(w http.ResponseWriter, r *http.Request) {
|
||||||
@@ -117,13 +117,15 @@ func CheckAppDisabled(storage storage.Storage) func(http.Handler) http.Handler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// PostEvents handle post events
|
||||||
type PostEvents struct{ storage storage.Storage }
|
type PostEvents struct{ storage storage.Storage }
|
||||||
|
|
||||||
|
// NewPostEvents return a new PostEvents handler
|
||||||
func NewPostEvents(storage storage.Storage) *PostEvents {
|
func NewPostEvents(storage storage.Storage) *PostEvents {
|
||||||
return &PostEvents{storage: storage}
|
return &PostEvents{storage: storage}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ServeHTTPC An event consists of a name and data (typically JSON) which may be sent to all subscribers to a particular channel or channels.
|
// ServeHTTP An event consists of a name and data (typically JSON) which may be sent to all subscribers to a particular channel or channels.
|
||||||
// This is conventionally known as triggering an event.
|
// This is conventionally known as triggering an event.
|
||||||
//
|
//
|
||||||
// The body should contain a Hash of parameters encoded as JSON where data parameter itself is JSON encoded.
|
// The body should contain a Hash of parameters encoded as JSON where data parameter itself is JSON encoded.
|
||||||
@@ -192,13 +194,15 @@ func (h *PostEvents) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetChannels handle get channels
|
||||||
type GetChannels struct{ storage storage.Storage }
|
type GetChannels struct{ storage storage.Storage }
|
||||||
|
|
||||||
|
// NewGetChannels return a new GetChannels handler
|
||||||
func NewGetChannels(storage storage.Storage) *GetChannels {
|
func NewGetChannels(storage storage.Storage) *GetChannels {
|
||||||
return &GetChannels{storage: storage}
|
return &GetChannels{storage: storage}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Allows fetching a hash of occupied channels (optionally filtered by prefix),
|
// ServeHTTP Allows fetching a hash of occupied channels (optionally filtered by prefix),
|
||||||
// and optionally one or more attributes for each channel.
|
// and optionally one or more attributes for each channel.
|
||||||
//
|
//
|
||||||
// Notes:
|
// Notes:
|
||||||
@@ -288,13 +292,15 @@ func (h *GetChannels) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetChannel handle get channel
|
||||||
type GetChannel struct{ storage storage.Storage }
|
type GetChannel struct{ storage storage.Storage }
|
||||||
|
|
||||||
|
// NewGetChannel return a new GetChannel handler
|
||||||
func NewGetChannel(storage storage.Storage) *GetChannel {
|
func NewGetChannel(storage storage.Storage) *GetChannel {
|
||||||
return &GetChannel{storage: storage}
|
return &GetChannel{storage: storage}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Fetch info for one channel
|
// ServeHTTP Fetch info for one channel
|
||||||
//
|
//
|
||||||
// Example:
|
// Example:
|
||||||
// {
|
// {
|
||||||
@@ -380,13 +386,15 @@ func (h *GetChannel) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetChannelUsers handle get users from a channel
|
||||||
type GetChannelUsers struct{ storage storage.Storage }
|
type GetChannelUsers struct{ storage storage.Storage }
|
||||||
|
|
||||||
|
// NewGetChannelUsers return a new GetChannelUsers handler
|
||||||
func NewGetChannelUsers(storage storage.Storage) *GetChannelUsers {
|
func NewGetChannelUsers(storage storage.Storage) *GetChannelUsers {
|
||||||
return &GetChannelUsers{storage: storage}
|
return &GetChannelUsers{storage: storage}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Allowed only for presence-channels
|
// ServeHTTP Allowed only for presence-channels
|
||||||
//
|
//
|
||||||
// Example:
|
// Example:
|
||||||
// {
|
// {
|
||||||
|
|||||||
@@ -142,7 +142,7 @@ func Test_getChannels_filter_by_presence_prefix_and_user_count(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// User count only alowed in Presence channels
|
// User count only allowed in Presence channels
|
||||||
func Test_getChannels_filter_by_private_prefix_and_info_user_count(t *testing.T) {
|
func Test_getChannels_filter_by_private_prefix_and_info_user_count(t *testing.T) {
|
||||||
appID := testApp.AppID
|
appID := testApp.AppID
|
||||||
|
|
||||||
|
|||||||
+15
-9
@@ -18,7 +18,7 @@ import (
|
|||||||
"ipe/subscription"
|
"ipe/subscription"
|
||||||
)
|
)
|
||||||
|
|
||||||
// An App
|
// Application represents a Pusher application
|
||||||
type Application struct {
|
type Application struct {
|
||||||
sync.RWMutex
|
sync.RWMutex
|
||||||
|
|
||||||
@@ -38,6 +38,7 @@ type Application struct {
|
|||||||
Stats *expvar.Map `json:"-"`
|
Stats *expvar.Map `json:"-"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewApplication returns a new Application
|
||||||
func NewApplication(
|
func NewApplication(
|
||||||
name,
|
name,
|
||||||
appID,
|
appID,
|
||||||
@@ -83,7 +84,7 @@ func (a *Application) Channels() []*channel.Channel {
|
|||||||
return channels
|
return channels
|
||||||
}
|
}
|
||||||
|
|
||||||
// Only Presence channels
|
// PresenceChannels Only Presence channels
|
||||||
func (a *Application) PresenceChannels() []*channel.Channel {
|
func (a *Application) PresenceChannels() []*channel.Channel {
|
||||||
a.RLock()
|
a.RLock()
|
||||||
defer a.RUnlock()
|
defer a.RUnlock()
|
||||||
@@ -99,7 +100,7 @@ func (a *Application) PresenceChannels() []*channel.Channel {
|
|||||||
return channels
|
return channels
|
||||||
}
|
}
|
||||||
|
|
||||||
// Only Private channels
|
// PrivateChannels Only Private channels
|
||||||
func (a *Application) PrivateChannels() []*channel.Channel {
|
func (a *Application) PrivateChannels() []*channel.Channel {
|
||||||
a.RLock()
|
a.RLock()
|
||||||
defer a.RUnlock()
|
defer a.RUnlock()
|
||||||
@@ -115,7 +116,7 @@ func (a *Application) PrivateChannels() []*channel.Channel {
|
|||||||
return channels
|
return channels
|
||||||
}
|
}
|
||||||
|
|
||||||
// Only Public channels
|
// PublicChannels Only Public channels
|
||||||
func (a *Application) PublicChannels() []*channel.Channel {
|
func (a *Application) PublicChannels() []*channel.Channel {
|
||||||
a.RLock()
|
a.RLock()
|
||||||
defer a.RUnlock()
|
defer a.RUnlock()
|
||||||
@@ -179,7 +180,7 @@ func (a *Application) Connect(conn *connection.Connection) {
|
|||||||
a.Stats.Add("TotalConnections", 1)
|
a.Stats.Add("TotalConnections", 1)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Find a Connection on this Application
|
// FindConnection Find a Connection on this Application
|
||||||
func (a *Application) FindConnection(socketID string) (*connection.Connection, error) {
|
func (a *Application) FindConnection(socketID string) (*connection.Connection, error) {
|
||||||
a.RLock()
|
a.RLock()
|
||||||
defer a.RUnlock()
|
defer a.RUnlock()
|
||||||
@@ -193,7 +194,7 @@ func (a *Application) FindConnection(socketID string) (*connection.Connection, e
|
|||||||
return nil, errors.New("connection not found")
|
return nil, errors.New("connection not found")
|
||||||
}
|
}
|
||||||
|
|
||||||
// DeleteChannel removes the Channel from Application
|
// RemoveChannel removes the Channel from Application
|
||||||
func (a *Application) RemoveChannel(c *channel.Channel) {
|
func (a *Application) RemoveChannel(c *channel.Channel) {
|
||||||
log.Infof("remove the Channel %s from Application %s", c.ID, a.Name)
|
log.Infof("remove the Channel %s from Application %s", c.ID, a.Name)
|
||||||
a.Lock()
|
a.Lock()
|
||||||
@@ -216,7 +217,7 @@ func (a *Application) RemoveChannel(c *channel.Channel) {
|
|||||||
a.Stats.Add("TotalChannels", -1)
|
a.Stats.Add("TotalChannels", -1)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add a new Channel to this APP
|
// AddChannel Add a new Channel to this APP
|
||||||
func (a *Application) AddChannel(c *channel.Channel) {
|
func (a *Application) AddChannel(c *channel.Channel) {
|
||||||
log.Infof("adding a new Channel %s to Application %s", c.ID, a.Name)
|
log.Infof("adding a new Channel %s to Application %s", c.ID, a.Name)
|
||||||
|
|
||||||
@@ -240,7 +241,7 @@ func (a *Application) AddChannel(c *channel.Channel) {
|
|||||||
a.Stats.Add("TotalChannels", 1)
|
a.Stats.Add("TotalChannels", 1)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Returns a Channel from this Application
|
// FindOrCreateChannelByChannelID Returns a Channel from this Application
|
||||||
// If not found then the Channel is created and added to this Application
|
// If not found then the Channel is created and added to this Application
|
||||||
func (a *Application) FindOrCreateChannelByChannelID(n string) *channel.Channel {
|
func (a *Application) FindOrCreateChannelByChannelID(n string) *channel.Channel {
|
||||||
c, err := a.FindChannelByChannelID(n)
|
c, err := a.FindChannelByChannelID(n)
|
||||||
@@ -270,7 +271,7 @@ func (a *Application) FindOrCreateChannelByChannelID(n string) *channel.Channel
|
|||||||
return c
|
return c
|
||||||
}
|
}
|
||||||
|
|
||||||
// Find the Channel by Channel ID
|
// FindChannelByChannelID Find the Channel by Channel ID
|
||||||
func (a *Application) FindChannelByChannelID(n string) (*channel.Channel, error) {
|
func (a *Application) FindChannelByChannelID(n string) (*channel.Channel, error) {
|
||||||
a.RLock()
|
a.RLock()
|
||||||
defer a.RUnlock()
|
defer a.RUnlock()
|
||||||
@@ -284,12 +285,16 @@ func (a *Application) FindChannelByChannelID(n string) (*channel.Channel, error)
|
|||||||
return nil, errors.New("channel does not exists")
|
return nil, errors.New("channel does not exists")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Publish an event into the channel
|
||||||
|
// skip the ignore connection
|
||||||
func (a *Application) Publish(c *channel.Channel, event events.Raw, ignore string) error {
|
func (a *Application) Publish(c *channel.Channel, event events.Raw, ignore string) error {
|
||||||
a.Stats.Add("TotalUniqueMessages", 1)
|
a.Stats.Add("TotalUniqueMessages", 1)
|
||||||
|
|
||||||
return c.Publish(event, ignore)
|
return c.Publish(event, ignore)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Unsubscribe unsubscribe the given connection from the channel
|
||||||
|
// remove the channel from the application if it is empty
|
||||||
func (a *Application) Unsubscribe(c *channel.Channel, conn *connection.Connection) error {
|
func (a *Application) Unsubscribe(c *channel.Channel, conn *connection.Connection) error {
|
||||||
err := c.Unsubscribe(conn)
|
err := c.Unsubscribe(conn)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -303,6 +308,7 @@ func (a *Application) Unsubscribe(c *channel.Channel, conn *connection.Connectio
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Subscribe the connection into the given channel
|
||||||
func (a *Application) Subscribe(c *channel.Channel, conn *connection.Connection, data string) error {
|
func (a *Application) Subscribe(c *channel.Channel, conn *connection.Connection, data string) error {
|
||||||
return c.Subscribe(conn, data)
|
return c.Subscribe(conn, data)
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-1
@@ -86,7 +86,7 @@ func (a *Application) TriggerChannelOccupiedHook(c *channel.Channel) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// channel_vacated
|
// TriggerChannelVacatedHook channel_vacated
|
||||||
// { "name": "channel_vacated", "channel": "test_channel" }
|
// { "name": "channel_vacated", "channel": "test_channel" }
|
||||||
func (a *Application) TriggerChannelVacatedHook(c *channel.Channel) {
|
func (a *Application) TriggerChannelVacatedHook(c *channel.Channel) {
|
||||||
event := newChannelVacatedHook(c)
|
event := newChannelVacatedHook(c)
|
||||||
@@ -98,6 +98,7 @@ func (a *Application) TriggerChannelVacatedHook(c *channel.Channel) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TriggerClientEventHook client_events
|
||||||
// {
|
// {
|
||||||
// "name": "client_event",
|
// "name": "client_event",
|
||||||
// "channel": "name of the channel the event was published on",
|
// "channel": "name of the channel the event was published on",
|
||||||
@@ -121,6 +122,7 @@ func (a *Application) TriggerClientEventHook(c *channel.Channel, s *subscription
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TriggerMemberAddedHook member_added
|
||||||
// {
|
// {
|
||||||
// "name": "member_added",
|
// "name": "member_added",
|
||||||
// "channel": "presence-your_channel_name",
|
// "channel": "presence-your_channel_name",
|
||||||
@@ -136,6 +138,7 @@ func (a *Application) TriggerMemberAddedHook(c *channel.Channel, s *subscription
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TriggerMemberRemovedHook member_removed
|
||||||
// {
|
// {
|
||||||
// "name": "member_removed",
|
// "name": "member_removed",
|
||||||
// "channel": "presence-your_channel_name",
|
// "channel": "presence-your_channel_name",
|
||||||
|
|||||||
+22
-11
@@ -18,8 +18,13 @@ import (
|
|||||||
"ipe/utils"
|
"ipe/utils"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// Option constructor function for Channel
|
||||||
type Option func(*Channel)
|
type Option func(*Channel)
|
||||||
|
|
||||||
|
// ListenerFunc listener function
|
||||||
type ListenerFunc func(*Channel, *subscription.Subscription)
|
type ListenerFunc func(*Channel, *subscription.Subscription)
|
||||||
|
|
||||||
|
// ClientEventListenerFunc listener for client events
|
||||||
type ClientEventListenerFunc func(*Channel, *subscription.Subscription, string, interface{})
|
type ClientEventListenerFunc func(*Channel, *subscription.Subscription, string, interface{})
|
||||||
|
|
||||||
// A Channel
|
// A Channel
|
||||||
@@ -51,30 +56,35 @@ func New(channelID string, options ...Option) *Channel {
|
|||||||
return c
|
return c
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WithMemberAddedListener appends the given ListenerFunc into the memberAddedListeners list
|
||||||
func WithMemberAddedListener(f ListenerFunc) func(*Channel) {
|
func WithMemberAddedListener(f ListenerFunc) func(*Channel) {
|
||||||
return func(c *Channel) {
|
return func(c *Channel) {
|
||||||
c.memberAddedListeners = append(c.memberAddedListeners, f)
|
c.memberAddedListeners = append(c.memberAddedListeners, f)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WithMemberRemovedListener appends the given ListenerFunc into the memberRemovedListeners list
|
||||||
func WithMemberRemovedListener(f ListenerFunc) func(*Channel) {
|
func WithMemberRemovedListener(f ListenerFunc) func(*Channel) {
|
||||||
return func(c *Channel) {
|
return func(c *Channel) {
|
||||||
c.memberRemovedListeners = append(c.memberRemovedListeners, f)
|
c.memberRemovedListeners = append(c.memberRemovedListeners, f)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WithChannelOccupiedListener appends the given ListenerFunc into the channelOccupiedListeners list
|
||||||
func WithChannelOccupiedListener(f ListenerFunc) func(*Channel) {
|
func WithChannelOccupiedListener(f ListenerFunc) func(*Channel) {
|
||||||
return func(c *Channel) {
|
return func(c *Channel) {
|
||||||
c.channelOccupiedListeners = append(c.channelOccupiedListeners, f)
|
c.channelOccupiedListeners = append(c.channelOccupiedListeners, f)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WithChannelVacatedListener appends the given ListenerFunc into the channelVacatedListeners list
|
||||||
func WithChannelVacatedListener(f ListenerFunc) func(*Channel) {
|
func WithChannelVacatedListener(f ListenerFunc) func(*Channel) {
|
||||||
return func(c *Channel) {
|
return func(c *Channel) {
|
||||||
c.channelVacatedListeners = append(c.channelVacatedListeners, f)
|
c.channelVacatedListeners = append(c.channelVacatedListeners, f)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WithClientEventListener appends the given ListenerFunc into the clientEventListeners list
|
||||||
func WithClientEventListener(f ClientEventListenerFunc) func(*Channel) {
|
func WithClientEventListener(f ClientEventListenerFunc) func(*Channel) {
|
||||||
return func(c *Channel) {
|
return func(c *Channel) {
|
||||||
c.clientEventListeners = append(c.clientEventListeners, f)
|
c.clientEventListeners = append(c.clientEventListeners, f)
|
||||||
@@ -95,32 +105,32 @@ func (c *Channel) Subscriptions() []*subscription.Subscription {
|
|||||||
return subscriptions
|
return subscriptions
|
||||||
}
|
}
|
||||||
|
|
||||||
// Return true if the Channel has at least one subscriber
|
// IsOccupied Return true if the Channel has at least one subscriber
|
||||||
func (c *Channel) IsOccupied() bool {
|
func (c *Channel) IsOccupied() bool {
|
||||||
return c.TotalSubscriptions() > 0
|
return c.TotalSubscriptions() > 0
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if the type of the Channel is presence or is private
|
// IsPresenceOrPrivate Check if the type of the Channel is presence or is private
|
||||||
func (c *Channel) IsPresenceOrPrivate() bool {
|
func (c *Channel) IsPresenceOrPrivate() bool {
|
||||||
return c.IsPresence() || c.IsPrivate()
|
return c.IsPresence() || c.IsPrivate()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if the type of the Channel is public
|
// IsPublic Check if the type of the Channel is public
|
||||||
func (c *Channel) IsPublic() bool {
|
func (c *Channel) IsPublic() bool {
|
||||||
return !c.IsPresenceOrPrivate()
|
return !c.IsPresenceOrPrivate()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if the type of the Channel is presence
|
// IsPresence Check if the type of the Channel is presence
|
||||||
func (c *Channel) IsPresence() bool {
|
func (c *Channel) IsPresence() bool {
|
||||||
return utils.IsPresenceChannel(c.ID)
|
return utils.IsPresenceChannel(c.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if the type of the Channel is private
|
// IsPrivate Check if the type of the Channel is private
|
||||||
func (c *Channel) IsPrivate() bool {
|
func (c *Channel) IsPrivate() bool {
|
||||||
return utils.IsPrivateChannel(c.ID)
|
return utils.IsPrivateChannel(c.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get the total of subscribers
|
// TotalSubscriptions Get the total of subscribers
|
||||||
func (c *Channel) TotalSubscriptions() int {
|
func (c *Channel) TotalSubscriptions() int {
|
||||||
c.RLock()
|
c.RLock()
|
||||||
defer c.RUnlock()
|
defer c.RUnlock()
|
||||||
@@ -128,7 +138,7 @@ func (c *Channel) TotalSubscriptions() int {
|
|||||||
return len(c.subscriptions)
|
return len(c.subscriptions)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get the total of users.
|
// TotalUsers Get the total of users.
|
||||||
func (c *Channel) TotalUsers() int {
|
func (c *Channel) TotalUsers() int {
|
||||||
c.RLock()
|
c.RLock()
|
||||||
defer c.RUnlock()
|
defer c.RUnlock()
|
||||||
@@ -142,7 +152,7 @@ func (c *Channel) TotalUsers() int {
|
|||||||
return len(total)
|
return len(total)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add a new subscriber to the Channel
|
// Subscribe Add a new subscriber to the Channel
|
||||||
func (c *Channel) Subscribe(conn *connection.Connection, channelData string) error {
|
func (c *Channel) Subscribe(conn *connection.Connection, channelData string) error {
|
||||||
log.Infof("Subscribing %s to Channel %s", conn.SocketID, c.ID)
|
log.Infof("Subscribing %s to Channel %s", conn.SocketID, c.ID)
|
||||||
|
|
||||||
@@ -219,7 +229,7 @@ func (c *Channel) IsSubscribed(conn *connection.Connection) bool {
|
|||||||
return exists
|
return exists
|
||||||
}
|
}
|
||||||
|
|
||||||
// Remove the subscriber from the Channel
|
// Unsubscribe Remove the subscriber from the Channel
|
||||||
// It destroy the Channel if the channels does not have any subscribers.
|
// It destroy the Channel if the channels does not have any subscribers.
|
||||||
func (c *Channel) Unsubscribe(conn *connection.Connection) error {
|
func (c *Channel) Unsubscribe(conn *connection.Connection) error {
|
||||||
log.Infof("unsubscribe %s from Channel %s", conn.SocketID, c.ID)
|
log.Infof("unsubscribe %s from Channel %s", conn.SocketID, c.ID)
|
||||||
@@ -254,7 +264,7 @@ func (c *Channel) Unsubscribe(conn *connection.Connection) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Publish a MemberAddedEvent to all subscriptions
|
// PublishMemberAddedEvent Publish a MemberAddedEvent to all subscriptions
|
||||||
func (c *Channel) PublishMemberAddedEvent(data string, subscription *subscription.Subscription) {
|
func (c *Channel) PublishMemberAddedEvent(data string, subscription *subscription.Subscription) {
|
||||||
c.RLock()
|
c.RLock()
|
||||||
defer c.RUnlock()
|
defer c.RUnlock()
|
||||||
@@ -266,7 +276,7 @@ func (c *Channel) PublishMemberAddedEvent(data string, subscription *subscriptio
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Publish a MemberRemovedEvent to all subscriptions
|
// PublishMemberRemovedEvent Publish a MemberRemovedEvent to all subscriptions
|
||||||
func (c *Channel) PublishMemberRemovedEvent(subscription *subscription.Subscription) {
|
func (c *Channel) PublishMemberRemovedEvent(subscription *subscription.Subscription) {
|
||||||
c.RLock()
|
c.RLock()
|
||||||
defer c.RUnlock()
|
defer c.RUnlock()
|
||||||
@@ -279,6 +289,7 @@ func (c *Channel) PublishMemberRemovedEvent(subscription *subscription.Subscript
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Publish messages to all Subscribers
|
// Publish messages to all Subscribers
|
||||||
|
// skip the ignore connection
|
||||||
func (c *Channel) Publish(event events.Raw, ignore string) error {
|
func (c *Channel) Publish(event events.Raw, ignore string) error {
|
||||||
c.RLock()
|
c.RLock()
|
||||||
defer c.RUnlock()
|
defer c.RUnlock()
|
||||||
|
|||||||
+4
-1
@@ -4,7 +4,7 @@
|
|||||||
|
|
||||||
package config
|
package config
|
||||||
|
|
||||||
// The config file
|
// File config file
|
||||||
type File struct {
|
type File struct {
|
||||||
Host string `yaml:"host"` // The host, eg: :8080 will start on 0.0.0.0:8080
|
Host string `yaml:"host"` // The host, eg: :8080 will start on 0.0.0.0:8080
|
||||||
SSL SSL `yaml:"ssl"`
|
SSL SSL `yaml:"ssl"`
|
||||||
@@ -12,6 +12,7 @@ type File struct {
|
|||||||
Apps []Application `yaml:"apps"`
|
Apps []Application `yaml:"apps"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SSL related configuration options
|
||||||
type SSL struct {
|
type SSL struct {
|
||||||
Enabled bool `yaml:"enabled"`
|
Enabled bool `yaml:"enabled"`
|
||||||
Host string `yaml:"host"`
|
Host string `yaml:"host"`
|
||||||
@@ -19,6 +20,7 @@ type SSL struct {
|
|||||||
CertFile string `yaml:"cert_file"`
|
CertFile string `yaml:"cert_file"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Application related configuration options
|
||||||
type Application struct {
|
type Application struct {
|
||||||
Name string `yaml:"name"`
|
Name string `yaml:"name"`
|
||||||
AppID string `yaml:"app_id"`
|
AppID string `yaml:"app_id"`
|
||||||
@@ -30,6 +32,7 @@ type Application struct {
|
|||||||
WebHooks Webhooks `yaml:"webhooks"`
|
WebHooks Webhooks `yaml:"webhooks"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Webhooks related configuration options
|
||||||
type Webhooks struct {
|
type Webhooks struct {
|
||||||
Enabled bool `yaml:"enabled"`
|
Enabled bool `yaml:"enabled"`
|
||||||
URL string `yaml:"url"`
|
URL string `yaml:"url"`
|
||||||
|
|||||||
@@ -25,7 +25,7 @@ type Connection struct {
|
|||||||
CreatedAt time.Time
|
CreatedAt time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new Subscriber
|
// New Create a new Subscriber
|
||||||
func New(socketID string, s Socket) *Connection {
|
func New(socketID string, s Socket) *Connection {
|
||||||
log.Infof("Creating a new Subscriber %+v", socketID)
|
log.Infof("Creating a new Subscriber %+v", socketID)
|
||||||
|
|
||||||
|
|||||||
+29
-14
@@ -12,6 +12,14 @@ import (
|
|||||||
"ipe/subscription"
|
"ipe/subscription"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// SubscribeData data for Subscribe event
|
||||||
|
type SubscribeData struct {
|
||||||
|
Channel string `json:"channel"`
|
||||||
|
Auth string `json:"auth,omitempty"`
|
||||||
|
ChannelData string `json:"channel_data,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Subscribe event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher:subscribe",
|
// "event": "pusher:subscribe",
|
||||||
// "data": {
|
// "data": {
|
||||||
@@ -20,27 +28,23 @@ import (
|
|||||||
// "channelData": "extra data"
|
// "channelData": "extra data"
|
||||||
// }
|
// }
|
||||||
// }
|
// }
|
||||||
type SubscribeData struct {
|
|
||||||
Channel string `json:"channel"`
|
|
||||||
Auth string `json:"auth,omitempty"`
|
|
||||||
ChannelData string `json:"channel_data,omitempty"`
|
|
||||||
}
|
|
||||||
|
|
||||||
type Subscribe struct {
|
type Subscribe struct {
|
||||||
Event string `json:"event"`
|
Event string `json:"event"`
|
||||||
Data SubscribeData `json:"data"`
|
Data SubscribeData `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new subscribe event with the specified channel and data
|
// NewSubscribe Create a new subscribe event with the specified channel and data
|
||||||
func NewSubscribe(channel, auth, channelData string) Subscribe {
|
func NewSubscribe(channel, auth, channelData string) Subscribe {
|
||||||
data := SubscribeData{Channel: channel, Auth: auth, ChannelData: channelData}
|
data := SubscribeData{Channel: channel, Auth: auth, ChannelData: channelData}
|
||||||
return Subscribe{Event: "pusher:subscribe", Data: data}
|
return Subscribe{Event: "pusher:subscribe", Data: data}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// UnsubscribeData
|
||||||
type UnsubscribeData struct {
|
type UnsubscribeData struct {
|
||||||
Channel string `json:"channel"`
|
Channel string `json:"channel"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Unsubscribe event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher:unsubscribe",
|
// "event": "pusher:unsubscribe",
|
||||||
// "data": {
|
// "data": {
|
||||||
@@ -58,6 +62,7 @@ func NewUnsubscribe(channel string) Unsubscribe {
|
|||||||
return Unsubscribe{Event: "pusher:unsubscribe", Data: data}
|
return Unsubscribe{Event: "pusher:unsubscribe", Data: data}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SubscriptionSucceeded event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher_internal:subscription_succeeded",
|
// "event": "pusher_internal:subscription_succeeded",
|
||||||
// "channel": "the channel"
|
// "channel": "the channel"
|
||||||
@@ -73,8 +78,7 @@ func NewSubscriptionSucceeded(channel, data string) SubscriptionSucceeded {
|
|||||||
return SubscriptionSucceeded{Event: "pusher_internal:subscription_succeeded", Channel: channel, Data: data}
|
return SubscriptionSucceeded{Event: "pusher_internal:subscription_succeeded", Channel: channel, Data: data}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Data Subscription Succeed
|
// SubscriptionSucceededPresenceData Data Subscription Succeed
|
||||||
|
|
||||||
// "{
|
// "{
|
||||||
// \"presence\": {
|
// \"presence\": {
|
||||||
// \"ids\": [\"11814b369700141b222a3f3791cec2d9\",\"71dd6a29da2a4833336d2a964becf820\"],
|
// \"ids\": [\"11814b369700141b222a3f3791cec2d9\",\"71dd6a29da2a4833336d2a964becf820\"],
|
||||||
@@ -97,6 +101,7 @@ type SubscriptionSucceededPresenceData struct {
|
|||||||
Count int `json:"count"`
|
Count int `json:"count"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewSubscriptionSucceedPresenceData returns new SubscriptionSucceededPresenceData
|
||||||
func NewSubscriptionSucceedPresenceData(subscriptions map[string]*subscription.Subscription) SubscriptionSucceededPresenceData {
|
func NewSubscriptionSucceedPresenceData(subscriptions map[string]*subscription.Subscription) SubscriptionSucceededPresenceData {
|
||||||
event := SubscriptionSucceededPresenceData{}
|
event := SubscriptionSucceededPresenceData{}
|
||||||
|
|
||||||
@@ -123,6 +128,7 @@ func NewSubscriptionSucceedPresenceData(subscriptions map[string]*subscription.S
|
|||||||
return event
|
return event
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Pong event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher:pong",
|
// "event": "pusher:pong",
|
||||||
// "data": {}
|
// "data": {}
|
||||||
@@ -132,11 +138,12 @@ type Pong struct {
|
|||||||
Data string `json:"data"`
|
Data string `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new pong event
|
// NewPong Create a new pong event
|
||||||
func NewPong() Pong {
|
func NewPong() Pong {
|
||||||
return Pong{Event: "pusher:pong", Data: "{}"}
|
return Pong{Event: "pusher:pong", Data: "{}"}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Ping event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher:ping",
|
// "event": "pusher:ping",
|
||||||
// "data": {}
|
// "data": {}
|
||||||
@@ -146,11 +153,12 @@ type Ping struct {
|
|||||||
Data string `json:"data"`
|
Data string `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new ping event
|
// NewPing Create a new ping event
|
||||||
func NewPing() Ping {
|
func NewPing() Ping {
|
||||||
return Ping{Event: "pusher:ping", Data: "{}"}
|
return Ping{Event: "pusher:ping", Data: "{}"}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Error event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher:error",
|
// "event": "pusher:error",
|
||||||
// "data": {
|
// "data": {
|
||||||
@@ -163,7 +171,7 @@ type Error struct {
|
|||||||
Data interface{} `json:"data"`
|
Data interface{} `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new error event
|
// NewError Create a new error event
|
||||||
// Pusher protocol is very strange in some parts
|
// Pusher protocol is very strange in some parts
|
||||||
// It send null in some errors.
|
// It send null in some errors.
|
||||||
func NewError(code int, message string) Error {
|
func NewError(code int, message string) Error {
|
||||||
@@ -183,6 +191,7 @@ func NewError(code int, message string) Error {
|
|||||||
return Error{Event: "pusher:error", Data: data}
|
return Error{Event: "pusher:error", Data: data}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ConnectionEstablished event
|
||||||
// {
|
// {
|
||||||
// "event" : "pusher:connection_established",
|
// "event" : "pusher:connection_established",
|
||||||
// "data" : {
|
// "data" : {
|
||||||
@@ -195,7 +204,7 @@ type ConnectionEstablished struct {
|
|||||||
Data string `json:"data"`
|
Data string `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new connection established event using the specified socketId
|
// NewConnectionEstablished Create a new connection established event using the specified socketId
|
||||||
func NewConnectionEstablished(socketID string) ConnectionEstablished {
|
func NewConnectionEstablished(socketID string) ConnectionEstablished {
|
||||||
b, err := json.Marshal(struct {
|
b, err := json.Marshal(struct {
|
||||||
SocketID string `json:"socket_id"`
|
SocketID string `json:"socket_id"`
|
||||||
@@ -211,6 +220,7 @@ func NewConnectionEstablished(socketID string) ConnectionEstablished {
|
|||||||
return ConnectionEstablished{Event: "pusher:connection_established", Data: string(b)}
|
return ConnectionEstablished{Event: "pusher:connection_established", Data: string(b)}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MemberAdded event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher_internal:member_added",
|
// "event": "pusher_internal:member_added",
|
||||||
// "channel": "presence-example-channel",
|
// "channel": "presence-example-channel",
|
||||||
@@ -222,10 +232,12 @@ type MemberAdded struct {
|
|||||||
Data string `json:"data"`
|
Data string `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewMemberAdded creates a new MemberAdded event
|
||||||
func NewMemberAdded(channel, data string) MemberAdded {
|
func NewMemberAdded(channel, data string) MemberAdded {
|
||||||
return MemberAdded{Event: "pusher_internal:member_added", Channel: channel, Data: data}
|
return MemberAdded{Event: "pusher_internal:member_added", Channel: channel, Data: data}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MemberRemoved event
|
||||||
// {
|
// {
|
||||||
// "event": "pusher_internal:member_removed",
|
// "event": "pusher_internal:member_removed",
|
||||||
// "channel": "presence-example-channel",
|
// "channel": "presence-example-channel",
|
||||||
@@ -237,6 +249,7 @@ type MemberRemoved struct {
|
|||||||
Data string `json:"data"`
|
Data string `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewMemberRemoved returns a new MemberRemoved event
|
||||||
func NewMemberRemoved(channel string, userID string) MemberRemoved {
|
func NewMemberRemoved(channel string, userID string) MemberRemoved {
|
||||||
data, err := json.Marshal(struct {
|
data, err := json.Marshal(struct {
|
||||||
UserID string `json:"user_id"`
|
UserID string `json:"user_id"`
|
||||||
@@ -251,6 +264,7 @@ func NewMemberRemoved(channel string, userID string) MemberRemoved {
|
|||||||
return MemberRemoved{Event: "pusher_internal:member_removed", Channel: channel, Data: string(data)}
|
return MemberRemoved{Event: "pusher_internal:member_removed", Channel: channel, Data: string(data)}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Raw event, usually used for client events
|
||||||
// {
|
// {
|
||||||
// "event": "client-?",
|
// "event": "client-?",
|
||||||
// "channel": "The channel",
|
// "channel": "The channel",
|
||||||
@@ -262,13 +276,14 @@ type Raw struct {
|
|||||||
Data json.RawMessage `json:"data"`
|
Data json.RawMessage `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Response event
|
||||||
type Response struct {
|
type Response struct {
|
||||||
Event string `json:"event"`
|
Event string `json:"event"`
|
||||||
Channel string `json:"channel"`
|
Channel string `json:"channel"`
|
||||||
Data interface{} `json:"data"`
|
Data interface{} `json:"data"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// The response event that is broadcasted to the client sockets
|
// NewResponse The response event that is broadcasted to the client sockets
|
||||||
func NewResponse(name, channel string, data interface{}) Response {
|
func NewResponse(name, channel string, data interface{}) Response {
|
||||||
return Response{Event: name, Channel: channel, Data: data}
|
return Response{Event: name, Channel: channel, Data: data}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,6 +4,8 @@ package mocks
|
|||||||
// used in the test suite
|
// used in the test suite
|
||||||
type MockSocket struct{}
|
type MockSocket struct{}
|
||||||
|
|
||||||
|
// WriteJSON always returns nil
|
||||||
|
// used in the test suite
|
||||||
func (s MockSocket) WriteJSON(i interface{}) error {
|
func (s MockSocket) WriteJSON(i interface{}) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -19,15 +19,18 @@ type Storage interface {
|
|||||||
AddApp(application *app.Application) error
|
AddApp(application *app.Application) error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// InMemory in memory implementation of Storage
|
||||||
type InMemory struct {
|
type InMemory struct {
|
||||||
sync.RWMutex
|
sync.RWMutex
|
||||||
Apps []*app.Application
|
Apps []*app.Application
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewInMemory returns an InMemory storage
|
||||||
func NewInMemory() Storage {
|
func NewInMemory() Storage {
|
||||||
return &InMemory{}
|
return &InMemory{}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// AddApp adds app into memory
|
||||||
func (db *InMemory) AddApp(application *app.Application) error {
|
func (db *InMemory) AddApp(application *app.Application) error {
|
||||||
db.Lock()
|
db.Lock()
|
||||||
defer db.Unlock()
|
defer db.Unlock()
|
||||||
|
|||||||
@@ -6,14 +6,14 @@ package subscription
|
|||||||
|
|
||||||
import "ipe/connection"
|
import "ipe/connection"
|
||||||
|
|
||||||
// A Channel Subscription
|
// Subscription A Channel Subscription
|
||||||
type Subscription struct {
|
type Subscription struct {
|
||||||
Connection *connection.Connection
|
Connection *connection.Connection
|
||||||
ID string
|
ID string
|
||||||
Data string
|
Data string
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a new Subscription
|
// New Create a new Subscription
|
||||||
func New(conn *connection.Connection, data string) *Subscription {
|
func New(conn *connection.Connection, data string) *Subscription {
|
||||||
return &Subscription{Connection: conn, Data: data}
|
return &Subscription{Connection: conn, Data: data}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -34,15 +34,17 @@ var upgrader = websocket.Upgrader{
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Websocket handler for real time websocket messages
|
||||||
type Websocket struct {
|
type Websocket struct {
|
||||||
storage storage.Storage
|
storage storage.Storage
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewWebsocket returns a new Websocket handler
|
||||||
func NewWebsocket(storage storage.Storage) *Websocket {
|
func NewWebsocket(storage storage.Storage) *Websocket {
|
||||||
return &Websocket{storage: storage}
|
return &Websocket{storage: storage}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Websocket GET /app/{key}
|
// ServeHTTP Websocket GET /app/{key}
|
||||||
func (h *Websocket) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
func (h *Websocket) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||||
conn, err := upgrader.Upgrade(w, r, nil)
|
conn, err := upgrader.Upgrade(w, r, nil)
|
||||||
defer func() {
|
defer func() {
|
||||||
|
|||||||
Reference in New Issue
Block a user