package services

import (
	config "astrology-api/configs"
	"astrology-api/constants"
	dto "astrology-api/dto/consultation"
	models "astrology-api/models/usermodel"
	"astrology-api/notifications"
	"astrology-api/repositories"
	"errors"
	"fmt"
	"log"
	"strings"
	"time"

	"gorm.io/gorm"
)

// The waiting queue (astrologer_waiting_queue).
//
// A customer who finds an astrologer busy joins the queue. When the
// astrologer frees up the queue is processed FIFO, across chat and call
// together because one astrologer serves one session at a time:
//
//	WAITING -> NOTIFIED -> CONNECTING -> CONNECTED -> COMPLETED
//
// The customer at the head is NOTIFIED and holds the astrologer for the
// notification window; nobody else can request that astrologer meanwhile
// (queueGate). Tapping Connect moves them to CONNECTING, and the consultation
// they then place is linked to the row, which follows that consultation from
// there: accepted is CONNECTED, ended is COMPLETED, missed is EXPIRED,
// rejected or cancelled is REJECTED.
//
// The API schedules nothing, so every timeout is applied lazily by
// advanceQueue — on every queue call, on /consultation/start and precheck,
// whenever a session ends, when the astrologer comes online, and in the admin
// sweep the panel's cron already drives. Every transition runs under a lock on
// the astrologer's row, so two of those running at once cannot notify two
// customers or reorder the queue.

const (
	defaultQueueMaxWaitingMinutes        = 30.0
	defaultQueueNotificationWindowSecs   = 60.0
	defaultQueueConnectionTimeoutSeconds = 30.0
	defaultQueueMaxCapacity              = 20.0
	queuePollIntervalSeconds             = 5
)

//////////////////////////////////////////////////////////////
// Configuration (admin panel > Settings > Consultation & Chat)
//////////////////////////////////////////////////////////////

func (s *consultationService) queueMaxWaitingMinutes() int {
	return positiveFlag(s.repository.GetSystemFlagFloat(constants.FlagQueueMaxWaitingMinutes, defaultQueueMaxWaitingMinutes), defaultQueueMaxWaitingMinutes)
}

func (s *consultationService) queueNotificationWindowSeconds() int {
	return positiveFlag(s.repository.GetSystemFlagFloat(constants.FlagQueueNotificationWindowSeconds, defaultQueueNotificationWindowSecs), defaultQueueNotificationWindowSecs)
}

func (s *consultationService) queueConnectionTimeoutSeconds() int {
	return positiveFlag(s.repository.GetSystemFlagFloat(constants.FlagQueueConnectionTimeoutSeconds, defaultQueueConnectionTimeoutSeconds), defaultQueueConnectionTimeoutSeconds)
}

// queueMaxCapacity is 0 for unlimited.
func (s *consultationService) queueMaxCapacity() int {

	value := int(s.repository.GetSystemFlagFloat(constants.FlagQueueMaxCapacity, defaultQueueMaxCapacity))

	if value < 0 {
		return 0
	}

	return value
}

// positiveFlag keeps a zero or negative timeout from expiring every row on
// sight.
func positiveFlag(value float64, fallback float64) int {

	if value <= 0 {
		return int(fallback)
	}

	return int(value)
}

//////////////////////////////////////////////////////////////
// Package-level entry points for the astrologer stack
//////////////////////////////////////////////////////////////

func queueEngine() *consultationService {
	return &consultationService{repository: repositories.NewConsultationRepository(config.DB)}
}

// KickWaitingQueue settles an astrologer's queue and hands the turn to the
// next customer if the astrologer is free. Called when the astrologer comes
// online. Never fails the caller.
func KickWaitingQueue(astrologerID uint) {

	if config.DB == nil || astrologerID == 0 {
		return
	}

	if _, err := queueEngine().advanceQueue(astrologerID); err != nil {
		log.Printf("waiting queue: advance astrologer %d: %v", astrologerID, err)
	}
}

// AstrologerIDByUserID maps the astrologer app's token (users.id) to the
// astrologers.id the queue is keyed on.
func AstrologerIDByUserID(astrologerUserID uint) (uint, error) {

	astrologer, err := queueEngine().repository.GetAstrologerByUserID(astrologerUserID)

	if err != nil {
		return 0, err
	}

	if astrologer == nil {
		return 0, errors.New("astrologer not found")
	}

	return astrologer.ID, nil
}

