llm-ready-data / whatsapp-service /pkg /instance /repository /instance_repository.go
validops-east-1's picture
restructure repo into production layout; add whatsapp-service, tests, ddl, docs, scripts; scrub hardcoded secrets
83d6851
Raw
History Blame Contribute Delete
11.8 kB
package instance_repository
import (
"context"
"fmt"
"time"
instance_model "agentdeck-whatsapp-service/pkg/instance/model"
"agentdeck-whatsapp-service/pkg/supabase"
"github.com/google/uuid"
)
type InstanceRepository interface {
Create(instance instance_model.Instance) (*instance_model.Instance, error)
GetInstanceByID(instanceId string) (*instance_model.Instance, error)
GetConnectedInstanceByID(instanceId string) (*instance_model.Instance, error)
GetInstanceByToken(token string) (*instance_model.Instance, error)
GetInstanceByName(name string) (*instance_model.Instance, error)
Update(*instance_model.Instance) error
UpdateConnected(userId string, status bool, disconnectReason string) error
UpdateQrcode(userId string, qr string) error
UpdateProxy(userId string, proxy string) error
UpdateJid(userId string, jid string) error
GetAllConnectedInstances() ([]*instance_model.Instance, error)
GetAllConnectedInstancesByClientName(clientName string) ([]*instance_model.Instance, error)
GetAll(clientName string) ([]*instance_model.Instance, error)
Delete(instanceId string) error
GetAdvancedSettings(instanceId string) (*instance_model.AdvancedSettings, error)
UpdateAdvancedSettings(instanceId string, settings *instance_model.AdvancedSettings) error
}
// instanceRow is the snake_case PostgREST representation of an instance row.
// The public API model keeps its camelCase JSON tags; this row maps only the DB.
type instanceRow struct {
Id string `json:"id"`
Name string `json:"name"`
Token string `json:"token"`
Webhook string `json:"webhook"`
RabbitmqEnable string `json:"rabbitmq_enable"`
WebSocketEnable string `json:"web_socket_enable"`
NatsEnable string `json:"nats_enable"`
Jid string `json:"jid"`
Qrcode string `json:"qrcode"`
Connected bool `json:"connected"`
Expiration int64 `json:"expiration"`
DisconnectReason string `json:"disconnect_reason"`
Events string `json:"events"`
OsName string `json:"os_name"`
Proxy string `json:"proxy"`
ClientName string `json:"client_name"`
CreatedAt *string `json:"created_at,omitempty"`
AlwaysOnline bool `json:"always_online"`
RejectCall bool `json:"reject_call"`
MsgRejectCall string `json:"msg_reject_call"`
ReadMessages bool `json:"read_messages"`
IgnoreGroups bool `json:"ignore_groups"`
IgnoreStatus bool `json:"ignore_status"`
}
type instanceRepository struct {
supa *supabase.Client
}
func toRow(instance instance_model.Instance) instanceRow {
return instanceRow{
Id: instance.Id,
Name: instance.Name,
Token: instance.Token,
Webhook: instance.Webhook,
RabbitmqEnable: instance.RabbitmqEnable,
WebSocketEnable: instance.WebSocketEnable,
NatsEnable: instance.NatsEnable,
Jid: instance.Jid,
Qrcode: instance.Qrcode,
Connected: instance.Connected,
Expiration: instance.Expiration,
DisconnectReason: instance.DisconnectReason,
Events: instance.Events,
OsName: instance.OsName,
Proxy: instance.Proxy,
ClientName: instance.ClientName,
AlwaysOnline: instance.AlwaysOnline,
RejectCall: instance.RejectCall,
MsgRejectCall: instance.MsgRejectCall,
ReadMessages: instance.ReadMessages,
IgnoreGroups: instance.IgnoreGroups,
IgnoreStatus: instance.IgnoreStatus,
}
}
func fromRow(r *instanceRow) *instance_model.Instance {
instance := &instance_model.Instance{
Id: r.Id,
Name: r.Name,
Token: r.Token,
Webhook: r.Webhook,
RabbitmqEnable: r.RabbitmqEnable,
WebSocketEnable: r.WebSocketEnable,
NatsEnable: r.NatsEnable,
Jid: r.Jid,
Qrcode: r.Qrcode,
Connected: r.Connected,
Expiration: r.Expiration,
DisconnectReason: r.DisconnectReason,
Events: r.Events,
OsName: r.OsName,
Proxy: r.Proxy,
ClientName: r.ClientName,
AlwaysOnline: r.AlwaysOnline,
RejectCall: r.RejectCall,
MsgRejectCall: r.MsgRejectCall,
ReadMessages: r.ReadMessages,
IgnoreGroups: r.IgnoreGroups,
IgnoreStatus: r.IgnoreStatus,
}
if r.CreatedAt != nil {
if t, err := time.Parse(time.RFC3339, *r.CreatedAt); err == nil {
instance.CreatedAt = t
}
}
return instance
}
func (i *instanceRepository) Create(instance instance_model.Instance) (*instance_model.Instance, error) {
if instance.Id == "" {
instance.Id = uuid.New().String()
}
row := toRow(instance)
ctx := context.Background()
var created []instanceRow
if err := i.supa.Table("wp_instances").Insert(ctx, row, "return=representation", &created); err != nil {
return nil, err
}
if len(created) == 0 {
return nil, fmt.Errorf("no instance row returned")
}
return fromRow(&created[0]), nil
}
func (i *instanceRepository) GetInstanceByToken(token string) (*instance_model.Instance, error) {
q := supabase.NewQuery().Eq("token", token).Limit(1)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
if len(rows) == 0 {
return nil, fmt.Errorf("instances: no row for token")
}
return fromRow(&rows[0]), nil
}
func (i *instanceRepository) GetInstanceByName(name string) (*instance_model.Instance, error) {
q := supabase.NewQuery().Eq("name", name).Limit(1)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
if len(rows) == 0 {
return nil, fmt.Errorf("instances not found for name")
}
return fromRow(&rows[0]), nil
}
func (i *instanceRepository) GetInstanceByID(instanceId string) (*instance_model.Instance, error) {
if _, err := uuid.Parse(instanceId); err != nil {
return nil, fmt.Errorf("invalid UUID format: %v", err)
}
q := supabase.NewQuery().Eq("id", instanceId).Limit(1)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
if len(rows) == 0 {
return nil, fmt.Errorf("instance not found")
}
return fromRow(&rows[0]), nil
}
func (i *instanceRepository) GetConnectedInstanceByID(instanceId string) (*instance_model.Instance, error) {
q := supabase.NewQuery().Eq("id", instanceId).Eq("connected", "true").Limit(1)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
if len(rows) == 0 {
return nil, fmt.Errorf("connected instance not found")
}
return fromRow(&rows[0]), nil
}
func (i *instanceRepository) Update(instance *instance_model.Instance) error {
row := toRow(*instance)
ctx := context.Background()
q := supabase.NewQuery().Eq("id", instance.Id)
return i.supa.Table("wp_instances").Update(ctx, q, row)
}
func (i *instanceRepository) UpdateConnected(userId string, connected bool, disconnectReason string) error {
ctx := context.Background()
q := supabase.NewQuery().Eq("id", userId)
body := map[string]interface{}{
"connected": connected,
"disconnect_reason": disconnectReason,
}
return i.supa.Table("wp_instances").Update(ctx, q, body)
}
func (i *instanceRepository) UpdateQrcode(userId string, qr string) error {
ctx := context.Background()
q := supabase.NewQuery().Eq("id", userId)
return i.supa.Table("wp_instances").Update(ctx, q, map[string]interface{}{"qrcode": qr})
}
func (i *instanceRepository) UpdateProxy(userId string, proxy string) error {
ctx := context.Background()
q := supabase.NewQuery().Eq("id", userId)
return i.supa.Table("wp_instances").Update(ctx, q, map[string]interface{}{"proxy": proxy})
}
func (i *instanceRepository) UpdateJid(userId string, jid string) error {
ctx := context.Background()
q := supabase.NewQuery().Eq("id", userId)
return i.supa.Table("wp_instances").Update(ctx, q, map[string]interface{}{"jid": jid})
}
func (i *instanceRepository) GetAllConnectedInstances() ([]*instance_model.Instance, error) {
q := supabase.NewQuery().Eq("connected", "true")
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
out := make([]*instance_model.Instance, 0, len(rows))
for idx := range rows {
out = append(out, fromRow(&rows[idx]))
}
return out, nil
}
func (i *instanceRepository) GetAllConnectedInstancesByClientName(clientName string) ([]*instance_model.Instance, error) {
q := supabase.NewQuery().Eq("connected", "true").Eq("client_name", clientName)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
out := make([]*instance_model.Instance, 0, len(rows))
for idx := range rows {
out = append(out, fromRow(&rows[idx]))
}
return out, nil
}
func (i *instanceRepository) GetAll(clientName string) ([]*instance_model.Instance, error) {
q := supabase.NewQuery().Eq("client_name", clientName)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
out := make([]*instance_model.Instance, 0, len(rows))
for idx := range rows {
out = append(out, fromRow(&rows[idx]))
}
return out, nil
}
func (i *instanceRepository) Delete(instanceId string) error {
ctx := context.Background()
// Cascade delete via a Supabase RPC, which deletes related rows and the
// instance atomically (see ddl/006_functions.sql).
var result interface{}
if err := i.supa.RPC(ctx, "wp_delete_instance", map[string]interface{}{"p_instance_id": instanceId}, &result); err != nil {
// Fallback: delete rows individually if the RPC is unavailable.
q := supabase.NewQuery().Eq("id", instanceId)
if derr := i.supa.Table("wp_instances").Delete(ctx, q); derr != nil {
return fmt.Errorf("failed to delete instance: %v", derr)
}
}
return nil
}
func (i *instanceRepository) GetAdvancedSettings(instanceId string) (*instance_model.AdvancedSettings, error) {
if _, err := uuid.Parse(instanceId); err != nil {
return nil, fmt.Errorf("invalid UUID format: %v", err)
}
q := supabase.NewQuery().
Eq("id", instanceId).
Select("always_online, reject_call, msg_reject_call, read_messages, ignore_groups, ignore_status").
Limit(1)
var rows []instanceRow
ctx := context.Background()
if err := i.supa.Table("wp_instances").Select(ctx, q, &rows); err != nil {
return nil, err
}
if len(rows) == 0 {
return nil, fmt.Errorf("instance not found")
}
return &instance_model.AdvancedSettings{
AlwaysOnline: rows[0].AlwaysOnline,
RejectCall: rows[0].RejectCall,
MsgRejectCall: rows[0].MsgRejectCall,
ReadMessages: rows[0].ReadMessages,
IgnoreGroups: rows[0].IgnoreGroups,
IgnoreStatus: rows[0].IgnoreStatus,
}, nil
}
func (i *instanceRepository) UpdateAdvancedSettings(instanceId string, settings *instance_model.AdvancedSettings) error {
if _, err := uuid.Parse(instanceId); err != nil {
return fmt.Errorf("invalid UUID format: %v", err)
}
ctx := context.Background()
q := supabase.NewQuery().Eq("id", instanceId)
body := map[string]interface{}{
"always_online": settings.AlwaysOnline,
"reject_call": settings.RejectCall,
"msg_reject_call": settings.MsgRejectCall,
"read_messages": settings.ReadMessages,
"ignore_groups": settings.IgnoreGroups,
"ignore_status": settings.IgnoreStatus,
}
return i.supa.Table("wp_instances").Update(ctx, q, body)
}
func NewInstanceRepository(supa *supabase.Client) InstanceRepository {
return &instanceRepository{supa: supa}
}