package httpapi import ( "context" "errors" "fmt" "net/http" "strings" "time" "calllinesystem/server/internal/domain" "calllinesystem/server/internal/model" "calllinesystem/server/internal/security" "gorm.io/gorm" ) func (s *Server) publicStatus(w http.ResponseWriter, r *http.Request) { token := strings.TrimSpace(r.PathValue("token")) if (len(token) < 40 || len(token) > 128) && !(s.config.Environment == "development" && strings.HasPrefix(token, "demo-visitor-") && len(token) <= 64) { writeError(w, &apiError{Status: http.StatusNotFound, Code: "STATUS_NOT_FOUND", Message: "排队状态不存在或链接已失效"}) return } var ticket model.QueueTicket err := s.db.WithContext(r.Context()).Where("public_token_hash = ?", security.HashToken(token)).First(&ticket).Error if err != nil { writeError(w, mapNotFound(err, "STATUS_NOT_FOUND", "排队状态不存在或链接已失效")) return } view, err := s.publicStatusView(r.Context(), ticket) if err != nil { writeError(w, err) return } writeJSON(w, http.StatusOK, view) } func (s *Server) publicProjects(w http.ResponseWriter, r *http.Request) { var projects []model.Project if err := s.db.WithContext(r.Context()). Where("status = ?", model.ProjectRunning). Order("name ASC").Find(&projects).Error; err != nil { writeError(w, err) return } views := make([]map[string]any, 0, len(projects)) for _, project := range projects { views = append(views, map[string]any{ "id": project.ID, "name": project.Name, "status": project.Status, "visitor_notice": project.VisitorNotice, "min_party_size": project.MinPartySize, "max_party_size": project.MaxPartySize, }) } writeJSON(w, http.StatusOK, map[string]any{"projects": views}) } func (s *Server) publicCreateTicket(w http.ResponseWriter, r *http.Request) { projectID := r.PathValue("id") if err := validateUUID(projectID); err != nil { writeError(w, err) return } if s.publicTicketLimiter != nil { if allowed, retry := s.publicTicketLimiter.allow("ip:" + publicQueryClientKey(r)); !allowed { writePublicTicketRateLimit(w, retry) return } } var actor model.User if err := s.db.WithContext(r.Context()).Where("username = ?", model.PublicVisitorUsername).First(&actor).Error; err != nil { writeError(w, err) return } s.createTicketForActor(w, r, actor.ID, true) } type publicPhoneQueryRequest struct { Phone string `json:"phone"` } // publicStatusByPhone is a temporary operational-test flow. It must be // replaced by OTP or an external identity provider before production use. func (s *Server) publicStatusByPhone(w http.ResponseWriter, r *http.Request) { if s.config.Environment == "production" { writeError(w, &apiError{Status: http.StatusNotFound, Code: "STATUS_NOT_FOUND", Message: "排队状态不存在或链接已失效"}) return } s.statusByPhone(w, r, true) } func (s *Server) statusByPhone(w http.ResponseWriter, r *http.Request, limitByClientIP bool) { if limitByClientIP && s.publicQueryLimiter != nil { if allowed, retry := s.publicQueryLimiter.allow("ip:" + publicQueryClientKey(r)); !allowed { writePublicQueryRateLimit(w, retry) return } } var input publicPhoneQueryRequest if err := decodeJSON(r, &input); err != nil { writeError(w, err) return } phone, err := security.NormalizePhone(input.Phone) if err != nil { writeError(w, &apiError{Status: http.StatusUnprocessableEntity, Code: "INVALID_PHONE", Message: "请输入有效的手机号"}) return } phoneDigest := s.cipher.Digest(phone) if s.publicQueryLimiter != nil { if allowed, retry := s.publicQueryLimiter.allow("phone:" + phoneDigest); !allowed { writePublicQueryRateLimit(w, retry) return } } var tickets []model.QueueTicket activeSessionStatuses := []string{"RUNNING", "PAUSED"} if err := s.db.WithContext(r.Context()).Model(&model.QueueTicket{}). Joins("JOIN projects ON projects.id = queue_tickets.project_id"). Joins("JOIN queue_sessions ON queue_sessions.id = queue_tickets.queue_session_id AND queue_sessions.project_id = queue_tickets.project_id"). Where("queue_tickets.phone_hmac = ? AND queue_tickets.status IN ? AND queue_sessions.status IN ? AND projects.status IN ?", phoneDigest, []string{model.TicketWaiting, model.TicketCalled, model.TicketArrived}, activeSessionStatuses, []string{model.ProjectRunning, model.ProjectPaused}). Where("queue_sessions.business_date = (? AT TIME ZONE projects.timezone)::date", s.now()). Order("projects.name ASC, queue_tickets.ticket_number ASC"). Find(&tickets).Error; err != nil { writeError(w, err) return } views := make([]map[string]any, 0, len(tickets)) for _, ticket := range tickets { view, err := s.publicStatusView(r.Context(), ticket) if err != nil { writeError(w, err) return } views = append(views, view) } writeJSON(w, http.StatusOK, map[string]any{"tickets": views}) } func publicQueryClientKey(r *http.Request) string { if ip := remoteIP(r); ip != nil { return *ip } return "unknown" } func writePublicQueryRateLimit(w http.ResponseWriter, retry time.Duration) { seconds := int((retry + time.Second - 1) / time.Second) if seconds < 1 { seconds = 1 } w.Header().Set("Retry-After", fmt.Sprintf("%d", seconds)) writeError(w, &apiError{Status: http.StatusTooManyRequests, Code: "PUBLIC_QUERY_RATE_LIMITED", Message: "查询次数过多,请稍后再试"}) } func writePublicTicketRateLimit(w http.ResponseWriter, retry time.Duration) { seconds := int((retry + time.Second - 1) / time.Second) if seconds < 1 { seconds = 1 } w.Header().Set("Retry-After", fmt.Sprintf("%d", seconds)) writeError(w, &apiError{Status: http.StatusTooManyRequests, Code: "PUBLIC_TICKET_RATE_LIMITED", Message: "取号次数过多,请稍后再试"}) } func (s *Server) publicStatusView(ctx context.Context, ticket model.QueueTicket) (map[string]any, error) { var project model.Project var session model.QueueSession if err := s.db.WithContext(ctx).First(&project, "id = ?", ticket.ProjectID).Error; err != nil { return nil, err } if err := s.db.WithContext(ctx).First(&session, "id = ? AND project_id = ?", ticket.QueueSessionID, ticket.ProjectID).Error; err != nil { return nil, err } experiencedPeople, err := s.displayedExperiencedPeople(ctx, project, &session) if err != nil { return nil, err } phoneSuffix, err := s.ticketPhoneLast4(ticket) if err != nil { return nil, err } ticketsAhead := 0 peopleAhead := 0 if ticket.Status == model.TicketWaiting { var totals ticketPeopleTotals if err := s.db.WithContext(ctx).Model(&model.QueueTicket{}). Select("count(*) AS ticket_count, COALESCE(sum(party_size), 0) AS people_count"). Where("project_id = ? AND queue_session_id = ? AND status = ? AND ticket_number < ?", ticket.ProjectID, ticket.QueueSessionID, model.TicketWaiting, ticket.TicketNumber).Scan(&totals).Error; err != nil { return nil, err } ticketsAhead = int(totals.TicketCount) peopleAhead = int(totals.PeopleCount) } var latestCalledTicket struct { DisplayNumber string `gorm:"column:display_number"` } latestCalledNumber := any(nil) if err := s.db.WithContext(ctx).Model(&model.QueueTicket{}). Select("display_number"). Where("project_id = ? AND queue_session_id = ? AND called_at IS NOT NULL", ticket.ProjectID, ticket.QueueSessionID). Order("called_at DESC, ticket_number DESC").First(&latestCalledTicket).Error; err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { return nil, err } else if latestCalledTicket.DisplayNumber != "" { latestCalledNumber = latestCalledTicket.DisplayNumber } eta, err := domain.CalculateETA(domain.ETAInput{ PeopleAhead: peopleAhead, IntervalPerPerson: time.Duration(project.ETAIntervalSeconds) * time.Second, Running: project.Status == model.ProjectRunning && session.Status == "RUNNING" && ticket.Status == model.TicketWaiting, }) if err != nil { return nil, err } if ticket.Status != model.TicketWaiting { eta = domain.ETAResult{Available: false, Reason: "ticket_not_waiting"} } return map[string]any{ "ticket_number": ticket.DisplayNumber, "display_number": ticket.DisplayNumber, "project_name": project.Name, "status": ticket.Status, "party_size": ticket.PartySize, "phone_last4": phoneSuffix, "estimated_wait": eta, "visitor_notice": project.VisitorNotice, "experienced_people": experiencedPeople, "last_updated_at": s.now(), "called_at": ticket.CalledAt, "ticket": map[string]any{ "display_number": ticket.DisplayNumber, "status": ticket.Status, "party_size": ticket.PartySize, "joined_at": ticket.JoinedAt, "called_at": ticket.CalledAt, "arrived_at": ticket.ArrivedAt, "completed_at": ticket.CompletedAt, "missed_at": ticket.MissedAt, }, "project": map[string]any{"id": project.ID, "name": project.Name, "status": project.Status}, "tickets_ahead": ticketsAhead, "people_ahead": peopleAhead, "queue_position": queuePosition(ticket.Status, ticketsAhead), "latest_called_number": latestCalledNumber, "eta": eta, "revision": session.Revision, "server_time": s.now(), }, nil } func queuePosition(status string, ticketsAhead int) any { if status != model.TicketWaiting { return nil } return ticketsAhead + 1 } // displayTicketDTO is the complete public ticket shape for a display. Keeping // this separate from staff ticketView makes it impossible to accidentally // serialize honorifics, encrypted contact fields, internal IDs or timestamps. type displayTicketDTO struct { TicketNumber string `json:"ticket_number"` DisplayNumber string `json:"display_number"` Status string `json:"status"` PartySize int `json:"party_size"` } type displayBatchDTO struct { BatchNumber int `json:"batch_number"` Sequence int `json:"sequence"` Status string `json:"status"` CallMode string `json:"call_mode"` RequestedCount int `json:"requested_count"` TicketCount int `json:"ticket_count"` PeopleCount int `json:"people_count"` CalledAt time.Time `json:"called_at"` Tickets []displayTicketDTO `json:"tickets"` } type displayProjectDTO struct { ID string `json:"id"` Name string `json:"name"` Status string `json:"status"` } type displaySnapshotDTO struct { ProjectName string `json:"project_name"` Status string `json:"status"` Project displayProjectDTO `json:"project"` Revision int64 `json:"revision"` WaitingCount int64 `json:"waiting_count"` WaitingTicketCount int64 `json:"waiting_ticket_count"` WaitingPeopleCount int64 `json:"waiting_people_count"` ExperiencedPeople int64 `json:"experienced_people"` CurrentBatch *displayBatchDTO `json:"current_batch"` RecentBatches []displayBatchDTO `json:"recent_batches"` EstimatedWait domain.ETAResult `json:"estimated_wait"` ServerTime time.Time `json:"server_time"` LastUpdatedAt time.Time `json:"last_updated_at"` } func normalizeDisplayProjectCode(value string) (string, bool) { code := strings.ToUpper(strings.TrimSpace(value)) if !projectCodePattern.MatchString(code) { return "", false } return code, true } func (s *Server) displaySnapshot(w http.ResponseWriter, r *http.Request) { identifier := strings.TrimSpace(r.PathValue("token")) var project model.Project var err error if code, ok := normalizeDisplayProjectCode(identifier); ok { err = s.db.WithContext(r.Context()).Where("code = ?", code).First(&project).Error } else { if len(identifier) < 40 || len(identifier) > 128 { writeError(w, &apiError{Status: http.StatusNotFound, Code: "DISPLAY_NOT_FOUND", Message: "公示屏绑定不存在"}) return } err = s.db.WithContext(r.Context()).Where("display_token_hash = ?", security.HashToken(identifier)).First(&project).Error } if err != nil { writeError(w, mapNotFound(err, "DISPLAY_NOT_FOUND", "公示屏绑定不存在")) return } session, err := s.currentQueueSession(s.db.WithContext(r.Context()), project) if errors.Is(err, gorm.ErrRecordNotFound) { experiencedPeople, metricErr := s.displayedExperiencedPeople(r.Context(), project, nil) if metricErr != nil { writeError(w, metricErr) return } writeJSON(w, http.StatusOK, displaySnapshotDTO{ ProjectName: project.Name, Status: project.Status, Project: displayProjectDTO{ID: project.ID, Name: project.Name, Status: project.Status}, ExperiencedPeople: experiencedPeople, RecentBatches: []displayBatchDTO{}, EstimatedWait: domain.ETAResult{Available: false, Reason: "queue_not_running"}, ServerTime: s.now(), LastUpdatedAt: project.UpdatedAt, }) return } if err != nil { writeError(w, err) return } waiting, err := queueTotals(s.db.WithContext(r.Context()), project.ID, session.ID, model.TicketWaiting) if err != nil { writeError(w, err) return } batch, err := s.currentDisplayBatch(r.Context(), project.ID, session.ID) if err != nil { writeError(w, err) return } recentBatches, err := s.recentDisplayBatches(r, project.ID, session.ID) if err != nil { writeError(w, err) return } var lastWaiting model.QueueTicket lastWaitingErr := s.db.WithContext(r.Context()).Select("id", "party_size"). Where("project_id = ? AND queue_session_id = ? AND status = ?", project.ID, session.ID, model.TicketWaiting). Order("ticket_number DESC").First(&lastWaiting).Error if lastWaitingErr != nil && !errors.Is(lastWaitingErr, gorm.ErrRecordNotFound) { writeError(w, lastWaitingErr) return } peopleAhead := int(waiting.PeopleCount) if lastWaitingErr == nil { peopleAhead = max(0, peopleAhead-lastWaiting.PartySize) } estimatedWait, err := domain.CalculateETA(domain.ETAInput{ PeopleAhead: peopleAhead, IntervalPerPerson: time.Duration(project.ETAIntervalSeconds) * time.Second, Running: project.Status == model.ProjectRunning && session.Status == "RUNNING" && waiting.TicketCount > 0, }) if err != nil { writeError(w, err) return } experiencedPeople, err := s.displayedExperiencedPeople(r.Context(), project, &session) if err != nil { writeError(w, err) return } lastUpdated := project.UpdatedAt if session.UpdatedAt.After(lastUpdated) { lastUpdated = session.UpdatedAt } writeJSON(w, http.StatusOK, displaySnapshotDTO{ ProjectName: project.Name, Status: project.Status, Project: displayProjectDTO{ID: project.ID, Name: project.Name, Status: project.Status}, Revision: session.Revision, WaitingCount: waiting.TicketCount, WaitingTicketCount: waiting.TicketCount, WaitingPeopleCount: waiting.PeopleCount, ExperiencedPeople: experiencedPeople, CurrentBatch: batch, RecentBatches: recentBatches, EstimatedWait: estimatedWait, ServerTime: s.now(), LastUpdatedAt: lastUpdated, }) } func (s *Server) currentDisplayBatch(ctx context.Context, projectID, sessionID string) (*displayBatchDTO, error) { var batch model.CallBatch err := s.db.WithContext(ctx).Where("project_id = ? AND queue_session_id = ? AND status = 'CALLED'", projectID, sessionID). Order("batch_sequence DESC").First(&batch).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, err } tickets, err := displayBatchTickets(s.db.WithContext(ctx), batch.ID, projectID, true) if err != nil { return nil, err } view := newDisplayBatchDTO(batch, tickets) return &view, nil } func (s *Server) recentDisplayBatches(r *http.Request, projectID, sessionID string) ([]displayBatchDTO, error) { var batches []model.CallBatch if err := s.db.WithContext(r.Context()).Where("project_id = ? AND queue_session_id = ? AND status = 'COMPLETED'", projectID, sessionID). Order("batch_sequence DESC").Limit(5).Find(&batches).Error; err != nil { return nil, err } views := make([]displayBatchDTO, 0, len(batches)) for _, batch := range batches { tickets, err := displayBatchTickets(s.db.WithContext(r.Context()), batch.ID, projectID, false) if err != nil { return nil, err } views = append(views, newDisplayBatchDTO(batch, tickets)) } return views, nil } func displayBatchTickets(db *gorm.DB, batchID, projectID string, calledOnly bool) ([]displayTicketDTO, error) { query := db.Table("queue_tickets"). Select("queue_tickets.display_number, queue_tickets.status, queue_tickets.party_size"). Joins("JOIN call_batch_tickets cbt ON cbt.ticket_id = queue_tickets.id AND cbt.project_id = queue_tickets.project_id"). Where("cbt.call_batch_id = ? AND cbt.project_id = ?", batchID, projectID). Order("cbt.position ASC") if calledOnly { query = query.Where("queue_tickets.status = ?", model.TicketCalled) } type row struct { DisplayNumber string Status string PartySize int } var rows []row if err := query.Scan(&rows).Error; err != nil { return nil, err } tickets := make([]displayTicketDTO, 0, len(rows)) for _, item := range rows { tickets = append(tickets, displayTicketDTO{ TicketNumber: item.DisplayNumber, DisplayNumber: item.DisplayNumber, Status: item.Status, PartySize: item.PartySize, }) } return tickets, nil } func newDisplayBatchDTO(batch model.CallBatch, tickets []displayTicketDTO) displayBatchDTO { return displayBatchDTO{ BatchNumber: batch.BatchSequence, Sequence: batch.BatchSequence, Status: batch.Status, CallMode: batch.CallMode, RequestedCount: batch.RequestedCount, TicketCount: batch.TicketCount, PeopleCount: batch.PeopleCount, CalledAt: batch.CalledAt, Tickets: tickets, } } func publicDisplayProjectView(project map[string]any) map[string]any { return map[string]any{ "id": project["id"], "name": project["name"], "status": project["status"], "waiting_count": project["waiting_count"], "waiting_ticket_count": project["waiting_ticket_count"], "waiting_people_count": project["waiting_people_count"], "issued_ticket_count": project["issued_ticket_count"], "latest_ticket_number": project["latest_ticket_number"], "experienced_people": project["experienced_people"], "current_batch": project["current_batch"], "estimated_wait": project["estimated_wait"], "last_updated_at": project["last_updated_at"], } } func (s *Server) displayOverview(w http.ResponseWriter, r *http.Request) { var projects []model.Project if err := s.db.WithContext(r.Context()).Order("name ASC").Find(&projects).Error; err != nil { writeError(w, err) return } views := make([]map[string]any, 0, len(projects)) for _, project := range projects { projection, _, _, _, _, err := s.adminProjectProjection(r.Context(), project) if err != nil { writeError(w, err) return } views = append(views, publicDisplayProjectView(projection)) } writeJSON(w, http.StatusOK, map[string]any{"projects": views, "server_time": s.now()}) } func (s *Server) events(w http.ResponseWriter, r *http.Request) { projectID := strings.TrimSpace(r.URL.Query().Get("project_id")) if err := validateUUID(projectID); err != nil { writeError(w, err) return } if err := s.authorizeProject(r.Context(), projectID); err != nil { writeError(w, err) return } flusher, ok := w.(http.Flusher) if !ok { writeError(w, fmt.Errorf("streaming unsupported")) return } w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache, no-transform") w.Header().Set("Connection", "keep-alive") w.Header().Set("X-Accel-Buffering", "no") _, _ = fmt.Fprint(w, "retry: 5000\n\n") flusher.Flush() events, unsubscribe := s.hub.subscribe(projectID) defer unsubscribe() keepAlive := time.NewTicker(15 * time.Second) defer keepAlive.Stop() for { select { case <-r.Context().Done(): return case body := <-events: _, _ = fmt.Fprintf(w, "event: queue\ndata: %s\n\n", body) flusher.Flush() case <-keepAlive.C: _, _ = fmt.Fprint(w, ": keep-alive\n\n") flusher.Flush() } } }