// AddToWaitingQueue is the astrologer adding a customer from their app. Same
// duplicate and capacity rules as the customer joining.
func AddToWaitingQueue(astrologerUserID uint, userID uint, waitingType string) (*dto.QueueEntryResponse, error) {

	astrologerID, err := AstrologerIDByUserID(astrologerUserID)

	if err != nil {
		return nil, err
	}

	medium := constants.MediumChat

	if strings.ToUpper(strings.TrimSpace(waitingType)) != "CHAT" {
		medium = constants.MediumAudio
	}

	return queueEngine().joinQueue(userID, dto.JoinQueueRequest{
		AstrologerID: astrologerID,
		Medium:       medium,
	}, true)
}

// RemoveFromWaitingQueue is the astrologer removing a customer: REJECTED.
func RemoveFromWaitingQueue(astrologerUserID uint, queueID uint) error {

	astrologerID, err := AstrologerIDByUserID(astrologerUserID)

	if err != nil {
		return err
	}

	s := queueEngine()

	row, err := s.repository.GetQueueEntry(queueID)

	if err != nil {
		return err
	}

	if row == nil || row.AstrologerID != astrologerID {
		return errors.New("waiting user not found")
	}

	return s.closeQueueEntry(row, constants.QueueStatusRejected, constants.QueueEndRemovedByAstrologer, false)
}

//////////////////////////////////////////////////////////////
// Customer: join
//////////////////////////////////////////////////////////////

func (s *consultationService) JoinQueue(userID uint, request dto.JoinQueueRequest) (*dto.QueueEntryResponse, error) {
	return s.joinQueue(userID, request, false)
}

func (s *consultationService) joinQueue(
	userID uint,
	request dto.JoinQueueRequest,
	byAstrologer bool,
) (*dto.QueueEntryResponse, error) {

	medium := strings.ToUpper(strings.TrimSpace(request.Medium))

	if !constants.IsValidMedium(medium) {
		return nil, errors.New("medium must be CHAT, AUDIO or VIDEO")
	}

	if userID == 0 {
		return nil, errors.New("user is required")
	}

	astrologer, err := s.repository.GetAstrologerByID(request.AstrologerID)

	if err != nil {
		return nil, err
	}

	if astrologer == nil {
		return nil, errors.New("astrologer not found")
	}

	if astrologer.UserID == userID {
		return nil, errors.New("you cannot join your own waiting queue")
	}

	if !astrologer.IsActive {
		return nil, errors.New("astrologer is not available")
	}

	name := astrologerName(astrologer)
	waitingType := waitingTypeFor(medium)

	// Settle timeouts and the hand-over first, so the duplicate check,
	// capacity and position below are all current.
	if _, err := s.advanceQueue(astrologer.ID); err != nil {
		return nil, err
	}

	tx := s.repository.Begin()

	if tx.Error != nil {
		return nil, tx.Error
	}

	committed := false

	defer func() {
		if !committed {
			tx.Rollback()
		}
	}()

	if err := s.repository.LockAstrologerQueue(tx, astrologer.ID); err != nil {
		return nil, err
	}

	rows, err := s.repository.ListQueueForAstrologer(tx, astrologer.ID)

	if err != nil {
		return nil, err
	}

	active := activeQueueRows(rows)

	//------------------------------------------------
	// Duplicate Request Protection
	//------------------------------------------------

	for index, row := range active {

		if row.UserID == userID && row.WaitingType == waitingType {
			return nil, &QueueDuplicateError{
				Message: fmt.Sprintf(
					"You are already in the waiting queue for Astrologer %s. Your current position is %d.",
					name,
					index+1,
				),
			}
		}
	}

	//------------------------------------------------
	// Nobody To Wait For
	//------------------------------------------------

	if !byAstrologer && len(active) == 0 {

		if reason, _ := s.availabilityReason(astrologer.ID, medium); reason == "" {
			return nil, fmt.Errorf(
				"Astrologer %s is available now. Please start the %s directly.",
				name,
				queueMediumWord(waitingType),
			)
		}
	}

	//------------------------------------------------
	// Capacity
	//------------------------------------------------

	if capacity := s.queueMaxCapacity(); capacity > 0 && len(active) >= capacity {
		return nil, fmt.Errorf(
			"The waiting queue for Astrologer %s is full. Please try again later.",
			name,
		)
	}

	now := time.Now()
	expires := now.Add(time.Duration(s.queueMaxWaitingMinutes()) * time.Minute)

	row := &models.AstrologerWaitingQueue{
		AstrologerID:  astrologer.ID,
		UserID:        userID,
		WaitingType:   waitingType,
		Medium:        medium,
		QueuePosition: len(active) + 1,
		Status:        constants.QueueStatusWaiting,
		RequestedAt:   &now,
		WaitExpiresAt: &expires,
		CreatedAt:     &now,
		UpdatedAt:     &now,
	}

	if err := s.repository.CreateQueueEntry(tx, row); err != nil {
		return nil, err
	}

	if err := tx.Commit().Error; err != nil {
		return nil, err
	}

	committed = true

	if !byAstrologer {
		notifications.Notify().CustomerWaiting(astrologer.ID, userID, waitingType, row.QueuePosition)
	}

	response, err := s.queueEntryView(row.ID)

	if err != nil {
		return nil, err
	}

	return response, nil
}

