Spaces:
Running
Running
restructure repo into production layout; add whatsapp-service, tests, ddl, docs, scripts; scrub hardcoded secrets
83d6851 | 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} | |
| } |