@@ -0,0 +1,62 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"neuralhunt/internal/data"
|
||||
)
|
||||
|
||||
func (s *Server) dispatchHostedCreditEvents(ctx context.Context) {
|
||||
if strings.TrimSpace(s.customerServiceInternalURL) == "" || strings.TrimSpace(s.internalServiceSecret) == "" {
|
||||
return
|
||||
}
|
||||
events, err := s.store.PendingHostedCreditEvents(ctx, 25)
|
||||
if err != nil {
|
||||
log.Printf("hosted credit outbox: %v", err)
|
||||
return
|
||||
}
|
||||
for _, ev := range events {
|
||||
if ev.Attempts > 0 && ev.Attempts%20 == 0 {
|
||||
log.Printf("hosted credit event %s still pending after %d attempts", ev.EventID, ev.Attempts)
|
||||
}
|
||||
if err := s.deliverHostedCreditEvent(ctx, ev); err != nil {
|
||||
_ = s.store.MarkHostedCreditEventAttempt(ctx, ev.EventID, err.Error())
|
||||
continue
|
||||
}
|
||||
_ = s.store.MarkHostedCreditEventDelivered(ctx, ev.EventID)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Server) deliverHostedCreditEvent(ctx context.Context, ev data.HostedCreditEvent) error {
|
||||
body, _ := json.Marshal(map[string]any{
|
||||
"event_id": ev.EventID,
|
||||
"worker_client_id": ev.WorkerClientID,
|
||||
"reward_client_id": ev.RewardClientID,
|
||||
"task_id": ev.TaskID,
|
||||
"seq": ev.Seq,
|
||||
"score": ev.Score,
|
||||
})
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, s.customerServiceInternalURL+"/internal/game/positive-tip", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Authorization", "Bearer "+s.internalServiceSecret)
|
||||
resp, err := s.customerServiceHTTP.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
||||
if resp.StatusCode/100 != 2 {
|
||||
return fmt.Errorf("customer service HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -47,6 +47,8 @@ type Server struct {
|
||||
lottery *guessLottery
|
||||
adminUser, adminPass, staticDir, artifactDir string
|
||||
internalServiceSecret string
|
||||
customerServiceInternalURL string
|
||||
customerServiceHTTP *http.Client
|
||||
upgrader websocket.Upgrader
|
||||
wsAllowedOrigins map[string]struct{}
|
||||
maxUserWS, maxLeaderboardWS int64
|
||||
@@ -65,14 +67,16 @@ func New(store *data.Store, a *auth.Manager, sm *settings.Manager, hub *wsx.Hub,
|
||||
lottery: newGuessLottery(func(d beaconDrawAudit) {
|
||||
_ = store.RecordBeaconDraw(context.Background(), d.TaskID, d.WindowEnd, d.BeaconID, d.BeaconRound, d.Randomness, d.Signature, d.BoostedPath, d.Tickets, d.Selected)
|
||||
}),
|
||||
adminUser: env("ADMIN_USER", "admin"),
|
||||
adminPass: env("ADMIN_PASSWORD", "change-me"),
|
||||
staticDir: env("STATIC_DIR", ""),
|
||||
artifactDir: artifactDir,
|
||||
internalServiceSecret: strings.TrimSpace(os.Getenv("CUSTOMER_SERVICE_SHARED_SECRET")),
|
||||
wsAllowedOrigins: parseOriginAllowlist(os.Getenv("WS_ALLOWED_ORIGINS")),
|
||||
maxUserWS: int64(envIntServer("WS_MAX_USER_CONNECTIONS", 5000)),
|
||||
maxLeaderboardWS: int64(envIntServer("WS_MAX_LEADERBOARD_CONNECTIONS", 500)),
|
||||
adminUser: env("ADMIN_USER", "admin"),
|
||||
adminPass: env("ADMIN_PASSWORD", "change-me"),
|
||||
staticDir: env("STATIC_DIR", ""),
|
||||
artifactDir: artifactDir,
|
||||
internalServiceSecret: strings.TrimSpace(os.Getenv("CUSTOMER_SERVICE_SHARED_SECRET")),
|
||||
customerServiceInternalURL: strings.TrimRight(strings.TrimSpace(os.Getenv("CUSTOMER_SERVICE_INTERNAL_URL")), "/"),
|
||||
customerServiceHTTP: &http.Client{Timeout: 5 * time.Second},
|
||||
wsAllowedOrigins: parseOriginAllowlist(os.Getenv("WS_ALLOWED_ORIGINS")),
|
||||
maxUserWS: int64(envIntServer("WS_MAX_USER_CONNECTIONS", 5000)),
|
||||
maxLeaderboardWS: int64(envIntServer("WS_MAX_LEADERBOARD_CONNECTIONS", 500)),
|
||||
}
|
||||
s.upgrader = websocket.Upgrader{CheckOrigin: s.checkWSOrigin, Subprotocols: []string{"neuralhunt.v1"}}
|
||||
return s
|
||||
@@ -851,6 +855,17 @@ func (s *Server) guess(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
s.runtime.MarkSQLiteWrite()
|
||||
// A positive hosted-worker tip means a genuine personal-best improvement,
|
||||
// not merely a submitted or lottery-selected guess. Record it durably and
|
||||
// let the scheduler deliver the idempotent credit event to Customer Service.
|
||||
if accepted.Improved {
|
||||
rewardOwner := s.store.RewardOwnerForWorker(r.Context(), c.ClientID)
|
||||
if _, created, err := s.store.RecordHostedPositiveCreditEvent(r.Context(), id, c.ClientID, rewardOwner, in.Seq, accepted.State.BestScore); err != nil {
|
||||
log.Printf("hosted positive-tip outbox: %v", err)
|
||||
} else if created {
|
||||
s.runtime.MarkSQLiteWrite()
|
||||
}
|
||||
}
|
||||
p.Score = round(p.Score, s.settings.Get().PublicScorePrecision)
|
||||
s.hub.PublishPoint(id, c.ClientID, p)
|
||||
}
|
||||
@@ -1825,6 +1840,7 @@ func (s *Server) Scheduler(ctx context.Context) {
|
||||
s.runDueActions(ctx)
|
||||
maintenance++
|
||||
if maintenance%5 == 0 {
|
||||
s.dispatchHostedCreditEvents(ctx)
|
||||
if err := s.store.EnsureActiveTasks(ctx, s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits); err != nil {
|
||||
log.Printf("scheduler: %v", err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user