-
Notifications
You must be signed in to change notification settings - Fork 179
Expand file tree
/
Copy pathmain.go
More file actions
141 lines (119 loc) · 2.65 KB
/
Copy pathmain.go
File metadata and controls
141 lines (119 loc) · 2.65 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
package main
import (
"fmt"
"log"
"net/http"
"github.com/gorilla/websocket"
"github.com/anthdm/hollywood/actor"
)
type server struct {
ctx *actor.Context
sessions map[string]*actor.PID
}
func newServer() actor.Receiver {
return &server{
sessions: make(map[string]*actor.PID),
}
}
func (s *server) Receive(ctx *actor.Context) {
switch msg := ctx.Message().(type) {
case actor.Started:
log.Println("Server started on port 8080")
s.serve()
s.ctx = ctx
_ = msg
case message:
s.broadcast(ctx.Sender(), msg)
default:
fmt.Printf("Server received %v\n", msg)
}
}
func (s *server) serve() {
go func() {
http.HandleFunc("/ws", s.handleWebsocket)
http.ListenAndServe(":8080", nil)
}()
}
func (s *server) handleWebsocket(w http.ResponseWriter, r *http.Request) {
fmt.Println("New connection")
upgrader := websocket.Upgrader{}
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
fmt.Println(err)
return
}
username := r.URL.Query().Get("username")
pid := s.ctx.SpawnChild(newUser(username, conn, s.ctx.PID()), username)
s.sessions[pid.GetID()] = pid
}
func (s *server) broadcast(sender *actor.PID, msg message) {
for _, pid := range s.sessions {
if !pid.Equals(sender) {
s.ctx.Send(pid, msg)
}
}
}
type message struct {
Content string `json:"content"`
Owner string `json:"owner"`
}
type user struct {
conn *websocket.Conn
ctx *actor.Context
serverPid *actor.PID
Name string
}
func newUser(name string, conn *websocket.Conn, serverPid *actor.PID) actor.Producer {
return func() actor.Receiver {
return &user{
Name: name,
conn: conn,
serverPid: serverPid,
}
}
}
func (u *user) Receive(ctx *actor.Context) {
switch msg := ctx.Message().(type) {
case actor.Started:
u.ctx = ctx
go u.listen()
case message:
u.send(&msg)
case actor.Stopped:
_ = msg
u.conn.Close()
default:
fmt.Printf("%s received %v\n", u.Name, msg)
}
}
func (u *user) listen() {
var msg message
for {
if err := u.conn.ReadJSON(&msg); err != nil {
fmt.Printf("Error reading message: %v\n", err)
return
}
msg.Owner = u.Name
go u.handleMessage(msg)
}
}
func (u *user) handleMessage(msg message) {
switch msg.Content {
case "exit": // Send exit message to stop the actor and close the websocket connection
u.ctx.Engine().Poison(u.ctx.PID())
default:
// Note that this is the server pid, so it will broadcast the message
u.ctx.Send(u.serverPid, msg)
}
}
func (u *user) send(msg *message) {
if err := u.conn.WriteJSON(msg); err != nil {
fmt.Printf("Error writing message: %v\n", err)
return
}
}
func main() {
engine := actor.NewEngine()
engine.Spawn(newServer, "server")
select {}
}