//////////////////////////////////////////////////////////////
// Customer: status
//////////////////////////////////////////////////////////////

// GetQueueStatus returns one row (any status) when queueID is set, otherwise
// every active row the customer holds. Polled by the app; each call also
// settles that astrologer's queue, so the countdowns do not wait on the cron.
func (s *consultationService) GetQueueStatus(userID uint, queueID uint) (*dto.QueueListResponse, error) {

	response := &dto.QueueListResponse{
		List:       []dto.QueueEntryResponse{},
		ServerTime: time.Now().Format(stampLayout),
	}

	if queueID > 0 {

		row, err := s.repository.GetQueueEntry(queueID)

		if err != nil {
			return nil, err
		}

		if row == nil || row.UserID != userID {
			return nil, errors.New("waiting request not found")
		}

		_, _ = s.advanceQueue(row.AstrologerID)

		view, err := s.queueEntryView(row.ID)

		if err != nil {
			return nil, err
		}

		response.List = append(response.List, *view)

		return response, nil
	}

	rows, err := s.repository.ListUserQueueEntries(userID)

	if err != nil {
		return nil, err
	}

	settled := map[uint]bool{}

	for _, row := range rows {

		if !settled[row.AstrologerID] {
			_, _ = s.advanceQueue(row.AstrologerID)
			settled[row.AstrologerID] = true
		}
	}

	// Re-read: settling may have expired or promoted some of them.
	rows, err = s.repository.ListUserQueueEntries(userID)

	if err != nil {
		return nil, err
	}

	for _, row := range rows {

		view, err := s.queueEntryView(row.ID)

		if err != nil {
			return nil, err
		}

		response.List = append(response.List, *view)
	}

	return response, nil
}

//////////////////////////////////////////////////////////////
// Customer: connect
//////////////////////////////////////////////////////////////

// ConnectQueue is the customer tapping Connect after their turn notification:
// NOTIFIED -> CONNECTING. They then have the connection timeout to place the
// request with /consultation/start, which links it to this row.
func (s *consultationService) ConnectQueue(userID uint, queueID uint) (*dto.QueueEntryResponse, error) {

	row, err := s.repository.GetQueueEntry(queueID)

	if err != nil {
		return nil, err
	}

	if row == nil || row.UserID != userID {
		return nil, errors.New("waiting request not found")
	}

	if _, err := s.advanceQueue(row.AstrologerID); err != nil {
		return nil, err
	}

	tx := s.repository.Begin()

	if tx.Error != nil {
		return nil, tx.Error
	}

	committed := false

	defer func() {
		if !committed {
			tx.Rollback()
		}
	}()

	if err := s.repository.LockAstrologerQueue(tx, row.AstrologerID); err != nil {
		return nil, err
	}

	rows, err := s.repository.ListQueueForAstrologer(tx, row.AstrologerID)

	if err != nil {
		return nil, err
	}

	current := findQueueRow(rows, queueID)

	if current == nil {

		// Not in play any more: it reached a terminal status.
		latest, err := s.repository.GetQueueEntry(queueID)

		if err != nil {
			return nil, err
		}

		status := ""

		if latest != nil {
			status = latest.Status
		}

		return nil, fmt.Errorf("this waiting request is no longer active (%s)", status)
	}

	switch current.Status {

	case constants.QueueStatusConnecting:
		// Already connecting: a double tap is not an error.

	case constants.QueueStatusNotified:

		now := time.Now()
		expires := now.Add(time.Duration(s.queueConnectionTimeoutSeconds()) * time.Second)

		if err := s.repository.UpdateQueueEntry(tx, current.ID, map[string]interface{}{
			"status":             constants.QueueStatusConnecting,
			"connecting_at":      now,
			"connect_expires_at": expires,
			"updated_at":         now,
		}); err != nil {
			return nil, err
		}

	case constants.QueueStatusWaiting:

		position := 0

		for index, active := range activeQueueRows(rows) {
			if active.ID == current.ID {
				position = index + 1
			}
		}

		return nil, fmt.Errorf("it is not your turn yet. Your current position is %d", position)

	default:
		return nil, fmt.Errorf("this waiting request is no longer active (%s)", current.Status)
	}

	if err := tx.Commit().Error; err != nil {
		return nil, err
	}

	committed = true

	return s.queueEntryView(queueID)
}

