diff --git a/PROJECTS/beginner/canary-token-generator/.env.example b/PROJECTS/beginner/canary-token-generator/.env.example index 4f97ac3f..e7aa0c74 100644 --- a/PROJECTS/beginner/canary-token-generator/.env.example +++ b/PROJECTS/beginner/canary-token-generator/.env.example @@ -20,7 +20,7 @@ NGINX_HOST_PORT=22784 PUBLIC_BASE_URL=https://canary.your.domain # Frontend build-time vars (baked into the bundle by Vite) -VITE_APP_TITLE=Canary Token Generator +VITE_APP_TITLE="Canary Token Generator" VITE_API_URL=/api # --------------------------------------------------------------------------- diff --git a/PROJECTS/beginner/canary-token-generator/.gitignore b/PROJECTS/beginner/canary-token-generator/.gitignore index f151d537..a0160c84 100644 --- a/PROJECTS/beginner/canary-token-generator/.gitignore +++ b/PROJECTS/beginner/canary-token-generator/.gitignore @@ -10,8 +10,7 @@ docs/ # Environment files (secrets) # ---------------------------------------------------------------------------- .env -.env.local -.env.*.local +.env.* !.env.example # ---------------------------------------------------------------------------- diff --git a/PROJECTS/beginner/canary-token-generator/backend/cmd/canary/main.go b/PROJECTS/beginner/canary-token-generator/backend/cmd/canary/main.go index e6010ed0..920383a6 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/cmd/canary/main.go +++ b/PROJECTS/beginner/canary-token-generator/backend/cmd/canary/main.go @@ -127,7 +127,14 @@ func run(configPath string) error { shutdownErr := gracefulShutdown(cfg, logger, srv, telemetry, rdb, db) logger.Info("waiting for in-flight notifications") - notifySvc.Wait() + notifyShutdownCtx, notifyShutdownCancel := context.WithTimeout( + context.Background(), + cfg.Server.ShutdownTimeout, + ) + if nErr := notifySvc.Shutdown(notifyShutdownCtx); nErr != nil { + logger.Warn("notify shutdown timed out", "error", nErr) + } + notifyShutdownCancel() wg.Wait() return shutdownErr } diff --git a/PROJECTS/beginner/canary-token-generator/backend/internal/event/integration_test.go b/PROJECTS/beginner/canary-token-generator/backend/internal/event/integration_test.go index 25be1ffc..be1187c1 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/internal/event/integration_test.go +++ b/PROJECTS/beginner/canary-token-generator/backend/internal/event/integration_test.go @@ -128,6 +128,7 @@ func setupIntgStack(t *testing.T) *intgStack { eventRepo, eventSvc, logger, + false, ) r := chi.NewRouter() @@ -412,8 +413,8 @@ func TestIntegration_ManagePageReturnsTokenAndEventsAndSilencedCount(t *testing. require.Equal(t, int64(3), manageResp.Data.EventsTotal, "all 3 events recorded (one sent + two deduped)") - require.Equal(t, int64(2), manageResp.Data.EventsSilencedActive, - "two triggers deduped within 15-min window") + require.Equal(t, int64(1), manageResp.Data.EventsSilencedActive, + "one unique IP silenced (same source for both dedup hits)") require.Len(t, manageResp.Data.Events, 3, "events page payload") require.False(t, manageResp.Data.Page.HasMore, diff --git a/PROJECTS/beginner/canary-token-generator/backend/internal/event/service.go b/PROJECTS/beginner/canary-token-generator/backend/internal/event/service.go index 5bd64765..b07f07b2 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/internal/event/service.go +++ b/PROJECTS/beginner/canary-token-generator/backend/internal/event/service.go @@ -8,7 +8,6 @@ import ( "errors" "fmt" "log/slog" - "strconv" "time" "github.com/redis/go-redis/v9" @@ -17,9 +16,9 @@ import ( ) const ( - dedupKeyPrefix = "dedup:trigger:" - defaultDedupTTL = 15 * time.Minute - dedupScanBatch = 100 + dedupKeyPrefix = "dedup:trigger:" + dedupActivePrefix = "dedup:active:" + defaultDedupTTL = 15 * time.Minute ) type Service struct { @@ -132,6 +131,17 @@ func (s *Service) dedupGate( s.logger.WarnContext(ctx, "dedup incr failed", "error", iErr, "key", key) } + trackKey := dedupActivePrefix + tokenID + if _, sErr := s.rdb.SAdd(ctx, trackKey, sourceIP).Result(); sErr != nil { + s.logger.WarnContext(ctx, "dedup track add", + "error", sErr, "key", trackKey) + } + if _, eErr := s.rdb.Expire( + ctx, trackKey, s.dedupTTL, + ).Result(); eErr != nil { + s.logger.WarnContext(ctx, "dedup track expire", + "error", eErr, "key", trackKey) + } return false } @@ -142,40 +152,14 @@ func (s *Service) CountActiveDedup( if s.rdb == nil { return 0, nil } - pattern := dedupKeyPrefix + tokenID + ":*" - var total int64 - var cursor uint64 - for { - keys, next, err := s.rdb.Scan( - ctx, cursor, pattern, dedupScanBatch, - ).Result() - if err != nil { - return 0, fmt.Errorf("dedup scan: %w", err) - } - for _, key := range keys { - v, gErr := s.rdb.Get(ctx, key).Result() - if errors.Is(gErr, redis.Nil) { - continue - } - if gErr != nil { - s.logger.WarnContext(ctx, "dedup count: get key", - "error", gErr, "key", key) - continue - } - n, pErr := strconv.ParseInt(v, 10, 64) - if pErr != nil { - continue - } - if n > 1 { - total += n - 1 - } - } - if next == 0 { - break - } - cursor = next + n, err := s.rdb.SCard(ctx, dedupActivePrefix+tokenID).Result() + if errors.Is(err, redis.Nil) { + return 0, nil } - return total, nil + if err != nil { + return 0, fmt.Errorf("dedup count: %w", err) + } + return n, nil } func (s *Service) RunRetentionLoop( diff --git a/PROJECTS/beginner/canary-token-generator/backend/internal/event/service_test.go b/PROJECTS/beginner/canary-token-generator/backend/internal/event/service_test.go index 1202edd3..6109e5c2 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/internal/event/service_test.go +++ b/PROJECTS/beginner/canary-token-generator/backend/internal/event/service_test.go @@ -582,8 +582,8 @@ func TestService_CountActiveDedup_CountsSilencedAcrossIPs(t *testing.T) { n, err := svc.CountActiveDedup(context.Background(), testTokenID) require.NoError(t, err) - require.Equal(t, int64(2+4), n, - "key1=3 (silenced 2) + key2=5 (silenced 4)") + require.Equal(t, int64(2), n, + "two distinct IPs were silenced (203.0.113.1 and 203.0.113.2)") } func TestService_CountActiveDedup_IgnoresOtherTokens(t *testing.T) { @@ -616,7 +616,8 @@ func TestService_CountActiveDedup_IgnoresOtherTokens(t *testing.T) { n, err := svc.CountActiveDedup(context.Background(), testTokenID) require.NoError(t, err) - require.Equal(t, int64(2), n, "only this token's keys counted") + require.Equal(t, int64(1), n, + "only this token's silenced IPs counted (one distinct IP)") } func TestService_CountActiveDedup_NilRedisReturnsZero(t *testing.T) { diff --git a/PROJECTS/beginner/canary-token-generator/backend/internal/notify/service.go b/PROJECTS/beginner/canary-token-generator/backend/internal/notify/service.go index 898326bc..3dc887f2 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/internal/notify/service.go +++ b/PROJECTS/beginner/canary-token-generator/backend/internal/notify/service.go @@ -12,14 +12,27 @@ import ( "github.com/CarterPerez-dev/cybersecurity-projects/canary-token-generator/backend/internal/event" ) -const defaultSendTimeout = 30 * time.Second +const ( + defaultSendTimeout = 30 * time.Second + defaultWorkers = 8 + defaultQueueSize = 256 +) type Service struct { senders map[string]Sender status StatusWriter logger *slog.Logger sendTimeout time.Duration - wg sync.WaitGroup + workers int + queue chan dispatchJob + workerWg sync.WaitGroup + jobWg sync.WaitGroup + closeOnce sync.Once +} + +type dispatchJob struct { + info event.NotifyInfo + evt *event.Event } type Option func(*Service) @@ -32,16 +45,40 @@ func WithSendTimeout(d time.Duration) Option { return func(s *Service) { s.sendTimeout = d } } +func WithMaxConcurrent(n int) Option { + return func(s *Service) { + if n > 0 { + s.workers = n + } + } +} + +func WithQueueSize(n int) Option { + return func(s *Service) { + if n > 0 { + s.queue = make(chan dispatchJob, n) + } + } +} + func NewService(status StatusWriter, opts ...Option) *Service { s := &Service{ senders: make(map[string]Sender), status: status, logger: slog.Default(), sendTimeout: defaultSendTimeout, + workers: defaultWorkers, } for _, o := range opts { o(s) } + if s.queue == nil { + s.queue = make(chan dispatchJob, defaultQueueSize) + } + for range s.workers { + s.workerWg.Add(1) + go s.worker() + } return s } @@ -55,15 +92,50 @@ func (s *Service) Register(senders ...Sender) { } func (s *Service) Notify(info event.NotifyInfo, evt *event.Event) { - s.wg.Add(1) - go func() { - defer s.wg.Done() - s.dispatch(info, evt) - }() + s.jobWg.Add(1) + select { + case s.queue <- dispatchJob{info: info, evt: evt}: + default: + s.jobWg.Done() + s.logger.Warn("notify: queue full, dropping", + "event_id", evt.ID, + "token_id", info.TokenID, + "channel", info.AlertChannel, + ) + s.markStatus( + context.Background(), + evt.ID, + event.NotifyFailed, + nil, + ) + } } func (s *Service) Wait() { - s.wg.Wait() + s.jobWg.Wait() +} + +func (s *Service) Shutdown(ctx context.Context) error { + s.closeOnce.Do(func() { close(s.queue) }) + done := make(chan struct{}) + go func() { + s.workerWg.Wait() + close(done) + }() + select { + case <-done: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +func (s *Service) worker() { + defer s.workerWg.Done() + for job := range s.queue { + s.dispatch(job.info, job.evt) + s.jobWg.Done() + } } func (s *Service) dispatch(info event.NotifyInfo, evt *event.Event) { diff --git a/PROJECTS/beginner/canary-token-generator/backend/internal/token/service.go b/PROJECTS/beginner/canary-token-generator/backend/internal/token/service.go index b402d5b8..e0301f6b 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/internal/token/service.go +++ b/PROJECTS/beginner/canary-token-generator/backend/internal/token/service.go @@ -14,6 +14,7 @@ import ( "github.com/go-playground/validator/v10" "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgconn" ) const ( @@ -22,6 +23,9 @@ const ( metadataDestinationURL = "destination_url" metadataIncludeKeys = "include_keys" + + pgUniqueViolationCode = "23505" + maxTokenIDAttempts = 5 ) var ( @@ -35,6 +39,7 @@ var ( "token: no generator registered for this type", ) ErrGenerateFailed = errors.New("token: artifact generation failed") + ErrValidation = errors.New("token: request validation failed") ) var allowedIncludeKeys = map[string]struct{}{ @@ -94,7 +99,7 @@ func (s *Service) Create( ) (*Token, Artifact, error) { if err := s.validate.Struct(req); err != nil { return nil, Artifact{}, fmt.Errorf( - "validate request: %w", err, + "%w: %w", ErrValidation, err, ) } if err := validateTypeMetadata(req.Type, req.Metadata); err != nil { @@ -108,16 +113,8 @@ func (s *Service) Create( ) } - id, err := generateTokenID() - if err != nil { - return nil, Artifact{}, fmt.Errorf( - "generate id: %w", err, - ) - } manageID := uuid.NewString() - tok := &Token{ - ID: id, ManageID: manageID, Type: req.Type, Memo: req.Memo, @@ -132,19 +129,45 @@ func (s *Service) Create( Metadata: normalizeMetadata(req.Metadata), } - art, err := gen.Generate(ctx, tok, s.baseURL) - if err != nil { - return nil, Artifact{}, fmt.Errorf( - "%w: %w", ErrGenerateFailed, err, - ) - } + var ( + art Artifact + insertErr error + ) + for attempt := range maxTokenIDAttempts { + id, err := generateTokenID() + if err != nil { + return nil, Artifact{}, fmt.Errorf( + "generate id (attempt %d): %w", attempt, err, + ) + } + tok.ID = id - if err := s.repo.Insert(ctx, tok); err != nil { - return nil, Artifact{}, fmt.Errorf( - "persist token: %w", err, - ) + art, err = gen.Generate(ctx, tok, s.baseURL) + if err != nil { + return nil, Artifact{}, fmt.Errorf( + "%w: %w", ErrGenerateFailed, err, + ) + } + + insertErr = s.repo.Insert(ctx, tok) + if insertErr == nil { + return tok, art, nil + } + if !isUniqueViolation(insertErr) { + return nil, Artifact{}, fmt.Errorf( + "persist token: %w", insertErr, + ) + } } - return tok, art, nil + return nil, Artifact{}, fmt.Errorf( + "persist token after %d id-collision retries: %w", + maxTokenIDAttempts, insertErr, + ) +} + +func isUniqueViolation(err error) bool { + var pgErr *pgconn.PgError + return errors.As(err, &pgErr) && pgErr.Code == pgUniqueViolationCode } func (s *Service) GetByID( diff --git a/PROJECTS/beginner/canary-token-generator/backend/internal/token/service_test.go b/PROJECTS/beginner/canary-token-generator/backend/internal/token/service_test.go index 5dc4258e..34dd5dbd 100644 --- a/PROJECTS/beginner/canary-token-generator/backend/internal/token/service_test.go +++ b/PROJECTS/beginner/canary-token-generator/backend/internal/token/service_test.go @@ -12,6 +12,7 @@ import ( "sync/atomic" "testing" + "github.com/jackc/pgx/v5/pgconn" "github.com/stretchr/testify/require" "github.com/CarterPerez-dev/cybersecurity-projects/canary-token-generator/backend/internal/event" @@ -358,6 +359,119 @@ func TestService_Create_DistinctIDsAcrossCalls(t *testing.T) { ) } +type collisionRepo struct { + mu sync.Mutex + collisions int + remaining int + inserted []*token.Token + byID map[string]*token.Token +} + +func newCollisionRepo(collisions int) *collisionRepo { + return &collisionRepo{ + remaining: collisions, + byID: map[string]*token.Token{}, + } +} + +func (r *collisionRepo) Insert(_ context.Context, t *token.Token) error { + r.mu.Lock() + defer r.mu.Unlock() + if r.remaining > 0 { + r.remaining-- + r.collisions++ + return &pgconn.PgError{Code: "23505"} + } + r.inserted = append(r.inserted, t) + r.byID[t.ID] = t + return nil +} + +func (r *collisionRepo) GetByID( + _ context.Context, id string, +) (*token.Token, error) { + r.mu.Lock() + defer r.mu.Unlock() + t, ok := r.byID[id] + if !ok { + return nil, token.ErrNotFound + } + return t, nil +} + +func (r *collisionRepo) GetByManageID( + _ context.Context, _ string, +) (*token.Token, error) { + return nil, token.ErrNotFound +} + +func (r *collisionRepo) IncrementTriggerCount( + _ context.Context, _ string, +) error { + return nil +} + +func (r *collisionRepo) DeleteByManageID( + _ context.Context, _ string, +) error { + return token.ErrNotFound +} + +func TestService_Create_RetriesOnTokenIDCollision(t *testing.T) { + repo := newCollisionRepo(2) + gen := &fakeGenerator{ + tokenType: token.TypeWebbug, + artifact: generators.Artifact{ + Kind: generators.KindURL, + URL: "https://canary.example.com/c/x", + }, + } + svc := token.NewService( + repo, + token.MapRegistry{token.TypeWebbug: gen}, + token.ServiceConfig{BaseURL: "https://canary.example.com"}, + ) + + tok, _, err := svc.Create(context.Background(), token.CreateRequest{ + Type: token.TypeWebbug, + Memo: "x", + AlertChannel: token.ChannelWebhook, + WebhookURL: "https://example.com/h", + }, "fp", "ip") + require.NoError(t, err) + require.NotNil(t, tok) + require.Equal(t, 2, repo.collisions, + "repo must have rejected exactly two prior IDs") + require.Len(t, repo.inserted, 1, "exactly one token persisted") + require.Equal(t, int32(3), gen.calls.Load(), + "generator called once per attempt (regenerates artifact per id)") +} + +func TestService_Create_GivesUpAfterMaxCollisions(t *testing.T) { + repo := newCollisionRepo(10) + gen := &fakeGenerator{ + tokenType: token.TypeWebbug, + artifact: generators.Artifact{ + Kind: generators.KindURL, + URL: "https://canary.example.com/c/x", + }, + } + svc := token.NewService( + repo, + token.MapRegistry{token.TypeWebbug: gen}, + token.ServiceConfig{BaseURL: "https://canary.example.com"}, + ) + + _, _, err := svc.Create(context.Background(), token.CreateRequest{ + Type: token.TypeWebbug, + Memo: "x", + AlertChannel: token.ChannelWebhook, + WebhookURL: "https://example.com/h", + }, "fp", "ip") + require.Error(t, err) + require.Contains(t, err.Error(), "id-collision retries") +} + func TestService_GetByID_NotFoundReturnsNilNil(t *testing.T) { repo := newFakeRepo() svc := token.NewService(repo, token.MapRegistry{}, diff --git a/PROJECTS/beginner/canary-token-generator/frontend/.npmrc b/PROJECTS/beginner/canary-token-generator/frontend/.npmrc new file mode 100644 index 00000000..1f80f77c --- /dev/null +++ b/PROJECTS/beginner/canary-token-generator/frontend/.npmrc @@ -0,0 +1,3 @@ +# ©AngelaMos | 2026 +# .npmrc +strict-dep-builds=false diff --git a/PROJECTS/beginner/canary-token-generator/frontend/index.html b/PROJECTS/beginner/canary-token-generator/frontend/index.html index 7f69fb26..db375a7c 100644 --- a/PROJECTS/beginner/canary-token-generator/frontend/index.html +++ b/PROJECTS/beginner/canary-token-generator/frontend/index.html @@ -9,20 +9,20 @@