Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 20 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,22 @@
# WalletOps Companion

Personal ops console for simulated wallet events: signed webhook ingest, alert rules, a retrying worker, and schema-checked AI summaries. Flutter client + Go API + Postgres.
Personal ops console for **simulated** wallet events. Partners post signed webhooks; Postgres stores them idempotently; a Go worker claims and processes them against alert rules; a Flutter client lists events and can request a schema-checked summary.

No real custody, keys, or on-chain sends.

## Why the backend matters

The interview-friendly core is reliability around ingest and jobs:

| Concern | Behavior |
|---------|----------|
| Forged webhooks | HMAC-SHA256 over the raw body (`X-Signature: sha256=<hex>`) |
| Duplicate delivery | Unique `(user_id, idempotency_key)`; replay returns the same event (`200`) |
| Concurrent workers | `FOR UPDATE SKIP LOCKED` claim — safe under two pollers |
| Crash mid-process | `claimed_at` lease; stale `processing` rows become claimable again |
| Retries | Failed attempts back off (capped at 60s) up to 5 tries |

Open `api/internal/webhook/handler_test.go` (`replay`) and `api/internal/worker/worker_test.go` (`TestConcurrentClaimNoDoubleProcess`, `TestReclaimExpiredProcessingLease`).

## Architecture

Expand Down Expand Up @@ -38,6 +54,7 @@ Copy `.env.example` → `.env` (compose already injects defaults for local):
# 1) API + Postgres
docker compose up --build -d
curl -s http://127.0.0.1:8080/v1/health
# expect status=ok, worker ticks, queue.by_status

# 2) Seed user + two signed events (maps user_ref=demo-user-1)
./scripts/seed_webhooks.sh
Expand All @@ -52,7 +69,7 @@ flutter run --dart-define=API_BASE=http://127.0.0.1:8080
# Android emulator: API_BASE=http://10.0.2.2:8080
```

In the app: sign in as `demo-user-1@walletops.local` / `ops-secret-1` → Events (status → processed) → open an event → **Explain** (mock AI).
In the app: sign in as `demo-user-1@walletops.local` / `ops-secret-1` → Events (status → processed) → open an event → **Explain** (mock AI by default).

## Tests

Expand All @@ -66,4 +83,4 @@ cd api && DATABASE_URL='postgres://walletops:walletops@localhost:5432/walletops?
cd mobile && flutter pub get && flutter analyze && flutter test --dart-define=API_BASE=http://127.0.0.1:8080
```

Module path under `api/` is a placeholder — change before publishing.
Module path: `github.com/omid/walletops/api`.
3 changes: 3 additions & 0 deletions api/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,9 @@ func main() {
WorkerSnapshot: func() any {
return wrk.Stats.Snapshot()
},
QueueSnapshot: func(ctx context.Context) (any, error) {
return eventStore.QueueStats(ctx)
},
})


Expand Down
77 changes: 72 additions & 5 deletions api/internal/events/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -230,11 +230,23 @@ func (s *Store) GetByID(ctx context.Context, id string) (Event, error) {
return e, nil
}

// ClaimLease is how long a processing row may sit before another worker can reclaim it.
const ClaimLease = 45 * time.Second

func (s *Store) ClaimNext(ctx context.Context, maxAttempts int) (Event, error) {
return s.ClaimNextWithLease(ctx, maxAttempts, ClaimLease)
}