//////////////////////////////////////////////////////////////
// Customer: cancel
//////////////////////////////////////////////////////////////

func (s *consultationService) CancelQueue(userID uint, queueID uint) (*dto.QueueEntryResponse, error) {

	row, err := s.repository.GetQueueEntry(queueID)

	if err != nil {
		return nil, err
	}

	if row == nil || row.UserID != userID {
		return nil, errors.New("waiting request not found")
	}

	if err := s.closeQueueEntry(row, constants.QueueStatusCancelled, constants.QueueEndCustomerCancelled, true); err != nil {
		return nil, err
	}

	return s.queueEntryView(queueID)
}

// closeQueueEntry moves an active row to a terminal status under the
// astrologer lock, then hands the turn on if it was holding it.
func (s *consultationService) closeQueueEntry(
	row *models.AstrologerWaitingQueue,
	status string,
	reason string,
	byCustomer bool,
) error {

	tx := s.repository.Begin()

	if tx.Error != nil {
		return tx.Error
	}

	committed := false

	defer func() {
		if !committed {
			tx.Rollback()
		}
	}()

	if err := s.repository.LockAstrologerQueue(tx, row.AstrologerID); err != nil {
		return err
	}

	rows, err := s.repository.ListQueueForAstrologer(tx, row.AstrologerID)

	if err != nil {
		return err
	}

	current := findQueueRow(rows, row.ID)

	if current == nil || !constants.IsActiveQueueStatus(current.Status) {

		latest, _ := s.repository.GetQueueEntry(row.ID)

		if latest != nil {
			return fmt.Errorf("this waiting request is already %s", strings.ToLower(latest.Status))
		}

		return errors.New("this waiting request is no longer active")
	}

	// The request is already ringing: that is cancelled through the
	// consultation, and this row follows it.
	if byCustomer && current.ConsultationID != nil {
		return errors.New("your consultation request is already in progress, cancel it with /consultation/cancel")
	}

	now := time.Now()

	if err := s.setQueueStatus(tx, current, status, reason, now); err != nil {
		return err
	}

	if err := s.renumberQueue(tx, rows); err != nil {
		return err
	}

	if err := tx.Commit().Error; err != nil {
		return err
	}

	committed = true

	_, _ = s.advanceQueue(row.AstrologerID)

	return nil
}

//////////////////////////////////////////////////////////////
// The engine
//////////////////////////////////////////////////////////////

type queueEvent struct {
	expired bool
	row     models.AstrologerWaitingQueue
}

