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} }