package redis import ( "fmt" "time" ) type PubSub struct { *baseClient } func (c *Client) PubSub() *PubSub { return &PubSub{ baseClient: &baseClient{ opt: c.opt, connPool: newSingleConnPool(c.connPool, nil, false), }, } } func (c *Client) Publish(channel, message string) *IntReq { req := NewIntReq("PUBLISH", channel, message) c.Process(req) return req } type Message struct { Channel string Payload string } type PMessage struct { Channel string Pattern string Payload string } type Subscription struct { Kind string Channel string Count int } func (c *PubSub) Receive() (interface{}, error) { return c.ReceiveTimeout(0) } func (c *PubSub) ReceiveTimeout(timeout time.Duration) (interface{}, error) { cn, err := c.conn() if err != nil { return nil, err } cn.readTimeout = timeout replyIface, err := NewIfaceSliceReq().ParseReply(cn.Rd) if err != nil { return nil, err } reply, ok := replyIface.([]interface{}) if !ok { return nil, fmt.Errorf("redis: unexpected reply type %T", replyIface) } switch msgName := reply[0].(string); msgName { case "subscribe", "unsubscribe", "psubscribe", "punsubscribe": return &Subscription{ Kind: msgName, Channel: reply[1].(string), Count: int(reply[2].(int64)), }, nil case "message": return &Message{ Channel: reply[1].(string), Payload: reply[2].(string), }, nil case "pmessage": return &PMessage{ Pattern: reply[1].(string), Channel: reply[2].(string), Payload: reply[3].(string), }, nil default: return nil, fmt.Errorf("redis: unsupported message name: %q", msgName) } } func (c *PubSub) subscribe(cmd string, channels ...string) error { cn, err := c.conn() if err != nil { return err } args := append([]string{cmd}, channels...) req := NewIfaceSliceReq(args...) return c.writeReq(cn, req) } func (c *PubSub) Subscribe(channels ...string) error { return c.subscribe("SUBSCRIBE", channels...) } func (c *PubSub) PSubscribe(patterns ...string) error { return c.subscribe("PSUBSCRIBE", patterns...) } func (c *PubSub) unsubscribe(cmd string, channels ...string) error { cn, err := c.conn() if err != nil { return err } args := append([]string{cmd}, channels...) req := NewIfaceSliceReq(args...) return c.writeReq(cn, req) } func (c *PubSub) Unsubscribe(channels ...string) error { return c.unsubscribe("UNSUBSCRIBE", channels...) } func (c *PubSub) PUnsubscribe(patterns ...string) error { return c.unsubscribe("PUNSUBSCRIBE", patterns...) }