// advanceQueue applies every timeout to one astrologer's queue, follows the
// linked consultations, and — if nobody holds the astrologer and they are free
// — notifies the next eligible customer, FIFO. Idempotent and safe to call
// concurrently.
func (s *consultationService) advanceQueue(astrologerID uint) (dto.QueueSweepResult, error) {

	result := dto.QueueSweepResult{Astrologers: 1}

	astrologer, err := s.repository.GetAstrologerByID(astrologerID)

	if err != nil {
		return result, err
	}

	tx := s.repository.Begin()

	if tx.Error != nil {
		return result, tx.Error
	}

	committed := false

	defer func() {
		if !committed {
			tx.Rollback()
		}
	}()

	if err := s.repository.LockAstrologerQueue(tx, astrologerID); err != nil {
		return result, err
	}

	rows, err := s.repository.ListQueueForAstrologer(tx, astrologerID)

	if err != nil {
		return result, err
	}

	if len(rows) == 0 {
		return result, nil
	}

	now := time.Now()

	var events []queueEvent
	var holder *models.AstrologerWaitingQueue

	// availability is cached per medium for this pass.
	availability := map[string]string{}

	available := func(medium string) bool {

		reason, seen := availability[medium]

		if !seen {
			reason, _ = s.availabilityReason(astrologerID, medium)
			availability[medium] = reason
		}

		return reason == ""
	}

	//------------------------------------------------
	// Timeouts And Linked Consultations
	//------------------------------------------------

	for index := range rows {

		row := &rows[index]

		switch row.Status {

		case constants.QueueStatusConnected:

			if row.ConsultationID == nil {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusCompleted, constants.QueueEndSessionCompleted, now); err != nil {
					return result, err
				}
				result.Completed++
				continue
			}

			consultation, err := s.repository.GetConsultationByID(*row.ConsultationID)

			if err != nil {
				return result, err
			}

			if consultation == nil || !constants.IsOpenStatus(consultation.Status) {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusCompleted, constants.QueueEndSessionCompleted, now); err != nil {
					return result, err
				}
				result.Completed++
			}

		case constants.QueueStatusConnecting:

			if row.ConsultationID != nil {

				consultation, err := s.repository.GetConsultationByID(*row.ConsultationID)

				if err != nil {
					return result, err
				}

				if consultation == nil {
					if err := s.setQueueStatus(tx, row, constants.QueueStatusRejected, constants.QueueEndRequestCancelled, now); err != nil {
						return result, err
					}
					result.Rejected++
					continue
				}

				if consultation.Status == constants.ConsultationStatusRequested {
					holder = row
					continue
				}

				status, reason := queueOutcomeFor(consultation)

				if status == constants.QueueStatusConnected {
					if err := s.repository.UpdateQueueEntry(tx, row.ID, map[string]interface{}{
						"status":       status,
						"connected_at": now,
						"updated_at":   now,
					}); err != nil {
						return result, err
					}
					row.Status = status
					continue
				}

				if status != "" {
					if err := s.setQueueStatus(tx, row, status, reason, now); err != nil {
						return result, err
					}
					countQueueOutcome(&result, status)
				}

				continue
			}

			if row.ConnectExpiresAt != nil && row.ConnectExpiresAt.Before(now) {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusExpired, constants.QueueEndConnectTimeout, now); err != nil {
					return result, err
				}
				result.Expired++
				events = append(events, queueEvent{expired: true, row: *row})
				continue
			}

			holder = row

		case constants.QueueStatusNotified:

			if row.NotifyExpiresAt != nil && row.NotifyExpiresAt.Before(now) {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusExpired, constants.QueueEndNotifyTimeout, now); err != nil {
					return result, err
				}
				result.Expired++
				events = append(events, queueEvent{expired: true, row: *row})
				continue
			}

			// The astrologer went offline or busy again before this
			// customer connected.
			if !available(queueMedium(row)) {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusRejected, constants.QueueEndAstrologerUnavailable, now); err != nil {
					return result, err
				}
				result.Rejected++
				continue
			}

			holder = row

		case constants.QueueStatusWaiting:

			if row.WaitExpiresAt != nil && row.WaitExpiresAt.Before(now) {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusExpired, constants.QueueEndWaitTimeout, now); err != nil {
					return result, err
				}
				result.Expired++
				events = append(events, queueEvent{expired: true, row: *row})
			}
		}
	}

	//------------------------------------------------
	// FIFO Hand-Over
	//------------------------------------------------
	//
	// Only when nobody holds the astrologer. The first WAITING customer whose
	// medium the astrologer currently takes is notified; one who cannot
	// connect right now is SKIPPED and the next is tried.

	if holder == nil && astrologer != nil {

		for index := range rows {

			row := &rows[index]

			if row.Status != constants.QueueStatusWaiting {
				continue
			}

			medium := queueMedium(row)

			if !available(medium) {
				continue
			}

			if reason := s.queueIneligibility(row.UserID, astrologer, medium); reason != "" {
				if err := s.setQueueStatus(tx, row, constants.QueueStatusSkipped, reason, now); err != nil {
					return result, err
				}
				result.Skipped++
				continue
			}

			expires := now.Add(time.Duration(s.queueNotificationWindowSeconds()) * time.Second)

			if err := s.repository.UpdateQueueEntry(tx, row.ID, map[string]interface{}{
				"status":            constants.QueueStatusNotified,
				"notified_at":       now,
				"notify_expires_at": expires,
				"updated_at":        now,
			}); err != nil {
				return result, err
			}

			row.Status = constants.QueueStatusNotified
			row.NotifiedAt = &now
			row.NotifyExpiresAt = &expires

			result.Notified++
			events = append(events, queueEvent{row: *row})

			break
		}
	}

	if err := s.renumberQueue(tx, rows); err != nil {
		return result, err
	}

	if err := tx.Commit().Error; err != nil {
		return result, err
	}

	committed = true

	//------------------------------------------------
	// Tell The Customers
	//------------------------------------------------

	window := s.queueNotificationWindowSeconds()

	for _, event := range events {

		if event.expired {
			notifications.Notify().QueueExpired(event.row.UserID, event.row.AstrologerID, event.row.WaitingType, event.row.ID)
			continue
		}

		notifications.Notify().QueueTurn(event.row.UserID, event.row.AstrologerID, event.row.WaitingType, event.row.ID, window)
	}

	return result, nil
}

