From 5b83fe9ca48b82417434354d5c905a99767bc3c4 Mon Sep 17 00:00:00 2001 From: Bretton Date: Thu, 13 Aug 2026 06:41:03 -0700 Subject: [PATCH] =?UTF-8?q?wip(task15):=20cycle=20J=20=E2=80=94=20config?= =?UTF-8?q?=20kill-switches,=20main=20wiring=20(noop=E2=86=92real=20enqueu?= =?UTF-8?q?er=20+=20delivery=20workers=20gated=20on=20OUTBOUND=5FWORKERS>0?= =?UTF-8?q?),=20admin=20queue=20surface,=20metrics?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- cmd/tidepool/main.go | 86 ++++++++++++++++++++++-- internal/config/config.go | 67 +++++++++++++++++++ internal/config/config_test.go | 69 ++++++++++++++++++++ internal/ingest/follow.go | 94 ++++++++++++++++++++++++++- internal/outbound/metrics.go | 14 ++++ internal/outbound/switches.go | 43 ++++++++++++ internal/outbound/switches_test.go | 68 +++++++++++++++++++ internal/outbound/worker.go | 7 +- internal/outbound/worker_test.go | 18 ++--- internal/store/interfaces.go | 11 ++++ internal/store/outbound_deliveries.go | 46 +++++++++++++ 11 files changed, 508 insertions(+), 15 deletions(-) create mode 100644 internal/outbound/metrics.go create mode 100644 internal/outbound/switches.go create mode 100644 internal/outbound/switches_test.go diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index a478931..b96bf9a 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -13,6 +13,7 @@ import ( "net/url" "os" "os/signal" + "sync" "syscall" "time" @@ -26,6 +27,7 @@ import ( "tidepool/internal/identity" "tidepool/internal/ingest" "tidepool/internal/materialize" + "tidepool/internal/outbound" "tidepool/internal/personas" "tidepool/internal/prune" "tidepool/internal/repo" @@ -39,6 +41,13 @@ const ( writeTimeout = 30 * time.Second idleTimeout = 2 * time.Minute shutdownTimeout = 15 * time.Second + + // outboundInboxTTL memoizes a community's resolved delivery inbox; a + // rotation is caught by the worker's cache-bypassing re-resolve on a 4xx. + outboundInboxTTL = time.Hour + // outboundWorkerIdle is how long a delivery worker sleeps when the queue is + // empty before polling ClaimNext again. + outboundWorkerIdle = time.Second ) func main() { @@ -410,6 +419,7 @@ func run(logger *slog.Logger) error { Backfill: backfill, Repos: repoManager, Sweeper: handler, + Deliveries: store.NewOutboundDeliveries(database), Logger: logger, }) if err != nil { @@ -470,7 +480,7 @@ func run(logger *slog.Logger) error { // wired end to end must not start accumulating it. var consumerDone <-chan struct{} if cfg.ConsumerEnabled { - consumerDone, err = startConsumer(ctx, cfg, database, repoManager, personasService, logger) + consumerDone, err = startConsumer(ctx, cfg, database, repoManager, personasService, apClient, personasService, logger) if err != nil { return err } @@ -587,6 +597,8 @@ func startConsumer( database *sql.DB, repoManager *repo.Manager, minter consume.ActorMinter, + apClient *ap.Client, + signers outbound.SignerProvider, logger *slog.Logger, ) (<-chan struct{}, error) { // The most SSRF-exposed egress in the bridge: the well-known host comes @@ -605,11 +617,54 @@ func startConsumer( return nil, fmt.Errorf("consumer: handle resolver: %w", err) } + // The outbound delivery pipe (task 15). The noop enqueuer is the default — + // the consumer still writes durable outbound state, but nothing federates + // — and it is swapped for the real enqueuer ONLY when OUTBOUND_WORKERS>0. + // This is the gate that keeps a not-yet-wired deployment (and the e2e, + // which runs with workers=0) from sending anything. + inboxes := outbound.NewInboxResolver(apClient, outboundInboxTTL) + var enqueuer consume.OutboundEnqueuer = consume.NewNoopEnqueuer(logger) + var worker *outbound.Worker + if cfg.OutboundWorkers > 0 { + realEnqueuer, err := outbound.NewEnqueuer(outbound.EnqueuerOptions{ + DB: database, + Translator: outbound.NewTranslator(cfg.APUserOrigin), + Inboxes: inboxes, + Actors: store.NewAPActors(database), + UserOrigin: cfg.APUserOrigin, + Logger: logger, + }) + if err != nil { + return nil, fmt.Errorf("consumer: outbound enqueuer: %w", err) + } + enqueuer = realEnqueuer + + worker, err = outbound.NewWorker(outbound.WorkerOptions{ + DB: database, + Actors: store.NewAPActors(database), + Prefs: store.NewFederationPrefs(database), + Signers: signers, + Inboxes: inboxes, + Sender: apClient, + Switches: outbound.ConfigSwitches{ + Disabled: cfg.OutboundDisabled, + Dry: cfg.OutboundDryRun, + DisabledHosts: cfg.OutboundDisabledHosts, + DisabledCommunities: cfg.OutboundDisabledCommunities, + DisabledActors: cfg.OutboundDisabledActors, + }, + Logger: logger, + }) + if err != nil { + return nil, fmt.Errorf("consumer: outbound worker: %w", err) + } + } + dispatcher, err := consume.NewDispatcher(consume.Options{ DB: database, Actors: minter, Resolver: resolver, - Enqueuer: consume.NewNoopEnqueuer(logger), + Enqueuer: enqueuer, // Reads committed records so a subject's community resolves for // mappings written before migration 016 filled community_did. Records: repoManager, @@ -645,13 +700,36 @@ func startConsumer( // events quietly stop arriving. consume.PublishMetrics(context.Background(), connector, state) - done := make(chan struct{}) + // The connector loop and every delivery worker share one WaitGroup, so the + // returned done channel closes only after ALL of them have drained on ctx + // cancellation — shutdown joins delivery the same way it joins the consumer. + var wg sync.WaitGroup + wg.Add(1) go func() { - defer close(done) + defer wg.Done() if err := connector.Start(ctx); err != nil { logger.Error("jetstream consumer stopped", "error", err) } }() + if worker != nil { + for i := 0; i < cfg.OutboundWorkers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + if err := worker.Run(ctx, outboundWorkerIdle); err != nil && !errors.Is(err, context.Canceled) { + logger.Error("outbound worker stopped", "error", err) + } + }() + } + logger.Info("outbound delivery workers started", + "count", cfg.OutboundWorkers, "dry_run", cfg.OutboundDryRun, "global_disabled", cfg.OutboundDisabled) + } + + done := make(chan struct{}) + go func() { + wg.Wait() + close(done) + }() logger.Info("jetstream consumer started", "url", subscribeURL) return done, nil } diff --git a/internal/config/config.go b/internal/config/config.go index d0f9c39..c4f82f1 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -126,6 +126,25 @@ type Config struct { // nowhere to dial would come up "healthy" and consume nothing, which is // the failure mode cursors and lag metrics exist to make impossible. JetstreamURL string + // OutboundWorkers is how many delivery workers to run (OUTBOUND_WORKERS, + // default 0 = OFF). Delivery starts ONLY when this is >0 AND + // ConsumerEnabled: until a deployment is wired end to end, the consumer + // still records outbound state but the noop enqueuer federates nothing. + OutboundWorkers int + // OutboundDryRun makes workers translate + log but POST nothing, leaving + // deliveries pending (OUTBOUND_DRY_RUN, default false). + OutboundDryRun bool + // OutboundDisabled is the global kill switch: every delivery is PARKED + // (stays pending, resumes when cleared), nothing federates + // (OUTBOUND_DISABLED, default false). + OutboundDisabled bool + // OutboundDisabledHosts / Communities / Actors are the scoped kill + // switches: any delivery whose inbox host, community AP id, or actor DID is + // in the set is parked (comma-separated OUTBOUND_DISABLED_HOSTS / + // OUTBOUND_DISABLED_COMMUNITIES / OUTBOUND_DISABLED_ACTORS). Empty = allow. + OutboundDisabledHosts map[string]struct{} + OutboundDisabledCommunities map[string]struct{} + OutboundDisabledActors map[string]struct{} // AdminToken is the bearer token protecting the /admin API (community // subscribe/unsubscribe/backfill). ADMIN_TOKEN; required in production, // dev default is a fixed, publicly known value. @@ -448,6 +467,25 @@ func Load(logger *slog.Logger) (*Config, error) { return nil, fmt.Errorf("config: JETSTREAM_URL is required when CONSUMER_ENABLED is set") } + // Outbound delivery (task 15), default OFF: workers start only when + // OUTBOUND_WORKERS>0 AND the consumer is on, so a not-yet-wired deployment + // keeps the noop enqueuer and federates nothing. + cfg.OutboundWorkers, err = intVarNonNegative(logger, "OUTBOUND_WORKERS", 0) + if err != nil { + return nil, err + } + cfg.OutboundDryRun, err = boolVarDefault(logger, "OUTBOUND_DRY_RUN", false) + if err != nil { + return nil, err + } + cfg.OutboundDisabled, err = boolVarDefault(logger, "OUTBOUND_DISABLED", false) + if err != nil { + return nil, err + } + cfg.OutboundDisabledHosts = parseSet(os.Getenv("OUTBOUND_DISABLED_HOSTS")) + cfg.OutboundDisabledCommunities = parseSet(os.Getenv("OUTBOUND_DISABLED_COMMUNITIES")) + cfg.OutboundDisabledActors = parseSet(os.Getenv("OUTBOUND_DISABLED_ACTORS")) + // Retention knobs for the task-11 pruners: same semantics as // FIREHOSE_RETENTION (real defaults everywhere, must be positive). cfg.TombstoneRetention, err = durationVar(logger, "TOMBSTONE_RETENTION", 720*time.Hour) @@ -600,6 +638,35 @@ func intVar(logger *slog.Logger, name string, fallback int) (int, error) { return parsed, nil } +// intVarNonNegative is intVar for a knob whose OFF value is 0: it accepts 0 (and +// any positive int), unlike intVar which treats 0 as invalid. Used for +// OUTBOUND_WORKERS, where 0 means "no delivery workers". +func intVarNonNegative(logger *slog.Logger, name string, fallback int) (int, error) { + raw := strings.TrimSpace(os.Getenv(name)) + if raw == "" { + logger.Info(name+" not set, using default", "value", fallback) + return fallback, nil + } + parsed, err := strconv.Atoi(raw) + if err != nil || parsed < 0 { + return 0, fmt.Errorf("config: %s must be a non-negative integer, got %q", name, raw) + } + return parsed, nil +} + +// parseSet splits a comma-separated environment value into a set, dropping +// blanks. An empty or unset value yields an empty (but non-nil) set, which every +// membership test reads as "nothing disabled". +func parseSet(raw string) map[string]struct{} { + set := make(map[string]struct{}) + for _, item := range strings.Split(raw, ",") { + if item = strings.TrimSpace(item); item != "" { + set[item] = struct{}{} + } + } + return set +} + // boolVar reports whether an environment variable is set to a truthy value // ("1", "true", "yes", case-insensitive). func boolVar(name string) bool { diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 29204d2..4006f63 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -23,6 +23,9 @@ func clearConfigEnv(t *testing.T) { "ALLOW_PRIVATE_FETCH", "ALLOW_DEV_REQUEST_CRAWL", "RELAY_HOSTS", "AP_USER_ORIGIN", "AP_HOST_FALLTHROUGH_DEV", "CONSUMER_ENABLED", "JETSTREAM_URL", + "OUTBOUND_WORKERS", "OUTBOUND_DRY_RUN", "OUTBOUND_DISABLED", + "OUTBOUND_DISABLED_HOSTS", "OUTBOUND_DISABLED_COMMUNITIES", + "OUTBOUND_DISABLED_ACTORS", } { t.Setenv(name, "") } @@ -524,3 +527,69 @@ func TestLoad_ConsumerEnabledRejectsATypo(t *testing.T) { "false, because a flag disabled by a typo is invisible") assert.Contains(t, err.Error(), "CONSUMER_ENABLED") } + +func TestLoad_OutboundDefaultsOff(t *testing.T) { + clearConfigEnv(t) + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + + assert.Equal(t, 0, cfg.OutboundWorkers, + "OUTBOUND_WORKERS defaults to 0: delivery is OFF until a deployment is wired end to end") + assert.False(t, cfg.OutboundDryRun, "OUTBOUND_DRY_RUN defaults to false") + assert.False(t, cfg.OutboundDisabled, "OUTBOUND_DISABLED defaults to false") + assert.Empty(t, cfg.OutboundDisabledHosts, "no hosts disabled by default") + assert.Empty(t, cfg.OutboundDisabledCommunities) + assert.Empty(t, cfg.OutboundDisabledActors) +} + +func TestLoad_OutboundWorkersParses(t *testing.T) { + clearConfigEnv(t) + t.Setenv("OUTBOUND_WORKERS", "4") + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + assert.Equal(t, 4, cfg.OutboundWorkers) +} + +func TestLoad_OutboundWorkersRejectsNegative(t *testing.T) { + clearConfigEnv(t) + t.Setenv("OUTBOUND_WORKERS", "-1") + + _, err := Load(discardLogger()) + require.Error(t, err, "a negative worker count is a config error, not silently 0") + assert.Contains(t, err.Error(), "OUTBOUND_WORKERS") +} + +func TestLoad_OutboundDryRunParses(t *testing.T) { + clearConfigEnv(t) + t.Setenv("OUTBOUND_DRY_RUN", "true") + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + assert.True(t, cfg.OutboundDryRun) +} + +func TestLoad_OutboundDisableListsParse(t *testing.T) { + clearConfigEnv(t) + t.Setenv("OUTBOUND_DISABLED", "1") + t.Setenv("OUTBOUND_DISABLED_HOSTS", "lemmy.world, sh.itjust.works") + t.Setenv("OUTBOUND_DISABLED_COMMUNITIES", "https://lemmy.world/c/tech") + t.Setenv("OUTBOUND_DISABLED_ACTORS", "did:plc:aaa,did:plc:bbb , ") + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + + assert.True(t, cfg.OutboundDisabled, "the global kill switch parses") + + assert.Contains(t, cfg.OutboundDisabledHosts, "lemmy.world") + assert.Contains(t, cfg.OutboundDisabledHosts, "sh.itjust.works", + "comma-separated hosts are trimmed and split into a set") + assert.Len(t, cfg.OutboundDisabledHosts, 2) + + assert.Contains(t, cfg.OutboundDisabledCommunities, "https://lemmy.world/c/tech") + + assert.Contains(t, cfg.OutboundDisabledActors, "did:plc:aaa") + assert.Contains(t, cfg.OutboundDisabledActors, "did:plc:bbb") + assert.Len(t, cfg.OutboundDisabledActors, 2, "blank entries are dropped, not stored as empty keys") +} diff --git a/internal/ingest/follow.go b/internal/ingest/follow.go index 069c2e4..c3b5d14 100644 --- a/internal/ingest/follow.go +++ b/internal/ingest/follow.go @@ -51,7 +51,12 @@ type AdminOptions struct { // Sweeper serves POST /admin/objects/sweep-deleted (optional; the // endpoint answers 501 when nil). See sweep.go. Sweeper DeleteSweeper - Logger *slog.Logger + // Deliveries serves the task-15 outbound queue endpoints (optional; those + // endpoints answer 501 when nil): GET /admin/outbound inspects the queue, + // POST /admin/outbound/redrive resets poisoned rows, POST + // /admin/outbound/cancel parks an actor's or community's pending work. + Deliveries store.OutboundDeliveries + Logger *slog.Logger } // Admin is the operator API driving the community subscription lifecycle, @@ -76,6 +81,7 @@ type Admin struct { backfill Backfiller repos RepoReemitter sweeper DeleteSweeper + deliveries store.OutboundDeliveries logger *slog.Logger // reconciler serves POST /admin/communities/reconcile; nil (the // endpoint answers 501) unless a follow list is configured. Set once @@ -118,6 +124,7 @@ func NewAdmin(opts AdminOptions) (*Admin, error) { backfill: opts.Backfill, repos: opts.Repos, sweeper: opts.Sweeper, + deliveries: opts.Deliveries, logger: logger, }, nil } @@ -138,10 +145,95 @@ func (a *Admin) Routes(r chi.Router) { r.Post("/communities/reconcile", a.handleReconcile) r.Post("/reemit", a.handleReemit) r.Post("/objects/sweep-deleted", a.handleSweepDeleted) + r.Get("/outbound", a.handleOutboundInspect) + r.Post("/outbound/redrive", a.handleOutboundRedrive) + r.Post("/outbound/cancel", a.handleOutboundCancel) r.Method(http.MethodGet, "/metrics", http.HandlerFunc(scopedMetrics)) }) } +// handleOutboundInspect reports the delivery queue depth by state — the +// operator's window on pending/poisoned/cancelled backlog (task 15). +func (a *Admin) handleOutboundInspect(w http.ResponseWriter, r *http.Request) { + if a.deliveries == nil { + http.Error(w, "outbound delivery is not configured", http.StatusNotImplemented) + return + } + counts, err := a.deliveries.CountsByState(r.Context()) + if err != nil { + a.logger.Error("outbound inspect failed", "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + byState := map[string]int{} + for state, n := range counts { + byState[string(state)] = n + } + w.Header().Set("Content-Type", "application/json; charset=utf-8") + _ = json.NewEncoder(w).Encode(map[string]any{"by_state": byState}) +} + +// outboundMutateRequest is the shared body for redrive/cancel: any subset of +// the optional filters. +type outboundMutateRequest struct { + Activity string `json:"activity"` + Community string `json:"community"` + Actor string `json:"actor"` +} + +// handleOutboundRedrive resets poisoned deliveries back to pending, optionally +// scoped to one activity or community. +func (a *Admin) handleOutboundRedrive(w http.ResponseWriter, r *http.Request) { + if a.deliveries == nil { + http.Error(w, "outbound delivery is not configured", http.StatusNotImplemented) + return + } + var req outboundMutateRequest + if r.Body != nil { + _ = json.NewDecoder(r.Body).Decode(&req) + } + redriven, err := a.deliveries.RedrivePoisoned(r.Context(), req.Activity, req.Community) + if err != nil { + a.logger.Error("outbound redrive failed", "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json; charset=utf-8") + _ = json.NewEncoder(w).Encode(map[string]any{"redriven": redriven}) +} + +// handleOutboundCancel parks an actor's or a community's PENDING deliveries as +// cancelled (consent withdrawal / community removal). Exactly one of actor or +// community must be given. +func (a *Admin) handleOutboundCancel(w http.ResponseWriter, r *http.Request) { + if a.deliveries == nil { + http.Error(w, "outbound delivery is not configured", http.StatusNotImplemented) + return + } + var req outboundMutateRequest + if r.Body != nil { + _ = json.NewDecoder(r.Body).Decode(&req) + } + if (req.Actor == "") == (req.Community == "") { + http.Error(w, `body must set exactly one of {"actor":"did:..."} or {"community":"https://..."}`, http.StatusBadRequest) + return + } + var cancelled int64 + var err error + if req.Actor != "" { + cancelled, err = a.deliveries.CancelForActor(r.Context(), req.Actor) + } else { + cancelled, err = a.deliveries.CancelForCommunity(r.Context(), req.Community) + } + if err != nil { + a.logger.Error("outbound cancel failed", "error", err) + http.Error(w, "internal error", http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json; charset=utf-8") + _ = json.NewEncoder(w).Encode(map[string]any{"cancelled": cancelled}) +} + // scopedMetrics writes the JSON expvar map filtered to tidepool's own // counters — the same wire format as expvar.Handler(), minus Go's global // cmdline/memstats. diff --git a/internal/outbound/metrics.go b/internal/outbound/metrics.go new file mode 100644 index 0000000..a97ea4a --- /dev/null +++ b/internal/outbound/metrics.go @@ -0,0 +1,14 @@ +package outbound + +import "expvar" + +// Delivery outcome counters, published under the tidepool_ prefix so they +// surface at /admin/metrics (ingest.Admin's scopedMetrics filters to that +// prefix). The worker bumps exactly one per terminal outcome, plus parked for +// each kill-switch/dry-run/causal-wait deferral. +var ( + metricDelivered = expvar.NewInt("tidepool_outbound_delivered") + metricPoisoned = expvar.NewInt("tidepool_outbound_poisoned") + metricCancelled = expvar.NewInt("tidepool_outbound_cancelled") + metricParked = expvar.NewInt("tidepool_outbound_parked") +) diff --git a/internal/outbound/switches.go b/internal/outbound/switches.go new file mode 100644 index 0000000..c6e5663 --- /dev/null +++ b/internal/outbound/switches.go @@ -0,0 +1,43 @@ +package outbound + +// ConfigSwitches is the config-backed Switches: the operator kill switches +// (decision 19) resolved from static configuration. It is the adapter main +// builds from config values and hands the Worker. +// +// A blocked delivery is PARKED by the worker (stays pending, resumes when the +// switch clears) — never poisoned or cancelled. The zero value (all fields +// empty) allows everything, so a deployment that sets no switches behaves +// exactly like AllowAll. +type ConfigSwitches struct { + // Disabled is the global kill switch: when true, every delivery is blocked. + Disabled bool + // Dry, when true, makes the worker translate + log but POST nothing. + Dry bool + // DisabledHosts / Communities / Actors block a delivery whose inbox host, + // community AP id, or actor DID is in the set. Nil sets block nothing. + DisabledHosts map[string]struct{} + DisabledCommunities map[string]struct{} + DisabledActors map[string]struct{} +} + +// OutboundAllowed reports whether a delivery in this scope may be sent: false +// when the global switch is engaged or the scope's host, community, or actor is +// in a disabled set. +func (s ConfigSwitches) OutboundAllowed(scope DeliveryScope) bool { + if s.Disabled { + return false + } + if _, blocked := s.DisabledHosts[scope.InboxHost]; blocked { + return false + } + if _, blocked := s.DisabledCommunities[scope.CommunityAPID]; blocked { + return false + } + if _, blocked := s.DisabledActors[scope.ActorDID]; blocked { + return false + } + return true +} + +// DryRun reports whether to translate + log without POSTing. +func (s ConfigSwitches) DryRun() bool { return s.Dry } diff --git a/internal/outbound/switches_test.go b/internal/outbound/switches_test.go new file mode 100644 index 0000000..10308f1 --- /dev/null +++ b/internal/outbound/switches_test.go @@ -0,0 +1,68 @@ +package outbound + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +// Cycle J: the config-backed kill switches. A blocked scope is what the worker +// parks on; the zero value must behave exactly like AllowAll. + +func set(items ...string) map[string]struct{} { + s := make(map[string]struct{}, len(items)) + for _, item := range items { + s[item] = struct{}{} + } + return s +} + +var switchScope = DeliveryScope{ + ActorDID: "did:plc:actor", + CommunityAPID: "https://lemmy.world/c/tech", + InboxHost: "lemmy.world", +} + +func TestConfigSwitches_EmptyIsAllowAll(t *testing.T) { + var s ConfigSwitches + assert.True(t, s.OutboundAllowed(switchScope), + "the zero value blocks nothing — a deployment with no switches federates normally") + assert.False(t, s.DryRun()) +} + +func TestConfigSwitches_GlobalDisableBlocksEverything(t *testing.T) { + s := ConfigSwitches{Disabled: true} + assert.False(t, s.OutboundAllowed(switchScope), "the global kill switch blocks all scopes") + assert.False(t, s.OutboundAllowed(DeliveryScope{ActorDID: "other", CommunityAPID: "x", InboxHost: "y"}), + "including scopes in no disabled set") +} + +func TestConfigSwitches_ScopedDisableBlocksOnlyThatScope(t *testing.T) { + cases := []struct { + name string + s ConfigSwitches + }{ + {"host", ConfigSwitches{DisabledHosts: set("lemmy.world")}}, + {"community", ConfigSwitches{DisabledCommunities: set("https://lemmy.world/c/tech")}}, + {"actor", ConfigSwitches{DisabledActors: set("did:plc:actor")}}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + assert.False(t, tc.s.OutboundAllowed(switchScope), + "a delivery whose %s is in the disabled set is blocked", tc.name) + // A scope that shares none of the disabled dimensions is allowed. + assert.True(t, tc.s.OutboundAllowed(DeliveryScope{ + ActorDID: "did:plc:other", + CommunityAPID: "https://lemmy.world/c/other", + InboxHost: "sh.itjust.works", + }), "an unrelated scope still delivers") + }) + } +} + +func TestConfigSwitches_DryRun(t *testing.T) { + s := ConfigSwitches{Dry: true} + assert.True(t, s.DryRun(), "dry-run is reported independently of allow/deny") + assert.True(t, s.OutboundAllowed(switchScope), + "dry-run does not block: the worker still claims, translates and logs — it just POSTs nothing") +} diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go index fd5b96e..6b49509 100644 --- a/internal/outbound/worker.go +++ b/internal/outbound/worker.go @@ -244,9 +244,11 @@ func (w *Worker) handle(ctx context.Context, delivery *store.OutboundDelivery) e return err } if blocked { - if _, err := w.deliveries.CancelForActor(ctx, activity.ActorDID); err != nil { + cancelled, err := w.deliveries.CancelForActor(ctx, activity.ActorDID) + if err != nil { return fmt.Errorf("cancel deliveries for %s: %w", activity.ActorDID, err) } + metricCancelled.Add(cancelled) return nil } } @@ -338,6 +340,7 @@ func (w *Worker) deliverSuccess(ctx context.Context, delivery *store.OutboundDel if !applied { return nil // a stale claim: another worker already recorded the outcome } + metricDelivered.Add(1) w.stampAccepted(ctx, activity) return w.voteCallback(ctx, activity) } @@ -413,6 +416,7 @@ func (w *Worker) poison(ctx context.Context, delivery *store.OutboundDelivery, c if err != nil { return fmt.Errorf("poison delivery %s: %w", delivery.ActivityID, err) } + metricPoisoned.Add(1) return nil } @@ -424,6 +428,7 @@ func (w *Worker) park(ctx context.Context, delivery *store.OutboundDelivery, cla if err != nil { return fmt.Errorf("park delivery %s: %w", delivery.ActivityID, err) } + metricParked.Add(1) return nil } diff --git a/internal/outbound/worker_test.go b/internal/outbound/worker_test.go index 7ad8d80..62672d2 100644 --- a/internal/outbound/worker_test.go +++ b/internal/outbound/worker_test.go @@ -107,10 +107,10 @@ func (f fakeSigners) SignerFor(context.Context, string) (*ap.Signer, error) { // fakeSwitches records the scopes it was consulted with. type fakeSwitches struct { - mu sync.Mutex - allow bool - dryRun bool - scopes []DeliveryScope + mu sync.Mutex + allow bool + dryRun bool + scopes []DeliveryScope } func (s *fakeSwitches) OutboundAllowed(scope DeliveryScope) bool { @@ -358,11 +358,11 @@ func TestWorker_ConsentAsymmetry(t *testing.T) { explanation: "a federation opt-out cancels an outward create", }, { - name: "disabled actor STILL delivers a delete", - disable: func(t *testing.T, conn *sql.DB) { seedWorkerActor(t, conn, false, false) }, - kind: "Delete", - wantState: store.DeliveryStateDelivered, - wantPosted: true, + name: "disabled actor STILL delivers a delete", + disable: func(t *testing.T, conn *sql.DB) { seedWorkerActor(t, conn, false, false) }, + kind: "Delete", + wantState: store.DeliveryStateDelivered, + wantPosted: true, explanation: "retraction asymmetry: a Delete goes out even for a disabled actor — it is the " + "only way an opted-out user takes down what is already federated", }, diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index a158f77..5477feb 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -547,4 +547,15 @@ type OutboundDeliveries interface { // (parent_unaccepted). Per-community serialization makes a lower-seq // delivery on the same line a causal ancestor. HasPoisonedPredecessor(ctx context.Context, orderingKey, targetInbox string, seq int64) (bool, error) + + // CountsByState returns the number of deliveries in each state — the + // operator queue-inspect (GET /admin/outbound). + CountsByState(ctx context.Context) (map[DeliveryState]int, error) + + // RedrivePoisoned resets poisoned deliveries back to pending for + // redelivery, clearing the lease and rescheduling now. activityID and + // orderingKey are optional filters (empty = no filter on that column); the + // attempt counter is reset so a redriven delivery gets a fresh budget. + // Returns how many rows were redriven. + RedrivePoisoned(ctx context.Context, activityID, orderingKey string) (int64, error) } diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go index 34a070f..1389962 100644 --- a/internal/store/outbound_deliveries.go +++ b/internal/store/outbound_deliveries.go @@ -278,6 +278,52 @@ func (r *postgresOutboundDeliveries) HasPoisonedPredecessor(ctx context.Context, return exists, nil } +func (r *postgresOutboundDeliveries) CountsByState(ctx context.Context) (map[DeliveryState]int, error) { + rows, err := r.db.QueryContext(ctx, + `SELECT state, COUNT(*) FROM outbound_deliveries GROUP BY state`) + if err != nil { + return nil, fmt.Errorf("count outbound_deliveries by state: %w", err) + } + defer rows.Close() + + counts := make(map[DeliveryState]int) + for rows.Next() { + var state string + var n int + if err := rows.Scan(&state, &n); err != nil { + return nil, fmt.Errorf("scan delivery state count: %w", err) + } + counts[DeliveryState(state)] = n + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate delivery state counts: %w", err) + } + return counts, nil +} + +func (r *postgresOutboundDeliveries) RedrivePoisoned(ctx context.Context, activityID, orderingKey string) (int64, error) { + // Empty filters pass through the NULLIF/COALESCE guard: a blank $2/$3 means + // "any row", so the sweep can target one activity, one community, or all + // poisoned deliveries. Reset attempts + next_attempt_at so a redriven + // delivery gets a fresh budget and is immediately claimable. + result, err := r.db.ExecContext(ctx, ` + UPDATE outbound_deliveries + SET state = 'pending', attempts = 0, claimed_until = NULL, next_attempt_at = now(), + updated_at = now() + WHERE state = 'poisoned' + AND ($1 = '' OR activity_id = $1) + AND ($2 = '' OR ordering_key = $2)`, + activityID, orderingKey) + if err != nil { + return 0, fmt.Errorf("redrive poisoned deliveries: %w", err) + } + affected, err := result.RowsAffected() + if err != nil { + return 0, fmt.Errorf("redrive poisoned deliveries: rows affected: %w", err) + } + return affected, nil +} + func scanOutboundDelivery(row rowScanner) (*OutboundDelivery, error) { var delivery OutboundDelivery var state string -- 2.51.2