func (s *Store) ClaimNextWithLease(ctx context.Context, maxAttempts int, lease time.Duration) (Event, error) {
if lease <= 0 {
lease = ClaimLease
}
secs := int(lease.Seconds())
var e Event
err := s.pool.QueryRow(ctx, `
UPDATE events
SET status = 'processing'
SET status = 'processing',
claimed_at = now()
WHERE id = (
SELECT id FROM events
WHERE status = 'pending'
Expand All @@ -243,13 +255,18 @@ func (s *Store) ClaimNext(ctx context.Context, maxAttempts int) (Event, error) {
AND attempt_count < $1
AND COALESCE(processed_at, received_at) <= now() - make_interval(secs => LEAST(60, GREATEST(2, attempt_count * 2)))
)
OR (
status = 'processing'
AND claimed_at IS NOT NULL
AND claimed_at <= now() - make_interval(secs => $2)
)
ORDER BY received_at ASC
LIMIT 1
FOR UPDATE SKIP LOCKED
)
RETURNING id::text, user_id::text, idempotency_key, type, payload, status,
attempt_count, last_error, matched_rule_id::text, received_at, processed_at
`, maxAttempts).Scan(
`, maxAttempts, secs).Scan(
&e.ID, &e.UserID, &e.IdempotencyKey, &e.Type, &e.Payload, &e.Status,
&e.AttemptCount, &e.LastError, &e.MatchedRuleID, &e.ReceivedAt, &e.ProcessedAt,
)
Expand All @@ -266,7 +283,8 @@ func (s *Store) ClaimByID(ctx context.Context, id string) (Event, error) {
var e Event
err := s.pool.QueryRow(ctx, `
UPDATE events
SET status = 'processing'
SET status = 'processing',
claimed_at = now()
WHERE id = $1 AND status IN ('pending', 'failed')
RETURNING id::text, user_id::text, idempotency_key, type, payload, status,
attempt_count, last_error, matched_rule_id::text, received_at, processed_at
Expand All @@ -289,7 +307,8 @@ func (s *Store) MarkProcessed(ctx context.Context, id string, matchedRuleID *str
SET status = 'processed',
matched_rule_id = $2,
processed_at = now(),
last_error = NULL
last_error = NULL,
claimed_at = NULL
WHERE id = $1 AND status = 'processing'
`, id, matchedRuleID)
if err != nil {
Expand All @@ -308,7 +327,8 @@ func (s *Store) MarkAttemptFailed(ctx context.Context, id, lastError string) (Ev
SET attempt_count = attempt_count + 1,
last_error = $2,
status = 'failed',
processed_at = now()
processed_at = now(),
claimed_at = NULL
WHERE id = $1 AND status = 'processing'
RETURNING id::text, user_id::text, idempotency_key, type, payload, status,
attempt_count, last_error, matched_rule_id::text, received_at, processed_at
Expand All @@ -324,3 +344,50 @@ func (s *Store) MarkAttemptFailed(ctx context.Context, id, lastError string) (Ev
}
return e, nil
}

type QueueStats struct {
ByStatus map[string]int64 `json:"by_status"`
OldestPendingS *float64 `json:"oldest_pending_seconds,omitempty"`
}

func (s *Store) QueueStats(ctx context.Context) (QueueStats, error) {
rows, err := s.pool.Query(ctx, `
SELECT status, count(*)::bigint
FROM events
GROUP BY status
`)
if err != nil {
return QueueStats{}, fmt.Errorf("queue counts: %w", err)
}
defer rows.Close()

out := QueueStats{ByStatus: map[string]int64{}}
for rows.Next() {
var status string
var n int64
if err := rows.Scan(&status, &n); err != nil {
return QueueStats{}, fmt.Errorf("scan queue count: %w", err)
}
out.ByStatus[status] = n
}
if err := rows.Err(); err != nil {
return QueueStats{}, err
}

var age float64
err = s.pool.QueryRow(ctx, `
SELECT EXTRACT(EPOCH FROM (now() - received_at))::float8
FROM events
WHERE status = 'pending'
ORDER BY received_at ASC
LIMIT 1
`).Scan(&age)
if errors.Is(err, pgx.ErrNoRows) {
return out, nil
}
if err != nil {
return QueueStats{}, fmt.Errorf("oldest pending: %w", err)
}
out.OldestPendingS = &age
return out, nil
}
13 changes: 13 additions & 0 deletions api/internal/httpapi/health.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package httpapi

import (
"context"
"net/http"

"github.com/jackc/pgx/v5/pgxpool"
Expand All @@ -9,6 +10,7 @@ import (
type HealthHandler struct {
Pool *pgxpool.Pool
WorkerSnapshot func() any
QueueSnapshot func(ctx context.Context) (any, error)
}

func (h HealthHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
Expand All @@ -23,5 +25,16 @@ func (h HealthHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if h.WorkerSnapshot != nil {
body["worker"] = h.WorkerSnapshot()
}
if h.QueueSnapshot != nil {
queue, err := h.QueueSnapshot(r.Context())
if err != nil {
WriteJSON(w, http.StatusServiceUnavailable, map[string]any{
"status": "unavailable",
"error": "queue_stats_failed",
})
return
}
body["queue"] = queue
}
WriteJSON(w, http.StatusOK, body)
}
102 changes: 102 additions & 0 deletions api/internal/worker/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"fmt"
"log/slog"
"os"
"sync"
"testing"
"time"

Expand Down Expand Up @@ -66,6 +67,107 @@ func TestProcessSuccess(t *testing.T) {
}
}

func TestConcurrentClaimNoDoubleProcess(t *testing.T) {
pool := testPool(t)
ctx := context.Background()
userID := createUser(t, pool)
eventStore := events.NewStore(pool)
ruleStore := rules.NewStore(pool)
const n = 40
for i := 0; i < n; i++ {
_, _, err := eventStore.CreatePending(ctx, events.CreateInput{
UserID: userID,
IdempotencyKey: fmt.Sprintf("evt_race_%d_%d", time.Now().UnixNano(), i),
Type: "tx_simulated",
Payload: []byte(`{"amount":1}`),
})
if err != nil {
t.Fatal(err)
}
}

w := worker.New(eventStore, ruleStore, slog.Default())
errCh := make(chan error, 2)
var wg sync.WaitGroup
run := func() {
defer wg.Done()
for {
ok, err := w.ProcessOne(ctx)
if err != nil {
errCh <- err
return
}
if !ok {
return
}
}
}
wg.Add(2)
go run()
go run()
wg.Wait()
close(errCh)
for err := range errCh {
if err != nil {
t.Fatal(err)
}
}

processed, err := eventStore.ListForUser(ctx, userID, "processed")
if err != nil {
t.Fatal(err)
}
if len(processed) != n {
t.Fatalf("processed=%d want %d", len(processed), n)
}
for _, status := range []string{"pending", "processing", "failed"} {
left, err := eventStore.ListForUser(ctx, userID, status)
if err != nil {
t.Fatal(err)
}
if len(left) != 0 {
t.Fatalf("leftover status=%s count=%d", status, len(left))
}
}
}

func TestReclaimExpiredProcessingLease(t *testing.T) {
pool := testPool(t)
ctx := context.Background()
userID := createUser(t, pool)
eventStore := events.NewStore(pool)

ev, _, err := eventStore.CreatePending(ctx, events.CreateInput{
UserID: userID,
IdempotencyKey: fmt.Sprintf("evt_lease_%d", time.Now().UnixNano()),
Type: "balance_drop",
Payload: []byte(`{"amount":10}`),
})
if err != nil {
t.Fatal(err)
}
if _, err := eventStore.ClaimByID(ctx, ev.ID); err != nil {
t.Fatal(err)
}

if _, err := pool.Exec(ctx, `
UPDATE events
SET claimed_at = now() - interval '2 minutes',
received_at = timestamptz '2000-01-01'
WHERE id = $1
`, ev.ID); err != nil {
t.Fatal(err)
}

reclaimed, err := eventStore.ClaimNextWithLease(ctx, worker.MaxAttempts, time.Second)
if err != nil {
t.Fatalf("reclaim: %v", err)
}
if reclaimed.ID != ev.ID {
t.Fatalf("reclaimed id=%s want %s", reclaimed.ID, ev.ID)
}
}

func TestProcessFailureIncrementsAttempts(t *testing.T) {
pool := testPool(t)
ctx := context.Background()
Expand Down
9 changes: 9 additions & 0 deletions api/migrations/00003_claimed_at.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
-- +goose Up
ALTER TABLE events
ADD COLUMN claimed_at timestamptz;

CREATE INDEX events_claim_queue_idx ON events (status, claimed_at, received_at);

-- +goose Down
DROP INDEX IF EXISTS events_claim_queue_idx;
ALTER TABLE events DROP COLUMN IF EXISTS claimed_at;
Loading
Loading