-
Notifications
You must be signed in to change notification settings - Fork 5
/
Copy pathpgsql_transport.go
129 lines (98 loc) · 2.8 KB
/
pgsql_transport.go
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
package gosumer
import (
"context"
"fmt"
"log"
"strings"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
type PgDatabase struct {
Host string
Port uint16
User string
Password string
Database string
TableName string
}
var pool *pgxpool.Pool
var _ Transport = (*PgDatabase)(nil)
func (database PgDatabase) connect() error {
var err error
pool, err = pgxpool.New(context.Background(), fmt.Sprintf("postgres://%s:%s@%s:%d/%s", database.User, database.Password, database.Host, database.Port, database.Database))
if err != nil {
return err
}
return nil
}
func (database PgDatabase) GetChannelName() string {
// Symfony uses the format "schema.table" for channel name
if strings.Contains(database.TableName, ".") {
return fmt.Sprintf(`"%s"`, strings.Replace(database.TableName, `"`, "", -1))
}
return database.TableName
}
func (database PgDatabase) listen(fn process, message any, sec int) error {
err := database.connect()
if err != nil {
return err
}
defer pool.Close()
database.listenEvery(sec, fn, message)
log.Printf("Successfully connected to the database!")
_, err = pool.Exec(context.Background(), fmt.Sprintf("LISTEN %s", database.GetChannelName()))
if err != nil {
return err
}
defer pool.Exec(context.Background(), fmt.Sprintf("UNLISTEN %s", database.GetChannelName()))
conn, err := pool.Acquire(context.Background())
if err != nil {
return err
}
defer conn.Release()
for {
_, err := conn.Conn().WaitForNotification(context.Background())
if err != nil {
return err
}
err = database.processMessage(fn, message)
if err != nil {
continue
}
}
}
func (database PgDatabase) listenEvery(seconds int, fn process, message any) {
delay := time.Duration(seconds) * time.Second
go func() error {
for {
<-time.After(delay)
_ = database.processMessage(fn, message)
}
}()
}
func (database PgDatabase) delete(id int) error {
_, err := pool.Exec(context.Background(), fmt.Sprintf("DELETE FROM %s WHERE id=$1", database.TableName), id)
if err != nil {
return err
}
return nil
}
func (database PgDatabase) processMessage(fn process, message any) error {
var messengerMessage MessengerMessage
row := pool.QueryRow(context.Background(), fmt.Sprintf("SELECT * FROM %s WHERE delivered_at IS NULL AND queue_name = 'go' ORDER BY id DESC LIMIT 1", database.TableName))
if err := row.Scan(&messengerMessage.ID, &messengerMessage.Body, &messengerMessage.Headers, &messengerMessage.QueueName, &messengerMessage.CreatedAt, &messengerMessage.AvailableAt, &messengerMessage.DeliveredAt); err != nil {
return err
}
msg, err := formatMessage(messengerMessage.Body, message)
if err != nil {
return err
}
e := make(chan error)
go fn(msg, e)
processErr := <-e
if processErr != nil {
return processErr
}
go database.delete(messengerMessage.ID)
return nil
}