// queueIneligibility is why a customer at the head of the queue cannot take
// their turn right now, or "".
func (s *consultationService) queueIneligibility(
	userID uint,
	astrologer *models.Astrologer,
	medium string,
) string {

	if open, err := s.repository.GetOpenConsultation(userID); err == nil && open != nil {
		return constants.QueueEndSessionOpenElsewhere
	}

	freeMinutes, err := s.freeChatEntitlement(userID, medium)

	if err != nil || freeMinutes > 0 {
		return ""
	}

	balance, err := s.walletBalance(userID)

	if err != nil {
		return ""
	}

	required := rateForMedium(astrologer, medium) * s.repository.GetSystemFlagFloat(
		constants.FlagMinConsultationMinutes,
		defaultMinConsultationMinutes,
	)

	if balance < required {
		return constants.QueueEndInsufficientBalance
	}

	return ""
}

// queueGate keeps the astrologer for the customer whose turn it is. Returns a
// reason and message when someone else holds them.
func (s *consultationService) queueGate(
	userID uint,
	astrologer *models.Astrologer,
) (*models.AstrologerWaitingQueue, string, string) {

	_, _ = s.advanceQueue(astrologer.ID)

	rows, err := s.repository.ListQueueForAstrologer(nil, astrologer.ID)

	if err != nil {
		return nil, "", ""
	}

	for index := range rows {

		row := &rows[index]

		if row.Status != constants.QueueStatusNotified && row.Status != constants.QueueStatusConnecting {
			continue
		}

		if row.UserID == userID {
			return row, "", ""
		}

		return nil, "QUEUE_RESERVED", fmt.Sprintf(
			"Astrologer %s is reserved for the next customer in the waiting queue. Join the queue to wait for your turn.",
			astrologerName(astrologer),
		)
	}

	return nil, "", ""
}

// linkQueueToConsultation attaches the consultation the turn-holder just
// placed to their queue row. Best effort: the consultation already exists.
func (s *consultationService) linkQueueToConsultation(row *models.AstrologerWaitingQueue, consultationID uint) {

	if row == nil {
		return
	}

	now := time.Now()

	updates := map[string]interface{}{
		"status":          constants.QueueStatusConnecting,
		"consultation_id": consultationID,
		"updated_at":      now,
	}

	if row.ConnectingAt == nil {
		updates["connecting_at"] = now
	}

	if err := s.repository.UpdateQueueEntry(nil, row.ID, updates); err != nil {
		log.Printf("waiting queue: link row %d to consultation %d: %v", row.ID, consultationID, err)
	}
}

// syncQueue moves the queue row a consultation was placed from, inside the
// consultation's own transaction. Errors are logged, never returned: a queue
// write must not roll back a bill, and MySQL keeps the transaction usable
// after a failed statement.
func (s *consultationService) syncQueue(tx *gorm.DB, consultation *models.Consultation) {

	status, reason := queueOutcomeFor(consultation)

	if status == "" {
		return
	}

	now := time.Now()

	if status == constants.QueueStatusConnected {

		if err := s.repository.CloseQueueForConsultation(tx, consultation.ID,
			[]string{constants.QueueStatusConnecting},
			map[string]interface{}{"status": status, "connected_at": now, "updated_at": now},
		); err != nil {
			log.Printf("waiting queue: consultation %d connected: %v", consultation.ID, err)
		}

		return
	}

	if err := s.repository.CloseQueueForConsultation(tx, consultation.ID,
		[]string{constants.QueueStatusConnected},
		map[string]interface{}{
			"status":     constants.QueueStatusCompleted,
			"end_reason": constants.QueueEndSessionCompleted,
			"ended_at":   now,
			"updated_at": now,
		},
	); err != nil {
		log.Printf("waiting queue: consultation %d completed: %v", consultation.ID, err)
	}

	if err := s.repository.CloseQueueForConsultation(tx, consultation.ID,
		[]string{constants.QueueStatusConnecting},
		map[string]interface{}{"status": status, "end_reason": reason, "ended_at": now, "updated_at": now},
	); err != nil {
		log.Printf("waiting queue: consultation %d closed: %v", consultation.ID, err)
	}
}

