@@ -173,3 +173,20 @@ CREATE TABLE IF NOT EXISTS customer_link_tokens (
|
||||
created_at INTEGER NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS customer_link_tokens_exp_idx ON customer_link_tokens(expires_at);
|
||||
|
||||
-- Durable outbox for hosted-worker positive-tip rewards. A "positive tip" is
|
||||
-- a signed guess that improves that worker's personal best score. Delivery to
|
||||
-- Customer Service is retried and the event id is idempotent there.
|
||||
CREATE TABLE IF NOT EXISTS hosted_credit_events (
|
||||
event_id TEXT PRIMARY KEY,
|
||||
worker_client_id TEXT NOT NULL REFERENCES clients(id) ON DELETE CASCADE,
|
||||
reward_client_id TEXT NOT NULL REFERENCES clients(id) ON DELETE CASCADE,
|
||||
task_id TEXT NOT NULL REFERENCES tasks(id) ON DELETE CASCADE,
|
||||
seq INTEGER NOT NULL,
|
||||
score REAL NOT NULL,
|
||||
created_at INTEGER NOT NULL,
|
||||
delivered_at INTEGER,
|
||||
attempts INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT NOT NULL DEFAULT ''
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS hosted_credit_events_pending_idx ON hosted_credit_events(delivered_at,created_at);
|
||||
|
||||
@@ -1784,3 +1784,72 @@ func (s *Store) ConsumeCustomerLinkToken(ctx context.Context, tokenHash string)
|
||||
}
|
||||
return clientID, nil
|
||||
}
|
||||
|
||||
// HostedCreditEvent is a durable outbox item for customer-service engagement
|
||||
// credits. Keeping it in the game database means a temporary Customer Service
|
||||
// outage does not silently lose a positive-tip reward.
|
||||
type HostedCreditEvent struct {
|
||||
EventID string `json:"event_id"`
|
||||
WorkerClientID string `json:"worker_client_id"`
|
||||
RewardClientID string `json:"reward_client_id"`
|
||||
TaskID string `json:"task_id"`
|
||||
Seq int64 `json:"seq"`
|
||||
Score float64 `json:"score"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
Attempts int `json:"attempts"`
|
||||
}
|
||||
|
||||
func (s *Store) RecordHostedPositiveCreditEvent(ctx context.Context, taskID, workerClientID, rewardClientID string, seq int64, score float64) (string, bool, error) {
|
||||
workerClientID = strings.TrimSpace(workerClientID)
|
||||
rewardClientID = strings.TrimSpace(rewardClientID)
|
||||
if workerClientID == "" || rewardClientID == "" || rewardClientID == workerClientID {
|
||||
return "", false, nil
|
||||
}
|
||||
// The identity delegation itself is the server-side proof that this is a
|
||||
// hosted/delegated worker rather than an ordinary browser/CLI identity.
|
||||
var owner string
|
||||
if err := s.DB.QueryRowContext(ctx, `SELECT owner_client_id FROM identity_delegations WHERE worker_client_id=?`, workerClientID).Scan(&owner); err != nil || owner != rewardClientID {
|
||||
return "", false, nil
|
||||
}
|
||||
eventID := fmt.Sprintf("%s:%s:%d", taskID, workerClientID, seq)
|
||||
res, err := s.DB.ExecContext(ctx, `INSERT OR IGNORE INTO hosted_credit_events(event_id,worker_client_id,reward_client_id,task_id,seq,score,created_at) VALUES(?,?,?,?,?,?,?)`, eventID, workerClientID, rewardClientID, taskID, seq, score, time.Now().UTC().UnixMilli())
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
n, err := res.RowsAffected()
|
||||
return eventID, n == 1, err
|
||||
}
|
||||
|
||||
func (s *Store) PendingHostedCreditEvents(ctx context.Context, limit int) ([]HostedCreditEvent, error) {
|
||||
if limit < 1 || limit > 200 {
|
||||
limit = 50
|
||||
}
|
||||
rows, err := s.DB.QueryContext(ctx, `SELECT event_id,worker_client_id,reward_client_id,task_id,seq,score,created_at,attempts FROM hosted_credit_events WHERE delivered_at IS NULL ORDER BY created_at LIMIT ?`, limit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []HostedCreditEvent
|
||||
for rows.Next() {
|
||||
var x HostedCreditEvent
|
||||
var created int64
|
||||
if err := rows.Scan(&x.EventID, &x.WorkerClientID, &x.RewardClientID, &x.TaskID, &x.Seq, &x.Score, &created, &x.Attempts); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
x.CreatedAt = time.UnixMilli(created).UTC()
|
||||
out = append(out, x)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (s *Store) MarkHostedCreditEventDelivered(ctx context.Context, eventID string) error {
|
||||
_, err := s.DB.ExecContext(ctx, `UPDATE hosted_credit_events SET delivered_at=?,last_error='' WHERE event_id=?`, time.Now().UTC().UnixMilli(), eventID)
|
||||
return err
|
||||
}
|
||||
func (s *Store) MarkHostedCreditEventAttempt(ctx context.Context, eventID, lastErr string) error {
|
||||
if len(lastErr) > 1000 {
|
||||
lastErr = lastErr[:1000]
|
||||
}
|
||||
_, err := s.DB.ExecContext(ctx, `UPDATE hosted_credit_events SET attempts=attempts+1,last_error=? WHERE event_id=?`, lastErr, eventID)
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user