// queueOutcomeFor maps a consultation's status onto its queue row.
func queueOutcomeFor(consultation *models.Consultation) (string, string) {

	switch consultation.Status {

	case constants.ConsultationStatusAccepted, constants.ConsultationStatusOngoing:
		return constants.QueueStatusConnected, ""

	case constants.ConsultationStatusCompleted:
		return constants.QueueStatusCompleted, constants.QueueEndSessionCompleted

	case constants.ConsultationStatusMissed:
		return constants.QueueStatusExpired, constants.QueueEndRequestMissed

	case constants.ConsultationStatusRejected:
		return constants.QueueStatusRejected, constants.QueueEndRequestRejected

	case constants.ConsultationStatusCancelled:

		if consultation.EndReason == constants.EndReasonCancelled {
			return constants.QueueStatusCancelled, constants.QueueEndRequestCancelled
		}

		return constants.QueueStatusRejected, constants.QueueEndRequestCancelled
	}

	return "", ""
}

func countQueueOutcome(result *dto.QueueSweepResult, status string) {

	switch status {
	case constants.QueueStatusExpired:
		result.Expired++
	case constants.QueueStatusRejected, constants.QueueStatusCancelled:
		result.Rejected++
	case constants.QueueStatusCompleted:
		result.Completed++
	}
}

// sweepQueue is the queue pass of the admin sweep.
func (s *consultationService) sweepQueue(limit int) (dto.QueueSweepResult, error) {

	total := dto.QueueSweepResult{}

	ids, err := s.repository.ListAstrologersWithQueue(limit)

	if err != nil {
		return total, err
	}

	var failures []string

	for _, id := range ids {

		result, err := s.advanceQueue(id)

		if err != nil {
			failures = append(failures, fmt.Sprintf("astrologer %d: %v", id, err))
			continue
		}

		total.Astrologers++
		total.Expired += result.Expired
		total.Notified += result.Notified
		total.Skipped += result.Skipped
		total.Rejected += result.Rejected
		total.Completed += result.Completed
	}

	if len(failures) > 0 {
		return total, errors.New(strings.Join(failures, "; "))
	}

	return total, nil
}

//////////////////////////////////////////////////////////////
// Helpers
//////////////////////////////////////////////////////////////

func (s *consultationService) setQueueStatus(
	tx *gorm.DB,
	row *models.AstrologerWaitingQueue,
	status string,
	reason string,
	now time.Time,
) error {

	updates := map[string]interface{}{
		"status":     status,
		"end_reason": reason,
		"ended_at":   now,
		"updated_at": now,
	}

	if status == constants.QueueStatusExpired {
		updates["expired_at"] = now
	}

	if err := s.repository.UpdateQueueEntry(tx, row.ID, updates); err != nil {
		return err
	}

	row.Status = status
	row.EndReason = reason
	row.EndedAt = &now

	return nil
}

// renumberQueue rewrites queue_position for the active rows, FIFO, chat and
// call together. Rows passed in are already in requested_at order.
func (s *consultationService) renumberQueue(tx *gorm.DB, rows []models.AstrologerWaitingQueue) error {

	position := 0

	for index := range rows {

		row := &rows[index]

		if !constants.IsActiveQueueStatus(row.Status) {
			continue
		}

		position++

		if row.QueuePosition == position {
			continue
		}

		if err := s.repository.UpdateQueueEntry(tx, row.ID, map[string]interface{}{
			"queue_position": position,
		}); err != nil {
			return err
		}

		row.QueuePosition = position
	}

	return nil
}

func activeQueueRows(rows []models.AstrologerWaitingQueue) []models.AstrologerWaitingQueue {

	active := make([]models.AstrologerWaitingQueue, 0, len(rows))

	for _, row := range rows {
		if constants.IsActiveQueueStatus(row.Status) {
			active = append(active, row)
		}
	}

	return active
}

func findQueueRow(rows []models.AstrologerWaitingQueue, id uint) *models.AstrologerWaitingQueue {

	for index := range rows {
		if rows[index].ID == id {
			return &rows[index]
		}
	}

	return nil
}

// queueMedium is the row's medium, defaulting a row written before the medium
// column existed from its waiting type.
func queueMedium(row *models.AstrologerWaitingQueue) string {

	if medium := strings.ToUpper(strings.TrimSpace(row.Medium)); constants.IsValidMedium(medium) {
		return medium
	}

	if row.WaitingType == "CHAT" {
		return constants.MediumChat
	}

	return constants.MediumAudio
}

func queueMediumWord(waitingType string) string {

	if waitingType == "CHAT" {
		return "chat"
	}

	return "call"
}

// QueueDuplicateError is the customer already holding a place for this
// astrologer and service type. The controller answers it with a 409.
type QueueDuplicateError struct {
	Message string
}

func (e *QueueDuplicateError) Error() string {
	return e.Message
}

//////////////////////////////////////////////////////////////
// Response
//////////////////////////////////////////////////////////////

func (s *consultationService) queueEntryView(queueID uint) (*dto.QueueEntryResponse, error) {

	row, err := s.repository.GetQueueEntry(queueID)

	if err != nil {
		return nil, err
	}

	if row == nil {
		return nil, errors.New("waiting request not found")
	}

	astrologer, err := s.repository.GetAstrologerByID(row.AstrologerID)

	if err != nil {
		return nil, err
	}

	rows, err := s.repository.ListQueueForAstrologer(nil, row.AstrologerID)

	if err != nil {
		return nil, err
	}

	active := activeQueueRows(rows)
	now := time.Now()

	view := &dto.QueueEntryResponse{
		QueueID:             row.ID,
		AstrologerID:        row.AstrologerID,
		WaitingType:         row.WaitingType,
		Medium:              queueMedium(row),
		Status:              row.Status,
		EndReason:           row.EndReason,
		RequestedAt:         formatStamp(row.RequestedAt),
		WaitExpiresAt:       formatStamp(row.WaitExpiresAt),
		NotifiedAt:          formatStamp(row.NotifiedAt),
		NotifyExpiresAt:     formatStamp(row.NotifyExpiresAt),
		ConnectingAt:        formatStamp(row.ConnectingAt),
		ConnectExpiresAt:    formatStamp(row.ConnectExpiresAt),
		ConnectedAt:         formatStamp(row.ConnectedAt),
		EndedAt:             formatStamp(row.EndedAt),
		ServerTime:          now.Format(stampLayout),
		PollIntervalSeconds: queuePollIntervalSeconds,
		NextAction:          "NONE",
	}

	name := "the astrologer"

	if astrologer != nil {
		name = astrologerName(astrologer)
		view.AstrologerName = name
		view.AstrologerImage = astrologer.ProfileImage
	}

	if row.ConsultationID != nil {
		view.ConsultationID = *row.ConsultationID
	}

	if constants.IsActiveQueueStatus(row.Status) {

		view.TotalInQueue = len(active)

		for index, item := range active {
			if item.ID == row.ID {
				view.Position = index + 1
			}
		}
	}

	word := queueMediumWord(row.WaitingType)

	switch row.Status {

	case constants.QueueStatusWaiting:
		view.NextAction = "WAIT"
		view.Message = fmt.Sprintf("You are in the waiting queue for Astrologer %s. Your current position is %d.", name, view.Position)

		if row.WaitExpiresAt != nil {
			view.SecondsRemaining = secondsUntil(*row.WaitExpiresAt, now)
		}

	case constants.QueueStatusNotified:
		view.NextAction = "CONNECT"
		view.Message = fmt.Sprintf("Astrologer %s is now available for %s. Please try to reconnect.", name, word)

		if row.NotifyExpiresAt != nil {
			view.SecondsRemaining = secondsUntil(*row.NotifyExpiresAt, now)
		}

	case constants.QueueStatusConnecting:

		if row.ConsultationID != nil {
			view.NextAction = "WAIT"
			view.Message = fmt.Sprintf("Your %s request has been sent to Astrologer %s.", word, name)
		} else {
			view.NextAction = "START_CONSULTATION"
			view.Message = fmt.Sprintf("Connecting to Astrologer %s. Start the %s now.", name, word)

			if row.ConnectExpiresAt != nil {
				view.SecondsRemaining = secondsUntil(*row.ConnectExpiresAt, now)
			}
		}

	case constants.QueueStatusConnected:
		view.Message = fmt.Sprintf("You are connected with Astrologer %s.", name)

	case constants.QueueStatusExpired:
		view.Message = fmt.Sprintf("Your waiting request for Astrologer %s has expired. Please try connecting again.", name)

	case constants.QueueStatusCancelled:
		view.Message = fmt.Sprintf("You left the waiting queue for Astrologer %s.", name)

	case constants.QueueStatusSkipped:
		view.Message = fmt.Sprintf("You were skipped in the waiting queue for Astrologer %s because you were not eligible to connect.", name)

	case constants.QueueStatusRejected:
		view.Message = fmt.Sprintf("Your waiting request for Astrologer %s could not be connected.", name)

	case constants.QueueStatusCompleted:
		view.Message = fmt.Sprintf("Your waiting request for Astrologer %s is complete.", name)
	}

	return view, nil
}
